diff --git a/README.md b/README.md index bd165c7..08e6d0e 100644 --- a/README.md +++ b/README.md @@ -909,6 +909,28 @@ All configuration is done via environment variables. Below is the complete refer | `ARR_BRIDGE_ENABLED` | `false` | Enable the fake qBittorrent API server | | `ARR_BRIDGE_PORT` | `8282` | Port for the *arr bridge (add as qBittorrent in Radarr/Sonarr) | +### 🔄 Provider Reconciliation (opt-in) + +After the initial historical library import performed in Radarr/Sonarr, +provider reconciliation detects newly completed files added directly to a +configured debrid provider. It compares snapshots in SQLite, creates +mount-backed symlinks in the existing Arr library, and asks Radarr/Sonarr to +rescan the affected movie or series. It does not copy media locally and does +not delete or repair provider content. + +| Variable | Default | Description | +| --- | --- | --- | +| `PROVIDER_RECONCILIATION_ENABLED` | `false` | Enable reconciliation | +| `PROVIDER_RECONCILIATION_RECENT_INTERVAL_MS` | `900000` | Recent snapshot interval | +| `PROVIDER_RECONCILIATION_FULL_INTERVAL_MS` | `21600000` | Full snapshot interval | +| `PROVIDER_RECONCILIATION_RADARR_URL` | — | Radarr API endpoint | +| `PROVIDER_RECONCILIATION_SONARR_URL` | — | Sonarr API endpoint | +| `PROVIDER_RECONCILIATION_MOUNT_BASE` | `/mnt/schrodrive` | Shared mount base | + +The provider adapter uses the common `DebridProvider` contract. See +[`docs/provider-reconciliation.md`](docs/provider-reconciliation.md) for +provider capability and live-validation status. + ### 📁 Organiser | Variable | Default | Description | diff --git a/docs/provider-reconciliation.md b/docs/provider-reconciliation.md new file mode 100644 index 0000000..2f3e5fc --- /dev/null +++ b/docs/provider-reconciliation.md @@ -0,0 +1,50 @@ +# Provider reconciliation + +SchröDrive's provider reconciliation layer is provider-agnostic. It consumes +the existing `DebridProvider` contract and never calls provider delete, repair, +or dead-scanner operations. + +The historical library import is a one-time operation performed through the +Radarr/Sonarr library import flow. This worker handles the ongoing case where +a new completed file is added directly to a provider after that library is +already configured. + +```text +provider listTorrents/fetchDirectories + -> normalized snapshot and local SQLite diff + -> mount-backed symlink in the Arr library + -> RescanMovie / RescanSeries + -> Radarr, Sonarr, Plex and Jellyfin see the same provider-backed file +``` + +The worker is disabled by default with +`PROVIDER_RECONCILIATION_ENABLED=false`. + +## Contract coverage + +Every provider registered by SchröDrive implements the common provider +contract used by the worker: stable torrent IDs, normalized torrent files, a +completed virtual directory tree, and a provider mount path. Recent polling is +an optimization; full polling plus the local snapshot diff is the source of +truth. + +| Provider | Contract adapter | Live reconciliation validation | +| --- | --- | --- | +| TorBox | common adapter | pending provider-backed test | +| Real-Debrid | common adapter | pending provider-backed test | +| AllDebrid | common adapter | validated with provider-backed Sonarr E2E | +| Premiumize | common adapter | pending provider-backed test | +| Debrid-Link | common adapter | pending provider-backed test | +| Deepbrid | common adapter | pending provider-backed test | +| Offcloud | common adapter | pending provider-backed test | +| Put.io | common adapter | pending provider-backed test | +| MegaDebrid | common adapter | pending provider-backed test | +| Seedr | common adapter | pending provider-backed test | +| PikPak | common adapter | pending provider-backed test | + +“Common adapter” means the provider satisfies SchröDrive's existing interface; +it does not claim that a live account has been tested. A provider is skipped +when it is not configured or cannot expose a complete read-only snapshot. + +The Arr library contains symlinks only. Provider content remains on the mount; +the reconciliation layer does not download or copy media to local storage. diff --git a/src/core/config.ts b/src/core/config.ts index ab0cffd..a9b15a7 100644 --- a/src/core/config.ts +++ b/src/core/config.ts @@ -59,6 +59,20 @@ export const config = { alldebridApiKey: process.env.ALLDEBRID_API_KEY || "", alldebridApiBase: process.env.ALLDEBRID_API_BASE || "https://api.alldebrid.com/v4", alldebridAgent: process.env.ALLDEBRID_AGENT || "schrodrive", + // Provider reconciliation is opt-in and disabled by default. + providerReconciliationEnabled: String(process.env.PROVIDER_RECONCILIATION_ENABLED ?? "false").toLowerCase() === "true", + providerReconciliationRecentIntervalMs: Number(process.env.PROVIDER_RECONCILIATION_RECENT_INTERVAL_MS || 900000), + providerReconciliationFullIntervalMs: Number(process.env.PROVIDER_RECONCILIATION_FULL_INTERVAL_MS || 21600000), + providerReconciliationRecentLimit: Number(process.env.PROVIDER_RECONCILIATION_RECENT_LIMIT || 30), + providerReconciliationRunFullOnStart: String(process.env.PROVIDER_RECONCILIATION_RUN_FULL_ON_START ?? "true").toLowerCase() !== "false", + providerReconciliationDryRun: String(process.env.PROVIDER_RECONCILIATION_DRY_RUN ?? "false").toLowerCase() === "true", + providerReconciliationRadarrUrl: process.env.PROVIDER_RECONCILIATION_RADARR_URL || "", + providerReconciliationRadarrApiKey: process.env.PROVIDER_RECONCILIATION_RADARR_API_KEY || "", + providerReconciliationSonarrUrl: process.env.PROVIDER_RECONCILIATION_SONARR_URL || "", + providerReconciliationSonarrApiKey: process.env.PROVIDER_RECONCILIATION_SONARR_API_KEY || "", + providerReconciliationMountBase: process.env.PROVIDER_RECONCILIATION_MOUNT_BASE || "/mnt/schrodrive", + providerReconciliationMoviesLibraryPath: process.env.PROVIDER_RECONCILIATION_MOVIES_LIBRARY_PATH || "", + providerReconciliationShowsLibraryPath: process.env.PROVIDER_RECONCILIATION_SHOWS_LIBRARY_PATH || "", // AllDebrid WebDAV (if supported) alldebridWebdavUrl: process.env.ALLDEBRID_WEBDAV_URL || "", alldebridWebdavUsername: process.env.ALLDEBRID_WEBDAV_USERNAME || "", diff --git a/src/index.ts b/src/index.ts index 2ee6c8b..6d617fc 100644 --- a/src/index.ts +++ b/src/index.ts @@ -12,6 +12,8 @@ import { getDb, closeDb, pruneOldEntries, pruneExpiredStrmCodes } from "./core/d import { startStrmServer, stopStrmServer } from "./services/strmService"; import { startCloudLinksBridge, stopCloudLinksBridge } from "./services/cloudLinks/bridge"; import { startArrBridge, stopArrBridge } from "./services/arrBridge"; +import { startProviderReconciliation } from "./services/providerReconciliationRuntime"; +import type { ProviderReconciliationWorker } from "./services/providerReconciliation"; const program = new Command(); program @@ -31,6 +33,7 @@ program } // Register graceful shutdown handlers + let providerReconciliationWorker: ProviderReconciliationWorker | undefined; const shutdown = () => { console.log(`[${new Date().toISOString()}][serve] Shutting down — unmounting FUSE drives...`); try { @@ -43,6 +46,7 @@ program stopStrmServer().catch(() => {}); stopCloudLinksBridge().catch(() => {}); stopArrBridge().catch(() => {}); + providerReconciliationWorker?.stop(); setTimeout(() => { console.log(`[${new Date().toISOString()}][serve] Closing database and exiting...`); closeDb(); @@ -77,7 +81,13 @@ program }); } - promises.push(mountVirtualDrive()); + // The provider reconciliation worker submits Arr rescans against this + // mount. Wait until mountVirtualDrive has established the visible paths + // before starting the worker below; otherwise Arr can reject the first + // scan as a missing file during FUSE startup. + await mountVirtualDrive().catch((err: any) => { + console.error(`[${new Date().toISOString()}][serve] Virtual drive mount failed (non-fatal): ${err?.message}`); + }); } if (config.runDeadScannerWatch) { @@ -97,6 +107,8 @@ program console.log("[serve] Starting media server watchlist poller (RUN_WATCHLIST_POLLER=true)"); startWatchlistPoller(); } + + providerReconciliationWorker = startProviderReconciliation(); // Start the main server startServer(); diff --git a/src/providers/alldebrid.ts b/src/providers/alldebrid.ts index 7e66ffc..36ed009 100644 --- a/src/providers/alldebrid.ts +++ b/src/providers/alldebrid.ts @@ -618,6 +618,30 @@ export class AllDebridProvider implements DebridProvider { return result; } + /** + * Returns completed virtual directories for an explicit magnet subset. + * This is read-only and lets provider reconciliation implement + * a bounded recent scan without downloading or mutating provider state. + */ + async fetchDirectoriesForIds(magnets: Array>): Promise { + const ids = magnets.map((magnet) => String(magnet.id)).filter(Boolean); + const fileTrees = await this.fetchFileTrees(ids); + return magnets.map((magnet) => { + const id = String(magnet.id); + const files: VirtualFile[] = (fileTrees.get(id) || []).map((file) => ({ + id: file.path, + name: file.path, + size: file.size, + })); + return { + id, + name: sanitiseName(magnet.filename || magnet.name || id), + originalName: magnet.filename || magnet.name || id, + files, + }; + }); + } + /** * Fetches the complete magnet list from AllDebrid and converts it into * virtual directories. Only includes fully downloaded magnets diff --git a/src/services/providerReconciliation.ts b/src/services/providerReconciliation.ts new file mode 100644 index 0000000..5899b9d --- /dev/null +++ b/src/services/providerReconciliation.ts @@ -0,0 +1,463 @@ +/** + * Provider-agnostic direct-file reconciliation for mounted debrid providers. + * + * Providers do not need a push change feed for this path: the adapter + * reconciles read-only torrent status plus completed file trees, emits stable + * direct-file events, and hands added/changed files to Arr. It is deliberately + * opt-in and separate from provider delete/repair lifecycle services. + */ + +import { createHash } from 'node:crypto'; +import { promises as fsp } from 'node:fs'; +import path from 'node:path'; +import { getDb } from '../core/db'; +import type { DebridProvider, TorrentInfo, VirtualDirectory } from '../providers'; +import { classifyTorrent } from '../core/mediaClassifier'; +import { parseMediaFilename } from './mediaParser'; + +export type DirectFileAction = 'added' | 'changed' | 'deleted'; +export type SourceCategory = 'Movies' | 'Shows'; +export type ArrKind = 'radarr' | 'sonarr'; + +export interface ProviderSnapshot { + provider?: string; + providerItemId: string; + name: string; + directoryName?: string; + status: string; + progress?: number; + files: Array<{ path: string; size: number }>; + observedAt: string; +} + +export interface DirectFileEvent { + provider: string; + providerItemId: string; + action: DirectFileAction; + path: string; + tree: Array<{ path: string; size: number }>; + sourceCategory: SourceCategory; + observedAt: string; + stableDedupeKey: string; +} + +export interface ArrRoute { + kind: ArrKind; + baseUrl: string; + apiKey: string; + /** Absolute path visible inside the target Arr container. */ + sourcePathPrefix?: string; + /** Import mode used by Arr for provider-backed paths. */ + importMode?: 'Move' | 'Copy'; + /** Optional Arr-visible library root. When set, media is exposed by symlink. */ + symlinkLibraryPath?: string; +} + +export interface ArrCommandResult { + commandId: string; + status: string; + result?: string; +} + +export interface ArrClient { + submitScan(route: ArrRoute, event: DirectFileEvent): Promise; + getCommand(route: ArrRoute, commandId: string): Promise; +} + +export interface IntakeStateStore { + getItem(providerItemId: string): IntakeState | undefined; + listItems(): IntakeState[]; + saveItem(state: IntakeState): void; + hasEvent(key: string): boolean; + saveEvent(event: DirectFileEvent, arr?: ArrCommandResult): void; + getCursor(): { recentAt?: string; fullAt?: string }; + saveCursor(mode: 'recent' | 'full', observedAt: string): void; +} + +export interface IntakeState { + provider?: string; + providerItemId: string; + fingerprint: string; + path: string; + tree: Array<{ path: string; size: number }>; + sourceCategory: SourceCategory; + lastAction: DirectFileAction; + commandId?: string; + terminalStatus?: string; + updatedAt: string; +} + +export class InMemoryIntakeStateStore implements IntakeStateStore { + private readonly items = new Map(); + private readonly events = new Map(); + private cursor: { recentAt?: string; fullAt?: string } = {}; + + getItem(id: string): IntakeState | undefined { return this.items.get(id); } + listItems(): IntakeState[] { return [...this.items.values()]; } + saveItem(state: IntakeState): void { this.items.set(state.providerItemId, state); } + hasEvent(key: string): boolean { return this.events.has(key); } + saveEvent(event: DirectFileEvent): void { this.events.set(event.stableDedupeKey, event); } + getCursor(): { recentAt?: string; fullAt?: string } { return { ...this.cursor }; } + saveCursor(mode: 'recent' | 'full', observedAt: string): void { this.cursor[mode === 'recent' ? 'recentAt' : 'fullAt'] = observedAt; } +} + +/** Persistent reconciliation state in the configured SchröDrive DB. */ +export class SqliteIntakeStateStore implements IntakeStateStore { + constructor() { + getDb().exec(`CREATE TABLE IF NOT EXISTS provider_reconciliation_intake ( + provider_item_id TEXT PRIMARY KEY, + state_json TEXT NOT NULL, + updated_at INTEGER NOT NULL + )`); + getDb().exec(`CREATE TABLE IF NOT EXISTS provider_reconciliation_events ( + dedupe_key TEXT PRIMARY KEY, + event_json TEXT NOT NULL, + arr_json TEXT, + created_at INTEGER NOT NULL + )`); + getDb().exec(`CREATE TABLE IF NOT EXISTS provider_reconciliation_cursor ( + name TEXT PRIMARY KEY, + observed_at TEXT NOT NULL + )`); + } + getItem(id: string): IntakeState | undefined { + const row = getDb().prepare('SELECT state_json FROM provider_reconciliation_intake WHERE provider_item_id = ?').get(id) as { state_json?: string } | undefined; + return row?.state_json ? JSON.parse(row.state_json) as IntakeState : undefined; + } + listItems(): IntakeState[] { + return (getDb().prepare('SELECT state_json FROM provider_reconciliation_intake').all() as Array<{ state_json: string }>) + .flatMap((row) => { try { return [JSON.parse(row.state_json) as IntakeState]; } catch { return []; } }); + } + saveItem(state: IntakeState): void { + getDb().prepare(`INSERT INTO provider_reconciliation_intake(provider_item_id,state_json,updated_at) + VALUES (?,?,?) ON CONFLICT(provider_item_id) DO UPDATE SET state_json=excluded.state_json,updated_at=excluded.updated_at`) + .run(state.providerItemId, JSON.stringify(state), Date.parse(state.updatedAt)); + } + hasEvent(key: string): boolean { + return !!getDb().prepare('SELECT 1 FROM provider_reconciliation_events WHERE dedupe_key = ?').get(key); + } + saveEvent(event: DirectFileEvent, arr?: ArrCommandResult): void { + getDb().prepare(`INSERT OR IGNORE INTO provider_reconciliation_events(dedupe_key,event_json,arr_json,created_at) + VALUES (?,?,?,?)`).run(event.stableDedupeKey, JSON.stringify(event), arr ? JSON.stringify(arr) : null, Date.parse(event.observedAt)); + } + getCursor(): { recentAt?: string; fullAt?: string } { + const rows = getDb().prepare('SELECT name, observed_at FROM provider_reconciliation_cursor').all() as Array<{ name: string; observed_at: string }>; + return Object.fromEntries(rows.map((row) => [row.name === 'recent' ? 'recentAt' : 'fullAt', row.observed_at])); + } + saveCursor(mode: 'recent' | 'full', observedAt: string): void { + getDb().prepare(`INSERT INTO provider_reconciliation_cursor(name,observed_at) VALUES (?,?) + ON CONFLICT(name) DO UPDATE SET observed_at=excluded.observed_at`).run(mode, observedAt); + } +} + +export class HttpArrClient implements ArrClient { + async submitScan(route: ArrRoute, event: DirectFileEvent): Promise { + let commandName = route.kind === 'radarr' ? 'DownloadedMoviesScan' : 'DownloadedEpisodesScan'; + const providerPath = route.sourcePathPrefix + ? `${route.sourcePathPrefix.replace(/\/$/, '')}/${event.path.replace(/^\/+/, '')}` + : event.path; + const scanPath = route.symlinkLibraryPath + ? await exposeAsSymlink(route, event, providerPath) + : providerPath; + let commandBody: Record = { name: commandName, path: scanPath, importMode: route.importMode || 'Copy' }; + if (route.symlinkLibraryPath) { + const entityId = await findArrEntityId(route, event); + commandName = route.kind === 'radarr' ? 'RescanMovie' : 'RescanSeries'; + commandBody = { name: commandName, [route.kind === 'radarr' ? 'movieId' : 'seriesId']: entityId }; + } + const response = await fetch(`${route.baseUrl.replace(/\/$/, '')}/api/v3/command`, { + method: 'POST', + headers: { 'X-Api-Key': route.apiKey, 'Content-Type': 'application/json' }, + body: JSON.stringify(commandBody), + }); + if (!response.ok) throw new Error(`Arr command submission failed: HTTP ${response.status}`); + const body = await response.json() as { id?: number; status?: string; result?: string }; + if (!body.id) throw new Error('Arr command response did not include an id'); + return { commandId: String(body.id), status: body.status || 'queued', result: body.result }; + } + + async getCommand(route: ArrRoute, commandId: string): Promise { + const response = await fetch(`${route.baseUrl.replace(/\/$/, '')}/api/v3/command/${encodeURIComponent(commandId)}`, { + headers: { 'X-Api-Key': route.apiKey }, + }); + if (!response.ok) throw new Error(`Arr command status failed: HTTP ${response.status}`); + const body = await response.json() as { id?: number; status?: string; result?: string }; + return { commandId, status: body.status || 'unknown', result: body.result }; + } +} + +function normalizedTitle(value: string): string { + return value.toLocaleLowerCase().replace(/[^\p{L}\p{N}]+/gu, ''); +} + +async function findArrEntityId(route: ArrRoute, event: DirectFileEvent): Promise { + const filename = path.basename(event.path); + const parsed = parseMediaFilename(filename, event.path); + const title = normalizedTitle(parsed.title || event.path.split('/').filter(Boolean).slice(-2, -1)[0] || ''); + const endpoint = route.kind === 'radarr' ? 'movie' : 'series'; + const response = await fetch(`${route.baseUrl.replace(/\/$/, '')}/api/v3/${endpoint}`, { + headers: { 'X-Api-Key': route.apiKey }, + }); + if (!response.ok) throw new Error(`Arr ${endpoint} lookup failed: HTTP ${response.status}`); + const records = await response.json() as Array<{ id?: number; title?: string }>; + const match = records.find((record) => record.id && normalizedTitle(record.title || '') === title); + if (!match?.id) throw new Error(`Arr ${endpoint} record not found for ${parsed.title || filename}`); + return match.id; +} + +function safeSegment(value: string): string { + return value.replace(/[\\/:*?"<>|]/g, ' ').replace(/\s+/g, ' ').trim() || 'Unknown'; +} + +/** + * Creates the Arr-facing library entry without copying provider data. The + * target is deliberately a symlink into the shared SchröDrive mount. + */ +async function exposeAsSymlink(route: ArrRoute, event: DirectFileEvent, providerPath: string): Promise { + const library = route.symlinkLibraryPath!; + const filename = path.basename(event.path); + const parsed = parseMediaFilename(filename, event.path); + const pathParts = event.path.split('/').filter(Boolean); + const title = safeSegment(parsed.title || pathParts[pathParts.length - 2] || path.parse(filename).name); + const directory = event.sourceCategory === 'Shows' + ? path.join(library, title, `Season ${parsed.season ?? 1}`) + : path.join(library, parsed.year ? `${title} (${parsed.year})` : title); + const destination = path.join(directory, filename); + await fsp.mkdir(directory, { recursive: true }); + const existing = await fsp.lstat(destination).catch(() => undefined); + if (existing?.isSymbolicLink()) { + const current = await fsp.readlink(destination); + if (current === providerPath) return destination; + // Multiple providers may expose the same title. Keep the first healthy + // provider-backed link instead of oscillating the library on every poll. + return destination; + } else if (existing) { + throw new Error(`Refusing to overwrite non-symlink Arr library entry: ${destination}`); + } + // Use the container-visible absolute mount path. A relative link would be + // resolved against the host bind source, which differs between SchröDrive + // and Arr containers. + await fsp.symlink(providerPath, destination); + return destination; +} + +export interface ProviderReadOnlySource { + listSnapshot(): Promise; + listRecentSnapshot?(limit: number): Promise; +} + +/** Adapts SchröDrive's common provider contract to reconciliation snapshots. */ +export class ProviderSnapshotSource implements ProviderReadOnlySource { + constructor(private readonly provider: Pick & { fetchDirectoriesForIds?: (torrents: TorrentInfo[]) => Promise }, private readonly namespaceIds = true) {} + + async listSnapshot(): Promise { + const observedAt = new Date().toISOString(); + const torrents = await this.provider.listTorrents(); + const directories = await this.provider.fetchDirectories(); + const trees = new Map(directories.map((directory) => [String(directory.id), directory])); + return this.toSnapshots(torrents, trees, observedAt); + } + + async listRecentSnapshot(limit: number): Promise { + const observedAt = new Date().toISOString(); + const torrents = (await this.provider.listTorrents()) + .sort((a, b) => (b.addedAt?.getTime() || 0) - (a.addedAt?.getTime() || 0)) + .slice(0, Math.max(0, limit)); + const directories = this.provider.fetchDirectoriesForIds + ? await this.provider.fetchDirectoriesForIds(torrents) + : await this.provider.fetchDirectories(); + return this.toSnapshots(torrents, new Map(directories.map((directory) => [String(directory.id), directory])), observedAt); + } + + private toSnapshots(torrents: TorrentInfo[], trees: Map, observedAt: string): ProviderSnapshot[] { + return torrents.map((torrent) => { + const directory = trees.get(String(torrent.id)); + return { provider: this.provider.id, providerItemId: this.namespaceIds ? `${this.provider.id}:${torrent.id}` : String(torrent.id), name: torrent.name, directoryName: directory?.name || directory?.originalName, status: torrent.status, progress: torrent.progress, + files: ((directory?.files?.length ? directory.files : torrent.files) || []).map((file) => ({ path: 'path' in file ? file.path : file.name, size: file.size })), observedAt }; + }); + } +} + +export interface IntakeOptions { + dryRun?: boolean; + maxAttempts?: number; + routeFor: (category: SourceCategory, provider?: string) => ArrRoute; + onEvent?: (event: DirectFileEvent) => Promise | void; + onReview?: (event: DirectFileEvent, error: Error) => Promise | void; +} + +const VIDEO_EXTENSIONS = new Set(['3g2', '3gp', 'avi', 'flv', 'mkv', 'mk3d', 'm4v', 'mov', 'mp2', 'mp4', 'mpe', 'mpeg', 'mpg', 'mpv', 'ts', 'm2ts', 'webm', 'wmv', 'ogm']); +// Keep subtitles and sidecar subtitle attachments with the video tree. +const SUBTITLE_EXTENSIONS = new Set(['ass', 'idx', 'mpsub', 'sbv', 'smi', 'srt', 'ssa', 'sub', 'sup', 'vtt']); + +export function isMediaFile(filePath: string): boolean { + const extension = filePath.split('.').pop()?.toLowerCase() || ''; + return VIDEO_EXTENSIONS.has(extension) || SUBTITLE_EXTENSIONS.has(extension); +} + +function categoryFor(snapshot: ProviderSnapshot): SourceCategory { + return classifyTorrent(snapshot.name, snapshot.files.map((file) => file.path)) === 'shows' ? 'Shows' : 'Movies'; +} + +function isFinished(snapshot: ProviderSnapshot): boolean { + const status = snapshot.status.toLowerCase(); + return snapshot.progress === undefined + ? ['finished', 'downloaded', 'completed', 'seeding', 'ready', 'cached'].includes(status) + : snapshot.progress >= 100 || ['finished', 'downloaded', 'completed', 'seeding', 'ready', 'cached'].includes(status); +} + +function fingerprint(snapshot: Pick): string { + return createHash('sha256').update(JSON.stringify({ id: snapshot.providerItemId, status: snapshot.status, files: snapshot.files })).digest('hex'); +} + +function eventKey(provider: string, id: string, action: DirectFileAction, fp: string): string { + return `${provider}:${id}:${action}:${fp}`; +} + +export class ProviderReconciliationIntake { + constructor( + private readonly source: ProviderReadOnlySource, + private readonly arr: ArrClient, + private readonly store: IntakeStateStore, + private readonly options: IntakeOptions, + ) {} + + async reconcile(mode: 'recent' | 'full' = 'full', recentLimit = 30): Promise { + await this.pollPendingCommands(); + const current = mode === 'recent' && this.source.listRecentSnapshot + ? await this.source.listRecentSnapshot(recentLimit) + : await this.source.listSnapshot(); + const seen = new Set(current.map((item) => item.providerItemId)); + const events: DirectFileEvent[] = []; + + for (const item of current) { + if (!isFinished(item) || item.files.length === 0) continue; + item.files = item.files.filter((file) => isMediaFile(file.path)); + if (item.files.length === 0) continue; + const prior = this.store.getItem(item.providerItemId); + const nextFingerprint = fingerprint(item); + const action: DirectFileAction | undefined = !prior || prior.lastAction === 'deleted' + ? 'added' : prior.fingerprint === nextFingerprint ? undefined : 'changed'; + if (!action) continue; + const event = this.makeEvent(item, action, nextFingerprint); + await this.dispatch(event, item, nextFingerprint); + events.push(event); + } + + // A missing status-list item is a removal, but an item still processing is not. + for (const previous of mode === 'full' ? this.store.listItems().filter((item) => item.lastAction !== 'deleted') : []) { + if (seen.has(previous.providerItemId)) continue; + const event: DirectFileEvent = { + provider: previous.provider || 'alldebrid', providerItemId: previous.providerItemId, action: 'deleted', + path: previous.path, tree: [], sourceCategory: previous.sourceCategory, + observedAt: new Date().toISOString(), stableDedupeKey: eventKey(previous.provider || 'alldebrid', previous.providerItemId, 'deleted', previous.fingerprint), + }; + if (!this.store.hasEvent(event.stableDedupeKey)) { + await this.options.onEvent?.(event); + this.store.saveEvent(event); + this.store.saveItem({ ...previous, tree: previous.tree || [], lastAction: 'deleted', updatedAt: event.observedAt }); + events.push(event); + } + } + this.store.saveCursor(mode, new Date().toISOString()); + return events; + } + + /** Reconciles Arr command state after a crash or an interrupted poll. */ + private async pollPendingCommands(): Promise { + for (const item of this.store.listItems()) { + if (!item.commandId || item.terminalStatus === 'completed' || item.terminalStatus === 'failed') continue; + try { + const command = await this.arr.getCommand(this.options.routeFor(item.sourceCategory, item.provider), item.commandId); + this.store.saveItem({ ...item, terminalStatus: command.status, updatedAt: new Date().toISOString() }); + if (command.status === 'failed') { + await this.options.onReview?.({ + provider: item.provider || 'alldebrid', providerItemId: item.providerItemId, action: item.lastAction, + path: item.path, tree: item.tree || [], sourceCategory: item.sourceCategory, + observedAt: new Date().toISOString(), + stableDedupeKey: eventKey(item.provider || 'alldebrid', item.providerItemId, item.lastAction, item.fingerprint), + }, new Error(`Arr command ${item.commandId} failed`)); + } + } catch { + // A transient status failure is retried on the next reconciliation. + } + } + } + + private makeEvent(item: ProviderSnapshot, action: DirectFileAction, fp: string): DirectFileEvent { + const provider = item.provider || 'alldebrid'; + const category = categoryFor(item); + const categoryDirectory = category.toLowerCase(); + const providerPath = item.files[0].path.replace(/^\/+/, ''); + const path = item.directoryName + ? `${categoryDirectory}/${item.directoryName.replace(/^\/+|\/+$/g, '')}/${providerPath}` + : providerPath; + return { + provider, providerItemId: item.providerItemId, action, + path, tree: item.files, sourceCategory: category, + observedAt: item.observedAt, stableDedupeKey: eventKey(provider, item.providerItemId, action, fp), + }; + } + + private async dispatch(event: DirectFileEvent, item: ProviderSnapshot, fp: string): Promise { + if (this.store.hasEvent(event.stableDedupeKey)) return; + await this.options.onEvent?.(event); + if (this.options.dryRun) { + this.store.saveEvent(event); + this.store.saveItem({ provider: event.provider, providerItemId: item.providerItemId, fingerprint: fp, path: event.path, tree: event.tree, sourceCategory: event.sourceCategory, lastAction: event.action, updatedAt: event.observedAt }); + return; + } + const route = this.options.routeFor(event.sourceCategory, event.provider); + let lastError: unknown; + for (let attempt = 1; attempt <= (this.options.maxAttempts || 3); attempt++) { + try { + const command = await this.arr.submitScan(route, event); + this.store.saveEvent(event, command); + this.store.saveItem({ provider: event.provider, providerItemId: item.providerItemId, fingerprint: fp, path: event.path, tree: event.tree, sourceCategory: event.sourceCategory, lastAction: event.action, commandId: command.commandId, terminalStatus: command.status, updatedAt: event.observedAt }); + return; + } catch (error) { + lastError = error; + } + } + const error = lastError instanceof Error ? lastError : new Error(String(lastError)); + await this.options.onReview?.(event, error); + throw error; + } +} + +/** + * Testable scheduler for the fork worker. It is intentionally not started by + * the application entry point; provider reconciliation must explicitly opt in. + */ +export class ProviderReconciliationWorker { + private recentTimer: ReturnType | undefined; + private fullTimer: ReturnType | undefined; + + constructor( + private readonly intake: ProviderReconciliationIntake | ProviderReconciliationIntake[], + private readonly intervals: { recentMs: number; fullMs: number; recentLimit?: number; runFullOnStart?: boolean }, + ) {} + + private intakes(): ProviderReconciliationIntake[] { return Array.isArray(this.intake) ? this.intake : [this.intake]; } + async runRecent(): Promise { return (await Promise.all(this.intakes().map((intake) => intake.reconcile('recent', this.intervals.recentLimit || 30)))).flat(); } + async runFull(): Promise { return (await Promise.all(this.intakes().map((intake) => intake.reconcile('full')))).flat(); } + + start(): void { + if (this.recentTimer || this.fullTimer) return; + this.runRecent().catch(() => undefined); + if (this.intervals.runFullOnStart !== false) this.runFull().catch(() => undefined); + this.recentTimer = setInterval(() => { this.runRecent().catch(() => undefined); }, this.intervals.recentMs); + this.fullTimer = setInterval(() => { this.runFull().catch(() => undefined); }, this.intervals.fullMs); + } + + stop(): void { + if (this.recentTimer) clearInterval(this.recentTimer); + if (this.fullTimer) clearInterval(this.fullTimer); + this.recentTimer = undefined; + this.fullTimer = undefined; + } + + isRunning(): boolean { + return !!this.recentTimer || !!this.fullTimer; + } +} diff --git a/src/services/providerReconciliationCapabilities.ts b/src/services/providerReconciliationCapabilities.ts new file mode 100644 index 0000000..4b732c3 --- /dev/null +++ b/src/services/providerReconciliationCapabilities.ts @@ -0,0 +1,51 @@ +/** + * Provider reconciliation capability contract. + * + * This is intentionally descriptive only: it does not start polling or alter + * any provider. A future generic worker can use it to opt in only when the + * provider exposes the required read-only operations. + */ +export interface ReconciliationCapabilities { + statusList: boolean; + fileTree: boolean; + recentSnapshot: boolean; + fullSnapshot: boolean; + changeDetection: boolean; + pushEvents: boolean; +} + +export type ReconciliationMode = 'disabled' | 'polling-hybrid' | 'polling-full-only' | 'push-only'; + +export interface ReconciliationCapabilityAssessment { + capabilities: ReconciliationCapabilities; + mode: ReconciliationMode; + reason: string; +} + +/** + * Computes the safe mode without assuming that providers share an API. + * Change detection is a worker-side snapshot diff and therefore requires both + * a status listing and a file tree. Recent polling is optional; full polling + * is the safe minimum for a polling worker. + */ +export function assessReconciliationCapabilities( + capabilities: ReconciliationCapabilities, +): ReconciliationCapabilityAssessment { + const hasPolling = capabilities.statusList && capabilities.fileTree && capabilities.fullSnapshot; + const hasRecent = hasPolling && capabilities.recentSnapshot; + const changeDetection = hasPolling && capabilities.changeDetection; + const hasPush = capabilities.pushEvents; + + if (!hasPolling && !hasPush) { + return { capabilities: { ...capabilities, changeDetection }, mode: 'disabled', reason: 'provider exposes neither a complete snapshot contract nor push events' }; + } + if (hasPolling) { + return { capabilities: { ...capabilities, recentSnapshot: hasRecent, changeDetection }, mode: hasRecent ? 'polling-hybrid' : 'polling-full-only', reason: hasRecent ? 'frequent recent polling plus periodic full polling; local snapshot diff is authoritative' : 'periodic full polling with local snapshot diff; recent polling is unavailable' }; + } + if (hasPush) { + return { capabilities: { ...capabilities, changeDetection }, mode: 'push-only', reason: 'provider exposes push events but not a complete polling snapshot contract' }; + } + // Keep an exhaustive fallback for future capability extensions. The current + // branches above make this unreachable, but it must remain a valid mode. + return { capabilities: { ...capabilities, recentSnapshot: hasRecent, changeDetection }, mode: 'polling-full-only', reason: 'fallback polling mode with worker-side snapshot diff' }; +} diff --git a/src/services/providerReconciliationRuntime.ts b/src/services/providerReconciliationRuntime.ts new file mode 100644 index 0000000..90af7c4 --- /dev/null +++ b/src/services/providerReconciliationRuntime.ts @@ -0,0 +1,62 @@ +import { config } from "../core/config"; +import { registry } from "../providers"; +import { parseMediaFilename } from "./mediaParser"; +import { recordOrganizerReview } from "./organizerReview"; +import { + ProviderSnapshotSource, + ProviderReconciliationIntake, + ProviderReconciliationWorker, + HttpArrClient, + SqliteIntakeStateStore, + type ArrRoute, +} from "./providerReconciliation"; + +export function providerReconciliationRoutes(providerId = "alldebrid"): { movies: ArrRoute; shows: ArrRoute } { + const sourcePathPrefix = `${config.providerReconciliationMountBase.replace(/\/$/, '')}/${providerId}`; + return { + movies: { kind: "radarr", baseUrl: config.providerReconciliationRadarrUrl, apiKey: config.providerReconciliationRadarrApiKey, sourcePathPrefix, importMode: "Copy", symlinkLibraryPath: config.providerReconciliationMoviesLibraryPath || undefined }, + shows: { kind: "sonarr", baseUrl: config.providerReconciliationSonarrUrl, apiKey: config.providerReconciliationSonarrApiKey, sourcePathPrefix, importMode: "Copy", symlinkLibraryPath: config.providerReconciliationShowsLibraryPath || undefined }, + }; +} + +/** Builds the opt-in provider reconciliation worker without changing provider lifecycle services. */ +export function createProviderReconciliationWorker(): ProviderReconciliationWorker | undefined { + if (!config.providerReconciliationRadarrUrl || !config.providerReconciliationRadarrApiKey || !config.providerReconciliationSonarrUrl || !config.providerReconciliationSonarrApiKey) { + console.warn("[provider-reconciliation] Radarr/Sonarr routes are incomplete; worker disabled"); + return undefined; + } + const arr = new HttpArrClient(); + const store = new SqliteIntakeStateStore(); + const intakes = config.providers.flatMap((providerId) => { + const provider = registry.get(providerId); + if (!provider || !provider.isConfigured()) return []; + const source = new ProviderSnapshotSource(provider); + return [new ProviderReconciliationIntake(source, arr, store, { + dryRun: config.providerReconciliationDryRun, + routeFor: (category, sourceProvider) => category === "Movies" ? providerReconciliationRoutes(sourceProvider || providerId).movies : providerReconciliationRoutes(sourceProvider || providerId).shows, + onReview: async (event, error) => { + const parsed = parseMediaFilename(event.path, event.path); + recordOrganizerReview(event.path, { ...parsed, status: parsed.status === "matched" ? "ambiguous" : parsed.status, reason: `Provider reconciliation (${event.provider}): ${error.message}` }); + }, + })]; + }); + if (intakes.length === 0) { + console.warn("[provider-reconciliation] no configured providers available; worker disabled"); + return undefined; + } + return new ProviderReconciliationWorker(intakes, { + recentMs: Math.max(1000, config.providerReconciliationRecentIntervalMs), fullMs: Math.max(1000, config.providerReconciliationFullIntervalMs), recentLimit: Math.max(1, config.providerReconciliationRecentLimit), runFullOnStart: config.providerReconciliationRunFullOnStart, + }); +} + +export function startProviderReconciliation(): ProviderReconciliationWorker | undefined { + if (!config.providerReconciliationEnabled) { + console.log("[provider-reconciliation] disabled (PROVIDER_RECONCILIATION_ENABLED=false)"); + return undefined; + } + const worker = createProviderReconciliationWorker(); + if (!worker) return undefined; + worker.start(); + console.log("[provider-reconciliation] started; provider intake is read-only and provider deletion/repair is not used"); + return worker; +} diff --git a/tests/e2e/arr-bridge/qbittorrent-api.test.ts b/tests/e2e/arr-bridge/qbittorrent-api.test.ts index 4cbcd2b..f96330a 100644 --- a/tests/e2e/arr-bridge/qbittorrent-api.test.ts +++ b/tests/e2e/arr-bridge/qbittorrent-api.test.ts @@ -100,7 +100,7 @@ describe('*arr bridge qBittorrent-compatible API', () => { form.append('urls', magnet); form.append('savepath', `/downloads/${category}`); form.append('category', category); - form.append('tags', 'cinecircle-test'); + form.append('tags', 'provider-reconciliation-test'); form.append('skip_checking', 'false'); form.append('paused', 'false'); form.append('sequentialDownload', 'false'); @@ -114,7 +114,7 @@ describe('*arr bridge qBittorrent-compatible API', () => { expect(res.status).toBe(200); const torrents = await (await fetch(`${BASE_URL}/api/v2/torrents/info`)).json(); expect(torrents).toEqual(expect.arrayContaining([ - expect.objectContaining({ hash, category, tags: 'cinecircle-test' }), + expect.objectContaining({ hash, category, tags: 'provider-reconciliation-test' }), ])); }); diff --git a/tests/e2e/provider-reconciliation-three-inputs.test.ts b/tests/e2e/provider-reconciliation-three-inputs.test.ts new file mode 100644 index 0000000..3298a24 --- /dev/null +++ b/tests/e2e/provider-reconciliation-three-inputs.test.ts @@ -0,0 +1,115 @@ +import { describe, expect, test, afterEach } from 'bun:test'; +import { mkdtemp, readlink, lstat, rm } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import path from 'node:path'; +import { classifyTorrent } from '../../src/core/mediaClassifier'; +import { parseMediaFilename } from '../../src/services/mediaParser'; +import { HttpArrClient, type ArrRoute, type DirectFileEvent } from '../../src/services/providerReconciliation'; + +const originalFetch = globalThis.fetch; + +function event(category: 'Movies' | 'Shows', path: string): DirectFileEvent { + return { + provider: 'alldebrid', providerItemId: 'fixture-item', action: 'added', path, + tree: [{ path, size: 1 }], sourceCategory: category, + observedAt: '2026-09-17T00:00:00.000Z', stableDedupeKey: `fixture:${category}:${path}`, + }; +} + +function route(kind: 'radarr' | 'sonarr', sourcePathPrefix?: string): ArrRoute { + return { kind, baseUrl: `http://${kind}.fixture`, apiKey: 'fixture-api-key', sourcePathPrefix }; +} + +afterEach(() => { globalThis.fetch = originalFetch; }); + +describe('provider reconciliation fixture E2E', () => { + test('A: historical library input parses and classifies without a write boundary', () => { + const movie = parseMediaFilename('Fixture.Movie (2020).mkv', 'Movies/Fixture.Movie (2020).mkv'); + const episode = parseMediaFilename('Fixture.Show S01E02.mkv', 'Shows/Fixture.Show S01E02.mkv'); + expect(movie.kind).toBe('movie'); + expect(movie.year).toBe(2020); + expect(episode.kind).toBe('episode'); + expect(episode.season).toBe(1); + expect(episode.episode).toBe(2); + expect(classifyTorrent('Fixture.Movie (2020).mkv')).toBe('movies'); + expect(classifyTorrent('Fixture.Show S01E02.mkv')).toBe('shows'); + }); + + test('B: Seerr fixture routes movie and episode requests to Arr scan commands', async () => { + const requests: Array<{ url: string; method?: string; headers: Record; body?: any }> = []; + globalThis.fetch = (async (input: RequestInfo | URL, init?: RequestInit) => { + requests.push({ url: String(input), method: init?.method, headers: Object.fromEntries(new Headers(init?.headers).entries()), body: init?.body ? JSON.parse(String(init.body)) : undefined }); + if (init?.method === 'POST') return new Response(JSON.stringify({ id: requests.length, status: 'queued' }), { status: 201, headers: { 'content-type': 'application/json' } }); + return new Response(JSON.stringify({ id: 1, status: 'completed', result: 'successful' }), { status: 200, headers: { 'content-type': 'application/json' } }); + }) as typeof fetch; + + const client = new HttpArrClient(); + const fixtureSeerrMovie = { mediaType: 'movie', tmdbId: 'fixture-movie-id', requestedPath: 'arr-visible/movie.mkv' }; + const fixtureSeerrTv = { mediaType: 'tv', tmdbId: 'fixture-tv-id', requestedPath: 'arr-visible/show.mkv' }; + await client.submitScan(route('radarr'), event('Movies', fixtureSeerrMovie.requestedPath)); + await client.submitScan(route('sonarr'), event('Shows', fixtureSeerrTv.requestedPath)); + await client.getCommand(route('radarr'), '1'); + await client.getCommand(route('sonarr'), '2'); + + expect(requests.filter((request) => request.method === 'POST').map((request) => request.body.name)).toEqual([ + 'DownloadedMoviesScan', 'DownloadedEpisodesScan', + ]); + expect(requests.filter((request) => request.method === 'POST').map((request) => request.body.path)).toEqual([ + 'arr-visible/movie.mkv', 'arr-visible/show.mkv', + ]); + expect(requests.every((request) => request.headers['x-api-key'] === 'fixture-api-key')).toBe(true); + expect(requests.filter((request) => request.method !== 'POST').map((request) => request.url)).toEqual([ + 'http://radarr.fixture/api/v3/command/1', 'http://sonarr.fixture/api/v3/command/2', + ]); + }); + + test('B2: prefixes the provider-relative path with the Arr-visible mount path', async () => { + let body: any; + globalThis.fetch = (async (_input: RequestInfo | URL, init?: RequestInit) => { + body = JSON.parse(String(init?.body)); + return new Response(JSON.stringify({ id: 1, status: 'queued' }), { status: 201, headers: { 'content-type': 'application/json' } }); + }) as typeof fetch; + + await new HttpArrClient().submitScan(route('sonarr', '/mnt/schrodrive/alldebrid'), event('Shows', 'Shows/Lanterns.S01E06.mkv')); + expect(body.path).toBe('/mnt/schrodrive/alldebrid/Shows/Lanterns.S01E06.mkv'); + expect(body.importMode).toBe('Copy'); + }); + + test('B3: exposes provider media as a symlink and never copies it locally', async () => { + let body: any; + globalThis.fetch = (async (_input: RequestInfo | URL, init?: RequestInit) => { + if (init?.body) body = JSON.parse(String(init.body)); + if (!init?.body) return new Response(JSON.stringify([{ id: 1, title: 'Lanterns' }]), { status: 200 }); + return new Response(JSON.stringify({ id: 1, status: 'queued' }), { status: 201, headers: { 'content-type': 'application/json' } }); + }) as typeof fetch; + const root = await mkdtemp(path.join(tmpdir(), 'provider-reconciliation-symlink-')); + try { + const routeWithLibrary = { ...route('sonarr', '/mnt/schrodrive/alldebrid'), symlinkLibraryPath: root }; + const item = event('Shows', 'shows/Lanterns/Lanterns.S01E06.Bad.Optics.mkv'); + await new HttpArrClient().submitScan(routeWithLibrary, item); + const link = path.join(root, 'Lanterns', 'Season 1', 'Lanterns.S01E06.Bad.Optics.mkv'); + expect((await lstat(link)).isSymbolicLink()).toBe(true); + expect(await readlink(link)).toBe('/mnt/schrodrive/alldebrid/shows/Lanterns/Lanterns.S01E06.Bad.Optics.mkv'); + expect(body.name).toBe('RescanSeries'); + expect(body.seriesId).toBe(1); + } finally { + await rm(root, { recursive: true, force: true }); + } + }); + + test('C: direct AllDebrid fixture reaches the same Arr completion boundary', async () => { + const completed: string[] = []; + const source = { async listSnapshot() { return [{ providerItemId: 'fixture-ad', name: 'Fixture.Show S01E02', status: 'finished', files: [{ path: 'arr-visible/show.mkv', size: 1 }], observedAt: '2026-09-17T00:00:00.000Z' }]; } }; + const store = new (await import('../../src/services/providerReconciliation')).InMemoryIntakeStateStore(); + const arr = { + async submitScan(_route: ArrRoute, item: DirectFileEvent) { completed.push(`${item.sourceCategory}:${item.path}`); return { commandId: 'fixture-command', status: 'queued' }; }, + async getCommand(_route: ArrRoute, commandId: string) { return { commandId, status: 'completed', result: 'successful' }; }, + }; + const { ProviderReconciliationIntake } = await import('../../src/services/providerReconciliation'); + const intake = new ProviderReconciliationIntake(source, arr, store, { routeFor: () => route('sonarr') }); + await intake.reconcile(); + await intake.reconcile(); + expect(completed).toEqual(['Shows:arr-visible/show.mkv']); + expect(store.getItem('fixture-ad')?.terminalStatus).toBe('completed'); + }); +}); diff --git a/tests/e2e/provider-reconciliation.test.ts b/tests/e2e/provider-reconciliation.test.ts new file mode 100644 index 0000000..f4e15e3 --- /dev/null +++ b/tests/e2e/provider-reconciliation.test.ts @@ -0,0 +1,249 @@ +import { describe, expect, test } from 'bun:test'; +import { closeDb } from '../../src/core/db'; +import { + ProviderReconciliationIntake, + InMemoryIntakeStateStore, + SqliteIntakeStateStore, + type ProviderSnapshot, + type ArrClient, + type ArrCommandResult, + type DirectFileEvent, + isMediaFile, + ProviderSnapshotSource, + ProviderReconciliationWorker, +} from '../../src/services/providerReconciliation'; + +function snapshot(id: string, name: string, files = [{ path: `${name}.mkv`, size: 10 }]): ProviderSnapshot { + return { providerItemId: id, name, status: 'finished', files, observedAt: '2026-09-17T00:00:00.000Z' }; +} + +class SequenceSource { + constructor(private readonly rounds: ProviderSnapshot[][]) {} + async listSnapshot(): Promise { return this.rounds.shift() || []; } +} + +class FakeArr implements ArrClient { + submitted: Array<{ kind: string; event: DirectFileEvent }> = []; + polled: string[] = []; + async submitScan(route: any, event: DirectFileEvent): Promise { + this.submitted.push({ kind: route.kind, event }); + return { commandId: `${route.kind}-${this.submitted.length}`, status: 'queued' }; + } + async getCommand(_route: any, commandId: string): Promise { + this.polled.push(commandId); + return { commandId, status: 'completed', result: 'successful' }; + } +} + +const routeFor = (category: 'Movies' | 'Shows') => ({ + kind: category === 'Movies' ? 'radarr' as const : 'sonarr' as const, + baseUrl: `http://${category.toLowerCase()}.test`, + apiKey: 'fixture-key', +}); + +describe('provider reconciliation direct intake', () => { + test('uses the common provider contract for a non-AllDebrid provider', async () => { + const source = new ProviderSnapshotSource({ + id: 'realdebrid', + async listTorrents() { return [{ id: 'rd-1', name: 'Example.Show S01E01', status: 'seeding', progress: 100, bytes: 10, files: [{ id: 'f1', name: 'Example.Show.S01E01.mkv', path: 'Example.Show.S01E01.mkv', size: 10, selected: true }] } as any]; }, + async fetchDirectories() { return [{ id: 'rd-1', name: 'Example.Show.S01E01', originalName: 'Example.Show S01E01', files: [{ id: 'f1', name: 'Example.Show.S01E01.mkv', size: 10 }] }]; }, + }); + const snapshots = await source.listSnapshot(); + expect(snapshots[0]).toMatchObject({ provider: 'realdebrid', providerItemId: 'realdebrid:rd-1', status: 'seeding' }); + expect(snapshots[0].files).toEqual([{ path: 'Example.Show.S01E01.mkv', size: 10 }]); + }); + + test('uses the existing AllDebrid client for bounded recent and full snapshots', async () => { + const recentRequests: string[][] = []; + let fullCalls = 0; + const provider = { + async listTorrents() { + return [ + { id: 'old', name: 'Old Movie', filename: 'Old Movie', status: 'finished', progress: 100, bytes: 1, files: [], addedAt: new Date('2026-09-16T00:00:00Z') }, + { id: 'new', name: 'New Movie', filename: 'New Movie', status: 'finished', progress: 100, bytes: 1, files: [], addedAt: new Date('2026-09-17T00:00:00Z') }, + ]; + }, + async fetchDirectories() { + fullCalls++; + return [{ id: 'old', name: 'Old Movie', originalName: 'Old Movie', files: [{ id: 'old.mkv', name: 'old.mkv', size: 1 }] }, { id: 'new', name: 'New Movie', originalName: 'New Movie', files: [{ id: 'new.mkv', name: 'new.mkv', size: 1 }] }]; + }, + async fetchDirectoriesForIds(items: Array<{ id: string }>) { + recentRequests.push(items.map((item) => item.id)); + return items.map((item) => ({ id: item.id, name: item.id, originalName: item.id, files: [{ id: `${item.id}.mkv`, name: `${item.id}.mkv`, size: 1 }] })); + }, + }; + const source = new ProviderSnapshotSource(provider as any, false); + const recent = await source.listRecentSnapshot!(1); + const full = await source.listSnapshot(); + expect(recent.map((item) => item.providerItemId)).toEqual(['new']); + expect(recentRequests).toEqual([['new']]); + expect(full.map((item) => item.providerItemId)).toEqual(['old', 'new']); + expect(fullCalls).toBe(1); + + const arr = new FakeArr(); + const intake = new ProviderReconciliationIntake( + new ProviderSnapshotSource(provider as any, false), arr, new InMemoryIntakeStateStore(), { routeFor }, + ); + const events = await intake.reconcile('recent', 1); + expect(events).toHaveLength(1); + expect(events[0].provider).toBe('alldebrid'); + expect(arr.submitted[0].kind).toBe('radarr'); + }); + + test('emits add and routes Movies to Radarr and Shows to Sonarr', async () => { + const arr = new FakeArr(); + const intake = new ProviderReconciliationIntake( + new SequenceSource([[snapshot('m-1', 'Movie (2026)'), snapshot('s-1', 'Show S01E02')]]), + arr, + new InMemoryIntakeStateStore(), + { routeFor }, + ); + + const events = await intake.reconcile(); + expect(events.map((event) => event.action)).toEqual(['added', 'added']); + expect(events.map((event) => event.sourceCategory)).toEqual(['Movies', 'Shows']); + expect(arr.submitted.map((item) => item.kind)).toEqual(['radarr', 'sonarr']); + expect(arr.submitted[0].event.path).toBe('Movie (2026).mkv'); + }); + + test('emits changed when a completed file tree changes, even after a missed round', async () => { + const store = new InMemoryIntakeStateStore(); + const arr = new FakeArr(); + const intake = new ProviderReconciliationIntake( + new SequenceSource([ + [snapshot('m-1', 'Movie (2026)', [{ path: 'Movie.mkv', size: 10 }])], + [snapshot('m-1', 'Movie (2026)', [{ path: 'Movie.mkv', size: 20 }])], + ]), + arr, store, { routeFor }, + ); + await intake.reconcile(); + const second = await intake.reconcile(); + expect(second).toHaveLength(1); + expect(second[0].action).toBe('changed'); + expect(second[0].stableDedupeKey).not.toBe((await intake.reconcile())[0]?.stableDedupeKey); + }); + + test('suppresses duplicate delivery and emits deletion from a missing status item', async () => { + const store = new InMemoryIntakeStateStore(); + const arr = new FakeArr(); + const intake = new ProviderReconciliationIntake( + new SequenceSource([[snapshot('m-1', 'Movie (2026)')], [snapshot('m-1', 'Movie (2026)')], []]), + arr, store, { routeFor }, + ); + expect((await intake.reconcile()).map((e) => e.action)).toEqual(['added']); + expect(await intake.reconcile()).toEqual([]); + expect((await intake.reconcile()).map((e) => e.action)).toEqual(['deleted']); + expect(arr.submitted).toHaveLength(1); + }); + + test('dry-run persists state without submitting to Arr', async () => { + const arr = new FakeArr(); + const store = new InMemoryIntakeStateStore(); + const intake = new ProviderReconciliationIntake( + new SequenceSource([[snapshot('m-1', 'Movie (2026)')]]), arr, store, + { routeFor, dryRun: true }, + ); + const events = await intake.reconcile(); + expect(events).toHaveLength(1); + expect(arr.submitted).toHaveLength(0); + expect(store.getItem('m-1')?.lastAction).toBe('added'); + }); + + test('keeps all supported subtitles and attachments in the event tree', async () => { + const names = ['video.mkv', 'captions.srt', 'captions.ass', 'captions.ssa', 'captions.sub', 'captions.vtt', 'captions.idx', 'captions.sup', 'captions.sbv', 'captions.mpsub', 'cover.jpg']; + const arr = new FakeArr(); + const intake = new ProviderReconciliationIntake( + new SequenceSource([[snapshot('m-1', 'Movie (2026)', names.map((path) => ({ path, size: 1 })))]]) , + arr, new InMemoryIntakeStateStore(), { routeFor }, + ); + const events = await intake.reconcile(); + expect(names.filter(isMediaFile)).toHaveLength(10); + expect(events[0].tree.map((file) => file.path)).toEqual(names.slice(0, 10)); + }); + + test('restarts from persisted state and polls the pending Arr command', async () => { + const store = new InMemoryIntakeStateStore(); + const firstArr = new FakeArr(); + await new ProviderReconciliationIntake( + new SequenceSource([[snapshot('m-1', 'Movie (2026)')]]), firstArr, store, { routeFor }, + ).reconcile(); + const secondArr = new FakeArr(); + const events = await new ProviderReconciliationIntake( + new SequenceSource([[snapshot('m-1', 'Movie (2026)')]]), secondArr, store, { routeFor }, + ).reconcile(); + expect(events).toEqual([]); + expect(secondArr.polled).toEqual(['radarr-1']); + expect(store.getItem('m-1')?.terminalStatus).toBe('completed'); + }); + + test('persists the direct event and cursor across a SQLite store restart', async () => { + const itemId = `sqlite-${Date.now()}-${Math.random().toString(16).slice(2)}`; + const firstArr = new FakeArr(); + const firstStore = new SqliteIntakeStateStore(); + const source = new SequenceSource([[snapshot(itemId, 'Movie (2026)')]]); + const first = new ProviderReconciliationIntake(source, firstArr, firstStore, { routeFor }); + const events = await first.reconcile('full'); + expect(events.some((event) => event.providerItemId === itemId && event.action === 'added')).toBe(true); + expect(firstStore.getCursor().fullAt).toBeDefined(); + expect(firstStore.hasEvent(events[0].stableDedupeKey)).toBe(true); + + closeDb(); + const secondStore = new SqliteIntakeStateStore(); + const secondArr = new FakeArr(); + const second = new ProviderReconciliationIntake( + new SequenceSource([[snapshot(itemId, 'Movie (2026)')]]), secondArr, secondStore, { routeFor }, + ); + expect(await second.reconcile('full')).toEqual([]); + expect(secondArr.submitted).toHaveLength(0); + expect(secondStore.getItem(itemId)?.lastAction).toBe('added'); + closeDb(); + }); + + test('retries transient Arr submission failures and preserves one correlation', async () => { + let attempts = 0; + const arr: ArrClient = { + async submitScan() { + attempts++; + if (attempts < 3) throw new Error('temporary'); + return { commandId: 'sonarr-1', status: 'queued' }; + }, + async getCommand(_route, commandId) { return { commandId, status: 'completed' }; }, + }; + const intake = new ProviderReconciliationIntake( + new SequenceSource([[snapshot('s-1', 'Show S01E02')]]), arr, new InMemoryIntakeStateStore(), + { routeFor, maxAttempts: 3 }, + ); + await intake.reconcile(); + expect(attempts).toBe(3); + }); + + test('hands permanent Arr failure to Review after bounded retries', async () => { + const reviewed: DirectFileEvent[] = []; + const arr: ArrClient = { + async submitScan() { throw new Error('permanent'); }, + async getCommand(_route, commandId) { return { commandId, status: 'failed' }; }, + }; + const intake = new ProviderReconciliationIntake( + new SequenceSource([[snapshot('m-1', 'Movie (2026)')]]), arr, new InMemoryIntakeStateStore(), + { routeFor, maxAttempts: 2, onReview: (event) => { reviewed.push(event); } }, + ); + await expect(intake.reconcile()).rejects.toThrow('permanent'); + expect(reviewed).toHaveLength(1); + expect(reviewed[0].action).toBe('added'); + }); + + test('starts recent and full polling at configured intervals and stops cleanly', async () => { + let calls = 0; + const intake = { async reconcile() { calls++; return []; } }; + const worker = new ProviderReconciliationWorker(intake as any, { recentMs: 10, fullMs: 15, recentLimit: 2 }); + worker.start(); + expect(worker.isRunning()).toBe(true); + await new Promise((resolve) => setTimeout(resolve, 25)); + worker.stop(); + const stoppedAt = calls; + expect(stoppedAt).toBeGreaterThanOrEqual(2); + expect(worker.isRunning()).toBe(false); + await new Promise((resolve) => setTimeout(resolve, 25)); + expect(calls).toBe(stoppedAt); + }); +}); diff --git a/tests/unit/services/provider-reconciliation-runtime.test.ts b/tests/unit/services/provider-reconciliation-runtime.test.ts new file mode 100644 index 0000000..53575ef --- /dev/null +++ b/tests/unit/services/provider-reconciliation-runtime.test.ts @@ -0,0 +1,46 @@ +import { afterEach, describe, expect, test } from 'bun:test'; +import { config } from '../../../src/core/config'; +import { providerReconciliationRoutes, startProviderReconciliation } from '../../../src/services/providerReconciliationRuntime'; + +const original = { + enabled: config.providerReconciliationEnabled, + radarrUrl: config.providerReconciliationRadarrUrl, + radarrKey: config.providerReconciliationRadarrApiKey, + sonarrUrl: config.providerReconciliationSonarrUrl, + sonarrKey: config.providerReconciliationSonarrApiKey, +}; + +afterEach(() => { + config.providerReconciliationEnabled = original.enabled; + config.providerReconciliationRadarrUrl = original.radarrUrl; + config.providerReconciliationRadarrApiKey = original.radarrKey; + config.providerReconciliationSonarrUrl = original.sonarrUrl; + config.providerReconciliationSonarrApiKey = original.sonarrKey; +}); + +describe('provider reconciliation runtime wiring', () => { + test('is disabled by default and does not construct a provider worker', () => { + config.providerReconciliationEnabled = false; + expect(startProviderReconciliation()).toBeUndefined(); + }); + + test('requires explicit complete Arr routes when enabled', () => { + config.providerReconciliationEnabled = true; + config.providerReconciliationRadarrUrl = 'http://radarr.test'; + config.providerReconciliationRadarrApiKey = 'radarr-fixture-key'; + config.providerReconciliationSonarrUrl = ''; + config.providerReconciliationSonarrApiKey = ''; + expect(startProviderReconciliation()).toBeUndefined(); + }); + + test('keeps explicit Movies/Radarr and Shows/Sonarr routing for any provider', () => { + config.providerReconciliationRadarrUrl = 'http://radarr.test/'; + config.providerReconciliationRadarrApiKey = 'radarr-fixture-key'; + config.providerReconciliationSonarrUrl = 'http://sonarr.test/'; + config.providerReconciliationSonarrApiKey = 'sonarr-fixture-key'; + expect(providerReconciliationRoutes('realdebrid')).toEqual({ + movies: { kind: 'radarr', baseUrl: 'http://radarr.test/', apiKey: 'radarr-fixture-key', sourcePathPrefix: '/mnt/schrodrive/realdebrid', importMode: 'Copy' }, + shows: { kind: 'sonarr', baseUrl: 'http://sonarr.test/', apiKey: 'sonarr-fixture-key', sourcePathPrefix: '/mnt/schrodrive/realdebrid', importMode: 'Copy' }, + }); + }); +});