perf(tasks): avoid rereading flow history during active writes (#153806)

This commit is contained in:
Peter Steinberger 2026-09-20 10:02:13 -07:00 • committed by GitHub
parent 0aca75722d
commit a0f9b0b766
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
9 changed files with 151 additions and 19 deletions

View file

@ -711,14 +711,7 @@ describe("registered task flow reconciliation", () => {
if (intervening === "lookup") {
measureLookup("pending", "lookup-run", "lookup-0-16");
measureLookup("pending-no-eligible", "ineligible-run", "ineligible");
const { db } = openOpenClawStateDatabase();
executeSqliteQuerySync(
db,
getNodeSqliteKysely<DB>(db)
.updateTable("flow_runs")
.set({ sync_mode: "managed" })
.where("flow_id", "=", "lookup-0-16"),
);
expect(deleteTaskFlowRecordById("lookup-0-16")).toBe(true);
measureLookup("pending-fresh", "lookup-run", "lookup-0-15");
} else if (intervening === "update") {
expect(
@ -777,7 +770,7 @@ describe("registered task flow reconciliation", () => {
expect(legacy.get(created.flowId)).toEqual(settled);
console.log("Run lookup flow snapshots", JSON.stringify(lookupReads));
expect(lookupReads.map(({ count }) => count)).toEqual([0, 0, 1, 0, 1, 0, 1]);
expect(lookupReads.map(({ rows }) => rows)).toEqual([0, 0, 35, 0, 35, 0, 35]);
expect(lookupReads.map(({ rows }) => rows)).toEqual([0, 0, 1, 0, 1, 0, 34]);
expect(
lookupReads.filter(({ count }) => count > 0).every(({ textBytes }) => textBytes > 0),
).toBe(true);

View file

@ -0,0 +1,114 @@
import { expectDefined } from "@openclaw/normalization-core";
import { afterEach, expect, it, vi } from "vitest";
import { createDeferred } from "../../test/helpers/promise.js";
import { runOpenClawStateWriteTransaction } from "../state/openclaw-state-db.js";
import { captureOpenClawStateWorkerContext } from "../state/openclaw-state-worker-context.js";
import { withOpenClawTestState } from "../test-utils/openclaw-test-state.js";
import {
ensureTaskFlowRegistryReady,
getTaskFlowById,
getTaskMirroredFlowIds,
readResidentTaskFlow,
runTaskFlowRegistryWorkerMutation,
} from "./task-flow-registry.js";
import { getTaskFlowRegistryStore } from "./task-flow-registry.store.js";
import { configureTaskFlowRegistryRuntime } from "./task-flow-registry.store.test-support.js";
import { resetTaskFlowRegistryForTests } from "./task-flow-registry.test-support.js";
import type { TaskFlowRecord } from "./task-flow-registry.types.js";
afterEach(() => vi.restoreAllMocks());
it.each(["update", "delete", "rollback"] as const)(
"refreshes only dirty flows during an overlapping worker %s",
async (operation) => {
await withOpenClawTestState({ layout: "state-only" }, async () => {
resetTaskFlowRegistryForTests({ persist: false });
const store = getTaskFlowRegistryStore();
const records: TaskFlowRecord[] = Array.from({ length: 20 }, (_, index) => ({
flowId: `flow-${index}`,
syncMode: "managed",
controllerId: "tests/projection",
ownerKey: "agent:main:main",
revision: 0,
status: index === 0 ? "running" : "succeeded",
notifyPolicy: "silent",
goal: "Retained flow",
stateJson: { payload: "x".repeat(4096) },
createdAt: 1,
updatedAt: 1,
}));
const initial = expectDefined(records[0], "initial flow");
const next = { ...initial, revision: 1, goal: "Updated flow" };
const release = createDeferred();
let mutation: Promise<void> | undefined;
try {
runOpenClawStateWriteTransaction(() => records.forEach((flow) => store.upsertFlow(flow)));
ensureTaskFlowRegistryReady();
const context = captureOpenClawStateWorkerContext();
const snapshots: number[] = [];
const load = store.loadSnapshot.bind(store);
vi.spyOn(store, "loadSnapshot").mockImplementation((...args) => {
const snapshot = load(...args);
snapshots.push(snapshot.flows.size);
return snapshot;
});
const events = vi.fn();
configureTaskFlowRegistryRuntime({ observers: { onEvent: events } });
const unrelated = readResidentTaskFlow("flow-1");
for (let index = 0; index < 5; index += 1) {
getTaskMirroredFlowIds(["flow-1"]);
}
expect(snapshots).toEqual([]);
mutation = runTaskFlowRegistryWorkerMutation(
{ flowId: initial.flowId, admission: context.admission },
() => release.promise,
async () =>
operation === "delete" ? undefined : operation === "update" ? next : initial,
);
const refresh = () => {
if (operation === "delete") {
store.deleteFlow(initial.flowId);
} else {
store.upsertFlow(next);
}
for (let index = 0; index < 5; index += 1) {
getTaskMirroredFlowIds(["flow-1"]);
}
expect(readResidentTaskFlow(initial.flowId)).toEqual(
operation === "delete" ? undefined : next,
);
expect(snapshots).toEqual(Array(5).fill(operation === "delete" ? 0 : 1));
expect(readResidentTaskFlow("flow-1")).toBe(unrelated);
expect(events).not.toHaveBeenCalled();
};
if (operation === "rollback") {
expect(() =>
runOpenClawStateWriteTransaction(() => {
refresh();
throw new Error("abort refresh");
}),
).toThrow("abort refresh");
expect(readResidentTaskFlow(initial.flowId)).toEqual(initial);
} else {
refresh();
}
release.resolve();
await mutation;
expect(getTaskFlowById(initial.flowId)).toEqual(
operation === "delete" ? undefined : operation === "update" ? next : initial,
);
expect(events).toHaveBeenCalledTimes(operation === "rollback" ? 0 : 1);
snapshots.length = 0;
for (let index = 0; index < 5; index += 1) {
getTaskMirroredFlowIds(["flow-1"]);
}
expect(snapshots).toEqual([]);
} finally {
release.resolve();
await mutation;
resetTaskFlowRegistryForTests({ persist: false });
}
});
},
);

View file

@ -11,6 +11,7 @@ import {
executeSqliteQueryTakeFirstSync,
getNodeSqliteKysely,
prepareSqliteQuerySync,
sqliteStringSet,
} from "../infra/kysely-sync.js";
import { normalizeSqliteNumber } from "../infra/sqlite-number.js";
import type { DB as OpenClawStateKyselyDatabase } from "../state/openclaw-state-db.generated.js";
@ -202,13 +203,22 @@ export function listTaskFlowRecordsForOwnerReadInDatabase(
return read(ownerKey).rows.map(rowToFlowRecord);
}
export function readTaskFlowRegistrySnapshot(db: DatabaseSync): TaskFlowRegistryStoreSnapshot {
const query = getFlowRegistryKysely(db)
export function readTaskFlowRegistrySnapshot(
db: DatabaseSync,
flowIds?: readonly string[],
): TaskFlowRegistryStoreSnapshot {
let query = getFlowRegistryKysely(db)
.selectFrom("flow_runs")
.select(FLOW_RUN_SELECT_COLUMNS)
.orderBy("created_at", "asc")
.orderBy("flow_id", "asc");
const flows = new Map<string, TaskFlowRecord>();
if (flowIds) {
if (flowIds.length === 0) {
return { flows };
}
query = query.where("flow_id", "in", sqliteStringSet(flowIds));
}
// Finish native reads before decoding so SQLite errors retain precedence.
for (const row of executeSqliteQuerySync(db, query).rows) {
flows.set(row.flow_id, rowToFlowRecord(row));

View file

@ -69,8 +69,10 @@ function withWriteTransaction(write: (database: FlowRegistryDatabase) => void) {
});
}
export function loadTaskFlowRegistryStateFromSqlite(): TaskFlowRegistryStoreSnapshot {
return readTaskFlowRegistrySnapshot(openFlowRegistryDatabase().db);
export function loadTaskFlowRegistryStateFromSqlite(
flowIds?: readonly string[],
): TaskFlowRegistryStoreSnapshot {
return readTaskFlowRegistrySnapshot(openFlowRegistryDatabase().db, flowIds);
}
/** Loads task flows without creating or migrating shared state. */

View file

@ -35,7 +35,7 @@ type TaskFlowRegistryStore = {
context: OpenClawStateWorkerContext,
flowId: string,
): Promise<TaskFlowRecord | undefined>;
loadSnapshot: () => TaskFlowRegistryStoreSnapshot;
loadSnapshot: (flowIds?: readonly string[]) => TaskFlowRegistryStoreSnapshot;
upsertFlow: (flow: TaskFlowRecord) => void;
syncMirroredTask: (
task: TaskFlowSyncInput,

View file

@ -29,7 +29,7 @@ export type TaskFlowRegistryUpdatePublication = {
publish: () => void;
};
/** Full task-flow registry snapshot used for persistence restore and replacement writes. */
/** Task-flow rows for a full restore or an explicitly scoped projection refresh. */
export type TaskFlowRegistryStoreSnapshot = {
flows: Map<string, TaskFlowRecord>;
};

View file

@ -233,10 +233,11 @@ export function ensureTaskFlowRegistryReady(options?: { refreshProjection?: bool
if (options?.refreshProjection === false || (!projectionDirty && dirtyFlowIds.size === 0)) {
return;
}
const restored = getTaskFlowRegistryStore().loadSnapshot();
const flowIds = projectionDirty ? undefined : [...dirtyFlowIds];
const restored = getTaskFlowRegistryStore().loadSnapshot(flowIds);
const previous = flows;
const next = new Map(previous);
for (const flowId of next.keys()) {
for (const flowId of flowIds ?? next.keys()) {
if (!restored.flows.has(flowId)) {
next.delete(flowId);
}

View file

@ -292,7 +292,7 @@ describe("asynchronous registry restoration", () => {
"refreshes a flow write pending %s after synchronous snapshot installation",
async (when) => {
const store = createInMemoryTaskFlowRegistryStore({ flows: new Map([[flow.flowId, flow]]) });
const loadSnapshot = vi.fn(() => store.loadSnapshot());
const loadSnapshot = vi.fn(store.loadSnapshot);
const release = createDeferred();
const context = captureOpenClawStateWorkerContext();
let pending: Promise<void> | undefined;
@ -325,6 +325,8 @@ describe("asynchronous registry restoration", () => {
currentStep: "pending mutation",
});
expect(loadSnapshot).toHaveBeenCalledTimes(2);
expect(loadSnapshot).toHaveBeenNthCalledWith(1);
expect(loadSnapshot).toHaveBeenNthCalledWith(2, [flow.flowId]);
} finally {
release.resolve();
await pending;

View file

@ -17,6 +17,7 @@ import type {
TaskFlowRegistryObservedUpdate,
TaskFlowRegistryStoreSnapshot,
} from "../tasks/task-flow-registry.store.types.js";
import type { TaskFlowRecord } from "../tasks/task-flow-registry.types.js";
import type { TaskInitialWorkerOperations } from "../tasks/task-initial-worker.types.js";
import { captureTaskCreationEventTarget } from "../tasks/task-registry-agent-event-target.js";
import {
@ -369,7 +370,16 @@ export function createInMemoryTaskFlowRegistryStore(
return {
withSnapshotAsync: async (_context, consume) => consume(structuredClone(state)),
readFlowAsync: async (_context, flowId) => structuredClone(state.flows.get(flowId)),
loadSnapshot: () => structuredClone(state),
loadSnapshot: (flowIds) => {
const flows = new Map<string, TaskFlowRecord>();
for (const flowId of flowIds ?? state.flows.keys()) {
const flow = state.flows.get(flowId);
if (flow) {
flows.set(flowId, structuredClone(flow));
}
}
return { flows };
},
upsertFlow: (flow) => {
state.flows.set(flow.flowId, structuredClone(flow));
},