Add Synology Download Station search and 4K import pipeline.

Admin can search DS, confirm TMDB/IMDb, then auto-move finished downloads into a fixed 4K library folder via the media-node.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
Jos Vooges | STH 2026-08-30 17:42:12 +02:00
parent 756b2adda8
commit 18eb702de2
17 changed files with 2202 additions and 17 deletions

View file

@ -0,0 +1,590 @@
"use client";
import { useCallback, useEffect, useMemo, useState } from "react";
import { Nav, useAuth } from "@/components/Nav";
interface NodeRow {
id: string;
name: string;
status: string;
}
interface DsConfig {
id?: string;
nodeId: string;
nodeName?: string;
baseUrl: string;
username: string;
downloadDestination: string;
downloadHostPath: string;
fourKHostPath: string;
enabled: boolean;
lastOkAt?: string | null;
lastError?: string | null;
}
interface SearchHit {
title: string;
size: number;
seeds: number;
peers: number;
uri: string | null;
module: string | null;
suggestQuery: string;
}
interface MetaHit {
tmdbId: number;
title: string;
year: number | null;
overview: string | null;
posterUrl: string | null;
imdbId: string | null;
}
interface JobRow {
id: string;
nodeId: string;
nodeName: string;
status: string;
sourceTitle: string;
movieTitle: string;
movieYear: number | null;
posterUrl: string | null;
tmdbId: number | null;
imdbId: string | null;
destFolderName: string | null;
importedPath: string | null;
error: string | null;
createdAt: string;
updatedAt: string;
}
function formatBytes(n: number): string {
if (!n || n <= 0) return "—";
const u = ["B", "KB", "MB", "GB", "TB"];
let i = 0;
let v = n;
while (v >= 1024 && i < u.length - 1) {
v /= 1024;
i++;
}
return `${v.toFixed(i === 0 ? 0 : 1)} ${u[i]}`;
}
function statusClass(s: string): string {
switch (s) {
case "IMPORTED":
return "badge badge-ok";
case "DOWNLOADING":
case "IMPORTING":
case "FINISHED":
return "badge badge-warn";
case "FAILED":
return "badge badge-danger";
default:
return "badge";
}
}
const emptyConfig = (nodeId = ""): DsConfig => ({
nodeId,
baseUrl: "",
username: "",
downloadDestination: "Downloads",
downloadHostPath: "/volume1/Downloads",
fourKHostPath: "/volume1/Media/Films/4K",
enabled: true,
});
export default function DownloadsPage() {
useAuth();
const [nodes, setNodes] = useState<NodeRow[]>([]);
const [configs, setConfigs] = useState<DsConfig[]>([]);
const [nodeId, setNodeId] = useState("");
const [cfg, setCfg] = useState<DsConfig>(emptyConfig());
const [password, setPassword] = useState("");
const [cfgMsg, setCfgMsg] = useState<string | null>(null);
const [cfgErr, setCfgErr] = useState<string | null>(null);
const [cfgBusy, setCfgBusy] = useState(false);
const [query, setQuery] = useState("");
const [searching, setSearching] = useState(false);
const [hits, setHits] = useState<SearchHit[]>([]);
const [searchErr, setSearchErr] = useState<string | null>(null);
const [pick, setPick] = useState<SearchHit | null>(null);
const [metaQ, setMetaQ] = useState("");
const [metaHits, setMetaHits] = useState<MetaHit[]>([]);
const [metaBusy, setMetaBusy] = useState(false);
const [selectedMeta, setSelectedMeta] = useState<MetaHit | null>(null);
const [jobBusy, setJobBusy] = useState(false);
const [jobErr, setJobErr] = useState<string | null>(null);
const [jobs, setJobs] = useState<JobRow[]>([]);
const loadJobs = useCallback(() => {
fetch("/api/v1/admin/downloads/jobs", { credentials: "include" })
.then((r) => r.json())
.then((d) => setJobs(d.jobs ?? []))
.catch(() => undefined);
}, []);
const loadAll = useCallback(() => {
Promise.all([
fetch("/api/v1/admin/nodes", { credentials: "include" }).then((r) => r.json()),
fetch("/api/v1/admin/downloads/configs", { credentials: "include" }).then((r) => r.json()),
]).then(([n, c]) => {
const nodeList: NodeRow[] = (n.nodes ?? []).map(
(x: { id: string; name: string; status: string }) => ({
id: x.id,
name: x.name,
status: x.status,
})
);
setNodes(nodeList);
const list: DsConfig[] = c.configs ?? [];
setConfigs(list);
const preferred = list[0]?.nodeId || nodeList[0]?.id || "";
setNodeId((prev) => prev || preferred);
});
loadJobs();
}, [loadJobs]);
useEffect(() => {
loadAll();
const t = setInterval(loadJobs, 15_000);
return () => clearInterval(t);
}, [loadAll, loadJobs]);
useEffect(() => {
if (!nodeId) return;
const existing = configs.find((c) => c.nodeId === nodeId);
setCfg(existing ? { ...existing } : emptyConfig(nodeId));
setPassword("");
setCfgMsg(null);
setCfgErr(null);
}, [nodeId, configs]);
const hasConfig = useMemo(() => configs.some((c) => c.nodeId === nodeId), [configs, nodeId]);
async function saveConfig() {
setCfgBusy(true);
setCfgMsg(null);
setCfgErr(null);
try {
const res = await fetch("/api/v1/admin/downloads/config", {
method: "PUT",
credentials: "include",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({
...cfg,
nodeId,
password: password || undefined,
}),
});
const data = await res.json().catch(() => ({}));
if (!res.ok) {
setCfgErr(data.error?.message ?? "Opslaan mislukt");
return;
}
setCfgMsg("Opgeslagen — verbinding OK");
setPassword("");
loadAll();
} finally {
setCfgBusy(false);
}
}
async function runSearch() {
if (!query.trim() || !nodeId) return;
setSearching(true);
setSearchErr(null);
setHits([]);
try {
const res = await fetch("/api/v1/admin/downloads/search", {
method: "POST",
credentials: "include",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ nodeId, q: query.trim() }),
});
const data = await res.json().catch(() => ({}));
if (!res.ok) {
setSearchErr(data.error?.message ?? "Zoeken mislukt");
return;
}
setHits(data.results ?? []);
} finally {
setSearching(false);
}
}
async function openMatch(hit: SearchHit) {
setPick(hit);
setSelectedMeta(null);
setJobErr(null);
setMetaQ(hit.suggestQuery || hit.title);
setMetaHits([]);
await searchMeta(hit.suggestQuery || hit.title);
}
async function searchMeta(q: string) {
const term = q.trim();
if (!term) return;
setMetaBusy(true);
try {
const res = await fetch(
`/api/v1/admin/metadata/search?type=movie&q=${encodeURIComponent(term)}`,
{ credentials: "include" }
);
const data = await res.json().catch(() => ({}));
setMetaHits(data.results ?? []);
} finally {
setMetaBusy(false);
}
}
async function confirmDownload() {
if (!pick || !selectedMeta || !pick.uri) {
setJobErr("Kies een match en zorg dat er een download-URI is");
return;
}
setJobBusy(true);
setJobErr(null);
try {
const res = await fetch("/api/v1/admin/downloads/jobs", {
method: "POST",
credentials: "include",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({
nodeId,
uri: pick.uri,
sourceTitle: pick.title,
tmdbId: selectedMeta.tmdbId,
imdbId: selectedMeta.imdbId ?? undefined,
movieTitle: selectedMeta.title,
movieYear: selectedMeta.year,
posterUrl: selectedMeta.posterUrl,
}),
});
const data = await res.json().catch(() => ({}));
if (!res.ok) {
setJobErr(data.error?.message ?? "Download starten mislukt");
return;
}
setPick(null);
loadJobs();
} finally {
setJobBusy(false);
}
}
async function retryImport(id: string) {
await fetch(`/api/v1/admin/downloads/jobs/${id}/retry-import`, {
method: "POST",
credentials: "include",
});
loadJobs();
}
return (
<>
<Nav />
<div className="container">
<div className="page-header">
<div>
<div className="page-kicker">Synology · 4K</div>
<h1>Downloads</h1>
</div>
</div>
<section className="stack-section">
<div className="section-head">
<h2 className="section-title">Download Station</h2>
<p className="muted">
Zoek via de NAS, bevestig IMDb/TMDB, en na finish gaat de film naar de vaste 4K-map.
</p>
</div>
<div className="card">
<div className="inline-form" style={{ marginBottom: "1rem" }}>
<select value={nodeId} onChange={(e) => setNodeId(e.target.value)}>
<option value="">Kies Synology-node…</option>
{nodes.map((n) => (
<option key={n.id} value={n.id}>
{n.name} ({n.status})
</option>
))}
</select>
</div>
<h3 className="card-subtitle">Configuratie</h3>
<div className="form-grid">
<label>
DSM URL
<input
value={cfg.baseUrl}
onChange={(e) => setCfg({ ...cfg, baseUrl: e.target.value })}
placeholder="https://192.168.x.x:5001"
/>
</label>
<label>
Gebruiker
<input
value={cfg.username}
onChange={(e) => setCfg({ ...cfg, username: e.target.value })}
autoComplete="username"
/>
</label>
<label>
Wachtwoord {hasConfig ? "(leeg = behouden)" : ""}
<input
type="password"
value={password}
onChange={(e) => setPassword(e.target.value)}
autoComplete="new-password"
/>
</label>
<label>
DS destination (share-relatief)
<input
value={cfg.downloadDestination}
onChange={(e) => setCfg({ ...cfg, downloadDestination: e.target.value })}
placeholder="Downloads"
/>
</label>
<label>
Download host-pad (node)
<input
value={cfg.downloadHostPath}
onChange={(e) => setCfg({ ...cfg, downloadHostPath: e.target.value })}
placeholder="/volume1/Downloads"
/>
</label>
<label>
4K bibliotheek host-pad
<input
value={cfg.fourKHostPath}
onChange={(e) => setCfg({ ...cfg, fourKHostPath: e.target.value })}
placeholder="/volume1/Media/Films/4K"
/>
</label>
</div>
<div className="inline-form" style={{ marginTop: "0.75rem" }}>
<button type="button" disabled={cfgBusy || !nodeId} onClick={saveConfig}>
{cfgBusy ? "Bezig…" : "Opslaan + testen"}
</button>
{cfgMsg && <span className="ok-text">{cfgMsg}</span>}
{cfgErr && <span className="danger-text">{cfgErr}</span>}
</div>
{cfg.lastError && <p className="muted danger-text">Laatste fout: {cfg.lastError}</p>}
</div>
</section>
<section className="stack-section">
<div className="section-head">
<h2 className="section-title">Zoeken</h2>
<p className="muted">Vereist BT Search-modules in Download Station op de NAS.</p>
</div>
<div className="card">
<div className="inline-form">
<input
value={query}
onChange={(e) => setQuery(e.target.value)}
placeholder="Film + 2160p / UHD…"
onKeyDown={(e) => e.key === "Enter" && void runSearch()}
/>
<button type="button" disabled={searching || !hasConfig} onClick={runSearch}>
{searching ? "Zoeken…" : "Zoeken"}
</button>
</div>
{searchErr && <p className="danger-text">{searchErr}</p>}
{!hasConfig && <p className="muted">Sla eerst de DS-configuratie op.</p>}
<div className="mobile-cards" style={{ marginTop: "1rem" }}>
{hits.map((h) => (
<article key={`${h.title}-${h.size}-${h.seeds}`} className="mobile-card">
<div className="mobile-card-title">{h.title}</div>
<div className="muted">
{formatBytes(h.size)} · {h.seeds} seeds · {h.module || "module"}
</div>
<button
type="button"
disabled={!h.uri}
onClick={() => void openMatch(h)}
style={{ marginTop: "0.5rem" }}
>
Kiezen
</button>
</article>
))}
</div>
<div className="desktop-table" style={{ marginTop: "1rem" }}>
{hits.length > 0 && (
<table>
<thead>
<tr>
<th>Titel</th>
<th>Grootte</th>
<th>Seeds</th>
<th></th>
</tr>
</thead>
<tbody>
{hits.map((h) => (
<tr key={`${h.title}-${h.size}-${h.seeds}`}>
<td>
<div>{h.title}</div>
<div className="muted">{h.module}</div>
</td>
<td>{formatBytes(h.size)}</td>
<td>
{h.seeds}/{h.peers}
</td>
<td>
<button type="button" disabled={!h.uri} onClick={() => void openMatch(h)}>
Kiezen
</button>
</td>
</tr>
))}
</tbody>
</table>
)}
</div>
</div>
</section>
<section className="stack-section">
<div className="section-head">
<h2 className="section-title">Jobs</h2>
</div>
<div className="card">
<div className="mobile-cards">
{jobs.map((j) => (
<article key={j.id} className="mobile-card">
<div className="mobile-card-head">
<span className={statusClass(j.status)}>{j.status}</span>
<span className="muted">{j.nodeName}</span>
</div>
<div className="mobile-card-title">
{j.movieTitle}
{j.movieYear ? ` (${j.movieYear})` : ""}
</div>
<div className="muted">{j.sourceTitle}</div>
{j.error && <div className="danger-text">{j.error}</div>}
{(j.status === "FAILED" || j.status === "FINISHED") && (
<button type="button" onClick={() => void retryImport(j.id)}>
Opnieuw importeren
</button>
)}
</article>
))}
</div>
<div className="desktop-table">
<table>
<thead>
<tr>
<th>Status</th>
<th>Film</th>
<th>Bron</th>
<th>Node</th>
<th></th>
</tr>
</thead>
<tbody>
{jobs.map((j) => (
<tr key={j.id}>
<td>
<span className={statusClass(j.status)}>{j.status}</span>
</td>
<td>
{j.movieTitle}
{j.movieYear ? ` (${j.movieYear})` : ""}
{j.error && <div className="danger-text">{j.error}</div>}
</td>
<td className="muted">{j.sourceTitle}</td>
<td>{j.nodeName}</td>
<td>
{(j.status === "FAILED" || j.status === "FINISHED") && (
<button type="button" onClick={() => void retryImport(j.id)}>
Retry import
</button>
)}
</td>
</tr>
))}
{jobs.length === 0 && (
<tr>
<td colSpan={5} className="muted">
Nog geen downloads.
</td>
</tr>
)}
</tbody>
</table>
</div>
</div>
</section>
</div>
{pick && (
<div className="modal-backdrop" onClick={() => setPick(null)}>
<div className="modal-panel" onClick={(e) => e.stopPropagation()}>
<div className="page-kicker">Bevestig match</div>
<h2 style={{ marginBottom: "0.5rem" }}>Wat is dit?</h2>
<p className="muted" style={{ marginBottom: "1rem" }}>
{pick.title}
</p>
<div className="inline-form">
<input
value={metaQ}
onChange={(e) => setMetaQ(e.target.value)}
placeholder="TMDB-zoek of tt…"
onKeyDown={(e) => e.key === "Enter" && void searchMeta(metaQ)}
/>
<button type="button" disabled={metaBusy} onClick={() => void searchMeta(metaQ)}>
{metaBusy ? "…" : "Zoek"}
</button>
</div>
<div className="match-list">
{metaHits.map((m) => (
<button
key={m.tmdbId}
type="button"
className={`match-item${selectedMeta?.tmdbId === m.tmdbId ? " selected" : ""}`}
onClick={() => setSelectedMeta(m)}
>
{m.posterUrl && (
// eslint-disable-next-line @next/next/no-img-element
<img src={m.posterUrl} alt="" width={40} height={60} />
)}
<span>
<strong>
{m.title}
{m.year ? ` (${m.year})` : ""}
</strong>
<small className="muted">
TMDB {m.tmdbId}
{m.imdbId ? ` · ${m.imdbId}` : ""}
</small>
</span>
</button>
))}
</div>
{jobErr && <p className="danger-text">{jobErr}</p>}
<div className="inline-form" style={{ marginTop: "1rem" }}>
<button type="button" disabled={jobBusy || !selectedMeta} onClick={() => void confirmDownload()}>
{jobBusy ? "Starten…" : "Bevestigen + downloaden"}
</button>
<button type="button" className="btn-logout" onClick={() => setPick(null)}>
Annuleren
</button>
</div>
</div>
</div>
)}
</>
);
}

View file

@ -833,6 +833,26 @@ tr.row-selected {
margin-bottom: 1rem;
}
.form-grid > label {
display: flex;
flex-direction: column;
gap: 0.35rem;
font-size: 0.85rem;
color: var(--text-muted);
}
.form-grid > label input {
margin-bottom: 0;
min-height: 44px;
background: var(--bg-deep);
border: 1px solid var(--border-dim);
color: var(--text-bright);
border-radius: 7px;
padding: 0.55rem 0.75rem;
font-family: var(--font-ui);
font-size: 0.9rem;
}
.form-span {
grid-column: 1 / -1;
}
@ -1203,3 +1223,98 @@ tr.row-selected {
transition-duration: 0.01ms !important;
}
}
.badge {
display: inline-flex;
align-items: center;
font-family: var(--font-mono);
font-size: 0.68rem;
letter-spacing: 0.04em;
text-transform: uppercase;
padding: 0.2rem 0.45rem;
border-radius: 4px;
border: 1px solid var(--border-dim);
color: var(--text-muted);
background: var(--bg-elevated);
}
.badge-ok {
color: var(--ok);
background: var(--ok-bg);
border-color: rgba(52, 211, 153, 0.35);
}
.badge-warn {
color: var(--warn);
background: rgba(251, 191, 36, 0.12);
border-color: rgba(251, 191, 36, 0.3);
}
.badge-danger {
color: var(--danger);
background: var(--danger-bg);
border-color: rgba(244, 63, 94, 0.35);
}
.ok-text {
color: var(--ok);
}
.danger-text {
color: var(--danger);
}
.modal-backdrop {
position: fixed;
inset: 0;
z-index: 80;
background: rgba(0, 0, 0, 0.65);
display: flex;
align-items: flex-end;
justify-content: center;
padding: 1rem;
}
@media (min-width: 720px) {
.modal-backdrop {
align-items: center;
}
}
.modal-panel {
width: min(560px, 100%);
max-height: 90vh;
overflow: auto;
background: var(--bg-panel);
border: 1px solid var(--border);
border-radius: var(--radius);
padding: 1.25rem;
box-shadow: var(--shadow-panel);
}
.match-list {
display: flex;
flex-direction: column;
gap: 0.5rem;
margin-top: 0.75rem;
max-height: 40vh;
overflow: auto;
}
.match-item {
display: flex;
gap: 0.75rem;
align-items: center;
text-align: left;
background: var(--bg-elevated);
border: 1px solid var(--border-dim);
border-radius: 8px;
padding: 0.5rem 0.75rem;
color: inherit;
cursor: pointer;
width: 100%;
}
.match-item.selected {
border-color: var(--accent);
background: var(--accent-dim);
}
.match-item span {
display: flex;
flex-direction: column;
gap: 0.15rem;
}
.match-item img {
border-radius: 4px;
object-fit: cover;
}

View file

@ -10,6 +10,7 @@ const desktopLinks = [
{ href: "/nodes", label: "Nodes" },
{ href: "/movies", label: "Movies" },
{ href: "/series", label: "Series" },
{ href: "/downloads", label: "Downloads" },
{ href: "/streams", label: "Streams" },
{ href: "/accounts", label: "Accounts" },
{ href: "/settings", label: "Instellingen" },
@ -20,9 +21,9 @@ const mobileLinks = [
{ href: "/dashboard", label: "Home", short: "Home" },
{ href: "/nodes", label: "Nodes", short: "Nodes" },
{ href: "/movies", label: "Movies", short: "Films" },
{ href: "/downloads", label: "Downloads", short: "DL" },
{ href: "/streams", label: "Streams", short: "Live" },
{ href: "/accounts", label: "Accounts", short: "Acc." },
{ href: "/settings", label: "Instellingen", short: "Set." },
];
export function Nav() {

View file

@ -0,0 +1,61 @@
-- CreateEnum
DO $$ BEGIN
CREATE TYPE "DownloadJobStatus" AS ENUM ('DOWNLOADING', 'FINISHED', 'IMPORTING', 'IMPORTED', 'FAILED');
EXCEPTION WHEN duplicate_object THEN null; END $$;
-- CreateTable
CREATE TABLE IF NOT EXISTS "synology_ds_configs" (
"id" TEXT NOT NULL,
"node_id" TEXT NOT NULL,
"base_url" TEXT NOT NULL,
"username" TEXT NOT NULL,
"password_enc" TEXT NOT NULL,
"download_destination" TEXT NOT NULL,
"download_host_path" TEXT NOT NULL,
"four_k_host_path" TEXT NOT NULL,
"enabled" BOOLEAN NOT NULL DEFAULT true,
"last_ok_at" TIMESTAMP(3),
"last_error" TEXT,
"created_at" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
"updated_at" TIMESTAMP(3) NOT NULL,
CONSTRAINT "synology_ds_configs_pkey" PRIMARY KEY ("id")
);
CREATE UNIQUE INDEX IF NOT EXISTS "synology_ds_configs_node_id_key" ON "synology_ds_configs"("node_id");
DO $$ BEGIN
ALTER TABLE "synology_ds_configs"
ADD CONSTRAINT "synology_ds_configs_node_id_fkey"
FOREIGN KEY ("node_id") REFERENCES "nodes"("id") ON DELETE CASCADE ON UPDATE CASCADE;
EXCEPTION WHEN duplicate_object THEN null; END $$;
-- CreateTable
CREATE TABLE IF NOT EXISTS "download_jobs" (
"id" TEXT NOT NULL,
"node_id" TEXT NOT NULL,
"status" "DownloadJobStatus" NOT NULL DEFAULT 'DOWNLOADING',
"source_uri" TEXT NOT NULL,
"source_title" TEXT NOT NULL,
"ds_task_id" TEXT,
"tmdb_id" INTEGER,
"imdb_id" TEXT,
"movie_title" TEXT NOT NULL,
"movie_year" INTEGER,
"poster_url" TEXT,
"dest_folder_name" TEXT,
"source_path" TEXT,
"imported_path" TEXT,
"error" TEXT,
"created_at" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
"updated_at" TIMESTAMP(3) NOT NULL,
CONSTRAINT "download_jobs_pkey" PRIMARY KEY ("id")
);
CREATE INDEX IF NOT EXISTS "download_jobs_status_updated_at_idx" ON "download_jobs"("status", "updated_at");
CREATE INDEX IF NOT EXISTS "download_jobs_node_id_status_idx" ON "download_jobs"("node_id", "status");
DO $$ BEGIN
ALTER TABLE "download_jobs"
ADD CONSTRAINT "download_jobs_node_id_fkey"
FOREIGN KEY ("node_id") REFERENCES "nodes"("id") ON DELETE CASCADE ON UPDATE CASCADE;
EXCEPTION WHEN duplicate_object THEN null; END $$;

View file

@ -35,6 +35,14 @@ enum LibraryShelfKind {
SERIES
}
enum DownloadJobStatus {
DOWNLOADING
FINISHED
IMPORTING
IMPORTED
FAILED
}
model User {
id String @id @default(uuid())
email String @unique
@ -140,11 +148,13 @@ model Node {
createdAt DateTime @default(now()) @map("created_at")
updatedAt DateTime @updatedAt @map("updated_at")
location Location? @relation(fields: [locationId], references: [id])
credentials NodeCredential?
mediaFiles MediaFile[]
location Location? @relation(fields: [locationId], references: [id])
credentials NodeCredential?
mediaFiles MediaFile[]
playbackSessions PlaybackSession[]
scanRoots NodeScanRoot[]
scanRoots NodeScanRoot[]
synologyDsConfig SynologyDsConfig?
downloadJobs DownloadJob[]
@@map("nodes")
}
@ -362,3 +372,54 @@ model PlaybackSession {
@@index([addonTokenId])
@@map("playback_sessions")
}
/** Synology Download Station credentials + fixed 4K import paths (per node). */
model SynologyDsConfig {
id String @id @default(uuid())
nodeId String @unique @map("node_id")
/** DSM base URL, e.g. https://nas.local:5001 */
baseUrl String @map("base_url")
username String
passwordEnc String @map("password_enc")
/** DS task destination (share-relative), e.g. Downloads/complete */
downloadDestination String @map("download_destination")
/** Absolute host path of that destination on the media-node */
downloadHostPath String @map("download_host_path")
/** Absolute host path for 4K movies library folder */
fourKHostPath String @map("four_k_host_path")
enabled Boolean @default(true)
lastOkAt DateTime? @map("last_ok_at")
lastError String? @map("last_error")
createdAt DateTime @default(now()) @map("created_at")
updatedAt DateTime @updatedAt @map("updated_at")
node Node @relation(fields: [nodeId], references: [id], onDelete: Cascade)
@@map("synology_ds_configs")
}
model DownloadJob {
id String @id @default(uuid())
nodeId String @map("node_id")
status DownloadJobStatus @default(DOWNLOADING)
sourceUri String @map("source_uri")
sourceTitle String @map("source_title")
dsTaskId String? @map("ds_task_id")
tmdbId Int? @map("tmdb_id")
imdbId String? @map("imdb_id")
movieTitle String @map("movie_title")
movieYear Int? @map("movie_year")
posterUrl String? @map("poster_url")
destFolderName String? @map("dest_folder_name")
sourcePath String? @map("source_path")
importedPath String? @map("imported_path")
error String?
createdAt DateTime @default(now()) @map("created_at")
updatedAt DateTime @updatedAt @map("updated_at")
node Node @relation(fields: [nodeId], references: [id], onDelete: Cascade)
@@index([status, updatedAt])
@@index([nodeId, status])
@@map("download_jobs")
}

View file

@ -10,6 +10,7 @@ import { registerNodeRoutes } from "./nodes/routes";
import { registerStremioRoutes } from "./stremio/routes";
import { registerAdminRoutes } from "./admin/routes";
import { registerInstallRoutes } from "./install/routes";
import { registerDownloadRoutes } from "./downloads/routes";
import { nodeConnectionManager } from "./websocket/manager";
import { MetadataService } from "./metadata/service";
import { toErrorResponse } from "./security/errors";
@ -71,6 +72,7 @@ async function main() {
await registerStremioRoutes(app, config);
await registerAdminRoutes(app, config);
await registerInstallRoutes(app, config);
registerDownloadRoutes(app, config);
app.get("/api/v1/node/connect", { websocket: true }, (socket) => {
void nodeConnectionManager.handleConnection(socket);
@ -179,6 +181,65 @@ async function main() {
console.warn("Could not ensure library shelves:", err);
}
try {
await prisma.$executeRawUnsafe(`
DO $$ BEGIN
CREATE TYPE "DownloadJobStatus" AS ENUM ('DOWNLOADING', 'FINISHED', 'IMPORTING', 'IMPORTED', 'FAILED');
EXCEPTION WHEN duplicate_object THEN null; END $$;
`);
await prisma.$executeRawUnsafe(`
CREATE TABLE IF NOT EXISTS "synology_ds_configs" (
"id" TEXT NOT NULL,
"node_id" TEXT NOT NULL,
"base_url" TEXT NOT NULL,
"username" TEXT NOT NULL,
"password_enc" TEXT NOT NULL,
"download_destination" TEXT NOT NULL,
"download_host_path" TEXT NOT NULL,
"four_k_host_path" TEXT NOT NULL,
"enabled" BOOLEAN NOT NULL DEFAULT true,
"last_ok_at" TIMESTAMP(3),
"last_error" TEXT,
"created_at" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
"updated_at" TIMESTAMP(3) NOT NULL,
CONSTRAINT "synology_ds_configs_pkey" PRIMARY KEY ("id")
)
`);
await prisma.$executeRawUnsafe(
`CREATE UNIQUE INDEX IF NOT EXISTS "synology_ds_configs_node_id_key" ON "synology_ds_configs"("node_id")`
);
await prisma.$executeRawUnsafe(`
CREATE TABLE IF NOT EXISTS "download_jobs" (
"id" TEXT NOT NULL,
"node_id" TEXT NOT NULL,
"status" "DownloadJobStatus" NOT NULL DEFAULT 'DOWNLOADING',
"source_uri" TEXT NOT NULL,
"source_title" TEXT NOT NULL,
"ds_task_id" TEXT,
"tmdb_id" INTEGER,
"imdb_id" TEXT,
"movie_title" TEXT NOT NULL,
"movie_year" INTEGER,
"poster_url" TEXT,
"dest_folder_name" TEXT,
"source_path" TEXT,
"imported_path" TEXT,
"error" TEXT,
"created_at" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
"updated_at" TIMESTAMP(3) NOT NULL,
CONSTRAINT "download_jobs_pkey" PRIMARY KEY ("id")
)
`);
await prisma.$executeRawUnsafe(
`CREATE INDEX IF NOT EXISTS "download_jobs_status_updated_at_idx" ON "download_jobs"("status", "updated_at")`
);
await prisma.$executeRawUnsafe(
`CREATE INDEX IF NOT EXISTS "download_jobs_node_id_status_idx" ON "download_jobs"("node_id", "status")`
);
} catch (err) {
console.warn("Could not ensure download tables:", err);
}
// Repair drifted episode.series_id so series.episodes and seasons.episodes stay aligned.
try {
await prisma.$executeRawUnsafe(`

View file

@ -0,0 +1,119 @@
import type { FastifyInstance } from "fastify";
import { requireAdmin } from "../auth/routes";
import type { Config } from "../config";
import { MetadataService } from "../metadata/service";
import { AppError } from "../security/errors";
import { DownloadService, startDownloadPoller } from "./service";
export function registerDownloadRoutes(app: FastifyInstance, config: Config): DownloadService {
const metadata = new MetadataService(config.TMDB_API_KEY);
const downloads = new DownloadService(config.SESSION_SECRET, metadata);
startDownloadPoller(downloads);
app.get("/api/v1/admin/downloads/configs", { preHandler: requireAdmin }, async () => {
const configs = await downloads.listConfigs();
return { configs };
});
app.put("/api/v1/admin/downloads/config", { preHandler: requireAdmin }, async (request) => {
const body = request.body as {
nodeId?: string;
baseUrl?: string;
username?: string;
password?: string;
downloadDestination?: string;
downloadHostPath?: string;
fourKHostPath?: string;
enabled?: boolean;
};
if (!body.nodeId || !body.baseUrl || !body.username) {
throw new AppError("INVALID_REQUEST", "nodeId, baseUrl en username verplicht", 400);
}
const row = await downloads.upsertConfig({
nodeId: body.nodeId,
baseUrl: body.baseUrl,
username: body.username,
password: body.password,
downloadDestination: body.downloadDestination ?? "",
downloadHostPath: body.downloadHostPath ?? "",
fourKHostPath: body.fourKHostPath ?? "",
enabled: body.enabled,
});
return { config: row };
});
app.post("/api/v1/admin/downloads/config/test", { preHandler: requireAdmin }, async (request) => {
const body = request.body as { nodeId?: string };
if (!body.nodeId) throw new AppError("INVALID_REQUEST", "nodeId verplicht", 400);
return downloads.testConfig(body.nodeId);
});
app.post("/api/v1/admin/downloads/search", { preHandler: requireAdmin }, async (request) => {
const body = request.body as { nodeId?: string; q?: string };
if (!body.nodeId || !body.q?.trim()) {
throw new AppError("INVALID_REQUEST", "nodeId en q verplicht", 400);
}
const results = await downloads.search(body.nodeId, body.q);
return { results };
});
app.get("/api/v1/admin/downloads/jobs", { preHandler: requireAdmin }, async (request) => {
const q = request.query as { limit?: string };
const limit = q.limit ? parseInt(q.limit, 10) : 50;
const jobs = await downloads.listJobs(limit);
return {
jobs: jobs.map((j) => ({
id: j.id,
nodeId: j.nodeId,
nodeName: j.node.name,
status: j.status,
sourceTitle: j.sourceTitle,
movieTitle: j.movieTitle,
movieYear: j.movieYear,
posterUrl: j.posterUrl,
tmdbId: j.tmdbId,
imdbId: j.imdbId,
destFolderName: j.destFolderName,
importedPath: j.importedPath,
error: j.error,
createdAt: j.createdAt,
updatedAt: j.updatedAt,
})),
};
});
app.post("/api/v1/admin/downloads/jobs", { preHandler: requireAdmin }, async (request) => {
const body = request.body as {
nodeId?: string;
uri?: string;
sourceTitle?: string;
tmdbId?: number;
imdbId?: string;
movieTitle?: string;
movieYear?: number | null;
posterUrl?: string | null;
};
if (!body.nodeId || !body.uri || !body.movieTitle) {
throw new AppError("INVALID_REQUEST", "nodeId, uri en movieTitle verplicht", 400);
}
const job = await downloads.createJob({
nodeId: body.nodeId,
uri: body.uri,
sourceTitle: body.sourceTitle || body.movieTitle,
tmdbId: body.tmdbId,
imdbId: body.imdbId,
movieTitle: body.movieTitle,
movieYear: body.movieYear,
posterUrl: body.posterUrl,
});
return { job };
});
app.post("/api/v1/admin/downloads/jobs/:id/retry-import", { preHandler: requireAdmin }, async (request) => {
const { id } = request.params as { id: string };
const job = await downloads.retryImport(id);
return { job };
});
return downloads;
}

View file

@ -0,0 +1,533 @@
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";
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(/\/+$/, "");
}
export type DsConfigInput = {
nodeId: string;
baseUrl: string;
username: string;
password?: string;
downloadDestination: string;
downloadHostPath: string;
fourKHostPath: string;
enabled?: boolean;
};
export class DownloadService {
constructor(
private readonly sessionSecret: string,
private readonly metadata: MetadataService
) {}
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,
fourKHostPath: r.fourKHostPath,
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 fourKHostPath = normalizeHostPath(input.fourKHostPath);
if (!username || !downloadDestination || !downloadHostPath || !fourKHostPath) {
throw new AppError("INVALID_REQUEST", "Vul alle verplichte velden in", 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,
fourKHostPath,
enabled: input.enabled ?? true,
lastOkAt: new Date(),
lastError: null,
},
update: {
baseUrl,
username,
passwordEnc,
downloadDestination,
downloadHostPath,
fourKHostPath,
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,
fourKHostPath: row.fourKHostPath,
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): Promise<
Array<{
title: string;
size: number;
seeds: number;
peers: number;
uri: string | null;
module: string | null;
suggestQuery: string;
}>
> {
const { client } = await this.clientForNode(nodeId);
let hits: DsSearchHit[];
try {
hits = await client.search(query);
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;
tmdbId?: number;
imdbId?: string;
movieTitle: string;
movieYear?: number | null;
posterUrl?: string | 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 film (titel)", 400);
}
if (!input.tmdbId && !input.imdbId) {
throw new AppError("INVALID_REQUEST", "Bevestig met TMDB- of IMDb-id", 400);
}
const { client, config } = await this.clientForNode(input.nodeId);
// Resolve full metadata when possible
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;
// Prefer IMDb resolve when only tt… is given
if (!tmdbId && imdbId) {
const hits = await this.metadata.searchCatalog("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;
}
}
// Ensure movie row exists so library matching is smoother after import
if (tmdbId) {
try {
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] movie upsert:", err);
}
}
const destFolderName = folderNameForMovie(title, year);
// Snapshot task ids before create so we can discover the new id
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",
sourceUri: uri,
sourceTitle: input.sourceTitle.trim() || title,
dsTaskId,
tmdbId: tmdbId ?? null,
imdbId,
movieTitle: title,
movieYear: year,
posterUrl,
destFolderName,
},
});
return job;
}
async listJobs(limit = 50) {
return prisma.downloadJob.findMany({
take: Math.min(100, Math.max(1, limit)),
orderBy: { createdAt: "desc" },
include: { node: { select: { id: true, name: true } } },
});
}
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;
}
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 destFolder = job.destFolderName || folderNameForMovie(job.movieTitle, job.movieYear);
const destDir = `${normalizeHostPath(config.fourKHostPath)}/${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 },
});
const ack = await nodeConnectionManager.importMediaAcked(
job.nodeId,
{
jobId: job.id,
sourcePath,
destDir,
movieTitle: job.movieTitle,
movieYear: job.movieYear,
qualityHint: "2160p",
},
120_000
);
if (!ack.ok) {
await prisma.downloadJob.update({
where: { id: jobId },
data: {
status: "FAILED",
error: ack.error || "Import mislukt of timeout",
},
});
return;
}
await prisma.downloadJob.update({
where: { id: jobId },
data: {
status: "IMPORTED",
importedPath: ack.importedPath ?? null,
error: null,
},
});
}
}
function isTaskDone(status: string): boolean {
return status === "finished" || status === "seeding";
}
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 = 30_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);
}

View file

@ -0,0 +1,319 @@
import http from "node:http";
import https from "node:https";
import { URL } from "node:url";
type DsApiResponse<T = unknown> = {
success: boolean;
data?: T;
error?: { code: number };
};
export type DsSearchHit = {
title: string;
date: number;
size: number;
seeds: number;
peers: number;
leechs: number;
download_uri?: string;
magnet?: string;
page?: string;
external?: boolean;
module_id?: string;
module_title?: string;
};
export type DsTask = {
id: string;
type: string;
username?: string;
title: string;
size: number;
status:
| "waiting"
| "downloading"
| "paused"
| "finishing"
| "finished"
| "hash_checking"
| "seeding"
| "filehosting_waiting"
| "extracting"
| "error"
| string;
status_extra?: { error_detail?: string; unzip_progress?: number };
additional?: {
detail?: { destination?: string; uri?: string; connected_leechers?: number };
transfer?: { size_downloaded?: number; speed_download?: number };
file?: Array<{ filename: string; size: number; size_downloaded?: number }>;
};
};
function requestJson(urlStr: string, opts?: { method?: string; body?: string }): Promise<unknown> {
return new Promise((resolve, reject) => {
const u = new URL(urlStr);
const lib = u.protocol === "https:" ? https : http;
const req = lib.request(
{
protocol: u.protocol,
hostname: u.hostname,
port: u.port || (u.protocol === "https:" ? 443 : 80),
path: `${u.pathname}${u.search}`,
method: opts?.method ?? "GET",
rejectUnauthorized: false,
timeout: 45_000,
headers: opts?.body
? {
"Content-Type": "application/x-www-form-urlencoded",
"Content-Length": Buffer.byteLength(opts.body),
}
: undefined,
},
(res) => {
const chunks: Buffer[] = [];
res.on("data", (c) => chunks.push(c));
res.on("end", () => {
const raw = Buffer.concat(chunks).toString("utf8");
try {
resolve(JSON.parse(raw));
} catch (err) {
reject(new Error(`Invalid JSON from Synology (${res.statusCode}): ${raw.slice(0, 200)}`));
}
});
}
);
req.on("error", reject);
req.on("timeout", () => {
req.destroy();
reject(new Error("Synology request timeout"));
});
if (opts?.body) req.write(opts.body);
req.end();
});
}
function qs(params: Record<string, string | number | undefined>): string {
const sp = new URLSearchParams();
for (const [k, v] of Object.entries(params)) {
if (v === undefined || v === "") continue;
sp.set(k, String(v));
}
return sp.toString();
}
export class SynologyDownloadStation {
private sid: string | null = null;
private sidAt = 0;
constructor(
private readonly baseUrl: string,
private readonly username: string,
private readonly password: string
) {}
private root(): string {
return this.baseUrl.replace(/\/+$/, "");
}
async login(force = false): Promise<void> {
if (!force && this.sid && Date.now() - this.sidAt < 10 * 60_000) return;
const url =
`${this.root()}/webapi/auth.cgi?` +
qs({
api: "SYNO.API.Auth",
version: 3,
method: "login",
account: this.username,
passwd: this.password,
session: "DownloadStation",
format: "sid",
});
const data = (await requestJson(url)) as DsApiResponse<{ sid: string }>;
if (!data.success || !data.data?.sid) {
throw new Error(`Synology login mislukt (code ${data.error?.code ?? "?"})`);
}
this.sid = data.data.sid;
this.sidAt = Date.now();
}
private async withSid<T>(fn: (sid: string) => Promise<T>): Promise<T> {
await this.login();
try {
return await fn(this.sid!);
} catch (err) {
// One retry after re-login (expired session)
await this.login(true);
return fn(this.sid!);
}
}
async testConnection(): Promise<void> {
await this.login(true);
const info = await this.withSid(async (sid) => {
const url =
`${this.root()}/webapi/DownloadStation/info.cgi?` +
qs({
api: "SYNO.DownloadStation.Info",
version: 1,
method: "getinfo",
_sid: sid,
});
return (await requestJson(url)) as DsApiResponse;
});
if (!info.success) {
throw new Error(`Download Station niet bereikbaar (code ${info.error?.code ?? "?"})`);
}
}
/** Search via DS BTSearch modules (must be configured on the NAS). */
async search(keyword: string, limit = 40): Promise<DsSearchHit[]> {
const q = keyword.trim();
if (!q) return [];
return this.withSid(async (sid) => {
const startUrl =
`${this.root()}/webapi/DownloadStation/btsearch.cgi?` +
qs({
api: "SYNO.DownloadStation.BTSearch",
version: 1,
method: "start",
keyword: q,
module: "enabled",
_sid: sid,
});
const started = (await requestJson(startUrl)) as DsApiResponse<{ taskid: string }>;
if (!started.success || !started.data?.taskid) {
throw new Error(
`DS-zoekopdracht starten mislukt (code ${started.error?.code ?? "?"}). Controleer BT Search-modules op de NAS.`
);
}
const taskId = started.data.taskid;
try {
const deadline = Date.now() + 25_000;
let items: DsSearchHit[] = [];
while (Date.now() < deadline) {
await sleep(900);
const listUrl =
`${this.root()}/webapi/DownloadStation/btsearch.cgi?` +
qs({
api: "SYNO.DownloadStation.BTSearch",
version: 1,
method: "list",
taskid: taskId,
offset: 0,
limit,
_sid: sid,
});
const listed = (await requestJson(listUrl)) as DsApiResponse<{
finished?: boolean;
items?: DsSearchHit[];
total?: number;
}>;
if (!listed.success) {
throw new Error(`DS-zoekresultaten ophalen mislukt (code ${listed.error?.code ?? "?"})`);
}
items = listed.data?.items ?? [];
if (listed.data?.finished) break;
}
return items;
} finally {
const cleanUrl =
`${this.root()}/webapi/DownloadStation/btsearch.cgi?` +
qs({
api: "SYNO.DownloadStation.BTSearch",
version: 1,
method: "clean",
taskid: taskId,
_sid: sid,
});
void requestJson(cleanUrl).catch(() => undefined);
}
});
}
async createTask(uri: string, destination: string): Promise<void> {
const dest = destination.replace(/^\/+/, "").replace(/\/+$/, "");
await this.withSid(async (sid) => {
const body = qs({
api: "SYNO.DownloadStation.Task",
version: 1,
method: "create",
uri,
destination: dest,
_sid: sid,
});
const url = `${this.root()}/webapi/DownloadStation/task.cgi`;
const res = (await requestJson(url, { method: "POST", body })) as DsApiResponse;
if (!res.success) {
throw new Error(`Download-taak aanmaken mislukt (code ${res.error?.code ?? "?"})`);
}
});
}
async listTasks(): Promise<DsTask[]> {
return this.withSid(async (sid) => {
const url =
`${this.root()}/webapi/DownloadStation/task.cgi?` +
qs({
api: "SYNO.DownloadStation.Task",
version: 1,
method: "list",
additional: "detail,transfer,file",
_sid: sid,
});
const res = (await requestJson(url)) as DsApiResponse<{ tasks?: DsTask[] }>;
if (!res.success) {
throw new Error(`Taken ophalen mislukt (code ${res.error?.code ?? "?"})`);
}
return res.data?.tasks ?? [];
});
}
async deleteTask(id: string, forceComplete = false): Promise<void> {
await this.withSid(async (sid) => {
const url =
`${this.root()}/webapi/DownloadStation/task.cgi?` +
qs({
api: "SYNO.DownloadStation.Task",
version: 1,
method: "delete",
id,
force_complete: forceComplete ? "true" : "false",
_sid: sid,
});
const res = (await requestJson(url)) as DsApiResponse;
if (!res.success) {
throw new Error(`Taak verwijderen mislukt (code ${res.error?.code ?? "?"})`);
}
});
}
}
function sleep(ms: number): Promise<void> {
return new Promise((r) => setTimeout(r, ms));
}
/** Strip common release tags so TMDB search works better. */
export function cleanReleaseQuery(raw: string): string {
let s = raw.trim();
s = s.replace(/\.(mkv|mp4|m4v|avi|ts|m2ts)$/i, "");
s = s.replace(/[._]+/g, " ");
s = s.replace(
/\b(2160p|1080p|720p|480p|4k|uhd|bluray|blu-?ray|web-?dl|webrip|hdtv|remux|proper|repack|extended|theatrical|hybrid|multi|dual|nl|dut|dutch|eng|english|x265|x264|h\.?265|h\.?264|hevc|av1|hdr10?\+?|dv|dovi|atmos|dts(?:-hd)?|truehd|aac|ac3|flac|pgs|subs?|sample)\b/gi,
" "
);
s = s.replace(/\[[^\]]*]/g, " ");
s = s.replace(/\([^)]*(?:19|20)\d{2}[^)]*\)/g, (m) => m); // keep year parens
s = s.replace(/\s+/g, " ").trim();
return s;
}
export function folderNameForMovie(title: string, year?: number | null): string {
const safe = title
.replace(/[<>:"/\\|?*\x00-\x1f]/g, "")
.replace(/\s+/g, " ")
.trim();
if (year && year > 1900) return `${safe} (${year})`;
return safe;
}

View file

@ -15,6 +15,8 @@ import type {
PlaybackSessionEndedPayload,
PlaybackSessionProgressPayload,
CreatePlaybackSessionAckPayload,
ImportMediaAckPayload,
ImportMediaPayload,
} from "@media-cluster/protocol";
interface NodeConnection {
@ -29,6 +31,11 @@ interface PendingAck {
timer: ReturnType<typeof setTimeout>;
}
interface PendingImportAck {
resolve: (result: ImportMediaAckPayload) => void;
timer: ReturnType<typeof setTimeout>;
}
interface SyncState {
syncId: string;
fileIds: Set<string>;
@ -40,6 +47,7 @@ class NodeConnectionManager {
private config: Config | null = null;
private offlineCheckInterval: ReturnType<typeof setInterval> | null = null;
private pendingSessionAcks = new Map<string, PendingAck>();
private pendingImportAcks = new Map<string, PendingImportAck>();
private syncState = new Map<string, SyncState>();
private syncQueues = new Map<string, Promise<void>>();
@ -135,6 +143,9 @@ class NodeConnectionManager {
case "CREATE_PLAYBACK_SESSION_ACK":
this.handlePlaybackSessionAck(message.payload as CreatePlaybackSessionAckPayload);
break;
case "IMPORT_MEDIA_ACK":
this.handleImportMediaAck(message.payload as ImportMediaAckPayload);
break;
default:
this.send(socket, createMessage("ERROR", { message: `Unknown type: ${message.type}` }));
}
@ -413,6 +424,14 @@ class NodeConnectionManager {
pending.resolve(payload.ok);
}
private handleImportMediaAck(payload: ImportMediaAckPayload): void {
const pending = this.pendingImportAcks.get(payload.jobId);
if (!pending) return;
clearTimeout(pending.timer);
this.pendingImportAcks.delete(payload.jobId);
pending.resolve(payload);
}
private async handlePlaybackEnded(payload: PlaybackSessionEndedPayload): Promise<void> {
await prisma.playbackSession.updateMany({
where: { id: payload.sessionId },
@ -476,6 +495,41 @@ class NodeConnectionManager {
return this.sendToNode(nodeId, "RESTART", { reason });
}
isOnline(nodeId: string): boolean {
const conn = this.connections.get(nodeId);
return !!conn && conn.socket.readyState === 1;
}
async importMediaAcked(
nodeId: string,
payload: ImportMediaPayload,
timeoutMs = 120_000
): Promise<ImportMediaAckPayload> {
if (!this.isOnline(nodeId)) {
return { jobId: payload.jobId, ok: false, error: "Node offline" };
}
const acked = new Promise<ImportMediaAckPayload>((resolve) => {
const timer = setTimeout(() => {
this.pendingImportAcks.delete(payload.jobId);
resolve({ jobId: payload.jobId, ok: false, error: "Import timeout" });
}, timeoutMs);
this.pendingImportAcks.set(payload.jobId, { resolve, timer });
});
const pushed = this.sendToNode(nodeId, "IMPORT_MEDIA", payload);
if (!pushed) {
const pending = this.pendingImportAcks.get(payload.jobId);
if (pending) {
clearTimeout(pending.timer);
this.pendingImportAcks.delete(payload.jobId);
}
return { jobId: payload.jobId, ok: false, error: "Kan niet naar node sturen" };
}
return acked;
}
createPlaybackSession(nodeId: string, payload: CreatePlaybackSessionPayload): boolean {
return this.sendToNode(nodeId, "CREATE_PLAYBACK_SESSION", payload);
}
@ -523,11 +577,6 @@ class NodeConnectionManager {
}
}
isOnline(nodeId: string): boolean {
const conn = this.connections.get(nodeId);
return !!conn && conn.socket.readyState === 1;
}
private send(socket: WebSocket, message: ControlMessage): void {
if (socket.readyState === 1) {
socket.send(JSON.stringify(message));

View file

@ -31,8 +31,7 @@ Opties:
--public-url= Publieke HTTPS-URL via NPM (poort ${PORT})
--movies= Synology-pad naar films — bootstrap; later via Admin uitbreiden
--series= Synology-pad naar series (optioneel)
--media-root= Host-root die gemount wordt (default /volume1) zodat nieuwe
mappen vanaf Admin zichtbaar zijn zonder herinstall
--media-root= Host-root RW-mount (default /volume1) voor scan + Downloads→4K import
--port= Host-poort (default 8080; bij meerdere nodes op 1 LAN: uniek)
--appdata= Appdata-map (default /volume1/docker/media-node)
--uid= Container-user (default 1026 = typische DSM-admin)
@ -165,7 +164,7 @@ VOLUMES=(
-v "$APPDATA/config:/etc/media-node"
-v "$APPDATA/data:/var/lib/media-node"
-v "$BIN_PATH:/usr/local/bin/media-node:ro"
-v "$MEDIA_ROOT:$MEDIA_ROOT:ro"
-v "$MEDIA_ROOT:$MEDIA_ROOT:rw"
)
if [ -d /etc/ssl/certs ]; then
VOLUMES+=(-v /etc/ssl/certs:/etc/ssl/certs:ro)
@ -216,6 +215,7 @@ echo "Klaar. Logs: docker logs -f media-node"
echo "Status: in Admin → Nodes (moet online worden)"
echo "NPM: $PUBLIC_URL → Synology-LAN-IP:$PORT (buffering uit, Force SSL)"
echo "Paden: Admin → Nodes → Edit (host-paden onder $MEDIA_ROOT)"
echo "Schrijven: media-root is RW (nodig voor Downloads → 4K import)"
echo "UID/GID: bij permission denied: --uid=/--gid= van 'id jouwuser'"
echo
echo "Na succes: credentials staan in $APPDATA — herstart zonder nieuwe token is ok."

View file

@ -24,7 +24,7 @@ import (
"github.com/sthmedia/media-node/internal/streaming"
)
var version = "1.2.6"
var version = "1.3.0"
func main() {
if len(os.Args) < 2 {

View file

@ -19,6 +19,7 @@ import (
"github.com/sthmedia/media-node/internal/config"
"github.com/sthmedia/media-node/internal/database"
"github.com/sthmedia/media-node/internal/health"
"github.com/sthmedia/media-node/internal/importmedia"
"github.com/sthmedia/media-node/internal/scanner"
"github.com/sthmedia/media-node/internal/streaming"
)
@ -388,11 +389,51 @@ func (c *Client) handleMessage(msgType string, payload json.RawMessage) {
}
case "RESTART":
go c.doRestart()
case "IMPORT_MEDIA":
go c.handleImportMedia(payload)
case "ERROR":
log.Printf("Master error: %s", string(payload))
}
}
func (c *Client) handleImportMedia(payload json.RawMessage) {
var p struct {
JobID string `json:"jobId"`
SourcePath string `json:"sourcePath"`
DestDir string `json:"destDir"`
MovieTitle string `json:"movieTitle"`
MovieYear *int `json:"movieYear"`
QualityHint string `json:"qualityHint"`
}
if err := json.Unmarshal(payload, &p); err != nil || p.JobID == "" {
log.Printf("IMPORT_MEDIA: invalid payload")
return
}
result, err := importmedia.ImportMovie(p.SourcePath, p.DestDir, p.MovieTitle, p.MovieYear, p.QualityHint)
ok := err == nil
errMsg := ""
imported := ""
if err != nil {
errMsg = err.Error()
log.Printf("IMPORT_MEDIA failed job=%s: %v", p.JobID, err)
} else {
imported = result.ImportedPath
log.Printf("IMPORT_MEDIA ok job=%s path=%s", p.JobID, imported)
c.scanner.IndexPath(imported)
}
ack := controlMessage("IMPORT_MEDIA_ACK", map[string]interface{}{
"jobId": p.JobID,
"ok": ok,
"importedPath": imported,
"error": errMsg,
})
if werr := c.write(ack); werr != nil {
log.Printf("failed to send IMPORT_MEDIA_ACK: %v", werr)
}
}
func (c *Client) applyMediaPaths(movies, series []string) {
// nil means "not provided" — keep current. Empty slice means clear (rare).
if movies == nil && series == nil {

View file

@ -0,0 +1,187 @@
package importmedia
import (
"fmt"
"io"
"os"
"path/filepath"
"regexp"
"strings"
)
var videoExt = map[string]bool{
".mkv": true, ".mp4": true, ".m4v": true,
".avi": true, ".ts": true, ".m2ts": true,
}
var sampleRe = regexp.MustCompile(`(?i)(?:^|[.\-_])(sample|trailer|preview)(?:[.\-_]|$)`)
type Result struct {
ImportedPath string
}
// ImportMovie moves the largest video from sourcePath into destDir as Title (Year) - 2160p.ext
// and cleans leftover files under the download folder when source was a directory.
func ImportMovie(sourcePath, destDir, title string, year *int, qualityHint string) (*Result, error) {
sourcePath = filepath.Clean(strings.TrimSpace(sourcePath))
destDir = filepath.Clean(strings.TrimSpace(destDir))
if sourcePath == "" || destDir == "" {
return nil, fmt.Errorf("source/dest pad ontbreekt")
}
video, err := findLargestVideo(sourcePath)
if err != nil {
return nil, err
}
if err := os.MkdirAll(destDir, 0o755); err != nil {
return nil, fmt.Errorf("doelmap aanmaken: %w", err)
}
ext := strings.ToLower(filepath.Ext(video))
base := sanitizeName(title)
if year != nil && *year > 1900 {
base = fmt.Sprintf("%s (%d)", base, *year)
}
q := strings.TrimSpace(qualityHint)
if q == "" {
q = "2160p"
}
destName := fmt.Sprintf("%s - %s%s", base, q, ext)
destPath := filepath.Join(destDir, destName)
if err := moveFile(video, destPath); err != nil {
return nil, err
}
// Cleanup download leftovers (only if source was a dir under a parent, or sibling junk)
cleanupSource(sourcePath, video)
return &Result{ImportedPath: destPath}, nil
}
func findLargestVideo(sourcePath string) (string, error) {
info, err := os.Stat(sourcePath)
if err != nil {
return "", fmt.Errorf("bron niet gevonden: %w", err)
}
var best string
var bestSize int64
consider := func(path string, size int64) {
if !videoExt[strings.ToLower(filepath.Ext(path))] {
return
}
name := filepath.Base(path)
if sampleRe.MatchString(name) {
return
}
if size > bestSize {
bestSize = size
best = path
}
}
if !info.IsDir() {
consider(sourcePath, info.Size())
if best == "" {
return "", fmt.Errorf("geen videobestand: %s", sourcePath)
}
return best, nil
}
err = filepath.Walk(sourcePath, func(path string, fi os.FileInfo, walkErr error) error {
if walkErr != nil || fi == nil || fi.IsDir() {
return nil
}
consider(path, fi.Size())
return nil
})
if err != nil {
return "", err
}
if best == "" {
return "", fmt.Errorf("geen videobestand in %s", sourcePath)
}
return best, nil
}
func moveFile(src, dst string) error {
if err := os.Rename(src, dst); err == nil {
return nil
}
// Cross-device: copy + remove
in, err := os.Open(src)
if err != nil {
return err
}
defer in.Close()
out, err := os.OpenFile(dst, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o644)
if err != nil {
return err
}
if _, err := io.Copy(out, in); err != nil {
out.Close()
_ = os.Remove(dst)
return err
}
if err := out.Close(); err != nil {
return err
}
return os.Remove(src)
}
func cleanupSource(sourcePath, movedVideo string) {
info, err := os.Stat(sourcePath)
if err != nil {
// source was a file that moved away
parent := filepath.Dir(sourcePath)
_ = removeEmptyDirs(parent)
return
}
if !info.IsDir() {
_ = os.Remove(sourcePath)
_ = removeEmptyDirs(filepath.Dir(sourcePath))
return
}
_ = filepath.Walk(sourcePath, func(path string, fi os.FileInfo, walkErr error) error {
if walkErr != nil || fi == nil {
return nil
}
if path == movedVideo {
return nil
}
if !fi.IsDir() {
_ = os.Remove(path)
}
return nil
})
_ = removeEmptyDirs(sourcePath)
}
func removeEmptyDirs(root string) error {
var dirs []string
_ = filepath.Walk(root, func(path string, fi os.FileInfo, err error) error {
if err == nil && fi != nil && fi.IsDir() {
dirs = append(dirs, path)
}
return nil
})
for i := len(dirs) - 1; i >= 0; i-- {
_ = os.Remove(dirs[i]) // only succeeds if empty
}
return nil
}
func sanitizeName(s string) string {
s = strings.TrimSpace(s)
replacer := strings.NewReplacer(
"<", "", ">", "", ":", "", `"`, "", "/", "", "\\", "",
"|", "", "?", "", "*", "",
)
s = replacer.Replace(s)
return strings.Join(strings.Fields(s), " ")
}

View file

@ -612,6 +612,32 @@ func (s *Scanner) indexFile(path string, info os.FileInfo, eventType, preferredT
s.onEvents([]LibraryEvent{event})
}
// IndexPath forces a library event for a newly imported video file.
func (s *Scanner) IndexPath(path string) {
info, err := os.Stat(path)
if err != nil || info.IsDir() || !isVideo(path) {
return
}
s.mu.Lock()
roots := s.cfg.Media.EffectiveRoots()
s.mu.Unlock()
rootType := map[string]string{}
rootShelf := map[string]string{}
for _, root := range roots {
kind := root.Kind
if kind == "series" {
kind = "episode"
}
clean := filepath.Clean(root.Path)
rootType[clean] = kind
if root.ShelfID != "" {
rootShelf[clean] = root.ShelfID
}
}
kind, shelf := s.detectRootMeta(path, rootType, rootShelf)
s.indexFile(path, info, "FILE_ADDED", kind, shelf, true)
}
func (s *Scanner) handleRemove(path string) {
files, err := s.store.FilesUnderPath(path)
if err != nil || len(files) == 0 {

View file

@ -23,6 +23,8 @@ export type ControlMessageType =
| "RESCAN"
| "CONFIG_UPDATE"
| "RESTART"
| "IMPORT_MEDIA"
| "IMPORT_MEDIA_ACK"
| "PLAYBACK_SESSION_ENDED"
| "PLAYBACK_SESSION_PROGRESS"
| "ERROR";
@ -110,6 +112,26 @@ export interface RestartPayload {
reason?: string;
}
/** Master → node: move finished Download Station output into 4K library folder. */
export interface ImportMediaPayload {
jobId: string;
/** Absolute host path of the finished download folder/file */
sourcePath: string;
/** Absolute host path for destination movie folder (Title (Year)) */
destDir: string;
movieTitle: string;
movieYear?: number | null;
/** Prefer 2160p in renamed filename */
qualityHint?: string;
}
export interface ImportMediaAckPayload {
jobId: string;
ok: boolean;
importedPath?: string;
error?: string;
}
export interface PlaybackSessionEndedPayload {
sessionId: string;
reason: string;

View file

@ -71,6 +71,9 @@ importers:
nanoid:
specifier: ^5.0.9
version: 5.1.16
prisma:
specifier: ^6.3.1
version: 6.19.3(typescript@5.9.3)
stremio-addon-sdk:
specifier: ^1.6.10
version: 1.6.10(supports-color@5.5.0)
@ -84,9 +87,6 @@ importers:
'@types/ws':
specifier: ^8.5.14
version: 8.18.1
prisma:
specifier: ^6.3.1
version: 6.19.3(typescript@5.9.3)
tsx:
specifier: ^4.19.2
version: 4.23.12