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.
This commit is contained in:
@@ -0,0 +1,42 @@
|
||||
"""V3 Pipeline Orchestrator — state machines, queues, leases, and feature flags.
|
||||
|
||||
Coordinates the multi-stage intelligence pipeline with explicit state transitions,
|
||||
idempotency keys, retry policies, dead-letter handling, and independent v2/v3
|
||||
routing behind feature flags.
|
||||
"""
|
||||
|
||||
from services.intelligence_pipeline_v3.orchestrator.feature_flags import (
|
||||
FeatureFlags,
|
||||
PipelineVersion,
|
||||
)
|
||||
from services.intelligence_pipeline_v3.orchestrator.leases import (
|
||||
Lease,
|
||||
LeaseExpiredError,
|
||||
LeaseManager,
|
||||
)
|
||||
from services.intelligence_pipeline_v3.orchestrator.queues import (
|
||||
QueueMessage,
|
||||
QueueName,
|
||||
QueueRouter,
|
||||
)
|
||||
from services.intelligence_pipeline_v3.orchestrator.state import (
|
||||
PipelineState,
|
||||
PipelineStateMachine,
|
||||
StageState,
|
||||
StateTransition,
|
||||
)
|
||||
|
||||
__all__ = [
|
||||
"FeatureFlags",
|
||||
"Lease",
|
||||
"LeaseExpiredError",
|
||||
"LeaseManager",
|
||||
"PipelineState",
|
||||
"PipelineStateMachine",
|
||||
"PipelineVersion",
|
||||
"QueueMessage",
|
||||
"QueueName",
|
||||
"QueueRouter",
|
||||
"StageState",
|
||||
"StateTransition",
|
||||
]
|
||||
@@ -0,0 +1,117 @@
|
||||
"""Feature flags for independent v2/v3 pipeline routing.
|
||||
|
||||
Supports per-agent, per-document-type, and percentage-based routing
|
||||
between pipeline versions. Both versions can run simultaneously.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import enum
|
||||
import hashlib
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Any
|
||||
from uuid import UUID
|
||||
|
||||
|
||||
class PipelineVersion(str, enum.Enum):
|
||||
"""Available pipeline versions."""
|
||||
|
||||
V2 = "v2"
|
||||
V3 = "v3"
|
||||
SHADOW = "shadow" # V3 runs alongside V2 but doesn't affect outputs
|
||||
|
||||
|
||||
@dataclass
|
||||
class FeatureFlags:
|
||||
"""Pipeline version routing with per-agent, per-document-type,
|
||||
and percentage-based controls.
|
||||
|
||||
Each flag can be overridden independently. The evaluation order:
|
||||
1. Agent-specific override (if set)
|
||||
2. Document-type override (if set)
|
||||
3. Percentage-based routing (deterministic by document_id)
|
||||
4. Default version
|
||||
"""
|
||||
|
||||
default_version: PipelineVersion = PipelineVersion.V2
|
||||
v3_enabled: bool = False
|
||||
shadow_enabled: bool = False
|
||||
v3_percentage: int = 0 # 0-100, percentage of documents routed to v3
|
||||
agent_overrides: dict[str, PipelineVersion] = field(default_factory=dict)
|
||||
document_type_overrides: dict[str, PipelineVersion] = field(
|
||||
default_factory=dict
|
||||
)
|
||||
excluded_document_types: set[str] = field(default_factory=set)
|
||||
|
||||
def resolve(
|
||||
self,
|
||||
document_id: str,
|
||||
agent_id: str | UUID | None = None,
|
||||
document_type: str | None = None,
|
||||
) -> PipelineVersion:
|
||||
"""Determine which pipeline version handles a document.
|
||||
|
||||
Resolution is deterministic for the same inputs.
|
||||
"""
|
||||
if not self.v3_enabled and not self.shadow_enabled:
|
||||
return PipelineVersion.V2
|
||||
|
||||
# Check excluded document types
|
||||
if document_type and document_type in self.excluded_document_types:
|
||||
return PipelineVersion.V2
|
||||
|
||||
# Agent-specific override
|
||||
agent_key = str(agent_id) if agent_id else None
|
||||
if agent_key and agent_key in self.agent_overrides:
|
||||
return self.agent_overrides[agent_key]
|
||||
|
||||
# Document-type override
|
||||
if document_type and document_type in self.document_type_overrides:
|
||||
return self.document_type_overrides[document_type]
|
||||
|
||||
# Shadow mode: run both
|
||||
if self.shadow_enabled:
|
||||
return PipelineVersion.SHADOW
|
||||
|
||||
# Percentage-based routing (deterministic hash)
|
||||
if self.v3_percentage > 0:
|
||||
bucket = self._hash_to_bucket(document_id)
|
||||
if bucket < self.v3_percentage:
|
||||
return PipelineVersion.V3
|
||||
|
||||
return self.default_version
|
||||
|
||||
def _hash_to_bucket(self, document_id: str) -> int:
|
||||
"""Deterministic hash to 0-99 bucket for percentage routing."""
|
||||
h = hashlib.sha256(document_id.encode()).hexdigest()
|
||||
return int(h[:8], 16) % 100
|
||||
|
||||
def is_v3_active(self) -> bool:
|
||||
"""Whether v3 processing is active in any form."""
|
||||
return self.v3_enabled or self.shadow_enabled or self.v3_percentage > 0
|
||||
|
||||
def set_agent_override(
|
||||
self, agent_id: str | UUID, version: PipelineVersion
|
||||
) -> None:
|
||||
"""Set a per-agent pipeline version override."""
|
||||
self.agent_overrides[str(agent_id)] = version
|
||||
|
||||
def clear_agent_override(self, agent_id: str | UUID) -> None:
|
||||
"""Remove a per-agent override."""
|
||||
self.agent_overrides.pop(str(agent_id), None)
|
||||
|
||||
def to_dict(self) -> dict[str, Any]:
|
||||
"""Serialize flags for API/config responses."""
|
||||
return {
|
||||
"default_version": self.default_version.value,
|
||||
"v3_enabled": self.v3_enabled,
|
||||
"shadow_enabled": self.shadow_enabled,
|
||||
"v3_percentage": self.v3_percentage,
|
||||
"agent_overrides": {
|
||||
k: v.value for k, v in self.agent_overrides.items()
|
||||
},
|
||||
"document_type_overrides": {
|
||||
k: v.value for k, v in self.document_type_overrides.items()
|
||||
},
|
||||
"excluded_document_types": list(self.excluded_document_types),
|
||||
}
|
||||
@@ -0,0 +1,149 @@
|
||||
"""Lease management for pipeline stage workers.
|
||||
|
||||
Leases ensure exactly-once processing semantics. A worker must acquire
|
||||
a lease before processing a stage. Expired leases allow re-processing
|
||||
by another worker.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from uuid import UUID, uuid4
|
||||
|
||||
|
||||
class LeaseExpiredError(Exception):
|
||||
"""Raised when an operation is attempted on an expired lease."""
|
||||
|
||||
def __init__(self, lease_id: UUID, expired_at: datetime) -> None:
|
||||
self.lease_id = lease_id
|
||||
self.expired_at = expired_at
|
||||
super().__init__(
|
||||
f"Lease {lease_id} expired at {expired_at.isoformat()}"
|
||||
)
|
||||
|
||||
|
||||
@dataclass
|
||||
class Lease:
|
||||
"""A time-bounded processing lease for a pipeline stage."""
|
||||
|
||||
lease_id: UUID
|
||||
run_id: UUID
|
||||
stage: str
|
||||
worker_id: str
|
||||
acquired_at: datetime
|
||||
expires_at: datetime
|
||||
released: bool = False
|
||||
renewed_count: int = 0
|
||||
|
||||
@property
|
||||
def is_expired(self) -> bool:
|
||||
"""Check if the lease has passed its expiry time."""
|
||||
return datetime.now(timezone.utc) >= self.expires_at
|
||||
|
||||
@property
|
||||
def is_active(self) -> bool:
|
||||
"""Check if the lease is currently active."""
|
||||
return not self.released and not self.is_expired
|
||||
|
||||
def renew(self, extension: timedelta) -> None:
|
||||
"""Extend the lease expiry.
|
||||
|
||||
Raises LeaseExpiredError if already expired.
|
||||
"""
|
||||
if self.is_expired:
|
||||
raise LeaseExpiredError(self.lease_id, self.expires_at)
|
||||
if self.released:
|
||||
raise LeaseExpiredError(self.lease_id, self.expires_at)
|
||||
self.expires_at = datetime.now(timezone.utc) + extension
|
||||
self.renewed_count += 1
|
||||
|
||||
def release(self) -> None:
|
||||
"""Mark the lease as released (work completed or abandoned)."""
|
||||
self.released = True
|
||||
|
||||
|
||||
@dataclass
|
||||
class LeaseManager:
|
||||
"""Manages leases for pipeline stage workers.
|
||||
|
||||
In production, this would use Redis or database-backed distributed locks.
|
||||
This implementation provides the lease lifecycle logic for testing.
|
||||
"""
|
||||
|
||||
default_ttl: timedelta = field(default_factory=lambda: timedelta(seconds=120))
|
||||
_active_leases: dict[tuple[UUID, str], Lease] = field(default_factory=dict)
|
||||
_all_leases: list[Lease] = field(default_factory=list)
|
||||
|
||||
def acquire(
|
||||
self,
|
||||
run_id: UUID,
|
||||
stage: str,
|
||||
worker_id: str,
|
||||
ttl: timedelta | None = None,
|
||||
) -> Lease | None:
|
||||
"""Attempt to acquire a lease for a (run_id, stage) pair.
|
||||
|
||||
Returns None if an active lease already exists for that pair.
|
||||
Expired leases are cleaned up and allow re-acquisition.
|
||||
"""
|
||||
key = (run_id, stage)
|
||||
existing = self._active_leases.get(key)
|
||||
|
||||
if existing is not None:
|
||||
if existing.is_active:
|
||||
return None # Already leased
|
||||
# Expired — clean up
|
||||
del self._active_leases[key]
|
||||
|
||||
lease = Lease(
|
||||
lease_id=uuid4(),
|
||||
run_id=run_id,
|
||||
stage=stage,
|
||||
worker_id=worker_id,
|
||||
acquired_at=datetime.now(timezone.utc),
|
||||
expires_at=datetime.now(timezone.utc) + (ttl or self.default_ttl),
|
||||
)
|
||||
self._active_leases[key] = lease
|
||||
self._all_leases.append(lease)
|
||||
return lease
|
||||
|
||||
def release(self, lease: Lease) -> None:
|
||||
"""Release a lease, making the slot available."""
|
||||
lease.release()
|
||||
key = (lease.run_id, lease.stage)
|
||||
if key in self._active_leases and self._active_leases[key] is lease:
|
||||
del self._active_leases[key]
|
||||
|
||||
def renew(self, lease: Lease, extension: timedelta | None = None) -> None:
|
||||
"""Renew an active lease. Raises LeaseExpiredError if expired."""
|
||||
lease.renew(extension or self.default_ttl)
|
||||
|
||||
def is_leased(self, run_id: UUID, stage: str) -> bool:
|
||||
"""Check if a (run_id, stage) pair has an active lease."""
|
||||
key = (run_id, stage)
|
||||
existing = self._active_leases.get(key)
|
||||
if existing is None:
|
||||
return False
|
||||
if not existing.is_active:
|
||||
del self._active_leases[key]
|
||||
return False
|
||||
return True
|
||||
|
||||
def active_count(self) -> int:
|
||||
"""Number of currently active leases."""
|
||||
# Clean up expired
|
||||
expired_keys = [
|
||||
k for k, v in self._active_leases.items() if not v.is_active
|
||||
]
|
||||
for k in expired_keys:
|
||||
del self._active_leases[k]
|
||||
return len(self._active_leases)
|
||||
|
||||
def get_expired(self) -> list[Lease]:
|
||||
"""Get all expired but unreleased leases (for recovery)."""
|
||||
return [
|
||||
lease
|
||||
for lease in self._active_leases.values()
|
||||
if lease.is_expired and not lease.released
|
||||
]
|
||||
@@ -0,0 +1,258 @@
|
||||
"""Bounded application parallelism for the v3 pipeline.
|
||||
|
||||
Provides async worker pools, specialist micro-batching, adjudicator
|
||||
semaphore with queue backpressure, and load-shedding rules that never
|
||||
drop safety-critical documents silently.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import enum
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Callable, Coroutine
|
||||
|
||||
|
||||
class LoadSheddingAction(str, enum.Enum):
|
||||
"""Actions when load shedding is triggered."""
|
||||
|
||||
QUEUE = "queue" # Re-queue for later processing
|
||||
REJECT = "reject" # Reject with error (non-safety-critical only)
|
||||
DEGRADE = "degrade" # Process with reduced quality (skip optional stages)
|
||||
|
||||
|
||||
class DocumentPriority(str, enum.Enum):
|
||||
"""Document priority classes for load shedding decisions."""
|
||||
|
||||
SAFETY_CRITICAL = "safety_critical" # Never silently dropped
|
||||
HIGH = "high"
|
||||
NORMAL = "normal"
|
||||
LOW = "low"
|
||||
|
||||
|
||||
@dataclass
|
||||
class WorkerPoolConfig:
|
||||
"""Configuration for an async worker pool."""
|
||||
|
||||
max_workers: int = 4
|
||||
batch_size: int = 8
|
||||
batch_timeout_ms: int = 100
|
||||
queue_max_depth: int = 500
|
||||
shed_threshold: float = 0.8 # Start shedding at 80% capacity
|
||||
|
||||
|
||||
@dataclass
|
||||
class WorkerStats:
|
||||
"""Runtime statistics for a worker pool."""
|
||||
|
||||
active_workers: int = 0
|
||||
queued_items: int = 0
|
||||
processed_total: int = 0
|
||||
shed_total: int = 0
|
||||
errors_total: int = 0
|
||||
avg_latency_ms: float = 0.0
|
||||
last_activity: datetime | None = None
|
||||
|
||||
|
||||
class AsyncWorkerPool:
|
||||
"""Configurable async worker pool with bounded concurrency.
|
||||
|
||||
Replaces the single sequential extraction loop with concurrent
|
||||
processing while respecting resource limits.
|
||||
"""
|
||||
|
||||
def __init__(self, config: WorkerPoolConfig | None = None) -> None:
|
||||
self.config = config or WorkerPoolConfig()
|
||||
self._semaphore = asyncio.Semaphore(self.config.max_workers)
|
||||
self._stats = WorkerStats()
|
||||
self._running = False
|
||||
self._tasks: set[asyncio.Task[Any]] = set()
|
||||
|
||||
@property
|
||||
def stats(self) -> WorkerStats:
|
||||
return self._stats
|
||||
|
||||
@property
|
||||
def is_running(self) -> bool:
|
||||
return self._running
|
||||
|
||||
@property
|
||||
def available_slots(self) -> int:
|
||||
"""Number of available worker slots."""
|
||||
return max(0, self.config.max_workers - self._stats.active_workers)
|
||||
|
||||
def should_shed_load(self) -> bool:
|
||||
"""Whether load shedding should be active."""
|
||||
if self.config.queue_max_depth <= 0:
|
||||
return False
|
||||
ratio = self._stats.queued_items / self.config.queue_max_depth
|
||||
return ratio >= self.config.shed_threshold
|
||||
|
||||
async def submit(
|
||||
self,
|
||||
coro_fn: Callable[..., Coroutine[Any, Any, Any]],
|
||||
*args: Any,
|
||||
document_id: str = "",
|
||||
priority: DocumentPriority = DocumentPriority.NORMAL,
|
||||
) -> LoadSheddingAction | None:
|
||||
"""Submit work to the pool.
|
||||
|
||||
Returns None on successful submission, or a LoadSheddingAction
|
||||
if load shedding was applied. Safety-critical documents are
|
||||
never silently rejected.
|
||||
"""
|
||||
if self.should_shed_load():
|
||||
if priority == DocumentPriority.SAFETY_CRITICAL:
|
||||
# Safety-critical: always queue, never shed
|
||||
pass
|
||||
elif priority == DocumentPriority.LOW:
|
||||
self._stats.shed_total += 1
|
||||
return LoadSheddingAction.REJECT
|
||||
else:
|
||||
self._stats.shed_total += 1
|
||||
return LoadSheddingAction.QUEUE
|
||||
|
||||
self._stats.queued_items += 1
|
||||
task = asyncio.create_task(self._run_with_semaphore(coro_fn, *args))
|
||||
self._tasks.add(task)
|
||||
task.add_done_callback(self._tasks.discard)
|
||||
return None
|
||||
|
||||
async def _run_with_semaphore(
|
||||
self,
|
||||
coro_fn: Callable[..., Coroutine[Any, Any, Any]],
|
||||
*args: Any,
|
||||
) -> Any:
|
||||
"""Execute work bounded by the semaphore."""
|
||||
async with self._semaphore:
|
||||
self._stats.active_workers += 1
|
||||
self._stats.queued_items = max(0, self._stats.queued_items - 1)
|
||||
start = datetime.now(timezone.utc)
|
||||
try:
|
||||
result = await coro_fn(*args)
|
||||
self._stats.processed_total += 1
|
||||
return result
|
||||
except Exception:
|
||||
self._stats.errors_total += 1
|
||||
raise
|
||||
finally:
|
||||
self._stats.active_workers -= 1
|
||||
elapsed = (
|
||||
datetime.now(timezone.utc) - start
|
||||
).total_seconds() * 1000
|
||||
# Rolling average
|
||||
n = self._stats.processed_total + self._stats.errors_total
|
||||
if n > 0:
|
||||
self._stats.avg_latency_ms = (
|
||||
self._stats.avg_latency_ms * (n - 1) + elapsed
|
||||
) / n
|
||||
self._stats.last_activity = datetime.now(timezone.utc)
|
||||
|
||||
async def start(self) -> None:
|
||||
"""Mark the pool as running."""
|
||||
self._running = True
|
||||
|
||||
async def shutdown(self, timeout: float = 30.0) -> None:
|
||||
"""Wait for all active tasks to complete."""
|
||||
self._running = False
|
||||
if self._tasks:
|
||||
await asyncio.wait(self._tasks, timeout=timeout)
|
||||
|
||||
|
||||
class AdjudicatorSemaphore:
|
||||
"""GPU-safe concurrency control for the 9B adjudicator.
|
||||
|
||||
Limits concurrent adjudication requests to match vLLM's max-num-seqs
|
||||
setting. Provides queue-depth monitoring and backpressure signaling.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
max_concurrent: int = 8,
|
||||
max_queued: int = 32,
|
||||
) -> None:
|
||||
self.max_concurrent = max_concurrent
|
||||
self.max_queued = max_queued
|
||||
self._semaphore = asyncio.Semaphore(max_concurrent)
|
||||
self._queued = 0
|
||||
self._active = 0
|
||||
self._total_processed = 0
|
||||
|
||||
@property
|
||||
def active_count(self) -> int:
|
||||
return self._active
|
||||
|
||||
@property
|
||||
def queued_count(self) -> int:
|
||||
return self._queued
|
||||
|
||||
@property
|
||||
def is_backpressured(self) -> bool:
|
||||
"""Whether the adjudicator queue is full."""
|
||||
return self._queued >= self.max_queued
|
||||
|
||||
async def acquire(self) -> bool:
|
||||
"""Acquire adjudicator access.
|
||||
|
||||
Returns False if backpressure prevents queuing.
|
||||
"""
|
||||
if self._queued >= self.max_queued:
|
||||
return False
|
||||
self._queued += 1
|
||||
await self._semaphore.acquire()
|
||||
self._queued -= 1
|
||||
self._active += 1
|
||||
return True
|
||||
|
||||
def release(self) -> None:
|
||||
"""Release adjudicator slot."""
|
||||
self._active -= 1
|
||||
self._total_processed += 1
|
||||
self._semaphore.release()
|
||||
|
||||
@property
|
||||
def utilization(self) -> float:
|
||||
"""Current GPU utilization fraction."""
|
||||
return self._active / self.max_concurrent if self.max_concurrent > 0 else 0.0
|
||||
|
||||
|
||||
@dataclass
|
||||
class MicroBatcher:
|
||||
"""Specialist micro-batching with configurable latency limits.
|
||||
|
||||
Accumulates items until batch_size is reached or timeout expires,
|
||||
then processes the batch together for efficiency.
|
||||
"""
|
||||
|
||||
batch_size: int = 16
|
||||
timeout_ms: int = 50
|
||||
_buffer: list[Any] = field(default_factory=list)
|
||||
_batch_count: int = 0
|
||||
|
||||
def add(self, item: Any) -> list[Any] | None:
|
||||
"""Add an item. Returns a full batch if ready, else None."""
|
||||
self._buffer.append(item)
|
||||
if len(self._buffer) >= self.batch_size:
|
||||
return self.flush()
|
||||
return None
|
||||
|
||||
def flush(self) -> list[Any]:
|
||||
"""Force-flush the current buffer as a batch."""
|
||||
batch = self._buffer[:]
|
||||
self._buffer.clear()
|
||||
if batch:
|
||||
self._batch_count += 1
|
||||
return batch
|
||||
|
||||
@property
|
||||
def pending_count(self) -> int:
|
||||
return len(self._buffer)
|
||||
|
||||
@property
|
||||
def total_batches(self) -> int:
|
||||
return self._batch_count
|
||||
|
||||
@property
|
||||
def is_empty(self) -> bool:
|
||||
return len(self._buffer) == 0
|
||||
@@ -0,0 +1,133 @@
|
||||
"""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())
|
||||
@@ -0,0 +1,187 @@
|
||||
"""Pipeline and stage state machines with explicit transitions and idempotency."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import enum
|
||||
import hashlib
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any
|
||||
from uuid import UUID, uuid4
|
||||
|
||||
|
||||
class PipelineState(str, enum.Enum):
|
||||
"""Top-level pipeline run states."""
|
||||
|
||||
PENDING = "pending"
|
||||
SEGMENTING = "segmenting"
|
||||
EXTRACTING = "extracting"
|
||||
RESOLVING = "resolving"
|
||||
VERIFYING = "verifying"
|
||||
ROUTING = "routing"
|
||||
ADJUDICATING = "adjudicating"
|
||||
IMPACT = "impact"
|
||||
PERSISTING = "persisting"
|
||||
COMPLETED = "completed"
|
||||
FAILED = "failed"
|
||||
DEAD_LETTER = "dead_letter"
|
||||
|
||||
|
||||
class StageState(str, enum.Enum):
|
||||
"""Per-stage execution states."""
|
||||
|
||||
QUEUED = "queued"
|
||||
LEASED = "leased"
|
||||
RUNNING = "running"
|
||||
SUCCEEDED = "succeeded"
|
||||
RETRYING = "retrying"
|
||||
FAILED = "failed"
|
||||
SKIPPED = "skipped"
|
||||
|
||||
|
||||
# Valid transitions for the pipeline state machine
|
||||
_PIPELINE_TRANSITIONS: dict[PipelineState, set[PipelineState]] = {
|
||||
PipelineState.PENDING: {PipelineState.SEGMENTING, PipelineState.FAILED},
|
||||
PipelineState.SEGMENTING: {PipelineState.EXTRACTING, PipelineState.FAILED},
|
||||
PipelineState.EXTRACTING: {PipelineState.RESOLVING, PipelineState.FAILED},
|
||||
PipelineState.RESOLVING: {PipelineState.VERIFYING, PipelineState.FAILED},
|
||||
PipelineState.VERIFYING: {PipelineState.ROUTING, PipelineState.FAILED},
|
||||
PipelineState.ROUTING: {
|
||||
PipelineState.ADJUDICATING,
|
||||
PipelineState.IMPACT,
|
||||
PipelineState.FAILED,
|
||||
},
|
||||
PipelineState.ADJUDICATING: {PipelineState.IMPACT, PipelineState.FAILED},
|
||||
PipelineState.IMPACT: {PipelineState.PERSISTING, PipelineState.FAILED},
|
||||
PipelineState.PERSISTING: {PipelineState.COMPLETED, PipelineState.FAILED},
|
||||
PipelineState.COMPLETED: set(),
|
||||
PipelineState.FAILED: {PipelineState.DEAD_LETTER, PipelineState.PENDING},
|
||||
PipelineState.DEAD_LETTER: set(),
|
||||
}
|
||||
|
||||
# Valid transitions for stage states
|
||||
_STAGE_TRANSITIONS: dict[StageState, set[StageState]] = {
|
||||
StageState.QUEUED: {StageState.LEASED, StageState.SKIPPED},
|
||||
StageState.LEASED: {StageState.RUNNING, StageState.QUEUED},
|
||||
StageState.RUNNING: {StageState.SUCCEEDED, StageState.RETRYING, StageState.FAILED},
|
||||
StageState.SUCCEEDED: set(),
|
||||
StageState.RETRYING: {StageState.QUEUED, StageState.FAILED},
|
||||
StageState.FAILED: set(),
|
||||
StageState.SKIPPED: set(),
|
||||
}
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class StateTransition:
|
||||
"""Immutable record of a state transition."""
|
||||
|
||||
transition_id: UUID
|
||||
run_id: UUID
|
||||
from_state: PipelineState | StageState
|
||||
to_state: PipelineState | StageState
|
||||
timestamp: datetime
|
||||
reason: str
|
||||
idempotency_key: str
|
||||
|
||||
|
||||
def _compute_idempotency_key(
|
||||
document_id: str, stage: str, attempt: int
|
||||
) -> str:
|
||||
"""Deterministic idempotency key from document, stage, and attempt."""
|
||||
raw = f"{document_id}:{stage}:{attempt}"
|
||||
return hashlib.sha256(raw.encode()).hexdigest()[:32]
|
||||
|
||||
|
||||
@dataclass
|
||||
class PipelineStateMachine:
|
||||
"""Manages state transitions for a single pipeline run.
|
||||
|
||||
Enforces valid transitions, records history, and generates
|
||||
idempotency keys for each stage attempt.
|
||||
"""
|
||||
|
||||
run_id: UUID = field(default_factory=uuid4)
|
||||
document_id: str = ""
|
||||
state: PipelineState = PipelineState.PENDING
|
||||
stage_states: dict[str, StageState] = field(default_factory=dict)
|
||||
stage_attempts: dict[str, int] = field(default_factory=dict)
|
||||
history: list[StateTransition] = field(default_factory=list)
|
||||
max_retries: int = 3
|
||||
metadata: dict[str, Any] = field(default_factory=dict)
|
||||
|
||||
def transition_pipeline(
|
||||
self, to_state: PipelineState, reason: str = ""
|
||||
) -> StateTransition:
|
||||
"""Advance the pipeline to a new state.
|
||||
|
||||
Raises ValueError if the transition is invalid.
|
||||
"""
|
||||
allowed = _PIPELINE_TRANSITIONS.get(self.state, set())
|
||||
if to_state not in allowed:
|
||||
raise ValueError(
|
||||
f"Invalid pipeline transition: {self.state.value} -> {to_state.value}"
|
||||
)
|
||||
|
||||
transition = StateTransition(
|
||||
transition_id=uuid4(),
|
||||
run_id=self.run_id,
|
||||
from_state=self.state,
|
||||
to_state=to_state,
|
||||
timestamp=datetime.now(timezone.utc),
|
||||
reason=reason,
|
||||
idempotency_key=_compute_idempotency_key(
|
||||
self.document_id, to_state.value, 0
|
||||
),
|
||||
)
|
||||
self.state = to_state
|
||||
self.history.append(transition)
|
||||
return transition
|
||||
|
||||
def transition_stage(
|
||||
self, stage: str, to_state: StageState, reason: str = ""
|
||||
) -> StateTransition:
|
||||
"""Advance a stage to a new state.
|
||||
|
||||
Raises ValueError if the transition is invalid.
|
||||
"""
|
||||
current = self.stage_states.get(stage, StageState.QUEUED)
|
||||
allowed = _STAGE_TRANSITIONS.get(current, set())
|
||||
if to_state not in allowed:
|
||||
raise ValueError(
|
||||
f"Invalid stage transition for '{stage}': "
|
||||
f"{current.value} -> {to_state.value}"
|
||||
)
|
||||
|
||||
attempt = self.stage_attempts.get(stage, 0)
|
||||
if to_state == StageState.RETRYING:
|
||||
attempt += 1
|
||||
self.stage_attempts[stage] = attempt
|
||||
|
||||
transition = StateTransition(
|
||||
transition_id=uuid4(),
|
||||
run_id=self.run_id,
|
||||
from_state=current,
|
||||
to_state=to_state,
|
||||
timestamp=datetime.now(timezone.utc),
|
||||
reason=reason,
|
||||
idempotency_key=_compute_idempotency_key(
|
||||
self.document_id, stage, attempt
|
||||
),
|
||||
)
|
||||
self.stage_states[stage] = to_state
|
||||
self.history.append(transition)
|
||||
return transition
|
||||
|
||||
def can_retry(self, stage: str) -> bool:
|
||||
"""Check whether the stage has retries remaining."""
|
||||
return self.stage_attempts.get(stage, 0) < self.max_retries
|
||||
|
||||
def should_dead_letter(self) -> bool:
|
||||
"""Check if the pipeline run should move to dead letter."""
|
||||
if self.state != PipelineState.FAILED:
|
||||
return False
|
||||
# Dead-letter if any stage exceeded max retries
|
||||
for stage, attempts in self.stage_attempts.items():
|
||||
if attempts >= self.max_retries:
|
||||
return True
|
||||
return False
|
||||
Reference in New Issue
Block a user