Files
openclaw/src/sessions/transcript-events.ts
T
Josh Lehman 0a8e3604ba refactor: flip sessions and transcripts to sqlite storage (#98236)
* refactor(sessions): migrate runtime storage to sqlite

* test(sessions): fix sqlite CI regressions

* test(sessions): align remaining sqlite fixtures

* fix(codex): require sqlite trajectory recorder

* test(sessions): align orphan recovery sqlite fixture

* test(sessions): align sqlite rebase fixtures

* fix(sessions): finish current-main integration of the sqlite flip

Resolve the whole-store SDK removal across its owner boundary: drop the
loadSessionStore re-export and the registry whole-store wrappers, wire
hasTrackedActiveSessionRun into gateway chat, complete the
preserveLockedHarnessIds cleanup contract, flip the codex thread-history
import to storePath targets, and port remaining main-side tests from
file-store helpers to session accessor reads.

* chore: drop committed pebbles log, revert plugin-inspector bump, refresh generated docs

Remove the 1.8k-line .pebbles/events.jsonl work log from the branch, restore
the plugin-inspector advisory lane to main's pinned 0.3.10 so the supply-chain
bump gets its own review, and regenerate docs_map, the plugin SDK API baseline,
and the export-surface ratchet for the merged tree.

* feat(sessions): keep archived transcripts by default with zstd cold storage

Codex-style retention: deleting or resetting a session archives its
transcript as a zstd-compressed JSONL artifact (plain when the runtime
lacks node:zlib zstd) and keeps it until the disk budget evicts oldest
first. resetArchiveRetention now governs both deleted and reset archives
and defaults to keep; maxDiskBytes defaults to 2gb so retention stays
bounded, with archives evicted before live sessions. The cron reaper
follows the same knob instead of deleting archives on its own timer.

* fix(state): converge agent DB migration lineages and bound database growth

Merge coherence: run both structure-gated legacy memory-schema repairs
(flip-lineage drop, main-lineage identity rebuild) before the flip
migration so pre-flip v1/v2 and pre-merge flip v1/v4 databases all
converge, and hoist foreign_keys=OFF outside the schema transaction
where the pragma was silently ignored and the v1 sessions rebuild
cascade-deleted session_entries.

Growth guards: fresh agent DBs enable auto_vacuum=INCREMENTAL, WAL
maintenance releases freed pages in bounded passes (never a blocking
full VACUUM), and doctor reports state/agent DB bloat from freelist
stats.

* fix(codex): resolve the store path for thread-history import via the SDK

The supervision catalog passed the legacy sessionFile locator to the
storePath-targeted transcript mirror; resolve the agent store path with
the session-store SDK helper instead of a runtime-object seam so test
fakes and headless callers need no extra surface. Drop the obsolete
missing-session-id preprocessing case: sessions rows are NOT NULL on
session_id and upsert repairs id-less patches at write time.

* fix(sessions): fail safe on malformed disk-budget config and doctor stat errors

A malformed explicit maxDiskBytes disables the budget instead of
falling back to the destructive 2gb default the user never chose, and
the doctor bloat check skips databases whose paths stat-fail instead of
aborting doctor.

* fix(sessions): complete sqlite conflict translations

* test(sqlite): align hardening checks with maintenance

* test(sessions): inspect compressed transcript archives

* fix(tests): await session seeds and drop unused helpers flagged by CI lint

The five unawaited writeSessionStoreSeed calls raced their SQLite seeds
against the assertions, failing compact shards; the bloat probe drops a
useless initializer and the merged tests drop now-unused helpers.

* test(sessions): type legacy proof events directly

* test(sessions): align hardening contracts

* perf(sessions): read usage transcript sizes from SQL aggregates

Usage/cost scans walked every session and materialized every transcript
event just to re-stringify it for a byte estimate — the #86718 stall
class reborn on the DB. readTranscriptStatsSync sums stored JSON bytes
in SQLite without loading a single row.

* fix(sessions): re-root foreign-root transcript paths onto the current sessions dir

Restored backups, moved OPENCLAW_STATE_DIR, and rehearsal copies carry
absolute sessionFile paths from the old root; the containment fallback
kept those foreign paths, so migration read (and would archive) files in
the original root and reported local copies missing. Re-root the
canonical agents/<id>/sessions suffix onto the current dir when the file
exists there; genuine cross-root layouts still fall through unchanged.

* test(agents): seed harness admission through sqlite

* fix(sqlite): close agent db on pragma setup failure

* fix(doctor): compact and retrofit incremental auto-vacuum after session import

The migration is the sanctioned offline window: post-import compact
reclaims import churn and applies auto_vacuum=INCREMENTAL to databases
created before the fresh-DB pragma existed, so runtime maintenance can
release pages in bounded passes on every install.

---------

Co-authored-by: Peter Steinberger <steipete@gmail.com>
2026-07-11 14:50:37 -07:00

177 lines
6.2 KiB
TypeScript

// Transcript event helpers serialize and trim session transcript events.
import { asPositiveSafeInteger } from "@openclaw/normalization-core/number-coercion";
import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce";
import { parseAgentSessionKey } from "../routing/session-key.js";
/** Storage-neutral identity for the session transcript that changed. */
export type SessionTranscriptUpdateTarget = {
agentId: string;
sessionId: string;
sessionKey: string;
};
type SessionTranscriptUpdateFields = {
sessionFile?: string;
target?: SessionTranscriptUpdateTarget;
sessionKey?: string;
agentId?: string;
sessionId?: string;
message?: unknown;
messageId?: string;
messageSeq?: number;
};
/** Normalized transcript update emitted after a session transcript changes. */
export type SessionTranscriptUpdate = Omit<SessionTranscriptUpdateFields, "sessionFile"> & {
target: SessionTranscriptUpdateTarget;
};
/** Internal transcript update that may identify a transcript without a file path. */
export type InternalSessionTranscriptUpdate = SessionTranscriptUpdateFields;
type SessionTranscriptListener = (update: SessionTranscriptUpdate) => void;
type InternalSessionTranscriptListener = (update: InternalSessionTranscriptUpdate) => void;
const SESSION_TRANSCRIPT_LISTENERS = new Set<SessionTranscriptListener>();
const INTERNAL_SESSION_TRANSCRIPT_LISTENERS = new Set<InternalSessionTranscriptListener>();
/** Registers a listener for normalized session transcript updates. */
export function onSessionTranscriptUpdate(listener: SessionTranscriptListener): () => void {
SESSION_TRANSCRIPT_LISTENERS.add(listener);
return () => {
SESSION_TRANSCRIPT_LISTENERS.delete(listener);
};
}
/** Registers an internal listener for identity-only or file-backed transcript updates. */
export function onInternalSessionTranscriptUpdate(
listener: InternalSessionTranscriptListener,
): () => void {
INTERNAL_SESSION_TRANSCRIPT_LISTENERS.add(listener);
return () => {
INTERNAL_SESSION_TRANSCRIPT_LISTENERS.delete(listener);
};
}
/** Emits a normalized transcript update to all registered listeners. */
export function emitSessionTranscriptUpdate(update: InternalSessionTranscriptUpdate): void {
const nextUpdate = normalizeSessionTranscriptUpdate(update, { allowIdentityOnly: true });
if (!nextUpdate) {
return;
}
const publicUpdate = projectPublicSessionTranscriptUpdate(nextUpdate);
if (publicUpdate) {
emitPublicSessionTranscriptUpdate(publicUpdate);
}
emitInternalTranscriptUpdate(nextUpdate);
}
/** Emits an internal transcript update, including identity-only updates. */
export function emitInternalSessionTranscriptUpdate(update: InternalSessionTranscriptUpdate): void {
const nextUpdate = normalizeSessionTranscriptUpdate(update, { allowIdentityOnly: true });
if (!nextUpdate) {
return;
}
emitInternalTranscriptUpdate(nextUpdate);
}
function normalizeSessionTranscriptUpdate(
update: InternalSessionTranscriptUpdate,
options: { allowIdentityOnly: boolean },
): InternalSessionTranscriptUpdate | undefined {
const normalized = {
sessionFile: update.sessionFile,
target: update.target,
sessionKey: update.sessionKey,
agentId: update.agentId,
sessionId: update.sessionId,
message: update.message,
messageId: update.messageId,
messageSeq: update.messageSeq,
};
const trimmed = normalizeOptionalString(normalized.sessionFile);
const target = normalizeUpdateTarget(normalized);
if (!trimmed && (!options.allowIdentityOnly || !target)) {
return undefined;
}
const messageSeq = asPositiveSafeInteger(normalized.messageSeq);
const sessionKey = normalizeOptionalString(normalized.sessionKey) ?? target?.sessionKey;
const agentId = normalizeOptionalString(normalized.agentId) ?? target?.agentId;
const sessionId = normalizeOptionalString(normalized.sessionId) ?? target?.sessionId;
return {
...(trimmed ? { sessionFile: trimmed } : {}),
...(target ? { target } : {}),
...(sessionKey ? { sessionKey } : {}),
...(agentId ? { agentId } : {}),
...(sessionId ? { sessionId } : {}),
...(normalized.message !== undefined ? { message: normalized.message } : {}),
...(normalizeOptionalString(normalized.messageId)
? { messageId: normalizeOptionalString(normalized.messageId) }
: {}),
...(messageSeq !== undefined ? { messageSeq } : {}),
};
}
function emitPublicSessionTranscriptUpdate(nextUpdate: SessionTranscriptUpdate): void {
for (const listener of SESSION_TRANSCRIPT_LISTENERS) {
try {
listener(nextUpdate);
} catch {
/* ignore */
}
}
}
function emitInternalTranscriptUpdate(nextUpdate: InternalSessionTranscriptUpdate): void {
for (const listener of INTERNAL_SESSION_TRANSCRIPT_LISTENERS) {
try {
listener(nextUpdate);
} catch {
/* ignore */
}
}
}
function projectPublicSessionTranscriptUpdate(
update: InternalSessionTranscriptUpdate,
): SessionTranscriptUpdate | undefined {
const target = update.target;
if (!target) {
return undefined;
}
return {
target,
...(update.sessionKey ? { sessionKey: update.sessionKey } : {}),
...(update.agentId ? { agentId: update.agentId } : {}),
...(update.sessionId ? { sessionId: update.sessionId } : {}),
...(update.message !== undefined ? { message: update.message } : {}),
...(update.messageId ? { messageId: update.messageId } : {}),
...(update.messageSeq !== undefined ? { messageSeq: update.messageSeq } : {}),
};
}
function normalizeUpdateTarget(update: {
agentId?: string;
sessionId?: string;
sessionKey?: string;
target?: SessionTranscriptUpdate["target"];
}): SessionTranscriptUpdateTarget | undefined {
const sessionKey =
normalizeOptionalString(update.target?.sessionKey) ??
normalizeOptionalString(update.sessionKey);
const agentId =
normalizeOptionalString(update.target?.agentId) ??
normalizeOptionalString(update.agentId) ??
(sessionKey ? parseAgentSessionKey(sessionKey)?.agentId : undefined);
const sessionId =
normalizeOptionalString(update.target?.sessionId) ?? normalizeOptionalString(update.sessionId);
if (!agentId || !sessionId || !sessionKey) {
return undefined;
}
return {
agentId,
sessionId,
sessionKey,
};
}