commit 2a72cda40dce8339a22f2a0084e86209a3df3893
parent 503d46b4352063b0b96b53881008a58243024ad8
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Wed, 29 Jul 2026 11:49:32 -0400
Two bugs in the sweep controls that would have shown up on first use
**The sweep would have refused to start.** `writeSettings` is async and neither
`startDigestSweep` nor `stopDigestSweep` awaited it. The loop reads
`sweepEnabled` at the top of its very first pass, so the un-awaited write raced
it: the sweep would log "disarmed in settings" and exit immediately, having
done nothing, while the flag it just wrote said it was armed. The stop path had
the mirror-image race against the boot-resume hook.
**"Stop" would have looked wedged.** Draining only the orchestrator means the
sweep stops after the current CHANNEL finishes — and the heaviest channel is 8.7
days of audio. Stop now drains the per-channel job as well, so it takes effect
within one video. Still a drain, not a cancel: the video in flight completes and
nothing part-generated is thrown away.
Neither was reachable by the type checker or the test suite, and both would have
been the very first thing an operator hit.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Diffstat:
2 files changed, 38 insertions(+), 8 deletions(-)
diff --git a/common/controller/digestSweep.ts b/common/controller/digestSweep.ts
@@ -62,7 +62,14 @@ const IDLE_POLL_MS = 60_000;
// corpus the sweep cannot make progress on, and both want a human.
const MAX_BARREN_PASSES = 3;
-type SweepLive = { jobId: string; startedAt: number };
+type SweepLive = {
+ jobId: string;
+ startedAt: number;
+ // The per-channel job the sweep is currently waiting on. Tracked so a stop
+ // can drain THAT too — without it, "stop" means "after the current channel
+ // finishes", and the heaviest channel is 8.7 days long.
+ channelJobId: string | null;
+};
type SweepSingleton = { live: SweepLive | null };
declare global {
@@ -101,6 +108,9 @@ export async function runDigestChannelJob(opts: {
channelSlug: string;
lane: DigestLaneChoice;
background?: boolean;
+ // Called with the job id as soon as it exists, so a caller can drain this
+ // specific job rather than only the thing that launched it.
+ onStarted?: (jobId: string) => void;
onDone?: () => void;
}): Promise<{ ok: boolean; error?: string }> {
const { paths, channelSlug, lane } = opts;
@@ -154,6 +164,7 @@ export async function runDigestChannelJob(opts: {
},
});
if (!result.ok) return { ok: false, error: result.error };
+ opts.onStarted?.(result.jobId);
// Wait for the channel to finish before returning: the sweep is sequential by
// design (one GPU), and the queue would serialize these anyway — awaiting
// makes that explicit and lets the loop re-plan against real results.
@@ -165,6 +176,7 @@ async function runSweepLoop(
paths: Paths,
lane: DigestLaneChoice,
channelSlugs: string[] | undefined,
+ live: SweepLive,
onLog: (msg: string) => void,
signal: AbortSignal,
drainSignal: AbortSignal,
@@ -226,7 +238,11 @@ async function runSweepLoop(
// Behind anything an operator clicks by hand: a sweep is weeks long and
// must never make a deliberate single-channel run wait for it.
background: true,
+ onStarted: (jobId) => {
+ live.channelJobId = jobId;
+ },
});
+ live.channelJobId = null;
if (!outcome.ok) {
// One channel failing to START is not the sweep failing.
onLog(`!! ${channel.channelSlug}: ${outcome.error ?? "failed to start"}`);
@@ -290,13 +306,20 @@ export async function startDigestSweep(
const settings = getSettings();
if (!settings.digest.sweepEnabled) {
- writeSettings({
+ // AWAITED, and it matters: the loop reads `sweepEnabled` at the top of its
+ // very first pass, so an un-awaited write races it and the sweep quits
+ // immediately with "disarmed in settings" — refusing to start at all.
+ await writeSettings({
...settings,
digest: { ...settings.digest, sweepEnabled: true },
});
}
- const live: SweepLive = { jobId: "", startedAt: Date.now() };
+ const live: SweepLive = {
+ jobId: "",
+ startedAt: Date.now(),
+ channelJobId: null,
+ };
const lane = opts.lane ?? "local";
const result = await runManagedFunction({
kind: DIGEST_SWEEP_KIND,
@@ -313,6 +336,7 @@ export async function startDigestSweep(
paths,
lane,
opts.channelSlugs,
+ live,
onLog,
signal,
ctx.drainSignal,
@@ -340,17 +364,23 @@ export async function resumeDigestSweepIfEnabled(
// Disarm and stop. Persisting the flag is the point: without it a restart would
// resurrect a sweep the operator had just stopped.
-export function stopDigestSweep(): boolean {
+export async function stopDigestSweep(): Promise<boolean> {
const settings = getSettings();
if (settings.digest.sweepEnabled) {
- writeSettings({
+ // Awaited so a restart cannot resurrect a sweep the operator just stopped.
+ await writeSettings({
...settings,
digest: { ...settings.digest, sweepEnabled: false },
});
}
const live = getSingleton().live;
if (!live) return false;
- // Drain, not cancel: the in-flight channel finishes rather than losing a
- // part-generated video, and the loop exits at its next check.
+ // Drain BOTH: the orchestrator, and the per-channel job it is waiting on.
+ // Draining only the orchestrator would mean "stop" waits for the current
+ // channel to finish — up to 8.7 days on the heaviest one — while the button
+ // sat there looking wedged. A drain is still graceful: the video in flight
+ // completes, the batch stops taking new ones, and nothing part-generated is
+ // thrown away.
+ if (live.channelJobId) getRegistry().requestDrain(live.channelJobId);
return getRegistry().requestDrain(live.jobId);
}
diff --git a/editor/app/jobs/actions.ts b/editor/app/jobs/actions.ts
@@ -231,7 +231,7 @@ export async function startDigestSweepAction(): Promise<DigestSweepResult> {
export async function stopDigestSweepAction(): Promise<DigestSweepResult> {
try {
- stopDigestSweep();
+ await stopDigestSweep();
revalidatePath("/");
revalidatePath("/jobs");
return { ok: true };