mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
refactor(core): retire pre-agent session migration (#163160)
* refactor(core): retire pre-agent session migration Retire automatic discovery, merge, relocation, and recovery for the January 2026 global session layout. Keep configured JSON stores and current per-agent ownership and repair contracts. Consolidate approval option splitting and derive audit scalar types from existing schemas. * refactor(core): finish retired session layout cutover Move remaining migration fixtures to supported configured and per-agent stores, preserving ownership, receipt, refusal, and lifecycle checks. Remove the orphaned first-line reader after its sole legacy adapter was retired.
This commit is contained in:
parent
61b940ae26
commit
acfd767f15
34 changed files with 252 additions and 1106 deletions
|
|
@ -2423,7 +2423,7 @@ src/infra/state-migrations.doctor.ts 3
|
|||
src/infra/state-migrations.exec-approvals.ts 1
|
||||
src/infra/state-migrations.fs.ts 1
|
||||
src/infra/state-migrations.legacy-session-store.ts 3
|
||||
src/infra/state-migrations.legacy-sessions.ts 2
|
||||
src/infra/state-migrations.legacy-sessions.ts 1
|
||||
src/infra/state-migrations.mcp-oauth.ts 2
|
||||
src/infra/state-migrations.meeting-transcripts-detection.ts 4
|
||||
src/infra/state-migrations.meeting-transcripts-files.ts 5
|
||||
|
|
@ -2431,7 +2431,7 @@ src/infra/state-migrations.meeting-transcripts-verify.ts 2
|
|||
src/infra/state-migrations.meeting-transcripts.ts 2
|
||||
src/infra/state-migrations.node-host.ts 2
|
||||
src/infra/state-migrations.runtime-state.ts 1
|
||||
src/infra/state-migrations.session-store.ts 9
|
||||
src/infra/state-migrations.session-store.ts 6
|
||||
src/infra/state-migrations.shared-auth-store.ts 1
|
||||
src/infra/state-migrations.state-dir.ts 1
|
||||
src/infra/state-migrations.storage.ts 2
|
||||
|
|
|
|||
|
|
@ -17,7 +17,7 @@ auth health, sandbox images, and plugin installs.
|
|||
|
||||
Doctor can migrate supported on-disk layouts into the current structure:
|
||||
|
||||
- Session rows and transcripts: import legacy `sessions.json` and JSONL history from `~/.openclaw/sessions/` or per-agent `sessions/` directories into `~/.openclaw/agents/<agentId>/agent/openclaw-agent.sqlite`
|
||||
- Session rows and transcripts: import legacy `sessions.json` and JSONL history from per-agent `sessions/` directories or explicitly configured stores into `~/.openclaw/agents/<agentId>/agent/openclaw-agent.sqlite`
|
||||
- Agent dir: from `~/.openclaw/agent/` to `~/.openclaw/agents/<agentId>/agent/`
|
||||
- WhatsApp auth state (Baileys): from legacy `~/.openclaw/credentials/*.json` (except `oauth.json`) to `~/.openclaw/credentials/whatsapp/<accountId>/...` (default account id: `default`)
|
||||
- Signed device identity: from `~/.openclaw/identity/device.json` into the `primary` `device_identities` row in `state/openclaw.sqlite`; Doctor also owns repair of invalid canonical rows; Gateway and node-host startup refuse an unimported identity instead of creating a replacement, and leave the separate device-auth file untouched
|
||||
|
|
@ -32,6 +32,8 @@ auth health, sandbox images, and plugin installs.
|
|||
|
||||
Repair prepares retained archive media normalization from read-only database snapshots before stopping a managed Gateway. Unchanged archives require no archive write transaction. Doctor records verified content and file identities in existing migration metadata, so another run at the same version skips parsing and reading unchanged archive copies. Imports, restores, changed files, and new versions invalidate those facts; canonical blob digests are still checked. Actual repairs retain stopped-writer authority and source revalidation.
|
||||
|
||||
Automatic discovery no longer imports the pre-agent `~/.openclaw/sessions/` layout. An explicit session-store path remains supported, including a configured path at that location.
|
||||
|
||||
Legacy session-file import and repair belong to Doctor. Gateway startup checks readiness without importing those files; runtime session access uses only SQLite. An unreadable legacy session index and its transcripts remain at their original paths, and repeated startups refuse readiness with the Doctor command for the active profile. Stop the Gateway, back up its state, repair the named source, and run `openclaw doctor --fix` before restarting it. The [targeted migration sequence](/cli/doctor#session-sqlite-migration) provides inspection and validation evidence. Current SQLite maintenance does not require legacy files to remain on disk.
|
||||
|
||||
When an unavailable plugin still needs legacy session files, Doctor retains those originals after verifying the core import. Startup accepts the retained files only when their session owners have matching verified imports. An unused configured agent does not need an empty database for another agent's history. Changed, unassigned, or unimported source rows still require repair before startup.
|
||||
|
|
|
|||
|
|
@ -1,5 +1,5 @@
|
|||
// Versioned metadata-only activity audit query payloads.
|
||||
import { type TProperties, type TSchema, Type } from "typebox";
|
||||
import { type Static, type TProperties, type TSchema, Type } from "typebox";
|
||||
import { closedObject } from "./closed-object.js";
|
||||
import { NonEmptyString } from "./primitives.js";
|
||||
|
||||
|
|
@ -525,31 +525,19 @@ type AuditActivityInboundMessageV1Terminal =
|
|||
status: "succeeded";
|
||||
outcome: "completed";
|
||||
errorCode?: never;
|
||||
reasonCode?:
|
||||
| "fast_abort"
|
||||
| "plugin_bound_handled"
|
||||
| "plugin_bound_unavailable"
|
||||
| "plugin_bound_declined"
|
||||
| "before_dispatch_handled"
|
||||
| "acp_dispatch_completed"
|
||||
| "acp_dispatch_empty"
|
||||
| "active_run_injected";
|
||||
reasonCode?: Static<typeof inboundCompletedReasonSchema>;
|
||||
}
|
||||
| {
|
||||
status: "blocked";
|
||||
outcome: "skipped";
|
||||
errorCode?: never;
|
||||
reasonCode?:
|
||||
| "duplicate"
|
||||
| "reply_operation_active"
|
||||
| "reply_operation_aborted"
|
||||
| "acp_dispatch_aborted";
|
||||
reasonCode?: Static<typeof inboundSkippedReasonSchema>;
|
||||
}
|
||||
| {
|
||||
status: "failed";
|
||||
outcome: "failed";
|
||||
errorCode: "message_processing_failed";
|
||||
reasonCode?: "acp_dispatch_failed" | "plugin_bound_error";
|
||||
reasonCode?: Static<typeof inboundFailureReasonSchema>;
|
||||
};
|
||||
export type AuditActivityInboundMessageV1 = AuditActivityMessageRecordBaseV1 & {
|
||||
eventType: "inbound_message";
|
||||
|
|
@ -573,21 +561,16 @@ type AuditActivityOutboundMessageV1Terminal =
|
|||
status: "blocked";
|
||||
outcome: "suppressed";
|
||||
errorCode?: never;
|
||||
reasonCode:
|
||||
| "cancelled_by_message_sending_hook"
|
||||
| "cancelled_by_reply_payload_sending_hook"
|
||||
| "empty_after_message_sending_hook"
|
||||
| "empty_after_reply_payload_sending_hook"
|
||||
| "no_visible_payload";
|
||||
reasonCode: Static<typeof outboundSuppressedReasonSchema>;
|
||||
failureStage?: never;
|
||||
deliveryKind?: never;
|
||||
}
|
||||
| {
|
||||
status: "failed";
|
||||
outcome: "failed";
|
||||
errorCode: "message_delivery_failed" | "message_delivery_partial_failure";
|
||||
errorCode: Static<typeof outboundFailureErrorSchema>;
|
||||
reasonCode?: never;
|
||||
failureStage: "platform_send" | "queue" | "unknown";
|
||||
failureStage: Static<typeof outboundFailureStageSchema>;
|
||||
deliveryKind?: "text" | "media" | "other";
|
||||
}
|
||||
| {
|
||||
|
|
@ -595,7 +578,7 @@ type AuditActivityOutboundMessageV1Terminal =
|
|||
outcome: "unknown";
|
||||
errorCode?: never;
|
||||
reasonCode?: never;
|
||||
failureStage: "platform_send" | "queue" | "unknown";
|
||||
failureStage: Static<typeof outboundFailureStageSchema>;
|
||||
deliveryKind?: never;
|
||||
};
|
||||
export type AuditActivityOutboundMessageV1 = AuditActivityMessageRecordBaseV1 & {
|
||||
|
|
|
|||
|
|
@ -83,7 +83,6 @@ export async function loadAndMaybeMigrateDoctorConfig(params: {
|
|||
observe: false,
|
||||
invocationPurpose: "doctor",
|
||||
repairPrefixedConfig: shouldRepair,
|
||||
recoverCorruptTargetStore: shouldRepair,
|
||||
doctorOnlyStateMigrations: shouldRepair,
|
||||
preparePluginMetadataSnapshot: true,
|
||||
...(params.agentDatabaseMigrationDiscovery
|
||||
|
|
|
|||
|
|
@ -114,23 +114,6 @@ describe("runDoctorConfigPreflight state migration input", () => {
|
|||
readConfigFileSnapshot.mockReset();
|
||||
});
|
||||
|
||||
it("passes explicit corrupt-target recovery to state migrations", async () => {
|
||||
await runDoctorConfigPreflight({
|
||||
...options,
|
||||
recoverCorruptTargetStore: true,
|
||||
});
|
||||
|
||||
expect(autoMigrateLegacyState).toHaveBeenCalledWith({
|
||||
cfg: { gateway: { mode: "local", port: 19091 } },
|
||||
configIncludedPaths: [],
|
||||
env: process.env,
|
||||
log: undefined,
|
||||
recoverCorruptTargetStore: true,
|
||||
doctorOnlyStateMigrations: undefined,
|
||||
onStepReceipt: expect.any(Function),
|
||||
});
|
||||
});
|
||||
|
||||
it("preserves a retired custom cron partition with invalid Gateway config", async () => {
|
||||
const sourceConfig = {
|
||||
gateway: { mode: "local", port: "not-a-port" },
|
||||
|
|
@ -211,7 +194,6 @@ describe("runDoctorConfigPreflight state migration input", () => {
|
|||
configIncludedPaths: includedPaths,
|
||||
env: process.env,
|
||||
log: undefined,
|
||||
recoverCorruptTargetStore: undefined,
|
||||
doctorOnlyStateMigrations: undefined,
|
||||
onStepReceipt: expect.any(Function),
|
||||
});
|
||||
|
|
|
|||
|
|
@ -265,7 +265,6 @@ async function runDoctorConfigPreflightOperation(
|
|||
...(pluginDoctorConfig ? { pluginDoctorConfig } : {}),
|
||||
configIncludedPaths: snapshot.includedPaths ?? [],
|
||||
env: process.env,
|
||||
recoverCorruptTargetStore: options.recoverCorruptTargetStore,
|
||||
doctorOnlyStateMigrations: options.doctorOnlyStateMigrations,
|
||||
invocationPurpose: options.invocationPurpose,
|
||||
...(options.agentDatabaseMigrationDiscovery
|
||||
|
|
|
|||
|
|
@ -188,18 +188,7 @@ export function resolveDoctorSessionSqliteTargets(params: {
|
|||
const candidates = discoversHistory
|
||||
? resolveAllAgentSessionStoreCandidateTargetsSync(params.cfg, { env: params.env })
|
||||
: resolveAllAgentSessionStoreTargetsSync(params.cfg, { env: params.env });
|
||||
const legacyStorePath = path.join(resolveStateDir(params.env), "sessions", "sessions.json");
|
||||
const legacyTargets =
|
||||
discoversHistory && fs.existsSync(legacyStorePath)
|
||||
? resolveSessionStoreTargets(params.cfg, { allAgents: true }, { env: params.env }).map(
|
||||
(target) => ({
|
||||
agentId: target.agentId,
|
||||
sqlitePath: resolveTargetSqlitePath(target, params.env),
|
||||
storePath: legacyStorePath,
|
||||
}),
|
||||
)
|
||||
: [];
|
||||
const targets = [...legacyTargets, ...candidates].map((target) => ({
|
||||
const targets = candidates.map((target) => ({
|
||||
target,
|
||||
sqlitePath: resolveTargetSqlitePath(target, params.env),
|
||||
}));
|
||||
|
|
|
|||
|
|
@ -130,14 +130,28 @@ describe("resumed Codex session binding migration", () => {
|
|||
},
|
||||
);
|
||||
|
||||
it.each(["default", "legacy-root"] as const)(
|
||||
it.each(["default", "configured"] as const)(
|
||||
"does not resurrect an imported session because an unrelated %s source is unimported",
|
||||
async (layout) => {
|
||||
await withOpenClawTestState({ label: `codex-mixed-source-${layout}` }, async (state) => {
|
||||
const { cfg, scope } = await seedDeferredPluginSessionSource(state, "external", "codex");
|
||||
await runDoctorSessionSqlite({ cfg, env: state.env, allAgents: true, mode: "import" });
|
||||
const { cfg, scope } = await seedDeferredPluginSessionSource(
|
||||
state,
|
||||
layout === "default" ? "external" : "default",
|
||||
"codex",
|
||||
);
|
||||
if (layout === "configured") {
|
||||
cfg.session = { store: state.path("configured-sessions/{agentId}/sessions.json") };
|
||||
}
|
||||
await runDoctorSessionSqlite({
|
||||
cfg,
|
||||
env: state.env,
|
||||
...(layout === "configured"
|
||||
? { store: scope.storePath, agent: "main" }
|
||||
: { allAgents: true }),
|
||||
mode: "import",
|
||||
});
|
||||
const directory =
|
||||
layout === "default" ? state.sessionsDir("main") : state.statePath("sessions");
|
||||
layout === "default" ? state.sessionsDir("main") : state.path("configured-sessions/main");
|
||||
fs.mkdirSync(directory, { recursive: true });
|
||||
const unrelatedStore = path.join(directory, "sessions.json");
|
||||
const unrelatedSource = JSON.stringify({
|
||||
|
|
|
|||
|
|
@ -145,11 +145,11 @@ describe("session sources needed by deferred plugin migrations", () => {
|
|||
"retains an ordinary import's $kind when another Doctor records pending work before unlink (unused agent: $unusedAgent)",
|
||||
async ({ kind, unusedAgent }) => {
|
||||
await withOpenClawTestState({ label: "deferred-plugin-archive-race" }, async (state) => {
|
||||
const { cfg, storePath, originals, scope } = await seedDeferredPluginSessionSource(
|
||||
state,
|
||||
unusedAgent ? "legacy-root" : "external",
|
||||
);
|
||||
const unusedDatabase = state.statePath("agents/ops/agent/openclaw-agent.sqlite");
|
||||
const { cfg, storePath, originals, scope } = await seedDeferredPluginSessionSource(state);
|
||||
const unusedDatabase = resolveSqliteTargetFromSessionStorePath(storePath, {
|
||||
agentId: "ops",
|
||||
env: state.env,
|
||||
}).path;
|
||||
if (unusedAgent) {
|
||||
cfg.agents = { ...cfg.agents, entries: { ...cfg.agents?.entries, ops: {} } };
|
||||
}
|
||||
|
|
@ -486,9 +486,9 @@ describe("session sources needed by deferred plugin migrations", () => {
|
|||
it.each([
|
||||
{ layout: "external", missingTranscript: false },
|
||||
{ layout: "default", missingTranscript: false },
|
||||
{ layout: "legacy-root", missingTranscript: false },
|
||||
{ layout: "legacy-root-custom-store", missingTranscript: false },
|
||||
{ layout: "legacy-root-with-unused-agent", missingTranscript: false },
|
||||
{ layout: "configured-root", missingTranscript: false },
|
||||
{ layout: "configured-home-store", missingTranscript: false },
|
||||
{ layout: "configured-root-with-unused-agent", missingTranscript: false },
|
||||
{ layout: "default", missingTranscript: true },
|
||||
{ layout: "relocated-interrupted", missingTranscript: false },
|
||||
] as const)(
|
||||
|
|
@ -497,17 +497,19 @@ describe("session sources needed by deferred plugin migrations", () => {
|
|||
await withOpenClawTestState({ label: "deferred-plugin-session-source" }, async (state) => {
|
||||
const { cfg, storePath, originals, scope } = await seedDeferredPluginSessionSource(
|
||||
state,
|
||||
layout === "legacy-root-with-unused-agent" || layout === "legacy-root-custom-store"
|
||||
layout === "configured-root" || layout === "configured-root-with-unused-agent"
|
||||
? "legacy-root"
|
||||
: layout === "relocated-interrupted"
|
||||
? "default"
|
||||
: layout,
|
||||
: layout === "configured-home-store"
|
||||
? "external"
|
||||
: layout === "relocated-interrupted"
|
||||
? "default"
|
||||
: layout,
|
||||
"fixture-plugin",
|
||||
missingTranscript ? "declared" : undefined,
|
||||
);
|
||||
if (layout === "legacy-root-custom-store") {
|
||||
scope.storePath = state.path("custom/sessions.json");
|
||||
cfg.session = { store: scope.storePath };
|
||||
if (layout === "configured-root" || layout === "configured-root-with-unused-agent") {
|
||||
scope.storePath = storePath;
|
||||
cfg.session = { store: storePath };
|
||||
}
|
||||
let foreignSource: { path: string; bytes: Buffer } | undefined;
|
||||
if (layout === "relocated-interrupted") {
|
||||
|
|
@ -521,8 +523,11 @@ describe("session sources needed by deferred plugin migrations", () => {
|
|||
fs.writeFileSync(storePath, JSON.stringify(entries));
|
||||
originals.set(storePath, fs.readFileSync(storePath));
|
||||
}
|
||||
const unusedDatabase = state.statePath("agents/ops/agent/openclaw-agent.sqlite");
|
||||
if (layout === "legacy-root-with-unused-agent") {
|
||||
const unusedDatabase = resolveSqliteTargetFromSessionStorePath(storePath, {
|
||||
agentId: "ops",
|
||||
env: state.env,
|
||||
}).path;
|
||||
if (layout === "configured-root-with-unused-agent") {
|
||||
cfg.agents = { ...cfg.agents, entries: { ...cfg.agents?.entries, ops: {} } };
|
||||
expect(fs.existsSync(unusedDatabase)).toBe(false);
|
||||
}
|
||||
|
|
@ -532,7 +537,7 @@ describe("session sources needed by deferred plugin migrations", () => {
|
|||
const run = () =>
|
||||
runDoctorSessionSqlite({ cfg, env: state.env, allAgents: true, mode: "import" });
|
||||
const imported = await run();
|
||||
if (layout === "legacy-root-with-unused-agent") {
|
||||
if (layout === "configured-root-with-unused-agent") {
|
||||
expect(fs.existsSync(unusedDatabase)).toBe(false);
|
||||
}
|
||||
expect(imported.totals.importedEntries).toBe(missingTranscript ? 3 : 2);
|
||||
|
|
@ -562,19 +567,19 @@ describe("session sources needed by deferred plugin migrations", () => {
|
|||
if (
|
||||
!missingTranscript &&
|
||||
(layout === "default" ||
|
||||
layout === "legacy-root" ||
|
||||
layout === "legacy-root-custom-store")
|
||||
layout === "configured-root" ||
|
||||
layout === "configured-home-store")
|
||||
) {
|
||||
const identities = new Map(
|
||||
[...originals.keys()].map((file) => [file, readMigrationArtifactIdentity(file)]),
|
||||
);
|
||||
await autoMigrateLegacyState({
|
||||
cfg:
|
||||
layout === "legacy-root-custom-store"
|
||||
? { ...cfg, session: { store: "~/custom/sessions.json" } }
|
||||
layout === "configured-home-store"
|
||||
? { ...cfg, session: { store: "~/external-sessions/sessions.json" } }
|
||||
: cfg,
|
||||
env:
|
||||
layout === "legacy-root-custom-store"
|
||||
layout === "configured-home-store"
|
||||
? { ...state.env, OPENCLAW_HOME: state.root, HOME: state.root }
|
||||
: state.env,
|
||||
homedir: () => state.home,
|
||||
|
|
@ -753,22 +758,27 @@ describe("session sources needed by deferred plugin migrations", () => {
|
|||
},
|
||||
);
|
||||
|
||||
it("admits mixed top-level sources only after the live owner has a verified import", async () => {
|
||||
it("admits a configured shared source only after the live owner has a verified import", async () => {
|
||||
await withOpenClawTestState({ label: "deferred-mixed-retained-owner" }, async (state) => {
|
||||
const { cfg, storePath, originals, scope } = await seedDeferredPluginSessionSource(
|
||||
state,
|
||||
"legacy-root",
|
||||
);
|
||||
const { cfg, storePath, originals, scope } = await seedDeferredPluginSessionSource(state);
|
||||
cfg.agents = { ownership: "explicit", entries: { main: {}, retired: {} } };
|
||||
const retiredPath = openOpenClawAgentDatabase({ agentId: "retired", env: state.env }).path;
|
||||
const retiredPath = openOpenClawAgentDatabase({
|
||||
agentId: "retired",
|
||||
env: state.env,
|
||||
path: resolveSqliteTargetFromSessionStorePath(storePath, {
|
||||
agentId: "retired",
|
||||
env: state.env,
|
||||
}).path,
|
||||
}).path;
|
||||
closeOpenClawAgentDatabasesForTest();
|
||||
const deletion = beginAgentDeletionJournal(
|
||||
{
|
||||
agentId: "retired",
|
||||
operationId: "delete-mixed-retired-owner",
|
||||
agentDir: path.dirname(retiredPath),
|
||||
agentDir: state.agentDir("retired"),
|
||||
sessionsDir: state.sessionsDir("retired"),
|
||||
workspaceDir: state.statePath("workspace-retired"),
|
||||
databasePaths: [retiredPath],
|
||||
deleteFiles: false,
|
||||
},
|
||||
{ env: state.env },
|
||||
|
|
@ -837,14 +847,11 @@ describe("session sources needed by deferred plugin migrations", () => {
|
|||
});
|
||||
});
|
||||
|
||||
it.each(["unimported-owner", "unassigned", "retired-owner", "malformed", "unreadable"] as const)(
|
||||
it.each(["unimported-owner", "retired-owner", "malformed", "unreadable"] as const)(
|
||||
"keeps readiness blocked for a retained source with %s state",
|
||||
async (kind) => {
|
||||
await withOpenClawTestState({ label: `deferred-readiness-${kind}` }, async (state) => {
|
||||
const { cfg, storePath } = await seedDeferredPluginSessionSource(
|
||||
state,
|
||||
kind === "unimported-owner" ? "external" : "legacy-root",
|
||||
);
|
||||
const { cfg, storePath } = await seedDeferredPluginSessionSource(state);
|
||||
cfg.agents = {
|
||||
ownership: "explicit",
|
||||
...(kind === "unimported-owner"
|
||||
|
|
@ -859,12 +866,7 @@ describe("session sources needed by deferred plugin migrations", () => {
|
|||
fs.unlinkSync(storePath);
|
||||
fs.mkdirSync(storePath);
|
||||
} else {
|
||||
const key =
|
||||
kind === "unimported-owner"
|
||||
? "agent:ops:waiting"
|
||||
: kind === "retired-owner"
|
||||
? "agent:retired:waiting"
|
||||
: "voice:unassigned";
|
||||
const key = kind === "unimported-owner" ? "agent:ops:waiting" : "agent:retired:waiting";
|
||||
source[key] = { sessionId: "waiting", updatedAt: 1 };
|
||||
fs.writeFileSync(storePath, JSON.stringify(source));
|
||||
}
|
||||
|
|
|
|||
|
|
@ -28,6 +28,7 @@ it.each([false, true])(
|
|||
async (interrupt) => {
|
||||
await withOpenClawTestState({ label: "shared-orphan-settlement" }, async (state) => {
|
||||
const { cfg, storePath } = await seedDeferredPluginSessionSource(state, "legacy-root");
|
||||
cfg.session = { store: storePath };
|
||||
cfg.agents = { ...cfg.agents, entries: { ...cfg.agents?.entries, ops: {} } };
|
||||
const entries: Record<string, unknown> = JSON.parse(fs.readFileSync(storePath, "utf8"));
|
||||
entries["agent:ops:kept"] = {
|
||||
|
|
|
|||
|
|
@ -57,12 +57,12 @@ function seedUnreadableSiblingDeletion(store: TestStore): void {
|
|||
|
||||
describe("runDoctorSessionSqlite", () => {
|
||||
it.each(["destination", "shared-state"])(
|
||||
"holds a top-level legacy import with %s orphaned WAL history",
|
||||
"holds a configured legacy import with %s orphaned WAL history",
|
||||
async (location) => {
|
||||
const stateDir = fs.realpathSync.native(
|
||||
autoCleanupTempDirs.make("doctor-held-legacy-store-"),
|
||||
);
|
||||
const storePath = path.join(stateDir, "sessions", "sessions.json");
|
||||
const storePath = path.join(stateDir, "agents", "main", "sessions", "sessions.json");
|
||||
const sqlitePath = path.join(stateDir, "agents", "main", "agent", "openclaw-agent.sqlite");
|
||||
const walPath =
|
||||
location === "destination"
|
||||
|
|
@ -77,6 +77,7 @@ describe("runDoctorSessionSqlite", () => {
|
|||
fs.writeFileSync(walPath, wal);
|
||||
const cfg: OpenClawConfig = {
|
||||
agents: { ownership: "explicit", entries: { main: {} } },
|
||||
session: { store: storePath },
|
||||
};
|
||||
const env = { ...process.env, OPENCLAW_STATE_DIR: stateDir };
|
||||
|
||||
|
|
@ -320,7 +321,7 @@ describe("runDoctorSessionSqlite", () => {
|
|||
).toBeUndefined();
|
||||
});
|
||||
|
||||
it("partitions the retired top-level store without guessing unscoped ownership", async () => {
|
||||
it("uses configured fixed-store ownership for an explicitly selected former global store", async () => {
|
||||
const stateDir = autoCleanupTempDirs.make("openclaw-doctor-retired-sessions-");
|
||||
const sessionDir = path.join(stateDir, "sessions");
|
||||
const storePath = path.join(sessionDir, "sessions.json");
|
||||
|
|
@ -339,7 +340,7 @@ describe("runDoctorSessionSqlite", () => {
|
|||
sessionId: "ops-session",
|
||||
updatedAt: 30,
|
||||
},
|
||||
"voice:ambiguous": { sessionId: "ambiguous-session", updatedAt: 40 },
|
||||
"agent:main:voice:opaque": { sessionId: "opaque-session", updatedAt: 40 },
|
||||
}),
|
||||
{ mode: 0o600 },
|
||||
);
|
||||
|
|
@ -354,9 +355,17 @@ describe("runDoctorSessionSqlite", () => {
|
|||
{ mode: 0o600 },
|
||||
);
|
||||
|
||||
const cfg = {
|
||||
agents: { ownership: "explicit" as const, entries: { main: {}, ops: {} } },
|
||||
};
|
||||
const agents = { ownership: "explicit" as const, entries: { main: {}, ops: {} } };
|
||||
const ignored = await runDoctorSessionSqlite({
|
||||
allAgents: true,
|
||||
cfg: { agents },
|
||||
env,
|
||||
mode: "dry-run",
|
||||
});
|
||||
expect(ignored.targets).toEqual([]);
|
||||
expect(fs.existsSync(storePath)).toBe(true);
|
||||
|
||||
const cfg = { agents, session: { store: storePath } };
|
||||
const report = await runDoctorSessionSqlite({
|
||||
allAgents: true,
|
||||
cfg,
|
||||
|
|
@ -366,18 +375,17 @@ describe("runDoctorSessionSqlite", () => {
|
|||
|
||||
expect(report.targets.map((target) => target.agentId)).toEqual(["main", "ops"]);
|
||||
expect(report.totals).toMatchObject({
|
||||
archivedLegacyStoreFiles: 0,
|
||||
importedEntries: 2,
|
||||
archivedLegacyStoreFiles: 1,
|
||||
importedEntries: 3,
|
||||
importedTranscriptEvents: 2,
|
||||
legacyEntries: 2,
|
||||
sqliteEntries: 2,
|
||||
legacyEntries: 3,
|
||||
sqliteEntries: 3,
|
||||
});
|
||||
for (const [agentId, sessionId] of [
|
||||
["main", "main-会議"],
|
||||
["ops", "ops-session"],
|
||||
] as const) {
|
||||
const agentStorePath = path.join(stateDir, "agents", agentId, "sessions", "sessions.json");
|
||||
const readScope = { agentId, env, storePath: agentStorePath };
|
||||
const readScope = { agentId, env, storePath };
|
||||
expect(
|
||||
loadExactSessionEntry({
|
||||
...readScope,
|
||||
|
|
@ -387,26 +395,13 @@ describe("runDoctorSessionSqlite", () => {
|
|||
expect(
|
||||
loadExactSessionEntry({
|
||||
...readScope,
|
||||
sessionKey: "voice:ambiguous",
|
||||
}),
|
||||
).toBeUndefined();
|
||||
sessionKey: `agent:${agentId}:voice:opaque`,
|
||||
})?.entry.sessionId,
|
||||
).toBe(agentId === "main" ? "opaque-session" : undefined);
|
||||
}
|
||||
expect(fs.existsSync(storePath)).toBe(true);
|
||||
expect(fs.existsSync(path.join(sessionDir, "main-session.jsonl"))).toBe(true);
|
||||
expect(fs.existsSync(path.join(sessionDir, "ops-session.jsonl"))).toBe(true);
|
||||
|
||||
const owned = await runDoctorSessionSqlite({
|
||||
allAgents: true,
|
||||
cfg: {
|
||||
...cfg,
|
||||
agents: { ...cfg.agents, defaults: { sessionStore: { agentId: "main" } } },
|
||||
},
|
||||
env,
|
||||
mode: "import",
|
||||
});
|
||||
expect(owned.totals.archivedLegacyStoreFiles).toBe(1);
|
||||
expect(owned.totals.importedEntries).toBe(3);
|
||||
expect(fs.existsSync(storePath)).toBe(false);
|
||||
expect(fs.existsSync(path.join(sessionDir, "main-session.jsonl"))).toBe(false);
|
||||
expect(fs.existsSync(path.join(sessionDir, "ops-session.jsonl"))).toBe(false);
|
||||
});
|
||||
|
||||
it.each([true, false])(
|
||||
|
|
|
|||
|
|
@ -447,27 +447,6 @@ describe("doctor legacy state migrations", () => {
|
|||
);
|
||||
});
|
||||
|
||||
it("repairs canonical headerless legacy transcript paths", async () => {
|
||||
const root = makeDoctorStateDir();
|
||||
const legacyDir = path.join(root, "sessions");
|
||||
const targetDir = path.join(root, "agents", "main", "sessions");
|
||||
fs.mkdirSync(targetDir, { recursive: true });
|
||||
fs.writeFileSync(path.join(targetDir, "legacy.jsonl"), '{"role":"user"}\n', "utf8");
|
||||
writeJson5(path.join(targetDir, "sessions.json"), {
|
||||
"agent:main:main": {
|
||||
sessionId: "legacy",
|
||||
sessionFile: path.join(legacyDir, "legacy.jsonl"),
|
||||
updatedAt: 10,
|
||||
},
|
||||
});
|
||||
|
||||
const detected = await detectLegacyStateMigrations({
|
||||
cfg: {},
|
||||
env: { OPENCLAW_STATE_DIR: root } as NodeJS.ProcessEnv,
|
||||
});
|
||||
expect(detected.sessions.hasLegacy).toBe(true);
|
||||
});
|
||||
|
||||
it("migrates legacy ACP metadata from retired custom-root agent stores", async () => {
|
||||
const root = makeDoctorStateDir();
|
||||
const customRoot = makeDoctorStateDir();
|
||||
|
|
|
|||
|
|
@ -173,8 +173,6 @@ function createLegacyStateMigrationDetectionResult(params?: {
|
|||
hasLegacy: false,
|
||||
},
|
||||
sessions: {
|
||||
legacyDir: "/tmp/state/sessions",
|
||||
legacyStorePath: "/tmp/state/sessions/sessions.json",
|
||||
targetDir: "/tmp/state/agents/main/sessions",
|
||||
targetStorePath: "/tmp/state/agents/main/sessions/sessions.json",
|
||||
hasLegacy: params?.hasLegacySessions ?? false,
|
||||
|
|
|
|||
|
|
@ -18,7 +18,6 @@ export type DoctorConfigPreflightOptions = {
|
|||
invocationPurpose?: LegacyStateMigrationInvocationPurpose;
|
||||
migrateLegacyConfig?: boolean;
|
||||
repairPrefixedConfig?: boolean;
|
||||
recoverCorruptTargetStore?: boolean;
|
||||
invalidConfigNote?: string | false;
|
||||
observe?: boolean;
|
||||
measure?: ConfigSnapshotReadMeasure;
|
||||
|
|
|
|||
|
|
@ -25,7 +25,6 @@ import {
|
|||
import { AGENT_DATABASE_PREFLIGHT_CONCURRENCY } from "../../state/openclaw-database-preflight-agent-scheduler.js";
|
||||
import { runTasksWithConcurrency } from "../../utils/run-with-concurrency.js";
|
||||
import { cloneEnvWithPlatformSemantics } from "../config-env-vars.js";
|
||||
import { resolveStateDir } from "../paths.js";
|
||||
import type { OpenClawConfig } from "../types.openclaw.js";
|
||||
import { migrateLegacyMainSessionKeys } from "./legacy-main-session-migration.js";
|
||||
import {
|
||||
|
|
@ -45,7 +44,6 @@ import {
|
|||
import {
|
||||
resolveAllAgentSessionStoreTargetsSync,
|
||||
resolveConfiguredAgentDatabaseTargets,
|
||||
resolveSessionStoreTargets,
|
||||
} from "./targets.js";
|
||||
|
||||
export type SessionStartupMigrationLogger = Record<"info" | "warn", (message: string) => void>;
|
||||
|
|
@ -64,23 +62,8 @@ export function assertSessionStoreMigrationComplete(params: {
|
|||
).filter(
|
||||
(target) => !target.agentId || !readAgentDatabaseAdmissionRefusal(target.agentId, { env }),
|
||||
);
|
||||
const legacyRootStore = path.join(resolveStateDir(env), "sessions", "sessions.json");
|
||||
const legacyTargets = fs.existsSync(legacyRootStore)
|
||||
? resolveSessionStoreTargets(params.cfg, { allAgents: true }, readOptions).map((target) => ({
|
||||
agentId: target.agentId,
|
||||
sqlitePath: resolveSqliteTargetFromSessionStorePath(target.storePath, {
|
||||
agentId: target.agentId,
|
||||
...readOptions,
|
||||
}).path,
|
||||
storePath: legacyRootStore,
|
||||
}))
|
||||
: [];
|
||||
const sources: readonly { agentId?: string; storePath: string; sqlitePath?: string }[] = [
|
||||
...(legacyTargets.length > 0 ? legacyTargets : [{ storePath: legacyRootStore }]),
|
||||
...targets,
|
||||
];
|
||||
const sourcesByPath = new Map<string, Array<(typeof sources)[number]>>();
|
||||
for (const target of sources) {
|
||||
const sourcesByPath = new Map<string, typeof targets>();
|
||||
for (const target of targets) {
|
||||
const sourcePath = path.resolve(target.storePath);
|
||||
sourcesByPath.set(sourcePath, [...(sourcesByPath.get(sourcePath) ?? []), target]);
|
||||
}
|
||||
|
|
@ -107,7 +90,7 @@ export function assertSessionStoreMigrationComplete(params: {
|
|||
});
|
||||
const legacyStore = legacySources.find(([storePath, candidates]) => {
|
||||
type SourceOwner = {
|
||||
target: { agentId: string; storePath: string; sqlitePath?: string };
|
||||
target: { agentId: string; storePath: string };
|
||||
destination: string;
|
||||
retained: boolean;
|
||||
imported: boolean;
|
||||
|
|
@ -117,12 +100,10 @@ export function assertSessionStoreMigrationComplete(params: {
|
|||
if (!target.agentId) {
|
||||
return true;
|
||||
}
|
||||
const destination =
|
||||
target.sqlitePath ??
|
||||
resolveSqliteTargetFromSessionStorePath(target.storePath, {
|
||||
agentId: target.agentId,
|
||||
...readOptions,
|
||||
}).path;
|
||||
const destination = resolveSqliteTargetFromSessionStorePath(target.storePath, {
|
||||
agentId: target.agentId,
|
||||
...readOptions,
|
||||
}).path;
|
||||
const deletion =
|
||||
classifyDeletion?.(storePath, target.agentId) ??
|
||||
classifyDeletion?.(destination, target.agentId);
|
||||
|
|
|
|||
|
|
@ -1958,7 +1958,6 @@ describe("doctor health contributions", () => {
|
|||
detected,
|
||||
config: cfg,
|
||||
legacySessionSurfaces,
|
||||
recoverCorruptTargetStore: false,
|
||||
});
|
||||
});
|
||||
|
||||
|
|
@ -2067,7 +2066,6 @@ describe("doctor health contributions", () => {
|
|||
config: cfg,
|
||||
doctorOnlyStateMigrations: true,
|
||||
legacySessionSurfaces,
|
||||
recoverCorruptTargetStore: true,
|
||||
});
|
||||
});
|
||||
|
||||
|
|
|
|||
|
|
@ -283,7 +283,6 @@ async function runLegacyStateHealth(ctx: DoctorHealthFlowContext): Promise<void>
|
|||
detected: legacyState,
|
||||
config: ctx.cfg,
|
||||
...(doctorOnlyStateMigrations ? { doctorOnlyStateMigrations: true } : {}),
|
||||
recoverCorruptTargetStore: ctx.options.repair === true || ctx.options.yes === true,
|
||||
legacySessionSurfaces,
|
||||
});
|
||||
recordDoctorHealthWarnings(
|
||||
|
|
|
|||
|
|
@ -175,10 +175,7 @@ function attachManagedImageRecordInDatabase(
|
|||
.where("attachment_id", "=", params.attachmentId)
|
||||
.where("session_key", "=", params.sessionKey),
|
||||
);
|
||||
if (!row) {
|
||||
return false;
|
||||
}
|
||||
if (row.cleanup_pending === 1) {
|
||||
if (!row || row.cleanup_pending === 1) {
|
||||
return false;
|
||||
}
|
||||
const current = managedImageRecordFromRow(row);
|
||||
|
|
|
|||
|
|
@ -122,8 +122,7 @@ async function runApprovalStoreOperation<T>(
|
|||
|
||||
function execute<Key extends Operation>(
|
||||
type: Key,
|
||||
input: OperatorApprovalWorkerOperations[Key]["input"],
|
||||
{ databaseOptions, assertCurrent, guard }: Options,
|
||||
{ databaseOptions, assertCurrent, guard, ...input }: Input<Key>,
|
||||
onCommitted?: (resolutionKey: string) => void,
|
||||
): Promise<OperatorApprovalWorkerOperations[Key]["output"]> {
|
||||
const context = captureOpenClawStateWorkerContext({
|
||||
|
|
@ -190,16 +189,13 @@ function execute<Key extends Operation>(
|
|||
}
|
||||
|
||||
export function insertOperatorApproval(params: Input<"operatorApprovals.insert">) {
|
||||
const { databaseOptions, assertCurrent, guard, ...input } = params;
|
||||
return execute("operatorApprovals.insert", input, { databaseOptions, assertCurrent, guard });
|
||||
return execute("operatorApprovals.insert", params);
|
||||
}
|
||||
export function getOperatorApprovalDetailed(params: Input<"operatorApprovals.get">) {
|
||||
const { databaseOptions, assertCurrent, guard, ...input } = params;
|
||||
return execute("operatorApprovals.get", input, { databaseOptions, assertCurrent, guard });
|
||||
return execute("operatorApprovals.get", params);
|
||||
}
|
||||
export function listPendingOperatorApprovals(params: Input<"operatorApprovals.pending"> = {}) {
|
||||
const { databaseOptions, assertCurrent, guard, ...input } = params;
|
||||
return execute("operatorApprovals.pending", input, { databaseOptions, assertCurrent, guard });
|
||||
return execute("operatorApprovals.pending", params);
|
||||
}
|
||||
export function resolveOperatorApproval(
|
||||
params: Input<"operatorApprovals.resolve"> & {
|
||||
|
|
@ -207,25 +203,17 @@ export function resolveOperatorApproval(
|
|||
onCommitted?: (resolutionKey: string) => void;
|
||||
},
|
||||
) {
|
||||
const { databaseOptions, assertCurrent, guard, onCommitted, ...input } = params;
|
||||
return execute(
|
||||
"operatorApprovals.resolve",
|
||||
input,
|
||||
{ databaseOptions, assertCurrent, guard },
|
||||
onCommitted,
|
||||
);
|
||||
const { onCommitted, ...input } = params;
|
||||
return execute("operatorApprovals.resolve", input, onCommitted);
|
||||
}
|
||||
export function forceDenyOperatorApproval(params: Input<"operatorApprovals.deny">) {
|
||||
const { databaseOptions, assertCurrent, guard, ...input } = params;
|
||||
return execute("operatorApprovals.deny", input, { databaseOptions, assertCurrent, guard });
|
||||
return execute("operatorApprovals.deny", params);
|
||||
}
|
||||
export function expireDueOperatorApprovals(params: Input<"operatorApprovals.expire">) {
|
||||
const { databaseOptions, assertCurrent, guard, ...input } = params;
|
||||
return execute("operatorApprovals.expire", input, { databaseOptions, assertCurrent, guard });
|
||||
return execute("operatorApprovals.expire", params);
|
||||
}
|
||||
export function consumeOperatorApprovalAllowOnce(params: Input<"operatorApprovals.consume">) {
|
||||
const { databaseOptions, assertCurrent, guard, ...input } = params;
|
||||
return execute("operatorApprovals.consume", input, { databaseOptions, assertCurrent, guard });
|
||||
return execute("operatorApprovals.consume", params);
|
||||
}
|
||||
|
||||
async function readApprovalStore<T>(
|
||||
|
|
@ -291,10 +279,5 @@ export function listCronStandingGrants(params: { limit?: number } & Options = {}
|
|||
}
|
||||
|
||||
export function revokeCronStandingGrant(params: Input<"operatorApprovals.revokeCronGrant">) {
|
||||
const { databaseOptions, assertCurrent, guard, ...input } = params;
|
||||
return execute("operatorApprovals.revokeCronGrant", input, {
|
||||
databaseOptions,
|
||||
assertCurrent,
|
||||
guard,
|
||||
});
|
||||
return execute("operatorApprovals.revokeCronGrant", params);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -192,12 +192,11 @@ export function createSessionRowProjectionArchive(params: {
|
|||
readPins.clear();
|
||||
pinCounts.clear();
|
||||
},
|
||||
describe(initial: records.Row | undefined) {
|
||||
if (initial?.entry?.archivedAt === undefined) {
|
||||
return initial;
|
||||
describe(row: records.Row | undefined) {
|
||||
if (row?.entry?.archivedAt === undefined) {
|
||||
return row;
|
||||
}
|
||||
const row = initial;
|
||||
if (records.ready(row) && row.entry.archivedAt !== undefined) {
|
||||
if (records.ready(row)) {
|
||||
const id = records.identity(row);
|
||||
materialized.delete(id);
|
||||
materialized.add(id);
|
||||
|
|
|
|||
|
|
@ -344,7 +344,7 @@ describe("runStartupSessionMigration", () => {
|
|||
);
|
||||
|
||||
it.each(["configured", "retired-root"] as const)(
|
||||
"preserves the %s legacy source and requires explicit Doctor import",
|
||||
"preserves the %s legacy source until a configured Doctor import",
|
||||
async (layout) => {
|
||||
const stateDir = fs.realpathSync.native(tempDirs.make("openclaw-legacy-session-startup-"));
|
||||
const env = { OPENCLAW_STATE_DIR: stateDir, OPENCLAW_PROFILE: "migration" };
|
||||
|
|
@ -362,6 +362,12 @@ describe("runStartupSessionMigration", () => {
|
|||
fs.mkdirSync(path.dirname(storePath), { recursive: true });
|
||||
fs.writeFileSync(storePath, original);
|
||||
|
||||
if (layout === "retired-root") {
|
||||
await expect(
|
||||
runStartupSessionMigration({ cfg, env, log: makeLog() }),
|
||||
).resolves.toBeUndefined();
|
||||
cfg.session = { store: storePath };
|
||||
}
|
||||
await expect(runStartupSessionMigration({ cfg, env, log: makeLog() })).rejects.toThrow(
|
||||
"openclaw --profile migration doctor --fix",
|
||||
);
|
||||
|
|
|
|||
|
|
@ -1,36 +0,0 @@
|
|||
import fs from "node:fs";
|
||||
import { StringDecoder } from "node:string_decoder";
|
||||
import { readFileWindowFullySync } from "@openclaw/fs-safe/advanced";
|
||||
|
||||
// StringDecoder preserves UTF-8 sequences split across chunks. Bound the scan
|
||||
// so a missing newline cannot read indefinitely.
|
||||
const HEADER_CHUNK_BYTES = 8192;
|
||||
const HEADER_MAX_CHARS = 1024 * 1024;
|
||||
|
||||
/** Reads the first newline-terminated line of a file, or undefined when there is none. */
|
||||
export function readFirstLineSync(filePath: string): string | undefined {
|
||||
const fd = fs.openSync(filePath, "r");
|
||||
try {
|
||||
const decoder = new StringDecoder("utf8");
|
||||
const chunk = Buffer.alloc(HEADER_CHUNK_BYTES);
|
||||
let carry = "";
|
||||
for (let position = 0; ;) {
|
||||
const bytesRead = readFileWindowFullySync(fd, chunk, position);
|
||||
if (bytesRead <= 0) {
|
||||
carry += decoder.end();
|
||||
return carry.length > 0 ? carry : undefined;
|
||||
}
|
||||
position += bytesRead;
|
||||
carry += decoder.write(chunk.subarray(0, bytesRead));
|
||||
const newline = carry.indexOf("\n");
|
||||
if (newline >= 0) {
|
||||
return carry.slice(0, newline);
|
||||
}
|
||||
if (carry.length > HEADER_MAX_CHARS) {
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
fs.closeSync(fd);
|
||||
}
|
||||
}
|
||||
|
|
@ -92,7 +92,7 @@ import {
|
|||
detectLegacyExecApprovals,
|
||||
migrateLegacyExecApprovals,
|
||||
} from "./state-migrations.exec-approvals.js";
|
||||
import { migrationFileExists, readSessionStoreJson5, safeReadDir } from "./state-migrations.fs.js";
|
||||
import { migrationFileExists, readSessionStoreJson5 } from "./state-migrations.fs.js";
|
||||
import {
|
||||
classifyLegacyOwnerFindings,
|
||||
tryResolveDoctorSessionMigrationAgentId,
|
||||
|
|
@ -178,7 +178,6 @@ import {
|
|||
listLegacySessionKeys,
|
||||
mergeSessionStoreAliasPlans,
|
||||
migrateLegacyAcpSessionMetadata,
|
||||
resolveStaleLegacySessionFile,
|
||||
resolveSessionStoreOwnership,
|
||||
type SessionStoreOwnership,
|
||||
} from "./state-migrations.session-store.js";
|
||||
|
|
@ -271,8 +270,6 @@ export async function detectLegacyStateMigrations(params: {
|
|||
const targetMainKey = normalizeOptionalString(params.cfg.session?.mainKey) ?? DEFAULT_MAIN_KEY;
|
||||
const targetScope = params.cfg.session?.scope;
|
||||
|
||||
const sessionsLegacyDir = path.join(stateDir, "sessions");
|
||||
const sessionsLegacyStorePath = path.join(sessionsLegacyDir, "sessions.json");
|
||||
const sessionsTargetDir = path.join(stateDir, "agents", targetAgentId, "sessions");
|
||||
const sessionsTargetStorePath = path.join(sessionsTargetDir, "sessions.json");
|
||||
const pluginConfig = params.pluginDoctorConfig ?? params.cfg;
|
||||
|
|
@ -317,11 +314,6 @@ export async function detectLegacyStateMigrations(params: {
|
|||
),
|
||||
};
|
||||
const { preserveForeignMainAliases } = sessionStoreOwnership;
|
||||
const hasLegacySessions =
|
||||
detectSessionFiles &&
|
||||
(migrationFileExists(sessionsLegacyStorePath) ||
|
||||
safeReadDir(sessionsLegacyDir).some((e) => e.isFile() && e.name.endsWith(".jsonl")));
|
||||
|
||||
const targetSessionParsed =
|
||||
detectSessionFiles && migrationFileExists(sessionsTargetStorePath)
|
||||
? readSessionStoreJson5(sessionsTargetStorePath)
|
||||
|
|
@ -341,18 +333,6 @@ export async function detectLegacyStateMigrations(params: {
|
|||
legacySessionSurfaces: legacySessionSurfaces.surfaces,
|
||||
})
|
||||
: [];
|
||||
const hasStaleSessionFiles =
|
||||
targetSessionParsed.ok &&
|
||||
Object.values(targetSessionParsed.store).some((entry) =>
|
||||
Boolean(
|
||||
resolveStaleLegacySessionFile({
|
||||
entry,
|
||||
legacyDir: sessionsLegacyDir,
|
||||
targetDir: sessionsTargetDir,
|
||||
}),
|
||||
),
|
||||
);
|
||||
|
||||
const targetAgentDir = migrationTarget?.dir;
|
||||
const targetAgentIdentity = targetAgentDir
|
||||
? resolveIdentityPathViaExistingAncestorSync(targetAgentDir)
|
||||
|
|
@ -593,13 +573,9 @@ export async function detectLegacyStateMigrations(params: {
|
|||
)
|
||||
).plans;
|
||||
|
||||
const sessionsHaveLegacy =
|
||||
Boolean(sessionMigrationAgentId) &&
|
||||
(hasLegacySessions || legacyKeys.length > 0 || hasStaleSessionFiles);
|
||||
const sessionsHaveLegacy = Boolean(sessionMigrationAgentId) && legacyKeys.length > 0;
|
||||
const agentDirHasLegacy = Boolean(migrationAgentId) && hasLegacyAgentDir;
|
||||
const deferredSessions =
|
||||
!sessionMigrationAgentId &&
|
||||
(hasLegacySessions || legacyKeys.length > 0 || hasStaleSessionFiles);
|
||||
const deferredSessions = !sessionMigrationAgentId && legacyKeys.length > 0;
|
||||
const deferredAgentDir = !migrationAgentId && hasLegacyAgentDir;
|
||||
const ownerFindings = classifyLegacyOwnerFindings({
|
||||
requiredWarnings: [...pluginPlanWarnings, ...legacySessionSurfaces.failures],
|
||||
|
|
@ -611,15 +587,9 @@ export async function detectLegacyStateMigrations(params: {
|
|||
doctorOnlyStateMigrations: params.doctorOnlyStateMigrations,
|
||||
});
|
||||
const preview: string[] = [];
|
||||
if (sessionsHaveLegacy && hasLegacySessions) {
|
||||
preview.push(`- Sessions: ${sessionsLegacyDir} → ${sessionsTargetDir}`);
|
||||
}
|
||||
if (sessionsHaveLegacy && legacyKeys.length > 0) {
|
||||
preview.push(`- Sessions: canonicalize legacy keys in ${sessionsTargetStorePath}`);
|
||||
}
|
||||
if (sessionsHaveLegacy && hasStaleSessionFiles) {
|
||||
preview.push(`- Sessions: repair migrated transcript paths in ${sessionsTargetStorePath}`);
|
||||
}
|
||||
if (agentDirHasLegacy) {
|
||||
preview.push(
|
||||
...legacyAgentSources.map(({ legacyDir }) => `- Agent dir: ${legacyDir} → ${targetAgentDir}`),
|
||||
|
|
@ -733,8 +703,6 @@ export async function detectLegacyStateMigrations(params: {
|
|||
oauthDir,
|
||||
pluginSessionStoreAgentIds,
|
||||
sessions: {
|
||||
legacyDir: sessionsLegacyDir,
|
||||
legacyStorePath: sessionsLegacyStorePath,
|
||||
targetDir: sessionsTargetDir,
|
||||
targetStorePath: sessionsTargetStorePath,
|
||||
hasLegacy: sessionsHaveLegacy,
|
||||
|
|
@ -1113,7 +1081,6 @@ type LegacyStateMigrationExecutionPlan = {
|
|||
agentDatabaseEndpoints?: LegacyStateMigrationEndpoint[];
|
||||
legacySessionStoreEndpoints?: LegacyStateMigrationEndpoint[];
|
||||
legacySessionStoreRefusal?: PreparedLegacyStateMigrationStep["refusal"];
|
||||
recoverCorruptTargetStore?: boolean;
|
||||
skipAgentScopedMigrations?: boolean;
|
||||
pluginStateMigrationInventory?: PluginDoctorStateMigrationInventory;
|
||||
deferPostSessionPluginMigrations?: boolean;
|
||||
|
|
@ -1330,7 +1297,7 @@ function buildLegacyStateMigrationSteps(
|
|||
pluginMigrationTargets,
|
||||
],
|
||||
sessions: [
|
||||
pathEndpoints(detected.sessions.legacyDir, detected.sessions.legacyStorePath),
|
||||
pathEndpoints(detected.sessions.targetDir, detected.sessions.targetStorePath),
|
||||
detected.sessions.hasLegacy,
|
||||
pathEndpoints(detected.sessions.targetDir, detected.sessions.targetStorePath),
|
||||
],
|
||||
|
|
@ -1603,10 +1570,9 @@ function buildLegacyStateMigrationSteps(
|
|||
if (repairSessionFiles) {
|
||||
finalSteps.push(
|
||||
finalStep("sessions", () =>
|
||||
migrateLegacySessions(detected, now, {
|
||||
migrateLegacySessions(detected, {
|
||||
cfg: params.sessionConfig ?? params.config,
|
||||
env,
|
||||
recoverCorruptTargetStore: params.recoverCorruptTargetStore,
|
||||
legacySessionSurfaces: params.legacySessionSurfaces,
|
||||
}),
|
||||
),
|
||||
|
|
@ -2456,7 +2422,6 @@ export async function runLegacyStateMigrations(params: {
|
|||
config?: OpenClawConfig;
|
||||
env?: NodeJS.ProcessEnv;
|
||||
now?: () => number;
|
||||
recoverCorruptTargetStore?: boolean;
|
||||
doctorOnlyStateMigrations?: boolean;
|
||||
onStepReceipt?: (receipt: LegacyStateMigrationStepReceipt) => void;
|
||||
legacySessionSurfaces: PreparedLegacySessionSurfaces;
|
||||
|
|
@ -2477,7 +2442,6 @@ export async function runLegacyStateMigrations(params: {
|
|||
config,
|
||||
env,
|
||||
now: params.now,
|
||||
recoverCorruptTargetStore: params.recoverCorruptTargetStore,
|
||||
pluginStateMigrationInventory,
|
||||
// The health contribution's later session repair consumes preflight's handoff.
|
||||
// This migration pass does not own that separate phase.
|
||||
|
|
@ -2551,7 +2515,6 @@ export async function autoMigrateLegacyState(params: {
|
|||
homedir?: () => string;
|
||||
log?: MigrationLogger;
|
||||
now?: () => number;
|
||||
recoverCorruptTargetStore?: boolean;
|
||||
doctorOnlyStateMigrations?: boolean;
|
||||
legacySessionSurfaces?: PreparedLegacySessionSurfaces;
|
||||
onStepReceipt?: (receipt: LegacyStateMigrationStepReceipt) => void;
|
||||
|
|
@ -2808,7 +2771,6 @@ async function executeLegacyStateMigrations(
|
|||
})),
|
||||
legacySessionStoreEndpoints: discoveredSessionStores.endpoints,
|
||||
legacySessionStoreRefusal,
|
||||
recoverCorruptTargetStore: params.recoverCorruptTargetStore,
|
||||
skipAgentScopedMigrations: hasCustomAgentDir,
|
||||
pluginStateMigrationInventory,
|
||||
legacySessionSurfaces,
|
||||
|
|
|
|||
|
|
@ -7,7 +7,6 @@ import {
|
|||
loadLegacySessionStore,
|
||||
saveLegacySessionStore,
|
||||
} from "./state-migrations.legacy-session-store.js";
|
||||
import { resolveStaleLegacySessionFile } from "./state-migrations.session-store.js";
|
||||
|
||||
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
|
||||
const DAY_MS = 24 * 60 * 60 * 1000;
|
||||
|
|
@ -180,26 +179,3 @@ it("normalizes compatibility writes before persistence", async () => {
|
|||
expect(persisted[MAIN_KEY]?.skillsSnapshot).not.toHaveProperty("resolvedSkills");
|
||||
expectNormalized(loadLegacySessionStore(storePath), "slack");
|
||||
});
|
||||
|
||||
it("repairs a stale session file whose header straddles the read chunk boundary", async () => {
|
||||
const sessionId = "sess-boundary-1";
|
||||
const legacyDir = path.join(root, "legacy-sessions");
|
||||
const targetDir = path.join(root, "sessions");
|
||||
await fs.mkdir(legacyDir);
|
||||
await fs.mkdir(targetDir);
|
||||
const legacySessionFile = path.join(legacyDir, `${sessionId}.jsonl`);
|
||||
const targetSessionFile = path.join(targetDir, `${sessionId}.jsonl`);
|
||||
// The three-byte character begins at 8191, splitting it across read chunks.
|
||||
const prefix = Buffer.from(`{"type":"session","id":"${sessionId}","pad":"`);
|
||||
await fs.writeFile(
|
||||
targetSessionFile,
|
||||
Buffer.concat([prefix, Buffer.from("a".repeat(8191 - prefix.length)), Buffer.from('中"}\n')]),
|
||||
);
|
||||
expect(
|
||||
resolveStaleLegacySessionFile({
|
||||
entry: { sessionId, sessionFile: legacySessionFile },
|
||||
legacyDir,
|
||||
targetDir,
|
||||
}),
|
||||
).toBe(targetSessionFile);
|
||||
});
|
||||
|
|
|
|||
|
|
@ -3,16 +3,13 @@ import fs from "node:fs";
|
|||
import path from "node:path";
|
||||
import { resolveInstallAgentDir } from "../agents/install-agent-dir.js";
|
||||
import type { SessionEntry } from "../config/sessions.js";
|
||||
import { resolveSessionStoreTargets } from "../config/sessions/targets.js";
|
||||
import type { OpenClawConfig } from "../config/types.openclaw.js";
|
||||
import { buildAgentMainSessionKey } from "../routing/session-key.js";
|
||||
import { readExistingAgentSchemaMeta } from "../state/openclaw-agent-db-schema-helpers.js";
|
||||
import { readDeferredPluginMigrations } from "./deferred-plugin-migrations.js";
|
||||
import { preserveDeferredPluginSessionSource } from "./deferred-plugin-session-sources.js";
|
||||
import { isErrno } from "./errors.js";
|
||||
import { openNodeSqliteDatabase } from "./node-sqlite.js";
|
||||
import { isPathInside } from "./path-guards.js";
|
||||
import { resolveTargetSqlitePath } from "./session-sqlite-migration-readers.js";
|
||||
import { resolveSqliteDatabaseFilePaths, SQLITE_SIDECAR_SUFFIXES } from "./sqlite-files.js";
|
||||
import { quoteSqliteIdentifier } from "./sqlite-schema-sql.js";
|
||||
import {
|
||||
|
|
@ -23,7 +20,6 @@ import {
|
|||
ensureMigrationDir,
|
||||
migrationFileExists,
|
||||
readSessionStoreJson5,
|
||||
safeReadDir,
|
||||
type SessionEntryLike,
|
||||
} from "./state-migrations.fs.js";
|
||||
import {
|
||||
|
|
@ -31,11 +27,8 @@ import {
|
|||
canonicalizeSessionStore,
|
||||
distinctSessionStoreAliasWarning,
|
||||
isAmbiguousSharedStoreKey,
|
||||
selectNewerSessionEntry,
|
||||
normalizeSessionEntry,
|
||||
pickLatestLegacyDirectEntry,
|
||||
removeDirIfEmpty,
|
||||
resolveStaleLegacySessionFile,
|
||||
saveSessionStoreStrict,
|
||||
unresolvedSessionStoreIdentityWarning,
|
||||
} from "./state-migrations.session-store.js";
|
||||
|
|
@ -127,21 +120,16 @@ export function inspectLegacyAgentDir(
|
|||
}
|
||||
}
|
||||
|
||||
function normalizeMergedSessionStore(
|
||||
merged: Record<string, SessionEntryLike>,
|
||||
protectedKeys: ReadonlySet<string>,
|
||||
): {
|
||||
function normalizeTargetSessionStore(target: Record<string, SessionEntryLike>): {
|
||||
store: Record<string, SessionEntry>;
|
||||
rejectedProtectedKeyCount: number;
|
||||
} {
|
||||
const store = Object.create(null) as Record<string, SessionEntry>;
|
||||
let rejectedProtectedKeyCount = 0;
|
||||
for (const [key, entry] of Object.entries(merged)) {
|
||||
for (const [key, entry] of Object.entries(target)) {
|
||||
const normalizedEntry = normalizeSessionEntry(entry, key);
|
||||
if (!normalizedEntry) {
|
||||
if (protectedKeys.has(key)) {
|
||||
rejectedProtectedKeyCount++;
|
||||
}
|
||||
rejectedProtectedKeyCount++;
|
||||
continue;
|
||||
}
|
||||
store[key] = normalizedEntry;
|
||||
|
|
@ -151,43 +139,28 @@ function normalizeMergedSessionStore(
|
|||
|
||||
export async function migrateLegacySessions(
|
||||
detected: LegacyStateDetection,
|
||||
now: () => number,
|
||||
options: {
|
||||
cfg: OpenClawConfig;
|
||||
env: NodeJS.ProcessEnv;
|
||||
recoverCorruptTargetStore?: boolean;
|
||||
legacySessionSurfaces: PreparedLegacySessionSurfaces;
|
||||
},
|
||||
): Promise<MigrationMessages> {
|
||||
const changes: string[] = [];
|
||||
const warnings: string[] = [];
|
||||
const recoverableWarnings: string[] = [];
|
||||
if (!detected.sessions.hasLegacy) {
|
||||
return { changes, warnings };
|
||||
}
|
||||
if (options.legacySessionSurfaces.failures.length > 0) {
|
||||
return {
|
||||
changes,
|
||||
warnings: [...options.legacySessionSurfaces.failures],
|
||||
};
|
||||
return { changes, warnings: [...options.legacySessionSurfaces.failures] };
|
||||
}
|
||||
const env = { ...options.env, OPENCLAW_STATE_DIR: detected.stateDir };
|
||||
const pending = readDeferredPluginMigrations({ env });
|
||||
// The shared legacy index imports into configured stores, not a database beside the index.
|
||||
const legacyTargets = resolveSessionStoreTargets(options.cfg, { allAgents: true }, { env }).map(
|
||||
(target) => ({
|
||||
agentId: target.agentId,
|
||||
sqlitePath: resolveTargetSqlitePath(target, env),
|
||||
storePath: detected.sessions.legacyStorePath,
|
||||
}),
|
||||
);
|
||||
if (
|
||||
[
|
||||
{ agentId: detected.targetAgentId, storePath: detected.sessions.targetStorePath },
|
||||
...legacyTargets,
|
||||
].some((target) =>
|
||||
preserveDeferredPluginSessionSource({ cfg: options.cfg, env, target, pending }),
|
||||
)
|
||||
preserveDeferredPluginSessionSource({
|
||||
cfg: options.cfg,
|
||||
env,
|
||||
target: { agentId: detected.targetAgentId, storePath: detected.sessions.targetStorePath },
|
||||
pending: readDeferredPluginMigrations({ env }),
|
||||
})
|
||||
) {
|
||||
return {
|
||||
changes,
|
||||
|
|
@ -197,21 +170,17 @@ export async function migrateLegacySessions(
|
|||
],
|
||||
};
|
||||
}
|
||||
const legacyParsed = migrationFileExists(detected.sessions.legacyStorePath)
|
||||
? readSessionStoreJson5(detected.sessions.legacyStorePath)
|
||||
: { store: {}, ok: true };
|
||||
if (!legacyParsed.ok) {
|
||||
warnings.push(
|
||||
`Legacy sessions store unreadable; left in place at ${detected.sessions.legacyStorePath}`,
|
||||
);
|
||||
return { changes, warnings };
|
||||
}
|
||||
|
||||
ensureMigrationDir(detected.sessions.targetDir);
|
||||
const targetParsed = migrationFileExists(detected.sessions.targetStorePath)
|
||||
? readSessionStoreJson5(detected.sessions.targetStorePath)
|
||||
: { store: {}, ok: true };
|
||||
const legacyStore = legacyParsed.store;
|
||||
if (!targetParsed.ok) {
|
||||
return {
|
||||
changes,
|
||||
warnings: [
|
||||
`Target sessions store unreadable; left untouched at ${detected.sessions.targetStorePath}. Repair the index, then rerun openclaw doctor --fix.`,
|
||||
],
|
||||
};
|
||||
}
|
||||
const targetStore = targetParsed.store;
|
||||
if (detected.sessions.targetStoreAliases.hasUnresolvedIdentity) {
|
||||
warnings.push(
|
||||
|
|
@ -229,22 +198,18 @@ export async function migrateLegacySessions(
|
|||
return { changes, warnings };
|
||||
}
|
||||
|
||||
const ambiguousAliasedKeys = new Set(
|
||||
[...Object.keys(targetStore), ...Object.keys(legacyStore)].filter(
|
||||
(key) =>
|
||||
isAmbiguousSharedStoreKey(key, detected.targetMainKey, detected.targetScope) ||
|
||||
(detected.sessions.preserveForeignMainAliases &&
|
||||
isLegacyDefaultMainAliasKey(key, detected.targetMainKey)),
|
||||
),
|
||||
const ambiguousAliasedKeys = Object.keys(targetStore).filter(
|
||||
(key) =>
|
||||
isAmbiguousSharedStoreKey(key, detected.targetMainKey, detected.targetScope) ||
|
||||
(detected.sessions.preserveForeignMainAliases &&
|
||||
isLegacyDefaultMainAliasKey(key, detected.targetMainKey)),
|
||||
);
|
||||
// Atomic replacement separates filesystem aliases. Defer the whole merge so
|
||||
// a later startup cannot treat each pathname as a different session owner.
|
||||
if (detected.sessions.targetStoreAliases.hasDistinctAliases) {
|
||||
warnings.push(
|
||||
ambiguousAliasedKeys.size > 0
|
||||
ambiguousAliasedKeys.length > 0
|
||||
? aliasedSessionStoreMigrationWarning({
|
||||
subject: "migration of",
|
||||
count: ambiguousAliasedKeys.size,
|
||||
count: ambiguousAliasedKeys.length,
|
||||
storePath: detected.sessions.targetStorePath,
|
||||
})
|
||||
: distinctSessionStoreAliasWarning(
|
||||
|
|
@ -255,7 +220,7 @@ export async function migrateLegacySessions(
|
|||
return { changes, warnings };
|
||||
}
|
||||
|
||||
const canonicalizedTarget = canonicalizeSessionStore({
|
||||
const canonicalized = canonicalizeSessionStore({
|
||||
store: targetStore,
|
||||
agentId: detected.targetAgentId,
|
||||
mainKey: detected.targetMainKey,
|
||||
|
|
@ -266,163 +231,20 @@ export async function migrateLegacySessions(
|
|||
preserveForeignMainAliases: detected.sessions.preserveForeignMainAliases,
|
||||
legacySessionSurfaces: options.legacySessionSurfaces.surfaces,
|
||||
});
|
||||
const canonicalizedLegacy = canonicalizeSessionStore({
|
||||
store: legacyStore,
|
||||
agentId: detected.targetAgentId,
|
||||
mainKey: detected.targetMainKey,
|
||||
scope: detected.targetScope,
|
||||
preserveCanonicalAgentOwner: true,
|
||||
preserveForeignMainAliases: detected.sessions.preserveForeignMainAliases,
|
||||
legacySessionSurfaces: options.legacySessionSurfaces.surfaces,
|
||||
});
|
||||
const targetKeys = new Set(Object.keys(canonicalizedTarget.store));
|
||||
const preservedLegacyForeignMainAliasCount = detected.sessions.preserveForeignMainAliases
|
||||
? Object.keys(legacyStore).filter((key) =>
|
||||
isLegacyDefaultMainAliasKey(key, detected.targetMainKey),
|
||||
).length
|
||||
: 0;
|
||||
|
||||
let repairedStaleSessionFiles = false;
|
||||
for (const entry of Object.values(canonicalizedTarget.store)) {
|
||||
const targetSessionFile = resolveStaleLegacySessionFile({
|
||||
entry,
|
||||
legacyDir: detected.sessions.legacyDir,
|
||||
targetDir: detected.sessions.targetDir,
|
||||
});
|
||||
if (targetSessionFile) {
|
||||
entry.sessionFile = targetSessionFile;
|
||||
repairedStaleSessionFiles = true;
|
||||
}
|
||||
}
|
||||
|
||||
const merged = Object.create(null) as Record<string, SessionEntryLike>;
|
||||
for (const [key, entry] of Object.entries(canonicalizedTarget.store)) {
|
||||
merged[key] = entry;
|
||||
}
|
||||
for (const [key, entry] of Object.entries(canonicalizedLegacy.store)) {
|
||||
merged[key] = selectNewerSessionEntry({
|
||||
existing: merged[key],
|
||||
incoming: entry,
|
||||
preferIncomingOnTie: false,
|
||||
});
|
||||
}
|
||||
|
||||
const mainKey = buildAgentMainSessionKey({
|
||||
agentId: detected.targetAgentId,
|
||||
mainKey: detected.targetMainKey,
|
||||
});
|
||||
let migratedDirectChatKey: string | undefined;
|
||||
if (!merged[mainKey]) {
|
||||
const latest = pickLatestLegacyDirectEntry(legacyStore, options.legacySessionSurfaces.surfaces);
|
||||
if (latest?.sessionId) {
|
||||
merged[mainKey] = latest;
|
||||
migratedDirectChatKey = mainKey;
|
||||
}
|
||||
}
|
||||
|
||||
const targetExists = migrationFileExists(detected.sessions.targetStorePath);
|
||||
let targetReadable = !targetExists || targetParsed.ok;
|
||||
if (!targetReadable) {
|
||||
if (options.recoverCorruptTargetStore) {
|
||||
const archivedTargetPath = `${detected.sessions.targetStorePath}.corrupt-${now()}`;
|
||||
try {
|
||||
fs.renameSync(detected.sessions.targetStorePath, archivedTargetPath);
|
||||
changes.push(`Archived corrupt target sessions store → ${archivedTargetPath}`);
|
||||
targetReadable = true;
|
||||
} catch (err) {
|
||||
warnings.push(
|
||||
`Target sessions store unreadable; failed to archive ${detected.sessions.targetStorePath}: ${String(err)}`,
|
||||
);
|
||||
}
|
||||
} else {
|
||||
warnings.push(
|
||||
`Target sessions store unreadable; left untouched to avoid overwriting at ${detected.sessions.targetStorePath}. Run openclaw doctor --fix to archive it and retry the legacy merge.`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
if (
|
||||
targetReadable &&
|
||||
(Object.keys(legacyStore).length > 0 || Object.keys(targetStore).length > 0)
|
||||
) {
|
||||
const normalized = normalizeMergedSessionStore(merged, targetKeys);
|
||||
if (normalized.rejectedProtectedKeyCount > 0) {
|
||||
warnings.push(
|
||||
`Refused legacy session migration because normalization rejected ${normalized.rejectedProtectedKeyCount} existing target session ${normalized.rejectedProtectedKeyCount === 1 ? "key" : "keys"}; left ${detected.sessions.targetStorePath} and ${detected.sessions.legacyStorePath} in place. Repair the conflicting rows, then rerun openclaw doctor --fix.`,
|
||||
);
|
||||
return { changes, warnings };
|
||||
}
|
||||
await saveSessionStoreStrict(detected.sessions.targetStorePath, normalized.store);
|
||||
if (migratedDirectChatKey) {
|
||||
changes.push(`Migrated latest direct-chat session → ${migratedDirectChatKey}`);
|
||||
}
|
||||
changes.push(`Merged sessions store → ${detected.sessions.targetStorePath}`);
|
||||
if (preservedLegacyForeignMainAliasCount > 0) {
|
||||
recoverableWarnings.push(
|
||||
`Preserved ${preservedLegacyForeignMainAliasCount} ambiguous session key(s) while importing legacy sessions into ${detected.sessions.targetStorePath}`,
|
||||
);
|
||||
}
|
||||
if (canonicalizedTarget.legacyKeys.length > 0) {
|
||||
changes.push(`Canonicalized ${canonicalizedTarget.legacyKeys.length} legacy session key(s)`);
|
||||
}
|
||||
if (repairedStaleSessionFiles) {
|
||||
changes.push("Repaired migrated session transcript paths");
|
||||
}
|
||||
}
|
||||
|
||||
if (!targetReadable) {
|
||||
const normalized = normalizeTargetSessionStore(canonicalized.store);
|
||||
if (normalized.rejectedProtectedKeyCount > 0) {
|
||||
warnings.push(
|
||||
`Refused legacy session migration because normalization rejected ${normalized.rejectedProtectedKeyCount} existing target session ${normalized.rejectedProtectedKeyCount === 1 ? "key" : "keys"}; left ${detected.sessions.targetStorePath} in place. Repair the conflicting rows, then rerun openclaw doctor --fix.`,
|
||||
);
|
||||
return { changes, warnings };
|
||||
}
|
||||
|
||||
const entries = safeReadDir(detected.sessions.legacyDir);
|
||||
for (const entry of entries) {
|
||||
if (!entry.isFile()) {
|
||||
continue;
|
||||
}
|
||||
if (entry.name === "sessions.json") {
|
||||
continue;
|
||||
}
|
||||
const from = path.join(detected.sessions.legacyDir, entry.name);
|
||||
let to = path.join(detected.sessions.targetDir, entry.name);
|
||||
if (migrationFileExists(to)) {
|
||||
const parsed = path.parse(entry.name);
|
||||
to = path.join(detected.sessions.targetDir, `${parsed.name}.legacy-${now()}${parsed.ext}`);
|
||||
}
|
||||
try {
|
||||
fs.renameSync(from, to);
|
||||
changes.push(`Moved ${entry.name} → agents/${detected.targetAgentId}/sessions`);
|
||||
} catch (err) {
|
||||
warnings.push(`Failed moving ${from}: ${String(err)}`);
|
||||
if (Object.keys(targetStore).length > 0) {
|
||||
await saveSessionStoreStrict(detected.sessions.targetStorePath, normalized.store);
|
||||
if (canonicalized.legacyKeys.length > 0) {
|
||||
changes.push(`Canonicalized ${canonicalized.legacyKeys.length} legacy session key(s)`);
|
||||
}
|
||||
}
|
||||
|
||||
try {
|
||||
if (migrationFileExists(detected.sessions.legacyStorePath)) {
|
||||
fs.rmSync(detected.sessions.legacyStorePath, { force: true });
|
||||
}
|
||||
} catch {
|
||||
// ignore
|
||||
}
|
||||
|
||||
removeDirIfEmpty(detected.sessions.legacyDir);
|
||||
const legacyLeft = safeReadDir(detected.sessions.legacyDir).filter((e) => e.isFile());
|
||||
if (legacyLeft.length > 0) {
|
||||
const backupDir = `${detected.sessions.legacyDir}.legacy-${now()}`;
|
||||
try {
|
||||
fs.renameSync(detected.sessions.legacyDir, backupDir);
|
||||
warnings.push(`Left legacy sessions at ${backupDir}`);
|
||||
} catch {
|
||||
// ignore
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
changes,
|
||||
warnings: [...warnings, ...recoverableWarnings],
|
||||
...(warnings.length === 0 && recoverableWarnings.length > 0 && changes.length > 0
|
||||
? { warningDisposition: "recoverable" as const }
|
||||
: {}),
|
||||
};
|
||||
return { changes, warnings };
|
||||
}
|
||||
|
||||
type SqliteFamilyPlan = {
|
||||
|
|
|
|||
|
|
@ -11,7 +11,6 @@ import {
|
|||
type ChannelIngressQueue,
|
||||
} from "../channels/message/ingress-queue.js";
|
||||
import { importLegacyChannelIngressEntries } from "../channels/message/ingress-queue.migration.js";
|
||||
import { resolveStateDir } from "../config/paths.js";
|
||||
import { resolveSessionStorePathCore } from "../config/sessions/paths.js";
|
||||
import { readSessionIdentityEvidenceBatch } from "../config/sessions/session-accessor.js";
|
||||
import { resolveSqliteTargetFromSessionStorePath } from "../config/sessions/session-sqlite-target.js";
|
||||
|
|
@ -63,25 +62,20 @@ function hasUnimportedSessionIdentity(params: {
|
|||
env: params.env,
|
||||
});
|
||||
const defaultStore = resolveSessionStorePathCore(undefined, { agentId, env: params.env });
|
||||
const legacyRootStore = path.join(resolveStateDir(params.env), "sessions", "sessions.json");
|
||||
const sources = new Map([
|
||||
[configuredStore, configuredStore],
|
||||
[defaultStore, defaultStore],
|
||||
[legacyRootStore, configuredStore],
|
||||
]);
|
||||
const sources = new Set([configuredStore, defaultStore]);
|
||||
let importedIdentity = false;
|
||||
let unimportedIdentity = false;
|
||||
for (const [storePath, destination] of sources) {
|
||||
for (const storePath of sources) {
|
||||
if (storePath.endsWith(".sqlite")) {
|
||||
continue;
|
||||
}
|
||||
const key = `${agentId}\0${storePath}\0${destination}`;
|
||||
const key = `${agentId}\0${storePath}`;
|
||||
let sourceEvidence = params.cache.get(key);
|
||||
if (sourceEvidence === undefined) {
|
||||
const before = fs.statSync(storePath, { throwIfNoEntry: false, bigint: true });
|
||||
sourceEvidence = { imported: false, sessionIds: new Set() };
|
||||
if (before) {
|
||||
const sqlitePath = resolveSqliteTargetFromSessionStorePath(destination, {
|
||||
const sqlitePath = resolveSqliteTargetFromSessionStorePath(storePath, {
|
||||
agentId,
|
||||
env: params.env,
|
||||
}).path;
|
||||
|
|
@ -90,7 +84,6 @@ function hasUnimportedSessionIdentity(params: {
|
|||
target: {
|
||||
agentId,
|
||||
storePath,
|
||||
...(storePath === legacyRootStore ? { sqlitePath } : {}),
|
||||
},
|
||||
sqlitePath,
|
||||
env: params.env,
|
||||
|
|
|
|||
|
|
@ -576,29 +576,6 @@ describe("Doctor with a deleted agent database", () => {
|
|||
assertSessionStoreMigrationComplete({ cfg, env, operation: "doctor" }),
|
||||
).not.toThrow();
|
||||
expect(fs.readFileSync(retainedStore, "utf8")).toBe(retainedStoreBytes);
|
||||
if (registered && location === "default") {
|
||||
const globalStore = path.join(stateDir, "sessions", "sessions.json");
|
||||
fs.mkdirSync(path.dirname(globalStore), { recursive: true });
|
||||
fs.writeFileSync(globalStore, retainedStoreBytes);
|
||||
const assertGlobalReady = () =>
|
||||
assertSessionStoreMigrationComplete({
|
||||
cfg: { agents: { ownership: "explicit", entries: { retired: {} } } },
|
||||
env,
|
||||
operation: "doctor",
|
||||
});
|
||||
expect(assertGlobalReady).not.toThrow();
|
||||
expect(fs.readFileSync(globalStore, "utf8")).toBe(retainedStoreBytes);
|
||||
fs.writeFileSync(
|
||||
globalStore,
|
||||
JSON.stringify({
|
||||
"agent:retired:legacy": { sessionId: "retired-legacy", updatedAt: 1 },
|
||||
"agent:unassigned:legacy": { sessionId: "unknown-legacy", updatedAt: 1 },
|
||||
}),
|
||||
);
|
||||
expect(assertGlobalReady).toThrow("Legacy session store requires migration");
|
||||
fs.writeFileSync(globalStore, "{}");
|
||||
expect(assertGlobalReady).toThrow("Legacy session store requires migration");
|
||||
}
|
||||
expect(fs.readFileSync(retainedPath).equals(before)).toBe(true);
|
||||
}
|
||||
expect(() => openOpenClawAgentDatabase({ agentId: "retired", env })).toThrow(
|
||||
|
|
|
|||
|
|
@ -43,7 +43,6 @@ import {
|
|||
prepareDeferredPluginSessionImportReader,
|
||||
preserveDeferredPluginSessionSource,
|
||||
} from "./deferred-plugin-session-sources.js";
|
||||
import { readFirstLineSync } from "./first-line-read.js";
|
||||
import { expandHomePrefix } from "./home-dir.js";
|
||||
import { importLegacyAcpSessionMetadata } from "./state-migrations.acp-session-metadata.js";
|
||||
import {
|
||||
|
|
@ -61,7 +60,6 @@ import {
|
|||
} from "./state-migrations.session-store-paths.js";
|
||||
import {
|
||||
isLegacyDefaultMainAliasKey,
|
||||
isLegacyGroupKey,
|
||||
resolveCanonicalAgentSessionOwner,
|
||||
isSurfaceGroupKey,
|
||||
type PreparedLegacySessionSurfaces,
|
||||
|
|
@ -200,42 +198,6 @@ function canonicalizeSessionKeyForAgent(params: {
|
|||
return normalizeSessionKeyPreservingOpaquePeerIds(`agent:${agentId}:${raw}`);
|
||||
}
|
||||
|
||||
export function pickLatestLegacyDirectEntry(
|
||||
store: Record<string, SessionEntryLike>,
|
||||
legacySessionSurfaces: PreparedLegacySessionSurfaces["surfaces"] = [],
|
||||
): SessionEntryLike | null {
|
||||
let best: SessionEntryLike | null = null;
|
||||
let bestUpdated = -1;
|
||||
for (const [key, entry] of Object.entries(store)) {
|
||||
if (!entry || typeof entry !== "object") {
|
||||
continue;
|
||||
}
|
||||
const normalized = key.trim();
|
||||
if (!normalized) {
|
||||
continue;
|
||||
}
|
||||
const normalizedLower = normalizeLowercaseStringOrEmpty(normalized);
|
||||
if (normalizedLower === "global") {
|
||||
continue;
|
||||
}
|
||||
if (normalizedLower.startsWith("agent:")) {
|
||||
continue;
|
||||
}
|
||||
if (normalizedLower.startsWith("subagent:")) {
|
||||
continue;
|
||||
}
|
||||
if (isLegacyGroupKey(normalized, legacySessionSurfaces) || isSurfaceGroupKey(normalized)) {
|
||||
continue;
|
||||
}
|
||||
const updatedAt = typeof entry.updatedAt === "number" ? entry.updatedAt : 0;
|
||||
if (updatedAt > bestUpdated) {
|
||||
bestUpdated = updatedAt;
|
||||
best = entry;
|
||||
}
|
||||
}
|
||||
return best;
|
||||
}
|
||||
|
||||
export function normalizeSessionEntry(
|
||||
entry: SessionEntryLike,
|
||||
sessionKey?: string,
|
||||
|
|
@ -252,25 +214,6 @@ export function normalizeSessionEntry(
|
|||
return normalized;
|
||||
}
|
||||
|
||||
export function selectNewerSessionEntry(params: {
|
||||
existing: SessionEntryLike | undefined;
|
||||
incoming: SessionEntryLike;
|
||||
preferIncomingOnTie?: boolean;
|
||||
}): SessionEntryLike {
|
||||
if (!params.existing) {
|
||||
return params.incoming;
|
||||
}
|
||||
const existingUpdated = asFiniteNumber(params.existing.updatedAt) ?? 0;
|
||||
const incomingUpdated = asFiniteNumber(params.incoming.updatedAt) ?? 0;
|
||||
if (incomingUpdated > existingUpdated) {
|
||||
return params.incoming;
|
||||
}
|
||||
if (incomingUpdated < existingUpdated) {
|
||||
return params.existing;
|
||||
}
|
||||
return params.preferIncomingOnTie ? params.incoming : params.existing;
|
||||
}
|
||||
|
||||
export function canonicalizeSessionStore(params: {
|
||||
store: Record<string, SessionEntryLike>;
|
||||
agentId: string;
|
||||
|
|
@ -359,72 +302,6 @@ export function distinctSessionStoreAliasWarning(subject: string, storePath: str
|
|||
return `Deferred ${subject} in aliased store ${storePath}; atomic replacement cannot update distinct filesystem aliases as one operation. Remove filesystem aliases or configure one canonical session.store path, then rerun openclaw doctor --fix`;
|
||||
}
|
||||
|
||||
export function resolveStaleLegacySessionFile(params: {
|
||||
entry: unknown;
|
||||
legacyDir: string;
|
||||
targetDir: string;
|
||||
}): string | undefined {
|
||||
if (!params.entry || typeof params.entry !== "object" || Array.isArray(params.entry)) {
|
||||
return undefined;
|
||||
}
|
||||
const entry = params.entry as SessionEntryLike;
|
||||
const rawSessionFile = entry.sessionFile;
|
||||
if (typeof rawSessionFile !== "string") {
|
||||
return undefined;
|
||||
}
|
||||
const legacySessionFile = path.isAbsolute(rawSessionFile)
|
||||
? path.resolve(rawSessionFile)
|
||||
: path.resolve(params.legacyDir, rawSessionFile);
|
||||
const relative = path.relative(path.resolve(params.legacyDir), legacySessionFile);
|
||||
if (
|
||||
relative.startsWith("..") ||
|
||||
path.isAbsolute(relative) ||
|
||||
migrationFileExists(legacySessionFile)
|
||||
) {
|
||||
return undefined;
|
||||
}
|
||||
const legacyBackupHasTranscript = safeReadDir(path.dirname(params.legacyDir)).some(
|
||||
(dirent) =>
|
||||
dirent.isDirectory() &&
|
||||
dirent.name.startsWith(`${path.basename(params.legacyDir)}.legacy-`) &&
|
||||
migrationFileExists(
|
||||
path.join(path.dirname(params.legacyDir), dirent.name, path.basename(legacySessionFile)),
|
||||
),
|
||||
);
|
||||
if (legacyBackupHasTranscript) {
|
||||
return undefined;
|
||||
}
|
||||
const parsed = path.parse(path.basename(legacySessionFile));
|
||||
const hasCollisionRename = safeReadDir(params.targetDir).some(
|
||||
(dirent) =>
|
||||
dirent.isFile() &&
|
||||
dirent.name.startsWith(`${parsed.name}.legacy-`) &&
|
||||
dirent.name.endsWith(parsed.ext),
|
||||
);
|
||||
if (hasCollisionRename) {
|
||||
return undefined;
|
||||
}
|
||||
const targetSessionFile = path.join(params.targetDir, path.basename(legacySessionFile));
|
||||
if (!migrationFileExists(targetSessionFile) || typeof entry.sessionId !== "string") {
|
||||
return undefined;
|
||||
}
|
||||
try {
|
||||
const firstLine = readFirstLineSync(targetSessionFile);
|
||||
const header = firstLine ? (JSON.parse(firstLine) as unknown) : undefined;
|
||||
if (!header || typeof header !== "object" || Array.isArray(header)) {
|
||||
return undefined;
|
||||
}
|
||||
if ((header as { type?: unknown }).type === "session") {
|
||||
return (header as { id?: unknown }).id === entry.sessionId ? targetSessionFile : undefined;
|
||||
}
|
||||
const canonicalFileName =
|
||||
path.basename(entry.sessionId) === entry.sessionId ? `${entry.sessionId}.jsonl` : undefined;
|
||||
return canonicalFileName === path.basename(targetSessionFile) ? targetSessionFile : undefined;
|
||||
} catch {
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
|
||||
function sessionStoreMayNeedCanonicalization(params: {
|
||||
store: Record<string, SessionEntryLike>;
|
||||
storeAgentIds: Iterable<string>;
|
||||
|
|
|
|||
|
|
@ -15,26 +15,6 @@ export function isSurfaceGroupKey(key: string): boolean {
|
|||
return key.includes(":group:") || key.includes(":channel:");
|
||||
}
|
||||
|
||||
export function isLegacyGroupKey(
|
||||
key: string,
|
||||
surfaces: PreparedLegacySessionSurfaces["surfaces"] = [],
|
||||
): boolean {
|
||||
const trimmed = key.trim();
|
||||
if (!trimmed) {
|
||||
return false;
|
||||
}
|
||||
const lower = normalizeLowercaseStringOrEmpty(trimmed);
|
||||
if (lower.startsWith("group:") || lower.startsWith("channel:")) {
|
||||
return true;
|
||||
}
|
||||
for (const surface of surfaces) {
|
||||
if (surface.isLegacyGroupSessionKey?.(trimmed)) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
export function isLegacyDefaultMainAliasKey(key: string, mainKey: string): boolean {
|
||||
const lower = normalizeLowercaseStringOrEmpty(key.trim());
|
||||
const canonicalMainKey = normalizeMainKey(mainKey);
|
||||
|
|
|
|||
|
|
@ -1383,16 +1383,13 @@ describe("state migrations", () => {
|
|||
expect(result.notices ?? []).not.toContain(
|
||||
"Deferred legacy agent/session migration: select an agent owner",
|
||||
);
|
||||
await expect(fs.readFile(legacyStorePath)).resolves.toEqual(legacyBytes);
|
||||
await expect(
|
||||
fs.readFile(
|
||||
path.join(stateDir, "agents", targetAgentId, "sessions", "legacy-session.jsonl"),
|
||||
"utf8",
|
||||
),
|
||||
fs.readFile(path.join(legacySessionsDir, "legacy-session.jsonl"), "utf8"),
|
||||
).resolves.toBe("{}\n");
|
||||
await expect(
|
||||
fs.readFile(path.join(stateDir, "agents", targetAgentId, "agent", "settings.json"), "utf8"),
|
||||
).resolves.toContain('"legacy":true');
|
||||
await expectMissingPath(path.join(legacySessionsDir, "sessions.json"));
|
||||
await expectMissingPath(legacyAgentDir);
|
||||
},
|
||||
);
|
||||
|
|
@ -2087,7 +2084,6 @@ describe("state migrations", () => {
|
|||
);
|
||||
expect(detectionCase.channelPairing.hasLegacy).toBe(true);
|
||||
expect(detectionCase.preview).toEqual([
|
||||
`- Sessions: ${path.join(detectionCase.stateDir, "sessions")} → ${path.join(detectionCase.stateDir, "agents", "worker-1", "sessions")}`,
|
||||
`- Sessions: canonicalize legacy keys in ${path.join(detectionCase.stateDir, "agents", "worker-1", "sessions", "sessions.json")}`,
|
||||
`- Agent dir: ${path.join(detectionCase.stateDir, "agent")} → ${path.join(detectionCase.stateDir, "agents", "worker-1", "agent")}`,
|
||||
"- Channel pairing state: legacy JSON files → shared SQLite state",
|
||||
|
|
@ -2095,7 +2091,7 @@ describe("state migrations", () => {
|
|||
]);
|
||||
});
|
||||
|
||||
it("runs legacy state migrations and canonicalizes the merged session store", async () => {
|
||||
it("runs legacy state migrations and canonicalizes the per-agent session store", async () => {
|
||||
const { root, stateDir, env, cfg } = await createLegacyStateFixture({ includePreKey: true });
|
||||
cfg.session = { ...cfg.session, mainKey: "Desk" };
|
||||
const targetStorePath = path.join(stateDir, "agents", "worker-1", "sessions", "sessions.json");
|
||||
|
|
@ -2118,20 +2114,14 @@ describe("state migrations", () => {
|
|||
"canonical-runtime",
|
||||
15,
|
||||
);
|
||||
await fs.writeFile(targetStorePath, `${JSON.stringify(targetStore, null, 2)}\n`, "utf8");
|
||||
cfg.session = { ...cfg.session, store: targetStorePath };
|
||||
const legacyStorePath = path.join(stateDir, "sessions", "sessions.json");
|
||||
const legacyStore = JSON.parse(await fs.readFile(legacyStorePath, "utf8")) as Record<
|
||||
string,
|
||||
unknown
|
||||
>;
|
||||
legacyStore["Agent:main:desk"] = { sessionId: "mixed-case-foreign", updatedAt: 40 };
|
||||
legacyStore["legacy-prototype"] = {
|
||||
targetStore["Agent:main:desk"] = { sessionId: "mixed-case-foreign", updatedAt: 40 };
|
||||
targetStore["legacy-prototype"] = {
|
||||
sessionId: "prototype-row",
|
||||
updatedAt: 10,
|
||||
sessionFile: "trace.jsonl",
|
||||
};
|
||||
await fs.writeFile(legacyStorePath, `${JSON.stringify(legacyStore, null, 2)}\n`, "utf8");
|
||||
await fs.writeFile(targetStorePath, `${JSON.stringify(targetStore, null, 2)}\n`, "utf8");
|
||||
cfg.session = { ...cfg.session, store: targetStorePath };
|
||||
|
||||
const detected = await detectLegacyStateMigrations({
|
||||
cfg,
|
||||
|
|
@ -2147,25 +2137,17 @@ describe("state migrations", () => {
|
|||
config: cfg,
|
||||
now: () => 1234,
|
||||
});
|
||||
expect(result.warnings).toStrictEqual([
|
||||
`Preserved 1 ambiguous session key(s) while importing legacy sessions into ${targetStorePath}`,
|
||||
]);
|
||||
expect(result.warnings).toStrictEqual([]);
|
||||
expect(result.changes).toEqual([
|
||||
"Migrated 2 chatapp/alpha allowFrom entries → shared SQLite state",
|
||||
`Moved MobileAuth auth creds.json → ${path.join(stateDir, "credentials", "mobileauth", "default", "creds.json")}`,
|
||||
`Moved MobileAuth auth pre-key-1.json → ${path.join(stateDir, "credentials", "mobileauth", "default", "pre-key-1.json")}`,
|
||||
`Migrated latest direct-chat session → agent:worker-1:desk`,
|
||||
`Merged sessions store → ${path.join(stateDir, "agents", "worker-1", "sessions", "sessions.json")}`,
|
||||
"Canonicalized 3 legacy session key(s)",
|
||||
"Moved trace.jsonl → agents/worker-1/sessions",
|
||||
"Canonicalized 4 legacy session key(s)",
|
||||
"Migrated 2 ACP session metadata rows → shared SQLite state",
|
||||
"Moved agent file settings.json → agents/worker-1/agent",
|
||||
]);
|
||||
expect(result.stepReceipts.find((receipt) => receipt.id === "sessions")).toMatchObject({
|
||||
outcome: "warning",
|
||||
warnings: [
|
||||
`Preserved 1 ambiguous session key(s) while importing legacy sessions into ${targetStorePath}`,
|
||||
],
|
||||
outcome: "completed",
|
||||
});
|
||||
expect(
|
||||
result.stepReceipts.find((receipt) => receipt.id === "acp-session-metadata"),
|
||||
|
|
@ -2180,7 +2162,6 @@ describe("state migrations", () => {
|
|||
"utf8",
|
||||
),
|
||||
) as Record<string, { sessionId: string; sessionFile?: string; acp?: unknown }>;
|
||||
expect(mergedStore["agent:worker-1:desk"]?.sessionId).toBe("legacy-direct");
|
||||
expect(mergedStore["group:mobile-room"]).toBeUndefined();
|
||||
expect(mergedStore["group:legacy-room"]).toBeUndefined();
|
||||
expect(mergedStore["agent:worker-1:unknown:group:mobile-room"]?.sessionId).toBe(
|
||||
|
|
@ -2198,11 +2179,12 @@ describe("state migrations", () => {
|
|||
expect(mergedStore["agent:worker-1:legacy-prototype"]).not.toHaveProperty("sessionFile");
|
||||
expect(mergedStore["agent:worker-1:acp:task"]?.acp).toBeUndefined();
|
||||
|
||||
await expect(fs.readFile(path.join(stateDir, "sessions", "trace.jsonl"), "utf8")).resolves.toBe(
|
||||
"{}\n",
|
||||
);
|
||||
await expect(
|
||||
fs.readFile(path.join(stateDir, "agents", "worker-1", "sessions", "trace.jsonl"), "utf8"),
|
||||
).resolves.toBe("{}\n");
|
||||
await expectMissingPath(path.join(stateDir, "sessions", "sessions.json"));
|
||||
await expectMissingPath(path.join(stateDir, "sessions", "trace.jsonl"));
|
||||
fs.readFile(path.join(stateDir, "sessions", "sessions.json"), "utf8"),
|
||||
).resolves.toContain("legacy-direct");
|
||||
|
||||
await expect(
|
||||
fs.readFile(path.join(stateDir, "agents", "worker-1", "agent", "settings.json"), "utf8"),
|
||||
|
|
@ -2228,13 +2210,14 @@ describe("state migrations", () => {
|
|||
await expectMissingPath(path.join(stateDir, "credentials", "chatapp-allowFrom.json"));
|
||||
});
|
||||
|
||||
it("canonicalizes parsed owners before removing the legacy store", async () => {
|
||||
it("preserves parsed owners while repairing the target main alias", async () => {
|
||||
const { root, stateDir, env } = createMigrationContext(await createTempDir());
|
||||
const legacyStorePath = path.join(stateDir, "sessions", "sessions.json");
|
||||
const legacyStorePath = path.join(stateDir, "agents", "main", "sessions", "sessions.json");
|
||||
await fs.mkdir(path.dirname(legacyStorePath), { recursive: true });
|
||||
await fs.writeFile(
|
||||
legacyStorePath,
|
||||
JSON.stringify({
|
||||
"agent:main:main": { sessionId: "main-session", updatedAt: 30 },
|
||||
"agent:archive:main": { sessionId: "archive-session", updatedAt: 20 },
|
||||
}),
|
||||
"utf8",
|
||||
|
|
@ -2252,9 +2235,10 @@ describe("state migrations", () => {
|
|||
string,
|
||||
{ sessionId: string }
|
||||
>;
|
||||
expect(store["agent:main:work"]?.sessionId).toBe("main-session");
|
||||
expect(store["agent:archive:work"]?.sessionId).toBe("archive-session");
|
||||
expect(store["agent:main:main"]).toBeUndefined();
|
||||
expect(store["agent:archive:main"]).toBeUndefined();
|
||||
await expectMissingPath(legacyStorePath);
|
||||
});
|
||||
|
||||
it("defers non-main owner merges across hard-linked stores", async () => {
|
||||
|
|
@ -2303,10 +2287,14 @@ describe("state migrations", () => {
|
|||
);
|
||||
});
|
||||
|
||||
it("defers an unambiguous legacy merge through a final store symlink", async () => {
|
||||
it("defers per-agent key repair through a final store symlink", async () => {
|
||||
const { root, stateDir, env } = createMigrationContext(await createTempDir());
|
||||
const outsideStorePath = path.join(root, "outside-sessions.json");
|
||||
await fs.writeFile(outsideStorePath, "{}\n", "utf8");
|
||||
await fs.writeFile(
|
||||
outsideStorePath,
|
||||
'{"task":{"sessionId":"legacy-task","updatedAt":10}}\n',
|
||||
"utf8",
|
||||
);
|
||||
const targetStorePath = path.join(stateDir, "agents", "main", "sessions", "sessions.json");
|
||||
await fs.mkdir(path.dirname(targetStorePath), { recursive: true });
|
||||
await fs.symlink(outsideStorePath, targetStorePath);
|
||||
|
|
@ -2325,7 +2313,7 @@ describe("state migrations", () => {
|
|||
const result = await runLegacyStateMigrations({ detected, config: cfg, now: () => 1234 });
|
||||
|
||||
expect((await fs.lstat(targetStorePath)).isSymbolicLink()).toBe(true);
|
||||
await expect(fs.readFile(outsideStorePath, "utf8")).resolves.toBe("{}\n");
|
||||
await expect(fs.readFile(outsideStorePath, "utf8")).resolves.toContain('"task"');
|
||||
await expect(fs.readFile(legacyStorePath, "utf8")).resolves.toContain("legacy-task");
|
||||
expect(result.warnings).toContain(
|
||||
`Deferred legacy session migration in final-component symlink store ${targetStorePath}; configure one canonical session.store path, then rerun openclaw doctor --fix`,
|
||||
|
|
@ -2336,7 +2324,7 @@ describe("state migrations", () => {
|
|||
const { root, stateDir, env } = createMigrationContext(await createTempDir());
|
||||
const targetStorePath = path.join(stateDir, "agents", "main", "sessions", "sessions.json");
|
||||
await fs.mkdir(path.dirname(targetStorePath), { recursive: true });
|
||||
await fs.writeFile(targetStorePath, "{}\n", "utf8");
|
||||
await fs.writeFile(targetStorePath, '{"task":{"sessionId":"legacy","updatedAt":10}}\n', "utf8");
|
||||
const configuredStorePath = path.join(root, "configured-sessions.json");
|
||||
await fs.writeFile(configuredStorePath, "{}\n", "utf8");
|
||||
const legacyStorePath = path.join(stateDir, "sessions", "sessions.json");
|
||||
|
|
@ -2371,14 +2359,14 @@ describe("state migrations", () => {
|
|||
expect.stringContaining("filesystem identity could not be established"),
|
||||
);
|
||||
await expect(fs.readFile(legacyStorePath, "utf8")).resolves.toContain("legacy");
|
||||
await expect(fs.readFile(targetStorePath, "utf8")).resolves.toBe("{}\n");
|
||||
await expect(fs.readFile(targetStorePath, "utf8")).resolves.toContain('"task"');
|
||||
});
|
||||
|
||||
it("keeps the legacy source when its store write fails", async () => {
|
||||
const { root, stateDir, env } = createMigrationContext(await createTempDir());
|
||||
const targetStorePath = path.join(stateDir, "agents", "main", "sessions", "sessions.json");
|
||||
await fs.mkdir(path.dirname(targetStorePath), { recursive: true });
|
||||
await fs.writeFile(targetStorePath, "{}\n", "utf8");
|
||||
await fs.writeFile(targetStorePath, '{"task":{"sessionId":"legacy","updatedAt":10}}\n', "utf8");
|
||||
const legacyStorePath = path.join(stateDir, "sessions", "sessions.json");
|
||||
await fs.mkdir(path.dirname(legacyStorePath), { recursive: true });
|
||||
await fs.writeFile(
|
||||
|
|
@ -2416,49 +2404,6 @@ describe("state migrations", () => {
|
|||
await expect(fs.readFile(legacyStorePath, "utf8")).resolves.toContain("legacy");
|
||||
});
|
||||
|
||||
it("preserves shared ownership through missing parent-symlink store paths", async () => {
|
||||
const { root, stateDir, env } = createMigrationContext(await createTempDir());
|
||||
const agentsDir = path.join(stateDir, "agents");
|
||||
await fs.mkdir(agentsDir, { recursive: true });
|
||||
const aliasAgentsDir = path.join(root, "agents-alias");
|
||||
await fs.symlink(agentsDir, aliasAgentsDir, "dir");
|
||||
const configuredStorePath = path.join(aliasAgentsDir, "ops", "sessions", "sessions.json");
|
||||
const targetStorePath = path.join(agentsDir, "ops", "sessions", "sessions.json");
|
||||
const legacyStorePath = path.join(stateDir, "sessions", "sessions.json");
|
||||
await fs.mkdir(path.dirname(legacyStorePath), { recursive: true });
|
||||
await fs.writeFile(
|
||||
legacyStorePath,
|
||||
JSON.stringify({
|
||||
"agent:main:work": { sessionId: "foreign-main", updatedAt: 10 },
|
||||
}),
|
||||
"utf8",
|
||||
);
|
||||
const cfg = {
|
||||
session: { mainKey: "work", store: configuredStorePath },
|
||||
agents: { list: [{ id: "ops", default: true }] },
|
||||
} as OpenClawConfig;
|
||||
const detected = await detectLegacyStateMigrations({
|
||||
cfg,
|
||||
env,
|
||||
homedir: () => root,
|
||||
pluginSessionStoreAgentIds: ["voice"],
|
||||
});
|
||||
expect(detected.sessions.preserveAmbiguousKeys).toBe(true);
|
||||
expect(detected.sessions.preserveForeignMainAliases).toBe(true);
|
||||
|
||||
await runLegacyStateMigrations({ detected, config: cfg, now: () => 1234 });
|
||||
|
||||
const store = JSON.parse(await fs.readFile(targetStorePath, "utf8")) as Record<
|
||||
string,
|
||||
{ sessionId: string }
|
||||
>;
|
||||
expect(store["agent:main:work"]?.sessionId).toBe("foreign-main");
|
||||
expect(store["agent:ops:work"]).toBeUndefined();
|
||||
await expect(fs.readFile(configuredStorePath, "utf8")).resolves.toBe(
|
||||
await fs.readFile(targetStorePath, "utf8"),
|
||||
);
|
||||
});
|
||||
|
||||
describe("aliased store ownership", () => {
|
||||
let configuredStorePath: string;
|
||||
let targetStorePath: string;
|
||||
|
|
@ -3131,7 +3076,7 @@ describe("state migrations", () => {
|
|||
);
|
||||
});
|
||||
|
||||
it("migrates existing and imported ACP metadata in one canonical session phase", async () => {
|
||||
it("migrates multi-owner ACP metadata in one canonical session phase", async () => {
|
||||
const { root, stateDir, env } = createMigrationContext(await createTempDir());
|
||||
const storeTemplate = path.join(
|
||||
stateDir,
|
||||
|
|
@ -3146,21 +3091,6 @@ describe("state migrations", () => {
|
|||
await fs.mkdir(path.dirname(storePath), { recursive: true });
|
||||
await fs.writeFile(
|
||||
storePath,
|
||||
JSON.stringify({
|
||||
"agent:main:existing": createLegacyAcpSessionEntry(
|
||||
"existing-main",
|
||||
20,
|
||||
"main",
|
||||
"existing-runtime",
|
||||
20,
|
||||
),
|
||||
}),
|
||||
"utf8",
|
||||
);
|
||||
const legacyStorePath = path.join(stateDir, "sessions", "sessions.json");
|
||||
await fs.mkdir(path.dirname(legacyStorePath), { recursive: true });
|
||||
await fs.writeFile(
|
||||
legacyStorePath,
|
||||
JSON.stringify({
|
||||
"agent:voice:main": createLegacyAcpSessionEntry(
|
||||
"voice-main",
|
||||
|
|
@ -3169,6 +3099,13 @@ describe("state migrations", () => {
|
|||
"voice-runtime",
|
||||
10,
|
||||
),
|
||||
"agent:main:existing": createLegacyAcpSessionEntry(
|
||||
"existing-main",
|
||||
20,
|
||||
"main",
|
||||
"existing-runtime",
|
||||
20,
|
||||
),
|
||||
}),
|
||||
"utf8",
|
||||
);
|
||||
|
|
@ -5268,90 +5205,6 @@ describe("state migrations", () => {
|
|||
]);
|
||||
});
|
||||
|
||||
it("preserves a corrupt target session store instead of overwriting it with legacy-only data", async () => {
|
||||
const { root, stateDir, env, cfg } = await createLegacyStateFixture();
|
||||
|
||||
const targetStorePath = path.join(stateDir, "agents", "worker-1", "sessions", "sessions.json");
|
||||
// target sessions.json is corrupt (trailing garbage → JSON5.parse fails) and
|
||||
// holds a target-only key that has no legacy counterpart.
|
||||
const corruptBytes = `${JSON.stringify({
|
||||
"agent:worker-1:desk:target-only": { sessionId: "target-only-session", updatedAt: 99 },
|
||||
})}\n<<<corrupt trailing garbage>>>`;
|
||||
await fs.writeFile(targetStorePath, corruptBytes, "utf8");
|
||||
|
||||
const detected = await detectLegacyStateMigrations({
|
||||
cfg,
|
||||
env,
|
||||
homedir: () => root,
|
||||
});
|
||||
const result = await runLegacyStateMigrations({
|
||||
detected,
|
||||
now: () => 1234,
|
||||
});
|
||||
|
||||
// The corrupt bytes must survive on disk (parse still fails after migration).
|
||||
const afterRaw = await fs.readFile(targetStorePath, "utf8");
|
||||
expect(afterRaw).toContain("corrupt trailing garbage");
|
||||
expect(afterRaw).toBe(corruptBytes);
|
||||
|
||||
// No "Merged sessions store" change was committed against the corrupt target.
|
||||
expect(result.changes.some((c) => c.startsWith("Merged sessions store"))).toBe(false);
|
||||
|
||||
// And no direct-chat migration is reported either: the legacy direct entry was
|
||||
// not saved (the target was left untouched), so doctor/startup logs must not
|
||||
// claim a session migration happened on this skip path.
|
||||
expect(result.changes.some((c) => c.startsWith("Migrated latest direct-chat session"))).toBe(
|
||||
false,
|
||||
);
|
||||
|
||||
// The user is warned that the target store was left untouched because it is unreadable.
|
||||
expect(result.warnings.some((w) => /unreadable|corrupt/i.test(w))).toBe(true);
|
||||
|
||||
// Legacy store is NOT deleted or renamed, so a later explicit doctor --fix
|
||||
// can retry the migration from the detector's normal legacy path.
|
||||
await expect(
|
||||
fs.readFile(path.join(stateDir, "sessions", "sessions.json"), "utf8"),
|
||||
).resolves.toContain("legacy-direct");
|
||||
await expect(fs.readFile(path.join(stateDir, "sessions", "trace.jsonl"), "utf8")).resolves.toBe(
|
||||
"{}\n",
|
||||
);
|
||||
});
|
||||
|
||||
it("archives a corrupt target session store before explicit recovery", async () => {
|
||||
const { root, stateDir, env, cfg } = await createLegacyStateFixture();
|
||||
|
||||
const targetStorePath = path.join(stateDir, "agents", "worker-1", "sessions", "sessions.json");
|
||||
const corruptBytes = `${JSON.stringify({
|
||||
"agent:worker-1:desk:target-only": { sessionId: "target-only-session", updatedAt: 99 },
|
||||
})}\n<<<corrupt trailing garbage>>>`;
|
||||
await fs.writeFile(targetStorePath, corruptBytes, "utf8");
|
||||
|
||||
const detected = await detectLegacyStateMigrations({
|
||||
cfg,
|
||||
env,
|
||||
homedir: () => root,
|
||||
});
|
||||
const result = await runLegacyStateMigrations({
|
||||
detected,
|
||||
now: () => 1234,
|
||||
recoverCorruptTargetStore: true,
|
||||
});
|
||||
|
||||
const archivedPath = `${targetStorePath}.corrupt-1234`;
|
||||
await expect(fs.readFile(archivedPath, "utf8")).resolves.toBe(corruptBytes);
|
||||
|
||||
const recoveredStore = JSON.parse(await fs.readFile(targetStorePath, "utf8")) as Record<
|
||||
string,
|
||||
{ sessionId?: string }
|
||||
>;
|
||||
expect(recoveredStore["agent:worker-1:desk"]?.sessionId).toBe("legacy-direct");
|
||||
expect(recoveredStore["agent:worker-1:desk:target-only"]).toBeUndefined();
|
||||
expect(result.changes).toContain(`Archived corrupt target sessions store → ${archivedPath}`);
|
||||
expect(result.changes).toContain(`Merged sessions store → ${targetStorePath}`);
|
||||
expect(result.warnings).toStrictEqual([]);
|
||||
await expectMissingPath(path.join(stateDir, "sessions", "sessions.json"));
|
||||
});
|
||||
|
||||
it("preserves a readable target store when normalization rejects an existing key", async () => {
|
||||
const { root, stateDir, env, cfg } = await createLegacyStateFixture();
|
||||
|
||||
|
|
@ -5359,7 +5212,7 @@ describe("state migrations", () => {
|
|||
const legacyStorePath = path.join(stateDir, "sessions", "sessions.json");
|
||||
const targetBytes = `${JSON.stringify(
|
||||
{
|
||||
"agent:worker-1:desk": { sessionId: "valid-session", updatedAt: 50 },
|
||||
desk: { sessionId: "valid-session", updatedAt: 50 },
|
||||
"agent:worker-1:desk:invalid": { sessionId: "../invalid", updatedAt: 60 },
|
||||
},
|
||||
null,
|
||||
|
|
@ -5390,135 +5243,8 @@ describe("state migrations", () => {
|
|||
);
|
||||
});
|
||||
|
||||
it("still filters invalid legacy-only rows", async () => {
|
||||
const { root, stateDir, env, cfg } = await createLegacyStateFixture();
|
||||
|
||||
const targetStorePath = path.join(stateDir, "agents", "worker-1", "sessions", "sessions.json");
|
||||
const legacyStorePath = path.join(stateDir, "sessions", "sessions.json");
|
||||
await fs.writeFile(
|
||||
legacyStorePath,
|
||||
`${JSON.stringify(
|
||||
{
|
||||
invalidLegacy: { sessionId: "../invalid", updatedAt: 100 },
|
||||
},
|
||||
null,
|
||||
2,
|
||||
)}\n`,
|
||||
"utf8",
|
||||
);
|
||||
|
||||
const detected = await detectLegacyStateMigrations({
|
||||
cfg,
|
||||
env,
|
||||
homedir: () => root,
|
||||
});
|
||||
const result = await runLegacyStateMigrations({
|
||||
detected,
|
||||
now: () => 1234,
|
||||
});
|
||||
|
||||
const afterStore = JSON.parse(await fs.readFile(targetStorePath, "utf8")) as Record<
|
||||
string,
|
||||
{ sessionId?: string }
|
||||
>;
|
||||
expect(afterStore["agent:worker-1:invalidLegacy"]).toBeUndefined();
|
||||
expect(Object.values(afterStore).some((entry) => entry.sessionId === "group-session")).toBe(
|
||||
true,
|
||||
);
|
||||
expect(result.changes).toContain(`Merged sessions store → ${targetStorePath}`);
|
||||
expect(result.warnings).toStrictEqual([]);
|
||||
await expectMissingPath(legacyStorePath);
|
||||
});
|
||||
|
||||
it("keeps a path-safe Unicode legacy session attached to its transcript", async () => {
|
||||
const { root, stateDir, env, cfg } = await createLegacyStateFixture();
|
||||
|
||||
const sessionId = "volume-main-हिन्दी-会議-000000";
|
||||
const transcriptName = `${sessionId}.jsonl`;
|
||||
const legacySessionsDir = path.join(stateDir, "sessions");
|
||||
const legacyStorePath = path.join(legacySessionsDir, "sessions.json");
|
||||
const targetSessionsDir = path.join(stateDir, "agents", "worker-1", "sessions");
|
||||
const targetStorePath = path.join(targetSessionsDir, "sessions.json");
|
||||
await fs.writeFile(
|
||||
legacyStorePath,
|
||||
`${JSON.stringify(
|
||||
{
|
||||
unicode: {
|
||||
sessionFile: path.join(legacySessionsDir, transcriptName),
|
||||
sessionId,
|
||||
updatedAt: 100,
|
||||
},
|
||||
},
|
||||
null,
|
||||
2,
|
||||
)}\n`,
|
||||
"utf8",
|
||||
);
|
||||
const transcript = `${JSON.stringify({ type: "session", sessionId })}\n`;
|
||||
await fs.writeFile(path.join(legacySessionsDir, transcriptName), transcript, "utf8");
|
||||
|
||||
const detected = await detectLegacyStateMigrations({
|
||||
cfg,
|
||||
env,
|
||||
homedir: () => root,
|
||||
});
|
||||
const result = await runLegacyStateMigrations({ detected, now: () => 1234 });
|
||||
|
||||
const migratedStore = JSON.parse(await fs.readFile(targetStorePath, "utf8")) as Record<
|
||||
string,
|
||||
{ sessionId?: string }
|
||||
>;
|
||||
expect(migratedStore["agent:worker-1:unicode"]?.sessionId).toBe(sessionId);
|
||||
await expect(fs.readFile(path.join(targetSessionsDir, transcriptName), "utf8")).resolves.toBe(
|
||||
transcript,
|
||||
);
|
||||
await expectMissingPath(legacyStorePath);
|
||||
expect(result.warnings).toStrictEqual([]);
|
||||
});
|
||||
|
||||
it("defers when an invalid legacy winner would replace an existing target key", async () => {
|
||||
const { root, stateDir, env, cfg } = await createLegacyStateFixture();
|
||||
|
||||
const targetStorePath = path.join(stateDir, "agents", "worker-1", "sessions", "sessions.json");
|
||||
const legacyStorePath = path.join(stateDir, "sessions", "sessions.json");
|
||||
const conflictKey = "agent:worker-1:desk:conflict";
|
||||
const targetBytes = `${JSON.stringify(
|
||||
{
|
||||
[conflictKey]: { sessionId: "target-session", updatedAt: 50 },
|
||||
},
|
||||
null,
|
||||
2,
|
||||
)}\n`;
|
||||
const legacyBytes = `${JSON.stringify(
|
||||
{
|
||||
[conflictKey]: { sessionId: "../invalid", updatedAt: 100 },
|
||||
},
|
||||
null,
|
||||
2,
|
||||
)}\n`;
|
||||
await fs.writeFile(targetStorePath, targetBytes, "utf8");
|
||||
await fs.writeFile(legacyStorePath, legacyBytes, "utf8");
|
||||
|
||||
const detected = await detectLegacyStateMigrations({
|
||||
cfg,
|
||||
env,
|
||||
homedir: () => root,
|
||||
});
|
||||
const result = await runLegacyStateMigrations({
|
||||
detected,
|
||||
now: () => 1234,
|
||||
});
|
||||
|
||||
await expect(fs.readFile(targetStorePath, "utf8")).resolves.toBe(targetBytes);
|
||||
await expect(fs.readFile(legacyStorePath, "utf8")).resolves.toBe(legacyBytes);
|
||||
expect(result.changes.some((change) => change.startsWith("Merged sessions store"))).toBe(false);
|
||||
expect(result.warnings).toContainEqual(
|
||||
expect.stringContaining("normalization rejected 1 existing target session key"),
|
||||
);
|
||||
});
|
||||
|
||||
it.runIf(process.platform !== "win32")(
|
||||
"preserves a key-shaped pending row while Doctor moves its matching legacy transcript",
|
||||
"preserves a key-shaped pending row during per-agent session key repair",
|
||||
async () => {
|
||||
const { root, stateDir, env, cfg } = await createLegacyStateFixture();
|
||||
|
||||
|
|
@ -5531,7 +5257,7 @@ describe("state migrations", () => {
|
|||
);
|
||||
const pendingKey = "agent:worker-1:desk";
|
||||
const transcriptKey = "agent:worker-1:desk:transcript";
|
||||
const ordinaryKey = "agent:worker-1:desk:ordinary";
|
||||
const ordinaryKey = "ordinary";
|
||||
const targetStore = {
|
||||
[pendingKey]: {
|
||||
sessionId: pendingKey,
|
||||
|
|
@ -5564,18 +5290,12 @@ describe("state migrations", () => {
|
|||
});
|
||||
|
||||
expect(result.warnings).toEqual([]);
|
||||
expect(result.changes).toContain(`Merged sessions store → ${targetStorePath}`);
|
||||
expect(result.changes).toContain("Moved trace.jsonl → agents/worker-1/sessions");
|
||||
expect(result.changes).toContain(`Moved ${pendingKey}.jsonl → agents/worker-1/sessions`);
|
||||
expect(result.changes).not.toContain("Rewrote migrated session transcript paths");
|
||||
await expect(
|
||||
fs.readFile(path.join(stateDir, "agents", "worker-1", "sessions", "trace.jsonl"), "utf8"),
|
||||
fs.readFile(path.join(stateDir, "sessions", "trace.jsonl"), "utf8"),
|
||||
).resolves.toBe("{}\n");
|
||||
await expect(
|
||||
fs.readFile(
|
||||
path.join(stateDir, "agents", "worker-1", "sessions", `${pendingKey}.jsonl`),
|
||||
"utf8",
|
||||
),
|
||||
fs.readFile(path.join(stateDir, "sessions", `${pendingKey}.jsonl`), "utf8"),
|
||||
).resolves.toBe('{"type":"session"}\n');
|
||||
|
||||
const afterStore = JSON.parse(await fs.readFile(targetStorePath, "utf8")) as Record<
|
||||
|
|
@ -5610,7 +5330,7 @@ describe("state migrations", () => {
|
|||
expect(afterStore[pendingKey]?.sessionId).toBeUndefined();
|
||||
expect(afterStore[transcriptKey]?.sessionId).toBe("trace");
|
||||
expect(afterStore[transcriptKey]?.sessionFile).toBeUndefined();
|
||||
expect(afterStore[ordinaryKey]?.sessionId).toBe("ordinary-session");
|
||||
expect(afterStore["agent:worker-1:ordinary"]?.sessionId).toBe("ordinary-session");
|
||||
|
||||
const firstBytes = await fs.readFile(targetStorePath, "utf8");
|
||||
const rerun = await rerunAutomaticMigrationAfterRestart({
|
||||
|
|
|
|||
|
|
@ -44,8 +44,6 @@ export type LegacyStateDetection = Pick<MigrationMessages, "warningDisposition"
|
|||
oauthDir: string;
|
||||
pluginSessionStoreAgentIds: readonly string[];
|
||||
sessions: {
|
||||
legacyDir: string;
|
||||
legacyStorePath: string;
|
||||
targetDir: string;
|
||||
targetStorePath: string;
|
||||
hasLegacy: boolean;
|
||||
|
|
|
|||
|
|
@ -5,11 +5,11 @@ import { assertSqliteFlipStartupRefusal } from "./sqlite-sessions-transcripts-fl
|
|||
function startupRefusal(command: string) {
|
||||
return {
|
||||
message: `gateway refused startup: legacy migration required (code=78 signal=null)
|
||||
Legacy session store requires migration: /qa/state/sessions/sessions.json. Run "${command}" against the same state/config before starting OpenClaw.`,
|
||||
Legacy session store requires migration: /qa/state/agents/main/sessions/sessions.json. Run "${command}" against the same state/config before starting OpenClaw.`,
|
||||
preservedSourceFiles: [
|
||||
"agents/main/sessions/sessions.json",
|
||||
"agents/main/sessions/archive-fixture/cold-archive.jsonl",
|
||||
"sessions/sessions.json",
|
||||
"agents/main/sessions/sqlite-legacy-main.jsonl",
|
||||
],
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -13,7 +13,7 @@ export function assertSqliteFlipStartupRefusal(
|
|||
expect.arrayContaining([
|
||||
"agents/main/sessions/sessions.json",
|
||||
"agents/main/sessions/archive-fixture/cold-archive.jsonl",
|
||||
"sessions/sessions.json",
|
||||
"agents/main/sessions/sqlite-legacy-main.jsonl",
|
||||
]),
|
||||
);
|
||||
}
|
||||
|
|
@ -29,7 +29,6 @@ export function assertSqliteFlipProofCore(report: SqliteFlipProofReport): void {
|
|||
);
|
||||
assertSqliteFlipStartupRefusal(report.startupRefusal);
|
||||
expect(refusalCheckpoint?.activeJsonl).toEqual(seededCheckpoint?.activeJsonl);
|
||||
expect(refusalCheckpoint?.legacyStateJsonl).toEqual(seededCheckpoint?.legacyStateJsonl);
|
||||
expect(refusalCheckpoint?.sqlite.sessionEntries).toBe(seededCheckpoint?.sqlite.sessionEntries);
|
||||
expect(refusalCheckpoint?.sqlite.transcriptEvents).toBe(
|
||||
seededCheckpoint?.sqlite.transcriptEvents,
|
||||
|
|
@ -43,21 +42,7 @@ export function assertSqliteFlipProofCore(report: SqliteFlipProofReport): void {
|
|||
)
|
||||
.every((checkpoint) => checkpoint.activeJsonl.length === 0),
|
||||
).toBe(true);
|
||||
expect(
|
||||
report.checkpoints.some(
|
||||
(checkpoint) =>
|
||||
checkpoint.label === "seeded-legacy-store" && checkpoint.legacyStateJsonl.length > 0,
|
||||
),
|
||||
).toBe(true);
|
||||
expect(
|
||||
report.checkpoints
|
||||
.filter(
|
||||
(checkpoint) =>
|
||||
checkpoint.label !== "seeded-legacy-store" &&
|
||||
checkpoint.label !== "after-startup-refusal",
|
||||
)
|
||||
.every((checkpoint) => checkpoint.legacyStateJsonl.length === 0),
|
||||
).toBe(true);
|
||||
expect(seededCheckpoint?.activeJsonl.length).toBeGreaterThan(0);
|
||||
expect(
|
||||
report.checkpoints.some(
|
||||
(checkpoint) =>
|
||||
|
|
|
|||
|
|
@ -436,7 +436,6 @@ function isBuiltCliEntrypoint(entrypoint: readonly string[]): boolean {
|
|||
function buildProofContext(stateDir: string, signal?: AbortSignal) {
|
||||
const agentDir = path.join(stateDir, "agents", AGENT_ID);
|
||||
const activeSessionsDir = path.join(agentDir, "sessions");
|
||||
const legacySessionsDir = path.join(stateDir, "sessions");
|
||||
return {
|
||||
activeSessionsDir,
|
||||
cleanups: new Set<ProofCleanup>(),
|
||||
|
|
@ -449,7 +448,6 @@ function buildProofContext(stateDir: string, signal?: AbortSignal) {
|
|||
deleteSessionKey: DELETE_SESSION_KEY,
|
||||
fullTurnAssistantText: FULL_TURN_ASSISTANT_TEXT,
|
||||
fullTurnSessionKey: FULL_TURN_SESSION_KEY,
|
||||
legacySessionsDir,
|
||||
legacySessionId: "sqlite-legacy-main",
|
||||
mockOpenAiRequestLog: path.join(stateDir, "mock-openai-requests.ndjson"),
|
||||
oldStateSessionKeys: [...OLD_STATE_SESSION_KEYS],
|
||||
|
|
@ -657,7 +655,6 @@ function ownProofChild(context: ProofContext, child: ProofChildProcess): () => P
|
|||
|
||||
async function seedLegacySessionStore(context: ProofContext): Promise<void> {
|
||||
await fs.mkdir(context.activeSessionsDir, { recursive: true });
|
||||
await fs.mkdir(context.legacySessionsDir, { recursive: true });
|
||||
await fs.mkdir(path.join(context.stateDir, "agent"), { recursive: true });
|
||||
const now = Date.now();
|
||||
const firstSharedSessionKey = expectDefined(
|
||||
|
|
@ -678,28 +675,24 @@ async function seedLegacySessionStore(context: ProofContext): Promise<void> {
|
|||
[secondSharedSessionKey]: legacyEntry("sqlite-shared-session", now - 3_000, {
|
||||
sessionFile: "sqlite-shared-b.jsonl",
|
||||
}),
|
||||
};
|
||||
const oldStateEntries = {
|
||||
main: legacyEntry(context.legacySessionId, now - 4_000),
|
||||
"+15551234567": legacyEntry("sqlite-old-direct", now - 5_000),
|
||||
"group:legacy-room": legacyEntry("sqlite-old-group", now - 6_000, {
|
||||
[context.resetSessionKey]: legacyEntry(context.legacySessionId, now - 4_000),
|
||||
"agent:main:+15551234567": legacyEntry("sqlite-old-direct", now - 5_000),
|
||||
"agent:main:unknown:group:legacy-room": legacyEntry("sqlite-old-group", now - 6_000, {
|
||||
groupChannel: "legacy-room",
|
||||
}),
|
||||
"partial-direct": legacyEntry("sqlite-partial-import", now - 7_000),
|
||||
"agent:main:partial-direct": legacyEntry("sqlite-partial-import", now - 7_000),
|
||||
};
|
||||
for (const [index, sessionKey] of SCALE_SESSION_KEYS.entries()) {
|
||||
entries[sessionKey] = legacyEntry(scaleSessionId(index), now - 20_000 - index);
|
||||
}
|
||||
await writeJsonFile(context.storePath, entries, 2);
|
||||
await writeJsonFile(path.join(context.legacySessionsDir, "sessions.json"), oldStateEntries, 2);
|
||||
await writeJsonFile(path.join(context.stateDir, "agent", "old-settings.json"), {
|
||||
source: "old-agent-layout",
|
||||
});
|
||||
const legacyDir = context.legacySessionsDir;
|
||||
const activeDir = context.activeSessionsDir;
|
||||
await writeMessageTranscript(legacyDir, context.legacySessionId, "sqlite-user-1", "legacy hello");
|
||||
await writeMessageTranscript(legacyDir, "sqlite-old-direct", "sqlite-old-direct-1", "old dm");
|
||||
await writeMessageTranscript(legacyDir, "sqlite-old-group", "sqlite-old-group-1", "old group");
|
||||
await writeMessageTranscript(activeDir, context.legacySessionId, "sqlite-user-1", "legacy hello");
|
||||
await writeMessageTranscript(activeDir, "sqlite-old-direct", "sqlite-old-direct-1", "old dm");
|
||||
await writeMessageTranscript(activeDir, "sqlite-old-group", "sqlite-old-group-1", "old group");
|
||||
await writeMessageTranscript(activeDir, "sqlite-delete-session", "sqlite-delete-1", "delete me");
|
||||
await writeMessageTranscript(
|
||||
activeDir,
|
||||
|
|
@ -741,10 +734,10 @@ async function seedLegacySessionStore(context: ProofContext): Promise<void> {
|
|||
]);
|
||||
}
|
||||
await writeJsonFile(
|
||||
path.join(context.legacySessionsDir, `${context.legacySessionId}.trajectory.jsonl`),
|
||||
path.join(context.activeSessionsDir, `${context.legacySessionId}.trajectory.jsonl`),
|
||||
{ type: "trajectory", sessionId: context.legacySessionId },
|
||||
);
|
||||
await writeJsonFile(path.join(context.legacySessionsDir, "old-orphan.deleted.jsonl"), {
|
||||
await writeJsonFile(path.join(context.activeSessionsDir, "old-orphan.deleted.jsonl"), {
|
||||
type: "event",
|
||||
id: "old-orphan",
|
||||
});
|
||||
|
|
@ -849,19 +842,17 @@ async function importProofSession(
|
|||
}
|
||||
|
||||
async function requireLegacyStartupRefusal(inst: OpenClawTestInstance, context: ProofContext) {
|
||||
const legacyStorePath = path.join(context.legacySessionsDir, "sessions.json");
|
||||
const legacyStorePath = context.storePath;
|
||||
const validStore = await fs.readFile(legacyStorePath);
|
||||
// Valid stores can migrate during startup. A refused source must remain visible
|
||||
// to the next startup instead of being moved outside migration discovery.
|
||||
// Startup must leave the refused per-agent index and every transcript intact
|
||||
// so the same source remains available for the next attempt and Doctor repair.
|
||||
await fs.writeFile(legacyStorePath, `${validStore.toString("utf8")}\n<<<invalid legacy store>>>`);
|
||||
const sources = new Map<string, Buffer>();
|
||||
for (const directory of [context.activeSessionsDir, context.legacySessionsDir]) {
|
||||
await walkFiles(directory, async (filePath) => {
|
||||
sources.set(filePath, await fs.readFile(filePath));
|
||||
});
|
||||
if (!sources.has(path.join(directory, "sessions.json"))) {
|
||||
throw new Error(`missing seeded legacy session store in ${directory}`);
|
||||
}
|
||||
await walkFiles(context.activeSessionsDir, async (filePath) => {
|
||||
sources.set(filePath, await fs.readFile(filePath));
|
||||
});
|
||||
if (!sources.has(legacyStorePath)) {
|
||||
throw new Error(`missing seeded legacy session store at ${legacyStorePath}`);
|
||||
}
|
||||
let message = "";
|
||||
for (const attempt of [1, 2]) {
|
||||
|
|
@ -2014,7 +2005,6 @@ async function captureCheckpoint(
|
|||
...(options.doctor ? { doctor: options.doctor } : {}),
|
||||
gatewayLogTail: tail(options.gatewayLogTail ?? ""),
|
||||
label,
|
||||
legacyStateJsonl: await inventoryActiveJsonl(context.legacySessionsDir),
|
||||
sqlite: readSqliteEvidence(context.agentDbPath, context.trackedSessionKeys),
|
||||
};
|
||||
}
|
||||
|
|
@ -2256,19 +2246,16 @@ function validateCheckpointInvariants(
|
|||
checkpoint: ProofCheckpoint,
|
||||
failures: string[],
|
||||
): void {
|
||||
if (checkpoint.label !== "seeded-legacy-store" && checkpoint.label !== "after-startup-refusal") {
|
||||
for (const [description, inventory] of [
|
||||
["active sessions directory", checkpoint.activeJsonl],
|
||||
["old sessions directory", checkpoint.legacyStateJsonl],
|
||||
] as const) {
|
||||
if (inventory.length > 0) {
|
||||
failures.push(
|
||||
`${checkpoint.label}: ${description} still has JSONL files: ${inventory
|
||||
.map((entry) => entry.path)
|
||||
.join(", ")}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
if (
|
||||
checkpoint.label !== "seeded-legacy-store" &&
|
||||
checkpoint.label !== "after-startup-refusal" &&
|
||||
checkpoint.activeJsonl.length > 0
|
||||
) {
|
||||
failures.push(
|
||||
`${checkpoint.label}: active sessions directory still has JSONL files: ${checkpoint.activeJsonl
|
||||
.map((entry) => entry.path)
|
||||
.join(", ")}`,
|
||||
);
|
||||
}
|
||||
const doctor = checkpoint.doctor;
|
||||
if (checkpoint.label.startsWith("after-doctor") && doctor?.code !== 0) {
|
||||
|
|
@ -2443,7 +2430,7 @@ function printCheckpoint(checkpoint: ProofCheckpoint): void {
|
|||
[
|
||||
`[sqlite-sessions-transcripts-flip-proof] ${checkpoint.label}`,
|
||||
` sqlite sessions=${checkpoint.sqlite.sessions} entries=${checkpoint.sqlite.sessionEntries} transcriptEvents=${checkpoint.sqlite.transcriptEvents}`,
|
||||
` activeJsonl=${checkpoint.activeJsonl.length} legacyStateJsonl=${checkpoint.legacyStateJsonl.length} archiveArtifacts=${checkpoint.archiveArtifacts.length}`,
|
||||
` activeJsonl=${checkpoint.activeJsonl.length} archiveArtifacts=${checkpoint.archiveArtifacts.length}`,
|
||||
checkpoint.doctor
|
||||
? ` doctor ${checkpoint.doctor.mode} code=${String(checkpoint.doctor.code)} totals=${JSON.stringify(
|
||||
checkpoint.doctor.totals ?? {},
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue