From 72c29cce7b6adc242973d81a942297dc1e67dd46 Mon Sep 17 00:00:00 2001 From: Luis Mendes Date: Tue, 31 Mar 2026 02:54:46 +0000 Subject: [PATCH] feat(events): wire all 38 domain event emissions into services Added event_bus.emit() calls for all 32 previously unwired EventType members across domain services. Each emission includes contextual data (plan_id, actor_name, tool_name, resource_id, etc.) for subscribers. Events are fire-and-forget with error handling to prevent emission failures from affecting service logic. Services modified: plan_lifecycle_service, decision_service, invariant services, actor runtime, tool execution, resource services, sandbox/checkpoint, ACMS pipeline, validation runner, session service. Added Behave tests verifying event emission for critical paths. ISSUES CLOSED: #923 --- features/event_emission_wiring.feature | 145 +++++ features/steps/actor_runtime_steps.py | 4 + features/steps/event_emission_wiring_steps.py | 544 ++++++++++++++++++ features/steps/m2_actor_tool_smoke_steps.py | 3 +- .../steps/safety_profile_enforcement_steps.py | 3 +- .../tool_lifecycle_coverage_boost_steps.py | 3 +- .../steps/tool_lifecycle_runtime_steps.py | 3 +- .../application/services/acms_pipeline.py | 29 + .../services/checkpoint_service.py | 101 +++- .../application/services/context_service.py | 27 +- .../application/services/decision_service.py | 107 +++- .../application/services/invariant_service.py | 84 ++- .../services/plan_lifecycle_service.py | 99 ++++ .../services/repo_indexing_service.py | 36 +- .../services/resource_handler_service.py | 31 +- .../application/services/session_service.py | 17 + .../services/validation_pipeline.py | 45 +- src/cleveragents/tool/actor_runtime.py | 61 +- src/cleveragents/tool/lifecycle.py | 58 +- 19 files changed, 1362 insertions(+), 38 deletions(-) create mode 100644 features/event_emission_wiring.feature create mode 100644 features/steps/event_emission_wiring_steps.py diff --git a/features/event_emission_wiring.feature b/features/event_emission_wiring.feature new file mode 100644 index 000000000..f40652153 --- /dev/null +++ b/features/event_emission_wiring.feature @@ -0,0 +1,145 @@ +Feature: Domain event emission wiring across services + As an infrastructure operator + I want all 38 EventType values to be emitted from the correct services + So that the audit log, observability, and reactive subscribers work correctly + + # --------------------------------------------------------------------------- + # Plan Lifecycle event emissions + # --------------------------------------------------------------------------- + + Scenario: PlanLifecycleService emits PLAN_STATE_CHANGED on complete_strategize + Given an in-memory PlanLifecycleService with a tracking event bus + And a plan in Strategize/PROCESSING phase + When I call complete_strategize on that plan + Then ew the audit log should contain a "plan.state_changed" event + And ew the event details should contain key "processing_state" + + Scenario: PlanLifecycleService emits PLAN_ERRORED on fail_strategize + Given an in-memory PlanLifecycleService with a tracking event bus + And a plan in Strategize/PROCESSING phase + When I call fail_strategize with error "Test error" + Then ew the audit log should contain a "plan.errored" event + And ew the event details should contain key "error_message" + + Scenario: PlanLifecycleService emits PLAN_ERRORED on fail_execute + Given an in-memory PlanLifecycleService with a tracking event bus + And a plan in Execute/PROCESSING phase + When I call fail_execute with error "Execution failed" + Then ew the audit log should contain a "plan.errored" event + + Scenario: PlanLifecycleService emits PLAN_ERRORED on fail_apply + Given an in-memory PlanLifecycleService with a tracking event bus + And a plan in Apply/PROCESSING phase + When I call fail_apply with error "Apply failed" + Then ew the audit log should contain a "plan.errored" event + + # --------------------------------------------------------------------------- + # Decision event emissions + # --------------------------------------------------------------------------- + + Scenario: DecisionService emits DECISION_APPROVED for non-correction decisions + Given an in-memory DecisionService with a tracking event bus + When I record a normal decision + Then ew the audit log should contain a "decision.approved" event + And ew the audit log should contain a "decision.created" event + + Scenario: DecisionService emits DECISION_CORRECTED for correction decisions + Given an in-memory DecisionService with a tracking event bus + When I record a correction decision + Then ew the audit log should contain a "decision.corrected" event + + Scenario: DecisionService emits DECISION_SUPERSEDED on mark_superseded + Given an in-memory DecisionService with a tracking event bus + And ew two decisions exist for the same plan + When ew I mark the first decision as superseded by the second + Then ew the audit log should contain a "decision.superseded" event + + # --------------------------------------------------------------------------- + # Invariant event emissions + # --------------------------------------------------------------------------- + + Scenario: InvariantService emits INVARIANT_ENFORCED and INVARIANT_RECONCILED + Given an InvariantService with a tracking event bus + And an active invariant exists + When I enforce invariants for a plan + Then ew the audit log should contain a "invariant.enforced" event + And ew the audit log should contain a "invariant.reconciled" event + + Scenario: InvariantService emits INVARIANT_VIOLATED for violated invariants + Given an InvariantService with a tracking event bus + And an active invariant exists with known ID + When I enforce invariants with violations for that invariant + Then ew the audit log should contain a "invariant.violated" event + + # --------------------------------------------------------------------------- + # Session event emissions + # --------------------------------------------------------------------------- + + Scenario: SessionService emits SESSION_MESSAGE_SENT on append_message + Given a PersistentSessionService with a tracking event bus + And a session exists + When I append a message to the session + Then ew the audit log should contain a "session.message_sent" event + And ew the event details should contain key "message_id" + + # --------------------------------------------------------------------------- + # Checkpoint / Sandbox event emissions + # --------------------------------------------------------------------------- + + Scenario: CheckpointService emits CHECKPOINT_CREATED on create_checkpoint + Given a CheckpointService with a tracking event bus + When I create a checkpoint + Then ew the audit log should contain a "checkpoint.created" event + And ew the event details should contain key "checkpoint_id" + + Scenario: CheckpointService emits SANDBOX_CREATED on register_sandbox + Given a CheckpointService with a tracking event bus + When I register a sandbox for a plan + Then ew the audit log should contain a "sandbox.created" event + + Scenario: CheckpointService emits SANDBOX_COMMITTED on mark_plan_applied + Given a CheckpointService with a tracking event bus + When I mark a plan as applied + Then ew the audit log should contain a "sandbox.committed" event + + # --------------------------------------------------------------------------- + # Validation event emissions + # --------------------------------------------------------------------------- + + Scenario: ValidationPipeline emits VALIDATION_STARTED and VALIDATION_PASSED + Given a ValidationPipeline with a tracking event bus and a passing validation + When I run the pipeline + Then ew the audit log should contain a "validation.started" event + And ew the audit log should contain a "validation.passed" event + + Scenario: ValidationPipeline emits VALIDATION_FAILED when required fails + Given a ValidationPipeline with a tracking event bus and a failing validation + When I run the pipeline + Then ew the audit log should contain a "validation.started" event + And ew the audit log should contain a "validation.failed" event + + # --------------------------------------------------------------------------- + # Tool Lifecycle event emissions + # --------------------------------------------------------------------------- + + Scenario: ToolRuntime emits TOOL_INVOKED and TOOL_COMPLETED on execute + Given a ToolRuntime with a tracking event bus and a registered tool + When I execute the tool + Then ew the audit log should contain a "tool.invoked" event + And ew the audit log should contain a "tool.completed" event + + Scenario: ToolRuntime emits TOOL_ERRORED when execution fails + Given a ToolRuntime with a tracking event bus and a failing tool + When I execute the failing tool expecting an error + Then ew the audit log should contain a "tool.invoked" event + And ew the audit log should contain a "tool.errored" event + + # --------------------------------------------------------------------------- + # Actor Runtime event emissions + # --------------------------------------------------------------------------- + + Scenario: ToolCallingRuntime emits ACTOR_INVOKED and ACTOR_COMPLETED + Given a ToolCallingRuntime with a tracking event bus and a simple LLM + When I run the tool loop with a simple prompt + Then ew the audit log should contain a "actor.invoked" event + And ew the audit log should contain a "actor.completed" event diff --git a/features/steps/actor_runtime_steps.py b/features/steps/actor_runtime_steps.py index 383015b37..b3b136f23 100644 --- a/features/steps/actor_runtime_steps.py +++ b/features/steps/actor_runtime_steps.py @@ -7,6 +7,7 @@ from typing import Any from behave import given, then, when +from cleveragents.infrastructure.events.reactive import ReactiveEventBus from cleveragents.tool.actor_context import ToolActorContext, ToolCallRecord from cleveragents.tool.actor_runtime import ( LLMResponse, @@ -346,6 +347,7 @@ def step_given_runtime(context: Any) -> None: registry=context.registry, runner=context.runner, llm_caller=context.llm_caller, + event_bus=ReactiveEventBus(), ) @@ -356,6 +358,7 @@ def step_given_runtime_max_iter(context: Any, n: int) -> None: runner=context.runner, llm_caller=context.llm_caller, max_iterations=n, + event_bus=ReactiveEventBus(), ) @@ -366,6 +369,7 @@ def step_given_runtime_with_router(context: Any) -> None: runner=context.runner, llm_caller=context.llm_caller, router=context.router, + event_bus=ReactiveEventBus(), ) diff --git a/features/steps/event_emission_wiring_steps.py b/features/steps/event_emission_wiring_steps.py new file mode 100644 index 000000000..982245b14 --- /dev/null +++ b/features/steps/event_emission_wiring_steps.py @@ -0,0 +1,544 @@ +"""Step definitions for event_emission_wiring.feature. + +Tests that domain services correctly emit events through the EventBus +for all 32 newly wired event types. +""" + +from __future__ import annotations + +from typing import Any +from unittest.mock import MagicMock + +from behave import given, then, when # type: ignore[import-untyped] +from behave.runner import Context # type: ignore[import-untyped] +from ulid import ULID + +from cleveragents.application.services.checkpoint_service import CheckpointService +from cleveragents.application.services.decision_service import DecisionService +from cleveragents.application.services.invariant_service import InvariantService +from cleveragents.application.services.plan_lifecycle_service import ( + PlanLifecycleService, +) +from cleveragents.application.services.session_service import ( + PersistentSessionService, +) +from cleveragents.application.services.validation_pipeline import ( + ValidationCommand, + ValidationPipeline, +) +from cleveragents.domain.models.core.invariant import InvariantScope +from cleveragents.domain.models.core.plan import PlanPhase, ProcessingState +from cleveragents.domain.models.core.session import MessageRole +from cleveragents.domain.models.core.tool import ValidationMode +from cleveragents.infrastructure.events.models import DomainEvent +from cleveragents.infrastructure.events.reactive import ReactiveEventBus + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + + +def _make_settings() -> Any: + mock = MagicMock() + mock.async_enabled = False + mock.cleveragents_max_context_size = 50 * 1024 * 1024 + return mock + + +def _get_events(ctx: Context) -> list[DomainEvent]: + return ctx.ew_event_bus.audit_log # type: ignore[no-any-return] + + +def _find_event(ctx: Context, event_type_value: str) -> DomainEvent | None: + for ev in _get_events(ctx): + if ev.event_type.value == event_type_value: + return ev + return None + + +# --------------------------------------------------------------------------- +# Plan Lifecycle steps +# --------------------------------------------------------------------------- + + +@given("an in-memory PlanLifecycleService with a tracking event bus") +def step_plan_service_with_bus(ctx: Context) -> None: + ctx.ew_event_bus = ReactiveEventBus() + ctx.ew_plan_svc = PlanLifecycleService( + settings=_make_settings(), + event_bus=ctx.ew_event_bus, + ) + + +def _create_plan_in_phase( + svc: PlanLifecycleService, + action_name: str, + phase: PlanPhase, +) -> str: + """Helper to create a plan and set it to a specific phase/processing state.""" + svc.create_action( + name=action_name, + description="Test action", + definition_of_done="Done", + strategy_actor="local/strategy", + execution_actor="local/execution", + ) + plan = svc.use_action(action_name) + plan.processing_state = ProcessingState.PROCESSING + plan.phase = phase + svc._commit_plan(plan) + return plan.identity.plan_id + + +@given("a plan in Strategize/PROCESSING phase") +def step_plan_strategize_processing(ctx: Context) -> None: + svc: PlanLifecycleService = ctx.ew_plan_svc + ctx.ew_plan_id = _create_plan_in_phase(svc, "local/ew-action", PlanPhase.STRATEGIZE) + ctx.ew_event_bus.clear_audit_log() + + +@given("a plan in Execute/PROCESSING phase") +def step_plan_execute_processing(ctx: Context) -> None: + svc: PlanLifecycleService = ctx.ew_plan_svc + ctx.ew_plan_id = _create_plan_in_phase( + svc, "local/ew-exec-action", PlanPhase.EXECUTE + ) + ctx.ew_event_bus.clear_audit_log() + + +@given("a plan in Apply/PROCESSING phase") +def step_plan_apply_processing(ctx: Context) -> None: + svc: PlanLifecycleService = ctx.ew_plan_svc + ctx.ew_plan_id = _create_plan_in_phase( + svc, "local/ew-apply-action", PlanPhase.APPLY + ) + ctx.ew_event_bus.clear_audit_log() + + +@when("I call complete_strategize on that plan") +def step_complete_strategize(ctx: Context) -> None: + ctx.ew_plan_svc.complete_strategize(ctx.ew_plan_id) + + +@when('I call fail_strategize with error "{error}"') +def step_fail_strategize(ctx: Context, error: str) -> None: + ctx.ew_plan_svc.fail_strategize(ctx.ew_plan_id, error) + + +@when('I call fail_execute with error "{error}"') +def step_fail_execute(ctx: Context, error: str) -> None: + ctx.ew_plan_svc.fail_execute(ctx.ew_plan_id, error) + + +@when('I call fail_apply with error "{error}"') +def step_fail_apply(ctx: Context, error: str) -> None: + ctx.ew_plan_svc.fail_apply(ctx.ew_plan_id, error) + + +# --------------------------------------------------------------------------- +# Decision steps +# --------------------------------------------------------------------------- + + +@given("an in-memory DecisionService with a tracking event bus") +def step_decision_service_with_bus(ctx: Context) -> None: + ctx.ew_event_bus = ReactiveEventBus() + ctx.ew_decision_svc = DecisionService(event_bus=ctx.ew_event_bus) + + +@when("I record a normal decision") +def step_record_normal_decision(ctx: Context) -> None: + plan_id = str(ULID()) + ctx.ew_decision_svc.record_decision( + plan_id=plan_id, + decision_type="strategy_choice", + question="Which strategy?", + chosen_option="Option A", + ) + + +@when("I record a correction decision") +def step_record_correction_decision(ctx: Context) -> None: + plan_id = str(ULID()) + ctx.ew_decision_svc.record_decision( + plan_id=plan_id, + decision_type="strategy_choice", + question="Which strategy?", + chosen_option="Option B", + is_correction=True, + corrects_decision_id=str(ULID()), + correction_reason="Correcting previous", + ) + + +@given("ew two decisions exist for the same plan") +def step_two_decisions(ctx: Context) -> None: + plan_id = str(ULID()) + ctx.ew_plan_id_for_decisions = plan_id + d1 = ctx.ew_decision_svc.record_decision( + plan_id=plan_id, + decision_type="strategy_choice", + question="Q1", + chosen_option="O1", + ) + d2 = ctx.ew_decision_svc.record_decision( + plan_id=plan_id, + decision_type="strategy_choice", + question="Q2", + chosen_option="O2", + ) + ctx.ew_decision_id_1 = d1.decision_id + ctx.ew_decision_id_2 = d2.decision_id + ctx.ew_event_bus.clear_audit_log() + + +@when("ew I mark the first decision as superseded by the second") +def step_mark_superseded(ctx: Context) -> None: + ctx.ew_decision_svc.mark_superseded(ctx.ew_decision_id_1, ctx.ew_decision_id_2) + + +# --------------------------------------------------------------------------- +# Invariant steps +# --------------------------------------------------------------------------- + + +@given("an InvariantService with a tracking event bus") +def step_invariant_service_with_bus(ctx: Context) -> None: + ctx.ew_event_bus = ReactiveEventBus() + ctx.ew_invariant_svc = InvariantService(event_bus=ctx.ew_event_bus) + + +@given("an active invariant exists") +def step_active_invariant(ctx: Context) -> None: + inv = ctx.ew_invariant_svc.add_invariant( + text="Must not delete files", + scope=InvariantScope.GLOBAL, + source_name="test", + ) + ctx.ew_invariant = inv + + +@given("an active invariant exists with known ID") +def step_active_invariant_known_id(ctx: Context) -> None: + inv = ctx.ew_invariant_svc.add_invariant( + text="Must not modify production data", + scope=InvariantScope.GLOBAL, + source_name="test", + ) + ctx.ew_invariant = inv + ctx.ew_invariant_id = inv.id + + +@when("I enforce invariants for a plan") +def step_enforce_invariants(ctx: Context) -> None: + ctx.ew_event_bus.clear_audit_log() + ctx.ew_invariant_svc.enforce_invariants( + plan_id=str(ULID()), + invariants=[ctx.ew_invariant], + ) + + +@when("I enforce invariants with violations for that invariant") +def step_enforce_with_violations(ctx: Context) -> None: + ctx.ew_event_bus.clear_audit_log() + ctx.ew_invariant_svc.enforce_invariants( + plan_id=str(ULID()), + invariants=[ctx.ew_invariant], + violated_invariant_ids=[ctx.ew_invariant_id], + ) + + +# --------------------------------------------------------------------------- +# Session steps +# --------------------------------------------------------------------------- + + +@given("a PersistentSessionService with a tracking event bus") +def step_session_service_with_bus(ctx: Context) -> None: + ctx.ew_event_bus = ReactiveEventBus() + mock_session_repo = MagicMock() + mock_message_repo = MagicMock() + mock_session_repo.get_by_id.return_value = MagicMock(updated_at=None) + mock_message_repo.count_for_session.return_value = 0 + ctx.ew_session_svc = PersistentSessionService( + session_repo=mock_session_repo, + message_repo=mock_message_repo, + event_bus=ctx.ew_event_bus, + ) + + +@given("a session exists") +def step_session_exists(ctx: Context) -> None: + session = ctx.ew_session_svc.create() + ctx.ew_session_id = session.session_id + ctx.ew_event_bus.clear_audit_log() + + +@when("I append a message to the session") +def step_append_message(ctx: Context) -> None: + ctx.ew_session_svc.append_message( + ctx.ew_session_id, + MessageRole.USER, + "Hello, world!", + ) + + +# --------------------------------------------------------------------------- +# Checkpoint / Sandbox steps +# --------------------------------------------------------------------------- + + +@given("a CheckpointService with a tracking event bus") +def step_checkpoint_service_with_bus(ctx: Context) -> None: + ctx.ew_event_bus = ReactiveEventBus() + ctx.ew_checkpoint_svc = CheckpointService(event_bus=ctx.ew_event_bus) + + +@when("I create a checkpoint") +def step_create_checkpoint(ctx: Context) -> None: + ctx.ew_checkpoint_svc.create_checkpoint( + plan_id=str(ULID()), + sandbox_ref="abc123", + reason="Test checkpoint", + ) + + +@when("I register a sandbox for a plan") +def step_register_sandbox(ctx: Context) -> None: + ctx.ew_checkpoint_svc.register_sandbox( + plan_id=str(ULID()), + sandbox_ref="sandbox-ref-abc", + ) + + +@when("I mark a plan as applied") +def step_mark_applied(ctx: Context) -> None: + ctx.ew_checkpoint_svc.mark_plan_applied(str(ULID())) + + +# --------------------------------------------------------------------------- +# Validation steps +# --------------------------------------------------------------------------- + + +@given("a ValidationPipeline with a tracking event bus and a passing validation") +def step_validation_passing(ctx: Context) -> None: + ctx.ew_event_bus = ReactiveEventBus() + + def _executor(name: str, config: dict[str, Any]) -> dict[str, Any]: + return {"passed": True, "message": "OK"} + + cmd = ValidationCommand( + validation_name="test-check", + resource_id="res-1", + resource_name="test-resource", + mode=ValidationMode.REQUIRED, + ) + ctx.ew_pipeline = ValidationPipeline( + commands=[cmd], + executor=_executor, + event_bus=ctx.ew_event_bus, + ) + + +@given("a ValidationPipeline with a tracking event bus and a failing validation") +def step_validation_failing(ctx: Context) -> None: + ctx.ew_event_bus = ReactiveEventBus() + + def _executor(name: str, config: dict[str, Any]) -> dict[str, Any]: + return {"passed": False, "message": "FAIL"} + + cmd = ValidationCommand( + validation_name="test-check-fail", + resource_id="res-2", + resource_name="test-resource-fail", + mode=ValidationMode.REQUIRED, + ) + ctx.ew_pipeline = ValidationPipeline( + commands=[cmd], + executor=_executor, + event_bus=ctx.ew_event_bus, + ) + + +@when("I run the pipeline") +def step_run_pipeline(ctx: Context) -> None: + ctx.ew_pipeline.run() + + +# --------------------------------------------------------------------------- +# Common assertion steps +# --------------------------------------------------------------------------- + + +@then('ew the audit log should contain a "{event_type}" event') +def step_ew_event_bus_received(ctx: Context, event_type: str) -> None: + ev = _find_event(ctx, event_type) + assert ev is not None, ( + f"Expected event '{event_type}' not found. " + f"Got: {[e.event_type.value for e in _get_events(ctx)]}" + ) + + +@then('ew the event details should contain key "{key}"') +def step_ew_event_details_key(ctx: Context, key: str) -> None: + events = _get_events(ctx) + found = False + for ev in events: + if key in ev.details: + found = True + break + assert found, ( + f"Key '{key}' not found in any event details. " + f"Events: {[(e.event_type.value, list(e.details.keys())) for e in events]}" + ) + + +# --------------------------------------------------------------------------- +# Tool Lifecycle steps +# --------------------------------------------------------------------------- + + +@given("a ToolRuntime with a tracking event bus and a registered tool") +def step_tool_runtime_with_bus(ctx: Context) -> None: + from cleveragents.tool.lifecycle import ( + ToolDescriptor, + ToolExecutionContext, + ToolResult, + ToolRuntime, + ) + + ctx.ew_event_bus = ReactiveEventBus() + + class _PassingTool: + def discover(self) -> ToolDescriptor: + return ToolDescriptor(name="test-tool", description="A test tool") + + def activate(self, ctx: ToolExecutionContext) -> None: + pass + + def execute( + self, params: dict[str, Any], ctx: ToolExecutionContext + ) -> ToolResult: + return ToolResult(success=True, data={"result": "ok"}) + + def deactivate(self, ctx: ToolExecutionContext) -> None: + pass + + from cleveragents.domain.models.core.tool import Tool, ToolCapability + + tool = Tool( + name="local/test-tool", + description="Test", + source="builtin", + tool_type="tool", + timeout=30, + capability=ToolCapability(), + ) + ctx.ew_tool_runtime = ToolRuntime( + tools={"local/test-tool": tool}, + instances={"local/test-tool": _PassingTool()}, + event_bus=ctx.ew_event_bus, + ) + ctx.ew_tool_ctx = ToolExecutionContext(plan_id="test-plan-001") + + +@when("I execute the tool") +def step_execute_tool(ctx: Context) -> None: + ctx.ew_tool_runtime.execute("local/test-tool", {}, ctx.ew_tool_ctx) + + +@given("a ToolRuntime with a tracking event bus and a failing tool") +def step_tool_runtime_failing(ctx: Context) -> None: + from cleveragents.tool.lifecycle import ( + ToolDescriptor, + ToolExecutionContext, + ToolResult, + ToolRuntime, + ) + + ctx.ew_event_bus = ReactiveEventBus() + + class _FailingTool: + def discover(self) -> ToolDescriptor: + return ToolDescriptor(name="fail-tool", description="A failing tool") + + def activate(self, ctx: ToolExecutionContext) -> None: + pass + + def execute( + self, params: dict[str, Any], ctx: ToolExecutionContext + ) -> ToolResult: + raise RuntimeError("Tool execution failed!") + + def deactivate(self, ctx: ToolExecutionContext) -> None: + pass + + from cleveragents.domain.models.core.tool import Tool, ToolCapability + + tool = Tool( + name="local/fail-tool", + description="Failing test", + source="builtin", + tool_type="tool", + timeout=30, + capability=ToolCapability(), + ) + ctx.ew_tool_runtime = ToolRuntime( + tools={"local/fail-tool": tool}, + instances={"local/fail-tool": _FailingTool()}, + event_bus=ctx.ew_event_bus, + ) + ctx.ew_tool_ctx = ToolExecutionContext(plan_id="test-plan-002") + + +@when("I execute the failing tool expecting an error") +def step_execute_failing_tool(ctx: Context) -> None: + import contextlib + + from cleveragents.tool.lifecycle import ToolExecutionError + + with contextlib.suppress(ToolExecutionError): + ctx.ew_tool_runtime.execute("local/fail-tool", {}, ctx.ew_tool_ctx) + + +# --------------------------------------------------------------------------- +# Actor Runtime steps +# --------------------------------------------------------------------------- + + +@given("a ToolCallingRuntime with a tracking event bus and a simple LLM") +def step_actor_runtime_with_bus(ctx: Context) -> None: + from cleveragents.tool.actor_runtime import ( + LLMResponse, + ToolCallingRuntime, + ) + from cleveragents.tool.registry import ToolRegistry + from cleveragents.tool.runner import ToolRunner + + ctx.ew_event_bus = ReactiveEventBus() + + registry = ToolRegistry() + runner = ToolRunner(registry=registry) + + class _SimpleLLM: + def invoke( + self, + prompt: str, + tool_schemas: list[dict[str, Any]], + tool_results: list[dict[str, Any]] | None = None, + actor_config: dict[str, Any] | None = None, + ) -> LLMResponse: + return LLMResponse(content="Hello!") + + ctx.ew_actor_runtime = ToolCallingRuntime( + registry=registry, + runner=runner, + llm_caller=_SimpleLLM(), + event_bus=ctx.ew_event_bus, + ) + + +@when("I run the tool loop with a simple prompt") +def step_run_tool_loop(ctx: Context) -> None: + ctx.ew_actor_runtime.run_tool_loop("Say hello") diff --git a/features/steps/m2_actor_tool_smoke_steps.py b/features/steps/m2_actor_tool_smoke_steps.py index 2559a0c50..f46429cd2 100644 --- a/features/steps/m2_actor_tool_smoke_steps.py +++ b/features/steps/m2_actor_tool_smoke_steps.py @@ -25,6 +25,7 @@ from cleveragents.domain.models.core.tool import ( ToolSource, ToolType, ) +from cleveragents.infrastructure.events.reactive import ReactiveEventBus from cleveragents.skills.context import SkillContext from cleveragents.skills.protocol import SkillDefinition, SkillMetadata from cleveragents.skills.registry import SkillRegistry @@ -312,7 +313,7 @@ def step_m2_create_tool_runtime(context: Context, name: str) -> None: timeout=300, ) instance = _M2MockToolInstance(name) - runtime = ToolRuntime() + runtime = ToolRuntime(event_bus=ReactiveEventBus()) runtime.register_tool(tool, instance) context._m2_runtime = runtime context._m2_mock_instance = instance diff --git a/features/steps/safety_profile_enforcement_steps.py b/features/steps/safety_profile_enforcement_steps.py index 5a7946ab9..cc97de17c 100644 --- a/features/steps/safety_profile_enforcement_steps.py +++ b/features/steps/safety_profile_enforcement_steps.py @@ -19,6 +19,7 @@ from cleveragents.domain.models.core.tool import ( ToolCapability, ToolSource, ) +from cleveragents.infrastructure.events.reactive import ReactiveEventBus from cleveragents.tool.context import ToolExecutionContext from cleveragents.tool.lifecycle import ( ToolAccessDeniedError, @@ -88,7 +89,7 @@ def _make_tool( @given('a registered tool "test/writer" that writes and is safe') def step_register_writer(context: Context) -> None: """Register a tool that writes but is not unsafe.""" - context.enforcement_runtime = ToolRuntime() + context.enforcement_runtime = ToolRuntime(event_bus=ReactiveEventBus()) tool = _make_tool("test/writer", writes=True, unsafe=False) desc = ToolDescriptor( name=tool.name, diff --git a/features/steps/tool_lifecycle_coverage_boost_steps.py b/features/steps/tool_lifecycle_coverage_boost_steps.py index 56a82bea2..fb25002c7 100644 --- a/features/steps/tool_lifecycle_coverage_boost_steps.py +++ b/features/steps/tool_lifecycle_coverage_boost_steps.py @@ -26,6 +26,7 @@ from cleveragents.domain.models.core.tool import ( ToolSource, ToolType, ) +from cleveragents.infrastructure.events.reactive import ReactiveEventBus from cleveragents.tool.context import ( CancellationToken, ToolCancelledError, @@ -202,7 +203,7 @@ def step_cb_cache_get_not_none(context: Context, plan_id: str, tool_name: str) - @given("I create a coverage-boost tool runtime") def step_create_cb_runtime(context: Context) -> None: - context.cb_runtime = ToolRuntime() + context.cb_runtime = ToolRuntime(event_bus=ReactiveEventBus()) context.cb_mocks: dict[str, _CBMockToolInstance] = {} context.cb_error = None diff --git a/features/steps/tool_lifecycle_runtime_steps.py b/features/steps/tool_lifecycle_runtime_steps.py index d2db68438..0ba084a3a 100644 --- a/features/steps/tool_lifecycle_runtime_steps.py +++ b/features/steps/tool_lifecycle_runtime_steps.py @@ -14,6 +14,7 @@ from cleveragents.domain.models.core.tool import ( ToolSource, ToolType, ) +from cleveragents.infrastructure.events.reactive import ReactiveEventBus from cleveragents.tool.context import ( BoundResource, CancellationToken, @@ -607,7 +608,7 @@ def _make_tool( @given("I create a tool runtime") def step_create_runtime(context: Context) -> None: - context.tool_runtime = ToolRuntime() + context.tool_runtime = ToolRuntime(event_bus=ReactiveEventBus()) context.mock_instances = dict[str, MockToolInstance]() diff --git a/src/cleveragents/application/services/acms_pipeline.py b/src/cleveragents/application/services/acms_pipeline.py index 515f83c5f..6c3f774c7 100644 --- a/src/cleveragents/application/services/acms_pipeline.py +++ b/src/cleveragents/application/services/acms_pipeline.py @@ -39,6 +39,7 @@ from pydantic import BaseModel, Field if TYPE_CHECKING: from cleveragents.config.settings import Settings from cleveragents.infrastructure.database.unit_of_work import UnitOfWork + from cleveragents.infrastructure.events.protocol import EventBus from cleveragents.application.services.acms_phase2 import ( GreedyKnapsackPacker, @@ -76,6 +77,8 @@ from cleveragents.domain.models.core.context_fragment import ( compute_context_hash, ) from cleveragents.domain.models.core.context_policy import ContextView +from cleveragents.infrastructure.events.models import DomainEvent +from cleveragents.infrastructure.events.types import EventType logger = structlog.get_logger(__name__) @@ -512,6 +515,7 @@ class ContextAssemblyPipeline(ACMSPipeline): *, settings: Settings | None = None, unit_of_work: UnitOfWork | None = None, + event_bus: EventBus | None = None, # Phase 1 — Strategy Orchestration strategy_selector: StrategySelector | None = None, budget_allocator: BudgetAllocator | None = None, @@ -561,6 +565,7 @@ class ContextAssemblyPipeline(ACMSPipeline): ) self._last_timings: StageTimings | None = None self._pipeline_logger = logger.bind(service="context_assembly_pipeline") + self._event_bus = event_bus @property def last_timings(self) -> StageTimings | None: @@ -695,6 +700,30 @@ class ContextAssemblyPipeline(ACMSPipeline): total_ms=round(total_ms, 3), ) + # Emit CONTEXT_BUILT event + if self._event_bus is not None: + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.CONTEXT_BUILT, + plan_id=plan_id, + details={ + "strategy": strategy_name, + "fragments_selected": len(final_fragments), + "total_tokens": total_tokens, + "budget_used": round(budget_used, 4), + "total_ms": round(total_ms, 3), + }, + ) + ) + except Exception: + self._pipeline_logger.warning( + "event_bus_emit_failed", + event_type="CONTEXT_BUILT", + plan_id=plan_id, + exc_info=True, + ) + return ContextPayload( plan_id=plan_id, fragments=final_fragments, diff --git a/src/cleveragents/application/services/checkpoint_service.py b/src/cleveragents/application/services/checkpoint_service.py index 4ed25d07c..3a7e6e2a0 100644 --- a/src/cleveragents/application/services/checkpoint_service.py +++ b/src/cleveragents/application/services/checkpoint_service.py @@ -102,6 +102,23 @@ class CheckpointService: plan_id: The plan to mark. """ self._plan_applied.add(plan_id) + # Emit SANDBOX_COMMITTED event when plan is applied + if self._event_bus is not None: + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.SANDBOX_COMMITTED, + plan_id=plan_id, + details={ + "sandbox_ref": self._plan_sandbox_refs.get(plan_id, ""), + }, + ) + ) + except Exception: + logger.warning( + "event_bus_emit_failed: SANDBOX_COMMITTED for plan %s", + plan_id, + ) def register_sandbox(self, plan_id: str, sandbox_ref: str) -> None: """Register a sandbox reference for a plan. @@ -111,6 +128,23 @@ class CheckpointService: sandbox_ref: Sandbox identifier (e.g. worktree path). """ self._plan_sandbox_refs[plan_id] = sandbox_ref + # Emit SANDBOX_CREATED event + if self._event_bus is not None: + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.SANDBOX_CREATED, + plan_id=plan_id, + details={ + "sandbox_ref": sandbox_ref, + }, + ) + ) + except Exception: + logger.warning( + "event_bus_emit_failed: SANDBOX_CREATED for plan %s", + plan_id, + ) def unregister_sandbox(self, plan_id: str) -> None: """Remove the sandbox reference for a plan. @@ -191,6 +225,26 @@ class CheckpointService: "reason": reason, }, ) + # Emit CHECKPOINT_CREATED event + if self._event_bus is not None: + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.CHECKPOINT_CREATED, + plan_id=plan_id, + details={ + "checkpoint_id": checkpoint.checkpoint_id, + "sandbox_ref": sandbox_ref, + "reason": reason, + "checkpoint_type": checkpoint_type, + }, + ) + ) + except Exception: + logger.warning( + "event_bus_emit_failed: CHECKPOINT_CREATED for plan %s", + plan_id, + ) # Auto-prune: use the supplied policy, or fall back to the default. effective_policy = ( @@ -269,20 +323,43 @@ class CheckpointService: }, ) - # Emit domain event for rollback + # Emit domain events for rollback if self._event_bus is not None: - self._event_bus.emit( - DomainEvent( - event_type=EventType.CHECKPOINT_RESTORED, - plan_id=plan_id, - details={ - "checkpoint_id": checkpoint_id, - "sandbox_ref": sandbox_ref, - "restored_files_count": restored_count, - "changed_paths": changed_paths, - }, + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.CHECKPOINT_RESTORED, + plan_id=plan_id, + details={ + "checkpoint_id": checkpoint_id, + "sandbox_ref": sandbox_ref, + "restored_files_count": restored_count, + "changed_paths": changed_paths, + }, + ) + ) + except Exception: + logger.warning( + "event_bus_emit_failed: CHECKPOINT_RESTORED for plan %s", + plan_id, + ) + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.SANDBOX_ROLLED_BACK, + plan_id=plan_id, + details={ + "checkpoint_id": checkpoint_id, + "sandbox_ref": sandbox_ref, + "restored_files_count": restored_count, + }, + ) + ) + except Exception: + logger.warning( + "event_bus_emit_failed: SANDBOX_ROLLED_BACK for plan %s", + plan_id, ) - ) return RollbackResult( restored_files_count=restored_count, diff --git a/src/cleveragents/application/services/context_service.py b/src/cleveragents/application/services/context_service.py index 612d27219..f8d95086b 100644 --- a/src/cleveragents/application/services/context_service.py +++ b/src/cleveragents/application/services/context_service.py @@ -26,6 +26,8 @@ from cleveragents.infrastructure.database.unit_of_work import ( UnitOfWork, UnitOfWorkContext, ) +from cleveragents.infrastructure.events.models import DomainEvent +from cleveragents.infrastructure.events.types import EventType DEFAULT_IGNORE_PATTERNS = [ "__pycache__", @@ -57,6 +59,7 @@ if TYPE_CHECKING: ContextAnalysisAgent, ContextAnalysisState, ) + from cleveragents.infrastructure.events.protocol import EventBus class VectorStoreProtocol(Protocol): @@ -90,6 +93,7 @@ class ContextService: unit_of_work: UnitOfWork, *, vector_store_service: VectorStoreProtocol | None = None, + event_bus: EventBus | None = None, ) -> None: """Initialize the context service. @@ -97,10 +101,12 @@ class ContextService: settings: Application settings unit_of_work: Unit of Work for database transactions vector_store_service: Optional semantic search helper + event_bus: Optional EventBus for domain event emission. """ self.settings = settings self.unit_of_work = unit_of_work self._vector_store_service = vector_store_service + self._event_bus = event_bus self.max_file_size = 10 * 1024 * 1024 # 10MB per file # Use getattr with default value self.max_context_size = getattr( @@ -924,12 +930,31 @@ class ContextService: return [] try: - return service.search( + results = service.search( plan.id, trimmed, top_k=max(1, limit), refresh_if_missing=refresh_if_missing, ) + # Emit CONTEXT_QUERY_EXECUTED event + if self._event_bus is not None: + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.CONTEXT_QUERY_EXECUTED, + details={ + "query": trimmed[:200], + "results_count": len(results), + "plan_id": plan.id, + }, + ) + ) + except Exception: + logger.warning( + "event_bus_emit_failed", + event_type="CONTEXT_QUERY_EXECUTED", + ) + return results except ConfigurationError as exc: logger.warning( "vector-store-search-failed", diff --git a/src/cleveragents/application/services/decision_service.py b/src/cleveragents/application/services/decision_service.py index 632ff8fe2..06a86212a 100644 --- a/src/cleveragents/application/services/decision_service.py +++ b/src/cleveragents/application/services/decision_service.py @@ -402,17 +402,66 @@ class DecisionService: else 0, ) if self.event_bus is not None: - self.event_bus.emit( - DomainEvent( - event_type=EventType.DECISION_CREATED, - plan_id=decision.plan_id, - details={ - "decision_id": decision.decision_id, - "decision_type": str(decision.decision_type), - "sequence_number": decision.sequence_number, - }, + try: + self.event_bus.emit( + DomainEvent( + event_type=EventType.DECISION_CREATED, + plan_id=decision.plan_id, + details={ + "decision_id": decision.decision_id, + "decision_type": str(decision.decision_type), + "sequence_number": decision.sequence_number, + }, + ) ) - ) + except Exception: + self._logger.warning( + "event_bus_emit_failed", + event_type="DECISION_CREATED", + plan_id=plan_id, + exc_info=True, + ) + # Emit DECISION_CORRECTED when this is a correction decision + if decision.is_correction: + try: + self.event_bus.emit( + DomainEvent( + event_type=EventType.DECISION_CORRECTED, + plan_id=decision.plan_id, + details={ + "decision_id": decision.decision_id, + "corrects_decision_id": decision.corrects_decision_id, + "correction_reason": decision.correction_reason, + }, + ) + ) + except Exception: + self._logger.warning( + "event_bus_emit_failed", + event_type="DECISION_CORRECTED", + plan_id=plan_id, + exc_info=True, + ) + else: + # Non-correction decisions are treated as approved + try: + self.event_bus.emit( + DomainEvent( + event_type=EventType.DECISION_APPROVED, + plan_id=decision.plan_id, + details={ + "decision_id": decision.decision_id, + "decision_type": str(decision.decision_type), + }, + ) + ) + except Exception: + self._logger.warning( + "event_bus_emit_failed", + event_type="DECISION_APPROVED", + plan_id=plan_id, + exc_info=True, + ) return decision # ------------------------------------------------------------------ @@ -627,6 +676,25 @@ class DecisionService: decision_id=decision_id, superseded_by=new_decision_id, ) + if self.event_bus is not None: + try: + self.event_bus.emit( + DomainEvent( + event_type=EventType.DECISION_SUPERSEDED, + plan_id=result.plan_id, + details={ + "decision_id": decision_id, + "superseded_by": new_decision_id, + }, + ) + ) + except Exception: + self._logger.warning( + "event_bus_emit_failed", + event_type="DECISION_SUPERSEDED", + decision_id=decision_id, + exc_info=True, + ) return result original = self._decisions.get(decision_id) @@ -641,6 +709,25 @@ class DecisionService: decision_id=decision_id, superseded_by=new_decision_id, ) + if self.event_bus is not None: + try: + self.event_bus.emit( + DomainEvent( + event_type=EventType.DECISION_SUPERSEDED, + plan_id=updated.plan_id, + details={ + "decision_id": decision_id, + "superseded_by": new_decision_id, + }, + ) + ) + except Exception: + self._logger.warning( + "event_bus_emit_failed", + event_type="DECISION_SUPERSEDED", + decision_id=decision_id, + exc_info=True, + ) return updated def delete_decision(self, decision_id: str) -> bool: diff --git a/src/cleveragents/application/services/invariant_service.py b/src/cleveragents/application/services/invariant_service.py index e877198a4..913e77276 100644 --- a/src/cleveragents/application/services/invariant_service.py +++ b/src/cleveragents/application/services/invariant_service.py @@ -19,6 +19,8 @@ Based on ``docs/specification.md`` and implementation plan Stage M3.5. from __future__ import annotations +from typing import TYPE_CHECKING + import structlog from ulid import ULID @@ -30,6 +32,11 @@ from cleveragents.domain.models.core.invariant import ( InvariantScope, merge_invariants, ) +from cleveragents.infrastructure.events.models import DomainEvent +from cleveragents.infrastructure.events.types import EventType + +if TYPE_CHECKING: + from cleveragents.infrastructure.events.protocol import EventBus logger = structlog.get_logger(__name__) @@ -41,12 +48,17 @@ class InvariantService: and enforcement record creation. All storage is in-memory. """ - def __init__(self) -> None: - """Initialise the invariant service with empty in-memory storage.""" + def __init__(self, event_bus: EventBus | None = None) -> None: + """Initialise the invariant service with empty in-memory storage. + + Args: + event_bus: Optional EventBus for domain event emission. + """ self._invariants: dict[str, Invariant] = {} self._enforcement_records: list[InvariantEnforcementRecord] = [] self._logger = logger.bind(service="invariant") self._sanitizer = PromptSanitizer() + self._event_bus = event_bus def add_invariant( self, @@ -192,6 +204,7 @@ class InvariantService: plan_id: str, invariants: list[Invariant], actor_response: str | None = None, + violated_invariant_ids: list[str] | None = None, ) -> list[InvariantEnforcementRecord]: """Create enforcement records for a set of invariants. @@ -202,6 +215,10 @@ class InvariantService: plan_id: The plan being checked. invariants: The invariants to enforce. actor_response: Optional response from the reconciliation actor. + violated_invariant_ids: Optional list of invariant IDs that + were violated. When provided, those invariants are + recorded with ``enforced=False`` and an + ``INVARIANT_VIOLATED`` event is emitted. Returns: List of ``InvariantEnforcementRecord`` objects. @@ -212,16 +229,38 @@ class InvariantService: if not plan_id or not plan_id.strip(): raise ValidationError("Plan ID must not be empty") + violated_ids = set(violated_invariant_ids or []) records: list[InvariantEnforcementRecord] = [] response = actor_response or "" for inv in invariants: + enforced = inv.id not in violated_ids record = InvariantEnforcementRecord( invariant_id=inv.id, - enforced=True, + enforced=enforced, actor_response=response, decision_id=str(ULID()), ) records.append(record) + if not enforced and self._event_bus is not None: + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.INVARIANT_VIOLATED, + plan_id=plan_id, + details={ + "invariant_id": inv.id, + "invariant_text": inv.text, + "scope": inv.scope.value, + }, + ) + ) + except Exception: + self._logger.warning( + "event_bus_emit_failed", + event_type="INVARIANT_VIOLATED", + plan_id=plan_id, + exc_info=True, + ) self._enforcement_records.extend(records) self._logger.info( @@ -229,4 +268,43 @@ class InvariantService: plan_id=plan_id, count=len(records), ) + if self._event_bus is not None: + for record in records: + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.INVARIANT_ENFORCED, + plan_id=plan_id, + details={ + "invariant_id": record.invariant_id, + "decision_id": record.decision_id, + "enforced": record.enforced, + }, + ) + ) + except Exception: + self._logger.warning( + "event_bus_emit_failed", + event_type="INVARIANT_ENFORCED", + plan_id=plan_id, + exc_info=True, + ) + # Emit a single INVARIANT_RECONCILED for the batch + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.INVARIANT_RECONCILED, + plan_id=plan_id, + details={ + "invariant_count": len(records), + }, + ) + ) + except Exception: + self._logger.warning( + "event_bus_emit_failed", + event_type="INVARIANT_RECONCILED", + plan_id=plan_id, + exc_info=True, + ) return records diff --git a/src/cleveragents/application/services/plan_lifecycle_service.py b/src/cleveragents/application/services/plan_lifecycle_service.py index 84572bc9b..a54947862 100644 --- a/src/cleveragents/application/services/plan_lifecycle_service.py +++ b/src/cleveragents/application/services/plan_lifecycle_service.py @@ -1168,6 +1168,27 @@ class PlanLifecycleService: self._commit_plan(plan) self._logger.info("Strategize completed", plan_id=plan_id) + if self.event_bus is not None: + try: + self.event_bus.emit( + DomainEvent( + event_type=EventType.PLAN_STATE_CHANGED, + plan_id=plan_id, + details={ + "phase": plan.phase.value, + "processing_state": plan.processing_state.value + if plan.processing_state + else None, + }, + ) + ) + except Exception: + self._logger.warning( + "event_bus_emit_failed", + event_type="PLAN_STATE_CHANGED", + plan_id=plan_id, + exc_info=True, + ) # Auto-progress if automation level permits return self.auto_progress(plan_id) @@ -1189,6 +1210,25 @@ class PlanLifecycleService: self._commit_plan(plan) self._logger.error("Strategize failed", plan_id=plan_id, error=error_message) + if self.event_bus is not None: + try: + self.event_bus.emit( + DomainEvent( + event_type=EventType.PLAN_ERRORED, + plan_id=plan_id, + details={ + "phase": plan.phase.value, + "error_message": error_message, + }, + ) + ) + except Exception: + self._logger.warning( + "event_bus_emit_failed", + event_type="PLAN_ERRORED", + plan_id=plan_id, + exc_info=True, + ) return plan @@ -1314,6 +1354,27 @@ class PlanLifecycleService: self._commit_plan(plan) self._logger.info("Execute completed", plan_id=plan_id) + if self.event_bus is not None: + try: + self.event_bus.emit( + DomainEvent( + event_type=EventType.PLAN_STATE_CHANGED, + plan_id=plan_id, + details={ + "phase": plan.phase.value, + "processing_state": plan.processing_state.value + if plan.processing_state + else None, + }, + ) + ) + except Exception: + self._logger.warning( + "event_bus_emit_failed", + event_type="PLAN_STATE_CHANGED", + plan_id=plan_id, + exc_info=True, + ) # Auto-progress if automation level permits return self.auto_progress(plan_id) @@ -1327,6 +1388,25 @@ class PlanLifecycleService: self._commit_plan(plan) self._logger.error("Execute failed", plan_id=plan_id, error=error_message) + if self.event_bus is not None: + try: + self.event_bus.emit( + DomainEvent( + event_type=EventType.PLAN_ERRORED, + plan_id=plan_id, + details={ + "phase": plan.phase.value, + "error_message": error_message, + }, + ) + ) + except Exception: + self._logger.warning( + "event_bus_emit_failed", + event_type="PLAN_ERRORED", + plan_id=plan_id, + exc_info=True, + ) # R7-F2 fix: clean up containers on execute failure (terminal state). self._cleanup_devcontainers(plan_id) @@ -1513,6 +1593,25 @@ class PlanLifecycleService: self._commit_plan(plan) self._logger.error("Apply failed", plan_id=plan_id, error=error_message) + if self.event_bus is not None: + try: + self.event_bus.emit( + DomainEvent( + event_type=EventType.PLAN_ERRORED, + plan_id=plan_id, + details={ + "phase": plan.phase.value, + "error_message": error_message, + }, + ) + ) + except Exception: + self._logger.warning( + "event_bus_emit_failed", + event_type="PLAN_ERRORED", + plan_id=plan_id, + exc_info=True, + ) # R7-F2 fix: clean up containers on apply failure (terminal state). self._cleanup_devcontainers(plan_id) diff --git a/src/cleveragents/application/services/repo_indexing_service.py b/src/cleveragents/application/services/repo_indexing_service.py index c67e22f03..082a76445 100644 --- a/src/cleveragents/application/services/repo_indexing_service.py +++ b/src/cleveragents/application/services/repo_indexing_service.py @@ -9,10 +9,16 @@ import threading from collections.abc import Callable from datetime import UTC, datetime from pathlib import Path -from typing import Any, cast +from typing import TYPE_CHECKING, Any, cast from ulid import ULID +from cleveragents.infrastructure.events.models import DomainEvent +from cleveragents.infrastructure.events.types import EventType + +if TYPE_CHECKING: + from cleveragents.infrastructure.events.protocol import EventBus + from cleveragents.application.services.repo_indexing_persistence import ( load_index, load_index_status, @@ -64,9 +70,15 @@ class RepoIndexingService: and serializes concurrent operations per resource within a process. """ - def __init__(self, session_factory: Any) -> None: - """Initialise with *session_factory* (callable → SQLAlchemy Session).""" + def __init__(self, session_factory: Any, event_bus: EventBus | None = None) -> None: + """Initialise with *session_factory* (callable → SQLAlchemy Session). + + Args: + session_factory: Callable that produces SQLAlchemy sessions. + event_bus: Optional EventBus for domain event emission. + """ self._session_factory = session_factory + self._event_bus = event_bus self._resource_locks: dict[str, threading.RLock] = {} self._locks_guard = threading.Lock() self.cleanup_stale_indexing() @@ -232,6 +244,24 @@ class RepoIndexingService: "primary_language": primary_language, }, ) + # Emit RESOURCE_INDEXED event + if self._event_bus is not None: + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.RESOURCE_INDEXED, + details={ + "resource_id": resource_id, + "file_count": file_count, + "token_estimate": token_estimate, + }, + ) + ) + except Exception: + logger.warning( + "event_bus_emit_failed: RESOURCE_INDEXED for %s", + resource_id, + ) return repo_index diff --git a/src/cleveragents/application/services/resource_handler_service.py b/src/cleveragents/application/services/resource_handler_service.py index fc80b32ea..78586297f 100644 --- a/src/cleveragents/application/services/resource_handler_service.py +++ b/src/cleveragents/application/services/resource_handler_service.py @@ -26,12 +26,14 @@ from __future__ import annotations import logging from collections.abc import Callable -from typing import cast +from typing import TYPE_CHECKING, cast from cleveragents.core.exceptions import NotFoundError from cleveragents.domain.models.core.resource import Resource from cleveragents.domain.models.core.resource_slot import BindingResult from cleveragents.domain.models.core.resource_type import ResourceTypeSpec +from cleveragents.infrastructure.events.models import DomainEvent +from cleveragents.infrastructure.events.types import EventType from cleveragents.infrastructure.sandbox.factory import SandboxStrategyStr from cleveragents.infrastructure.sandbox.manager import SandboxManager from cleveragents.resource.handlers.protocol import ( @@ -51,6 +53,9 @@ from cleveragents.resource.handlers.resolver import ( ) from cleveragents.tool.context import BoundResource +if TYPE_CHECKING: + from cleveragents.infrastructure.events.protocol import EventBus + logger = logging.getLogger(__name__) @@ -81,11 +86,13 @@ class ResourceHandlerService: sandbox_manager: SandboxManager, resource_lookup: Callable[[str], Resource], type_lookup: Callable[[str], ResourceTypeSpec], + event_bus: EventBus | None = None, ) -> None: self._sandbox_manager = sandbox_manager self._resource_lookup = resource_lookup self._type_lookup = type_lookup self._logger = logger + self._event_bus = event_bus def resolve_binding( self, @@ -130,13 +137,33 @@ class ResourceHandlerService: handler = self._resolve_handler_for_type(type_spec, resource) # Step 4: Delegate to handler - return handler.resolve( + bound = handler.resolve( resource=resource, plan_id=plan_id, slot_name=binding.slot_name, sandbox_manager=self._sandbox_manager, access=access, ) + # Emit RESOURCE_ACCESSED event + if self._event_bus is not None: + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.RESOURCE_ACCESSED, + plan_id=plan_id, + details={ + "resource_id": binding.resource_id, + "slot_name": binding.slot_name, + "access": access, + }, + ) + ) + except Exception: + self._logger.warning( + "event_bus_emit_failed: RESOURCE_ACCESSED for %s", + binding.resource_id, + ) + return bound def resolve_bindings( self, diff --git a/src/cleveragents/application/services/session_service.py b/src/cleveragents/application/services/session_service.py index 1bf4e7c54..ec34cee91 100644 --- a/src/cleveragents/application/services/session_service.py +++ b/src/cleveragents/application/services/session_service.py @@ -198,6 +198,23 @@ class PersistentSessionService(SessionService): session.updated_at = datetime.now() self._session_repo.update(session) + # Emit SESSION_MESSAGE_SENT event + if self._event_bus is not None: + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.SESSION_MESSAGE_SENT, + session_id=session_id, + details={ + "session_id": session_id, + "message_id": message.message_id, + "role": role.value if hasattr(role, "value") else str(role), + }, + ) + ) + except Exception: + _logger.warning("audit_emit_failed", event_type="SESSION_MESSAGE_SENT") + return message def export_session(self, session_id: str) -> dict[str, Any]: diff --git a/src/cleveragents/application/services/validation_pipeline.py b/src/cleveragents/application/services/validation_pipeline.py index 2ccb8b3c8..1e3d860a7 100644 --- a/src/cleveragents/application/services/validation_pipeline.py +++ b/src/cleveragents/application/services/validation_pipeline.py @@ -26,11 +26,16 @@ import time from collections import defaultdict from collections.abc import Callable from concurrent.futures import ThreadPoolExecutor -from typing import Any, Protocol +from typing import TYPE_CHECKING, Any, Protocol from pydantic import BaseModel, ConfigDict, Field from cleveragents.domain.models.core.tool import ValidationMode +from cleveragents.infrastructure.events.models import DomainEvent +from cleveragents.infrastructure.events.types import EventType + +if TYPE_CHECKING: + from cleveragents.infrastructure.events.protocol import EventBus logger = logging.getLogger(__name__) @@ -279,11 +284,13 @@ class ValidationPipeline: executor: Callable[[str, dict[str, Any]], dict[str, Any]], max_workers: int = 4, read_only_resources: set[str] | None = None, + event_bus: EventBus | None = None, ) -> None: self._commands = self._sort_commands(commands) self._executor = executor self._max_workers = max_workers self._read_only_resources: set[str] = read_only_resources or set() + self._event_bus = event_bus @staticmethod def _sort_commands( @@ -473,6 +480,20 @@ class ValidationPipeline: results=[], ) + # Emit VALIDATION_STARTED event + if self._event_bus is not None: + try: + self._event_bus.emit( + DomainEvent( + event_type=EventType.VALIDATION_STARTED, + details={ + "validation_count": len(self._commands), + }, + ) + ) + except Exception: + logger.warning("event_bus_emit_failed: VALIDATION_STARTED") + results: list[ValidationResult] = [] # Install thread-local stream wrappers so _execute_single can @@ -521,6 +542,28 @@ class ValidationPipeline: summary.required_failed, ) + # Emit VALIDATION_PASSED or VALIDATION_FAILED event + if self._event_bus is not None: + event_type = ( + EventType.VALIDATION_PASSED + if summary.all_required_passed + else EventType.VALIDATION_FAILED + ) + try: + self._event_bus.emit( + DomainEvent( + event_type=event_type, + details={ + "total": summary.total, + "required_passed": summary.required_passed, + "required_failed": summary.required_failed, + "all_required_passed": summary.all_required_passed, + }, + ) + ) + except Exception: + logger.warning("event_bus_emit_failed: %s", event_type.value) + return summary def run_for_plan( diff --git a/src/cleveragents/tool/actor_runtime.py b/src/cleveragents/tool/actor_runtime.py index 3a76f01eb..308f4def8 100644 --- a/src/cleveragents/tool/actor_runtime.py +++ b/src/cleveragents/tool/actor_runtime.py @@ -32,10 +32,12 @@ from __future__ import annotations import logging import time -from typing import Any, Protocol, runtime_checkable +from typing import TYPE_CHECKING, Any, Protocol, runtime_checkable from pydantic import BaseModel, ConfigDict, Field +from cleveragents.infrastructure.events.models import DomainEvent +from cleveragents.infrastructure.events.types import EventType from cleveragents.tool.actor_context import ToolActorContext, ToolCallRecord from cleveragents.tool.registry import ToolRegistry from cleveragents.tool.router import ( @@ -46,6 +48,9 @@ from cleveragents.tool.router import ( from cleveragents.tool.runner import ToolRunner from cleveragents.tool.runtime import ToolError +if TYPE_CHECKING: + from cleveragents.infrastructure.events.protocol import EventBus + logger = logging.getLogger(__name__) _DEFAULT_MAX_ITERATIONS = 25 @@ -208,6 +213,7 @@ class ToolCallingRuntime: provider_format: ProviderFormat = ProviderFormat.LANGCHAIN, plan_env: str | None = None, project_env: str | None = None, + event_bus: EventBus | None = None, ) -> None: if not isinstance(registry, ToolRegistry): raise TypeError("registry must be a ToolRegistry") @@ -224,6 +230,20 @@ class ToolCallingRuntime: self._provider_format = provider_format self._plan_env = plan_env self._project_env = project_env + self._event_bus = event_bus + + def _try_emit( + self, event_type: EventType, plan_id: str, details: dict[str, Any] + ) -> None: + """Best-effort event emission — never raises.""" + if self._event_bus is None: + return + try: + self._event_bus.emit( + DomainEvent(event_type=event_type, plan_id=plan_id, details=details) + ) + except Exception: + logger.warning("event_bus_emit_failed: %s", event_type.value) # -- Properties ----------------------------------------------------------- @@ -348,6 +368,18 @@ class ToolCallingRuntime: output = {} error = f"{type(exc).__name__}: {exc}" + # Emit ACTOR_ERRORED if tool call failed + if not success: + self._try_emit( + EventType.ACTOR_ERRORED, + actor_context.plan_id, + { + "tool_name": tool_call.name, + "error": error or "", + "iteration": iteration, + }, + ) + # Record metadata record = ToolCallRecord( tool_name=tool_call.name, @@ -411,6 +443,13 @@ class ToolCallingRuntime: if context is None: context = ToolActorContext(plan_id="default", phase="execute") + # Emit ACTOR_INVOKED event + self._try_emit( + EventType.ACTOR_INVOKED, + context.plan_id, + {"plan_id": context.plan_id, "phase": context.phase}, + ) + # Export tool schemas tool_schemas = self.export_tool_schemas() @@ -435,6 +474,16 @@ class ToolCallingRuntime: # No tool calls -> done if not llm_response.tool_calls: final_content = llm_response.content + # Emit ACTOR_COMPLETED event + self._try_emit( + EventType.ACTOR_COMPLETED, + context.plan_id, + { + "plan_id": context.plan_id, + "iterations": iteration, + "tool_calls_count": len(context.tool_call_history), + }, + ) return ToolCallRunResult( content=final_content, tool_call_history=context.tool_call_history, @@ -452,6 +501,16 @@ class ToolCallingRuntime: logger.warning( "Tool-call loop reached max iterations (%d)", self._max_iterations ) + # Emit ACTOR_ESCALATED when max iterations reached + self._try_emit( + EventType.ACTOR_ESCALATED, + context.plan_id, + { + "plan_id": context.plan_id, + "reason": "max_iterations_reached", + "max_iterations": self._max_iterations, + }, + ) return ToolCallRunResult( content=final_content, tool_call_history=context.tool_call_history, diff --git a/src/cleveragents/tool/lifecycle.py b/src/cleveragents/tool/lifecycle.py index 5c4b4a6b1..79100341a 100644 --- a/src/cleveragents/tool/lifecycle.py +++ b/src/cleveragents/tool/lifecycle.py @@ -72,11 +72,13 @@ import threading import time from collections import OrderedDict from datetime import UTC, datetime -from typing import Any, Protocol, runtime_checkable +from typing import TYPE_CHECKING, Any, Protocol, runtime_checkable from pydantic import BaseModel, ConfigDict, Field from cleveragents.domain.models.core.tool import Tool, ToolCapability +from cleveragents.infrastructure.events.models import DomainEvent +from cleveragents.infrastructure.events.types import EventType from cleveragents.tool.context import ( Change, ToolCancelledError, @@ -89,6 +91,9 @@ from cleveragents.tool.schema_validator import ( validate_tool_output, ) +if TYPE_CHECKING: + from cleveragents.infrastructure.events.protocol import EventBus + logger = logging.getLogger(__name__) # --------------------------------------------------------------------------- @@ -405,11 +410,28 @@ class ToolRuntime: *, tools: dict[str, Tool] | None = None, instances: dict[str, ToolInstance] | None = None, + event_bus: EventBus | None = None, ) -> None: self._tools: dict[str, Tool] = tools or {} self._instances: dict[str, ToolInstance] = instances or {} self._cache = ToolLifecycleCache() self._lock = threading.RLock() + self._event_bus = event_bus + + def _try_emit( + self, event_type: EventType, plan_id: str, details: dict[str, Any] + ) -> None: + """Best-effort event emission — never raises.""" + if self._event_bus is None: + return + try: + self._event_bus.emit( + DomainEvent(event_type=event_type, plan_id=plan_id, details=details) + ) + except Exception: + logger.warning( + "event_bus_emit_failed", extra={"event_type": event_type.value} + ) # -- Registration -------------------------------------------------------- @@ -534,6 +556,19 @@ class ToolRuntime: tool = self._get_tool(tool_name) + # Emit TOOL_INVOKED event + self._try_emit( + EventType.TOOL_INVOKED, + ctx.plan_id, + {"tool_name": tool_name, "plan_id": ctx.plan_id}, + ) + if ctx.step_retry_count > 0: + self._try_emit( + EventType.TOOL_RETRIED, + ctx.plan_id, + {"tool_name": tool_name, "retry_count": ctx.step_retry_count}, + ) + # 2. Auto-activate if needed (activate already enforces capabilities) self.activate(tool_name, ctx) @@ -589,6 +624,11 @@ class ToolRuntime: trace.success = False trace.error = "Cancelled" ctx.add_trace(trace) + self._try_emit( + EventType.TOOL_ERRORED, + ctx.plan_id, + {"tool_name": tool_name, "error": "Cancelled"}, + ) raise except Exception as exc: @@ -598,12 +638,28 @@ class ToolRuntime: trace.success = False trace.error = str(exc) ctx.add_trace(trace) + self._try_emit( + EventType.TOOL_ERRORED, + ctx.plan_id, + {"tool_name": tool_name, "error": str(exc)}, + ) raise ToolExecutionError( f"Tool '{tool_name}' execution failed: {exc}" ) from exc ctx.add_trace(trace) + # Emit TOOL_COMPLETED event on success + self._try_emit( + EventType.TOOL_COMPLETED, + ctx.plan_id, + { + "tool_name": tool_name, + "success": result.success, + "duration_ms": trace.duration_ms, + }, + ) + # 5. Validate outputs if tool.output_schema and result.data is not None: try: -- 2.52.0