- shared/ = portable Android-origin sources vendored from deferred/desktop-server (app/build.gradle.kts srcDir repointed; PlaybackController.kt excluded as Android-only) - backend/ = bundled-lite engine (SQLite + inline queue); .venv symlinked from the old checkout, PYTHONPATH pins THIS backend's code over any editable install - repoRoot() resolves this project dir (env SHONAR_REPO still wins); desktop-dev.sh watches shared/ + backend/ - Verified: :app:compileKotlin + :app:test green (23 tests); engine boots on :8010, self-migrates, /healthz ok
141 lines
5.2 KiB
Python
141 lines
5.2 KiB
Python
"""In-process job runner (desktop bundled-lite engine).
|
|
|
|
``queue_backend=inline`` replaces the arq/Redis transport with a single
|
|
asyncio consumer inside the uvicorn process: DB rows stay the source of
|
|
truth (ProcessingJob), this just runs the work. One job at a time —
|
|
local faster-whisper is multi-GB per pass, mirroring the worker's
|
|
``max_jobs=1`` rule.
|
|
|
|
Transient failures retry in-process with backoff up to MAX_TRIES (the
|
|
same budget arq gives via ``max_tries``); the final failure is recorded
|
|
by the task body itself (see processing._fail).
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import contextlib
|
|
import logging
|
|
|
|
from shonar.db.models import JobType
|
|
from shonar.services import processing
|
|
|
|
logger = logging.getLogger("shonar.inline_queue")
|
|
|
|
RETRY_DELAY_SECONDS = 5.0
|
|
# Desktop engine: hard-delete expired soft-deletes once the app has been up
|
|
# for a day, then daily. Delayed first run keeps startup snappy.
|
|
RETENTION_INTERVAL_SECONDS = 24 * 3600.0
|
|
|
|
_queue: asyncio.Queue[tuple[str, str, int]] | None = None
|
|
_consumer: asyncio.Task | None = None
|
|
_retention: asyncio.Task | None = None
|
|
|
|
|
|
def _get_queue() -> asyncio.Queue[tuple[str, str, int]]:
|
|
global _queue
|
|
if _queue is None:
|
|
_queue = asyncio.Queue()
|
|
return _queue
|
|
|
|
|
|
async def start() -> None:
|
|
"""Start the consumer and re-run anything the DB says is pending."""
|
|
global _consumer, _retention
|
|
_get_queue()
|
|
# Inline engine = single process: any row still marked `running` at
|
|
# startup is a corpse from the previous process (the worker died with
|
|
# it). sweep_stale's 2h live-worker grace — correct for multi-worker
|
|
# arq deployments — would starve these jobs, so reclaim them first.
|
|
from sqlalchemy import update
|
|
|
|
from shonar.db.models import JobStatus, ProcessingJob
|
|
from shonar.db.session import session_factory
|
|
|
|
async with session_factory()() as s:
|
|
await s.execute(
|
|
update(ProcessingJob)
|
|
.where(ProcessingJob.status == JobStatus.running)
|
|
.values(status=JobStatus.queued, started_at=None)
|
|
)
|
|
await s.commit()
|
|
if _consumer is None or _consumer.done():
|
|
_consumer = asyncio.create_task(_consume(), name="shonar-inline-queue")
|
|
if _retention is None or _retention.done():
|
|
_retention = asyncio.create_task(_retention_loop(), name="shonar-retention")
|
|
# Crash recovery: queued rows (and orphaned running rows requeued by
|
|
# the sweep) go back on the in-process queue.
|
|
count = await processing.sweep_stale()
|
|
if count:
|
|
logger.info("inline queue startup sweep requeued %d jobs", count)
|
|
|
|
|
|
async def _retention_loop() -> None:
|
|
"""Daily hard-delete sweep for the desktop engine (no arq cron here).
|
|
|
|
First pass a few minutes after start (the app may only run for hours
|
|
at a time, so a full-day initial sleep could starve the sweep), then
|
|
once a day while running.
|
|
"""
|
|
from shonar.services import retention
|
|
|
|
await asyncio.sleep(120.0)
|
|
while True:
|
|
try:
|
|
purged = await retention.sweep_deleted()
|
|
if purged["recordings"] or purged["users"]:
|
|
logger.info("inline retention sweep: %s", purged)
|
|
except asyncio.CancelledError:
|
|
raise
|
|
except Exception: # noqa: BLE001 — the loop must survive any failure
|
|
logger.exception("inline retention sweep failed")
|
|
await asyncio.sleep(RETENTION_INTERVAL_SECONDS)
|
|
|
|
|
|
async def stop() -> None:
|
|
global _consumer, _retention
|
|
if _consumer is not None:
|
|
_consumer.cancel()
|
|
with contextlib.suppress(BaseException): # noqa: BLE001 — shutdown is best-effort
|
|
await _consumer
|
|
_consumer = None
|
|
if _retention is not None:
|
|
_retention.cancel()
|
|
with contextlib.suppress(BaseException): # noqa: BLE001
|
|
await _retention
|
|
_retention = None
|
|
|
|
|
|
async def enqueue(job_type: JobType, recording_id: str) -> None:
|
|
await _get_queue().put((job_type.value, str(recording_id), 1))
|
|
|
|
|
|
async def _consume() -> None:
|
|
q = _get_queue()
|
|
while True:
|
|
job_value, recording_id, attempt = await q.get()
|
|
ctx = {"job_try": attempt}
|
|
try:
|
|
if job_value == JobType.transcribe.value:
|
|
await processing.run_transcribe(ctx, recording_id)
|
|
else:
|
|
await processing.run_summarize(ctx, recording_id)
|
|
except processing.ProviderTransientError as e:
|
|
if attempt < processing.MAX_TRIES:
|
|
logger.warning(
|
|
"inline job %s:%s transient failure (%s); retry %d/%d",
|
|
job_value, recording_id, e, attempt + 1, processing.MAX_TRIES,
|
|
)
|
|
await asyncio.sleep(RETRY_DELAY_SECONDS)
|
|
await q.put((job_value, recording_id, attempt + 1))
|
|
else:
|
|
logger.error(
|
|
"inline job %s:%s failed after %d tries: %s",
|
|
job_value, recording_id, attempt, e,
|
|
)
|
|
except asyncio.CancelledError:
|
|
raise
|
|
except Exception: # noqa: BLE001 — consumer must survive any task crash
|
|
logger.exception("inline job %s:%s crashed", job_value, recording_id)
|
|
finally:
|
|
q.task_done()
|