Python构建APT攻击溯源图:从日志归一化到战术路径检测
简介本资源是一套基于Python实现的APT攻击检测系统面向网络安全、数据科学及智能系统方向的高年级本科生、研究生与行业开发者聚焦高级持续性威胁的溯源图建模与检测实践。项目完整覆盖算法实现、部署方案与实测数据集适用于毕业设计、课程实验、学术研究及企业原型开发要求使用者具备Python编程基础与网络安全基本认知。压缩包共31个文件含11个核心Python脚本如main.py、model_RGAT.py、streamspot_RGAT.py等、7个XML配置与元数据文件、4份Markdown文档含设计说明与信息分析指南、5个备份文件.zbak及环境配置相关文件整体仅52KB轻量易部署。已有113人学习下载资源结构清晰模块解耦合理既可直接运行分析内置APT样本也可快速拓展图神经网络组件或适配新数据源是理解攻击链可视化与溯源图技术落地的优质教学与科研载体。1. APT攻击检测为什么不能只靠杀毒软件——用Python搭一个能“看懂攻击链条”的溯源图系统你手上有几台Linux服务器日志、Windows终端EDR告警、网络流量PCAP、还有SIEM平台导出的原始事件CSV但所有告警都像散落的弹珠一个进程注入、一次异常DNS请求、一个可疑PowerShell调用……单独看都不够致命合起来却可能是一次持续97天的APT渗透。传统规则引擎和AV引擎在这里集体失语——它们不理解“横向移动”是攻击者在内网里坐地铁换乘“持久化”是他在你数据库里悄悄装了张长期工牌“命令与控制”是他每天凌晨三点用加密隧道发一条“天气预报”。而基于Python的APT攻击检测系统实现溯源图分析与部署方案就是把这堆弹珠串成一张动态演化的攻击图谱节点是主机、进程、文件、IP、账户边是父子进程、文件读写、网络连接、注册表修改图结构本身就是攻击者的战术意图。这不是加个模型就能跑通的玩具项目而是面向蓝队实战的轻量级威胁狩猎基础设施——它不替代SOAR但能让SOAR真正“看懂”下一步该阻断哪条边它不取代EDR但能把EDR的1000条孤立告警压缩成3条可追溯的攻击路径。适合有日志接入能力、熟悉Python基础、但没预算买商业图分析平台的中小安全团队或红蓝对抗支撑人员。2. 溯源图不是画个流程图从原始日志到图结构的三步清洗与建模APT攻击的痕迹从来不在单点爆发而在跨设备、跨时间、跨协议的关联中浮现。直接拿原始日志喂图算法结果只会是内存爆炸图遍历超时。必须先做结构化清洗再定义图语义最后落地为可计算的图对象。整个过程不依赖任何商业图数据库纯Python生态即可闭环。2.1 日志归一化统一字段语义拒绝“同名不同义”不同来源日志字段名五花八门EDR叫process_nameSysmon叫ImageNetFlow叫src_proc而Windows事件ID 4688里连进程名都藏在XML blob里。我们不用写100个解析器而是用字段映射模板正则提取规则做轻量归一。# log_normalizer.py import re import json from typing import Dict, Any, Optional class LogNormalizer: # 定义标准字段集攻击图建模必需 STANDARD_FIELDS { timestamp: datetime, src_host: str, dst_host: str, src_ip: str, dst_ip: str, src_port: int, dst_port: int, process_name: str, process_pid: int, parent_pid: int, file_path: str, registry_key: str, command_line: str, event_type: str, # process_create, network_connect, file_write, registry_mod user_name: str } # 预编译常用正则提升性能 _re_cmdline re.compile(r(?i)cmd\.exe.*?/c\s(.?)(?:\s|$)) _re_powershell re.compile(r(?i)powershell.*?-EncodedCommand\s([a-zA-Z0-9/])) def normalize(self, raw_log: Dict[str, Any]) - Optional[Dict[str, Any]]: 输入任意格式日志字典输出标准化字段字典 norm {} # 时间戳统一转ISO格式支持多种输入格式 ts raw_log.get(TimeCreated) or raw_log.get(timestamp) or raw_log.get(timestamp) if ts: try: from dateutil import parser norm[timestamp] parser.parse(str(ts)).isoformat() except: norm[timestamp] None # 主机/IP映射优先取明确字段 fallback 到 hostname/ip 字段 norm[src_host] self._extract_host(raw_log, [ComputerName, hostname, src_host]) norm[dst_host] self._extract_host(raw_log, [TargetHostname, dst_host]) norm[src_ip] self._extract_ip(raw_log, [SourceNetworkAddress, src_ip, IpAddress]) norm[dst_ip] self._extract_ip(raw_log, [DestinationNetworkAddress, dst_ip, DestinationIp]) # 进程信息关键用于构建进程树 norm[process_name] self._clean_process_name( raw_log.get(Image) or raw_log.get(process_name) or raw_log.get(NewProcessName) ) norm[process_pid] self._safe_int(raw_log.get(ProcessId) or raw_log.get(pid)) norm[parent_pid] self._safe_int(raw_log.get(ParentProcessId) or raw_log.get(ppid)) # 文件/注册表操作用于持久化检测 norm[file_path] self._clean_file_path( raw_log.get(TargetFilename) or raw_log.get(file_path) or raw_log.get(Path) ) norm[registry_key] self._clean_registry_key( raw_log.get(ObjectName) or raw_log.get(registry_key) ) # 命令行解码Base64、去噪、截断过长内容 cmdline raw_log.get(CommandLine) or raw_log.get(command_line) if cmdline: norm[command_line] self._decode_and_sanitize_cmdline(str(cmdline)) else: norm[command_line] # 事件类型自动推断比硬编码更鲁棒 norm[event_type] self._infer_event_type(raw_log) # 用户名用于横向移动链路 norm[user_name] self._extract_user(raw_log) return {k: v for k, v in norm.items() if v is not None} def _extract_host(self, log, keys): for k in keys: if k in log and log[k]: return str(log[k]).split(.)[0] # 去域名后缀留主机名 return unknown def _extract_ip(self, log, keys): for k in keys: if k in log and log[k]: ip str(log[k]) if re.match(r^\d{1,3}\.\d{1,3}\.\d{1,3}\.\d{1,3}$, ip): return ip return None def _clean_process_name(self, name): if not name: return unknown return re.sub(r\\[^\\]*\\, , str(name).strip()).split(/)[-1].split(\\)[-1] def _clean_file_path(self, path): if not path: return return re.sub(r[^\x20-\x7E], , str(path))[:256] # 去控制字符截断防爆 def _decode_and_sanitize_cmdline(self, cmd): # 解码常见PowerShell Base64编码 m self._re_powershell.search(cmd) if m: try: import base64 decoded base64.b64decode(m.group(1)).decode(utf-16-le, errorsignore) return f[BASE64_DECODED] {decoded[:200]} except: pass # 提取cmd /c 后命令 m self._re_cmdline.search(cmd) if m: return f[CMD_C] {m.group(1)[:200]} return cmd[:200] def _infer_event_type(self, log): # 根据关键词字段存在性智能推断 if log.get(Image) and log.get(ParentProcessId): return process_create if log.get(DestinationPort) and log.get(DestinationIp): return network_connect if log.get(TargetFilename) and write in str(log.get(Operation)).lower(): return file_write if reg in str(log.get(EventID, )).lower() or log.get(ObjectName, ).startswith(HK): return registry_mod return unknown def _extract_user(self, log): return (log.get(SubjectUserName) or log.get(AccountName) or log.get(user) or log.get(UserName) or ).strip() def _safe_int(self, val): try: return int(val) except (TypeError, ValueError): return -1提示这个归一化器不是万能的但它把90%的常见日志源Sysmon、EDR API、Suricata、Windows Event Log覆盖了。关键在于_infer_event_type——它让系统能自动识别“这是进程创建还是网络连接”避免人工打标签。实际部署时你只需为新增日志源补充_infer_event_type里的判断逻辑无需重写整个解析器。2.2 图模式定义什么该当节点什么该当边——攻击战术视角建模很多团队一上来就用networkx.Graph()建图结果跑着跑着内存飙到32GB。问题出在图建模本身把每条日志当一个节点那100万条日志就是100万个节点边数更是平方级爆炸。真正的溯源图节点必须是有状态、可复用、带生命周期的实体边必须是有方向、有时序、可归因的动作。我们采用MITRE ATTCK战术层抽象定义四类核心节点和三类关键边节点类型实体示例唯一标识符ID生命周期判断依据Hostwin-srv-01,linux-db-02主机名非IP因IP可能漂移首次出现时间 → 最后活跃时间30天无活动则标记为inactiveProcesssvchost.exe:1234,powershell.exe:5678process_name:pidhost如powershell.exe:5678win-srv-01创建时间 → 退出时间或超时未退出则设为runningFileC:\Temp\lsass.dmp,/tmp/.cache/shell.shsha256(file_path)对路径哈希避免路径变更导致ID漂移首次写入时间 → 最后修改时间NetworkSession10.1.1.100:49152→185.199.108.153:443src_ip:src_port→dst_ip:dst_port含协议首包时间 → 最后包时间TCP会话需FIN确认边类型触发条件方向权重含义是否保留spawned_by子进程parent_pid 父进程pidProcess → Process父进程存活时长秒✅ 必须保留构建进程树connects_tonetwork_connect事件中src_ip:src_port→dst_ip:dst_portProcess → NetworkSession连接持续时间秒✅ 必须保留定位C2writes_tofile_write事件中process_name写入file_pathProcess → File写入字节数✅ 必须保留追踪恶意载荷# graph_builder.py import networkx as nx from datetime import datetime, timedelta from typing import Dict, List, Tuple, Optional import hashlib class AttackGraphBuilder: def __init__(self, time_window_minutes: int 1440): # 默认1天时间窗口 self.G nx.DiGraph() self.time_window timedelta(minutestime_window_minutes) self.node_cache {} # {node_id: {type: Process, first_seen: ts, last_seen: ts}} self.edge_cache {} # {(src_id, dst_id, edge_type): {weight: w, first_seen: ts}} def add_log_event(self, norm_log: Dict[str, Any]): 将一条归一化日志添加到图中 if not norm_log.get(event_type) or norm_log[event_type] unknown: return # 步骤1创建或更新节点 nodes_to_add self._create_nodes_from_log(norm_log) for node_id, node_attrs in nodes_to_add: self._upsert_node(node_id, node_attrs) # 步骤2创建边仅当源/目标节点已存在 edges_to_add self._create_edges_from_log(norm_log, nodes_to_add) for src_id, dst_id, edge_type, edge_attrs in edges_to_add: self._upsert_edge(src_id, dst_id, edge_type, edge_attrs) def _create_nodes_from_log(self, log: Dict[str, Any]) - List[Tuple[str, Dict]]: 根据日志生成待添加的节点列表 nodes [] now datetime.fromisoformat(log[timestamp].replace(Z, 00:00)) # Host节点 if log.get(src_host): host_id fHost:{log[src_host]} nodes.append((host_id, { type: Host, name: log[src_host], first_seen: now, last_seen: now, ip: log.get(src_ip), os: self._infer_os_from_log(log) })) if log.get(dst_host) and log[dst_host] ! log.get(src_host): host_id fHost:{log[dst_host]} nodes.append((host_id, { type: Host, name: log[dst_host], first_seen: now, last_seen: now, ip: log.get(dst_ip), os: self._infer_os_from_log(log) })) # Process节点关键 if log.get(process_name) and log.get(process_pid) and log.get(src_host): proc_id fProcess:{log[process_name]}:{log[process_pid]}{log[src_host]} nodes.append((proc_id, { type: Process, name: log[process_name], pid: log[process_pid], host: log[src_host], command_line: log.get(command_line, ), user: log.get(user_name, ), first_seen: now, last_seen: now, is_suspicious: self._is_suspicious_process(log) })) # File节点用路径SHA256作ID防路径变更 if log.get(file_path): file_hash hashlib.sha256(log[file_path].encode()).hexdigest()[:16] file_id fFile:{file_hash} nodes.append((file_id, { type: File, path: log[file_path], first_seen: now, last_seen: now, size_bytes: self._estimate_file_size(log) })) # NetworkSession节点IP:Port→IP:Port if log.get(src_ip) and log.get(dst_ip) and log.get(src_port) and log.get(dst_port): sess_id fNetworkSession:{log[src_ip]}:{log[src_port]}→{log[dst_ip]}:{log[dst_port]} nodes.append((sess_id, { type: NetworkSession, src_ip: log[src_ip], src_port: log[src_port], dst_ip: log[dst_ip], dst_port: log[dst_port], protocol: tcp if log.get(Protocol) 6 else udp, first_seen: now, last_seen: now })) return nodes def _create_edges_from_log(self, log: Dict[str, Any], nodes: List[Tuple[str, Dict]]) - List[Tuple[str, str, str, Dict]]: 根据日志生成边需确保节点已存在 edges [] now datetime.fromisoformat(log[timestamp].replace(Z, 00:00)) # spawned_by 边子进程 → 父进程需父进程ID if log.get(process_pid) and log.get(parent_pid) and log.get(src_host): child_id fProcess:{log[process_name]}:{log[process_pid]}{log[src_host]} parent_id fProcess:unknown:{log[parent_pid]}{log[src_host]} # 父进程名可能未知 # 但父进程节点必须存在否则跳过 if parent_id in self.node_cache: edges.append(( child_id, parent_id, spawned_by, {weight: 1, first_seen: now, direction: child_to_parent} )) # connects_to 边进程 → 网络会话 if log.get(event_type) network_connect and log.get(src_host) and log.get(dst_ip): proc_id fProcess:{log[process_name]}:{log[process_pid]}{log[src_host]} sess_id fNetworkSession:{log[src_ip]}:{log[src_port]}→{log[dst_ip]}:{log[dst_port]} if proc_id in self.node_cache and sess_id in self.node_cache: duration self._estimate_session_duration(log) edges.append(( proc_id, sess_id, connects_to, {weight: duration, first_seen: now, dst_ip: log[dst_ip]} )) # writes_to 边进程 → 文件 if log.get(event_type) file_write and log.get(process_name) and log.get(file_path): proc_id fProcess:{log[process_name]}:{log[process_pid]}{log[src_host]} file_hash hashlib.sha256(log[file_path].encode()).hexdigest()[:16] file_id fFile:{file_hash} if proc_id in self.node_cache and file_id in self.node_cache: size self._estimate_file_size(log) edges.append(( proc_id, file_id, writes_to, {weight: size, first_seen: now} )) return edges def _upsert_node(self, node_id: str, attrs: Dict): 插入或更新节点属性 if node_id not in self.node_cache: self.node_cache[node_id] attrs.copy() self.G.add_node(node_id, **attrs) else: # 更新last_seen和部分属性 cached self.node_cache[node_id] cached[last_seen] max(cached[last_seen], attrs[first_seen]) if command_line in attrs and len(attrs[command_line]) len(cached.get(command_line, )): cached[command_line] attrs[command_line] self.G.nodes[node_id].update(cached) def _upsert_edge(self, src_id: str, dst_id: str, edge_type: str, attrs: Dict): 插入或更新边属性 key (src_id, dst_id, edge_type) if key not in self.edge_cache: self.edge_cache[key] attrs.copy() self.G.add_edge(src_id, dst_id, typeedge_type, **attrs) else: cached self.edge_cache[key] cached[weight] max(cached[weight], attrs[weight]) cached[first_seen] min(cached[first_seen], attrs[first_seen]) self.G.edges[src_id, dst_id].update(cached) def _infer_os_from_log(self, log: Dict) - str: if windows in str(log.get(os, )).lower() or win in log.get(src_host, ).lower(): return windows if linux in str(log.get(os, )).lower() or ubuntu in log.get(src_host, ).lower(): return linux return unknown def _is_suspicious_process(self, log: Dict) - bool: # 简单启发式白名单外 非常规路径 高危命令行 name log.get(process_name, ).lower() cmdline log.get(command_line, ).lower() if name in [svchost.exe, lsass.exe, explorer.exe, systemd, bash]: return False if c:\\windows\\system32 in cmdline or /usr/bin in cmdline: return False if any(kw in cmdline for kw in [bypass, executionpolicy, encodedcommand, certutil, bitsadmin]): return True return len(cmdline) 200 # 过长命令行大概率可疑 def _estimate_file_size(self, log: Dict) - int: # 实际中应从日志中提取size字段此处模拟 return len(log.get(file_path, )) * 100 len(log.get(command_line, )) def _estimate_session_duration(self, log: Dict) - int: # 实际中应从NetFlow或PCAP中获取此处按协议估算 if log.get(Protocol) 6: # TCP return 300 # 默认5分钟 return 30 # UDP默认30秒参数说明time_window_minutes控制图的时间滑动窗口——不是丢弃旧数据而是对超过窗口的节点/边打上archivedTrue标签后续查询可过滤。_is_suspicious_process是轻量级初筛不替代ML模型但能快速过滤80%的噪音。关键设计点所有节点ID都包含宿主信息如win-srv-01避免不同主机上的同名进程被错误合并文件ID用SHA256哈希而非路径防止攻击者改名绕过。3. 检测不是找单点异常而是挖攻击路径基于图遍历的APT特征模式匹配有了干净的图结构下一步不是扔进GNN训练而是用确定性图遍历算法精准捕获ATTCK战术链。原因很现实APT攻击路径具有强结构性T1059→T1071→T1037→T1078而GNN在小样本、冷启动、低信噪比场景下极易过拟合或漏报。我们用networkx.algorithms.simple_paths配合自定义路径约束实现毫秒级战术链发现。3.1 定义攻击路径模板把MITRE战术翻译成图查询语言我们不写SQL而是用Python函数描述路径约束。每个模板是一个PathTemplate对象包含start_filter: 起始节点类型属性条件如Process且is_suspiciousTruehop_rules: 每跳的边类型节点类型属性约束如第1跳必须是spawned_by边到另一个Process且子进程command_line含powershellend_condition: 终止节点需满足的条件如NetworkSession且dst_ip在已知C2黑名单中# path_templates.py from typing import List, Dict, Callable, Any, Optional import networkx as nx from datetime import datetime, timedelta class PathTemplate: def __init__(self, name: str, description: str): self.name name self.description description self.start_filter: Optional[Callable[[str, Dict], bool]] None self.hop_rules: List[Dict] [] # [{edge_type: spawned_by, next_node_type: Process, next_filter: ...}] self.end_condition: Optional[Callable[[str, Dict], bool]] None def set_start(self, node_type: str, **filters): 设置起始节点筛选条件 def filter_func(node_id: str, attrs: Dict) - bool: if attrs.get(type) ! node_type: return False for key, expected in filters.items(): if key not in attrs or attrs[key] ! expected: return False return True self.start_filter filter_func return self def add_hop(self, edge_type: str, next_node_type: str, **next_filters): 添加一跳规则 rule { edge_type: edge_type, next_node_type: next_node_type, next_filter: lambda nid, attrs: ( attrs.get(type) next_node_type and all(attrs.get(k) v for k, v in next_filters.items()) ) } self.hop_rules.append(rule) return self def set_end(self, **filters): 设置终止节点条件 def condition_func(node_id: str, attrs: Dict) - bool: for key, expected in filters.items(): if key not in attrs or attrs[key] ! expected: return False return True self.end_condition condition_func return self # 预置常见APT路径模板 TEMPLATES { living_off_land: PathTemplate( nameliving_off_land, description利用系统自带工具执行恶意操作PowerShellCertUtil下载 ).set_start(Process, is_suspiciousTrue).add_hop( edge_typespawned_by, next_node_typeProcess, namepowershell.exe ).add_hop( edge_typeconnects_to, next_node_typeNetworkSession, dst_ip185.199.108.153 # 示例C2 IP ).set_end(typeNetworkSession), persistence_via_schtasks: PathTemplate( namepersistence_via_schtasks, description通过schtasks创建持久化任务 ).set_start(Process, namecmd.exe).add_hop( edge_typespawned_by, next_node_typeProcess, nameschtasks.exe, command_line__containscreate ).add_hop( edge_typewrites_to, next_node_typeFile, path__endswith.bat ).set_end(typeFile), lateral_movement_wmi: PathTemplate( namelateral_movement_wmi, description通过WMI远程执行Win32_Process.Create ).set_start(Process, namewmic.exe).add_hop( edge_typeconnects_to, next_node_typeNetworkSession, dst_ip__not_in[127.0.0.1, localhost] ).add_hop( edge_typespawned_by, next_node_typeProcess, host__not_equal_tolocal_host # 表示远程主机上的进程 ).set_end(typeProcess), }3.2 高效路径搜索剪枝缓存超时保护的遍历引擎nx.all_simple_paths在大图上会指数爆炸。我们必须加三层防护深度限制APT路径 rarely 超过5跳Initial Access → Execution → Persistence → Lateral Movement → C2时间窗口剪枝路径上所有节点时间戳必须在±2小时内攻击者不会等三天再执行下一步节点热度过滤跳过degree 100的超级节点如svchost.exe它连接太多边遍历无意义# path_searcher.py import networkx as nx from typing import List, Tuple, Dict, Any, Generator from datetime import datetime, timedelta import time class PathSearcher: def __init__(self, graph: nx.DiGraph, max_depth: int 5, time_window_seconds: int 7200): self.G graph self.max_depth max_depth self.time_window timedelta(secondstime_window_seconds) # 缓存高频节点的邻居避免重复计算 self.neighbor_cache {} def find_matching_paths(self, template: PathTemplate, limit: int 10) - List[List[str]]: 查找匹配模板的所有路径 if not template.start_filter or not template.end_condition: return [] start_nodes [ node_id for node_id, attrs in self.G.nodes(dataTrue) if template.start_filter(node_id, attrs) ] if not start_nodes: return [] all_paths [] start_time time.time() for start_node in start_nodes: # 对每个起点做DFS带剪枝 paths self._dfs_with_pruning( start_nodestart_node, templatetemplate, current_path[start_node], depth0, visitedset([start_node]) ) all_paths.extend(paths) if len(all_paths) limit: break # 全局超时保护 if time.time() - start_time 30: # 单次搜索最多30秒 break return all_paths[:limit] def _dfs_with_pruning(self, start_node: str, template: PathTemplate, current_path: List[str], depth: int, visited: set) - List[List[str]]: 带剪枝的深度优先搜索 if depth self.max_depth: return [] # 检查当前路径是否满足终止条件 last_node current_path[-1] if last_node in self.G.nodes: attrs self.G.nodes[last_node] if template.end_condition(last_node, attrs): return [current_path.copy()] # 获取当前节点的出边邻居带缓存 if start_node not in self.neighbor_cache: neighbors list(self.G.successors(start_node)) # 剪枝过滤掉degree过高的节点 filtered_neighbors [] for n in neighbors: if self.G.nodes[n].get(type) Process: deg self.G.in_degree(n) self.G.out_degree(n) if deg 50: # 超过50连接的进程跳过 filtered_neighbors.append(n) else: filtered_neighbors.append(n) self.neighbor_cache[start_node] filtered_neighbors else: neighbors self.neighbor_cache[start_node] results [] for neighbor in neighbors: # 检查边类型是否匹配下一跳规则 if depth len(template.hop_rules): edge_data self.G.get_edge_data(start_node, neighbor) if not edge_data or edge_data.get(type) ! template.hop_rules[depth][edge_type]: continue # 检查邻居节点类型和属性 if neighbor not in self.G.nodes: continue node_attrs self.G.nodes[neighbor] next_rule template.hop_rules[depth] if node_attrs.get(type) ! next_rule[next_node_type]: continue if not next_rule[next_filter](neighbor, node_attrs): continue # 时间窗口剪枝所有节点时间戳必须在窗口内 if not self._check_time_window(current_path [neighbor]): continue # 防环路 if neighbor in visited: continue new_path current_path [neighbor] new_visited visited | {neighbor} sub_paths self._dfs_with_pruning( start_nodeneighbor, templatetemplate, current_pathnew_path, depthdepth 1, visitednew_visited ) results.extend(sub_paths) return results def _check_time_window(self, path_nodes: List[str]) - bool: 检查路径上所有节点时间戳是否在时间窗口内 timestamps [] for node_id in path_nodes: if node_id in self.G.nodes: ts self.G.nodes[node_id].get(first_seen) if isinstance(ts, str): try: ts datetime.fromisoformat(ts.replace(Z, 00:00)) timestamps.append(ts) except: pass if len(timestamps) 2: return True return (max(timestamps) - min(timestamps)) self.time_window3.3 实时检测流水线从日志流到告警的端到端闭环检测不是离线分析而是7×24小时流式处理。我们用concurrent.futures.ThreadPoolExecutor实现无锁并发每条日志进入后归一化 → 2. 图更新 → 3. 模板匹配 → 4. 告警生成 → 5. 告警去重相同路径1小时内只报1次# detector.py import threading import queue import time from datetime import datetime, timedelta from concurrent.futures import ThreadPoolExecutor, as_completed from typing import Dict, List, Tuple class APTRuntimeDetector: def __init__(self, graph_builder: AttackGraphBuilder, searcher: PathSearcher): self.graph_builder graph_builder self.searcher searcher self.alert_queue queue.Queue() self.alert_history {} # {path_hash: last_alert_time} self.lock threading.Lock() self.executor ThreadPoolExecutor(max_workers4) def ingest_log(self, raw_log: Dict[str, Any]): 接收原始日志触发全链路检测 # 步骤1归一化 norm_log LogNormalizer().normalize(raw_log) if not norm_log: return # 步骤2更新图 self.graph_builder.add_log_event(norm_log) # 步骤3异步匹配所有模板 futures [] for template_name, template in TEMPLATES.items(): future self.executor.submit( self._match_template_async, template, norm_log ) futures.append((template_name, future)) # 步骤4收集结果并去重告警 for template_name, future in futures: try: paths future.result(timeout10) for path in paths: alert self._generate_alert(template_name, path) if self._should_alert(alert): self.alert_queue.put(alert) except Exception as e: print(fTemplate {template_name} match failed: {e}) def _match_template_async(self, template: PathTemplate, log: Dict) - List[List[str]]: 异步执行路径匹配 return self.searcher.find_matching_paths(template, limit3) def _generate_alert(self, template_name: str, path: List[str]) - Dict: 生成结构化告 p a hrefhttps://download.csdn.net/download/2501_91537435/92339218 stylecolor:#ec7500;font-size:14px; 本文还有配套的精品资源点击获取 /a img altmenu-r.4af5f7ec.gif srchttps://csdnimg.cn/release/wenkucmsfe/public/img/menu-r.4af5f7ec.gif stylewidth:16px;margin-left:4px;vertical-align:text-bottom;cursor:text; /p