From d5bd4b3446c853ee151fcf9741a865aab6d445a5 Mon Sep 17 00:00:00 2001 From: davidgut1982 Date: Sun, 31 May 2026 02:17:31 +0000 Subject: [PATCH 1/2] chore: commit deployed runtime (*_local modules, migrations, units, config example) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Captures the actual running implementation that was previously untracked: - src/server_local.py: asyncpg-based MCP server (replaces Supabase server.py) - sse_server_local.py: HTTP/SSE transport wrapper for server_local - daemon/{worker,watcher,chunker,embedder}_local.py: asyncpg-backed indexing pipeline - migrations/020_vector_search_indexer.sql: full schema (pgvector tables + functions) - systemd/: three user service units (worker, watcher, http) with path notes - config.yaml.example: redacted template (password → CHANGE_ME, placeholder paths) - .gitignore: adds config.yaml, *.bak*, embed_cache.json exclusions - requirements.txt: adds asyncpg>=0.29.0 (required by all _local modules) - SETUP.md: end-to-end install guide (PostgreSQL/pgvector, venv, config, units, MCP entry) No secrets committed. config.yaml remains gitignored. Co-Authored-By: Claude Opus 4.8 (1M context) --- .gitignore | 8 +- SETUP.md | 214 ++++++++++++++++ config.yaml.example | 35 +++ daemon/chunker_local.py | 76 ++++++ daemon/embedder_local.py | 23 ++ daemon/watcher_local.py | 113 ++++++++ daemon/worker_local.py | 139 ++++++++++ migrations/020_vector_search_indexer.sql | 146 +++++++++++ requirements.txt | 3 + src/server_local.py | 312 ++++++++++++++++++++++- sse_server_local.py | 105 ++++++++ systemd/vector-indexer-mcp-http.service | 22 ++ systemd/vector-indexer-watcher.service | 22 ++ systemd/vector-indexer-worker.service | 22 ++ 14 files changed, 1232 insertions(+), 8 deletions(-) create mode 100644 SETUP.md create mode 100644 config.yaml.example create mode 100644 daemon/chunker_local.py create mode 100644 daemon/embedder_local.py create mode 100644 daemon/watcher_local.py create mode 100644 daemon/worker_local.py create mode 100644 migrations/020_vector_search_indexer.sql create mode 100644 sse_server_local.py create mode 100644 systemd/vector-indexer-mcp-http.service create mode 100644 systemd/vector-indexer-watcher.service create mode 100644 systemd/vector-indexer-worker.service diff --git a/.gitignore b/.gitignore index 2d7a4d1..096895e 100644 --- a/.gitignore +++ b/.gitignore @@ -1,6 +1,11 @@ # Environment files .env +# Runtime config (secrets) — copy config.yaml.example → config.yaml and fill in +config.yaml +config.yaml.bak* +*.bak* + # Python __pycache__/ *.py[cod] @@ -34,10 +39,11 @@ env/ *.swp *.swo -# Model cache +# Model cache / embedding cache .cache/ models/ sentence_transformers/ +embed_cache.json # Testing .pytest_cache/ diff --git a/SETUP.md b/SETUP.md new file mode 100644 index 0000000..6cc5518 --- /dev/null +++ b/SETUP.md @@ -0,0 +1,214 @@ +# Vector Indexer MCP — Deployment Setup Guide + +This document describes the **actual running deployment** — the `*_local.py` modules, +migrations, and systemd units that ship in this repo. The original `src/server.py` / +`daemon/*.py` modules (Supabase backend) remain in tree but are not used by the +production deployment. + +--- + +## Prerequisites + +- Python 3.10+ +- PostgreSQL 14+ with **pgvector** extension +- systemd (user session) for service management + +--- + +## 1. PostgreSQL / pgvector + +### Install pgvector (if not already installed) + +```bash +# Debian/Ubuntu +sudo apt install postgresql-14-pgvector # or postgresql-15-pgvector + +# Build from source (any distro) +git clone https://github.com/pgvector/pgvector.git +cd pgvector && make && sudo make install +``` + +### Create role and database + +```sql +-- as postgres superuser +CREATE ROLE vectoruser WITH LOGIN PASSWORD 'your-strong-password'; +CREATE DATABASE vectorindex OWNER vectoruser; + +-- connect to vectorindex and enable the extension +\c vectorindex +CREATE EXTENSION IF NOT EXISTS vector; +GRANT ALL ON SCHEMA public TO vectoruser; +``` + +### Apply the migration + +```bash +psql -U vectoruser -d vectorindex -f migrations/020_vector_search_indexer.sql +``` + +This creates: +- `file_metadata`, `file_chunks`, `file_embeddings`, `index_queue`, `index_stats` tables +- GIN index for full-text search +- HNSW index for vector cosine similarity +- `search_hybrid()` and `get_index_health()` stored functions + +--- + +## 2. Python environment + +```bash +cd /path/to/vector-indexer-mcp + +python3 -m venv venv +source venv/bin/activate + +pip install -r requirements.txt +``` + +Key runtime dependencies (all in `requirements.txt`): +- `asyncpg` — PostgreSQL driver used by all `*_local.py` modules +- `sentence-transformers` + `torch` — embedding generation +- `watchdog` — file system monitoring +- `mcp` — MCP server protocol +- `uvicorn`, `starlette`, `sse-starlette` — HTTP/SSE transport (`sse_server_local.py`) + +> **Note on pyproject.toml**: The `[project.scripts]` entries (`vector-indexer-worker`, +> `vector-indexer-mcp`) point at the old Supabase modules (`daemon.worker:main`, +> `src.server:main`). Do **not** use them for the local deployment — use the systemd +> units or the module invocations below instead. + +--- + +## 3. Configuration + +```bash +cp config.yaml.example config.yaml +# Edit config.yaml and set: +# database.password → your PostgreSQL password (same as step 1) +# watcher.paths → list of absolute directory paths to monitor +``` + +**Important**: `config.yaml` is gitignored and must never be committed — it contains +your database password. The modules load it from a **hardcoded path**: + +``` +/home/david/vector-indexer-mcp/config.yaml # in src/server_local.py +/home/david/vector-indexer-mcp/config.yaml # in daemon/worker_local.py +/home/david/vector-indexer-mcp/config.yaml # in daemon/watcher_local.py +``` + +This path is hardcoded in the source (known limitation — see Recommended Follow-ups). +For a different checkout location, update `CONFIG_PATH` in each `*_local.py` file or +create a symlink. + +### watcher.paths gotcha + +The `watcher.paths` list accepts **any** absolute path, including paths outside the +project root. The daemon watches each listed directory recursively. Paths that do not +exist at startup are skipped with a warning (they can be added later by restarting the +service). + +--- + +## 4. Systemd user services + +The `systemd/` directory contains the three unit files. Copy them to your user systemd +directory and adjust the `WorkingDirectory` and `ExecStart` paths to match your +checkout location before enabling. + +```bash +# Adjust paths inside each .service file first, then: +cp systemd/vector-indexer-worker.service ~/.config/systemd/user/ +cp systemd/vector-indexer-watcher.service ~/.config/systemd/user/ +cp systemd/vector-indexer-mcp-http.service ~/.config/systemd/user/ + +systemctl --user daemon-reload + +# Start worker first (embedder / queue processor) +systemctl --user enable --now vector-indexer-worker + +# Start file watcher (enqueues changed files) +systemctl --user enable --now vector-indexer-watcher + +# Start HTTP/SSE server (optional — for non-stdio MCP clients) +systemctl --user enable --now vector-indexer-mcp-http +``` + +Check status / logs: + +```bash +systemctl --user status vector-indexer-worker +journalctl --user -u vector-indexer-worker -f +``` + +### Service startup order + +`vector-indexer-worker` → `vector-indexer-mcp-http` (declared in `After=`). +`vector-indexer-watcher` is independent of the other two. + +--- + +## 5. Claude MCP stdio entry + +For direct stdio usage (the typical Claude Code / claude-mpm setup), add this to your +`mcpServers` configuration: + +```json +{ + "vector-indexer-mcp": { + "command": "/path/to/vector-indexer-mcp/venv/bin/python", + "args": ["-m", "src.server_local"], + "cwd": "/path/to/vector-indexer-mcp" + } +} +``` + +The server reads `config.yaml` at import time (hardcoded path — see note above), so +the `cwd` must match the checkout location, **or** the `CONFIG_PATH` constant in +`src/server_local.py` must be updated. + +--- + +## 6. Initial bulk index + +After the services are running, queue all files in a watched directory for indexing: + +```bash +# Via MCP tool (from a Claude session): +# search_hybrid → reindex_path(path="/your/project", recursive=true, force=true) + +# Or directly via psql: +INSERT INTO index_queue(file_path, event_type, status) +SELECT file_path, 'modify', 'pending' +FROM file_metadata; +``` + +Or use `scripts/bulk_index.py` if it exists. + +--- + +## 7. Verify the index + +```bash +source venv/bin/activate +python scripts/verify_index.py +``` + +Or via MCP `index_status` tool — returns `total_files`, `total_chunks`, +`total_embeddings`, and `queue.pending`. + +--- + +## Recommended Follow-ups (not done in this reconciliation) + +1. **De-hardcode `CONFIG_PATH`** — read from `VECTOR_INDEXER_CONFIG` env var with a + fallback, so the service works from any checkout location without editing source. +2. **Fix `pyproject.toml` scripts** — `vector-indexer-worker` and `vector-indexer-mcp` + point at the old Supabase modules; update them to point at `daemon.worker_local:main` + and `src.server_local:main` (or remove them and rely on the systemd units). +3. **Prune dead dependencies** — `supabase`, `postgrest-py`, and `tiktoken` are listed + in `requirements.txt` / `pyproject.toml` but not used by the `*_local.py` runtime. + Removing them reduces install size and eliminates a large dependency surface. +4. **Add `asyncpg` to `pyproject.toml` dependencies** — currently only in + `requirements.txt`; the two should stay in sync. diff --git a/config.yaml.example b/config.yaml.example new file mode 100644 index 0000000..e1f4285 --- /dev/null +++ b/config.yaml.example @@ -0,0 +1,35 @@ +database: + host: localhost + port: 5432 + name: vectorindex + user: vectoruser + # Replace with your actual database password — never commit the real value + password: CHANGE_ME + +watcher: + # List the absolute paths you want the watcher to monitor. + # These are machine-specific; adjust for your environment. + # Note: paths outside the project root (e.g. /home/user/other-repo) are + # supported — the watcher follows them correctly. + paths: + - /path/to/your/project-one + - /path/to/your/project-two + debounce_ms: 500 + extensions: [.py, .md, .txt, .yaml, .yml, .json, .toml, .sh, .sql, .ts, .js, .tsx, .jsx] + exclude_patterns: [__pycache__, .git, node_modules, .venv, venv, .cache, dist, build, .next, __snapshots__] + +embedding: + model: paraphrase-multilingual-MiniLM-L12-v2 + batch_size: 32 + device: cpu # set to "cuda" if a GPU is available + +chunking: + chunk_size_tokens: 500 + overlap_tokens: 50 + +queue: + poll_interval_ms: 500 + batch_size: 50 + +server: + port: 5577 diff --git a/daemon/chunker_local.py b/daemon/chunker_local.py new file mode 100644 index 0000000..565278f --- /dev/null +++ b/daemon/chunker_local.py @@ -0,0 +1,76 @@ +"""Token-aware text chunker for local asyncpg-based vector indexer.""" +from dataclasses import dataclass +from typing import List + + +@dataclass +class TextChunk: + text: str + start_line: int + end_line: int + chunk_index: int + token_count: int + + +class TextChunker: + def __init__(self, chunk_size_tokens: int = 500, overlap_tokens: int = 50): + self.chunk_size = chunk_size_tokens + self.overlap = overlap_tokens + try: + import tiktoken + self.enc = tiktoken.get_encoding("cl100k_base") + self._use_tiktoken = True + except Exception: + self._use_tiktoken = False + + def _count_tokens(self, text: str) -> int: + if self._use_tiktoken: + return len(self.enc.encode(text)) + return len(text.split()) + + def chunk_text(self, text: str, file_path: str = "") -> List[TextChunk]: + lines = text.splitlines(keepends=True) + chunks = [] + current_lines = [] + current_tokens = 0 + chunk_idx = 0 + start_line = 1 + + for line_num, line in enumerate(lines, 1): + line_tokens = self._count_tokens(line) + if current_tokens + line_tokens > self.chunk_size and current_lines: + chunk_text = "".join(current_lines) + chunks.append(TextChunk( + text=chunk_text, + start_line=start_line, + end_line=line_num - 1, + chunk_index=chunk_idx, + token_count=current_tokens + )) + chunk_idx += 1 + # Overlap: keep last N tokens worth of lines + overlap_lines = [] + overlap_tokens = 0 + for prev_line in reversed(current_lines): + t = self._count_tokens(prev_line) + if overlap_tokens + t > self.overlap: + break + overlap_lines.insert(0, prev_line) + overlap_tokens += t + current_lines = overlap_lines + [line] + current_tokens = overlap_tokens + line_tokens + start_line = line_num - len(overlap_lines) + else: + current_lines.append(line) + current_tokens += line_tokens + + if current_lines: + chunks.append(TextChunk( + text="".join(current_lines), + start_line=start_line, + end_line=len(lines), + chunk_index=chunk_idx, + token_count=current_tokens + )) + + return chunks diff --git a/daemon/embedder_local.py b/daemon/embedder_local.py new file mode 100644 index 0000000..355cc6f --- /dev/null +++ b/daemon/embedder_local.py @@ -0,0 +1,23 @@ +"""Embedding generator for local asyncpg-based vector indexer.""" +from typing import List + + +class EmbeddingGenerator: + def __init__(self, model_name: str = "paraphrase-multilingual-MiniLM-L12-v2", device: str = "cpu"): + self.model_name = model_name + self.device = device + self._model = None + + def _get_model(self): + if self._model is None: + from sentence_transformers import SentenceTransformer + self._model = SentenceTransformer(self.model_name, device=self.device) + return self._model + + def embed_batch(self, texts: List[str]) -> List[List[float]]: + model = self._get_model() + embeddings = model.encode(texts, batch_size=32, show_progress_bar=False) + return [e.tolist() for e in embeddings] + + def embed_single(self, text: str) -> List[float]: + return self.embed_batch([text])[0] diff --git a/daemon/watcher_local.py b/daemon/watcher_local.py new file mode 100644 index 0000000..5d90bb0 --- /dev/null +++ b/daemon/watcher_local.py @@ -0,0 +1,113 @@ +"""Asyncpg-based file watcher daemon for local PostgreSQL vector indexer.""" +import asyncio +import asyncpg +import logging +import time +import yaml +from pathlib import Path +from watchdog.observers import Observer +from watchdog.events import FileSystemEventHandler + +logger = logging.getLogger(__name__) + +CONFIG_PATH = "/home/david/vector-indexer-mcp/config.yaml" + + +class IndexEventHandler(FileSystemEventHandler): + def __init__(self, queue, extensions, exclude_patterns): + self.queue = queue + self.extensions = extensions + self.exclude_patterns = exclude_patterns + self._debounce = {} + + def _should_index(self, path: str) -> bool: + p = Path(path) + if any(ex in path for ex in self.exclude_patterns): + return False + return p.suffix in self.extensions + + def _enqueue(self, path: str, event_type: str): + if not self._should_index(path): + return + now = time.time() + key = (path, event_type) + self._debounce[key] = now + try: + self.queue.put_nowait((path, event_type, now)) + except asyncio.QueueFull: + pass + + def on_created(self, event): + if not event.is_directory: + self._enqueue(event.src_path, 'create') + + def on_modified(self, event): + if not event.is_directory: + self._enqueue(event.src_path, 'modify') + + def on_deleted(self, event): + if not event.is_directory: + self._enqueue(event.src_path, 'delete') + + def on_moved(self, event): + if not event.is_directory: + self._enqueue(event.src_path, 'delete') + self._enqueue(event.dest_path, 'create') + + +class FileWatcherDaemon: + def __init__(self, config_path: str = CONFIG_PATH): + with open(config_path) as f: + self.config = yaml.safe_load(f) + db = self.config['database'] + self.dsn = f"postgresql://{db['user']}:{db['password']}@{db['host']}:{db['port']}/{db['name']}" + w = self.config.get('watcher', {}) + self.paths = w.get('paths', []) + self.extensions = set(w.get('extensions', ['.py', '.md'])) + self.exclude_patterns = w.get('exclude_patterns', ['__pycache__', '.git', 'node_modules']) + self.debounce_ms = w.get('debounce_ms', 500) + self.queue = asyncio.Queue(maxsize=10000) + + async def start(self): + self.pool = await asyncpg.create_pool(self.dsn) + observer = Observer() + handler = IndexEventHandler(self.queue, self.extensions, self.exclude_patterns) + for path in self.paths: + p = Path(path) + if p.exists(): + observer.schedule(handler, str(p), recursive=True) + logger.info(f"Watching: {path}") + else: + logger.warning(f"Watch path does not exist: {path}") + observer.start() + logger.info("FileWatcherDaemon started") + await self._flush_loop() + + async def _flush_loop(self): + debounce_s = self.debounce_ms / 1000 + pending = {} + while True: + try: + path, event_type, ts = await asyncio.wait_for(self.queue.get(), timeout=debounce_s) + pending[(path, event_type)] = ts + except asyncio.TimeoutError: + pass + now = time.time() + to_flush = [(p, e) for (p, e), ts in list(pending.items()) if now - ts >= debounce_s] + if to_flush: + async with self.pool.acquire() as conn: + for path, event_type in to_flush: + await conn.execute( + "INSERT INTO index_queue(file_path, event_type, status) VALUES($1,$2,'pending') " + "ON CONFLICT DO NOTHING", + path, event_type + ) + del pending[(path, event_type)] + + +if __name__ == "__main__": + import sys + logging.basicConfig(level=logging.INFO) + config = sys.argv[1] if len(sys.argv) > 1 else CONFIG_PATH + daemon = FileWatcherDaemon(config) + asyncio.run(daemon.start()) diff --git a/daemon/worker_local.py b/daemon/worker_local.py new file mode 100644 index 0000000..23e08e0 --- /dev/null +++ b/daemon/worker_local.py @@ -0,0 +1,139 @@ +"""Asyncpg-based indexing worker for local PostgreSQL vector search. + +Processes index_queue table events and generates embeddings for file content. +""" +import asyncio +import asyncpg +import hashlib +import logging +import yaml +from pathlib import Path +from .chunker_local import TextChunker +from .embedder_local import EmbeddingGenerator + +logger = logging.getLogger(__name__) + +INDEXABLE_EXTENSIONS = { + '.py', '.md', '.txt', '.json', '.yaml', '.yml', '.toml', + '.js', '.ts', '.tsx', '.jsx', '.html', '.css', '.sql', + '.sh', '.bash', '.ini', '.cfg' +} + +CONFIG_PATH = "/home/david/vector-indexer-mcp/config.yaml" + + +class IndexingWorker: + def __init__(self, config_path: str = CONFIG_PATH): + with open(config_path) as f: + self.config = yaml.safe_load(f) + db = self.config['database'] + self.dsn = f"postgresql://{db['user']}:{db['password']}@{db['host']}:{db['port']}/{db['name']}" + chunk_cfg = self.config.get('chunking', {}) + self.chunker = TextChunker( + chunk_size_tokens=chunk_cfg.get('chunk_size_tokens', 500), + overlap_tokens=chunk_cfg.get('overlap_tokens', 50) + ) + emb_cfg = self.config.get('embedding', {}) + self.embedder = EmbeddingGenerator( + model_name=emb_cfg.get('model', 'paraphrase-multilingual-MiniLM-L12-v2'), + device=emb_cfg.get('device', 'cpu') + ) + self.pool = None + + async def start(self): + self.pool = await asyncpg.create_pool(self.dsn) + logger.info("IndexingWorker started") + await self._process_loop() + + async def _process_loop(self): + poll_ms = self.config.get('queue', {}).get('poll_interval_ms', 500) + batch_size = self.config.get('queue', {}).get('batch_size', 50) + while True: + try: + await self._process_batch(batch_size) + except Exception as e: + logger.error(f"Worker error: {e}") + await asyncio.sleep(poll_ms / 1000) + + async def _process_batch(self, batch_size: int): + async with self.pool.acquire() as conn: + rows = await conn.fetch( + "UPDATE index_queue SET status='processing' WHERE id IN " + "(SELECT id FROM index_queue WHERE status='pending' LIMIT $1 FOR UPDATE SKIP LOCKED) " + "RETURNING id, file_path, event_type", + batch_size + ) + for row in rows: + try: + if row['event_type'] == 'delete': + await self._handle_delete(conn, row['file_path']) + else: + await self._handle_upsert(conn, row['file_path']) + await conn.execute( + "UPDATE index_queue SET status='done', processed_at=NOW() WHERE id=$1", + row['id'] + ) + except Exception as e: + logger.error(f"Failed to index {row['file_path']}: {e}") + await conn.execute( + "UPDATE index_queue SET status='error' WHERE id=$1", + row['id'] + ) + + async def _handle_upsert(self, conn, file_path: str): + p = Path(file_path) + if not p.exists() or p.suffix not in INDEXABLE_EXTENSIONS: + return + try: + content = p.read_text(encoding='utf-8', errors='ignore') + except Exception as e: + logger.warning(f"Cannot read {file_path}: {e}") + return + file_hash = hashlib.sha256(content.encode()).hexdigest() + existing = await conn.fetchrow( + "SELECT id, file_hash FROM file_metadata WHERE file_path=$1", file_path + ) + if existing and existing['file_hash'] == file_hash: + return # unchanged + # Upsert metadata + if existing: + file_id = existing['id'] + await conn.execute( + "UPDATE file_metadata SET file_hash=$1, file_size=$2, updated_at=NOW(), indexed_at=NOW() WHERE id=$3", + file_hash, len(content), file_id + ) + await conn.execute("DELETE FROM file_chunks WHERE file_id=$1", file_id) + else: + file_id = await conn.fetchval( + "INSERT INTO file_metadata(file_path, file_hash, file_size) VALUES($1,$2,$3) RETURNING id", + file_path, file_hash, len(content) + ) + # Chunk and embed + chunks = self.chunker.chunk_text(content, file_path) + if not chunks: + return + texts = [c.text for c in chunks] + embeddings = self.embedder.embed_batch(texts) + for chunk, embedding in zip(chunks, embeddings): + chunk_id = await conn.fetchval( + "INSERT INTO file_chunks(file_id, chunk_index, chunk_text, chunk_start_line, chunk_end_line, token_count) " + "VALUES($1,$2,$3,$4,$5,$6) RETURNING id", + file_id, chunk.chunk_index, chunk.text, chunk.start_line, chunk.end_line, chunk.token_count + ) + await conn.execute( + "INSERT INTO file_embeddings(chunk_id, embedding) VALUES($1,$2::vector)", + chunk_id, str(embedding) + ) + logger.info(f"Indexed {file_path}: {len(chunks)} chunks") + + async def _handle_delete(self, conn, file_path: str): + await conn.execute("DELETE FROM file_metadata WHERE file_path=$1", file_path) + logger.info(f"Removed from index: {file_path}") + + +if __name__ == "__main__": + import sys + logging.basicConfig(level=logging.INFO) + config = sys.argv[1] if len(sys.argv) > 1 else CONFIG_PATH + worker = IndexingWorker(config) + asyncio.run(worker.start()) diff --git a/migrations/020_vector_search_indexer.sql b/migrations/020_vector_search_indexer.sql new file mode 100644 index 0000000..0da3f39 --- /dev/null +++ b/migrations/020_vector_search_indexer.sql @@ -0,0 +1,146 @@ +-- Enable pgvector +CREATE EXTENSION IF NOT EXISTS vector; + +-- File metadata table +CREATE TABLE IF NOT EXISTS file_metadata ( + id SERIAL PRIMARY KEY, + file_path TEXT UNIQUE NOT NULL, + file_hash TEXT, + file_size INTEGER, + indexed_at TIMESTAMPTZ DEFAULT NOW(), + updated_at TIMESTAMPTZ DEFAULT NOW() +); + +-- File chunks table +CREATE TABLE IF NOT EXISTS file_chunks ( + id SERIAL PRIMARY KEY, + file_id INTEGER REFERENCES file_metadata(id) ON DELETE CASCADE, + chunk_index INTEGER NOT NULL, + chunk_text TEXT NOT NULL, + chunk_start_line INTEGER, + chunk_end_line INTEGER, + token_count INTEGER, + fts_vector tsvector GENERATED ALWAYS AS (to_tsvector('english', chunk_text)) STORED +); + +-- File embeddings table (384-dim for paraphrase-multilingual-MiniLM-L12-v2) +CREATE TABLE IF NOT EXISTS file_embeddings ( + id SERIAL PRIMARY KEY, + chunk_id INTEGER REFERENCES file_chunks(id) ON DELETE CASCADE, + embedding vector(384) +); + +-- Index queue table +CREATE TABLE IF NOT EXISTS index_queue ( + id SERIAL PRIMARY KEY, + file_path TEXT NOT NULL, + event_type TEXT NOT NULL, + status TEXT DEFAULT 'pending', + created_at TIMESTAMPTZ DEFAULT NOW(), + processed_at TIMESTAMPTZ +); + +-- Index stats table +CREATE TABLE IF NOT EXISTS index_stats ( + id SERIAL PRIMARY KEY, + stat_time TIMESTAMPTZ DEFAULT NOW(), + total_files INTEGER, + total_chunks INTEGER, + total_embeddings INTEGER, + queue_pending INTEGER +); + +-- Indexes +CREATE INDEX IF NOT EXISTS file_chunks_fts_idx ON file_chunks USING GIN(fts_vector); +CREATE INDEX IF NOT EXISTS file_embeddings_hnsw_idx ON file_embeddings USING hnsw(embedding vector_cosine_ops); +CREATE INDEX IF NOT EXISTS file_metadata_path_idx ON file_metadata(file_path); +CREATE INDEX IF NOT EXISTS index_queue_status_idx ON index_queue(status); + +-- Hybrid search function +CREATE OR REPLACE FUNCTION search_hybrid( + query_text TEXT, + query_embedding vector(384), + result_limit INTEGER DEFAULT 20, + alpha FLOAT DEFAULT 0.5 +) +RETURNS TABLE( + chunk_id INTEGER, + file_path TEXT, + chunk_text TEXT, + chunk_start_line INTEGER, + chunk_end_line INTEGER, + fts_rank FLOAT, + vector_similarity FLOAT, + combined_score FLOAT +) AS $$ +BEGIN + RETURN QUERY + WITH fts_results AS ( + SELECT + fc.id, + fm.file_path, + fc.chunk_text, + fc.chunk_start_line, + fc.chunk_end_line, + ts_rank(fc.fts_vector, plainto_tsquery('english', query_text))::FLOAT as fts_rank, + 0.0::FLOAT as vector_sim + FROM file_chunks fc + JOIN file_metadata fm ON fc.file_id = fm.id + WHERE fc.fts_vector @@ plainto_tsquery('english', query_text) + ), + vector_results AS ( + SELECT + fc.id, + fm.file_path, + fc.chunk_text, + fc.chunk_start_line, + fc.chunk_end_line, + 0.0::FLOAT as fts_rank, + (1 - (fe.embedding <=> query_embedding))::FLOAT as vector_sim + FROM file_embeddings fe + JOIN file_chunks fc ON fe.chunk_id = fc.id + JOIN file_metadata fm ON fc.file_id = fm.id + ), + combined AS ( + SELECT + COALESCE(f.id, v.id) as id, + COALESCE(f.file_path, v.file_path) as file_path, + COALESCE(f.chunk_text, v.chunk_text) as chunk_text, + COALESCE(f.chunk_start_line, v.chunk_start_line) as chunk_start_line, + COALESCE(f.chunk_end_line, v.chunk_end_line) as chunk_end_line, + COALESCE(f.fts_rank, 0.0) as fts_rank, + COALESCE(v.vector_sim, 0.0) as vector_sim + FROM fts_results f + FULL OUTER JOIN vector_results v ON f.id = v.id + ) + SELECT + c.id, + c.file_path, + c.chunk_text, + c.chunk_start_line, + c.chunk_end_line, + c.fts_rank, + c.vector_sim, + (alpha * c.vector_sim + (1-alpha) * c.fts_rank)::FLOAT as combined_score + FROM combined c + ORDER BY combined_score DESC + LIMIT result_limit; +END; +$$ LANGUAGE plpgsql; + +-- Index health function +CREATE OR REPLACE FUNCTION get_index_health() +RETURNS JSON AS $$ +DECLARE + result JSON; +BEGIN + SELECT json_build_object( + 'total_files', (SELECT COUNT(*) FROM file_metadata), + 'total_chunks', (SELECT COUNT(*) FROM file_chunks), + 'total_embeddings', (SELECT COUNT(*) FROM file_embeddings), + 'pending_queue', (SELECT COUNT(*) FROM index_queue WHERE status = 'pending'), + 'last_indexed', (SELECT MAX(indexed_at) FROM file_metadata) + ) INTO result; + RETURN result; +END; +$$ LANGUAGE plpgsql; diff --git a/requirements.txt b/requirements.txt index 06cb541..0b0f8a9 100644 --- a/requirements.txt +++ b/requirements.txt @@ -25,6 +25,9 @@ pyyaml>=6.0 # MCP server mcp>=1.0.0 +# Local PostgreSQL backend (used by *_local.py modules) +asyncpg>=0.29.0 + # HTTP/SSE Transport uvicorn>=0.27.0 starlette>=0.36.0 diff --git a/src/server_local.py b/src/server_local.py index c435b48..aa5bef9 100644 --- a/src/server_local.py +++ b/src/server_local.py @@ -1,12 +1,310 @@ +#!/usr/bin/env python3 """ -Backwards-compatibility shim. +Vector Indexer MCP Server - Local PostgreSQL backend via asyncpg. -Early installs configured ~/.mcp.json to invoke `-m src.server_local`. -The correct module is `src.server`, but this shim ensures both work. - -See: https://github.com/davidgut1982/vector-indexer-mcp/issues/1 +Tools: +- search_semantic: Vector similarity search +- search_lexical: Full-text search (FTS) +- search_hybrid: Combined semantic + lexical search +- index_status: Get index health statistics +- reindex_path: Force reindex a file/directory +- get_file_chunks: View chunks for a specific file +- search_similar_files: Find files similar to a given file """ -from src.server import main # noqa: F401 + +import asyncio +import asyncpg +import json +import logging +import yaml +from pathlib import Path +from typing import Any, Optional +from mcp.server import Server +from mcp.server.stdio import stdio_server +from mcp.types import Tool, TextContent + +logger = logging.getLogger(__name__) + +CONFIG_PATH = "/home/david/vector-indexer-mcp/config.yaml" + +with open(CONFIG_PATH) as f: + config = yaml.safe_load(f) + +db_cfg = config['database'] +DSN = f"postgresql://{db_cfg['user']}:{db_cfg['password']}@{db_cfg['host']}:{db_cfg['port']}/{db_cfg['name']}" + +_pool = None +_embedder = None + +INDEXABLE_EXTENSIONS = { + '.py', '.md', '.txt', '.json', '.yaml', '.yml', '.toml', + '.js', '.ts', '.tsx', '.jsx', '.html', '.css', '.sql', + '.sh', '.bash', '.ini', '.cfg' +} + + +async def get_pool(): + global _pool + if _pool is None: + _pool = await asyncpg.create_pool(DSN) + return _pool + + +def get_embedder(): + global _embedder + if _embedder is None: + from sentence_transformers import SentenceTransformer + emb_cfg = config.get('embedding', {}) + _embedder = SentenceTransformer( + emb_cfg.get('model', 'paraphrase-multilingual-MiniLM-L12-v2'), + device=emb_cfg.get('device', 'cpu') + ) + return _embedder + + +server = Server("vector-indexer-mcp") + + +@server.list_tools() +async def list_tools(): + return [ + Tool( + name="search_semantic", + description="Search indexed files using semantic similarity", + inputSchema={ + "type": "object", + "properties": { + "query": {"type": "string"}, + "limit": {"type": "integer", "default": 20}, + "threshold": {"type": "number", "default": 0.5}, + "paths": {"type": "array", "items": {"type": "string"}} + }, + "required": ["query"] + } + ), + Tool( + name="search_lexical", + description="Search indexed files using full-text search", + inputSchema={ + "type": "object", + "properties": { + "query": {"type": "string"}, + "limit": {"type": "integer", "default": 20}, + "paths": {"type": "array", "items": {"type": "string"}} + }, + "required": ["query"] + } + ), + Tool( + name="search_hybrid", + description="Search using both semantic and lexical methods", + inputSchema={ + "type": "object", + "properties": { + "query": {"type": "string"}, + "limit": {"type": "integer", "default": 20}, + "alpha": {"type": "number", "default": 0.5}, + "paths": {"type": "array", "items": {"type": "string"}} + }, + "required": ["query"] + } + ), + Tool( + name="index_status", + description="Get current index health and statistics", + inputSchema={"type": "object", "properties": {}} + ), + Tool( + name="reindex_path", + description="Force reindex of a directory or file", + inputSchema={ + "type": "object", + "properties": { + "path": {"type": "string"}, + "recursive": {"type": "boolean", "default": True}, + "force": {"type": "boolean", "default": False} + }, + "required": ["path"] + } + ), + Tool( + name="get_file_chunks", + description="Get all indexed chunks for a specific file", + inputSchema={ + "type": "object", + "properties": { + "file_path": {"type": "string"} + }, + "required": ["file_path"] + } + ), + Tool( + name="search_similar_files", + description="Find files similar to a given file", + inputSchema={ + "type": "object", + "properties": { + "file_path": {"type": "string"}, + "limit": {"type": "integer", "default": 10} + }, + "required": ["file_path"] + } + ), + ] + + +@server.call_tool() +async def call_tool(name: str, arguments: dict) -> list[TextContent]: + pool = await get_pool() + result = await _dispatch(name, arguments, pool) + return [TextContent(type="text", text=json.dumps(result, default=str))] + + +async def _dispatch(name: str, args: dict, pool) -> Any: + if name == "search_semantic": + return await _search_semantic(pool, **args) + elif name == "search_lexical": + return await _search_lexical(pool, **args) + elif name == "search_hybrid": + return await _search_hybrid(pool, **args) + elif name == "index_status": + return await _index_status(pool) + elif name == "reindex_path": + return await _reindex_path(pool, **args) + elif name == "get_file_chunks": + return await _get_file_chunks(pool, **args) + elif name == "search_similar_files": + return await _search_similar_files(pool, **args) + else: + return {"error": f"Unknown tool: {name}"} + + +async def _search_semantic(pool, query: str, limit: int = 20, threshold: float = 0.5, paths=None): + embedder = get_embedder() + q_emb = embedder.encode(query).tolist() + where = "WHERE 1=1" + params = [str(q_emb), limit] + if paths: + conditions = " OR ".join([f"fm.file_path LIKE ${len(params)+i+1}" for i in range(len(paths))]) + where += f" AND ({conditions})" + params.extend([f"{p}%" for p in paths]) + sql = f""" + SELECT fm.file_path, fc.chunk_text, fc.chunk_start_line, fc.chunk_end_line, + 1 - (fe.embedding <=> $1::vector) as similarity + FROM file_embeddings fe + JOIN file_chunks fc ON fe.chunk_id = fc.id + JOIN file_metadata fm ON fc.file_id = fm.id + {where} + ORDER BY fe.embedding <=> $1::vector + LIMIT $2 + """ + async with pool.acquire() as conn: + rows = await conn.fetch(sql, *params) + return [dict(r) for r in rows if r['similarity'] >= threshold] + + +async def _search_lexical(pool, query: str, limit: int = 20, paths=None): + where = "WHERE fc.fts_vector @@ plainto_tsquery('english', $1)" + params = [query, limit] + if paths: + conditions = " OR ".join([f"fm.file_path LIKE ${len(params)+i+1}" for i in range(len(paths))]) + where += f" AND ({conditions})" + params.extend([f"{p}%" for p in paths]) + sql = f""" + SELECT fm.file_path, fc.chunk_text, fc.chunk_start_line, fc.chunk_end_line, + ts_rank(fc.fts_vector, plainto_tsquery('english', $1)) as rank + FROM file_chunks fc + JOIN file_metadata fm ON fc.file_id = fm.id + {where} + ORDER BY rank DESC LIMIT $2 + """ + async with pool.acquire() as conn: + rows = await conn.fetch(sql, *params) + return [dict(r) for r in rows] + + +async def _search_hybrid(pool, query: str, limit: int = 20, alpha: float = 0.5, paths=None): + embedder = get_embedder() + q_emb = embedder.encode(query).tolist() + async with pool.acquire() as conn: + rows = await conn.fetch( + "SELECT * FROM search_hybrid($1, $2::vector, $3, $4)", + query, str(q_emb), limit, alpha + ) + results = [dict(r) for r in rows] + if paths: + results = [r for r in results if any(r['file_path'].startswith(p) for p in paths)] + return results[:limit] + + +async def _index_status(pool): + async with pool.acquire() as conn: + health = await conn.fetchrow("SELECT get_index_health() as h") + pending = await conn.fetchval("SELECT COUNT(*) FROM index_queue WHERE status='pending'") + processing = await conn.fetchval("SELECT COUNT(*) FROM index_queue WHERE status='processing'") + return { + "health": json.loads(health['h']) if health else {}, + "queue": {"pending": pending, "processing": processing} + } + + +async def _reindex_path(pool, path: str, recursive: bool = True, force: bool = False): + p = Path(path) + if not p.exists(): + return {"error": f"Path not found: {path}", "queued": 0} + if p.is_file(): + files = [str(p)] if p.suffix in INDEXABLE_EXTENSIONS else [] + else: + glob_fn = p.rglob if recursive else p.glob + files = [str(f) for f in glob_fn("*") if f.is_file() and f.suffix in INDEXABLE_EXTENSIONS] + async with pool.acquire() as conn: + for f in files: + await conn.execute( + "INSERT INTO index_queue(file_path, event_type, status) VALUES($1,$2,'pending')", + f, 'modify' if force else 'create' + ) + return {"queued": len(files), "path": path} + + +async def _get_file_chunks(pool, file_path: str): + async with pool.acquire() as conn: + rows = await conn.fetch( + "SELECT fc.chunk_index, fc.chunk_text, fc.chunk_start_line, fc.chunk_end_line, " + "fc.token_count, fm.indexed_at " + "FROM file_chunks fc JOIN file_metadata fm ON fc.file_id=fm.id " + "WHERE fm.file_path=$1 ORDER BY fc.chunk_index", + file_path + ) + return [dict(r) for r in rows] + + +async def _search_similar_files(pool, file_path: str, limit: int = 10): + async with pool.acquire() as conn: + rows = await conn.fetch(""" + WITH source_emb AS ( + SELECT fe.embedding + FROM file_embeddings fe + JOIN file_chunks fc ON fe.chunk_id = fc.id + JOIN file_metadata fm ON fc.file_id = fm.id + WHERE fm.file_path = $1 + LIMIT 1 + ) + SELECT DISTINCT fm.file_path, + 1 - (fe.embedding <=> (SELECT embedding FROM source_emb)) as similarity + FROM file_embeddings fe + JOIN file_chunks fc ON fe.chunk_id = fc.id + JOIN file_metadata fm ON fc.file_id = fm.id + WHERE fm.file_path != $1 + ORDER BY similarity DESC LIMIT $2 + """, file_path, limit) + return [dict(r) for r in rows] + + +async def main(): + logging.basicConfig(level=logging.INFO) + async with stdio_server() as (read_stream, write_stream): + await server.run(read_stream, write_stream, server.create_initialization_options()) + if __name__ == "__main__": - main() + asyncio.run(main()) diff --git a/sse_server_local.py b/sse_server_local.py new file mode 100644 index 0000000..146c4d8 --- /dev/null +++ b/sse_server_local.py @@ -0,0 +1,105 @@ +#!/usr/bin/env python3 +""" +HTTP/SSE Transport Wrapper for vector-indexer-mcp (local PostgreSQL backend). + +Endpoints: + GET /sse - SSE connection for MCP protocol + POST /messages/ - Message endpoint for MCP protocol + GET /health - Health check endpoint + +Usage: + python sse_server_local.py + python sse_server_local.py --port 5577 + MCP_SSE_PORT=5577 python sse_server_local.py +""" + +import argparse +import asyncio +import logging +import os +import sys +from pathlib import Path + +# Add project root to sys.path +_project_root = str(Path(__file__).resolve().parent) +if _project_root not in sys.path: + sys.path.insert(0, _project_root) + +import uvicorn +from starlette.applications import Starlette +from starlette.middleware import Middleware +from starlette.middleware.cors import CORSMiddleware +from starlette.routing import Route, Mount +from starlette.responses import JSONResponse, Response +from mcp.server.sse import SseServerTransport + +logging.basicConfig( + level=logging.INFO, + format='%(asctime)s - %(name)s - %(levelname)s - %(message)s' +) +logger = logging.getLogger("vector-indexer-mcp-sse-local") + + +class _SseResponse(Response): + async def __call__(self, scope, receive, send): + pass # SSE transport already handled the response + + +def create_app(): + """Create Starlette app wrapping server_local's MCP server.""" + # Import the local server module to get its Server instance + from src.server_local import server as mcp_server + + sse = SseServerTransport("/messages/") + + async def handle_sse(request): + logger.info(f"SSE connection from {request.client.host if request.client else 'unknown'}") + async with sse.connect_sse( + request.scope, request.receive, request._send + ) as streams: + await mcp_server.run( + streams[0], streams[1], mcp_server.create_initialization_options() + ) + return _SseResponse() + + async def health(request): + return JSONResponse({ + "status": "healthy", + "server": "vector-indexer-mcp", + "transport": "sse-local" + }) + + routes = [ + Route("/health", endpoint=health, methods=["GET"]), + Route("/sse", endpoint=handle_sse, methods=["GET"]), + Mount("/messages/", app=sse.handle_post_message), + ] + + middleware = [ + Middleware( + CORSMiddleware, + allow_origins=["*"], + allow_credentials=True, + allow_methods=["GET", "POST", "OPTIONS"], + allow_headers=["*"], + ) + ] + + return Starlette(routes=routes, middleware=middleware) + + +def main(): + parser = argparse.ArgumentParser() + parser.add_argument("--port", "-p", type=int, + default=int(os.environ.get("MCP_SSE_PORT", "5577"))) + parser.add_argument("--host", default=os.environ.get("MCP_SSE_HOST", "127.0.0.1")) + args = parser.parse_args() + + logger.info("Initializing vector-indexer-mcp (local PostgreSQL)...") + app = create_app() + logger.info(f"Starting SSE server on {args.host}:{args.port}") + uvicorn.run(app, host=args.host, port=args.port, log_level="info") + + +if __name__ == "__main__": + main() diff --git a/systemd/vector-indexer-mcp-http.service b/systemd/vector-indexer-mcp-http.service new file mode 100644 index 0000000..b6708f0 --- /dev/null +++ b/systemd/vector-indexer-mcp-http.service @@ -0,0 +1,22 @@ +# vector-indexer-mcp-http.service +# Install: cp this file to ~/.config/systemd/user/ +# Then: systemctl --user daemon-reload && systemctl --user enable --now vector-indexer-mcp-http +# +# NOTE: Paths below (WorkingDirectory, ExecStart) are machine-specific. +# Adjust them to match your actual checkout and venv location. + +[Unit] +Description=Vector Indexer MCP HTTP/SSE Server (local asyncpg) +After=vector-indexer-worker.service + +[Service] +Type=simple +WorkingDirectory=/home/david/vector-indexer-mcp +ExecStart=/home/david/vector-indexer-mcp/venv/bin/python sse_server_local.py --port 5577 --host 127.0.0.1 +Restart=on-failure +RestartSec=5 +StandardOutput=journal +StandardError=journal + +[Install] +WantedBy=default.target diff --git a/systemd/vector-indexer-watcher.service b/systemd/vector-indexer-watcher.service new file mode 100644 index 0000000..1f40245 --- /dev/null +++ b/systemd/vector-indexer-watcher.service @@ -0,0 +1,22 @@ +# vector-indexer-watcher.service +# Install: cp this file to ~/.config/systemd/user/ +# Then: systemctl --user daemon-reload && systemctl --user enable --now vector-indexer-watcher +# +# NOTE: Paths below (WorkingDirectory, ExecStart) are machine-specific. +# Adjust them to match your actual checkout and venv location. + +[Unit] +Description=Vector Indexer File Watcher (local asyncpg) +After=network.target + +[Service] +Type=simple +WorkingDirectory=/home/david/vector-indexer-mcp +ExecStart=/home/david/vector-indexer-mcp/venv/bin/python -m daemon.watcher_local +Restart=on-failure +RestartSec=5 +StandardOutput=journal +StandardError=journal + +[Install] +WantedBy=default.target diff --git a/systemd/vector-indexer-worker.service b/systemd/vector-indexer-worker.service new file mode 100644 index 0000000..b7806ae --- /dev/null +++ b/systemd/vector-indexer-worker.service @@ -0,0 +1,22 @@ +# vector-indexer-worker.service +# Install: cp this file to ~/.config/systemd/user/ +# Then: systemctl --user daemon-reload && systemctl --user enable --now vector-indexer-worker +# +# NOTE: Paths below (WorkingDirectory, ExecStart) are machine-specific. +# Adjust them to match your actual checkout and venv location. + +[Unit] +Description=Vector Indexer Worker (local asyncpg) +After=network.target + +[Service] +Type=simple +WorkingDirectory=/home/david/vector-indexer-mcp +ExecStart=/home/david/vector-indexer-mcp/venv/bin/python -m daemon.worker_local +Restart=on-failure +RestartSec=5 +StandardOutput=journal +StandardError=journal + +[Install] +WantedBy=default.target From 4caee451194e6eadf2f2074bb1b32ad9448f5a66 Mon Sep 17 00:00:00 2001 From: Version Control Agent Date: Mon, 1 Jun 2026 16:08:38 -0500 Subject: [PATCH 2/2] fix: ruff lint and format Co-Authored-By: Claude Sonnet 4.6 --- daemon/chunker_local.py | 34 ++++++---- daemon/embedder_local.py | 8 ++- daemon/watcher_local.py | 44 +++++++----- daemon/worker_local.py | 84 +++++++++++++++-------- src/server_local.py | 141 +++++++++++++++++++++++++-------------- sse_server_local.py | 39 ++++++----- 6 files changed, 223 insertions(+), 127 deletions(-) diff --git a/daemon/chunker_local.py b/daemon/chunker_local.py index 565278f..889038a 100644 --- a/daemon/chunker_local.py +++ b/daemon/chunker_local.py @@ -1,4 +1,5 @@ """Token-aware text chunker for local asyncpg-based vector indexer.""" + from dataclasses import dataclass from typing import List @@ -18,6 +19,7 @@ def __init__(self, chunk_size_tokens: int = 500, overlap_tokens: int = 50): self.overlap = overlap_tokens try: import tiktoken + self.enc = tiktoken.get_encoding("cl100k_base") self._use_tiktoken = True except Exception: @@ -40,13 +42,15 @@ def chunk_text(self, text: str, file_path: str = "") -> List[TextChunk]: line_tokens = self._count_tokens(line) if current_tokens + line_tokens > self.chunk_size and current_lines: chunk_text = "".join(current_lines) - chunks.append(TextChunk( - text=chunk_text, - start_line=start_line, - end_line=line_num - 1, - chunk_index=chunk_idx, - token_count=current_tokens - )) + chunks.append( + TextChunk( + text=chunk_text, + start_line=start_line, + end_line=line_num - 1, + chunk_index=chunk_idx, + token_count=current_tokens, + ) + ) chunk_idx += 1 # Overlap: keep last N tokens worth of lines overlap_lines = [] @@ -65,12 +69,14 @@ def chunk_text(self, text: str, file_path: str = "") -> List[TextChunk]: current_tokens += line_tokens if current_lines: - chunks.append(TextChunk( - text="".join(current_lines), - start_line=start_line, - end_line=len(lines), - chunk_index=chunk_idx, - token_count=current_tokens - )) + chunks.append( + TextChunk( + text="".join(current_lines), + start_line=start_line, + end_line=len(lines), + chunk_index=chunk_idx, + token_count=current_tokens, + ) + ) return chunks diff --git a/daemon/embedder_local.py b/daemon/embedder_local.py index 355cc6f..cf2a782 100644 --- a/daemon/embedder_local.py +++ b/daemon/embedder_local.py @@ -1,9 +1,14 @@ """Embedding generator for local asyncpg-based vector indexer.""" + from typing import List class EmbeddingGenerator: - def __init__(self, model_name: str = "paraphrase-multilingual-MiniLM-L12-v2", device: str = "cpu"): + def __init__( + self, + model_name: str = "paraphrase-multilingual-MiniLM-L12-v2", + device: str = "cpu", + ): self.model_name = model_name self.device = device self._model = None @@ -11,6 +16,7 @@ def __init__(self, model_name: str = "paraphrase-multilingual-MiniLM-L12-v2", de def _get_model(self): if self._model is None: from sentence_transformers import SentenceTransformer + self._model = SentenceTransformer(self.model_name, device=self.device) return self._model diff --git a/daemon/watcher_local.py b/daemon/watcher_local.py index 5d90bb0..fa92664 100644 --- a/daemon/watcher_local.py +++ b/daemon/watcher_local.py @@ -1,12 +1,14 @@ """Asyncpg-based file watcher daemon for local PostgreSQL vector indexer.""" + import asyncio -import asyncpg import logging import time -import yaml from pathlib import Path -from watchdog.observers import Observer + +import asyncpg +import yaml from watchdog.events import FileSystemEventHandler +from watchdog.observers import Observer logger = logging.getLogger(__name__) @@ -39,33 +41,35 @@ def _enqueue(self, path: str, event_type: str): def on_created(self, event): if not event.is_directory: - self._enqueue(event.src_path, 'create') + self._enqueue(event.src_path, "create") def on_modified(self, event): if not event.is_directory: - self._enqueue(event.src_path, 'modify') + self._enqueue(event.src_path, "modify") def on_deleted(self, event): if not event.is_directory: - self._enqueue(event.src_path, 'delete') + self._enqueue(event.src_path, "delete") def on_moved(self, event): if not event.is_directory: - self._enqueue(event.src_path, 'delete') - self._enqueue(event.dest_path, 'create') + self._enqueue(event.src_path, "delete") + self._enqueue(event.dest_path, "create") class FileWatcherDaemon: def __init__(self, config_path: str = CONFIG_PATH): with open(config_path) as f: self.config = yaml.safe_load(f) - db = self.config['database'] + db = self.config["database"] self.dsn = f"postgresql://{db['user']}:{db['password']}@{db['host']}:{db['port']}/{db['name']}" - w = self.config.get('watcher', {}) - self.paths = w.get('paths', []) - self.extensions = set(w.get('extensions', ['.py', '.md'])) - self.exclude_patterns = w.get('exclude_patterns', ['__pycache__', '.git', 'node_modules']) - self.debounce_ms = w.get('debounce_ms', 500) + w = self.config.get("watcher", {}) + self.paths = w.get("paths", []) + self.extensions = set(w.get("extensions", [".py", ".md"])) + self.exclude_patterns = w.get( + "exclude_patterns", ["__pycache__", ".git", "node_modules"] + ) + self.debounce_ms = w.get("debounce_ms", 500) self.queue = asyncio.Queue(maxsize=10000) async def start(self): @@ -88,25 +92,31 @@ async def _flush_loop(self): pending = {} while True: try: - path, event_type, ts = await asyncio.wait_for(self.queue.get(), timeout=debounce_s) + path, event_type, ts = await asyncio.wait_for( + self.queue.get(), timeout=debounce_s + ) pending[(path, event_type)] = ts except asyncio.TimeoutError: pass now = time.time() - to_flush = [(p, e) for (p, e), ts in list(pending.items()) if now - ts >= debounce_s] + to_flush = [ + (p, e) for (p, e), ts in list(pending.items()) if now - ts >= debounce_s + ] if to_flush: async with self.pool.acquire() as conn: for path, event_type in to_flush: await conn.execute( "INSERT INTO index_queue(file_path, event_type, status) VALUES($1,$2,'pending') " "ON CONFLICT DO NOTHING", - path, event_type + path, + event_type, ) del pending[(path, event_type)] if __name__ == "__main__": import sys + logging.basicConfig(level=logging.INFO) config = sys.argv[1] if len(sys.argv) > 1 else CONFIG_PATH daemon = FileWatcherDaemon(config) diff --git a/daemon/worker_local.py b/daemon/worker_local.py index 23e08e0..5e1bc60 100644 --- a/daemon/worker_local.py +++ b/daemon/worker_local.py @@ -2,21 +2,39 @@ Processes index_queue table events and generates embeddings for file content. """ + import asyncio -import asyncpg import hashlib import logging -import yaml from pathlib import Path + +import asyncpg +import yaml + from .chunker_local import TextChunker from .embedder_local import EmbeddingGenerator logger = logging.getLogger(__name__) INDEXABLE_EXTENSIONS = { - '.py', '.md', '.txt', '.json', '.yaml', '.yml', '.toml', - '.js', '.ts', '.tsx', '.jsx', '.html', '.css', '.sql', - '.sh', '.bash', '.ini', '.cfg' + ".py", + ".md", + ".txt", + ".json", + ".yaml", + ".yml", + ".toml", + ".js", + ".ts", + ".tsx", + ".jsx", + ".html", + ".css", + ".sql", + ".sh", + ".bash", + ".ini", + ".cfg", } CONFIG_PATH = "/home/david/vector-indexer-mcp/config.yaml" @@ -26,17 +44,17 @@ class IndexingWorker: def __init__(self, config_path: str = CONFIG_PATH): with open(config_path) as f: self.config = yaml.safe_load(f) - db = self.config['database'] + db = self.config["database"] self.dsn = f"postgresql://{db['user']}:{db['password']}@{db['host']}:{db['port']}/{db['name']}" - chunk_cfg = self.config.get('chunking', {}) + chunk_cfg = self.config.get("chunking", {}) self.chunker = TextChunker( - chunk_size_tokens=chunk_cfg.get('chunk_size_tokens', 500), - overlap_tokens=chunk_cfg.get('overlap_tokens', 50) + chunk_size_tokens=chunk_cfg.get("chunk_size_tokens", 500), + overlap_tokens=chunk_cfg.get("overlap_tokens", 50), ) - emb_cfg = self.config.get('embedding', {}) + emb_cfg = self.config.get("embedding", {}) self.embedder = EmbeddingGenerator( - model_name=emb_cfg.get('model', 'paraphrase-multilingual-MiniLM-L12-v2'), - device=emb_cfg.get('device', 'cpu') + model_name=emb_cfg.get("model", "paraphrase-multilingual-MiniLM-L12-v2"), + device=emb_cfg.get("device", "cpu"), ) self.pool = None @@ -46,8 +64,8 @@ async def start(self): await self._process_loop() async def _process_loop(self): - poll_ms = self.config.get('queue', {}).get('poll_interval_ms', 500) - batch_size = self.config.get('queue', {}).get('batch_size', 50) + poll_ms = self.config.get("queue", {}).get("poll_interval_ms", 500) + batch_size = self.config.get("queue", {}).get("batch_size", 50) while True: try: await self._process_batch(batch_size) @@ -61,23 +79,22 @@ async def _process_batch(self, batch_size: int): "UPDATE index_queue SET status='processing' WHERE id IN " "(SELECT id FROM index_queue WHERE status='pending' LIMIT $1 FOR UPDATE SKIP LOCKED) " "RETURNING id, file_path, event_type", - batch_size + batch_size, ) for row in rows: try: - if row['event_type'] == 'delete': - await self._handle_delete(conn, row['file_path']) + if row["event_type"] == "delete": + await self._handle_delete(conn, row["file_path"]) else: - await self._handle_upsert(conn, row['file_path']) + await self._handle_upsert(conn, row["file_path"]) await conn.execute( "UPDATE index_queue SET status='done', processed_at=NOW() WHERE id=$1", - row['id'] + row["id"], ) except Exception as e: logger.error(f"Failed to index {row['file_path']}: {e}") await conn.execute( - "UPDATE index_queue SET status='error' WHERE id=$1", - row['id'] + "UPDATE index_queue SET status='error' WHERE id=$1", row["id"] ) async def _handle_upsert(self, conn, file_path: str): @@ -85,7 +102,7 @@ async def _handle_upsert(self, conn, file_path: str): if not p.exists() or p.suffix not in INDEXABLE_EXTENSIONS: return try: - content = p.read_text(encoding='utf-8', errors='ignore') + content = p.read_text(encoding="utf-8", errors="ignore") except Exception as e: logger.warning(f"Cannot read {file_path}: {e}") return @@ -93,20 +110,24 @@ async def _handle_upsert(self, conn, file_path: str): existing = await conn.fetchrow( "SELECT id, file_hash FROM file_metadata WHERE file_path=$1", file_path ) - if existing and existing['file_hash'] == file_hash: + if existing and existing["file_hash"] == file_hash: return # unchanged # Upsert metadata if existing: - file_id = existing['id'] + file_id = existing["id"] await conn.execute( "UPDATE file_metadata SET file_hash=$1, file_size=$2, updated_at=NOW(), indexed_at=NOW() WHERE id=$3", - file_hash, len(content), file_id + file_hash, + len(content), + file_id, ) await conn.execute("DELETE FROM file_chunks WHERE file_id=$1", file_id) else: file_id = await conn.fetchval( "INSERT INTO file_metadata(file_path, file_hash, file_size) VALUES($1,$2,$3) RETURNING id", - file_path, file_hash, len(content) + file_path, + file_hash, + len(content), ) # Chunk and embed chunks = self.chunker.chunk_text(content, file_path) @@ -118,11 +139,17 @@ async def _handle_upsert(self, conn, file_path: str): chunk_id = await conn.fetchval( "INSERT INTO file_chunks(file_id, chunk_index, chunk_text, chunk_start_line, chunk_end_line, token_count) " "VALUES($1,$2,$3,$4,$5,$6) RETURNING id", - file_id, chunk.chunk_index, chunk.text, chunk.start_line, chunk.end_line, chunk.token_count + file_id, + chunk.chunk_index, + chunk.text, + chunk.start_line, + chunk.end_line, + chunk.token_count, ) await conn.execute( "INSERT INTO file_embeddings(chunk_id, embedding) VALUES($1,$2::vector)", - chunk_id, str(embedding) + chunk_id, + str(embedding), ) logger.info(f"Indexed {file_path}: {len(chunks)} chunks") @@ -133,6 +160,7 @@ async def _handle_delete(self, conn, file_path: str): if __name__ == "__main__": import sys + logging.basicConfig(level=logging.INFO) config = sys.argv[1] if len(sys.argv) > 1 else CONFIG_PATH worker = IndexingWorker(config) diff --git a/src/server_local.py b/src/server_local.py index aa5bef9..c384020 100644 --- a/src/server_local.py +++ b/src/server_local.py @@ -13,15 +13,16 @@ """ import asyncio -import asyncpg import json import logging -import yaml from pathlib import Path -from typing import Any, Optional +from typing import Any + +import asyncpg +import yaml from mcp.server import Server from mcp.server.stdio import stdio_server -from mcp.types import Tool, TextContent +from mcp.types import TextContent, Tool logger = logging.getLogger(__name__) @@ -30,16 +31,31 @@ with open(CONFIG_PATH) as f: config = yaml.safe_load(f) -db_cfg = config['database'] +db_cfg = config["database"] DSN = f"postgresql://{db_cfg['user']}:{db_cfg['password']}@{db_cfg['host']}:{db_cfg['port']}/{db_cfg['name']}" _pool = None _embedder = None INDEXABLE_EXTENSIONS = { - '.py', '.md', '.txt', '.json', '.yaml', '.yml', '.toml', - '.js', '.ts', '.tsx', '.jsx', '.html', '.css', '.sql', - '.sh', '.bash', '.ini', '.cfg' + ".py", + ".md", + ".txt", + ".json", + ".yaml", + ".yml", + ".toml", + ".js", + ".ts", + ".tsx", + ".jsx", + ".html", + ".css", + ".sql", + ".sh", + ".bash", + ".ini", + ".cfg", } @@ -54,10 +70,11 @@ def get_embedder(): global _embedder if _embedder is None: from sentence_transformers import SentenceTransformer - emb_cfg = config.get('embedding', {}) + + emb_cfg = config.get("embedding", {}) _embedder = SentenceTransformer( - emb_cfg.get('model', 'paraphrase-multilingual-MiniLM-L12-v2'), - device=emb_cfg.get('device', 'cpu') + emb_cfg.get("model", "paraphrase-multilingual-MiniLM-L12-v2"), + device=emb_cfg.get("device", "cpu"), ) return _embedder @@ -77,10 +94,10 @@ async def list_tools(): "query": {"type": "string"}, "limit": {"type": "integer", "default": 20}, "threshold": {"type": "number", "default": 0.5}, - "paths": {"type": "array", "items": {"type": "string"}} + "paths": {"type": "array", "items": {"type": "string"}}, }, - "required": ["query"] - } + "required": ["query"], + }, ), Tool( name="search_lexical", @@ -90,10 +107,10 @@ async def list_tools(): "properties": { "query": {"type": "string"}, "limit": {"type": "integer", "default": 20}, - "paths": {"type": "array", "items": {"type": "string"}} + "paths": {"type": "array", "items": {"type": "string"}}, }, - "required": ["query"] - } + "required": ["query"], + }, ), Tool( name="search_hybrid", @@ -104,15 +121,15 @@ async def list_tools(): "query": {"type": "string"}, "limit": {"type": "integer", "default": 20}, "alpha": {"type": "number", "default": 0.5}, - "paths": {"type": "array", "items": {"type": "string"}} + "paths": {"type": "array", "items": {"type": "string"}}, }, - "required": ["query"] - } + "required": ["query"], + }, ), Tool( name="index_status", description="Get current index health and statistics", - inputSchema={"type": "object", "properties": {}} + inputSchema={"type": "object", "properties": {}}, ), Tool( name="reindex_path", @@ -122,21 +139,19 @@ async def list_tools(): "properties": { "path": {"type": "string"}, "recursive": {"type": "boolean", "default": True}, - "force": {"type": "boolean", "default": False} + "force": {"type": "boolean", "default": False}, }, - "required": ["path"] - } + "required": ["path"], + }, ), Tool( name="get_file_chunks", description="Get all indexed chunks for a specific file", inputSchema={ "type": "object", - "properties": { - "file_path": {"type": "string"} - }, - "required": ["file_path"] - } + "properties": {"file_path": {"type": "string"}}, + "required": ["file_path"], + }, ), Tool( name="search_similar_files", @@ -145,10 +160,10 @@ async def list_tools(): "type": "object", "properties": { "file_path": {"type": "string"}, - "limit": {"type": "integer", "default": 10} + "limit": {"type": "integer", "default": 10}, }, - "required": ["file_path"] - } + "required": ["file_path"], + }, ), ] @@ -179,13 +194,17 @@ async def _dispatch(name: str, args: dict, pool) -> Any: return {"error": f"Unknown tool: {name}"} -async def _search_semantic(pool, query: str, limit: int = 20, threshold: float = 0.5, paths=None): +async def _search_semantic( + pool, query: str, limit: int = 20, threshold: float = 0.5, paths=None +): embedder = get_embedder() q_emb = embedder.encode(query).tolist() where = "WHERE 1=1" params = [str(q_emb), limit] if paths: - conditions = " OR ".join([f"fm.file_path LIKE ${len(params)+i+1}" for i in range(len(paths))]) + conditions = " OR ".join( + [f"fm.file_path LIKE ${len(params) + i + 1}" for i in range(len(paths))] + ) where += f" AND ({conditions})" params.extend([f"{p}%" for p in paths]) sql = f""" @@ -200,14 +219,16 @@ async def _search_semantic(pool, query: str, limit: int = 20, threshold: float = """ async with pool.acquire() as conn: rows = await conn.fetch(sql, *params) - return [dict(r) for r in rows if r['similarity'] >= threshold] + return [dict(r) for r in rows if r["similarity"] >= threshold] async def _search_lexical(pool, query: str, limit: int = 20, paths=None): where = "WHERE fc.fts_vector @@ plainto_tsquery('english', $1)" params = [query, limit] if paths: - conditions = " OR ".join([f"fm.file_path LIKE ${len(params)+i+1}" for i in range(len(paths))]) + conditions = " OR ".join( + [f"fm.file_path LIKE ${len(params) + i + 1}" for i in range(len(paths))] + ) where += f" AND ({conditions})" params.extend([f"{p}%" for p in paths]) sql = f""" @@ -223,28 +244,39 @@ async def _search_lexical(pool, query: str, limit: int = 20, paths=None): return [dict(r) for r in rows] -async def _search_hybrid(pool, query: str, limit: int = 20, alpha: float = 0.5, paths=None): +async def _search_hybrid( + pool, query: str, limit: int = 20, alpha: float = 0.5, paths=None +): embedder = get_embedder() q_emb = embedder.encode(query).tolist() async with pool.acquire() as conn: rows = await conn.fetch( "SELECT * FROM search_hybrid($1, $2::vector, $3, $4)", - query, str(q_emb), limit, alpha + query, + str(q_emb), + limit, + alpha, ) results = [dict(r) for r in rows] if paths: - results = [r for r in results if any(r['file_path'].startswith(p) for p in paths)] + results = [ + r for r in results if any(r["file_path"].startswith(p) for p in paths) + ] return results[:limit] async def _index_status(pool): async with pool.acquire() as conn: health = await conn.fetchrow("SELECT get_index_health() as h") - pending = await conn.fetchval("SELECT COUNT(*) FROM index_queue WHERE status='pending'") - processing = await conn.fetchval("SELECT COUNT(*) FROM index_queue WHERE status='processing'") + pending = await conn.fetchval( + "SELECT COUNT(*) FROM index_queue WHERE status='pending'" + ) + processing = await conn.fetchval( + "SELECT COUNT(*) FROM index_queue WHERE status='processing'" + ) return { - "health": json.loads(health['h']) if health else {}, - "queue": {"pending": pending, "processing": processing} + "health": json.loads(health["h"]) if health else {}, + "queue": {"pending": pending, "processing": processing}, } @@ -256,12 +288,17 @@ async def _reindex_path(pool, path: str, recursive: bool = True, force: bool = F files = [str(p)] if p.suffix in INDEXABLE_EXTENSIONS else [] else: glob_fn = p.rglob if recursive else p.glob - files = [str(f) for f in glob_fn("*") if f.is_file() and f.suffix in INDEXABLE_EXTENSIONS] + files = [ + str(f) + for f in glob_fn("*") + if f.is_file() and f.suffix in INDEXABLE_EXTENSIONS + ] async with pool.acquire() as conn: for f in files: await conn.execute( "INSERT INTO index_queue(file_path, event_type, status) VALUES($1,$2,'pending')", - f, 'modify' if force else 'create' + f, + "modify" if force else "create", ) return {"queued": len(files), "path": path} @@ -273,14 +310,15 @@ async def _get_file_chunks(pool, file_path: str): "fc.token_count, fm.indexed_at " "FROM file_chunks fc JOIN file_metadata fm ON fc.file_id=fm.id " "WHERE fm.file_path=$1 ORDER BY fc.chunk_index", - file_path + file_path, ) return [dict(r) for r in rows] async def _search_similar_files(pool, file_path: str, limit: int = 10): async with pool.acquire() as conn: - rows = await conn.fetch(""" + rows = await conn.fetch( + """ WITH source_emb AS ( SELECT fe.embedding FROM file_embeddings fe @@ -296,14 +334,19 @@ async def _search_similar_files(pool, file_path: str, limit: int = 10): JOIN file_metadata fm ON fc.file_id = fm.id WHERE fm.file_path != $1 ORDER BY similarity DESC LIMIT $2 - """, file_path, limit) + """, + file_path, + limit, + ) return [dict(r) for r in rows] async def main(): logging.basicConfig(level=logging.INFO) async with stdio_server() as (read_stream, write_stream): - await server.run(read_stream, write_stream, server.create_initialization_options()) + await server.run( + read_stream, write_stream, server.create_initialization_options() + ) if __name__ == "__main__": diff --git a/sse_server_local.py b/sse_server_local.py index 146c4d8..6ad181c 100644 --- a/sse_server_local.py +++ b/sse_server_local.py @@ -14,7 +14,6 @@ """ import argparse -import asyncio import logging import os import sys @@ -25,17 +24,16 @@ if _project_root not in sys.path: sys.path.insert(0, _project_root) -import uvicorn -from starlette.applications import Starlette -from starlette.middleware import Middleware -from starlette.middleware.cors import CORSMiddleware -from starlette.routing import Route, Mount -from starlette.responses import JSONResponse, Response -from mcp.server.sse import SseServerTransport +import uvicorn # noqa: E402 +from mcp.server.sse import SseServerTransport # noqa: E402 +from starlette.applications import Starlette # noqa: E402 +from starlette.middleware import Middleware # noqa: E402 +from starlette.middleware.cors import CORSMiddleware # noqa: E402 +from starlette.responses import JSONResponse, Response # noqa: E402 +from starlette.routing import Mount, Route # noqa: E402 logging.basicConfig( - level=logging.INFO, - format='%(asctime)s - %(name)s - %(levelname)s - %(message)s' + level=logging.INFO, format="%(asctime)s - %(name)s - %(levelname)s - %(message)s" ) logger = logging.getLogger("vector-indexer-mcp-sse-local") @@ -53,7 +51,9 @@ def create_app(): sse = SseServerTransport("/messages/") async def handle_sse(request): - logger.info(f"SSE connection from {request.client.host if request.client else 'unknown'}") + logger.info( + f"SSE connection from {request.client.host if request.client else 'unknown'}" + ) async with sse.connect_sse( request.scope, request.receive, request._send ) as streams: @@ -63,11 +63,13 @@ async def handle_sse(request): return _SseResponse() async def health(request): - return JSONResponse({ - "status": "healthy", - "server": "vector-indexer-mcp", - "transport": "sse-local" - }) + return JSONResponse( + { + "status": "healthy", + "server": "vector-indexer-mcp", + "transport": "sse-local", + } + ) routes = [ Route("/health", endpoint=health, methods=["GET"]), @@ -90,8 +92,9 @@ async def health(request): def main(): parser = argparse.ArgumentParser() - parser.add_argument("--port", "-p", type=int, - default=int(os.environ.get("MCP_SSE_PORT", "5577"))) + parser.add_argument( + "--port", "-p", type=int, default=int(os.environ.get("MCP_SSE_PORT", "5577")) + ) parser.add_argument("--host", default=os.environ.get("MCP_SSE_HOST", "127.0.0.1")) args = parser.parse_args()