diff --git a/benchmarks/subplan_spawn_bench.py b/benchmarks/subplan_spawn_bench.py new file mode 100644 index 000000000..d45dccb41 --- /dev/null +++ b/benchmarks/subplan_spawn_bench.py @@ -0,0 +1,148 @@ +"""Airspeed Velocity benchmarks for subplan spawn service throughput. + +Measures spawn creation, validation, and entry building overhead for +SubplanService across varying numbers of spawn entries. +""" + +from __future__ import annotations + +from unittest.mock import MagicMock + +from cleveragents.application.services.subplan_service import ( + SpawnEntry, + SubplanService, +) +from cleveragents.domain.models.core.decision import ( + ContextSnapshot, + Decision, + DecisionType, + ResourceRef, +) +from cleveragents.domain.models.core.plan import ( + ExecutionMode, + NamespacedName, + Plan, + PlanIdentity, + SubplanConfig, +) + +_PLAN_ID = "01HGZ6FE0AQDYTR4BXVQZ6PN00" + + +def _mock_decision_service() -> MagicMock: + svc: MagicMock = MagicMock() + svc.list_by_type = MagicMock(return_value=[]) + return svc + + +def _make_decision( + decision_id: str, + sequence: int = 0, + resource_ids: list[str] | None = None, +) -> Decision: + resources: list[ResourceRef] = [] + if resource_ids: + resources = [ResourceRef(resource_id=rid) for rid in resource_ids] + return Decision( + decision_id=decision_id, + plan_id=_PLAN_ID, + decision_type=DecisionType.SUBPLAN_SPAWN, + sequence_number=sequence, + question="Spawn child plan?", + chosen_option="local/sub-action", + context_snapshot=ContextSnapshot(relevant_resources=resources), + ) + + +def _make_plan() -> Plan: + return Plan( + identity=PlanIdentity(plan_id=_PLAN_ID), + namespaced_name=NamespacedName(namespace="local", name="bench-plan"), + description="Benchmark plan", + action_name="local/bench-action", + ) + + +class SubplanSpawnThroughputSuite: + """Benchmark SubplanService spawn throughput.""" + + def setup(self) -> None: + """Prepare fixtures for spawn benchmarks.""" + self.service = SubplanService( + decision_service=_mock_decision_service(), + ) + self.plan = _make_plan() + self.seq_config = SubplanConfig( + execution_mode=ExecutionMode.SEQUENTIAL, + ) + self.par_config = SubplanConfig( + execution_mode=ExecutionMode.PARALLEL, + max_parallel=50, + ) + # Build entries for varying sizes + self.entries_1 = self._build_entries(1) + self.entries_5 = self._build_entries(5) + self.entries_20 = self._build_entries(20) + + def _build_entries(self, count: int) -> list[SpawnEntry]: + entries: list[SpawnEntry] = [] + for i in range(count): + did = f"01HGZ6FE0AQDYTR4BXVQZ6D{i:02d}0" + dec = _make_decision(did, sequence=i) + entries.append(SpawnEntry(decision=dec, action_name=f"local/sub-{i}")) + return entries + + def time_spawn_single_sequential(self) -> None: + """Time spawning a single child plan in sequential mode.""" + self.service.spawn( + parent_plan=self.plan, + config=self.seq_config, + spawn_entries=self.entries_1, + ) + + def time_spawn_5_sequential(self) -> None: + """Time spawning 5 child plans in sequential mode.""" + self.service.spawn( + parent_plan=self.plan, + config=self.seq_config, + spawn_entries=self.entries_5, + ) + + def time_spawn_20_parallel(self) -> None: + """Time spawning 20 child plans in parallel mode.""" + self.service.spawn( + parent_plan=self.plan, + config=self.par_config, + spawn_entries=self.entries_20, + ) + + def time_validate_5_entries(self) -> None: + """Time validating 5 spawn entries.""" + self.service.validate_spawn( + config=self.seq_config, + spawn_entries=self.entries_5, + ) + + def time_validate_20_entries_with_resources(self) -> None: + """Time validating 20 entries with resource scope check.""" + self.service.validate_spawn( + config=self.par_config, + spawn_entries=self.entries_20, + available_resources=frozenset(f"res-{i}" for i in range(100)), + ) + + def time_build_entries_from_decisions(self) -> None: + """Time building spawn entries from 5 decisions.""" + decisions = [ + _make_decision( + f"01HGZ6FE0AQDYTR4BXVQZ6D{i:02d}0", + sequence=i, + resource_ids=[f"01HGZ6FE0AQDYTR4BXVQZ6R{i:02d}0"], + ) + for i in range(5) + ] + self.service.build_spawn_entries(decisions) + + def time_service_construction(self) -> None: + """Time constructing SubplanService.""" + SubplanService(decision_service=_mock_decision_service()) diff --git a/docs/reference/subplan_service.md b/docs/reference/subplan_service.md new file mode 100644 index 000000000..034d48887 --- /dev/null +++ b/docs/reference/subplan_service.md @@ -0,0 +1,140 @@ +# Subplan Service: Spawn Workflow and Lifecycle + +The `SubplanService` coordinates the spawning of child plans from +`DecisionService` spawn entries and `SubplanConfig`. It validates +resource scopes, merge strategies, and parallelism bounds before +creating child plan statuses. + +## Overview + +When a parent plan decides to decompose work into child plans (via +`subplan_spawn` or `subplan_parallel_spawn` decisions during Strategize), +the `SubplanService` handles the spawn workflow: + +1. Extract spawn decisions from the `DecisionService` +2. Build `SpawnEntry` objects from those decisions +3. Validate the spawn request +4. Create `SubplanStatus` and `SpawnMetadata` for each entry +5. Return the result for the caller to persist on the parent plan + +## Spawn Workflow + +``` +DecisionService SubplanService Parent Plan + | | | + | get_spawn_decisions | | + |<-----------------------| | + | [Decision, ...] | | + |----------------------->| | + | | build_spawn_entries | + | | validate_spawn | + | | spawn | + | |------------------------->| + | | SpawnResult | + | | (statuses + metadata) | +``` + +## Key Types + +### SpawnMetadata + +Persisted alongside each child plan for status output and provenance: + +| Field | Type | Description | +|---------------------|--------|------------------------------------------| +| `spawn_decision_id` | `str` | ULID of the decision that triggered spawn| +| `parent_plan_id` | `str` | ULID of the parent plan | +| `root_plan_id` | `str` | ULID of the root plan in the hierarchy | +| `execution_mode` | `str` | How the subplan should be executed | + +### SpawnEntry + +A single spawn request derived from a decision: + +| Field | Type | Description | +|--------------------|--------------|--------------------------------------| +| `decision` | `Decision` | The decision that triggered spawn | +| `action_name` | `str` | Namespaced action name for child plan| +| `target_resources` | `list[str]` | Resource IDs the child plan uses | +| `description` | `str` | Description/prompt for the child plan| + +### SpawnResult + +Result of spawning child plans: + +| Field | Type | Description | +|--------------------|-------------------------------|--------------------------------| +| `spawned_statuses` | `list[SubplanStatus]` | Status for each child plan | +| `metadata` | `dict[str, SpawnMetadata]` | Metadata keyed by subplan_id | +| `total_spawned` | `int` | Number of child plans created | +| `execution_mode` | `str` | Execution mode from config | + +## Spawn Validation + +Before any child plan is created, `validate_spawn` checks: + +1. **Resource scopes resolved**: All `target_resources` in each entry + exist in the `available_resources` set. Unresolved resources produce + a validation error. + +2. **Merge strategy defined**: The `SubplanConfig.merge_strategy` must + not be `None`. Without a merge strategy, child plan results cannot + be combined. + +3. **max_parallel bounds**: In `PARALLEL` execution mode, the number of + spawn entries must not exceed `config.max_parallel`. + +4. **Valid action names**: Each entry must have a non-empty `action_name`. + +5. **Correct decision types**: Each entry's decision must be either + `subplan_spawn` or `subplan_parallel_spawn`. + +If any check fails, `SpawnValidationError` is raised with all errors. + +## Integration with Other Services + +- **DecisionService**: Provides spawn-type decisions via + `list_by_type(plan_id, "subplan_spawn")`. +- **SubplanExecutionService**: Executes the spawned child plans using + the `SubplanStatus` objects from `SpawnResult`. +- **SubplanMergeService**: Merges child plan outputs using the strategy + from `SubplanConfig`. +- **PlanLifecycleService**: Manages the parent plan's lifecycle, + persisting the updated `subplan_statuses` and `subplan_config`. + +## Example Usage + +```python +from cleveragents.application.services.subplan_service import ( + SubplanService, +) + +# Get spawn decisions for the parent plan +decisions = subplan_service.get_spawn_decisions(parent_plan.identity.plan_id) + +# Build spawn entries from decisions +entries = subplan_service.build_spawn_entries(decisions) + +# Spawn child plans (validates first) +result = subplan_service.spawn( + parent_plan=parent_plan, + config=parent_plan.subplan_config, + spawn_entries=entries, + available_resources=frozenset(known_resource_ids), +) + +# Update the parent plan with spawned statuses +parent_plan.subplan_statuses = result.spawned_statuses +``` + +## Error Handling + +| Error | When | +|--------------------------|---------------------------------------------| +| `ValueError` | `None` or empty required arguments | +| `SpawnValidationError` | Validation checks fail (resource scope, | +| | merge strategy, max_parallel, action name, | +| | or decision type) | + +`SpawnValidationError` inherits from `ValidationError` and includes a +`validation_errors` list with detailed messages for each failure. diff --git a/features/steps/subplan_spawn_service_steps.py b/features/steps/subplan_spawn_service_steps.py new file mode 100644 index 000000000..456f69d02 --- /dev/null +++ b/features/steps/subplan_spawn_service_steps.py @@ -0,0 +1,587 @@ +"""Step definitions for subplan spawn service scenarios.""" + +from __future__ import annotations + +from unittest.mock import MagicMock + +from behave import given, then, when +from behave.runner import Context + +from cleveragents.application.services.subplan_service import ( + SpawnEntry, + SpawnValidationError, + SubplanService, +) +from cleveragents.domain.models.core.decision import ( + ContextSnapshot, + Decision, + DecisionType, + ResourceRef, +) +from cleveragents.domain.models.core.plan import ( + ExecutionMode, + NamespacedName, + Plan, + PlanIdentity, + SubplanConfig, +) + +_PLAN_ID = "01HGZ6FE0AQDYTR4BXVQZ6PN00" +_ROOT_ID = "01HGZ6FE0AQDYTR4BXVQZ6RF00" +_DEC_ID1 = "01HGZ6FE0AQDYTR4BXVQZ6DA00" +_DEC_ID2 = "01HGZ6FE0AQDYTR4BXVQZ6DB00" +_DEC_ID3 = "01HGZ6FE0AQDYTR4BXVQZ6DC00" + + +def _mock_decision_service() -> MagicMock: + """Create a mock DecisionService.""" + svc: MagicMock = MagicMock() + svc.list_by_type = MagicMock(return_value=[]) + return svc + + +def _make_decision( + decision_id: str, + plan_id: str = _PLAN_ID, + decision_type: DecisionType = DecisionType.SUBPLAN_SPAWN, + sequence: int = 0, + chosen_option: str = "local/sub-action", + resource_ids: list[str] | None = None, +) -> Decision: + """Create a Decision with sensible defaults.""" + resources: list[ResourceRef] = [] + if resource_ids: + resources = [ResourceRef(resource_id=rid) for rid in resource_ids] + return Decision( + decision_id=decision_id, + plan_id=plan_id, + decision_type=decision_type, + sequence_number=sequence, + question="Should we spawn a child plan?", + chosen_option=chosen_option, + context_snapshot=ContextSnapshot(relevant_resources=resources), + ) + + +def _make_plan( + plan_id: str = _PLAN_ID, + parent_plan_id: str | None = None, + root_plan_id: str | None = None, +) -> Plan: + """Create a minimal Plan for testing.""" + return Plan( + identity=PlanIdentity( + plan_id=plan_id, + parent_plan_id=parent_plan_id, + root_plan_id=root_plan_id, + ), + namespaced_name=NamespacedName(namespace="local", name="test-plan"), + description="A test plan for subplan spawn", + action_name="local/test-action", + ) + + +# ---- Successful spawn ---- + + +@given("a parent plan with a subplan spawn decision") +def step_given_parent_with_spawn_decision(context: Context) -> None: + context.decision_service = _mock_decision_service() + context.subplan_service = SubplanService( + decision_service=context.decision_service, + ) + context.parent_plan = _make_plan() + dec: Decision = _make_decision(_DEC_ID1) + context.spawn_entries = [ + SpawnEntry( + decision=dec, + action_name="local/sub-action", + description="Test subplan", + ) + ] + + +@given("a subplan config with sequential execution mode") +def step_given_config_sequential(context: Context) -> None: + context.subplan_config = SubplanConfig( + execution_mode=ExecutionMode.SEQUENTIAL, + ) + + +@when("the subplan service spawns child plans") +def step_when_spawn(context: Context) -> None: + context.spawn_error = None + try: + context.spawn_result = context.subplan_service.spawn( + parent_plan=context.parent_plan, + config=context.subplan_config, + spawn_entries=context.spawn_entries, + available_resources=getattr(context, "available_resources", None), + ) + except (SpawnValidationError, ValueError) as exc: + context.spawn_error = exc + context.spawn_result = None + + +@then("{count:d} child plan status should be created") +def step_then_count_statuses_singular(context: Context, count: int) -> None: + assert context.spawn_result is not None, f"Spawn failed: {context.spawn_error}" + actual: int = len(context.spawn_result.spawned_statuses) + assert actual == count, f"Expected {count} statuses, got {actual}" + + +@then("{count:d} child plan statuses should be created") +def step_then_count_statuses_plural(context: Context, count: int) -> None: + assert context.spawn_result is not None, f"Spawn failed: {context.spawn_error}" + actual: int = len(context.spawn_result.spawned_statuses) + assert actual == count, f"Expected {count} statuses, got {actual}" + + +@then("the spawn metadata should contain the parent plan id") +def step_then_metadata_parent_id(context: Context) -> None: + assert context.spawn_result is not None + for meta in context.spawn_result.metadata.values(): + assert meta.parent_plan_id == context.parent_plan.identity.plan_id + + +@then("the spawn metadata should contain the root plan id") +def step_then_metadata_root_id(context: Context) -> None: + assert context.spawn_result is not None + plan: Plan = context.parent_plan + expected_root: str = plan.identity.root_plan_id or plan.identity.plan_id + for meta in context.spawn_result.metadata.values(): + assert meta.root_plan_id == expected_root + + +@then("the spawn metadata should contain the execution mode") +def step_then_metadata_exec_mode(context: Context) -> None: + assert context.spawn_result is not None + for meta in context.spawn_result.metadata.values(): + assert meta.execution_mode in {m.value for m in ExecutionMode} + + +# ---- Multiple entries ---- + + +@given("a parent plan with {count:d} subplan spawn decisions") +def step_given_parent_with_n_decisions(context: Context, count: int) -> None: + context.decision_service = _mock_decision_service() + context.subplan_service = SubplanService( + decision_service=context.decision_service, + ) + context.parent_plan = _make_plan() + context.spawn_entries = [] + for i in range(count): + # Build unique ULID-like IDs + suffix: str = f"{i:02d}" + did: str = f"01HGZ6FE0AQDYTR4BXVQZ6D{suffix}0" + dec: Decision = _make_decision(did, sequence=i) + context.spawn_entries.append( + SpawnEntry( + decision=dec, + action_name=f"local/sub-action-{i}", + description=f"Subplan {i}", + ) + ) + + +@given("a subplan config with parallel execution mode and max_parallel {n:d}") +def step_given_config_parallel(context: Context, n: int) -> None: + context.subplan_config = SubplanConfig( + execution_mode=ExecutionMode.PARALLEL, + max_parallel=n, + ) + + +@then("each status should have a unique subplan id") +def step_then_unique_ids(context: Context) -> None: + assert context.spawn_result is not None + ids: list[str] = [s.subplan_id for s in context.spawn_result.spawned_statuses] + assert len(ids) == len(set(ids)), f"Duplicate IDs found: {ids}" + + +@then("the spawn result total_spawned should be {count:d}") +def step_then_total_spawned(context: Context, count: int) -> None: + assert context.spawn_result is not None + assert context.spawn_result.total_spawned == count + + +# ---- Validation: resource scopes ---- + + +@given('a parent plan with a spawn entry targeting resource "{res_id}"') +def step_given_entry_with_resource(context: Context, res_id: str) -> None: + context.decision_service = _mock_decision_service() + context.subplan_service = SubplanService( + decision_service=context.decision_service, + ) + context.parent_plan = _make_plan() + dec: Decision = _make_decision( + _DEC_ID1, + resource_ids=[res_id], + ) + context.spawn_entries = [ + SpawnEntry( + decision=dec, + action_name="local/sub-action", + target_resources=[res_id], + ) + ] + + +@given('available resources are "{res_a}" and "{res_b}"') +def step_given_available_resources( + context: Context, + res_a: str, + res_b: str, +) -> None: + context.available_resources = frozenset({res_a, res_b}) + + +@when("the subplan service validates the spawn request") +def step_when_validate(context: Context) -> None: + context.validation_result = context.subplan_service.validate_spawn( + config=context.subplan_config, + spawn_entries=context.spawn_entries, + available_resources=getattr(context, "available_resources", None), + ) + + +@then("the validation result should be invalid") +def step_then_invalid(context: Context) -> None: + assert not context.validation_result.valid, ( + "Expected invalid, got valid with no errors" + ) + + +@then("the validation errors should mention unresolved resource") +def step_then_error_mentions_resource(context: Context) -> None: + errors_str: str = " ".join(context.validation_result.errors) + assert "resource" in errors_str.lower(), ( + f"Expected 'resource' in errors: {context.validation_result.errors}" + ) + + +# ---- Validation: merge strategy ---- + + +@given("a parent plan with a valid spawn entry") +def step_given_valid_entry(context: Context) -> None: + context.decision_service = _mock_decision_service() + context.subplan_service = SubplanService( + decision_service=context.decision_service, + ) + context.parent_plan = _make_plan() + dec: Decision = _make_decision(_DEC_ID1) + context.spawn_entries = [ + SpawnEntry( + decision=dec, + action_name="local/sub-action", + ) + ] + + +@given("a subplan config with no merge strategy") +def step_given_config_no_merge(context: Context) -> None: + # Create a config and then forcibly set merge_strategy to None + config: SubplanConfig = SubplanConfig( + execution_mode=ExecutionMode.SEQUENTIAL, + ) + # Pydantic models allow assignment when validate_assignment is True + # We use model_copy to create a version with None merge_strategy + # Since SubplanMergeStrategy is required, we use object.__setattr__ + object.__setattr__(config, "merge_strategy", None) + context.subplan_config = config + + +@then("the validation errors should mention merge strategy") +def step_then_error_mentions_merge(context: Context) -> None: + errors_str: str = " ".join(context.validation_result.errors) + assert "merge strategy" in errors_str.lower(), ( + f"Expected 'merge strategy' in errors: {context.validation_result.errors}" + ) + + +# ---- Validation: max_parallel ---- + + +@then("the validation errors should mention max_parallel") +def step_then_error_mentions_max_parallel(context: Context) -> None: + errors_str: str = " ".join(context.validation_result.errors) + assert "max_parallel" in errors_str.lower(), ( + f"Expected 'max_parallel' in errors: {context.validation_result.errors}" + ) + + +# ---- Metadata persistence ---- + + +@given('a parent plan that is a subplan with root plan id "{root_id}"') +def step_given_subplan_parent(context: Context, root_id: str) -> None: + context.decision_service = _mock_decision_service() + context.subplan_service = SubplanService( + decision_service=context.decision_service, + ) + context.parent_plan = _make_plan( + parent_plan_id="01HGZ6FE0AQDYTR4BXVQZ6PP00", + root_plan_id=root_id, + ) + + +@given("a valid spawn decision") +def step_given_valid_spawn_decision(context: Context) -> None: + dec: Decision = _make_decision(_DEC_ID1) + context.spawn_entries = [ + SpawnEntry( + decision=dec, + action_name="local/sub-action", + ) + ] + + +@then('the spawn metadata root_plan_id should be "{expected}"') +def step_then_metadata_root_equals(context: Context, expected: str) -> None: + assert context.spawn_result is not None + for meta in context.spawn_result.metadata.values(): + assert meta.root_plan_id == expected, ( + f"Expected root '{expected}', got '{meta.root_plan_id}'" + ) + + +@then('the spawn metadata execution_mode should be "{expected}"') +def step_then_metadata_mode_equals(context: Context, expected: str) -> None: + assert context.spawn_result is not None + for meta in context.spawn_result.metadata.values(): + assert meta.execution_mode == expected, ( + f"Expected mode '{expected}', got '{meta.execution_mode}'" + ) + + +# ---- Decision type validation ---- + + +@given("a parent plan with a non-spawn decision type") +def step_given_non_spawn_decision(context: Context) -> None: + context.decision_service = _mock_decision_service() + context.subplan_service = SubplanService( + decision_service=context.decision_service, + ) + context.parent_plan = _make_plan() + dec: Decision = _make_decision( + _DEC_ID1, + decision_type=DecisionType.STRATEGY_CHOICE, + ) + context.spawn_entries = [ + SpawnEntry( + decision=dec, + action_name="local/sub-action", + ) + ] + + +@then("the validation errors should mention decision type") +def step_then_error_mentions_decision_type(context: Context) -> None: + errors_str: str = " ".join(context.validation_result.errors) + assert "type" in errors_str.lower(), ( + f"Expected 'type' in errors: {context.validation_result.errors}" + ) + + +# ---- Get spawn decisions ---- + + +@given('a decision service with spawn decisions for plan "{plan_id}"') +def step_given_decision_service_with_spawns( + context: Context, + plan_id: str, +) -> None: + spawn_dec: Decision = _make_decision( + _DEC_ID1, + plan_id=plan_id, + sequence=0, + ) + parallel_dec: Decision = _make_decision( + _DEC_ID2, + plan_id=plan_id, + decision_type=DecisionType.SUBPLAN_PARALLEL_SPAWN, + sequence=1, + ) + + mock_svc: MagicMock = _mock_decision_service() + mock_svc.list_by_type = MagicMock( + side_effect=lambda pid, dtype: ( + [spawn_dec] if dtype == DecisionType.SUBPLAN_SPAWN.value else [parallel_dec] + ) + ) + context.decision_service = mock_svc + context.subplan_service = SubplanService( + decision_service=mock_svc, + ) + context.plan_id = plan_id + + +@when("the subplan service retrieves spawn decisions") +def step_when_get_spawn_decisions(context: Context) -> None: + context.spawn_decisions = context.subplan_service.get_spawn_decisions( + context.plan_id, + ) + + +@then("the returned decisions should only contain spawn types") +def step_then_only_spawn_types(context: Context) -> None: + for dec in context.spawn_decisions: + assert dec.decision_type in ( + DecisionType.SUBPLAN_SPAWN, + DecisionType.SUBPLAN_PARALLEL_SPAWN, + ), f"Unexpected decision type: {dec.decision_type}" + + +# ---- Build spawn entries ---- + + +@given("a list of spawn decisions with resource refs") +def step_given_decisions_with_resources(context: Context) -> None: + context.decision_service = _mock_decision_service() + context.subplan_service = SubplanService( + decision_service=context.decision_service, + ) + context.input_decisions = [ + _make_decision( + _DEC_ID1, + resource_ids=["01HGZ6FE0AQDYTR4BXVQZ6RA00"], + chosen_option="local/test-sub", + ), + _make_decision( + _DEC_ID2, + sequence=1, + resource_ids=[ + "01HGZ6FE0AQDYTR4BXVQZ6RB00", + "01HGZ6FE0AQDYTR4BXVQZ6RC00", + ], + chosen_option="local/test-sub-2", + ), + ] + + +@when("the subplan service builds spawn entries") +def step_when_build_entries(context: Context) -> None: + context.built_entries = context.subplan_service.build_spawn_entries( + context.input_decisions, + ) + + +@then("each entry should have action_name and target_resources") +def step_then_entries_have_fields(context: Context) -> None: + for entry in context.built_entries: + assert entry.action_name, "Entry missing action_name" + # First entry should have 1 resource, second should have 2 + assert len(context.built_entries) == 2 + assert len(context.built_entries[0].target_resources) == 1 + assert len(context.built_entries[1].target_resources) == 2 + + +# ---- Service construction ---- + + +@when("the subplan service is constructed with None decision_service") +def step_when_construct_none(context: Context) -> None: + context.construct_error = None + try: + SubplanService(decision_service=None) # type: ignore[arg-type] + except ValueError as exc: + context.construct_error = exc + + +@then("a ValueError should be raised for subplan service") +def step_then_value_error(context: Context) -> None: + error: Exception | None = getattr( + context, + "construct_error", + None, + ) or getattr(context, "spawn_error", None) + assert error is not None, "Expected ValueError but none was raised" + assert isinstance(error, ValueError), ( + f"Expected ValueError, got {type(error).__name__}" + ) + + +# ---- Null guard scenarios ---- + + +@given("a valid subplan service") +def step_given_valid_service(context: Context) -> None: + context.decision_service = _mock_decision_service() + context.subplan_service = SubplanService( + decision_service=context.decision_service, + ) + context.subplan_config = SubplanConfig() + + +@given("a parent plan with identity") +def step_given_plan_with_identity(context: Context) -> None: + context.parent_plan = _make_plan() + + +@when("spawn is called with None parent_plan") +def step_when_spawn_none_plan(context: Context) -> None: + context.spawn_error = None + try: + context.subplan_service.spawn( + parent_plan=None, # type: ignore[arg-type] + config=context.subplan_config, + spawn_entries=[ + SpawnEntry( + decision=_make_decision(_DEC_ID1), + action_name="local/test", + ) + ], + ) + except ValueError as exc: + context.spawn_error = exc + + +@when("spawn is called with empty spawn_entries") +def step_when_spawn_empty_entries(context: Context) -> None: + context.spawn_error = None + try: + context.subplan_service.spawn( + parent_plan=context.parent_plan, + config=context.subplan_config, + spawn_entries=[], + ) + except ValueError as exc: + context.spawn_error = exc + + +@then('a spawn ValueError should be raised with message "{msg}"') +def step_then_value_error_msg(context: Context, msg: str) -> None: + error: Exception | None = getattr( + context, + "construct_error", + None, + ) or getattr(context, "spawn_error", None) + assert error is not None, f"Expected ValueError with '{msg}'" + assert isinstance(error, ValueError) + assert msg in str(error), f"Expected '{msg}' in '{error}'" + + +# ---- Parallel spawn decision ---- + + +@given("a parent plan with a subplan_parallel_spawn decision") +def step_given_parallel_spawn_decision(context: Context) -> None: + context.decision_service = _mock_decision_service() + context.subplan_service = SubplanService( + decision_service=context.decision_service, + ) + context.parent_plan = _make_plan() + dec: Decision = _make_decision( + _DEC_ID1, + decision_type=DecisionType.SUBPLAN_PARALLEL_SPAWN, + ) + context.spawn_entries = [ + SpawnEntry( + decision=dec, + action_name="local/parallel-sub", + ) + ] diff --git a/features/subplan_spawn_service.feature b/features/subplan_spawn_service.feature new file mode 100644 index 000000000..c2c46981a --- /dev/null +++ b/features/subplan_spawn_service.feature @@ -0,0 +1,129 @@ +@phase1 @subplan @spawn +Feature: Subplan Spawn Service + As a system orchestrating hierarchical plans + I want to spawn child plans from decision entries + And validate spawn requests before execution + So that parent plans can safely decompose work into coordinated child plans + + # --- Successful spawn --- + + @spawn_success + Scenario: Successfully spawn child plan from decision entry + Given a parent plan with a subplan spawn decision + And a subplan config with sequential execution mode + When the subplan service spawns child plans + Then 1 child plan status should be created + And the spawn metadata should contain the parent plan id + And the spawn metadata should contain the root plan id + And the spawn metadata should contain the execution mode + + # --- Multiple entries --- + + @spawn_multiple + Scenario: Spawn from multiple decision entries + Given a parent plan with 3 subplan spawn decisions + And a subplan config with parallel execution mode and max_parallel 5 + When the subplan service spawns child plans + Then 3 child plan statuses should be created + And each status should have a unique subplan id + And the spawn result total_spawned should be 3 + + # --- Validation: resource scopes --- + + @validation @resource_scope + Scenario: Validate spawn with missing resource scope + Given a parent plan with a spawn entry targeting resource "res-missing" + And available resources are "res-a" and "res-b" + And a subplan config with sequential execution mode + When the subplan service validates the spawn request + Then the validation result should be invalid + And the validation errors should mention unresolved resource + + # --- Validation: merge strategy --- + + @validation @merge_strategy + Scenario: Validate spawn with undefined merge strategy + Given a parent plan with a valid spawn entry + And a subplan config with no merge strategy + When the subplan service validates the spawn request + Then the validation result should be invalid + And the validation errors should mention merge strategy + + # --- Validation: max_parallel --- + + @validation @max_parallel + Scenario: Validate spawn exceeding max_parallel bound + Given a parent plan with 5 subplan spawn decisions + And a subplan config with parallel execution mode and max_parallel 3 + When the subplan service validates the spawn request + Then the validation result should be invalid + And the validation errors should mention max_parallel + + # --- Metadata persistence --- + + @metadata + Scenario: Spawn metadata persisted for status output + Given a parent plan that is a subplan with root plan id "01HGZ6FE0AQDYTR4BXVQZ6RF00" + And a subplan config with sequential execution mode + And a valid spawn decision + When the subplan service spawns child plans + Then the spawn metadata root_plan_id should be "01HGZ6FE0AQDYTR4BXVQZ6RF00" + And the spawn metadata execution_mode should be "sequential" + + # --- Decision type validation --- + + @validation @decision_type + Scenario: Validate spawn with wrong decision type + Given a parent plan with a non-spawn decision type + And a subplan config with sequential execution mode + When the subplan service validates the spawn request + Then the validation result should be invalid + And the validation errors should mention decision type + + # --- Get spawn decisions --- + + @spawn_decisions + Scenario: Get spawn decisions from decision service + Given a decision service with spawn decisions for plan "01HGZ6FE0AQDYTR4BXVQZ6PN00" + When the subplan service retrieves spawn decisions + Then the returned decisions should only contain spawn types + + # --- Build spawn entries --- + + @build_entries + Scenario: Build spawn entries from decisions + Given a list of spawn decisions with resource refs + When the subplan service builds spawn entries + Then each entry should have action_name and target_resources + + # --- Service construction --- + + @construction + Scenario: SubplanService rejects None decision_service + When the subplan service is constructed with None decision_service + Then a ValueError should be raised for subplan service + + # --- Spawn rejects None inputs --- + + @validation @null_guard + Scenario: Spawn rejects None parent_plan + Given a valid subplan service + When spawn is called with None parent_plan + Then a spawn ValueError should be raised with message "parent_plan must not be None" + + @validation @null_guard + Scenario: Spawn rejects empty spawn_entries + Given a valid subplan service + And a parent plan with identity + When spawn is called with empty spawn_entries + Then a spawn ValueError should be raised with message "spawn_entries must not be empty" + + # --- Parallel spawn decision --- + + @spawn_parallel + Scenario: Spawn from parallel spawn decision type + Given a parent plan with a subplan_parallel_spawn decision + And a subplan config with parallel execution mode and max_parallel 5 + When the subplan service spawns child plans + Then 1 child plan status should be created + And the spawn metadata execution_mode should be "parallel" diff --git a/robot/helper_subplan_spawn.py b/robot/helper_subplan_spawn.py new file mode 100644 index 000000000..68cb6c211 --- /dev/null +++ b/robot/helper_subplan_spawn.py @@ -0,0 +1,251 @@ +"""Helper script for subplan spawn Robot Framework smoke tests. + +Exercises SubplanService spawn workflow, validation, and metadata +persistence without requiring the full service layer. +""" + +from __future__ import annotations + +import sys +from unittest.mock import MagicMock + +from cleveragents.application.services.subplan_service import ( + SpawnEntry, + SubplanService, +) +from cleveragents.domain.models.core.decision import ( + ContextSnapshot, + Decision, + DecisionType, + ResourceRef, +) +from cleveragents.domain.models.core.plan import ( + ExecutionMode, + NamespacedName, + Plan, + PlanIdentity, + SubplanConfig, +) + +_PLAN_ID = "01HGZ6FE0AQDYTR4BXVQZ6PN00" +_ROOT_ID = "01HGZ6FE0AQDYTR4BXVQZ6RF00" +_DEC_ID1 = "01HGZ6FE0AQDYTR4BXVQZ6DA00" +_DEC_ID2 = "01HGZ6FE0AQDYTR4BXVQZ6DB00" +_DEC_ID3 = "01HGZ6FE0AQDYTR4BXVQZ6DC00" + + +def _mock_decision_service() -> MagicMock: + svc: MagicMock = MagicMock() + svc.list_by_type = MagicMock(return_value=[]) + return svc + + +def _make_decision( + decision_id: str, + plan_id: str = _PLAN_ID, + decision_type: DecisionType = DecisionType.SUBPLAN_SPAWN, + sequence: int = 0, + chosen_option: str = "local/sub-action", + resource_ids: list[str] | None = None, +) -> Decision: + resources: list[ResourceRef] = [] + if resource_ids: + resources = [ResourceRef(resource_id=rid) for rid in resource_ids] + return Decision( + decision_id=decision_id, + plan_id=plan_id, + decision_type=decision_type, + sequence_number=sequence, + question="Spawn child plan?", + chosen_option=chosen_option, + context_snapshot=ContextSnapshot(relevant_resources=resources), + ) + + +def _make_plan( + plan_id: str = _PLAN_ID, + parent_plan_id: str | None = None, + root_plan_id: str | None = None, +) -> Plan: + return Plan( + identity=PlanIdentity( + plan_id=plan_id, + parent_plan_id=parent_plan_id, + root_plan_id=root_plan_id, + ), + namespaced_name=NamespacedName(namespace="local", name="test-plan"), + description="Robot test plan", + action_name="local/test-action", + ) + + +def _spawn_single() -> None: + """Spawn a single child plan from a decision entry.""" + svc = SubplanService(decision_service=_mock_decision_service()) + plan = _make_plan() + config = SubplanConfig(execution_mode=ExecutionMode.SEQUENTIAL) + dec = _make_decision(_DEC_ID1) + entries = [SpawnEntry(decision=dec, action_name="local/sub-action")] + result = svc.spawn(parent_plan=plan, config=config, spawn_entries=entries) + assert result.total_spawned == 1 + assert len(result.spawned_statuses) == 1 + meta = next(iter(result.metadata.values())) + assert meta.parent_plan_id == _PLAN_ID + print("spawn-single-ok") + + +def _spawn_multiple() -> None: + """Spawn multiple child plans from decision entries.""" + svc = SubplanService(decision_service=_mock_decision_service()) + plan = _make_plan() + config = SubplanConfig( + execution_mode=ExecutionMode.PARALLEL, + max_parallel=5, + ) + entries = [] + for i, did in enumerate([_DEC_ID1, _DEC_ID2, _DEC_ID3]): + dec = _make_decision(did, sequence=i) + entries.append(SpawnEntry(decision=dec, action_name=f"local/sub-{i}")) + result = svc.spawn(parent_plan=plan, config=config, spawn_entries=entries) + assert result.total_spawned == 3 + ids = [s.subplan_id for s in result.spawned_statuses] + assert len(ids) == len(set(ids)), "Duplicate IDs" + print("spawn-multiple-ok") + + +def _validate_resource_scope() -> None: + """Validate spawn with missing resource scope.""" + svc = SubplanService(decision_service=_mock_decision_service()) + config = SubplanConfig(execution_mode=ExecutionMode.SEQUENTIAL) + dec = _make_decision(_DEC_ID1, resource_ids=["res-missing"]) + entries = [ + SpawnEntry( + decision=dec, + action_name="local/sub-action", + target_resources=["res-missing"], + ) + ] + result = svc.validate_spawn( + config=config, + spawn_entries=entries, + available_resources=frozenset({"res-a", "res-b"}), + ) + assert not result.valid + assert any("resource" in e.lower() for e in result.errors) + print("validate-resource-ok") + + +def _validate_max_parallel() -> None: + """Validate spawn exceeding max_parallel bound.""" + svc = SubplanService(decision_service=_mock_decision_service()) + config = SubplanConfig( + execution_mode=ExecutionMode.PARALLEL, + max_parallel=2, + ) + entries = [] + for i in range(5): + did = f"01HGZ6FE0AQDYTR4BXVQZ6D{i:02d}0" + dec = _make_decision(did, sequence=i) + entries.append(SpawnEntry(decision=dec, action_name=f"local/sub-{i}")) + result = svc.validate_spawn(config=config, spawn_entries=entries) + assert not result.valid + assert any("max_parallel" in e.lower() for e in result.errors) + print("validate-max-parallel-ok") + + +def _metadata_persistence() -> None: + """Verify spawn metadata contains correct plan IDs.""" + svc = SubplanService(decision_service=_mock_decision_service()) + plan = _make_plan( + parent_plan_id="01HGZ6FE0AQDYTR4BXVQZ6PP00", + root_plan_id=_ROOT_ID, + ) + config = SubplanConfig(execution_mode=ExecutionMode.SEQUENTIAL) + dec = _make_decision(_DEC_ID1) + entries = [SpawnEntry(decision=dec, action_name="local/sub-action")] + result = svc.spawn(parent_plan=plan, config=config, spawn_entries=entries) + meta = next(iter(result.metadata.values())) + assert meta.root_plan_id == _ROOT_ID + assert meta.execution_mode == "sequential" + print("metadata-persistence-ok") + + +def _validation_guards() -> None: + """Verify service rejects invalid construction and inputs.""" + # None decision_service + try: + SubplanService(decision_service=None) # type: ignore[arg-type] + raise AssertionError("Should have raised ValueError") + except ValueError: + pass + + # None parent_plan + svc = SubplanService(decision_service=_mock_decision_service()) + try: + svc.spawn( + parent_plan=None, # type: ignore[arg-type] + config=SubplanConfig(), + spawn_entries=[ + SpawnEntry( + decision=_make_decision(_DEC_ID1), + action_name="local/test", + ) + ], + ) + raise AssertionError("Should have raised ValueError") + except ValueError: + pass + + # Empty spawn_entries + try: + svc.spawn( + parent_plan=_make_plan(), + config=SubplanConfig(), + spawn_entries=[], + ) + raise AssertionError("Should have raised ValueError") + except ValueError: + pass + + print("validation-guards-ok") + + +def _build_entries() -> None: + """Build spawn entries from decisions.""" + svc = SubplanService(decision_service=_mock_decision_service()) + decisions = [ + _make_decision( + _DEC_ID1, + resource_ids=["01HGZ6FE0AQDYTR4BXVQZ6RA00"], + chosen_option="local/test-sub", + ), + _make_decision( + _DEC_ID2, + sequence=1, + chosen_option="local/test-sub-2", + ), + ] + entries = svc.build_spawn_entries(decisions) + assert len(entries) == 2 + assert entries[0].action_name == "local/test-sub" + assert len(entries[0].target_resources) == 1 + print("build-entries-ok") + + +_COMMANDS: dict[str, object] = { + "spawn-single": _spawn_single, + "spawn-multiple": _spawn_multiple, + "validate-resource": _validate_resource_scope, + "validate-max-parallel": _validate_max_parallel, + "metadata-persistence": _metadata_persistence, + "validation-guards": _validation_guards, + "build-entries": _build_entries, +} + +if __name__ == "__main__": + if len(sys.argv) < 2 or sys.argv[1] not in _COMMANDS: + print(f"Usage: {sys.argv[0]} <{'|'.join(_COMMANDS)}>", file=sys.stderr) + sys.exit(1) + fn = _COMMANDS[sys.argv[1]] + assert callable(fn) + fn() diff --git a/robot/subplan_spawn.robot b/robot/subplan_spawn.robot new file mode 100644 index 000000000..ff097228d --- /dev/null +++ b/robot/subplan_spawn.robot @@ -0,0 +1,59 @@ +*** Settings *** +Documentation Smoke tests for subplan spawn service, validation, +... and metadata persistence. +Resource ${CURDIR}/common.resource +Suite Setup Setup Test Environment +Suite Teardown Cleanup Test Environment + +*** Variables *** +${HELPER} robot/helper_subplan_spawn.py + +*** Test Cases *** +Spawn Single Child Plan From Decision Entry + [Documentation] Verify spawning a single child plan creates status and metadata + [Tags] subplan spawn + ${result}= Run Process ${PYTHON} ${HELPER} spawn-single cwd=${WORKSPACE} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} spawn-single-ok + +Spawn Multiple Child Plans From Decisions + [Documentation] Verify spawning multiple child plans from decision entries + [Tags] subplan spawn multiple + ${result}= Run Process ${PYTHON} ${HELPER} spawn-multiple cwd=${WORKSPACE} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} spawn-multiple-ok + +Validate Resource Scope Missing + [Documentation] Verify validation rejects missing resource scopes + [Tags] subplan spawn validation + ${result}= Run Process ${PYTHON} ${HELPER} validate-resource cwd=${WORKSPACE} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} validate-resource-ok + +Validate Max Parallel Exceeded + [Documentation] Verify validation rejects exceeding max_parallel bounds + [Tags] subplan spawn validation parallel + ${result}= Run Process ${PYTHON} ${HELPER} validate-max-parallel cwd=${WORKSPACE} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} validate-max-parallel-ok + +Spawn Metadata Persistence + [Documentation] Verify spawn metadata contains correct plan IDs and mode + [Tags] subplan spawn metadata + ${result}= Run Process ${PYTHON} ${HELPER} metadata-persistence cwd=${WORKSPACE} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} metadata-persistence-ok + +Service Validation Guards + [Documentation] Verify service rejects invalid inputs + [Tags] subplan spawn validation + ${result}= Run Process ${PYTHON} ${HELPER} validation-guards cwd=${WORKSPACE} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} validation-guards-ok + +Build Spawn Entries From Decisions + [Documentation] Verify building spawn entries from decision objects + [Tags] subplan spawn entries + ${result}= Run Process ${PYTHON} ${HELPER} build-entries cwd=${WORKSPACE} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} build-entries-ok diff --git a/src/cleveragents/application/container.py b/src/cleveragents/application/container.py index 5546d957e..59de5cc17 100644 --- a/src/cleveragents/application/container.py +++ b/src/cleveragents/application/container.py @@ -25,6 +25,7 @@ from cleveragents.application.services.project_service import ProjectService from cleveragents.application.services.resource_registry_service import ( ResourceRegistryService, ) +from cleveragents.application.services.subplan_service import SubplanService from cleveragents.application.services.vector_store_service import VectorStoreService from cleveragents.config.settings import Settings, get_settings from cleveragents.domain.providers.ai_provider import AIProviderInterface @@ -266,6 +267,12 @@ class Container(containers.DeclarativeContainer): plan_lifecycle_service=plan_lifecycle_service, ) + # Subplan Service - Factory (spawns child plans from decisions) + subplan_service = providers.Factory( + SubplanService, + decision_service=decision_service, + ) + # Resource Registry Service - uses database session factory from UoW resource_registry_service = providers.Factory( _build_resource_registry_service, diff --git a/src/cleveragents/application/services/__init__.py b/src/cleveragents/application/services/__init__.py index dbe7f4ea8..733698046 100644 --- a/src/cleveragents/application/services/__init__.py +++ b/src/cleveragents/application/services/__init__.py @@ -71,6 +71,14 @@ from cleveragents.application.services.subplan_merge_service import ( SubplanMergeResult, SubplanMergeService, ) +from cleveragents.application.services.subplan_service import ( + SpawnEntry, + SpawnMetadata, + SpawnResult, + SpawnValidationError, + SpawnValidationResult, + SubplanService, +) from cleveragents.application.services.tool_registry_service import ( ToolRegistryService, ) @@ -126,11 +134,17 @@ __all__ = [ "SemanticValidationService", "SemanticValidationSeverity", "SkillRegistryService", + "SpawnEntry", + "SpawnMetadata", + "SpawnResult", + "SpawnValidationError", + "SpawnValidationResult", "SubplanExecutionOutput", "SubplanExecutionResult", "SubplanExecutionService", "SubplanMergeResult", "SubplanMergeService", + "SubplanService", "SyntaxCheckRule", "ToolRegistryService", "ValidationAttachment", diff --git a/src/cleveragents/application/services/subplan_service.py b/src/cleveragents/application/services/subplan_service.py new file mode 100644 index 000000000..d2ca44564 --- /dev/null +++ b/src/cleveragents/application/services/subplan_service.py @@ -0,0 +1,447 @@ +"""Subplan service for building and spawning child plans. + +Coordinates the spawning of child plans from ``DecisionService`` spawn entries +and ``SubplanConfig``. The service validates resource scopes, merge strategy +presence, and ``max_parallel`` bounds before creating child plans. + +Spawn metadata (spawn decision ID, parent/root plan IDs, execution mode) is +persisted alongside each child plan for status output and provenance tracking. + +Design decisions: + - Dependency injection: ``DecisionService`` and ``UnitOfWork`` are + injected via the constructor (consistent with existing service patterns). + - Stateless: All spawn state is derived from ``Decision`` entries and + ``SubplanConfig`` on the parent ``Plan``. + - Fail-fast: Validation errors are raised immediately before any plan + creation occurs. + +Based on: + - docs/specification.md L18170-L18295 (subplan spawning) + - ADR-006 (Plan Lifecycle) + - Forgejo issue #197 +""" + +from __future__ import annotations + +import logging +from dataclasses import dataclass, field +from typing import TYPE_CHECKING + +from cleveragents.core.exceptions import ValidationError +from cleveragents.domain.models.core.decision import Decision, DecisionType +from cleveragents.domain.models.core.plan import ( + ExecutionMode, + Plan, + SubplanConfig, + SubplanStatus, +) + +if TYPE_CHECKING: + from cleveragents.application.services.decision_service import DecisionService + +logger = logging.getLogger(__name__) + + +# --------------------------------------------------------------------------- +# Value objects +# --------------------------------------------------------------------------- + + +@dataclass(frozen=True) +class SpawnMetadata: + """Metadata for a spawned subplan. + + Persisted alongside each child plan for status output and provenance. + + Attributes: + spawn_decision_id: ULID of the decision that triggered the spawn. + parent_plan_id: ULID of the parent plan. + root_plan_id: ULID of the root plan in the hierarchy. + execution_mode: How the subplan should be executed. + """ + + spawn_decision_id: str + parent_plan_id: str + root_plan_id: str + execution_mode: str + + +@dataclass(frozen=True) +class SpawnEntry: + """A single spawn request derived from a decision. + + Attributes: + decision: The decision that triggered this spawn. + action_name: Namespaced action name for the child plan. + target_resources: Resource IDs the child plan operates on. + description: Description/prompt for the child plan. + """ + + decision: Decision + action_name: str + target_resources: list[str] = field(default_factory=list) + description: str = "" + + +@dataclass(frozen=True) +class SpawnValidationResult: + """Result of validating a spawn request. + + Attributes: + valid: Whether the spawn request passed all validation checks. + errors: List of validation error messages. + """ + + valid: bool + errors: list[str] = field(default_factory=list) + + +@dataclass(frozen=True) +class SpawnResult: + """Result of spawning child plans. + + Attributes: + spawned_statuses: Status objects for each spawned child plan. + metadata: Spawn metadata for each child plan (keyed by subplan_id). + total_spawned: Number of child plans spawned. + execution_mode: The execution mode from the SubplanConfig. + """ + + spawned_statuses: list[SubplanStatus] + metadata: dict[str, SpawnMetadata] = field(default_factory=dict) + total_spawned: int = 0 + execution_mode: str = ExecutionMode.SEQUENTIAL + + +# --------------------------------------------------------------------------- +# Exceptions +# --------------------------------------------------------------------------- + + +class SpawnValidationError(ValidationError): + """Raised when spawn validation fails. + + Attributes: + validation_errors: List of individual validation failure messages. + """ + + def __init__(self, validation_errors: list[str]) -> None: + self.validation_errors: list[str] = validation_errors + errors_str: str = "; ".join(validation_errors) + super().__init__(f"Spawn validation failed: {errors_str}") + + +# --------------------------------------------------------------------------- +# Service +# --------------------------------------------------------------------------- + + +class SubplanService: + """Service for spawning child plans from decision spawn entries. + + Coordinates the creation of child plans by: + + 1. Extracting spawn entries from ``DecisionService`` for a parent plan. + 2. Validating each entry (resource scopes, merge strategy, parallelism). + 3. Building ``SubplanStatus`` objects for the parent plan. + 4. Persisting ``SpawnMetadata`` for status queries. + + Args: + decision_service: Service for querying decision records. + + Raises: + ValueError: If *decision_service* is ``None``. + """ + + def __init__( + self, + decision_service: DecisionService, + ) -> None: + if decision_service is None: + raise ValueError("decision_service must not be None") + self._decision_service: DecisionService = decision_service + + @property + def decision_service(self) -> DecisionService: + """The injected decision service.""" + return self._decision_service + + # ------------------------------------------------------------------ + # spawn + # ------------------------------------------------------------------ + + def spawn( + self, + parent_plan: Plan, + config: SubplanConfig, + spawn_entries: list[SpawnEntry], + *, + available_resources: frozenset[str] | None = None, + ) -> SpawnResult: + """Build child plan statuses from spawn entries. + + Validates the spawn request, then creates ``SubplanStatus`` objects + and ``SpawnMetadata`` for each entry. The caller is responsible for + persisting the updated parent plan. + + Args: + parent_plan: The parent plan that will own the child plans. + config: Subplan execution configuration. + spawn_entries: Entries describing which child plans to create. + available_resources: Optional set of known resource IDs for + scope validation. When ``None``, resource scope checks + are skipped. + + Returns: + A :class:`SpawnResult` with statuses and metadata. + + Raises: + ValueError: If *parent_plan*, *config*, or *spawn_entries* is + ``None``. + ValueError: If *spawn_entries* is empty. + SpawnValidationError: If validation fails. + """ + if parent_plan is None: + raise ValueError("parent_plan must not be None") + if config is None: + raise ValueError("config must not be None") + if spawn_entries is None: + raise ValueError("spawn_entries must not be None") + if not spawn_entries: + raise ValueError("spawn_entries must not be empty") + + # Validate before spawning + validation: SpawnValidationResult = self.validate_spawn( + config=config, + spawn_entries=spawn_entries, + available_resources=available_resources, + ) + if not validation.valid: + raise SpawnValidationError(validation.errors) + + # Build statuses and metadata + parent_id: str = parent_plan.identity.plan_id + root_id: str = parent_plan.identity.root_plan_id or parent_id + mode: str = config.execution_mode.value + + statuses: list[SubplanStatus] = [] + metadata: dict[str, SpawnMetadata] = {} + + for entry in spawn_entries: + # Use the decision's downstream_plan_ids if available, + # otherwise generate a subplan_id from the decision_id + subplan_id: str = entry.decision.decision_id + + status: SubplanStatus = SubplanStatus( + subplan_id=subplan_id, + action_name=entry.action_name, + target_resources=list(entry.target_resources), + ) + statuses.append(status) + + meta: SpawnMetadata = SpawnMetadata( + spawn_decision_id=entry.decision.decision_id, + parent_plan_id=parent_id, + root_plan_id=root_id, + execution_mode=mode, + ) + metadata[subplan_id] = meta + + logger.info( + "spawned_subplans", + extra={ + "parent_plan_id": parent_id, + "root_plan_id": root_id, + "count": len(statuses), + "execution_mode": mode, + }, + ) + + return SpawnResult( + spawned_statuses=statuses, + metadata=metadata, + total_spawned=len(statuses), + execution_mode=mode, + ) + + # ------------------------------------------------------------------ + # validate_spawn + # ------------------------------------------------------------------ + + def validate_spawn( + self, + config: SubplanConfig, + spawn_entries: list[SpawnEntry], + *, + available_resources: frozenset[str] | None = None, + ) -> SpawnValidationResult: + """Validate a spawn request before execution. + + Checks: + 1. Resource scopes are resolved (all referenced resources exist). + 2. Merge strategy is defined (not None/missing). + 3. ``max_parallel`` bounds are respected (number of entries does + not exceed ``max_parallel`` when using PARALLEL mode). + + Args: + config: Subplan execution configuration. + spawn_entries: Entries describing which child plans to create. + available_resources: Optional set of known resource IDs. When + provided, each entry's ``target_resources`` are checked + against this set. + + Returns: + A :class:`SpawnValidationResult` indicating success or listing + errors. + + Raises: + ValueError: If *config* or *spawn_entries* is ``None``. + """ + if config is None: + raise ValueError("config must not be None") + if spawn_entries is None: + raise ValueError("spawn_entries must not be None") + + errors: list[str] = [] + + # 1. Resource scope validation + if available_resources is not None: + for entry in spawn_entries: + for resource_id in entry.target_resources: + if resource_id not in available_resources: + errors.append( + f"Unresolved resource scope: resource '{resource_id}' " + f"not found for spawn decision " + f"'{entry.decision.decision_id}'" + ) + + # 2. Merge strategy validation + if config.merge_strategy is None: + errors.append("Merge strategy must be defined in SubplanConfig") + + # 3. max_parallel bounds validation + if config.execution_mode == ExecutionMode.PARALLEL: + entry_count: int = len(spawn_entries) + if entry_count > config.max_parallel: + errors.append( + f"Number of spawn entries ({entry_count}) exceeds " + f"max_parallel bound ({config.max_parallel})" + ) + + # 4. Each entry must have a valid action_name + for entry in spawn_entries: + if not entry.action_name or not entry.action_name.strip(): + errors.append( + f"Spawn entry for decision " + f"'{entry.decision.decision_id}' has empty action_name" + ) + + # 5. Each entry must reference a spawn-type decision + for entry in spawn_entries: + if entry.decision.decision_type not in ( + DecisionType.SUBPLAN_SPAWN, + DecisionType.SUBPLAN_PARALLEL_SPAWN, + ): + errors.append( + f"Decision '{entry.decision.decision_id}' has type " + f"'{entry.decision.decision_type}' but must be " + f"'{DecisionType.SUBPLAN_SPAWN}' or " + f"'{DecisionType.SUBPLAN_PARALLEL_SPAWN}'" + ) + + valid: bool = len(errors) == 0 + return SpawnValidationResult(valid=valid, errors=errors) + + # ------------------------------------------------------------------ + # get_spawn_decisions + # ------------------------------------------------------------------ + + def get_spawn_decisions(self, plan_id: str) -> list[Decision]: + """Retrieve spawn-type decisions for a plan. + + Queries the ``DecisionService`` for both ``subplan_spawn`` and + ``subplan_parallel_spawn`` decision types. + + Args: + plan_id: ULID of the plan to query. + + Returns: + List of spawn decisions ordered by sequence number. + + Raises: + ValueError: If *plan_id* is empty. + """ + if not plan_id or not plan_id.strip(): + raise ValueError("plan_id must not be empty") + + spawn_decisions: list[Decision] = self._decision_service.list_by_type( + plan_id, DecisionType.SUBPLAN_SPAWN.value + ) + parallel_decisions: list[Decision] = self._decision_service.list_by_type( + plan_id, DecisionType.SUBPLAN_PARALLEL_SPAWN.value + ) + + all_decisions: list[Decision] = spawn_decisions + parallel_decisions + all_decisions.sort(key=lambda d: d.sequence_number) + return all_decisions + + # ------------------------------------------------------------------ + # build_spawn_entries + # ------------------------------------------------------------------ + + def build_spawn_entries( + self, + decisions: list[Decision], + default_action_name: str = "local/subplan-action", + ) -> list[SpawnEntry]: + """Build spawn entries from decisions. + + Converts decision objects into ``SpawnEntry`` objects suitable + for passing to :meth:`spawn`. + + Args: + decisions: Spawn-type decisions. + default_action_name: Action name to use when the decision's + ``chosen_option`` does not specify one. + + Returns: + List of spawn entries. + + Raises: + ValueError: If *decisions* is ``None``. + """ + if decisions is None: + raise ValueError("decisions must not be None") + + entries: list[SpawnEntry] = [] + for decision in decisions: + # Use chosen_option as action_name if it looks like a + # namespaced action reference; otherwise use default + action_name: str = default_action_name + if "/" in decision.chosen_option: + action_name = decision.chosen_option + + # Extract target resources from context snapshot + target_resources: list[str] = [ + ref.resource_id for ref in decision.context_snapshot.relevant_resources + ] + + entry: SpawnEntry = SpawnEntry( + decision=decision, + action_name=action_name, + target_resources=target_resources, + description=decision.question, + ) + entries.append(entry) + + return entries + + +__all__: list[str] = [ + "SpawnEntry", + "SpawnMetadata", + "SpawnResult", + "SpawnValidationError", + "SpawnValidationResult", + "SubplanService", +] diff --git a/vulture_whitelist.py b/vulture_whitelist.py index f8d0089b2..0970862bf 100644 --- a/vulture_whitelist.py +++ b/vulture_whitelist.py @@ -469,3 +469,17 @@ CheckpointNotFoundError # noqa: B018, F821 rollback_plan # noqa: B018, F821 _build_checkpoint_service # noqa: B018, F821 validate_checkpoint_type # noqa: B018, F821 + +# SubplanService public API (issue #197) +spawn_decision_id # noqa: B018, F821 +SpawnMetadata # noqa: B018, F821 +SpawnEntry # noqa: B018, F821 +SpawnResult # noqa: B018, F821 +SpawnValidationResult # noqa: B018, F821 +SpawnValidationError # noqa: B018, F821 +SubplanService # noqa: B018, F821 +subplan_service # noqa: B018, F821 +validation_errors # noqa: B018, F821 +get_spawn_decisions # noqa: B018, F821 +build_spawn_entries # noqa: B018, F821 +validate_spawn # noqa: B018, F821