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
180 changes: 150 additions & 30 deletions apps/hub/src/credential-expiry-sweep.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -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.
Expand Down Expand Up @@ -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 = {
Expand Down Expand Up @@ -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 [];
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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<void> {
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`.
Expand All @@ -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(
Expand Down Expand Up @@ -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(),
},
);
}
}

Expand Down
Loading
Loading