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"},
    },
}