"""Correlate RDG 303/303-failed with workstation rdp.login.success (1149).""" from __future__ import annotations from datetime import timedelta, timezone from sqlalchemy import select from sqlalchemy.orm import Session from sqlalchemy.orm.attributes import flag_modified from app.models import Event from app.services.host_sessions import ( SESSION_TERMINATED_AT_KEY, _details_dict, _event_login_user, ) from app.services.rdg_client_host import find_windows_host_by_ipv4 from app.services.rdg_session_flap import ( RDG_END_TYPES, RDG_SUCCESS_TYPE, _event_user, event_internal_ip, find_rdg_success_before_end, ) SESSION_CLOSED_BY_RDG_AT_KEY = "session_closed_by_rdg_at" SESSION_CLOSED_BY_RDG_EVENT_ID_KEY = "session_closed_by_rdg_event_id" USER_ENRICHED_FROM_RDG_EVENT_ID_KEY = "user_enriched_from_rdg_event_id" # RCM 1149 via RD Gateway often has empty Param1/Param2; RDG 302 has the account. RDG_LOGIN_USER_ENRICH_WINDOW = timedelta(minutes=5) WORKSTATION_LOGIN_TYPE = "rdp.login.success" def _as_utc(dt): if dt is None: return None if dt.tzinfo is None: return dt.replace(tzinfo=timezone.utc) return dt.astimezone(timezone.utc) def normalize_sam_account(user: str) -> str: text = (user or "").strip() if "\\" in text: return text.split("\\")[-1].strip().casefold() if "@" in text: return text.split("@")[0].strip().casefold() return text.casefold() def users_match_rdg(login_user: str, rdg_user: str) -> bool: left = normalize_sam_account(login_user) right = normalize_sam_account(rdg_user) return bool(left and right and left == right) def _login_user_missing(event: Event) -> bool: details = _details_dict(event) for key in ("user", "username"): val = details.get(key) if val is not None and str(val).strip() not in ("", "-"): return False return True def _apply_rdg_user_to_login(login_event: Event, *, rdg_event: Event, rdg_user: str) -> None: details = dict(_details_dict(login_event)) details["user"] = rdg_user details[USER_ENRICHED_FROM_RDG_EVENT_ID_KEY] = rdg_event.id login_event.details = details flag_modified(login_event, "details") summary = (login_event.summary or "").strip() if summary.startswith("RCM 1149") and rdg_user not in summary: rest = summary[len("RCM 1149") :].strip() login_event.summary = f"RCM 1149 {rdg_user} {rest}".strip() def find_rdg_success_for_workstation_login(db: Session, login_event: Event) -> Event | None: """Nearest RDG 302 for this workstation IP within the enrich window.""" if login_event.type != WORKSTATION_LOGIN_TYPE: return None host = login_event.host if host is None or not (host.ipv4 or "").strip(): return None workstation_ip = host.ipv4.strip() login_at = _as_utc(login_event.occurred_at) window_start = login_at - RDG_LOGIN_USER_ENRICH_WINDOW window_end = login_at + RDG_LOGIN_USER_ENRICH_WINDOW candidates = db.scalars( select(Event) .where( Event.type == RDG_SUCCESS_TYPE, Event.occurred_at >= window_start, Event.occurred_at <= window_end, Event.id != login_event.id, ) .order_by(Event.occurred_at.desc()) ).all() best: Event | None = None best_delta: timedelta | None = None for rdg in candidates: if event_internal_ip(rdg) != workstation_ip: continue if not _event_user(rdg): continue delta = abs(_as_utc(rdg.occurred_at) - login_at) if best is None or best_delta is None or delta < best_delta: best = rdg best_delta = delta return best def enrich_workstation_login_user_from_rdg(db: Session, login_event: Event) -> Event | None: """Fill details.user when RCM 1149 EventLog left Param1/Param2 empty (seen on some Win10 Pro).""" if login_event.type != WORKSTATION_LOGIN_TYPE: return None if not _login_user_missing(login_event): return None rdg = find_rdg_success_for_workstation_login(db, login_event) if rdg is None: return None rdg_user = _event_user(rdg) if not rdg_user: return None _apply_rdg_user_to_login(login_event, rdg_event=rdg, rdg_user=rdg_user) return login_event def enrich_empty_login_from_rdg_success(db: Session, rdg_success_event: Event) -> Event | None: """Backfill empty workstation 1149 when RDG 302 is ingested after it.""" if rdg_success_event.type != RDG_SUCCESS_TYPE: return None internal_ip = event_internal_ip(rdg_success_event) rdg_user = _event_user(rdg_success_event) if not internal_ip or not rdg_user: return None client_host = find_windows_host_by_ipv4(db, internal_ip) if client_host is None: return None rdg_at = _as_utc(rdg_success_event.occurred_at) window_start = rdg_at - RDG_LOGIN_USER_ENRICH_WINDOW window_end = rdg_at + RDG_LOGIN_USER_ENRICH_WINDOW candidates = db.scalars( select(Event) .where( Event.host_id == client_host.id, Event.type == WORKSTATION_LOGIN_TYPE, Event.occurred_at >= window_start, Event.occurred_at <= window_end, Event.id != rdg_success_event.id, ) .order_by(Event.occurred_at.desc()) ).all() for login in candidates: if not _login_user_missing(login): continue _apply_rdg_user_to_login(login, rdg_event=rdg_success_event, rdg_user=rdg_user) return login return None def event_closed_by_rdg(event: Event) -> bool: details = _details_dict(event) at = details.get(SESSION_CLOSED_BY_RDG_AT_KEY) return at is not None and str(at).strip() != "" def mark_login_closed_by_rdg(login_event: Event, *, rdg_end_event: Event) -> None: details = dict(_details_dict(login_event)) details[SESSION_CLOSED_BY_RDG_AT_KEY] = rdg_end_event.occurred_at.isoformat() details[SESSION_CLOSED_BY_RDG_EVENT_ID_KEY] = rdg_end_event.id login_event.details = details flag_modified(login_event, "details") def _login_already_closed(login_event: Event) -> bool: from app.services.rdp_session_logoff import event_closed_by_logoff details = _details_dict(login_event) if details.get("session_terminated") is True: return True at = details.get(SESSION_TERMINATED_AT_KEY) if at is not None and str(at).strip() != "": return True if event_closed_by_rdg(login_event): return True return event_closed_by_logoff(login_event) def find_workstation_login_for_rdg_end(db: Session, rdg_end_event: Event) -> Event | None: if rdg_end_event.type not in RDG_END_TYPES: return None internal_ip = event_internal_ip(rdg_end_event) if not internal_ip: return None client_host = find_windows_host_by_ipv4(db, internal_ip) if client_host is None: return None rdg_user = _event_user(rdg_end_event) if not rdg_user: return None end_at = rdg_end_event.occurred_at candidates = db.scalars( select(Event) .where( Event.host_id == client_host.id, Event.type == WORKSTATION_LOGIN_TYPE, Event.occurred_at <= end_at, Event.id != rdg_end_event.id, ) .order_by(Event.occurred_at.desc()) ).all() for login in candidates: if _login_already_closed(login): continue login_user = (_event_login_user(login) or "").strip() if login_user in ("", "-"): login_user = "" # RCM 1149 may lack user (empty Param1); still close by workstation IP + open session. if login_user and not users_match_rdg(login_user, rdg_user): continue return login return None def find_rdg_end_after_workstation_login(db: Session, login_event: Event) -> Event | None: """Runtime lookup for historical events without persisted close flag.""" if login_event.type != WORKSTATION_LOGIN_TYPE: return None if _login_already_closed(login_event): return None host = login_event.host if host is None or not host.ipv4: return None workstation_ip = host.ipv4.strip() login_user = _event_login_user(login_event) if not login_user: return None login_at = login_event.occurred_at candidates = db.scalars( select(Event) .where( Event.type.in_(RDG_END_TYPES), Event.occurred_at >= login_at, Event.id != login_event.id, ) .order_by(Event.occurred_at.asc()) ).all() for end in candidates: if event_internal_ip(end) != workstation_ip: continue if not users_match_rdg(_event_user(end), login_user): continue if find_rdg_success_before_end(db, end) is not None: continue return end return None def resolve_workstation_login_closed(db: Session, login_event: Event) -> bool: if event_closed_by_rdg(login_event): return True return find_rdg_end_after_workstation_login(db, login_event) is not None def close_workstation_session_for_rdg_end(db: Session, rdg_end_event: Event) -> Event | None: """On RDG disconnect, mark matching workstation login as session-closed.""" if rdg_end_event.type not in RDG_END_TYPES: return None if find_rdg_success_before_end(db, rdg_end_event) is not None: return None login = find_workstation_login_for_rdg_end(db, rdg_end_event) if login is None: return None mark_login_closed_by_rdg(login, rdg_end_event=rdg_end_event) return login