Files
admin/host-stacks/vds-kzntsv/tg-digest/http_worker.py
vitya 9da36ddb27 docs(runbook): tg-digest VDS deploy runbook + stack artifacts (task:1311)
Ранбук первого деплоя паттерна 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.
2026-08-28 23:33:47 +03:00

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()