Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 36 additions & 0 deletions tests/v1/kv_offload/tiering/test_async_lookup.py
Original file line number Diff line number Diff line change
Expand Up @@ -299,6 +299,42 @@ def test_mark_miss_flips_cached_verdict_without_reprobing(self):
assert "reqA" not in mgr._req_keys
mgr.shutdown()

@pytest.mark.parametrize("flushed", [True, False], ids=["in_flight", "pending"])
def test_mark_miss_skips_newer_unresolved_probe(self, flushed: bool):
"""A failed load reported for an older generation of a key must not
touch a newer probe of that key. Sequence: request A resolves K and
promotes it, A finishes (cleanup deletes the RESOLVED entry), request B
re-probes K, then A's load fails and the tier calls mark_miss([K]).
Forcing B's PENDING / IN_FLIGHT entry to RESOLVED made flush() or
drain_results() assert on the scheduler thread."""
mgr = InMemoryLookupManager(existing_keys={_key(1)})
try:
assert mgr.lookup(_key(1), _ctx("A")) is None
mgr.flush()
mgr._results_ready.wait()
mgr._results_ready.clear()
assert mgr.lookup(_key(1), _ctx("A")) is True
mgr.cleanup("A")
assert _key(1) not in mgr._lookup_state

assert mgr.lookup(_key(1), _ctx("B")) is None
if flushed:
mgr.flush()
mgr._results_ready.wait()
mgr._results_ready.clear()
expected = LookupPhase.IN_FLIGHT if flushed else LookupPhase.PENDING
assert mgr._lookup_state[_key(1)].phase is expected

mgr.mark_miss([_key(1)])
if not flushed:
mgr.flush()
mgr._results_ready.wait()
mgr._results_ready.clear()
# B's own probe result is delivered.
assert mgr.lookup(_key(1), _ctx("B")) is True
finally:
mgr.shutdown()

def test_enqueue_once_invariant_enforced(self):
"""A key is enqueued for probing exactly once, so drain_results() may
receive at most one result per key. Normal operation resolves a key a
Expand Down
27 changes: 27 additions & 0 deletions tests/v1/kv_offload/tiering/test_fs_tier.py
Original file line number Diff line number Diff line change
Expand Up @@ -604,6 +604,33 @@ def test_failed_load_corrects_verdict_and_removes_corrupt_file(
assert lookup_and_wait(tier, [key(1)], ctx=fresh) == [LookupResult.MISS]


def test_failed_load_after_abort_does_not_break_newer_probe(fs_tier):
"""A promotion for request A fails after A was aborted and request B
re-probed the same key. The failure must not flip B's in-flight probe:
B's next lookup drains its own result (MISS: the file is gone) instead of
asserting on the scheduler thread."""
tier, _ = fs_tier
k = key(1)
tier.submit_store(make_job(1, [k], [0]))
assert all(r.success for r in drain(tier))
ctx_a, ctx_b = ReqContext(req_id="A"), ReqContext(req_id="B")
assert lookup_and_wait(tier, [k], ctx=ctx_a) == [LookupResult.HIT]

os.remove(tier.file_mapper.get_file_name(k))
tier.submit_load(make_job(2, [k], [1], is_promotion=True))
tier.on_request_finished(ctx_a)

assert tier.lookup(k, ctx_b) == LookupResult.RETRY
tier.on_schedule_end(ScheduleEndContext(new_req_ids=[], preempted_req_ids=()))
deadline = time.monotonic() + 1.0
while tier._lookup_manager._pending_results.empty():
assert time.monotonic() < deadline
time.sleep(0.01)

assert [r.success for r in drain(tier)] == [False]
assert tier.lookup(k, ctx_b) == LookupResult.MISS


@pytest.mark.parametrize("use_c_ext", [True, False])
def test_batched_partial_load_failure_keeps_loaded_blocks(
fs_tier, monkeypatch, use_c_ext
Expand Down
8 changes: 5 additions & 3 deletions vllm/v1/kv_offload/tiering/async_lookup.py
Original file line number Diff line number Diff line change
Expand Up @@ -216,12 +216,14 @@ def drain_results(self) -> None:
def mark_miss(self, keys: Collection[OffloadKey]) -> None:
"""Force the cached verdict for ``keys`` to False after a failed load, so
the scheduler stops re-issuing the doomed promotion (livelock, #49176).
Keys with no cached entry are skipped."""
Keys with no cached entry are skipped, and so are keys whose entry is
a newer probe still PENDING or IN_FLIGHT: its own result will land,
and forcing it to RESOLVED here would trip the phase asserts in
flush() / drain_results() on the scheduler thread."""
for key in keys:
state = self._lookup_state.get(key)
if state is not None:
if state is not None and state.phase is LookupPhase.RESOLVED:
state.result = False
state.phase = LookupPhase.RESOLVED

def cleanup(self, req_id: str) -> None:
"""Release request references, retaining in-flight lookups.
Expand Down
Loading