From 172c31d759cde44522cc2d089ec9a22a3d0d61e6 Mon Sep 17 00:00:00 2001 From: YusefSyed <211442445+YusefSyed@users.noreply.github.com> Date: Fri, 4 Sep 2026 16:13:03 -0400 Subject: [PATCH 1/3] fix(controller): release completed local process records --- agentlightning/controller/local_reconciler.py | 6 +- tests/controller/test_local_reconciler.py | 114 +++++++++++++++++- 2 files changed, 114 insertions(+), 6 deletions(-) diff --git a/agentlightning/controller/local_reconciler.py b/agentlightning/controller/local_reconciler.py index 847bfe736..a6fa070b6 100644 --- a/agentlightning/controller/local_reconciler.py +++ b/agentlightning/controller/local_reconciler.py @@ -166,7 +166,8 @@ async def _reconcile_once(self, *, spawn_queued: bool = True) -> None: await self._patch(rollout.rollout_id, RolloutState.RUNNING, last_attempt_id=item.attempt_id) continue - await self._finish_proc(rollout, item) + if await self._finish_proc(rollout, item): + self._rid_to_proc.pop(rollout.rollout_id) now = time.monotonic() for rollout_id, item in list(self._rid_to_proc.items()): @@ -178,8 +179,9 @@ async def _reconcile_once(self, *, spawn_queued: bool = True) -> None: timeout is not None and (now - item.spawned_at) > timeout and await self._kill_process_group(rollout_id, item) + and await self._patch(rollout_id, RolloutState.FAILED, "local subprocess timed out") ): - await self._patch(rollout_id, RolloutState.FAILED, "local subprocess timed out") + self._rid_to_proc.pop(rollout_id) async def _finish_proc(self, rollout: Rollout, item: Proc) -> bool: if rollout.status.state == RolloutState.QUEUING: diff --git a/tests/controller/test_local_reconciler.py b/tests/controller/test_local_reconciler.py index de0d1bda4..1601b1849 100644 --- a/tests/controller/test_local_reconciler.py +++ b/tests/controller/test_local_reconciler.py @@ -3,14 +3,15 @@ """Unit tests for local subprocess reconciliation.""" import asyncio -from unittest.mock import AsyncMock +import time +from unittest.mock import AsyncMock, MagicMock import httpx import pytest from omegaconf import OmegaConf from agentlightning.client import AgentLightningAsyncClient -from agentlightning.controller.local_reconciler import LocalReconciler +from agentlightning.controller.local_reconciler import LocalReconciler, Proc from agentlightning.schemas import Rollout, RolloutConfig, RolloutLifecycleStatus, RolloutState @@ -22,11 +23,13 @@ def _response(json: object) -> httpx.Response: ) -def _reconciler(*, state: RolloutState = RolloutState.QUEUING) -> tuple[LocalReconciler, AsyncMock]: +def _reconciler( + *, state: RolloutState = RolloutState.QUEUING, timeout_seconds: int = 3600 +) -> tuple[LocalReconciler, AsyncMock]: rollout = Rollout( rollout_id="rollout-1", input={"question": "1 + 1"}, - config=RolloutConfig(), + config=RolloutConfig(timeout_seconds=timeout_seconds), status=RolloutLifecycleStatus(state=state, created_at=1.0, updated_at=1.0), ) api = AsyncMock(spec=AgentLightningAsyncClient) @@ -100,3 +103,106 @@ async def test_shutdown_still_fails_running_rollout_without_local_process() -> N "state": "failed", "error_message": "local subprocess is not running", } + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + ("returncode", "expected_status"), + [ + (0, {"state": "succeeded", "last_attempt_id": "attempt-1"}), + (1, {"state": "failed", "error_message": "subprocess exited with code 1"}), + ], +) +async def test_reconcile_removes_completed_process_after_terminal_patch( + returncode: int, expected_status: dict[str, str] +) -> None: + reconciler, api = _reconciler(state=RolloutState.RUNNING) + proc = MagicMock(spec=asyncio.subprocess.Process) + proc.returncode = returncode + reconciler._rid_to_proc["rollout-1"] = Proc("attempt-1", proc, spawned_at=time.monotonic()) + + await reconciler._reconcile_once() + + assert "rollout-1" not in reconciler._rid_to_proc + assert api.patch.await_args.kwargs["json"]["status"] == expected_status + + +@pytest.mark.asyncio +async def test_reconcile_keeps_completed_process_when_terminal_patch_fails_then_retries() -> None: + reconciler, api = _reconciler(state=RolloutState.RUNNING) + proc = MagicMock(spec=asyncio.subprocess.Process) + proc.returncode = 0 + reconciler._rid_to_proc["rollout-1"] = Proc("attempt-1", proc, spawned_at=time.monotonic()) + api.patch.side_effect = [ + httpx.Response(500, request=httpx.Request("PATCH", "http://server/api/rollouts/rollout-1")), + _response({}), + ] + + await reconciler._reconcile_once() + + assert "rollout-1" in reconciler._rid_to_proc + + await reconciler._reconcile_once() + + assert "rollout-1" not in reconciler._rid_to_proc + assert api.patch.await_count == 2 + + +@pytest.mark.asyncio +async def test_reconcile_keeps_running_process() -> None: + reconciler, api = _reconciler(state=RolloutState.RUNNING) + proc = MagicMock(spec=asyncio.subprocess.Process) + proc.returncode = None + reconciler._rid_to_proc["rollout-1"] = Proc("attempt-1", proc, spawned_at=time.monotonic()) + + await reconciler._reconcile_once() + + assert reconciler._rid_to_proc["rollout-1"].proc is proc + api.patch.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_reconcile_removes_timed_out_process_after_failure_patch(monkeypatch: pytest.MonkeyPatch) -> None: + reconciler, api = _reconciler(state=RolloutState.RUNNING, timeout_seconds=1) + proc = MagicMock(spec=asyncio.subprocess.Process) + proc.returncode = None + reconciler._rid_to_proc["rollout-1"] = Proc("attempt-1", proc, spawned_at=time.monotonic() - 2) + kill_process_group = AsyncMock(return_value=True) + monkeypatch.setattr(reconciler, "_kill_process_group", kill_process_group) + + await reconciler._reconcile_once() + + assert "rollout-1" not in reconciler._rid_to_proc + kill_process_group.assert_awaited_once() + assert api.patch.await_args.kwargs["json"]["status"] == { + "state": "failed", + "error_message": "local subprocess timed out", + } + + +@pytest.mark.asyncio +async def test_reconcile_keeps_timed_out_process_when_failure_patch_fails_then_retries( + monkeypatch: pytest.MonkeyPatch, +) -> None: + reconciler, api = _reconciler(state=RolloutState.RUNNING, timeout_seconds=1) + proc = MagicMock(spec=asyncio.subprocess.Process) + proc.returncode = None + reconciler._rid_to_proc["rollout-1"] = Proc("attempt-1", proc, spawned_at=time.monotonic() - 2) + kill_process_group = AsyncMock(return_value=True) + monkeypatch.setattr(reconciler, "_kill_process_group", kill_process_group) + api.patch.side_effect = [ + httpx.Response(500, request=httpx.Request("PATCH", "http://server/api/rollouts/rollout-1")), + _response({}), + ] + + await reconciler._reconcile_once() + + assert "rollout-1" in reconciler._rid_to_proc + + # A successful kill has completed; retry reporting through the exited path. + proc.returncode = -9 + await reconciler._reconcile_once() + + assert "rollout-1" not in reconciler._rid_to_proc + assert api.patch.await_count == 2 + kill_process_group.assert_awaited_once() From 2013a23ce9ffe9451edabf436e68b92b2f97e108 Mon Sep 17 00:00:00 2001 From: YusefSyed <211442445+YusefSyed@users.noreply.github.com> Date: Fri, 4 Sep 2026 16:46:44 -0400 Subject: [PATCH 2/3] fix(controller): preserve timeout outcomes across retries --- agentlightning/controller/local_reconciler.py | 2 ++ tests/controller/test_local_reconciler.py | 31 ++++++++++++++++--- 2 files changed, 28 insertions(+), 5 deletions(-) diff --git a/agentlightning/controller/local_reconciler.py b/agentlightning/controller/local_reconciler.py index a6fa070b6..d510e5f24 100644 --- a/agentlightning/controller/local_reconciler.py +++ b/agentlightning/controller/local_reconciler.py @@ -188,6 +188,8 @@ async def _finish_proc(self, rollout: Rollout, item: Proc) -> bool: patched = await self._patch(rollout.rollout_id, RolloutState.RUNNING, last_attempt_id=item.attempt_id) if not patched: return False + if item.killed: + return await self._patch(rollout.rollout_id, RolloutState.FAILED, "local subprocess timed out") if item.proc.returncode == 0: return await self._patch(rollout.rollout_id, RolloutState.SUCCEEDED, last_attempt_id=item.attempt_id) return await self._patch( diff --git a/tests/controller/test_local_reconciler.py b/tests/controller/test_local_reconciler.py index 1601b1849..a58cfa6c2 100644 --- a/tests/controller/test_local_reconciler.py +++ b/tests/controller/test_local_reconciler.py @@ -111,6 +111,7 @@ async def test_shutdown_still_fails_running_rollout_without_local_process() -> N [ (0, {"state": "succeeded", "last_attempt_id": "attempt-1"}), (1, {"state": "failed", "error_message": "subprocess exited with code 1"}), + (-9, {"state": "failed", "error_message": "subprocess exited with code -9"}), ], ) async def test_reconcile_removes_completed_process_after_terminal_patch( @@ -188,7 +189,11 @@ async def test_reconcile_keeps_timed_out_process_when_failure_patch_fails_then_r proc = MagicMock(spec=asyncio.subprocess.Process) proc.returncode = None reconciler._rid_to_proc["rollout-1"] = Proc("attempt-1", proc, spawned_at=time.monotonic() - 2) - kill_process_group = AsyncMock(return_value=True) + async def kill_process_group(_: str, item: Proc) -> bool: + item.killed = True + proc.returncode = -9 + return True + monkeypatch.setattr(reconciler, "_kill_process_group", kill_process_group) api.patch.side_effect = [ httpx.Response(500, request=httpx.Request("PATCH", "http://server/api/rollouts/rollout-1")), @@ -199,10 +204,26 @@ async def test_reconcile_keeps_timed_out_process_when_failure_patch_fails_then_r assert "rollout-1" in reconciler._rid_to_proc - # A successful kill has completed; retry reporting through the exited path. - proc.returncode = -9 await reconciler._reconcile_once() assert "rollout-1" not in reconciler._rid_to_proc - assert api.patch.await_count == 2 - kill_process_group.assert_awaited_once() + assert [call.kwargs["json"]["status"] for call in api.patch.await_args_list] == [ + {"state": "failed", "error_message": "local subprocess timed out"}, + {"state": "failed", "error_message": "local subprocess timed out"}, + ] + + +@pytest.mark.asyncio +async def test_reconcile_finishes_queued_completed_process_after_running_transition() -> None: + reconciler, api = _reconciler() + proc = MagicMock(spec=asyncio.subprocess.Process) + proc.returncode = 0 + reconciler._rid_to_proc["rollout-1"] = Proc("attempt-1", proc, spawned_at=time.monotonic()) + + await reconciler._reconcile_once() + + assert "rollout-1" not in reconciler._rid_to_proc + assert [call.kwargs["json"]["status"] for call in api.patch.await_args_list] == [ + {"state": "running", "last_attempt_id": "attempt-1"}, + {"state": "succeeded", "last_attempt_id": "attempt-1"}, + ] From 5a007709e2851d80f9cae910098436774786392a Mon Sep 17 00:00:00 2001 From: YusefSyed <211442445+YusefSyed@users.noreply.github.com> Date: Mon, 14 Sep 2026 14:36:41 -0400 Subject: [PATCH 3/3] fix(controller): preserve history in timeout retry fix --- agentlightning/controller/local_reconciler.py | 7 +- tests/controller/test_local_reconciler.py | 220 +++++++++++------- 2 files changed, 141 insertions(+), 86 deletions(-) diff --git a/agentlightning/controller/local_reconciler.py b/agentlightning/controller/local_reconciler.py index d510e5f24..fcbd3a3ab 100644 --- a/agentlightning/controller/local_reconciler.py +++ b/agentlightning/controller/local_reconciler.py @@ -166,8 +166,7 @@ async def _reconcile_once(self, *, spawn_queued: bool = True) -> None: await self._patch(rollout.rollout_id, RolloutState.RUNNING, last_attempt_id=item.attempt_id) continue - if await self._finish_proc(rollout, item): - self._rid_to_proc.pop(rollout.rollout_id) + await self._finish_proc(rollout, item) now = time.monotonic() for rollout_id, item in list(self._rid_to_proc.items()): @@ -179,15 +178,15 @@ async def _reconcile_once(self, *, spawn_queued: bool = True) -> None: timeout is not None and (now - item.spawned_at) > timeout and await self._kill_process_group(rollout_id, item) - and await self._patch(rollout_id, RolloutState.FAILED, "local subprocess timed out") ): - self._rid_to_proc.pop(rollout_id) + await self._patch(rollout_id, RolloutState.FAILED, "local subprocess timed out") async def _finish_proc(self, rollout: Rollout, item: Proc) -> bool: if rollout.status.state == RolloutState.QUEUING: patched = await self._patch(rollout.rollout_id, RolloutState.RUNNING, last_attempt_id=item.attempt_id) if not patched: return False + # Normal timeouts reach here; shutdown kills happen after final reconciliation. if item.killed: return await self._patch(rollout.rollout_id, RolloutState.FAILED, "local subprocess timed out") if item.proc.returncode == 0: diff --git a/tests/controller/test_local_reconciler.py b/tests/controller/test_local_reconciler.py index a58cfa6c2..1a67ac456 100644 --- a/tests/controller/test_local_reconciler.py +++ b/tests/controller/test_local_reconciler.py @@ -3,8 +3,10 @@ """Unit tests for local subprocess reconciliation.""" import asyncio +import signal import time -from unittest.mock import AsyncMock, MagicMock +from typing import cast +from unittest.mock import AsyncMock, Mock import httpx import pytest @@ -24,7 +26,9 @@ def _response(json: object) -> httpx.Response: def _reconciler( - *, state: RolloutState = RolloutState.QUEUING, timeout_seconds: int = 3600 + *, + state: RolloutState = RolloutState.QUEUING, + timeout_seconds: int = 3600, ) -> tuple[LocalReconciler, AsyncMock]: rollout = Rollout( rollout_id="rollout-1", @@ -44,6 +48,26 @@ def _reconciler( return LocalReconciler(api, config), api +class _Process: + def __init__(self, returncode: int | None) -> None: + self.returncode = returncode + self.pid = 1234 + self.wait = AsyncMock() + + +def _proc(*, returncode: int | None, killed: bool = False, attempt_id: str = "attempt-1") -> tuple[Proc, _Process]: + process = _Process(returncode) + return ( + Proc( + attempt_id=attempt_id, + proc=cast(asyncio.subprocess.Process, process), + spawned_at=time.monotonic(), + killed=killed, + ), + process, + ) + + @pytest.mark.asyncio async def test_shutdown_does_not_spawn_queued_rollouts(monkeypatch: pytest.MonkeyPatch) -> None: reconciler, api = _reconciler() @@ -105,125 +129,157 @@ async def test_shutdown_still_fails_running_rollout_without_local_process() -> N } +@pytest.mark.asyncio +@pytest.mark.parametrize("failed_terminal_patches", [0, 1, 2]) +@pytest.mark.parametrize("returncode", [-9, 0]) +async def test_timeout_reconciliation_retries_and_retains_process_record( + monkeypatch: pytest.MonkeyPatch, + failed_terminal_patches: int, + returncode: int, +) -> None: + reconciler, api = _reconciler(state=RolloutState.RUNNING, timeout_seconds=1) + rollout = Rollout.model_validate(api.get.return_value.json()[0]) + item, process = _proc(returncode=None) + reconciler._rid_to_proc[rollout.rollout_id] = item + item.spawned_at = 0.0 + api.patch.side_effect = [ + *[httpx.ConnectError("server unavailable") for _ in range(failed_terminal_patches)], + _response({}), + ] + killpg = Mock() + monkeypatch.setattr("agentlightning.controller.local_reconciler.os.killpg", killpg) + + async def complete_wait() -> None: + process.returncode = returncode + + process.wait.side_effect = complete_wait + spawn_for = AsyncMock(return_value=True) + monkeypatch.setattr(reconciler, "_spawn_for", spawn_for) + + for _ in range(failed_terminal_patches + 1): + await reconciler._reconcile_once() + + assert item.killed + assert reconciler._rid_to_proc[rollout.rollout_id] is item + assert item.attempt_id == "attempt-1" + process.wait.assert_awaited_once_with() + killpg.assert_called_once_with(process.pid, signal.SIGKILL) + spawn_for.assert_not_awaited() + assert [call.kwargs["json"]["status"] for call in api.patch.await_args_list] == [ + {"state": "failed", "error_message": "local subprocess timed out"} + ] * (failed_terminal_patches + 1) + + api.get.return_value = _response([]) + await reconciler._reconcile_once() + assert reconciler._rid_to_proc[rollout.rollout_id] is item + + @pytest.mark.asyncio @pytest.mark.parametrize( - ("returncode", "expected_status"), + ("returncode", "expected"), [ (0, {"state": "succeeded", "last_attempt_id": "attempt-1"}), (1, {"state": "failed", "error_message": "subprocess exited with code 1"}), (-9, {"state": "failed", "error_message": "subprocess exited with code -9"}), ], ) -async def test_reconcile_removes_completed_process_after_terminal_patch( - returncode: int, expected_status: dict[str, str] +@pytest.mark.parametrize("failed_terminal_patches", [0, 1]) +async def test_unmarked_terminal_process_retries_and_retains_record( + returncode: int, + expected: dict[str, str], + failed_terminal_patches: int, ) -> None: reconciler, api = _reconciler(state=RolloutState.RUNNING) - proc = MagicMock(spec=asyncio.subprocess.Process) - proc.returncode = returncode - reconciler._rid_to_proc["rollout-1"] = Proc("attempt-1", proc, spawned_at=time.monotonic()) - - await reconciler._reconcile_once() - - assert "rollout-1" not in reconciler._rid_to_proc - assert api.patch.await_args.kwargs["json"]["status"] == expected_status - - -@pytest.mark.asyncio -async def test_reconcile_keeps_completed_process_when_terminal_patch_fails_then_retries() -> None: - reconciler, api = _reconciler(state=RolloutState.RUNNING) - proc = MagicMock(spec=asyncio.subprocess.Process) - proc.returncode = 0 - reconciler._rid_to_proc["rollout-1"] = Proc("attempt-1", proc, spawned_at=time.monotonic()) + rollout = Rollout.model_validate(api.get.return_value.json()[0]) + item, _ = _proc(returncode=returncode) + reconciler._rid_to_proc[rollout.rollout_id] = item api.patch.side_effect = [ - httpx.Response(500, request=httpx.Request("PATCH", "http://server/api/rollouts/rollout-1")), + *[httpx.ConnectError("server unavailable") for _ in range(failed_terminal_patches)], _response({}), ] - await reconciler._reconcile_once() - - assert "rollout-1" in reconciler._rid_to_proc + for _ in range(failed_terminal_patches + 1): + await reconciler._reconcile_once() + assert reconciler._rid_to_proc[rollout.rollout_id] is item + assert [call.kwargs["json"]["status"] for call in api.patch.await_args_list] == [expected] * ( + failed_terminal_patches + 1 + ) + api.get.return_value = _response([]) await reconciler._reconcile_once() - - assert "rollout-1" not in reconciler._rid_to_proc - assert api.patch.await_count == 2 + assert reconciler._rid_to_proc[rollout.rollout_id] is item @pytest.mark.asyncio -async def test_reconcile_keeps_running_process() -> None: +@pytest.mark.parametrize("patch_fails", [False, True]) +async def test_run_stops_then_shutdown_kills_after_final_reconcile( + monkeypatch: pytest.MonkeyPatch, + patch_fails: bool, +) -> None: reconciler, api = _reconciler(state=RolloutState.RUNNING) - proc = MagicMock(spec=asyncio.subprocess.Process) - proc.returncode = None - reconciler._rid_to_proc["rollout-1"] = Proc("attempt-1", proc, spawned_at=time.monotonic()) + item, process = _proc(returncode=None) + reconciler._rid_to_proc["rollout-1"] = item + events: list[str] = [] - await reconciler._reconcile_once() + async def get(*args: object, **kwargs: object) -> httpx.Response: + del args, kwargs + events.append("reconcile") + reconciler.stop() + return _response([Rollout.model_validate(api.get.return_value.json()[0]).model_dump(mode="json")]) - assert reconciler._rid_to_proc["rollout-1"].proc is proc - api.patch.assert_not_awaited() + async def complete_wait() -> None: + process.returncode = -9 + process.wait.side_effect = complete_wait + killpg = Mock(side_effect=lambda *args: events.append("kill")) + monkeypatch.setattr("agentlightning.controller.local_reconciler.os.killpg", killpg) + api.get.side_effect = get + if patch_fails: + api.patch.side_effect = httpx.ConnectError("server unavailable") -@pytest.mark.asyncio -async def test_reconcile_removes_timed_out_process_after_failure_patch(monkeypatch: pytest.MonkeyPatch) -> None: - reconciler, api = _reconciler(state=RolloutState.RUNNING, timeout_seconds=1) - proc = MagicMock(spec=asyncio.subprocess.Process) - proc.returncode = None - reconciler._rid_to_proc["rollout-1"] = Proc("attempt-1", proc, spawned_at=time.monotonic() - 2) - kill_process_group = AsyncMock(return_value=True) - monkeypatch.setattr(reconciler, "_kill_process_group", kill_process_group) + await reconciler.run() - await reconciler._reconcile_once() - - assert "rollout-1" not in reconciler._rid_to_proc - kill_process_group.assert_awaited_once() + assert events == ["reconcile", "reconcile", "kill"] + assert item.killed + assert reconciler._rid_to_proc["rollout-1"] is item assert api.patch.await_args.kwargs["json"]["status"] == { "state": "failed", - "error_message": "local subprocess timed out", + "error_message": "local controller shutdown", } @pytest.mark.asyncio -async def test_reconcile_keeps_timed_out_process_when_failure_patch_fails_then_retries( +async def test_timeout_during_run_retries_original_reason_during_shutdown_without_rekill( monkeypatch: pytest.MonkeyPatch, ) -> None: reconciler, api = _reconciler(state=RolloutState.RUNNING, timeout_seconds=1) - proc = MagicMock(spec=asyncio.subprocess.Process) - proc.returncode = None - reconciler._rid_to_proc["rollout-1"] = Proc("attempt-1", proc, spawned_at=time.monotonic() - 2) - async def kill_process_group(_: str, item: Proc) -> bool: - item.killed = True - proc.returncode = -9 - return True - - monkeypatch.setattr(reconciler, "_kill_process_group", kill_process_group) - api.patch.side_effect = [ - httpx.Response(500, request=httpx.Request("PATCH", "http://server/api/rollouts/rollout-1")), - _response({}), - ] + rollout = Rollout.model_validate(api.get.return_value.json()[0]) + item, process = _proc(returncode=None) + item.spawned_at = 0.0 + reconciler._rid_to_proc[rollout.rollout_id] = item + events: list[str] = [] - await reconciler._reconcile_once() + async def get(*args: object, **kwargs: object) -> httpx.Response: + del args, kwargs + events.append("reconcile") + reconciler.stop() + return _response([rollout.model_dump(mode="json")]) - assert "rollout-1" in reconciler._rid_to_proc + async def complete_wait() -> None: + process.returncode = -9 - await reconciler._reconcile_once() + process.wait.side_effect = complete_wait + killpg = Mock(side_effect=lambda *args: events.append("kill")) + monkeypatch.setattr("agentlightning.controller.local_reconciler.os.killpg", killpg) + api.get.side_effect = get + api.patch.side_effect = [httpx.ConnectError("server unavailable"), _response({})] + + await reconciler.run() - assert "rollout-1" not in reconciler._rid_to_proc + assert events == ["reconcile", "kill", "reconcile"] + assert reconciler._rid_to_proc[rollout.rollout_id] is item + process.wait.assert_awaited_once_with() assert [call.kwargs["json"]["status"] for call in api.patch.await_args_list] == [ {"state": "failed", "error_message": "local subprocess timed out"}, {"state": "failed", "error_message": "local subprocess timed out"}, ] - - -@pytest.mark.asyncio -async def test_reconcile_finishes_queued_completed_process_after_running_transition() -> None: - reconciler, api = _reconciler() - proc = MagicMock(spec=asyncio.subprocess.Process) - proc.returncode = 0 - reconciler._rid_to_proc["rollout-1"] = Proc("attempt-1", proc, spawned_at=time.monotonic()) - - await reconciler._reconcile_once() - - assert "rollout-1" not in reconciler._rid_to_proc - assert [call.kwargs["json"]["status"] for call in api.patch.await_args_list] == [ - {"state": "running", "last_attempt_id": "attempt-1"}, - {"state": "succeeded", "last_attempt_id": "attempt-1"}, - ]