fix(sandbox): hold the teardown lease on every del: stop, and pin the claims that had no test

90936b49 said `_held_teardown_lease` wrapped "both" `_backend.destroy()` call
sites. There are three. `_drop_unhealthy_sandbox` marks `del:` and then blocks on
the same unbounded stop, and it untracks *before* claiming, so `_renew_owned_leases`
cannot see the id either — nothing refreshed the marker. Reproduced against a real
redis: the peer's `take()` succeeds 1.0s into a 2.5s stop. Same window, third path.

That miss came from the habit the rest of this commit addresses: a property
asserted in prose, with no test that could falsify it. Auditing every load-bearing
claim in this feature — AGENTS.md, the store docstrings, the provider's design
comments — against the test that would go red turned up several more, each
verified by mutating the code and watching the suite stay green.

Tests that could not fail:

- `test_reconcile_fails_closed_when_ownership_unknown` reached the grace gate, not
  the claim. A bare MagicMock answers `owner()` with a truthy mock, so the
  container read as peer-owned and deferred; `claim()` was never called. It stayed
  green with `_claim_ownership` failing open. Adding the grace ahead of the claim
  is what hollowed it out — inserting a gate can silently disarm the tests for
  the gate behind it.
- `test_adoption_grace_restarts_when_a_live_owner_republishes` never distinguished
  reset from pause. Those diverge only on a *second* lapse, which it never drove,
  so it passed with the reset deleted.

Claims with no test at all, each now pinned (mutation → red, per test):

- `destroy()`, `_evict_oldest_warm`, `_reclaim_warm_pool_sandbox`,
  `_register_created_sandbox` and `shutdown()`'s warm loop were each the one
  untested sibling of an "every path does X" enumeration. `shutdown()` was never
  driven with a non-empty warm pool, so a loop bypassing the ownership claim —
  stopping a live peer's container on our exit — went unnoticed.
- Renewal's unknown-is-not-lost rule, the single deliberate exception to
  fail-closed. Inverting it drops every active and warm sandbox on every instance
  the moment the store blinks.
- Both hops of the stream-bridge redis inference. Deleting either left the suite
  green while every config.yaml-native multi-instance deployment silently fell
  back to memory — #4206 reopened on exactly the deployments the inference exists
  for.

Claims narrowed instead, because they promised more than the code delivers:

- "run against both backends ... cannot drift" — CI provisions no redis, so the
  merge gate runs the memory tier only and the Lua never executes there.
- "Every destroy path claims before untracking" — `_drop_unhealthy_sandbox`
  untracks first, deliberately, under its `expected_info` TOCTOU guard.
- "Atomic: concurrent claims from different instances cannot both succeed" — true
  via Lua on redis, vacuous on the single-instance memory store, and pinned by
  neither, since the contract suite drives sequential calls. A concurrency test
  against the memory store would make the claim look covered while the mechanism
  that carries it still never runs in CI.
This commit is contained in:
Aari
2026-07-17 01:31:15 +08:00
parent 90936b49aa
commit 0d2377b2a2
4 changed files with 377 additions and 10 deletions
+3 -3
View File
@@ -362,13 +362,13 @@ Proxied through nginx: `/api/langgraph/*` → Gateway LangGraph-compatible runti
- **Cross-instance ownership store** (`aio_sandbox/ownership/`, #4206): gateway instances sharing a container backend coordinate container ownership through a pluggable lease store, selected by `sandbox.ownership.type` (`memory` | `redis`) and resolved like `stream_bridge` (`factory.py`, lazy per-branch import, `redis` optional extra, `DEER_FLOW_SANDBOX_OWNERSHIP_REDIS_URL` env escape hatch; a set `DEER_FLOW_STREAM_BRIDGE_REDIS_URL` implies a multi-instance deployment and infers `redis`). `memory` is single-instance only and declares `supports_cross_process = False`.
- **A lease answers "who reaps this container", not "who may use it".** That splits the interface in two: `take()` transfers ownership on the **acquire** path (a container is deterministic per user/thread, so consecutive turns legitimately land on different instances — a conditional claim there would strand the thread until the previous lease expired), while `claim()` succeeds only if the container is unowned or already ours and gates every **adopt/reap** path. `release()` never clears a peer's lease.
- **A lease carries a state, and that is what makes the destroy window safe.** `own:` = responsible for this container; `del:` = tearing it down (`claim(..., for_destroy=True)`). `take()` is refused against a `del:` lease, so a container cannot be re-acquired between a destroy path's claim and its container stop. Without the two states an unconditional `take()` would silently overwrite the destroyer's claim and the peer's stop would land on a container the new owner had already handed to an agent — i.e. #4206 again. That pairing is what replaced the previous same-host `flock` guard, which is gone; Redis makes the scope genuinely multi-instance instead of same-host. A destroyer that dies mid-stop leaves a `del:` marker that lapses with the TTL. On the acquire path a refused take raises `SandboxBeingDestroyedError`: the reuse/reclaim paths drop the container and cold-start, and the discover path propagates (falling through to create would collide with the not-yet-removed container name).
- **The `del:` state has to be *held* for the stop, not just written before it.** The two states alone do not make `flock` redundant — a held lock cannot expire, whereas a lease can, and `claim(..., for_destroy=True)` writes the marker with the ordinary lease TTL. Nothing else refreshes it: `renew()` extends only `own:` and deliberately reports a teardown as `LOST`, and the destroy paths drop the sandbox from the maps `_renew_owned_leases` iterates. So a container stop that outlived the TTL let the marker lapse, a peer's `take()` succeeded against the still-running container, and the stop then landed on the turn that had just been handed it — the exact window `del:` exists to close, reopened by its own expiry. `_held_teardown_lease` wraps both `_backend.destroy()` call sites and re-claims the marker every `renewal_interval_seconds` until the stop returns. This needs no abnormal backend: the schema bounds only `renewal_interval_seconds` (> 0) and `ttl_multiplier` (>= 2), so a legal config puts the TTL below a normal container stop, and `LocalContainerBackend._stop_container` passes no `timeout` to `subprocess.run`, so a wedged daemon blocks unbounded even at the default 120s TTL. The TTL stays finite on purpose — the heartbeat dies with the process, so a destroyer that crashes mid-stop still releases the container one TTL later instead of marking it undestroyable forever. Raising a separate teardown TTL instead would only be sufficient if it were bounded above every backend's real stop deadline.
- **Fail-closed both directions.** Establishment: a sandbox whose ownership cannot be published is never handed out (a just-created container is destroyed rather than leaked as an adoptable orphan) — acquiring raises `OwnershipBackendError`, matching the stream bridge's fail-hard v1 policy. Reaping: a store that cannot answer is treated as peer-owned, so an outage never turns live peer containers into orphans. Every destroy path claims **before** untracking, so a refused claim cannot leave a container running and invisible to every map.
- **The `del:` state has to be *held* for the stop, not just written before it.** The two states alone do not make `flock` redundant — a held lock cannot expire, whereas a lease can, and `claim(..., for_destroy=True)` writes the marker with the ordinary lease TTL. Nothing else refreshes it: `renew()` extends only `own:` and deliberately reports a teardown as `LOST`, and the destroy paths drop the sandbox from the maps `_renew_owned_leases` iterates. So a container stop that outlived the TTL let the marker lapse, a peer's `take()` succeeded against the still-running container, and the stop then landed on the turn that had just been handed it — the exact window `del:` exists to close, reopened by its own expiry. `_held_teardown_lease` wraps **every** `del:`-marked stop — `_destroy_warm_entry`, `destroy()`, and `_drop_unhealthy_sandbox` — and re-claims the marker every `renewal_interval_seconds` until the stop returns. `_drop_unhealthy_sandbox` needs it most: it untracks *before* claiming (under its `expected_info` TOCTOU guard), so `_renew_owned_leases` cannot see the id either. This needs no abnormal backend: the schema bounds only `renewal_interval_seconds` (> 0) and `ttl_multiplier` (>= 2), so a legal config puts the TTL below a normal container stop, and `LocalContainerBackend._stop_container` passes no `timeout` to `subprocess.run`, so a wedged daemon blocks unbounded even at the default 120s TTL. The TTL stays finite on purpose — the heartbeat dies with the process, so a destroyer that crashes mid-stop still releases the container one TTL later instead of marking it undestroyable forever. Raising a separate teardown TTL instead would only be sufficient if it were bounded above every backend's real stop deadline.
- **Fail-closed both directions.** Establishment: a sandbox whose ownership cannot be published is never handed out (a just-created container is destroyed rather than leaked as an adoptable orphan) — acquiring raises `OwnershipBackendError`, matching the stream bridge's fail-hard v1 policy. Reaping: a store that cannot answer is treated as peer-owned, so an outage never turns live peer containers into orphans. **Renewal is the deliberate exception**: an unanswerable store there means *unknown*, not lost, so `_refresh_ownership` keeps the sandbox and retries — failing closed on that path would evict every live sandbox on every instance the moment the store blinked. The TTL still bounds how long a genuinely dead owner holds a lease. Both paths that stop a container they still track — `destroy()` and `_destroy_warm_entry` claim **before** untracking, so a refused claim cannot leave a container running and untracked. (`_drop_unhealthy_sandbox` untracks first, under its `expected_info` TOCTOU guard, then claims before the stop; a refused claim there leaves the container to the next reconcile, which re-adopts it after the grace.)
- **Renewal is independent of `idle_timeout`** (`_start_lease_renewal`, own daemon thread; TTL = `renewal_interval_seconds × ttl_multiplier`). Renewal used to ride on the idle checker, which `__init__` only starts when `idle_timeout > 0` — so `idle_timeout: 0` ("keep warm VMs until shutdown", a documented config) let every lease lapse. Liveness and reaping must not share a switch. Renewal covers warm entries as well as active ones; losing a lease drops the sandbox from this instance's maps **without touching the container** (`_forget_lost_sandbox`) — destroying it there would be the very cross-instance kill this store prevents.
- **`renew()` distinguishes lapsed from lost** (`RenewOutcome`), and the two must not be collapsed. `LAPSED` means the lease is simply absent — nobody took it — so `_refresh_ownership` re-establishes it; `LOST` means a peer holds it and it is never re-taken. Treating an absent lease as lost meant a Redis restart without persistence (every key gone) evicted every in-flight sandbox on every instance at once.
- **An absent lease means the same thing on both paths, and reconciliation must say so too.** The `LAPSED` rule above only covers an owner renewing its *own* lease; on its own it does not make state loss safe, because reconciliation reads the same absent key as "orphan, adopt". After a Redis flush (restart without persistence, or eviction under `maxmemory`) every owner is alive and merely pre-renewal-tick, so whichever instance reconciles first would adopt every live container, each real owner's next renewal would report `LOST`, and it would drop a sandbox mid-turn for the adopter to idle-destroy — #4206 through the back door. `_adoptable_after_grace` closes it: an untracked container must be seen unowned (`owner()`, a read-only peek — the atomic `claim()` is still what actually gates adoption) across a full lease TTL before it can be adopted, tracked per container in `_unowned_since`. That rebuilds the delay the flush erased — a live owner republishes within one renewal interval, shorter than the TTL by construction (`ttl_multiplier >= 2`) — while a genuinely crashed owner never republishes, so its containers are still adopted one grace later rather than leaking. A republished lease **resets** the grace; a pausing-only timer would still expire over a live owner's lease. The grace is skipped when `supports_cross_process` is `False`: no peer can hold a lease such a store would show us, so single-instance deployments keep instant orphan cleanup, and a grace could not help a multi-worker gateway on `memory` anyway (peers are invisible to each other's leases with or without it).
- **The `memory` store is single-instance only** and says so via `supports_cross_process = False`; the provider logs a warning at startup when the configured store cannot see peers. A multi-worker gateway on `memory` has no cross-process coordination at all — same contract as `stream_bridge`'s memory backend. This is why the redis inference matters: it reads `app_config.stream_bridge` **and** the env var, in the same order the bridge's own resolver does, so any deployment already pointing the bridge at Redis (i.e. every multi-instance one) gets a redis ownership store without extra config.
- `get()` stays a pure in-memory lookup and must never call the store (that is blocking filesystem/network IO on the event loop); anchored by `tests/blocking_io/test_aio_sandbox_get.py`, which injects a deliberately-blocking probe store so the anchor keeps its teeth regardless of the configured backend. Tests: `tests/test_sandbox_ownership_store.py` (store contract, run against both backends — the redis tier is opt-in via `DEER_FLOW_TEST_REDIS_URL` and self-skips; there is no fake-redis tier because the redis exclusion lives in Lua a fake would not execute) and `tests/test_sandbox_orphan_reconciliation.py` (provider behaviour, two providers sharing one store).
- `get()` stays a pure in-memory lookup and must never call the store (that is blocking filesystem/network IO on the event loop); anchored by `tests/blocking_io/test_aio_sandbox_get.py`, which injects a deliberately-blocking probe store so the anchor keeps its teeth regardless of the configured backend. Tests: `tests/test_sandbox_ownership_store.py` (store contract, defined once for **both** backends — but the redis tier is `@pytest.mark.integration` + opt-in via `DEER_FLOW_TEST_REDIS_URL` and self-skips, and **CI provisions no redis**, so the merge gate runs the memory tier only and the Lua scripts never execute there; drift between the backends is caught only when the suite runs against a live redis. There is no fake-redis tier because the redis exclusion lives in Lua a fake would not execute) and `tests/test_sandbox_orphan_reconciliation.py` (provider behaviour, two providers sharing one store).
- `BoxliteProvider` (`packages/harness/deerflow/community/boxlite/`) - BoxLite micro-VM isolation. The `boxlite` runtime is optional (`deerflow-harness[boxlite]`) and lazy-imported only when this provider is selected. The provider owns one private asyncio event loop on a daemon thread because BoxLite handles are loop-affine; sync `Sandbox` calls marshal onto that loop with `run_coroutine_threadsafe`.
Boxes are named deterministically from `user_id:thread_id`, released into an in-process warm pool after each agent turn, and reclaimed only by the same user/thread. Warm-pool health checks use a short explicit timeout and forward that timeout through both BoxLite `exec(timeout=...)` and the private-loop `.result(timeout)` bridge so a hung VM cannot pin the per-thread acquire lock indefinitely.
`sandbox.replicas` caps active + warm VMs per gateway process; if capacity is exhausted, only warm-pool VMs are evicted. `sandbox.idle_timeout` stops idle warm VMs after the configured seconds. `reset()` is intentionally a lightweight registry clear for `reset_sandbox_provider()` and does not close boxes, stop the idle reaper, or close the private loop; full teardown remains `shutdown()`.
@@ -1130,7 +1130,11 @@ class AioSandboxProvider(WarmPoolLifecycleMixin[SandboxInfo], SandboxProvider):
# in which case stopping it is the cross-instance kill again.
if self._claim_ownership(sandbox_id, for_destroy=True):
try:
self._backend.destroy(info)
# Held like the other two stop paths: this one untracks before
# claiming, so `_renew_owned_leases` cannot see the id either
# and nothing else would refresh the marker.
with self._held_teardown_lease(sandbox_id):
self._backend.destroy(info)
except Exception as e:
logger.warning(f"Error destroying unhealthy sandbox {sandbox_id}: {e}")
self._release_ownership(sandbox_id)
@@ -111,8 +111,16 @@ class SandboxOwnershipStore(abc.ABC):
def claim(self, sandbox_id: str, *, for_destroy: bool = False) -> bool:
"""Take ownership of *sandbox_id* only if it is unowned or already ours.
Atomic: concurrent claims from different instances cannot both succeed.
Gates every adopt/reap path.
Exclusive: succeeds only when the container is unowned or already ours,
which is what gates every adopt/reap path.
The read-modify-write must not interleave. On redis that is Lua (one
script, server-side); the memory store serializes on a process-local lock
and is single-instance anyway, so "different instances" cannot arise
there. Note what is *not* verified: the contract suite drives sequential
calls, so it pins the exclusion predicate, not the atomicity and CI
runs the memory tier only, so the Lua that carries it never executes on
the merge gate.
Args:
for_destroy: mark the lease as a teardown in progress, so a
@@ -755,6 +755,12 @@ def test_reconcile_fails_closed_when_ownership_unknown():
worker = _make_provider_for_reconciliation(worker_id="worker-b")
worker._ownership = MagicMock()
# Configure what the adoption grace reads, or it short-circuits before the
# claim: a bare MagicMock answers `owner()` with a truthy mock, which reads
# as "peer-owned" and defers — so the assertion below would pass without the
# fail-closed branch ever running (it did exactly that until this was fixed).
worker._ownership.supports_cross_process = True
worker._ownership.owner.return_value = None
worker._ownership.claim.side_effect = OwnershipBackendError("store down")
info = SandboxInfo(
sandbox_id="unknown2",
@@ -764,8 +770,16 @@ def test_reconcile_fails_closed_when_ownership_unknown():
)
worker._backend.list_running.return_value = [info]
worker._reconcile_orphans()
# Unowned for a full grace, so the container is adoptable and the only thing
# left standing between it and the warm pool is the claim.
aio_mod = importlib.import_module("deerflow.community.aio_sandbox.aio_sandbox_provider")
now = time.time()
with patch.object(aio_mod.time, "time", return_value=now):
worker._reconcile_orphans()
with patch.object(aio_mod.time, "time", return_value=now + compute_lease_ttl(worker._ownership_config) + 1):
worker._reconcile_orphans()
assert worker._ownership.claim.called, "the fail-closed claim branch was never reached; this test guards nothing"
assert "unknown2" not in worker._warm_pool
worker._backend.destroy.assert_not_called()
@@ -801,6 +815,112 @@ def test_init_always_starts_lease_renewal(monkeypatch, idle_timeout):
assert ("idle" in started) is (idle_timeout > 0)
def test_renewal_keeps_the_sandbox_when_the_store_cannot_answer():
"""The one deliberate exception to fail-closed, and it had no test.
Everywhere else an unanswerable store means "not ours". Renewal is the
opposite on purpose: `_refresh_ownership` returns True on an
`OwnershipBackendError` because unknown is not lost, and the TTL still bounds
how long a genuinely dead owner holds the lease. Invert it and a Redis outage
makes every instance drop every active and warm sandbox at once the same
fleet-wide eviction the LAPSED/LOST split exists to prevent, which is pinned
only for the flushed-store path, never for a raising one.
"""
from deerflow.community.aio_sandbox.ownership import OwnershipBackendError
worker = _make_provider_for_reconciliation(worker_id="worker-a")
worker._ownership = MagicMock()
worker._ownership.renew.side_effect = OwnershipBackendError("store down")
info = SandboxInfo(
sandbox_id="live02",
sandbox_url="http://localhost:8080",
container_name="deer-flow-sandbox-live02",
created_at=time.time(),
)
sandbox = MagicMock()
worker._sandboxes["live02"] = sandbox
worker._sandbox_infos["live02"] = info
worker._thread_sandboxes[("u1", "t1")] = "live02"
worker._warm_pool["warm02"] = (info, time.time())
worker._renew_owned_leases()
assert worker._ownership.renew.called, "renewal never reached the store; this test guards nothing"
assert "live02" in worker._sandboxes, "a store outage evicted a live sandbox nobody had taken"
assert ("u1", "t1") in worker._thread_sandboxes
assert "warm02" in worker._warm_pool, "a store outage dropped a warm entry nobody had taken"
sandbox.close.assert_not_called()
worker._backend.destroy.assert_not_called()
def test_load_config_carries_the_stream_bridge_section():
"""Hop 1 of the "no extra config for multi-instance" promise.
The redis inference reads `app_config.stream_bridge`, so `_load_config` has to
carry it. Nothing pinned this: the only test that drives `__init__`
monkeypatches `_load_config` wholesale and omits the key entirely, so deleting
it here left every test green while every config.yaml-native multi-instance
deployment silently fell back to `memory` #4206 reopened on exactly the
deployments the inference exists for.
"""
from deerflow.config.stream_bridge_config import StreamBridgeConfig
aio_mod = importlib.import_module("deerflow.community.aio_sandbox.aio_sandbox_provider")
provider = aio_mod.AioSandboxProvider.__new__(aio_mod.AioSandboxProvider)
bridge = StreamBridgeConfig(type="redis", redis_url="redis://bridge:6379/0")
app_config = MagicMock()
app_config.stream_bridge = bridge
app_config.sandbox = MagicMock(ownership=None, image=None, port=None, container_prefix=None, idle_timeout=600, replicas=3, mounts=[], environment={})
with patch.object(aio_mod, "get_app_config", return_value=app_config):
loaded = provider._load_config()
assert loaded["stream_bridge"] is bridge, "_load_config dropped the stream_bridge section the redis inference reads"
def test_init_infers_redis_ownership_from_a_redis_stream_bridge():
"""Hop 2: `__init__` must actually feed the bridge into the resolver.
Drives the real `__init__` against a real `AppConfig`-shaped object rather
than stubbing `_load_config`, because the defect would be in the wiring
between them the same reason `test_init_always_starts_lease_renewal` drives
`__init__` instead of calling `_start_lease_renewal` directly.
"""
from deerflow.config.stream_bridge_config import StreamBridgeConfig
aio_mod = importlib.import_module("deerflow.community.aio_sandbox.aio_sandbox_provider")
app_config = MagicMock()
app_config.stream_bridge = StreamBridgeConfig(type="redis", redis_url="redis://bridge:6379/0")
# No sandbox.ownership section at all: the deployment never configured one.
app_config.sandbox = MagicMock(ownership=None, image=None, port=None, container_prefix=None, idle_timeout=600, replicas=3, mounts=[], environment={})
built: list = []
def fake_store(config, *, owner_id=None):
built.append(config)
store = MagicMock()
store.supports_cross_process = True
return store
with (
patch.object(aio_mod, "get_app_config", return_value=app_config),
patch.object(aio_mod, "make_sandbox_ownership_store", side_effect=fake_store),
patch.object(aio_mod.AioSandboxProvider, "_create_backend", lambda self: MagicMock()),
patch.object(aio_mod.AioSandboxProvider, "_reconcile_orphans", lambda self: None),
patch.object(aio_mod.AioSandboxProvider, "_register_signal_handlers", lambda self: None),
patch.object(aio_mod.AioSandboxProvider, "_start_lease_renewal", lambda self: None),
patch.object(aio_mod.AioSandboxProvider, "_start_idle_checker", lambda self: None),
patch.object(aio_mod.atexit, "register", lambda *a, **k: None),
):
aio_mod.AioSandboxProvider()
assert len(built) == 1
assert built[0].type == "redis", "a redis stream bridge did not infer a redis ownership store; multi-instance deployments silently fall back to memory"
assert built[0].redis_url == "redis://bridge:6379/0"
def test_renewal_loop_refreshes_owned_leases():
"""The renewal thread actually renews (the loop body, not just its wiring)."""
from deerflow.config.sandbox_config import SandboxOwnershipConfig
@@ -1061,9 +1181,12 @@ def test_adoption_grace_expires_so_a_truly_orphaned_container_is_still_adopted()
def test_adoption_grace_restarts_when_a_live_owner_republishes():
"""A republished lease must reset the grace, not just pause it.
Otherwise an owner that republishes late still loses its container: the
adopter's grace, started before the republish, would expire regardless and
hand it a container that now has a live owner.
Reset and pause only diverge on a **second** lapse. Pause leaves the original
timestamp behind, so the next time the lease drops the grace is already spent
and the adopter takes a live container with no wait at all. Stopping at "A
republished, B defers" would prove nothing — a paused timer defers there too,
because the container simply reads as owned. So the second lapse is the whole
test; without it this passes with the reset deleted.
"""
aio_mod = importlib.import_module("deerflow.community.aio_sandbox.aio_sandbox_provider")
shared = _make_shared_ownership_store()
@@ -1092,6 +1215,18 @@ def test_adoption_grace_restarts_when_a_live_owner_republishes():
assert "live01" not in worker_b._warm_pool, "a stale grace expired over a lease a live owner had already republished"
assert shared.owner("live01") == "worker-a"
# The republish must have cleared B's timer, not merely paused it. A second
# blip drops the key again: B has to serve a *fresh* full grace, which A's
# next renewal tick will beat. A paused timer would still hold the original
# start, so B would adopt A's live container instantly, with no grace at all.
assert "live01" not in worker_b._unowned_since, "the republish left a stale grace timer behind"
shared._leases.clear()
with patch.object(aio_mod.time, "time", return_value=now + ttl + 2):
worker_b._reconcile_orphans()
assert "live01" not in worker_b._warm_pool, "a grace timer left over from before the republish expired instantly on the next lapse"
def test_acquire_refuses_a_container_a_peer_is_destroying():
"""#4206's last window: `take()` must not overrun a destroyer's claim.
@@ -1182,6 +1317,226 @@ def test_teardown_marker_is_held_for_a_stop_that_outlives_the_lease_ttl():
assert worker_b._ownership.take("doomed1") is True
def test_unhealthy_drop_holds_the_teardown_marker_for_its_stop():
"""The third `del:`-marked stop path needs the same hold as the other two.
`_drop_unhealthy_sandbox` claims for destroy and then blocks on the backend
stop exactly like `_destroy_warm_entry` and `destroy()`. Its sibling test
`test_unhealthy_sandbox_owned_by_peer_is_not_destroyed` pins the *gate* a
peer-owned container is not stopped but never lets the marker **expire**
during an in-flight stop, which is the same blind spot that hid this window
on the other two paths. It untracks before claiming, so `_renew_owned_leases`
cannot see the id either: nothing refreshes the marker unless the stop holds
it.
"""
from deerflow.config.sandbox_config import SandboxOwnershipConfig
lease_ttl = 0.15
shared = _make_shared_ownership_store(ttl_seconds=lease_ttl)
worker_a = _make_provider_for_reconciliation(worker_id="worker-a", store=shared)
worker_b = _make_provider_for_reconciliation(worker_id="worker-b", store=shared)
worker_a._ownership_config = SandboxOwnershipConfig(renewal_interval_seconds=0.05, ttl_multiplier=3.0)
info = SandboxInfo(
sandbox_id="sick01",
sandbox_url="http://localhost:8080",
container_name="deer-flow-sandbox-sick01",
created_at=time.time(),
)
stop_entered = threading.Event()
release_stop = threading.Event()
def slow_destroy(entry):
stop_entered.set()
release_stop.wait(timeout=5)
worker_a._backend.destroy = MagicMock(side_effect=slow_destroy)
worker_a._sandboxes["sick01"] = MagicMock()
worker_a._sandbox_infos["sick01"] = info
dropper = threading.Thread(
target=lambda: worker_a._drop_unhealthy_sandbox("sick01", "health check failed"),
daemon=True,
)
dropper.start()
try:
assert stop_entered.wait(timeout=5), "the drop never reached the backend stop"
deadline = time.time() + lease_ttl * 4
while time.time() < deadline:
assert not worker_b._ownership.take("sick01"), "a peer took a container whose unhealthy-drop stop was still in flight"
time.sleep(0.02)
finally:
release_stop.set()
dropper.join(timeout=5)
assert shared.owner("sick01") is None, "the teardown marker outlived the stop that justified it"
def test_destroy_holds_the_teardown_marker_for_its_stop():
"""The third of the three `del:`-marked stops, and the one with no test.
`_destroy_warm_entry` and `_drop_unhealthy_sandbox` each have a held-marker
test; `destroy()` is wrapped but nothing pins the wrap, so deleting it goes
unnoticed. "Every path does X" claims keep leaving exactly one sibling
untested this is that sibling.
"""
from deerflow.config.sandbox_config import SandboxOwnershipConfig
lease_ttl = 0.15
shared = _make_shared_ownership_store(ttl_seconds=lease_ttl)
worker_a = _make_provider_for_reconciliation(worker_id="worker-a", store=shared)
worker_b = _make_provider_for_reconciliation(worker_id="worker-b", store=shared)
worker_a._ownership_config = SandboxOwnershipConfig(renewal_interval_seconds=0.05, ttl_multiplier=3.0)
info = SandboxInfo(
sandbox_id="doomed3",
sandbox_url="http://localhost:8080",
container_name="deer-flow-sandbox-doomed3",
created_at=time.time(),
)
stop_entered = threading.Event()
release_stop = threading.Event()
def slow_destroy(entry):
stop_entered.set()
release_stop.wait(timeout=5)
worker_a._backend.destroy = MagicMock(side_effect=slow_destroy)
worker_a._sandboxes["doomed3"] = MagicMock()
worker_a._sandbox_infos["doomed3"] = info
destroyer = threading.Thread(target=lambda: worker_a.destroy("doomed3"), daemon=True)
destroyer.start()
try:
assert stop_entered.wait(timeout=5), "destroy never reached the backend stop"
deadline = time.time() + lease_ttl * 4
while time.time() < deadline:
assert not worker_b._ownership.take("doomed3"), "a peer took a container whose destroy() stop was still in flight"
time.sleep(0.02)
finally:
release_stop.set()
destroyer.join(timeout=5)
assert shared.owner("doomed3") is None, "the teardown marker outlived the stop that justified it"
def test_evict_keeps_the_warm_entry_when_the_claim_is_refused():
"""Replica eviction must not pop before it knows the container is going away.
The sibling of `test_refused_idle_destroy_keeps_the_warm_entry`, which pins
the same rule for the idle path. Popping first on a refused claim loses the
container: still running, owned by a peer, and no longer in any of our maps
so nothing here would ever reap or reclaim it.
"""
shared = _make_shared_ownership_store()
worker_a = _make_provider_for_reconciliation(worker_id="worker-a", store=shared)
worker_b = _make_provider_for_reconciliation(worker_id="worker-b", store=shared)
info = SandboxInfo(
sandbox_id="peer01",
sandbox_url="http://localhost:8080",
container_name="deer-flow-sandbox-peer01",
created_at=time.time(),
)
worker_a._warm_pool["peer01"] = (info, time.time() - 5)
# A peer owns it, so A's eviction claim is refused.
worker_b._publish_ownership("peer01")
assert worker_a._evict_oldest_warm() is None
worker_a._backend.destroy.assert_not_called()
assert "peer01" in worker_a._warm_pool, "evicting popped a container it was refused permission to destroy"
assert shared.owner("peer01") == "worker-b"
def test_reclaim_drops_a_container_a_peer_is_destroying():
"""The warm-pool half of the acquire-side teardown refusal.
`test_cached_sandbox_being_destroyed_is_dropped_not_reused` pins the
in-process reuse path and `test_acquire_refuses_a_container_a_peer_is_destroying`
the discover path; reclaim is the third and had no test. It must not raise
(the caller falls through to a cold start) and must not leave the doomed
container in the warm pool.
"""
shared = _make_shared_ownership_store()
worker_a = _make_provider_for_reconciliation(worker_id="worker-a", store=shared)
worker_b = _make_provider_for_reconciliation(worker_id="worker-b", store=shared)
info = SandboxInfo(
sandbox_id="dying02",
sandbox_url="http://localhost:8080",
container_name="deer-flow-sandbox-dying02",
created_at=time.time(),
)
worker_a._warm_pool["dying02"] = (info, time.time())
worker_a._check_tracked_sandbox_alive = MagicMock(return_value=True)
# B's reaper marks the teardown; its stop is in flight.
assert worker_b._claim_ownership("dying02", for_destroy=True) is True
reclaimed = worker_a._reclaim_warm_pool_sandbox("t1", "dying02", user_id="u1")
assert reclaimed is None, "reclaimed a container a peer is tearing down"
assert "dying02" not in worker_a._warm_pool
assert "dying02" not in worker_a._sandboxes
worker_a._backend.destroy.assert_not_called()
def test_created_sandbox_is_rolled_back_when_a_peer_is_destroying_its_id():
"""Rollback must cover a teardown marker, not just a store outage.
`test_ownership_rollback_on_create_closes_the_client_it_drops` drives this
path with `OwnershipBackendError` only. The comment says the teardown case is
reachable too a peer that died mid-stop leaves a `del:` marker until its
TTL lapses and without rollback the container we just started is leaked as
an adoptable orphan.
"""
shared = _make_shared_ownership_store()
worker_a = _make_provider_for_reconciliation(worker_id="worker-a", store=shared)
worker_b = _make_provider_for_reconciliation(worker_id="worker-b", store=shared)
info = SandboxInfo(
sandbox_id="fresh01",
sandbox_url="http://localhost:8080",
container_name="deer-flow-sandbox-fresh01",
created_at=time.time(),
)
# A peer's teardown marker is still on this id when we finish creating.
assert worker_b._claim_ownership("fresh01", for_destroy=True) is True
with pytest.raises(SandboxBeingDestroyedError):
worker_a._register_created_sandbox("t1", "fresh01", info, user_id="u1")
worker_a._backend.destroy.assert_called_once_with(info)
assert "fresh01" not in worker_a._sandboxes, "a container we could not own was handed out anyway"
def test_shutdown_does_not_stop_a_peers_warm_container():
"""Shutdown is a reap path and must be gated like every other one.
Nothing drove `shutdown()` with a non-empty warm pool, so a loop that called
`_backend.destroy` directly skipping the ownership claim would go
unnoticed. On a multi-instance gateway that is #4206 on the shutdown path:
our exit stops a container a live peer is serving.
"""
shared = _make_shared_ownership_store()
worker_a = _make_provider_for_reconciliation(worker_id="worker-a", store=shared)
worker_b = _make_provider_for_reconciliation(worker_id="worker-b", store=shared)
mine = SandboxInfo(sandbox_id="mine01", sandbox_url="http://localhost:8080", container_name="c-mine01", created_at=time.time())
theirs = SandboxInfo(sandbox_id="peer02", sandbox_url="http://localhost:8081", container_name="c-peer02", created_at=time.time())
worker_a._warm_pool["mine01"] = (mine, time.time())
worker_a._publish_ownership("mine01")
worker_a._warm_pool["peer02"] = (theirs, time.time())
worker_b._publish_ownership("peer02") # a live peer owns this one
worker_a.shutdown()
destroyed = {call.args[0].sandbox_id for call in worker_a._backend.destroy.call_args_list}
assert "mine01" in destroyed, "shutdown left our own warm container running"
assert "peer02" not in destroyed, "shutdown stopped a container a live peer owns"
assert shared.owner("peer02") == "worker-b"
def test_teardown_heartbeat_stops_when_the_stop_returns():
"""A finite TTL must survive the fix, or a crashed destroyer leaks forever.