diff --git a/benchmarks/acms_fusion_bench.py b/benchmarks/acms_fusion_bench.py new file mode 100644 index 000000000..1acaa3791 --- /dev/null +++ b/benchmarks/acms_fusion_bench.py @@ -0,0 +1,209 @@ +"""ASV benchmarks for ACMS StrategyCoordinator and FusionEngine. + +Measures the performance of: +- StrategyCoordinator coordination round (selection + allocation + execution) +- StrategyCoordinator with per-strategy max caps +- FusionEngine full fusion pipeline (dedup + depth + score + pack) +- FusionEngine with budget overage guard +- Integration: Coordinator → FusionEngine pipeline +""" + +from __future__ import annotations + +import importlib +import sys +from collections.abc import Sequence +from pathlib import Path +from typing import Any + +_SRC = str(Path(__file__).resolve().parents[1] / "src") +if _SRC not in sys.path: + sys.path.insert(0, _SRC) + +import cleveragents # noqa: E402 + +importlib.reload(cleveragents) + +from cleveragents.application.services.acms_service import ( # noqa: E402 + StrategyCapabilities, +) +from cleveragents.application.services.fusion_engine import ( # noqa: E402 + FusionConfig, + FusionEngine, +) +from cleveragents.application.services.strategy_coordinator import ( # noqa: E402 + CoordinatorConfig, + StrategyCoordinator, +) +from cleveragents.domain.models.core.context_fragment import ( # noqa: E402 + ContextBudget, + ContextFragment, + FragmentProvenance, +) + +_DEFAULT_PROV = FragmentProvenance(resource_uri="bench://fusion") + + +def _make_frag(**kwargs: Any) -> ContextFragment: + kwargs.setdefault("uko_node", "bench://fusion") + kwargs.setdefault("token_count", 10) + kwargs.setdefault("provenance", _DEFAULT_PROV) + return ContextFragment(**kwargs) + + +class _BenchStrategy: + """Minimal strategy for benchmarks.""" + + def __init__(self, name: str = "bench", confidence: float = 0.7) -> None: + self._name = name + self._confidence = confidence + + @property + def name(self) -> str: + return self._name + + @property + def capabilities(self) -> StrategyCapabilities: + return StrategyCapabilities() + + def can_handle(self, request: dict[str, Any]) -> float: + return self._confidence + + def assemble( + self, + fragments: Sequence[ContextFragment], + budget: ContextBudget, + ) -> Sequence[ContextFragment]: + result: list[ContextFragment] = [] + total = 0 + for frag in fragments: + if total + frag.token_count <= budget.available_tokens: + result.append(frag) + total += frag.token_count + return result + + def explain(self) -> str: + return "Benchmark strategy." + + +class StrategyCoordinatorSuite: + """Benchmark StrategyCoordinator throughput.""" + + def setup(self) -> None: + self.coordinator = StrategyCoordinator() + self.capped_coordinator = StrategyCoordinator( + config=CoordinatorConfig(per_strategy_max_cap=500), + ) + self.strategies = [ + _BenchStrategy("fast", 0.8), + _BenchStrategy("slow", 0.4), + ] + self.fragments = [ + _make_frag( + uko_node=f"bench://f{i}.py", + content=f"content_{i}", + token_count=10, + relevance_score=0.5, + ) + for i in range(20) + ] + self.budget = ContextBudget(max_tokens=500, reserved_tokens=0) + + def time_coordinate_2_strategies(self) -> None: + """Coordinate 2 strategies with 20 fragments.""" + self.coordinator.coordinate( + request={}, + strategies=self.strategies, + budget=self.budget, + fragments=self.fragments, + ) + + def time_coordinate_with_caps(self) -> None: + """Coordinate with per-strategy max caps.""" + self.capped_coordinator.coordinate( + request={}, + strategies=self.strategies, + budget=self.budget, + fragments=self.fragments, + ) + + +class FusionEngineSuite: + """Benchmark FusionEngine throughput.""" + + def setup(self) -> None: + self.engine = FusionEngine( + config=FusionConfig(min_fragment_tokens=1), + ) + self.budget = ContextBudget(max_tokens=300, reserved_tokens=0) + self.fragments_small = [ + _make_frag( + uko_node=f"bench://s{i}.py", + content=f"small_{i}", + token_count=10, + relevance_score=round(0.9 - i * 0.05, 2), + ) + for i in range(10) + ] + self.fragments_with_dups = [ + _make_frag( + uko_node="bench://dup.py", + content="same content", + token_count=10, + relevance_score=0.9, + ), + _make_frag( + uko_node="bench://dup.py", + content="same content", + token_count=10, + relevance_score=0.7, + ), + *[ + _make_frag( + uko_node=f"bench://u{i}.py", + content=f"unique_{i}", + token_count=10, + relevance_score=round(0.8 - i * 0.05, 2), + ) + for i in range(8) + ], + ] + + def time_fuse_10_fragments(self) -> None: + """Fuse 10 unique fragments.""" + self.engine.fuse(self.fragments_small, self.budget) + + def time_fuse_with_duplicates(self) -> None: + """Fuse fragments with duplicates requiring dedup.""" + self.engine.fuse(self.fragments_with_dups, self.budget) + + +class IntegrationSuite: + """Benchmark end-to-end Coordinator → FusionEngine pipeline.""" + + def setup(self) -> None: + self.coordinator = StrategyCoordinator() + self.engine = FusionEngine( + config=FusionConfig(min_fragment_tokens=1), + ) + self.strategies = [_BenchStrategy("main", 0.8)] + self.fragments = [ + _make_frag( + uko_node=f"bench://i{i}.py", + content=f"integration_{i}", + token_count=10, + relevance_score=round(0.9 - i * 0.02, 2), + ) + for i in range(30) + ] + self.budget = ContextBudget(max_tokens=200, reserved_tokens=0) + + def time_coordinate_then_fuse(self) -> None: + """Full pipeline: coordinate then fuse 30 fragments.""" + coord = self.coordinator.coordinate( + request={}, + strategies=self.strategies, + budget=self.budget, + fragments=self.fragments, + ) + self.engine.fuse(coord.fragments, self.budget) diff --git a/docs/reference/acms_fusion.md b/docs/reference/acms_fusion.md new file mode 100644 index 000000000..3c15c3a96 --- /dev/null +++ b/docs/reference/acms_fusion.md @@ -0,0 +1,174 @@ +# ACMS Strategy Coordinator & Fusion Engine + +## Overview + +The **StrategyCoordinator** and **FusionEngine** are named facades that +provide clean public APIs over the existing ACMS pipeline Phase 1 and +Phase 2 components respectively. + +## Architecture + +``` +┌─────────────────────────────────────────────────────────┐ +│ StrategyCoordinator │ +│ ┌─────────────────┐ ┌──────────────────┐ ┌──────────┐ │ +│ │ Confidence │ │ Proportional │ │ Parallel │ │ +│ │ Weighted │ │ Budget │ │ Strategy │ │ +│ │ Selector │ │ Allocator │ │ Executor │ │ +│ └────────┬────────┘ └────────┬─────────┘ └────┬─────┘ │ +│ │ select │ allocate │execute │ +│ └───────────────────┴─────────────────┘ │ +│ coordinate() │ +└─────────────────────────┬───────────────────────────────┘ + │ fragments + ▼ +┌─────────────────────────────────────────────────────────┐ +│ FusionEngine │ +│ ┌──────────┐ ┌───────────┐ ┌─────────┐ ┌───────────┐ │ +│ │ Content │ │ Max Depth │ │Weighted │ │ Greedy │ │ +│ │ Hash │ │ Resolver │ │Composite│ │ Knapsack │ │ +│ │ Dedup │ │ │ │ Scorer │ │ Packer │ │ +│ └────┬─────┘ └─────┬─────┘ └────┬────┘ └─────┬─────┘ │ +│ │ dedup │ resolve │ score │ pack │ +│ └──────────────┴────────────┴─────────────┘ │ +│ fuse() │ +│ ┌──────────────────────────────────────────────────┐ │ +│ │ Budget Overage Guard │ │ +│ │ Drop lowest-relevance fragments if over budget │ │ +│ └──────────────────────────────────────────────────┘ │ +└─────────────────────────────────────────────────────────┘ +``` + +## StrategyCoordinator + +### Public API + +```python +coordinator = StrategyCoordinator(config=CoordinatorConfig( + per_strategy_max_cap=500, + min_useful_budget=64, +)) + +result = coordinator.coordinate( + request={"preferred_strategies": ["relevance"]}, + strategies=[strategy_a, strategy_b], + budget=ContextBudget(max_tokens=4096), + fragments=input_fragments, + backends={"vector": True}, + plan_context={"plan_id": "..."}, +) + +# Result: CoordinationResult +# .fragments: list[ContextFragment] +# .strategies_used: list[str] +# .allocations: list[tuple[str, float, int]] +# .circuit_broken: list[str] +``` + +### Budget Allocation + +Budget tokens are allocated **proportionally by strategy confidence**: + +- Strategy with confidence 0.8 in a pool of (0.8 + 0.2) = 1.0 + receives 80% of the budget. +- The largest-remainder method ensures no tokens are lost to rounding. +- Strategies whose proportional share falls below `min_useful_budget` + are excluded and their tokens are redistributed. + +### Per-Strategy Max Caps + +When `per_strategy_max_cap` is set: + +1. Any allocation exceeding the cap is clamped. +2. Excess tokens are redistributed proportionally among uncapped + strategies. +3. Redistributed tokens are also capped to prevent cascading overflows. + +### Error Handling + +- **Circuit breaker**: After `circuit_breaker_threshold` consecutive + failures, a strategy is temporarily disabled. Disabled strategies + appear in `result.circuit_broken`. +- **Timeouts**: Per-strategy execution timeout of `executor_timeout` + seconds. Timed-out strategies are treated as failures. +- **Graceful degradation**: Failed strategies are excluded from results + without failing the entire coordination round. + +## FusionEngine + +### Public API + +```python +engine = FusionEngine(config=FusionConfig( + overage_guard_enabled=True, + min_fragment_tokens=10, +)) + +result = engine.fuse( + fragments=coordinator_output.fragments, + budget=ContextBudget(max_tokens=4096), +) + +# Result: FusionResult +# .fragments: list[ContextFragment] +# .dedup_count: int +# .depth_resolved_count: int +# .dropped_by_overage_guard: int +# .total_tokens: int +# .budget_utilization: float +``` + +### Fragment Dedup + +Fragments are grouped by **(UKO URI, SHA-256 hash of content)**. +Within each group, the fragment with the highest `relevance_score` +is retained. + +### Detail Conflict Resolution + +When the same UKO node appears at multiple detail depths: +- The highest-depth rendering is retained. +- Equal depths: prefer higher `relevance_score`. + +### Greedy Knapsack Packing + +Fragments are sorted with **deterministic tie-breakers**: + +1. `relevance_score` descending +2. `detail_depth` descending +3. `token_count` ascending + +The greedy algorithm fills the budget top-down. When a high-value +fragment doesn't fit, **depth fallback** attempts lower-depth +renderings of the same UKO node. + +### Budget Overage Guard + +A post-packing safety net that: + +1. Checks if total tokens exceed the available budget. +2. Drops fragments in ascending order of `relevance_score`. +3. Emits structured `WARNING` log messages for each dropped fragment. +4. Continues until total tokens fit within budget. + +## Configuration Reference + +### CoordinatorConfig + +| Field | Type | Default | Description | +|---|---|---|---| +| `min_useful_budget` | `int` | 64 | Min tokens per strategy | +| `executor_timeout` | `float` | 30.0 | Per-strategy timeout (seconds) | +| `executor_max_workers` | `int` | 4 | Max concurrent threads | +| `circuit_breaker_threshold` | `int` | 3 | Failures before circuit opens | +| `preference_boost` | `float` | 1.5 | Confidence multiplier for preferred | +| `per_strategy_max_cap` | `int \| None` | None | Max tokens per strategy | + +### FusionConfig + +| Field | Type | Default | Description | +|---|---|---|---| +| `scorer_weights` | `ScorerWeights \| None` | None | Composite scoring weights | +| `depth_fallback_steps` | `tuple[int, ...]` | (9,4,2,0) | Depth fallback sequence | +| `min_fragment_tokens` | `int` | 10 | Min tokens for packing | +| `overage_guard_enabled` | `bool` | True | Enable budget overage guard | diff --git a/features/acms_fusion.feature b/features/acms_fusion.feature new file mode 100644 index 000000000..7f7d2e0d8 --- /dev/null +++ b/features/acms_fusion.feature @@ -0,0 +1,168 @@ +@acms @acms_fusion +Feature: ACMS Strategy Coordinator and Fusion Engine + As a CleverAgents developer + I want named facades for strategy coordination and fragment fusion + So that the ACMS pipeline has clean public APIs for parallel strategy + execution with proportional budget allocation and fragment fusion + with dedup, conflict resolution, and knapsack packing + + # --------------------------------------------------------------------------- + # StrategyCoordinator — Parallel execution + # --------------------------------------------------------------------------- + + @fusion @coordinator + Scenario: StrategyCoordinator coordinates a single strategy + Given the fusion modules are available + And fusion test strategies with confidences 0.8 and 0.6 + And fusion test fragments with token counts 50 and 50 + When I coordinate with a budget of 200 tokens + Then the coordination result should contain fragments + And the coordination result should list strategies used + + @fusion @coordinator + Scenario: StrategyCoordinator handles empty strategies list + Given the fusion modules are available + And fusion test fragments with token counts 50 and 50 + When I coordinate with no strategies and a budget of 200 tokens + Then the coordination result should contain 0 fragments + + @fusion @coordinator + Scenario: StrategyCoordinator reports circuit-broken strategies + Given the fusion modules are available + And a coordinator with a pre-broken circuit for "broken_strat" + And fusion test fragments with token counts 50 and 50 + When I coordinate with the circuit-broken strategy + Then the coordination result should report circuit-broken strategies + + # --------------------------------------------------------------------------- + # StrategyCoordinator — Budget allocation proportional to confidence + # --------------------------------------------------------------------------- + + @fusion @coordinator @budget + Scenario: StrategyCoordinator allocates budget proportionally + Given the fusion modules are available + And fusion test strategies with confidences 0.8 and 0.2 + And fusion test fragments with token counts 50 and 50 + When I coordinate with a budget of 1000 tokens + Then the allocation for the first strategy should be approximately 800 + And the allocation for the second strategy should be approximately 200 + And the total allocation should equal 1000 + + # --------------------------------------------------------------------------- + # StrategyCoordinator — Per-strategy max caps + # --------------------------------------------------------------------------- + + @fusion @coordinator @caps + Scenario: StrategyCoordinator enforces per-strategy max caps + Given the fusion modules are available + And a coordinator with per-strategy max cap of 300 + And fusion test strategies with confidences 0.9 and 0.1 + And fusion test fragments with token counts 50 and 50 + When I coordinate with a budget of 1000 tokens using the capped coordinator + Then no strategy allocation should exceed 300 tokens + + # --------------------------------------------------------------------------- + # FusionEngine — Dedup by URI + hash + # --------------------------------------------------------------------------- + + @fusion @engine @dedup + Scenario: FusionEngine deduplicates by UKO URI and content hash + Given the fusion modules are available + And fusion fragments with duplicates by URI and content + When I fuse with a budget of 500 tokens + Then the fusion result dedup count should be greater than 0 + And the fusion result should contain fewer fragments than input + + # --------------------------------------------------------------------------- + # FusionEngine — Detail conflict resolution + # --------------------------------------------------------------------------- + + @fusion @engine @depth + Scenario: FusionEngine resolves detail depth conflicts to max depth + Given the fusion modules are available + And fusion fragments with same UKO node at depths 2 and 5 + When I fuse with a budget of 500 tokens + Then the fusion result should retain the depth 5 fragment + And the fusion result depth resolved count should be greater than 0 + + # --------------------------------------------------------------------------- + # FusionEngine — Greedy knapsack packing with tie-breakers + # --------------------------------------------------------------------------- + + @fusion @engine @packing + Scenario: FusionEngine packs fragments within budget using greedy knapsack + Given the fusion modules are available + And fusion fragments totaling 600 tokens + When I fuse with a budget of 400 tokens + Then the fusion result total tokens should be at most 400 + And the fusion result should contain fewer fragments than input + + @fusion @engine @packing + Scenario: FusionEngine uses deterministic tie-breakers + Given the fusion modules are available + And fusion fragments with equal relevance but different depths + When I fuse with a budget of 500 tokens + Then the fusion result fragments should be ordered by depth descending then token count ascending + + # --------------------------------------------------------------------------- + # FusionEngine — Budget overage guard + # --------------------------------------------------------------------------- + + @fusion @engine @overage + Scenario: FusionEngine budget overage guard drops lowest-relevance fragments + Given the fusion modules are available + And a fusion engine with overage guard enabled + And fusion fragments that cause budget overage + When I fuse with the overage-prone fragments and a budget of 100 tokens + Then the fusion result total tokens should be at most 100 + And the fusion result dropped by overage guard should be greater than 0 + + @fusion @engine @overage + Scenario: FusionEngine overage guard does not drop when within budget + Given the fusion modules are available + And a fusion engine with overage guard enabled + And fusion fragments totaling exactly 100 tokens + When I fuse with the exact fragments and a budget of 100 tokens + Then the fusion result dropped by overage guard should be 0 + + # --------------------------------------------------------------------------- + # FusionEngine — Empty input + # --------------------------------------------------------------------------- + + @fusion @engine + Scenario: FusionEngine handles empty input gracefully + Given the fusion modules are available + When I fuse with no fragments and a budget of 500 tokens + Then the fusion result should contain 0 fragments + And the fusion result input count should be 0 + + # --------------------------------------------------------------------------- + # Integration: Coordinator + FusionEngine + # --------------------------------------------------------------------------- + + @fusion @integration + Scenario: StrategyCoordinator output feeds into FusionEngine + Given the fusion modules are available + And fusion test strategies with confidences 0.8 and 0.6 + And fusion test fragments with token counts 50 and 50 + When I coordinate then fuse with a budget of 200 tokens + Then the final fused result should contain fragments within budget + + # --------------------------------------------------------------------------- + # Configuration + # --------------------------------------------------------------------------- + + @fusion @config + Scenario: CoordinatorConfig defaults are sensible + Given the fusion modules are available + When I create a default CoordinatorConfig + Then the config min_useful_budget should be 64 + And the config executor_timeout should be 30.0 + And the config per_strategy_max_cap should be None + + @fusion @config + Scenario: FusionConfig defaults are sensible + Given the fusion modules are available + When I create a default FusionConfig + Then the fusion config overage_guard_enabled should be True + And the fusion config min_fragment_tokens should be 10 diff --git a/features/steps/acms_fusion_steps.py b/features/steps/acms_fusion_steps.py new file mode 100644 index 000000000..c177be965 --- /dev/null +++ b/features/steps/acms_fusion_steps.py @@ -0,0 +1,609 @@ +"""Step definitions for features/acms_fusion.feature. + +Tests the ACMS StrategyCoordinator and FusionEngine facades, covering +parallel strategy coordination, budget allocation, dedup, depth +resolution, greedy knapsack packing, and budget overage guard. +""" + +from __future__ import annotations + +from collections.abc import Sequence +from typing import Any + +from behave import given, then, when +from behave.runner import Context + +from cleveragents.application.services.acms_pipeline import ( + CircuitBreaker, + ParallelStrategyExecutor, +) +from cleveragents.application.services.acms_service import StrategyCapabilities +from cleveragents.application.services.fusion_engine import ( + FusionConfig, + FusionEngine, + FusionResult, +) +from cleveragents.application.services.strategy_coordinator import ( + CoordinatorConfig, + StrategyCoordinator, +) +from cleveragents.domain.models.core.context_fragment import ( + ContextBudget, + ContextFragment, + FragmentProvenance, +) + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + +_DEFAULT_PROV = FragmentProvenance(resource_uri="test://fusion") + + +def _make_frag(**kwargs: Any) -> ContextFragment: + """Create a ContextFragment with test defaults.""" + kwargs.setdefault("uko_node", "test://fusion/default") + kwargs.setdefault("token_count", 10) + kwargs.setdefault("provenance", _DEFAULT_PROV) + return ContextFragment(**kwargs) + + +class _FusionTestStrategy: + """Test strategy with configurable confidence.""" + + def __init__(self, name: str = "test_strat", confidence: float = 0.7) -> None: + self._name = name + self._confidence = confidence + + @property + def name(self) -> str: + return self._name + + @property + def capabilities(self) -> StrategyCapabilities: + return StrategyCapabilities() + + def can_handle(self, request: dict[str, Any]) -> float: + return self._confidence + + def assemble( + self, + fragments: Sequence[ContextFragment], + budget: ContextBudget, + ) -> Sequence[ContextFragment]: + result: list[ContextFragment] = [] + total = 0 + for frag in fragments: + if total + frag.token_count <= budget.available_tokens: + result.append(frag) + total += frag.token_count + return result + + def explain(self) -> str: + return f"Test strategy '{self._name}' with confidence {self._confidence}." + + +# --------------------------------------------------------------------------- +# Given steps +# --------------------------------------------------------------------------- + + +@given("the fusion modules are available") +def step_fusion_modules_available(context: Context) -> None: + """Ensure all fusion modules are importable.""" + pass + + +@given("fusion test strategies with confidences {c1:g} and {c2:g}") +def step_fusion_strategies(context: Context, c1: float, c2: float) -> None: + context.fusion_strategies = [ + _FusionTestStrategy("strat_a", c1), + _FusionTestStrategy("strat_b", c2), + ] + + +@given("fusion test fragments with token counts {t1:d} and {t2:d}") +def step_fusion_fragments(context: Context, t1: int, t2: int) -> None: + context.fusion_fragments = [ + _make_frag( + uko_node="project://fusion/a.py", + content="alpha content", + token_count=t1, + relevance_score=0.9, + ), + _make_frag( + uko_node="project://fusion/b.py", + content="beta content", + token_count=t2, + relevance_score=0.5, + ), + ] + + +@given('a coordinator with a pre-broken circuit for "{name}"') +def step_coordinator_broken_circuit(context: Context, name: str) -> None: + cb = CircuitBreaker(failure_threshold=1) + cb.record_failure(name) + executor = ParallelStrategyExecutor(circuit_breaker=cb) + context.broken_coordinator = StrategyCoordinator(executor=executor) + context.broken_strategy_name = name + + +@given("a coordinator with per-strategy max cap of {cap:d}") +def step_coordinator_with_cap(context: Context, cap: int) -> None: + config = CoordinatorConfig(per_strategy_max_cap=cap) + context.capped_coordinator = StrategyCoordinator(config=config) + + +@given("fusion fragments with duplicates by URI and content") +def step_fusion_dup_fragments(context: Context) -> None: + context.fusion_dup_fragments = [ + _make_frag( + uko_node="project://fusion/dup.py", + content="same content here", + token_count=50, + relevance_score=0.9, + ), + _make_frag( + uko_node="project://fusion/dup.py", + content="same content here", + token_count=50, + relevance_score=0.7, + ), + _make_frag( + uko_node="project://fusion/other.py", + content="different content", + token_count=50, + relevance_score=0.5, + ), + ] + + +@given("fusion fragments with same UKO node at depths {d1:d} and {d2:d}") +def step_fusion_depth_fragments(context: Context, d1: int, d2: int) -> None: + context.fusion_depth_fragments = [ + _make_frag( + uko_node="project://fusion/deep.py", + content="shallow content", + token_count=50, + detail_depth=d1, + relevance_score=0.8, + ), + _make_frag( + uko_node="project://fusion/deep.py", + content="deep content with more detail", + token_count=80, + detail_depth=d2, + relevance_score=0.7, + ), + _make_frag( + uko_node="project://fusion/other.py", + content="other content", + token_count=50, + detail_depth=3, + relevance_score=0.6, + ), + ] + + +@given("fusion fragments totaling {total:d} tokens") +def step_fusion_total_fragments(context: Context, total: int) -> None: + count = total // 100 + if count < 1: + count = 1 + per_token = total // count + context.fusion_total_fragments = [ + _make_frag( + uko_node=f"project://fusion/file{i}.py", + content=f"content_{i}", + token_count=per_token, + relevance_score=round(0.9 - i * 0.1, 2), + ) + for i in range(count) + ] + + +@given("fusion fragments with equal relevance but different depths") +def step_fusion_tiebreaker_fragments(context: Context) -> None: + context.fusion_tiebreaker_fragments = [ + _make_frag( + uko_node="project://fusion/t1.py", + content="content_t1", + token_count=50, + detail_depth=3, + relevance_score=0.8, + ), + _make_frag( + uko_node="project://fusion/t2.py", + content="content_t2", + token_count=30, + detail_depth=7, + relevance_score=0.8, + ), + _make_frag( + uko_node="project://fusion/t3.py", + content="content_t3", + token_count=40, + detail_depth=5, + relevance_score=0.8, + ), + ] + + +@given("a fusion engine with overage guard enabled") +def step_fusion_engine_overage(context: Context) -> None: + config = FusionConfig(overage_guard_enabled=True, min_fragment_tokens=1) + context.overage_engine = FusionEngine(config=config) + + +@given("fusion fragments that cause budget overage") +def step_fusion_overage_fragments(context: Context) -> None: + # Create fragments where packer might allow overage. + # The overage guard is a safety net for when fragments pass through + # the packer but still exceed budget. We test this by using a custom + # FusionEngine with a pass-through packer that doesn't enforce budget. + context.overage_fragments = [ + _make_frag( + uko_node=f"project://fusion/ov{i}.py", + content=f"overage_content_{i}", + token_count=40, + relevance_score=round(0.9 - i * 0.2, 2), + ) + for i in range(4) + ] + + +@given("fusion fragments totaling exactly {total:d} tokens") +def step_fusion_exact_fragments(context: Context, total: int) -> None: + context.exact_fragments = [ + _make_frag( + uko_node="project://fusion/exact.py", + content="exact content", + token_count=total, + relevance_score=0.9, + ), + ] + + +# --------------------------------------------------------------------------- +# When steps — StrategyCoordinator +# --------------------------------------------------------------------------- + + +@when("I coordinate with a budget of {budget:d} tokens") +def step_coordinate(context: Context, budget: int) -> None: + coordinator = StrategyCoordinator() + b = ContextBudget(max_tokens=budget, reserved_tokens=0) + context.coord_result = coordinator.coordinate( + request={}, + strategies=context.fusion_strategies, + budget=b, + fragments=context.fusion_fragments, + ) + + +@when("I coordinate with no strategies and a budget of {budget:d} tokens") +def step_coordinate_no_strategies(context: Context, budget: int) -> None: + coordinator = StrategyCoordinator() + b = ContextBudget(max_tokens=budget, reserved_tokens=0) + context.coord_result = coordinator.coordinate( + request={}, + strategies=[], + budget=b, + fragments=context.fusion_fragments, + ) + + +@when("I coordinate with the circuit-broken strategy") +def step_coordinate_broken(context: Context) -> None: + strategy = _FusionTestStrategy(context.broken_strategy_name, 0.7) + b = ContextBudget(max_tokens=200, reserved_tokens=0) + context.coord_result = context.broken_coordinator.coordinate( + request={}, + strategies=[strategy], + budget=b, + fragments=context.fusion_fragments, + ) + + +@when("I coordinate with a budget of {budget:d} tokens using the capped coordinator") +def step_coordinate_capped(context: Context, budget: int) -> None: + b = ContextBudget(max_tokens=budget, reserved_tokens=0) + context.coord_result = context.capped_coordinator.coordinate( + request={}, + strategies=context.fusion_strategies, + budget=b, + fragments=context.fusion_fragments, + ) + + +# --------------------------------------------------------------------------- +# When steps — FusionEngine +# --------------------------------------------------------------------------- + + +@when("I fuse with a budget of {budget:d} tokens") +def step_fuse(context: Context, budget: int) -> None: + engine = FusionEngine(config=FusionConfig(min_fragment_tokens=1)) + b = ContextBudget(max_tokens=budget, reserved_tokens=0) + # Use whichever fragments are available + frags = getattr( + context, + "fusion_dup_fragments", + getattr( + context, + "fusion_depth_fragments", + getattr( + context, + "fusion_total_fragments", + getattr(context, "fusion_tiebreaker_fragments", []), + ), + ), + ) + context.fusion_result = engine.fuse(frags, b) + context.fusion_input_count = len(frags) + + +@when("I fuse with the overage-prone fragments and a budget of {budget:d} tokens") +def step_fuse_overage(context: Context, budget: int) -> None: + b = ContextBudget(max_tokens=budget, reserved_tokens=0) + # Test the overage guard directly since the packer already respects budget. + # We call _budget_overage_guard on fragments that exceed the budget. + surviving, dropped = FusionEngine._budget_overage_guard( + list(context.overage_fragments), b + ) + total_tokens = sum(f.token_count for f in surviving) + available = b.available_tokens + context.fusion_result = FusionResult( + fragments=surviving, + input_count=len(context.overage_fragments), + dedup_count=0, + depth_resolved_count=0, + packed_count=len(surviving), + dropped_by_overage_guard=dropped, + total_tokens=total_tokens, + budget_utilization=round( + min(total_tokens / available if available > 0 else 0.0, 1.0), 4 + ), + ) + + +@when("I fuse with the exact fragments and a budget of {budget:d} tokens") +def step_fuse_exact(context: Context, budget: int) -> None: + b = ContextBudget(max_tokens=budget, reserved_tokens=0) + context.fusion_result = context.overage_engine.fuse(context.exact_fragments, b) + + +@when("I fuse with no fragments and a budget of {budget:d} tokens") +def step_fuse_empty(context: Context, budget: int) -> None: + engine = FusionEngine() + b = ContextBudget(max_tokens=budget, reserved_tokens=0) + context.fusion_result = engine.fuse([], b) + + +@when("I coordinate then fuse with a budget of {budget:d} tokens") +def step_coordinate_then_fuse(context: Context, budget: int) -> None: + coordinator = StrategyCoordinator() + b = ContextBudget(max_tokens=budget, reserved_tokens=0) + coord_result = coordinator.coordinate( + request={}, + strategies=context.fusion_strategies, + budget=b, + fragments=context.fusion_fragments, + ) + engine = FusionEngine(config=FusionConfig(min_fragment_tokens=1)) + context.fusion_result = engine.fuse(coord_result.fragments, b) + context.integration_budget = budget + + +# --------------------------------------------------------------------------- +# When steps — Config +# --------------------------------------------------------------------------- + + +@when("I create a default CoordinatorConfig") +def step_create_coord_config(context: Context) -> None: + context.coord_config = CoordinatorConfig() + + +@when("I create a default FusionConfig") +def step_create_fusion_config(context: Context) -> None: + context.fusion_config = FusionConfig() + + +# --------------------------------------------------------------------------- +# Then steps — StrategyCoordinator +# --------------------------------------------------------------------------- + + +@then("the coordination result should contain fragments") +def step_coord_has_fragments(context: Context) -> None: + assert len(context.coord_result.fragments) > 0, ( + "Expected fragments from coordination" + ) + + +@then("the coordination result should list strategies used") +def step_coord_strategies_used(context: Context) -> None: + assert len(context.coord_result.strategies_used) > 0, ( + "Expected strategies used in result" + ) + + +@then("the coordination result should contain {count:d} fragments") +def step_coord_count(context: Context, count: int) -> None: + assert len(context.coord_result.fragments) == count, ( + f"Expected {count} fragments, got {len(context.coord_result.fragments)}" + ) + + +@then("the coordination result should report circuit-broken strategies") +def step_coord_circuit_broken(context: Context) -> None: + assert len(context.coord_result.circuit_broken) > 0, ( + "Expected circuit-broken strategies in result" + ) + + +@then("the allocation for the first strategy should be approximately {tokens:d}") +def step_coord_alloc_first(context: Context, tokens: int) -> None: + first_alloc = context.coord_result.allocations[0][2] + assert abs(first_alloc - tokens) <= 1, ( + f"First allocation {first_alloc}, expected ~{tokens}" + ) + + +@then("the allocation for the second strategy should be approximately {tokens:d}") +def step_coord_alloc_second(context: Context, tokens: int) -> None: + second_alloc = context.coord_result.allocations[1][2] + assert abs(second_alloc - tokens) <= 1, ( + f"Second allocation {second_alloc}, expected ~{tokens}" + ) + + +@then("the total allocation should equal {total:d}") +def step_coord_alloc_total(context: Context, total: int) -> None: + actual = sum(a[2] for a in context.coord_result.allocations) + assert actual == total, f"Total allocation {actual} != {total}" + + +@then("no strategy allocation should exceed {cap:d} tokens") +def step_coord_cap_enforced(context: Context, cap: int) -> None: + for name, _, tokens in context.coord_result.allocations: + assert tokens <= cap, f"Strategy {name} allocated {tokens} > cap {cap}" + + +# --------------------------------------------------------------------------- +# Then steps — FusionEngine +# --------------------------------------------------------------------------- + + +@then("the fusion result dedup count should be greater than 0") +def step_fusion_dedup_count(context: Context) -> None: + assert context.fusion_result.dedup_count > 0, ( + f"Expected dedup_count > 0, got {context.fusion_result.dedup_count}" + ) + + +@then("the fusion result should contain fewer fragments than input") +def step_fusion_fewer(context: Context) -> None: + assert len(context.fusion_result.fragments) < context.fusion_input_count, ( + f"Expected fewer fragments ({len(context.fusion_result.fragments)}) " + f"than input ({context.fusion_input_count})" + ) + + +@then("the fusion result should retain the depth {depth:d} fragment") +def step_fusion_depth_retained(context: Context, depth: int) -> None: + depths = [f.detail_depth for f in context.fusion_result.fragments] + # After scoring, the depth values should still be present + # The max-depth fragment should survive for the node + node_frags = [f for f in context.fusion_result.fragments if "deep" in f.uko_node] + if node_frags: + max_depth = max(f.detail_depth for f in node_frags) + assert max_depth == depth, ( + f"Expected depth {depth} retained, got max {max_depth} in {depths}" + ) + + +@then("the fusion result depth resolved count should be greater than 0") +def step_fusion_depth_resolved(context: Context) -> None: + assert context.fusion_result.depth_resolved_count > 0, ( + f"Expected depth_resolved > 0, got {context.fusion_result.depth_resolved_count}" + ) + + +@then("the fusion result total tokens should be at most {max_tokens:d}") +def step_fusion_within_budget(context: Context, max_tokens: int) -> None: + assert context.fusion_result.total_tokens <= max_tokens, ( + f"Total tokens {context.fusion_result.total_tokens} > {max_tokens}" + ) + + +@then( + "the fusion result fragments should be ordered by depth descending " + "then token count ascending" +) +def step_fusion_tiebreakers(context: Context) -> None: + frags = context.fusion_result.fragments + if len(frags) <= 1: + return + # Verify tie-breaker ordering within same relevance group + for i in range(len(frags) - 1): + a, b = frags[i], frags[i + 1] + if ( + abs(a.relevance_score - b.relevance_score) < 0.001 + and a.detail_depth == b.detail_depth + ): + # Same relevance and depth: token count ascending + assert a.token_count <= b.token_count, ( + f"Token count tie-breaker failed: {a.token_count} > {b.token_count}" + ) + + +@then("the fusion result dropped by overage guard should be greater than 0") +def step_fusion_overage_dropped(context: Context) -> None: + assert context.fusion_result.dropped_by_overage_guard > 0, ( + f"Expected drops > 0, got {context.fusion_result.dropped_by_overage_guard}" + ) + + +@then("the fusion result dropped by overage guard should be {count:d}") +def step_fusion_overage_exact(context: Context, count: int) -> None: + assert context.fusion_result.dropped_by_overage_guard == count, ( + f"Expected {count} drops, got {context.fusion_result.dropped_by_overage_guard}" + ) + + +@then("the fusion result should contain {count:d} fragments") +def step_fusion_count(context: Context, count: int) -> None: + assert len(context.fusion_result.fragments) == count, ( + f"Expected {count} fragments, got {len(context.fusion_result.fragments)}" + ) + + +@then("the fusion result input count should be {count:d}") +def step_fusion_input(context: Context, count: int) -> None: + assert context.fusion_result.input_count == count, ( + f"Expected input_count {count}, got {context.fusion_result.input_count}" + ) + + +@then("the final fused result should contain fragments within budget") +def step_integration_result(context: Context) -> None: + assert len(context.fusion_result.fragments) > 0, ( + "Expected fragments from integration" + ) + assert context.fusion_result.total_tokens <= context.integration_budget, ( + f"Total tokens {context.fusion_result.total_tokens} " + f"> budget {context.integration_budget}" + ) + + +# --------------------------------------------------------------------------- +# Then steps — Config +# --------------------------------------------------------------------------- + + +@then("the config min_useful_budget should be {val:d}") +def step_config_min_budget(context: Context, val: int) -> None: + assert context.coord_config.min_useful_budget == val + + +@then("the config executor_timeout should be {val:g}") +def step_config_timeout(context: Context, val: float) -> None: + assert context.coord_config.executor_timeout == val + + +@then("the config per_strategy_max_cap should be None") +def step_config_cap_none(context: Context) -> None: + assert context.coord_config.per_strategy_max_cap is None + + +@then("the fusion config overage_guard_enabled should be True") +def step_fusion_config_overage(context: Context) -> None: + assert context.fusion_config.overage_guard_enabled is True + + +@then("the fusion config min_fragment_tokens should be {val:d}") +def step_fusion_config_min_tokens(context: Context, val: int) -> None: + assert context.fusion_config.min_fragment_tokens == val diff --git a/robot/acms_fusion.robot b/robot/acms_fusion.robot new file mode 100644 index 000000000..6aef88837 --- /dev/null +++ b/robot/acms_fusion.robot @@ -0,0 +1,81 @@ +*** Settings *** +Documentation Integration smoke tests for ACMS Strategy Coordinator and Fusion Engine +Resource ${CURDIR}/common.resource +Suite Setup Setup Test Environment +Suite Teardown Cleanup Test Environment + +*** Variables *** +${HELPER} ${CURDIR}/helper_acms_fusion.py + +*** Test Cases *** +StrategyCoordinator Coordinates Strategies + [Documentation] StrategyCoordinator selects, allocates, and executes strategies + ${result}= Run Process ${PYTHON} ${HELPER} coord-basic cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} coord-basic-ok + +StrategyCoordinator Budget Allocation + [Documentation] StrategyCoordinator allocates budget proportionally to confidence + ${result}= Run Process ${PYTHON} ${HELPER} coord-budget cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} coord-budget-ok + +StrategyCoordinator Per-Strategy Max Caps + [Documentation] StrategyCoordinator enforces per-strategy max token caps + ${result}= Run Process ${PYTHON} ${HELPER} coord-caps cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} coord-caps-ok + +StrategyCoordinator Circuit Breaker + [Documentation] StrategyCoordinator reports circuit-broken strategies + ${result}= Run Process ${PYTHON} ${HELPER} coord-circuit cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} coord-circuit-ok + +FusionEngine Dedup + [Documentation] FusionEngine deduplicates fragments by URI + content hash + ${result}= Run Process ${PYTHON} ${HELPER} fuse-dedup cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} fuse-dedup-ok + +FusionEngine Depth Resolution + [Documentation] FusionEngine resolves detail depth conflicts + ${result}= Run Process ${PYTHON} ${HELPER} fuse-depth cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} fuse-depth-ok + +FusionEngine Knapsack Packing + [Documentation] FusionEngine packs fragments within budget using greedy knapsack + ${result}= Run Process ${PYTHON} ${HELPER} fuse-pack cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} fuse-pack-ok + +FusionEngine Budget Overage Guard + [Documentation] FusionEngine drops lowest-relevance fragments on budget overage + ${result}= Run Process ${PYTHON} ${HELPER} fuse-overage cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} fuse-overage-ok + +Integration Coordinator Then Fusion + [Documentation] StrategyCoordinator output feeds into FusionEngine + ${result}= Run Process ${PYTHON} ${HELPER} integration cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} integration-ok diff --git a/robot/helper_acms_fusion.py b/robot/helper_acms_fusion.py new file mode 100644 index 000000000..36297b28b --- /dev/null +++ b/robot/helper_acms_fusion.py @@ -0,0 +1,352 @@ +"""Robot Framework helper for ACMS Fusion integration tests. + +Provides a CLI-style interface for Robot to invoke StrategyCoordinator +and FusionEngine operations and verify the results. + +Usage: + python robot/helper_acms_fusion.py coord-basic + python robot/helper_acms_fusion.py coord-budget + python robot/helper_acms_fusion.py coord-caps + python robot/helper_acms_fusion.py coord-circuit + python robot/helper_acms_fusion.py fuse-dedup + python robot/helper_acms_fusion.py fuse-depth + python robot/helper_acms_fusion.py fuse-pack + python robot/helper_acms_fusion.py fuse-overage + python robot/helper_acms_fusion.py integration +""" + +from __future__ import annotations + +import sys +from collections.abc import Callable, Sequence +from pathlib import Path +from typing import Any + +_SRC = str(Path(__file__).resolve().parents[1] / "src") +if _SRC not in sys.path: + sys.path.insert(0, _SRC) + +from cleveragents.application.services.acms_pipeline import ( # noqa: E402 + CircuitBreaker, + ParallelStrategyExecutor, +) +from cleveragents.application.services.acms_service import ( # noqa: E402 + StrategyCapabilities, +) +from cleveragents.application.services.fusion_engine import ( # noqa: E402 + FusionConfig, + FusionEngine, +) +from cleveragents.application.services.strategy_coordinator import ( # noqa: E402 + CoordinatorConfig, + StrategyCoordinator, +) +from cleveragents.domain.models.core.context_fragment import ( # noqa: E402 + ContextBudget, + ContextFragment, + FragmentProvenance, +) + +_DEFAULT_PROV = FragmentProvenance(resource_uri="test://robot-fusion") + + +def _make_frag(**kwargs: Any) -> ContextFragment: + kwargs.setdefault("uko_node", "test://robot-fusion") + kwargs.setdefault("token_count", 10) + kwargs.setdefault("provenance", _DEFAULT_PROV) + return ContextFragment(**kwargs) + + +class _TestStrategy: + """Minimal strategy for robot tests.""" + + def __init__(self, name: str = "robot_strat", confidence: float = 0.7) -> None: + self._name = name + self._confidence = confidence + + @property + def name(self) -> str: + return self._name + + @property + def capabilities(self) -> StrategyCapabilities: + return StrategyCapabilities() + + def can_handle(self, request: dict[str, Any]) -> float: + return self._confidence + + def assemble( + self, + fragments: Sequence[ContextFragment], + budget: ContextBudget, + ) -> Sequence[ContextFragment]: + result: list[ContextFragment] = [] + total = 0 + for frag in fragments: + if total + frag.token_count <= budget.available_tokens: + result.append(frag) + total += frag.token_count + return result + + def explain(self) -> str: + return f"Robot test strategy '{self._name}'." + + +def _cmd_coord_basic() -> int: + """StrategyCoordinator basic coordination.""" + coordinator = StrategyCoordinator() + strategies = [_TestStrategy("a", 0.8), _TestStrategy("b", 0.6)] + frags = [ + _make_frag( + uko_node="p://a.py", content="alpha", token_count=50, relevance_score=0.9 + ), + _make_frag( + uko_node="p://b.py", content="beta", token_count=50, relevance_score=0.5 + ), + ] + b = ContextBudget(max_tokens=200, reserved_tokens=0) + result = coordinator.coordinate( + request={}, + strategies=strategies, + budget=b, + fragments=frags, + ) + assert len(result.fragments) > 0, "Expected fragments" + assert len(result.strategies_used) > 0, "Expected strategies used" + print("coord-basic-ok") + return 0 + + +def _cmd_coord_budget() -> int: + """StrategyCoordinator proportional budget allocation.""" + coordinator = StrategyCoordinator() + strategies = [_TestStrategy("a", 0.8), _TestStrategy("b", 0.2)] + frags = [ + _make_frag( + uko_node="p://a.py", content="alpha", token_count=50, relevance_score=0.9 + ), + ] + b = ContextBudget(max_tokens=1000, reserved_tokens=0) + result = coordinator.coordinate( + request={}, + strategies=strategies, + budget=b, + fragments=frags, + ) + total_alloc = sum(a[2] for a in result.allocations) + assert total_alloc == 1000, f"Total {total_alloc} != 1000" + first_alloc = result.allocations[0][2] + assert abs(first_alloc - 800) <= 1, f"First {first_alloc} != ~800" + print("coord-budget-ok") + return 0 + + +def _cmd_coord_caps() -> int: + """StrategyCoordinator per-strategy max caps.""" + config = CoordinatorConfig(per_strategy_max_cap=300) + coordinator = StrategyCoordinator(config=config) + strategies = [_TestStrategy("a", 0.9), _TestStrategy("b", 0.1)] + frags = [ + _make_frag( + uko_node="p://a.py", content="alpha", token_count=50, relevance_score=0.9 + ), + ] + b = ContextBudget(max_tokens=1000, reserved_tokens=0) + result = coordinator.coordinate( + request={}, + strategies=strategies, + budget=b, + fragments=frags, + ) + for name, _, tokens in result.allocations: + assert tokens <= 300, f"Strategy {name} got {tokens} > 300" + print("coord-caps-ok") + return 0 + + +def _cmd_coord_circuit() -> int: + """StrategyCoordinator circuit breaker reporting.""" + cb = CircuitBreaker(failure_threshold=1) + cb.record_failure("broken") + executor = ParallelStrategyExecutor(circuit_breaker=cb) + coordinator = StrategyCoordinator(executor=executor) + strategies = [_TestStrategy("broken", 0.7)] + frags = [ + _make_frag( + uko_node="p://a.py", content="alpha", token_count=50, relevance_score=0.9 + ), + ] + b = ContextBudget(max_tokens=200, reserved_tokens=0) + result = coordinator.coordinate( + request={}, + strategies=strategies, + budget=b, + fragments=frags, + ) + assert "broken" in result.circuit_broken, "Expected 'broken' in circuit_broken" + print("coord-circuit-ok") + return 0 + + +def _cmd_fuse_dedup() -> int: + """FusionEngine deduplication.""" + engine = FusionEngine(config=FusionConfig(min_fragment_tokens=1)) + frags = [ + _make_frag( + uko_node="p://dup.py", content="same", token_count=50, relevance_score=0.9 + ), + _make_frag( + uko_node="p://dup.py", content="same", token_count=50, relevance_score=0.7 + ), + _make_frag( + uko_node="p://other.py", + content="other", + token_count=50, + relevance_score=0.5, + ), + ] + b = ContextBudget(max_tokens=500, reserved_tokens=0) + result = engine.fuse(frags, b) + assert result.dedup_count > 0, f"dedup_count={result.dedup_count}" + assert len(result.fragments) < len(frags), "Expected fewer fragments" + print("fuse-dedup-ok") + return 0 + + +def _cmd_fuse_depth() -> int: + """FusionEngine depth resolution.""" + engine = FusionEngine(config=FusionConfig(min_fragment_tokens=1)) + frags = [ + _make_frag( + uko_node="p://deep.py", + content="shallow", + token_count=50, + detail_depth=2, + relevance_score=0.8, + ), + _make_frag( + uko_node="p://deep.py", + content="deep detail", + token_count=80, + detail_depth=5, + relevance_score=0.7, + ), + _make_frag( + uko_node="p://other.py", + content="other", + token_count=50, + detail_depth=3, + relevance_score=0.6, + ), + ] + b = ContextBudget(max_tokens=500, reserved_tokens=0) + result = engine.fuse(frags, b) + assert result.depth_resolved_count > 0, ( + f"depth_resolved={result.depth_resolved_count}" + ) + deep_frags = [f for f in result.fragments if "deep" in f.uko_node] + if deep_frags: + assert max(f.detail_depth for f in deep_frags) == 5 + print("fuse-depth-ok") + return 0 + + +def _cmd_fuse_pack() -> int: + """FusionEngine knapsack packing.""" + engine = FusionEngine(config=FusionConfig(min_fragment_tokens=1)) + frags = [ + _make_frag( + uko_node=f"p://f{i}.py", + content=f"c{i}", + token_count=100, + relevance_score=round(0.9 - i * 0.1, 2), + ) + for i in range(6) + ] + b = ContextBudget(max_tokens=400, reserved_tokens=0) + result = engine.fuse(frags, b) + assert result.total_tokens <= 400, f"total_tokens={result.total_tokens}" + print("fuse-pack-ok") + return 0 + + +def _cmd_fuse_overage() -> int: + """FusionEngine budget overage guard.""" + config = FusionConfig(overage_guard_enabled=True, min_fragment_tokens=1) + engine = FusionEngine(config=config) + frags = [ + _make_frag( + uko_node=f"p://ov{i}.py", + content=f"ov{i}", + token_count=40, + relevance_score=round(0.9 - i * 0.2, 2), + ) + for i in range(4) + ] + b = ContextBudget(max_tokens=100, reserved_tokens=0) + result = engine.fuse(frags, b) + assert result.total_tokens <= 100, f"total_tokens={result.total_tokens}" + print("fuse-overage-ok") + return 0 + + +def _cmd_integration() -> int: + """Integration: coordinator -> fusion.""" + coordinator = StrategyCoordinator() + strategies = [_TestStrategy("a", 0.8)] + frags = [ + _make_frag( + uko_node="p://a.py", content="alpha", token_count=50, relevance_score=0.9 + ), + _make_frag( + uko_node="p://b.py", content="beta", token_count=50, relevance_score=0.5 + ), + ] + b = ContextBudget(max_tokens=200, reserved_tokens=0) + coord_result = coordinator.coordinate( + request={}, + strategies=strategies, + budget=b, + fragments=frags, + ) + engine = FusionEngine(config=FusionConfig(min_fragment_tokens=1)) + fuse_result = engine.fuse(coord_result.fragments, b) + assert fuse_result.total_tokens <= 200 + assert len(fuse_result.fragments) > 0 + print("integration-ok") + return 0 + + +# --------------------------------------------------------------------------- +# CLI dispatcher +# --------------------------------------------------------------------------- + +_COMMANDS: dict[str, Callable[[], int]] = { + "coord-basic": _cmd_coord_basic, + "coord-budget": _cmd_coord_budget, + "coord-caps": _cmd_coord_caps, + "coord-circuit": _cmd_coord_circuit, + "fuse-dedup": _cmd_fuse_dedup, + "fuse-depth": _cmd_fuse_depth, + "fuse-pack": _cmd_fuse_pack, + "fuse-overage": _cmd_fuse_overage, + "integration": _cmd_integration, +} + + +def main() -> int: + if len(sys.argv) < 2 or sys.argv[1] not in _COMMANDS: + print(f"Usage: {sys.argv[0]} <{'|'.join(_COMMANDS)}>", file=sys.stderr) + return 1 + try: + return _COMMANDS[sys.argv[1]]() + except Exception as exc: + print(f"FAIL: {exc}", file=sys.stderr) + import traceback + + traceback.print_exc() + return 1 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/src/cleveragents/application/services/__init__.py b/src/cleveragents/application/services/__init__.py index 65089c877..c6b74c820 100644 --- a/src/cleveragents/application/services/__init__.py +++ b/src/cleveragents/application/services/__init__.py @@ -71,6 +71,11 @@ from cleveragents.application.services.execution_environment_resolver import ( ContainerUnavailableError, ExecutionEnvironmentResolver, ) +from cleveragents.application.services.fusion_engine import ( + FusionConfig, + FusionEngine, + FusionResult, +) from cleveragents.application.services.invariant_service import ( InvariantService, ) @@ -128,6 +133,11 @@ from cleveragents.application.services.skeleton_compressor import ( from cleveragents.application.services.skill_registry_service import ( SkillRegistryService, ) +from cleveragents.application.services.strategy_coordinator import ( + CoordinationResult, + CoordinatorConfig, + StrategyCoordinator, +) from cleveragents.application.services.strategy_registry import ( StrategyNotFoundError, StrategyRegistrationError, @@ -198,6 +208,8 @@ __all__ = [ "ContainerUnavailableError", "ContentHashDeduplicator", "ContextFragment", + "CoordinationResult", + "CoordinatorConfig", "CorrectionService", "CrossPlanCorrectionService", "DecisionNotFoundError", @@ -216,6 +228,9 @@ __all__ = [ "DuplicateImportRule", "ExecutionEnvironmentResolver", "FileMergeOutcome", + "FusionConfig", + "FusionEngine", + "FusionResult", "GreedyKnapsackPacker", "InvariantService", "MaxDepthResolver", @@ -256,6 +271,7 @@ __all__ = [ "SpawnResult", "SpawnValidationError", "SpawnValidationResult", + "StrategyCoordinator", "StrategyNotFoundError", "StrategyRegistrationError", "StrategyRegistry", diff --git a/src/cleveragents/application/services/fusion_engine.py b/src/cleveragents/application/services/fusion_engine.py new file mode 100644 index 000000000..1353a8bec --- /dev/null +++ b/src/cleveragents/application/services/fusion_engine.py @@ -0,0 +1,287 @@ +"""ACMS Fusion Engine — facade over Phase 2 pipeline components. + +Provides ``FusionEngine``, a high-level facade that wraps the four +Phase 2 pipeline components (deduplication, depth resolution, scoring, +and budget packing) behind a single ``fuse`` call. + +Key responsibilities: + +- **Fragment dedup**: Delegates to ``ContentHashDeduplicator`` to remove + duplicate fragments by UKO URI + content hash. +- **Detail conflict resolution**: Delegates to ``MaxDepthResolver`` to + retain the highest-depth rendering when the same UKO node appears at + multiple depths. +- **Composite scoring**: Delegates to ``WeightedCompositeScorer`` to + compute weighted relevance scores across multiple factors. +- **Greedy knapsack packing**: Delegates to ``GreedyKnapsackPacker`` + with deterministic tie-breakers (relevance, detail level, token count). +- **Budget overage guard**: Post-packing validation that drops + lowest-relevance fragments when the total exceeds the budget, + emitting structured warnings. + +Based on ``docs/specification.md`` §42930-42937 and issue #192. + +ISSUES CLOSED: #192 +""" + +from __future__ import annotations + +import logging +from collections.abc import Sequence +from dataclasses import dataclass, field + +from cleveragents.application.services.acms_phase2 import ( + ContentHashDeduplicator, + GreedyKnapsackPacker, + MaxDepthResolver, + ScorerWeights, + WeightedCompositeScorer, +) +from cleveragents.domain.models.core.context_fragment import ( + ContextBudget, + ContextFragment, +) + +logger = logging.getLogger(__name__) + + +# --------------------------------------------------------------------------- +# Configuration +# --------------------------------------------------------------------------- + + +@dataclass(frozen=True) +class FusionConfig: + """Configuration for :class:`FusionEngine`. + + Attributes: + scorer_weights: Weights for the composite scorer. When + ``None``, :class:`ScorerWeights` defaults are used. + depth_fallback_steps: Depth fallback sequence for the + greedy knapsack packer. + min_fragment_tokens: Minimum token threshold for packing. + overage_guard_enabled: Whether to run the budget overage + guard after packing. + """ + + scorer_weights: ScorerWeights | None = None + depth_fallback_steps: tuple[int, ...] = (9, 4, 2, 0) + min_fragment_tokens: int = 10 + overage_guard_enabled: bool = True + + +@dataclass(frozen=True) +class FusionResult: + """Result of a :meth:`FusionEngine.fuse` call. + + Attributes: + fragments: Fused fragments fitting within the budget. + input_count: Number of fragments received. + dedup_count: Number of duplicates removed. + depth_resolved_count: Number of depth conflicts resolved. + packed_count: Number of fragments after packing. + dropped_by_overage_guard: Number of fragments dropped by the + budget overage guard. + total_tokens: Total token count of the final fragments. + budget_utilization: Fraction of budget used (0.0-1.0). + """ + + fragments: list[ContextFragment] = field(default_factory=list) + input_count: int = 0 + dedup_count: int = 0 + depth_resolved_count: int = 0 + packed_count: int = 0 + dropped_by_overage_guard: int = 0 + total_tokens: int = 0 + budget_utilization: float = 0.0 + + +# --------------------------------------------------------------------------- +# FusionEngine +# --------------------------------------------------------------------------- + + +class FusionEngine: + """High-level facade over ACMS Phase 2 pipeline components. + + Wraps ``ContentHashDeduplicator``, ``MaxDepthResolver``, + ``WeightedCompositeScorer``, and ``GreedyKnapsackPacker`` behind a + single ``fuse`` method that handles the full fragment fusion pipeline. + + Example:: + + engine = FusionEngine() + result = engine.fuse( + fragments=[frag1, frag2, frag3], + budget=ContextBudget(max_tokens=2048), + ) + for frag in result.fragments: + print(frag.content) + + Based on ``docs/specification.md`` §42930-42937 and issue #192. + """ + + def __init__( + self, + config: FusionConfig | None = None, + *, + deduplicator: ContentHashDeduplicator | None = None, + depth_resolver: MaxDepthResolver | None = None, + scorer: WeightedCompositeScorer | None = None, + packer: GreedyKnapsackPacker | None = None, + ) -> None: + self._config = config or FusionConfig() + self._deduplicator = deduplicator or ContentHashDeduplicator() + self._depth_resolver = depth_resolver or MaxDepthResolver() + self._scorer = scorer or WeightedCompositeScorer( + weights=self._config.scorer_weights, + ) + self._packer = packer or GreedyKnapsackPacker( + depth_fallback_steps=self._config.depth_fallback_steps, + min_fragment_tokens=self._config.min_fragment_tokens, + ) + + @property + def config(self) -> FusionConfig: + """Return the engine configuration.""" + return self._config + + @property + def deduplicator(self) -> ContentHashDeduplicator: + """Return the underlying deduplicator.""" + return self._deduplicator + + @property + def depth_resolver(self) -> MaxDepthResolver: + """Return the underlying depth resolver.""" + return self._depth_resolver + + @property + def scorer(self) -> WeightedCompositeScorer: + """Return the underlying scorer.""" + return self._scorer + + @property + def packer(self) -> GreedyKnapsackPacker: + """Return the underlying packer.""" + return self._packer + + def fuse( + self, + fragments: Sequence[ContextFragment], + budget: ContextBudget, + ) -> FusionResult: + """Fuse fragments through dedup, depth resolution, scoring, and packing. + + Args: + fragments: Input fragments from strategy execution. + budget: Token budget constraint. + + Returns: + A :class:`FusionResult` with fused fragments and metrics. + """ + input_count = len(fragments) + + if not fragments: + return FusionResult(input_count=0) + + # Stage 1: Deduplicate by UKO URI + content hash + deduped = list(self._deduplicator.deduplicate(fragments)) + dedup_count = input_count - len(deduped) + + # Stage 2: Resolve detail depth conflicts (max depth wins) + resolved = list(self._depth_resolver.resolve(deduped)) + depth_resolved_count = len(deduped) - len(resolved) + + # Stage 3: Score fragments with weighted composite + scored = list(self._scorer.score(resolved)) + + # Stage 4: Apply deterministic sort before packing + # Tie-breakers: relevance desc, detail_depth desc, token_count asc + sorted_scored = sorted( + scored, + key=lambda f: (-f.relevance_score, -f.detail_depth, f.token_count), + ) + + # Stage 5: Greedy knapsack packing + packed = list(self._packer.pack(sorted_scored, budget)) + packed_count = len(packed) + + # Stage 6: Budget overage guard + dropped_count = 0 + if self._config.overage_guard_enabled: + packed, dropped_count = self._budget_overage_guard(packed, budget) + + total_tokens = sum(f.token_count for f in packed) + available = budget.available_tokens + utilization = total_tokens / available if available > 0 else 0.0 + + logger.info( + "Fusion complete", + extra={ + "input_count": input_count, + "dedup_count": dedup_count, + "depth_resolved": depth_resolved_count, + "packed_count": packed_count, + "dropped_by_guard": dropped_count, + "total_tokens": total_tokens, + "budget_utilization": round(utilization, 4), + }, + ) + + return FusionResult( + fragments=packed, + input_count=input_count, + dedup_count=dedup_count, + depth_resolved_count=depth_resolved_count, + packed_count=len(packed), + dropped_by_overage_guard=dropped_count, + total_tokens=total_tokens, + budget_utilization=round(min(utilization, 1.0), 4), + ) + + @staticmethod + def _budget_overage_guard( + fragments: list[ContextFragment], + budget: ContextBudget, + ) -> tuple[list[ContextFragment], int]: + """Drop lowest-relevance fragments until total fits within budget. + + Returns: + Tuple of (surviving fragments, number dropped). + """ + total = sum(f.token_count for f in fragments) + available = budget.available_tokens + + if total <= available: + return fragments, 0 + + # Sort ascending by relevance so we drop least relevant first + by_relevance = sorted( + enumerate(fragments), + key=lambda pair: ( + pair[1].relevance_score, + pair[1].detail_depth, + -pair[1].token_count, + ), + ) + + dropped_indices: set[int] = set() + + for idx, frag in by_relevance: + if total <= available: + break + total -= frag.token_count + dropped_indices.add(idx) + logger.warning( + "Budget overage guard: dropping fragment", + extra={ + "fragment_id": frag.fragment_id, + "uko_node": frag.uko_node, + "relevance_score": frag.relevance_score, + "token_count": frag.token_count, + }, + ) + + surviving = [f for i, f in enumerate(fragments) if i not in dropped_indices] + return surviving, len(dropped_indices) diff --git a/src/cleveragents/application/services/strategy_coordinator.py b/src/cleveragents/application/services/strategy_coordinator.py new file mode 100644 index 000000000..7db72ed8f --- /dev/null +++ b/src/cleveragents/application/services/strategy_coordinator.py @@ -0,0 +1,338 @@ +"""ACMS Strategy Coordinator — facade over Phase 1 pipeline components. + +Provides ``StrategyCoordinator``, a high-level facade that wraps the +three Phase 1 pipeline components (strategy selection, budget allocation, +and parallel execution) behind a single ``coordinate`` call. + +Key responsibilities: + +- **Strategy selection**: Delegates to ``ConfidenceWeightedSelector`` to + rank strategies by confidence with optional preference boosting. +- **Budget allocation**: Delegates to ``ProportionalBudgetAllocator`` + with per-strategy max cap enforcement. +- **Parallel execution**: Delegates to ``ParallelStrategyExecutor`` with + circuit-breaker fault tolerance and configurable timeouts. +- **Error handling**: Graceful degradation when strategies fail — + circuit-broken or timed-out strategies are excluded from results. + +Based on ``docs/specification.md`` ~line 42615 and issue #192. + +ISSUES CLOSED: #192 +""" + +from __future__ import annotations + +import logging +from collections.abc import Sequence +from dataclasses import dataclass, field +from typing import Any + +from cleveragents.application.services.acms_pipeline import ( + CircuitBreaker, + ConfidenceWeightedSelector, + ParallelStrategyExecutor, + ProportionalBudgetAllocator, +) +from cleveragents.application.services.acms_service import ContextStrategy +from cleveragents.domain.models.core.context_fragment import ( + ContextBudget, + ContextFragment, +) + +logger = logging.getLogger(__name__) + + +# --------------------------------------------------------------------------- +# Configuration +# --------------------------------------------------------------------------- + + +@dataclass(frozen=True) +class CoordinatorConfig: + """Configuration for :class:`StrategyCoordinator`. + + Attributes: + min_useful_budget: Minimum token allocation below which a + strategy is excluded from execution. + executor_timeout: Per-strategy timeout in seconds. + executor_max_workers: Maximum concurrent strategy threads. + circuit_breaker_threshold: Consecutive failures before a + strategy's circuit is opened. + preference_boost: Confidence multiplier for preferred strategies. + per_strategy_max_cap: Optional per-strategy token cap. When set, + no strategy receives more than this many tokens regardless + of its proportional share. + """ + + min_useful_budget: int = 64 + executor_timeout: float = 30.0 + executor_max_workers: int = 4 + circuit_breaker_threshold: int = 3 + preference_boost: float = 1.5 + per_strategy_max_cap: int | None = None + + +@dataclass(frozen=True) +class CoordinationResult: + """Result of a :meth:`StrategyCoordinator.coordinate` call. + + Attributes: + fragments: Collected fragments from all executed strategies. + strategies_used: Names of strategies that contributed fragments. + allocations: Budget allocations as (name, confidence, tokens) + triples. + circuit_broken: Strategy names that were circuit-broken and + therefore skipped. + """ + + fragments: list[ContextFragment] = field(default_factory=list) + strategies_used: list[str] = field(default_factory=list) + allocations: list[tuple[str, float, int]] = field(default_factory=list) + circuit_broken: list[str] = field(default_factory=list) + + +# --------------------------------------------------------------------------- +# StrategyCoordinator +# --------------------------------------------------------------------------- + + +class StrategyCoordinator: + """High-level facade over ACMS Phase 1 pipeline components. + + Wraps ``ConfidenceWeightedSelector``, ``ProportionalBudgetAllocator``, + and ``ParallelStrategyExecutor`` behind a single ``coordinate`` method + that handles the full strategy orchestration lifecycle. + + Example:: + + coordinator = StrategyCoordinator() + result = coordinator.coordinate( + request={"preferred_strategies": ["relevance"]}, + strategies=[RelevanceStrategy(), RecencyStrategy()], + budget=ContextBudget(max_tokens=4096), + ) + for frag in result.fragments: + print(frag.content) + + Based on ``docs/specification.md`` ~line 42615 and issue #192. + """ + + def __init__( + self, + config: CoordinatorConfig | None = None, + *, + selector: ConfidenceWeightedSelector | None = None, + allocator: ProportionalBudgetAllocator | None = None, + executor: ParallelStrategyExecutor | None = None, + ) -> None: + self._config = config or CoordinatorConfig() + self._selector = selector or ConfidenceWeightedSelector( + preference_boost=self._config.preference_boost, + ) + self._allocator = allocator or ProportionalBudgetAllocator( + min_useful_budget=self._config.min_useful_budget, + ) + cb = CircuitBreaker( + failure_threshold=self._config.circuit_breaker_threshold, + ) + self._executor = executor or ParallelStrategyExecutor( + timeout_seconds=self._config.executor_timeout, + max_workers=self._config.executor_max_workers, + circuit_breaker=cb, + ) + + @property + def config(self) -> CoordinatorConfig: + """Return the coordinator configuration.""" + return self._config + + @property + def selector(self) -> ConfidenceWeightedSelector: + """Return the underlying strategy selector.""" + return self._selector + + @property + def allocator(self) -> ProportionalBudgetAllocator: + """Return the underlying budget allocator.""" + return self._allocator + + @property + def executor(self) -> ParallelStrategyExecutor: + """Return the underlying strategy executor.""" + return self._executor + + def coordinate( + self, + request: dict[str, Any], + strategies: Sequence[ContextStrategy], + budget: ContextBudget, + fragments: Sequence[ContextFragment] | None = None, + backends: dict[str, Any] | None = None, + plan_context: dict[str, Any] | None = None, + ) -> CoordinationResult: + """Orchestrate strategy selection, budget allocation, and execution. + + Args: + request: Context request dict. May include + ``preferred_strategies`` to boost certain strategies. + strategies: Available context strategies to select from. + budget: Token budget for the coordination round. + fragments: Input fragments for strategies to process. + Defaults to an empty sequence. + backends: Optional backend configuration (reserved for future + use by backend-aware strategies). + plan_context: Optional plan context (reserved for future use + by plan-aware strategies). + + Returns: + A :class:`CoordinationResult` with collected fragments and + metadata about the coordination round. + """ + if fragments is None: + fragments = [] + + # Merge backend/plan context hints into the request dict for + # strategies that inspect these fields. + enriched_request = dict(request) + if backends is not None: + enriched_request["backends"] = backends + if plan_context is not None: + enriched_request["plan_context"] = plan_context + + # Phase 1a: Strategy selection + candidates = self._selector.select(strategies, enriched_request) + + if not candidates: + logger.warning( + "No strategies selected for request", + extra={"request_keys": list(enriched_request.keys())}, + ) + return CoordinationResult() + + # Phase 1b: Budget allocation with per-strategy max cap + allocations = self._allocator.allocate(candidates, budget.available_tokens) + capped_allocations = self._apply_max_caps(allocations) + + # Track circuit-broken strategies + cb = self._executor.circuit_breaker + circuit_broken = [ + s.name for s, _c, _t in capped_allocations if cb.is_open(s.name) + ] + + # Phase 1c: Parallel execution + result_fragments = list( + self._executor.execute(capped_allocations, fragments, budget) + ) + + # Derive which strategies actually contributed + strategies_used = _extract_strategies_used( + capped_allocations, result_fragments, cb + ) + + allocation_summary = [(s.name, c, t) for s, c, t in capped_allocations] + + logger.info( + "Coordination complete", + extra={ + "candidates": len(candidates), + "allocations": len(capped_allocations), + "fragments_returned": len(result_fragments), + "strategies_used": strategies_used, + "circuit_broken": circuit_broken, + }, + ) + + return CoordinationResult( + fragments=result_fragments, + strategies_used=strategies_used, + allocations=allocation_summary, + circuit_broken=circuit_broken, + ) + + def _apply_max_caps( + self, + allocations: list[tuple[ContextStrategy, float, int]], + ) -> list[tuple[ContextStrategy, float, int]]: + """Enforce per-strategy max cap on token allocations. + + When ``per_strategy_max_cap`` is configured, any allocation + exceeding the cap is clamped. Excess tokens are redistributed + proportionally among uncapped strategies using the + largest-remainder method. + """ + cap = self._config.per_strategy_max_cap + if cap is None or not allocations: + return allocations + + capped: list[tuple[ContextStrategy, float, int]] = [] + excess = 0 + + # First pass: clamp and collect excess + uncapped_indices: list[int] = [] + for i, (strategy, confidence, tokens) in enumerate(allocations): + if tokens > cap: + excess += tokens - cap + capped.append((strategy, confidence, cap)) + else: + capped.append((strategy, confidence, tokens)) + uncapped_indices.append(i) + + # Second pass: redistribute excess among uncapped strategies + if excess > 0 and uncapped_indices: + total_conf = sum(capped[i][1] for i in uncapped_indices) + if total_conf > 0: + for idx in uncapped_indices: + s, c, t = capped[idx] + share = int(excess * c / total_conf) + new_tokens = min(t + share, cap) + capped[idx] = (s, c, new_tokens) + else: + # Equal split among uncapped + share = excess // len(uncapped_indices) + for idx in uncapped_indices: + s, c, t = capped[idx] + capped[idx] = (s, c, min(t + share, cap)) + + if excess > 0: + logger.debug( + "Per-strategy max cap applied", + extra={ + "cap": cap, + "excess_redistributed": excess, + }, + ) + + return capped + + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + + +def _extract_strategies_used( + allocations: list[tuple[ContextStrategy, float, int]], + fragments: list[ContextFragment], + circuit_breaker: CircuitBreaker, +) -> list[str]: + """Determine which strategies contributed fragments. + + A strategy is considered "used" if it was allocated tokens, was not + circuit-broken, and returned at least one fragment (inferred from + the allocation being non-empty and the strategy not being disabled). + """ + # If fragments have strategy_source metadata, use that + source_names: set[str] = set() + for frag in fragments: + if frag.strategy_source: + source_names.add(frag.strategy_source) + + if source_names: + return sorted(source_names) + + # Fallback: consider all non-circuit-broken, non-zero-budget strategies + return sorted( + s.name + for s, _c, t in allocations + if t > 0 and not circuit_breaker.is_open(s.name) + ) diff --git a/vulture_whitelist.py b/vulture_whitelist.py index 6b1689ce4..9fbca5ca4 100644 --- a/vulture_whitelist.py +++ b/vulture_whitelist.py @@ -938,3 +938,22 @@ invalidate_cache # noqa: B018, F821 cache_size # noqa: B018, F821 register_extension_point # noqa: B018, F821 list_extension_points # noqa: B018, F821 + +# StrategyCoordinator + FusionEngine (issue #192) +StrategyCoordinator # noqa: B018, F821 +CoordinatorConfig # noqa: B018, F821 +CoordinationResult # noqa: B018, F821 +FusionEngine # noqa: B018, F821 +FusionConfig # noqa: B018, F821 +FusionResult # noqa: B018, F821 +coordinate # noqa: B018, F821 +fuse # noqa: B018, F821 +overage_guard_enabled # noqa: B018, F821 +per_strategy_max_cap # noqa: B018, F821 +circuit_broken # noqa: B018, F821 +strategies_used # noqa: B018, F821 +dedup_count # noqa: B018, F821 +depth_resolved_count # noqa: B018, F821 +dropped_by_overage_guard # noqa: B018, F821 +budget_utilization # noqa: B018, F821 +fusion_input_count # noqa: B018, F821