stremio/apps/master-api/src/downloads/service.ts
Jos Vooges | STH bd21b01956 Add full series download flow with episode and season import.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-08 02:30:40 +02:00

1520 lines
49 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

import { prisma } from "../database/client";
import { decryptSecret, encryptSecret } from "../security/crypto";
import { AppError } from "../security/errors";
import { nodeConnectionManager } from "../websocket/manager";
import { MetadataService } from "../metadata/service";
import {
SynologyDownloadStation,
cleanReleaseQuery,
folderNameForMovie,
type DsSearchHit,
type DsTask,
} from "../synology/download-station";
import {
OpenSubtitlesClient,
extractReleaseTags,
stripQualityTokens,
type SubtitleHit,
type SubtitleSource,
} from "../opensubtitles/client";
import { OpenSubtitlesOrgClient } from "../opensubtitles/org-client";
import { createOpenSubtitlesClientFromDb, isOpenSubtitlesEnabled } from "../settings/opensubtitles";
import type { DownloadLibraryTarget, DownloadSeriesMode, SynologyDsConfig } from "@prisma/client";
import {
folderNameForShow,
parseEpisodesJson,
stringifyEpisodes,
type SeriesEpisodeTrack,
} from "./series-track";
export type LibraryTarget = DownloadLibraryTarget;
export type SeriesMode = DownloadSeriesMode;
const LIBRARY_TARGETS: LibraryTarget[] = ["FILMS", "FOUR_K", "SERIES"];
const SERIES_MODES: SeriesMode[] = ["EPISODE", "SEASON"];
export function isLibraryTarget(v: unknown): v is LibraryTarget {
return typeof v === "string" && LIBRARY_TARGETS.includes(v as LibraryTarget);
}
export function isSeriesMode(v: unknown): v is SeriesMode {
return typeof v === "string" && SERIES_MODES.includes(v as SeriesMode);
}
function normalizeBaseUrl(raw: string): string {
let u = raw.trim().replace(/\/+$/, "");
if (!/^https?:\/\//i.test(u)) u = `https://${u}`;
return u;
}
function normalizeHostPath(p: string): string {
return p.trim().replace(/\/+$/, "");
}
/** Append quality hint only when the user did not already specify one. */
export function searchQueryForTarget(query: string, target: LibraryTarget): string {
const q = query.trim();
if (!q) return q;
if (/\b(2160p|1080p|720p|480p|4k|uhd)\b/i.test(q)) return q;
if (target === "FOUR_K") return `${q} 2160p`;
if (target === "FILMS") return `${q} 1080p`;
return q;
}
function qualityHintForTarget(target: LibraryTarget): string {
if (target === "FOUR_K") return "2160p";
if (target === "FILMS") return "1080p";
return "1080p";
}
function libraryHostPath(config: SynologyDsConfig, target: LibraryTarget): string {
if (target === "FILMS") return normalizeHostPath(config.filmsHostPath || "");
if (target === "SERIES") return normalizeHostPath(config.seriesHostPath || "");
return normalizeHostPath(config.fourKHostPath);
}
function requireLibraryPath(config: SynologyDsConfig, target: LibraryTarget): string {
const path = libraryHostPath(config, target);
if (!path) {
const label =
target === "FILMS" ? "Films (1080)" : target === "SERIES" ? "Series" : "4K";
throw new AppError(
"NOT_CONFIGURED",
`Bibliotheekpad voor ${label} ontbreekt — vul dit in onder Instellingen → Download Station`,
400
);
}
return path;
}
export type DsConfigInput = {
nodeId: string;
baseUrl: string;
username: string;
password?: string;
downloadDestination: string;
downloadHostPath: string;
filmsHostPath: string;
fourKHostPath: string;
seriesHostPath: string;
enabled?: boolean;
};
export class DownloadService {
constructor(
private readonly sessionSecret: string,
private readonly metadata: MetadataService
) {}
private async opensubtitles(): Promise<OpenSubtitlesClient | null> {
return createOpenSubtitlesClientFromDb(this.sessionSecret);
}
private orgClient(): OpenSubtitlesOrgClient {
return new OpenSubtitlesOrgClient();
}
private async downloadSubtitleFile(
fileId: number,
source: SubtitleSource
): Promise<{ content: Buffer; fileName: string; source: SubtitleSource }> {
if (source === "org") {
const file = await this.orgClient().downloadFile(fileId);
return { ...file, source: "org" };
}
const osClient = await this.opensubtitles();
if (!osClient) {
throw new AppError(
"OPENSUBTITLES",
"OpenSubtitles.com niet geconfigureerd — kies een OS.org-hit of vul API-key + login in",
400
);
}
const file = await osClient.downloadFile(fileId);
return { ...file, source: "com" };
}
async hasOpenSubtitles(): Promise<boolean> {
return isOpenSubtitlesEnabled();
}
async searchSubtitles(opts: {
imdbId?: string;
tmdbId?: number;
languages?: string;
releaseHint?: string;
movieTitle?: string;
movieYear?: number | null;
/** Fallback voor OS.org (bijv. serie S01E02) als IMDb ontbreekt of episode-specifiek. */
query?: string;
seasonNumber?: number;
episodeNumber?: number;
mediaKind?: "movie" | "episode";
}): Promise<{
imdbId: string | null;
results: SubtitleHit[];
configured: boolean;
}> {
const enabled = await isOpenSubtitlesEnabled();
if (!enabled) {
return { imdbId: null, results: [], configured: false };
}
const comClient = await this.opensubtitles();
const orgClient = this.orgClient();
const configured = true;
const isEpisode =
opts.mediaKind === "episode" ||
(opts.seasonNumber != null &&
opts.episodeNumber != null &&
opts.episodeNumber > 0);
let imdbId = opts.imdbId?.trim() || null;
if (!imdbId && opts.tmdbId) {
imdbId = isEpisode
? await this.metadata.resolveSeriesImdbId(opts.tmdbId)
: await this.metadata.resolveMovieImdbId(opts.tmdbId);
}
const languages = opts.languages?.trim() || "nl";
let query = opts.query?.trim() || undefined;
if (
!query &&
isEpisode &&
opts.movieTitle &&
opts.seasonNumber != null &&
opts.episodeNumber != null
) {
query = `${opts.movieTitle} S${String(opts.seasonNumber).padStart(2, "0")}E${String(opts.episodeNumber).padStart(2, "0")}`;
}
// Zonder IMDb: alleen OS.org query (als aanwezig)
if (!imdbId && !query && !(isEpisode && opts.tmdbId)) {
return { imdbId: null, results: [], configured };
}
const tasks: Array<Promise<SubtitleHit[]>> = [];
const safe = (label: string, p: Promise<SubtitleHit[]>) =>
p.catch((err) => {
console.warn(`[downloads] subtitle search ${label}:`, err instanceof Error ? err.message : err);
return [] as SubtitleHit[];
});
if (comClient && (imdbId || (isEpisode && opts.tmdbId))) {
tasks.push(
safe(
"com",
comClient.search({
imdbId: imdbId || undefined,
languages,
type: isEpisode ? "episode" : "movie",
seasonNumber: opts.seasonNumber,
episodeNumber: opts.episodeNumber,
parentTmdbId: isEpisode ? opts.tmdbId : undefined,
})
)
);
}
if (imdbId) {
tasks.push(safe("org/imdb", orgClient.search({ imdbId, languages })));
}
if (query) {
tasks.push(safe("org/query", orgClient.search({ query, languages })));
}
if (tasks.length === 0) {
return { imdbId, results: [], configured };
}
const batches = await Promise.all(tasks);
const merged = mergeSubtitleHits(batches.flat());
const hintRaw = opts.releaseHint || "";
const hint = stripQualityTokens(hintRaw).toLowerCase();
const scored = merged
.map((r) => {
let score = r.downloadCount;
if (r.language === "nl" || r.language.startsWith("nl") || r.language === "dut") {
score += 1_000_000;
}
if (r.source === "org") score += 8;
if (r.source === "com") score += 5;
if (hint) {
const hay = `${r.release} ${r.fileName}`.toLowerCase();
const tokens = hint.split(/[^a-z0-9]+/).filter((t) => t.length > 2);
for (const t of tokens) {
if (hay.includes(t)) score += 50;
}
}
return { r, score };
})
.sort((a, b) => b.score - a.score)
.map((x) => x.r);
return { imdbId, results: scored.slice(0, 40), configured };
}
/** TV/client: resolve IMDb (+ episode query) from media file, then search OS.com/OS.org. */
async searchSubtitlesForMediaFile(mediaFileId: string, languages = "nl,en") {
const file = await prisma.mediaFile.findUnique({
where: { id: mediaFileId },
include: {
movie: { select: { imdbId: true, tmdbId: true, title: true, year: true } },
episode: {
select: {
imdbId: true,
tmdbId: true,
title: true,
seasonNumber: true,
episodeNumber: true,
series: { select: { imdbId: true, tmdbId: true, title: true, year: true } },
},
},
},
});
if (!file) throw new AppError("NOT_FOUND", "Mediabestand niet gevonden", 404);
const imdbId =
file.movie?.imdbId ||
file.episode?.imdbId ||
file.episode?.series?.imdbId ||
undefined;
const tmdbId = file.movie?.tmdbId || file.episode?.tmdbId || file.episode?.series?.tmdbId || undefined;
const movieTitle = file.movie?.title || file.episode?.series?.title || undefined;
const movieYear = file.movie?.year ?? file.episode?.series?.year ?? null;
let query: string | undefined;
if (file.episode && !file.episode.imdbId) {
const s = String(file.episode.seasonNumber).padStart(2, "0");
const e = String(file.episode.episodeNumber).padStart(2, "0");
const name = file.episode.series.title;
query = `${name} S${s}E${e}`;
}
const base = await this.searchSubtitles({
imdbId,
tmdbId: tmdbId ?? undefined,
languages,
releaseHint: file.releaseName,
movieTitle,
movieYear,
query,
});
return {
...base,
title: movieTitle || null,
seasonNumber: file.episode?.seasonNumber ?? null,
episodeNumber: file.episode?.episodeNumber ?? null,
};
}
async downloadSubtitleForClient(fileId: number, source: SubtitleSource) {
if (!Number.isFinite(fileId) || fileId <= 0) {
throw new AppError("INVALID_REQUEST", "fileId ongeldig", 400);
}
const src: SubtitleSource = source === "com" ? "com" : "org";
const file = await this.downloadSubtitleFile(fileId, src);
const lower = file.fileName.toLowerCase();
const format = lower.endsWith(".vtt")
? "vtt"
: lower.endsWith(".ass") || lower.endsWith(".ssa")
? "ass"
: "srt";
return {
fileId,
source: src,
fileName: file.fileName,
format,
contentBase64: file.content.toString("base64"),
};
}
private async clientForNode(nodeId: string): Promise<{
client: SynologyDownloadStation;
config: Awaited<ReturnType<typeof prisma.synologyDsConfig.findUniqueOrThrow>>;
}> {
const config = await prisma.synologyDsConfig.findUnique({ where: { nodeId } });
if (!config || !config.enabled) {
throw new AppError("NOT_CONFIGURED", "Download Station niet geconfigureerd voor deze node", 400);
}
const password = decryptSecret(config.passwordEnc, this.sessionSecret);
const client = new SynologyDownloadStation(config.baseUrl, config.username, password);
return { client, config };
}
async listConfigs() {
const rows = await prisma.synologyDsConfig.findMany({
include: { node: { select: { id: true, name: true, status: true } } },
orderBy: { updatedAt: "desc" },
});
return rows.map((r) => ({
id: r.id,
nodeId: r.nodeId,
nodeName: r.node.name,
nodeStatus: r.node.status,
baseUrl: r.baseUrl,
username: r.username,
hasPassword: true,
downloadDestination: r.downloadDestination,
downloadHostPath: r.downloadHostPath,
filmsHostPath: r.filmsHostPath,
fourKHostPath: r.fourKHostPath,
seriesHostPath: r.seriesHostPath,
enabled: r.enabled,
lastOkAt: r.lastOkAt,
lastError: r.lastError,
}));
}
async upsertConfig(input: DsConfigInput) {
const node = await prisma.node.findUnique({ where: { id: input.nodeId } });
if (!node) throw new AppError("NOT_FOUND", "Node niet gevonden", 404);
const existing = await prisma.synologyDsConfig.findUnique({ where: { nodeId: input.nodeId } });
if (!input.password?.trim() && !existing) {
throw new AppError("INVALID_REQUEST", "Wachtwoord verplicht voor nieuwe configuratie", 400);
}
const baseUrl = normalizeBaseUrl(input.baseUrl);
const username = input.username.trim();
const downloadDestination = input.downloadDestination.trim().replace(/^\/+/, "").replace(/\/+$/, "");
const downloadHostPath = normalizeHostPath(input.downloadHostPath);
const filmsHostPath = normalizeHostPath(input.filmsHostPath);
const fourKHostPath = normalizeHostPath(input.fourKHostPath);
const seriesHostPath = normalizeHostPath(input.seriesHostPath);
if (
!username ||
!downloadDestination ||
!downloadHostPath ||
!filmsHostPath ||
!fourKHostPath ||
!seriesHostPath
) {
throw new AppError(
"INVALID_REQUEST",
"Vul alle verplichte velden in (inclusief Films-, 4K- en Series-pad)",
400
);
}
const passwordEnc = input.password?.trim()
? encryptSecret(input.password.trim(), this.sessionSecret)
: existing!.passwordEnc;
const password = decryptSecret(passwordEnc, this.sessionSecret);
const client = new SynologyDownloadStation(baseUrl, username, password);
try {
await client.testConnection();
} catch (err) {
const message = err instanceof Error ? err.message : "Verbinding mislukt";
throw new AppError("DS_ERROR", message, 400);
}
const row = await prisma.synologyDsConfig.upsert({
where: { nodeId: input.nodeId },
create: {
nodeId: input.nodeId,
baseUrl,
username,
passwordEnc,
downloadDestination,
downloadHostPath,
filmsHostPath,
fourKHostPath,
seriesHostPath,
enabled: input.enabled ?? true,
lastOkAt: new Date(),
lastError: null,
},
update: {
baseUrl,
username,
passwordEnc,
downloadDestination,
downloadHostPath,
filmsHostPath,
fourKHostPath,
seriesHostPath,
enabled: input.enabled ?? true,
lastOkAt: new Date(),
lastError: null,
},
});
return {
id: row.id,
nodeId: row.nodeId,
baseUrl: row.baseUrl,
username: row.username,
downloadDestination: row.downloadDestination,
downloadHostPath: row.downloadHostPath,
filmsHostPath: row.filmsHostPath,
fourKHostPath: row.fourKHostPath,
seriesHostPath: row.seriesHostPath,
enabled: row.enabled,
lastOkAt: row.lastOkAt,
};
}
async testConfig(nodeId: string) {
const { client } = await this.clientForNode(nodeId);
await client.testConnection();
await prisma.synologyDsConfig.update({
where: { nodeId },
data: { lastOkAt: new Date(), lastError: null },
});
return { ok: true };
}
async search(
nodeId: string,
query: string,
target: LibraryTarget = "FOUR_K"
): Promise<
Array<{
title: string;
size: number;
seeds: number;
peers: number;
uri: string | null;
module: string | null;
suggestQuery: string;
}>
> {
const { client, config } = await this.clientForNode(nodeId);
requireLibraryPath(config, target);
const searchQ = searchQueryForTarget(query, target);
let hits: DsSearchHit[];
try {
hits = await client.search(searchQ);
await prisma.synologyDsConfig.update({
where: { nodeId },
data: { lastOkAt: new Date(), lastError: null },
});
} catch (err) {
const message = err instanceof Error ? err.message : "Zoeken mislukt";
await prisma.synologyDsConfig.update({
where: { nodeId },
data: { lastError: message },
});
throw new AppError("DS_ERROR", message, 400);
}
return hits.map((h) => {
const uri = h.download_uri || h.magnet || null;
return {
title: h.title,
size: h.size ?? 0,
seeds: h.seeds ?? 0,
peers: h.peers ?? h.leechs ?? 0,
uri,
module: h.module_title || h.module_id || null,
suggestQuery: cleanReleaseQuery(h.title),
};
});
}
async createJob(input: {
nodeId: string;
uri: string;
sourceTitle: string;
libraryTarget: LibraryTarget;
seriesMode?: SeriesMode | null;
seasonNumber?: number | null;
episodeNumber?: number | null;
episodeTitle?: string | null;
episodes?: Array<{ episodeNumber: number; episodeTitle?: string }>;
tmdbId?: number;
imdbId?: string;
movieTitle: string;
movieYear?: number | null;
posterUrl?: string | null;
subtitleFileId?: number | null;
subtitleLang?: string | null;
subtitleSource?: SubtitleSource | null;
}) {
const uri = input.uri.trim();
if (!uri) throw new AppError("INVALID_REQUEST", "Download-URI ontbreekt", 400);
if (!input.movieTitle.trim()) {
throw new AppError("INVALID_REQUEST", "Bevestig eerst de titel", 400);
}
if (!input.tmdbId && !input.imdbId) {
throw new AppError("INVALID_REQUEST", "Bevestig met TMDB- of IMDb-id", 400);
}
if (!isLibraryTarget(input.libraryTarget)) {
throw new AppError("INVALID_REQUEST", "Ongeldig bibliotheekdoel", 400);
}
const isSeries = input.libraryTarget === "SERIES";
let seriesMode: SeriesMode | null = null;
if (isSeries) {
if (!isSeriesMode(input.seriesMode)) {
throw new AppError("INVALID_REQUEST", "seriesMode verplicht (EPISODE of SEASON)", 400);
}
seriesMode = input.seriesMode;
if (!input.seasonNumber || input.seasonNumber < 0) {
throw new AppError("INVALID_REQUEST", "Seizoen verplicht", 400);
}
if (seriesMode === "EPISODE" && (!input.episodeNumber || input.episodeNumber <= 0)) {
throw new AppError("INVALID_REQUEST", "Aflevering verplicht", 400);
}
}
const { client, config } = await this.clientForNode(input.nodeId);
requireLibraryPath(config, input.libraryTarget);
let tmdbId = input.tmdbId;
let imdbId = input.imdbId?.trim() || null;
let title = input.movieTitle.trim();
let year = input.movieYear ?? null;
let posterUrl = input.posterUrl ?? null;
if (!tmdbId && imdbId) {
const hits = await this.metadata.searchCatalog(isSeries ? "series" : "movie", imdbId);
const hit = hits[0];
if (hit) {
tmdbId = hit.tmdbId;
if (hit.imdbId) imdbId = hit.imdbId;
title = hit.title;
year = hit.year;
posterUrl = hit.posterUrl;
}
}
if (tmdbId && !imdbId) {
imdbId = isSeries
? await this.metadata.resolveSeriesImdbId(tmdbId)
: await this.metadata.resolveMovieImdbId(tmdbId);
}
if (tmdbId) {
try {
if (isSeries) {
await prisma.series.upsert({
where: { tmdbId },
create: {
title,
year,
posterUrl,
tmdbId,
imdbId,
},
update: {
title,
year: year ?? undefined,
posterUrl: posterUrl ?? undefined,
...(imdbId ? { imdbId } : {}),
},
});
} else {
await prisma.movie.upsert({
where: { tmdbId },
create: {
title,
year,
posterUrl,
tmdbId,
imdbId,
},
update: {
title,
year: year ?? undefined,
posterUrl: posterUrl ?? undefined,
...(imdbId ? { imdbId } : {}),
},
});
}
} catch (err) {
console.warn("[downloads] catalog upsert:", err);
}
}
let episodeTracks: SeriesEpisodeTrack[] = [];
let episodeNumber = input.episodeNumber ?? null;
let episodeTitle = input.episodeTitle?.trim() || null;
const seasonNumber = isSeries ? input.seasonNumber! : null;
if (isSeries && seriesMode === "SEASON" && tmdbId) {
const season = await this.metadata.getTvSeason(tmdbId, seasonNumber!);
const fromMeta = season?.episodes ?? [];
const picked =
input.episodes && input.episodes.length > 0
? input.episodes
: fromMeta.map((e) => ({ episodeNumber: e.episodeNumber, episodeTitle: e.title }));
if (picked.length === 0) {
throw new AppError("INVALID_REQUEST", "Geen afleveringen voor dit seizoen", 400);
}
episodeTracks = picked.map((e) => ({
episodeNumber: e.episodeNumber,
episodeTitle: e.episodeTitle || fromMeta.find((m) => m.episodeNumber === e.episodeNumber)?.title || "",
importedPath: null,
subtitlePath: null,
subtitleError: null,
}));
} else if (isSeries && seriesMode === "EPISODE") {
episodeTracks = [
{
episodeNumber: episodeNumber!,
episodeTitle: episodeTitle || "",
importedPath: null,
subtitlePath: null,
subtitleError: null,
},
];
}
const destFolderName = isSeries
? folderNameForShow(title, year)
: folderNameForMovie(title, year);
const releaseTags = extractReleaseTags(input.sourceTitle, title, year) || null;
// Season packs: no single subtitle at create. Episode: optional.
const allowSubtitle = !isSeries || seriesMode === "EPISODE";
const subtitleFileId =
allowSubtitle && input.subtitleFileId != null && Number.isFinite(input.subtitleFileId)
? Math.floor(input.subtitleFileId)
: null;
const subtitleLang = subtitleFileId
? (input.subtitleLang?.trim().toLowerCase() || "nl")
: null;
const subtitleSource: SubtitleSource | null = subtitleFileId
? input.subtitleSource === "org"
? "org"
: "com"
: null;
let subtitleContent: string | null = null;
let subtitleFileName: string | null = null;
if (subtitleFileId && subtitleSource) {
try {
const file = await this.downloadSubtitleFile(subtitleFileId, subtitleSource);
subtitleContent = file.content.toString("utf8");
subtitleFileName = file.fileName;
if (!subtitleContent.trim()) {
throw new AppError("OPENSUBTITLES", "Ondertitelbestand is leeg", 400);
}
} catch (err) {
const message =
err instanceof AppError
? err.message
: err instanceof Error
? err.message
: "Ondertitel downloaden mislukt";
throw new AppError("OPENSUBTITLES", message, 400);
}
}
let beforeIds = new Set<string>();
try {
const before = await client.listTasks();
beforeIds = new Set(before.map((t) => t.id));
} catch {
/* optional */
}
try {
await client.createTask(uri, config.downloadDestination);
} catch (err) {
const message = err instanceof Error ? err.message : "Taak aanmaken mislukt";
throw new AppError("DS_ERROR", message, 400);
}
let dsTaskId: string | null = null;
try {
await sleep(800);
const after = await client.listTasks();
const created = after.find((t) => !beforeIds.has(t.id) && uriLikelyMatch(t, uri, input.sourceTitle));
const fallback = after.find((t) => !beforeIds.has(t.id));
dsTaskId = (created ?? fallback)?.id ?? null;
} catch {
/* keep null — poller can still match by title */
}
const job = await prisma.downloadJob.create({
data: {
nodeId: input.nodeId,
status: "DOWNLOADING",
libraryTarget: input.libraryTarget,
seriesMode,
sourceUri: uri,
sourceTitle: input.sourceTitle.trim() || title,
dsTaskId,
tmdbId: tmdbId ?? null,
imdbId,
movieTitle: title,
movieYear: year,
posterUrl,
seasonNumber,
episodeNumber,
episodeTitle,
episodesJson: episodeTracks.length ? stringifyEpisodes(episodeTracks) : null,
destFolderName,
releaseTags,
subtitleFileId,
subtitleLang,
subtitleSource,
subtitleContent,
subtitleError: null,
subtitlePath: null,
},
});
if (subtitleContent && !isSeries) {
try {
let movieId: string | null = null;
if (tmdbId) {
const movie = await prisma.movie.findUnique({ where: { tmdbId } });
movieId = movie?.id ?? null;
}
await prisma.librarySubtitle.create({
data: {
movieId,
tmdbId: tmdbId ?? null,
imdbId,
language: subtitleLang || "nl",
format: subtitleFileName?.toLowerCase().endsWith(".ass") ? "ass" : "srt",
content: subtitleContent,
fileName: subtitleFileName,
downloadJobId: job.id,
},
});
} catch (err) {
console.warn("[downloads] early library subtitle:", err);
}
}
return job;
}
async listJobs(limit = 50) {
const jobs = await prisma.downloadJob.findMany({
take: Math.min(200, Math.max(1, limit)),
orderBy: { createdAt: "desc" },
include: { node: { select: { id: true, name: true } } },
});
const downloading = jobs.filter((j) => j.status === "DOWNLOADING");
const progressByJobId = new Map<string, JobProgress>();
const byNode = new Map<string, typeof downloading>();
for (const job of downloading) {
const list = byNode.get(job.nodeId) ?? [];
list.push(job);
byNode.set(job.nodeId, list);
}
for (const [nodeId, nodeJobs] of byNode) {
try {
const { client } = await this.clientForNode(nodeId);
const tasks = await client.listTasks();
const byId = new Map(tasks.map((t) => [t.id, t]));
for (const job of nodeJobs) {
const task = resolveTask(job.dsTaskId, job.sourceTitle, job.sourceUri, byId, tasks);
if (!task) continue;
progressByJobId.set(job.id, progressFromTask(task));
if (!job.dsTaskId && task.id) {
void prisma.downloadJob
.update({ where: { id: job.id }, data: { dsTaskId: task.id } })
.catch(() => undefined);
}
}
} catch (err) {
console.warn(
`[downloads] live progress node ${nodeId}:`,
err instanceof Error ? err.message : err
);
}
}
return jobs.map((j) => ({
...j,
progress: progressByJobId.get(j.id) ?? null,
}));
}
async deleteJob(jobId: string) {
const job = await prisma.downloadJob.findUnique({ where: { id: jobId } });
if (!job) throw new AppError("NOT_FOUND", "Job niet gevonden", 404);
if (job.status !== "IMPORTED") {
throw new AppError("INVALID_REQUEST", "Alleen voltooide jobs kunnen uit de lijst", 400);
}
await prisma.downloadJob.delete({ where: { id: jobId } });
return { ok: true };
}
async clearCompletedJobs() {
const result = await prisma.downloadJob.deleteMany({
where: { status: "IMPORTED" },
});
return { deleted: result.count };
}
async retryImport(jobId: string) {
const job = await prisma.downloadJob.findUnique({ where: { id: jobId } });
if (!job) throw new AppError("NOT_FOUND", "Job niet gevonden", 404);
if (!["FINISHED", "FAILED", "IMPORTING"].includes(job.status)) {
throw new AppError("INVALID_REQUEST", "Job is niet klaar voor import", 400);
}
await prisma.downloadJob.update({
where: { id: jobId },
data: { status: "FINISHED", error: null },
});
await this.importJob(jobId);
return prisma.downloadJob.findUniqueOrThrow({ where: { id: jobId } });
}
/** Poll active DS tasks and kick off imports when finished. */
async pollOnce(): Promise<void> {
const configs = await prisma.synologyDsConfig.findMany({ where: { enabled: true } });
for (const config of configs) {
try {
await this.pollNode(config.nodeId);
} catch (err) {
const message = err instanceof Error ? err.message : "Poll mislukt";
console.warn(`[downloads] poll node ${config.nodeId}:`, message);
await prisma.synologyDsConfig.update({
where: { id: config.id },
data: { lastError: message },
});
}
}
}
private async pollNode(nodeId: string): Promise<void> {
const active = await prisma.downloadJob.findMany({
where: { nodeId, status: { in: ["DOWNLOADING", "FINISHED", "IMPORTING"] } },
});
if (active.length === 0) return;
const { client, config } = await this.clientForNode(nodeId);
const tasks = await client.listTasks();
const byId = new Map(tasks.map((t) => [t.id, t]));
for (const job of active) {
if (job.status === "IMPORTING") continue;
if (job.status === "FINISHED") {
await this.importJob(job.id);
continue;
}
const task = resolveTask(job.dsTaskId, job.sourceTitle, job.sourceUri, byId, tasks);
if (!task) continue;
if (!job.dsTaskId) {
await prisma.downloadJob.update({
where: { id: job.id },
data: { dsTaskId: task.id },
});
}
if (task.status === "error") {
await prisma.downloadJob.update({
where: { id: job.id },
data: {
status: "FAILED",
error: task.status_extra?.error_detail || "Download Station: error",
},
});
continue;
}
// Keep polling for finish; progress is read live via listJobs.
if (isTaskDone(task.status)) {
const sourcePath = guessSourcePath(config.downloadHostPath, task, job.sourceTitle);
await prisma.downloadJob.update({
where: { id: job.id },
data: {
status: "FINISHED",
sourcePath,
error: null,
},
});
await this.importJob(job.id);
}
}
await prisma.synologyDsConfig.update({
where: { nodeId },
data: { lastOkAt: new Date(), lastError: null },
});
}
async importJob(jobId: string): Promise<void> {
const job = await prisma.downloadJob.findUnique({ where: { id: jobId } });
if (!job) return;
if (job.status === "IMPORTED" || job.status === "IMPORTING") return;
const config = await prisma.synologyDsConfig.findUnique({ where: { nodeId: job.nodeId } });
if (!config) {
await prisma.downloadJob.update({
where: { id: jobId },
data: { status: "FAILED", error: "Geen DS-config" },
});
return;
}
const isSeries = job.libraryTarget === "SERIES";
const destFolder =
job.destFolderName ||
(isSeries
? folderNameForShow(job.movieTitle, job.movieYear)
: folderNameForMovie(job.movieTitle, job.movieYear));
let destRoot: string;
try {
destRoot = requireLibraryPath(config, job.libraryTarget);
} catch (err) {
const message = err instanceof AppError ? err.message : "Bibliotheekpad ontbreekt";
await prisma.downloadJob.update({
where: { id: jobId },
data: { status: "FAILED", error: message },
});
return;
}
const destDir = `${destRoot}/${destFolder}`;
const sourcePath =
job.sourcePath || `${normalizeHostPath(config.downloadHostPath)}/${job.sourceTitle}`;
if (!nodeConnectionManager.isOnline(job.nodeId)) {
await prisma.downloadJob.update({
where: { id: jobId },
data: { status: "FINISHED", error: "Node offline — import wacht" },
});
return;
}
await prisma.downloadJob.update({
where: { id: jobId },
data: { status: "IMPORTING", sourcePath, destFolderName: destFolder, error: null },
});
let tracks = parseEpisodesJson(job.episodesJson);
if (isSeries && tracks.length === 0 && job.episodeNumber) {
tracks = [
{
episodeNumber: job.episodeNumber,
episodeTitle: job.episodeTitle || "",
importedPath: null,
subtitlePath: null,
subtitleError: null,
},
];
}
let subtitle:
| { language: string; extension?: string; contentBase64: string }
| undefined;
let subtitleText: string | null = job.subtitleContent;
let subtitleFileName: string | null = null;
let subtitleError: string | null = null;
// Only attach subtitle on import for movies / single-episode series
const attachSubtitleNow = !isSeries || job.seriesMode === "EPISODE";
if (attachSubtitleNow && (job.subtitleFileId || job.subtitleContent)) {
try {
if (!subtitleText?.trim() && job.subtitleFileId) {
const source = job.subtitleSource === "org" ? "org" : "com";
const file = await this.downloadSubtitleFile(job.subtitleFileId, source);
subtitleText = file.content.toString("utf8");
subtitleFileName = file.fileName;
await prisma.downloadJob.update({
where: { id: jobId },
data: { subtitleContent: subtitleText, subtitleError: null },
});
}
if (!subtitleText?.trim()) {
throw new Error("Geen ondertitelinhoud beschikbaar");
}
const ext = subtitleFileName?.toLowerCase().endsWith(".ass") ? ".ass" : ".srt";
subtitle = {
language: job.subtitleLang || "nl",
extension: ext,
contentBase64: Buffer.from(subtitleText, "utf8").toString("base64"),
};
} catch (err) {
subtitleError =
err instanceof Error ? err.message : "Subtitle download mislukt";
console.warn(`[downloads] subtitle job=${jobId}:`, subtitleError);
}
}
const releaseTags =
job.releaseTags ||
extractReleaseTags(job.sourceTitle, job.movieTitle, job.movieYear) ||
undefined;
const ack = await nodeConnectionManager.importMediaAcked(
job.nodeId,
{
jobId: job.id,
sourcePath,
destDir,
movieTitle: job.movieTitle,
movieYear: job.movieYear,
releaseTags,
qualityHint: qualityHintForTarget(job.libraryTarget),
kind: isSeries ? "series" : "movie",
seasonNumber: isSeries ? job.seasonNumber ?? undefined : undefined,
episodes: isSeries
? tracks.map((t) => ({
episodeNumber: t.episodeNumber,
episodeTitle: t.episodeTitle || undefined,
}))
: undefined,
subtitle: attachSubtitleNow ? subtitle : undefined,
},
isSeries && job.seriesMode === "SEASON" ? 300_000 : 120_000
);
if (!ack.ok) {
await prisma.downloadJob.update({
where: { id: jobId },
data: {
status: "FAILED",
error: ack.error || "Import mislukt of timeout",
subtitleError,
},
});
return;
}
if (isSeries) {
const byEp = new Map((ack.episodes ?? []).map((e) => [e.episodeNumber, e]));
tracks = tracks.map((t) => {
const hit = byEp.get(t.episodeNumber);
if (!hit) return t;
return {
...t,
importedPath: hit.importedPath || t.importedPath,
subtitlePath: hit.subtitlePath || t.subtitlePath,
};
});
// Also add unexpected imported episodes
for (const ep of ack.episodes ?? []) {
if (!tracks.some((t) => t.episodeNumber === ep.episodeNumber)) {
tracks.push({
episodeNumber: ep.episodeNumber,
episodeTitle: "",
importedPath: ep.importedPath,
subtitlePath: ep.subtitlePath || null,
subtitleError: null,
});
}
}
const importedCount = tracks.filter((t) => t.importedPath).length;
const missing = tracks.filter((t) => !t.importedPath).map((t) => `E${String(t.episodeNumber).padStart(2, "0")}`);
if (job.seriesMode === "EPISODE" && subtitle && !ack.subtitlePath && !tracks[0]?.subtitlePath) {
subtitleError =
subtitleError ||
"Aflevering geïmporteerd, maar node schreef geen .srt — update media-node";
if (tracks[0]) tracks[0].subtitleError = subtitleError;
}
let statusNote: string | null = null;
if (importedCount === 0) {
await prisma.downloadJob.update({
where: { id: jobId },
data: {
status: "FAILED",
error: "Geen afleveringen geïmporteerd",
episodesJson: stringifyEpisodes(tracks),
},
});
return;
}
if (missing.length > 0) {
statusNote = `Geïmporteerd ${importedCount}/${tracks.length}; mist: ${missing.join(", ")}`;
} else if (job.seriesMode === "SEASON") {
const needSubs = tracks.filter((t) => t.importedPath && !t.subtitlePath).length;
if (needSubs > 0) {
statusNote = `${importedCount} afleveringen klaar — kies nog ondertitels (${needSubs})`;
}
}
await prisma.downloadJob.update({
where: { id: jobId },
data: {
status: "IMPORTED",
importedPath: ack.importedPath ?? tracks.find((t) => t.importedPath)?.importedPath ?? null,
subtitlePath:
ack.subtitlePath ?? tracks.find((t) => t.subtitlePath)?.subtitlePath ?? null,
subtitleError,
episodesJson: stringifyEpisodes(tracks),
error: statusNote || subtitleError,
},
});
return;
}
if (subtitle && !ack.subtitlePath) {
subtitleError =
subtitleError ||
"Film geïmporteerd, maar node schreef geen .srt — update media-node (1.3.2+) of gebruik ‘Ondertitel opnieuw’";
} else if (!subtitle && job.subtitleFileId) {
subtitleError =
subtitleError || "Ondertitel niet meegestuurd naar node (download mislukt)";
}
await prisma.downloadJob.update({
where: { id: jobId },
data: {
status: "IMPORTED",
importedPath: ack.importedPath ?? null,
subtitlePath: ack.subtitlePath ?? null,
subtitleError,
error: subtitleError,
},
});
if (subtitleText && !subtitleError) {
try {
let movieId: string | null = null;
if (job.tmdbId) {
const movie = await prisma.movie.findUnique({ where: { tmdbId: job.tmdbId } });
movieId = movie?.id ?? null;
}
const existing = await prisma.librarySubtitle.findFirst({
where: { downloadJobId: job.id },
});
if (!existing) {
await prisma.librarySubtitle.create({
data: {
movieId,
tmdbId: job.tmdbId,
imdbId: job.imdbId,
language: job.subtitleLang || "nl",
format: subtitle?.extension?.replace(/^\./, "") || "srt",
content: subtitleText,
fileName: subtitleFileName,
downloadJobId: job.id,
},
});
}
} catch (err) {
console.warn(`[downloads] save library subtitle job=${jobId}:`, err);
}
}
}
/** Attach / replace subtitle for one episode of an imported series job. */
async attachEpisodeSubtitle(
jobId: string,
episodeNumber: number,
opts: {
subtitleFileId: number;
subtitleSource?: SubtitleSource | null;
subtitleLang?: string | null;
}
) {
const job = await prisma.downloadJob.findUnique({ where: { id: jobId } });
if (!job) throw new AppError("NOT_FOUND", "Job niet gevonden", 404);
if (job.libraryTarget !== "SERIES") {
throw new AppError("INVALID_REQUEST", "Alleen voor series-jobs", 400);
}
const tracks = parseEpisodesJson(job.episodesJson);
const track = tracks.find((t) => t.episodeNumber === episodeNumber);
if (!track?.importedPath) {
throw new AppError("INVALID_REQUEST", "Aflevering nog niet geïmporteerd", 400);
}
const source: SubtitleSource = opts.subtitleSource === "org" ? "org" : "com";
const file = await this.downloadSubtitleFile(Math.floor(opts.subtitleFileId), source);
const content = file.content.toString("utf8");
if (!content.trim()) throw new AppError("OPENSUBTITLES", "Ondertitelbestand is leeg", 400);
const ext = file.fileName?.toLowerCase().endsWith(".ass") ? ".ass" : ".srt";
const lang = (opts.subtitleLang || "nl").trim().toLowerCase() || "nl";
const ack = await nodeConnectionManager.writeSubtitleAcked(
job.nodeId,
{
jobId: `${job.id}:E${episodeNumber}`,
videoPath: track.importedPath,
subtitle: {
language: lang,
extension: ext,
contentBase64: Buffer.from(content, "utf8").toString("base64"),
},
},
60_000
);
if (!ack.ok) {
track.subtitleError = ack.error || "Ondertitel schrijven mislukt";
await prisma.downloadJob.update({
where: { id: jobId },
data: { episodesJson: stringifyEpisodes(tracks), error: track.subtitleError },
});
throw new AppError("SUBTITLE", track.subtitleError, 400);
}
track.subtitlePath = ack.subtitlePath || null;
track.subtitleError = null;
const needSubs = tracks.filter((t) => t.importedPath && !t.subtitlePath).length;
await prisma.downloadJob.update({
where: { id: jobId },
data: {
episodesJson: stringifyEpisodes(tracks),
error:
needSubs > 0
? `${tracks.filter((t) => t.importedPath).length} afleveringen — nog ${needSubs} zonder sub`
: null,
},
});
return prisma.downloadJob.findUniqueOrThrow({ where: { id: jobId } });
}
/** Re-download / re-pick and write sidecar next to already-imported video. */
async retrySubtitle(
jobId: string,
opts?: {
subtitleFileId?: number | null;
subtitleSource?: SubtitleSource | null;
subtitleLang?: string | null;
}
) {
const job = await prisma.downloadJob.findUnique({ where: { id: jobId } });
if (!job) throw new AppError("NOT_FOUND", "Job niet gevonden", 404);
if (!job.importedPath) {
throw new AppError("INVALID_REQUEST", "Nog geen geïmporteerd pad — importeer eerst de film", 400);
}
const pickId =
opts?.subtitleFileId != null && Number.isFinite(opts.subtitleFileId)
? Math.floor(opts.subtitleFileId)
: null;
const pickSource: SubtitleSource | null = pickId
? opts?.subtitleSource === "org"
? "org"
: "com"
: null;
let content = pickId ? null : job.subtitleContent;
let fileName: string | null = null;
let source: SubtitleSource | null = pickSource || (job.subtitleSource === "org" ? "org" : job.subtitleFileId ? "com" : null);
let fileId = pickId || job.subtitleFileId;
const lang = (opts?.subtitleLang || job.subtitleLang || "nl").trim().toLowerCase() || "nl";
// Nieuwe keuze of opnieuw downloaden
if (pickId && pickSource) {
const file = await this.downloadSubtitleFile(pickId, pickSource);
content = file.content.toString("utf8");
fileName = file.fileName;
source = pickSource;
fileId = pickId;
} else if (!content?.trim() && fileId && source) {
try {
const file = await this.downloadSubtitleFile(fileId, source);
content = file.content.toString("utf8");
fileName = file.fileName;
} catch (err) {
const message = err instanceof Error ? err.message : "Download mislukt";
// Als .com faalt: forceer opnieuw kiezen via UI (OS.org)
if (source === "com") {
throw new AppError(
"OPENSUBTITLES",
`${message} — kies opnieuw een OS.org-ondertitel bij deze job.`,
400
);
}
throw err instanceof AppError ? err : new AppError("OPENSUBTITLES", message, 400);
}
}
if (!content?.trim()) {
throw new AppError(
"OPENSUBTITLES",
"Geen ondertitelinhoud — kies een OS.org-hit bij deze job",
400
);
}
await prisma.downloadJob.update({
where: { id: jobId },
data: {
subtitleFileId: fileId,
subtitleSource: source,
subtitleLang: lang,
subtitleContent: content,
subtitleError: null,
},
});
const ext = fileName?.toLowerCase().endsWith(".ass") ? ".ass" : ".srt";
const ack = await nodeConnectionManager.writeSubtitleAcked(
job.nodeId,
{
jobId: job.id,
videoPath: job.importedPath,
subtitle: {
language: lang,
extension: ext,
contentBase64: Buffer.from(content, "utf8").toString("base64"),
},
},
60_000
);
if (!ack.ok || !ack.subtitlePath) {
const message = ack.error || "Node schreef geen ondertitel (update media-node?)";
await prisma.downloadJob.update({
where: { id: jobId },
data: { subtitleError: message, error: message },
});
throw new AppError("SUBTITLE", message, 400);
}
await prisma.downloadJob.update({
where: { id: jobId },
data: {
subtitlePath: ack.subtitlePath,
subtitleError: null,
error: null,
},
});
try {
let movieId: string | null = null;
if (job.tmdbId) {
const movie = await prisma.movie.findUnique({ where: { tmdbId: job.tmdbId } });
movieId = movie?.id ?? null;
}
const existing = await prisma.librarySubtitle.findFirst({
where: { downloadJobId: job.id },
});
if (existing) {
await prisma.librarySubtitle.update({
where: { id: existing.id },
data: {
content,
fileName,
language: lang,
format: ext.replace(/^\./, ""),
},
});
} else {
await prisma.librarySubtitle.create({
data: {
movieId,
tmdbId: job.tmdbId,
imdbId: job.imdbId,
language: lang,
format: ext.replace(/^\./, ""),
content,
fileName,
downloadJobId: job.id,
},
});
}
} catch (err) {
console.warn(`[downloads] library subtitle retry job=${jobId}:`, err);
}
return prisma.downloadJob.findUniqueOrThrow({
where: { id: jobId },
include: { node: { select: { id: true, name: true } } },
});
}
}
function mergeSubtitleHits(hits: SubtitleHit[]): SubtitleHit[] {
// Bewaar hits van .com én .org apart (andere catalogs), dedupe alleen binnen dezelfde bron.
const byId = new Map<string, SubtitleHit>();
for (const hit of hits) {
const idKey = `${hit.source}:${hit.fileId}`;
const prev = byId.get(idKey);
if (!prev || hit.downloadCount > prev.downloadCount) {
byId.set(idKey, hit);
}
}
return [...byId.values()];
}
function isTaskDone(status: string): boolean {
return status === "finished" || status === "seeding";
}
export type JobProgress = {
percent: number | null;
bytesDownloaded: number | null;
bytesTotal: number | null;
/** Bytes per second from Download Station */
speedBps: number | null;
dsStatus: string | null;
};
function progressFromTask(task: DsTask): JobProgress {
const size = typeof task.size === "number" && task.size > 0 ? task.size : 0;
const downloaded = task.additional?.transfer?.size_downloaded ?? 0;
const speed = task.additional?.transfer?.speed_download ?? 0;
let percent: number | null = null;
if (size > 0) {
percent = Math.min(100, Math.round((downloaded / size) * 1000) / 10);
} else if (isTaskDone(task.status) || task.status === "finishing") {
percent = 100;
}
return {
percent,
bytesDownloaded: downloaded > 0 ? downloaded : downloaded === 0 && size > 0 ? 0 : null,
bytesTotal: size > 0 ? size : null,
speedBps: speed > 0 ? speed : speed === 0 ? 0 : null,
dsStatus: task.status || null,
};
}
function uriLikelyMatch(task: DsTask, uri: string, sourceTitle: string): boolean {
const detailUri = task.additional?.detail?.uri || "";
if (uri && detailUri && (detailUri === uri || detailUri.includes(uri.slice(0, 48)))) return true;
if (sourceTitle && task.title && titlesSimilar(task.title, sourceTitle)) return true;
return false;
}
function resolveTask(
dsTaskId: string | null,
sourceTitle: string,
sourceUri: string,
byId: Map<string, DsTask>,
tasks: DsTask[]
): DsTask | undefined {
if (dsTaskId && byId.has(dsTaskId)) return byId.get(dsTaskId);
return tasks.find((t) => uriLikelyMatch(t, sourceUri, sourceTitle));
}
function titlesSimilar(a: string, b: string): boolean {
const na = a.toLowerCase().replace(/[^a-z0-9]+/g, "");
const nb = b.toLowerCase().replace(/[^a-z0-9]+/g, "");
if (!na || !nb) return false;
return na.includes(nb.slice(0, Math.min(24, nb.length))) || nb.includes(na.slice(0, Math.min(24, na.length)));
}
function guessSourcePath(downloadHostPath: string, task: DsTask, sourceTitle: string): string {
const root = normalizeHostPath(downloadHostPath);
// Prefer first file path if DS returns nested names under destination
const files = task.additional?.file ?? [];
if (files.length === 1 && files[0]?.filename && !files[0].filename.includes("/")) {
// Single file at destination root — source is the file itself
return `${root}/${files[0].filename}`;
}
// Multi-file / folder: use task title as folder name (common DS behaviour)
const folder = (task.title || sourceTitle).replace(/[<>:"/\\|?*\x00-\x1f]/g, "").trim();
return `${root}/${folder}`;
}
function sleep(ms: number): Promise<void> {
return new Promise((r) => setTimeout(r, ms));
}
let pollerStarted = false;
export function startDownloadPoller(service: DownloadService, intervalMs = 10_000): void {
if (pollerStarted) return;
pollerStarted = true;
const tick = () => {
void service.pollOnce().catch((err) => console.warn("[downloads] poll:", err));
};
setTimeout(tick, 8_000);
setInterval(tick, intervalMs);
}