supermemory/packages/pipecat-sdk-python/src/supermemory_pipecat/utils.py
Dhravya 03773c4f2e
feat(python-sdks): SDK-level cross-source memory deduplication (#1532)
## Stack Context

Part 2 of a 3-PR stack moving memory deduplication into the SDKs. See `sdk-dedup/tools-ts` (parent) for the full context and the TypeScript implementation this mirrors.

## What?

Port the normalized, priority-ordered (`static > dynamic > search`) profile deduplication into the Python SDKs.

- Each request injects one **owned memory block that replaces** the prior block rather than accumulating.
- Dedup is **request-local** (no shared state), so it stays correct under concurrency.

Covers OpenAI, Agent Framework (middleware + context provider), Cartesia, and Pipecat.

## Why?

Keeps the Python SDKs at behavioral parity with the TypeScript SDK so all integrations deduplicate memory the same way.

## Testing

- OpenAI: 31 passed, 11 skipped (live)
- Agent Framework: 59 passed
- Cartesia: 8 passed
- Pipecat: 8 passed

🤖 Generated with [Claude Code](https://claude.com/claude-code)

<!-- CURSOR_SUMMARY -->
---

> [!NOTE]
> **Medium Risk**
> Changes memory formatting and system-prompt injection across multiple SDK integrations; incorrect dedup or replacement could alter LLM context, but there is no auth or data-store risk.
>
> **Overview**
> Ports **normalized cross-source memory deduplication** and **replace-not-append injection** into the Python OpenAI, Agent Framework, Cartesia, and Pipecat packages so they match the TypeScript SDK behavior.
>
> **Deduplication** uses request-local keys: strip optional `[YYYY-MM-DD]` prefixes, normalize whitespace, and compare with `casefold`, with priority **static → dynamic → search**. In **`query` mode**, profile static/dynamic are excluded from dedup input so facts that only appear in search (or overlap profile) are not dropped before formatting.
>
> **Injection** no longer appends memory text every turn. OpenAI and Agent Framework middleware **strip prior owned `<supermemory context="user-memories" readonly>` blocks** and **replace** them once per request while keeping the caller’s system instructions; extra system messages lose stale blocks only. New helpers (`strip`/`replace`/`wrap`) live in each package’s utils.
>
> Tests cover normalized fact variants, query-mode search retention, and stale block replacement.
>
> <sup>Reviewed by [Cursor Bugbot](https://cursor.com/bugbot) for commit 42f308b224. Bugbot is set up for automated code reviews on this repo. Configure [here](https://www.cursor.com/dashboard/bugbot).</sup>
<!-- /CURSOR_SUMMARY -->
2026-09-01 06:10:36 +00:00

192 lines
6 KiB
Python

"""Utility functions for Supermemory Pipecat integration."""
import re
from datetime import datetime, timezone
from typing import Any, Dict, List, Union
_DYNAMIC_DATE_PREFIX = re.compile(
r"^\s*(?:\[recent\]\s*)?(?:\[\d{4}-\d{2}-\d{2}\]\s*)?",
re.IGNORECASE,
)
_USER_MEMORIES_TAG_PATTERN = re.compile(
r"<\s*/?\s*user_memories\b[^>]*>",
re.IGNORECASE,
)
def escape_memory_delimiters(text: str) -> str:
"""Neutralize reserved memory-wrapper tags inside formatted content."""
return _USER_MEMORIES_TAG_PATTERN.sub(
lambda match: match.group(0).replace("<", "&lt;").replace(">", "&gt;"),
text,
)
def get_last_user_message(messages: List[Dict[str, Any]]) -> str | None:
"""Extract the last user message content from a list of messages."""
for msg in reversed(messages):
content = msg.get("content")
if msg.get("role") == "user" and isinstance(content, str):
return content
return None
def format_relative_time(iso_timestamp: str) -> str:
"""Convert ISO timestamp to relative time string.
Format rules:
- [just now] - within 30 minutes
- [Xmins ago] - 30-60 minutes
- [X hrs ago] - less than 1 day
- [Xd ago] - less than 1 week
- [X Jul] - more than 1 week, same year
- [X Jul, 2023] - different year
"""
try:
dt = datetime.fromisoformat(iso_timestamp.replace("Z", "+00:00"))
now = datetime.now(timezone.utc)
diff = now - dt
seconds = diff.total_seconds()
minutes = seconds / 60
hours = seconds / 3600
days = seconds / 86400
if minutes < 30:
return "just now"
elif minutes < 60:
return f"{int(minutes)}mins ago"
elif hours < 24:
return f"{int(hours)} hrs ago"
elif days < 7:
return f"{int(days)}d ago"
elif dt.year == now.year:
return f"{dt.day} {dt.strftime('%b')}"
else:
return f"{dt.day} {dt.strftime('%b')}, {dt.year}"
except Exception:
return ""
def _field(item: Any, *names: str, default: Any = None) -> Any:
"""Read a field from a dict or pydantic/SDK model.
Accepts camelCase and snake_case names so helpers work with both raw JSON
dicts and Stainless-generated response models.
"""
if item is None:
return default
if isinstance(item, dict):
for name in names:
if name in item and item[name] is not None:
return item[name]
return default
for name in names:
value = getattr(item, name, None)
if value is not None:
return value
return default
def deduplicate_memories(
static: List[str],
dynamic: List[str],
search_results: List[Any],
) -> Dict[str, Union[List[str], List[Any]]]:
"""Deduplicate memories. Priority: static > dynamic > search.
Args:
static: List of static memory strings.
dynamic: List of dynamic memory strings.
search_results: Search result dicts or pydantic models with a memory field.
"""
seen: set[str] = set()
def comparison_key(memory: str) -> str:
# Dynamic profile entries are date-labelled by the API while search
# results contain the same memory without that presentation prefix.
without_prefix = _DYNAMIC_DATE_PREFIX.sub("", memory.strip())
return " ".join(without_prefix.split()).casefold()
def unique_strings(memories: List[str]) -> List[str]:
out: List[str] = []
for m in memories:
if not isinstance(m, str):
continue
key = comparison_key(m)
if key and key not in seen:
seen.add(key)
out.append(m)
return out
def unique_search(results: List[Any]) -> List[Any]:
out: List[Any] = []
for r in results:
# v4 search.memories/hybrid uses `memory` or `chunk`.
memory = (
r if isinstance(r, str) else _field(r, "memory", "chunk", "content", default="")
)
if not isinstance(memory, str):
memory = ""
memory = memory.strip()
key = comparison_key(memory)
if key and key not in seen:
seen.add(key)
out.append(r)
return out
return {
"static": unique_strings(static),
"dynamic": unique_strings(dynamic),
"search_results": unique_search(search_results),
}
def format_memories_to_text(
memories: Dict[str, Union[List[str], List[Any]]],
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.
Search results include temporal context (e.g., '3d ago') from updatedAt.
"""
sections = []
static = memories["static"]
dynamic = memories["dynamic"]
search_results = memories["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")
lines = []
for item in search_results:
if isinstance(item, str):
lines.append(f"- {item}")
continue
memory = _field(item, "memory", "chunk", "content", default="")
updated_at = _field(item, "updatedAt", "updated_at", default="")
time_str = format_relative_time(updated_at) if updated_at else ""
if time_str:
lines.append(f"- [{time_str}] {memory}")
else:
lines.append(f"- {memory}")
sections.append("\n".join(lines))
if not sections:
return ""
return f"{system_prompt}\n" + "\n\n".join(sections)