fix(sessions): back off automatic maintenance replans (#156583)

Coalesce invalidated automatic maintenance behind the per-store write-quiet
timer and pause after three consecutive rejections. Preserve pending keys
and commit authority. In warn mode, capture the age fact without creating
a reclamation operation.

Refs #153257. Thanks to @abuegab1-spec for tracing the generation-change
loop and suggesting the explicit warn-mode guard.

Validated with failing baseline regressions, 86 owner/sibling tests,
isolated Gateway write-load and persisted-state proof, and Codex review.
This commit is contained in:
Peter Steinberger 2026-09-23 08:07:29 -07:00 • committed by GitHub
parent 60d25b947d
commit 20445202d3
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 305 additions and 162 deletions

View file

@ -279,6 +279,11 @@ Ordinary entry writes also arm background maintenance at the next age boundary,
with a periodic recheck every 30 minutes while the store remains open. This lets
eligible sessions age out without further traffic. Writes that cannot change
age or count maintenance outcomes skip candidate scans.
If writes invalidate an automatic maintenance plan, its replacement waits for
a quiet window after the last write (one second, then two seconds). Three
consecutive invalidations pause automatic retries and log the cause; a new
entry write can schedule another attempt. `warn` mode captures the maintenance
age fact without constructing or dispatching automatic reclamation.
`maxEntries` defaults to 5000 unarchived session rows. Archived rows do not consume
the cap. Existing explicit limits remain unchanged.

View file

@ -4,6 +4,7 @@ import { DatabaseSync } from "node:sqlite";
import { setImmediate as yieldToEventLoop } from "node:timers/promises";
import { afterEach, expect, it, vi } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js";
import * as logging from "../../logging/logger.js";
import * as agentDatabase from "../../state/openclaw-agent-db.js";
import {
closeOpenClawAgentDatabaseByPath,
@ -64,6 +65,86 @@ function createStore(pruneAfterMs = 1_000, key = sessionKey) {
return { database, request, scope, storePath, updatedAt };
}
it("captures warn-mode age facts without constructing or dispatching reclamation", async () => {
const { request, storePath } = createStore();
request.maintenanceConfig.mode = "warn";
const capture = vi.spyOn(ageFacts, "captureSessionEntryMaintenanceAgeFact");
const plans = vi.spyOn(reclamation, "createSessionMaintenancePlanningOperation");
kickSessionEntryMaintenanceAfterWrite(request);
await yieldToEventLoop();
expect(capture).toHaveBeenCalledTimes(1);
expect(plans).not.toHaveBeenCalled();
expect(reclamation.runSqliteSessionReclamation).not.toHaveBeenCalled();
await vi.advanceTimersByTimeAsync(30 * 60 * 1_000);
expect(loadSessionEntry({ sessionKey, storePath })?.archivedAt).toBeUndefined();
expect(capture).toHaveBeenCalledTimes(1);
});
it("coalesces rejected plans behind write quiet, caps retries, and resumes on a new write", async () => {
const { request, scope, storePath, updatedAt } = createStore();
const victimKey = "agent:main:replan-victim";
runOpenClawAgentWriteTransaction((owner) => {
writeSessionEntry(owner, victimKey, { sessionId: "victim", updatedAt: updatedAt - 2_000 });
}, scope);
const logger = logging.getChildLogger({ subsystem: "session-sqlite" });
vi.spyOn(logging, "getChildLogger").mockReturnValue(logger);
const warn = vi.spyOn(logger, "warn");
const createPlan = reclamation.createSessionMaintenancePlanningOperation;
let rejections = 0;
const plans = vi
.spyOn(reclamation, "createSessionMaintenancePlanningOperation")
.mockImplementation((params) => {
const operation = createPlan(params);
if (rejections < 3) {
queueMicrotask(() => {
rejections += 1;
runOpenClawAgentWriteTransaction((owner) => {
writeSessionEntry(owner, sessionKey, {
sessionId: "age-kick",
updatedAt: Date.now(),
label: `write-${rejections}`,
});
}, scope);
kickSessionEntryMaintenanceAfterWrite(request);
});
}
return operation;
});
kickSessionEntryMaintenanceAfterWrite(request);
await yieldToEventLoop();
expect(plans).toHaveBeenCalledTimes(1);
expect(loadSessionEntry({ sessionKey: victimKey, storePath })?.archivedAt).toBeUndefined();
for (let write = 0; write < 5; write += 1) {
await vi.advanceTimersByTimeAsync(999);
kickSessionEntryMaintenanceAfterWrite(request);
}
expect(plans).toHaveBeenCalledTimes(1);
await vi.advanceTimersByTimeAsync(1_000);
await yieldToEventLoop();
expect(plans).toHaveBeenCalledTimes(2);
await vi.advanceTimersByTimeAsync(2_000);
await yieldToEventLoop();
expect(plans).toHaveBeenCalledTimes(3);
expect(warn).toHaveBeenCalledWith(
"SQLite automatic session maintenance paused after repeated input changes",
expect.objectContaining({ rejections: 3, error: expect.any(Error) }),
);
await vi.advanceTimersByTimeAsync(30 * 60 * 1_000);
expect(plans).toHaveBeenCalledTimes(3);
expect(loadSessionEntry({ sessionKey: victimKey, storePath })?.archivedAt).toBeUndefined();
kickSessionEntryMaintenanceAfterWrite(request);
await yieldToEventLoop();
expect(plans).toHaveBeenCalledTimes(3);
await vi.advanceTimersByTimeAsync(1_000);
await yieldToEventLoop();
expect(plans).toHaveBeenCalledTimes(4);
expect(loadSessionEntry({ sessionKey: victimKey, storePath })).toMatchObject({
archiveReason: "age-retention",
});
});
it("archives an entry at its age boundary without another write", async () => {
const { request, storePath } = createStore();
kickSessionEntryMaintenanceAfterWrite(request);
@ -332,6 +413,8 @@ it.each(
kickSessionEntryMaintenanceAfterWrite(request);
await yieldToEventLoop();
expect(changed).toBe(true);
await vi.advanceTimersByTimeAsync(1_000);
await yieldToEventLoop();
expect(loadSessionEntry({ sessionKey: victimKey, storePath })).toMatchObject({
archivedAt: expect.any(Number),
archiveReason: mutation === "insert" ? "active-session-cap" : "age-retention",
@ -430,6 +513,10 @@ it.runIf(process.platform !== "win32").each(["before preparation", "after prepar
try {
kickSessionEntryMaintenanceAfterWrite(request);
await yieldToEventLoop();
if (when === "after preparation") {
await vi.advanceTimersByTimeAsync(1_000);
await yieldToEventLoop();
}
expect(dispatch).toHaveBeenCalledTimes(1);
expect(replaced).toBe(false);
expect(loadSessionEntry({ sessionKey, storePath })?.archivedAt).toBeUndefined();

View file

@ -48,12 +48,16 @@ type SessionEntryMaintenanceOwner = SessionEntryMaintenanceRequest & {
database: OpenClawAgentDatabase;
generation: number;
running: boolean;
rejections: number;
retryDelayMs?: number;
immediate?: ReturnType<typeof setImmediate>;
timer?: ReturnType<typeof setTimeout>;
unregisterClose?: () => void;
};
const maintenanceByStore = new Map<string, SessionEntryMaintenanceOwner>();
const MAINTENANCE_WRITE_QUIET_MS = 1_000;
const MAX_MAINTENANCE_REJECTIONS = 3;
/** Coalesce automatic logical maintenance outside ordinary entry-write latency. */
export function kickSessionEntryMaintenanceAfterWrite(
@ -72,7 +76,14 @@ export function kickSessionEntryMaintenanceAfterWrite(
owner.activeSessionKeys.add(params.activeSessionKey);
Object.assign(owner, params, { generation: owner.generation + 1 });
if (!owner.running) {
scheduleImmediateMaintenance(databasePath, owner);
if (owner.retryDelayMs !== undefined) {
if (owner.rejections >= MAX_MAINTENANCE_REJECTIONS) {
owner.rejections = 0;
}
scheduleMaintenanceAfterWriteQuiet(databasePath, owner);
} else {
scheduleImmediateMaintenance(databasePath, owner);
}
}
return;
}
@ -85,6 +96,7 @@ export function kickSessionEntryMaintenanceAfterWrite(
database,
generation: 1,
running: false,
rejections: 0,
};
maintenanceByStore.set(databasePath, created);
created.unregisterClose = registerNodeSqliteDisposeCallback(database.db, () =>
@ -115,6 +127,26 @@ function scheduleImmediateMaintenance(
});
}
function scheduleMaintenanceAfterWriteQuiet(
databasePath: string,
owner: SessionEntryMaintenanceOwner,
): void {
owner.running = false;
owner.retryDelayMs = MAINTENANCE_WRITE_QUIET_MS * 2 ** Math.max(0, owner.rejections - 1);
if (owner.timer) {
// A write restarts this one-shot quiet window; no polling or competing retry owner.
owner.timer.refresh();
return;
}
owner.timer = setTimeout(() => {
owner.timer = undefined;
owner.retryDelayMs = undefined;
owner.running = true;
void runPendingMaintenance(databasePath, owner);
}, owner.retryDelayMs);
owner.timer.unref();
}
async function runPendingMaintenance(
databasePath: string,
owner: SessionEntryMaintenanceOwner,
@ -123,176 +155,195 @@ async function runPendingMaintenance(
maintenanceByStore.get(databasePath) === owner &&
owner.database.db.isOpen &&
getOpenClawAgentDatabaseIfOpen(toDatabaseOptions(owner.scope)) === owner.database;
while (isCurrent()) {
const generation = owner.generation;
const activeSessionKeys = [...owner.activeSessionKeys];
owner.activeSessionKeys.clear();
let nextMaintenanceAt: number | undefined = Infinity;
let planningChanged = false;
try {
const prepared = await runExclusiveSqliteSessionWrite(
if (!isCurrent()) {
retireMaintenanceOwner(databasePath, owner);
return;
}
const generation = owner.generation;
const activeSessionKeys = [...owner.activeSessionKeys];
owner.activeSessionKeys.clear();
let nextMaintenanceAt: number | undefined = Infinity;
let planningChanged = false;
try {
const prepared = await runExclusiveSqliteSessionWrite(
owner.scope,
async () => {
// The writer queue can outlive the handle that admitted this owner.
// Check inside the acquired lane so an evicted owner cannot reopen the path.
if (!isCurrent()) {
return undefined;
}
const maintenance = owner.maintenanceConfig
? normalizeResolvedMaintenanceConfigInput(owner.maintenanceConfig)
: resolveMaintenanceConfig();
const ageCapture = captureSessionEntryMaintenanceAgeFact(owner.database.db, maintenance);
if (maintenance.mode === "warn") {
return { maintenance, ageCapture, operation: undefined };
}
// Cold planning stays off-thread. Only an already-owned, current fact can
// justify the compact count read before dispatching a no-op Worker request.
if (
ageCapture.fact &&
isOpenClawAgentDatabasePathCurrent(owner.database) &&
canSkipSessionEntryMaintenanceInDatabase(owner.database, { maintenance })
) {
return { maintenance, ageCapture, operation: undefined };
}
const operation = createSessionMaintenancePlanningOperation({
databaseOptions: toDatabaseOptions(owner.scope),
input: {
ageFact: ageCapture.fact,
activeSessionKeys,
archiveDirectory: owner.archiveDirectory,
maintenance,
preservation: null,
storePath: owner.storePath,
},
});
return { maintenance, operation, ageCapture };
},
"session.maintenance.plan",
);
if (!prepared) {
retireMaintenanceOwner(databasePath, owner);
return;
}
const { maintenance, operation } = prepared;
if (maintenance.mode === "warn") {
if (isCurrent() && owner.generation !== generation) {
scheduleMaintenanceAfterWriteQuiet(databasePath, owner);
} else {
retireMaintenanceOwner(databasePath, owner);
}
return;
}
let { ageCapture } = prepared;
const assertInputsCurrent = () => {
if (!isCurrent()) {
throw new Error("SQLite automatic maintenance owner retired");
}
if (
owner.generation !== generation ||
(operation &&
operation.input.preservation !== null &&
!isDeepStrictEqual(
operation.input.preservation,
captureSessionMaintenancePreservation(operation.input.storePath),
))
) {
planningChanged = true;
throw new Error("SQLite automatic maintenance inputs changed before commit");
}
};
const assertCurrent = () => {
assertInputsCurrent();
if (!isSessionEntryMaintenanceAgeCaptureCurrent(owner.database.db, ageCapture)) {
planningChanged = true;
throw new Error("SQLite automatic maintenance age fact changed before commit");
}
if (!operation && !isOpenClawAgentDatabasePathCurrent(owner.database)) {
planningChanged = true;
throw new Error("SQLite automatic maintenance database path changed");
}
};
const runPlanning = () => {
if (!operation) {
assertCurrent();
return Promise.resolve({
kind: "maintenance-plan" as const,
value: emptySessionEntryMaintenancePlan(),
});
}
return runSqliteSessionReclamation({
diagnostics: { kind: "maintenance-plan" },
assertCommitAllowed: assertCurrent,
onWorkerResult: (result) => {
if (result.kind === "maintenance-plan" && isCurrent()) {
adoptSessionEntryMaintenanceAgeFact(owner.database.db, ageCapture, result.ageFact);
}
},
forceInProcess: false,
plan: operation,
});
};
let result = await runPlanning();
if (!operation) {
assertCurrent();
}
if (operation && result.kind === "maintenance-preservation-required") {
await runExclusiveSqliteSessionWrite(
owner.scope,
async () => {
// The writer queue can outlive the handle that admitted this owner.
// Check inside the acquired lane so an evicted owner cannot reopen the path.
if (!isCurrent()) {
return undefined;
}
const maintenance = owner.maintenanceConfig
? normalizeResolvedMaintenanceConfigInput(owner.maintenanceConfig)
: resolveMaintenanceConfig();
const ageCapture = captureSessionEntryMaintenanceAgeFact(owner.database.db, maintenance);
// Cold planning stays off-thread. Only an already-owned, current fact can
// justify the compact count read before dispatching a no-op Worker request.
if (
ageCapture.fact &&
maintenance.mode === "enforce" &&
isOpenClawAgentDatabasePathCurrent(owner.database) &&
canSkipSessionEntryMaintenanceInDatabase(owner.database, { maintenance })
) {
return { maintenance, ageCapture, operation: undefined };
}
const operation = createSessionMaintenancePlanningOperation({
databaseOptions: toDatabaseOptions(owner.scope),
input: {
ageFact: ageCapture.fact,
activeSessionKeys,
archiveDirectory: owner.archiveDirectory,
maintenance,
preservation: null,
storePath: owner.storePath,
},
});
return { maintenance, operation, ageCapture };
assertInputsCurrent();
// The in-process transaction also invalidates facts on rollback. Replan
// from current owner state only after the explicit preservation rollback.
ageCapture = captureSessionEntryMaintenanceAgeFact(
owner.database.db,
operation.input.maintenance,
);
operation.input.ageFact = ageCapture.fact;
operation.input.preservation = captureSessionMaintenancePreservation(
operation.input.storePath,
);
},
"session.maintenance.plan",
);
if (!prepared) {
break;
}
const { maintenance, operation } = prepared;
let { ageCapture } = prepared;
const assertInputsCurrent = () => {
if (!isCurrent()) {
throw new Error("SQLite automatic maintenance owner retired");
}
if (
owner.generation !== generation ||
(operation &&
operation.input.preservation !== null &&
!isDeepStrictEqual(
operation.input.preservation,
captureSessionMaintenancePreservation(operation.input.storePath),
))
) {
planningChanged = true;
throw new Error("SQLite automatic maintenance inputs changed before commit");
}
};
const assertCurrent = () => {
assertInputsCurrent();
if (!isSessionEntryMaintenanceAgeCaptureCurrent(owner.database.db, ageCapture)) {
planningChanged = true;
throw new Error("SQLite automatic maintenance age fact changed before commit");
}
if (!operation && !isOpenClawAgentDatabasePathCurrent(owner.database)) {
planningChanged = true;
throw new Error("SQLite automatic maintenance database path changed");
}
};
const runPlanning = () => {
if (!operation) {
assertCurrent();
return Promise.resolve({
kind: "maintenance-plan" as const,
value: emptySessionEntryMaintenancePlan(),
});
}
return runSqliteSessionReclamation({
diagnostics: { kind: "maintenance-plan" },
assertCommitAllowed: assertCurrent,
onWorkerResult: (result) => {
if (result.kind === "maintenance-plan" && isCurrent()) {
adoptSessionEntryMaintenanceAgeFact(owner.database.db, ageCapture, result.ageFact);
}
},
forceInProcess: false,
plan: operation,
});
};
let result =
maintenance.mode === "warn"
? { kind: "maintenance-plan" as const, value: emptySessionEntryMaintenancePlan() }
: await runPlanning();
if (!operation) {
assertCurrent();
}
if (operation && result.kind === "maintenance-preservation-required") {
await runExclusiveSqliteSessionWrite(
owner.scope,
async () => {
assertInputsCurrent();
// The in-process transaction also invalidates facts on rollback. Replan
// from current owner state only after the explicit preservation rollback.
ageCapture = captureSessionEntryMaintenanceAgeFact(
owner.database.db,
operation.input.maintenance,
);
operation.input.ageFact = ageCapture.fact;
operation.input.preservation = captureSessionMaintenancePreservation(
operation.input.storePath,
);
},
"session.maintenance.plan",
);
result = await runPlanning();
}
if (result.kind !== "maintenance-plan") {
throw new Error("SQLite automatic maintenance returned another operation's result");
}
const plan = result.value;
await finalizeSessionEntryMaintenancePlansAfterWriterReleaseBestEffort(owner.scope, [plan], {
isCurrent,
});
if (isCurrent() && owner.generation === generation) {
nextMaintenanceAt = readSessionEntryMaintenanceNextAgeAt(owner.database, maintenance);
}
} catch (error) {
if (planningChanged && isCurrent()) {
owner.generation += 1;
activeSessionKeys.forEach((key) => owner.activeSessionKeys.add(key));
} else {
getChildLogger({ subsystem: "session-sqlite" }).warn(
"SQLite automatic session maintenance failed",
{ error, path: databasePath },
);
}
result = await runPlanning();
}
// Any write during awaited planning/finalization increments the generation.
// Keep this owner alive so that write gets a fresh maintenance snapshot.
if (!isCurrent()) {
break;
if (result.kind !== "maintenance-plan") {
throw new Error("SQLite automatic maintenance returned another operation's result");
}
if (owner.generation === generation) {
if (nextMaintenanceAt === undefined) {
break;
}
const plan = result.value;
await finalizeSessionEntryMaintenancePlansAfterWriterReleaseBestEffort(owner.scope, [plan], {
isCurrent,
});
owner.rejections = 0;
if (isCurrent() && owner.generation === generation) {
nextMaintenanceAt = readSessionEntryMaintenanceNextAgeAt(owner.database, maintenance);
}
} catch (error) {
if (planningChanged && isCurrent()) {
activeSessionKeys.forEach((key) => owner.activeSessionKeys.add(key));
owner.rejections += 1;
owner.running = false;
owner.timer = setTimeout(
() => {
owner.timer = undefined;
owner.running = true;
void runPendingMaintenance(databasePath, owner);
},
// Bound relative delays too: Node clamps overflowed timeouts to 1 ms.
Math.max(
1,
Math.min(SESSION_ENTRY_MAINTENANCE_INTERVAL_MS, nextMaintenanceAt - Date.now()),
),
);
owner.timer.unref();
owner.retryDelayMs = MAINTENANCE_WRITE_QUIET_MS;
if (owner.rejections >= MAX_MAINTENANCE_REJECTIONS) {
getChildLogger({ subsystem: "session-sqlite" }).warn(
"SQLite automatic session maintenance paused after repeated input changes",
{ error, path: databasePath, rejections: owner.rejections },
);
} else {
scheduleMaintenanceAfterWriteQuiet(databasePath, owner);
}
return;
}
getChildLogger({ subsystem: "session-sqlite" }).warn(
"SQLite automatic session maintenance failed",
{ error, path: databasePath },
);
}
retireMaintenanceOwner(databasePath, owner);
// Writes during finalization also coalesce behind the next quiet window.
if (!isCurrent()) {
retireMaintenanceOwner(databasePath, owner);
return;
}
if (owner.generation === generation) {
if (nextMaintenanceAt === undefined) {
retireMaintenanceOwner(databasePath, owner);
return;
}
owner.running = false;
owner.timer = setTimeout(
() => {
owner.timer = undefined;
owner.running = true;
void runPendingMaintenance(databasePath, owner);
},
// Bound relative delays too: Node clamps overflowed timeouts to 1 ms.
Math.max(1, Math.min(SESSION_ENTRY_MAINTENANCE_INTERVAL_MS, nextMaintenanceAt - Date.now())),
);
owner.timer.unref();
return;
}
scheduleMaintenanceAfterWriteQuiet(databasePath, owner);
}