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())