tmux-amq/tests/test_client_adapter.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

202 lines
7.2 KiB
Python

import asyncio
import json
import pytest
from tamq.client import (
CoordinationEngineAdapter,
TamqClient,
TamqClientError,
TamqProtocolError,
TamqTargetUnavailable,
WakeRequest,
)
from tamq.service import Service
from tamq.store import Store
async def _start_service(path, store):
service = Service(path, store)
server = await asyncio.start_unix_server(service.handle, path=str(path))
return service, server
def test_coordination_adapter_negotiates_and_deduplicates_wake(tmp_path, monkeypatch):
monkeypatch.setattr("tamq.service.validate_targets", lambda repos: None)
async def run():
path = tmp_path / "tamq.sock"
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint(
"ep-1", 42, "tamq", ["target"], "manual",
delivery_ack_mode="acknowledged",
)
service, server = await _start_service(path, store)
adapter = CoordinationEngineAdapter(TamqClient(path))
wake = WakeRequest(
lease_id="lease-7",
trigger_id="trigger-9",
target_repo="target",
prompt="Continue task T1.",
)
try:
first = await adapter.wake(wake)
repeated = await adapter.wake(wake)
assert repeated.message_id == first.message_id
assert first.deduplicated is False
assert repeated.deduplicated is True
row = await adapter.receipt(first.message_id)
assert row["client_id"] == "coordination-engine"
assert row["idempotency_key"] == "lease-7"
assert row["correlation_id"] == "trigger-9"
assert json.loads(row["envelope_metadata"])["trigger_id"] == "trigger-9"
assert len(store.list()) == 1
assert store.line_state("ep-1", "coordination-engine")["messages"] == 1
await adapter.acknowledge(first.message_id)
assert (await adapter.receipt(first.message_id))["state"] == "acknowledged"
finally:
server.close()
await server.wait_closed()
service.close()
asyncio.run(run())
def test_client_reconnects_and_recovers_receipt_from_durable_store(tmp_path, monkeypatch):
monkeypatch.setattr("tamq.service.validate_targets", lambda repos: None)
async def run():
path = tmp_path / "tamq.sock"
database = tmp_path / "queue.sqlite3"
store = Store(database)
store.register_endpoint("ep", 42, "tamq", ["target"], "manual")
first_service, first_server = await _start_service(path, store)
client = TamqClient(path)
sent = await client.send(
sender_repo="coordination-engine",
target_repo="target",
body="wake",
endpoint_id="ep",
idempotency_key="lease-restart",
)
first_server.close()
await first_server.wait_closed()
first_service.close()
path.unlink(missing_ok=True)
reopened = Store(database)
second_service, second_server = await _start_service(path, reopened)
try:
receipt = await client.message(sent["message_id"])
assert receipt["state"] == "pending"
repeated = await client.send(
sender_repo="coordination-engine",
target_repo="target",
body="wake",
endpoint_id="ep",
idempotency_key="lease-restart",
)
assert repeated["deduplicated"] is True
finally:
second_server.close()
await second_server.wait_closed()
second_service.close()
asyncio.run(run())
def test_adapter_rejects_incompatible_protocol_and_disappeared_endpoint(
tmp_path, monkeypatch
):
monkeypatch.setattr("tamq.service.validate_targets", lambda repos: None)
async def run():
path = tmp_path / "tamq.sock"
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint("ep", 42, "tamq", ["target"], "manual")
service, server = await _start_service(path, store)
try:
with pytest.raises(TamqProtocolError):
await TamqClient(path, protocol="2.0").negotiate()
store.disconnect_endpoint("ep")
adapter = CoordinationEngineAdapter(TamqClient(path))
with pytest.raises(TamqTargetUnavailable):
await adapter.wake(
WakeRequest("lease", "target", "wake", "trigger")
)
finally:
server.close()
await server.wait_closed()
service.close()
asyncio.run(run())
def test_repeated_wake_rebinds_same_message_after_endpoint_replacement(
tmp_path, monkeypatch
):
monkeypatch.setattr("tamq.service.validate_targets", lambda repos: None)
async def run():
path = tmp_path / "tamq.sock"
store = Store(tmp_path / "queue.sqlite3")
store.register_endpoint("old", 41, "tamq-old", ["target"], "manual")
service, server = await _start_service(path, store)
adapter = CoordinationEngineAdapter(TamqClient(path))
wake = WakeRequest("lease-rebind", "target", "wake", "trigger")
try:
first = await adapter.wake(wake)
store.disconnect_endpoint("old")
store.register_endpoint("new", 42, "tamq-new", ["target"], "manual")
repeated = await adapter.wake(wake)
assert repeated.message_id == first.message_id
assert repeated.endpoint_id == "new"
assert repeated.deduplicated is True
assert store.message(first.message_id)["endpoint_id"] == "new"
rebound = store.protocol_events(
message_id=first.message_id, event_type="message.rebound"
)
assert len(rebound) == 1
finally:
server.close()
await server.wait_closed()
service.close()
asyncio.run(run())
def test_idempotency_key_conflict_is_rejected(tmp_path, monkeypatch):
monkeypatch.setattr("tamq.service.validate_targets", lambda repos: None)
async def run():
path = tmp_path / "tamq.sock"
store = Store(tmp_path / "queue.sqlite3")
service, server = await _start_service(path, store)
client = TamqClient(path)
try:
first = await client.send(
sender_repo="coordination-engine",
target_repo="one",
body="first",
idempotency_key="same",
)
with pytest.raises(TamqClientError, match="different message"):
await client.send(
sender_repo="coordination-engine",
target_repo="two",
body="second",
idempotency_key="same",
)
store.db.execute(
"UPDATE messages SET state='failed',attempt_count=4 WHERE message_id=?",
(first["message_id"],),
)
store.db.commit()
retried = await client.retry(first["message_id"])
assert retried["state"] == "pending"
finally:
server.close()
await server.wait_closed()
service.close()
asyncio.run(run())