tmux-amq/tests/test_service_delivery.py
tegwick 6d2ccc7760
Some checks failed
tamq-ci / test (push) Failing after 5s
feat: complete reliable coordination adapter
Assistant: codex
Assistant-Model: gpt-5.6-sol
Assistant-Session: 01a03397-4d51-7fd1-8ff2-946eb22ea2bc
2026-08-26 08:11:09 +02:00

314 lines
11 KiB
Python

import json
from tamq.control import PaneDisplay
from tamq.service import Service
from tamq.store import Store
class FakeControl:
instances = []
def __init__(self, session):
self.session = session
self.injected = []
self.placed = []
self.submitted = []
self.__class__.instances.append(self)
def start(self):
return None
def session_exists(self, expected_pid=None):
return True
def inject(self, window, text):
self.injected.append((window, text))
def submit(self, window, text):
self.submitted.append((window, text))
def place(self, window, text):
self.placed.append((window, text))
def pane_display(self, window):
return PaneDisplay(
tty_path=f"/dev/pts/{window.rsplit(':', 1)[-1]}",
cursor_x=0,
cursor_y=0,
pane_height=24,
alternate_on=False,
)
def close(self):
return None
def test_service_delivery_injects_pending_messages(tmp_path, monkeypatch):
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint("tmux-amq-9-boot", 9, "tamq", ["repo-b"], "pane")
message_id = store.add("repo-a", "repo-b", "continue")
monkeypatch.setattr("tamq.service.ControlModeClient", FakeControl)
Service(store=store)._deliver_once()
assert FakeControl.instances[-1].injected == [("tamq:repo-b", "From:repo-a: continue")]
assert store.list()[0]["message_id"] == message_id
assert store.list()[0]["state"] == "injected"
def test_service_never_injects_manual_endpoint_messages(tmp_path, monkeypatch):
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint("tmux-amq-9-boot", 9, "tamq", ["repo-b"], "manual")
store.add("repo-a", "repo-b", "continue")
monkeypatch.setattr("tamq.service.ControlModeClient", FakeControl)
Service(store=store)._deliver_once()
assert FakeControl.instances[-1].injected == []
assert store.list()[0]["state"] == "pending"
def test_service_pushy_mode_places_once_without_enter(tmp_path, monkeypatch):
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint("tmux-amq-9-boot", 9, "tamq", ["repo-b"], "pushy")
message_id = store.add("repo-a", "repo-b", "first\nsecond")
monkeypatch.setattr("tamq.service.ControlModeClient", FakeControl)
service = Service(store=store, input_grace=0)
service._deliver_once()
control = FakeControl.instances[-1]
service._deliver_once()
assert control.placed == [
(
"tamq:repo-b",
"From:repo-a: first\nFrom:repo-a: second",
)
]
assert control.submitted == []
assert store.list()[0]["state"] == "injected"
def test_service_trigger_mode_submits_exactly_once(tmp_path, monkeypatch):
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint("ep", 9, "tamq", ["repo-b"], "trigger")
store.add("repo-a", "repo-b", "go", provenance="operator_input")
monkeypatch.setattr("tamq.service.ControlModeClient", FakeControl)
service = Service(store=store, input_grace=0)
service._deliver_once()
control = FakeControl.instances[-1]
service._deliver_once()
assert control.submitted == [("tamq:repo-b", "From:repo-a/o: go")]
assert control.placed == []
lifecycle = store.protocol_events(message_id=store.list()[0]["message_id"])
assert [event["event_type"] for event in lifecycle] == [
"message.accepted",
"delivery.attempted",
"delivery.injected",
]
assert lifecycle[-1]["delivery_mode"] == "trigger"
def test_failed_pushy_submission_remains_pending(tmp_path, monkeypatch):
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint("tmux-amq-9-boot", 9, "tamq", ["repo-b"], "pushy")
store.add("repo-a", "repo-b", "continue")
class FailingControl(FakeControl):
def place(self, window, text):
raise RuntimeError("tmux rejected input")
monkeypatch.setattr("tamq.service.ControlModeClient", FailingControl)
Service(store=store, input_grace=0)._deliver_once()
assert store.list()[0]["state"] == "pending"
failure = store.protocol_events(event_type="delivery.failed")[0]
assert failure["delivery_mode"] == "pushy"
detail = json.loads(failure["detail"])
assert detail["reason"] == "RuntimeError"
assert detail["attempt"] == 1
assert detail["attempt_limit"] == 4
def test_input_delivery_waits_for_new_endpoint_grace(tmp_path, monkeypatch):
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint("ep", 9, "tamq", ["repo-b"], "trigger")
store.add("repo-a", "repo-b", "first")
monkeypatch.setattr("tamq.service.ControlModeClient", FakeControl)
service = Service(store=store, input_grace=10)
service._deliver_once()
assert FakeControl.instances[-1].submitted == []
assert store.list()[0]["state"] == "pending"
def test_default_injected_policy_completes_output_delivery_once(tmp_path, monkeypatch):
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint("tmux-amq-9-boot", 9, "tamq", ["repo-b"], "output")
message_id = store.add("repo-a", "repo-b", "continue")
writes = []
monkeypatch.setattr("tamq.service.ControlModeClient", FakeControl)
monkeypatch.setattr(
"tamq.service.write_terminal_output",
lambda path, text: writes.append((path, text)),
)
service = Service(store=store)
service._deliver_once()
service._deliver_once()
assert writes == [
(
"/dev/pts/repo-b",
"\r\nFrom:repo-a: continue\r\n",
)
]
row = store.list()[0]
assert row["state"] == "injected"
assert row["displayed_at"] is not None
assert store.db.execute("SELECT COUNT(*) FROM leases").fetchone()[0] == 0
def test_failed_terminal_output_remains_undisplayed_and_retryable(tmp_path, monkeypatch):
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint("tmux-amq-9-boot", 9, "tamq", ["repo-b"], "output")
store.add("repo-a", "repo-b", "continue")
attempts = []
monkeypatch.setattr("tamq.service.ControlModeClient", FakeControl)
def fail_once(path, text):
attempts.append((path, text))
if len(attempts) == 1:
raise OSError("temporary failure")
monkeypatch.setattr("tamq.service.write_terminal_output", fail_once)
service = Service(store=store, retry_backoff=(0,))
service._deliver_once()
assert store.list()[0]["displayed_at"] is None
service._deliver_once()
assert len(attempts) == 2
assert store.list()[0]["displayed_at"] is not None
def test_acknowledged_policy_waits_then_redelivers_same_message(tmp_path, monkeypatch):
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint(
"ep", 9, "tamq", ["repo-b"], "trigger",
delivery_ack_mode="acknowledged",
delivery_max_attempts=3,
ack_timeout_seconds=10,
)
message_id = store.add("repo-a", "repo-b", "confirm")
monkeypatch.setattr("tamq.service.ControlModeClient", FakeControl)
service = Service(store=store, input_grace=0, retry_backoff=(0,))
service._deliver_once()
first_control = FakeControl.instances[-1]
assert store.message(message_id)["state"] == "awaiting_ack"
assert store.message(message_id)["attempt_count"] == 1
assert len(first_control.submitted) == 1
service._deliver_once()
assert FakeControl.instances[-1].submitted == []
store.db.execute(
"UPDATE messages SET next_attempt_at=0 WHERE message_id=?", (message_id,)
)
store.db.commit()
service._deliver_once()
second_control = FakeControl.instances[-1]
assert store.message(message_id)["attempt_count"] == 2
assert len(second_control.submitted) == 1
assert store.acknowledge(message_id) is True
service._deliver_once()
assert FakeControl.instances[-1].submitted == []
def test_delivery_attempt_cap_is_terminal_and_operator_retry_resets(tmp_path, monkeypatch):
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint(
"ep", 9, "tamq", ["repo-b"], "pushy", delivery_max_attempts=2
)
message_id = store.add("repo-a", "repo-b", "fail")
class FailingControl(FakeControl):
def place(self, window, text):
raise OSError("no pane")
monkeypatch.setattr("tamq.service.ControlModeClient", FailingControl)
service = Service(store=store, input_grace=0, retry_backoff=(0,))
service._deliver_once()
service._deliver_once()
service._deliver_once()
row = store.message(message_id)
assert row["state"] == "failed"
assert row["attempt_count"] == 2
assert row["last_failure_reason"] == "OSError"
assert store.retry_message(message_id) is True
assert store.message(message_id)["state"] == "pending"
assert store.message(message_id)["attempt_count"] == 0
def test_ack_timeout_exhaustion_is_terminal_but_late_ack_wins(tmp_path, monkeypatch):
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint(
"ep", 9, "tamq", ["repo-b"], "trigger",
delivery_ack_mode="acknowledged",
delivery_max_attempts=2,
ack_timeout_seconds=10,
)
message_id = store.add("repo-a", "repo-b", "confirm")
monkeypatch.setattr("tamq.service.ControlModeClient", FakeControl)
service = Service(store=store, input_grace=0, retry_backoff=(0,))
service._deliver_once()
for _ in range(2):
store.db.execute(
"UPDATE messages SET next_attempt_at=0 WHERE message_id=?", (message_id,)
)
store.db.commit()
service._deliver_once()
row = store.message(message_id)
assert row["state"] == "failed"
assert row["attempt_count"] == 2
assert row["last_failure_reason"] == "ack_timeout"
assert store.acknowledge(message_id) is True
assert store.message(message_id)["state"] == "acknowledged"
assert store.protocol_events(message_id=message_id)[-1]["outcome"] == "late_acknowledged"
def test_expired_delivery_lease_consumes_attempt_and_schedules_retry(tmp_path):
store = Store(tmp_path / "queue.sqlite3")
message_id = store.add("repo-a", "repo-b", "lease")
claim = store.claim_delivery(
message_id,
"ep",
attempt_limit=2,
retry_delay=7,
delivery_mode="output",
ttl=1,
now=10,
)
assert claim is not None
assert store.expire_delivery_leases(now=12) == 1
row = store.message(message_id)
assert row["state"] == "pending"
assert row["attempt_count"] == 1
assert row["last_failure_reason"] == "lease_expired"
assert row["next_attempt_at"] == 19
def test_service_disconnects_disappeared_tmux_endpoint(tmp_path, monkeypatch):
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint("tmux-amq-9-boot", 9, "tamq", ["repo-b"])
class MissingControl(FakeControl):
def session_exists(self, expected_pid=None):
assert expected_pid == 9
return False
monkeypatch.setattr("tamq.service.ControlModeClient", MissingControl)
Service(store=store)._deliver_once()
assert store.endpoints() == []