From 7d17ecc03ef1698586a547c324a6dd71931dbb35 Mon Sep 17 00:00:00 2001 From: 3metaJun <251347867+3metaJun@users.noreply.github.com> Date: Fri, 11 Sep 2026 03:03:04 +0800 Subject: [PATCH] fix: harden upstream sync recovery --- README.md | 28 ++- scripts/sync-upstream-transaction.mjs | 278 +++++++++++++++++++++----- scripts/sync-upstream.mjs | 9 +- scripts/sync-upstream.test.mjs | 156 +++++++++++++-- 4 files changed, 394 insertions(+), 77 deletions(-) diff --git a/README.md b/README.md index bca1ec6..8f572e7 100644 --- a/README.md +++ b/README.md @@ -276,13 +276,27 @@ content still matches the previous manifest. Pass `--apply --force` after reviewing a diff when overwriting or removing a locally changed file is intentional. -Only one `--apply` can run against a target at a time. Writes and removals are -staged and journaled before the target is changed; an ordinary failure rolls -the whole refresh back. If the process is interrupted, the next `--apply` -recovers the unfinished transaction before rebuilding the sync plan. Dry runs -and `check-upstream` refuse to inspect a target while a sync is active or needs -recovery. The transaction prevents persistent partial refreshes, but readers -that ignore the sync lock can still observe files changing during the commit. +All sync commands use one exclusive lock per target, so concurrent `--apply`, +dry-run, and `check-upstream` commands are serialized. The lock records the +host, platform, and process start identity. A lock from another runtime (for +example Windows versus WSL), or one whose owner cannot be verified, is left in +place and the command explains how to inspect and remove it after confirming +that no sync is running. + +Writes and removals are staged and journaled before the target is changed; an +ordinary failure rolls the whole refresh back. If the process is interrupted, +the next `--apply` recovers the unfinished transaction before rebuilding the +sync plan. If a user changed a target during an unfinished rollback, mstack +moves that file to `.mstack-sync-upstream/recovery//.user` +and restores the previous version; the command prints both paths. A +`COMMITTED` marker is authoritative, so later user edits are preserved while +transaction sidecars are cleaned up. Dry runs and `check-upstream` refuse to +inspect a target while a sync is active or needs recovery. + +On Windows, directory fsync is best-effort because Node cannot portably flush a +directory handle. The transaction provides process-crash recovery, but it does +not provide a power-loss durability guarantee on that platform. Readers that +ignore the sync lock can still observe files changing during the commit. `--apply` writes `profiles/upstream-manifest.json`. Its hashes describe the transformed upstream baseline, not package integrity. Files configured with diff --git a/scripts/sync-upstream-transaction.mjs b/scripts/sync-upstream-transaction.mjs index 4574cc7..a48ae5a 100644 --- a/scripts/sync-upstream-transaction.mjs +++ b/scripts/sync-upstream-transaction.mjs @@ -16,7 +16,9 @@ import { unlinkSync, writeFileSync, } from "node:fs"; +import { execFileSync } from "node:child_process"; import { createHash, randomUUID } from "node:crypto"; +import { hostname } from "node:os"; import { basename, dirname, isAbsolute, posix, relative, resolve, sep } from "node:path"; const LOCK_NAME = ".mstack-sync-upstream.lock"; @@ -24,6 +26,101 @@ const STATE_NAME = ".mstack-sync-upstream"; const SIDECAR_NAME = ".mstack-sync-upstream-tx"; const JOURNAL_VERSION = 1; const HASH_PATTERN = /^[0-9a-f]{64}$/; +const RUNTIME_HOST = hostname(); + +function processStartIdentity(pid) { + if (!Number.isSafeInteger(pid) || pid < 1) return undefined; + if (process.platform === "linux") { + try { + const stat = readFileSync(`/proc/${pid}/stat`, "utf8"); + const commandEnd = stat.lastIndexOf(")"); + if (commandEnd === -1) return undefined; + const fields = stat.slice(commandEnd + 2).trim().split(/\s+/); + return fields[19] ? `linux:${fields[19]}` : undefined; + } catch { + return undefined; + } + } + if (process.platform === "win32") { + try { + const start = execFileSync( + "powershell.exe", + [ + "-NoProfile", + "-NonInteractive", + "-Command", + `$p = Get-Process -Id ${pid} -ErrorAction SilentlyContinue; if ($null -ne $p) { $p.StartTime.ToFileTimeUtc() }`, + ], + { encoding: "utf8", stdio: ["ignore", "pipe", "ignore"] }, + ).trim(); + return start ? `windows:${start}` : undefined; + } catch { + return undefined; + } + } + if (process.platform === "darwin") { + try { + const start = execFileSync("ps", ["-p", String(pid), "-o", "lstart="], { + encoding: "utf8", + stdio: ["ignore", "pipe", "ignore"], + }).trim(); + return start ? `darwin:${start}` : undefined; + } catch { + return undefined; + } + } + return undefined; +} + +function processLiveness(pid) { + try { + process.kill(pid, 0); + return "alive"; + } catch (error) { + if (error.code === "ESRCH") return "dead"; + return "unknown"; + } +} + +function ownerLiveness(owner) { + const liveness = processLiveness(owner.pid); + if (liveness === "dead") return "dead"; + if (owner.host !== RUNTIME_HOST || owner.platform !== process.platform) return "unknown"; + if (typeof owner.start !== "string" || owner.start.length === 0) return "unknown"; + const actualStart = processStartIdentity(owner.pid); + if (!actualStart) return "unknown"; + return actualStart === owner.start ? "alive" : "dead"; +} + +function currentOwner() { + return { + version: 2, + pid: process.pid, + token: randomUUID(), + host: RUNTIME_HOST, + platform: process.platform, + start: processStartIdentity(process.pid) ?? null, + }; +} + +function assertOwnerShape(owner, label) { + if (!owner || !Number.isSafeInteger(owner.pid) || owner.pid < 1 || + typeof owner.token !== "string" || owner.token.length === 0) { + throw new Error(`Invalid upstream sync ${label}.`); + } + if (owner.version !== undefined && ![1, 2].includes(owner.version)) { + throw new Error(`Invalid upstream sync ${label}.`); + } + if (owner.host !== undefined && typeof owner.host !== "string") { + throw new Error(`Invalid upstream sync ${label}.`); + } + if (owner.platform !== undefined && typeof owner.platform !== "string") { + throw new Error(`Invalid upstream sync ${label}.`); + } + if (owner.start !== undefined && owner.start !== null && typeof owner.start !== "string") { + throw new Error(`Invalid upstream sync ${label}.`); + } +} function digest(value) { return createHash("sha256").update(value).digest("hex"); @@ -190,6 +287,20 @@ function assertTargetDoesNotOverlap(targets, target, label) { if (overlap) throw new Error(`Overlapping targets in ${label}: ${overlap} and ${target}`); } +function assertReservedTransactionTarget(target, label) { + if (target.split("/").some((part) => [LOCK_NAME, STATE_NAME, SIDECAR_NAME].includes(part))) { + throw new Error(`Upstream sync ${label} uses a reserved transaction path: ${target}`); + } +} + +function targetRootIdentity(path) { + let normalized = path.replaceAll("\\", "/"); + const windowsDrive = normalized.match(/^([A-Za-z]):(?:\/|$)/); + if (windowsDrive) normalized = `/mnt/${windowsDrive[1].toLowerCase()}${normalized.slice(2)}`; + normalized = normalized.replace(/\/+$/, ""); + return process.platform === "win32" || normalized.startsWith("/mnt/") ? normalized.toLowerCase() : normalized; +} + function assertFingerprintShape(value, label) { if (!value || typeof value !== "object" || !["absent", "file"].includes(value.kind)) { throw new Error(`Invalid upstream sync transaction ${label}.`); @@ -233,12 +344,10 @@ function readJournal(targetRoot, transactionPath, id) { if (!journal || typeof journal !== "object" || journal.version !== JOURNAL_VERSION || journal.id !== id) { throw new Error(`Invalid upstream sync transaction journal at ${journalPath}.`); } - if (!journal.owner || !Number.isSafeInteger(journal.owner.pid) || journal.owner.pid < 1 || - typeof journal.owner.token !== "string" || journal.owner.token.length === 0) { - throw new Error(`Invalid upstream sync transaction owner at ${journalPath}.`); - } - if (journal.targetRoot !== realpathSync(targetRoot)) { - throw new Error(`Upstream sync transaction journal targets a different checkout: ${journalPath}`); + assertOwnerShape(journal.owner, `transaction owner at ${journalPath}`); + if (typeof journal.targetRoot !== "string" || journal.targetRoot.length === 0 || + targetRootIdentity(journal.targetRoot) !== targetRootIdentity(realpathSync(targetRoot))) { + throw new Error(`Invalid upstream sync transaction target root at ${journalPath}.`); } if (!Array.isArray(journal.createdDirectories) || !Array.isArray(journal.operations) || !journal.operations.length) { throw new Error(`Invalid upstream sync transaction journal at ${journalPath}.`); @@ -247,6 +356,7 @@ function readJournal(targetRoot, transactionPath, id) { const createdDirectories = new Set(); for (const [index, path] of journal.createdDirectories.entries()) { normalizedRelativePath(path, `createdDirectories[${index}]`); + assertReservedTransactionTarget(path, `createdDirectories[${index}]`); resolveInside(targetRoot, path, `createdDirectories[${index}]`); if (createdDirectories.has(path)) throw new Error(`Duplicate directory in upstream sync transaction: ${path}`); createdDirectories.add(path); @@ -259,6 +369,7 @@ function readJournal(targetRoot, transactionPath, id) { throw new Error(`Invalid upstream sync transaction operation ${index}.`); } normalizedRelativePath(operation.target, `operations[${index}].target`); + assertReservedTransactionTarget(operation.target, `operations[${index}].target`); resolveInside(targetRoot, operation.target, `operations[${index}].target`); assertTargetDoesNotOverlap(targets, operation.target, "upstream sync transaction"); targets.push(operation.target); @@ -360,21 +471,47 @@ function cleanupAbandonedJournalWrite(targetRoot, transactionPath) { } function transactionPaths(targetRoot, journal) { - return journal.operations.map((operation) => ({ + return journal.operations.map((operation, index) => ({ ...operation, + index, targetPath: resolveInside(targetRoot, operation.target, "operation target"), stagePath: operation.stage === null ? undefined : resolveInside(targetRoot, operation.stage, "operation stage"), backupPath: resolveInside(targetRoot, operation.backup, "operation backup"), })); } -function rollbackOperation(operation) { +function quarantineUnknownTarget(targetRoot, transactionPath, operation) { + const { stateRoot } = transactionRoots(targetRoot); + const recoveryRoot = resolve(stateRoot, "recovery", basename(transactionPath)); + resolveInside( + targetRoot, + relative(targetRoot, recoveryRoot).split(sep).join("/"), + "upstream sync recovery directory", + ); + mkdirSync(recoveryRoot, { recursive: true, mode: 0o700 }); + fsyncDirectory(dirname(recoveryRoot)); + fsyncDirectory(recoveryRoot); + const recoveryPath = resolve(recoveryRoot, `${operation.index}.user`); + if (lstatIfPresent(recoveryPath)) { + throw new Error(`Recovery file already exists: ${recoveryPath}`); + } + durableRename(operation.targetPath, recoveryPath); + return { path: recoveryPath, target: operation.target }; +} + +function rollbackOperation(targetRoot, transactionPath, operation) { const backup = lstatIfPresent(operation.backupPath); - const current = fingerprint(operation.targetPath); + let current = fingerprint(operation.targetPath); + let quarantined; if (operation.before.kind === "file") { if (backup) { assertFingerprint(operation.backupPath, operation.before, "Upstream sync backup"); + if (current.kind !== "absent" && + !fingerprintsEqual(current, operation.before) && !fingerprintsEqual(current, operation.after)) { + quarantined = quarantineUnknownTarget(targetRoot, transactionPath, operation); + current = { kind: "absent" }; + } if (current.kind === "absent") { durableRename(operation.backupPath, operation.targetPath); } else if (fingerprintsEqual(current, operation.after)) { @@ -390,26 +527,32 @@ function rollbackOperation(operation) { } } else { if (backup) throw new Error(`Unexpected backup for a newly created target: ${operation.backupPath}`); - if (fingerprintsEqual(current, operation.after)) durableUnlink(operation.targetPath); - else if (current.kind !== "absent") { - throw new Error(`Refusing to remove an unknown file while rolling back: ${operation.targetPath}`); + if (fingerprintsEqual(current, operation.after)) { + durableUnlink(operation.targetPath); + } else if (current.kind !== "absent") { + quarantined = quarantineUnknownTarget(targetRoot, transactionPath, operation); + current = { kind: "absent" }; } } if (operation.stagePath) removeExpectedFile(operation.stagePath, operation.after, "Upstream sync staged file"); assertFingerprint(operation.targetPath, operation.before, "Rolled back upstream sync target"); + return quarantined; } -function rollbackJournal(targetRoot, journal) { +function rollbackJournal(targetRoot, transactionPath, journal) { const errors = []; + const quarantined = []; for (const operation of transactionPaths(targetRoot, journal).reverse()) { try { - rollbackOperation(operation); + const result = rollbackOperation(targetRoot, transactionPath, operation); + if (result) quarantined.push(result); } catch (error) { errors.push(`${operation.target}: ${error.message}`); } } if (errors.length) throw new Error(errors.join("\n")); + return quarantined; } function cleanupTransaction(targetRoot, transactionPath, journal) { @@ -452,15 +595,6 @@ function validateCommittedTargets(targetRoot, journal) { } } -function processIsAlive(pid) { - try { - process.kill(pid, 0); - return true; - } catch (error) { - return error.code !== "ESRCH"; - } -} - function lockOwnerText(owner) { return JSON.stringify(owner); } @@ -482,7 +616,17 @@ function createOwnedLock(lockPath, owner) { let published = false; try { writeDurableFile(temporaryPath, `${lockOwnerText(owner)}\n`); - linkSync(temporaryPath, lockPath); + try { + linkSync(temporaryPath, lockPath); + } catch (error) { + if (["EPERM", "EOPNOTSUPP", "ENOTSUP", "EXDEV"].includes(error.code)) { + throw new Error( + `Upstream sync requires hard links for its lock on this filesystem: ${lockPath}`, + { cause: error }, + ); + } + throw error; + } published = true; fsyncDirectory(dirname(lockPath)); durableUnlink(temporaryPath); @@ -509,14 +653,21 @@ function removeDeadLockInitializers(targetRoot) { const prefix = `${LOCK_NAME}.`; for (const entry of readdirSync(targetRoot, { withFileTypes: true })) { if (!entry.isFile() || !entry.name.startsWith(prefix) || !entry.name.endsWith(".tmp")) continue; - const match = entry.name.match(/^\.mstack-sync-upstream\.lock\.([1-9]\d*)\.([0-9a-f-]+)\.tmp$/); - if (!match) throw new Error(`Invalid upstream sync lock initializer: ${resolve(targetRoot, entry.name)}`); - const pid = Number(match[1]); - if (!Number.isSafeInteger(pid) || processIsAlive(pid)) { - throw new Error(`Another upstream sync is initializing its lock: ${resolve(targetRoot, entry.name)}`); + const match = entry.name.match(/^\.mstack-sync-upstream\.lock\.(?:takeover\.)?([1-9]\d*)\.([0-9a-f-]+)\.tmp$/); + if (!match) continue; + const initializerPath = resolve(targetRoot, entry.name); + let initializer; + try { + initializer = JSON.parse(readRegularFile(initializerPath, "upstream sync lock initializer").toString("utf8")); + assertOwnerShape(initializer, `lock initializer at ${initializerPath}`); + } catch (error) { + throw new Error(`Invalid upstream sync lock initializer: ${initializerPath}`, { cause: error }); + } + if (ownerLiveness(initializer) !== "dead") { + throw new Error(`Another upstream sync is initializing its lock: ${initializerPath}`); } try { - durableUnlink(resolve(targetRoot, entry.name)); + durableUnlink(initializerPath); } catch (error) { if (error.code !== "ENOENT") throw error; } @@ -534,8 +685,8 @@ function acquireTakeoverGuard(guardPath, owner) { } catch (error) { throw new Error(`Invalid upstream sync takeover claim: ${reclaimedPath}`, { cause: error }); } - if (!reclaimedOwner || !Number.isSafeInteger(reclaimedOwner.pid) || reclaimedOwner.pid < 1 || - processIsAlive(reclaimedOwner.pid)) { + assertOwnerShape(reclaimedOwner, `takeover claim at ${reclaimedPath}`); + if (ownerLiveness(reclaimedOwner) !== "dead") { throw new Error(`Another upstream sync is recovering the lock: ${reclaimedPath}`); } durableUnlink(reclaimedPath); @@ -551,8 +702,12 @@ function acquireTakeoverGuard(guardPath, owner) { } catch (error) { throw new Error(`Invalid or initializing upstream sync takeover guard: ${guardPath}`, { cause: error }); } - if (!current || current.version !== 1 || !Number.isSafeInteger(current.pid) || current.pid < 1 || - typeof current.token !== "string" || current.token.length === 0 || processIsAlive(current.pid)) { + try { + assertOwnerShape(current, `takeover guard at ${guardPath}`); + } catch (error) { + throw new Error(`Invalid upstream sync takeover guard: ${guardPath}`, { cause: error }); + } + if (ownerLiveness(current) !== "dead") { throw new Error(`Another upstream sync is recovering the lock: ${guardPath}`); } const reclaimedPath = `${guardPath}.${owner.token}.reclaimed`; @@ -576,7 +731,7 @@ function acquireTakeoverGuard(guardPath, owner) { export function acquireUpstreamSyncLock(targetRoot, { recoverStale = true } = {}) { const lockPath = resolve(targetRoot, LOCK_NAME); - const owner = { version: 1, pid: process.pid, token: randomUUID() }; + const owner = currentOwner(); removeDeadLockInitializers(targetRoot); try { return { lockPath, owner: createOwnedLock(lockPath, owner) }; @@ -590,13 +745,18 @@ export function acquireUpstreamSyncLock(targetRoot, { recoverStale = true } = {} } catch (error) { throw new Error(`Invalid or initializing upstream sync lock: ${lockPath}`, { cause: error }); } - if (!current || current.version !== 1 || !Number.isSafeInteger(current.pid) || current.pid < 1 || - typeof current.token !== "string" || current.token.length === 0) { - throw new Error(`Invalid upstream sync lock: ${lockPath}`); + try { + assertOwnerShape(current, `lock at ${lockPath}`); + } catch (error) { + throw new Error(`Invalid upstream sync lock: ${lockPath}`, { cause: error }); } - if (processIsAlive(current.pid)) { + const currentLiveness = ownerLiveness(current); + if (currentLiveness === "alive") { throw new Error(`Another live upstream sync owns the lock: ${lockPath}`); } + if (currentLiveness === "unknown") { + throw new Error(`Cannot verify the upstream sync lock owner; verify no sync is running, then remove ${lockPath} and retry.`); + } if (!recoverStale) { throw new Error(`Target has an active or interrupted upstream sync; rerun sync-upstream with --apply: ${lockPath}`); } @@ -611,13 +771,18 @@ export function acquireUpstreamSyncLock(targetRoot, { recoverStale = true } = {} } catch (error) { throw new Error(`Invalid or initializing upstream sync lock: ${lockPath}`, { cause: error }); } - if (!guardedCurrent || guardedCurrent.version !== 1 || !Number.isSafeInteger(guardedCurrent.pid) || - guardedCurrent.pid < 1 || typeof guardedCurrent.token !== "string" || guardedCurrent.token.length === 0) { - throw new Error(`Invalid upstream sync lock: ${lockPath}`); + try { + assertOwnerShape(guardedCurrent, `lock at ${lockPath}`); + } catch (error) { + throw new Error(`Invalid upstream sync lock: ${lockPath}`, { cause: error }); } - if (processIsAlive(guardedCurrent.pid)) { + const guardedLiveness = ownerLiveness(guardedCurrent); + if (guardedLiveness === "alive") { throw new Error(`Another live upstream sync owns the lock: ${lockPath}`); } + if (guardedLiveness === "unknown") { + throw new Error(`Cannot verify the upstream sync lock owner; verify no sync is running, then remove ${lockPath} and retry.`); + } durableUnlink(lockPath); try { return { lockPath, owner: createOwnedLock(lockPath, owner) }; @@ -656,6 +821,7 @@ export function assertUpstreamSyncTargetReadable(targetRoot, heldLock) { export function recoverUpstreamSyncTransactions(targetRoot, lock) { assertLockOwned(lock); let recovered = 0; + const quarantined = []; for (const id of listTransactions(targetRoot)) { assertLockOwned(lock); normalizedRelativePath(id, "transaction id"); @@ -673,16 +839,23 @@ export function recoverUpstreamSyncTransactions(targetRoot, lock) { if (journal.owner.pid === process.pid && journal.owner.token !== lock.owner.token) { throw new Error(`Upstream sync transaction owner token does not match the current lock: ${transactionPath}`); } - if (journal.owner.pid !== process.pid && processIsAlive(journal.owner.pid)) { + const transactionLiveness = ownerLiveness(journal.owner); + if (transactionLiveness === "alive" && journal.owner.token !== lock.owner.token) { throw new Error(`Upstream sync transaction is still owned by live process ${journal.owner.pid}: ${transactionPath}`); } - if (markerExists(transactionPath, "COMMITTED")) validateCommittedTargets(targetRoot, journal); - else rollbackJournal(targetRoot, journal); + if (transactionLiveness === "unknown") { + throw new Error( + `Cannot verify the upstream sync transaction owner; verify no sync is running, then inspect ${transactionPath}.`, + ); + } + if (!markerExists(transactionPath, "COMMITTED")) { + quarantined.push(...rollbackJournal(targetRoot, transactionPath, journal).map((item) => ({ id, ...item }))); + } cleanupTransaction(targetRoot, transactionPath, journal); recovered += 1; } assertLockOwned(lock); - return recovered; + return { count: recovered, quarantined }; } function missingParentDirectories(targetRoot, targetPath) { @@ -709,9 +882,7 @@ function buildJournal(targetRoot, id, changes, lock) { throw new Error(`Invalid upstream sync change ${index}.`); } const target = normalizedRelativePath(change.target, `change ${index} target`); - if (target.split("/").some((part) => [LOCK_NAME, STATE_NAME, SIDECAR_NAME].includes(part))) { - throw new Error(`Upstream sync target uses a reserved transaction path: ${target}`); - } + assertReservedTransactionTarget(target, `change ${index} target`); assertTargetDoesNotOverlap(targets, target, "upstream sync transaction"); targets.push(target); const targetPath = resolveInside(targetRoot, target, `change ${index} target`); @@ -849,6 +1020,7 @@ export function applyUpstreamSyncTransaction(targetRoot, changes, lock) { const journal = buildJournal(targetRoot, id, changes, lock); const { transactionsRoot } = transactionRoots(targetRoot); const transactionPath = resolve(transactionsRoot, id); + let quarantined = []; try { createJournal(targetRoot, journal); assertLockOwned(lock); @@ -862,7 +1034,7 @@ export function applyUpstreamSyncTransaction(targetRoot, changes, lock) { } try { if (existsSync(resolve(transactionPath, "journal.json"))) { - rollbackJournal(targetRoot, journal); + quarantined = rollbackJournal(targetRoot, transactionPath, journal); cleanupTransaction(targetRoot, transactionPath, journal); } else if (existsSync(transactionPath)) { cleanupAbandonedJournalWrite(targetRoot, transactionPath); @@ -873,6 +1045,10 @@ export function applyUpstreamSyncTransaction(targetRoot, changes, lock) { { cause: error }, ); } + if (quarantined.length) { + const details = quarantined.map(({ target, path }) => `Preserved ${target} at ${path}.`).join("\n"); + throw new Error(`${error.message}\n${details}`, { cause: error }); + } throw error; } } diff --git a/scripts/sync-upstream.mjs b/scripts/sync-upstream.mjs index b74122f..4837743 100644 --- a/scripts/sync-upstream.mjs +++ b/scripts/sync-upstream.mjs @@ -62,10 +62,13 @@ const releaseReadLockAtExit = () => { if (!apply) process.once("exit", releaseReadLockAtExit); try { -const recoveredTransactions = apply ? recoverUpstreamSyncTransactions(targetRoot, syncLock) : 0; +const recovery = apply ? recoverUpstreamSyncTransactions(targetRoot, syncLock) : { count: 0, quarantined: [] }; if (!apply) assertUpstreamSyncTargetReadable(targetRoot, syncLock); -if (recoveredTransactions) { - console.log(`Recovered ${recoveredTransactions} interrupted upstream sync transaction(s).`); +if (recovery.count) { + console.log(`Recovered ${recovery.count} interrupted upstream sync transaction(s).`); +} +for (const { id, target, path } of recovery.quarantined) { + console.warn(`Preserved the user-modified ${target} from transaction ${id} at ${path}.`); } if (!existsSync(sourceRoot)) throw new Error(`Upstream checkout not found: ${sourceRoot}`); const upstreamsPath = resolveInside(targetRoot, "profiles/upstreams.json", "upstream profile"); diff --git a/scripts/sync-upstream.test.mjs b/scripts/sync-upstream.test.mjs index dac6935..57d32b0 100644 --- a/scripts/sync-upstream.test.mjs +++ b/scripts/sync-upstream.test.mjs @@ -6,13 +6,14 @@ import { mkdtempSync, readFileSync, readdirSync, + realpathSync, rmSync, unlinkSync, writeFileSync, } from "node:fs"; import { spawn, spawnSync } from "node:child_process"; import { createHash } from "node:crypto"; -import { tmpdir } from "node:os"; +import { hostname, tmpdir } from "node:os"; import { join, relative, resolve } from "node:path"; import test from "node:test"; @@ -766,6 +767,50 @@ test("recovers a stale transaction after the process exits between a rename and } }); +test("quarantines user edits before rolling back an interrupted transaction", () => { + const { root, source, target } = fixture(); + try { + applyBaseline(source, target); + advanceSource(source, target, () => { + write(join(source, "automations", "benny", "README.md"), "changed pstack automation\n"); + }, "prepare interrupted sync with user edit"); + + const interrupted = runWithFaults( + synchronizer, + ["--source", source, "--target", target, "--apply", "--force"], + root, + [{ + method: "renameSync", + phase: "after", + fromIncludes: "/.mstack-sync-upstream-tx/", + toEndsWith: "/automations/benny/README.md", + action: "exit", + exitCode: 86, + }], + ); + assert.equal(interrupted.status, 86, interrupted.stderr); + + const transactionsRoot = join(target, ".mstack-sync-upstream", "transactions"); + const [transactionId] = readdirSync(transactionsRoot); + write(join(target, "automations", "benny", "README.md"), "user edit during recovery\n"); + + const recovered = run(synchronizer, ["--source", source, "--target", target, "--apply", "--force"]); + assert.equal(recovered.status, 0, recovered.stderr); + assert.match(recovered.stderr, new RegExp(`automations/benny/README\\.md.*${transactionId}.*recovery`)); + assert.equal( + readFileSync(join(target, "automations", "benny", "README.md"), "utf8"), + "changed mstack automation\n", + ); + assert.equal( + readFileSync(join(target, ".mstack-sync-upstream", "recovery", transactionId, "0.user"), "utf8"), + "user edit during recovery\n", + ); + assert.equal(existsSync(transactionsRoot) ? readdirSync(transactionsRoot).length : 0, 0); + } finally { + rmSync(root, { recursive: true, force: true }); + } +}); + test("cleans up a committed transaction without rolling it back after the committed marker rename", () => { const { root, source, target } = fixture(); try { @@ -802,30 +847,22 @@ test("cleans up a committed transaction without rolling it back after the commit "committed mstack automation\n", ); - const recovered = runWithFaults( - synchronizer, - ["--source", source, "--target", target, "--apply", "--force"], - root, - [{ - method: "renameSync", - phase: "before", - fromIncludes: "/.mstack-sync-upstream-tx/", - toEndsWith: "/automations/benny/README.md", - message: "committed transaction attempted rollback", - }], - ); - assert.equal(recovered.status, 0, recovered.stderr); + write(join(target, "automations", "benny", "README.md"), "user edit after commit\n"); + + const recovered = run(synchronizer, ["--source", source, "--target", target, "--apply"]); + assert.notEqual(recovered.status, 0); assert.doesNotMatch(recovered.stderr, /committed transaction attempted rollback/); assert.equal( readFileSync(join(target, "automations", "benny", "README.md"), "utf8"), - "committed mstack automation\n", + "user edit after commit\n", ); assert.equal( readFileSync(join(target, "docs", "guide", "02-meta-mode.md"), "utf8"), "Committed /meta-mode guide for mstack.\nHarness confirms.\n", ); const checked = run(checker, ["--source", source, "--target", target, "--strict"]); - assert.equal(checked.status, 0, checked.stderr); + assert.notEqual(checked.status, 0); + assert.match(checked.stderr, /automations\/benny\/README\.md/); assertNoTransactionState(target); } finally { rmSync(root, { recursive: true, force: true }); @@ -1000,6 +1037,25 @@ test("takes over dead sync locks and stale takeover guards", () => { } }); +test("takes over a lock whose PID was reused by a different process instance", () => { + const { root, source, target } = fixture(); + try { + write(join(target, ".mstack-sync-upstream.lock"), `${JSON.stringify({ + version: 2, + pid: process.pid, + token: "reused-pid-owner", + host: hostname(), + platform: process.platform, + start: "different-process-start", + })}\n`); + const result = run(synchronizer, ["--source", source, "--target", target, "--apply"]); + assert.equal(result.status, 0, result.stderr); + assertNoTransactionState(target); + } finally { + rmSync(root, { recursive: true, force: true }); + } +}); + for (const [kind, description] of [["malformed", "malformed"], ["escaping", "path-escaping"]]) { test(`refuses a ${description} interrupted transaction journal without mutating files`, () => { const { root, source, target } = fixture(); @@ -1056,3 +1112,71 @@ for (const [kind, description] of [["malformed", "malformed"], ["escaping", "pat } }); } + +test("rejects a journal operation that targets reserved transaction state", () => { + const { root, source, target } = fixture(); + try { + applyBaseline(source, target); + const transactionId = "invalid-reserved"; + const transaction = transactionDirectory(target, transactionId); + mkdirSync(transaction, { recursive: true }); + write(join(transaction, "COMMITTING"), "\n"); + write(join(transaction, "journal.json"), `${JSON.stringify({ + version: 1, + id: transactionId, + owner: { pid: 2147483647, token: "reserved-test-owner" }, + targetRoot: realpathSync(target), + createdDirectories: [], + operations: [{ + kind: "write", + target: ".mstack-sync-upstream.lock", + stage: ".mstack-sync-upstream-tx/invalid-reserved/0.new", + backup: ".mstack-sync-upstream-tx/invalid-reserved/0.old", + before: { kind: "absent" }, + after: { kind: "file", sha256: sha256("reserved replacement\n"), mode: 0o600 }, + }], + }, null, 2)}\n`); + + const before = snapshotTree(target); + const result = run(synchronizer, ["--source", source, "--target", target, "--apply"]); + assert.notEqual(result.status, 0); + assert.match(result.stderr, /reserved transaction path/i); + assert.deepEqual(snapshotTree(target), before); + } finally { + rmSync(root, { recursive: true, force: true }); + } +}); + +test("rejects a transaction journal copied from a different checkout", () => { + const { root, source, target } = fixture(); + try { + applyBaseline(source, target); + const transactionId = "invalid-other-target"; + const transaction = transactionDirectory(target, transactionId); + mkdirSync(transaction, { recursive: true }); + write(join(transaction, "COMMITTING"), "\n"); + write(join(transaction, "journal.json"), `${JSON.stringify({ + version: 1, + id: transactionId, + owner: { pid: 2147483647, token: "other-target-owner" }, + targetRoot: join(root, "other-checkout"), + createdDirectories: [], + operations: [{ + kind: "manifest", + target: "profiles/upstream-manifest.json", + stage: `.mstack-sync-upstream-tx/${transactionId}/0.new`, + backup: `.mstack-sync-upstream-tx/${transactionId}/0.old`, + before: { kind: "absent" }, + after: { kind: "file", sha256: sha256("other checkout\n"), mode: 0o600 }, + }], + }, null, 2)}\n`); + + const before = snapshotTree(target); + const result = run(synchronizer, ["--source", source, "--target", target, "--apply"]); + assert.notEqual(result.status, 0); + assert.match(result.stderr, /target root|different checkout/i); + assert.deepEqual(snapshotTree(target), before); + } finally { + rmSync(root, { recursive: true, force: true }); + } +});