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 @@ -465,11 +465,11 @@ describe("ProviderRuntimeIngestion", () => {
});

it.each([
{ delivery: "buffered", enableLegacyTokenStreaming: false },
{ delivery: "streamed", enableLegacyTokenStreaming: true },
{ delivery: "buffered", responseStreamingMode: "paragraph" as const },
{ delivery: "streamed", responseStreamingMode: "token" as const },
])("settles OpenCode aborted turns and saves $delivery assistant text", async (settings) => {
const harness = await createHarness({
serverSettings: { enableLegacyTokenStreaming: settings.enableLegacyTokenStreaming },
serverSettings: { responseStreamingMode: settings.responseStreamingMode },
});
const threadId = asThreadId("thread-1");
const turnId = asTurnId("opencode-aborted-turn");
Expand Down Expand Up @@ -521,7 +521,7 @@ describe("ProviderRuntimeIngestion", () => {
"finalizes old buffered text on late %s without stopping the newer turn",
async (terminalType) => {
const harness = await createHarness({
serverSettings: { enableLegacyTokenStreaming: false },
serverSettings: { responseStreamingMode: "paragraph" },
});
const threadId = asThreadId("thread-1");
const oldTurnId = asTurnId("old-buffered-turn");
Expand Down Expand Up @@ -605,7 +605,7 @@ describe("ProviderRuntimeIngestion", () => {
{ source: "an unspecified turn", turnId: undefined },
])("ignores late OpenCode aborts for $source across newer turns", async (lateAbort) => {
const harness = await createHarness({
serverSettings: { enableLegacyTokenStreaming: true },
serverSettings: { responseStreamingMode: "token" },
});
const threadId = asThreadId("thread-1");
const stoppedTurnId = asTurnId("opencode-stopped-turn");
Expand Down Expand Up @@ -2712,7 +2712,7 @@ describe("ProviderRuntimeIngestion", () => {
});

it("keeps streaming while an async question is pending", async () => {
const harness = await createHarness({ serverSettings: { enableLegacyTokenStreaming: true } });
const harness = await createHarness({ serverSettings: { responseStreamingMode: "token" } });
const base = {
provider: ProviderDriverKind.make("codex"),
createdAt: "2026-01-01T00:00:00.000Z",
Expand Down Expand Up @@ -2954,7 +2954,7 @@ describe("ProviderRuntimeIngestion", () => {
});

it("starts a new streaming assistant message segment after approval", async () => {
const harness = await createHarness({ serverSettings: { enableLegacyTokenStreaming: true } });
const harness = await createHarness({ serverSettings: { responseStreamingMode: "token" } });
const startedAt = "2026-03-28T07:00:00.000Z";
const pausedAt = "2026-03-28T07:00:01.000Z";
const resumedAt = "2026-03-28T07:00:02.000Z";
Expand Down Expand Up @@ -3061,7 +3061,7 @@ describe("ProviderRuntimeIngestion", () => {
});

it("streams assistant deltas when thread.turn.start requests streaming mode", async () => {
const harness = await createHarness({ serverSettings: { enableLegacyTokenStreaming: true } });
const harness = await createHarness({ serverSettings: { responseStreamingMode: "token" } });
const now = "2026-01-01T00:00:00.000Z";

await Effect.runPromise(
Expand Down Expand Up @@ -3235,6 +3235,62 @@ describe("ProviderRuntimeIngestion", () => {
);
});

it("holds every paragraph until completion in turn mode", async () => {
const harness = await createHarness({ serverSettings: { responseStreamingMode: "turn" } });
const now = "2026-01-01T00:00:00.000Z";
const codex = ProviderDriverKind.make("codex");
const threadId = asThreadId("thread-1");
const turnId = asTurnId("turn-wait-mode");
const itemId = asItemId("item-wait-mode");

await harness.emitAndDrain([
{
type: "turn.started",
eventId: asEventId("evt-wait-started"),
provider: codex,
createdAt: now,
threadId,
turnId,
},
]);
harness.advanceClock(1_000);
await harness.emitAndDrain([
{
type: "content.delta",
eventId: asEventId("evt-wait-delta"),
provider: codex,
createdAt: now,
threadId,
turnId,
itemId,
payload: {
streamKind: "assistant_text",
delta: "First paragraph.\n\nSecond paragraph.\n\n",
},
},
]);
const messageText = async () =>
(await harness.readModel()).threads
.find((t) => t.id === threadId)
?.messages.find((m: ProviderRuntimeTestMessage) => m.id === `assistant:${itemId}`)?.text;
// Paragraph mode would have delivered both paragraphs by now.
expect(await messageText()).toBeUndefined();

await harness.emitAndDrain([
{
type: "item.completed",
eventId: asEventId("evt-wait-completed"),
provider: codex,
createdAt: now,
threadId,
turnId,
itemId,
payload: { itemType: "assistant_message", status: "completed" },
},
]);
expect(await messageText()).toBe("First paragraph.\n\nSecond paragraph.\n\n");
});

it("holds paragraphs that finish inside the pacing window and lands them together", async () => {
const harness = await createHarness();
const codex = ProviderDriverKind.make("codex");
Expand Down
48 changes: 27 additions & 21 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,5 @@
import {
ApprovalRequestId,
type AssistantDeliveryMode,
CommandId,
MessageId,
type OrchestrationEvent,
Expand All @@ -14,7 +13,9 @@ import {
TurnId,
type OrchestrationCheckpointSummary,
type OrchestrationThreadActivity,
type ProjectId,
type ProviderRuntimeEvent,
type ResponseStreamingMode,
RuntimeRequestId,
} from "@t3tools/contracts";
import * as Cache from "effect/Cache";
Expand Down Expand Up @@ -1167,7 +1168,19 @@ const make = Effect.gen(function* () {
});
});

const appendBufferedAssistantText = (messageId: MessageId, delta: string, atMillis: number) =>
const resolveResponseStreamingMode = (projectId: ProjectId) =>
Effect.map(
serverSettingsService.getSettings,
(settings) => resolveProjectSettings(settings, projectId).settings.responseStreamingMode,
);

// `mode` is "turn" or "paragraph"; token mode never buffers.
const appendBufferedAssistantText = (
messageId: MessageId,
delta: string,
mode: Exclude<ResponseStreamingMode, "token">,
atMillis: number,
) =>
Cache.getOption(bufferedAssistantTextByMessageId, messageId).pipe(
Effect.flatMap((existingText) =>
Effect.gen(function* () {
Expand All @@ -1176,9 +1189,13 @@ const make = Effect.gen(function* () {
onSome: (text) => `${text}${delta}`,
});

// Deliver finished paragraphs and closed code blocks early so the
// user sees progress without token-by-token repaints.
const { ready, rest } = splitBufferedAssistantText(nextText);
// Paragraph mode delivers finished paragraphs and closed code blocks
// early so the user sees progress without token-by-token repaints.
// Turn mode holds everything until the turn finishes or pauses.
const { ready, rest } =
mode === "paragraph"
? splitBufferedAssistantText(nextText)
: { ready: "", rest: nextText };
Comment thread
coderabbitai[bot] marked this conversation as resolved.
const lastDeliveredAt = Option.getOrUndefined(
yield* Cache.getOption(lastAssistantDeliveryAtByMessageId, messageId),
);
Expand Down Expand Up @@ -1757,19 +1774,14 @@ const make = Effect.gen(function* () {
yield* rememberAssistantMessageId(thread.id, turnId, assistantMessageId);
}

const assistantDeliveryMode: AssistantDeliveryMode = yield* Effect.map(
serverSettingsService.getSettings,
(settings) =>
resolveProjectSettings(settings, thread.projectId).settings.enableLegacyTokenStreaming
? "streaming"
: "buffered",
);
if (assistantDeliveryMode === "buffered") {
const streamingMode = yield* resolveResponseStreamingMode(thread.projectId);
if (streamingMode !== "token") {
// Pace on the server clock. OpenCode stamps every delta of a part
// with the part's start time, so the event time cannot measure gaps.
const spillChunk = yield* appendBufferedAssistantText(
assistantMessageId,
assistantDelta,
streamingMode,
yield* Clock.currentTimeMillis,
);
if (spillChunk.length > 0) {
Expand Down Expand Up @@ -1807,15 +1819,9 @@ const make = Effect.gen(function* () {
turnId: pauseForUserTurnId,
streamingOnly: true,
});
const assistantDeliveryMode: AssistantDeliveryMode = yield* Effect.map(
serverSettingsService.getSettings,
(settings) =>
resolveProjectSettings(settings, thread.projectId).settings.enableLegacyTokenStreaming
? "streaming"
: "buffered",
);
const streamingMode = yield* resolveResponseStreamingMode(thread.projectId);
const flushedMessageIds =
assistantDeliveryMode === "buffered"
streamingMode !== "token"
? yield* flushBufferedAssistantMessagesForTurn({
event,
threadId: thread.id,
Expand Down
Loading
Loading