Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -4247,6 +4247,67 @@ describe("ProviderRuntimeIngestion", () => {
expect(message?.streaming).toBe(false);
});

it("shows buffered prose before a tool starts and preserves later assistant text", async () => {
const harness = await createHarness();
const threadId = asThreadId("thread-1");
const turnId = asTurnId("turn-tool-boundary");
const itemId = asItemId("item-tool-boundary-assistant");
const now = "2026-01-01T00:00:00.000Z";
const base = { provider: ProviderDriverKind.make("claude"), createdAt: now, threadId, turnId };

await harness.emitAndDrain([
{ ...base, type: "turn.started", eventId: asEventId("tool-boundary-turn") },
{
...base,
type: "content.delta",
eventId: asEventId("tool-boundary-prose"),
itemId,
payload: { streamKind: "assistant_text", delta: "I will inspect the server" },
},
{
...base,
type: "item.started",
eventId: asEventId("tool-boundary-start"),
itemId: asItemId("tool-boundary-call"),
payload: { itemType: "command_execution", status: "inProgress", title: "Inspect" },
},
]);

const midThread = (await harness.readModel()).threads.find((entry) => entry.id === threadId);
const midMessage = midThread?.messages.find(
(entry: ProviderRuntimeTestMessage) => entry.id === `assistant:${itemId}`,
);
expect(midMessage).toMatchObject({ text: "I will inspect the server", streaming: true });
expect(
midThread?.activities.some(
(activity: ProviderRuntimeTestActivity) => activity.kind === "tool.started",
),
).toBe(true);

await harness.emitAndDrain([
{
...base,
type: "content.delta",
eventId: asEventId("tool-boundary-later-prose"),
itemId,
payload: { streamKind: "assistant_text", delta: " and report back." },
},
{
...base,
type: "item.completed",
eventId: asEventId("tool-boundary-prose-completed"),
itemId,
payload: { itemType: "assistant_message", status: "completed" },
},
]);
const finalThread = (await harness.readModel()).threads.find((entry) => entry.id === threadId);
expect(
finalThread?.messages.find(
(entry: ProviderRuntimeTestMessage) => entry.id === `assistant:${itemId}`,
),
).toMatchObject({ text: "I will inspect the server and report back.", streaming: false });
});

it("flushes and completes buffered assistant text when an approval request opens", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";
Expand Down
15 changes: 15 additions & 0 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3323,6 +3323,21 @@ const make = Effect.gen(function* () {
}

const activities = runtimeEventToActivities(activityEvent, taskTitle);
// A tool start is a natural prose boundary. A provider may begin a tool
// before it completes the assistant item, leaving an otherwise valid
// single paragraph buffered while tool activity is already visible.
if (activities.some((activity) => activity.kind === "tool.started")) {
const turnId = toTurnId(event.turnId);
if (turnId) {
yield* flushBufferedAssistantMessagesForTurn({
event,
threadId: thread.id,
turnId,
createdAt: now,
commandTag: "assistant-delta-flush-on-tool-started",
});
}
}
yield* Effect.forEach(activities, (activity) =>
providerCommandId(event, "thread-activity-append").pipe(
Effect.flatMap((commandId) =>
Expand Down
Loading