From eba3248d86491b06eec043988b7d133331b0dcac Mon Sep 17 00:00:00 2001 From: Alessandro Pogliaghi Date: Fri, 17 Jul 2026 13:39:22 +0100 Subject: [PATCH 1/2] feat(cloud): add native steering to agent adapters --- .../claude/claude-agent.streamed-text.test.ts | 58 ++++++++ .../agent/src/adapters/claude/claude-agent.ts | 31 ++++- packages/agent/src/adapters/claude/types.ts | 1 + .../app-server-client.test.ts | 9 +- .../codex-app-server/app-server-client.ts | 19 ++- .../codex-app-server-agent.test.ts | 126 +++++++++++++++++- .../codex-app-server-agent.ts | 52 +++++--- .../codex-app-server/turn-controller.ts | 5 - 8 files changed, 273 insertions(+), 28 deletions(-) diff --git a/packages/agent/src/adapters/claude/claude-agent.streamed-text.test.ts b/packages/agent/src/adapters/claude/claude-agent.streamed-text.test.ts index 96775fa040..4bd3d8e13f 100644 --- a/packages/agent/src/adapters/claude/claude-agent.streamed-text.test.ts +++ b/packages/agent/src/adapters/claude/claude-agent.streamed-text.test.ts @@ -243,6 +243,64 @@ describe("ClaudeAcpAgent.prompt — streamed assistant text wiring", () => { ]); }); + it("keeps the original turn open until a pending steer is consumed", async () => { + const { agent, client } = makeAgent(); + const sessionId = "s-steer-ordering"; + const { query, input } = installFakeSession(agent, sessionId); + + const promptPromise = agent.prompt({ + sessionId, + prompt: [{ type: "text", text: "use orange" }], + }); + let promptSettled = false; + void promptPromise.then(() => { + promptSettled = true; + }); + await tick(); + await echoUserMessage(query, input); + + const steerResult = await agent.prompt({ + sessionId, + prompt: [{ type: "text", text: "use green instead" }], + _meta: { steer: true }, + }); + expect(steerResult._meta).toEqual({ steer: true }); + + await send(query, assistantMessage(sessionId, "msg_orange", "ORANGE")); + await send(query, resultSuccess(sessionId)); + expect(promptSettled).toBe(false); + + await echoUserMessage(query, input); + await send(query, assistantMessage(sessionId, "msg_green", "GREEN")); + await send(query, resultSuccess(sessionId)); + + await expect(promptPromise).resolves.toMatchObject({ + stopReason: "end_turn", + }); + expect(messageChunkTexts(client.sessionUpdate.mock.calls)).toEqual([ + "ORANGE", + "GREEN", + ]); + }); + + it("declines an explicit steer after the active turn has ended", async () => { + const { agent } = makeAgent(); + const sessionId = "s-expired-steer"; + installFakeSession(agent, sessionId); + + await expect( + agent.prompt({ + sessionId, + prompt: [{ type: "text", text: "too late" }], + _meta: { steer: true }, + }), + ).resolves.toMatchObject({ _meta: { steer: false } }); + + const session = (agent as unknown as { session: { turnQueue: unknown[] } }) + .session; + expect(session.turnQueue).toHaveLength(0); + }); + it("reconnects a disconnected signed-commit server before the turn", async () => { const { agent } = makeAgent(); const sessionId = "s-heal"; diff --git a/packages/agent/src/adapters/claude/claude-agent.ts b/packages/agent/src/adapters/claude/claude-agent.ts index 402d8dbdd4..983567d8a5 100644 --- a/packages/agent/src/adapters/claude/claude-agent.ts +++ b/packages/agent/src/adapters/claude/claude-agent.ts @@ -482,12 +482,20 @@ export class ClaudeAcpAgent extends BaseAcpAgent { const hasInFlightTurns = this.session.activeTurn !== null || this.session.turnQueue.length > 0; - if (hasInFlightTurns && isSteerMeta(params._meta)) { + const isSteer = isSteerMeta(params._meta); + if (hasInFlightTurns && isSteer) { // Fold into the running turn (promptToClaude tagged it priority:"next"); // the benign end_turn is ignored by clients, which key off _meta.steer. + const owner = + this.session.activeTurn ?? + this.session.turnQueue.find((turn) => !turn.settled); + owner?.pendingSteerUuids.add(promptUuid); this.session.input.push(userMessage); await this.broadcastUserMessage(params); - return { stopReason: "end_turn" }; + return { stopReason: "end_turn", _meta: { steer: true } }; + } + if (isSteer) { + return { stopReason: "end_turn", _meta: { steer: false } }; } if (!hasInFlightTurns && !isLocalOnlyCommand) { @@ -507,6 +515,7 @@ export class ClaudeAcpAgent extends BaseAcpAgent { const turn: Turn = { promptUuid, + pendingSteerUuids: new Set(), isLocalOnlyCommand, commandName: commandMatch?.[1], broadcast: () => this.broadcastUserMessage(params), @@ -1059,6 +1068,21 @@ export class ClaudeAcpAgent extends BaseAcpAgent { }, ); + if ( + !isTaskNotification && + session.activeTurn && + session.activeTurn.pendingSteerUuids.size > 0 + ) { + this.logger.debug( + "Deferring turn completion until pending steers are consumed", + { + sessionId, + pendingSteers: session.activeTurn.pendingSteerUuids.size, + }, + ); + break; + } + if ( (message as { stop_reason?: string }).stop_reason === "refusal" ) { @@ -1198,6 +1222,9 @@ export class ClaudeAcpAgent extends BaseAcpAgent { // active one first), then drops from the feed. Runs before the // cancelled guard so a turn enqueued after a cancel still starts. if (message.type === "user" && "uuid" in message && message.uuid) { + if (session.activeTurn?.pendingSteerUuids.delete(message.uuid)) { + break; + } const queued = session.turnQueue.find( (t) => t.promptUuid === message.uuid && !t.settled, ); diff --git a/packages/agent/src/adapters/claude/types.ts b/packages/agent/src/adapters/claude/types.ts index 80078dfa8f..37d46ac0d6 100644 --- a/packages/agent/src/adapters/claude/types.ts +++ b/packages/agent/src/adapters/claude/types.ts @@ -42,6 +42,7 @@ export type BackgroundTerminal = /** One in-flight `prompt()` call, settled by the session's consumer. */ export type Turn = { promptUuid: string; + pendingSteerUuids: Set; isLocalOnlyCommand: boolean; commandName?: string; /** Invoked once at activation, matching the pre-consumer broadcast timing. */ diff --git a/packages/agent/src/adapters/codex-app-server/app-server-client.test.ts b/packages/agent/src/adapters/codex-app-server/app-server-client.test.ts index 9501f019ce..a46cc96091 100644 --- a/packages/agent/src/adapters/codex-app-server/app-server-client.test.ts +++ b/packages/agent/src/adapters/codex-app-server/app-server-client.test.ts @@ -4,7 +4,7 @@ import { createBidirectionalStreams, type StreamPair, } from "../../utils/streams"; -import { AppServerClient } from "./app-server-client"; +import { AppServerClient, AppServerRequestError } from "./app-server-client"; interface RpcMessage { id?: number | string; @@ -86,7 +86,12 @@ describe("AppServerClient", () => { error: { code: -32001, message: "Server overloaded; retry later." }, }); - await expect(pending).rejects.toThrow("Server overloaded; retry later."); + const error = await pending.catch((requestError: unknown) => requestError); + expect(error).toBeInstanceOf(AppServerRequestError); + expect(error).toMatchObject({ + code: -32001, + message: "Server overloaded; retry later.", + }); await client.close(); }); diff --git a/packages/agent/src/adapters/codex-app-server/app-server-client.ts b/packages/agent/src/adapters/codex-app-server/app-server-client.ts index ff32050f5c..7d911173e2 100644 --- a/packages/agent/src/adapters/codex-app-server/app-server-client.ts +++ b/packages/agent/src/adapters/codex-app-server/app-server-client.ts @@ -24,6 +24,17 @@ export interface AppServerRpc { close(): Promise; } +export class AppServerRequestError extends Error { + constructor( + readonly code: number, + message: string, + readonly data?: unknown, + ) { + super(message); + this.name = "AppServerRequestError"; + } +} + /** * Bidirectional newline-delimited JSON-RPC client for the native Codex `app-server` subprocess. * Transport-agnostic via a {@link StreamPair} so tests can drive it over in-memory streams. @@ -173,7 +184,13 @@ export class AppServerClient implements AppServerRpc { } this.pending.delete(message.id); if (message.error) { - call.reject(new Error(message.error.message)); + call.reject( + new AppServerRequestError( + message.error.code, + message.error.message, + message.error.data, + ), + ); } else { call.resolve(message.result); } diff --git a/packages/agent/src/adapters/codex-app-server/codex-app-server-agent.test.ts b/packages/agent/src/adapters/codex-app-server/codex-app-server-agent.test.ts index 9761794839..7062d98f70 100644 --- a/packages/agent/src/adapters/codex-app-server/codex-app-server-agent.test.ts +++ b/packages/agent/src/adapters/codex-app-server/codex-app-server-agent.test.ts @@ -10,6 +10,7 @@ import type { AppServerClientHandlers, AppServerRpc, } from "./app-server-client"; +import { AppServerRequestError } from "./app-server-client"; import { CodexAppServerAgent } from "./codex-app-server-agent"; import { sandboxPolicyFor } from "./session-config"; @@ -48,6 +49,9 @@ function makeStubRpc(responses: Record) { }; } const response = responses[method]; + if (response instanceof Error) { + throw response; + } return ( typeof response === "function" ? await response(params) @@ -2133,7 +2137,12 @@ describe("CodexAppServerAgent", () => { prompt: [{ type: "text", text: "more context" }], } as unknown as PromptRequest); - // The single turn/completed resolves both the original and the folded prompt. + await expect(second).resolves.toMatchObject({ + stopReason: "end_turn", + _meta: { steer: true }, + }); + + // The original prompt remains the sole owner of turn completion. stub.emit("thread/tokenUsage/updated", { tokenUsage: { last: { @@ -2145,12 +2154,11 @@ describe("CodexAppServerAgent", () => { }, }); stub.emit("turn/completed", { turn: { status: "completed" } }); - const [firstResult, secondResult] = await Promise.all([first, second]); + const firstResult = await first; expect(firstResult).toMatchObject({ stopReason: "end_turn", usage: { totalTokens: 45 }, }); - expect(secondResult).toEqual({ stopReason: "end_turn" }); const steer = stub.requests.find((r) => r.method === "turn/steer"); expect(steer?.params).toMatchObject({ @@ -2164,6 +2172,118 @@ describe("CodexAppServerAgent", () => { ); }); + it("rejects a failed turn/steer without echoing or acknowledging it", async () => { + const stub = makeStubRpc({ + "thread/start": { thread: { id: "t" } }, + "turn/start": { turn: { id: "turn_1" } }, + "turn/steer": new Error("steer transport failed"), + }); + const { client, sessionUpdates } = makeFakeClient(); + const agent = new CodexAppServerAgent(client, { + processOptions: { binaryPath: "/x/codex" }, + rpcFactory: stub.factory, + }); + + await agent.newSession({ cwd: "/r" } as unknown as NewSessionRequest); + const first = agent.prompt({ + sessionId: "t", + prompt: [{ type: "text", text: "one" }], + } as unknown as PromptRequest); + stub.emit("turn/started", { threadId: "t", turn: { id: "turn_1" } }); + + await expect( + agent.prompt({ + sessionId: "t", + prompt: [{ type: "text", text: "lost steer" }], + _meta: { steer: true }, + } as unknown as PromptRequest), + ).rejects.toThrow("steer transport failed"); + expect(sessionUpdates).not.toContainEqual( + expect.objectContaining({ + update: expect.objectContaining({ + sessionUpdate: "user_message_chunk", + content: { type: "text", text: "lost steer" }, + }), + }), + ); + + stub.emit("turn/completed", { turn: { status: "completed" } }); + await first; + }); + + it("declines a stale turn/steer so the caller can queue it normally", async () => { + const stub = makeStubRpc({ + "thread/start": { thread: { id: "t" } }, + "turn/start": { turn: { id: "turn_1" } }, + "turn/steer": new AppServerRequestError( + -32600, + "expected active turn id `turn_1` but found `turn_2`", + ), + }); + const { client, sessionUpdates } = makeFakeClient(); + const agent = new CodexAppServerAgent(client, { + processOptions: { binaryPath: "/x/codex" }, + rpcFactory: stub.factory, + }); + + await agent.newSession({ cwd: "/r" } as unknown as NewSessionRequest); + const first = agent.prompt({ + sessionId: "t", + prompt: [{ type: "text", text: "one" }], + } as unknown as PromptRequest); + stub.emit("turn/started", { threadId: "t", turn: { id: "turn_1" } }); + + await expect( + agent.prompt({ + sessionId: "t", + prompt: [{ type: "text", text: "queue me" }], + _meta: { steer: true }, + } as unknown as PromptRequest), + ).resolves.toMatchObject({ _meta: { steer: false } }); + expect(sessionUpdates).not.toContainEqual( + expect.objectContaining({ + update: expect.objectContaining({ + sessionUpdate: "user_message_chunk", + content: { type: "text", text: "queue me" }, + }), + }), + ); + + stub.emit("turn/completed", { turn: { status: "completed" } }); + await first; + }); + + it("declines an explicit steer after the active turn has ended", async () => { + const stub = makeStubRpc({ + "thread/start": { thread: { id: "t" } }, + }); + const { client, sessionUpdates } = makeFakeClient(); + const agent = new CodexAppServerAgent(client, { + processOptions: { binaryPath: "/x/codex" }, + rpcFactory: stub.factory, + }); + + await agent.newSession({ cwd: "/r" } as unknown as NewSessionRequest); + await expect( + agent.prompt({ + sessionId: "t", + prompt: [{ type: "text", text: "too late" }], + _meta: { steer: true }, + } as unknown as PromptRequest), + ).resolves.toMatchObject({ _meta: { steer: false } }); + expect( + stub.requests.filter((request) => request.method === "turn/start"), + ).toHaveLength(0); + expect(sessionUpdates).not.toContainEqual( + expect.objectContaining({ + update: expect.objectContaining({ + sessionUpdate: "user_message_chunk", + content: { type: "text", text: "too late" }, + }), + }), + ); + }); + it("refreshes the live turnId from each turn/steer response", async () => { const stub = makeStubRpc({ "thread/start": { thread: { id: "t" } }, diff --git a/packages/agent/src/adapters/codex-app-server/codex-app-server-agent.ts b/packages/agent/src/adapters/codex-app-server/codex-app-server-agent.ts index 823a23c6e3..5029f2127a 100644 --- a/packages/agent/src/adapters/codex-app-server/codex-app-server-agent.ts +++ b/packages/agent/src/adapters/codex-app-server/codex-app-server-agent.ts @@ -50,6 +50,7 @@ import { resolveSpokenNarration } from "../session-meta"; import { AppServerClient, type AppServerClientHandlers, + AppServerRequestError, type AppServerRpc, } from "./app-server-client"; import { handleServerRequest } from "./approvals"; @@ -88,6 +89,16 @@ import { parseStructuredOutput } from "./structured-output"; import { TurnController } from "./turn-controller"; import { UsageTracker } from "./usage-tracker"; +function isStaleTurnSteerError(error: unknown): boolean { + if (!(error instanceof AppServerRequestError) || error.code !== -32600) { + return false; + } + return ( + error.message === "no active turn to steer" || + /^expected active turn id `.*` but found `.*`$/.test(error.message) + ); +} + type AppServerSessionMeta = { // The host sends either a plain string or the Claude-style `{ append }` form. systemPrompt?: string | { append?: string }; @@ -753,32 +764,43 @@ export class CodexAppServerAgent extends BaseAcpAgent { if (dropped > 0) { this.logger.warn("Dropped non-text/non-image prompt blocks", { dropped }); } - // Echo the user prompt (codex emits none), for fresh turns and steering alike. - this.broadcastUserInput(params.prompt); - if (this.turns.isRunning) { // A turn is already running: fold the message in via turn/steer (precondition: the // active turnId). Refresh from the response's rotated turnId so a later steer/interrupt // still targets the live turn (no turn/started is re-emitted for a steer). - const steerRes = await this.rpc - .request<{ turnId?: string }>(APP_SERVER_METHODS.TURN_STEER, { - threadId: this.threadId, - input, - expectedTurnId: this.turns.activeTurnId, - }) - .catch((err) => { - this.logger.warn("turn/steer failed", err); - return undefined; - }); + let steerRes: { turnId?: string }; + try { + steerRes = await this.rpc.request<{ turnId?: string }>( + APP_SERVER_METHODS.TURN_STEER, + { + threadId: this.threadId, + input, + expectedTurnId: this.turns.activeTurnId, + }, + ); + } catch (error) { + if ( + (params._meta as { steer?: unknown } | undefined)?.steer === true && + isStaleTurnSteerError(error) + ) { + return { stopReason: "end_turn", _meta: { steer: false } }; + } + throw error; + } this.turns.onSteered(steerRes?.turnId); - const response = await this.turns.awaitCompletion(); - return { stopReason: response.stopReason }; + this.broadcastUserInput(params.prompt); + return { stopReason: "end_turn", _meta: { steer: true } }; + } + if ((params._meta as { steer?: unknown } | undefined)?.steer === true) { + return { stopReason: "end_turn", _meta: { steer: false } }; } if (this.turns.isPending) { // A turn is pending but has no turnId yet, so we can't steer; fail fast. throw new Error("prompt() called while a turn is already in progress"); } + // Codex does not echo user input, so emit it only once delivery can proceed. + this.broadcastUserInput(params.prompt); const response = await this.runTurn(input); return this.maybeOfferPlanImplementation(response); } diff --git a/packages/agent/src/adapters/codex-app-server/turn-controller.ts b/packages/agent/src/adapters/codex-app-server/turn-controller.ts index e6f6483a9a..90e7bcb21b 100644 --- a/packages/agent/src/adapters/codex-app-server/turn-controller.ts +++ b/packages/agent/src/adapters/codex-app-server/turn-controller.ts @@ -48,11 +48,6 @@ export class TurnController { if (typeof id === "string") this.turnId = id; } - /** Await the in-flight turn's completion (the steer path reuses the original). */ - awaitCompletion(): Promise { - return this.completion ?? Promise.resolve({ stopReason: "end_turn" }); - } - /** Atomically claim the pending turn (clears the slot + turnId synchronously), or undefined if already claimed. */ claim(): PendingTurn | undefined { const pending = this.pending; From a04d2b18844afebd8b3a9e431771f325495f4418 Mon Sep 17 00:00:00 2001 From: Alessandro Pogliaghi Date: Fri, 17 Jul 2026 13:39:31 +0100 Subject: [PATCH 2/2] feat(cloud): negotiate idempotent steer delivery --- .../agent/src/server/agent-server.test.ts | 258 +++++++++++ packages/agent/src/server/agent-server.ts | 430 +++++++++++------- packages/agent/src/server/schemas.ts | 2 + .../api-client/src/posthog-client.test.ts | 13 +- packages/api-client/src/posthog-client.ts | 22 +- packages/shared/src/sessions.test.ts | 9 +- packages/shared/src/sessions.ts | 7 +- 7 files changed, 569 insertions(+), 172 deletions(-) diff --git a/packages/agent/src/server/agent-server.test.ts b/packages/agent/src/server/agent-server.test.ts index 087f2b55b7..bce20573d8 100644 --- a/packages/agent/src/server/agent-server.test.ts +++ b/packages/agent/src/server/agent-server.test.ts @@ -2267,6 +2267,106 @@ describe("AgentServer HTTP Mode", () => { expect(prompt).toHaveBeenCalledTimes(4); }, 20000); + it("steers an active turn without emitting a separate turn completion", async () => { + const s = createServer(); + await s.start(); + const prompt = vi.fn(async () => ({ + stopReason: "end_turn", + _meta: { steer: true }, + })); + const broadcastTurnComplete = vi.fn(); + const resetTurnMessages = vi.fn(); + const serverInternals = s as unknown as { + activeOwnedTurnCount: number; + broadcastTurnComplete: typeof broadcastTurnComplete; + session: { + clientConnection: { prompt: typeof prompt }; + logWriter: { resetTurnMessages: typeof resetTurnMessages }; + }; + }; + serverInternals.activeOwnedTurnCount = 1; + serverInternals.broadcastTurnComplete = broadcastTurnComplete; + serverInternals.session.clientConnection.prompt = prompt; + serverInternals.session.logWriter.resetTurnMessages = resetTurnMessages; + + const response = await fetch(`http://localhost:${port}/command`, { + method: "POST", + headers: { + Authorization: `Bearer ${createToken()}`, + "Content-Type": "application/json", + }, + body: JSON.stringify({ + jsonrpc: "2.0", + id: "steer-1", + method: "user_message", + params: { content: "change direction", steer: true }, + }), + }); + + expect(response.status).toBe(200); + await expect(response.json()).resolves.toMatchObject({ + result: { stopReason: "steered", steered: true }, + }); + expect(prompt).toHaveBeenCalledWith( + expect.objectContaining({ + _meta: expect.objectContaining({ steer: true }), + }), + ); + expect(broadcastTurnComplete).not.toHaveBeenCalled(); + expect(resetTurnMessages).not.toHaveBeenCalled(); + }, 20000); + + it("declines steering without blocking on a fallback normal turn", async () => { + const s = createServer(); + await s.start(); + const prompt = vi.fn(); + const broadcastTurnComplete = vi.fn(); + const resetTurnMessages = vi.fn(); + const serverInternals = s as unknown as { + activeOwnedTurnCount: number; + broadcastTurnComplete: typeof broadcastTurnComplete; + session: { + clientConnection: { prompt: typeof prompt }; + logWriter: { resetTurnMessages: typeof resetTurnMessages }; + }; + }; + serverInternals.activeOwnedTurnCount = 1; + prompt.mockImplementationOnce(async () => { + serverInternals.activeOwnedTurnCount = 0; + return { stopReason: "end_turn", _meta: { steer: false } }; + }); + serverInternals.broadcastTurnComplete = broadcastTurnComplete; + serverInternals.session.clientConnection.prompt = prompt; + serverInternals.session.logWriter.resetTurnMessages = resetTurnMessages; + + const response = await fetch(`http://localhost:${port}/command`, { + method: "POST", + headers: { + Authorization: `Bearer ${createToken()}`, + "Content-Type": "application/json", + }, + body: JSON.stringify({ + jsonrpc: "2.0", + id: "steer-race", + method: "user_message", + params: { content: "continue normally", steer: true }, + }), + }); + + expect(response.status).toBe(200); + await expect(response.json()).resolves.toMatchObject({ + result: { stopReason: "steer_declined", steered: false }, + }); + expect(prompt).toHaveBeenCalledTimes(1); + expect(prompt.mock.calls[0]?.[0]).toEqual( + expect.objectContaining({ + _meta: expect.objectContaining({ steer: true }), + }), + ); + expect(resetTurnMessages).not.toHaveBeenCalled(); + expect(broadcastTurnComplete).not.toHaveBeenCalled(); + }, 20000); + it("redelivers a messageId whose first delivery failed before producing a turn", async () => { const s = createServer(); await s.start(); @@ -2307,6 +2407,163 @@ describe("AgentServer HTTP Mode", () => { expect(prompt).toHaveBeenCalledTimes(2); }, 20000); + it("keeps a recoverable delivery committed across an ambiguous retry", async () => { + const s = createServer(); + await s.start(); + const prompt = vi + .fn() + .mockRejectedValue(new Error("API Error: The operation timed out.")); + const serverInternals = s as unknown as { + session: { clientConnection: { prompt: typeof prompt } }; + }; + serverInternals.session.clientConnection.prompt = prompt; + + const token = createToken(); + const send = async (requestId: string) => + fetch(`http://localhost:${port}/command`, { + method: "POST", + headers: { + Authorization: `Bearer ${token}`, + "Content-Type": "application/json", + }, + body: JSON.stringify({ + jsonrpc: "2.0", + id: requestId, + method: "user_message", + params: { + content: "do the thing", + messageId: "m-recoverable", + }, + }), + }); + + const first = await send("first-attempt"); + await expect(first.json()).resolves.toMatchObject({ + result: { stopReason: "error_recoverable" }, + }); + expect(prompt).toHaveBeenCalledTimes(1); + + const retry = await send("ambiguous-retry"); + await expect(retry.json()).resolves.toMatchObject({ + result: { stopReason: "duplicate_delivery", duplicate: true }, + }); + expect(prompt).toHaveBeenCalledTimes(1); + }, 20000); + + it("shares a failed in-flight messageId outcome with concurrent retries", async () => { + const s = createServer(); + await s.start(); + let rejectFirstDelivery!: (error: Error) => void; + const prompt = vi + .fn() + .mockImplementationOnce( + () => + new Promise((_resolve, reject) => { + rejectFirstDelivery = reject; + }), + ) + .mockResolvedValueOnce({ stopReason: "end_turn" }); + const serverInternals = s as unknown as { + logger: { info: (...args: unknown[]) => void }; + session: { clientConnection: { prompt: typeof prompt } }; + }; + serverInternals.session.clientConnection.prompt = prompt; + const infoLog = vi.spyOn(serverInternals.logger, "info"); + + const token = createToken(); + const send = async (requestId: string) => + fetch(`http://localhost:${port}/command`, { + method: "POST", + headers: { + Authorization: `Bearer ${token}`, + "Content-Type": "application/json", + }, + body: JSON.stringify({ + jsonrpc: "2.0", + id: requestId, + method: "user_message", + params: { content: "do the thing", messageId: "m-concurrent" }, + }), + }); + + const firstResponse = send("first-attempt"); + await vi.waitFor(() => expect(prompt).toHaveBeenCalledTimes(1)); + + let retrySettled = false; + const retryResponse = send("concurrent-retry").finally(() => { + retrySettled = true; + }); + await vi.waitFor(() => { + expect(infoLog).toHaveBeenCalledWith( + "Awaiting in-flight user_message delivery", + { messageId: "m-concurrent" }, + ); + expect(prompt).toHaveBeenCalledTimes(1); + expect(retrySettled).toBe(false); + }); + + rejectFirstDelivery(new Error("sdk connection lost")); + const [first, retry] = await Promise.all([firstResponse, retryResponse]); + await expect(first.json()).resolves.toMatchObject({ + error: { message: "sdk connection lost" }, + }); + await expect(retry.json()).resolves.toMatchObject({ + error: { message: "sdk connection lost" }, + }); + expect(prompt).toHaveBeenCalledTimes(1); + }, 20000); + + it("keeps an accepted messageId committed when teardown clears the active session", async () => { + const s = createServer(); + await s.start(); + let finishPrompt!: (result: { stopReason: "end_turn" }) => void; + const prompt = vi.fn( + () => + new Promise<{ stopReason: "end_turn" }>((resolve) => { + finishPrompt = resolve; + }), + ); + const serverInternals = s as unknown as { + session: { clientConnection: { prompt: typeof prompt } } | null; + }; + const acceptedSession = serverInternals.session; + if (!acceptedSession) throw new Error("expected active test session"); + acceptedSession.clientConnection.prompt = prompt; + + const token = createToken(); + const send = async (requestId: string) => + fetch(`http://localhost:${port}/command`, { + method: "POST", + headers: { + Authorization: `Bearer ${token}`, + "Content-Type": "application/json", + }, + body: JSON.stringify({ + jsonrpc: "2.0", + id: requestId, + method: "user_message", + params: { content: "do the thing", messageId: "m-teardown" }, + }), + }); + + const firstResponse = send("first-attempt"); + await vi.waitFor(() => expect(prompt).toHaveBeenCalledTimes(1)); + + serverInternals.session = null; + finishPrompt({ stopReason: "end_turn" }); + const first = await firstResponse; + await expect(first.json()).resolves.toMatchObject({ + result: { stopReason: "end_turn" }, + }); + + serverInternals.session = acceptedSession; + const retry = await send("retry"); + await expect(retry.json()).resolves.toMatchObject({ + result: { stopReason: "duplicate_delivery", duplicate: true }, + }); + expect(prompt).toHaveBeenCalledTimes(1); + }, 20000); + // Shared plumbing for the relay-echo tests: install a controllable // prompt, stub the log writer so relayAgentResponse has an answer to // relay, and spy on the relay_message client call. @@ -2447,6 +2704,7 @@ describe("AgentServer HTTP Mode", () => { expect(runStarted?.notification?.params).toMatchObject({ runId: "test-run-id", taskId: "test-task-id", + steering: "native", }); // Agent reports its semver so clients can gate UI features // against agent capabilities (e.g. `>=0.40.1`). The exact value diff --git a/packages/agent/src/server/agent-server.ts b/packages/agent/src/server/agent-server.ts index 208c643c4b..c25336e948 100644 --- a/packages/agent/src/server/agent-server.ts +++ b/packages/agent/src/server/agent-server.ts @@ -298,6 +298,15 @@ function isManualCompactPrompt(prompt: ContentBlock[]): boolean { return /^\/compact(?:\s|$)/.test(promptBlocksToText(prompt).trimStart()); } +function extractSteeringCapability(result: unknown): string | undefined { + const steering = ( + result as { + agentCapabilities?: { _meta?: { posthog?: { steering?: unknown } } }; + } + )?.agentCapabilities?._meta?.posthog?.steering; + return typeof steering === "string" ? steering : undefined; +} + interface LocalSkillPromptContext { /** Set when the message is a bare `/skill` invocation the adapter should strip. */ skillName?: string; @@ -377,6 +386,8 @@ export class AgentServer { private preSessionEvents: Record[] = []; private deliveredMessageIds = new Set(); private pendingCompactContinuationMessageIds = new Set(); + private inFlightMessageDeliveries = new Map>(); + private activeOwnedTurnCount = 0; private pendingPermissions = new Map< string, { @@ -937,49 +948,51 @@ export class AgentServer { switch (method) { case POSTHOG_NOTIFICATIONS.USER_MESSAGE: case "user_message": { - this.logger.debug("Received user_message command", { - hasContent: - typeof params.content === "string" - ? params.content.trim().length > 0 - : Array.isArray(params.content) && params.content.length > 0, - artifactCount: Array.isArray(params.artifacts) - ? params.artifacts.length - : 0, - }); - const builtPrompt = await this.buildPromptFromContentAndArtifacts({ - content: params.content as string | ContentBlock[] | undefined, - artifacts: Array.isArray(params.artifacts) - ? (params.artifacts as TaskRunArtifact[]) - : [], - taskId: this.session.payload.task_id, - runId: this.session.payload.run_id, - }); - const prompt = builtPrompt.prompt; - if (prompt.length === 0) { - throw new Error("User message cannot be empty"); - } - + const commandSession = this.session; const messageId = typeof params.messageId === "string" && params.messageId ? params.messageId : undefined; + const inFlightDelivery = messageId + ? this.inFlightMessageDeliveries.get(messageId) + : undefined; + if (inFlightDelivery) { + this.logger.info("Awaiting in-flight user_message delivery", { + messageId, + }); + return await inFlightDelivery; + } + let retryCompactContinuation = false; - if (messageId) { - if (this.deliveredMessageIds.has(messageId)) { - if (this.pendingCompactContinuationMessageIds.has(messageId)) { - retryCompactContinuation = true; - this.logger.info("Retrying pending compact continuation", { - messageId, - }); - } else { - this.logger.info("Duplicate user_message delivery ignored", { - messageId, - }); - return { stopReason: "duplicate_delivery", duplicate: true }; - } + if (messageId && this.deliveredMessageIds.has(messageId)) { + if (this.pendingCompactContinuationMessageIds.has(messageId)) { + retryCompactContinuation = true; + this.logger.info("Retrying pending compact continuation", { + messageId, + }); } else { - this.deliveredMessageIds.add(messageId); + this.logger.info("Duplicate user_message delivery ignored", { + messageId, + }); + return { stopReason: "duplicate_delivery", duplicate: true }; } + } + + let resolveDelivery: (result: unknown) => void = () => {}; + let rejectDelivery: (error: unknown) => void = () => {}; + const deliveryOutcome = new Promise((resolve, reject) => { + resolveDelivery = resolve; + rejectDelivery = reject; + }); + void deliveryOutcome.catch(() => {}); + if (messageId) { + this.inFlightMessageDeliveries.set(messageId, deliveryOutcome); + } + let deliveryCommitted = retryCompactContinuation; + const commitDelivery = (): void => { + deliveryCommitted = true; + if (!messageId) return; + this.deliveredMessageIds.add(messageId); if (this.deliveredMessageIds.size > 500) { const oldest = this.deliveredMessageIds.values().next().value; if (oldest !== undefined) { @@ -987,137 +1000,213 @@ export class AgentServer { this.pendingCompactContinuationMessageIds.delete(oldest); } } - } - this.logger.debug("Built user_message prompt", { - blockTypes: prompt.map((block) => block.type), - }); - const promptPreview = promptBlocksToText(prompt); + }; - this.logger.debug( - `Processing user message (detectedPrUrl=${this.detectedPrUrl ?? "none"}): ${promptPreview.substring(0, 100)}...`, - ); + try { + this.logger.debug("Received user_message command", { + hasContent: + typeof params.content === "string" + ? params.content.trim().length > 0 + : Array.isArray(params.content) && params.content.length > 0, + artifactCount: Array.isArray(params.artifacts) + ? params.artifacts.length + : 0, + }); + const builtPrompt = await this.buildPromptFromContentAndArtifacts({ + content: params.content as string | ContentBlock[] | undefined, + artifacts: Array.isArray(params.artifacts) + ? (params.artifacts as TaskRunArtifact[]) + : [], + taskId: commandSession.payload.task_id, + runId: commandSession.payload.run_id, + }); + const prompt = builtPrompt.prompt; + if (prompt.length === 0) { + throw new Error("User message cannot be empty"); + } - this.session.logWriter.resetTurnMessages(this.session.payload.run_id); + this.logger.debug("Built user_message prompt", { + blockTypes: prompt.map((block) => block.type), + }); + const promptPreview = promptBlocksToText(prompt); - // Resolve before buildDetectedPrContext so a warm auto-publish upgrade - // also flips the detected-PR context to its push variant. - const autoPublishUpgrade = await this.resolveWarmAutoPublishUpgrade(); - const hostContext = [ - ...(autoPublishUpgrade ? [autoPublishUpgrade] : []), - ...(this.detectedPrUrl - ? [this.buildDetectedPrContext(this.detectedPrUrl)] - : []), - ]; - const promptMeta: Record = { - ...(builtPrompt.meta ?? {}), - ...(hostContext.length > 0 - ? { prContext: hostContext.join("\n\n") } - : {}), - }; + this.logger.debug( + `Processing user message (detectedPrUrl=${this.detectedPrUrl ?? "none"}): ${promptPreview.substring(0, 100)}...`, + ); - const manualCompactPrompt = isManualCompactPrompt(prompt); - const acpSessionId = this.session.acpSessionId; - const continueAfterCompaction = (): Promise => - this.promptWithUpstreamRetry({ - sessionId: acpSessionId, - prompt: [ - hiddenTextBlock( - "Compaction is complete. Continue working on the task from the compacted context, following the user's instructions from the /compact command.", - ), - ], - }); + // Resolve before buildDetectedPrContext so a warm auto-publish upgrade + // also flips the detected-PR context to its push variant. + const autoPublishUpgrade = await this.resolveWarmAutoPublishUpgrade(); + const hostContext = [ + ...(autoPublishUpgrade ? [autoPublishUpgrade] : []), + ...(this.detectedPrUrl + ? [this.buildDetectedPrContext(this.detectedPrUrl)] + : []), + ]; + const promptMeta: Record = { + ...(builtPrompt.meta ?? {}), + ...(hostContext.length > 0 + ? { prContext: hostContext.join("\n\n") } + : {}), + }; - let compactCommandCompleted = retryCompactContinuation; - let result: PromptResponse; - this.suppressAdapterTurnComplete = - manualCompactPrompt || retryCompactContinuation; - try { - if (retryCompactContinuation) { - result = await continueAfterCompaction(); - if (messageId) { - this.pendingCompactContinuationMessageIds.delete(messageId); + if (params.steer === true) { + if (this.activeOwnedTurnCount > 0) { + const result = await commandSession.clientConnection.prompt({ + sessionId: commandSession.acpSessionId, + prompt, + _meta: { ...promptMeta, steer: true }, + }); + const accepted = + (result._meta as { steer?: unknown } | undefined)?.steer === + true; + if (accepted) { + commitDelivery(); + const outcome = { stopReason: "steered", steered: true }; + resolveDelivery(outcome); + return outcome; + } } - } else { - result = await this.session.clientConnection.prompt({ - sessionId: this.session.acpSessionId, - prompt, - ...(Object.keys(promptMeta).length > 0 - ? { _meta: promptMeta } - : {}), + const outcome = { + stopReason: "steer_declined", + steered: false, + }; + resolveDelivery(outcome); + return outcome; + } + + commandSession.logWriter.resetTurnMessages( + commandSession.payload.run_id, + ); + + const manualCompactPrompt = isManualCompactPrompt(prompt); + const acpSessionId = commandSession.acpSessionId; + const continueAfterCompaction = (): Promise => + this.promptWithUpstreamRetry({ + sessionId: acpSessionId, + prompt: [ + hiddenTextBlock( + "Compaction is complete. Continue working on the task from the compacted context, following the user's instructions from the /compact command.", + ), + ], }); - if (result.stopReason === "end_turn" && manualCompactPrompt) { - compactCommandCompleted = true; - if (messageId) { - this.pendingCompactContinuationMessageIds.add(messageId); - } - // `/compact` is an SDK-local command, so without a follow-up the - // cloud run reports completion before the model resumes the task. - this.recordTurnUsage(result.usage); - result = await continueAfterCompaction(); + let result: PromptResponse; + this.suppressAdapterTurnComplete = + manualCompactPrompt || retryCompactContinuation; + try { + if (retryCompactContinuation) { + result = await this.runOwnedTurn(continueAfterCompaction); if (messageId) { this.pendingCompactContinuationMessageIds.delete(messageId); } + } else { + result = await this.runOwnedTurn(() => { + const promptResult = commandSession.clientConnection.prompt({ + sessionId: commandSession.acpSessionId, + prompt, + ...(Object.keys(promptMeta).length > 0 + ? { _meta: promptMeta } + : {}), + }); + if (!promptResult) { + throw new Error("Agent connection did not accept the prompt"); + } + return promptResult; + }); + + if (result.stopReason === "end_turn" && manualCompactPrompt) { + commitDelivery(); + if (messageId) { + this.pendingCompactContinuationMessageIds.add(messageId); + } + // `/compact` is an SDK-local command, so without a follow-up the + // cloud run reports completion before the model resumes the task. + this.recordTurnUsage(result.usage); + result = await this.runOwnedTurn(continueAfterCompaction); + if (messageId) { + this.pendingCompactContinuationMessageIds.delete(messageId); + } + } } + } catch (error) { + await commandSession.logWriter.flushAll(); + const { recoverable } = await this.handleTurnFailure( + commandSession.payload, + "followup", + error, + ); + if (!recoverable) { + throw error; + } + commitDelivery(); + const outcome = { stopReason: "error_recoverable" }; + resolveDelivery(outcome); + return outcome; + } finally { + this.suppressAdapterTurnComplete = false; } - } catch (error) { - if (messageId && !compactCommandCompleted) { - this.deliveredMessageIds.delete(messageId); - } - await this.session.logWriter.flushAll(); - const { recoverable } = await this.handleTurnFailure( - this.session.payload, - "followup", - error, - ); - if (!recoverable) { - throw error; - } - return { stopReason: "error_recoverable" }; - } finally { - this.suppressAdapterTurnComplete = false; - } + commitDelivery(); - this.logger.debug("User message completed", { - stopReason: result.stopReason, - }); + this.logger.debug("User message completed", { + stopReason: result.stopReason, + }); - if (result.stopReason === "end_turn") { - void this.syncCloudBranchMetadata(this.session.payload); - } + if (result.stopReason === "end_turn") { + void this.syncCloudBranchMetadata(commandSession.payload); + } - this.recordTurnUsage(result.usage); - this.broadcastTurnComplete(result.stopReason); + this.recordTurnUsage(result.usage); + this.broadcastTurnComplete(result.stopReason); - if (result.stopReason === "end_turn") { - // Relay the response to Slack. For follow-ups this is the primary - // delivery path — the HTTP caller only handles reactions. Echo the - // initiating message's id so the backend can attribute the answer. - this.relayAgentResponse(this.session.payload, messageId).catch( - (err) => - this.logger.debug("Failed to relay follow-up response", err), - ); - } + if (result.stopReason === "end_turn") { + // Relay the response to Slack. For follow-ups this is the primary + // delivery path — the HTTP caller only handles reactions. Echo the + // initiating message's id so the backend can attribute the answer. + this.relayAgentResponse(commandSession.payload, messageId).catch( + (err) => + this.logger.debug("Failed to relay follow-up response", err), + ); + } - // Flush logs and include the assistant's response text so callers - // (e.g. Slack follow-up forwarding) can extract it without racing - // against async log persistence to object storage. - let assistantMessage: string | undefined; - try { - await this.session.logWriter.flush(this.session.payload.run_id, { - coalesce: true, - }); - assistantMessage = this.session.logWriter.getFullAgentResponse( - this.session.payload.run_id, - ); - } catch { - this.logger.debug("Failed to extract assistant message from logs"); - } + // Flush logs and include the assistant's response text so callers + // (e.g. Slack follow-up forwarding) can extract it without racing + // against async log persistence to object storage. + let assistantMessage: string | undefined; + try { + await commandSession.logWriter.flush( + commandSession.payload.run_id, + { + coalesce: true, + }, + ); + assistantMessage = commandSession.logWriter.getFullAgentResponse( + commandSession.payload.run_id, + ); + } catch { + this.logger.debug("Failed to extract assistant message from logs"); + } - return { - stopReason: result.stopReason, - ...(assistantMessage && { assistant_message: assistantMessage }), - }; + const outcome = { + stopReason: result.stopReason, + ...(assistantMessage && { assistant_message: assistantMessage }), + }; + resolveDelivery(outcome); + return outcome; + } catch (error) { + if (messageId && !deliveryCommitted) { + this.deliveredMessageIds.delete(messageId); + } + rejectDelivery(error); + throw error; + } finally { + if ( + messageId && + this.inFlightMessageDeliveries.get(messageId) === deliveryOutcome + ) { + this.inFlightMessageDeliveries.delete(messageId); + } + } } case POSTHOG_NOTIFICATIONS.CANCEL: @@ -1467,10 +1556,11 @@ export class AgentServer { clientStream, ); - await clientConnection.initialize({ + const initializeResult = await clientConnection.initialize({ protocolVersion: PROTOCOL_VERSION, clientCapabilities: {}, }); + const steering = extractSteeringCapability(initializeResult); const runState = preTaskRun?.state as Record | undefined; // Preserve native Codex modes for cloud runs so they behave the same as @@ -1624,6 +1714,7 @@ export class AgentServer { runId: payload.run_id, taskId: payload.task_id, agentVersion: this.config.version ?? packageJson.version, + ...(steering ? { steering } : {}), }, }; this.broadcastEvent({ @@ -1686,6 +1777,15 @@ export class AgentServer { return { classification: classifyAgentError(message), message }; } + private async runOwnedTurn(operation: () => Promise): Promise { + this.activeOwnedTurnCount += 1; + try { + return await operation(); + } finally { + this.activeOwnedTurnCount -= 1; + } + } + /** * Send an initial/resume turn prompt, absorbing transient upstream * failures with a bounded number of retries. These turns run unattended @@ -1900,12 +2000,18 @@ export class AgentServer { }); this.session.logWriter.resetTurnMessages(payload.run_id); + const acpSessionId = this.session.acpSessionId; + if (!acpSessionId) { + throw new Error("Agent session is missing its ACP session ID"); + } - const result = await this.promptWithUpstreamRetry({ - sessionId: this.session.acpSessionId, - prompt: initialPrompt, - ...(initialPromptMeta ? { _meta: initialPromptMeta } : {}), - }); + const result = await this.runOwnedTurn(() => + this.promptWithUpstreamRetry({ + sessionId: acpSessionId, + prompt: initialPrompt, + ...(initialPromptMeta ? { _meta: initialPromptMeta } : {}), + }), + ); this.logger.debug("Initial task message completed", { stopReason: result.stopReason, @@ -2112,12 +2218,18 @@ export class AgentServer { const builtPrompt = await buildPrompt(); this.session.logWriter.resetTurnMessages(payload.run_id); + const acpSessionId = this.session.acpSessionId; + if (!acpSessionId) { + throw new Error("Agent session is missing its ACP session ID"); + } - const result = await this.promptWithUpstreamRetry({ - sessionId: this.session.acpSessionId, - prompt: builtPrompt.prompt, - ...(builtPrompt.meta ? { _meta: builtPrompt.meta } : {}), - }); + const result = await this.runOwnedTurn(() => + this.promptWithUpstreamRetry({ + sessionId: acpSessionId, + prompt: builtPrompt.prompt, + ...(builtPrompt.meta ? { _meta: builtPrompt.meta } : {}), + }), + ); this.logger.debug(`${logLabel} completed`, { stopReason: result.stopReason, diff --git a/packages/agent/src/server/schemas.ts b/packages/agent/src/server/schemas.ts index fd3f2789fe..faf79588fd 100644 --- a/packages/agent/src/server/schemas.ts +++ b/packages/agent/src/server/schemas.ts @@ -65,6 +65,8 @@ export const userMessageParamsSchema = z ]) .optional(), artifacts: z.array(z.record(z.string(), z.unknown())).optional(), + messageId: z.string().min(1).optional(), + steer: z.boolean().optional(), }) .refine( (params) => { diff --git a/packages/api-client/src/posthog-client.test.ts b/packages/api-client/src/posthog-client.test.ts index 3e6509ac08..472a8eeb74 100644 --- a/packages/api-client/src/posthog-client.test.ts +++ b/packages/api-client/src/posthog-client.test.ts @@ -1443,7 +1443,7 @@ describe("PostHogAPIClient", () => { }); }); - it("returns the entries collected so far when a later page fails", async () => { + it("marks entries collected before a failed page as incomplete", async () => { const fetch = vi .fn() .mockResolvedValueOnce(page(makeEntries(50, "a"), true)) @@ -1455,11 +1455,14 @@ describe("PostHogAPIClient", () => { }); const client = makeClient(fetch); - const result = await client.getTaskRunSessionLogs("task-1", "run-1", { - limit: 100000, - }); + const result = await client.getTaskRunSessionLogsResult( + "task-1", + "run-1", + { limit: 100000 }, + ); - expect(result).toHaveLength(50); + expect(result).toEqual({ entries: expect.any(Array), complete: false }); + expect(result.entries).toHaveLength(50); expect(fetch).toHaveBeenCalledTimes(2); }); diff --git a/packages/api-client/src/posthog-client.ts b/packages/api-client/src/posthog-client.ts index 5478769308..29c687a4fb 100644 --- a/packages/api-client/src/posthog-client.ts +++ b/packages/api-client/src/posthog-client.ts @@ -131,6 +131,11 @@ export const CLOUD_USAGE_LIMIT_ERROR_MESSAGE = "Cloud usage limit reached"; export const SESSION_LOGS_MAX_PAGE_SIZE = 5000; +export interface TaskRunSessionLogsResult { + entries: StoredLogEntry[]; + complete: boolean; +} + /** Thrown when the backend rejects a cloud run with a 429 usage-limit error. */ export class CloudUsageLimitError extends Error { limitType: UsageLimitType; @@ -3157,6 +3162,15 @@ export class PostHogAPIClient { runId: string, options?: { limit?: number; after?: string }, ): Promise { + return (await this.getTaskRunSessionLogsResult(taskId, runId, options)) + .entries; + } + + async getTaskRunSessionLogsResult( + taskId: string, + runId: string, + options?: { limit?: number; after?: string }, + ): Promise { const maxEntries = options?.limit ?? SESSION_LOGS_MAX_PAGE_SIZE; const entries: StoredLogEntry[] = []; try { @@ -3187,21 +3201,21 @@ export class PostHogAPIClient { log.warn( `Failed to fetch session logs page at offset ${offset}: ${response.status} ${response.statusText}`, ); - break; + return { entries, complete: false }; } const page = (await response.json()) as StoredLogEntry[]; entries.push(...page); const hasMore = response.headers.get("X-Has-More") === "true"; if (!hasMore || page.length === 0) { - break; + return { entries, complete: true }; } offset += page.length; } - return entries; + return { entries, complete: false }; } catch (err) { log.warn("Failed to fetch task run session logs", err); - return entries; + return { entries, complete: false }; } } diff --git a/packages/shared/src/sessions.test.ts b/packages/shared/src/sessions.test.ts index 85d736fbc3..8706578afe 100644 --- a/packages/shared/src/sessions.test.ts +++ b/packages/shared/src/sessions.test.ts @@ -95,15 +95,20 @@ describe("sessionSupportsNativeSteer", () => { { isCloud: false, steering: "interrupt-resend", adapter: "claude" }, false, ], - // Cloud runs queue/resend; they never steer locally regardless of capability. + // Cloud runs steer only when the sandbox explicitly advertises support. [ "cloud claude native", { isCloud: true, steering: "native", adapter: "claude" }, - false, + true, ], [ "cloud codex native", { isCloud: true, steering: "native", adapter: "codex" }, + true, + ], + [ + "cloud without capability", + { isCloud: true, steering: undefined, adapter: "claude" }, false, ], ])("%s", (_label, session, expected) => { diff --git a/packages/shared/src/sessions.ts b/packages/shared/src/sessions.ts index 493bb5b5ac..fd4fca99c3 100644 --- a/packages/shared/src/sessions.ts +++ b/packages/shared/src/sessions.ts @@ -62,6 +62,9 @@ export interface AgentSession { promptStartedAt: number | null; currentPromptId?: number | null; logUrl?: string; + /** Full cloud transcript entry count across the resume chain. */ + cloudTranscriptEntryCount?: number; + /** Leaf-run cursor used to reconcile live cloud log updates. */ processedLineCount?: number; framework?: "claude"; adapter?: Adapter; @@ -222,7 +225,7 @@ export function resolveBypassRevertMode( * Whether a mid-turn message can be folded into the running turn (steered) * rather than interrupt-and-resent. Decided by the adapter's negotiated * `steering` capability: "native" folds (claude, codex app-server); - * "interrupt-resend" (legacy) does not. Cloud runs never steer locally. + * "interrupt-resend" (legacy) does not. * * Fallback: if `steering` is unset (a start path that predates capability * plumbing), Claude is still treated as native — it has always steered — so the @@ -231,7 +234,7 @@ export function resolveBypassRevertMode( export function sessionSupportsNativeSteer( session: Pick, ): boolean { - if (session.isCloud) return false; if (session.steering === "native") return true; + if (session.isCloud) return false; return session.steering == null && session.adapter === "claude"; }