Some checks failed
tamq-ci / test (push) Failing after 5s
Assistant: codex Assistant-Model: gpt-5.6-sol Assistant-Session: 01a03397-4d51-7fd1-8ff2-946eb22ea2bc
202 lines
7.2 KiB
Python
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())
|