mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
perf(logbook): avoid transferring unused frame metadata (#153162)
This commit is contained in:
parent
c5ad7506ab
commit
325e9ca2f9
5 changed files with 157 additions and 31 deletions
|
|
@ -1,12 +1,17 @@
|
|||
import path from "node:path";
|
||||
import { fileURLToPath, pathToFileURL } from "node:url";
|
||||
import type {
|
||||
OpenClawPluginApi,
|
||||
OpenClawPluginNodeInvokePolicy,
|
||||
OpenClawPluginService,
|
||||
} from "openclaw/plugin-sdk/plugin-entry";
|
||||
import { createCapturedPluginRegistration } from "openclaw/plugin-sdk/plugin-test-runtime";
|
||||
import { resolveRuntimeWorkerUrl } from "openclaw/plugin-sdk/process-runtime";
|
||||
import { useAutoCleanupTempDirTracker } from "openclaw/plugin-sdk/test-env";
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
import plugin from "./index.js";
|
||||
import { dayKeyFor } from "./src/day.js";
|
||||
import { logbookSqliteBackendEntrypoint } from "./src/sqlite-backend-entrypoint.test-support.js";
|
||||
import { LogbookStore } from "./src/store.js";
|
||||
|
||||
type PolicyContext = Parameters<OpenClawPluginNodeInvokePolicy["handle"]>[0];
|
||||
|
|
@ -16,19 +21,110 @@ function registerLogbook(runtimeSource = fileURLToPath(new URL("./index.ts", imp
|
|||
captured.api.pluginConfig = { captureEnabled: false };
|
||||
const policies: OpenClawPluginNodeInvokePolicy[] = [];
|
||||
const services: OpenClawPluginService[] = [];
|
||||
const methods: Array<{ method: string; options: unknown }> = [];
|
||||
const methods: Array<{
|
||||
method: string;
|
||||
handler: Parameters<OpenClawPluginApi["registerGatewayMethod"]>[1];
|
||||
options: unknown;
|
||||
}> = [];
|
||||
captured.api.registerNodeInvokePolicy = (policy) => policies.push(policy);
|
||||
captured.api.registerService = (service) => services.push(service);
|
||||
captured.api.registerGatewayMethod = (method, _handler, options) => {
|
||||
methods.push({ method, options });
|
||||
captured.api.registerGatewayMethod = (method, handler, options) => {
|
||||
methods.push({ method, handler, options });
|
||||
};
|
||||
plugin.register({ ...captured.api, runtimeSource });
|
||||
return { policies, services, methods };
|
||||
}
|
||||
|
||||
afterEach(() => vi.restoreAllMocks());
|
||||
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
|
||||
|
||||
describe("logbook gateway methods", () => {
|
||||
it("returns ordered frame metadata through the registered service and SQLite worker", async () => {
|
||||
const stateDir = tempDirs.make("logbook-range-rpc-");
|
||||
const { methods, services } = registerLogbook();
|
||||
const service = services[0]!;
|
||||
const handler = methods.find((entry) => entry.method === "logbook.frames")!.handler;
|
||||
const context = {
|
||||
config: {},
|
||||
stateDir,
|
||||
logger: { info() {}, warn() {}, error() {}, debug() {} },
|
||||
};
|
||||
const store = await LogbookStore.open(
|
||||
path.join(stateDir, "logbook"),
|
||||
resolveRuntimeWorkerUrl(logbookSqliteBackendEntrypoint),
|
||||
);
|
||||
vi.spyOn(LogbookStore, "open").mockResolvedValueOnce(store);
|
||||
try {
|
||||
await service.start(context);
|
||||
const startMs = Date.now();
|
||||
const day = dayKeyFor(startMs);
|
||||
const ids: number[] = [];
|
||||
for (const [capturedAtMs, idle] of [
|
||||
[startMs, false],
|
||||
[startMs, true],
|
||||
[startMs + 2000, false],
|
||||
] as const) {
|
||||
ids.push(
|
||||
await store.insertFrame({
|
||||
capturedAtMs,
|
||||
day,
|
||||
path: store.frameFilePath(day, capturedAtMs),
|
||||
screenIndex: 0,
|
||||
byteSize: 9,
|
||||
contentHash: `synthetic-${ids.length}`,
|
||||
idle,
|
||||
}),
|
||||
);
|
||||
}
|
||||
await store.createBatch({
|
||||
day,
|
||||
startMs,
|
||||
endMs: startMs + 3000,
|
||||
frameIds: [ids[0]!, ids[2]!],
|
||||
});
|
||||
const call = async (endMs: number) => {
|
||||
const params = { startMs, endMs };
|
||||
const respond = vi.fn();
|
||||
await handler({
|
||||
req: { type: "req", id: "range", method: "logbook.frames", params },
|
||||
params,
|
||||
client: null,
|
||||
isWebchatConnect: () => false,
|
||||
respond,
|
||||
get context(): never {
|
||||
throw new Error("The frame range handler does not use Gateway request context");
|
||||
},
|
||||
});
|
||||
expect(respond).toHaveBeenCalledTimes(1);
|
||||
return respond.mock.calls[0];
|
||||
};
|
||||
expect(await call(startMs + 2000)).toEqual([
|
||||
true,
|
||||
{
|
||||
frames: [
|
||||
{ id: ids[0], capturedAtMs: startMs, idle: false },
|
||||
{ id: ids[1], capturedAtMs: startMs, idle: true },
|
||||
],
|
||||
},
|
||||
]);
|
||||
expect(await call(startMs)).toEqual([true, { frames: [] }]);
|
||||
expect(await store.frameById(ids[0]!)).toMatchObject({
|
||||
day,
|
||||
path: store.frameFilePath(day, startMs),
|
||||
width: undefined,
|
||||
height: undefined,
|
||||
screenIndex: 0,
|
||||
byteSize: 9,
|
||||
});
|
||||
} finally {
|
||||
try {
|
||||
await service.stop?.(context);
|
||||
} finally {
|
||||
await store.close();
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
it("keeps only process-wide status independent of the authenticated profile", () => {
|
||||
const { methods } = registerLogbook();
|
||||
expect(methods.find((entry) => entry.method === "logbook.status")?.options).toEqual({
|
||||
|
|
|
|||
|
|
@ -226,11 +226,7 @@ export default definePluginEntry({
|
|||
registerWrite("logbook.frames", async (params) => {
|
||||
const startMs = readNumberParam(params, "startMs");
|
||||
const endMs = readNumberParam(params, "endMs");
|
||||
const frames = (await requireService().framesInRange(startMs, endMs)).map((frame) => ({
|
||||
id: frame.id,
|
||||
capturedAtMs: frame.capturedAtMs,
|
||||
idle: frame.idle,
|
||||
}));
|
||||
const frames = await requireService().framesInRange(startMs, endMs);
|
||||
return { frames };
|
||||
});
|
||||
|
||||
|
|
|
|||
|
|
@ -28,7 +28,10 @@ export type LogbookOperations = {
|
|||
unbatchedActiveFrames: Operation<{ limit: number }, Pick<LogbookFrame, "id" | "capturedAtMs">[]>;
|
||||
countUnbatchedActiveFrames: Operation<undefined, number>;
|
||||
frameById: Operation<{ id: number }, LogbookFrame | null>;
|
||||
framesInRange: Operation<{ startMs: number; endMs: number }, LogbookFrame[]>;
|
||||
framesInRange: Operation<
|
||||
{ startMs: number; endMs: number },
|
||||
Pick<LogbookFrame, "id" | "capturedAtMs" | "idle">[]
|
||||
>;
|
||||
createBatch: Operation<LogbookBatchInput, number>;
|
||||
setBatchStatus: Operation<
|
||||
{ batchId: number; status: LogbookBatchStatus; error?: string; model?: string },
|
||||
|
|
|
|||
|
|
@ -326,7 +326,7 @@ describe("Logbook native statement and read budgets", () => {
|
|||
});
|
||||
}
|
||||
|
||||
it("selects ordered pending frames without decoding unused text or narrowing full readers", () => {
|
||||
it("selects ordered pending and range metadata without decoding unused text or narrowing full readers", () => {
|
||||
const { backend } = openBackend();
|
||||
insertCandidate(backend, 3000);
|
||||
insertCandidate(backend, 1000);
|
||||
|
|
@ -353,6 +353,19 @@ describe("Logbook native statement and read budgets", () => {
|
|||
expect(pending).toEqual(expected.slice(0, limit));
|
||||
}
|
||||
|
||||
reads.frameRows = 0;
|
||||
reads.frameTextBytes = 0;
|
||||
const range = backend.execute({ type: "framesInRange", input: { startMs: 500, endMs: 3000 } });
|
||||
expect(reads.frameRows).toBe(5);
|
||||
expect.soft(reads.frameTextBytes).toBe(0);
|
||||
expect(range).toEqual([
|
||||
{ id: 5, capturedAtMs: 500, idle: true },
|
||||
{ id: 6, capturedAtMs: 750, idle: false },
|
||||
{ id: 2, capturedAtMs: 1000, idle: false },
|
||||
{ id: 3, capturedAtMs: 1000, idle: false },
|
||||
{ id: 4, capturedAtMs: 2000, idle: false },
|
||||
]);
|
||||
|
||||
const fullFrame = {
|
||||
id: 6,
|
||||
capturedAtMs: 750,
|
||||
|
|
@ -365,9 +378,6 @@ describe("Logbook native statement and read budgets", () => {
|
|||
idle: false,
|
||||
};
|
||||
expect(backend.execute({ type: "frameById", input: { id: 6 } })).toEqual(fullFrame);
|
||||
expect(backend.execute({ type: "framesInRange", input: { startMs: 750, endMs: 751 } })).toEqual(
|
||||
[fullFrame],
|
||||
);
|
||||
expect(backend.execute({ type: "batchFrames", input: { batchId: 1 } })).toEqual([fullFrame]);
|
||||
expect(backend.execute({ type: "sampledBatchFrames", input: { batchId: 1 } })).toEqual([
|
||||
fullFrame,
|
||||
|
|
@ -404,20 +414,18 @@ describe("Logbook native statement and read budgets", () => {
|
|||
});
|
||||
expect(frames).toHaveLength(count);
|
||||
if (count > 0) {
|
||||
expect(frames).toEqual(
|
||||
expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
id: 1,
|
||||
capturedAtMs: 0,
|
||||
day,
|
||||
path: `captures/${"nested/".repeat(12)}0.jpg`,
|
||||
width: 640,
|
||||
height: 480,
|
||||
byteSize: 10,
|
||||
screenIndex: 0,
|
||||
idle: false,
|
||||
}),
|
||||
]),
|
||||
expect(backend.execute({ type: "frameById", input: { id: 1 } })).toEqual(
|
||||
expect.objectContaining({
|
||||
id: 1,
|
||||
capturedAtMs: 0,
|
||||
day,
|
||||
path: `captures/${"nested/".repeat(12)}0.jpg`,
|
||||
width: 640,
|
||||
height: 480,
|
||||
byteSize: 10,
|
||||
screenIndex: 0,
|
||||
idle: false,
|
||||
}),
|
||||
);
|
||||
}
|
||||
},
|
||||
|
|
@ -446,7 +454,7 @@ describe("Logbook native statement and read budgets", () => {
|
|||
}
|
||||
let originalError: unknown;
|
||||
try {
|
||||
backend.execute({ type: "framesInRange", input: { startMs: 0, endMs: 120_000 } });
|
||||
backend.execute({ type: "frameById", input: { id: 1 } });
|
||||
} catch (error) {
|
||||
originalError = error;
|
||||
}
|
||||
|
|
@ -454,6 +462,17 @@ describe("Logbook native statement and read budgets", () => {
|
|||
throw new Error("Expected the full frame reader to reject the unsafe integer fixture");
|
||||
}
|
||||
expect(originalError).toMatchObject({ code: "ERR_OUT_OF_RANGE" });
|
||||
let rangeError: unknown;
|
||||
try {
|
||||
backend.execute({ type: "framesInRange", input: { startMs: 0, endMs: 120_000 } });
|
||||
} catch (error) {
|
||||
rangeError = error;
|
||||
}
|
||||
expect(rangeError).toMatchObject({
|
||||
code: "ERR_OUT_OF_RANGE",
|
||||
name: originalError.name,
|
||||
message: originalError.message,
|
||||
});
|
||||
let pendingError: unknown;
|
||||
try {
|
||||
backend.execute({ type: "unbatchedActiveFrames", input: { limit: 1 } });
|
||||
|
|
|
|||
|
|
@ -306,11 +306,23 @@ class LogbookDatabaseStore {
|
|||
return row ? toFrame(row) : null;
|
||||
}
|
||||
|
||||
framesInRange(startMs: number, endMs: number): LogbookFrame[] {
|
||||
framesInRange(
|
||||
startMs: number,
|
||||
endMs: number,
|
||||
): Pick<LogbookFrame, "id" | "capturedAtMs" | "idle">[] {
|
||||
return executeSqliteQuerySync(
|
||||
this.db,
|
||||
this.framesQuery.where("captured_at_ms", ">=", startMs).where("captured_at_ms", "<", endMs),
|
||||
).rows.map(toFrame);
|
||||
this.framesQuery
|
||||
.clearSelect()
|
||||
// Keep native integer decoding and overflow errors for unused numeric fields.
|
||||
.select(["id", "captured_at_ms", "screen_index", "width", "height", "byte_size", "idle"])
|
||||
.where("captured_at_ms", ">=", startMs)
|
||||
.where("captured_at_ms", "<", endMs),
|
||||
).rows.map((row) => ({
|
||||
id: row.id,
|
||||
capturedAtMs: row.captured_at_ms,
|
||||
idle: row.idle === 1,
|
||||
}));
|
||||
}
|
||||
|
||||
createBatch(params: LogbookBatchInput): number {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue