From 35130e5528a52161630fd1c93572e1b2b7cd911b Mon Sep 17 00:00:00 2001 From: Christophe Besson Date: Tue, 11 Aug 2026 23:11:36 +0200 Subject: feat(node): audit logging + local admin UI rewrite Add SQLite audit store for legal compliance (LCEN/DSA): logs user IP, actions (handshake, file download/upload/delete, stream, chat), and timestamps. Retention: 1 year, with cleanup method. WebRTC transport now logs all user actions to the audit store with remote IP extraction from the ICE transport. Local web UI rewritten as a proper admin dashboard: - Stats cards (groups, files, peers) - Connected peers table with IP, username, group, state - Group cards with file listings and shared directory info - Audit log page with event/user filtering - Dark theme, responsive, auto-refresh - JSON API: /api/status, /api/groups, /api/peers, /api/audit, /api/config Co-Authored-By: Claude Opus 4.6 --- packages/meshbay-node/src/meshbay_node/audit.py | 148 ++++++ packages/meshbay-node/src/meshbay_node/daemon.py | 17 + .../src/meshbay_node/transport/webrtc_server.py | 59 ++- packages/meshbay-node/src/meshbay_node/ui/app.py | 542 ++++++++++++++++----- 4 files changed, 631 insertions(+), 135 deletions(-) create mode 100644 packages/meshbay-node/src/meshbay_node/audit.py (limited to 'packages/meshbay-node/src') diff --git a/packages/meshbay-node/src/meshbay_node/audit.py b/packages/meshbay-node/src/meshbay_node/audit.py new file mode 100644 index 0000000..6346f16 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/audit.py @@ -0,0 +1,148 @@ +""" +MeshBay Node — SQLite audit log for legal compliance. + +Logs user actions with IP address, timestamp, and details. +Required by LCEN (France), EU e-Commerce Directive, and DSA +for hosting service operators. + +Retention: 1 year minimum. Cleanup is caller's responsibility. +""" + +import logging +import time +from dataclasses import dataclass +from pathlib import Path + +import aiosqlite + +log = logging.getLogger(__name__) + +_SCHEMA = """ +CREATE TABLE IF NOT EXISTS audit_log ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + timestamp REAL NOT NULL, + user_id TEXT NOT NULL, + username TEXT NOT NULL DEFAULT '', + ip TEXT NOT NULL DEFAULT '', + event TEXT NOT NULL, + group_id TEXT NOT NULL DEFAULT '', + detail TEXT NOT NULL DEFAULT '' +); +CREATE INDEX IF NOT EXISTS idx_audit_ts ON audit_log(timestamp); +CREATE INDEX IF NOT EXISTS idx_audit_user ON audit_log(user_id); +CREATE INDEX IF NOT EXISTS idx_audit_event ON audit_log(event); +""" + +EVENTS = { + "connect", + "disconnect", + "handshake", + "file_download", + "file_upload", + "file_delete", + "stream_video", + "chat_message", + "chat_history", + "index_sync", + "auth_failed", +} + +RETENTION_DAYS = 365 + + +@dataclass +class AuditEntry: + id: int + timestamp: float + user_id: str + username: str + ip: str + event: str + group_id: str + detail: str + + +class AuditStore: + """Async SQLite audit log.""" + + def __init__(self, db_path: Path): + self._db_path = db_path + self._db: aiosqlite.Connection | None = None + + async def open(self) -> None: + self._db_path.parent.mkdir(parents=True, exist_ok=True) + self._db = await aiosqlite.connect(str(self._db_path)) + await self._db.executescript(_SCHEMA) + await self._db.commit() + + async def close(self) -> None: + if self._db: + await self._db.close() + self._db = None + + async def log_event( + self, + user_id: str, + event: str, + ip: str = "", + username: str = "", + group_id: str = "", + detail: str = "", + ) -> None: + if not self._db: + return + await self._db.execute( + "INSERT INTO audit_log (timestamp, user_id, username, ip, event, group_id, detail) " + "VALUES (?, ?, ?, ?, ?, ?, ?)", + (time.time(), user_id, username, ip, event, group_id, detail), + ) + await self._db.commit() + + async def get_entries( + self, + since: float = 0, + limit: int = 200, + user_id: str | None = None, + event: str | None = None, + ) -> list[AuditEntry]: + conditions = ["timestamp > ?"] + params: list = [since] + if user_id: + conditions.append("user_id = ?") + params.append(user_id) + if event: + conditions.append("event = ?") + params.append(event) + params.append(limit) + + where = " AND ".join(conditions) + cursor = await self._db.execute( + f"SELECT id, timestamp, user_id, username, ip, event, group_id, detail " + f"FROM audit_log WHERE {where} ORDER BY timestamp DESC LIMIT ?", + params, + ) + rows = await cursor.fetchall() + return [ + AuditEntry( + id=r[0], timestamp=r[1], user_id=r[2], username=r[3], + ip=r[4], event=r[5], group_id=r[6], detail=r[7], + ) + for r in rows + ] + + async def entry_count(self) -> int: + if not self._db: + return 0 + cursor = await self._db.execute("SELECT COUNT(*) FROM audit_log") + row = await cursor.fetchone() + return row[0] + + async def cleanup(self, retention_days: int = RETENTION_DAYS) -> int: + """Delete entries older than retention_days. Returns count deleted.""" + if not self._db: + return 0 + cutoff = time.time() - (retention_days * 86400) + cursor = await self._db.execute( + "DELETE FROM audit_log WHERE timestamp < ?", (cutoff,)) + await self._db.commit() + return cursor.rowcount diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index 5b6e770..5851b34 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -33,6 +33,7 @@ import uvicorn from meshbay_common import MNP_VERSION from meshbay_common.protocol import MNP +from meshbay_node.audit import AuditStore from meshbay_node.chat.store import ChatStore from meshbay_node.config import Config, load_config, write_example_config from meshbay_node.hub_client import HubClient, HubConfig @@ -109,6 +110,7 @@ class NodeDaemon: self._webrtc = None self._denylist = Denylist() if Denylist else None self._chat_stores: dict[str, ChatStore] = {} + self._audit_store: AuditStore | None = None self._indexers: list[DirectoryIndexer] = [] self._tasks: list[asyncio.Task] = [] self._hub: HubClient | None = None @@ -192,6 +194,12 @@ class NodeDaemon: groups_ctx[gid]["chat_store"] = store log.info("Chat stores opened: %d groups", len(self._chat_stores)) + # 4b. Audit store (legal compliance — IP + action logging) + audit_db = data_dir / "audit.db" + self._audit_store = AuditStore(db_path=audit_db) + await self._audit_store.open() + log.info("Audit store opened: %s", audit_db) + # 5. Denylist denylist = self._denylist @@ -210,6 +218,7 @@ class NodeDaemon: self._webrtc._ctx["chat_store"] = first.get("chat_store") self._webrtc._ctx["hub_ws"] = _WsSender(hub) self._webrtc._ctx["node_user_id"] = session.user_id + self._webrtc._ctx["audit_store"] = self._audit_store log.info("WebRTC transport ready") else: log.warning("WebRTC not available (aiortc not installed)") @@ -315,6 +324,11 @@ class NodeDaemon: group_cfg.http_port, group_cfg.name) # 10. Local web UI + self._state["groups_ctx"] = groups_ctx + self._state["config"] = self._config + self._state["audit_store"] = self._audit_store + self._state["webrtc"] = self._webrtc + self._state["hub"] = hub from meshbay_node.ui import create_ui_app ui_app = create_ui_app(self._state) ui_cfg = uvicorn.Config( @@ -412,6 +426,9 @@ class NodeDaemon: if self._webrtc: await self._webrtc.close_all() + if self._audit_store: + await self._audit_store.close() + for store in self._chat_stores.values(): await store.close() diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py index ab95550..e692c80 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py @@ -124,16 +124,35 @@ class _DataChannelBuffer: yield msgpack.unpackb(msg_bytes, raw=False) +def _get_remote_ip(pc: RTCPeerConnection) -> str: + """Best-effort extraction of the remote peer IP from the ICE transport.""" + try: + dtls = pc.sctp and pc.sctp.transport + ice = dtls and dtls.transport + conn = ice and ice._connection + if conn and hasattr(conn, '_nominated') and conn._nominated: + for pair in conn._nominated.values(): + return pair.remote_candidate.host + if conn and conn.remote_candidates: + return conn.remote_candidates[0].host + except Exception: + pass + return "" + + class WebRTCPeerSession: """One WebRTC peer connection, handling MNP over a DataChannel.""" - def __init__(self, pc: RTCPeerConnection, node_ctx: dict): + def __init__(self, pc: RTCPeerConnection, node_ctx: dict, peer_id: str = ""): self._pc = pc self._ctx = node_ctx self._channel: RTCDataChannel | None = None self._buffer = _DataChannelBuffer() self._user_id: str | None = None self._group_id: str | None = None + self._peer_id: str = peer_id + self._remote_ip: str = "" + self._username: str = "" def _setup_channel(self, channel: RTCDataChannel) -> None: self._channel = channel @@ -178,6 +197,20 @@ class WebRTCPeerSession: log.error("Error handling %s on DataChannel: %s", mtype, e) self._send({"type": "error", "detail": str(e)}) + def _audit(self, event: str, detail: str = "") -> None: + audit = self._ctx.get("audit_store") + if audit and self._user_id: + if not self._remote_ip: + self._remote_ip = _get_remote_ip(self._pc) + asyncio.ensure_future(audit.log_event( + user_id=self._user_id, + event=event, + ip=self._remote_ip, + username=self._username, + group_id=self._group_id or "", + detail=detail, + )) + def _do_handshake(self, msg: dict) -> None: token = msg.get("token", "") group_id = msg.get("group_id", "") @@ -185,6 +218,7 @@ class WebRTCPeerSession: decoded = jwt.decode(token, self._ctx["hub_pk_pem"], algorithms=["EdDSA"]) except Exception as e: self._send({"type": "error", "detail": f"Invalid JWT: {e}"}) + self._audit_auth_failed(group_id, str(e)) return denylist = self._ctx.get("denylist") @@ -202,6 +236,7 @@ class WebRTCPeerSession: self._user_id = decoded["sub"] self._group_id = group_id + self._username = decoded.get("username", "") peers = self._ctx.get("_peers") if peers is not None: @@ -214,6 +249,19 @@ class WebRTCPeerSession: "v": MNP_VERSION, "node_pk": pk_to_b64(self._ctx["sk_node"].public_key()), }) + self._audit("handshake") + + def _audit_auth_failed(self, group_id: str, reason: str) -> None: + audit = self._ctx.get("audit_store") + if audit: + self._remote_ip = _get_remote_ip(self._pc) + asyncio.ensure_future(audit.log_event( + user_id="unknown", + event="auth_failed", + ip=self._remote_ip, + group_id=group_id, + detail=reason, + )) def _group_ctx(self) -> dict: if "groups" in self._ctx and self._group_id: @@ -274,6 +322,8 @@ class WebRTCPeerSession: file_hash, ) self._send(chunk_data) + if chunk_index == 0: + self._audit("file_download", entry.name) def _do_stream_segment(self, msg: dict) -> None: ctx = self._group_ctx() @@ -365,6 +415,7 @@ class WebRTCPeerSession: pass self._send({"type": "ack", "v": MNP_VERSION}) + self._audit("chat_message") def _do_chat_history(self, msg: dict) -> None: chat_store = self._ctx.get("chat_store") @@ -441,6 +492,7 @@ class WebRTCPeerSession: final_path = shared_root / safe_name tmp_path.rename(final_path) log.info("Upload complete: %s (%d chunks)", safe_name, total_chunks) + self._audit("file_upload", safe_name) def _do_file_delete(self, msg: dict) -> None: ctx = self._group_ctx() @@ -463,6 +515,7 @@ class WebRTCPeerSession: if file_path.exists(): file_path.unlink() log.info("File deleted: %s", entry.name) + self._audit("file_delete", entry.name) ctx["index"].remove_entry(file_id) self._send({ @@ -549,6 +602,7 @@ class WebRTCPeerSession: "file_id": file_id, }) log.info("Streamed %s: %d segments", entry.name, index) + self._audit("stream_video", entry.name) def _send(self, obj: dict) -> None: if self._channel and self._channel.readyState == "open": @@ -558,6 +612,7 @@ class WebRTCPeerSession: self._channel.readyState if self._channel else "none") async def close(self) -> None: + self._audit("disconnect") peers = self._ctx.get("_peers") if peers and self._user_id: peers.pop(self._user_id, None) @@ -639,7 +694,7 @@ class WebRTCTransport: iceServers=[RTCIceServer(urls=s) for s in self._stun] if self._stun else [] ) pc = RTCPeerConnection(configuration=config) - session = WebRTCPeerSession(pc, self._ctx) + session = WebRTCPeerSession(pc, self._ctx, peer_id=peer_id) self._sessions[peer_id] = session @pc.on("datachannel") diff --git a/packages/meshbay-node/src/meshbay_node/ui/app.py b/packages/meshbay-node/src/meshbay_node/ui/app.py index a63e28f..5e77ed8 100644 --- a/packages/meshbay-node/src/meshbay_node/ui/app.py +++ b/packages/meshbay-node/src/meshbay_node/ui/app.py @@ -1,138 +1,171 @@ """ -MeshBay Node — local web UI (localhost:18000). +MeshBay Node — local administration web UI (localhost:18000). -Minimal FastAPI app providing: - GET / → status page (HTML) - GET /api/status → JSON status - GET /api/files → JSON file list from the group index - GET /api/config → JSON config summary (no secrets) +FastAPI app providing: + - Dashboard: node status, connected peers, group overview + - Groups: file listing, shared directory info + - Peers: connected WebRTC/QUIC clients + - Audit log: IP + action log for legal compliance + - API endpoints for all data (JSON) Served only on 127.0.0.1 — not exposed to the network. +No authentication required (localhost only). """ -import asyncio import json import logging -from typing import TYPE_CHECKING +import time +from pathlib import Path -from fastapi import FastAPI, WebSocket, WebSocketDisconnect +from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Query from fastapi.responses import HTMLResponse from meshbay_node import __version__ -if TYPE_CHECKING: - from meshbay_node.chat.store import ChatStore - from meshbay_node.indexer import GroupIndex - log = logging.getLogger(__name__) def create_ui_app(state: dict) -> FastAPI: - """ - Create the UI FastAPI app. - - state dict is updated by the daemon and read by UI endpoints: - state["status"] : str — "starting" | "running" | "error" - state["group_id"] : str - state["group_name"] : str - state["hub_url"] : str - state["username"] : str - state["node_port"] : int - state["index"] : GroupIndex | None - state["endpoint_hint"]: str | None - """ app = FastAPI( - title="MeshBay Node UI", + title="MeshBay Node Admin", version=__version__, - docs_url=None, # disable Swagger on local UI + docs_url=None, redoc_url=None, ) + # ── JSON API ───────────────────────────────────────────────────────────── + @app.get("/api/status") async def api_status(): - index = state.get("index") + indexes = state.get("indexes", {}) + total_files = sum(idx.count for idx in indexes.values()) + groups_ctx = state.get("groups_ctx", {}) + webrtc = state.get("webrtc") return { - "version": __version__, - "status": state.get("status", "starting"), - "hub_url": state.get("hub_url", ""), - "username": state.get("username", ""), - "group_id": state.get("group_id", ""), - "group_name": state.get("group_name", ""), - "node_port": state.get("node_port", 0), + "version": __version__, + "status": state.get("status", "starting"), + "hub_url": state.get("hub_url", ""), + "username": state.get("username", ""), + "node_port": state.get("node_port", 0), + "quic_port": state.get("quic_port", 0), "endpoint_hint": state.get("endpoint_hint"), - "file_count": index.count if index else 0, - "index_version": index.version if index else 0, + "group_count": len(groups_ctx), + "total_files": total_files, + "webrtc_peers": webrtc.active_peers if webrtc else 0, } - @app.get("/api/files") - async def api_files(): - index = state.get("index") - if not index: + @app.get("/api/groups") + async def api_groups(): + groups_ctx = state.get("groups_ctx", {}) + config = state.get("config") + result = [] + for gid, ctx in groups_ctx.items(): + cfg = None + if config: + cfg = next((g for g in config.groups if g.id == gid), None) + idx = ctx.get("index") + result.append({ + "id": gid, + "name": cfg.name if cfg else gid[:8], + "shared_dir": str(ctx.get("shared_root", "")), + "visibility": cfg.visibility if cfg else "private", + "file_count": idx.count if idx else 0, + "index_version": idx.version if idx else 0, + }) + return {"groups": result} + + @app.get("/api/groups/{group_id}/files") + async def api_group_files(group_id: str): + groups_ctx = state.get("groups_ctx", {}) + ctx = groups_ctx.get(group_id) + if not ctx: + return {"files": []} + idx = ctx.get("index") + if not idx: return {"files": []} return { "files": [ { - "id": e.id[:16] + "…", - "name": e.name, - "path": e.path, - "size": e.size, - "type": e.type, - "duration": e.duration, + "id": e.id, + "name": e.name, + "path": e.path, + "size": e.size, + "type": e.type, + "added_at": e.added_at, } - for e in index.entries + for e in idx.entries ] } - @app.get("/", response_class=HTMLResponse) - async def root(): - index = state.get("index") - status = state.get("status", "starting") - file_count = index.count if index else 0 - status_color = {"running": "#22c55e", "error": "#ef4444"}.get(status, "#f59e0b") - - files_html = "" - if index: - rows = "".join( - f"{e.name}{e.type}" - f"{e.size // 1024} KB{e.path or '/'}" - for e in index.entries - ) - files_html = f""" -

Files ({file_count})

- - - {rows} -
NameTypeSizePath
""" - - return f""" - - - - MeshBay Node - - - - -

🔗 MeshBay Node {status}

-

- Hub: {state.get("hub_url", "—")}  |  - User: {state.get("username", "—")}  |  - Group: {state.get("group_name") or state.get("group_id") or "—"}  |  - Port: {state.get("node_port", "—")} -

-

Endpoint: {state.get("endpoint_hint") or "unknown"}

- {files_html} -
- MeshBay Node v{__version__} — JSON status - — JSON files — Chat - -""" + @app.get("/api/peers") + async def api_peers(): + webrtc = state.get("webrtc") + if not webrtc: + return {"peers": []} + peers = [] + for pid, session in list(webrtc._sessions.items()): + from meshbay_node.transport.webrtc_server import _get_remote_ip + peers.append({ + "peer_id": pid, + "user_id": session._user_id or "", + "username": session._username or "", + "group_id": session._group_id or "", + "remote_ip": session._remote_ip or _get_remote_ip(session._pc), + "state": session._pc.connectionState, + }) + return {"peers": peers} + + @app.get("/api/audit") + async def api_audit( + since: float = 0, + limit: int = 200, + user_id: str | None = Query(default=None), + event: str | None = Query(default=None), + ): + audit = state.get("audit_store") + if not audit: + return {"entries": []} + entries = await audit.get_entries( + since=since, limit=limit, user_id=user_id, event=event) + return { + "entries": [ + { + "id": e.id, + "timestamp": e.timestamp, + "user_id": e.user_id, + "username": e.username, + "ip": e.ip, + "event": e.event, + "group_id": e.group_id, + "detail": e.detail, + } + for e in entries + ] + } + + @app.get("/api/config") + async def api_config(): + config = state.get("config") + if not config: + return {} + return { + "hub_url": config.hub.url, + "username": config.hub.username, + "node_port": config.node.port, + "quic_port": config.node.quic_port, + "http_port": config.node.http_port, + "ui_port": config.node.ui_port, + "data_dir": str(config.data_dir), + "groups": [ + { + "id": g.id, + "name": g.name, + "shared_dir": g.shared_dir, + "visibility": g.visibility, + } + for g in config.groups + ], + } # ── Chat endpoints ─────────────────────────────────────────────────────── @@ -159,7 +192,6 @@ def create_ui_app(state: dict) -> FastAPI: @app.websocket("/ws/chat") async def chat_websocket(ws: WebSocket): - """WebSocket for real-time chat push to the local UI.""" await ws.accept() _chat_subscribers.append(ws) try: @@ -171,7 +203,6 @@ def create_ui_app(state: dict) -> FastAPI: _chat_subscribers.remove(ws) async def broadcast_chat_to_ui(msg: dict) -> None: - """Push a chat message to all connected UI WebSocket clients.""" payload = json.dumps(msg) dead = [] for ws in _chat_subscribers: @@ -184,42 +215,287 @@ def create_ui_app(state: dict) -> FastAPI: app.broadcast_chat = broadcast_chat_to_ui - @app.get("/chat", response_class=HTMLResponse) - async def chat_page(): - return f""" + # ── HTML UI ────────────────────────────────────────────────────────────── + + @app.get("/", response_class=HTMLResponse) + async def root(): + return _render_page(state) + + @app.get("/audit", response_class=HTMLResponse) + async def audit_page(): + return _render_audit_page() + + return app + + +def _fmt_size(n: int) -> str: + if n < 1024: + return f"{n} B" + if n < 1024 * 1024: + return f"{n / 1024:.1f} KB" + if n < 1024 * 1024 * 1024: + return f"{n / (1024 * 1024):.1f} MB" + return f"{n / (1024 * 1024 * 1024):.2f} GB" + + +def _render_page(state: dict) -> str: + status = state.get("status", "starting") + indexes = state.get("indexes", {}) + groups_ctx = state.get("groups_ctx", {}) + config = state.get("config") + webrtc = state.get("webrtc") + total_files = sum(idx.count for idx in indexes.values()) + peer_count = webrtc.active_peers if webrtc else 0 + status_color = {"running": "#22c55e", "error": "#ef4444"}.get(status, "#f59e0b") + + # Groups section + groups_html = "" + for gid, ctx in groups_ctx.items(): + cfg = None + if config: + cfg = next((g for g in config.groups if g.id == gid), None) + idx = ctx.get("index") + name = cfg.name if cfg else gid[:8] + shared = ctx.get("shared_root", "") + vis = cfg.visibility if cfg else "private" + fcount = idx.count if idx else 0 + total_size = sum(e.size for e in idx.entries) if idx else 0 + + file_rows = "" + if idx: + for e in sorted(idx.entries, key=lambda x: x.name): + file_rows += ( + f"{e.name}{e.type}" + f"{_fmt_size(e.size)}{e.path or '/'}" + ) + + groups_html += f""" +
+

{name} + {vis} +

+

Directory: {shared}

+

Files: {fcount} — Total: {_fmt_size(total_size)}

+

ID: {gid}

+
File list + + + {file_rows} +
NameTypeSizePath
+
+
""" + + # Peers section + peers_html = "" + if webrtc: + for pid, session in list(webrtc._sessions.items()): + from meshbay_node.transport.webrtc_server import _get_remote_ip + ip = session._remote_ip or _get_remote_ip(session._pc) + peers_html += ( + f"{session._username or session._user_id or '—'}" + f"{ip or '—'}" + f"{session._group_id[:8] if session._group_id else '—'}" + f"{session._pc.connectionState}" + ) + if not peers_html: + peers_html = 'No connected peers' + + return f""" - - MeshBay Chat - + +MeshBay Node Admin + + -

MeshBay Chat

-
-

Back to status

- +
+

MeshBay Node {status}

+ + + +
+
{len(groups_ctx)}
Groups
+
{total_files}
Files
+
{peer_count}
Connected Peers
+
+
{state.get("username", "—")}
+
User
+
+
+ +

Connected Peers

+ + + {peers_html} +
UserIPGroupState
+ +

Groups

+ {groups_html or '

No groups configured

'} + +

Node Configuration

+
+

Hub: {state.get("hub_url", "—")}

+

QUIC port: {state.get("quic_port", "—")} — + TCP port: {state.get("node_port", "—")}

+

Node ID: {state.get("endpoint_hint") or "—"}

+
+ + +
+ """ - return app + +def _render_audit_page() -> str: + return """ + + + +MeshBay Node — Audit Log + + + + +
+

Audit Log

+ + +
+ + + +
+ +
+ + + +
TimeEventUserIPGroupDetail
+
+ + +""" -- cgit v1.2.3