diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index c22cf3e53bea..c4f57726f20f 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -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"); @@ -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"); @@ -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"); @@ -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", @@ -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"; @@ -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( @@ -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"); diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 8126537bd9f7..d7ae589047bf 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -1,6 +1,5 @@ import { ApprovalRequestId, - type AssistantDeliveryMode, CommandId, MessageId, type OrchestrationEvent, @@ -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"; @@ -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, + atMillis: number, + ) => Cache.getOption(bufferedAssistantTextByMessageId, messageId).pipe( Effect.flatMap((existingText) => Effect.gen(function* () { @@ -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 }; const lastDeliveredAt = Option.getOrUndefined( yield* Cache.getOption(lastAssistantDeliveryAtByMessageId, messageId), ); @@ -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) { @@ -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, diff --git a/apps/web/src/components/settings/SettingsPanels.tsx b/apps/web/src/components/settings/SettingsPanels.tsx index 67a20b31ecd2..a0d713f58bee 100644 --- a/apps/web/src/components/settings/SettingsPanels.tsx +++ b/apps/web/src/components/settings/SettingsPanels.tsx @@ -38,6 +38,7 @@ import { MIN_PANEL_ANIMATION_DURATION_MS, MIN_PROMPT_FONT_SIZE, MIN_SIDEBAR_AUTO_SETTLE_AFTER_DAYS, + type ResponseStreamingMode, MIN_TERMINAL_FONT_SIZE, type QuitConfirmationMode, } from "@t3tools/contracts/settings"; @@ -94,6 +95,15 @@ import { isMacPlatform } from "../../lib/utils"; import { EMPTY_SERVER_PROVIDERS } from "../../state/server"; import { useArchivedThreadSnapshots } from "../../lib/archivedThreadsState"; import { formatRelativeTimeLabel } from "../../timestampFormat"; +import { + AlertDialog, + AlertDialogClose, + AlertDialogDescription, + AlertDialogFooter, + AlertDialogHeader, + AlertDialogPopup, + AlertDialogTitle, +} from "../ui/alert-dialog"; import { Button } from "../ui/button"; import { Collapsible, CollapsiblePanel, CollapsibleTrigger } from "../ui/collapsible"; import { @@ -168,6 +178,18 @@ const ENVIRONMENT_IDENTIFICATION_LABELS: Record = { + turn: "Wait for the full response", + paragraph: "Show finished paragraphs", + token: "Token by token (legacy)", +}; + +const RESPONSE_STREAMING_MODE_DESCRIPTIONS: Record = { + turn: "Text appears once the agent finishes its turn.", + paragraph: "Each paragraph or code block appears as soon as it is complete.", + token: "Every token repaints the message as it arrives. Slower and harder to read.", +}; + const TIMESTAMP_FORMAT_LABELS = { locale: "System default", "12-hour": "12-hour", @@ -564,9 +586,8 @@ export function useSettingsRestore(onRestored?: () => void) { ...(settings.contextWindowMeterEnabled !== DEFAULT_UNIFIED_SETTINGS.contextWindowMeterEnabled ? ["Context window indicator"] : []), - ...(settings.enableLegacyTokenStreaming !== - DEFAULT_UNIFIED_SETTINGS.enableLegacyTokenStreaming - ? ["Stream token by token"] + ...(settings.responseStreamingMode !== DEFAULT_UNIFIED_SETTINGS.responseStreamingMode + ? ["Response streaming"] : []), ...(settings.enableProviderUpdateChecks !== DEFAULT_UNIFIED_SETTINGS.enableProviderUpdateChecks @@ -640,7 +661,7 @@ export function useSettingsRestore(onRestored?: () => void) { settings.fontSizeTerminal, settings.glassOpacity, settings.panelAnimationDurationMs, - settings.enableLegacyTokenStreaming, + settings.responseStreamingMode, settings.enableProviderUpdateChecks, settings.continueThreadsAfterServerUpdate, settings.sidebarAutoSettleAfterDays, @@ -744,7 +765,7 @@ export function useSettingsRestore(onRestored?: () => void) { sidebarCompactThreadRows: DEFAULT_UNIFIED_SETTINGS.sidebarCompactThreadRows, sidebarAutoSettleAfterDays: DEFAULT_UNIFIED_SETTINGS.sidebarAutoSettleAfterDays, sidebarAutoSettleOnMerge: DEFAULT_UNIFIED_SETTINGS.sidebarAutoSettleOnMerge, - enableLegacyTokenStreaming: DEFAULT_UNIFIED_SETTINGS.enableLegacyTokenStreaming, + responseStreamingMode: DEFAULT_UNIFIED_SETTINGS.responseStreamingMode, enableProviderUpdateChecks: DEFAULT_UNIFIED_SETTINGS.enableProviderUpdateChecks, continueThreadsAfterServerUpdate: DEFAULT_UNIFIED_SETTINGS.continueThreadsAfterServerUpdate, backgroundActivity: DEFAULT_UNIFIED_SETTINGS.backgroundActivity, @@ -797,6 +818,44 @@ export function useSettingsRestore(onRestored?: () => void) { }; } +/** + * Gate in front of the legacy token-by-token mode. The primary action steers + * the user to paragraph streaming; the legacy path is the quiet option. + */ +function TokenStreamingWarningDialog({ + open, + onOpenChange, + onConfirm, + onUseParagraphs, +}: { + open: boolean; + onOpenChange: (open: boolean) => void; + onConfirm: () => void; + onUseParagraphs: () => void; +}) { + return ( + + + + Token by token is a worse experience + + Token streaming repaints the message on every delta. It is slower, harder to read, and + costs more CPU on every connected device. This mode stays only for backwards + compatibility. Use paragraph streaming instead. + + + + + }>Cancel + + + + + ); +} + function BackgroundActivityAdvancedDialog({ open, onOpenChange, @@ -2014,7 +2073,6 @@ function AutoSettleDaysInput({ const LEGACY_FEATURE_TARGET_IDS: ReadonlySet = new Set([ "legacy-plan-mode", "legacy-context-window-indicator", - "legacy-token-streaming", "legacy-sidebar", ]); @@ -2082,35 +2140,6 @@ function LegacyFeaturesSection() { /> } /> - { - if (!checked) { - updateSettings({ enableLegacyTokenStreaming: false }); - return; - } - void (async () => { - const api = readLocalApi(); - const confirmed = await (api ?? ensureLocalApi()).dialogs.confirm( - [ - "Turn on token-by-token output?", - "It is significantly slower than the default buffered output and hurts the reading experience. This switch exists only for backwards compatibility.", - ].join("\n"), - ); - if (confirmed) updateSettings({ enableLegacyTokenStreaming: true }); - })(); - }} - aria-label="Stream token by token (legacy)" - /> - } - /> 0; const [backgroundActivityDialogOpen, setBackgroundActivityDialogOpen] = useState(false); + const [tokenStreamingWarningOpen, setTokenStreamingWarningOpen] = useState(false); + const mixedResponseStreamingMode = useScopedSettingsMixed(["responseStreamingMode"]); const lastEnabledProjectGroupingMode = useRef( readLastEnabledProjectGroupingMode(), ); @@ -2387,6 +2418,76 @@ export function GeneralSettingsPanel() { } /> + + updateSettings({ + responseStreamingMode: DEFAULT_UNIFIED_SETTINGS.responseStreamingMode, + }) + } + /> + ) : null + } + control={ + <> + + { + updateSettings({ responseStreamingMode: "token" }); + setTokenStreamingWarningOpen(false); + }} + onUseParagraphs={() => { + updateSettings({ responseStreamingMode: "paragraph" }); + setTokenStreamingWarningOpen(false); + }} + /> + + } + /> { expect(isSettingsSearchScopeAvailable(updates.scope, "environment")).toBe(true); expect(isSettingsSearchScopeAvailable(updates.scope, "all")).toBe(true); expect(isSettingsSearchScopeAvailable(updates.scope, "project")).toBe(false); - const streaming = getSettingsSearchTargetScope("legacy-token-streaming")!; + const streaming = getSettingsSearchTargetScope("response-streaming")!; expect(streaming.scope).toBe("project-defaults"); expect(isSettingsSearchScopeAvailable(streaming.scope, "project")).toBe(true); for (const id of ["legacy-plan-mode", "legacy-context-window-indicator", "legacy-sidebar"]) { diff --git a/apps/web/src/components/settings/settingsSearch.ts b/apps/web/src/components/settings/settingsSearch.ts index e6961badcfc5..1fe96d2b4df5 100644 --- a/apps/web/src/components/settings/settingsSearch.ts +++ b/apps/web/src/components/settings/settingsSearch.ts @@ -261,6 +261,13 @@ export const SETTINGS_SEARCH_ITEMS = [ to: "/settings/general", searchTerms: ["timestamp clock locale system browser os 12 hour 24 hour"], }, + { + id: "response-streaming", + title: "Response streaming", + to: "/settings/general", + scope: "project-defaults", + searchTerms: ["output token paragraph buffered wait turn legacy"], + }, { id: "hide-whitespace-changes", title: "Hide whitespace changes", @@ -398,13 +405,6 @@ export const SETTINGS_SEARCH_ITEMS = [ to: "/settings/general", searchTerms: ["composer meter usage tokens circle old"], }, - { - id: "legacy-token-streaming", - title: "Stream token by token (legacy)", - to: "/settings/general", - scope: "project-defaults", - searchTerms: ["response output old compatibility"], - }, { id: "legacy-sidebar", title: "Sidebar (legacy)", diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 774a630ad44b..7b1fd2f96597 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -143,8 +143,6 @@ export const ProviderRequestKind = Schema.Literals([ "mcp-elicitation", ]); export type ProviderRequestKind = typeof ProviderRequestKind.Type; -export const AssistantDeliveryMode = Schema.Literals(["buffered", "streaming"]); -export type AssistantDeliveryMode = typeof AssistantDeliveryMode.Type; export const ProviderApprovalDecision = Schema.Literals([ "accept", "acceptForSession", diff --git a/packages/contracts/src/settings.ts b/packages/contracts/src/settings.ts index 9d3eaf877af8..1262a303ba41 100644 --- a/packages/contracts/src/settings.ts +++ b/packages/contracts/src/settings.ts @@ -952,6 +952,15 @@ export type BackgroundActivitySettings = typeof BackgroundActivitySettings.Type; * background activity, theme. UI, search and the write planner derive * eligibility from this list, so adding a key here is the whole opt-in. */ +/** + * How assistant text reaches clients while a turn runs. + * - `turn`: hold the whole message until the turn finishes or pauses. + * - `paragraph`: deliver each finished paragraph or closed code block. + * - `token`: forward every provider delta. Legacy, kept for compatibility. + */ +export const ResponseStreamingMode = Schema.Literals(["turn", "paragraph", "token"]); +export type ResponseStreamingMode = typeof ResponseStreamingMode.Type; + export const PROJECT_SCOPED_SERVER_SETTING_KEYS = [ "defaultModelSelection", "defaultRuntimeMode", @@ -968,7 +977,7 @@ export const PROJECT_SCOPED_SERVER_SETTING_KEYS = [ "sidebarAutoSettleOnMerge", "sidebarAutoSettleAfterDays", "continueThreadsAfterServerUpdate", - "enableLegacyTokenStreaming", + "responseStreamingMode", ] as const; export type ProjectScopedServerSettingKey = (typeof PROJECT_SCOPED_SERVER_SETTING_KEYS)[number]; @@ -993,16 +1002,17 @@ export const ProjectSettingsOverrides = Schema.Struct({ sidebarAutoSettleOnMerge: Schema.optionalKey(Schema.Boolean), sidebarAutoSettleAfterDays: Schema.optionalKey(Schema.NullOr(SidebarAutoSettleAfterDays)), continueThreadsAfterServerUpdate: Schema.optionalKey(Schema.Boolean), - enableLegacyTokenStreaming: Schema.optionalKey(Schema.Boolean), + responseStreamingMode: Schema.optionalKey(ResponseStreamingMode), } satisfies Record); export type ProjectSettingsOverrides = typeof ProjectSettingsOverrides.Type; export const ServerSettings = Schema.Struct({ - // Legacy token-by-token assistant output. Deliberately a fresh key (was + // How assistant text reaches clients during a turn. Deliberately a fresh + // key (was `enableLegacyTokenStreaming`, before that // `enableAssistantStreaming`): decoding drops the old key, so everyone, - // including prior opt-ins, resets to the buffered default. - enableLegacyTokenStreaming: Schema.Boolean.pipe( - Schema.withDecodingDefault(Effect.succeed(false)), + // including prior token-streaming opt-ins, resets to the paragraph default. + responseStreamingMode: ResponseStreamingMode.pipe( + Schema.withDecodingDefault(Effect.succeed("paragraph" as const)), ), enableProviderUpdateChecks: Schema.Boolean.pipe(Schema.withDecodingDefault(Effect.succeed(true))), // Retain the update-era key; recovery now needs an environment-owned opt-in. @@ -1337,7 +1347,7 @@ const OpenCodeSettingsPatch = Schema.Struct({ export const ServerSettingsPatch = Schema.Struct({ // Server settings - enableLegacyTokenStreaming: Schema.optionalKey(Schema.Boolean), + responseStreamingMode: Schema.optionalKey(ResponseStreamingMode), enableProviderUpdateChecks: Schema.optionalKey(Schema.Boolean), continueThreadsAfterServerUpdate: Schema.optionalKey(Schema.Boolean), enableAgentBrowserAccess: Schema.optionalKey(Schema.Boolean),