#!/usr/bin/env python3 """Агрегация истории внешних TCP-соединений между запусками пайплайна. Читает host_facts/*.json, извлекает внешних пиров (по логике detect_external_deps), обновляет connection_history.json и удаляет записи старше ttl_days дней. Использование: python3 scripts/aggregate_connections.py python3 scripts/aggregate_connections.py --facts-dir host_facts --history connection_history.json --ttl-days 30 """ import argparse import json import sys from datetime import datetime, timedelta, timezone from pathlib import Path # Добавляем директорию скрипта в path для импорта generate_report sys.path.insert(0, str(Path(__file__).parent)) from generate_report import ( # noqa: E402 _classify_external_ip, _get_listening_ports, build_ip_to_host, load_facts, parse_conntrack, parse_connections, parse_ports, ) def _empty_history(now: datetime) -> dict: return { "version": 1, "last_updated": now.isoformat(), "peers": [], } def _load_history(path: Path, now: datetime) -> dict: if not path.exists(): return _empty_history(now) try: data = json.loads(path.read_text(encoding="utf-8")) if not isinstance(data, dict) or "peers" not in data: print(f"Предупреждение: {path} имеет неверный формат, сброс.", file=sys.stderr) return _empty_history(now) return data except Exception as exc: print(f"Предупреждение: не удалось прочитать {path}: {exc}", file=sys.stderr) return _empty_history(now) def _peer_key(to_host: str, to_port: int, to_service: str, proto: str, subnet: str, peer_addr: str) -> tuple: return (to_host, to_port, to_service, proto, subnet, peer_addr) def _upsert_peer(peers_index: dict, hostname: str, port: int, service: str, proto: str, ext_type: str, subnet: str, peer_addr: str, now: datetime) -> None: """Добавляет или обновляет запись пира в индексе.""" if ext_type != "local": return # сохраняем только локальные пиры key = _peer_key(hostname, port, service, proto, subnet, peer_addr) now_iso = now.isoformat() if key in peers_index: peers_index[key]["last_seen"] = now_iso else: peers_index[key] = { "to_host": hostname, "to_port": port, "to_service": service, "proto": proto, "ext_type": ext_type, "subnet": subnet, "peer_addr": peer_addr, "first_seen": now_iso, "last_seen": now_iso, } def aggregate(facts_dir: Path, history_path: Path, ttl_days: int) -> None: hosts = load_facts(facts_dir) if not hosts: print(f"Предупреждение: в '{facts_dir}' нет JSON-файлов, история не обновлена.", file=sys.stderr) return ip_to_host = build_ip_to_host(hosts) now = datetime.now(timezone.utc) history = _load_history(history_path, now) # Строим индекс существующих записей key → peer_dict peers_index: dict = {} for peer in history.get("peers", []): key = _peer_key( peer["to_host"], peer["to_port"], peer.get("to_service", "—"), peer.get("proto", "tcp"), peer["subnet"], peer["peer_addr"], ) peers_index[key] = peer # Обходим все хосты и извлекаем внешних пиров for hostname, f in hosts.items(): listening = _get_listening_ports(f) # Маппинг порт→процесс для обогащения conntrack port_to_proc: dict = {} for p in parse_ports(f.get("ports_raw", "")): if p["process"] != "—" and p["port"].isdigit(): port_to_proc[p["port"]] = p["process"] # Источник 1: ss conns_ss = parse_connections(f.get("connections_raw", "")) # Источник 2: conntrack (обогащаем именами процессов) conns_ct = parse_conntrack(f.get("conntrack_raw", "")) for c in conns_ct: if c["process"] == "—" and c["local_port"] in port_to_proc: c["process"] = port_to_proc[c["local_port"]] # Дедупликация ss_seen = {(c["local_port"], c["peer_addr"], c["peer_port"]) for c in conns_ss} conns_ct_new = [ c for c in conns_ct if (c["local_port"], c["peer_addr"], c["peer_port"]) not in ss_seen ] all_conns = conns_ss + conns_ct_new for conn in all_conns: lp = conn["local_port"] if not lp.isdigit(): continue local_port_int = int(lp) if local_port_int not in listening: continue peer_addr = conn["peer_addr"] if peer_addr in ip_to_host: continue ext_type, subnet = _classify_external_ip(peer_addr) service = conn["process"] proto = conn.get("proto", "tcp") _upsert_peer(peers_index, hostname, local_port_int, service, proto, ext_type, subnet, peer_addr, now) # Удаляем записи старше ttl_days cutoff = now - timedelta(days=ttl_days) retained = [] expired = 0 for peer in peers_index.values(): try: last_seen = datetime.fromisoformat(peer["last_seen"]) # Для сравнения делаем оба значения timezone-aware if last_seen.tzinfo is None: last_seen = last_seen.replace(tzinfo=timezone.utc) if last_seen >= cutoff: retained.append(peer) else: expired += 1 except Exception: retained.append(peer) # при ошибке парсинга оставляем history["version"] = 1 history["last_updated"] = now.isoformat() history["peers"] = retained history_path.write_text(json.dumps(history, ensure_ascii=False, indent=2), encoding="utf-8") print( f"connection_history.json — {len(retained)} пиров" f" (обновлено: {len(peers_index)}, удалено по TTL: {expired})" ) def main() -> None: parser = argparse.ArgumentParser(description="Агрегация истории внешних TCP-соединений") parser.add_argument("--facts-dir", default="host_facts", help="Директория с host_facts/*.json") parser.add_argument("--history", default="connection_history.json", help="Файл истории") parser.add_argument("--ttl-days", type=int, default=30, help="Срок хранения пира в днях") args = parser.parse_args() facts_dir = Path(args.facts_dir) history_path = Path(args.history) if not facts_dir.exists(): sys.exit(f"Ошибка: директория '{facts_dir}' не найдена") aggregate(facts_dir, history_path, args.ttl_days) if __name__ == "__main__": main()