Ранбук первого деплоя паттерна sched+worker (модель для yt-digest): стек tg-digest (воркер internal, mem_limit 256m), runtime-регистрация задачи POST /tasks, env-контракт, smoke/verify/rollback, gotchas (сессия Telethon, TZ UTC, timeoutMs). Артефакты: compose source-of-truth, Dockerfile + http_worker.py (черновики для коммита в victor/tg-digest). task:1311 разблокирован (1305 done) -> ready.
64 lines
2.6 KiB
Python
64 lines
2.6 KiB
Python
#!/usr/bin/env python3
|
|
"""tg-digest HTTP-воркер для sched (http runner, simple mode).
|
|
|
|
GET /healthz -> 200 {"ok": true} (healthcheck стека)
|
|
POST /run -> проверка x-sched-api-key (WORKER_API_KEY); запуск
|
|
`python -m src.worker`; 200 {ok, runId, summary} | 5xx {error}.
|
|
|
|
Simple mode: sched шлёт POST и ждёт ответ до config.timeoutMs; любой 2xx = succeeded.
|
|
runId приходит в заголовке x-sched-run-id (пишется в ответ, в лог).
|
|
Заголовки x-sched-api-key от sched — auth per-task (task config.auth.apiKey).
|
|
"""
|
|
import json
|
|
import os
|
|
import subprocess
|
|
import sys
|
|
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
|
|
|
WORKER_API_KEY = os.environ.get("WORKER_API_KEY", "")
|
|
WORKER_CMD = [sys.executable, "-m", "src.worker"]
|
|
WORKER_TIMEOUT = float(os.environ.get("WORKER_TIMEOUT_S") or 3600)
|
|
PORT = int(os.environ.get("PORT") or 8080)
|
|
|
|
|
|
class Handler(BaseHTTPRequestHandler):
|
|
def log_message(self, *args): # тишина в stdout (логи — по runId)
|
|
pass
|
|
|
|
def _send(self, code, obj):
|
|
body = json.dumps(obj, ensure_ascii=False).encode("utf-8")
|
|
self.send_response(code)
|
|
self.send_header("Content-Type", "application/json; charset=utf-8")
|
|
self.send_header("Content-Length", str(len(body)))
|
|
self.end_headers()
|
|
self.wfile.write(body)
|
|
|
|
def do_GET(self):
|
|
if self.path == "/healthz":
|
|
self._send(200, {"ok": True})
|
|
else:
|
|
self._send(404, {"error": "not found"})
|
|
|
|
def do_POST(self):
|
|
if self.path != "/run":
|
|
return self._send(404, {"error": "not found"})
|
|
if WORKER_API_KEY and self.headers.get("x-sched-api-key") != WORKER_API_KEY:
|
|
return self._send(401, {"error": "unauthorized"})
|
|
run_id = self.headers.get("x-sched-run-id", "?")
|
|
try:
|
|
r = subprocess.run(WORKER_CMD, capture_output=True, text=True,
|
|
env=os.environ, timeout=WORKER_TIMEOUT)
|
|
except subprocess.TimeoutExpired:
|
|
return self._send(504, {"error": "worker timeout", "runId": run_id})
|
|
if r.returncode != 0:
|
|
return self._send(500, {"error": (r.stderr or r.stdout)[-500:], "runId": run_id})
|
|
try:
|
|
summary = json.loads(r.stdout)
|
|
except ValueError:
|
|
summary = {"stdout": r.stdout[-500:]}
|
|
self._send(200, {"ok": True, "runId": run_id, "summary": summary})
|
|
|
|
|
|
if __name__ == "__main__":
|
|
ThreadingHTTPServer(("0.0.0.0", PORT), Handler).serve_forever()
|