7ebbb9a9d1
Agent-style report body for SAC aggregation; shared stats cards; RDP ban note in docs Co-authored-by: Cursor <cursoragent@cursor.com>
371 lines
12 KiB
Python
371 lines
12 KiB
Python
"""SAC-generated daily reports from ingested events (F-NOT-05 / notif-32)."""
|
||
|
||
from __future__ import annotations
|
||
|
||
import html
|
||
import logging
|
||
import uuid
|
||
from collections import Counter
|
||
from dataclasses import dataclass
|
||
from datetime import datetime, timedelta, timezone
|
||
from typing import Any
|
||
from zoneinfo import ZoneInfo
|
||
|
||
from sqlalchemy import select
|
||
from sqlalchemy.orm import Session
|
||
|
||
from app.config import Settings, get_settings
|
||
from app.models import Event, Host
|
||
from app.services.daily_report_format import (
|
||
RDP_BAN_TYPES,
|
||
SSH_BAN_TYPES,
|
||
build_report_body,
|
||
enrich_stats_for_storage,
|
||
)
|
||
from app.services.ingest import ingest_event
|
||
from app.services.notify_dispatch import notify_daily_report
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
REPORT_TYPES = {
|
||
"ssh-monitor": "report.daily.ssh",
|
||
"rdp-login-monitor": "report.daily.rdp",
|
||
}
|
||
|
||
SSH_SUCCESS = "ssh.login.success"
|
||
SSH_FAILED = "ssh.login.failed"
|
||
SSH_SUDO = "privilege.sudo.command"
|
||
RDP_SUCCESS = "rdp.login.success"
|
||
RDP_FAILED = "rdp.login.failed"
|
||
SESSION_LOGIND_NEW = "session.logind.new"
|
||
|
||
|
||
@dataclass(frozen=True)
|
||
class DailyReportResult:
|
||
host_id: int
|
||
hostname: str
|
||
report_type: str
|
||
event_id: str
|
||
created: bool
|
||
skipped_reason: str | None = None
|
||
|
||
|
||
def _now_in_tz(settings: Settings) -> datetime:
|
||
tz = ZoneInfo(settings.sac_daily_report_timezone.strip() or "Europe/Moscow")
|
||
return datetime.now(timezone.utc).astimezone(tz)
|
||
|
||
|
||
def _report_type_for_host(host: Host) -> str | None:
|
||
return REPORT_TYPES.get((host.product or "").strip())
|
||
|
||
|
||
def _day_start_in_tz(settings: Settings, ref: datetime | None = None) -> datetime:
|
||
local = ref or _now_in_tz(settings)
|
||
start = local.replace(hour=0, minute=0, second=0, microsecond=0)
|
||
return start.astimezone(timezone.utc)
|
||
|
||
|
||
def host_has_report_today(db: Session, host_id: int, report_type: str, settings: Settings) -> bool:
|
||
since = _day_start_in_tz(settings)
|
||
row = db.scalar(
|
||
select(Event.id)
|
||
.where(
|
||
Event.host_id == host_id,
|
||
Event.type == report_type,
|
||
Event.received_at >= since,
|
||
)
|
||
.limit(1)
|
||
)
|
||
return row is not None
|
||
|
||
|
||
def _ip_from_details(details: dict[str, Any] | None) -> str:
|
||
if not details:
|
||
return ""
|
||
for key in ("source_ip", "ip_address", "ip"):
|
||
val = details.get(key)
|
||
if val is not None and str(val).strip() not in ("", "-"):
|
||
return str(val).strip()
|
||
return ""
|
||
|
||
|
||
def _user_from_details(details: dict[str, Any] | None) -> str:
|
||
if not details:
|
||
return ""
|
||
for key in ("user", "username"):
|
||
val = details.get(key)
|
||
if val is not None and str(val).strip():
|
||
return str(val).strip()
|
||
return ""
|
||
|
||
|
||
def _format_user_line_ssh(details: dict[str, Any] | None) -> str | None:
|
||
user = _user_from_details(details)
|
||
if not user:
|
||
return None
|
||
if not details:
|
||
return f"👤 {user}"
|
||
tty = str(details.get("tty") or details.get("pts") or "").strip()
|
||
since = str(details.get("since") or details.get("session_since") or "").strip()
|
||
ip = _ip_from_details(details)
|
||
parts = [f"👤 {user}"]
|
||
if tty:
|
||
parts.append(f"| {tty}")
|
||
if since:
|
||
parts.append(f"| с {since}")
|
||
if ip:
|
||
parts.append(f"| 🌐 {ip}")
|
||
return " ".join(parts) if len(parts) > 1 else f"👤 {user}"
|
||
|
||
|
||
def _collect_active_users_ssh(events: list[Event]) -> list[str]:
|
||
lines: list[str] = []
|
||
seen: set[str] = set()
|
||
for ev in events:
|
||
if ev.type != SESSION_LOGIND_NEW:
|
||
continue
|
||
details = ev.details if isinstance(ev.details, dict) else None
|
||
line = _format_user_line_ssh(details)
|
||
if line and line not in seen:
|
||
lines.append(line)
|
||
seen.add(line)
|
||
if lines:
|
||
return lines
|
||
for ev in events:
|
||
if ev.type != SSH_SUCCESS:
|
||
continue
|
||
details = ev.details if isinstance(ev.details, dict) else None
|
||
user = _user_from_details(details)
|
||
if not user:
|
||
continue
|
||
key = user.lower()
|
||
if key in seen:
|
||
continue
|
||
seen.add(key)
|
||
ip = _ip_from_details(details)
|
||
line = f"👤 {user}"
|
||
if ip:
|
||
line += f" | 🌐 {ip}"
|
||
lines.append(line)
|
||
return lines
|
||
|
||
|
||
def _collect_active_users_rdp(events: list[Event]) -> list[str]:
|
||
users: list[str] = []
|
||
seen: set[str] = set()
|
||
for ev in events:
|
||
if ev.type != RDP_SUCCESS:
|
||
continue
|
||
details = ev.details if isinstance(ev.details, dict) else None
|
||
user = _user_from_details(details)
|
||
if not user:
|
||
continue
|
||
key = user.lower()
|
||
if key in seen:
|
||
continue
|
||
seen.add(key)
|
||
users.append(f"👤 {user}")
|
||
return users
|
||
|
||
|
||
def _aggregate_ssh(events: list[Event]) -> dict[str, Any]:
|
||
ok = failed = sudo = 0
|
||
failed_ips: Counter[str] = Counter()
|
||
for ev in events:
|
||
if ev.type == SSH_SUCCESS:
|
||
ok += 1
|
||
elif ev.type == SSH_FAILED:
|
||
failed += 1
|
||
ip = _ip_from_details(ev.details if isinstance(ev.details, dict) else None)
|
||
if ip:
|
||
failed_ips[ip] += 1
|
||
elif ev.type == SSH_SUDO:
|
||
sudo += 1
|
||
top_ips = [f"{ip} — {cnt}" for ip, cnt in failed_ips.most_common(5)]
|
||
active_users = _collect_active_users_ssh(events)
|
||
return {
|
||
"successful_ssh": ok,
|
||
"failed_ssh": failed,
|
||
"sudo_commands": sudo,
|
||
"active_bans": sum(1 for e in events if e.type in SSH_BAN_TYPES),
|
||
"top_failed_ips": top_ips,
|
||
"active_users": active_users,
|
||
}
|
||
|
||
|
||
def _aggregate_rdp(events: list[Event]) -> dict[str, Any]:
|
||
ok = failed = 0
|
||
failed_ips: Counter[str] = Counter()
|
||
for ev in events:
|
||
details = ev.details if isinstance(ev.details, dict) else None
|
||
if ev.type == RDP_SUCCESS:
|
||
ok += 1
|
||
elif ev.type == RDP_FAILED:
|
||
failed += 1
|
||
ip = _ip_from_details(details)
|
||
if ip:
|
||
failed_ips[ip] += 1
|
||
active_users = _collect_active_users_rdp(events)
|
||
unique = []
|
||
seen: set[str] = set()
|
||
for line in active_users:
|
||
u = line.replace("👤", "").strip().split("|")[0].strip()
|
||
if u.lower() not in seen:
|
||
seen.add(u.lower())
|
||
unique.append(u)
|
||
return {
|
||
"rdp_success": ok,
|
||
"rdp_failed": failed,
|
||
"active_bans": sum(1 for e in events if e.type in RDP_BAN_TYPES),
|
||
"top_failed_ips": [f"{ip} — {cnt}" for ip, cnt in failed_ips.most_common(5)],
|
||
"active_users": active_users,
|
||
"unique_users": unique,
|
||
}
|
||
|
||
|
||
def _build_report_body_ssh(host: Host, stats: dict[str, Any], when_local: datetime) -> str:
|
||
enriched = enrich_stats_for_storage("ssh", stats)
|
||
return build_report_body("ssh", host, enriched, when_local, sac_generated=True)
|
||
|
||
|
||
def _build_report_body_rdp(host: Host, stats: dict[str, Any], when_local: datetime) -> str:
|
||
enriched = enrich_stats_for_storage("windows", stats)
|
||
return build_report_body("windows", host, enriched, when_local, sac_generated=True)
|
||
|
||
|
||
def _details_from_body(body: str, stats: dict[str, Any]) -> dict[str, Any]:
|
||
escaped = html.escape(body)
|
||
# UI (Vue) допускает div; для единообразия — обёртка без br (Telegram читает через sanitize)
|
||
report_html = f'<div class="agent-report">{escaped.replace(chr(10), "<br>")}</div>'
|
||
return {
|
||
"stats": stats,
|
||
"report_body": body,
|
||
"report_format": "plain",
|
||
"report_html": report_html,
|
||
"generated_by": "sac",
|
||
}
|
||
|
||
|
||
def _events_last_24h(db: Session, host_id: int) -> list[Event]:
|
||
since = datetime.now(timezone.utc) - timedelta(hours=24)
|
||
return list(
|
||
db.scalars(
|
||
select(Event)
|
||
.where(Event.host_id == host_id, Event.occurred_at >= since)
|
||
.order_by(Event.occurred_at.desc())
|
||
).all()
|
||
)
|
||
|
||
|
||
def generate_daily_report_for_host(db: Session, host: Host, settings: Settings | None = None) -> DailyReportResult | None:
|
||
cfg = settings or get_settings()
|
||
report_type = _report_type_for_host(host)
|
||
if report_type is None:
|
||
return None
|
||
|
||
if cfg.sac_daily_report_skip_if_agent_sent and host_has_report_today(db, host.id, report_type, cfg):
|
||
return DailyReportResult(
|
||
host_id=host.id,
|
||
hostname=host.hostname,
|
||
report_type=report_type,
|
||
event_id="",
|
||
created=False,
|
||
skipped_reason="agent_report_exists_today",
|
||
)
|
||
|
||
events = _events_last_24h(db, host.id)
|
||
if not events and cfg.sac_daily_report_require_activity:
|
||
return DailyReportResult(
|
||
host_id=host.id,
|
||
hostname=host.hostname,
|
||
report_type=report_type,
|
||
event_id="",
|
||
created=False,
|
||
skipped_reason="no_events_24h",
|
||
)
|
||
|
||
when_local = _now_in_tz(cfg)
|
||
if report_type == "report.daily.ssh":
|
||
stats = enrich_stats_for_storage("ssh", _aggregate_ssh(events))
|
||
body = _build_report_body_ssh(host, stats, when_local)
|
||
title = "Ежедневный отчёт SSH"
|
||
summary = (
|
||
f"SSH 24ч: успех {stats['successful_logins']}, неудач {stats['failed_logins']}, "
|
||
f"sudo {stats['sudo_commands']}"
|
||
)
|
||
else:
|
||
stats = enrich_stats_for_storage("windows", _aggregate_rdp(events))
|
||
body = _build_report_body_rdp(host, stats, when_local)
|
||
title = "Ежедневный отчёт Windows"
|
||
summary = (
|
||
f"RDP 24ч: успех {stats['successful_logins']}, неудач {stats['failed_logins']}, "
|
||
f"банов {stats['active_bans']}"
|
||
)
|
||
|
||
event_id = str(uuid.uuid4())
|
||
day_key = when_local.strftime("%Y-%m-%d")
|
||
payload = {
|
||
"schema_version": "1.0",
|
||
"event_id": event_id,
|
||
"occurred_at": datetime.now(timezone.utc).isoformat(),
|
||
"source": {
|
||
"product": host.product or "unknown",
|
||
"product_version": host.product_version or "sac",
|
||
},
|
||
"host": {
|
||
"hostname": host.hostname,
|
||
"os_family": host.os_family or ("windows" if report_type == "report.daily.rdp" else "linux"),
|
||
"display_name": host.display_name,
|
||
"ipv4": host.ipv4,
|
||
},
|
||
"category": "report",
|
||
"type": report_type,
|
||
"severity": "info",
|
||
"title": title,
|
||
"summary": summary,
|
||
"details": _details_from_body(body, stats),
|
||
"dedup_key": f"sac|{host.id}|{report_type}|{day_key}",
|
||
}
|
||
|
||
event, created = ingest_event(db, payload)
|
||
if created:
|
||
notify_daily_report(event, db=db)
|
||
return DailyReportResult(
|
||
host_id=host.id,
|
||
hostname=host.hostname,
|
||
report_type=report_type,
|
||
event_id=event.event_id,
|
||
created=created,
|
||
skipped_reason=None if created else "duplicate_event_id",
|
||
)
|
||
|
||
|
||
def run_daily_reports(db: Session, settings: Settings | None = None, *, force: bool = False) -> list[DailyReportResult]:
|
||
cfg = settings or get_settings()
|
||
if not cfg.sac_daily_report_enabled:
|
||
logger.info("daily report disabled (SAC_DAILY_REPORT_ENABLED=false)")
|
||
return []
|
||
|
||
local_now = _now_in_tz(cfg)
|
||
if not force and local_now.hour < cfg.sac_daily_report_hour:
|
||
logger.info(
|
||
"daily report skipped: local hour %s < SAC_DAILY_REPORT_HOUR=%s",
|
||
local_now.hour,
|
||
cfg.sac_daily_report_hour,
|
||
)
|
||
return []
|
||
|
||
hosts = db.scalars(select(Host).order_by(Host.hostname)).all()
|
||
results: list[DailyReportResult] = []
|
||
for host in hosts:
|
||
try:
|
||
res = generate_daily_report_for_host(db, host, cfg)
|
||
if res is not None:
|
||
results.append(res)
|
||
except Exception:
|
||
logger.exception("daily report failed host_id=%s hostname=%s", host.id, host.hostname)
|
||
db.commit()
|
||
created_n = sum(1 for r in results if r.created)
|
||
logger.info("daily report finished hosts=%s created=%s", len(results), created_n)
|
||
return results
|