test(dreamer): threshold filter + last_dream_at relocation regression tests

Tests for Finding 2 and Finding 3 (code-level):

- TestThresholdFilter (tests/dreamer/test_dream_scheduler.py):
  * Mixed levels below explicit threshold: 30 explicit + 40 deductive
    + 10 inductive → no trigger (core regression, buggy count would trigger)
  * Explicit-only at threshold: 60 explicit → triggers
  * Contradiction excluded: 100 contradiction + 10 explicit → no trigger
    (confirms positive == "explicit" filter excludes all dreamer output)

- TestEnqueueDreamMetadataShape (tests/deriver/test_enqueue_dream.py):
  * AsyncMock-patched update_collection_internal_metadata verifies
    enqueue writes last_dream_document_count but NOT last_dream_at

- TestLastDreamAtCompletionWrite (tests/dreamer/test_dreamer_integration.py):
  * Happy path: run_dream returns DreamResult → last_dream_at written
  * Failure path: run_dream returns None → last_dream_at absent
  * Exception path: run_dream raises → last_dream_at absent,
    process_dream swallows exception (queue-processed semantics preserved)

Docstring on check_and_schedule_dream tightened: "document threshold"
-> "explicit-observation threshold" to reflect filter semantics.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
lilyplasticlabs 2026-04-17 15:30:45 -04:00
parent 1487cd269f
commit b89997cb67
4 changed files with 358 additions and 3 deletions

View File

@ -231,11 +231,11 @@ async def check_and_schedule_dream(
collection: models.Collection, collection: models.Collection,
) -> bool: ) -> bool:
""" """
Check if a collection has reached the document threshold and schedule a timer-based dream. Check if a collection has reached the explicit-observation threshold and schedule a timer-based dream.
This function only schedules a timer-based dream if: This function only schedules a timer-based dream if:
1. Dreams are enabled 1. Dreams are enabled
2. Document threshold is reached 2. Explicit-observation threshold is reached (dreamer output does not count)
3. Minimum hours between dreams have passed 3. Minimum hours between dreams have passed
4. No dream is already scheduled for this collection 4. No dream is already scheduled for this collection

View File

@ -0,0 +1,68 @@
"""Regression tests for `enqueue_dream` metadata write shape.
Finding 3 (code-level) moves the `last_dream_at` timestamp write from
enqueue time to dream-completion time (in `process_dream`). These tests
verify that `enqueue_dream` no longer writes `last_dream_at` and still
writes `last_dream_document_count`.
"""
from unittest.mock import AsyncMock, patch
import pytest
from src import schemas
from src.deriver.enqueue import enqueue_dream
class TestEnqueueDreamMetadataShape:
"""Regression tests for Finding 3 code-level: `last_dream_at` relocation."""
@pytest.mark.asyncio
async def test_update_data_omits_last_dream_at(self):
"""`enqueue_dream` must NOT write `last_dream_at`.
Moved to completion (in `process_dream`) so duplicate enqueues can't
reset the 8-hour guard and failed dreams don't falsely advance it.
"""
with (
patch(
"src.deriver.enqueue.crud.update_collection_internal_metadata",
new_callable=AsyncMock,
) as mock_update,
# Short-circuit the dedup / insert paths so we only exercise the
# metadata update call. db_session.scalar returns False for both
# the in-progress and pending checks; db_session.execute is a no-op.
patch(
"src.deriver.enqueue.tracked_db",
) as mock_db_ctx,
):
mock_session = AsyncMock()
mock_session.scalar = AsyncMock(return_value=False)
mock_session.execute = AsyncMock()
mock_db_ctx.return_value.__aenter__.return_value = mock_session
await enqueue_dream(
workspace_name="test_workspace",
observer="alice",
observed="bob",
dream_type=schemas.DreamType.OMNI,
document_count=42,
session_name=None,
)
assert (
mock_update.called
), "update_collection_internal_metadata must be called"
call_kwargs = mock_update.call_args.kwargs
update_data = call_kwargs["update_data"]
dream_metadata = update_data["dream"]
assert (
"last_dream_document_count" in dream_metadata
), "last_dream_document_count should still be written at enqueue"
assert dream_metadata["last_dream_document_count"] == 42
assert "last_dream_at" not in dream_metadata, (
"last_dream_at must NOT be written at enqueue — it now writes "
"at dream completion in process_dream (orchestrator.py)."
)

View File

@ -4,8 +4,14 @@ from typing import Any
from unittest.mock import AsyncMock, patch from unittest.mock import AsyncMock, patch
import pytest import pytest
from sqlalchemy.ext.asyncio import AsyncSession
from src.dreamer.dream_scheduler import DreamScheduler, set_dream_scheduler from src import models
from src.dreamer.dream_scheduler import (
DreamScheduler,
check_and_schedule_dream,
set_dream_scheduler,
)
from src.schemas import DreamType from src.schemas import DreamType
from src.utils.work_unit import construct_work_unit_key from src.utils.work_unit import construct_work_unit_key
@ -411,6 +417,127 @@ class TestDocumentCountAtExecutionTime:
assert captured_document_count == CURRENT_DOC_COUNT assert captured_document_count == CURRENT_DOC_COUNT
class TestThresholdFilter:
"""Regression tests for Finding 2: threshold must count only explicit-level docs.
Previously the threshold counted all documents in a collection, including
dreamer output (deductive/inductive/contradiction). This created a feedback
loop where each dream's output inflated the trigger for the next dream.
The fix filters the count to `level == "explicit"` only.
"""
async def _make_collection(
self,
db_session: AsyncSession,
sample_data: tuple[models.Workspace, models.Peer],
) -> models.Collection:
"""Helper: create a Collection in the test workspace with no dream metadata."""
workspace, peer = sample_data
collection = models.Collection(
observer=peer.name,
observed=peer.name,
workspace_name=workspace.name,
internal_metadata={},
)
db_session.add(collection)
await db_session.commit()
return collection
async def _insert_doc(
self,
db_session: AsyncSession,
collection: models.Collection,
level: str,
) -> None:
"""Helper: insert one Document at the given level."""
db_session.add(
models.Document(
content="test",
level=level,
workspace_name=collection.workspace_name,
observer=collection.observer,
observed=collection.observed,
)
)
@pytest.mark.asyncio
async def test_mixed_levels_below_explicit_threshold(
self,
dream_scheduler: DreamScheduler,
db_session: AsyncSession,
sample_data: tuple[models.Workspace, models.Peer],
):
"""30 explicit + 40 deductive + 10 inductive → should NOT trigger.
Total doc count = 80 (would trigger under the buggy unfiltered count),
but explicit count = 30 < threshold 50, so the correct behavior is to
NOT schedule a dream. This is the core regression: the fix must reject
this scenario.
"""
collection = await self._make_collection(db_session, sample_data)
for _ in range(30):
await self._insert_doc(db_session, collection, "explicit")
for _ in range(40):
await self._insert_doc(db_session, collection, "deductive")
for _ in range(10):
await self._insert_doc(db_session, collection, "inductive")
await db_session.commit()
with patch.object(dream_scheduler, "schedule_dream", new_callable=AsyncMock):
scheduled = await check_and_schedule_dream(db_session, collection)
assert scheduled is False, (
"Threshold should filter on explicit level only — dreamer output "
"(deductive/inductive) must not count toward the trigger."
)
@pytest.mark.asyncio
async def test_explicit_only_at_threshold(
self,
dream_scheduler: DreamScheduler,
db_session: AsyncSession,
sample_data: tuple[models.Workspace, models.Peer],
):
"""60 explicit + 0 derived → should trigger (60 ≥ threshold 50)."""
collection = await self._make_collection(db_session, sample_data)
for _ in range(60):
await self._insert_doc(db_session, collection, "explicit")
await db_session.commit()
with patch.object(
dream_scheduler, "schedule_dream", new_callable=AsyncMock
) as mock_schedule:
scheduled = await check_and_schedule_dream(db_session, collection)
assert scheduled is True
assert mock_schedule.called, "schedule_dream should fire when threshold met"
@pytest.mark.asyncio
async def test_contradiction_excluded_from_count(
self,
dream_scheduler: DreamScheduler,
db_session: AsyncSession,
sample_data: tuple[models.Workspace, models.Peer],
):
"""Contradiction-level docs are dreamer output — must not count.
100 contradictions + 10 explicit explicit=10 < threshold=50, no trigger.
Confirms the positive `== "explicit"` filter excludes contradiction by
construction (same as deductive/inductive).
"""
collection = await self._make_collection(db_session, sample_data)
for _ in range(100):
await self._insert_doc(db_session, collection, "contradiction")
for _ in range(10):
await self._insert_doc(db_session, collection, "explicit")
await db_session.commit()
with patch.object(dream_scheduler, "schedule_dream", new_callable=AsyncMock):
scheduled = await check_and_schedule_dream(db_session, collection)
assert scheduled is False
class TestEnqueueCancelsDreamsCorrectly: class TestEnqueueCancelsDreamsCorrectly:
"""Integration test verifying the full flow of message enqueue cancelling dreams.""" """Integration test verifying the full flow of message enqueue cancelling dreams."""

View File

@ -0,0 +1,160 @@
"""Integration tests for the dream completion write.
Finding 3 (code-level) relocates `last_dream_at` from enqueue time to
dream-completion time (in `process_dream`). These tests exercise the real
Postgres JSONB merge via `tracked_db` to verify the write lands in the
collection's internal_metadata on successful dreams — and critically,
does NOT land on failures or exceptions.
"""
from typing import Any
from unittest.mock import AsyncMock, patch
import pytest
import pytest_asyncio
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from src import models
from src.dreamer.orchestrator import DreamResult, process_dream
from src.schemas import DreamType
from src.utils.queue_payload import DreamPayload
@pytest_asyncio.fixture
async def seeded_collection(
db_session: AsyncSession,
sample_data: tuple[models.Workspace, models.Peer],
) -> models.Collection:
"""Create a Collection with an empty dream metadata dict."""
workspace, peer = sample_data
collection = models.Collection(
observer=peer.name,
observed=peer.name,
workspace_name=workspace.name,
internal_metadata={},
)
db_session.add(collection)
await db_session.commit()
await db_session.refresh(collection)
return collection
def _make_dream_result() -> DreamResult:
"""Build a minimal non-null DreamResult for happy-path tests."""
return DreamResult(
run_id="test_run_01",
specialists_run=["deduction", "induction"],
deduction_success=True,
induction_success=True,
surprisal_enabled=False,
surprisal_conclusion_count=0,
total_iterations=3,
total_duration_ms=1234.5,
input_tokens=100,
output_tokens=50,
)
async def _get_dream_metadata(
db_session: AsyncSession, collection: models.Collection
) -> dict[str, Any]:
"""Re-fetch collection and return its internal_metadata['dream'] dict (or {})."""
await db_session.refresh(collection)
stmt = select(models.Collection).where(models.Collection.id == collection.id)
refreshed = (await db_session.execute(stmt)).scalar_one()
dream_meta: dict[str, Any] = refreshed.internal_metadata.get("dream", {})
return dream_meta
class TestLastDreamAtCompletionWrite:
"""Regression tests for Finding 3: `last_dream_at` written at completion."""
@pytest.mark.asyncio
async def test_happy_path_writes_last_dream_at(
self,
db_session: AsyncSession,
seeded_collection: models.Collection,
):
"""Non-null DreamResult → `last_dream_at` is set in internal_metadata."""
payload = DreamPayload(
dream_type=DreamType.OMNI,
observer=seeded_collection.observer,
observed=seeded_collection.observed,
)
with patch(
"src.dreamer.orchestrator.run_dream",
new=AsyncMock(return_value=_make_dream_result()),
):
await process_dream(payload, seeded_collection.workspace_name)
dream_meta = await _get_dream_metadata(db_session, seeded_collection)
assert (
"last_dream_at" in dream_meta
), "process_dream must write last_dream_at when run_dream returns a result"
# ISO format sanity check: contains 'T' separator and either 'Z' or '+'
assert "T" in dream_meta["last_dream_at"]
@pytest.mark.asyncio
async def test_failure_path_leaves_last_dream_at_null(
self,
db_session: AsyncSession,
seeded_collection: models.Collection,
):
"""run_dream returns None → `last_dream_at` stays absent.
Lenient success criteria: the guard only advances on completion of a
non-null DreamResult. Failed runs (None return) must not count.
"""
payload = DreamPayload(
dream_type=DreamType.OMNI,
observer=seeded_collection.observer,
observed=seeded_collection.observed,
)
with patch(
"src.dreamer.orchestrator.run_dream",
new=AsyncMock(return_value=None),
):
await process_dream(payload, seeded_collection.workspace_name)
dream_meta = await _get_dream_metadata(db_session, seeded_collection)
assert "last_dream_at" not in dream_meta, (
"last_dream_at must NOT be written when run_dream returns None "
"(failed dream). The guard should not falsely advance."
)
@pytest.mark.asyncio
async def test_exception_path_leaves_last_dream_at_null(
self,
db_session: AsyncSession,
seeded_collection: models.Collection,
):
"""run_dream raises → `last_dream_at` stays absent.
`process_dream` catches exceptions (logs + marks task processed without
re-raising) so the queue worker doesn't get stuck retrying. The guard
write must not happen in the exception path it's inside the
`if result is not None` block, which never executes if an exception
bypassed the assignment.
"""
payload = DreamPayload(
dream_type=DreamType.OMNI,
observer=seeded_collection.observer,
observed=seeded_collection.observed,
)
with patch(
"src.dreamer.orchestrator.run_dream",
new=AsyncMock(side_effect=RuntimeError("simulated specialist crash")),
):
# process_dream swallows exceptions internally; no re-raise expected
await process_dream(payload, seeded_collection.workspace_name)
dream_meta = await _get_dream_metadata(db_session, seeded_collection)
assert "last_dream_at" not in dream_meta, (
"last_dream_at must NOT be written when run_dream raises. "
"process_dream swallows the exception but the guard write must "
"not occur."
)