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
« 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)."""
3from __future__ import annotations
5import asyncio
6import json
7from abc import ABC, abstractmethod
8from dataclasses import asdict, dataclass
9from typing import Any
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
21logger = LoggingConfig.get_logger(__name__)
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
29@dataclass
30class ChangeEvent:
31 """Webhook event enqueued for durable processing."""
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
41 def to_payload(self) -> dict[str, Any]:
42 return asdict(self)
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 )
57class QueueBackend(ABC):
58 """Abstract webhook event queue."""
60 @abstractmethod
61 async def enqueue(self, event: ChangeEvent) -> str:
62 """Enqueue an event. Returns a message/job id string."""
64 @abstractmethod
65 async def close(self) -> None:
66 """Release queue resources."""
69class SQLiteChangeEventQueue(QueueBackend):
70 """SQLite-backed durable queue using the WS-4.1 jobs table."""
72 def __init__(self, job_queue: SQLiteJobQueue):
73 self._job_queue = job_queue
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)
90 @property
91 def job_queue(self) -> SQLiteJobQueue:
92 return self._job_queue
94 async def close(self) -> None:
95 return None
98class QueueBackendManager:
99 """Factory and lifecycle for the webhook queue backend."""
101 _backend: QueueBackend | None = None
102 _job_queue: SQLiteJobQueue | None = None
103 _engine = None
104 _queue_db_op_lock: asyncio.Lock | None = None
106 @classmethod
107 async def initialize(cls) -> QueueBackend:
108 if cls._backend is not None:
109 return cls._backend
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)
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
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
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
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
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
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
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)