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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
2 changes: 2 additions & 0 deletions config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,7 @@ export interface AppConfig {
SUPABASE_WORKER_CANARY_REQUIRED: boolean;
SUPABASE_WORKER_CANARY_USER_IDS: Set<string>;
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;
Expand Down Expand Up @@ -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'
Expand Down
49 changes: 49 additions & 0 deletions docs/CONNECTION_GENERATIONS_V1.md
Original file line number Diff line number Diff line change
@@ -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
2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
68 changes: 68 additions & 0 deletions src/connection-http.ts
Original file line number Diff line number Diff line change
@@ -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<string, string> = {
"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<Record<string, unknown>> {
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<never>((_, 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<Uint8Array> | 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<string, unknown>;
} 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();
}
}
}
213 changes: 213 additions & 0 deletions src/connection-lifecycle.ts
Original file line number Diff line number Diff line change
@@ -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<string, unknown>;
connections: ConnectionGuard[];
};
export type ConnectionRpc = (
name: string,
params: Record<string, unknown>,
retrySafe: boolean,
) => Promise<unknown>;
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<string, unknown> {
if (!value || typeof value !== "object" || Array.isArray(value))
throw Error("connection_response_invalid");
return value as Record<string, unknown>;
}
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<ConnectionSnapshot> {
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<string, unknown>,
accountId: string | null = null,
): Promise<ConnectionGuard> {
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<ConnectionProvider, ConnectionGuard>();
private busy = new Set<ConnectionProvider>();
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<string, unknown>,
accountId: string | null = null,
): Promise<ConnectionGuard> {
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;
}
}
Loading
Loading