Skip to content
2 changes: 1 addition & 1 deletion .changeset/21498-cli-compose-migration-recovery.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,3 @@
---
"@objectstack/cli": patch
---
Expand All @@ -10,4 +10,4 @@

- **The `os migrate` data commands** (`recorded-by`, `resume`, `value-shapes`, `summary-nulls`, `files-to-references`, `meta --stored`, `audit-metadata-bodies`, `os storage orphans`) now boot with the plugin. A run interrupted before any of its chunks committed now resumes to completion. A command booted over an interrupted run also warns about that run on stderr first.
- **Every `os serve` boot** (and so `os start` and `os dev`, which spawn it) composes the plugin beside `PlatformObjectsPlugin`, which registers the journal the scan reads. An interrupted run is reported once at boot, with the `os migrate resume --run <id>` command that resumes it. Nothing is resumed automatically. A database with no interrupted run prints nothing. A config that composes its own `new MigrationRecoveryPlugin()` keeps that instance.
- **Still refused:** a `recorded-by` run that had committed a chunk before it was interrupted, or that was started with a non-default `--chunk-size`. `resume` now reaches the runner for these runs, and the runner refuses them with `PLAN_CHANGED`. Re-running `os migrate recorded-by --apply` converts whatever rows still hold the sentinel.
- **A run that had committed a chunk, or that was started with a non-default `--chunk-size`,** reaches the runner too. The runner fix that lets it resume is in the `@objectstack/core` entry for #21528.
15 changes: 15 additions & 0 deletions .changeset/21528-core-resume-started-over-plan.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
---
'@objectstack/core': patch
---

fix(core): a resumed migration run is compared against the chunk plan it started over, so `os migrate resume` completes an interrupted `recorded-by` run that had committed a chunk or was started with a non-default `--chunk-size` (#21528)

Clause-②: no

`runMigrationJournal` recomputed a resumed run's chunk plan from the rows `load()` returned at resume time, at the plan's current chunk size, and refused `PLAN_CHANGED` when that plan's hash differed from the one `run_started` recorded. Two kinds of interrupted run could differ. A plan whose `load()` selects only the work still to do, which `recorded-by`'s plan does, returns fewer rows once a chunk has committed. And the plan handed back for a resume carries its own chunk size, not the one the run was started with. So `os migrate resume` listed such a run as `resumable: true`, and `os migrate resume --run <id> --yes` then refused it.

A resume now reads the chunk plan back from the journal's `run_started` record:

- **Identity.** The plan's id and step names are hashed with the recorded chunk boundaries and compared with the recorded hash. A plan whose id or steps changed is still refused `PLAN_CHANGED`. The run resumes at the chunk size it started with.
- **Rows.** Each step's rows are bound to that chunk plan. If `load()` returns every row the run started over, each chunk's rows are where the journal put them, as before. If it returns exactly the rows of the chunks not yet committed, those rows go, in order, to those chunks. Any other row count is refused `PLAN_CHANGED`, and the message names the step.
- **Unwind.** If a chunk fails after a resume that bound its rows the second way, the runner compensates the chunks this process committed, newest first. It then stops at the newest chunk an earlier process committed and journals `run_failed`, because `load()` no longer returns that chunk's rows. It does not compensate other rows in their place, and the run ends `failed`.
246 changes: 202 additions & 44 deletions packages/cli/src/commands/migrate/resume.recorded-by.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,12 +34,22 @@
* composed into every served boot adds no line when there is nothing to
* report.
*
* ⛔ The interruption is in chunk 0, before any chunk committed. A run that
* had committed a chunk is refused `PLAN_CHANGED` on resume — `recorded-by`'s
* `load()` selects only rows still holding the sentinel, so the chunk plan it
* recomputes no longer hashes to the one the journal recorded. That is a
* defect of the plan's shape, not of this composition, and it is reported on
* its own; this file does not pin around it.
* ## The run's identity survives its own progress (#21528)
*
* `recorded-by`'s `load()` selects only the rows still holding the sentinel,
* so it shrinks as the run commits chunks, and the plan the owner registers
* for resume carries the default chunk size, not the one the run was started
* with. Both used to change the chunk plan a resume recomputed, and the
* runner refused `PLAN_CHANGED` a run the list had just called resumable.
* Three more interrupted runs, made the same way:
*
* 4. killed in chunk 1 after chunk 0 committed (203 rows at the default
* size): it resumes to completion, and chunk 0 is not run again;
* 5. started with a chunk size of 2 (the `--chunk-size 2` an operator
* passes) and killed in chunk 0: it resumes with the journal's size, two
* chunks, not the one chunk the registered plan's default would make;
* 6. the control: a run started by a plan whose step had another name — a
* changed plan — is still refused `PLAN_CHANGED`, and nothing is written.
*
* Every boot runs in a hook: a case only reads what a boot printed or wrote.
*/
Expand Down Expand Up @@ -97,9 +107,11 @@ function childEnv(overrides: Record<string, string | undefined>): Record<string,
/**
* The crash. Boots the data stack over the fixture database, writes the
* sentinel rows, then runs the recorded-by plan under the real runner with a
* forward that kills the process from inside chunk 0's transaction. The plan
* keeps the owner's id, step names and chunk size, so the journal records the
* hash the owner's registered plan computes on resume.
* forward that runs the owner's forward for every chunk before
* `FIXTURE_KILL_CHUNK` and kills the process from inside that chunk's
* transaction. The plan keeps the owner's id; `FIXTURE_CHUNK_SIZE` starts it
* at another size, as `--chunk-size` does, and `FIXTURE_STEP_NAME` renames its
* step, which makes it a different plan from the one the owner registers.
*/
const CRASH_CHILD = `
const rt = await import('@objectstack/runtime');
Expand All @@ -126,13 +138,20 @@ for (const [i, id] of ids.entries()) {
operation_type: 'create', recorded_by: process.env.FIXTURE_SENTINEL, recorded_at: at,
}, { context: { isSystem: true } });
}
const owned = mp.createRecordedBySentinelPlan();
const chunkSize = process.env.FIXTURE_CHUNK_SIZE ? Number(process.env.FIXTURE_CHUNK_SIZE) : undefined;
const killChunk = Number(process.env.FIXTURE_KILL_CHUNK ?? '0');
const owned = mp.createRecordedBySentinelPlan(chunkSize === undefined ? {} : { chunkSize });
const step = owned.steps[0];
const crashing = { ...owned, steps: [{ ...step, forward: async (_rows, ctx) => {
process.stderr.write('[fixture] run ' + ctx.runId + ' killed in chunk ' + ctx.chunkIndex + '\\n');
process.kill(process.pid, 'SIGKILL');
await new Promise(() => {});
} }] };
const crashing = { ...owned, steps: [{
...step,
...(process.env.FIXTURE_STEP_NAME ? { name: process.env.FIXTURE_STEP_NAME } : {}),
forward: async (rows, ctx, engine) => {
if (ctx.chunkIndex !== killChunk) return step.forward(rows, ctx, engine);
process.stderr.write('[fixture] run ' + ctx.runId + ' killed in chunk ' + ctx.chunkIndex + '\\n');
process.kill(process.pid, 'SIGKILL');
await new Promise(() => {});
},
}] };
await core.runMigrationJournal(ql, crashing);
process.stderr.write('[fixture] the run finished, so nothing was interrupted\\n');
process.exit(3);
Expand Down Expand Up @@ -232,69 +251,154 @@ async function readRows(dbFile: string, sql: string, bindings: unknown[] = []):
}
}

let root: string;
let runId: string;
let resumeList: { payload: any; exitCode: number };
let resumeAct: { payload: any; exitCode: number };
let resumed: ProjectCopy;
let scanOutput: string;
let freshOutput: string;

beforeAll(async () => {
root = mkdtempSync(join(tmpdir(), 'os-21498-'));
const origin = makeProject(root, 'origin');
interface InterruptedFixture {
/** The interrupted run's id, as the killed child printed it. */
runId: string;
/** The database the crash left behind. Copy it; never resume over it. */
origin: ProjectCopy;
}

// ── the crash ──────────────────────────────────────────────────────────────
/**
* Seed `ids` as sentinel rows and kill a recorded-by run in chunk `killChunk`.
* See `CRASH_CHILD` for what the other two options change about the run.
*/
function interruptRun(
name: string,
ids: readonly string[],
opts: { killChunk?: number; chunkSize?: number; stepName?: string } = {},
): InterruptedFixture {
const origin = makeProject(root, name);
const killChunk = opts.killChunk ?? 0;
const crash = spawnSync(process.execPath, ['--input-type=module', '-e', CRASH_CHILD], {
cwd: CLI_ROOT,
env: childEnv({
OS_ARTIFACT_PATH: join(origin.dir, 'dist', 'objectstack.json'),
FIXTURE_PROJECT: origin.dir,
FIXTURE_DB: origin.dbFile,
FIXTURE_IDS: JSON.stringify(SENTINEL_IDS),
FIXTURE_IDS: JSON.stringify(ids),
FIXTURE_SENTINEL: RECORDED_BY_SENTINEL,
FIXTURE_KILL_CHUNK: String(killChunk),
FIXTURE_CHUNK_SIZE: opts.chunkSize === undefined ? undefined : String(opts.chunkSize),
FIXTURE_STEP_NAME: opts.stepName,
}),
encoding: 'utf8',
timeout: CHILD_BUDGET_MS,
});
const killed = /\[fixture\] run (\S+) killed in chunk 0/.exec(crash.stderr ?? '');
const killed = new RegExp(`\\[fixture\\] run (\\S+) killed in chunk ${killChunk}\\b`).exec(crash.stderr ?? '');
if (crash.signal !== 'SIGKILL' || !killed) {
throw new Error(
`the fixture did not leave an interrupted run (status ${crash.status}, signal ${crash.signal})\n${crash.stderr}`,
`the fixture '${name}' did not leave an interrupted run (status ${crash.status}, signal ${crash.signal})\n${crash.stderr}`,
);
}
runId = killed[1];
return { runId: killed[1], origin };
}

// Each consumer gets its own copy: a resume concludes the run, a serve boot
// writes rows of its own.
resumed = makeProject(root, 'resumed');
const scanned = makeProject(root, 'scanned');
cpSync(join(origin.dir, 'data'), join(resumed.dir, 'data'), { recursive: true });
cpSync(join(origin.dir, 'data'), join(scanned.dir, 'data'), { recursive: true });
const fresh = makeProject(root, 'fresh');
/** A fresh copy of the database a crash left behind. */
function copyOf(fixture: InterruptedFixture, name: string): ProjectCopy {
const copy = makeProject(root, name);
cpSync(join(fixture.origin.dir, 'data'), join(copy.dir, 'data'), { recursive: true });
return copy;
}

// ── 1. the real command ────────────────────────────────────────────────────
/** Each `argv` through the real `os migrate resume --json`, in order, over `project`. */
async function resumeOver(project: ProjectCopy, ...argvs: string[][]): Promise<Array<{ payload: any; exitCode: number }>> {
const out: Array<{ payload: any; exitCode: number }> = [];
const savedEnv: Record<string, string | undefined> = {};
const savedCwd = process.cwd();
for (const key of [...OVERRIDING_ENV, 'OS_ARTIFACT_PATH'] as const) savedEnv[key] = process.env[key];
try {
for (const key of OVERRIDING_ENV) delete process.env[key];
process.env.OS_ARTIFACT_PATH = join(resumed.dir, 'dist', 'objectstack.json');
process.chdir(resumed.dir);
const url = `file:${resumed.dbFile}`;
resumeList = await resumeJson(['--database-url', url]);
resumeAct = await resumeJson(['--run', runId, '--yes', '--database-url', url]);
process.env.OS_ARTIFACT_PATH = join(project.dir, 'dist', 'objectstack.json');
process.chdir(project.dir);
const url = `file:${project.dbFile}`;
for (const argv of argvs) out.push(await resumeJson([...argv, '--database-url', url]));
} finally {
process.chdir(savedCwd);
for (const [key, value] of Object.entries(savedEnv)) {
if (value === undefined) delete process.env[key];
else process.env[key] = value;
}
}
return out;
}

/** Journal kinds and chunk indices of one run, in `seq` order. */
async function journalOf(dbFile: string, id: string): Promise<Array<{ kind: string; chunk_index: number | null }>> {
return readRows(dbFile, 'SELECT kind, chunk_index FROM sys_migration_journal WHERE run_id = ? ORDER BY seq', [id]);
}

/** How many of `ids` still hold the sentinel. */
async function sentinelCount(dbFile: string, ids: readonly string[]): Promise<number> {
const rows = await readRows(
dbFile,
`SELECT COUNT(*) AS n FROM sys_metadata_history WHERE recorded_by = ? AND id IN (${ids.map(() => '?').join(', ')})`,
[RECORDED_BY_SENTINEL, ...ids],
);
return Number(rows[0].n);
}

/** One more row than the default chunk size holds, plus two: chunk 0 is 200 rows, chunk 1 is 3. */
const COMMITTED_IDS = Array.from({ length: 203 }, (_, i) => `h_21528_c${String(i).padStart(3, '0')}`);
const SIZED_IDS = ['h_21528_s_a', 'h_21528_s_b', 'h_21528_s_c'];
const CHANGED_IDS = ['h_21528_x_a', 'h_21528_x_b', 'h_21528_x_c'];

let root: string;
let runId: string;
let resumeList: { payload: any; exitCode: number };
let resumeAct: { payload: any; exitCode: number };
let resumed: ProjectCopy;
let scanOutput: string;
let freshOutput: string;

/** The #21528 runs: each fixture, its resumed copy, and what list and act answered. */
interface ResumedFixture extends InterruptedFixture {
copy: ProjectCopy;
list?: { payload: any; exitCode: number };
act: { payload: any; exitCode: number };
}
let afterCommit: ResumedFixture;
let sized: ResumedFixture;
let changed: ResumedFixture;

beforeAll(async () => {
root = mkdtempSync(join(tmpdir(), 'os-21498-'));

// ── the crash ──────────────────────────────────────────────────────────────
const interrupted = interruptRun('origin', SENTINEL_IDS);
runId = interrupted.runId;

// Each consumer gets its own copy: a resume concludes the run, a serve boot
// writes rows of its own.
resumed = copyOf(interrupted, 'resumed');
const scanned = copyOf(interrupted, 'scanned');
const fresh = makeProject(root, 'fresh');

// ── 1. the real command ────────────────────────────────────────────────────
[resumeList, resumeAct] = await resumeOver(resumed, [], ['--run', runId, '--yes']);

// ── 2 and 3. the real serve boot ───────────────────────────────────────────
scanOutput = await serveBoot(scanned);
freshOutput = await serveBoot(fresh);

// ── 4, 5 and 6. runs whose recomputed chunk plan differs (#21528) ───────────
const resumeFixture = async (
fixture: InterruptedFixture,
name: string,
withList: boolean,
): Promise<ResumedFixture> => {
const copy = copyOf(fixture, name);
const act = ['--run', fixture.runId, '--yes'];
if (!withList) return { ...fixture, copy, act: (await resumeOver(copy, act))[0] };
const [list, done] = await resumeOver(copy, [], act);
return { ...fixture, copy, list, act: done };
};
afterCommit = await resumeFixture(interruptRun('after-commit', COMMITTED_IDS, { killChunk: 1 }), 'after-commit-resumed', true);
sized = await resumeFixture(interruptRun('sized', SIZED_IDS, { chunkSize: 2 }), 'sized-resumed', true);
changed = await resumeFixture(
interruptRun('changed', CHANGED_IDS, { stepName: 'sys_metadata_history.recorded_by: an earlier step' }),
'changed-resumed',
false,
);
}, HOOK_TIMEOUT_MS);

afterAll(() => {
Expand Down Expand Up @@ -350,3 +454,57 @@ describe('os serve reports an interrupted run through the boot scan (#21498)', (
expect(freshOutput).not.toMatch(/journal scan failed/i);
});
});

describe("os migrate resume completes a recorded-by run whose remaining rows shrank or whose chunk size was not the default (#21528)", () => {
it('lists a run killed after a committed chunk as resumable, chunk 0 known committed', () => {
const listed = afterCommit.list!.payload.interrupted.find((r: any) => r.runId === afterCommit.runId);
expect(listed, JSON.stringify(afterCommit.list!.payload)).toBeDefined();
expect(listed.committedChunks).toEqual([0]);
expect(listed.unknownChunks).toEqual([1]);
expect(listed.resumable).toBe(true);
});

it('resumes that run to completion without running its committed chunk again', async () => {
expect(afterCommit.act.exitCode, JSON.stringify(afterCommit.act.payload)).toBe(0);
expect(afterCommit.act.payload).toMatchObject({
runId: afterCommit.runId, status: 'completed', chunksTotal: 2, chunksCommitted: 2,
});
expect(await sentinelCount(afterCommit.copy.dbFile, COMMITTED_IDS)).toBe(0);

const journal = await journalOf(afterCommit.copy.dbFile, afterCommit.runId);
// Chunk 0 committed before the kill and is started exactly once; chunk 1
// is started by the killed run and again by the resume.
expect(journal.filter((e) => e.kind === 'chunk_started').map((e) => e.chunk_index)).toEqual([0, 1, 1]);
expect(journal.filter((e) => e.kind === 'chunk_done').map((e) => e.chunk_index)).toEqual([0, 1]);
expect(journal.at(-1)!.kind).toBe('run_done');
});

it('lists a run started with a chunk size of 2 as resumable', () => {
const listed = sized.list!.payload.interrupted.find((r: any) => r.runId === sized.runId);
expect(listed, JSON.stringify(sized.list!.payload)).toBeDefined();
expect(listed.unknownChunks).toEqual([0]);
expect(listed.resumable).toBe(true);
});

it('resumes that run with the chunk size it started with, from the journal', async () => {
expect(sized.act.exitCode, JSON.stringify(sized.act.payload)).toBe(0);
// Three rows at size 2 are two chunks; the registered plan's default size
// would have made one.
expect(sized.act.payload).toMatchObject({
runId: sized.runId, status: 'completed', chunksTotal: 2, chunksCommitted: 2,
});
expect(await sentinelCount(sized.copy.dbFile, SIZED_IDS)).toBe(0);
const journal = await journalOf(sized.copy.dbFile, sized.runId);
expect(journal.filter((e) => e.kind === 'chunk_done').map((e) => e.chunk_index)).toEqual([0, 1]);
expect(journal.at(-1)!.kind).toBe('run_done');
});

it('still refuses PLAN_CHANGED a run another plan started (the control), and writes nothing', async () => {
expect(changed.act.exitCode, JSON.stringify(changed.act.payload)).toBe(1);
expect(changed.act.payload.error).toMatch(/^Refused \(PLAN_CHANGED\)/);
expect(await sentinelCount(changed.copy.dbFile, CHANGED_IDS)).toBe(CHANGED_IDS.length);
expect((await journalOf(changed.copy.dbFile, changed.runId)).map((e) => e.kind)).toEqual([
'run_started', 'chunk_started',
]);
});
});
Loading
Loading