294 lines
12 KiB
Python
294 lines
12 KiB
Python
import asyncio
|
|
import hashlib
|
|
import json
|
|
import shlex
|
|
import shutil
|
|
from pathlib import Path
|
|
|
|
from sqlmodel import Session
|
|
|
|
from app.config import Settings
|
|
from app.models import ImportJob, Source, utcnow
|
|
from app.providers.base import DownloadOptions, Provider, ProviderError
|
|
from app.services.jellyfin import refresh_jellyfin
|
|
from app.services.paths import jellyfin_profile_target, tmp_target
|
|
|
|
|
|
class ImportCancelled(RuntimeError):
|
|
pass
|
|
|
|
|
|
def file_md5(path: Path, chunk_size: int = 1024 * 1024) -> str:
|
|
digest = hashlib.md5() # nosec B324 - integrity check, not security/auth
|
|
with path.open("rb") as fh:
|
|
while chunk := fh.read(chunk_size):
|
|
digest.update(chunk)
|
|
return digest.hexdigest()
|
|
|
|
|
|
def local_md5_manifest(base_dir: Path) -> dict[str, str]:
|
|
base_dir = base_dir.resolve()
|
|
manifest: dict[str, str] = {}
|
|
if not base_dir.exists():
|
|
return manifest
|
|
for path in sorted(p for p in base_dir.rglob("*") if p.is_file()):
|
|
if path.is_symlink():
|
|
raise ProviderError(f"Refusing to hash symlink: {path.relative_to(base_dir)}")
|
|
manifest[path.relative_to(base_dir).as_posix()] = file_md5(path)
|
|
return manifest
|
|
|
|
|
|
def job_title(source: Source | None) -> str:
|
|
if source is None:
|
|
return "media"
|
|
return source.title or source.external_id or "media"
|
|
|
|
|
|
def ensure_not_cancelled(session: Session, job: ImportJob) -> None:
|
|
session.refresh(job)
|
|
if job.status in {"cancelling", "cancelled"}:
|
|
raise ImportCancelled("Import abgebrochen")
|
|
|
|
|
|
async def remote_md5_manifest(remote_dir: str, settings: Settings) -> dict[str, str]:
|
|
if not settings.rsync_target or not settings.rsync_ssh_key:
|
|
return {}
|
|
ssh_cmd = [
|
|
"ssh",
|
|
"-i",
|
|
str(settings.rsync_ssh_key),
|
|
"-o",
|
|
"BatchMode=yes",
|
|
"-o",
|
|
"StrictHostKeyChecking=accept-new",
|
|
settings.rsync_target,
|
|
"cd " + shlex.quote(remote_dir) + " && find . -type f -print0 | sort -z | xargs -0 md5sum",
|
|
]
|
|
proc = await asyncio.create_subprocess_exec(
|
|
*ssh_cmd,
|
|
stdout=asyncio.subprocess.PIPE,
|
|
stderr=asyncio.subprocess.PIPE,
|
|
)
|
|
stdout, stderr = await proc.communicate()
|
|
if proc.returncode != 0:
|
|
raise ProviderError(f"Jellyfin md5 verification failed: {stderr.decode(errors='replace')[:500]}")
|
|
manifest: dict[str, str] = {}
|
|
for line in stdout.decode(errors="replace").splitlines():
|
|
if not line.strip():
|
|
continue
|
|
checksum, _, rel = line.partition(" ")
|
|
if not rel:
|
|
checksum, _, rel = line.partition(" ")
|
|
manifest[rel.removeprefix("./")] = checksum.strip()
|
|
return manifest
|
|
|
|
|
|
def remote_dir_from_target(target_path: str | None, settings: Settings) -> str | None:
|
|
if not target_path or not settings.rsync_target:
|
|
return None
|
|
prefix = f"{settings.rsync_target}:"
|
|
if target_path.startswith(prefix):
|
|
return target_path[len(prefix):]
|
|
return None
|
|
|
|
|
|
async def verify_jellyfin_transfer(local_dir: Path, target_path: str | None, settings: Settings) -> bool:
|
|
local_manifest = local_md5_manifest(local_dir)
|
|
if not local_manifest:
|
|
raise ProviderError("No local files available for md5 verification")
|
|
remote_dir = remote_dir_from_target(target_path, settings)
|
|
if remote_dir is None:
|
|
# No remote rsync target: the local media root is the final target.
|
|
return local_md5_manifest(local_dir) == local_manifest
|
|
remote_manifest = await remote_md5_manifest(remote_dir, settings)
|
|
if remote_manifest != local_manifest:
|
|
missing = sorted(set(local_manifest) - set(remote_manifest))[:5]
|
|
changed = sorted(k for k in local_manifest.keys() & remote_manifest.keys() if local_manifest[k] != remote_manifest[k])[:5]
|
|
raise ProviderError(f"Jellyfin md5 verification mismatch; missing={missing}, changed={changed}")
|
|
return True
|
|
|
|
|
|
async def cleanup_verified_download_data(tmp_dir: Path, target_dir: Path, target_path: str | None, settings: Settings) -> None:
|
|
if tmp_dir.exists():
|
|
shutil.rmtree(tmp_dir)
|
|
# When rsync copied to a remote Jellyfin host, local media_root is staging and can be removed after md5 proof.
|
|
if remote_dir_from_target(target_path, settings) is not None and target_dir.exists():
|
|
media_root = settings.media_root.resolve()
|
|
resolved = target_dir.resolve()
|
|
if resolved != media_root and media_root in resolved.parents:
|
|
shutil.rmtree(resolved)
|
|
|
|
|
|
|
|
def copy_tree_contents(src: Path, dest: Path) -> list[Path]:
|
|
dest.mkdir(parents=True, exist_ok=True)
|
|
copied: list[Path] = []
|
|
for path in src.iterdir():
|
|
if path.is_symlink():
|
|
raise ProviderError(f"Refusing to copy symlink from import output: {path.name}")
|
|
target = dest / path.name
|
|
if path.is_dir():
|
|
for child in path.rglob("*"):
|
|
if child.is_symlink():
|
|
raise ProviderError(f"Refusing to copy symlink from import output: {child.relative_to(src)}")
|
|
if target.exists():
|
|
shutil.rmtree(target)
|
|
shutil.copytree(path, target)
|
|
copied.extend(p for p in target.rglob("*") if p.is_file())
|
|
else:
|
|
shutil.copy2(path, target)
|
|
copied.append(target)
|
|
return copied
|
|
|
|
|
|
async def sync_to_jellyfin_vm(local_dir: Path, settings: Settings) -> str | None:
|
|
"""Copy the completed local import directory to Jellyfin's media VM via rsync."""
|
|
if not settings.rsync_target or not settings.rsync_ssh_key:
|
|
return None
|
|
|
|
local_dir = local_dir.resolve()
|
|
media_root = settings.media_root.resolve()
|
|
if media_root != local_dir and media_root not in local_dir.parents:
|
|
raise ProviderError(f"Refusing to sync path outside media root: {local_dir}")
|
|
|
|
relative = local_dir.relative_to(media_root)
|
|
remote_root = Path(str(settings.rsync_remote_root)).as_posix().rstrip("/") or "/jellyfin"
|
|
remote_dir = f"{remote_root}/{relative.as_posix()}"
|
|
ssh_cmd = [
|
|
"ssh",
|
|
"-i",
|
|
str(settings.rsync_ssh_key),
|
|
"-o",
|
|
"BatchMode=yes",
|
|
"-o",
|
|
"StrictHostKeyChecking=accept-new",
|
|
]
|
|
|
|
mkdir_proc = await asyncio.create_subprocess_exec(
|
|
*ssh_cmd,
|
|
settings.rsync_target,
|
|
f"mkdir -p {shlex.quote(remote_dir)}",
|
|
stdout=asyncio.subprocess.PIPE,
|
|
stderr=asyncio.subprocess.PIPE,
|
|
)
|
|
_, mkdir_stderr = await mkdir_proc.communicate()
|
|
if mkdir_proc.returncode != 0:
|
|
raise ProviderError(f"Jellyfin transfer mkdir failed: {mkdir_stderr.decode(errors='replace')[:500]}")
|
|
|
|
rsync_proc = await asyncio.create_subprocess_exec(
|
|
"rsync",
|
|
"-a",
|
|
"--delete",
|
|
"-s",
|
|
"-e",
|
|
" ".join(shlex.quote(part) for part in ssh_cmd),
|
|
f"{local_dir}/",
|
|
f"{settings.rsync_target}:{remote_dir}/",
|
|
stdout=asyncio.subprocess.PIPE,
|
|
stderr=asyncio.subprocess.PIPE,
|
|
)
|
|
_, rsync_stderr = await rsync_proc.communicate()
|
|
if rsync_proc.returncode != 0:
|
|
raise ProviderError(f"Jellyfin transfer rsync failed: {rsync_stderr.decode(errors='replace')[:500]}")
|
|
|
|
return f"{settings.rsync_target}:{remote_dir}"
|
|
|
|
|
|
def update_job(session: Session, job: ImportJob, **fields) -> None:
|
|
for key, value in fields.items():
|
|
setattr(job, key, value)
|
|
job.updated_at = utcnow()
|
|
session.add(job)
|
|
session.commit()
|
|
session.refresh(job)
|
|
|
|
|
|
async def run_import_job(job_id: int, engine, providers: list[Provider], settings: Settings) -> None:
|
|
tmp_dir: Path | None = None
|
|
target_dir: Path | None = None
|
|
with Session(engine) as session:
|
|
job = session.get(ImportJob, job_id)
|
|
if job is None:
|
|
return
|
|
source = session.get(Source, job.source_id)
|
|
if source is None:
|
|
update_job(session, job, status="failed", error="Source missing")
|
|
return
|
|
provider = next((p for p in providers if p.name == source.provider or p.can_handle(source.url)), None)
|
|
if provider is None:
|
|
update_job(session, job, status="failed", error="No provider available for source")
|
|
return
|
|
try:
|
|
ensure_not_cancelled(session, job)
|
|
last_reported_progress = 0.0
|
|
|
|
def report_download_progress(provider_progress: float) -> None:
|
|
nonlocal last_reported_progress
|
|
ensure_not_cancelled(session, job)
|
|
mapped = 0.15 + (max(0.0, min(provider_progress, 1.0)) * 0.60)
|
|
if mapped - last_reported_progress >= 0.01 or mapped >= 0.75:
|
|
last_reported_progress = mapped
|
|
update_job(session, job, status="downloading", progress=round(mapped, 3))
|
|
|
|
update_job(session, job, status="downloading", progress=0.15, error=None)
|
|
title = job_title(source)
|
|
stored_metadata = json.loads(source.metadata_json or "{}")
|
|
import_request = stored_metadata.get("import_request", {}) if isinstance(stored_metadata, dict) else {}
|
|
provider_metadata = stored_metadata.get("metadata", stored_metadata) if isinstance(stored_metadata, dict) else {}
|
|
tmp_dir = tmp_target(settings.tmp_dir, job.id or 0, title)
|
|
result = await provider.download(
|
|
source.url,
|
|
tmp_dir,
|
|
DownloadOptions(
|
|
max_height=settings.default_max_height,
|
|
target_library=job.target_library,
|
|
max_bytes=settings.max_download_bytes,
|
|
),
|
|
progress_callback=report_download_progress,
|
|
)
|
|
ensure_not_cancelled(session, job)
|
|
update_job(session, job, status="postprocessing", progress=0.75)
|
|
target_dir = jellyfin_profile_target(
|
|
settings.media_root,
|
|
job.target_library,
|
|
import_request.get("title") or title,
|
|
year=import_request.get("year"),
|
|
series_title=import_request.get("series_title"),
|
|
season=import_request.get("season"),
|
|
episode=import_request.get("episode"),
|
|
episode_title=import_request.get("episode_title"),
|
|
channel=provider_metadata.get("uploader") or provider_metadata.get("channel") or source.provider,
|
|
external_id=source.external_id,
|
|
)
|
|
update_job(session, job, status="copying", progress=0.82)
|
|
copied = copy_tree_contents(tmp_dir, target_dir)
|
|
if not copied and not result.output_files:
|
|
raise ProviderError("No files were produced by the import")
|
|
ensure_not_cancelled(session, job)
|
|
remote_target = await sync_to_jellyfin_vm(target_dir, settings)
|
|
final_target = remote_target or str(target_dir)
|
|
update_job(session, job, target_path=final_target, status="refreshing", progress=0.9)
|
|
verify_ok = await verify_jellyfin_transfer(target_dir, final_target, settings)
|
|
if verify_ok:
|
|
await cleanup_verified_download_data(tmp_dir, target_dir, final_target, settings)
|
|
refreshed = await refresh_jellyfin(settings)
|
|
note = None if refreshed else "Import completed; Jellyfin refresh skipped because no API key is configured"
|
|
update_job(session, job, status="done", progress=1.0, error=note)
|
|
except ImportCancelled:
|
|
if tmp_dir and tmp_dir.exists():
|
|
shutil.rmtree(tmp_dir)
|
|
update_job(session, job, status="cancelled", error="Import abgebrochen; temporäre Download-Dateien gelöscht")
|
|
except Exception as exc: # keep background task failures visible in DB/API
|
|
update_job(session, job, status="failed", error=str(exc)[:1000])
|
|
|
|
|
|
def schedule_import(job_id: int, engine, providers: list[Provider], settings: Settings) -> None:
|
|
asyncio.create_task(run_import_job(job_id, engine, providers, settings))
|
|
|
|
|
|
def metadata_json(metadata, *, import_request: dict | None = None) -> str:
|
|
payload = metadata.model_dump()
|
|
if import_request is not None:
|
|
payload = {"metadata": payload, "import_request": import_request}
|
|
return json.dumps(payload, ensure_ascii=False)
|