- 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
79 lines
2.6 KiB
Python
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
|