"use server"; import { revalidatePath } from "next/cache"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { saveSettings } from "../settings/saveSettings"; import { laneRootFromScope } from "yt-dlp-transcript-common/lib/laneMigration"; import { isDefaultChannelPriority } from "yt-dlp-transcript-common/lib/channelPriority"; import { operationsForLane } from "yt-dlp-transcript-common/lib/operations"; import type { AutoQueueKind } from "yt-dlp-transcript-common/lib/autoQueueTypes"; import { drainAutoRunner, startAutoRunner, startAutoRunnerBlockedReason, } from "yt-dlp-transcript-common/controller/autoRunner"; import { pruneJobLogs } from "yt-dlp-transcript-common/jobs/listJobs"; import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; import { readJobMeta } from "yt-dlp-transcript-common/jobs/jobMeta"; import type { JobSpec } from "yt-dlp-transcript-common/jobs/jobSpec"; import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; import { runJobSpec } from "./runJobSpec"; import { stuckJobIds } from "./active/buildActiveJobs"; export async function cancelJobAction(id: string): Promise<{ ok: boolean }> { const ok = getRegistry().cancel(id); revalidatePath("/jobs"); return { ok }; } // Escape hatch for a WEDGED slot (see registry.forceRelease): unconditionally // free the scheduler slot for this id — even if its record is already terminal // or evicted — SIGKILLing any still-running child. Unblocks a queue whose head // is stuck (the live payload re-reads the slot on the next poll). export async function forceReleaseJobAction( id: string, ): Promise<{ ok: boolean }> { const ok = getRegistry().forceRelease(id); revalidatePath("/jobs"); return { ok }; } // Reap every stuck slot the live payload reports: force-release each. The hard // cases were already healed by the build that listed them; forceRelease is a // no-op on a freed slot (registry.forceRelease). Returns how many were released. export async function reapStuckJobsAction(): Promise<{ count: number }> { const registry = getRegistry(); const stuckIds = await stuckJobIds(); let count = 0; for (const id of stuckIds) { if (registry.forceRelease(id)) count++; } revalidatePath("/jobs"); return { count }; } // Soft-cancel: let the batch's in-flight sub-operations finish, start no new // ones, then complete and release the queue for the next job. export async function drainJobAction(id: string): Promise<{ ok: boolean }> { const ok = getRegistry().requestDrain(id); revalidatePath("/jobs"); return { ok }; } // Spin-down: drain every running job and cancel every queued one in a single // pass. requestDrain() drains running jobs and, for queued jobs, falls through // to cancel() (see registry.ts), so one loop covers both. Returns how many jobs // were acted on. export async function drainAllAction(): Promise<{ count: number }> { const registry = getRegistry(); const targets = registry .list() .filter((j) => j.status === "running" || j.status === "queued"); let count = 0; for (const j of targets) { if (registry.requestDrain(j.id)) count++; } revalidatePath("/jobs"); return { count }; } // Granular control: nudge a QUEUED job up (-1) or down (+1) within its queue. // Only queued jobs move; the running head is never displaced. export async function reorderJobAction( id: string, dir: -1 | 1, ): Promise<{ ok: boolean }> { const ok = getRegistry().reorder(id, dir); revalidatePath("/jobs"); return { ok }; } // Granular control: jump a QUEUED job to the front of its queue (it runs next // when the current job releases the queue). export async function promoteJobAction(id: string): Promise<{ ok: boolean }> { const ok = getRegistry().promote(id); revalidatePath("/jobs"); return { ok }; } // Resolve a job's replay descriptor: prefer the live registry record, fall back // to the on-disk meta sidecar so a failed job the 100-job cap evicted (or one // from before a restart) can still be retried. Null if the kind isn't // replayable (no spec was ever attached). async function resolveSpec(id: string): Promise { const live = getRegistry().get(id)?.spec; if (live) return live; const meta = await readJobMeta(getPaths(), id); return meta?.spec ?? null; } // One-click retry: re-run a (typically failed) job from its stored spec. As on // any replay, bucket jobs re-derive from the channel's CURRENT state. The // re-run jumps ahead of other queued work (run next) by reusing the same // promote the reorder buttons use — a no-op if it starts immediately. export async function retryJobAction(id: string): Promise { const spec = await resolveSpec(id); if (!spec) { return { ok: false, error: "This job can't be retried — it has no replay descriptor.", }; } const res = await runJobSpec(spec); if (res.ok) getRegistry().promote(res.jobId); revalidatePath("/jobs"); return res; } // Retry every currently-listed failed job that has a spec. Returns how many were // re-launched (kinds without a replay descriptor are skipped) and the new jobs' // ids — the button reads the count; `pnpm ops job retry-failed --wait` follows // the ids. The new job's stream is cancelled here: nothing consumes it (the // button never did — a dropped stream leaves the on-disk log running). export async function retryAllFailedAction(): Promise<{ count: number; jobIds: string[]; }> { const registry = getRegistry(); const failed = registry .list() .filter((j): j is typeof j & { spec: JobSpec } => Boolean(j.status === "failed" && j.spec), ); const jobIds: string[] = []; for (const j of failed) { const res = await runJobSpec(j.spec); if (res.ok) { registry.promote(res.jobId); void res.stream.cancel(); jobIds.push(res.jobId); } } revalidatePath("/jobs"); return { count: jobIds.length, jobIds }; } // ARM A LANE: write its scope as a tree, switch it on, and start its runner. // // This replaces `startDigestSweepAction(channels)` and // `startBackfillSweepAction(kinds, channels)`. What those did was persist a // flag plus a scope in one write and then launch a corpus walk; what this does // is persist `enabled: true` plus a `root` in one write and then start the lane // runner. The one-write property is why it is still a single action: a scope // written separately from the switch means a restart between the two writes // resurrects a deliberately-bounded run as a corpus-wide one, which is the bug // the sweeps' own start/stop pair documented at length. // // THE SCOPE IS THE TREE, and it is built by the SAME function the settings // migration uses (`laneRootFromScope`), so arming a lane here and migrating a // sweep that was armed with the same scope produce byte-identical roots. There // is one leaf builder; a second one is how the two would drift. // // An ABSENT scope keeps the stored tree — that is the difference between "arm // this lane on these channels" and "switch this lane back on". An empty scope // object is NOT absent: `{}` means the whole corpus, exactly as an empty // `sweepChannels` did. export type ArmLaneResult = { ok: boolean; jobId?: string; error?: string }; export async function armLaneAction( lane: AutoQueueKind, scope?: { channels?: string[]; operations?: string[] }, ): Promise { try { const settings = getSettings(); // ONE WRITER OF `root` WHILE A MODEL EXISTS. A scope IS a tree — the whole // point of this action is that it writes one — so while the channel // priority document says anything, arming WITH a scope would write a tree // the runner does not dispatch from (`laneDispatchRoot` compiles) and that // the next priority save overwrites. Refused with the same message the // policy editor gives, rather than accepted and silently ignored. // // Arming with NO scope is untouched: it keeps the stored tree, which is // "switch this lane back on" and says nothing about priority. if (scope && !isDefaultChannelPriority(settings.channelPriority)) { return { ok: false, error: "This lane's rules are generated from the channel priorities. " + "Set the channels' tiers on /channels, then arm the lane without a " + "scope — a tree written here is not what the runner dispatches from.", }; } // VALIDATED HERE, because the sanitizer does not. // // An unknown operation id survives a settings write and then matches // nothing: a leaf naming it draws no list at all, so a lane armed on it // reports itself running and does exactly no work. That is the worst // failure this console can have — the operator clicks, and watches a lane // hold at zero forever. So: drop the unknowns, arm what is left, and SAY // which were dropped. const known = new Set(operationsForLane(lane, settings).map((o) => o.id)); const wanted = scope?.operations ?? []; const operations = wanted.filter((id) => known.has(id)); const dropped = wanted.filter((id) => !known.has(id)); if (wanted.length > 0 && operations.length === 0) { return { ok: false, error: `No enabled operation is named by this scope (${dropped.join(", ")}). ` + `A lane armed on it would run forever without doing anything.`, }; } const policy = settings.autoQueue[lane]; await saveSettings({ autoQueue: { ...settings.autoQueue, [lane]: { ...policy, enabled: true, // TICKING EVERY OPERATION IS NOT THE SAME AS NAMING TODAY'S THREE, // and `available` is what says so — the collapse lives in the leaf // builder, with the settings migration, rather than being written // once here and once there. root: scope ? laneRootFromScope(lane, { channels: scope.channels, operations, available: [...known], }) : policy.root, }, }, }); // ARM = ENABLE + START, because the two were one act on the button this // replaces. `startAutoRunner` is a no-op on a lane already running. const jobId = await startAutoRunner(lane); revalidatePath("/jobs"); revalidatePath("/operations/[id]", "page"); if (jobId === null) { return { ok: false, error: startAutoRunnerBlockedReason(lane) ?? "The lane could not be started (see job logs).", }; } return dropped.length > 0 ? { ok: true, jobId, error: `Ignored ${dropped.length} operation(s) this build does not have enabled: ${dropped.join(", ")}.`, } : { ok: true, jobId }; } catch (e) { return { ok: false, error: (e as Error).message }; } } // DISARM: switch the lane off and let the runner finish what it is holding. // // A DRAIN, not a stop, and that is the sweeps' behaviour kept: the unit in // flight completes instead of being thrown away. The `enabled: false` is what // stops a restart bringing it back — the boot hook resumes every enabled lane. // // The TREE IS LEFT ALONE. A disarm that cleared the scope would be the bug the // old `stopDigestSweep` had to grow a second clause for; here the scope is a // tree an operator authored, and it must survive being switched off. export async function disarmLaneAction( lane: AutoQueueKind, ): Promise { try { const settings = getSettings(); await saveSettings({ autoQueue: { ...settings.autoQueue, [lane]: { ...settings.autoQueue[lane], enabled: false }, }, }); drainAutoRunner(lane); revalidatePath("/jobs"); revalidatePath("/operations/[id]", "page"); return { ok: true }; } catch (e) { return { ok: false, error: (e as Error).message }; } } // Retention scopes offered by the ClearLogsMenu. "all" clears every finished // job's log; the day-scopes clear anything older than that. Running/queued jobs // are never deleted (see pruneJobLogs). const DAY_MS = 24 * 60 * 60 * 1000; export type ClearLogsScope = "7d" | "30d" | "90d" | "all"; export async function clearFinishedLogsAction( scope: ClearLogsScope, ): Promise<{ deleted: number }> { const opts = scope === "all" ? { all: true } : { olderThanMs: { "7d": 7, "30d": 30, "90d": 90 }[scope] * DAY_MS }; const { deleted } = await pruneJobLogs(getPaths(), opts); revalidatePath("/jobs"); return { deleted }; }