diff --git a/backend/shonar/services/processing.py b/backend/shonar/services/processing.py index 35da1af..6f77c85 100644 --- a/backend/shonar/services/processing.py +++ b/backend/shonar/services/processing.py @@ -435,6 +435,21 @@ async def run_transcribe(ctx: dict, recording_id: str) -> None: await _fail(session, rec, job, str(e), ctx, e) await session.commit() 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( session, rec, result.text, result.segments, result.language, provider.name, getattr(result, "model", ""), @@ -549,6 +564,25 @@ async def run_summarize(ctx: dict, recording_id: str) -> None: await session.commit() return 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, result.model, tone=job.tone) job.status = JobStatus.succeeded diff --git a/backend/tests/test_ai_pipeline.py b/backend/tests/test_ai_pipeline.py index 166535c..7b7a18e 100644 --- a/backend/tests/test_ai_pipeline.py +++ b/backend/tests/test_ai_pipeline.py @@ -251,6 +251,45 @@ async def test_config_error_fails_fast(client, monkeypatch): 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 -----------------------------------------------------------------