import json import os import re import threading import csv from io import StringIO from collections import Counter from contextlib import asynccontextmanager from datetime import datetime from pathlib import Path from fastapi import FastAPI, Query, Request from fastapi.responses import HTMLResponse, JSONResponse, PlainTextResponse from fastapi.templating import Jinja2Templates from app import security from app import firewall from app import vps from app import docker_ops LOGS_DIR = Path(os.environ.get("LOGS_DIR", "/logs")) CONFIG_DIR = Path(os.environ.get("CONFIG_DIR", "/config")) PROXIES_FILE = CONFIG_DIR / "proxies.json" @asynccontextmanager async def lifespan(app: FastAPI): if os.environ.get("WEBHOOK_URL"): t = threading.Thread(target=security.alert_loop, daemon=True) t.start() yield app = FastAPI(title="Apache Proxy Monitor", lifespan=lifespan) templates = Jinja2Templates(directory=Path(__file__).parent / "templates") API_TOKEN = os.environ.get("API_TOKEN", "") @app.middleware("http") async def token_auth(request: Request, call_next): if API_TOKEN and request.url.path.startswith("/api/"): if request.headers.get("X-API-Token") != API_TOKEN: return JSONResponse({"error": "unauthorized"}, status_code=401) return await call_next(request) # Combined log format, with optional "%v:%p" vhost prefix (other_vhosts_access.log) ACCESS_RE = re.compile( r"^(?:(?P\S+:\d+)\s+)?" r"(?P\S+)\s+\S+\s+\S+\s+" r"\[(?P[^\]]+)\]\s+" r'"(?P\S+)\s+(?P\S+)(?:\s+(?P[^"]*))?"\s+' r"(?P\d{3})\s+(?P\S+)\s+" r'"(?P[^"]*)"\s+"(?P[^"]*)"' ) ERROR_RE = re.compile( r"\[(?P[^\]]+)\]\s+\[(?P[^\]]+)\]\s+\[pid \d+[^\]]*\]\s*(?P.*)" ) def load_vhosts(): if not PROXIES_FILE.exists(): return [] try: return json.loads(PROXIES_FILE.read_text()) except Exception: return [] def vhost_from_filename(path): name = path.name for suffix in ("-access.log", "_access.log"): if name.endswith(suffix): return name[: -len(suffix)] return name.replace(".log", "") def parse_ts(s): try: return datetime.strptime(s, "%d/%b/%Y:%H:%M:%S %z") except Exception: return None def read_requests(max_lines=200): requests = [] for path in sorted(LOGS_DIR.glob("*access*.log")): fname_vhost = vhost_from_filename(path) try: lines = path.read_text(errors="replace").splitlines() except Exception: continue for line in lines: m = ACCESS_RE.match(line) if not m: continue dt = parse_ts(m.group("ts")) requests.append( { "vhost": m.group("vhost") or fname_vhost, "ip": m.group("ip"), "ts": m.group("ts"), "time": dt.timestamp() if dt else 0, "method": m.group("method"), "path": m.group("path"), "proto": m.group("proto") or "", "status": m.group("status"), "bytes": m.group("bytes"), "referer": m.group("referer"), "agent": m.group("agent"), } ) requests.sort(key=lambda r: r["time"], reverse=True) return requests[:max_lines] @app.get("/", response_class=HTMLResponse) async def index(request: Request): return templates.TemplateResponse(request, "index.html", {"request": request}) @app.get("/api/vhosts") async def api_vhosts(): return load_vhosts() @app.get("/api/logs") async def api_logs(lines: int = Query(200, ge=1, le=2000)): return read_requests(lines) @app.get("/api/logs/ip/{ip}") async def api_logs_ip(ip: str): reqs = read_requests(max_lines=100000) return [r for r in reqs if r["ip"] == ip][:1000] @app.get("/api/vps/logs/ip/{ip}") async def api_vps_logs_ip(ip: str): return vps.fetch_logs_ip(ip) @app.get("/api/docker/info") async def api_docker_info(): return docker_ops.host_info() @app.get("/api/docker/containers") async def api_docker_containers(): return docker_ops.list_containers() @app.get("/api/docker/images") async def api_docker_images(): return docker_ops.list_images() @app.post("/api/docker/containers/{cid}/{action}") async def api_docker_action(cid: str, action: str): return docker_ops.container_action(cid, action) @app.get("/api/docker/containers/{cid}/stats") async def api_docker_stats(cid: str): return docker_ops.container_stats(cid) @app.get("/api/errors") async def api_errors(lines: int = Query(50, ge=1, le=500)): path = LOGS_DIR / "error.log" if not path.exists(): return [] try: raw = path.read_text(errors="replace").splitlines() except Exception: return [] return raw[-lines:][::-1] @app.get("/api/stats") async def api_stats(): reqs = read_requests(max_lines=100000) vhosts = load_vhosts() status_counts = Counter(r["status"] for r in reqs) ip_counts = Counter(r["ip"] for r in reqs) vhost_counts = Counter(r["vhost"] for r in reqs) return { "total": len(reqs), "unique_ips": len(ip_counts), "status": dict(status_counts.most_common()), "top_ips": [{"ip": ip, "count": c} for ip, c in ip_counts.most_common(15)], "per_vhost": [{"vhost": v, "count": c} for v, c in vhost_counts.most_common()], "errors_4xx": sum(c for s, c in status_counts.items() if s.startswith("4")), "errors_5xx": sum(c for s, c in status_counts.items() if s.startswith("5")), "vhosts": vhosts, } @app.get("/api/security") async def api_security(hours: int = Query(168, ge=1, le=720)): return security.security_summary(hours=hours) @app.get("/api/security/ip/{ip}") async def api_security_ip(ip: str): return security.ip_attempts(ip) @app.get("/api/geo/{ip}") async def api_geo(ip: str): return security.geolocate([ip]).get(ip, {}) @app.post("/api/block/{ip}") async def api_block_ip(ip: str): results = [] for s in firewall.load_servers(): if s.get("host") == "local": continue results.append({"server": s.get("name"), "result": firewall.block_ip(s, ip)}) results.append({"server": "ubuntu VPS (.2)", "result": vps.block_ip(ip)}) return results @app.get("/api/threats/export") async def api_threats_export(): sec = security.security_summary(hours=720) buf = StringIO() w = csv.writer(buf) w.writerow(["ip", "country", "cc", "org", "asn", "attempts", "first_seen", "last_seen"]) for a in sec.get("attempts", []): w.writerow([a["ip"], a["country"], a["cc"], a["org"], a["as"], a["count"], a["first"], a["last"]]) return PlainTextResponse(buf.getvalue(), media_type="text/csv", headers={"Content-Disposition": "attachment; filename=threats.csv"}) @app.get("/api/security/nodes") async def api_security_nodes(): return security.registered_nodes() @app.get("/api/security/status") async def api_security_status(): return security.compromise_status() @app.get("/api/firewall") async def api_firewall(): results = [firewall.firewall_status(s) for s in firewall.load_servers()] try: agent_fw = vps.fetch_firewall() agent_names = {a.get("name") for a in agent_fw} results = [r for r in results if r.get("name") not in agent_names] + agent_fw except Exception: pass return results @app.post("/api/firewall/{server_id}/block") async def api_firewall_block(server_id: int, data: dict): servers = firewall.load_servers() if server_id < 0 or server_id >= len(servers): raise HTTPException(404, "Server not found") return firewall.block_ip(servers[server_id], data.get("ip", "")) @app.post("/api/firewall/{server_id}/unblock") async def api_firewall_unblock(server_id: int, data: dict): servers = firewall.load_servers() if server_id < 0 or server_id >= len(servers): raise HTTPException(404, "Server not found") return firewall.unblock_ip(servers[server_id], data.get("ip", "")) @app.post("/api/firewall/{server_id}/toggle") async def api_firewall_toggle(server_id: int, data: dict): servers = firewall.load_servers() if server_id < 0 or server_id >= len(servers): raise HTTPException(404, "Server not found") t = data.get("type") enabled = bool(data.get("enabled", True)) if t == "country": return firewall.toggle_country(servers[server_id], enabled) if t == "scanner": return firewall.toggle_scanner(servers[server_id], enabled) raise HTTPException(400, "Unknown toggle type") @app.get("/api/vps/logs") async def api_vps_logs(lines: int = Query(500, ge=1, le=5000)): return vps.fetch_logs(lines) @app.get("/api/vps/stats") async def api_vps_stats(): return vps.fetch_stats()