commit cc7b95b3d5a63c9c23b961779b9706c36899c851
parent a587e2abcc7e3147a0147095cf32bdbebc4b7ca5
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Thu, 17 Sep 2026 13:42:31 -0400
relocate: an in-flight move stops the lanes instead of racing them
The omnimirror incident (2026-09-13) in one sentence: a relocate copy
started while an auto-digest batch was running, a unit wrote a sidecar into
`data/` after rsync was already past that directory, and the verify refused a
131 GB transfer at its last step.
Both guards that could have seen it run ONCE, at the start of a run:
assertChannelMediaReachable in runOperationBatch, and the reservation in the
auto runner. A relocation marker is the signal every other guard in the system
already honours, so both places now re-ask it at the point they would do work.
operationBatch.next() returns null on a marker, which ENDS the run cleanly
(mechanic 2 in that file: zero is a hold, null is a stop) — the batch is
per-channel, so nothing further in it could have been dispatched anyway.
The auto runner records a skip, not a failure, so the platform backoff is
untouched. The trade-off is stated at the site: the `finally` marks the video
completed for the session, so one unit per channel is retired early and the
relocation finishing re-enables the channel on the next sweep. That is a far
cheaper loss than a copy refused after it has moved everything.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
2 files changed, 57 insertions(+), 1 deletion(-)
diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts
@@ -80,6 +80,7 @@ import { downloadQueueKey } from "../lib/queueKeys";
import { isGateHeld } from "../lib/pauseGates";
import {
inspectChannelMedia,
+ readRelocationMarker,
type ChannelMediaStatus,
} from "../lib/channelMedia";
import {
@@ -1646,6 +1647,33 @@ async function runLoop(
const { pick, channelSlug, unitPlatform } = picked;
let result: UnitResult = { outcome: "failed" };
try {
+ // A RELOCATION THAT STARTED BETWEEN THE PICK AND THE RUN.
+ //
+ // The pending lists are built from snapshots and the reservation above is
+ // taken before this runs; a relocate job can begin in that window. Its
+ // rsync is copying the very `data/` this unit is about to write into —
+ // the omnimirror shape (2026-09-13), where a sidecar written mid-copy
+ // left the verify refusing a 131 GB transfer. A marker is the signal, the
+ // same one every other guard already honours.
+ //
+ // TRADE-OFF, ACCEPTED: the `finally` below marks the video completed for
+ // the session, so this unit is not offered again until the runner
+ // restarts or the snapshot regenerates. The window is one unit per
+ // channel per session, and the relocation finishing re-enables the
+ // channel on the next sweep — which is a far cheaper loss than a copy
+ // refused at its last step.
+ const marker = await readRelocationMarker(paths, channelSlug);
+ if (marker) {
+ onLog(
+ `Auto-${kind}: skipping ${channelSlug}/${pick.videoId} — a relocation ` +
+ `(${marker.direction}) to ${marker.target} is in flight.`,
+ );
+ // `return` still runs the `finally`, which is what releases the
+ // reservation and folds this outcome — a skip, not a failure, so it
+ // never touches the platform backoff.
+ result = { outcome: "skipped" };
+ return;
+ }
result = isOperationLane(kind)
? await runOperationPick(picked, runSignal)
: await launchUnit({
diff --git a/common/controller/operationBatch.ts b/common/controller/operationBatch.ts
@@ -38,7 +38,10 @@ import { open, readFile, readdir } from "node:fs/promises";
import type { Paths } from "../lib/paths";
import { getSettings, type SiteSettings } from "../lib/settings";
import { isGateHeld } from "../lib/pauseGates";
-import { assertChannelMediaReachable } from "../lib/channelMedia";
+import {
+ assertChannelMediaReachable,
+ readRelocationMarker,
+} from "../lib/channelMedia";
import { runPool } from "../jobs/concurrentRunner";
import type { TaskTracker } from "../jobs/taskHooks";
import type { JobProgress } from "../jobs/registry";
@@ -1696,6 +1699,31 @@ export async function runOperationBatch(
cursor++;
continue;
}
+ // GUARD 5: A RELOCATION THAT STARTED WHILE THIS BATCH WAS RUNNING.
+ //
+ // The reachability guard above runs ONCE, at start. The omnimirror
+ // incident (2026-09-13) is what that misses: a relocate copy began after
+ // this batch did, a unit wrote a sidecar into `data/` while rsync was
+ // already past that directory, and the copy's verify then refused. A
+ // marker is the one signal that says "these bytes are being copied right
+ // now" — every other guard in the system already treats a channel
+ // carrying one as in-transition, and this is the same rule applied per
+ // pull rather than per job.
+ //
+ // Returning null ENDS the run cleanly (see mechanic 2 in the header: zero
+ // is a hold, null is a stop). That is deliberate — the batch is
+ // per-channel, so nothing further in this run could be dispatched anyway,
+ // and the next sweep re-derives everything from disk once the move is
+ // done.
+ const marker = await readRelocationMarker(opts.paths, opts.channelSlug);
+ if (marker) {
+ log(
+ `Stopping: a relocation (${marker.direction}) of ${opts.channelSlug} ` +
+ `to ${marker.target} is in flight — the rest of this batch is left ` +
+ `for the next pass.`,
+ );
+ return null;
+ }
// RE-DERIVED FROM DISK, every pull. A restart, a concurrent lane, or a
// share that landed while this job ran are all visible here. The run
// re-asks the same question; this pull is what lets a non-dispatch cost a