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
135 changes: 135 additions & 0 deletions apps/web/src/components/BackgroundQueuedMessages.tsx
Original file line number Diff line number Diff line change
@@ -0,0 +1,135 @@
import { useEffect, useLayoutEffect, useMemo, useRef } from "react";
import { useParams } from "@tanstack/react-router";
import { useShallow } from "zustand/react/shallow";
import { parseScopedThreadKey } from "@t3tools/client-runtime/environment";
import { derivePendingRequests } from "@t3tools/client-runtime/pending-requests";

import { useComposerDraftStore } from "../composerDraftStore";
import { useClientSettingsHydrated } from "../hooks/useSettings";
import { activeQueueThreadKey } from "../lib/backgroundQueueActiveThread";
import { sendBackgroundQueuedMessage } from "../lib/sendBackgroundQueuedMessage";
import {
isQueuedMessageDue,
latestCompletedToolActivityId,
useQueuedMessageStore,
useQueuedMessages,
} from "../queuedMessageStore";
import { derivePhase } from "../session-logic";
import { useEnvironments } from "../state/environments";
import { useServerConfigs, useThread, useThreadShell } from "../state/entities";
import { useEnvironmentThread } from "../state/threads";
import { resolveThreadRouteTarget } from "../threadRoutes";

export function BackgroundQueueCoordinator() {
const target = useParams({ strict: false, select: resolveThreadRouteTarget });
const draftId = target?.kind === "draft" ? target.draftId : null;
const draftThread = useComposerDraftStore((store) =>
draftId ? store.getDraftSession(draftId) : null,
);
return (
<BackgroundQueuedMessages activeThreadKey={activeQueueThreadKey(target ?? null, draftThread)} />
);
}

/** Watches every thread that has a queue, including ones the chat view is not showing. */
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 provider = config?.providers.find(
(entry) => entry.instanceId === message?.sendOptions?.modelSelection.instanceId,
);
const providerReady =
provider?.enabled === true &&
provider.installed === true &&
provider.availability !== "unavailable" &&
provider.status === "ready";
const blocked =
(active && !sending) ||
rewinding ||
!connected ||
!hydrated ||
!thread ||
!shell ||
detail.status !== "live" ||
pending.approvals.length > 0 ||
pending.userInputs.length > 0 ||
!providerReady;
const attempted = useRef<string | null>(null);
const live = useRef({ blocked, phase, latestToolActivityId });
useLayoutEffect(() => {
live.current = { blocked, phase, latestToolActivityId };
}, [blocked, phase, latestToolActivityId]);

useEffect(() => {
if (!ref || !message?.sendOptions || blocked) return;
if (
!isQueuedMessageDue({
message,
phase,
latestToolActivityId,
})
) {
return;
}
const boundary = `${message.id}\0${phase}\0${latestToolActivityId ?? ""}\0${thread?.latestTurn?.turnId ?? ""}`;
if (attempted.current === boundary) return;
attempted.current = boundary;
const options = message.sendOptions;
void sendBackgroundQueuedMessage(
ref,
message,
options,
() => {
const current = live.current;
return (
!current.blocked &&
isQueuedMessageDue({
message,
phase: current.phase,
latestToolActivityId: current.latestToolActivityId,
})
);
},
() => live.current.latestToolActivityId,
);
}, [blocked, latestToolActivityId, message, phase, ref, thread?.latestTurn?.turnId]);

return null;
}
29 changes: 25 additions & 4 deletions 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 @@ -7657,6 +7662,22 @@ export default function ChatView(props: ChatViewProps) {
previewAnnotations: [...composerPreviewAnnotations],
reviewComments: [...composerReviewComments],
submissionIntent,
sendOptions: {
modelSelection: ctxSelectedModelSelection,
runtimeMode,
interactionMode: sendInteractionMode,
promptEffort: resolvePromptInjectedEffort(
getProviderModelCapabilities(
ctxSelectedProviderModels,
ctxSelectedModel,
ctxSelectedProvider,
),
ctxSelectedPromptEffort,
),
...(localCheckoutBranchMismatch
? { branch: localCheckoutBranchMismatch.currentBranch }
: {}),
},
queuedAfterToolActivityId: latestCompletedToolActivityId(threadActivities),
createdAt: new Date().toISOString(),
});
Expand Down Expand Up @@ -8619,9 +8640,9 @@ export default function ChatView(props: ChatViewProps) {
}
};

// Sends the oldest queued message once it is due: a tool call finished
// after it was queued, or the turn ended. Only one leaves per boundary; the
// take inside onSend re-anchors the rest.
// Sends the oldest queued message for the open thread. Other threads are
// dispatched by BackgroundQueueCoordinator, which stays mounted across navigation.
// Only one leaves per boundary; the take inside onSend re-anchors the rest.
const sendQueuedMessage = useEffectEvent((message: QueuedComposerMessage) => {
void onSend(undefined, message.submissionIntent, undefined, message);
});
Expand Down
31 changes: 31 additions & 0 deletions apps/web/src/lib/backgroundQueueActiveThread.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
import { scopeThreadRef } from "@t3tools/client-runtime/environment";
import { ThreadId } from "@t3tools/contracts";
import { describe, expect, it } from "vite-plus/test";

import { DraftId } from "../composerDraftStore";
import { resolveThreadRouteTarget } from "../threadRoutes";
import { activeQueueThreadKey } from "./backgroundQueueActiveThread";

describe("activeQueueThreadKey", () => {
it("uses the promoted server thread while the route is still a draft", () => {
const promotedTo = scopeThreadRef("env-2" as never, ThreadId.make("server-thread"));

expect(
activeQueueThreadKey(resolveThreadRouteTarget({ draftId: DraftId.make("draft-1") }), {
environmentId: "env-1" as never,
threadId: ThreadId.make("draft-thread"),
promotedTo,
}),
).toBe("env-2:server-thread");
});

it("leaves an unpromoted draft to the chat view", () => {
expect(
activeQueueThreadKey(resolveThreadRouteTarget({ draftId: DraftId.make("draft-1") }), {
environmentId: "env-1" as never,
threadId: ThreadId.make("draft-thread"),
promotedTo: null,
}),
).toBeNull();
});
});
17 changes: 17 additions & 0 deletions apps/web/src/lib/backgroundQueueActiveThread.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
import { scopedThreadKey } from "@t3tools/client-runtime/environment";

import { resolveActiveThreadRouteRef, type ThreadRouteTarget } from "../threadRoutes";

type DraftRouteState = Parameters<typeof resolveActiveThreadRouteRef>[1];

/**
* The thread whose queue the open chat owns. A draft route that has already
* promoted still shows that server thread, so Stop must stay on the chat path.
*/
export function activeQueueThreadKey(
target: ThreadRouteTarget | null,
draftThread: DraftRouteState,
): string | null {
const ref = resolveActiveThreadRouteRef(target, draftThread);
return ref ? scopedThreadKey(ref) : null;
}
Loading
Loading