from fastapi import FastAPI, UploadFile, File, Form, Request from fastapi.responses import FileResponse, JSONResponse from contextlib import asynccontextmanager import os import uuid import asyncio import gradio as gr # ============================== # LOGGER + QUEUE # ============================== from utils.logger import logger from utils.job_queue import start_worker, create_job, get_job from ingestion.resolver import resolve_input # ============================== # AUTH SYSTEM (HARDENED IMPORT) # ============================== try: from auth.routes import router as auth_router from auth.database import Base, engine AUTH_ENABLED = True except Exception as e: logger.error(f"[AUTH BOOT FAILED] {e}") AUTH_ENABLED = False # ============================== # CORE PIPELINE # ============================== from utils.transcription import transcribe_video from utils.srt import generate_srt from utils.render import render_subtitles from utils.highlights import detect_highlights from utils.viral_scorer import score_clip from utils.director import rewrite_script, viral_score from utils.engagement import simulate_retention from utils.platform import adapt_platform from utils.persona import predict_audience from utils.clipper import create_clips from utils.autonomous_engine import run_autonomous_engine # ============================== # PUBLISHER # ============================== from publisher.publisher_ai import autonomous_loop from publisher.scheduler_engine import init_scheduler from publisher.platform_dispatcher import dispatch_publish from publisher.bulk import execute as bulk_execute from publisher.metadata_engine import generate_metadata from publisher.thumbnail_engine import generate_thumbnail # ============================== # SAFE bcrypt SHIELD (DO NOT FAIL BOOT) # ============================== try: import bcrypt except Exception as e: logger.warning(f"[bcrypt warning ignored] {e}") # ============================== # INIT # ============================== UPLOAD_DIR = "jobs" os.makedirs(UPLOAD_DIR, exist_ok=True) # ========================================================= # SAFE DB INITIALIZATION # ========================================================= def init_database_safe(): if not AUTH_ENABLED: logger.warning("Auth disabled - skipping DB init") return try: Base.metadata.create_all(bind=engine) logger.info("Database initialized successfully") except Exception as e: logger.error(f"Database init failed (non-fatal): {e}") # ========================================================= # LIFECYCLE # ========================================================= @asynccontextmanager async def lifespan(app: FastAPI): logger.info("Starting Basyx Whisper V10.1") # DB init (NON-FATAL) init_database_safe() # workers (must not block boot) try: start_worker() except Exception as e: logger.error(f"Worker failed: {e}") # scheduler try: init_scheduler() except Exception as e: logger.error(f"Scheduler failed: {e}") # autonomous engine (isolated task) try: asyncio.create_task(autonomous_loop()) except Exception as e: logger.error(f"Autonomous engine failed: {e}") yield logger.info("Shutdown complete") # ========================================================= # APP # ========================================================= app = FastAPI( title="Basyx Whisper V10.1 Autonomous Operator", lifespan=lifespan, ) # AUTH ROUTER (only if available) if AUTH_ENABLED: app.include_router(auth_router) # ========================================================= # TASKS # ========================================================= VALID_TASKS = { "autonomous", "auto-publish", "publish", "bulk-publish", "generate-metadata", "generate-thumbnail", "schedule-post", "transcribe", "subtitles", "render", "highlights", "viral-score", "strategy", "batch", "clips", } def normalize_task(task: str): task = task.lower().replace("_", "-") if task not in VALID_TASKS: raise Exception(f"Unknown task: {task}") return task # ========================================================= # SAFE INPUT RESOLVER # ========================================================= async def safe_resolve(file, source): try: if not file and not source: return None upload_file = file if isinstance(file, UploadFile) else None return await asyncio.to_thread(resolve_input, source, upload_file) except Exception as e: logger.error(f"resolve_input failed: {e}") return None # ========================================================= # EXECUTION ENGINE # ========================================================= async def execute_task(video_path, task, payload=None, webhook=None): payload = payload or {} if task == "bulk-publish": return await bulk_execute(payload), None if task not in ["bulk-publish", "schedule-post"] and not video_path: return {"error": "No valid input resolved"}, None if task == "autonomous": return await asyncio.to_thread(run_autonomous_engine, video_path), None if task == "auto-publish": auto = await asyncio.to_thread(run_autonomous_engine, video_path) return await dispatch_publish(variants=auto.get("all_variants", [])), None if task == "publish": return await dispatch_publish(video_path=video_path, payload=payload), None if task == "generate-metadata": return generate_metadata(video_path), None if task == "generate-thumbnail": output_path = os.path.join(UPLOAD_DIR, f"{uuid.uuid4()}.jpg") thumb = generate_thumbnail(video_path, output=output_path) return {"thumbnail": thumb}, output_path if task == "batch": job_id = create_job(video_path, webhook=webhook) return {"status": "queued", "job_id": job_id}, None if task == "transcribe": words = await asyncio.to_thread(transcribe_video, video_path) return {"words": words}, None if task == "subtitles": words = await asyncio.to_thread(transcribe_video, video_path) return {"srt": generate_srt(words)}, None if task == "render": words = await asyncio.to_thread(transcribe_video, video_path) srt = generate_srt(words) output = os.path.join(UPLOAD_DIR, f"{uuid.uuid4()}_render.mp4") await asyncio.to_thread(render_subtitles, video_path, srt, output) return {"status": "render_complete"}, output if task == "highlights": words = await asyncio.to_thread(transcribe_video, video_path) highlights = detect_highlights(words) or [] clips = create_clips(video_path, highlights) return {"clips_created": len(clips)}, (clips[0] if clips else None) if task == "clips": words = await asyncio.to_thread(transcribe_video, video_path) highlights = detect_highlights(words) or [] return {"clips": create_clips(video_path, highlights)}, None if task == "viral-score": words = await asyncio.to_thread(transcribe_video, video_path) segments = detect_highlights(words) or [] return {"scores": [score_clip(s) for s in segments]}, None if task == "strategy": words = await asyncio.to_thread(transcribe_video, video_path) script = rewrite_script(words) persona = predict_audience(words) curve = simulate_retention(words) return { "hook": script["hook"], "persona": persona, "viral_score": viral_score(curve), "platforms": { "tiktok": adapt_platform(script, "tiktok"), "reels": adapt_platform(script, "reels"), }, }, None return {"error": "Task execution failed"}, None # ========================================================= # ROUTER # ========================================================= @app.post("/execute/{task_name}") async def execute_router( request: Request, task_name: str, file: UploadFile = File(None), url_input: str = Form(None), source: str = Form(None), webhook: str = Form(None), ): try: task = normalize_task(task_name) payload = {} if request.headers.get("content-type", "").startswith("application/json"): payload = await request.json() video_path = await safe_resolve(file, url_input or source) result, output = await execute_task(video_path, task, payload, webhook) if output and isinstance(output, str) and os.path.exists(output): return FileResponse(output) return {"task": task, "result": result} except Exception as e: logger.exception(e) return JSONResponse({"error": str(e)}, status_code=500) # ========================================================= # HEALTH # ========================================================= @app.get("/api/health") def health(): return { "status": "online", "version": "V10.1", "auth_enabled": AUTH_ENABLED } @app.get("/api/status/{job_id}") def status(job_id: str): return get_job(job_id) or {"error": "Job not found"} # ========================================================= # GRADIO UI # ========================================================= async def ui_handler(video, task, webhook, url_input): source = url_input or video video_path = await safe_resolve(video, source) result, output = await execute_task( video_path, normalize_task(task), {}, webhook, ) return str(result), output with gr.Blocks() as demo: gr.Markdown("# 🚀 Basyx Whisper V10.1 Stable Operator") video_input = gr.Video() url_input = gr.Textbox(label="Video URL") task_dropdown = gr.Dropdown( choices=list(VALID_TASKS), value="autonomous", ) webhook_input = gr.Textbox(label="Webhook") run_btn = gr.Button("Execute") output_box = gr.Textbox() video_output = gr.Video() run_btn.click( ui_handler, [video_input, task_dropdown, webhook_input, url_input], [output_box, video_output], ) app = gr.mount_gradio_app(app, demo, path="/")