File size: 4,742 Bytes
4ec3855
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
"""
Job Progress Tracker
Manages per-job state, progress, and stage tracking.
Thread-safe for background task usage.
"""

import time
import threading
import json
from pathlib import Path
from typing import Optional, Dict
from dataclasses import dataclass, field, asdict
from enum import Enum


class JobStatus(str, Enum):
    QUEUED = "queued"
    PROCESSING = "processing"
    COMPLETED = "completed"
    FAILED = "failed"
    CANCELLED = "cancelled"


STAGE_NAMES = [
    "downloading",
    "extracting_audio",
    "separating_vocals",
    "transcribing",
    "profiling_speakers",
    "translating",
    "generating_tts",
    "assembling_audio",
    "mixing_audio",
    "merging_video",
    "generating_subtitles",
]


@dataclass
class JobState:
    job_id: str
    status: JobStatus = JobStatus.QUEUED
    current_stage: str = ""
    stage_number: int = 0
    total_stages: int = len(STAGE_NAMES)
    progress_percent: int = 0
    elapsed_time: float = 0
    estimated_remaining: str = "calculating..."
    start_time: float = 0
    error_message: str = ""
    video_title: str = ""
    target_language: str = ""
    output_path: str = ""
    subtitle_paths: Dict = field(default_factory=dict)

    def to_dict(self):
        d = asdict(self)
        d["status"] = self.status.value
        d["elapsed_time_str"] = _format_time(self.elapsed_time)
        return d


def _format_time(seconds: float) -> str:
    if seconds < 60:
        return f"{int(seconds)}s"
    elif seconds < 3600:
        return f"{int(seconds // 60)}m {int(seconds % 60)}s"
    else:
        h = int(seconds // 3600)
        m = int((seconds % 3600) // 60)
        return f"{h}h {m}m"


class ProgressTracker:
    """Thread-safe job progress tracker."""

    def __init__(self):
        self._jobs: Dict[str, JobState] = {}
        self._lock = threading.Lock()

    def create_job(self, job_id: str, target_language: str = "") -> JobState:
        with self._lock:
            state = JobState(
                job_id=job_id,
                target_language=target_language,
                start_time=time.time()
            )
            self._jobs[job_id] = state
            return state

    def get_job(self, job_id: str) -> Optional[JobState]:
        with self._lock:
            return self._jobs.get(job_id)

    def update_stage(self, job_id: str, stage_name: str):
        with self._lock:
            job = self._jobs.get(job_id)
            if not job:
                return
            job.status = JobStatus.PROCESSING
            job.current_stage = stage_name
            if stage_name in STAGE_NAMES:
                job.stage_number = STAGE_NAMES.index(stage_name) + 1
            job.elapsed_time = time.time() - job.start_time
            # Estimate remaining
            if job.stage_number > 0:
                per_stage = job.elapsed_time / job.stage_number
                remaining_stages = job.total_stages - job.stage_number
                est = per_stage * remaining_stages
                job.estimated_remaining = f"~{_format_time(est)}"
            job.progress_percent = int(job.stage_number / job.total_stages * 100)

    def update_progress(self, job_id: str, percent: int):
        with self._lock:
            job = self._jobs.get(job_id)
            if job:
                # Blend stage progress into overall progress
                base = (job.stage_number - 1) / job.total_stages * 100
                stage_contrib = percent / job.total_stages
                job.progress_percent = min(int(base + stage_contrib), 99)
                job.elapsed_time = time.time() - job.start_time

    def complete_job(self, job_id: str, output_path: str, subtitle_paths: dict = None):
        with self._lock:
            job = self._jobs.get(job_id)
            if job:
                job.status = JobStatus.COMPLETED
                job.progress_percent = 100
                job.current_stage = "done"
                job.output_path = output_path
                job.subtitle_paths = subtitle_paths or {}
                job.elapsed_time = time.time() - job.start_time
                job.estimated_remaining = "0s"

    def fail_job(self, job_id: str, error: str):
        with self._lock:
            job = self._jobs.get(job_id)
            if job:
                job.status = JobStatus.FAILED
                job.error_message = error
                job.elapsed_time = time.time() - job.start_time
                job.estimated_remaining = ""

    def list_jobs(self) -> list:
        with self._lock:
            return [j.to_dict() for j in self._jobs.values()]

    def remove_job(self, job_id: str):
        with self._lock:
            self._jobs.pop(job_id, None)


# Global singleton
tracker = ProgressTracker()