From 1f6beed0d2074ed7aa0a7fd02f4015fa3101e97c Mon Sep 17 00:00:00 2001 From: Athena Kaminsky Date: Fri, 7 Aug 2026 09:41:09 -0500 Subject: [PATCH] feat(security): SSRF guard on fetch_url + single-use tool-approval tokens Two tool/agent-layer hardening changes: * fetch_url now resolves the target host and refuses to connect if any resolved address is loopback, private (RFC1918/ULA), link-local (incl. the 169.254.169.254 cloud-metadata endpoint), multicast, reserved, or unspecified. IPv4-mapped IPv6 is unwrapped first, and the guard re-runs on every redirect hop so a public URL cannot 302 its way to an internal target. * /chat/approve now requires a single-use token minted when the stream pauses for approval and delivered only in that stream's tool_request event, compared in constant time. Previously the pending approval was keyed solely on a client-supplied conversation_id, so anyone who could enumerate a conversation_id could approve another client's pending action. The frontend threads the token from the tool_request event into the approve call. Co-authored-by: Cursor --- interface/web/src/Chatbot.jsx | 11 ++++- synapse/chat.py | 21 +++++++--- synapse/main.py | 11 ++++- synapse/tools.py | 77 +++++++++++++++++++++++++++++++---- 4 files changed, 104 insertions(+), 16 deletions(-) diff --git a/interface/web/src/Chatbot.jsx b/interface/web/src/Chatbot.jsx index 9583794..8188a6e 100644 --- a/interface/web/src/Chatbot.jsx +++ b/interface/web/src/Chatbot.jsx @@ -18,6 +18,7 @@ export function Chatbot({ conversationId, setConversationId, onConversationChang const [images, setImages] = useState([]); // {name, b64} for vision models const [activeTool, setActiveTool] = useState(null); // playbook tool currently running const [pendingApproval, setPendingApproval] = useState(null); // [{name, arguments}] awaiting yes/no + const [approvalToken, setApprovalToken] = useState(null); // single-use token authorizing /chat/approve const [editingIdx, setEditingIdx] = useState(null); // user message being edited const [editText, setEditText] = useState(""); const [listening, setListening] = useState(false); // mic dictation active @@ -286,7 +287,11 @@ export function Chatbot({ conversationId, setConversationId, onConversationChang continue; } if (pendingEventType === "tool_request") { - try { setPendingApproval(JSON.parse(payload)); } catch { /* ignore */ } + try { + const parsed = JSON.parse(payload); + setPendingApproval(parsed.actions || []); + setApprovalToken(parsed.token || null); + } catch { /* ignore */ } pendingEventType = null; continue; } @@ -401,13 +406,15 @@ export function Chatbot({ conversationId, setConversationId, onConversationChang // Approve or deny the pending action tool(s); the open chat stream resumes. const resolveApproval = async (approve) => { const req = pendingApproval || []; + const token = approvalToken; setPendingApproval(null); + setApprovalToken(null); const decisions = {}; req.forEach(a => { decisions[a.name] = approve; }); try { await fetch(`${API_BASE}/chat/approve`, { method: "POST", headers: { "Content-Type": "application/json" }, - body: JSON.stringify({ conversation_id: conversationId, decisions }), + body: JSON.stringify({ conversation_id: conversationId, token, decisions }), }); } catch { /* ignore */ } }; diff --git a/synapse/chat.py b/synapse/chat.py index 8e43c29..e479407 100644 --- a/synapse/chat.py +++ b/synapse/chat.py @@ -3,6 +3,7 @@ from __future__ import annotations import asyncio import json as _json import logging +import secrets import threading from typing import AsyncGenerator, Dict, List, Optional, Any @@ -171,12 +172,20 @@ async def _run_tool_loop(manager, messages, model, tool_schemas, temperature, nu action_calls = [c for c in calls if _tools.is_action(c.get("function", {}).get("name", ""))] if policy == "ask" and action_calls: event = asyncio.Event() - pending_approvals[conversation_id] = {"event": event, "decisions": {}} - yield "__approve__" + _json.dumps([ - {"name": c.get("function", {}).get("name", ""), - "arguments": c.get("function", {}).get("arguments")} - for c in action_calls - ]) + # Single-use capability token, delivered only to the client that owns + # this stream. /chat/approve requires it, so knowing the (guessable, + # enumerable) conversation_id is no longer enough to approve someone + # else's pending action. + token = secrets.token_urlsafe(32) + pending_approvals[conversation_id] = {"event": event, "decisions": {}, "token": token} + yield "__approve__" + _json.dumps({ + "token": token, + "actions": [ + {"name": c.get("function", {}).get("name", ""), + "arguments": c.get("function", {}).get("arguments")} + for c in action_calls + ], + }) try: await asyncio.wait_for(event.wait(), timeout=_APPROVAL_TIMEOUT) decisions = pending_approvals[conversation_id]["decisions"] diff --git a/synapse/main.py b/synapse/main.py index 839e4d0..07d54e5 100644 --- a/synapse/main.py +++ b/synapse/main.py @@ -3,6 +3,7 @@ from __future__ import annotations import asyncio as _asyncio import json as _json +import secrets as _secrets import uuid as _uuid from collections import Counter as _Counter from typing import Any, AsyncGenerator, Dict, List, Optional, Tuple @@ -508,12 +509,20 @@ async def chat_stream_endpoint(payload: Dict[str, Any]): @app.post("/chat/approve") async def chat_approve(payload: Dict[str, Any] = Body(...)): """Resolve a pending per-call tool approval. `decisions` maps tool name -> - bool; the awaiting chat stream resumes and runs the approved actions.""" + bool; the awaiting chat stream resumes and runs the approved actions. + + The `token` (issued in the stream's tool_request event) is required: without + it, anyone who can guess/enumerate a conversation_id could approve another + client's pending action. Compared in constant time.""" conversation_id = payload.get("conversation_id") or "" + token = payload.get("token") or "" decisions = payload.get("decisions") or {} waiter = _chat.pending_approvals.get(conversation_id) if not waiter: raise HTTPException(status_code=404, detail="no pending approval for this conversation") + expected = waiter.get("token") or "" + if not token or not _secrets.compare_digest(str(token), str(expected)): + raise HTTPException(status_code=403, detail="invalid or missing approval token") waiter["decisions"] = {k: bool(v) for k, v in decisions.items()} waiter["event"].set() return {"status": "resumed"} diff --git a/synapse/tools.py b/synapse/tools.py index 92b1ecc..a25b452 100644 --- a/synapse/tools.py +++ b/synapse/tools.py @@ -60,20 +60,83 @@ async def _web_search(query: str = "", **_) -> str: return res or "(no results)" +_FETCH_MAX_REDIRECTS = 5 + + +def _ip_is_blocked(ip: str) -> bool: + """True if an address is one an outbound fetch has no business reaching: + loopback, RFC1918/ULA private, link-local (incl. 169.254.169.254 cloud + metadata), multicast, reserved, or unspecified. IPv4-mapped IPv6 is unwrapped + first so ::ffff:127.0.0.1 can't sneak a loopback past the check.""" + import ipaddress + try: + addr = ipaddress.ip_address(ip.split("%")[0]) # drop any IPv6 zone id + except ValueError: + return True # unparseable -> refuse rather than guess + mapped = getattr(addr, "ipv4_mapped", None) + if mapped is not None: + addr = mapped + return ( + addr.is_loopback or addr.is_private or addr.is_link_local + or addr.is_multicast or addr.is_reserved or addr.is_unspecified + ) + + +def _ssrf_guard(host: str) -> str | None: + """Resolve a hostname and return an error string if ANY of its A/AAAA + records is a blocked address, else None. Checking every answer stops a name + from smuggling one private record alongside a public one. + + ponytail: this validates then httpx re-resolves on connect, so a sub-second + DNS-rebind could still slip a private address through the TOCTOU gap. That's + an advanced attack against a playbook-gated, single-user tool; pin the + connection to the resolved IP if this ever faces untrusted callers.""" + import socket + if not host: + return "missing host" + try: + infos = socket.getaddrinfo(host, None) + except socket.gaierror as e: + return f"cannot resolve host: {e}" + ips = {info[4][0] for info in infos} + if not ips: + return "host did not resolve" + blocked = [ip for ip in ips if _ip_is_blocked(ip)] + if blocked: + return f"refusing to fetch a private/loopback/link-local address ({', '.join(sorted(blocked))})" + return None + + async def _fetch_url(url: str = "", **_) -> str: import re import httpx + from urllib.parse import urlparse, urljoin url = (url or "").strip() if not url.startswith(("http://", "https://")): return json.dumps({"error": "url must start with http:// or https://"}) - # ponytail: no SSRF allow/deny-list — local single-user assistant, and the - # tool only runs when a playbook explicitly grants fetch_url. Add host - # filtering if this ever serves multiple/untrusted users. + # SSRF guard: validate the host of the initial URL AND every redirect hop + # against the private/loopback/link-local block-list before connecting, so a + # granted fetch_url can't be steered at 127.0.0.1:11434, cloud metadata, or + # LAN hosts — and a public URL can't 302 its way there either. try: - async with httpx.AsyncClient(timeout=15.0, follow_redirects=True) as c: - r = await c.get(url, headers={"User-Agent": "NexusOS/1.0"}) - r.raise_for_status() - html = r.text + async with httpx.AsyncClient(timeout=15.0, follow_redirects=False) as c: + for _ in range(_FETCH_MAX_REDIRECTS + 1): + parsed = urlparse(url) + if parsed.scheme not in ("http", "https"): + return json.dumps({"error": "only http(s) URLs are allowed"}) + err = _ssrf_guard(parsed.hostname or "") + if err: + return json.dumps({"error": f"blocked: {err}"}) + r = await c.get(url, headers={"User-Agent": "NexusOS/1.0"}) + location = r.headers.get("location") + if r.is_redirect and location: + url = urljoin(url, location) + continue + r.raise_for_status() + html = r.text + break + else: + return json.dumps({"error": "too many redirects"}) except Exception as e: return json.dumps({"error": f"fetch failed: {e}"}) text = re.sub(r"(?is)<(script|style).*?", " ", html)