diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts b/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts index 64bd8d3072..89356743ff 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.test.ts @@ -47,6 +47,7 @@ import type { ProviderAdapterError } from "../Errors.ts"; import * as ProviderAdapterRegistry from "../Services/ProviderAdapterRegistry.ts"; import * as ProviderService from "../Services/ProviderService.ts"; import * as ProviderEventLoggers from "../Layers/ProviderEventLoggers.ts"; +import { type EventNdjsonLogger } from "../Layers/EventNdjsonLogger.ts"; import { makeProviderServiceLive } from "../Layers/ProviderService.ts"; import { ProviderSessionDirectoryLive } from "../Layers/ProviderSessionDirectory.ts"; import { makeAdapterRegistryMock } from "../testUtils/providerAdapterRegistryMock.ts"; @@ -11584,4 +11585,92 @@ describe("PrimeAgentDaemonAdapter", () => { }), ).pipe(Effect.provide(testLayer)), ); + it.effect( + "records bounded diagnostic classification and suppresses arbitrary private errors from native and canonical logs on SessionClosed", + () => + Effect.scoped( + Effect.gen(function* () { + const loggedNativeEntries: unknown[] = []; + const mockNativeLogger: EventNdjsonLogger = { + filePath: "/mock/native.ndjson", + write: (entry) => + Effect.sync(() => { + loggedNativeEntries.push(entry); + }), + close: () => Effect.void, + }; + const captures = makeCaptures(); + const adapter = yield* makePrimeAgentDaemonAdapter(decodeSettings({}), manager, { + instanceId, + runtimeFactory: fakeRuntimeFactory(captures), + nativeEventLogger: mockNativeLogger, + }); + const subscription = yield* subscribe(adapter); + yield* adapter.startSession({ threadId, cwd: process.cwd(), runtimeMode: "full-access" }); + yield* awaitObservedType(subscription.observed, "thread.started"); + + yield* offer(captures, { + _tag: "SessionClosed", + error: "PRIVATE_SECRET_DO_NOT_LOG /private/server/credentials.key", + diagnostic: { + reason: "ingress-capacity", + connectionGeneration: 2, + proofEpoch: 1, + }, + }); + const exited = yield* awaitObservedType(subscription.observed, "session.exited"); + + expect(exited).toMatchObject({ + payload: { + exitKind: "error", + reason: "Prime Agent session closed unexpectedly.", + }, + }); + + // Verify diagnostic classification is recorded in native logger + const closedEntry = loggedNativeEntries.find( + (entry) => + typeof entry === "object" && + entry !== null && + "event" in entry && + typeof (entry as { event: unknown }).event === "object" && + (entry as { event: unknown }).event !== null && + "method" in (entry as { event: { method: unknown } }).event && + (entry as { event: { method: unknown } }).event.method === "SessionClosed", + ) as + | { + readonly event: { + readonly kind: string; + readonly provider: string; + readonly threadId: ThreadId; + readonly method: string; + readonly diagnostic?: unknown; + }; + } + | undefined; + expect(closedEntry).toBeDefined(); + expect(closedEntry?.event).toMatchObject({ + kind: "notification", + provider: "primeAgent", + threadId, + method: "SessionClosed", + diagnostic: { + reason: "ingress-capacity", + connectionGeneration: 2, + proofEpoch: 1, + }, + }); + + // Verify arbitrary private error strings remain strictly absent from both native and canonical logs + const serializedNative = encodeUnknownJson(loggedNativeEntries); + const serializedCanonical = encodeUnknownJson(subscription.events); + expect(serializedNative).not.toContain("PRIVATE_SECRET_DO_NOT_LOG"); + expect(serializedNative).not.toContain("/private/server/credentials.key"); + expect(serializedCanonical).not.toContain("PRIVATE_SECRET_DO_NOT_LOG"); + expect(serializedCanonical).not.toContain("/private/server/credentials.key"); + + yield* Fiber.interrupt(subscription.fiber); + }), + ).pipe(Effect.provide(testLayer)), + ); }); diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts b/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts index fded6c197c..6627b1305c 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonAdapter.ts @@ -1431,6 +1431,9 @@ export function makePrimeAgentDaemonAdapter( provider: PROVIDER, threadId, method: event._tag, + ...(event._tag === "SessionClosed" && event.diagnostic !== undefined + ? { diagnostic: event.diagnostic } + : {}), }, }, threadId, diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonEvents.test.ts b/apps/server/src/provider/prime/PrimeAgentDaemonEvents.test.ts index 9e3c3b70f6..1a7e865f87 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonEvents.test.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonEvents.test.ts @@ -1241,6 +1241,18 @@ describe("PrimeAgentDaemonEvents", () => { expect(decodePrimeAgentDaemonEvent({ type: "closed", error: "daemon exited" })).toEqual({ _tag: "SessionClosed", error: "daemon exited", + diagnostic: { reason: "provider-closed" }, + }); + expect( + decodePrimeAgentDaemonEvent({ + type: "closed", + error: "daemon exited", + diagnostic: { reason: "mcp-restore", connectionGeneration: 3, proofEpoch: 4 }, + }), + ).toEqual({ + _tag: "SessionClosed", + error: "daemon exited", + diagnostic: { reason: "provider-closed" }, }); }); diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonEvents.ts b/apps/server/src/provider/prime/PrimeAgentDaemonEvents.ts index 261153791f..5b155bc85b 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonEvents.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonEvents.ts @@ -28,6 +28,24 @@ const serviceTier = Schema.NullOr( const queueMode = Schema.Literals(["all", "one-at-a-time"]); const stopReason = Schema.Literals(["stop", "length", "toolUse", "error", "aborted"]); +export const PrimeSessionClosedDiagnosticReason = Schema.Literals([ + "proof-lost", + "ingress-capacity", + "snapshot-reconciliation", + "mcp-restore", + "provider-closed", +]); +export type PrimeSessionClosedDiagnosticReason = typeof PrimeSessionClosedDiagnosticReason.Type; + +const nonNegativeInt = Schema.Int.check(Schema.isGreaterThanOrEqualTo(0)); + +export const PrimeSessionClosedDiagnostic = Schema.Struct({ + reason: PrimeSessionClosedDiagnosticReason, + connectionGeneration: Schema.optional(nonNegativeInt), + proofEpoch: Schema.optional(nonNegativeInt), +}); +export type PrimeSessionClosedDiagnostic = typeof PrimeSessionClosedDiagnostic.Type; + const textContent = Schema.Struct({ type: Schema.Literal("text"), text: Schema.String, @@ -650,7 +668,10 @@ export const PrimeAgentDaemonConnectionEvent = Schema.Union([ error: Schema.optional(Schema.String), }), Schema.Struct({ type: Schema.Literal("heartbeats_changed") }), - Schema.Struct({ type: Schema.Literal("closed"), error: Schema.optional(Schema.String) }), + Schema.Struct({ + type: Schema.Literal("closed"), + error: Schema.optional(Schema.String), + }), ]); export type PrimeAgentDaemonConnectionEvent = typeof PrimeAgentDaemonConnectionEvent.Type; @@ -1387,7 +1408,11 @@ export type PrimeDaemonEvent = ( readonly error?: string | undefined; } | { readonly _tag: "HeartbeatsChanged" } - | { readonly _tag: "SessionClosed"; readonly error?: string | undefined } + | { + readonly _tag: "SessionClosed"; + readonly error?: string | undefined; + readonly diagnostic?: PrimeSessionClosedDiagnostic | undefined; + } | { readonly _tag: "Ignored"; readonly reason: "unknown-event" | "malformed-event"; @@ -2045,7 +2070,11 @@ function mapPrimeAgentDaemonConnectionEvent( case "heartbeats_changed": return { _tag: "HeartbeatsChanged" }; case "closed": - return { _tag: "SessionClosed", error: optionalBounded(event.error) }; + return { + _tag: "SessionClosed", + error: optionalBounded(event.error), + diagnostic: { reason: "provider-closed" }, + }; } } diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.test.ts b/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.test.ts index c5131ea617..05751b3105 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.test.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.test.ts @@ -1419,6 +1419,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }, ]); expect( @@ -1478,6 +1479,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }, ]); }), @@ -1680,6 +1682,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent event ingress exceeded its bounded capacity.", + diagnostic: expect.objectContaining({ reason: "ingress-capacity" }), }, ]); return; @@ -1888,6 +1891,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }, ]); expect(test.captures.order.filter((step) => step === "snapshot")).toHaveLength(1); @@ -1935,6 +1939,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect((yield* collectEvents(runtime, 1))[0]).toEqual({ _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "snapshot-reconciliation", proofEpoch: 2 }), }); expect(snapshotReads).toBe(2); }), @@ -2006,6 +2011,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect(recoveredEvents[3]).toEqual({ _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }); expect( test.captures.connectionCalls.filter((call) => call.method === "submitCorrelatedPrompt"), @@ -2511,6 +2517,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect(published.at(-1)).toEqual({ _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }); expect(published.filter((event) => event._tag === "SessionClosed")).toHaveLength(1); const afterTerminal = yield* runtime @@ -2620,6 +2627,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect(published.at(-1)).toEqual({ _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }); }), ); @@ -2712,6 +2720,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "ingress-capacity" }), }, ]); }), @@ -2788,6 +2797,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect((yield* collectEvents(runtime, 1))[0]).toEqual({ _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "ingress-capacity" }), }); yield* Fiber.join(callbacks); expect(verificationCalls).toBe(2); @@ -2874,6 +2884,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect((yield* collectEvents(runtime, 1))[0]).toEqual({ _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }); expect(verificationCalls).toBe(2); expect(resourceCalls).toBe(2); @@ -3060,6 +3071,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect((yield* collectEvents(runtime, 1))[0]).toEqual({ _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }); expect(replacementCalls).toBe(2); }), @@ -3104,6 +3116,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect(published.at(-1)).toEqual({ _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }); expect(published.filter((event) => event._tag === "SessionClosed")).toHaveLength(1); }), @@ -3528,6 +3541,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }, ]); const afterTerminal = yield* runtime @@ -3592,6 +3606,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }, ]); }), @@ -3816,6 +3831,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }, ]); }), @@ -4169,6 +4185,12 @@ describe("PrimeAgentDaemonSessionRuntime", () => { "ConnectionStatus", "SessionClosed", ]); + expect(events[2]).toMatchObject({ + _tag: "SessionClosed", + error: + "Pylon browser tools could not be restored after the Prime Agent daemon reconnected.", + diagnostic: expect.objectContaining({ reason: "mcp-restore" }), + }); const prompt = yield* Fiber.join(promptFiber); expect(prompt).toMatchObject({ _tag: "Failure", @@ -4179,6 +4201,41 @@ describe("PrimeAgentDaemonSessionRuntime", () => { ), ); + it.effect( + "classifies internally synthesized MCP reconnect failure without snapshot as mcp-restore", + () => + Effect.scoped( + Effect.gen(function* () { + const { emit, make } = fixture(); + const runtime = yield* make(undefined, undefined, undefined, undefined, { + ownerId: "pylon:mcp-restore-test", + server: { + name: "t3-code", + type: "http", + url: "http://127.0.0.1:4321/mcp/mcp-restore-test", + headers: { Authorization: "Bearer scoped-secret" }, + }, + }); + const eventsFiber = yield* collectEvents(runtime, 3).pipe(Effect.forkChild); + + yield* Effect.promise(() => emit({ type: "connection_status", status: "reconnecting" })); + 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", + ]); + expect(events[2]).toMatchObject({ + _tag: "SessionClosed", + error: "Prime Agent reconnected without restoring Pylon's scoped browser tools.", + diagnostic: expect.objectContaining({ reason: "mcp-restore" }), + }); + }), + ), + ); + it.effect("fails closed before snapshot when the daemon cannot own scoped MCP servers", () => Effect.gen(function* () { const { captures, make } = fixture({ omitMcpSupport: true }); @@ -5560,6 +5617,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { _tag: "SessionClosed", error: "Pylon's managed provider extension could not be verified after Prime Agent reconnected.", + diagnostic: expect.objectContaining({ reason: "snapshot-reconciliation" }), }); expect( test.captures.connectionCalls.filter((call) => call.method === "getToolDefinition"), @@ -5623,6 +5681,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }, ]); const extra = yield* runtime.events.pipe( @@ -5687,6 +5746,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }, ]); const extra = yield* runtime.events.pipe( @@ -6212,6 +6272,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect(published.at(-1)).toEqual({ _tag: "SessionClosed", error: "Prime Agent event ingress exceeded its bounded capacity.", + diagnostic: expect.objectContaining({ reason: "ingress-capacity" }), }); expect(published.filter((event) => event._tag === "SessionClosed")).toHaveLength(1); const promptError = yield* runtime.prompt({ text: "must remain failed" }).pipe(Effect.flip); @@ -6259,6 +6320,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect(published.at(-1)).toEqual({ _tag: "SessionClosed", error: "Prime Agent event ingress exceeded its bounded capacity.", + diagnostic: expect.objectContaining({ reason: "ingress-capacity" }), }); const answer = yield* Fiber.join(asking); expect(answer).toMatchObject({ @@ -6607,6 +6669,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect(published.at(-1)).toEqual({ _tag: "SessionClosed", error: "Prime Agent event ingress exceeded its bounded capacity.", + diagnostic: expect.objectContaining({ reason: "ingress-capacity" }), }); expect(verificationCalls).toBe(1); const prompt = yield* Fiber.join(promptFiber); @@ -6677,6 +6740,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect((yield* collectEvents(runtime, 1))[0]).toEqual({ _tag: "SessionClosed", error: "Prime Agent event ingress exceeded its bounded capacity.", + diagnostic: expect.objectContaining({ reason: "ingress-capacity" }), }); yield* Fiber.join(callbacks); expect(yield* Fiber.join(prompt)).toMatchObject({ @@ -8918,6 +8982,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { _tag: "SessionClosed", error: "Pylon's managed provider extension could not be verified after Prime Agent reconnected.", + diagnostic: expect.objectContaining({ reason: "snapshot-reconciliation" }), }); }), ), @@ -9053,6 +9118,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }, ]); expect(snapshotReads).toBe(2); @@ -9099,6 +9165,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }, ]); yield* Fiber.join(closing); @@ -9144,6 +9211,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect((yield* collectEvents(runtime, 1))[0]).toEqual({ _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }); yield* Effect.promise(() => test.emit({ @@ -9231,6 +9299,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }, ]); yield* Fiber.join(closing); @@ -9395,7 +9464,11 @@ describe("PrimeAgentDaemonSessionRuntime", () => { ]).then(() => undefined), ); expect(yield* collectEvents(runtime, 1)).toEqual([ - { _tag: "SessionClosed", error: undefined }, + { + _tag: "SessionClosed", + error: undefined, + diagnostic: expect.objectContaining({ reason: "provider-closed" }), + }, ]); const callsBefore = [...test.captures.connectionCalls]; expect(yield* runtime.getInputQueueStatus.pipe(Effect.flip)).toMatchObject({ @@ -9610,6 +9683,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect((yield* collectEvents(runtime, 1))[0]).toEqual({ _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }); expect(test.captures.commands.filter((command) => command.type === "list")).toHaveLength(0); const callsBefore = [...test.captures.connectionCalls]; @@ -9717,6 +9791,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { expect((yield* collectEvents(runtime, 1))[0]).toEqual({ _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }); test.setCorrelatedPromptLifecycleProof(true); yield* Effect.promise(() => @@ -11291,6 +11366,7 @@ describe("PrimeAgentDaemonSessionRuntime", () => { { _tag: "SessionClosed", error: "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnostic: expect.objectContaining({ reason: "proof-lost" }), }, ]); }), diff --git a/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.ts b/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.ts index c0d10b7a12..5710300c71 100644 --- a/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.ts +++ b/apps/server/src/provider/prime/PrimeAgentDaemonSessionRuntime.ts @@ -64,6 +64,7 @@ import { decodePrimeAgentPromptLifecycleSubmitResult, primeAgentDaemonImageDigest, PRIME_AGENT_DAEMON_MESSAGE_TEXT_MAX_CHARS, + type PrimeSessionClosedDiagnosticReason, PRIME_AGENT_DAEMON_TRANSCRIPT_MAX_MESSAGES, primeAgentPromptLifecycleCanAdvance, primeAgentPromptLifecycleIsSame, @@ -1884,7 +1885,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo if (workerRecoveryCorrelatedProofIsCurrent(recovery)) return false; recovery.provisionalSnapshot = undefined; settleReconnectResolution(recovery.resolution.generation, false); - void failCorrelatedProofRecovery().catch(() => undefined); + void failCorrelatedProofRecovery(undefined, "proof-lost").catch(() => undefined); return true; }; @@ -1905,7 +1906,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo recovery.provisionalSnapshot = undefined; } settleReconnectResolution(generation, false); - void failCorrelatedProofRecovery().catch(() => undefined); + void failCorrelatedProofRecovery(undefined, "proof-lost").catch(() => undefined); return false; } if ( @@ -1931,7 +1932,9 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo } const settled = settleReconnectResolution(generation, reconciled); if (settled && correlatedPromptLifecycleAvailable && !reconciled) { - void failCorrelatedProofRecovery().catch(() => undefined); + void failCorrelatedProofRecovery(undefined, "snapshot-reconciliation").catch( + () => undefined, + ); } if ( settled && @@ -2772,6 +2775,10 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo const terminalEvent = { _tag: "SessionClosed", error: "Prime Agent event ingress exceeded its bounded capacity.", + diagnostic: { + reason: "ingress-capacity", + connectionGeneration, + }, } satisfies PrimeDaemonEvent; return failActivePrivateSideQuestions().pipe( Effect.andThen( @@ -3488,6 +3495,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedProofIngressEpoch?: number, ordinaryIngressFence?: OrdinaryIngressFence, providerRouteRetirement?: ProviderRouteRetirement, + synthesizedDiagnosticReason?: PrimeSessionClosedDiagnosticReason, ) => { if ( !ordinaryIngressFenceIsCurrent(ordinaryIngressFence) || @@ -3513,6 +3521,15 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedPromptLifecycle: correlatedPromptLifecycleAvailable, }), ); + if (decoded._tag === "SessionClosed" && synthesizedDiagnosticReason !== undefined) { + decoded = { + ...decoded, + diagnostic: { + ...decoded.diagnostic, + reason: synthesizedDiagnosticReason, + }, + }; + } if ( correlatedPromptLifecycleAvailable && decoded._tag === "SessionResynced" && @@ -3891,6 +3908,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedProofIngressEpoch?: number, ordinaryIngressFence?: OrdinaryIngressFence, providerRouteRetirement?: ProviderRouteRetirement, + synthesizedDiagnosticReason?: PrimeSessionClosedDiagnosticReason, ) => Effect.suspend(() => { if ( @@ -3911,6 +3929,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedProofIngressEpoch, ordinaryIngressFence, providerRouteRetirement, + synthesizedDiagnosticReason, ); } const workerCloseRaw = Predicate.isObject(raw) && "type" in raw && raw.type === "closed"; @@ -3930,6 +3949,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedProofIngressEpoch, ordinaryIngressFence, providerRouteRetirement, + synthesizedDiagnosticReason, ), ), ); @@ -3939,12 +3959,19 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo proofEpoch: number | undefined, ordinaryIngressFence?: OrdinaryIngressFence, providerRouteRetirement?: ProviderRouteRetirement, + synthesizedDiagnosticReason?: PrimeSessionClosedDiagnosticReason, ) => !ordinaryIngressFenceIsCurrent(ordinaryIngressFence) || !providerRouteRetirementIsCurrent(providerRouteRetirement) ? Effect.void : proofEpoch === undefined || correlatedPromptLifecycleProofFenceIsCurrent(proofEpoch) - ? routeRawEvent(raw, proofEpoch, ordinaryIngressFence, providerRouteRetirement) + ? routeRawEvent( + raw, + proofEpoch, + ordinaryIngressFence, + providerRouteRetirement, + synthesizedDiagnosticReason, + ) : Effect.fail(CORRELATED_PROOF_FENCE_RETIRED); let mcpRecoveryTail = Promise.resolve(); let mcpRecoveryPending = false; @@ -4027,6 +4054,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedProofIngressEpoch?: number, ordinaryIngressFence?: OrdinaryIngressFence, providerRouteRetirement?: ProviderRouteRetirement, + synthesizedDiagnosticReason?: PrimeSessionClosedDiagnosticReason, ): Promise => { if ( !ordinaryIngressFenceIsCurrent(ordinaryIngressFence) || @@ -4075,6 +4103,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedProofIngressEpoch, ordinaryIngressFence, providerRouteRetirement, + synthesizedDiagnosticReason, ), ); } @@ -4109,6 +4138,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo undefined, ordinaryIngressFence, providerRouteRetirement, + "mcp-restore", ); return; } @@ -4178,6 +4208,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo undefined, ordinaryIngressFence, providerRouteRetirement, + "mcp-restore", ); return; } @@ -4194,6 +4225,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedProofIngressEpoch, ordinaryIngressFence, providerRouteRetirement, + synthesizedDiagnosticReason, ); } }); @@ -4382,6 +4414,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedProofIngressEpoch?: number, ordinaryIngressFence?: OrdinaryIngressFence, providerRouteRetirement?: ProviderRouteRetirement, + synthesizedDiagnosticReason?: PrimeSessionClosedDiagnosticReason, ): Promise => { if ( !ordinaryIngressFenceIsCurrent(ordinaryIngressFence) || @@ -4438,6 +4471,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedProofIngressEpoch, ordinaryIngressFence, providerRouteRetirement, + synthesizedDiagnosticReason, ); } const routeEffect = Effect.gen(function* () { @@ -4474,6 +4508,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo undefined, ordinaryIngressFence, providerRouteRetirement, + "snapshot-reconciliation", ); return; } @@ -4507,6 +4542,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo undefined, ordinaryIngressFence, providerRouteRetirement, + "snapshot-reconciliation", ); return; } @@ -4522,6 +4558,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedProofIngressEpoch, ordinaryIngressFence, providerRouteRetirement, + synthesizedDiagnosticReason, ), ); if ( @@ -4683,7 +4720,11 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo }); const failCorrelatedProofRecovery = ( error = "Prime Agent correlated prompt capability proof was lost during recovery.", + diagnosticReason?: PrimeSessionClosedDiagnosticReason, ): Promise => { + const capturedProofEpoch = + activeWorkerRecovery?.correlatedProofEpoch ?? + (correlatedProofEpoch >= 0 ? correlatedProofEpoch : undefined); correlatedProofRecoveryPending = false; if (correlatedProofRecoveryFailed) return Promise.resolve(); correlatedProofRecoveryFailed = true; @@ -4701,25 +4742,38 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo if (activeWorkerRecovery === workerRecovery) activeWorkerRecovery = undefined; } settleReconnectResolution(connectionGeneration, false); + const effectiveReason: PrimeSessionClosedDiagnosticReason = diagnosticReason ?? "proof-lost"; const terminal = { _tag: "SessionClosed", error, + diagnostic: { + reason: effectiveReason, + ...(connectionGeneration >= 0 ? { connectionGeneration } : {}), + ...(capturedProofEpoch !== undefined && capturedProofEpoch >= 0 + ? { proofEpoch: capturedProofEpoch } + : {}), + }, } satisfies PrimeDaemonEvent; return runPromise( failActivePrivateSideQuestions().pipe(Effect.andThen(offerRuntimeEvent(terminal))), ); }; - const routeProvedCorrelatedRawEvent = (raw: unknown, proofEpoch: number): Promise => + const routeProvedCorrelatedRawEvent = ( + raw: unknown, + proofEpoch: number, + synthesizedDiagnosticReason?: PrimeSessionClosedDiagnosticReason, + ): Promise => routeManagedAwareRawEvent( raw, proofEpoch, undefined, correlatedProviderRouteRetirement, + synthesizedDiagnosticReason, ).catch((cause) => { if (cause !== CORRELATED_PROOF_FENCE_RETIRED) return Promise.reject(cause); if (initializing && initializationOverflow) return Promise.resolve(); return proofEpoch === correlatedProofEpoch - ? failCorrelatedProofRecovery() + ? failCorrelatedProofRecovery(undefined, "proof-lost") : Promise.resolve(); }); let ordinaryRawRouteTail = Promise.resolve(); @@ -4760,7 +4814,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedProofRouteCount >= MAX_CORRELATED_PROOF_ROUTES || correlatedProofRouteWeight + weight > MAX_CORRELATED_PROOF_ROUTE_WEIGHT ) { - return failCorrelatedProofRecovery(); + return failCorrelatedProofRecovery(undefined, "ingress-capacity"); } correlatedProofRouteCount += 1; correlatedProofRouteWeight += weight; @@ -4812,13 +4866,14 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo mcpAttached = false; mcpRecoveryPending = true; } - if (pendingStrictWorkerRecovery) return failCorrelatedProofRecovery(); + if (pendingStrictWorkerRecovery) + return failCorrelatedProofRecovery(undefined, "proof-lost"); const providerRouteRetirement = correlatedProviderRouteRetirement; return serializeCorrelatedProofRoute(raw, () => routeManagedAwareRawEvent(raw, undefined, undefined, providerRouteRetirement), ); } - if (rawType === "closed") return failCorrelatedProofRecovery(); + if (rawType === "closed") return failCorrelatedProofRecovery(undefined, "provider-closed"); if (rawType === "session_replaced") { const pendingStrictWorkerRecovery = activeWorkerRecovery?.correlatedProofEpoch !== undefined; @@ -4834,9 +4889,10 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo mcpAttached = false; mcpRecoveryPending = true; } - if (pendingStrictWorkerRecovery) return failCorrelatedProofRecovery(); + if (pendingStrictWorkerRecovery) + return failCorrelatedProofRecovery(undefined, "proof-lost"); const proofEpoch = captureCorrelatedPromptLifecycleProofFence(); - if (proofEpoch === undefined) return failCorrelatedProofRecovery(); + if (proofEpoch === undefined) return failCorrelatedProofRecovery(undefined, "proof-lost"); const providerRouteRetirement = correlatedProviderRouteRetirement; return serializeCorrelatedProofRoute(raw, () => { const replacementRoute = Effect.gen(function* () { @@ -4849,7 +4905,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo } if (!correlatedPromptLifecycleProofFenceIsCurrent(proofEpoch)) { if (proofEpoch === correlatedProofEpoch) { - yield* Effect.promise(() => failCorrelatedProofRecovery()); + yield* Effect.promise(() => failCorrelatedProofRecovery(undefined, "proof-lost")); } return; } @@ -4865,7 +4921,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo }); if (!correlatedPromptLifecycleProofFenceIsCurrent(proofEpoch)) { if (proofEpoch === correlatedProofEpoch) { - yield* Effect.promise(() => failCorrelatedProofRecovery()); + yield* Effect.promise(() => failCorrelatedProofRecovery(undefined, "proof-lost")); } return; } @@ -4884,7 +4940,9 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo if (proofEpoch !== correlatedProofEpoch) return; if (!snapshotResolution.snapshotPublished) { settleReconnectResolution(snapshotResolution.generation, false); - yield* Effect.promise(() => failCorrelatedProofRecovery()); + yield* Effect.promise(() => + failCorrelatedProofRecovery(undefined, "snapshot-reconciliation"), + ); return; } correlatedProofRecoveryPending = false; @@ -4892,7 +4950,9 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo Effect.catch(() => proofEpoch !== correlatedProofEpoch ? Effect.void - : Effect.promise(() => failCorrelatedProofRecovery()), + : Effect.promise(() => + failCorrelatedProofRecovery(undefined, "snapshot-reconciliation"), + ), ), ); return runPromise( @@ -4924,21 +4984,21 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo ? reconnectResolution : undefined; const proofEpoch = captureCorrelatedPromptLifecycleProofFence(); - if (proofEpoch === undefined) return failCorrelatedProofRecovery(); + if (proofEpoch === undefined) return failCorrelatedProofRecovery(undefined, "proof-lost"); return serializeCorrelatedProofRoute(raw, async () => { if (!correlatedPromptLifecycleProofFenceIsCurrent(proofEpoch)) { await routeProvedCorrelatedRawEvent(raw, proofEpoch); return; } if (connectionStatus === "connected" && correlatedProofRecoveryPending) { - await failCorrelatedProofRecovery(); + await failCorrelatedProofRecovery(undefined, "proof-lost"); return; } await routeProvedCorrelatedRawEvent(raw, proofEpoch); if (rawType !== "session_resynced" || proofEpoch !== correlatedProofEpoch) return; if (snapshotResolution !== undefined && !snapshotResolution.snapshotPublished) { settleReconnectResolution(snapshotResolution.generation, false); - await failCorrelatedProofRecovery(); + await failCorrelatedProofRecovery(undefined, "snapshot-reconciliation"); return; } correlatedProofRecoveryPending = false; @@ -6762,7 +6822,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo ) { return false; } - await failCorrelatedProofRecovery(); + await failCorrelatedProofRecovery(undefined, "proof-lost"); return true; }; const routeWorkerRecoveryTerminal = async ( @@ -6770,26 +6830,34 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo ) => { if (recovery.terminalFallbackRouted) return; recovery.terminalFallbackRouted = true; + const isSynthesizedFallback = recovery.fallbackCloseRaw === undefined; const fallback = recovery.fallbackCloseRaw ?? ({ type: "closed", error: "Prime Agent replacement-worker recovery ended before Pylon could verify it.", } as const); + const synthesizedReason: PrimeSessionClosedDiagnosticReason | undefined = + isSynthesizedFallback ? "snapshot-reconciliation" : undefined; if (recovery.correlatedProofEpoch === undefined) { await routeManagedAwareRawEvent( fallback, undefined, recovery.ordinaryIngressFence, recovery.ordinaryIngressFence, + synthesizedReason, ); return; } if (!correlatedPromptLifecycleProofFenceIsCurrent(recovery.correlatedProofEpoch)) { - await failCorrelatedProofRecovery(); + await failCorrelatedProofRecovery(undefined, "proof-lost"); return; } - await routeProvedCorrelatedRawEvent(fallback, recovery.correlatedProofEpoch); + await routeProvedCorrelatedRawEvent( + fallback, + recovery.correlatedProofEpoch, + synthesizedReason, + ); }; const finishSuccessfulWorkerCloseRecovery = ( recovery: NonNullable, @@ -6837,7 +6905,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedWorkerProofEpoch === undefined || !correlatedPromptLifecycleProofFenceIsCurrent(correlatedWorkerProofEpoch)) ) { - return failCorrelatedProofRecovery(); + return failCorrelatedProofRecovery(undefined, "proof-lost"); } const routeWorkerRaw = (workerRaw: unknown): Promise => correlatedWorkerProofEpoch === undefined @@ -7018,7 +7086,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo if (correlatedPromptLifecycleAvailable) { quiescenceCorrelatedProofEpoch = captureCorrelatedPromptLifecycleProofFence(); if (quiescenceCorrelatedProofEpoch === undefined) { - yield* Effect.promise(() => failCorrelatedProofRecovery()); + yield* Effect.promise(() => failCorrelatedProofRecovery(undefined, "proof-lost")); return yield* runtimeError( "rlm-quiescence", "request-failed", @@ -7047,9 +7115,9 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo Effect.suspend(() => { if (workerRecoveryCorrelatedProofIsCurrent(recovery)) return Effect.void; rejectRetiredWorkerRecoveryProof(recovery); - return Effect.promise(() => failCorrelatedProofRecovery()).pipe( - Effect.flatMap(() => workerRecoveryFailure()), - ); + return Effect.promise(() => + failCorrelatedProofRecovery(undefined, "proof-lost"), + ).pipe(Effect.flatMap(() => workerRecoveryFailure())); }); const adoptConcurrentWorkerRecovery = Effect.gen(function* () { const closeAttempt = workerCloseRecoveryAttempt; @@ -7165,7 +7233,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo correlatedPromptLifecycleAvailable && correlatedRecoveryProofEpoch === undefined ) { - yield* Effect.promise(() => failCorrelatedProofRecovery()); + yield* Effect.promise(() => failCorrelatedProofRecovery(undefined, "proof-lost")); return yield* workerRecoveryFailure(); } recoveryState ??= { @@ -7275,7 +7343,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo (publicationProofEpoch === undefined || !correlatedPromptLifecycleProofFenceIsCurrent(publicationProofEpoch)) ) { - yield* Effect.promise(() => failCorrelatedProofRecovery()); + yield* Effect.promise(() => failCorrelatedProofRecovery(undefined, "proof-lost")); return yield* workerRecoveryFailure(); } const usage = subtractCumulativeUsage(currentUsage, rlmTurnUsageBaseline); @@ -7290,7 +7358,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo ).pipe( Effect.catch((cause) => cause === CORRELATED_PROOF_FENCE_RETIRED - ? Effect.promise(() => failCorrelatedProofRecovery()).pipe( + ? Effect.promise(() => failCorrelatedProofRecovery(undefined, "proof-lost")).pipe( Effect.flatMap(() => workerRecoveryFailure()), ) : Effect.fail(cause), @@ -7445,7 +7513,7 @@ export const makePrimeAgentDaemonSessionRuntime = Effect.fn("makePrimeAgentDaemo const CORRELATED_PROOF_UNAVAILABLE_ERROR = "Prime Agent correlated prompt capability proof is unavailable for the current attachment."; const correlatedProofUnavailable = (operation: "prompt" | "abort") => - Effect.promise(() => failCorrelatedProofRecovery()).pipe( + Effect.promise(() => failCorrelatedProofRecovery(undefined, "proof-lost")).pipe( Effect.flatMap(() => runtimeError(operation, "request-failed", CORRELATED_PROOF_UNAVAILABLE_ERROR), ),