Coverage for src/qdrant_loader_core/graph/schema/utils.py: 100%

39 statements  

« prev     ^ index     » next       coverage.py v7.15.0, created at 2026-07-20 10:12 +0000

1from __future__ import annotations 

2 

3import logging 

4 

5from ..models import CoreEdgeType, CoreNodeLabel 

6from .versioning import get_version_query, set_version_query 

7 

8logger = logging.getLogger(__name__) 

9 

10 

11async def _apply_v1(graph_store) -> None: 

12 await _ensure_indexes(graph_store) 

13 await apply(graph_store) 

14 

15 

16MIGRATIONS = { 

17 1: _apply_v1, 

18} 

19LATEST_VERSION = max(MIGRATIONS) 

20 

21 

22# ------------------------ 

23# ENTRY POINT 

24# ------------------------ 

25async def init_schema(graph_store): 

26 current_version = await _get_current_version(graph_store) 

27 

28 logger.info( 

29 "Current graph schema version=%s", 

30 current_version, 

31 ) 

32 

33 for version in range(current_version + 1, LATEST_VERSION + 1): 

34 logger.info( 

35 "Applying graph schema migration v%s", 

36 version, 

37 ) 

38 

39 await MIGRATIONS[version](graph_store) 

40 await _set_version(graph_store, version) 

41 

42 logger.info( 

43 "Graph schema migration v%s applied successfully", 

44 version, 

45 ) 

46 

47 

48# ------------------------ 

49# INDEXES 

50# ------------------------ 

51async def _ensure_indexes(graph_store): 

52 queries = [f"CREATE INDEX ON :{label.value}(id)" for label in CoreNodeLabel] 

53 

54 for query in queries: 

55 try: 

56 await graph_store.query_cypher(query, {}) 

57 except Exception as exc: 

58 if "already exists" in str(exc).lower(): 

59 continue 

60 raise 

61 

62 

63# ------------------------ 

64# VERSION 

65# ------------------------ 

66async def _get_current_version(graph_store) -> int: 

67 result = await graph_store.query_cypher( 

68 get_version_query(), 

69 {}, 

70 ) 

71 

72 if not result: 

73 return 0 

74 

75 return int(result[0][0]) 

76 

77 

78async def _set_version( 

79 graph_store, 

80 version: int, 

81): 

82 await graph_store.query_cypher( 

83 set_version_query(), 

84 {"version": version}, 

85 ) 

86 

87 

88async def apply(graph_store): 

89 """ 

90 Initialize canonical graph schema metadata. 

91 

92 These nodes are informational metadata used to 

93 document supported node labels and edge types. 

94 """ 

95 

96 for label in CoreNodeLabel: 

97 await graph_store.query_cypher( 

98 """ 

99 MERGE (:SchemaNodeType {name: $name}) 

100 """, 

101 {"name": label.value}, 

102 ) 

103 

104 for edge_type in CoreEdgeType: 

105 await graph_store.query_cypher( 

106 """ 

107 MERGE (:SchemaEdgeType {name: $name}) 

108 """, 

109 {"name": edge_type.value}, 

110 )