Skip to content

Commit 10454b3

Browse files
fix(core): os migrate resume completes a recorded-by run that committed a chunk or used a non-default --chunk-size (#21528) (#21554)
Fixes #21528 Clause-②: no ## What changed `runMigrationJournal` (`packages/core/src/utils/migration-journal.ts`) recomputed a resumed run's chunk plan from the rows `load()` returns at resume time, at the resumed plan's chunk size, and refused `PLAN_CHANGED` when that plan's hash differed from the one `run_started` recorded. A resume now reads the chunk plan back from `run_started`, which has carried it since the runner's first commit (ADR-0119 D2 item 2: "carrying the plan hash and chunk plan"): - **Identity, the one place the journal hashes.** `hashMigrationPlan` is unchanged. On a resume it hashes the RECORDED chunk boundaries with the plan's declared id and step names, so it compares the plan against what the run started over. The run's chunk size comes back from the journal. A plan whose id or steps changed still refuses `PLAN_CHANGED`. - **Rows.** Per step, with N rows started over and K of them in committed chunks: a `load()` that returns N rows binds positionally, as before. One that returns exactly N − K rows (it selects only the remaining work, as `recorded-by`'s does) binds those rows, in order, to the chunks not yet committed. Any other count refuses `PLAN_CHANGED` and names the step. Every resume that passed before binds exactly as before (N rows, the same boundaries). - **Unwind after such a resume.** A chunk an earlier process committed, whose rows `load()` no longer returns, cannot be compensated. The unwind compensates this process's chunks newest-first, then halts with `run_failed` at that chunk with a `reason`. It does not hand `compensate()` other rows and journal a clean unwind. - A `run_started` with no recorded chunk plan (only a hand-written journal) keeps today's check: reproduce the recorded hash from the current rows. No published member is added. `MigrationPlan`, `MigrationPlanStep`, `MigrationJournalEvent` and the `run_started` payload keep their shapes. `MigrationPlanStep.load` and `RunMigrationJournalOptions.chunkSize` gain TSDoc for what a resume does with them. The plan (`recorded-by-sentinel.ts`), `os migrate resume`'s source and the list mode's `resumable` are untouched. No new error code: both new refusals are `PLAN_CHANGED`, the code the same inputs drew before. ## The public door, before and after `packages/cli/src/commands/migrate/resume.recorded-by.integration.test.ts` now makes three more interrupted runs the way a crash does (a child process runs the recorded-by plan under the real runner and is SIGKILLed inside a chunk's transaction) and drives the real `os migrate resume --json`: - **(a) 203 sentinel rows, killed in chunk 1 after chunk 0 committed.** The list says `resumable: true`, `committedChunks: [0]`, `unknownChunks: [1]`. Before (core built from `1ac7308d7a`): `--run RUN_ID --yes` exited 1 with `Refused (PLAN_CHANGED): ... plan hash 037ebfcc70ee54096b8eb2a6aed8cd49 does not match the journal's f36863e5fee4bb4e40093c66a0b3096f`. After: exit 0, `completed`, 2 of 2 chunks, no row left holding the sentinel, `chunk_started` indices `[0, 1, 1]` (chunk 0 is not run again). - **(b) 3 rows started at chunk size 2, killed in chunk 0.** The list says `resumable: true`. Before: exit 1, `PLAN_CHANGED`. After: exit 0, `completed`, `chunksTotal: 2` (the journal's size; the registered plan's default 200 would make one chunk). - **(c) The control.** A run started by a plan whose step had another name: exit 1, `Refused (PLAN_CHANGED)`, before and after. The sentinel rows and the journal are untouched. Before the fix that file read 2 failed / 7 passed (the two acts above); after, 9 passed. ## Tests (at `14e5287fa8`, this branch's head) - `packages/core/src/utils/migration-journal.test.ts`: 29 passed (23 before + 6 new: a shrinking load resumes after a committed chunk; a non-shrinking load resumes positionally from a runner-written journal; the journal's chunk size wins over the plan's; the changed-plan control, three ways (step renamed, plan id changed, step added), each refused with `code: 'PLAN_CHANGED'` and zero journal writes; a row count that is neither N nor N − K is refused, naming the step; the unwind halt). The crash helper runs the REAL runner and stops a forward inside its chunk, so every resume reads a journal the runner wrote. - `packages/metadata-protocol/src/migrations/recorded-by-sentinel.test.ts`: 8 passed (1 new: the real plan, started at size 2 and killed after chunk 0 committed, resumes with the plan the owner registers at its default size, 2 of 2 chunks). - The door file above: 9 passed (`--project integration`, run locally because this diff edits that file). - Typecheck: `@objectstack/core` (with `check:test-typecheck`: 4 files / 4 errors held, unchanged), `@objectstack/metadata-protocol`, `@objectstack/cli` (test layer: 3 files / 28 errors held, unchanged), all exit 0. - Full suites on `ae27f00812` (this diff before the merge of `main`, which touched none of these files): `@objectstack/core` 79 files / 2220 tests passed; `@objectstack/metadata-protocol` 205 files passed, 3 skipped / 3152 tests passed, 19 skipped. `@objectstack/cli`'s unit tier: only the tier-partition pin (`test/vitest-tiers-partition.test.ts`, 22 passed); this diff changes no CLI source. - Gates: `node scripts/pm/dispatch-gates.mjs --commands` on `14e5287fa8` derived 67 commands; all 67 ran, reconciled with `--ran` (each line carrying its exit code): 0 NOT MEASURED, 66 exit 0, and one exit 1, `node scripts/check-empty-changeset.mjs --base origin/main`, explained in the next section. Four roster gates the derivation flags as sharing a directory with these paths also ran green: `check-changeset-fixed`, `check:authz-resolver`, `check:error-code-casing`, `check:filter-alias-parity`. - Lint, narrowed: the population is the 4 changed `.ts` files, none ignored by `eslint.config.mjs`. `eslint --no-inline-config --format json` read 4 files, 0 errors, 0 warnings. That config never enables type-aware linting (no `parserOptions.project`), so this diff cannot move the verdict of any file it does not touch. The repo-wide `pnpm lint` is CI's. ## Reverse verification and ablation - **Reverse verification**, with the fix committed: `migration-journal.ts` restored to `1ac7308d7a` in the working tree only, blob `df8d8d5009` checked equal to the base's, then core rebuilt and `node scripts/ablation-dist-preflight.mjs @objectstack/core planResumedRun --absent` passed. Red as predicted: core unit 4 failed / 25 passed (the four resume pins; the control and the positional-resume pin stay green, as they should on both trees), the plan pin 1 failed / 7 passed, the door 2 failed / 7 passed (a and b; the control green). Restored with `git checkout HEAD --` under an EXIT/INT/TERM trap, proven by the blob (`361d75e7ee` == HEAD) and an empty `git diff HEAD`, then rebuilt, with the preflight in default mode showing the marker back in 4 built files and a clean tree. - **Ablation of the unwind guard**, which reverse verification cannot isolate (the old runner refuses before reaching it): `node scripts/ablation-replace.mjs` planted the naive positional fallback (`rowsByChunk.get(c.index) ?? rowsByStep[...].slice(offset, offset + length)`). The landing was shown by the anchor count 1 → 0 and the blob change. The unwind pin went red: `expected 'compensated' to be 'failed'`, which is the run handing chunk 0 other rows and journalling a clean unwind. Restored by the tool: blob == HEAD, `git diff HEAD` empty. The core unit suite imports the runner by relative path, so no build leg applies. ## The pending #21498 changeset: a correction to confirm This PR changes `.changeset/21498-cli-compose-migration-recovery.md`, which it did not add. That note's "Still refused" bullet said the runner refuses these two kinds of run with `PLAN_CHANGED`. This change makes that false, and both notes are still pending, so they would ship in one release. The bullet now says these runs reach the runner too and points to the `@objectstack/core` entry for #21528. `check-empty-changeset` stays red on this by design: it is the gate's DELIBERATE CORRECTION class, and its remedy is to say so here and get the correction confirmed. Restoring the old bullet would publish a sentence this PR makes false. **Please confirm the correction.** The new changeset (`.changeset/21528-core-resume-started-over-plan.md`) is an `@objectstack/core` patch with `Clause-②: no`. No member is added to a published contract. Resume now accepts the runs its list mode already advertises as resumable, as ADR-0119 D2 item 5 declares. ## Acceptance notes - **The list's `resumable` is plan presence only** (`resume.ts`: `Boolean(plans?.get(r.planId))`). So a run whose plan genuinely changed, or whose rows moved, is still listed `resumable: true` and then refused. The triage ruling keeps the list's wording out of this card. I read this from source; I did not measure the list for the control run. - **A forward resume skips compensated chunks.** The forward loop skips every chunk with a `chunk_done`, and a chunk that was committed and then compensated has one. So resuming forward a run whose in-run unwind failed partway would skip the chunks that unwind had undone. For `recorded-by` (a shrinking `load()`), that run's row count matches neither binding, so it is refused `PLAN_CHANGED`, as it was before. Only a plan whose `load()` does not shrink would take the skip, and no such plan is registered on this tree. I read this from source and did not measure it. Carrier: none. --- _Generated by [Claude Code](https://claude.ai/code/session_01DDZNkDVwPQnevTFcYE47H3)_ --------- Co-authored-by: Claude <noreply@anthropic.com>
1 parent ce53218 commit 10454b3

6 files changed

Lines changed: 647 additions & 79 deletions

File tree

‎.changeset/21498-cli-compose-migration-recovery.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,4 +10,4 @@ Clause-②: no
1010

1111
- **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.
1212
- **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.
13-
- **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.
13+
- **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.
Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,15 @@
1+
---
2+
'@objectstack/core': patch
3+
---
4+
5+
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)
6+
7+
Clause-②: no
8+
9+
`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.
10+
11+
A resume now reads the chunk plan back from the journal's `run_started` record:
12+
13+
- **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.
14+
- **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.
15+
- **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`.

‎packages/cli/src/commands/migrate/resume.recorded-by.integration.test.ts‎

Lines changed: 202 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -34,12 +34,22 @@
3434
* composed into every served boot adds no line when there is nothing to
3535
* report.
3636
*
37-
* ⛔ The interruption is in chunk 0, before any chunk committed. A run that
38-
* had committed a chunk is refused `PLAN_CHANGED` on resume — `recorded-by`'s
39-
* `load()` selects only rows still holding the sentinel, so the chunk plan it
40-
* recomputes no longer hashes to the one the journal recorded. That is a
41-
* defect of the plan's shape, not of this composition, and it is reported on
42-
* its own; this file does not pin around it.
37+
* ## The run's identity survives its own progress (#21528)
38+
*
39+
* `recorded-by`'s `load()` selects only the rows still holding the sentinel,
40+
* so it shrinks as the run commits chunks, and the plan the owner registers
41+
* for resume carries the default chunk size, not the one the run was started
42+
* with. Both used to change the chunk plan a resume recomputed, and the
43+
* runner refused `PLAN_CHANGED` a run the list had just called resumable.
44+
* Three more interrupted runs, made the same way:
45+
*
46+
* 4. killed in chunk 1 after chunk 0 committed (203 rows at the default
47+
* size): it resumes to completion, and chunk 0 is not run again;
48+
* 5. started with a chunk size of 2 (the `--chunk-size 2` an operator
49+
* passes) and killed in chunk 0: it resumes with the journal's size, two
50+
* chunks, not the one chunk the registered plan's default would make;
51+
* 6. the control: a run started by a plan whose step had another name — a
52+
* changed plan — is still refused `PLAN_CHANGED`, and nothing is written.
4353
*
4454
* Every boot runs in a hook: a case only reads what a boot printed or wrote.
4555
*/
@@ -97,9 +107,11 @@ function childEnv(overrides: Record<string, string | undefined>): Record<string,
97107
/**
98108
* The crash. Boots the data stack over the fixture database, writes the
99109
* sentinel rows, then runs the recorded-by plan under the real runner with a
100-
* forward that kills the process from inside chunk 0's transaction. The plan
101-
* keeps the owner's id, step names and chunk size, so the journal records the
102-
* hash the owner's registered plan computes on resume.
110+
* forward that runs the owner's forward for every chunk before
111+
* `FIXTURE_KILL_CHUNK` and kills the process from inside that chunk's
112+
* transaction. The plan keeps the owner's id; `FIXTURE_CHUNK_SIZE` starts it
113+
* at another size, as `--chunk-size` does, and `FIXTURE_STEP_NAME` renames its
114+
* step, which makes it a different plan from the one the owner registers.
103115
*/
104116
const CRASH_CHILD = `
105117
const rt = await import('@objectstack/runtime');
@@ -126,13 +138,20 @@ for (const [i, id] of ids.entries()) {
126138
operation_type: 'create', recorded_by: process.env.FIXTURE_SENTINEL, recorded_at: at,
127139
}, { context: { isSystem: true } });
128140
}
129-
const owned = mp.createRecordedBySentinelPlan();
141+
const chunkSize = process.env.FIXTURE_CHUNK_SIZE ? Number(process.env.FIXTURE_CHUNK_SIZE) : undefined;
142+
const killChunk = Number(process.env.FIXTURE_KILL_CHUNK ?? '0');
143+
const owned = mp.createRecordedBySentinelPlan(chunkSize === undefined ? {} : { chunkSize });
130144
const step = owned.steps[0];
131-
const crashing = { ...owned, steps: [{ ...step, forward: async (_rows, ctx) => {
132-
process.stderr.write('[fixture] run ' + ctx.runId + ' killed in chunk ' + ctx.chunkIndex + '\\n');
133-
process.kill(process.pid, 'SIGKILL');
134-
await new Promise(() => {});
135-
} }] };
145+
const crashing = { ...owned, steps: [{
146+
...step,
147+
...(process.env.FIXTURE_STEP_NAME ? { name: process.env.FIXTURE_STEP_NAME } : {}),
148+
forward: async (rows, ctx, engine) => {
149+
if (ctx.chunkIndex !== killChunk) return step.forward(rows, ctx, engine);
150+
process.stderr.write('[fixture] run ' + ctx.runId + ' killed in chunk ' + ctx.chunkIndex + '\\n');
151+
process.kill(process.pid, 'SIGKILL');
152+
await new Promise(() => {});
153+
},
154+
}] };
136155
await core.runMigrationJournal(ql, crashing);
137156
process.stderr.write('[fixture] the run finished, so nothing was interrupted\\n');
138157
process.exit(3);
@@ -232,69 +251,154 @@ async function readRows(dbFile: string, sql: string, bindings: unknown[] = []):
232251
}
233252
}
234253

235-
let root: string;
236-
let runId: string;
237-
let resumeList: { payload: any; exitCode: number };
238-
let resumeAct: { payload: any; exitCode: number };
239-
let resumed: ProjectCopy;
240-
let scanOutput: string;
241-
let freshOutput: string;
242-
243-
beforeAll(async () => {
244-
root = mkdtempSync(join(tmpdir(), 'os-21498-'));
245-
const origin = makeProject(root, 'origin');
254+
interface InterruptedFixture {
255+
/** The interrupted run's id, as the killed child printed it. */
256+
runId: string;
257+
/** The database the crash left behind. Copy it; never resume over it. */
258+
origin: ProjectCopy;
259+
}
246260

247-
// ── the crash ──────────────────────────────────────────────────────────────
261+
/**
262+
* Seed `ids` as sentinel rows and kill a recorded-by run in chunk `killChunk`.
263+
* See `CRASH_CHILD` for what the other two options change about the run.
264+
*/
265+
function interruptRun(
266+
name: string,
267+
ids: readonly string[],
268+
opts: { killChunk?: number; chunkSize?: number; stepName?: string } = {},
269+
): InterruptedFixture {
270+
const origin = makeProject(root, name);
271+
const killChunk = opts.killChunk ?? 0;
248272
const crash = spawnSync(process.execPath, ['--input-type=module', '-e', CRASH_CHILD], {
249273
cwd: CLI_ROOT,
250274
env: childEnv({
251275
OS_ARTIFACT_PATH: join(origin.dir, 'dist', 'objectstack.json'),
252276
FIXTURE_PROJECT: origin.dir,
253277
FIXTURE_DB: origin.dbFile,
254-
FIXTURE_IDS: JSON.stringify(SENTINEL_IDS),
278+
FIXTURE_IDS: JSON.stringify(ids),
255279
FIXTURE_SENTINEL: RECORDED_BY_SENTINEL,
280+
FIXTURE_KILL_CHUNK: String(killChunk),
281+
FIXTURE_CHUNK_SIZE: opts.chunkSize === undefined ? undefined : String(opts.chunkSize),
282+
FIXTURE_STEP_NAME: opts.stepName,
256283
}),
257284
encoding: 'utf8',
258285
timeout: CHILD_BUDGET_MS,
259286
});
260-
const killed = /\[fixture\] run (\S+) killed in chunk 0/.exec(crash.stderr ?? '');
287+
const killed = new RegExp(`\\[fixture\\] run (\\S+) killed in chunk ${killChunk}\\b`).exec(crash.stderr ?? '');
261288
if (crash.signal !== 'SIGKILL' || !killed) {
262289
throw new Error(
263-
`the fixture did not leave an interrupted run (status ${crash.status}, signal ${crash.signal})\n${crash.stderr}`,
290+
`the fixture '${name}' did not leave an interrupted run (status ${crash.status}, signal ${crash.signal})\n${crash.stderr}`,
264291
);
265292
}
266-
runId = killed[1];
293+
return { runId: killed[1], origin };
294+
}
267295

268-
// Each consumer gets its own copy: a resume concludes the run, a serve boot
269-
// writes rows of its own.
270-
resumed = makeProject(root, 'resumed');
271-
const scanned = makeProject(root, 'scanned');
272-
cpSync(join(origin.dir, 'data'), join(resumed.dir, 'data'), { recursive: true });
273-
cpSync(join(origin.dir, 'data'), join(scanned.dir, 'data'), { recursive: true });
274-
const fresh = makeProject(root, 'fresh');
296+
/** A fresh copy of the database a crash left behind. */
297+
function copyOf(fixture: InterruptedFixture, name: string): ProjectCopy {
298+
const copy = makeProject(root, name);
299+
cpSync(join(fixture.origin.dir, 'data'), join(copy.dir, 'data'), { recursive: true });
300+
return copy;
301+
}
275302

276-
// ── 1. the real command ────────────────────────────────────────────────────
303+
/** Each `argv` through the real `os migrate resume --json`, in order, over `project`. */
304+
async function resumeOver(project: ProjectCopy, ...argvs: string[][]): Promise<Array<{ payload: any; exitCode: number }>> {
305+
const out: Array<{ payload: any; exitCode: number }> = [];
277306
const savedEnv: Record<string, string | undefined> = {};
278307
const savedCwd = process.cwd();
279308
for (const key of [...OVERRIDING_ENV, 'OS_ARTIFACT_PATH'] as const) savedEnv[key] = process.env[key];
280309
try {
281310
for (const key of OVERRIDING_ENV) delete process.env[key];
282-
process.env.OS_ARTIFACT_PATH = join(resumed.dir, 'dist', 'objectstack.json');
283-
process.chdir(resumed.dir);
284-
const url = `file:${resumed.dbFile}`;
285-
resumeList = await resumeJson(['--database-url', url]);
286-
resumeAct = await resumeJson(['--run', runId, '--yes', '--database-url', url]);
311+
process.env.OS_ARTIFACT_PATH = join(project.dir, 'dist', 'objectstack.json');
312+
process.chdir(project.dir);
313+
const url = `file:${project.dbFile}`;
314+
for (const argv of argvs) out.push(await resumeJson([...argv, '--database-url', url]));
287315
} finally {
288316
process.chdir(savedCwd);
289317
for (const [key, value] of Object.entries(savedEnv)) {
290318
if (value === undefined) delete process.env[key];
291319
else process.env[key] = value;
292320
}
293321
}
322+
return out;
323+
}
324+
325+
/** Journal kinds and chunk indices of one run, in `seq` order. */
326+
async function journalOf(dbFile: string, id: string): Promise<Array<{ kind: string; chunk_index: number | null }>> {
327+
return readRows(dbFile, 'SELECT kind, chunk_index FROM sys_migration_journal WHERE run_id = ? ORDER BY seq', [id]);
328+
}
329+
330+
/** How many of `ids` still hold the sentinel. */
331+
async function sentinelCount(dbFile: string, ids: readonly string[]): Promise<number> {
332+
const rows = await readRows(
333+
dbFile,
334+
`SELECT COUNT(*) AS n FROM sys_metadata_history WHERE recorded_by = ? AND id IN (${ids.map(() => '?').join(', ')})`,
335+
[RECORDED_BY_SENTINEL, ...ids],
336+
);
337+
return Number(rows[0].n);
338+
}
339+
340+
/** One more row than the default chunk size holds, plus two: chunk 0 is 200 rows, chunk 1 is 3. */
341+
const COMMITTED_IDS = Array.from({ length: 203 }, (_, i) => `h_21528_c${String(i).padStart(3, '0')}`);
342+
const SIZED_IDS = ['h_21528_s_a', 'h_21528_s_b', 'h_21528_s_c'];
343+
const CHANGED_IDS = ['h_21528_x_a', 'h_21528_x_b', 'h_21528_x_c'];
344+
345+
let root: string;
346+
let runId: string;
347+
let resumeList: { payload: any; exitCode: number };
348+
let resumeAct: { payload: any; exitCode: number };
349+
let resumed: ProjectCopy;
350+
let scanOutput: string;
351+
let freshOutput: string;
352+
353+
/** The #21528 runs: each fixture, its resumed copy, and what list and act answered. */
354+
interface ResumedFixture extends InterruptedFixture {
355+
copy: ProjectCopy;
356+
list?: { payload: any; exitCode: number };
357+
act: { payload: any; exitCode: number };
358+
}
359+
let afterCommit: ResumedFixture;
360+
let sized: ResumedFixture;
361+
let changed: ResumedFixture;
362+
363+
beforeAll(async () => {
364+
root = mkdtempSync(join(tmpdir(), 'os-21498-'));
365+
366+
// ── the crash ──────────────────────────────────────────────────────────────
367+
const interrupted = interruptRun('origin', SENTINEL_IDS);
368+
runId = interrupted.runId;
369+
370+
// Each consumer gets its own copy: a resume concludes the run, a serve boot
371+
// writes rows of its own.
372+
resumed = copyOf(interrupted, 'resumed');
373+
const scanned = copyOf(interrupted, 'scanned');
374+
const fresh = makeProject(root, 'fresh');
375+
376+
// ── 1. the real command ────────────────────────────────────────────────────
377+
[resumeList, resumeAct] = await resumeOver(resumed, [], ['--run', runId, '--yes']);
294378

295379
// ── 2 and 3. the real serve boot ───────────────────────────────────────────
296380
scanOutput = await serveBoot(scanned);
297381
freshOutput = await serveBoot(fresh);
382+
383+
// ── 4, 5 and 6. runs whose recomputed chunk plan differs (#21528) ───────────
384+
const resumeFixture = async (
385+
fixture: InterruptedFixture,
386+
name: string,
387+
withList: boolean,
388+
): Promise<ResumedFixture> => {
389+
const copy = copyOf(fixture, name);
390+
const act = ['--run', fixture.runId, '--yes'];
391+
if (!withList) return { ...fixture, copy, act: (await resumeOver(copy, act))[0] };
392+
const [list, done] = await resumeOver(copy, [], act);
393+
return { ...fixture, copy, list, act: done };
394+
};
395+
afterCommit = await resumeFixture(interruptRun('after-commit', COMMITTED_IDS, { killChunk: 1 }), 'after-commit-resumed', true);
396+
sized = await resumeFixture(interruptRun('sized', SIZED_IDS, { chunkSize: 2 }), 'sized-resumed', true);
397+
changed = await resumeFixture(
398+
interruptRun('changed', CHANGED_IDS, { stepName: 'sys_metadata_history.recorded_by: an earlier step' }),
399+
'changed-resumed',
400+
false,
401+
);
298402
}, HOOK_TIMEOUT_MS);
299403

300404
afterAll(() => {
@@ -350,3 +454,57 @@ describe('os serve reports an interrupted run through the boot scan (#21498)', (
350454
expect(freshOutput).not.toMatch(/journal scan failed/i);
351455
});
352456
});
457+
458+
describe("os migrate resume completes a recorded-by run whose remaining rows shrank or whose chunk size was not the default (#21528)", () => {
459+
it('lists a run killed after a committed chunk as resumable, chunk 0 known committed', () => {
460+
const listed = afterCommit.list!.payload.interrupted.find((r: any) => r.runId === afterCommit.runId);
461+
expect(listed, JSON.stringify(afterCommit.list!.payload)).toBeDefined();
462+
expect(listed.committedChunks).toEqual([0]);
463+
expect(listed.unknownChunks).toEqual([1]);
464+
expect(listed.resumable).toBe(true);
465+
});
466+
467+
it('resumes that run to completion without running its committed chunk again', async () => {
468+
expect(afterCommit.act.exitCode, JSON.stringify(afterCommit.act.payload)).toBe(0);
469+
expect(afterCommit.act.payload).toMatchObject({
470+
runId: afterCommit.runId, status: 'completed', chunksTotal: 2, chunksCommitted: 2,
471+
});
472+
expect(await sentinelCount(afterCommit.copy.dbFile, COMMITTED_IDS)).toBe(0);
473+
474+
const journal = await journalOf(afterCommit.copy.dbFile, afterCommit.runId);
475+
// Chunk 0 committed before the kill and is started exactly once; chunk 1
476+
// is started by the killed run and again by the resume.
477+
expect(journal.filter((e) => e.kind === 'chunk_started').map((e) => e.chunk_index)).toEqual([0, 1, 1]);
478+
expect(journal.filter((e) => e.kind === 'chunk_done').map((e) => e.chunk_index)).toEqual([0, 1]);
479+
expect(journal.at(-1)!.kind).toBe('run_done');
480+
});
481+
482+
it('lists a run started with a chunk size of 2 as resumable', () => {
483+
const listed = sized.list!.payload.interrupted.find((r: any) => r.runId === sized.runId);
484+
expect(listed, JSON.stringify(sized.list!.payload)).toBeDefined();
485+
expect(listed.unknownChunks).toEqual([0]);
486+
expect(listed.resumable).toBe(true);
487+
});
488+
489+
it('resumes that run with the chunk size it started with, from the journal', async () => {
490+
expect(sized.act.exitCode, JSON.stringify(sized.act.payload)).toBe(0);
491+
// Three rows at size 2 are two chunks; the registered plan's default size
492+
// would have made one.
493+
expect(sized.act.payload).toMatchObject({
494+
runId: sized.runId, status: 'completed', chunksTotal: 2, chunksCommitted: 2,
495+
});
496+
expect(await sentinelCount(sized.copy.dbFile, SIZED_IDS)).toBe(0);
497+
const journal = await journalOf(sized.copy.dbFile, sized.runId);
498+
expect(journal.filter((e) => e.kind === 'chunk_done').map((e) => e.chunk_index)).toEqual([0, 1]);
499+
expect(journal.at(-1)!.kind).toBe('run_done');
500+
});
501+
502+
it('still refuses PLAN_CHANGED a run another plan started (the control), and writes nothing', async () => {
503+
expect(changed.act.exitCode, JSON.stringify(changed.act.payload)).toBe(1);
504+
expect(changed.act.payload.error).toMatch(/^Refused \(PLAN_CHANGED\)/);
505+
expect(await sentinelCount(changed.copy.dbFile, CHANGED_IDS)).toBe(CHANGED_IDS.length);
506+
expect((await journalOf(changed.copy.dbFile, changed.runId)).map((e) => e.kind)).toEqual([
507+
'run_started', 'chunk_started',
508+
]);
509+
});
510+
});

0 commit comments

Comments
 (0)