From 212291ea94d97045255f7f2020945b6e256e0d57 Mon Sep 17 00:00:00 2001 From: David Kifer Date: Wed, 24 Jun 2026 13:37:17 -0400 Subject: [PATCH] feat: implement job resumption, custom queue reordering, and auto-venv start script setup --- app/main.py | 159 ++++++++++++++++++++++++- app/upscaler.py | 294 ++++++++++++++++++++++++++++------------------ context/README.md | 58 +++++++++ start.py | 90 ++++++++++++-- static/app.js | 107 ++++++++++++++++- static/styles.css | 1 + 6 files changed, 569 insertions(+), 140 deletions(-) create mode 100644 context/README.md diff --git a/app/main.py b/app/main.py index 40c43e1..8f2e652 100644 --- a/app/main.py +++ b/app/main.py @@ -27,21 +27,97 @@ app.add_middleware( ) # In-memory databases +class CustomJobQueue: + def __init__(self): + self.queue = [] + self.lock = threading.Lock() + self.condition = threading.Condition(self.lock) + + def put(self, job_id: str): + with self.lock: + if job_id not in self.queue: + self.queue.append(job_id) + self.condition.notify() + + def get(self) -> str: + with self.lock: + while not self.queue: + self.condition.wait() + return self.queue.pop(0) + + def remove(self, job_id: str) -> bool: + with self.lock: + if job_id in self.queue: + self.queue.remove(job_id) + return True + return False + + def get_all(self) -> List[str]: + with self.lock: + return list(self.queue) + + def reorder(self, job_ids: List[str]): + with self.lock: + valid_ids = [jid for jid in job_ids if jid in self.queue] + missing_ids = [jid for jid in self.queue if jid not in valid_ids] + self.queue = valid_ids + missing_ids + + def task_done(self): + pass + + def empty(self) -> bool: + with self.lock: + return len(self.queue) == 0 + + def qsize(self) -> int: + with self.lock: + return len(self.queue) + jobs_db: Dict[str, upscaler.UpscaleJob] = {} ws_connections: Dict[str, List[WebSocket]] = {} preview_db: Dict[str, Dict[str, str]] = {} # preview_id -> {orig, upscaled} -# FIFO queue for upscaling jobs to prevent GPU memory overload -job_queue = queue.Queue() +# Custom thread-safe queue for upscaling jobs to support reordering & cancellation +job_queue = CustomJobQueue() queue_lock = threading.Lock() current_running_job_id = None main_loop = None +JOBS_FILE = os.path.join(upscaler.BASE_DIR, "jobs.json") + +def load_jobs_db(): + global jobs_db + if os.path.exists(JOBS_FILE): + try: + with open(JOBS_FILE, "r") as f: + data = json.load(f) + for job_id, job_data in data.items(): + job = upscaler.UpscaleJob.from_dict(job_data) + # Automatically put queued items back in the queue + if job.status == "queued": + job_queue.put(job_id) + # Mark active items as interrupted so they can be resumed + elif job.status in ["analyzing", "extracting", "upscaling", "assembling"]: + job.status = "interrupted" + job.eta = "Interrupted" + jobs_db[job_id] = job + except Exception as e: + print(f"Error loading jobs database: {e}") + +def save_jobs_db(): + try: + with open(JOBS_FILE, "w") as f: + data = {job_id: job.to_dict() for job_id, job in jobs_db.items()} + json.dump(data, f, indent=4) + except Exception as e: + print(f"Error saving jobs database: {e}") + @app.on_event("startup") def startup_event(): global main_loop main_loop = asyncio.get_event_loop() + load_jobs_db() global_webhook_url = None @@ -64,6 +140,7 @@ def send_webhook_notification(url: str, payload: dict): # Broadcast updates to websockets and webhooks def broadcast_progress(job_id: str, data: dict): + save_jobs_db() job = jobs_db.get(job_id) if job: data["is_preview"] = getattr(job, "is_preview", False) @@ -344,6 +421,7 @@ def start_upscale(req: StartUpscaleRequest): jobs_db[job_id] = job job_queue.put(job_id) + save_jobs_db() # Broadcast initial queued progress broadcast_progress(job_id, { @@ -386,6 +464,9 @@ def cancel_job(job_id: str): if not job: raise HTTPException(status_code=404, detail="Job not found.") + # Remove from queue if it was queued + job_queue.remove(job_id) + job.cancel() # Broadcast cancellation status broadcast_progress(job_id, { @@ -395,6 +476,7 @@ def cancel_job(job_id: str): "total_frames": job.total_frames, "eta": "N/A" }) + save_jobs_db() return {"job_id": job_id, "status": "cancelled"} @app.post("/api/preview/generate") @@ -486,7 +568,22 @@ def safe_delete_file(file_path: str): @app.get("/api/jobs") def list_jobs(): - """List details of all submitted jobs""" + """List details of all submitted jobs in queue-sorted order""" + active_id = current_running_job_id + queued_ids = job_queue.get_all() + + # Sort active first, then queued in order, then history by start time descending + def get_sort_key(job): + if job.job_id == active_id: + return (0, 0) + elif job.job_id in queued_ids: + return (1, queued_ids.index(job.job_id)) + else: + t = job.start_time if job.start_time is not None else 0 + return (2, -t) + + sorted_jobs = sorted(jobs_db.values(), key=get_sort_key) + return [ { "job_id": job.job_id, @@ -500,9 +597,10 @@ def list_jobs(): "scale": job.scale, "output_file": os.path.basename(job.output_file) if job.output_file else None, "video_path": job.video_path, - "is_preview": getattr(job, "is_preview", False) + "is_preview": getattr(job, "is_preview", False), + "queue_position": queued_ids.index(job.job_id) if job.job_id in queued_ids else -1 if job.job_id == active_id else None } - for job in jobs_db.values() + for job in sorted_jobs ] @app.delete("/api/jobs/{job_id}") @@ -512,6 +610,9 @@ def delete_job(job_id: str): if not job: raise HTTPException(status_code=404, detail="Job not found.") + # Remove from queue if it is queued + job_queue.remove(job_id) + # Safely delete original preview video if present for ext in [".mp4", ".mkv", ".avi", ".mov", ".webm"]: orig_prev_path = os.path.join(upscaler.OUTPUT_DIR, f"original_{job_id}{ext}") @@ -534,6 +635,8 @@ def delete_job(job_id: str): # Delete from in-memory db if job_id in jobs_db: del jobs_db[job_id] + + save_jobs_db() return {"job_id": job_id, "status": "purged"} @@ -562,11 +665,57 @@ def purge_all_jobs(): # Reset in-memory database jobs_db.clear() + # Re-initialize custom queue + global job_queue + job_queue = CustomJobQueue() + # Reset upload metadata file save_upload_metadata({}) + save_jobs_db() + return {"status": "all purged"} +class ReorderQueueRequest(BaseModel): + job_ids: List[str] + +@app.post("/api/queue/reorder") +def reorder_queue(req: ReorderQueueRequest): + """Reorder the job queue""" + job_queue.reorder(req.job_ids) + save_jobs_db() + return {"status": "success", "queue": job_queue.get_all()} + +@app.get("/api/queue") +def get_queue(): + """Get the current job queue order""" + return {"queue": job_queue.get_all()} + +@app.post("/api/upscale/resume/{job_id}") +def resume_job(job_id: str): + """Resume an interrupted/failed upscale job""" + job = jobs_db.get(job_id) + if not job: + raise HTTPException(status_code=404, detail="Job not found.") + + # Re-queue the job + job.status = "queued" + job.error = None + job.eta = "Queued for resume..." + + job_queue.put(job_id) + save_jobs_db() + + broadcast_progress(job_id, { + "status": "queued", + "progress": job.progress, + "current_frame": job.current_frame, + "total_frames": job.total_frames, + "eta": "Queued for resume..." + }) + + return {"job_id": job_id, "status": "queued"} + # Websocket endpoint for real-time progress updates @app.websocket("/ws/progress/{job_id}") async def websocket_progress(websocket: WebSocket, job_id: str): diff --git a/app/upscaler.py b/app/upscaler.py index 23c8b2a..950d6d9 100644 --- a/app/upscaler.py +++ b/app/upscaler.py @@ -78,6 +78,31 @@ class UpscaleJob: self._is_cancelled = False self._lock = threading.Lock() + def to_dict(self) -> dict: + """Serialize job attributes, excluding internal thread/process resources.""" + return {k: v for k, v in self.__dict__.items() if not k.startswith('_')} + + @classmethod + def from_dict(cls, data: dict) -> 'UpscaleJob': + """Deserialize job from dictionary, reconstructing internal locks and processes.""" + job = cls( + job_id=data.get('job_id'), + video_path=data.get('video_path'), + model=data.get('model'), + scale=data.get('scale', 4), + tile_size=data.get('tile_size', 256), + preserve_audio=data.get('preserve_audio', True), + webhook_url=data.get('webhook_url'), + transcode_format=data.get('transcode_format', 'mp4'), + is_preview=data.get('is_preview', False) + ) + for k, v in data.items(): + setattr(job, k, v) + job._processes = [] + job._is_cancelled = False + job._lock = threading.Lock() + return job + def update_status(self, status: str, progress: float = None, current_frame: int = None, eta: str = None, error: str = None): with self._lock: self.status = status @@ -256,133 +281,168 @@ def run_upscale_pipeline(job: UpscaleJob, on_progress_update: Callable[[str, Dic except Exception as cut_err: print(f"Error cutting original preview video: {cut_err}") - # Step 1: Extract Frames - job.update_status("extracting", progress=10) - on_progress_update(job.job_id, {"status": "extracting", "progress": 10}) + # Step 1: Extract Frames (Support Skipping on Resume) + skip_extraction = False + if os.path.exists(input_frames_dir): + extracted_files = sorted([f for f in os.listdir(input_frames_dir) if f.startswith("frame_")]) + if len(extracted_files) > 0: + skip_extraction = True + print(f"Job {job.job_id}: Found existing input frames ({len(extracted_files)} frames). Skipping extraction step.") + job.total_frames = len(extracted_files) - # High quality JPG frames to balance disk usage and speed - extract_cmd = ["ffmpeg", "-y"] - if job.ss is not None: - extract_cmd.extend(["-ss", str(job.ss)]) - if job.t is not None: - extract_cmd.extend(["-t", str(job.t)]) - extract_cmd.extend(["-i", job.video_path]) - - # Apply unsharp pre-filter if enabled - if getattr(job, "unsharp", False): - extract_cmd.extend(["-vf", "unsharp"]) + if not skip_extraction: + job.update_status("extracting", progress=10) + on_progress_update(job.job_id, {"status": "extracting", "progress": 10}) - extract_cmd.extend([ - "-q:v", "2", - os.path.join(input_frames_dir, "frame_%08d.jpg") - ]) - - p_extract = job.run_command(extract_cmd) - stdout, stderr = p_extract.communicate() - job.cleanup_process(p_extract) - - if p_extract.returncode != 0: - raise RuntimeError(f"FFmpeg frame extraction failed: {stderr}") + # High quality JPG frames to balance disk usage and speed + extract_cmd = ["ffmpeg", "-y"] + if job.ss is not None: + extract_cmd.extend(["-ss", str(job.ss)]) + if job.t is not None: + extract_cmd.extend(["-t", str(job.t)]) + extract_cmd.extend(["-i", job.video_path]) - # Count actual frames extracted - extracted_files = sorted([f for f in os.listdir(input_frames_dir) if f.startswith("frame_")]) - actual_total = len(extracted_files) - if actual_total == 0: - raise RuntimeError("No frames extracted from video") - - job.total_frames = actual_total - - # Step 2: Upscale Frames - job.update_status("upscaling", progress=20, current_frame=0) - on_progress_update(job.job_id, {"status": "upscaling", "progress": 20, "current_frame": 0, "total_frames": actual_total}) - - current_tile_size = job.tile_size - while True: - # Launch Real-ESRGAN on directory - upscale_cmd = [ - BIN_PATH, - "-i", input_frames_dir, - "-o", output_frames_dir, - "-n", job.model, - "-s", str(job.scale), - "-t", str(current_tile_size), - "-f", "jpg" - ] - if getattr(job, "gpu_ids", None) is not None: - upscale_cmd.extend(["-g", str(job.gpu_ids)]) - if getattr(job, "tta", False): - upscale_cmd.append("-x") + # Apply unsharp pre-filter if enabled + if getattr(job, "unsharp", False): + extract_cmd.extend(["-vf", "unsharp"]) - upscale_start_time = time.time() - p_upscale = job.run_command(upscale_cmd) + extract_cmd.extend([ + "-q:v", "2", + os.path.join(input_frames_dir, "frame_%08d.jpg") + ]) - # Monitor thread for output files - while p_upscale.poll() is None: + p_extract = job.run_command(extract_cmd) + stdout, stderr = p_extract.communicate() + job.cleanup_process(p_extract) + + if p_extract.returncode != 0: + raise RuntimeError(f"FFmpeg frame extraction failed: {stderr}") + + # Count actual frames extracted + extracted_files = sorted([f for f in os.listdir(input_frames_dir) if f.startswith("frame_")]) + actual_total = len(extracted_files) + if actual_total == 0: + raise RuntimeError("No frames extracted from video") + + job.total_frames = actual_total + else: + actual_total = job.total_frames + + # Step 2: Upscale Frames (Support Resuming by Skipping already upscaled frames) + if os.path.exists(output_frames_dir): + output_files = os.listdir(output_frames_dir) + skipped_frames = 0 + for f in output_files: + if f.startswith("frame_") and f.endswith(".jpg"): + out_path = os.path.join(output_frames_dir, f) + if os.path.exists(out_path) and os.path.getsize(out_path) > 0: + in_path = os.path.join(input_frames_dir, f) + if os.path.exists(in_path): + try: + os.remove(in_path) + skipped_frames += 1 + except Exception as ex: + print(f"Error removing resumed frame {in_path}: {ex}") + if skipped_frames > 0: + print(f"Job {job.job_id}: Skipping {skipped_frames} already upscaled frames.") + + remaining_inputs = len(os.listdir(input_frames_dir)) if os.path.exists(input_frames_dir) else 0 + + if remaining_inputs == 0: + print(f"Job {job.job_id}: All frames already upscaled. Skipping upscaling step.") + job.update_status("upscaling", progress=80.0, current_frame=actual_total) + on_progress_update(job.job_id, {"status": "upscaling", "progress": 80.0, "current_frame": actual_total, "total_frames": actual_total}) + else: + job.update_status("upscaling", progress=20, current_frame=actual_total - remaining_inputs) + on_progress_update(job.job_id, {"status": "upscaling", "progress": 20, "current_frame": actual_total - remaining_inputs, "total_frames": actual_total}) + + current_tile_size = job.tile_size + while True: + # Launch Real-ESRGAN on directory + upscale_cmd = [ + BIN_PATH, + "-i", input_frames_dir, + "-o", output_frames_dir, + "-n", job.model, + "-s", str(job.scale), + "-t", str(current_tile_size), + "-f", "jpg" + ] + if getattr(job, "gpu_ids", None) is not None: + upscale_cmd.extend(["-g", str(job.gpu_ids)]) + if getattr(job, "tta", False): + upscale_cmd.append("-x") + + upscale_start_time = time.time() + p_upscale = job.run_command(upscale_cmd) + + # Monitor thread for output files + while p_upscale.poll() is None: + if job._is_cancelled: + return + + processed_files = len(os.listdir(output_frames_dir)) + progress_pct = 20.0 + (float(processed_files) / actual_total) * 60.0 # upscaling is 20% to 80% + + # Estimate ETA + elapsed = time.time() - upscale_start_time + this_run_processed = processed_files - (actual_total - remaining_inputs) + if this_run_processed > 0: + sec_per_frame = elapsed / this_run_processed + rem_frames = actual_total - processed_files + eta_sec = rem_frames * sec_per_frame + + # Format ETA + if eta_sec > 60: + eta_str = f"{int(eta_sec // 60)}m {int(eta_sec % 60)}s" + else: + eta_str = f"{int(eta_sec)}s" + else: + eta_str = "Calculating..." + + job.update_status("upscaling", progress=progress_pct, current_frame=processed_files, eta=eta_str) + on_progress_update(job.job_id, { + "status": "upscaling", + "progress": progress_pct, + "current_frame": processed_files, + "total_frames": actual_total, + "eta": eta_str + }) + time.sleep(0.5) + + stdout, stderr = p_upscale.communicate() + job.cleanup_process(p_upscale) + if job._is_cancelled: return - processed_files = len(os.listdir(output_frames_dir)) - progress_pct = 20.0 + (float(processed_files) / actual_total) * 60.0 # upscaling is 20% to 80% - - # Estimate ETA - elapsed = time.time() - upscale_start_time - if processed_files > 0: - sec_per_frame = elapsed / processed_files - rem_frames = actual_total - processed_files - eta_sec = rem_frames * sec_per_frame + if p_upscale.returncode != 0: + err_msg = (stdout or "") + "\n" + (stderr or "") + is_alloc_error = any(x in err_msg.lower() for x in ["vkallocatememory", "out of memory", "allocation", "vram", "failed to allocate"]) - # Format ETA - if eta_sec > 60: - eta_str = f"{int(eta_sec // 60)}m {int(eta_sec % 60)}s" - else: - eta_str = f"{int(eta_sec)}s" + if is_alloc_error: + if current_tile_size <= 0: + next_tile_size = 256 + else: + next_tile_size = current_tile_size // 2 + + if next_tile_size >= 32: + print(f"Job {job.job_id}: Real-ESRGAN failed with VRAM allocation error. Retrying with tile size halved from {current_tile_size} to {next_tile_size}.") + current_tile_size = next_tile_size + + # Clean up only output frames that we attempted to upscale in this run + for filename in os.listdir(input_frames_dir): + out_path = os.path.join(output_frames_dir, filename) + if os.path.exists(out_path): + try: + os.unlink(out_path) + except Exception: + pass + continue + + raise RuntimeError(f"Real-ESRGAN failed with exit code {p_upscale.returncode}: {err_msg}") else: - eta_str = "Calculating..." - - job.update_status("upscaling", progress=progress_pct, current_frame=processed_files, eta=eta_str) - on_progress_update(job.job_id, { - "status": "upscaling", - "progress": progress_pct, - "current_frame": processed_files, - "total_frames": actual_total, - "eta": eta_str - }) - time.sleep(0.5) - - stdout, stderr = p_upscale.communicate() - job.cleanup_process(p_upscale) - - if job._is_cancelled: - return - - if p_upscale.returncode != 0: - err_msg = (stdout or "") + "\n" + (stderr or "") - is_alloc_error = any(x in err_msg.lower() for x in ["vkallocatememory", "out of memory", "allocation", "vram", "failed to allocate"]) - - if is_alloc_error: - if current_tile_size <= 0: - next_tile_size = 256 - else: - next_tile_size = current_tile_size // 2 - - if next_tile_size >= 32: - print(f"Job {job.job_id}: Real-ESRGAN failed with VRAM allocation error. Retrying with tile size halved from {current_tile_size} to {next_tile_size}.") - current_tile_size = next_tile_size - # Clean up output frames directory before retrying - for filename in os.listdir(output_frames_dir): - file_path = os.path.join(output_frames_dir, filename) - try: - if os.path.isfile(file_path) or os.path.islink(file_path): - os.unlink(file_path) - elif os.path.isdir(file_path): - shutil.rmtree(file_path) - except Exception as cleanup_err: - print(f"Error cleaning file {file_path}: {cleanup_err}") - continue - - raise RuntimeError(f"Real-ESRGAN failed with exit code {p_upscale.returncode}: {err_msg}") - else: - break + break # Final validation of upscale output processed_files = len(os.listdir(output_frames_dir)) diff --git a/context/README.md b/context/README.md new file mode 100644 index 0000000..7cf3b0f --- /dev/null +++ b/context/README.md @@ -0,0 +1,58 @@ +# 🧠 AI Video Upscaler - Project Context & Reference + +Welcome! This document acts as a memory reference folder for AI agents working on this project. It details the system architecture, code organization, pipeline stages, and custom feature implementations (Resumption & Queue Management). + +--- + +## 🏗️ System Architecture + +The AI Video Upscaler is a full-stack, single-user application designed to upscale videos locally using Vulkan GPU-accelerated AI models. + +### Technology Stack +1. **FastAPI (Python Backend)**: Handles API endpoints, serving static assets, WebSockets, background threads, and orchestrates subprocess execution. +2. **Real-ESRGAN ncnn-vulkan (AI Engine)**: A compiled C++ Vulkan binary located in `realesrgan-bin/realesrgan-ncnn-vulkan`. +3. **FFmpeg & FFprobe (Multimedia Toolkit)**: Splitting videos into high-quality JPEG frames, querying metadata, extracting audio streams, and re-muxing audio/subtitles back into the upscaled output. +4. **Vanilla HTML/CSS/JS (Frontend)**: Cyberpunk-themed web dashboard with real-time progress logging, side-by-side comparative preview, and queue management. + +--- + +## 📁 Key File Map + +- **`app/main.py`**: API route definitions, WebSocket orchestrator, custom thread-safe job queue, and JSON database state preservation (`jobs.json`). +- **`app/upscaler.py`**: Core pipeline orchestrator (`run_upscale_pipeline`). Executes external FFmpeg and Real-ESRGAN commands. Contains the resume logic. +- **`static/index.html`**, **`static/styles.css`**, **`static/app.js`**: Frontend interface, WebSocket handlers, interactive comparison slider, zoom/pan tool, and queue reordering actions. +- **`jobs.json`**: Persisted file storing details of all jobs (automatically generated at root). + +--- + +## 🔄 Upscale Pipeline Lifecycle + +A typical upscale job follows these sequential steps: +1. **Queued (`queued`)**: Added to the FIFO queue. +2. **Analyzing (`analyzing`)**: GPU lock acquired. Querying stream specs with `ffprobe`. +3. **Extracting (`extracting`)**: FFmpeg extracts video frames into `temp//input_frames/frame_%08d.jpg`. +4. **Upscaling (`upscaling`)**: Real-ESRGAN binary upscales images to `temp//output_frames/frame_%08d.jpg`. +5. **Assembling (`assembling`)**: FFmpeg merges upscaled frames with original audio/subtitles. +6. **Completed (`completed`)** / **Failed (`failed`)** / **Cancelled (`cancelled`)** / **Interrupted (`interrupted`)**. + +--- + +## ⚡ Custom Enhancements + +### 1. Job Resumption (`interrupted` status) +- **Persistency**: The status of all jobs is stored in `jobs.json` at the root. On application restart, any incomplete/active job is automatically loaded in the `interrupted` status. +- **Skip Extraction**: On resume, if `temp//input_frames` contains frames, the extraction step is bypassed. +- **Incremental Upscaling**: The upscaler inspects `output_frames/` and checks for already upscaled frames. It removes corresponding files from `input_frames/`, meaning the AI model only processes the remaining un-upscaled frames. +- **Fast Assembly**: If all frames are already upscaled, it skips the upscaling step entirely and directly runs the FFmpeg reassembly. + +### 2. Queue Management (Cancel & Reorder) +- **Custom Queue (`CustomJobQueue`)**: A thread-safe, list-backed FIFO queue replacing the standard `queue.Queue`. It allows: + - Querying the active queue list (`GET /api/queue`). + - Swapping queued jobs and re-ordering (`POST /api/queue/reorder`). + - Graceful removal upon cancellation. +- **Interactive UI Arrows**: Arrows are displayed next to queued items in the dashboard to move jobs up and down, triggering the API to swap their order on-the-fly. + +### 3. Automatic Virtual Environment Setup (`start.py`) +- **Self-Sufficiency**: Running `python3 start.py` automatically checks for a local virtual environment (`venv/` or `.venv/`). +- **Auto-Provisioning**: If no virtual environment is found, `start.py` will initialize one in `venv/`, upgrade `pip`, install all dependencies listed in `requirements.txt`, mark the Real-ESRGAN binary as executable (`chmod +x`), and create necessary folders (`uploads/`, `outputs/`, `temp/`). +- **Seamless Launch**: It then automatically launches the server process using the newly created environment interpreter. diff --git a/start.py b/start.py index cc9620c..695ec8e 100755 --- a/start.py +++ b/start.py @@ -21,18 +21,83 @@ import webbrowser BASE_DIR = os.path.dirname(os.path.abspath(__file__)) PID_FILE = os.path.join(BASE_DIR, ".uvicorn.pid") LOG_FILE = os.path.join(BASE_DIR, "server.log") -# Detect the right Python interpreter: -# 1. If running inside an activated venv, use that interpreter -# 2. Otherwise look for venv/ then .venv/ in the project -if sys.prefix != sys.base_prefix: - # Already inside an activated venv — use it - VENV_PYTHON = sys.executable -else: - VENV_PYTHON = os.path.join(BASE_DIR, "venv", "bin", "python") - if not os.path.isfile(VENV_PYTHON): - VENV_PYTHON = os.path.join(BASE_DIR, ".venv", "bin", "python") - if not os.path.isfile(VENV_PYTHON): - VENV_PYTHON = sys.executable # last resort: system python + +VENV_PYTHON = sys.executable # default fallback + +def ensure_venv(): + """Detect or create the virtual environment and install dependencies.""" + global VENV_PYTHON + + # 1. Already inside an activated venv — use it + if sys.prefix != sys.base_prefix: + VENV_PYTHON = sys.executable + return + + venv_dir = os.path.join(BASE_DIR, "venv") + dot_venv_dir = os.path.join(BASE_DIR, ".venv") + + # Check if either venv/ or .venv/ exists with a python interpreter + selected_venv = None + if os.path.isdir(os.path.join(venv_dir, "bin")): + selected_venv = venv_dir + elif os.path.isdir(os.path.join(dot_venv_dir, "bin")): + selected_venv = dot_venv_dir + + if selected_venv: + python_exe = os.path.join(selected_venv, "bin", "python") + if os.path.isfile(python_exe): + VENV_PYTHON = python_exe + return + + # No valid venv found — let's build it! + print("⚙️ Virtual environment not detected. Initializing setup...") + print(f" Creating virtual environment at: {venv_dir}") + try: + import venv + venv.create(venv_dir, with_pip=True) + print("✅ Virtual environment created.") + except Exception as e: + print(f"❌ Failed to create virtual environment via 'venv' module: {e}") + print(" Attempting subprocess fallback...") + try: + subprocess.run([sys.executable, "-m", "venv", venv_dir], check=True) + print("✅ Virtual environment created (fallback).") + except Exception as err: + print(f"❌ Subprocess fallback failed: {err}") + sys.exit(1) + + VENV_PYTHON = os.path.join(venv_dir, "bin", "python") + pip_exe = os.path.join(venv_dir, "bin", "pip") + + # Install dependencies + requirements_file = os.path.join(BASE_DIR, "requirements.txt") + if os.path.isfile(requirements_file): + print("📦 Installing Python dependencies from requirements.txt...") + try: + # Upgrade pip + subprocess.run([pip_exe, "install", "--upgrade", "pip", "--quiet"], check=True) + # Install requirements + subprocess.run([pip_exe, "install", "-r", requirements_file], check=True) + print("✅ All Python dependencies installed successfully.") + except Exception as e: + print(f"❌ Error installing dependencies: {e}") + sys.exit(1) + else: + print("⚠ requirements.txt not found. Skipping dependency installation.") + + # Configure Real-ESRGAN binary permissions + binary_path = os.path.join(BASE_DIR, "realesrgan-bin", "realesrgan-ncnn-vulkan") + if os.path.isfile(binary_path): + try: + os.chmod(binary_path, 0o755) + print("✅ Real-ESRGAN binary marked as executable.") + except Exception as e: + print(f"⚠ Failed to change binary permissions: {e}") + + # Create working directories + for d in ["uploads", "outputs", "temp"]: + os.makedirs(os.path.join(BASE_DIR, d), exist_ok=True) + print("✅ Working directories verified.") def _is_running(pid: int) -> bool: @@ -153,4 +218,5 @@ if __name__ == "__main__": help="Don't open the browser automatically") args = parser.parse_args() + ensure_venv() start(args.host, args.port, open_browser=not args.no_browser) diff --git a/static/app.js b/static/app.js index 3726b48..21d1e2a 100644 --- a/static/app.js +++ b/static/app.js @@ -1276,6 +1276,37 @@ async function loadGallery() { } } +async function moveQueueItem(jobId, direction) { + try { + const res = await fetch("/api/queue"); + if (!res.ok) return; + const data = await res.json(); + const queue = data.queue; + const index = queue.indexOf(jobId); + if (index === -1) return; + + const newIndex = index + direction; + if (newIndex < 0 || newIndex >= queue.length) return; + + // Swap + const temp = queue[index]; + queue[index] = queue[newIndex]; + queue[newIndex] = temp; + + const reorderRes = await fetch("/api/queue/reorder", { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ job_ids: queue }) + }); + + if (reorderRes.ok) { + loadQueue(); + } + } catch (err) { + console.error("Error reordering queue:", err); + } +} + async function loadQueue() { const queueList = document.getElementById("queue-list"); if (!queueList) return; @@ -1287,7 +1318,7 @@ async function loadQueue() { const jobs = await res.json(); const displayJobs = jobs.filter(job => - ["queued", "analyzing", "extracting", "upscaling", "assembling", "completed", "failed", "cancelled"].includes(job.status) + ["queued", "analyzing", "extracting", "upscaling", "assembling", "completed", "failed", "cancelled", "interrupted"].includes(job.status) ); if (displayJobs.length === 0) { @@ -1318,13 +1349,40 @@ async function loadQueue() { if (job.status === "failed" && job.error) { metaHtml += ` Error: ${job.error}`; - } else if (job.status !== "queued" && job.status !== "failed" && job.status !== "cancelled" && job.status !== "completed") { + } else if (job.status !== "queued" && job.status !== "failed" && job.status !== "cancelled" && job.status !== "completed" && job.status !== "interrupted") { metaHtml += ` ETA: ${job.eta || 'Calculating...'}`; } let actionsHtml = ""; - if (["queued", "analyzing", "extracting", "upscaling", "assembling"].includes(job.status)) { + let reorderHtml = ""; + + if (job.status === "queued" && job.queue_position !== null && job.queue_position !== undefined) { + const queuedJobs = displayJobs.filter(j => j.status === "queued"); + const isFirst = job.queue_position === 0; + const isLast = job.queue_position === queuedJobs.length - 1; + + reorderHtml = ` + + + `; + } + + if (job.status === "interrupted") { actionsHtml = ` + + + `; + } else if (["queued", "analyzing", "extracting", "upscaling", "assembling"].includes(job.status)) { + actionsHtml = ` + ${reorderHtml} @@ -1366,9 +1424,44 @@ async function loadQueue() { `; + const upBtn = item.querySelector(".btn-move-up"); + if (upBtn) { + upBtn.addEventListener("click", async (e) => { + e.stopPropagation(); + await moveQueueItem(job.job_id, -1); + }); + } + + const downBtn = item.querySelector(".btn-move-down"); + if (downBtn) { + downBtn.addEventListener("click", async (e) => { + e.stopPropagation(); + await moveQueueItem(job.job_id, 1); + }); + } + + const resumeBtn = item.querySelector(".btn-resume-job"); + if (resumeBtn) { + resumeBtn.addEventListener("click", async (e) => { + e.stopPropagation(); + try { + const resumeRes = await fetch(`/api/upscale/resume/${job.job_id}`, { method: "POST" }); + if (resumeRes.ok) { + loadQueue(); + } else { + const err = await resumeRes.json(); + alert("Failed to resume job: " + formatFetchError(err, "unknown error")); + } + } catch (err) { + console.error("Error resuming job:", err); + } + }); + } + const abortBtn = item.querySelector(".btn-abort-job"); if (abortBtn) { - abortBtn.addEventListener("click", async () => { + abortBtn.addEventListener("click", async (e) => { + e.stopPropagation(); if (confirm("Are you sure you want to abort this upscaling job?")) { try { const abortRes = await fetch(`/api/upscale/cancel/${job.job_id}`, { method: "POST" }); @@ -1384,7 +1477,8 @@ async function loadQueue() { const deleteBtn = item.querySelector(".btn-delete-job"); if (deleteBtn) { - deleteBtn.addEventListener("click", async () => { + deleteBtn.addEventListener("click", async (e) => { + e.stopPropagation(); if (confirm("Are you sure you want to clear this job from history?")) { try { const delRes = await fetch(`/api/jobs/${job.job_id}`, { method: "DELETE" }); @@ -1401,7 +1495,8 @@ async function loadQueue() { const viewBtn = item.querySelector(".btn-view-job"); if (viewBtn) { - viewBtn.addEventListener("click", () => { + viewBtn.addEventListener("click", (e) => { + e.stopPropagation(); currentJobId = job.job_id; if (socket) { try { diff --git a/static/styles.css b/static/styles.css index 0ba568f..a014897 100644 --- a/static/styles.css +++ b/static/styles.css @@ -656,6 +656,7 @@ input:checked + .slider:before { .status-badge.assembling { background: rgba(0, 240, 255, 0.1); border: 1px solid var(--primary); color: var(--primary); } .status-badge.completed { background: rgba(0, 255, 135, 0.1); border: 1px solid var(--success); color: var(--success); } .status-badge.failed { background: rgba(255, 0, 85, 0.1); border: 1px solid var(--danger); color: var(--danger); } +.status-badge.interrupted { background: rgba(255, 255, 255, 0.08); border: 1px solid var(--text-muted); color: var(--text-muted); } .bar-container { width: 100%;