diff --git a/.gitignore b/.gitignore index 1c9e5635..3f3c4f42 100644 --- a/.gitignore +++ b/.gitignore @@ -41,4 +41,5 @@ yarn-error.log* *.pem .claude .venv +.arch __pycache__ diff --git a/packages/pipecat-sdk-python/Agents.md b/packages/pipecat-sdk-python/Agents.md new file mode 100644 index 00000000..0da9f607 --- /dev/null +++ b/packages/pipecat-sdk-python/Agents.md @@ -0,0 +1,256 @@ +# agents.md + +Agent knowledge file for Supermemory Pipecat SDK. This describes how Supermemory integrates with Pipecat for memory-enhanced voice AI. + +## What This Package Does + +Supermemory Pipecat SDK provides a `FrameProcessor` that intercepts conversation frames, retrieves relevant memories from Supermemory, injects them into the LLM context, and stores new messages. + +## Package Structure + +``` +packages/pipecat-sdk-python/ +├── src/supermemory_pipecat/ +│ ├── __init__.py # Exports: SupermemoryPipecatService +│ ├── service.py # Main FrameProcessor class +│ ├── exceptions.py # Error hierarchy +│ └── utils.py # Helpers: get_last_user_message, deduplicate_memories +├── pyproject.toml # Package: supermemory-pipecat +└── README.md +``` + +## Core Class: SupermemoryPipecatService + +Extends `pipecat.processors.frame_processor.FrameProcessor`. + +### Constructor Parameters + +| Parameter | Type | Required | Default | +|-----------|------|----------|---------| +| `user_id` | str | Yes | - | +| `session_id` | str | No | None | +| `api_key` | str | No | `SUPERMEMORY_API_KEY` env | +| `params` | InputParams | No | defaults | +| `base_url` | str | No | `https://api.supermemory.ai` | + +### InputParams Configuration + +```python +InputParams( + search_limit=10, # Max memories to retrieve + search_threshold=0.1, # Similarity threshold (0.0-1.0) + mode="full", # "profile" | "query" | "full" + add_memory="always", # "always" | "never" + add_as_system_message=True, + system_prompt="Based on previous conversations, I recall:\n\n", + position=1, +) +``` + +### Key Mappings + +- `user_id` → `container_tag` (direct, no transformation) +- `session_id` → `custom_id` (in Supermemory API) + +## How It Works + +### Pipeline Position + +``` +[Transport] → [STT] → [UserContext] → [SupermemoryPipecatService] → [LLM] → [TTS] → [Output] +``` + +### Frame Processing Flow + +1. **Intercept**: Catches `LLMContextFrame`, `OpenAILLMContextFrame`, `LLMMessagesFrame` +2. **Extract**: Gets last user message from context +3. **Track**: Stores message in `_conversation_history` (clean, no injections) +4. **Retrieve**: Calls Supermemory `/v4/profile` API +5. **Inject**: Adds formatted memories to context as system message +6. **Store**: Sends last user message to Supermemory (background, non-blocking) +7. **Push**: Forwards enhanced frame downstream + +### Supermemory API Integration + +**Retrieval** - `POST /v4/profile`: +```json +{ + "containerTag": "user-123", + "q": "What's the weather?", + "limit": 10, + "threshold": 0.1 +} +``` + +**Response**: +```json +{ + "profile": { + "static": ["User lives in SF", "Prefers Celsius"], + "dynamic": ["Recently asked about weather"] + }, + "searchResults": { + "results": [{"memory": "User likes sunny weather"}] + } +} +``` + +**Storage** - via `supermemory.memories.add()`: +```python +{ + "content": "User: What's the weather?", + "container_tags": ["user-123"], + "custom_id": "session-456", + "metadata": {"platform": "pipecat"} +} +``` + +## Memory Modes + +| Mode | Static Profile | Dynamic Profile | Search Results | +|------|----------------|-----------------|----------------| +| `profile` | Yes | Yes | No | +| `query` | No | No | Yes | +| `full` | Yes | Yes | Yes | + +## Instance State + +```python +self.user_id: str # User identifier +self.container_tag: str # Same as user_id +self.session_id: Optional[str] # Session grouping +self._conversation_history: List[Dict] # Clean message history +self._last_query: Optional[str] # Dedup tracking +self._supermemory_client # Supermemory SDK client +``` + +## Error Handling + +- Memory retrieval failures: Log warning, continue without memories +- Memory storage failures: Log error, don't crash pipeline +- Frame processing errors: Log error, pass original frame through + +## Sample Usage + +```python +from supermemory_pipecat import SupermemoryPipecatService + +# Create service +memory = SupermemoryPipecatService( + api_key=os.getenv("SUPERMEMORY_API_KEY"), + user_id="user-123", + session_id="conv-456", + params=SupermemoryPipecatService.InputParams( + mode="full", + add_memory="always", + ), +) + +# Add to pipeline +pipeline = Pipeline([ + transport.input(), + stt, + context_aggregator.user(), + memory, # ← Intercepts here, retrieves/injects/stores + llm, + tts, + transport.output(), + context_aggregator.assistant(), +]) +``` + +## Full Working Example + +```python +"""Pipecat + Supermemory Voice Bot""" + +import os +from fastapi import FastAPI, WebSocket +from pipecat.pipeline.pipeline import Pipeline +from pipecat.pipeline.runner import PipelineRunner +from pipecat.pipeline.task import PipelineParams, PipelineTask +from pipecat.processors.aggregators.openai_llm_context import OpenAILLMContext +from pipecat.services.openai.llm import OpenAILLMService +from pipecat.services.openai.tts import OpenAITTSService +from pipecat.services.openai.stt import OpenAISTTService +from pipecat.transports.websocket.fastapi import ( + FastAPIWebsocketParams, + FastAPIWebsocketTransport, +) +from pipecat.audio.vad.silero import SileroVADAnalyzer +from pipecat.serializers.protobuf import ProtobufFrameSerializer + +from supermemory_pipecat import SupermemoryPipecatService + +app = FastAPI() + +SYSTEM_PROMPT = """You are a helpful voice assistant with memory. +Keep responses brief. Your output will be converted to audio.""" + +@app.websocket("/ws") +async def websocket_endpoint(websocket: WebSocket): + await websocket.accept() + + transport = FastAPIWebsocketTransport( + websocket=websocket, + params=FastAPIWebsocketParams( + audio_in_enabled=True, + audio_out_enabled=True, + vad_enabled=True, + vad_analyzer=SileroVADAnalyzer(), + vad_audio_passthrough=True, + serializer=ProtobufFrameSerializer(), + ), + ) + + stt = OpenAISTTService(api_key=os.getenv("OPENAI_API_KEY")) + llm = OpenAILLMService(api_key=os.getenv("OPENAI_API_KEY"), model="gpt-4o-mini") + tts = OpenAITTSService(api_key=os.getenv("OPENAI_API_KEY"), voice="alloy") + + # Supermemory integration + memory = SupermemoryPipecatService( + user_id="test-user", + session_id="voice-session", + ) + + context = OpenAILLMContext([{"role": "system", "content": SYSTEM_PROMPT}]) + context_aggregator = llm.create_context_aggregator(context) + + pipeline = Pipeline([ + transport.input(), + stt, + context_aggregator.user(), + memory, + llm, + tts, + transport.output(), + context_aggregator.assistant(), + ]) + + task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True)) + runner = PipelineRunner(handle_sigint=False) + await runner.run(task) + +if __name__ == "__main__": + import uvicorn + uvicorn.run(app, host="0.0.0.0", port=8001) +``` + +## Key Files Reference + +- `service.py:46` - `SupermemoryPipecatService` class definition +- `service.py:96` - `__init__()` constructor +- `service.py:157` - `_retrieve_memories()` API call +- `service.py:218` - `_store_message()` storage logic +- `service.py:264` - `_enhance_context_with_memories()` injection +- `service.py:328` - `process_frame()` main entry point + +## Differences from Mem0 + +| Aspect | Mem0 | Supermemory | +|--------|------|-------------| +| Identity | `user_id`, `agent_id`, `run_id` | `user_id` only (= container_tag) | +| Retrieval | `memory.search()` | `/v4/profile` (static + dynamic + search) | +| Storage | Full conversation | Last user message only | +| Metadata | `{"platform": "pipecat"}` | `{"platform": "pipecat"}` | +| Session | N/A | `session_id` → `custom_id` | diff --git a/packages/pipecat-sdk-python/README.md b/packages/pipecat-sdk-python/README.md new file mode 100644 index 00000000..7f2b8d70 --- /dev/null +++ b/packages/pipecat-sdk-python/README.md @@ -0,0 +1,174 @@ +# Supermemory Pipecat SDK + +Memory-enhanced conversational AI pipelines with [Supermemory](https://supermemory.ai) and [Pipecat](https://github.com/pipecat-ai/pipecat). + +## Installation + +```bash +pip install supermemory-pipecat +``` + +## Quick Start + +```python +import os +from pipecat.pipeline.pipeline import Pipeline +from pipecat.services.openai import OpenAILLMService, OpenAIUserContextAggregator +from supermemory_pipecat import SupermemoryPipecatService + +# Create memory service +memory = SupermemoryPipecatService( + api_key=os.getenv("SUPERMEMORY_API_KEY"), + user_id="user-123", # Required: used as container_tag + session_id="conversation-456", # Optional: groups memories by session +) + +# Create pipeline with memory +pipeline = Pipeline([ + transport.input(), + stt, + user_context, + memory, # Automatically retrieves and injects relevant memories + llm, + transport.output(), +]) +``` + +## Configuration + +### Parameters + +| Parameter | Type | Required | Description | +| ------------ | ----------- | -------- | ---------------------------------------------------------- | +| `user_id` | str | **Yes** | User identifier - used as container_tag for memory scoping | +| `session_id` | str | No | Session/conversation ID for grouping memories | +| `api_key` | str | No | Supermemory API key (or set `SUPERMEMORY_API_KEY` env var) | +| `params` | InputParams | No | Advanced configuration | +| `base_url` | str | No | Custom API endpoint | + +### Advanced Configuration + +```python +from supermemory_pipecat import SupermemoryPipecatService + +memory = SupermemoryPipecatService( + user_id="user-123", + session_id="conv-456", + params=SupermemoryPipecatService.InputParams( + search_limit=10, # Max memories to retrieve + search_threshold=0.1, # Similarity threshold + mode="full", # "profile", "query", or "full" + add_memory="always", # "always" or "never" + add_as_system_message=True, + system_prompt="Based on previous conversations, I recall:\n\n", + ), +) +``` + +### Memory Modes + +| Mode | Static Profile | Dynamic Profile | Search Results | +| ----------- | -------------- | --------------- | -------------- | +| `"profile"` | Yes | Yes | No | +| `"query"` | No | No | Yes | +| `"full"` | Yes | Yes | Yes | + +## How It Works + +1. **Intercepts context frames** - Listens for `LLMContextFrame` in the pipeline +2. **Tracks conversation** - Maintains clean conversation history (no injected memories) +3. **Retrieves memories** - Queries `/v4/profile` API with user's message +4. **Injects memories** - Formats and adds to LLM context as system message +5. **Stores messages** - Sends last user message to Supermemory (background, non-blocking) + +### What Gets Stored + +Only the last user message is sent to Supermemory: + +``` +User: What's the weather like today? +``` + +Stored as: + +```json +{ + "content": "User: What's the weather like today?", + "container_tags": ["user-123"], + "custom_id": "conversation-456", + "metadata": { "platform": "pipecat" } +} +``` + +## Full Example + +```python +import asyncio +import os +from fastapi import FastAPI, WebSocket +from pipecat.pipeline.pipeline import Pipeline +from pipecat.pipeline.task import PipelineTask +from pipecat.pipeline.runner import PipelineRunner +from pipecat.services.openai import ( + OpenAILLMService, + OpenAIUserContextAggregator, +) +from pipecat.transports.network.fastapi_websocket import ( + FastAPIWebsocketTransport, + FastAPIWebsocketParams, +) +from supermemory_pipecat import SupermemoryPipecatService + +app = FastAPI() + +@app.websocket("/chat") +async def websocket_endpoint(websocket: WebSocket): + await websocket.accept() + + transport = FastAPIWebsocketTransport( + websocket=websocket, + params=FastAPIWebsocketParams(audio_out_enabled=True), + ) + + user_context = OpenAIUserContextAggregator() + + # Supermemory memory service + memory = SupermemoryPipecatService( + user_id="alice", + session_id="session-123", + ) + + llm = OpenAILLMService( + api_key=os.getenv("OPENAI_API_KEY"), + model="gpt-4", + ) + + pipeline = Pipeline([ + transport.input(), + user_context, + memory, + llm, + transport.output(), + ]) + + runner = PipelineRunner() + task = PipelineTask(pipeline) + await runner.run(task) +``` + +## Conversation History + +Access the tracked conversation (without injected memories): + +```python +# Get conversation history +history = memory.get_conversation_history() +# [{"role": "user", "content": "Hello"}, {"role": "user", "content": "What's my name?"}] + +# Clear history +memory.clear_conversation_history() +``` + +## License + +MIT diff --git a/packages/pipecat-sdk-python/pyproject.toml b/packages/pipecat-sdk-python/pyproject.toml new file mode 100644 index 00000000..72baa642 --- /dev/null +++ b/packages/pipecat-sdk-python/pyproject.toml @@ -0,0 +1,80 @@ +[build-system] +requires = ["hatchling"] +build-backend = "hatchling.build" + +[project] +name = "supermemory-pipecat" +version = "0.1.0" +description = "Supermemory integration for Pipecat - memory-enhanced conversational AI pipelines" +readme = "README.md" +license = "MIT" +requires-python = ">=3.8.1" +authors = [ + { name = "Supermemory", email = "support@supermemory.ai" } +] +keywords = [ + "supermemory", + "pipecat", + "memory", + "conversational-ai", + "llm", + "voice-ai", +] +classifiers = [ + "Development Status :: 4 - Beta", + "Intended Audience :: Developers", + "License :: OSI Approved :: MIT License", + "Programming Language :: Python :: 3", + "Programming Language :: Python :: 3.8", + "Programming Language :: Python :: 3.9", + "Programming Language :: Python :: 3.10", + "Programming Language :: Python :: 3.11", + "Programming Language :: Python :: 3.12", + "Topic :: Scientific/Engineering :: Artificial Intelligence", +] +dependencies = [ + "pipecat-ai>=0.0.98", + "supermemory>=3.16.0", + "pydantic>=2.10.0", + "aiohttp>=3.11.0", + "loguru>=0.7.3", +] + +[project.optional-dependencies] +dev = [ + "pytest>=8.3.5", + "pytest-asyncio>=0.24.0", + "mypy>=1.14.1", + "black>=24.8.0", + "isort>=5.13.2", +] + +[project.urls] +Homepage = "https://supermemory.ai" +Documentation = "https://docs.supermemory.ai" +Repository = "https://github.com/supermemoryai/supermemory" + +[tool.hatch.build.targets.wheel] +packages = ["src/supermemory_pipecat"] + +[tool.hatch.build.targets.sdist] +include = [ + "/src", + "/tests", + "/README.md", + "/LICENSE", +] + +[tool.black] +line-length = 100 +target-version = ["py38"] + +[tool.isort] +profile = "black" +line_length = 100 + +[tool.mypy] +python_version = "3.8" +warn_return_any = true +warn_unused_configs = true +disallow_untyped_defs = true diff --git a/packages/pipecat-sdk-python/src/supermemory_pipecat/__init__.py b/packages/pipecat-sdk-python/src/supermemory_pipecat/__init__.py new file mode 100644 index 00000000..35b888cc --- /dev/null +++ b/packages/pipecat-sdk-python/src/supermemory_pipecat/__init__.py @@ -0,0 +1,59 @@ +"""Supermemory Pipecat SDK - Memory-enhanced conversational AI pipelines. + +This package provides seamless integration between Supermemory and Pipecat, +enabling persistent memory and context enhancement for voice AI applications. + +Example: + ```python + from supermemory_pipecat import SupermemoryPipecatService + + # Create memory service + memory = SupermemoryPipecatService( + api_key=os.getenv("SUPERMEMORY_API_KEY"), + user_id="user-123", + ) + + # Add to Pipecat pipeline + pipeline = Pipeline([ + transport.input(), + stt, + user_context, + memory, # Automatically retrieves and injects memories + llm, + transport.output(), + ]) + ``` +""" + +from .service import SupermemoryPipecatService +from .exceptions import ( + SupermemoryPipecatError, + ConfigurationError, + MemoryRetrievalError, + MemoryStorageError, + APIError, + NetworkError, +) +from .utils import ( + get_last_user_message, + deduplicate_memories, + format_memories_to_text, +) + +__version__ = "0.1.0" + +__all__ = [ + # Main service + "SupermemoryPipecatService", + # Exceptions + "SupermemoryPipecatError", + "ConfigurationError", + "MemoryRetrievalError", + "MemoryStorageError", + "APIError", + "NetworkError", + # Utilities + "get_last_user_message", + "deduplicate_memories", + "format_memories_to_text", +] diff --git a/packages/pipecat-sdk-python/src/supermemory_pipecat/exceptions.py b/packages/pipecat-sdk-python/src/supermemory_pipecat/exceptions.py new file mode 100644 index 00000000..1de8094f --- /dev/null +++ b/packages/pipecat-sdk-python/src/supermemory_pipecat/exceptions.py @@ -0,0 +1,58 @@ +"""Custom exceptions for Supermemory Pipecat integration.""" + +from typing import Optional + + +class SupermemoryPipecatError(Exception): + """Base exception for all Supermemory Pipecat errors.""" + + def __init__(self, message: str, original_error: Optional[Exception] = None): + super().__init__(message) + self.message = message + self.original_error = original_error + + def __str__(self) -> str: + if self.original_error: + return f"{self.message}: {self.original_error}" + return self.message + + +class ConfigurationError(SupermemoryPipecatError): + """Raised when there are configuration issues (e.g., missing API key, invalid params).""" + + +class MemoryRetrievalError(SupermemoryPipecatError): + """Raised when memory retrieval operations fail.""" + + +class MemoryStorageError(SupermemoryPipecatError): + """Raised when memory storage operations fail.""" + + +class APIError(SupermemoryPipecatError): + """Raised when Supermemory API requests fail.""" + + def __init__( + self, + message: str, + status_code: Optional[int] = None, + response_text: Optional[str] = None, + original_error: Optional[Exception] = None, + ): + super().__init__(message, original_error) + self.status_code = status_code + self.response_text = response_text + + def __str__(self) -> str: + parts = [self.message] + if self.status_code: + parts.append(f"Status: {self.status_code}") + if self.response_text: + parts.append(f"Response: {self.response_text}") + if self.original_error: + parts.append(f"Cause: {self.original_error}") + return " | ".join(parts) + + +class NetworkError(SupermemoryPipecatError): + """Raised when network operations fail.""" diff --git a/packages/pipecat-sdk-python/src/supermemory_pipecat/service.py b/packages/pipecat-sdk-python/src/supermemory_pipecat/service.py new file mode 100644 index 00000000..4fc08b15 --- /dev/null +++ b/packages/pipecat-sdk-python/src/supermemory_pipecat/service.py @@ -0,0 +1,405 @@ +"""Supermemory Pipecat service integration. + +This module provides a memory service that integrates with Supermemory to store +and retrieve conversational memories, enhancing LLM context with relevant +historical information. +""" + +import asyncio +import os +from typing import Any, Dict, List, Optional + +from loguru import logger +from pydantic import BaseModel, Field + +from pipecat.frames.frames import Frame, LLMContextFrame, LLMMessagesFrame +from pipecat.processors.aggregators.llm_context import LLMContext +from pipecat.processors.aggregators.openai_llm_context import ( + OpenAILLMContext, + OpenAILLMContextFrame, +) +from pipecat.processors.frame_processor import FrameDirection, FrameProcessor + +from .exceptions import ( + APIError, + ConfigurationError, + MemoryRetrievalError, + MemoryStorageError, +) +from .utils import ( + deduplicate_memories, + format_memories_to_text, + get_last_user_message, +) + +try: + import aiohttp +except ImportError: + aiohttp = None # type: ignore + +try: + import supermemory +except ImportError: + supermemory = None # type: ignore + + +class SupermemoryPipecatService(FrameProcessor): + """A memory service that integrates Supermemory with Pipecat pipelines. + + This service intercepts message frames in the pipeline, retrieves relevant + memories from Supermemory, enhances the context, and optionally stores + new conversations. + + Example: + ```python + from supermemory_pipecat import SupermemoryPipecatService + + memory = SupermemoryPipecatService( + api_key=os.getenv("SUPERMEMORY_API_KEY"), + user_id="user-123", + ) + + pipeline = Pipeline([ + transport.input(), + stt, + user_context, + memory, # Memory service enhances context here + llm, + transport.output(), + ]) + ``` + """ + + class InputParams(BaseModel): + """Configuration parameters for Supermemory Pipecat service. + + Parameters: + search_limit: Maximum number of memories to retrieve per query. + search_threshold: Minimum similarity threshold for memory retrieval. + system_prompt: Prefix text for memory context messages. + add_as_system_message: Whether to add memories as system messages. + position: Position to insert memory messages in context. + mode: Memory retrieval mode - "profile", "query", or "full". + add_memory: When to store memories - "always" or "never". + """ + + search_limit: int = Field(default=10, ge=1) + search_threshold: float = Field(default=0.1, ge=0.0, le=1.0) + system_prompt: str = Field(default="Based on previous conversations, I recall:\n\n") + add_as_system_message: bool = Field(default=True) + position: int = Field(default=1) + mode: str = Field(default="full") # "profile", "query", "full" + add_memory: str = Field(default="always") # "always", "never" + + def __init__( + self, + *, + api_key: Optional[str] = None, + user_id: str, + session_id: Optional[str] = None, + params: Optional[InputParams] = None, + base_url: Optional[str] = None, + ): + """Initialize the Supermemory Pipecat service. + + Args: + api_key: The API key for Supermemory. Falls back to SUPERMEMORY_API_KEY env var. + user_id: The user ID - used as container_tag for memory scoping (REQUIRED). + session_id: Session/conversation ID for grouping memories (optional). + params: Configuration parameters for memory retrieval and storage. + base_url: Optional custom base URL for Supermemory API. + + Raises: + ConfigurationError: If API key is missing or user_id not provided. + """ + super().__init__() + + # Get API key + self.api_key = api_key or os.getenv("SUPERMEMORY_API_KEY") + if not self.api_key: + raise ConfigurationError( + "API key is required. Provide api_key parameter or set SUPERMEMORY_API_KEY environment variable." + ) + + # user_id is required and used directly as container_tag + if not user_id: + raise ConfigurationError("user_id is required") + + self.user_id = user_id + self.container_tag = user_id # container_tag = user_id directly + self.session_id = session_id # optional session/conversation ID + + # Configuration + self.params = params or SupermemoryPipecatService.InputParams() + self.base_url = base_url or "https://api.supermemory.ai" + + # Initialize Supermemory client for storage operations + self._supermemory_client = None + if supermemory is not None: + try: + self._supermemory_client = supermemory.Supermemory(api_key=self.api_key) + except Exception as e: + logger.warning(f"Failed to initialize Supermemory client: {e}") + + # Track conversation history separately (clean, no injected memories) + self._conversation_history: List[Dict[str, str]] = [] + + # Track last query to avoid duplicate processing + self._last_query: Optional[str] = None + + logger.info( + f"Initialized SupermemoryPipecatService with " + f"user_id={user_id}, session_id={session_id}" + ) + + async def _retrieve_memories(self, query: str) -> Dict[str, Any]: + """Retrieve relevant memories from Supermemory. + + Args: + query: The query to search for relevant memories. + + Returns: + Dictionary containing profile and search results. + """ + try: + logger.debug(f"Retrieving memories for query: {query[:100]}...") + + payload: Dict[str, Any] = { + "containerTag": self.container_tag, + } + + # Add query for search modes + if self.params.mode != "profile" and query: + payload["q"] = query + payload["limit"] = self.params.search_limit + payload["threshold"] = self.params.search_threshold + + if aiohttp is None: + raise MemoryRetrievalError( + "aiohttp is required for memory retrieval. Install with: pip install aiohttp" + ) + + async with aiohttp.ClientSession() as session: + async with session.post( + f"{self.base_url}/v4/profile", + headers={ + "Content-Type": "application/json", + "Authorization": f"Bearer {self.api_key}", + }, + json=payload, + ) as response: + if not response.ok: + error_text = await response.text() + raise APIError( + "Supermemory profile search failed", + status_code=response.status, + response_text=error_text, + ) + + data = await response.json() + logger.debug( + f"Retrieved memories - static: {len(data.get('profile', {}).get('static', []))}, " + f"dynamic: {len(data.get('profile', {}).get('dynamic', []))}, " + f"search: {len(data.get('searchResults', {}).get('results', []))}" + ) + return data + + except aiohttp.ClientError as e: + logger.error(f"Network error retrieving memories: {e}") + raise MemoryRetrievalError("Network error during memory retrieval", e) + except APIError: + raise + except Exception as e: + logger.error(f"Error retrieving memories: {e}") + raise MemoryRetrievalError("Failed to retrieve memories", e) + + async def _store_message(self, message: Dict[str, str]) -> None: + """Store a single message in Supermemory. + + Args: + message: Message dict with 'role' and 'content' keys. + """ + if self.params.add_memory != "always": + return + + if self._supermemory_client is None: + logger.warning("Supermemory client not initialized, skipping memory storage") + return + + try: + content = message.get("content", "") + if not content or not isinstance(content, str): + return + + # Format: "User: message content" or "Assistant: message content" + role = message.get("role", "user").capitalize() + formatted_content = f"{role}: {content}" + + logger.debug(f"Storing message to Supermemory: {formatted_content[:100]}...") + + # Build storage params + add_params: Dict[str, Any] = { + "content": formatted_content, + "container_tags": [self.container_tag], + "metadata": {"platform": "pipecat"}, + } + if self.session_id: + add_params["custom_id"] = self.session_id + + # Store asynchronously + try: + await self._supermemory_client.memories.add(**add_params) + except TypeError: + # Sync client fallback + self._supermemory_client.memories.add(**add_params) + + logger.debug("Successfully stored message in Supermemory") + + except Exception as e: + # Don't fail the pipeline on storage errors + logger.error(f"Error storing message in Supermemory: {e}") + + def _enhance_context_with_memories( + self, + context: LLMContext, + query: str, + memories_data: Dict[str, Any], + ) -> None: + """Enhance the LLM context with relevant memories. + + Args: + context: The LLM context to enhance. + query: The query used for retrieval. + memories_data: Raw memory data from Supermemory API. + """ + # Skip if same query (avoid duplicate processing) + if self._last_query == query: + return + + self._last_query = query + + # Extract and deduplicate memories + profile = memories_data.get("profile", {}) + search_results = memories_data.get("searchResults", {}) + + deduplicated = deduplicate_memories( + static=profile.get("static", []), + dynamic=profile.get("dynamic", []), + search_results=search_results.get("results", []), + ) + + # Check if we have any memories + total_memories = ( + len(deduplicated["static"]) + + len(deduplicated["dynamic"]) + + len(deduplicated["search_results"]) + ) + + if total_memories == 0: + logger.debug("No memories found to inject") + return + + # Format memories based on mode + include_profile = self.params.mode in ("profile", "full") + include_search = self.params.mode in ("query", "full") + + memory_text = format_memories_to_text( + deduplicated, + system_prompt=self.params.system_prompt, + include_static=include_profile, + include_dynamic=include_profile, + include_search=include_search, + ) + + if not memory_text: + return + + # Inject memories into context + if self.params.add_as_system_message: + context.add_message({"role": "system", "content": memory_text}) + else: + context.add_message({"role": "user", "content": memory_text}) + + logger.debug(f"Enhanced context with {total_memories} memories") + + async def process_frame(self, frame: Frame, direction: FrameDirection) -> None: + """Process incoming frames, intercept context frames for memory integration. + + Args: + frame: The incoming frame to process. + direction: The direction of frame flow in the pipeline. + """ + await super().process_frame(frame, direction) + + context = None + messages = None + + # Handle different frame types + if isinstance(frame, (LLMContextFrame, OpenAILLMContextFrame)): + context = frame.context + elif isinstance(frame, LLMMessagesFrame): + messages = frame.messages + context = LLMContext(messages) + + if context: + try: + # Get messages from context + context_messages = context.get_messages() + + # Find latest user message for memory query + latest_user_message = get_last_user_message(context_messages) + + if latest_user_message: + # Track the user message in our conversation history (clean) + user_msg = {"role": "user", "content": latest_user_message} + + # Only add if it's a new message (not already tracked) + if ( + not self._conversation_history + or self._conversation_history[-1].get("content") != latest_user_message + ): + self._conversation_history.append(user_msg) + + # Retrieve memories from Supermemory + try: + memories_data = await self._retrieve_memories(latest_user_message) + + # Enhance context with memories + self._enhance_context_with_memories( + context, latest_user_message, memories_data + ) + except (MemoryRetrievalError, APIError) as e: + # Log but don't fail the pipeline + logger.warning(f"Memory retrieval failed, continuing without memories: {e}") + + # Store the last user message (runs in background, non-blocking) + asyncio.create_task(self._store_message(user_msg)) + + # Pass the frame downstream + if messages is not None: + # For LLMMessagesFrame, create new frame with enhanced messages + await self.push_frame(LLMMessagesFrame(context.get_messages())) + else: + # For context frames, pass the enhanced frame + await self.push_frame(frame) + + except Exception as e: + logger.error(f"Error processing frame with Supermemory: {e}") + # Still pass the original frame through on error + await self.push_frame(frame) + else: + # Non-context frames pass through unchanged + await self.push_frame(frame, direction) + + def get_conversation_history(self) -> List[Dict[str, str]]: + """Get the tracked conversation history (without injected memories). + + Returns: + List of message dicts with 'role' and 'content'. + """ + return self._conversation_history.copy() + + def clear_conversation_history(self) -> None: + """Clear the tracked conversation history.""" + self._conversation_history.clear() + self._last_query = None diff --git a/packages/pipecat-sdk-python/src/supermemory_pipecat/utils.py b/packages/pipecat-sdk-python/src/supermemory_pipecat/utils.py new file mode 100644 index 00000000..7d76be87 --- /dev/null +++ b/packages/pipecat-sdk-python/src/supermemory_pipecat/utils.py @@ -0,0 +1,133 @@ +"""Utility functions for Supermemory Pipecat integration.""" + +from typing import Any, Dict, List, Optional + + +def get_last_user_message(messages: List[Dict[str, Any]]) -> Optional[str]: + """ + Extract the last user message from a list of messages. + + Args: + messages: List of message dictionaries with 'role' and 'content' keys + + Returns: + The content of the last user message, or None if not found + """ + for message in reversed(messages): + if message.get("role") == "user": + content = message.get("content", "") + if isinstance(content, str): + return content + elif isinstance(content, list): + # Handle content that is an array of content parts + text_parts = [] + for part in content: + if isinstance(part, dict) and part.get("type") == "text": + text_parts.append(part.get("text", "")) + elif isinstance(part, str): + text_parts.append(part) + return " ".join(text_parts) + return None + + +def deduplicate_memories( + static: Optional[List[Any]] = None, + dynamic: Optional[List[Any]] = None, + search_results: Optional[List[Any]] = None, +) -> Dict[str, List[str]]: + """ + Deduplicates memory items across sources. + Priority: Static > Dynamic > Search Results. + + Args: + static: List of static profile memories + dynamic: List of dynamic profile memories + search_results: List of search result memories + + Returns: + Dictionary with deduplicated 'static', 'dynamic', and 'search_results' lists + """ + static_items = static or [] + dynamic_items = dynamic or [] + search_items = search_results or [] + + def extract_memory_text(item: Any) -> Optional[str]: + if isinstance(item, dict): + item = item.get("memory") + if isinstance(item, str): + trimmed = item.strip() + return trimmed or None + return None + + static_memories: List[str] = [] + seen_memories: set = set() + + for item in static_items: + memory = extract_memory_text(item) + if memory is not None: + static_memories.append(memory) + seen_memories.add(memory) + + dynamic_memories: List[str] = [] + for item in dynamic_items: + memory = extract_memory_text(item) + if memory is not None and memory not in seen_memories: + dynamic_memories.append(memory) + seen_memories.add(memory) + + search_memories: List[str] = [] + for item in search_items: + memory = extract_memory_text(item) + if memory is not None and memory not in seen_memories: + search_memories.append(memory) + seen_memories.add(memory) + + return { + "static": static_memories, + "dynamic": dynamic_memories, + "search_results": search_memories, + } + + +def format_memories_to_text( + memories: Dict[str, List[str]], + system_prompt: str = "Based on previous conversations, I recall:\n\n", + include_static: bool = True, + include_dynamic: bool = True, + include_search: bool = True, +) -> str: + """ + Format deduplicated memories into a text string for injection. + + Args: + memories: Dictionary with 'static', 'dynamic', 'search_results' lists + system_prompt: Prefix text for the memory content + include_static: Whether to include static profile memories + include_dynamic: Whether to include dynamic profile memories + include_search: Whether to include search result memories + + Returns: + Formatted memory text string + """ + sections = [] + + static = memories.get("static", []) + dynamic = memories.get("dynamic", []) + search_results = memories.get("search_results", []) + + if include_static and static: + sections.append("## User Profile (Persistent)") + sections.append("\n".join(f"- {item}" for item in static)) + + if include_dynamic and dynamic: + sections.append("## Recent Context") + sections.append("\n".join(f"- {item}" for item in dynamic)) + + if include_search and search_results: + sections.append("## Relevant Memories") + sections.append("\n".join(f"- {item}" for item in search_results)) + + if not sections: + return "" + + return f"{system_prompt}\n" + "\n\n".join(sections) diff --git a/packages/pipecat-sdk-python/tests/__init__.py b/packages/pipecat-sdk-python/tests/__init__.py new file mode 100644 index 00000000..e20d656a --- /dev/null +++ b/packages/pipecat-sdk-python/tests/__init__.py @@ -0,0 +1 @@ +"""Tests for Supermemory Pipecat SDK."""