diff --git a/test/test_api_v3_on_demand_socket.py b/test/test_api_v3_on_demand_socket.py index 534c6a91..ae4c66e8 100644 --- a/test/test_api_v3_on_demand_socket.py +++ b/test/test_api_v3_on_demand_socket.py @@ -306,8 +306,11 @@ class TestNoDisplayListening: resp = api_v3_client.post(START_URL, json={"plugin_id": "clock"}) assert resp.status_code == 200 outcome = self._outcome() - assert sent == [resp.get_json()["data"]["request_id"]] - assert old not in sent + # The old start is cancelled before the new one is sent, and + # cancel() waits out a send of it in flight: nothing of it lands + # after the new one. + new = resp.get_json()["data"]["request_id"] + assert sent[-1] == new and sent.count(new) == 1 assert outcome["status"] == "idle" and outcome["last_event"] == "superseded" def test_a_service_that_will_not_start_is_an_error(self, api_v3_client, service, clock): @@ -478,13 +481,20 @@ class TestRealSocket: assert data["transport"] == "socket" assert [x.request_id for x in live.drain()] == [data["request_id"]] - def test_a_display_that_went_away_is_an_error(self, api_v3_client, service, live, + def test_a_display_that_went_away_is_given_up_on(self, api_v3_client, service, live, monkeypatch): + from web_interface import on_demand_dispatch monkeypatch.setattr(f"{DISPLAY}.ON_DEMAND_SOCKET_WAIT_RUNNING_SECONDS", 0.3) + monkeypatch.setattr(on_demand_dispatch, "RETRY_INTERVAL", 0.01) live.close() resp = api_v3_client.post(START_URL, json={"plugin_id": "weather"}) - assert resp.status_code == 503 + # The service still reads as running: answered at once, and the + # web process gives up once the wait is over. + assert resp.status_code == 202 assert resp.get_json()["data"]["socket_error"] == "no_socket" + d = on_demand_dispatch.current() + assert _until(lambda: not d.pending()) + assert d.status()["error"] == "start-timeout" assert _mailbox_writes(service["cache"]) == [] def test_a_full_queue_is_reported_not_mailed(self, api_v3_client, service, monkeypatch): diff --git a/test/test_on_demand_dispatch.py b/test/test_on_demand_dispatch.py index b57da654..2dc9aaa6 100644 --- a/test/test_on_demand_dispatch.py +++ b/test/test_on_demand_dispatch.py @@ -204,6 +204,29 @@ class TestOneAtATime: status = d.status() assert status["status"] == "idle" and status["last_event"] == "requested-stop" + def test_a_cancel_waits_for_a_send_in_flight(self, make): + # Whatever the caller sends after cancel() must land after the + # cancelled start, not race it to the display. + in_flight, release = threading.Event(), threading.Event() + + def slow(payload): + in_flight.set() + release.wait(5) + raise _not_listening() + + send = FakeSend(slow) + d = make(send) + d.submit(_payload("old")) + assert in_flight.wait(5) + done = threading.Event() + threading.Thread(target=lambda: (d.cancel("superseded"), done.set()), + daemon=True).start() + assert not done.wait(0.1), "cancel returned while the old send was in flight" + release.set() + assert done.wait(5) + assert _settled(d) + assert send.sent == ["old"] + def test_a_cancel_with_nothing_pending_does_nothing(self, make): d = make(FakeSend("ack")) assert d.cancel() is None diff --git a/web_interface/on_demand_dispatch.py b/web_interface/on_demand_dispatch.py index 288ea52a..92cfe3f3 100644 --- a/web_interface/on_demand_dispatch.py +++ b/web_interface/on_demand_dispatch.py @@ -63,6 +63,11 @@ class OnDemandDispatcher: self._clock = clock self._wall = wall_clock self._lock = threading.Lock() + # Held by the worker while it reads the pending start and sends it, + # so cancel() can wait out a send already in flight: a cancelled (or + # superseded) start must not reach the display after the request + # that replaced it. Never taken while holding _lock. + self._send_lock = threading.Lock() self._wake = threading.Event() # Bumped by every submit and cancel: a send that started under an # older generation does not report its result as the current one. @@ -106,6 +111,10 @@ class OnDemandDispatcher: self._pending = None self._finish(self._describe('idle', pending, last_event=reason)) self._wake.set() + # A send of it may be in flight (at most the client's timeout, 1 s): + # let it finish, so whatever the caller sends next lands after it. + with self._send_lock: + pass logger.info("On-demand start %s cancelled before the display took it (%s)", pending.get('request_id'), reason) return pending.get('request_id') @@ -152,13 +161,14 @@ class OnDemandDispatcher: def _run(self) -> None: while True: - with self._lock: - payload, generation = self._pending, self._generation - deadline = self._deadline - if payload is None: - self._thread = None - return - outcome, error = self._attempt(payload) + with self._send_lock: + with self._lock: + payload, generation = self._pending, self._generation + deadline = self._deadline + if payload is None: + self._thread = None + return + outcome, error = self._attempt(payload) with self._lock: if generation != self._generation: continue # superseded or cancelled meanwhile