Coverage for src/lexigram/admin/integrations/search_sync.py: 91%
64 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-21 14:56 +0800
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-21 14:56 +0800
1from __future__ import annotations
3from typing import TYPE_CHECKING, Any
5from lexigram.logging import get_logger
7if TYPE_CHECKING:
8 from lexigram.admin.data.data_source import QueryResult
9 from lexigram.contracts.search import SearchableSpec
11_log = get_logger(__name__)
14class SearchSyncDataSourceWrapper:
15 """Wraps an IDataSource to sync CRUD operations to the search index.
17 Intercepts create/update/delete and calls the search engine after
18 the underlying data source operation succeeds. Errors from the
19 search engine are logged at DEBUG level and never propagated, so
20 a search backend outage does not break admin CRUD.
21 """
23 def __init__(self, inner: Any, search_engine: Any, spec: SearchableSpec) -> None:
24 self._inner = inner
25 self._search_engine = search_engine
26 self._index_name = spec.index_name
27 self._fields = spec.fields
29 async def find_one(self, item_id: Any) -> Any | None:
30 return await self._inner.find_one(item_id)
32 async def find_many(self, query: Any) -> QueryResult:
33 return await self._inner.find_many(query)
35 async def count(self, query: Any) -> int:
36 return await self._inner.count(query)
38 async def create(self, data: dict[str, Any]) -> Any:
39 entity = await self._inner.create(data)
40 await self._index_entity(entity)
41 return entity
43 async def update(self, item_id: Any, data: dict[str, Any]) -> Any:
44 entity = await self._inner.update(item_id, data)
45 await self._index_entity(entity)
46 return entity
48 async def delete(self, item_id: Any) -> bool:
49 ok = await self._inner.delete(item_id)
50 if ok:
51 await self._remove_document(item_id)
52 return ok
54 async def bulk_create(self, items: list[dict[str, Any]]) -> list[Any]:
55 return await self._inner.bulk_create(items)
57 async def bulk_update(self, ids: list[Any], data: dict[str, Any]) -> int:
58 return await self._inner.bulk_update(ids, data)
60 async def bulk_delete(self, ids: list[Any]) -> int:
61 return await self._inner.bulk_delete(ids)
63 async def _index_entity(self, entity: Any) -> None:
64 """Upsert *entity* into the search index."""
65 doc = self._to_document(entity)
66 if doc is None or not self._index_name:
67 return
68 try:
69 result = await self._search_engine.index(self._index_name, [doc])
70 if hasattr(result, "is_err") and result.is_err():
71 _log.debug("search.index_failed", index=self._index_name)
72 except Exception:
73 _log.debug("search.index_error", index=self._index_name, exc_info=True)
75 async def _remove_document(self, item_id: Any) -> None:
76 """Delete the document with *item_id* from the search index."""
77 if not self._index_name:
78 return
79 try:
80 result = await self._search_engine.delete(self._index_name, str(item_id))
81 if hasattr(result, "is_err") and result.is_err():
82 _log.debug("search.delete_failed", index=self._index_name)
83 except Exception:
84 _log.debug("search.delete_error", index=self._index_name, exc_info=True)
86 def _to_document(self, entity: Any) -> dict[str, Any] | None:
87 """Build the search document dict from *entity*.
89 Handles both ``dict`` and model object entities. The document
90 always includes an ``id`` key extracted from the entity, plus
91 the fields declared in ``self._fields``.
93 Returns ``None`` when no ID can be extracted.
94 """
95 if isinstance(entity, dict):
96 doc_id = entity.get("id")
97 doc = {f: entity.get(f) for f in self._fields}
98 else:
99 doc_id = getattr(entity, "id", None)
100 doc = {f: getattr(entity, f, None) for f in self._fields}
101 if doc_id is None:
102 return None
103 doc["id"] = str(doc_id)
104 return doc