Skip to content

[Core][KVConnector] Keep the offloaded KV tier across pause/sleep - #57810

Open
aoshen02 wants to merge 1 commit into
vllm-project:mainfrom
aoshen02:fix/pause-keeps-connector-cache
Open

aoshen02 wants to merge 1 commit into
vllm-project:mainfrom
aoshen02:fix/pause-keeps-connector-cache

Conversation

@aoshen02

@aoshen02 aoshen02 commented Sep 20, 2026 •

Copy link
Copy Markdown
Contributor

Summary

pause/sleep exist to hand GPU memory back. An external KV connector tier
(CPU/SSD offload) is the one cache that survives that handover, and the
pause/sleep cascade evicts it. For an RL rollout that time-shares the GPU
(sleep(1) → train → wake_up()), that tier is what makes resumption cheap:
without it the resume re-prefills every in-flight sample.

Whether the tier is still valid depends on the caller replacing the weights,
which the engine cannot know at pause time. Block hashes do not cover the
weights, so the default stays "evict" and this adds an opt-out:

llm.sleep(level=1, clear_connector_cache=False)   # same weights on wake-up

Threaded through pause_scheduler/sleep to _reset_caches(reset_connector=),
the Python and Rust clients, and /pause and /sleep. Default True
everywhere, so nothing changes unless a caller opts out. The Rust client sends
the argument only when opting out, so a new frontend still drives an older
engine.

Measured

2x GB200, slurm job 30001, on the nightly image vllm/vllm-openai:nightly
(0.29.1rc1.dev422), DeepSeek-V4.1-Flash, this branch's files bind-mounted
over the image one by one (mounting the whole package would hide the compiled
_C that ships in it).

Server, on the head node's 4 GPUs:

VLLM_SERVER_DEV_MODE=1 vllm serve $MODEL --served-model-name dsv41 \
  --tensor-parallel-size 4 --enable-prefix-caching --max-model-len 65536 \
  --gpu-memory-utilization 0.85 --port 8712 \
  --kv-transfer-config '{"kv_connector":"OffloadingConnector","kv_role":"kv_both",
    "kv_connector_extra_config":{"block_size":128,"cpu_bytes_to_use":34359738368}}'

Client, from the second node: 16 concurrent completions of ~19.5k tokens each
(311,744 tokens queried in total), then /sleep?level=1 with and without
&clear_connector_cache=false, then /wake_up, then the same 16 prompts
again. Full script in the details block below.

replay after wake-up default clear_connector_cache=false
external prefix-cache hits 0 / 311,744 311,296 / 311,744
kv_offload_load_bytes_total 0 1.31 GB
kv_offload_store_bytes_total 1.31 GB (all recomputed) 0
p50 latency 4.421 s 0.490 s

The counters are the evidence: with the tier cleared the replay stores the same
1.31 GB all over again, with it kept it reads exactly that 1.31 GB back. The
latency column is indicative only — the two arms run in sequence against one
server and share prompts, so the second arm starts warmer.

Two things about this run that a reviewer should know:

  • The patch set mounted into the image also carried [KV Offload] Let unfinished requests keep the chunks they will resume from #57813, so the connector
    config above additionally had "pin_in_flight_chunks": true, which is that
    PR's key and not part of this one. It cannot have moved these numbers: the
    run stores 1.22 GiB into a 32 GiB tier, so nothing is ever evicted and
    pinning has nothing to do. Reproducing this PR alone means dropping that key,
    as the command above does.
  • The client's prompt generator counts words, not tokens, so --prompt-tokens 4096 produces ~19.5k tokens per request. The absolute latencies depend on
    that; the hit/miss counters do not.
Client script
# SPDX-License-Identifier: Apache-2.0
"""Does the offloaded KV tier survive a sleep, and does a resume reuse it?

Sequence per arm:
  1. warm N distinct long prefixes (they land in the CPU tier),
  2. sleep(level=1) -- with or without clear_connector_cache,
  3. wake_up,
  4. replay the same prefixes and read the external-cache hit counters.

With the tier kept, step 4 hits the offloaded blocks; with it cleared, step 4
re-prefills from scratch. The difference is the point of PR #57810.
"""

import argparse
import json
import time
from concurrent.futures import ThreadPoolExecutor

import requests


def metrics(base: str) -> dict[str, float]:
    out: dict[str, float] = {}
    for line in requests.get(f"{base}/metrics", timeout=30).text.splitlines():
        if line.startswith("#") or " " not in line:
            continue
        name, _, value = line.rpartition(" ")
        if "prefix_cache" in name or "offload" in name or "connector" in name:
            try:
                out[name] = out.get(name, 0.0) + float(value)
            except ValueError:
                pass
    return out


def prompt(idx: int, tokens: int) -> str:
    # Distinct prefixes, long enough to fill several offload chunks.
    return f"session-{idx}: " + " ".join(f"w{idx}_{i}" for i in range(tokens))


def fire(base: str, idx: int, tokens: int, max_tokens: int = 8) -> float:
    t0 = time.perf_counter()
    r = requests.post(
        f"{base}/v1/completions",
        json={
            "model": "dsv41",
            "prompt": prompt(idx, tokens),
            "max_tokens": max_tokens,
            "temperature": 0.0,
        },
        timeout=600,
    )
    r.raise_for_status()
    return time.perf_counter() - t0


def run_arm(base: str, keep_tier: bool, n: int, tokens: int) -> dict:
    with ThreadPoolExecutor(max_workers=n) as pool:
        warm = list(pool.map(lambda i: fire(base, i, tokens), range(n)))
    before = metrics(base)

    params = {"level": 1, "mode": "abort"}
    if keep_tier:
        params["clear_connector_cache"] = "false"
    requests.post(f"{base}/sleep", params=params, timeout=600).raise_for_status()
    requests.post(f"{base}/wake_up", timeout=600).raise_for_status()

    with ThreadPoolExecutor(max_workers=n) as pool:
        replay = list(pool.map(lambda i: fire(base, i, tokens), range(n)))
    after = metrics(base)

    delta = {k: after.get(k, 0.0) - before.get(k, 0.0) for k in set(before) | set(after)}
    return {
        "keep_tier": keep_tier,
        "warm_p50_s": sorted(warm)[len(warm) // 2],
        "replay_p50_s": sorted(replay)[len(replay) // 2],
        "metric_delta": {k: v for k, v in sorted(delta.items()) if v},
    }


def main() -> None:
    ap = argparse.ArgumentParser()
    ap.add_argument("--port", type=int, required=True)
    ap.add_argument("--host", default="127.0.0.1")
    ap.add_argument("--requests", type=int, default=16)
    ap.add_argument("--prompt-tokens", type=int, default=4096)
    ap.add_argument("--out", required=True)
    args = ap.parse_args()

    base = f"http://{args.host}:{args.port}"
    result = {
        "cleared": run_arm(base, False, args.requests, args.prompt_tokens),
        "kept": run_arm(base, True, args.requests, args.prompt_tokens),
    }
    result["verdict"] = (
        "tier reused on resume"
        if result["kept"]["replay_p50_s"] < result["cleared"]["replay_p50_s"]
        else "no measurable difference"
    )
    print(json.dumps(result, indent=2))
    with open(args.out, "w") as f:
        json.dump(result, f, indent=2)


if __name__ == "__main__":
    main()

Not a duplicate

_reset_caches / the connector eviction on pause is untouched by any open PR.
#57296 is KV event payloads, #56754 is request draining, #49789 is
graph/runtime state across weight reloads.

Test plan

pytest tests/v1/core/test_scheduler.py -q                          # 177 passed
pytest tests/v1/kv_connector/unit/offloading_connector -q          # 467 passed
pytest tests/v1/engine/test_engine_core.py -q -k "pause or sleep_forwards"  # 10 passed
cargo test -p vllm-server -- pause_route sleep_route                # 5 passed
pre-commit run --files <changed>                                   # passed

The rest of test_engine_core.py starts a real engine and does not run on my
box (another process holds the GPU: "Free memory on device cuda:0 (92.93/139.8
GiB) ... less than desired GPU memory utilization"). CI covers it.

test_pause_clear_connector_cache_opt_out covers both values through
EngineCore.pause_scheduler; test_pause_forwards_connector_flag covers the
immediate and deferred EngineCoreProc paths and that an omitted argument
clears; sleep_route_sends_connector_opt_out_only_when_requested pins the Rust
wire format.

No model-quality impact: sampling, weights and attention are untouched.

Review rounds

Adversarial automated review, several rounds.

Round 1 rejected the original design (removing the eviction outright):
block hashes exclude the weight version and weight updates do not reset
connectors, so that eviction is what prevents stale KV — hence the opt-out. It
then found a P1 in the Rust client, which sent three positional arguments
unconditionally and broke /pause and /sleep against an older engine.

A later round found a P1 that this revision fixes, and it is the reason the
connector file is in this diff.
Keeping the tier crashed the very case the
feature exists for:

sleep(level=1, mode="keep", clear_connector_cache=False) -> wake_up()
AssertionError at offloading/scheduler.py:1792

_reset_caches preempts the running requests, and preempted_req_ids is only
consumed by the next schedule(). Nothing is scheduled while the engine
sleeps, so that set survives the sleep and lands on the step that resumes the
request — and with the tier kept, that step already carries the request's
resume load. The preemption flush asserted that a preempted request can
only hold stores, which is true when preemption and resumption are separate
steps, and false here.

The flush now takes the store jobs and leaves the load alone: a load reads no
freed block, so it is not what the flush is there to settle. Default clearing
never reached this, because with the tier gone there is no load to find —
which is also why the completed-request measurement above could not have
caught it. test_preempted_request_resumed_in_the_same_step covers both
scheduling modes and fails on the previous revision.

The opt-out is now refused where it is not safe

The flag reaches whatever connector is configured, and the same-step resume is
a state every connector that tracks per-request transfers has to expect. Only
OffloadingConnector does, after the fix above. MooncakeStoreConnector also
implements reset_cache, also consumes preempted_req_ids, and pops
load_specs[req_id] and _unfinished_requests[req_id] for them
(mooncake/store/scheduler.py:226-231) — the entries its new-request loop
then asserts on four lines later.

Rather than change a second connector I cannot exercise here, the opt-out is
fail-closed: connectors declare supports_retained_cache_on_pause (default
False, True on OffloadingConnector, the conjunction of its children on
MultiConnector), and a pause that asks to keep a cache the connector cannot
keep is refused before anything is paused:

MooncakeStoreConnector cannot keep its cache across a pause: the resume lands
in the step that reports the pause's preemptions, which this connector does
not expect. Drop clear_connector_cache=False.

Making Mooncake handle it is a separate, testable change; until then its users
get an error at the call instead of a crash at the wake-up.

Two related notes:

The same round also asked for, and this revision adds, the /pause opt-out
wire test, a sleep(level>=1) threading test, the one-line note on why the
Python client may always send the argument (client and engine ship together,
so they cannot disagree), and the parameter's documentation on
EngineCoreProc.pause_scheduler, which is the docstring MP clients read.

A further round, told about all of the above, checked for any path where the
flag defaults to False by accident, the reverse gRPC skew (old frontend, new
engine), whether the typed bool query parameter changes the accepted request
format (False/false/1/0), and whether the deferred
engine_idle_callback closure can run with a flag a later pause changed. No
defects.

AI assistance

Developed with AI assistance (Claude). I reviewed every changed line, ran the
tests above, and I am the accountable submitter.

🤖 Generated with Claude Code

@claude claude Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Claude Code Review

This pull request is from a fork — automated review is disabled. A repository maintainer can comment @claude review to run a one-time review.

@aoshen02

Copy link
Copy Markdown
Contributor Author

/ci run

@github-actions

Copy link
Copy Markdown

✅ Triggered Buildkite CI #90079 for commit 68610e535619.

@mergify mergify Bot added the scheduler label Sep 20, 2026
@aoshen02
aoshen02 force-pushed the fix/pause-keeps-connector-cache branch from 68610e5 to 24d8396 Compare September 20, 2026 15:17
@aoshen02

Copy link
Copy Markdown
Contributor Author

/ci run

@mergify mergify Bot added the frontend label Sep 20, 2026
@github-actions

Copy link
Copy Markdown

❌ This PR is 1 commit behind upstream main. Your branch must contain every commit currently on upstream main. No new CI build was started. Merge or rebase onto the latest main, then rerun /ci run. To test this branch at your own risk, use /ci run --allow-stale.

@aoshen02
aoshen02 force-pushed the fix/pause-keeps-connector-cache branch from 24d8396 to 758382b Compare September 20, 2026 15:20
@aoshen02

Copy link
Copy Markdown
Contributor Author

/ci run

@github-actions

Copy link
Copy Markdown

❌ This PR is 1 commit behind upstream main. Your branch must contain every commit currently on upstream main. No new CI build was started. Merge or rebase onto the latest main, then rerun /ci run. To test this branch at your own risk, use /ci run --allow-stale.

@mergify

mergify Bot commented Sep 20, 2026

Copy link
Copy Markdown
Contributor

Documentation preview: https://vllm--57810.org.readthedocs.build/en/57810/

@mergify mergify Bot added the documentation Improvements or additions to documentation label Sep 20, 2026
@aoshen02

Copy link
Copy Markdown
Contributor Author

/ci run

@aoshen02
aoshen02 force-pushed the fix/pause-keeps-connector-cache branch from 758382b to b3315a4 Compare September 20, 2026 16:05
@aoshen02
aoshen02 requested a review from BugenZhao as a code owner September 20, 2026 16:05
@github-actions

Copy link
Copy Markdown

❌ This PR is 1 commit behind upstream main. Your branch must contain every commit currently on upstream main. No new CI build was started. Merge or rebase onto the latest main, then rerun /ci run. To test this branch at your own risk, use /ci run --allow-stale.

@mergify mergify Bot added the rust label Sep 20, 2026
@aoshen02

Copy link
Copy Markdown
Contributor Author

/ci run

@aoshen02
aoshen02 force-pushed the fix/pause-keeps-connector-cache branch from b3315a4 to 37c675d Compare September 20, 2026 16:39
@mergify

mergify Bot commented Sep 21, 2026

Copy link
Copy Markdown
Contributor

This pull request has merge conflicts that must be resolved before it can be
merged. Please rebase the PR, @aoshen02.

https://docs.github.com/en/pull-requests/collaborating-with-pull-requests/working-with-forks/syncing-a-fork

@mergify mergify Bot added the needs-rebase label Sep 21, 2026
@aoshen02
aoshen02 force-pushed the fix/pause-keeps-connector-cache branch from bcfbd48 to 9ce24e4 Compare September 21, 2026 01:32
@aoshen02

Copy link
Copy Markdown
Contributor Author

/ci run

@github-actions

Copy link
Copy Markdown

✅ Triggered Buildkite CI #90115 for commit 9ce24e453dbf.

@aoshen02
aoshen02 force-pushed the fix/pause-keeps-connector-cache branch from 9ce24e4 to a94ff5c Compare September 21, 2026 02:21
@mergify mergify Bot added the kv-connector label Sep 21, 2026
@aoshen02

Copy link
Copy Markdown
Contributor Author

/ci run

@mergify mergify Bot removed the needs-rebase label Sep 21, 2026
@github-actions

Copy link
Copy Markdown

❌ This PR is 3 commits behind upstream main. Your branch must contain every commit currently on upstream main. No new CI build was started. Merge or rebase onto the latest main, then rerun /ci run. To test this branch at your own risk, use /ci run --allow-stale.

An external KV connector tier is the one cache that survives a
pause/sleep: the GPU blocks are discarded, the offloaded copy is not.
Today the pause/release cascade always evicts it, so a caller that only
wants the GPU back (an RL rollout time-sharing the device with training)
loses the very thing that makes resuming cheap and has to re-prefill.

Add clear_connector_cache to pause_generation()/sleep() and thread it to
EngineCore._reset_caches(). It defaults to True, so behavior is
unchanged: block hashes do not cover the weights, and weight updates do
not reset connectors, so callers that replace the weights must keep
evicting the tier or they would serve stale KV. Callers that leave the
weights alone can now opt out.

Signed-off-by: Ao Shen <aoshen@inferact.ai>
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

Signed-off-by: aoshen02 <aoshen@inferact.ai>
@aoshen02
aoshen02 force-pushed the fix/pause-keeps-connector-cache branch from a94ff5c to 6f470a0 Compare September 21, 2026 02:39
@aoshen02

Copy link
Copy Markdown
Contributor Author

/ci run

@github-actions

Copy link
Copy Markdown

❌ This PR is 3 commits behind upstream main. Your branch must contain every commit currently on upstream main. No new CI build was started. Merge or rebase onto the latest main, then rerun /ci run. To test this branch at your own risk, use /ci run --allow-stale.

@mergify

mergify Bot commented Sep 30, 2026

Copy link
Copy Markdown
Contributor

This pull request has merge conflicts that must be resolved before it can be
merged. Please rebase the PR, @aoshen02.

https://docs.github.com/en/pull-requests/collaborating-with-pull-requests/working-with-forks/syncing-a-fork

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant