diff --git a/src/dreamer/orchestrator.py b/src/dreamer/orchestrator.py index 14678ff6..3857706a 100644 --- a/src/dreamer/orchestrator.py +++ b/src/dreamer/orchestrator.py @@ -326,14 +326,26 @@ DREAM: {payload.dream_type} documents for {workspace_name}/{payload.observer}/{p # Write last_dream_at at completion (not enqueue) so duplicate # enqueues can't reset the 8h guard. Lenient: any non-null result. + # Read-modify-write: update_collection_internal_metadata uses a + # top-level JSONB || merge, so passing {"dream": {"last_dream_at": ...}} + # alone would replace the whole "dream" subkey and drop + # last_dream_document_count (written by enqueue_dream). now_iso = datetime.now(timezone.utc).isoformat() async with tracked_db("dream.last_dream_at_write") as db: + collection = await crud.get_collection( + db, + workspace_name, + observer=payload.observer, + observed=payload.observed, + ) + dream_meta = dict(collection.internal_metadata.get("dream", {})) + dream_meta["last_dream_at"] = now_iso await crud.update_collection_internal_metadata( db, workspace_name, payload.observer, payload.observed, - update_data={"dream": {"last_dream_at": now_iso}}, + update_data={"dream": dream_meta}, ) except Exception as e: diff --git a/tests/dreamer/test_dreamer_integration.py b/tests/dreamer/test_dreamer_integration.py index c87353f4..86083bdb 100644 --- a/tests/dreamer/test_dreamer_integration.py +++ b/tests/dreamer/test_dreamer_integration.py @@ -125,6 +125,53 @@ class TestLastDreamAtCompletionWrite: "(failed dream). The guard should not falsely advance." ) + @pytest.mark.asyncio + async def test_completion_preserves_last_dream_document_count( + self, + db_session: AsyncSession, + sample_data: tuple[models.Workspace, models.Peer], + ): + """Completion write must NOT drop sibling keys in the `dream` sub-object. + + `update_collection_internal_metadata` uses a top-level JSONB `||` merge, + which replaces the whole `"dream"` key. Without a read-modify-write, the + completion write (`last_dream_at`) would wipe `last_dream_document_count` + that enqueue set — causing the next `check_and_schedule_dream` to read 0 + as the baseline and let a fresh dream trigger after the 8h guard expires + even without any new explicit documents. + """ + workspace, peer = sample_data + # Pre-seed collection as if enqueue_dream already wrote the baseline. + collection = models.Collection( + observer=peer.name, + observed=peer.name, + workspace_name=workspace.name, + internal_metadata={"dream": {"last_dream_document_count": 42}}, + ) + db_session.add(collection) + await db_session.commit() + await db_session.refresh(collection) + + payload = DreamPayload( + dream_type=DreamType.OMNI, + observer=collection.observer, + observed=collection.observed, + ) + + with patch( + "src.dreamer.orchestrator.run_dream", + new=AsyncMock(return_value=_make_dream_result()), + ): + await process_dream(payload, collection.workspace_name) + + dream_meta = await _get_dream_metadata(db_session, collection) + assert "last_dream_at" in dream_meta, "last_dream_at must be written" + assert dream_meta.get("last_dream_document_count") == 42, ( + "last_dream_document_count from enqueue must be preserved across " + "the completion write — top-level JSONB || would drop it without " + "the read-modify-write in process_dream." + ) + @pytest.mark.asyncio async def test_exception_path_leaves_last_dream_at_null( self,