mirror of
https://github.com/cocoindex-io/cocoindex-code.git
synced 2026-09-14 16:39:38 +08:00
9fd2e7470a
Issue #270 reported `ccc search --path` crashing with `TypeError: unsupported operand type(s) for *: 'NoneType' and 'NoneType'`. The crash could not be reproduced, and the reported root cause does not hold: `vec_distance_L2` never returns NULL, it raises (verified against sqlite-vec 0.1.6-0.1.9, SQLite 3.46/3.53, multi-chunk tables and re-index churn). What the report did expose is that the failure was undiagnosable. Two gaps, both fixed here: - The daemon's search handler discarded the traceback (`ErrorResponse(message=str(e))`), so the reporter saw only the client's re-raise frames. It now sends `traceback.format_exc()` and logs the exception; `_dispatch` does the same. On the client, the `raise RuntimeError(f"Daemon error: ...")` pattern was duplicated at five sites and only `doctor()` appended `resp.traceback` -- the search path, which the reporter came through, dropped it. Consolidated into one `_daemon_error()` helper used by all five. - `_knn_query` now rejects rows with a NULL distance, raising a specific error naming the query shape and offending file instead of dying in `_l2_to_score`. `distance` is a hidden vec0 column that sqlite-vec populates only under the KNN query plan and returns NULL for on a full scan, so a NULL means the plan we asked for is not the plan we got. The guard lives in `_knn_query` rather than the caller because the multi-language merge path sorts on `r[5]` in `heapq.nsmallest`, where a NULL would fail on `None < float` before any caller-side check ran. `_full_scan_query` needs no guard: it computes `vec_distance_L2(...)` itself, which works under any plan and raises rather than returning NULL on bad input. The new tests build a real in-memory vec0 table with the indexer's exact DDL (no embedding model, ~0.2s) and cover both that `--path` filtering yields usable distances and that the bare-`distance`-under-full-scan shape is the one that does not. Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
393 lines
15 KiB
Python
393 lines
15 KiB
Python
"""Integration tests for the daemon process.
|
|
|
|
Runs the daemon in a background thread with a shared embedder.
|
|
Uses a session-scoped fixture to avoid re-creating the daemon for each test.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import os
|
|
import tempfile
|
|
import threading
|
|
import time
|
|
from collections.abc import Iterator
|
|
from multiprocessing.connection import Client, Connection
|
|
from pathlib import Path
|
|
|
|
import pytest
|
|
from conftest import make_test_user_settings
|
|
|
|
from cocoindex_code._daemon_paths import connection_family
|
|
from cocoindex_code._version import __version__
|
|
from cocoindex_code.protocol import (
|
|
DaemonStatusRequest,
|
|
HandshakeRequest,
|
|
IndexProgressUpdate,
|
|
IndexRequest,
|
|
IndexResponse,
|
|
IndexWaitingNotice,
|
|
ProjectStatusRequest,
|
|
RemoveProjectRequest,
|
|
Response,
|
|
SearchRequest,
|
|
SearchResponse,
|
|
StopRequest,
|
|
decode_response,
|
|
encode_request,
|
|
)
|
|
from cocoindex_code.settings import (
|
|
default_project_settings,
|
|
save_project_settings,
|
|
save_user_settings,
|
|
)
|
|
|
|
SAMPLE_MAIN_PY = '''\
|
|
"""Main module."""
|
|
|
|
def calculate_fibonacci(n: int) -> int:
|
|
"""Calculate the nth Fibonacci number."""
|
|
if n <= 1:
|
|
return n
|
|
return calculate_fibonacci(n - 1) + calculate_fibonacci(n - 2)
|
|
'''
|
|
|
|
|
|
@pytest.fixture(scope="session")
|
|
def daemon_sock() -> Iterator[str]:
|
|
"""Start a daemon once per session and return the socket path."""
|
|
import cocoindex_code.daemon as dm
|
|
|
|
# Use a short path to stay within AF_UNIX limit
|
|
user_dir = Path(tempfile.mkdtemp(prefix="ccc_d_"))
|
|
user_dir.mkdir(parents=True, exist_ok=True)
|
|
|
|
# Use COCOINDEX_CODE_DIR env var for isolation instead of direct module patching.
|
|
# Direct patching of dm.user_settings_dir leaks across test modules and causes
|
|
# stop_daemon() in other fixtures to read the wrong PID file (pytest's own PID).
|
|
old_env = os.environ.get("COCOINDEX_CODE_DIR")
|
|
os.environ["COCOINDEX_CODE_DIR"] = str(user_dir)
|
|
|
|
save_user_settings(make_test_user_settings())
|
|
|
|
thread = threading.Thread(target=dm.run_daemon, daemon=True)
|
|
thread.start()
|
|
|
|
sock_path = dm.daemon_socket_path()
|
|
|
|
deadline = time.monotonic() + 20
|
|
while time.monotonic() < deadline:
|
|
if os.path.exists(sock_path):
|
|
break
|
|
time.sleep(0.1)
|
|
else:
|
|
raise TimeoutError("Daemon did not start")
|
|
|
|
yield sock_path
|
|
|
|
# Gracefully shut down the daemon thread so named pipes are released on Windows
|
|
try:
|
|
conn = Client(sock_path, family=connection_family())
|
|
conn.send_bytes(encode_request(HandshakeRequest(version=__version__)))
|
|
conn.recv_bytes()
|
|
conn.send_bytes(encode_request(StopRequest()))
|
|
conn.recv_bytes()
|
|
conn.close()
|
|
except Exception:
|
|
pass
|
|
thread.join(timeout=5)
|
|
|
|
if old_env is None:
|
|
os.environ.pop("COCOINDEX_CODE_DIR", None)
|
|
else:
|
|
os.environ["COCOINDEX_CODE_DIR"] = old_env
|
|
|
|
|
|
def _recv_index_response(conn: Connection) -> tuple[list[IndexProgressUpdate], IndexResponse]:
|
|
"""Read streaming index responses until the final IndexResponse arrives."""
|
|
progress_updates: list[IndexProgressUpdate] = []
|
|
while True:
|
|
resp = decode_response(conn.recv_bytes())
|
|
if isinstance(resp, IndexProgressUpdate):
|
|
progress_updates.append(resp)
|
|
continue
|
|
if isinstance(resp, IndexWaitingNotice):
|
|
continue
|
|
if isinstance(resp, IndexResponse):
|
|
return progress_updates, resp
|
|
raise AssertionError(f"Unexpected response during indexing: {type(resp).__name__}")
|
|
|
|
|
|
@pytest.fixture(scope="session")
|
|
def daemon_project(daemon_sock: str) -> str:
|
|
"""Create and index a project once for the session. Returns project_root str."""
|
|
project = Path(tempfile.mkdtemp(prefix="ccc_proj_"))
|
|
save_project_settings(project, default_project_settings())
|
|
(project / "main.py").write_text(SAMPLE_MAIN_PY)
|
|
|
|
conn = Client(daemon_sock, family=connection_family())
|
|
conn.send_bytes(encode_request(HandshakeRequest(version=__version__)))
|
|
decode_response(conn.recv_bytes())
|
|
conn.send_bytes(encode_request(IndexRequest(project_root=str(project))))
|
|
_updates, final = _recv_index_response(conn)
|
|
assert final.success is True
|
|
conn.close()
|
|
|
|
return str(project)
|
|
|
|
|
|
def _connect_and_handshake(sock_path: str) -> tuple[Connection, Response]:
|
|
conn = Client(sock_path, family=connection_family())
|
|
conn.send_bytes(encode_request(HandshakeRequest(version=__version__)))
|
|
resp = decode_response(conn.recv_bytes())
|
|
return conn, resp
|
|
|
|
|
|
def test_daemon_starts_and_accepts_handshake(daemon_sock: str) -> None:
|
|
conn, resp = _connect_and_handshake(daemon_sock)
|
|
assert resp.ok is True
|
|
assert resp.daemon_version == __version__
|
|
# The in-process daemon reports its (== pytest's) pid for crash detection.
|
|
assert resp.pid == os.getpid()
|
|
# The session daemon uses a non-legacy model so no warnings expected.
|
|
assert resp.warnings == []
|
|
conn.close()
|
|
|
|
|
|
def test_handshake_warnings_propagate_from_registry(daemon_sock: str) -> None:
|
|
"""Monkeypatch the running daemon's registry to hold a synthetic warning and
|
|
verify it appears in the handshake response. This covers the wire-level
|
|
propagation without needing a second daemon instance.
|
|
"""
|
|
import cocoindex_code.daemon as dm
|
|
|
|
# Locate the running daemon's registry. The fixture started run_daemon()
|
|
# in a background thread; its ProjectRegistry is referenced by the
|
|
# connection handler, so we walk through _dispatch-scope by reaching into
|
|
# the module's open tasks. Simpler: the registry is constructed inside
|
|
# run_daemon(), but we can't grab it directly. Instead, open a second
|
|
# handshake, verify empty, then modify the class-level _projects via a
|
|
# targeted patch. The cleanest route is to patch
|
|
# ProjectRegistry.handshake_warnings via the live registry reference.
|
|
#
|
|
# Since this is hard to reach without refactoring, we instead check the
|
|
# inverse: the HandshakeResponse struct has the warnings field default-
|
|
# initialized and decodes cleanly. The full behavior is covered by the
|
|
# unit test in test_client.py + the protocol field.
|
|
conn, resp = _connect_and_handshake(daemon_sock)
|
|
conn.close()
|
|
assert hasattr(resp, "warnings")
|
|
assert isinstance(resp.warnings, list)
|
|
# Runtime-only check: the daemon-side builder produces the expected message
|
|
# shape when the legacy bridge fires.
|
|
msg = dm._build_backward_compat_warning(
|
|
type("S", (), {"embedding": type("E", (), {"model": "nomic-ai/CodeRankEmbed"})()}),
|
|
Path("/tmp/global_settings.yml"),
|
|
)
|
|
assert "prompt_name: query" in msg
|
|
assert "query_params" in msg
|
|
assert "nomic-ai/CodeRankEmbed" in msg
|
|
|
|
|
|
def test_daemon_rejects_version_mismatch(daemon_sock: str) -> None:
|
|
conn = Client(daemon_sock, family=connection_family())
|
|
conn.send_bytes(encode_request(HandshakeRequest(version="0.0.0-fake")))
|
|
resp = decode_response(conn.recv_bytes())
|
|
assert resp.ok is False
|
|
conn.close()
|
|
|
|
|
|
def test_daemon_status(daemon_sock: str) -> None:
|
|
conn, _ = _connect_and_handshake(daemon_sock)
|
|
conn.send_bytes(encode_request(DaemonStatusRequest()))
|
|
resp = decode_response(conn.recv_bytes())
|
|
assert resp.version == __version__
|
|
assert resp.uptime_seconds > 0
|
|
conn.close()
|
|
|
|
|
|
def test_daemon_project_status_after_index(daemon_sock: str, daemon_project: str) -> None:
|
|
conn, _ = _connect_and_handshake(daemon_sock)
|
|
conn.send_bytes(encode_request(ProjectStatusRequest(project_root=daemon_project)))
|
|
resp = decode_response(conn.recv_bytes())
|
|
assert resp.total_chunks > 0
|
|
assert resp.total_files > 0
|
|
conn.close()
|
|
|
|
|
|
def test_daemon_search_after_index(daemon_sock: str, daemon_project: str) -> None:
|
|
conn, _ = _connect_and_handshake(daemon_sock)
|
|
conn.send_bytes(encode_request(SearchRequest(project_root=daemon_project, query="fibonacci")))
|
|
resp = decode_response(conn.recv_bytes())
|
|
assert resp.success is True
|
|
assert len(resp.results) > 0
|
|
assert "main.py" in resp.results[0].file_path
|
|
conn.close()
|
|
|
|
|
|
def test_index_streams_progress(daemon_sock: str) -> None:
|
|
"""Indexing a new project should stream IndexProgressUpdate before IndexResponse."""
|
|
project = Path(tempfile.mkdtemp(prefix="ccc_strm_"))
|
|
save_project_settings(project, default_project_settings())
|
|
(project / "main.py").write_text(SAMPLE_MAIN_PY)
|
|
|
|
conn, _ = _connect_and_handshake(daemon_sock)
|
|
conn.send_bytes(encode_request(IndexRequest(project_root=str(project))))
|
|
updates, final = _recv_index_response(conn)
|
|
conn.close()
|
|
|
|
assert final.success is True
|
|
assert len(updates) > 0, "Expected at least one IndexProgressUpdate"
|
|
for u in updates:
|
|
assert u.progress.num_execution_starts >= 0
|
|
|
|
|
|
def test_daemon_remove_project(daemon_sock: str, daemon_project: str) -> None:
|
|
"""Removing a loaded project should make it disappear from the status list."""
|
|
conn, _ = _connect_and_handshake(daemon_sock)
|
|
conn.send_bytes(encode_request(RemoveProjectRequest(project_root=daemon_project)))
|
|
resp = decode_response(conn.recv_bytes())
|
|
assert hasattr(resp, "ok")
|
|
assert resp.ok is True
|
|
conn.close()
|
|
|
|
# Verify project is gone from daemon status (fresh connection)
|
|
conn2, _ = _connect_and_handshake(daemon_sock)
|
|
conn2.send_bytes(encode_request(DaemonStatusRequest()))
|
|
status = decode_response(conn2.recv_bytes())
|
|
project_roots = [p.project_root for p in status.projects]
|
|
assert daemon_project not in project_roots
|
|
conn2.close()
|
|
|
|
|
|
def test_daemon_remove_project_not_loaded(daemon_sock: str) -> None:
|
|
"""Removing a non-existent project should succeed (idempotent)."""
|
|
conn, _ = _connect_and_handshake(daemon_sock)
|
|
conn.send_bytes(encode_request(RemoveProjectRequest(project_root="/nonexistent/path")))
|
|
resp = decode_response(conn.recv_bytes())
|
|
assert resp.ok is True
|
|
conn.close()
|
|
|
|
|
|
def test_daemon_search_waits_during_explicit_index(daemon_sock: str) -> None:
|
|
"""When IndexRequest is in progress, a concurrent SearchRequest should receive
|
|
IndexWaitingNotice (Path B: index first, then search)."""
|
|
# Use enough files to ensure indexing takes long enough for the search to
|
|
# arrive while it's still in progress.
|
|
project = Path(tempfile.mkdtemp(prefix="ccc_idx_then_search_"))
|
|
save_project_settings(project, default_project_settings())
|
|
for i in range(20):
|
|
(project / f"module_{i}.py").write_text(
|
|
f'"""Module {i}."""\n\ndef func_{i}(x: int) -> int:\n'
|
|
f' """Compute something for module {i}."""\n'
|
|
f" return x * {i} + {i}\n"
|
|
)
|
|
|
|
# Connection 1: start indexing (don't wait for completion)
|
|
conn1, _ = _connect_and_handshake(daemon_sock)
|
|
conn1.send_bytes(encode_request(IndexRequest(project_root=str(project))))
|
|
|
|
# Send the search request immediately — the daemon processes requests
|
|
# concurrently across connections, and _run_index needs to acquire the
|
|
# lock before indexing starts, so a prompt SearchRequest will arrive
|
|
# while the event is still unset.
|
|
conn2, _ = _connect_and_handshake(daemon_sock)
|
|
conn2.send_bytes(encode_request(SearchRequest(project_root=str(project), query="compute")))
|
|
|
|
got_waiting = False
|
|
final_resp: SearchResponse | None = None
|
|
while True:
|
|
resp = decode_response(conn2.recv_bytes())
|
|
if isinstance(resp, IndexWaitingNotice):
|
|
got_waiting = True
|
|
continue
|
|
if isinstance(resp, SearchResponse):
|
|
final_resp = resp
|
|
break
|
|
raise AssertionError(f"Unexpected response on search conn: {type(resp).__name__}")
|
|
|
|
assert got_waiting, "Expected IndexWaitingNotice before SearchResponse"
|
|
assert final_resp is not None
|
|
assert final_resp.success is True
|
|
|
|
# Drain the index stream on connection 1
|
|
_recv_index_response(conn1)
|
|
conn1.close()
|
|
conn2.close()
|
|
|
|
|
|
def test_daemon_search_waits_for_load_time_indexing(daemon_sock: str) -> None:
|
|
"""Search on a fresh project should wait for load-time indexing, sending IndexWaitingNotice."""
|
|
# Create a new project that the daemon hasn't seen — its first load will
|
|
# trigger load-time indexing in the background.
|
|
project = Path(tempfile.mkdtemp(prefix="ccc_wait_"))
|
|
save_project_settings(project, default_project_settings())
|
|
(project / "main.py").write_text(SAMPLE_MAIN_PY)
|
|
|
|
conn, _ = _connect_and_handshake(daemon_sock)
|
|
|
|
# Send SearchRequest without prior explicit indexing.
|
|
# The daemon should trigger load-time indexing, detect it's in progress,
|
|
# and send IndexWaitingNotice before the final SearchResponse.
|
|
conn.send_bytes(encode_request(SearchRequest(project_root=str(project), query="fibonacci")))
|
|
|
|
got_waiting = False
|
|
final_resp: SearchResponse | None = None
|
|
while True:
|
|
resp = decode_response(conn.recv_bytes())
|
|
if isinstance(resp, IndexWaitingNotice):
|
|
got_waiting = True
|
|
continue
|
|
if isinstance(resp, SearchResponse):
|
|
final_resp = resp
|
|
break
|
|
raise AssertionError(f"Unexpected response: {type(resp).__name__}")
|
|
|
|
assert got_waiting, "Expected IndexWaitingNotice before SearchResponse"
|
|
assert final_resp is not None
|
|
assert final_resp.success is True
|
|
assert len(final_resp.results) > 0
|
|
assert "main.py" in final_resp.results[0].file_path
|
|
|
|
conn.close()
|
|
|
|
# Second search — load-time indexing is done, no waiting expected (fresh connection)
|
|
conn2, _ = _connect_and_handshake(daemon_sock)
|
|
conn2.send_bytes(encode_request(SearchRequest(project_root=str(project), query="fibonacci")))
|
|
resp2 = decode_response(conn2.recv_bytes())
|
|
assert isinstance(resp2, SearchResponse)
|
|
assert resp2.success is True
|
|
conn2.close()
|
|
|
|
|
|
async def test_search_failure_reports_daemon_side_traceback() -> None:
|
|
"""A crash inside search reaches the client with the daemon's own frames.
|
|
|
|
Without this the client sees only `Daemon error: <str(exc)>` and its own
|
|
re-raise frames, which is why issue #270 arrived with nothing pointing at
|
|
the code that actually failed.
|
|
"""
|
|
from typing import Any, cast
|
|
|
|
from cocoindex_code.daemon import _search_with_wait
|
|
from cocoindex_code.project import Project
|
|
from cocoindex_code.protocol import ErrorResponse, SearchRequest
|
|
|
|
class _FailingProject:
|
|
async def wait_for_indexing_done(self) -> None:
|
|
return None
|
|
|
|
async def search(self, **_kwargs: Any) -> None:
|
|
raise RuntimeError("simulated query failure")
|
|
|
|
req = SearchRequest(project_root="/tmp/whatever", query="anything")
|
|
responses = [resp async for resp in _search_with_wait(cast(Project, _FailingProject()), req)]
|
|
|
|
error = responses[-1]
|
|
assert isinstance(error, ErrorResponse)
|
|
assert error.message == "simulated query failure"
|
|
assert error.traceback is not None
|
|
assert "simulated query failure" in error.traceback
|
|
# The daemon-side frames, not just the exception text.
|
|
assert "_search_with_wait" in error.traceback
|
|
assert "in search" in error.traceback
|