From d736564eb1b4d09b0016b71a188dab1eb2a1e669 Mon Sep 17 00:00:00 2001 From: Quanzheng Long Date: Thu, 7 May 2026 10:31:18 -0700 Subject: [PATCH] test(langgraph): de-flake heartbeat progress test (#7735) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Summary De-flake `test_arun_with_retry_timeout_observer_emits_progress_on_heartbeat` — the test was hitting a CI-runner-load-sensitive race where the idle-timeout watchdog could fire before the task body's first await ran. ## Root cause `_TimedAttemptScope.__init__` sets `_last_progress = time.monotonic()` immediately, but the watchdog itself doesn't start polling until *after* `wrap_config` and task scheduling. Under heavy CI load that gap can grow large enough that: ``` T₀ scope.__init__() → _last_progress = T₀ … some scheduling slack … Tₙ watchdog runs, computes remaining = T₀ + 0.2 − Tₙ ≤ 0 → TimeoutError fires ``` The error reports `elapsed: 0.000s` because `elapsed` is measured from the post-scheduling `start` (≈Tₙ), not from `_last_progress` (T₀). The previous test set `idle_timeout=0.2s`, which left almost no headroom for that scheduling slack. ## Fix (test-side only — no production change) - **Heartbeat at task-body entry**: `runtime.heartbeat()` is now called before the first `await asyncio.sleep(...)`, which resets `_last_progress` to "now" the moment the task body actually starts running. This eliminates the scope-init-to-first-await gap as a flake source. - **Idle timeout 0.2s → 1.0s**: gives ~5× headroom over the ~400ms task duration, so scheduling pressure stays comfortably within budget. ## Why test-side instead of fixing the production race The proper production fix would be to set `_last_progress` at watchdog-entry time rather than at scope-init time. That's a behaviour change in the retry/timeout machinery and out of scope for a flaky-test fix. The two test-side defenses make this particular test stable without touching production semantics; the underlying race in `_TimedAttemptScope` is worth a separate follow-up. ## Test plan - 10/10 repeated local runs pass: ``` uv run pytest tests/test_retry.py::test_arun_with_retry_timeout_observer_emits_progress_on_heartbeat --count=10 ``` - All assertions still meaningful: still verifies start/finish events, at least one progress event, rate-limited progress count (≤ total events), and per-event metadata (task_name, attempt, idle_timeout_secs, progress_at). Co-authored-by: Cursor --- libs/langgraph/tests/test_retry.py | 21 +++++++++++++++++---- 1 file changed, 17 insertions(+), 4 deletions(-) diff --git a/libs/langgraph/tests/test_retry.py b/libs/langgraph/tests/test_retry.py index 8ee515e8e..f5d4d74e0 100644 --- a/libs/langgraph/tests/test_retry.py +++ b/libs/langgraph/tests/test_retry.py @@ -1674,15 +1674,28 @@ async def test_arun_with_retry_timeout_observer_tracks_attempts(): async def test_arun_with_retry_timeout_observer_emits_progress_on_heartbeat(): events: list = [] + # `_TimedAttemptScope.__init__` sets `_last_progress` to `time.monotonic()`, + # but the watchdog itself doesn't start running until after `wrap_config` + # and task scheduling — under CI load that gap can be large enough to eat + # the entire idle window before the task body's first await even runs. We + # defend against that by: + # 1. Using a generous idle_timeout so scheduling slack stays well within it. + # 2. Calling `runtime.heartbeat()` BEFORE the first sleep, which resets + # `_last_progress` to "now" the moment the task body actually starts. + idle_timeout_s = 1.0 + class HeartbeatProc: async def ainvoke(self, input, config): runtime = config[CONF][CONFIG_KEY_RUNTIME] + runtime.heartbeat() # reset the idle clock at task-body entry for _ in range(8): await asyncio.sleep(0.05) runtime.heartbeat() return "ok" - task = _make_task(HeartbeatProc(), timeout=_idle_timeout(0.2), name="heartbeat") + task = _make_task( + HeartbeatProc(), timeout=_idle_timeout(idle_timeout_s), name="heartbeat" + ) task.config[CONF][CONFIG_KEY_TIMED_ATTEMPT_OBSERVER] = events.append assert await arun_with_retry(task, retry_policy=None) == "ok" @@ -1691,13 +1704,13 @@ async def test_arun_with_retry_timeout_observer_emits_progress_on_heartbeat(): assert by_event[-1] == "finish" progress = [ev for ev in events if ev.event == "progress"] assert progress, "expected at least one progress event from heartbeat" - # Rate limit is `idle_timeout / 4` = 0.05s; with 8 heartbeats spaced ~0.05s - # we should see at most ~one progress event per heartbeat (well below 8). + # Rate limit is `idle_timeout / 4` = 0.25s; with the task running for + # ~400ms we expect 1–2 progress events (well below the 9 heartbeats). assert len(progress) <= len(by_event) for ev in progress: assert ev.context.task_name == "heartbeat" assert ev.context.attempt == 1 - assert ev.context.idle_timeout_secs == 0.2 + assert ev.context.idle_timeout_secs == idle_timeout_s assert isinstance(ev.progress_at, datetime)