import os import sys import json import uuid import struct import socket import sqlite3 import argparse import asyncio import secrets import warnings from datetime import datetime, timezone, timedelta from typing import Dict, Any, List, Optional # Suppress cryptography / pgpy deprecation notices for a clean terminal output warnings.filterwarnings("ignore") import pgpy from pgpy.constants import ( PubKeyAlgorithm, KeyFlags, HashAlgorithm, SymmetricKeyAlgorithm, CompressionAlgorithm ) from fastapi import FastAPI, HTTPException import uvicorn CONFIG_FILE_NAME = "server_config.json" DEFAULT_DB_FILE = "logar_state.db" EVALUATION_WINDOW_HOURS = 12 RUN_THRESHOLD = 4 app = FastAPI(title="LOGAR Cloud Ingestion & Hermes Hub", version="2.0.0") # Global context holding server state SERVER_STATE: Dict[str, Any] = {} def generate_server_keypair(server_name: str): """Generates an OpenPGP RSA 2048 key with encryption capability.""" key = pgpy.PGPKey.new(PubKeyAlgorithm.RSAEncryptOrSign, 2048) uid = pgpy.PGPUID.new(server_name) key.add_uid( uid, usage={KeyFlags.EncryptCommunications, KeyFlags.EncryptStorage}, hashes=[HashAlgorithm.SHA256], ciphers=[SymmetricKeyAlgorithm.AES256], compression=[CompressionAlgorithm.Uncompressed] ) private_key_armored = str(key) public_key_armored = str(key.pubkey) fingerprint = str(key.pubkey.fingerprint) return private_key_armored, public_key_armored, fingerprint def load_or_init_config(config_path: str = CONFIG_FILE_NAME) -> Dict[str, Any]: """Loads existing server_config.json or creates a new one on first run.""" if os.path.exists(config_path): print(f"[*] Loading server configuration from: {os.path.abspath(config_path)}") with open(config_path, "r", encoding="utf-8") as f: config = json.load(f) return config print(f"[!] Config '{config_path}' not found. Initializing first-run configuration...") server_name = "LOGAR-Cloud-Hub" private_key, public_key, fingerprint = generate_server_keypair(server_name) auth_token = secrets.token_hex(24) config = { "server_name": server_name, "tcp_host": "0.0.0.0", "tcp_port": 9443, "hermes_host": "0.0.0.0", "hermes_port": 8443, "auth_token": auth_token, "db_path": DEFAULT_DB_FILE, "evaluation_window_hours": EVALUATION_WINDOW_HOURS, "min_persistence_runs": RUN_THRESHOLD, "server_fingerprint": fingerprint, "public_key": public_key, "private_key": private_key } with open(config_path, "w", encoding="utf-8") as f: json.dump(config, f, indent=2) print(f"[+] Successfully generated new server config and OpenPGP keypair.") print(f"[+] Server Encryption Fingerprint: {fingerprint}") print(f"[+] Saved to: {os.path.abspath(config_path)}") return config def create_client_config( server_host: str, server_port: int, output_path: str, config_path: str = CONFIG_FILE_NAME ) -> Dict[str, Any]: """Creates a client configuration file containing the server address, auth token, and encryption-only key/fingerprint.""" server_conf = load_or_init_config(config_path) client_conf = { "server_host": server_host, "server_port": server_port, "server_fingerprint": server_conf["server_fingerprint"], "server_public_key": server_conf["public_key"], "auth_token": server_conf["auth_token"] } out_dir = os.path.dirname(os.path.abspath(output_path)) if out_dir and not os.path.exists(out_dir): os.makedirs(out_dir, exist_ok=True) with open(output_path, "w", encoding="utf-8") as f: json.dump(client_conf, f, indent=2) print(f"[+] Client configuration successfully written to: {os.path.abspath(output_path)}") print(f" - Server Target: {server_host}:{server_port}") print(f" - Encryption Fingerprint: {server_conf['server_fingerprint']}") return client_conf def init_db(db_path: str): """Initializes the SQLite schema for multi-run temporal tracking.""" conn = sqlite3.connect(db_path) conn.execute(""" CREATE TABLE IF NOT EXISTS active_issues ( fingerprint TEXT PRIMARY KEY, site_name TEXT, server TEXT, signature TEXT, severity TEXT, message TEXT, os_type TEXT, first_seen TEXT, last_seen TEXT, run_count INTEGER, status TEXT, last_run_id TEXT ) """) conn.execute(""" CREATE TABLE IF NOT EXISTS ingest_runs ( run_id TEXT PRIMARY KEY, site_name TEXT, server TEXT, timestamp TEXT, log_count INTEGER ) """) conn.commit() conn.close() def process_ingested_logs(payload: Dict[str, Any], db_path: str, window_hours: int, min_runs: int) -> Dict[str, Any]: """ Evaluates candidate issues against the 12-hour evaluation window and 4-run rule. Zero-state clients send raw candidate entries; this engine handles temporal state. """ client_server = payload.get("server", "unknown-host") site_name = payload.get("site_name") or (client_server.split(".", 1)[1] if "." in client_server else "default") logs = payload.get("logs", []) run_id = str(uuid.uuid4()) now = datetime.now(timezone.utc) now_iso = now.isoformat() conn = sqlite3.connect(db_path) cursor = conn.cursor() # Record the batch run cursor.execute( "INSERT INTO ingest_runs (run_id, site_name, server, timestamp, log_count) VALUES (?, ?, ?, ?, ?)", (run_id, site_name, client_server, now_iso, len(logs)) ) processed_count = 0 promoted_to_verified = 0 for log in logs: severity = str(log.get("severity", "WARNING")).upper() # Edge forwarder filter safeguard: retain INFO to ERROR / CRITICAL; strip verbose debug noise if severity in ["DEBUG", "TRACE"]: continue # Errors are always passed immediately; the 4-run rule only concerns warnings is_error = severity in ["ERROR", "CRITICAL", "FATAL"] signature = log.get("signature", "unknown") server = log.get("server", client_server) message = log.get("message", "") os_type = log.get("os_type", "unknown") fp = f"{site_name}:{server}:{signature}" cursor.execute( "SELECT run_count, first_seen, last_seen, status, last_run_id FROM active_issues WHERE fingerprint = ?", (fp,) ) row = cursor.fetchone() if row: run_count, first_seen_str, last_seen_str, current_status, last_run_id = row try: last_seen_dt = datetime.fromisoformat(last_seen_str) except Exception: last_seen_dt = now # 12-hour evaluation window expiry check if (now - last_seen_dt) > timedelta(hours=window_hours): # Window elapsed: reset to new cycle new_runs = 1 new_first_seen = now_iso new_status = "VERIFIED" if is_error else "TRANSIENT" else: # Same run guard: only increment count once per distinct run batch if last_run_id != run_id: new_runs = run_count + 1 else: new_runs = run_count new_first_seen = first_seen_str # 4-run rule applies to warnings; errors are always passed immediately as VERIFIED new_status = "VERIFIED" if (is_error or new_runs >= min_runs) else "TRANSIENT" if new_status == "VERIFIED" and current_status != "VERIFIED": promoted_to_verified += 1 cursor.execute(""" UPDATE active_issues SET run_count = ?, last_seen = ?, first_seen = ?, status = ?, last_run_id = ?, message = ?, severity = ? WHERE fingerprint = ? """, (new_runs, now_iso, new_first_seen, new_status, run_id, message, severity, fp)) else: initial_status = "VERIFIED" if (is_error or 1 >= min_runs) else "TRANSIENT" if initial_status == "VERIFIED": promoted_to_verified += 1 cursor.execute(""" INSERT INTO active_issues (fingerprint, site_name, server, signature, severity, message, os_type, first_seen, last_seen, run_count, status, last_run_id) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, (fp, site_name, server, signature, severity, message, os_type, now_iso, now_iso, 1, initial_status, run_id)) processed_count += 1 conn.commit() conn.close() return { "status": "success", "run_id": run_id, "processed": processed_count, "promoted_verified": promoted_to_verified } async def handle_socket_client(reader: asyncio.StreamReader, writer: asyncio.StreamWriter): """ Authenticated TCP socket handler. Protocol: - 4-byte big-endian prefix: payload length - Payload: JSON with auth_token and encrypted_payload (OpenPGP ASCII armored) - Response: 4-byte length + JSON confirmation """ addr = writer.get_extra_info("peername") try: # Read 4-byte length prefix length_bytes = await reader.readexactly(4) length = struct.unpack(">I", length_bytes)[0] if length <= 0 or length > 10 * 1024 * 1024: # 10MB limit raise ValueError(f"Invalid frame size: {length}") payload_bytes = await reader.readexactly(length) envelope = json.loads(payload_bytes.decode("utf-8")) # Authenticate socket client expected_token = SERVER_STATE["config"]["auth_token"] provided_token = envelope.get("auth_token") if not secrets.compare_digest(str(provided_token), str(expected_token)): err_msg = json.dumps({"status": "error", "message": "Authentication failed"}).encode("utf-8") writer.write(struct.pack(">I", len(err_msg)) + err_msg) await writer.drain() writer.close() await writer.wait_closed() return # Decrypt payload using server's OpenPGP private key encrypted_armored = envelope.get("encrypted_payload", "") pgp_msg = pgpy.PGPMessage.from_blob(encrypted_armored) priv_key = SERVER_STATE["private_key_obj"] decrypted_obj = priv_key.decrypt(pgp_msg) decrypted_json_str = decrypted_obj.message log_payload = json.loads(decrypted_json_str) # Ingest and apply 12h window / 4-run rule res = process_ingested_logs( log_payload, db_path=SERVER_STATE["config"]["db_path"], window_hours=SERVER_STATE["config"]["evaluation_window_hours"], min_runs=SERVER_STATE["config"]["min_persistence_runs"] ) resp_bytes = json.dumps(res).encode("utf-8") writer.write(struct.pack(">I", len(resp_bytes)) + resp_bytes) await writer.drain() except Exception as e: err = json.dumps({"status": "error", "message": str(e)}).encode("utf-8") try: writer.write(struct.pack(">I", len(err)) + err) await writer.drain() except Exception: pass finally: writer.close() try: await writer.wait_closed() except Exception: pass @app.get("/api/hermes/report") def get_hermes_report(): """ Agentic Integration endpoint: Consumed by Hermes to fetch anomalies that have persisted across the 12-hour evaluation window and satisfied the 4-run rule. """ db_path = SERVER_STATE["config"]["db_path"] window_hours = SERVER_STATE["config"]["evaluation_window_hours"] min_runs = SERVER_STATE["config"]["min_persistence_runs"] now = datetime.now(timezone.utc) conn = sqlite3.connect(db_path) cursor = conn.cursor() cursor.execute(""" SELECT fingerprint, site_name, server, signature, severity, message, os_type, first_seen, last_seen, run_count, status FROM active_issues WHERE status = 'VERIFIED' """) rows = cursor.fetchall() conn.close() report = [] for r in rows: last_seen_dt = datetime.fromisoformat(r[8]) # Only return anomalies active within the evaluation window if (now - last_seen_dt) <= timedelta(hours=window_hours): report.append({ "fingerprint": r[0], "site": r[1], "server": r[2], "signature": r[3], "severity": r[4], "message": r[5], "os_type": r[6], "first_seen": r[7], "last_seen": r[8], "consecutive_runs": r[9], "evaluation_window": f"{window_hours}h", "verified": True, "status": r[10] }) return report @app.get("/api/hermes/all") def get_all_issues(): """Diagnostic endpoint to inspect both transient candidate blips and verified anomalies.""" db_path = SERVER_STATE["config"]["db_path"] conn = sqlite3.connect(db_path) cursor = conn.cursor() cursor.execute(""" SELECT fingerprint, site_name, server, signature, severity, message, os_type, first_seen, last_seen, run_count, status FROM active_issues """) rows = cursor.fetchall() conn.close() return [ { "fingerprint": r[0], "site": r[1], "server": r[2], "signature": r[3], "severity": r[4], "message": r[5], "os_type": r[6], "first_seen": r[7], "last_seen": r[8], "run_count": r[9], "status": r[10] } for r in rows ] @app.get("/health") def health_check(): return { "status": "healthy", "server_name": SERVER_STATE["config"]["server_name"], "fingerprint": SERVER_STATE["config"]["server_fingerprint"], "tcp_port": SERVER_STATE["config"]["tcp_port"], "hermes_port": SERVER_STATE["config"]["hermes_port"] } async def run_server(): """Runs the TCP socket listener and the Hermes REST API concurrently.""" config = SERVER_STATE["config"] tcp_host = config["tcp_host"] tcp_port = int(config["tcp_port"]) hermes_host = config["hermes_host"] hermes_port = int(config["hermes_port"]) # Start TCP Socket Server tcp_server = await asyncio.start_server(handle_socket_client, tcp_host, tcp_port) print(f"[*] LOGAR TCP Socket Server listening on {tcp_host}:{tcp_port}") # Start FastAPI / Uvicorn server for Hermes uv_config = uvicorn.Config(app, host=hermes_host, port=hermes_port, log_level="warning") uv_server = uvicorn.Server(uv_config) print(f"[*] Hermes Reporting API available at http://{hermes_host}:{hermes_port}/api/hermes/report") await asyncio.gather( tcp_server.serve_forever(), uv_server.serve() ) def main(): parser = argparse.ArgumentParser(description="LOGAR Cloud Hub & TCP Socket Ingestion Server") parser.add_argument("--config", default=CONFIG_FILE_NAME, help="Path to server_config.json") parser.add_argument("--create-client-config", action="store_true", help="Generate a client config with encryption-only fingerprint and server address") parser.add_argument("--client-out", default="client_config.json", help="Output file path for generated client config") parser.add_argument("--server-host", default="127.0.0.1", help="Server address to embed in client config") parser.add_argument("--server-port", type=int, default=None, help="TCP port to embed in client config") args = parser.parse_args() config = load_or_init_config(args.config) init_db(config["db_path"]) # Load OpenPGP private key into memory priv_key_obj, _ = pgpy.PGPKey.from_blob(config["private_key"]) SERVER_STATE["config"] = config SERVER_STATE["private_key_obj"] = priv_key_obj if args.create_client_config: port = args.server_port or config["tcp_port"] create_client_config( server_host=args.server_host, server_port=port, output_path=args.client_out, config_path=args.config ) sys.exit(0) print("=" * 60) print(f" LOGAR Server Hub: {config['server_name']}") print(f" Encryption Fingerprint: {config['server_fingerprint']}") print(f" Evaluation Window: {config['evaluation_window_hours']} hours | 4-Run Rule: Warnings | Immediate Pass: Errors") print("=" * 60) try: asyncio.run(run_server()) except KeyboardInterrupt: print("\n[!] Server shutting down.") if __name__ == "__main__": main()