Coverage for src/qdrant_loader/webhooks/queue_backend.py: 88%

97 statements  

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

1"""Persistent queue backend for webhook events (WS-4).""" 

2 

3from __future__ import annotations 

4 

5import asyncio 

6import json 

7from abc import ABC, abstractmethod 

8from dataclasses import asdict, dataclass 

9from typing import Any 

10 

11from qdrant_loader.config import get_settings 

12from qdrant_loader.core.state.session import ( 

13 create_tables, 

14 dispose_engine, 

15 initialize_engine_and_session, 

16) 

17from qdrant_loader.core.worker.job_types import JobType 

18from qdrant_loader.core.worker.queue import SQLiteJobQueue 

19from qdrant_loader.utils.logging import LoggingConfig 

20 

21logger = LoggingConfig.get_logger(__name__) 

22 

23# Re-export for backward compatibility 

24SINGLE_UPSERT = JobType.SINGLE_UPSERT.value 

25SINGLE_DELETE = JobType.SINGLE_DELETE.value 

26FULL_SCAN = "FULL_SCAN" # Not yet in JobType; reserved for future use 

27 

28 

29@dataclass 

30class ChangeEvent: 

31 """Webhook event enqueued for durable processing.""" 

32 

33 source: str 

34 source_type: str | None 

35 project_id: str | None 

36 operation: str 

37 entity_id: str | None = None 

38 payload: Any = None 

39 force: bool = False 

40 

41 def to_payload(self) -> dict[str, Any]: 

42 return asdict(self) 

43 

44 @classmethod 

45 def from_payload(cls, data: dict[str, Any]) -> ChangeEvent: 

46 return cls( 

47 source=data["source"], 

48 source_type=data["source_type"], 

49 project_id=data.get("project_id"), 

50 operation=data["operation"], 

51 entity_id=data.get("entity_id"), 

52 payload=data.get("payload"), 

53 force=bool(data.get("force", False)), 

54 ) 

55 

56 

57class QueueBackend(ABC): 

58 """Abstract webhook event queue.""" 

59 

60 @abstractmethod 

61 async def enqueue(self, event: ChangeEvent) -> str: 

62 """Enqueue an event. Returns a message/job id string.""" 

63 

64 @abstractmethod 

65 async def close(self) -> None: 

66 """Release queue resources.""" 

67 

68 

69class SQLiteChangeEventQueue(QueueBackend): 

70 """SQLite-backed durable queue using the WS-4.1 jobs table.""" 

71 

72 def __init__(self, job_queue: SQLiteJobQueue): 

73 self._job_queue = job_queue 

74 

75 async def enqueue(self, event: ChangeEvent) -> str: 

76 job = await self._job_queue.enqueue( 

77 event.operation, 

78 event.to_payload(), 

79 ) 

80 logger.info( 

81 "Enqueued webhook event", 

82 job_id=job.id, 

83 operation=event.operation, 

84 source_type=event.source_type, 

85 source=event.source, 

86 entity_id=event.entity_id, 

87 ) 

88 return str(job.id) 

89 

90 @property 

91 def job_queue(self) -> SQLiteJobQueue: 

92 return self._job_queue 

93 

94 async def close(self) -> None: 

95 return None 

96 

97 

98class QueueBackendManager: 

99 """Factory and lifecycle for the webhook queue backend.""" 

100 

101 _backend: QueueBackend | None = None 

102 _job_queue: SQLiteJobQueue | None = None 

103 _engine = None 

104 _queue_db_op_lock: asyncio.Lock | None = None 

105 

106 @classmethod 

107 async def initialize(cls) -> QueueBackend: 

108 if cls._backend is not None: 

109 return cls._backend 

110 

111 settings = get_settings() 

112 state_config = settings.global_config.state_management 

113 engine, session_factory = initialize_engine_and_session(state_config) 

114 await create_tables(engine) 

115 

116 cls._engine = engine 

117 if cls._queue_db_op_lock is None: 

118 cls._queue_db_op_lock = asyncio.Lock() 

119 cls._job_queue = SQLiteJobQueue( 

120 session_factory, 

121 db_op_lock=cls._queue_db_op_lock, 

122 ) 

123 cls._backend = SQLiteChangeEventQueue(cls._job_queue) 

124 logger.info( 

125 "Initialized persistent webhook queue", 

126 database_path=state_config.database_path, 

127 ) 

128 return cls._backend 

129 

130 @classmethod 

131 def get_backend(cls) -> QueueBackend: 

132 if cls._backend is None: 

133 raise RuntimeError( 

134 "Webhook queue is not initialized. Call QueueBackendManager.initialize() " 

135 "during server startup." 

136 ) 

137 return cls._backend 

138 

139 @classmethod 

140 def get_job_queue(cls) -> SQLiteJobQueue: 

141 if cls._job_queue is None: 

142 raise RuntimeError("Webhook job queue is not initialized.") 

143 return cls._job_queue 

144 

145 @classmethod 

146 async def shutdown(cls) -> None: 

147 if cls._engine is not None: 

148 await dispose_engine(cls._engine) 

149 cls._engine = None 

150 cls._queue_db_op_lock = None 

151 cls._job_queue = None 

152 cls._backend = None 

153 

154 @classmethod 

155 def set_backend( 

156 cls, backend: QueueBackend, job_queue: SQLiteJobQueue | None = None 

157 ) -> None: 

158 """Override backend for testing.""" 

159 cls._backend = backend 

160 cls._job_queue = job_queue 

161 

162 @classmethod 

163 def reset(cls) -> None: 

164 """Reset manager state (for tests).""" 

165 cls._engine = None 

166 cls._queue_db_op_lock = None 

167 cls._job_queue = None 

168 cls._backend = None 

169 

170 

171def parse_job_payload(job) -> ChangeEvent: 

172 """Deserialize a Job row into a ChangeEvent.""" 

173 data = json.loads(job.payload_json) 

174 return ChangeEvent.from_payload(data)