Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 8 additions & 1 deletion .env.example
Original file line number Diff line number Diff line change
@@ -1,5 +1,12 @@
PORT=3010
DATABASE_URL=REPLACE_WITH_SETTLEMENTS_DATABASE_URL
WORKBENCH_DATABASE_URL=REPLACE_WITH_WORKBENCH_RUNTIME_DATABASE_URL
WORKBENCH_MIGRATION_DATABASE_URL=REPLACE_WITH_WORKBENCH_MIGRATION_DATABASE_URL
WORKBENCH_DATABASE_ROLE_PASSWORD=REPLACE_WITH_WORKBENCH_DATABASE_ROLE_PASSWORD
SETTLEMENTS_DATABASE_URL=REPLACE_WITH_SETTLEMENTS_RUNTIME_DATABASE_URL
LEGACY_CONNECTORS_DATABASE_URL=REPLACE_WITH_LEGACY_CONNECTORS_RUNTIME_DATABASE_URL
WORKBENCH_JOB_RECEIPT_STORE=postgres
WORKBENCH_JOB_RECEIPT_RETENTION_DAYS=365
WORKBENCH_METRICS_TOKEN=REPLACE_WITH_DEDICATED_SCRAPE_TOKEN
REDIS_URL=REPLACE_WITH_REDIS_URL
LEDGER_CORE_URL=REPLACE_WITH_LEDGER_CORE_URL
OIDC_ISSUER=REPLACE_WITH_HTTPS_IDENTITY_ISSUER
Expand Down
4 changes: 3 additions & 1 deletion .github/workflows/guardian.yml
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,9 @@ jobs:
else
git clone --depth 1 https://github.com/mavulahq/settlements ../settlements
fi
if git ls-remote --exit-code --heads https://github.com/mavulahq/legacy-connectors "$LEGACY_CONNECTORS_REF"; then
if git ls-remote --exit-code --heads https://github.com/mavulahq/legacy-connectors "$branch"; then
git clone --depth 1 --branch "$branch" https://github.com/mavulahq/legacy-connectors ../legacy-connectors
elif git ls-remote --exit-code --heads https://github.com/mavulahq/legacy-connectors "$LEGACY_CONNECTORS_REF"; then
git clone --depth 1 --branch "$LEGACY_CONNECTORS_REF" https://github.com/mavulahq/legacy-connectors ../legacy-connectors
else
git clone --depth 1 https://github.com/mavulahq/legacy-connectors ../legacy-connectors
Expand Down
59 changes: 47 additions & 12 deletions contracts/openapi/workbench.public.v1.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ info:
name: GNU Affero General Public License v3.0 only
identifier: AGPL-3.0-only
servers:
- url: https://workbench.mavula.dev
- url: https://workbench.mavula.io
description: Workbench public endpoint
security: [{ bearerAuth: [] }]
tags:
Expand All @@ -27,6 +27,9 @@ paths:
description: Submit tenant-scoped payment work. The tenant is always derived from the access token.
tags: [Jobs]
x-mavula-permissions: [workbench.jobs.write]
parameters:
- { $ref: '#/components/parameters/IdempotencyKey' }
- { $ref: '#/components/parameters/CorrelationId' }
requestBody:
required: true
content:
Expand All @@ -35,9 +38,7 @@ paths:
examples:
payment_capture:
value:
queue: payments
type: PAYMENT_CAPTURE
max_attempts: 3
payload:
idempotency_key: payment_capture_20260715_001
correlation_id: checkout_20260715_001
Expand All @@ -52,6 +53,8 @@ paths:
'400': { $ref: '#/components/responses/BadRequest' }
'401': { $ref: '#/components/responses/Unauthorized' }
'403': { $ref: '#/components/responses/Forbidden' }
'409': { $ref: '#/components/responses/Conflict' }
'503': { $ref: '#/components/responses/ServiceUnavailable' }
/api/jobs/{jobId}:
get:
operationId: getJob
Expand Down Expand Up @@ -289,18 +292,40 @@ components:
additionalProperties: false
properties: { limit: { type: integer, minimum: 1, maximum: 500 } }
CreateJob:
oneOf:
- $ref: '#/components/schemas/CreatePaymentCaptureJob'
- $ref: '#/components/schemas/CreatePaymentDisbursementJob'
- $ref: '#/components/schemas/CreatePaymentSettlementJob'
- $ref: '#/components/schemas/CreatePaymentReconciliationJob'
discriminator: { propertyName: type }
CreatePaymentCaptureJob:
type: object
additionalProperties: false
required: [type, payload]
properties:
type: { type: string, enum: [PAYMENT_CAPTURE, PAYMENT_DISBURSEMENT, PAYMENT_SETTLEMENT, PAYMENT_RECONCILIATION] }
queue: { type: string, const: payments, default: payments }
payload:
oneOf:
- $ref: '#/components/schemas/PaymentStartPayload'
- $ref: '#/components/schemas/PaymentSettlementPayload'
- $ref: '#/components/schemas/PaymentReconciliationPayload'
max_attempts: { type: integer, minimum: 1, default: 3 }
type: { type: string, const: PAYMENT_CAPTURE }
payload: { $ref: '#/components/schemas/PaymentStartPayload' }
CreatePaymentDisbursementJob:
type: object
additionalProperties: false
required: [type, payload]
properties:
type: { type: string, const: PAYMENT_DISBURSEMENT }
payload: { $ref: '#/components/schemas/PaymentStartPayload' }
CreatePaymentSettlementJob:
type: object
additionalProperties: false
required: [type, payload]
properties:
type: { type: string, const: PAYMENT_SETTLEMENT }
payload: { $ref: '#/components/schemas/PaymentSettlementPayload' }
CreatePaymentReconciliationJob:
type: object
additionalProperties: false
required: [type, payload]
properties:
type: { type: string, const: PAYMENT_RECONCILIATION }
payload: { $ref: '#/components/schemas/PaymentReconciliationPayload' }
WorkerJob:
type: object
required: [id, queue, type, tenant_id, payload, status, attempts, max_attempts, created_at, updated_at]
Expand Down Expand Up @@ -397,12 +422,19 @@ components:
payload: { type: object, additionalProperties: true }
PlatformStatus:
type: object
additionalProperties: false
required: [status, service, version, environment, started_at, uptime_seconds, dependencies, worker, schedules, queues]
properties:
service: { type: string }
status: { type: string, enum: [ok, degraded, down] }
timestamp: { type: string, format: date-time }
version: { type: string }
environment: { type: string }
started_at: { type: string, format: date-time }
uptime_seconds: { type: integer, minimum: 0 }
worker: { type: object, additionalProperties: true }
dependencies: { type: object, additionalProperties: true }
schedules: { type: array, items: { $ref: '#/components/schemas/ScheduledJob' } }
queues: { type: array, items: { $ref: '#/components/schemas/QueueStats' } }
HttpError:
type: object
required: [statusCode, message]
Expand Down Expand Up @@ -432,3 +464,6 @@ components:
Conflict:
description: Idempotency key, state, or delivery reference conflicts with existing state
content: { application/json: { schema: { $ref: '#/components/schemas/HttpError' } } }
ServiceUnavailable:
description: Durable receipt storage or queue execution is temporarily unavailable
content: { application/json: { schema: { $ref: '#/components/schemas/HttpError' } } }
7 changes: 6 additions & 1 deletion package.json
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,9 @@
"start:dev": "ts-node-dev --respawn --transpile-only src/main.ts",
"test": "NODE_OPTIONS=--experimental-vm-modules jest --runInBand",
"test:e2e": "NODE_OPTIONS=--experimental-vm-modules jest --config jest.e2e.config.js --runInBand",
"test:all": "NODE_OPTIONS=--experimental-vm-modules jest --runInBand && NODE_OPTIONS=--experimental-vm-modules jest --config jest.e2e.config.js --runInBand"
"test:all": "NODE_OPTIONS=--experimental-vm-modules jest --runInBand && NODE_OPTIONS=--experimental-vm-modules jest --config jest.e2e.config.js --runInBand",
"prisma:migrate": "prisma migrate deploy --schema prisma/schema.prisma",
"database:provision-role": "node scripts/provision-database-role.mjs"
},
"dependencies": {
"@mavula/legacy-connectors": "workspace:*",
Expand All @@ -27,15 +29,18 @@
"jose": "6.2.3",
"bullmq": "^5.0.0",
"ioredis": "^5.3.2",
"pg": "^8.13.1",
"reflect-metadata": "^0.1.13",
"rxjs": "^7.8.0"
},
"devDependencies": {
"@nestjs/testing": "^10.0.0",
"@types/jest": "^29.0.0",
"@types/multer": "^1.4.12",
"@types/pg": "^8.11.10",
"@types/supertest": "^6.0.2",
"jest": "^29.0.0",
"prisma": "5.7.0",
"supertest": "^6.3.3",
"ts-jest": "^29.0.0",
"ts-node-dev": "^2.0.0",
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
CREATE SCHEMA IF NOT EXISTS workbench;

DO $$
BEGIN
IF NOT EXISTS (SELECT 1 FROM pg_roles WHERE rolname = 'workbench_app') THEN
CREATE ROLE workbench_app NOLOGIN;
END IF;
IF NOT EXISTS (SELECT 1 FROM pg_roles WHERE rolname = 'workbench_maintenance') THEN
CREATE ROLE workbench_maintenance NOLOGIN;
END IF;
ALTER ROLE workbench_app NOSUPERUSER NOCREATEDB NOCREATEROLE NOINHERIT NOBYPASSRLS;
ALTER ROLE workbench_maintenance NOSUPERUSER NOCREATEDB NOCREATEROLE NOINHERIT NOBYPASSRLS;
END
$$;

GRANT USAGE ON SCHEMA workbench TO workbench_app, workbench_maintenance;

CREATE TABLE workbench.job_submission_receipts (
id text PRIMARY KEY,
"tenantId" text NOT NULL,
operation text NOT NULL CHECK (operation = 'create-job'),
"keyDigest" text NOT NULL CHECK ("keyDigest" ~ '^[a-f0-9]{64}$'),
"requestHash" text NOT NULL CHECK ("requestHash" ~ '^[a-f0-9]{64}$'),
"jobId" text NOT NULL,
"correlationId" text NOT NULL,
"actorId" text NOT NULL,
state text NOT NULL CHECK (state IN ('PENDING','COMPLETED')),
"httpStatus" integer CHECK ("httpStatus" BETWEEN 200 AND 299),
"responseBody" jsonb,
"completedAt" timestamp(3),
"expiresAt" timestamp(3) NOT NULL,
"createdAt" timestamp(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
"updatedAt" timestamp(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE ("tenantId", operation, "keyDigest"),
UNIQUE ("tenantId", "jobId"),
CHECK (
(state = 'PENDING' AND "httpStatus" IS NULL AND "responseBody" IS NULL AND "completedAt" IS NULL)
OR
(state = 'COMPLETED' AND "httpStatus" IS NOT NULL AND "responseBody" IS NOT NULL AND "completedAt" IS NOT NULL)
)
);
CREATE INDEX job_submission_receipts_expires_idx ON workbench.job_submission_receipts("expiresAt");
CREATE INDEX job_submission_receipts_tenant_created_idx ON workbench.job_submission_receipts("tenantId", "createdAt");

CREATE OR REPLACE FUNCTION workbench.current_tenant_id()
RETURNS text LANGUAGE sql STABLE
AS $$ SELECT NULLIF(current_setting('app.current_tenant_id', true), '') $$;
REVOKE ALL ON FUNCTION workbench.current_tenant_id() FROM PUBLIC;
GRANT EXECUTE ON FUNCTION workbench.current_tenant_id() TO workbench_app;

ALTER TABLE workbench.job_submission_receipts ENABLE ROW LEVEL SECURITY;
ALTER TABLE workbench.job_submission_receipts FORCE ROW LEVEL SECURITY;
CREATE POLICY tenant_isolation ON workbench.job_submission_receipts TO workbench_app
USING ("tenantId" = workbench.current_tenant_id())
WITH CHECK ("tenantId" = workbench.current_tenant_id());
CREATE POLICY maintenance_access ON workbench.job_submission_receipts TO workbench_maintenance
USING (true) WITH CHECK (true);

GRANT SELECT, INSERT ON workbench.job_submission_receipts TO workbench_app;
GRANT UPDATE (state, "httpStatus", "responseBody", "completedAt", "updatedAt")
ON workbench.job_submission_receipts TO workbench_app;
GRANT SELECT, DELETE ON workbench.job_submission_receipts TO workbench_maintenance;

CREATE OR REPLACE FUNCTION workbench.delete_expired_job_submission_receipt(
tenant_id text, receipt_operation text, key_digest text
) RETURNS integer
LANGUAGE plpgsql SECURITY DEFINER
SET search_path = pg_catalog, workbench
AS $$
DECLARE deleted_count integer;
BEGIN
DELETE FROM workbench.job_submission_receipts
WHERE "tenantId" = tenant_id AND operation = receipt_operation
AND "keyDigest" = key_digest AND "expiresAt" <= now();
GET DIAGNOSTICS deleted_count = ROW_COUNT;
RETURN deleted_count;
END
$$;

CREATE OR REPLACE FUNCTION workbench.cleanup_expired_job_submission_receipts(batch_limit integer)
RETURNS integer
LANGUAGE plpgsql SECURITY DEFINER
SET search_path = pg_catalog, workbench
AS $$
DECLARE deleted_count integer;
BEGIN
WITH expired AS (
SELECT id FROM workbench.job_submission_receipts
WHERE "expiresAt" <= now() ORDER BY "expiresAt" LIMIT batch_limit
FOR UPDATE SKIP LOCKED
)
DELETE FROM workbench.job_submission_receipts receipt
USING expired WHERE receipt.id = expired.id;
GET DIAGNOSTICS deleted_count = ROW_COUNT;
RETURN deleted_count;
END
$$;

ALTER FUNCTION workbench.delete_expired_job_submission_receipt(text, text, text) OWNER TO workbench_maintenance;
ALTER FUNCTION workbench.cleanup_expired_job_submission_receipts(integer) OWNER TO workbench_maintenance;
REVOKE ALL ON FUNCTION workbench.delete_expired_job_submission_receipt(text, text, text) FROM PUBLIC;
REVOKE ALL ON FUNCTION workbench.cleanup_expired_job_submission_receipts(integer) FROM PUBLIC;
GRANT EXECUTE ON FUNCTION workbench.delete_expired_job_submission_receipt(text, text, text) TO workbench_app;
GRANT EXECUTE ON FUNCTION workbench.cleanup_expired_job_submission_receipts(integer) TO workbench_app;
1 change: 1 addition & 0 deletions prisma/migrations/migration_lock.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
provider = "postgresql"
28 changes: 28 additions & 0 deletions prisma/schema.prisma
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
datasource db {
provider = "postgresql"
url = env("DATABASE_URL")

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Use the documented migration URL for Prisma

With the new environment split, .env.example no longer defines DATABASE_URL and instead documents WORKBENCH_MIGRATION_DATABASE_URL for privileged migrations, but this datasource still reads only DATABASE_URL. In deployments following the new contract, pnpm prisma:migrate will fail before creating the workbench schema/roles, or it will accidentally run with a runtime URL if an old DATABASE_URL remains. Point the Prisma datasource/script at WORKBENCH_MIGRATION_DATABASE_URL or explicitly bridge it in the migrate command.

Useful? React with 👍 / 👎.

}

model JobSubmissionReceipt {
id String @id
tenantId String
operation String
keyDigest String
requestHash String
jobId String
correlationId String
actorId String
state String
httpStatus Int?
responseBody Json?
completedAt DateTime?
expiresAt DateTime
createdAt DateTime @default(now())
updatedAt DateTime @default(now()) @updatedAt

@@unique([tenantId, operation, keyDigest])
@@unique([tenantId, jobId])
@@index([expiresAt])
@@index([tenantId, createdAt])
@@map("job_submission_receipts")
}
6 changes: 6 additions & 0 deletions scripts/check-openapi.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -25,4 +25,10 @@ if (summaries.length !== operationIds.length || permissions.length !== operation
for (const schema of ['CreateJob', 'WorkerJob', 'LegacyBatchReceipt', 'LegacyRejectionReport', 'PlatformStatus']) {
if (!source.includes(` ${schema}:`)) throw new Error(`OpenAPI schema missing: ${schema}`);
}
const declaredSchemas = new Set(
[...source.matchAll(/^ ([A-Za-z][A-Za-z0-9]+):$/gm)].map((match) => match[1]),
);
for (const reference of source.matchAll(/\$ref: '#\/components\/schemas\/([A-Za-z][A-Za-z0-9]+)'/g)) {
if (!declaredSchemas.has(reference[1])) throw new Error(`OpenAPI schema reference missing: ${reference[1]}`);
}
console.log(`workbench OpenAPI covers ${paths.size} public routes`);
34 changes: 34 additions & 0 deletions scripts/provision-database-role.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
#!/usr/bin/env node
import pg from 'pg';

const databaseUrl = required('WORKBENCH_MIGRATION_DATABASE_URL');
const password = required('WORKBENCH_DATABASE_ROLE_PASSWORD');
if (password.length < 16 || password.startsWith('REPLACE_WITH_')) {
throw new Error('WORKBENCH_DATABASE_ROLE_PASSWORD must be a non-placeholder secret of at least 16 characters');
}
const pool = new pg.Pool({ connectionString: withoutSchema(databaseUrl) });
try {
await pool.query(`ALTER ROLE workbench_app WITH LOGIN PASSWORD '${password.replaceAll("'", "''")}'`);
const { rows } = await pool.query(
'SELECT rolcanlogin, rolsuper, rolcreatedb, rolcreaterole, rolinherit, rolbypassrls FROM pg_roles WHERE rolname = $1',
['workbench_app'],
);
const role = rows[0];
if (!role?.rolcanlogin || role.rolsuper || role.rolcreatedb || role.rolcreaterole || role.rolinherit || role.rolbypassrls) {
throw new Error('workbench_app role attributes do not satisfy the runtime policy');
}
} finally {
await pool.end();
}

function required(name) {
const value = process.env[name]?.trim();
if (!value) throw new Error(`${name} is required`);
return value;
}

function withoutSchema(value) {
const url = new URL(value);
url.searchParams.delete('schema');
return url.toString();
}
4 changes: 4 additions & 0 deletions src/app.module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,11 +21,15 @@ import { PermissionsGuard } from './auth/permissions.guard';
import { ServiceTokenService } from './auth/service-token.service';
import { LegacyBatchesController } from './controllers/legacy-batches.controller';
import { LegacyBatchRuntimeService } from './worker/legacy-batch-runtime.service';
import { JobSubmissionService } from './idempotency/job-submission.service';
import { MetricsTokenGuard } from './auth/metrics-token.guard';

@Module({
controllers: [StatusController, JobsController, LegacyBatchesController],
providers: [
JobStoreService,
JobSubmissionService,
MetricsTokenGuard,
PlatformStatusService,
JobHandlersService,
PaymentOutboxPublisherService,
Expand Down
24 changes: 24 additions & 0 deletions src/auth/metrics-token.guard.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
import { CanActivate, ExecutionContext, Injectable, UnauthorizedException } from '@nestjs/common';
import { createHash, timingSafeEqual } from 'node:crypto';
import { getRuntimeConfig } from '../utils/runtime-config';

@Injectable()
export class MetricsTokenGuard implements CanActivate {
private readonly expected = getRuntimeConfig().metricsToken;

canActivate(context: ExecutionContext): boolean {
if (!this.expected) {
throw new UnauthorizedException('Metrics scrape token is not configured');
}
const authorization = context.switchToHttp().getRequest().headers.authorization;
if (typeof authorization !== 'string' || !authorization.startsWith('Bearer ')) {
throw new UnauthorizedException('Metrics scrape token is required');
}
const expectedDigest = createHash('sha256').update(this.expected).digest();
const suppliedDigest = createHash('sha256').update(authorization.slice(7)).digest();
if (!timingSafeEqual(expectedDigest, suppliedDigest)) {
throw new UnauthorizedException('Invalid metrics scrape token');
}
return true;
}
}
Loading