项目文件夹

文件
wehub-resource-sync c889a57b6b
Test Suites / Build CI Environment (push) Has been cancelled
Test Suites / Basic Tests (push) Has been cancelled
Test Suites / End-to-End Tests (push) Has been cancelled
Test Suites / CLI Tests (push) Has been cancelled
Test Suites / Slow End-to-End Tests (push) Has been cancelled
Test Suites / Graph Database Tests (push) Has been cancelled
Test Suites / Vector DB Tests (push) Has been cancelled
Test Suites / Temporal Graph Test (push) Has been cancelled
Test Suites / Search Test on Different DBs (push) Has been cancelled
Test Suites / Example Tests (push) Has been cancelled
Test Suites / Notebook Tests (push) Has been cancelled
Test Suites / OS and Python Tests Ubuntu (push) Has been cancelled
Test Suites / OS and Python Tests Extended (push) Has been cancelled
Test Suites / LLM Test Suite (push) Has been cancelled
Test Suites / S3 File Storage Test (push) Has been cancelled
Test Suites / Run Integration Tests (push) Has been cancelled
Test Suites / MCP Tests (push) Has been cancelled
Test Suites / Docker Compose Test (push) Has been cancelled
Test Suites / Docker CI test (push) Has been cancelled
Test Suites / Relational DB Migration Tests (push) Has been cancelled
Test Suites / Distributed Cognee Test (push) Has been cancelled
Test Suites / DB Examples Tests (push) Has been cancelled
Test Suites / Test Completion Status (push) Has been cancelled
Test Suites / Claude Code Review (push) Has been cancelled
Test Suites / basic checks (push) Has been cancelled
build | Build and Push Cognee MCP Docker Image to dockerhub / docker-build-and-push (push) Has been cancelled
Scorecard supply-chain security / Scorecard analysis (push) Has been cancelled
build | Build and Push Docker Image to dockerhub / docker-build-and-push (push) Has been cancelled
Weighted Edges Tests / Test Weighted Edges Core Functionality (3.11) (push) Has been cancelled
Weighted Edges Tests / Test Weighted Edges Core Functionality (3.12) (push) Has been cancelled
Weighted Edges Tests / Test Weighted Edges with Different Graph Databases (kuzu, kuzu) (push) Has been cancelled
Weighted Edges Tests / Test Weighted Edges with Different Graph Databases (neo4j, neo4j) (push) Has been cancelled
Weighted Edges Tests / Test Weighted Edges Examples (push) Has been cancelled
Weighted Edges Tests / Code Quality for Weighted Edges (push) Has been cancelled
chore: import upstream snapshot with attribution
2026-07-13 13:02:24 +08:00

74 行
2.9 KiB
Python

"""Mem0 memory source.
Reads a Mem0 export and yields COGX memory records. Accepts the shapes
produced by the Mem0 platform export API and by the OSS ``get_all()`` call:
- a plain JSON list of memory objects
- ``{"results": [...]}`` / ``{"memories": [...]}`` wrappers
- already-parsed Python lists/dicts (for live-API integration: fetch with the
``mem0ai`` client yourself and pass the response in)
Each memory becomes a :class:`COGXMemory` with scope taken from
``user_id``/``agent_id``/``run_id`` and timestamps preserved.
"""
import json
from pathlib import Path
from typing import Any, AsyncIterator, Dict, List, Union
from cognee.modules.migration.cogx import COGXMemory, COGXRecord, COGXScope, parse_timestamp
from cognee.modules.migration.sources.base import MemorySource
_CONTENT_KEYS = ("memory", "text", "data", "content")
class Mem0Source(MemorySource):
source_system = "mem0"
def __init__(self, data: Union[str, Path, List[Any], Dict[str, Any]], mode: str = "re-derive"):
super().__init__(mode=mode)
self._data = data
def _load_raw(self) -> List[Dict[str, Any]]:
data = self._data
if isinstance(data, (str, Path)):
data = json.loads(Path(data).read_text(encoding="utf-8"))
if isinstance(data, dict):
for key in ("results", "memories", "items"):
if isinstance(data.get(key), list):
data = data[key]
break
else:
raise ValueError(
"Unrecognized Mem0 export shape: expected a list or a dict "
"with a 'results'/'memories' key."
)
if not isinstance(data, list):
raise ValueError("Unrecognized Mem0 export shape: expected a list of memories.")
return [item for item in data if isinstance(item, dict)]
async def records(self) -> AsyncIterator[COGXRecord]:
for index, item in enumerate(self._load_raw()):
content = next(
(item[key] for key in _CONTENT_KEYS if isinstance(item.get(key), str)), None
)
if not content:
continue
categories = item.get("categories") or []
if isinstance(categories, str):
categories = [categories]
yield COGXMemory(
external_system=self.source_system,
external_id=str(item.get("id") or f"mem0-{index}"),
content=content,
categories=[str(category) for category in categories],
scope=COGXScope(
user_id=item.get("user_id"),
agent_id=item.get("agent_id"),
run_id=item.get("run_id"),
),
created_at=parse_timestamp(item.get("created_at")),
updated_at=parse_timestamp(item.get("updated_at")),
metadata={"mem0_metadata": item.get("metadata")} if item.get("metadata") else {},
)