diff --git a/docs/providers.md b/docs/providers.md index 7595240bd..537acd5ed 100644 --- a/docs/providers.md +++ b/docs/providers.md @@ -44,6 +44,8 @@ OpenCode owns user message IDs. Do not pass Paseo-generated IDs to OpenCode prom Every provider adapter owns its canonical user-message timeline rows. When a foreground prompt is accepted, the adapter must emit exactly one `user_message` timeline item for that submitted prompt, using the same message ID it gives to or receives from the provider runtime. Optimistic client messages are UI-only and provider transcript echoes are optional; neither is allowed to be the only source of truth. If the provider later echoes the same submitted user message, dedupe it only within the active turn. Prefer provider-visible message IDs, but ACP runtimes may omit that ID or replace it with a provider-owned one; in that case suppress only echo chunks whose accumulated text is a prefix of the active submitted prompt. Do not perform global transcript text dedupe. +Submitted user-message rows preserve both identities: `messageId` is the provider-visible ID and the optional `clientMessageId` is the Paseo ID from `AgentRunOptions`. Attach `clientMessageId` only to the canonical row for that foreground submission; provider history and externally initiated user rows do not have a Paseo client ID. + Draft metadata lookups should avoid creating provider sessions when the upstream provider has top-level APIs for that metadata. Prefer `AgentClient.fetchCatalog`, `listCommands`, or `listFeatures` over creating a scratch `AgentSession`; scratch sessions can show up as empty native sessions in provider import/history UIs. `fetchCatalog` is the single discovery API for models and modes — provider implementations may use one process, separate upstream calls, or static data internally, but callers outside the provider do not get separate runtime model/mode probes. Draft feature and command listing must use the explicit draft model only; if no model is selected yet, return no metadata instead of resolving a default model through catalog discovery. Provider session import has its own contract. The picker calls `listImportableSessions` and receives rows only: provider handle, cwd, title, prompt previews, and last activity. Import calls `importSession({ providerHandleId, cwd })` for the selected row and must not call listing again. The provider returns the resumed session, storage config, persistence handle, and hydrated timeline for that one native session; `AgentManager.importProviderSession` seeds the daemon timeline and publishes the Paseo agent only after it is ready. diff --git a/docs/timeline-sync.md b/docs/timeline-sync.md index e90abf412..e96250f0d 100644 --- a/docs/timeline-sync.md +++ b/docs/timeline-sync.md @@ -83,6 +83,10 @@ the existing optimistic-message rules, then restores any unmatched prompts after history. This keeps late history before a newly submitted prompt without duplicating an acknowledged prompt. +Canonical submitted user rows carry the provider's `messageId` and Paseo's optional +`clientMessageId`. Clients reconcile optimistic prompts by `clientMessageId`. Content matching is +limited to the dated compatibility path for daemon timelines created before that field existed. + ## Relevant code - Server live stream forwarding: `packages/server/src/server/session.ts` diff --git a/packages/app/src/timeline/session-stream-reducers.test.ts b/packages/app/src/timeline/session-stream-reducers.test.ts index 05048e2d5..1b84d2e15 100644 --- a/packages/app/src/timeline/session-stream-reducers.test.ts +++ b/packages/app/src/timeline/session-stream-reducers.test.ts @@ -538,7 +538,7 @@ describe("processTimelineResponse", () => { expect(result.error).toBe(null); }); - it("reconciles an optimistic user message during an after-page response", () => { + it("reconciles a legacy canonical user message by content during an after-page response", () => { const existingCursor: TimelineCursor = { epoch: "epoch-1", startSeq: 1, @@ -559,7 +559,7 @@ describe("processTimelineResponse", () => { item: { type: "user_message", text: "sent while catching up", - messageId: "optimistic-after", + messageId: "canonical-after", }, }, ], @@ -568,7 +568,42 @@ describe("processTimelineResponse", () => { const userMessages = result.tail.filter((item) => item.kind === "user_message"); expect(userMessages).toHaveLength(1); - expect(userMessages[0]?.id).toBe("optimistic-after"); + expect(userMessages[0]?.id).toBe("canonical-after"); + expect(userMessages[0]?.optimistic).toBeUndefined(); + }); + + it("reconciles an optimistic user message by client message id", () => { + const optimistic = makeOptimisticUserMessage("local presentation", "client-message"); + + const result = processTimelineResponse({ + ...baseTimelineInput, + currentTail: [optimistic], + currentCursor: { epoch: "epoch-1", startSeq: 1, endSeq: 1 }, + payload: { + ...baseTimelineInput.payload, + epoch: "epoch-1", + entries: [ + { + ...makeTimelineEntry(2, "provider presentation", "user_message"), + item: { + type: "user_message", + text: "provider presentation", + messageId: "provider-message", + clientMessageId: "client-message", + }, + }, + ], + }, + }); + + const userMessages = result.tail.filter((item) => item.kind === "user_message"); + expect(userMessages).toEqual([ + expect.objectContaining({ + id: "provider-message", + clientMessageId: "client-message", + text: "local presentation", + }), + ]); expect(userMessages[0]?.optimistic).toBeUndefined(); }); @@ -1178,7 +1213,7 @@ describe("processTimelineResponse", () => { ]); }); - it("does not match equal prompt text when canonical message ids differ", () => { + it("does not match equal prompt text when canonical client message ids differ", () => { const prompt = makeOptimisticUserMessage("continue", "local-prompt"); const result = processTimelineResponse({ @@ -1197,6 +1232,7 @@ describe("processTimelineResponse", () => { type: "user_message", text: "continue", messageId: "remote-prompt", + clientMessageId: "remote-client-prompt", }, }, ], diff --git a/packages/app/src/timeline/session-stream-reducers.ts b/packages/app/src/timeline/session-stream-reducers.ts index 854832762..1ae2708f6 100644 --- a/packages/app/src/timeline/session-stream-reducers.ts +++ b/packages/app/src/timeline/session-stream-reducers.ts @@ -312,6 +312,7 @@ function mergeCanonicalUserWithLocalPresentation( return { kind: "user_message", id: canonical.id, + ...(canonical.clientMessageId ? { clientMessageId: canonical.clientMessageId } : {}), text: local.text, timestamp: local.timestamp, ...(local.images && local.images.length > 0 ? { images: local.images } : {}), @@ -321,6 +322,27 @@ function mergeCanonicalUserWithLocalPresentation( }; } +interface CanonicalUserMessageIdentity { + messageId?: string; + clientMessageId?: string; + text: string; +} + +function matchesOptimisticUserMessageIdentity( + canonical: CanonicalUserMessageIdentity, + optimistic: UserMessageItem, +): boolean { + if (canonical.clientMessageId !== undefined) { + return canonical.clientMessageId === optimistic.id; + } + if (canonical.messageId === optimistic.id) { + return true; + } + // COMPAT(userMessageClientId): added in v0.2.0, remove after 2027-01-20 once + // the supported daemon floor emits clientMessageId on submitted user messages. + return canonical.text.length > 0 && canonical.text === optimistic.text; +} + function reconcileLocalUserPresentationAfterReplace(params: { canonicalTail: StreamItem[]; previousTail: StreamItem[]; @@ -352,6 +374,21 @@ function reconcileLocalUserPresentationAfterReplace(params: { return false; } if (local.item.optimistic) { + const canonical = params.canonicalTail[index]; + if (canonical?.kind !== "user_message") { + return false; + } + const identityMatches = matchesOptimisticUserMessageIdentity( + { + messageId: canonical.id, + clientMessageId: canonical.clientMessageId, + text: canonical.text, + }, + local.item, + ); + if (canonical.clientMessageId !== undefined || identityMatches) { + return identityMatches; + } return ordinal >= local.ordinal; } return params.canonicalTail[index]?.id === local.item.id; @@ -773,10 +810,7 @@ function matchesOptimisticUserMessage(params: { if (event.type !== "timeline" || event.item.type !== "user_message") { return false; } - if (event.item.messageId !== undefined) { - return event.item.messageId === params.optimistic.id; - } - return event.item.text.length > 0 && event.item.text === params.optimistic.text; + return matchesOptimisticUserMessageIdentity(event.item, params.optimistic); } function acknowledgeOptimisticUserMessage(params: { diff --git a/packages/app/src/types/stream.ts b/packages/app/src/types/stream.ts index 13728a70e..21d2726d3 100644 --- a/packages/app/src/types/stream.ts +++ b/packages/app/src/types/stream.ts @@ -86,6 +86,7 @@ export type UserMessageImageAttachment = AttachmentMetadata; export interface UserMessageItem { kind: "user_message"; id: string; + clientMessageId?: string; text: string; timestamp: Date; optimistic?: true; @@ -238,6 +239,7 @@ function markThoughtReady(item: ThoughtItem): ThoughtItem { function buildUserMessageItem(input: { id: string; + clientMessageId?: string; text: string; timestamp: Date; optimistic?: UserMessageItem | null; @@ -246,6 +248,7 @@ function buildUserMessageItem(input: { return { kind: "user_message", id: input.id, + ...(input.clientMessageId ? { clientMessageId: input.clientMessageId } : {}), text: input.optimistic.text, timestamp: input.optimistic.timestamp, ...(input.optimistic.images && input.optimistic.images.length > 0 @@ -260,6 +263,7 @@ function buildUserMessageItem(input: { return { kind: "user_message", id: input.id, + ...(input.clientMessageId ? { clientMessageId: input.clientMessageId } : {}), text: input.text, timestamp: input.timestamp, }; @@ -351,6 +355,7 @@ function appendUserMessage( text: string, timestamp: Date, messageId?: string, + clientMessageId?: string, ): StreamItem[] { const { chunk, hasContent } = normalizeChunk(text); if (!hasContent) { @@ -366,6 +371,7 @@ function appendUserMessage( const nextItem = buildUserMessageItem({ id: entryId, + clientMessageId, text: chunk, timestamp, optimistic, @@ -833,7 +839,9 @@ function reduceTimelineEvent( const item = event.item; switch (item.type) { case "user_message": - return finalizeActiveThoughts(appendUserMessage(state, item.text, timestamp, item.messageId)); + return finalizeActiveThoughts( + appendUserMessage(state, item.text, timestamp, item.messageId, item.clientMessageId), + ); case "assistant_message": return finalizeActiveThoughts( appendAssistantMessage( diff --git a/packages/protocol/src/agent-types.ts b/packages/protocol/src/agent-types.ts index 54b5275cf..756ea3fd1 100644 --- a/packages/protocol/src/agent-types.ts +++ b/packages/protocol/src/agent-types.ts @@ -338,7 +338,7 @@ export interface CompactionTimelineItem { } export type AgentTimelineItem = - | { type: "user_message"; text: string; messageId?: string } + | { type: "user_message"; text: string; messageId?: string; clientMessageId?: string } | { type: "assistant_message"; text: string; messageId?: string } | { type: "reasoning"; text: string } | ToolCallTimelineItem diff --git a/packages/protocol/src/messages.ts b/packages/protocol/src/messages.ts index a566f8e5c..f4b419864 100644 --- a/packages/protocol/src/messages.ts +++ b/packages/protocol/src/messages.ts @@ -570,6 +570,7 @@ export const AgentTimelineItemPayloadSchema: z.ZodType row.item)).toContainEqual({ type: "user_message", text: "hello from composer", - messageId: "msg-client-1", + messageId: "provider-message-1", + clientMessageId: "msg-client-1", }); } finally { await manager.flush().catch(() => undefined); diff --git a/packages/server/src/server/agent/agent-prompt.test.ts b/packages/server/src/server/agent/agent-prompt.test.ts index f42061698..ba3c64f5e 100644 --- a/packages/server/src/server/agent/agent-prompt.test.ts +++ b/packages/server/src/server/agent/agent-prompt.test.ts @@ -190,7 +190,7 @@ test("sendPromptToAgent forwards the client message id as run options", async () expect(streamAgentSpy).toHaveBeenCalledWith("agent-1", "hello", { outputSchema: { type: "object" }, - messageId: "msg-client-1", + clientMessageId: "msg-client-1", }); }); diff --git a/packages/server/src/server/agent/agent-prompt.ts b/packages/server/src/server/agent/agent-prompt.ts index 96b4c9c29..22bd8cdc2 100644 --- a/packages/server/src/server/agent/agent-prompt.ts +++ b/packages/server/src/server/agent/agent-prompt.ts @@ -198,7 +198,7 @@ export async function sendPromptToAgent( } const runOptions = params.messageId - ? { ...params.runOptions, messageId: params.messageId } + ? { ...params.runOptions, clientMessageId: params.messageId } : params.runOptions; return await startAgentRun(params.agentManager, params.agentId, params.prompt, params.logger, { diff --git a/packages/server/src/server/agent/agent-sdk-types.ts b/packages/server/src/server/agent/agent-sdk-types.ts index c5b8e37b1..6836cd8bc 100644 --- a/packages/server/src/server/agent/agent-sdk-types.ts +++ b/packages/server/src/server/agent/agent-sdk-types.ts @@ -199,7 +199,7 @@ export interface AgentRunOptions { outputSchema?: unknown; resumeFrom?: AgentPersistenceHandle; maxThinkingTokens?: number; - messageId?: string; + clientMessageId?: string; } export interface AgentUsage { @@ -368,7 +368,7 @@ export interface CompactionTimelineItem { } export type AgentTimelineItem = - | { type: "user_message"; text: string; messageId?: string } + | { type: "user_message"; text: string; messageId?: string; clientMessageId?: string } | { type: "assistant_message"; text: string; messageId?: string } | { type: "reasoning"; text: string } | ToolCallTimelineItem diff --git a/packages/server/src/server/agent/create-agent/create.test.ts b/packages/server/src/server/agent/create-agent/create.test.ts index e5c90efc2..89cf7bab6 100644 --- a/packages/server/src/server/agent/create-agent/create.test.ts +++ b/packages/server/src/server/agent/create-agent/create.test.ts @@ -75,7 +75,7 @@ test("session create forwards clientMessageId to the initial prompt run options" }); expect(streamAgent).toHaveBeenCalledWith("agent-1", "hello from create", { - messageId: "msg-create-1", + clientMessageId: "msg-create-1", }); }); diff --git a/packages/server/src/server/agent/create-agent/create.ts b/packages/server/src/server/agent/create-agent/create.ts index 182edc833..b31a858f3 100644 --- a/packages/server/src/server/agent/create-agent/create.ts +++ b/packages/server/src/server/agent/create-agent/create.ts @@ -271,7 +271,7 @@ async function resolveSessionCreateAgent( input.outputSchema || clientMessageId ? { ...(input.outputSchema ? { outputSchema: input.outputSchema } : {}), - ...(clientMessageId ? { messageId: clientMessageId } : {}), + ...(clientMessageId ? { clientMessageId } : {}), } : undefined; const workspaceId = setupContinuation ? createdWorkspaceId : input.workspaceId; diff --git a/packages/server/src/server/agent/providers/acp-agent.test.ts b/packages/server/src/server/agent/providers/acp-agent.test.ts index d95ebbd65..201d4e0e4 100644 --- a/packages/server/src/server/agent/providers/acp-agent.test.ts +++ b/packages/server/src/server/agent/providers/acp-agent.test.ts @@ -2329,7 +2329,9 @@ describe("ACPAgentSession", () => { events.push(event); }); - const { turnId } = await session.startTurn("hello", { messageId: "msg-client-1" }); + const { turnId } = await session.startTurn("hello", { + clientMessageId: "msg-client-1", + }); expect(prompt).toHaveBeenCalledWith({ sessionId: "session-1", @@ -2343,7 +2345,12 @@ describe("ACPAgentSession", () => { type: "timeline", provider: "claude-acp", turnId, - item: { type: "user_message", text: "hello", messageId: "msg-client-1" }, + item: { + type: "user_message", + text: "hello", + messageId: "msg-client-1", + clientMessageId: "msg-client-1", + }, }, ]); @@ -2362,7 +2369,7 @@ describe("ACPAgentSession", () => { events.push(event); }); - await session.startTurn("hello", { messageId: "msg-client-1" }); + await session.startTurn("hello", { clientMessageId: "msg-client-1" }); await session.sessionUpdate({ sessionId: "session-1", update: { @@ -2389,7 +2396,7 @@ describe("ACPAgentSession", () => { events.push(event); }); - await session.startTurn("hello", { messageId: "msg-client-1" }); + await session.startTurn("hello", { clientMessageId: "msg-client-1" }); await session.sessionUpdate({ sessionId: "session-1", update: { @@ -2404,7 +2411,12 @@ describe("ACPAgentSession", () => { { type: "timeline", provider: "claude-acp", - item: { type: "user_message", text: "hello", messageId: "msg-client-1" }, + item: { + type: "user_message", + text: "hello", + messageId: "msg-client-1", + clientMessageId: "msg-client-1", + }, turnId: expect.any(String), }, ]); @@ -2428,7 +2440,7 @@ describe("ACPAgentSession", () => { events.push(event); }); - await session.startTurn("first", { messageId: "msg-client-1" }); + await session.startTurn("first", { clientMessageId: "msg-client-1" }); await session.sessionUpdate({ sessionId: "session-1", update: { @@ -2440,7 +2452,7 @@ describe("ACPAgentSession", () => { await Promise.resolve(); await Promise.resolve(); - await session.startTurn("second", { messageId: "msg-client-2" }); + await session.startTurn("second", { clientMessageId: "msg-client-2" }); await session.sessionUpdate({ sessionId: "session-1", update: { @@ -2455,13 +2467,23 @@ describe("ACPAgentSession", () => { { type: "timeline", provider: "claude-acp", - item: { type: "user_message", text: "first", messageId: "msg-client-1" }, + item: { + type: "user_message", + text: "first", + messageId: "msg-client-1", + clientMessageId: "msg-client-1", + }, turnId: expect.any(String), }, { type: "timeline", provider: "claude-acp", - item: { type: "user_message", text: "second", messageId: "msg-client-2" }, + item: { + type: "user_message", + text: "second", + messageId: "msg-client-2", + clientMessageId: "msg-client-2", + }, turnId: expect.any(String), }, ]); @@ -2479,7 +2501,7 @@ describe("ACPAgentSession", () => { events.push(event); }); - await session.startTurn("hello", { messageId: "msg-client-1" }); + await session.startTurn("hello", { clientMessageId: "msg-client-1" }); await session.sessionUpdate({ sessionId: "session-1", update: { @@ -2495,7 +2517,12 @@ describe("ACPAgentSession", () => { { type: "timeline", provider: "claude-acp", - item: { type: "user_message", text: "hello", messageId: "msg-client-1" }, + item: { + type: "user_message", + text: "hello", + messageId: "msg-client-1", + clientMessageId: "msg-client-1", + }, turnId: expect.any(String), }, ]); diff --git a/packages/server/src/server/agent/providers/acp-agent.ts b/packages/server/src/server/agent/providers/acp-agent.ts index 6da37dad6..ddb28b01d 100644 --- a/packages/server/src/server/agent/providers/acp-agent.ts +++ b/packages/server/src/server/agent/providers/acp-agent.ts @@ -1474,13 +1474,13 @@ export class ACPAgentSession implements AgentSession, ACPClient { } const turnId = randomUUID(); - const messageId = options?.messageId ?? randomUUID(); + const messageId = options?.clientMessageId ?? randomUUID(); this.activeForegroundTurnId = turnId; this.fallbackAssistantMessageId = null; this.activeSubmittedUserMessage = null; this.emitBootstrapThreadEvent(); this.pushEvent({ type: "turn_started", provider: this.provider, turnId }); - this.emitSubmittedUserMessage(prompt, messageId, turnId); + this.emitSubmittedUserMessage(prompt, messageId, turnId, options?.clientMessageId); void this.connection .prompt({ @@ -2687,6 +2687,7 @@ export class ACPAgentSession implements AgentSession, ACPClient { prompt: AgentPromptInput, messageId: string, turnId: string, + clientMessageId?: string, ): void { const text = extractPromptText(prompt); if (text.trim().length === 0) { @@ -2698,7 +2699,12 @@ export class ACPAgentSession implements AgentSession, ACPClient { type: "timeline", provider: this.provider, turnId, - item: { type: "user_message", text, messageId }, + item: { + type: "user_message", + text, + messageId, + ...(clientMessageId ? { clientMessageId } : {}), + }, }); } diff --git a/packages/server/src/server/agent/providers/claude/agent.ts b/packages/server/src/server/agent/providers/claude/agent.ts index 4a07fd4fe..6497a9e96 100644 --- a/packages/server/src/server/agent/providers/claude/agent.ts +++ b/packages/server/src/server/agent/providers/claude/agent.ts @@ -2043,7 +2043,7 @@ class ClaudeAgentSession implements AgentSession { async startTurn( prompt: AgentPromptInput, - _options?: AgentRunOptions, + options?: AgentRunOptions, ): Promise<{ turnId: string }> { if (this.closed) { throw new Error("Claude session is closed"); @@ -2109,7 +2109,7 @@ class ClaudeAgentSession implements AgentSession { this.input.push(sdkMessage); setTimeout(() => { if (this.activeForegroundTurnId === turnId) { - this.emitSubmittedUserMessage(sdkMessage, turnId); + this.emitSubmittedUserMessage(sdkMessage, turnId, options?.clientMessageId); } }, 0); } catch (error) { @@ -3686,6 +3686,7 @@ class ClaudeAgentSession implements AgentSession { private emitSubmittedUserMessage( message: Extract, turnId: string, + clientMessageId?: string, ): void { const events: AgentStreamEvent[] = []; this.appendUserMessageEvents(message, events); @@ -3695,7 +3696,11 @@ class ClaudeAgentSession implements AgentSession { this.foregroundHasVisibleActivity = true; for (const event of events) { if (event.type === "timeline") { - this.notifySubscribers({ ...event, turnId }); + const item = + event.item.type === "user_message" && clientMessageId + ? { ...event.item, clientMessageId } + : event.item; + this.notifySubscribers({ ...event, item, turnId }); } else { this.notifySubscribers(event); } diff --git a/packages/server/src/server/agent/providers/codex-app-server-agent.test.ts b/packages/server/src/server/agent/providers/codex-app-server-agent.test.ts index 4c05495fa..a9df432f4 100644 --- a/packages/server/src/server/agent/providers/codex-app-server-agent.test.ts +++ b/packages/server/src/server/agent/providers/codex-app-server-agent.test.ts @@ -1014,6 +1014,30 @@ describe("Codex app-server provider", () => { await session.close(); }); + test("correlates a Codex user message with the submitting client message", async () => { + const appServer = createFakeCodexAppServer(); + const session = new CodexAppServerAgentSession( + createConfig({ cwd: "/workspace/project" }), + null, + createTestLogger(), + async () => appServer.child, + ); + + await session.startTurn("remember this", { clientMessageId: "client-message" }); + const userMessage = waitForNextTimelineItem(session, "user_message"); + emitCodexUserMessage(appServer, { id: "codex-message", text: "remember this" }); + + await expect(userMessage).resolves.toMatchObject({ + item: { + type: "user_message", + messageId: "codex-message", + clientMessageId: "client-message", + }, + }); + appServer.completeTurn(); + await session.close(); + }); + test("configures Codex app-server to use a custom provider base URL", async () => { const capturedRequests = await runCustomCodexProviderTurn( "codex-iisb", diff --git a/packages/server/src/server/agent/providers/codex-app-server-agent.ts b/packages/server/src/server/agent/providers/codex-app-server-agent.ts index 165e1a30d..ac3017dac 100644 --- a/packages/server/src/server/agent/providers/codex-app-server-agent.ts +++ b/packages/server/src/server/agent/providers/codex-app-server-agent.ts @@ -3096,6 +3096,7 @@ export class CodexAppServerAgentSession implements AgentSession { private readonly subscribers = new Set<(event: AgentStreamEvent) => void>(); private nextTurnOrdinal = 0; private activeForegroundTurnId: string | null = null; + private activeClientMessageId: string | null = null; private cachedRuntimeInfo: AgentRuntimeInfo | null = null; private serviceTier: "fast" | null = null; private planModeEnabled = false; @@ -3820,6 +3821,7 @@ export class CodexAppServerAgentSession implements AgentSession { const turnId = this.createTurnId(); this.activeForegroundTurnId = turnId; + this.activeClientMessageId = options?.clientMessageId ?? null; this.currentTurnId = null; try { @@ -3835,6 +3837,7 @@ export class CodexAppServerAgentSession implements AgentSession { await this.client.request("turn/start", turnStart.params, TURN_START_TIMEOUT_MS); } catch (error) { this.activeForegroundTurnId = null; + this.activeClientMessageId = null; throw error; } @@ -4238,6 +4241,7 @@ export class CodexAppServerAgentSession implements AgentSession { this.pendingSubAgentNotificationsByThreadId.clear(); this.subscribers.clear(); this.activeForegroundTurnId = null; + this.activeClientMessageId = null; if (this.client) { await this.client.dispose(); } @@ -5287,6 +5291,7 @@ export class CodexAppServerAgentSession implements AgentSession { }); } this.activeForegroundTurnId = null; + this.activeClientMessageId = null; this.pendingSubAgentNotificationsByThreadId.clear(); this.resetTurnTrackingState(); } @@ -5825,7 +5830,11 @@ export class CodexAppServerAgentSession implements AgentSession { if (!this.rememberCodexUserMessageTurn(timelineItem.messageId)) { return; } - this.emitEvent({ type: "timeline", provider: CODEX_PROVIDER, item: timelineItem }); + const item = this.activeClientMessageId + ? { ...timelineItem, clientMessageId: this.activeClientMessageId } + : timelineItem; + this.activeClientMessageId = null; + this.emitEvent({ type: "timeline", provider: CODEX_PROVIDER, item }); } private warnUnknownNotificationMethod(method: string, params: unknown): void { diff --git a/packages/server/src/server/agent/providers/mock-load-test-agent.ts b/packages/server/src/server/agent/providers/mock-load-test-agent.ts index d371ef1ac..33b9ddfac 100644 --- a/packages/server/src/server/agent/providers/mock-load-test-agent.ts +++ b/packages/server/src/server/agent/providers/mock-load-test-agent.ts @@ -606,7 +606,7 @@ export class MockLoadTestAgentSession implements AgentSession { async startTurn( prompt: AgentPromptInput, - _options?: AgentRunOptions, + options?: AgentRunOptions, ): Promise<{ turnId: string }> { if (this.activeTurn) { throw new Error("Mock load-test provider already has an active turn"); @@ -648,6 +648,7 @@ export class MockLoadTestAgentSession implements AgentSession { type: "user_message", text: promptToText(prompt), messageId: userMessageId, + ...(options?.clientMessageId ? { clientMessageId: options.clientMessageId } : {}), }, }); }, 0); diff --git a/packages/server/src/server/agent/providers/omp/agent.ts b/packages/server/src/server/agent/providers/omp/agent.ts index 1e002c9f4..28ca8784d 100644 --- a/packages/server/src/server/agent/providers/omp/agent.ts +++ b/packages/server/src/server/agent/providers/omp/agent.ts @@ -896,6 +896,7 @@ export class OmpAgentSession implements AgentSession { private activeAskUserDialog: ActiveAskUserDialog | null = null; private pendingCombinedAskUserResponse: PendingCombinedAskUserResponse | null = null; private activeTurnId: string | null = null; + private activeClientMessageId: string | null = null; private activeAssistantMessageId: string | null = null; private activeTurnStarted = false; private activeTurnHasUserMessage = false; @@ -974,7 +975,7 @@ export class OmpAgentSession implements AgentSession { }); } - async startTurn(prompt: AgentPromptInput, _options?: AgentRunOptions): Promise { + async startTurn(prompt: AgentPromptInput, options?: AgentRunOptions): Promise { if (this.activeTurnId) { throw new Error("An OMP turn is already active"); } @@ -983,6 +984,7 @@ export class OmpAgentSession implements AgentSession { const turnId = randomUUID(); this.live = true; this.activeTurnId = turnId; + this.activeClientMessageId = options?.clientMessageId ?? null; this.activeAssistantMessageId = null; this.activeTurnStarted = false; this.activeTurnHasUserMessage = false; @@ -1010,6 +1012,7 @@ export class OmpAgentSession implements AgentSession { return; } this.activeTurnId = null; + this.activeClientMessageId = null; this.activeTurnStarted = false; this.activeTurnHasUserMessage = false; this.activeAssistantMessageId = null; @@ -1143,6 +1146,7 @@ export class OmpAgentSession implements AgentSession { } if (turnId && this.activeTurnId === turnId) { this.activeTurnId = null; + this.activeClientMessageId = null; this.activeTurnStarted = false; this.activeTurnHasUserMessage = false; this.activeAssistantMessageId = null; @@ -1329,6 +1333,7 @@ export class OmpAgentSession implements AgentSession { item: { type: "user_message", text: promptText, + ...(this.activeClientMessageId ? { clientMessageId: this.activeClientMessageId } : {}), }, }); } @@ -1764,6 +1769,7 @@ export class OmpAgentSession implements AgentSession { } const turnId = this.activeTurnId; this.activeTurnId = null; + this.activeClientMessageId = null; this.activeTurnStarted = false; this.activeTurnHasUserMessage = false; this.clearNoTurnBuffers(); @@ -2003,6 +2009,7 @@ export class OmpAgentSession implements AgentSession { } const nativeMessage = event.message as OmpAgentMessage & { id?: unknown; entryId?: unknown }; const messageId = readNativeMessageId(nativeMessage); + const clientMessageId = this.activeClientMessageId; const emitUserMessage = (resolvedMessageId?: string): void => { if (resolvedMessageId) { // OMP re-emits user message_end frames for entries it has already @@ -2021,6 +2028,7 @@ export class OmpAgentSession implements AgentSession { type: "user_message", text, ...(resolvedMessageId ? { messageId: resolvedMessageId } : {}), + ...(clientMessageId ? { clientMessageId } : {}), }, }); }; @@ -2085,6 +2093,7 @@ export class OmpAgentSession implements AgentSession { private completeTurn(turnId: string | undefined, messages: OmpAgentMessage[]): void { this.activeTurnId = null; + this.activeClientMessageId = null; this.activeAssistantMessageId = null; this.activeTurnStarted = false; this.activeTurnHasUserMessage = false; diff --git a/packages/server/src/server/agent/providers/opencode-agent.ts b/packages/server/src/server/agent/providers/opencode-agent.ts index aabd1b627..0eeca0e7f 100644 --- a/packages/server/src/server/agent/providers/opencode-agent.ts +++ b/packages/server/src/server/agent/providers/opencode-agent.ts @@ -1707,6 +1707,7 @@ export interface OpenCodeEventTranslationState { cwd?: string; messageRoles: Map; pendingUserMessageText?: string | null; + pendingClientMessageId?: string | null; emittedUserMessageIds?: Set; accumulatedUsage: AgentUsage; sessionTotalCostUsd?: number; @@ -2370,7 +2371,12 @@ function appendOpenCodeUserMessageUpdated( events.push({ type: "timeline", provider: "opencode", - item: { type: "user_message", text, messageId: info.id }, + item: { + type: "user_message", + text, + messageId: info.id, + ...(state.pendingClientMessageId ? { clientMessageId: state.pendingClientMessageId } : {}), + }, }); } @@ -2895,6 +2901,7 @@ class OpenCodeAgentSession implements AgentSession { /** Tracks the role of each message by ID to distinguish user from assistant messages */ private messageRoles = new Map(); private pendingUserMessageText: string | null = null; + private pendingClientMessageId: string | null = null; private emittedUserMessageIds = new Set(); /** Tracks streamed textual part IDs to suppress final full-text echoes from OpenCode. */ private streamedPartKeys = new Set(); @@ -3103,6 +3110,7 @@ class OpenCodeAgentSession implements AgentSession { const parts = buildOpenCodePromptParts(prompt); this.pendingUserMessageText = buildOpenCodeUserTimelineText(prompt); + this.pendingClientMessageId = options?.clientMessageId ?? null; this.suppressAssistantMessagesUntilIdle.active = false; const model = this.parseModel(this.config.model); const thinkingOptionId = this.config.thinkingOptionId; @@ -3645,6 +3653,7 @@ class OpenCodeAgentSession implements AgentSession { this.subAgentsByCallId.clear(); this.subAgentCallIdByChildSessionId.clear(); this.pendingUserMessageText = null; + this.pendingClientMessageId = null; this.abortController = null; this.notifySubscribers({ type: "turn_started", provider: "opencode" }, turnId); return turnId; @@ -3670,6 +3679,7 @@ class OpenCodeAgentSession implements AgentSession { this.runningToolCalls.clear(); } this.pendingUserMessageText = null; + this.pendingClientMessageId = null; this.activeForegroundTurnId = null; this.activeForegroundTurnSource = null; this.abortController = null; @@ -4068,6 +4078,7 @@ class OpenCodeAgentSession implements AgentSession { cwd: this.config.cwd, messageRoles: this.messageRoles, pendingUserMessageText: this.pendingUserMessageText, + pendingClientMessageId: this.pendingClientMessageId, emittedUserMessageIds: this.emittedUserMessageIds, accumulatedUsage: this.accumulatedUsage, sessionTotalCostUsd: this.sessionTotalCostUsd, diff --git a/packages/server/src/server/agent/providers/pi/agent.ts b/packages/server/src/server/agent/providers/pi/agent.ts index 7a18f6ebb..e87ff42ec 100644 --- a/packages/server/src/server/agent/providers/pi/agent.ts +++ b/packages/server/src/server/agent/providers/pi/agent.ts @@ -258,6 +258,7 @@ interface PiCapturedEntry extends PiCapturedUserMessageEntry { interface PendingPiUserMessage { text: string; turnId: string | undefined; + clientMessageId?: string; } interface PendingExtensionResult { @@ -1179,6 +1180,7 @@ export class PiRpcAgentSession implements AgentSession { private activeAskUserDialog: ActiveAskUserDialog | null = null; private pendingCombinedAskUserResponse: PendingCombinedAskUserResponse | null = null; private activeTurnId: string | null = null; + private activeClientMessageId: string | null = null; private activeAssistantMessageId: string | null = null; private activeTurnStarted = false; private activeNoTurnPromptText: string | null = null; @@ -1240,7 +1242,7 @@ export class PiRpcAgentSession implements AgentSession { }); } - async startTurn(prompt: AgentPromptInput, _options?: AgentRunOptions): Promise { + async startTurn(prompt: AgentPromptInput, options?: AgentRunOptions): Promise { if (this.activeTurnId) { throw new Error("A Pi turn is already active"); } @@ -1248,6 +1250,7 @@ export class PiRpcAgentSession implements AgentSession { const payload = convertPromptInput(prompt, { model: this.state.model }); const turnId = randomUUID(); this.activeTurnId = turnId; + this.activeClientMessageId = options?.clientMessageId ?? null; this.activeAssistantMessageId = null; this.activeTurnStarted = false; this.activePromptRequestId = null; @@ -1278,6 +1281,7 @@ export class PiRpcAgentSession implements AgentSession { return; } this.activeTurnId = null; + this.activeClientMessageId = null; this.activeTurnStarted = false; this.activeAssistantMessageId = null; this.clearNoTurnBuffers(); @@ -1393,6 +1397,7 @@ export class PiRpcAgentSession implements AgentSession { await this.runtimeSession.abort(); if (turnId && this.activeTurnId === turnId) { this.activeTurnId = null; + this.activeClientMessageId = null; this.activeTurnStarted = false; this.activeAssistantMessageId = null; this.clearNoTurnBuffers(); @@ -1574,6 +1579,7 @@ export class PiRpcAgentSession implements AgentSession { item: { type: "user_message", text: promptText, + ...(this.activeClientMessageId ? { clientMessageId: this.activeClientMessageId } : {}), }, }); } @@ -1796,6 +1802,7 @@ export class PiRpcAgentSession implements AgentSession { type: "user_message", text: pending.text, messageId: entry.id, + ...(pending.clientMessageId ? { clientMessageId: pending.clientMessageId } : {}), }, }); } @@ -1961,6 +1968,7 @@ export class PiRpcAgentSession implements AgentSession { } const turnId = this.activeTurnId; this.activeTurnId = null; + this.activeClientMessageId = null; this.activeTurnStarted = false; this.clearNoTurnBuffers(); this.emit({ @@ -2170,7 +2178,11 @@ export class PiRpcAgentSession implements AgentSession { if (!text) { return; } - this.pendingUserMessages.push({ text, turnId }); + this.pendingUserMessages.push({ + text, + turnId, + ...(this.activeClientMessageId ? { clientMessageId: this.activeClientMessageId } : {}), + }); void this.requestEntryCapture("message_end").catch((error: unknown) => { const message = error instanceof Error ? error.message : String(error); this.emit({ @@ -2221,6 +2233,7 @@ export class PiRpcAgentSession implements AgentSession { private completeTurn(turnId: string | undefined, messages: PiAgentMessage[]): void { this.activeTurnId = null; + this.activeClientMessageId = null; this.activeAssistantMessageId = null; this.activeTurnStarted = false; this.clearNoTurnBuffers();