mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-09-18 10:36:17 +00:00
* fix(sandbox): claim ownership before readiness-timeout destroy (#4248) When a freshly-created sandbox fails to become ready within the 60s timeout, _create_sandbox / _create_sandbox_async destroyed the container with a bare self._backend.destroy(info) call. Ownership is published by _register_created_sandbox only after the readiness gate, so for the whole timeout window the container ran unowned — exactly the state a peer gateway startup reconciliation is built to adopt across. With the claim skipped, a peer could adopt the not-yet-ready Pod and this instance subsequent stop landed on the turn the peer had just handed it: a cross-instance kill of an active turn. Route the destroy through a new _destroy_unready_sandbox helper that first claims the teardown lease (claim(..., for_destroy=True) writes the del: marker) and holds it via _held_teardown_lease for the duration of the stop, matching the ownership guard every other reap path already uses (_destroy_warm_entry, _drop_unhealthy_reserved). Fail closed if a peer already holds the lease or the ownership store cannot answer: leave the container for the peer own reconciliation rather than stopping it from underneath an active turn. Fixes #4248. * fix(sandbox): reserve local teardown around readiness-timeout destroy Review feedback (AnnaSuSu) on #4505: the ownership claim is only the cross-instance half of the guard — claim() succeeds against our own lease by design, so in the window between the readiness timeout and the claim a same-process _reconcile_orphans (idle checker, every 60s) can adopt the unready container into _warm_pool; the subsequent claim still succeeds and the stop lands on an entry this instance has just adopted, leaving a dead warm entry for the next reclaim to hand out. Wrap the still-untracked check, claim, and stop in _reserve_local_teardown / _finish_local_teardown, with the predicate checking the id is absent from the active and warm maps — the same pairing _destroy_warm_entry already uses. Add an interleaving regression test mirroring test_reconcile_does_not_adopt_a_container_this_instance_is_tearing_down: reconcile runs while the destroy thread is parked after reserving but before its `del:` claim lands, and must not adopt. A mirror test asserts a genuinely unowned container is still adopted when no teardown is in flight, so the reservation cannot over-block reconciliation. * style(sandbox): apply ruff format to readiness-timeout destroy changes --------- Co-authored-by: now-ing <24534365+now-ing@users.noreply.github.com>
This commit is contained in:
parent
f37c734406
commit
2a96341889
@ -1924,6 +1924,61 @@ class AioSandboxProvider(WarmPoolLifecycleMixin[SandboxInfo], SandboxProvider):
|
|||||||
await asyncio.to_thread(_unlock_file, lock_file)
|
await asyncio.to_thread(_unlock_file, lock_file)
|
||||||
await asyncio.to_thread(lock_file.close)
|
await asyncio.to_thread(lock_file.close)
|
||||||
|
|
||||||
|
def _destroy_unready_sandbox(self, sandbox_id: str, info: SandboxInfo) -> None:
|
||||||
|
"""Tear down a freshly-created container whose readiness check failed.
|
||||||
|
|
||||||
|
The container was started by the backend but never reached ready, so it
|
||||||
|
never entered ``_register_created_sandbox`` and the ownership store has
|
||||||
|
no lease for it yet. For the full readiness timeout (60s) it runs
|
||||||
|
unowned, which is exactly the window a peer gateway's startup
|
||||||
|
reconciliation is built to adopt across (#4206). Without a claim, a peer
|
||||||
|
that adopts the not-yet-ready Pod and then has this instance's stop land
|
||||||
|
on it is a cross-instance kill that interrupts an active turn (#4248).
|
||||||
|
|
||||||
|
Claim the teardown lease first so this reap path is gated by the same
|
||||||
|
ownership guard as every other destroy (``_destroy_warm_entry``,
|
||||||
|
``_drop_unhealthy_reserved``); fail closed (leave the container for the
|
||||||
|
peer to reap via its own reconciliation) if a peer already owns it or
|
||||||
|
the ownership store cannot answer.
|
||||||
|
|
||||||
|
The claim alone is only the cross-**instance** half: it succeeds against
|
||||||
|
our own lease by design, so it says nothing about this process. The
|
||||||
|
same-process half is the local teardown reservation, taken first and
|
||||||
|
held across the whole path — between the readiness timeout and the
|
||||||
|
claim, ``_reconcile_orphans`` (idle checker, every 60s) can see this
|
||||||
|
container running, untracked, and past its recovery grace, and park it
|
||||||
|
in ``_warm_pool``; the subsequent claim would still succeed and the
|
||||||
|
stop would land on an entry this instance has just adopted, leaving a
|
||||||
|
dead warm entry for the next reclaim to hand out. The predicate checks
|
||||||
|
the id is absent from both the active and warm maps; the reservation
|
||||||
|
makes that check and the teardown mark one critical section, so no
|
||||||
|
adopt/acquire can slip between them (same pairing as
|
||||||
|
``_destroy_warm_entry``).
|
||||||
|
"""
|
||||||
|
if not self._reserve_local_teardown(
|
||||||
|
sandbox_id,
|
||||||
|
lambda: sandbox_id not in self._sandboxes and sandbox_id not in self._sandbox_infos and sandbox_id not in self._warm_pool,
|
||||||
|
):
|
||||||
|
logger.warning(
|
||||||
|
"Not destroying unready sandbox %s: adopted or being torn down by this instance",
|
||||||
|
sandbox_id,
|
||||||
|
)
|
||||||
|
return
|
||||||
|
try:
|
||||||
|
if not self._claim_ownership(sandbox_id, for_destroy=True):
|
||||||
|
logger.warning(
|
||||||
|
"Not destroying unready sandbox %s: owned by another instance or ownership unavailable",
|
||||||
|
sandbox_id,
|
||||||
|
)
|
||||||
|
return
|
||||||
|
try:
|
||||||
|
with self._held_teardown_lease(sandbox_id):
|
||||||
|
self._backend.destroy(info)
|
||||||
|
except Exception as e:
|
||||||
|
logger.warning(f"Error destroying unready sandbox {sandbox_id}: {e}")
|
||||||
|
finally:
|
||||||
|
self._finish_local_teardown(sandbox_id)
|
||||||
|
|
||||||
def _create_sandbox(self, thread_id: str | None, sandbox_id: str, *, user_id: str | None = None) -> str:
|
def _create_sandbox(self, thread_id: str | None, sandbox_id: str, *, user_id: str | None = None) -> str:
|
||||||
"""Create a new sandbox via the backend.
|
"""Create a new sandbox via the backend.
|
||||||
|
|
||||||
@ -1960,7 +2015,11 @@ class AioSandboxProvider(WarmPoolLifecycleMixin[SandboxInfo], SandboxProvider):
|
|||||||
|
|
||||||
# Wait for sandbox to be ready
|
# Wait for sandbox to be ready
|
||||||
if not wait_for_sandbox_ready(info.sandbox_url, timeout=60):
|
if not wait_for_sandbox_ready(info.sandbox_url, timeout=60):
|
||||||
self._backend.destroy(info)
|
# The container is running but unowned: ownership is published by
|
||||||
|
# ``_register_created_sandbox`` after this gate. Claim the teardown
|
||||||
|
# lease before stopping it so a peer cannot adopt the not-yet-ready
|
||||||
|
# Pod in the meantime (#4248).
|
||||||
|
self._destroy_unready_sandbox(sandbox_id, info)
|
||||||
raise RuntimeError(f"Sandbox {sandbox_id} failed to become ready within timeout at {info.sandbox_url}")
|
raise RuntimeError(f"Sandbox {sandbox_id} failed to become ready within timeout at {info.sandbox_url}")
|
||||||
|
|
||||||
return self._register_created_sandbox(thread_id, sandbox_id, info, user_id=effective_user_id)
|
return self._register_created_sandbox(thread_id, sandbox_id, info, user_id=effective_user_id)
|
||||||
@ -1991,7 +2050,11 @@ class AioSandboxProvider(WarmPoolLifecycleMixin[SandboxInfo], SandboxProvider):
|
|||||||
|
|
||||||
# Wait for sandbox to be ready without blocking the event loop.
|
# Wait for sandbox to be ready without blocking the event loop.
|
||||||
if not await wait_for_sandbox_ready_async(info.sandbox_url, timeout=60):
|
if not await wait_for_sandbox_ready_async(info.sandbox_url, timeout=60):
|
||||||
await asyncio.to_thread(self._backend.destroy, info)
|
# The container is running but unowned: ownership is published by
|
||||||
|
# ``_register_created_sandbox`` after this gate. Claim the teardown
|
||||||
|
# lease before stopping it so a peer cannot adopt the not-yet-ready
|
||||||
|
# Pod in the meantime (#4248).
|
||||||
|
await asyncio.to_thread(self._destroy_unready_sandbox, sandbox_id, info)
|
||||||
raise RuntimeError(f"Sandbox {sandbox_id} failed to become ready within timeout at {info.sandbox_url}")
|
raise RuntimeError(f"Sandbox {sandbox_id} failed to become ready within timeout at {info.sandbox_url}")
|
||||||
|
|
||||||
# Registration publishes ownership (blocking store IO), so it is offloaded
|
# Registration publishes ownership (blocking store IO), so it is offloaded
|
||||||
|
|||||||
@ -5,6 +5,7 @@ import contextlib
|
|||||||
import hashlib
|
import hashlib
|
||||||
import importlib
|
import importlib
|
||||||
import stat
|
import stat
|
||||||
|
import threading
|
||||||
from types import SimpleNamespace
|
from types import SimpleNamespace
|
||||||
from unittest.mock import MagicMock, patch
|
from unittest.mock import MagicMock, patch
|
||||||
|
|
||||||
@ -1071,3 +1072,222 @@ def test_aio_forced_collision_never_overwrites_active_tenant(
|
|||||||
assert len(create_calls) == 1
|
assert len(create_calls) == 1
|
||||||
assert provider.acquire("thread-a", user_id="user-a") == sandbox_id
|
assert provider.acquire("thread-a", user_id="user-a") == sandbox_id
|
||||||
assert provider._sandbox_infos[sandbox_id] is info_a
|
assert provider._sandbox_infos[sandbox_id] is info_a
|
||||||
|
|
||||||
|
|
||||||
|
# --- #4248 regression: readiness-timeout destroy ownership ---
|
||||||
|
|
||||||
|
|
||||||
|
def _make_unready_destroy_provider(tmp_path, *, sandbox_id, base_url, monkeypatch, aio_mod):
|
||||||
|
"""Provider wired so ``_create_sandbox`` reaches the readiness-timeout branch.
|
||||||
|
|
||||||
|
``wait_for_sandbox_ready`` always returns False; the backend records what the
|
||||||
|
destroy path did. Mirrors the fixtures used by the warm-replica eviction
|
||||||
|
test, minus the warm pool.
|
||||||
|
"""
|
||||||
|
provider = _make_provider(tmp_path)
|
||||||
|
provider._lock = aio_mod.threading.Lock()
|
||||||
|
provider._config = {"replicas": 3}
|
||||||
|
provider._thread_locks = {}
|
||||||
|
provider._warm_pool = {}
|
||||||
|
provider._sandbox_infos = {}
|
||||||
|
provider._thread_sandboxes = {}
|
||||||
|
provider._last_activity = {}
|
||||||
|
provider._active_sandbox_identity = {}
|
||||||
|
provider._warm_pool_identity = {}
|
||||||
|
unready_info = aio_mod.SandboxInfo(sandbox_id=sandbox_id, sandbox_url=base_url)
|
||||||
|
provider._backend = SimpleNamespace(
|
||||||
|
create=MagicMock(return_value=unready_info),
|
||||||
|
destroy=MagicMock(),
|
||||||
|
)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
aio_mod.AioSandboxProvider,
|
||||||
|
"_get_extra_mounts",
|
||||||
|
lambda _self, _thread_id, *, user_id=None: [],
|
||||||
|
)
|
||||||
|
return provider, unready_info
|
||||||
|
|
||||||
|
|
||||||
|
def test_create_sandbox_claims_ownership_before_readiness_timeout_destroy(tmp_path, monkeypatch):
|
||||||
|
"""#4248: a readiness-timeout destroy must run under a `del:` teardown lease.
|
||||||
|
|
||||||
|
Before #4248 the unready container was reaped with a bare ``destroy`` call.
|
||||||
|
Ownership is published by ``_register_created_sandbox`` only after the
|
||||||
|
readiness gate, so for up to 60s the container ran unowned and a peer could
|
||||||
|
adopt it; the subsequent stop landed on whatever turn the peer had handed it.
|
||||||
|
"""
|
||||||
|
aio_mod = importlib.import_module("deerflow.community.aio_sandbox.aio_sandbox_provider")
|
||||||
|
provider, unready_info = _make_unready_destroy_provider(
|
||||||
|
tmp_path,
|
||||||
|
sandbox_id="unready",
|
||||||
|
base_url="http://unready",
|
||||||
|
monkeypatch=monkeypatch,
|
||||||
|
aio_mod=aio_mod,
|
||||||
|
)
|
||||||
|
monkeypatch.setattr(aio_mod, "wait_for_sandbox_ready", lambda _url, *, timeout=60: False)
|
||||||
|
|
||||||
|
# The heartbeat releases the teardown lease on exit, so the destroy call is
|
||||||
|
# the only place we can observe the `del:` state. Snapshot the lease at
|
||||||
|
# the instant destroy runs.
|
||||||
|
destroy_snapshots: list = []
|
||||||
|
|
||||||
|
def destroy_spy(info):
|
||||||
|
destroy_snapshots.append(provider._ownership._leases.get(info.sandbox_id))
|
||||||
|
|
||||||
|
provider._backend.destroy.side_effect = destroy_spy
|
||||||
|
|
||||||
|
with pytest.raises(RuntimeError, match="failed to become ready"):
|
||||||
|
provider._create_sandbox("thread-4248", "unready", user_id="user-4248")
|
||||||
|
|
||||||
|
provider._backend.destroy.assert_called_once_with(unready_info)
|
||||||
|
assert destroy_snapshots, "destroy must run inside the held teardown lease"
|
||||||
|
lease = destroy_snapshots[0]
|
||||||
|
assert lease is not None, "teardown lease must be held while destroy runs"
|
||||||
|
assert lease.owner_id == provider._owner_id
|
||||||
|
assert lease.destroying is True, "destroy must run under a `del:` teardown lease"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.anyio
|
||||||
|
async def test_create_sandbox_async_claims_ownership_before_readiness_timeout_destroy(tmp_path, monkeypatch):
|
||||||
|
"""#4248 (async path): same teardown-lease guard on the async readiness branch."""
|
||||||
|
aio_mod = importlib.import_module("deerflow.community.aio_sandbox.aio_sandbox_provider")
|
||||||
|
provider, unready_info = _make_unready_destroy_provider(
|
||||||
|
tmp_path,
|
||||||
|
sandbox_id="unready-async",
|
||||||
|
base_url="http://unready-async",
|
||||||
|
monkeypatch=monkeypatch,
|
||||||
|
aio_mod=aio_mod,
|
||||||
|
)
|
||||||
|
|
||||||
|
async def fake_wait_async(_url, *, timeout=60, poll_interval=1.0):
|
||||||
|
return False
|
||||||
|
|
||||||
|
monkeypatch.setattr(aio_mod, "wait_for_sandbox_ready_async", fake_wait_async)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
aio_mod,
|
||||||
|
"wait_for_sandbox_ready",
|
||||||
|
lambda *_a, **_k: (_ for _ in ()).throw(AssertionError("sync readiness should not be used")),
|
||||||
|
)
|
||||||
|
|
||||||
|
destroy_snapshots: list = []
|
||||||
|
|
||||||
|
def destroy_spy(info):
|
||||||
|
destroy_snapshots.append(provider._ownership._leases.get(info.sandbox_id))
|
||||||
|
|
||||||
|
provider._backend.destroy.side_effect = destroy_spy
|
||||||
|
|
||||||
|
with pytest.raises(RuntimeError, match="failed to become ready"):
|
||||||
|
await provider._create_sandbox_async("thread-4248-async", "unready-async", user_id="user-4248-async")
|
||||||
|
|
||||||
|
provider._backend.destroy.assert_called_once_with(unready_info)
|
||||||
|
assert destroy_snapshots, "destroy must run inside the held teardown lease"
|
||||||
|
lease = destroy_snapshots[0]
|
||||||
|
assert lease is not None
|
||||||
|
assert lease.owner_id == provider._owner_id
|
||||||
|
assert lease.destroying is True, "destroy must run under a `del:` teardown lease"
|
||||||
|
|
||||||
|
|
||||||
|
def test_create_sandbox_skips_destroy_when_unready_sandbox_owned_by_peer(tmp_path, monkeypatch):
|
||||||
|
"""#4248 fail-closed: if a peer already owns the unready container, do not stop it.
|
||||||
|
|
||||||
|
The lease refuses our teardown claim, so the container is left for the peer
|
||||||
|
to reap via its own reconciliation. Stopping it anyway would be the
|
||||||
|
cross-instance kill this guard exists to prevent.
|
||||||
|
"""
|
||||||
|
aio_mod = importlib.import_module("deerflow.community.aio_sandbox.aio_sandbox_provider")
|
||||||
|
provider, unready_info = _make_unready_destroy_provider(
|
||||||
|
tmp_path,
|
||||||
|
sandbox_id="peer-owned",
|
||||||
|
base_url="http://peer-owned",
|
||||||
|
monkeypatch=monkeypatch,
|
||||||
|
aio_mod=aio_mod,
|
||||||
|
)
|
||||||
|
monkeypatch.setattr(aio_mod, "wait_for_sandbox_ready", lambda _url, *, timeout=60: False)
|
||||||
|
|
||||||
|
# Every claim refuses: peer holds the lease (or the store cannot answer).
|
||||||
|
provider._ownership.claim = lambda _sid, *, for_destroy=False: False
|
||||||
|
|
||||||
|
with pytest.raises(RuntimeError, match="failed to become ready"):
|
||||||
|
provider._create_sandbox("thread-peer", "peer-owned", user_id="user-peer")
|
||||||
|
|
||||||
|
provider._backend.destroy.assert_not_called()
|
||||||
|
|
||||||
|
|
||||||
|
def test_reconcile_does_not_adopt_a_container_whose_unready_teardown_is_reserved(tmp_path, monkeypatch):
|
||||||
|
"""#4248 follow-up: the readiness-timeout destroy must hold the local
|
||||||
|
reservation, not just the cross-instance claim.
|
||||||
|
|
||||||
|
The claim succeeds against our own lease by design, so without
|
||||||
|
``_reserve_local_teardown`` there is a window — readiness failed, claim not
|
||||||
|
yet written — in which the idle checker's ``_reconcile_orphans`` sees the
|
||||||
|
container running, untracked, and past its recovery grace, and adopts it
|
||||||
|
into ``_warm_pool``. The claim then still succeeds (the lease is ours) and
|
||||||
|
the stop lands on an entry this instance has just adopted, leaving a dead
|
||||||
|
warm entry for the next reclaim to hand out. This is the same interleaving
|
||||||
|
shape as ``test_reconcile_does_not_adopt_a_container_this_instance_is_tearing_down``
|
||||||
|
in ``test_sandbox_orphan_reconciliation.py``.
|
||||||
|
"""
|
||||||
|
aio_mod = importlib.import_module("deerflow.community.aio_sandbox.aio_sandbox_provider")
|
||||||
|
provider, unready_info = _make_unready_destroy_provider(
|
||||||
|
tmp_path,
|
||||||
|
sandbox_id="unready-race",
|
||||||
|
base_url="http://unready-race",
|
||||||
|
monkeypatch=monkeypatch,
|
||||||
|
aio_mod=aio_mod,
|
||||||
|
)
|
||||||
|
provider._unowned_since = {}
|
||||||
|
provider._backend.list_running = MagicMock(return_value=[unready_info])
|
||||||
|
monkeypatch.setattr(aio_mod, "wait_for_sandbox_ready", lambda _url, *, timeout=60: False)
|
||||||
|
|
||||||
|
# Park the destroy thread after it has reserved the local teardown but
|
||||||
|
# before the `del:` claim lands — the exact window reconcile would adopt in.
|
||||||
|
at_claim, let_claim = threading.Event(), threading.Event()
|
||||||
|
real_claim = provider._claim_ownership
|
||||||
|
|
||||||
|
def gated_claim(sandbox_id, *, for_destroy=False):
|
||||||
|
if for_destroy:
|
||||||
|
at_claim.set()
|
||||||
|
assert let_claim.wait(timeout=5)
|
||||||
|
return real_claim(sandbox_id, for_destroy=for_destroy)
|
||||||
|
|
||||||
|
provider._claim_ownership = gated_claim
|
||||||
|
reaper = threading.Thread(
|
||||||
|
target=lambda: provider._destroy_unready_sandbox("unready-race", unready_info),
|
||||||
|
daemon=True,
|
||||||
|
)
|
||||||
|
reaper.start()
|
||||||
|
try:
|
||||||
|
assert at_claim.wait(timeout=5), "the unready destroy never reached its claim"
|
||||||
|
# Reserved locally, still running, untracked, and the `del:` marker is
|
||||||
|
# not written yet — exactly the shape reconcile would have adopted
|
||||||
|
# before the reservation wrapped this path.
|
||||||
|
provider._reconcile_orphans()
|
||||||
|
assert "unready-race" not in provider._warm_pool, "reconcile adopted a container this instance is tearing down"
|
||||||
|
finally:
|
||||||
|
let_claim.set()
|
||||||
|
reaper.join(timeout=5)
|
||||||
|
|
||||||
|
# The reservation is released once the stop returns, and the destroy did run.
|
||||||
|
provider._backend.destroy.assert_called_once_with(unready_info)
|
||||||
|
assert provider._local_teardown == set(), "a teardown reservation outlived the stop it guarded"
|
||||||
|
|
||||||
|
|
||||||
|
def test_reconcile_adopts_unready_container_when_no_teardown_is_in_flight(tmp_path, monkeypatch):
|
||||||
|
"""Mirror of the interleaving test: with no destroy running, the same
|
||||||
|
not-yet-registered container *is* adoptable, so the guard above cannot
|
||||||
|
over-block legitimate reconciliation of a container whose creator crashed
|
||||||
|
before the readiness gate.
|
||||||
|
"""
|
||||||
|
aio_mod = importlib.import_module("deerflow.community.aio_sandbox.aio_sandbox_provider")
|
||||||
|
provider, unready_info = _make_unready_destroy_provider(
|
||||||
|
tmp_path,
|
||||||
|
sandbox_id="adoptable",
|
||||||
|
base_url="http://adoptable",
|
||||||
|
monkeypatch=monkeypatch,
|
||||||
|
aio_mod=aio_mod,
|
||||||
|
)
|
||||||
|
provider._unowned_since = {}
|
||||||
|
provider._backend.list_running = MagicMock(return_value=[unready_info])
|
||||||
|
|
||||||
|
provider._reconcile_orphans()
|
||||||
|
|
||||||
|
assert "adoptable" in provider._warm_pool, "reconcile must still adopt a genuinely unowned container"
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user