mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
feat: report session list CPU time in diagnostics (#149291)
* feat(gateway): measure synchronous session-list CPU * fix(ci): repair native wrapper and skill-read checks Restore three newly imported modules in the shared native PR wrapper inventory. Observe fs-safe Root.read through its current final-fence hook while preserving all aggregate byte, status, and inventory assertions. Both defects reproduced before repair. All 24 plugin-skill bundle cases pass; nine focused native-wrapper cases pass on Blacksmith Testbox with unchanged limits. Independent review through P2 is clean. * fix(ci): refresh measured UI budget and cold worktree fixture Record the approved 362706-byte startup baseline under a 355 KiB cap, preserving the 512-byte growth and 64-byte variance allowances. Exact base and merge builds emitted byte-identical UI assets; hosted CI measured the same size. Include the newly imported bundled plugin source utility in cold worktree fixtures, matching the repair already landed in #109506. Preserve all assertions and timeouts. * test(slack): align Socket Mode fixtures with SDK acknowledgements Give the shared receiver mock real EventEmitter listener ownership and an asynchronous acknowledgement sender, preserving the production envelope guard. Narrow WebSocket ArrayBuffer data before converting it so the loopback test typechecks without copying an existing Buffer. The former ingress-cleanup failure reproduces before this repair; all 120 selected Slack tests and the extension test-type graph pass afterward.
This commit is contained in:
parent
12a483edfb
commit
8f463a22fc
12 changed files with 614 additions and 221 deletions
|
|
@ -1,5 +1,5 @@
|
|||
{
|
||||
"startupJsGzipBytes": 362012,
|
||||
"reason": "Maintainer-approved 354 KiB cap after integrated scoped catalog invalidation (#147845), remote workspaces (#147773), and update-status feedback (#147883). Controlled build measured 362012 B; deferring the new workspace labels saved only 64 B and still failed the previous limit. Existing 512 B growth and 64 B build variance allowances remain unchanged.",
|
||||
"updatedAt": "2026-09-14"
|
||||
"startupJsGzipBytes": 362706,
|
||||
"reason": "Maintainer-approved correction in #149291 for the existing startup bundle: exact base 7b6094ec7b7d and CI merge b524940307fc emitted byte-identical Control UI output at 362706 B with Vite 8.2.2/Pako 3.0.1 on Node 24.21.0. CI run 35173887678 also measured 362706 B. The 512 B growth and 64 B build variance allowances remain unchanged.",
|
||||
"updatedAt": "2026-09-17"
|
||||
}
|
||||
|
|
|
|||
|
|
@ -249,6 +249,25 @@ admission before the handler. The `response` phase includes the synchronous resp
|
|||
durations, not CPU time or proof of client receipt. No query text or session
|
||||
contents are included.
|
||||
|
||||
The same record includes fractional-millisecond current-thread CPU measurements
|
||||
for synchronous work: `storeLoadThreadCpuMs`, `prepareThreadCpuMs`,
|
||||
`rowThreadCpuMs`, `cacheSelectionThreadCpuMs`, `cachePublicationThreadCpuMs`, and
|
||||
`responseThreadCpuMs`. Preparation and row totals accumulate synchronous
|
||||
chunks, excluding yields and intervening microtasks. Row CPU also includes final
|
||||
list construction. Cache selection and publication finish before the response
|
||||
callback is measured; response CPU excludes network waits. Measurements finish
|
||||
before this diagnostic record is published or logged.
|
||||
Registry readiness waits and worker CPU are not included in these measurements.
|
||||
These are selected inclusive CPU intervals, including same-thread native work and
|
||||
garbage collection, not SQL-only CPU or a complete request CPU total.
|
||||
|
||||
Unvisited measurements are omitted. Hits and followers retain their own selection
|
||||
and response CPU without inheriting the producer's store, projection, or publication
|
||||
work. If a CPU counter read fails, all CPU fields are omitted for that request;
|
||||
its result and elapsed diagnostics are preserved. Existing activation and the
|
||||
one-second warning threshold are unchanged, so missing slow records do not account
|
||||
for CPU consumed by faster requests.
|
||||
|
||||
### WS log style
|
||||
|
||||
`openclaw gateway` supports a per-gateway style switch:
|
||||
|
|
|
|||
|
|
@ -1,3 +1,4 @@
|
|||
import { EventEmitter } from "node:events";
|
||||
import fs from "node:fs";
|
||||
import path from "node:path";
|
||||
import type { ChannelRuntimeSurface } from "openclaw/plugin-sdk/channel-contract";
|
||||
|
|
@ -514,11 +515,9 @@ vi.mock("@slack/bolt", () => {
|
|||
requestListener = (...args: unknown[]) => slackTestState.httpRequestListenerMock(...args);
|
||||
}
|
||||
class SocketModeReceiver {
|
||||
client = {
|
||||
...slackClient,
|
||||
on: vi.fn(),
|
||||
off: vi.fn(),
|
||||
};
|
||||
client = Object.assign(new EventEmitter(), slackClient, {
|
||||
send: vi.fn<(envelopeId: string) => Promise<void>>().mockResolvedValue(undefined),
|
||||
});
|
||||
|
||||
constructor(args: { logger?: { error: (...args: unknown[]) => void } }) {
|
||||
slackTestState.socketModeLogger = args.logger;
|
||||
|
|
|
|||
|
|
@ -505,7 +505,11 @@ describe("createSlackBoltApp", () => {
|
|||
const acknowledgements: string[] = [];
|
||||
socketServer.on("connection", (socket) => {
|
||||
socket.on("message", (data) => {
|
||||
const bytes = Array.isArray(data) ? Buffer.concat(data) : Buffer.from(data);
|
||||
const bytes = Array.isArray(data)
|
||||
? Buffer.concat(data)
|
||||
: data instanceof ArrayBuffer
|
||||
? Buffer.from(data)
|
||||
: data;
|
||||
acknowledgements.push(JSON.parse(bytes.toString("utf8")).envelope_id);
|
||||
});
|
||||
socket.send(JSON.stringify({ type: "hello" }));
|
||||
|
|
|
|||
|
|
@ -50,8 +50,8 @@ const CONTROL_UI_LOCALE_GZIP_BYTES = 300 * KIB;
|
|||
const controlUiPerformanceBudgets = {
|
||||
startupJsRequests: 18,
|
||||
startupCssRequests: 1,
|
||||
// 354 KiB approved for scoped catalog invalidation, remote workspaces, and update-status feedback.
|
||||
startupJsGzipBytes: 354 * KIB,
|
||||
// 355 KiB approved in #149291 for the measured existing bundle; growth allowances stay fixed.
|
||||
startupJsGzipBytes: 355 * KIB,
|
||||
// Keep 45 KiB advisory: tiny integrated changes must not exhaust the budget.
|
||||
// The fixed 50 KiB ceiling bounds accumulation of small changes.
|
||||
startupCssGzipBytes: 50 * KIB,
|
||||
|
|
|
|||
|
|
@ -17,6 +17,12 @@ const METRICS = [
|
|||
"writerExecutionMs",
|
||||
"completionDelayMs",
|
||||
"handlerElapsedMs",
|
||||
"storeLoadThreadCpuMs",
|
||||
"prepareThreadCpuMs",
|
||||
"rowThreadCpuMs",
|
||||
"cacheSelectionThreadCpuMs",
|
||||
"cachePublicationThreadCpuMs",
|
||||
"responseThreadCpuMs",
|
||||
"prepareSyncMs",
|
||||
"rowSyncMs",
|
||||
"yieldWaitMs",
|
||||
|
|
|
|||
|
|
@ -250,67 +250,80 @@ export async function respondWithCachedSessionList(params: {
|
|||
run: () => Promise<SessionsListResult>;
|
||||
diagnostics?: SessionListDiagnostics;
|
||||
}): Promise<void> {
|
||||
const workKey = sessionListWorkKey(params.request, params.client, params.config);
|
||||
const state = sessionListState(params.context, params.config);
|
||||
const modelCatalogRevision = readSessionListModelCatalogFence(params.modelCatalog);
|
||||
// Activity windows and child retention expire without mutations; hidden paginated rows
|
||||
// prevent deriving a safe deadline, so only concurrent temporal requests share work.
|
||||
// Rejected and off-page candidates can change live/goal state without a store write.
|
||||
// Searches and active-only reads may coalesce in flight, but cannot reuse completed pages.
|
||||
const cacheCompleted =
|
||||
params.request.activeMinutes === undefined &&
|
||||
params.request.activeOnly !== true &&
|
||||
!params.request.spawnedBy &&
|
||||
!params.request.search?.trim();
|
||||
const completed = cacheCompleted
|
||||
? readCompletedSessionList(state, workKey, modelCatalogRevision)
|
||||
: undefined;
|
||||
if (completed) {
|
||||
params.diagnostics?.setCacheRole("completed-hit");
|
||||
params.diagnostics?.setSelectedRowCount(completed.count);
|
||||
params.respond(true, completed, undefined);
|
||||
return;
|
||||
}
|
||||
const pending = state.inFlight.get(workKey);
|
||||
if (pending?.modelCatalogRevision === modelCatalogRevision) {
|
||||
params.diagnostics?.setCacheRole("in-flight-follower", pending.workTrace);
|
||||
const result = await pending.promise;
|
||||
params.diagnostics?.setSelectedRowCount(result.count);
|
||||
params.respond(true, result, undefined);
|
||||
return;
|
||||
}
|
||||
|
||||
// A request may share only work begun at the same fence. A transition during projection
|
||||
// leaves current callers intact but fences every later caller and cache write.
|
||||
params.diagnostics?.setCacheRole("projection-owner");
|
||||
const promise = Promise.resolve()
|
||||
.then(params.run)
|
||||
.then((result) => {
|
||||
if (
|
||||
cacheCompleted &&
|
||||
sessionListsByContext.get(params.context) === state &&
|
||||
matchesSessionListFence(state, readSessionListFence(params.context)) &&
|
||||
readSessionListModelCatalogFence(params.modelCatalog) === modelCatalogRevision
|
||||
) {
|
||||
const now = Date.now();
|
||||
const expiresAt = resolveSessionListExpiration(result, now);
|
||||
if (expiresAt === undefined || expiresAt > now) {
|
||||
state.completed.delete(workKey);
|
||||
state.completed.set(workKey, { modelCatalogRevision, result, expiresAt });
|
||||
pruneMapToMaxSize(state.completed, SESSIONS_LIST_COMPLETED_CACHE_LIMIT);
|
||||
}
|
||||
}
|
||||
return result;
|
||||
});
|
||||
const operation = { modelCatalogRevision, promise, workTrace: params.diagnostics?.trace };
|
||||
state.inFlight.set(workKey, operation);
|
||||
let selectionCpu = params.diagnostics?.startSyncCpu();
|
||||
try {
|
||||
const result = await promise;
|
||||
params.diagnostics?.setSelectedRowCount(result.count);
|
||||
params.respond(true, result, undefined);
|
||||
} finally {
|
||||
if (state.inFlight.get(workKey) === operation) {
|
||||
state.inFlight.delete(workKey);
|
||||
const workKey = sessionListWorkKey(params.request, params.client, params.config);
|
||||
const state = sessionListState(params.context, params.config);
|
||||
const modelCatalogRevision = readSessionListModelCatalogFence(params.modelCatalog);
|
||||
// Activity windows and child retention expire without mutations; hidden paginated rows
|
||||
// prevent deriving a safe deadline, so only concurrent temporal requests share work.
|
||||
// Rejected and off-page candidates can change live/goal state without a store write.
|
||||
// Searches and active-only reads may coalesce in flight, but cannot reuse completed pages.
|
||||
const cacheCompleted =
|
||||
params.request.activeMinutes === undefined &&
|
||||
params.request.activeOnly !== true &&
|
||||
!params.request.spawnedBy &&
|
||||
!params.request.search?.trim();
|
||||
const completed = cacheCompleted
|
||||
? readCompletedSessionList(state, workKey, modelCatalogRevision)
|
||||
: undefined;
|
||||
const pending = completed ? undefined : state.inFlight.get(workKey);
|
||||
params.diagnostics?.finishSyncCpu("cacheSelectionThreadCpuMs", selectionCpu);
|
||||
selectionCpu = undefined;
|
||||
if (completed) {
|
||||
params.diagnostics?.setCacheRole("completed-hit");
|
||||
params.diagnostics?.setSelectedRowCount(completed.count);
|
||||
params.respond(true, completed, undefined);
|
||||
return;
|
||||
}
|
||||
if (pending?.modelCatalogRevision === modelCatalogRevision) {
|
||||
params.diagnostics?.setCacheRole("in-flight-follower", pending.workTrace);
|
||||
const result = await pending.promise;
|
||||
params.diagnostics?.setSelectedRowCount(result.count);
|
||||
params.respond(true, result, undefined);
|
||||
return;
|
||||
}
|
||||
|
||||
// A request may share only work begun at the same fence. A transition during projection
|
||||
// leaves current callers intact but fences every later caller and cache write.
|
||||
params.diagnostics?.setCacheRole("projection-owner");
|
||||
const promise = Promise.resolve()
|
||||
.then(params.run)
|
||||
.then((result) => {
|
||||
const publicationCpu = params.diagnostics?.startSyncCpu();
|
||||
try {
|
||||
if (
|
||||
cacheCompleted &&
|
||||
sessionListsByContext.get(params.context) === state &&
|
||||
matchesSessionListFence(state, readSessionListFence(params.context)) &&
|
||||
readSessionListModelCatalogFence(params.modelCatalog) === modelCatalogRevision
|
||||
) {
|
||||
const now = Date.now();
|
||||
const expiresAt = resolveSessionListExpiration(result, now);
|
||||
if (expiresAt === undefined || expiresAt > now) {
|
||||
state.completed.delete(workKey);
|
||||
state.completed.set(workKey, { modelCatalogRevision, result, expiresAt });
|
||||
pruneMapToMaxSize(state.completed, SESSIONS_LIST_COMPLETED_CACHE_LIMIT);
|
||||
}
|
||||
}
|
||||
return result;
|
||||
} finally {
|
||||
params.diagnostics?.finishSyncCpu("cachePublicationThreadCpuMs", publicationCpu);
|
||||
}
|
||||
});
|
||||
const operation = { modelCatalogRevision, promise, workTrace: params.diagnostics?.trace };
|
||||
state.inFlight.set(workKey, operation);
|
||||
try {
|
||||
const result = await promise;
|
||||
params.diagnostics?.setSelectedRowCount(result.count);
|
||||
params.respond(true, result, undefined);
|
||||
} finally {
|
||||
if (state.inFlight.get(workKey) === operation) {
|
||||
state.inFlight.delete(workKey);
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
// Normal selection closes before any response or wait; this also covers selection errors.
|
||||
params.diagnostics?.finishSyncCpu("cacheSelectionThreadCpuMs", selectionCpu);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -25,6 +25,13 @@ type Phase =
|
|||
| "response"
|
||||
| "handlerExit";
|
||||
type CacheRole = "unreached" | "completed-hit" | "in-flight-follower" | "projection-owner";
|
||||
type SynchronousCpuMetric =
|
||||
| "storeLoadThreadCpuMs"
|
||||
| "prepareThreadCpuMs"
|
||||
| "rowThreadCpuMs"
|
||||
| "cacheSelectionThreadCpuMs"
|
||||
| "cachePublicationThreadCpuMs"
|
||||
| "responseThreadCpuMs";
|
||||
const sessionListDiagnostics = channel("openclaw.session.list");
|
||||
|
||||
export type SessionListDiagnostics = NonNullable<ReturnType<typeof startSessionListDiagnostics>>;
|
||||
|
|
@ -53,6 +60,30 @@ function startSessionListDiagnostics(
|
|||
| undefined;
|
||||
let selectedRowCount: number | undefined;
|
||||
let responseOutcome: "none" | "ok" | "error" | "threw" = "none";
|
||||
let cpuMetrics: Partial<Record<SynchronousCpuMetric, number>> | undefined = {};
|
||||
const startSyncCpu = (): NodeJS.CpuUsage | undefined => {
|
||||
if (!cpuMetrics) {
|
||||
return undefined;
|
||||
}
|
||||
try {
|
||||
return process.threadCpuUsage();
|
||||
} catch {
|
||||
cpuMetrics = undefined;
|
||||
return undefined;
|
||||
}
|
||||
};
|
||||
const finishSyncCpu = (metric: SynchronousCpuMetric, started: NodeJS.CpuUsage | undefined) => {
|
||||
if (!started || !cpuMetrics) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
const used = process.threadCpuUsage(started);
|
||||
cpuMetrics[metric] = (cpuMetrics[metric] ?? 0) + (used.user + used.system) / 1_000;
|
||||
} catch {
|
||||
// Failed probes omit CPU totals for this request without replacing its result.
|
||||
cpuMetrics = undefined;
|
||||
}
|
||||
};
|
||||
const mark = (next: Phase) => {
|
||||
checkpoint = performance.now();
|
||||
timing.mark(phase);
|
||||
|
|
@ -61,6 +92,8 @@ function startSessionListDiagnostics(
|
|||
return {
|
||||
trace,
|
||||
mark,
|
||||
startSyncCpu,
|
||||
finishSyncCpu,
|
||||
get projection() {
|
||||
return projection;
|
||||
},
|
||||
|
|
@ -82,12 +115,14 @@ function startSessionListDiagnostics(
|
|||
respond: ((...args) => {
|
||||
mark("response");
|
||||
responseOutcome = args[0] ? "ok" : "error";
|
||||
const responseCpu = startSyncCpu();
|
||||
try {
|
||||
return respond(...args);
|
||||
} catch (error) {
|
||||
responseOutcome = "threw";
|
||||
throw error;
|
||||
} finally {
|
||||
finishSyncCpu("responseThreadCpuMs", responseCpu);
|
||||
mark("handlerExit");
|
||||
}
|
||||
}) satisfies RespondFn,
|
||||
|
|
@ -116,6 +151,7 @@ function startSessionListDiagnostics(
|
|||
handlerElapsedMs: Math.round(handlerElapsedMs),
|
||||
cacheRole,
|
||||
phaseDurationsMs,
|
||||
...cpuMetrics,
|
||||
...(projection
|
||||
? Object.fromEntries(
|
||||
Object.entries(projection).map(([key, value]) => [key, Math.round(value)]),
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ import { channel } from "node:diagnostics_channel";
|
|||
import { performance } from "node:perf_hooks";
|
||||
import { isMainThread, threadId } from "node:worker_threads";
|
||||
import { afterEach, beforeEach, expect, test, vi } from "vitest";
|
||||
import * as registryRead from "../../agents/subagents/registry/subagent-registry-read.js";
|
||||
import { upsertSessionEntryCore } from "../../config/sessions/session-accessor.js";
|
||||
import {
|
||||
areDiagnosticsEnabledForProcess,
|
||||
|
|
@ -16,36 +17,71 @@ import {
|
|||
import { createDeferredCore } from "../../shared/deferred.js";
|
||||
import { withOpenClawTestState } from "../../test-utils/openclaw-test-state.js";
|
||||
import * as titleReader from "../session-transcript-title-reader.js";
|
||||
import * as rowProjection from "../session-utils-row.js";
|
||||
import * as sessionUtils from "../session-utils.js";
|
||||
import {
|
||||
identifiedClient,
|
||||
listSessions,
|
||||
requestContext,
|
||||
seedSessions,
|
||||
sessionReadHandlers,
|
||||
} from "./sessions-read-cache.test-support.js";
|
||||
import { sessionLog } from "./sessions-shared.js";
|
||||
import { sessionSubscriptionHandlers } from "./sessions-subscriptions.js";
|
||||
import type { RespondFn } from "./types.js";
|
||||
|
||||
const scheduler = vi.hoisted(() => ({ onYield: undefined as (() => Promise<void>) | undefined }));
|
||||
const scheduler = vi.hoisted(() => ({
|
||||
onYield: undefined as (() => Promise<void>) | undefined,
|
||||
afterYield: undefined as (() => void) | undefined,
|
||||
}));
|
||||
vi.mock("node:timers/promises", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("node:timers/promises")>();
|
||||
return {
|
||||
...actual,
|
||||
setImmediate: async (...args: Parameters<typeof actual.setImmediate>) => {
|
||||
const result = await actual.setImmediate(...args);
|
||||
await scheduler.onYield?.();
|
||||
return result;
|
||||
setImmediate: (...args: Parameters<typeof actual.setImmediate>) => {
|
||||
const pause = (async () => {
|
||||
const result = await actual.setImmediate(...args);
|
||||
await scheduler.onYield?.();
|
||||
return result;
|
||||
})();
|
||||
const afterYield = scheduler.afterYield;
|
||||
if (afterYield) {
|
||||
// Register unrelated work after the consumer's reaction, before its awaited continuation.
|
||||
queueMicrotask(() => {
|
||||
void pause.then(afterYield);
|
||||
});
|
||||
}
|
||||
return pause;
|
||||
},
|
||||
};
|
||||
});
|
||||
|
||||
let previousDiagnostics: boolean;
|
||||
let clock: number;
|
||||
let cpu: NodeJS.CpuUsage;
|
||||
let cpuProbeFailure: Error | undefined;
|
||||
const threadCpuProbe = vi.fn<(previous?: NodeJS.CpuUsage) => NodeJS.CpuUsage>();
|
||||
const producerCpuFields = [
|
||||
"storeLoadThreadCpuMs",
|
||||
"prepareThreadCpuMs",
|
||||
"rowThreadCpuMs",
|
||||
"cachePublicationThreadCpuMs",
|
||||
] as const;
|
||||
const threadCpuFields = [...producerCpuFields, "cacheSelectionThreadCpuMs", "responseThreadCpuMs"];
|
||||
let records: Array<{ trace: DiagnosticTraceContext | undefined; fields: Record<string, unknown> }>;
|
||||
beforeEach(() => {
|
||||
previousDiagnostics = areDiagnosticsEnabledForProcess();
|
||||
setDiagnosticsEnabledForProcess(true);
|
||||
clock = 0;
|
||||
cpu = { user: 0, system: 0 };
|
||||
cpuProbeFailure = undefined;
|
||||
threadCpuProbe.mockReset().mockImplementation((previous = { user: 0, system: 0 }) => {
|
||||
if (cpuProbeFailure) {
|
||||
throw cpuProbeFailure;
|
||||
}
|
||||
return { user: cpu.user - previous.user, system: cpu.system - previous.system };
|
||||
});
|
||||
vi.spyOn(process, "threadCpuUsage").mockImplementation(threadCpuProbe);
|
||||
records = [];
|
||||
vi.spyOn(sessionLog, "isEnabled").mockReturnValue(true);
|
||||
vi.spyOn(sessionLog, "warn").mockImplementation((message, fields) => {
|
||||
|
|
@ -56,20 +92,69 @@ beforeEach(() => {
|
|||
});
|
||||
afterEach(() => {
|
||||
scheduler.onYield = undefined;
|
||||
scheduler.afterYield = undefined;
|
||||
setDiagnosticsEnabledForProcess(previousDiagnostics);
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
function controlProjectionClock() {
|
||||
function expectNoCpuFields(record: unknown, fields: readonly string[] = threadCpuFields) {
|
||||
for (const field of fields) {
|
||||
expect(record).not.toHaveProperty(field);
|
||||
}
|
||||
}
|
||||
|
||||
function controlProjectionWork(hooks?: { afterPreparation?: () => void; afterRow?: () => void }) {
|
||||
vi.spyOn(performance, "now").mockImplementation(() => clock);
|
||||
const load = sessionUtils.loadCombinedSessionStoreForGatewayCore;
|
||||
vi.spyOn(sessionUtils, "loadCombinedSessionStoreForGatewayCore").mockImplementation((...args) => {
|
||||
try {
|
||||
return load(...args);
|
||||
} finally {
|
||||
cpu.user += 650;
|
||||
cpu.system += 100;
|
||||
}
|
||||
});
|
||||
const readRowInputs = rowProjection.readSessionRowInputs;
|
||||
vi.spyOn(rowProjection, "readSessionRowInputs").mockImplementation((...args) => {
|
||||
try {
|
||||
return readRowInputs(...args);
|
||||
} finally {
|
||||
cpu.user += 750;
|
||||
cpu.system += 250;
|
||||
}
|
||||
});
|
||||
const materializeRow = rowProjection.materializeSessionRow;
|
||||
vi.spyOn(rowProjection, "materializeSessionRow").mockImplementation((...args) => {
|
||||
try {
|
||||
return materializeRow(...args);
|
||||
} finally {
|
||||
cpu.user += 750;
|
||||
cpu.system += 250;
|
||||
}
|
||||
});
|
||||
const presentRow = rowProjection.presentSessionRow;
|
||||
vi.spyOn(rowProjection, "presentSessionRow").mockImplementation((...args) => {
|
||||
try {
|
||||
return presentRow(...args);
|
||||
} finally {
|
||||
cpu.user += 750;
|
||||
cpu.system += 250;
|
||||
hooks?.afterRow?.();
|
||||
}
|
||||
});
|
||||
const read = titleReader.readSessionTitleFieldsFromTranscriptBatch;
|
||||
return vi
|
||||
.spyOn(titleReader, "readSessionTitleFieldsFromTranscriptBatch")
|
||||
.mockImplementation((...args) => {
|
||||
const result = read(...args);
|
||||
// Advance only the synchronous preparation interval; the real yield remains in production.
|
||||
clock += 20;
|
||||
return result;
|
||||
try {
|
||||
return read(...args);
|
||||
} finally {
|
||||
// Charge real synchronous work independently of the instrumentation's probe count.
|
||||
clock += 20;
|
||||
cpu.user += 1_250;
|
||||
cpu.system += 250;
|
||||
hooks?.afterPreparation?.();
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
|
|
@ -87,7 +172,7 @@ test.each(["channel-only", "slow-warning"])("attributes %s operations", async (m
|
|||
clock += catalogDelay;
|
||||
return undefined;
|
||||
};
|
||||
const projection = controlProjectionClock();
|
||||
const projection = controlProjectionWork();
|
||||
const trace = createDiagnosticTraceContext();
|
||||
const events: unknown[] = [];
|
||||
const diagnostics = channel("openclaw.session.list");
|
||||
|
|
@ -104,7 +189,11 @@ test.each(["channel-only", "slow-warning"])("attributes %s operations", async (m
|
|||
client,
|
||||
context,
|
||||
isWebchatConnect: () => true,
|
||||
respond: (...response) => responses.push(response),
|
||||
respond: (...response) => {
|
||||
cpu.user += 750;
|
||||
cpu.system += 375;
|
||||
responses.push(response);
|
||||
},
|
||||
});
|
||||
expect(responses).toEqual([[true, { subscribed: true, list: listed }, undefined, undefined]]);
|
||||
expect(context.subscribeSessionEvents).toHaveBeenCalledWith(client.connId);
|
||||
|
|
@ -118,6 +207,9 @@ test.each(["channel-only", "slow-warning"])("attributes %s operations", async (m
|
|||
handlerElapsedMs: 20 + catalogDelay,
|
||||
cacheRole: "projection-owner",
|
||||
prepareSyncMs: 20,
|
||||
storeLoadThreadCpuMs: 0.75,
|
||||
prepareThreadCpuMs: 1.5,
|
||||
rowThreadCpuMs: 3,
|
||||
projectionPasses: 1,
|
||||
selectedRowCount: 1,
|
||||
handlerOutcome: "returned",
|
||||
|
|
@ -127,11 +219,13 @@ test.each(["channel-only", "slow-warning"])("attributes %s operations", async (m
|
|||
operation: "sessions.subscribe",
|
||||
handlerElapsedMs: catalogDelay,
|
||||
cacheRole: "completed-hit",
|
||||
responseThreadCpuMs: 1.125,
|
||||
selectedRowCount: 1,
|
||||
handlerOutcome: "returned",
|
||||
responseOutcome: "ok",
|
||||
});
|
||||
expect(events[1]).not.toHaveProperty("projectionPasses");
|
||||
expectNoCpuFields(events[1], producerCpuFields);
|
||||
const serialized = JSON.stringify(events);
|
||||
for (const privateValue of [
|
||||
client.connId,
|
||||
|
|
@ -150,69 +244,168 @@ test.each(["channel-only", "slow-warning"])("attributes %s operations", async (m
|
|||
} finally {
|
||||
diagnostics.unsubscribe(collect);
|
||||
}
|
||||
threadCpuProbe.mockClear();
|
||||
await listSessions({ client, context, request });
|
||||
expect(events).toHaveLength(2);
|
||||
});
|
||||
});
|
||||
|
||||
test("captures a fast failed projection while preserving the original error", async () => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async () => {
|
||||
const context = requestContext(await seedSessions());
|
||||
setDiagnosticsEnabledForProcess(false);
|
||||
vi.spyOn(performance, "now").mockImplementation(() => clock);
|
||||
const failure = new Error("synthetic-private-projection-error");
|
||||
vi.spyOn(titleReader, "readSessionTitleFieldsFromTranscriptBatch").mockImplementation(() => {
|
||||
clock += 25;
|
||||
throw failure;
|
||||
});
|
||||
const events: unknown[] = [];
|
||||
const diagnostics = channel("openclaw.session.list");
|
||||
const collect = (event: unknown) => events.push(event);
|
||||
diagnostics.subscribe(collect);
|
||||
try {
|
||||
await expect(
|
||||
listSessions({
|
||||
client: identifiedClient("owner@example.com"),
|
||||
context,
|
||||
request: { agentId: "main", limit: 1, includeDerivedTitles: true },
|
||||
}),
|
||||
).rejects.toBe(failure);
|
||||
expect(events).toHaveLength(1);
|
||||
expect(events[0]).toMatchObject({
|
||||
operation: "sessions.list",
|
||||
handlerElapsedMs: 25,
|
||||
cacheRole: "projection-owner",
|
||||
handlerOutcome: "threw",
|
||||
responseOutcome: "none",
|
||||
});
|
||||
expect(JSON.stringify(events)).not.toContain(failure.message);
|
||||
expect(sessionLog.warn).not.toHaveBeenCalled();
|
||||
} finally {
|
||||
diagnostics.unsubscribe(collect);
|
||||
if (warn) {
|
||||
expect(threadCpuProbe).toHaveBeenCalled();
|
||||
} else {
|
||||
expect(threadCpuProbe).not.toHaveBeenCalled();
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
test("separates producer work, follower wait, and completed hits under their own request traces", async () => {
|
||||
test.each([
|
||||
{ stage: "projection", cpuFailure: "none" },
|
||||
{ stage: "response", cpuFailure: "none" },
|
||||
{ stage: "projection", cpuFailure: "start" },
|
||||
{ stage: "projection", cpuFailure: "finish" },
|
||||
] as const)(
|
||||
"preserves a fast $stage error when CPU probe failure is $cpuFailure",
|
||||
async ({ stage, cpuFailure }) => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async () => {
|
||||
const context = requestContext(await seedSessions());
|
||||
setDiagnosticsEnabledForProcess(false);
|
||||
vi.spyOn(performance, "now").mockImplementation(() => clock);
|
||||
if (cpuFailure !== "none") {
|
||||
controlProjectionWork();
|
||||
}
|
||||
if (cpuFailure === "start") {
|
||||
cpuProbeFailure = new Error("synthetic CPU probe failure");
|
||||
}
|
||||
const failure = new Error("synthetic-private-projection-error");
|
||||
const fail = () => {
|
||||
clock += 25;
|
||||
cpu.user += 250;
|
||||
cpu.system += 125;
|
||||
if (cpuFailure === "finish") {
|
||||
cpuProbeFailure = new Error("synthetic CPU probe failure");
|
||||
}
|
||||
throw failure;
|
||||
};
|
||||
if (stage === "projection") {
|
||||
vi.spyOn(titleReader, "readSessionTitleFieldsFromTranscriptBatch").mockImplementation(fail);
|
||||
}
|
||||
const events: unknown[] = [];
|
||||
const diagnostics = channel("openclaw.session.list");
|
||||
const collect = (event: unknown) => events.push(event);
|
||||
diagnostics.subscribe(collect);
|
||||
try {
|
||||
await expect(
|
||||
stage === "projection"
|
||||
? listSessions({
|
||||
client: identifiedClient("owner@example.com"),
|
||||
context,
|
||||
request: { agentId: "main", limit: 1, includeDerivedTitles: true },
|
||||
})
|
||||
: sessionReadHandlers["sessions.list"]!({
|
||||
req: { type: "req", id: "private-request", method: "sessions.list" },
|
||||
params: { agentId: "main", limit: 1 },
|
||||
client: identifiedClient("owner@example.com"),
|
||||
context,
|
||||
respond: fail,
|
||||
isWebchatConnect: () => true,
|
||||
}),
|
||||
).rejects.toBe(failure);
|
||||
expect(events).toHaveLength(1);
|
||||
expect(events[0]).toMatchObject({
|
||||
operation: "sessions.list",
|
||||
handlerElapsedMs: 25,
|
||||
cacheRole: "projection-owner",
|
||||
handlerOutcome: "threw",
|
||||
responseOutcome: stage === "projection" ? "none" : "threw",
|
||||
});
|
||||
if (cpuFailure === "none") {
|
||||
expect(events[0]).toMatchObject({
|
||||
[stage === "projection" ? "prepareThreadCpuMs" : "responseThreadCpuMs"]: 0.375,
|
||||
});
|
||||
} else {
|
||||
expect(threadCpuProbe.mock.results).toContainEqual({
|
||||
type: "throw",
|
||||
value: cpuProbeFailure,
|
||||
});
|
||||
if (cpuFailure === "finish") {
|
||||
expect(
|
||||
threadCpuProbe.mock.results.some(
|
||||
(reading) =>
|
||||
reading.type === "return" && reading.value.user + reading.value.system > 0,
|
||||
),
|
||||
).toBe(true);
|
||||
}
|
||||
expectNoCpuFields(events[0]);
|
||||
}
|
||||
if (stage === "projection") {
|
||||
expectNoCpuFields(events[0], [
|
||||
"rowThreadCpuMs",
|
||||
"cachePublicationThreadCpuMs",
|
||||
"responseThreadCpuMs",
|
||||
]);
|
||||
}
|
||||
expect(JSON.stringify(events)).not.toContain(failure.message);
|
||||
expect(sessionLog.warn).not.toHaveBeenCalled();
|
||||
} finally {
|
||||
diagnostics.unsubscribe(collect);
|
||||
}
|
||||
});
|
||||
},
|
||||
);
|
||||
|
||||
test("separates producer CPU from registry readiness, yielded work, and followers", async () => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async () => {
|
||||
const config = await seedSessions();
|
||||
const context = requestContext(config);
|
||||
const client = identifiedClient("owner@example.com");
|
||||
const request = { agentId: "main", limit: 1, includeDerivedTitles: true };
|
||||
const request = { agentId: "main", limit: 2, includeDerivedTitles: true };
|
||||
const catalog = vi.fn(async () => undefined);
|
||||
context.readPreparedGatewayModelCatalog = catalog;
|
||||
const projection = controlProjectionClock();
|
||||
let projectedRows = 0;
|
||||
const rowsAtYield: number[] = [];
|
||||
const projection = controlProjectionWork({
|
||||
afterRow: () => {
|
||||
projectedRows++;
|
||||
clock += 20;
|
||||
},
|
||||
});
|
||||
context.workerPlacementDiskSpaceReader = {
|
||||
read: () => undefined,
|
||||
version: () => {
|
||||
cpu.user += 1_250;
|
||||
return 0;
|
||||
},
|
||||
};
|
||||
const entered = createDeferredCore();
|
||||
const release = createDeferredCore();
|
||||
const registryEntered = createDeferredCore();
|
||||
const releaseRegistry = createDeferredCore();
|
||||
const prepareRegistry = registryRead.prepareSubagentSessionListReadIndex;
|
||||
vi.spyOn(registryRead, "prepareSubagentSessionListReadIndex").mockImplementation(
|
||||
async (...args) => {
|
||||
const work = await prepareRegistry(...args);
|
||||
registryEntered.resolve();
|
||||
await releaseRegistry.promise;
|
||||
return work;
|
||||
},
|
||||
);
|
||||
scheduler.onYield = async () => {
|
||||
rowsAtYield.push(projectedRows);
|
||||
entered.resolve();
|
||||
await release.promise;
|
||||
};
|
||||
const unrelatedWork = createDeferredCore();
|
||||
scheduler.afterYield = () => {
|
||||
cpu.user += 900_000;
|
||||
cpu.system += 100_000;
|
||||
unrelatedWork.resolve();
|
||||
};
|
||||
const ownerTrace = createDiagnosticTraceContext();
|
||||
const followerTrace = createDiagnosticTraceContext();
|
||||
const owner = runWithDiagnosticTraceContext(ownerTrace, () =>
|
||||
listSessions({ client, context, request }),
|
||||
);
|
||||
await registryEntered.promise;
|
||||
cpu.user += 400_000;
|
||||
cpu.system += 100_000;
|
||||
releaseRegistry.resolve();
|
||||
await entered.promise;
|
||||
const follower = runWithDiagnosticTraceContext(followerTrace, () =>
|
||||
listSessions({ client, context, request }),
|
||||
|
|
@ -224,7 +417,10 @@ test("separates producer work, follower wait, and completed hits under their own
|
|||
clock += 1_500;
|
||||
release.resolve();
|
||||
const [owned, followed] = await Promise.all([owner, follower]);
|
||||
await unrelatedWork.promise;
|
||||
expect(followed).toBe(owned);
|
||||
expect(rowsAtYield).toEqual([0, 1]);
|
||||
expect(owned.sessions).toHaveLength(2);
|
||||
expect(projection).toHaveBeenCalledOnce();
|
||||
expect(records).toHaveLength(2);
|
||||
const ownerRecord = records.find((record) => record.trace?.traceId === ownerTrace.traceId);
|
||||
|
|
@ -239,18 +435,24 @@ test("separates producer work, follower wait, and completed hits under their own
|
|||
threadId,
|
||||
isMainThread,
|
||||
prepareSyncMs: 20,
|
||||
rowSyncMs: 0,
|
||||
storeLoadThreadCpuMs: 0.75,
|
||||
prepareThreadCpuMs: 1.5,
|
||||
rowThreadCpuMs: 6,
|
||||
cacheSelectionThreadCpuMs: 1.25,
|
||||
cachePublicationThreadCpuMs: 1.25,
|
||||
rowSyncMs: 40,
|
||||
yieldWaitMs: 1_500,
|
||||
yieldCount: 1,
|
||||
yieldCount: 2,
|
||||
projectionPasses: 1,
|
||||
selectedRowCount: 1,
|
||||
selectedRowCount: 2,
|
||||
},
|
||||
});
|
||||
expect(followerRecord).toMatchObject({
|
||||
trace: followerTrace,
|
||||
fields: {
|
||||
cacheRole: "in-flight-follower",
|
||||
selectedRowCount: 1,
|
||||
cacheSelectionThreadCpuMs: 1.25,
|
||||
selectedRowCount: 2,
|
||||
workTraceId: ownerTrace.traceId,
|
||||
workSpanId: ownerTrace.spanId,
|
||||
},
|
||||
|
|
@ -258,6 +460,7 @@ test("separates producer work, follower wait, and completed hits under their own
|
|||
expect(followerRecord?.fields).not.toHaveProperty("prepareSyncMs");
|
||||
expect(followerRecord?.fields).not.toHaveProperty("yieldWaitMs");
|
||||
expect(followerRecord?.fields).not.toHaveProperty("projectionPasses");
|
||||
expectNoCpuFields(followerRecord?.fields, producerCpuFields);
|
||||
|
||||
scheduler.onYield = undefined;
|
||||
catalog.mockImplementation(async () => {
|
||||
|
|
@ -273,11 +476,12 @@ test("separates producer work, follower wait, and completed hits under their own
|
|||
expect(records).toHaveLength(3);
|
||||
expect(records[2]).toMatchObject({
|
||||
trace: hitTrace,
|
||||
fields: { cacheRole: "completed-hit", selectedRowCount: 1 },
|
||||
fields: { cacheRole: "completed-hit", selectedRowCount: 2 },
|
||||
});
|
||||
expect(records[2]?.fields).not.toHaveProperty("projectionPasses");
|
||||
expect(records[2]?.fields).not.toHaveProperty("workTraceId");
|
||||
expect(records[2]?.fields).not.toHaveProperty("rowSyncMs");
|
||||
expectNoCpuFields(records[2]?.fields, producerCpuFields);
|
||||
});
|
||||
});
|
||||
|
||||
|
|
@ -299,7 +503,7 @@ test("accumulates the bounded visibility repairs without counting yielded waits
|
|||
);
|
||||
}
|
||||
updatedAt.mockRestore();
|
||||
controlProjectionClock();
|
||||
controlProjectionWork();
|
||||
let pass = 0;
|
||||
scheduler.onYield = async () => {
|
||||
const name = ["first", "second", "third"][pass++];
|
||||
|
|
@ -310,6 +514,7 @@ test("accumulates the bounded visibility repairs without counting yielded waits
|
|||
);
|
||||
}
|
||||
clock += 500;
|
||||
cpu.user += 300_000;
|
||||
};
|
||||
const result = await listSessions({
|
||||
client: identifiedClient("viewer@example.com"),
|
||||
|
|
@ -324,6 +529,9 @@ test("accumulates the bounded visibility repairs without counting yielded waits
|
|||
rowRepairCount: 2,
|
||||
fullReloadCount: 1,
|
||||
prepareSyncMs: 80,
|
||||
storeLoadThreadCpuMs: 1.5,
|
||||
prepareThreadCpuMs: 6,
|
||||
rowThreadCpuMs: 12,
|
||||
rowSyncMs: 0,
|
||||
yieldWaitMs: 2_000,
|
||||
yieldCount: 4,
|
||||
|
|
@ -332,44 +540,77 @@ test("accumulates the bounded visibility repairs without counting yielded waits
|
|||
});
|
||||
});
|
||||
|
||||
test.each(["disabled", "sink-disabled", "sink-throws", "disabled-during-request"])(
|
||||
"preserves the response when diagnostics are %s",
|
||||
async (mode) => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async () => {
|
||||
const context = requestContext(await seedSessions());
|
||||
if (mode === "disabled") {
|
||||
test.each([
|
||||
"disabled",
|
||||
"sink-disabled",
|
||||
"sink-throws",
|
||||
"disabled-during-request",
|
||||
"cpu-start-throws",
|
||||
"cpu-finish-throws",
|
||||
])("preserves the response when diagnostics are %s", async (mode) => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async () => {
|
||||
const context = requestContext(await seedSessions());
|
||||
threadCpuProbe.mockClear();
|
||||
const cpuThrows = mode === "cpu-start-throws" || mode === "cpu-finish-throws";
|
||||
controlProjectionWork({
|
||||
afterPreparation: () => {
|
||||
if (mode === "cpu-finish-throws") {
|
||||
cpuProbeFailure = new Error("synthetic CPU probe failure");
|
||||
}
|
||||
},
|
||||
});
|
||||
if (mode === "cpu-start-throws") {
|
||||
cpuProbeFailure = new Error("synthetic CPU probe failure");
|
||||
}
|
||||
if (mode === "disabled") {
|
||||
setDiagnosticsEnabledForProcess(false);
|
||||
}
|
||||
if (mode === "sink-disabled") {
|
||||
vi.mocked(sessionLog.isEnabled).mockReturnValue(false);
|
||||
}
|
||||
if (mode === "sink-throws") {
|
||||
vi.mocked(sessionLog.warn).mockImplementation(() => {
|
||||
throw new Error("synthetic sink failure");
|
||||
});
|
||||
}
|
||||
context.readPreparedGatewayModelCatalog = async () => {
|
||||
clock += 1_100;
|
||||
if (mode === "disabled-during-request") {
|
||||
setDiagnosticsEnabledForProcess(false);
|
||||
}
|
||||
if (mode === "sink-disabled") {
|
||||
vi.mocked(sessionLog.isEnabled).mockReturnValue(false);
|
||||
}
|
||||
if (mode === "sink-throws") {
|
||||
vi.mocked(sessionLog.warn).mockImplementation(() => {
|
||||
throw new Error("synthetic sink failure");
|
||||
});
|
||||
}
|
||||
vi.spyOn(performance, "now").mockImplementation(() => clock);
|
||||
context.readPreparedGatewayModelCatalog = async () => {
|
||||
clock += 1_100;
|
||||
if (mode === "disabled-during-request") {
|
||||
setDiagnosticsEnabledForProcess(false);
|
||||
}
|
||||
return undefined;
|
||||
};
|
||||
const result = await listSessions({
|
||||
client: identifiedClient("owner@example.com"),
|
||||
context,
|
||||
request: { agentId: "main", limit: 1 },
|
||||
});
|
||||
expect(result.sessions).toHaveLength(1);
|
||||
if (mode === "sink-throws") {
|
||||
expect(sessionLog.warn).toHaveBeenCalledOnce();
|
||||
} else {
|
||||
expect(sessionLog.warn).not.toHaveBeenCalled();
|
||||
}
|
||||
return undefined;
|
||||
};
|
||||
const result = await listSessions({
|
||||
client: identifiedClient("owner@example.com"),
|
||||
context,
|
||||
request: { agentId: "main", limit: 1, includeDerivedTitles: true },
|
||||
});
|
||||
},
|
||||
);
|
||||
expect(result.sessions).toHaveLength(1);
|
||||
if (mode === "sink-throws" || cpuThrows) {
|
||||
expect(sessionLog.warn).toHaveBeenCalledOnce();
|
||||
} else {
|
||||
expect(sessionLog.warn).not.toHaveBeenCalled();
|
||||
}
|
||||
if (mode === "disabled" || mode === "sink-disabled") {
|
||||
expect(threadCpuProbe).not.toHaveBeenCalled();
|
||||
}
|
||||
if (cpuThrows) {
|
||||
expect(threadCpuProbe.mock.results).toContainEqual({
|
||||
type: "throw",
|
||||
value: cpuProbeFailure,
|
||||
});
|
||||
if (mode === "cpu-finish-throws") {
|
||||
expect(
|
||||
threadCpuProbe.mock.results.some(
|
||||
(reading) => reading.type === "return" && reading.value.user + reading.value.system > 0,
|
||||
),
|
||||
).toBe(true);
|
||||
}
|
||||
expect(records).toHaveLength(1);
|
||||
expectNoCpuFields(records[0]?.fields);
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
test("preserves the original projection error even when its slow diagnostic sink throws", async () => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async () => {
|
||||
|
|
@ -378,6 +619,8 @@ test("preserves the original projection error even when its slow diagnostic sink
|
|||
const failure = new Error("synthetic projection failure");
|
||||
vi.spyOn(titleReader, "readSessionTitleFieldsFromTranscriptBatch").mockImplementation(() => {
|
||||
clock += 1_500;
|
||||
cpu.user += 500;
|
||||
cpu.system += 125;
|
||||
throw failure;
|
||||
});
|
||||
vi.mocked(sessionLog.warn).mockImplementation(() => {
|
||||
|
|
@ -391,5 +634,9 @@ test("preserves the original projection error even when its slow diagnostic sink
|
|||
}),
|
||||
).rejects.toBe(failure);
|
||||
expect(sessionLog.warn).toHaveBeenCalledOnce();
|
||||
expect(sessionLog.warn).toHaveBeenCalledWith(
|
||||
"slow session list",
|
||||
expect.objectContaining({ prepareThreadCpuMs: 0.625, handlerOutcome: "threw" }),
|
||||
);
|
||||
});
|
||||
});
|
||||
|
|
|
|||
|
|
@ -272,13 +272,19 @@ export const sessionReadHandlers: GatewayRequestHandlers = {
|
|||
if (!loaded) {
|
||||
const loadedStore = measureDiagnosticsTimelineSpanSync(
|
||||
"gateway.sessions.list.store_load",
|
||||
() =>
|
||||
loadCombinedSessionStoreForGatewayCore(cfg, {
|
||||
agentId: p.agentId,
|
||||
configuredAgentsOnly,
|
||||
projection: "list",
|
||||
...(p.activeOnly === true ? { preserveSentinelOwners: true } : {}),
|
||||
}),
|
||||
() => {
|
||||
const storeCpu = diagnostics?.startSyncCpu();
|
||||
try {
|
||||
return loadCombinedSessionStoreForGatewayCore(cfg, {
|
||||
agentId: p.agentId,
|
||||
configuredAgentsOnly,
|
||||
projection: "list",
|
||||
...(p.activeOnly === true ? { preserveSentinelOwners: true } : {}),
|
||||
});
|
||||
} finally {
|
||||
diagnostics?.finishSyncCpu("storeLoadThreadCpuMs", storeCpu);
|
||||
}
|
||||
},
|
||||
{
|
||||
config: cfg,
|
||||
phase: "sessions.list",
|
||||
|
|
@ -337,6 +343,7 @@ export const sessionReadHandlers: GatewayRequestHandlers = {
|
|||
cfg,
|
||||
workStartedAt,
|
||||
projectionTiming,
|
||||
cpuTiming: diagnostics,
|
||||
durableStorePath,
|
||||
...(entryFilter ? { entryFilter } : {}),
|
||||
storePath,
|
||||
|
|
|
|||
|
|
@ -62,6 +62,14 @@ export type SessionListProjectionTiming = {
|
|||
yieldCount: number;
|
||||
};
|
||||
|
||||
type SessionListCpuTiming = {
|
||||
startSyncCpu: () => NodeJS.CpuUsage | undefined;
|
||||
finishSyncCpu: (
|
||||
metric: "prepareThreadCpuMs" | "rowThreadCpuMs",
|
||||
started: NodeJS.CpuUsage | undefined,
|
||||
) => void;
|
||||
};
|
||||
|
||||
type SessionSelectionScope =
|
||||
| { opts: SessionsListParams; targetsBySessionKey: GatewayStoredSessionTargets }
|
||||
| {
|
||||
|
|
@ -336,6 +344,7 @@ export async function listSessionsFromStoreAsync(
|
|||
params: ListSessionsFromStoreParams & {
|
||||
workStartedAt?: number;
|
||||
projectionTiming?: SessionListProjectionTiming;
|
||||
cpuTiming?: SessionListCpuTiming;
|
||||
},
|
||||
): Promise<SessionsListResult> {
|
||||
const stateContext = captureOpenClawStateWorkerContext();
|
||||
|
|
@ -347,6 +356,7 @@ export async function listSessionsFromStoreAsync(
|
|||
return withPinnedActivePluginRegistryWorkspaceDir(() =>
|
||||
withSessionProjectionWorkBudget(async (budget) => {
|
||||
const timing = params.projectionTiming;
|
||||
const cpuTiming = params.cpuTiming;
|
||||
let syncStartedAt = timing ? performance.now() : 0;
|
||||
let syncPhase: "prepareSyncMs" | "rowSyncMs" | undefined = "prepareSyncMs";
|
||||
const yieldIfNeeded = (): Promise<void> | undefined => {
|
||||
|
|
@ -379,8 +389,17 @@ export async function listSessionsFromStoreAsync(
|
|||
budget.yieldIfNeeded,
|
||||
);
|
||||
// Each chunk shares roster facts, then releases them before another request can run.
|
||||
let step = withAgentRosterFactsBatch(cfg, () => preparation.next());
|
||||
while (!step.done) {
|
||||
let step: ReturnType<typeof preparation.next>;
|
||||
while (true) {
|
||||
const chunkCpu = cpuTiming?.startSyncCpu();
|
||||
try {
|
||||
step = withAgentRosterFactsBatch(cfg, () => preparation.next());
|
||||
} finally {
|
||||
cpuTiming?.finishSyncCpu("prepareThreadCpuMs", chunkCpu);
|
||||
}
|
||||
if (step.done) {
|
||||
break;
|
||||
}
|
||||
if (step.value) {
|
||||
const checkpoint = performance.now();
|
||||
if (timing) {
|
||||
|
|
@ -396,13 +415,15 @@ export async function listSessionsFromStoreAsync(
|
|||
if (pause) {
|
||||
await pause;
|
||||
}
|
||||
step = withAgentRosterFactsBatch(cfg, () => preparation.next());
|
||||
}
|
||||
const list = step.value;
|
||||
const sessions: GatewaySessionRow[] = [];
|
||||
const includeTranscriptFields = list.includeDerivedTitles || list.includeLastMessage;
|
||||
const transcriptFields = includeTranscriptFields
|
||||
? readScopedSessionTitleFieldsFromTranscriptBatch(
|
||||
let transcriptFields: ReturnType<typeof readScopedSessionTitleFieldsFromTranscriptBatch>;
|
||||
if (includeTranscriptFields) {
|
||||
const transcriptCpu = cpuTiming?.startSyncCpu();
|
||||
try {
|
||||
transcriptFields = readScopedSessionTitleFieldsFromTranscriptBatch(
|
||||
list.entries.slice(0, list.transcriptFieldRows).flatMap(([key, entry]) => {
|
||||
if (!entry.sessionId) {
|
||||
return [];
|
||||
|
|
@ -417,8 +438,13 @@ export async function listSessionsFromStoreAsync(
|
|||
},
|
||||
];
|
||||
}),
|
||||
)
|
||||
: [];
|
||||
);
|
||||
} finally {
|
||||
cpuTiming?.finishSyncCpu("prepareThreadCpuMs", transcriptCpu);
|
||||
}
|
||||
} else {
|
||||
transcriptFields = [];
|
||||
}
|
||||
// Optional transcript reads can spend the remaining budget even for an empty page.
|
||||
const checkpoint = performance.now();
|
||||
if (timing) {
|
||||
|
|
@ -433,56 +459,67 @@ export async function listSessionsFromStoreAsync(
|
|||
let transcriptFieldIndex = 0;
|
||||
for (let nextRowIndex = 0; nextRowIndex < list.entries.length;) {
|
||||
// Release roster facts before a pause so resumed rows observe current entries.
|
||||
const pause = withAgentRosterFactsBatch(cfg, () => {
|
||||
while (nextRowIndex < list.entries.length) {
|
||||
const i = nextRowIndex++;
|
||||
const [key, entry] = expectDefined(list.entries[i], "entries entry at i");
|
||||
const target = expectDefined(targetsBySessionKey.get(key), "session row owner");
|
||||
const { inputs, presentation } = readSessionRowInputs({
|
||||
cfg,
|
||||
storePath: target.storeTarget.storePath, // Aggregate paths are display-only.
|
||||
store,
|
||||
modelSource: target,
|
||||
key: target.storeKey ?? key,
|
||||
entry,
|
||||
agentId: target.agentId,
|
||||
modelCatalog: params.modelCatalog,
|
||||
now: list.now,
|
||||
storeChildSessionLinksByKey: list.storeChildSessionLinksByKey,
|
||||
excludedChildKeys: list.excludedChildKeys,
|
||||
rowContext: list.rowContext,
|
||||
configuredAgentIds: list.configuredAgentIds,
|
||||
skipTranscriptUsageFallback: true,
|
||||
lightweightListRow: true,
|
||||
});
|
||||
const row = presentSessionRow(materializeSessionRow(inputs), presentation);
|
||||
row.key = key;
|
||||
if (entry?.sessionId && i < list.transcriptFieldRows && includeTranscriptFields) {
|
||||
const { firstUserMessage, lastMessagePreview } = expectDefined(
|
||||
transcriptFields[transcriptFieldIndex++],
|
||||
"batched transcript fields at transcriptFieldIndex",
|
||||
);
|
||||
if (list.includeDerivedTitles) {
|
||||
row.derivedTitle = deriveSessionTitle(entry, firstUserMessage, row.displayName);
|
||||
let pause: Promise<void> | undefined;
|
||||
const rowCpu = cpuTiming?.startSyncCpu();
|
||||
try {
|
||||
pause = withAgentRosterFactsBatch(cfg, () => {
|
||||
while (nextRowIndex < list.entries.length) {
|
||||
const i = nextRowIndex++;
|
||||
const [key, entry] = expectDefined(list.entries[i], "entries entry at i");
|
||||
const target = expectDefined(targetsBySessionKey.get(key), "session row owner");
|
||||
const { inputs, presentation } = readSessionRowInputs({
|
||||
cfg,
|
||||
storePath: target.storeTarget.storePath, // Aggregate paths are display-only.
|
||||
store,
|
||||
modelSource: target,
|
||||
key: target.storeKey ?? key,
|
||||
entry,
|
||||
agentId: target.agentId,
|
||||
modelCatalog: params.modelCatalog,
|
||||
now: list.now,
|
||||
storeChildSessionLinksByKey: list.storeChildSessionLinksByKey,
|
||||
excludedChildKeys: list.excludedChildKeys,
|
||||
rowContext: list.rowContext,
|
||||
configuredAgentIds: list.configuredAgentIds,
|
||||
skipTranscriptUsageFallback: true,
|
||||
lightweightListRow: true,
|
||||
});
|
||||
const row = presentSessionRow(materializeSessionRow(inputs), presentation);
|
||||
row.key = key;
|
||||
if (entry?.sessionId && i < list.transcriptFieldRows && includeTranscriptFields) {
|
||||
const { firstUserMessage, lastMessagePreview } = expectDefined(
|
||||
transcriptFields[transcriptFieldIndex++],
|
||||
"batched transcript fields at transcriptFieldIndex",
|
||||
);
|
||||
if (list.includeDerivedTitles) {
|
||||
row.derivedTitle = deriveSessionTitle(entry, firstUserMessage, row.displayName);
|
||||
}
|
||||
if (list.includeLastMessage && lastMessagePreview) {
|
||||
row.lastMessagePreview = lastMessagePreview;
|
||||
}
|
||||
}
|
||||
if (list.includeLastMessage && lastMessagePreview) {
|
||||
row.lastMessagePreview = lastMessagePreview;
|
||||
sessions.push(row);
|
||||
const rowPause = nextRowIndex < list.entries.length ? yieldIfNeeded() : undefined;
|
||||
if (rowPause) {
|
||||
return rowPause;
|
||||
}
|
||||
}
|
||||
sessions.push(row);
|
||||
const rowPause = nextRowIndex < list.entries.length ? yieldIfNeeded() : undefined;
|
||||
if (rowPause) {
|
||||
return rowPause;
|
||||
}
|
||||
}
|
||||
return undefined;
|
||||
});
|
||||
return undefined;
|
||||
});
|
||||
} finally {
|
||||
cpuTiming?.finishSyncCpu("rowThreadCpuMs", rowCpu);
|
||||
}
|
||||
if (pause) {
|
||||
await pause;
|
||||
}
|
||||
}
|
||||
|
||||
return buildSessionsListResult(params, list, sessions);
|
||||
const resultCpu = cpuTiming?.startSyncCpu();
|
||||
try {
|
||||
return buildSessionsListResult(params, list, sessions);
|
||||
} finally {
|
||||
cpuTiming?.finishSyncCpu("rowThreadCpuMs", resultCpu);
|
||||
}
|
||||
} finally {
|
||||
if (timing && syncPhase) {
|
||||
timing[syncPhase] += performance.now() - syncStartedAt;
|
||||
|
|
|
|||
|
|
@ -4,6 +4,7 @@ import { startGatewayBenchDiagnostics } from "../../scripts/lib/gateway-bench-di
|
|||
|
||||
it("aggregates captured work with bounded memory and drops content-bearing fields", () => {
|
||||
const source = channel("openclaw.session.write");
|
||||
const lists = channel("openclaw.session.list");
|
||||
const finish = startGatewayBenchDiagnostics();
|
||||
let result;
|
||||
try {
|
||||
|
|
@ -18,6 +19,17 @@ it("aggregates captured work with bounded memory and drops content-bearing field
|
|||
agentId: "private-agent",
|
||||
error: "private-error",
|
||||
});
|
||||
lists.publish({
|
||||
operation: "sessions.list",
|
||||
responseOutcome: "ok",
|
||||
cacheRole: "projection-owner",
|
||||
storeLoadThreadCpuMs: 0.75 * (index + 1),
|
||||
prepareThreadCpuMs: 1.5 * (index + 1),
|
||||
rowThreadCpuMs: 2.25 * (index + 1),
|
||||
cacheSelectionThreadCpuMs: 0.125 * (index + 1),
|
||||
cachePublicationThreadCpuMs: 0.25 * (index + 1),
|
||||
responseThreadCpuMs: 0.375 * (index + 1),
|
||||
});
|
||||
}
|
||||
for (let index = 0; index < 300; index += 1) {
|
||||
source.publish({ operation: `session.synthetic-${index}`, outcome: "ok", elapsedMs: 1 });
|
||||
|
|
@ -26,7 +38,7 @@ it("aggregates captured work with bounded memory and drops content-bearing field
|
|||
result = finish();
|
||||
}
|
||||
expect(result.groups).toHaveLength(256);
|
||||
expect(result.droppedEvents).toBe(45);
|
||||
expect(result.droppedEvents).toBe(46);
|
||||
expect(result.collectionErrors).toBe(0);
|
||||
expect(result.groups[0]).toMatchObject({
|
||||
operation: "session-entry.patch",
|
||||
|
|
@ -36,6 +48,19 @@ it("aggregates captured work with bounded memory and drops content-bearing field
|
|||
writerExecutionMs: { count: 2, total: 4, max: 2 },
|
||||
},
|
||||
});
|
||||
expect(result.groups[1]).toMatchObject({
|
||||
operation: "sessions.list",
|
||||
cacheRole: "projection-owner",
|
||||
count: 2,
|
||||
metrics: {
|
||||
storeLoadThreadCpuMs: { count: 2, total: 2.25, max: 1.5 },
|
||||
prepareThreadCpuMs: { count: 2, total: 4.5, max: 3 },
|
||||
rowThreadCpuMs: { count: 2, total: 6.75, max: 4.5 },
|
||||
cacheSelectionThreadCpuMs: { count: 2, total: 0.375, max: 0.25 },
|
||||
cachePublicationThreadCpuMs: { count: 2, total: 0.75, max: 0.5 },
|
||||
responseThreadCpuMs: { count: 2, total: 1.125, max: 0.75 },
|
||||
},
|
||||
});
|
||||
expect(JSON.stringify(result)).not.toContain("private");
|
||||
source.publish({ operation: "session-entry.patch", outcome: "ok", queueWaitMs: 999 });
|
||||
expect(result.groups[0]).toMatchObject({ count: 2 });
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue