""" orchestrator.py - Durable task orchestrator with state machine, checkpoints, and resume-after-sleep. Each video goes through a pipeline state machine: DISCOVERED → DOWNLOADING → TRANSCRIBING → TRANSLATING → VOICING → RENDERING → QA → UPLOADING → PUBLISHED All state transitions persist to HF Dataset via StateManager, so work survives Space sleep/restart. Lease-based claiming prevents duplicate processing. Integration: Used by agent_loop.py to drive the autonomous pipeline. """ import json import os import threading import time import uuid from datetime import datetime, timedelta from typing import Any, Dict, List, Optional # ================================================================ # Pipeline Stage State Machine # ================================================================ STAGES = { "discovered": {"next": "downloading", "rollback": None}, "downloading": {"next": "transcribing", "rollback": "discovered"}, "transcribing": {"next": "translating", "rollback": "downloading"}, "translating": {"next": "voicing", "rollback": "transcribing"}, "voicing": {"next": "rendering", "rollback": "translating"}, "rendering": {"next": "qa", "rollback": "voicing"}, "qa": {"next": "uploading", "rollback": "voicing"}, # QA fail → re-voice "uploading": {"next": "published", "rollback": "rendering"}, "published": {"next": None, "rollback": None}, # Terminal "failed": {"next": None, "rollback": None}, # Terminal } TERMINAL_STAGES = {"published", "failed"} LEASE_DURATION_MINUTES = 10 MAX_RETRIES = 3 # Instance identifier (unique per Space restart) INSTANCE_ID = os.getenv("HOSTNAME", f"instance-{uuid.uuid4().hex[:8]}") class TaskOrchestrator: """Durable task orchestrator with state machine and checkpoint support.""" def __init__(self, state, brain=None): self.state = state # StateManager self.brain = brain # NexBrain (optional, for retry decisions) self._lock = threading.Lock() self._ensure_task_queue() # ================================================================ # Public API # ================================================================ def add_task(self, video_id: str, title: str, url: str, priority: int = 0) -> bool: """Add a new video to the task queue. Returns True if added, False if video already exists in queue or was already processed. """ with self._lock: task_queue = self._get_queue() # Skip if already processed if self.state and hasattr(self.state, "_state"): if video_id in self.state._state.get("processed_videos", {}): print(f"[ORCHESTRATOR] Skipping {video_id}: already processed") return False # Skip if already in queue for task in task_queue: if task["video_id"] == video_id: print(f"[ORCHESTRATOR] Skipping {video_id}: already in queue (stage={task['stage']})") return False # Skip if permanently failed if self.state and hasattr(self.state, "_state"): failed = self.state._state.get("failed_videos", {}).get(video_id, {}) if failed.get("retry_count", 0) >= MAX_RETRIES: print(f"[ORCHESTRATOR] Skipping {video_id}: permanently failed ({failed.get('retry_count', 0)} retries)") return False now = datetime.now().isoformat() task = { "video_id": video_id, "title": title, "url": url, "stage": "discovered", "priority": priority, "status": "pending", "claimed_at": None, "lease_until": None, "claimed_by": None, "retry_count": 0, "checkpoint": { "download_path": None, "audio_path": None, "transcript": None, "translated": None, "voice_path": None, "final_path": None, }, "created_at": now, "updated_at": now, "error": None, "stage_failed": None, } task_queue.append(task) self._save_queue(task_queue) print(f"[ORCHESTRATOR] Added task: {video_id} ({title})") return True def claim_next_task(self) -> Optional[Dict[str, Any]]: """Claim the next pending task with lease-based locking. Returns the task dict or None if no tasks available. Lease duration: 10 minutes. After that, task can be reclaimed. """ with self._lock: # First, release any expired leases self._release_expired_leases(task_queue_only=True) task_queue = self._get_queue() now = datetime.now() # Sort by priority (higher first), then by creation time (older first) pending = [ t for t in task_queue if t["status"] == "pending" and t["stage"] not in TERMINAL_STAGES ] pending.sort(key=lambda t: (-t["priority"], t["created_at"])) if not pending: return None task = pending[0] lease_until = now + timedelta(minutes=LEASE_DURATION_MINUTES) # Update task in queue for i, t in enumerate(task_queue): if t["video_id"] == task["video_id"]: task_queue[i]["status"] = "claimed" task_queue[i]["claimed_at"] = now.isoformat() task_queue[i]["lease_until"] = lease_until.isoformat() task_queue[i]["claimed_by"] = INSTANCE_ID task_queue[i]["updated_at"] = now.isoformat() task = task_queue[i] break self._save_queue(task_queue) print(f"[ORCHESTRATOR] Claimed task: {task['video_id']} (stage={task['stage']}, lease until {lease_until.isoformat()})") return dict(task) # Return a copy def advance_stage(self, video_id: str, checkpoint_data: Optional[Dict] = None) -> bool: """Advance a task to the next pipeline stage. Args: video_id: The video task ID checkpoint_data: Optional dict to merge into task checkpoint Returns: True if stage was advanced, False otherwise """ with self._lock: task_queue = self._get_queue() for i, task in enumerate(task_queue): if task["video_id"] != video_id: continue current_stage = task["stage"] if current_stage in TERMINAL_STAGES: print(f"[ORCHESTRATOR] Cannot advance {video_id}: in terminal stage '{current_stage}'") return False stage_info = STAGES.get(current_stage) if not stage_info or not stage_info["next"]: print(f"[ORCHESTRATOR] Cannot advance {video_id}: no next stage for '{current_stage}'") return False next_stage = stage_info["next"] task_queue[i]["stage"] = next_stage task_queue[i]["updated_at"] = datetime.now().isoformat() task_queue[i]["error"] = None task_queue[i]["stage_failed"] = None # If we reached published, mark as completed if next_stage == "published": task_queue[i]["status"] = "completed" task_queue[i]["lease_until"] = None task_queue[i]["claimed_by"] = None # Merge checkpoint data if checkpoint_data: current_checkpoint = task_queue[i].get("checkpoint", {}) current_checkpoint.update(checkpoint_data) task_queue[i]["checkpoint"] = current_checkpoint self._save_queue(task_queue) print(f"[ORCHESTRATOR] Advanced {video_id}: {current_stage} → {next_stage}") return True print(f"[ORCHESTRATOR] Task not found: {video_id}") return False def fail_task(self, video_id: str, error: str, stage_failed: Optional[str] = None) -> bool: """Mark a task as failed with error details. The task is NOT removed from the queue. It can be retried later. """ with self._lock: task_queue = self._get_queue() for i, task in enumerate(task_queue): if task["video_id"] != video_id: continue retry_count = task.get("retry_count", 0) + 1 task_queue[i]["retry_count"] = retry_count task_queue[i]["error"] = error[:500] task_queue[i]["stage_failed"] = stage_failed or task["stage"] task_queue[i]["updated_at"] = datetime.now().isoformat() if retry_count >= MAX_RETRIES: task_queue[i]["status"] = "failed" task_queue[i]["stage"] = "failed" task_queue[i]["lease_until"] = None task_queue[i]["claimed_by"] = None print(f"[ORCHESTRATOR] Task {video_id} PERMANENTLY FAILED after {retry_count} attempts: {error[:100]}") else: # Reset to pending for retry task_queue[i]["status"] = "pending" task_queue[i]["lease_until"] = None task_queue[i]["claimed_by"] = None # Rollback to previous stage if possible stage_info = STAGES.get(task["stage"], {}) rollback = stage_info.get("rollback") if rollback: task_queue[i]["stage"] = rollback print(f"[ORCHESTRATOR] Task {video_id} failed (attempt {retry_count}/{MAX_RETRIES}): {error[:100]}") self._save_queue(task_queue) # Also update failed_videos in state if self.state and hasattr(self.state, "_state"): if "failed_videos" not in self.state._state: self.state._state["failed_videos"] = {} self.state._state["failed_videos"][video_id] = { "error": error[:500], "stage_failed": stage_failed or task["stage"], "retry_count": retry_count, "last_failed_at": datetime.now().isoformat(), "first_failed_at": task.get("created_at", datetime.now().isoformat()), } return True print(f"[ORCHESTRATOR] Task not found: {video_id}") return False def retry_task(self, video_id: str) -> bool: """Move a failed task back to the queue for retry. Returns True if the task was queued for retry, False if max retries reached. """ with self._lock: task_queue = self._get_queue() for i, task in enumerate(task_queue): if task["video_id"] != video_id: continue retry_count = task.get("retry_count", 0) if retry_count >= MAX_RETRIES: print(f"[ORCHESTRATOR] Cannot retry {video_id}: max retries ({MAX_RETRIES}) reached") return False # Exponential backoff: wait 2^retry_count minutes before retry backoff_minutes = 2 ** retry_count backoff_until = datetime.now() + timedelta(minutes=backoff_minutes) task_queue[i]["status"] = "pending" task_queue[i]["stage"] = "discovered" # Start from scratch task_queue[i]["lease_until"] = None task_queue[i]["claimed_by"] = None task_queue[i]["checkpoint"] = { "download_path": None, "audio_path": None, "transcript": None, "translated": None, "voice_path": None, "final_path": None, } task_queue[i]["updated_at"] = datetime.now().isoformat() # Store backoff_until for smart scheduling self._save_queue(task_queue) print(f"[ORCHESTRATOR] Retry queued for {video_id} (attempt {retry_count + 1}, backoff {backoff_minutes}min)") return True return False def get_task(self, video_id: str) -> Optional[Dict[str, Any]]: """Get a specific task by video ID.""" task_queue = self._get_queue() for task in task_queue: if task["video_id"] == video_id: return dict(task) return None def get_active_tasks(self) -> List[Dict[str, Any]]: """Get all tasks currently being processed (claimed).""" task_queue = self._get_queue() now = datetime.now() active = [] for task in task_queue: if task["status"] == "claimed": # Check if lease expired lease_until = task.get("lease_until") if lease_until: try: if datetime.fromisoformat(lease_until) < now: continue # Expired, not truly active except (ValueError, TypeError): pass active.append(dict(task)) return active def get_pending_tasks(self) -> List[Dict[str, Any]]: """Get tasks waiting in queue, sorted by priority (high first) then creation time.""" task_queue = self._get_queue() pending = [dict(t) for t in task_queue if t["status"] == "pending" and t["stage"] not in TERMINAL_STAGES] pending.sort(key=lambda t: (-t["priority"], t["created_at"])) return pending def get_completed_tasks(self) -> List[Dict[str, Any]]: """Get completed tasks.""" task_queue = self._get_queue() return [dict(t) for t in task_queue if t["status"] == "completed"] def get_failed_tasks(self) -> List[Dict[str, Any]]: """Get failed tasks.""" task_queue = self._get_queue() return [dict(t) for t in task_queue if t["status"] == "failed"] def resume_interrupted(self) -> List[Dict[str, Any]]: """Find tasks with expired leases and release them for re-claiming. Returns list of released tasks. """ with self._lock: return self._release_expired_leases(task_queue_only=False) def cleanup_completed(self, max_age_hours: int = 24) -> int: """Remove completed tasks older than max_age_hours. Returns the number of tasks removed. """ with self._lock: task_queue = self._get_queue() now = datetime.now() cutoff = now - timedelta(hours=max_age_hours) to_remove = [] for i, task in enumerate(task_queue): if task["status"] != "completed": continue updated = task.get("updated_at", "") try: updated_dt = datetime.fromisoformat(updated) if updated_dt < cutoff: to_remove.append(i) except (ValueError, TypeError): pass for idx in reversed(to_remove): task_queue.pop(idx) if to_remove: self._save_queue(task_queue) print(f"[ORCHESTRATOR] Cleaned up {len(to_remove)} completed tasks") return len(to_remove) def get_queue_summary(self) -> Dict[str, Any]: """Get a summary of the task queue.""" task_queue = self._get_queue() summary = { "total": len(task_queue), "pending": 0, "claimed": 0, "completed": 0, "failed": 0, "by_stage": {}, } for task in task_queue: status = task.get("status", "unknown") if status in summary: summary[status] += 1 stage = task.get("stage", "unknown") summary["by_stage"][stage] = summary["by_stage"].get(stage, 0) + 1 return summary # ================================================================ # Internal Methods # ================================================================ def _ensure_task_queue(self): """Ensure task_queue key exists in state.""" if self.state and hasattr(self.state, "_state"): if "task_queue" not in self.state._state: self.state._state["task_queue"] = [] self.state.save() def _get_queue(self) -> List[Dict]: """Get the task queue from state.""" if self.state and hasattr(self.state, "_state"): return self.state._state.get("task_queue", []) return [] def _save_queue(self, task_queue: List[Dict]): """Save the task queue to state.""" if self.state and hasattr(self.state, "_state"): self.state._state["task_queue"] = task_queue # Save to HF Dataset (async-friendly: just mark for save) try: self.state.save() except Exception as e: print(f"[ORCHESTRATOR] Warning: could not save state: {e}") def _release_expired_leases(self, task_queue_only: bool = False) -> List[Dict[str, Any]]: """Release tasks with expired leases. Must be called with _lock held. Returns list of released tasks. """ task_queue = self._get_queue() now = datetime.now() released = [] for i, task in enumerate(task_queue): if task["status"] != "claimed": continue lease_until = task.get("lease_until") if not lease_until: continue try: lease_dt = datetime.fromisoformat(lease_until) if lease_dt < now: task_queue[i]["status"] = "pending" task_queue[i]["lease_until"] = None task_queue[i]["claimed_by"] = None task_queue[i]["updated_at"] = now.isoformat() released.append(dict(task_queue[i])) print(f"[ORCHESTRATOR] Released expired lease for {task['video_id']}") except (ValueError, TypeError): # Invalid date format, release task_queue[i]["status"] = "pending" task_queue[i]["lease_until"] = None released.append(dict(task_queue[i])) if released: self._save_queue(task_queue) return released