| from __future__ import annotations |
|
|
| import ipaddress |
| import mimetypes |
| import shutil |
| import socket |
| from dataclasses import replace |
| from pathlib import Path |
| from urllib.parse import unquote, urlparse |
| from urllib.request import Request, urlopen |
|
|
| from renderer.core.config import Settings |
| from renderer.core.models import AIReelsRequest, RenderRequest, Scene |
| from renderer.core.utils import safe_filename |
|
|
|
|
| class IngestError(ValueError): |
| pass |
|
|
|
|
| class AssetIngestor: |
| """Resolve local paths and remote URLs into job-local files.""" |
|
|
| def __init__(self, settings: Settings) -> None: |
| self.settings = settings |
|
|
| def resolve_render_request(self, request: RenderRequest, workdir: Path) -> RenderRequest: |
| return replace( |
| request, |
| scenes=[ |
| replace(scene, media=str(self.resolve(scene.media, workdir / "inputs", f"scene_{idx:03d}"))) |
| for idx, scene in enumerate(request.scenes) |
| ], |
| voiceover=self.resolve_optional(request.voiceover, workdir / "inputs", "voiceover"), |
| background_music=self.resolve_optional(request.background_music, workdir / "inputs", "music"), |
| ) |
|
|
| def resolve_ai_reels_request(self, request: AIReelsRequest, workdir: Path) -> AIReelsRequest: |
| return replace( |
| request, |
| voiceover=str(self.resolve(request.voiceover, workdir / "inputs", "voiceover")), |
| assets=[str(self.resolve(asset, workdir / "inputs", f"asset_{idx:03d}")) for idx, asset in enumerate(request.assets)], |
| background_music=self.resolve_optional(request.background_music, workdir / "inputs", "music"), |
| ) |
|
|
| def resolve_optional(self, value: str | None, directory: Path, stem: str) -> str | None: |
| if not value: |
| return None |
| return str(self.resolve(value, directory, stem)) |
|
|
| def resolve(self, value: str, directory: Path, stem: str) -> Path: |
| directory.mkdir(parents=True, exist_ok=True) |
| if is_remote_url(value): |
| return self.download(value, directory, stem) |
| path = Path(value) |
| if not path.exists(): |
| raise IngestError(f"Asset does not exist: {value}") |
| return path |
|
|
| def download(self, url: str, directory: Path, stem: str) -> Path: |
| parsed = urlparse(url) |
| if parsed.scheme not in {"http", "https"} or not parsed.hostname: |
| raise IngestError("Only http and https asset URLs are supported") |
| if not self.settings.allow_private_asset_urls: |
| _reject_private_host(parsed.hostname) |
| request = Request(url, headers={"User-Agent": "basyx-ffmpeg-renderer/1.0"}) |
| with urlopen(request, timeout=self.settings.download_timeout_seconds) as response: |
| content_length = response.headers.get("Content-Length") |
| if content_length and int(content_length) > self.settings.max_download_bytes: |
| raise IngestError("Remote asset exceeds MAX_DOWNLOAD_BYTES") |
| suffix = _suffix_from_response(url, response.headers.get("Content-Type")) |
| target = directory / safe_filename(f"{stem}{suffix}") |
| total = 0 |
| with target.open("wb") as output: |
| while True: |
| chunk = response.read(1024 * 1024) |
| if not chunk: |
| break |
| total += len(chunk) |
| if total > self.settings.max_download_bytes: |
| target.unlink(missing_ok=True) |
| raise IngestError("Remote asset exceeds MAX_DOWNLOAD_BYTES") |
| output.write(chunk) |
| return target |
|
|
|
|
| def stage_upload(source: Path, uploads_dir: Path, filename: str) -> Path: |
| uploads_dir.mkdir(parents=True, exist_ok=True) |
| target = uploads_dir / safe_filename(filename) |
| if source.resolve() != target.resolve(): |
| shutil.copy2(source, target) |
| return target |
|
|
|
|
| def is_remote_url(value: str) -> bool: |
| return urlparse(value).scheme in {"http", "https"} |
|
|
|
|
| def _suffix_from_response(url: str, content_type: str | None) -> str: |
| path_suffix = Path(unquote(urlparse(url).path)).suffix |
| if path_suffix: |
| return path_suffix[:16] |
| if content_type: |
| guessed = mimetypes.guess_extension(content_type.split(";", 1)[0].strip()) |
| if guessed: |
| return guessed |
| return ".bin" |
|
|
|
|
| def _reject_private_host(hostname: str) -> None: |
| try: |
| addresses = socket.getaddrinfo(hostname, None) |
| except socket.gaierror as exc: |
| raise IngestError(f"Could not resolve host: {hostname}") from exc |
| for address in addresses: |
| ip = ipaddress.ip_address(address[4][0]) |
| if ip.is_private or ip.is_loopback or ip.is_link_local or ip.is_multicast: |
| raise IngestError("Private, loopback, link-local, and multicast asset hosts are not allowed") |
|
|