Initial version
This commit is contained in:
@@ -0,0 +1,135 @@
|
||||
"""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")
|
||||
Reference in New Issue
Block a user