import os import sys import json import socket import ssl import struct import argparse import subprocess import urllib.request import warnings from datetime import datetime, timezone, timedelta from typing import Optional, Dict, Any, List # Suppress cryptography / pgpy deprecation notices warnings.filterwarnings("ignore") import pgpy CONFIG_FILE_NAME = "client_config.json" STATE_FILE_NAME = "client_state.json" def is_cert_expiring_soon(cert_path: str, threshold_days: int = 30) -> bool: """Checks if client certificate at cert_path is expiring within threshold_days.""" if not os.path.exists(cert_path): return True try: from cryptography import x509 with open(cert_path, "r", encoding="utf-8") as f: cert = x509.load_pem_x509_certificate(f.read().encode("utf-8")) expiry = getattr(cert, "not_valid_after_utc", None) if expiry is None: expiry = cert.not_valid_after.replace(tzinfo=timezone.utc) now = datetime.now(timezone.utc) return expiry <= (now + timedelta(days=threshold_days)) except Exception: return True def enroll_client_if_needed( hub_url: str, enrollment_secret: str, cert_dir: str, client_id: str, hostname: str, os_type: str = "linux", force_renew: bool = False, threshold_days: int = 30 ): """Bootstraps client enrollment if certificates are missing or expiring soon.""" os.makedirs(cert_dir, exist_ok=True) ca_path = os.path.join(cert_dir, "ca.crt") cert_path = os.path.join(cert_dir, "client.crt") key_path = os.path.join(cert_dir, "client.key") if not force_renew and os.path.exists(ca_path) and os.path.exists(cert_path) and os.path.exists(key_path): if not is_cert_expiring_soon(cert_path, threshold_days=threshold_days): return True print(f"[*] Client certificate at {cert_path} is expiring within {threshold_days} days. Auto-renewing...") action_name = "re-enrolling" if os.path.exists(cert_path) else "enrolling" print(f"[*] Bootstrapping client {action_name} with LOGAR Hub at {hub_url}...") enroll_endpoint = f"{hub_url.rstrip('/')}/api/client/enroll" payload = { "client_id": client_id, "hostname": hostname, "os": os_type, "enrollment_secret": enrollment_secret } req = urllib.request.Request( enroll_endpoint, data=json.dumps(payload).encode("utf-8"), headers={"Content-Type": "application/json"} ) with urllib.request.urlopen(req, timeout=10) as resp: if resp.status != 200: raise RuntimeError(f"Enrollment failed with status code {resp.status}") data = json.loads(resp.read().decode("utf-8")) with open(ca_path, "w", encoding="utf-8") as f: f.write(data["ca_cert"]) with open(cert_path, "w", encoding="utf-8") as f: f.write(data["client_cert"]) with open(key_path, "w", encoding="utf-8") as f: f.write(data["client_key"]) try: os.chmod(key_path, 0o600) except Exception: pass print(f"[+] Client certificates updated successfully in {os.path.abspath(cert_dir)}") return True def get_tls_socket(hub_host: str, hub_port: int, cert_dir: str): """Establishes an mTLS connection with the LOGAR hub using client certificates.""" ca_path = os.path.join(cert_dir, "ca.crt") cert_path = os.path.join(cert_dir, "client.crt") key_path = os.path.join(cert_dir, "client.key") if not (os.path.exists(ca_path) and os.path.exists(cert_path) and os.path.exists(key_path)): raise FileNotFoundError(f"mTLS certificates not found in '{cert_dir}'. Enroll client first.") ctx = ssl.create_default_context(ssl.Purpose.SERVER_AUTH, cafile=ca_path) ctx.load_cert_chain(certfile=cert_path, keyfile=key_path) ctx.minimum_version = ssl.TLSVersion.TLSv1_3 ctx.check_hostname = False raw_sock = socket.create_connection((hub_host, hub_port), timeout=15) return ctx.wrap_socket(raw_sock, server_hostname=hub_host) def get_state_path(config_path: str, custom_state_path: Optional[str] = None) -> str: if custom_state_path: return custom_state_path config_dir = os.path.dirname(os.path.abspath(config_path)) return os.path.join(config_dir, STATE_FILE_NAME) def load_state(state_path: str) -> dict: if os.path.exists(state_path): try: with open(state_path, "r", encoding="utf-8") as f: return json.load(f) except Exception as e: print(f"[!] Warning: Failed to read state file '{state_path}': {e}") return {} return {} def save_state(state_path: str, state: dict): try: temp_path = f"{state_path}.tmp" with open(temp_path, "w", encoding="utf-8") as f: json.dump(state, f, indent=2) os.replace(temp_path, state_path) except Exception as e: print(f"[!] Warning: Could not save client state to '{state_path}': {e}") def commit_state(state: dict, state_path: str): if "new_last_cursor" in state: val = state.pop("new_last_cursor") if val: state["last_cursor"] = val if "new_last_timestamp_us" in state: val = state.pop("new_last_timestamp_us") if val: state["last_timestamp_us"] = val if "new_sent_cursors" in state: state["sent_cursors"] = state.pop("new_sent_cursors") save_state(state_path, state) def load_config(config_path: str = CONFIG_FILE_NAME): if not os.path.exists(config_path): raise FileNotFoundError( f"Client configuration file not found at: {config_path}\n" f"Generate one from the server using: python Server.py --create-client-config --client-out {config_path}" ) with open(config_path, "r", encoding="utf-8") as f: return json.load(f) def get_machine_identifier() -> str: """ Returns the hostname of the machine sending the logs, and appends the network/DNS domain if available. """ # 1. Try fully-qualified domain name (FQDN) fqdn = socket.getfqdn() if fqdn and "." in fqdn and not fqdn.startswith("localhost"): return fqdn hostname = socket.gethostname() # 2. Check /etc/resolv.conf domain or search directive try: if os.path.exists("/etc/resolv.conf"): with open("/etc/resolv.conf", "r", encoding="utf-8") as f: for line in f: parts = line.strip().split() if parts and parts[0] in ["domain", "search"] and len(parts) > 1: domain = parts[1] if domain and not domain.startswith("."): return f"{hostname}.{domain}" except Exception: pass # 3. Try reverse DNS lookup try: host_ip = socket.gethostbyname(hostname) canonical_name = socket.gethostbyaddr(host_ip)[0] if canonical_name and "." in canonical_name and not canonical_name.startswith("localhost"): return canonical_name except Exception: pass return hostname def get_recent_linux_logs(hours: int = 24, state: Optional[dict] = None) -> list: """ Collects info, warnings, and errors from systemd journalctl over the lookback window. Edge Filtering: Retains INFO, WARNING, and ERROR. Strips DEBUG (priority 7) and skips events older than lookback window (default: 24h). State Tracking: Skips events older than lookback window (default 24h) and events that have already been sent in previous runs. """ last_cursor = None last_timestamp_us = 0.0 sent_cursors = set() if state: last_cursor = state.get("last_cursor") try: last_timestamp_us = float(state.get("last_timestamp_us", 0)) except (ValueError, TypeError): last_timestamp_us = 0.0 sent_cursors = set(state.get("sent_cursors", [])) cmd = ["journalctl", "--since", f"{hours} hours ago", "-p", "info", "--output=json"] result = None if last_cursor: cmd_with_cursor = ["journalctl", "--since", f"{hours} hours ago", "--after-cursor", str(last_cursor), "-p", "info", "--output=json"] try: res = subprocess.run(cmd_with_cursor, capture_output=True, text=True, check=False) if res.returncode == 0: result = res except FileNotFoundError: print("[!] journalctl command not found. Ensure this script runs on a systemd-enabled Linux system.") return [] except Exception: pass if result is None: try: result = subprocess.run(cmd, capture_output=True, text=True, check=False) except FileNotFoundError: print("[!] journalctl command not found. Ensure this script runs on a systemd-enabled Linux system.") return [] except Exception as e: print(f"[!] Error running journalctl: {e}") return [] logs = [] machine_id = get_machine_identifier() cutoff_epoch_us = (datetime.now(timezone.utc) - timedelta(hours=hours)).timestamp() * 1_000_000 newest_cursor = None newest_timestamp_us = last_timestamp_us collected_cursors = [] for line in result.stdout.splitlines(): line_str = line.strip() if not line_str: continue try: entry = json.loads(line_str) entry_cursor = entry.get("__CURSOR") entry_ts_us_raw = entry.get("__REALTIME_TIMESTAMP") entry_ts_us = 0.0 if entry_ts_us_raw: try: entry_ts_us = float(entry_ts_us_raw) except (ValueError, TypeError): pass # 1. Skip entries older than lookback window (default: 24h) if entry_ts_us and entry_ts_us < cutoff_epoch_us: continue # 2. Skip already sent events if entry_cursor and (entry_cursor in sent_cursors or entry_cursor == last_cursor): continue if last_timestamp_us > 0 and entry_ts_us > 0 and entry_ts_us < last_timestamp_us: continue # Advance newest tracking for new entries if entry_cursor: newest_cursor = entry_cursor collected_cursors.append(entry_cursor) if entry_ts_us > newest_timestamp_us: newest_timestamp_us = entry_ts_us priority = int(entry.get("PRIORITY", "6")) # Priority 0: Emerg, 1: Alert, 2: Crit, 3: Err (-> ERROR) # Priority 4: Warning, 5: Notice (-> WARNING) # Priority 6: Info (-> INFO) # Priority 7: Debug (skip) if priority > 6: continue if priority <= 3: sev = "ERROR" elif priority in (4, 5): sev = "WARNING" else: sev = "INFO" logs.append({ "server": machine_id, "os_type": "linux", "signature": entry.get("SYSLOG_IDENTIFIER", "unknown"), "severity": sev, "message": entry.get("MESSAGE", "")[:2048] }) except (json.JSONDecodeError, ValueError): continue if state is not None: state["new_last_cursor"] = newest_cursor or last_cursor state["new_last_timestamp_us"] = max(newest_timestamp_us, last_timestamp_us) state["new_sent_cursors"] = (list(sent_cursors) + collected_cursors)[-1000:] state["last_run_timestamp"] = datetime.now(timezone.utc).isoformat() return logs def send_encrypted_logs_over_socket(config: dict, logs: list): """ Streams logs to the LOGAR hub. Uses mutual TLS 1.3 (mTLS) with client certificates if available, or falls back to OpenPGP encrypted envelope over TCP. """ server_host = config["server_host"] server_port = int(config["server_port"]) cert_dir = config.get("cert_dir", "certs") enrollment_secret = config.get("enrollment_secret") machine_id = get_machine_identifier() # Attempt automatic enrollment bootstrap if certs are missing and secret is provided hub_url = None if enrollment_secret: hermes_host = config.get("hermes_host", server_host) hermes_port = config.get("hermes_port", 8443) hub_url = f"http://{hermes_host}:{hermes_port}" try: enroll_client_if_needed(hub_url, enrollment_secret, cert_dir, machine_id, machine_id, os_type="linux") except Exception as e: print(f"[!] Warning: Enrollment bootstrap failed: {e}") ca_path = os.path.join(cert_dir, "ca.crt") cert_path = os.path.join(cert_dir, "client.crt") key_path = os.path.join(cert_dir, "client.key") has_mtls_certs = os.path.exists(ca_path) and os.path.exists(cert_path) and os.path.exists(key_path) if has_mtls_certs: print(f"[*] Connecting to LOGAR server at {server_host}:{server_port} over mTLS (TLS 1.3)...") sock = None try: sock = get_tls_socket(server_host, server_port, cert_dir) except (ssl.SSLError, ssl.CertificateError, ConnectionResetError) as tls_err: if enrollment_secret and hub_url: print(f"[!] TLS handshake error ({tls_err}). Re-enrolling with LOGAR Hub...") try: enroll_client_if_needed(hub_url, enrollment_secret, cert_dir, machine_id, machine_id, os_type="linux", force_renew=True) sock = get_tls_socket(server_host, server_port, cert_dir) except Exception as retry_err: print(f"[!] Re-enrollment or reconnection retry failed: {retry_err}") raise else: raise with sock: payload = { "server": machine_id, "timestamp": datetime.now(timezone.utc).isoformat(), "logs": logs } payload_bytes = json.dumps(payload).encode("utf-8") frame = struct.pack(">I", len(payload_bytes)) + payload_bytes sock.sendall(frame) resp_len_bytes = sock.recv(4) if not resp_len_bytes: raise ConnectionError("Server closed mTLS connection without response.") resp_len = struct.unpack(">I", resp_len_bytes)[0] resp_bytes = bytearray() while len(resp_bytes) < resp_len: chunk = sock.recv(min(4096, resp_len - len(resp_bytes))) if not chunk: break resp_bytes.extend(chunk) response = json.loads(resp_bytes.decode("utf-8")) print(f"[+] Server response: {response}") return response # Fallback to OpenPGP envelope over plain TCP socket auth_token = config.get("auth_token", "") pub_key_armored = config.get("server_public_key") if not pub_key_armored: raise ValueError("No server public key or mTLS certificates available for connection.") expected_fp = config.get("server_fingerprint", "").replace(" ", "").upper() pub_key, _ = pgpy.PGPKey.from_blob(pub_key_armored) actual_fp = str(pub_key.fingerprint).replace(" ", "").upper() if expected_fp and actual_fp != expected_fp: raise ValueError(f"Server fingerprint mismatch! Expected {expected_fp}, but key has {actual_fp}.") payload = { "server": machine_id, "timestamp": datetime.now(timezone.utc).isoformat(), "logs": logs } payload_json = json.dumps(payload) pgp_msg = pgpy.PGPMessage.new(payload_json) encrypted_msg = pub_key.encrypt(pgp_msg) encrypted_armored = str(encrypted_msg) envelope = { "auth_token": auth_token, "timestamp": datetime.now(timezone.utc).isoformat(), "encrypted_payload": encrypted_armored } envelope_bytes = json.dumps(envelope).encode("utf-8") print(f"[*] Connecting to LOGAR server at {server_host}:{server_port} over secure TCP socket...") with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as sock: sock.settimeout(15.0) sock.connect((server_host, server_port)) frame = struct.pack(">I", len(envelope_bytes)) + envelope_bytes sock.sendall(frame) resp_len_bytes = sock.recv(4) if not resp_len_bytes: raise ConnectionError("Server closed connection without response.") resp_len = struct.unpack(">I", resp_len_bytes)[0] resp_bytes = bytearray() while len(resp_bytes) < resp_len: chunk = sock.recv(min(4096, resp_len - len(resp_bytes))) if not chunk: break resp_bytes.extend(chunk) response = json.loads(resp_bytes.decode("utf-8")) print(f"[+] Server response: {response}") return response def main(): parser = argparse.ArgumentParser(description="LOGAR Linux Edge Log Forwarder with State Tracking") parser.add_argument("--config", default=CONFIG_FILE_NAME, help="Path to client_config.json") parser.add_argument("--hours", type=int, default=24, help="Lookback window in hours for journalctl logs (default: 24)") parser.add_argument("--state-file", default=None, help="Path to state tracking file (default: client_state.json next to config)") parser.add_argument("--no-state", action="store_true", help="Disable state tracking and send all events matching lookback window") args = parser.parse_args() try: config = load_config(args.config) except Exception as e: print(f"[!] Configuration error: {e}") sys.exit(1) state_path = get_state_path(args.config, args.state_file) state = None if args.no_state else load_state(state_path) machine_id = get_machine_identifier() print(f"[*] Edge Forwarder Node: {machine_id}") if state and ("last_cursor" in state or "last_timestamp_us" in state): print(f"[*] State tracking active: resuming after previous cursor/timestamp (state file: {state_path})") elif not args.no_state: print(f"[*] State tracking initialized (state file: {state_path})") print(f"[*] Scanning Linux journalctl for unsent entries (last {args.hours} hours)...") candidate_logs = get_recent_linux_logs(hours=args.hours, state=state) print(f"[*] Found {len(candidate_logs)} unsent candidate entries (INFO to ERROR, entries > {args.hours}h and already-sent skipped).") if not candidate_logs: print("[*] No new unsent events to transmit.") if state is not None: commit_state(state, state_path) return try: resp = send_encrypted_logs_over_socket(config, candidate_logs) if state is not None and resp and resp.get("status") == "success": commit_state(state, state_path) print(f"[+] State successfully committed to {state_path}") except Exception as e: print(f"[!] Failed to stream logs to server: {e}") sys.exit(1) if __name__ == "__main__": main()