diff --git a/apps/server/src/provider/Layers/PrimeAgentAdapter.test.ts b/apps/server/src/provider/Layers/PrimeAgentAdapter.test.ts index 2009a6ef3..b344aea87 100644 --- a/apps/server/src/provider/Layers/PrimeAgentAdapter.test.ts +++ b/apps/server/src/provider/Layers/PrimeAgentAdapter.test.ts @@ -16,6 +16,7 @@ import * as Schema from "effect/Schema"; import * as Stream from "effect/Stream"; import { + EnvironmentId, PrimeAgentSettings, ProviderDriverKind, ProviderInstanceId, @@ -25,6 +26,7 @@ import { } from "@t3tools/contracts"; import { ServerConfig } from "../../config.ts"; +import * as McpProviderSession from "../../mcp/McpProviderSession.ts"; import { makePrimeAgentAdapter, parsePrimeAgentResumeMarker, @@ -126,6 +128,19 @@ exec ${process.execPath} ${mockAgentPath} "$@" assert.equal(rejectedLaunchArgs._tag, "Failure"); const threadId = ThreadId.make("resume/thread"); + const mcpSession = { + providerSessionId: "provider-session-prime-acp-test", + threadId, + environmentId: EnvironmentId.make("environment-prime-acp-test"), + providerInstanceId: ProviderInstanceId.make("primeAgent"), + endpoint: "http://127.0.0.1:4321/mcp/provider-session-prime-acp-test", + authorizationHeader: "Bearer scoped-secret", + expiresAt: 4_000_000_000_000, + }; + yield* Effect.acquireRelease( + Effect.sync(() => McpProviderSession.setMcpProviderSession(mcpSession)), + () => Effect.sync(() => McpProviderSession.clearMcpProviderSession(threadId)), + ); const startupWarning = yield* Deferred.make(); const unavailableResources = yield* Deferred.make< @@ -316,6 +331,25 @@ exec ${process.execPath} ${mockAgentPath} "$@" .map((line) => JSON.parse(line) as unknown), ), ); + const newSessionEntry = requestEntries.find( + (entry) => + typeof entry === "object" && + entry !== null && + "method" in entry && + entry.method === "session/new", + ); + assert.isDefined(newSessionEntry); + const newSessionParams = ( + newSessionEntry as { readonly params?: { readonly mcpServers?: unknown } } + ).params; + assert.deepEqual(newSessionParams?.mcpServers, [ + { + name: "t3-code", + type: "http", + url: mcpSession.endpoint, + headers: [{ name: "Authorization", value: mcpSession.authorizationHeader }], + }, + ]); assert.isTrue( requestEntries.some( (entry) => diff --git a/apps/server/src/provider/Layers/PrimeAgentAdapter.ts b/apps/server/src/provider/Layers/PrimeAgentAdapter.ts index 9e051f69d..352487181 100644 --- a/apps/server/src/provider/Layers/PrimeAgentAdapter.ts +++ b/apps/server/src/provider/Layers/PrimeAgentAdapter.ts @@ -27,6 +27,7 @@ import type * as EffectAcpSchema from "effect-acp/schema"; import { resolveAttachmentPath } from "../../attachmentStore.ts"; import { ServerConfig } from "../../config.ts"; +import * as McpProviderSession from "../../mcp/McpProviderSession.ts"; import { ProviderAdapterProcessError, ProviderAdapterRequestError, @@ -518,6 +519,7 @@ export function makePrimeAgentAdapter( provider: PROVIDER, threadId: input.threadId, }); + const mcpSession = McpProviderSession.readMcpProviderSession(input.threadId); const acp = yield* makePrimeAgentAcpRuntime({ primeAgentSettings, ...(options?.environment ? { environment: options.environment } : {}), @@ -526,6 +528,23 @@ export function makePrimeAgentAdapter( sessionDir, continueSession: parsePrimeAgentResumeMarker(input.resumeCursor), model, + ...(mcpSession === undefined + ? {} + : { + mcpServers: [ + { + type: "http" as const, + name: "t3-code", + url: mcpSession.endpoint, + headers: [ + { + name: "Authorization", + value: mcpSession.authorizationHeader, + }, + ], + }, + ], + }), clientInfo: { name: "pylon", version: "0.0.0" }, observeSessionUpdate: (notification) => { const update = parsePrimeAgentAcpTerminalUpdate(notification); diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index 688424cab..36e0eb89d 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -16,6 +16,7 @@ import type { } from "@t3tools/contracts"; import { ApprovalRequestId, + EnvironmentId, EventId, ProviderDriverKind, ProviderInstanceId, @@ -29,6 +30,7 @@ import { import { createModelSelection } from "@t3tools/shared/model"; import { it, assert, describe, vi } from "@effect/vitest"; +import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; @@ -59,6 +61,7 @@ import * as ProviderEventLoggers from "./ProviderEventLoggers.ts"; import { ProviderSessionDirectoryLive } from "./ProviderSessionDirectory.ts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import * as ProviderSessionRuntime from "../../persistence/ProviderSessionRuntime.ts"; +import * as McpProviderSession from "../../mcp/McpProviderSession.ts"; import { makeSqlitePersistenceLive, SqlitePersistenceMemory, @@ -2551,40 +2554,46 @@ validation.layer("ProviderServiceLive validation", (it) => { describe("agent browser access", () => { const revokedThreads: Array = []; + const makeAgentBrowserProviderLayer = ( + enableAgentBrowserAccess: boolean, + codex: ReturnType, + options: NonNullable[0]>, + ) => { + const providerAdapterLayer = Layer.succeed( + ProviderAdapterRegistry.ProviderAdapterRegistry, + makeAdapterRegistryMock({ [CODEX_DRIVER]: codex.adapter }), + ); + const runtimeRepositoryLayer = ProviderSessionRuntime.layer.pipe( + Layer.provide(SqlitePersistenceMemory), + ); + const directoryLayer = ProviderSessionDirectoryLive.pipe(Layer.provide(runtimeRepositoryLayer)); + return makeProviderServiceLive(options).pipe( + Layer.provide(providerAdapterLayer), + Layer.provide(directoryLayer), + Layer.provide(ServerSettings.ServerSettingsService.layerTest({ enableAgentBrowserAccess })), + Layer.provide(serverConfigTestLayer), + Layer.provide(AnalyticsService.layerTest), + Layer.provide( + Layer.succeed( + ProviderEventLoggers.ProviderEventLoggers, + ProviderEventLoggers.NoOpProviderEventLoggers, + ), + ), + ); + }; + const startSessionWith = (enableAgentBrowserAccess: boolean, threadId: ThreadId) => Effect.gen(function* () { const issued: Array = []; const codex = makeFakeCodexAdapter(); - const providerAdapterLayer = Layer.succeed( - ProviderAdapterRegistry.ProviderAdapterRegistry, - makeAdapterRegistryMock({ [CODEX_DRIVER]: codex.adapter }), - ); - const runtimeRepositoryLayer = ProviderSessionRuntime.layer.pipe( - Layer.provide(SqlitePersistenceMemory), - ); - const directoryLayer = ProviderSessionDirectoryLive.pipe( - Layer.provide(runtimeRepositoryLayer), - ); - const providerLayer = makeProviderServiceLive({ + const providerLayer = makeAgentBrowserProviderLayer(enableAgentBrowserAccess, codex, { issueMcpCredential: (request) => Effect.sync(() => { issued.push(request.threadId); return undefined; }), revokeMcpCredential: (revoked) => Effect.sync(() => void revokedThreads.push(revoked)), - }).pipe( - Layer.provide(providerAdapterLayer), - Layer.provide(directoryLayer), - Layer.provide(ServerSettings.ServerSettingsService.layerTest({ enableAgentBrowserAccess })), - Layer.provide(serverConfigTestLayer), - Layer.provide(AnalyticsService.layerTest), - Layer.provide( - Layer.succeed( - ProviderEventLoggers.ProviderEventLoggers, - ProviderEventLoggers.NoOpProviderEventLoggers, - ), - ), - ); + }); yield* Effect.gen(function* () { const provider = yield* ProviderService.ProviderService; @@ -2599,6 +2608,17 @@ describe("agent browser access", () => { return issued; }); + const issuedBrowserCredential = (threadId: ThreadId) => ({ + config: { + environmentId: EnvironmentId.make("environment-browser-test"), + threadId, + providerSessionId: `provider-session-${threadId}`, + providerInstanceId: codexInstanceId, + endpoint: `http://127.0.0.1:4321/mcp/provider-session-${threadId}`, + authorizationHeader: "Bearer scoped-secret", + }, + }); + // Credential issuance is the observable that matters: it is the only place a // credential is minted, and `/mcp` accepts nothing else, so withholding it is // what actually denies every provider and external MCP client. @@ -2633,4 +2653,80 @@ describe("agent browser access", () => { assert.deepEqual(issued, [threadId]); }).pipe(Effect.provide(NodeServices.layer)), ); + + it.effect("revokes the MCP credential when an adapter session exits on its own", () => + Effect.gen(function* () { + const threadId = asThreadId("thread-browser-terminal"); + const codex = makeFakeCodexAdapter(); + const revoked = yield* Deferred.make(); + const providerLayer = makeAgentBrowserProviderLayer(true, codex, { + issueMcpCredential: (request) => Effect.succeed(issuedBrowserCredential(request.threadId)), + revokeMcpCredential: (revokedThreadId) => + Deferred.succeed(revoked, revokedThreadId).pipe(Effect.asVoid), + }); + + yield* Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + yield* provider.startSession(threadId, { + provider: CODEX_DRIVER, + providerInstanceId: codexInstanceId, + threadId, + runtimeMode: "full-access", + }); + assert.isDefined(McpProviderSession.readMcpProviderSession(threadId)); + yield* codex.stopSession(threadId); + yield* Effect.yieldNow; + codex.emit({ + type: "session.exited", + eventId: asEventId("evt-browser-terminal"), + provider: CODEX_DRIVER, + createdAt: "2026-01-01T00:00:00.000Z", + threadId, + payload: { exitKind: "error" }, + }); + + assert.equal(yield* Deferred.await(revoked), threadId); + assert.isUndefined(McpProviderSession.readMcpProviderSession(threadId)); + }).pipe(Effect.provide(providerLayer)); + }).pipe(Effect.provide(NodeServices.layer)), + ); + + it.effect("revokes the MCP credential even when explicit adapter stop fails", () => + Effect.gen(function* () { + const threadId = asThreadId("thread-browser-stop-failure"); + const codex = makeFakeCodexAdapter(); + const revoked = yield* Deferred.make(); + codex.stopSession.mockImplementationOnce(() => + Effect.fail( + new ProviderAdapterRequestError({ + provider: CODEX_DRIVER, + method: "stopSession", + detail: "synthetic stop failure", + }), + ), + ); + const providerLayer = makeAgentBrowserProviderLayer(true, codex, { + issueMcpCredential: (request) => Effect.succeed(issuedBrowserCredential(request.threadId)), + revokeMcpCredential: (revokedThreadId) => + Deferred.succeed(revoked, revokedThreadId).pipe(Effect.asVoid), + }); + + yield* Effect.gen(function* () { + const provider = yield* ProviderService.ProviderService; + yield* provider.startSession(threadId, { + provider: CODEX_DRIVER, + providerInstanceId: codexInstanceId, + threadId, + runtimeMode: "full-access", + }); + assert.isDefined(McpProviderSession.readMcpProviderSession(threadId)); + + const stopped = yield* provider.stopSession({ threadId }).pipe(Effect.result); + + assert.equal(stopped._tag, "Failure"); + assert.equal(yield* Deferred.await(revoked), threadId); + assert.isUndefined(McpProviderSession.readMcpProviderSession(threadId)); + }).pipe(Effect.provide(providerLayer)); + }).pipe(Effect.provide(NodeServices.layer)), + ); }); diff --git a/apps/server/src/provider/Layers/ProviderService.ts b/apps/server/src/provider/Layers/ProviderService.ts index a088486c1..6adb0c463 100644 --- a/apps/server/src/provider/Layers/ProviderService.ts +++ b/apps/server/src/provider/Layers/ProviderService.ts @@ -101,7 +101,7 @@ export interface ProviderServiceLiveOptions { * test see whether a credential was requested at all. */ readonly issueMcpCredential?: typeof McpSessionRegistry.issueActiveMcpCredential; - /** Same seam as `issueMcpCredential`, for observing the deny path's revoke. */ + /** Same seam as `issueMcpCredential`, for observing session credential revocation. */ readonly revokeMcpCredential?: typeof McpSessionRegistry.revokeActiveMcpThread; } @@ -305,9 +305,12 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( return credential; }); const clearMcpSession = (threadId: ThreadId) => - McpSessionRegistry.revokeActiveMcpThread(threadId).pipe( - Effect.tap(() => Effect.sync(() => McpProviderSession.clearMcpProviderSession(threadId))), + revokeMcpCredential(threadId).pipe( + Effect.ensuring(Effect.sync(() => McpProviderSession.clearMcpProviderSession(threadId))), ); + const clearAllMcpSessions = McpSessionRegistry.revokeAllActiveMcpCredentials().pipe( + Effect.ensuring(Effect.sync(() => McpProviderSession.clearAllMcpProviderSessions())), + ); const publishRuntimeEvent = (event: ProviderRuntimeEvent): Effect.Effect => Effect.succeed(event).pipe( @@ -363,22 +366,6 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( }); }); - const processRuntimeEvent = ( - source: { - readonly instanceId: ProviderInstanceId; - readonly provider: ProviderDriverKind; - }, - event: ProviderRuntimeEvent, - ): Effect.Effect => - Effect.sync(() => correlateRuntimeEventWithInstance(source, event)).pipe( - Effect.flatMap((canonicalEvent) => - increment(providerRuntimeEventsTotal, { - provider: canonicalEvent.provider, - eventType: canonicalEvent.type, - }).pipe(Effect.andThen(publishRuntimeEvent(canonicalEvent))), - ), - ); - // `subscribedAdapters` is our source-of-truth for "which instance adapters // are currently wired into the runtime event bus". It both tracks the set // of live subscriptions (so `reconcileInstanceSubscriptions` can diff and @@ -395,6 +382,36 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( Effect.map((map) => Array.from(map.entries())), ); + const processRuntimeEvent = ( + source: { + readonly instanceId: ProviderInstanceId; + readonly provider: ProviderDriverKind; + readonly adapter: ProviderAdapterShape; + }, + event: ProviderRuntimeEvent, + ): Effect.Effect => + Effect.sync(() => correlateRuntimeEventWithInstance(source, event)).pipe( + Effect.flatMap((canonicalEvent) => + increment(providerRuntimeEventsTotal, { + provider: canonicalEvent.provider, + eventType: canonicalEvent.type, + }).pipe( + Effect.andThen(publishRuntimeEvent(canonicalEvent)), + Effect.andThen( + Effect.gen(function* () { + if (canonicalEvent.type !== "session.exited") return; + const mcpSession = McpProviderSession.readMcpProviderSession(canonicalEvent.threadId); + if (mcpSession?.providerInstanceId !== source.instanceId) return; + const currentAdapters = yield* Ref.get(subscribedAdapters); + if (currentAdapters.get(source.instanceId) !== source.adapter) return; + const stillActive = yield* source.adapter.hasSession(canonicalEvent.threadId); + if (!stillActive) yield* clearMcpSession(canonicalEvent.threadId); + }), + ), + ), + ), + ); + // Rebuild the map of id → adapter from the registry and fork a new event // subscription for every instance that is either brand new or whose adapter // identity changed (indicating the underlying `ProviderInstance` was torn @@ -418,6 +435,7 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( { instanceId: id, provider: adapter.provider, + adapter, }, event, ), @@ -1630,9 +1648,12 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( "provider.thread_id": input.threadId, }); if (routed.isActive) { - yield* routed.adapter.stopSession(routed.threadId); + yield* routed.adapter + .stopSession(routed.threadId) + .pipe(Effect.ensuring(clearMcpSession(input.threadId))); + } else { + yield* clearMcpSession(input.threadId); } - yield* clearMcpSession(input.threadId); yield* directory.upsert({ threadId: input.threadId, provider: routed.adapter.provider, @@ -1853,9 +1874,10 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( }), ), ).pipe(Effect.asVoid); - yield* Effect.forEach(currentAdapters, ([, adapter]) => adapter.stopAll()).pipe(Effect.asVoid); - yield* McpSessionRegistry.revokeAllActiveMcpCredentials(); - McpProviderSession.clearAllMcpProviderSessions(); + yield* Effect.forEach(currentAdapters, ([, adapter]) => adapter.stopAll()).pipe( + Effect.asVoid, + Effect.ensuring(clearAllMcpSessions), + ); const bindings = yield* directory.listBindings().pipe(Effect.orElseSucceed(() => [])); yield* Effect.forEach(bindings, (binding) => Effect.gen(function* () { diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts b/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts index 051db355f..1bc7d3938 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts @@ -5,6 +5,7 @@ import * as NodeServices from "@effect/platform-node/NodeServices"; import { describe, expect, it } from "@effect/vitest"; import { ApprovalRequestId, + EnvironmentId, PrimeAgentSettings, PROVIDER_SESSION_AGENT_MESSAGE_MAX_CHARS, ProviderDriverKind, @@ -28,6 +29,7 @@ import * as Stream from "effect/Stream"; import * as TestClock from "effect/testing/TestClock"; import { ServerConfig } from "../../config.ts"; +import * as McpProviderSession from "../../mcp/McpProviderSession.ts"; import type { ProviderAdapterError } from "../Errors.ts"; import { attachmentRelativePath } from "../../attachmentStore.ts"; import type { @@ -930,6 +932,45 @@ describe("PrimeAgentDaemonAdapter", () => { }), ).pipe(Effect.provide(testLayer)), ); + it.effect("passes the thread-scoped Pylon MCP server into the daemon runtime", () => + Effect.scoped( + Effect.gen(function* () { + const captures = makeCaptures(); + const mcpSession = { + providerSessionId: "provider-session-prime-test", + threadId, + environmentId: EnvironmentId.make("environment-prime-test"), + providerInstanceId: instanceId, + endpoint: "http://127.0.0.1:4321/mcp/provider-session-prime-test", + authorizationHeader: "Bearer scoped-secret", + expiresAt: 4_000_000_000_000, + }; + yield* Effect.acquireRelease( + Effect.sync(() => McpProviderSession.setMcpProviderSession(mcpSession)), + () => Effect.sync(() => McpProviderSession.clearMcpProviderSession(threadId)), + ); + const adapter = yield* makePrimeAgentDaemonAdapter(decodeSettings({}), manager, { + instanceId, + runtimeFactory: fakeRuntimeFactory(captures), + }); + + yield* adapter.startSession({ threadId, cwd: process.cwd(), runtimeMode: "full-access" }); + + expect(captures.runtimeInputs).toHaveLength(1); + expect(captures.runtimeInputs[0]?.mcpServer).toEqual({ + ownerId: `pylon:${mcpSession.providerSessionId}`, + server: { + name: "t3-code", + type: "http", + url: mcpSession.endpoint, + headers: { Authorization: mcpSession.authorizationHeader }, + }, + }); + yield* adapter.stopSession(threadId); + }), + ).pipe(Effect.provide(testLayer)), + ); + it.effect("fails closed when the loaded managed extension source changes", () => Effect.scoped( Effect.gen(function* () { diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts b/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts index 100ebad4a..1329894d4 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts @@ -53,6 +53,7 @@ import * as SynchronizedRef from "effect/SynchronizedRef"; import { resolveAttachmentPath } from "../../attachmentStore.ts"; import { ServerConfig } from "../../config.ts"; +import * as McpProviderSession from "../../mcp/McpProviderSession.ts"; import { resolveProviderHomePath } from "../../pathExpansion.ts"; import { ProviderAdapterProcessError, @@ -2348,10 +2349,24 @@ export function makePrimeAgentDaemonAdapter( scopeTransferred ? Effect.void : Scope.close(sessionScope, Exit.void), ); const agentDir = primeAgentSettings.agentHomePath.trim(); + const mcpSession = McpProviderSession.readMcpProviderSession(input.threadId); const runtime = yield* runtimeFactory({ manager, cwd, sessionDir, + ...(mcpSession === undefined + ? {} + : { + mcpServer: { + ownerId: `pylon:${mcpSession.providerSessionId}`, + server: { + name: "t3-code", + type: "http" as const, + url: mcpSession.endpoint, + headers: { Authorization: mcpSession.authorizationHeader }, + }, + }, + }), ...(agentDir.length === 0 ? {} : { agentDir: resolveProviderHomePath(agentDir) }), ...(model === "default" ? {} : { model }), ...(approvalRequired diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonBridge.ts b/apps/server/src/provider/prime/PrimeAgentDaemonBridge.ts index 3449d39b7..6b934f01f 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonBridge.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonBridge.ts @@ -68,6 +68,14 @@ export interface PrimeAgentDaemonPromptOptions { readonly signal?: AbortSignal; } +/** Public Prime Agent shape used to attach Pylon's scoped HTTP MCP server. */ +export interface PrimeAgentDaemonAcpMcpServer { + readonly name: string; + readonly type: "http"; + readonly url: string; + readonly headers: Record; +} + export type PrimeAgentDaemonExtensionUiResponse = | { readonly value: string } | { readonly confirmed: boolean } @@ -156,6 +164,15 @@ export interface PrimeAgentDaemonAgentConnection { ) => Promise; readonly getCommands: () => Promise; readonly getResourceSnapshot: () => Promise; + readonly supportsAcpMcpServers?: () => boolean; + readonly replaceAcpMcpServers?: ( + servers: ReadonlyArray, + ownerId: string, + ) => Promise; + readonly releaseAcpMcpServers?: ( + ownerId: string, + serverNames: ReadonlyArray, + ) => Promise; readonly getModelCatalog?: () => Promise; readonly getAvailableModels?: () => Promise; readonly reload: () => Promise; diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.test.ts b/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.test.ts index 433e21e82..6d45e1968 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.test.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.test.ts @@ -32,6 +32,7 @@ import { primeAgentLiveActivityToolLabel, sanitizePrimeAgentLiveActivityMessages, type PrimeAgentDaemonSessionRuntime, + type PrimeAgentDaemonSessionRuntimeInput, } from "./PrimeAgentDaemonSessionRuntime.ts"; import { PRIME_AGENT_ACP_RESUME_CURSOR } from "./PrimeAgentResumeCursor.ts"; @@ -119,6 +120,10 @@ function fixture(options?: { readonly duringResourceSnapshot?: ReadonlyArray; readonly afterSnapshotEvent?: unknown; readonly attachFailure?: boolean; + readonly omitMcpSupport?: boolean; + readonly mcpSupported?: boolean; + readonly replaceMcpImpl?: () => Promise; + readonly releaseMcpImpl?: () => Promise; readonly resourceSnapshot?: unknown; readonly commands?: unknown; readonly modelCatalog?: unknown; @@ -273,6 +278,13 @@ function fixture(options?: { if (options?.omitQueueMutation === true) { Object.defineProperty(this, "mutateQueuedMessage", { value: undefined }); } + if (options?.omitMcpSupport === true) { + Object.defineProperties(this, { + supportsAcpMcpServers: { value: undefined }, + replaceAcpMcpServers: { value: undefined }, + releaseAcpMcpServers: { value: undefined }, + }); + } if (options?.authoritativeRlmChildren === undefined) { Object.defineProperty(this, "getRlmChildSnapshots", { value: undefined }); } @@ -478,6 +490,23 @@ function fixture(options?: { } ); } + supportsAcpMcpServers(): boolean { + captures.connectionCalls.push({ method: "supportsAcpMcpServers", args: [] }); + return options?.mcpSupported ?? true; + } + replaceAcpMcpServers(servers: ReadonlyArray, ownerId: string): Promise { + captures.order.push("replace-mcp"); + captures.connectionCalls.push({ method: "replaceAcpMcpServers", args: [servers, ownerId] }); + return options?.replaceMcpImpl?.() ?? Promise.resolve(undefined); + } + releaseAcpMcpServers(ownerId: string, serverNames: ReadonlyArray): Promise { + captures.order.push("release-mcp"); + captures.connectionCalls.push({ + method: "releaseAcpMcpServers", + args: [ownerId, serverNames], + }); + return options?.releaseMcpImpl?.() ?? Promise.resolve(undefined); + } reload(): Promise { captures.connectionCalls.push({ method: "reload", args: [] }); return options?.reloadImpl?.() ?? Promise.resolve(undefined); @@ -568,6 +597,7 @@ function fixture(options?: { ); } dispose(): Promise { + captures.order.push("dispose"); captures.disposeCount += 1; return Promise.resolve(undefined); } @@ -602,6 +632,7 @@ function fixture(options?: { extensions?: ReadonlyArray, requiredExtension?: { readonly path: string; readonly markerCommand: string }, resumeSessionId?: string, + mcpServer?: PrimeAgentDaemonSessionRuntimeInput["mcpServer"], ) => makePrimeAgentDaemonSessionRuntime({ manager, @@ -616,6 +647,7 @@ function fixture(options?: { : { disableExtensionDiscovery: true, disableAutoReconnect: true, requiredExtension }), ...(resumeCursor === undefined ? {} : { resumeCursor }), ...(resumeSessionId === undefined ? {} : { resumeSessionId }), + ...(mcpServer === undefined ? {} : { mcpServer }), }); const emit = (event: unknown) => Promise.resolve(listener?.(event)); const emitWatch = (event: unknown) => Promise.resolve(watcherListener?.(event)); @@ -692,6 +724,203 @@ describe("PrimeAgentDaemonSessionRuntime", () => { }), ); + it.effect("attaches Pylon's scoped MCP server before the initial snapshot and releases it", () => + Effect.gen(function* () { + const { captures, make } = fixture(); + const mcpServer = { + ownerId: "pylon:provider-session-1", + server: { + name: "t3-code", + type: "http" as const, + url: "http://127.0.0.1:4321/mcp/provider-session-1", + headers: { Authorization: "Bearer scoped-secret" }, + }, + }; + + yield* Effect.scoped( + Effect.gen(function* () { + yield* make(undefined, undefined, undefined, undefined, mcpServer); + expect(captures.order).toEqual([ + "request-recovery", + "attach", + "replace-mcp", + "subscribe", + "snapshot", + ]); + expect( + captures.connectionCalls.find((call) => call.method === "replaceAcpMcpServers"), + ).toEqual({ + method: "replaceAcpMcpServers", + args: [[mcpServer.server], mcpServer.ownerId], + }); + }), + ); + + expect( + captures.connectionCalls.find((call) => call.method === "releaseAcpMcpServers"), + ).toEqual({ + method: "releaseAcpMcpServers", + args: [mcpServer.ownerId, [mcpServer.server.name]], + }); + expect(captures.order.indexOf("release-mcp")).toBeLessThan(captures.order.indexOf("dispose")); + expect(captures.disposeCount).toBe(1); + expect(captures.closeCount).toBe(1); + }), + ); + + it.effect("does not replace live MCP ownership for an ordinary resync snapshot", () => + Effect.scoped( + Effect.gen(function* () { + const { captures, emit, make } = fixture(); + const runtime = yield* make(undefined, undefined, undefined, undefined, { + ownerId: "pylon:provider-session-1", + server: { + name: "t3-code", + type: "http", + url: "http://127.0.0.1:4321/mcp/provider-session-1", + headers: { Authorization: "Bearer scoped-secret" }, + }, + }); + const eventsFiber = yield* collectEvents(runtime, 2).pipe(Effect.forkChild); + + yield* Effect.promise(() => emit({ type: "session_resynced", snapshot: snapshot(5) })); + + expect((yield* Fiber.join(eventsFiber)).map((event) => event._tag)).toEqual([ + "SessionResynced", + "SessionResynced", + ]); + expect( + captures.connectionCalls.filter((call) => call.method === "replaceAcpMcpServers"), + ).toHaveLength(1); + }), + ), + ); + + it.effect("reattaches Pylon's scoped MCP server before publishing a daemon resync", () => + Effect.scoped( + Effect.gen(function* () { + const { captures, emit, make } = fixture(); + const runtime = yield* make(undefined, undefined, undefined, undefined, { + ownerId: "pylon:provider-session-1", + server: { + name: "t3-code", + type: "http", + url: "http://127.0.0.1:4321/mcp/provider-session-1", + headers: { Authorization: "Bearer scoped-secret" }, + }, + }); + const eventsFiber = yield* collectEvents(runtime, 4).pipe(Effect.forkChild); + + yield* Effect.promise(() => + emit({ type: "connection_status", status: "reconnecting", error: "private" }), + ); + yield* Effect.promise(() => emit({ type: "session_resynced", snapshot: snapshot(5) })); + yield* Effect.promise(() => emit({ type: "connection_status", status: "connected" })); + + const events = yield* Fiber.join(eventsFiber); + expect(events.map((event) => event._tag)).toEqual([ + "SessionResynced", + "ConnectionStatus", + "SessionResynced", + "ConnectionStatus", + ]); + expect( + captures.connectionCalls.filter((call) => call.method === "replaceAcpMcpServers"), + ).toHaveLength(2); + }), + ), + ); + + it.effect("blocks a waiting prompt when scoped MCP recovery fails", () => + Effect.scoped( + Effect.gen(function* () { + let replacements = 0; + let signalReplacementStarted: () => void = () => undefined; + const replacementStarted = new Promise((resolve) => { + signalReplacementStarted = resolve; + }); + let rejectRecovery: (reason?: unknown) => void = () => undefined; + const { captures, emit, make } = fixture({ + replaceMcpImpl: () => { + replacements += 1; + if (replacements === 1) return Promise.resolve(undefined); + signalReplacementStarted(); + return new Promise((_, reject) => { + rejectRecovery = reject; + }); + }, + }); + const runtime = yield* make(undefined, undefined, undefined, undefined, { + ownerId: "pylon:provider-session-1", + server: { + name: "t3-code", + type: "http", + url: "http://127.0.0.1:4321/mcp/provider-session-1", + headers: { Authorization: "Bearer scoped-secret" }, + }, + }); + const eventsFiber = yield* collectEvents(runtime, 3).pipe(Effect.forkChild); + + yield* Effect.promise(() => + emit({ type: "connection_status", status: "reconnecting", error: "private" }), + ); + const resyncFiber = yield* Effect.promise(() => + emit({ type: "session_resynced", snapshot: snapshot(5) }), + ).pipe(Effect.forkChild); + yield* Effect.promise(() => replacementStarted); + const promptFiber = yield* runtime + .prompt({ text: "must not run without browser tools" }) + .pipe(Effect.result, Effect.forkChild); + yield* Effect.yieldNow; + rejectRecovery(new Error("still streaming")); + yield* Fiber.join(resyncFiber); + yield* Effect.promise(() => emit({ type: "connection_status", status: "connected" })); + + const events = yield* Fiber.join(eventsFiber); + expect(events.map((event) => event._tag)).toEqual([ + "SessionResynced", + "ConnectionStatus", + "SessionClosed", + ]); + const prompt = yield* Fiber.join(promptFiber); + expect(prompt).toMatchObject({ + _tag: "Failure", + failure: { operation: "configure-mcp", reason: "request-failed" }, + }); + expect(captures.connectionCalls.filter((call) => call.method === "prompt")).toEqual([]); + }), + ), + ); + + it.effect("fails closed before snapshot when the daemon cannot own scoped MCP servers", () => + Effect.gen(function* () { + const { captures, make } = fixture({ omitMcpSupport: true }); + const result = yield* Effect.scoped( + make(undefined, undefined, undefined, undefined, { + ownerId: "pylon:provider-session-1", + server: { + name: "t3-code", + type: "http", + url: "http://127.0.0.1:4321/mcp/provider-session-1", + headers: { Authorization: "Bearer scoped-secret" }, + }, + }).pipe(Effect.result), + ); + + expect(result).toMatchObject({ + _tag: "Failure", + failure: { + _tag: "PrimeAgentDaemonSessionRuntimeError", + operation: "configure-mcp", + reason: "incompatible-api", + }, + }); + expect(captures.order).not.toContain("snapshot"); + expect(captures.disposeCount).toBe(1); + expect(captures.closeCount).toBe(1); + }), + ); + it.effect("projects a bounded path-free session resource catalog", () => Effect.gen(function* () { const { make } = fixture({ diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.ts b/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.ts index 830f54021..b8f089acc 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.ts @@ -34,6 +34,7 @@ import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; import { + type PrimeAgentDaemonAcpMcpServer, type PrimeAgentDaemonAgentConnection, type PrimeAgentDaemonExtensionUiResponse, type PrimeAgentDaemonImage, @@ -631,6 +632,7 @@ const runtimeErrorOperation = Schema.Literals([ "configure-client", "create-session", "attach-session", + "configure-mcp", "initial-snapshot", "verify-extension", "reload-resources", @@ -704,6 +706,11 @@ export interface PrimeAgentDaemonSessionRuntimeInput { readonly path: string; readonly markerCommand: string; }; + /** Pylon-owned, thread-scoped MCP server attached only for this live provider session. */ + readonly mcpServer?: { + readonly ownerId: string; + readonly server: PrimeAgentDaemonAcpMcpServer; + }; readonly resumeCursor?: unknown; /** Private stable native id selected from the server-owned identity sidecar. */ readonly resumeSessionId?: string; @@ -1263,10 +1270,72 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo ), }).pipe(Effect.onError(() => completeUnattachedOwnedSession)); - const closeAttachedSession = Effect.promise(async () => { - await connection?.dispose().catch(() => undefined); - client.close(); + let mcpAttached = false; + const releaseMcpServer = Effect.suspend(() => { + const configured = input.mcpServer; + const release = connection?.releaseAcpMcpServers; + if (!mcpAttached || configured === undefined || !Predicate.isFunction(release)) { + return Effect.void; + } + mcpAttached = false; + return Effect.tryPromise({ + try: () => + release + .call(connection, configured.ownerId, [configured.server.name]) + .then(() => undefined), + catch: () => + runtimeError( + "configure-mcp", + "request-failed", + "Could not release Pylon's scoped MCP server from the Prime Agent session.", + ), + }).pipe( + Effect.catch((error) => + Effect.logWarning("Could not release Prime Agent's Pylon MCP session.", { + operation: error.operation, + reason: error.reason, + }), + ), + ); + }); + const closeAttachedSession = releaseMcpServer.pipe( + Effect.andThen( + Effect.promise(async () => { + await connection?.dispose().catch(() => undefined); + client.close(); + }), + ), + ); + const configureMcpServer = Effect.gen(function* () { + const configured = input.mcpServer; + if (configured === undefined) return; + const supports = connection?.supportsAcpMcpServers; + const replace = connection?.replaceAcpMcpServers; + const release = connection?.releaseAcpMcpServers; + if ( + !Predicate.isFunction(supports) || + !Predicate.isFunction(replace) || + !Predicate.isFunction(release) || + supports.call(connection) !== true + ) { + return yield* runtimeError( + "configure-mcp", + "incompatible-api", + "The installed Prime Agent daemon cannot attach Pylon's scoped MCP server. Upgrade Prime Agent or disable agent browser access for this session.", + ); + } + yield* Effect.tryPromise({ + try: () => replace.call(connection, [configured.server], configured.ownerId), + catch: () => + runtimeError( + "configure-mcp", + "request-failed", + "Prime Agent rejected Pylon's scoped MCP server configuration.", + ), + }); + mcpAttached = true; }); + yield* configureMcpServer.pipe(Effect.onError(() => closeAttachedSession)); let verifiedInventory: | readonly [typeof resourceSnapshotSchema.Type, typeof commandsSchema.Type] | undefined; @@ -1487,6 +1556,98 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo handlePrivateSideQuestionEvent(raw).pipe( Effect.flatMap((handled) => (handled ? Effect.void : offerDecoded(raw))), ); + let mcpRecoveryTail = Promise.resolve(); + let mcpRecoveryPending = false; + let mcpRecoveryFailed = false; + const routeMcpAwareRawEvent = (raw: unknown): Promise => { + const rawType = + typeof raw === "object" && raw !== null && "type" in raw && typeof raw.type === "string" + ? raw.type + : undefined; + if ( + input.mcpServer === undefined || + (rawType !== "session_resynced" && rawType !== "connection_status" && rawType !== "closed") + ) { + return runPromise(routeRawEvent(raw)); + } + const delivery = mcpRecoveryTail.then(() => + runPromise( + Effect.gen(function* () { + if (rawType === "connection_status") { + const status = + "status" in (raw as object) && + typeof (raw as { readonly status?: unknown }).status === "string" + ? (raw as { readonly status: string }).status + : undefined; + if (status === "reconnecting") { + mcpAttached = false; + mcpRecoveryPending = true; + } + if (status === "connected" && mcpRecoveryPending) { + mcpRecoveryPending = false; + mcpRecoveryFailed = true; + yield* routeRawEvent({ + type: "closed", + error: "Prime Agent reconnected without restoring Pylon's scoped browser tools.", + }); + return; + } + if (status === "connected" && mcpRecoveryFailed) return; + yield* routeRawEvent(raw); + return; + } + if (rawType === "session_resynced" && mcpRecoveryPending && !mcpRecoveryFailed) { + const restored = yield* configureMcpServer.pipe( + Effect.as(true), + Effect.orElseSucceed(() => false), + ); + mcpRecoveryPending = false; + if (!restored) { + mcpRecoveryFailed = true; + yield* routeRawEvent({ + type: "closed", + error: + "Pylon browser tools could not be restored after the Prime Agent daemon reconnected.", + }); + return; + } + } + if (rawType === "closed") { + mcpRecoveryPending = false; + mcpRecoveryFailed = true; + } + if (!mcpRecoveryFailed || rawType === "closed") yield* routeRawEvent(raw); + }), + ), + ); + mcpRecoveryTail = delivery.catch(() => undefined); + return delivery; + }; + const awaitMcpRecovery = Effect.suspend(() => + input.mcpServer === undefined + ? Effect.void + : Effect.tryPromise({ + try: () => mcpRecoveryTail, + catch: () => + runtimeError( + "configure-mcp", + "request-failed", + "Could not restore Pylon browser tools after the Prime Agent daemon reconnected.", + ), + }).pipe( + Effect.andThen( + Effect.suspend(() => + mcpRecoveryFailed + ? runtimeError( + "configure-mcp", + "request-failed", + "Pylon browser tools are unavailable after the Prime Agent daemon reconnected.", + ) + : Effect.void, + ), + ), + ), + ); // The initial snapshot reserves one queue slot. Admission is cumulative for // the whole initialization phase: draining a batch never reopens capacity. @@ -1501,7 +1662,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo bufferedEvents.push(event); return; } - return runPromise(routeRawEvent(event)); + return routeMcpAwareRawEvent(event); }); const rawSnapshot = yield* Effect.tryPromise({ @@ -2588,6 +2749,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo promptInput: PrimeAgentDaemonPromptInput, ) { yield* ensureOpen("prompt"); + yield* awaitMcpRecovery; const images = yield* validateImages("prompt", promptInput.images); yield* validatePromptContent("prompt", promptInput.text, images); yield* callVoid("prompt", () => @@ -2603,6 +2765,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo promptInput: PrimeAgentDaemonPromptInput, ) { yield* ensureOpen("steer"); + yield* awaitMcpRecovery; const images = yield* validateImages("steer", promptInput.images); yield* validatePromptContent("steer", promptInput.text, images); const method = yield* requireMethod("steer", connection!.steer); @@ -2613,6 +2776,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo promptInput: PrimeAgentDaemonPromptInput, ) { yield* ensureOpen("follow-up"); + yield* awaitMcpRecovery; const images = yield* validateImages("follow-up", promptInput.images); yield* validatePromptContent("follow-up", promptInput.text, images); const method = yield* requireMethod("follow-up", connection!.followUp); @@ -3267,6 +3431,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo bestEffortAbortSideQuestion(nativeId, active), ).pipe( Effect.andThen(failActivePrivateSideQuestions()), + Effect.andThen(releaseMcpServer), Effect.andThen( Effect.tryPromise({ try: () => connection!.dispose(), diff --git a/docs/internals/prime-agent-daemon-parity.md b/docs/internals/prime-agent-daemon-parity.md index 71fb02c3d..e66a77646 100644 --- a/docs/internals/prime-agent-daemon-parity.md +++ b/docs/internals/prime-agent-daemon-parity.md @@ -55,8 +55,8 @@ are never used to synthesize assistant prose. | `abortBranchSummary` | Intentionally folded into Stop | Pylon does not start a standalone native branch-summary operation; stopping the owning turn remains authoritative. | | `setAutoRetryEnabled`, `abortRetry` | Deferred | Retry lifecycle is observed safely, but enabled state has no authoritative readback and the setter writes shared provider settings. Stop already cancels the owning turn. A distinct retry control needs truthful session state and receipts. | | `reload`, `getCommands` | Integrated | Full-access sessions can reload while idle and show bounded safe command metadata. Supervised sessions fail closed. | -| `acquireSessionInputPause` | Intentionally redundant | Every Pylon prompt and resource reload is serialized by the per-thread adapter lock, and reload is admitted only while the owned session is idle. Pylon does not expose a second ACP input path that could bypass that lock, so holding a native lease would add reconnect failure modes without fencing additional work. Revisit with a session-scoped MCP bridge or another independent input owner. | -| `supportsAcpMcpServers`, `replaceAcpMcpServers`, `releaseAcpMcpServers` | Deferred with per-thread MCP ownership | Pylon does not currently send ACP `mcpServers` or expose its per-thread MCP bridge for Prime Agent. Mirroring Prime's owner-scoped replacement API before Pylon has one authoritative MCP configuration and cleanup lifecycle would create split-brain server state. Existing Prime-owned MCP settings continue to load through normal resource reload. | +| `acquireSessionInputPause` | Intentionally redundant | Every Pylon prompt and resource reload is serialized by the per-thread adapter lock, and reload is admitted only while the owned session is idle. Pylon's scoped MCP server is attached before the first snapshot and released only after turn ownership ends, so it is never replaced during live input. Holding a native lease would add reconnect failure modes without fencing additional work. Revisit if Pylon supports live MCP configuration changes. | +| `supportsAcpMcpServers`, `replaceAcpMcpServers`, `releaseAcpMcpServers` | Integrated with scoped Pylon ownership | `McpProviderSession` is the single per-thread source of truth. Before the first daemon snapshot, Pylon replaces one `t3-code` HTTP server under the stable owner `pylon:` and fails closed if Prime cannot own it. After a daemon reconnect, Pylon reclaims that ownership before publishing the resynced session; if it cannot, the session closes instead of continuing without browser tools. Session teardown releases only that owner and server name before disposing the connection. ACP fallback sends the same scoped server in `session/new`. Browser-disabled sessions send nothing, and Prime-owned MCP settings and catalogs remain private. | | `getResourceSnapshot` | Partially integrated by safe outcome | Commands and safe skill/prompt metadata are decoded internally. Native paths, diagnostics, extensions, themes, packages, and MCP configuration are not sent to clients. | | `respondToExtensionUiRequest` | Partially integrated by safe outcome | Select, confirm, and input dialogs plus bounded notifications, status, and widgets are correlated without exposing native request envelopes. Submitted free-form input uses a transient provider RPC and is redacted from durable activities. Editor replacement is cancelled because its prefill may contain sensitive model or tool material that cannot safely enter Pylon's synchronized event stream. | | `getRlmChildSnapshots`, `getRlmMaxDepthStatus`, `setRlmMaxDepth`, `cancelRlmChild`, `sendAgentMessage` | Integrated | On 0.8.0 the authoritative roster atomically replaces Pylon's private cache after strict bounded decoding; older versions retain the event-derived roster. Canonical Pylon task IDs resolve through that private live roster. Messaging is ephemeral. | diff --git a/docs/user/providers-prime-agent.md b/docs/user/providers-prime-agent.md index dd8e724ce..69ee75bda 100644 --- a/docs/user/providers-prime-agent.md +++ b/docs/user/providers-prime-agent.md @@ -208,14 +208,21 @@ Daemon-backed threads resume the exact Prime transcript selected for that Pylon transcript is removed or its private identity cannot be verified, Pylon reports a resume failure instead of silently opening a blank or merely recent Prime session. +## Browser access + +When **Settings → Integrations → Browser → Allow agent browser access** is enabled, new Prime Agent +sessions receive Pylon's thread-scoped preview tools. This works in daemon-backed sessions and ACP +compatibility mode. The scoped connection is removed when the provider session stops. Turning browser +access off withholds both the tools and their instructions; it does not affect browser tabs you control. + ## Current Limitations - Prime Agent 0.8.0 has no daemon-native or operating-system sandbox policy. Supervised mode gates tool admission but does not restrict an approved tool. - Authentication is managed in Prime Agent, not Pylon. -- Plan mode, provider-conversation rollback, general per-item queue editing or reordering, and - Pylon's per-thread MCP bridge are not supported yet. Pylon integrates Prime Agent 0.8.0's mutation - API only for removing a lane's sole item. With multiple count-only items, clients cannot identify a +- Plan mode, provider-conversation rollback, and general per-item queue editing or reordering are + not supported yet. Pylon integrates Prime Agent 0.8.0's mutation API only for removing a lane's + sole item. With multiple count-only items, clients cannot identify a specific target safely without exposing queued text, and ambiguous mutations are never retried. - Pylon does not present live Prime reasoning streams, durable or historical child-session transcripts, cost breakdowns, goal mutations, heartbeats, saved-session history, or native package or MCP catalogs as first-class features. Active children have only the bounded **Live activity** view described above.