Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 27 additions & 10 deletions kernel_ai/prometheus_setup.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
"""Prometheus metrics registration (optional dependency)."""
import os
import threading
import time

from flask import Response, g, jsonify, request
Expand All @@ -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."""
Expand All @@ -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():
Expand Down
17 changes: 11 additions & 6 deletions kernel_ai/services/core_observability.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
Expand Down Expand Up @@ -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):
Expand All @@ -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}
Expand Down Expand Up @@ -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()


Expand All @@ -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()


Expand Down Expand Up @@ -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()
136 changes: 36 additions & 100 deletions kernel_ai/services/crypto_security.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

from __future__ import annotations

import logging
import os
import random
import subprocess
Expand All @@ -10,6 +11,11 @@

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__)


def infer_crypto_protocol(local_port, remote_port, process_name):
"""Infer protocol/algorithm pair from ports and process hints."""
Expand Down Expand Up @@ -77,7 +83,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()]
Expand Down Expand Up @@ -210,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 Exception:
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 Exception:
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 Exception:
disk = None
try:
net = psutil.net_io_counters()
except Exception:
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):
Expand Down Expand Up @@ -419,7 +342,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}
Expand All @@ -437,7 +361,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)
Expand Down Expand Up @@ -475,7 +399,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
Expand Down Expand Up @@ -509,7 +433,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)
Expand Down Expand Up @@ -609,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
Expand Down Expand Up @@ -685,7 +609,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)
Expand Down Expand Up @@ -718,13 +642,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))
Expand All @@ -740,7 +664,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 = [
Expand All @@ -756,7 +680,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")
Expand All @@ -771,7 +695,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

Expand Down Expand Up @@ -910,7 +834,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
Expand Down Expand Up @@ -945,7 +869,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))]
Expand Down Expand Up @@ -1027,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,
)
Loading
Loading