S.H.O.N.A.R._Desktop_Companion/backend/tests/test_inline_queue.py
avi 11bb13cb3d Fix re-summarize running one voice behind: commit the job row before poking the inline queue.
The reprocess endpoint poked transport_enqueue while its request session
was still uncommitted; the inline worker woke instantly, read the job row
with its own session, and summarized with the PREVIOUS click's tone
(TRACE: click Funny -> worker read 'dry, witty'; click Neutral -> ran
'funny'). The UI disabled the voice chips for the whole wrong-voice run,
which read as 'spinner does nothing until I leave and come back'.

Rows are the real queue: commit before the poke. Regression test pins
worker tone to the requested tone and Neutral clearing a prior voice;
verified red without the fix, green with it.
2026-09-18 17:03:37 -05:00

197 lines
8.1 KiB
Python

"""Inline queue tests (desktop bundled-lite engine): with
``queue_backend=inline`` an upload runs transcribe→summarize automatically
in-process — no arq, no Redis. Providers are fakes (see test_ai_pipeline).
"""
from __future__ import annotations
import asyncio
import uuid
from sqlalchemy import select
from shonar.core.config import get_settings
from shonar.db import session as db_session
from shonar.db.models import JobStatus, JobType, ProcessingJob
from shonar.services import inline_queue, processing
from .test_ai_pipeline import FakeLlm, FakeTranscriber, upload_recording, use_fakes, user_tokens
def use_inline(monkeypatch):
"""Force queue_backend=inline for everything resolved via processing."""
inline_settings = get_settings().model_copy(update={"queue_backend": "inline"})
monkeypatch.setattr(processing, "get_settings", lambda: inline_settings)
async def _wait_terminal(recording_id: str, want: int = 2, timeout: float = 20.0):
"""Poll the DB until `want` jobs reach a terminal state."""
deadline = asyncio.get_running_loop().time() + timeout
rows: list[ProcessingJob] = []
while asyncio.get_running_loop().time() < deadline:
async with db_session._session_factory() as s:
rows = list(
(
await s.scalars(
select(ProcessingJob).where(
ProcessingJob.recording_id == uuid.UUID(recording_id)
)
)
).all()
)
if sum(
1
for j in rows
if j.status in (JobStatus.succeeded, JobStatus.failed, JobStatus.skipped)
) >= want:
return {j.job_type: j.status for j in rows}
await asyncio.sleep(0.1)
raise AssertionError(
f"jobs did not reach terminal state in {timeout}s: "
f"{[(j.job_type, j.status) for j in rows]}"
)
async def test_inline_backend_processes_upload_without_arq(client, monkeypatch):
use_fakes(monkeypatch, FakeTranscriber(), FakeLlm())
use_inline(monkeypatch)
await inline_queue.start()
try:
token = await user_tokens(client, email="inline@shonar.dev")
rec = await upload_recording(client, token, client_id="inline-1")
statuses = await _wait_terminal(rec["id"], want=2)
assert statuses[JobType.transcribe] == JobStatus.succeeded
assert statuses[JobType.summarize] == JobStatus.succeeded
h = {"Authorization": f"Bearer {token}"}
r = await client.get(f"/api/v1/recordings/{rec['id']}", headers=h)
assert r.json()["processing_status"] == "completed"
t = await client.get(f"/api/v1/recordings/{rec['id']}/transcript", headers=h)
assert t.json()["text"] == "hello world from the meeting"
s = await client.get(f"/api/v1/recordings/{rec['id']}/summary", headers=h)
assert s.json()["content"]["short"] == "Standup happened."
finally:
await inline_queue.stop()
async def test_inline_retries_transient_failure(client, monkeypatch):
stt = FakeTranscriber()
state = {"tries": 0}
original = stt.transcribe
async def flaky(audio, mime, *, language_hint=None):
state["tries"] += 1
if state["tries"] == 1:
raise processing.ProviderTransientError("temporary hiccup")
return await original(audio, mime, language_hint=language_hint)
stt.transcribe = flaky
use_fakes(monkeypatch, stt, FakeLlm())
use_inline(monkeypatch)
monkeypatch.setattr(inline_queue, "RETRY_DELAY_SECONDS", 0.05)
await inline_queue.start()
try:
token = await user_tokens(client, email="inline2@shonar.dev")
rec = await upload_recording(client, token, client_id="inline-2")
statuses = await _wait_terminal(rec["id"], want=2)
assert statuses[JobType.transcribe] == JobStatus.succeeded
assert state["tries"] >= 2 # the retry actually happened
finally:
await inline_queue.stop()
async def test_reprocess_reruns_in_place_with_model(client, monkeypatch):
"""POST /reprocess?job=transcribe&model= re-runs the pipeline on the
SAME recording (no re-upload) and persists the model override."""
stt = FakeTranscriber()
use_fakes(monkeypatch, stt, FakeLlm())
use_inline(monkeypatch)
await inline_queue.start()
try:
token = await user_tokens(client, email="reproc@shonar.dev")
h = {"Authorization": f"Bearer {token}"}
rec = await upload_recording(client, token, client_id="reproc-1")
await _wait_terminal(rec["id"], want=2)
calls_before = stt.calls
r = await client.post(
f"/api/v1/recordings/{rec['id']}/reprocess?job=transcribe&model=large-v3",
headers=h,
)
assert r.status_code == 200, r.text
statuses = await _wait_terminal(rec["id"], want=2)
assert statuses[JobType.transcribe] == JobStatus.succeeded
assert stt.calls == calls_before + 1 # re-ran, no new recording
# Same recording id, model override persisted on the row.
r = await client.get(f"/api/v1/recordings/{rec['id']}", headers=h)
assert r.json()["id"] == rec["id"]
assert r.json()["transcription_model"] == "large-v3"
finally:
await inline_queue.stop()
async def test_reprocess_summarize_uses_requested_tone(client, monkeypatch):
"""Regression: the inline worker used to read the job row BEFORE the
reprocess request committed, so every re-summarize ran with the
PREVIOUS click's tone (click Funny -> 'dry, witty' summary; click
Neutral -> 'funny' summary). The endpoint now commits before poking
the queue; this pins the worker's tone to the tone in the request."""
llm = FakeLlm()
use_fakes(monkeypatch, FakeTranscriber(), llm)
use_inline(monkeypatch)
await inline_queue.start()
try:
token = await user_tokens(client, email="reproc-tone@shonar.dev")
h = {"Authorization": f"Bearer {token}"}
rec = await upload_recording(client, token, client_id="reproc-tone-1")
await _wait_terminal(rec["id"], want=2)
# Re-summarize asking for a NEW voice.
r = await client.post(
f"/api/v1/recordings/{rec['id']}/reprocess"
"?job=summarize&tone=sarcastic",
headers=h,
)
assert r.status_code == 200, r.text
await _wait_terminal(rec["id"], want=2)
# The LLM call for this job must have carried 'sarcastic', not
# whatever the first (auto) pass used.
assert llm.seen_tones[-1] == "sarcastic", llm.seen_tones
s = await client.get(f"/api/v1/recordings/{rec['id']}/summary", headers=h)
assert s.json()["tone"] == "sarcastic"
# And Neutral (tone=None) must CLEAR the voice, not inherit it.
r = await client.post(
f"/api/v1/recordings/{rec['id']}/reprocess?job=summarize",
headers=h,
)
assert r.status_code == 200, r.text
await _wait_terminal(rec["id"], want=2)
assert llm.seen_tones[-1] is None, llm.seen_tones
s = await client.get(f"/api/v1/recordings/{rec['id']}/summary", headers=h)
assert s.json()["tone"] is None
finally:
await inline_queue.stop()
async def test_reprocess_validation(client, monkeypatch):
token = await user_tokens(client, email="reproc2@shonar.dev")
h = {"Authorization": f"Bearer {token}"}
rec = await upload_recording(client, token, client_id="reproc-2")
# Nothing transcribed yet -> summarize refuses.
r = await client.post(
f"/api/v1/recordings/{rec['id']}/reprocess?job=summarize", headers=h)
assert r.status_code == 409
# Bad model name -> 422 (never silently substituted).
r = await client.post(
f"/api/v1/recordings/{rec['id']}/reprocess?job=transcribe&model=not-a-model",
headers=h)
assert r.status_code == 422
# Not the owner -> 404.
other = await user_tokens(client, email="reproc3@shonar.dev")
r = await client.post(
f"/api/v1/recordings/{rec['id']}/reprocess?job=transcribe",
headers={"Authorization": f"Bearer {other}"})
assert r.status_code == 404