Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739import { McpOutcome } from "./mcp-outcome";import { MCP_CONTACT_INTERVAL_MS } from "../shared/mcp-health";import { mcpContactFailure, needsMcpMaintenance, mcpRecoveryDecision,} from "./mcp-health";import type { Agent } from "agents";import type * as Cause from "effect/Cause";import { normalizeMcpCatalog } from "./mcp-catalog";import * as Effect from "effect/Effect";import * as Semaphore from "effect/Semaphore";import { mcpObservation, type McpConnection, type McpObservation,} from "../shared/mcp";import { agentCall, agentValidation } from "./agent-io";import type { McpConnections } from "./mcp-connections";import type { McpOAuthStorage } from "./mcp-oauth-storage";import { McpHeaders, mcpHeaderFetch } from "./mcp-headers";import { OperationFailure } from "./operation-result";import type { ExtensionDiagnosticSink } from "../shared/extension-diagnostics";import { mcpDiagnosticOutcome, mcpCauseOutcome, recordMcpDiagnostic, type McpDiagnostic,} from "./mcp-diagnostics";
export class McpRuntime { private readonly outcomes = new Map<string, McpOutcome>(); private readonly maintenanceLock = Semaphore.makeUnsafe(1); private readonly operations = new Map< string, { lock: Semaphore.Semaphore; users: number } >();
private readonly restoring = new Map< string, { attempt: McpConnection; native: Agent["mcp"]["mcpConnections"][string]; diagnostic: McpDiagnostic; } >();
constructor( private readonly connections: McpConnections, private readonly native: Agent["mcp"], private readonly waitUntil: (work: Promise<unknown>) => void, private readonly oauth: McpOAuthStorage, private readonly headers: McpHeaders, private readonly maintenance: (enabled: boolean) => Promise<void>, private readonly diagnostics?: ExtensionDiagnosticSink, ) { this.native.onServerStateChanged(() => this.projectRestored()); }
add(value: unknown) { return Effect.gen({ self: this }, function* () { const connection = yield* agentValidation(() => this.connections.add(value), ); yield* this.syncMaintenance(); return yield* this.connect(connection.id, connection.revision); }); }
setEnabled(id: string, revision: number, enabled: unknown) { return Effect.gen({ self: this }, function* () { const connection = yield* agentValidation(() => this.connections.setEnabled(id, revision, enabled), ); if (connection.enabled) { yield* this.syncMaintenance(); return yield* this.connect(id, connection.revision); } this.oauth.retire(id); this.launch(this.syncMaintenance()); this.launch(this.exclusive(id, this.disconnect(id))); return connection; }); }
remove(id: string, revision: number) { return agentValidation(() => { const ref = this.connections.credentialReference(id); const result = this.connections.remove(id, revision); this.oauth.retire(id); this.launch(this.headers.discard(ref)); this.launch(this.syncMaintenance()); this.launch(this.exclusive(id, this.disconnect(id))); return result; }); }
initializeCredentials() { return this.headers .prune(this.connections.credentialReferences()) .pipe(Effect.andThen(this.syncMaintenance())); }
setHeaders(id: string, revision: number, value: unknown) { return Effect.gen({ self: this }, function* () { const current = yield* agentValidation(() => this.connections.get(id)); if (current.authMode !== "headers") return yield* Effect.fail(new OperationFailure("mcp_headers_invalid")); const prepared = yield* this.headers.prepare(id, value); const updated = yield* agentValidation(() => this.connections.replaceHeaders( id, revision, prepared.ref, prepared.names, ), ).pipe( Effect.onError(() => this.headers .discard(prepared.ref) .pipe(Effect.catchCause(() => Effect.void)), ), ); this.launch(this.headers.discard(updated.previous)); return updated.connection.enabled ? yield* this.connect(id, updated.connection.revision) : updated.connection; }); }
connect( id: string, revision: number, operation: McpDiagnostic["operation"] = "mcp-connect", ) { return agentValidation(() => { const attempt = this.connections.beginConnect(id, revision); this.launch( this.exclusive( id, this.connectAttempt(attempt, undefined, { operation, startedAt: Date.now(), }), ), ); return attempt; }); }
refresh(id: string, revision: number) { return this.connect(id, revision, "mcp-refresh"); }
private syncMaintenance() { return this.maintenanceLock.withPermits(1)( agentCall(() => this.maintenance(this.connections.list().some(needsMcpMaintenance)), ), ); }
maintain() { return Effect.forEach( this.connections.list().filter((item) => item.enabled), (item) => this.exclusive( item.id, Effect.gen({ self: this }, function* () { const current = this.connections.find(item.id); if (!current?.enabled || current.revision !== item.revision) return; if ( current.state === "error" && current.lastError === "connection_failed" ) { if (mcpRecoveryDecision(current) === "connect") { const attempt = this.connections.beginConnect( current.id, current.revision, true, ); yield* this.connectAttempt(attempt); } return; } if ( current.state !== "ready" || (current.health.lastCheckedAt && Date.now() - Date.parse(current.health.lastCheckedAt) < MCP_CONTACT_INTERVAL_MS) ) return; const native = this.native.mcpConnections[current.id]; const diagnostic: McpDiagnostic = { operation: "mcp-refresh", startedAt: Date.now(), }; if (!native || native.connectionState === "failed") { this.observe( current, { state: "error", lastError: "connection_failed", capabilities: current.capabilities, }, diagnostic, ); return; } const error = yield* Effect.tryPromise({ try: (signal) => native.client.ping({ signal, timeout: 10_000 }), catch: mcpContactFailure, }).pipe( Effect.as(null), Effect.catch((failure) => Effect.succeed(failure)), ); if (error) { const latest = this.connections.find(current.id); const lastError = latest?.lastError === "authentication_failed" ? latest.lastError : error; this.observe( current, { state: "error", lastError, capabilities: current.capabilities, }, diagnostic, ); } else this.connections.contact(current.id, current.revision); }), ), { concurrency: "unbounded", discard: true }, ); }
// The SDK restores transports independently. Observe each exact restored // connection; never put unrelated server locks behind an installation-wide wait. recover() { for (const item of this.connections.list()) { if (!item.enabled) continue; const native = this.native.mcpConnections[item.id]; // Persisted SDK transports have already started native restoration. For // application-owned transports, startup shares periodic retry admission. if (!native && mcpRecoveryDecision(item) !== "connect") continue; const attempt = this.connections.beginConnect( item.id, item.revision, true, ); if (native) this.restoring.set(item.id, { attempt, native, diagnostic: { operation: "mcp-connect", startedAt: Date.now() }, }); else this.launch(this.exclusive(item.id, this.connectAttempt(attempt))); } this.projectRestored(); // The SDK returns failed discovery to `connected`; its completion barrier // distinguishes that terminal outcome from initialization before discovery. // This observer owns no server lock and never delays independent projection // or cleanup. this.launch( agentCall(() => this.native.waitForConnections()).pipe( Effect.andThen(Effect.sync(() => this.projectRestored(true))), ), ); for (const server of this.native.listServers()) { if (!this.connections.find(server.id)?.enabled) this.launch(this.exclusive(server.id, this.disconnect(server.id))); } // A crash can leave credentials after the SDK registration and desired row // have already been removed. Recover cleanup from private vault references. this.launch( this.oauth.storedServers().pipe( Effect.flatMap((ids) => Effect.forEach( ids, (id) => this.exclusive( id, Effect.suspend(() => this.connections.find(id)?.enabled ? Effect.void : this.disconnect(id), ), ), { concurrency: "unbounded" }, ), ), ), ); }
private projectRestored(settled = false) { for (const [id, restored] of this.restoring) { const current = this.connections.find(id); if ( !current?.enabled || current.revision !== restored.attempt.revision || this.native.mcpConnections[id] !== restored.native ) { this.restoring.delete(id); continue; } const state = restored.native.connectionState; if (state === "ready") { this.observe(restored.attempt, this.catalog(id), restored.diagnostic); this.restoring.delete(id); } else if (settled && state === "connected") { this.observe( restored.attempt, { state: "error", lastError: "connection_failed", capabilities: restored.attempt.capabilities, }, restored.diagnostic, // Recovery retains its existing retry admission; diagnostics know // discovery was the stage that failed after transport connected. "discovery", ); this.restoring.delete(id); } else if (state === "authenticating" || state === "failed") { this.observe( restored.attempt, { state: state === "authenticating" ? "authenticating" : "error", lastError: state === "authenticating" ? "authentication_required" : "connection_failed", capabilities: restored.attempt.capabilities, }, restored.diagnostic, ); this.restoring.delete(id); } } }
private launch<E>(work: Effect.Effect<unknown, E>) { this.waitUntil( Effect.runPromise(work.pipe(Effect.catchCause(() => Effect.void))), ); }
private exclusive<A, E>(id: string, work: Effect.Effect<A, E>) { return Effect.suspend(() => { const entry = this.operations.get(id) ?? { lock: Semaphore.makeUnsafe(1), users: 0, }; entry.users++; this.operations.set(id, entry); return entry.lock .withPermits(1)(work) .pipe( Effect.ensuring( Effect.sync(() => { if (--entry.users === 0) this.operations.delete(id); }), ), ); }); }
authorize(id: string, revision: number, callbackUrl: string) { return Effect.gen({ self: this }, function* () { const attempt = yield* agentValidation(() => { const item = this.connections.get(id); if (item.authMode !== "oauth") throw new Error("OAuth is not configured"); return this.connections.beginConnect(id, revision); }); return yield* this.exclusive( id, Effect.gen({ self: this }, function* () { yield* this.connectAttempt(attempt, callbackUrl, { operation: "mcp-auth", startedAt: Date.now(), }); const current = this.connections.find(id); if (!current?.enabled || current.revision !== attempt.revision) return null; if (current.state === "ready") return { state: "ready" as const }; const url = this.native .listServers() .find((server) => server.id === id)?.auth_url; return url ? { state: "authorize" as const, url } : null; }), ); }); }
callback(request: Request) { const diagnostic: McpDiagnostic = { operation: "mcp-auth", startedAt: Date.now(), }; let accepted: McpConnection | undefined; const id = new URL(request.url).searchParams.get("state")?.split(".")[1] ?? ""; return this.exclusive( id, Effect.gen({ self: this }, function* () { const current = this.connections.find(id); if ( !current?.enabled || current.authMode !== "oauth" || !this.native.isCallbackRequest(request) ) return false; const provider = this.native.mcpConnections[id]?.options.transport.authProvider; const state = new URL(request.url).searchParams.get("state")!; if ( !provider || !(yield* agentCall(() => provider.checkState(state))).valid ) return false; accepted = current; const result = yield* agentCall(() => this.native.handleCallbackRequest(request), ); const latest = this.connections.find(id); if (!latest?.enabled || latest.revision !== current.revision) { yield* this.disconnect(id); return false; } const observation = result.authSuccess ? yield* this.connectRegistered(current) : ({ state: "error", lastError: "authentication_required", capabilities: current.capabilities, } satisfies McpObservation); this.observe(current, observation, diagnostic); const observed = this.connections.find(id); return ( !!observed?.enabled && observed.revision === current.revision && observed.state === "ready" ); }).pipe( Effect.onError((cause) => Effect.sync(() => { if (accepted) this.reportCurrentFailure( accepted, diagnostic, cause, "authentication", ); }), ), ), ); }
private disconnect(id: string) { this.outcomes.delete(id); const diagnostic: McpDiagnostic = { operation: "mcp-disconnect", startedAt: Date.now(), }; return this.oauth.clear(id).pipe( Effect.andThen(agentCall(() => this.native.removeServer(id))), Effect.tap(() => Effect.sync(() => { if (!this.connections.find(id)?.enabled) recordMcpDiagnostic(this.diagnostics, id, diagnostic, { status: "completed", }); }), ), Effect.onError((cause) => Effect.sync(() => { if (!this.connections.find(id)?.enabled) recordMcpDiagnostic( this.diagnostics, id, diagnostic, mcpCauseOutcome(cause, "connection"), ); }), ), ); }
private connectAttempt( attempt: McpConnection, callbackUrl?: string, diagnostic: McpDiagnostic = { operation: "mcp-connect", startedAt: Date.now(), }, ) { return Effect.gen({ self: this }, function* () { const current = this.connections.find(attempt.id); if (!current?.enabled || current.revision !== attempt.revision) return; const redirect = callbackUrl ?? this.native.listServers().find((server) => server.id === attempt.id) ?.callback_url; const observation = attempt.authMode === "oauth" && !redirect ? ({ state: "authenticating", capabilities: attempt.capabilities, lastError: "authentication_required", } satisfies McpObservation) : yield* ( attempt.authMode === "headers" ? this.connectHeaders(attempt) : this.connectServer(attempt, redirect) ).pipe( Effect.catch((error) => Effect.succeed<McpObservation>({ state: "error", lastError: error instanceof OperationFailure && error.code === "mcp_credentials_unreadable" ? "credentials_unreadable" : "connection_failed", capabilities: attempt.capabilities, }), ), ); this.observe(attempt, observation, diagnostic); }).pipe( Effect.onError((cause) => Effect.sync(() => this.reportCurrentFailure(attempt, diagnostic, cause, "connection"), ), ), ); }
private reportCurrentFailure( attempt: McpConnection, diagnostic: McpDiagnostic, cause: Cause.Cause<unknown>, failure: "authentication" | "connection", ) { const current = this.connections.find(attempt.id); if (current?.enabled && current.revision === attempt.revision) recordMcpDiagnostic( this.diagnostics, attempt.id, diagnostic, mcpCauseOutcome(cause, failure), ); }
private observe( attempt: McpConnection, observation: McpObservation, diagnostic: McpDiagnostic = { operation: "mcp-refresh", startedAt: Date.now(), }, diagnosticFailure?: "discovery", ) { const validated = mcpObservation.safeParse(observation); const outcome: McpObservation = validated.success ? validated.data : { state: "error", lastError: "discovery_failed", capabilities: attempt.capabilities, }; const applied = this.connections.observe( attempt.id, attempt.revision, outcome, ); if (applied) { recordMcpDiagnostic( this.diagnostics, attempt.id, diagnostic, diagnosticFailure ? { status: "error", failure: diagnosticFailure } : mcpDiagnosticOutcome(outcome), ); this.launch(this.syncMaintenance()); } }
private connectHeaders(connection: McpConnection) { return Effect.gen({ self: this }, function* () { yield* agentCall(() => this.native.removeServer(connection.id)); if (!connection.headerNames.length) return { state: "authenticating", lastError: "authentication_required", capabilities: connection.capabilities, } satisfies McpObservation; const ref = this.connections.credentialReference(connection.id); const credentials = yield* this.headers.resolve(ref); const outcome = new McpOutcome(); this.outcomes.set(connection.id, outcome); const transport = mcpHeaderFetch( connection.endpoint, credentials, () => { const current = this.connections.find(connection.id); return !!current?.enabled && current.revision === connection.revision; }, () => { const current = this.connections.find(connection.id); if (current?.state === "ready") this.observe(connection, { state: "error", lastError: "authentication_failed", capabilities: current.capabilities, }); }, ); // Agents 0.22's compatibility connect() is the public in-memory path. // registerServer() persists transport options; secret connections instead // restore from our encrypted reference, never a native SQL header record. const started = Date.now(); const success = yield* agentCall(() => this.native.connect(connection.endpoint, { reconnect: { id: connection.id }, transport: { type: "auto", fetch: outcome.fetch(transport) }, }), ).pipe( Effect.as(true), Effect.catch(() => Effect.succeed(false)), ); if (!success) return { state: "error", // The pinned compatibility API owns discovery with a 15s timeout. // It flattens that timeout too; elapsed admission distinguishes it // from a prompt protocol rejection without parsing provider text. lastError: outcome.error(Date.now() - started, 15_000), capabilities: connection.capabilities, } satisfies McpObservation; return this.catalog(connection.id); }); }
private connectServer(connection: McpConnection, callbackUrl?: string) { return Effect.gen({ self: this }, function* () { const previous = this.native .listServers() .find((server) => server.id === connection.id); const sameCallback = previous?.callback_url === callbackUrl; this.oauth.retire(connection.id); if (previous?.callback_url && !sameCallback) yield* this.oauth.clear(connection.id); yield* agentCall(() => this.native.removeServer(connection.id)); const outcome = new McpOutcome(); this.outcomes.set(connection.id, outcome); yield* agentCall(() => this.native.registerServer(connection.id, { name: connection.name, url: connection.endpoint, callbackUrl, clientId: sameCallback ? (previous?.client_id ?? undefined) : undefined, transport: { type: "auto", fetch: outcome.fetch() }, retry: { maxAttempts: 1 }, }), ); return yield* this.connectRegistered(connection); }); }
private connectRegistered(connection: McpConnection) { return Effect.gen({ self: this }, function* () { this.outcomes.get(connection.id)?.reset(); const result = yield* agentCall(() => this.native.connectToServer(connection.id), ); if (result.state === "authenticating") return { state: "authenticating", lastError: "authentication_required", capabilities: connection.capabilities, } satisfies McpObservation; if (result.state !== "connected") return { state: "error", lastError: this.outcomes.get(connection.id)?.error() ?? "connection_failed", capabilities: connection.capabilities, } satisfies McpObservation; return yield* this.discover(connection); }); }
private discover(connection: McpConnection) { return Effect.gen({ self: this }, function* () { const outcome = this.outcomes.get(connection.id); outcome?.reset(); const started = Date.now(); const timeoutMs = 10_000; const discovery = yield* agentCall(() => this.native.discoverIfConnected(connection.id, { timeoutMs }), ); if (!discovery?.success) return { state: "error", lastError: outcome?.error(Date.now() - started, timeoutMs) ?? "connection_failed", capabilities: connection.capabilities, } satisfies McpObservation; return this.catalog(connection.id); }); }
private catalog(id: string): McpObservation { const filter = { serverId: id }; return { state: "ready", lastError: null, ...normalizeMcpCatalog( id, { tools: this.native.listTools(filter), resources: this.native.listResources(filter), prompts: this.native.listPrompts(filter), }, this.connections.get(id).capabilities, ), } satisfies McpObservation; }}