"""Transport request-reply tối giản trên redis-py (bàn giao 14, 04/08/2026). Tự viết thay vì rq/arq/celery — QUYẾT ĐỊNH của Cowork: deps chỉ thêm `redis`, ~100 dòng hiểu được toàn bộ. Giao thức: API : LPUSH pc:jobs:recommend '{"job_id": ..., "payload": ...}' BLPOP pc:result:{job_id} (timeout → 503 phía route) worker : BRPOP pc:jobs:recommend → app.engine.handle_job(payload) LPUSH pc:result:{job_id} + EXPIRE (TTL dọn rác nếu API đã bỏ đi) Mode theo env `REDIS_URL` — cùng nếp `DATABASE_URL` (app/db.py): không set → "inprocess", app chạy Y HỆT bản không queue; set → "queue". KHÔNG âm thầm fallback in-process khi queue chết — che chết worker tệ hơn lỗi rõ (BRIEF 1.4). CHƯA có priority queue (realtime/TV FullVision §5.4 sẽ cần) — chừa chỗ bằng tên list theo LOẠI job (`pc:jobs:recommend`) thay vì một list chung; thêm loại job mới = thêm list mới, không đổi giao thức. Bàn giao 18 (05/08) dùng đúng chỗ chừa đó: queue thứ hai `pc:jobs:scan` cho CV worker (scripts/cv_worker.py) — các hàm nhận `jobs_key`/`key` tuỳ chọn, mặc định giữ nguyên queue recommend nên mọi caller cũ không đổi một chữ. Mọi lệnh Redis ở đây là BLOCKING (redis-py sync) — route phải gọi `submit` qua `run_in_executor`, nhất quán với cách engine đang được gọi (BRIEF 2.4). Các hàm phía worker (`serve_one`/`beat`/`clear_heartbeat`) nhận `client` tuỳ chọn để test tiêm connection fakeredis THỨ HAI — hai connection như hai process thật, còn worker thật dùng client của module sau `setup()`. """ from __future__ import annotations import json import os import time import uuid import redis JOBS_KEY = "pc:jobs:recommend" RESULT_KEY = "pc:result:{job_id}" HEARTBEAT_KEY = "pc:worker:heartbeat" # CV scan (bàn giao 18): queue + heartbeat RIÊNG — engine worker chết không # được làm health/scan tưởng CV worker chết và ngược lại. RESULT_KEY dùng # chung được vì job_id là uuid — không đụng nhau giữa hai loại job. SCAN_JOBS_KEY = "pc:jobs:scan" SCAN_HEARTBEAT_KEY = "pc:worker:heartbeat:scan" # Analyze clip (bàn giao 24): job DÀI (phút) nên KHÔNG đi đường request-reply # `submit` (timeout 30s là cho job ngắn) — API `enqueue` fire-and-forget rồi # poll status key; worker cập nhật status trong lúc chạy. Cùng CV worker với # scan nên dùng chung SCAN_HEARTBEAT_KEY, chỉ thêm list job mới (đúng chỗ # chừa "thêm loại job = thêm list" ở trên). ANALYZE_JOBS_KEY = "pc:jobs:analyze" ANALYZE_STATUS_KEY = "pc:analyze:{analyze_id}" ANALYZE_STATUS_TTL_S = 3600 # kết quả sống 1h — đủ cho Danh đọc, tự dọn rác # Segment video đa cú (lát A1, 13/08): job scan-toàn-video, cùng giao thức # fire-and-forget + status key như analyze (cùng CV worker, cùng heartbeat). # Status/tiến độ per VIDEO ở VIDEO_STATUS_KEY; kết quả BỀN (danh sách cú + # JSON per cú) nằm trên FILE ngay khi có — Redis chỉ là kênh tiến độ, vì # kết quả job từng MẤT THẬT khi chỉ nằm Redis TTL 1h (BG29b, 13/08). SEGMENT_JOBS_KEY = "pc:jobs:segment" VIDEO_STATUS_KEY = "pc:video:{video_id}" RESULT_TTL_S = 60 # kết quả mồ côi (API đã timeout bỏ đi) tự bốc hơi HEARTBEAT_TTL_S = 15 # worker SET mỗi ~5s → chết là health thấy trong ≤15s DEFAULT_TIMEOUT_S = 30.0 # redis-py 8.x đặt socket_timeout MẶC ĐỊNH 5s cho MỌI lệnh (_defaults.py — # đời 7.x là None). Hệ quả đã ĐO trên Redis thật 04/08: BLPOP/BRPOP block # server-side >= 5s là client tự ném TimeoutError đúng mốc 5.0s, trước khi # server kịp trả nil (fakeredis không có tầng socket nên test không thấy). # Chống bằng cách chờ theo LÁT < SOCKET_TIMEOUT_S qua `_bwait` — KHÔNG BAO # GIỜ gọi thẳng blpop/brpop với timeout dài. socket_timeout pin tường minh # trong setup() để đổi mặc định thư viện không âm thầm đổi hành vi ở đây. SOCKET_TIMEOUT_S = 5.0 WAIT_SLICE_S = 2.0 class QueueDown(Exception): """Không nói chuyện được với Redis (mất kết nối / chưa bật).""" class QueueTimeout(Exception): """Job đã đẩy nhưng không có kết quả trong hạn — worker chết hoặc kẹt.""" _client: redis.Redis | None = None _mode: str = "inprocess" def setup(url: str | None = None, client: redis.Redis | None = None) -> str: """Khởi tạo transport từ `url`/env `REDIS_URL` (chạy thật) hoặc `client` tiêm sẵn (test dùng fakeredis). Trả về mode. Idempotent như db.setup(). Mode thứ ba **"local"** (14/08/2026, việc B — Space HF một container): không có `REDIS_URL` mà env `POOLCOACH_LOCAL_CV` bật → transport chạy trên fakeredis TRONG TIẾN TRÌNH và app tự nuôi một luồng CV worker (app/localcv.py). Vì sao fakeredis chứ không tự viết bộ nhớ tạm: nó là ĐÚNG backend mà 400+ test đang chạy qua chính module này từ bàn giao 14 — thêm một bản hiện thực thứ hai là thêm một chỗ để lệch. Vì sao KHÔNG gọi nó là "queue": mode queue nghĩa là `/api/recommend` phải đi engine worker; ở "local" thì engine vẫn IN-PROCESS như bản không queue, chỉ tầng CV đi hàng đợi. Thiếu fakeredis → về "inprocess", app không chết. """ global _client, _mode teardown() if client is not None: _client, _mode = client, "queue" return _mode url = url or os.environ.get("REDIS_URL") if url: _client = redis.Redis.from_url(url, decode_responses=True, socket_timeout=SOCKET_TIMEOUT_S) _mode = "queue" return _mode if os.environ.get("POOLCOACH_LOCAL_CV", "").strip().lower() not in ( "", "0", "false", "no"): try: import fakeredis except ImportError: print("[poolcoach] POOLCOACH_LOCAL_CV bat nhung thieu goi " "fakeredis -- ve mode inprocess", flush=True) return _mode _client = fakeredis.FakeStrictRedis(decode_responses=True) _mode = "local" return _mode def teardown() -> None: """Về mode inprocess (đóng connection nếu có).""" global _client, _mode if _client is not None: try: _client.close() except Exception: # noqa: BLE001 — đóng lúc Redis đã chết vẫn phải xong pass _client = None _mode = "inprocess" def mode() -> str: return _mode def cv_enabled() -> bool: """Tầng CV (scan/analyze/segment) có đường chạy không — dùng cho route thay cho `mode() == "queue"`: cả "queue" (worker process) lẫn "local" (luồng nhúng) đều có worker phục vụ, chỉ "inprocess" là không.""" return _mode in ("queue", "local") def timeout_s() -> float: """Hạn chờ kết quả — đọc env MỖI lần (test override bằng monkeypatch).""" return float(os.environ.get("POOLCOACH_QUEUE_TIMEOUT_S", DEFAULT_TIMEOUT_S)) def _require(client: redis.Redis | None = None) -> redis.Redis: c = client or _client if c is None: raise RuntimeError("transport chưa setup — mode inprocess không có " "client Redis (check mode() trước)") return c def _bwait(op, key: str, timeout: float): """BLPOP/BRPOP chờ tối đa `timeout`s bằng nhiều LÁT ngắn. Mỗi lát < SOCKET_TIMEOUT_S để server luôn kịp trả nil trước khi client đứt socket (xem chú thích SOCKET_TIMEOUT_S). Trả item hoặc None khi hết hạn; lỗi Redis ném nguyên cho caller phân loại. """ deadline = time.monotonic() + timeout while True: remaining = deadline - time.monotonic() if remaining <= 0: return None # kẹp >= 0.05: BLPOP timeout 0 nghĩa là "chờ vô hạn" — không bao giờ # được truyền 0 xuống Redis (tràn deadline tối đa 50ms, vô hại) item = op(key, timeout=max(min(WAIT_SLICE_S, remaining), 0.05)) if item is not None: return item # ------------------------------------------------------------- phía API def submit(payload: dict, jobs_key: str = JOBS_KEY) -> dict: """Đẩy MỘT job và chờ kết quả (blocking — gọi trong executor). Trả reply dict của worker (`{"ok": ...}` — xem engine.handle_job). Ném QueueDown khi Redis không nói chuyện được, QueueTimeout khi quá hạn `timeout_s()` không có kết quả (worker chết/kẹt — 503 phía route, KHÔNG fallback in-process). `jobs_key` chọn LOẠI job (mặc định recommend; scan dùng SCAN_JOBS_KEY) — cùng giao thức, khác list. """ c = _require() job_id = uuid.uuid4().hex limit = timeout_s() try: c.lpush(jobs_key, json.dumps({"job_id": job_id, "payload": payload})) item = _bwait(c.blpop, RESULT_KEY.format(job_id=job_id), limit) except redis.RedisError as e: raise QueueDown(f"{type(e).__name__}: {e}") from e if item is None: raise QueueTimeout(f"không có kết quả sau {limit:g}s (job {job_id})") return json.loads(item[1]) def enqueue(payload: dict, jobs_key: str = JOBS_KEY, job_id: str | None = None) -> str: """Đẩy MỘT job KHÔNG chờ kết quả (bàn giao 24 — job dài kiểu analyze). Cùng wire format với `submit` ({"job_id", "payload"}) nên `serve_one` phía worker không đổi một chữ; khác là không BLPOP reply — tiến độ/kết quả đi đường status key (`set_analyze_status`). `job_id` truyền vào được để route dùng LUÔN nó làm id public (client poll GET bằng id này). Ném QueueDown khi Redis không nói chuyện được. """ c = _require() job_id = job_id or uuid.uuid4().hex try: c.lpush(jobs_key, json.dumps({"job_id": job_id, "payload": payload})) except redis.RedisError as e: raise QueueDown(f"{type(e).__name__}: {e}") from e return job_id def set_analyze_status(analyze_id: str, data: dict, client: redis.Redis | None = None) -> None: """Ghi status job analyze (SET + TTL — key tự bốc hơi sau 1h). `data` là dict JSON-thuần {"status": queued|running|done|error, ...} — shape do route/worker thống nhất, transport không phán. Ném QueueDown khi Redis chết — caller quyết nuốt hay không (worker nuốt để job dài không chết vì một nhịp Redis rớt; route thì 503). """ try: _require(client).set(ANALYZE_STATUS_KEY.format(analyze_id=analyze_id), json.dumps(data), ex=ANALYZE_STATUS_TTL_S) except redis.RedisError as e: raise QueueDown(f"{type(e).__name__}: {e}") from e def get_analyze_status(analyze_id: str) -> dict | None: """Đọc status job analyze — None nếu id lạ hoặc key đã hết TTL. Ném QueueDown khi Redis chết (route trả 503, không đổ oan 404).""" try: raw = _require().get(ANALYZE_STATUS_KEY.format(analyze_id=analyze_id)) except redis.RedisError as e: raise QueueDown(f"{type(e).__name__}: {e}") from e return json.loads(raw) if raw else None def set_video_status(video_id: str, data: dict, client: redis.Redis | None = None) -> None: """Status job segment per VIDEO (lát A1) — cùng giao ước set_analyze_status: SET + TTL, shape do route/worker thống nhất, ném QueueDown khi Redis chết (worker nuốt, route 503).""" try: _require(client).set(VIDEO_STATUS_KEY.format(video_id=video_id), json.dumps(data), ex=ANALYZE_STATUS_TTL_S) except redis.RedisError as e: raise QueueDown(f"{type(e).__name__}: {e}") from e def get_video_status(video_id: str) -> dict | None: """Đọc status job segment — None nếu id lạ hoặc hết TTL (danh sách cú BỀN nằm trên file, route tự đọc tiếp — Redis chỉ là kênh tiến độ).""" try: raw = _require().get(VIDEO_STATUS_KEY.format(video_id=video_id)) except redis.RedisError as e: raise QueueDown(f"{type(e).__name__}: {e}") from e return json.loads(raw) if raw else None def ping() -> bool: """Redis còn sống? — cho /api/health, nuốt lỗi (False là câu trả lời).""" try: return bool(_require().ping()) except redis.RedisError: return False def worker_alive(key: str = HEARTBEAT_KEY) -> bool: """Heartbeat worker còn hạn? — cho /api/health, nuốt lỗi như ping(). `key` chọn worker (mặc định engine; CV dùng SCAN_HEARTBEAT_KEY).""" try: return _require().exists(key) == 1 except redis.RedisError: return False def worker_alive_or_raise(key: str = HEARTBEAT_KEY) -> bool: """Heartbeat check TRƯỚC enqueue (fast-fail, bàn giao 19) — KHÁC `worker_alive`: Redis chết phải NÉM QueueDown để route trả đúng message "Redis" như đường submit (nuốt thành False sẽ đổ oan cho worker). Lệnh Redis blocking như mọi lệnh khác — route gọi qua executor.""" try: return _require().exists(key) == 1 except redis.RedisError as e: raise QueueDown(f"{type(e).__name__}: {e}") from e # ----------------------------------------------------------- phía worker def serve_one(handler, timeout_s: float = 5.0, client: redis.Redis | None = None, jobs_key: str = JOBS_KEY) -> bool: """MỘT vòng BRPOP → ``handler(payload)`` → reply + TTL. Trả False nếu hết `timeout_s` không có job — vòng ngoài của worker nhân đó beat heartbeat / bắt Ctrl+C. Job kẹt lại từ lúc worker chết vẫn được xử lý khi worker sống lại (reply mồ côi tự hết hạn theo RESULT_TTL_S) — chấp nhận tốn vài giây sim thay vì thêm deadline vào giao thức. `jobs_key` chọn LOẠI job như bên `submit`. """ c = _require(client) item = _bwait(c.brpop, jobs_key, timeout_s) if item is None: return False job = json.loads(item[1]) reply = handler(job["payload"]) key = RESULT_KEY.format(job_id=job["job_id"]) pipe = c.pipeline() pipe.lpush(key, json.dumps(reply)) pipe.expire(key, RESULT_TTL_S) pipe.execute() return True def beat(client: redis.Redis | None = None, key: str = HEARTBEAT_KEY) -> None: """Worker báo sống — CHỈ gọi sau khi JIT/model-load xong (bối cảnh #3: health không được báo alive lúc còn đang khởi động). `key` là heartbeat của TỪNG worker (engine mặc định; CV dùng SCAN_HEARTBEAT_KEY).""" _require(client).set(key, "alive", ex=HEARTBEAT_TTL_S) def clear_heartbeat(client: redis.Redis | None = None, key: str = HEARTBEAT_KEY) -> None: """Worker thoát sạch (Ctrl+C) — xoá heartbeat để health thấy chết NGAY, không đợi TTL. Redis chết lúc này thì thôi — TTL sẽ dọn.""" try: _require(client).delete(key) except redis.RedisError: pass