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