Files
leocrm/tests/test_outbox.py
Agent Zero 0260f3410d Phase 7: Plugin-Gate, Event-Envelope, Pro-Handler Outbox-Verarbeitung
7.1 Plugin-Gate korrigiert:
- require_active_plugin nutzt current_user fuer tenant_id statt current_setting()
- Keine neue DB-Session mehr — nutzt bestehende get_db Dependency
- Fail-closed bei Fehlern

7.4 Einheitlicher Event-Envelope:
- Sauberes Envelope mit event_id, event_name, tenant_id, aggregate_type, aggregate_id, occurred_at, correlation_id, schema_version, data
- Keine _-Praefixe mehr im payload
- Handler empfangen envelope statt rohes payload

7.6 Verarbeitung pro Handler:
- Globaler consumer_inbox Check entfernt
- Pro-Handler Idempotency: outbox_deliveries pruefen ob Handler bereits erfolgreich
- Bereits erfolgreiche Handler werden uebersprungen
- consumer_inbox pro Handler geschrieben

7.7 no_handlers: Bereits implementiert (terminaler Status)
7.8 Cron-Jobs: Bereits mit Redis SET NX Locking implementiert

Tests: 23/23 Outbox-Tests bestanden
2026-08-03 15:20:06 +02:00

215 lines
6.9 KiB
Python

"""Tests for the transactional outbox pattern.
Covers:
- enqueue_outbox_event inserts rows with status='pending'
- process_outbox_batch publishes events to the in-process bus
- Retry logic with exponential backoff
- Max attempts → 'failed' status
- Empty batch returns 0
"""
from __future__ import annotations
import uuid
import pytest
from sqlalchemy import text
from sqlalchemy.ext.asyncio import async_sessionmaker, AsyncSession
from app.core.event_bus import get_event_bus
from app.core.outbox import enqueue_outbox_event, process_outbox_batch
@pytest.mark.asyncio
async def test_enqueue_outbox_event_inserts_pending_row(db_session):
"""enqueue_outbox_event inserts a row with status='pending'."""
tenant_id = uuid.uuid4()
await enqueue_outbox_event(
db_session, tenant_id, "contact.created",
{"contact_id": "abc-123", "tenant_id": str(tenant_id)},
)
await db_session.flush()
rows = (
await db_session.execute(
text("SELECT event_name, status, payload FROM event_outbox WHERE tenant_id = :tid"),
{"tid": str(tenant_id)},
)
).fetchall()
assert len(rows) == 1
assert rows[0][0] == "contact.created"
assert rows[0][1] == "pending"
assert rows[0][2]["contact_id"] == "abc-123"
@pytest.mark.asyncio
async def test_process_outbox_batch_publishes_events(
db_session, session_factory: async_sessionmaker[AsyncSession],
):
"""process_outbox_batch publishes pending events and marks them 'published'."""
tenant_id = uuid.uuid4()
received_events: list[tuple[str, dict]] = []
async def _handler(payload: dict) -> None:
received_events.append(("test.event", payload))
bus = get_event_bus()
bus.subscribe("test.event", _handler)
try:
await enqueue_outbox_event(
db_session, tenant_id, "test.event",
{"key": "value"},
)
await db_session.flush()
await db_session.commit()
# Use a separate session to simulate the worker
async with session_factory() as worker_session:
count = await process_outbox_batch(worker_session, batch_size=10, tenant_ids=[tenant_id])
assert count == 1
assert len(received_events) == 1
assert received_events[0][1]["data"]["key"] == "value"
# Verify the event is marked as published
rows = (
await db_session.execute(
text("SELECT status FROM event_outbox WHERE tenant_id = :tid"),
{"tid": str(tenant_id)},
)
).fetchall()
assert rows[0][0] == "published"
finally:
bus.unsubscribe("test.event", _handler)
@pytest.mark.asyncio
async def test_process_outbox_batch_empty_returns_zero(
session_factory: async_sessionmaker[AsyncSession],
):
"""process_outbox_batch returns 0 when no pending events exist."""
async with session_factory() as worker_session:
count = await process_outbox_batch(worker_session, batch_size=10)
assert count == 0
@pytest.mark.asyncio
async def test_process_outbox_batch_retry_on_failure(
db_session, session_factory: async_sessionmaker[AsyncSession],
):
"""When a handler raises, the event is retried with exponential backoff."""
tenant_id = uuid.uuid4()
async def _failing_handler(payload: dict) -> None:
raise RuntimeError("Handler failure")
bus = get_event_bus()
bus.subscribe("test.failing", _failing_handler)
try:
await enqueue_outbox_event(
db_session, tenant_id, "test.failing",
{"attempt": 1},
)
await db_session.flush()
await db_session.commit()
async with session_factory() as worker_session:
count = await process_outbox_batch(worker_session, batch_size=10, tenant_ids=[tenant_id])
assert count == 0 # nothing was successfully published
# Verify the event is back to 'pending' with attempts=1 and a retry time
rows = (
await db_session.execute(
text(
"SELECT status, attempts, next_retry_at "
"FROM event_outbox WHERE tenant_id = :tid"
),
{"tid": str(tenant_id)},
)
).fetchall()
assert rows[0][0] == "pending"
assert rows[0][1] == 1
assert rows[0][2] is not None
finally:
bus.unsubscribe("test.failing", _failing_handler)
@pytest.mark.asyncio
async def test_process_outbox_batch_max_attempts_marks_failed(
db_session, session_factory: async_sessionmaker[AsyncSession],
):
"""After max_attempts failures, the event is marked as 'failed'."""
tenant_id = uuid.uuid4()
async def _always_fails(payload: dict) -> None:
raise RuntimeError("Always fails")
bus = get_event_bus()
bus.subscribe("test.maxfail", _always_fails)
try:
# Insert an event that already has attempts = max_attempts - 1
await db_session.execute(
text(
"INSERT INTO event_outbox (tenant_id, event_name, payload, status, attempts, max_attempts) "
"VALUES (:tid, 'test.maxfail', CAST(:payload AS JSONB), 'pending', 4, 5)"
),
{"tid": str(tenant_id), "payload": '{"k": "v"}'},
)
await db_session.flush()
await db_session.commit()
async with session_factory() as worker_session:
count = await process_outbox_batch(worker_session, batch_size=10, tenant_ids=[tenant_id])
assert count == 0
rows = (
await db_session.execute(
text("SELECT status, attempts FROM event_outbox WHERE tenant_id = :tid"),
{"tid": str(tenant_id)},
)
).fetchall()
assert rows[0][0] == "failed"
finally:
bus.unsubscribe("test.maxfail", _always_fails)
@pytest.mark.asyncio
async def test_enqueue_multiple_events_and_batch_size(
db_session, session_factory: async_sessionmaker[AsyncSession],
):
"""Multiple events are enqueued and batch_size limits processing."""
tenant_id = uuid.uuid4()
received: list[str] = []
async def _handler(payload: dict) -> None:
received.append(payload.get("idx", "?"))
bus = get_event_bus()
bus.subscribe("test.batch", _handler)
try:
for i in range(5):
await enqueue_outbox_event(
db_session, tenant_id, "test.batch",
{"idx": str(i)},
)
await db_session.flush()
await db_session.commit()
async with session_factory() as worker_session:
count = await process_outbox_batch(worker_session, batch_size=3, tenant_ids=[tenant_id])
assert count == 3
assert len(received) == 3
# Process the remaining 2
async with session_factory() as worker_session:
count2 = await process_outbox_batch(worker_session, batch_size=3, tenant_ids=[tenant_id])
assert count2 == 2
assert len(received) == 5
finally:
bus.unsubscribe("test.batch", _handler)