From 2f946506bb89c156269232252ec614aad30e8a51 Mon Sep 17 00:00:00 2001 From: Vincent Koc <25068+vincentkoc@users.noreply.github.com> Date: Sun, 12 Jul 2026 04:33:15 +0200 Subject: [PATCH] fix(sessions): restore SQLite terminal freshness --- CHANGELOG.md | 1 + ...-session-transcript-schema-baseline.sha256 | 2 +- scripts/generate-kysely-types.mjs | 16 +- src/commands/doctor-session-sqlite.test.ts | 24 +++ src/commands/doctor-session-sqlite.ts | 15 ++ src/config/sessions/lifecycle.test.ts | 166 ++++++++++++++++ src/config/sessions/lifecycle.ts | 52 ++--- .../session-accessor.sqlite-contract.ts | 2 + .../sessions/session-accessor.sqlite.ts | 177 +++++++++++++++--- src/config/sessions/session-accessor.test.ts | 50 +++++ src/config/sessions/session-accessor.ts | 2 + src/state/openclaw-agent-db.generated.d.ts | 2 + src/state/openclaw-agent-db.test.ts | 78 +++++++- src/state/openclaw-agent-db.ts | 41 +++- src/state/openclaw-agent-schema.generated.ts | 2 + src/state/openclaw-agent-schema.sql | 2 + 16 files changed, 559 insertions(+), 73 deletions(-) create mode 100644 src/config/sessions/lifecycle.test.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index 399fdee7235b..be5dbc0c38c5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -30,6 +30,7 @@ Docs: https://docs.openclaw.ai ### Fixes +- **SQLite terminal session recovery:** track physical transcript mutation time in the agent database so killed or timed-out main sessions rotate when transcript writes outlive the registry update, while preserving legacy transcript mtimes during doctor import. - **Gateway chat typecheck:** import chat event types from their owning protocol schema after the retired aggregate type module was removed, restoring full project typechecks. - **Packaged Crabbox commands:** include the lease-freshness helper imported by the published wrapper so `crabbox:*` commands do not fail with `ERR_MODULE_NOT_FOUND` in npm installs. - **Plugin session catalogs:** reject unknown catalog filters, report catalogs as plugin capabilities, and preserve them in SDK registration captures instead of silently returning empty results or classifying catalog-only plugins as capability-free. diff --git a/docs/.generated/sqlite-session-transcript-schema-baseline.sha256 b/docs/.generated/sqlite-session-transcript-schema-baseline.sha256 index 5e723d90f2ea..5d339e202540 100644 --- a/docs/.generated/sqlite-session-transcript-schema-baseline.sha256 +++ b/docs/.generated/sqlite-session-transcript-schema-baseline.sha256 @@ -1 +1 @@ -4bb5f50a60c3664656d56f2be8d9884d48287f084c03959120b2825991b23f21 sqlite-session-transcript-schema-baseline.sql +5afc3eb779d73bf3a38ab2924d3d3a2095ce6e928a791d47fef7c27747cfb0f7 sqlite-session-transcript-schema-baseline.sql diff --git a/scripts/generate-kysely-types.mjs b/scripts/generate-kysely-types.mjs index 1fe16dff6e1d..ceaefbcd5ef4 100644 --- a/scripts/generate-kysely-types.mjs +++ b/scripts/generate-kysely-types.mjs @@ -51,14 +51,14 @@ function columnBaseType(columnTypeLocal) { function columnType(column, primaryKeyColumnCount) { const baseType = columnBaseType(String(column.type ?? "")); - const generated = - column.dflt_value != null || - (primaryKeyColumnCount === 1 && - Number(column.pk) > 0 && - String(column.type ?? "") - .toUpperCase() - .includes("INT")); - const nullable = Number(column.notnull) !== 1 && !generated; + const generatedPrimaryKey = + primaryKeyColumnCount === 1 && + Number(column.pk) > 0 && + String(column.type ?? "") + .toUpperCase() + .includes("INT"); + const generated = column.dflt_value != null || generatedPrimaryKey; + const nullable = Number(column.notnull) !== 1 && !generatedPrimaryKey; const valueType = nullable ? `${baseType} | null` : baseType; return generated ? `Generated<${valueType}>` : valueType; } diff --git a/src/commands/doctor-session-sqlite.test.ts b/src/commands/doctor-session-sqlite.test.ts index ae9892b6b95e..3e611ab5025f 100644 --- a/src/commands/doctor-session-sqlite.test.ts +++ b/src/commands/doctor-session-sqlite.test.ts @@ -6,6 +6,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { loadExactSqliteSessionEntry, loadSqliteTranscriptEventsSync, + readSqliteTranscriptStatsSync, upsertSqliteSessionEntry, } from "../config/sessions/session-accessor.sqlite.js"; import { requireNodeSqlite } from "../infra/node-sqlite.js"; @@ -156,6 +157,29 @@ describe("runDoctorSessionSqlite", () => { expect(message?.message?.content).toEqual([{ type: "text", text: "legacy string" }]); }); + it("preserves the legacy transcript mtime as the SQLite mutation watermark", async () => { + const store = createLegacyStore(); + const transcriptMtimeMs = 1_700_000_000_000; + const transcriptMtime = new Date(transcriptMtimeMs); + fs.utimesSync(store.transcriptPath, transcriptMtime, transcriptMtime); + + const report = await runDoctorSessionSqlite({ + env: store.env, + mode: "import", + store: store.storePath, + }); + + expect(report.totals).toMatchObject({ importedEntries: 1, issues: 0 }); + expect( + readSqliteTranscriptStatsSync({ + agentId: "main", + sessionId: "session-1", + sessionKey: "agent:main:main", + storePath: store.storePath, + }).lastMutationAtMs, + ).toBe(transcriptMtimeMs); + }); + it("preserves a same-generation canonical harness owner during legacy import", async () => { const store = createLegacyStore({ entryOverrides: { lifecycleRevision: "rev-1" }, diff --git a/src/commands/doctor-session-sqlite.ts b/src/commands/doctor-session-sqlite.ts index d6d26292dfc9..dc29e558a75a 100644 --- a/src/commands/doctor-session-sqlite.ts +++ b/src/commands/doctor-session-sqlite.ts @@ -446,6 +446,7 @@ async function importLegacySessionRecord( report: DoctorSessionSqliteTargetReport, ): Promise { const result = countTranscriptEvents(record); + const transcriptMtimeMs = readLegacyTranscriptMtimeMs(record); if (result.status === "missing") { if (markAlreadyMigratedTranscript(target, record, report)) { return; @@ -473,6 +474,7 @@ async function importLegacySessionRecord( ...(record.transcriptPath ? { readTranscriptEvents: createTranscriptEventPrefixReader(record.transcriptPath) } : {}), + ...(transcriptMtimeMs !== undefined ? { transcriptMtimeMs } : {}), }); report.importedEntries += 1; report.importedTranscriptEvents += imported.transcriptEvents; @@ -491,6 +493,7 @@ async function importLegacySessionRecord( ...(record.transcriptPath && result.status === "ok" ? { readTranscriptEvents: createTranscriptEventReader(record.transcriptPath) } : {}), + ...(transcriptMtimeMs !== undefined ? { transcriptMtimeMs } : {}), }); report.importedEntries += 1; report.importedTranscriptEvents += imported.transcriptEvents; @@ -878,6 +881,18 @@ function countTranscriptEvents( return countTranscriptEventsForPath(record.transcriptPath); } +function readLegacyTranscriptMtimeMs(record: LegacySessionRecord): number | undefined { + if (!record.transcriptPath) { + return undefined; + } + try { + const mtimeMs = Math.floor(fs.statSync(record.transcriptPath).mtimeMs); + return Number.isFinite(mtimeMs) && mtimeMs >= 0 ? mtimeMs : undefined; + } catch { + return undefined; + } +} + function listUnreferencedJsonlFiles( storePath: string, referencedPaths: readonly string[], diff --git a/src/config/sessions/lifecycle.test.ts b/src/config/sessions/lifecycle.test.ts new file mode 100644 index 000000000000..c4ad9e8bffb1 --- /dev/null +++ b/src/config/sessions/lifecycle.test.ts @@ -0,0 +1,166 @@ +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { + closeOpenClawAgentDatabasesForTest, + resolveOpenClawAgentSqlitePath, +} from "../../state/openclaw-agent-db.js"; +import { closeOpenClawStateDatabaseForTest } from "../../state/openclaw-state-db.js"; +import { + hasTerminalMainSessionTranscriptNewerThanRegistry, + hasTerminalMainSessionTranscriptNewerThanRegistrySync, +} from "./lifecycle.js"; +import { appendTranscriptEvent, loadSessionEntry, upsertSessionEntry } from "./session-accessor.js"; +import type { SessionEntry } from "./types.js"; + +describe("terminal main session transcript freshness", () => { + let stateDir: string; + let storePath: string; + + beforeEach(() => { + stateDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-session-lifecycle-")); + storePath = path.join(stateDir, "agents", "main", "sessions", "sessions.json"); + }); + + afterEach(() => { + closeOpenClawAgentDatabasesForTest(); + closeOpenClawStateDatabaseForTest(); + fs.rmSync(stateDir, { recursive: true, force: true }); + }); + + async function createEntry(params: { + sessionFile?: string; + sessionKey?: string; + status: SessionEntry["status"]; + updatedAt: number; + }): Promise<{ entry: SessionEntry; sessionKey: string }> { + const sessionKey = params.sessionKey ?? "agent:main:main"; + const sessionId = `session-${params.status ?? "ended"}-${sessionKey.replaceAll(":", "-")}`; + const sessionEntry = { + sessionFile: + params.sessionFile ?? + `sqlite:main:${sessionId}:${resolveOpenClawAgentSqlitePath({ + agentId: "main", + env: { OPENCLAW_STATE_DIR: stateDir }, + })}`, + sessionId, + status: params.status, + updatedAt: params.updatedAt, + }; + await upsertSessionEntry({ agentId: "main", sessionKey, storePath }, sessionEntry); + await appendTranscriptEvent( + { agentId: "main", sessionId, sessionKey, storePath }, + { + type: "custom", + timestamp: "1970-01-01T00:00:00.001Z", + }, + ); + const storedEntry = loadSessionEntry({ agentId: "main", sessionKey, storePath }); + if (!storedEntry) { + throw new Error("expected session entry"); + } + return { + entry: { + ...storedEntry, + sessionFile: sessionEntry.sessionFile, + status: params.status, + updatedAt: params.updatedAt, + }, + sessionKey, + }; + } + + function check(entry: SessionEntry, sessionKey: string): boolean { + return hasTerminalMainSessionTranscriptNewerThanRegistrySync({ + agentId: "main", + entry, + sessionKey, + storePath, + }); + } + + it("uses the physical SQLite mutation watermark instead of event timestamps", async () => { + const registryTimestampMs = Date.now() - 10_000; + const { entry, sessionKey } = await createEntry({ + status: "killed", + updatedAt: registryTimestampMs, + }); + + expect(entry.updatedAt).toBe(registryTimestampMs); + expect(check(entry, sessionKey)).toBe(true); + await expect( + hasTerminalMainSessionTranscriptNewerThanRegistry({ + agentId: "main", + entry, + sessionKey, + storePath, + }), + ).resolves.toBe(true); + }); + + it.each(["done", "failed"] as const)("keeps %s terminal sessions reusable", async (status) => { + const { entry, sessionKey } = await createEntry({ + status, + updatedAt: Date.now() - 10_000, + }); + + expect(check(entry, sessionKey)).toBe(false); + }); + + it("uses SQLite freshness for entries that still contain legacy transcript paths", async () => { + const { entry, sessionKey } = await createEntry({ + sessionFile: path.join(stateDir, "legacy-session.jsonl"), + status: "killed", + updatedAt: Date.now() - 10_000, + }); + + expect(check(entry, sessionKey)).toBe(true); + }); + + it("does not rotate after a same-millisecond registry write observes the mutation", async () => { + const now = 1_700_000_000_000; + const dateNow = vi.spyOn(Date, "now").mockReturnValue(now); + const { entry, sessionKey } = await createEntry({ + status: "killed", + updatedAt: now, + }); + expect(check(entry, sessionKey)).toBe(true); + + await upsertSessionEntry({ agentId: "main", sessionKey, storePath }, entry); + const refreshed = loadSessionEntry({ agentId: "main", sessionKey, storePath }); + dateNow.mockRestore(); + + if (!refreshed) { + throw new Error("expected refreshed session entry"); + } + expect(check(refreshed, sessionKey)).toBe(false); + }); + + it("does not rotate non-main sessions or rows newer than the transcript", async () => { + const nonMain = await createEntry({ + sessionKey: "agent:main:other", + status: "killed", + updatedAt: Date.now() - 10_000, + }); + const newerRegistry = await createEntry({ + status: "timeout", + updatedAt: Date.now() + 10_000, + }); + await upsertSessionEntry( + { agentId: "main", sessionKey: newerRegistry.sessionKey, storePath }, + newerRegistry.entry, + ); + const refreshedRegistry = loadSessionEntry({ + agentId: "main", + sessionKey: newerRegistry.sessionKey, + storePath, + }); + if (!refreshedRegistry) { + throw new Error("expected refreshed registry entry"); + } + + expect(check(nonMain.entry, nonMain.sessionKey)).toBe(false); + expect(check(refreshedRegistry, newerRegistry.sessionKey)).toBe(false); + }); +}); diff --git a/src/config/sessions/lifecycle.ts b/src/config/sessions/lifecycle.ts index 2c6d499f09a9..00f4279300d7 100644 --- a/src/config/sessions/lifecycle.ts +++ b/src/config/sessions/lifecycle.ts @@ -1,6 +1,5 @@ // Session lifecycle timestamps prefer store metadata and fall back to transcript headers. import fs from "node:fs"; -import fsp from "node:fs/promises"; import { asDateTimestampMs } from "../../shared/number-coercion.js"; import { canonicalizeMainSessionAlias } from "./main-session.js"; import { @@ -8,6 +7,7 @@ import { resolveSessionFilePathOptions, type SessionFilePathOptions, } from "./paths.js"; +import { readTranscriptStatsSync } from "./session-accessor.js"; import { isTerminalSessionStatus, type SessionEntry, type SessionScope } from "./types.js"; type SessionLifecycleEntry = Pick< @@ -242,13 +242,13 @@ export function resolveTerminalMainSessionTranscriptRegistryCheck( return { sessionId, registryTimestampMs }; } -function isTranscriptMtimeNewerThanRegistry(params: { - transcriptMtimeMs: number; +function isTranscriptMutationNewerThanRegistry(params: { + transcriptMutationAtMs: number; registryTimestampMs: number; }): boolean { - const transcriptMtimeMs = Math.floor(params.transcriptMtimeMs); + const transcriptMutationAtMs = Math.floor(params.transcriptMutationAtMs); const registryTimestampMs = Math.floor(params.registryTimestampMs); - return Number.isFinite(transcriptMtimeMs) && transcriptMtimeMs > registryTimestampMs; + return Number.isFinite(transcriptMutationAtMs) && transcriptMutationAtMs > registryTimestampMs; } export function hasTerminalMainSessionTranscriptNewerThanRegistrySync( @@ -258,16 +258,20 @@ export function hasTerminalMainSessionTranscriptNewerThanRegistrySync( if (!check) { return false; } - const pathOptions = resolveSessionFilePathOptions({ - agentId: params.agentId, - storePath: params.storePath, - }); try { - const sessionFile = resolveSessionFilePath(check.sessionId, params.entry, pathOptions); - const stats = fs.statSync(sessionFile); - return isTranscriptMtimeNewerThanRegistry({ - transcriptMtimeMs: stats.mtimeMs, - registryTimestampMs: check.registryTimestampMs, + // Runtime transcripts are SQLite-only. Legacy-looking sessionFile values still + // resolve through agent/session/store scope, so a file stat would read stale state. + const stats = readTranscriptStatsSync({ + agentId: params.agentId, + sessionId: check.sessionId, + storePath: params.storePath, + }); + if (stats.lastMutationAtMs === undefined) { + return false; + } + return isTranscriptMutationNewerThanRegistry({ + transcriptMutationAtMs: stats.lastMutationAtMs, + registryTimestampMs: stats.lastObservedMutationAtMs ?? check.registryTimestampMs, }); } catch { return false; @@ -277,23 +281,5 @@ export function hasTerminalMainSessionTranscriptNewerThanRegistrySync( export async function hasTerminalMainSessionTranscriptNewerThanRegistry( params: TerminalMainSessionTranscriptRegistryParams, ): Promise { - const check = resolveTerminalMainSessionTranscriptRegistryCheck(params); - if (!check) { - return false; - } - const pathOptions = resolveSessionFilePathOptions({ - agentId: params.agentId, - storePath: params.storePath, - }); - try { - // Session admission owns this bounded stat as the terminal-main reconciliation gate. - const sessionFile = resolveSessionFilePath(check.sessionId, params.entry, pathOptions); - const stats = await fsp.stat(sessionFile); - return isTranscriptMtimeNewerThanRegistry({ - transcriptMtimeMs: stats.mtimeMs, - registryTimestampMs: check.registryTimestampMs, - }); - } catch { - return false; - } + return hasTerminalMainSessionTranscriptNewerThanRegistrySync(params); } diff --git a/src/config/sessions/session-accessor.sqlite-contract.ts b/src/config/sessions/session-accessor.sqlite-contract.ts index 160961e75c61..f04a741e3479 100644 --- a/src/config/sessions/session-accessor.sqlite-contract.ts +++ b/src/config/sessions/session-accessor.sqlite-contract.ts @@ -63,6 +63,8 @@ export type TranscriptEvent = unknown; export type SessionTranscriptStats = { eventCount: number; + lastMutationAtMs?: number; + lastObservedMutationAtMs?: number; maxSeq: number; sizeBytes: number; }; diff --git a/src/config/sessions/session-accessor.sqlite.ts b/src/config/sessions/session-accessor.sqlite.ts index 7a2ff0f3a61f..55fe35252898 100644 --- a/src/config/sessions/session-accessor.sqlite.ts +++ b/src/config/sessions/session-accessor.sqlite.ts @@ -297,6 +297,7 @@ export type SqliteSessionImportRowsParams = { sessionKey: string; entry: SessionEntry; readTranscriptEvents?: (append: (event: TranscriptEvent) => void) => void; + transcriptMtimeMs?: number; }; /** Summary of rows written by an internal doctor/migration import. */ @@ -1393,8 +1394,21 @@ export function readSqliteTranscriptStatsSync( ]) .where("session_id", "=", resolved.sessionId), ); + const session = executeSqliteQueryTakeFirstSync( + database.db, + db + .selectFrom("sessions") + .select(["transcript_observed_at", "transcript_updated_at"]) + .where("session_id", "=", resolved.sessionId), + ); return { eventCount: row?.event_count ?? 0, + ...(session?.transcript_updated_at !== null && session?.transcript_updated_at !== undefined + ? { lastMutationAtMs: session.transcript_updated_at } + : {}), + ...(session?.transcript_observed_at !== null && session?.transcript_observed_at !== undefined + ? { lastObservedMutationAtMs: session.transcript_observed_at } + : {}), maxSeq: row?.max_seq ?? 0, sizeBytes: row?.size_bytes ?? 0, }; @@ -1608,6 +1622,9 @@ export async function deleteSqliteTranscript(scope: SessionTranscriptReadScope): let deleted = false; runOpenClawAgentWriteTransaction((database) => { deleted = deleteSqliteTranscriptEventsInTransaction(database, resolved.sessionId); + if (deleted) { + touchTranscriptMutationInTransaction(database, resolved.sessionId); + } }, toDatabaseOptions(resolved)); return deleted; }); @@ -1689,12 +1706,25 @@ export async function importSqliteSessionRows( if (existingEventJson.has(eventJson)) { return; } - if (appendTranscriptEventInTransaction(database, transcriptScope, event)) { + if ( + appendTranscriptEventInTransaction(database, transcriptScope, event, { + touchMutation: false, + }) + ) { existingEventJson.add(eventJson); transcriptEvents += 1; } }); } + if (params.transcriptMtimeMs !== undefined) { + advanceTranscriptMutationAtInTransaction( + database, + params.entry.sessionId, + params.transcriptMtimeMs, + ); + } else if (transcriptEvents > 0) { + touchTranscriptMutationInTransaction(database, params.entry.sessionId); + } }, toDatabaseOptions(resolved)); return { sessionId: params.entry.sessionId, @@ -1745,9 +1775,7 @@ export async function appendSqliteTranscriptEvents( const resolved = resolveSqliteTranscriptScope(scope); await runExclusiveSqliteSessionWrite(resolved, async () => { runOpenClawAgentWriteTransaction((database) => { - for (const event of events) { - appendTranscriptEventInTransaction(database, resolved, event); - } + appendTranscriptEventsInTransaction(database, resolved, events); }, toDatabaseOptions(resolved)); }); } @@ -3897,7 +3925,15 @@ function writeSessionEntry( const db = getSessionKysely(database.db); const normalizedEntry = normalizeSqliteSessionEntryTimestamp(entry); const updatedAt = normalizedEntry.updatedAt; - const sessionRow = bindSqliteSessionRoot({ entry: normalizedEntry, sessionKey, updatedAt }); + // Registry writes snapshot the current transcript watermark so recovery can + // distinguish same-millisecond transcript writes before and after this row. + const transcriptObservedAt = + readTranscriptMutationStateInTransaction(database, normalizedEntry.sessionId).updatedAt ?? + updatedAt; + const sessionRow = { + ...bindSqliteSessionRoot({ entry: normalizedEntry, sessionKey, updatedAt }), + transcript_observed_at: transcriptObservedAt, + }; executeSqliteQuerySync( database.db, db @@ -3907,6 +3943,7 @@ function writeSessionEntry( conflict.column("session_id").doUpdateSet({ session_key: sessionKey, session_scope: sessionRow.session_scope, + transcript_observed_at: transcriptObservedAt, updated_at: updatedAt, started_at: sessionRow.started_at, ended_at: sessionRow.ended_at, @@ -4134,6 +4171,67 @@ function readNextTranscriptSeq(database: OpenClawAgentDatabase, sessionId: strin return maxSeq + 1; } +function normalizeTranscriptMutationAtMs(value: number): number | undefined { + const timestamp = Math.floor(value); + return Number.isFinite(timestamp) && timestamp >= 0 ? timestamp : undefined; +} + +function readTranscriptMutationStateInTransaction( + database: OpenClawAgentDatabase, + sessionId: string, +): { observedAt: number | null; updatedAt: number | null } { + const db = getSessionKysely(database.db); + const row = executeSqliteQueryTakeFirstSync( + database.db, + db + .selectFrom("sessions") + .select(["transcript_observed_at", "transcript_updated_at"]) + .where("session_id", "=", sessionId), + ); + return { + observedAt: row?.transcript_observed_at ?? null, + updatedAt: row?.transcript_updated_at ?? null, + }; +} + +function advanceTranscriptMutationAtInTransaction( + database: OpenClawAgentDatabase, + sessionId: string, + value: number, + options: { strictly?: boolean } = {}, +): void { + const transcriptUpdatedAt = normalizeTranscriptMutationAtMs(value); + if (transcriptUpdatedAt === undefined) { + return; + } + const state = readTranscriptMutationStateInTransaction(database, sessionId); + const next = options.strictly + ? Math.max(transcriptUpdatedAt, (state.updatedAt ?? -1) + 1, (state.observedAt ?? -1) + 1) + : Math.max(transcriptUpdatedAt, state.updatedAt ?? 0); + if (state.updatedAt !== null && state.updatedAt >= next) { + return; + } + const db = getSessionKysely(database.db); + executeSqliteQuerySync( + database.db, + db + .updateTable("sessions") + .set({ transcript_updated_at: next }) + .where("session_id", "=", sessionId), + ); +} + +function touchTranscriptMutationInTransaction( + database: OpenClawAgentDatabase, + sessionId: string, +): void { + const now = normalizeTranscriptMutationAtMs(Date.now()); + if (now === undefined) { + return; + } + advanceTranscriptMutationAtInTransaction(database, sessionId, now, { strictly: true }); +} + function deleteSqliteTranscriptEventsInTransaction( database: OpenClawAgentDatabase, sessionId: string, @@ -4490,16 +4588,18 @@ function writeSqliteParentForkTranscriptInTransaction( ...(params.source.appendMode ? { appendMode: params.source.appendMode } : {}), } : null; - appendTranscriptEventInTransaction(database, targetScope, { - ...createSessionTranscriptHeader({ - cwd: params.source.cwd, - sessionId: targetScope.sessionId, - }), - parentSession: params.parentSessionFile, - }); - for (const event of [...pathEntries, ...labelEntries, ...(leafEntry ? [leafEntry] : [])]) { - appendTranscriptEventInTransaction(database, targetScope, event); - } + appendTranscriptEventsInTransaction(database, targetScope, [ + { + ...createSessionTranscriptHeader({ + cwd: params.source.cwd, + sessionId: targetScope.sessionId, + }), + parentSession: params.parentSessionFile, + }, + ...pathEntries, + ...labelEntries, + ...(leafEntry ? [leafEntry] : []), + ]); } function forkSqliteParentTranscriptInTransaction( @@ -4679,20 +4779,13 @@ function forkSqliteCheckpointTranscriptInTransaction( sessionKey: params.targetSessionKey, }; const sessionFile = formatSqliteSessionMarkerForScope(targetScope); - appendTranscriptEventInTransaction( - database, - targetScope, + appendTranscriptEventsInTransaction(database, targetScope, [ createSessionTranscriptHeader({ cwd: readTranscriptHeaderCwd(selected.rows), sessionId, }), - ); - for (const event of selected.rows) { - if (isSessionTranscriptHeader(event)) { - continue; - } - appendTranscriptEventInTransaction(database, targetScope, event); - } + ...selected.rows.filter((event) => !isSessionTranscriptHeader(event)), + ]); return { status: "created", sessionId, @@ -4834,7 +4927,7 @@ function appendTranscriptEventInTransaction( database: OpenClawAgentDatabase, scope: ResolvedTranscriptScope, event: TranscriptEvent, - options: { dedupeByMessageIdempotency?: boolean } = {}, + options: { dedupeByMessageIdempotency?: boolean; touchMutation?: boolean } = {}, ): boolean { const db = getSessionKysely(database.db); const createdAt = readEventTimestamp(event) ?? Date.now(); @@ -4864,6 +4957,9 @@ function appendTranscriptEventInTransaction( created_at: createdAt, }), ); + if (options.touchMutation !== false) { + touchTranscriptMutationInTransaction(database, scope.sessionId); + } if (!identity) { return true; } @@ -4897,6 +4993,27 @@ function appendTranscriptEventInTransaction( return true; } +function appendTranscriptEventsInTransaction( + database: OpenClawAgentDatabase, + scope: ResolvedTranscriptScope, + events: readonly TranscriptEvent[], +): number { + let appended = 0; + for (const event of events) { + if ( + appendTranscriptEventInTransaction(database, scope, event, { + touchMutation: false, + }) + ) { + appended += 1; + } + } + if (appended > 0) { + touchTranscriptMutationInTransaction(database, scope.sessionId); + } + return appended; +} + function appendTranscriptEventRowInTransaction( database: OpenClawAgentDatabase, scope: ResolvedTranscriptScope, @@ -5041,8 +5158,11 @@ function replaceSqliteTranscriptEventsInTransaction( resolved: ResolvedTranscriptScope, events: readonly TranscriptEvent[], ): void { - deleteSqliteTranscriptEventsInTransaction(database, resolved.sessionId); + const deleted = deleteSqliteTranscriptEventsInTransaction(database, resolved.sessionId); if (events.length === 0) { + if (deleted) { + touchTranscriptMutationInTransaction(database, resolved.sessionId); + } return; } ensureTranscriptSessionRoot(database, resolved, readEventTimestamp(events[0]) ?? Date.now()); @@ -5058,6 +5178,9 @@ function replaceSqliteTranscriptEventsInTransaction( seq += 1; } } + if (deleted || seq > 0) { + touchTranscriptMutationInTransaction(database, resolved.sessionId); + } } function readTranscriptIdentityByEventId( diff --git a/src/config/sessions/session-accessor.test.ts b/src/config/sessions/session-accessor.test.ts index 2ee99476ef26..b24b9a203b7b 100644 --- a/src/config/sessions/session-accessor.test.ts +++ b/src/config/sessions/session-accessor.test.ts @@ -24,6 +24,7 @@ import { patchSessionEntryTarget, persistSessionResetLifecycle, persistSessionTranscriptTurn, + readTranscriptStatsSync, readSessionUpdatedAt, recordInboundSessionMeta, replaceSessionEntry, @@ -2424,6 +2425,55 @@ describe("session accessor seam", () => { ).toBe(`sqlite:main:session-1:${path.join(tempDir, "openclaw-agent.sqlite")}`); }); + it("tracks replacement and deletion transcript mutations", async () => { + const dateNow = vi.spyOn(Date, "now").mockReturnValue(1_700_000_000_000); + const scope = { + agentId: "main", + sessionId: "session-1", + sessionKey: "agent:main:main", + storePath, + }; + await upsertSessionEntry(scope, { + sessionId: scope.sessionId, + updatedAt: 10, + }); + await replaceSqliteTranscriptEvents(scope, [ + { sessionId: scope.sessionId, type: "session" }, + { timestamp: "1970-01-01T00:00:00.001Z", type: "custom" }, + ]); + + const replaced = readTranscriptStatsSync(scope); + expect(replaced).toMatchObject({ + eventCount: 2, + lastMutationAtMs: expect.any(Number), + }); + expect(replaced.lastMutationAtMs).toBeGreaterThanOrEqual(1_700_000_000_000); + + await importSqliteSessionRows({ + agentId: scope.agentId, + entry: { + sessionId: scope.sessionId, + updatedAt: 10, + }, + sessionKey: scope.sessionKey, + storePath: scope.storePath, + transcriptMtimeMs: 1_600_000_000_000, + }); + const imported = readTranscriptStatsSync(scope); + expect(imported.lastMutationAtMs).toBe(replaced.lastMutationAtMs); + expect(imported.lastObservedMutationAtMs).toBe(replaced.lastMutationAtMs); + + await replaceSqliteTranscriptEvents(scope, []); + + const cleared = readTranscriptStatsSync(scope); + dateNow.mockRestore(); + expect(cleared).toMatchObject({ + eventCount: 0, + lastMutationAtMs: expect.any(Number), + }); + expect(cleared.lastMutationAtMs).toBeGreaterThan(imported.lastMutationAtMs ?? 0); + }); + it("resolves an explicit read transcript file without agent identity", () => { const explicitSessionFile = path.join(tempDir, "explicit-read-session.jsonl"); diff --git a/src/config/sessions/session-accessor.ts b/src/config/sessions/session-accessor.ts index 6165d46baae1..223ce37a4254 100644 --- a/src/config/sessions/session-accessor.ts +++ b/src/config/sessions/session-accessor.ts @@ -289,6 +289,8 @@ export type TranscriptEvent = unknown; export type SessionTranscriptStats = { eventCount: number; + lastMutationAtMs?: number; + lastObservedMutationAtMs?: number; maxSeq: number; sizeBytes: number; }; diff --git a/src/state/openclaw-agent-db.generated.d.ts b/src/state/openclaw-agent-db.generated.d.ts index 60d745cbf666..8ac04f8b795c 100644 --- a/src/state/openclaw-agent-db.generated.d.ts +++ b/src/state/openclaw-agent-db.generated.d.ts @@ -138,6 +138,8 @@ export interface Sessions { spawned_by: string | null; started_at: number | null; status: string | null; + transcript_observed_at: Generated; + transcript_updated_at: Generated; updated_at: number; } diff --git a/src/state/openclaw-agent-db.test.ts b/src/state/openclaw-agent-db.test.ts index edc5e810757a..1b99c15354f1 100644 --- a/src/state/openclaw-agent-db.test.ts +++ b/src/state/openclaw-agent-db.test.ts @@ -1027,6 +1027,82 @@ describe("openclaw agent database", () => { }); }); + it("adds transcript mutation watermarks to v4 session tables", () => { + const stateDir = createTempStateDir(); + const databasePath = path.join( + stateDir, + "agents", + "worker-1", + "agent", + "openclaw-agent.sqlite", + ); + fs.mkdirSync(path.dirname(databasePath), { recursive: true }); + const currentSchema = fs.readFileSync( + new URL("./openclaw-agent-schema.sql", import.meta.url), + "utf8", + ); + const previousSchema = currentSchema.replace( + [ + " transcript_updated_at INTEGER DEFAULT NULL,\n", + " transcript_observed_at INTEGER DEFAULT NULL,\n", + ].join(""), + "", + ); + expect(previousSchema).not.toBe(currentSchema); + const { DatabaseSync } = requireNodeSqlite(); + const db = new DatabaseSync(databasePath); + db.exec(previousSchema); + db.exec(` + INSERT INTO schema_meta + (meta_key, role, schema_version, agent_id, app_version, created_at, updated_at) + VALUES ('primary', 'agent', 4, 'worker-1', NULL, 1, 1); + INSERT INTO sessions + (session_id, session_key, created_at, updated_at) + VALUES ('session-1', 'agent:worker-1:main', 10, 20); + INSERT INTO sessions + (session_id, session_key, created_at, updated_at) + VALUES ('session-2', 'agent:worker-1:other', 10, 20); + INSERT INTO transcript_events + (session_id, seq, event_json, created_at) + VALUES ('session-1', 0, '{"type":"custom"}', 1); + PRAGMA user_version = 4; + `); + db.close(); + + const database = openOpenClawAgentDatabase({ + agentId: "worker-1", + env: { OPENCLAW_STATE_DIR: stateDir }, + }); + const columns = database.db.prepare("PRAGMA table_info(sessions)").all() as Array<{ + name?: unknown; + }>; + + expect(columns.map((column) => column.name)).toEqual( + expect.arrayContaining(["transcript_observed_at", "transcript_updated_at"]), + ); + expect( + database.db + .prepare( + "SELECT transcript_observed_at, transcript_updated_at FROM sessions WHERE session_id = ?", + ) + .get("session-1"), + ).toEqual({ + transcript_observed_at: 20, + transcript_updated_at: expect.any(Number), + }); + expect( + database.db + .prepare( + "SELECT transcript_observed_at, transcript_updated_at FROM sessions WHERE session_id = ?", + ) + .get("session-2"), + ).toEqual({ + transcript_observed_at: null, + transcript_updated_at: null, + }); + expect(readSqliteNumberPragma(database.db, "user_version")).toBe(OPENCLAW_AGENT_SCHEMA_VERSION); + }); + it("inspects registered database ownership without mutating the database", () => { const stateDir = createTempStateDir(); const database = openOpenClawAgentDatabase({ @@ -1132,7 +1208,7 @@ describe("openclaw agent database", () => { env: { OPENCLAW_STATE_DIR: stateDir }, }); - expect(readSqliteNumberPragma(database.db, "user_version")).toBe(4); + expect(readSqliteNumberPragma(database.db, "user_version")).toBe(OPENCLAW_AGENT_SCHEMA_VERSION); const session = database.db .prepare( ` diff --git a/src/state/openclaw-agent-db.ts b/src/state/openclaw-agent-db.ts index 182d45ac1bbe..94ebc7525f21 100644 --- a/src/state/openclaw-agent-db.ts +++ b/src/state/openclaw-agent-db.ts @@ -46,10 +46,10 @@ export { resolveOpenClawAgentSqlitePath } from "./openclaw-agent-db.paths.js"; * per pathname, protected with private file modes, and registered in the shared * OpenClaw state database for discovery and maintenance. */ -// v4 = session/transcript flip (branch lineage). Main's v2 memory-identity +// v5 = transcript mutation watermark. The v4 session/transcript flip and main's v2 memory-identity // change is folded in structure-gated (migrateMemoryIndexSourcesIdentity), so // v2 main DBs and pre-merge v4 flip DBs both converge on this schema. -export const OPENCLAW_AGENT_SCHEMA_VERSION = 4; +export const OPENCLAW_AGENT_SCHEMA_VERSION = 5; const OPENCLAW_AGENT_DB_DIR_MODE = 0o700; const OPENCLAW_AGENT_DB_FILE_MODE = 0o600; const OPENCLAW_AGENT_DB_SLOW_OPEN_MS = 1_000; @@ -155,6 +155,28 @@ function dropLegacyMemoryIndexSchema(db: DatabaseSync): void { `); } +function backfillTranscriptMutationWatermarks(db: DatabaseSync): void { + const transcriptTable = db + .prepare("SELECT 1 AS ok FROM sqlite_master WHERE type = 'table' AND name = ?") + .get("transcript_events") as { ok?: unknown } | undefined; + if (transcriptTable?.ok !== 1) { + return; + } + db.prepare( + ` + UPDATE sessions + SET + transcript_updated_at = COALESCE(transcript_updated_at, ?), + transcript_observed_at = COALESCE(transcript_observed_at, updated_at) + WHERE EXISTS ( + SELECT 1 + FROM transcript_events + WHERE transcript_events.session_id = sessions.session_id + ) + `, + ).run(Date.now()); +} + function migrateOpenClawAgentSchema(db: DatabaseSync): void { const userVersion = readSqliteUserVersion(db); if (userVersion >= OPENCLAW_AGENT_SCHEMA_VERSION) { @@ -164,7 +186,17 @@ function migrateOpenClawAgentSchema(db: DatabaseSync): void { db.exec("DROP INDEX IF EXISTS idx_agent_transcript_events_session;"); } const columns = readSqliteSessionColumns(db); - if (userVersion > 1 || !columns) { + if (columns && !columns.has("transcript_updated_at")) { + db.exec("ALTER TABLE sessions ADD COLUMN transcript_updated_at INTEGER DEFAULT NULL;"); + } + if (columns && !columns.has("transcript_observed_at")) { + db.exec("ALTER TABLE sessions ADD COLUMN transcript_observed_at INTEGER DEFAULT NULL;"); + } + if (!columns) { + return; + } + if (userVersion > 1) { + backfillTranscriptMutationWatermarks(db); return; } const copyColumns = [ @@ -232,6 +264,8 @@ function migrateOpenClawAgentSchema(db: DatabaseSync): void { session_scope TEXT NOT NULL DEFAULT 'conversation' CHECK (session_scope IN ('conversation', 'shared-main', 'group', 'channel')), created_at INTEGER NOT NULL, updated_at INTEGER NOT NULL, + transcript_updated_at INTEGER DEFAULT NULL, + transcript_observed_at INTEGER DEFAULT NULL, started_at INTEGER, ended_at INTEGER, status TEXT CHECK (status IS NULL OR status IN ('running', 'done', 'failed', 'killed', 'timeout')), @@ -252,6 +286,7 @@ function migrateOpenClawAgentSchema(db: DatabaseSync): void { DROP TABLE sessions; ALTER TABLE sessions_new RENAME TO sessions; `); + backfillTranscriptMutationWatermarks(db); } function parseMigratedSessionEntry(value: unknown): MigratedSessionEntry | null { diff --git a/src/state/openclaw-agent-schema.generated.ts b/src/state/openclaw-agent-schema.generated.ts index 829f5e65e5a8..7b58064f30d6 100644 --- a/src/state/openclaw-agent-schema.generated.ts +++ b/src/state/openclaw-agent-schema.generated.ts @@ -19,6 +19,8 @@ CREATE TABLE IF NOT EXISTS sessions ( session_scope TEXT NOT NULL DEFAULT 'conversation' CHECK (session_scope IN ('conversation', 'shared-main', 'group', 'channel')), created_at INTEGER NOT NULL, updated_at INTEGER NOT NULL, + transcript_updated_at INTEGER DEFAULT NULL, + transcript_observed_at INTEGER DEFAULT NULL, started_at INTEGER, ended_at INTEGER, status TEXT CHECK (status IS NULL OR status IN ('running', 'done', 'failed', 'killed', 'timeout')), diff --git a/src/state/openclaw-agent-schema.sql b/src/state/openclaw-agent-schema.sql index 78d3d14a4c45..c38f7b1b4366 100644 --- a/src/state/openclaw-agent-schema.sql +++ b/src/state/openclaw-agent-schema.sql @@ -14,6 +14,8 @@ CREATE TABLE IF NOT EXISTS sessions ( session_scope TEXT NOT NULL DEFAULT 'conversation' CHECK (session_scope IN ('conversation', 'shared-main', 'group', 'channel')), created_at INTEGER NOT NULL, updated_at INTEGER NOT NULL, + transcript_updated_at INTEGER DEFAULT NULL, + transcript_observed_at INTEGER DEFAULT NULL, started_at INTEGER, ended_at INTEGER, status TEXT CHECK (status IS NULL OR status IN ('running', 'done', 'failed', 'killed', 'timeout')),