import path from "node:path"; import { open, readdir, stat, readFile, rm } from "node:fs/promises"; import type { Paths } from "../lib/paths"; import { mapConcurrent } from "../lib/concurrency"; import { getRegistry, type JobRecord, type JobStatus } from "./registry"; // One page of job rows, each a couple of small sidecar reads. const ENTRY_READ_CONCURRENCY = 32; import { metaPath, readJobMeta } from "./jobMeta"; import { jobSpecDetail } from "./jobDetail"; import { jobIdTime } from "./ulid"; export type JobListEntry = { id: string; kind?: string; channelSlug?: string; videoId?: string; queueKey?: string; status: JobStatus | "archived"; queuedAt: number; startedAt?: number; endedAt?: number; exitCode?: number; inRegistry: boolean; // True when the job carries a replay descriptor (JobSpec) and can therefore // be re-run from it (Retry). See common/jobs/jobSpec.ts. replayable: boolean; logPath: string; logSize: number; // A short phrase naming what this particular job is for, when its kind alone // does not say (a fetch-window job's requester and clip). See jobSpecDetail. detail?: string; // Why a `cancelled` job was cancelled, when something other than a person // pressing Cancel decided it — today only the boot pass (bootQueuedJobs.ts), // e.g. "server restarted; the scheduler re-derives syncs". From the sidecar; // absent on every other job. cancelReason?: string; }; export type JobsPage = { entries: JobListEntry[]; // More jobs exist beyond this page (i.e. raise `limit` to see them). hasMore: boolean; // Total jobs known on disk + in the registry, independent of the page size. total: number; }; const TERMINAL_STATUSES: ReadonlySet = new Set([ "done", "failed", "cancelled", ]); // Default page size for the /jobs list. The page grows this via ?limit=. export const DEFAULT_JOBS_LIMIT = 50; // Automatic on-finish retention (see maybePruneJobLogs): keep the newest N jobs // and drop anything older than the max age, whichever bites first. Generous — // the point is to bound the directory (including the constant refresh-report // churn), not to be aggressive. const RETENTION_KEEP_LAST = 500; const RETENTION_MAX_AGE_MS = 30 * 24 * 60 * 60 * 1000; // 30 days // The list-page is `force-dynamic`, so job finishes could fire a readdir on // every finalize. Throttle so a burst of jobs prunes at most once per window. const PRUNE_THROTTLE_MS = 60_000; async function readLogIds(paths: Paths): Promise> { try { const files = await readdir(paths.jobsDir); return new Set( files.filter((n) => n.endsWith(".log")).map((n) => n.replace(/\.log$/, "")), ); } catch { return new Set(); } } // Build one JobListEntry, preferring the live registry record and falling back // to the on-disk `.meta.json` sidecar (or the log mtime for old logs with no // sidecar). `hasLog` says whether an `.log` file exists on disk. async function buildEntry( paths: Paths, id: string, live_: JobRecord | undefined, hasLog: boolean, ): Promise { const logPath = hasLog ? path.join(paths.jobsDir, `${id}.log`) : live_?.logPath ?? path.join(paths.jobsDir, `${id}.log`); let logSize = 0; let mtime = 0; if (hasLog) { try { const s = await stat(logPath); logSize = s.size; mtime = s.mtimeMs; } catch { /* log vanished between listing and stat; treat as size 0 */ } } if (live_) { return { id, kind: live_.kind, channelSlug: live_.channelSlug, videoId: live_.videoId, queueKey: live_.queueKey, status: live_.status, queuedAt: live_.queuedAt, startedAt: live_.startedAt, endedAt: live_.endedAt, exitCode: live_.exitCode, inRegistry: true, replayable: Boolean(live_.spec), detail: jobSpecDetail(live_.kind, live_.spec), logPath, logSize, }; } // Not in the registry: recover what we can from the sidecar so an evicted or // post-restart job still shows its kind/channel/status/duration. With no // sidecar (old logs), fall back to the mtime-based minimal entry. A // non-terminal sidecar status means the job isn't actually live (it's not in // the registry), so report it as "archived". const meta = await readJobMeta(paths, id); if (meta) { return { id, kind: meta.kind, channelSlug: meta.channelSlug, videoId: meta.videoId, queueKey: meta.queueKey, status: TERMINAL_STATUSES.has(meta.status) ? meta.status : "archived", queuedAt: meta.queuedAt, startedAt: meta.startedAt, endedAt: meta.endedAt, exitCode: meta.exitCode, inRegistry: false, replayable: Boolean(meta.spec), detail: jobSpecDetail(meta.kind, meta.spec), // Only on a job that did end `cancelled`: a reason beside any other // status would explain something that did not happen. ...(meta.status === "cancelled" && meta.cancelReason ? { cancelReason: meta.cancelReason } : {}), logPath, logSize, }; } return { id, status: "archived", queuedAt: mtime, startedAt: mtime, endedAt: mtime, inRegistry: false, replayable: false, logPath, logSize, }; } // A page of the jobs list, newest first. Merges the in-memory registry with the // leftover `.log` files from previous server lifetimes. Ordering and paging use // jobIdTime(id) — decoded from the id itself (ULID, or the legacy `-` // scheme) — so only the returned page is stat-ed / sidecar-read, NOT the whole // directory. Registry records use their authoritative queuedAt as the sort key. export async function listAllJobs( paths: Paths, opts: { limit?: number; before?: number } = {}, ): Promise { const registry = getRegistry(); const live = new Map(); for (const r of registry.list()) live.set(r.id, r); const logIds = await readLogIds(paths); // Candidate set = every log file ∪ every registry record. Time is the sort // key: registry queuedAt when known, else decoded from the id. const candidates = new Map< string, { id: string; time: number; hasLog: boolean } >(); for (const id of logIds) { candidates.set(id, { id, time: jobIdTime(id), hasLog: true }); } for (const r of live.values()) { const existing = candidates.get(r.id); if (existing) existing.time = r.queuedAt; else candidates.set(r.id, { id: r.id, time: r.queuedAt, hasLog: false }); } let sorted = Array.from(candidates.values()).sort((a, b) => b.time - a.time); if (typeof opts.before === "number") { const before = opts.before; sorted = sorted.filter((c) => c.time < before); } const limit = opts.limit ?? DEFAULT_JOBS_LIMIT; const slice = sorted.slice(0, limit + 1); const hasMore = slice.length > limit; const page = hasMore ? slice.slice(0, limit) : slice; // One page's worth of per-job sidecar reads, concurrently rather than one at // a time. mapConcurrent preserves input order, so the sort above still holds. const entries = await mapConcurrent(page, ENTRY_READ_CONCURRENCY, (c) => buildEntry(paths, c.id, live.get(c.id), c.hasLog), ); return { entries, hasMore, total: candidates.size }; } // Resolve a single job by id for the detail page — avoids listing every job. // Null when neither the registry nor an `.log` file knows the id. export async function getJobEntry( paths: Paths, id: string, ): Promise { const live_ = getRegistry().get(id); let hasLog = false; try { await stat(path.join(paths.jobsDir, `${id}.log`)); hasLog = true; } catch { /* no log file */ } if (!live_ && !hasLog) return null; return buildEntry(paths, id, live_, hasLog); } export async function readLogChunk( logPath: string, fromBytes: number, ): Promise<{ content: string; nextOffset: number }> { let raw: string; try { raw = await readFile(logPath, "utf8"); } catch { return { content: "", nextOffset: fromBytes }; } if (fromBytes >= raw.length) { return { content: "", nextOffset: raw.length }; } return { content: raw.slice(fromBytes), nextOffset: raw.length }; } // THE LAST `lines` LINES OF A LOG, read from its END. A transcription job's log // runs to megabytes; an agent asking "what is it doing" wants the last screen // of it, not a download of the whole file (`readLogChunk` from 0 is that). At // most `maxBytes` are read: a log whose last window holds fewer lines answers // with what it has, and the first line of a mid-file window is dropped, so a // half line is never handed back as a whole one. Missing file → []. export async function readLogTail( logPath: string, lines: number, maxBytes = 256 * 1024, ): Promise { if (lines <= 0) return []; let handle: Awaited> | undefined; try { handle = await open(logPath, "r"); const { size } = await handle.stat(); const start = Math.max(0, size - maxBytes); const length = size - start; if (length <= 0) return []; const buf = Buffer.alloc(length); await handle.read(buf, 0, length, start); const all = buf.toString("utf8").split("\n"); if (start > 0) all.shift(); while (all.length && all[all.length - 1] === "") all.pop(); return all.slice(-lines); } catch { return []; } finally { await handle?.close().catch(() => {}); } } // Delete `.log` + `.meta.json` pairs by retention policy. Never touches a job // the registry currently reports running or queued. `all` clears every finished // job; otherwise `keepLast` drops the tail past the newest N and `olderThanMs` // drops anything older than the cutoff (union — either condition deletes). export async function pruneJobLogs( paths: Paths, opts: { keepLast?: number; olderThanMs?: number; all?: boolean; now?: number }, ): Promise<{ deleted: number }> { const protectedIds = new Set( getRegistry() .list() .filter((r) => r.status === "running" || r.status === "queued") .map((r) => r.id), ); const logIds = await readLogIds(paths); const items = Array.from(logIds) .map((id) => ({ id, t: jobIdTime(id) })) .sort((a, b) => b.t - a.t); const now = opts.now ?? Date.now(); const cutoff = opts.olderThanMs != null ? now - opts.olderThanMs : null; let deleted = 0; for (let i = 0; i < items.length; i++) { const { id, t } = items[i]; if (protectedIds.has(id)) continue; let doDelete = Boolean(opts.all); if (!doDelete && opts.keepLast != null && i >= opts.keepLast) doDelete = true; // t === 0 means an unparseable id (age unknown); leave age-based pruning to // keepLast so we never delete something whose age we can't establish. if (!doDelete && cutoff != null && t > 0 && t < cutoff) doDelete = true; if (!doDelete) continue; await rm(path.join(paths.jobsDir, `${id}.log`), { force: true }); await rm(metaPath(paths, id), { force: true }); deleted++; } return { deleted }; } // Throttled automatic retention, called after a job finalizes. Best-effort: a // prune failure must never affect the job. Bounds the `.jobs` directory so the // list stays fast without any manual clearing. let lastPruneAt = 0; export async function maybePruneJobLogs(paths: Paths): Promise { const now = Date.now(); if (now - lastPruneAt < PRUNE_THROTTLE_MS) return; lastPruneAt = now; try { await pruneJobLogs(paths, { keepLast: RETENTION_KEEP_LAST, olderThanMs: RETENTION_MAX_AGE_MS, now, }); } catch { /* best-effort */ } }