diff --git a/src/mcpd/src/main.ts b/src/mcpd/src/main.ts index 02a92a5..4b72346 100644 --- a/src/mcpd/src/main.ts +++ b/src/mcpd/src/main.ts @@ -474,16 +474,32 @@ async function main(): Promise { }, }, secretRefResolver: secretResolverBridge, + }, { + // Cache-transition events go through pino so BACKEND_UNREACHABLE / + // BACKEND_RECOVERED land in ErrorLogBuffer and `mcpctl errors`. + log: { + warn: (obj: Record, msg: string): void => { app.log.warn(obj, msg); }, + info: (obj: Record, msg: string): void => { app.log.info(obj, msg); }, + }, }); const secretService = new SecretService(secretRepo, secretBackendService); const secretMigrateService = new SecretMigrateService(secretRepo, secretBackendService); const secretBackendRotator = new SecretBackendRotator({ backends: secretBackendService, secrets: secretService, + log: { + error: (obj: Record, msg: string): void => { app.log.error(obj, msg); }, + warn: (msg: string): void => { app.log.warn(msg); }, + }, }); const secretBackendRotatorLoop = new SecretBackendRotatorLoop({ backends: secretBackendService, rotator: secretBackendRotator, + log: { + info: (msg: string): void => { app.log.info(`[rotator] ${msg}`); }, + warn: (msg: string): void => { app.log.warn(`[rotator] ${msg}`); }, + error: (obj: Record, msg: string): void => { app.log.error(obj, msg); }, + }, }); const llmAdapters = new LlmAdapterRegistry(); // LlmService takes the adapter registry so create()/update() can run an diff --git a/src/mcpd/src/services/secret-backend-rotator-loop.ts b/src/mcpd/src/services/secret-backend-rotator-loop.ts index 2fae8ce..0cbef82 100644 --- a/src/mcpd/src/services/secret-backend-rotator-loop.ts +++ b/src/mcpd/src/services/secret-backend-rotator-loop.ts @@ -26,7 +26,11 @@ export interface SecretBackendRotatorLoopDeps { /** Override in tests. */ setTimeout?: (cb: () => void, ms: number) => NodeJS.Timeout; clearTimeout?: (t: NodeJS.Timeout) => void; - log?: { info: (msg: string) => void; warn: (msg: string) => void }; + log?: { + info: (msg: string) => void; + warn: (msg: string) => void; + error: (obj: Record, msg: string) => void; + }; } const DEFAULT_INTERVAL_MS = 24 * 3600 * 1000; @@ -36,7 +40,7 @@ export class SecretBackendRotatorLoop { private readonly timers = new Map(); private readonly setT: (cb: () => void, ms: number) => NodeJS.Timeout; private readonly clearT: (t: NodeJS.Timeout) => void; - private readonly log: { info: (msg: string) => void; warn: (msg: string) => void }; + private readonly log: NonNullable; private stopped = false; constructor(private readonly deps: SecretBackendRotatorLoopDeps) { @@ -44,9 +48,11 @@ export class SecretBackendRotatorLoop { this.clearT = deps.clearTimeout ?? ((t) => global.clearTimeout(t)); this.log = deps.log ?? { // eslint-disable-next-line no-console - info: (m) => console.log(`[rotator] ${m}`), + info: (m: string): void => { console.log(`[rotator] ${m}`); }, // eslint-disable-next-line no-console - warn: (m) => console.warn(`[rotator] ${m}`), + warn: (m: string): void => { console.warn(`[rotator] ${m}`); }, + // eslint-disable-next-line no-console + error: (obj: Record, m: string): void => { console.error(JSON.stringify({ level: 'fatal', ...obj, message: m })); }, }; } @@ -70,13 +76,10 @@ export class SecretBackendRotatorLoop { this.deps.rotator.healthCheck(b.id) .then((res) => { if (!res.ok) { - // eslint-disable-next-line no-console - console.error(JSON.stringify({ - level: 'fatal', - kind: 'BACKEND_TOKEN_DEAD', - backend: b.name, - message: res.message ?? 'unknown', - })); + this.log.error( + { kind: 'BACKEND_TOKEN_DEAD', backend: b.name }, + res.message ?? 'unknown', + ); this.log.warn(`backend '${b.name}' health check failed: ${res.message ?? 'unknown'}`); } }) diff --git a/src/mcpd/src/services/secret-backend-rotator.service.ts b/src/mcpd/src/services/secret-backend-rotator.service.ts index 0a8ed11..27a8989 100644 --- a/src/mcpd/src/services/secret-backend-rotator.service.ts +++ b/src/mcpd/src/services/secret-backend-rotator.service.ts @@ -53,18 +53,37 @@ export interface TokenMeta { rotatable?: boolean; } +/** + * Structured logger. Must be a real pino-shaped logger in production: the + * `BACKEND_TOKEN_DEAD` fatals below used to go out via bare `console.error`, + * which bypasses the pino multistream feeding `ErrorLogBuffer` — so the one + * failure `mcpctl errors` exists to surface was the one it never saw. + */ +export interface RotatorLog { + error(obj: Record, msg: string): void; + warn(msg: string): void; +} + export interface SecretBackendRotatorDeps { backends: SecretBackendService; secrets: SecretService; fetch?: typeof globalThis.fetch; now?: () => Date; + log?: RotatorLog; } export class SecretBackendRotator { private readonly now: () => Date; + private readonly log: RotatorLog; constructor(private readonly deps: SecretBackendRotatorDeps) { this.now = deps.now ?? (() => new Date()); + this.log = deps.log ?? { + // eslint-disable-next-line no-console + error: (obj: Record, msg: string): void => { console.error(JSON.stringify({ level: 'fatal', ...obj, message: msg })); }, + // eslint-disable-next-line no-console + warn: (msg: string): void => { console.warn(msg); }, + }; } /** True iff this backend is a wizard-provisioned token-auth openbao with rotation enabled. */ @@ -144,15 +163,16 @@ export class SecretBackendRotator { : err; const wrappedMsg = wrapped instanceof Error ? wrapped.message : String(wrapped); await this.recordError(backendId, meta, wrappedMsg); - // Loud, structured log so the operator sees it in `kubectl logs deploy/mcpd`. - // eslint-disable-next-line no-console - console.error(JSON.stringify({ - level: 'fatal', - kind: tokenDead ? 'BACKEND_TOKEN_DEAD' : 'BACKEND_ROTATION_FAILED', - backend: backend.name, - url: cfg.url, - message: wrappedMsg, - })); + // Loud and structured, through pino so it also lands in ErrorLogBuffer + // and therefore in `mcpctl errors` — not just in `kubectl logs`. + this.log.error( + { + kind: tokenDead ? 'BACKEND_TOKEN_DEAD' : 'BACKEND_ROTATION_FAILED', + backend: backend.name, + url: cfg.url, + }, + wrappedMsg, + ); throw wrapped; } @@ -164,7 +184,7 @@ export class SecretBackendRotator { // Log but don't fail the rotation — the new token is already live. const msg = err instanceof Error ? err.message : String(err); // eslint-disable-next-line no-console - console.warn(`rotation: revoke old accessor '${oldAccessor}' on backend '${backend.name}' failed (continuing): ${msg}`); + this.log.warn(`rotation: revoke old accessor '${oldAccessor}' on backend '${backend.name}' failed (continuing): ${msg}`); } } @@ -249,7 +269,7 @@ export class SecretBackendRotator { } catch (inner) { // Don't mask the original error — just log the DB failure. // eslint-disable-next-line no-console - console.warn(`rotation: failed to persist lastRotationError (${message}): ${inner instanceof Error ? inner.message : String(inner)}`); + this.log.warn(`rotation: failed to persist lastRotationError (${message}): ${inner instanceof Error ? inner.message : String(inner)}`); } } } diff --git a/src/mcpd/src/services/secret-backend.service.ts b/src/mcpd/src/services/secret-backend.service.ts index fcf5d97..c1578a1 100644 --- a/src/mcpd/src/services/secret-backend.service.ts +++ b/src/mcpd/src/services/secret-backend.service.ts @@ -2,6 +2,7 @@ import type { SecretBackend } from '@prisma/client'; import type { ISecretBackendRepository } from '../repositories/secret-backend.repository.js'; import type { SecretBackendDriver } from './secret-backends/types.js'; import { createDriver, type DriverFactoryDeps } from './secret-backends/factory.js'; +import { CachingSecretBackendDriver, type CachingDriverOptions, type CacheStats } from './secret-backends/caching.js'; import { NotFoundError, ConflictError } from './mcp-server.service.js'; export class SecretBackendInUseError extends Error { @@ -17,6 +18,7 @@ export class SecretBackendService { constructor( private readonly repo: ISecretBackendRepository, private readonly driverDeps: DriverFactoryDeps, + private readonly cacheOpts: CachingDriverOptions = {}, ) {} async list(): Promise { @@ -87,12 +89,34 @@ export class SecretBackendService { this.driverCache.delete(id); } - /** Get the driver for a given backend id, creating + caching on first call. */ + /** + * Get the driver for a given backend id, creating + caching on first call. + * + * Remote backends are wrapped in `CachingSecretBackendDriver` so a backend + * outage degrades to "serving last known-good" instead of failing every + * caller. `plaintext` is deliberately NOT wrapped: its `read()` is an + * identity function over the DB row passed in by the caller, so a value cache + * there would serve pre-update data with nothing to invalidate it. + * + * Config changes invalidate for free — `update()` and `delete()` drop this + * map, and the value cache lives inside the driver instance, which is right: + * if `url`/`mount`/`pathPrefix` change, the cached names now mean something + * different. + */ driverFor(backend: SecretBackend): SecretBackendDriver { const cached = this.driverCache.get(backend.id); if (cached) return cached; - const driver = createDriver(backend, this.driverDeps); + const base = createDriver(backend, this.driverDeps); + const driver = backend.type === 'plaintext' + ? base + : new CachingSecretBackendDriver(base, { ...this.cacheOpts, backendName: backend.name }); this.driverCache.set(backend.id, driver); return driver; } + + /** Cache state for a backend, for the health endpoint. Never exposes values. */ + cacheStatsFor(backend: SecretBackend): CacheStats | undefined { + const driver = this.driverFor(backend); + return driver instanceof CachingSecretBackendDriver ? driver.stats() : undefined; + } } diff --git a/src/mcpd/src/services/secret-backends/caching.ts b/src/mcpd/src/services/secret-backends/caching.ts new file mode 100644 index 0000000..22b4276 --- /dev/null +++ b/src/mcpd/src/services/secret-backends/caching.ts @@ -0,0 +1,198 @@ +/** + * Caching + stale-while-error decorator for any `SecretBackendDriver`. + * + * ## Why this exists + * + * `SecretService.resolveData()` calls `driver.read()` on *every* use, and every + * consumer funnels through it: server env resolution, LLM api keys, chat, git + * providers, code repos, webhooks. With a remote backend that means one network + * round-trip per secret per call, and — worse — any OpenBao blip propagates + * straight through. An instance that restarts during a blip fails env + * resolution, gets marked ERROR, and enters a 30s×5-then-5min backoff + * (`instance.service.ts`), so a few seconds of backend unavailability turns + * into minutes of degraded service. + * + * ## Semantics + * + * - **Fresh** (age < ttlMs): served from memory, no network. + * - **Stale-while-error**: past the TTL we always try the backend first. If it + * answers, we refresh. If it fails *as a transport failure* + * (`SecretBackendUnavailableError`), we serve the last known-good value + * instead of throwing. This is the part that actually stops the ERROR storm. + * - **`SecretNotFoundError` evicts and rethrows.** Never served stale — that + * would resurrect a deliberately deleted or revoked credential, which is + * strictly worse than an outage. + * - **Any other error rethrows, without stale.** A 403 that survives a token + * refresh means our grants were revoked; papering over it with cached data is + * exactly how an upstream OpenBao re-init once went unnoticed for four days. + * - **No negative caching.** A miss must re-check; retry/backoff already lives + * in the driver. + * + * The stale window is deliberately unbounded. A cap would mean a long outage + * eventually takes mcpd down anyway, which defeats the purpose, and the + * revoked-credential case is already handled definitively by `SecretNotFound`. + * + * ## What this does NOT fix + * + * A cold cache during an outage. If mcpd restarts while the backend is + * unreachable, nothing has a last-known-good value and secret-bearing servers + * fail to start — honestly, with a loud error. That is the correct behaviour: + * booting a server with an empty credential is the failure mode that had + * gitea-mcp reporting healthy while every authed call failed. The mitigation is + * to warm this cache at boot, not to invent a value. + * + * Values live in heap in cleartext for the TTL, so the map is bounded (LRU) and + * values are never logged. + */ +import type { SecretBackendDriver, SecretData, ExternalRef } from './types.js'; +import { SecretNotFoundError, SecretBackendUnavailableError } from './types.js'; + +export interface CachingDriverLog { + warn(obj: Record, msg: string): void; + info(obj: Record, msg: string): void; +} + +export interface CachingDriverOptions { + /** How long a value is served without consulting the backend. */ + ttlMs?: number; + /** LRU bound — these are plaintext credentials held in memory. */ + maxEntries?: number; + /** Backend name, for log context only. */ + backendName?: string; + now?: () => number; + log?: CachingDriverLog; +} + +interface CacheEntry { + data: SecretData; + fetchedAt: number; + /** Set when we last served this past its TTL because the backend was down. */ + staleSince: number | undefined; +} + +export const DEFAULT_CACHE_TTL_MS = 300_000; +export const DEFAULT_CACHE_MAX_ENTRIES = 500; + +const NOOP_LOG: CachingDriverLog = { warn: () => undefined, info: () => undefined }; + +export interface CacheStats { + entries: number; + servingStale: number; + oldestStaleSince: number | undefined; +} + +export class CachingSecretBackendDriver implements SecretBackendDriver { + readonly kind: string; + + private readonly entries = new Map(); + private readonly ttlMs: number; + private readonly maxEntries: number; + private readonly backendName: string; + private readonly nowFn: () => number; + private readonly log: CachingDriverLog; + + constructor(private readonly inner: SecretBackendDriver, opts: CachingDriverOptions = {}) { + this.kind = `cached:${inner.kind}`; + this.ttlMs = opts.ttlMs ?? DEFAULT_CACHE_TTL_MS; + this.maxEntries = opts.maxEntries ?? DEFAULT_CACHE_MAX_ENTRIES; + this.backendName = opts.backendName ?? inner.kind; + this.nowFn = opts.now ?? ((): number => Date.now()); + this.log = opts.log ?? NOOP_LOG; + } + + async read(input: { name: string; externalRef: ExternalRef; data: SecretData }): Promise { + const now = this.nowFn(); + const cached = this.entries.get(input.name); + + if (cached !== undefined && now - cached.fetchedAt < this.ttlMs) { + this.touch(input.name, cached); + return cached.data; + } + + try { + const data = await this.inner.read(input); + if (cached?.staleSince !== undefined) { + // Edge-triggered: only on the transition back to healthy. + this.log.info( + { kind: 'BACKEND_RECOVERED', backend: this.backendName, secret: input.name, + staleForMs: now - cached.staleSince }, + `secret backend '${this.backendName}' recovered; '${input.name}' is live again`, + ); + } + this.store(input.name, { data, fetchedAt: now, staleSince: undefined }); + return data; + } catch (err) { + if (err instanceof SecretNotFoundError) { + // Definitive. Drop the stale copy so we can never hand it out later. + this.entries.delete(input.name); + throw err; + } + if (!(err instanceof SecretBackendUnavailableError) || cached === undefined) { + throw err; + } + if (cached.staleSince === undefined) { + cached.staleSince = now; + this.log.warn( + { kind: 'BACKEND_UNREACHABLE', backend: this.backendName, secret: input.name, + ageMs: now - cached.fetchedAt, reason: err.message }, + `secret backend '${this.backendName}' unreachable; serving cached '${input.name}'`, + ); + } + this.touch(input.name, cached); + return cached.data; + } + } + + async write(input: { name: string; data: SecretData }): Promise<{ externalRef: ExternalRef; storedData: SecretData }> { + const result = await this.inner.write(input); + // Cache what a subsequent read() would return — the values just written — + // not `storedData`, which remote drivers deliberately leave empty. + this.store(input.name, { data: input.data, fetchedAt: this.nowFn(), staleSince: undefined }); + return result; + } + + async delete(input: { name: string; externalRef: ExternalRef }): Promise { + await this.inner.delete(input); + this.entries.delete(input.name); + } + + async list(): Promise> { + return this.inner.list(); + } + + async healthCheck(): Promise<{ ok: boolean; detail?: string }> { + return this.inner.healthCheck?.() ?? { ok: true, detail: 'no probe' }; + } + + async authCheck(): Promise<{ ok: boolean; detail?: string }> { + return this.inner.authCheck?.() ?? { ok: true, detail: 'no probe' }; + } + + /** Cache state for the backend health endpoint. Never exposes values. */ + stats(): CacheStats { + let servingStale = 0; + let oldestStaleSince: number | undefined; + for (const e of this.entries.values()) { + if (e.staleSince === undefined) continue; + servingStale++; + if (oldestStaleSince === undefined || e.staleSince < oldestStaleSince) oldestStaleSince = e.staleSince; + } + return { entries: this.entries.size, servingStale, oldestStaleSince }; + } + + /** Move an entry to the MRU end of the insertion-ordered Map. */ + private touch(name: string, entry: CacheEntry): void { + this.entries.delete(name); + this.entries.set(name, entry); + } + + private store(name: string, entry: CacheEntry): void { + this.entries.delete(name); + this.entries.set(name, entry); + while (this.entries.size > this.maxEntries) { + const oldest = this.entries.keys().next(); + if (oldest.done === true) break; + this.entries.delete(oldest.value); + } + } +} diff --git a/src/mcpd/src/services/secret-backends/openbao.ts b/src/mcpd/src/services/secret-backends/openbao.ts index 93e902c..34ddedc 100644 --- a/src/mcpd/src/services/secret-backends/openbao.ts +++ b/src/mcpd/src/services/secret-backends/openbao.ts @@ -28,6 +28,7 @@ */ import { readFile } from 'node:fs/promises'; import type { SecretBackendDriver, SecretData, ExternalRef, SecretRefResolver } from './types.js'; +import { SecretNotFoundError, SecretBackendUnavailableError } from './types.js'; /** Best-effort read of a response body for error messages. Empty on parse failure. */ async function bodyText(res: Response): Promise { @@ -77,10 +78,23 @@ export interface OpenBaoDriverDeps { readServiceAccountToken?: (path: string) => Promise; /** Clock for cache TTL — overridable in tests. */ now?: () => number; + /** Per-request timeout. Without one, an unreachable OpenBao hangs every caller. */ + timeoutMs?: number; + /** Total attempts for retryable failures (network / 5xx / 429). 1 disables retry. */ + maxAttempts?: number; + /** Base for exponential backoff between retries; full jitter is applied. */ + backoffBaseMs?: number; + /** Test seam — real sleeps would make the retry tests take seconds. */ + sleep?: (ms: number) => Promise; } const SA_TOKEN_DEFAULT_PATH = '/var/run/secrets/kubernetes.io/serviceaccount/token'; const TOKEN_RENEW_GRACE_MS = 60_000; +const DEFAULT_TIMEOUT_MS = 5_000; +const DEFAULT_MAX_ATTEMPTS = 3; +const DEFAULT_BACKOFF_BASE_MS = 200; +/** Statuses worth retrying: the backend is up but cannot answer right now. */ +const RETRYABLE_STATUS = new Set([429, 500, 502, 503, 504]); export class OpenBaoDriver implements SecretBackendDriver { readonly kind = 'openbao'; @@ -98,6 +112,10 @@ export class OpenBaoDriver implements SecretBackendDriver { private readonly resolver: SecretRefResolver | undefined; private readonly readSaToken: (path: string) => Promise; private readonly nowFn: () => number; + private readonly timeoutMs: number; + private readonly maxAttempts: number; + private readonly backoffBaseMs: number; + private readonly sleep: (ms: number) => Promise; // Cached vault token + when (epoch ms) it should be considered expired and refetched. private cachedToken: string | undefined; @@ -131,13 +149,19 @@ export class OpenBaoDriver implements SecretBackendDriver { if (deps.secretRefResolver !== undefined) this.resolver = deps.secretRefResolver; this.readSaToken = deps.readServiceAccountToken ?? ((path) => readFile(path, 'utf-8').then((s) => s.trim())); this.nowFn = deps.now ?? (() => Date.now()); + this.timeoutMs = deps.timeoutMs ?? DEFAULT_TIMEOUT_MS; + this.maxAttempts = deps.maxAttempts ?? DEFAULT_MAX_ATTEMPTS; + this.backoffBaseMs = deps.backoffBaseMs ?? DEFAULT_BACKOFF_BASE_MS; + this.sleep = deps.sleep ?? ((ms: number): Promise => new Promise((r) => { setTimeout(r, ms); })); } async read(input: { name: string; externalRef: ExternalRef; data: SecretData }): Promise { const path = this.pathFor(input.name); const res = await this.request('GET', `/v1/${this.mount}/data/${path}`); if (res.status === 404) { - throw new Error(`OpenBao: secret '${input.name}' not found at ${path}`); + // Definitive answer, not a transport failure — the caching decorator + // must evict rather than serve a stale value here. + throw new SecretNotFoundError(`OpenBao: secret '${input.name}' not found at ${path}`); } if (!res.ok) throw new Error(`OpenBao read ${path}: HTTP ${res.status} ${await bodyText(res)}`); const body = await res.json() as { data?: { data?: SecretData } }; @@ -174,10 +198,50 @@ export class OpenBaoDriver implements SecretBackendDriver { })); } + /** + * LIVENESS. Deliberately unauthenticated: `sys/health` needs no token, and + * routing it through `request()` (as this used to) took a login first — so an + * expired role reported as "OpenBao is down", and every probe cost a login. + * + * OpenBao encodes its state in the status code, so map it rather than + * collapsing everything to ok/not-ok. + */ async healthCheck(): Promise<{ ok: boolean; detail?: string }> { try { - const res = await this.request('GET', '/v1/sys/health'); - return { ok: res.ok, detail: `HTTP ${res.status}` }; + const headers: Record = {}; + if (this.namespace !== undefined) headers['X-Vault-Namespace'] = this.namespace; + const res = await this.fetchImpl(`${this.url}/v1/sys/health`, { + method: 'GET', + headers, + signal: AbortSignal.timeout(this.timeoutMs), + }); + switch (res.status) { + case 200: return { ok: true, detail: 'active' }; + case 429: return { ok: true, detail: 'standby' }; + case 472: case 473: return { ok: true, detail: 'replication secondary' }; + case 501: return { ok: false, detail: 'not initialized' }; + case 503: return { ok: false, detail: 'sealed' }; + default: return { ok: res.ok, detail: `HTTP ${String(res.status)}` }; + } + } catch (err) { + return { ok: false, detail: err instanceof Error ? err.message : String(err) }; + } + } + + /** + * READINESS. Exercises the capability we actually depend on — read/list under + * `//` — using the credentials we hold. + * + * `list()` rather than `auth/token/lookup-self` on purpose: lookup-self only + * proves the token exists, not that its policy still grants anything. The + * four-day outage in e51b924 was exactly a live token whose grants had been + * dropped by an upstream re-init. The existing read policy already permits + * this call, so it needs no bao-side change. + */ + async authCheck(): Promise<{ ok: boolean; detail?: string }> { + try { + await this.list(); + return { ok: true, detail: `readable at ${this.mount}/${this.pathPrefix}` }; } catch (err) { return { ok: false, detail: err instanceof Error ? err.message : String(err) }; } @@ -206,11 +270,28 @@ export class OpenBaoDriver implements SecretBackendDriver { const loginUrl = `${this.url}/v1/auth/${this.k8sAuthMount}/login`; const headers: Record = { 'Content-Type': 'application/json' }; if (this.namespace !== undefined) headers['X-Vault-Namespace'] = this.namespace; - const res = await this.fetchImpl(loginUrl, { - method: 'POST', - headers, - body: JSON.stringify({ role: this.k8sRole, jwt }), - }); + // Bounded like every other call: a hung login is indistinguishable from a + // hung read to the caller, and this one used to have no timeout at all. + let res: Response; + try { + res = await this.fetchImpl(loginUrl, { + method: 'POST', + headers, + body: JSON.stringify({ role: this.k8sRole, jwt }), + signal: AbortSignal.timeout(this.timeoutMs), + }); + } catch (err) { + throw new SecretBackendUnavailableError( + `OpenBao kubernetes login (role=${this.k8sRole!}): ${err instanceof Error ? err.message : String(err)}`, + { cause: err }, + ); + } + if (RETRYABLE_STATUS.has(res.status)) { + throw new SecretBackendUnavailableError( + `OpenBao kubernetes login (role=${this.k8sRole!}): HTTP ${String(res.status)}`, + { lastStatus: res.status }, + ); + } if (!res.ok) { const text = await res.text().catch(() => ''); throw new Error(`OpenBao kubernetes login (role=${this.k8sRole!}): HTTP ${String(res.status)} ${text}`); @@ -229,30 +310,77 @@ export class OpenBaoDriver implements SecretBackendDriver { return clientToken; } - private async request(method: string, path: string, body?: unknown): Promise { - const token = await this.getToken(); + /** Build a fresh RequestInit — headers must not be shared across attempts. */ + private buildInit(method: string, token: string, body?: unknown): RequestInit { const headers: Record = { 'X-Vault-Token': token }; if (this.namespace !== undefined) headers['X-Vault-Namespace'] = this.namespace; if (body !== undefined) headers['Content-Type'] = 'application/json'; - - const init: RequestInit = { method, headers }; + const init: RequestInit = { method, headers, signal: AbortSignal.timeout(this.timeoutMs) }; if (body !== undefined) init.body = JSON.stringify(body); + return init; + } - const res = await this.fetchImpl(`${this.url}${path}`, init); + /** Full-jitter exponential backoff, so concurrent callers don't resonate. */ + private backoffFor(attempt: number): number { + return Math.random() * this.backoffBaseMs * Math.pow(2, attempt - 1); + } - // If the cached token expired between cache-check and request (k8s clock - // skew, server-side revocation, etc.), purge cache and retry once. - if (res.status === 403 && this.cachedToken !== undefined) { - this.cachedToken = undefined; - this.cachedTokenExpiresAt = 0; - const fresh = await this.getToken(); - const retryHeaders: Record = { 'X-Vault-Token': fresh }; - if (this.namespace !== undefined) retryHeaders['X-Vault-Namespace'] = this.namespace; - if (body !== undefined) retryHeaders['Content-Type'] = 'application/json'; - const retryInit: RequestInit = { method, headers: retryHeaders }; - if (body !== undefined) retryInit.body = JSON.stringify(body); - return this.fetchImpl(`${this.url}${path}`, retryInit); + private async request(method: string, path: string, body?: unknown): Promise { + const url = `${this.url}${path}`; + let lastStatus: number | undefined; + let lastErr: unknown; + + for (let attempt = 1; attempt <= this.maxAttempts; attempt++) { + let res: Response; + try { + const token = await this.getToken(); + res = await this.fetchImpl(url, this.buildInit(method, token, body)); + } catch (err) { + // Network failure, DNS failure, or our own AbortSignal firing. + lastErr = err; + if (attempt < this.maxAttempts) { + await this.sleep(this.backoffFor(attempt)); + continue; + } + throw new SecretBackendUnavailableError( + `OpenBao ${method} ${path}: ${err instanceof Error ? err.message : String(err)} (after ${String(attempt)} attempt(s))`, + { cause: err }, + ); + } + + // If the cached token expired between cache-check and request (k8s clock + // skew, server-side revocation, etc.), purge cache and retry once. This + // is deliberately OUTSIDE the retry budget: it is a credential refresh, + // not a backend-unavailable condition, and it must stay single-shot so a + // genuinely revoked grant fails loudly instead of looping. + if (res.status === 403 && this.cachedToken !== undefined) { + this.cachedToken = undefined; + this.cachedTokenExpiresAt = 0; + const fresh = await this.getToken(); + return this.fetchImpl(url, this.buildInit(method, fresh, body)); + } + + // The backend is up but cannot answer right now — 503 is also what a + // sealed OpenBao returns, which used to be an immediate hard failure. + if (RETRYABLE_STATUS.has(res.status) && attempt < this.maxAttempts) { + lastStatus = res.status; + await this.sleep(this.backoffFor(attempt)); + continue; + } + if (RETRYABLE_STATUS.has(res.status)) { + throw new SecretBackendUnavailableError( + `OpenBao ${method} ${path}: HTTP ${String(res.status)} after ${String(attempt)} attempt(s)`, + { lastStatus: res.status }, + ); + } + + return res; } - return res; + + /* c8 ignore next 5 -- unreachable: every loop exit above returns or throws */ + throw new SecretBackendUnavailableError( + `OpenBao ${method} ${path}: exhausted ${String(this.maxAttempts)} attempt(s)`, + lastErr !== undefined ? { cause: lastErr, ...(lastStatus !== undefined ? { lastStatus } : {}) } : (lastStatus !== undefined ? { lastStatus } : {}), + ); } } diff --git a/src/mcpd/src/services/secret-backends/types.ts b/src/mcpd/src/services/secret-backends/types.ts index bab41d5..8d91ce9 100644 --- a/src/mcpd/src/services/secret-backends/types.ts +++ b/src/mcpd/src/services/secret-backends/types.ts @@ -46,8 +46,24 @@ export interface SecretBackendDriver { /** List everything the backend knows about. Used for migration + drift detection. */ list(): Promise>; - /** Optional: health probe. Used by `mcpctl describe secretbackend`. */ + /** + * Optional LIVENESS probe: is the backend reachable at all? + * + * Must NOT require authentication — the whole point is to separate "the + * backend is down/sealed" from "our credentials stopped working". Compare + * `authCheck()`, which is the readiness half. + */ healthCheck?(): Promise<{ ok: boolean; detail?: string }>; + + /** + * Optional READINESS probe: can we actually read through this backend with + * the credentials we hold? + * + * A backend that answers `healthCheck()` but fails here is the exact shape of + * the incident where a re-initialised OpenBao left mcpd holding valid-looking + * tokens that granted nothing. Reporting one signal for both hides it. + */ + authCheck?(): Promise<{ ok: boolean; detail?: string }>; } /** Stored config for a SecretBackend row; dispatched on `type`. */ @@ -66,3 +82,40 @@ export interface BackendRow { export interface SecretRefResolver { resolve(secretName: string, key: string): Promise; } + +/** + * The backend gave a definitive answer: this secret (or key) does not exist. + * + * Callers may treat this as final. The caching decorator EVICTS on this and + * never serves a stale value for it — serving stale here would resurrect a + * deliberately deleted or revoked credential, which is strictly worse than an + * outage. + */ +export class SecretNotFoundError extends Error { + constructor(message: string, options?: { cause?: unknown }) { + super(message, options); + this.name = 'SecretNotFoundError'; + } +} + +/** + * The backend could not be reached or did not answer: DNS/TCP failure, request + * timeout, or an exhausted retry budget against 5xx/429. + * + * This is the ONLY error the caching decorator will serve a stale value for. + * The distinction has to be typed rather than string-matched: a mis-classified + * "not found" would resurrect deleted secrets, and a mis-classified auth + * failure would silently paper over a backend whose grants were revoked — the + * failure mode that let an OpenBao re-init break every secret write for four + * days (commit e51b924). + */ +export class SecretBackendUnavailableError extends Error { + /** HTTP status of the last attempt, when the failure was an HTTP response. */ + readonly lastStatus: number | undefined; + + constructor(message: string, options?: { cause?: unknown; lastStatus?: number }) { + super(message, options?.cause !== undefined ? { cause: options.cause } : undefined); + this.name = 'SecretBackendUnavailableError'; + this.lastStatus = options?.lastStatus; + } +} diff --git a/src/mcpd/tests/secret-backends.test.ts b/src/mcpd/tests/secret-backends.test.ts index 0f86075..43785a2 100644 --- a/src/mcpd/tests/secret-backends.test.ts +++ b/src/mcpd/tests/secret-backends.test.ts @@ -1,6 +1,7 @@ import { describe, it, expect, vi } from 'vitest'; import { PlaintextDriver } from '../src/services/secret-backends/plaintext.js'; import { OpenBaoDriver } from '../src/services/secret-backends/openbao.js'; +import { SecretNotFoundError, SecretBackendUnavailableError } from '../src/services/secret-backends/types.js'; describe('PlaintextDriver', () => { const driver = new PlaintextDriver({ listAllPlaintext: async () => [{ name: 'a', data: { k: 'v' } }] }); @@ -242,3 +243,91 @@ describe('OpenBaoDriver', () => { }); }); }); + +describe('OpenBaoDriver: resilience', () => { + const resolver = { resolve: vi.fn(async () => 'test-vault-token') }; + /** No real sleeping — otherwise the backoff tests take seconds. */ + const noSleep = async (): Promise => undefined; + + function driverWith(fetchFn: ReturnType, opts: Record = {}): OpenBaoDriver { + return new OpenBaoDriver( + { url: 'http://bao.example:8200', tokenSecretRef: { name: 'bao', key: 'token' } }, + { fetch: fetchFn as unknown as typeof fetch, secretRefResolver: resolver, sleep: noSleep, ...opts }, + ); + } + + it('maps a 404 read to SecretNotFoundError', async () => { + const fetchFn = vi.fn(async () => new Response('', { status: 404 })); + await expect(driverWith(fetchFn).read({ name: 'gone', externalRef: '', data: {} })) + .rejects.toThrow(SecretNotFoundError); + }); + + it('purges the token cache and retries once on 403 — outside the retry budget', async () => { + // This path existed but was never covered; it is the revocation/re-init case. + let n = 0; + const fetchFn = vi.fn(async () => { + n++; + if (n === 1) return new Response('', { status: 403 }); + return new Response(JSON.stringify({ data: { data: { token: 'ok' } } }), { status: 200 }); + }); + const d = driverWith(fetchFn); + await expect(d.read({ name: 's', externalRef: '', data: {} })).resolves.toEqual({ token: 'ok' }); + expect(fetchFn).toHaveBeenCalledTimes(2); + }); + + it('retries a 503 (sealed) and succeeds', async () => { + let n = 0; + const fetchFn = vi.fn(async () => { + n++; + if (n < 3) return new Response('', { status: 503 }); + return new Response(JSON.stringify({ data: { data: { token: 'ok' } } }), { status: 200 }); + }); + await expect(driverWith(fetchFn).read({ name: 's', externalRef: '', data: {} })) + .resolves.toEqual({ token: 'ok' }); + expect(fetchFn).toHaveBeenCalledTimes(3); + }); + + it('throws SecretBackendUnavailableError once the retry budget is exhausted', async () => { + const fetchFn = vi.fn(async () => new Response('', { status: 503 })); + await expect(driverWith(fetchFn, { maxAttempts: 3 }).read({ name: 's', externalRef: '', data: {} })) + .rejects.toThrow(SecretBackendUnavailableError); + expect(fetchFn).toHaveBeenCalledTimes(3); + }); + + it('classifies a network/abort failure as SecretBackendUnavailableError', async () => { + const fetchFn = vi.fn(async () => { throw new DOMException('timed out', 'TimeoutError'); }); + await expect(driverWith(fetchFn, { maxAttempts: 2 }).read({ name: 's', externalRef: '', data: {} })) + .rejects.toThrow(SecretBackendUnavailableError); + expect(fetchFn).toHaveBeenCalledTimes(2); + }); + + it('passes an AbortSignal on every request', async () => { + const fetchFn = vi.fn(async () => new Response(JSON.stringify({ data: { data: {} } }), { status: 200 })); + await driverWith(fetchFn, { timeoutMs: 1234 }).read({ name: 's', externalRef: '', data: {} }); + const [, init] = fetchFn.mock.calls[0] as [unknown, RequestInit]; + expect(init.signal).toBeInstanceOf(AbortSignal); + }); + + it('healthCheck is unauthenticated and maps OpenBao status codes', async () => { + const cases: Array<[number, boolean, string]> = [ + [200, true, 'active'], + [429, true, 'standby'], + [501, false, 'not initialized'], + [503, false, 'sealed'], + ]; + for (const [status, ok, detail] of cases) { + const fetchFn = vi.fn(async () => new Response('', { status })); + const result = await driverWith(fetchFn).healthCheck(); + expect(result).toEqual({ ok, detail }); + // The whole point of the split: no token is minted for a liveness probe. + const [, init] = fetchFn.mock.calls[0] as [unknown, RequestInit]; + expect((init.headers as Record)['X-Vault-Token']).toBeUndefined(); + } + }); + + it('authCheck reports false when the token can no longer list', async () => { + const fetchFn = vi.fn(async () => new Response('', { status: 403 })); + const result = await driverWith(fetchFn).authCheck(); + expect(result.ok).toBe(false); + }); +}); diff --git a/src/mcpd/tests/secret-cache.test.ts b/src/mcpd/tests/secret-cache.test.ts new file mode 100644 index 0000000..fbbf047 --- /dev/null +++ b/src/mcpd/tests/secret-cache.test.ts @@ -0,0 +1,198 @@ +import { describe, it, expect, vi } from 'vitest'; +import { + CachingSecretBackendDriver, + type CachingDriverLog, +} from '../src/services/secret-backends/caching.js'; +import { + SecretNotFoundError, + SecretBackendUnavailableError, + type SecretBackendDriver, + type SecretData, +} from '../src/services/secret-backends/types.js'; + +/** Minimal fake backing driver whose read() behaviour the tests drive. */ +function makeInner(overrides: Partial = {}): SecretBackendDriver & { + read: ReturnType; + write: ReturnType; + delete: ReturnType; +} { + return { + kind: 'fake', + read: vi.fn(async () => ({ token: 'live' } as SecretData)), + write: vi.fn(async () => ({ externalRef: 'ref', storedData: {} as SecretData })), + delete: vi.fn(async () => undefined), + list: vi.fn(async () => []), + ...overrides, + } as never; +} + +function makeLog(): CachingDriverLog & { warns: Array>; infos: Array> } { + const warns: Array> = []; + const infos: Array> = []; + return { warns, infos, warn: (o) => { warns.push(o); }, info: (o) => { infos.push(o); } }; +} + +const REQ = { name: 'gitea-creds', externalRef: 'secret/mcpctl/gitea-creds', data: {} }; + +describe('CachingSecretBackendDriver', () => { + it('serves from cache within the TTL without touching the backend', async () => { + const inner = makeInner(); + let now = 1_000; + const d = new CachingSecretBackendDriver(inner, { ttlMs: 5_000, now: () => now }); + + expect(await d.read(REQ)).toEqual({ token: 'live' }); + now += 4_999; + expect(await d.read(REQ)).toEqual({ token: 'live' }); + + expect(inner.read).toHaveBeenCalledTimes(1); + }); + + it('refetches once the TTL has elapsed', async () => { + const inner = makeInner(); + let now = 1_000; + const d = new CachingSecretBackendDriver(inner, { ttlMs: 5_000, now: () => now }); + + await d.read(REQ); + now += 5_001; + await d.read(REQ); + + expect(inner.read).toHaveBeenCalledTimes(2); + }); + + it('serves the stale value when the backend is unavailable', async () => { + const inner = makeInner(); + let now = 1_000; + const log = makeLog(); + const d = new CachingSecretBackendDriver(inner, { ttlMs: 1_000, now: () => now, log, backendName: 'bao' }); + + await d.read(REQ); + inner.read.mockRejectedValue(new SecretBackendUnavailableError('bao down')); + now += 10_000; + + // This is the whole point: no throw, so instance.service never marks ERROR. + expect(await d.read(REQ)).toEqual({ token: 'live' }); + expect(log.warns[0]?.kind).toBe('BACKEND_UNREACHABLE'); + }); + + it('logs BACKEND_UNREACHABLE only on the transition, not on every stale read', async () => { + const inner = makeInner(); + let now = 1_000; + const log = makeLog(); + const d = new CachingSecretBackendDriver(inner, { ttlMs: 1_000, now: () => now, log }); + + await d.read(REQ); + inner.read.mockRejectedValue(new SecretBackendUnavailableError('bao down')); + for (let i = 0; i < 5; i++) { now += 2_000; await d.read(REQ); } + + expect(log.warns.filter((w) => w.kind === 'BACKEND_UNREACHABLE')).toHaveLength(1); + }); + + it('logs BACKEND_RECOVERED once the backend answers again', async () => { + const inner = makeInner(); + let now = 1_000; + const log = makeLog(); + const d = new CachingSecretBackendDriver(inner, { ttlMs: 1_000, now: () => now, log }); + + await d.read(REQ); + inner.read.mockRejectedValue(new SecretBackendUnavailableError('bao down')); + now += 2_000; + await d.read(REQ); + + inner.read.mockResolvedValue({ token: 'rotated' }); + now += 2_000; + expect(await d.read(REQ)).toEqual({ token: 'rotated' }); + expect(log.infos.filter((i) => i.kind === 'BACKEND_RECOVERED')).toHaveLength(1); + }); + + it('NEVER serves stale for a deleted secret — evicts and rethrows', async () => { + // Regression guard. Serving stale here would resurrect a revoked + // credential, which is strictly worse than an outage. + const inner = makeInner(); + let now = 1_000; + const d = new CachingSecretBackendDriver(inner, { ttlMs: 1_000, now: () => now }); + + await d.read(REQ); + inner.read.mockRejectedValue(new SecretNotFoundError('gone')); + now += 2_000; + + await expect(d.read(REQ)).rejects.toThrow(SecretNotFoundError); + expect(d.stats().entries).toBe(0); + + // And the entry really is gone — a later unavailable error has nothing to serve. + inner.read.mockRejectedValue(new SecretBackendUnavailableError('bao down')); + await expect(d.read(REQ)).rejects.toThrow(SecretBackendUnavailableError); + }); + + it('does not serve stale for a non-transport error (e.g. revoked grants)', async () => { + const inner = makeInner(); + let now = 1_000; + const d = new CachingSecretBackendDriver(inner, { ttlMs: 1_000, now: () => now }); + + await d.read(REQ); + inner.read.mockRejectedValue(new Error('OpenBao read: HTTP 403 permission denied')); + now += 2_000; + + await expect(d.read(REQ)).rejects.toThrow(/403/); + }); + + it('rethrows on a cold cache even when the backend is unavailable', async () => { + const inner = makeInner({ read: vi.fn(async () => { throw new SecretBackendUnavailableError('bao down'); }) as never }); + const d = new CachingSecretBackendDriver(inner); + + await expect(d.read(REQ)).rejects.toThrow(SecretBackendUnavailableError); + }); + + it('write() refreshes the cache so a read-after-write does not lag', async () => { + const inner = makeInner(); + const d = new CachingSecretBackendDriver(inner, { ttlMs: 60_000 }); + + await d.read(REQ); + await d.write({ name: REQ.name, data: { token: 'brand-new' } }); + + expect(await d.read(REQ)).toEqual({ token: 'brand-new' }); + expect(inner.read).toHaveBeenCalledTimes(1); + }); + + it('delete() evicts', async () => { + const inner = makeInner(); + const d = new CachingSecretBackendDriver(inner, { ttlMs: 60_000 }); + + await d.read(REQ); + await d.delete({ name: REQ.name, externalRef: REQ.externalRef }); + + expect(d.stats().entries).toBe(0); + }); + + it('bounds the map with an LRU eviction', async () => { + const inner = makeInner(); + const d = new CachingSecretBackendDriver(inner, { ttlMs: 60_000, maxEntries: 2 }); + + await d.read({ ...REQ, name: 'a' }); + await d.read({ ...REQ, name: 'b' }); + await d.read({ ...REQ, name: 'a' }); // 'a' becomes most-recently-used + await d.read({ ...REQ, name: 'c' }); // evicts 'b' + + expect(d.stats().entries).toBe(2); + inner.read.mockClear(); + await d.read({ ...REQ, name: 'a' }); + expect(inner.read).not.toHaveBeenCalled(); // 'a' survived + await d.read({ ...REQ, name: 'b' }); + expect(inner.read).toHaveBeenCalledTimes(1); // 'b' was evicted + }); + + it('reports stale count and age via stats()', async () => { + const inner = makeInner(); + let now = 1_000; + const d = new CachingSecretBackendDriver(inner, { ttlMs: 1_000, now: () => now }); + + await d.read({ ...REQ, name: 'a' }); + await d.read({ ...REQ, name: 'b' }); + expect(d.stats()).toMatchObject({ entries: 2, servingStale: 0 }); + + inner.read.mockRejectedValue(new SecretBackendUnavailableError('down')); + now += 2_000; + await d.read({ ...REQ, name: 'a' }); + + expect(d.stats()).toMatchObject({ entries: 2, servingStale: 1, oldestStaleSince: 3_000 }); + }); +});