From 95739b8976c30909e54dcb17b2b4048eebe93efe Mon Sep 17 00:00:00 2001 From: Ayobami Haastrup <47716486+AyobamiH@users.noreply.github.com> Date: Tue, 8 Sep 2026 11:31:11 -0700 Subject: [PATCH 1/5] Fail closed on held legacy queue revisions --- package.json | 2 +- src/publication-ledger.ts | 1 + src/supabase-worker.ts | 52 ++++++++++++++++++++++++++++++- src/worker-claims.ts | 1 + test/legacy-revision-hold.test.ts | 34 ++++++++++++++++++++ 5 files changed, 88 insertions(+), 2 deletions(-) create mode 100644 test/legacy-revision-hold.test.ts diff --git a/package.json b/package.json index 7be1961..3aa6c11 100644 --- a/package.json +++ b/package.json @@ -6,7 +6,7 @@ "scripts": { "build": "tsx scripts/build.ts", "typecheck": "tsc --noEmit --project tsconfig.json", - "test": "npm run build && node dist/test/security-hardening.test.js && node dist/test/source-ssrf.test.js && node dist/test/browser-collector-ingest.test.js && node dist/test/slot-scheduler.test.js && node dist/test/daily-inventory-planner.test.js && node dist/test/refresh-queue-finalization.test.js && node dist/test/publish-summary.test.js && node dist/test/threads-refresh.test.js && node dist/test/linkedin-refresh.test.js && node dist/test/instagram-image-timeout.test.js && node dist/test/cost-control.test.js && node dist/test/recovery-scheduler.test.js && node dist/test/social-connector.test.js && node dist/test/meta-publication-boundary.test.js && node dist/test/cloudflare-health.test.js && node dist/test/exclusive-run-gate.test.js && node dist/test/runtime-scope.test.js && node dist/test/tenant-platform-policy.test.js && node dist/test/supabase-client-retry.test.js && node dist/test/worker-claims.test.js && node dist/test/publication-ledger.test.js && node dist/test/publication-outcome.test.js && node dist/test/publication-executor.test.js && node dist/test/provider-single-dispatch.test.js && node dist/test/publication-receipts.test.js", + "test": "npm run build && node dist/test/security-hardening.test.js && node dist/test/source-ssrf.test.js && node dist/test/browser-collector-ingest.test.js && node dist/test/slot-scheduler.test.js && node dist/test/daily-inventory-planner.test.js && node dist/test/refresh-queue-finalization.test.js && node dist/test/publish-summary.test.js && node dist/test/threads-refresh.test.js && node dist/test/linkedin-refresh.test.js && node dist/test/instagram-image-timeout.test.js && node dist/test/cost-control.test.js && node dist/test/recovery-scheduler.test.js && node dist/test/social-connector.test.js && node dist/test/meta-publication-boundary.test.js && node dist/test/cloudflare-health.test.js && node dist/test/exclusive-run-gate.test.js && node dist/test/runtime-scope.test.js && node dist/test/tenant-platform-policy.test.js && node dist/test/supabase-client-retry.test.js && node dist/test/worker-claims.test.js && node dist/test/publication-ledger.test.js && node dist/test/publication-outcome.test.js && node dist/test/publication-executor.test.js && node dist/test/provider-single-dispatch.test.js && node dist/test/publication-receipts.test.js && node dist/test/legacy-revision-hold.test.js", "smoke:dist": "node dist/src/cli.js status", "ci": "npm run typecheck && npm test && npm run smoke:dist", "dev": "tsx src/agent.ts", diff --git a/src/publication-ledger.ts b/src/publication-ledger.ts index 4fa5d92..16317ef 100644 --- a/src/publication-ledger.ts +++ b/src/publication-ledger.ts @@ -17,6 +17,7 @@ export const REQUIRED_PUBLICATION_CAPABILITIES = [ 'publication-exact-history-receipt-v1', 'publication-queue-compatibility-fence-v1', 'publication-provenance-snapshot-v1', + 'publication-legacy-queue-hold-v1', ] as const; export interface PublicationSchemaContract { diff --git a/src/supabase-worker.ts b/src/supabase-worker.ts index c436436..75dea10 100644 --- a/src/supabase-worker.ts +++ b/src/supabase-worker.ts @@ -172,6 +172,7 @@ interface QueueItemRow { retry_of?: string | null; recovery_execution_id?: string | null; recovery_reason?: string | null; + legacy_revision_hold_id?: string | null; } interface PublishHistoryRow { @@ -471,6 +472,26 @@ class WorkerJobError extends Error { } } +const LEGACY_QUEUE_REVISION_HELD_CODE = 'legacy_queue_revision_held'; +const LEGACY_QUEUE_REVISION_HELD_MESSAGE = + 'This historical queue revision is protected because its publication outcome cannot be proven safely.'; +const LEGACY_QUEUE_REVISION_HELD_NEXT_ACTION = + 'Do not retry, edit, skip, release, or recreate this post. Preserve it for operator review.'; + +function assertQueueRevisionNotHeld( + row: Pick +): void { + if (!row.legacy_revision_hold_id) return; + throw new WorkerJobError( + LEGACY_QUEUE_REVISION_HELD_CODE, + LEGACY_QUEUE_REVISION_HELD_MESSAGE, + { + queueItemId: row.id, + next_action: LEGACY_QUEUE_REVISION_HELD_NEXT_ACTION, + } + ); +} + interface SourceRecordRow { id: string; user_id: string; @@ -3591,6 +3612,7 @@ async function assertOpenAIRepairAllowedForPublish( async function publishQueueRow(job: AgentJobRow, row: QueueItemRow, _settings: UserSettingsRow): Promise { if (row.user_id !== job.user_id) throw new WorkerJobError('publication_tenant_mismatch'); + assertQueueRevisionNotHeld(row); // Reload per item. publish_all must not retain the first item's stale credentials/settings. const tenant = await loadTenantContext(job.user_id); return withTenantRuntime(tenant, async () => { @@ -3724,6 +3746,13 @@ async function handleSkipSlot(job: AgentJobRow): Promise { throw new WorkerJobError('missing_queue_target', 'queue_item_id or slot_index is required'); } + const targets = await supabaseSelect('queue_items', { + select: 'id,user_id,legacy_revision_hold_id', + filters, + limit: 100, + }); + targets.forEach(assertQueueRevisionNotHeld); + const rows = await supabaseUpdate('queue_items', { status: 'skipped', error_message: null, @@ -3764,6 +3793,13 @@ async function handleReleaseSlot(job: AgentJobRow): Promise { throw new WorkerJobError('missing_queue_target', 'queue_item_id or slot_index is required'); } + const targets = await supabaseSelect('queue_items', { + select: 'id,user_id,legacy_revision_hold_id', + filters, + limit: 100, + }); + targets.forEach(assertQueueRevisionNotHeld); + const rows = await supabaseDelete('queue_items', { filters, returning: true, @@ -4248,7 +4284,7 @@ async function loadPriorFailedRecovery(row: QueueItemRow): Promise { const settingsByUser = await loadAutomationSettingsByUser(); const dueRows = await supabaseSelect('queue_items', { - select: 'id,user_id,platform,status,slot_index,scheduled_for,scheduled_local_date,scheduled_timezone,retry_of,recovery_execution_id', + select: 'id,user_id,platform,status,slot_index,scheduled_for,scheduled_local_date,scheduled_timezone,retry_of,recovery_execution_id,legacy_revision_hold_id', filters: [ { column: 'status', operator: 'in', value: ['pending', 'ready'] }, { column: 'scheduled_for', operator: 'lte', value: now.toISOString() }, @@ -4258,6 +4294,19 @@ async function enqueueDuePublishJobs(stats: SchedulerStats, now: Date): Promise< }); for (const row of dueRows) { + if (row.legacy_revision_hold_id) { + incrementSchedulerSkip(stats, LEGACY_QUEUE_REVISION_HELD_CODE); + await recordAutomationSkip( + row.user_id, + 'publish_now', + LEGACY_QUEUE_REVISION_HELD_CODE, + LEGACY_QUEUE_REVISION_HELD_MESSAGE, + LEGACY_QUEUE_REVISION_HELD_NEXT_ACTION, + { queueItemId: row.id, platform: row.platform } + ); + continue; + } + const settings = settingsByUser.get(row.user_id); if (!settings || !settings.automation_publish_enabled) { incrementSchedulerSkip(stats, 'publish_automation_disabled'); @@ -4889,6 +4938,7 @@ export function startSupabaseWorkerLoop(log = logger): { stop: () => void } | un } export const __test__ = { + assertQueueRevisionNotHeld, publishQueueRow, stalePublishJobResult, handlePublishAll, diff --git a/src/worker-claims.ts b/src/worker-claims.ts index dfcbae0..b84ad80 100644 --- a/src/worker-claims.ts +++ b/src/worker-claims.ts @@ -11,6 +11,7 @@ export const REQUIRED_WORKER_CLAIM_CAPABILITIES = [ 'source-angle-atomic-commit-v1', 'angle-queue-atomic-commit-v1', 'queue-angle-identity-v1', + 'legacy-queue-revision-hold-v1', ] as const; export interface WorkerSchemaContract { diff --git a/test/legacy-revision-hold.test.ts b/test/legacy-revision-hold.test.ts new file mode 100644 index 0000000..63cafad --- /dev/null +++ b/test/legacy-revision-hold.test.ts @@ -0,0 +1,34 @@ +import assert from 'node:assert/strict'; +import test from 'node:test'; + +import { __test__ } from '../src/supabase-worker'; + +test('held historical revisions fail before any provider or mutation path', () => { + assert.doesNotThrow(() => { + __test__.assertQueueRevisionNotHeld({ + id: 'queue-safe', + legacy_revision_hold_id: null, + }); + }); + + assert.throws( + () => { + __test__.assertQueueRevisionNotHeld({ + id: 'queue-held', + legacy_revision_hold_id: 'legacy-hold', + }); + }, + (error: unknown) => { + const candidate = error as { + code?: string; + message?: string; + context?: Record; + }; + return candidate.code === 'legacy_queue_revision_held' + && candidate.message?.includes('publication outcome cannot be proven safely') === true + && candidate.context?.queueItemId === 'queue-held' + && candidate.context?.next_action === + 'Do not retry, edit, skip, release, or recreate this post. Preserve it for operator review.'; + } + ); +}); From d4f5c9d05dfa0a8a26fc1ff4f48a1bd633dbcfa2 Mon Sep 17 00:00:00 2001 From: Ayobami Haastrup <47716486+AyobamiH@users.noreply.github.com> Date: Tue, 8 Sep 2026 11:57:35 -0700 Subject: [PATCH 2/5] Add fail-closed production canary controls --- .../workflows/deploy-cloudflare-worker.yml | 30 ++++++- config.ts | 16 ++++ package.json | 2 +- src/canary-policy.ts | 68 +++++++++++++++ src/cloudflare-worker.ts | 32 +++++++ src/publication-executor.ts | 9 +- src/supabase-worker.ts | 84 +++++++++++++++++- test/canary-policy.test.ts | 87 +++++++++++++++++++ test/cloudflare-health.test.ts | 30 +++++++ wrangler.toml | 12 ++- 10 files changed, 361 insertions(+), 9 deletions(-) create mode 100644 src/canary-policy.ts create mode 100644 test/canary-policy.test.ts diff --git a/.github/workflows/deploy-cloudflare-worker.yml b/.github/workflows/deploy-cloudflare-worker.yml index cb0bc08..bc850b0 100644 --- a/.github/workflows/deploy-cloudflare-worker.yml +++ b/.github/workflows/deploy-cloudflare-worker.yml @@ -33,11 +33,25 @@ jobs: EXPECTED_SHA: ${{ inputs.expected_sha }} RELEASE_SHA: ${{ github.sha }} ROLLOUT_PREFLIGHT_CONFIRMED: ${{ inputs.rollout_preflight_confirmed }} + CANARY_USER_IDS: ${{ secrets.SUPABASE_WORKER_CANARY_USER_IDS }} run: | set -eu test "$ROLLOUT_PREFLIGHT_CONFIRMED" = 'true' test "$EXPECTED_SHA" = "$RELEASE_SHA" echo "$EXPECTED_SHA" | grep -Eq '^[0-9a-f]{40}$' + CANARY_USER_IDS="$CANARY_USER_IDS" node <<'NODE' + const ids = String(process.env.CANARY_USER_IDS || '') + .split(',') + .map(value => value.trim()) + .filter(Boolean); + if (ids.length < 1 || ids.length > 3) { + throw new Error('SUPABASE_WORKER_CANARY_USER_IDS must contain 1-3 tenant UUIDs'); + } + const uuid = /^[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i; + if (new Set(ids.map(value => value.toLowerCase())).size !== ids.length || ids.some(value => !uuid.test(value))) { + throw new Error('SUPABASE_WORKER_CANARY_USER_IDS must contain unique tenant UUIDs'); + } + NODE - uses: actions/checkout@v4 with: ref: ${{ github.sha }} @@ -48,7 +62,21 @@ jobs: cache: npm - run: npm ci - run: npm run ci - - run: npx wrangler deploy --tag "${{ github.sha }}" + - name: Deploy the allowlisted, provider-disabled canary env: CLOUDFLARE_API_TOKEN: ${{ secrets.CLOUDFLARE_API_TOKEN }} CLOUDFLARE_ACCOUNT_ID: ${{ secrets.CLOUDFLARE_ACCOUNT_ID }} + CANARY_USER_IDS: ${{ secrets.SUPABASE_WORKER_CANARY_USER_IDS }} + run: | + set -eu + secrets_file="$RUNNER_TEMP/worker-canary-secrets.json" + trap 'rm -f "$secrets_file"' EXIT + CANARY_USER_IDS="$CANARY_USER_IDS" SECRETS_FILE="$secrets_file" node <<'NODE' + const fs = require('node:fs'); + fs.writeFileSync( + process.env.SECRETS_FILE, + JSON.stringify({ SUPABASE_WORKER_CANARY_USER_IDS: process.env.CANARY_USER_IDS }), + { mode: 0o600 } + ); + NODE + npx wrangler deploy --tag "${{ github.sha }}" --secrets-file "$secrets_file" diff --git a/config.ts b/config.ts index a799fff..b8088bf 100644 --- a/config.ts +++ b/config.ts @@ -109,6 +109,10 @@ export interface AppConfig { CREDENTIAL_ENCRYPTION_KEY: string; SUPABASE_WORKER_POLL_INTERVAL_MS: number; SUPABASE_WORKER_BATCH_SIZE: number; + SUPABASE_WORKER_CANARY_REQUIRED: boolean; + SUPABASE_WORKER_CANARY_USER_IDS: Set; + SUPABASE_WORKER_GENERATION_ENABLED: boolean; + SUPABASE_PROVIDER_DISPATCH_ENABLED: boolean; DAILY_INVENTORY_PLANNER_ENABLED: boolean; DAILY_INVENTORY_PLANNER_START_LOCAL_DATE: string; } @@ -241,6 +245,18 @@ function buildBaseConfig(): AppConfig { CREDENTIAL_ENCRYPTION_KEY: IS_SOCIAL_CONNECTOR_ONLY ? '' : process.env.CREDENTIAL_ENCRYPTION_KEY || '', SUPABASE_WORKER_POLL_INTERVAL_MS: Number.parseInt(process.env.SUPABASE_WORKER_POLL_INTERVAL_MS || '10000', 10), SUPABASE_WORKER_BATCH_SIZE: Number.parseInt(process.env.SUPABASE_WORKER_BATCH_SIZE || '10', 10), + SUPABASE_WORKER_CANARY_REQUIRED: + (process.env.NODE_ENV || 'development') === 'production' + || parseBooleanEnv(process.env.SUPABASE_WORKER_CANARY_REQUIRED, false), + SUPABASE_WORKER_CANARY_USER_IDS: toSubSet(process.env.SUPABASE_WORKER_CANARY_USER_IDS, ''), + SUPABASE_WORKER_GENERATION_ENABLED: parseBooleanEnv( + process.env.SUPABASE_WORKER_GENERATION_ENABLED, + (process.env.NODE_ENV || 'development') !== 'production' + ), + SUPABASE_PROVIDER_DISPATCH_ENABLED: parseBooleanEnv( + process.env.SUPABASE_PROVIDER_DISPATCH_ENABLED, + (process.env.NODE_ENV || 'development') !== 'production' + ), DAILY_INVENTORY_PLANNER_ENABLED: parseBooleanEnv(process.env.DAILY_INVENTORY_PLANNER_ENABLED, false), DAILY_INVENTORY_PLANNER_START_LOCAL_DATE: process.env.DAILY_INVENTORY_PLANNER_START_LOCAL_DATE || '', }; diff --git a/package.json b/package.json index 3aa6c11..5c59d97 100644 --- a/package.json +++ b/package.json @@ -6,7 +6,7 @@ "scripts": { "build": "tsx scripts/build.ts", "typecheck": "tsc --noEmit --project tsconfig.json", - "test": "npm run build && node dist/test/security-hardening.test.js && node dist/test/source-ssrf.test.js && node dist/test/browser-collector-ingest.test.js && node dist/test/slot-scheduler.test.js && node dist/test/daily-inventory-planner.test.js && node dist/test/refresh-queue-finalization.test.js && node dist/test/publish-summary.test.js && node dist/test/threads-refresh.test.js && node dist/test/linkedin-refresh.test.js && node dist/test/instagram-image-timeout.test.js && node dist/test/cost-control.test.js && node dist/test/recovery-scheduler.test.js && node dist/test/social-connector.test.js && node dist/test/meta-publication-boundary.test.js && node dist/test/cloudflare-health.test.js && node dist/test/exclusive-run-gate.test.js && node dist/test/runtime-scope.test.js && node dist/test/tenant-platform-policy.test.js && node dist/test/supabase-client-retry.test.js && node dist/test/worker-claims.test.js && node dist/test/publication-ledger.test.js && node dist/test/publication-outcome.test.js && node dist/test/publication-executor.test.js && node dist/test/provider-single-dispatch.test.js && node dist/test/publication-receipts.test.js && node dist/test/legacy-revision-hold.test.js", + "test": "npm run build && node dist/test/security-hardening.test.js && node dist/test/source-ssrf.test.js && node dist/test/browser-collector-ingest.test.js && node dist/test/slot-scheduler.test.js && node dist/test/daily-inventory-planner.test.js && node dist/test/refresh-queue-finalization.test.js && node dist/test/publish-summary.test.js && node dist/test/threads-refresh.test.js && node dist/test/linkedin-refresh.test.js && node dist/test/instagram-image-timeout.test.js && node dist/test/cost-control.test.js && node dist/test/recovery-scheduler.test.js && node dist/test/social-connector.test.js && node dist/test/meta-publication-boundary.test.js && node dist/test/cloudflare-health.test.js && node dist/test/canary-policy.test.js && node dist/test/exclusive-run-gate.test.js && node dist/test/runtime-scope.test.js && node dist/test/tenant-platform-policy.test.js && node dist/test/supabase-client-retry.test.js && node dist/test/worker-claims.test.js && node dist/test/publication-ledger.test.js && node dist/test/publication-outcome.test.js && node dist/test/publication-executor.test.js && node dist/test/provider-single-dispatch.test.js && node dist/test/publication-receipts.test.js && node dist/test/legacy-revision-hold.test.js", "smoke:dist": "node dist/src/cli.js status", "ci": "npm run typecheck && npm test && npm run smoke:dist", "dev": "tsx src/agent.ts", diff --git a/src/canary-policy.ts b/src/canary-policy.ts new file mode 100644 index 0000000..a0eb592 --- /dev/null +++ b/src/canary-policy.ts @@ -0,0 +1,68 @@ +export type RolloutControlledJobKind = + | 'fetch_sources' + | 'refresh_queue' + | 'publish_now' + | 'publish_all' + | 'skip_slot' + | 'release_slot'; + +export interface RolloutPolicy { + canaryRequired: boolean; + canaryUserIds: ReadonlySet; + generationEnabled: boolean; + providerDispatchEnabled: boolean; +} + +const GENERATION_JOB_KINDS = new Set([ + 'fetch_sources', + 'refresh_queue', +]); + +const PUBLICATION_JOB_KINDS = new Set([ + 'publish_now', + 'publish_all', +]); + +const OPERATOR_JOB_KINDS: RolloutControlledJobKind[] = [ + 'skip_slot', + 'release_slot', +]; + +export function allowedCanaryUserIds(policy: RolloutPolicy): string[] | undefined { + if (!policy.canaryRequired) return undefined; + return [...policy.canaryUserIds] + .map(userId => userId.trim().toLowerCase()) + .filter(Boolean) + .sort(); +} + +export function canaryAllowsUser(policy: RolloutPolicy, userId: string): boolean { + if (!policy.canaryRequired) return true; + const normalized = userId.trim().toLowerCase(); + return [...policy.canaryUserIds].some( + allowed => allowed.trim().toLowerCase() === normalized + ); +} + +export function runnableJobKinds(policy: RolloutPolicy): RolloutControlledJobKind[] { + return [ + ...(policy.generationEnabled ? [...GENERATION_JOB_KINDS] : []), + ...(policy.providerDispatchEnabled ? [...PUBLICATION_JOB_KINDS] : []), + ...OPERATOR_JOB_KINDS, + ]; +} + +export function rolloutBlockCode( + policy: RolloutPolicy, + userId: string, + kind: RolloutControlledJobKind +): string | undefined { + if (!canaryAllowsUser(policy, userId)) return 'rollout_canary_tenant_blocked'; + if (GENERATION_JOB_KINDS.has(kind) && !policy.generationEnabled) { + return 'rollout_generation_disabled'; + } + if (PUBLICATION_JOB_KINDS.has(kind) && !policy.providerDispatchEnabled) { + return 'rollout_provider_dispatch_disabled'; + } + return undefined; +} diff --git a/src/cloudflare-worker.ts b/src/cloudflare-worker.ts index 0db1540..912e07c 100644 --- a/src/cloudflare-worker.ts +++ b/src/cloudflare-worker.ts @@ -15,6 +15,10 @@ interface Env { SERVICE_ROLE_KEY?: string; CREDENTIAL_ENCRYPTION_KEY?: string; SUPABASE_WORKER_BATCH_SIZE?: string; + SUPABASE_WORKER_CANARY_REQUIRED?: string; + SUPABASE_WORKER_CANARY_USER_IDS?: string; + SUPABASE_WORKER_GENERATION_ENABLED?: string; + SUPABASE_PROVIDER_DISPATCH_ENABLED?: string; DAILY_INVENTORY_PLANNER_ENABLED?: string; DAILY_INVENTORY_PLANNER_START_LOCAL_DATE?: string; HTTP_TIMEOUT_MS?: string; @@ -87,8 +91,25 @@ function publicationCapabilities(): Record item.trim().toLowerCase()) + .filter(Boolean) + ).size; +} + function healthPayload(env: Env): Record { const metadata = env.CF_VERSION_METADATA; + const production = env.NODE_ENV === 'production'; + const canaryRequired = production || booleanBinding(env.SUPABASE_WORKER_CANARY_REQUIRED, false); + const allowedTenantCount = canaryUserCount(env.SUPABASE_WORKER_CANARY_USER_IDS); return { ok: true, liveness: 'ok', @@ -104,6 +125,17 @@ function healthPayload(env: Env): Record { appliedSchema: 'unverified', }, publicationCapabilities: publicationCapabilities(), + rollout: { + tenantScope: canaryRequired + ? allowedTenantCount > 0 + ? 'allowlisted' + : 'blocked_empty_allowlist' + : 'unrestricted', + canaryRequired, + allowedTenantCount, + generationEnabled: booleanBinding(env.SUPABASE_WORKER_GENERATION_ENABLED, !production), + providerDispatchEnabled: booleanBinding(env.SUPABASE_PROVIDER_DISPATCH_ENABLED, !production), + }, executionGate: scheduledTickGate.snapshot(), }; } diff --git a/src/publication-executor.ts b/src/publication-executor.ts index 40ed525..d5d8997 100644 --- a/src/publication-executor.ts +++ b/src/publication-executor.ts @@ -352,13 +352,20 @@ export async function reconcilePublication( } /** Also recovers orphaned publish_all attempts after their parent job has already ended. */ -export async function recoverStalePublications(now = Date.now()): Promise { +export async function recoverStalePublications( + now = Date.now(), + allowedUserIds?: readonly string[] +): Promise { + if (allowedUserIds?.length === 0) return 0; try { await ledger.assertPublicationLedgerContract(); const intents = await supabaseSelect('publication_intents', { filters: [ { column: 'state', operator: 'in', value: ['claimed', 'dispatching'] }, { column: 'updated_at', operator: 'lte', value: new Date(now - 180_000).toISOString() }, + ...(allowedUserIds + ? [{ column: 'user_id', operator: 'in' as const, value: [...allowedUserIds] }] + : []), ], order: 'updated_at.asc', limit: 50, diff --git a/src/supabase-worker.ts b/src/supabase-worker.ts index 75dea10..fc7a910 100644 --- a/src/supabase-worker.ts +++ b/src/supabase-worker.ts @@ -8,6 +8,12 @@ import * as instagram from './instagram'; import * as linkedin from './linkedin'; import * as logger from './logger'; import { activePlatformsFromSettings } from './platform-settings'; +import { + allowedCanaryUserIds, + rolloutBlockCode, + runnableJobKinds, + type RolloutPolicy, +} from './canary-policy'; import * as workerClaims from './worker-claims'; import { executePublication, reconcilePublication, recoverStalePublications, PublicationPreflightError } from './publication-executor'; import * as threads from './threads'; @@ -540,6 +546,29 @@ const SUPPORTED_JOB_KINDS = new Set([ 'release_slot', ]); +function rolloutPolicy(): RolloutPolicy { + return { + canaryRequired: config.SUPABASE_WORKER_CANARY_REQUIRED, + canaryUserIds: config.SUPABASE_WORKER_CANARY_USER_IDS, + generationEnabled: config.SUPABASE_WORKER_GENERATION_ENABLED, + providerDispatchEnabled: config.SUPABASE_PROVIDER_DISPATCH_ENABLED, + }; +} + +function rolloutAllowedUserIds(): string[] | undefined { + return allowedCanaryUserIds(rolloutPolicy()); +} + +function assertRolloutAllowsJob(userId: string, kind: JobKind): void { + const code = rolloutBlockCode(rolloutPolicy(), userId, kind); + if (!code) return; + throw new WorkerJobError(code, code, { + kind, + userId, + next_action: 'Change the reviewed rollout controls before this job may execute.', + }); +} + const ACTIVE_QUEUE_STATUSES = ['pending', 'ready', 'publishing']; const ACTIVE_ANGLE_STATUSES: AngleRecordStatus[] = ['unused', 'in_progress']; const ANGLE_EXTRACTION_TIMEOUT_MS = 12_000; @@ -2115,9 +2144,21 @@ function assertSupportedJobKind(kind: string): asserts kind is JobKind { } async function listPendingJobs(): Promise { + const policy = rolloutPolicy(); + const allowedUserIds = allowedCanaryUserIds(policy); + if (allowedUserIds?.length === 0) return []; + const kinds = runnableJobKinds(policy); + if (!kinds.length) return []; + return supabaseSelect('agent_jobs', { select: '*', - filters: [{ column: 'status', operator: 'eq', value: 'pending' }], + filters: [ + { column: 'status', operator: 'eq', value: 'pending' }, + { column: 'kind', operator: 'in', value: kinds }, + ...(allowedUserIds + ? [{ column: 'user_id', operator: 'in' as const, value: allowedUserIds }] + : []), + ], order: 'created_at.asc', limit: Math.max(1, Math.min(config.SUPABASE_WORKER_BATCH_SIZE || 10, 50)), }); @@ -3441,6 +3482,9 @@ async function queueFromBankedAngles( } async function handleRefreshQueue(job: AgentJobRow, tenant: TenantContext): Promise { + if (!config.SUPABASE_WORKER_GENERATION_ENABLED) { + throw new WorkerJobError('rollout_generation_disabled'); + } const sources = await supabaseSelect('user_sources', { select: '*', filters: [{ column: 'user_id', operator: 'eq', value: job.user_id }], @@ -3620,6 +3664,9 @@ async function publishQueueRow(job: AgentJobRow, row: QueueItemRow, _settings: U userId: job.user_id, queueItemId: row.id, platform: row.platform, }, { prepare: async payload => { + if (!config.SUPABASE_PROVIDER_DISPATCH_ENABLED) { + throw new PublicationPreflightError('rollout_provider_dispatch_disabled'); + } // Disabled hosted publishers are rejected before token refresh or paid media work. if (payload.platform === 'threads' || payload.platform === 'instagram') { throw new PublicationPreflightError('legacy_meta_publication_disabled'); @@ -3938,11 +3985,17 @@ async function recordAutomationSkip( } async function enqueueDueFetchJobs(stats: SchedulerStats, now: Date): Promise { + if (!config.SUPABASE_WORKER_GENERATION_ENABLED) return; + const allowedUserIds = rolloutAllowedUserIds(); + if (allowedUserIds?.length === 0) return; const settingsRows = await supabaseSelect('user_settings', { select: '*', filters: [ { column: 'automation_enabled', operator: 'eq', value: true }, { column: 'automation_fetch_enabled', operator: 'eq', value: true }, + ...(allowedUserIds + ? [{ column: 'user_id', operator: 'in' as const, value: allowedUserIds }] + : []), ], limit: 200, }); @@ -4089,11 +4142,17 @@ async function recordDailyInventoryAlert( } async function enqueueDueSlotFillJobs(stats: SchedulerStats, now: Date): Promise { + if (!config.SUPABASE_WORKER_GENERATION_ENABLED) return; + const allowedUserIds = rolloutAllowedUserIds(); + if (allowedUserIds?.length === 0) return; const settingsRows = await supabaseSelect('user_settings', { select: '*', filters: [ { column: 'automation_enabled', operator: 'eq', value: true }, { column: 'automation_publish_enabled', operator: 'eq', value: true }, + ...(allowedUserIds + ? [{ column: 'user_id', operator: 'in' as const, value: allowedUserIds }] + : []), ], limit: 200, }); @@ -4257,9 +4316,16 @@ async function enqueueDueSlotFillJobs(stats: SchedulerStats, now: Date): Promise } async function loadAutomationSettingsByUser(): Promise> { + const allowedUserIds = rolloutAllowedUserIds(); + if (allowedUserIds?.length === 0) return new Map(); const rows = await supabaseSelect('user_settings', { select: '*', - filters: [{ column: 'automation_enabled', operator: 'eq', value: true }], + filters: [ + { column: 'automation_enabled', operator: 'eq', value: true }, + ...(allowedUserIds + ? [{ column: 'user_id', operator: 'in' as const, value: allowedUserIds }] + : []), + ], limit: 500, }); return new Map(rows.map(row => [row.user_id, row])); @@ -4282,12 +4348,18 @@ async function loadPriorFailedRecovery(row: QueueItemRow): Promise { + if (!config.SUPABASE_PROVIDER_DISPATCH_ENABLED) return; + const allowedUserIds = rolloutAllowedUserIds(); + if (allowedUserIds?.length === 0) return; const settingsByUser = await loadAutomationSettingsByUser(); const dueRows = await supabaseSelect('queue_items', { select: 'id,user_id,platform,status,slot_index,scheduled_for,scheduled_local_date,scheduled_timezone,retry_of,recovery_execution_id,legacy_revision_hold_id', filters: [ { column: 'status', operator: 'in', value: ['pending', 'ready'] }, { column: 'scheduled_for', operator: 'lte', value: now.toISOString() }, + ...(allowedUserIds + ? [{ column: 'user_id', operator: 'in' as const, value: allowedUserIds }] + : []), ], order: 'scheduled_for.asc', limit: 100, @@ -4754,12 +4826,17 @@ async function stalePublishJobResult(job: AgentJobRow, _logs: WorkerLogRow[]): P } async function cleanupStaleRunningJobs(stats: SchedulerStats, now: Date): Promise { + const allowedUserIds = rolloutAllowedUserIds(); + if (allowedUserIds?.length === 0) return; const cutoff = addMinutesIso(now, -STALE_RUNNING_JOB_MINUTES); const jobs = await supabaseSelect('agent_jobs', { select: '*', filters: [ { column: 'status', operator: 'eq', value: 'running' }, { column: 'started_at', operator: 'lte', value: cutoff }, + ...(allowedUserIds + ? [{ column: 'user_id', operator: 'in' as const, value: allowedUserIds }] + : []), ], order: 'started_at.asc', limit: 50, @@ -4829,7 +4906,7 @@ export async function runSupabaseAutomationScheduler(): Promise skipped: {}, }; const now = new Date(); - await recoverStalePublications(now.getTime()); + await recoverStalePublications(now.getTime(), rolloutAllowedUserIds()); await cleanupStaleRunningJobs(stats, now); await enqueueDueFetchJobs(stats, now); await enqueueDueSlotFillJobs(stats, now); @@ -4839,6 +4916,7 @@ export async function runSupabaseAutomationScheduler(): Promise async function handleClaimedJob(job: AgentJobRow): Promise { assertSupportedJobKind(job.kind); + assertRolloutAllowsJob(job.user_id, job.kind); await assertTenantEntitlement(job); const tenant = await loadTenantContext(job.user_id); const kind = job.kind as JobKind; diff --git a/test/canary-policy.test.ts b/test/canary-policy.test.ts new file mode 100644 index 0000000..f92d7f5 --- /dev/null +++ b/test/canary-policy.test.ts @@ -0,0 +1,87 @@ +import assert from 'node:assert/strict'; +import * as fs from 'node:fs'; +import * as path from 'node:path'; + +import { + allowedCanaryUserIds, + canaryAllowsUser, + rolloutBlockCode, + runnableJobKinds, + type RolloutPolicy, +} from '../src/canary-policy'; + +function policy(overrides: Partial = {}): RolloutPolicy { + return { + canaryRequired: true, + canaryUserIds: new Set(), + generationEnabled: false, + providerDispatchEnabled: false, + ...overrides, + }; +} + +function main(): void { + const blocked = policy(); + assert.deepEqual(allowedCanaryUserIds(blocked), []); + assert.equal(canaryAllowsUser(blocked, 'tenant-a'), false); + assert.equal(rolloutBlockCode(blocked, 'tenant-a', 'skip_slot'), 'rollout_canary_tenant_blocked'); + + const allowlisted = policy({ canaryUserIds: new Set(['tenant-a']) }); + assert.equal(canaryAllowsUser(allowlisted, 'tenant-a'), true); + assert.equal(canaryAllowsUser(allowlisted, 'tenant-b'), false); + assert.deepEqual(runnableJobKinds(allowlisted), ['skip_slot', 'release_slot']); + assert.equal(rolloutBlockCode(allowlisted, 'tenant-a', 'refresh_queue'), 'rollout_generation_disabled'); + assert.equal(rolloutBlockCode(allowlisted, 'tenant-a', 'publish_now'), 'rollout_provider_dispatch_disabled'); + assert.equal(rolloutBlockCode(allowlisted, 'tenant-a', 'skip_slot'), undefined); + + const generationOnly = policy({ + canaryUserIds: new Set(['tenant-a']), + generationEnabled: true, + }); + assert.deepEqual(runnableJobKinds(generationOnly), [ + 'fetch_sources', + 'refresh_queue', + 'skip_slot', + 'release_slot', + ]); + assert.equal(rolloutBlockCode(generationOnly, 'tenant-a', 'refresh_queue'), undefined); + assert.equal(rolloutBlockCode(generationOnly, 'tenant-a', 'publish_now'), 'rollout_provider_dispatch_disabled'); + + const dispatchOnly = policy({ + canaryUserIds: new Set(['tenant-a']), + providerDispatchEnabled: true, + }); + assert.equal(rolloutBlockCode(dispatchOnly, 'tenant-a', 'publish_now'), undefined); + assert.equal(rolloutBlockCode(dispatchOnly, 'tenant-a', 'refresh_queue'), 'rollout_generation_disabled'); + + const unrestricted = policy({ + canaryRequired: false, + canaryUserIds: new Set(), + generationEnabled: true, + providerDispatchEnabled: true, + }); + assert.equal(allowedCanaryUserIds(unrestricted), undefined); + assert.equal(canaryAllowsUser(unrestricted, 'any-tenant'), true); + + const root = path.resolve(__dirname, '../..'); + const wrangler = fs.readFileSync(path.join(root, 'wrangler.toml'), 'utf8'); + const productionConfig = wrangler.split('[env.collector_staging]')[0]; + assert.match(productionConfig, /SUPABASE_WORKER_BATCH_SIZE = "1"/); + assert.match(productionConfig, /SUPABASE_WORKER_CANARY_REQUIRED = "true"/); + assert.match(productionConfig, /SUPABASE_WORKER_GENERATION_ENABLED = "false"/); + assert.match(productionConfig, /SUPABASE_PROVIDER_DISPATCH_ENABLED = "false"/); + assert.match(productionConfig, /DAILY_INVENTORY_PLANNER_ENABLED = "false"/); + assert.match(productionConfig, /REDDIT_CONNECTOR_ENABLED = "false"/); + assert.match(productionConfig, /required = \["SUPABASE_WORKER_CANARY_USER_IDS"\]/); + + const workflow = fs.readFileSync( + path.join(root, '.github/workflows/deploy-cloudflare-worker.yml'), + 'utf8' + ); + assert.match(workflow, /SUPABASE_WORKER_CANARY_USER_IDS/); + assert.match(workflow, /--secrets-file/); + + console.log('Canary rollout policy tests passed.'); +} + +main(); diff --git a/test/cloudflare-health.test.ts b/test/cloudflare-health.test.ts index dcd1e7b..fd95f3a 100644 --- a/test/cloudflare-health.test.ts +++ b/test/cloudflare-health.test.ts @@ -63,6 +63,36 @@ async function main(): Promise { assert.equal(capabilities.x.state, 'tenant_scoped'); }); + await test('health reports rollout controls without exposing tenant identifiers', async () => { + const body = await health({ + SUPABASE_WORKER_CANARY_REQUIRED: 'true', + SUPABASE_WORKER_CANARY_USER_IDS: '11111111-1111-4111-8111-111111111111,22222222-2222-4222-8222-222222222222,11111111-1111-4111-8111-111111111111', + SUPABASE_WORKER_GENERATION_ENABLED: 'false', + SUPABASE_PROVIDER_DISPATCH_ENABLED: 'false', + }); + + assert.deepEqual(body.rollout, { + tenantScope: 'allowlisted', + canaryRequired: true, + allowedTenantCount: 2, + generationEnabled: false, + providerDispatchEnabled: false, + }); + assert.equal(JSON.stringify(body).includes('11111111-1111-4111-8111-111111111111'), false); + }); + + await test('required canary with no allowlist is visibly blocked', async () => { + const body = await health({ + NODE_ENV: 'production', + SUPABASE_WORKER_CANARY_USER_IDS: '', + }); + assert.equal(body.rollout.tenantScope, 'blocked_empty_allowlist'); + assert.equal(body.rollout.canaryRequired, true); + assert.equal(body.rollout.allowedTenantCount, 0); + assert.equal(body.rollout.generationEnabled, false); + assert.equal(body.rollout.providerDispatchEnabled, false); + }); + await test('non-SHA Cloudflare version tags are not reported as canonical Git SHAs', async () => { const body = await health({ CF_VERSION_METADATA: { diff --git a/wrangler.toml b/wrangler.toml index d1bb994..0c1242a 100644 --- a/wrangler.toml +++ b/wrangler.toml @@ -9,10 +9,16 @@ crons = ["*/1 * * * *"] [version_metadata] binding = "CF_VERSION_METADATA" +[secrets] +required = ["SUPABASE_WORKER_CANARY_USER_IDS"] + [vars] NODE_ENV = "production" -SUPABASE_WORKER_BATCH_SIZE = "5" -DAILY_INVENTORY_PLANNER_ENABLED = "true" +SUPABASE_WORKER_BATCH_SIZE = "1" +SUPABASE_WORKER_CANARY_REQUIRED = "true" +SUPABASE_WORKER_GENERATION_ENABLED = "false" +SUPABASE_PROVIDER_DISPATCH_ENABLED = "false" +DAILY_INVENTORY_PLANNER_ENABLED = "false" DAILY_INVENTORY_PLANNER_START_LOCAL_DATE = "2026-07-14" HTTP_TIMEOUT_MS = "45000" OPENAI_IMAGE_MODEL = "gpt-image-2" @@ -21,7 +27,7 @@ CLOUDINARY_FOLDER = "social-agent/instagram" COLLECTOR_INGEST_ENABLED = "false" COLLECTOR_INGEST_WRITE_ENABLED = "false" COLLECTOR_INGEST_ENV = "production" -REDDIT_CONNECTOR_ENABLED = "true" +REDDIT_CONNECTOR_ENABLED = "false" REDDIT_CONNECTOR_MAX_POSTS_PER_RUN = "2" REDDIT_CONNECTOR_PAIRING_TTL_SECONDS = "600" APP_ALLOWED_ORIGINS = "https://www.oneclickpostfactory.com,https://oneclickpostfactory.woeinvests.workers.dev,https://oneclickwebsitefactory.tailwaggingwebdesign.com" From 0cbcf69c655ae9fab9e1d0007470f30856ebc5d1 Mon Sep 17 00:00:00 2001 From: Ayobami Haastrup <47716486+AyobamiH@users.noreply.github.com> Date: Tue, 8 Sep 2026 11:59:15 -0700 Subject: [PATCH 3/5] Update production config fixtures for canary controls --- test/security-hardening.test.ts | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/test/security-hardening.test.ts b/test/security-hardening.test.ts index e74cea0..779efde 100644 --- a/test/security-hardening.test.ts +++ b/test/security-hardening.test.ts @@ -441,6 +441,10 @@ async function main(): Promise { DAILY_INVENTORY_PLANNER_ENABLED: false, DAILY_INVENTORY_PLANNER_START_LOCAL_DATE: '', SUPABASE_WORKER_BATCH_SIZE: 10, + SUPABASE_WORKER_CANARY_REQUIRED: true, + SUPABASE_WORKER_CANARY_USER_IDS: new Set(), + SUPABASE_WORKER_GENERATION_ENABLED: false, + SUPABASE_PROVIDER_DISPATCH_ENABLED: false, }); assert.ok(issues.some(issue => issue.includes('COOKIE_SECURE'))); @@ -514,6 +518,10 @@ async function main(): Promise { DAILY_INVENTORY_PLANNER_ENABLED: false, DAILY_INVENTORY_PLANNER_START_LOCAL_DATE: '', SUPABASE_WORKER_BATCH_SIZE: 10, + SUPABASE_WORKER_CANARY_REQUIRED: true, + SUPABASE_WORKER_CANARY_USER_IDS: new Set(), + SUPABASE_WORKER_GENERATION_ENABLED: false, + SUPABASE_PROVIDER_DISPATCH_ENABLED: false, }); assert.ok(issues.some(issue => issue.includes('SUPABASE_SERVICE_ROLE_KEY'))); From a208f6c7e8201afe439c1b838410fd8461d22f1e Mon Sep 17 00:00:00 2001 From: Ayobami Haastrup <47716486+AyobamiH@users.noreply.github.com> Date: Tue, 8 Sep 2026 12:00:35 -0700 Subject: [PATCH 4/5] Assert tenant-filtered publication recovery --- test/publish-summary.test.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/test/publish-summary.test.ts b/test/publish-summary.test.ts index ddd042b..d341ff2 100644 --- a/test/publish-summary.test.ts +++ b/test/publish-summary.test.ts @@ -61,7 +61,7 @@ test('production stale publication recovery has no log/status reconciliation fal assert.ok(worker.includes( 'return reconcilePublication({ userId: job.user_id, queueItemId: row.id, platform: row.platform });' )); - assert.ok(worker.includes('await recoverStalePublications(now.getTime());')); + assert.ok(worker.includes('await recoverStalePublications(now.getTime(), rolloutAllowedUserIds());')); }); test('unknown publication guidance never authorises a blind retry', () => { From 263c2af33415538a6ea659ca7277e57425801de4 Mon Sep 17 00:00:00 2001 From: Ayobami Haastrup <47716486+AyobamiH@users.noreply.github.com> Date: Tue, 8 Sep 2026 12:01:54 -0700 Subject: [PATCH 5/5] Update recovery orchestration contract test --- test/publication-executor.test.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/test/publication-executor.test.ts b/test/publication-executor.test.ts index 7426388..89dcb9a 100644 --- a/test/publication-executor.test.ts +++ b/test/publication-executor.test.ts @@ -228,7 +228,7 @@ async function main() { assert.ok(!publish.includes("supabaseInsert")); assert.ok(!publish.includes("status: 'failed'")); assert.ok(!worker.includes('findPublishHistoryForQueueItem')); - assert.ok(worker.includes('await recoverStalePublications(now.getTime())')); + assert.ok(worker.includes('await recoverStalePublications(now.getTime(), rolloutAllowedUserIds())')); }); } main().catch(error => { console.error(error); process.exitCode = 1; });