From 07ce297323a90324f02f787803f902af39f0f37f Mon Sep 17 00:00:00 2001 From: samtruman <39094339+samtruman@users.noreply.github.com> Date: Fri, 25 Sep 2026 11:27:33 +0000 Subject: [PATCH 1/8] feat(cinecircle): wire opt-in AllDebrid reconciliation worker --- ...inecircle-test-reconciliation.override.yml | 24 ++ src/core/config.ts | 10 + src/index.ts | 6 + src/services/cinecircleAlldebridIntake.ts | 368 ++++++++++++++++++ src/services/cinecircleAlldebridRuntime.ts | 72 ++++ .../providerReconciliationCapabilities.ts | 51 +++ tests/e2e/cinecircle-alldebrid-intake.test.ts | 237 +++++++++++ .../cinecircle-alldebrid-runtime.test.ts | 46 +++ 8 files changed, 814 insertions(+) create mode 100644 deploy/cinecircle-test-reconciliation.override.yml create mode 100644 src/services/cinecircleAlldebridIntake.ts create mode 100644 src/services/cinecircleAlldebridRuntime.ts create mode 100644 src/services/providerReconciliationCapabilities.ts create mode 100644 tests/e2e/cinecircle-alldebrid-intake.test.ts create mode 100644 tests/unit/services/cinecircle-alldebrid-runtime.test.ts diff --git a/deploy/cinecircle-test-reconciliation.override.yml b/deploy/cinecircle-test-reconciliation.override.yml new file mode 100644 index 0000000..6ca2cc5 --- /dev/null +++ b/deploy/cinecircle-test-reconciliation.override.yml @@ -0,0 +1,24 @@ +# Test-only override for the CineCircle AllDebrid reconciliation worker. +# Apply with the existing cinecircle-test compose file only. +services: + schrodrive-test: + image: schrodrive:cinecircle-alldebrid-reconciliation + environment: + CINECIRCLE_ALLDEBRID_RECONCILIATION_ENABLED: "true" + CINECIRCLE_ALLDEBRID_RECENT_INTERVAL_MS: "60000" + CINECIRCLE_ALLDEBRID_FULL_INTERVAL_MS: "300000" + CINECIRCLE_ALLDEBRID_RECENT_LIMIT: "30" + CINECIRCLE_ALLDEBRID_DRY_RUN: "false" + CINECIRCLE_RADARR_URL: http://radarr-test:7878 + CINECIRCLE_RADARR_API_KEY: ${CINECIRCLE_RADARR_API_KEY:-} + CINECIRCLE_SONARR_URL: http://sonarr-test:8989 + CINECIRCLE_SONARR_API_KEY: ${CINECIRCLE_SONARR_API_KEY:-} + # Direct AllDebrid polling is the only input under test here. + RUN_WEBHOOK: "false" + RUN_POLLER: "false" + RUN_WATCHLIST_POLLER: "false" + RUN_DEAD_SCANNER: "false" + RUN_DEAD_SCANNER_WATCH: "false" + ENABLE_REPAIR: "false" + PREEMPTIVE_REPAIR: "false" + RUN_ORGANIZER_WATCH: "false" diff --git a/src/core/config.ts b/src/core/config.ts index ab0cffd..6c45e55 100644 --- a/src/core/config.ts +++ b/src/core/config.ts @@ -59,6 +59,16 @@ 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", + // CineCircle direct AllDebrid reconciliation is opt-in and disabled by default. + cineCircleAlldebridReconciliationEnabled: String(process.env.CINECIRCLE_ALLDEBRID_RECONCILIATION_ENABLED ?? "false").toLowerCase() === "true", + cineCircleAlldebridRecentIntervalMs: Number(process.env.CINECIRCLE_ALLDEBRID_RECENT_INTERVAL_MS || 900000), + cineCircleAlldebridFullIntervalMs: Number(process.env.CINECIRCLE_ALLDEBRID_FULL_INTERVAL_MS || 21600000), + cineCircleAlldebridRecentLimit: Number(process.env.CINECIRCLE_ALLDEBRID_RECENT_LIMIT || 30), + cineCircleAlldebridDryRun: String(process.env.CINECIRCLE_ALLDEBRID_DRY_RUN ?? "false").toLowerCase() === "true", + cineCircleRadarrUrl: process.env.CINECIRCLE_RADARR_URL || "", + cineCircleRadarrApiKey: process.env.CINECIRCLE_RADARR_API_KEY || "", + cineCircleSonarrUrl: process.env.CINECIRCLE_SONARR_URL || "", + cineCircleSonarrApiKey: process.env.CINECIRCLE_SONARR_API_KEY || "", // 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..c027428 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 { startCineCircleAllDebridReconciliation } from "./services/cinecircleAlldebridRuntime"; +import type { CineCircleAllDebridReconciliationWorker } from "./services/cinecircleAlldebridIntake"; const program = new Command(); program @@ -31,6 +33,7 @@ program } // Register graceful shutdown handlers + let cineCircleReconciliationWorker: CineCircleAllDebridReconciliationWorker | 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(() => {}); + cineCircleReconciliationWorker?.stop(); setTimeout(() => { console.log(`[${new Date().toISOString()}][serve] Closing database and exiting...`); closeDb(); @@ -97,6 +101,8 @@ program console.log("[serve] Starting media server watchlist poller (RUN_WATCHLIST_POLLER=true)"); startWatchlistPoller(); } + + cineCircleReconciliationWorker = startCineCircleAllDebridReconciliation(); // Start the main server startServer(); diff --git a/src/services/cinecircleAlldebridIntake.ts b/src/services/cinecircleAlldebridIntake.ts new file mode 100644 index 0000000..a3a8794 --- /dev/null +++ b/src/services/cinecircleAlldebridIntake.ts @@ -0,0 +1,368 @@ +/** + * CineCircle fork-only AllDebrid direct-file intake. + * + * AllDebrid has no push change feed in the integration used by SchröDrive, so + * this adapter reconciles read-only magnet status plus completed file trees, + * emits stable direct-file events, and hands added/changed files to Arr. + * It is deliberately not wired into the generic provider lifecycle. + */ + +import { createHash } from 'node:crypto'; +import { getDb } from '../core/db'; +import type { AllDebridProvider } from '../providers/alldebrid'; +import type { TorrentInfo, VirtualDirectory } from '../providers'; +import { classifyTorrent } from '../core/mediaClassifier'; + +export type DirectFileAction = 'added' | 'changed' | 'deleted'; +export type SourceCategory = 'Movies' | 'Shows'; +export type ArrKind = 'radarr' | 'sonarr'; + +export interface AllDebridSnapshot { + providerItemId: string; + name: string; + status: string; + files: Array<{ path: string; size: number }>; + observedAt: string; +} + +export interface DirectFileEvent { + provider: 'alldebrid'; + 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; +} + +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 { + 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 fork state; creates only its own table in the configured test DB. */ +export class SqliteIntakeStateStore implements IntakeStateStore { + constructor() { + getDb().exec(`CREATE TABLE IF NOT EXISTS cinecircle_alldebrid_intake ( + provider_item_id TEXT PRIMARY KEY, + state_json TEXT NOT NULL, + updated_at INTEGER NOT NULL + )`); + getDb().exec(`CREATE TABLE IF NOT EXISTS cinecircle_alldebrid_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 cinecircle_alldebrid_cursor ( + name TEXT PRIMARY KEY, + observed_at TEXT NOT NULL + )`); + } + getItem(id: string): IntakeState | undefined { + const row = getDb().prepare('SELECT state_json FROM cinecircle_alldebrid_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 cinecircle_alldebrid_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 cinecircle_alldebrid_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 cinecircle_alldebrid_events WHERE dedupe_key = ?').get(key); + } + saveEvent(event: DirectFileEvent, arr?: ArrCommandResult): void { + getDb().prepare(`INSERT OR IGNORE INTO cinecircle_alldebrid_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 cinecircle_alldebrid_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 cinecircle_alldebrid_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 { + const commandName = route.kind === 'radarr' ? 'DownloadedMoviesScan' : 'DownloadedEpisodesScan'; + 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({ name: commandName, path: event.path, importMode: 'Move' }), + }); + 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 }; + } +} + +export interface AllDebridReadOnlySource { + listSnapshot(): Promise; + listRecentSnapshot?(limit: number): Promise; +} + +/** Uses only the existing provider's status and completed-directory methods. */ +export class AllDebridProviderSource implements AllDebridReadOnlySource { + constructor(private readonly provider: Pick) {} + + 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 = await this.provider.fetchDirectoriesForIds(torrents); + return this.toSnapshots(torrents, new Map(directories.map((directory) => [String(directory.id), directory])), observedAt); + } + + private toSnapshots(torrents: TorrentInfo[], trees: Map, observedAt: string): AllDebridSnapshot[] { + return torrents.map((torrent) => { + const directory = trees.get(String(torrent.id)); + return { providerItemId: String(torrent.id), name: torrent.name, status: torrent.status, + files: (directory?.files || []).map((file) => ({ path: file.name, size: file.size })), observedAt }; + }); + } +} + +export interface IntakeOptions { + dryRun?: boolean; + maxAttempts?: number; + routeFor: (category: SourceCategory) => 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: AllDebridSnapshot): SourceCategory { + return classifyTorrent(snapshot.name, snapshot.files.map((file) => file.path)) === 'shows' ? 'Shows' : 'Movies'; +} + +function fingerprint(snapshot: Pick): string { + return createHash('sha256').update(JSON.stringify({ id: snapshot.providerItemId, status: snapshot.status, files: snapshot.files })).digest('hex'); +} + +function eventKey(id: string, action: DirectFileAction, fp: string): string { + return `alldebrid:${id}:${action}:${fp}`; +} + +export class CineCircleAllDebridIntake { + constructor( + private readonly source: AllDebridReadOnlySource, + 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 (item.status !== 'finished' || 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: 'alldebrid', providerItemId: previous.providerItemId, action: 'deleted', + path: previous.path, tree: [], sourceCategory: previous.sourceCategory, + observedAt: new Date().toISOString(), stableDedupeKey: eventKey(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.commandId); + this.store.saveItem({ ...item, terminalStatus: command.status, updatedAt: new Date().toISOString() }); + if (command.status === 'failed') { + await this.options.onReview?.({ + provider: 'alldebrid', providerItemId: item.providerItemId, action: item.lastAction, + path: item.path, tree: item.tree || [], sourceCategory: item.sourceCategory, + observedAt: new Date().toISOString(), + stableDedupeKey: eventKey(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: AllDebridSnapshot, action: DirectFileAction, fp: string): DirectFileEvent { + const category = categoryFor(item); + return { + provider: 'alldebrid', providerItemId: item.providerItemId, action, + path: item.files[0].path, tree: item.files, sourceCategory: category, + observedAt: item.observedAt, stableDedupeKey: eventKey(item.providerItemId, action, fp), + }; + } + + private async dispatch(event: DirectFileEvent, item: AllDebridSnapshot, 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({ 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); + 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({ 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; CineCircle wiring must explicitly opt in. + */ +export class CineCircleAllDebridReconciliationWorker { + private recentTimer: ReturnType | undefined; + private fullTimer: ReturnType | undefined; + + constructor( + private readonly intake: CineCircleAllDebridIntake, + private readonly intervals: { recentMs: number; fullMs: number; recentLimit?: number }, + ) {} + + runRecent(): Promise { return this.intake.reconcile('recent', this.intervals.recentLimit || 30); } + runFull(): Promise { return this.intake.reconcile('full'); } + + start(): void { + if (this.recentTimer || this.fullTimer) return; + this.runRecent().catch(() => undefined); + 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/cinecircleAlldebridRuntime.ts b/src/services/cinecircleAlldebridRuntime.ts new file mode 100644 index 0000000..b680814 --- /dev/null +++ b/src/services/cinecircleAlldebridRuntime.ts @@ -0,0 +1,72 @@ +import { config } from "../core/config"; +import { registry } from "../providers"; +import { AllDebridProvider } from "../providers/alldebrid"; +import { parseMediaFilename } from "./mediaParser"; +import { recordOrganizerReview } from "./organizerReview"; +import { + AllDebridProviderSource, + CineCircleAllDebridIntake, + CineCircleAllDebridReconciliationWorker, + HttpArrClient, + SqliteIntakeStateStore, + type ArrRoute, +} from "./cinecircleAlldebridIntake"; + +export function cineCircleReconciliationRoutes(): { movies: ArrRoute; shows: ArrRoute } { + return { + movies: { kind: "radarr", baseUrl: config.cineCircleRadarrUrl, apiKey: config.cineCircleRadarrApiKey }, + shows: { kind: "sonarr", baseUrl: config.cineCircleSonarrUrl, apiKey: config.cineCircleSonarrApiKey }, + }; +} + +/** Builds the opt-in direct AllDebrid worker without changing provider lifecycle services. */ +export function createCineCircleAllDebridReconciliationWorker(): CineCircleAllDebridReconciliationWorker | undefined { + const provider = registry.get("alldebrid"); + if (!(provider instanceof AllDebridProvider)) { + console.warn("[cinecircle-reconciliation] AllDebrid provider is not registered; worker disabled"); + return undefined; + } + if (!provider.isConfigured()) { + console.warn("[cinecircle-reconciliation] AllDebrid credentials are not configured; worker disabled"); + return undefined; + } + + const routes = cineCircleReconciliationRoutes(); + if (!routes.movies.baseUrl || !routes.movies.apiKey || !routes.shows.baseUrl || !routes.shows.apiKey) { + console.warn("[cinecircle-reconciliation] Radarr/Sonarr routes are incomplete; worker disabled"); + return undefined; + } + + const source = new AllDebridProviderSource(provider); + const arr = new HttpArrClient(); + const store = new SqliteIntakeStateStore(); + const intake = new CineCircleAllDebridIntake(source, arr, store, { + dryRun: config.cineCircleAlldebridDryRun, + routeFor: (category) => category === "Movies" ? routes.movies : routes.shows, + onReview: async (event, error) => { + const parsed = parseMediaFilename(event.path, event.path); + recordOrganizerReview(event.path, { + ...parsed, + status: parsed.status === "matched" ? "ambiguous" : parsed.status, + reason: `AllDebrid direct intake: ${error.message}`, + }); + }, + }); + return new CineCircleAllDebridReconciliationWorker(intake, { + recentMs: Math.max(1000, config.cineCircleAlldebridRecentIntervalMs), + fullMs: Math.max(1000, config.cineCircleAlldebridFullIntervalMs), + recentLimit: Math.max(1, config.cineCircleAlldebridRecentLimit), + }); +} + +export function startCineCircleAllDebridReconciliation(): CineCircleAllDebridReconciliationWorker | undefined { + if (!config.cineCircleAlldebridReconciliationEnabled) { + console.log("[cinecircle-reconciliation] disabled (CINECIRCLE_ALLDEBRID_RECONCILIATION_ENABLED=false)"); + return undefined; + } + const worker = createCineCircleAllDebridReconciliationWorker(); + if (!worker) return undefined; + worker.start(); + console.log("[cinecircle-reconciliation] started; direct AllDebrid intake is read-only and provider deletion/repair is not used"); + return worker; +} 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/tests/e2e/cinecircle-alldebrid-intake.test.ts b/tests/e2e/cinecircle-alldebrid-intake.test.ts new file mode 100644 index 0000000..43549fd --- /dev/null +++ b/tests/e2e/cinecircle-alldebrid-intake.test.ts @@ -0,0 +1,237 @@ +import { describe, expect, test } from 'bun:test'; +import { closeDb } from '../../src/core/db'; +import { + CineCircleAllDebridIntake, + InMemoryIntakeStateStore, + SqliteIntakeStateStore, + type AllDebridSnapshot, + type ArrClient, + type ArrCommandResult, + type DirectFileEvent, + isMediaFile, + AllDebridProviderSource, + CineCircleAllDebridReconciliationWorker, +} from '../../src/services/cinecircleAlldebridIntake'; + +function snapshot(id: string, name: string, files = [{ path: `${name}.mkv`, size: 10 }]): AllDebridSnapshot { + return { providerItemId: id, name, status: 'finished', files, observedAt: '2026-09-17T00:00:00.000Z' }; +} + +class SequenceSource { + constructor(private readonly rounds: AllDebridSnapshot[][]) {} + 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('CineCircle AllDebrid direct intake', () => { + 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 AllDebridProviderSource(provider as any); + 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 CineCircleAllDebridIntake( + new AllDebridProviderSource(provider as any), 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 CineCircleAllDebridIntake( + 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 CineCircleAllDebridIntake( + 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 CineCircleAllDebridIntake( + 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 CineCircleAllDebridIntake( + 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 CineCircleAllDebridIntake( + 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 CineCircleAllDebridIntake( + new SequenceSource([[snapshot('m-1', 'Movie (2026)')]]), firstArr, store, { routeFor }, + ).reconcile(); + const secondArr = new FakeArr(); + const events = await new CineCircleAllDebridIntake( + 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 firstArr = new FakeArr(); + const firstStore = new SqliteIntakeStateStore(); + const source = new SequenceSource([[snapshot('sqlite-1', 'Movie (2026)')]]); + const first = new CineCircleAllDebridIntake(source, firstArr, firstStore, { routeFor }); + const events = await first.reconcile('full'); + expect(events).toHaveLength(1); + 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 CineCircleAllDebridIntake( + new SequenceSource([[snapshot('sqlite-1', 'Movie (2026)')]]), secondArr, secondStore, { routeFor }, + ); + expect(await second.reconcile('full')).toEqual([]); + expect(secondArr.submitted).toHaveLength(0); + expect(secondStore.getItem('sqlite-1')?.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 CineCircleAllDebridIntake( + 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 CineCircleAllDebridIntake( + 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 CineCircleAllDebridReconciliationWorker(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/cinecircle-alldebrid-runtime.test.ts b/tests/unit/services/cinecircle-alldebrid-runtime.test.ts new file mode 100644 index 0000000..3c48605 --- /dev/null +++ b/tests/unit/services/cinecircle-alldebrid-runtime.test.ts @@ -0,0 +1,46 @@ +import { afterEach, describe, expect, test } from 'bun:test'; +import { config } from '../../../src/core/config'; +import { cineCircleReconciliationRoutes, startCineCircleAllDebridReconciliation } from '../../../src/services/cinecircleAlldebridRuntime'; + +const original = { + enabled: config.cineCircleAlldebridReconciliationEnabled, + radarrUrl: config.cineCircleRadarrUrl, + radarrKey: config.cineCircleRadarrApiKey, + sonarrUrl: config.cineCircleSonarrUrl, + sonarrKey: config.cineCircleSonarrApiKey, +}; + +afterEach(() => { + config.cineCircleAlldebridReconciliationEnabled = original.enabled; + config.cineCircleRadarrUrl = original.radarrUrl; + config.cineCircleRadarrApiKey = original.radarrKey; + config.cineCircleSonarrUrl = original.sonarrUrl; + config.cineCircleSonarrApiKey = original.sonarrKey; +}); + +describe('CineCircle AllDebrid runtime wiring', () => { + test('is disabled by default and does not construct a provider worker', () => { + config.cineCircleAlldebridReconciliationEnabled = false; + expect(startCineCircleAllDebridReconciliation()).toBeUndefined(); + }); + + test('requires explicit complete Arr routes when enabled', () => { + config.cineCircleAlldebridReconciliationEnabled = true; + config.cineCircleRadarrUrl = 'http://radarr.test'; + config.cineCircleRadarrApiKey = 'radarr-fixture-key'; + config.cineCircleSonarrUrl = ''; + config.cineCircleSonarrApiKey = ''; + expect(startCineCircleAllDebridReconciliation()).toBeUndefined(); + }); + + test('keeps explicit Movies/Radarr and Shows/Sonarr routing', () => { + config.cineCircleRadarrUrl = 'http://radarr.test/'; + config.cineCircleRadarrApiKey = 'radarr-fixture-key'; + config.cineCircleSonarrUrl = 'http://sonarr.test/'; + config.cineCircleSonarrApiKey = 'sonarr-fixture-key'; + expect(cineCircleReconciliationRoutes()).toEqual({ + movies: { kind: 'radarr', baseUrl: 'http://radarr.test/', apiKey: 'radarr-fixture-key' }, + shows: { kind: 'sonarr', baseUrl: 'http://sonarr.test/', apiKey: 'sonarr-fixture-key' }, + }); + }); +}); From d5d08d550827c5a350db95cdec96eed05c0a951f Mon Sep 17 00:00:00 2001 From: samtruman <39094339+samtruman@users.noreply.github.com> Date: Fri, 25 Sep 2026 14:44:54 +0000 Subject: [PATCH 2/8] test(cinecircle): isolate Arr test mount paths --- ...inecircle-test-reconciliation.override.yml | 42 +++++++++++++++++++ 1 file changed, 42 insertions(+) diff --git a/deploy/cinecircle-test-reconciliation.override.yml b/deploy/cinecircle-test-reconciliation.override.yml index 6ca2cc5..45aaf4a 100644 --- a/deploy/cinecircle-test-reconciliation.override.yml +++ b/deploy/cinecircle-test-reconciliation.override.yml @@ -22,3 +22,45 @@ services: ENABLE_REPAIR: "false" PREEMPTIVE_REPAIR: "false" RUN_ORGANIZER_WATCH: "false" + volumes: + - type: bind + source: /home/samtruman/docker/cinecircle-test/schrodrive + target: /mnt/schrodrive + bind: + propagation: rshared + + radarr-test: + volumes: + - type: bind + source: /home/samtruman/docker/cinecircle-test/radarr/config + target: /config + - type: bind + source: /home/samtruman/docker/cinecircle-test + target: /mnt/cinecircle-test + - type: bind + source: /home/samtruman/docker/cinecircle-test/schrodrive + target: /mnt/schrodrive + read_only: true + bind: + propagation: rslave + - type: bind + source: /home/samtruman/docker/cinecircle-test/schrodrive/downloads + target: /mnt/schrodrive/downloads + + sonarr-test: + volumes: + - type: bind + source: /home/samtruman/docker/cinecircle-test/sonarr/config + target: /config + - type: bind + source: /home/samtruman/docker/cinecircle-test + target: /mnt/cinecircle-test + - type: bind + source: /home/samtruman/docker/cinecircle-test/schrodrive + target: /mnt/schrodrive + read_only: true + bind: + propagation: rslave + - type: bind + source: /home/samtruman/docker/cinecircle-test/schrodrive/downloads + target: /mnt/schrodrive/downloads From 1d1d136a1883fba9c3f95aae364b65d40265e580 Mon Sep 17 00:00:00 2001 From: samtruman <39094339+samtruman@users.noreply.github.com> Date: Fri, 25 Sep 2026 14:56:31 +0000 Subject: [PATCH 3/8] test(cinecircle): pass AllDebrid credential to isolated runtime --- deploy/cinecircle-test-reconciliation.override.yml | 1 + 1 file changed, 1 insertion(+) diff --git a/deploy/cinecircle-test-reconciliation.override.yml b/deploy/cinecircle-test-reconciliation.override.yml index 45aaf4a..0884e72 100644 --- a/deploy/cinecircle-test-reconciliation.override.yml +++ b/deploy/cinecircle-test-reconciliation.override.yml @@ -13,6 +13,7 @@ services: CINECIRCLE_RADARR_API_KEY: ${CINECIRCLE_RADARR_API_KEY:-} CINECIRCLE_SONARR_URL: http://sonarr-test:8989 CINECIRCLE_SONARR_API_KEY: ${CINECIRCLE_SONARR_API_KEY:-} + ALLDEBRID_API_KEY: ${ALLDEBRID_API_KEY:-} # Direct AllDebrid polling is the only input under test here. RUN_WEBHOOK: "false" RUN_POLLER: "false" From 56cb68a2da6efb8e4ef479f77f4ff6dc5d6c2ea3 Mon Sep 17 00:00:00 2001 From: samtruman <39094339+samtruman@users.noreply.github.com> Date: Fri, 25 Sep 2026 15:43:46 +0000 Subject: [PATCH 4/8] fix(cinecircle): use ready mount paths and copy imports --- ...inecircle-test-reconciliation.override.yml | 12 ++- src/core/config.ts | 2 + src/index.ts | 8 +- src/services/cinecircleAlldebridIntake.ts | 19 +++- src/services/cinecircleAlldebridRuntime.ts | 4 +- tests/e2e/cinecircle-alldebrid-intake.test.ts | 9 +- tests/e2e/cinecircle-three-inputs.test.ts | 90 +++++++++++++++++++ .../cinecircle-alldebrid-runtime.test.ts | 8 +- 8 files changed, 138 insertions(+), 14 deletions(-) create mode 100644 tests/e2e/cinecircle-three-inputs.test.ts diff --git a/deploy/cinecircle-test-reconciliation.override.yml b/deploy/cinecircle-test-reconciliation.override.yml index 0884e72..1828f95 100644 --- a/deploy/cinecircle-test-reconciliation.override.yml +++ b/deploy/cinecircle-test-reconciliation.override.yml @@ -13,6 +13,8 @@ services: CINECIRCLE_RADARR_API_KEY: ${CINECIRCLE_RADARR_API_KEY:-} CINECIRCLE_SONARR_URL: http://sonarr-test:8989 CINECIRCLE_SONARR_API_KEY: ${CINECIRCLE_SONARR_API_KEY:-} + CINECIRCLE_ALLDEBRID_ARR_PATH: /mnt/schrodrive/alldebrid + CINECIRCLE_ALLDEBRID_ARR_IMPORT_MODE: Copy ALLDEBRID_API_KEY: ${ALLDEBRID_API_KEY:-} # Direct AllDebrid polling is the only input under test here. RUN_WEBHOOK: "false" @@ -29,6 +31,12 @@ services: target: /mnt/schrodrive bind: propagation: rshared + - type: bind + source: /home/samtruman/cinecircle-test-clean/schrodrive/data + target: /data + - type: bind + source: /home/samtruman/cinecircle-test-clean/schrodrive/config + target: /config radarr-test: volumes: @@ -51,10 +59,10 @@ services: sonarr-test: volumes: - type: bind - source: /home/samtruman/docker/cinecircle-test/sonarr/config + source: /home/samtruman/cinecircle-test-clean/sonarr/config target: /config - type: bind - source: /home/samtruman/docker/cinecircle-test + source: /home/samtruman/cinecircle-test-clean target: /mnt/cinecircle-test - type: bind source: /home/samtruman/docker/cinecircle-test/schrodrive diff --git a/src/core/config.ts b/src/core/config.ts index 6c45e55..cf130c4 100644 --- a/src/core/config.ts +++ b/src/core/config.ts @@ -69,6 +69,8 @@ export const config = { cineCircleRadarrApiKey: process.env.CINECIRCLE_RADARR_API_KEY || "", cineCircleSonarrUrl: process.env.CINECIRCLE_SONARR_URL || "", cineCircleSonarrApiKey: process.env.CINECIRCLE_SONARR_API_KEY || "", + cineCircleAlldebridArrPath: process.env.CINECIRCLE_ALLDEBRID_ARR_PATH || "/mnt/schrodrive/alldebrid", + cineCircleAlldebridArrImportMode: process.env.CINECIRCLE_ALLDEBRID_ARR_IMPORT_MODE || "Copy", // 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 c027428..2a3b0cc 100644 --- a/src/index.ts +++ b/src/index.ts @@ -81,7 +81,13 @@ program }); } - promises.push(mountVirtualDrive()); + // The AllDebrid reconciliation worker submits Arr scans 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) { diff --git a/src/services/cinecircleAlldebridIntake.ts b/src/services/cinecircleAlldebridIntake.ts index a3a8794..1dc1916 100644 --- a/src/services/cinecircleAlldebridIntake.ts +++ b/src/services/cinecircleAlldebridIntake.ts @@ -20,6 +20,7 @@ export type ArrKind = 'radarr' | 'sonarr'; export interface AllDebridSnapshot { providerItemId: string; name: string; + directoryName?: string; status: string; files: Array<{ path: string; size: number }>; observedAt: string; @@ -40,6 +41,10 @@ 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'; } export interface ArrCommandResult { @@ -141,10 +146,13 @@ export class SqliteIntakeStateStore implements IntakeStateStore { export class HttpArrClient implements ArrClient { async submitScan(route: ArrRoute, event: DirectFileEvent): Promise { const commandName = route.kind === 'radarr' ? 'DownloadedMoviesScan' : 'DownloadedEpisodesScan'; + const path = route.sourcePathPrefix + ? `${route.sourcePathPrefix.replace(/\/$/, '')}/${event.path.replace(/^\/+/, '')}` + : event.path; 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({ name: commandName, path: event.path, importMode: 'Move' }), + body: JSON.stringify({ name: commandName, path, importMode: route.importMode || 'Copy' }), }); 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 }; @@ -191,7 +199,7 @@ export class AllDebridProviderSource implements AllDebridReadOnlySource { private toSnapshots(torrents: TorrentInfo[], trees: Map, observedAt: string): AllDebridSnapshot[] { return torrents.map((torrent) => { const directory = trees.get(String(torrent.id)); - return { providerItemId: String(torrent.id), name: torrent.name, status: torrent.status, + return { providerItemId: String(torrent.id), name: torrent.name, directoryName: directory?.name, status: torrent.status, files: (directory?.files || []).map((file) => ({ path: file.name, size: file.size })), observedAt }; }); } @@ -298,9 +306,14 @@ export class CineCircleAllDebridIntake { private makeEvent(item: AllDebridSnapshot, action: DirectFileAction, fp: string): DirectFileEvent { 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: 'alldebrid', providerItemId: item.providerItemId, action, - path: item.files[0].path, tree: item.files, sourceCategory: category, + path, tree: item.files, sourceCategory: category, observedAt: item.observedAt, stableDedupeKey: eventKey(item.providerItemId, action, fp), }; } diff --git a/src/services/cinecircleAlldebridRuntime.ts b/src/services/cinecircleAlldebridRuntime.ts index b680814..51088ad 100644 --- a/src/services/cinecircleAlldebridRuntime.ts +++ b/src/services/cinecircleAlldebridRuntime.ts @@ -14,8 +14,8 @@ import { export function cineCircleReconciliationRoutes(): { movies: ArrRoute; shows: ArrRoute } { return { - movies: { kind: "radarr", baseUrl: config.cineCircleRadarrUrl, apiKey: config.cineCircleRadarrApiKey }, - shows: { kind: "sonarr", baseUrl: config.cineCircleSonarrUrl, apiKey: config.cineCircleSonarrApiKey }, + movies: { kind: "radarr", baseUrl: config.cineCircleRadarrUrl, apiKey: config.cineCircleRadarrApiKey, sourcePathPrefix: config.cineCircleAlldebridArrPath, importMode: config.cineCircleAlldebridArrImportMode === "Move" ? "Move" : "Copy" }, + shows: { kind: "sonarr", baseUrl: config.cineCircleSonarrUrl, apiKey: config.cineCircleSonarrApiKey, sourcePathPrefix: config.cineCircleAlldebridArrPath, importMode: config.cineCircleAlldebridArrImportMode === "Move" ? "Move" : "Copy" }, }; } diff --git a/tests/e2e/cinecircle-alldebrid-intake.test.ts b/tests/e2e/cinecircle-alldebrid-intake.test.ts index 43549fd..6f1be94 100644 --- a/tests/e2e/cinecircle-alldebrid-intake.test.ts +++ b/tests/e2e/cinecircle-alldebrid-intake.test.ts @@ -166,12 +166,13 @@ describe('CineCircle AllDebrid direct intake', () => { }); 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('sqlite-1', 'Movie (2026)')]]); + const source = new SequenceSource([[snapshot(itemId, 'Movie (2026)')]]); const first = new CineCircleAllDebridIntake(source, firstArr, firstStore, { routeFor }); const events = await first.reconcile('full'); - expect(events).toHaveLength(1); + 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); @@ -179,11 +180,11 @@ describe('CineCircle AllDebrid direct intake', () => { const secondStore = new SqliteIntakeStateStore(); const secondArr = new FakeArr(); const second = new CineCircleAllDebridIntake( - new SequenceSource([[snapshot('sqlite-1', 'Movie (2026)')]]), secondArr, secondStore, { routeFor }, + new SequenceSource([[snapshot(itemId, 'Movie (2026)')]]), secondArr, secondStore, { routeFor }, ); expect(await second.reconcile('full')).toEqual([]); expect(secondArr.submitted).toHaveLength(0); - expect(secondStore.getItem('sqlite-1')?.lastAction).toBe('added'); + expect(secondStore.getItem(itemId)?.lastAction).toBe('added'); closeDb(); }); diff --git a/tests/e2e/cinecircle-three-inputs.test.ts b/tests/e2e/cinecircle-three-inputs.test.ts new file mode 100644 index 0000000..59d8121 --- /dev/null +++ b/tests/e2e/cinecircle-three-inputs.test.ts @@ -0,0 +1,90 @@ +import { describe, expect, test, afterEach } from 'bun:test'; +import { classifyTorrent } from '../../src/core/mediaClassifier'; +import { parseMediaFilename } from '../../src/services/mediaParser'; +import { HttpArrClient, type ArrRoute, type DirectFileEvent } from '../../src/services/cinecircleAlldebridIntake'; + +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('CineCircle three-input 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('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/cinecircleAlldebridIntake')).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 { CineCircleAllDebridIntake } = await import('../../src/services/cinecircleAlldebridIntake'); + const intake = new CineCircleAllDebridIntake(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/unit/services/cinecircle-alldebrid-runtime.test.ts b/tests/unit/services/cinecircle-alldebrid-runtime.test.ts index 3c48605..8b38c53 100644 --- a/tests/unit/services/cinecircle-alldebrid-runtime.test.ts +++ b/tests/unit/services/cinecircle-alldebrid-runtime.test.ts @@ -8,6 +8,8 @@ const original = { radarrKey: config.cineCircleRadarrApiKey, sonarrUrl: config.cineCircleSonarrUrl, sonarrKey: config.cineCircleSonarrApiKey, + alldebridArrPath: config.cineCircleAlldebridArrPath, + alldebridArrImportMode: config.cineCircleAlldebridArrImportMode, }; afterEach(() => { @@ -16,6 +18,8 @@ afterEach(() => { config.cineCircleRadarrApiKey = original.radarrKey; config.cineCircleSonarrUrl = original.sonarrUrl; config.cineCircleSonarrApiKey = original.sonarrKey; + config.cineCircleAlldebridArrPath = original.alldebridArrPath; + config.cineCircleAlldebridArrImportMode = original.alldebridArrImportMode; }); describe('CineCircle AllDebrid runtime wiring', () => { @@ -39,8 +43,8 @@ describe('CineCircle AllDebrid runtime wiring', () => { config.cineCircleSonarrUrl = 'http://sonarr.test/'; config.cineCircleSonarrApiKey = 'sonarr-fixture-key'; expect(cineCircleReconciliationRoutes()).toEqual({ - movies: { kind: 'radarr', baseUrl: 'http://radarr.test/', apiKey: 'radarr-fixture-key' }, - shows: { kind: 'sonarr', baseUrl: 'http://sonarr.test/', apiKey: 'sonarr-fixture-key' }, + movies: { kind: 'radarr', baseUrl: 'http://radarr.test/', apiKey: 'radarr-fixture-key', sourcePathPrefix: '/mnt/schrodrive/alldebrid', importMode: 'Copy' }, + shows: { kind: 'sonarr', baseUrl: 'http://sonarr.test/', apiKey: 'sonarr-fixture-key', sourcePathPrefix: '/mnt/schrodrive/alldebrid', importMode: 'Copy' }, }); }); }); From e4bf9945176110a594ffb18fff285eb75995df79 Mon Sep 17 00:00:00 2001 From: samtruman <39094339+samtruman@users.noreply.github.com> Date: Fri, 25 Sep 2026 16:14:10 +0000 Subject: [PATCH 5/8] fix(cinecircle): expose AllDebrid media through symlink rescans --- ...inecircle-test-reconciliation.override.yml | 10 ++- src/core/config.ts | 3 + src/services/cinecircleAlldebridIntake.ts | 77 +++++++++++++++++-- src/services/cinecircleAlldebridRuntime.ts | 5 +- tests/e2e/cinecircle-three-inputs.test.ts | 25 ++++++ 5 files changed, 111 insertions(+), 9 deletions(-) diff --git a/deploy/cinecircle-test-reconciliation.override.yml b/deploy/cinecircle-test-reconciliation.override.yml index 1828f95..9704003 100644 --- a/deploy/cinecircle-test-reconciliation.override.yml +++ b/deploy/cinecircle-test-reconciliation.override.yml @@ -6,8 +6,9 @@ services: environment: CINECIRCLE_ALLDEBRID_RECONCILIATION_ENABLED: "true" CINECIRCLE_ALLDEBRID_RECENT_INTERVAL_MS: "60000" - CINECIRCLE_ALLDEBRID_FULL_INTERVAL_MS: "300000" - CINECIRCLE_ALLDEBRID_RECENT_LIMIT: "30" + CINECIRCLE_ALLDEBRID_FULL_INTERVAL_MS: "86400000" + CINECIRCLE_ALLDEBRID_RECENT_LIMIT: "1" + CINECIRCLE_ALLDEBRID_RUN_FULL_ON_START: "false" CINECIRCLE_ALLDEBRID_DRY_RUN: "false" CINECIRCLE_RADARR_URL: http://radarr-test:7878 CINECIRCLE_RADARR_API_KEY: ${CINECIRCLE_RADARR_API_KEY:-} @@ -15,6 +16,8 @@ services: CINECIRCLE_SONARR_API_KEY: ${CINECIRCLE_SONARR_API_KEY:-} CINECIRCLE_ALLDEBRID_ARR_PATH: /mnt/schrodrive/alldebrid CINECIRCLE_ALLDEBRID_ARR_IMPORT_MODE: Copy + CINECIRCLE_ALLDEBRID_MOVIES_LIBRARY_PATH: /mnt/cinecircle-test/radarr/movies + CINECIRCLE_ALLDEBRID_SHOWS_LIBRARY_PATH: /mnt/cinecircle-test/sonarr/shows ALLDEBRID_API_KEY: ${ALLDEBRID_API_KEY:-} # Direct AllDebrid polling is the only input under test here. RUN_WEBHOOK: "false" @@ -37,6 +40,9 @@ services: - type: bind source: /home/samtruman/cinecircle-test-clean/schrodrive/config target: /config + - type: bind + source: /home/samtruman/cinecircle-test-clean + target: /mnt/cinecircle-test radarr-test: volumes: diff --git a/src/core/config.ts b/src/core/config.ts index cf130c4..dfd05dc 100644 --- a/src/core/config.ts +++ b/src/core/config.ts @@ -64,6 +64,7 @@ export const config = { cineCircleAlldebridRecentIntervalMs: Number(process.env.CINECIRCLE_ALLDEBRID_RECENT_INTERVAL_MS || 900000), cineCircleAlldebridFullIntervalMs: Number(process.env.CINECIRCLE_ALLDEBRID_FULL_INTERVAL_MS || 21600000), cineCircleAlldebridRecentLimit: Number(process.env.CINECIRCLE_ALLDEBRID_RECENT_LIMIT || 30), + cineCircleAlldebridRunFullOnStart: String(process.env.CINECIRCLE_ALLDEBRID_RUN_FULL_ON_START ?? "true").toLowerCase() !== "false", cineCircleAlldebridDryRun: String(process.env.CINECIRCLE_ALLDEBRID_DRY_RUN ?? "false").toLowerCase() === "true", cineCircleRadarrUrl: process.env.CINECIRCLE_RADARR_URL || "", cineCircleRadarrApiKey: process.env.CINECIRCLE_RADARR_API_KEY || "", @@ -71,6 +72,8 @@ export const config = { cineCircleSonarrApiKey: process.env.CINECIRCLE_SONARR_API_KEY || "", cineCircleAlldebridArrPath: process.env.CINECIRCLE_ALLDEBRID_ARR_PATH || "/mnt/schrodrive/alldebrid", cineCircleAlldebridArrImportMode: process.env.CINECIRCLE_ALLDEBRID_ARR_IMPORT_MODE || "Copy", + cineCircleAlldebridMoviesLibraryPath: process.env.CINECIRCLE_ALLDEBRID_MOVIES_LIBRARY_PATH || "", + cineCircleAlldebridShowsLibraryPath: process.env.CINECIRCLE_ALLDEBRID_SHOWS_LIBRARY_PATH || "", // AllDebrid WebDAV (if supported) alldebridWebdavUrl: process.env.ALLDEBRID_WEBDAV_URL || "", alldebridWebdavUsername: process.env.ALLDEBRID_WEBDAV_USERNAME || "", diff --git a/src/services/cinecircleAlldebridIntake.ts b/src/services/cinecircleAlldebridIntake.ts index 1dc1916..73f1a69 100644 --- a/src/services/cinecircleAlldebridIntake.ts +++ b/src/services/cinecircleAlldebridIntake.ts @@ -8,10 +8,13 @@ */ import { createHash } from 'node:crypto'; +import { promises as fsp } from 'node:fs'; +import path from 'node:path'; import { getDb } from '../core/db'; import type { AllDebridProvider } from '../providers/alldebrid'; import type { TorrentInfo, VirtualDirectory } from '../providers'; import { classifyTorrent } from '../core/mediaClassifier'; +import { parseMediaFilename } from './mediaParser'; export type DirectFileAction = 'added' | 'changed' | 'deleted'; export type SourceCategory = 'Movies' | 'Shows'; @@ -45,6 +48,8 @@ export interface ArrRoute { 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 { @@ -145,14 +150,23 @@ export class SqliteIntakeStateStore implements IntakeStateStore { export class HttpArrClient implements ArrClient { async submitScan(route: ArrRoute, event: DirectFileEvent): Promise { - const commandName = route.kind === 'radarr' ? 'DownloadedMoviesScan' : 'DownloadedEpisodesScan'; - const path = route.sourcePathPrefix + 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({ name: commandName, path, importMode: route.importMode || 'Copy' }), + 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 }; @@ -170,6 +184,59 @@ export class HttpArrClient implements ArrClient { } } +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; + await fsp.unlink(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 AllDebridReadOnlySource { listSnapshot(): Promise; listRecentSnapshot?(limit: number): Promise; @@ -354,7 +421,7 @@ export class CineCircleAllDebridReconciliationWorker { constructor( private readonly intake: CineCircleAllDebridIntake, - private readonly intervals: { recentMs: number; fullMs: number; recentLimit?: number }, + private readonly intervals: { recentMs: number; fullMs: number; recentLimit?: number; runFullOnStart?: boolean }, ) {} runRecent(): Promise { return this.intake.reconcile('recent', this.intervals.recentLimit || 30); } @@ -363,7 +430,7 @@ export class CineCircleAllDebridReconciliationWorker { start(): void { if (this.recentTimer || this.fullTimer) return; this.runRecent().catch(() => undefined); - this.runFull().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); } diff --git a/src/services/cinecircleAlldebridRuntime.ts b/src/services/cinecircleAlldebridRuntime.ts index 51088ad..1b0094b 100644 --- a/src/services/cinecircleAlldebridRuntime.ts +++ b/src/services/cinecircleAlldebridRuntime.ts @@ -14,8 +14,8 @@ import { export function cineCircleReconciliationRoutes(): { movies: ArrRoute; shows: ArrRoute } { return { - movies: { kind: "radarr", baseUrl: config.cineCircleRadarrUrl, apiKey: config.cineCircleRadarrApiKey, sourcePathPrefix: config.cineCircleAlldebridArrPath, importMode: config.cineCircleAlldebridArrImportMode === "Move" ? "Move" : "Copy" }, - shows: { kind: "sonarr", baseUrl: config.cineCircleSonarrUrl, apiKey: config.cineCircleSonarrApiKey, sourcePathPrefix: config.cineCircleAlldebridArrPath, importMode: config.cineCircleAlldebridArrImportMode === "Move" ? "Move" : "Copy" }, + movies: { kind: "radarr", baseUrl: config.cineCircleRadarrUrl, apiKey: config.cineCircleRadarrApiKey, sourcePathPrefix: config.cineCircleAlldebridArrPath, importMode: config.cineCircleAlldebridArrImportMode === "Move" ? "Move" : "Copy", symlinkLibraryPath: config.cineCircleAlldebridMoviesLibraryPath || undefined }, + shows: { kind: "sonarr", baseUrl: config.cineCircleSonarrUrl, apiKey: config.cineCircleSonarrApiKey, sourcePathPrefix: config.cineCircleAlldebridArrPath, importMode: config.cineCircleAlldebridArrImportMode === "Move" ? "Move" : "Copy", symlinkLibraryPath: config.cineCircleAlldebridShowsLibraryPath || undefined }, }; } @@ -56,6 +56,7 @@ export function createCineCircleAllDebridReconciliationWorker(): CineCircleAllDe recentMs: Math.max(1000, config.cineCircleAlldebridRecentIntervalMs), fullMs: Math.max(1000, config.cineCircleAlldebridFullIntervalMs), recentLimit: Math.max(1, config.cineCircleAlldebridRecentLimit), + runFullOnStart: config.cineCircleAlldebridRunFullOnStart, }); } diff --git a/tests/e2e/cinecircle-three-inputs.test.ts b/tests/e2e/cinecircle-three-inputs.test.ts index 59d8121..132a101 100644 --- a/tests/e2e/cinecircle-three-inputs.test.ts +++ b/tests/e2e/cinecircle-three-inputs.test.ts @@ -1,4 +1,7 @@ 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/cinecircleAlldebridIntake'; @@ -72,6 +75,28 @@ describe('CineCircle three-input fixture E2E', () => { 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(), 'cinecircle-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' }]; } }; From 6d3d259c32ef5c86b72f16ac959d4a0ed81700c1 Mon Sep 17 00:00:00 2001 From: samtruman <39094339+samtruman@users.noreply.github.com> Date: Fri, 25 Sep 2026 16:48:16 +0000 Subject: [PATCH 6/8] refactor: generalize provider reconciliation --- README.md | 20 +++ ...provider-reconciliation-test.override.yml} | 33 +++-- docs/provider-reconciliation.md | 45 ++++++ src/core/config.ts | 29 ++-- src/index.ts | 12 +- src/providers/alldebrid.ts | 24 ++++ src/services/cinecircleAlldebridRuntime.ts | 73 ---------- ...ridIntake.ts => providerReconciliation.ts} | 131 ++++++++++-------- src/services/providerReconciliationRuntime.ts | 62 +++++++++ ...vider-reconciliation-three-inputs.test.ts} | 10 +- ...est.ts => provider-reconciliation.test.ts} | 59 ++++---- .../cinecircle-alldebrid-runtime.test.ts | 50 ------- .../provider-reconciliation-runtime.test.ts | 46 ++++++ 13 files changed, 346 insertions(+), 248 deletions(-) rename deploy/{cinecircle-test-reconciliation.override.yml => provider-reconciliation-test.override.yml} (66%) create mode 100644 docs/provider-reconciliation.md delete mode 100644 src/services/cinecircleAlldebridRuntime.ts rename src/services/{cinecircleAlldebridIntake.ts => providerReconciliation.ts} (72%) create mode 100644 src/services/providerReconciliationRuntime.ts rename tests/e2e/{cinecircle-three-inputs.test.ts => provider-reconciliation-three-inputs.test.ts} (93%) rename tests/e2e/{cinecircle-alldebrid-intake.test.ts => provider-reconciliation.test.ts} (80%) delete mode 100644 tests/unit/services/cinecircle-alldebrid-runtime.test.ts create mode 100644 tests/unit/services/provider-reconciliation-runtime.test.ts diff --git a/README.md b/README.md index bd165c7..7b0ce11 100644 --- a/README.md +++ b/README.md @@ -909,6 +909,26 @@ 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) + +Provider reconciliation discovers completed files already present in configured +debrid providers, compares snapshots in SQLite, creates mount-backed symlinks, +and asks Radarr/Sonarr to rescan the existing library. 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/deploy/cinecircle-test-reconciliation.override.yml b/deploy/provider-reconciliation-test.override.yml similarity index 66% rename from deploy/cinecircle-test-reconciliation.override.yml rename to deploy/provider-reconciliation-test.override.yml index 9704003..9a2ff42 100644 --- a/deploy/cinecircle-test-reconciliation.override.yml +++ b/deploy/provider-reconciliation-test.override.yml @@ -1,25 +1,24 @@ -# Test-only override for the CineCircle AllDebrid reconciliation worker. +# Test-only override for the provider reconciliation worker. # Apply with the existing cinecircle-test compose file only. services: schrodrive-test: - image: schrodrive:cinecircle-alldebrid-reconciliation + image: schrodrive:provider-reconciliation environment: - CINECIRCLE_ALLDEBRID_RECONCILIATION_ENABLED: "true" - CINECIRCLE_ALLDEBRID_RECENT_INTERVAL_MS: "60000" - CINECIRCLE_ALLDEBRID_FULL_INTERVAL_MS: "86400000" - CINECIRCLE_ALLDEBRID_RECENT_LIMIT: "1" - CINECIRCLE_ALLDEBRID_RUN_FULL_ON_START: "false" - CINECIRCLE_ALLDEBRID_DRY_RUN: "false" - CINECIRCLE_RADARR_URL: http://radarr-test:7878 - CINECIRCLE_RADARR_API_KEY: ${CINECIRCLE_RADARR_API_KEY:-} - CINECIRCLE_SONARR_URL: http://sonarr-test:8989 - CINECIRCLE_SONARR_API_KEY: ${CINECIRCLE_SONARR_API_KEY:-} - CINECIRCLE_ALLDEBRID_ARR_PATH: /mnt/schrodrive/alldebrid - CINECIRCLE_ALLDEBRID_ARR_IMPORT_MODE: Copy - CINECIRCLE_ALLDEBRID_MOVIES_LIBRARY_PATH: /mnt/cinecircle-test/radarr/movies - CINECIRCLE_ALLDEBRID_SHOWS_LIBRARY_PATH: /mnt/cinecircle-test/sonarr/shows + PROVIDER_RECONCILIATION_ENABLED: "true" + PROVIDER_RECONCILIATION_RECENT_INTERVAL_MS: "60000" + PROVIDER_RECONCILIATION_FULL_INTERVAL_MS: "86400000" + PROVIDER_RECONCILIATION_RECENT_LIMIT: "1" + PROVIDER_RECONCILIATION_RUN_FULL_ON_START: "false" + PROVIDER_RECONCILIATION_DRY_RUN: "false" + PROVIDER_RECONCILIATION_RADARR_URL: http://radarr-test:7878 + PROVIDER_RECONCILIATION_RADARR_API_KEY: ${CINECIRCLE_RADARR_API_KEY:-} + PROVIDER_RECONCILIATION_SONARR_URL: http://sonarr-test:8989 + PROVIDER_RECONCILIATION_SONARR_API_KEY: ${CINECIRCLE_SONARR_API_KEY:-} + PROVIDER_RECONCILIATION_MOUNT_BASE: /mnt/schrodrive + PROVIDER_RECONCILIATION_MOVIES_LIBRARY_PATH: /mnt/cinecircle-test/radarr/movies + PROVIDER_RECONCILIATION_SHOWS_LIBRARY_PATH: /mnt/cinecircle-test/sonarr/shows ALLDEBRID_API_KEY: ${ALLDEBRID_API_KEY:-} - # Direct AllDebrid polling is the only input under test here. + # Direct provider polling is the only input under test here. RUN_WEBHOOK: "false" RUN_POLLER: "false" RUN_WATCHLIST_POLLER: "false" diff --git a/docs/provider-reconciliation.md b/docs/provider-reconciliation.md new file mode 100644 index 0000000..66326af --- /dev/null +++ b/docs/provider-reconciliation.md @@ -0,0 +1,45 @@ +# 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. + +```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 dfd05dc..a9b15a7 100644 --- a/src/core/config.ts +++ b/src/core/config.ts @@ -59,21 +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", - // CineCircle direct AllDebrid reconciliation is opt-in and disabled by default. - cineCircleAlldebridReconciliationEnabled: String(process.env.CINECIRCLE_ALLDEBRID_RECONCILIATION_ENABLED ?? "false").toLowerCase() === "true", - cineCircleAlldebridRecentIntervalMs: Number(process.env.CINECIRCLE_ALLDEBRID_RECENT_INTERVAL_MS || 900000), - cineCircleAlldebridFullIntervalMs: Number(process.env.CINECIRCLE_ALLDEBRID_FULL_INTERVAL_MS || 21600000), - cineCircleAlldebridRecentLimit: Number(process.env.CINECIRCLE_ALLDEBRID_RECENT_LIMIT || 30), - cineCircleAlldebridRunFullOnStart: String(process.env.CINECIRCLE_ALLDEBRID_RUN_FULL_ON_START ?? "true").toLowerCase() !== "false", - cineCircleAlldebridDryRun: String(process.env.CINECIRCLE_ALLDEBRID_DRY_RUN ?? "false").toLowerCase() === "true", - cineCircleRadarrUrl: process.env.CINECIRCLE_RADARR_URL || "", - cineCircleRadarrApiKey: process.env.CINECIRCLE_RADARR_API_KEY || "", - cineCircleSonarrUrl: process.env.CINECIRCLE_SONARR_URL || "", - cineCircleSonarrApiKey: process.env.CINECIRCLE_SONARR_API_KEY || "", - cineCircleAlldebridArrPath: process.env.CINECIRCLE_ALLDEBRID_ARR_PATH || "/mnt/schrodrive/alldebrid", - cineCircleAlldebridArrImportMode: process.env.CINECIRCLE_ALLDEBRID_ARR_IMPORT_MODE || "Copy", - cineCircleAlldebridMoviesLibraryPath: process.env.CINECIRCLE_ALLDEBRID_MOVIES_LIBRARY_PATH || "", - cineCircleAlldebridShowsLibraryPath: process.env.CINECIRCLE_ALLDEBRID_SHOWS_LIBRARY_PATH || "", + // 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 2a3b0cc..6d617fc 100644 --- a/src/index.ts +++ b/src/index.ts @@ -12,8 +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 { startCineCircleAllDebridReconciliation } from "./services/cinecircleAlldebridRuntime"; -import type { CineCircleAllDebridReconciliationWorker } from "./services/cinecircleAlldebridIntake"; +import { startProviderReconciliation } from "./services/providerReconciliationRuntime"; +import type { ProviderReconciliationWorker } from "./services/providerReconciliation"; const program = new Command(); program @@ -33,7 +33,7 @@ program } // Register graceful shutdown handlers - let cineCircleReconciliationWorker: CineCircleAllDebridReconciliationWorker | undefined; + let providerReconciliationWorker: ProviderReconciliationWorker | undefined; const shutdown = () => { console.log(`[${new Date().toISOString()}][serve] Shutting down — unmounting FUSE drives...`); try { @@ -46,7 +46,7 @@ program stopStrmServer().catch(() => {}); stopCloudLinksBridge().catch(() => {}); stopArrBridge().catch(() => {}); - cineCircleReconciliationWorker?.stop(); + providerReconciliationWorker?.stop(); setTimeout(() => { console.log(`[${new Date().toISOString()}][serve] Closing database and exiting...`); closeDb(); @@ -81,7 +81,7 @@ program }); } - // The AllDebrid reconciliation worker submits Arr scans against this + // 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. @@ -108,7 +108,7 @@ program startWatchlistPoller(); } - cineCircleReconciliationWorker = startCineCircleAllDebridReconciliation(); + 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/cinecircleAlldebridRuntime.ts b/src/services/cinecircleAlldebridRuntime.ts deleted file mode 100644 index 1b0094b..0000000 --- a/src/services/cinecircleAlldebridRuntime.ts +++ /dev/null @@ -1,73 +0,0 @@ -import { config } from "../core/config"; -import { registry } from "../providers"; -import { AllDebridProvider } from "../providers/alldebrid"; -import { parseMediaFilename } from "./mediaParser"; -import { recordOrganizerReview } from "./organizerReview"; -import { - AllDebridProviderSource, - CineCircleAllDebridIntake, - CineCircleAllDebridReconciliationWorker, - HttpArrClient, - SqliteIntakeStateStore, - type ArrRoute, -} from "./cinecircleAlldebridIntake"; - -export function cineCircleReconciliationRoutes(): { movies: ArrRoute; shows: ArrRoute } { - return { - movies: { kind: "radarr", baseUrl: config.cineCircleRadarrUrl, apiKey: config.cineCircleRadarrApiKey, sourcePathPrefix: config.cineCircleAlldebridArrPath, importMode: config.cineCircleAlldebridArrImportMode === "Move" ? "Move" : "Copy", symlinkLibraryPath: config.cineCircleAlldebridMoviesLibraryPath || undefined }, - shows: { kind: "sonarr", baseUrl: config.cineCircleSonarrUrl, apiKey: config.cineCircleSonarrApiKey, sourcePathPrefix: config.cineCircleAlldebridArrPath, importMode: config.cineCircleAlldebridArrImportMode === "Move" ? "Move" : "Copy", symlinkLibraryPath: config.cineCircleAlldebridShowsLibraryPath || undefined }, - }; -} - -/** Builds the opt-in direct AllDebrid worker without changing provider lifecycle services. */ -export function createCineCircleAllDebridReconciliationWorker(): CineCircleAllDebridReconciliationWorker | undefined { - const provider = registry.get("alldebrid"); - if (!(provider instanceof AllDebridProvider)) { - console.warn("[cinecircle-reconciliation] AllDebrid provider is not registered; worker disabled"); - return undefined; - } - if (!provider.isConfigured()) { - console.warn("[cinecircle-reconciliation] AllDebrid credentials are not configured; worker disabled"); - return undefined; - } - - const routes = cineCircleReconciliationRoutes(); - if (!routes.movies.baseUrl || !routes.movies.apiKey || !routes.shows.baseUrl || !routes.shows.apiKey) { - console.warn("[cinecircle-reconciliation] Radarr/Sonarr routes are incomplete; worker disabled"); - return undefined; - } - - const source = new AllDebridProviderSource(provider); - const arr = new HttpArrClient(); - const store = new SqliteIntakeStateStore(); - const intake = new CineCircleAllDebridIntake(source, arr, store, { - dryRun: config.cineCircleAlldebridDryRun, - routeFor: (category) => category === "Movies" ? routes.movies : routes.shows, - onReview: async (event, error) => { - const parsed = parseMediaFilename(event.path, event.path); - recordOrganizerReview(event.path, { - ...parsed, - status: parsed.status === "matched" ? "ambiguous" : parsed.status, - reason: `AllDebrid direct intake: ${error.message}`, - }); - }, - }); - return new CineCircleAllDebridReconciliationWorker(intake, { - recentMs: Math.max(1000, config.cineCircleAlldebridRecentIntervalMs), - fullMs: Math.max(1000, config.cineCircleAlldebridFullIntervalMs), - recentLimit: Math.max(1, config.cineCircleAlldebridRecentLimit), - runFullOnStart: config.cineCircleAlldebridRunFullOnStart, - }); -} - -export function startCineCircleAllDebridReconciliation(): CineCircleAllDebridReconciliationWorker | undefined { - if (!config.cineCircleAlldebridReconciliationEnabled) { - console.log("[cinecircle-reconciliation] disabled (CINECIRCLE_ALLDEBRID_RECONCILIATION_ENABLED=false)"); - return undefined; - } - const worker = createCineCircleAllDebridReconciliationWorker(); - if (!worker) return undefined; - worker.start(); - console.log("[cinecircle-reconciliation] started; direct AllDebrid intake is read-only and provider deletion/repair is not used"); - return worker; -} diff --git a/src/services/cinecircleAlldebridIntake.ts b/src/services/providerReconciliation.ts similarity index 72% rename from src/services/cinecircleAlldebridIntake.ts rename to src/services/providerReconciliation.ts index 73f1a69..5899b9d 100644 --- a/src/services/cinecircleAlldebridIntake.ts +++ b/src/services/providerReconciliation.ts @@ -1,18 +1,17 @@ /** - * CineCircle fork-only AllDebrid direct-file intake. + * Provider-agnostic direct-file reconciliation for mounted debrid providers. * - * AllDebrid has no push change feed in the integration used by SchröDrive, so - * this adapter reconciles read-only magnet status plus completed file trees, - * emits stable direct-file events, and hands added/changed files to Arr. - * It is deliberately not wired into the generic provider lifecycle. + * 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 { AllDebridProvider } from '../providers/alldebrid'; -import type { TorrentInfo, VirtualDirectory } from '../providers'; +import type { DebridProvider, TorrentInfo, VirtualDirectory } from '../providers'; import { classifyTorrent } from '../core/mediaClassifier'; import { parseMediaFilename } from './mediaParser'; @@ -20,17 +19,19 @@ export type DirectFileAction = 'added' | 'changed' | 'deleted'; export type SourceCategory = 'Movies' | 'Shows'; export type ArrKind = 'radarr' | 'sonarr'; -export interface AllDebridSnapshot { +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: 'alldebrid'; + provider: string; providerItemId: string; action: DirectFileAction; path: string; @@ -74,6 +75,7 @@ export interface IntakeStateStore { } export interface IntakeState { + provider?: string; providerItemId: string; fingerprint: string; path: string; @@ -99,51 +101,51 @@ export class InMemoryIntakeStateStore implements IntakeStateStore { saveCursor(mode: 'recent' | 'full', observedAt: string): void { this.cursor[mode === 'recent' ? 'recentAt' : 'fullAt'] = observedAt; } } -/** Persistent fork state; creates only its own table in the configured test DB. */ +/** Persistent reconciliation state in the configured SchröDrive DB. */ export class SqliteIntakeStateStore implements IntakeStateStore { constructor() { - getDb().exec(`CREATE TABLE IF NOT EXISTS cinecircle_alldebrid_intake ( + 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 cinecircle_alldebrid_events ( + 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 cinecircle_alldebrid_cursor ( + 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 cinecircle_alldebrid_intake WHERE provider_item_id = ?').get(id) as { state_json?: string } | 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 cinecircle_alldebrid_intake').all() as Array<{ state_json: string }>) + 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 cinecircle_alldebrid_intake(provider_item_id,state_json,updated_at) + 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 cinecircle_alldebrid_events WHERE dedupe_key = ?').get(key); + 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 cinecircle_alldebrid_events(dedupe_key,event_json,arr_json,created_at) + 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 cinecircle_alldebrid_cursor').all() as Array<{ name: string; observed_at: 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 cinecircle_alldebrid_cursor(name,observed_at) VALUES (?,?) + 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); } } @@ -226,7 +228,9 @@ async function exposeAsSymlink(route: ArrRoute, event: DirectFileEvent, provider if (existing?.isSymbolicLink()) { const current = await fsp.readlink(destination); if (current === providerPath) return destination; - await fsp.unlink(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}`); } @@ -237,16 +241,16 @@ async function exposeAsSymlink(route: ArrRoute, event: DirectFileEvent, provider return destination; } -export interface AllDebridReadOnlySource { - listSnapshot(): Promise; - listRecentSnapshot?(limit: number): Promise; +export interface ProviderReadOnlySource { + listSnapshot(): Promise; + listRecentSnapshot?(limit: number): Promise; } -/** Uses only the existing provider's status and completed-directory methods. */ -export class AllDebridProviderSource implements AllDebridReadOnlySource { - constructor(private readonly provider: Pick) {} +/** 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 { + async listSnapshot(): Promise { const observedAt = new Date().toISOString(); const torrents = await this.provider.listTorrents(); const directories = await this.provider.fetchDirectories(); @@ -254,20 +258,22 @@ export class AllDebridProviderSource implements AllDebridReadOnlySource { return this.toSnapshots(torrents, trees, observedAt); } - async listRecentSnapshot(limit: number): Promise { + 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 = await this.provider.fetchDirectoriesForIds(torrents); + 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): AllDebridSnapshot[] { + private toSnapshots(torrents: TorrentInfo[], trees: Map, observedAt: string): ProviderSnapshot[] { return torrents.map((torrent) => { const directory = trees.get(String(torrent.id)); - return { providerItemId: String(torrent.id), name: torrent.name, directoryName: directory?.name, status: torrent.status, - files: (directory?.files || []).map((file) => ({ path: file.name, size: file.size })), observedAt }; + 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 }; }); } } @@ -275,7 +281,7 @@ export class AllDebridProviderSource implements AllDebridReadOnlySource { export interface IntakeOptions { dryRun?: boolean; maxAttempts?: number; - routeFor: (category: SourceCategory) => ArrRoute; + routeFor: (category: SourceCategory, provider?: string) => ArrRoute; onEvent?: (event: DirectFileEvent) => Promise | void; onReview?: (event: DirectFileEvent, error: Error) => Promise | void; } @@ -289,21 +295,28 @@ export function isMediaFile(filePath: string): boolean { return VIDEO_EXTENSIONS.has(extension) || SUBTITLE_EXTENSIONS.has(extension); } -function categoryFor(snapshot: AllDebridSnapshot): SourceCategory { +function categoryFor(snapshot: ProviderSnapshot): SourceCategory { return classifyTorrent(snapshot.name, snapshot.files.map((file) => file.path)) === 'shows' ? 'Shows' : 'Movies'; } -function fingerprint(snapshot: Pick): string { +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(id: string, action: DirectFileAction, fp: string): string { - return `alldebrid:${id}:${action}:${fp}`; +function eventKey(provider: string, id: string, action: DirectFileAction, fp: string): string { + return `${provider}:${id}:${action}:${fp}`; } -export class CineCircleAllDebridIntake { +export class ProviderReconciliationIntake { constructor( - private readonly source: AllDebridReadOnlySource, + private readonly source: ProviderReadOnlySource, private readonly arr: ArrClient, private readonly store: IntakeStateStore, private readonly options: IntakeOptions, @@ -318,7 +331,7 @@ export class CineCircleAllDebridIntake { const events: DirectFileEvent[] = []; for (const item of current) { - if (item.status !== 'finished' || item.files.length === 0) continue; + 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); @@ -335,9 +348,9 @@ export class CineCircleAllDebridIntake { for (const previous of mode === 'full' ? this.store.listItems().filter((item) => item.lastAction !== 'deleted') : []) { if (seen.has(previous.providerItemId)) continue; const event: DirectFileEvent = { - provider: 'alldebrid', providerItemId: previous.providerItemId, action: 'deleted', + provider: previous.provider || 'alldebrid', providerItemId: previous.providerItemId, action: 'deleted', path: previous.path, tree: [], sourceCategory: previous.sourceCategory, - observedAt: new Date().toISOString(), stableDedupeKey: eventKey(previous.providerItemId, 'deleted', previous.fingerprint), + 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); @@ -355,14 +368,14 @@ export class CineCircleAllDebridIntake { 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.commandId); + 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: 'alldebrid', providerItemId: item.providerItemId, action: item.lastAction, + 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.providerItemId, item.lastAction, item.fingerprint), + stableDedupeKey: eventKey(item.provider || 'alldebrid', item.providerItemId, item.lastAction, item.fingerprint), }, new Error(`Arr command ${item.commandId} failed`)); } } catch { @@ -371,7 +384,8 @@ export class CineCircleAllDebridIntake { } } - private makeEvent(item: AllDebridSnapshot, action: DirectFileAction, fp: string): DirectFileEvent { + 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(/^\/+/, ''); @@ -379,27 +393,27 @@ export class CineCircleAllDebridIntake { ? `${categoryDirectory}/${item.directoryName.replace(/^\/+|\/+$/g, '')}/${providerPath}` : providerPath; return { - provider: 'alldebrid', providerItemId: item.providerItemId, action, + provider, providerItemId: item.providerItemId, action, path, tree: item.files, sourceCategory: category, - observedAt: item.observedAt, stableDedupeKey: eventKey(item.providerItemId, action, fp), + observedAt: item.observedAt, stableDedupeKey: eventKey(provider, item.providerItemId, action, fp), }; } - private async dispatch(event: DirectFileEvent, item: AllDebridSnapshot, fp: string): Promise { + 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({ providerItemId: item.providerItemId, fingerprint: fp, path: event.path, tree: event.tree, sourceCategory: event.sourceCategory, lastAction: event.action, updatedAt: event.observedAt }); + 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); + 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({ 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 }); + 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; @@ -413,19 +427,20 @@ export class CineCircleAllDebridIntake { /** * Testable scheduler for the fork worker. It is intentionally not started by - * the application entry point; CineCircle wiring must explicitly opt in. + * the application entry point; provider reconciliation must explicitly opt in. */ -export class CineCircleAllDebridReconciliationWorker { +export class ProviderReconciliationWorker { private recentTimer: ReturnType | undefined; private fullTimer: ReturnType | undefined; constructor( - private readonly intake: CineCircleAllDebridIntake, + private readonly intake: ProviderReconciliationIntake | ProviderReconciliationIntake[], private readonly intervals: { recentMs: number; fullMs: number; recentLimit?: number; runFullOnStart?: boolean }, ) {} - runRecent(): Promise { return this.intake.reconcile('recent', this.intervals.recentLimit || 30); } - runFull(): Promise { return this.intake.reconcile('full'); } + 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; 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/cinecircle-three-inputs.test.ts b/tests/e2e/provider-reconciliation-three-inputs.test.ts similarity index 93% rename from tests/e2e/cinecircle-three-inputs.test.ts rename to tests/e2e/provider-reconciliation-three-inputs.test.ts index 132a101..bf7bfdb 100644 --- a/tests/e2e/cinecircle-three-inputs.test.ts +++ b/tests/e2e/provider-reconciliation-three-inputs.test.ts @@ -4,7 +4,7 @@ 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/cinecircleAlldebridIntake'; +import { HttpArrClient, type ArrRoute, type DirectFileEvent } from '../../src/services/providerReconciliation'; const originalFetch = globalThis.fetch; @@ -22,7 +22,7 @@ function route(kind: 'radarr' | 'sonarr', sourcePathPrefix?: string): ArrRoute { afterEach(() => { globalThis.fetch = originalFetch; }); -describe('CineCircle three-input fixture E2E', () => { +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'); @@ -100,13 +100,13 @@ describe('CineCircle three-input fixture E2E', () => { 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/cinecircleAlldebridIntake')).InMemoryIntakeStateStore(); + 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 { CineCircleAllDebridIntake } = await import('../../src/services/cinecircleAlldebridIntake'); - const intake = new CineCircleAllDebridIntake(source, arr, store, { routeFor: () => route('sonarr') }); + 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']); diff --git a/tests/e2e/cinecircle-alldebrid-intake.test.ts b/tests/e2e/provider-reconciliation.test.ts similarity index 80% rename from tests/e2e/cinecircle-alldebrid-intake.test.ts rename to tests/e2e/provider-reconciliation.test.ts index 6f1be94..f4e15e3 100644 --- a/tests/e2e/cinecircle-alldebrid-intake.test.ts +++ b/tests/e2e/provider-reconciliation.test.ts @@ -1,25 +1,25 @@ import { describe, expect, test } from 'bun:test'; import { closeDb } from '../../src/core/db'; import { - CineCircleAllDebridIntake, + ProviderReconciliationIntake, InMemoryIntakeStateStore, SqliteIntakeStateStore, - type AllDebridSnapshot, + type ProviderSnapshot, type ArrClient, type ArrCommandResult, type DirectFileEvent, isMediaFile, - AllDebridProviderSource, - CineCircleAllDebridReconciliationWorker, -} from '../../src/services/cinecircleAlldebridIntake'; + ProviderSnapshotSource, + ProviderReconciliationWorker, +} from '../../src/services/providerReconciliation'; -function snapshot(id: string, name: string, files = [{ path: `${name}.mkv`, size: 10 }]): AllDebridSnapshot { +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: AllDebridSnapshot[][]) {} - async listSnapshot(): Promise { return this.rounds.shift() || []; } + constructor(private readonly rounds: ProviderSnapshot[][]) {} + async listSnapshot(): Promise { return this.rounds.shift() || []; } } class FakeArr implements ArrClient { @@ -41,7 +41,18 @@ const routeFor = (category: 'Movies' | 'Shows') => ({ apiKey: 'fixture-key', }); -describe('CineCircle AllDebrid direct intake', () => { +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; @@ -61,7 +72,7 @@ describe('CineCircle AllDebrid direct intake', () => { 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 AllDebridProviderSource(provider as any); + 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']); @@ -70,8 +81,8 @@ describe('CineCircle AllDebrid direct intake', () => { expect(fullCalls).toBe(1); const arr = new FakeArr(); - const intake = new CineCircleAllDebridIntake( - new AllDebridProviderSource(provider as any), arr, new InMemoryIntakeStateStore(), { routeFor }, + 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); @@ -81,7 +92,7 @@ describe('CineCircle AllDebrid direct intake', () => { test('emits add and routes Movies to Radarr and Shows to Sonarr', async () => { const arr = new FakeArr(); - const intake = new CineCircleAllDebridIntake( + const intake = new ProviderReconciliationIntake( new SequenceSource([[snapshot('m-1', 'Movie (2026)'), snapshot('s-1', 'Show S01E02')]]), arr, new InMemoryIntakeStateStore(), @@ -98,7 +109,7 @@ describe('CineCircle AllDebrid direct intake', () => { 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 CineCircleAllDebridIntake( + 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 }])], @@ -115,7 +126,7 @@ describe('CineCircle AllDebrid direct intake', () => { test('suppresses duplicate delivery and emits deletion from a missing status item', async () => { const store = new InMemoryIntakeStateStore(); const arr = new FakeArr(); - const intake = new CineCircleAllDebridIntake( + const intake = new ProviderReconciliationIntake( new SequenceSource([[snapshot('m-1', 'Movie (2026)')], [snapshot('m-1', 'Movie (2026)')], []]), arr, store, { routeFor }, ); @@ -128,7 +139,7 @@ describe('CineCircle AllDebrid direct intake', () => { test('dry-run persists state without submitting to Arr', async () => { const arr = new FakeArr(); const store = new InMemoryIntakeStateStore(); - const intake = new CineCircleAllDebridIntake( + const intake = new ProviderReconciliationIntake( new SequenceSource([[snapshot('m-1', 'Movie (2026)')]]), arr, store, { routeFor, dryRun: true }, ); @@ -141,7 +152,7 @@ describe('CineCircle AllDebrid direct intake', () => { 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 CineCircleAllDebridIntake( + const intake = new ProviderReconciliationIntake( new SequenceSource([[snapshot('m-1', 'Movie (2026)', names.map((path) => ({ path, size: 1 })))]]) , arr, new InMemoryIntakeStateStore(), { routeFor }, ); @@ -153,11 +164,11 @@ describe('CineCircle AllDebrid direct intake', () => { test('restarts from persisted state and polls the pending Arr command', async () => { const store = new InMemoryIntakeStateStore(); const firstArr = new FakeArr(); - await new CineCircleAllDebridIntake( + await new ProviderReconciliationIntake( new SequenceSource([[snapshot('m-1', 'Movie (2026)')]]), firstArr, store, { routeFor }, ).reconcile(); const secondArr = new FakeArr(); - const events = await new CineCircleAllDebridIntake( + const events = await new ProviderReconciliationIntake( new SequenceSource([[snapshot('m-1', 'Movie (2026)')]]), secondArr, store, { routeFor }, ).reconcile(); expect(events).toEqual([]); @@ -170,7 +181,7 @@ describe('CineCircle AllDebrid direct intake', () => { const firstArr = new FakeArr(); const firstStore = new SqliteIntakeStateStore(); const source = new SequenceSource([[snapshot(itemId, 'Movie (2026)')]]); - const first = new CineCircleAllDebridIntake(source, firstArr, firstStore, { routeFor }); + 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(); @@ -179,7 +190,7 @@ describe('CineCircle AllDebrid direct intake', () => { closeDb(); const secondStore = new SqliteIntakeStateStore(); const secondArr = new FakeArr(); - const second = new CineCircleAllDebridIntake( + const second = new ProviderReconciliationIntake( new SequenceSource([[snapshot(itemId, 'Movie (2026)')]]), secondArr, secondStore, { routeFor }, ); expect(await second.reconcile('full')).toEqual([]); @@ -198,7 +209,7 @@ describe('CineCircle AllDebrid direct intake', () => { }, async getCommand(_route, commandId) { return { commandId, status: 'completed' }; }, }; - const intake = new CineCircleAllDebridIntake( + const intake = new ProviderReconciliationIntake( new SequenceSource([[snapshot('s-1', 'Show S01E02')]]), arr, new InMemoryIntakeStateStore(), { routeFor, maxAttempts: 3 }, ); @@ -212,7 +223,7 @@ describe('CineCircle AllDebrid direct intake', () => { async submitScan() { throw new Error('permanent'); }, async getCommand(_route, commandId) { return { commandId, status: 'failed' }; }, }; - const intake = new CineCircleAllDebridIntake( + const intake = new ProviderReconciliationIntake( new SequenceSource([[snapshot('m-1', 'Movie (2026)')]]), arr, new InMemoryIntakeStateStore(), { routeFor, maxAttempts: 2, onReview: (event) => { reviewed.push(event); } }, ); @@ -224,7 +235,7 @@ describe('CineCircle AllDebrid direct intake', () => { 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 CineCircleAllDebridReconciliationWorker(intake as any, { recentMs: 10, fullMs: 15, recentLimit: 2 }); + 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)); diff --git a/tests/unit/services/cinecircle-alldebrid-runtime.test.ts b/tests/unit/services/cinecircle-alldebrid-runtime.test.ts deleted file mode 100644 index 8b38c53..0000000 --- a/tests/unit/services/cinecircle-alldebrid-runtime.test.ts +++ /dev/null @@ -1,50 +0,0 @@ -import { afterEach, describe, expect, test } from 'bun:test'; -import { config } from '../../../src/core/config'; -import { cineCircleReconciliationRoutes, startCineCircleAllDebridReconciliation } from '../../../src/services/cinecircleAlldebridRuntime'; - -const original = { - enabled: config.cineCircleAlldebridReconciliationEnabled, - radarrUrl: config.cineCircleRadarrUrl, - radarrKey: config.cineCircleRadarrApiKey, - sonarrUrl: config.cineCircleSonarrUrl, - sonarrKey: config.cineCircleSonarrApiKey, - alldebridArrPath: config.cineCircleAlldebridArrPath, - alldebridArrImportMode: config.cineCircleAlldebridArrImportMode, -}; - -afterEach(() => { - config.cineCircleAlldebridReconciliationEnabled = original.enabled; - config.cineCircleRadarrUrl = original.radarrUrl; - config.cineCircleRadarrApiKey = original.radarrKey; - config.cineCircleSonarrUrl = original.sonarrUrl; - config.cineCircleSonarrApiKey = original.sonarrKey; - config.cineCircleAlldebridArrPath = original.alldebridArrPath; - config.cineCircleAlldebridArrImportMode = original.alldebridArrImportMode; -}); - -describe('CineCircle AllDebrid runtime wiring', () => { - test('is disabled by default and does not construct a provider worker', () => { - config.cineCircleAlldebridReconciliationEnabled = false; - expect(startCineCircleAllDebridReconciliation()).toBeUndefined(); - }); - - test('requires explicit complete Arr routes when enabled', () => { - config.cineCircleAlldebridReconciliationEnabled = true; - config.cineCircleRadarrUrl = 'http://radarr.test'; - config.cineCircleRadarrApiKey = 'radarr-fixture-key'; - config.cineCircleSonarrUrl = ''; - config.cineCircleSonarrApiKey = ''; - expect(startCineCircleAllDebridReconciliation()).toBeUndefined(); - }); - - test('keeps explicit Movies/Radarr and Shows/Sonarr routing', () => { - config.cineCircleRadarrUrl = 'http://radarr.test/'; - config.cineCircleRadarrApiKey = 'radarr-fixture-key'; - config.cineCircleSonarrUrl = 'http://sonarr.test/'; - config.cineCircleSonarrApiKey = 'sonarr-fixture-key'; - expect(cineCircleReconciliationRoutes()).toEqual({ - movies: { kind: 'radarr', baseUrl: 'http://radarr.test/', apiKey: 'radarr-fixture-key', sourcePathPrefix: '/mnt/schrodrive/alldebrid', importMode: 'Copy' }, - shows: { kind: 'sonarr', baseUrl: 'http://sonarr.test/', apiKey: 'sonarr-fixture-key', sourcePathPrefix: '/mnt/schrodrive/alldebrid', importMode: 'Copy' }, - }); - }); -}); 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' }, + }); + }); +}); From 57444121f537da61a61acf51936b0c3c6b389749 Mon Sep 17 00:00:00 2001 From: samtruman <39094339+samtruman@users.noreply.github.com> Date: Fri, 25 Sep 2026 16:51:34 +0000 Subject: [PATCH 7/8] chore: remove local test stack references --- .../provider-reconciliation-test.override.yml | 80 ------------------- tests/e2e/arr-bridge/qbittorrent-api.test.ts | 4 +- ...ovider-reconciliation-three-inputs.test.ts | 2 +- 3 files changed, 3 insertions(+), 83 deletions(-) delete mode 100644 deploy/provider-reconciliation-test.override.yml diff --git a/deploy/provider-reconciliation-test.override.yml b/deploy/provider-reconciliation-test.override.yml deleted file mode 100644 index 9a2ff42..0000000 --- a/deploy/provider-reconciliation-test.override.yml +++ /dev/null @@ -1,80 +0,0 @@ -# Test-only override for the provider reconciliation worker. -# Apply with the existing cinecircle-test compose file only. -services: - schrodrive-test: - image: schrodrive:provider-reconciliation - environment: - PROVIDER_RECONCILIATION_ENABLED: "true" - PROVIDER_RECONCILIATION_RECENT_INTERVAL_MS: "60000" - PROVIDER_RECONCILIATION_FULL_INTERVAL_MS: "86400000" - PROVIDER_RECONCILIATION_RECENT_LIMIT: "1" - PROVIDER_RECONCILIATION_RUN_FULL_ON_START: "false" - PROVIDER_RECONCILIATION_DRY_RUN: "false" - PROVIDER_RECONCILIATION_RADARR_URL: http://radarr-test:7878 - PROVIDER_RECONCILIATION_RADARR_API_KEY: ${CINECIRCLE_RADARR_API_KEY:-} - PROVIDER_RECONCILIATION_SONARR_URL: http://sonarr-test:8989 - PROVIDER_RECONCILIATION_SONARR_API_KEY: ${CINECIRCLE_SONARR_API_KEY:-} - PROVIDER_RECONCILIATION_MOUNT_BASE: /mnt/schrodrive - PROVIDER_RECONCILIATION_MOVIES_LIBRARY_PATH: /mnt/cinecircle-test/radarr/movies - PROVIDER_RECONCILIATION_SHOWS_LIBRARY_PATH: /mnt/cinecircle-test/sonarr/shows - ALLDEBRID_API_KEY: ${ALLDEBRID_API_KEY:-} - # Direct provider polling is the only input under test here. - RUN_WEBHOOK: "false" - RUN_POLLER: "false" - RUN_WATCHLIST_POLLER: "false" - RUN_DEAD_SCANNER: "false" - RUN_DEAD_SCANNER_WATCH: "false" - ENABLE_REPAIR: "false" - PREEMPTIVE_REPAIR: "false" - RUN_ORGANIZER_WATCH: "false" - volumes: - - type: bind - source: /home/samtruman/docker/cinecircle-test/schrodrive - target: /mnt/schrodrive - bind: - propagation: rshared - - type: bind - source: /home/samtruman/cinecircle-test-clean/schrodrive/data - target: /data - - type: bind - source: /home/samtruman/cinecircle-test-clean/schrodrive/config - target: /config - - type: bind - source: /home/samtruman/cinecircle-test-clean - target: /mnt/cinecircle-test - - radarr-test: - volumes: - - type: bind - source: /home/samtruman/docker/cinecircle-test/radarr/config - target: /config - - type: bind - source: /home/samtruman/docker/cinecircle-test - target: /mnt/cinecircle-test - - type: bind - source: /home/samtruman/docker/cinecircle-test/schrodrive - target: /mnt/schrodrive - read_only: true - bind: - propagation: rslave - - type: bind - source: /home/samtruman/docker/cinecircle-test/schrodrive/downloads - target: /mnt/schrodrive/downloads - - sonarr-test: - volumes: - - type: bind - source: /home/samtruman/cinecircle-test-clean/sonarr/config - target: /config - - type: bind - source: /home/samtruman/cinecircle-test-clean - target: /mnt/cinecircle-test - - type: bind - source: /home/samtruman/docker/cinecircle-test/schrodrive - target: /mnt/schrodrive - read_only: true - bind: - propagation: rslave - - type: bind - source: /home/samtruman/docker/cinecircle-test/schrodrive/downloads - target: /mnt/schrodrive/downloads 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 index bf7bfdb..3298a24 100644 --- a/tests/e2e/provider-reconciliation-three-inputs.test.ts +++ b/tests/e2e/provider-reconciliation-three-inputs.test.ts @@ -82,7 +82,7 @@ describe('provider reconciliation fixture E2E', () => { 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(), 'cinecircle-symlink-')); + 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'); From 0b378b0d8010ecc6177aca7f2f13f0875bd76b00 Mon Sep 17 00:00:00 2001 From: samtruman <39094339+samtruman@users.noreply.github.com> Date: Fri, 25 Sep 2026 16:57:27 +0000 Subject: [PATCH 8/8] docs: clarify ongoing provider intake scope --- README.md | 10 ++++++---- docs/provider-reconciliation.md | 5 +++++ 2 files changed, 11 insertions(+), 4 deletions(-) diff --git a/README.md b/README.md index 7b0ce11..08e6d0e 100644 --- a/README.md +++ b/README.md @@ -911,10 +911,12 @@ All configuration is done via environment variables. Below is the complete refer ### 🔄 Provider Reconciliation (opt-in) -Provider reconciliation discovers completed files already present in configured -debrid providers, compares snapshots in SQLite, creates mount-backed symlinks, -and asks Radarr/Sonarr to rescan the existing library. It does not copy media -locally and does not delete or repair provider content. +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 | | --- | --- | --- | diff --git a/docs/provider-reconciliation.md b/docs/provider-reconciliation.md index 66326af..2f3e5fc 100644 --- a/docs/provider-reconciliation.md +++ b/docs/provider-reconciliation.md @@ -4,6 +4,11 @@ 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