Spaces:
Sleeping
Sleeping
File size: 5,338 Bytes
db4ba8d 518aad9 8b949ea 518aad9 8b949ea 518aad9 8b949ea db4ba8d | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 | """
TradeFlow AI — Celery Application
PRD §0.2 Invariant #6: Celery queues must be strictly prioritized.
Queue order: critical > high > default > low
Enterprise tier MUST use critical/high queues.
"""
try:
from celery import Celery
from kombu import Exchange, Queue
except Exception: # pragma: no cover - provide lightweight fallbacks for tests
class _DummyConf:
def __init__(self):
self.task_queues = ()
self.task_default_queue = None
self.task_default_exchange = None
self.task_default_routing_key = None
self.beat_schedule = {}
self.task_routes = {}
def update(self, *a, **k):
return None
class Celery: # minimal stand-in
def __init__(self, *a, **k):
self.conf = _DummyConf()
def task(self, *targs, **tkwargs):
def _decorator(fn):
return fn
return _decorator
class Exchange:
def __init__(self, *a, **k):
pass
class Queue:
def __init__(self, *a, **k):
pass
from ..config import settings
from celery.signals import worker_process_init
import asyncio
_worker_loop = None
@worker_process_init.connect
def init_celery_worker(**kwargs):
global _worker_loop
from ..dependencies import init_supabase
_worker_loop = asyncio.new_event_loop()
asyncio.set_event_loop(_worker_loop)
_worker_loop.run_until_complete(init_supabase())
def get_worker_loop():
global _worker_loop
if _worker_loop is None:
_worker_loop = asyncio.new_event_loop()
asyncio.set_event_loop(_worker_loop)
return _worker_loop
# ── Celery app ────────────────────────────────────────────────────────────────
celery_app = Celery(
"tradeflow",
broker=settings.REDIS_URL,
backend=settings.REDIS_URL,
include=[
"src.tasks.ocr_tasks",
"src.tasks.submit_tasks",
"src.tasks.learning_tasks",
],
)
# ── Configuration ─────────────────────────────────────────────────────────────
celery_app.conf.update(
task_serializer="json",
result_serializer="json",
accept_content=["json"],
timezone="Asia/Jakarta",
enable_utc=True,
task_soft_time_limit=settings.CELERY_TASK_SOFT_TIME_LIMIT,
task_time_limit=settings.CELERY_TASK_TIME_LIMIT,
task_acks_late=True,
worker_prefetch_multiplier=1,
task_track_started=True,
result_expires=86400, # 24h
broker_connection_retry_on_startup=True,
)
# ── Priority Queues ───────────────────────────────────────────────────────────
# INVARIANT: Enterprise tier MUST use critical/high queues
default_exchange = Exchange("tradeflow", type="direct")
celery_app.conf.task_queues = (
Queue("critical", default_exchange, routing_key="critical"), # blockchain anchoring
Queue("high", default_exchange, routing_key="high"), # OCR + extraction (enterprise)
Queue("default", default_exchange, routing_key="default"), # standard processing
Queue("low", default_exchange, routing_key="low"), # retraining, reporting
)
celery_app.conf.task_default_queue = "default"
celery_app.conf.task_default_exchange = "tradeflow"
celery_app.conf.task_default_routing_key = "default"
# Dead letter queue for exhausted retries
celery_app.conf.task_queues += (
Queue("dlq", default_exchange, routing_key="dlq"),
)
# ── Task routes ───────────────────────────────────────────────────────────────
celery_app.conf.task_routes = {
"src.tasks.ocr_tasks.preprocess_document": {"queue": "high"},
"src.tasks.ocr_tasks.run_ocr": {"queue": "high"},
"src.tasks.ocr_tasks.extract_fields": {"queue": "high"},
"src.tasks.ocr_tasks.validate_fields": {"queue": "default"},
"src.tasks.ocr_tasks.recommend_hs": {"queue": "default"},
"src.tasks.ocr_tasks.assess_risk": {"queue": "default"},
"src.tasks.submit_tasks.submit_to_ceisa": {"queue": "high"},
"src.tasks.submit_tasks.anchor_blockchain": {"queue": "critical"},
"src.tasks.submit_tasks.process_ceisa_response": {"queue": "high"},
"src.tasks.learning_tasks.retrain_predictor": {"queue": "low"},
"src.tasks.learning_tasks.refresh_btki_embeddings": {"queue": "low"},
}
# ── Celery Beat schedule (periodic tasks) ─────────────────────────────────────
celery_app.conf.beat_schedule = {
"refresh-btki-monthly": {
"task": "src.tasks.learning_tasks.refresh_btki_embeddings",
"schedule": 30 * 24 * 60 * 60, # 30 days in seconds
"options": {"queue": "low"},
},
"check-retrain-trigger-daily": {
"task": "src.tasks.learning_tasks.check_retrain_trigger",
"schedule": 24 * 60 * 60, # daily
"options": {"queue": "low"},
},
}
|