From 866c2bfd0fe6882be16f12ab5c51d00f5583be2d Mon Sep 17 00:00:00 2001 From: Ayobami Haastrup <47716486+AyobamiH@users.noreply.github.com> Date: Thu, 1 Oct 2026 15:48:02 +0100 Subject: [PATCH 1/2] Prepare exact M3 connection-generation consumer from reviewed canonical source Verify content-addressed source transport, full Worker CI and no production credentials. No deployment or provider activation. The source-only preparer will be removed from the actual implementation candidate before integration. --- .github/workflows/prepare-m3-connections.yml | 66 ++++++++++++++++++++ 1 file changed, 66 insertions(+) create mode 100644 .github/workflows/prepare-m3-connections.yml diff --git a/.github/workflows/prepare-m3-connections.yml b/.github/workflows/prepare-m3-connections.yml new file mode 100644 index 0000000..8d6843c --- /dev/null +++ b/.github/workflows/prepare-m3-connections.yml @@ -0,0 +1,66 @@ +name: Prepare and prove M3 connection consumer +on: + push: + branches: [codex/m3-connections-20261001] + paths: ['.github/workflows/prepare-m3-connections.yml'] +permissions: + contents: read +jobs: + prepare: + if: github.repository == 'OneClickPostFactory/social-agents' + runs-on: ubuntu-latest + timeout-minutes: 6 + permissions: + contents: write + steps: + - uses: actions/checkout@11d5960a326750d5838078e36cf38b85af677262 + with: + ref: '66b17ef976f842c254769a58efc60c17d5e6ab65' + persist-credentials: false + - uses: actions/setup-node@49933ea5288caeca8642d1e84afbd3f7d6820020 + with: + node-version: '24' + - name: Verify exact source transport and permitted patch paths + env: + GH_TOKEN: ${{ github.token }} + run: | + node --input-type=module <<'JS' + import assert from 'node:assert/strict'; + import {createHash} from 'node:crypto'; + import {inflateSync} from 'node:zlib'; + import {execFileSync} from 'node:child_process'; + import {writeFileSync} from 'node:fs'; + assert.equal(execFileSync('git',['rev-parse','HEAD'],{encoding:'utf8'}).trim(),'66b17ef976f842c254769a58efc60c17d5e6ab65'); + let encoded=''; + for(const sha of ['59f48a7f5576c6daa1c7e084ba8bec00d911e7b6','07c0bf67587fa675b5621872b79c800a3160acbe']){ + const r=await fetch(`https://api.github.com/repos/OneClickPostFactory/social-agents/git/blobs/${sha}`,{headers:{Authorization:`Bearer ${process.env.GH_TOKEN}`,Accept:'application/vnd.github+json'},redirect:'error',signal:AbortSignal.timeout(15000)}); + assert.equal(r.status,200);const v=await r.json();assert.equal(v.sha,sha);assert.equal(v.encoding,'base64'); + const b=Buffer.from(v.content,'base64');assert.equal(createHash('sha1').update(`blob ${b.length}\0`).update(b).digest('hex'),sha);encoded+=b.toString('utf8'); + } + assert.equal(encoded.length,18272);assert.match(encoded,/^[A-Za-z0-9+/]+={0,2}$/); + const patch=inflateSync(Buffer.from(encoded,'base64'),{maxOutputLength:131072}); + assert.equal(patch.length,49786);assert.equal(createHash('sha256').update(patch).digest('hex'),'e0bbd1a2d1f02220f94bc32623ce471f07838e384bd4929c5e3658c3a6ab8b8d'); + const allowed=['AGENTS.md','config.ts','docs/CONNECTION_GENERATIONS_V1.md','package.json','src/connection-http.ts','src/connection-lifecycle.ts','src/connection-runtime.ts','src/supabase-worker.ts','src/worker-readiness.ts','test/connection-runtime.test.ts','test/security-hardening.test.ts','wrangler.toml']; + const paths=[...patch.toString('utf8').matchAll(/^diff --git a\/(\S+) b\/(\S+)$/gm)].map(m=>{assert.equal(m[1],m[2]);return m[1]});assert.deepEqual(paths.sort(),allowed.sort()); + writeFileSync('/tmp/m3-worker.patch',patch,{mode:0o600});writeFileSync('/tmp/m3-worker-paths.json',JSON.stringify(allowed)); + execFileSync('git',['apply','--check','/tmp/m3-worker.patch'],{stdio:'pipe'});execFileSync('git',['apply','/tmp/m3-worker.patch'],{stdio:'pipe'}); + console.log('Exact connection consumer source applied; runtime credentials and production requests=0.'); + JS + - run: npm ci --ignore-scripts --no-audit + - name: Full Worker typing, regressions and executable smoke + run: | + npm run ci + git diff --check + - name: Persist exact tested source blobs without branch mutation + env: + GH_TOKEN: ${{ github.token }} + run: | + node --input-type=module <<'JS' + import assert from 'node:assert/strict';import {readFileSync} from 'node:fs';import {createHash} from 'node:crypto'; + for(const path of JSON.parse(readFileSync('/tmp/m3-worker-paths.json','utf8'))){ + const b=readFileSync(path);assert.ok(b.length<600000); + const sha=createHash('sha1').update(`blob ${b.length}\0`).update(b).digest('hex'); + const r=await fetch('https://api.github.com/repos/OneClickPostFactory/social-agents/git/blobs',{method:'POST',headers:{Authorization:`Bearer ${process.env.GH_TOKEN}`,Accept:'application/vnd.github+json','Content-Type':'application/json'},redirect:'error',signal:AbortSignal.timeout(15000),body:JSON.stringify({content:b.toString('utf8'),encoding:'utf-8'})}); + assert.equal(r.status,201);const v=await r.json();assert.equal(v.sha,sha);console.log(JSON.stringify({path,sha,source_base:'66b17ef976f842c254769a58efc60c17d5e6ab65'})); + } + JS From 2b905b3304e35c4366a37769a72200dd40e4e63e Mon Sep 17 00:00:00 2001 From: Ayobami Haastrup <47716486+AyobamiH@users.noreply.github.com> Date: Thu, 1 Oct 2026 15:52:23 +0100 Subject: [PATCH 2/2] Implement request-scoped credential generations and exact publication-account fencing Promote source tested in cloud run36879225624/job110426431012: full Worker CI including15 new connection runtime tests passed. Require the paired schema before managed snapshots; preserve stable CAS operation identity on response loss; reject mismatched queue destinations before new publication intent. Preserve generation/publishing disabled and all old provider boundaries. Paired native database and disconnect-race acceptance still required. No production deployment, secret export or social publication. --- .github/workflows/prepare-m3-connections.yml | 66 ------ AGENTS.md | 9 + config.ts | 2 + docs/CONNECTION_GENERATIONS_V1.md | 49 +++++ package.json | 2 +- src/connection-http.ts | 68 ++++++ src/connection-lifecycle.ts | 213 +++++++++++++++++++ src/connection-runtime.ts | 108 ++++++++++ src/supabase-worker.ts | 54 +++-- src/worker-readiness.ts | 9 +- test/connection-runtime.test.ts | 96 +++++++++ test/security-hardening.test.ts | 2 + wrangler.toml | 1 + 13 files changed, 596 insertions(+), 83 deletions(-) delete mode 100644 .github/workflows/prepare-m3-connections.yml create mode 100644 docs/CONNECTION_GENERATIONS_V1.md create mode 100644 src/connection-http.ts create mode 100644 src/connection-lifecycle.ts create mode 100644 src/connection-runtime.ts create mode 100644 test/connection-runtime.test.ts diff --git a/.github/workflows/prepare-m3-connections.yml b/.github/workflows/prepare-m3-connections.yml deleted file mode 100644 index 8d6843c..0000000 --- a/.github/workflows/prepare-m3-connections.yml +++ /dev/null @@ -1,66 +0,0 @@ -name: Prepare and prove M3 connection consumer -on: - push: - branches: [codex/m3-connections-20261001] - paths: ['.github/workflows/prepare-m3-connections.yml'] -permissions: - contents: read -jobs: - prepare: - if: github.repository == 'OneClickPostFactory/social-agents' - runs-on: ubuntu-latest - timeout-minutes: 6 - permissions: - contents: write - steps: - - uses: actions/checkout@11d5960a326750d5838078e36cf38b85af677262 - with: - ref: '66b17ef976f842c254769a58efc60c17d5e6ab65' - persist-credentials: false - - uses: actions/setup-node@49933ea5288caeca8642d1e84afbd3f7d6820020 - with: - node-version: '24' - - name: Verify exact source transport and permitted patch paths - env: - GH_TOKEN: ${{ github.token }} - run: | - node --input-type=module <<'JS' - import assert from 'node:assert/strict'; - import {createHash} from 'node:crypto'; - import {inflateSync} from 'node:zlib'; - import {execFileSync} from 'node:child_process'; - import {writeFileSync} from 'node:fs'; - assert.equal(execFileSync('git',['rev-parse','HEAD'],{encoding:'utf8'}).trim(),'66b17ef976f842c254769a58efc60c17d5e6ab65'); - let encoded=''; - for(const sha of ['59f48a7f5576c6daa1c7e084ba8bec00d911e7b6','07c0bf67587fa675b5621872b79c800a3160acbe']){ - const r=await fetch(`https://api.github.com/repos/OneClickPostFactory/social-agents/git/blobs/${sha}`,{headers:{Authorization:`Bearer ${process.env.GH_TOKEN}`,Accept:'application/vnd.github+json'},redirect:'error',signal:AbortSignal.timeout(15000)}); - assert.equal(r.status,200);const v=await r.json();assert.equal(v.sha,sha);assert.equal(v.encoding,'base64'); - const b=Buffer.from(v.content,'base64');assert.equal(createHash('sha1').update(`blob ${b.length}\0`).update(b).digest('hex'),sha);encoded+=b.toString('utf8'); - } - assert.equal(encoded.length,18272);assert.match(encoded,/^[A-Za-z0-9+/]+={0,2}$/); - const patch=inflateSync(Buffer.from(encoded,'base64'),{maxOutputLength:131072}); - assert.equal(patch.length,49786);assert.equal(createHash('sha256').update(patch).digest('hex'),'e0bbd1a2d1f02220f94bc32623ce471f07838e384bd4929c5e3658c3a6ab8b8d'); - const allowed=['AGENTS.md','config.ts','docs/CONNECTION_GENERATIONS_V1.md','package.json','src/connection-http.ts','src/connection-lifecycle.ts','src/connection-runtime.ts','src/supabase-worker.ts','src/worker-readiness.ts','test/connection-runtime.test.ts','test/security-hardening.test.ts','wrangler.toml']; - const paths=[...patch.toString('utf8').matchAll(/^diff --git a\/(\S+) b\/(\S+)$/gm)].map(m=>{assert.equal(m[1],m[2]);return m[1]});assert.deepEqual(paths.sort(),allowed.sort()); - writeFileSync('/tmp/m3-worker.patch',patch,{mode:0o600});writeFileSync('/tmp/m3-worker-paths.json',JSON.stringify(allowed)); - execFileSync('git',['apply','--check','/tmp/m3-worker.patch'],{stdio:'pipe'});execFileSync('git',['apply','/tmp/m3-worker.patch'],{stdio:'pipe'}); - console.log('Exact connection consumer source applied; runtime credentials and production requests=0.'); - JS - - run: npm ci --ignore-scripts --no-audit - - name: Full Worker typing, regressions and executable smoke - run: | - npm run ci - git diff --check - - name: Persist exact tested source blobs without branch mutation - env: - GH_TOKEN: ${{ github.token }} - run: | - node --input-type=module <<'JS' - import assert from 'node:assert/strict';import {readFileSync} from 'node:fs';import {createHash} from 'node:crypto'; - for(const path of JSON.parse(readFileSync('/tmp/m3-worker-paths.json','utf8'))){ - const b=readFileSync(path);assert.ok(b.length<600000); - const sha=createHash('sha1').update(`blob ${b.length}\0`).update(b).digest('hex'); - const r=await fetch('https://api.github.com/repos/OneClickPostFactory/social-agents/git/blobs',{method:'POST',headers:{Authorization:`Bearer ${process.env.GH_TOKEN}`,Accept:'application/vnd.github+json','Content-Type':'application/json'},redirect:'error',signal:AbortSignal.timeout(15000),body:JSON.stringify({content:b.toString('utf8'),encoding:'utf-8'})}); - assert.equal(r.status,201);const v=await r.json();assert.equal(v.sha,sha);console.log(JSON.stringify({path,sha,source_base:'66b17ef976f842c254769a58efc60c17d5e6ab65'})); - } - JS diff --git a/AGENTS.md b/AGENTS.md index a60ad34..b370126 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -301,3 +301,12 @@ For Threads-specific work: For X-specific work: - `src/x.ts` + +## M3.4 paired connection lifecycle + +Read `docs/CONNECTION_GENERATIONS_V1.md` before changing SaaS credential writes, +refresh/callback logic or publication account selection. The new guard is default +off and requires the matching application/schema release. Captured generations +are not optional when enabled; never fall back to ordinary writes after a guarded +failure. UI source, schema, customer accounts and original local post-once are +not owned by this repository. diff --git a/config.ts b/config.ts index b8088bf..e750aed 100644 --- a/config.ts +++ b/config.ts @@ -112,6 +112,7 @@ export interface AppConfig { SUPABASE_WORKER_CANARY_REQUIRED: boolean; SUPABASE_WORKER_CANARY_USER_IDS: Set; SUPABASE_WORKER_GENERATION_ENABLED: boolean; + CONNECTION_LIFECYCLE_ENABLED: boolean; SUPABASE_PROVIDER_DISPATCH_ENABLED: boolean; DAILY_INVENTORY_PLANNER_ENABLED: boolean; DAILY_INVENTORY_PLANNER_START_LOCAL_DATE: string; @@ -249,6 +250,7 @@ function buildBaseConfig(): AppConfig { (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, ''), + CONNECTION_LIFECYCLE_ENABLED: process.env.CONNECTION_LIFECYCLE_ENABLED === 'true', SUPABASE_WORKER_GENERATION_ENABLED: parseBooleanEnv( process.env.SUPABASE_WORKER_GENERATION_ENABLED, (process.env.NODE_ENV || 'development') !== 'production' diff --git a/docs/CONNECTION_GENERATIONS_V1.md b/docs/CONNECTION_GENERATIONS_V1.md new file mode 100644 index 0000000..92ac91b --- /dev/null +++ b/docs/CONNECTION_GENERATIONS_V1.md @@ -0,0 +1,49 @@ +# M3.4: paired connection ownership, default off + +This publisher is paired with the app/schema repository's new additive +`connection-generations-v1` contract, migration `20261001154000`. +`CONNECTION_LIFECYCLE_ENABLED` must be exactly `true`; default is false. Do not +activate one consumer before the other consumer and schema are accepted. + +The trusted runtime loads encrypted credentials and provider generation/revision +in one atomic database snapshot. Each tenant execution owns an AsyncLocalStorage +session; callbacks never consult another tenant's mutable connection guard. Refresh +and verification writes compare the captured owner, generation and credential +revision and the database's ciphertext fingerprint. A lost reply can replay only +the identical operation/ciphertext, not construct a different refresh intent. +Replacing or disconnecting advances the owner generation; a late callback or +refresh cannot restore that session. Raw service-role writes to adopted credential +groups are refused. Administrative break-glass changes are detected by fingerprint +comparison and require new verification; this is not a claim of defeating DB owners. + +The existing X callback is implemented in the web repository with PKCE, one-use +state and generation-bound commit. LinkedIn read verification uses its supported +userinfo endpoint and checks the saved member URN. This proves an account identifier, +not a person's real-world identity or publishing permission. Required provider +scopes and actual owner account acceptance remain release gates. + +Queue destination bindings cover the saved draft/schedule hash, provider account +and generation. Unbound new work fails before creating an intent; the owner can +verify and review the destination in the app. The database also rechecks the +binding at begin-dispatch, preventing a disconnect that wins that boundary from +being ignored. Already-dispatched/unknown attempts retain their exact receipts +and cannot be reassigned or blindly resent. An external request already authorised +and in flight cannot be withdrawn by a later local disconnect; provider revocation +and remote deletion are separate operations, not promises made by this change. + +The pure `connection-lifecycle.ts` and `connection-http.ts` files are copied from +the app's shared contract and compared byte-for-byte in private paired CI. No +private app source or production data is copied to this public repository. +Native paired tests use explicit synthetic accounts, actual PostgreSQL and +intercepted external transport, never paid model requests or live posts. + +Existing retired Threads/Instagram publishers remain disabled, Facebook remains +paused, and generation/publishing switches are not changed. Managed OpenAI +credentials do not fall back to a global key after the owner clears their key. +The later M4 budget/content and M5 provider acceptance work are separate. + +Sources checked 1 October 2026: +- OAuth security: https://www.rfc-editor.org/rfc/rfc9700.html +- PostgreSQL row locking: https://www.postgresql.org/docs/17/explicit-locking.html +- X PKCE and refresh permissions: https://docs.x.com/fundamentals/authentication/oauth-2-0/authorization-code +- LinkedIn account endpoint/scopes: https://learn.microsoft.com/en-us/linkedin/consumer/integrations/self-serve/sign-in-with-linkedin-v2 diff --git a/package.json b/package.json index 320bbd4..db4c6df 100644 --- a/package.json +++ b/package.json @@ -6,7 +6,7 @@ "scripts": { "build": "node --import 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/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 && node dist/test/http-cancellation.test.js && node dist/test/scheduler-isolation.test.js && node dist/test/agent-jobs.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 && node dist/test/http-cancellation.test.js && node dist/test/scheduler-isolation.test.js && node dist/test/agent-jobs.test.js && node dist/test/connection-runtime.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/connection-http.ts b/src/connection-http.ts new file mode 100644 index 0000000..af88278 --- /dev/null +++ b/src/connection-http.ts @@ -0,0 +1,68 @@ +// Narrow credential-bearing transport shared with the canonical Worker. No +// redirect, unbounded body, implicit retry, arbitrary endpoint or raw-error log. +const endpoints: Record = { + "https://api.x.com/2/oauth2/token": "POST", + "https://api.x.com/2/users/me": "GET", + "https://api.linkedin.com/v2/userinfo": "GET", + "https://graph.threads.net/v1.0/me?fields=id": "GET", +}; +export async function connectionJson( + url: string, + init: RequestInit, + fetchImpl: typeof fetch = fetch, +): Promise> { + if (endpoints[url] !== (init.method || "GET")) throw Error("connection_endpoint_rejected"); + if (init.signal?.aborted) throw Error("connection_request_cancelled"); + const controller = new AbortController(); + let stop: (error: Error) => void = () => {}; + const stopped = new Promise((_, reject) => { + stop = reject; + }); + const cancel = () => { + controller.abort(); + stop(Error("connection_request_unverified")); + }; + const timer = setTimeout(cancel, 10000); + init.signal?.addEventListener("abort", cancel, { once: true }); + let reader: ReadableStreamDefaultReader | undefined; + try { + if (init.signal?.aborted) throw Error("connection_request_cancelled"); + const r = await Promise.race([ + fetchImpl(url, { ...init, redirect: "manual", signal: controller.signal }), + stopped, + ]); + if (r.status < 200 || r.status >= 300 || !r.body) { + void r.body?.cancel().catch(() => {}); + throw Error("connection_request_unverified"); + } + reader = r.body.getReader(); + let size = 0, + text = ""; + const decoder = new TextDecoder("utf-8", { fatal: true }); + for (;;) { + const next = await Promise.race([reader.read(), stopped]); + if (next.done) break; + size += next.value.byteLength; + if (size > 65536) throw Error("connection_response_unverified"); + text += decoder.decode(next.value, { stream: true }); + } + if (controller.signal.aborted) throw Error("connection_request_unverified"); + const value: unknown = JSON.parse(text + decoder.decode()); + if (!value || typeof value !== "object" || Array.isArray(value)) + throw Error("connection_response_unverified"); + return value as Record; + } catch { + // No provider body, token-bearing URL, secret or private exception escapes. + throw Error( + "Connection verification failed. Reconnect with the correct account and permissions.", + ); + } finally { + clearTimeout(timer); + init.signal?.removeEventListener("abort", cancel); + controller.abort(); + if (reader) { + void reader.cancel().catch(() => {}); + reader.releaseLock(); + } + } +} diff --git a/src/connection-lifecycle.ts b/src/connection-lifecycle.ts new file mode 100644 index 0000000..6649041 --- /dev/null +++ b/src/connection-lifecycle.ts @@ -0,0 +1,213 @@ +// Shared deterministic contract. Kept byte-identical in the canonical publisher +// as src/connection-lifecycle.ts; paired CI verifies the source hash. +export const CONNECTION_PROVIDERS = ["openai", "x", "threads", "linkedin", "instagram"] as const; +export type ConnectionProvider = (typeof CONNECTION_PROVIDERS)[number]; +export type ConnectionGuard = { + provider: ConnectionProvider; + generation: number; + revision: number; + state: "stored_not_verified" | "verified" | "disconnected" | "connecting" | "needs_reconnect"; + account_id: string | null; + verified_at: string | null; +}; +export type ConnectionSnapshot = { + schema: "ocpf.connection-snapshot.v1"; + user_id: string; + credentials: Record; + connections: ConnectionGuard[]; +}; +export type ConnectionRpc = ( + name: string, + params: Record, + retrySafe: boolean, +) => Promise; +export type ConnectionMutationKind = + | "replace" + | "disconnect" + | "refresh" + | "verify" + | "failure" + | "callback"; +const uuid = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i; +const account = /^[A-Za-z0-9:_-]{1,256}$/; +export const CREDENTIAL_FIELDS = [ + "openai_api_key", + "threads_token", + "instagram_token", + "instagram_account_id", + "facebook_page_id", + "facebook_page_access_token", + "linkedin_token", + "linkedin_person_urn", + "linkedin_refresh_token", + "linkedin_client_id", + "linkedin_client_secret", + "meta_access_token", + "x_client_id", + "x_client_secret", + "x_oauth2_access_token", + "x_oauth2_refresh_token", +] as const; +export function connectionProvider(field: string): ConnectionProvider { + if (!(CREDENTIAL_FIELDS as readonly string[]).includes(field)) + throw Error("connection_field_invalid"); + if (field.startsWith("x_")) return "x"; + if (field.startsWith("linkedin_")) return "linkedin"; + if (field === "threads_token") return "threads"; + if (field === "openai_api_key") return "openai"; + return "instagram"; +} +function record(value: unknown): Record { + if (!value || typeof value !== "object" || Array.isArray(value)) + throw Error("connection_response_invalid"); + return value as Record; +} +function counter(value: unknown, minimum: number): value is number { + return ( + typeof value === "number" && + Number.isSafeInteger(value) && + value >= minimum && + value < Number.MAX_SAFE_INTEGER + ); +} +export function parseConnectionGuard(value: unknown): ConnectionGuard { + const r = record(value); + if ( + !CONNECTION_PROVIDERS.includes(r.provider as ConnectionProvider) || + !counter(r.generation, 1) || + !counter(r.revision, 0) || + !["stored_not_verified", "verified", "disconnected", "connecting", "needs_reconnect"].includes( + String(r.state), + ) || + !(r.account_id === null || (typeof r.account_id === "string" && account.test(r.account_id))) || + !( + r.verified_at === null || + (typeof r.verified_at === "string" && Number.isFinite(Date.parse(r.verified_at))) + ) || + (r.state === "verified" && (!r.account_id || !r.verified_at)) + ) + throw Error("connection_response_invalid"); + // Only fixed safe fields ever reach a browser; reject unknown data rather than + // spreading an arbitrary provider/database response into status objects. + return { + provider: r.provider as ConnectionProvider, + generation: r.generation, + revision: r.revision, + state: r.state as ConnectionGuard["state"], + account_id: r.account_id as string | null, + verified_at: r.verified_at as string | null, + }; +} +export async function captureConnections( + rpc: ConnectionRpc, + userId: string, +): Promise { + if (!uuid.test(userId)) throw Error("connection_identity_required"); + const raw = record(await rpc("capture_connection_snapshot", { p_user_id: userId }, true)); + if ( + raw.schema !== "ocpf.connection-snapshot.v1" || + raw.user_id !== userId || + !Array.isArray(raw.connections) || + raw.connections.length !== CONNECTION_PROVIDERS.length + ) + throw Error("connection_snapshot_invalid"); + const credentials = record(raw.credentials); + if (credentials.user_id !== userId) throw Error("connection_snapshot_invalid"); + const connections = raw.connections.map(parseConnectionGuard); + if (new Set(connections.map((c) => c.provider)).size !== CONNECTION_PROVIDERS.length) + throw Error("connection_snapshot_invalid"); + return { schema: "ocpf.connection-snapshot.v1", user_id: userId, credentials, connections }; +} +export async function mutateConnection( + rpc: ConnectionRpc, + userId: string, + guard: ConnectionGuard, + operationId: string, + kind: ConnectionMutationKind, + patch: Record, + accountId: string | null = null, +): Promise { + parseConnectionGuard(guard); + if (!uuid.test(userId) || !uuid.test(operationId)) throw Error("connection_identity_required"); + const expectedGeneration = guard.generation + (["replace", "disconnect"].includes(kind) ? 1 : 0); + const value = record( + await rpc( + "mutate_connection", + { + p_user_id: userId, + p_provider: guard.provider, + p_generation: guard.generation, + p_revision: guard.revision, + p_operation_id: operationId, + p_kind: kind, + p_patch: patch, + p_account_id: accountId, + }, + true, + ), + ); + const next = parseConnectionGuard(value.connection); + if ( + value.schema !== "ocpf.connection-mutation.v1" || + value.user_id !== userId || + value.operation_id !== operationId || + next.provider !== guard.provider || + next.generation !== expectedGeneration || + next.revision !== guard.revision + 1 || + (kind === "disconnect" && next.state !== "disconnected") || + (["callback", "verify"].includes(kind) && + (next.state !== "verified" || next.account_id !== accountId)) || + (kind === "refresh" && + (next.state !== "stored_not_verified" || next.account_id !== guard.account_id)) || + (kind === "replace" && (next.state !== "stored_not_verified" || next.account_id !== null)) || + (kind === "failure" && next.state !== "needs_reconnect") + ) + throw Error("connection_mutation_unverified"); + return next; +} +export class ConnectionSession { + private readonly guards = new Map(); + private busy = new Set(); + private readonly rpc: ConnectionRpc; + readonly snapshot: ConnectionSnapshot; + constructor(rpc: ConnectionRpc, snapshot: ConnectionSnapshot) { + this.rpc = rpc; + this.snapshot = snapshot; + for (const c of snapshot.connections) this.guards.set(c.provider, Object.freeze({ ...c })); + } + guard(provider: ConnectionProvider): ConnectionGuard { + const result = this.guards.get(provider); + if (!result) throw Error("connection_snapshot_invalid"); + return result; + } + async apply( + provider: ConnectionProvider, + kind: "refresh" | "verify" | "failure", + patch: Record, + accountId: string | null = null, + ): Promise { + if (this.busy.has(provider)) throw Error("connection_local_mutation_in_progress"); + this.busy.add(provider); + try { + const next = await mutateConnection( + this.rpc, + this.snapshot.user_id, + this.guard(provider), + crypto.randomUUID(), + kind, + patch, + accountId, + ); + this.guards.set(provider, Object.freeze(next)); + return next; + } finally { + this.busy.delete(provider); + } + } + assertAvailable(provider: ConnectionProvider): ConnectionGuard { + const guard = this.guard(provider); + if (["disconnected", "connecting"].includes(guard.state)) + throw Error("connection_not_available"); + return guard; + } +} diff --git a/src/connection-runtime.ts b/src/connection-runtime.ts new file mode 100644 index 0000000..9a8c708 --- /dev/null +++ b/src/connection-runtime.ts @@ -0,0 +1,108 @@ +import { AsyncLocalStorage } from 'node:async_hooks'; +import config from '../config'; +import { supabaseUpdate, type SupabaseMutationOptions } from './supabase-client'; +import { captureConnections, ConnectionSession, type ConnectionSnapshot, type ConnectionRpc, type ConnectionProvider } from './connection-lifecycle'; + +const scope = new AsyncLocalStorage(); +export const CONNECTION_CAPABILITIES = [ + 'connection-atomic-snapshot-v1', 'connection-owner-generation-v1', 'connection-refresh-cas-v1', + 'connection-legacy-write-denial-v1', 'connection-callback-generation-v1', 'connection-publication-binding-v1', +] as const; +// This narrow adapter intentionally does not inherit the old client's redirect, +// unbounded body or exception formatting. No automatic external-token retry. +export const connectionRpc: ConnectionRpc = async (name, params, retrySafe) => { + if (!['capture_connection_snapshot','mutate_connection','get_connection_generation_contract','check_queue_connection_binding'].includes(name)) + throw Error('connection_rpc_not_allowed'); + const origin = new URL(config.SUPABASE_URL); + const local = process.env.PUBLICATION_DATABASE_TEST === 'local-only' && ['127.0.0.1','localhost'].includes(origin.hostname); + if ((!local && origin.protocol !== 'https:') || origin.username || origin.password || origin.search || origin.hash || origin.pathname !== '/') + throw Error('connection_database_origin_invalid'); + const body = JSON.stringify(params); + if (body.length > 65536) throw Error('connection_rpc_input_invalid'); + const once = async (): Promise => { + const controller = new AbortController(); + let stop: (error: Error) => void = () => {}; + const cancelled = new Promise((_, reject) => { stop = reject; }); + const timer = setTimeout(() => { controller.abort(); stop(Error('connection_transport_uncertain')); }, 5000); + let reader: ReadableStreamDefaultReader | undefined; + try { + let response: Response; + try { + response = await Promise.race([fetch(`${origin.origin}/rest/v1/rpc/${name}`, { + method:'POST', redirect:'manual', signal:controller.signal, + headers:{Authorization:`Bearer ${config.SUPABASE_SERVICE_ROLE_KEY}`, apikey:config.SUPABASE_SERVICE_ROLE_KEY, 'Content-Type':'application/json'}, body, + }), cancelled]); + } catch { throw Error('connection_transport_uncertain'); } + if (response.status !== 200 || !response.body) { + void response.body?.cancel().catch(() => {}); + throw Error('connection_rpc_rejected'); + } + reader = response.body.getReader(); + let bytes = 0, text = ''; const decoder = new TextDecoder('utf-8',{fatal:true}); + for (;;) { + const next = await Promise.race([reader.read(),cancelled]); if (next.done) break; + bytes += next.value.byteLength; if (bytes > 262144) throw Error('connection_response_unverified'); + text += decoder.decode(next.value,{stream:true}); + } + if (controller.signal.aborted) throw Error('connection_transport_uncertain'); + return JSON.parse(text + decoder.decode()) as unknown; + } finally { + clearTimeout(timer); controller.abort(); + if (reader) { void reader.cancel().catch(() => {}); reader.releaseLock(); } + } + }; + try { return await once(); } catch (error) { + // Exact stable body is replayed at most once and only for a contract declared + // retry-safe and an unclassified/lost transport. 4xx or invalid JSON is not retried. + if (retrySafe && error instanceof Error && error.message === 'connection_transport_uncertain') { + try { return await once(); } catch { throw Error('connection_changed_or_unverified'); } + } + throw Error('connection_changed_or_unverified'); + } +}; +export async function loadConnectionSnapshot(userId: string): Promise { + if (!config.CONNECTION_LIFECYCLE_ENABLED) throw Error('connection_lifecycle_disabled'); + const c = await connectionRpc('get_connection_generation_contract', {}, false) as Record; + if (!c || c.contract !== 'connection-generations-v1' || c.migration !== '20261001154000' + || !Array.isArray(c.capabilities) || !CONNECTION_CAPABILITIES.every(x => (c.capabilities as unknown[]).includes(x))) + throw Error('connection_schema_unverified'); + return captureConnections(connectionRpc,userId); +} +export function withConnectionSnapshot(snapshot: ConnectionSnapshot, fn: () => Promise): Promise { + return scope.run(new ConnectionSession(connectionRpc,snapshot), fn); +} +export function currentConnectionSession(userId?: string): ConnectionSession { + const session = scope.getStore(); + if (!session || (userId !== undefined && session.snapshot.user_id !== userId)) throw Error('connection_scope_missing_or_wrong_owner'); + return session; +} +export function assertProviderAvailable(provider: ConnectionProvider): void { + if (config.CONNECTION_LIFECYCLE_ENABLED) currentConnectionSession().assertAvailable(provider); +} +export async function connectionCredentialUpdate(table: string, patch: Record, options: SupabaseMutationOptions): Promise { + if (table !== 'user_credentials') throw Error('connection_table_invalid'); + if (!config.CONNECTION_LIFECYCLE_ENABLED) return supabaseUpdate(table,patch,options); + const f=options.filters; + if (!f || f.length!==1 || f[0].column!=='user_id' || f[0].operator!=='eq' || typeof f[0].value!=='string') throw Error('connection_owner_required'); + const session=currentConnectionSession(f[0].value); + const prefixes=new Set(Object.keys(patch).map(k=>k.split('_')[0])); + if (prefixes.size!==1) throw Error('connection_mixed_provider_patch'); + const provider=Array.from(prefixes)[0] as ConnectionProvider; + if (!['x','threads','linkedin'].includes(provider)) throw Error('connection_provider_unsupported'); + session.assertAvailable(provider); + const encrypted=Object.keys(patch).some(k=>k.endsWith('_enc')); + const status=patch[`${provider}_verification_status`]; + const kind=encrypted?'refresh':status==='verified'?'verify':'failure'; + const accountId=kind==='verify' ? (provider==='linkedin'?config.LINKEDIN_PERSON_URN:patch[`${provider}_account_id`]) : null; + if (accountId!==null && typeof accountId!=='string') throw Error('connection_account_required'); + await session.apply(provider,kind,patch,accountId as string|null); + // All callers consume this only as confirmation; never fabricate a credential row. + return []; +} + +export async function assertQueueConnectionReady(userId:string,queueId:string):Promise { + if(!config.CONNECTION_LIFECYCLE_ENABLED)return; + const r=await connectionRpc('check_queue_connection_binding',{p_user_id:userId,p_queue_id:queueId},true) as Record; + if(!r || r.schema!=='ocpf.queue-connection-binding.v1'||r.user_id!==userId||r.queue_id!==queueId + || !['bound','existing_intent'].includes(String(r.state))) throw Error('connection_changed_review_destination_before_publish'); +} diff --git a/src/supabase-worker.ts b/src/supabase-worker.ts index 6870596..5b6f6d7 100644 --- a/src/supabase-worker.ts +++ b/src/supabase-worker.ts @@ -1,3 +1,6 @@ +import { connectionCredentialUpdate, loadConnectionSnapshot, withConnectionSnapshot, assertProviderAvailable, assertQueueConnectionReady } from './connection-runtime'; +import type { ConnectionSnapshot } from './connection-lifecycle'; +import { connectionJson } from './connection-http'; import { AgentJobOwnershipLostError, assertAgentJobsContract, assertActiveAgentJob, claimAgentJob, enqueueScheduledAgentJob, finishAgentJob, scheduledOperationKey, withAgentJobLease, type AgentJobRow } from './agent-jobs'; import * as crypto from 'node:crypto'; @@ -298,6 +301,7 @@ interface AngleRecordRow { } interface TenantContext { + connectionSnapshot?: ConnectionSnapshot; userId: string; settings: UserSettingsRow; credentials: TenantCredentials; @@ -1661,10 +1665,9 @@ async function loadTenantContext(userId: string): Promise { limit: 1, }))[0] || {}; - const credentialRow = (await supabaseSelect('user_credentials', { - select: '*', - filters: [{ column: 'user_id', operator: 'eq', value: userId }], - limit: 1, + const connectionSnapshot = config.CONNECTION_LIFECYCLE_ENABLED ? await loadConnectionSnapshot(userId) : undefined; + const credentialRow = connectionSnapshot ? connectionSnapshot.credentials as TenantCredentialRow : (await supabaseSelect('user_credentials', { + select: '*', filters: [{ column: 'user_id', operator: 'eq', value: userId }], limit: 1, }))[0]; const credentials = decryptTenantCredentials(credentialRow); @@ -1674,6 +1677,7 @@ async function loadTenantContext(userId: string): Promise { userId, settings, credentials, + connectionSnapshot, activePlatforms, }; } @@ -1759,7 +1763,7 @@ async function markXCredentialVerified( userId: string, verification: Awaited> ): Promise { - await supabaseUpdate('user_credentials', { + await connectionCredentialUpdate('user_credentials', { x_verified_at: nowIso(), x_last_verification_failed_at: null, x_verification_status: 'verified', @@ -1784,7 +1788,7 @@ async function markXCredentialFailure(userId: string, error: unknown): Promise { try { + assertProviderAvailable('x'); const verification = await x.verifyCredentials(); await markXCredentialVerified(userId, verification); return verification.authMode; @@ -1810,8 +1815,9 @@ async function verifyXCredentialForPublish(userId: string): Promise { try { + assertProviderAvailable('threads'); const preparation = await threads.prepareAccessTokenForPublish(); - await supabaseUpdate('user_credentials', { + await connectionCredentialUpdate('user_credentials', { threads_verified_at: nowIso(), threads_last_verification_failed_at: null, threads_verification_status: 'verified', @@ -1829,7 +1835,7 @@ async function prepareThreadsCredentialForPublish(userId: string): Promise const context = isPlatformPublishError(error) ? platformErrorContext(error) : { normalized_error_code: 'verification_failed', user_message: publicError(error) }; - await supabaseUpdate('user_credentials', { + await connectionCredentialUpdate('user_credentials', { threads_last_verification_failed_at: nowIso(), threads_verification_status: 'needs_reconnect', threads_verification_error: publicError(error), @@ -1845,6 +1851,7 @@ async function prepareThreadsCredentialForPublish(userId: string): Promise } async function refreshLinkedInCredentialForPublish(userId: string): Promise { + assertProviderAvailable('linkedin'); if (!linkedin.shouldRefreshAccessToken()) return; try { @@ -1853,7 +1860,7 @@ async function refreshLinkedInCredentialForPublish(userId: string): Promise { - await supabaseUpdate('user_credentials', { + await connectionCredentialUpdate('user_credentials', { linkedin_verified_at: nowIso(), linkedin_last_verification_failed_at: null, linkedin_verification_status: 'verified', @@ -2019,10 +2026,17 @@ function tenantCloudinaryFolder(baseFolder: string, userId: string): string { } async function withTenantRuntime(tenant: TenantContext, fn: () => Promise): Promise { + if (config.CONNECTION_LIFECYCLE_ENABLED) { + if (!tenant.connectionSnapshot) throw Error('connection_snapshot_required'); + return withConnectionSnapshot(tenant.connectionSnapshot, () => withTenantRuntimeImpl(tenant, fn)); + } + return withTenantRuntimeImpl(tenant, fn); +} +async function withTenantRuntimeImpl(tenant: TenantContext, fn: () => Promise): Promise { const previous = snapshotConfig(); const restoreThreadsTokenPersistence = threads.setTokenPersistence(async tokens => { const expiresAt = secondsFromNowIso(tokens.expiresIn); - await supabaseUpdate('user_credentials', { + await connectionCredentialUpdate('user_credentials', { threads_token_enc: encryptCredential(tokens.accessToken), ...(expiresAt ? { threads_expires_at: expiresAt } : {}), }, { @@ -2055,7 +2069,7 @@ async function withTenantRuntime(tenant: TenantContext, fn: () => Promise) if (refreshTokenExpiresAt) { patch.linkedin_refresh_token_expires_at = refreshTokenExpiresAt; } - await supabaseUpdate('user_credentials', patch, { + await connectionCredentialUpdate('user_credentials', patch, { filters: [{ column: 'user_id', operator: 'eq', value: tenant.userId }], }); await writeWorkerLog(tenant.userId, 'info', 'linkedin_oauth2_tokens_refreshed', { @@ -2081,7 +2095,7 @@ async function withTenantRuntime(tenant: TenantContext, fn: () => Promise) patch.x_expires_at = expiresAt; } - await supabaseUpdate('user_credentials', patch, { + await connectionCredentialUpdate('user_credentials', patch, { filters: [{ column: 'user_id', operator: 'eq', value: tenant.userId }], }); await writeWorkerLog(tenant.userId, 'info', 'x_oauth2_tokens_refreshed', { @@ -2091,7 +2105,7 @@ async function withTenantRuntime(tenant: TenantContext, fn: () => Promise) expiresAtUpdated: Boolean(expiresAt), }); }); - config.OPENAI_API_KEY = tenant.credentials.openaiApiKey || previous.OPENAI_API_KEY || ''; + config.OPENAI_API_KEY = tenant.credentials.openaiApiKey || (config.CONNECTION_LIFECYCLE_ENABLED ? '' : previous.OPENAI_API_KEY) || ''; config.OPENAI_MODEL = tenant.settings.ai_model || 'gpt-4o-mini'; config.OPENAI_IMAGE_MODEL = previous.OPENAI_IMAGE_MODEL || config.OPENAI_IMAGE_MODEL || 'gpt-image-2'; config.CLOUDINARY_FOLDER = tenantCloudinaryFolder(previous.CLOUDINARY_FOLDER || config.CLOUDINARY_FOLDER, tenant.userId); @@ -3594,6 +3608,7 @@ async function publishQueueRow(job: AgentJobRow, row: QueueItemRow, _settings: U // 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 () => { + await assertQueueConnectionReady(job.user_id,row.id); const execution = await executePublication({ userId: job.user_id, queueItemId: row.id, platform: row.platform, }, { @@ -3609,13 +3624,21 @@ async function publishQueueRow(job: AgentJobRow, row: QueueItemRow, _settings: U let providerAccountRef: string; if (payload.platform === 'x') { if (payload.text.trim().length > 280) throw new PublicationPreflightError('x_text_too_long'); - const verification = await x.verifyCredentials(); + assertProviderAvailable('x'); + const verification = await x.verifyCredentials(); + if (config.CONNECTION_LIFECYCLE_ENABLED) await markXCredentialVerified(job.user_id, verification); providerAccountRef = verification.accountId; } else { if (!config.LINKEDIN_TOKEN || !config.LINKEDIN_PERSON_URN) { throw new PublicationPreflightError('linkedin_not_connected'); } await refreshLinkedInCredentialForPublish(job.user_id); + if (config.CONNECTION_LIFECYCLE_ENABLED) { + const me = await connectionJson('https://api.linkedin.com/v2/userinfo', { headers:{Authorization:`Bearer ${config.LINKEDIN_TOKEN}`} }); + if (typeof me.sub !== 'string' || !/^[A-Za-z0-9_-]{1,200}$/.test(me.sub) + || `urn:li:person:${me.sub}` !== config.LINKEDIN_PERSON_URN) throw new PublicationPreflightError('linkedin_account_mismatch'); + await markLinkedInCredentialVerified(job.user_id); + } providerAccountRef = config.LINKEDIN_PERSON_URN; } // Recheck entitlement/pause/enablement immediately before the dispatch boundary. @@ -4871,6 +4894,7 @@ export function startSupabaseWorkerLoop(log = logger): { stop: () => void } | un } export const __test__ = { + loadTenantContext, withTenantRuntime, markXCredentialVerified, refreshLinkedInCredentialForPublish, extractSourceBankWithJobTimeout, assertQueueRevisionNotHeld, publishQueueRow, diff --git a/src/worker-readiness.ts b/src/worker-readiness.ts index c9e489b..3436fb2 100644 --- a/src/worker-readiness.ts +++ b/src/worker-readiness.ts @@ -4,6 +4,7 @@ import { timingSafeEqual } from 'node:crypto'; // provider clients or loggers here: liveness must never initialise execution. export interface ReadinessEnv { WORKER_TICK_TOKEN?: string; + CONNECTION_LIFECYCLE_ENABLED?: string; SUPABASE_URL?: string; SUPABASE_SERVICE_ROLE_KEY?: string; SUPABASE_SECRET_KEY?: string; @@ -34,6 +35,11 @@ export const READINESS_CONTRACTS = [ }, ] as const; +export const CONNECTION_READINESS_CONTRACT = { + name:'connections',rpc:'get_connection_generation_contract',contract:'connection-generations-v1', + capabilities:['connection-atomic-snapshot-v1','connection-owner-generation-v1','connection-refresh-cas-v1', + 'connection-legacy-write-denial-v1','connection-callback-generation-v1','connection-publication-binding-v1'], +} as const; const UUID = /^[a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12}$/i; const SHA = /^(?:[a-f0-9]{40}|[a-f0-9]{64})$/i; const HEADERS = { 'Cache-Control': 'no-store, private, max-age=0', Vary: 'Authorization, X-OCPF-Readiness-Nonce' }; @@ -125,7 +131,7 @@ export async function handleWorkerReadiness( request.signal.addEventListener('abort', abort, { once: true }); const timeout = setTimeout(abort, TIMEOUT_MS); try { - const dependencies = await Promise.all(READINESS_CONTRACTS.map(async expected => { + const dependencies = await Promise.all([...READINESS_CONTRACTS,...(env.CONNECTION_LIFECYCLE_ENABLED === 'true'?[CONNECTION_READINESS_CONTRACT]:[])].map(async expected => { let response: Response | undefined; try { // workerd rejects redirect:error before transport. manual plus the exact @@ -142,6 +148,7 @@ export async function handleWorkerReadiness( const capabilities = body.capabilities; let compatible = body.contract === expected.contract && Array.isArray(capabilities) && expected.capabilities.every(c => capabilities.includes(c)); + if (expected.name === 'connections') compatible = compatible && body.migration === '20261001154000'; if (expected.name === 'publication') compatible = compatible && body.migration === '20260907054000' && body.lock_order_migration === '20260913061000'; return { name: expected.name, state: compatible ? 'verified' : 'contract_mismatch' }; diff --git a/test/connection-runtime.test.ts b/test/connection-runtime.test.ts new file mode 100644 index 0000000..55965be --- /dev/null +++ b/test/connection-runtime.test.ts @@ -0,0 +1,96 @@ +import assert from 'node:assert/strict'; +import { test, afterEach } from 'node:test'; +import config from '../config'; +import { CONNECTION_PROVIDERS, type ConnectionGuard, type ConnectionSnapshot } from '../src/connection-lifecycle'; +import { connectionRpc, CONNECTION_CAPABILITIES, loadConnectionSnapshot, withConnectionSnapshot, currentConnectionSession, connectionCredentialUpdate, assertQueueConnectionReady } from '../src/connection-runtime'; +import { __test__ as worker } from '../src/supabase-worker'; +import { installScopedConfig, runWithRuntimeScope } from '../src/runtime-scope'; +import { handleWorkerReadiness, READINESS_CONTRACTS, CONNECTION_READINESS_CONTRACT } from '../src/worker-readiness'; +const user='11111111-1111-4111-8111-111111111111', other='22222222-2222-4222-8222-222222222222'; +const guard=(provider:ConnectionGuard['provider']):ConnectionGuard=>({provider,generation:2,revision:1,state:'verified',account_id:'12345',verified_at:'2026-01-01T00:00:00Z'}); +const snapshot=(id=user):ConnectionSnapshot=>({schema:'ocpf.connection-snapshot.v1',user_id:id,credentials:{user_id:id},connections:CONNECTION_PROVIDERS.map(guard)}); +const original=globalThis.fetch; +Object.assign(config,{SUPABASE_URL:'https://database.invalid',SUPABASE_SERVICE_ROLE_KEY:'private-test-key',CREDENTIAL_ENCRYPTION_KEY:'fixture',CONNECTION_LIFECYCLE_ENABLED:true}); +installScopedConfig(config); +afterEach(()=>{globalThis.fetch=original;config.CONNECTION_LIFECYCLE_ENABLED=true;}); +function mock(handler:(name:string,body:any,init:RequestInit)=>Response|Promise) { + const calls:Array<{name:string;body:any;init:RequestInit}>=[]; + globalThis.fetch=async(input,init)=>{ + const url=new URL(String(input));assert.equal(url.origin,'https://database.invalid'); + const body=JSON.parse(String(init?.body||'{}'));const name=url.pathname.split('/').pop()!; + calls.push({name,body,init:init!});return handler(name,body,init!); + }; + return calls; +} +function response(p:any) { + return Response.json({schema:'ocpf.connection-mutation.v1',operation_id:p.p_operation_id,user_id:p.p_user_id, + connection:{...guard(p.p_provider),revision:p.p_revision+1,generation:p.p_generation,state:p.p_kind==='verify'?'verified':p.p_kind==='failure'?'needs_reconnect':'stored_not_verified',account_id:p.p_account_id||'12345'}}); +} +test('native capture requires the exact full additive schema first',async()=>{ + const calls=mock(n=>n==='get_connection_generation_contract'?Response.json({contract:'connection-generations-v1',migration:'20261001154000',capabilities:CONNECTION_CAPABILITIES}):Response.json(snapshot())); + assert.equal((await loadConnectionSnapshot(user)).user_id,user);assert.deepEqual(calls.map(c=>c.name),['get_connection_generation_contract','capture_connection_snapshot']); + calls.forEach(c=>{assert.equal(c.init.redirect,'manual');assert.ok(c.init.signal);}); +}); +test('missing schema capability blocks before adopting a tenant',async()=>{ + const calls=mock(()=>Response.json({contract:'connection-generations-v1',migration:'20261001154000',capabilities:CONNECTION_CAPABILITIES.slice(1)})); + await assert.rejects(loadConnectionSnapshot(user),/schema_unverified/);assert.equal(calls.length,1); +}); +test('missing or different tenant scope cannot persist any credential',async()=>{ + const calls=mock(()=>{throw Error('must not read');}); + await assert.rejects(connectionCredentialUpdate('user_credentials',{x_oauth2_access_token_enc:'opaque'},{filters:[{column:'user_id',operator:'eq',value:user}]}),/scope_missing/); + await withConnectionSnapshot(snapshot(),async()=>await assert.rejects(connectionCredentialUpdate('user_credentials',{x_oauth2_access_token_enc:'opaque'},{filters:[{column:'user_id',operator:'eq',value:other}]}),/scope_missing/)); + assert.equal(calls.length,0); +}); +test('refresh updates use the captured generation/revision and never raw table mutation',async()=>{ + const calls=mock((_n,p)=>response(p)); + await withConnectionSnapshot(snapshot(),async()=>{ + await connectionCredentialUpdate('user_credentials',{x_oauth2_access_token_enc:'enc:v1:opaque'},{filters:[{column:'user_id',operator:'eq',value:user}]}); + assert.equal(currentConnectionSession().guard('x').revision,2); + }); + assert.equal(calls[0].name,'mutate_connection');assert.equal(calls[0].body.p_generation,2);assert.equal(calls[0].body.p_revision,1); + assert.equal(calls[0].body.p_kind,'refresh');assert.equal(calls.length,1); +}); +test('lost mutation transport reuses the identical operation and ciphertext exactly once',async()=>{ + let first=true;const calls=mock((_n,p)=>{if(first){first=false;throw new TypeError('lost');}return response(p);}); + await withConnectionSnapshot(snapshot(),()=>connectionCredentialUpdate('user_credentials',{x_oauth2_access_token_enc:'enc:v1:opaque'},{filters:[{column:'user_id',operator:'eq',value:user}]})); + assert.equal(calls.length,2);assert.deepEqual(calls[0].body,calls[1].body); +}); +for(const code of [301,400,401,403,500]) test(`HTTP ${code} rejection neither leaks private error nor retries`,async()=>{ + const calls=mock(()=>new Response('private-token-in-error',{status:code})); + await assert.rejects(connectionRpc('mutate_connection',{},true),e=>e instanceof Error&&!e.message.includes('private-token'));assert.equal(calls.length,1); +}); +test('interleaved real runtime scopes retain independent guards and global-key removal',async()=>{ + mock((_n,p)=>response(p));config.OPENAI_API_KEY='old-global-key'; + await Promise.all([user,other].map(id=>runWithRuntimeScope(async()=>{ + const s=snapshot(id); + await worker.withTenantRuntime({userId:id,settings:{},credentials:{},activePlatforms:[],connectionSnapshot:s},async()=>{ + assert.equal(config.OPENAI_API_KEY,'');await new Promise(r=>setImmediate(r)); + assert.equal(currentConnectionSession().snapshot.user_id,id); + await connectionCredentialUpdate('user_credentials',{x_oauth2_access_token_enc:'encrypted-'+id},{filters:[{column:'user_id',operator:'eq',value:id}]}); + }); + }))); + assert.equal(config.OPENAI_API_KEY,'old-global-key');assert.throws(()=>currentConnectionSession(),/scope_missing/); +}); +test('new queue without current account binding is rejected before an intent is created',async()=>{ + const calls=mock((_n,p)=>Response.json({schema:'ocpf.queue-connection-binding.v1',user_id:p.p_user_id,queue_id:p.p_queue_id,state:'review_required'})); + await assert.rejects(assertQueueConnectionReady(user,other),/review_destination/);assert.equal(calls.length,1); +}); +test('receipt recovery remains available without authorising a new mismatched destination',async()=>{ + mock((_n,p)=>Response.json({schema:'ocpf.queue-connection-binding.v1',user_id:p.p_user_id,queue_id:p.p_queue_id,state:'existing_intent'})); + await assertQueueConnectionReady(user,other); +}); +test('flag off does not call adoption or queue guard endpoints',async()=>{ + const calls=mock(()=>Response.json({}));config.CONNECTION_LIFECYCLE_ENABLED=false; + await assertQueueConnectionReady(user,other);await assert.rejects(loadConnectionSnapshot(user),/disabled/);assert.equal(calls.length,0); +}); +test('authenticated readiness conditionally proves the complete connection schema without adoption',async()=>{ + const id='33333333-3333-4333-8333-333333333333',sha='a'.repeat(40),names:string[]=[]; + const req=new Request('https://worker.invalid/readyz',{headers:{Authorization:'Bearer operator','X-OCPF-Readiness-Nonce':other,'X-OCPF-Expected-Version':id,'X-OCPF-Expected-Sha':sha}}); + const env={WORKER_TICK_TOKEN:'operator',SUPABASE_URL:'https://database.invalid',SUPABASE_SERVICE_ROLE_KEY:'fixture',CREDENTIAL_ENCRYPTION_KEY:'fixture',CONNECTION_LIFECYCLE_ENABLED:'true',CF_VERSION_METADATA:{id,tag:sha,timestamp:'2026-01-01T00:00:00Z'}}; + const f:typeof fetch=async(input)=>{ + const name=new URL(String(input)).pathname.split('/').pop()!;names.push(name); + const c=[...READINESS_CONTRACTS,CONNECTION_READINESS_CONTRACT].find(c=>c.rpc===name)!;assert.ok(c); + return Response.json({contract:c.contract,capabilities:c.capabilities,...(c.name==='publication'?{migration:'20260907054000',lock_order_migration:'20260913061000'}:c.name==='connections'?{migration:'20261001154000'}:{})}); + }; + assert.equal((await handleWorkerReadiness(req,env,f)).status,200);assert.equal(names.length,4);assert.ok(!names.includes('capture_connection_snapshot')); +}); diff --git a/test/security-hardening.test.ts b/test/security-hardening.test.ts index 779efde..c03e33e 100644 --- a/test/security-hardening.test.ts +++ b/test/security-hardening.test.ts @@ -445,6 +445,7 @@ async function main(): Promise { SUPABASE_WORKER_CANARY_USER_IDS: new Set(), SUPABASE_WORKER_GENERATION_ENABLED: false, SUPABASE_PROVIDER_DISPATCH_ENABLED: false, + CONNECTION_LIFECYCLE_ENABLED: false, }); assert.ok(issues.some(issue => issue.includes('COOKIE_SECURE'))); @@ -522,6 +523,7 @@ async function main(): Promise { SUPABASE_WORKER_CANARY_USER_IDS: new Set(), SUPABASE_WORKER_GENERATION_ENABLED: false, SUPABASE_PROVIDER_DISPATCH_ENABLED: false, + CONNECTION_LIFECYCLE_ENABLED: false, }); assert.ok(issues.some(issue => issue.includes('SUPABASE_SERVICE_ROLE_KEY'))); diff --git a/wrangler.toml b/wrangler.toml index 0c1242a..c7b8eb9 100644 --- a/wrangler.toml +++ b/wrangler.toml @@ -18,6 +18,7 @@ SUPABASE_WORKER_BATCH_SIZE = "1" SUPABASE_WORKER_CANARY_REQUIRED = "true" SUPABASE_WORKER_GENERATION_ENABLED = "false" SUPABASE_PROVIDER_DISPATCH_ENABLED = "false" +CONNECTION_LIFECYCLE_ENABLED = "false" DAILY_INVENTORY_PLANNER_ENABLED = "false" DAILY_INVENTORY_PLANNER_START_LOCAL_DATE = "2026-07-14" HTTP_TIMEOUT_MS = "45000"