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
« 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"""
6from __future__ import annotations
8import asyncio
9import signal
10from pathlib import Path
12import click
13from click.types import Choice
14from click.types import Path as ClickPath
16from qdrant_loader.config.workspace import validate_workspace_flags
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
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))
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
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
114 # Validate flag combinations
115 validate_workspace_flags(workspace, config, env)
117 # Setup workspace / logging
118 workspace_config = None
119 if workspace:
120 workspace_config = setup_workspace(workspace)
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__)
130 # Load configuration (required before get_global_config / get_settings)
131 load_config_with_workspace(workspace_config, config, env)
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)
141 # Setup graceful shutdown
142 stop_event = asyncio.Event()
144 def _handle_signal(signum, _frame):
145 logger.info("serve.signal_received", signum=signum)
146 stop_event.set()
148 signal.signal(signal.SIGINT, _handle_signal)
149 signal.signal(signal.SIGTERM, _handle_signal)
151 logger.info("serve.config_loaded")
152 config_obj = get_global_config()
153 settings = get_settings()
155 state_manager = None
156 logger.info("serve.state_manager_init")
157 state_manager = StateManager(config_obj.state_management)
158 await state_manager.initialize()
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 )
167 logger.info("serve.qdrant_init")
168 qdrant_manager = QdrantManager(settings)
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 )
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)
195 orchestrator = PipelineOrchestrator(settings, pipeline_components, project_manager)
197 logger.info("serve.handler_init")
199 job_handler = IngestionJobHandler(
200 orchestrator=orchestrator,
201 session_factory=session_factory,
202 )
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 )
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 )
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)
236 logger.info("serve.starting")
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
249 logger.info("serve.http_server_init", host=host, port=port)
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()
257 async def worker_pool_task():
258 """Drain queue on job arrival or visibility timeout reclaim (every ~60s lease).
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.
264 Timeout (lease_seconds) ensures expired RUNNING jobs get reclaimed.
265 """
266 pending_event = job_queue.notify()
267 lease_seconds = worker_runtime.lease_seconds
269 while not stop_event.is_set():
270 processed = await worker_pool.run_until_empty()
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
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())
295 done, _ = await asyncio.wait(
296 {scheduler_runner, worker_runner, stop_waiter, http_runner},
297 return_when=asyncio.FIRST_COMPLETED,
298 )
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")