Coverage for src/lexigram/admin/integrations/search_sync.py: 0%

64 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-24 23:18 +0800

1from __future__ import annotations 

2 

3from typing import TYPE_CHECKING, Any 

4 

5from lexigram.logging import get_logger 

6 

7if TYPE_CHECKING: 

8 from lexigram.admin.data.data_source import QueryResult 

9 from lexigram.contracts.search import SearchableSpec 

10 

11_log = get_logger(__name__) 

12 

13 

14class SearchSyncDataSourceWrapper: 

15 """Wraps an IDataSource to sync CRUD operations to the search index. 

16 

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

22 

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 

28 

29 async def find_one(self, item_id: Any) -> Any | None: 

30 return await self._inner.find_one(item_id) 

31 

32 async def find_many(self, query: Any) -> QueryResult: 

33 return await self._inner.find_many(query) 

34 

35 async def count(self, query: Any) -> int: 

36 return await self._inner.count(query) 

37 

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 

42 

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 

47 

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 

53 

54 async def bulk_create(self, items: list[dict[str, Any]]) -> list[Any]: 

55 return await self._inner.bulk_create(items) 

56 

57 async def bulk_update(self, ids: list[Any], data: dict[str, Any]) -> int: 

58 return await self._inner.bulk_update(ids, data) 

59 

60 async def bulk_delete(self, ids: list[Any]) -> int: 

61 return await self._inner.bulk_delete(ids) 

62 

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) 

74 

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) 

85 

86 def _to_document(self, entity: Any) -> dict[str, Any] | None: 

87 """Build the search document dict from *entity*. 

88 

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``. 

92 

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