From f7e0ae2aff80c7fb181f8ddd030287ccab8c5e1b Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Tue, 15 Sep 2026 11:34:46 +0530 Subject: [PATCH] test(supermemory): pin shutdown flush waits for in-flight write, never re-sends Answers the open review ask on #109359 with the corrected contract: a hung remote add_memory is bounded by the SDK timeout (empirically 1.03s at timeout=1.0, max_retries=0), so shutdown's flush waiting on the capture lock is bounded, not forever. The guard: while the worker owns the write (blocked inside add_memory), a concurrent shutdown must wait and must not re-send the same pending batch. Mutation-checked: lock -> nullcontext makes this test and the session-switch interleaving test fail. --- .../memory/test_supermemory_provider.py | 31 +++++++++++++++++++ 1 file changed, 31 insertions(+) diff --git a/tests/plugins/memory/test_supermemory_provider.py b/tests/plugins/memory/test_supermemory_provider.py index 95cf700f14..7276d8022e 100644 --- a/tests/plugins/memory/test_supermemory_provider.py +++ b/tests/plugins/memory/test_supermemory_provider.py @@ -259,6 +259,37 @@ def test_failed_switch_flush_is_retried_at_shutdown(provider, frozen_capture_clo assert provider._pending_turns == [] +def test_shutdown_waits_for_inflight_write_and_does_not_resend(provider, monkeypatch): + """A hung remote write bounds the shutdown flush: it waits for the in-flight owner (the SDK call + is timeout-bounded, never an unconditioned wait) and never re-sends a batch another thread owns. + Regression for the review ask on #109359: without the lock, a concurrent flush re-snapshots the + same pending batch and duplicates the remote append.""" + entered, release = threading.Event(), threading.Event() + real_add = provider._client.add_memory + + def slow_add(content, metadata=None, **kwargs): + entered.set() + assert release.wait(timeout=2) + return real_add(content, metadata=metadata, **kwargs) + + monkeypatch.setattr(provider._client, "add_memory", slow_add) + worker = threading.Thread(target=provider.sync_turn, args=("A", "a"), kwargs={"session_id": "session-1"}) + worker.start() + assert entered.wait(timeout=2) # worker owns the write and is blocked inside add_memory + + flusher = threading.Thread(target=provider.shutdown, name="shutdown-flusher") + flusher.start() + flusher.join(timeout=0.2) + assert flusher.is_alive() # shutdown's flush waits for the capture lock, it does not duplicate the write + + release.set() + worker.join(timeout=2) + flusher.join(timeout=2) + assert not worker.is_alive() and not flusher.is_alive() + assert len(provider._client.add_calls) == 1 # A went out exactly once despite the concurrent shutdown + assert provider._pending_turns == [] + + def test_sync_turn_drops_inline_image_payloads(provider, frozen_capture_clock): blob = "A" * 4096 provider.sync_turn(f"describe this data:image/png;base64,{blob}", "a screenshot", session_id="session-1")