Skip to content
This repository was archived by the owner on Aug 6, 2026. It is now read-only.
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
88 changes: 81 additions & 7 deletions packages/agent/src/server/agent-server.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down Expand Up @@ -2782,20 +2859,17 @@ 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;
};
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")
Expand Down
52 changes: 42 additions & 10 deletions packages/agent/src/server/agent-server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -4109,23 +4125,25 @@ ${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,
sessionRegistered: this.session.logWriter.isRegistered(payload.run_id),
});
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(
Expand Down Expand Up @@ -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(),
Expand Down
30 changes: 16 additions & 14 deletions packages/agent/src/server/question-relay.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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),
},
};
Expand All @@ -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),
Expand All @@ -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,
);
Expand All @@ -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),
},
};
Expand All @@ -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),
},
};
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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),
Expand Down
59 changes: 59 additions & 0 deletions packages/agent/src/session-log-writer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 });
Expand Down
Loading
Loading