S.H.O.N.A.R._Desktop_Companion/backend/shonar/worker.py
avi 76c867fca4 Standalone Shonar Desktop: vendor portable sources + local engine; decouple from ~/Projects/Shonar
- 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
2026-09-14 17:14:54 -05:00

79 lines
2.6 KiB
Python

"""arq worker entrypoint (M7): runs the AI pipeline tasks.
Run: arq shonar.worker.WorkerSettings
"""
from __future__ import annotations
import logging
from arq import cron
from arq.connections import RedisSettings
from shonar.core.config import get_settings
from shonar.db.session import dispose_engine, get_engine
from shonar.services import processing
logger = logging.getLogger("shonar.worker")
async def run_transcribe(ctx: dict, recording_id: str) -> None:
await processing.run_transcribe(ctx, recording_id)
async def run_summarize(ctx: dict, recording_id: str) -> None:
await processing.run_summarize(ctx, recording_id)
async def sweep(ctx: dict) -> None: # noqa: ARG001 — arq cron signature
count = await processing.sweep_stale()
if count:
logger.info("sweep re-enqueued %d stale jobs", count)
async def retention_sweep(ctx: dict) -> None: # noqa: ARG001 — arq cron signature
from shonar.services import retention
purged = await retention.sweep_deleted()
if purged["recordings"] or purged["users"]:
logger.info("retention sweep: %s", purged)
async def startup(ctx: dict) -> None:
get_engine()
# Crash recovery before accepting new work: jobs stuck `running` and
# queued rows the transport missed go back through arq.
count = await processing.sweep_stale()
if count:
logger.info("startup sweep re-enqueued %d stale jobs", count)
async def shutdown(ctx: dict) -> None: # noqa: ARG001
await dispose_engine()
def _redis() -> RedisSettings:
return RedisSettings.from_dsn(get_settings().redis_url)
class WorkerSettings:
functions = [run_transcribe, run_summarize, sweep, retention_sweep]
# Transport-loss backstop beyond the startup sweep: anything still
# queued (missed enqueue, dead worker between runs) goes back through
# arq every 5 minutes. Rows are the queue; this just pokes.
cron_jobs = [
cron(sweep, minute={0, 5, 10, 15, 20, 25, 30, 35, 40, 45, 50, 55}),
# Hard-delete expired soft-deletes once a day (3:17 local, off-peak).
cron(retention_sweep, hour=3, minute=17),
]
on_startup = startup
on_shutdown = shutdown
redis_settings = _redis()
# Retry budget for transient provider failures; the tasks themselves
# mark jobs failed on the last try (see processing.MAX_TRIES).
max_tries = 3
# One job per worker process at a time. Local faster-whisper is
# multi-GB per concurrent pass; default max_jobs=10 let one worker run
# ~8 transcriptions at once and OOM-swap the box. Scale by running more
# worker processes, never by raising this.
max_jobs = 1