mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-10-07 15:32:21 +02:00
fix(security): harden Python service boundaries
This commit is contained in:
+13
-3
@@ -3249,7 +3249,7 @@ async def _stream_llm_inner(url: str, model: str, messages: List[Dict], temperat
|
||||
yield f'event: error\ndata: {json.dumps({"error": "Network error", "status": 502, "fallback_eligible": False})}\n\n'
|
||||
except Exception as e:
|
||||
logger.error(f"Ollama stream error: {e}")
|
||||
yield f'event: error\ndata: {json.dumps({"error": str(e), "status": 502, "fallback_eligible": False})}\n\n'
|
||||
yield f'event: error\ndata: {json.dumps({"error": _stream_failure_message(e), "status": 502, "fallback_eligible": False})}\n\n'
|
||||
return
|
||||
|
||||
# ── Anthropic streaming ──
|
||||
@@ -3402,7 +3402,7 @@ async def _stream_llm_inner(url: str, model: str, messages: List[Dict], temperat
|
||||
yield f'event: error\ndata: {json.dumps({"error": "Network error", "status": 502, "fallback_eligible": False})}\n\n'
|
||||
except Exception as e:
|
||||
logger.error(f"Anthropic stream error: {e}")
|
||||
yield f'event: error\ndata: {json.dumps({"error": str(e), "status": 502, "fallback_eligible": False})}\n\n'
|
||||
yield f'event: error\ndata: {json.dumps({"error": _stream_failure_message(e), "status": 502, "fallback_eligible": False})}\n\n'
|
||||
return
|
||||
|
||||
# ── OpenAI-compatible streaming ──
|
||||
@@ -3872,7 +3872,17 @@ async def _stream_llm_inner(url: str, model: str, messages: List[Dict], temperat
|
||||
yield f'event: error\ndata: {json.dumps({"error": "Network error", "status": 502, "fallback_eligible": False})}\n\n'
|
||||
except Exception as e:
|
||||
logger.error(f"Stream error: {e}")
|
||||
yield f'event: error\ndata: {json.dumps({"error": str(e), "status": 502, "fallback_eligible": False})}\n\n'
|
||||
yield f'event: error\ndata: {json.dumps({"error": _stream_failure_message(e), "status": 502, "fallback_eligible": False})}\n\n'
|
||||
|
||||
|
||||
def _stream_failure_message(error: BaseException) -> str:
|
||||
"""Client-facing text for an unexpected streaming failure.
|
||||
|
||||
The raw exception can carry request URLs, local paths or provider internals;
|
||||
callers log it server-side and stream only this generic message, like the
|
||||
named transport failures above it.
|
||||
"""
|
||||
return f"Model stream failed ({type(error).__name__})"
|
||||
|
||||
|
||||
def _summarize_stream_error(err_chunk: Optional[str]) -> str:
|
||||
|
||||
@@ -0,0 +1,32 @@
|
||||
"""Per-principal extraction directories for email attachments.
|
||||
|
||||
Shared by the HTTP email routes and the email MCP server, which both extract
|
||||
attachments under MAIL_ATTACHMENTS_DIR.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
from pathlib import Path
|
||||
|
||||
|
||||
def attachment_scope_dir(root, folder, uid, *, owner, account_id) -> Path:
|
||||
"""Return the extraction directory for one message's attachments.
|
||||
|
||||
IMAP UIDs are small per-mailbox counters and folder names are server- or
|
||||
caller-supplied (`/`-delimited hierarchies, absolute or `..` segments), so
|
||||
neither may become a path. The directory is one hex segment derived from
|
||||
(owner, account, folder, uid): distinct principals and mailboxes never
|
||||
share it, and no input can steer it out of *root*. Containment is checked
|
||||
after resolution, so a symlinked entry cannot redirect it either.
|
||||
"""
|
||||
scope = json.dumps(
|
||||
[str(owner or ""), str(account_id or ""), str(folder or ""), str(uid or "")],
|
||||
ensure_ascii=False,
|
||||
)
|
||||
base = Path(root).resolve()
|
||||
target = (base / hashlib.sha256(scope.encode("utf-8")).hexdigest()[:32]).resolve()
|
||||
if target.parent != base:
|
||||
raise ValueError("attachment directory escapes the extraction root")
|
||||
return target
|
||||
+9
-2
@@ -6,11 +6,14 @@ writable, and storage is local-first. Served by ``GET /api/ready`` and suitable
|
||||
for an orchestrator readiness probe (200 only when every critical check passes).
|
||||
"""
|
||||
|
||||
import logging
|
||||
import os
|
||||
import uuid
|
||||
from datetime import datetime
|
||||
from typing import Dict
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def check_readiness() -> Dict[str, object]:
|
||||
"""Run the readiness checks and return a JSON-serialisable report.
|
||||
@@ -33,7 +36,10 @@ def check_readiness() -> Dict[str, object]:
|
||||
conn.execute(sql_text("SELECT 1"))
|
||||
checks["database"] = {"ok": True}
|
||||
except Exception as e:
|
||||
checks["database"] = {"ok": False, "error": str(e)}
|
||||
# The raw driver error can carry the DB host/user/path; keep it in the
|
||||
# server log and give the client only the exception type.
|
||||
logger.warning("Readiness database check failed: %s", e)
|
||||
checks["database"] = {"ok": False, "error_type": type(e).__name__}
|
||||
|
||||
# Data directory present and writable — home must be able to hold its own data.
|
||||
try:
|
||||
@@ -44,7 +50,8 @@ def check_readiness() -> Dict[str, object]:
|
||||
os.remove(probe)
|
||||
checks["data_dir"] = {"ok": True, "path": DATA_DIR}
|
||||
except Exception as e:
|
||||
checks["data_dir"] = {"ok": False, "error": str(e)}
|
||||
logger.warning("Readiness data_dir check failed: %s", e)
|
||||
checks["data_dir"] = {"ok": False, "error_type": type(e).__name__}
|
||||
|
||||
# Local-first: storage stays on the home machine (informational, never fatal).
|
||||
local_first = (
|
||||
|
||||
@@ -193,3 +193,17 @@ def strip_think(text: str, *, prose: bool = False, prompt_echo: bool = True) ->
|
||||
# from `src.research_utils` working while delegating to the central impl.
|
||||
def strip_thinking(text: str) -> str:
|
||||
return strip_think(text or "", prose=False, prompt_echo=True)
|
||||
|
||||
|
||||
_CLOSED_THINK_OPEN_RE = re.compile(r"<think(?:ing)?>", re.IGNORECASE)
|
||||
_CLOSED_THINK_CLOSE_RE = re.compile(r"</think(?:ing)?>", re.IGNORECASE)
|
||||
|
||||
|
||||
def strip_closed_think_blocks(text: str) -> str:
|
||||
"""Remove closed ``<think>``/``<thinking>`` blocks, leaving everything else.
|
||||
|
||||
Same result as ``re.sub(r'<think(?:ing)?>[\\s\\S]*?</think(?:ing)?>', '', text,
|
||||
flags=re.I)`` but forward-only (see _sub_delimited), so an unclosed opener
|
||||
flood in model output stays O(n) instead of O(n^2).
|
||||
"""
|
||||
return _sub_delimited(text or "", _CLOSED_THINK_OPEN_RE, _CLOSED_THINK_CLOSE_RE, lambda _inner: "")
|
||||
|
||||
@@ -236,10 +236,10 @@ async def do_manage_calendar(content: str, owner: Optional[str] = None, *, impor
|
||||
text = str(raw).strip().lower()
|
||||
if text in {"none", "no", "off", "false"}:
|
||||
return None
|
||||
m = re.search(r"(\d+)\s*(?:minutes?|mins?|m)\b", text)
|
||||
m = re.search(r"(?<!\d)(\d+)\s*(?:minutes?|mins?|m)\b", text)
|
||||
if m:
|
||||
return max(0, int(m.group(1)))
|
||||
m = re.search(r"(\d+)\s*(?:hours?|hrs?|h)\b", text)
|
||||
m = re.search(r"(?<!\d)(\d+)\s*(?:hours?|hrs?|h)\b", text)
|
||||
if m:
|
||||
return max(0, int(m.group(1)) * 60)
|
||||
if text.isdigit():
|
||||
@@ -251,7 +251,7 @@ async def do_manage_calendar(content: str, owner: Optional[str] = None, *, impor
|
||||
if minutes_before is None:
|
||||
return desc
|
||||
reminder_only = re.compile(
|
||||
r"^\s*(?:remind(?:er)?|alarm)\s*:?\s*\d+\s*"
|
||||
r"^\s*(?:remind(?:er)?|alarm)\s*(?::\s*)?\d+\s*"
|
||||
r"(?:minutes?|mins?|m|hours?|hrs?|h)\b.*$",
|
||||
re.I,
|
||||
)
|
||||
@@ -497,8 +497,8 @@ async def do_manage_calendar(content: str, owner: Optional[str] = None, *, impor
|
||||
delta = None
|
||||
if dur:
|
||||
import re as _re_d
|
||||
h = _re_d.search(r'(\d+)\s*(?:h|hr|hours?)', dur)
|
||||
m = _re_d.search(r'(\d+)\s*(?:m|min|minutes?)', dur)
|
||||
h = _re_d.search(r'(?<!\d)(\d+)\s*(?:h|hr|hours?)', dur)
|
||||
m = _re_d.search(r'(?<!\d)(\d+)\s*(?:m|min|minutes?)', dur)
|
||||
secs = (int(h.group(1)) * 3600 if h else 0) + (int(m.group(1)) * 60 if m else 0)
|
||||
if secs > 0:
|
||||
delta = timedelta(seconds=secs)
|
||||
|
||||
+1
-1
@@ -322,7 +322,7 @@ async def do_manage_notes(content: str, owner: Optional[str] = None) -> Dict:
|
||||
if looks_like_reminder:
|
||||
temporal = re.search(
|
||||
r"\b(?:today|tonight|tomorrow|tmrw|yesterday)\b(?:\s+(?:at\s+)?\d{1,2}(?::\d{2})?\s*(?:am|pm)?)?"
|
||||
r"|\b\d{1,2}(?::\d{2})?\s*(?:am|pm)?\s+(?:today|tonight|tomorrow|tmrw|yesterday)\b"
|
||||
r"|\b\d{1,2}(?::\d{2})?(?:\s*(?:am|pm))?\s+(?:today|tonight|tomorrow|tmrw|yesterday)\b"
|
||||
r"|\bin\s+\d+\s*(?:hour|hr|minute|min|day)s?\b",
|
||||
lower_combined,
|
||||
)
|
||||
|
||||
+70
-9
@@ -63,19 +63,80 @@ INTERNAL_UPLOAD_URL_RE = re.compile(
|
||||
r"([0-9a-fA-F]{32}(?:\.[A-Za-z0-9]+)?)"
|
||||
r"(?=$|[\s\"'<>\[\](){},;!?:&#]|\.(?![A-Za-z0-9]))"
|
||||
)
|
||||
PDF_SOURCE_UPLOAD_RE = re.compile(
|
||||
r"<!--\s*pdf(?:_form)?_source\b[^>]*\bupload_id="
|
||||
r"[\"']([0-9a-fA-F]{32}(?:\.[A-Za-z0-9]+)?)[\"'][^>]*-->",
|
||||
# `<!--\s*pdf(?:_form)?_source\b[^>]*\bupload_id=["'](id)["'][^>]*-->` and
|
||||
# `\[Attachment:[^\]\r\n]*\|\s*id=(id)(?:\s*\||\s*\])` are matched by the
|
||||
# forward-only scanners below. As single regexes, every opener in a run with no
|
||||
# closing `>` / `]` rescanned to the end of that run: O(n^2) on chat content,
|
||||
# which is unbounded on persisted assistant output (CodeQL py/polynomial-redos).
|
||||
_PDF_SOURCE_OPEN_RE = re.compile(r"<!--\s*pdf(?:_form)?_source\b", re.IGNORECASE)
|
||||
_PDF_SOURCE_ID_RE = re.compile(
|
||||
r"\bupload_id=[\"']([0-9a-fA-F]{32}(?:\.[A-Za-z0-9]+)?)[\"']",
|
||||
re.IGNORECASE,
|
||||
)
|
||||
ATTACHMENT_REFERENCE_LINE_RE = re.compile(
|
||||
r"\[Attachment:[^\]\r\n]*\|\s*id="
|
||||
r"([0-9a-fA-F]{32}(?:\.[A-Za-z0-9]+)?)"
|
||||
r"(?:\s*\||\s*\])",
|
||||
_ATTACHMENT_REFERENCE_OPEN_RE = re.compile(r"\[Attachment:", re.IGNORECASE)
|
||||
_ATTACHMENT_REFERENCE_STOP_RE = re.compile(r"[\]\r\n]")
|
||||
_ATTACHMENT_REFERENCE_TAIL_RE = re.compile(
|
||||
r"\|\s*id=([0-9a-fA-F]{32}(?:\.[A-Za-z0-9]+)?)(?:\s*\||\s*\])",
|
||||
re.IGNORECASE,
|
||||
)
|
||||
|
||||
|
||||
def _pdf_source_upload_ids(value: str) -> list[str]:
|
||||
"""IDs from `<!-- pdf_source ... upload_id="<id>" ... -->` comments.
|
||||
|
||||
Every opener before the next `>` shares that `>`: the comment matches only
|
||||
if `--` sits right before it, and then with the last `upload_id=` (the
|
||||
greedy `[^>]*` backtracks from the right). Otherwise all of those openers
|
||||
fail together, so the scan resumes after the `>`.
|
||||
"""
|
||||
found: list[str] = []
|
||||
pos = 0
|
||||
while True:
|
||||
opener = _PDF_SOURCE_OPEN_RE.search(value, pos)
|
||||
if opener is None:
|
||||
return found
|
||||
close = value.find(">", opener.end())
|
||||
if close < 0:
|
||||
return found
|
||||
if value[close - 2:close] == "--":
|
||||
last = None
|
||||
for last in _PDF_SOURCE_ID_RE.finditer(value, opener.end(), close):
|
||||
pass
|
||||
if last is not None:
|
||||
found.append(last.group(1))
|
||||
pos = close + 1
|
||||
|
||||
|
||||
def _attachment_reference_ids(value: str) -> list[str]:
|
||||
"""IDs from `[Attachment: name | id=<id> | ...]` reference lines.
|
||||
|
||||
The label scan stops at the first `]`/CR/LF, so every opener before that
|
||||
stop shares it and can only use the last `|` (scanning right to left)
|
||||
whose tail matches. If none does, all of those openers fail together and
|
||||
the scan resumes at the stop.
|
||||
"""
|
||||
found: list[str] = []
|
||||
pos = 0
|
||||
while True:
|
||||
opener = _ATTACHMENT_REFERENCE_OPEN_RE.search(value, pos)
|
||||
if opener is None:
|
||||
return found
|
||||
stop = _ATTACHMENT_REFERENCE_STOP_RE.search(value, opener.end())
|
||||
end = stop.start() if stop else len(value)
|
||||
tail = None
|
||||
pipe = value.rfind("|", opener.end(), end)
|
||||
while pipe >= 0:
|
||||
tail = _ATTACHMENT_REFERENCE_TAIL_RE.match(value, pipe)
|
||||
if tail is not None:
|
||||
break
|
||||
pipe = value.rfind("|", opener.end(), pipe)
|
||||
if tail is None:
|
||||
pos = end
|
||||
else:
|
||||
found.append(tail.group(1))
|
||||
pos = tail.end()
|
||||
|
||||
|
||||
def is_valid_upload_id(upload_id: str) -> bool:
|
||||
"""Return True when *upload_id* matches the canonical uploads.json id format."""
|
||||
return UPLOAD_ID_RE.fullmatch(upload_id or "") is not None
|
||||
@@ -110,8 +171,8 @@ def extract_internal_upload_ids(value: Any) -> set[str]:
|
||||
return set()
|
||||
return (
|
||||
set(INTERNAL_UPLOAD_URL_RE.findall(value))
|
||||
| set(PDF_SOURCE_UPLOAD_RE.findall(value))
|
||||
| set(ATTACHMENT_REFERENCE_LINE_RE.findall(value))
|
||||
| set(_pdf_source_upload_ids(value))
|
||||
| set(_attachment_reference_ids(value))
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -154,3 +154,52 @@ def check_outbound_url(
|
||||
if not saw_ip:
|
||||
return False, "host does not resolve to an IP"
|
||||
return True, "ok"
|
||||
|
||||
|
||||
|
||||
class OutboundAddressBlocked(PermissionError):
|
||||
"""A non-HTTP outbound host resolves into a disallowed address range."""
|
||||
|
||||
|
||||
def connect_outbound_tcp(
|
||||
host: str,
|
||||
port: int,
|
||||
*,
|
||||
timeout=socket._GLOBAL_DEFAULT_TIMEOUT,
|
||||
block_private: bool = False,
|
||||
source_address=None,
|
||||
resolver: Optional[Callable[..., list]] = None,
|
||||
) -> socket.socket:
|
||||
"""Open a TCP connection to *host* under the outbound address policy.
|
||||
|
||||
For raw TCP clients (IMAP/SMTP). The host is resolved exactly once, every
|
||||
resolved address is judged with the same policy as check_outbound_url, and
|
||||
the socket connects only to those already-judged addresses. A second
|
||||
lookup at connect time (what socket.create_connection(host) does) would let
|
||||
a rebinding name pass the check with a public answer and then connect to
|
||||
metadata/private space. TLS callers keep wrapping the returned socket with
|
||||
``server_hostname=host``, so SNI and certificate checks are unchanged.
|
||||
"""
|
||||
resolve = resolver or socket.getaddrinfo
|
||||
infos = resolve(host, port, 0, socket.SOCK_STREAM)
|
||||
for _family, _type, _proto, _canon, sockaddr in infos:
|
||||
ip = ipaddress.ip_address(str(sockaddr[0]).split("%")[0])
|
||||
reason = _classify(ip, block_private=block_private)
|
||||
if reason:
|
||||
raise OutboundAddressBlocked(reason)
|
||||
last_error: Optional[OSError] = None
|
||||
for family, socktype, proto, _canon, sockaddr in infos:
|
||||
sock = socket.socket(family, socktype, proto)
|
||||
try:
|
||||
if timeout is not socket._GLOBAL_DEFAULT_TIMEOUT: # same contract as create_connection
|
||||
sock.settimeout(timeout)
|
||||
if source_address:
|
||||
sock.bind(source_address)
|
||||
sock.connect(sockaddr)
|
||||
return sock
|
||||
except OSError as exc:
|
||||
last_error = exc
|
||||
sock.close()
|
||||
if last_error is not None:
|
||||
raise last_error
|
||||
raise socket.gaierror(socket.EAI_NONAME, f"{host} did not resolve to an address")
|
||||
|
||||
Reference in New Issue
Block a user