Skip to content

Commit d3daa1e

Browse files
authored
refactor(api): route every async dispatch through one off-loop helper (#162)
ADR-0007 bounded how long an unreachable Redis can block a dispatch and moved one call site off the event loop, leaving the other 33 bounded but still inline. This finishes that follow-up and makes the property enforceable rather than remembered. Add backend/core/task_dispatch.py with dispatch_task and dispatch_task_best_effort, and route every dispatch reachable from a request handler through it. That covers the 33 remaining sites plus two helpers that were synchronous but only ever called from async code. enqueue_model_preparation is one, and it blocked twice, since it also wrote download progress to Redis inline. _enqueue_push_channel_refresh is the other, and its own docstring said it existed to keep API paths fast, which the blocking publish undermined. The helper uses asyncio.to_thread rather than Starlette's run_in_threadpool, for two reasons beyond taste. calendar_service is imported by the worker, whose image ships no ASGI stack, so a Starlette import there would not resolve. And the loop's default executor is separate from the anyio limiter that serves sync route handlers, so a Redis outage can no longer consume the threads those handlers need. That closes the threadpool residual ADR-0007 recorded rather than merely shrinking it. Best-effort dispatch is now a named function rather than a bare except at each site, so the choice to swallow is visible where it is made instead of being inferred from a try block. Worker-side code still dispatches inline, deliberately, having no event loop to protect. A test walks the tree and fails on any send_task, apply_async or delay inside an async def, naming file and line. Verified by reintroducing one inline dispatch, which it caught at documents.py:123. Align the tests on one seam while here. Four files reached for a module's celery_app re-export, which is the same import-time binding mistake that made the suite slow, and they broke as soon as the import moved. They now patch the shared app object. Two more stubbed enqueue_model_preparation with sync lambdas that are now awaited. Behaviour is unchanged and re-measured against an unreachable broker: 6.03s to fail a dispatch, with 60 concurrent no-op requests at a 0.9ms median and a 0.00s worst case, matching the run before the mechanism swap. 1151 tests pass. Refs: docs/adr/0007-bounded-fail-fast-task-dispatch.md, docs/ARCHITECTURE.md
1 parent f625d6d commit d3daa1e

31 files changed

Lines changed: 258 additions & 94 deletions

backend/api/services/health_service.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@
1414
from backend.celery_app import celery_app
1515
from backend.core.db import sync_engine
1616
from backend.core.redis import get_redis_url
17+
from backend.core.task_dispatch import dispatch_task
1718
from backend.preload_models import check_model_status
1819
from backend.utils.config_manager import async_get_system_api_keys, config_manager
1920
from backend.utils.deployment_warnings import get_deployment_warnings
@@ -528,7 +529,7 @@ async def _get_device_component(worker_status: str) -> tuple[dict[str, Any], boo
528529
)
529530

530531
try:
531-
task = celery_app.send_task("backend.worker.tasks.get_worker_device_status")
532+
task = await dispatch_task("backend.worker.tasks.get_worker_device_status")
532533
payload = await asyncio.to_thread(task.get, timeout=5)
533534
except Exception: # noqa: BLE001
534535
return (

backend/api/v1/endpoints/backup.py

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
RESTORE_LOCK_TTL_SECONDS,
3030
)
3131
from backend.core.redis import REDIS_URL
32+
from backend.core.task_dispatch import dispatch_task
3233
from backend.models.user import User
3334
from backend.utils.path_manager import PathManager
3435
from backend.utils.rate_limit import enforce_upload_concurrency
@@ -92,7 +93,7 @@ async def _dispatch_restore(
9293
)
9394

9495
try:
95-
task = celery_app.send_task(
96+
task = await dispatch_task(
9697
"backend.worker.tasks.restore_backup_task",
9798
kwargs={
9899
"zip_path": str(archive_path),
@@ -138,7 +139,7 @@ async def export_backup(
138139
try:
139140
# Trigger Celery task
140141
# Uses send_task to avoid importing the task function directly (bypasses heavy imports).
141-
task = celery_app.send_task(
142+
task = await dispatch_task(
142143
"backend.worker.tasks.create_backup_task",
143144
kwargs={
144145
"include_audio": include_audio,

backend/api/v1/endpoints/cli_oauth.py

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@
2323
from sqlalchemy.ext.asyncio import AsyncSession
2424

2525
from backend.api.deps import get_current_admin_user, get_current_user, get_db
26-
from backend.celery_app import celery_app
26+
from backend.core.task_dispatch import dispatch_task
2727
from backend.models.cli_oauth import (
2828
CliOAuthCredential,
2929
CliOAuthCredentialStatus,
@@ -363,7 +363,7 @@ async def get_codex_models(
363363
models=[CliCodexModel(**model) for model in cached], source="live"
364364
)
365365
# Warm the cache for next time (runs in worker-io); serve the fallback now.
366-
celery_app.send_task("backend.worker.tasks.refresh_codex_models_task")
366+
await dispatch_task("backend.worker.tasks.refresh_codex_models_task")
367367
return CliCodexModelsRead(models=_CODEX_FALLBACK_MODELS, source="fallback")
368368

369369

@@ -400,7 +400,7 @@ async def refresh_codex_models(
400400
stale list it was pressed to replace.
401401
"""
402402
await codex_oauth.clear_model_catalog()
403-
celery_app.send_task("backend.worker.tasks.refresh_codex_models_task")
403+
await dispatch_task("backend.worker.tasks.refresh_codex_models_task")
404404
return CliCodexModelsRead(models=_CODEX_FALLBACK_MODELS, source="fallback")
405405

406406

@@ -476,7 +476,7 @@ async def start_cli_oauth(
476476
# verification URL + code. Never block the request — a slow login would
477477
# otherwise trip a proxy timeout (Cloudflare 520).
478478
await codex_oauth.clear_login_state(current_user.id)
479-
celery_app.send_task(
479+
await dispatch_task(
480480
"backend.worker.tasks.codex_device_login_task",
481481
args=[current_user.id],
482482
)

backend/api/v1/endpoints/documents.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@
99

1010
from backend.api.deps import get_current_user, get_db
1111
from backend.api.error_handling import sanitized_http_exception
12-
from backend.celery_app import celery_app
12+
from backend.core.task_dispatch import dispatch_task
1313
from backend.models.document import Document, DocumentStatus
1414
from backend.models.recording import Recording
1515
from backend.models.recording_public import DocumentPublicRead, serialize_document
@@ -119,7 +119,7 @@ async def upload_document(
119119
await db.commit()
120120
await db.refresh(document)
121121

122-
task = celery_app.send_task(
122+
task = await dispatch_task(
123123
"backend.worker.tasks.process_document_task", args=[document.id]
124124
)
125125
from backend.models.task import register_task_ownership

backend/api/v1/endpoints/notes_templates.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,7 @@
1515
from sqlmodel import select
1616

1717
from backend.api.deps import get_current_user, get_db
18-
from backend.celery_app import celery_app
18+
from backend.core.task_dispatch import dispatch_task
1919
from backend.models.notes_template import (
2020
NotesTemplate,
2121
NotesTemplateCreate,
@@ -414,7 +414,7 @@ async def generate_notes_structure(
414414

415415
job_id = new_job_id()
416416
await publish_job_async(job_id, {"status": STATUS_PENDING})
417-
task = celery_app.send_task(
417+
task = await dispatch_task(
418418
"backend.worker.tasks.generate_notes_structure_task",
419419
args=[job_id, current_user.id, brief],
420420
)

backend/api/v1/endpoints/recordings/helpers.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@
1212
from sqlmodel import select
1313

1414
import backend.api.v1.endpoints.recordings as recordings_module
15+
from backend.core.task_dispatch import dispatch_task
1516
from backend.models.calendar import CalendarConnection, CalendarEvent, CalendarSource
1617
from backend.models.chat import ChatMessage
1718
from backend.models.context_chunk import ContextChunk
@@ -700,7 +701,7 @@ async def _requeue_for_processing(
700701
await db.commit()
701702
await db.refresh(recording)
702703

703-
task = recordings_module.celery_app.send_task(
704+
task = await dispatch_task(
704705
"backend.worker.tasks.process_recording_task",
705706
args=[recording.id, True, engine_override],
706707
)

backend/api/v1/endpoints/recordings/routes_actions.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66

77
import backend.api.v1.endpoints.recordings as recordings_module
88
from backend.api.deps import get_current_user, get_db
9+
from backend.core.task_dispatch import dispatch_task
910
from backend.models.calendar import CalendarEvent
1011
from backend.models.recording import RecordingStatus, RecordingUpdate
1112
from backend.models.recording_public import RecordingPublicRead, serialize_recording
@@ -298,7 +299,7 @@ async def infer_speakers_for_recording(
298299
await db.commit()
299300
await db.refresh(recording)
300301

301-
task = recordings_module.celery_app.send_task(
302+
task = await dispatch_task(
302303
"backend.worker.tasks.infer_speakers_task", args=[recording.id]
303304
)
304305
recording.celery_task_id = task.id

backend/api/v1/endpoints/recordings/routes_capture.py

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@
1010
import backend.api.v1.endpoints.recordings as recordings_module
1111
from backend.api.deps import get_current_recording_client_user, get_db
1212
from backend.api.error_handling import sanitized_http_exception
13+
from backend.core.task_dispatch import dispatch_task
1314
from backend.models.pipeline import RecordingAudioChunk, RecordingAudioWindowManifest
1415
from backend.models.recording import (
1516
CaptureSourceReportCreate,
@@ -221,7 +222,7 @@ async def upload_segment(
221222
"enable_live_transcription"
222223
):
223224
try:
224-
recordings_module.celery_app.send_task(
225+
await dispatch_task(
225226
"backend.processing.live_transcribe.transcribe_segment_live_task",
226227
args=[recording.id, sequence],
227228
)
@@ -234,7 +235,7 @@ async def upload_segment(
234235
)
235236
elif segment_suffix in BROWSER_AUDIO_SEGMENT_SUFFIXES:
236237
try:
237-
recordings_module.celery_app.send_task(
238+
await dispatch_task(
238239
"backend.processing.segment_transcode.transcode_segment_task",
239240
args=[recording.id, sequence],
240241
)
@@ -466,7 +467,7 @@ async def finalize_upload(
466467
await db.commit()
467468
await db.refresh(recording)
468469

469-
task = recordings_module.celery_app.send_task(
470+
task = await dispatch_task(
470471
"backend.worker.tasks.process_recording_task", args=[recording.id]
471472
)
472473
recording.celery_task_id = task.id
@@ -476,7 +477,7 @@ async def finalize_upload(
476477

477478
await register_task_ownership(db, task.id, recording.user_id)
478479

479-
proxy_task = recordings_module.celery_app.send_task(
480+
proxy_task = await dispatch_task(
480481
"backend.worker.tasks.generate_proxy_task", args=[recording.id]
481482
)
482483
if proxy_task:

backend/api/v1/endpoints/recordings/routes_import_upload.py

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@
1414
import backend.api.v1.endpoints.recordings as recordings_module
1515
from backend.api.deps import get_current_user, get_db
1616
from backend.api.error_handling import sanitized_http_exception
17+
from backend.core.task_dispatch import dispatch_task
1718
from backend.models.pipeline import RecordingAudioChunk, RecordingAudioWindowManifest
1819
from backend.models.recording import ClientStatus, Recording, RecordingStatus
1920
from backend.models.recording_public import RecordingPublicRead, serialize_recording
@@ -172,7 +173,7 @@ async def import_audio(
172173
await db.commit()
173174

174175
# Trigger processing task
175-
task = recordings_module.celery_app.send_task(
176+
task = await dispatch_task(
176177
"backend.worker.tasks.process_recording_task", args=[recording.id]
177178
)
178179
recording.celery_task_id = task.id
@@ -184,7 +185,7 @@ async def import_audio(
184185

185186
# Trigger proxy generation task
186187
if not recording.proxy_path:
187-
proxy_task = recordings_module.celery_app.send_task(
188+
proxy_task = await dispatch_task(
188189
"backend.worker.tasks.generate_proxy_task", args=[recording.id]
189190
)
190191
if proxy_task:
@@ -438,7 +439,7 @@ async def finalize_chunked_import(
438439
await db.commit()
439440
await db.refresh(recording)
440441

441-
task = recordings_module.celery_app.send_task(
442+
task = await dispatch_task(
442443
"backend.worker.tasks.process_recording_task", args=[recording.id]
443444
)
444445
recording.celery_task_id = task.id
@@ -449,7 +450,7 @@ async def finalize_chunked_import(
449450
await register_task_ownership(db, task.id, recording.user_id)
450451

451452
if not recording.proxy_path:
452-
proxy_task = recordings_module.celery_app.send_task(
453+
proxy_task = await dispatch_task(
453454
"backend.worker.tasks.generate_proxy_task", args=[recording.id]
454455
)
455456
if proxy_task:
@@ -538,7 +539,7 @@ async def upload_recording(
538539
await db.commit()
539540
await db.refresh(recording)
540541

541-
task = recordings_module.celery_app.send_task(
542+
task = await dispatch_task(
542543
"backend.worker.tasks.process_recording_task", args=[recording.id]
543544
)
544545
recording.celery_task_id = task.id
@@ -549,7 +550,7 @@ async def upload_recording(
549550
await register_task_ownership(db, task.id, recording.user_id)
550551

551552
if not recording.proxy_path:
552-
proxy_task = recordings_module.celery_app.send_task(
553+
proxy_task = await dispatch_task(
553554
"backend.worker.tasks.generate_proxy_task", args=[recording.id]
554555
)
555556
if proxy_task:

backend/api/v1/endpoints/speakers/routes_global.py

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@
1313

1414
import backend.api.v1.endpoints.speakers as speakers_module
1515
from backend.api.deps import get_current_user, get_db
16+
from backend.core.task_dispatch import dispatch_task
1617
from backend.models.people_tag_schemas import PeopleTagRead
1718
from backend.models.recording import (
1819
Recording,
@@ -545,7 +546,7 @@ async def recalibrate_voiceprint(
545546
)
546547
continue
547548

548-
task = speakers_module.celery_app.send_task(
549+
task = await dispatch_task(
549550
"backend.worker.tasks.extract_embedding_task",
550551
args=[target_audio, segs, device_str, hf_token],
551552
)
@@ -729,7 +730,7 @@ async def split_speaker(
729730
seg_tuples = [(s.start, s.end) for s in segments]
730731
target_audio = select_recording_audio_for_embedding(rec)
731732

732-
task = speakers_module.celery_app.send_task(
733+
task = await dispatch_task(
733734
"backend.worker.tasks.extract_embedding_task",
734735
args=[target_audio, seg_tuples, device_str, hf_token],
735736
)
@@ -801,7 +802,7 @@ async def split_speaker(
801802

802803
if remaining_seg_tuples:
803804
target_audio = select_recording_audio_for_embedding(rec)
804-
task = speakers_module.celery_app.send_task(
805+
task = await dispatch_task(
805806
"backend.worker.tasks.extract_embedding_task",
806807
args=[target_audio, remaining_seg_tuples, device_str, hf_token],
807808
)

0 commit comments

Comments
 (0)