Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
213 changes: 213 additions & 0 deletions apps/web/src/components/BackgroundQueuedMessages.test.tsx
Original file line number Diff line number Diff line change
@@ -0,0 +1,213 @@
import { act, create, type ReactTestRenderer } from "react-test-renderer";
import { afterEach, beforeEach, describe, expect, it, vi } from "vite-plus/test";
import { EnvironmentId, ThreadId, ProviderInstanceId } from "@t3tools/contracts";
import { scopeThreadRef, scopedThreadKey } from "@t3tools/client-runtime/environment";
import { useComposerDraftStore } from "../composerDraftStore";
import { useQueuedMessageStore } from "../queuedMessageStore";
import { BackgroundQueuedMessages } from "./BackgroundQueuedMessages";

const io = vi.hoisted(() => ({
threads: new Map<string, unknown>(),
prepare: vi.fn(),
start: vi.fn(),
connected: true,
shellPresent: true,
config: {},
}));
vi.mock("../state/entities", () => ({
useThread: (ref: { threadId: string }) => io.threads.get(ref.threadId),
useServerConfigs: () => new Map([[EnvironmentId.make("env"), io.config]]),
useThreadShell: () => (io.shellPresent ? {} : null),
}));
vi.mock("../state/threads", () => ({ useEnvironmentThread: () => ({ status: "live" }) }));
vi.mock("../state/environments", () => ({
useEnvironments: () => ({
environments: [
{
environmentId: EnvironmentId.make("env"),
connection: { phase: io.connected ? "connected" : "reconnecting" },
},
],
}),
}));
vi.mock("../hooks/useSettings", () => ({ useClientSettingsHydrated: () => true }));
vi.mock("../lib/sendBackgroundQueuedMessage", () => ({
sendBackgroundQueuedMessage: async (...args: unknown[]) => io.start(...args),
}));

const ref = scopeThreadRef(EnvironmentId.make("env"), ThreadId.make("a"));
const key = scopedThreadKey(ref);
let root: ReactTestRenderer | undefined;
const render = async (activeThreadKey: string | null) => {
await act(async () => {
const element = <BackgroundQueuedMessages activeThreadKey={activeThreadKey} />;
if (root) root.update(element);
else root = create(element);
});
};
function enqueue(prompt = "follow up") {
return useQueuedMessageStore.getState().enqueue(key, {
prompt,
images: [],
files: [],
terminalContexts: [],
previewAnnotations: [],
reviewComments: [],
submissionIntent: "foreground",
queuedAfterToolActivityId: null,
createdAt: "2026-09-22T00:00:00Z",
sendOptions: {
modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5" },
runtimeMode: "full-access",
interactionMode: "default",
promptEffort: null,
},
});
}
beforeEach(() => {
(globalThis as { IS_REACT_ACT_ENVIRONMENT?: boolean }).IS_REACT_ACT_ENVIRONMENT = true;
useQueuedMessageStore.setState({
queuesByThreadKey: {},
backgroundSendsByThreadKey: {},
drainGeneration: 0,
});
io.connected = true;
io.shellPresent = true;
io.threads.set("a", { session: { status: "running" }, activities: [] });
io.start.mockReset().mockImplementation(async (_ref, message, _options, canSend) => {
if (!canSend()) return false;
return Boolean(useQueuedMessageStore.getState().take(key, message.id, null));
});
});
afterEach(async () => {
if (root) await act(() => root!.unmount());
root = undefined;
});
describe("background queued messages", () => {
it("sends a completed thread's queue after navigating away, without selecting it again", async () => {
enqueue();
await render(key);
await render("another-thread");
expect(io.start).not.toHaveBeenCalled();
io.threads.set("a", { session: { status: "ready" }, activities: [] });
await render("another-thread");
expect(io.start).toHaveBeenCalledOnce();
expect(io.start.mock.calls[0]?.[0]).toEqual(ref);
expect(useQueuedMessageStore.getState().queuesByThreadKey[key]).toBeUndefined();
});
it("leaves selected and disconnected threads to wait", async () => {
enqueue();
io.threads.set("a", { session: { status: "ready" }, activities: [] });
await render(key);
expect(io.start).not.toHaveBeenCalled();
io.connected = false;
await render(null);
expect(io.start).not.toHaveBeenCalled();
io.connected = true;
io.shellPresent = true;
await render(null);
expect(io.start).toHaveBeenCalledOnce();
});
it("advances after a text-only turn even when there is no new tool activity", async () => {
enqueue("first");
enqueue("second");
io.threads.set("a", { session: { status: "ready" }, activities: [] });
await render(null);
expect(io.start).toHaveBeenCalledOnce();
io.threads.set("a", { session: { status: "running" }, activities: [] });
await render(null);
expect(io.start).toHaveBeenCalledOnce();
io.threads.set("a", { session: { status: "ready" }, activities: [] });
await render(null);
expect(io.start).toHaveBeenCalledTimes(2);
});
it("waits for checkpoint rewind to finish", async () => {
enqueue();
useComposerDraftStore.setState({ rewindingThreadKeys: new Set([key]) });
io.threads.set("a", { session: { status: "ready" }, activities: [] });
await render(null);
expect(io.start).not.toHaveBeenCalled();
await act(() => {
useComposerDraftStore.setState({ rewindingThreadKeys: new Set() });
});
expect(io.start).toHaveBeenCalledOnce();
});
it.each(["approval.requested", "user-input.requested"])(
"waits for %s to resolve",
async (kind) => {
enqueue();
io.threads.set("a", {
session: { status: "ready" },
activities: [
{
id: "request",
kind,
createdAt: "2026-09-22T00:00:00Z",
payload: {
requestId: "req-1",
questions: [
{
id: "q1",
header: "name",
question: "name?",
options: [],
allowCustomAnswer: true,
multiSelect: false,
},
],
},
},
],
});
await render(null);
expect(io.start).not.toHaveBeenCalled();
io.threads.set("a", { session: { status: "ready" }, activities: [] });
await render(null);
expect(io.start).toHaveBeenCalledOnce();
},
);
it("does not miss a boundary that arrives while dispatch is awaiting its receipt", async () => {
enqueue("first");
enqueue("second");
let release!: () => void;
const receipt = new Promise<void>((resolve) => {
release = resolve;
});
io.start.mockImplementationOnce(async (_ref, message) => {
useQueuedMessageStore.getState().take(key, message.id, null);
await receipt;
return true;
});
io.threads.set("a", {
session: { status: "ready" },
activities: [],
latestTurn: { turnId: "old" },
});
await render(null);
io.threads.set("a", {
session: { status: "ready" },
activities: [],
latestTurn: { turnId: "new" },
});
await render(null);
await act(async () => release());
expect(io.start).toHaveBeenCalledTimes(2);
});
it("retries when the shell arrives after the detail", async () => {
enqueue();
io.shellPresent = false;
io.threads.set("a", { session: { status: "ready" }, activities: [] });
await render(null);
expect(io.start).not.toHaveBeenCalled();
io.shellPresent = true;
await render(null);
expect(io.start).toHaveBeenCalledOnce();
});
it("does not send a held message", async () => {
const message = enqueue();
useQueuedMessageStore.getState().holdAtFront(key, message);
io.threads.set("a", { session: { status: "ready" }, activities: [] });
await render(null);
expect(io.start).not.toHaveBeenCalled();
});
});
139 changes: 139 additions & 0 deletions apps/web/src/components/BackgroundQueuedMessages.tsx
Original file line number Diff line number Diff line change
@@ -0,0 +1,139 @@
import { useEffect, useLayoutEffect, useMemo, useRef, useState } from "react";
import { useComposerDraftStore } from "../composerDraftStore";
import { useShallow } from "zustand/react/shallow";
import { parseScopedThreadKey, scopedThreadKey } from "@t3tools/client-runtime/environment";
import { useParams } from "@tanstack/react-router";
import { resolveThreadRouteTarget } from "../threadRoutes";
import { derivePendingRequests } from "@t3tools/client-runtime/pending-requests";
import {
useQueuedMessageStore,
useQueuedMessages,
isQueuedMessageDue,
latestCompletedToolActivityId,
} from "../queuedMessageStore";
import { useThread, useThreadShell, useServerConfigs } from "../state/entities";
import { useEnvironmentThread } from "../state/threads";
import { useEnvironments } from "../state/environments";
import { useClientSettingsHydrated } from "../hooks/useSettings";
import { derivePhase } from "../session-logic";
import { sendBackgroundQueuedMessage } from "../lib/sendBackgroundQueuedMessage";

export function BackgroundQueueCoordinator() {
const target = useParams({ strict: false, select: resolveThreadRouteTarget });
return (
<BackgroundQueuedMessages
activeThreadKey={target?.kind === "server" ? scopedThreadKey(target.threadRef) : null}
/>
);
}

/** Subscribe only to threads with queued work, independently of the selected route. */
export function BackgroundQueuedMessages({ activeThreadKey }: { activeThreadKey: string | null }) {
const keys = useQueuedMessageStore(
useShallow((state) => [
...new Set([
...Object.keys(state.queuesByThreadKey),
...Object.keys(state.backgroundSendsByThreadKey),
]),
]),
);
return keys.map((threadKey) => (
<BackgroundThreadQueue
key={threadKey}
threadKey={threadKey}
active={activeThreadKey === threadKey}
/>
));
}

function BackgroundThreadQueue({ threadKey, active }: { threadKey: string; active: boolean }) {
const ref = useMemo(() => parseScopedThreadKey(threadKey), [threadKey]);
const thread = useThread(ref);
const shell = useThreadShell(ref);
const sending = useQueuedMessageStore(
(state) => state.backgroundSendsByThreadKey[threadKey] !== undefined,
);
const detail = useEnvironmentThread(ref?.environmentId ?? null, ref?.threadId ?? null);
const { environments } = useEnvironments();
const configs = useServerConfigs();
const hydrated = useClientSettingsHydrated();
const message = useQueuedMessages(threadKey)[0];
const rewinding = useComposerDraftStore((state) => state.rewindingThreadKeys.has(threadKey));
const config = ref ? configs.get(ref.environmentId) : undefined;
const phase = derivePhase(thread?.session ?? null);
const latestToolActivityId = latestCompletedToolActivityId(thread?.activities ?? []);
const pending = derivePendingRequests(thread?.activities ?? []);
const connected = environments.some(
(environment) =>
environment.environmentId === ref?.environmentId &&
environment.connection.phase === "connected",
);
const blocked =
(active && !sending) ||
rewinding ||
!connected ||
!hydrated ||
!thread ||
!shell ||
detail.status !== "live" ||
pending.approvals.length > 0 ||
pending.userInputs.length > 0;
const busy = useRef(false);
const [completedAttempt, setCompletedAttempt] = useState<object | null>(null);
const lastSentBoundary = useRef<string | null>(null);
const boundary = JSON.stringify([phase, latestToolActivityId, thread?.latestTurn?.turnId]);
const live = useRef({ blocked, phase, latestToolActivityId, message, config, boundary });
const mounted = useRef(false);
useLayoutEffect(() => {
live.current = { blocked, phase, latestToolActivityId, message, config, boundary };
}, [blocked, phase, latestToolActivityId, message, config, boundary]);
useEffect(() => {
mounted.current = true;
return () => {
mounted.current = false;
};
}, []);
useEffect(() => {
if (lastSentBoundary.current !== boundary) lastSentBoundary.current = null;
if (
!ref ||
!config ||
!message?.sendOptions ||
blocked ||
busy.current ||
completedAttempt === live.current ||
lastSentBoundary.current === boundary
)
return;
const canSend = () => {
const current = live.current;
return (
mounted.current &&
!current.blocked &&
isQueuedMessageDue({
message,
phase: current.phase,
latestToolActivityId: current.latestToolActivityId,
})
);
};
if (!canSend()) return;
busy.current = true;
const attempt = live.current;
void sendBackgroundQueuedMessage(
ref,
message,
message.sendOptions,
canSend,
() => live.current.latestToolActivityId,
)
.then((sent) => {
if (sent) lastSentBoundary.current = boundary;
})
.finally(() => {
busy.current = false;
if (mounted.current) setCompletedAttempt(attempt);
});
}, [ref, message, blocked, boundary, config, completedAttempt]);
return null;
}
20 changes: 19 additions & 1 deletion apps/web/src/components/ChatView.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -3201,7 +3201,7 @@ export default function ChatView(props: ChatViewProps) {
localDispatchStartedAt,
latestUserMessageAt,
isPreparingWorktree: isLocallyPreparingWorktree,
isSendBusy,
isSendBusy: isLocalSendBusy,
backgroundSubmissionPending,
} = useLocalDispatchState({
activeThread,
Expand All @@ -3211,6 +3211,11 @@ export default function ChatView(props: ChatViewProps) {
activePendingUserInput: activePendingUserInput?.requestId ?? null,
threadError,
});
const isBackgroundQueueSending = useQueuedMessageStore(
(state) =>
activeThreadKey !== null && state.backgroundSendsByThreadKey[activeThreadKey] !== undefined,
);
const isSendBusy = isLocalSendBusy || isBackgroundQueueSending;
const optimisticCompactionMessage = optimisticUserMessages.at(-1);
const pendingCompactionMessage =
isSendBusy &&
Expand Down Expand Up @@ -7651,6 +7656,19 @@ export default function ChatView(props: ChatViewProps) {
previewAnnotations: [...composerPreviewAnnotations],
reviewComments: [...composerReviewComments],
submissionIntent,
sendOptions: {
modelSelection: ctxSelectedModelSelection,
runtimeMode,
interactionMode: sendInteractionMode,
promptEffort: resolvePromptInjectedEffort(
getProviderModelCapabilities(
ctxSelectedProviderModels,
ctxSelectedModel,
ctxSelectedProvider,
),
ctxSelectedPromptEffort,
),
},
queuedAfterToolActivityId: latestCompletedToolActivityId(threadActivities),
createdAt: new Date().toISOString(),
});
Expand Down
Loading
Loading