Files
openclaw/scripts/e2e/lib/session-log-mentions.ts
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

262 lines
7.1 KiB
TypeScript

// Session Log Mentions script supports OpenClaw repository automation.
import fs from "node:fs/promises";
import path from "node:path";
import { DatabaseSync } from "node:sqlite";
import { readPositiveIntEnv } from "./env-limits.mjs";
type SessionLogMentionLimits = {
fileMaxBytes: number;
totalMaxBytes: number;
};
type SessionLogNeedles = Record<string, string>;
const DEFAULT_FILE_MAX_BYTES = 4 * 1024 * 1024;
const DEFAULT_TOTAL_MAX_BYTES = 16 * 1024 * 1024;
export function readSessionLogMentionLimits(
env: NodeJS.ProcessEnv = process.env,
): SessionLogMentionLimits {
return {
fileMaxBytes: readPositiveIntEnv(
"OPENCLAW_SESSION_LOG_MENTION_FILE_MAX_BYTES",
DEFAULT_FILE_MAX_BYTES,
env,
),
totalMaxBytes: readPositiveIntEnv(
"OPENCLAW_SESSION_LOG_MENTION_TOTAL_MAX_BYTES",
DEFAULT_TOTAL_MAX_BYTES,
env,
),
};
}
function taggedError(message: string, code: string) {
return Object.assign(new Error(message), { code });
}
function countOccurrences(haystack: string, needle: string): number {
if (!needle) {
return 0;
}
let count = 0;
let offset = 0;
for (;;) {
const next = haystack.indexOf(needle, offset);
if (next < 0) {
return count;
}
count += 1;
offset = next + needle.length;
}
}
function createCounts(needles: SessionLogNeedles): Record<string, number> {
return Object.fromEntries(Object.keys(needles).map((key) => [key, 0]));
}
function recordRole(record: unknown): string | undefined {
if (!record || typeof record !== "object") {
return undefined;
}
const candidate = record as { message?: unknown; role?: unknown };
if (typeof candidate.role === "string") {
return candidate.role;
}
if (!candidate.message || typeof candidate.message !== "object") {
return undefined;
}
const message = candidate.message as { role?: unknown };
return typeof message.role === "string" ? message.role : undefined;
}
function collectStringLeaves(value: unknown, output: string[]) {
if (typeof value === "string") {
output.push(value);
return;
}
if (Array.isArray(value)) {
for (const item of value) {
collectStringLeaves(item, output);
}
return;
}
if (!value || typeof value !== "object") {
return;
}
for (const item of Object.values(value)) {
collectStringLeaves(item, output);
}
}
function sessionLogScanText(line: string): string | null {
const trimmed = line.trim();
if (!trimmed) {
return null;
}
try {
const record = JSON.parse(trimmed) as unknown;
if (recordRole(record) === "user") {
return null;
}
const strings: string[] = [];
collectStringLeaves(record, strings);
return strings.join("\n");
} catch {
return line;
}
}
function assertWithinLimit(params: {
byteCount: number;
filePath?: string;
label: string;
limit: number;
}) {
if (params.byteCount <= params.limit) {
return;
}
const source = params.filePath ? ` ${params.filePath}` : "";
throw taggedError(
`session log mention scan exceeded ${params.label} limit${source}: ${params.byteCount} > ${params.limit}`,
"ETOOBIG",
);
}
export async function countSessionLogMentions(params: {
limits?: SessionLogMentionLimits;
needles: SessionLogNeedles;
sessionsDir: string;
}): Promise<Record<string, number>> {
const limits = params.limits ?? readSessionLogMentionLimits();
const counts = createCounts(params.needles);
const addCounts = (nextCounts: Record<string, number>) => {
for (const [key, count] of Object.entries(nextCounts)) {
counts[key] = (counts[key] ?? 0) + count;
}
};
let files: string[];
try {
files = await fs.readdir(params.sessionsDir);
} catch {
files = [];
}
let totalBytes = 0;
for (const file of files.filter((candidate) => candidate.endsWith(".jsonl")).toSorted()) {
const filePath = path.join(params.sessionsDir, file);
const stat = await fs.stat(filePath).catch(() => null);
if (!stat?.isFile()) {
continue;
}
assertWithinLimit({
byteCount: stat.size,
filePath,
label: "per-file",
limit: limits.fileMaxBytes,
});
totalBytes += stat.size;
assertWithinLimit({
byteCount: totalBytes,
label: "total",
limit: limits.totalMaxBytes,
});
const raw = await fs.readFile(filePath, "utf8").catch(() => "");
const actualBytes = Buffer.byteLength(raw, "utf8");
assertWithinLimit({
byteCount: actualBytes,
filePath,
label: "per-file",
limit: limits.fileMaxBytes,
});
for (const line of raw.split(/\r?\n/u)) {
const scanText = sessionLogScanText(line);
if (scanText === null) {
continue;
}
for (const [key, needle] of Object.entries(params.needles)) {
counts[key] += countOccurrences(scanText, needle);
}
}
}
addCounts(
await countSqliteTranscriptMentions({
limits,
needles: params.needles,
sessionsDir: params.sessionsDir,
startingBytes: totalBytes,
}),
);
return counts;
}
function resolveAgentSqlitePathFromSessionsDir(sessionsDir: string): string | null {
if (path.basename(sessionsDir) !== "sessions") {
return null;
}
return path.join(path.dirname(sessionsDir), "agent", "openclaw-agent.sqlite");
}
async function countSqliteTranscriptMentions(params: {
limits: SessionLogMentionLimits;
needles: SessionLogNeedles;
sessionsDir: string;
startingBytes: number;
}): Promise<Record<string, number>> {
const counts = createCounts(params.needles);
const sqlitePath = resolveAgentSqlitePathFromSessionsDir(params.sessionsDir);
if (!sqlitePath) {
return counts;
}
const stat = await fs.stat(sqlitePath).catch(() => null);
if (!stat?.isFile()) {
return counts;
}
let totalBytes = params.startingBytes;
let db: DatabaseSync | null = null;
try {
db = new DatabaseSync(sqlitePath, { readOnly: true });
const hasTranscriptEvents = db
.prepare("SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'transcript_events'")
.get();
if (!hasTranscriptEvents) {
return counts;
}
const rows = db.prepare("SELECT event_json FROM transcript_events ORDER BY session_id, seq");
for (const row of rows.iterate() as Iterable<{ event_json?: unknown }>) {
if (typeof row.event_json !== "string") {
continue;
}
const byteCount = Buffer.byteLength(row.event_json, "utf8");
assertWithinLimit({
byteCount,
filePath: sqlitePath,
label: "per-file",
limit: params.limits.fileMaxBytes,
});
totalBytes += byteCount;
assertWithinLimit({
byteCount: totalBytes,
label: "total",
limit: params.limits.totalMaxBytes,
});
const scanText = sessionLogScanText(row.event_json);
if (scanText === null) {
continue;
}
for (const [key, needle] of Object.entries(params.needles)) {
counts[key] += countOccurrences(scanText, needle);
}
}
return counts;
} catch (error) {
if (error && typeof error === "object" && "code" in error) {
throw error;
}
return counts;
} finally {
db?.close();
}
}