| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112 |
- 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
|