feat(memory): Switch embedder from sentence-transformers to fastembed — 2GB to 100MB install

Co-Authored-By: @memory <memory@aipass>
This commit is contained in:
AIOSAI
2026-05-07 17:43:17 -07:00
co-authored by @memory
parent dd52c8aad8
commit 7dc84c3e31
3 changed files with 53 additions and 258 deletions
@@ -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:
@@ -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}")
+30 -91
View File
@@ -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()