pipecat-sdk

This commit is contained in:
Prasanna721 2026-01-09 16:57:28 -08:00
parent 4850856920
commit 7aea611775
9 changed files with 1167 additions and 0 deletions

1
.gitignore vendored
View file

@ -41,4 +41,5 @@ yarn-error.log*
*.pem
.claude
.venv
.arch
__pycache__

View file

@ -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` |

View file

@ -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

View file

@ -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

View file

@ -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",
]

View file

@ -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."""

View file

@ -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

View file

@ -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)

View file

@ -0,0 +1 @@
"""Tests for Supermemory Pipecat SDK."""