Spaces:
Sleeping
Sleeping
| """ | |
| Video Download Middleware API | |
| ============================== | |
| Flow: | |
| 1. POST /download → triggers GitHub Action → returns job_id | |
| 2. GET /status/{job_id} → returns pending / processing / done / failed | |
| 3. GET /file/{job_id} → streams the video file (only when done) | |
| n8n webhooks: | |
| - POST /n8n/callback → n8n calls this when the video is ready (from "downlode video1" node) | |
| """ | |
| import os | |
| import uuid | |
| import time | |
| import asyncio | |
| import logging | |
| from pathlib import Path | |
| from typing import Optional | |
| import httpx | |
| import aiofiles | |
| from fastapi import FastAPI, HTTPException, BackgroundTasks, Request | |
| from fastapi.responses import StreamingResponse, JSONResponse | |
| from pydantic import BaseModel | |
| # ─── Config ──────────────────────────────────────────────────────────────────── | |
| GITHUB_TOKEN = os.getenv("GITHUB_TOKEN", "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") | |
| GITHUB_REPO = os.getenv("GITHUB_REPO", "expher510/AutoClip-Pipeline") | |
| WORKFLOW_FILE = os.getenv("WORKFLOW_FILE", "video-download.yml") | |
| # Public URL of THIS api (so n8n can call back) | |
| # e.g. https://my-api.hf.space or https://xxxxx.ngrok.io | |
| PUBLIC_URL = os.getenv("PUBLIC_URL", "https://your-api-public-url.com") | |
| # اختياري - لو موجود الـ API هو اللي يشغل n8n | |
| # لو None، n8n هو اللي يبعت callback للـ API | |
| N8N_TRIGGER_URL = os.getenv("N8N_TRIGGER_URL", None) | |
| VIDEOS_DIR = Path("videos") | |
| VIDEOS_DIR.mkdir(exist_ok=True) | |
| # ─── In-memory job store ──────────────────────────────────────────────────────── | |
| # job_id → {status, filename, file_path, video_url, error, created_at, run_id} | |
| jobs: dict[str, dict] = {} | |
| # ─── App ─────────────────────────────────────────────────────────────────────── | |
| app = FastAPI( | |
| title="EGDownloader API", | |
| description="Bridges client ↔ n8n/GitHub Action video pipeline", | |
| version="1.0.0", | |
| docs_url="/", | |
| redoc_url="/redoc", | |
| ) | |
| logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s") | |
| log = logging.getLogger(__name__) | |
| # ══════════════════════════════════════════════════════════════════════════════ | |
| # 1. REQUEST DOWNLOAD | |
| # ══════════════════════════════════════════════════════════════════════════════ | |
| class DownloadRequest(BaseModel): | |
| url: str | |
| cookies_txt: Optional[str] = "" # optional – paste Netscape cookie text | |
| async def download_shorts(body: DownloadRequest): | |
| """ | |
| نفس /download بالظبط، بس بيحوّل Shorts URL لـ watch URL الأول. | |
| https://youtube.com/shorts/VIDEO_ID → https://youtube.com/watch?v=VIDEO_ID | |
| """ | |
| if "/shorts/" in body.url: | |
| video_id = body.url.split("/shorts/")[1].split("?")[0] | |
| body.url = f"https://www.youtube.com/watch?v={video_id}" | |
| return await start_download(body) | |
| async def start_download(body: DownloadRequest): | |
| """ | |
| طريقتين للشغل: | |
| 1) Push mode (الـ default - اللي عندك دلوقتي): | |
| N8N_TRIGGER_URL = None | |
| n8n هو اللي يشغل GitHub Action وبعدين يبعت callback على /n8n/callback/{job_id} | |
| الـ client بس يعمل polling على /status/{job_id} | |
| 2) Trigger mode: | |
| N8N_TRIGGER_URL = "https://..." | |
| الـ API هو اللي يبعت trigger لـ n8n، وn8n يرد بـ callback | |
| """ | |
| job_id = str(uuid.uuid4()) | |
| jobs[job_id] = { | |
| "status": "pending", | |
| "video_url": body.url, | |
| "filename": None, | |
| "file_path": None, | |
| "run_id": None, | |
| "error": None, | |
| "created_at": time.time(), | |
| "mode": "trigger" if N8N_TRIGGER_URL else "push", | |
| } | |
| # ── Trigger mode: الـ API يشغل n8n ──────────────────────────────────────── | |
| if N8N_TRIGGER_URL: | |
| payload = { | |
| "url": body.url, | |
| "cookies_txt": body.cookies_txt or "", | |
| "n8n_webhook": f"{PUBLIC_URL}/n8n/callback/{job_id}", | |
| } | |
| try: | |
| async with httpx.AsyncClient(timeout=15) as client: | |
| r = await client.post(N8N_TRIGGER_URL, json=payload) | |
| r.raise_for_status() | |
| jobs[job_id]["status"] = "processing" | |
| log.info(f"[{job_id}] triggered n8n → {r.status_code}") | |
| except Exception as exc: | |
| jobs[job_id]["status"] = "failed" | |
| jobs[job_id]["error"] = str(exc) | |
| log.error(f"[{job_id}] failed to trigger n8n: {exc}") | |
| raise HTTPException(status_code=502, detail=f"Could not trigger pipeline: {exc}") | |
| # ── Push mode: n8n هو اللي هيبعت callback لما يخلص ─────────────────────── | |
| else: | |
| jobs[job_id]["status"] = "pending" | |
| log.info(f"[{job_id}] job created (push mode) - waiting for n8n callback") | |
| return { | |
| "job_id": job_id, | |
| "status": jobs[job_id]["status"], | |
| "mode": jobs[job_id]["mode"], | |
| "callback_url": f"{PUBLIC_URL}/n8n/callback/{job_id}", | |
| } | |
| # ══════════════════════════════════════════════════════════════════════════════ | |
| # 2. POLL STATUS | |
| # ══════════════════════════════════════════════════════════════════════════════ | |
| async def get_status(job_id: str): | |
| """ | |
| Returns one of: pending | processing | done | failed | |
| When done, also returns filename & download_url. | |
| """ | |
| job = jobs.get(job_id) | |
| if not job: | |
| raise HTTPException(status_code=404, detail="Job not found") | |
| resp = { | |
| "job_id": job_id, | |
| "status": job["status"], | |
| "video_url": job["video_url"], | |
| "created_at": job["created_at"], | |
| } | |
| if job["status"] == "done": | |
| resp["filename"] = job["filename"] | |
| resp["download_url"] = f"{PUBLIC_URL}/file/{job_id}" | |
| if job["status"] == "failed": | |
| resp["error"] = job["error"] | |
| return resp | |
| # ══════════════════════════════════════════════════════════════════════════════ | |
| # 3. DOWNLOAD FILE | |
| # ══════════════════════════════════════════════════════════════════════════════ | |
| async def download_file(job_id: str): | |
| job = jobs.get(job_id) | |
| if not job: | |
| raise HTTPException(status_code=404, detail="Job not found") | |
| if job["status"] != "done": | |
| raise HTTPException(status_code=425, detail=f"Not ready yet – status: {job['status']}") | |
| file_path: Path = job["file_path"] | |
| if not file_path or not file_path.exists(): | |
| raise HTTPException(status_code=410, detail="File no longer available on server") | |
| filename = job["filename"] or file_path.name | |
| async def file_streamer(): | |
| async with aiofiles.open(file_path, "rb") as f: | |
| while chunk := await f.read(1024 * 512): # 512 KB chunks | |
| yield chunk | |
| return StreamingResponse( | |
| file_streamer(), | |
| media_type="video/mp4", | |
| headers={"Content-Disposition": f'attachment; filename="{filename}"'}, | |
| ) | |
| # ══════════════════════════════════════════════════════════════════════════════ | |
| # 4. n8n CALLBACK (n8n POSTs the finished binary here) | |
| # This replaces the "downlode video1" webhook that was in n8n | |
| # ══════════════════════════════════════════════════════════════════════════════ | |
| async def n8n_callback(job_id: str, request: Request): | |
| """ | |
| n8n calls this endpoint (instead of its own 'downlode video1' webhook). | |
| It sends either: | |
| a) A JSON body with {status, run_id, filename, video_url} → we then | |
| fetch the artifact from GitHub ourselves, OR | |
| b) Raw binary (the video file) directly in the body. | |
| """ | |
| job = jobs.get(job_id) | |
| if not job: | |
| log.warning(f"[{job_id}] callback for unknown job") | |
| return {"ok": False, "reason": "unknown job"} | |
| content_type = request.headers.get("content-type", "") | |
| # ── Case A: JSON notification (status + run_id) ────────────────────────── | |
| if "application/json" in content_type: | |
| data = await request.json() | |
| log.info(f"[{job_id}] JSON callback: {data}") | |
| if data.get("status") == "success": | |
| run_id = data.get("run_id") | |
| filename = data.get("filename", f"{job_id}.mp4") | |
| jobs[job_id]["run_id"] = run_id | |
| jobs[job_id]["filename"] = filename | |
| # Fetch the artifact binary from GitHub in the background | |
| asyncio.create_task(_fetch_github_artifact(job_id, run_id, filename)) | |
| else: | |
| jobs[job_id]["status"] = "failed" | |
| jobs[job_id]["error"] = data.get("error", "n8n reported failure") | |
| return {"ok": True} | |
| # ── Case B: Raw binary video ────────────────────────────────────────────── | |
| body = await request.body() | |
| if not body: | |
| jobs[job_id]["status"] = "failed" | |
| jobs[job_id]["error"] = "empty callback body" | |
| return {"ok": False} | |
| filename = request.headers.get("x-filename", f"{job_id}.mp4") | |
| file_path = VIDEOS_DIR / f"{job_id}_{filename}" | |
| async with aiofiles.open(file_path, "wb") as f: | |
| await f.write(body) | |
| jobs[job_id]["status"] = "done" | |
| jobs[job_id]["filename"] = filename | |
| jobs[job_id]["file_path"] = file_path | |
| log.info(f"[{job_id}] saved {file_path} ({len(body)/1024/1024:.1f} MB)") | |
| return {"ok": True} | |
| # ══════════════════════════════════════════════════════════════════════════════ | |
| # 5. BACKGROUND: fetch artifact from GitHub (used in Case A above) | |
| # ══════════════════════════════════════════════════════════════════════════════ | |
| async def _fetch_github_artifact(job_id: str, run_id: str, filename: str): | |
| """Polls GitHub until the artifact is ready, then downloads it.""" | |
| artifacts_url = f"https://api.github.com/repos/{GITHUB_REPO}/actions/runs/{run_id}/artifacts" | |
| headers = { | |
| "Authorization": f"token {GITHUB_TOKEN}", | |
| "Accept": "application/vnd.github+json", | |
| } | |
| for attempt in range(20): # up to ~5 minutes | |
| await asyncio.sleep(15) | |
| try: | |
| async with httpx.AsyncClient(timeout=30) as client: | |
| r = await client.get(artifacts_url, headers=headers) | |
| r.raise_for_status() | |
| artifacts = r.json().get("artifacts", []) | |
| if not artifacts: | |
| log.info(f"[{job_id}] no artifacts yet (attempt {attempt+1})") | |
| continue | |
| artifact = artifacts[0] | |
| dl_url = artifact["archive_download_url"] | |
| # Download the zip (GitHub wraps artifacts in a zip) | |
| async with httpx.AsyncClient( | |
| timeout=300, | |
| follow_redirects=True, | |
| headers=headers, | |
| ) as client: | |
| resp = await client.get(dl_url) | |
| resp.raise_for_status() | |
| # Save & unzip | |
| zip_path = VIDEOS_DIR / f"{job_id}.zip" | |
| async with aiofiles.open(zip_path, "wb") as f: | |
| await f.write(resp.content) | |
| # Extract the video | |
| import zipfile | |
| with zipfile.ZipFile(zip_path, "r") as z: | |
| names = z.namelist() | |
| video_name = next((n for n in names if n.endswith(".mp4")), names[0]) | |
| z.extract(video_name, VIDEOS_DIR) | |
| extracted = VIDEOS_DIR / video_name | |
| zip_path.unlink(missing_ok=True) | |
| jobs[job_id]["status"] = "done" | |
| jobs[job_id]["filename"] = video_name | |
| jobs[job_id]["file_path"] = extracted | |
| log.info(f"[{job_id}] artifact saved → {extracted}") | |
| return | |
| except Exception as exc: | |
| log.error(f"[{job_id}] artifact fetch error: {exc}") | |
| jobs[job_id]["status"] = "failed" | |
| jobs[job_id]["error"] = "artifact never became available" | |
| # ══════════════════════════════════════════════════════════════════════════════ | |
| # 6. EXTRA UTILS | |
| # ══════════════════════════════════════════════════════════════════════════════ | |
| async def list_jobs(): | |
| return { | |
| jid: {k: v for k, v in j.items() if k != "file_path"} | |
| for jid, j in jobs.items() | |
| } | |
| async def delete_job(job_id: str): | |
| job = jobs.pop(job_id, None) | |
| if not job: | |
| raise HTTPException(status_code=404, detail="Job not found") | |
| if job.get("file_path") and Path(job["file_path"]).exists(): | |
| Path(job["file_path"]).unlink() | |
| return {"deleted": job_id} | |
| async def health(): | |
| return {"status": "ok", "jobs": len(jobs)} |