Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit 78f836fda3c7175864e7c1ab7376984f010ac310
parent 43620a86ed8c1fa649143891e0445161906aaced
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Thu, 25 Jun 2026 00:15:00 -0400

Phase 1: per-channel keep-latest retention + deletion pinning

Adds the backend for a per-channel "keep latest N videos" retention rule, the
first phase of a larger video-persistence subsystem. UI lands in a later phase;
for now keepLatest is set via a channel's config.json.

Retention (rolling window):
- ChannelConfig.keepLatest (0/undefined = off; clamped to KEEP_LATEST_MAX).
- common/controller/keptVideos.ts: computeKeptVideoIds() returns the newest N
  video ids by upload_date (metadata.info.json, with a YYYYMMDD_ dir-name
  fallback). Unit-tested (keptVideos.test.ts).
- cleanAudioFromTranscribed skips dirs in the kept window (like do-not-clean).
- channelSnapshot folds the kept set into the cleanup-exclusion set so reclaim
  estimates and the transcribed/multiple-format/foreign buckets exclude kept
  videos; records snapshot.keptCount.

Deletion pinning:
- common/controller/checkKeptDeleted.ts re-probes only the kept window via
  runAvailabilityCheck (onlyIds + recheck-non-deleted) and permanently pins any
  video that is gone from source (deleted/private/members_only) with a
  do-not-clean marker, so it survives rolling out of the window.
- checkKeptDeletedAction managed job (kind "check-kept-deleted"), wired into the
  job re-run dispatcher and job-kind labels.

Scheduling:
- syncScheduler.keepLatestCheckIntervalMinutes (default daily) + per-channel
  lastKeptCheckAt state.
- runSchedulerTick enqueues the kept-deletion check for due keepLatest>0
  channels on the per-channel local queue: suppressed during quiet hours, capped
  per tick, and skipped for a channel just synced this tick.

Typecheck clean (common + editor); keptVideos unit tests pass.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

Diffstat:
Mcommon/controller/channelSnapshot.ts | 25++++++++++++++++++++++---
Acommon/controller/checkKeptDeleted.ts | 97+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/cleanAudioFromTranscribed.ts | 16++++++++++++++++
Acommon/controller/keptVideos.test.ts | 113+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/controller/keptVideos.ts | 66++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/syncSchedulerState.ts | 6++++++
Mcommon/lib/channelConfig.ts | 21+++++++++++++++++++++
Mcommon/lib/settings.ts | 14++++++++++++++
Meditor/CHANGELOG.md | 1+
Meditor/app/channels/[slug]/whisperActions.ts | 31+++++++++++++++++++++++++++++++
Meditor/app/jobs/jobKindLabels.ts | 2++
Meditor/app/jobs/runJobSpec.ts | 3+++
Meditor/app/scheduler/runTick.ts | 34++++++++++++++++++++++++++++++++++
Meditor/app/settings/actions.ts | 3+++
14 files changed, 429 insertions(+), 3 deletions(-)

diff --git a/common/controller/channelSnapshot.ts b/common/controller/channelSnapshot.ts @@ -26,6 +26,7 @@ import { extractVideoId } from "../ytdlp/runYtdlp"; import { reconcileVideoDirs } from "./reconcileVideoDirs"; import { loadFailedTranscriptions } from "./failedTranscriptions"; import { readChannelConfig } from "./channels"; +import { computeKeptVideoIds } from "./keptVideos"; import { loadMaybeMissing } from "./quickAvailabilityCheck"; export type AvailabilitySnapshot = { @@ -82,6 +83,11 @@ export type ChannelSnapshot = { }; undownloadedIds: string[]; excludedFromDownload?: ExcludedFromDownload; + // Count of videos in the channel's keep-latest window (ChannelConfig.keepLatest). + // These are protected from the Clean-audio sweep (folded into the cleanup- + // exclusion set alongside do-not-clean markers) and targeted for source-video + // persistence. Optional: older snapshots lack it; readers must default to 0. + keptCount?: number; availability?: AvailabilitySnapshot; // Known videos absent from the channel's most recent fresh flat-playlist // fetch (the quick availability check). The maybe-missing.json sidecar is the @@ -214,6 +220,13 @@ export async function generateChannelSnapshot( const targetAudioFile = config?.audioFormat ? `audio.${config.audioFormat}` : null; + // The keep-latest window: newest N videos protected from cleanup. Computed + // once and folded into the cleanup-exclusion set below. + const keptIds = await computeKeptVideoIds({ + paths, + channelSlug: slug, + keepLatest: config?.keepLatest ?? 0, + }); const videoDirNames = dirEntries .filter((d) => d.isDirectory()) .map((d) => d.name); @@ -307,6 +320,11 @@ export async function generateChannelSnapshot( for (const v of perVideo) { if (v.doNotClean) doNotCleanIds.add(v.id); } + // Everything shielded from the Clean-audio sweep: explicit do-not-clean markers + // plus the rolling keep-latest window. The cleanup buckets/reclaim estimates + // below exclude this whole set so they match what cleanAudioFromTranscribed + // will actually remove. + const protectedFromCleanup = new Set<string>([...doNotCleanIds, ...keptIds]); const noTranscript: string[] = []; const downloadedNoTranscript: string[] = []; @@ -355,7 +373,7 @@ export async function generateChannelSnapshot( // Reclaim estimate for the wrong-format sweep: every non-target audio file // in a cleanable (not do-not-clean) dir. Covers orphans (no target) and the // extras counted in multipleAudioFormats — the sweep removes them all. - if (targetAudioFile && !doNotCleanIds.has(id)) { + if (targetAudioFile && !protectedFromCleanup.has(id)) { for (const name of files.audioFiles) { if (name !== targetAudioFile) { foreignAudioBytes += audioSizes[name] ?? 0; @@ -366,7 +384,7 @@ export async function generateChannelSnapshot( targetAudioFile && files.audioFiles.includes(targetAudioFile) && files.audioFiles.length > 1 && - !doNotCleanIds.has(id) + !protectedFromCleanup.has(id) ) { multipleAudioFormats.push(id); for (const name of files.audioFiles) { @@ -378,7 +396,7 @@ export async function generateChannelSnapshot( if ( files.hasWhisper && files.audioFiles.length > 0 && - !doNotCleanIds.has(id) + !protectedFromCleanup.has(id) ) { transcribedWithAudio.push(id); for (const name of files.audioFiles) { @@ -490,6 +508,7 @@ export async function generateChannelSnapshot( }, undownloadedIds, excludedFromDownload, + keptCount: keptIds.size, availability, ...(maybeMissing ? { maybeMissing } : {}), cleanupBytes: { diff --git a/common/controller/checkKeptDeleted.ts b/common/controller/checkKeptDeleted.ts @@ -0,0 +1,97 @@ +import path from "node:path"; +import type { Paths } from "../lib/paths"; +import { EXCLUDED_FROM_DOWNLOAD, type Availability } from "../lib/availability"; +import { resolveEffectiveAvailability } from "../lib/availability-server"; +import { loadDoNotClean, setDoNotClean } from "../lib/doNotClean-server"; +import { readChannelConfig } from "./channels"; +import { computeKeptVideoIds } from "./keptVideos"; +import { runAvailabilityCheck } from "./checkAvailability"; + +// Periodic "are my kept videos still there?" pass. For a channel's keep-latest +// window, re-probe source availability and pin (do-not-clean) any video found +// permanently gone (deleted / private / members_only). A pinned video survives +// even after it rolls out of the rolling window — it's now irreplaceable, so we +// keep its media forever. Reuses runAvailabilityCheck (onlyIds = the kept set) +// for the actual yt-dlp probing. + +export type CheckKeptDeletedOptions = { + channelSlug: string; + paths: Paths; + onLog?: (msg: string) => void; + signal?: AbortSignal; + concurrency?: number; +}; + +export type CheckKeptDeletedResult = { + // Size of the keep-latest window inspected. + kept: number; + // Videos that probed as permanently gone (EXCLUDED_FROM_DOWNLOAD). + deleted: number; + // Videos newly pinned with the do-not-clean marker this run. + pinned: number; +}; + +function isGone(a: Availability | null): boolean { + return ( + a !== null && + (EXCLUDED_FROM_DOWNLOAD as ReadonlyArray<Availability>).includes(a) + ); +} + +export async function checkKeptDeleted({ + channelSlug, + paths, + onLog, + signal, + concurrency, +}: CheckKeptDeletedOptions): Promise<CheckKeptDeletedResult> { + const log = onLog ?? ((m: string) => console.log(m)); + const config = await readChannelConfig(paths, channelSlug); + const keepLatest = config?.keepLatest ?? 0; + if (keepLatest <= 0) { + log(`${channelSlug}: keep-latest disabled — nothing to check.`); + return { kept: 0, deleted: 0, pinned: 0 }; + } + + const keptIds = await computeKeptVideoIds({ paths, channelSlug, keepLatest }); + if (keptIds.size === 0) { + log(`${channelSlug}: no videos in the keep-latest window.`); + return { kept: 0, deleted: 0, pinned: 0 }; + } + const onlyIds = [...keptIds]; + log(`${channelSlug}: checking ${onlyIds.length} kept video(s) for deletion…`); + + // Probe only the kept set, re-checking everything not already known-gone. + await runAvailabilityCheck({ + channelSlug, + paths, + mode: "recheck-non-deleted", + onlyIds, + concurrency, + onLog: log, + signal, + }); + + const dataDir = path.join(paths.channelsDir, channelSlug, "data"); + let deleted = 0; + let pinned = 0; + for (const id of onlyIds) { + if (signal?.aborted) break; + const videoDir = path.join(dataDir, id); + const availability = await resolveEffectiveAvailability(videoDir); + if (!isGone(availability)) continue; + deleted++; + // Only count/log a NEW pin — leave an existing marker (and its note) intact. + const already = await loadDoNotClean(videoDir); + if (already) continue; + const note = `deleted from source (keep-latest): ${availability}`; + await setDoNotClean(videoDir, true, note); + pinned++; + log(`Pinned ${id} as do-not-clean (${availability}).`); + } + + log( + `${channelSlug}: ${deleted} of ${onlyIds.length} kept video(s) gone from source; pinned ${pinned}.`, + ); + return { kept: onlyIds.length, deleted, pinned }; +} diff --git a/common/controller/cleanAudioFromTranscribed.ts b/common/controller/cleanAudioFromTranscribed.ts @@ -2,6 +2,8 @@ import path from "node:path"; import fs from "fs-extra"; import type { Paths } from "../lib/paths"; import { isDoNotClean } from "../lib/doNotClean-server"; +import { readChannelConfig } from "./channels"; +import { computeKeptVideoIds } from "./keptVideos"; const { pathExists, readdir, remove } = fs; @@ -33,6 +35,15 @@ export async function cleanAudioFromTranscribed({ } const dirs = await readdir(dataDir); + // The rolling keep-latest window is protected from cleanup just like the + // explicit do-not-clean marker (the source media of recent videos is kept). + const config = await readChannelConfig(paths, channelSlug); + const keptIds = await computeKeptVideoIds({ + paths, + channelSlug, + keepLatest: config?.keepLatest ?? 0, + }); + let cleanedDirs = 0; let removedFiles = 0; let skipped = 0; @@ -50,6 +61,11 @@ export async function cleanAudioFromTranscribed({ !e.endsWith(".part"), ); if (audioFiles.length === 0) continue; + if (keptIds.has(id)) { + log(`Skipped ${id} (in keep-latest window)`); + skipped++; + continue; + } if (await isDoNotClean(videoDir)) { log(`Skipped ${id} (marked do not clean)`); skipped++; diff --git a/common/controller/keptVideos.test.ts b/common/controller/keptVideos.test.ts @@ -0,0 +1,113 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import type { Paths } from "../lib/paths"; +import { computeKeptVideoIds } from "./keptVideos"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test common/controller/keptVideos.test.ts + +// Seed channels/<slug>/data/<id>/metadata.info.json for each [id, uploadDate]. +// A null uploadDate omits the metadata file so the dir-name fallback is exercised. +async function seedChannel( + channelsDir: string, + slug: string, + videos: ReadonlyArray<[string, string | null]>, +): Promise<void> { + for (const [id, uploadDate] of videos) { + const dir = path.join(channelsDir, slug, "data", id); + await mkdir(dir, { recursive: true }); + if (uploadDate !== null) { + await writeFile( + path.join(dir, "metadata.info.json"), + JSON.stringify({ id, upload_date: uploadDate }), + ); + } + } +} + +async function withPaths(fn: (paths: Paths) => Promise<void>): Promise<void> { + const dir = await mkdtemp(path.join(tmpdir(), "ttb-kept-")); + const paths = { channelsDir: path.join(dir, "channels") } as Paths; + try { + await fn(paths); + } finally { + await rm(dir, { recursive: true, force: true }); + } +} + +test("keeps the newest N by upload_date", async () => { + await withPaths(async (paths) => { + await seedChannel(paths.channelsDir, "ch", [ + ["a", "20240101"], + ["b", "20240301"], + ["c", "20240201"], + ["d", "20231231"], + ]); + const kept = await computeKeptVideoIds({ + paths, + channelSlug: "ch", + keepLatest: 2, + }); + assert.deepEqual([...kept].sort(), ["b", "c"]); + }); +}); + +test("keepLatest <= 0 keeps nothing", async () => { + await withPaths(async (paths) => { + await seedChannel(paths.channelsDir, "ch", [["a", "20240101"]]); + assert.equal( + (await computeKeptVideoIds({ paths, channelSlug: "ch", keepLatest: 0 })) + .size, + 0, + ); + }); +}); + +test("keepLatest larger than the catalog keeps all", async () => { + await withPaths(async (paths) => { + await seedChannel(paths.channelsDir, "ch", [ + ["a", "20240101"], + ["b", "20240201"], + ]); + const kept = await computeKeptVideoIds({ + paths, + channelSlug: "ch", + keepLatest: 10, + }); + assert.deepEqual([...kept].sort(), ["a", "b"]); + }); +}); + +test("falls back to a YYYYMMDD_ dir-name prefix when metadata is absent", async () => { + await withPaths(async (paths) => { + await seedChannel(paths.channelsDir, "ch", [ + ["20240301_newest", null], + ["20240101_older", null], + ["nodate", null], + ]); + const kept = await computeKeptVideoIds({ + paths, + channelSlug: "ch", + keepLatest: 1, + }); + assert.deepEqual([...kept], ["20240301_newest"]); + }); +}); + +test("missing channel data dir yields an empty set", async () => { + await withPaths(async (paths) => { + assert.equal( + ( + await computeKeptVideoIds({ + paths, + channelSlug: "ghost", + keepLatest: 5, + }) + ).size, + 0, + ); + }); +}); diff --git a/common/controller/keptVideos.ts b/common/controller/keptVideos.ts @@ -0,0 +1,66 @@ +import path from "node:path"; +import { readdir } from "node:fs/promises"; +import type { Paths } from "../lib/paths"; +import { loadRawMetadataFromDir } from "../lib/transcripts-server"; + +// Rolling "keep-latest" window computation. Given a channel's keepLatest config, +// returns the ids of the newest N videos (by upload date). Used by both the +// cleanup protection (cleanAudioFromTranscribed / channelSnapshot) and the +// per-download persistence rule (downloadOneManaged). Kept pure-ish (just fs +// reads) so it can be reused everywhere the kept set is needed. + +export type ComputeKeptOptions = { + paths: Paths; + channelSlug: string; + keepLatest: number; +}; + +// List a channel's data-dir video ids (directories, skipping dotfiles). Mirrors +// editor's readDataDirVideoIds, duplicated here so common doesn't depend on the +// editor's server-only module. +export async function listChannelVideoIds( + paths: Paths, + channelSlug: string, +): Promise<string[]> { + const dataDir = path.join(paths.channelsDir, channelSlug, "data"); + try { + const entries = await readdir(dataDir, { withFileTypes: true }); + return entries + .filter((e) => e.isDirectory() && !e.name.startsWith(".")) + .map((e) => e.name); + } catch { + return []; + } +} + +// Recency sort key for a video: upload_date (YYYYMMDD) from metadata.info.json, +// falling back to a YYYYMMDD_ prefix on the dir name, else "" (sorts oldest). +async function uploadKey(videoDir: string, id: string): Promise<string> { + const meta = await loadRawMetadataFromDir(videoDir); + if (meta?.upload_date) return meta.upload_date; + return id.match(/^(\d{8})(?:_|$)/)?.[1] ?? ""; +} + +// The set of the newest `keepLatest` video ids for a channel (by upload date, +// newest first; ties broken by id descending for determinism). Empty when +// keepLatest <= 0 or the channel has no data dir. +export async function computeKeptVideoIds({ + paths, + channelSlug, + keepLatest, +}: ComputeKeptOptions): Promise<Set<string>> { + if (!Number.isFinite(keepLatest) || keepLatest <= 0) return new Set(); + const ids = await listChannelVideoIds(paths, channelSlug); + if (ids.length === 0) return new Set(); + const dataDir = path.join(paths.channelsDir, channelSlug, "data"); + const keyed = await Promise.all( + ids.map(async (id) => ({ + id, + key: await uploadKey(path.join(dataDir, id), id), + })), + ); + keyed.sort((a, b) => + a.key === b.key ? b.id.localeCompare(a.id) : b.key.localeCompare(a.key), + ); + return new Set(keyed.slice(0, Math.floor(keepLatest)).map((k) => k.id)); +} diff --git a/common/jobs/syncSchedulerState.ts b/common/jobs/syncSchedulerState.ts @@ -24,6 +24,10 @@ export type ChannelSyncState = { // Outcome of the last reconciled scheduled sync, for the observability panel. lastOutcome: "ok" | "failed" | null; lastOutcomeAt: number | null; + // Epoch ms the scheduler last queued a keep-latest deletion check for this + // channel. Drives the keepLatestCheckIntervalMinutes cadence independently of + // the sync cadence. null = never run. + lastKeptCheckAt: number | null; }; export type SchedulerSkip = { slug: string; reason: string }; @@ -50,6 +54,7 @@ export function emptyChannelSyncState(): ChannelSyncState { lastQueuedAt: null, lastOutcome: null, lastOutcomeAt: null, + lastKeptCheckAt: null, }; } @@ -139,6 +144,7 @@ function coerceChannelState(value: unknown): ChannelSyncState { ? r.lastOutcome : null, lastOutcomeAt: num(r.lastOutcomeAt), + lastKeptCheckAt: num(r.lastKeptCheckAt), }; } diff --git a/common/lib/channelConfig.ts b/common/lib/channelConfig.ts @@ -24,6 +24,15 @@ export type ChannelConfig = { url?: string; audioFormat?: AudioFormat; keepSourceVideo?: boolean; + // Keep-latest retention/persistence window. The newest N videos (by upload + // date) are protected from the Clean-audio sweep AND have their source video + // persisted to the saved-video store (see common/controller/keptVideos.ts and + // the per-download persistence rule in downloadOneManaged.ts). Semantics: + // undefined / 0 -> disabled + // > 0 -> keep newest N (clamped to KEEP_LATEST_MAX) + // A kept video later found deleted-from-source is pinned permanently via the + // do-not-clean marker so it survives even after it rolls out of the window. + keepLatest?: number; ytdlpExtraArgs?: string[]; subLangs?: string; lastSyncedAt?: string; @@ -61,6 +70,9 @@ export const CHANNEL_SLEEP_BETWEEN_DOWNLOADS_MAX_SECONDS = 600; export const SYNC_INTERVAL_MIN_MINUTES = 1; export const SYNC_INTERVAL_MAX_MINUTES = 44640; +// Upper bound for a channel's keep-latest window. 0/undefined disables it. +export const KEEP_LATEST_MAX = 100000; + export const AUDIO_CHECK_INTERVAL_DEFAULT_SECONDS = 60; export const AUDIO_CHECK_INTERVAL_MIN_SECONDS = 10; export const AUDIO_CHECK_INTERVAL_MAX_SECONDS = 600; @@ -116,6 +128,15 @@ export function parseChannelConfig(raw: unknown): ChannelConfig | null { config.keepSourceVideo = r.keepSourceVideo; } if ( + typeof r.keepLatest === "number" && + Number.isFinite(r.keepLatest) && + r.keepLatest >= 0 + ) { + // 0 is the "disabled" sentinel preserved as-is; positives clamp to the cap. + config.keepLatest = + r.keepLatest === 0 ? 0 : clampInt(r.keepLatest, 1, KEEP_LATEST_MAX); + } + if ( Array.isArray(r.ytdlpExtraArgs) && r.ytdlpExtraArgs.every((x) => typeof x === "string") ) { diff --git a/common/lib/settings.ts b/common/lib/settings.ts @@ -142,6 +142,13 @@ export type SyncSchedulerSettings = { // [SYNC_HEARTBEAT_MIN_SECONDS, SYNC_HEARTBEAT_MAX_SECONDS]. The env var // SYNC_HEARTBEAT_SECONDS overrides this at runtime. See SCHEDULED_SYNC.md. heartbeatSeconds: number; + // Cadence (minutes) for the scheduled keep-latest deletion check. For each + // channel with ChannelConfig.keepLatest > 0, the tick re-probes the kept + // window for source deletion (checkKeptDeletedAction) at most this often and + // pins any gone videos as do-not-clean. Clamped into the sync-interval window; + // default daily. The check shares the same concurrency cap and quiet-hours + // window as scheduled syncs. See editor/app/scheduler/runTick.ts. + keepLatestCheckIntervalMinutes: number; }; export type SocialLink = { @@ -200,6 +207,7 @@ export const SYNC_SCHEDULER_MAX_CONCURRENT_DEFAULT = 2; export const SYNC_SCHEDULER_MAX_CONCURRENT_MAX = 16; export const SYNC_SCHEDULER_BACKOFF_BASE_DEFAULT_MINUTES = 30; export const SYNC_SCHEDULER_BACKOFF_MAX_DEFAULT_MINUTES = 1440; +export const KEEP_LATEST_CHECK_DEFAULT_INTERVAL_MINUTES = 1440; // Internal-heartbeat cadence bounds. 0 means "off" (use an external cron // heartbeat); any other value is clamped into [MIN, MAX] seconds. The floor @@ -218,6 +226,7 @@ export function defaultSyncScheduler(): SyncSchedulerSettings { backoffBaseMinutes: SYNC_SCHEDULER_BACKOFF_BASE_DEFAULT_MINUTES, backoffMaxMinutes: SYNC_SCHEDULER_BACKOFF_MAX_DEFAULT_MINUTES, heartbeatSeconds: SYNC_HEARTBEAT_DEFAULT_SECONDS, + keepLatestCheckIntervalMinutes: KEEP_LATEST_CHECK_DEFAULT_INTERVAL_MINUTES, }; } @@ -290,6 +299,11 @@ export function sanitizeSyncScheduler(value: unknown): SyncSchedulerSettings { ), ), heartbeatSeconds: clampHeartbeatSeconds(r.heartbeatSeconds), + keepLatestCheckIntervalMinutes: clampPositiveInt( + r.keepLatestCheckIntervalMinutes, + d.keepLatestCheckIntervalMinutes, + SYNC_INTERVAL_MAX_MINUTES, + ), }; } diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog ## [Unreleased] +- **New per-channel "keep latest N" retention rule that protects recent videos from the Clean-audio sweep and pins any that get deleted from their source (backend; UI lands in a later change).** A channel can set `keepLatest` (in `config.json` for now) to shield its newest N videos — by upload date, a rolling window — from the **Clean audio** cleanup: those dirs are skipped just like a `do-not-clean` marker, and the snapshot's reclaim estimates/cleanup buckets exclude them (a new `keptCount` is recorded). Because a kept video can later be **deleted from its source** (YouTube etc.) and become irreplaceable, a new **kept-deletion check** re-probes just the kept window's availability (reusing `runAvailabilityCheck` with `onlyIds` + `recheck-non-deleted`) and **permanently pins** any video found `deleted`/`private`/`members_only` with a `do-not-clean` marker, so it survives even after it rolls out of the window. The check runs on demand via the new `check-kept-deleted` managed job (`checkKeptDeletedAction`, re-runnable/bookmarkable) and automatically from the sync scheduler on its own cadence (new `syncScheduler.keepLatestCheckIntervalMinutes`, default daily; per-channel `lastKeptCheckAt` state; suppressed during quiet hours, capped per tick, and skipped for a channel just synced this tick). This is phase 1 of a larger **video-persistence** subsystem (per-download persistence rules, a separate saved-video store, and backups follow). See `common/lib/channelConfig.ts` (`keepLatest`), the new `common/controller/keptVideos.ts` (`computeKeptVideoIds` + tests) and `common/controller/checkKeptDeleted.ts`, `common/controller/cleanAudioFromTranscribed.ts`, `common/controller/channelSnapshot.ts`, and the scheduler wiring in `common/lib/settings.ts`, `common/jobs/syncSchedulerState.ts`, and `editor/app/scheduler/runTick.ts`. - **The monitor widget can now carry optional control buttons (Pause/Resume Transcriptions, Drain all).** The `/widget` view stays read-only by default, but a new **Show control buttons** option in the widget builder (URL flag `controls=1`) adds an interactive row at the top with the same **Pause Transcriptions** toggle (between-segment GPU release) and **Drain all** as the full app — so a pinned/iframe monitor can pause the GPU or wind work down without opening the editor. The controls stay visible even when the widget is otherwise idle (so you can pause preemptively), and the widget refetches its worker payload on a pause/resume so the toggle flips immediately instead of waiting for the next poll. Defaults keep the widget control-free, so existing links render unchanged. See `editor/app/widget/lib/config.ts`, the new `editor/app/widget/components/WidgetControls.tsx`, `editor/app/widget/components/MonitorWidget.tsx`, and `editor/app/widget/builder/components/WidgetBuilder.tsx`. - **"Pause Transcriptions" now frees the GPU between parakeet segments instead of running the in-flight video to completion.** The global pause control (renamed from "Pause all" → **Pause Transcriptions** / **Resume Transcriptions**) used to only stop handing out new worker slots — any in-flight transcription kept running until its whole file was done, so the GPU stayed busy. Pausing now *also* sends a graceful between-segment stop to any partial-capable in-flight job: a **parakeet** worker finishes the current window, caches it (`win-NNNN.json`), and exits `paused` (a skip, not a failure), so the GPU frees within one segment and the video resumes from its cached windows on the next run — the same mechanism as the per-worker **Stop & keep progress** button, now wired into the global pause. Non-parakeet engines (whisper-cpp, chough) keep prior behavior: they stop taking new work but run their in-flight file to completion. The button is also surfaced on the **Active Jobs** screen (`/jobs/active`) next to **Drain all**, not just the Workers page. `resumeAll()` restores each worker's pre-pause state as before. See `common/jobs/workerPool.ts` (`pauseAll`), the new `editor/app/jobs/components/PauseTranscriptionsButton.tsx` (shared by both screens), `editor/app/workers/components/WorkersView.tsx`, and `editor/app/jobs/active/page.tsx`. - **"Drain all" (and the per-runner Drain) no longer hangs the auto-transcribe runner.** Draining could leave the runner stuck "draining" forever — appearing to freeze the app — whenever one of its per-video transcription units was *parked* in the worker pool waiting for a free slot at the moment drain fired (much more likely now that auto units yield slots to manual transcriptions). The runner forwarded only its hard-cancel signal — never the soft `drainSignal` — to a parked unit's `pool.acquire()`, so a soft drain could never unblock it: the unit's promise never settled, the runner's in-flight count never reached zero, and its drain loop spun indefinitely (hard **Cancel**/**Stop** always worked, since that signal *was* forwarded). The runner now threads `ctx.drainSignal` into each transcription unit, so a parked (not-yet-started) unit unblocks and is skipped on drain while a unit already transcribing finishes normally — correct drain semantics, and the runner finalizes promptly. See `common/controller/autoRunner.ts` (the `launchUnit` drainSignal wiring) and the new parked-unit drain regression test in `editor/e2e/auto-queue.spec.ts`. diff --git a/editor/app/channels/[slug]/whisperActions.ts b/editor/app/channels/[slug]/whisperActions.ts @@ -25,6 +25,7 @@ import { } from "yt-dlp-transcript-common/controller/failedTranscodings"; import { clearFailedTranscriptions } from "yt-dlp-transcript-common/controller/failedTranscriptions"; import { cleanAudioFromTranscribed } from "yt-dlp-transcript-common/controller/cleanAudioFromTranscribed"; +import { checkKeptDeleted } from "yt-dlp-transcript-common/controller/checkKeptDeleted"; import { cleanExtraAudioFormats } from "yt-dlp-transcript-common/controller/cleanExtraAudioFormats"; import { removeWrongFormatAudio } from "yt-dlp-transcript-common/controller/removeWrongFormatAudio"; import { verifyTranscripts } from "yt-dlp-transcript-common/controller/verifyTranscripts"; @@ -367,6 +368,36 @@ export async function cleanAudioAction( }); } +// Re-probe the channel's keep-latest window for source deletion and pin any +// gone videos (do-not-clean) so they survive even after rolling out of the +// window. Mirrors cleanAudioAction's managed-job shape. Also driven by the sync +// scheduler (editor/app/scheduler/runTick.ts). +export async function checkKeptDeletedAction( + slug: string, + queueKey?: string, +): Promise<StreamActionResult> { + const paths = getPaths(); + return runManagedFunction({ + kind: "check-kept-deleted", + queueKey: resolveQueueKey(channelQueueKey(slug), queueKey), + paths, + channelSlug: slug, + spec: { kind: "check-kept-deleted", slug, params: { queueKey } }, + fn: async (onLog, signal) => { + const result = await checkKeptDeleted({ + channelSlug: slug, + paths, + onLog, + signal, + }); + onLog( + `Kept-deletion check: inspected ${result.kept}, ${result.deleted} gone from source, pinned ${result.pinned}.`, + ); + revalidatePath(`/channels/${slug}`); + }, + }); +} + export type VerifyResult = | { ok: true; duplicates: string[]; missing: string[] } | { ok: false; error: string }; diff --git a/editor/app/jobs/jobKindLabels.ts b/editor/app/jobs/jobKindLabels.ts @@ -14,6 +14,8 @@ const JOB_KIND_LABELS: Record<string, string> = { "download-missing-subs": "Download missing subs", "import-one": "Import video", "retry-bucket": "Retry", + "clean-audio-transcribed": "Clean audio", + "check-kept-deleted": "Check kept videos", sync: "Sync", }; diff --git a/editor/app/jobs/runJobSpec.ts b/editor/app/jobs/runJobSpec.ts @@ -14,6 +14,7 @@ import { syncAction, } from "../channels/[slug]/pipelineActions"; import { + checkKeptDeletedAction, cleanAudioAction, cleanExtraAudioFormatsAction, clearFailedTranscodingsAction, @@ -128,6 +129,8 @@ export async function runJobSpec(spec: JobSpec): Promise<StreamActionResult> { return removeWrongFormatAudioAction(spec.slug, queueKey); case "clean-audio-transcribed": return cleanAudioAction(spec.slug, queueKey); + case "check-kept-deleted": + return checkKeptDeletedAction(spec.slug, queueKey); default: return { ok: false, error: `Cannot re-run job kind: ${spec.kind}` }; } diff --git a/editor/app/scheduler/runTick.ts b/editor/app/scheduler/runTick.ts @@ -19,6 +19,7 @@ import { type SchedulerState, } from "yt-dlp-transcript-common/jobs/syncSchedulerState"; import { syncAction } from "../channels/[slug]/pipelineActions"; +import { checkKeptDeletedAction } from "../channels/[slug]/whisperActions"; export type SchedulerTickResult = { ok: boolean; @@ -30,6 +31,8 @@ export type SchedulerTickResult = { running: number; // Channels for which a sync was queued this tick. queued: string[]; + // Channels for which a keep-latest deletion check was queued this tick. + keptChecksQueued: string[]; // Channels deliberately held back, with a reason. skipped: SchedulerSkip[]; }; @@ -58,6 +61,7 @@ export async function runSchedulerTick(): Promise<SchedulerTickResult> { reason: "tick already running", running: 0, queued: [], + keptChecksQueued: [], skipped: [], }; } @@ -85,6 +89,7 @@ export async function runSchedulerTick(): Promise<SchedulerTickResult> { reason: "scheduler disabled", running, queued: [], + keptChecksQueued: [], skipped: [], }; } @@ -131,6 +136,34 @@ export async function runSchedulerTick(): Promise<SchedulerTickResult> { scheduler.quietHoursStart, scheduler.quietHoursEnd, ); + + // Keep-latest deletion checks run on their own cadence, on the per-channel + // local queue (channel:<slug>) — separate from the platform download queue a + // sync uses, so the two don't serialize against each other. Suppressed during + // quiet hours, capped per tick like syncs, and skipped for any channel we + // just queued a sync for (avoid hitting the source twice in one tick). + const keptChecksQueued: string[] = []; + if (!quiet) { + const queuedThisTick = new Set(queued); + const intervalMs = scheduler.keepLatestCheckIntervalMinutes * 60_000; + const dueForKeptCheck = channels.filter((c) => { + if (!c.config.keepLatest || c.config.keepLatest <= 0) return false; + if (queuedThisTick.has(c.slug)) return false; + const last = channelState(state, c.slug).lastKeptCheckAt ?? 0; + return now - last >= intervalMs; + }); + for (const c of dueForKeptCheck.slice(0, scheduler.maxConcurrentSyncs)) { + const result = await checkKeptDeletedAction(c.slug); + if (!result.ok) { + skipped.push({ slug: c.slug, reason: `kept-check: ${result.error}` }); + continue; + } + channelState(state, c.slug).lastKeptCheckAt = now; + keptChecksQueued.push(c.slug); + void result.stream.cancel(); + } + } + recordRun(state, { at: now, queued, skipped }); await writeSchedulerState(paths, state); @@ -140,6 +173,7 @@ export async function runSchedulerTick(): Promise<SchedulerTickResult> { reason: quiet ? "quiet hours" : undefined, running, queued, + keptChecksQueued, skipped, }; } finally { diff --git a/editor/app/settings/actions.ts b/editor/app/settings/actions.ts @@ -144,6 +144,9 @@ export async function saveSettingsAction( backoffBaseMinutes: intOrNaN("syncSchedulerBackoffBaseMinutes"), backoffMaxMinutes: intOrNaN("syncSchedulerBackoffMaxMinutes"), heartbeatSeconds: intOrNaN("syncSchedulerHeartbeatSeconds"), + keepLatestCheckIntervalMinutes: intOrNaN( + "syncSchedulerKeepLatestCheckIntervalMinutes", + ), }; let socialInput: unknown;