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.
197 lines
8.1 KiB
Python
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
|