Multi-stage evidence-grounded inference architecture replacing the monolithic 9B model extraction pipeline. CPU-first specialist services handle routine extraction while the 9B vLLM model is preserved for semantic adjudication of ambiguous cases. Key components: - Capability-aware inference gateway (OpenAI-compatible + Ollama) - Endpoint registry with DB migrations and REST API - Sentence-aware document segmenter (property tests) - Deterministic financial parsing with offset integrity - Symbol resolution with ambiguity detection - Specialist service (GLiNER2, dynamic batching, K8s deployment) - Company-specific sentiment (FinBERT, calibration) - Retrieval-based novelty and duplicate detection - Confidence calibration pipeline - Deterministic routing engine (property tests) - 9B adjudication layer with VRAM gating - Stock-specific impact model (features, labels, baseline, trained) - Pipeline orchestrator (state machine, queues, leases, feature flags) - Bounded parallelism (async workers, semaphore, load shedding) - Observability (tracing, metrics, alerts) - Compatibility adapter (v3→v2 golden mapping tests) - Shadow/canary promotion framework - Active learning and fine-tuning pipeline Test results: 1,161 tests pass, ruff lint clean. All 282 spec tasks completed.
134 lines
4.2 KiB
Python
134 lines
4.2 KiB
Python
"""Queue definitions and routing for the v3 intelligence pipeline.
|
|
|
|
Provides fast-path, adjudication, persistence, and review queues with
|
|
backpressure and dead-letter support.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import enum
|
|
from dataclasses import dataclass, field
|
|
from datetime import datetime, timezone
|
|
from typing import Any
|
|
from uuid import UUID, uuid4
|
|
|
|
|
|
class QueueName(str, enum.Enum):
|
|
"""Named queues in the v3 pipeline topology."""
|
|
|
|
INCOMING = "intelligence.v3.incoming"
|
|
FAST_PATH = "intelligence.v3.fast"
|
|
ADJUDICATION = "intelligence.v3.adjudication"
|
|
PERSISTENCE = "intelligence.v3.persist"
|
|
REVIEW = "intelligence.v3.review"
|
|
DEAD_LETTER = "intelligence.v3.dead_letter"
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class QueueMessage:
|
|
"""Immutable message envelope for queue transport."""
|
|
|
|
message_id: UUID
|
|
queue: QueueName
|
|
run_id: UUID
|
|
document_id: str
|
|
payload: dict[str, Any]
|
|
enqueued_at: datetime
|
|
attempt: int = 0
|
|
idempotency_key: str = ""
|
|
priority: int = 0
|
|
|
|
@classmethod
|
|
def create(
|
|
cls,
|
|
queue: QueueName,
|
|
run_id: UUID,
|
|
document_id: str,
|
|
payload: dict[str, Any] | None = None,
|
|
priority: int = 0,
|
|
idempotency_key: str = "",
|
|
) -> QueueMessage:
|
|
return cls(
|
|
message_id=uuid4(),
|
|
queue=queue,
|
|
run_id=run_id,
|
|
document_id=document_id,
|
|
payload=payload or {},
|
|
enqueued_at=datetime.now(timezone.utc),
|
|
priority=priority,
|
|
idempotency_key=idempotency_key,
|
|
)
|
|
|
|
|
|
@dataclass
|
|
class QueueRouter:
|
|
"""In-memory queue router with backpressure and depth tracking.
|
|
|
|
In production, this would be backed by Redis lists or a dedicated
|
|
message broker. This implementation provides the queue routing logic
|
|
and depth-based backpressure for testing and single-process usage.
|
|
"""
|
|
|
|
max_depth: int = 1000
|
|
_queues: dict[QueueName, list[QueueMessage]] = field(default_factory=dict)
|
|
_processed_keys: set[str] = field(default_factory=set)
|
|
|
|
def __post_init__(self) -> None:
|
|
for q in QueueName:
|
|
if q not in self._queues:
|
|
self._queues[q] = []
|
|
|
|
def enqueue(self, message: QueueMessage) -> bool:
|
|
"""Add a message to its designated queue.
|
|
|
|
Returns False if backpressure is triggered (queue full) or
|
|
if the idempotency key was already processed.
|
|
"""
|
|
if message.idempotency_key and message.idempotency_key in self._processed_keys:
|
|
return False # Duplicate — idempotent reject
|
|
|
|
queue = self._queues.setdefault(message.queue, [])
|
|
if len(queue) >= self.max_depth:
|
|
return False # Backpressure
|
|
|
|
queue.append(message)
|
|
return True
|
|
|
|
def dequeue(self, queue: QueueName) -> QueueMessage | None:
|
|
"""Pop the next message from a queue (FIFO). Returns None if empty."""
|
|
q = self._queues.get(queue, [])
|
|
if not q:
|
|
return None
|
|
msg = q.pop(0)
|
|
if msg.idempotency_key:
|
|
self._processed_keys.add(msg.idempotency_key)
|
|
return msg
|
|
|
|
def depth(self, queue: QueueName) -> int:
|
|
"""Current depth of the given queue."""
|
|
return len(self._queues.get(queue, []))
|
|
|
|
def is_saturated(self, queue: QueueName) -> bool:
|
|
"""Whether the queue has reached max depth (backpressure active)."""
|
|
return self.depth(queue) >= self.max_depth
|
|
|
|
def move_to_dead_letter(self, message: QueueMessage) -> QueueMessage:
|
|
"""Move a failed message to the dead-letter queue."""
|
|
dlq_msg = QueueMessage(
|
|
message_id=uuid4(),
|
|
queue=QueueName.DEAD_LETTER,
|
|
run_id=message.run_id,
|
|
document_id=message.document_id,
|
|
payload={**message.payload, "original_queue": message.queue.value},
|
|
enqueued_at=datetime.now(timezone.utc),
|
|
attempt=message.attempt,
|
|
idempotency_key="", # DLQ messages get new identity
|
|
priority=message.priority,
|
|
)
|
|
self._queues.setdefault(QueueName.DEAD_LETTER, []).append(dlq_msg)
|
|
return dlq_msg
|
|
|
|
def total_depth(self) -> int:
|
|
"""Sum of all queue depths."""
|
|
return sum(len(q) for q in self._queues.values())
|