from __future__ import annotations from dataclasses import dataclass, field @dataclass class FakeRepository: active: object | None = None embeddings: dict[str, tuple[float, ...]] = field(default_factory=dict) publications: list[object] = field(default_factory=list) def load_active_snapshot(self, source_type, source_uid): return self.active def load_embedding(self, embedding_hash): return self.embeddings.get(embedding_hash) def activate(self, publication): self.active = publication.snapshot self.embeddings.update(publication.embeddings) self.publications.append(publication) class RecordingEmbedder: profile_key = "qwen:1024:test" def __init__(self, fail=False): self.calls = [] self.fail = fail def embed(self, texts): self.calls.extend(texts) if self.fail: raise RuntimeError("embedding unavailable") return [[float(index), 1.0] for index, _text in enumerate(texts)] def _snapshot(version: int, owner: str): from app.core.knowledge.point_builder import build_knowledge_snapshot return build_knowledge_snapshot( "DataFlow", { "uid": "01900000-0000-7000-8000-000000001001", "version": version, "name": "sync", "purpose": "sync customer data", "owner": owner, "business_domain_uid": "01900000-0000-7000-8000-000000001002", }, ) def test_sync_reuses_unchanged_embeddings_and_activates_once(): from app.core.knowledge.sync import KnowledgeSyncService repository = FakeRepository() first_embedder = RecordingEmbedder() service = KnowledgeSyncService(repository=repository, embedder=first_embedder) first_result = service.sync(_snapshot(1, "team-a")) assert first_result.status == "canonical_active" assert len(first_embedder.calls) == 3 second_embedder = RecordingEmbedder() service = KnowledgeSyncService(repository=repository, embedder=second_embedder) second_result = service.sync(_snapshot(2, "team-b")) assert second_result.status == "canonical_active" assert second_result.diff_counts == { "added": 0, "modified": 1, "deleted": 0, "unchanged": 2, } assert second_embedder.calls == ["team-b"] assert repository.active.source_revision == 2 assert len(repository.publications) == 2 def test_sync_failure_keeps_old_active_snapshot(): import pytest from app.core.knowledge.sync import KnowledgeSyncService repository = FakeRepository() KnowledgeSyncService(repository=repository, embedder=RecordingEmbedder()).sync( _snapshot(1, "team-a") ) with pytest.raises(RuntimeError, match="embedding unavailable"): KnowledgeSyncService( repository=repository, embedder=RecordingEmbedder(fail=True) ).sync(_snapshot(2, "team-b")) assert repository.active.source_revision == 1 assert len(repository.publications) == 1 def test_sync_same_revision_and_snapshot_is_idempotent(): from app.core.knowledge.sync import KnowledgeSyncService repository = FakeRepository() embedder = RecordingEmbedder() service = KnowledgeSyncService(repository=repository, embedder=embedder) snapshot = _snapshot(1, "team-a") service.sync(snapshot) result = service.sync(snapshot) assert result.status == "canonical_active" assert result.embedded_chunk_count == 0 assert result.reused_embedding_count == 3 assert len(repository.publications) == 1