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
2 changes: 2 additions & 0 deletions kernel_ai/contracts/api_contracts.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ class ProcessesRealtimeResponse(TypedDict):
syscalls_interception: list
network_tracing: list
security_hooks: list
semantic_ops: list
neural_graph: dict
memory_visual: dict
meta: dict
Expand Down Expand Up @@ -161,6 +162,7 @@ def validate_processes_realtime_response(payload: Any) -> None:
_expect_key(data, "syscalls_interception", list, "processes_realtime")
_expect_key(data, "network_tracing", list, "processes_realtime")
_expect_key(data, "security_hooks", list, "processes_realtime")
_expect_key(data, "semantic_ops", list, "processes_realtime")
_expect_key(data, "neural_graph", dict, "processes_realtime")
_expect_key(data, "memory_visual", dict, "processes_realtime")
_expect_key(data, "meta", dict, "processes_realtime")
Expand Down
2 changes: 2 additions & 0 deletions kernel_ai/http/api_handlers/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
"""Domain-oriented API handler modules."""

83 changes: 83 additions & 0 deletions kernel_ai/http/api_handlers/kernel.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
"""Kernel/observability API handlers."""

from datetime import datetime

import psutil
from flask import current_app, request
from flask import jsonify

from kernel_ai.http.common import api_json, build_error_payload
from kernel_ai.sentry_helpers import capture_exception
from kernel_ai.services import core_observability as _core_observability_service
from kernel_ai.services import telemetry_orchestration as _telemetry
from kernel_ai.state import get_state_container


def syscalls_realtime():
return api_json(
lambda: {
"timestamp": datetime.now().isoformat(),
"syscalls": _telemetry.get_real_system_calls(),
"cpu_usage": psutil.cpu_percent(interval=1),
"memory_usage": psutil.virtual_memory().percent,
"system_info": _core_observability_service.get_system_info(),
}
)


def kernel_data():
return api_json(
lambda: {
"timestamp": datetime.now().isoformat(),
"syscalls": _telemetry.get_real_system_calls(),
"subsystems": _telemetry.get_kernel_subsystem_status(),
"processes": len(psutil.pids()),
"system_stats": {
"cpu_count": psutil.cpu_count(),
"memory_total": psutil.virtual_memory().total,
"disk_usage": psutil.disk_usage("/").percent,
},
}
)


def process_kernel_map():
return api_json(_telemetry.get_process_kernel_map)


def nginx_files():
return api_json(lambda: {"files": _telemetry.get_nginx_open_files()})


def get_execution_context():
return api_json(
lambda: _telemetry.get_execution_context_data(
exec_context_prev=get_state_container(current_app).exec_context_prev
)
)


def kernel_dna():
return api_json(_telemetry.get_kernel_dna_data)


def sentry_test():
"""Temporary endpoint for manual Sentry verification."""
if not current_app.config.get("SENTRY_TEST_ENDPOINT_ENABLED", False):
return jsonify(build_error_payload("Not found", "not_found")), 404

expected_token = current_app.config.get("SENTRY_TEST_ENDPOINT_TOKEN", "")
incoming_token = (request.headers.get("X-Sentry-Test-Token") or "").strip()
if expected_token and incoming_token != expected_token:
return jsonify(build_error_payload("Forbidden", "forbidden")), 403

mode = (request.args.get("mode") or "capture").strip().lower()
if mode == "raise":
raise RuntimeError("Manual sentry raise test")

try:
raise RuntimeError("Manual sentry capture test")
except RuntimeError as exc:
capture_exception(exc, where="api.sentry_test", extra={"mode": mode})

return jsonify({"ok": True, "mode": mode, "message": "Sentry event captured"})
60 changes: 60 additions & 0 deletions kernel_ai/http/api_handlers/network_system.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
"""Network/system API handlers."""

from flask import current_app, request

from kernel_ai.collectors import proc_fs as _proc_fs
from kernel_ai.http.common import api_json
from kernel_ai.services import devices as _devices_service
from kernel_ai.services import network as _network_service
from kernel_ai.services import system_view as _system_view_service
from kernel_ai.state import get_state_container


def active_connections():
return api_json(lambda: {"connections": _network_service.get_active_connections()})


def traceroute_info():
def _payload():
state = get_state_container(current_app)
remote_ip = request.args.get("ip", "").strip()
if not remote_ip:
raise ValueError("Missing 'ip' query parameter")
return _network_service.get_traceroute_info(
remote_ip,
traceroute_cache=state.traceroute_cache,
cache_ttl_seconds=state.traceroute_cache_ttl_seconds,
)

return api_json(_payload, exception_statuses=[(ValueError, 400)])


def network_stack_realtime():
return api_json(
lambda: _network_service.get_network_stack_realtime(
network_stack_prev=get_state_container(current_app).network_stack_prev
)
)


def devices_realtime():
return api_json(
lambda: _devices_service.get_devices_realtime(
devices_prev=get_state_container(current_app).devices_prev,
read_diskstats_fn=_proc_fs.read_diskstats,
read_interrupt_lines_fn=_proc_fs.read_interrupt_lines,
read_tty_irq_total_fn=_proc_fs.read_tty_irq_total,
)
)


def filesystem_blocks():
return api_json(
lambda: _system_view_service.get_filesystem_blocks(
filesystem_prev=get_state_container(current_app).filesystem_prev
)
)


def isolation_context():
return api_json(_system_view_service.get_isolation_context)
84 changes: 84 additions & 0 deletions kernel_ai/http/api_handlers/processes.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
"""Process-centric API handlers."""

from datetime import datetime
from flask import request

from kernel_ai.http.common import api_json
from kernel_ai.services import process_inspect as _process_inspect_service
from kernel_ai.services import process_timeline as _process_timeline_service
from kernel_ai.services import processes as _processes_service


def get_processes():
return api_json(lambda: {"processes": _processes_service.get_processes_basic_data()})


def get_process_threads(pid):
return api_json(lambda: _process_inspect_service.get_process_threads_info(pid))


def get_process_cpu(pid):
return api_json(lambda: _process_inspect_service.get_process_cpu_info(pid))


def get_process_fds(pid):
return api_json(lambda: _process_inspect_service.get_process_fds_info(pid))


def get_processes_detailed():
return api_json(lambda: {"processes": _processes_service.get_processes_detailed_data()})


def get_ipc_links():
def _payload():
max_pairs = request.args.get("max_pairs", default=120, type=int)
max_nodes = request.args.get("max_nodes", default=24, type=int)
max_pairs = max(20, min(300, max_pairs))
max_nodes = max(8, min(64, max_nodes))
return _process_inspect_service.get_ipc_links_summary(max_pairs=max_pairs, max_nodes=max_nodes)

return api_json(_payload)


def get_proc_matrix():
def _payload():
matrix = _processes_service.get_proc_matrix_data()
return {"matrix": matrix, "timestamp": datetime.now().isoformat()}

return api_json(_payload)


def get_proc_timeline():
return api_json(
lambda: _process_timeline_service.get_proc_timeline_data(
request.args.get("pid", type=int),
window_s=request.args.get("window_s", default=30, type=int),
),
exception_statuses=[(ValueError, 400), (ProcessLookupError, 404)],
)


def get_proc_timeline_branches():
def _payload():
limit = request.args.get("limit", default=6, type=int)
events = request.args.get("events", default=10, type=int)
window_s = request.args.get("window_s", default=30, type=int)
return _process_timeline_service.get_proc_timeline_branches_data(
limit=limit,
events=events,
window_s=window_s,
)

return api_json(_payload)


def processes_realtime():
return api_json(_processes_service.collect_processes_realtime)


def get_proc_graph():
return api_json(_processes_service.get_proc_graph_data, error_extra={"nodes": [], "edges": []})


def get_process_files():
return api_json(lambda: {"curves": [], "timestamp": datetime.now().isoformat()}, error_extra={"curves": []})
46 changes: 46 additions & 0 deletions kernel_ai/http/api_handlers/security_logs.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
"""Security/crypto/frontend-log API handlers."""

from flask import current_app, jsonify, request

from kernel_ai.http.common import api_json
from kernel_ai.services import crypto_security as _crypto_security_service
from kernel_ai.services import frontend_logs as _frontend_logs_service
from kernel_ai.state import get_state_container


def _write_frontend_event(event_payload):
state = get_state_container(current_app)
_frontend_logs_service.write_frontend_event(
event_payload=event_payload,
frontend_log_file=state.frontend_log_file,
write_lock=state.frontend_log_write_lock,
)


def crypto_realtime():
state = get_state_container(current_app)
return api_json(
lambda: _crypto_security_service.collect_crypto_realtime(
crypto_prev=state.crypto_prev,
entropy_prev=state.entropy_prev,
)
)


def security_realtime():
state = get_state_container(current_app)
return api_json(
lambda: _crypto_security_service.collect_security_realtime(
security_prev=state.security_prev
)
)


def ingest_frontend_logs():
if request.method == "OPTIONS":
return ("", 204)
response_payload, status_code = _frontend_logs_service.ingest_frontend_logs_payload(
payload=request.get_json(silent=True),
write_event_fn=_write_frontend_event,
)
return jsonify(response_payload), status_code
74 changes: 74 additions & 0 deletions kernel_ai/logging_helpers.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
"""Shared backend logging helpers."""

from __future__ import annotations

import logging
from collections.abc import Mapping, Sequence

try:
from flask import g, has_request_context
except ImportError: # pragma: no cover
g = None

def has_request_context() -> bool: # type: ignore[override]
return False


_SENSITIVE_KEYS = {"authorization", "cookie", "set-cookie", "password", "passwd", "secret", "token", "api_key"}
_MASK = "***"
_MAX_STR_LEN = 512


def _truncate_string(value: str, max_len: int = _MAX_STR_LEN) -> str:
if len(value) <= max_len:
return value
return value[:max_len] + "...<truncated>"


def sanitize_payload(value):
"""Recursively mask sensitive fields and cap string lengths."""
if isinstance(value, str):
return _truncate_string(value)
if isinstance(value, Mapping):
out = {}
for key, val in value.items():
key_str = str(key)
if key_str.lower() in _SENSITIVE_KEYS:
out[key_str] = _MASK
else:
out[key_str] = sanitize_payload(val)
return out
if isinstance(value, Sequence) and not isinstance(value, (bytes, bytearray, str)):
return [sanitize_payload(x) for x in value]
return value


def log_event(
logger: logging.Logger,
level: int | str,
message: str,
*,
event_dataset: str | None = None,
component: str | None = None,
operation: str | None = None,
event_data: dict | None = None,
) -> None:
"""Emit structured event with sanitized context."""
if isinstance(level, str):
level_int = getattr(logging, level.upper(), logging.INFO)
else:
level_int = int(level)

payload = sanitize_payload(event_data or {})
if event_dataset:
payload.setdefault("event.dataset", event_dataset)
if component:
payload.setdefault("component", component)
if operation:
payload.setdefault("operation", operation)
if has_request_context() and g is not None:
request_id = getattr(g, "request_id", None)
if request_id:
payload.setdefault("request.id", request_id)

logger.log(level_int, message, extra={"event_data": payload})
27 changes: 27 additions & 0 deletions kernel_ai/sentry_helpers.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
"""Lightweight helpers for optional Sentry reporting."""

from __future__ import annotations

from collections.abc import Mapping


def capture_exception(exc: BaseException, *, where: str | None = None, extra: Mapping | None = None) -> None:
"""Send an exception to Sentry if SDK is available/initialized.

This helper never raises to keep exception paths safe.
"""
try:
import sentry_sdk
except Exception:
return

try:
with sentry_sdk.push_scope() as scope:
if where:
scope.set_tag("where", str(where))
if extra:
for key, value in extra.items():
scope.set_extra(str(key), value)
sentry_sdk.capture_exception(exc)
except Exception:
return
Loading
Loading