Files
2026-08-30 20:44:00 +02:00

136 lines
5.6 KiB
Python

"""journal.py — SQLite saga journal: idempotency, statuses, durability (plan §10)."""
import pytest
from mcp_server.journal import ActionStatus, OperationStatus, SagaJournal
@pytest.fixture
async def fresh_journal(tmp_path):
journal = SagaJournal(tmp_path / "saga.sqlite")
await journal.connect()
yield journal
await journal.close()
async def test_connect_is_idempotent_safe_for_new_db(fresh_journal):
record = await fresh_journal.begin_operation("complete_step", "op-1")
assert record.status is OperationStatus.PLANNED
assert record.kind == "complete_step"
assert record.actions == []
async def test_begin_generates_operation_id_when_none_given(fresh_journal):
record = await fresh_journal.begin_operation("complete_step")
assert record.operation_id
async def test_record_and_mark_actions(fresh_journal):
await fresh_journal.begin_operation("complete_step", "op-1")
action = await fresh_journal.record_action(
"op-1", "decrement_container", {"item_id": 12, "subitem_id": 31, "qty_stored": 48.0}
)
assert action.status is ActionStatus.PLANNED
assert action.payload["subitem_id"] == 31
marked = await fresh_journal.mark_action(
"op-1", action.id, ActionStatus.DONE, response={"ok": True}
)
assert marked.status is ActionStatus.DONE
assert marked.response == {"ok": True}
async def test_get_operation_returns_actions_in_order(fresh_journal):
await fresh_journal.begin_operation("complete_step", "op-1")
first = await fresh_journal.record_action("op-1", "decrement_container", {"subitem_id": 31})
second = await fresh_journal.record_action("op-1", "finish_step", {"step_id": 9})
third = await fresh_journal.record_action("op-1", "post_comment", {"body": "done"})
record = await fresh_journal.get_operation("op-1")
assert [a.id for a in record.actions] == [first.id, second.id, third.id]
assert [a.action for a in record.actions] == [
"decrement_container",
"finish_step",
"post_comment",
]
async def test_has_completed_only_after_finish(fresh_journal):
await fresh_journal.begin_operation("complete_step", "op-1")
action = await fresh_journal.record_action("op-1", "decrement_container", {})
await fresh_journal.mark_action("op-1", action.id, ActionStatus.DONE)
assert await fresh_journal.has_completed("op-1") is False
await fresh_journal.finish_operation("op-1", OperationStatus.COMPLETED)
assert await fresh_journal.has_completed("op-1") is True
async def test_begin_with_completed_operation_id_is_idempotent_replay(fresh_journal):
"""Plan §13: double execution of the same operation_id is idempotent."""
await fresh_journal.begin_operation("complete_step", "op-1")
await fresh_journal.record_action("op-1", "decrement_container", {"subitem_id": 31})
await fresh_journal.finish_operation("op-1", OperationStatus.COMPLETED)
replayed = await fresh_journal.begin_operation("complete_step", "op-1")
assert replayed.status is OperationStatus.COMPLETED
assert len(replayed.actions) == 1
record = await fresh_journal.get_operation("op-1")
assert len(record.actions) == 1, "replay must not append new actions"
async def test_begin_with_inprogress_operation_id_resumes_it(fresh_journal):
await fresh_journal.begin_operation("complete_step", "op-1")
await fresh_journal.record_action("op-1", "decrement_container", {"subitem_id": 31})
resumed = await fresh_journal.begin_operation("complete_step", "op-1")
assert resumed.status is OperationStatus.PLANNED
assert len(resumed.actions) == 1
async def test_failed_and_compensated_statuses_persist(fresh_journal):
await fresh_journal.begin_operation("complete_step", "op-1")
action = await fresh_journal.record_action("op-1", "finish_step", {"step_id": 9})
await fresh_journal.mark_action("op-1", action.id, ActionStatus.FAILED)
await fresh_journal.finish_operation("op-1", OperationStatus.COMPENSATED)
record = await fresh_journal.get_operation("op-1")
assert record.status is OperationStatus.COMPENSATED
assert record.actions[0].status is ActionStatus.FAILED
assert await fresh_journal.has_completed("op-1") is False
async def test_partial_failure_status_roundtrip(fresh_journal):
"""Plan §10: compensation failure ⇒ partial_failure in the journal."""
await fresh_journal.begin_operation("complete_step", "op-1")
await fresh_journal.finish_operation("op-1", OperationStatus.PARTIAL_FAILURE)
record = await fresh_journal.get_operation("op-1")
assert record.status is OperationStatus.PARTIAL_FAILURE
async def test_journal_persists_across_reconnect(tmp_path):
"""The journal doubles as the audit log — it must survive restarts."""
path = tmp_path / "saga.sqlite"
journal = SagaJournal(path)
await journal.connect()
await journal.begin_operation("complete_step", "op-1")
action = await journal.record_action("op-1", "decrement_container", {"subitem_id": 31})
await journal.mark_action("op-1", action.id, ActionStatus.DONE, response={"qty": 48.0})
await journal.finish_operation("op-1", OperationStatus.COMPLETED)
await journal.close()
reopened = SagaJournal(path)
await reopened.connect()
try:
record = await reopened.get_operation("op-1")
assert record.kind == "complete_step"
assert record.status is OperationStatus.COMPLETED
assert record.actions[0].response == {"qty": 48.0}
assert await reopened.has_completed("op-1") is True
finally:
await reopened.close()
async def test_get_operation_of_unknown_id_raises(fresh_journal):
with pytest.raises(Exception):
await fresh_journal.get_operation("nope")