import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { getJobEntry, listAllJobs, readLogTail, type JobListEntry, } from "yt-dlp-transcript-common/jobs/listJobs"; import { getRegistry, type JobRecord } from "yt-dlp-transcript-common/jobs/registry"; // THE JOBS, AS AN AGENT READS THEM — `GET /api/ops/job/` and // `GET /api/ops/jobs`. // // Read-only shaping over the two places a job is known, the same two the /jobs // page reads: the live registry (`getRegistry()`) and the `.jobs` directory // (`getJobEntry` / `listAllJobs`, which fall back to the `.meta.json` // sidecar). No rule lives here, and nothing is constructed: no channel stats, // no disk gate, no runner status — which is the difference from the // `/api/jobs/active` view, a UI poll that builds the whole live payload (and // walks channel stats for progress) on every request. export type OpsQueueView = { // The scheduler queue the job waits on ("" = runs in parallel, never queued). key: string; // 0 = running at the head, 1.. = queued that many places back. position: number; // How many jobs wait on that queue in all. queued: number; // The job holding the queue's head, when it is not this one. head?: { id: string; kind: string; channelSlug?: string; startedAt?: number }; }; export type OpsJobView = Omit & { // Only while the job is queued or running. queue?: OpsQueueView; // Live jobs only: what the runner last reported. progress?: JobRecord["progress"]; tasks?: number; draining?: boolean; background?: boolean; }; export function queueOf(id: string): OpsQueueView | undefined { const registry = getRegistry(); const record = registry.get(id); if (!record || (record.status !== "queued" && record.status !== "running")) { return undefined; } const position = registry.positionInQueue(id); if (position < 0) return undefined; const snap = registry.listQueues().find((q) => q.name === record.queueKey); const head = snap?.running; return { key: record.queueKey, position, queued: snap?.queued.length ?? 0, ...(head && head.id !== id ? { head: { id: head.id, kind: head.kind, ...(head.channelSlug ? { channelSlug: head.channelSlug } : {}), ...(head.startedAt ? { startedAt: head.startedAt } : {}), }, } : {}), }; } export function viewOf(entry: JobListEntry): OpsJobView { const { logPath: _logPath, ...rest } = entry; void _logPath; const record = getRegistry().get(entry.id); const queue = queueOf(entry.id); return { ...rest, ...(queue ? { queue } : {}), ...(record?.progress ? { progress: record.progress } : {}), ...(record?.tasks?.length ? { tasks: record.tasks.length } : {}), ...(record?.draining ? { draining: true } : {}), ...(record?.background ? { background: true } : {}), }; } // One job, or null when neither the registry nor `.jobs/` knows the id. export async function readJob( id: string, tail: number, ): Promise<{ job: OpsJobView; tail?: string[] } | null> { const entry = await getJobEntry(getPaths(), id); if (!entry) return null; return { job: viewOf(entry), ...(tail > 0 ? { tail: await readLogTail(entry.logPath, tail) } : {}), }; } export type JobsFilter = { active?: boolean; failed?: boolean; kind?: string; slug?: string; limit: number; }; // How far back a FILTERED list looks. A page of the newest `limit` jobs // filtered afterwards would answer "no failed syncs" whenever the last fifty // jobs were refresh-reports; reading this many sidecars (32 at a time) is what // the /jobs page's own "Load more" costs at its ceiling. export const JOBS_SCAN = 2000; export async function listJobs( filter: JobsFilter, ): Promise<{ jobs: OpsJobView[]; scanned: number; total: number }> { const match = (e: { status: string; kind?: string; channelSlug?: string }) => (!filter.failed || e.status === "failed") && (!filter.kind || e.kind === filter.kind) && (!filter.slug || e.channelSlug === filter.slug); if (filter.active) { // The live head, straight from the registry: queue order, running first. const registry = getRegistry(); const live = registry .list() .filter((r) => r.status === "running" || r.status === "queued") .filter(match); const paths = getPaths(); const entries = await Promise.all(live.map((r) => getJobEntry(paths, r.id))); const views = entries .filter((e): e is JobListEntry => e !== null) .map(viewOf); views.sort( (a, b) => (a.queue?.key ?? "").localeCompare(b.queue?.key ?? "") || (a.queue?.position ?? 0) - (b.queue?.position ?? 0), ); return { jobs: views.slice(0, filter.limit), scanned: live.length, total: views.length, }; } const filtered = Boolean(filter.failed || filter.kind || filter.slug); const page = await listAllJobs(getPaths(), { limit: filtered ? JOBS_SCAN : filter.limit, }); const hits = page.entries.filter(match); return { jobs: hits.slice(0, filter.limit).map(viewOf), scanned: page.entries.length, total: page.total, }; }