fix(dreamer): preserve last_dream_document_count in completion write

CodeRabbit caught this: update_collection_internal_metadata uses a
top-level JSONB `||` merge, so passing {"dream": {"last_dream_at": ...}}
replaces the entire "dream" subkey and drops last_dream_document_count
that was written by enqueue_dream.

Symptom: after every completed dream, the baseline drops to 0. Next
check_and_schedule_dream reads documents_since_last_dream as
current_count - 0 = current_count, so any collection with >= 50
explicit observations can re-trigger immediately once the 8h guard
expires, even with no new raw material.

Fix: read-modify-write. Fetch current collection, merge last_dream_at
into the existing "dream" dict, write the merged dict back. Preserves
sibling keys (current: last_dream_document_count; future-proof for
telemetry fields that might land in PR 4).

Regression test added to tests/dreamer/test_dreamer_integration.py:
pre-seeds {"dream": {"last_dream_document_count": 42}}, runs
process_dream, asserts both last_dream_at is written AND
last_dream_document_count == 42 is preserved.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
lilyplasticlabs 2026-04-17 17:43:28 -04:00
parent b89997cb67
commit c8fe40a329
2 changed files with 60 additions and 1 deletions

View File

@ -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 # Write last_dream_at at completion (not enqueue) so duplicate
# enqueues can't reset the 8h guard. Lenient: any non-null result. # 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() now_iso = datetime.now(timezone.utc).isoformat()
async with tracked_db("dream.last_dream_at_write") as db: 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( await crud.update_collection_internal_metadata(
db, db,
workspace_name, workspace_name,
payload.observer, payload.observer,
payload.observed, payload.observed,
update_data={"dream": {"last_dream_at": now_iso}}, update_data={"dream": dream_meta},
) )
except Exception as e: except Exception as e:

View File

@ -125,6 +125,53 @@ class TestLastDreamAtCompletionWrite:
"(failed dream). The guard should not falsely advance." "(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 @pytest.mark.asyncio
async def test_exception_path_leaves_last_dream_at_null( async def test_exception_path_leaves_last_dream_at_null(
self, self,