Spaces:
Sleeping
Sleeping
File size: 6,400 Bytes
116524e | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 | """DeduplicationManager — coordinates similarity detection and consolidation."""
from __future__ import annotations
import logging
from typing import TYPE_CHECKING, Any, Dict, List, Optional
from ..protocols.deduplication import DeduplicationConfig
from .detector import SimilarityDetector
from .operations import (
ConsolidationOperation,
ConsolidationOpType,
DeleteOp,
KeepOp,
MergeOp,
UpdateOp,
apply_consolidation_operations,
)
from .prompts import format_pair_for_logging, generate_similarity_report
if TYPE_CHECKING:
from ..core.skillbook import Skillbook
logger = logging.getLogger(__name__)
class DeduplicationManager:
"""Manages similarity detection and feeds info to SkillManager.
Coordinates:
1. Computing / updating embeddings for skills
2. Detecting similar skill pairs
3. Generating similarity reports for the SkillManager prompt
4. Parsing and applying consolidation operations
Satisfies :class:`DeduplicationManagerLike` via ``get_similarity_report``.
"""
def __init__(self, config: DeduplicationConfig | None = None) -> None:
self.config = config or DeduplicationConfig()
self.detector = SimilarityDetector(self.config)
# ------------------------------------------------------------------
# DeduplicationManagerLike interface
# ------------------------------------------------------------------
def get_similarity_report(self, skillbook: "Skillbook") -> Optional[str]:
"""Generate a similarity report for the SkillManager prompt.
Should be called **before** the SkillManager runs.
Returns:
Formatted report, or ``None`` if no similar pairs found or
deduplication is disabled.
"""
if not self.config.enabled:
return None
self.detector.ensure_embeddings(skillbook)
similar_pairs = self.detector.detect_similar_pairs(skillbook)
if len(similar_pairs) < self.config.min_pairs_to_report:
if similar_pairs:
logger.debug(
"Found %d similar pairs, below threshold of %d",
len(similar_pairs),
self.config.min_pairs_to_report,
)
return None
logger.info("Found %d similar skill pairs", len(similar_pairs))
for skill_a, skill_b, similarity in similar_pairs:
logger.debug(format_pair_for_logging(skill_a, skill_b, similarity))
return generate_similarity_report(similar_pairs)
# ------------------------------------------------------------------
# Consolidation operation parsing / application
# ------------------------------------------------------------------
def parse_consolidation_operations(
self, response_data: Dict[str, Any]
) -> List[ConsolidationOperation]:
"""Parse consolidation operations from SkillManager response data."""
operations: List[ConsolidationOperation] = []
raw_ops = response_data.get("consolidation_operations", [])
if not isinstance(raw_ops, list):
logger.warning("consolidation_operations is not a list")
return operations
for raw_op in raw_ops:
if not isinstance(raw_op, dict):
continue
raw_type = raw_op.get("type", "").upper()
try:
op_type = ConsolidationOpType(raw_type)
except ValueError:
logger.warning("Unknown consolidation operation type: %r", raw_type)
continue
try:
if op_type is ConsolidationOpType.MERGE:
operations.append(
MergeOp(
source_ids=raw_op.get("source_ids", []),
merged_content=raw_op.get("merged_content", ""),
keep_id=raw_op.get("keep_id", ""),
reasoning=raw_op.get("reasoning", ""),
)
)
elif op_type is ConsolidationOpType.DELETE:
operations.append(
DeleteOp(
skill_id=raw_op.get("skill_id", ""),
reasoning=raw_op.get("reasoning", ""),
)
)
elif op_type is ConsolidationOpType.KEEP:
operations.append(
KeepOp(
skill_ids=raw_op.get("skill_ids", []),
differentiation=raw_op.get("differentiation", ""),
reasoning=raw_op.get("reasoning", ""),
)
)
elif op_type is ConsolidationOpType.UPDATE:
operations.append(
UpdateOp(
skill_id=raw_op.get("skill_id", ""),
new_content=raw_op.get("new_content", ""),
reasoning=raw_op.get("reasoning", ""),
)
)
except Exception as e:
logger.warning(
"Failed to parse consolidation operation (%s): %s",
type(e).__name__,
e,
)
logger.info("Parsed %d consolidation operations", len(operations))
return operations
def apply_operations(
self,
operations: List[ConsolidationOperation],
skillbook: "Skillbook",
) -> None:
"""Apply consolidation operations to the skillbook."""
if not operations:
return
logger.info("Applying %d consolidation operations", len(operations))
apply_consolidation_operations(operations, skillbook)
def apply_operations_from_response(
self,
response_data: Dict[str, Any],
skillbook: "Skillbook",
) -> List[ConsolidationOperation]:
"""Parse and apply consolidation operations in one step."""
operations = self.parse_consolidation_operations(response_data)
self.apply_operations(operations, skillbook)
return operations
|