Coverage for src/qdrant_loader/cli/commands/serve_cmd.py: 33%

136 statements  

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

1""" 

2qdrant-loader serve CLI command 

3Wires config, queue, worker pool, scheduler, HTTP server. Graceful SIGTERM/SIGINT. 

4""" 

5 

6from __future__ import annotations 

7 

8import asyncio 

9import signal 

10from pathlib import Path 

11 

12import click 

13from click.types import Choice 

14from click.types import Path as ClickPath 

15 

16from qdrant_loader.config.workspace import validate_workspace_flags 

17 

18# Max time to wait for background tasks (scheduler, workers, uvicorn) to stop 

19# gracefully on shutdown before cancelling them outright. Webhook requests are 

20# enqueue-only (they don't block on ingestion), so this only needs to cover 

21# uvicorn's own connection drain; kept short to leave headroom for 

22# state_manager.dispose() before an orchestrator's SIGKILL grace period expires. 

23SHUTDOWN_TIMEOUT_SECONDS = 10 

24 

25 

26@click.command( 

27 "serve", help="Run the qdrant-loader service (scheduler, workers, HTTP server)" 

28) 

29@click.option( 

30 "--workspace", 

31 type=ClickPath(path_type=Path), 

32 help="Workspace directory containing config.yaml and .env files.", 

33) 

34@click.option( 

35 "--config", 

36 type=ClickPath(exists=True, path_type=Path), 

37 help="Path to config file.", 

38) 

39@click.option( 

40 "--env", 

41 type=ClickPath(exists=True, path_type=Path), 

42 help="Path to .env file.", 

43) 

44@click.option( 

45 "--log-level", 

46 type=Choice( 

47 ["DEBUG", "INFO", "WARNING", "ERROR", "CRITICAL"], case_sensitive=False 

48 ), 

49 default="INFO", 

50 help="Set the logging level.", 

51) 

52@click.option( 

53 "--host", 

54 default="127.0.0.1", 

55 show_default=True, 

56 help="Webhook HTTP server host.", 

57) 

58@click.option( 

59 "--port", 

60 default=8081, 

61 type=int, 

62 show_default=True, 

63 help="Webhook HTTP server port.", 

64) 

65def serve_cmd( 

66 workspace: Path | None, 

67 config: Path | None, 

68 env: Path | None, 

69 log_level: str, 

70 host: str, 

71 port: int, 

72): 

73 """Run the qdrant-loader service (scheduler + worker pool + webhook HTTP server).""" 

74 asyncio.run(_serve_main(workspace, config, env, log_level, host, port)) 

75 

76 

77async def _serve_main( 

78 workspace: Path | None, 

79 config: Path | None, 

80 env: Path | None, 

81 log_level: str, 

82 host: str = "127.0.0.1", 

83 port: int = 8081, 

84): 

85 import uvicorn 

86 

87 from qdrant_loader.cli.config_loader import ( 

88 load_config_with_workspace, 

89 setup_workspace, 

90 ) 

91 from qdrant_loader.config import get_global_config, get_settings 

92 from qdrant_loader.core.pipeline.config import PipelineConfig 

93 from qdrant_loader.core.pipeline.factory import PipelineComponentsFactory 

94 from qdrant_loader.core.pipeline.orchestrator import PipelineOrchestrator 

95 from qdrant_loader.core.project_manager import ProjectManager 

96 from qdrant_loader.core.qdrant_manager import QdrantManager 

97 from qdrant_loader.core.state.state_manager import StateManager 

98 from qdrant_loader.core.worker.handlers import IngestionJobHandler 

99 from qdrant_loader.core.worker.job_types import JobType 

100 from qdrant_loader.core.worker.pool import QueueWorkerPool 

101 from qdrant_loader.core.worker.queue import SQLiteJobQueue 

102 from qdrant_loader.core.worker.scheduler import IncrementalPullScheduler 

103 from qdrant_loader.utils.logging import LoggingConfig 

104 from qdrant_loader.webhooks.auth import ( 

105 WEBHOOK_AUTH_NOT_CONFIGURED_MESSAGE, 

106 webhook_auth_configured, 

107 ) 

108 from qdrant_loader.webhooks.queue_backend import ( 

109 QueueBackendManager, 

110 SQLiteChangeEventQueue, 

111 ) 

112 from qdrant_loader.webhooks.server import app as webhook_app 

113 

114 # Validate flag combinations 

115 validate_workspace_flags(workspace, config, env) 

116 

117 # Setup workspace / logging 

118 workspace_config = None 

119 if workspace: 

120 workspace_config = setup_workspace(workspace) 

121 

122 log_file = ( 

123 str(workspace_config.logs_path / "serve.log") 

124 if workspace_config 

125 else "qdrant-loader.log" 

126 ) 

127 LoggingConfig.setup(level=log_level, format="console", file=log_file) 

128 logger = LoggingConfig.get_logger(__name__) 

129 

130 # Load configuration (required before get_global_config / get_settings) 

131 load_config_with_workspace(workspace_config, config, env) 

132 

133 # Fail closed on missing webhook auth BEFORE any DB engine/uvicorn/asyncio 

134 # resources are created, so misconfiguration exits cleanly instead of 

135 # surfacing as an unhandled SystemExit deep inside uvicorn's lifespan 

136 # (which otherwise races with StateManager teardown and leaves a dangling 

137 # aiosqlite task at interpreter shutdown). 

138 if not webhook_auth_configured(): 

139 raise click.ClickException(WEBHOOK_AUTH_NOT_CONFIGURED_MESSAGE) 

140 

141 # Setup graceful shutdown 

142 stop_event = asyncio.Event() 

143 

144 def _handle_signal(signum, _frame): 

145 logger.info("serve.signal_received", signum=signum) 

146 stop_event.set() 

147 

148 signal.signal(signal.SIGINT, _handle_signal) 

149 signal.signal(signal.SIGTERM, _handle_signal) 

150 

151 logger.info("serve.config_loaded") 

152 config_obj = get_global_config() 

153 settings = get_settings() 

154 

155 state_manager = None 

156 logger.info("serve.state_manager_init") 

157 state_manager = StateManager(config_obj.state_management) 

158 await state_manager.initialize() 

159 

160 logger.info("serve.queue_init") 

161 session_factory = state_manager.session_factory 

162 job_queue = SQLiteJobQueue( 

163 session_factory, 

164 db_op_lock=state_manager.queue_db_op_lock, 

165 ) 

166 

167 logger.info("serve.qdrant_init") 

168 qdrant_manager = QdrantManager(settings) 

169 

170 logger.info("serve.pipeline_init") 

171 pipeline_factory = PipelineComponentsFactory() 

172 concurrency = config_obj.concurrency 

173 pipeline_config = PipelineConfig( 

174 max_chunk_workers=concurrency.max_chunk_workers, 

175 max_embed_workers=concurrency.max_embed_workers, 

176 max_upsert_workers=concurrency.max_upsert_workers, 

177 queue_size=concurrency.queue_size, 

178 upsert_batch_size=concurrency.upsert_batch_size, 

179 ) 

180 pipeline_components = pipeline_factory.create_components( 

181 settings, 

182 pipeline_config, 

183 qdrant_manager, 

184 state_manager=state_manager, 

185 ) 

186 

187 logger.info("serve.project_manager_init") 

188 project_manager = ProjectManager( 

189 projects_config=settings.projects_config, 

190 global_collection_name=settings.global_config.qdrant.collection_name, 

191 ) 

192 async with session_factory() as session: 

193 await project_manager.initialize(session) 

194 

195 orchestrator = PipelineOrchestrator(settings, pipeline_components, project_manager) 

196 

197 logger.info("serve.handler_init") 

198 

199 job_handler = IngestionJobHandler( 

200 orchestrator=orchestrator, 

201 session_factory=session_factory, 

202 ) 

203 

204 worker_runtime = config_obj.workers.runtime 

205 logger.info( 

206 "serve.pool_init", 

207 worker_count=worker_runtime.worker_count, 

208 lease_seconds=worker_runtime.lease_seconds, 

209 max_attempts=worker_runtime.max_attempts, 

210 retry_backoff_base_seconds=worker_runtime.retry_backoff_base_seconds, 

211 ) 

212 worker_pool = QueueWorkerPool( 

213 queue=job_queue, 

214 handler=job_handler, 

215 worker_count=worker_runtime.worker_count, 

216 lease_seconds=worker_runtime.lease_seconds, 

217 max_attempts=worker_runtime.max_attempts, 

218 retry_backoff_base_seconds=worker_runtime.retry_backoff_base_seconds, 

219 job_types=[JobType.BULK_INGEST.value, JobType.INCREMENTAL_PULL.value], 

220 ) 

221 

222 logger.info("serve.scheduler_init") 

223 schedule = config_obj.workers.schedules.incremental_pull 

224 scheduler = IncrementalPullScheduler( 

225 queue=job_queue, 

226 projects_config=settings.projects_config, 

227 schedule=schedule, 

228 ) 

229 

230 # Share the already-initialised job_queue with QueueBackendManager so that: 

231 # 1. The webhook HTTP layer enqueues jobs on the same SQLiteJobQueue instance. 

232 # 2. job_queue.notify() fires immediately when /ingest enqueues a BULK_INGEST 

233 # job, waking the worker pool without waiting for the lease timeout. 

234 QueueBackendManager.set_backend(SQLiteChangeEventQueue(job_queue), job_queue) 

235 

236 logger.info("serve.starting") 

237 

238 # --- HTTP server (webhook endpoints) --- 

239 http_server_config = uvicorn.Config( 

240 webhook_app, 

241 host=host, 

242 port=port, 

243 log_level=log_level.lower(), 

244 ) 

245 http_server = uvicorn.Server(http_server_config) 

246 # serve_cmd owns signal handling; prevent uvicorn from overriding SIGINT/SIGTERM. 

247 http_server.install_signal_handlers = lambda: None 

248 

249 logger.info("serve.http_server_init", host=host, port=port) 

250 

251 async def scheduler_task(): 

252 await scheduler.run(stop_event) 

253 # When scheduler is disabled it returns immediately; keep this task 

254 # alive so asyncio.wait() only exits on stop_event or a real failure. 

255 await stop_event.wait() 

256 

257 async def worker_pool_task(): 

258 """Drain queue on job arrival or visibility timeout reclaim (every ~60s lease). 

259 

260 Instead of polling every 1s, await job_queue.notify() which signals when: 

261 - A new job is enqueued (from scheduler or trigger). 

262 - A job is released for retry. 

263 

264 Timeout (lease_seconds) ensures expired RUNNING jobs get reclaimed. 

265 """ 

266 pending_event = job_queue.notify() 

267 lease_seconds = worker_runtime.lease_seconds 

268 

269 while not stop_event.is_set(): 

270 processed = await worker_pool.run_until_empty() 

271 

272 # If no jobs were processed, wait for notification (with timeout for reclaim) 

273 if processed == 0: 

274 pending_event.clear() 

275 try: 

276 await asyncio.wait_for( 

277 pending_event.wait(), 

278 timeout=lease_seconds, 

279 ) 

280 except TimeoutError: 

281 # Timeout triggers visibility timeout reclaim: 

282 # claim_next() will pick up expired RUNNING jobs 

283 pass 

284 

285 scheduler_runner = None 

286 worker_runner = None 

287 stop_waiter = None 

288 http_runner = None 

289 try: 

290 scheduler_runner = asyncio.create_task(scheduler_task()) 

291 worker_runner = asyncio.create_task(worker_pool_task()) 

292 stop_waiter = asyncio.create_task(stop_event.wait()) 

293 http_runner = asyncio.create_task(http_server.serve()) 

294 

295 done, _ = await asyncio.wait( 

296 {scheduler_runner, worker_runner, stop_waiter, http_runner}, 

297 return_when=asyncio.FIRST_COMPLETED, 

298 ) 

299 

300 # Check for background task failures 

301 for task in (scheduler_runner, worker_runner, http_runner): 

302 if task in done and (exc := task.exception()) is not None: 

303 logger.error( 

304 "serve.background_task_failed", 

305 error=str(exc), 

306 error_type=type(exc).__name__, 

307 ) 

308 raise exc 

309 except Exception as exc: 

310 logger.error("serve.loop_error", error=str(exc)) 

311 raise 

312 finally: 

313 logger.info("serve.shutting_down") 

314 # Signal uvicorn to stop gracefully so its lifespan cleanup (webhook worker 

315 # teardown, QueueBackendManager.shutdown) runs before we dispose StateManager. 

316 http_server.should_exit = True 

317 for task in (scheduler_runner, worker_runner, stop_waiter): 

318 if task is not None: 

319 task.cancel() 

320 remaining_tasks = [ 

321 t 

322 for t in (scheduler_runner, worker_runner, stop_waiter, http_runner) 

323 if t is not None 

324 ] 

325 try: 

326 await asyncio.wait_for( 

327 asyncio.gather(*remaining_tasks, return_exceptions=True), 

328 timeout=SHUTDOWN_TIMEOUT_SECONDS, 

329 ) 

330 except TimeoutError: 

331 # uvicorn didn't honour should_exit (e.g. stuck request or lifespan 

332 # hook); cancel everything outright so the process can still exit. 

333 logger.warning("serve.shutdown_timeout", timeout=SHUTDOWN_TIMEOUT_SECONDS) 

334 for task in remaining_tasks: 

335 task.cancel() 

336 await asyncio.gather(*remaining_tasks, return_exceptions=True) 

337 if state_manager is not None: 

338 await state_manager.dispose() 

339 logger.info("serve.shutdown_complete")