feat(di): wire decision services #461

Merged
freemo merged 1 commits from feature/m4-decision-di into master 2026-02-27 03:22:22 +00:00
11 changed files with 1103 additions and 0 deletions
+214
View File
@@ -0,0 +1,214 @@
"""ASV benchmarks for DecisionService DI resolution and operations.
Measures:
- Container resolution of DecisionService
- Decision recording through the service layer
- Decision tree retrieval through the service layer
"""
from __future__ import annotations
import sys
from datetime import UTC, datetime
from pathlib import Path
from unittest.mock import MagicMock
try:
from cleveragents.application.services.decision_service import DecisionService
from cleveragents.domain.models.core.action import Action, ActionState
from cleveragents.domain.models.core.decision import Decision, DecisionType
from cleveragents.domain.models.core.plan import (
NamespacedName,
Plan,
PlanIdentity,
PlanPhase,
PlanTimestamps,
ProcessingState,
)
from cleveragents.infrastructure.database.models import Base
from cleveragents.infrastructure.database.repositories import (
ActionRepository,
DecisionRepository,
LifecyclePlanRepository,
)
from cleveragents.infrastructure.database.unit_of_work import UnitOfWork
except ModuleNotFoundError:
sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src"))
from cleveragents.application.services.decision_service import DecisionService
from cleveragents.domain.models.core.action import Action, ActionState
from cleveragents.domain.models.core.decision import Decision, DecisionType
from cleveragents.domain.models.core.plan import (
NamespacedName,
Plan,
PlanIdentity,
PlanPhase,
PlanTimestamps,
ProcessingState,
)
from cleveragents.infrastructure.database.models import Base
from cleveragents.infrastructure.database.repositories import (
ActionRepository,
DecisionRepository,
LifecyclePlanRepository,
)
from cleveragents.infrastructure.database.unit_of_work import UnitOfWork
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker
_PLAN_ID = "01HV00000000000000DIBENCH1"
def _make_uow() -> UnitOfWork:
"""Create a UoW backed by an in-memory SQLite database."""
engine = create_engine("sqlite:///:memory:", echo=False)
Base.metadata.create_all(engine)
uow = UnitOfWork.__new__(UnitOfWork)
uow.database_url = "sqlite:///:memory:"
uow._engine = engine
uow._session_factory = sessionmaker(
bind=engine,
expire_on_commit=False,
autoflush=False,
autocommit=False,
)
uow._database_initialized = True
uow._prompt_for_migration = None
return uow
def _seed_prerequisites(uow: UnitOfWork) -> None:
"""Create the action + plan needed to satisfy FK constraints."""
with uow.transaction() as ctx:
action_repo = ctx.actions
action = Action(
namespaced_name=NamespacedName.parse("local/bench-di-action"),
description="Benchmark action",
definition_of_done="Done",
strategy_actor="local/s",
execution_actor="local/e",
state=ActionState.AVAILABLE,
created_at=datetime(2026, 1, 1, tzinfo=UTC),
updated_at=datetime(2026, 1, 1, tzinfo=UTC),
)
action_repo.create(action)
with uow.transaction() as ctx:
plan_repo = ctx.lifecycle_plans
now = datetime(2026, 3, 1, tzinfo=UTC)
plan = Plan(
identity=PlanIdentity(plan_id=_PLAN_ID, attempt=1),
namespaced_name=NamespacedName(namespace="local", name="bench-di-plan"),
action_name="local/bench-di-action",
description="Benchmark DI plan",
definition_of_done="Done",
phase=PlanPhase.STRATEGIZE,
processing_state=ProcessingState.PROCESSING,
strategy_actor="local/s",
execution_actor="local/e",
timestamps=PlanTimestamps(created_at=now, updated_at=now),
created_by="bench",
tags=[],
reusable=True,
read_only=False,
)
plan_repo.create(plan)
def _make_decision(
seq: int = 0,
parent_id: str | None = None,
dtype: DecisionType = DecisionType.PROMPT_DEFINITION,
) -> Decision:
return Decision(
plan_id=_PLAN_ID,
parent_decision_id=parent_id,
sequence_number=seq,
decision_type=dtype,
question="Benchmark question?",
chosen_option="Benchmark option",
)
class TimeDecisionServiceResolution:
"""Benchmark DecisionService resolution from the DI container."""
timeout = 60
def setup(self) -> None:
import os
os.environ["CLEVERAGENTS_DATABASE_URL"] = "sqlite:///:memory:"
from cleveragents.application.container import get_container, reset_container
reset_container()
self.container = get_container()
def time_decision_service_resolution(self) -> None:
self.container.decision_service()
def teardown(self) -> None:
import os
from cleveragents.application.container import reset_container
reset_container()
os.environ.pop("CLEVERAGENTS_DATABASE_URL", None)
class TimeDecisionRecord:
"""Benchmark decision recording through DecisionService."""
timeout = 60
def setup(self) -> None:
self.uow = _make_uow()
_seed_prerequisites(self.uow)
self.svc = DecisionService(settings=MagicMock(), unit_of_work=self.uow)
self._seq = 100
def time_decision_record(self) -> None:
self._seq += 1
d = _make_decision(
seq=self._seq,
dtype=DecisionType.STRATEGY_CHOICE,
)
self.svc.record_decision(d)
class TimeDecisionTreeRetrieval:
"""Benchmark decision tree retrieval through DecisionService."""
timeout = 60
def setup(self) -> None:
self.uow = _make_uow()
_seed_prerequisites(self.uow)
self.svc = DecisionService(settings=MagicMock(), unit_of_work=self.uow)
# Build a small tree
root = _make_decision(seq=0)
self.svc.record_decision(root)
self.root_id = root.decision_id
seq = 1
for _ in range(3):
child = _make_decision(
seq=seq,
parent_id=root.decision_id,
dtype=DecisionType.STRATEGY_CHOICE,
)
self.svc.record_decision(child)
for _ in range(2):
seq += 1
gc = _make_decision(
seq=seq,
parent_id=child.decision_id,
dtype=DecisionType.IMPLEMENTATION_CHOICE,
)
self.svc.record_decision(gc)
seq += 1
def time_decision_tree_retrieval(self) -> None:
self.svc.get_decision_tree(self.root_id)
+108
View File
@@ -0,0 +1,108 @@
# Dependency Injection Reference
The CleverAgents DI container is built on `dependency-injector`'s
`DeclarativeContainer`. All service registrations live in
`src/cleveragents/application/container.py`.
## Container Overview
| Provider | Type | Description |
|-------------------------------|-----------|-----------------------------------------|
| `settings` | Singleton | Application settings |
| `provider_registry` | Singleton | AI provider registry |
| `database_url` | Callable | Database URL resolution |
| `unit_of_work` | Factory | Unit of Work (new per request) |
| `ai_provider` | Singleton | AI provider instance |
| `project_service` | Factory | Project management service |
| `vector_store_service` | Factory | Vector store operations |
| `context_service` | Factory | Context management service |
| `actor_service` | Factory | Actor configuration service |
| `actor_registry` | Factory | Actor registry with provider binding |
| `plan_service` | Factory | Legacy plan management (deprecated) |
| **`decision_service`** | Factory | Decision tree recording service |
| **`plan_lifecycle_service`** | Factory | V3 four-phase plan lifecycle |
| `resource_registry_service` | Factory | Resource registry operations |
| `namespaced_project_repo` | Factory | Namespaced project repository |
| `project_resource_link_repo` | Factory | Project-resource link repository |
| `stream_router` | Singleton | Reactive stream routing |
| `langgraph_bridge` | Singleton | LangGraph integration bridge |
| `route_bridge` | Singleton | Reactive route bridge |
## Decision Service Registration
The `DecisionService` is registered as a **Factory** provider, meaning
each resolution returns a new instance. It depends on:
- `settings` — Application settings (Singleton)
- `unit_of_work` — Unit of Work (Factory)
```python
decision_service = providers.Factory(
DecisionService,
settings=settings,
unit_of_work=unit_of_work,
)
```
`DecisionService` wraps `DecisionRepository` (accessible via
`UnitOfWorkContext.decisions`) and provides:
| Method | Description |
|-------------------------|--------------------------------------------|
| `record_decision()` | Persist a new decision |
| `get_decision()` | Retrieve a decision by ID |
| `get_decisions_for_plan()` | Get all decisions for a plan |
| `get_decision_tree()` | BFS traversal of the decision tree |
| `get_path_to_root()` | Walk from a decision up to the root |
| `mark_superseded()` | Mark a decision as superseded |
| `list_by_type()` | List decisions by type for a plan |
## Plan Lifecycle Service Registration
The `PlanLifecycleService` is registered as a **Factory** provider
with the `DecisionService` injected:
```python
plan_lifecycle_service = providers.Factory(
PlanLifecycleService,
settings=settings,
unit_of_work=unit_of_work,
decision_service=decision_service,
)
```
When a `DecisionService` is wired, the lifecycle service automatically
records decisions during phase transitions:
- **start_strategize** → records a `strategy_choice` decision
- **start_execute** → records an `implementation_choice` decision
If no `DecisionService` is provided (e.g. in legacy tests), decision
recording is silently skipped.
## Usage Example
```python
from cleveragents.application.container import get_container
container = get_container()
# Resolve services
decision_svc = container.decision_service()
lifecycle_svc = container.plan_lifecycle_service()
# Record a decision manually
from cleveragents.domain.models.core.decision import Decision, DecisionType
decision = Decision(
plan_id="01HV...",
sequence_number=0,
decision_type=DecisionType.STRATEGY_CHOICE,
question="Which approach?",
chosen_option="Approach A",
)
decision_svc.record_decision(decision)
# Retrieve decisions
decisions = decision_svc.get_decisions_for_plan("01HV...")
```
+36
View File
@@ -0,0 +1,36 @@
Feature: Decision DI Wiring
As a developer
I want DecisionService wired into the DI container
So that decision recording is automatic during plan lifecycle transitions
@phase1
Scenario: DI container resolves DecisionService
Given decdi- the application container is initialized
When decdi- I request a DecisionService from the container
Then decdi- the container should return a valid DecisionService instance
@phase1
Scenario: DecisionService records decision during plan strategize
Given decdi- a PlanLifecycleService with a DecisionService
And decdi- an action and plan exist in strategize phase
When decdi- I start the strategize phase
Then decdi- a strategy_choice decision should be recorded for the plan
@phase1
Scenario: DecisionService records decision during plan execute
Given decdi- a PlanLifecycleService with a DecisionService
And decdi- an action and plan exist in execute phase
When decdi- I start the execute phase
Then decdi- an implementation_choice decision should be recorded for the plan
@phase1
Scenario: DecisionRepository is accessible via DI
Given decdi- the application container is initialized
When decdi- I open a UnitOfWork transaction from the container
Then decdi- the transaction context should expose a DecisionRepository
@phase1
Scenario: DecisionService scoping is correct
Given decdi- the application container is initialized
When decdi- I request two DecisionService instances from the container
Then decdi- the instances should be different objects because it is a Factory
+244
View File
@@ -0,0 +1,244 @@
"""Step definitions for decision_di_wiring.feature.
Covers DI container wiring of DecisionService and its integration with
PlanLifecycleService for automatic decision recording.
All step text uses the ``decdi-`` prefix to avoid collisions with
existing step definitions.
"""
from __future__ import annotations
import os
from unittest.mock import MagicMock
from behave import given, then, when
from behave.runner import Context
# -------------------------------------------------------------------
# Helpers
# -------------------------------------------------------------------
def _ensure_test_db_env(context: Context) -> None:
"""Set up a test database URL if not already set."""
if not hasattr(context, "_decdi_env_saved"):
context._decdi_env_saved = os.environ.get("CLEVERAGENTS_DATABASE_URL")
os.environ["CLEVERAGENTS_DATABASE_URL"] = "sqlite:///:memory:"
def _restore_env(context: Context) -> None:
"""Restore environment variables."""
saved = getattr(context, "_decdi_env_saved", None)
if saved is None:
os.environ.pop("CLEVERAGENTS_DATABASE_URL", None)
else:
os.environ["CLEVERAGENTS_DATABASE_URL"] = saved
def _cleanup_container(context: Context) -> None:
"""Reset the global container after the scenario."""
from cleveragents.application.container import reset_container
reset_container()
_restore_env(context)
# -------------------------------------------------------------------
# Scenario: DI container resolves DecisionService
# -------------------------------------------------------------------
@given("decdi- the application container is initialized")
def step_init_container(context: Context) -> None:
from cleveragents.application.container import get_container, reset_container
_ensure_test_db_env(context)
reset_container()
context.container = get_container()
context.add_cleanup(_cleanup_container, context)
@when("decdi- I request a DecisionService from the container")
def step_request_decision_service(context: Context) -> None:
context.decision_service = context.container.decision_service()
@then("decdi- the container should return a valid DecisionService instance")
def step_check_decision_service(context: Context) -> None:
from cleveragents.application.services.decision_service import DecisionService
assert isinstance(context.decision_service, DecisionService), (
f"Expected DecisionService, got {type(context.decision_service).__name__}"
)
# -------------------------------------------------------------------
# Scenario: DecisionService records decision during plan strategize
# -------------------------------------------------------------------
@given("decdi- a PlanLifecycleService with a DecisionService")
def step_create_lifecycle_with_decision(context: Context) -> None:
"""Build an in-memory PlanLifecycleService wired with DecisionService."""
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker
from cleveragents.application.services.decision_service import DecisionService
from cleveragents.application.services.plan_lifecycle_service import (
PlanLifecycleService,
)
from cleveragents.infrastructure.database.models import Base
from cleveragents.infrastructure.database.unit_of_work import UnitOfWork
# Build a fresh in-memory DB
engine = create_engine("sqlite:///:memory:", echo=False)
Base.metadata.create_all(engine)
# Create a UoW that skips migration
uow = UnitOfWork.__new__(UnitOfWork)
uow.database_url = "sqlite:///:memory:"
uow._engine = engine
uow._session_factory = sessionmaker(
bind=engine,
expire_on_commit=False,
autoflush=False,
autocommit=False,
)
uow._database_initialized = True
uow._prompt_for_migration = None
mock_settings = MagicMock()
mock_settings.database_url = "sqlite:///:memory:"
decision_svc = DecisionService(settings=mock_settings, unit_of_work=uow)
lifecycle_svc = PlanLifecycleService(
settings=mock_settings,
unit_of_work=uow,
decision_service=decision_svc,
)
context.decision_svc = decision_svc
context.lifecycle_svc = lifecycle_svc
context.uow = uow
@given("decdi- an action and plan exist in strategize phase")
def step_create_action_and_plan_strategize(context: Context) -> None:
context.action = context.lifecycle_svc.create_action(
name="local/decdi-strategize-test",
description="Test action for strategize",
definition_of_done="Done",
strategy_actor="local/s",
execution_actor="local/e",
)
context.plan = context.lifecycle_svc.use_action("local/decdi-strategize-test")
context.plan_id = context.plan.identity.plan_id
@when("decdi- I start the strategize phase")
def step_start_strategize(context: Context) -> None:
context.lifecycle_svc.start_strategize(context.plan_id)
@then("decdi- a strategy_choice decision should be recorded for the plan")
def step_check_strategize_decision(context: Context) -> None:
decisions = context.decision_svc.get_decisions_for_plan(context.plan_id)
strategy_decisions = [
d for d in decisions if str(d.decision_type) == "strategy_choice"
]
assert len(strategy_decisions) >= 1, (
f"Expected at least 1 strategy_choice decision, got {len(strategy_decisions)}"
)
# -------------------------------------------------------------------
# Scenario: DecisionService records decision during plan execute
# -------------------------------------------------------------------
@given("decdi- an action and plan exist in execute phase")
def step_create_action_and_plan_execute(context: Context) -> None:
context.lifecycle_svc.create_action(
name="local/decdi-execute-test",
description="Test action for execute",
definition_of_done="Done",
strategy_actor="local/s",
execution_actor="local/e",
)
plan = context.lifecycle_svc.use_action("local/decdi-execute-test")
plan_id = plan.identity.plan_id
# Progress to execute phase: start -> complete strategize -> execute_plan
context.lifecycle_svc.start_strategize(plan_id)
context.lifecycle_svc.complete_strategize(plan_id)
# After complete_strategize, auto_progress may transition to execute
# depending on the automation profile. Fetch the latest state.
plan = context.lifecycle_svc.get_plan(plan_id)
from cleveragents.domain.models.core.plan import PlanPhase
if plan.phase == PlanPhase.STRATEGIZE:
context.lifecycle_svc.execute_plan(plan_id)
context.plan = context.lifecycle_svc.get_plan(plan_id)
context.plan_id = plan_id
@when("decdi- I start the execute phase")
def step_start_execute(context: Context) -> None:
context.lifecycle_svc.start_execute(context.plan_id)
@then("decdi- an implementation_choice decision should be recorded for the plan")
def step_check_execute_decision(context: Context) -> None:
decisions = context.decision_svc.get_decisions_for_plan(context.plan_id)
impl_decisions = [
d for d in decisions if str(d.decision_type) == "implementation_choice"
]
assert len(impl_decisions) >= 1, (
f"Expected at least 1 implementation_choice decision, got {len(impl_decisions)}"
)
# -------------------------------------------------------------------
# Scenario: DecisionRepository is accessible via DI
# -------------------------------------------------------------------
@when("decdi- I open a UnitOfWork transaction from the container")
def step_open_uow_transaction(context: Context) -> None:
uow = context.container.unit_of_work()
uow.init_database()
with uow.transaction() as ctx:
context.uow_ctx = ctx
context.has_decisions = hasattr(ctx, "decisions")
context.decisions_type_name = type(ctx.decisions).__name__
@then("decdi- the transaction context should expose a DecisionRepository")
def step_check_decision_repo(context: Context) -> None:
assert context.has_decisions, (
"UnitOfWorkContext does not have 'decisions' attribute"
)
assert context.decisions_type_name == "DecisionRepository", (
f"Expected DecisionRepository, got {context.decisions_type_name}"
)
# -------------------------------------------------------------------
# Scenario: DecisionService scoping is correct
# -------------------------------------------------------------------
@when("decdi- I request two DecisionService instances from the container")
def step_request_two_instances(context: Context) -> None:
context.svc_a = context.container.decision_service()
context.svc_b = context.container.decision_service()
@then("decdi- the instances should be different objects because it is a Factory")
def step_check_different_instances(context: Context) -> None:
assert context.svc_a is not context.svc_b, (
"Factory provider should return different instances per call"
)
+23
View File
@@ -0,0 +1,23 @@
*** Settings ***
Documentation Smoke tests for Decision DI wiring
Resource ${CURDIR}/common.resource
Suite Setup Setup Test Environment
Suite Teardown Cleanup Test Environment
*** Variables ***
${HELPER_SCRIPT} robot/helper_decision_di.py
*** Test Cases ***
Verify Decision DI Resolution
[Documentation] Verify DecisionService can be resolved from the DI container
[Tags] di decision smoke
${result}= Run Process ${PYTHON} ${HELPER_SCRIPT} resolve-service cwd=${WORKSPACE} timeout=30s
Should Be Equal As Integers ${result.rc} 0
Should Contain ${result.stdout} resolve-service-ok
Verify Decision Recording Integration
[Documentation] Verify DecisionService records decisions during lifecycle transitions
[Tags] di decision lifecycle smoke
${result}= Run Process ${PYTHON} ${HELPER_SCRIPT} record-integration cwd=${WORKSPACE} timeout=30s
Should Be Equal As Integers ${result.rc} 0
Should Contain ${result.stdout} record-integration-ok
+163
View File
@@ -0,0 +1,163 @@
"""Helper script for Robot Framework decision DI wiring smoke tests.
Usage:
python robot/helper_decision_di.py <subcommand>
Subcommands:
resolve-service Verify DecisionService resolves from DI container
record-integration Verify decision recording during lifecycle transitions
"""
from __future__ import annotations
import os
import sys
from pathlib import Path
from unittest.mock import MagicMock
# Ensure src is importable when run from workspace root
sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src"))
from sqlalchemy import create_engine # noqa: I001
from sqlalchemy.orm import sessionmaker
from cleveragents.application.services.decision_service import DecisionService
from cleveragents.application.services.plan_lifecycle_service import (
PlanLifecycleService,
)
from cleveragents.infrastructure.database.models import Base
from cleveragents.infrastructure.database.unit_of_work import UnitOfWork
# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------
def _make_uow() -> UnitOfWork:
"""Create a UoW backed by an in-memory SQLite database."""
engine = create_engine("sqlite:///:memory:", echo=False)
Base.metadata.create_all(engine)
uow = UnitOfWork.__new__(UnitOfWork)
uow.database_url = "sqlite:///:memory:"
uow._engine = engine
uow._session_factory = sessionmaker(
bind=engine,
expire_on_commit=False,
autoflush=False,
autocommit=False,
)
uow._database_initialized = True
uow._prompt_for_migration = None
return uow
# ---------------------------------------------------------------------------
# Subcommands
# ---------------------------------------------------------------------------
def _resolve_service() -> None:
"""Verify DecisionService can be instantiated from the container."""
os.environ["CLEVERAGENTS_DATABASE_URL"] = "sqlite:///:memory:"
try:
from cleveragents.application.container import get_container, reset_container
reset_container()
container = get_container()
svc = container.decision_service()
assert isinstance(svc, DecisionService), (
f"Expected DecisionService, got {type(svc).__name__}"
)
# Verify it's a Factory (different instances per call)
svc2 = container.decision_service()
assert svc is not svc2, "Factory should return different instances"
reset_container()
finally:
os.environ.pop("CLEVERAGENTS_DATABASE_URL", None)
print("resolve-service-ok")
def _record_integration() -> None:
"""Verify decision recording during plan lifecycle transitions."""
uow = _make_uow()
mock_settings = MagicMock()
mock_settings.database_url = "sqlite:///:memory:"
decision_svc = DecisionService(settings=mock_settings, unit_of_work=uow)
lifecycle_svc = PlanLifecycleService(
settings=mock_settings,
unit_of_work=uow,
decision_service=decision_svc,
)
# Create action + plan
lifecycle_svc.create_action(
name="local/robot-di-test",
description="Robot DI integration test",
definition_of_done="Pass",
strategy_actor="local/s",
execution_actor="local/e",
)
plan = lifecycle_svc.use_action("local/robot-di-test")
plan_id = plan.identity.plan_id
# Start strategize -> should record a strategy_choice decision
lifecycle_svc.start_strategize(plan_id)
decisions = decision_svc.get_decisions_for_plan(plan_id)
strategy_decisions = [
d for d in decisions if str(d.decision_type) == "strategy_choice"
]
assert len(strategy_decisions) >= 1, (
f"Expected strategy_choice decision, got {len(strategy_decisions)}"
)
# Complete strategize and progress to execute
lifecycle_svc.complete_strategize(plan_id)
plan = lifecycle_svc.get_plan(plan_id)
from cleveragents.domain.models.core.plan import PlanPhase
if plan.phase == PlanPhase.STRATEGIZE:
lifecycle_svc.execute_plan(plan_id)
# Start execute -> should record an implementation_choice decision
lifecycle_svc.start_execute(plan_id)
decisions = decision_svc.get_decisions_for_plan(plan_id)
impl_decisions = [
d for d in decisions if str(d.decision_type) == "implementation_choice"
]
assert len(impl_decisions) >= 1, (
f"Expected implementation_choice decision, got {len(impl_decisions)}"
)
print("record-integration-ok")
# ---------------------------------------------------------------------------
# Dispatcher
# ---------------------------------------------------------------------------
_COMMANDS = {
"resolve-service": _resolve_service,
"record-integration": _record_integration,
}
def main() -> None:
if len(sys.argv) < 2:
raise SystemExit(f"Usage: {sys.argv[0]} <{'|'.join(_COMMANDS)}>")
command = sys.argv[1]
handler = _COMMANDS.get(command)
if handler is None:
raise SystemExit(f"Unknown command: {command}")
handler()
if __name__ == "__main__":
main()
+19
View File
@@ -12,6 +12,10 @@ from cleveragents.actor.registry import ActorRegistry
from cleveragents.application.reactive_registry_adapter import register_registry_agents
from cleveragents.application.services.actor_service import ActorService
from cleveragents.application.services.context_service import ContextService
from cleveragents.application.services.decision_service import DecisionService
from cleveragents.application.services.plan_lifecycle_service import (
PlanLifecycleService,
)
from cleveragents.application.services.plan_service import PlanService
from cleveragents.application.services.project_service import ProjectService
from cleveragents.application.services.resource_registry_service import (
@@ -207,6 +211,21 @@ class Container(containers.DeclarativeContainer):
actor_service=actor_service,
)
# Decision Service - Factory (wraps DecisionRepository via UoW)
decision_service = providers.Factory(
DecisionService,
settings=settings,
unit_of_work=unit_of_work,
)
# Plan Lifecycle Service - Factory (v3 four-phase lifecycle)
plan_lifecycle_service = providers.Factory(
PlanLifecycleService,
settings=settings,
unit_of_work=unit_of_work,
decision_service=decision_service,
)
# Resource Registry Service - uses database session factory from UoW
resource_registry_service = providers.Factory(
_build_resource_registry_service,
@@ -12,6 +12,9 @@ from cleveragents.application.services.config_service import (
from cleveragents.application.services.correction_service import (
CorrectionService,
)
from cleveragents.application.services.decision_service import (
DecisionService,
)
from cleveragents.application.services.invariant_service import (
InvariantService,
)
@@ -54,6 +57,7 @@ __all__ = [
"ConfigLevel",
"ConfigService",
"CorrectionService",
"DecisionService",
"DefaultValidationRunner",
"InvariantService",
"PersistentSessionService",
@@ -0,0 +1,210 @@
"""Application-layer service for decision tree operations.
``DecisionService`` wraps the :class:`DecisionRepository` behind a
thin application-layer façade, adding structured logging, transaction
management via :class:`UnitOfWork`, and a consistent API surface for
callers (e.g. ``PlanLifecycleService``).
Based on:
- ADR-007 (Repository Pattern / Unit of Work)
- ADR-033 (Decision Recording Protocol)
- Forgejo issue #173 (Wire decision services into DI)
"""
from __future__ import annotations
from typing import TYPE_CHECKING
import structlog
from cleveragents.domain.models.core.decision import Decision
if TYPE_CHECKING:
from cleveragents.infrastructure.database.unit_of_work import UnitOfWork
logger = structlog.get_logger(__name__)
class DecisionService:
"""Coordinate decision persistence through the Unit of Work.
All public methods open a transaction, delegate to
:class:`DecisionRepository`, and commit. Callers do not need to
manage sessions directly.
Args:
settings: Application settings (unused currently but kept for
consistency with the service constructor convention).
unit_of_work: A :class:`UnitOfWork` instance used to open
transactional scopes.
"""
def __init__(
self,
settings: object,
unit_of_work: UnitOfWork,
) -> None:
self.settings = settings
self.unit_of_work = unit_of_work
self._logger = logger.bind(service="decision")
# ------------------------------------------------------------------
# record_decision
# ------------------------------------------------------------------
def record_decision(self, decision: Decision) -> Decision:
"""Persist a new decision.
Opens a UnitOfWork transaction, creates the decision via the
repository, and commits.
Args:
decision: The :class:`Decision` domain object to persist.
Returns:
The same :class:`Decision` instance (pass-through).
Raises:
DuplicateDecisionError: If a decision with the same ID exists.
DatabaseError: On transient or unexpected DB errors.
"""
self._logger.info(
"recording_decision",
decision_id=decision.decision_id,
plan_id=decision.plan_id,
decision_type=str(decision.decision_type),
sequence_number=decision.sequence_number,
)
with self.unit_of_work.transaction() as ctx:
ctx.decisions.create(decision)
self._logger.debug(
"decision_recorded",
decision_id=decision.decision_id,
)
return decision
# ------------------------------------------------------------------
# get_decision
# ------------------------------------------------------------------
def get_decision(self, decision_id: str) -> Decision | None:
"""Retrieve a single decision by its ULID.
Args:
decision_id: ULID of the decision.
Returns:
The :class:`Decision` or ``None`` if not found.
"""
self._logger.debug("getting_decision", decision_id=decision_id)
with self.unit_of_work.transaction() as ctx:
return ctx.decisions.get(decision_id)
# ------------------------------------------------------------------
# get_decisions_for_plan
# ------------------------------------------------------------------
def get_decisions_for_plan(self, plan_id: str) -> list[Decision]:
"""Return all decisions for a plan ordered by sequence number.
Args:
plan_id: ULID of the plan.
Returns:
List of :class:`Decision` instances.
"""
self._logger.debug("getting_decisions_for_plan", plan_id=plan_id)
with self.unit_of_work.transaction() as ctx:
return ctx.decisions.get_by_plan(plan_id)
# ------------------------------------------------------------------
# get_decision_tree
# ------------------------------------------------------------------
def get_decision_tree(self, root_id: str) -> list[Decision]:
"""Retrieve the full decision tree rooted at *root_id* (BFS order).
Args:
root_id: ULID of the root decision.
Returns:
List of :class:`Decision` instances (BFS order).
Raises:
DecisionNotFoundError: If the root decision does not exist.
"""
self._logger.debug("getting_decision_tree", root_id=root_id)
with self.unit_of_work.transaction() as ctx:
return ctx.decisions.get_tree(root_id)
# ------------------------------------------------------------------
# get_path_to_root
# ------------------------------------------------------------------
def get_path_to_root(self, decision_id: str) -> list[Decision]:
"""Walk from *decision_id* up to the root (leaf → root order).
Args:
decision_id: ULID of the starting decision.
Returns:
List of :class:`Decision` instances (leaf first).
Raises:
DecisionNotFoundError: If the starting decision is missing.
"""
self._logger.debug("getting_path_to_root", decision_id=decision_id)
with self.unit_of_work.transaction() as ctx:
return ctx.decisions.get_path_to_root(decision_id)
# ------------------------------------------------------------------
# mark_superseded
# ------------------------------------------------------------------
def mark_superseded(
self,
decision_id: str,
new_decision_id: str,
) -> Decision:
"""Mark a decision as superseded by a new one.
Args:
decision_id: ULID of the decision to mark.
new_decision_id: ULID of the replacement decision.
Returns:
Updated :class:`Decision`.
Raises:
DecisionNotFoundError: If *decision_id* does not exist.
"""
self._logger.info(
"marking_superseded",
decision_id=decision_id,
new_decision_id=new_decision_id,
)
with self.unit_of_work.transaction() as ctx:
return ctx.decisions.update_superseded_by(decision_id, new_decision_id)
# ------------------------------------------------------------------
# list_by_type
# ------------------------------------------------------------------
def list_by_type(self, plan_id: str, decision_type: str) -> list[Decision]:
"""List decisions of a given type for a plan.
Args:
plan_id: ULID of the plan.
decision_type: String value of a :class:`DecisionType` enum
member (e.g. ``"strategy_choice"``).
Returns:
List of :class:`Decision` instances.
"""
self._logger.debug(
"listing_by_type",
plan_id=plan_id,
decision_type=decision_type,
)
with self.unit_of_work.transaction() as ctx:
return ctx.decisions.list_by_type(plan_id, decision_type)
@@ -82,6 +82,7 @@ from cleveragents.domain.models.core.plan import (
)
if TYPE_CHECKING:
from cleveragents.application.services.decision_service import DecisionService
from cleveragents.config.settings import Settings
from cleveragents.infrastructure.database.unit_of_work import UnitOfWork
@@ -147,6 +148,7 @@ class PlanLifecycleService:
self,
settings: Settings,
unit_of_work: UnitOfWork | None = None,
decision_service: DecisionService | None = None,
):
"""Initialize the plan lifecycle service.
@@ -157,20 +159,71 @@ class PlanLifecycleService:
``ActionRepository`` and ``LifecyclePlanRepository``.
When ``None``, the service falls back to in-memory
storage for backward compatibility with existing tests.
decision_service: Optional :class:`DecisionService` for
automatically recording decisions during phase
transitions. When ``None``, decision recording is
silently skipped.
"""
self.settings = settings
self.unit_of_work = unit_of_work
self.decision_service = decision_service
self._logger = logger.bind(service="plan_lifecycle")
# In-memory fallback storage (used only when no UoW is provided)
self._actions: dict[str, Action] = {}
self._plans: dict[str, Plan] = {}
# Decision sequence counter per plan (monotonic within a plan)
self._decision_seq: dict[str, int] = {}
@property
def _persisted(self) -> bool:
"""Return True when a UnitOfWork is wired for persistence."""
return self.unit_of_work is not None
def _next_seq(self, plan_id: str) -> int:
"""Return the next decision sequence number for *plan_id*."""
seq = self._decision_seq.get(plan_id, -1) + 1
self._decision_seq[plan_id] = seq
return seq
def _try_record_decision(
self,
plan_id: str,
decision_type: str,
question: str,
chosen_option: str,
parent_decision_id: str | None = None,
) -> None:
"""Record a decision if DecisionService is available.
Failures are logged but never propagated — decision recording
must not block lifecycle transitions.
"""
if self.decision_service is None:
return
from cleveragents.domain.models.core.decision import Decision as DecisionModel
from cleveragents.domain.models.core.decision import DecisionType as DT
try:
decision = DecisionModel(
plan_id=plan_id,
parent_decision_id=parent_decision_id,
sequence_number=self._next_seq(plan_id),
decision_type=DT(decision_type),
question=question,
chosen_option=chosen_option,
)
self.decision_service.record_decision(decision)
except Exception:
self._logger.warning(
"decision_recording_failed",
plan_id=plan_id,
decision_type=decision_type,
exc_info=True,
)
def _persist_action_create(self, action: Action, ctx: Any) -> None:
"""Persist a new action via the repository.
@@ -687,6 +740,13 @@ class PlanLifecycleService:
self._commit_plan(plan)
self._logger.info("Strategize started", plan_id=plan_id)
self._try_record_decision(
plan_id=plan_id,
decision_type="strategy_choice",
question="Which strategy should the plan follow?",
chosen_option=f"Begin strategize phase for plan {plan_id}",
)
return plan
def complete_strategize(self, plan_id: str) -> Plan:
@@ -808,6 +868,13 @@ class PlanLifecycleService:
self._commit_plan(plan)
self._logger.info("Execute started", plan_id=plan_id)
self._try_record_decision(
plan_id=plan_id,
decision_type="implementation_choice",
question="How should the plan be executed?",
chosen_option=f"Begin execute phase for plan {plan_id}",
)
return plan
def complete_execute(self, plan_id: str) -> Plan:
+15
View File
@@ -323,3 +323,18 @@ list_models # noqa: B018, F821
check_plan_budget # noqa: B018, F821
check_daily_budget # noqa: B018, F821
get_daily_spend # noqa: B018, F821
# Decision DI wiring — public API (M4, #173)
DecisionService # noqa: B018, F821
decision_service # noqa: B018, F821
plan_lifecycle_service # noqa: B018, F821
record_decision # noqa: B018, F821
get_decision # noqa: B018, F821
get_decisions_for_plan # noqa: B018, F821
get_decision_tree # noqa: B018, F821
get_path_to_root # noqa: B018, F821
mark_superseded # noqa: B018, F821
list_by_type # noqa: B018, F821
_try_record_decision # noqa: B018, F821
_next_seq # noqa: B018, F821
_decision_seq # noqa: B018, F821