From c9028519cae611d2a66706dfd2b1ffd4ad0babba Mon Sep 17 00:00:00 2001 From: lyon Date: Sat, 20 Jun 2026 12:01:30 +0800 Subject: [PATCH] fix: overlap workbench session include query --- internal/cloud/server-workbench-http.test.ts | 75 ++++++++++++++++++++ internal/cloud/server-workbench-http.ts | 9 ++- 2 files changed, 83 insertions(+), 1 deletion(-) diff --git a/internal/cloud/server-workbench-http.test.ts b/internal/cloud/server-workbench-http.test.ts index 9cd35280..d456cb6a 100644 --- a/internal/cloud/server-workbench-http.test.ts +++ b/internal/cloud/server-workbench-http.test.ts @@ -987,6 +987,72 @@ test("workbench session list uses durable projection checkpoint without trace/re } }); +test("workbench session list overlaps include lookup with the page query", async () => { + const makeSession = (id, traceId, updatedAt) => ({ + id, + projectId: "prj_hwpod_workbench", + agentId: "hwlab-code-agent", + status: "running", + ownerUserId: ACTOR.id, + conversationId: id.replace(/^ses_/u, "cnv_"), + threadId: `thread-${id}`, + lastTraceId: traceId, + updatedAt, + session: { + sessionStatus: "running", + lastTraceId: traceId, + messages: [{ role: "user", text: id, traceId, status: "sent" }] + } + }); + const included = makeSession("ses_parallel_include", "trc_parallel_include", "2026-06-19T14:00:00.000Z"); + const firstPage = makeSession("ses_parallel_page_a", "trc_parallel_page_a", "2026-06-19T16:00:00.000Z"); + const secondPage = makeSession("ses_parallel_page_b", "trc_parallel_page_b", "2026-06-19T15:00:00.000Z"); + const facts = emptyFacts(); + for (const session of [included, firstPage, secondPage]) mergeFacts(facts, buildDurableFactsForSession({ session, status: "running" })); + const accessController = { + store: {}, + async ensureBootstrap() {}, + async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } + }; + let pageReleased = false; + let releasePage = () => {}; + const pageGate = new Promise((resolve) => { releasePage = resolve; }); + const queries = []; + const runtimeStore = { + async queryWorkbenchFacts(params = {}) { + queries.push({ ...params, pageReleased }); + if (params.sessionOrder === "updated_desc") await pageGate; + const filtered = filterFacts(facts, params); + return { + facts: filtered, + count: Object.values(filtered).reduce((sum, rows) => sum + rows.length, 0), + persistence: { adapter: "test-durable-workbench-facts", durable: true } + }; + } + }; + const server = createCloudApiServer({ accessController, runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() }); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + + try { + const { port } = server.address(); + const pending = getJson(port, `/v1/workbench/sessions?includeSessionId=${encodeURIComponent(included.id)}&limit=1`); + await waitForCondition(() => queries.some((query) => query.sessionId === included.id) && queries.some((query) => query.sessionOrder === "updated_desc")); + const includeIndex = queries.findIndex((query) => query.sessionId === included.id); + assert.ok(includeIndex >= 0); + assert.equal(queries[includeIndex].pageReleased, false); + pageReleased = true; + releasePage(); + const sessions = await pending; + assert.equal(sessions.status, 200); + assert.equal(sessions.body.sessions[0].sessionId, included.id); + assert.equal(sessions.body.hasMore, true); + } finally { + pageReleased = true; + releasePage(); + await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); + } +}); + test("workbench read model reads terminal session status from durable facts", async () => { const traceStore = createCodeAgentTraceStore(); const traceId = "trc_workbench_durable_terminal_after_memory_running"; @@ -1272,6 +1338,15 @@ async function getJson(port, path) { }; } +async function waitForCondition(predicate, timeoutMs = 500) { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (predicate()) return; + await new Promise((resolve) => setTimeout(resolve, 5)); + } + assert.fail("condition was not met before timeout"); +} + function createDurableFactsRuntimeStore({ sessions = [], queryError = null, queries = [] } = {}) { const facts = emptyFacts(); for (const item of sessions) mergeFacts(facts, buildDurableFactsForSession(item)); diff --git a/internal/cloud/server-workbench-http.ts b/internal/cloud/server-workbench-http.ts index c273fbd7..21f6f002 100644 --- a/internal/cloud/server-workbench-http.ts +++ b/internal/cloud/server-workbench-http.ts @@ -323,12 +323,13 @@ async function handleWorkbenchSessionList(response, url, options, actor) { const includeRouteId = includeSessionId ?? safeConversationId(url.searchParams.get("includeSessionId")); const readModel = createWorkbenchReadModel(options, actor); const pageQuery = { ownerUserId: actor.role === "admin" ? undefined : actor.id, limit: limit + offset + 1, families: WORKBENCH_SESSION_LIST_PAGE_FAMILIES, sessionOrder: "updated_desc" }; + const includeQueryPromise = includeRouteId ? queryWorkbenchSessionInclude(readModel, includeRouteId, perf) : null; const pageResult = await (perf ? perf.measure("workbench_session_page_query", () => readModel.queryFacts(pageQuery)) : readModel.queryFacts(pageQuery)); if (pageResult.error) return sendJson(response, 503, workbenchProjectionStoreError(pageResult.error)); let pageFacts = pageResult.facts; let naturalPage = visibleFactSessions(pageFacts, actor).slice(offset, offset + limit + 1); if (includeRouteId && !naturalPage.some((session) => factSessionMatchesRouteId(session, includeRouteId))) { - const includeResult = await (perf ? perf.measure("workbench_session_include_query", () => queryFactsByRouteId(readModel, includeRouteId, { families: WORKBENCH_SESSION_LIST_PAGE_FAMILIES })) : queryFactsByRouteId(readModel, includeRouteId, { families: WORKBENCH_SESSION_LIST_PAGE_FAMILIES })); + const includeResult = await includeQueryPromise; if (includeResult.error) return sendJson(response, 503, workbenchProjectionStoreError(includeResult.error)); const included = visibleFactSessions(includeResult.facts, actor)[0] ?? null; if (included) { @@ -356,6 +357,12 @@ async function handleWorkbenchSessionList(response, url, options, actor) { }); } +function queryWorkbenchSessionInclude(readModel, includeRouteId, perf = null) { + const run = () => queryFactsByRouteId(readModel, includeRouteId, { families: WORKBENCH_SESSION_LIST_PAGE_FAMILIES }); + const promise = perf ? perf.measure("workbench_session_include_query", run) : run(); + return promise.catch((error) => ({ error })); +} + function sessionMatchesRouteId(session, routeId) { const sessionId = safeSessionId(routeId); if (sessionId && session?.id === sessionId) return true;