"""REST API(stdlib http.server 实现,零第三方依赖)。 接口覆盖 api-design.md 中与本模块相关的部分: - 健康检查:GET /healthz、GET /readyz - 规则管理:GET/POST /api/v1/rules、GET/PUT/DELETE /api/v1/rules/{id}、 POST /api/v1/rules/{id}/enable|disable - 事件查询:GET /api/v1/events - 告警查询与管理:GET /api/v1/alerts、GET /api/v1/alerts/{id}、 POST /api/v1/alerts/{id}/ack、POST /api/v1/alerts/{id}/close、 POST /api/v1/alerts/batch-ack - 内部样本接入:POST /api/v1/ingest(供无 Kafka 环境端到端联调) """ from __future__ import annotations import json import re import threading import time import uuid from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from typing import Any, Callable, Dict, List, Optional, Tuple from urllib.parse import urlparse from .app import DetectorApp from .models import MetricRule, MetricSample # 路由条目:(method, path_regex, handler) Route = Tuple[str, str, Callable] class APIServer: def __init__(self, app: DetectorApp, addr: str = "0.0.0.0:8080") -> None: self.app = app host, port = self._split_addr(addr) self.server = ThreadingHTTPServer((host, port), self._handler_factory()) self.routes: List[Route] = [ ("GET", r"^/healthz$", self.handle_healthz), ("GET", r"^/readyz$", self.handle_readyz), ("GET", r"^/api/v1/rules$", self.handle_list_rules), ("POST", r"^/api/v1/rules$", self.handle_create_rule), ("GET", r"^/api/v1/rules/([^/]+)$", self.handle_get_rule), ("PUT", r"^/api/v1/rules/([^/]+)$", self.handle_update_rule), ("DELETE", r"^/api/v1/rules/([^/]+)$", self.handle_delete_rule), ("POST", r"^/api/v1/rules/([^/]+)/enable$", self.handle_enable_rule), ("POST", r"^/api/v1/rules/([^/]+)/disable$", self.handle_disable_rule), ("GET", r"^/api/v1/events$", self.handle_list_events), ("GET", r"^/api/v1/alerts$", self.handle_list_alerts), ("POST", r"^/api/v1/alerts/batch-ack$", self.handle_batch_ack), ("GET", r"^/api/v1/alerts/([^/]+)$", self.handle_get_alert), ("POST", r"^/api/v1/alerts/([^/]+)/ack$", self.handle_ack_alert), ("POST", r"^/api/v1/alerts/([^/]+)/close$", self.handle_close_alert), ("POST", r"^/api/v1/ingest$", self.handle_ingest), ] @staticmethod def _split_addr(addr: str) -> Tuple[str, int]: if ":" in addr: host, port = addr.rsplit(":", 1) return host, int(port) return addr, 8080 def _handler_factory(self) -> type: api = self class Handler(BaseHTTPRequestHandler): server_version = "ThresholdEventDetector/0.1" def do_GET(self): # noqa: N802 self._dispatch("GET") def do_POST(self): # noqa: N802 self._dispatch("POST") def do_PUT(self): # noqa: N802 self._dispatch("PUT") def do_DELETE(self): # noqa: N802 self._dispatch("DELETE") def _dispatch(self, method: str) -> None: parsed = urlparse(self.path) for route_method, pattern, handler in api.routes: if route_method != method: continue m = re.match(pattern, parsed.path) if not m: continue try: handler(self, *m.groups()) except _APIError as exc: self._write_json(exc.status, {"code": exc.code, "message": exc.message, "data": None}) except Exception as exc: # noqa: BLE001 self._write_json(500, {"code": 50000, "message": f"internal error: {exc}", "data": None}) return self._write_json(404, {"code": 40400, "message": "not found", "data": None}) def _read_json(self) -> Dict[str, Any]: length = int(self.headers.get("Content-Length") or 0) if length <= 0: return {} raw = self.rfile.read(length) return json.loads(raw.decode("utf-8")) def _query(self) -> Dict[str, str]: parsed = urlparse(self.path) result: Dict[str, str] = {} if parsed.query: for pair in parsed.query.split("&"): if "=" in pair: k, v = pair.split("=", 1) result[k] = v else: result[pair] = "" return result def _write_json(self, status: int, body: Dict[str, Any]) -> None: data = json.dumps(body, ensure_ascii=False).encode("utf-8") self.send_response(status) self.send_header("Content-Type", "application/json; charset=utf-8") self.send_header("Content-Length", str(len(data))) self.end_headers() self.wfile.write(data) def log_message(self, format: str, *args) -> None: # noqa: A002 # 静默默认请求日志,便于测试输出整洁;生产可重写为 logging。 pass return Handler def serve_forever(self) -> None: self.server.serve_forever() def shutdown(self) -> None: self.server.shutdown() # ------------------------------------------------------------------ # # 健康检查 # ------------------------------------------------------------------ # def handle_healthz(self, handler) -> None: handler._write_json(200, {"code": 0, "message": "ok", "data": {"status": "ok"}}) def handle_readyz(self, handler) -> None: handler._write_json(200, {"code": 0, "message": "ok", "data": {"status": "ok", "deps": {"rule_store": "up"}}}) # ------------------------------------------------------------------ # # 规则管理 # ------------------------------------------------------------------ # def handle_list_rules(self, handler) -> None: q = handler._query() rules = self.app.rule_loader.get_rules() if q.get("metric"): rules = [r for r in rules if r.metric == q["metric"]] if q.get("enabled") in ("true", "false"): enabled = q["enabled"] == "true" rules = [r for r in rules if r.enabled == enabled] if q.get("severity"): rules = [r for r in rules if r.severity == q["severity"]] page = int(q.get("page", "1")) page_size = min(int(q.get("page_size", "20")), 200) total = len(rules) start = (page - 1) * page_size items = [r.to_dict() for r in rules[start : start + page_size]] handler._write_json(200, {"code": 0, "message": "ok", "data": {"total": total, "items": items}}) def handle_get_rule(self, handler, rule_id: str) -> None: rule = self.app.rule_loader.get_rule(rule_id) if rule is None: raise _APIError(404, 40400, f"rule {rule_id} not found") handler._write_json(200, {"code": 0, "message": "ok", "data": rule.to_dict()}) def handle_create_rule(self, handler) -> None: data = handler._read_json() # api-design:POST /rules 允许不传 rule_id,服务端自动生成。 if not data.get("rule_id"): data["rule_id"] = f"r-{uuid.uuid4().hex[:12]}" rule = MetricRule.from_dict(data) created = self.app.rule_store.create_rule(rule) self.app.rule_loader.load() handler._write_json(200, {"code": 0, "message": "ok", "data": created.to_dict()}) def handle_update_rule(self, handler, rule_id: str) -> None: data = handler._read_json() data["rule_id"] = rule_id updated = self.app.rule_store.update_rule(MetricRule.from_dict(data)) self.app.rule_loader.load() handler._write_json(200, {"code": 0, "message": "ok", "data": updated.to_dict()}) def handle_delete_rule(self, handler, rule_id: str) -> None: deleted = self.app.rule_store.delete_rule(rule_id) self.app.rule_loader.load() handler._write_json(200, {"code": 0, "message": "ok", "data": {"deleted": deleted}}) def handle_enable_rule(self, handler, rule_id: str) -> None: rule = self._set_rule_enabled(rule_id, True) handler._write_json(200, {"code": 0, "message": "ok", "data": {"enabled": rule.enabled}}) def handle_disable_rule(self, handler, rule_id: str) -> None: rule = self._set_rule_enabled(rule_id, False) handler._write_json(200, {"code": 0, "message": "ok", "data": {"enabled": rule.enabled}}) def _set_rule_enabled(self, rule_id: str, enabled: bool) -> MetricRule: rule = self.app.rule_loader.get_rule(rule_id) if rule is None: raise _APIError(404, 40400, f"rule {rule_id} not found") rule.enabled = enabled self.app.rule_store.update_rule(rule) self.app.rule_loader.load() return rule # ------------------------------------------------------------------ # # 事件 / 告警查询与管理 # ------------------------------------------------------------------ # def handle_list_events(self, handler) -> None: q = handler._query() total, items = self.app.event_store.list_events( host_id=q.get("host_id"), rule_id=q.get("rule_id"), status=q.get("status"), severity=q.get("severity"), start=_parse_float(q.get("from")), end=_parse_float(q.get("to")), page=int(q.get("page", "1")), page_size=min(int(q.get("page_size", "20")), 200), ) handler._write_json(200, {"code": 0, "message": "ok", "data": {"total": total, "items": [e.to_dict() for e in items]}}) def handle_list_alerts(self, handler) -> None: q = handler._query() total, items = self.app.alert_store.list_alerts( severity=q.get("severity"), status=q.get("status"), ack_status=q.get("ack_status"), start=_parse_float(q.get("from")), end=_parse_float(q.get("to")), page=int(q.get("page", "1")), page_size=min(int(q.get("page_size", "20")), 200), ) handler._write_json(200, {"code": 0, "message": "ok", "data": {"total": total, "items": [a.to_dict() for a in items]}}) def handle_get_alert(self, handler, alert_id: str) -> None: alert = self.app.alert_store.get(alert_id) if alert is None: raise _APIError(404, 40400, f"alert {alert_id} not found") handler._write_json(200, {"code": 0, "message": "ok", "data": alert.to_dict()}) def handle_ack_alert(self, handler, alert_id: str) -> None: body = handler._read_json() alert = self.app.alert_store.ack(alert_id, body.get("ack_by") or "system") if alert is None: raise _APIError(404, 40400, f"alert {alert_id} not found") handler._write_json(200, {"code": 0, "message": "ok", "data": alert.to_dict()}) def handle_close_alert(self, handler, alert_id: str) -> None: alert = self.app.alert_store.close(alert_id) if alert is None: raise _APIError(404, 40400, f"alert {alert_id} not found") handler._write_json(200, {"code": 0, "message": "ok", "data": alert.to_dict()}) def handle_batch_ack(self, handler) -> None: body = handler._read_json() ids = body.get("alert_ids") or [] acked = 0 for alert_id in ids: if self.app.alert_store.ack(str(alert_id)) is not None: acked += 1 handler._write_json(200, {"code": 0, "message": "ok", "data": {"acked": acked}}) # ------------------------------------------------------------------ # # 样本接入(端到端联调) # ------------------------------------------------------------------ # def handle_ingest(self, handler) -> None: body = handler._read_json() host_id = str(body.get("host_id") or "") if not host_id: raise _APIError(400, 40001, "host_id is required") host_group = body.get("host_group") samples = [MetricSample.from_dict(s) for s in (body.get("samples") or [])] events = self.app.ingest(host_id, host_group, samples) handler._write_json(200, {"code": 0, "message": "ok", "data": {"events": [e.to_dict() for e in events]}}) class _APIError(Exception): def __init__(self, status: int, code: int, message: str) -> None: super().__init__(message) self.status = status self.code = code self.message = message def _parse_float(value: Optional[str]) -> Optional[float]: if value in (None, ""): return None try: return float(value) except ValueError: return None def serve(app: DetectorApp, addr: str = "0.0.0.0:8080") -> APIServer: api = APIServer(app, addr) api.serve_forever() return api