Moves git-compatible content from Nextcloud (Botomir/Projekte/Angefangen/Kino-Projekt) into this repo per cleanup request, under nextcloud-import/ to avoid colliding with the actively maintained backend/frontend/worker layout: - status/planning docs (Botomir-Status.md, Plan-Botomir.md, Projektbeschreibung.md) - docs/ (deployment.md, provider-policy.md) - backend-nextcloud-variant/ (older app/ layout, kept for reference only) - plugin.video.xstream/, script.module.xstreamscraper/ (vendored Kodi addon source) A third-party Firefox extension folder was intentionally left in Nextcloud (unrelated vendor code, not project source).
112 lines
4.9 KiB
Python
112 lines
4.9 KiB
Python
import asyncio
|
|
import json
|
|
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
|
|
|
|
|
|
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():
|
|
target = dest / path.name
|
|
if path.is_dir():
|
|
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
|
|
|
|
|
|
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:
|
|
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:
|
|
last_reported_progress = 0.0
|
|
|
|
def report_download_progress(provider_progress: float) -> None:
|
|
nonlocal last_reported_progress
|
|
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)
|
|
title = source.title or source.external_id or "media"
|
|
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,
|
|
)
|
|
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,
|
|
)
|
|
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")
|
|
update_job(session, job, target_path=str(target_dir), status="refreshing", progress=0.9)
|
|
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 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)
|