From af927eaee96ee3a79ca614510aa5008cde86bec6 Mon Sep 17 00:00:00 2001 From: zhangchi47 Date: Fri, 14 Aug 2026 12:10:15 +0800 Subject: [PATCH] fix(consolidation): preserve temporal fields across dedup CREATE dedup folds source ids and merged text directly into the existing observation, bypassing the normal CREATE writer. Carry all four source-derived temporal fields into that fold and merge them null-safely. UPDATE dedup folds the rewritten observation into the survivor before deleting the redundant row. Merge temporal fields from both rows so the survivor keeps the complete source interval. The ordinary UPDATE path now passes event_date and applies the same rules: event_date and occurred_start use the earliest non-null value; occurred_end and mentioned_at use the latest non-null value. Apply the same merge in external memory stores. Document why event_date uses the store snapshot in that branch: MemoryFact does not carry event_date. Add regression coverage for PostgreSQL CREATE/UPDATE fold SQL, ordinary UPDATE parameters, external-store folds, and null/min/max semantics. The external-store ordinary UPDATE test covers the newly added event_date merge branch. Validation: - 45 consolidation dedup tests - 15 UPDATE/integrity tests - 7 source aggregation tests - Existing real-LLM temporal consolidation integration test --- .../engine/consolidation/consolidator.py | 156 ++++++++++++++++-- .../tests/test_consolidation_dedup.py | 110 ++++++++++++ .../test_integrity_violation_not_retried.py | 109 ++++++++++++ 3 files changed, 361 insertions(+), 14 deletions(-) diff --git a/hindsight-api-slim/hindsight_api/engine/consolidation/consolidator.py b/hindsight-api-slim/hindsight_api/engine/consolidation/consolidator.py index d72091c3df..8e6c3c0bae 100644 --- a/hindsight-api-slim/hindsight_api/engine/consolidation/consolidator.py +++ b/hindsight-api-slim/hindsight_api/engine/consolidation/consolidator.py @@ -237,6 +237,16 @@ class _DedupOutcome: best_text: str = "" +@dataclass(frozen=True) +class _TemporalFields: + """The four temporal columns that must survive an observation merge.""" + + event_date: datetime | None + occurred_start: datetime | None + occurred_end: datetime | None + mentioned_at: datetime | None + + async def _dedup_adjudicate( pool: DatabaseBackend, memory_engine: "MemoryEngine", @@ -322,6 +332,7 @@ async def _dedup_reconcile_create( create_text: str, create_source_ids: list[uuid.UUID], tags: list[str] | None, + source_temporal_fields: _TemporalFields, txn=None, ) -> str | None: """Semantic dedup for a single CREATE (create-time, focused 1-by-1). @@ -362,6 +373,26 @@ async def _dedup_reconcile_create( SET text = $1, source_memory_ids = (SELECT array_agg(DISTINCT e) FROM unnest(source_memory_ids || $2::uuid[]) e), proof_count = (SELECT count(DISTINCT e) FROM unnest(source_memory_ids || $2::uuid[]) e), + event_date = CASE + WHEN $5 IS NULL THEN event_date + WHEN event_date IS NULL THEN $5 + ELSE LEAST(event_date, $5) + END, + occurred_start = CASE + WHEN $6 IS NULL THEN occurred_start + WHEN occurred_start IS NULL THEN $6 + ELSE LEAST(occurred_start, $6) + END, + occurred_end = CASE + WHEN $7 IS NULL THEN occurred_end + WHEN occurred_end IS NULL THEN $7 + ELSE GREATEST(occurred_end, $7) + END, + mentioned_at = CASE + WHEN $8 IS NULL THEN mentioned_at + WHEN mentioned_at IS NULL THEN $8 + ELSE GREATEST(mentioned_at, $8) + END, updated_at = now(){search_vector_clause} WHERE id = $3::uuid AND text = $4 RETURNING id @@ -370,6 +401,10 @@ async def _dedup_reconcile_create( live_source_ids, uuid.UUID(outcome.best_id), outcome.best_text, + source_temporal_fields.event_date, + source_temporal_fields.occurred_start, + source_temporal_fields.occurred_end, + source_temporal_fields.mentioned_at, ) if folded is None: # The twin vanished (or was rewritten) during the connection-free LLM window. @@ -382,7 +417,15 @@ async def _dedup_reconcile_create( return None else: await _reconcile_merge_via_store( - store, conn, memory_engine, bank_id, outcome.best_id, outcome.merged_text, live_source_ids, txn=txn + store, + conn, + memory_engine, + bank_id, + outcome.best_id, + outcome.merged_text, + live_source_ids, + source_temporal_fields, + txn=txn, ) return outcome.best_id @@ -426,7 +469,8 @@ async def _dedup_reconcile_update( # Fold the updated observation's live sources into the twin (keeping the twin's embedding, as # in the create path) then delete the now-redundant updated row. The all_strict/any tag match # guarantees twin and updated share scope, so dropping the updated row's tags loses no - # visibility. Temporal fields follow the surviving twin (minimal scope; matches create). + # visibility. Temporal fields are merged with the surviving twin using the same source-field + # rules as ordinary observation updates, so folding cannot discard dates from either row. # The fold + delete share one short transaction so the twin gains the sources exactly as the # redundant row is removed; the slow adjudication above already ran connection-free. store = get_memories() @@ -469,6 +513,26 @@ async def _dedup_reconcile_update( proof_count = ( SELECT count(DISTINCT e) FROM unnest(t.source_memory_ids || $6::uuid[]) e ), + event_date = CASE + WHEN u.event_date IS NULL THEN t.event_date + WHEN t.event_date IS NULL THEN u.event_date + ELSE LEAST(t.event_date, u.event_date) + END, + occurred_start = CASE + WHEN u.occurred_start IS NULL THEN t.occurred_start + WHEN t.occurred_start IS NULL THEN u.occurred_start + ELSE LEAST(t.occurred_start, u.occurred_start) + END, + occurred_end = CASE + WHEN u.occurred_end IS NULL THEN t.occurred_end + WHEN t.occurred_end IS NULL THEN u.occurred_end + ELSE GREATEST(t.occurred_end, u.occurred_end) + END, + mentioned_at = CASE + WHEN u.mentioned_at IS NULL THEN t.mentioned_at + WHEN t.mentioned_at IS NULL THEN u.mentioned_at + ELSE GREATEST(t.mentioned_at, u.mentioned_at) + END, updated_at = now(){search_vector_clause} FROM {fq_table("memory_units")} u WHERE t.id = $2::uuid AND u.id = $3::uuid AND t.text = $4 AND u.text = $5 @@ -493,8 +557,22 @@ async def _dedup_reconcile_update( live_u_sources = await _filter_live_source_memories(conn, bank_id, updated_sources) if not live_u_sources: return + updated_temporal_fields = _TemporalFields( + event_date=updated_obs[0].event_date, + occurred_start=updated_obs[0].occurred_start, + occurred_end=updated_obs[0].occurred_end, + mentioned_at=updated_obs[0].mentioned_at, + ) await _reconcile_merge_via_store( - store, conn, memory_engine, bank_id, outcome.best_id, outcome.merged_text, live_u_sources, txn=txn + store, + conn, + memory_engine, + bank_id, + outcome.best_id, + outcome.merged_text, + live_u_sources, + updated_temporal_fields, + txn=txn, ) await _execute_delete_action(conn, bank_id, updated_id, txn=txn) logger.info( @@ -993,6 +1071,16 @@ def _merge_max(a: "datetime | str | None", b: "datetime | str | None") -> "datet return a if b is None else b if a is None else max(a, b) +def _merge_temporal_fields(left: _TemporalFields, right: _TemporalFields) -> _TemporalFields: + """Merge two observation temporal snapshots using the source aggregation rules.""" + return _TemporalFields( + event_date=_merge_min(left.event_date, right.event_date), + occurred_start=_merge_min(left.occurred_start, right.occurred_start), + occurred_end=_merge_max(left.occurred_end, right.occurred_end), + mentioned_at=_merge_max(left.mentioned_at, right.mentioned_at), + ) + + async def _reconcile_merge_via_store( store, conn, @@ -1001,6 +1089,7 @@ async def _reconcile_merge_via_store( observation_id: str, merged_text: str, add_source_ids: list, + add_temporal_fields: _TemporalFields, txn=None, ) -> None: """Dedup merge for a store that owns its rows: fold the extra source facts and the merged text @@ -1012,6 +1101,13 @@ async def _reconcile_merge_via_store( if cur is None: return merged_sources = list(dict.fromkeys([*(cur.source_memory_ids or []), *(str(s) for s in add_source_ids)])) + current_temporal_fields = _TemporalFields( + event_date=cur.event_date, + occurred_start=cur.occurred_start, + occurred_end=cur.occurred_end, + mentioned_at=cur.mentioned_at, + ) + merged_temporal_fields = _merge_temporal_fields(current_temporal_fields, add_temporal_fields) embeddings = await embedding_utils.generate_embeddings_batch(memory_engine.embeddings, [merged_text]) await store.upsert_observation( conn=conn, @@ -1025,10 +1121,10 @@ async def _reconcile_merge_via_store( tags=list(cur.tags or []), proof_count=len(merged_sources), source_memory_ids=merged_sources, - event_date=cur.event_date, - occurred_start=cur.occurred_start, - occurred_end=cur.occurred_end, - mentioned_at=cur.mentioned_at, + event_date=merged_temporal_fields.event_date, + occurred_start=merged_temporal_fields.occurred_start, + occurred_end=merged_temporal_fields.occurred_end, + mentioned_at=merged_temporal_fields.mentioned_at, created_at=cur.created_at, ), ) @@ -2134,6 +2230,7 @@ async def _process_memory_batch( new_text=update.text, observations=union_observations, source_fact_tags=agg.tags, + source_event_date=agg.event_date, source_occurred_start=agg.occurred_start, source_occurred_end=agg.occurred_end, source_mentioned_at=agg.mentioned_at, @@ -2195,6 +2292,8 @@ async def _process_memory_batch( # Semantic near-duplicate reconciliation: merge this CREATE into an existing # near-identical observation (LLM-adjudicated, 1-by-1) instead of inserting a dup. if dedup_enabled: + # This fold bypasses the ordinary CREATE writer, so carry the source-derived + # temporal fields explicitly or the new evidence would be lost. merged_into = await _dedup_reconcile_create( pool, memory_engine, @@ -2204,6 +2303,12 @@ async def _process_memory_batch( create.text, create_source_ids, agg.tags, + _TemporalFields( + event_date=agg.event_date, + occurred_start=agg.occurred_start, + occurred_end=agg.occurred_end, + mentioned_at=agg.mentioned_at, + ), txn=txn, ) if merged_into is not None: @@ -2338,6 +2443,7 @@ async def _execute_update_action( new_text: str, observations: list["MemoryFact"], source_fact_tags: list[str] | None = None, + source_event_date: datetime | None = None, source_occurred_start: datetime | None = None, source_occurred_end: datetime | None = None, source_mentioned_at: datetime | None = None, @@ -2348,7 +2454,8 @@ async def _execute_update_action( Update an existing observation. Extends source_memory_ids with all contributing memories, updates temporal fields - (LEAST for occurred_start, GREATEST for occurred_end / mentioned_at), and merges tags. + (LEAST for event_date / occurred_start, GREATEST for occurred_end / mentioned_at), and + merges tags. The embedding is computed off-connection (a slow embedder must never pin a pooled connection); the liveness check + UPDATE + history + observation_sources sync then run @@ -2425,11 +2532,28 @@ async def _execute_update_action( embedding = $2::vector, source_memory_ids = $3, proof_count = $4, - tags = $9, + tags = $10, updated_at = now(), - occurred_start = LEAST(occurred_start, COALESCE($6, occurred_start)), - occurred_end = GREATEST(occurred_end, COALESCE($7, occurred_end)), - mentioned_at = GREATEST(mentioned_at, COALESCE($8, mentioned_at)){search_vector_clause} + event_date = CASE + WHEN $6 IS NULL THEN event_date + WHEN event_date IS NULL THEN $6 + ELSE LEAST(event_date, $6) + END, + occurred_start = CASE + WHEN $7 IS NULL THEN occurred_start + WHEN occurred_start IS NULL THEN $7 + ELSE LEAST(occurred_start, $7) + END, + occurred_end = CASE + WHEN $8 IS NULL THEN occurred_end + WHEN occurred_end IS NULL THEN $8 + ELSE GREATEST(occurred_end, $8) + END, + mentioned_at = CASE + WHEN $9 IS NULL THEN mentioned_at + WHEN mentioned_at IS NULL THEN $9 + ELSE GREATEST(mentioned_at, $9) + END{search_vector_clause} WHERE id = $5 """, new_text, @@ -2437,6 +2561,7 @@ async def _execute_update_action( source_ids, len(source_ids), uuid.UUID(observation_id), + source_event_date, source_occurred_start, source_occurred_end, source_mentioned_at, @@ -2458,7 +2583,10 @@ async def _execute_update_action( else: # Upsert overwrites the whole observation, so start from its current state (fetched # from the store) and apply the same merge the SQL does — LEAST/GREATEST on the times - # — while preserving fields the update never touches (event_date, created_at). + # — while preserving fields the update does not touch (created_at). + # MemoryFact does not expose event_date, so merge that field from the store snapshot + # and the source aggregation. The other temporal fields remain sourced from the + # recall model, where MemoryFact does expose them. current = await store.get_memories( conn=conn, fq_table=fq_table, bank_id=bank_id, unit_ids=[observation_id] ) @@ -2475,7 +2603,7 @@ async def _execute_update_action( tags=merged_tags, proof_count=len(source_ids), source_memory_ids=[str(s) for s in source_ids], - event_date=cur.event_date if cur else None, + event_date=_merge_min(cur.event_date if cur else None, source_event_date), occurred_start=_merge_min(model.occurred_start, source_occurred_start), occurred_end=_merge_max(model.occurred_end, source_occurred_end), mentioned_at=_merge_max(model.mentioned_at, source_mentioned_at), diff --git a/hindsight-api-slim/tests/test_consolidation_dedup.py b/hindsight-api-slim/tests/test_consolidation_dedup.py index 2a240ca802..8e443c55c7 100644 --- a/hindsight-api-slim/tests/test_consolidation_dedup.py +++ b/hindsight-api-slim/tests/test_consolidation_dedup.py @@ -10,6 +10,7 @@ import uuid from contextlib import asynccontextmanager from dataclasses import dataclass +from datetime import datetime, timezone from unittest.mock import DEFAULT, AsyncMock, patch import pytest @@ -22,7 +23,10 @@ _dedup_reconcile_update, _DedupDecision, _duplicate_create_target, + _merge_temporal_fields, _norm_obs_text, + _reconcile_merge_via_store, + _TemporalFields, ) from hindsight_api.engine.memories import RecallArms from hindsight_api.engine.search.types import RetrievalResult @@ -195,6 +199,12 @@ def _ctx(threshold: float = 0.97): create_text="YouTube content in Uzbek is very rich.", create_source_ids=[uuid.uuid4()], tags=["t1"], + source_temporal_fields=_TemporalFields( + event_date=datetime(2024, 1, 2, tzinfo=timezone.utc), + occurred_start=datetime(2023, 1, 2, tzinfo=timezone.utc), + occurred_end=datetime(2024, 1, 3, tzinfo=timezone.utc), + mentioned_at=datetime(2024, 1, 4, tzinfo=timezone.utc), + ), ) return kwargs, conn, llm @@ -332,9 +342,25 @@ async def test_dedup_llm_merge_folds_into_twin() -> None: assert result == _TWIN_ID # merged into the twin; caller skips the CREATE conn.fetchval.assert_awaited_once() # fold is a RETURNING-gated UPDATE args = conn.fetchval.await_args.args + fold_sql = args[0] assert args[1] == "Uzbek content on YouTube is very rich." # merged text persisted assert args[2] == kwargs["create_source_ids"] # new (live) source facts folded in assert args[3] == uuid.UUID(_TWIN_ID) # onto the twin row + assert "event_date = CASE" in fold_sql + assert "occurred_start = CASE" in fold_sql + assert "occurred_end = CASE" in fold_sql + assert "mentioned_at = CASE" in fold_sql + assert "ELSE LEAST(event_date, $5)" in fold_sql + assert "ELSE LEAST(occurred_start, $6)" in fold_sql + assert "ELSE GREATEST(occurred_end, $7)" in fold_sql + assert "ELSE GREATEST(mentioned_at, $8)" in fold_sql + assert args[4] == "Uzbek content on YouTube is described as very rich." + assert args[5:] == ( + kwargs["source_temporal_fields"].event_date, + kwargs["source_temporal_fields"].occurred_start, + kwargs["source_temporal_fields"].occurred_end, + kwargs["source_temporal_fields"].mentioned_at, + ) async def test_dedup_llm_merge_sanitizes_text_before_write() -> None: @@ -411,10 +437,19 @@ async def test_dedup_update_merge_folds_into_twin_and_deletes_updated() -> None: # LIVE sources (snapshotted via fetchrow, filtered FOR SHARE). conn.fetchval.assert_awaited_once() fold_args = conn.fetchval.await_args.args + fold_sql = fold_args[0] assert fold_args[1] == "Uzbek YouTube content is very rich and growing." # merged text on the twin assert fold_args[2] == uuid.UUID(_TWIN_ID) # survivor = the twin assert fold_args[3] == uuid.UUID(_UPDATED_ID) # folded-from = the updated row assert fold_args[6] == conn.fetchrow_result["source_memory_ids"] # only live updated-row sources + assert "event_date = CASE" in fold_sql + assert "occurred_start = CASE" in fold_sql + assert "occurred_end = CASE" in fold_sql + assert "mentioned_at = CASE" in fold_sql + assert "ELSE LEAST(t.event_date, u.event_date)" in fold_sql + assert "ELSE LEAST(t.occurred_start, u.occurred_start)" in fold_sql + assert "ELSE GREATEST(t.occurred_end, u.occurred_end)" in fold_sql + assert "ELSE GREATEST(t.mentioned_at, u.mentioned_at)" in fold_sql # Then the updated row is deleted: DELETE of the row, and DELETE of its observation_history # (no longer cascaded from memory_units — that FK was dropped). assert conn.execute.await_count == 2 @@ -580,6 +615,81 @@ async def test_dedup_update_all_updated_sources_deleted_skips_fold_and_delete() conn.execute.assert_not_called() +def test_merge_temporal_fields_uses_min_max_and_preserves_non_null_values() -> None: + early = datetime(2023, 1, 1, tzinfo=timezone.utc) + late = datetime(2024, 1, 1, tzinfo=timezone.utc) + survivor = _TemporalFields( + event_date=late, + occurred_start=late, + occurred_end=early, + mentioned_at=early, + ) + incoming = _TemporalFields( + event_date=early, + occurred_start=None, + occurred_end=late, + mentioned_at=None, + ) + + merged = _merge_temporal_fields(survivor, incoming) + + assert merged == _TemporalFields( + event_date=early, + occurred_start=late, + occurred_end=late, + mentioned_at=early, + ) + + +async def test_store_dedup_fold_merges_temporal_fields() -> None: + early = datetime(2023, 1, 1, tzinfo=timezone.utc) + late = datetime(2024, 1, 1, tzinfo=timezone.utc) + store = types.SimpleNamespace( + get_memories=AsyncMock( + return_value=[ + types.SimpleNamespace( + source_memory_ids=["existing-source"], + event_date=late, + occurred_start=late, + occurred_end=early, + mentioned_at=early, + tags=["t1"], + created_at=late, + ) + ] + ), + upsert_observation=AsyncMock(), + ) + memory_engine = types.SimpleNamespace(embeddings=object()) + incoming = _TemporalFields( + event_date=early, + occurred_start=None, + occurred_end=late, + mentioned_at=None, + ) + + with patch( + "hindsight_api.engine.consolidation.consolidator.embedding_utils.generate_embeddings_batch", + new=AsyncMock(return_value=[[0.1, 0.2, 0.3]]), + ): + await _reconcile_merge_via_store( + store, + conn=object(), + memory_engine=memory_engine, + bank_id="bank1", + observation_id=_TWIN_ID, + merged_text="merged text", + add_source_ids=[uuid.UUID("55555555-5555-4555-8555-555555555555")], + add_temporal_fields=incoming, + ) + + record = store.upsert_observation.await_args.kwargs["record"] + assert record.event_date == early + assert record.occurred_start == late + assert record.occurred_end == late + assert record.mentioned_at == early + + # ── _process_memory_batch create-contract (created vs skipped) ──────────────── diff --git a/hindsight-api-slim/tests/test_integrity_violation_not_retried.py b/hindsight-api-slim/tests/test_integrity_violation_not_retried.py index fc0cf1cdd4..6651f45e79 100644 --- a/hindsight-api-slim/tests/test_integrity_violation_not_retried.py +++ b/hindsight-api-slim/tests/test_integrity_violation_not_retried.py @@ -15,6 +15,7 @@ import json import uuid from contextlib import ExitStack +from datetime import datetime, timezone from types import SimpleNamespace from unittest.mock import AsyncMock, MagicMock, patch @@ -310,6 +311,114 @@ async def test_update_action_writes_history_when_row_present(): append_mock.assert_called_once() +@pytest.mark.asyncio +async def test_update_action_merges_event_date_with_source_temporal_fields() -> None: + """The PostgreSQL UPDATE path must carry event_date through the same min merge as CREATE.""" + from hindsight_api.engine.consolidation import consolidator + + observation_id = str(uuid.uuid4()) + source_ids = [uuid.uuid4()] + source_event_date = datetime(2023, 1, 1, tzinfo=timezone.utc) + source_occurred_start = datetime(2023, 2, 1, tzinfo=timezone.utc) + source_occurred_end = datetime(2024, 2, 1, tzinfo=timezone.utc) + source_mentioned_at = datetime(2024, 3, 1, tzinfo=timezone.utc) + + conn = AsyncMock() + conn.execute_rows_affected = AsyncMock(return_value=1) + conn.transaction = MagicMock(return_value=_AsyncNullCtx(None)) + memory_engine = MagicMock() + memory_engine._backend.ops.uses_observation_sources_table = False + append_mock = AsyncMock() + + with _patch_update_action_deps(consolidator, conn, source_ids, append_mock): + result = await consolidator._execute_update_action( + pool=MagicMock(), + memory_engine=memory_engine, + bank_id="bank-x", + source_memory_ids=source_ids, + observation_id=observation_id, + new_text="new observation text", + observations=[_observation_fact(observation_id)], + source_fact_tags=["scope_b"], + source_event_date=source_event_date, + source_occurred_start=source_occurred_start, + source_occurred_end=source_occurred_end, + source_mentioned_at=source_mentioned_at, + ) + + assert result is not None + update_args = conn.execute_rows_affected.await_args.args + update_sql = update_args[0] + assert "event_date = CASE" in update_sql + assert "ELSE LEAST(event_date, $6)" in update_sql + assert update_args[6:10] == ( + source_event_date, + source_occurred_start, + source_occurred_end, + source_mentioned_at, + ) + + +@pytest.mark.asyncio +async def test_update_action_store_branch_merges_event_date() -> None: + """The external-store UPDATE path must also min-merge event_date.""" + from hindsight_api.engine.consolidation import consolidator + + observation_id = str(uuid.uuid4()) + source_ids = [uuid.uuid4()] + early = datetime(2023, 1, 1, tzinfo=timezone.utc) + late = datetime(2024, 1, 1, tzinfo=timezone.utc) + current = SimpleNamespace( + event_date=late, + created_at=late, + tags=["scope_a"], + source_memory_ids=[str(uuid.uuid4())], + ) + store = SimpleNamespace( + writes_memory_rows_in_sql_for=lambda bank_id: False, + get_memories=AsyncMock(return_value=[current]), + upsert_observation=AsyncMock(), + ) + conn = AsyncMock() + conn.transaction = MagicMock(return_value=_AsyncNullCtx(None)) + memory_engine = MagicMock() + memory_engine._backend.ops.uses_observation_sources_table = False + + with ExitStack() as stack: + stack.enter_context(patch("hindsight_api.config.get_config", _fake_config)) + stack.enter_context( + patch.object(consolidator, "acquire_with_retry", MagicMock(return_value=_AsyncNullCtx(conn))) + ) + stack.enter_context(patch.object(consolidator, "get_memories", MagicMock(return_value=store))) + stack.enter_context(patch.object(consolidator, "_any_live_source_memory", AsyncMock(return_value=True))) + stack.enter_context( + patch.object(consolidator, "_filter_live_source_memories", AsyncMock(return_value=source_ids)) + ) + stack.enter_context( + patch.object( + consolidator.embedding_utils, + "generate_embeddings_batch", + AsyncMock(return_value=[[0.1, 0.2, 0.3]]), + ) + ) + stack.enter_context(patch.object(consolidator, "_append_observation_history", AsyncMock())) + result = await consolidator._execute_update_action( + pool=MagicMock(), + memory_engine=memory_engine, + bank_id="bank-x", + source_memory_ids=source_ids, + observation_id=observation_id, + new_text="new observation text", + observations=[_observation_fact(observation_id)], + source_fact_tags=["scope_b"], + source_event_date=early, + ) + + assert result is not None + record = store.upsert_observation.await_args.kwargs["record"] + assert record.event_date == early + + @pytest.mark.parametrize( "status, expected", [