From 3190351ec67c99dbbc70b16930c4d902a8db88b4 Mon Sep 17 00:00:00 2001 From: Aleksei Fedorov Date: Sat, 4 Apr 2026 11:57:53 +0000 Subject: [PATCH 1/2] declarative route registartion --- kernel_ai/prometheus_setup.py | 37 ++++++++++++++++------- kernel_ai/services/core_observability.py | 17 +++++++---- kernel_ai/services/crypto_security.py | 38 +++++++++++++----------- kernel_ai/services/network.py | 7 +++-- kernel_ai/services/processes.py | 20 +++++++++++++ kernel_ai/services/processes_runtime.py | 34 ++++++++++++--------- kernel_ai/services/system_view.py | 9 ++++-- 7 files changed, 110 insertions(+), 52 deletions(-) diff --git a/kernel_ai/prometheus_setup.py b/kernel_ai/prometheus_setup.py index a2caef0..b166c1a 100644 --- a/kernel_ai/prometheus_setup.py +++ b/kernel_ai/prometheus_setup.py @@ -1,5 +1,6 @@ """Prometheus metrics registration (optional dependency).""" import os +import threading import time from flask import Response, g, jsonify, request @@ -19,6 +20,31 @@ _PROMETHEUS_AVAILABLE = False CONTENT_TYPE_LATEST = "text/plain; version=0.0.4; charset=utf-8" +_REQUEST_COUNT = None +_REQUEST_LATENCY = None +_METRICS_LOCK = threading.Lock() + + +def _get_or_create_http_metrics(): + """Create Prometheus metric objects once per process.""" + global _REQUEST_COUNT, _REQUEST_LATENCY + if _REQUEST_COUNT is not None and _REQUEST_LATENCY is not None: + return _REQUEST_COUNT, _REQUEST_LATENCY + with _METRICS_LOCK: + if _REQUEST_COUNT is None: + _REQUEST_COUNT = Counter( + "http_requests_total", + "Total HTTP requests", + ["method", "endpoint", "status"], + ) + if _REQUEST_LATENCY is None: + _REQUEST_LATENCY = Histogram( + "http_request_duration_seconds", + "HTTP request latency in seconds", + buckets=(0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, float("inf")), + ) + return _REQUEST_COUNT, _REQUEST_LATENCY + def init_prometheus(app): """Register before/after request hooks and /metrics on the given Flask app.""" @@ -32,16 +58,7 @@ def prometheus_metrics_disabled(): return - request_count = Counter( - "http_requests_total", - "Total HTTP requests", - ["method", "endpoint", "status"], - ) - request_latency = Histogram( - "http_request_duration_seconds", - "HTTP request latency in seconds", - buckets=(0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, float("inf")), - ) + request_count, request_latency = _get_or_create_http_metrics() @app.before_request def _prometheus_before_request(): diff --git a/kernel_ai/services/core_observability.py b/kernel_ai/services/core_observability.py index 4cc5c08..30910d2 100644 --- a/kernel_ai/services/core_observability.py +++ b/kernel_ai/services/core_observability.py @@ -2,11 +2,14 @@ from __future__ import annotations +import logging import platform import sys import psutil +logger = logging.getLogger(__name__) + def get_system_info(): """Get system information.""" @@ -87,7 +90,7 @@ def get_kernel_subsystem_status(): with open("/proc/loadavg", "r", encoding="utf-8", errors="ignore") as f: loadavg = f.read().strip().split() running_processes = int(float(loadavg[3].split("/")[0])) - except Exception: + except (OSError, ValueError, IndexError, psutil.Error): running_processes = len(psutil.pids()) if "psutil" in sys.modules else 50 subsystems["process_scheduler"] = {"status": "active", "usage": scheduler_usage, "processes": running_processes} except (IOError, ValueError, KeyError): @@ -104,7 +107,7 @@ def get_kernel_subsystem_status(): parts = line.split() fs_usage = min(100, max(20, int(parts[5]) // 100)) if len(parts) >= 6 else 60 break - except Exception: + except (OSError, ValueError, IndexError): fs_usage = 60 fs_processes = max(5, min(50, mount_count * 2)) subsystems["file_system"] = {"status": "active", "usage": fs_usage, "processes": fs_processes} @@ -136,14 +139,15 @@ def get_kernel_subsystem_status(): tcp_connections = len([line for line in f if line.strip() and not line.startswith("sl")]) network_usage = min(100, max(20, tcp_connections // 10)) network_processes = max(8, min(50, tcp_connections // 5)) - except Exception: + except (OSError, ValueError, IndexError): pass subsystems["network_stack"] = {"status": "active", "usage": network_usage, "processes": network_processes} except (IOError, ValueError): subsystems["network_stack"] = {"status": "active", "usage": 50, "processes": 12} return subsystems - except Exception: + except (OSError, ValueError, KeyError, psutil.Error) as exc: + logger.debug("Failed to build kernel subsystem status, using mock: %s", exc) return get_mock_kernel_subsystems() @@ -167,7 +171,7 @@ def get_process_kernel_map(openai_available=False, openai_module=None): if openai_module is None or not hasattr(openai_module, "api_key") or not openai_module.api_key: return get_mock_process_kernel_map() return get_mock_process_kernel_map() - except Exception: + except (AttributeError, OSError, ValueError): return get_mock_process_kernel_map() @@ -210,5 +214,6 @@ def get_nginx_open_files(): else: files.append({"path": file.path, "type": "other"}) return files[:10] - except Exception: + except (psutil.Error, OSError, ValueError) as exc: + logger.debug("Failed to read nginx open files, using mock: %s", exc) return get_mock_nginx_files() diff --git a/kernel_ai/services/crypto_security.py b/kernel_ai/services/crypto_security.py index 5c11ae9..a3099ad 100644 --- a/kernel_ai/services/crypto_security.py +++ b/kernel_ai/services/crypto_security.py @@ -2,6 +2,7 @@ from __future__ import annotations +import logging import os import random import subprocess @@ -10,6 +11,8 @@ import psutil +logger = logging.getLogger(__name__) + def infer_crypto_protocol(local_port, remote_port, process_name): """Infer protocol/algorithm pair from ports and process hints.""" @@ -77,7 +80,7 @@ def parse_proc_crypto_entries(): try: with open("/proc/crypto", "r", encoding="utf-8", errors="ignore") as f: raw = f.read() - except Exception: + except OSError: return entries blocks = [block.strip() for block in raw.split("\n\n") if block.strip()] @@ -214,7 +217,7 @@ def read_sysctl_int(path, default=0): with open(path, "r", encoding="utf-8", errors="ignore") as f: raw = f.read().strip() return int(raw or default) - except Exception: + except (OSError, ValueError, TypeError): return int(default) @@ -226,7 +229,7 @@ def read_proc_interrupt_total(): parts = line.strip().split() if len(parts) >= 2: return int(parts[1]) - except Exception: + except (OSError, ValueError): return 0 return 0 @@ -240,11 +243,11 @@ def collect_entropy_cloud_status(entropy_prev): write_threshold = read_sysctl_int("/proc/sys/kernel/random/write_wakeup_threshold", 64) try: disk = psutil.disk_io_counters() - except Exception: + except psutil.Error: disk = None try: net = psutil.net_io_counters() - except Exception: + except psutil.Error: net = None intr_total = read_proc_interrupt_total() prev_ts = entropy_prev.get("timestamp") @@ -419,7 +422,8 @@ def collect_crypto_realtime(crypto_prev, entropy_prev=None, callbacks=None): try: connections = psutil.net_connections(kind="inet") - except Exception: + except psutil.Error as exc: + logger.debug("Failed to read net connections for crypto realtime: %s", exc) connections = [] tls_ports = {443, 8443, 9443, 6443} @@ -437,7 +441,7 @@ def collect_crypto_realtime(crypto_prev, entropy_prev=None, callbacks=None): if pid_i: try: process_name = psutil.Process(pid_i).name().lower() - except Exception: + except psutil.Error: process_name = f"pid-{pid_i}" tls_listener_by_port[int(local_port)] = {"pid": pid_i, "process": process_name} tls_listener_names.add(process_name) @@ -475,7 +479,7 @@ def collect_crypto_realtime(crypto_prev, entropy_prev=None, callbacks=None): try: proc = psutil.Process(pid_i) process_name = proc.name() - except Exception: + except psutil.Error: process_name = f"pid-{pid_i}" else: unknown_pid_flows += 1 @@ -509,7 +513,7 @@ def collect_crypto_realtime(crypto_prev, entropy_prev=None, callbacks=None): for proc in psutil.process_iter(attrs=["pid", "name"]): try: name = str(proc.info.get("name", "")).lower() - except Exception: + except (psutil.Error, KeyError, TypeError): continue if any(token in name for token in ["nginx", "sshd", "curl", "openssl", "kube", "vpn", "python"]): protocol, algorithm = infer_crypto_protocol_fn(0, 0, name) @@ -685,7 +689,7 @@ def classify_trust(score): ) except (psutil.NoSuchProcess, psutil.AccessDenied, psutil.ZombieProcess): continue - except Exception: + except (psutil.Error, KeyError, TypeError, ValueError): continue process_rows.sort(key=lambda p: (p["risk_score"], p["mem_percent"], p["threads"]), reverse=True) @@ -718,13 +722,13 @@ def classify_trust(score): try: listen_ports = len([c for c in psutil.net_connections(kind="inet") if str(getattr(c, "status", "") or "") == "LISTEN"]) - except Exception: + except psutil.Error: listen_ports = 0 try: with open("/proc/modules", "r", encoding="utf-8", errors="ignore") as f: loaded_modules = sum(1 for _ in f) - except Exception: + except OSError: loaded_modules = 0 ptrace_processes = sum(1 for p in process_rows if any(tok in p.get("name", "") for tok in ptrace_like)) @@ -740,7 +744,7 @@ def classify_trust(score): timeout=1.8, ).strip() setuid_bins = int(out or 0) - except Exception: + except (subprocess.SubprocessError, OSError, ValueError): setuid_bins = 0 attack_surface = [ @@ -756,7 +760,7 @@ def _read_text(path): try: with open(path, "r", encoding="utf-8", errors="ignore") as f: return str(f.read().strip()) - except Exception: + except OSError: return "" apparmor_raw = _read_text("/sys/module/apparmor/parameters/enabled") @@ -771,7 +775,7 @@ def _read_text(path): lsm_list_raw = _read_text("/sys/kernel/security/lsm") active_lsms = [x.strip() for x in lsm_list_raw.split(",")] if lsm_list_raw else [] stacking_enabled = len([x for x in active_lsms if x in {"selinux", "apparmor", "bpf"}]) > 1 - except Exception: + except (OSError, AttributeError, TypeError): active_lsms = [] stacking_enabled = False @@ -910,7 +914,7 @@ def _read_text(path): seccomp_mode = "filter" else: seccomp_mode = "unknown" - except Exception: + except OSError: pass seccomp_counts[seccomp_mode] = seccomp_counts.get(seccomp_mode, 0) + 1 @@ -945,7 +949,7 @@ def _read_text(path): continue try: cap_eff_val = int(cap_eff_hex, 16) - except Exception: + except ValueError: continue all_caps = [all_capabilities_map.get(bit, f"CAP_{bit}") for bit in range(41) if (cap_eff_val & (1 << bit))] diff --git a/kernel_ai/services/network.py b/kernel_ai/services/network.py index 33ca58a..86b6ee6 100644 --- a/kernel_ai/services/network.py +++ b/kernel_ai/services/network.py @@ -3,6 +3,7 @@ from __future__ import annotations import ipaddress +import logging import re import subprocess import time @@ -13,6 +14,8 @@ from kernel_ai.services.infra_utils import resolve_binary from kernel_ai.state import NETWORK_STACK_PREV, TRACEROUTE_CACHE, TRACEROUTE_CACHE_TTL_SECONDS +logger = logging.getLogger(__name__) + def get_active_connections(): """Get active network connections.""" @@ -48,7 +51,7 @@ def hex_to_ip(hex_str): } ) return connections[:20] - except Exception: + except (OSError, ValueError, IndexError): return get_mock_active_connections() @@ -99,7 +102,7 @@ def _get_default_iface(): for iface in pernic.keys(): if iface != "lo": return iface - except Exception: + except (psutil.Error, OSError): pass return "lo" diff --git a/kernel_ai/services/processes.py b/kernel_ai/services/processes.py index 4c3b917..35ea378 100644 --- a/kernel_ai/services/processes.py +++ b/kernel_ai/services/processes.py @@ -15,6 +15,26 @@ from kernel_ai.services import processes_runtime as _runtime +def get_processes_basic_data() -> list[dict]: + """Collect lightweight process list for generic process table endpoint.""" + processes = [] + for proc in psutil.process_iter(["pid", "name", "status", "memory_info"]): + try: + memory_info = proc.info.get("memory_info") + memory_mb = (memory_info.rss / 1024 / 1024) if memory_info else 0.0 + processes.append( + { + "pid": proc.info["pid"], + "name": proc.info.get("name"), + "status": proc.info.get("status"), + "memory_mb": round(memory_mb, 1), + } + ) + except (psutil.NoSuchProcess, psutil.AccessDenied): + continue + return processes + + def get_proc_matrix_data() -> list[dict]: """Build Matrix view data (processes and resource usage).""" processes = [] diff --git a/kernel_ai/services/processes_runtime.py b/kernel_ai/services/processes_runtime.py index 896fc4c..63ec4c7 100644 --- a/kernel_ai/services/processes_runtime.py +++ b/kernel_ai/services/processes_runtime.py @@ -7,10 +7,13 @@ from __future__ import annotations from datetime import datetime +import logging import os import psutil +logger = logging.getLogger(__name__) + def get_processes_detailed_data() -> list[dict]: """Collect detailed process list for process visualization UI.""" @@ -68,7 +71,8 @@ def get_processes_detailed_data() -> list[dict]: ) except (psutil.NoSuchProcess, psutil.AccessDenied, psutil.ZombieProcess): continue - except Exception: + except (OSError, ValueError, TypeError, KeyError) as exc: + logger.debug("Skipping process in detailed scan due to unexpected data: %s", exc) continue return processes @@ -81,8 +85,8 @@ def _parse_meminfo_kb(): parts = line.split() if len(parts) >= 2 and parts[1].isdigit(): out[parts[0].rstrip(":")] = int(parts[1]) - except Exception: - pass + except OSError as exc: + logger.debug("Failed to read /proc/meminfo: %s", exc) return out @@ -109,7 +113,7 @@ def _build_memory_visual_rows(meminfo_kb, syscall_nodes, vm, swap): if mt <= 1: try: mt = max(1, int(getattr(vm, "total", 0) / 1024)) - except Exception: + except (TypeError, ValueError): mt = 1 seed_base = (mt % 100000) + int(mi.get("Active", 0) or 0) % 50000 @@ -220,7 +224,7 @@ def collect_processes_realtime(): try: with open("/sys/kernel/security/lsm", "r", encoding="utf-8", errors="ignore") as f: lsm_raw = str(f.read().strip()) - except Exception: + except OSError: lsm_raw = "" active_lsms = [x.strip() for x in lsm_raw.split(",") if x.strip()] @@ -228,7 +232,7 @@ def collect_processes_realtime(): try: with open("/proc/sys/kernel/yama/ptrace_scope", "r", encoding="utf-8", errors="ignore") as f: yama_scope = str(f.read().strip()) - except Exception: + except OSError: yama_scope = "" syscall_nodes = [] @@ -246,11 +250,11 @@ def collect_processes_realtime(): threads = int(proc.info.get("num_threads") or 0) try: rss = int(getattr(proc.memory_info(), "rss", 0) or 0) - except Exception: + except (psutil.Error, OSError, TypeError, ValueError): rss = 0 try: fd_count = int(proc.num_fds() or 0) - except Exception: + except (psutil.Error, OSError, TypeError, ValueError): fd_count = 0 seccomp_mode = "unknown" with open(f"/proc/{pid}/status", "r", encoding="utf-8", errors="ignore") as f: @@ -281,7 +285,8 @@ def collect_processes_realtime(): "rss_bytes": rss, } ) - except Exception: + except (psutil.Error, OSError, ValueError, TypeError, KeyError) as exc: + logger.debug("Skipping process in realtime collection: %s", exc) continue syscall_nodes.sort(key=lambda x: x.get("syscall_pressure", 0), reverse=True) syscall_nodes = syscall_nodes[:14] @@ -297,7 +302,7 @@ def collect_processes_realtime(): raddr = getattr(conn, "raddr", None) if raddr and len(raddr) >= 1: remote_ip = str(raddr[0]) - except Exception: + except (AttributeError, TypeError, ValueError): remote_ip = "" status = str(getattr(conn, "status", "") or "").upper() bucket = network_nodes.get(pid) @@ -305,7 +310,7 @@ def collect_processes_realtime(): proc_name = "unknown" try: proc_name = psutil.Process(pid).name() - except Exception: + except psutil.Error: proc_name = "unknown" bucket = {"pid": pid, "name": proc_name, "connections": 0, "remote_ips": set(), "states": {}} network_nodes[pid] = bucket @@ -314,8 +319,8 @@ def collect_processes_realtime(): bucket["remote_ips"].add(remote_ip) if status: bucket["states"][status] = bucket["states"].get(status, 0) + 1 - except Exception: - pass + except (psutil.Error, OSError) as exc: + logger.debug("Failed to sample net connections in processes realtime: %s", exc) network_tracing = [] for _, row in network_nodes.items(): @@ -444,7 +449,8 @@ def _add_edge(src_pid, dst_pid, edge_type, weight): meminfo_kb = _parse_meminfo_kb() strip_rows, mem_summary = _build_memory_visual_rows(meminfo_kb, syscall_nodes, vm, swap) memory_visual = {"layout": "strips", "rows": strip_rows, "summary": mem_summary} - except Exception: + except (psutil.Error, OSError, ValueError, TypeError) as exc: + logger.debug("Failed to build memory visual payload: %s", exc) memory_visual = { "layout": "strips", "rows": [], diff --git a/kernel_ai/services/system_view.py b/kernel_ai/services/system_view.py index f410260..f8f1ef0 100644 --- a/kernel_ai/services/system_view.py +++ b/kernel_ai/services/system_view.py @@ -2,6 +2,7 @@ from __future__ import annotations +import logging import os import re import time @@ -12,12 +13,14 @@ from kernel_ai.collectors import proc_fs as _proc_fs from kernel_ai.state import FILESYSTEM_PREV +logger = logging.getLogger(__name__) + def get_filesystem_blocks(): now = time.time() try: usage = psutil.disk_usage("/") - except Exception: + except psutil.Error: usage = None used_percent = float(usage.percent) if usage else 0.0 @@ -68,8 +71,8 @@ def get_filesystem_blocks(): break else: activity_counts["root"] += 1 - except Exception: - pass + except (psutil.Error, OSError) as exc: + logger.debug("Failed to sample process open files for filesystem heatmap: %s", exc) weighted = [] for z in zone_defs: From b8691bbe525b4a49f4fcd218c63bdeb4f7e05809 Mon Sep 17 00:00:00 2001 From: Aleksei Fedorov Date: Sat, 4 Apr 2026 12:32:24 +0000 Subject: [PATCH 2/2] Layers cleanup --- kernel_ai/services/crypto_security.py | 106 +++++--------------------- kernel_ai/services/network.py | 55 ++++++------- kernel_ai/services/system_view.py | 11 +-- 3 files changed, 54 insertions(+), 118 deletions(-) diff --git a/kernel_ai/services/crypto_security.py b/kernel_ai/services/crypto_security.py index a3099ad..e5b6ebe 100644 --- a/kernel_ai/services/crypto_security.py +++ b/kernel_ai/services/crypto_security.py @@ -11,6 +11,9 @@ import psutil +from kernel_ai.services.crypto import entropy as _entropy_service +from kernel_ai.services.crypto import security_pipeline as _security_pipeline + logger = logging.getLogger(__name__) @@ -213,98 +216,15 @@ def has_token(tokens): def read_sysctl_int(path, default=0): - try: - with open(path, "r", encoding="utf-8", errors="ignore") as f: - raw = f.read().strip() - return int(raw or default) - except (OSError, ValueError, TypeError): - return int(default) + return _entropy_service.read_sysctl_int(path, default=default) def read_proc_interrupt_total(): - try: - with open("/proc/stat", "r", encoding="utf-8", errors="ignore") as f: - for line in f: - if line.startswith("intr "): - parts = line.strip().split() - if len(parts) >= 2: - return int(parts[1]) - except (OSError, ValueError): - return 0 - return 0 + return _entropy_service.read_proc_interrupt_total() def collect_entropy_cloud_status(entropy_prev): - """Collect Linux random subsystem entropy status and source activity.""" - now = time.time() - entropy_bits = read_sysctl_int("/proc/sys/kernel/random/entropy_avail", 0) - pool_size_bits = read_sysctl_int("/proc/sys/kernel/random/poolsize", 256) - read_threshold = read_sysctl_int("/proc/sys/kernel/random/read_wakeup_threshold", 128) - write_threshold = read_sysctl_int("/proc/sys/kernel/random/write_wakeup_threshold", 64) - try: - disk = psutil.disk_io_counters() - except psutil.Error: - disk = None - try: - net = psutil.net_io_counters() - except psutil.Error: - net = None - intr_total = read_proc_interrupt_total() - prev_ts = entropy_prev.get("timestamp") - dt = max(now - prev_ts, 0.001) if prev_ts else None - disk_read_now = int(getattr(disk, "read_bytes", 0) or 0) - disk_write_now = int(getattr(disk, "write_bytes", 0) or 0) - net_sent_now = int(getattr(net, "bytes_sent", 0) or 0) - net_recv_now = int(getattr(net, "bytes_recv", 0) or 0) - if dt: - disk_delta = max((disk_read_now - int(entropy_prev.get("disk_read_bytes") or disk_read_now)) + (disk_write_now - int(entropy_prev.get("disk_write_bytes") or disk_write_now)), 0) - net_delta = max((net_sent_now - int(entropy_prev.get("net_sent_bytes") or net_sent_now)) + (net_recv_now - int(entropy_prev.get("net_recv_bytes") or net_recv_now)), 0) - intr_delta = max(intr_total - int(entropy_prev.get("interrupt_total") or intr_total), 0) - else: - disk_delta = 0 - net_delta = 0 - intr_delta = 0 - entropy_prev["timestamp"] = now - entropy_prev["disk_read_bytes"] = disk_read_now - entropy_prev["disk_write_bytes"] = disk_write_now - entropy_prev["net_sent_bytes"] = net_sent_now - entropy_prev["net_recv_bytes"] = net_recv_now - entropy_prev["interrupt_total"] = intr_total - - def scale_intensity(rate_value, scale): - return int(max(0, min(100, (float(rate_value) / float(scale)) * 100.0))) - - disk_rate = (disk_delta / dt) if dt else 0 - net_rate = (net_delta / dt) if dt else 0 - intr_rate = (intr_delta / dt) if dt else 0 - irq_intensity = scale_intensity(intr_rate, 25000) - disk_intensity = scale_intensity(disk_rate, 80 * 1024 * 1024) - net_intensity = scale_intensity(net_rate, 120 * 1024 * 1024) - hwrng_intensity = 68 if entropy_bits > max(read_threshold, 128) else 34 - sources = [ - {"source": "interrupt timing", "intensity": irq_intensity, "status": "active" if irq_intensity >= 25 else "low"}, - {"source": "disk IO", "intensity": disk_intensity, "status": "active" if disk_intensity >= 18 else "low"}, - {"source": "network timing", "intensity": net_intensity, "status": "active" if net_intensity >= 18 else "low"}, - {"source": "hardware RNG", "intensity": hwrng_intensity, "status": "active" if hwrng_intensity >= 50 else "limited"}, - ] - source_avg = int(sum(s["intensity"] for s in sources) / max(len(sources), 1)) - entropy_pct = max(0.0, min(1.0, float(entropy_bits) / max(float(pool_size_bits), 1.0))) - particle_density = max(16, min(84, int(18 + entropy_pct * 42 + source_avg * 0.35))) - key_birth_rate = round(0.6 + entropy_pct * 9.4 + source_avg * 0.06, 2) - crng_state = "ready" if entropy_bits >= max(read_threshold, 128) else "warming" - random_state = "stable" if entropy_bits >= max(write_threshold, 64) else "refilling" - return { - "entropy_pool_bits": int(entropy_bits), - "entropy_pool_size_bits": int(pool_size_bits), - "crng_state": crng_state, - "random_subsystem_state": random_state, - "particle_density": int(particle_density), - "key_birth_rate_est": float(key_birth_rate), - "sources": sources, - "read_wakeup_threshold": int(read_threshold), - "write_wakeup_threshold": int(write_threshold), - "mode": "live-heuristic", - } + return _entropy_service.collect_entropy_cloud_status(entropy_prev) def collect_algorithm_requesters(items, kernel_clients): @@ -613,7 +533,7 @@ def collect_crypto_realtime(crypto_prev, entropy_prev=None, callbacks=None): } -def collect_security_realtime(security_prev): +def _collect_security_realtime_legacy(security_prev): """ Stage-1 security subsystem telemetry: - Threat decision pipeline @@ -1031,3 +951,15 @@ def _read_text(path): "mode": "live-heuristic-v2", }, } + + +# Transitional wrapper: keep public API stable while implementation lives in +# dedicated submodule. +def collect_security_realtime(security_prev): + return _security_pipeline.collect_security_realtime( + security_prev=security_prev, + psutil_module=psutil, + subprocess_module=subprocess, + random_module=random, + os_module=os, + ) diff --git a/kernel_ai/services/network.py b/kernel_ai/services/network.py index 86b6ee6..7c61c18 100644 --- a/kernel_ai/services/network.py +++ b/kernel_ai/services/network.py @@ -199,7 +199,8 @@ def _get_ss_tcp_metrics(): return {} -def get_network_stack_realtime(): +def get_network_stack_realtime(network_stack_prev=None): + network_stack_prev = NETWORK_STACK_PREV if network_stack_prev is None else network_stack_prev now = time.time() iface = _get_default_iface() pernic = psutil.net_io_counters(pernic=True) @@ -236,7 +237,7 @@ def get_network_stack_realtime(): except OSError: established = 0 - prev_ts = NETWORK_STACK_PREV["timestamp"] + prev_ts = network_stack_prev["timestamp"] dt = max(0.001, now - prev_ts) if prev_ts else 1.0 def rate(curr, prev): @@ -244,10 +245,10 @@ def rate(curr, prev): return 0.0 return max(0.0, (curr - prev) / dt) - retrans_per_sec = rate(retrans_total, NETWORK_STACK_PREV["tcpext_retrans"]) - ip_in_per_sec = rate(ip_in_total, NETWORK_STACK_PREV["ip_in"]) - ip_out_per_sec = rate(ip_out_total, NETWORK_STACK_PREV["ip_out"]) - ip_drop_per_sec = rate(ip_discards_total, NETWORK_STACK_PREV["ip_discards"]) + retrans_per_sec = rate(retrans_total, network_stack_prev["tcpext_retrans"]) + ip_in_per_sec = rate(ip_in_total, network_stack_prev["ip_in"]) + ip_out_per_sec = rate(ip_out_total, network_stack_prev["ip_out"]) + ip_drop_per_sec = rate(ip_discards_total, network_stack_prev["ip_discards"]) rx_per_sec = 0.0 tx_per_sec = 0.0 @@ -255,21 +256,21 @@ def rate(curr, prev): rx_bytes = iface_stats.bytes_recv if iface_stats else 0 tx_bytes = iface_stats.bytes_sent if iface_stats else 0 iface_drops = (iface_stats.dropin + iface_stats.dropout) if iface_stats else 0 - if NETWORK_STACK_PREV["iface_rx"] is not None: - rx_per_sec = max(0.0, (rx_bytes - NETWORK_STACK_PREV["iface_rx"]) / dt) - if NETWORK_STACK_PREV["iface_tx"] is not None: - tx_per_sec = max(0.0, (tx_bytes - NETWORK_STACK_PREV["iface_tx"]) / dt) - if NETWORK_STACK_PREV["iface_drops"] is not None: - iface_drop_per_sec = max(0.0, (iface_drops - NETWORK_STACK_PREV["iface_drops"]) / dt) - - NETWORK_STACK_PREV["timestamp"] = now - NETWORK_STACK_PREV["tcpext_retrans"] = retrans_total - NETWORK_STACK_PREV["ip_in"] = ip_in_total - NETWORK_STACK_PREV["ip_out"] = ip_out_total - NETWORK_STACK_PREV["ip_discards"] = ip_discards_total - NETWORK_STACK_PREV["iface_rx"] = rx_bytes - NETWORK_STACK_PREV["iface_tx"] = tx_bytes - NETWORK_STACK_PREV["iface_drops"] = iface_drops + if network_stack_prev["iface_rx"] is not None: + rx_per_sec = max(0.0, (rx_bytes - network_stack_prev["iface_rx"]) / dt) + if network_stack_prev["iface_tx"] is not None: + tx_per_sec = max(0.0, (tx_bytes - network_stack_prev["iface_tx"]) / dt) + if network_stack_prev["iface_drops"] is not None: + iface_drop_per_sec = max(0.0, (iface_drops - network_stack_prev["iface_drops"]) / dt) + + network_stack_prev["timestamp"] = now + network_stack_prev["tcpext_retrans"] = retrans_total + network_stack_prev["ip_in"] = ip_in_total + network_stack_prev["ip_out"] = ip_out_total + network_stack_prev["ip_discards"] = ip_discards_total + network_stack_prev["iface_rx"] = rx_bytes + network_stack_prev["iface_tx"] = tx_bytes + network_stack_prev["iface_drops"] = iface_drops packets_per_sec = ip_in_per_sec + ip_out_per_sec drop_ratio = (ip_drop_per_sec / packets_per_sec) if packets_per_sec > 0 else 0.0 @@ -375,7 +376,9 @@ def get_route_hint(remote_ip): return {"remote_ip": remote_ip, "tool": "ip-route", "reached": False, "hop_count": 0, "hops": [], "note": "Route hint lookup timed out"} -def get_traceroute_info(remote_ip, max_hops=8): +def get_traceroute_info(remote_ip, max_hops=8, traceroute_cache=None, cache_ttl_seconds=None): + traceroute_cache = TRACEROUTE_CACHE if traceroute_cache is None else traceroute_cache + cache_ttl_seconds = TRACEROUTE_CACHE_TTL_SECONDS if cache_ttl_seconds is None else cache_ttl_seconds try: target_ip = ipaddress.ip_address(remote_ip) if target_ip.is_loopback or target_ip.is_unspecified: @@ -384,8 +387,8 @@ def get_traceroute_info(remote_ip, max_hops=8): raise ValueError("Invalid IP address") now = time.time() - cached = TRACEROUTE_CACHE.get(remote_ip) - if cached and (now - cached["timestamp"]) < TRACEROUTE_CACHE_TTL_SECONDS: + cached = traceroute_cache.get(remote_ip) + if cached and (now - cached["timestamp"]) < cache_ttl_seconds: return cached["data"] traceroute_cmd = resolve_binary("traceroute") @@ -400,7 +403,7 @@ def get_traceroute_info(remote_ip, max_hops=8): tool = "tracepath" else: data = get_route_hint(remote_ip) - TRACEROUTE_CACHE[remote_ip] = {"timestamp": now, "data": data} + traceroute_cache[remote_ip] = {"timestamp": now, "data": data} return data try: @@ -433,5 +436,5 @@ def get_traceroute_info(remote_ip, max_hops=8): reached = any(h.get("target") == remote_ip for h in hops) data = {"remote_ip": remote_ip, "tool": tool, "reached": reached, "hop_count": len(hops), "hops": hops[:max_hops], "note": None} - TRACEROUTE_CACHE[remote_ip] = {"timestamp": now, "data": data} + traceroute_cache[remote_ip] = {"timestamp": now, "data": data} return data diff --git a/kernel_ai/services/system_view.py b/kernel_ai/services/system_view.py index f8f1ef0..0dc0d4d 100644 --- a/kernel_ai/services/system_view.py +++ b/kernel_ai/services/system_view.py @@ -16,7 +16,8 @@ logger = logging.getLogger(__name__) -def get_filesystem_blocks(): +def get_filesystem_blocks(filesystem_prev=None): + filesystem_prev = FILESYSTEM_PREV if filesystem_prev is None else filesystem_prev now = time.time() try: usage = psutil.disk_usage("/") @@ -30,8 +31,8 @@ def get_filesystem_blocks(): io = psutil.disk_io_counters() write_bytes = int(io.write_bytes) if io else 0 - prev_ts = FILESYSTEM_PREV["timestamp"] - prev_write = FILESYSTEM_PREV["write_bytes"] + prev_ts = filesystem_prev["timestamp"] + prev_write = filesystem_prev["write_bytes"] dt = max(0.001, now - prev_ts) if prev_ts else 1.0 write_bps = 0.0 if prev_write is None else max(0.0, (write_bytes - prev_write) / dt) @@ -150,8 +151,8 @@ def get_filesystem_blocks(): ) writing_blocks = sum(1 for b in blocks if b["state"] == "writing") - FILESYSTEM_PREV["timestamp"] = now - FILESYSTEM_PREV["write_bytes"] = write_bytes + filesystem_prev["timestamp"] = now + filesystem_prev["write_bytes"] = write_bytes inode_pressure_global = 0 if zones: