* docs(plan-04): spec + plan for pipeline error transparency (#131) speckit spec/plan/research/data-model/contract/quickstart for plan-04. Grounds the fix in the real code map: shared failure-event builder (backend/core/failure.py) feeding tasks.py + dub_pipeline.py + dub_core.py, non-empty reason guarantee, sanitized diagnostic block, frontend renderer with docs deeplink. Closes-target: #131 (children #122, #63). Design only — no code changes yet. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * feat(pipeline): structured, non-empty failure events + logged tracebacks (#131) plan-04 backend: no more silent "unknown error". A shared failure helper guarantees a non-empty reason at every emit site and a sanitized, copyable diagnostic block. - backend/core/failure.py: build_failure()/build_failure_event() (reason falls back to the exception class name), sanitize() (reuses the logging_filter HF-token regex + redacts *TOKEN*/*KEY*/*SECRET* env values + home→~), diagnostic() (reuses the env capture), classify() reusing the error_docs_map 5-class taxonomy for the docs deeplink + hint. - core/tasks.py worker: structured event instead of bare str(e); keeps the logged traceback. - services/dub_pipeline.py: enrich download/extract error yields; ADD the missing outer `except Exception` (the #122 path — unhandled ingest errors were never surfaced with stage context); surface the previously-silent demucs/scene/thumbnail degradations as non-fatal `warning` events. - api/routers/batch.py: guaranteed non-empty batch failure reason. SSE payload is additive (legacy `error`/`stage`/`detail` keys preserved), so existing frontends keep working and already show the specific reason. Tests (TDD, fail-before/pass-after): 14 cases — non-empty-reason guarantee, redaction, diagnostic sanitization, and the 3 Test-matrix triggers (worker / extract / url). 483 passed, 0 regressions. Closes #131. Refs #122, #63. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * feat(dub-ui): show specific cause + docs deeplink + copyable diagnostic (#131) plan-04 frontend. The backend now sends a structured, non-empty failure; surface it to the user instead of "extract: unknown error". - dubSlice: DubFailure type + dubFailure state/setter. - useDubWorkflow: capture the structured failure on the SSE error event (reason/error_class/stage/hint/docs_topic/diagnostic); clear on new runs. - DubTab: DubFailureNotice renders the actionable hint, an "Open docs" deeplink (via the existing errorDocsMap classifier), and a "Copy diagnostic" button — shown beneath the error badge in both failure banners. typecheck + build clean; 66 frontend tests pass. Refs #131, #122, #63. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * fix(failure): annotate intentional best-effort excepts (CodeQL) The new security workflow's CodeQL flagged 5 bare `except: pass` blocks. All are deliberate best-effort guards (sanitize/diagnostic must never throw on the failure path; the test cancels the worker to tear it down). Added explanatory comments per CodeQL's py/empty-except rule. No behavior change. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
148 lines
5.6 KiB
Python
148 lines
5.6 KiB
Python
import asyncio
|
|
import time
|
|
import json
|
|
import logging
|
|
|
|
from core import job_store
|
|
from core import failure
|
|
|
|
logger = logging.getLogger("omnivoice.tasks")
|
|
|
|
|
|
class TaskManager:
|
|
"""In-memory task dispatcher with SQLite-backed metadata.
|
|
|
|
The dispatcher itself (queue + worker + listeners) stays in-memory for
|
|
speed, but every state transition and every SSE event is mirrored to
|
|
`jobs` / `job_events`. That means:
|
|
|
|
- clients can reconnect via `/tasks/stream/{id}?after_seq=N` and catch up
|
|
- restart recovers: orphaned `running` jobs are flipped to `failed`
|
|
- `GET /jobs` works across restarts
|
|
"""
|
|
|
|
def __init__(self):
|
|
self.queue = None
|
|
self.active_tasks = {}
|
|
|
|
def _init_queue(self):
|
|
if self.queue is None:
|
|
self.queue = asyncio.Queue()
|
|
|
|
async def add_task(self, task_id, task_type, func, *args, project_id=None, meta=None, **kwargs):
|
|
self._init_queue()
|
|
task_obj = {
|
|
"status": "pending",
|
|
"type": task_type,
|
|
"created_at": time.time(),
|
|
"history": [],
|
|
"listeners": [],
|
|
"listeners_lock": asyncio.Lock(),
|
|
"error": None,
|
|
"cancelled": False,
|
|
}
|
|
self.active_tasks[task_id] = task_obj
|
|
try:
|
|
job_store.create(task_id, type=task_type, project_id=project_id, meta=meta)
|
|
except Exception:
|
|
logger.exception("job_store.create failed (non-fatal); in-memory task still runs")
|
|
await self.queue.put((task_id, func, args, kwargs))
|
|
|
|
def cancel_task(self, task_id):
|
|
if task_id in self.active_tasks:
|
|
self.active_tasks[task_id]["cancelled"] = True
|
|
return True
|
|
return False
|
|
|
|
def is_cancelled(self, task_id):
|
|
t = self.active_tasks.get(task_id)
|
|
return t["cancelled"] if t else False
|
|
|
|
async def add_listener(self, task_id, q):
|
|
t = self.active_tasks.get(task_id)
|
|
if not t:
|
|
return False
|
|
async with t["listeners_lock"]:
|
|
t["listeners"].append(q)
|
|
return True
|
|
|
|
async def remove_listener(self, task_id, q):
|
|
t = self.active_tasks.get(task_id)
|
|
if not t:
|
|
return
|
|
async with t["listeners_lock"]:
|
|
if q in t["listeners"]:
|
|
t["listeners"].remove(q)
|
|
|
|
async def _push_event(self, task_id, event_str):
|
|
t = self.active_tasks.get(task_id)
|
|
if t is None:
|
|
return
|
|
if event_str is not None:
|
|
t["history"].append(event_str)
|
|
try:
|
|
seq = job_store.append_event(task_id, event_str)
|
|
# Stash the seq on the in-memory copy too, mainly for tests.
|
|
t.setdefault("event_seqs", []).append(seq)
|
|
except Exception:
|
|
# Never let disk writes break the live stream.
|
|
logger.exception("job_store.append_event failed; event delivered to listeners only")
|
|
# Snapshot listeners under lock so concurrent add/remove can't mutate mid-iteration.
|
|
async with t["listeners_lock"]:
|
|
listeners = list(t["listeners"])
|
|
for q in listeners:
|
|
await q.put(event_str)
|
|
|
|
async def worker(self):
|
|
self._init_queue()
|
|
while True:
|
|
task_id, func, args, kwargs = await self.queue.get()
|
|
t = self.active_tasks.get(task_id)
|
|
if not t:
|
|
self.queue.task_done()
|
|
continue
|
|
|
|
t["status"] = "running"
|
|
try:
|
|
job_store.mark_running(task_id)
|
|
except Exception:
|
|
logger.exception("job_store.mark_running failed (non-fatal)")
|
|
try:
|
|
import inspect
|
|
res = func(*args, **kwargs)
|
|
if inspect.isasyncgen(res):
|
|
async for update in res:
|
|
if t.get("cancelled"):
|
|
await self._push_event(task_id, f"data: {json.dumps({'type': 'cancelled'})}\n\n")
|
|
t["status"] = "cancelled"
|
|
try: job_store.mark_cancelled(task_id)
|
|
except Exception: logger.exception("job_store.mark_cancelled failed")
|
|
break
|
|
await self._push_event(task_id, update)
|
|
elif inspect.iscoroutine(res):
|
|
await res
|
|
if t["status"] != "cancelled":
|
|
t["status"] = "done"
|
|
try: job_store.mark_done(task_id)
|
|
except Exception: logger.exception("job_store.mark_done failed")
|
|
except Exception as e:
|
|
logger.exception("Task %s failed", task_id)
|
|
t["status"] = "failed"
|
|
# plan-04 (#131): structured, non-empty failure event instead of
|
|
# a bare str(e) (which is empty/cryptic for many exception types).
|
|
evt = failure.build_failure_event(e, stage="task", context={"task_id": task_id})
|
|
t["error"] = evt["reason"]
|
|
try:
|
|
job_store.mark_failed(task_id, evt["reason"])
|
|
except Exception:
|
|
logger.exception("job_store.mark_failed failed")
|
|
try:
|
|
await self._push_event(task_id, f"data: {json.dumps(evt)}\n\n")
|
|
except Exception as push_err:
|
|
logger.warning("Failed to push error event for %s: %s", task_id, push_err)
|
|
finally:
|
|
await self._push_event(task_id, None) # EOF
|
|
self.queue.task_done()
|
|
|
|
task_manager = TaskManager()
|