diff --git a/docs/api-reference/veryfront/agent.md b/docs/api-reference/veryfront/agent.md index a3ffe71d81..807f6dcd82 100644 --- a/docs/api-reference/veryfront/agent.md +++ b/docs/api-reference/veryfront/agent.md @@ -586,7 +586,7 @@ Input delivered to a hosted agent-service detached execution callback. | `buildRootOwnedChildRunResultHint` | Builds root owned child run result hint. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/child-run/result-summary.ts#L348) | | `buildRootOwnedChildRunResultText` | Builds root owned child run result text. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/child-run/result-summary.ts#L333) | | `buildRootOwnedDelegatedFindingsInstruction` | Builds root owned delegated findings instruction. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/conversation/delegation-policy.ts#L34) | -| `buildRuntimeAgentControlPlaneStreamRequestFromInvocation` | Builds runtime agent control plane stream request from invocation. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/agent-invocation-contract.ts#L452) | +| `buildRuntimeAgentControlPlaneStreamRequestFromInvocation` | Builds runtime agent control plane stream request from invocation. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/agent-invocation-contract.ts#L453) | | `buildRuntimeAvailableSkillsPromptBlock` | Builds a bounded, injection-safe runtime available-skills prompt. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/skill-prompt.ts#L457) | | `buildRuntimeLoadedSkillResponse` | Build a bounded loaded-skill response and fail closed on invalid metadata. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/skill-metadata.ts#L1097) | | `buildRuntimeSkillDefinition` | Build a bounded, immutable runtime skill definition. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/skill-metadata.ts#L913) | @@ -884,10 +884,10 @@ Input delivered to a hosted agent-service detached execution callback. | `parseHostedAgentServiceConfig` | Configuration used by parse hosted agent service. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/service/config.ts#L167) | | `parseHostedChatRequestFromRequest` | Request payload for parse hosted chat request from. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/hosted/chat-request-parser.ts#L593) | | `parseRuntimeAgentMarkdownDefinition` | Definition for parse runtime agent markdown. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/agent-definition.ts#L192) | -| `parseRuntimeAgentRunInvocation` | Parses runtime agent run invocation. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/agent-invocation-contract.ts#L479) | +| `parseRuntimeAgentRunInvocation` | Parses runtime agent run invocation. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/agent-invocation-contract.ts#L481) | | `parseRuntimeAgentRunInvocationAgentServiceChatRequestFromRequest` | Request payload for parse runtime agent run invocation hosted chat request from. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/hosted/chat-request-parser.ts#L632) | | `parseRuntimeAgentRunInvocationHostedChatRequestFromRequest` | Request payload for parse runtime agent run invocation hosted chat request from. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/hosted/chat-request-parser.ts#L632) | -| `parseRuntimeAgentRunInvocationOrError` | Error shape for parse runtime agent run invocation or. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/agent-invocation-contract.ts#L486) | +| `parseRuntimeAgentRunInvocationOrError` | Error shape for parse runtime agent run invocation or. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/agent-invocation-contract.ts#L488) | | `parseRuntimeSkillDocument` | Parses a bounded runtime skill document and fails closed on invalid input. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/skill-metadata.ts#L561) | | `parseRuntimeSkillMetadata` | Parses bounded runtime skill metadata and fails closed on invalid input. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/skill-metadata.ts#L569) | | `parseToolInputObject` | Parses tool input object. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/streaming/tool-input.ts#L135) | @@ -1020,7 +1020,7 @@ Input delivered to a hosted agent-service detached execution callback. | Name | Description | Source | | ------------------------------------ | ------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------ | -| `AgentRuntime` | Implement agent runtime. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/index.ts#L1164) | +| `AgentRuntime` | Implement agent runtime. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/index.ts#L1197) | | `AgentRuntimeMessageConversionError` | Error shape for agent runtime message conversion. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/runtime/message-adapter.ts#L138) | | `AgentServiceAuthError` | Error shape for hosted service auth. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/service/auth.ts#L14) | | `AppendConversationRunEventsError` | Error shape for append conversation run events. | [source](https://github.com/veryfront/veryfront-code/blob/main/src/agent/conversation/durable-append-errors.ts#L4) | diff --git a/src/agent/runtime/agent-invocation-contract.test.ts b/src/agent/runtime/agent-invocation-contract.test.ts index fa59ad2751..d1fc63b93a 100644 --- a/src/agent/runtime/agent-invocation-contract.test.ts +++ b/src/agent/runtime/agent-invocation-contract.test.ts @@ -5,6 +5,7 @@ import { buildRuntimeAgentControlPlaneStreamRequestFromInvocation, parseRuntimeAgentRunInvocation, parseRuntimeAgentRunInvocationOrError, + type RuntimeAgentControlPlaneStreamRequest, RuntimeAgentRunInvocationSchema, } from "../index.ts"; import { DEFAULT_LIMITS } from "#veryfront/security/input-validation/types.ts"; @@ -67,6 +68,21 @@ function createInvocation(overrides: Record = {}) { } describe("agent/runtime-agent-invocation-contract", () => { + it("keeps the legacy control-plane request shape source-compatible", () => { + const request: RuntimeAgentControlPlaneStreamRequest = { + agentId: "builder", + threadId: conversationId, + runId: "run_legacy_1", + messages: [], + tools: [], + context: [], + runtimeTargetKind: "main_branch", + agentSource: { type: "branch", branch: "main" }, + }; + + assertEquals(request.messageId, undefined); + }); + it("exports the control-plane runtime agent invocation schema from veryfront/agent", () => { const parsed = RuntimeAgentRunInvocationSchema.parse(createInvocation()); @@ -412,6 +428,7 @@ describe("agent/runtime-agent-invocation-contract", () => { agentId: "builder", threadId: conversationId, runId: "run_child_1", + messageId, parentRunId: "run_root_1", messages: parsed.messages, tools: parsed.tools, diff --git a/src/agent/runtime/agent-invocation-contract.ts b/src/agent/runtime/agent-invocation-contract.ts index 3c83cf8e85..95b4ba4f68 100644 --- a/src/agent/runtime/agent-invocation-contract.ts +++ b/src/agent/runtime/agent-invocation-contract.ts @@ -430,6 +430,7 @@ export type RuntimeAgentControlPlaneStreamRequest = { agentId: RuntimeAgentRunContext["agentId"]; threadId: RuntimeAgentRunContext["conversationId"]; runId: RuntimeAgentRunContext["runId"]; + messageId?: RuntimeAgentRunContext["messageId"]; taskId?: string; parentRunId?: Exclude; messages: RuntimeAgentRunInvocation["messages"]; @@ -456,6 +457,7 @@ export function buildRuntimeAgentControlPlaneStreamRequestFromInvocation( agentId: input.run.agentId, threadId: input.run.conversationId, runId: input.run.runId, + messageId: input.run.messageId, ...(input.taskId ? { taskId: input.taskId } : {}), ...(input.run.parentRunId ? { parentRunId: input.run.parentRunId } : {}), messages: input.messages, diff --git a/src/agent/runtime/index.ts b/src/agent/runtime/index.ts index 1fb14850e8..1c8a04813e 100644 --- a/src/agent/runtime/index.ts +++ b/src/agent/runtime/index.ts @@ -107,6 +107,8 @@ import { getRuntimeProviderReplayCheckpointMessageId, getRuntimeProviderReplayCheckpointPersister, getRuntimeProviderReplayCheckpoints, + getRuntimeProviderReplayCheckpointTurnComplete, + getRuntimeProviderReplayCheckpointTurnFailed, getRuntimeProviderTools, getRuntimeSourceIntegrationPolicy, getRuntimeToolExposureCheckpoint, @@ -690,6 +692,9 @@ async function persistToolExposureCheckpointBeforeContinuation(input: { type RuntimeProviderReplayCheckpointEmission = { state: ProviderReplayCheckpointEmissionState | undefined; persist: ((checkpoint: ProviderReplayCheckpoint) => void | Promise) | undefined; + complete: (() => void | Promise) | undefined; + fail: (() => void | Promise) | undefined; + failed: boolean; required: boolean; }; @@ -707,13 +712,36 @@ function resolveRuntimeProviderReplayCheckpointEmission( ? createProviderReplayCheckpointEmissionState({ messageId, existingCheckpoint }) : undefined, persist: getRuntimeProviderReplayCheckpointPersister(config), + complete: getRuntimeProviderReplayCheckpointTurnComplete(config), + fail: getRuntimeProviderReplayCheckpointTurnFailed(config), + failed: false, required: isRuntimeProviderReplayCheckpointPersistenceRequired(config), }; } +async function failProviderReplayCheckpointTurn( + emission: RuntimeProviderReplayCheckpointEmission, +): Promise { + if (emission.failed) return; + emission.failed = true; + await emission.fail?.(); +} + async function persistProviderReplayCheckpointAfterTurn(input: { emission: RuntimeProviderReplayCheckpointEmission; providerMetadata: Record | undefined; +}): Promise { + try { + await persistProviderReplayCheckpointAfterTurnUnsafe(input); + } catch (error) { + await failProviderReplayCheckpointTurn(input.emission); + throw error; + } +} + +async function persistProviderReplayCheckpointAfterTurnUnsafe(input: { + emission: RuntimeProviderReplayCheckpointEmission; + providerMetadata: Record | undefined; }): Promise { if (!input.emission.state) { if (input.emission.required) { @@ -721,13 +749,17 @@ async function persistProviderReplayCheckpointAfterTurn(input: { detail: "provider replay checkpoint message identity is required", }); } + await input.emission.complete?.(); return; } const checkpoint = captureProviderReplayCheckpoint( input.emission.state, input.providerMetadata, ); - if (!checkpoint) return; + if (!checkpoint) { + await input.emission.complete?.(); + return; + } if (!input.emission.persist) { if (input.emission.required) { throw DURABLE_RUN_EVENT_PERSISTENCE_FAILED.create({ @@ -737,6 +769,7 @@ async function persistProviderReplayCheckpointAfterTurn(input: { return; } await input.emission.persist(checkpoint); + await input.emission.complete?.(); } function isToolVisibleForStep(toolName: string, plan: ToolExposurePlan): boolean { @@ -1306,6 +1339,9 @@ export class AgentRuntime { const requestedModel = transport.requestedModel; const resolvedModelString = transport.resolvedModelString; const supportsToolCalling = supportsModelRuntimeToolCalling(transport.languageModel); + const providerReplayCheckpointEmission = resolveRuntimeProviderReplayCheckpointEmission( + this.config, + ); debugRuntimeModelRemap(requestedModel, resolvedModelString); return withSpan("agent.generate", async (span) => { @@ -1344,6 +1380,7 @@ export class AgentRuntime { context, runRuntimeContext, supportsToolCalling, + providerReplayCheckpointEmission, resolvedModelString, transport.languageModel, transport.headers, @@ -1360,6 +1397,9 @@ export class AgentRuntime { ) ), ); + }).catch(async (error) => { + await failProviderReplayCheckpointTurn(providerReplayCheckpointEmission); + throw error; }); } @@ -1426,6 +1466,9 @@ export class AgentRuntime { // Determine inference mode from the resolved model object, not the string. const isLocal = isLocalModelRuntime(languageModel); const supportsToolCalling = supportsModelRuntimeToolCalling(languageModel); + const providerReplayCheckpointEmission = resolveRuntimeProviderReplayCheckpointEmission( + this.config, + ); // Eagerly verify the model runtime is available. For local models this // checks that @huggingface/transformers can be imported. Must happen @@ -1487,6 +1530,7 @@ export class AgentRuntime { context, runRuntimeContext, supportsToolCalling, + providerReplayCheckpointEmission, resolvedModelString, languageModel, transport.headers, @@ -1516,6 +1560,13 @@ export class AgentRuntime { }); closeSSEStream(controller); } catch (error) { + try { + await failProviderReplayCheckpointTurn(providerReplayCheckpointEmission); + } catch (failureHookError) { + logger.debug("Provider replay failure hook rejected", { + error: failureHookError, + }); + } if (isAbortError(error, streamAbortSignal)) { closeSSEStream(controller); return; @@ -1559,6 +1610,7 @@ export class AgentRuntime { runtimeContext: Record | undefined, runRuntimeContext: AgentRunRuntimeContext, supportsToolCalling: boolean, + providerReplayCheckpointEmission: RuntimeProviderReplayCheckpointEmission, modelString?: string, resolvedModel?: ModelRuntime, headers?: HeadersInit, @@ -1583,9 +1635,6 @@ export class AgentRuntime { getRuntimeProviderReplayCheckpoints(this.config), { activeProvider: getActiveProviderReplayProvider(languageModel) }, ); - const providerReplayCheckpointEmission = resolveRuntimeProviderReplayCheckpointEmission( - this.config, - ); const totalUsage = { promptTokens: 0, completionTokens: 0, totalTokens: 0 }; if (!supportsToolCalling && this.config.tools) { @@ -2236,6 +2285,7 @@ export class AgentRuntime { runtimeContext: Record | undefined, runRuntimeContext: AgentRunRuntimeContext, supportsToolCalling: boolean, + providerReplayCheckpointEmission: RuntimeProviderReplayCheckpointEmission, modelString?: string, resolvedModel?: ModelRuntime, headers?: HeadersInit, @@ -2258,9 +2308,6 @@ export class AgentRuntime { getRuntimeProviderReplayCheckpoints(this.config), { activeProvider: getActiveProviderReplayProvider(languageModel) }, ); - const providerReplayCheckpointEmission = resolveRuntimeProviderReplayCheckpointEmission( - this.config, - ); const totalUsage = { promptTokens: 0, completionTokens: 0, totalTokens: 0 }; if (!supportsToolCalling && this.config.tools) { diff --git a/src/agent/runtime/provider-replay-emission.test.ts b/src/agent/runtime/provider-replay-emission.test.ts index 17378d9b2b..0485e38484 100644 --- a/src/agent/runtime/provider-replay-emission.test.ts +++ b/src/agent/runtime/provider-replay-emission.test.ts @@ -175,6 +175,9 @@ describe("provider replay checkpoint emission", () => { checkpoints.push(checkpoint); operations.push("persist:done"); }, + __vfProviderReplayCheckpointTurnComplete: () => { + operations.push("turn:complete"); + }, } as AgentConfig & RuntimeToolFilterConfig; const assistant = agent(config); @@ -185,6 +188,9 @@ describe("provider replay checkpoint emission", () => { } assertEquals(operations.indexOf("persist:done") < operations.indexOf("model:2"), true); + const completionIndex = operations.indexOf("turn:complete"); + assertEquals(completionIndex >= 0, true); + assertEquals(completionIndex < operations.indexOf("model:2"), true); assertEquals(checkpoints.length, 2); assertEquals(checkpoints[0]?.providerMessageBlockCounts, [2]); assertEquals(checkpoints[1]?.providerMessageBlockCounts, [2, 1]); @@ -230,6 +236,104 @@ describe("provider replay checkpoint emission", () => { assertEquals(model.callCount, 1); }); + it("closes the provider turn when no replay checkpoint is required", async () => { + let completedTurns = 0; + const model = scriptedModel([{ + text: "done", + providerMetadata: metadata([{ type: "text", text: "done" }]), + }], { + modelId: "anthropic/provider-replay-turn-boundary", + provider: "anthropic", + only: "generate", + }); + const config = { + id: "provider-replay-turn-boundary", + model: "anthropic/provider-replay-turn-boundary", + system: "Answer.", + skills: false, + maxSteps: 1, + resolveModelTransport: () => ({ model }), + __vfProviderReplayCheckpointTurnComplete: () => { + completedTurns++; + }, + __vfPersistProviderReplayCheckpoint: () => { + throw new Error("checkpoint persister must stay unused"); + }, + } as AgentConfig & RuntimeToolFilterConfig; + + await agent(config).generate({ input: "Answer" }); + + assertEquals(completedTurns, 1); + assertEquals(model.callCount, 1); + }); + + it("stops execution when the checkpoint persister rejects", async () => { + let failedTurns = 0; + const model = scriptedModel([{ + text: "done", + providerMetadata: metadata([{ + type: "thinking", + thinking: "", + signature: SIGNATURE, + }, { type: "text", text: "done" }]), + }], { + modelId: "anthropic/rejected-provider-replay-persister", + provider: "anthropic", + only: "generate", + }); + const config = { + id: "rejected-provider-replay-persister", + model: "anthropic/rejected-provider-replay-persister", + system: "Answer.", + skills: false, + maxSteps: 1, + resolveModelTransport: () => ({ model }), + __vfProviderReplayCheckpointMessageId: MESSAGE_ID, + __vfPersistProviderReplayCheckpoint: () => + Promise.reject(new Error("checkpoint sink rejected")), + __vfProviderReplayCheckpointTurnFailed: () => { + failedTurns++; + }, + } as AgentConfig & RuntimeToolFilterConfig; + + await assertRejects( + () => agent(config).generate({ input: "Answer" }), + Error, + "checkpoint sink rejected", + ); + assertEquals(model.callCount, 1); + assertEquals(failedTurns, 1); + }); + + it("fails the provider turn when streaming aborts before checkpoint capture", async () => { + let failedTurns = 0; + const model = scriptedModel([() => { + throw new Error("provider stream failed"); + }], { + modelId: "anthropic/failed-provider-replay-stream", + provider: "anthropic", + only: "stream", + }); + const config = { + id: "failed-provider-replay-stream", + model: "anthropic/failed-provider-replay-stream", + system: "Answer.", + skills: false, + maxSteps: 1, + resolveModelTransport: () => ({ model }), + __vfProviderReplayCheckpointMessageId: MESSAGE_ID, + __vfProviderReplayCheckpointTurnFailed: () => { + failedTurns++; + }, + } as AgentConfig & RuntimeToolFilterConfig; + + const stream = await agent(config).stream({ input: "Answer" }); + const body = await stream.toDataStreamResponse().text(); + + assertEquals(body.includes("provider stream failed"), true); + assertEquals(failedTurns, 1); + }); + it("required replay checkpoint persistence rejects a missing durable message identity", async () => { const model = scriptedModel([{ text: "done" }], { modelId: "anthropic/missing-provider-replay-message-id", diff --git a/src/agent/runtime/runtime-tool-config.ts b/src/agent/runtime/runtime-tool-config.ts index 275b2618b2..2003f6cdca 100644 --- a/src/agent/runtime/runtime-tool-config.ts +++ b/src/agent/runtime/runtime-tool-config.ts @@ -25,6 +25,8 @@ export type RuntimeToolFilterConfig = AgentConfig & { __vfPersistProviderReplayCheckpoint?: ( checkpoint: ProviderReplayCheckpoint, ) => void | Promise; + __vfProviderReplayCheckpointTurnComplete?: () => void | Promise; + __vfProviderReplayCheckpointTurnFailed?: () => void | Promise; __vfProviderReplayCheckpointPersistenceRequired?: boolean; __vfPersistToolExposureCheckpoint?: ( checkpoint: ToolExposureCheckpoint, @@ -147,6 +149,22 @@ export function getRuntimeProviderReplayCheckpointPersister( return typeof value === "function" ? value : undefined; } +/** Return the trusted hook that closes one provider response boundary. */ +export function getRuntimeProviderReplayCheckpointTurnComplete( + config: AgentConfig, +): (() => void | Promise) | undefined { + const value = (config as RuntimeToolFilterConfig).__vfProviderReplayCheckpointTurnComplete; + return typeof value === "function" ? value : undefined; +} + +/** Return the trusted hook that aborts one provider response boundary. */ +export function getRuntimeProviderReplayCheckpointTurnFailed( + config: AgentConfig, +): (() => void | Promise) | undefined { + const value = (config as RuntimeToolFilterConfig).__vfProviderReplayCheckpointTurnFailed; + return typeof value === "function" ? value : undefined; +} + /** Return whether provider replay state must be durable before continuation. */ export function isRuntimeProviderReplayCheckpointPersistenceRequired( config: AgentConfig, diff --git a/src/internal-agents/provider-replay-checkpoint-persister.test.ts b/src/internal-agents/provider-replay-checkpoint-persister.test.ts new file mode 100644 index 0000000000..d0898daa41 --- /dev/null +++ b/src/internal-agents/provider-replay-checkpoint-persister.test.ts @@ -0,0 +1,139 @@ +import "#veryfront/schemas/_test-setup.ts"; +import { + assertEquals, + assertInstanceOf, + assertRejects, + assertStringIncludes, +} from "#veryfront/testing/assert.ts"; +import { describe, it } from "#veryfront/testing/bdd.ts"; +import type { ProviderReplayCheckpoint } from "#veryfront/agent/runtime/provider-replay.ts"; +import { VeryfrontError } from "#veryfront/errors"; +import { createRunScopedProviderReplayCheckpointPersister } from "./provider-replay-checkpoint-persister.ts"; + +const RUN_ID = "run_checkpoint_1"; +const MESSAGE_ID = "10000000-1000-4000-8000-100000000001"; + +function checkpoint(): ProviderReplayCheckpoint { + return { + version: 1, + messageId: MESSAGE_ID, + provider: "anthropic", + providerBlocks: [{ + type: "provider-block", + provider: "anthropic", + block: { type: "thinking", thinking: "", signature: "" }, + }], + providerBlockPositions: [0], + providerMessageBlockCounts: [1], + totalPartCount: 1, + }; +} + +describe("run-scoped provider replay checkpoint persistence", () => { + it("keeps persistence pending until the exact-run append is acknowledged", async () => { + let capturedUrl: string | undefined; + let capturedInit: RequestInit | undefined; + let acknowledge: ((response: Response) => void) | undefined; + const responseGate = new Promise((resolve) => { + acknowledge = resolve; + }); + const persist = createRunScopedProviderReplayCheckpointPersister({ + apiUrl: "https://api.example.test/api", + runId: RUN_ID, + runEventAppendToken: "", + fetch: (input, init) => { + capturedUrl = String(input); + capturedInit = init; + return responseGate; + }, + }); + if (!persist) throw new Error("Expected a checkpoint persister"); + + let settled = false; + const persistence = persist(checkpoint()).then(() => { + settled = true; + }); + await Promise.resolve(); + + assertEquals(settled, false); + assertEquals(capturedUrl, `https://api.example.test/api/runs/${RUN_ID}/events`); + assertEquals(capturedInit?.method, "POST"); + assertEquals(new Headers(capturedInit?.headers).get("Authorization"), "Bearer "); + const body = JSON.parse(String(capturedInit?.body)) as { + events: Array>; + }; + assertEquals(body.events[0]?.type, "AGENT_RUN_PROVIDER_REPLAY_CHECKPOINT"); + assertEquals(body.events[0]?.messageId, MESSAGE_ID); + + acknowledge?.(Response.json({ appended_count: 1 })); + await persistence; + assertEquals(settled, true); + }); + + it("fails closed without exposing private response data or the credential", async () => { + const token = "private-test-token-that-must-not-appear"; + const persist = createRunScopedProviderReplayCheckpointPersister({ + apiUrl: "https://api.example.test", + runId: RUN_ID, + runEventAppendToken: token, + fetch: () => Promise.resolve(new Response("private response", { status: 503 })), + }); + if (!persist) throw new Error("Expected a checkpoint persister"); + + const error = await assertRejects( + () => persist(checkpoint()), + VeryfrontError, + "status 503", + ); + assertInstanceOf(error, VeryfrontError); + assertStringIncludes(error.slug, "durable-run-event-persistence-failed"); + assertEquals(String(error).includes(token), false); + assertEquals(String(error).includes("private response"), false); + }); + + it("aborts an in-flight append when the run is cancelled", async () => { + let requestSignal: AbortSignal | null | undefined; + const persist = createRunScopedProviderReplayCheckpointPersister({ + apiUrl: "https://api.example.test", + runId: RUN_ID, + runEventAppendToken: "", + fetch: (input, init) => { + requestSignal = new Request(input, init).signal; + return new Promise((_resolve, reject) => { + requestSignal?.addEventListener("abort", () => reject(requestSignal?.reason), { + once: true, + }); + }); + }, + }); + if (!persist) throw new Error("Expected a checkpoint persister"); + const controller = new AbortController(); + const reason = new DOMException("run cancelled", "AbortError"); + + const persistence = persist(checkpoint(), controller.signal); + await Promise.resolve(); + controller.abort(reason); + + let caught: unknown; + try { + await persistence; + } catch (error) { + caught = error; + } + assertEquals(caught, reason); + assertEquals(requestSignal?.aborted, true); + }); + + it("does not create a writer for a missing or malformed credential", () => { + for (const runEventAppendToken of [null, " token-with-whitespace "]) { + assertEquals( + createRunScopedProviderReplayCheckpointPersister({ + apiUrl: "https://api.example.test", + runId: RUN_ID, + runEventAppendToken, + }), + undefined, + ); + } + }); +}); diff --git a/src/internal-agents/provider-replay-checkpoint-persister.ts b/src/internal-agents/provider-replay-checkpoint-persister.ts new file mode 100644 index 0000000000..f6f7f9029d --- /dev/null +++ b/src/internal-agents/provider-replay-checkpoint-persister.ts @@ -0,0 +1,135 @@ +import { + createProviderReplayCheckpointEvent, + type ProviderReplayCheckpoint, +} from "#veryfront/agent/runtime/provider-replay.ts"; +import { DURABLE_RUN_EVENT_PERSISTENCE_FAILED } from "#veryfront/errors"; +import { + createVeryfrontApiRequestUrlResolver, + type VeryfrontApiRequestUrlResolver, +} from "#veryfront/platform/adapters/veryfront-api-url.ts"; +import { createOriginBoundOutboundFetch } from "#veryfront/security/http/outbound-fetch.ts"; + +const DEFAULT_PROVIDER_REPLAY_APPEND_TIMEOUT_MS = 15_000; +const MAX_RUN_EVENT_APPEND_TOKEN_BYTES = 4 * 1024; + +type Fetch = typeof globalThis.fetch; + +// Capture credential-touching intrinsics before project code can mutate the +// shared realm. The opaque writer token must never be passed to ambient +// prototype methods after project discovery. +const NativeTextEncoder = TextEncoder; +const apply = Reflect.apply; +const jsonStringify = JSON.stringify; +const stringTrim = String.prototype.trim; +const textEncoderEncode = NativeTextEncoder.prototype.encode; +const typedArrayPrototype = Object.getPrototypeOf(Uint8Array.prototype); +const typedArrayByteLengthGetterCandidate = Object.getOwnPropertyDescriptor( + typedArrayPrototype, + "byteLength", +)?.get; +const textEncoder = new NativeTextEncoder(); + +if (typeof typedArrayByteLengthGetterCandidate !== "function") { + throw new TypeError("Required Uint8Array byteLength intrinsic is unavailable"); +} +const typedArrayByteLengthGetter = typedArrayByteLengthGetterCandidate; + +function snapshotFetch(fetchImpl: Fetch): Fetch { + return (input, init) => apply(fetchImpl, undefined, [input, init]) as Promise; +} + +function isValidRunEventAppendToken(token: string | null | undefined): token is string { + if ( + typeof token !== "string" || token.length === 0 || + apply(stringTrim, token, []) !== token + ) { + return false; + } + const encoded = apply(textEncoderEncode, textEncoder, [token]) as Uint8Array; + return (apply(typedArrayByteLengthGetter, encoded, []) as number) <= + MAX_RUN_EVENT_APPEND_TOKEN_BYTES; +} + +function getAbortReason(signal: AbortSignal): unknown { + return signal.reason ?? new DOMException("This operation was aborted", "AbortError"); +} + +function persistenceFailure(detail: string) { + return DURABLE_RUN_EVENT_PERSISTENCE_FAILED.create({ detail }); +} + +/** Trusted host callback that durably appends checkpoints before continuation. */ +export type ProviderReplayCheckpointPersister = ( + checkpoint: ProviderReplayCheckpoint, + abortSignal?: AbortSignal, +) => Promise; + +/** Create an exact-run checkpoint writer backed by the API's durable append route. */ +export function createRunScopedProviderReplayCheckpointPersister(input: { + apiUrl: string; + runId: string; + runEventAppendToken: string | null | undefined; + timeoutMs?: number; + /** Explicit host-owned transport for tests. */ + fetch?: Fetch; +}): ProviderReplayCheckpointPersister | undefined { + if (!isValidRunEventAppendToken(input.runEventAppendToken)) return undefined; + + const token = input.runEventAppendToken; + const timeoutMs = input.timeoutMs ?? DEFAULT_PROVIDER_REPLAY_APPEND_TIMEOUT_MS; + const fetchImpl = input.fetch + ? snapshotFetch(input.fetch) + : createOriginBoundOutboundFetch(input.apiUrl); + const resolveApiUrl: VeryfrontApiRequestUrlResolver = createVeryfrontApiRequestUrlResolver( + input.apiUrl, + ); + const url = resolveApiUrl(`/runs/${encodeURIComponent(input.runId)}/events`); + + return async (checkpoint, abortSignal) => { + if (abortSignal?.aborted) throw getAbortReason(abortSignal); + + const controller = new AbortController(); + const timeoutError = persistenceFailure("Provider replay checkpoint persistence timed out"); + const timeout = setTimeout(() => controller.abort(timeoutError), timeoutMs); + const onAbort = () => controller.abort(getAbortReason(abortSignal!)); + abortSignal?.addEventListener("abort", onAbort, { once: true }); + + try { + const response = await fetchImpl(url, { + method: "POST", + headers: { + Accept: "application/json", + Authorization: `Bearer ${token}`, + "Cache-Control": "no-store", + "Content-Type": "application/json; charset=utf-8", + }, + body: apply(jsonStringify, JSON, [{ + events: [createProviderReplayCheckpointEvent(checkpoint)], + }]) as string, + cache: "no-store", + signal: controller.signal, + }); + if (!response.ok) { + await response.body?.cancel().catch(() => undefined); + throw persistenceFailure( + `Provider replay checkpoint append failed with status ${response.status}`, + ); + } + await response.body?.cancel().catch(() => undefined); + if (abortSignal?.aborted) throw getAbortReason(abortSignal); + } catch (error) { + if (abortSignal?.aborted) throw getAbortReason(abortSignal); + if (controller.signal.aborted) throw controller.signal.reason ?? timeoutError; + if ( + typeof error === "object" && error !== null && "slug" in error && + error.slug === DURABLE_RUN_EVENT_PERSISTENCE_FAILED.slug + ) { + throw error; + } + throw persistenceFailure("Provider replay checkpoint append failed"); + } finally { + clearTimeout(timeout); + abortSignal?.removeEventListener("abort", onAbort); + } + }; +} diff --git a/src/internal-agents/run-stream.test.ts b/src/internal-agents/run-stream.test.ts index f81cae6760..749ab36073 100644 --- a/src/internal-agents/run-stream.test.ts +++ b/src/internal-agents/run-stream.test.ts @@ -32,12 +32,15 @@ import type { RemoteToolSource, Tool } from "#veryfront/tool"; import { __resetLoggerConfigForTests, type LogEntry } from "#veryfront/utils/logger/logger.ts"; import type { AgentRunEventSink } from "#veryfront/runtime/model-call-context.ts"; import { getActiveRunEventSinks } from "#veryfront/runtime/run-event-sink-context.ts"; +import type { ProviderReplayCheckpoint } from "#veryfront/agent/runtime/provider-replay.ts"; import { AgentRunSessionManager } from "./session-manager.ts"; import { buildMergedTools, createRuntimeAgentStreamResponse, getExplicitlyDeniedToolNames, MODEL_CALL_CONTEXT_SSE_EVENT_NAME, + PROVIDER_REPLAY_PROTOCOL_HEADER, + PROVIDER_REPLAY_TURN_COMPLETE_SSE_EVENT_NAME, } from "./run-stream.ts"; function parseSseFrames(body: string): Array<{ event: string; data: unknown }> { @@ -536,6 +539,489 @@ describe("internal-agents/run-stream", () => { assertEquals(capturedMaxOutputTokens, 1200); }); + it("persists replay checkpoints before emitting the private turn boundary", async () => { + const sessionManager = new AgentRunSessionManager(); + const messageId = crypto.randomUUID(); + const checkpoint: ProviderReplayCheckpoint = { + version: 1, + messageId, + provider: "anthropic", + providerBlocks: [{ + type: "provider-block", + provider: "anthropic", + block: { type: "thinking", thinking: "", signature: "signed-private-block" }, + }], + providerBlockPositions: [0], + providerMessageBlockCounts: [1], + totalPartCount: 2, + }; + let capturedConfig: + | (Agent["config"] & { + __vfProviderReplayCheckpoints?: readonly ProviderReplayCheckpoint[]; + __vfProviderReplayCheckpointMessageId?: string; + __vfPersistProviderReplayCheckpoint?: ( + value: ProviderReplayCheckpoint, + ) => void | Promise; + __vfProviderReplayCheckpointTurnComplete?: () => void | Promise; + }) + | undefined; + const persistedCheckpoints: ProviderReplayCheckpoint[] = []; + const agent = { + id: "test", + config: { + id: "test", + model: "anthropic/claude-opus-4-8", + system: "test", + }, + } as unknown as Agent; + const input = { + agentId: "test", + threadId: crypto.randomUUID(), + runId: "run_1", + messageId, + messages: [], + tools: [], + context: [], + } as Parameters[0]; + + const response = await createRuntimeAgentStreamResponse(input, agent, { + sessionManager, + providerReplayCheckpointEmissionEnabled: true, + providerReplayCheckpoints: [], + persistProviderReplayCheckpoint: (value) => { + persistedCheckpoints.push(value); + return Promise.resolve(); + }, + createRuntime: (runtimeAgent) => { + capturedConfig = runtimeAgent.config as typeof capturedConfig; + return { + stream: async () => { + await capturedConfig?.__vfPersistProviderReplayCheckpoint?.(checkpoint); + await capturedConfig?.__vfProviderReplayCheckpointTurnComplete?.(); + return new ReadableStream({ + start(controller) { + controller.close(); + }, + }); + }, + }; + }, + }); + const frames = parseSseFrames(await response.text()); + + assertEquals(response.headers.get(PROVIDER_REPLAY_PROTOCOL_HEADER), "1"); + assertEquals(capturedConfig?.__vfProviderReplayCheckpointMessageId, messageId); + assertEquals(capturedConfig?.__vfProviderReplayCheckpoints, undefined); + assertEquals(persistedCheckpoints, [checkpoint]); + assertEquals(JSON.stringify(frames).includes("signed-private-block"), false); + assertEquals( + frames.some((frame) => frame.event === PROVIDER_REPLAY_TURN_COMPLETE_SSE_EVENT_NAME), + true, + ); + assertEquals( + frames.findIndex((frame) => frame.event === PROVIDER_REPLAY_TURN_COMPLETE_SSE_EVENT_NAME) < + frames.findIndex((frame) => frame.event === "RunError"), + true, + ); + }); + + it("continues checkpoint emission after the host gate is disabled", async () => { + const sessionManager = new AgentRunSessionManager(); + const messageId = crypto.randomUUID(); + const checkpoint: ProviderReplayCheckpoint = { + version: 1, + messageId, + provider: "anthropic", + providerBlocks: [{ + type: "provider-block", + provider: "anthropic", + block: { type: "thinking", thinking: "", signature: "existing-signature" }, + }], + providerBlockPositions: [0], + providerMessageBlockCounts: [1], + totalPartCount: 1, + }; + let persistCheckpoint: + | ((value: ProviderReplayCheckpoint) => void | Promise) + | undefined; + let completeProviderReplayTurn: (() => void | Promise) | undefined; + const persistedCheckpoints: ProviderReplayCheckpoint[] = []; + const agent = { + id: "test", + config: { + id: "test", + model: "anthropic/claude-opus-4-8", + system: "test", + }, + } as unknown as Agent; + + const response = await createRuntimeAgentStreamResponse( + { + threadId: crypto.randomUUID(), + runId: "run_1", + messageId, + messages: [], + tools: [], + context: [], + }, + agent, + { + sessionManager, + providerReplayCheckpointEmissionEnabled: false, + providerReplayCheckpoints: [checkpoint], + persistProviderReplayCheckpoint: (value) => { + persistedCheckpoints.push(value); + return Promise.resolve(); + }, + createRuntime: (runtimeAgent) => { + persistCheckpoint = (runtimeAgent.config as Agent["config"] & { + __vfPersistProviderReplayCheckpoint?: ( + value: ProviderReplayCheckpoint, + ) => void | Promise; + __vfProviderReplayCheckpointTurnComplete?: () => void | Promise; + }).__vfPersistProviderReplayCheckpoint; + completeProviderReplayTurn = (runtimeAgent.config as Agent["config"] & { + __vfProviderReplayCheckpointTurnComplete?: () => void | Promise; + }).__vfProviderReplayCheckpointTurnComplete; + return { + stream: async () => { + await persistCheckpoint?.(checkpoint); + await completeProviderReplayTurn?.(); + return new ReadableStream({ + start(controller) { + controller.close(); + }, + }); + }, + }; + }, + }, + ); + const frames = parseSseFrames(await response.text()); + + assertEquals(typeof persistCheckpoint, "function"); + assertEquals(typeof completeProviderReplayTurn, "function"); + assertEquals(persistedCheckpoints, [checkpoint]); + assertEquals(JSON.stringify(frames).includes("existing-signature"), false); + }); + + it("holds tool dispatch until durable persistence and the turn boundary complete", async () => { + const sessionManager = new AgentRunSessionManager(); + const messageId = crypto.randomUUID(); + const checkpoint: ProviderReplayCheckpoint = { + version: 1, + messageId, + provider: "anthropic", + providerBlocks: [{ + type: "provider-block", + provider: "anthropic", + block: { type: "thinking", thinking: "", signature: "ordered-signature" }, + }], + providerBlockPositions: [0], + providerMessageBlockCounts: [1], + totalPartCount: 1, + }; + let runtimeConfig: + | (Agent["config"] & { + __vfPersistProviderReplayCheckpoint?: ( + value: ProviderReplayCheckpoint, + ) => void | Promise; + __vfProviderReplayCheckpointTurnComplete?: () => void | Promise; + }) + | undefined; + const operations: string[] = []; + const agent = { + id: "test", + config: { + id: "test", + model: "anthropic/claude-opus-4-8", + system: "test", + }, + } as unknown as Agent; + + const response = await createRuntimeAgentStreamResponse( + { + threadId: crypto.randomUUID(), + runId: "run_1", + messageId, + messages: [], + tools: [], + context: [], + }, + agent, + { + sessionManager, + providerReplayCheckpointEmissionEnabled: true, + persistProviderReplayCheckpoint: async () => { + operations.push("persist:start"); + await Promise.resolve(); + operations.push("persist:done"); + }, + createRuntime: (runtimeAgent) => { + runtimeConfig = runtimeAgent.config as typeof runtimeConfig; + return { + stream: async () => + new ReadableStream({ + start(controller) { + controller.enqueue( + new TextEncoder().encode( + 'data: {"type":"step-start"}\n\ndata: {"type":"tool-input-start","toolCallId":"tool-1","toolName":"lookup"}\n\ndata: {"type":"tool-input-available","toolCallId":"tool-1","toolName":"lookup","input":{}}\n\n', + ), + ); + setTimeout(async () => { + await runtimeConfig?.__vfPersistProviderReplayCheckpoint?.(checkpoint); + await runtimeConfig?.__vfProviderReplayCheckpointTurnComplete?.(); + controller.close(); + }, 0); + }, + }), + }; + }, + }, + ); + const frames = parseSseFrames(await response.text()); + + assertEquals(operations, ["persist:start", "persist:done"]); + assertEquals( + frames.findIndex((frame) => frame.event === PROVIDER_REPLAY_TURN_COMPLETE_SSE_EVENT_NAME) < + frames.findIndex((frame) => frame.event === "ToolCallEnd"), + true, + ); + assertEquals(JSON.stringify(frames).includes("ordered-signature"), false); + }); + + it("fails closed when checkpoint emission has no runtime message identity", async () => { + const sessionManager = new AgentRunSessionManager(); + const agent = { + id: "test", + config: { + id: "test", + model: "anthropic/claude-opus-4-8", + system: "test", + }, + } as unknown as Agent; + + await assertRejects( + () => + createRuntimeAgentStreamResponse( + { + threadId: crypto.randomUUID(), + runId: "run_1", + messages: [], + tools: [], + context: [], + }, + agent, + { + sessionManager, + providerReplayCheckpointEmissionEnabled: true, + }, + ), + Error, + "Provider replay checkpoint emission requires a runtime message identity", + ); + }); + + it("fails before provider execution when no trusted checkpoint writer is available", async () => { + const sessionManager = new AgentRunSessionManager(); + let runtimeCreated = false; + const agent = { + id: "test", + config: { + id: "test", + model: "anthropic/claude-opus-4-8", + system: "test", + }, + } as unknown as Agent; + + await assertRejects( + () => + createRuntimeAgentStreamResponse( + { + threadId: crypto.randomUUID(), + runId: "run_1", + messageId: crypto.randomUUID(), + messages: [], + tools: [], + context: [], + }, + agent, + { + sessionManager, + providerReplayCheckpointEmissionEnabled: true, + createRuntime: () => { + runtimeCreated = true; + throw new Error("runtime must not be created"); + }, + }, + ), + Error, + "trusted run-event append token", + ); + assertEquals(runtimeCreated, false); + }); + + it("settles an open replay turn when the provider stream fails", async () => { + const sessionManager = new AgentRunSessionManager(); + const agent = { + id: "test", + config: { + id: "test", + model: "anthropic/claude-opus-4-8", + system: "test", + }, + } as unknown as Agent; + + const response = await createRuntimeAgentStreamResponse( + { + threadId: crypto.randomUUID(), + runId: "run_1", + messageId: crypto.randomUUID(), + messages: [], + tools: [], + context: [], + }, + agent, + { + sessionManager, + providerReplayCheckpointEmissionEnabled: true, + persistProviderReplayCheckpoint: () => Promise.resolve(), + createRuntime: () => ({ + stream: async () => + new ReadableStream({ + start(controller) { + controller.enqueue( + new TextEncoder().encode( + 'data: {"type":"step-start"}\n\ndata: {"type":"error","error":"provider stream failed"}\n\n', + ), + ); + controller.close(); + }, + }), + }), + }, + ); + + const body = await response.text(); + + assertStringIncludes(body, "event: RunError"); + assertEquals(body.includes("event: RunFinished"), false); + }); + + it("releases a pending tool boundary when the runtime turn fails", async () => { + const sessionManager = new AgentRunSessionManager(); + let failProviderReplayTurn: (() => void | Promise) | undefined; + const agent = { + id: "test", + config: { + id: "test", + model: "anthropic/claude-opus-4-8", + system: "test", + }, + } as unknown as Agent; + + const response = await createRuntimeAgentStreamResponse( + { + threadId: crypto.randomUUID(), + runId: "run_1", + messageId: crypto.randomUUID(), + messages: [], + tools: [], + context: [], + }, + agent, + { + sessionManager, + providerReplayCheckpointEmissionEnabled: true, + persistProviderReplayCheckpoint: () => Promise.resolve(), + createRuntime: (runtimeAgent) => { + failProviderReplayTurn = (runtimeAgent.config as Agent["config"] & { + __vfProviderReplayCheckpointTurnFailed?: () => void | Promise; + }).__vfProviderReplayCheckpointTurnFailed; + return { + stream: async () => + new ReadableStream({ + start(controller) { + controller.enqueue( + new TextEncoder().encode( + 'data: {"type":"step-start"}\n\ndata: {"type":"tool-input-start","toolCallId":"tool-1","toolName":"lookup"}\n\ndata: {"type":"tool-input-available","toolCallId":"tool-1","toolName":"lookup","input":{}}\n\n', + ), + ); + setTimeout(async () => { + await failProviderReplayTurn?.(); + controller.enqueue( + new TextEncoder().encode( + 'data: {"type":"error","error":"provider stream failed"}\n\n', + ), + ); + controller.close(); + }, 0); + }, + }), + }; + }, + }, + ); + + const body = await response.text(); + + assertStringIncludes(body, "event: RunError"); + assertEquals(body.includes("event: RunFinished"), false); + }); + + it("aborts a pending replay boundary when the run is cancelled", async () => { + const sessionManager = new AgentRunSessionManager(); + let runtimeCancelCalls = 0; + const runId = "run_cancel_pending_replay"; + const agent = { + id: "test", + config: { + id: "test", + model: "anthropic/claude-opus-4-8", + system: "test", + }, + } as unknown as Agent; + const response = await createRuntimeAgentStreamResponse( + { + threadId: crypto.randomUUID(), + runId, + messageId: crypto.randomUUID(), + messages: [], + tools: [], + context: [], + }, + agent, + { + sessionManager, + providerReplayCheckpointEmissionEnabled: true, + persistProviderReplayCheckpoint: () => Promise.resolve(), + createRuntime: () => ({ + stream: async () => + new ReadableStream({ + start(controller) { + controller.enqueue( + new TextEncoder().encode( + 'data: {"type":"step-start"}\n\ndata: {"type":"tool-input-start","toolCallId":"tool-1","toolName":"lookup"}\n\ndata: {"type":"tool-input-available","toolCallId":"tool-1","toolName":"lookup","input":{}}\n\n', + ), + ); + }, + cancel() { + runtimeCancelCalls++; + }, + }), + }), + }, + ); + + const body = response.text(); + await new Promise((resolve) => setTimeout(resolve, 0)); + assertEquals(sessionManager.cancelRun(runId), true); + await body; + + assertEquals(runtimeCancelCalls, 1); + assertEquals(sessionManager.getRunStatus(runId), null); + }); + it("composes the runtime system prompt with project, environment, and tool context", async () => { const sessionManager = new AgentRunSessionManager(); let capturedAgent: Agent | undefined; diff --git a/src/internal-agents/run-stream.ts b/src/internal-agents/run-stream.ts index 9bc998a76f..e5fd739ec2 100644 --- a/src/internal-agents/run-stream.ts +++ b/src/internal-agents/run-stream.ts @@ -73,6 +73,9 @@ import { composeInternalAgentRunSystemPrompt } from "./run-system-prompt.ts"; import type { RuntimeRunAgentInput } from "./schema.ts"; import { serverLogger } from "#veryfront/utils"; import { compareStrings } from "#veryfront/utils/compare.ts"; +import { type ProviderReplayCheckpoint } from "#veryfront/agent/runtime/provider-replay.ts"; +import { DURABLE_RUN_EVENT_PERSISTENCE_FAILED } from "#veryfront/errors"; +import type { ProviderReplayCheckpointPersister } from "./provider-replay-checkpoint-persister.ts"; const getAnyObjectSchema = defineSchema((v) => v.record(v.string(), v.unknown())); const anyObjectSchema = lazySchema(getAnyObjectSchema) as Schema>; @@ -88,6 +91,8 @@ const INTERNAL_AGENT_RUNTIME_HEARTBEAT_FRAME = new TextEncoder().encode( * folding it into the run's public event sequence. */ export const MODEL_CALL_CONTEXT_SSE_EVENT_NAME = "AgentRunModelCallContext"; +export const PROVIDER_REPLAY_TURN_COMPLETE_SSE_EVENT_NAME = "AgentRunProviderReplayTurnComplete"; +export const PROVIDER_REPLAY_PROTOCOL_HEADER = "X-Veryfront-Provider-Replay-Protocol"; type RuntimeFilteredAgent = Agent & { config: Agent["config"] & { @@ -180,6 +185,9 @@ export interface RuntimeAgentStreamExecutionDeps { abortSignal?: AbortSignal, ) => Promise>; }; + providerReplayCheckpointEmissionEnabled?: boolean; + providerReplayCheckpoints?: readonly ProviderReplayCheckpoint[]; + persistProviderReplayCheckpoint?: ProviderReplayCheckpointPersister; } function createInjectedStudioTool( @@ -779,6 +787,77 @@ function createModelCallContextRelay( }; } +type ProviderReplayPrivateFrame = { + event: string; + payload: Record; +}; + +function createProviderReplayCheckpointRelay(): { + complete: (messageId: string) => Promise; + fail: () => Promise; + takeCompletedTurn: () => Promise; + hasCompletedTurn: () => boolean; +} { + const buffered: ProviderReplayPrivateFrame[] = []; + let resolvePending: + | ((frames: ProviderReplayPrivateFrame[]) => void) + | undefined; + let rejectPending: ((error: Error) => void) | undefined; + let terminalError: Error | undefined; + + const takeReadyTurn = (): ProviderReplayPrivateFrame[] | undefined => { + const boundaryIndex = buffered.findIndex((frame) => + frame.event === PROVIDER_REPLAY_TURN_COMPLETE_SSE_EVENT_NAME + ); + return boundaryIndex >= 0 ? buffered.splice(0, boundaryIndex + 1) : undefined; + }; + const resolveIfReady = () => { + if (!resolvePending) return; + const frames = takeReadyTurn(); + if (!frames) return; + const resolve = resolvePending; + resolvePending = undefined; + rejectPending = undefined; + resolve(frames); + }; + + return { + complete: async (messageId) => { + buffered.push({ + event: PROVIDER_REPLAY_TURN_COMPLETE_SSE_EVENT_NAME, + payload: { + type: "AGENT_RUN_PROVIDER_REPLAY_TURN_COMPLETE", + messageId, + }, + }); + resolveIfReady(); + }, + fail: async () => { + if (terminalError) return; + terminalError = new Error("Provider replay turn failed before its boundary"); + buffered.splice(0); + const reject = rejectPending; + resolvePending = undefined; + rejectPending = undefined; + reject?.(terminalError); + }, + takeCompletedTurn: () => { + if (terminalError) return Promise.reject(terminalError); + const frames = takeReadyTurn(); + if (frames) return Promise.resolve(frames); + if (resolvePending) { + return Promise.reject(new Error("Provider replay turn boundary already has a waiter")); + } + return new Promise((resolve, reject) => { + resolvePending = resolve; + rejectPending = reject; + }); + }, + hasCompletedTurn: () => + buffered.some((frame) => frame.event === PROVIDER_REPLAY_TURN_COMPLETE_SSE_EVENT_NAME), + }; +} + export async function createRuntimeAgentStreamResponse( input: RuntimeRunAgentInput, agent: Agent, @@ -802,6 +881,8 @@ export async function createRuntimeAgentStreamResponse( let closeSandbox = createIdempotentAsyncCleanup(); const timing = createAgentRunEventTimingAnchor(); const modelCallContextRelay = createModelCallContextRelay(timing); + const providerReplayCheckpointRelay = createProviderReplayCheckpointRelay(); + let shouldEmitProviderReplayCheckpoints = false; try { const forwardedAllowedRemoteToolNames = getAllowedRemoteToolNames(input.forwardedProps); const sourceAllowedRemoteToolNames = getAgentAllowedRemoteToolNames(agent); @@ -953,6 +1034,23 @@ export async function createRuntimeAgentStreamResponse( }); }; const systemPrompt = await resolveSystemPrompt(); + const providerReplayCheckpoints = deps.providerReplayCheckpoints ?? []; + if (deps.providerReplayCheckpointEmissionEnabled === true && !input.messageId) { + throw new Error( + "Provider replay checkpoint emission requires a runtime message identity", + ); + } + shouldEmitProviderReplayCheckpoints = Boolean( + input.messageId && + (deps.providerReplayCheckpointEmissionEnabled === true || + providerReplayCheckpoints.some((checkpoint) => checkpoint.messageId === input.messageId)), + ); + if (shouldEmitProviderReplayCheckpoints && !deps.persistProviderReplayCheckpoint) { + throw DURABLE_RUN_EVENT_PERSISTENCE_FAILED.create({ + detail: + "A trusted run-event append token is required to persist a private provider replay checkpoint", + }); + } const runtimeAgent: RuntimeFilteredAgent = { ...agent, config: { @@ -966,6 +1064,24 @@ export async function createRuntimeAgentStreamResponse( ...(forwardedIntegrationToolDefs !== undefined ? { __vfForwardedIntegrationToolDefs: forwardedIntegrationToolDefs } : {}), + ...(providerReplayCheckpoints.length > 0 + ? { __vfProviderReplayCheckpoints: providerReplayCheckpoints } + : {}), + ...(shouldEmitProviderReplayCheckpoints + ? { + __vfProviderReplayCheckpointMessageId: input.messageId, + __vfPersistProviderReplayCheckpoint: (checkpoint: ProviderReplayCheckpoint) => + deps.persistProviderReplayCheckpoint!(checkpoint, abortSignal), + __vfProviderReplayCheckpointPersistenceRequired: true, + } + : {}), + ...(shouldEmitProviderReplayCheckpoints + ? { + __vfProviderReplayCheckpointTurnComplete: () => + providerReplayCheckpointRelay.complete(input.messageId!), + __vfProviderReplayCheckpointTurnFailed: providerReplayCheckpointRelay.fail, + } + : {}), }, }; const runtime = deps.createRuntime?.(runtimeAgent, mergedTools) ?? @@ -1133,6 +1249,7 @@ export async function createRuntimeAgentStreamResponse( threadId: input.threadId, agentId: agent.id, }); + void providerReplayCheckpointRelay.fail(); void cancelReaderOnce(cancellationError); }; @@ -1151,6 +1268,49 @@ export async function createRuntimeAgentStreamResponse( modelCallContextRelay.attach((event) => enqueueIfAttached(MODEL_CALL_CONTEXT_SSE_EVENT_NAME, event) ); + const enqueueProviderReplayFrame = (frame: ProviderReplayPrivateFrame) => { + if (!clientAttached) { + throw new Error( + "Provider replay checkpoint stream detached before persistence", + ); + } + try { + controller.enqueue( + formatAgUiEvent(frame.event, frame.payload), + ); + } catch { + clientAttached = false; + throw new Error( + "Provider replay checkpoint stream detached before persistence", + ); + } + }; + let providerReplayStepOpen = false; + const flushProviderReplayTurn = async () => { + if (!providerReplayStepOpen) return; + for (const frame of await providerReplayCheckpointRelay.takeCompletedTurn()) { + enqueueProviderReplayFrame(frame); + } + providerReplayStepOpen = false; + }; + const emitMappedEvent = async (mappedEvent: { + event: string; + payload: Record; + }) => { + if (mappedEvent.event === "StepStarted") { + await flushProviderReplayTurn(); + providerReplayStepOpen = shouldEmitProviderReplayCheckpoints; + } + if (mappedEvent.event === "ToolCallEnd") { + await flushProviderReplayTurn(); + } + if (mappedEvent.event === "RunError" && providerReplayStepOpen) { + await providerReplayCheckpointRelay.fail(); + providerReplayStepOpen = false; + } + prepareToolResultIfNeeded(mappedEvent.event, mappedEvent.payload); + enqueueIfAttached(mappedEvent.event, mappedEvent.payload); + }; heartbeatTimer = setInterval( enqueueHeartbeatIfAttached, INTERNAL_AGENT_RUNTIME_HEARTBEAT_INTERVAL_MS, @@ -1184,8 +1344,7 @@ export async function createRuntimeAgentStreamResponse( for (const event of parsed.events) { for (const mappedEvent of mapRuntimeEventToAgUi(state, event)) { - prepareToolResultIfNeeded(mappedEvent.event, mappedEvent.payload); - enqueueIfAttached(mappedEvent.event, mappedEvent.payload); + await emitMappedEvent(mappedEvent); } } } @@ -1195,12 +1354,17 @@ export async function createRuntimeAgentStreamResponse( const trailingEvents = parseSseJsonEvents(`${remainder}\n\n`); for (const event of trailingEvents.events) { for (const mappedEvent of mapRuntimeEventToAgUi(state, event)) { - prepareToolResultIfNeeded(mappedEvent.event, mappedEvent.payload); - enqueueIfAttached(mappedEvent.event, mappedEvent.payload); + await emitMappedEvent(mappedEvent); } } throwIfAborted(); + await flushProviderReplayTurn(); + while (providerReplayCheckpointRelay.hasCompletedTurn()) { + for (const frame of await providerReplayCheckpointRelay.takeCompletedTurn()) { + enqueueProviderReplayFrame(frame); + } + } for (const mappedEvent of finalizeRunEvents(state, completedResponse)) { enqueueIfAttached(mappedEvent.event, mappedEvent.payload); @@ -1369,6 +1533,7 @@ export async function createRuntimeAgentStreamResponse( "Content-Type": "text/event-stream", "Cache-Control": "no-cache", Connection: "keep-alive", + ...(shouldEmitProviderReplayCheckpoints ? { [PROVIDER_REPLAY_PROTOCOL_HEADER]: "1" } : {}), }, }); } diff --git a/src/internal-agents/schema.test.ts b/src/internal-agents/schema.test.ts index f491100ca3..2caf1cbbe2 100644 --- a/src/internal-agents/schema.test.ts +++ b/src/internal-agents/schema.test.ts @@ -76,6 +76,31 @@ describe("internal-agents/schema", () => { ); }); + it("preserves signed top-level provider replay checkpoints", () => { + const checkpoints = [{ + version: 1, + messageId: "assistant-message-1", + provider: "anthropic", + providerBlocks: [], + providerBlockPositions: [], + totalPartCount: 1, + }]; + const parsed = getInternalAgentStreamRequestSchema().parse({ + agentId: "agent_1", + threadId: "10000000-1000-4000-8000-100000000001", + runId: "run_1", + ...MAIN_BRANCH_TARGET, + agentSource: { type: "branch", branch: "main" }, + messages: [], + serverResolvedProviderReplayCheckpoints: checkpoints, + }); + + assertEquals( + toRuntimeRunAgentInput(parsed).serverResolvedProviderReplayCheckpoints, + checkpoints, + ); + }); + it("rejects oversized injected tool parameters", () => { assertThrows( () => diff --git a/src/internal-agents/schema.ts b/src/internal-agents/schema.ts index 675fa0574e..4f0df44cd8 100644 --- a/src/internal-agents/schema.ts +++ b/src/internal-agents/schema.ts @@ -71,6 +71,7 @@ export const getInternalAgentControlPlaneStreamRequestSchema = defineSchema((v) agentId: getAgentIdSchema(), threadId: v.string().uuid(), runId: getRunIdSchema(), + messageId: v.string().uuid().optional(), taskId: getRuntimeAgentTaskIdSchema().optional(), parentRunId: getRunIdSchema().optional(), state: v.unknown().optional(), @@ -96,6 +97,7 @@ export const getInternalAgentControlPlaneStreamRequestSchema = defineSchema((v) (value) => value === undefined || isWithinJsonSizeLimit(value, MAX_FORWARDED_PROPS_BYTES), { message: "forwardedProps must be less than 192 KB" }, ), + serverResolvedProviderReplayCheckpoints: v.unknown().optional(), }).strict().superRefine((input, ctx) => { validateRuntimeAgentTargetSelection(input, ctx); validateRuntimeAgentSourceTargetBinding(input, ctx); @@ -404,6 +406,7 @@ export function toRuntimeRunAgentInput( return { threadId: input.threadId, runId: input.runId, + ...(input.messageId ? { messageId: input.messageId } : {}), ...(input.taskId ? { taskId: input.taskId } : {}), ...(input.parentRunId ? { parentRunId: input.parentRunId } : {}), ...(input.state !== undefined ? { state: input.state } : {}), @@ -412,6 +415,11 @@ export function toRuntimeRunAgentInput( tools: input.tools, context: input.context, ...(input.forwardedProps ? { forwardedProps: input.forwardedProps } : {}), + ...(input.serverResolvedProviderReplayCheckpoints !== undefined + ? { + serverResolvedProviderReplayCheckpoints: input.serverResolvedProviderReplayCheckpoints, + } + : {}), } as RuntimeRunAgentInput; } @@ -436,6 +444,8 @@ export type RuntimeContextItem = AgUiRuntimeContextItem; export type RuntimeRunAgentInput = AgUiRuntimeRequest & { taskId?: string; allowDelegation?: boolean; + messageId?: string; + serverResolvedProviderReplayCheckpoints?: unknown; }; export type InternalAgentStreamRequest = InferSchema< ReturnType diff --git a/src/server/handlers/request/agent-stream.handler.test.ts b/src/server/handlers/request/agent-stream.handler.test.ts index 2c803ab13c..ee44e9d7bc 100644 --- a/src/server/handlers/request/agent-stream.handler.test.ts +++ b/src/server/handlers/request/agent-stream.handler.test.ts @@ -16,6 +16,7 @@ import { type RuntimeRemoteToolConfig, } from "#veryfront/agent/runtime/mcp-server-tool-sources.ts"; import { getRuntimeSourceIntegrationPolicy } from "#veryfront/agent/runtime/runtime-tool-config.ts"; +import type { ProviderReplayCheckpoint } from "#veryfront/agent/runtime/provider-replay.ts"; import { dynamicTool } from "#veryfront/tool"; import { markRemoteToolProvenance } from "#veryfront/tool/remote-tool-provenance.ts"; import { assertEquals, assertExists, assertStringIncludes } from "#veryfront/testing/assert.ts"; @@ -469,6 +470,106 @@ describe("server/handlers/request/agent-stream.handler", () => { assertStringIncludes(text, "event: RunFinished"); }); + it("binds the exact-run credential to acknowledged checkpoint persistence", async () => { + const messageId = "10000000-1000-4000-8000-100000000002"; + const checkpoint: ProviderReplayCheckpoint = { + version: 1, + messageId, + provider: "anthropic", + providerBlocks: [{ + type: "provider-block", + provider: "anthropic", + block: { type: "thinking", thinking: "", signature: "" }, + }], + providerBlockPositions: [0], + providerMessageBlockCounts: [1], + totalPartCount: 1, + }; + let factoryInput: + | { + runId: string; + runEventAppendToken: string | null | undefined; + } + | undefined; + const operations: string[] = []; + const handler = createTestAgentStreamHandler({ + ensureProjectDiscovery: () => Promise.resolve(createEmptyDiscoveryResult()), + getAgent: (id) => id === "assistant-1" ? createAgent("assistant-1") : undefined, + getAllAgentIds: () => ["assistant-1"], + sessionManager: new AgentRunSessionManager(), + providerReplayCheckpointEmissionEnabled: true, + createRunScopedProviderReplayCheckpointPersister: (input) => { + factoryInput = input; + return async (value) => { + operations.push("append:start"); + await Promise.resolve(); + assertEquals(value, checkpoint); + operations.push("append:acknowledged"); + }; + }, + createRuntime: (runtimeAgent) => ({ + stream: async (_messages, _context, callbacks) => { + const config = runtimeAgent.config as Agent["config"] & { + __vfPersistProviderReplayCheckpoint?: ( + value: ProviderReplayCheckpoint, + ) => void | Promise; + __vfProviderReplayCheckpointTurnComplete?: () => void | Promise; + }; + operations.push("runtime:before-persist"); + await config.__vfPersistProviderReplayCheckpoint?.(checkpoint); + operations.push("runtime:after-persist"); + await config.__vfProviderReplayCheckpointTurnComplete?.(); + callbacks?.onFinish?.({ + text: "", + messages: [], + toolCalls: [], + status: "completed", + usage: undefined, + metadata: { finishReason: "stop" }, + }); + return new ReadableStream({ + start(controller) { + controller.enqueue(encodeDataStreamEvent({ type: "step-start" })); + controller.enqueue(encodeDataStreamEvent({ type: "step-end" })); + controller.close(); + }, + }); + }, + }), + }); + const body = createAgentStreamRequestBody({ + serverResolvedProviderReplayCheckpoints: [checkpoint], + }); + const { jws, publicKeyPem } = await createControlPlaneSignature(body, { + requestId: "run_1", + }); + + const result = await handler.handle( + new Request("https://example.com/api/control-plane/runs/run_1/stream", { + method: "POST", + headers: { + "content-type": "application/json", + "x-veryfront-control-plane-jws": jws, + "x-veryfront-run-event-token": "", + }, + body, + }), + createCtx(publicKeyPem), + ); + + assertExists(result.response); + assertEquals(result.response.status, 200); + await result.response.text(); + assertEquals(factoryInput?.runId, "run_1"); + assertEquals(factoryInput?.runEventAppendToken, ""); + assertEquals(operations, [ + "runtime:before-persist", + "append:start", + "append:acknowledged", + "runtime:after-persist", + ]); + }); + it("accepts the public control-plane stream route", async () => { const handler = createTestAgentStreamHandler({ ensureProjectDiscovery: async () => createEmptyDiscoveryResult(), diff --git a/src/server/handlers/request/agent-stream.handler.ts b/src/server/handlers/request/agent-stream.handler.ts index 75aaa84911..be952fd4d8 100644 --- a/src/server/handlers/request/agent-stream.handler.ts +++ b/src/server/handlers/request/agent-stream.handler.ts @@ -95,10 +95,14 @@ import { prepareDeclarativeConfigContext } from "#veryfront/config/declarative-e import { normalizeSourceIntegrationPolicy } from "#veryfront/integrations/source-policy.ts"; import { runWithExactSourceIntegrationPolicy } from "#veryfront/integrations/source-policy-context.ts"; import { compareStrings } from "#veryfront/utils/compare.ts"; +import { isProviderReplayCheckpointEmissionEnabled } from "#veryfront/agent/hosted/chat-preparation.ts"; +import { getServerResolvedProviderReplayCheckpoints } from "#veryfront/agent/hosted/runtime-request-config.ts"; +import { RUN_EVENT_APPEND_TOKEN_HEADER } from "#veryfront/agent/hosted/chat-request-parser.ts"; import { FSAdapterWrapper } from "#veryfront/platform/adapters/fs/wrapper.ts"; import { MultiProjectFSAdapter } from "#veryfront/platform/adapters/fs/veryfront/multi-project-adapter.ts"; import { runWithoutRequestContext } from "#veryfront/platform/adapters/fs/veryfront/request-context.ts"; import type { SourceSnapshotFreshnessOptions } from "#veryfront/platform/adapters/base.ts"; +import { createRunScopedProviderReplayCheckpointPersister } from "#veryfront/internal-agents/provider-replay-checkpoint-persister.ts"; export interface AgentStreamHandlerDeps extends RuntimeAgentDiscoveryDeps, RuntimeAgentStreamExecutionDeps { @@ -106,6 +110,8 @@ export interface AgentStreamHandlerDeps getLocalTools?: (agentId: string) => RuntimeAgentStreamExecutionDeps["localTools"]; loadAgentSourceEnvironment?: AgentSourceEnvironmentLoader; normalizeSourceIntegrationPolicy?: typeof normalizeSourceIntegrationPolicy; + createRunScopedProviderReplayCheckpointPersister?: + typeof createRunScopedProviderReplayCheckpointPersister; } type AgentSourceTargetIdentity = Pick< @@ -128,6 +134,8 @@ const defaultDeps: AgentStreamHandlerDeps = { loadAgentSourceEnvironment: resolveAgentSourceEnvironment, getLocalTools: (agentId) => getDiscoveredHostTools({ agentId }) as RuntimeAgentStreamExecutionDeps["localTools"], + providerReplayCheckpointEmissionEnabled: isProviderReplayCheckpointEmissionEnabled(), + createRunScopedProviderReplayCheckpointPersister, }; const logger = serverLogger.component("agent-stream-handler"); const IntrinsicReflectApply = Reflect.apply; @@ -1003,6 +1011,7 @@ export class AgentStreamHandler extends BaseHandler { expectedSubject: payload.runId, expectedSurface: "studio", }); + const runEventAppendToken = req.headers.get(RUN_EVENT_APPEND_TOKEN_HEADER); assertAgentSourceMatchesHostedTarget(ctx, payload); const apiAuthToken = payload.credentials?.authToken || ctx.proxyToken || ""; if (payload.agentSource.type === "environment" && !apiAuthToken) { @@ -1151,6 +1160,24 @@ export class AgentStreamHandler extends BaseHandler { toRuntimeRunAgentInput(payload), runtimeBaseAgent as Agent, ); + const providerReplayCheckpoints = + getServerResolvedProviderReplayCheckpoints({ + forwardedProps: runtimeInput.forwardedProps, + ...(runtimeInput.serverResolvedProviderReplayCheckpoints !== undefined + ? { + serverResolvedProviderReplayCheckpoints: + runtimeInput.serverResolvedProviderReplayCheckpoints, + } + : {}), + serverEnvelopeVerified: true, + }); + const veryfrontApiUrl = resolveVeryfrontApiBaseUrlFromHostEnv(); + const persistProviderReplayCheckpoint = this.deps + .createRunScopedProviderReplayCheckpointPersister?.({ + apiUrl: veryfrontApiUrl, + runId: runtimeInput.runId, + runEventAppendToken, + }); const localTools = this.deps.getLocalTools?.(runtimeBaseAgent.id); const platformRuntimeAgent = await withVeryfrontPlatformRemoteTools({ agent: runtimeBaseAgent as Agent, @@ -1180,8 +1207,10 @@ export class AgentStreamHandler extends BaseHandler { createRuntimeAgentStreamResponse(runtimeInput, runtimeAgent, { ...this.deps, localTools, + providerReplayCheckpoints, + persistProviderReplayCheckpoint, projectAgentSandbox: { - apiUrl: resolveVeryfrontApiBaseUrlFromHostEnv(), + apiUrl: veryfrontApiUrl, authToken: projectRuntimeToken || undefined, branchId: payload.runtimeTargetBranchId, projectId: projectScopedContext.projectId ?? null, diff --git a/tests/integration/semantic-unit-boundary/src/internal-agents/provider-replay-checkpoint-persister-intrinsics.test.ts b/tests/integration/semantic-unit-boundary/src/internal-agents/provider-replay-checkpoint-persister-intrinsics.test.ts new file mode 100644 index 0000000000..95258ea8ae --- /dev/null +++ b/tests/integration/semantic-unit-boundary/src/internal-agents/provider-replay-checkpoint-persister-intrinsics.test.ts @@ -0,0 +1,62 @@ +// This security boundary test intentionally mutates shared-realm prototypes, +// so it belongs in the semantic integration suite rather than a unit module. +import { assertEquals } from "#veryfront/testing/assert.ts"; +import { describe, it } from "#veryfront/testing/bdd.ts"; +import type { ProviderReplayCheckpoint } from "#veryfront/agent/runtime/provider-replay.ts"; +import { createRunScopedProviderReplayCheckpointPersister } from "#veryfront/internal-agents/provider-replay-checkpoint-persister.ts"; + +const RUN_ID = "run_checkpoint_1"; +const MESSAGE_ID = "10000000-1000-4000-8000-100000000001"; +const testApply = Reflect.apply; + +function checkpoint(): ProviderReplayCheckpoint { + return { + version: 1, + messageId: MESSAGE_ID, + provider: "anthropic", + providerBlocks: [{ + type: "provider-block", + provider: "anthropic", + block: { type: "thinking", thinking: "", signature: "" }, + }], + providerBlockPositions: [0], + providerMessageBlockCounts: [1], + totalPartCount: 1, + }; +} + +describe("run-scoped provider replay checkpoint intrinsic boundary", () => { + it("keeps credential validation on captured intrinsics after project poisoning", async () => { + const token = "writer-token-that-must-stay-private"; + const nativeTrim = String.prototype.trim; + const nativeStringValueOf = String.prototype.valueOf; + const nativeEncode = TextEncoder.prototype.encode; + let observedTokenCalls = 0; + + try { + String.prototype.trim = function () { + const value = testApply(nativeStringValueOf, this, []) as string; + if (value === token) observedTokenCalls++; + return testApply(nativeTrim, this, []) as string; + }; + TextEncoder.prototype.encode = (function (this: TextEncoder, value = "") { + if (value === token) observedTokenCalls++; + return testApply(nativeEncode, this, [value]) as Uint8Array; + }) as typeof TextEncoder.prototype.encode; + + const persist = createRunScopedProviderReplayCheckpointPersister({ + apiUrl: "https://api.example.test", + runId: RUN_ID, + runEventAppendToken: token, + fetch: () => Promise.resolve(Response.json({ appended_count: 1 })), + }); + if (!persist) throw new Error("Expected a checkpoint persister"); + await persist(checkpoint()); + } finally { + String.prototype.trim = nativeTrim; + TextEncoder.prototype.encode = nativeEncode; + } + + assertEquals(observedTokenCalls, 0); + }); +});