136 lines
5.6 KiB
Python
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")
|