diff --git a/apps/hub/src/credential-expiry-sweep.ts b/apps/hub/src/credential-expiry-sweep.ts index 60e894418..86bfeef2c 100644 --- a/apps/hub/src/credential-expiry-sweep.ts +++ b/apps/hub/src/credential-expiry-sweep.ts @@ -28,9 +28,11 @@ import { type } from "arktype"; import { credential, principal, provider } from "@intx/db/schema"; import type { DB } from "@intx/db"; import { getLogger } from "@intx/log"; +import { reportError } from "@corbits/error-sink"; import { deliverCredentialMail, findDueCredentialExpiries, + type CredentialExpiredNotification, type ExpiringCredential, type NotifyDeliveryDeps, } from "@corbits/notify"; @@ -55,6 +57,27 @@ const MCP_REFRESH_LEAD_MS = POLL_INTERVAL_MS; const MCP_OAUTH_CLIENT_NAME = "Corbits Workbench"; const log = getLogger(["hub", "credential-expiry-sweep"]); +// Mail-first-claim-after (see `mailThenClaimExpiry`) leaves a still-due +// credential `active` for as long as it keeps failing to notify anyone — +// on purpose, so a transient mail outage or a not-yet-provisioned +// recipient gets more than one tick to resolve. Left unbounded, though, +// `active` stops meaning "usable": the credential is genuinely dead at +// the provider, and the rest of the system (anything selecting an +// `active` credential to actually use) would keep trusting a row that +// will never work again. Past these ages the sweep gives up on ever +// notifying anyone and claims the expiry anyway, so the stored status +// stops asserting something false — reported loudly, since nobody was +// told. +// +// The two failure shapes get different budgets: a mail outage is a +// transient, systemic condition most likely to resolve within a day if +// it's going to resolve at all, while a credential with zero active +// recipients is more often a slow-moving tenant-provisioning gap (a new +// tenant, a departed user being replaced) that deserves more real-world +// time before writing off the notification as unreachable. +const MAX_UNMAILED_CREDENTIAL_AGE_MS = 24 * 60 * 60 * 1000; +const MAX_UNOWNED_CREDENTIAL_AGE_MS = 7 * 24 * 60 * 60 * 1000; + // The OAuth-connected providers whose tokens expire, and how to name // each in the notification. A second such provider generalizes this // list, never a second parallel sweep. @@ -89,6 +112,10 @@ export type RefreshableMcpCredential = { readonly tokens: OAuthTokens; readonly clientInformation?: OAuthClientInformationMixed; readonly recipients: readonly { tenantId: string; principalId: string }[]; + /** The credential's real `expiresAt` column — when it became due for + * this reconnect-nudge fallback, used to bound how long an unmailed + * expiry is left `active` (see `MAX_UNMAILED_CREDENTIAL_AGE_MS`). */ + readonly expiresAt: Date; }; export type CredentialExpirySweepStore = { @@ -215,7 +242,7 @@ export function createDrizzleCredentialExpirySweepStore( ); const due = rows.filter( - (row) => + (row): row is typeof row & { expiresAt: Date } => row.expiresAt !== null && row.expiresAt.getTime() <= cutoff.getTime(), ); if (due.length === 0) return []; @@ -266,6 +293,7 @@ export function createDrizzleCredentialExpirySweepStore( slug: mcpSlugOf(row.providerName), name: parsedMetadata.name ?? row.providerName, serverUrl: row.apiBaseUrl ?? parsedMetadata.url, + expiresAt: row.expiresAt, tokens: { access_token: accessToken, token_type: "bearer", @@ -336,8 +364,93 @@ export type CredentialExpirySweepDeps = { }; /** - * One sweep: claim-and-mail every Hugging-Face-style credential expiry - * due at `now`, then refresh (or, on refresh failure, claim-and-mail) + * Mail a credential-expiry reconnect nudge, then claim the credential as + * `expired` only once that mail has gone out (or is already out — a + * `credential-expired` event dedupes on `credentialId` alone, so a + * redelivery from a retried tick is a harmless no-op, never a second + * mail). Claiming is deliberately the LAST step: `active` is the durable + * "still needs a decision" state, and this function leaves a credential + * `active` whenever that decision (mail sent, credential marked expired) + * did not fully land, so the next tick's own due-scan is the retry — + * no separate retry queue is needed. A thrown mail error is reported + * and swallowed here (never `active` → `expired` with the notification + * lost), and never propagates out to abort sibling candidates in the + * same tick. + * + * That "stay active and retry" grace is bounded by `dueSince` (see + * `MAX_UNMAILED_CREDENTIAL_AGE_MS` / `MAX_UNOWNED_CREDENTIAL_AGE_MS`): + * once a credential has been due for longer than its budget with the + * notification still unsent, this claims the expiry anyway rather than + * asserting `active` forever for a credential that is genuinely dead at + * the provider — reported loudly, since giving up here means nobody was + * ever told. + */ +async function mailThenClaimExpiry( + deps: CredentialExpirySweepDeps, + now: Date, + target: { + readonly credentialId: string; + readonly tenantId: string; + readonly recipients: readonly { tenantId: string; principalId: string }[]; + readonly dueSince: Date; + }, + event: CredentialExpiredNotification, +): Promise { + const ageMs = now.getTime() - target.dueSince.getTime(); + + if (target.recipients.length === 0) { + if (ageMs < MAX_UNOWNED_CREDENTIAL_AGE_MS) { + log.warn`credential ${target.credentialId} has no active recipient to notify; leaving it active until one exists`; + return; + } + reportError( + new Error( + "credential expiry never had an active recipient to notify; claiming it as expired unnotified", + ), + { + operation: "credential-expiry-sweep.abandon-unowned", + tenantId: target.tenantId, + extra: { + credentialId: target.credentialId, + dueSince: target.dueSince.toISOString(), + }, + }, + ); + await deps.store.claimExpiry(target.credentialId, now); + return; + } + + try { + await deliverCredentialMail(deps.notify, event); + } catch (err) { + if (ageMs < MAX_UNMAILED_CREDENTIAL_AGE_MS) { + reportError(err, { + operation: "credential-expiry-sweep.deliver", + tenantId: target.tenantId, + extra: { credentialId: target.credentialId }, + }); + return; + } + reportError(err, { + operation: "credential-expiry-sweep.abandon-unmailed", + tenantId: target.tenantId, + extra: { + credentialId: target.credentialId, + dueSince: target.dueSince.toISOString(), + }, + }); + await deps.store.claimExpiry(target.credentialId, now); + return; + } + const claimed = await deps.store.claimExpiry(target.credentialId, now); + // False means another replica already claimed this exact expiry + // between the mail attempt and this claim — not an error. + if (!claimed) return; +} + +/** + * One sweep: mail-then-claim every Hugging-Face-style credential expiry + * due at `now`, then refresh (or, on refresh failure, mail-then-claim) * every MCP oauth_token credential nearing its own real expiry. Exported * (rather than kept as a closure) so a test can drive a single, * deterministic pass without waiting on `setInterval`. @@ -352,16 +465,21 @@ export async function tickCredentialExpirySweep( const due = findDueCredentialExpiries(candidates, now); for (const { credential: expiring, event } of due) { - const claimed = await deps.store.claimExpiry(expiring.credentialId, now); - // False means another replica already claimed this exact expiry - // between `loadActiveCandidates` and this claim — not an error. - if (!claimed) continue; - - if (expiring.recipients.length === 0) { - log.warn`credential ${expiring.credentialId} expired with no active recipient to notify`; - continue; - } - await deliverCredentialMail(deps.notify, event); + // Guaranteed defined and parseable for anything `findDueCredentialExpiries` + // returned; `now` is an unreachable fallback, never a real fallback value. + const dueSince = + expiring.expiresAt !== undefined ? new Date(expiring.expiresAt) : now; + await mailThenClaimExpiry( + deps, + now, + { + credentialId: expiring.credentialId, + tenantId: expiring.tenantId, + recipients: expiring.recipients, + dueSince, + }, + event, + ); } const refreshable = await deps.store.loadRefreshableMcpCandidates( @@ -392,23 +510,25 @@ export async function tickCredentialExpirySweep( } log.warn`mcp oauth refresh failed for credential ${candidate.credentialId} ("${candidate.name}"): ${result.message}`; - const claimed = await deps.store.claimExpiry(candidate.credentialId, now); - // False means another replica already claimed this exact expiry — not an error. - if (!claimed) continue; - - if (candidate.recipients.length === 0) { - log.warn`credential ${candidate.credentialId} expired with no active recipient to notify`; - continue; - } - await deliverCredentialMail(deps.notify, { - kind: "credential-expired", - tenantId: candidate.tenantId, - credentialId: candidate.credentialId, - providerId: `mcp:${candidate.slug}`, - providerLabel: candidate.name, - recipients: [...candidate.recipients], - createdAt: now.toISOString(), - }); + await mailThenClaimExpiry( + deps, + now, + { + credentialId: candidate.credentialId, + tenantId: candidate.tenantId, + recipients: candidate.recipients, + dueSince: candidate.expiresAt, + }, + { + kind: "credential-expired", + tenantId: candidate.tenantId, + credentialId: candidate.credentialId, + providerId: `mcp:${candidate.slug}`, + providerLabel: candidate.name, + recipients: [...candidate.recipients], + createdAt: now.toISOString(), + }, + ); } } diff --git a/apps/hub/test/credential-expiry-sweep.test.ts b/apps/hub/test/credential-expiry-sweep.test.ts index 290bbfd2c..360630fe4 100644 --- a/apps/hub/test/credential-expiry-sweep.test.ts +++ b/apps/hub/test/credential-expiry-sweep.test.ts @@ -56,6 +56,7 @@ type McpCredentialRow = { accessToken: string; status: "active" | "expired"; recipients: readonly { tenantId: string; principalId: string }[]; + expiresAt: Date; }; function mcpCredential( @@ -72,6 +73,7 @@ function mcpCredential( accessToken: "access_old", status: "active", recipients: [{ tenantId: "tnt_1", principalId: "prn_1" }], + expiresAt: new Date("2026-08-13T11:00:00.000Z"), ...overrides, }; } @@ -136,6 +138,7 @@ function inMemoryStore( slug: r.slug, name: r.name, serverUrl: r.serverUrl, + expiresAt: r.expiresAt, tokens: { access_token: r.accessToken, token_type: "bearer", @@ -183,6 +186,43 @@ function notifyDeps(): NotifyDeliveryDeps & { }; } +/** A notify deps whose `mail` throws — simulates a Postgres blip after a + * credential has already been found due for expiry. */ +function throwingNotifyDeps(error: Error): NotifyDeliveryDeps & { + mailed: { tenantId: string; principalId: string }[]; +} { + return { + mailed: [], + mail: async () => { + throw error; + }, + addressing: { + inbox: (recipient) => `${recipient.principalId}@inbox.invalid`, + from: (kind) => `${kind}@notify.invalid`, + }, + dispatch: createInMemoryNotifyDispatchStore(), + sinks: createSinkRegistry(), + }; +} + +/** A notify deps whose `mail` always dedupes — simulates a retry tick + * re-mailing a credential whose notification already went out. */ +function dedupingNotifyDeps(): NotifyDeliveryDeps & { + mailed: { tenantId: string; principalId: string }[]; +} { + return { + mailed: [], + mail: async (items) => + items.map((item) => ({ messageKey: item.externalId, id: null })), + addressing: { + inbox: (recipient) => `${recipient.principalId}@inbox.invalid`, + from: (kind) => `${kind}@notify.invalid`, + }, + dispatch: createInMemoryNotifyDispatchStore(), + sinks: createSinkRegistry(), + }; +} + const now = new Date("2026-08-13T12:00:00.000Z"); describe("tickCredentialExpirySweep", () => { @@ -203,7 +243,7 @@ describe("tickCredentialExpirySweep", () => { ]); }); - test("a credential with no active recipients is claimed but never mailed", async () => { + test("a credential with no active recipients is never claimed, staying due for a later tick", async () => { const store = inMemoryStore([credential({ recipients: [] })]); const notify = notifyDeps(); @@ -214,10 +254,92 @@ describe("tickCredentialExpirySweep", () => { now: () => now, }); - expect(store.claims).toEqual(["cred_1"]); + expect(store.claims).toEqual([]); expect(notify.mailed).toEqual([]); }); + test("a mail failure leaves the credential unclaimed instead of expiring it silently", async () => { + const store = inMemoryStore([credential()]); + const notify = throwingNotifyDeps(new Error("postgres blip")); + + await tickCredentialExpirySweep({ + store, + notify, + hubUrl: HUB_URL, + now: () => now, + }); + + expect(store.claims).toEqual([]); + + // The next tick sees it still `active` and can finish the job. + const retryNotify = notifyDeps(); + await tickCredentialExpirySweep({ + store, + notify: retryNotify, + hubUrl: HUB_URL, + now: () => now, + }); + + expect(store.claims).toEqual(["cred_1"]); + expect(retryNotify.mailed).toEqual([ + { tenantId: "tnt_1", principalId: "prn_1" }, + ]); + }); + + test("a mail that dedupes (already sent by a prior tick) still claims the credential", async () => { + const store = inMemoryStore([credential()]); + const notify = dedupingNotifyDeps(); + + await tickCredentialExpirySweep({ + store, + notify, + hubUrl: HUB_URL, + now: () => now, + }); + + expect(store.claims).toEqual(["cred_1"]); + }); + + test("one candidate's mail failure does not abandon the rest of the tick", async () => { + const store = inMemoryStore([ + credential({ credentialId: "cred_1" }), + credential({ credentialId: "cred_2" }), + ]); + const mailed: { tenantId: string; principalId: string }[] = []; + const notify: NotifyDeliveryDeps = { + mail: async (items, opts) => { + if (items[0]?.externalId === "cred_1") { + throw new Error("postgres blip"); + } + return items.map((item, index) => { + const id = `mail-${index}`; + mailed.push({ + tenantId: item.tenantId, + principalId: item.principalId, + }); + opts?.enqueue?.({ id, item }); + return { messageKey: item.externalId, id }; + }); + }, + addressing: { + inbox: (recipient) => `${recipient.principalId}@inbox.invalid`, + from: (kind) => `${kind}@notify.invalid`, + }, + dispatch: createInMemoryNotifyDispatchStore(), + sinks: createSinkRegistry(), + }; + + await tickCredentialExpirySweep({ + store, + notify, + hubUrl: HUB_URL, + now: () => now, + }); + + expect(store.claims).toEqual(["cred_2"]); + expect(mailed).toEqual([{ tenantId: "tnt_1", principalId: "prn_1" }]); + }); + test("a credential not yet expired is neither claimed nor mailed", async () => { const store = inMemoryStore([ credential({ expiresAt: "2026-08-13T13:00:00.000Z" }), @@ -252,6 +374,45 @@ describe("tickCredentialExpirySweep", () => { expect(store.claims).toEqual([]); expect(notify.mailed).toEqual([]); }); + + test("a credential with no recipient past the unowned-age bound is claimed anyway, unnotified", async () => { + // 8 days before `now` — past MAX_UNOWNED_CREDENTIAL_AGE_MS (7 days). + const store = inMemoryStore([ + credential({ + recipients: [], + expiresAt: "2026-08-05T00:00:00.000Z", + }), + ]); + const notify = notifyDeps(); + + await tickCredentialExpirySweep({ + store, + notify, + hubUrl: HUB_URL, + now: () => now, + }); + + expect(store.claims).toEqual(["cred_1"]); + expect(notify.mailed).toEqual([]); + }); + + test("a mail failure past the unmailed-age bound is claimed anyway, unnotified", async () => { + // 36 hours before `now` — past MAX_UNMAILED_CREDENTIAL_AGE_MS (24h). + const store = inMemoryStore([ + credential({ expiresAt: "2026-08-11T00:00:00.000Z" }), + ]); + const notify = throwingNotifyDeps(new Error("postgres blip")); + + await tickCredentialExpirySweep({ + store, + notify, + hubUrl: HUB_URL, + now: () => now, + }); + + expect(store.claims).toEqual(["cred_1"]); + expect(notify.mailed).toEqual([]); + }); }); describe("tickCredentialExpirySweep — MCP oauth_token refresh (CL-6207)", () => { @@ -331,6 +492,59 @@ describe("tickCredentialExpirySweep — MCP oauth_token refresh (CL-6207)", () = ]); }); + test("a failed refresh with no active recipients is never claimed, staying due for a later tick", async () => { + const store = inMemoryStore([], [mcpCredential({ recipients: [] })]); + const notify = notifyDeps(); + const refreshMcpTokens = async (): Promise => ({ + ok: false, + message: "invalid_grant: refresh token revoked", + }); + + await tickCredentialExpirySweep({ + store, + notify, + hubUrl: HUB_URL, + refreshMcpTokens, + now: () => now, + }); + + expect(store.mcpClaims).toEqual([]); + expect(notify.mailed).toEqual([]); + }); + + test("a failed refresh whose reconnect mail also fails leaves the credential unclaimed", async () => { + const store = inMemoryStore([], [mcpCredential()]); + const notify = throwingNotifyDeps(new Error("postgres blip")); + const refreshMcpTokens = async (): Promise => ({ + ok: false, + message: "invalid_grant: refresh token revoked", + }); + + await tickCredentialExpirySweep({ + store, + notify, + hubUrl: HUB_URL, + refreshMcpTokens, + now: () => now, + }); + + expect(store.mcpClaims).toEqual([]); + + const retryNotify = notifyDeps(); + await tickCredentialExpirySweep({ + store, + notify: retryNotify, + hubUrl: HUB_URL, + refreshMcpTokens, + now: () => now, + }); + + expect(store.mcpClaims).toEqual(["cred_mcp_1"]); + expect(retryNotify.mailed).toEqual([ + { tenantId: "tnt_1", principalId: "prn_1" }, + ]); + }); + test("an api_key mcp credential is never a refresh candidate", async () => { const store = inMemoryStore( [], @@ -361,4 +575,57 @@ describe("tickCredentialExpirySweep — MCP oauth_token refresh (CL-6207)", () = expect(store.mcpRefreshes).toEqual([]); expect(store.mcpClaims).toEqual([]); }); + + test("a failed refresh with no recipient past the unowned-age bound is claimed anyway, unnotified", async () => { + // 8 days before `now` — past MAX_UNOWNED_CREDENTIAL_AGE_MS (7 days). + const store = inMemoryStore( + [], + [ + mcpCredential({ + recipients: [], + expiresAt: new Date("2026-08-05T00:00:00.000Z"), + }), + ], + ); + const notify = notifyDeps(); + const refreshMcpTokens = async (): Promise => ({ + ok: false, + message: "invalid_grant: refresh token revoked", + }); + + await tickCredentialExpirySweep({ + store, + notify, + hubUrl: HUB_URL, + refreshMcpTokens, + now: () => now, + }); + + expect(store.mcpClaims).toEqual(["cred_mcp_1"]); + expect(notify.mailed).toEqual([]); + }); + + test("a failed refresh whose reconnect mail keeps failing past the unmailed-age bound is claimed anyway", async () => { + // 36 hours before `now` — past MAX_UNMAILED_CREDENTIAL_AGE_MS (24h). + const store = inMemoryStore( + [], + [mcpCredential({ expiresAt: new Date("2026-08-11T00:00:00.000Z") })], + ); + const notify = throwingNotifyDeps(new Error("postgres blip")); + const refreshMcpTokens = async (): Promise => ({ + ok: false, + message: "invalid_grant: refresh token revoked", + }); + + await tickCredentialExpirySweep({ + store, + notify, + hubUrl: HUB_URL, + refreshMcpTokens, + now: () => now, + }); + + expect(store.mcpClaims).toEqual(["cred_mcp_1"]); + expect(notify.mailed).toEqual([]); + }); }); diff --git a/packages/notify/src/deliver.ts b/packages/notify/src/deliver.ts index 8212e1ec4..505137d32 100644 --- a/packages/notify/src/deliver.ts +++ b/packages/notify/src/deliver.ts @@ -6,6 +6,13 @@ // Fan-out to external sinks is queued strictly after the mail commits, one // dispatch row per (mail row, enabled sink). A sink is never called from this // path: the mail is the durable record, and a copy of it is the worker's job. +// +// That queuing step is NOT atomic with the mail write, and nothing +// reconciles the two if a crash or a throw lands between them (CL-7238): +// a mail row can end up committed with no dispatch row ever queued for it. +// Closing that gap needs a seam `@corbits/mailbox` doesn't expose today +// (an in-transaction hook, or a way to read back an existing row's id by +// key) — see CL-7238 for why this can't be fixed from this package alone. import { parseNotificationEvent, type ApprovalNotification,