Whisper / utils /job_queue.py
basyx's picture
Update utils/job_queue.py
9a9e289 verified
Raw
History Blame
5.13 kB
import threading
import uuid
from queue import Queue
import traceback
import time
from .logger import logger
from .validators import validate_video
from .transcription import transcribe_video
from .director import rewrite_script, viral_score
from .engagement import simulate_retention
from .platform import adapt_platform
from .persona import predict_audience
from .clipper import create_clip
# =====================================================
# GLOBAL STATE
# =====================================================
jobs = {}
queue = Queue()
WORKER_STARTED = False
# =====================================================
# CREATE DIRECTOR JOB
# =====================================================
def create_job(video_path: str, webhook: str | None = None):
job_id = str(uuid.uuid4())
jobs[job_id] = {
"id": job_id,
"status": "queued",
"stage": "waiting",
"progress": 0,
"video": video_path,
# V7 INTELLIGENCE OUTPUTS
"viral_score": None,
"persona": None,
"hook": None,
"platforms": [],
"strategy_summary": None,
"clips": [],
"webhook": webhook,
"error": None,
"created_at": time.time(),
}
queue.put(job_id)
logger.info(f"[V7] Director job queued: {job_id}")
return job_id
# =====================================================
# GET JOB
# =====================================================
def get_job(job_id: str):
return jobs.get(job_id)
# =====================================================
# UPDATE HELPERS
# =====================================================
def update(job_id, **kwargs):
if job_id in jobs:
jobs[job_id].update(kwargs)
# =====================================================
# V7 AUTONOMOUS DIRECTOR WORKER
# =====================================================
def worker():
logger.info("[V7] Autonomous Viral Director started")
while True:
job_id = queue.get()
job = jobs[job_id]
try:
# -----------------------------
# 1. TRANSCRIPTION
# -----------------------------
update(job_id, status="processing", stage="transcribing", progress=10)
words = transcribe_video(job["video"])
# -----------------------------
# 2. SCRIPT RECONSTRUCTION
# -----------------------------
update(job_id, stage="rewriting narrative", progress=25)
script = rewrite_script(words)
# -----------------------------
# 3. AUDIENCE MODELING
# -----------------------------
persona = predict_audience(words)
# -----------------------------
# 4. ENGAGEMENT SIMULATION
# -----------------------------
update(job_id, stage="simulating audience", progress=45)
curve = simulate_retention(words)
v_score = viral_score(curve)
# -----------------------------
# 5. PLATFORM STRATEGY
# -----------------------------
update(job_id, stage="platform adaptation", progress=65)
tiktok = adapt_platform(script, "tiktok")
reels = adapt_platform(script, "reels")
platforms = ["tiktok", "reels"]
# -----------------------------
# 6. SINGLE BEST OUTPUT (DIRECTOR DECISION)
# -----------------------------
update(job_id, stage="rendering final cut", progress=85)
clip = create_clip(
job["video"],
words[0]["start"],
words[-1]["end"],
0
)
# -----------------------------
# 7. FINAL DIRECTOR OUTPUT
# -----------------------------
update(job_id,
status="completed",
stage="director finished",
progress=100,
viral_score=v_score,
persona=persona,
hook=script["hook"],
platforms=platforms,
strategy_summary={
"curve_peak": max(curve),
"avg_curve": sum(curve) / len(curve),
"decision": "auto-selected best full narrative cut"
},
clips=[clip])
logger.info(f"[V7] Director output complete: {job_id}")
except Exception as e:
logger.error(traceback.format_exc())
update(job_id,
status="failed",
stage="error",
error=str(e))
finally:
queue.task_done()
# =====================================================
# START WORKER (SINGLETON SAFE)
# =====================================================
def start_worker():
global WORKER_STARTED
if WORKER_STARTED:
return
WORKER_STARTED = True
t = threading.Thread(target=worker, daemon=True)
t.start()
logger.info("[V7] Worker initialized")