From cc759dc7bc1703868b79718f777c73aa806ca90f Mon Sep 17 00:00:00 2001 From: Vojta Bartos Date: Wed, 22 Jul 2026 15:13:57 +0200 Subject: [PATCH] Relay background-turn responses to Slack A background (task-notification) turn resolves no tracked prompt, so the prompt-resolution relay sites never fire for it: when the agent ends its turn to wait on a background process, the answer it produces on wake-up never reaches Slack, and the agent has no way to know it was dropped. Trigger the relay from the ACP transport tap when a clean _posthog/background_turn_complete passes through, and scope every relay to the prose no prior relay has delivered (background turns append to the same turn buffer as the tracked turn before them, so an unscoped read would re-post already-delivered prose, and a silent wake-up would re-post the previous answer). --- .../agent/src/server/agent-server.test.ts | 88 +++++++++++++++++-- packages/agent/src/server/agent-server.ts | 52 ++++++++--- .../agent/src/server/question-relay.test.ts | 30 ++++--- packages/agent/src/session-log-writer.test.ts | 59 +++++++++++++ packages/agent/src/session-log-writer.ts | 42 ++++++++- 5 files changed, 236 insertions(+), 35 deletions(-) diff --git a/packages/agent/src/server/agent-server.test.ts b/packages/agent/src/server/agent-server.test.ts index 99e0c67a9d..aa701f5a33 100644 --- a/packages/agent/src/server/agent-server.test.ts +++ b/packages/agent/src/server/agent-server.test.ts @@ -479,6 +479,83 @@ describe("AgentServer HTTP Mode", () => { testServer.session = null; }); + describe("background turn relay", () => { + // A background (task-notification) turn resolves no tracked prompt, so + // the prompt-resolution relay sites never fire — the transport tap must + // relay its prose when _posthog/background_turn_complete passes through. + const stubRelaySession = (parts: string[] | undefined) => { + const testServer = createServer() as unknown as { + session: unknown; + posthogAPI: PostHogAPIClient; + handleAcpTransportMessage(message: unknown): void; + }; + const relaySpy = vi + .spyOn(testServer.posthogAPI, "relayMessage") + .mockResolvedValue(undefined); + testServer.session = { + payload: { task_id: "test-task-id", run_id: "test-run-id" }, + sseController: null, + logWriter: { + flush: vi.fn().mockResolvedValue(undefined), + takeUnrelayedAgentResponseParts: vi.fn().mockReturnValue(parts), + isRegistered: vi.fn().mockReturnValue(true), + }, + }; + return { testServer, relaySpy }; + }; + + it("relays the background turn's prose on background_turn_complete", async () => { + const { testServer, relaySpy } = stubRelaySession([ + "fetch finished", + "the actual answer", + ]); + + testServer.handleAcpTransportMessage({ + jsonrpc: "2.0", + method: "_posthog/background_turn_complete", + params: { sessionId: "s-1", stopReason: "end_turn" }, + }); + + await vi.waitFor(() => expect(relaySpy).toHaveBeenCalledTimes(1)); + expect(relaySpy).toHaveBeenCalledWith( + "test-task-id", + "test-run-id", + "fetch finished\n\nthe actual answer", + ["fetch finished", "the actual answer"], + undefined, + ); + testServer.session = null; + }); + + it("does not relay a background turn that did not end cleanly", async () => { + const { testServer, relaySpy } = stubRelaySession(["half an answer"]); + + testServer.handleAcpTransportMessage({ + jsonrpc: "2.0", + method: "_posthog/background_turn_complete", + params: { sessionId: "s-1", stopReason: "refusal" }, + }); + + await new Promise((resolve) => setImmediate(resolve)); + expect(relaySpy).not.toHaveBeenCalled(); + testServer.session = null; + }); + + it("does not relay a background turn that produced no new prose", async () => { + const { testServer, relaySpy } = stubRelaySession(undefined); + + testServer.handleAcpTransportMessage({ + jsonrpc: "2.0", + method: "_posthog/background_turn_complete", + params: { sessionId: "s-1", stopReason: "end_turn" }, + }); + + await new Promise((resolve) => setImmediate(resolve)); + expect(relaySpy).not.toHaveBeenCalled(); + testServer.session = null; + }); + }); + describe("GET /health", () => { it("returns ok status with active session", async () => { await createServer().start(); @@ -2782,8 +2859,9 @@ describe("AgentServer HTTP Mode", () => { session: { clientConnection: { prompt: typeof prompt }; logWriter: { - getFullAgentResponse: (runId: string) => string | undefined; - getAgentResponseParts: (runId: string) => string[]; + takeUnrelayedAgentResponseParts: ( + runId: string, + ) => string[] | undefined; }; }; posthogAPI: PostHogAPIClient; @@ -2791,11 +2869,7 @@ describe("AgentServer HTTP Mode", () => { serverInternals.session.clientConnection.prompt = prompt; vi.spyOn( serverInternals.session.logWriter, - "getFullAgentResponse", - ).mockReturnValue("final answer"); - vi.spyOn( - serverInternals.session.logWriter, - "getAgentResponseParts", + "takeUnrelayedAgentResponseParts", ).mockReturnValue(["final answer"]); const relaySpy = vi .spyOn(serverInternals.posthogAPI, "relayMessage") diff --git a/packages/agent/src/server/agent-server.ts b/packages/agent/src/server/agent-server.ts index 3a1853f7b9..e7e715c161 100644 --- a/packages/agent/src/server/agent-server.ts +++ b/packages/agent/src/server/agent-server.ts @@ -264,6 +264,22 @@ export function isTurnCompleteNotification(message: unknown): boolean { ); } +export function backgroundTurnCompleteStopReason( + message: unknown, +): string | null { + if ( + typeof message !== "object" || + message === null || + (message as { method?: unknown }).method !== + POSTHOG_NOTIFICATIONS.BACKGROUND_TURN_COMPLETE + ) { + return null; + } + const stopReason = (message as { params?: { stopReason?: unknown } }).params + ?.stopReason; + return typeof stopReason === "string" ? stopReason : "end_turn"; +} + interface SseController { send: (data: unknown) => void; close: () => void; @@ -4109,8 +4125,17 @@ ${signedCommitInstructions}${prLinkInstructions}${shellEfficiencyInstructions} }); } - const message = this.session.logWriter.getFullAgentResponse(payload.run_id); - if (!message) { + // Ordered assistant text blocks (one per message between tool calls), + // scoped to what no prior relay has delivered — a background turn appends + // to the same buffer as the tracked turn before it, so an unscoped read + // would re-post already-delivered prose. The backend picks the last entry + // — the post-last-tool-use answer — so Slack no longer sees the "Let me + // check…" narration. `message` stays as the joined fallback for backends + // that don't understand `text_parts`. + const messageParts = this.session.logWriter.takeUnrelayedAgentResponseParts( + payload.run_id, + ); + if (!messageParts) { this.logger.debug("No agent message found for Slack relay", { taskId: payload.task_id, runId: payload.run_id, @@ -4118,14 +4143,7 @@ ${signedCommitInstructions}${prLinkInstructions}${shellEfficiencyInstructions} }); return; } - - // Ordered assistant text blocks (one per message between tool calls). - // The backend picks the last entry — the post-last-tool-use answer — so - // Slack no longer sees the "Let me check…" narration. `message` stays as - // the joined fallback for backends that don't understand `text_parts`. - const messageParts = this.session.logWriter.getAgentResponseParts( - payload.run_id, - ); + const message = messageParts.join("\n\n"); try { await this.posthogAPI.relayMessage( @@ -4513,6 +4531,20 @@ ${signedCommitInstructions}${prLinkInstructions}${shellEfficiencyInstructions} } this.adapterEmittedTurnComplete = true; } + // A background (task-notification) turn resolves no tracked prompt, so + // the prompt-resolution relay sites never fire for it — without this, the + // prose it produced (e.g. results of a finished background process) never + // reaches Slack. Mirror the tracked path's gating on a clean end_turn. + if ( + backgroundTurnCompleteStopReason(message) === "end_turn" && + this.session + ) { + this.relayAgentResponse(this.session.payload).catch((error) => + this.logger.debug("Failed to relay background-turn response", { + error, + }), + ); + } const event = { type: "notification", timestamp: new Date().toISOString(), diff --git a/packages/agent/src/server/question-relay.test.ts b/packages/agent/src/server/question-relay.test.ts index bf1603a70a..3b5a3b3c1e 100644 --- a/packages/agent/src/server/question-relay.test.ts +++ b/packages/agent/src/server/question-relay.test.ts @@ -524,7 +524,9 @@ describe("Question relay", () => { payload: TEST_PAYLOAD, logWriter: { flush: vi.fn().mockResolvedValue(undefined), - getFullAgentResponse: vi.fn().mockReturnValue("agent response"), + takeUnrelayedAgentResponseParts: vi + .fn() + .mockReturnValue(["agent response"]), isRegistered: vi.fn().mockReturnValue(true), }, }; @@ -545,8 +547,7 @@ describe("Question relay", () => { payload: TEST_PAYLOAD, logWriter: { flush: vi.fn().mockResolvedValue(undefined), - getFullAgentResponse: vi.fn().mockReturnValue("agent response"), - getAgentResponseParts: vi + takeUnrelayedAgentResponseParts: vi .fn() .mockReturnValue(["first part", "agent response"]), isRegistered: vi.fn().mockReturnValue(true), @@ -559,7 +560,7 @@ describe("Question relay", () => { expect(relaySpy).toHaveBeenCalledWith( "test-task-id", "test-run-id", - "agent response", + "first part\n\nagent response", ["first part", "agent response"], undefined, ); @@ -574,8 +575,9 @@ describe("Question relay", () => { payload: TEST_PAYLOAD, logWriter: { flush: vi.fn().mockResolvedValue(undefined), - getFullAgentResponse: vi.fn().mockReturnValue("agent response"), - getAgentResponseParts: vi.fn().mockReturnValue(["agent response"]), + takeUnrelayedAgentResponseParts: vi + .fn() + .mockReturnValue(["agent response"]), isRegistered: vi.fn().mockReturnValue(true), }, }; @@ -601,7 +603,7 @@ describe("Question relay", () => { payload: TEST_PAYLOAD, logWriter: { flush: vi.fn().mockResolvedValue(undefined), - getFullAgentResponse: vi.fn().mockReturnValue(null), + takeUnrelayedAgentResponseParts: vi.fn().mockReturnValue(undefined), isRegistered: vi.fn().mockReturnValue(true), }, }; @@ -636,7 +638,7 @@ describe("Question relay", () => { clientConnection: { prompt: promptSpy }, logWriter: { flushAll: vi.fn().mockResolvedValue(undefined), - getFullAgentResponse: vi.fn().mockReturnValue(null), + takeUnrelayedAgentResponseParts: vi.fn().mockReturnValue(undefined), resetTurnMessages: vi.fn(), appendRawLine: vi.fn(), flush: vi.fn().mockResolvedValue(undefined), @@ -681,7 +683,7 @@ describe("Question relay", () => { clientConnection: { prompt: promptSpy }, logWriter: { flushAll: vi.fn().mockResolvedValue(undefined), - getFullAgentResponse: vi.fn().mockReturnValue(null), + takeUnrelayedAgentResponseParts: vi.fn().mockReturnValue(undefined), resetTurnMessages: vi.fn(), appendRawLine: vi.fn(), flush: vi.fn().mockResolvedValue(undefined), @@ -716,7 +718,7 @@ describe("Question relay", () => { clientConnection: { prompt: promptSpy }, logWriter: { flushAll: vi.fn().mockResolvedValue(undefined), - getFullAgentResponse: vi.fn().mockReturnValue(null), + takeUnrelayedAgentResponseParts: vi.fn().mockReturnValue(undefined), resetTurnMessages: vi.fn(), appendRawLine: vi.fn(), flush: vi.fn().mockResolvedValue(undefined), @@ -751,7 +753,7 @@ describe("Question relay", () => { clientConnection: { prompt: promptSpy }, logWriter: { flushAll: vi.fn().mockResolvedValue(undefined), - getFullAgentResponse: vi.fn().mockReturnValue(null), + takeUnrelayedAgentResponseParts: vi.fn().mockReturnValue(undefined), resetTurnMessages: vi.fn(), appendRawLine: vi.fn(), flush: vi.fn().mockResolvedValue(undefined), @@ -789,7 +791,7 @@ describe("Question relay", () => { clientConnection: { prompt: promptSpy }, logWriter: { flushAll: vi.fn().mockResolvedValue(undefined), - getFullAgentResponse: vi.fn().mockReturnValue(null), + takeUnrelayedAgentResponseParts: vi.fn().mockReturnValue(undefined), resetTurnMessages: vi.fn(), appendRawLine: vi.fn(), flush: vi.fn().mockResolvedValue(undefined), @@ -841,7 +843,7 @@ describe("Question relay", () => { clientConnection: { prompt: promptSpy }, logWriter: { flushAll: vi.fn().mockResolvedValue(undefined), - getFullAgentResponse: vi.fn().mockReturnValue(null), + takeUnrelayedAgentResponseParts: vi.fn().mockReturnValue(undefined), resetTurnMessages: vi.fn(), appendRawLine: vi.fn(), flush: vi.fn().mockResolvedValue(undefined), @@ -890,7 +892,7 @@ describe("Question relay", () => { clientConnection: { prompt: promptSpy }, logWriter: { flushAll: vi.fn().mockResolvedValue(undefined), - getFullAgentResponse: vi.fn().mockReturnValue(null), + takeUnrelayedAgentResponseParts: vi.fn().mockReturnValue(undefined), resetTurnMessages: vi.fn(), appendRawLine: vi.fn(), flush: vi.fn().mockResolvedValue(undefined), diff --git a/packages/agent/src/session-log-writer.test.ts b/packages/agent/src/session-log-writer.test.ts index efb2b91092..3deaab8c05 100644 --- a/packages/agent/src/session-log-writer.test.ts +++ b/packages/agent/src/session-log-writer.test.ts @@ -460,6 +460,65 @@ describe("SessionLogWriter", () => { expect(logWriter.getAgentResponseParts(sessionId)).toBeUndefined(); }); + it("takeUnrelayedAgentResponseParts returns only prose since the previous take", async () => { + const sessionId = "s1"; + logWriter.register(sessionId, { taskId: "t1", runId: sessionId }); + + // Tracked turn: prose, relayed once. + logWriter.appendRawLine( + sessionId, + makeSessionUpdate("agent_message", { + content: { type: "text", text: "I'll wait for the fetch." }, + }), + ); + expect(logWriter.takeUnrelayedAgentResponseParts(sessionId)).toEqual([ + "I'll wait for the fetch.", + ]); + + // A background turn appends to the same buffer without a reset; only + // its own prose must relay, not the already-delivered tracked prose. + logWriter.appendRawLine( + sessionId, + makeSessionUpdate("agent_message", { + content: { type: "text", text: "Fetch done — here's the answer." }, + }), + ); + expect(logWriter.takeUnrelayedAgentResponseParts(sessionId)).toEqual([ + "Fetch done — here's the answer.", + ]); + + // A silent background turn (no new prose) has nothing to relay. + expect( + logWriter.takeUnrelayedAgentResponseParts(sessionId), + ).toBeUndefined(); + }); + + it("takeUnrelayedAgentResponseParts starts over after a turn reset", async () => { + const sessionId = "s1"; + logWriter.register(sessionId, { taskId: "t1", runId: sessionId }); + + logWriter.appendRawLine( + sessionId, + makeSessionUpdate("agent_message", { + content: { type: "text", text: "first turn" }, + }), + ); + expect(logWriter.takeUnrelayedAgentResponseParts(sessionId)).toEqual([ + "first turn", + ]); + + logWriter.resetTurnMessages(sessionId); + logWriter.appendRawLine( + sessionId, + makeSessionUpdate("agent_message", { + content: { type: "text", text: "second turn" }, + }), + ); + expect(logWriter.takeUnrelayedAgentResponseParts(sessionId)).toEqual([ + "second turn", + ]); + }); + it("persisted log does not contain stale entries when chunks are superseded", async () => { const sessionId = "s1"; logWriter.register(sessionId, { taskId: "t1", runId: sessionId }); diff --git a/packages/agent/src/session-log-writer.ts b/packages/agent/src/session-log-writer.ts index 7941b8487e..e655df4956 100644 --- a/packages/agent/src/session-log-writer.ts +++ b/packages/agent/src/session-log-writer.ts @@ -44,6 +44,8 @@ interface SessionState { chunkBuffer?: ChunkBuffer; lastAgentMessage?: string; currentTurnMessages: string[]; + /** Count of currentTurnMessages entries already delivered by a relay. */ + relayedTurnMessageCount: number; toolUpdateCache: Map; pendingRawInputSnapshots: Map; } @@ -106,6 +108,7 @@ export class SessionLogWriter { this.sessions.set(sessionId, { context, currentTurnMessages: [], + relayedTurnMessageCount: 0, toolUpdateCache: new Map(), pendingRawInputSnapshots: new Map(), }); @@ -523,10 +526,8 @@ export class SessionLogWriter { /** * Returns the ordered assistant text blocks for the current turn — one entry * per message between tool calls. The last entry is the text after the final - * tool_use (the actual answer to the user). - * - * The Slack relay uses this so the backend can post only the last block - * instead of every interim "Let me check…" narration. + * tool_use (the actual answer to the user). Non-consuming; the Slack relay + * uses takeUnrelayedAgentResponseParts instead. */ getAgentResponseParts(sessionId: string): string[] | undefined { const session = this.sessions.get(sessionId); @@ -549,7 +550,40 @@ export class SessionLogWriter { const session = this.sessions.get(sessionId); if (session) { session.currentTurnMessages = []; + session.relayedTurnMessageCount = 0; + } + } + + /** + * Returns the assistant text blocks accumulated since the last take (or the + * last turn reset) and advances the relay cursor past them. + * + * The Slack relay uses this instead of getAgentResponseParts because not + * every turn resets the buffer: a background (task-notification) turn + * appends to the same buffer as the tracked turn before it, and relaying + * the full buffer would re-post prose that was already delivered. + */ + takeUnrelayedAgentResponseParts(sessionId: string): string[] | undefined { + const session = this.sessions.get(sessionId); + if (!session) return undefined; + + if (session.chunkBuffer) { + this.logger.warn( + "takeUnrelayedAgentResponseParts called with non-empty chunk buffer", + { + sessionId, + bufferedLength: session.chunkBuffer.text.length, + }, + ); } + + const start = Math.min( + session.relayedTurnMessageCount, + session.currentTurnMessages.length, + ); + const parts = session.currentTurnMessages.slice(start); + session.relayedTurnMessageCount = session.currentTurnMessages.length; + return parts.length > 0 ? parts : undefined; } private extractAgentMessageText(