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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion docs/agents/sessions.md
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ On mount `useAgentSession` loads the history, then `GET /api/sessions/[id]/state
- The sidebar polls `/api/agent/running` every 2.5 s while the tab is visible; the session-list response is the initial fallback.
- `invalidateSessionListCache()` bumps the generation but **keeps** the previous scan, fresh only while its generation matches. Callers needing only metadata (search hits to sidebar rows) pass `listAllSessions({ allowStale: true })` to read it while it rebuilds in the background, accepting that a seconds-old session is missing.
- `useAgentSession` treats per-session SSE as primary and opens it before each prompt. `prompt_done` completes the UI stage and notification at once, but the stream stays open for the next prompt: the selected session's while selected (`scheduleEventStreamClose()` skips it), any other's for a 30-second grace window. `agent_start` cancels the close timer; `agent_settled` finishes extension-injected runs that have no wrapper-level `prompt_done` and starts a fresh grace window. Never close on the first `agent_end`: retries, compaction and extension-queued messages continue the same logical prompt.
- While a run is active, `useAgentSession` polls `GET /api/agent/[id]` and reconciles on `visibilitychange` / `online`, for terminal events missed by background tabs or half-open connections.
- While a run is active, `useAgentSession` polls `GET /api/agent/[id]` and reconciles on `visibilitychange` / `online`, for terminal events missed by background tabs or half-open connections. Busy replies also refresh context usage; each completed assistant message requests an immediate usage read between model calls. Usage reads from mount, message completion, reconciliation and `agent_end` share a monotonic request id and check the current session/run and mounted hook. A reply applies only when its id is newer than the last applied one, so a delayed read cannot overwrite newer usage, and a newer read that fails does not discard an older one still in flight. Each mount marks the reads started before it as stale.
- Prompt runs carry a monotonic run id; late SSE or reconciliation answers from an old run must be ignored, or they resurrect stale streaming bubbles.
- Every SSE (re)connection is gated on `sessionHookMountedRef`. Under React Strict Mode (`next dev`) the mount-only effect's cleanup clears it and restores it only after the warm-session effect re-runs, so that effect must re-assert it before `maintainEventsConnected()`, or a dev tab never opens its stream.

Expand Down
103 changes: 103 additions & 0 deletions hooks/context-usage.test.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
import assert from "node:assert/strict";
import { readFile } from "node:fs/promises";
import test from "node:test";
import { Script, createContext } from "node:vm";
import ts from "typescript";

const source = ts.createSourceFile("useAgentSession.ts", await readFile(new URL("./useAgentSession.ts", import.meta.url), "utf8"), ts.ScriptTarget.Latest, true);
const nodes = [];
function visit(node) { nodes.push(node); ts.forEachChild(node, visit); }
visit(source);
function callback(name) {
const node = nodes.find((node) => ts.isVariableDeclaration(node) && node.name.getText(source) === name);
assert.ok(node, `missing ${name}`);
return new Script(ts.transpileModule(`(${node.initializer.arguments[0].getText(source)})`, { compilerOptions: { target: ts.ScriptTarget.ESNext } }).outputText);
}
function setup() {
const writes = [];
const requests = [];
const context = createContext({
sessionHookMountedRef: { current: true }, sessionIdRef: { current: "a" },
promptRunIdRef: { current: 1 }, contextUsageRequestIdRef: { current: 0 }, contextUsageAppliedIdRef: { current: 0 },
agentRunningRef: { current: true }, sdkAgentActiveRef: { current: true }, rpcPromptPendingRef: { current: true },
setContextUsage: (value) => writes.push(value),
fetch: (url) => { const request = Promise.withResolvers(); requests.push({ ...request, url }); return request.promise; },
syncLiveModel() {}, setIsCompacting() {}, setAutoCompactionEnabled() {}, setQueuedMessages() {},
normalizeQueuedMessages: (value) => value, finishPromptWithoutStream: () => { throw new Error("busy run settled"); },
});
context.applyContextUsage = callback("applyContextUsage").runInContext(context);
context.refreshContextUsage = callback("refreshContextUsage").runInContext(context);
context.reconcileAgentState = callback("reconcileAgentState").runInContext(context);
const reply = (index, usage, busy = true) => requests[index].resolve(Response.json({ running: busy, state: { contextUsage: usage, isStreaming: busy, isPromptRunning: busy } }));
return { context, writes, requests, reply };
}
const usage = (tokens) => ({ percent: tokens / 100, contextWindow: 10_000, tokens });

test("a busy reconciliation refreshes context usage without settling the run", async () => {
const state = setup();
const pending = state.context.reconcileAgentState("a");
state.reply(0, usage(100));
await pending;
assert.deepEqual(state.writes, [usage(100)]);
assert.equal(state.context.agentRunningRef.current, true);
});

test("a newer assistant usage read wins over a delayed poll", async () => {
const state = setup();
const poll = state.context.reconcileAgentState("a");
const refresh = state.context.refreshContextUsage("a");
state.reply(1, usage(200)); await refresh;
state.reply(0, usage(100)); await poll;
assert.deepEqual(state.writes, [usage(200)]);
});

test("a failed newer read does not discard an older poll's usage", async () => {
const state = setup();
const poll = state.context.reconcileAgentState("a");
const refresh = state.context.refreshContextUsage("a");
state.requests[1].resolve(new Response("Unavailable", { status: 503 })); await refresh;
state.reply(0, usage(100)); await poll;
assert.deepEqual(state.writes, [usage(100)]);
});

test("usage replies from an old run, another session or an unmounted hook are ignored", async () => {
for (const invalidate of [
(ctx) => { ctx.promptRunIdRef.current++; },
(ctx) => { ctx.sessionIdRef.current = "b"; },
(ctx) => { ctx.sessionHookMountedRef.current = false; },
(ctx) => { ctx.contextUsageAppliedIdRef.current = ctx.contextUsageRequestIdRef.current; },
]) {
const state = setup();
const pending = state.context.refreshContextUsage("a");
invalidate(state.context); state.reply(0, usage(100)); await pending;
assert.deepEqual(state.writes, []);
}
});

test("failed reads preserve usage and a subsequent read can recover", async () => {
const state = setup();
const failed = state.context.refreshContextUsage("a");
state.requests[0].resolve(new Response("Unavailable", { status: 503 })); await failed;
const broken = state.context.refreshContextUsage("a");
state.requests[1].reject(new TypeError("offline")); await broken;
assert.deepEqual(state.writes, []);
const recovered = state.context.refreshContextUsage("a");
state.reply(2, null); await recovered;
assert.deepEqual(state.writes, [null]);
});

test("only completed assistant messages trigger an immediate usage read", () => {
const messageEnd = nodes.find((node) => ts.isCaseClause(node) && node.expression.getText(source) === '"message_end"');
const script = new Script(ts.transpileModule(`(() => { switch(event.type) { ${messageEnd.getText(source)} } })()`, { compilerOptions: { target: ts.ScriptTarget.ESNext } }).outputText);
for (const role of ["assistant", "toolResult", "user", "system"]) {
const reads = [];
script.runInNewContext({
event: { type: "message_end", message: { role, content: [] } },
agentRunningRef: { current: true }, sessionIdRef: { current: "a" }, optimisticUserMessageKeyRef: { current: null },
isSystemMessageEvent: (event) => event.message.role === "system",
normalizeToolCalls: (message) => message, userMessageKey: () => "user", setMessages() {}, dispatch() {}, setAgentPhase() {},
refreshContextUsage: (sid) => reads.push(sid),
});
assert.deepEqual(reads, role === "assistant" ? ["a"] : []);
}
});
58 changes: 47 additions & 11 deletions hooks/useAgentSession.ts
Original file line number Diff line number Diff line change
Expand Up @@ -350,6 +350,10 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
const [liveThinkingLevel, setLiveThinkingLevel] = useState<ConcreteThinkingLevel | null>(null);
const [retryInfo, setRetryInfo] = useState<{ attempt: number; maxAttempts: number; errorMessage?: string } | null>(null);
const [contextUsage, setContextUsage] = useState<{ percent: number | null; contextWindow: number; tokens: number | null } | null>(null);
const contextUsageRequestIdRef = useRef(0);
// Highest request id whose reply was applied. A reply applies only when it
// is newer, so a failed newer read never discards an older good one.
const contextUsageAppliedIdRef = useRef(0);
const [systemPrompt, setSystemPrompt] = useState<string | null>(null);
const [forkingEntryId, setForkingEntryId] = useState<string | null>(null);
const [currentModelOverride, setCurrentModelOverride] = useState<{ provider: string; modelId: string } | null>(null);
Expand Down Expand Up @@ -575,6 +579,26 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
} satisfies SessionStatsInfo;
}, [messages, sessionStatsOverride, contextUsage, data?.context.messages, data?.filePath, data?.totalActiveMs, data?.stats, session?.id, session?.name]);

const applyContextUsage = useCallback((state: AgentStateResponse | undefined, sid: string, runId: number, requestId: number) => {
if (!sessionHookMountedRef.current || sessionIdRef.current !== sid
|| promptRunIdRef.current !== runId || requestId <= contextUsageAppliedIdRef.current) return;
contextUsageAppliedIdRef.current = requestId;
if (state?.contextUsage !== undefined) setContextUsage(state.contextUsage ?? null);
}, []);

const refreshContextUsage = useCallback(async (sid: string) => {
const runId = promptRunIdRef.current;
const requestId = ++contextUsageRequestIdRef.current;
try {
const res = await fetch(`/api/agent/${encodeURIComponent(sid)}`);
if (!res.ok) return;
const data = await res.json() as { state?: AgentStateResponse };
applyContextUsage(data.state, sid, runId, requestId);
} catch {
// A later message or the running-state poll retries the usage read.
}
}, [applyContextUsage]);

const loadSession = useCallback(async (sid: string, showLoading = false, includeState = false, options?: { force?: boolean }) => {
// Single-flight: concurrent reads for the same session (mount + SSE settle +
// reconcile) share one request unless the caller forces a fresh read.
Expand Down Expand Up @@ -689,6 +713,8 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
if (!includeState) return null;

try {
const runId = promptRunIdRef.current;
const usageRequestId = ++contextUsageRequestIdRef.current;
const stateRes = await fetch(`/api/sessions/${encodeURIComponent(sid)}/state`);
if (!stateRes.ok) throw new Error(`HTTP ${stateRes.status}`);
const agentState = await stateRes.json() as { running: boolean; state?: AgentStateResponse };
Expand All @@ -697,7 +723,7 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
const liveState = agentState.state;
syncLiveModel(liveState);
if (liveState) {
if (liveState.contextUsage !== undefined) setContextUsage(liveState.contextUsage ?? null);
applyContextUsage(liveState, sid, runId, usageRequestId);
if (liveState.systemPrompt !== undefined) setSystemPrompt(liveState.systemPrompt ?? null);
if (liveState.extensionStatuses !== undefined) setExtensionStatuses(liveState.extensionStatuses ?? []);
if (liveState.extensionWidgets !== undefined) setExtensionWidgets(liveState.extensionWidgets ?? []);
Expand All @@ -723,7 +749,7 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
if (loadFlightsRef.current.get(flightKey) === flight) loadFlightsRef.current.delete(flightKey);
});
return await flight;
}, [setToolPresetState, syncLiveModel]);
}, [applyContextUsage, setToolPresetState, syncLiveModel]);

const loadContext = useCallback(async (sid: string, leafId: string | null, before?: string | null, options?: { tail?: number; signal?: AbortSignal }) => {
try {
Expand Down Expand Up @@ -1280,6 +1306,7 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
const reconcileAgentState = useCallback(async (sid: string) => {
if (!agentRunningRef.current || sessionIdRef.current !== sid) return;
const runId = promptRunIdRef.current;
const usageRequestId = ++contextUsageRequestIdRef.current;
try {
const res = await fetch(`/api/agent/${encodeURIComponent(sid)}`);
if (!res.ok) return;
Expand All @@ -1289,6 +1316,7 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
// flight) — everything in it is stale, drop it.
if (sessionIdRef.current !== sid || promptRunIdRef.current !== runId) return;
const state = data.state;
applyContextUsage(state, sid, runId, usageRequestId);
syncLiveModel(state);
// Mirror compaction state unconditionally: a missed compaction_end
// would otherwise leave the "Stop compaction" UI stuck. No state
Expand All @@ -1305,7 +1333,6 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
}
if (!agentRunningRef.current) return;
if (state) {
if (state.contextUsage !== undefined) setContextUsage(state.contextUsage ?? null);
if (state.systemPrompt !== undefined) setSystemPrompt(state.systemPrompt ?? null);
if (state.extensionStatuses !== undefined) setExtensionStatuses(state.extensionStatuses ?? []);
if (state.extensionWidgets !== undefined) setExtensionWidgets(state.extensionWidgets ?? []);
Expand All @@ -1314,7 +1341,7 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
} catch {
// Network still down — the next poll / visibility / online tick retries.
}
}, [finishPromptWithoutStream, syncLiveModel]);
}, [applyContextUsage, finishPromptWithoutStream, syncLiveModel]);

// Recovery net for missed SSE events: while the agent is running, verify
// against the server periodically and whenever the tab returns to the
Expand Down Expand Up @@ -1380,12 +1407,16 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
setRetryInfo(null);
dispatch({ type: "end" });
if (sessionIdRef.current) {
loadSession(sessionIdRef.current);
fetch(`/api/agent/${encodeURIComponent(sessionIdRef.current)}`)
.then((r) => r.json())
.then((d: { state?: AgentStateResponse }) => {
const sid = sessionIdRef.current;
const runId = promptRunIdRef.current;
const usageRequestId = ++contextUsageRequestIdRef.current;
loadSession(sid);
fetch(`/api/agent/${encodeURIComponent(sid)}`)
.then((r) => r.ok ? r.json() : null)
.then((d: { state?: AgentStateResponse } | null) => {
if (!d || !sessionHookMountedRef.current || sessionIdRef.current !== sid || promptRunIdRef.current !== runId) return;
syncLiveModel(d.state);
if (d.state?.contextUsage !== undefined) setContextUsage(d.state.contextUsage ?? null);
applyContextUsage(d.state, sid, runId, usageRequestId);
if (d.state?.systemPrompt !== undefined) setSystemPrompt(d.state.systemPrompt ?? null);
if (d.state?.extensionStatuses !== undefined) setExtensionStatuses(d.state.extensionStatuses ?? []);
if (d.state?.extensionWidgets !== undefined) setExtensionWidgets(d.state.extensionWidgets ?? []);
Expand Down Expand Up @@ -1506,6 +1537,10 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
});
} else if (completed) {
setMessages((prev) => [...prev, normalizeToolCalls(completed)]);
if (completed.role === "assistant") {
const sid = sessionIdRef.current;
if (sid) void refreshContextUsage(sid);
}
}
dispatch({ type: "end" });
setAgentPhase({ kind: "waiting_model" });
Expand Down Expand Up @@ -1617,7 +1652,7 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
setExtensionDialogs((queue) => removeExtensionUiRequest(queue, event.id as string));
break;
}
}, [addNotice, cancelEventStreamGrace, handleExtensionUiRequest, loadSession, notifyPromptStage, onAgentEnd, scheduleEventStreamClose, scrollToBottom, settleUiStage, syncLiveModel]);
}, [addNotice, applyContextUsage, cancelEventStreamGrace, handleExtensionUiRequest, loadSession, notifyPromptStage, onAgentEnd, refreshContextUsage, scheduleEventStreamClose, scrollToBottom, settleUiStage, syncLiveModel]);
handleAgentEventRef.current = handleAgentEvent;

const handleSend = useCallback(async (message: string, images?: AttachedImage[]) => {
Expand Down Expand Up @@ -2396,6 +2431,8 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
// Load session on mount
useEffect(() => {
sessionHookMountedRef.current = true;
// Usage reads started before a remount are stale.
contextUsageAppliedIdRef.current = contextUsageRequestIdRef.current;
if (session) {
sessionIdRef.current = session.id;
// Snapshot fast path: show the cached history window immediately, then
Expand Down Expand Up @@ -2451,7 +2488,6 @@ export function useAgentSession(opts: UseAgentSessionOptions) {
}
if (agentState?.state) {
if (agentState.state.isCompacting !== undefined) setIsCompacting(agentState.state.isCompacting);
if (agentState.state.contextUsage !== undefined) setContextUsage(agentState.state.contextUsage ?? null);
if (agentState.state.systemPrompt !== undefined) setSystemPrompt(agentState.state.systemPrompt ?? null);
if (agentState.state.extensionStatuses !== undefined) setExtensionStatuses(agentState.state.extensionStatuses ?? []);
if (agentState.state.extensionWidgets !== undefined) setExtensionWidgets(agentState.state.extensionWidgets ?? []);
Expand Down