Merge pull request #1692 from pikasTech/codex/hwlab-1691-session-list

fix: overlap workbench session include query
This commit is contained in:
Lyon
2026-06-20 12:02:23 +08:00
committed by GitHub
2 changed files with 83 additions and 1 deletions
@@ -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));
+8 -1
View File
@@ -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;