From b4f1da85e40aa0d2f05307cef79618c612d06be2 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Mon, 3 Aug 2026 12:29:19 -0700 Subject: [PATCH] fix(gateway): expire temporal rows in cached session lists (#118875) --- .../server-methods/sessions-list-cache.ts | 37 ++++- .../sessions-read-cache.test.ts | 152 +++++++++++++++++- 2 files changed, 183 insertions(+), 6 deletions(-) diff --git a/src/gateway/server-methods/sessions-list-cache.ts b/src/gateway/server-methods/sessions-list-cache.ts index c1065c08039a..87866c5abebf 100644 --- a/src/gateway/server-methods/sessions-list-cache.ts +++ b/src/gateway/server-methods/sessions-list-cache.ts @@ -2,6 +2,7 @@ import type { SessionsListParams } from "../../../packages/gateway-protocol/src/ import type { OpenClawConfig } from "../../config/types.openclaw.js"; import { readAgentRunIndexVersion } from "../../infra/agent-run-registry.js"; import { isGatewayAdmin } from "../session-sharing.js"; +import type { SessionsListResult } from "../session-utils.types.js"; import { gatewayClientSessionCreator } from "./gateway-client-identity.js"; import { readSessionsMutationVersion } from "./session-change-event.js"; import type { GatewayClient, GatewayRequestContext, RespondFn } from "./types.js"; @@ -10,8 +11,8 @@ type SessionListFence = { agentRunIndexVersion: number; sessionsMutationVersion: number; }; -type SessionListOperation = SessionListFence & { promise: Promise }; -type SessionListCompleted = SessionListFence & { result: unknown }; +type SessionListOperation = SessionListFence & { promise: Promise }; +type SessionListCompleted = SessionListFence & { expiresAt?: number; result: SessionsListResult }; type SessionListState = { completed: Map; config: OpenClawConfig; @@ -78,13 +79,32 @@ function rememberCompletedSessionList( } } +function resolveSessionListExpiration(result: SessionsListResult): number | null | undefined { + let expiresAt: number | undefined; + for (const session of result.sessions) { + // Running durations tick continuously, and a retained child can sit outside + // this page, leaving no authoritative child expiration to cache safely. + if (session.hasActiveSubagentRun || session.childSessions?.length) { + return null; + } + const statusExpiration = session.agentStatus?.expiresAt; + if ( + statusExpiration !== undefined && + (expiresAt === undefined || statusExpiration < expiresAt) + ) { + expiresAt = statusExpiration; + } + } + return expiresAt; +} + export async function respondWithCachedSessionList(params: { client: GatewayClient | null; config: OpenClawConfig; context: GatewayRequestContext; request: SessionsListParams; respond: RespondFn; - run: () => Promise; + run: () => Promise; }): Promise { const workKey = sessionListWorkKey(params.request, params.client); const state = sessionListState(params.context, params.config); @@ -95,7 +115,11 @@ export async function respondWithCachedSessionList(params: { // prevent deriving a safe deadline, so only concurrent temporal requests share work. const cacheCompleted = params.request.activeMinutes === undefined && !params.request.spawnedBy; const completed = cacheCompleted ? state.completed.get(workKey) : undefined; - if (completed && matchesSessionListFence(completed, fence)) { + if ( + completed && + matchesSessionListFence(completed, fence) && + (completed.expiresAt === undefined || completed.expiresAt > Date.now()) + ) { params.respond(true, completed.result, undefined); return; } @@ -111,7 +135,10 @@ export async function respondWithCachedSessionList(params: { .then(params.run) .then((result) => { if (cacheCompleted && matchesSessionListFence(readSessionListFence(params.context), fence)) { - rememberCompletedSessionList(state, workKey, { ...fence, result }); + const expiresAt = resolveSessionListExpiration(result); + if (expiresAt !== null && (expiresAt === undefined || expiresAt > Date.now())) { + rememberCompletedSessionList(state, workKey, { ...fence, result, expiresAt }); + } } return result; }); diff --git a/src/gateway/server-methods/sessions-read-cache.test.ts b/src/gateway/server-methods/sessions-read-cache.test.ts index 47fda1d7d781..9cc74e72bfea 100644 --- a/src/gateway/server-methods/sessions-read-cache.test.ts +++ b/src/gateway/server-methods/sessions-read-cache.test.ts @@ -1,5 +1,9 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import type { SessionsListParams } from "../../../packages/gateway-protocol/src/index.js"; +import { + addSubagentRunForTests, + resetSubagentRegistryForTests, +} from "../../agents/subagent-registry.test-helpers.js"; import { loadSessionEntry, replaceSessionEntry, @@ -9,6 +13,7 @@ import type { OpenClawConfig } from "../../config/types.openclaw.js"; import { resetAgentEventsForTest } from "../../infra/agent-events.js"; import { clearAgentRunContext, registerAgentRunContext } from "../../infra/agent-run-registry.js"; import { withOpenClawTestState } from "../../test-utils/openclaw-test-state.js"; +import type { GatewaySessionRow } from "../session-utils.types.js"; import type { GatewayClient, GatewayRequestContext, RespondFn } from "./types.js"; const loader = vi.hoisted(() => ({ @@ -90,7 +95,7 @@ async function listSessions(params: { return responses[0]?.[1] as { count: number; nextOffset: number | null; - sessions: Array<{ hasActiveRun?: boolean; key: string }>; + sessions: GatewaySessionRow[]; totalCount: number; }; } @@ -235,6 +240,151 @@ describe("sessions.list single-flight", () => { }); }); + it("expires completed rows at the earliest projected agent-status deadline", async () => { + await withOpenClawTestState({ scenario: "minimal" }, async () => { + const clock = vi.spyOn(Date, "now").mockReturnValue(1_000); + const config = await seedSessions(); + for (const [name, expiresAt] of [ + ["active", 1_100], + ["draft", 1_200], + ] as const) { + const scope = { agentId: "main", sessionKey: `agent:main:${name}` }; + const entry = loadSessionEntry(scope); + if (!entry) { + throw new Error(`Missing seeded session ${scope.sessionKey}`); + } + await replaceSessionEntry(scope, { + ...entry, + agentStatus: { note: `${name} needs attention`, expiresAt }, + }); + } + const context = requestContext(config); + const client = identifiedClient("owner@example.com"); + const request = { agentId: "main", archived: "all" as const, limit: 100 }; + + const first = await listSessions({ client, context, request }); + expect( + first.sessions.find((session) => session.key === "agent:main:active")?.agentStatus, + ).toMatchObject({ expiresAt: 1_100 }); + expect( + first.sessions.find((session) => session.key === "agent:main:draft")?.agentStatus, + ).toMatchObject({ expiresAt: 1_200 }); + + clock.mockReturnValue(1_099); + expect(await listSessions({ client, context, request })).toBe(first); + expect(loader.calls).toHaveBeenCalledTimes(1); + + clock.mockReturnValue(1_100); + const expired = await Promise.all( + Array.from({ length: 8 }, () => listSessions({ client, context, request })), + ); + expect(expired.every((result) => result === expired[0])).toBe(true); + expect( + expired[0]?.sessions.find((session) => session.key === "agent:main:active")?.agentStatus, + ).toBeUndefined(); + expect( + expired[0]?.sessions.find((session) => session.key === "agent:main:draft")?.agentStatus, + ).toMatchObject({ expiresAt: 1_200 }); + expect(loader.calls).toHaveBeenCalledTimes(2); + + clock.mockReturnValue(1_199); + expect(await listSessions({ client, context, request })).toBe(expired[0]); + + clock.mockReturnValue(1_200); + const allExpired = await listSessions({ client, context, request }); + expect( + allExpired.sessions.find((session) => session.key === "agent:main:draft")?.agentStatus, + ).toBeUndefined(); + expect(loader.calls).toHaveBeenCalledTimes(3); + }); + }); + + it("expires retained child links when the child is outside the visible page", async () => { + await withOpenClawTestState({ scenario: "minimal" }, async () => { + const { clock, config } = await seedSessionsWithActivityTimes(); + const parentSessionKey = "agent:main:active"; + const childSessionKey = "agent:main:zzz-child"; + await upsertSessionEntry( + { agentId: "main", sessionKey: childSessionKey }, + { + sessionId: "completed-hidden-child", + endedAt: 400, + parentSessionKey, + spawnedBy: parentSessionKey, + status: "done", + updatedAt: 400, + visibility: "shared", + }, + ); + const context = requestContext(config); + const client = identifiedClient("owner@example.com"); + const request = { agentId: "main", archived: "all" as const, limit: 1 }; + + clock.mockReturnValue(1_800_400); + const retained = await listSessions({ client, context, request }); + expect(retained.sessions.map((session) => session.key)).toEqual([parentSessionKey]); + expect(retained.sessions[0]?.childSessions).toEqual([childSessionKey]); + + clock.mockReturnValue(1_800_401); + const expired = await listSessions({ client, context, request }); + expect(expired.sessions.map((session) => session.key)).toEqual([parentSessionKey]); + expect(expired.sessions[0]?.childSessions).toBeUndefined(); + expect(loader.calls).toHaveBeenCalledTimes(2); + }); + }); + + it("refreshes live subagent runtimes while retaining concurrent single-flight", async () => { + await withOpenClawTestState({ scenario: "minimal" }, async () => { + const now = 1_800_000_000_000; + const clock = vi.spyOn(Date, "now").mockReturnValue(now); + const config = await seedSessions(); + const runId = "sessions-list-cache-live-subagent"; + addSubagentRunForTests({ + runId, + childSessionKey: "agent:main:active", + controllerSessionKey: "agent:main:draft", + requesterSessionKey: "agent:main:draft", + requesterDisplayKey: "main", + task: "prove session runtime freshness", + cleanup: "keep", + createdAt: now - 1_000, + startedAt: now - 1_000, + }); + registerAgentRunContext(runId, { + agentId: "main", + projectSessionActive: true, + sessionId: "main-active", + sessionKey: "agent:main:active", + }); + try { + const context = requestContext(config); + const client = identifiedClient("owner@example.com"); + const request = { agentId: "main", archived: "all" as const, limit: 1 }; + + const first = await listSessions({ client, context, request }); + expect(first.sessions[0]).toMatchObject({ + key: "agent:main:active", + hasActiveSubagentRun: true, + runtimeMs: 1_000, + }); + + clock.mockReturnValue(now + 250); + const fresh = await Promise.all( + Array.from({ length: 8 }, () => listSessions({ client, context, request })), + ); + expect(fresh.every((result) => result === fresh[0])).toBe(true); + expect(fresh[0]?.sessions[0]).toMatchObject({ + hasActiveSubagentRun: true, + runtimeMs: 1_250, + }); + expect(loader.calls).toHaveBeenCalledTimes(2); + } finally { + clearAgentRunContext(runId); + resetSubagentRegistryForTests({ persist: false }); + } + }); + }); + it.each([ { description: "the last visible row crosses the inclusive activity cutoff",