Spaces:
Sleeping
Sleeping
| """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 | |