Files
stonks-oracle/services/intelligence_pipeline_v3/orchestrator/queues.py
T
Celes Renata a72f336ad1 feat: Intelligence Pipeline v3 — full implementation
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.
2026-07-13 02:14:59 +00:00

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())