diff --git a/benchmarks/acms_pipeline_phase2_bench.py b/benchmarks/acms_pipeline_phase2_bench.py new file mode 100644 index 000000000..e59485084 --- /dev/null +++ b/benchmarks/acms_pipeline_phase2_bench.py @@ -0,0 +1,162 @@ +"""ASV benchmarks for ACMS pipeline Phase 2 — Fragment Fusion components. + +Measures the performance of: +- ContentHashDeduplicator with varying fragment counts +- MaxDepthResolver with varying fragment counts +- WeightedCompositeScorer with varying fragment counts +- GreedyKnapsackPacker with varying fragment counts and budgets +- ScoredFragment creation +""" + +from __future__ import annotations + +import importlib +import sys +from pathlib import Path + +_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_phase2 import ( # noqa: E402 + ContentHashDeduplicator, + GreedyKnapsackPacker, + MaxDepthResolver, + WeightedCompositeScorer, +) +from cleveragents.domain.models.core.context_fragment import ( # noqa: E402 + ContextBudget, + ContextFragment, + FragmentProvenance, + ScoredFragment, +) + +_DEFAULT_PROV = FragmentProvenance(resource_uri="bench://phase2") + + +def _make_fragments(n: int, *, unique_nodes: bool = True) -> list[ContextFragment]: + """Create *n* fragments for benchmarks.""" + frags: list[ContextFragment] = [] + for i in range(n): + node = f"bench://file{i}" if unique_nodes else "bench://same-node" + frags.append( + ContextFragment( + uko_node=node, + content=f"content-{i}", + relevance_score=round(0.1 + (i % 10) * 0.09, 2), + token_count=50 + (i % 5) * 10, + detail_depth=i % 10, + provenance=_DEFAULT_PROV, + ) + ) + return frags + + +class ContentHashDeduplicatorSuite: + """Benchmark ContentHashDeduplicator throughput.""" + + params: list[int] = [10, 100, 500] + param_names: list[str] = ["fragment_count"] + + def setup(self, n: int) -> None: + self.dedup = ContentHashDeduplicator() + # Half the fragments are duplicates + base = _make_fragments(n // 2, unique_nodes=True) + self.fragments = base + [ + ContextFragment( + uko_node=f.uko_node, + content=f.content, + relevance_score=max(0.0, f.relevance_score - 0.1), + token_count=f.token_count, + detail_depth=f.detail_depth, + provenance=_DEFAULT_PROV, + ) + for f in base + ] + + def time_deduplicate(self, n: int) -> None: + self.dedup.deduplicate(self.fragments) + + +class MaxDepthResolverSuite: + """Benchmark MaxDepthResolver throughput.""" + + params: list[int] = [10, 100, 500] + param_names: list[str] = ["fragment_count"] + + def setup(self, n: int) -> None: + self.resolver = MaxDepthResolver() + # Multiple depths per node + self.fragments: list[ContextFragment] = [] + for i in range(n): + node_idx = i % (n // 3 + 1) + self.fragments.append( + ContextFragment( + uko_node=f"bench://file{node_idx}", + content=f"content-{i}", + relevance_score=0.5, + token_count=50, + detail_depth=i % 10, + provenance=_DEFAULT_PROV, + ) + ) + + def time_resolve(self, n: int) -> None: + self.resolver.resolve(self.fragments) + + +class WeightedCompositeScorerSuite: + """Benchmark WeightedCompositeScorer throughput.""" + + params: list[int] = [10, 100, 500] + param_names: list[str] = ["fragment_count"] + + def setup(self, n: int) -> None: + self.scorer = WeightedCompositeScorer() + self.fragments = _make_fragments(n) + + def time_score(self, n: int) -> None: + self.scorer.score(self.fragments) + + def time_score_detailed(self, n: int) -> None: + self.scorer.score_detailed(self.fragments) + + +class GreedyKnapsackPackerSuite: + """Benchmark GreedyKnapsackPacker throughput.""" + + params: list[int] = [10, 100, 500] + param_names: list[str] = ["fragment_count"] + + def setup(self, n: int) -> None: + self.packer = GreedyKnapsackPacker() + self.fragments = _make_fragments(n) + # Budget that fits about half the fragments + total = sum(f.token_count for f in self.fragments) + self.budget = ContextBudget(max_tokens=total // 2, reserved_tokens=0) + + def time_pack(self, n: int) -> None: + self.packer.pack(self.fragments, self.budget) + + +class ScoredFragmentSuite: + """Benchmark ScoredFragment creation.""" + + def setup(self) -> None: + self.fragment = ContextFragment( + uko_node="bench://file", + content="benchmark content", + token_count=10, + provenance=_DEFAULT_PROV, + ) + + def time_scored_fragment_creation(self) -> None: + ScoredFragment( + fragment=self.fragment, + composite_score=0.85, + score_components={"relevance": 0.9, "hierarchy": 0.8}, + ) diff --git a/features/acms_pipeline_phase2.feature b/features/acms_pipeline_phase2.feature new file mode 100644 index 000000000..44e696cd9 --- /dev/null +++ b/features/acms_pipeline_phase2.feature @@ -0,0 +1,323 @@ +@phase2 @acms @acms_pipeline_phase2 +Feature: ACMS Pipeline Phase 2 — Fragment Fusion Components + As a CleverAgents developer + I want production-grade Phase 2 pipeline components + So that context assembly deduplicates, resolves depths, scores, and packs fragments + + # =========================================================================== + # ScoredFragment — Domain Model + # =========================================================================== + + @scored_fragment + Scenario: Create a ScoredFragment with component breakdown + Given a context fragment with uko_node "project://app/main.py" and content "hello" and token_count 10 for phase2 + When I create a ScoredFragment with composite_score 0.85 and components: + | component | value | + | relevance | 0.9 | + | hierarchy | 0.8 | + | quality | 0.7 | + | recency | 1.0 | + Then the ScoredFragment composite_score should be 0.85 + And the ScoredFragment should have component "relevance" with value 0.9 + And the ScoredFragment should have component "hierarchy" with value 0.8 + And the ScoredFragment should have component "quality" with value 0.7 + And the ScoredFragment should have component "recency" with value 1.0 + And the ScoredFragment should reference the original fragment + + @scored_fragment @validation + Scenario: ScoredFragment rejects composite score above 1.0 + Given a context fragment with uko_node "project://app/main.py" and content "hello" and token_count 10 for phase2 + When I create a ScoredFragment with invalid composite_score 1.5 + Then a phase2 validation error should be raised + + @scored_fragment @validation + Scenario: ScoredFragment rejects negative composite score + Given a context fragment with uko_node "project://app/main.py" and content "hello" and token_count 10 for phase2 + When I create a ScoredFragment with invalid composite_score -0.1 + Then a phase2 validation error should be raised + + # =========================================================================== + # ContentHashDeduplicator + # =========================================================================== + + @deduplicator + Scenario: Deduplicate fragments with identical content for same UKO node + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | hello world | 0.5 | 10 | 3 | + | project://app/main.py | hello world | 0.9 | 10 | 3 | + When I deduplicate the fragments + Then 1 fragment should remain after deduplication + And the surviving fragment should have relevance score 0.9 + + @deduplicator + Scenario: Deduplicate keeps distinct content for same UKO node + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | hello world | 0.5 | 10 | 3 | + | project://app/main.py | goodbye world | 0.9 | 12 | 3 | + When I deduplicate the fragments + Then 2 fragments should remain after deduplication + + @deduplicator + Scenario: Deduplicate keeps fragments from different UKO nodes + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | hello world | 0.5 | 10 | 3 | + | project://app/io.py | hello world | 0.9 | 10 | 3 | + When I deduplicate the fragments + Then 2 fragments should remain after deduplication + + @deduplicator + Scenario: Deduplicate with empty fragment list + Given an empty phase2 fragment list + When I deduplicate the fragments + Then 0 fragments should remain after deduplication + + @deduplicator + Scenario: Deduplicate three copies retains highest scored + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | abc | 0.3 | 5 | 1 | + | project://app/main.py | abc | 0.7 | 5 | 1 | + | project://app/main.py | abc | 0.5 | 5 | 1 | + When I deduplicate the fragments + Then 1 fragment should remain after deduplication + And the surviving fragment should have relevance score 0.7 + + # =========================================================================== + # MaxDepthResolver + # =========================================================================== + + @depth_resolver + Scenario: Resolve depth conflicts to maximum depth per node + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | brief summary | 0.8 | 10 | 2 | + | project://app/main.py | full source | 0.8 | 100 | 9 | + When I resolve depth conflicts + Then 1 fragment should remain after depth resolution + And the surviving fragment should have detail depth 9 + + @depth_resolver + Scenario: Resolve preserves fragments for distinct nodes + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | alpha | 0.8 | 10 | 2 | + | project://app/io.py | beta | 0.7 | 20 | 5 | + When I resolve depth conflicts + Then 2 fragments should remain after depth resolution + + @depth_resolver + Scenario: Resolve same depth keeps higher relevance + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | alpha | 0.3 | 10 | 5 | + | project://app/main.py | beta | 0.9 | 20 | 5 | + When I resolve depth conflicts + Then 1 fragment should remain after depth resolution + And the surviving fragment should have relevance score 0.9 + + @depth_resolver + Scenario: Resolve with empty fragment list + Given an empty phase2 fragment list + When I resolve depth conflicts + Then 0 fragments should remain after depth resolution + + @depth_resolver + Scenario: Resolve preserves original fragment order + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/a.py | aaa | 0.5 | 10 | 3 | + | project://app/b.py | bbb | 0.6 | 10 | 3 | + | project://app/c.py | ccc | 0.7 | 10 | 3 | + When I resolve depth conflicts + Then 3 fragments should remain after depth resolution + And the resolved fragment order should be "aaa", "bbb", "ccc" + + # =========================================================================== + # WeightedCompositeScorer + # =========================================================================== + + @scorer + Scenario: Score fragments with default weights + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | hello | 0.8 | 10 | 3 | + When I score the fragments with default weights + Then the scored fragment should have a composite relevance score + And the scored fragment metadata should contain original relevance "0.8" + + @scorer + Scenario: Score with custom weights + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | hello | 1.0 | 10 | 3 | + And scorer weights: relevance=1.0, hierarchy=0.0, quality=0.0, recency=0.0 + When I score the fragments with custom weights + Then the scored fragment relevance score should be 1.0 + + @scorer + Scenario: Score with metadata-provided hierarchy weight + Given a phase2 fragment with metadata: + | uko_node | content | score | tokens | depth | hierarchy_weight | strategy_quality | recency_bonus | + | project://app/main.py | hello | 0.8 | 10 | 3 | 0.9 | 0.7 | 1.0 | + When I score the fragments with default weights + Then the scored fragment metadata should contain component "hierarchy" near 0.9 + + @scorer + Scenario: Score empty fragment list + Given an empty phase2 fragment list + When I score the fragments with default weights + Then 0 scored fragments should be returned + + @scorer + Scenario: Detailed scoring returns ScoredFragment objects + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | hello | 0.8 | 10 | 3 | + | project://app/io.py | world | 0.6 | 15 | 5 | + When I score the fragments using score_detailed + Then 2 ScoredFragment objects should be returned + And each ScoredFragment should have a composite_score between 0.0 and 1.0 + And each ScoredFragment should have score_components + + @scorer + Scenario: Score composite is clamped to 0.0-1.0 + Given a phase2 fragment with metadata: + | uko_node | content | score | tokens | depth | hierarchy_weight | strategy_quality | recency_bonus | + | project://app/main.py | hello | 1.0 | 10 | 3 | 1.0 | 1.0 | 1.0 | + When I score the fragments with default weights + Then the scored fragment relevance score should be at most 1.0 + + # =========================================================================== + # GreedyKnapsackPacker + # =========================================================================== + + @packer + Scenario: Pack fragments within budget + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | alpha | 0.9 | 100 | 5 | + | project://app/io.py | beta | 0.7 | 100 | 3 | + | project://app/util.py | gamma | 0.5 | 100 | 2 | + And a phase2 budget with max_tokens 250 and reserved_tokens 0 + When I pack the fragments into budget + Then 2 fragments should be packed + And the packed total tokens should be at most 250 + + @packer + Scenario: Pack all fragments when budget is sufficient + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | alpha | 0.9 | 50 | 5 | + | project://app/io.py | beta | 0.7 | 50 | 3 | + And a phase2 budget with max_tokens 200 and reserved_tokens 0 + When I pack the fragments into budget + Then 2 fragments should be packed + + @packer + Scenario: Pack empty fragment list + Given an empty phase2 fragment list + And a phase2 budget with max_tokens 1000 and reserved_tokens 0 + When I pack the fragments into budget + Then 0 fragments should be packed + + @packer + Scenario: Pack excludes fragments below minimum token threshold + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | a | 0.9 | 5 | 3 | + | project://app/io.py | beta | 0.7 | 50 | 3 | + And a phase2 budget with max_tokens 1000 and reserved_tokens 0 + When I pack the fragments into budget + Then 1 fragment should be packed + And the packed fragment uko_node should be "project://app/io.py" + + @packer + Scenario: Pack with depth fallback uses lower-depth alternative + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | full source code | 0.9 | 500 | 9 | + | project://app/main.py | brief summary | 0.5 | 50 | 2 | + | project://app/io.py | io module | 0.7 | 100 | 5 | + And a phase2 budget with max_tokens 200 and reserved_tokens 0 + When I pack the fragments into budget + Then the packed fragments should include uko_node "project://app/main.py" + And the packed fragment for "project://app/main.py" should have depth 2 + + @packer + Scenario: Pack respects reserved tokens + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | alpha | 0.9 | 100 | 5 | + | project://app/io.py | beta | 0.7 | 100 | 3 | + And a phase2 budget with max_tokens 250 and reserved_tokens 100 + When I pack the fragments into budget + Then 1 fragment should be packed + And the packed total tokens should be at most 150 + + @packer + Scenario: Pack prefers higher scored fragments + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | alpha | 0.3 | 100 | 5 | + | project://app/io.py | beta | 0.9 | 100 | 3 | + And a phase2 budget with max_tokens 150 and reserved_tokens 0 + When I pack the fragments into budget + Then 1 fragment should be packed + And the packed fragment uko_node should be "project://app/io.py" + + @packer + Scenario: Single fragment over entire budget is excluded + Given the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | huge content | 0.9 | 10000 | 9 | + And a phase2 budget with max_tokens 100 and reserved_tokens 0 + When I pack the fragments into budget + Then 0 fragments should be packed + + # =========================================================================== + # Pipeline Integration — DI injection of Phase 2 components + # =========================================================================== + + @pipeline_integration + Scenario: Inject ContentHashDeduplicator into pipeline + Given the ACMS pipeline with a ContentHashDeduplicator + And the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | hello world | 0.5 | 10 | 3 | + | project://app/main.py | hello world | 0.9 | 10 | 3 | + When I assemble context through the pipeline + Then the pipeline output should contain 1 fragment + + @pipeline_integration + Scenario: Inject MaxDepthResolver into pipeline + Given the ACMS pipeline with a MaxDepthResolver + And the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | brief | 0.8 | 10 | 2 | + | project://app/main.py | detailed | 0.8 | 50 | 9 | + When I assemble context through the pipeline + Then the pipeline output should contain 1 fragment + + @pipeline_integration + Scenario: Inject WeightedCompositeScorer into pipeline + Given the ACMS pipeline with a WeightedCompositeScorer + And the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | hello | 0.8 | 10 | 3 | + When I assemble context through the pipeline + Then the pipeline output fragments should have updated relevance scores + + @pipeline_integration + Scenario: Inject GreedyKnapsackPacker into pipeline + Given the ACMS pipeline with a GreedyKnapsackPacker + And the following phase2 fragments: + | uko_node | content | score | tokens | depth | + | project://app/main.py | alpha | 0.9 | 100 | 5 | + | project://app/io.py | beta | 0.7 | 100 | 3 | + | project://app/util.py | gamma | 0.5 | 100 | 2 | + And a phase2 budget with max_tokens 250 and reserved_tokens 0 + When I assemble context through the pipeline with budget + Then the pipeline output total tokens should be at most 250 diff --git a/features/steps/acms_pipeline_phase2_steps.py b/features/steps/acms_pipeline_phase2_steps.py new file mode 100644 index 000000000..354c303b2 --- /dev/null +++ b/features/steps/acms_pipeline_phase2_steps.py @@ -0,0 +1,455 @@ +"""Step definitions for features/acms_pipeline_phase2.feature. + +Tests the ACMS pipeline Phase 2 (Fragment Fusion) components directly +in-memory — no database required. +""" + +from __future__ import annotations + +from typing import Any + +from behave import given, then, when +from behave.runner import Context +from pydantic import ValidationError + +from cleveragents.application.services.acms_phase2 import ( + ContentHashDeduplicator, + GreedyKnapsackPacker, + MaxDepthResolver, + ScorerWeights, + WeightedCompositeScorer, +) +from cleveragents.application.services.acms_service import ACMSPipeline +from cleveragents.domain.models.core.context_fragment import ( + ContextBudget, + ContextFragment, + FragmentProvenance, + ScoredFragment, +) + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + +_DEFAULT_PROVENANCE = FragmentProvenance(resource_uri="test://phase2") + + +def _make_phase2_fragment(**kwargs: Any) -> ContextFragment: + """Create a ContextFragment with sensible Phase 2 test defaults.""" + kwargs.setdefault("uko_node", "test://default") + kwargs.setdefault("token_count", 0) + kwargs.setdefault("provenance", _DEFAULT_PROVENANCE) + return ContextFragment(**kwargs) + + +# --------------------------------------------------------------------------- +# ScoredFragment — Given / When / Then +# --------------------------------------------------------------------------- + + +@given( + 'a context fragment with uko_node "{uko}" and content "{content}" ' + "and token_count {tokens:d} for phase2" +) +def step_given_fragment_for_phase2( + context: Context, uko: str, content: str, tokens: int +) -> None: + context.phase2_fragment = _make_phase2_fragment( + uko_node=uko, content=content, token_count=tokens + ) + + +@when("I create a ScoredFragment with composite_score {score:g} and components:") +def step_create_scored_fragment(context: Context, score: float) -> None: + components: dict[str, float] = {} + for row in context.table: + components[row["component"]] = float(row["value"]) + context.scored_fragment = ScoredFragment( + fragment=context.phase2_fragment, + composite_score=score, + score_components=components, + ) + + +@when("I create a ScoredFragment with invalid composite_score {score:g}") +def step_create_scored_fragment_invalid(context: Context, score: float) -> None: + context.phase2_error = None + try: + ScoredFragment( + fragment=context.phase2_fragment, + composite_score=score, + score_components={}, + ) + except (ValidationError, ValueError) as exc: + context.phase2_error = exc + + +@then("the ScoredFragment composite_score should be {expected:g}") +def step_scored_fragment_composite(context: Context, expected: float) -> None: + assert context.scored_fragment.composite_score == expected + + +@then('the ScoredFragment should have component "{name}" with value {val:g}') +def step_scored_fragment_component(context: Context, name: str, val: float) -> None: + assert context.scored_fragment.score_components[name] == val + + +@then("the ScoredFragment should reference the original fragment") +def step_scored_fragment_ref(context: Context) -> None: + assert context.scored_fragment.fragment is context.phase2_fragment + + +@then("a phase2 validation error should be raised") +def step_phase2_validation_error(context: Context) -> None: + assert context.phase2_error is not None + + +# --------------------------------------------------------------------------- +# Fragment list construction — Given steps +# --------------------------------------------------------------------------- + + +@given("the following phase2 fragments:") +def step_given_phase2_fragments(context: Context) -> None: + frags: list[ContextFragment] = [] + for row in context.table: + frags.append( + _make_phase2_fragment( + uko_node=row["uko_node"], + content=row["content"], + relevance_score=float(row["score"]), + token_count=int(row["tokens"]), + detail_depth=int(row["depth"]), + ) + ) + context.phase2_fragments = frags + + +@given("an empty phase2 fragment list") +def step_given_empty_phase2(context: Context) -> None: + context.phase2_fragments = [] + + +@given("a phase2 fragment with metadata:") +def step_given_phase2_fragment_with_metadata(context: Context) -> None: + row = context.table[0] + meta: dict[str, str] = {} + if "hierarchy_weight" in row.headings: + meta["hierarchy_weight"] = row["hierarchy_weight"] + if "strategy_quality" in row.headings: + meta["strategy_quality"] = row["strategy_quality"] + if "recency_bonus" in row.headings: + meta["recency_bonus"] = row["recency_bonus"] + + context.phase2_fragments = [ + _make_phase2_fragment( + uko_node=row["uko_node"], + content=row["content"], + relevance_score=float(row["score"]), + token_count=int(row["tokens"]), + detail_depth=int(row["depth"]), + metadata=meta, + ) + ] + + +@given("a phase2 budget with max_tokens {max_t:d} and reserved_tokens {res_t:d}") +def step_given_phase2_budget(context: Context, max_t: int, res_t: int) -> None: + context.phase2_budget = ContextBudget(max_tokens=max_t, reserved_tokens=res_t) + + +@given( + "scorer weights: relevance={rel:g}, hierarchy={hier:g}, " + "quality={qual:g}, recency={rec:g}" +) +def step_given_scorer_weights( + context: Context, rel: float, hier: float, qual: float, rec: float +) -> None: + context.scorer_weights = ScorerWeights( + relevance=rel, hierarchy=hier, quality=qual, recency=rec + ) + + +# --------------------------------------------------------------------------- +# ContentHashDeduplicator — When / Then +# --------------------------------------------------------------------------- + + +@when("I deduplicate the fragments") +def step_deduplicate(context: Context) -> None: + dedup = ContentHashDeduplicator() + context.phase2_result = list(dedup.deduplicate(context.phase2_fragments)) + + +@then("{count:d} fragment should remain after deduplication") +def step_dedup_count_singular(context: Context, count: int) -> None: + assert len(context.phase2_result) == count, ( + f"Expected {count}, got {len(context.phase2_result)}" + ) + + +@then("{count:d} fragments should remain after deduplication") +def step_dedup_count_plural(context: Context, count: int) -> None: + assert len(context.phase2_result) == count, ( + f"Expected {count}, got {len(context.phase2_result)}" + ) + + +@then("the surviving fragment should have relevance score {expected:g}") +def step_surviving_relevance(context: Context, expected: float) -> None: + assert len(context.phase2_result) >= 1 + assert context.phase2_result[0].relevance_score == expected + + +# --------------------------------------------------------------------------- +# MaxDepthResolver — When / Then +# --------------------------------------------------------------------------- + + +@when("I resolve depth conflicts") +def step_resolve_depth(context: Context) -> None: + resolver = MaxDepthResolver() + context.phase2_result = list(resolver.resolve(context.phase2_fragments)) + + +@then("{count:d} fragment should remain after depth resolution") +def step_depth_count_singular(context: Context, count: int) -> None: + assert len(context.phase2_result) == count, ( + f"Expected {count}, got {len(context.phase2_result)}" + ) + + +@then("{count:d} fragments should remain after depth resolution") +def step_depth_count_plural(context: Context, count: int) -> None: + assert len(context.phase2_result) == count, ( + f"Expected {count}, got {len(context.phase2_result)}" + ) + + +@then("the surviving fragment should have detail depth {expected:d}") +def step_surviving_depth(context: Context, expected: int) -> None: + assert len(context.phase2_result) >= 1 + assert context.phase2_result[0].detail_depth == expected + + +@then('the resolved fragment order should be "{c1}", "{c2}", "{c3}"') +def step_resolved_order(context: Context, c1: str, c2: str, c3: str) -> None: + contents = [f.content for f in context.phase2_result] + assert contents == [c1, c2, c3], f"Expected [{c1}, {c2}, {c3}], got {contents}" + + +# --------------------------------------------------------------------------- +# WeightedCompositeScorer — When / Then +# --------------------------------------------------------------------------- + + +@when("I score the fragments with default weights") +def step_score_default(context: Context) -> None: + scorer = WeightedCompositeScorer() + context.phase2_result = list(scorer.score(context.phase2_fragments)) + + +@when("I score the fragments with custom weights") +def step_score_custom(context: Context) -> None: + scorer = WeightedCompositeScorer(weights=context.scorer_weights) + context.phase2_result = list(scorer.score(context.phase2_fragments)) + + +@when("I score the fragments using score_detailed") +def step_score_detailed(context: Context) -> None: + scorer = WeightedCompositeScorer() + context.scored_fragment_list = scorer.score_detailed(context.phase2_fragments) + + +@then("the scored fragment should have a composite relevance score") +def step_scored_has_composite(context: Context) -> None: + assert len(context.phase2_result) >= 1 + score = context.phase2_result[0].relevance_score + assert 0.0 <= score <= 1.0 + + +@then('the scored fragment metadata should contain original relevance "{val}"') +def step_scored_original_relevance(context: Context, val: str) -> None: + meta = context.phase2_result[0].metadata + assert meta.get("_original_relevance") == val + + +@then("the scored fragment relevance score should be {expected:g}") +def step_scored_exact_relevance(context: Context, expected: float) -> None: + assert len(context.phase2_result) >= 1 + actual = context.phase2_result[0].relevance_score + assert abs(actual - expected) < 1e-4, f"Expected {expected}, got {actual}" + + +@then('the scored fragment metadata should contain component "{name}" near {val:g}') +def step_scored_component_near(context: Context, name: str, val: float) -> None: + meta = context.phase2_result[0].metadata + key = f"_score_{name}" + actual = float(meta[key]) + assert abs(actual - val) < 0.01, f"Expected ~{val}, got {actual}" + + +@then("{count:d} scored fragments should be returned") +def step_scored_count(context: Context, count: int) -> None: + assert len(context.phase2_result) == count + + +@then("the scored fragment relevance score should be at most {max_val:g}") +def step_scored_at_most(context: Context, max_val: float) -> None: + assert context.phase2_result[0].relevance_score <= max_val + + +@then("{count:d} ScoredFragment objects should be returned") +def step_scored_fragment_objects_count(context: Context, count: int) -> None: + assert len(context.scored_fragment_list) == count + + +@then("each ScoredFragment should have a composite_score between 0.0 and 1.0") +def step_scored_fragment_range(context: Context) -> None: + for sf in context.scored_fragment_list: + assert 0.0 <= sf.composite_score <= 1.0 + + +@then("each ScoredFragment should have score_components") +def step_scored_fragment_has_components(context: Context) -> None: + for sf in context.scored_fragment_list: + assert len(sf.score_components) > 0 + + +# --------------------------------------------------------------------------- +# GreedyKnapsackPacker — When / Then +# --------------------------------------------------------------------------- + + +@when("I pack the fragments into budget") +def step_pack(context: Context) -> None: + packer = GreedyKnapsackPacker() + context.phase2_result = list( + packer.pack(context.phase2_fragments, context.phase2_budget) + ) + + +@then("{count:d} fragment should be packed") +def step_packed_count_singular(context: Context, count: int) -> None: + assert len(context.phase2_result) == count, ( + f"Expected {count}, got {len(context.phase2_result)}" + ) + + +@then("{count:d} fragments should be packed") +def step_packed_count_plural(context: Context, count: int) -> None: + assert len(context.phase2_result) == count, ( + f"Expected {count}, got {len(context.phase2_result)}" + ) + + +@then("the packed total tokens should be at most {max_tokens:d}") +def step_packed_total_at_most(context: Context, max_tokens: int) -> None: + total = sum(f.token_count for f in context.phase2_result) + assert total <= max_tokens, f"Expected <= {max_tokens}, got {total}" + + +@then('the packed fragment uko_node should be "{expected}"') +def step_packed_single_node(context: Context, expected: str) -> None: + assert len(context.phase2_result) >= 1 + assert context.phase2_result[0].uko_node == expected + + +@then('the packed fragments should include uko_node "{expected}"') +def step_packed_includes_node(context: Context, expected: str) -> None: + nodes = {f.uko_node for f in context.phase2_result} + assert expected in nodes, f"Expected {expected} in {nodes}" + + +@then('the packed fragment for "{node}" should have depth {depth:d}') +def step_packed_fragment_depth(context: Context, node: str, depth: int) -> None: + for frag in context.phase2_result: + if frag.uko_node == node: + assert frag.detail_depth == depth, ( + f"Expected depth {depth} for {node}, got {frag.detail_depth}" + ) + return + msg = f"No fragment found for node {node}" + raise AssertionError(msg) + + +# --------------------------------------------------------------------------- +# Pipeline Integration — Given / When / Then +# --------------------------------------------------------------------------- + +_TEST_PLAN_ID = "01JQTESTPN00000000000000BB" + + +@given("the ACMS pipeline with a ContentHashDeduplicator") +def step_pipeline_with_dedup(context: Context) -> None: + context.phase2_pipeline = ACMSPipeline( + deduplicator=ContentHashDeduplicator(), + ) + + +@given("the ACMS pipeline with a MaxDepthResolver") +def step_pipeline_with_resolver(context: Context) -> None: + context.phase2_pipeline = ACMSPipeline( + depth_resolver=MaxDepthResolver(), + ) + + +@given("the ACMS pipeline with a WeightedCompositeScorer") +def step_pipeline_with_scorer(context: Context) -> None: + context.phase2_pipeline = ACMSPipeline( + scorer=WeightedCompositeScorer(), + ) + + +@given("the ACMS pipeline with a GreedyKnapsackPacker") +def step_pipeline_with_packer(context: Context) -> None: + context.phase2_pipeline = ACMSPipeline( + packer=GreedyKnapsackPacker(), + ) + + +@when("I assemble context through the pipeline") +def step_assemble_pipeline(context: Context) -> None: + budget = getattr( + context, "phase2_budget", ContextBudget(max_tokens=100000, reserved_tokens=0) + ) + payload = context.phase2_pipeline.assemble( + plan_id=_TEST_PLAN_ID, + fragments=context.phase2_fragments, + budget=budget, + ) + context.phase2_payload = payload + + +@when("I assemble context through the pipeline with budget") +def step_assemble_pipeline_with_budget(context: Context) -> None: + payload = context.phase2_pipeline.assemble( + plan_id=_TEST_PLAN_ID, + fragments=context.phase2_fragments, + budget=context.phase2_budget, + ) + context.phase2_payload = payload + + +@then("the pipeline output should contain {count:d} fragment") +def step_pipeline_fragment_count_singular(context: Context, count: int) -> None: + assert len(context.phase2_payload.fragments) == count, ( + f"Expected {count}, got {len(context.phase2_payload.fragments)}" + ) + + +@then("the pipeline output should contain {count:d} fragments") +def step_pipeline_fragment_count_plural(context: Context, count: int) -> None: + assert len(context.phase2_payload.fragments) == count + + +@then("the pipeline output fragments should have updated relevance scores") +def step_pipeline_scorer_applied(context: Context) -> None: + for frag in context.phase2_payload.fragments: + # Scorer adds _original_relevance metadata + assert "_original_relevance" in frag.metadata + + +@then("the pipeline output total tokens should be at most {max_tokens:d}") +def step_pipeline_total_at_most(context: Context, max_tokens: int) -> None: + assert context.phase2_payload.total_tokens <= max_tokens diff --git a/robot/acms_pipeline_phase2.robot b/robot/acms_pipeline_phase2.robot new file mode 100644 index 000000000..9d9befb2c --- /dev/null +++ b/robot/acms_pipeline_phase2.robot @@ -0,0 +1,57 @@ +*** Settings *** +Documentation Integration smoke tests for ACMS pipeline Phase 2 components +Resource ${CURDIR}/common.resource +Suite Setup Setup Test Environment +Suite Teardown Cleanup Test Environment + +*** Variables *** +${HELPER} ${CURDIR}/helper_acms_pipeline_phase2.py + +*** Test Cases *** +ContentHashDeduplicator Removes Duplicates + [Documentation] Deduplicate fragments with identical content, keeping highest scored + ${result}= Run Process ${PYTHON} ${HELPER} dedup-basic cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} phase2-dedup-ok + +MaxDepthResolver Resolves Depth Conflicts + [Documentation] Resolve same UKO node at multiple depths to maximum + ${result}= Run Process ${PYTHON} ${HELPER} depth-resolve cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} phase2-depth-ok + +WeightedCompositeScorer Scores Fragments + [Documentation] Score fragments with configurable weighted factors + ${result}= Run Process ${PYTHON} ${HELPER} scorer-default cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} phase2-scorer-ok + +GreedyKnapsackPacker Packs Within Budget + [Documentation] Pack fragments using greedy knapsack within token budget + ${result}= Run Process ${PYTHON} ${HELPER} packer-budget cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} phase2-packer-ok + +ScoredFragment Model Validation + [Documentation] Create ScoredFragment model and verify fields + ${result}= Run Process ${PYTHON} ${HELPER} scored-fragment cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} phase2-scored-fragment-ok + +Pipeline Integration With Phase2 Components + [Documentation] Inject all Phase 2 components into ACMSPipeline + ${result}= Run Process ${PYTHON} ${HELPER} pipeline-integration cwd=${WORKSPACE} + Log ${result.stdout} + Log ${result.stderr} + Should Be Equal As Integers ${result.rc} 0 + Should Contain ${result.stdout} phase2-pipeline-ok diff --git a/robot/helper_acms_pipeline_phase2.py b/robot/helper_acms_pipeline_phase2.py new file mode 100644 index 000000000..d6c6fa81d --- /dev/null +++ b/robot/helper_acms_pipeline_phase2.py @@ -0,0 +1,181 @@ +"""Robot Framework helper for ACMS pipeline Phase 2 integration tests. + +Provides a CLI-style interface for Robot to invoke Phase 2 pipeline +components and verify the results. Exit code 0 = success, 1 = failure. + +Usage: + python robot/helper_acms_pipeline_phase2.py dedup-basic + python robot/helper_acms_pipeline_phase2.py depth-resolve + python robot/helper_acms_pipeline_phase2.py scorer-default + python robot/helper_acms_pipeline_phase2.py packer-budget + python robot/helper_acms_pipeline_phase2.py scored-fragment + python robot/helper_acms_pipeline_phase2.py pipeline-integration +""" + +from __future__ import annotations + +import sys +from collections.abc import Callable +from pathlib import Path + +_SRC = str(Path(__file__).resolve().parents[1] / "src") +if _SRC not in sys.path: + sys.path.insert(0, _SRC) + +from cleveragents.application.services.acms_phase2 import ( # noqa: E402 + ContentHashDeduplicator, + GreedyKnapsackPacker, + MaxDepthResolver, + WeightedCompositeScorer, +) +from cleveragents.application.services.acms_service import ACMSPipeline # noqa: E402 +from cleveragents.domain.models.core.context_fragment import ( # noqa: E402 + ContextBudget, + ContextFragment, + FragmentProvenance, + ScoredFragment, +) + +_DEFAULT_PROV = FragmentProvenance(resource_uri="test://robot-phase2") +_PLAN_ID = "01JQTESTPN00000000000000CC" + + +def _frag( + uko_node: str = "test://default", + content: str = "test content", + score: float = 0.5, + tokens: int = 10, + depth: int = 0, +) -> ContextFragment: + return ContextFragment( + uko_node=uko_node, + content=content, + relevance_score=score, + token_count=tokens, + detail_depth=depth, + provenance=_DEFAULT_PROV, + ) + + +def _cmd_dedup_basic() -> int: + """Deduplicate identical fragments and keep highest scored.""" + dedup = ContentHashDeduplicator() + fragments = [ + _frag(uko_node="project://a.py", content="hello", score=0.5), + _frag(uko_node="project://a.py", content="hello", score=0.9), + _frag(uko_node="project://b.py", content="hello", score=0.7), + ] + result = dedup.deduplicate(fragments) + assert len(result) == 2, f"Expected 2, got {len(result)}" + # Find the fragment for a.py — should be the higher-scored one + a_frags = [f for f in result if f.uko_node == "project://a.py"] + assert len(a_frags) == 1 + assert a_frags[0].relevance_score == 0.9 + print("phase2-dedup-ok") + return 0 + + +def _cmd_depth_resolve() -> int: + """Resolve depth conflicts to max depth per node.""" + resolver = MaxDepthResolver() + fragments = [ + _frag(uko_node="project://a.py", content="brief", depth=2, tokens=10), + _frag(uko_node="project://a.py", content="full", depth=9, tokens=100), + _frag(uko_node="project://b.py", content="b", depth=5, tokens=50), + ] + result = resolver.resolve(fragments) + assert len(result) == 2, f"Expected 2, got {len(result)}" + a_frags = [f for f in result if f.uko_node == "project://a.py"] + assert a_frags[0].detail_depth == 9 + print("phase2-depth-ok") + return 0 + + +def _cmd_scorer_default() -> int: + """Score fragments with default weights.""" + scorer = WeightedCompositeScorer() + fragments = [ + _frag(uko_node="project://a.py", content="hello", score=0.8), + _frag(uko_node="project://b.py", content="world", score=0.4), + ] + result = scorer.score(fragments) + assert len(result) == 2 + for f in result: + assert 0.0 <= f.relevance_score <= 1.0 + assert "_original_relevance" in f.metadata + print("phase2-scorer-ok") + return 0 + + +def _cmd_packer_budget() -> int: + """Pack fragments into budget.""" + packer = GreedyKnapsackPacker() + fragments = [ + _frag(uko_node="project://a.py", content="alpha", score=0.9, tokens=100), + _frag(uko_node="project://b.py", content="beta", score=0.7, tokens=100), + _frag(uko_node="project://c.py", content="gamma", score=0.5, tokens=100), + ] + budget = ContextBudget(max_tokens=250, reserved_tokens=0) + result = packer.pack(fragments, budget) + assert len(result) == 2, f"Expected 2, got {len(result)}" + total = sum(f.token_count for f in result) + assert total <= 250, f"Total {total} exceeds budget 250" + print("phase2-packer-ok") + return 0 + + +def _cmd_scored_fragment() -> int: + """Create and verify ScoredFragment model.""" + frag = _frag(uko_node="project://a.py", content="test", score=0.8) + scored = ScoredFragment( + fragment=frag, + composite_score=0.85, + score_components={"relevance": 0.9, "hierarchy": 0.8}, + ) + assert scored.composite_score == 0.85 + assert scored.score_components["relevance"] == 0.9 + assert scored.fragment is frag + print("phase2-scored-fragment-ok") + return 0 + + +def _cmd_pipeline_integration() -> int: + """Inject Phase 2 components into ACMSPipeline.""" + pipeline = ACMSPipeline( + deduplicator=ContentHashDeduplicator(), + depth_resolver=MaxDepthResolver(), + scorer=WeightedCompositeScorer(), + packer=GreedyKnapsackPacker(), + ) + fragments = [ + _frag(uko_node="project://a.py", content="hello", score=0.9, tokens=50), + _frag(uko_node="project://a.py", content="hello", score=0.5, tokens=50), + _frag(uko_node="project://b.py", content="world", score=0.7, tokens=50), + ] + budget = ContextBudget(max_tokens=200, reserved_tokens=0) + payload = pipeline.assemble( + plan_id=_PLAN_ID, + fragments=fragments, + budget=budget, + ) + assert len(payload.fragments) >= 1 + assert payload.total_tokens <= 200 + print("phase2-pipeline-ok") + return 0 + + +_COMMANDS: dict[str, Callable[[], int]] = { + "dedup-basic": _cmd_dedup_basic, + "depth-resolve": _cmd_depth_resolve, + "scorer-default": _cmd_scorer_default, + "packer-budget": _cmd_packer_budget, + "scored-fragment": _cmd_scored_fragment, + "pipeline-integration": _cmd_pipeline_integration, +} + + +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) + sys.exit(_COMMANDS[sys.argv[1]]()) diff --git a/src/cleveragents/application/services/__init__.py b/src/cleveragents/application/services/__init__.py index 2947ab239..310f667ee 100644 --- a/src/cleveragents/application/services/__init__.py +++ b/src/cleveragents/application/services/__init__.py @@ -3,6 +3,13 @@ Contains service classes that orchestrate business operations. """ +from cleveragents.application.services.acms_phase2 import ( + ContentHashDeduplicator, + GreedyKnapsackPacker, + MaxDepthResolver, + ScorerWeights, + WeightedCompositeScorer, +) from cleveragents.application.services.autonomy_controller import ( AutonomyController, ) @@ -167,6 +174,7 @@ __all__ = [ "ConfigLevel", "ConfigService", "ContainerUnavailableError", + "ContentHashDeduplicator", "ContextFragment", "CorrectionService", "CrossPlanCorrectionService", @@ -186,7 +194,9 @@ __all__ = [ "DuplicateImportRule", "ExecutionEnvironmentResolver", "FileMergeOutcome", + "GreedyKnapsackPacker", "InvariantService", + "MaxDepthResolver", "MergeConflictError", "MissingImportRule", "MissingSymbolRule", @@ -199,6 +209,7 @@ __all__ = [ "ResolvedValue", "RuntimeExecuteActor", "RuntimeExecuteResult", + "ScorerWeights", "SemanticCheckResult", "SemanticEmbeddingStrategy", "SemanticRuleRegistry", @@ -233,6 +244,7 @@ __all__ = [ "ValidationResult", "ValidationRunner", "ValidationSummary", + "WeightedCompositeScorer", "create_default_registry", "enforce_permission", "get_default_permission_service", diff --git a/src/cleveragents/application/services/acms_phase2.py b/src/cleveragents/application/services/acms_phase2.py new file mode 100644 index 000000000..708f21ca0 --- /dev/null +++ b/src/cleveragents/application/services/acms_phase2.py @@ -0,0 +1,466 @@ +"""ACMS pipeline Phase 2 — Fragment Fusion component implementations. + +Provides the four production-grade Phase 2 components for the ACMS context +assembly pipeline, replacing the no-op defaults from ``acms_service.py``: + +4. **ContentHashDeduplicator** — Groups fragments by UKO node URI, hashes + content within groups to detect duplicates, retains the fragment with + the highest ``relevance_score``. +5. **MaxDepthResolver** — When the same UKO node appears at multiple detail + depths, retains the highest-depth rendering per node. +6. **WeightedCompositeScorer** — Computes a composite relevance score from + configurable weighted factors (relevance, hierarchy, quality, recency). +7. **GreedyKnapsackPacker** — Selects fragments that fit within the token + budget using a greedy knapsack algorithm with depth fallback. + +All components satisfy the v1 Protocol signatures defined in +``acms_service.py`` and can be injected into ``ACMSPipeline`` via +constructor dependency injection. + +Based on ``docs/specification.md`` §42930-42937. +""" + +from __future__ import annotations + +import hashlib +import logging +from collections import defaultdict +from collections.abc import Sequence +from typing import Final + +from cleveragents.domain.models.core.context_fragment import ( + ContextBudget, + ContextFragment, + ScoredFragment, +) + +logger = logging.getLogger(__name__) + +# --------------------------------------------------------------------------- +# Default configuration constants (from spec §43004-43012) +# --------------------------------------------------------------------------- + +DEFAULT_RELEVANCE_WEIGHT: Final[float] = 0.4 +DEFAULT_HIERARCHY_WEIGHT: Final[float] = 0.3 +DEFAULT_QUALITY_WEIGHT: Final[float] = 0.2 +DEFAULT_RECENCY_WEIGHT: Final[float] = 0.1 + +DEFAULT_DEPTH_FALLBACK_STEPS: Final[tuple[int, ...]] = (9, 4, 2, 0) +DEFAULT_MIN_FRAGMENT_TOKENS: Final[int] = 10 + + +# --------------------------------------------------------------------------- +# 4. ContentHashDeduplicator (FragmentDeduplicator protocol) +# --------------------------------------------------------------------------- + + +class ContentHashDeduplicator: + """Remove duplicate fragments using content hashing. + + Groups fragments by ``uko_node`` URI; within each group, computes a + SHA-256 hash of the content to detect duplicates. When duplicates are + found, the fragment with the highest ``relevance_score`` is retained. + + Implements ``FragmentDeduplicator`` protocol from ``acms_service.py``. + + Based on ``docs/specification.md`` §42934 — ``ContentHashDeduplicator``. + """ + + def deduplicate( + self, + fragments: Sequence[ContextFragment], + ) -> Sequence[ContextFragment]: + """Return deduplicated fragments, retaining highest-scored duplicates.""" + if not fragments: + return list(fragments) + + # Group by (uko_node, content_hash) + groups: dict[tuple[str, str], ContextFragment] = {} + dedup_count = 0 + + for frag in fragments: + content_hash = _content_hash(frag.content) + key = (frag.uko_node, content_hash) + + if key in groups: + dedup_count += 1 + existing = groups[key] + if frag.relevance_score > existing.relevance_score: + groups[key] = frag + else: + groups[key] = frag + + result = list(groups.values()) + if dedup_count > 0: + logger.info( + "Deduplicated fragments", + extra={ + "dedup_count": dedup_count, + "input_count": len(fragments), + "output_count": len(result), + }, + ) + return result + + +# --------------------------------------------------------------------------- +# 5. MaxDepthResolver (DetailDepthResolver protocol) +# --------------------------------------------------------------------------- + + +class MaxDepthResolver: + """Resolve detail depth conflicts by retaining the highest depth per node. + + When the same UKO node appears at multiple detail depths across strategy + outputs, retains the highest-depth rendering. Fragments for distinct + UKO nodes pass through unchanged. + + Implements ``DetailDepthResolver`` protocol from ``acms_service.py``. + + Based on ``docs/specification.md`` §42935 — ``MaxDepthResolver``. + """ + + def resolve( + self, + fragments: Sequence[ContextFragment], + ) -> Sequence[ContextFragment]: + """Return fragments with depth conflicts resolved to maximum depth.""" + if not fragments: + return list(fragments) + + # Group by uko_node, keep max depth per node + best_by_node: dict[str, ContextFragment] = {} + + for frag in fragments: + node = frag.uko_node + if node in best_by_node: + existing = best_by_node[node] + if frag.detail_depth > existing.detail_depth: + best_by_node[node] = frag + elif ( + frag.detail_depth == existing.detail_depth + and frag.relevance_score > existing.relevance_score + ): + # Same depth: prefer higher relevance + best_by_node[node] = frag + else: + best_by_node[node] = frag + + # Preserve original order for fragments that survived + survived_ids = {f.fragment_id for f in best_by_node.values()} + result = [f for f in fragments if f.fragment_id in survived_ids] + + resolved_count = len(fragments) - len(result) + if resolved_count > 0: + logger.info( + "Resolved depth conflicts", + extra={ + "resolved_count": resolved_count, + "input_count": len(fragments), + "output_count": len(result), + }, + ) + return result + + +# --------------------------------------------------------------------------- +# 6. WeightedCompositeScorer (FragmentScorer protocol) +# --------------------------------------------------------------------------- + + +class ScorerWeights: + """Configurable scoring weight parameters. + + Based on ``docs/specification.md`` §43004-43008. + """ + + __slots__ = ("hierarchy", "quality", "recency", "relevance") + + def __init__( + self, + *, + relevance: float = DEFAULT_RELEVANCE_WEIGHT, + hierarchy: float = DEFAULT_HIERARCHY_WEIGHT, + quality: float = DEFAULT_QUALITY_WEIGHT, + recency: float = DEFAULT_RECENCY_WEIGHT, + ) -> None: + self.relevance = relevance + self.hierarchy = hierarchy + self.quality = quality + self.recency = recency + + +class WeightedCompositeScorer: + """Score fragments with a composite of weighted factors. + + Computes:: + + composite = (relevance_score * relevance_weight + + hierarchy_weight_factor * hierarchy_weight + + strategy_quality * quality_weight + + recency_bonus * recency_weight) + + Factor values are sourced from fragment metadata when available, with + sensible defaults (0.5 for hierarchy, 0.5 for quality, 1.0 for recency). + + The composite is clamped to [0.0, 1.0] and stored in the fragment's + ``relevance_score`` field (since the v1 pipeline uses ``ContextFragment`` + throughout, not ``ScoredFragment``). The original relevance score and + component breakdown are preserved in ``metadata``. + + Implements ``FragmentScorer`` protocol from ``acms_service.py``. + + Based on ``docs/specification.md`` §42936 — ``WeightedCompositeScorer``. + """ + + def __init__(self, weights: ScorerWeights | None = None) -> None: + self._weights = weights or ScorerWeights() + + @property + def weights(self) -> ScorerWeights: + """Return the configured scoring weights.""" + return self._weights + + def score( + self, + fragments: Sequence[ContextFragment], + ) -> Sequence[ContextFragment]: + """Return fragments re-scored with composite relevance scores.""" + if not fragments: + return list(fragments) + + scored: list[ContextFragment] = [] + for frag in fragments: + scored_frag = self._score_fragment(frag) + scored.append(scored_frag) + + logger.info( + "Scored fragments", + extra={ + "count": len(scored), + "min_score": min(f.relevance_score for f in scored), + "max_score": max(f.relevance_score for f in scored), + }, + ) + return scored + + def score_detailed( + self, + fragments: Sequence[ContextFragment], + ) -> list[ScoredFragment]: + """Score fragments and return ``ScoredFragment`` wrappers. + + This is the spec-aligned interface that produces ``ScoredFragment`` + objects with full component breakdowns. Not used by the v1 pipeline + but available for callers that need detailed scoring. + """ + result: list[ScoredFragment] = [] + for frag in fragments: + components = self._compute_components(frag) + composite = self._composite_from_components(components) + result.append( + ScoredFragment( + fragment=frag, + composite_score=composite, + score_components=components, + ) + ) + return result + + def _score_fragment(self, frag: ContextFragment) -> ContextFragment: + """Compute composite score and return a new fragment with it.""" + components = self._compute_components(frag) + composite = self._composite_from_components(components) + + # Preserve component breakdown in metadata + new_meta = dict(frag.metadata) + new_meta["_original_relevance"] = str(frag.relevance_score) + for key, val in components.items(): + new_meta[f"_score_{key}"] = str(round(val, 6)) + + return ContextFragment( + fragment_id=frag.fragment_id, + uko_node=frag.uko_node, + content=frag.content, + detail_depth=frag.detail_depth, + token_count=frag.token_count, + relevance_score=composite, + provenance=frag.provenance, + tier=frag.tier, + metadata=new_meta, + created_at=frag.created_at, + ) + + def _compute_components(self, frag: ContextFragment) -> dict[str, float]: + """Extract individual scoring factor values from a fragment.""" + meta = frag.metadata + + relevance = frag.relevance_score + hierarchy = float(meta.get("hierarchy_weight", "0.5")) + quality = float(meta.get("strategy_quality", "0.5")) + recency = float(meta.get("recency_bonus", "1.0")) + + return { + "relevance": relevance, + "hierarchy": hierarchy, + "quality": quality, + "recency": recency, + } + + def _composite_from_components(self, components: dict[str, float]) -> float: + """Compute the weighted composite from component values.""" + w = self._weights + raw = ( + components["relevance"] * w.relevance + + components["hierarchy"] * w.hierarchy + + components["quality"] * w.quality + + components["recency"] * w.recency + ) + return max(0.0, min(1.0, round(raw, 6))) + + +# --------------------------------------------------------------------------- +# 7. GreedyKnapsackPacker (BudgetPacker protocol) +# --------------------------------------------------------------------------- + + +class GreedyKnapsackPacker: + """Pack fragments into the token budget using a greedy knapsack algorithm. + + Sorts fragments by ``relevance_score`` descending, then greedily adds + fragments until the budget is exhausted. Fragments below + ``min_fragment_tokens`` are excluded. + + When a high-value fragment does not fit at its current depth, the packer + attempts **depth fallback**: it looks for alternative renderings of the + same UKO node at progressively lower depths from the input set. This + allows the packer to include a lower-depth version of important content + when the full version is too large. + + Implements ``BudgetPacker`` protocol from ``acms_service.py``. + + Based on ``docs/specification.md`` §42937 — ``GreedyKnapsackPacker``. + """ + + def __init__( + self, + *, + depth_fallback_steps: tuple[int, ...] = DEFAULT_DEPTH_FALLBACK_STEPS, + min_fragment_tokens: int = DEFAULT_MIN_FRAGMENT_TOKENS, + ) -> None: + self._depth_fallback_steps = depth_fallback_steps + self._min_fragment_tokens = min_fragment_tokens + + @property + def depth_fallback_steps(self) -> tuple[int, ...]: + """Return the configured depth fallback sequence.""" + return self._depth_fallback_steps + + @property + def min_fragment_tokens(self) -> int: + """Return the minimum fragment token threshold.""" + return self._min_fragment_tokens + + def pack( + self, + fragments: Sequence[ContextFragment], + budget: ContextBudget, + ) -> Sequence[ContextFragment]: + """Pack fragments into budget using greedy knapsack with depth fallback.""" + if not fragments: + return list(fragments) + + available = budget.available_tokens + + # Pre-filter fragments below minimum token threshold + eligible = [f for f in fragments if f.token_count >= self._min_fragment_tokens] + + # Build a lookup for depth fallback: uko_node → {depth → fragment} + node_depth_map: dict[str, dict[int, ContextFragment]] = defaultdict(dict) + for frag in eligible: + node = frag.uko_node + depth = frag.detail_depth + # Keep highest-scored fragment per (node, depth) pair + if depth not in node_depth_map[node] or ( + frag.relevance_score > node_depth_map[node][depth].relevance_score + ): + node_depth_map[node][depth] = frag + + # Sort by relevance_score descending for greedy selection + sorted_frags = sorted(eligible, key=lambda f: f.relevance_score, reverse=True) + + packed: list[ContextFragment] = [] + used_tokens = 0 + packed_nodes: set[str] = set() # Track which nodes are already packed + fallback_count = 0 + + for frag in sorted_frags: + if frag.uko_node in packed_nodes: + # Node already packed (possibly via fallback) + continue + + if used_tokens + frag.token_count <= available: + # Fragment fits directly + packed.append(frag) + used_tokens += frag.token_count + packed_nodes.add(frag.uko_node) + else: + # Try depth fallback — look for lower-depth versions + fallback = self._try_depth_fallback( + frag.uko_node, + available - used_tokens, + node_depth_map, + ) + if fallback is not None: + packed.append(fallback) + used_tokens += fallback.token_count + packed_nodes.add(frag.uko_node) + fallback_count += 1 + + logger.info( + "Packed fragments into budget", + extra={ + "input_count": len(fragments), + "eligible_count": len(eligible), + "packed_count": len(packed), + "used_tokens": used_tokens, + "available_tokens": available, + "budget_utilization": round( + used_tokens / available if available > 0 else 0.0, 4 + ), + "fallback_count": fallback_count, + }, + ) + return packed + + def _try_depth_fallback( + self, + uko_node: str, + remaining_budget: int, + node_depth_map: dict[str, dict[int, ContextFragment]], + ) -> ContextFragment | None: + """Try progressively lower depths for a fragment that doesn't fit.""" + depths = node_depth_map.get(uko_node, {}) + if not depths: + return None + + for target_depth in self._depth_fallback_steps: + if target_depth in depths: + candidate = depths[target_depth] + if ( + candidate.token_count <= remaining_budget + and candidate.token_count >= self._min_fragment_tokens + ): + return candidate + + return None + + +# --------------------------------------------------------------------------- +# Internal helpers +# --------------------------------------------------------------------------- + + +def _content_hash(content: str) -> str: + """Compute SHA-256 hash of fragment content for deduplication.""" + return hashlib.sha256(content.encode("utf-8")).hexdigest() diff --git a/src/cleveragents/domain/models/core/__init__.py b/src/cleveragents/domain/models/core/__init__.py index 5bdded228..f87f6a38d 100644 --- a/src/cleveragents/domain/models/core/__init__.py +++ b/src/cleveragents/domain/models/core/__init__.py @@ -68,6 +68,7 @@ from cleveragents.domain.models.core.context_fragment import ( ContextFragment, ContextPayload, FragmentProvenance, + ScoredFragment, build_provenance_map, compute_context_hash, ) @@ -438,6 +439,7 @@ __all__ = [ "SafetyProfileProvenance", "SafetyProfileRef", "SandboxStrategy", + "ScoredFragment", "Session", "SessionExportError", "SessionImportError", diff --git a/src/cleveragents/domain/models/core/context_fragment.py b/src/cleveragents/domain/models/core/context_fragment.py index b90f70146..ab208188c 100644 --- a/src/cleveragents/domain/models/core/context_fragment.py +++ b/src/cleveragents/domain/models/core/context_fragment.py @@ -230,6 +230,33 @@ class ContextPayload(BaseModel, frozen=True): return self.total_tokens <= self.budget.available_tokens +class ScoredFragment(BaseModel, frozen=True): + """A ``ContextFragment`` annotated with a composite score. + + Produced by the ``FragmentScorer`` pipeline component and consumed by + the ``BudgetPacker``. The ``composite_score`` combines relevance, + hierarchy weight, strategy quality, and recency into a single ranking + value. ``score_components`` preserves the individual factor + contributions for observability. + + Based on ``docs/specification.md`` ~line 42825. + """ + + fragment: ContextFragment + composite_score: float = Field( + ge=0.0, + le=1.0, + description="Combined ranking score (0.0-1.0)", + ) + score_components: dict[str, float] = Field( + default_factory=dict, + description=( + 'Breakdown of scoring factors, e.g. {"relevance": 0.9, ' + '"hierarchy": 0.8, "quality": 0.7, "recency": 1.0}' + ), + ) + + def compute_context_hash(fragments: tuple[ContextFragment, ...]) -> str: """Compute a SHA-256 hash over the ordered fragment contents. diff --git a/vulture_whitelist.py b/vulture_whitelist.py index 486dc8e97..183c373e8 100644 --- a/vulture_whitelist.py +++ b/vulture_whitelist.py @@ -801,3 +801,14 @@ _word_density # noqa: B018, F821 _jaccard_similarity # noqa: B018, F821 _uri_segments # noqa: B018, F821 _max_proximity # noqa: B018, F821 + +# ACMS Pipeline Phase 2 — Fragment Fusion components (public API) +ContentHashDeduplicator # noqa: B018, F821 +MaxDepthResolver # noqa: B018, F821 +WeightedCompositeScorer # noqa: B018, F821 +GreedyKnapsackPacker # noqa: B018, F821 +ScorerWeights # noqa: B018, F821 +ScoredFragment # noqa: B018, F821 +score_detailed # noqa: B018, F821 +depth_fallback_steps # noqa: B018, F821 +min_fragment_tokens # noqa: B018, F821