diff --git a/src/aipass/memory/apps/handlers/vector/embed_subprocess.py b/src/aipass/memory/apps/handlers/vector/embed_subprocess.py index cac58ebd..a303cdf5 100644 --- a/src/aipass/memory/apps/handlers/vector/embed_subprocess.py +++ b/src/aipass/memory/apps/handlers/vector/embed_subprocess.py @@ -1,16 +1,16 @@ # =================== AIPass ==================== # Name: embed_subprocess.py # Description: Embedding Subprocess Handler -# Version: 1.0.0 +# Version: 2.0.0 # Created: 2026-03-12 -# Modified: 2026-03-12 +# Modified: 2026-05-07 # ============================================= """ Embedding Subprocess Handler -Called via subprocess from rollover orchestrator to ensure sentence-transformers -and torch run in the memory-specific venv (AIPASS_MEMORY_PYTHON). +Called via subprocess from rollover orchestrator to generate embeddings +using fastembed (ONNX runtime, no torch dependency). Input: JSON on stdin with texts to encode Output: JSON on stdout with embeddings @@ -21,7 +21,7 @@ import json def main(): - """Process embedding request from stdin JSON""" + """Process embedding request from stdin JSON.""" try: input_data = json.load(sys.stdin) texts = input_data.get("texts", []) @@ -30,40 +30,19 @@ def main(): print(json.dumps({"success": True, "embeddings": [], "count": 0, "dimension": 384})) return - # Import here — runs in memory venv where these are installed - from sentence_transformers import SentenceTransformer - import torch + from fastembed import TextEmbedding - model = SentenceTransformer("all-MiniLM-L6-v2") + model = TextEmbedding("all-MiniLM-L6-v2") - use_gpu = torch.cuda.is_available() - if use_gpu: - model = model.to("cuda") - batch_size = 64 - else: - batch_size = 16 - - # Pre-sort by length (reduces padding waste) sorted_pairs = sorted(enumerate(texts), key=lambda x: len(x[1])) sorted_indices, sorted_texts = zip(*sorted_pairs) - # Encode - embeddings = model.encode( - list(sorted_texts), - batch_size=batch_size, - convert_to_tensor=False, - normalize_embeddings=True, - show_progress_bar=False, - ) + embeddings = list(model.embed(list(sorted_texts))) - # Restore original order ordered = [None] * len(texts) for orig_idx, sorted_idx in enumerate(sorted_indices): ordered[sorted_idx] = embeddings[orig_idx].tolist() - if use_gpu: - torch.cuda.empty_cache() - print(json.dumps({"success": True, "embeddings": ordered, "count": len(ordered), "dimension": 384})) except Exception as e: diff --git a/src/aipass/memory/apps/handlers/vector/embedder.py b/src/aipass/memory/apps/handlers/vector/embedder.py index d01e8e53..b4af4ef5 100644 --- a/src/aipass/memory/apps/handlers/vector/embedder.py +++ b/src/aipass/memory/apps/handlers/vector/embedder.py @@ -1,31 +1,22 @@ # =================== AIPass ==================== # Name: embedder.py # Description: Vector Embedding Handler -# Version: 0.2.0 +# Version: 0.3.0 # Created: 2025-11-16 -# Modified: 2026-03-06 +# Modified: 2026-05-07 # ============================================= """ Vector Embedding Handler -Generates semantic embeddings using sentence-transformers/all-MiniLM-L6-v2. -Implements production best practices from research. +Generates semantic embeddings using fastembed/all-MiniLM-L6-v2 (ONNX). Purpose: Convert text memories into 384-dimensional vectors for semantic search. Optimized for batch processing (100 lines during rollover). -Best Practices Applied: - - Pre-sort by length (30% padding reduction) - - Built-in normalization (L2 distance requirement) - - GPU memory cleanup (prevent VRAM leaks) - - Batch size optimization (64 GPU, 16 CPU) - - Singleton pattern (model loaded once) - Dependencies (optional): - - sentence-transformers - - torch + - fastembed """ from typing import List, Dict, Any @@ -35,147 +26,58 @@ from aipass.memory.apps.handlers.json import json_handler logger = get_system_logger() -# No service imports - handlers are pure workers (3-tier architecture) -# No module imports (handler independence) - - -# ============================================================================= -# EMBEDDING SERVICE (Singleton) -# ============================================================================= - class EmbeddingService: - """ - Production-ready embedding service - - Implements best practices: - - Batch size optimization (64 GPU, 16 CPU) - - Pre-sorting by length (reduces padding waste 30%) - - Built-in normalization (critical for L2 distance) - - GPU memory cleanup (prevents VRAM leaks) - """ + """Embedding service using fastembed (ONNX runtime, no torch dependency).""" def __init__(self, model_name: str = "all-MiniLM-L6-v2"): - """ - Initialize embedding service - - Args: - model_name: HuggingFace model identifier - - Raises: - ImportError: If sentence-transformers or torch are not installed - """ - # Late imports (heavy optional dependencies) try: - import torch - from sentence_transformers import SentenceTransformer + from fastembed import TextEmbedding except ImportError as e: logger.info(f"[embedder] Optional ML dependencies not available: {e}") - raise ImportError( - f"Embedding requires sentence-transformers and torch. " - f"Install with: pip install sentence-transformers torch. " - f"Original error: {e}" - ) + raise ImportError(f"Embedding requires fastembed. Install with: pip install fastembed. Original error: {e}") self.model_name = model_name - self.model = SentenceTransformer(model_name) - - # GPU optimization if available - self.use_gpu = torch.cuda.is_available() - if self.use_gpu: - self.model = self.model.to("cuda") - self.batch_size = 64 - else: - self.batch_size = 16 - - self.dimension = 384 # all-MiniLM-L6-v2 output dimension + self.model = TextEmbedding(model_name) + self.dimension = 384 def encode_batch(self, texts: List[str]) -> Dict[str, Any]: - """ - Encode batch of texts with all optimizations - - Best practices applied: - 1. Pre-sort by length (reduces padding waste) - 2. Batch processing (optimal batch size) - 3. Built-in normalization (L2 distance requirement) - 4. GPU cleanup (prevent VRAM leaks) - - Args: - texts: List of text strings to encode - - Returns: - Dict with embeddings and metadata - """ - import torch - + """Encode texts to embeddings with pre-sort optimization.""" if not texts: return {"embeddings": [], "count": 0, "dimension": self.dimension} - # Pre-sort by length (reduces padding waste by 30%) sorted_pairs = sorted(enumerate(texts), key=lambda x: len(x[1])) sorted_indices: list[int] = [p[0] for p in sorted_pairs] sorted_text_list: list[str] = [p[1] for p in sorted_pairs] - # Encode with optimal settings - embeddings = self.model.encode( - sorted_text_list, - batch_size=self.batch_size, - convert_to_tensor=False, # Return numpy for Chroma - normalize_embeddings=True, # Critical for L2 distance - show_progress_bar=False, - ) + embeddings = list(self.model.embed(sorted_text_list)) - # Restore original order ordered_embeddings: List[Any] = [None] * len(texts) for original_idx, sorted_idx in enumerate(sorted_indices): ordered_embeddings[sorted_idx] = embeddings[original_idx] - # Cleanup GPU memory if used - if self.use_gpu: - torch.cuda.empty_cache() - return {"embeddings": ordered_embeddings, "count": len(ordered_embeddings), "dimension": self.dimension} -# Global service instance (singleton pattern) _embedding_service = None def _get_service() -> EmbeddingService: - """ - Get or create embedding service singleton - - Lazy initialization - model loaded on first use - """ global _embedding_service if _embedding_service is None: _embedding_service = EmbeddingService() return _embedding_service -# ============================================================================= -# PUBLIC API -# ============================================================================= - - def encode_batch(texts: List[str]) -> Dict[str, Any]: """ - Encode batch of texts to embeddings - - This is the main public API. Delegates to singleton service - to avoid reloading the model. + Encode batch of texts to embeddings. Args: texts: List of text strings to encode Returns: Dict with embeddings and metadata - - Example: - result = encode_batch(["memory 1", "memory 2"]) - if result['success']: - embeddings = result['embeddings'] - # Each embedding is 384-dim numpy array """ if not texts: return {"success": True, "embeddings": [], "count": 0, "message": "No texts provided"} @@ -197,45 +99,27 @@ def encode_batch(texts: List[str]) -> Dict[str, Any]: def encode_memories(memories: List[Dict[str, Any]]) -> Dict[str, Any]: """ - Encode memory entries to embeddings - - Extracts text from memory entries and encodes them. - Preserves original memory structure for metadata. + Encode memory entries to embeddings. Args: memories: List of memory entry dicts (from extraction) Returns: Dict with embeddings and original memories - - Example: - memories = [{"content": "...", "timestamp": "..."}] - result = encode_memories(memories) - embeddings = result['embeddings'] - original = result['memories'] """ if not memories: return {"success": True, "embeddings": [], "memories": [], "count": 0, "message": "No memories provided"} - # Extract text content from memories texts = [] for memory in memories: - # Try common fields for text content - text = ( - memory.get("content") - or memory.get("text") - or memory.get("message") - or str(memory) # Fallback to string representation - ) + text = memory.get("content") or memory.get("text") or memory.get("message") or str(memory) texts.append(text) - # Encode texts encode_result = encode_batch(texts) if not encode_result["success"]: return encode_result - # Combine embeddings with original memories return { "success": True, "embeddings": encode_result["embeddings"], @@ -246,12 +130,7 @@ def encode_memories(memories: List[Dict[str, Any]]) -> Dict[str, Any]: def get_model_info() -> Dict[str, Any]: - """ - Get embedding model information - - Returns: - Dict with model metadata - """ + """Get embedding model information.""" try: service = _get_service() @@ -259,8 +138,6 @@ def get_model_info() -> Dict[str, Any]: "success": True, "model_name": service.model_name, "dimension": service.dimension, - "batch_size": service.batch_size, - "gpu_enabled": service.use_gpu, } except Exception as e: logger.warning(f"[embedder] Failed to get model info: {e}") diff --git a/src/aipass/memory/tests/test_vector.py b/src/aipass/memory/tests/test_vector.py index 38a48cb7..63f2287d 100644 --- a/src/aipass/memory/tests/test_vector.py +++ b/src/aipass/memory/tests/test_vector.py @@ -2,7 +2,7 @@ # META DATA HEADER # Name: tests/test_vector.py # Date: 2026-04-03 -# Version: 1.0.0 +# Version: 2.0.0 # Category: memory/tests # ============================================= @@ -10,12 +10,12 @@ Covers: - vector/embedder.py EmbeddingService class (init, encode_batch with - pre-sort by length and order restoration, GPU cleanup path) + pre-sort by length and order restoration) - vector/embedder.py Public API functions (encode_batch, encode_memories, get_model_info) - vector/embedder.py Singleton management (_get_service, global reset) -All tests use mocks/tmp_path -- no live sentence-transformers, torch, or GPU access. +All tests use mocks/tmp_path -- no live fastembed or ONNX access. """ import sys @@ -29,7 +29,7 @@ pytest.importorskip("chromadb") # --------------------------------------------------------------------------- -# Import helper -- torch and sentence_transformers must be mocked +# Import helper -- fastembed must be mocked # --------------------------------------------------------------------------- @@ -39,14 +39,10 @@ def _import_embedder(monkeypatch): Returns: Tuple of (embedder module, dict of mock objects) """ - mock_torch = MagicMock() - mock_torch.cuda.is_available.return_value = False - monkeypatch.setitem(sys.modules, "torch", mock_torch) - - mock_st = MagicMock() + mock_fastembed = MagicMock() mock_model = MagicMock() - mock_st.SentenceTransformer.return_value = mock_model - monkeypatch.setitem(sys.modules, "sentence_transformers", mock_st) + mock_fastembed.TextEmbedding.return_value = mock_model + monkeypatch.setitem(sys.modules, "fastembed", mock_fastembed) # Clear cached module for fresh import sys.modules.pop("aipass.memory.apps.handlers.vector.embedder", None) @@ -57,8 +53,7 @@ def _import_embedder(monkeypatch): from aipass.memory.apps.handlers.vector import embedder return embedder, { - "torch": mock_torch, - "st": mock_st, + "fastembed": mock_fastembed, "model": mock_model, } @@ -90,8 +85,8 @@ class TestPublicEncodeBatch: embedder, mocks = _import_embedder(monkeypatch) _reset_globals(embedder) - fake_embeddings = np.array([[0.1, 0.2, 0.3], [0.4, 0.5, 0.6]]) - mocks["model"].encode.return_value = fake_embeddings + fake_embeddings = [np.array([0.1, 0.2, 0.3]), np.array([0.4, 0.5, 0.6])] + mocks["model"].embed.return_value = iter(fake_embeddings) result = embedder.encode_batch(["hello world", "test text"]) @@ -103,7 +98,7 @@ class TestPublicEncodeBatch: embedder, mocks = _import_embedder(monkeypatch) _reset_globals(embedder) - mocks["model"].encode.side_effect = RuntimeError("CUDA out of memory") + mocks["model"].embed.side_effect = RuntimeError("ONNX runtime error") result = embedder.encode_batch(["some text"]) @@ -114,7 +109,7 @@ class TestPublicEncodeBatch: embedder, mocks = _import_embedder(monkeypatch) _reset_globals(embedder) - mocks["st"].SentenceTransformer.side_effect = RuntimeError("Model not found") + mocks["fastembed"].TextEmbedding.side_effect = RuntimeError("Model not found") result = embedder.encode_batch(["some text"]) @@ -145,8 +140,8 @@ class TestPublicEncodeMemories: embedder, mocks = _import_embedder(monkeypatch) _reset_globals(embedder) - fake_embeddings = np.array([[0.1, 0.2]]) - mocks["model"].encode.return_value = fake_embeddings + fake_embeddings = [np.array([0.1, 0.2])] + mocks["model"].embed.return_value = iter(fake_embeddings) memories = [{"content": "Important observation", "timestamp": "2026-01-01"}] result = embedder.encode_memories(memories) @@ -154,8 +149,7 @@ class TestPublicEncodeMemories: assert result["success"] is True assert result["count"] == 1 assert result["memories"] == memories - # Verify the model was called with extracted text - call_args = mocks["model"].encode.call_args + call_args = mocks["model"].embed.call_args texts_passed = call_args[0][0] assert texts_passed == ["Important observation"] @@ -163,15 +157,15 @@ class TestPublicEncodeMemories: embedder, mocks = _import_embedder(monkeypatch) _reset_globals(embedder) - fake_embeddings = np.array([[0.1, 0.2]]) - mocks["model"].encode.return_value = fake_embeddings + fake_embeddings = [np.array([0.1, 0.2])] + mocks["model"].embed.return_value = iter(fake_embeddings) memories = [{"text": "Session summary", "date": "2026-02-01"}] result = embedder.encode_memories(memories) assert result["success"] is True assert result["count"] == 1 - call_args = mocks["model"].encode.call_args + call_args = mocks["model"].embed.call_args texts_passed = call_args[0][0] assert texts_passed == ["Session summary"] @@ -179,16 +173,15 @@ class TestPublicEncodeMemories: embedder, mocks = _import_embedder(monkeypatch) _reset_globals(embedder) - fake_embeddings = np.array([[0.1, 0.2]]) - mocks["model"].encode.return_value = fake_embeddings + fake_embeddings = [np.array([0.1, 0.2])] + mocks["model"].embed.return_value = iter(fake_embeddings) memories = [{"arbitrary_key": "value123", "number": 42}] result = embedder.encode_memories(memories) assert result["success"] is True assert result["count"] == 1 - # The fallback is str(memory) which includes the full dict repr - call_args = mocks["model"].encode.call_args + call_args = mocks["model"].embed.call_args texts_passed = call_args[0][0] assert "arbitrary_key" in texts_passed[0] @@ -196,7 +189,7 @@ class TestPublicEncodeMemories: embedder, mocks = _import_embedder(monkeypatch) _reset_globals(embedder) - mocks["model"].encode.side_effect = RuntimeError("Encoding crashed") + mocks["model"].embed.side_effect = RuntimeError("Encoding crashed") memories = [{"content": "test memory"}] result = embedder.encode_memories(memories) @@ -207,8 +200,8 @@ class TestPublicEncodeMemories: embedder, mocks = _import_embedder(monkeypatch) _reset_globals(embedder) - fake_embeddings = np.array([[0.1], [0.2], [0.3]]) - mocks["model"].encode.return_value = fake_embeddings + fake_embeddings = [np.array([0.1]), np.array([0.2]), np.array([0.3])] + mocks["model"].embed.return_value = iter(fake_embeddings) memories = [ {"content": "first"}, @@ -221,11 +214,10 @@ class TestPublicEncodeMemories: assert result["count"] == 3 assert result["memories"] is memories - call_args = mocks["model"].encode.call_args + call_args = mocks["model"].embed.call_args texts_passed = call_args[0][0] assert texts_passed[0] == "first" assert texts_passed[1] == "second" - # Third falls back to str() assert "third" in texts_passed[2] @@ -246,14 +238,12 @@ class TestPublicGetModelInfo: assert result["success"] is True assert result["model_name"] == "all-MiniLM-L6-v2" assert result["dimension"] == 384 - assert result["batch_size"] == 16 # CPU batch size (GPU is mocked off) - assert result["gpu_enabled"] is False def test_service_init_failure_returns_error(self, monkeypatch): embedder, mocks = _import_embedder(monkeypatch) _reset_globals(embedder) - mocks["st"].SentenceTransformer.side_effect = ImportError("no model") + mocks["fastembed"].TextEmbedding.side_effect = ImportError("no model") result = embedder.get_model_info() @@ -273,28 +263,21 @@ class TestEmbeddingServiceEncodeBatch: embedder, mocks = _import_embedder(monkeypatch) _reset_globals(embedder) - # Track what texts the model receives (should be sorted by length) received_texts: list[Any] = [] - def fake_encode(texts, **kwargs): + def fake_embed(texts): received_texts.append(list(texts)) - # Return embeddings matching the sorted input length - return np.array([[float(i)] * 3 for i in range(len(texts))]) + return iter([np.array([float(i)] * 3) for i in range(len(texts))]) - mocks["model"].encode.side_effect = fake_encode + mocks["model"].embed.side_effect = fake_embed service = embedder.EmbeddingService() texts = ["long text here", "ab", "medium text"] result = service.encode_batch(texts) - # Model should receive texts sorted by length assert received_texts[0] == ["ab", "medium text", "long text here"] - # But returned embeddings should be in original order assert result["count"] == 3 - # Index 0 was "long text here" (sorted position 2) -> embedding [2,2,2] - # Index 1 was "ab" (sorted position 0) -> embedding [0,0,0] - # Index 2 was "medium text" (sorted position 1) -> embedding [1,1,1] embs = result["embeddings"] np.testing.assert_array_equal(embs[0], [2.0, 2.0, 2.0]) np.testing.assert_array_equal(embs[1], [0.0, 0.0, 0.0]) @@ -315,54 +298,10 @@ class TestEmbeddingServiceEncodeBatch: embedder, mocks = _import_embedder(monkeypatch) _reset_globals(embedder) - mocks["model"].encode.return_value = np.array([[0.5, 0.6]]) + mocks["model"].embed.return_value = iter([np.array([0.5, 0.6])]) service = embedder.EmbeddingService() result = service.encode_batch(["only one"]) assert result["count"] == 1 np.testing.assert_array_equal(result["embeddings"][0], [0.5, 0.6]) - - -# =========================================================================== -# Tests: EmbeddingService -- GPU path -# =========================================================================== - - -class TestEmbeddingServiceGPU: - """Test EmbeddingService GPU detection and cleanup.""" - - def test_gpu_enabled_sets_larger_batch_size(self, monkeypatch): - embedder, mocks = _import_embedder(monkeypatch) - _reset_globals(embedder) - - mocks["torch"].cuda.is_available.return_value = True - - service = embedder.EmbeddingService() - - assert service.use_gpu is True - assert service.batch_size == 64 - - def test_gpu_cache_cleared_after_encode(self, monkeypatch): - embedder, mocks = _import_embedder(monkeypatch) - _reset_globals(embedder) - - mocks["torch"].cuda.is_available.return_value = True - mocks["model"].encode.return_value = np.array([[0.1, 0.2]]) - - service = embedder.EmbeddingService() - service.encode_batch(["test text"]) - - mocks["torch"].cuda.empty_cache.assert_called_once() - - def test_cpu_does_not_clear_gpu_cache(self, monkeypatch): - embedder, mocks = _import_embedder(monkeypatch) - _reset_globals(embedder) - - mocks["torch"].cuda.is_available.return_value = False - mocks["model"].encode.return_value = np.array([[0.1, 0.2]]) - - service = embedder.EmbeddingService() - service.encode_batch(["test text"]) - - mocks["torch"].cuda.empty_cache.assert_not_called()