Raw provider crash fails the job instead of escaping into the sweep requeue loop

A non-AIError from transcribe/summarize (av.InvalidDataError on corrupt
audio, seen live with 'My recording 63') escaped the consumer, left the
row 'running', and got requeued at every engine restart forever. Both
runners now catch it, fail the row with the exception type in the
message, and keep a usable transcript (summarize completes with a
'Summary failed' note). Regression tests pin both paths (proven red
without the fix).
This commit is contained in:
avi 2026-09-19 15:24:50 -05:00
commit 514d75fe92
2 changed files with 73 additions and 0 deletions

View file

@ -435,6 +435,21 @@ async def run_transcribe(ctx: dict, recording_id: str) -> None:
await _fail(session, rec, job, str(e), ctx, e) await _fail(session, rec, job, str(e), ctx, e)
await session.commit() await session.commit()
return return
except asyncio.CancelledError:
raise
except Exception as e: # noqa: BLE001 — a raw provider crash (e.g. av.InvalidDataError
# on corrupt audio) must FAIL the job, not escape: an escaping
# exception leaves the row 'running' and the startup sweep
# requeues it forever (crash-loop seen live Sep 19).
logger.exception("transcribe provider raised for recording %s", rec.id)
job.status = JobStatus.failed
job.error = f"Transcription crashed: {type(e).__name__}: {e}"[:500]
job.stage = None
job.finished_at = utcnow()
rec.processing_status = ProcessingStatus.failed
rec.processing_error = job.error
await session.commit()
return
await store_transcript( await store_transcript(
session, rec, result.text, result.segments, result.language, session, rec, result.text, result.segments, result.language,
provider.name, getattr(result, "model", ""), provider.name, getattr(result, "model", ""),
@ -549,6 +564,25 @@ async def run_summarize(ctx: dict, recording_id: str) -> None:
await session.commit() await session.commit()
return return
result, result_provider, fallback_note = rescue result, result_provider, fallback_note = rescue
except asyncio.CancelledError:
raise
except Exception as e: # noqa: BLE001 — same reason as transcribe: a raw
# provider crash must fail the row, not escape into the sweep loop.
logger.exception("summarize provider raised for recording %s", rec.id)
job.status = JobStatus.failed
job.error = f"Summarize crashed: {type(e).__name__}: {e}"[:500]
job.stage = None
job.finished_at = utcnow()
# A usable transcript survives: complete with the failure note
# (mirrors _fail's transcript-preserving rule).
if await latest_transcript_text(session, rec.id):
rec.processing_status = ProcessingStatus.completed
rec.processing_error = f"Summary failed: {job.error}"
else:
rec.processing_status = ProcessingStatus.failed
rec.processing_error = job.error
await session.commit()
return
await store_summary(session, rec, result.to_dict(), result_provider, await store_summary(session, rec, result.to_dict(), result_provider,
result.model, tone=job.tone) result.model, tone=job.tone)
job.status = JobStatus.succeeded job.status = JobStatus.succeeded

View file

@ -251,6 +251,45 @@ async def test_config_error_fails_fast(client, monkeypatch):
assert stt.calls == 1 assert stt.calls == 1
async def test_raw_provider_crash_fails_job_no_escape(client, monkeypatch):
"""A non-AIError from the provider (e.g. av.InvalidDataError on corrupt
audio) must FAIL the row, not escape: an escaping exception leaves the
job 'running' and the startup sweep requeues it forever."""
stt = FakeTranscriber(fail=ValueError("Invalid data found when processing input"))
use_fakes(monkeypatch, stt, FakeLlm())
token = await user_tokens(client)
rec = await upload_recording(client, token, client_id="m7-crash-1")
await processing.run_transcribe({"job_try": 1}, rec["id"]) # must NOT raise
jobs = await jobs_for(rec["id"])
assert jobs[JobType.transcribe].status == JobStatus.failed
assert "ValueError" in (jobs[JobType.transcribe].error or "")
h = {"Authorization": f"Bearer {token}"}
r = await client.get(f"/api/v1/recordings/{rec['id']}", headers=h)
assert r.json()["processing_status"] == "failed"
async def test_raw_summarize_crash_keeps_transcript(client, monkeypatch):
use_fakes(monkeypatch, FakeTranscriber("fine text"), FakeLlm())
token = await user_tokens(client)
rec = await upload_recording(client, token, client_id="m7-crash-2")
await processing.run_transcribe({}, rec["id"])
class ExplodingLlm(FakeLlm):
async def summarize(self, transcript, **kw):
raise RuntimeError("llm segfaulted")
use_fakes(monkeypatch, FakeTranscriber("fine text"), ExplodingLlm())
await processing.run_summarize({"job_try": 1}, rec["id"]) # must NOT raise
jobs = await jobs_for(rec["id"])
assert jobs[JobType.summarize].status == JobStatus.failed
h = {"Authorization": f"Bearer {token}"}
r = await client.get(f"/api/v1/recordings/{rec['id']}", headers=h)
body = r.json()
# Transcript survived: recording completes with a summary-failure note.
assert body["processing_status"] == "completed"
assert "Summary failed" in body["processing_error"]
# --- versioning ----------------------------------------------------------------- # --- versioning -----------------------------------------------------------------