From e20fa4b16e5f10680eef2b2b58a9714363e29a49 Mon Sep 17 00:00:00 2001 From: ClapFy Date: Wed, 5 Aug 2026 12:32:59 +0300 Subject: [PATCH 1/2] feat: add active-turn message controls --- .../features/settings/SettingsRouteScreen.tsx | 29 +++ .../src/features/threads/ThreadComposer.tsx | 10 +- .../features/threads/ThreadDetailScreen.tsx | 3 + .../features/threads/ThreadRouteScreen.tsx | 1 + .../src/persistence/mobile-preferences.ts | 9 + apps/mobile/src/state/thread-outbox-model.ts | 18 +- apps/mobile/src/state/thread-outbox.test.ts | 25 ++ .../src/state/use-thread-composer-state.ts | 11 +- .../src/state/use-thread-outbox-drain.ts | 1 + .../Layers/ProjectionPipeline.test.ts | 97 +++++++ .../src/provider/Layers/CodexAdapter.test.ts | 27 ++ .../Layers/CodexSessionRuntime.test.ts | 45 +++- .../provider/Layers/CodexSessionRuntime.ts | 97 +++++-- apps/web/src/components/ChatView.tsx | 227 +++++++++++++++- apps/web/src/components/chat/ChatComposer.tsx | 47 +++- .../chat/ComposerPrimaryActions.tsx | 136 +++++----- .../components/settings/SettingsPanels.tsx | 53 ++++ .../settings/settingsSearch.test.ts | 6 +- .../src/components/settings/settingsSearch.ts | 5 + apps/web/src/webThreadOutbox.test.ts | 155 +++++++++++ apps/web/src/webThreadOutbox.ts | 242 ++++++++++++++++++ packages/contracts/src/settings.test.ts | 16 ++ packages/contracts/src/settings.ts | 7 + 23 files changed, 1176 insertions(+), 91 deletions(-) create mode 100644 apps/web/src/webThreadOutbox.test.ts create mode 100644 apps/web/src/webThreadOutbox.ts diff --git a/apps/mobile/src/features/settings/SettingsRouteScreen.tsx b/apps/mobile/src/features/settings/SettingsRouteScreen.tsx index 8547859adde..8636f2d0c84 100644 --- a/apps/mobile/src/features/settings/SettingsRouteScreen.tsx +++ b/apps/mobile/src/features/settings/SettingsRouteScreen.tsx @@ -8,6 +8,7 @@ import { NativeStackScreenOptions } from "../../native/StackHeader"; import { SymbolView } from "../../components/AppSymbol"; import * as Effect from "effect/Effect"; import { AsyncResult } from "effect/unstable/reactivity"; +import { DEFAULT_ACTIVE_TURN_MESSAGE_BEHAVIOR } from "@t3tools/contracts/settings"; import { useCallback, useEffect, useMemo, useRef, useState, useSyncExternalStore } from "react"; import { Alert, Linking, Platform, Pressable, ScrollView, View } from "react-native"; import { useSafeAreaInsets } from "react-native-safe-area-context"; @@ -530,8 +531,36 @@ function ConfiguredSettingsRouteScreen() { } function GeneralSettingsSection() { + const preferencesResult = useAtomValue(mobilePreferencesAtom); + const savePreferences = useAtomSet(updateMobilePreferencesAtom); + const activeTurnMessageBehavior = AsyncResult.isSuccess(preferencesResult) + ? (preferencesResult.value.activeTurnMessageBehavior ?? DEFAULT_ACTIVE_TURN_MESSAGE_BEHAVIOR) + : DEFAULT_ACTIVE_TURN_MESSAGE_BEHAVIOR; + return ( + + Alert.alert( + "Messages while working", + "Steer adds the message to the active turn. Queue waits and sends messages one at a time after the current turn finishes.", + [ + { + text: "Steer", + onPress: () => savePreferences({ activeTurnMessageBehavior: "steer" }), + }, + { + text: "Queue", + onPress: () => savePreferences({ activeTurnMessageBehavior: "queue" }), + }, + { text: "Cancel", style: "cancel" }, + ], + ) + } + /> ); diff --git a/apps/mobile/src/features/threads/ThreadComposer.tsx b/apps/mobile/src/features/threads/ThreadComposer.tsx index aab896efe03..42c9e4b17ba 100644 --- a/apps/mobile/src/features/threads/ThreadComposer.tsx +++ b/apps/mobile/src/features/threads/ThreadComposer.tsx @@ -8,6 +8,7 @@ import type { RuntimeMode, ServerConfig as T3ServerConfig, } from "@t3tools/contracts"; +import type { ActiveTurnMessageBehavior } from "@t3tools/contracts/settings"; import { detectComposerTrigger, replaceTextRange, @@ -101,6 +102,7 @@ export interface ThreadComposerProps { readonly serverConfig: T3ServerConfig | null; readonly queueCount: number; readonly activeThreadBusy: boolean; + readonly activeTurnMessageBehavior: ActiveTurnMessageBehavior; readonly environmentId: EnvironmentId; readonly projectCwd: string | null; readonly editorRef?: RefObject; @@ -310,9 +312,13 @@ export const ThreadComposer = memo(function ThreadComposer(props: ThreadComposer props.selectedThread.session?.status === "starting"; const sendLabel = - props.connectionState !== "connected" || props.activeThreadBusy || props.queueCount > 0 + props.connectionState !== "connected" || props.queueCount > 0 ? "Queue" - : "Send"; + : props.activeThreadBusy + ? props.activeTurnMessageBehavior === "queue" + ? "Queue" + : "Steer" + : "Send"; const currentModelSelection = props.selectedThread.modelSelection; const currentRuntimeMode = props.selectedThread.runtimeMode; const currentInteractionMode = props.selectedThread.interactionMode ?? "default"; diff --git a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx index 5cb04290f66..7abe8006f77 100644 --- a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx +++ b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx @@ -14,6 +14,7 @@ import type { ServerConfig as T3ServerConfig, ThreadId, } from "@t3tools/contracts"; +import type { ActiveTurnMessageBehavior } from "@t3tools/contracts/settings"; import * as Haptics from "expo-haptics"; import { memo, useCallback, useEffect, useLayoutEffect, useMemo, useRef, useState } from "react"; import { Platform, View, type GestureResponderEvent } from "react-native"; @@ -62,6 +63,7 @@ export interface ThreadDetailScreenProps { /** Message sync status for the selected thread (drives the composer status pill). */ readonly threadSyncStatus?: EnvironmentThreadStatus; readonly activeThreadBusy: boolean; + readonly activeTurnMessageBehavior: ActiveTurnMessageBehavior; readonly environmentId: EnvironmentId; readonly projectWorkspaceRoot: string | null; readonly threadCwd: string | null; @@ -430,6 +432,7 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread serverConfig={props.serverConfig} queueCount={props.selectedThreadQueueCount} activeThreadBusy={props.activeThreadBusy} + activeTurnMessageBehavior={props.activeTurnMessageBehavior} environmentId={props.environmentId} projectCwd={props.projectWorkspaceRoot} bottomInset={composerBottomInset} diff --git a/apps/mobile/src/features/threads/ThreadRouteScreen.tsx b/apps/mobile/src/features/threads/ThreadRouteScreen.tsx index 7fb4740ddce..3e3bf6be6fd 100644 --- a/apps/mobile/src/features/threads/ThreadRouteScreen.tsx +++ b/apps/mobile/src/features/threads/ThreadRouteScreen.tsx @@ -767,6 +767,7 @@ function ThreadRouteContent( connectionStateLabel={routeConnectionState} threadSyncStatus={selectedThreadDetailState.status} activeThreadBusy={composer.activeThreadBusy} + activeTurnMessageBehavior={composer.activeTurnMessageBehavior} environmentId={selectedThread.environmentId} projectWorkspaceRoot={selectedThreadProject?.workspaceRoot ?? null} threadCwd={selectedThreadCwd} diff --git a/apps/mobile/src/persistence/mobile-preferences.ts b/apps/mobile/src/persistence/mobile-preferences.ts index 9a5ed82b3b8..20efd11357f 100644 --- a/apps/mobile/src/persistence/mobile-preferences.ts +++ b/apps/mobile/src/persistence/mobile-preferences.ts @@ -6,6 +6,7 @@ import * as Ref from "effect/Ref"; import * as Schema from "effect/Schema"; import * as Semaphore from "effect/Semaphore"; import type { SidebarProjectGroupingMode } from "@t3tools/contracts"; +import type { ActiveTurnMessageBehavior } from "@t3tools/contracts/settings"; import * as MobileDatabase from "./mobile-database"; import * as MobileSecureStorage from "./mobile-secure-storage"; @@ -15,6 +16,7 @@ const PREFERENCES_KEY = "t3code.preferences"; const PREFERENCES_FALLBACK_KEY = "t3code.preferences.fallback"; export interface Preferences { + readonly activeTurnMessageBehavior?: ActiveTurnMessageBehavior; readonly liveActivitiesEnabled?: boolean; readonly baseFontSize?: number; readonly terminalFontSize?: number | null; @@ -74,6 +76,7 @@ export class MobilePreferencesStore extends Context.Service< function sanitizePreferences(parsed: Preferences): Preferences { const preferences: { + activeTurnMessageBehavior?: ActiveTurnMessageBehavior; liveActivitiesEnabled?: boolean; baseFontSize?: number; terminalFontSize?: number | null; @@ -87,6 +90,12 @@ function sanitizePreferences(parsed: Preferences): Preferences { threadListV2Enabled?: boolean; } = {}; + if ( + parsed.activeTurnMessageBehavior === "steer" || + parsed.activeTurnMessageBehavior === "queue" + ) { + preferences.activeTurnMessageBehavior = parsed.activeTurnMessageBehavior; + } if (typeof parsed.liveActivitiesEnabled === "boolean") { preferences.liveActivitiesEnabled = parsed.liveActivitiesEnabled; } diff --git a/apps/mobile/src/state/thread-outbox-model.ts b/apps/mobile/src/state/thread-outbox-model.ts index 3ba61be3872..7ec8b77dd7b 100644 --- a/apps/mobile/src/state/thread-outbox-model.ts +++ b/apps/mobile/src/state/thread-outbox-model.ts @@ -16,12 +16,16 @@ import { type RuntimeMode as RuntimeModeType, } from "@t3tools/contracts"; import * as Schema from "effect/Schema"; +import { + ActiveTurnMessageBehavior, + type ActiveTurnMessageBehavior as ActiveTurnMessageBehaviorType, +} from "@t3tools/contracts/settings"; import { DraftComposerImageAttachmentSchema } from "../lib/composer-image-schema"; import type { DraftComposerImageAttachment } from "../lib/composerImages"; import { scopedThreadKey } from "../lib/scopedEntities"; -const THREAD_OUTBOX_SCHEMA_VERSION = 3; +const THREAD_OUTBOX_SCHEMA_VERSION = 4; const THREAD_OUTBOX_MAX_RETRY_DELAY_MS = 16_000; const QueuedThreadCreationSchema = Schema.Struct({ @@ -37,7 +41,7 @@ const QueuedThreadCreationSchema = Schema.Struct({ }); export const QueuedThreadMessageSchema = Schema.Struct({ - schemaVersion: Schema.Literals([1, 2, THREAD_OUTBOX_SCHEMA_VERSION]), + schemaVersion: Schema.Literals([1, 2, 3, THREAD_OUTBOX_SCHEMA_VERSION]), environmentId: EnvironmentId, threadId: ThreadId, messageId: MessageId, @@ -47,6 +51,7 @@ export const QueuedThreadMessageSchema = Schema.Struct({ modelSelection: Schema.optional(ModelSelection), runtimeMode: Schema.optional(RuntimeMode), interactionMode: Schema.optional(ProviderInteractionMode), + activeTurnMessageBehavior: Schema.optional(ActiveTurnMessageBehavior), // Present when the queued item creates a brand-new thread (pending task) // instead of appending a turn to an existing one. creation: Schema.optional(QueuedThreadCreationSchema), @@ -76,6 +81,11 @@ export interface QueuedThreadMessage { readonly modelSelection?: ModelSelectionType; readonly runtimeMode?: RuntimeModeType; readonly interactionMode?: ProviderInteractionModeType; + /** + * Snapshot of the send preference at enqueue time. Older persisted mobile + * outbox entries omit this and retain the historical queue behavior. + */ + readonly activeTurnMessageBehavior?: ActiveTurnMessageBehaviorType; readonly creation?: QueuedThreadCreation; readonly createdAt: string; } @@ -154,6 +164,7 @@ export function resolveThreadOutboxDeliveryAction(input: { readonly shellStatus: EnvironmentShellStatus; readonly environmentConnected: boolean; readonly threadBusy: boolean; + readonly activeTurnMessageBehavior?: ActiveTurnMessageBehaviorType; }): ThreadOutboxDeliveryAction { if (input.isCreation) { // A pending task creates its thread on delivery. If the thread already @@ -169,7 +180,8 @@ export function resolveThreadOutboxDeliveryAction(input: { if (!input.threadExists) { return input.shellStatus === "live" ? "remove" : "wait"; } - return input.environmentConnected && !input.threadBusy ? "send" : "wait"; + const canSendWhileBusy = input.activeTurnMessageBehavior === "steer" || !input.threadBusy; + return input.environmentConnected && canSendWhileBusy ? "send" : "wait"; } /** diff --git a/apps/mobile/src/state/thread-outbox.test.ts b/apps/mobile/src/state/thread-outbox.test.ts index 89f8b26798b..f0d42d82c51 100644 --- a/apps/mobile/src/state/thread-outbox.test.ts +++ b/apps/mobile/src/state/thread-outbox.test.ts @@ -487,6 +487,31 @@ describe("thread outbox", () => { ).toBe("send"); }); + it("waits behind active work in queue mode and dispatches into it in steer mode", () => { + const input = { + isCreation: false, + threadExists: true, + shellStatus: "live" as const, + environmentConnected: true, + threadBusy: true, + }; + + // Omitted preserves the behavior of outbox entries written by older mobile builds. + expect(resolveThreadOutboxDeliveryAction(input)).toBe("wait"); + expect( + resolveThreadOutboxDeliveryAction({ + ...input, + activeTurnMessageBehavior: "queue", + }), + ).toBe("wait"); + expect( + resolveThreadOutboxDeliveryAction({ + ...input, + activeTurnMessageBehavior: "steer", + }), + ).toBe("send"); + }); + it("sends queued creations once connected and live, removing already-created ones", () => { expect( resolveThreadOutboxDeliveryAction({ diff --git a/apps/mobile/src/state/use-thread-composer-state.ts b/apps/mobile/src/state/use-thread-composer-state.ts index b09aadf7e6b..938cd29d479 100644 --- a/apps/mobile/src/state/use-thread-composer-state.ts +++ b/apps/mobile/src/state/use-thread-composer-state.ts @@ -1,5 +1,6 @@ import { useAtomValue } from "@effect/atom-react"; import { useCallback, useEffect, useMemo } from "react"; +import { AsyncResult } from "effect/unstable/reactivity"; import { CommandId, @@ -10,6 +11,7 @@ import { type RuntimeMode, type ThreadId, } from "@t3tools/contracts"; +import { DEFAULT_ACTIVE_TURN_MESSAGE_BEHAVIOR } from "@t3tools/contracts/settings"; import { safeErrorLogAttributes } from "@t3tools/client-runtime/errors"; import { deriveActiveWorkStartedAt } from "@t3tools/shared/orchestrationTiming"; @@ -41,6 +43,7 @@ import { useSelectedThreadDetail } from "../state/use-thread-detail"; import { useThreadSelection } from "../state/use-thread-selection"; import { enqueueThreadOutboxMessage } from "./thread-outbox"; import { useThreadOutboxMessages } from "./use-thread-outbox"; +import { mobilePreferencesAtom } from "./preferences"; export function appendReviewCommentToDraft(input: { readonly environmentId: EnvironmentId; @@ -78,6 +81,10 @@ export function useThreadComposerState() { const selectedThreadDetail = useSelectedThreadDetail(); const composerDrafts = useAtomValue(composerDraftsAtom); const queuedMessagesByThreadKey = useThreadOutboxMessages(); + const preferencesResult = useAtomValue(mobilePreferencesAtom); + const activeTurnMessageBehavior = AsyncResult.isSuccess(preferencesResult) + ? (preferencesResult.value.activeTurnMessageBehavior ?? DEFAULT_ACTIVE_TURN_MESSAGE_BEHAVIOR) + : DEFAULT_ACTIVE_TURN_MESSAGE_BEHAVIOR; useEffect(() => { ensureComposerDraftsLoaded(); @@ -164,6 +171,7 @@ export function useThreadComposerState() { modelSelection: draft.modelSelection ?? thread.modelSelection, runtimeMode: draft.runtimeMode ?? thread.runtimeMode, interactionMode: draft.interactionMode ?? thread.interactionMode, + activeTurnMessageBehavior, createdAt: metadata.createdAt, }); clearComposerDraftContent(threadKey); @@ -179,7 +187,7 @@ export function useThreadComposerState() { ); }); return messageId; - }, [selectedThreadDetail, selectedThreadShell]); + }, [activeTurnMessageBehavior, selectedThreadDetail, selectedThreadShell]); const onChangeDraftMessage = useCallback( (value: string) => { @@ -308,6 +316,7 @@ export function useThreadComposerState() { modelSelection, runtimeMode, interactionMode, + activeTurnMessageBehavior, activeThreadBusy, onChangeDraftMessage, onPickDraftImages, diff --git a/apps/mobile/src/state/use-thread-outbox-drain.ts b/apps/mobile/src/state/use-thread-outbox-drain.ts index d06a4098aab..8a7f01875ae 100644 --- a/apps/mobile/src/state/use-thread-outbox-drain.ts +++ b/apps/mobile/src/state/use-thread-outbox-drain.ts @@ -315,6 +315,7 @@ export function useThreadOutboxDrain(): void { shellStatus, environmentConnected: environment?.connectionState === "connected", threadBusy: thread?.session?.status === "running" || thread?.session?.status === "starting", + activeTurnMessageBehavior: nextQueuedMessage.activeTurnMessageBehavior, }); if (deliveryAction === "wait") { continue; diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts index 926182a3ef0..49a2833d25f 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts @@ -2529,6 +2529,103 @@ it.layer(makeProjectionPipelinePrefixedTestLayer("t3-pending-turn-terminal-test- assert.deepEqual(pendingRows, []); }), ); + + it.effect("reconciles a steer with the already-running turn", () => + Effect.gen(function* () { + const projectionPipeline = yield* OrchestrationProjectionPipeline; + const eventStore = yield* OrchestrationEventStore; + const sql = yield* SqlClient.SqlClient; + const threadId = ThreadId.make("thread-steer-reaffirmed"); + const turnId = TurnId.make("turn-steer-reaffirmed"); + const initialMessageId = MessageId.make("message-steer-initial"); + const steeredMessageId = MessageId.make("message-steer-follow-up"); + + const appendTurnStartRequested = ( + eventSuffix: string, + messageId: MessageId, + createdAt: string, + ) => + eventStore.append({ + type: "thread.turn-start-requested", + eventId: EventId.make(`evt-steer-${eventSuffix}`), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: createdAt, + commandId: CommandId.make(`cmd-steer-${eventSuffix}`), + causationEventId: null, + correlationId: CorrelationId.make(`cmd-steer-${eventSuffix}`), + metadata: {}, + payload: { + threadId, + messageId, + runtimeMode: "full-access", + createdAt, + }, + }); + + const appendRunningSession = (eventSuffix: string, updatedAt: string) => + eventStore.append({ + type: "thread.session-set", + eventId: EventId.make(`evt-steer-${eventSuffix}`), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: updatedAt, + commandId: CommandId.make(`cmd-steer-${eventSuffix}`), + causationEventId: null, + correlationId: CorrelationId.make(`cmd-steer-${eventSuffix}`), + metadata: {}, + payload: { + threadId, + session: { + threadId, + status: "running", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: turnId, + lastError: null, + updatedAt, + }, + }, + }); + + yield* appendTurnStartRequested( + "initial-request", + initialMessageId, + "2026-02-26T15:00:00.000Z", + ); + yield* appendRunningSession("initial-running", "2026-02-26T15:00:01.000Z"); + yield* appendTurnStartRequested( + "follow-up-request", + steeredMessageId, + "2026-02-26T15:00:02.000Z", + ); + // A successful Codex turn/steer reaffirms the same active turn id. + yield* appendRunningSession("steer-reaffirmed", "2026-02-26T15:00:03.000Z"); + + yield* projectionPipeline.bootstrap; + + const turnRows = yield* sql<{ + readonly turnId: string | null; + readonly pendingMessageId: string | null; + readonly state: string; + }>` + SELECT + turn_id AS "turnId", + pending_message_id AS "pendingMessageId", + state + FROM projection_turns + WHERE thread_id = ${threadId} + ORDER BY row_id + `; + assert.deepEqual(turnRows, [ + { + turnId, + pendingMessageId: initialMessageId, + state: "running", + }, + ]); + }), + ); }, ); diff --git a/apps/server/src/provider/Layers/CodexAdapter.test.ts b/apps/server/src/provider/Layers/CodexAdapter.test.ts index 7b8fbec5666..5c7f336f26e 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.test.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.test.ts @@ -723,6 +723,33 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => { }), ); + it.effect("maps a reaffirmed turn/started event for a successful steer", () => + Effect.gen(function* () { + const { adapter, runtime } = yield* startLifecycleRuntime(); + const firstEventFiber = yield* Stream.runHead(adapter.streamEvents).pipe(Effect.forkChild); + + yield* runtime.emit({ + id: asEventId("evt-turn-steered"), + kind: "notification", + provider: ProviderDriverKind.make("codex"), + threadId: asThreadId("thread-1"), + createdAt: "2026-01-01T00:00:00.000Z", + method: "turn/started", + message: "Codex turn steered.", + turnId: asTurnId("turn-1"), + } satisfies ProviderEvent); + + const firstEvent = yield* Fiber.join(firstEventFiber); + + NodeAssert.equal(firstEvent._tag, "Some"); + if (firstEvent._tag !== "Some") { + return; + } + NodeAssert.equal(firstEvent.value.type, "turn.started"); + NodeAssert.equal(firstEvent.value.turnId, "turn-1"); + }), + ); + it.effect("maps retryable Codex error notifications to runtime.warning", () => Effect.gen(function* () { const { adapter, runtime } = yield* startLifecycleRuntime(); diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts index d7346a0e0db..8ae374c8b05 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.test.ts @@ -4,7 +4,7 @@ import { it } from "@effect/vitest"; import * as Effect from "effect/Effect"; import * as Schema from "effect/Schema"; import { describe } from "vite-plus/test"; -import { DEFAULT_MODEL, ThreadId } from "@t3tools/contracts"; +import { DEFAULT_MODEL, ThreadId, TurnId } from "@t3tools/contracts"; import * as CodexErrors from "effect-codex-app-server/errors"; import * as CodexRpc from "effect-codex-app-server/rpc"; @@ -15,10 +15,12 @@ import { } from "../CodexDeveloperInstructions.ts"; import { codexSessionAppServerArgs } from "./codexLaunchArgs.ts"; import { + buildTurnSteerParams, buildTurnStartParams, hasConfiguredMcpServer, isRecoverableThreadResumeError, openCodexThread, + resolveCodexSteeringTurnId, } from "./CodexSessionRuntime.ts"; const isCodexAppServerRequestError = Schema.is(CodexErrors.CodexAppServerRequestError); @@ -246,6 +248,47 @@ describe("buildTurnStartParams", () => { }); }); +describe("buildTurnSteerParams", () => { + it("reuses the active provider turn only while the session is running", () => { + const activeTurnId = TurnId.make("provider-turn-active"); + + NodeAssert.equal(resolveCodexSteeringTurnId({ status: "running", activeTurnId }), activeTurnId); + NodeAssert.equal(resolveCodexSteeringTurnId({ status: "ready", activeTurnId }), undefined); + }); + + it("targets the active turn and preserves text and image input", () => { + const activeTurnId = TurnId.make("provider-turn-active"); + + NodeAssert.deepStrictEqual( + buildTurnSteerParams({ + threadId: "provider-thread-1", + activeTurnId, + prompt: "Change direction", + attachments: [ + { + type: "image", + url: "data:image/png;base64,abc", + }, + ], + }), + { + threadId: "provider-thread-1", + expectedTurnId: activeTurnId, + input: [ + { + type: "text", + text: "Change direction", + }, + { + type: "image", + url: "data:image/png;base64,abc", + }, + ], + }, + ); + }); +}); + describe("buildCodexDeveloperInstructions", () => { it("appends runtime info after the mode instructions", () => { const instructions = buildCodexDeveloperInstructions("default", { diff --git a/apps/server/src/provider/Layers/CodexSessionRuntime.ts b/apps/server/src/provider/Layers/CodexSessionRuntime.ts index 67108dd4dbb..bad66ff53d8 100644 --- a/apps/server/src/provider/Layers/CodexSessionRuntime.ts +++ b/apps/server/src/provider/Layers/CodexSessionRuntime.ts @@ -88,6 +88,7 @@ export type CodexTurnStartParamsWithCollaborationMode = typeof CodexTurnStartParamsWithCollaborationMode.Type; export type CodexResumeCursor = typeof CodexResumeCursorSchema.Type; +type CodexTurnUserInput = EffectCodexSchema.V2TurnStartParams__UserInput; type CodexServiceTier = NonNullable; type CodexThreadItem = | EffectCodexSchema.V2ThreadReadResponse["thread"]["turns"][number]["items"][number] @@ -358,6 +359,48 @@ function buildCodexCollaborationMode(input: { }; } +function buildCodexTurnInput(input: { + readonly prompt?: string; + readonly attachments?: ReadonlyArray<{ + readonly type: "image"; + readonly url: string; + }>; +}): ReadonlyArray { + const turnInput: Array = []; + if (input.prompt) { + turnInput.push({ + type: "text", + text: input.prompt, + }); + } + for (const attachment of input.attachments ?? []) { + turnInput.push(attachment); + } + return turnInput; +} + +export function buildTurnSteerParams(input: { + readonly threadId: string; + readonly activeTurnId: TurnId; + readonly prompt?: string; + readonly attachments?: ReadonlyArray<{ + readonly type: "image"; + readonly url: string; + }>; +}): EffectCodexSchema.V2TurnSteerParams { + return { + threadId: input.threadId, + expectedTurnId: input.activeTurnId, + input: buildCodexTurnInput(input), + }; +} + +export function resolveCodexSteeringTurnId( + session: Pick, +): TurnId | undefined { + return session.status === "running" ? session.activeTurnId : undefined; +} + export function buildTurnStartParams(input: { readonly threadId: string; readonly runtimeMode: RuntimeMode; @@ -374,17 +417,6 @@ export function buildTurnStartParams(input: { CodexTurnStartParamsWithCollaborationMode, CodexErrors.CodexAppServerProtocolParseError > { - const turnInput: Array = []; - if (input.prompt) { - turnInput.push({ - type: "text", - text: input.prompt, - }); - } - for (const attachment of input.attachments ?? []) { - turnInput.push(attachment); - } - const config = runtimeModeToThreadConfig(input.runtimeMode); const collaborationMode = buildCodexCollaborationMode({ ...(input.interactionMode ? { interactionMode: input.interactionMode } : {}), @@ -394,7 +426,7 @@ export function buildTurnStartParams(input: { return decodeCodexTurnStartParamsWithCollaborationMode({ threadId: input.threadId, - input: turnInput, + input: buildCodexTurnInput(input), approvalPolicy: config.approvalPolicy, approvalsReviewer: config.approvalsReviewer, sandboxPolicy: runtimeModeToTurnSandboxPolicy(input.runtimeMode), @@ -1289,9 +1321,44 @@ export const makeCodexSessionRuntime = ( ), ); } - const normalizedModel = normalizeCodexModelSlug( - input.model ?? (yield* Ref.get(sessionRef)).model, - ); + const session = yield* Ref.get(sessionRef); + const steeringTurnId = resolveCodexSteeringTurnId(session); + if (steeringTurnId) { + const response = yield* client.request( + "turn/steer", + buildTurnSteerParams({ + threadId: providerThreadId, + activeTurnId: steeringTurnId, + ...(input.input ? { prompt: input.input } : {}), + ...(input.attachments ? { attachments: input.attachments } : {}), + }), + ); + const turnId = TurnId.make(response.turnId); + yield* updateSession(sessionRef, { + status: "running", + activeTurnId: turnId, + }); + // Codex does not emit a second turn/started notification when + // turn/steer keeps the existing turn alive. Reaffirm the active + // turn so orchestration can reconcile the steer request with the + // running turn instead of leaving a stale pending-turn row. + yield* emitEvent({ + kind: "notification", + threadId: options.threadId, + method: "turn/started", + turnId, + message: "Codex turn steered.", + }); + const resumedProviderThreadId = currentProviderThreadId(yield* Ref.get(sessionRef)); + return { + threadId: options.threadId, + turnId, + ...(resumedProviderThreadId + ? { resumeCursor: { threadId: resumedProviderThreadId } } + : {}), + } satisfies ProviderTurnStartResult; + } + const normalizedModel = normalizeCodexModelSlug(input.model ?? session.model); const params = yield* buildTurnStartParams({ threadId: providerThreadId, runtimeMode: options.runtimeMode, diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index ad704d9e1b2..c0c54a6b882 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -164,10 +164,19 @@ import { nextProjectScriptId, projectScriptIdFromCommand, } from "~/projectScripts"; -import { newDraftId, newMessageId, newThreadId } from "~/lib/utils"; +import { newCommandId, newDraftId, newMessageId, newThreadId } from "~/lib/utils"; import { getProviderModelCapabilities, resolveSelectableProvider } from "../providerModels"; import { NO_PROVIDER_MODEL_SELECTION } from "../providerInstances"; import { useClientSettings, useEnvironmentSettings } from "../hooks/useSettings"; +import { + beginWebThreadOutboxDispatch, + EMPTY_WEB_THREAD_OUTBOX_QUEUE, + finishWebThreadOutboxDispatch, + shouldDrainWebThreadOutbox, + shouldQueueWebThreadMessage, + useWebThreadOutboxStore, + webThreadOutboxKey, +} from "../webThreadOutbox"; import { useNowMinute } from "../hooks/useNowMinute"; import { useNewThreadHandler } from "../hooks/useHandleNewThread"; import { resolveAppModelSelectionForInstance } from "../modelSelection"; @@ -1226,6 +1235,12 @@ function ChatViewContent(props: ChatViewProps) { (store) => store.threadLastVisitedAtById[routeThreadKey], ); const settings = useEnvironmentSettings(environmentId); + const activeThreadOutboxQueue = useWebThreadOutboxStore( + (state) => + state.queuesByThreadKey[webThreadOutboxKey(environmentId, props.threadId)] ?? + EMPTY_WEB_THREAD_OUTBOX_QUEUE, + ); + const pausedOutboxMessageIds = useWebThreadOutboxStore((state) => state.pausedMessageIds); // New-thread defaults live in the primary environment's settings.json (the // settings UI never writes to remote environments), so read them from the // primary server rather than the thread's environment. @@ -4580,9 +4595,16 @@ function ChatViewContent(props: ChatViewProps) { }), ); }; + const shouldQueueCurrentMessage = shouldQueueWebThreadMessage({ + activeTurnMessageBehavior: settings.activeTurnMessageBehavior, + hasQueuedMessages: activeThreadOutboxQueue.length > 0, + isSendBusy, + isServerThread, + phase, + }); if ( !activeThread || - isSendBusy || + (isSendBusy && !shouldQueueCurrentMessage) || isConnecting || threadDetailLoading || activeEnvironmentUnavailable || @@ -4722,6 +4744,98 @@ function ChatViewContent(props: ChatViewProps) { return; } + if (shouldQueueCurrentMessage) { + sendInFlightRef.current = true; + const composerImagesSnapshot = [...composerImages]; + const composerTerminalContextsSnapshot = [...sendableComposerTerminalContexts]; + const composerElementContextsSnapshot = [...composerElementContexts]; + const composerPreviewAnnotationsSnapshot = [...composerPreviewAnnotations]; + const composerReviewCommentsSnapshot: ReviewCommentContext[] = [...composerReviewComments]; + const messageTextWithContexts = appendElementContextsToPrompt( + appendTerminalContextsToPrompt(promptForSend, composerTerminalContextsSnapshot), + composerElementContextsSnapshot, + ); + const messageTextWithPreviewAnnotations = composerPreviewAnnotationsSnapshot.reduce( + (text, annotation) => appendPreviewAnnotationPrompt(text, annotation), + messageTextWithContexts, + ); + const messageTextForSend = appendReviewCommentsToPrompt( + messageTextWithPreviewAnnotations, + composerReviewCommentsSnapshot, + ); + const outgoingMessageText = formatOutgoingPrompt({ + provider: ctxSelectedProvider, + model: ctxSelectedModel, + models: ctxSelectedProviderModels, + effort: ctxSelectedPromptEffort, + text: messageTextForSend || IMAGE_ONLY_BOOTSTRAP_PROMPT, + }); + const attachmentsResult = await settlePromise(() => + Promise.all( + composerImagesSnapshot.map(async (image) => ({ + type: "image" as const, + name: image.name, + mimeType: image.mimeType, + sizeBytes: image.sizeBytes, + dataUrl: await readFileAsDataUrl(image.file), + })), + ), + ); + if (attachmentsResult._tag === "Failure") { + const error = squashAtomCommandFailure(attachmentsResult); + setThreadError( + threadIdForSend, + error instanceof Error ? error.message : "Failed to prepare the queued message.", + ); + sendInFlightRef.current = false; + return; + } + + const messageId = newMessageId(); + const createdAt = new Date().toISOString(); + const { durable } = useWebThreadOutboxStore.getState().enqueue({ + environmentId, + threadId: threadIdForSend, + messageId, + commandId: newCommandId(), + text: outgoingMessageText, + attachments: attachmentsResult.value, + modelSelection: ctxSelectedModelSelection, + runtimeMode, + interactionMode, + createdAt, + }); + promptRef.current = ""; + clearComposerDraftContent(composerDraftTarget); + composerRef.current?.resetCursorState(); + setThreadError(threadIdForSend, null); + if (expiredTerminalContextCount > 0) { + const toastCopy = buildExpiredTerminalContextToastCopy( + expiredTerminalContextCount, + "omitted", + ); + toastManager.add( + stackedThreadToast({ + type: "warning", + title: toastCopy.title, + description: toastCopy.description, + }), + ); + } + if (!durable) { + toastManager.add( + stackedThreadToast({ + type: "warning", + title: "Message queued for this session", + description: + "Browser storage could not save the queue, so this message will not survive a reload.", + }), + ); + } + sendInFlightRef.current = false; + return; + } + sendInFlightRef.current = true; if (isDraftHeroState && activeThreadKey) { let resolveDockStarted: (() => void) | undefined; @@ -5001,6 +5115,112 @@ function ChatViewContent(props: ChatViewProps) { } }; + const nextQueuedMessage = activeThreadOutboxQueue[0] ?? null; + const queuedMessagePaused = + nextQueuedMessage !== null && Boolean(pausedOutboxMessageIds[nextQueuedMessage.messageId]); + + useEffect(() => { + if ( + !nextQueuedMessage || + !activeThread || + nextQueuedMessage.environmentId !== environmentId || + nextQueuedMessage.threadId !== activeThread.id || + !shouldDrainWebThreadOutbox({ + phase, + isSendBusy, + isConnecting, + environmentUnavailable: activeEnvironmentUnavailable, + paused: queuedMessagePaused, + }) || + !beginWebThreadOutboxDispatch(nextQueuedMessage.messageId) + ) { + return; + } + + const deliver = async () => { + beginLocalDispatch({ preparingWorktree: false }); + const settingsResult = await persistThreadSettingsForNextTurn({ + threadId: nextQueuedMessage.threadId, + createdAt: nextQueuedMessage.createdAt, + modelSelection: nextQueuedMessage.modelSelection, + runtimeMode: nextQueuedMessage.runtimeMode, + interactionMode: nextQueuedMessage.interactionMode, + }); + const startResult = + settingsResult._tag === "Failure" + ? settingsResult + : await startThreadTurn({ + environmentId: nextQueuedMessage.environmentId, + input: { + commandId: nextQueuedMessage.commandId, + threadId: nextQueuedMessage.threadId, + message: { + messageId: nextQueuedMessage.messageId, + role: "user", + text: nextQueuedMessage.text, + attachments: nextQueuedMessage.attachments, + }, + modelSelection: nextQueuedMessage.modelSelection, + titleSeed: activeThread.title, + runtimeMode: nextQueuedMessage.runtimeMode, + interactionMode: nextQueuedMessage.interactionMode, + createdAt: nextQueuedMessage.createdAt, + }, + }); + + if (startResult._tag === "Failure") { + useWebThreadOutboxStore.getState().pause(nextQueuedMessage.messageId); + resetLocalDispatch(); + if (!isAtomCommandInterrupted(startResult)) { + const error = squashAtomCommandFailure(startResult); + setThreadError( + nextQueuedMessage.threadId, + error instanceof Error ? error.message : "Failed to send the queued message.", + ); + } + return; + } + + const { durable } = useWebThreadOutboxStore.getState().remove(nextQueuedMessage); + if (!durable) { + toastManager.add( + stackedThreadToast({ + type: "warning", + title: "Queued message sent", + description: + "Browser storage could not save the queue update. Its stable command ID prevents a duplicate turn if it reappears after reload.", + }), + ); + } + }; + + void deliver().finally(() => { + finishWebThreadOutboxDispatch(nextQueuedMessage.messageId); + }); + }, [ + activeEnvironmentUnavailable, + activeThread, + beginLocalDispatch, + environmentId, + isConnecting, + isSendBusy, + nextQueuedMessage, + persistThreadSettingsForNextTurn, + phase, + queuedMessagePaused, + resetLocalDispatch, + setThreadError, + startThreadTurn, + ]); + + const retryQueuedMessages = useCallback(() => { + if (!nextQueuedMessage) { + return; + } + useWebThreadOutboxStore.getState().retry(nextQueuedMessage.messageId); + setThreadError(nextQueuedMessage.threadId, null); + }, [nextQueuedMessage, setThreadError]); + const onInterrupt = async () => { if (!activeThread) return; const result = await interruptThreadTurn({ @@ -5989,6 +6209,8 @@ function ChatViewContent(props: ChatViewProps) { isSendBusy={isSendBusy} sendDisabledReason={threadDetailLoading ? "Messages loading" : null} isPreparingWorktree={isPreparingWorktree} + queuedMessageCount={activeThreadOutboxQueue.length} + queuedMessagesPaused={queuedMessagePaused} environmentUnavailable={activeEnvironmentUnavailableState} activePendingApproval={activePendingApproval} pendingApprovals={pendingApprovals} @@ -6024,6 +6246,7 @@ function ChatViewContent(props: ChatViewProps) { composerTerminalContextsRef={composerTerminalContextsRef} composerElementContextsRef={composerElementContextsRef} onSend={onSend} + onRetryQueuedMessages={retryQueuedMessages} onInterrupt={onInterrupt} onImplementPlanInNewThread={onImplementPlanInNewThread} onRespondToApproval={onRespondToApproval} diff --git a/apps/web/src/components/chat/ChatComposer.tsx b/apps/web/src/components/chat/ChatComposer.tsx index 201fb79c566..cc15086a78e 100644 --- a/apps/web/src/components/chat/ChatComposer.tsx +++ b/apps/web/src/components/chat/ChatComposer.tsx @@ -445,6 +445,7 @@ const ComposerFooterPrimaryActions = memo(function ComposerFooterPrimaryActions( isConnecting: boolean; isEnvironmentUnavailable: boolean; hasSendableContent: boolean; + activeTurnMessageBehavior: UnifiedSettings["activeTurnMessageBehavior"]; preserveComposerFocusOnPointerDown?: boolean; onPreviousPendingQuestion: () => void; onInterrupt: () => void; @@ -473,6 +474,7 @@ const ComposerFooterPrimaryActions = memo(function ComposerFooterPrimaryActions( isEnvironmentUnavailable={props.isEnvironmentUnavailable} isPreparingWorktree={props.isPreparingWorktree} hasSendableContent={props.hasSendableContent} + activeTurnMessageBehavior={props.activeTurnMessageBehavior} preserveComposerFocusOnPointerDown={props.preserveComposerFocusOnPointerDown ?? false} onPreviousPendingQuestion={props.onPreviousPendingQuestion} onInterrupt={props.onInterrupt} @@ -551,6 +553,8 @@ export interface ChatComposerProps { isSendBusy: boolean; sendDisabledReason: string | null; isPreparingWorktree: boolean; + queuedMessageCount: number; + queuedMessagesPaused: boolean; environmentUnavailable: { readonly label: string; readonly connection: EnvironmentConnectionPresentation; @@ -610,6 +614,7 @@ export interface ChatComposerProps { // Callbacks onSend: (e?: { preventDefault: () => void }) => void; + onRetryQueuedMessages: () => void; onInterrupt: () => void; onImplementPlanInNewThread: () => void; onRespondToApproval: ( @@ -663,6 +668,8 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) isSendBusy, sendDisabledReason, isPreparingWorktree, + queuedMessageCount, + queuedMessagesPaused, environmentUnavailable, activePendingApproval, pendingApprovals, @@ -697,6 +704,7 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) composerTerminalContextsRef, composerElementContextsRef, onSend, + onRetryQueuedMessages, onInterrupt, onImplementPlanInNewThread, onRespondToApproval, @@ -1180,6 +1188,8 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) isComposerCollapsedMobile && !isComposerApprovalState && pendingUserInputs.length === 0; const composerFooterHasWideActions = showPlanFollowUpPrompt || activePendingProgress !== null; + const isPrimarySendBusy = + isSendBusy && !(phase === "running" && settings.activeTurnMessageBehavior === "queue"); const showPlanSidebarToggle = Boolean(activePlan || sidebarProposedPlan || planSidebarOpen); const composerFooterActionLayoutKey = useMemo(() => { if (activePendingProgress) { @@ -1270,15 +1280,19 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) [activePendingIsResponding, activePendingProgress, activePendingResolvedAnswers], ); const collapsedComposerPrimaryActionDisabled = - phase === "running" || - isSendBusy || + isPrimarySendBusy || isSendDisabled || isConnecting || noProviderAvailable || projectSelectionRequired || environmentUnavailable !== null || !composerSendState.hasSendableContent; - const collapsedComposerPrimaryActionLabel = "Send message"; + const collapsedComposerPrimaryActionLabel = + phase === "running" + ? settings.activeTurnMessageBehavior === "queue" + ? "Queue message" + : "Steer active turn" + : "Send message"; const showMobilePendingAnswerActions = isMobileViewport && !isComposerCollapsedMobile && pendingPrimaryAction !== null; @@ -2832,6 +2846,7 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) } isPreparingWorktree={false} hasSendableContent={false} + activeTurnMessageBehavior={settings.activeTurnMessageBehavior} preserveComposerFocusOnPointerDown onPreviousPendingQuestion={onPreviousActivePendingUserInputQuestion} onInterrupt={handleInterruptPrimaryAction} @@ -3115,6 +3130,7 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) } isPreparingWorktree={false} hasSendableContent={false} + activeTurnMessageBehavior={settings.activeTurnMessageBehavior} preserveComposerFocusOnPointerDown onPreviousPendingQuestion={onPreviousActivePendingUserInputQuestion} onInterrupt={handleInterruptPrimaryAction} @@ -3237,7 +3253,7 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) isRunning={phase === "running"} showPlanFollowUpPrompt={pendingUserInputs.length === 0 && showPlanFollowUpPrompt} promptHasText={prompt.trim().length > 0} - isSendBusy={isSendBusy} + isSendBusy={isPrimarySendBusy} sendDisabledReason={sendDisabledReason} isConnecting={isConnecting} isEnvironmentUnavailable={ @@ -3247,6 +3263,7 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) } isPreparingWorktree={isPreparingWorktree} hasSendableContent={composerSendState.hasSendableContent} + activeTurnMessageBehavior={settings.activeTurnMessageBehavior} preserveComposerFocusOnPointerDown={isMobileViewport} onPreviousPendingQuestion={onPreviousActivePendingUserInputQuestion} onInterrupt={handleInterruptPrimaryAction} @@ -3255,6 +3272,28 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) )} + {queuedMessageCount > 0 ? ( +
+ + {queuedMessageCount} queued message{queuedMessageCount === 1 ? "" : "s"} will send + one at a time. + + {queuedMessagesPaused ? ( + + ) : null} +
+ ) : null} diff --git a/apps/web/src/components/chat/ComposerPrimaryActions.tsx b/apps/web/src/components/chat/ComposerPrimaryActions.tsx index 504b7e1cc44..c3949b537c3 100644 --- a/apps/web/src/components/chat/ComposerPrimaryActions.tsx +++ b/apps/web/src/components/chat/ComposerPrimaryActions.tsx @@ -6,6 +6,7 @@ import { StageBackdropButtonArt, useSidebarStageBackdropVariant } from "../Sideb import { Button } from "../ui/button"; import { Menu, MenuItem, MenuPopup, MenuTrigger } from "../ui/menu"; import { Spinner } from "../ui/spinner"; +import type { ActiveTurnMessageBehavior } from "@t3tools/contracts/settings"; interface PendingActionState { questionIndex: number; @@ -27,6 +28,7 @@ interface ComposerPrimaryActionsProps { isEnvironmentUnavailable: boolean; isPreparingWorktree: boolean; hasSendableContent: boolean; + activeTurnMessageBehavior: ActiveTurnMessageBehavior; preserveComposerFocusOnPointerDown?: boolean; onPreviousPendingQuestion: () => void; onInterrupt: () => void; @@ -67,6 +69,7 @@ export const ComposerPrimaryActions = memo(function ComposerPrimaryActions({ isEnvironmentUnavailable, isPreparingWorktree, hasSendableContent, + activeTurnMessageBehavior, preserveComposerFocusOnPointerDown = false, onPreviousPendingQuestion, onInterrupt, @@ -132,19 +135,78 @@ export const ComposerPrimaryActions = memo(function ComposerPrimaryActions({ ); } + const sendButton = ( + + ); + if (isRunning) { return ( - +
+ + {sendButton} +
); } @@ -202,55 +264,5 @@ export const ComposerPrimaryActions = memo(function ComposerPrimaryActions({ ); } - return ( - - ); + return sendButton; }); diff --git a/apps/web/src/components/settings/SettingsPanels.tsx b/apps/web/src/components/settings/SettingsPanels.tsx index e3e0a22a0c3..ace47aaa108 100644 --- a/apps/web/src/components/settings/SettingsPanels.tsx +++ b/apps/web/src/components/settings/SettingsPanels.tsx @@ -31,6 +31,8 @@ import { squashAtomCommandFailure, } from "@t3tools/client-runtime/state/runtime"; import { + type ActiveTurnMessageBehavior, + DEFAULT_ACTIVE_TURN_MESSAGE_BEHAVIOR, DEFAULT_ENVIRONMENT_IDENTIFICATION_MODE, DEFAULT_UNIFIED_SETTINGS, type EnvironmentIdentificationMode, @@ -600,6 +602,9 @@ export function useSettingsRestore(onRestored?: () => void) { ...(settings.timestampFormat !== DEFAULT_UNIFIED_SETTINGS.timestampFormat ? ["Time format"] : []), + ...(settings.activeTurnMessageBehavior !== DEFAULT_UNIFIED_SETTINGS.activeTurnMessageBehavior + ? ["Messages while working"] + : []), ...(settings.sidebarThreadPreviewCount !== DEFAULT_UNIFIED_SETTINGS.sidebarThreadPreviewCount ? ["Visible threads"] : []), @@ -654,6 +659,7 @@ export function useSettingsRestore(onRestored?: () => void) { isTextGenerationModelDirty, isBackgroundActivityDirty, settings.autoOpenPlanSidebar, + settings.activeTurnMessageBehavior, settings.confirmThreadArchive, settings.confirmThreadDelete, settings.addProjectBaseDirectory, @@ -693,6 +699,7 @@ export function useSettingsRestore(onRestored?: () => void) { setTheme("system"); updateSettings({ timestampFormat: DEFAULT_UNIFIED_SETTINGS.timestampFormat, + activeTurnMessageBehavior: DEFAULT_UNIFIED_SETTINGS.activeTurnMessageBehavior, wordWrap: DEFAULT_UNIFIED_SETTINGS.wordWrap, diffIgnoreWhitespace: DEFAULT_UNIFIED_SETTINGS.diffIgnoreWhitespace, environmentIdentificationMode: DEFAULT_UNIFIED_SETTINGS.environmentIdentificationMode, @@ -1647,6 +1654,52 @@ export function GeneralSettingsPanel() { return ( + + updateSettings({ + activeTurnMessageBehavior: DEFAULT_ACTIVE_TURN_MESSAGE_BEHAVIOR, + }) + } + /> + ) : null + } + control={ + + } + /> + { it("matches normalized title substrings", () => { expect(searchSettings(" WORD WRAP ", ITEMS).map((item) => item.id)).toEqual(["word-wrap"]); - expect(searchSettings("work")).toEqual([]); + expect(searchSettings("work").map((item) => item.id)).toEqual(["messages-while-working"]); }); it("keeps catalog order for multiple title matches", () => { @@ -65,6 +65,10 @@ describe("searchSettings", () => { }); it("serves anchor props to panels from the catalog", () => { + expect(searchableSetting("messages-while-working")).toEqual({ + id: "messages-while-working", + title: "Messages while working", + }); expect(searchableSetting("word-wrap")).toEqual({ id: "word-wrap", title: "Word wrap" }); expect(searchableSetting("archive")).toEqual({ id: "archive", title: "Archived threads" }); }); diff --git a/apps/web/src/components/settings/settingsSearch.ts b/apps/web/src/components/settings/settingsSearch.ts index 1ba231a5835..6e904096b25 100644 --- a/apps/web/src/components/settings/settingsSearch.ts +++ b/apps/web/src/components/settings/settingsSearch.ts @@ -85,6 +85,11 @@ export const SETTINGS_SEARCH_ITEMS = [ title: "Word wrap", to: "/settings/appearance", }, + { + id: "messages-while-working", + title: "Messages while working", + to: "/settings/general", + }, { id: "project-grouping", title: "Project grouping", diff --git a/apps/web/src/webThreadOutbox.test.ts b/apps/web/src/webThreadOutbox.test.ts new file mode 100644 index 00000000000..efdee665fbc --- /dev/null +++ b/apps/web/src/webThreadOutbox.test.ts @@ -0,0 +1,155 @@ +import { + CommandId, + EnvironmentId, + MessageId, + ProviderInstanceId, + ThreadId, +} from "@t3tools/contracts"; +import { afterEach, describe, expect, it } from "vite-plus/test"; + +import { + beginWebThreadOutboxDispatch, + finishWebThreadOutboxDispatch, + shouldDrainWebThreadOutbox, + shouldQueueWebThreadMessage, + useWebThreadOutboxStore, + webThreadOutboxKey, + writeWebThreadOutboxStorageForTest, + type QueuedWebThreadMessage, +} from "./webThreadOutbox"; + +const environmentId = EnvironmentId.make("environment-test"); +const threadId = ThreadId.make("thread-test"); + +function message(index: number): QueuedWebThreadMessage { + return { + environmentId, + threadId, + messageId: MessageId.make(`message-${index}`), + commandId: CommandId.make(`command-${index}`), + text: `Message ${index}`, + attachments: [], + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5", + }, + runtimeMode: "full-access", + interactionMode: "default", + createdAt: new Date(1_700_000_000_000 + index).toISOString(), + }; +} + +function resetOutbox(): void { + writeWebThreadOutboxStorageForTest(""); +} + +afterEach(resetOutbox); + +describe("web thread outbox", () => { + it("keeps an unbounded FIFO per thread", () => { + const store = useWebThreadOutboxStore.getState(); + for (let index = 0; index < 100; index += 1) { + store.enqueue(message(index)); + } + + const queue = + useWebThreadOutboxStore.getState().queuesByThreadKey[ + webThreadOutboxKey(environmentId, threadId) + ]; + expect(queue).toHaveLength(100); + expect(queue?.map((entry) => entry.messageId)).toEqual( + Array.from({ length: 100 }, (_, index) => MessageId.make(`message-${index}`)), + ); + }); + + it("deduplicates stable message ids and removes only the delivered head", () => { + const first = message(1); + const second = message(2); + const store = useWebThreadOutboxStore.getState(); + store.enqueue(first); + store.enqueue(second); + store.enqueue({ ...first, text: "Updated" }); + store.remove(first); + + const queue = + useWebThreadOutboxStore.getState().queuesByThreadKey[ + webThreadOutboxKey(environmentId, threadId) + ]; + expect(queue?.map((entry) => entry.messageId)).toEqual([second.messageId]); + }); + + it("permits only one dispatcher for a stable message id", () => { + const queued = message(3); + expect(beginWebThreadOutboxDispatch(queued.messageId)).toBe(true); + expect(beginWebThreadOutboxDispatch(queued.messageId)).toBe(false); + finishWebThreadOutboxDispatch(queued.messageId); + expect(beginWebThreadOutboxDispatch(queued.messageId)).toBe(true); + finishWebThreadOutboxDispatch(queued.messageId); + }); + + it("drains only from a ready, connected, unpaused thread", () => { + expect( + shouldDrainWebThreadOutbox({ + phase: "ready", + isSendBusy: false, + isConnecting: false, + environmentUnavailable: false, + paused: false, + }), + ).toBe(true); + expect( + shouldDrainWebThreadOutbox({ + phase: "running", + isSendBusy: false, + isConnecting: false, + environmentUnavailable: false, + paused: false, + }), + ).toBe(false); + expect( + shouldDrainWebThreadOutbox({ + phase: "ready", + isSendBusy: false, + isConnecting: false, + environmentUnavailable: false, + paused: true, + }), + ).toBe(false); + }); + + it("queues active-turn messages only when queue mode or an existing FIFO requires it", () => { + const activeThread = { + isServerThread: true, + phase: "running" as const, + isSendBusy: false, + hasQueuedMessages: false, + }; + + expect( + shouldQueueWebThreadMessage({ + ...activeThread, + activeTurnMessageBehavior: "queue", + }), + ).toBe(true); + expect( + shouldQueueWebThreadMessage({ + ...activeThread, + activeTurnMessageBehavior: "steer", + }), + ).toBe(false); + expect( + shouldQueueWebThreadMessage({ + ...activeThread, + activeTurnMessageBehavior: "steer", + hasQueuedMessages: true, + }), + ).toBe(true); + expect( + shouldQueueWebThreadMessage({ + ...activeThread, + activeTurnMessageBehavior: "queue", + isServerThread: false, + }), + ).toBe(false); + }); +}); diff --git a/apps/web/src/webThreadOutbox.ts b/apps/web/src/webThreadOutbox.ts new file mode 100644 index 00000000000..f823d9edb1f --- /dev/null +++ b/apps/web/src/webThreadOutbox.ts @@ -0,0 +1,242 @@ +import { + CommandId, + EnvironmentId, + MessageId, + ModelSelection, + ProviderInteractionMode, + RuntimeMode, + ThreadId, + type UploadChatAttachment, +} from "@t3tools/contracts"; +import type { ActiveTurnMessageBehavior } from "@t3tools/contracts/settings"; +import { scopedThreadKey, scopeThreadRef } from "@t3tools/client-runtime/environment"; +import * as Schema from "effect/Schema"; +import { create } from "zustand"; + +import { createMemoryStorage, type StateStorage } from "./lib/storage"; + +export const WEB_THREAD_OUTBOX_STORAGE_KEY = "t3code:thread-outbox:v1"; +const WEB_THREAD_OUTBOX_STORAGE_VERSION = 1; + +const QueuedWebImageAttachment = Schema.Struct({ + type: Schema.Literal("image"), + name: Schema.String, + mimeType: Schema.String, + sizeBytes: Schema.Number, + dataUrl: Schema.String, +}); + +const QueuedWebThreadMessageSchema = Schema.Struct({ + environmentId: EnvironmentId, + threadId: ThreadId, + messageId: MessageId, + commandId: CommandId, + text: Schema.String, + attachments: Schema.Array(QueuedWebImageAttachment), + modelSelection: ModelSelection, + runtimeMode: RuntimeMode, + interactionMode: ProviderInteractionMode, + createdAt: Schema.String, +}); + +export interface QueuedWebThreadMessage { + readonly environmentId: EnvironmentId; + readonly threadId: ThreadId; + readonly messageId: MessageId; + readonly commandId: CommandId; + readonly text: string; + readonly attachments: ReadonlyArray; + readonly modelSelection: ModelSelection; + readonly runtimeMode: RuntimeMode; + readonly interactionMode: ProviderInteractionMode; + readonly createdAt: string; +} + +const PersistedWebThreadOutboxState = Schema.Struct({ + queuesByThreadKey: Schema.Record(Schema.String, Schema.Array(QueuedWebThreadMessageSchema)), +}); +const decodePersistedState = Schema.decodeUnknownSync(PersistedWebThreadOutboxState); + +export function webThreadOutboxKey(environmentId: EnvironmentId, threadId: ThreadId): string { + return scopedThreadKey(scopeThreadRef(environmentId, threadId)); +} + +function readQueue( + queues: Record>, + threadKey: string, +): ReadonlyArray { + return Object.hasOwn(queues, threadKey) ? (queues[threadKey] ?? []) : []; +} + +function resolveBaseStorage(): { storage: StateStorage; durable: boolean } { + try { + if (typeof localStorage !== "undefined") { + return { storage: localStorage, durable: true }; + } + } catch { + // Sandboxed browsers can reject access to the localStorage property itself. + } + return { storage: createMemoryStorage(), durable: false }; +} + +const { storage: baseOutboxStorage, durable: storageIsDurable } = resolveBaseStorage(); + +function persistQueues(queues: Record>): { + written: boolean; + durable: boolean; +} { + try { + baseOutboxStorage.setItem( + WEB_THREAD_OUTBOX_STORAGE_KEY, + JSON.stringify({ + version: WEB_THREAD_OUTBOX_STORAGE_VERSION, + state: { queuesByThreadKey: queues }, + }), + ); + return { written: true, durable: storageIsDurable }; + } catch (error) { + console.error("[THREAD-OUTBOX] Could not persist queued messages.", error); + return { written: false, durable: false }; + } +} + +function readPersistedQueues(): Record> | null { + try { + const raw = baseOutboxStorage.getItem(WEB_THREAD_OUTBOX_STORAGE_KEY); + if (typeof raw !== "string" || raw.length === 0) { + return null; + } + const parsed: unknown = JSON.parse(raw); + const state = (parsed as { state?: unknown } | null)?.state; + return state ? decodePersistedState(state).queuesByThreadKey : null; + } catch { + return null; + } +} + +interface WebThreadOutboxState { + readonly queuesByThreadKey: Record>; + readonly pausedMessageIds: Readonly>; + readonly enqueue: (message: QueuedWebThreadMessage) => { durable: boolean }; + readonly remove: (message: QueuedWebThreadMessage) => { durable: boolean }; + readonly pause: (messageId: MessageId) => void; + readonly retry: (messageId: MessageId) => void; +} + +export const useWebThreadOutboxStore = create()((set, get) => ({ + queuesByThreadKey: {}, + pausedMessageIds: {}, + enqueue: (message) => { + const threadKey = webThreadOutboxKey(message.environmentId, message.threadId); + const queues = get().queuesByThreadKey; + const queue = readQueue(queues, threadKey); + const nextQueue = [ + ...queue.filter((candidate) => candidate.messageId !== message.messageId), + message, + ]; + const next = { ...queues, [threadKey]: nextQueue }; + const persisted = persistQueues(next); + // Even when browser storage is blocked or full, retain the message for the + // current session. The caller reports that it is not reload-safe. + set({ queuesByThreadKey: next }); + return { durable: persisted.written && persisted.durable }; + }, + remove: (message) => { + const threadKey = webThreadOutboxKey(message.environmentId, message.threadId); + const queues = get().queuesByThreadKey; + const nextQueue = readQueue(queues, threadKey).filter( + (candidate) => candidate.messageId !== message.messageId, + ); + const next = { ...queues }; + if (nextQueue.length === 0) { + delete next[threadKey]; + } else { + next[threadKey] = nextQueue; + } + const persisted = persistQueues(next); + const pausedMessageIds = { ...get().pausedMessageIds }; + delete pausedMessageIds[message.messageId]; + // The command id is stable, so a removal that fails to persist can only + // cause an idempotent acknowledgement after reload, never a second turn. + set({ queuesByThreadKey: next, pausedMessageIds }); + return { durable: persisted.written && persisted.durable }; + }, + pause: (messageId) => { + set((state) => ({ + pausedMessageIds: { ...state.pausedMessageIds, [messageId]: true }, + })); + }, + retry: (messageId) => { + set((state) => { + if (!state.pausedMessageIds[messageId]) { + return state; + } + const pausedMessageIds = { ...state.pausedMessageIds }; + delete pausedMessageIds[messageId]; + return { pausedMessageIds }; + }); + }, +})); + +export const EMPTY_WEB_THREAD_OUTBOX_QUEUE: ReadonlyArray = []; + +{ + const persisted = readPersistedQueues(); + if (persisted) { + useWebThreadOutboxStore.setState({ queuesByThreadKey: persisted }); + } +} + +const dispatchingMessageIds = new Set(); + +export function beginWebThreadOutboxDispatch(messageId: MessageId): boolean { + if (dispatchingMessageIds.has(messageId)) { + return false; + } + dispatchingMessageIds.add(messageId); + return true; +} + +export function finishWebThreadOutboxDispatch(messageId: MessageId): void { + dispatchingMessageIds.delete(messageId); +} + +export function shouldDrainWebThreadOutbox(input: { + readonly phase: "disconnected" | "connecting" | "ready" | "running"; + readonly isSendBusy: boolean; + readonly isConnecting: boolean; + readonly environmentUnavailable: boolean; + readonly paused: boolean; +}): boolean { + return ( + input.phase === "ready" && + !input.isSendBusy && + !input.isConnecting && + !input.environmentUnavailable && + !input.paused + ); +} + +export function shouldQueueWebThreadMessage(input: { + readonly activeTurnMessageBehavior: ActiveTurnMessageBehavior; + readonly hasQueuedMessages: boolean; + readonly isSendBusy: boolean; + readonly isServerThread: boolean; + readonly phase: "disconnected" | "connecting" | "ready" | "running"; +}): boolean { + return ( + input.isServerThread && + (input.hasQueuedMessages || + (input.activeTurnMessageBehavior === "queue" && + (input.phase === "running" || input.isSendBusy))) + ); +} + +export function writeWebThreadOutboxStorageForTest(raw: string): void { + baseOutboxStorage.setItem(WEB_THREAD_OUTBOX_STORAGE_KEY, raw); + useWebThreadOutboxStore.setState({ + queuesByThreadKey: readPersistedQueues() ?? {}, + pausedMessageIds: {}, + }); + dispatchingMessageIds.clear(); +} diff --git a/packages/contracts/src/settings.test.ts b/packages/contracts/src/settings.test.ts index 5bd22e95f20..d9335cef4d6 100644 --- a/packages/contracts/src/settings.test.ts +++ b/packages/contracts/src/settings.test.ts @@ -67,6 +67,22 @@ describe("ClientSettings environment identification", () => { }); }); +describe("ClientSettings messages while working", () => { + it("defaults to steering and accepts both delivery behaviors", () => { + expect(decodeClientSettings({}).activeTurnMessageBehavior).toBe("steer"); + expect( + decodeClientSettingsPatch({ activeTurnMessageBehavior: "steer" }).activeTurnMessageBehavior, + ).toBe("steer"); + expect( + decodeClientSettingsPatch({ activeTurnMessageBehavior: "queue" }).activeTurnMessageBehavior, + ).toBe("queue"); + }); + + it("rejects unsupported delivery behaviors", () => { + expect(() => decodeClientSettingsPatch({ activeTurnMessageBehavior: "send-later" })).toThrow(); + }); +}); + describe("ClientSettings sidebar v2", () => { it("defaults the beta off with a three-day auto-settle threshold", () => { const settings = decodeClientSettings({}); diff --git a/packages/contracts/src/settings.ts b/packages/contracts/src/settings.ts index cbb547b95fb..ad6f41ae437 100644 --- a/packages/contracts/src/settings.ts +++ b/packages/contracts/src/settings.ts @@ -102,6 +102,9 @@ export const DEFAULT_TERMINAL_FONT_SIZE: TerminalFontSize = 12; export const EnvironmentIdentificationMode = Schema.Literals(["artwork", "pill", "none"]); export type EnvironmentIdentificationMode = typeof EnvironmentIdentificationMode.Type; export const DEFAULT_ENVIRONMENT_IDENTIFICATION_MODE: EnvironmentIdentificationMode = "artwork"; +export const ActiveTurnMessageBehavior = Schema.Literals(["steer", "queue"]); +export type ActiveTurnMessageBehavior = typeof ActiveTurnMessageBehavior.Type; +export const DEFAULT_ACTIVE_TURN_MESSAGE_BEHAVIOR: ActiveTurnMessageBehavior = "steer"; /** * A user-chosen font family (a single name or a comma-separated list). Empty @@ -111,6 +114,9 @@ export const FontFamilyPreference = Schema.String.check(Schema.isMaxLength(200)) export type FontFamilyPreference = typeof FontFamilyPreference.Type; export const ClientSettingsSchema = Schema.Struct({ + activeTurnMessageBehavior: ActiveTurnMessageBehavior.pipe( + Schema.withDecodingDefault(Effect.succeed(DEFAULT_ACTIVE_TURN_MESSAGE_BEHAVIOR)), + ), autoOpenPlanSidebar: Schema.Boolean.pipe(Schema.withDecodingDefault(Effect.succeed(false))), confirmThreadArchive: Schema.Boolean.pipe(Schema.withDecodingDefault(Effect.succeed(false))), confirmThreadDelete: Schema.Boolean.pipe(Schema.withDecodingDefault(Effect.succeed(true))), @@ -747,6 +753,7 @@ export const ServerSettingsPatch = Schema.Struct({ export type ServerSettingsPatch = typeof ServerSettingsPatch.Type; export const ClientSettingsPatch = Schema.Struct({ + activeTurnMessageBehavior: Schema.optionalKey(ActiveTurnMessageBehavior), autoOpenPlanSidebar: Schema.optionalKey(Schema.Boolean), confirmThreadArchive: Schema.optionalKey(Schema.Boolean), confirmThreadDelete: Schema.optionalKey(Schema.Boolean), From 56e5d17e0e3b0ff48f14bcaecc451e65c69d4233 Mon Sep 17 00:00:00 2001 From: ClapFy Date: Wed, 5 Aug 2026 13:06:49 +0300 Subject: [PATCH 2/2] fix: address active-turn message review feedback --- apps/mobile/src/state/preferences.test.ts | 27 +++++++++ apps/mobile/src/state/preferences.ts | 41 +++++++++++++- apps/mobile/src/state/thread-outbox-model.ts | 14 +++++ apps/mobile/src/state/thread-outbox.test.ts | 23 ++++++++ .../src/state/use-thread-composer-state.ts | 10 +++- .../src/state/use-thread-outbox-drain.ts | 10 +++- apps/web/src/webThreadOutbox.test.ts | 52 ++++++++++++++++++ apps/web/src/webThreadOutbox.ts | 55 ++++++++++++++++++- 8 files changed, 224 insertions(+), 8 deletions(-) diff --git a/apps/mobile/src/state/preferences.test.ts b/apps/mobile/src/state/preferences.test.ts index c53594eb230..6bb638147eb 100644 --- a/apps/mobile/src/state/preferences.test.ts +++ b/apps/mobile/src/state/preferences.test.ts @@ -24,6 +24,7 @@ vi.mock("../lib/runtime", async () => { import type { Preferences } from "../persistence/mobile-preferences"; import { + awaitActiveTurnMessageBehavior, createMobilePreferencesState, MobilePreferencesLoadError, MobilePreferencesSaveError, @@ -62,6 +63,32 @@ function makePreferencesState( } describe("mobile preferences state", () => { + it("waits for the persisted active-turn behavior before sending", async () => { + const pendingLoad = deferred(); + const state = makePreferencesState({ + load: Effect.promise(() => pendingLoad.promise), + savePatch: (patch) => Effect.succeed(patch), + }); + const registry = AtomRegistry.make(); + const unmount = registry.mount(state.preferencesAtom); + + let settled = false; + const behaviorPromise = awaitActiveTurnMessageBehavior(registry, state.preferencesAtom).then( + (behavior) => { + settled = true; + return behavior; + }, + ); + await Promise.resolve(); + expect(settled).toBe(false); + + pendingLoad.resolve({ activeTurnMessageBehavior: "queue" }); + await expect(behaviorPromise).resolves.toBe("queue"); + + unmount(); + registry.dispose(); + }); + it.effect("shares one preference load across consumers", () => Effect.gen(function* () { const load = vi.fn(() => Promise.resolve({ baseFontSize: 17 })); diff --git a/apps/mobile/src/state/preferences.ts b/apps/mobile/src/state/preferences.ts index d173cf55be5..f303cd60677 100644 --- a/apps/mobile/src/state/preferences.ts +++ b/apps/mobile/src/state/preferences.ts @@ -1,6 +1,8 @@ import * as Effect from "effect/Effect"; -import { AsyncResult, Atom } from "effect/unstable/reactivity"; +import { AsyncResult, Atom, AtomRegistry } from "effect/unstable/reactivity"; +import { DEFAULT_ACTIVE_TURN_MESSAGE_BEHAVIOR } from "@t3tools/contracts/settings"; +import type { ActiveTurnMessageBehavior } from "@t3tools/contracts/settings"; import { MobilePreferencesStore, type Preferences } from "../persistence/mobile-preferences"; import * as Runtime from "../lib/runtime"; @@ -122,3 +124,40 @@ export const mobilePreferencesState = createMobilePreferencesState(mobilePrefere export const mobilePreferencesAtom = mobilePreferencesState.preferencesAtom; export const updateMobilePreferencesAtom = mobilePreferencesState.updatePreferencesAtom; + +function settledActiveTurnMessageBehavior( + result: AsyncResult.AsyncResult, +): ActiveTurnMessageBehavior | null { + if (result.waiting) { + return null; + } + return AsyncResult.isSuccess(result) + ? (result.value.activeTurnMessageBehavior ?? DEFAULT_ACTIVE_TURN_MESSAGE_BEHAVIOR) + : DEFAULT_ACTIVE_TURN_MESSAGE_BEHAVIOR; +} + +/** + * Reads the send behavior from the settled preference snapshot. A composer can + * render before the device preference read finishes, so capturing its + * render-time fallback would steer a message that the user intended to queue. + */ +export function awaitActiveTurnMessageBehavior( + registry: AtomRegistry.AtomRegistry, + preferencesAtom: Atom.Atom>, +): Promise { + const current = settledActiveTurnMessageBehavior(registry.get(preferencesAtom)); + if (current !== null) { + return Promise.resolve(current); + } + + return new Promise((resolve) => { + const unsubscribe = registry.subscribe(preferencesAtom, (result) => { + const behavior = settledActiveTurnMessageBehavior(result); + if (behavior === null) { + return; + } + unsubscribe(); + resolve(behavior); + }); + }); +} diff --git a/apps/mobile/src/state/thread-outbox-model.ts b/apps/mobile/src/state/thread-outbox-model.ts index 7ec8b77dd7b..35343ca7998 100644 --- a/apps/mobile/src/state/thread-outbox-model.ts +++ b/apps/mobile/src/state/thread-outbox-model.ts @@ -158,6 +158,20 @@ export function threadOutboxRetryDelayMs(attempt: number): number { export type ThreadOutboxDeliveryAction = "wait" | "remove" | "send"; +export function shouldDeferConfirmedThreadOutboxDelivery(input: { + readonly deliveryAction: ThreadOutboxDeliveryAction; + readonly isCreation: boolean; + readonly threadBusy: boolean; + readonly activeTurnMessageBehavior?: ActiveTurnMessageBehaviorType; +}): boolean { + return ( + input.deliveryAction === "send" && + !input.isCreation && + input.threadBusy && + input.activeTurnMessageBehavior !== "steer" + ); +} + export function resolveThreadOutboxDeliveryAction(input: { readonly isCreation: boolean; readonly threadExists: boolean; diff --git a/apps/mobile/src/state/thread-outbox.test.ts b/apps/mobile/src/state/thread-outbox.test.ts index f0d42d82c51..8dd8aaad356 100644 --- a/apps/mobile/src/state/thread-outbox.test.ts +++ b/apps/mobile/src/state/thread-outbox.test.ts @@ -18,6 +18,7 @@ import { resolveThreadOutboxDeliveryAction, resolveThreadOutboxFailureAction, resolveQueuedThreadSettings, + shouldDeferConfirmedThreadOutboxDelivery, shouldRetryThreadOutboxDelivery, threadOutboxRetryDelayMs, type QueuedThreadMessage, @@ -512,6 +513,28 @@ describe("thread outbox", () => { ).toBe("send"); }); + it("keeps steer delivery eligible when a thread becomes busy during persistence", () => { + const input = { + deliveryAction: "send" as const, + isCreation: false, + threadBusy: true, + }; + + expect(shouldDeferConfirmedThreadOutboxDelivery(input)).toBe(true); + expect( + shouldDeferConfirmedThreadOutboxDelivery({ + ...input, + activeTurnMessageBehavior: "queue", + }), + ).toBe(true); + expect( + shouldDeferConfirmedThreadOutboxDelivery({ + ...input, + activeTurnMessageBehavior: "steer", + }), + ).toBe(false); + }); + it("sends queued creations once connected and live, removing already-created ones", () => { expect( resolveThreadOutboxDeliveryAction({ diff --git a/apps/mobile/src/state/use-thread-composer-state.ts b/apps/mobile/src/state/use-thread-composer-state.ts index 938cd29d479..8308773806b 100644 --- a/apps/mobile/src/state/use-thread-composer-state.ts +++ b/apps/mobile/src/state/use-thread-composer-state.ts @@ -43,7 +43,7 @@ import { useSelectedThreadDetail } from "../state/use-thread-detail"; import { useThreadSelection } from "../state/use-thread-selection"; import { enqueueThreadOutboxMessage } from "./thread-outbox"; import { useThreadOutboxMessages } from "./use-thread-outbox"; -import { mobilePreferencesAtom } from "./preferences"; +import { awaitActiveTurnMessageBehavior, mobilePreferencesAtom } from "./preferences"; export function appendReviewCommentToDraft(input: { readonly environmentId: EnvironmentId; @@ -145,6 +145,10 @@ export function useThreadComposerState() { return null; } + const sendBehavior = await awaitActiveTurnMessageBehavior( + appAtomRegistry, + mobilePreferencesAtom, + ); const threadKey = scopedThreadKey(selectedThreadShell.environmentId, selectedThreadShell.id); const draft = getComposerDraftSnapshot(threadKey); const thread = selectedThreadDetail ?? selectedThreadShell; @@ -171,7 +175,7 @@ export function useThreadComposerState() { modelSelection: draft.modelSelection ?? thread.modelSelection, runtimeMode: draft.runtimeMode ?? thread.runtimeMode, interactionMode: draft.interactionMode ?? thread.interactionMode, - activeTurnMessageBehavior, + activeTurnMessageBehavior: sendBehavior, createdAt: metadata.createdAt, }); clearComposerDraftContent(threadKey); @@ -187,7 +191,7 @@ export function useThreadComposerState() { ); }); return messageId; - }, [activeTurnMessageBehavior, selectedThreadDetail, selectedThreadShell]); + }, [selectedThreadDetail, selectedThreadShell]); const onChangeDraftMessage = useCallback( (value: string) => { diff --git a/apps/mobile/src/state/use-thread-outbox-drain.ts b/apps/mobile/src/state/use-thread-outbox-drain.ts index 8a7f01875ae..eab20a63df5 100644 --- a/apps/mobile/src/state/use-thread-outbox-drain.ts +++ b/apps/mobile/src/state/use-thread-outbox-drain.ts @@ -32,6 +32,7 @@ import { resolveThreadOutboxDeliveryAction, resolveThreadOutboxFailureAction, resolveQueuedThreadSettings, + shouldDeferConfirmedThreadOutboxDelivery, threadOutboxRetryDelayMs, type QueuedThreadCreation, type QueuedThreadMessage, @@ -376,7 +377,14 @@ export function useThreadOutboxDrain(): void { ); const freshThreadBusy = freshThread?.session?.status === "running" || freshThread?.session?.status === "starting"; - if (deliveryAction === "send" && creation === undefined && freshThreadBusy) { + if ( + shouldDeferConfirmedThreadOutboxDelivery({ + deliveryAction, + isCreation: creation !== undefined, + threadBusy: freshThreadBusy, + activeTurnMessageBehavior: nextQueuedMessage.activeTurnMessageBehavior, + }) + ) { return true; } return deliveryAction === "remove" diff --git a/apps/web/src/webThreadOutbox.test.ts b/apps/web/src/webThreadOutbox.test.ts index efdee665fbc..568d09cf464 100644 --- a/apps/web/src/webThreadOutbox.test.ts +++ b/apps/web/src/webThreadOutbox.test.ts @@ -43,6 +43,17 @@ function resetOutbox(): void { writeWebThreadOutboxStorageForTest(""); } +function persistedOutbox(messages: ReadonlyArray): string { + return JSON.stringify({ + version: 1, + state: { + queuesByThreadKey: { + [webThreadOutboxKey(environmentId, threadId)]: messages, + }, + }, + }); +} + afterEach(resetOutbox); describe("web thread outbox", () => { @@ -78,6 +89,47 @@ describe("web thread outbox", () => { expect(queue?.map((entry) => entry.messageId)).toEqual([second.messageId]); }); + it("merges another tab's durable messages before enqueueing", () => { + const first = message(1); + const second = message(2); + const third = message(3); + const store = useWebThreadOutboxStore.getState(); + store.enqueue(first); + + writeWebThreadOutboxStorageForTest(persistedOutbox([first, second]), { + syncStore: false, + }); + store.enqueue(third); + + const queue = + useWebThreadOutboxStore.getState().queuesByThreadKey[ + webThreadOutboxKey(environmentId, threadId) + ]; + expect(queue?.map((entry) => entry.messageId)).toEqual([ + first.messageId, + second.messageId, + third.messageId, + ]); + }); + + it("preserves another tab's messages when removing a delivered message", () => { + const first = message(1); + const second = message(2); + const store = useWebThreadOutboxStore.getState(); + store.enqueue(first); + + writeWebThreadOutboxStorageForTest(persistedOutbox([first, second]), { + syncStore: false, + }); + store.remove(first); + + const queue = + useWebThreadOutboxStore.getState().queuesByThreadKey[ + webThreadOutboxKey(environmentId, threadId) + ]; + expect(queue?.map((entry) => entry.messageId)).toEqual([second.messageId]); + }); + it("permits only one dispatcher for a stable message id", () => { const queued = message(3); expect(beginWebThreadOutboxDispatch(queued.messageId)).toBe(true); diff --git a/apps/web/src/webThreadOutbox.ts b/apps/web/src/webThreadOutbox.ts index f823d9edb1f..c0752bc7fee 100644 --- a/apps/web/src/webThreadOutbox.ts +++ b/apps/web/src/webThreadOutbox.ts @@ -68,6 +68,30 @@ function readQueue( return Object.hasOwn(queues, threadKey) ? (queues[threadKey] ?? []) : []; } +function mergeQueues( + ...sources: ReadonlyArray>> +): Record> { + const merged: Record> = {}; + const threadKeys = new Set(sources.flatMap((source) => Object.keys(source))); + for (const threadKey of threadKeys) { + const messagesById = new Map(); + for (const source of sources) { + for (const message of readQueue(source, threadKey)) { + messagesById.set(message.messageId, message); + } + } + const queue = [...messagesById.values()].sort( + (left, right) => + left.createdAt.localeCompare(right.createdAt) || + String(left.messageId).localeCompare(String(right.messageId)), + ); + if (queue.length > 0) { + merged[threadKey] = queue; + } + } + return merged; +} + function resolveBaseStorage(): { storage: StateStorage; durable: boolean } { try { if (typeof localStorage !== "undefined") { @@ -114,6 +138,13 @@ function readPersistedQueues(): Record>, +): Record> { + const persisted = readPersistedQueues(); + return persisted === null ? queues : mergeQueues(queues, persisted); +} + interface WebThreadOutboxState { readonly queuesByThreadKey: Record>; readonly pausedMessageIds: Readonly>; @@ -128,7 +159,10 @@ export const useWebThreadOutboxStore = create()((set, get) pausedMessageIds: {}, enqueue: (message) => { const threadKey = webThreadOutboxKey(message.environmentId, message.threadId); - const queues = get().queuesByThreadKey; + // Another tab may have updated the shared outbox since this store last + // rendered. Merge its durable snapshot before applying this mutation so a + // full-key localStorage write cannot discard the other tab's messages. + const queues = mergeWithPersistedQueues(get().queuesByThreadKey); const queue = readQueue(queues, threadKey); const nextQueue = [ ...queue.filter((candidate) => candidate.messageId !== message.messageId), @@ -143,7 +177,7 @@ export const useWebThreadOutboxStore = create()((set, get) }, remove: (message) => { const threadKey = webThreadOutboxKey(message.environmentId, message.threadId); - const queues = get().queuesByThreadKey; + const queues = mergeWithPersistedQueues(get().queuesByThreadKey); const nextQueue = readQueue(queues, threadKey).filter( (candidate) => candidate.messageId !== message.messageId, ); @@ -187,6 +221,15 @@ export const EMPTY_WEB_THREAD_OUTBOX_QUEUE: ReadonlyArray { + if (event.key !== WEB_THREAD_OUTBOX_STORAGE_KEY) { + return; + } + useWebThreadOutboxStore.setState({ queuesByThreadKey: readPersistedQueues() ?? {} }); + }); +} + const dispatchingMessageIds = new Set(); export function beginWebThreadOutboxDispatch(messageId: MessageId): boolean { @@ -232,8 +275,14 @@ export function shouldQueueWebThreadMessage(input: { ); } -export function writeWebThreadOutboxStorageForTest(raw: string): void { +export function writeWebThreadOutboxStorageForTest( + raw: string, + options?: { readonly syncStore?: boolean }, +): void { baseOutboxStorage.setItem(WEB_THREAD_OUTBOX_STORAGE_KEY, raw); + if (options?.syncStore === false) { + return; + } useWebThreadOutboxStore.setState({ queuesByThreadKey: readPersistedQueues() ?? {}, pausedMessageIds: {},