190 lines
7.2 KiB
Python
190 lines
7.2 KiB
Python
#!/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()
|