#!/usr/bin/env python3
"""
SWORD Platform — Standalone Linux Sensor Agent (Zero-Pip)
Diseñado para correr en servidores Linux (Debian, Ubuntu, CentOS, RHEL) junto a Cowrie u OpenCanary.
No requiere dependencias de terceros (utiliza exclusivamente la biblioteca estándar de Python 3).
"""

import os
import sys
import time
import json
import logging
import argparse
import threading
import urllib.request
import urllib.error
from datetime import datetime, timezone
from typing import Dict, Any, List, Optional

# Configuración de Logging
logging.basicConfig(
    level=logging.INFO,
    format="[%(asctime)s] [%(levelname)s] [SWORD-Agent] %(message)s",
    datefmt="%Y-%m-%d %H:%M:%S"
)
logger = logging.getLogger("sword_agent")


# --- Adaptadores Integrados de Telemetría ---

def parse_cowrie_event(raw: Dict[str, Any]) -> Optional[Dict[str, Any]]:
    event_id = raw.get("eventid", "")
    timestamp_str = raw.get("timestamp")
    if timestamp_str:
        try:
            dt = datetime.fromisoformat(timestamp_str.replace("Z", "+00:00"))
        except ValueError:
            dt = datetime.now(timezone.utc)
    else:
        dt = datetime.now(timezone.utc)

    source_ip = raw.get("src_ip", "0.0.0.0")
    source_port = raw.get("src_port")
    dst_port = raw.get("dst_port", 22)
    service = "ssh"
    if dst_port == 23 or "telnet" in event_id:
        service = "telnet"

    event_type = "connection"
    username = raw.get("username")
    password = raw.get("password")
    command = None
    payload = None
    tags = ["cowrie", service]

    if event_id == "cowrie.login.failed":
        event_type = "login_attempt"
        tags.append("login_failed")
    elif event_id == "cowrie.login.success":
        event_type = "login_success_fake"
        tags.append("interactive_access")
    elif event_id == "cowrie.command.input":
        event_type = "command_executed"
        command = raw.get("input")
        tags.append("command")
    elif event_id == "cowrie.session.file_download":
        event_type = "file_download_attempt"
        payload = raw.get("url") or raw.get("shasum")
        tags.append("payload")
    elif event_id == "cowrie.client.version":
        event_type = "connection"
        payload = raw.get("version")

    return {
        "timestamp": dt.isoformat(),
        "service": service,
        "event_type": event_type,
        "source_ip": source_ip,
        "source_port": source_port,
        "destination_port": dst_port,
        "username": username,
        "password": password,
        "command": command,
        "payload": payload,
        "tags": tags,
        "raw_event": raw
    }


def parse_opencanary_event(raw: Dict[str, Any]) -> Optional[Dict[str, Any]]:
    logdata = raw.get("logdata", {})
    if not isinstance(logdata, dict):
        logdata = {}

    timestamp_str = raw.get("utc_time")
    if timestamp_str:
        try:
            dt = datetime.fromisoformat(timestamp_str.replace(" ", "T") + "+00:00")
        except ValueError:
            dt = datetime.now(timezone.utc)
    else:
        dt = datetime.now(timezone.utc)

    source_ip = raw.get("src_host", "0.0.0.0")
    source_port = raw.get("src_port")
    dst_port = raw.get("dst_port", 0)
    logtype = raw.get("logtype")

    service = "unknown"
    event_type = "connection"
    username = None
    password = None
    command = None
    payload = None
    tags = ["opencanary"]

    if logtype in (2000, 2001):
        service = "http"
        event_type = "web_probe"
        payload = logdata.get("PATH")
        username = logdata.get("USERNAME")
        password = logdata.get("PASSWORD")
        tags.append("http")
    elif logtype == 1001:
        service = "ftp"
        event_type = "login_attempt"
        username = logdata.get("USERNAME")
        password = logdata.get("PASSWORD")
        tags.append("ftp")
    elif logtype == 5001:
        service = "smb"
        event_type = "smb_probe"
        username = logdata.get("USER")
        payload = logdata.get("SHARE")
        tags.append("smb")
    elif logtype == 4001:
        service = "portscan"
        event_type = "portscan_detected"
        tags.append("portscan")
    elif logtype in (7001, 8001, 10001, 14001):
        service = "database"
        event_type = "db_probe"
        username = logdata.get("USERNAME")
        password = logdata.get("PASSWORD")
        tags.append("db")
    else:
        service = f"logtype_{logtype}"
        event_type = "service_probe"

    return {
        "timestamp": dt.isoformat(),
        "service": service,
        "event_type": event_type,
        "source_ip": source_ip,
        "source_port": source_port,
        "destination_port": dst_port,
        "username": username,
        "password": password,
        "command": command,
        "payload": payload,
        "tags": tags,
        "raw_event": raw
    }


# --- Cliente HTTP Nativo (urllib) ---

def http_post_json(url: str, token: str, data: Dict[str, Any], timeout: int = 10) -> tuple[bool, Dict[str, Any], Optional[str]]:
    headers = {
        "Authorization": f"Bearer {token}",
        "Content-Type": "application/json",
        "User-Agent": "SWORD-Linux-Agent/2.0"
    }
    json_bytes = json.dumps(data).encode("utf-8")
    req = urllib.request.Request(url, data=json_bytes, headers=headers, method="POST")
    try:
        with urllib.request.urlopen(req, timeout=timeout) as response:
            body = response.read().decode("utf-8")
            res_json = json.loads(body) if body else {}
            return True, res_json, None
    except urllib.error.HTTPError as e:
        err_msg = e.read().decode("utf-8", errors="replace")
        return False, {}, f"HTTP {e.code}: {err_msg[:200]}"
    except Exception as e:
        return False, {}, f"Conexión fallida: {str(e)}"


# --- Agente Principal ---

class StandaloneLinuxAgent:
    def __init__(
        self,
        collector_url: str,
        token: str,
        adapter_type: str = "cowrie",
        log_file: Optional[str] = None,
        buffer_file: str = "/var/log/sword-agent/buffer.json",
        batch_size: int = 15,
        heartbeat_interval: int = 60
    ):
        self.collector_url = collector_url.rstrip("/")
        self.token = token
        self.adapter_type = adapter_type
        self.log_file = log_file
        self.buffer_file = buffer_file
        self.batch_size = batch_size
        self.heartbeat_interval = heartbeat_interval
        self.buffer: List[Dict[str, Any]] = []
        self.running = False
        self._ensure_buffer_dir()
        self._load_buffer()

    def _ensure_buffer_dir(self):
        try:
            d = os.path.dirname(self.buffer_file)
            if d and not os.path.exists(d):
                os.makedirs(d, exist_ok=True)
        except Exception:
            self.buffer_file = os.path.expanduser("~/.sword_buffer.json")

    def _load_buffer(self):
        if os.path.exists(self.buffer_file):
            try:
                with open(self.buffer_file, "r", encoding="utf-8") as f:
                    content = f.read().strip()
                    if content:
                        self.buffer = json.loads(content)
                        logger.info(f"Búfer local cargado con {len(self.buffer)} eventos pendientes.")
            except Exception as e:
                logger.warning(f"Error cargando búfer: {e}")
                self.buffer = []

    def _save_buffer(self):
        try:
            with open(self.buffer_file, "w", encoding="utf-8") as f:
                json.dump(self.buffer, f)
        except Exception as e:
            logger.error(f"Error guardando búfer local: {e}")

    def send_heartbeat(self) -> bool:
        url = f"{self.collector_url}/api/ingest/heartbeat"
        payload = {
            "agent_version": "2.0.0-standalone",
            "metrics": {
                "buffer_len": len(self.buffer),
                "adapter": self.adapter_type,
                "pid": os.getpid()
            }
        }
        ok, res, err = http_post_json(url, self.token, payload, timeout=6)
        if ok:
            logger.debug(f"Heartbeat exitoso. Estado del sensor: {res.get('status')}")
            return True
        else:
            logger.warning(f"Falla de heartbeat: {err}")
            return False

    def flush_events(self) -> bool:
        if not self.buffer:
            return True

        batch = self.buffer[:self.batch_size]
        url = f"{self.collector_url}/api/ingest/events"
        payload = {"events": batch}

        ok, res, err = http_post_json(url, self.token, payload, timeout=10)
        if ok and res.get("success"):
            logger.info(f"Lote de {len(batch)} eventos transmitido con éxito al Collector.")
            self.buffer = self.buffer[len(batch):]
            self._save_buffer()
            return True
        else:
            logger.warning(f"Error transmitiendo lote ({err}). Eventos retenidos en búfer local ({len(self.buffer)} total).")
            return False

    def enqueue_raw_line(self, line: str):
        line = line.strip()
        if not line:
            return
        try:
            raw = json.loads(line)
        except Exception:
            return

        if self.adapter_type == "cowrie":
            ev = parse_cowrie_event(raw)
        elif self.adapter_type == "opencanary":
            ev = parse_opencanary_event(raw)
        else:
            ev = {
                "timestamp": datetime.now(timezone.utc).isoformat(),
                "service": "custom",
                "event_type": "probe",
                "source_ip": raw.get("src_ip", "0.0.0.0"),
                "tags": ["custom"],
                "raw_event": raw
            }

        if ev:
            self.buffer.append(ev)
            self._save_buffer()
            if len(self.buffer) >= self.batch_size:
                self.flush_events()

    def run_heartbeat_loop(self):
        while self.running:
            self.send_heartbeat()
            # Intento periódico de vaciar búfer
            if self.buffer:
                self.flush_events()
            for _ in range(self.heartbeat_interval):
                if not self.running:
                    break
                time.sleep(1)

    def follow_file(self):
        """Monitoreo en tiempo real estilo tail -F sobre el archivo de log del señuelo."""
        if not self.log_file:
            logger.error("No se especificó archivo de log para observar.")
            return

        logger.info(f"Iniciando observador en tiempo real sobre: {self.log_file}")
        while self.running:
            if not os.path.exists(self.log_file):
                logger.warning(f"Archivo {self.log_file} no encontrado. Esperando a que el honeypot lo cree...")
                time.sleep(5)
                continue

            try:
                with open(self.log_file, "r", encoding="utf-8", errors="replace") as f:
                    # Mover puntero al final del archivo para procesar solo nuevos eventos
                    f.seek(0, os.SEEK_END)
                    logger.info(f"Observador anclado al final de {self.log_file}. Esperando actividad hostil...")

                    while self.running:
                        line = f.readline()
                        if line:
                            self.enqueue_raw_line(line)
                        else:
                            # Comprobar si el archivo fue rotado o truncado
                            try:
                                current_size = os.path.getsize(self.log_file)
                                if current_size < f.tell():
                                    logger.info("Rotación detectada en archivo de log. Reabriendo...")
                                    break
                            except OSError:
                                break
                            time.sleep(0.5)
            except Exception as e:
                logger.error(f"Error leyendo archivo de log: {e}")
                time.sleep(3)

    def start(self):
        self.running = True
        logger.info(f"=== SWORD Standalone Linux Agent Iniciado ===")
        logger.info(f"• Collector URL: {self.collector_url}")
        logger.info(f"• Adaptador: {self.adapter_type}")
        logger.info(f"• Log File: {self.log_file}")

        # Enviar primer heartbeat de arranque
        self.send_heartbeat()

        # Iniciar loop de heartbeat en hilo dedicado
        hb_thread = threading.Thread(target=self.run_heartbeat_loop, daemon=True)
        hb_thread.start()

        # Iniciar observador de log en hilo principal
        try:
            self.follow_file()
        except KeyboardInterrupt:
            logger.info("Detención solicitada por el usuario.")
        finally:
            self.running = False
            self.flush_events()
            logger.info("Agente detenido limpiamente.")


def main():
    parser = argparse.ArgumentParser(description="SWORD Standalone Linux Sensor Agent")
    parser.add_argument("--collector", default=os.getenv("SWORD_COLLECTOR_URL", "http://localhost:8000"), help="URL base del Collector SWORD")
    parser.add_argument("--token", default=os.getenv("SWORD_SENSOR_TOKEN", ""), help="Token de autenticación del sensor")
    parser.add_argument("--adapter", default=os.getenv("SWORD_ADAPTER_TYPE", "cowrie"), choices=["cowrie", "opencanary", "custom"], help="Tipo de adaptador")
    parser.add_argument("--log-file", default=os.getenv("SWORD_LOG_FILE"), help="Ruta al log del honeypot (ej. /var/log/cowrie/cowrie.json)")
    parser.add_argument("--buffer-file", default=os.getenv("SWORD_BUFFER_FILE", "/var/log/sword-agent/buffer.json"), help="Ruta al búfer local")
    parser.add_argument("--batch-size", type=int, default=int(os.getenv("SWORD_BATCH_SIZE", "15")), help="Tamaño de lote")
    parser.add_argument("--heartbeat-interval", type=int, default=int(os.getenv("SWORD_HEARTBEAT_INTERVAL", "60")), help="Segundos entre heartbeats")

    args = parser.parse_args()

    if not args.token:
        logger.error("Error: Se requiere el token del sensor (--token o variable SWORD_SENSOR_TOKEN)")
        sys.exit(1)

    # Autodetección de log_file si no se especificó
    if not args.log_file:
        if args.adapter == "cowrie":
            candidates = [
                "/var/log/cowrie/cowrie.json",
                "/home/cowrie/cowrie/var/log/cowrie/cowrie.json",
                os.path.expanduser("~/cowrie/var/log/cowrie/cowrie.json")
            ]
            for c in candidates:
                if os.path.exists(c):
                    args.log_file = c
                    break
            if not args.log_file:
                args.log_file = "/var/log/cowrie/cowrie.json"
        elif args.adapter == "opencanary":
            candidates = [
                "/var/tmp/opencanary.log",
                "/var/log/opencanary/opencanary.log",
                "/var/log/opencanary.log"
            ]
            for c in candidates:
                if os.path.exists(c):
                    args.log_file = c
                    break
            if not args.log_file:
                args.log_file = "/var/tmp/opencanary.log"

    agent = StandaloneLinuxAgent(
        collector_url=args.collector,
        token=args.token,
        adapter_type=args.adapter,
        log_file=args.log_file,
        buffer_file=args.buffer_file,
        batch_size=args.batch_size,
        heartbeat_interval=args.heartbeat_interval
    )
    agent.start()


if __name__ == "__main__":
    main()
