LLMVideoPipeline/pipeline/queue_manager.py

138 lines
5.2 KiB
Python
Raw Normal View History

2026-08-16 09:39:36 -06:00
"""
Persistent SQLite Queue & Checkpoint Manager for 24/7 Autonomous Operation.
Provides job scheduling, per-step progress tracking, crash resumption, and job status reporting.
"""
import os
import sqlite3
import json
import time
import logging
from typing import List, Dict, Any, Optional
logger = logging.getLogger("YTFactory.Queue")
class QueueManager:
def __init__(self, db_path: str = "output/factory_queue.db"):
self.db_path = db_path
os.makedirs(os.path.dirname(os.path.abspath(self.db_path)), exist_ok=True)
self._init_db()
def _get_conn(self) -> sqlite3.Connection:
conn = sqlite3.connect(self.db_path)
conn.row_factory = sqlite3.Row
return conn
def _init_db(self):
with self._get_conn() as conn:
conn.execute("""
CREATE TABLE IF NOT EXISTS video_jobs (
job_id TEXT PRIMARY KEY,
topic TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'QUEUED',
progress_step TEXT NOT NULL DEFAULT 'PENDING',
duration INTEGER,
created_at REAL NOT NULL,
started_at REAL,
completed_at REAL,
error_message TEXT,
job_dir TEXT,
output_bundle TEXT
)
""")
conn.commit()
def add_job(self, topic: str, duration: Optional[int] = None) -> str:
"""Enqueue a new topic for processing."""
job_id = f"{int(time.time())}_{abs(hash(topic)) % 100000}"
now = time.time()
with self._get_conn() as conn:
conn.execute("""
INSERT INTO video_jobs (job_id, topic, status, progress_step, duration, created_at)
VALUES (?, ?, 'QUEUED', 'PENDING', ?, ?)
""", (job_id, topic.strip(), duration, now))
conn.commit()
logger.info(f"[QUEUE] Added job '{topic[:40]}' (ID: {job_id})")
return job_id
def get_next_job(self) -> Optional[Dict[str, Any]]:
"""Retrieve and lock the next pending or retriable job."""
with self._get_conn() as conn:
row = conn.execute("""
SELECT * FROM video_jobs
WHERE status = 'QUEUED'
ORDER BY created_at ASC
LIMIT 1
""").fetchone()
if row:
job = dict(row)
conn.execute("""
UPDATE video_jobs
SET status = 'PROCESSING', started_at = ?
WHERE job_id = ?
""", (time.time(), job["job_id"]))
conn.commit()
return job
return None
def update_progress(self, job_id: str, step: str, job_dir: Optional[str] = None):
"""Update current active processing step."""
with self._get_conn() as conn:
conn.execute("""
UPDATE video_jobs
SET progress_step = ?, job_dir = COALESCE(?, job_dir)
WHERE job_id = ?
""", (step, job_dir, job_id))
conn.commit()
def mark_completed(self, job_id: str, output_bundle: Dict[str, Any]):
"""Mark job successfully completed."""
with self._get_conn() as conn:
conn.execute("""
UPDATE video_jobs
SET status = 'COMPLETED', progress_step = 'DONE', completed_at = ?, output_bundle = ?
WHERE job_id = ?
""", (time.time(), json.dumps(output_bundle, ensure_ascii=False), job_id))
conn.commit()
logger.info(f"[QUEUE] Job {job_id} marked COMPLETED")
def mark_failed(self, job_id: str, error_msg: str):
"""Mark job as failed with error details."""
with self._get_conn() as conn:
conn.execute("""
UPDATE video_jobs
SET status = 'FAILED', error_message = ?, completed_at = ?
WHERE job_id = ?
""", (error_msg, time.time(), job_id))
conn.commit()
logger.error(f"[QUEUE] Job {job_id} marked FAILED: {error_msg}")
def list_jobs(self, limit: int = 25) -> List[Dict[str, Any]]:
"""List recent jobs with their statuses."""
with self._get_conn() as conn:
rows = conn.execute("""
SELECT job_id, topic, status, progress_step, duration, created_at, completed_at, error_message
FROM video_jobs
ORDER BY created_at DESC
LIMIT ?
""", (limit,)).fetchall()
return [dict(r) for r in rows]
def retry_failed(self) -> int:
"""Reset failed jobs back to QUEUED state."""
with self._get_conn() as conn:
cur = conn.execute("""
UPDATE video_jobs
SET status = 'QUEUED', progress_step = 'PENDING', error_message = NULL
WHERE status = 'FAILED'
""")
conn.commit()
return cur.rowcount
def has_pending_jobs(self) -> bool:
"""Check if any jobs are waiting in queue."""
with self._get_conn() as conn:
count = conn.execute("SELECT COUNT(*) FROM video_jobs WHERE status = 'QUEUED'").fetchone()[0]
return count > 0