diff --git a/mem0/graphs/neptune/neptunedb.py b/mem0/graphs/neptune/neptunedb.py index 18a5e139e..e5ebf87a0 100644 --- a/mem0/graphs/neptune/neptunedb.py +++ b/mem0/graphs/neptune/neptunedb.py @@ -1,7 +1,6 @@ import logging import uuid -from datetime import datetime -import pytz +from datetime import datetime, timezone from .base import NeptuneBase @@ -114,7 +113,7 @@ class MemoryGraph(NeptuneBase): "name": destination, "type": destination_type, "user_id": user_id, - "created_at": datetime.now(pytz.timezone("US/Pacific")).isoformat(), + "created_at": datetime.now(timezone.utc).isoformat(), } self.vector_store.insert( vectors=[dest_embedding], @@ -189,7 +188,7 @@ class MemoryGraph(NeptuneBase): "name": source, "type": source_type, "user_id": user_id, - "created_at": datetime.now(pytz.timezone("US/Pacific")).isoformat(), + "created_at": datetime.now(timezone.utc).isoformat(), } self.vector_store.insert( vectors=[source_embedding], @@ -316,14 +315,14 @@ class MemoryGraph(NeptuneBase): "name": source, "type": source_type, "user_id": user_id, - "created_at": datetime.now(pytz.timezone("US/Pacific")).isoformat(), + "created_at": datetime.now(timezone.utc).isoformat(), } destination_id = str(uuid.uuid4()) destination_payload = { "name": destination, "type": destination_type, "user_id": user_id, - "created_at": datetime.now(pytz.timezone("US/Pacific")).isoformat(), + "created_at": datetime.now(timezone.utc).isoformat(), } self.vector_store.insert( vectors=[source_embedding, dest_embedding], diff --git a/mem0/memory/main.py b/mem0/memory/main.py index cd59ba643..e7ac65f3c 100644 --- a/mem0/memory/main.py +++ b/mem0/memory/main.py @@ -8,10 +8,9 @@ import os import uuid import warnings from copy import deepcopy -from datetime import datetime +from datetime import datetime, timezone from typing import Any, Dict, Optional -import pytz from pydantic import ValidationError from mem0.configs.base import MemoryConfig, MemoryItem @@ -51,6 +50,19 @@ warnings.filterwarnings("ignore", category=DeprecationWarning, message=".*swigva logger = logging.getLogger(__name__) +def _normalize_iso_timestamp_to_utc(timestamp: Optional[str]) -> Optional[str]: + """Normalize timezone-aware ISO timestamps to UTC without rewriting naive values.""" + if not timestamp: + return timestamp + try: + parsed = datetime.fromisoformat(timestamp) + except ValueError: + return timestamp + if parsed.tzinfo is None: + return timestamp + return parsed.astimezone(timezone.utc).isoformat() + + # Fields that hold runtime auth/connection objects and must be preserved. # These are non-serializable objects (e.g. AWSV4SignerAuth, RequestsHttpConnection) # needed by clients like OpenSearch — not sensitive strings to redact. @@ -637,7 +649,10 @@ class Memory(MemoryBase): updated_metadata["agent_id"] = metadata["agent_id"] if metadata.get("run_id"): updated_metadata["run_id"] = metadata["run_id"] - updated_metadata["updated_at"] = datetime.now(pytz.timezone("US/Pacific")).isoformat() + updated_metadata["created_at"] = _normalize_iso_timestamp_to_utc( + updated_metadata.get("created_at") + ) + updated_metadata["updated_at"] = datetime.now(timezone.utc).isoformat() self.vector_store.update( vector_id=memory_id, @@ -700,8 +715,8 @@ class Memory(MemoryBase): id=memory.id, memory=memory.payload.get("data", ""), hash=memory.payload.get("hash"), - created_at=memory.payload.get("created_at"), - updated_at=memory.payload.get("updated_at"), + created_at=_normalize_iso_timestamp_to_utc(memory.payload.get("created_at")), + updated_at=_normalize_iso_timestamp_to_utc(memory.payload.get("updated_at")), ).model_dump() for key in promoted_payload_keys: @@ -803,8 +818,8 @@ class Memory(MemoryBase): id=mem.id, memory=mem.payload.get("data", ""), hash=mem.payload.get("hash"), - created_at=mem.payload.get("created_at"), - updated_at=mem.payload.get("updated_at"), + created_at=_normalize_iso_timestamp_to_utc(mem.payload.get("created_at")), + updated_at=_normalize_iso_timestamp_to_utc(mem.payload.get("updated_at")), ).model_dump(exclude={"score"}) for key in promoted_payload_keys: @@ -1035,8 +1050,8 @@ class Memory(MemoryBase): id=mem.id, memory=mem.payload.get("data", ""), hash=mem.payload.get("hash"), - created_at=mem.payload.get("created_at"), - updated_at=mem.payload.get("updated_at"), + created_at=_normalize_iso_timestamp_to_utc(mem.payload.get("created_at")), + updated_at=_normalize_iso_timestamp_to_utc(mem.payload.get("updated_at")), score=mem.score, ).model_dump() @@ -1145,7 +1160,7 @@ class Memory(MemoryBase): metadata = metadata or {} metadata["data"] = data metadata["hash"] = hashlib.md5(data.encode()).hexdigest() - metadata["created_at"] = datetime.now(pytz.timezone("US/Pacific")).isoformat() + metadata["created_at"] = datetime.now(timezone.utc).isoformat() self.vector_store.insert( vectors=[embeddings], @@ -1217,8 +1232,8 @@ class Memory(MemoryBase): new_metadata["data"] = data new_metadata["hash"] = hashlib.md5(data.encode()).hexdigest() - new_metadata["created_at"] = existing_memory.payload.get("created_at") - new_metadata["updated_at"] = datetime.now(pytz.timezone("US/Pacific")).isoformat() + new_metadata["created_at"] = _normalize_iso_timestamp_to_utc(existing_memory.payload.get("created_at")) + new_metadata["updated_at"] = datetime.now(timezone.utc).isoformat() # Preserve session identifiers from existing memory only if not provided in new metadata if "user_id" not in new_metadata and "user_id" in existing_memory.payload: @@ -1661,7 +1676,10 @@ class AsyncMemory(MemoryBase): updated_metadata["agent_id"] = meta["agent_id"] if meta.get("run_id"): updated_metadata["run_id"] = meta["run_id"] - updated_metadata["updated_at"] = datetime.now(pytz.timezone("US/Pacific")).isoformat() + updated_metadata["created_at"] = _normalize_iso_timestamp_to_utc( + updated_metadata.get("created_at") + ) + updated_metadata["updated_at"] = datetime.now(timezone.utc).isoformat() await asyncio.to_thread( self.vector_store.update, @@ -1747,8 +1765,8 @@ class AsyncMemory(MemoryBase): id=memory.id, memory=memory.payload.get("data", ""), hash=memory.payload.get("hash"), - created_at=memory.payload.get("created_at"), - updated_at=memory.payload.get("updated_at"), + created_at=_normalize_iso_timestamp_to_utc(memory.payload.get("created_at")), + updated_at=_normalize_iso_timestamp_to_utc(memory.payload.get("updated_at")), ).model_dump() for key in promoted_payload_keys: @@ -1855,8 +1873,8 @@ class AsyncMemory(MemoryBase): id=mem.id, memory=mem.payload.get("data", ""), hash=mem.payload.get("hash"), - created_at=mem.payload.get("created_at"), - updated_at=mem.payload.get("updated_at"), + created_at=_normalize_iso_timestamp_to_utc(mem.payload.get("created_at")), + updated_at=_normalize_iso_timestamp_to_utc(mem.payload.get("updated_at")), ).model_dump(exclude={"score"}) for key in promoted_payload_keys: @@ -2096,8 +2114,8 @@ class AsyncMemory(MemoryBase): id=mem.id, memory=mem.payload.get("data", ""), hash=mem.payload.get("hash"), - created_at=mem.payload.get("created_at"), - updated_at=mem.payload.get("updated_at"), + created_at=_normalize_iso_timestamp_to_utc(mem.payload.get("created_at")), + updated_at=_normalize_iso_timestamp_to_utc(mem.payload.get("updated_at")), score=mem.score, ).model_dump() @@ -2211,7 +2229,7 @@ class AsyncMemory(MemoryBase): metadata = metadata or {} metadata["data"] = data metadata["hash"] = hashlib.md5(data.encode()).hexdigest() - metadata["created_at"] = datetime.now(pytz.timezone("US/Pacific")).isoformat() + metadata["created_at"] = datetime.now(timezone.utc).isoformat() await asyncio.to_thread( self.vector_store.insert, @@ -2301,8 +2319,8 @@ class AsyncMemory(MemoryBase): new_metadata["data"] = data new_metadata["hash"] = hashlib.md5(data.encode()).hexdigest() - new_metadata["created_at"] = existing_memory.payload.get("created_at") - new_metadata["updated_at"] = datetime.now(pytz.timezone("US/Pacific")).isoformat() + new_metadata["created_at"] = _normalize_iso_timestamp_to_utc(existing_memory.payload.get("created_at")) + new_metadata["updated_at"] = datetime.now(timezone.utc).isoformat() # Preserve session identifiers from existing memory only if not provided in new metadata if "user_id" not in new_metadata and "user_id" in existing_memory.payload: diff --git a/mem0/vector_stores/redis.py b/mem0/vector_stores/redis.py index 81617cbd1..fc2048e35 100644 --- a/mem0/vector_stores/redis.py +++ b/mem0/vector_stores/redis.py @@ -1,10 +1,9 @@ import json import logging -from datetime import datetime +from datetime import datetime, timezone from functools import reduce import numpy as np -import pytz import redis from redis.commands.search.query import Query from redisvl.index import SearchIndex @@ -164,12 +163,12 @@ class RedisDB(VectorStoreBase): "hash": result["hash"], "data": result["memory"], "created_at": datetime.fromtimestamp( - int(result["created_at"]), tz=pytz.timezone("US/Pacific") + int(result["created_at"]), tz=timezone.utc ).isoformat(timespec="microseconds"), **( { "updated_at": datetime.fromtimestamp( - int(result["updated_at"]), tz=pytz.timezone("US/Pacific") + int(result["updated_at"]), tz=timezone.utc ).isoformat(timespec="microseconds") } if "updated_at" in result @@ -207,13 +206,13 @@ class RedisDB(VectorStoreBase): payload = { "hash": result["hash"], "data": result["memory"], - "created_at": datetime.fromtimestamp(int(result["created_at"]), tz=pytz.timezone("US/Pacific")).isoformat( + "created_at": datetime.fromtimestamp(int(result["created_at"]), tz=timezone.utc).isoformat( timespec="microseconds" ), **( { "updated_at": datetime.fromtimestamp( - int(result["updated_at"]), tz=pytz.timezone("US/Pacific") + int(result["updated_at"]), tz=timezone.utc ).isoformat(timespec="microseconds") } if "updated_at" in result @@ -271,12 +270,12 @@ class RedisDB(VectorStoreBase): "hash": result["hash"], "data": result["memory"], "created_at": datetime.fromtimestamp( - int(result["created_at"]), tz=pytz.timezone("US/Pacific") + int(result["created_at"]), tz=timezone.utc ).isoformat(timespec="microseconds"), **( { "updated_at": datetime.fromtimestamp( - int(result["updated_at"]), tz=pytz.timezone("US/Pacific") + int(result["updated_at"]), tz=timezone.utc ).isoformat(timespec="microseconds") } if result.__dict__.get("updated_at") diff --git a/tests/memory/test_main.py b/tests/memory/test_main.py index 3c3c2b2e1..3b178f243 100644 --- a/tests/memory/test_main.py +++ b/tests/memory/test_main.py @@ -1,9 +1,10 @@ import logging +from datetime import datetime, timezone from unittest.mock import MagicMock import pytest -from mem0.memory.main import AsyncMemory, Memory +from mem0.memory.main import AsyncMemory, Memory, _normalize_iso_timestamp_to_utc def _setup_mocks(mocker): @@ -127,3 +128,78 @@ class TestAsyncAddToVectorStoreErrors: assert result == [] assert "Empty response from LLM, no memories to extract" in caplog.text assert mock_capture_event.call_count == 1 + + +def _build_memory_instance(mocker, memory_cls): + _setup_mocks(mocker) + mocker.patch("mem0.memory.main.SQLiteManager", mocker.MagicMock()) + mocker.patch("mem0.memory.main.MEM0_TELEMETRY", False) + memory = memory_cls() + memory.config = mocker.MagicMock() + memory.config.custom_fact_extraction_prompt = None + memory.config.custom_update_memory_prompt = None + memory.api_version = "v1.1" + memory.vector_store = mocker.MagicMock() + memory.db = mocker.MagicMock() + return memory + + +def _assert_utc_timestamp(timestamp: str): + parsed = datetime.fromisoformat(timestamp) + assert parsed.tzinfo == timezone.utc + assert parsed.utcoffset().total_seconds() == 0 + + +def test_create_memory_uses_utc_timestamps(mocker): + memory = _build_memory_instance(mocker, Memory) + memory._create_memory("new memory", {"new memory": [0.1, 0.2, 0.3]}, metadata={}) + payload = memory.vector_store.insert.call_args.kwargs["payloads"][0] + _assert_utc_timestamp(payload["created_at"]) + + +def test_update_memory_uses_utc_timestamps(mocker): + memory = _build_memory_instance(mocker, Memory) + memory.vector_store.get.return_value = MagicMock( + payload={"data": "old memory", "created_at": "2026-03-17T17:00:00-07:00"} + ) + memory._update_memory("memory-id", "new memory", {"new memory": [0.1, 0.2, 0.3]}, metadata={}) + payload = memory.vector_store.update.call_args.kwargs["payload"] + assert payload["created_at"] == "2026-03-18T00:00:00+00:00" + _assert_utc_timestamp(payload["updated_at"]) + + +@pytest.mark.asyncio +async def test_async_create_memory_uses_utc_timestamps(mocker): + memory = _build_memory_instance(mocker, AsyncMemory) + await memory._create_memory("new memory", {"new memory": [0.1, 0.2, 0.3]}, metadata={}) + payload = memory.vector_store.insert.call_args.kwargs["payloads"][0] + _assert_utc_timestamp(payload["created_at"]) + + +@pytest.mark.asyncio +async def test_async_update_memory_uses_utc_timestamps(mocker): + memory = _build_memory_instance(mocker, AsyncMemory) + memory.vector_store.get.return_value = MagicMock( + payload={"data": "old memory", "created_at": "2026-03-17T17:00:00-07:00"} + ) + await memory._update_memory("memory-id", "new memory", {"new memory": [0.1, 0.2, 0.3]}, metadata={}) + payload = memory.vector_store.update.call_args.kwargs["payload"] + assert payload["created_at"] == "2026-03-18T00:00:00+00:00" + _assert_utc_timestamp(payload["updated_at"]) + + +def test_normalize_iso_timestamp_to_utc_preserves_naive_values(): + assert _normalize_iso_timestamp_to_utc("2026-03-18T00:00:00") == "2026-03-18T00:00:00" + + +def test_normalize_iso_timestamp_to_utc_converts_pacific(): + result = _normalize_iso_timestamp_to_utc("2026-03-17T17:00:00-07:00") + assert result == "2026-03-18T00:00:00+00:00" + + +def test_normalize_iso_timestamp_to_utc_handles_none(): + assert _normalize_iso_timestamp_to_utc(None) is None + + +def test_normalize_iso_timestamp_to_utc_handles_empty(): + assert _normalize_iso_timestamp_to_utc("") == "" diff --git a/tests/memory/test_neptune_memory.py b/tests/memory/test_neptune_memory.py index be711164f..c491b99af 100644 --- a/tests/memory/test_neptune_memory.py +++ b/tests/memory/test_neptune_memory.py @@ -1,8 +1,11 @@ import unittest +from datetime import datetime, timezone from unittest.mock import MagicMock, patch + import pytest -from mem0.graphs.neptune.neptunedb import MemoryGraph + from mem0.graphs.neptune.base import NeptuneBase +from mem0.graphs.neptune.neptunedb import MemoryGraph class TestNeptuneMemory(unittest.TestCase): @@ -289,6 +292,25 @@ class TestNeptuneMemory(unittest.TestCase): # Check the result self.assertEqual(result, mock_query_result) + def test_add_new_entities_payloads_use_utc_timestamps(self): + """Test that Neptune vector-store payloads use UTC timestamps.""" + self.memory_graph._add_new_entities_cypher( + source="alice", + source_embedding=[0.1, 0.2], + source_type="person", + destination="bob", + dest_embedding=[0.3, 0.4], + destination_type="person", + relationship="KNOWS", + user_id=self.user_id, + ) + + _, kwargs = self.mock_vector_store.insert.call_args + for payload in kwargs["payloads"]: + parsed = datetime.fromisoformat(payload["created_at"]) + self.assertEqual(parsed.tzinfo, timezone.utc) + self.assertEqual(parsed.utcoffset().total_seconds(), 0) + def test_search_graph_db(self): """Test the _search_graph_db method.""" # Mock node list