diff --git a/bin/sync.py b/bin/sync.py index a3bef49..ab56a4d 100644 --- a/bin/sync.py +++ b/bin/sync.py @@ -18,6 +18,7 @@ bin/db-compare.sh could never work there. A fresh machine clones first (git clone nexus-core), then runs this. """ import argparse +import contextlib import os import re import shutil @@ -162,7 +163,7 @@ def dump_db() -> bool: if not DB.exists(): return False try: - with sqlite3.connect(f"file:{DB}?mode=ro", uri=True) as conn: + with contextlib.closing(sqlite3.connect(f"file:{DB}?mode=ro", uri=True)) as conn: ext = _vec0_extension() if ext: try: @@ -232,9 +233,12 @@ def compare(db: Path = DB, dump: Path = DB_SQL) -> str: if not db.exists(): return "no-live" try: - with sqlite3.connect(f"file:{db}?mode=ro", uri=True) as live_conn: + # closing(), not sqlite3's context manager: that one commits without + # closing, and restore_db() unlinks the DB right after calling this - + # Windows fails that unlink while any handle is still open. + with contextlib.closing(sqlite3.connect(f"file:{db}?mode=ro", uri=True)) as live_conn: live = _state(live_conn) - with sqlite3.connect(":memory:") as dump_conn: + with contextlib.closing(sqlite3.connect(":memory:")) as dump_conn: dump_conn.executescript(dump.read_text(encoding="utf-8")) backup = _state(dump_conn) except (sqlite3.Error, OSError): @@ -287,7 +291,7 @@ def restore_db() -> None: for suffix in ("", "-wal", "-shm"): Path(str(DB) + suffix).unlink(missing_ok=True) try: - with sqlite3.connect(DB) as conn: + with contextlib.closing(sqlite3.connect(DB)) as conn, conn: conn.executescript(DB_SQL.read_text(encoding="utf-8")) print("Memory DB restored (conversations + history + facts).") except sqlite3.Error as exc: diff --git a/synapse/memory/store.py b/synapse/memory/store.py index 297b557..da1276e 100644 --- a/synapse/memory/store.py +++ b/synapse/memory/store.py @@ -251,6 +251,8 @@ class PersistentMemoryStore: "DELETE FROM settings WHERE key IN ('anthropic_api_key', 'escalation_model')" ) + self._sweep_orphan_msg_vectors(conn) + conn.commit() conn.close() @@ -788,6 +790,29 @@ class PersistentMemoryStore: except Exception: pass + def _sweep_orphan_msg_vectors(self, conn) -> None: + """One-time repair for databases written before delete_conversation + cleaned up after itself: drop vectors whose message is already gone.""" + try: + ids = [r["message_id"] for r in conn.execute( + "SELECT v.message_id FROM message_vectors v " + "LEFT JOIN messages m ON m.id = v.message_id WHERE m.id IS NULL" + ).fetchall()] + if ids: + conn.execute( + "DELETE FROM message_vectors WHERE message_id NOT IN " + "(SELECT id FROM messages)" + ) + if self.vec_enabled and conn.execute( + "SELECT 1 FROM sqlite_master WHERE name = 'vec_messages'" + ).fetchone(): + for message_id in ids: + conn.execute( + "DELETE FROM vec_messages WHERE rowid = ?", (message_id,) + ) + except Exception: + pass + def _backfill_vec_msgs(self, conn, dim: int) -> None: """Index any message_vectors rows missing from vec_messages.""" try: @@ -1215,4 +1240,4 @@ class PersistentMemoryStore: from ..nexus_config import MEMORY_DB DB_PATH = MEMORY_DB -store = PersistentMemoryStore(DB_PATH) \ No newline at end of file +store = PersistentMemoryStore(DB_PATH) diff --git a/synapse/nexus_config.py b/synapse/nexus_config.py index 841ae70..4615248 100644 --- a/synapse/nexus_config.py +++ b/synapse/nexus_config.py @@ -109,6 +109,36 @@ def path(name: str) -> Path: raise KeyError(f"Unknown config path name: {name}") # --- Settings class and exported instance --- +def _normalize_ollama_host(raw: str) -> str: + """Turn an OLLAMA_HOST value into a URL a client can actually connect to. + + OLLAMA_HOST is Ollama's *server bind* variable, and the common way to expose + Ollama on a LAN is `OLLAMA_HOST=0.0.0.0:11434`. Taken literally as a client + base URL that is unusable twice over: 0.0.0.0 means "every local interface" + to a listener but is not a destination, and there is no scheme for httpx to + parse. The result was a silent empty model list, because list_models() + catches everything and returns []. + + So: supply the scheme when it's missing, and rewrite wildcard binds to + loopback. An explicit host is left alone — someone pointing at a real remote + Ollama means it. + """ + host = (raw or "").strip().rstrip("/") + if not host: + return "http://127.0.0.1:11434" + if "://" not in host: + host = f"http://{host}" + scheme, _, rest = host.partition("://") + hostport = rest.split("/", 1)[0] + name, sep, port = hostport.rpartition(":") + if not sep: # no port given + name, port = hostport, "" + # 0.0.0.0 and :: are bind-any; from a client they mean "this machine". + if name.strip("[]") in ("0.0.0.0", "::", ""): + name = "127.0.0.1" + return f"{scheme}://{name}:{port}" if port else f"{scheme}://{name}" + + class Settings: """ Lightweight settings container. Use `settings` instance for runtime access, @@ -131,8 +161,15 @@ class Settings: self.ollama_log: Path = OLLAMA_LOG self.chat_log: Path = CHAT_LOG - # Env overrides - self.ollama_host: str = os.getenv("OLLAMA_HOST", "http://127.0.0.1:11434") + # Env overrides. Two values from one variable, because OLLAMA_HOST means + # two different things: where a server should LISTEN, and where a client + # should CONNECT. `ollama_bind` keeps the user's literal intent for a + # serve we spawn (0.0.0.0 to expose it on the LAN); `ollama_host` is the + # connectable form for our own requests. + self.ollama_bind: str = os.getenv("OLLAMA_HOST", "") or "127.0.0.1:11434" + self.ollama_host: str = _normalize_ollama_host( + os.getenv("OLLAMA_HOST", "http://127.0.0.1:11434") + ) self.ollama_timeout: int = int(os.getenv("OLLAMA_TIMEOUT", "120")) def as_dict(self) -> Dict[str, Any]: diff --git a/synapse/ollama_manager.py b/synapse/ollama_manager.py index 489b1ea..0281d23 100644 --- a/synapse/ollama_manager.py +++ b/synapse/ollama_manager.py @@ -200,6 +200,140 @@ def _chat_options(temperature: float | None, num_gpu: int | None, num_ctx: int | return opts +_THINK_OPEN = "" +_THINK_CLOSE = "" + + +def _partial_tag_tail(text: str, tag: str) -> int: + """Length of the longest suffix of `text` that could be the start of `tag`. + + A tag can arrive split across stream chunks (""), so that much + of the tail has to be held back rather than emitted. + """ + for n in range(min(len(tag) - 1, len(text)), 0, -1): + if text.endswith(tag[:n]): + return n + return 0 + + +class ThinkStripper: + """Removes a reasoning model's spans from a token stream. + + Ollama routes reasoning into `message.thinking` only when asked to think. + We ask for `think: False` because the reasoning is pure latency here — but + deepseek-r1 and friends emit `` inline in `content` anyway, so the + whole internal monologue reached the chat window, closing tags and all. + + Two shapes show up in practice. A well-formed span is dropped whole. A + *stray* closing tag with no opening — which is what actually shipped — is at + least removed, so the user does not see a literal `` in the reply. + Text already streamed before it cannot be recalled; `strip_think()` handles + that case properly for callers that have the complete message. + """ + + def __init__(self) -> None: + self._buf = "" + self._inside = False + + def feed(self, chunk: str) -> str: + self._buf += chunk + out: list[str] = [] + while True: + if self._inside: + end = self._buf.find(_THINK_CLOSE) + if end == -1: + keep = _partial_tag_tail(self._buf, _THINK_CLOSE) + self._buf = self._buf[len(self._buf) - keep:] if keep else "" + break + self._buf = self._buf[end + len(_THINK_CLOSE):] + self._inside = False + continue + + start = self._buf.find(_THINK_OPEN) + stray = self._buf.find(_THINK_CLOSE) + # A stray close before any open: drop the tag, keep going. + if stray != -1 and (start == -1 or stray < start): + out.append(self._buf[:stray]) + self._buf = self._buf[stray + len(_THINK_CLOSE):] + continue + if start == -1: + keep = max( + _partial_tag_tail(self._buf, _THINK_OPEN), + _partial_tag_tail(self._buf, _THINK_CLOSE), + ) + if keep: + out.append(self._buf[:len(self._buf) - keep]) + self._buf = self._buf[len(self._buf) - keep:] + else: + out.append(self._buf) + self._buf = "" + break + out.append(self._buf[:start]) + self._buf = self._buf[start + len(_THINK_OPEN):] + self._inside = True + return "".join(out) + + def flush(self) -> str: + """Whatever is left once the stream ends. An unterminated is + reasoning that never closed, so it is dropped rather than shown.""" + rest = "" if self._inside else self._buf + self._buf = "" + return rest + + +def strip_think(text: str) -> str: + """Remove reasoning from a complete message. + + Unlike the streaming case this sees everything, so an unmatched closing tag + can be handled the way it was meant: every token before it was reasoning, + and the answer is what follows. + """ + if not text or _THINK_CLOSE not in text and _THINK_OPEN not in text: + return text + import re + cleaned = re.sub(r".*?", "", text, flags=re.S) + if _THINK_CLOSE in cleaned: # stray close: the answer follows it + cleaned = cleaned.rsplit(_THINK_CLOSE, 1)[1] + cleaned = re.sub(r".*\Z", "", cleaned, flags=re.S) # never closed + return cleaned.strip() + + +async def _raise_for_ollama(r: httpx.Response) -> None: + """raise_for_status(), but say what Ollama actually said. + + Ollama answers every failure with {"error": "..."} — a model that isn't + pulled, a request too large for VRAM, a cloud model retired upstream — and + httpx's default message discards the body, leaving the user with: + + Client error '410 Gone' for url 'http://127.0.0.1:11434/api/chat' + + when the body held "glm-4.6 was retired at 2026-06-16". Same exception type + as before so existing handlers are unaffected; only the message improves. + """ + if r.is_success: + return + # A streamed response has no body loaded yet; reading it is what makes the + # error message available at all. + try: + await r.aread() + except Exception: + pass + detail = "" + try: + body = r.json() + if isinstance(body, dict): + detail = str(body.get("error") or "").strip() + except Exception: + detail = (r.text or "").strip() + if not detail: + r.raise_for_status() # nothing to add — keep httpx's wording + raise httpx.HTTPStatusError( + f"Ollama {r.status_code} from {r.request.url.path}: {detail[:400]}", + request=r.request, + response=r, + ) + + class OllamaManager: def __init__(self, runtime_dir=None): self.process = None @@ -258,7 +392,9 @@ class OllamaManager: duplicated verbatim in two methods, so the Windows carve-out below had to be fixed in both places or the two paths would disagree.""" env = os.environ.copy() - env["OLLAMA_HOST"] = self._api_base + # The bind value, not the connect value: a user who set 0.0.0.0 to reach + # Ollama from another machine must still get a server that listens there. + env["OLLAMA_HOST"] = settings.ollama_bind # Every platform uses the project's own model store. Windows used to be # exempt, because the installer pulled with a bare `ollama pull` into # %USERPROFILE%\.ollama and a NexusOS-spawned serve pointed elsewhere @@ -459,7 +595,7 @@ class OllamaManager: elapsed = time.perf_counter() - start _log.info("generate completed model=%s status=%d elapsed=%.3fs", model, r.status_code, elapsed) - r.raise_for_status() + await _raise_for_ollama(r) return r.json().get("response", "") except Exception as e: @@ -484,7 +620,7 @@ class OllamaManager: "stream": True, }), ) as response: - response.raise_for_status() + await _raise_for_ollama(response) async for line in response.aiter_lines(): if not line.strip(): continue @@ -547,8 +683,12 @@ class OllamaManager: async with httpx.AsyncClient(timeout=300.0) as client: r = await client.post(f"{self._api_base}/api/chat", json=body) elapsed = time.perf_counter() - start - r.raise_for_status() + await _raise_for_ollama(r) message = r.json().get("message", {}) + # A reasoning model puts its monologue in `content` even with + # think off, so strip it before anyone reads the answer. + if isinstance(message, dict) and message.get("content"): + message["content"] = strip_think(message["content"]) # Tool callers need the whole message (tool_calls); others want content. return message if tools else message.get("content", "") except Exception as e: @@ -576,7 +716,7 @@ class OllamaManager: f"{self._api_base}/api/embeddings", json=body, ) - r.raise_for_status() + await _raise_for_ollama(r) vec = r.json().get("embedding") return vec if vec else None except Exception as e: @@ -588,7 +728,7 @@ class OllamaManager: try: async with httpx.AsyncClient(timeout=5.0) as client: r = await client.get(f"{self._api_base}/api/tags") - r.raise_for_status() + await _raise_for_ollama(r) return [m["name"] for m in r.json().get("models", [])] except Exception: return [] @@ -629,7 +769,7 @@ class OllamaManager: try: async with httpx.AsyncClient(timeout=10.0) as client: r = await client.post(f"{self._api_base}/api/show", json={"model": model}) - r.raise_for_status() + await _raise_for_ollama(r) info = r.json().get("model_info", {}) or {} block_count = next( (v for k, v in info.items() if k.endswith(".block_count")), None @@ -679,7 +819,8 @@ class OllamaManager: f"{self._api_base}/api/chat", json=body, ) as response: - response.raise_for_status() + await _raise_for_ollama(response) + thinking = ThinkStripper() async for line in response.aiter_lines(): if not line.strip(): continue @@ -687,8 +828,13 @@ class OllamaManager: data = json.loads(line) token = data.get("message", {}).get("content", "") if token: - yield token + token = thinking.feed(token) + if token: + yield token if data.get("done", False): + tail = thinking.flush() # held-back partial tag + if tail: + yield tail elapsed = time.perf_counter() - start _log.info("chat stream completed model=%s elapsed=%.3fs", model, elapsed) eval_count = data.get("eval_count", 0) diff --git a/tests/test_documents.py b/tests/test_documents.py index b3c37a0..5f63702 100644 --- a/tests/test_documents.py +++ b/tests/test_documents.py @@ -136,6 +136,27 @@ def test_conversation_recall_uses_vec_and_matches_brute_force(): asyncio.run(run()) +def test_startup_sweeps_pre_existing_orphan_vectors(): + """Databases written before delete_conversation cleaned up after itself are + repaired the next time the store opens them.""" + import json + path = Path(tempfile.mkdtemp()) / "t.db" + s = PersistentMemoryStore(path) + s.create_conversation("c1") + mid = s.add_message("c1", "user", "lego star wars") + conn = s._connect() + conn.execute("INSERT INTO message_vectors (message_id, embedding) VALUES (?, ?)", + (mid, json.dumps([1.0, 0.0]))) + conn.execute("DELETE FROM messages WHERE id = ?", (mid,)) # the old leaky delete + conn.commit() + conn.close() + + reopened = PersistentMemoryStore(path) + conn = reopened._connect() + assert conn.execute("SELECT COUNT(*) FROM message_vectors").fetchone()[0] == 0 + conn.close() + + def test_projects_scope_documents_and_survive_delete(): s = _store() diff --git a/tests/test_smoke.py b/tests/test_smoke.py index 80314ac..f4a1858 100644 --- a/tests/test_smoke.py +++ b/tests/test_smoke.py @@ -217,16 +217,23 @@ def test_icon_source_requires_real_allowed_file_boundary(tmp_path): def test_ollama_stream_propagates_transport_errors(monkeypatch): + """A failing stream must surface, not be swallowed into an empty reply — + and it must carry Ollama's own explanation, since that is the only part the + user can act on. The response here is a real httpx.Response because the + error path reads the body, which a stubbed raise_for_status never exercised.""" + import httpx + class FailingResponse: async def __aenter__(self): - return self + return httpx.Response( + 503, + json={"error": "Ollama unavailable"}, + request=httpx.Request("POST", "http://127.0.0.1:11434/api/chat"), + ) async def __aexit__(self, *args): return False - def raise_for_status(self): - raise RuntimeError("Ollama unavailable") - class FailingClient: def __init__(self, **kwargs): pass @@ -246,7 +253,7 @@ def test_ollama_stream_propagates_transport_errors(monkeypatch): async for _ in OllamaManager()._chat_stream([], "model", 0): pass - with pytest.raises(RuntimeError, match="Ollama unavailable"): + with pytest.raises(httpx.HTTPStatusError, match="Ollama unavailable"): import asyncio asyncio.run(consume()) @@ -264,13 +271,17 @@ def test_sync_compare_detects_direction(tmp_path): """The guard that stops a stale box from overwriting the other's chats. Both backup and restore refuse to run when this says the wrong thing, so a silent break here loses conversation history.""" + import contextlib import sqlite3 as sq sync = _load_sync() db, dump = tmp_path / "memory.db", tmp_path / "memory.db.sql" def write(rows): db.unlink(missing_ok=True) - with sq.connect(db) as conn: + # closing() then the connection itself: sqlite3's own context manager + # commits but never closes, and Windows refuses to unlink a file that + # still has an open handle. + with contextlib.closing(sq.connect(db)) as conn, conn: # updated_at REAL, matching the production schema in store.py. A TEXT # column here hid a real TypeError for months: the comparison in # _extra() ran str-vs-str in the test and str-vs-float in the field. @@ -547,3 +558,94 @@ def test_preview_iframe_cannot_navigate_to_a_network_url(): assert "encodeURIComponent(doc)" in markdown assert "src={frameUrl}" in markdown assert "srcDoc={doc}" not in markdown + + +def test_ollama_failures_surface_the_reason_not_just_the_status(): + """Ollama answers every failure with {"error": "..."} and httpx's default + message throws it away. A user hitting a retired cloud model saw + "Client error '410 Gone' for url ..." when the body said exactly why.""" + import asyncio + import httpx + import pytest + from synapse.ollama_manager import _raise_for_ollama + + req = httpx.Request("POST", "http://127.0.0.1:11434/api/chat") + + retired = httpx.Response(410, json={"error": "glm-4.6 was retired at 2026-06-16"}, request=req) + with pytest.raises(httpx.HTTPStatusError) as ei: + asyncio.run(_raise_for_ollama(retired)) + assert "retired" in str(ei.value) and "410" in str(ei.value) + + # The common case, not just the exotic one. + missing = httpx.Response(404, json={"error": "model 'foo' not found"}, request=req) + with pytest.raises(httpx.HTTPStatusError) as ei: + asyncio.run(_raise_for_ollama(missing)) + assert "model 'foo' not found" in str(ei.value) + + # No usable body -> keep httpx's own wording rather than inventing one. + blank = httpx.Response(500, content=b"", request=req) + with pytest.raises(httpx.HTTPStatusError): + asyncio.run(_raise_for_ollama(blank)) + + # Success stays silent. + asyncio.run(_raise_for_ollama(httpx.Response(200, json={"ok": True}, request=req))) + + +def test_ollama_host_is_normalized_for_clients_but_not_for_binding(): + """OLLAMA_HOST is Ollama's *bind* variable, and `OLLAMA_HOST=0.0.0.0:11434` + is the normal way to expose it on a LAN. Used verbatim as a client base URL + it is unusable — no scheme, and 0.0.0.0 is not a destination — and every + request failed into list_models()'s bare `except: return []`, so the model + picker just went empty with no error anywhere.""" + from synapse.nexus_config import _normalize_ollama_host as norm + + assert norm("0.0.0.0:11434") == "http://127.0.0.1:11434" + assert norm("[::]:11434") == "http://127.0.0.1:11434" + assert norm("http://0.0.0.0:11434/") == "http://127.0.0.1:11434" + assert norm("127.0.0.1:11434") == "http://127.0.0.1:11434" # scheme supplied + assert norm("") == "http://127.0.0.1:11434" + # A real remote is deliberate — leave it alone. + assert norm("https://ollama.lan:11434") == "https://ollama.lan:11434" + assert norm("192.168.1.50:11434") == "http://192.168.1.50:11434" + + +def test_spawned_serve_keeps_the_users_bind_address(monkeypatch): + """Normalizing for the client must not quietly un-expose a server we spawn.""" + import importlib + from synapse import nexus_config + monkeypatch.setenv("OLLAMA_HOST", "0.0.0.0:11434") + reloaded = importlib.reload(nexus_config) + try: + assert reloaded.settings.ollama_bind == "0.0.0.0:11434" # listens everywhere + assert reloaded.settings.ollama_host == "http://127.0.0.1:11434" # we connect here + finally: + monkeypatch.delenv("OLLAMA_HOST", raising=False) + importlib.reload(nexus_config) + + +def test_think_blocks_never_reach_the_reply(): + """A reasoning model emits inline in `content` even with think off, + and the whole internal monologue reached the chat window — including the + literal closing tags. Streaming has to cope with a tag split across chunks, + and with the shape actually observed: a stray and no opening.""" + from synapse.ollama_manager import ThinkStripper, strip_think + + def stream(chunks): + s = ThinkStripper() + return "".join(s.feed(c) for c in chunks) + s.flush() + + assert stream(["hello ", "", "noise", "", "world"]) == "hello world" + # tag split across chunk boundaries + assert stream(["axb"]) == "ab" + # never closed -> it was all reasoning + assert stream(["keep", "", "runs off the end"]) == "keep" + # stray close, no open: at minimum the tag itself must not be shown + assert "" not in stream(["reasoning...", "", "the answer"]) + # ordinary text is untouched, including angle brackets + assert stream(["a < b ", "and c > d"]) == "a < b and c > d" + + # With the complete message the stray-close case can be handled properly: + # everything before it was reasoning. + assert strip_think("rambling\n\nThe answer") == "The answer" + assert strip_think("abc") == "ac" + assert strip_think("no tags here") == "no tags here"