from __future__ import annotations import asyncio import pytest from cmvr_edge_ai.server.invocation import ( CompletionDisposition, DuplicateInvocationError, DuplicateInvocationRouteError, DuplicateInvocationWaitError, InvocationBroker, InvocationBrokerClosedError, InvocationCancelledError, InvocationNotFoundError, InvocationQueueFullError, InvocationTimeoutError, UnknownInvocationRouteError, ) def _broker() -> InvocationBroker[dict[str, object], dict[str, object]]: broker: InvocationBroker[dict[str, object], dict[str, object]] = ( InvocationBroker() ) broker.register_route("detect.mobile_phone", capacity=2) broker.register_route("detect.analog_gauge", capacity=1) return broker def test_routes_are_registered_explicitly_and_validated() -> None: broker: InvocationBroker[object, object] = InvocationBroker() broker.register_route(" detect.mobile_phone ", capacity=2) assert broker.categories == ("detect.mobile_phone",) assert broker.route_stats("detect.mobile_phone").capacity == 2 with pytest.raises(DuplicateInvocationRouteError, match="already registered"): broker.register_route("detect.mobile_phone", capacity=1) with pytest.raises(ValueError, match="positive integer"): broker.register_route("bad.zero", capacity=0) with pytest.raises(ValueError, match="positive integer"): broker.register_route("bad.bool", capacity=True) def test_submit_receive_complete_and_wait_round_trip() -> None: async def exercise(): # type: ignore[no-untyped-def] broker = _broker() submitted = await broker.submit( "detect.mobile_phone", "request-1", {"image": b"jpeg"}, ) received = await broker.receive("detect.mobile_phone") assert received is submitted assert received.category == "detect.mobile_phone" assert received.request_id == "request-1" assert received.payload == {"image": b"jpeg"} assert received.submitted_at_ns > 0 waiter = asyncio.create_task(broker.wait("request-1", timeout_s=1)) await asyncio.sleep(0) disposition = await broker.complete( "request-1", {"count": 1}, invocation_token=submitted.invocation_token, ) response = await waiter stats = broker.route_stats("detect.mobile_phone") await broker.close() return disposition, response, stats disposition, response, stats = asyncio.run(exercise()) assert disposition is CompletionDisposition.ACCEPTED assert response == {"count": 1} assert stats.queued == 0 assert stats.pending == 0 def test_category_consumers_never_compete_for_another_route_queue() -> None: async def exercise(): # type: ignore[no-untyped-def] broker = _broker() gauge_receiver = asyncio.create_task( broker.receive("detect.analog_gauge") ) phone_receiver = asyncio.create_task( broker.receive("detect.mobile_phone") ) await asyncio.sleep(0) await broker.submit("detect.mobile_phone", "phone-1", {"frame": 1}) phone_request = await asyncio.wait_for(phone_receiver, timeout=1) await asyncio.sleep(0) assert gauge_receiver.done() is False await broker.complete( "phone-1", {"boxes": []}, invocation_token=phone_request.invocation_token, ) await broker.wait("phone-1", timeout_s=1) await broker.close() with pytest.raises(InvocationBrokerClosedError): await gauge_receiver return phone_request request = asyncio.run(exercise()) assert request.category == "detect.mobile_phone" assert request.request_id == "phone-1" def test_route_capacity_bounds_all_outstanding_work_independently() -> None: async def exercise(): # type: ignore[no-untyped-def] broker = _broker() await broker.submit("detect.analog_gauge", "gauge-1", {"frame": 1}) with pytest.raises(InvocationQueueFullError) as captured: await broker.submit("detect.analog_gauge", "gauge-2", {"frame": 2}) assert captured.value.category == "detect.analog_gauge" assert captured.value.capacity == 1 # A full gauge route does not consume the phone route's capacity. await broker.submit("detect.mobile_phone", "phone-1", {"frame": 3}) assert broker.route_stats("detect.analog_gauge").pending == 1 assert broker.route_stats("detect.mobile_phone").pending == 1 assert await broker.cancel("gauge-1") is True await broker.submit("detect.analog_gauge", "gauge-2", {"frame": 2}) next_gauge = await broker.receive("detect.analog_gauge") await broker.cancel("gauge-2") await broker.cancel("phone-1") await broker.close() return next_gauge request = asyncio.run(exercise()) assert request.request_id == "gauge-2" def test_request_ids_are_global_across_routes_and_recent_history() -> None: async def exercise(): # type: ignore[no-untyped-def] broker = _broker() submitted = await broker.submit( "detect.mobile_phone", "global-id", {"frame": 1} ) with pytest.raises(DuplicateInvocationError, match="global-id"): await broker.submit( "detect.analog_gauge", "global-id", {"frame": 2}, ) await broker.complete( "global-id", {"boxes": []}, invocation_token=submitted.invocation_token, ) assert await broker.wait("global-id") == {"boxes": []} with pytest.raises(DuplicateInvocationError, match="global-id"): await broker.submit( "detect.analog_gauge", "global-id", {"frame": 3}, ) await broker.close() asyncio.run(exercise()) def test_out_of_order_results_are_correlated_by_request_id() -> None: async def exercise(): # type: ignore[no-untyped-def] broker = _broker() first = await broker.submit( "detect.mobile_phone", "first", {"frame": 1} ) second = await broker.submit( "detect.mobile_phone", "second", {"frame": 2} ) assert (await broker.receive("detect.mobile_phone")).request_id == "first" assert (await broker.receive("detect.mobile_phone")).request_id == "second" first_waiter = asyncio.create_task(broker.wait("first", timeout_s=1)) second_waiter = asyncio.create_task(broker.wait("second", timeout_s=1)) await asyncio.sleep(0) await broker.complete( "second", {"owner": "second"}, invocation_token=second.invocation_token, ) await broker.complete( "first", {"owner": "first"}, invocation_token=first.invocation_token, ) results = await asyncio.gather(first_waiter, second_waiter) await broker.close() return results assert asyncio.run(exercise()) == [ {"owner": "first"}, {"owner": "second"}, ] def test_duplicate_and_unknown_result_deliveries_are_non_fatal() -> None: async def exercise(): # type: ignore[no-untyped-def] broker = _broker() submitted = await broker.submit( "detect.mobile_phone", "request-1", {"frame": 1} ) accepted = await broker.complete( "request-1", {"boxes": [1]}, invocation_token=submitted.invocation_token, ) duplicate_before_wait = await broker.fail( "request-1", RuntimeError("duplicate"), invocation_token=submitted.invocation_token, ) response = await broker.wait("request-1") duplicate_after_wait = await broker.complete( "request-1", {"boxes": [2]}, invocation_token=submitted.invocation_token, ) unknown = await broker.complete( "not-submitted", {"boxes": []}, invocation_token="unknown-token", ) await broker.close() return ( accepted, duplicate_before_wait, duplicate_after_wait, unknown, response, ) accepted, duplicate_before, duplicate_after, unknown, response = asyncio.run( exercise() ) assert accepted is CompletionDisposition.ACCEPTED assert duplicate_before is CompletionDisposition.DUPLICATE assert duplicate_after is CompletionDisposition.DUPLICATE assert unknown is CompletionDisposition.UNKNOWN assert response == {"boxes": [1]} def test_pipeline_failure_is_delivered_to_the_correct_waiter() -> None: async def exercise(): # type: ignore[no-untyped-def] broker = _broker() expected = ValueError("invalid gauge ellipse") submitted = await broker.submit( "detect.analog_gauge", "gauge-1", {"frame": 1} ) waiter = asyncio.create_task(broker.wait("gauge-1", timeout_s=1)) await asyncio.sleep(0) assert ( await broker.fail( "gauge-1", expected, invocation_token=submitted.invocation_token, ) is CompletionDisposition.ACCEPTED ) with pytest.raises(ValueError, match="invalid gauge ellipse") as captured: await waiter assert captured.value is expected duplicate = await broker.fail( "gauge-1", RuntimeError("again"), invocation_token=submitted.invocation_token, ) await broker.close() return duplicate assert asyncio.run(exercise()) is CompletionDisposition.DUPLICATE def test_timeout_removes_queued_request_and_classifies_late_result() -> None: async def exercise(): # type: ignore[no-untyped-def] broker = _broker() submitted = await broker.submit( "detect.analog_gauge", "slow", {"frame": 1} ) with pytest.raises(InvocationTimeoutError, match="slow"): await broker.wait("slow", timeout_s=0.01) assert broker.route_stats("detect.analog_gauge").queued == 0 assert broker.route_stats("detect.analog_gauge").pending == 0 late = await broker.complete( "slow", {"reading": 4.2}, invocation_token=submitted.invocation_token, ) # The cancelled queue entry cannot poison the next request. await broker.submit("detect.analog_gauge", "next", {"frame": 2}) next_request = await asyncio.wait_for( broker.receive("detect.analog_gauge"), timeout=1, ) await broker.cancel("next") await broker.close() return late, next_request late, next_request = asyncio.run(exercise()) assert late is CompletionDisposition.LATE assert next_request.request_id == "next" def test_stale_generation_cannot_complete_a_reused_request_id() -> None: async def exercise(): # type: ignore[no-untyped-def] broker: InvocationBroker[dict[str, object], dict[str, object]] = ( InvocationBroker(terminal_history_size=1) ) broker.register_route("detect.mobile_phone", capacity=1) old_request = await broker.submit( "detect.mobile_phone", "reused-id", {"generation": 1} ) assert await broker.receive("detect.mobile_phone") is old_request with pytest.raises(InvocationTimeoutError): await broker.wait("reused-id", timeout_s=0.001) assert broker.route_stats("detect.mobile_phone").pending == 1 with pytest.raises(InvocationQueueFullError): await broker.submit( "detect.mobile_phone", "while-model-busy", {"generation": 9} ) assert ( await broker.complete( old_request.request_id, {"generation": 1}, invocation_token=old_request.invocation_token, ) is CompletionDisposition.LATE ) assert broker.route_stats("detect.mobile_phone").pending == 0 eviction = await broker.submit( "detect.mobile_phone", "evict-history", {"generation": 0} ) assert await broker.receive("detect.mobile_phone") is eviction await broker.complete( eviction.request_id, {"generation": 0}, invocation_token=eviction.invocation_token, ) assert await broker.wait(eviction.request_id) == {"generation": 0} new_request = await broker.submit( "detect.mobile_phone", "reused-id", {"generation": 2} ) assert await broker.receive("detect.mobile_phone") is new_request stale = await broker.complete( old_request.request_id, {"generation": 1}, invocation_token=old_request.invocation_token, ) assert stale is CompletionDisposition.STALE accepted = await broker.complete( new_request.request_id, {"generation": 2}, invocation_token=new_request.invocation_token, ) response = await broker.wait(new_request.request_id) await broker.close() return accepted, response accepted, response = asyncio.run(exercise()) assert accepted is CompletionDisposition.ACCEPTED assert response == {"generation": 2} def test_consumed_timeout_remains_waitable_after_terminal_history_eviction() -> None: async def exercise(): # type: ignore[no-untyped-def] broker: InvocationBroker[dict[str, object], dict[str, object]] = ( InvocationBroker(terminal_history_size=1) ) broker.register_route("detect.mobile_phone", capacity=1) broker.register_route("detect.analog_gauge", capacity=1) slow = await broker.submit( "detect.mobile_phone", "slow", {"generation": 1} ) assert await broker.receive("detect.mobile_phone") is slow with pytest.raises(InvocationTimeoutError): await broker.wait("slow", timeout_s=0.001) # Churn unrelated terminal history while the consumed model request is # deliberately retained to keep its route capacity reserved. await broker.submit("detect.analog_gauge", "evict", {"generation": 0}) assert await broker.cancel("evict") is True with pytest.raises(InvocationTimeoutError, match="slow"): await broker.wait("slow") assert ( await broker.complete( "slow", {"generation": 1}, invocation_token=slow.invocation_token, ) is CompletionDisposition.LATE ) assert broker.route_stats("detect.mobile_phone").pending == 0 await broker.close() asyncio.run(exercise()) def test_explicit_cancel_wakes_waiter_and_rejects_late_result() -> None: async def exercise(): # type: ignore[no-untyped-def] broker = _broker() submitted = await broker.submit( "detect.mobile_phone", "cancel-me", {"frame": 1} ) waiter = asyncio.create_task(broker.wait("cancel-me")) await asyncio.sleep(0) assert await broker.cancel("cancel-me") is True assert await broker.cancel("cancel-me") is False with pytest.raises(InvocationCancelledError, match="cancel-me"): await waiter late = await broker.fail( "cancel-me", RuntimeError("too late"), invocation_token=submitted.invocation_token, ) await broker.close() return late assert asyncio.run(exercise()) is CompletionDisposition.LATE def test_cancelling_transport_wait_cancels_correlation_but_not_other_routes() -> None: async def exercise(): # type: ignore[no-untyped-def] broker = _broker() disconnected_request = await broker.submit( "detect.mobile_phone", "disconnected", {"frame": 1} ) healthy_request = await broker.submit( "detect.analog_gauge", "healthy", {"frame": 2} ) assert ( await broker.receive("detect.mobile_phone") is disconnected_request ) disconnected = asyncio.create_task(broker.wait("disconnected")) await asyncio.sleep(0) disconnected.cancel() with pytest.raises(asyncio.CancelledError): await disconnected assert ( await broker.complete( "disconnected", {"boxes": []}, invocation_token=disconnected_request.invocation_token, ) is CompletionDisposition.LATE ) assert ( await broker.complete( "healthy", {"reading": 2.0}, invocation_token=healthy_request.invocation_token, ) is CompletionDisposition.ACCEPTED ) healthy = await broker.wait("healthy", timeout_s=1) await broker.close() return healthy assert asyncio.run(exercise()) == {"reading": 2.0} def test_only_one_waiter_can_own_a_request() -> None: async def exercise(): # type: ignore[no-untyped-def] broker = _broker() submitted = await broker.submit( "detect.mobile_phone", "one-waiter", {"frame": 1} ) first = asyncio.create_task(broker.wait("one-waiter")) await asyncio.sleep(0) with pytest.raises(DuplicateInvocationWaitError, match="already has"): await broker.wait("one-waiter") await broker.complete( "one-waiter", {"boxes": []}, invocation_token=submitted.invocation_token, ) response = await first await broker.close() return response assert asyncio.run(exercise()) == {"boxes": []} def test_close_is_idempotent_and_wakes_receivers_and_waiters() -> None: async def exercise(): # type: ignore[no-untyped-def] broker = _broker() receiver = asyncio.create_task(broker.receive("detect.analog_gauge")) submitted = await broker.submit( "detect.mobile_phone", "pending", {"frame": 1} ) waiter = asyncio.create_task(broker.wait("pending")) await asyncio.sleep(0) await broker.close() await broker.close() with pytest.raises(InvocationBrokerClosedError): await receiver with pytest.raises(InvocationBrokerClosedError): await waiter assert broker.closed is True assert broker.pending_count == 0 assert ( await broker.complete( "pending", {"boxes": []}, invocation_token=submitted.invocation_token, ) is CompletionDisposition.LATE ) with pytest.raises(InvocationBrokerClosedError): await broker.submit("detect.mobile_phone", "new", {"frame": 2}) with pytest.raises(InvocationBrokerClosedError): await broker.receive("detect.mobile_phone") with pytest.raises(InvocationBrokerClosedError): broker.register_route("talk.asr", capacity=1) asyncio.run(exercise()) def test_unknown_routes_requests_and_invalid_timeout_fail_near_the_caller() -> None: async def exercise(): # type: ignore[no-untyped-def] broker = _broker() with pytest.raises(UnknownInvocationRouteError, match="available"): await broker.submit("talk.asr", "request-1", {}) with pytest.raises(UnknownInvocationRouteError, match="available"): await broker.receive("talk.asr") with pytest.raises(InvocationNotFoundError, match="missing"): await broker.wait("missing") submitted = await broker.submit("detect.mobile_phone", "request-2", {}) with pytest.raises(ValueError, match="finite positive"): await broker.wait("request-2", timeout_s=float("inf")) with pytest.raises(TypeError, match="must be an Exception"): await broker.fail( "request-2", "not-an-error", # type: ignore[arg-type] invocation_token=submitted.invocation_token, ) await broker.cancel("request-2") await broker.close() asyncio.run(exercise())