Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
"""add claimed_at to outbox_event

Revision ID: add_claimed_at_outbox_rev
Revises: outbox_event_table_rev
Create Date: 2026-08-20 00:00:00.000000

"""
from alembic import op
import sqlalchemy as sa


# revision identifiers, used by Alembic.
revision = "add_claimed_at_outbox_rev"
down_revision = "outbox_event_table_rev"
branch_labels = None
depends_on = None


def upgrade() -> None:
op.add_column(
"event_outbox",
sa.Column(
"claimed_at",
sa.DateTime(),
nullable=True,
),
)


def downgrade() -> None:
op.drop_column("event_outbox", "claimed_at")
5 changes: 5 additions & 0 deletions quantara/web_app/db/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -161,6 +161,10 @@ class Vault(Base):
DateTime, nullable=False, default=func.now(), onupdate=func.now()
)

__table_args__ = (
UniqueConstraint("user_id", "symbol", name="uq_vault_user_symbol"),
)


class TransactionStatus(PyEnum):
"""
Expand Down Expand Up @@ -243,6 +247,7 @@ class OutboxEvent(Base):
status = Column(String, nullable=False, default="pending") # pending, processing, processed, failed
retry_count = Column(Integer, nullable=False, default=0)
error_message = Column(String, nullable=True)
claimed_at = Column(DateTime, nullable=True)
created_at = Column(DateTime, nullable=False, default=func.now())
updated_at = Column(
DateTime, nullable=False, default=func.now(), onupdate=func.now()
Expand Down
77 changes: 58 additions & 19 deletions quantara/web_app/tasks/outbox_relay.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,11 @@
import os
import asyncio
import json
import uuid
import sentry_sdk
from datetime import datetime, timedelta
from celery import Celery
from sqlalchemy import and_, or_
from sqlalchemy.orm import Session
from web_app.db.database import SessionLocal, init_db
from web_app.db.models import OutboxEvent, Position, Status, Transaction, TransactionStatus
Expand All @@ -31,19 +33,35 @@
enable_utc=True,
)

STALE_PROCESSING_INTERVAL_MINUTES = 5


def _is_valid_uuid(value: str) -> bool:
try:
uuid.UUID(value)
return True
except ValueError:
return False


@celery_app.task(bind=True, max_retries=5, default_retry_delay=10)
def process_position_opened_task(self, event_id: str):
"""
Celery task that consumes the PositionOpened event from the outbox.
"""
logger.info("processing_position_opened_task_started", event_id=event_id)


if not _is_valid_uuid(event_id):
logger.error("invalid_event_id_format", event_id=event_id)
return

event_uuid = uuid.UUID(event_id)

init_db()
db: Session = SessionLocal()
try:
# 1. Fetch the outbox event
event = db.query(OutboxEvent).filter(OutboxEvent.id == event_id).first()
event = db.query(OutboxEvent).filter(OutboxEvent.id == event_uuid).first()
if not event:
logger.error("outbox_event_not_found", event_id=event_id)
return
Expand Down Expand Up @@ -107,15 +125,16 @@
except Exception as exc:
db.rollback()
logger.exception("outbox_event_processing_failed", event_id=event_id, error=str(exc))

# Update event status to failed and increment retry
try:
with SessionLocal() as fail_session:
evt = fail_session.query(OutboxEvent).filter(OutboxEvent.id == event_id).first()
evt = fail_session.query(OutboxEvent).filter(OutboxEvent.id == event_uuid).first()
if evt:
evt.status = "failed"
evt.retry_count += 1
evt.error_message = str(exc)
evt.claimed_at = None
fail_session.commit()
except Exception as update_err:
logger.error("failed_to_update_outbox_event_status", error=str(update_err))
Expand All @@ -133,15 +152,15 @@

def process_pending_events(self):
"""
Scans event_outbox for pending/failed events and publishes them to Celery.
Also flags events older than 24h with a Sentry warning.
Scans event_outbox for pending/failed/stale-processing events and
publishes them to Celery. Uses an atomic claim to prevent double-dispatch.
"""
logger.info("outbox_relay_scan_started")
db: Session = SessionLocal()
try:
# Check for events older than 24h that are not processed
cutoff_24h_naive = datetime.now() - timedelta(hours=24)

old_events = db.query(OutboxEvent).filter(
OutboxEvent.status != "processed",
OutboxEvent.created_at < cutoff_24h_naive
Expand All @@ -152,19 +171,44 @@
logger.warning("outbox_event_older_than_24h", event_id=str(event.id), created_at=str(event.created_at))
sentry_sdk.capture_message(msg, level="warning")

# Fetch pending or failed events
pending_events = db.query(OutboxEvent).filter(
# Reclaim stale "processing" events whose claim has expired
stale_cutoff = datetime.now() - timedelta(minutes=STALE_PROCESSING_INTERVAL_MINUTES)
reclaimed = db.query(OutboxEvent).filter(
OutboxEvent.status == "processing",
OutboxEvent.claimed_at.isnot(None),
OutboxEvent.claimed_at < stale_cutoff,
OutboxEvent.retry_count < self.max_retries,
).update(
{"status": "pending", "claimed_at": None},
synchronize_session="fetch",
)
if reclaimed:
logger.info("reclaimed_stale_processing_events", count=reclaimed)
db.commit()

# Fetch pending or failed events (now includes freshly reclaimed ones)
candidate_events = db.query(OutboxEvent).filter(
OutboxEvent.status.in_(["pending", "failed"]),
OutboxEvent.retry_count < self.max_retries
OutboxEvent.retry_count < self.max_retries,
).all()

if not pending_events:
if not candidate_events:
logger.info("no_pending_outbox_events")
return

for event in pending_events:
# Mark as processing
event.status = "processing"
for event in candidate_events:
# Atomic claim: only transition to processing if still pending/failed
claimed = db.query(OutboxEvent).filter(
OutboxEvent.id == event.id,
OutboxEvent.status.in_(["pending", "failed"]),
).update(
{"status": "processing", "claimed_at": datetime.now()},
synchronize_session="fetch",
)
if not claimed:
logger.info("outbox_event_already_claimed", event_id=str(event.id))
continue

db.commit()

# Publish to Celery
Expand All @@ -175,8 +219,3 @@
logger.exception("outbox_relay_scan_failed", error=str(e))
finally:
db.close()


if __name__ == "__main__":
relay = OutboxRelay()
relay.process_pending_events()
13 changes: 11 additions & 2 deletions quantara/web_app/tests/test_outbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -54,17 +54,26 @@ def test_relay_worker_dispatches_task():
mock_event.created_at = datetime.now()

mock_db = MagicMock()
# Query returns old events (empty) then pending events (mock_event)
mock_query = MagicMock()
mock_db.query.return_value = mock_query
mock_filter = MagicMock()
mock_query.filter.return_value = mock_filter
mock_filter.all.side_effect = [[], [mock_event]]

update_calls = []

def track_update(values, **kwargs):
update_calls.append(values)
for k, v in values.items():
setattr(mock_event, k, v)
return 1

mock_filter.update.side_effect = track_update

with patch("web_app.tasks.outbox_relay.SessionLocal", return_value=mock_db), \
patch("web_app.tasks.outbox_relay.process_position_opened_task.delay") as mock_delay, \
patch("web_app.tasks.outbox_relay.sentry_sdk.capture_message") as mock_sentry:

relay = OutboxRelay(max_retries=5)
relay.process_pending_events()

Expand Down
Loading
Loading