-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathworker_daemon.py
More file actions
159 lines (132 loc) · 5.48 KB
/
Copy pathworker_daemon.py
File metadata and controls
159 lines (132 loc) · 5.48 KB
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
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
"""
worker_daemon.py
================
Standalone worker process — no HTTP server, just task processing.
In a microservices deployment you run multiple replicas of this daemon
alongside a single API server. This separation means you can scale
workers independently of the API layer (a common FAANG interview design point).
Architecture
------------
[Client] ──HTTP──► [API container] ──ZADD──► [Redis]
│
ZPOPMIN
│
▼
[Worker container ×N]
(this file, replicated)
Run locally
-----------
python worker_daemon.py
In Docker
---------
docker-compose up worker
"""
from __future__ import annotations
import logging
import os
import random
import signal
import sys
import time
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] worker — %(message)s",
datefmt="%H:%M:%S",
)
logger = logging.getLogger("worker_daemon")
# ── Job registry (same as api/main.py) ───────────────────────────────────
def job_payment(order_id: str, amount: float) -> dict:
time.sleep(random.uniform(0.05, 0.2))
if random.random() < 0.15:
raise ConnectionError(f"Payment gateway timeout for order {order_id}")
return {"order_id": order_id, "amount": amount, "charged": True}
def job_email(recipient: str, subject: str, body: str = "") -> dict:
time.sleep(random.uniform(0.01, 0.06))
if random.random() < 0.05:
raise RuntimeError(f"SMTP relay error for {recipient}")
return {"to": recipient, "subject": subject, "delivered": True}
def job_report(report_id: str, rows: int = 1000) -> dict:
time.sleep(random.uniform(0.1, 0.4))
return {"report_id": report_id, "rows": rows, "status": "generated"}
def job_sync(resource: str) -> dict:
time.sleep(random.uniform(0.05, 0.15))
return {"resource": resource, "synced": True}
JOB_REGISTRY = {
"payment": job_payment,
"email": job_email,
"report": job_report,
"sync": job_sync,
}
def main() -> None:
redis_url = os.getenv("REDIS_URL", "redis://localhost:6379/0")
num_workers = int(os.getenv("WORKERS", "4"))
logger.info("═" * 50)
logger.info(" Task Queue Worker Daemon starting")
logger.info(f" Redis : {redis_url}")
logger.info(f" Workers: {num_workers}")
logger.info("═" * 50)
from core.queue import DeadLetterQueue
from core.worker import WorkerPool
from core.scheduler import Scheduler, ExponentialBackoffWithJitter
from core.dispatcher import Dispatcher
from monitoring.metrics import MetricsCollector
from plugins.middleware import RateLimiter, StructuredLogger, build_pipeline
try:
from plugins.redis_backend import RedisQueue, RedisStorage, RedisDLQ
queue = RedisQueue(redis_url=redis_url, job_registry=JOB_REGISTRY)
storage = RedisStorage(redis_url=redis_url)
dlq = RedisDLQ(redis_url=redis_url)
logger.info("✓ Connected to Redis")
except Exception as e:
logger.warning(f"Redis unavailable ({e}), using in-memory fallback")
from core.queue import PriorityTaskQueue
from plugins.storage import InMemoryStorage
queue = PriorityTaskQueue(maxsize=1000)
storage = InMemoryStorage()
dlq = DeadLetterQueue()
mc = MetricsCollector()
pool = WorkerPool(
num_workers = num_workers,
on_success = lambda t: (storage.save(t), mc.on_success(t)),
on_failure = lambda t: (storage.save(t), mc.on_failure(t)),
on_retry = lambda t: (storage.save(t), mc.on_retry(t)),
)
scheduler = Scheduler(
enqueue_fn = queue.enqueue,
strategy = ExponentialBackoffWithJitter(base=0.5, cap=30.0),
dlq_fn = dlq.enqueue,
)
dispatcher = Dispatcher(
queue = queue,
pool = pool,
scheduler = scheduler,
middleware = build_pipeline(RateLimiter(200.0), StructuredLogger()),
dlq = dlq,
poll_interval = 0.05,
max_inflight = num_workers * 2,
)
# ── Graceful shutdown on SIGTERM / SIGINT ─────────────────
def shutdown(signum, frame):
logger.info(f"Signal {signum} received — shutting down gracefully …")
dispatcher.stop(timeout=10.0)
pool.shutdown(wait=True)
report = mc.report()
logger.info(f"Final stats: {report['totals']}")
sys.exit(0)
signal.signal(signal.SIGTERM, shutdown)
signal.signal(signal.SIGINT, shutdown)
# ── Run forever ───────────────────────────────────────────
with dispatcher:
logger.info("✓ Dispatcher running — waiting for tasks …")
while True:
time.sleep(10)
report = mc.report()
logger.info(
f"[heartbeat] done={report['totals']['success']} "
f"failed={report['totals']['failure']} "
f"retried={report['totals']['retried']} "
f"queue_depth={queue.size()} "
f"tps={report['throughput']['tasks_per_sec']}"
)
if __name__ == "__main__":
main()