From d8f48302da711156b687a6299111b6c338720ab9 Mon Sep 17 00:00:00 2001 From: ChethanUK Date: Sun, 16 Aug 2026 18:25:01 +0200 Subject: [PATCH] fix(dreamer): record outcome and error metrics in collection dream metadata (#1003) --- src/dreamer/orchestrator.py | 83 ++++++- src/dreamer/specialists.py | 2 + tests/dreamer/test_dreamer_integration.py | 268 ++++++++++++++++++++-- 3 files changed, 336 insertions(+), 17 deletions(-) diff --git a/src/dreamer/orchestrator.py b/src/dreamer/orchestrator.py index 00f001f9..9f8ceb62 100644 --- a/src/dreamer/orchestrator.py +++ b/src/dreamer/orchestrator.py @@ -44,6 +44,8 @@ from src.utils.queue_payload import DreamPayload logger = logging.getLogger(__name__) +DREAM_FAILURE_ALERT_THRESHOLD = 3 + @dataclass class DreamResult: @@ -67,6 +69,11 @@ class DreamResult: input_tokens: int output_tokens: int + # Error details & created counts (honcho#1003) + created_observation_count: int = 0 + deduction_error_class: str | None = None + induction_error_class: str | None = None + async def run_dream( workspace_name: str, @@ -191,6 +198,10 @@ async def run_dream( total_iterations = 0 total_input_tokens = 0 total_output_tokens = 0 + # Track specialist error classes (honcho#1003) + deduction_error_class: str | None = None + induction_error_class: str | None = None + try: # Run deduction specialist (manages its own DB sessions) logger.info(f"[{run_id}] Running deduction specialist") @@ -217,6 +228,7 @@ async def run_dream( # propagate so the worker can shut down. SpecialistExecutionError # is no longer raised by `src/`, but the catch is broad enough to # cover provider/DB/tool/validation errors. + deduction_error_class = type(e).__name__ logger.error(f"[{run_id}] Deduction specialist failed: {e}", exc_info=True) accumulate_metric(task_name, "deduction_error", str(e), "blob") @@ -241,6 +253,7 @@ async def run_dream( ) induction_success = induction_result.success except Exception as e: + induction_error_class = type(e).__name__ logger.error(f"[{run_id}] Induction specialist failed: {e}", exc_info=True) accumulate_metric(task_name, "induction_error", str(e), "blob") @@ -297,6 +310,10 @@ async def run_dream( except Exception: # pragma: no cover - telemetry must not raise logger.debug("Failed to emit DreamRunEvent", exc_info=True) + total_created_observations = ( + deduction_result.created_observation_count if deduction_result else 0 + ) + (induction_result.created_observation_count if induction_result else 0) + return DreamResult( run_id=run_id, specialists_run=["deduction", "induction"], @@ -308,6 +325,9 @@ async def run_dream( total_duration_ms=duration_ms, input_tokens=total_input_tokens, output_tokens=total_output_tokens, + created_observation_count=total_created_observations, + deduction_error_class=deduction_error_class, + induction_error_class=induction_error_class, ) @@ -522,7 +542,8 @@ DREAM: {payload.dream_type} documents for {workspace_name}/{payload.observer}/{p + f"duration={result.total_duration_ms:.0f}ms" ) - # Both guard fields advance together only on successful consolidation. + # Both guard fields advance together on every dream run (fail-forward), + # while detailed outcome/error metrics are stored in dream metadata (honcho#1003). now_iso = datetime.now(timezone.utc).isoformat() async with tracked_db("dream.guard_pair_write") as db: collection = await crud.get_collection( @@ -542,6 +563,66 @@ DREAM: {payload.dream_type} documents for {workspace_name}/{payload.observer}/{p dream_meta = dict(collection.internal_metadata.get("dream", {})) dream_meta["last_dream_at"] = now_iso dream_meta["last_dream_document_count"] = current_explicit_count + + # Determine outcome: "success", "partial", or "failed" + succ_count = sum( + 1 + for s in ( + result.deduction_success, + result.induction_success, + ) + if s + ) + if succ_count == 2: + outcome = "success" + elif succ_count == 1: + outcome = "partial" + else: + outcome = "failed" + + dream_meta["last_dream_outcome"] = outcome + dream_meta["last_dream_run_id"] = result.run_id + dream_meta["last_dream_created_observations"] = ( + result.created_observation_count + ) + + # Form distinct error classes string (joined with comma, <= 2 entries) + err_classes = list( + dict.fromkeys( + cls + for cls in ( + result.deduction_error_class, + result.induction_error_class, + ) + if cls + ) + ) + dream_meta["last_dream_error_class"] = ( + ",".join(err_classes) if err_classes else None + ) + + if outcome in ("success", "partial"): + dream_meta["consecutive_failed_dreams"] = 0 + dream_meta["last_successful_dream_at"] = now_iso + dream_meta["last_successful_dream_document_count"] = ( + current_explicit_count + ) + else: + consecutive = ( + dream_meta.get("consecutive_failed_dreams", 0) + 1 + ) + dream_meta["consecutive_failed_dreams"] = consecutive + + logger.warning( + f"Dream cycle failed for {workspace_name}/{payload.observer}/{payload.observed}: " + f"run_id={result.run_id}, consecutive={consecutive}, errors={dream_meta['last_dream_error_class']}" + ) + if consecutive == DREAM_FAILURE_ALERT_THRESHOLD: + logger.error( + f"Dream failure threshold reached ({consecutive} consecutive failures) " + f"for {workspace_name}/{payload.observer}/{payload.observed}" + ) + await crud.update_collection_internal_metadata( db, workspace_name, diff --git a/src/dreamer/specialists.py b/src/dreamer/specialists.py index 28ff4577..0006a52d 100644 --- a/src/dreamer/specialists.py +++ b/src/dreamer/specialists.py @@ -66,6 +66,7 @@ class SpecialistResult: duration_ms: float success: bool content: str + created_observation_count: int = 0 # Tool names to exclude when peer card creation is disabled @@ -456,6 +457,7 @@ If you update it, send the full deduplicated list and remove stale entries. duration_ms=duration_ms, success=True, content=response.content, + created_observation_count=created_observation_count, ) except BaseException as e: # BaseException (not Exception) — asyncio.CancelledError doesn't diff --git a/tests/dreamer/test_dreamer_integration.py b/tests/dreamer/test_dreamer_integration.py index fb9b1894..a4dcf577 100644 --- a/tests/dreamer/test_dreamer_integration.py +++ b/tests/dreamer/test_dreamer_integration.py @@ -24,6 +24,7 @@ from src.dreamer.dream_scheduler import ( set_dream_scheduler, ) from src.dreamer.orchestrator import DreamResult, process_dream +from src.dreamer.specialists import SpecialistResult from src.schemas import ( DreamType, ResolvedConfiguration, @@ -54,19 +55,28 @@ async def seeded_collection( return collection -def _make_dream_result() -> DreamResult: - """Build a minimal non-null DreamResult for happy-path tests.""" +def _make_dream_result( + deduction_success: bool = True, + induction_success: bool = True, + created_count: int = 0, + deduction_error_class: str | None = None, + induction_error_class: str | None = None, +) -> DreamResult: + """Build a minimal non-null DreamResult for happy-path and outcome tests.""" return DreamResult( run_id="test_run_01", specialists_run=["deduction", "induction"], - deduction_success=True, - induction_success=True, + deduction_success=deduction_success, + induction_success=induction_success, surprisal_enabled=False, surprisal_conclusion_count=0, total_iterations=3, total_duration_ms=1234.5, input_tokens=100, output_tokens=50, + created_observation_count=created_count, + deduction_error_class=deduction_error_class, + induction_error_class=induction_error_class, ) @@ -104,19 +114,19 @@ class TestLastDreamAtCompletionWrite: 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" + assert "last_dream_at" in dream_meta, ( + "process_dream must write last_dream_at when run_dream returns a result" + ) # Must be a tz-aware UTC ISO timestamp. A naive datetime.now().isoformat() # would pass a loose "T in string" check but corrupt the 8h guard math # against tz-aware now() comparisons downstream. parsed = datetime.fromisoformat(dream_meta["last_dream_at"]) - assert ( - parsed.tzinfo is not None - ), f"last_dream_at must be timezone-aware, got {dream_meta['last_dream_at']!r}" - assert parsed.utcoffset() == timedelta( - 0 - ), f"last_dream_at must be UTC, got offset {parsed.utcoffset()}" + assert parsed.tzinfo is not None, ( + f"last_dream_at must be timezone-aware, got {dream_meta['last_dream_at']!r}" + ) + assert parsed.utcoffset() == timedelta(0), ( + f"last_dream_at must be UTC, got offset {parsed.utcoffset()}" + ) @pytest.mark.asyncio async def test_failure_path_leaves_last_dream_at_null( @@ -582,9 +592,9 @@ class TestGuardPairCoherence: "pre-Loop-4 the baseline was consumed at enqueue time and a " "silent failure would lock out retries on the same corpus." ) - assert ( - "last_dream_at" not in dream_meta - ), "Failed dream must not advance last_dream_at either." + assert "last_dream_at" not in dream_meta, ( + "Failed dream must not advance last_dream_at either." + ) with patch.object( _scheduler, "schedule_dream", new_callable=AsyncMock @@ -597,3 +607,229 @@ class TestGuardPairCoherence: "last_dream_at, no pending queue item." ) assert mock_schedule.called, "schedule_dream must be invoked on the retry path." + + +class TestDreamOutcomeMetadata: + """Regression & behavior tests for honcho#1003: metadata recording on dream outcome.""" + + @pytest.mark.parametrize( + ("deduction_succ", "induction_succ", "created_count", "expected_outcome"), + [ + (True, True, 5, "success"), + (True, False, 3, "partial"), + (False, True, 1, "partial"), + (False, False, 0, "failed"), + (True, True, 0, "success"), + ], + ) + @pytest.mark.asyncio + async def test_dream_outcome_classification_and_guard_advancement( + self, + db_session: AsyncSession, + seeded_collection: models.Collection, + deduction_succ: bool, + induction_succ: bool, + created_count: int, + expected_outcome: str, + ): + """Verify outcome classification, created count recording, and guard field advancement.""" + payload = DreamPayload( + dream_type=DreamType.OMNI, + observer=seeded_collection.observer, + observed=seeded_collection.observed, + ) + fake_result = _make_dream_result( + deduction_success=deduction_succ, + induction_success=induction_succ, + created_count=created_count, + ) + + with patch( + "src.dreamer.orchestrator.run_dream", + new=AsyncMock(return_value=fake_result), + ): + await process_dream(payload, seeded_collection.workspace_name) + + dream_meta = await _get_dream_metadata(db_session, seeded_collection) + assert dream_meta.get("last_dream_outcome") == expected_outcome + assert dream_meta.get("last_dream_created_observations") == created_count + assert "last_dream_at" in dream_meta + assert "last_dream_document_count" in dream_meta + + if expected_outcome in ("success", "partial"): + assert dream_meta.get("consecutive_failed_dreams") == 0 + assert "last_successful_dream_at" in dream_meta + assert "last_successful_dream_document_count" in dream_meta + else: + assert dream_meta.get("consecutive_failed_dreams") == 1 + assert "last_successful_dream_at" not in dream_meta + + @pytest.mark.asyncio + async def test_streak_failure_logging_and_recovery( + self, + db_session: AsyncSession, + seeded_collection: models.Collection, + caplog: pytest.LogCaptureFixture, + ): + """Assert consecutive failure counter increments, warning/error logs fire correctly, and recovery resets counter.""" + payload = DreamPayload( + dream_type=DreamType.OMNI, + observer=seeded_collection.observer, + observed=seeded_collection.observed, + ) + failed_result = _make_dream_result( + deduction_success=False, induction_success=False, created_count=0 + ) + + with patch( + "src.dreamer.orchestrator.run_dream", + new=AsyncMock(return_value=failed_result), + ): + for _ in range(4): + await process_dream(payload, seeded_collection.workspace_name) + + dream_meta = await _get_dream_metadata(db_session, seeded_collection) + assert dream_meta.get("consecutive_failed_dreams") == 4 + assert dream_meta.get("last_dream_outcome") == "failed" + + # Logs check: 4 WARNINGs, exactly 1 ERROR (at threshold = 3) + warnings = [ + r + for r in caplog.records + if r.levelname == "WARNING" and "Dream cycle failed" in r.message + ] + errors = [ + r + for r in caplog.records + if r.levelname == "ERROR" and "Dream failure threshold reached" in r.message + ] + assert len(warnings) == 4 + assert len(errors) == 1 + + # Recovery with a successful dream + success_result = _make_dream_result( + deduction_success=True, induction_success=True, created_count=2 + ) + with patch( + "src.dreamer.orchestrator.run_dream", + new=AsyncMock(return_value=success_result), + ): + await process_dream(payload, seeded_collection.workspace_name) + + dream_meta_after = await _get_dream_metadata(db_session, seeded_collection) + assert dream_meta_after.get("consecutive_failed_dreams") == 0 + assert dream_meta_after.get("last_dream_outcome") == "success" + assert "last_successful_dream_at" in dream_meta_after + + @pytest.mark.asyncio + async def test_error_class_formatting_and_deduplication( + self, + db_session: AsyncSession, + seeded_collection: models.Collection, + ): + """Verify error classes are recorded correctly when specialists raise exceptions in run_dream.""" + from src.dreamer.orchestrator import run_dream, SPECIALISTS + + mock_deduction = AsyncMock() + mock_induction = AsyncMock() + + succ_specialist_result = SpecialistResult( + run_id="test_spec_01", + specialist_type="induction", + iterations=1, + tool_calls_count=1, + input_tokens=10, + output_tokens=10, + duration_ms=100.0, + success=True, + content="Induction ok", + created_observation_count=0, + ) + + # Deduction raises ValueError, Induction succeeds + mock_deduction.run.side_effect = ValueError("Deduction error") + mock_induction.run.return_value = succ_specialist_result + + with patch.dict( + SPECIALISTS, {"deduction": mock_deduction, "induction": mock_induction} + ): + res1 = await run_dream( + workspace_name=seeded_collection.workspace_name, + observer=seeded_collection.observer, + observed=seeded_collection.observed, + ) + assert res1 is not None + assert res1.deduction_error_class == "ValueError" + assert res1.induction_error_class is None + + # Both raise ValueError (deduplication check) + mock_deduction.run.side_effect = ValueError("Error 1") + mock_induction.run.side_effect = ValueError("Error 2") + with patch.dict( + SPECIALISTS, {"deduction": mock_deduction, "induction": mock_induction} + ): + res2 = await run_dream( + workspace_name=seeded_collection.workspace_name, + observer=seeded_collection.observer, + observed=seeded_collection.observed, + ) + assert res2 is not None + assert res2.deduction_error_class == "ValueError" + assert res2.induction_error_class == "ValueError" + + # Deduction raises KeyError, Induction raises TypeError + mock_deduction.run.side_effect = KeyError("Key error") + mock_induction.run.side_effect = TypeError("Type error") + with patch.dict( + SPECIALISTS, {"deduction": mock_deduction, "induction": mock_induction} + ): + res3 = await run_dream( + workspace_name=seeded_collection.workspace_name, + observer=seeded_collection.observer, + observed=seeded_collection.observed, + ) + assert res3 is not None + assert res3.deduction_error_class == "KeyError" + assert res3.induction_error_class == "TypeError" + + @pytest.mark.asyncio + async def test_backward_compatibility_existing_metadata( + self, + db_session: AsyncSession, + sample_data: tuple[models.Workspace, models.Peer], + ): + """Legacy collection metadata without outcome keys updates cleanly on first failure.""" + workspace, peer = sample_data + collection = models.Collection( + observer=peer.name, + observed=peer.name, + workspace_name=workspace.name, + internal_metadata={ + "dream": { + "last_dream_at": "2026-01-01T00:00:00", + "last_dream_document_count": 10, + } + }, + ) + db_session.add(collection) + await db_session.commit() + await db_session.refresh(collection) + + payload = DreamPayload( + dream_type=DreamType.OMNI, + observer=peer.name, + observed=peer.name, + ) + failed_result = _make_dream_result( + deduction_success=False, induction_success=False, created_count=0 + ) + + with patch( + "src.dreamer.orchestrator.run_dream", + new=AsyncMock(return_value=failed_result), + ): + await process_dream(payload, workspace.name) + + dream_meta = await _get_dream_metadata(db_session, collection) + assert dream_meta.get("consecutive_failed_dreams") == 1 + assert dream_meta.get("last_dream_outcome") == "failed"