mirror of
https://github.com/openclaw/openclaw.git
synced 2026-07-28 17:31:19 +00:00
refactor(sessions): split SQLite session accessor (#107653)
* refactor(sessions): split SQLite session accessor * fix(sessions): keep extracted helpers module-private
This commit is contained in:
committed by
GitHub
parent
04e523a73b
commit
ecb6eac684
186
src/config/sessions/session-accessor.sqlite-archive.ts
Normal file
186
src/config/sessions/session-accessor.sqlite-archive.ts
Normal file
@@ -0,0 +1,186 @@
|
||||
import { randomUUID } from "node:crypto";
|
||||
import fs from "node:fs";
|
||||
import path from "node:path";
|
||||
import {
|
||||
encodeSessionArchiveContent,
|
||||
readSessionArchiveContentSync,
|
||||
SESSION_ARCHIVE_ZSTD_SUFFIX,
|
||||
} from "./archive-compression.js";
|
||||
import { formatSessionArchiveTimestamp } from "./artifacts.js";
|
||||
import type { SessionLifecycleArchivedTranscript } from "./session-accessor.sqlite-contract.js";
|
||||
|
||||
export type SqliteSessionStateDeletePlan = {
|
||||
archiveDirectory: string;
|
||||
archiveTranscript: boolean;
|
||||
content: string;
|
||||
hadTranscriptState: boolean;
|
||||
reason: "deleted" | "reset";
|
||||
sessionId: string;
|
||||
};
|
||||
|
||||
export type MaterializedSqliteSessionStateDeletePlan = SqliteSessionStateDeletePlan & {
|
||||
archivedTranscript: SessionLifecycleArchivedTranscript | null;
|
||||
};
|
||||
|
||||
function resolveSqliteTranscriptArchivePath(params: {
|
||||
archiveDirectory: string;
|
||||
reason: "deleted" | "reset";
|
||||
sessionId: string;
|
||||
nowMs?: number;
|
||||
}): string {
|
||||
const archiveDirectory = path.resolve(params.archiveDirectory);
|
||||
const archivePath = path.resolve(
|
||||
archiveDirectory,
|
||||
`${params.sessionId}.jsonl.${params.reason}.${formatSessionArchiveTimestamp(params.nowMs)}`,
|
||||
);
|
||||
if (path.dirname(archivePath) !== archiveDirectory) {
|
||||
throw new Error(`Cannot archive SQLite transcript outside ${archiveDirectory}`);
|
||||
}
|
||||
return archivePath;
|
||||
}
|
||||
|
||||
function findMatchingSqliteTranscriptArchive(params: {
|
||||
archiveDirectory: string;
|
||||
content: string;
|
||||
reason: "deleted" | "reset";
|
||||
sessionId: string;
|
||||
}): string | null {
|
||||
let entries: string[];
|
||||
try {
|
||||
entries = fs.readdirSync(params.archiveDirectory);
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
const prefix = `${params.sessionId}.jsonl.${params.reason}.`;
|
||||
for (const entry of entries) {
|
||||
if (!entry.startsWith(prefix)) {
|
||||
continue;
|
||||
}
|
||||
const archivePath = path.join(params.archiveDirectory, entry);
|
||||
const compressed = entry.endsWith(SESSION_ARCHIVE_ZSTD_SUFFIX);
|
||||
try {
|
||||
const stat = fs.statSync(archivePath);
|
||||
if (!stat.isFile()) {
|
||||
continue;
|
||||
}
|
||||
if (!compressed && stat.size !== Buffer.byteLength(params.content, "utf8")) {
|
||||
continue;
|
||||
}
|
||||
if (readSessionArchiveContentSync(archivePath) === params.content) {
|
||||
return archivePath;
|
||||
}
|
||||
} catch {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function writeSqliteTranscriptArchive(params: {
|
||||
archiveDirectory: string;
|
||||
content: string;
|
||||
reason: "deleted" | "reset";
|
||||
sessionId: string;
|
||||
}): string {
|
||||
fs.mkdirSync(params.archiveDirectory, { recursive: true });
|
||||
const existing = findMatchingSqliteTranscriptArchive(params);
|
||||
if (existing) {
|
||||
return existing;
|
||||
}
|
||||
// Archives are the long-lived cold tier; compress when the runtime can so
|
||||
// keep-forever retention stays cheap. Plain JSONL is the Bun/older fallback.
|
||||
const encoded = encodeSessionArchiveContent(params.content);
|
||||
for (let attempt = 0; attempt < 10; attempt += 1) {
|
||||
const archivePath = `${resolveSqliteTranscriptArchivePath({
|
||||
archiveDirectory: params.archiveDirectory,
|
||||
reason: params.reason,
|
||||
sessionId: params.sessionId,
|
||||
nowMs: Date.now() + attempt,
|
||||
})}${encoded.suffix}`;
|
||||
if (fs.existsSync(archivePath)) {
|
||||
continue;
|
||||
}
|
||||
const tempPath = `${archivePath}.${randomUUID()}.tmp`;
|
||||
try {
|
||||
fs.writeFileSync(tempPath, encoded.bytes, { flag: "wx", mode: 0o600 });
|
||||
fsyncRegularFile(tempPath);
|
||||
fs.renameSync(tempPath, archivePath);
|
||||
fsyncDirectory(params.archiveDirectory);
|
||||
return archivePath;
|
||||
} catch (error) {
|
||||
fs.rmSync(tempPath, { force: true });
|
||||
if ((error as { code?: unknown })?.code === "EEXIST") {
|
||||
continue;
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
throw new Error(`Could not create SQLite transcript archive for ${params.sessionId}`);
|
||||
}
|
||||
|
||||
function fsyncRegularFile(filePath: string): void {
|
||||
const fd = fs.openSync(filePath, "r");
|
||||
try {
|
||||
fs.fsyncSync(fd);
|
||||
} finally {
|
||||
fs.closeSync(fd);
|
||||
}
|
||||
}
|
||||
|
||||
function fsyncDirectory(dirPath: string): void {
|
||||
let fd: number | undefined;
|
||||
try {
|
||||
fd = fs.openSync(dirPath, "r");
|
||||
fs.fsyncSync(fd);
|
||||
} catch {
|
||||
// Directory fsync is not available on every supported platform/filesystem.
|
||||
} finally {
|
||||
if (fd !== undefined) {
|
||||
fs.closeSync(fd);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Runs duplicate probing, archive write, rename, and fsync outside SQLite
|
||||
// write transactions; deletion later consumes this durable proof.
|
||||
export function materializeSqliteSessionStateDeletePlans(
|
||||
plans: readonly SqliteSessionStateDeletePlan[],
|
||||
): MaterializedSqliteSessionStateDeletePlan[] {
|
||||
return dedupeSqliteSessionStateDeletePlans(plans).map((plan) => {
|
||||
const archivedTranscript =
|
||||
plan.archiveTranscript && plan.content.length > 0
|
||||
? {
|
||||
archivedPath: writeSqliteTranscriptArchive({
|
||||
archiveDirectory: plan.archiveDirectory,
|
||||
content: plan.content,
|
||||
reason: plan.reason,
|
||||
sessionId: plan.sessionId,
|
||||
}),
|
||||
sourcePath: path.join(plan.archiveDirectory, `${plan.sessionId}.jsonl`),
|
||||
}
|
||||
: null;
|
||||
return Object.assign({}, plan, { archivedTranscript });
|
||||
});
|
||||
}
|
||||
|
||||
// Multiple removed entries can point at one transcript session. If any owner
|
||||
// asked to keep an archive, the shared row gets exported once.
|
||||
function dedupeSqliteSessionStateDeletePlans(
|
||||
plans: readonly SqliteSessionStateDeletePlan[],
|
||||
): SqliteSessionStateDeletePlan[] {
|
||||
const deduped = new Map<string, SqliteSessionStateDeletePlan>();
|
||||
for (const plan of plans) {
|
||||
const existing = deduped.get(plan.sessionId);
|
||||
if (!existing) {
|
||||
deduped.set(plan.sessionId, plan);
|
||||
continue;
|
||||
}
|
||||
if (existing.content !== plan.content || existing.reason !== plan.reason) {
|
||||
throw new Error(`Conflicting SQLite transcript archive plans for ${plan.sessionId}`);
|
||||
}
|
||||
if (!existing.archiveTranscript && plan.archiveTranscript) {
|
||||
deduped.set(plan.sessionId, { ...existing, archiveTranscript: true });
|
||||
}
|
||||
}
|
||||
return [...deduped.values()];
|
||||
}
|
||||
108
src/config/sessions/session-accessor.sqlite-disk-budget.ts
Normal file
108
src/config/sessions/session-accessor.sqlite-disk-budget.ts
Normal file
@@ -0,0 +1,108 @@
|
||||
import type { SessionDiskBudgetSweepResult } from "./disk-budget.js";
|
||||
import { shouldPreserveMaintenanceEntry } from "./store-maintenance.js";
|
||||
import type { SessionEntry } from "./types.js";
|
||||
|
||||
export type SqliteSessionRowBytes = {
|
||||
entryBytesByKey: Map<string, number>;
|
||||
trajectoryBytesBySessionId: Map<string, number>;
|
||||
transcriptBytesBySessionId: Map<string, number>;
|
||||
};
|
||||
|
||||
function getSessionStateBytes(rowBytes: SqliteSessionRowBytes, sessionId: string): number {
|
||||
return (
|
||||
(rowBytes.transcriptBytesBySessionId.get(sessionId) ?? 0) +
|
||||
(rowBytes.trajectoryBytesBySessionId.get(sessionId) ?? 0)
|
||||
);
|
||||
}
|
||||
|
||||
function getEntryUpdatedAt(entry?: SessionEntry): number {
|
||||
return entry?.updatedAt ?? Number.NEGATIVE_INFINITY;
|
||||
}
|
||||
|
||||
export function enforceSqliteSessionDiskBudget(params: {
|
||||
collectStateIds: (entry: SessionEntry) => readonly string[];
|
||||
maintenance: { maxDiskBytes: number | null; highWaterBytes: number | null };
|
||||
onRemoveEntry?: (removed: { key: string; entry: SessionEntry }) => void;
|
||||
preserveKeys?: ReadonlySet<string>;
|
||||
rowBytes: SqliteSessionRowBytes;
|
||||
store: Record<string, SessionEntry>;
|
||||
}): SessionDiskBudgetSweepResult | null {
|
||||
const { maxDiskBytes, highWaterBytes } = params.maintenance;
|
||||
if (maxDiskBytes == null || highWaterBytes == null) {
|
||||
return null;
|
||||
}
|
||||
|
||||
let totalBytes = 0;
|
||||
const entryBytesByKey = new Map<string, number>();
|
||||
const sessionIdsByKey = new Map<string, readonly string[]>();
|
||||
const sessionIdRefCounts = new Map<string, number>();
|
||||
// State rows can be shared through usage-family references. Count each
|
||||
// session once; release its bytes only after the last entry is removed.
|
||||
for (const [key, entry] of Object.entries(params.store)) {
|
||||
const entryBytes = params.rowBytes.entryBytesByKey.get(key) ?? 0;
|
||||
const sessionIds = params.collectStateIds(entry);
|
||||
entryBytesByKey.set(key, entryBytes);
|
||||
sessionIdsByKey.set(key, sessionIds);
|
||||
totalBytes += entryBytes;
|
||||
for (const sessionId of sessionIds) {
|
||||
sessionIdRefCounts.set(sessionId, (sessionIdRefCounts.get(sessionId) ?? 0) + 1);
|
||||
}
|
||||
}
|
||||
for (const sessionId of sessionIdRefCounts.keys()) {
|
||||
totalBytes += getSessionStateBytes(params.rowBytes, sessionId);
|
||||
}
|
||||
|
||||
const totalBytesBefore = totalBytes;
|
||||
if (totalBytes <= maxDiskBytes) {
|
||||
return {
|
||||
totalBytesBefore,
|
||||
totalBytesAfter: totalBytes,
|
||||
removedFiles: 0,
|
||||
removedEntries: 0,
|
||||
freedBytes: 0,
|
||||
maxBytes: maxDiskBytes,
|
||||
highWaterBytes,
|
||||
overBudget: false,
|
||||
};
|
||||
}
|
||||
|
||||
let removedEntries = 0;
|
||||
const keys = Object.keys(params.store).toSorted(
|
||||
(left, right) => getEntryUpdatedAt(params.store[left]) - getEntryUpdatedAt(params.store[right]),
|
||||
);
|
||||
for (const key of keys) {
|
||||
if (totalBytes <= highWaterBytes) {
|
||||
break;
|
||||
}
|
||||
const entry = params.store[key];
|
||||
if (
|
||||
!entry ||
|
||||
shouldPreserveMaintenanceEntry({ key, entry, preserveKeys: params.preserveKeys })
|
||||
) {
|
||||
continue;
|
||||
}
|
||||
params.onRemoveEntry?.({ key, entry });
|
||||
delete params.store[key];
|
||||
removedEntries += 1;
|
||||
totalBytes -= entryBytesByKey.get(key) ?? 0;
|
||||
for (const sessionId of sessionIdsByKey.get(key) ?? []) {
|
||||
const nextRefCount = (sessionIdRefCounts.get(sessionId) ?? 0) - 1;
|
||||
if (nextRefCount > 0) {
|
||||
sessionIdRefCounts.set(sessionId, nextRefCount);
|
||||
continue;
|
||||
}
|
||||
sessionIdRefCounts.delete(sessionId);
|
||||
totalBytes -= getSessionStateBytes(params.rowBytes, sessionId);
|
||||
}
|
||||
}
|
||||
return {
|
||||
totalBytesBefore,
|
||||
totalBytesAfter: totalBytes,
|
||||
removedFiles: 0,
|
||||
removedEntries,
|
||||
freedBytes: Math.max(0, totalBytesBefore - totalBytes),
|
||||
maxBytes: maxDiskBytes,
|
||||
highWaterBytes,
|
||||
overBudget: true,
|
||||
};
|
||||
}
|
||||
160
src/config/sessions/session-accessor.sqlite-identity.ts
Normal file
160
src/config/sessions/session-accessor.sqlite-identity.ts
Normal file
@@ -0,0 +1,160 @@
|
||||
import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce";
|
||||
import { emitSessionIdentityMutation } from "../../sessions/session-lifecycle-events.js";
|
||||
import type { SessionEntry } from "./types.js";
|
||||
|
||||
type SqliteSessionEntryRemovalIdentity = {
|
||||
expectedEntry?: SessionEntry;
|
||||
sessionKey: string;
|
||||
};
|
||||
|
||||
type SqliteProjectedLifecycleIdentityMutation = {
|
||||
removals: Array<{
|
||||
expectedEntry: SessionEntry;
|
||||
sessionKey: string;
|
||||
}>;
|
||||
upsertedEntries: Array<{
|
||||
entry: SessionEntry;
|
||||
expectedEntry: SessionEntry | undefined;
|
||||
sessionKey: string;
|
||||
}>;
|
||||
};
|
||||
|
||||
function toSessionIdentityTarget(entry: SessionEntry | undefined, sessionKeys: readonly string[]) {
|
||||
const sessionId = normalizeOptionalString(entry?.sessionId);
|
||||
return { ...(sessionId ? { sessionId } : {}), sessionKeys };
|
||||
}
|
||||
|
||||
function emitCommittedSessionEntryRemoval(sessionKey: string, entry?: SessionEntry): void {
|
||||
emitSessionIdentityMutation({
|
||||
kind: "delete",
|
||||
previous: toSessionIdentityTarget(entry, [sessionKey]),
|
||||
});
|
||||
}
|
||||
|
||||
export function emitCommittedSessionEntryRemovals(
|
||||
removals: readonly SqliteSessionEntryRemovalIdentity[],
|
||||
): void {
|
||||
const emittedKeys = new Set<string>();
|
||||
for (const removal of removals) {
|
||||
if (emittedKeys.has(removal.sessionKey)) {
|
||||
continue;
|
||||
}
|
||||
emittedKeys.add(removal.sessionKey);
|
||||
emitCommittedSessionEntryRemoval(removal.sessionKey, removal.expectedEntry);
|
||||
}
|
||||
}
|
||||
|
||||
export function emitCommittedSessionEntryChange(params: {
|
||||
currentKey: string;
|
||||
currentEntry: SessionEntry;
|
||||
previousKey: string;
|
||||
previousEntry: SessionEntry;
|
||||
}): void {
|
||||
const previous = toSessionIdentityTarget(params.previousEntry, [params.previousKey]);
|
||||
const current = toSessionIdentityTarget(params.currentEntry, [params.currentKey]);
|
||||
const moved = params.previousKey !== params.currentKey;
|
||||
if (!moved && previous.sessionId === current.sessionId) {
|
||||
return;
|
||||
}
|
||||
emitSessionIdentityMutation({
|
||||
kind: moved ? "move" : "replace",
|
||||
previous,
|
||||
current,
|
||||
});
|
||||
}
|
||||
|
||||
export function emitCommittedSessionIdentityDiff(
|
||||
previous: ReadonlyMap<string, SessionEntry>,
|
||||
current: ReadonlyMap<string, SessionEntry>,
|
||||
): void {
|
||||
const currentKeysBySessionId = new Map<string, string[]>();
|
||||
for (const [sessionKey, entry] of current) {
|
||||
const sessionId = normalizeOptionalString(entry.sessionId);
|
||||
if (sessionId) {
|
||||
currentKeysBySessionId.set(sessionId, [
|
||||
...(currentKeysBySessionId.get(sessionId) ?? []),
|
||||
sessionKey,
|
||||
]);
|
||||
}
|
||||
}
|
||||
|
||||
const movedKeysByCurrentKey = new Map<string, string[]>();
|
||||
const handledPreviousKeys = new Set<string>();
|
||||
const handledCurrentKeys = new Set<string>();
|
||||
for (const [sessionKey, entry] of previous) {
|
||||
if (current.has(sessionKey)) {
|
||||
continue;
|
||||
}
|
||||
const sessionId = normalizeOptionalString(entry.sessionId);
|
||||
const currentKeys = sessionId ? currentKeysBySessionId.get(sessionId) : undefined;
|
||||
if (currentKeys?.length !== 1) {
|
||||
continue;
|
||||
}
|
||||
const [currentKey] = currentKeys;
|
||||
if (!currentKey) {
|
||||
continue;
|
||||
}
|
||||
movedKeysByCurrentKey.set(currentKey, [
|
||||
...(movedKeysByCurrentKey.get(currentKey) ?? []),
|
||||
sessionKey,
|
||||
]);
|
||||
handledPreviousKeys.add(sessionKey);
|
||||
handledCurrentKeys.add(currentKey);
|
||||
}
|
||||
for (const [currentKey, previousKeys] of movedKeysByCurrentKey) {
|
||||
const currentEntry = current.get(currentKey);
|
||||
if (currentEntry) {
|
||||
emitSessionIdentityMutation({
|
||||
kind: "move",
|
||||
previous: toSessionIdentityTarget(currentEntry, previousKeys),
|
||||
current: toSessionIdentityTarget(currentEntry, [currentKey]),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
for (const [sessionKey, previousEntry] of previous) {
|
||||
const currentEntry = current.get(sessionKey);
|
||||
if (currentEntry) {
|
||||
handledCurrentKeys.add(sessionKey);
|
||||
emitCommittedSessionEntryChange({
|
||||
currentEntry,
|
||||
currentKey: sessionKey,
|
||||
previousEntry,
|
||||
previousKey: sessionKey,
|
||||
});
|
||||
} else if (!handledPreviousKeys.has(sessionKey)) {
|
||||
emitCommittedSessionEntryRemoval(sessionKey, previousEntry);
|
||||
}
|
||||
}
|
||||
|
||||
for (const [sessionKey, currentEntry] of current) {
|
||||
if (handledCurrentKeys.has(sessionKey)) {
|
||||
continue;
|
||||
}
|
||||
emitSessionIdentityMutation({
|
||||
kind: "create",
|
||||
previous: { sessionKeys: [] },
|
||||
current: toSessionIdentityTarget(currentEntry, [sessionKey]),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
export function emitCommittedLifecycleIdentityMutations(params: {
|
||||
projected: SqliteProjectedLifecycleIdentityMutation;
|
||||
removedSessionKeys: readonly string[];
|
||||
}): void {
|
||||
const removedKeys = new Set(params.removedSessionKeys);
|
||||
const previous = new Map(
|
||||
params.projected.removals
|
||||
.filter((removal) => removedKeys.has(removal.sessionKey))
|
||||
.map((removal) => [removal.sessionKey, removal.expectedEntry]),
|
||||
);
|
||||
const current = new Map<string, SessionEntry>();
|
||||
for (const upsert of params.projected.upsertedEntries) {
|
||||
if (!current.has(upsert.sessionKey) && upsert.expectedEntry) {
|
||||
previous.set(upsert.sessionKey, upsert.expectedEntry);
|
||||
}
|
||||
current.set(upsert.sessionKey, upsert.entry);
|
||||
}
|
||||
emitCommittedSessionIdentityDiff(previous, current);
|
||||
}
|
||||
341
src/config/sessions/session-accessor.sqlite-parent-fork.ts
Normal file
341
src/config/sessions/session-accessor.sqlite-parent-fork.ts
Normal file
@@ -0,0 +1,341 @@
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { derivePromptTokens, normalizeUsage } from "../../agents/usage.js";
|
||||
import type {
|
||||
SessionParentForkDecision,
|
||||
TranscriptEvent,
|
||||
} from "./session-accessor.sqlite-contract.js";
|
||||
import { createSessionTranscriptHeader } from "./transcript-header.js";
|
||||
import {
|
||||
isSessionTranscriptLeafControl,
|
||||
mergeSessionTranscriptVisiblePathWithOpaqueAppendPath,
|
||||
scanSessionTranscriptTree,
|
||||
selectSessionTranscriptTreePathNodes,
|
||||
} from "./transcript-tree.js";
|
||||
import type { SessionEntry } from "./types.js";
|
||||
import { resolveFreshSessionTotalTokens, resolveSessionTotalTokens } from "./types.js";
|
||||
|
||||
export type SqliteParentForkSourceTranscript = {
|
||||
appendMode?: "side";
|
||||
appendParentId: string | null;
|
||||
branchEntries: TranscriptEvent[];
|
||||
cwd?: string;
|
||||
labelsToWrite: Array<{ targetId: string; label: string; timestamp: string }>;
|
||||
leafId: string | null;
|
||||
preserveLeafControl: boolean;
|
||||
};
|
||||
|
||||
type SqliteTranscriptParentTokenEstimate = {
|
||||
kind: "exact-context" | "legacy-or-bytes";
|
||||
tokens: number;
|
||||
};
|
||||
|
||||
const DEFAULT_PARENT_FORK_MAX_TOKENS = 100_000;
|
||||
|
||||
function formatParentForkTooLargeMessage(params: {
|
||||
parentTokens: number;
|
||||
maxTokens: number;
|
||||
}): string {
|
||||
return (
|
||||
`Parent context is too large to fork (${params.parentTokens}/${params.maxTokens} tokens); ` +
|
||||
"starting with isolated context instead."
|
||||
);
|
||||
}
|
||||
|
||||
export function resolveSqliteParentForkDecision(
|
||||
parentEntry: SessionEntry,
|
||||
transcriptEstimate?: SqliteTranscriptParentTokenEstimate,
|
||||
): SessionParentForkDecision {
|
||||
const maxTokens = DEFAULT_PARENT_FORK_MAX_TOKENS;
|
||||
const parentTokens =
|
||||
resolveFreshSessionTotalTokens(parentEntry) ??
|
||||
(transcriptEstimate?.kind === "exact-context"
|
||||
? transcriptEstimate.tokens
|
||||
: maxPositiveTokenCount(transcriptEstimate?.tokens, resolveSessionTotalTokens(parentEntry)));
|
||||
if (typeof parentTokens === "number" && parentTokens > maxTokens) {
|
||||
return {
|
||||
status: "skip",
|
||||
reason: "parent-too-large",
|
||||
maxTokens,
|
||||
parentTokens,
|
||||
message: formatParentForkTooLargeMessage({ parentTokens, maxTokens }),
|
||||
};
|
||||
}
|
||||
return {
|
||||
status: "fork",
|
||||
maxTokens,
|
||||
...(typeof parentTokens === "number" ? { parentTokens } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
export function estimateSqliteTranscriptPromptTokens(
|
||||
events: readonly TranscriptEvent[],
|
||||
): SqliteTranscriptParentTokenEstimate | undefined {
|
||||
let byteEstimate = 0;
|
||||
let latestUsageEstimate: number | undefined;
|
||||
let latestUsageEstimateIsExactContext = false;
|
||||
let trailingBytes = 0;
|
||||
for (const event of selectParentForkTokenEstimateEvents(events)) {
|
||||
const serializedBytes = Buffer.byteLength(JSON.stringify(event)) + 1;
|
||||
byteEstimate += serializedBytes;
|
||||
if (!isRecord(event)) {
|
||||
if (latestUsageEstimate !== undefined) {
|
||||
trailingBytes += serializedBytes;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
const message = isRecord(event.message) ? event.message : undefined;
|
||||
const usageRaw = isRecord(message?.usage)
|
||||
? message.usage
|
||||
: isRecord(event.usage)
|
||||
? event.usage
|
||||
: undefined;
|
||||
if (!usageRaw) {
|
||||
if (latestUsageEstimate !== undefined) {
|
||||
trailingBytes += serializedBytes;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
const contextUsage = readTranscriptContextUsage(usageRaw);
|
||||
if (contextUsage?.state === "unavailable") {
|
||||
latestUsageEstimate = undefined;
|
||||
latestUsageEstimateIsExactContext = false;
|
||||
trailingBytes = 0;
|
||||
continue;
|
||||
}
|
||||
if (contextUsage?.state === "available") {
|
||||
latestUsageEstimate = normalizePositiveTokenCount(contextUsage.totalTokens);
|
||||
latestUsageEstimateIsExactContext = true;
|
||||
trailingBytes = 0;
|
||||
continue;
|
||||
}
|
||||
const usage = normalizeUsage(usageRaw);
|
||||
const promptTokens = normalizePositiveTokenCount(
|
||||
derivePromptTokens({
|
||||
input: usage?.input,
|
||||
cacheRead: usage?.cacheRead,
|
||||
cacheWrite: usage?.cacheWrite,
|
||||
}),
|
||||
);
|
||||
const outputTokens = normalizePositiveTokenCount(usage?.output) ?? 0;
|
||||
const totalTokens =
|
||||
promptTokens === undefined
|
||||
? undefined
|
||||
: normalizePositiveTokenCount(promptTokens + outputTokens);
|
||||
if (typeof totalTokens === "number") {
|
||||
latestUsageEstimate = totalTokens;
|
||||
latestUsageEstimateIsExactContext = false;
|
||||
trailingBytes = 0;
|
||||
}
|
||||
}
|
||||
if (latestUsageEstimate !== undefined) {
|
||||
const tokens = normalizePositiveTokenCount(latestUsageEstimate + Math.ceil(trailingBytes / 4));
|
||||
return tokens === undefined
|
||||
? undefined
|
||||
: {
|
||||
kind: latestUsageEstimateIsExactContext ? "exact-context" : "legacy-or-bytes",
|
||||
tokens,
|
||||
};
|
||||
}
|
||||
const tokens = normalizePositiveTokenCount(Math.ceil(byteEstimate / 4));
|
||||
return tokens === undefined ? undefined : { kind: "legacy-or-bytes", tokens };
|
||||
}
|
||||
|
||||
function selectParentForkTokenEstimateEvents(
|
||||
events: readonly TranscriptEvent[],
|
||||
): TranscriptEvent[] {
|
||||
const entries = events.filter((entry) => !(isRecord(entry) && entry.type === "session"));
|
||||
const tree = scanSessionTranscriptTree(entries);
|
||||
const visiblePath = selectSessionTranscriptTreePathNodes(tree, tree.leafId);
|
||||
const appendPath = selectSessionTranscriptTreePathNodes(tree, tree.appendParentId);
|
||||
return mergeSessionTranscriptVisiblePathWithOpaqueAppendPath({
|
||||
visiblePath,
|
||||
appendPath,
|
||||
appendParentId: tree.appendParentId,
|
||||
}).nodes.flatMap((node) => node.entry);
|
||||
}
|
||||
|
||||
function isRecord(value: unknown): value is Record<string, unknown> {
|
||||
return typeof value === "object" && value !== null && !Array.isArray(value);
|
||||
}
|
||||
|
||||
function normalizePositiveTokenCount(value: unknown): number | undefined {
|
||||
return typeof value === "number" && Number.isFinite(value) && value > 0
|
||||
? Math.floor(value)
|
||||
: undefined;
|
||||
}
|
||||
|
||||
function maxPositiveTokenCount(...values: Array<number | undefined>): number | undefined {
|
||||
let max: number | undefined;
|
||||
for (const value of values) {
|
||||
const normalized = normalizePositiveTokenCount(value);
|
||||
if (normalized !== undefined && (max === undefined || normalized > max)) {
|
||||
max = normalized;
|
||||
}
|
||||
}
|
||||
return max;
|
||||
}
|
||||
|
||||
function readTranscriptContextUsage(
|
||||
usageRaw: Record<string, unknown>,
|
||||
): { state: "available"; totalTokens: number } | { state: "unavailable" } | undefined {
|
||||
const contextUsage = usageRaw.contextUsage;
|
||||
if (!isRecord(contextUsage)) {
|
||||
return undefined;
|
||||
}
|
||||
if (contextUsage.state === "unavailable") {
|
||||
return { state: "unavailable" };
|
||||
}
|
||||
if (contextUsage.state !== "available") {
|
||||
return undefined;
|
||||
}
|
||||
const totalTokens = normalizePositiveTokenCount(contextUsage.totalTokens);
|
||||
return totalTokens === undefined ? undefined : { state: "available", totalTokens };
|
||||
}
|
||||
|
||||
export function resolveSqliteParentForkSourceTranscript(
|
||||
fileEntries: readonly TranscriptEvent[],
|
||||
): SqliteParentForkSourceTranscript | null {
|
||||
if (fileEntries.length === 0) {
|
||||
return null;
|
||||
}
|
||||
const header = fileEntries.find(
|
||||
(entry): entry is Record<string, unknown> => isRecord(entry) && entry.type === "session",
|
||||
);
|
||||
const entries = fileEntries.filter((entry) => !(isRecord(entry) && entry.type === "session"));
|
||||
const tree = scanSessionTranscriptTree(entries);
|
||||
const visiblePath = selectSessionTranscriptTreePathNodes(tree, tree.leafId);
|
||||
const appendPath = selectSessionTranscriptTreePathNodes(tree, tree.appendParentId);
|
||||
const mergedPath = mergeSessionTranscriptVisiblePathWithOpaqueAppendPath({
|
||||
visiblePath,
|
||||
appendPath,
|
||||
appendParentId: tree.appendParentId,
|
||||
});
|
||||
const branchEntries = mergedPath.nodes.flatMap((node) => {
|
||||
if (!isRecord(node.entry)) {
|
||||
return [];
|
||||
}
|
||||
const parentId = node.selectedParentId;
|
||||
return [node.entry.parentId === parentId ? node.entry : { ...node.entry, parentId }];
|
||||
});
|
||||
const pathEntryIds = new Set(
|
||||
branchEntries.flatMap((entry) =>
|
||||
isRecord(entry) && typeof entry.id === "string" ? [entry.id] : [],
|
||||
),
|
||||
);
|
||||
const lastLeafUpdateNode = tree.nodes.findLast((node) => node.leafId !== undefined);
|
||||
return {
|
||||
appendParentId: mergedPath.appendParentId,
|
||||
...(lastLeafUpdateNode?.appendMode ? { appendMode: lastLeafUpdateNode.appendMode } : {}),
|
||||
branchEntries,
|
||||
cwd: typeof header?.cwd === "string" ? header.cwd : undefined,
|
||||
labelsToWrite: collectBranchLabels({ allEntries: entries, pathEntryIds }),
|
||||
leafId: tree.leafId,
|
||||
preserveLeafControl: isSessionTranscriptLeafControl(lastLeafUpdateNode?.entry),
|
||||
};
|
||||
}
|
||||
|
||||
function collectBranchLabels(params: {
|
||||
allEntries: readonly TranscriptEvent[];
|
||||
pathEntryIds: Set<string>;
|
||||
}): Array<{ targetId: string; label: string; timestamp: string }> {
|
||||
return params.allEntries.flatMap((entry) =>
|
||||
isRecord(entry) &&
|
||||
entry.type === "label" &&
|
||||
typeof entry.label === "string" &&
|
||||
typeof entry.targetId === "string" &&
|
||||
typeof entry.id === "string" &&
|
||||
!params.pathEntryIds.has(entry.id) &&
|
||||
params.pathEntryIds.has(entry.targetId) &&
|
||||
typeof entry.timestamp === "string"
|
||||
? [{ targetId: entry.targetId, label: entry.label, timestamp: entry.timestamp }]
|
||||
: [],
|
||||
);
|
||||
}
|
||||
|
||||
function generateEntryId(existingIds: Set<string>): string {
|
||||
for (let attempt = 0; attempt < 100; attempt += 1) {
|
||||
const id = randomUUID().slice(0, 8);
|
||||
if (!existingIds.has(id)) {
|
||||
existingIds.add(id);
|
||||
return id;
|
||||
}
|
||||
}
|
||||
const id = randomUUID();
|
||||
existingIds.add(id);
|
||||
return id;
|
||||
}
|
||||
|
||||
function buildLabelEntries(params: {
|
||||
labelsToWrite: Array<{ targetId: string; label: string; timestamp: string }>;
|
||||
pathEntryIds: Set<string>;
|
||||
lastEntryId: string | null;
|
||||
}): TranscriptEvent[] {
|
||||
let parentId = params.lastEntryId;
|
||||
return params.labelsToWrite.map(({ targetId, label, timestamp }) => {
|
||||
const entry = {
|
||||
type: "label",
|
||||
id: generateEntryId(params.pathEntryIds),
|
||||
parentId,
|
||||
timestamp,
|
||||
targetId,
|
||||
label,
|
||||
};
|
||||
parentId = entry.id;
|
||||
return entry;
|
||||
});
|
||||
}
|
||||
|
||||
function hasAssistantEntry(entries: readonly TranscriptEvent[]): boolean {
|
||||
return entries.some(
|
||||
(entry) =>
|
||||
isRecord(entry) &&
|
||||
entry.type === "message" &&
|
||||
isRecord(entry.message) &&
|
||||
entry.message.role === "assistant",
|
||||
);
|
||||
}
|
||||
|
||||
export function buildSqliteForkedChildTranscriptEvents(params: {
|
||||
parentSessionFile: string;
|
||||
source: SqliteParentForkSourceTranscript;
|
||||
targetSessionId: string;
|
||||
}): TranscriptEvent[] {
|
||||
const header = {
|
||||
...createSessionTranscriptHeader({ cwd: params.source.cwd, sessionId: params.targetSessionId }),
|
||||
parentSession: params.parentSessionFile,
|
||||
};
|
||||
if (!params.source.preserveLeafControl && !hasAssistantEntry(params.source.branchEntries)) {
|
||||
return [header];
|
||||
}
|
||||
|
||||
const pathEntryIds = new Set(
|
||||
params.source.branchEntries.flatMap((entry) =>
|
||||
isRecord(entry) && typeof entry.id === "string" ? [entry.id] : [],
|
||||
),
|
||||
);
|
||||
const lastPathEntry = params.source.branchEntries.at(-1);
|
||||
const lastPathEntryId =
|
||||
isRecord(lastPathEntry) && typeof lastPathEntry.id === "string" ? lastPathEntry.id : null;
|
||||
const labelEntries = buildLabelEntries({
|
||||
labelsToWrite: params.source.labelsToWrite,
|
||||
pathEntryIds,
|
||||
lastEntryId: lastPathEntryId,
|
||||
});
|
||||
const leafEntry = params.source.preserveLeafControl
|
||||
? {
|
||||
type: "leaf",
|
||||
id: generateEntryId(pathEntryIds),
|
||||
parentId: (labelEntries.at(-1) as { id?: string } | undefined)?.id ?? lastPathEntryId,
|
||||
timestamp: new Date().toISOString(),
|
||||
targetId: params.source.leafId,
|
||||
appendParentId: params.source.appendParentId,
|
||||
...(params.source.appendMode ? { appendMode: params.source.appendMode } : {}),
|
||||
}
|
||||
: null;
|
||||
return [
|
||||
header,
|
||||
...params.source.branchEntries,
|
||||
...labelEntries,
|
||||
...(leafEntry ? [leafEntry] : []),
|
||||
];
|
||||
}
|
||||
275
src/config/sessions/session-accessor.sqlite-read.ts
Normal file
275
src/config/sessions/session-accessor.sqlite-read.ts
Normal file
@@ -0,0 +1,275 @@
|
||||
import { sql } from "kysely";
|
||||
import {
|
||||
executeSqliteQuerySync,
|
||||
executeSqliteQueryTakeFirstSync,
|
||||
iterateSqliteQuerySync,
|
||||
} from "../../infra/kysely-sync.js";
|
||||
import { extractAssistantVisibleText } from "../../shared/chat-message-content.js";
|
||||
import { isTranscriptOnlyOpenClawAssistantModel } from "../../shared/transcript-only-openclaw-assistant.js";
|
||||
import {
|
||||
openOpenClawAgentDatabase,
|
||||
type OpenClawAgentDatabase,
|
||||
} from "../../state/openclaw-agent-db.js";
|
||||
import type {
|
||||
LatestTranscriptAssistantMessage,
|
||||
LatestTranscriptAssistantText,
|
||||
SessionTranscriptReadScope,
|
||||
SessionTranscriptStats,
|
||||
TranscriptEvent,
|
||||
} from "./session-accessor.sqlite-contract.js";
|
||||
import { normalizeSqliteNumber } from "./session-accessor.sqlite-normalize.js";
|
||||
import {
|
||||
getSessionKysely,
|
||||
resolveSqliteTranscriptReadScope,
|
||||
toDatabaseOptions,
|
||||
} from "./session-accessor.sqlite-scope.js";
|
||||
|
||||
export type SqliteTranscriptSnapshotRow = {
|
||||
eventJson: string;
|
||||
seq: number;
|
||||
};
|
||||
|
||||
/** Loads raw transcript events from the additive SQLite transcript store. */
|
||||
export async function loadSqliteTranscriptEvents(
|
||||
scope: SessionTranscriptReadScope,
|
||||
): Promise<TranscriptEvent[]> {
|
||||
return loadSqliteTranscriptEventsSync(scope);
|
||||
}
|
||||
|
||||
/** Loads raw transcript events synchronously from the additive SQLite transcript store. */
|
||||
export function loadSqliteTranscriptEventsSync(
|
||||
scope: SessionTranscriptReadScope,
|
||||
): TranscriptEvent[] {
|
||||
const resolved = resolveSqliteTranscriptReadScope(scope);
|
||||
const database = openOpenClawAgentDatabase(toDatabaseOptions(resolved));
|
||||
return loadSqliteTranscriptEventsFromDatabase(database, resolved.sessionId);
|
||||
}
|
||||
|
||||
export function loadSqliteTranscriptEventsFromDatabase(
|
||||
database: OpenClawAgentDatabase,
|
||||
sessionId: string,
|
||||
): TranscriptEvent[] {
|
||||
const db = getSessionKysely(database.db);
|
||||
const rows = executeSqliteQuerySync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("transcript_events")
|
||||
.select(["event_json"])
|
||||
.where("session_id", "=", sessionId)
|
||||
.orderBy("seq", "asc"),
|
||||
).rows;
|
||||
return rows.map((row) => JSON.parse(row.event_json) as TranscriptEvent);
|
||||
}
|
||||
|
||||
export function readSqliteTranscriptSnapshot(
|
||||
database: OpenClawAgentDatabase,
|
||||
sessionId: string,
|
||||
): { events: TranscriptEvent[]; rows: SqliteTranscriptSnapshotRow[] } {
|
||||
const db = getSessionKysely(database.db);
|
||||
const rows = executeSqliteQuerySync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("transcript_events")
|
||||
.select(["event_json", "seq"])
|
||||
.where("session_id", "=", sessionId)
|
||||
.orderBy("seq", "asc"),
|
||||
).rows;
|
||||
return {
|
||||
events: rows.map((row) => JSON.parse(row.event_json) as TranscriptEvent),
|
||||
rows: rows.map((row) => ({
|
||||
eventJson: row.event_json,
|
||||
seq: normalizeSqliteNumber(row.seq),
|
||||
})),
|
||||
};
|
||||
}
|
||||
|
||||
function sqliteTranscriptJsonlByteSize() {
|
||||
return /* kysely-allow-raw: JSONL size includes event bytes plus newline separators. */ sql<number>`COALESCE(SUM(LENGTH(CAST(event_json AS BLOB))), 0)
|
||||
+ CASE WHEN COUNT(*) > 0 THEN COUNT(*) - 1 ELSE 0 END`.as("size_bytes");
|
||||
}
|
||||
|
||||
/** Reads transcript freshness and byte size without materializing event rows. */
|
||||
export function readSqliteTranscriptStatsSync(
|
||||
scope: SessionTranscriptReadScope,
|
||||
): SessionTranscriptStats {
|
||||
const resolved = resolveSqliteTranscriptReadScope(scope);
|
||||
const database = openOpenClawAgentDatabase(toDatabaseOptions(resolved));
|
||||
const db = getSessionKysely(database.db);
|
||||
const row = executeSqliteQueryTakeFirstSync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("transcript_events")
|
||||
.select((eb) => [
|
||||
eb.fn.count<number>("seq").as("event_count"),
|
||||
eb.fn.max<number>("seq").as("max_seq"),
|
||||
sqliteTranscriptJsonlByteSize(),
|
||||
])
|
||||
.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,
|
||||
};
|
||||
}
|
||||
|
||||
export function readTranscriptEventJsonSetInTransaction(
|
||||
database: OpenClawAgentDatabase,
|
||||
sessionId: string,
|
||||
): Set<string> {
|
||||
const db = getSessionKysely(database.db);
|
||||
const rows = executeSqliteQuerySync(
|
||||
database.db,
|
||||
db.selectFrom("transcript_events").select("event_json").where("session_id", "=", sessionId),
|
||||
).rows;
|
||||
return new Set(rows.map((row) => row.event_json));
|
||||
}
|
||||
|
||||
/** Reads the latest visible assistant text from SQLite transcript rows in reverse order. */
|
||||
export function loadLatestSqliteAssistantText(
|
||||
scope: SessionTranscriptReadScope,
|
||||
options: { includeTranscriptOnlyOpenClawAssistant?: boolean } = {},
|
||||
): LatestTranscriptAssistantText | undefined {
|
||||
const resolved = resolveSqliteTranscriptReadScope(scope);
|
||||
const database = openOpenClawAgentDatabase(toDatabaseOptions(resolved));
|
||||
const db = getSessionKysely(database.db);
|
||||
const rows = iterateSqliteQuerySync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("transcript_events as te")
|
||||
.innerJoin("transcript_event_identities as ti", (join) =>
|
||||
join.onRef("ti.session_id", "=", "te.session_id").onRef("ti.seq", "=", "te.seq"),
|
||||
)
|
||||
.select("te.event_json as event_json")
|
||||
.where("te.session_id", "=", resolved.sessionId)
|
||||
.where("ti.event_type", "=", "message")
|
||||
.orderBy("ti.seq", "desc"),
|
||||
);
|
||||
for (const row of rows) {
|
||||
const latest = parseLatestAssistantMessageEvent(row.event_json, options);
|
||||
if (!latest) {
|
||||
continue;
|
||||
}
|
||||
const text = parseLatestAssistantText(latest);
|
||||
if (text) {
|
||||
return text;
|
||||
}
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
function parseLatestAssistantText(
|
||||
latest: LatestTranscriptAssistantMessage,
|
||||
): LatestTranscriptAssistantText | undefined {
|
||||
const message = latest.message as { timestamp?: unknown };
|
||||
const text = extractAssistantVisibleText(latest.message)?.trim();
|
||||
if (!text) {
|
||||
return undefined;
|
||||
}
|
||||
return {
|
||||
...(latest.id ? { id: latest.id } : {}),
|
||||
text,
|
||||
...(typeof message.timestamp === "number" && Number.isFinite(message.timestamp)
|
||||
? { timestamp: message.timestamp }
|
||||
: {}),
|
||||
};
|
||||
}
|
||||
|
||||
function parseLatestAssistantMessageEvent(
|
||||
raw: string,
|
||||
options: { includeTranscriptOnlyOpenClawAssistant?: boolean } = {},
|
||||
): LatestTranscriptAssistantMessage | undefined {
|
||||
let parsed: {
|
||||
id?: unknown;
|
||||
message?: { model?: unknown; provider?: unknown; role?: unknown; timestamp?: unknown };
|
||||
};
|
||||
try {
|
||||
parsed = JSON.parse(raw) as typeof parsed;
|
||||
} catch {
|
||||
return undefined;
|
||||
}
|
||||
const message = parsed.message;
|
||||
if (!message || message.role !== "assistant") {
|
||||
return undefined;
|
||||
}
|
||||
if (
|
||||
!options.includeTranscriptOnlyOpenClawAssistant &&
|
||||
isTranscriptOnlyOpenClawAssistantModel(message.provider, message.model)
|
||||
) {
|
||||
return undefined;
|
||||
}
|
||||
return {
|
||||
...(typeof parsed.id === "string" && parsed.id.trim() ? { id: parsed.id } : {}),
|
||||
message,
|
||||
};
|
||||
}
|
||||
|
||||
/** Finds the newest transcript record accepted by the matcher without parsing older rows. */
|
||||
export function findSqliteTranscriptEvent(
|
||||
scope: SessionTranscriptReadScope,
|
||||
match: (event: TranscriptEvent) => boolean,
|
||||
): { event: TranscriptEvent } | undefined {
|
||||
const resolved = resolveSqliteTranscriptReadScope(scope);
|
||||
const database = openOpenClawAgentDatabase(toDatabaseOptions(resolved));
|
||||
return findSqliteTranscriptEventInDatabase(database, resolved.sessionId, match);
|
||||
}
|
||||
|
||||
export function findSqliteTranscriptEventInDatabase(
|
||||
database: OpenClawAgentDatabase,
|
||||
sessionId: string,
|
||||
match: (event: TranscriptEvent) => boolean,
|
||||
): { event: TranscriptEvent } | undefined {
|
||||
const db = getSessionKysely(database.db);
|
||||
const rows = executeSqliteQuerySync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("transcript_events")
|
||||
.select(["event_json"])
|
||||
.where("session_id", "=", sessionId)
|
||||
.orderBy("seq", "desc"),
|
||||
).rows;
|
||||
for (const row of rows) {
|
||||
try {
|
||||
const event = JSON.parse(row.event_json) as TranscriptEvent;
|
||||
if (match(event)) {
|
||||
return { event };
|
||||
}
|
||||
} catch {
|
||||
// Malformed rows are skipped, matching transcript index tolerance.
|
||||
}
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
export function readTranscriptEventMessage(
|
||||
event: TranscriptEvent,
|
||||
): Record<string, unknown> | undefined {
|
||||
if (!event || typeof event !== "object" || Array.isArray(event)) {
|
||||
return undefined;
|
||||
}
|
||||
const message = (event as { message?: unknown }).message;
|
||||
return message && typeof message === "object" && !Array.isArray(message)
|
||||
? (message as Record<string, unknown>)
|
||||
: undefined;
|
||||
}
|
||||
|
||||
export function readTranscriptEventId(event: TranscriptEvent): string | undefined {
|
||||
if (!event || typeof event !== "object" || Array.isArray(event)) {
|
||||
return undefined;
|
||||
}
|
||||
const id = (event as { id?: unknown }).id;
|
||||
return typeof id === "string" && id.trim() ? id : undefined;
|
||||
}
|
||||
262
src/config/sessions/session-accessor.sqlite-scope.ts
Normal file
262
src/config/sessions/session-accessor.sqlite-scope.ts
Normal file
@@ -0,0 +1,262 @@
|
||||
import path from "node:path";
|
||||
import { getNodeSqliteKysely } from "../../infra/kysely-sync.js";
|
||||
import { getChildLogger } from "../../logging/logger.js";
|
||||
import {
|
||||
DEFAULT_AGENT_ID,
|
||||
normalizeAgentId,
|
||||
parseAgentSessionKey,
|
||||
resolveAgentIdFromSessionKey,
|
||||
} from "../../routing/session-key.js";
|
||||
import { runQueuedStoreWrite, type StoreWriterQueue } from "../../shared/store-writer-queue.js";
|
||||
import type { DB as OpenClawAgentKyselyDatabase } from "../../state/openclaw-agent-db.generated.js";
|
||||
import {
|
||||
resolveOpenClawAgentSqlitePath,
|
||||
type OpenClawAgentDatabaseOptions,
|
||||
} from "../../state/openclaw-agent-db.js";
|
||||
import type {
|
||||
SessionAccessScope,
|
||||
SessionTranscriptReadScope,
|
||||
SessionTranscriptWriteScope,
|
||||
} from "./session-accessor.sqlite-contract.js";
|
||||
import { resolveSqliteTargetFromSessionStorePath } from "./session-sqlite-target.js";
|
||||
import { formatSqliteSessionFileMarker } from "./sqlite-marker.js";
|
||||
import { normalizeStoreSessionKey } from "./store-entry.js";
|
||||
import type { SessionEntry } from "./types.js";
|
||||
|
||||
type SessionSqliteDatabase = Pick<
|
||||
OpenClawAgentKyselyDatabase,
|
||||
| "conversations"
|
||||
| "session_conversations"
|
||||
| "session_entries"
|
||||
| "session_routes"
|
||||
| "sessions"
|
||||
| "trajectory_runtime_events"
|
||||
| "transcript_event_identities"
|
||||
| "transcript_events"
|
||||
>;
|
||||
|
||||
export type ResolvedSqliteScope = {
|
||||
agentId: string;
|
||||
env?: NodeJS.ProcessEnv;
|
||||
path?: string;
|
||||
sessionKey: string;
|
||||
};
|
||||
|
||||
export type ResolvedSqliteReadScope = {
|
||||
agentId: string;
|
||||
env?: NodeJS.ProcessEnv;
|
||||
path?: string;
|
||||
sessionKey?: string;
|
||||
};
|
||||
|
||||
export type ResolvedTranscriptScope = ResolvedSqliteScope & {
|
||||
sessionId: string;
|
||||
};
|
||||
|
||||
type ResolvedTranscriptReadScope = ResolvedSqliteReadScope & {
|
||||
sessionId: string;
|
||||
};
|
||||
|
||||
const SQLITE_SESSION_SLOW_WRITE_MS = 1_000;
|
||||
const SQLITE_SESSION_WRITER_QUEUES = new Map<string, StoreWriterQueue>();
|
||||
|
||||
export function getSessionKysely(database: import("node:sqlite").DatabaseSync) {
|
||||
return getNodeSqliteKysely<SessionSqliteDatabase>(database);
|
||||
}
|
||||
|
||||
export async function runExclusiveSqliteSessionWrite<T>(
|
||||
scope: Pick<ResolvedSqliteReadScope, "agentId" | "env" | "path">,
|
||||
fn: () => Promise<T>,
|
||||
): Promise<T> {
|
||||
const databaseOptions = toDatabaseOptions(scope);
|
||||
const storePath = resolveOpenClawAgentSqlitePath(databaseOptions);
|
||||
const startedAt = Date.now();
|
||||
try {
|
||||
const result = await runQueuedStoreWrite({
|
||||
queues: SQLITE_SESSION_WRITER_QUEUES,
|
||||
storePath,
|
||||
label: "runExclusiveSqliteSessionWrite",
|
||||
fn,
|
||||
});
|
||||
const elapsedMs = Date.now() - startedAt;
|
||||
if (elapsedMs >= SQLITE_SESSION_SLOW_WRITE_MS) {
|
||||
getChildLogger({ subsystem: "session-sqlite" }).warn("slow SQLite session write", {
|
||||
agentId: scope.agentId,
|
||||
elapsedMs,
|
||||
storePath,
|
||||
});
|
||||
}
|
||||
return result;
|
||||
} catch (error) {
|
||||
getChildLogger({ subsystem: "session-sqlite" }).warn("SQLite session write failed", {
|
||||
agentId: scope.agentId,
|
||||
elapsedMs: Date.now() - startedAt,
|
||||
error,
|
||||
storePath,
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
export function resolveSqliteScope(
|
||||
scope: Pick<SessionAccessScope, "agentId" | "env" | "sessionKey" | "storePath">,
|
||||
): ResolvedSqliteScope {
|
||||
const scopedAgentId = resolveExplicitSqliteAgentId(scope);
|
||||
const storeTarget = scope.storePath
|
||||
? resolveSqliteTargetFromSessionStorePath(scope.storePath, { agentId: scopedAgentId })
|
||||
: undefined;
|
||||
const agentId = resolveSqliteAgentId({
|
||||
scopedAgentId,
|
||||
sessionKey: scope.sessionKey,
|
||||
storeAgentId: storeTarget?.agentId,
|
||||
useDefaultAgentForUnownedStore: Boolean(
|
||||
storeTarget?.path && !storeTarget.agentId && !scopedAgentId,
|
||||
),
|
||||
});
|
||||
if (!agentId) {
|
||||
throw new Error("Cannot resolve SQLite session scope without an agent id");
|
||||
}
|
||||
return {
|
||||
agentId,
|
||||
...(scope.env ? { env: scope.env } : {}),
|
||||
...(storeTarget ? { path: storeTarget.path } : {}),
|
||||
sessionKey: normalizeSqliteSessionKey(scope.sessionKey),
|
||||
};
|
||||
}
|
||||
|
||||
export function resolveSqliteReadScope(
|
||||
scope: Pick<SessionTranscriptReadScope, "agentId" | "env" | "sessionKey" | "storePath">,
|
||||
): ResolvedSqliteReadScope {
|
||||
const sessionKey = scope.sessionKey ? normalizeSqliteSessionKey(scope.sessionKey) : undefined;
|
||||
const scopedAgentId = resolveExplicitSqliteAgentId({ ...scope, sessionKey });
|
||||
const storeTarget = scope.storePath
|
||||
? resolveSqliteTargetFromSessionStorePath(scope.storePath, { agentId: scopedAgentId })
|
||||
: undefined;
|
||||
const agentId = resolveSqliteAgentId({
|
||||
scopedAgentId,
|
||||
sessionKey,
|
||||
storeAgentId: storeTarget?.agentId,
|
||||
useDefaultAgentForUnownedStore: Boolean(
|
||||
storeTarget?.path && !storeTarget.agentId && !scopedAgentId,
|
||||
),
|
||||
});
|
||||
if (!agentId) {
|
||||
throw new Error("Cannot resolve SQLite transcript read scope without an agent id");
|
||||
}
|
||||
return {
|
||||
agentId,
|
||||
...(scope.env ? { env: scope.env } : {}),
|
||||
...(storeTarget ? { path: storeTarget.path } : {}),
|
||||
...(sessionKey ? { sessionKey } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
function resolveExplicitSqliteAgentId(params: {
|
||||
agentId?: string;
|
||||
sessionKey?: string;
|
||||
}): string | undefined {
|
||||
return params.agentId
|
||||
? normalizeAgentId(params.agentId)
|
||||
: parseAgentSessionKey(params.sessionKey)?.agentId;
|
||||
}
|
||||
|
||||
export function resolveSqliteStoreScope(
|
||||
storePath: string,
|
||||
options: { agentId?: string } = {},
|
||||
): ResolvedSqliteScope {
|
||||
return resolveSqliteScope({
|
||||
...(options.agentId ? { agentId: options.agentId } : {}),
|
||||
sessionKey: "",
|
||||
storePath,
|
||||
});
|
||||
}
|
||||
|
||||
function resolveSqliteAgentId(params: {
|
||||
scopedAgentId?: string;
|
||||
sessionKey?: string;
|
||||
storeAgentId?: string;
|
||||
useDefaultAgentForUnownedStore?: boolean;
|
||||
}): string | undefined {
|
||||
const scopedAgentId = params.scopedAgentId ? normalizeAgentId(params.scopedAgentId) : undefined;
|
||||
if (scopedAgentId && params.storeAgentId && scopedAgentId !== params.storeAgentId) {
|
||||
throw new Error(
|
||||
`SQLite session store path belongs to agent ${params.storeAgentId}; requested agent ${scopedAgentId}.`,
|
||||
);
|
||||
}
|
||||
const resolved =
|
||||
scopedAgentId ??
|
||||
params.storeAgentId ??
|
||||
(params.sessionKey !== undefined ? resolveAgentIdFromSessionKey(params.sessionKey) : undefined);
|
||||
return resolved ?? (params.useDefaultAgentForUnownedStore ? DEFAULT_AGENT_ID : undefined);
|
||||
}
|
||||
|
||||
export function resolveSqliteTranscriptArchiveDirectory(
|
||||
scope: Pick<ResolvedSqliteReadScope, "agentId" | "env" | "path">,
|
||||
): string {
|
||||
const databasePath = resolveOpenClawAgentSqlitePath(toDatabaseOptions(scope));
|
||||
const databaseDir = path.dirname(databasePath);
|
||||
if (path.basename(databaseDir) !== "agent") {
|
||||
return databaseDir;
|
||||
}
|
||||
return path.join(path.dirname(databaseDir), "sessions");
|
||||
}
|
||||
|
||||
export function resolveSqliteTranscriptScope(
|
||||
scope: Pick<
|
||||
SessionTranscriptWriteScope,
|
||||
"agentId" | "env" | "sessionId" | "sessionKey" | "storePath"
|
||||
>,
|
||||
): ResolvedTranscriptScope {
|
||||
if (!scope.sessionId) {
|
||||
throw new Error(
|
||||
`Cannot resolve SQLite transcript scope without a session id: ${scope.sessionKey}`,
|
||||
);
|
||||
}
|
||||
if (!scope.sessionKey) {
|
||||
throw new Error(
|
||||
`Cannot resolve SQLite transcript scope without a session key: ${scope.sessionId}`,
|
||||
);
|
||||
}
|
||||
return {
|
||||
...resolveSqliteScope({ ...scope, sessionKey: scope.sessionKey }),
|
||||
sessionId: scope.sessionId,
|
||||
};
|
||||
}
|
||||
|
||||
export function resolveSqliteTranscriptReadScope(
|
||||
scope: Pick<
|
||||
SessionTranscriptReadScope,
|
||||
"agentId" | "env" | "sessionId" | "sessionKey" | "storePath"
|
||||
>,
|
||||
): ResolvedTranscriptReadScope {
|
||||
return {
|
||||
...resolveSqliteReadScope(scope),
|
||||
sessionId: scope.sessionId,
|
||||
};
|
||||
}
|
||||
|
||||
export function toDatabaseOptions(
|
||||
scope: Pick<ResolvedSqliteReadScope, "agentId" | "env" | "path">,
|
||||
): OpenClawAgentDatabaseOptions {
|
||||
return {
|
||||
agentId: scope.agentId,
|
||||
...(scope.env ? { env: scope.env } : {}),
|
||||
...(scope.path ? { path: scope.path } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
export function normalizeSqliteSessionKey(sessionKey: string): string {
|
||||
return normalizeStoreSessionKey(sessionKey);
|
||||
}
|
||||
|
||||
export function cloneSessionEntry(entry: SessionEntry): SessionEntry {
|
||||
return structuredClone(entry);
|
||||
}
|
||||
|
||||
export function formatSqliteSessionMarkerForScope(scope: ResolvedTranscriptScope): string {
|
||||
return formatSqliteSessionFileMarker({
|
||||
agentId: scope.agentId,
|
||||
sessionId: scope.sessionId,
|
||||
storePath: scope.path ?? resolveOpenClawAgentSqlitePath(toDatabaseOptions(scope)),
|
||||
});
|
||||
}
|
||||
103
src/config/sessions/session-accessor.sqlite-session-row.ts
Normal file
103
src/config/sessions/session-accessor.sqlite-session-row.ts
Normal file
@@ -0,0 +1,103 @@
|
||||
import {
|
||||
normalizeSqliteChatType,
|
||||
normalizeSqliteText,
|
||||
} from "./session-accessor.sqlite-normalize.js";
|
||||
import { bindSessionEntryProvenance } from "./session-accessor.sqlite-provenance.js";
|
||||
import { normalizeSqliteStatus } from "./session-accessor.sqlite-status.js";
|
||||
import type { SessionEntry } from "./types.js";
|
||||
|
||||
export function normalizeSqliteSessionEntryTimestamp(entry: SessionEntry): SessionEntry {
|
||||
if (typeof entry.updatedAt === "number" && Number.isFinite(entry.updatedAt)) {
|
||||
return entry;
|
||||
}
|
||||
const updatedAt =
|
||||
typeof entry.sessionStartedAt === "number" && Number.isFinite(entry.sessionStartedAt)
|
||||
? entry.sessionStartedAt
|
||||
: Date.now();
|
||||
return { ...entry, updatedAt };
|
||||
}
|
||||
|
||||
export function bindSqliteSessionRoot(params: {
|
||||
entry: SessionEntry;
|
||||
sessionKey: string;
|
||||
updatedAt: number;
|
||||
}) {
|
||||
const updatedAt = Number.isFinite(params.entry.updatedAt)
|
||||
? params.entry.updatedAt
|
||||
: params.updatedAt;
|
||||
return {
|
||||
session_id: params.entry.sessionId,
|
||||
session_key: params.sessionKey,
|
||||
session_scope: resolveSqliteSessionScope(params.entry, params.sessionKey),
|
||||
created_at: resolveSqliteSessionCreatedAt(params.entry, updatedAt),
|
||||
updated_at: updatedAt,
|
||||
...bindSessionEntryProvenance(params.entry),
|
||||
started_at: finiteSqliteNumber(params.entry.startedAt),
|
||||
ended_at: finiteSqliteNumber(params.entry.endedAt),
|
||||
status: normalizeSqliteStatus(params.entry.status),
|
||||
chat_type: normalizeSqliteChatType(params.entry.chatType),
|
||||
channel: resolveSqliteSessionChannel(params.entry),
|
||||
account_id: resolveSqliteSessionAccountId(params.entry),
|
||||
primary_conversation_id: null,
|
||||
model_provider: normalizeSqliteText(params.entry.modelProvider),
|
||||
model: normalizeSqliteText(params.entry.model),
|
||||
agent_harness_id: normalizeSqliteText(params.entry.agentHarnessId),
|
||||
parent_session_key: normalizeSqliteText(params.entry.parentSessionKey),
|
||||
spawned_by: normalizeSqliteText(params.entry.spawnedBy),
|
||||
display_name: resolveSqliteSessionDisplayName(params.entry),
|
||||
};
|
||||
}
|
||||
|
||||
function resolveSqliteSessionScope(
|
||||
entry: Pick<SessionEntry, "chatType">,
|
||||
sessionKey: string,
|
||||
): "conversation" | "shared-main" | "group" | "channel" {
|
||||
const chatType = normalizeSqliteChatType(entry.chatType);
|
||||
const normalizedKey = sessionKey.trim().toLowerCase();
|
||||
if (chatType === "direct" && (normalizedKey === "main" || normalizedKey.endsWith(":main"))) {
|
||||
return "shared-main";
|
||||
}
|
||||
if (chatType === "group" || chatType === "channel") {
|
||||
return chatType;
|
||||
}
|
||||
return "conversation";
|
||||
}
|
||||
|
||||
function resolveSqliteSessionCreatedAt(entry: SessionEntry, updatedAt: number): number {
|
||||
for (const candidate of [entry.sessionStartedAt, entry.startedAt, entry.updatedAt, updatedAt]) {
|
||||
if (typeof candidate === "number" && Number.isFinite(candidate) && candidate >= 0) {
|
||||
return candidate;
|
||||
}
|
||||
}
|
||||
return updatedAt;
|
||||
}
|
||||
|
||||
function finiteSqliteNumber(value: unknown): number | null {
|
||||
return typeof value === "number" && Number.isFinite(value) ? value : null;
|
||||
}
|
||||
|
||||
function resolveSqliteSessionChannel(entry: SessionEntry): string | null {
|
||||
return (
|
||||
normalizeSqliteText(entry.channel) ??
|
||||
normalizeSqliteText(entry.deliveryContext?.channel) ??
|
||||
normalizeSqliteText(entry.lastChannel) ??
|
||||
normalizeSqliteText(entry.origin?.provider)
|
||||
);
|
||||
}
|
||||
|
||||
function resolveSqliteSessionAccountId(entry: SessionEntry): string | null {
|
||||
return (
|
||||
normalizeSqliteText(entry.deliveryContext?.accountId) ??
|
||||
normalizeSqliteText(entry.lastAccountId) ??
|
||||
normalizeSqliteText(entry.origin?.accountId)
|
||||
);
|
||||
}
|
||||
|
||||
function resolveSqliteSessionDisplayName(entry: SessionEntry): string | null {
|
||||
return (
|
||||
normalizeSqliteText(entry.displayName) ??
|
||||
normalizeSqliteText(entry.label) ??
|
||||
normalizeSqliteText(entry.subject) ??
|
||||
normalizeSqliteText(entry.groupId)
|
||||
);
|
||||
}
|
||||
174
src/config/sessions/session-accessor.sqlite-transcript-state.ts
Normal file
174
src/config/sessions/session-accessor.sqlite-transcript-state.ts
Normal file
@@ -0,0 +1,174 @@
|
||||
import {
|
||||
executeSqliteQuerySync,
|
||||
executeSqliteQueryTakeFirstSync,
|
||||
} from "../../infra/kysely-sync.js";
|
||||
import type { OpenClawAgentDatabase } from "../../state/openclaw-agent-db.js";
|
||||
import { normalizeSqliteNumber } from "./session-accessor.sqlite-normalize.js";
|
||||
import { getSessionKysely, type ResolvedTranscriptScope } from "./session-accessor.sqlite-scope.js";
|
||||
import { deleteSessionTranscriptIndexInTransaction } from "./session-transcript-index.js";
|
||||
|
||||
export function ensureTranscriptSessionRoot(
|
||||
database: OpenClawAgentDatabase,
|
||||
scope: ResolvedTranscriptScope,
|
||||
updatedAt: number,
|
||||
): void {
|
||||
const db = getSessionKysely(database.db);
|
||||
executeSqliteQuerySync(
|
||||
database.db,
|
||||
db
|
||||
.insertInto("sessions")
|
||||
.values({
|
||||
session_id: scope.sessionId,
|
||||
session_key: scope.sessionKey,
|
||||
session_scope: "conversation",
|
||||
created_at: updatedAt,
|
||||
updated_at: updatedAt,
|
||||
})
|
||||
.onConflict((conflict) =>
|
||||
conflict.column("session_id").doUpdateSet({
|
||||
session_key: scope.sessionKey,
|
||||
updated_at: updatedAt,
|
||||
}),
|
||||
),
|
||||
);
|
||||
writeTranscriptSessionRoute(database, {
|
||||
sessionId: scope.sessionId,
|
||||
sessionKey: scope.sessionKey,
|
||||
updatedAt,
|
||||
});
|
||||
}
|
||||
|
||||
export function writeSessionRoute(
|
||||
database: OpenClawAgentDatabase,
|
||||
params: { sessionId: string; sessionKey: string; updatedAt: number },
|
||||
): void {
|
||||
const db = getSessionKysely(database.db);
|
||||
executeSqliteQuerySync(
|
||||
database.db,
|
||||
db
|
||||
.insertInto("session_routes")
|
||||
.values({
|
||||
session_key: params.sessionKey,
|
||||
session_id: params.sessionId,
|
||||
updated_at: params.updatedAt,
|
||||
})
|
||||
.onConflict((conflict) =>
|
||||
conflict.column("session_key").doUpdateSet({
|
||||
session_id: params.sessionId,
|
||||
updated_at: params.updatedAt,
|
||||
}),
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
function writeTranscriptSessionRoute(
|
||||
database: OpenClawAgentDatabase,
|
||||
params: { sessionId: string; sessionKey: string; updatedAt: number },
|
||||
): void {
|
||||
const db = getSessionKysely(database.db);
|
||||
const existing = executeSqliteQueryTakeFirstSync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("session_routes")
|
||||
.select("session_id")
|
||||
.where("session_key", "=", params.sessionKey),
|
||||
);
|
||||
// Late transcript-only appends may create routes, but cannot move a current
|
||||
// session key back to an older transcript id.
|
||||
if (existing && existing.session_id !== params.sessionId) {
|
||||
return;
|
||||
}
|
||||
writeSessionRoute(database, params);
|
||||
}
|
||||
|
||||
export function readNextTranscriptSeq(database: OpenClawAgentDatabase, sessionId: string): number {
|
||||
const db = getSessionKysely(database.db);
|
||||
const row = executeSqliteQueryTakeFirstSync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("transcript_events")
|
||||
.select((eb) => eb.fn.max<number | bigint>("seq").as("max_seq"))
|
||||
.where("session_id", "=", sessionId),
|
||||
);
|
||||
const maxSeq =
|
||||
row?.max_seq === null || row?.max_seq === undefined ? -1 : normalizeSqliteNumber(row.max_seq);
|
||||
return maxSeq + 1;
|
||||
}
|
||||
|
||||
function normalizeTranscriptMutationAtMs(value: number): number | undefined {
|
||||
const timestamp = Math.floor(value);
|
||||
return Number.isFinite(timestamp) && timestamp >= 0 ? timestamp : undefined;
|
||||
}
|
||||
|
||||
export 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,
|
||||
};
|
||||
}
|
||||
|
||||
export 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),
|
||||
);
|
||||
}
|
||||
|
||||
export function touchTranscriptMutationInTransaction(
|
||||
database: OpenClawAgentDatabase,
|
||||
sessionId: string,
|
||||
): void {
|
||||
const now = normalizeTranscriptMutationAtMs(Date.now());
|
||||
if (now !== undefined) {
|
||||
advanceTranscriptMutationAtInTransaction(database, sessionId, now, { strictly: true });
|
||||
}
|
||||
}
|
||||
|
||||
export function deleteSqliteTranscriptEventsInTransaction(
|
||||
database: OpenClawAgentDatabase,
|
||||
sessionId: string,
|
||||
): boolean {
|
||||
const db = getSessionKysely(database.db);
|
||||
executeSqliteQuerySync(
|
||||
database.db,
|
||||
db.deleteFrom("transcript_event_identities").where("session_id", "=", sessionId),
|
||||
);
|
||||
const result = executeSqliteQuerySync(
|
||||
database.db,
|
||||
db.deleteFrom("transcript_events").where("session_id", "=", sessionId),
|
||||
);
|
||||
// FTS rows have no FK onto transcript_events; clear them in this transaction.
|
||||
deleteSessionTranscriptIndexInTransaction(database.db, sessionId);
|
||||
return (result.numAffectedRows ?? 0n) > 0n;
|
||||
}
|
||||
460
src/config/sessions/session-accessor.sqlite-transcript-store.ts
Normal file
460
src/config/sessions/session-accessor.sqlite-transcript-store.ts
Normal file
@@ -0,0 +1,460 @@
|
||||
import type { AgentMessage } from "../../agents/runtime/index.js";
|
||||
import { redactTranscriptMessage } from "../../agents/transcript-redact.js";
|
||||
import {
|
||||
executeSqliteQuerySync,
|
||||
executeSqliteQueryTakeFirstSync,
|
||||
} from "../../infra/kysely-sync.js";
|
||||
import { redactSecrets } from "../../logging/redact.js";
|
||||
import type { OpenClawAgentDatabase } from "../../state/openclaw-agent-db.js";
|
||||
import type {
|
||||
TranscriptEvent,
|
||||
TranscriptMessageAppendOptions,
|
||||
} from "./session-accessor.sqlite-contract.js";
|
||||
import {
|
||||
findSqliteTranscriptEventInDatabase,
|
||||
loadSqliteTranscriptEventsFromDatabase,
|
||||
readTranscriptEventId,
|
||||
readTranscriptEventMessage,
|
||||
} from "./session-accessor.sqlite-read.js";
|
||||
import { getSessionKysely, type ResolvedTranscriptScope } from "./session-accessor.sqlite-scope.js";
|
||||
import {
|
||||
deleteSqliteTranscriptEventsInTransaction,
|
||||
ensureTranscriptSessionRoot,
|
||||
readNextTranscriptSeq,
|
||||
touchTranscriptMutationInTransaction,
|
||||
} from "./session-accessor.sqlite-transcript-state.js";
|
||||
import { indexAppendedTranscriptEventInTransaction } from "./session-transcript-index.js";
|
||||
import { createSessionTranscriptHeader } from "./transcript-header.js";
|
||||
import {
|
||||
isSessionTranscriptLeafControl,
|
||||
parseSessionTranscriptTreeEntry,
|
||||
} from "./transcript-tree.js";
|
||||
import { resolveVisibleTranscriptAppendParentId } from "./transcript-visible-events.js";
|
||||
|
||||
export function appendTranscriptEventInTransaction(
|
||||
database: OpenClawAgentDatabase,
|
||||
scope: ResolvedTranscriptScope,
|
||||
event: TranscriptEvent,
|
||||
options: { dedupeByMessageIdempotency?: boolean; touchMutation?: boolean } = {},
|
||||
): boolean {
|
||||
const db = getSessionKysely(database.db);
|
||||
const createdAt = readEventTimestamp(event) ?? Date.now();
|
||||
ensureTranscriptSessionRoot(database, scope, createdAt);
|
||||
const identity = readTranscriptEventIdentity(event);
|
||||
if (identity && readTranscriptIdentityByEventId(database, scope.sessionId, identity.eventId)) {
|
||||
return false;
|
||||
}
|
||||
if (
|
||||
identity?.messageIdempotencyKey &&
|
||||
options.dedupeByMessageIdempotency &&
|
||||
readTranscriptIdentityByMessageIdempotencyKey(
|
||||
database,
|
||||
scope.sessionId,
|
||||
identity.messageIdempotencyKey,
|
||||
)
|
||||
) {
|
||||
return false;
|
||||
}
|
||||
const seq = readNextTranscriptSeq(database, scope.sessionId);
|
||||
executeSqliteQuerySync(
|
||||
database.db,
|
||||
db.insertInto("transcript_events").values({
|
||||
session_id: scope.sessionId,
|
||||
seq,
|
||||
event_json: JSON.stringify(event),
|
||||
created_at: createdAt,
|
||||
}),
|
||||
);
|
||||
if (options.touchMutation !== false) {
|
||||
touchTranscriptMutationInTransaction(database, scope.sessionId);
|
||||
}
|
||||
indexAppendedTranscriptEventInTransaction(database.db, {
|
||||
sessionId: scope.sessionId,
|
||||
seq,
|
||||
event,
|
||||
eventId: identity?.eventId ?? null,
|
||||
createdAt,
|
||||
});
|
||||
if (!identity) {
|
||||
return true;
|
||||
}
|
||||
// Caller-checked appends may retain a duplicate key in the payload, but the
|
||||
// identity index can point at only one row.
|
||||
const indexedMessageIdempotencyKey =
|
||||
identity.messageIdempotencyKey &&
|
||||
!options.dedupeByMessageIdempotency &&
|
||||
readTranscriptIdentityByMessageIdempotencyKey(
|
||||
database,
|
||||
scope.sessionId,
|
||||
identity.messageIdempotencyKey,
|
||||
)
|
||||
? undefined
|
||||
: identity.messageIdempotencyKey;
|
||||
executeSqliteQuerySync(
|
||||
database.db,
|
||||
db
|
||||
.insertInto("transcript_event_identities")
|
||||
.values({
|
||||
session_id: scope.sessionId,
|
||||
event_id: identity.eventId,
|
||||
seq,
|
||||
event_type: identity.eventType,
|
||||
parent_id: identity.parentId,
|
||||
message_idempotency_key: indexedMessageIdempotencyKey,
|
||||
created_at: createdAt,
|
||||
})
|
||||
.onConflict((conflict) => conflict.columns(["session_id", "event_id"]).doNothing()),
|
||||
);
|
||||
return true;
|
||||
}
|
||||
|
||||
export 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,
|
||||
event: TranscriptEvent,
|
||||
seq: number,
|
||||
state: { seenEventIds: Set<string>; seenMessageIdempotencyKeys: Set<string> },
|
||||
): boolean {
|
||||
const db = getSessionKysely(database.db);
|
||||
const createdAt = readEventTimestamp(event) ?? Date.now();
|
||||
const identity = readTranscriptEventIdentity(event);
|
||||
if (identity && state.seenEventIds.has(identity.eventId)) {
|
||||
return false;
|
||||
}
|
||||
executeSqliteQuerySync(
|
||||
database.db,
|
||||
db.insertInto("transcript_events").values({
|
||||
session_id: scope.sessionId,
|
||||
seq,
|
||||
event_json: JSON.stringify(event),
|
||||
created_at: createdAt,
|
||||
}),
|
||||
);
|
||||
indexAppendedTranscriptEventInTransaction(database.db, {
|
||||
sessionId: scope.sessionId,
|
||||
seq,
|
||||
event,
|
||||
eventId: identity?.eventId ?? null,
|
||||
createdAt,
|
||||
});
|
||||
if (!identity) {
|
||||
return true;
|
||||
}
|
||||
state.seenEventIds.add(identity.eventId);
|
||||
const indexedMessageIdempotencyKey =
|
||||
identity.messageIdempotencyKey &&
|
||||
!state.seenMessageIdempotencyKeys.has(identity.messageIdempotencyKey)
|
||||
? identity.messageIdempotencyKey
|
||||
: undefined;
|
||||
if (indexedMessageIdempotencyKey) {
|
||||
state.seenMessageIdempotencyKeys.add(indexedMessageIdempotencyKey);
|
||||
}
|
||||
executeSqliteQuerySync(
|
||||
database.db,
|
||||
db.insertInto("transcript_event_identities").values({
|
||||
session_id: scope.sessionId,
|
||||
event_id: identity.eventId,
|
||||
seq,
|
||||
event_type: identity.eventType,
|
||||
parent_id: identity.parentId,
|
||||
message_idempotency_key: indexedMessageIdempotencyKey,
|
||||
created_at: createdAt,
|
||||
}),
|
||||
);
|
||||
return true;
|
||||
}
|
||||
|
||||
export function ensureTranscriptHeader(
|
||||
database: OpenClawAgentDatabase,
|
||||
scope: ResolvedTranscriptScope,
|
||||
cwd: string | undefined,
|
||||
now: number,
|
||||
): void {
|
||||
const db = getSessionKysely(database.db);
|
||||
const existing = executeSqliteQueryTakeFirstSync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("transcript_events")
|
||||
.select("seq")
|
||||
.where("session_id", "=", scope.sessionId)
|
||||
.limit(1),
|
||||
);
|
||||
if (existing) {
|
||||
return;
|
||||
}
|
||||
appendTranscriptEventInTransaction(
|
||||
database,
|
||||
scope,
|
||||
createSessionTranscriptHeader({ cwd, sessionId: scope.sessionId }),
|
||||
);
|
||||
ensureTranscriptSessionRoot(database, scope, now);
|
||||
}
|
||||
|
||||
export function readActiveTranscriptAppendParentId(
|
||||
database: OpenClawAgentDatabase,
|
||||
sessionId: string,
|
||||
): string | null {
|
||||
const db = getSessionKysely(database.db);
|
||||
const latest = executeSqliteQueryTakeFirstSync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("transcript_event_identities as ti")
|
||||
.innerJoin("transcript_events as te", (join) =>
|
||||
join.onRef("te.session_id", "=", "ti.session_id").onRef("te.seq", "=", "ti.seq"),
|
||||
)
|
||||
.select(["ti.event_type", "te.event_json"])
|
||||
.where("ti.session_id", "=", sessionId)
|
||||
.orderBy("ti.seq", "desc")
|
||||
.limit(1),
|
||||
);
|
||||
if (!latest) {
|
||||
return null;
|
||||
}
|
||||
try {
|
||||
const event = JSON.parse(latest.event_json) as unknown;
|
||||
const treeEntry = parseSessionTranscriptTreeEntry(event);
|
||||
if (!treeEntry) {
|
||||
return resolveVisibleTranscriptAppendParentId(
|
||||
loadSqliteTranscriptEventsFromDatabase(database, sessionId),
|
||||
);
|
||||
}
|
||||
if (latest.event_type !== "leaf") {
|
||||
return treeEntry.appendParentId;
|
||||
}
|
||||
const leafReferencesKnown =
|
||||
treeEntry.leafId !== undefined &&
|
||||
transcriptTreeReferenceExists(database, sessionId, treeEntry.leafId) &&
|
||||
transcriptTreeReferenceExists(database, sessionId, treeEntry.appendParentId);
|
||||
if (isSessionTranscriptLeafControl(event) && leafReferencesKnown) {
|
||||
return treeEntry.appendParentId;
|
||||
}
|
||||
} catch {
|
||||
// Fall through to the tolerant full-tree resolver.
|
||||
}
|
||||
return resolveVisibleTranscriptAppendParentId(
|
||||
loadSqliteTranscriptEventsFromDatabase(database, sessionId),
|
||||
);
|
||||
}
|
||||
|
||||
function transcriptTreeReferenceExists(
|
||||
database: OpenClawAgentDatabase,
|
||||
sessionId: string,
|
||||
eventId: string | null,
|
||||
): boolean {
|
||||
return (
|
||||
eventId === null || readTranscriptIdentityByEventId(database, sessionId, eventId) !== undefined
|
||||
);
|
||||
}
|
||||
|
||||
export function replaceSqliteTranscriptEventsInTransaction(
|
||||
database: OpenClawAgentDatabase,
|
||||
resolved: ResolvedTranscriptScope,
|
||||
events: readonly TranscriptEvent[],
|
||||
): void {
|
||||
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());
|
||||
let seq = 0;
|
||||
const seenEventIds = new Set<string>();
|
||||
const seenMessageIdempotencyKeys = new Set<string>();
|
||||
for (const event of events) {
|
||||
if (
|
||||
appendTranscriptEventRowInTransaction(database, resolved, event, seq, {
|
||||
seenEventIds,
|
||||
seenMessageIdempotencyKeys,
|
||||
})
|
||||
) {
|
||||
seq += 1;
|
||||
}
|
||||
}
|
||||
if (deleted || seq > 0) {
|
||||
touchTranscriptMutationInTransaction(database, resolved.sessionId);
|
||||
}
|
||||
}
|
||||
|
||||
export function readTranscriptIdentityByEventId(
|
||||
database: OpenClawAgentDatabase,
|
||||
sessionId: string,
|
||||
eventId: string,
|
||||
): { eventId: string; seq: number } | undefined {
|
||||
const db = getSessionKysely(database.db);
|
||||
const row = executeSqliteQueryTakeFirstSync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("transcript_event_identities")
|
||||
.select(["event_id", "seq"])
|
||||
.where("session_id", "=", sessionId)
|
||||
.where("event_id", "=", eventId),
|
||||
);
|
||||
return row ? { eventId: row.event_id, seq: row.seq } : undefined;
|
||||
}
|
||||
|
||||
function readTranscriptIdentityByMessageIdempotencyKey(
|
||||
database: OpenClawAgentDatabase,
|
||||
sessionId: string,
|
||||
idempotencyKey: string,
|
||||
): { eventId: string; seq: number } | undefined {
|
||||
const db = getSessionKysely(database.db);
|
||||
const row = executeSqliteQueryTakeFirstSync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("transcript_event_identities")
|
||||
.select(["event_id", "seq"])
|
||||
.where("session_id", "=", sessionId)
|
||||
.where("message_idempotency_key", "=", idempotencyKey)
|
||||
.orderBy("seq", "desc")
|
||||
.limit(1),
|
||||
);
|
||||
return row ? { eventId: row.event_id, seq: row.seq } : undefined;
|
||||
}
|
||||
|
||||
function readTranscriptMessageByIdempotencyKey(
|
||||
database: OpenClawAgentDatabase,
|
||||
scope: ResolvedTranscriptScope,
|
||||
idempotencyKey: string,
|
||||
): { messageId: string; message: unknown } | undefined {
|
||||
const identity = readTranscriptIdentityByMessageIdempotencyKey(
|
||||
database,
|
||||
scope.sessionId,
|
||||
idempotencyKey,
|
||||
);
|
||||
return identity ? readTranscriptMessageByIdentity(database, scope, identity) : undefined;
|
||||
}
|
||||
|
||||
export function readTranscriptMessageByScopedIdempotencyKey(
|
||||
database: OpenClawAgentDatabase,
|
||||
scope: ResolvedTranscriptScope,
|
||||
idempotencyKey: string,
|
||||
lookup: TranscriptMessageAppendOptions<unknown>["idempotencyLookup"],
|
||||
): { messageId: string; message: unknown } | undefined {
|
||||
if (lookup !== "scan-assistant") {
|
||||
return readTranscriptMessageByIdempotencyKey(database, scope, idempotencyKey);
|
||||
}
|
||||
const found = findSqliteTranscriptEventInDatabase(database, scope.sessionId, (event) => {
|
||||
const message = readTranscriptEventMessage(event);
|
||||
return message?.role === "assistant" && message.idempotencyKey === idempotencyKey;
|
||||
});
|
||||
if (!found) {
|
||||
return undefined;
|
||||
}
|
||||
const message = readTranscriptEventMessage(found.event);
|
||||
return message
|
||||
? { messageId: readTranscriptEventId(found.event) ?? idempotencyKey, message }
|
||||
: undefined;
|
||||
}
|
||||
|
||||
export function readTranscriptMessageByEventId(
|
||||
database: OpenClawAgentDatabase,
|
||||
scope: ResolvedTranscriptScope,
|
||||
eventId: string,
|
||||
): { messageId: string; message: unknown } | undefined {
|
||||
const identity = readTranscriptIdentityByEventId(database, scope.sessionId, eventId);
|
||||
return identity ? readTranscriptMessageByIdentity(database, scope, identity) : undefined;
|
||||
}
|
||||
|
||||
function readTranscriptMessageByIdentity(
|
||||
database: OpenClawAgentDatabase,
|
||||
scope: ResolvedTranscriptScope,
|
||||
identity: { eventId: string; seq: number },
|
||||
): { messageId: string; message: unknown } | undefined {
|
||||
const db = getSessionKysely(database.db);
|
||||
const eventRow = executeSqliteQueryTakeFirstSync(
|
||||
database.db,
|
||||
db
|
||||
.selectFrom("transcript_events")
|
||||
.select(["event_json"])
|
||||
.where("session_id", "=", scope.sessionId)
|
||||
.where("seq", "=", identity.seq),
|
||||
);
|
||||
if (!eventRow) {
|
||||
return undefined;
|
||||
}
|
||||
const event = JSON.parse(eventRow.event_json) as { message?: unknown };
|
||||
return { messageId: identity.eventId, message: event.message };
|
||||
}
|
||||
|
||||
function readTranscriptEventIdentity(event: unknown):
|
||||
| {
|
||||
eventId: string;
|
||||
eventType: string | null;
|
||||
parentId: string | null;
|
||||
messageIdempotencyKey: string | null;
|
||||
}
|
||||
| undefined {
|
||||
if (!event || typeof event !== "object" || Array.isArray(event)) {
|
||||
return undefined;
|
||||
}
|
||||
const record = event as Record<string, unknown>;
|
||||
const eventId = typeof record.id === "string" && record.id.trim() ? record.id.trim() : undefined;
|
||||
return eventId
|
||||
? {
|
||||
eventId,
|
||||
eventType: typeof record.type === "string" ? record.type : null,
|
||||
parentId: typeof record.parentId === "string" ? record.parentId : null,
|
||||
messageIdempotencyKey: readMessageIdempotencyKey(record.message),
|
||||
}
|
||||
: undefined;
|
||||
}
|
||||
|
||||
export function readMessageIdempotencyKey(message: unknown): string | null {
|
||||
if (!message || typeof message !== "object" || Array.isArray(message)) {
|
||||
return null;
|
||||
}
|
||||
const value = (message as { idempotencyKey?: unknown }).idempotencyKey;
|
||||
return typeof value === "string" && value.trim() ? value.trim() : null;
|
||||
}
|
||||
|
||||
function readEventTimestamp(event: unknown): number | undefined {
|
||||
if (!event || typeof event !== "object" || Array.isArray(event)) {
|
||||
return undefined;
|
||||
}
|
||||
const value = (event as { timestamp?: unknown }).timestamp;
|
||||
if (typeof value === "number" && Number.isFinite(value)) {
|
||||
return value;
|
||||
}
|
||||
if (typeof value !== "string" || !value.trim()) {
|
||||
return undefined;
|
||||
}
|
||||
const parsed = Date.parse(value);
|
||||
return Number.isFinite(parsed) ? parsed : undefined;
|
||||
}
|
||||
|
||||
export function redactTranscriptMessageForStorage<TMessage>(
|
||||
message: TMessage,
|
||||
options: Pick<TranscriptMessageAppendOptions<TMessage>, "config">,
|
||||
): TMessage {
|
||||
return isTranscriptAgentMessage(message)
|
||||
? (redactTranscriptMessage(message, options.config) as TMessage)
|
||||
: redactSecrets(message);
|
||||
}
|
||||
|
||||
function isTranscriptAgentMessage(value: unknown): value is AgentMessage {
|
||||
return (
|
||||
typeof value === "object" &&
|
||||
value !== null &&
|
||||
!Array.isArray(value) &&
|
||||
typeof (value as { role?: unknown }).role === "string"
|
||||
);
|
||||
}
|
||||
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user