#!/usr/bin/env python3 """Kafka topic and consumer lag smoke checks for the Go native production stack.""" from __future__ import annotations import argparse import json import subprocess import sys from dataclasses import asdict, dataclass DEFAULT_HOST = "114.55.58.251" DEFAULT_USER = "root" DEFAULT_BOOTSTRAP = "127.0.0.1:9092" DEFAULT_KAFKA_BIN = "/opt/kafka/current/bin" DEFAULT_TOPICS = [ "vehicle.raw.go.gb32960.v1", "vehicle.raw.go.jt808.v1", "vehicle.raw.go.yutong-mqtt.v1", "vehicle.event.go.unified.v1", ] DEFAULT_GROUPS = [ "go-history-writer", "go-stat-writer", "go-realtime-api", ] @dataclass(frozen=True) class Check: name: str status: str message: str @dataclass(frozen=True) class ConsumerLag: group: str topic: str partition: int current_offset: int | None log_end_offset: int lag: int | None def parse_topic_list(output: str) -> set[str]: return {line.strip() for line in output.splitlines() if line.strip()} def topic_checks(existing: set[str], required: list[str]) -> list[Check]: checks: list[Check] = [] for topic in required: if topic in existing: checks.append(Check("topic." + topic, "pass", "present")) else: checks.append(Check("topic." + topic, "fail", "missing")) return checks def parse_consumer_group_describe(output: str) -> list[ConsumerLag]: rows: list[ConsumerLag] = [] for line in output.splitlines(): parts = line.split() if len(parts) < 6 or parts[0] == "GROUP": continue group, topic = parts[0], parts[1] try: partition = int(parts[2]) current = parse_optional_int(parts[3]) log_end = int(parts[4]) lag = parse_optional_int(parts[5]) except ValueError: continue rows.append(ConsumerLag(group, topic, partition, current, log_end, lag)) return rows def parse_optional_int(value: str) -> int | None: value = value.strip() if value == "-": return None return int(value) def group_lag_check(group: str, rows: list[ConsumerLag], max_lag: int) -> Check: group_rows = [row for row in rows if row.group == group] if not group_rows: return Check("group." + group, "fail", "no rows") concrete_lags = [row.lag for row in group_rows if row.lag is not None] if not concrete_lags: return Check("group." + group, "fail", "no concrete lag rows") max_seen = max(concrete_lags) topics = sorted({row.topic for row in group_rows}) status = "pass" if max_seen <= max_lag else "fail" return Check( "group." + group, status, "max_lag=" + str(max_seen) + "; threshold=" + str(max_lag) + "; topics=" + ",".join(topics), ) def ssh(host: str, user: str, command: str, timeout: float) -> str: target = user + "@" + host if user else host completed = subprocess.run( ["ssh", "-o", "StrictHostKeyChecking=no", "-o", "UserKnownHostsFile=/dev/null", target, command], check=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, timeout=timeout, ) return completed.stdout def remote_topic_command(kafka_bin: str, bootstrap: str) -> str: return kafka_bin.rstrip("/") + "/kafka-topics.sh --bootstrap-server " + bootstrap + " --list" def remote_group_command(kafka_bin: str, bootstrap: str, group: str) -> str: return kafka_bin.rstrip("/") + "/kafka-consumer-groups.sh --bootstrap-server " + bootstrap + " --describe --group " + group def run(args: argparse.Namespace) -> tuple[str, list[Check]]: topic_output = ssh(args.host, args.user, remote_topic_command(args.kafka_bin, args.bootstrap), args.timeout) checks = topic_checks(parse_topic_list(topic_output), DEFAULT_TOPICS) for group in DEFAULT_GROUPS: group_output = ssh(args.host, args.user, remote_group_command(args.kafka_bin, args.bootstrap, group), args.timeout) checks.append(group_lag_check(group, parse_consumer_group_describe(group_output), args.max_lag)) status = "fail" if any(check.status == "fail" for check in checks) else "pass" return status, checks def parse_args(argv: list[str]) -> argparse.Namespace: parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--host", default=DEFAULT_HOST) parser.add_argument("--user", default=DEFAULT_USER) parser.add_argument("--bootstrap", default=DEFAULT_BOOTSTRAP) parser.add_argument("--kafka-bin", default=DEFAULT_KAFKA_BIN) parser.add_argument("--max-lag", type=int, default=100) parser.add_argument("--timeout", type=float, default=20.0) return parser.parse_args(argv) def main(argv: list[str]) -> int: args = parse_args(argv) try: status, checks = run(args) except Exception as exc: status = "fail" checks = [Check("kafka.ssh", "fail", str(exc))] print(json.dumps({ "status": status, "host": args.host, "bootstrap": args.bootstrap, "checks": [asdict(check) for check in checks], }, ensure_ascii=False, indent=2)) return 0 if status == "pass" else 1 if __name__ == "__main__": raise SystemExit(main(sys.argv[1:]))