310 lines
12 KiB
Python
310 lines
12 KiB
Python
import os
|
|
import sys
|
|
import json
|
|
import socket
|
|
import struct
|
|
import argparse
|
|
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
|
|
|
|
try:
|
|
import win32evtlog
|
|
except ImportError:
|
|
win32evtlog = None
|
|
|
|
CONFIG_FILE_NAME = "client_config.json"
|
|
STATE_FILE_NAME = "client_state.json"
|
|
|
|
|
|
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_record_number" in state:
|
|
val = state.pop("new_last_record_number")
|
|
if val:
|
|
state["last_record_number"] = val
|
|
if "new_sent_record_ids" in state:
|
|
state["sent_record_ids"] = state.pop("new_sent_record_ids")
|
|
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.
|
|
"""
|
|
fqdn = socket.getfqdn()
|
|
if fqdn and "." in fqdn and not fqdn.startswith("localhost"):
|
|
return fqdn
|
|
|
|
hostname = socket.gethostname()
|
|
user_dns_domain = os.environ.get("USERDNSDOMAIN")
|
|
if user_dns_domain and user_dns_domain.lower() != hostname.lower():
|
|
return f"{hostname}.{user_dns_domain.lower()}"
|
|
|
|
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_windows_logs(hours: int = 24, state: Optional[dict] = None) -> list:
|
|
"""
|
|
Scans the Windows Application Event Log backwards for events within the window.
|
|
Edge Filtering: Retains INFO, WARNING, and ERROR. Drops Audit and Debug noise.
|
|
State Tracking: Skips events older than lookback window (default 24h) and events
|
|
that have already been sent in previous runs.
|
|
"""
|
|
if win32evtlog is None:
|
|
print("[!] pywin32 is not installed or not running on Windows. Returning mock/empty candidate list.")
|
|
return []
|
|
|
|
server = "localhost"
|
|
log_type = "Application"
|
|
flags = win32evtlog.EVENTLOG_BACKWARDS_READ | win32evtlog.EVENTLOG_SEQUENTIAL_READ
|
|
|
|
try:
|
|
hand = win32evtlog.OpenEventLog(server, log_type)
|
|
except Exception as e:
|
|
print(f"[!] Error opening Windows event log: {e}")
|
|
return []
|
|
|
|
logs = []
|
|
cutoff_time = datetime.now() - timedelta(hours=hours)
|
|
machine_id = get_machine_identifier()
|
|
|
|
last_record_number = 0
|
|
sent_record_ids = set()
|
|
if state:
|
|
last_record_number = int(state.get("last_record_number", 0))
|
|
sent_record_ids = set(state.get("sent_record_ids", []))
|
|
|
|
# Windows Event Log EventTypes:
|
|
# 1: EVENTLOG_ERROR_TYPE -> ERROR
|
|
# 2: EVENTLOG_WARNING_TYPE -> WARNING
|
|
# 4: EVENTLOG_INFORMATION_TYPE -> INFO
|
|
# Excludes: 8 (Audit Success), 16 (Audit Failure), and other verbose noise
|
|
sev_map = {
|
|
1: "ERROR",
|
|
2: "WARNING",
|
|
4: "INFO"
|
|
}
|
|
|
|
newest_record_number = 0
|
|
collected_record_ids = []
|
|
|
|
while True:
|
|
events = win32evtlog.ReadEventLog(hand, flags, 0)
|
|
if not events:
|
|
break
|
|
|
|
for event in events:
|
|
rec_num = int(event.RecordNumber)
|
|
if newest_record_number == 0:
|
|
newest_record_number = rec_num
|
|
|
|
# 1. Skip entries older than lookback window (default: 24h)
|
|
if event.TimeGenerated < cutoff_time:
|
|
break
|
|
|
|
# 2. Skip already sent events if we've reached records <= last_record_number
|
|
# (unless the log was cleared and numbers wrapped, i.e. newest_record_number < last_record_number)
|
|
if last_record_number > 0 and newest_record_number >= last_record_number:
|
|
if rec_num <= last_record_number:
|
|
break
|
|
|
|
rec_id = f"{rec_num}:{event.TimeGenerated.isoformat()}"
|
|
if rec_id in sent_record_ids:
|
|
continue
|
|
|
|
# Filter: upload everything from INFO to ERROR only
|
|
if event.EventType in sev_map:
|
|
msg = " ".join(event.StringInserts) if event.StringInserts else "Event Log Entry"
|
|
logs.append({
|
|
"server": machine_id,
|
|
"os_type": "windows",
|
|
"signature": event.SourceName or "Windows-Event",
|
|
"severity": sev_map[event.EventType],
|
|
"message": msg[:2048] # Limit message length
|
|
})
|
|
collected_record_ids.append(rec_id)
|
|
|
|
if events[-1].TimeGenerated < cutoff_time:
|
|
break
|
|
if last_record_number > 0 and newest_record_number >= last_record_number and events[-1].RecordNumber <= last_record_number:
|
|
break
|
|
|
|
win32evtlog.CloseEventLog(hand)
|
|
|
|
if state is not None:
|
|
target_rec = max(newest_record_number, last_record_number)
|
|
state["new_last_record_number"] = target_rec
|
|
state["new_sent_record_ids"] = (list(sent_record_ids) + collected_record_ids)[-1000:]
|
|
state["last_run_timestamp"] = datetime.now(timezone.utc).isoformat()
|
|
|
|
return logs
|
|
|
|
|
|
def send_encrypted_logs_over_socket(config: dict, logs: list):
|
|
"""
|
|
Encrypts the payload using the server's OpenPGP public key and streams
|
|
over an authenticated TCP socket. Zero local state is maintained on the client.
|
|
"""
|
|
server_host = config["server_host"]
|
|
server_port = int(config["server_port"])
|
|
auth_token = config["auth_token"]
|
|
pub_key_armored = config["server_public_key"]
|
|
expected_fp = config.get("server_fingerprint", "").replace(" ", "").upper()
|
|
|
|
# Load and verify server public key
|
|
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}."
|
|
)
|
|
|
|
machine_id = get_machine_identifier()
|
|
|
|
# Prepare zero-state candidate batch
|
|
payload = {
|
|
"server": machine_id,
|
|
"timestamp": datetime.now(timezone.utc).isoformat(),
|
|
"logs": logs
|
|
}
|
|
payload_json = json.dumps(payload)
|
|
|
|
# Encrypt payload with server's encryption-only key
|
|
pgp_msg = pgpy.PGPMessage.new(payload_json)
|
|
encrypted_msg = pub_key.encrypt(pgp_msg)
|
|
encrypted_armored = str(encrypted_msg)
|
|
|
|
# Envelope with socket authentication header
|
|
envelope = {
|
|
"auth_token": auth_token,
|
|
"timestamp": datetime.now(timezone.utc).isoformat(),
|
|
"encrypted_payload": encrypted_armored
|
|
}
|
|
envelope_bytes = json.dumps(envelope).encode("utf-8")
|
|
|
|
# Connect over TCP socket and transmit with 4-byte length prefix framing
|
|
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))
|
|
|
|
# Send frame: length (4 bytes big-endian) + envelope
|
|
frame = struct.pack(">I", len(envelope_bytes)) + envelope_bytes
|
|
sock.sendall(frame)
|
|
|
|
# Receive response length
|
|
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 Windows 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 event 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_record_number" in state:
|
|
print(f"[*] State tracking active: resuming from record #{state['last_record_number']} (state file: {state_path})")
|
|
elif not args.no_state:
|
|
print(f"[*] State tracking initialized (state file: {state_path})")
|
|
|
|
print(f"[*] Scanning Windows Application event log for unsent entries (last {args.hours} hours)...")
|
|
candidate_logs = get_recent_windows_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()
|