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
14 changes: 14 additions & 0 deletions docs/features/kv_offloading_usage.md
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,7 @@ vllm serve <model> \
| `store_threshold` | no | `0` | single-tier | Min lookups before a block is offloaded. Values ≥ 2 are rejected by `TieringOffloadingSpec`. |
| `max_tracker_size` | no | `64000` | single-tier | Max entries in the lookup tracker. |
| `secondary_tiers` | no | `[]` | multi-tier | List of secondary tier configs (see below). |
| `tier_poll_interval_s` | no | `0.001` | multi-tier | Pause between control-plane rounds for tiers that must be serviced between engine steps (currently `p2p`). See [Threading model](#threading-model). Set to `0` to disable the thread and service those tiers once per engine step instead. |
| `offload_prompt_only` | no | `true` | both | If `true`, only prompt (prefill) blocks are offloaded; decode blocks are skipped. |
| `self_describing_kv_events` | no | `false` | both | Opt-in. When `true` *and* KV cache events are enabled (`--kv-events-config` with `enable_kv_cache_events`), the connector emits self-describing block-granular `BlockStored`/`BlockRemoved` payloads (constituent block hashes, whole-chunk `token_ids`, per-block `block_size`, parent hash, LoRA + group/cache-spec metadata) instead of the placeholder fallback, so external KV-event consumers can index offloaded blocks. Inert unless events are enabled. With `TieringOffloadingSpec`, a CPU promotion is self-describing when a local request observes its primary-tier `HIT` before event translation; otherwise its stored event may retain the placeholder, while a later `HIT` can backfill metadata for removal. Pending-removal/re-promotion races and externally initiated promotions may also produce placeholders, and consumers must ignore removals for unknown hashes. Partial recurrent tails emit the hash-aligned portion from the physical block start through the tail boundary. Other sliding-window/SSM chunks keep the placeholder fallback. In chunk mode (`block_size` > GPU block size, or `blocks_per_chunk` > 1), overlapping chunks re-announce shared per-block hashes, so consumers must reference-count (deduplicate) repeated store/remove announcements. |
| `spec_module_path` | no | — | both | Python import path for a custom `OffloadingSpec` not in the built-in registry. Required only when `spec_name` is not built-in (advanced). |
Expand Down Expand Up @@ -150,6 +151,19 @@ The filesystem and object-store tiers can publish hash-only `BlockStored` KV eve

Set the optional `locality` tier field to `LOCAL` or `REMOTE` to describe the tier's storage location relative to the publishing vLLM instance. `LOCAL` marks storage local to that instance, while `REMOTE` marks storage that is not local to it. When the setting is omitted, locality is unspecified. vLLM does not infer it from the tier type, so an `obj` tier is not implicitly `REMOTE`. A KV event includes `locality` only when the tier explicitly configures it. This metadata describes the tier property without implying that a consumer can already route requests to its blocks.

### Threading model

Tier methods run in the scheduler process, under the tiering manager's lock, so they are never entered concurrently. They are not always called from the same thread, though.

Most tiers only ever see the scheduler thread, which holds the lock for the length of one engine step. A tier whose counterpart is not driven by the engine — the `p2p` tier, answering a remote peer — needs servicing while the engine is busy running the model, because otherwise a peer's lookup or fetch waits for the next step boundary. On a saturated rank that wait can be seconds, dwarfing the transfer it gates. Such tiers set `serves_external_requests`, and the tiering manager runs one thread that, once per `tier_poll_interval_s`, polls them for finished jobs and lets them serve external requests. A tier that implements `serve_external_requests()` without setting the flag is served once per engine step only.

Two consequences worth knowing:

- Any tier can be reached from that thread indirectly, because serving a peer looks keys up through the tiering manager, which fans out to the other tiers. Tier state must therefore not assume a particular thread, only that the manager lock is held.
- Work the thread deliberately leaves alone is batched per step: a promotion it initiates is submitted at the next `on_schedule_end()`, not immediately.

Setting `tier_poll_interval_s` to `0` disables the thread; opted-in tiers then get serviced once per engine step, as they were before it existed.

### Filesystem (FS)

The filesystem tier (`type: "fs"`) writes blocks to a filesystem directory.
Expand Down
123 changes: 123 additions & 0 deletions tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -1017,6 +1017,129 @@ def test_scheduler_reports_lookup_sync_delay(request_runner):
assert reduced[f"{_ConnectorMetricName.LOOKUP_SYNC_DELAY}_sum"] > 0


def test_on_schedule_end_is_the_last_manager_call_of_a_step(request_runner):
"""on_schedule_end must run after every other manager call of the step.

Callers bound a step's exclusive access to the manager by releasing at
on_schedule_end, so anything issued after it — notably prepare_store, which
evicts — would fall outside that window.
"""
block_size = 4
runner = request_runner(
block_size=block_size,
num_gpu_blocks=8,
async_scheduling=False,
)
runner.manager.prepare_store.side_effect = lambda keys, req_context: (
generate_store_output([])
)

connector_scheduler = runner.connector_scheduler
build_meta = connector_scheduler.build_connector_meta
windows: list[list[str]] = []

def spy(scheduler_output):
first = len(runner.manager.mock_calls)
meta = build_meta(scheduler_output)
# Entering and exiting manager.lock brackets the step; only the state
# operations it guards are ordered relative to on_schedule_end.
windows.append(
[
name
for name, _, _ in runner.manager.mock_calls[first:]
if not name.startswith("lock.")
]
)
return meta

connector_scheduler.build_connector_meta = spy

runner.new_request(token_ids=[0] * (block_size * 2))
runner.run(decoded_tokens=[EOS_TOKEN_ID])

assert windows, "build_connector_meta was never called"
assert any("prepare_store" in w for w in windows), (
"step must exercise the calls that follow on_schedule_end today"
)
assert any("on_request_finished" in w for w in windows)
for names in windows:
assert names[-1] == "on_schedule_end", (
f"on_schedule_end must be the last manager call of a step, got {names}"
)


def test_step_releases_the_manager_before_the_model_runs(request_runner):
"""The step-scoped region must close when build_connector_meta returns.

Everything after it -- the model future above all -- runs with the manager
free, which is the only window another thread has to use it.
"""
runner = request_runner(
block_size=4,
num_gpu_blocks=10,
async_scheduling=False,
)
runner.manager.prepare_store.side_effect = lambda keys, req_context: (
generate_store_output([])
)

runner.new_request(token_ids=[1] * 4)
runner.run(decoded_tokens=[EOS_TOKEN_ID])

scheduler = runner.connector_scheduler
assert scheduler._step_region is None
lock = runner.manager.lock
assert lock.__enter__.call_count > 0
assert lock.__enter__.call_count == lock.__exit__.call_count


def test_has_pending_push_work_does_not_open_the_step_region(request_runner):
"""It must take a short region, never the step-scoped one.

The engine asks this before schedule() and returns early when the answer is
False, so a step-scoped acquire here would never reach the release in
build_connector_meta and would hold the manager for good.
"""
runner = request_runner(
block_size=4,
num_gpu_blocks=10,
async_scheduling=False,
)
scheduler = runner.connector_scheduler
runner.manager.has_pending_work.return_value = False

assert scheduler.has_pending_push_work() is False
assert scheduler._step_region is None
lock = runner.manager.lock
assert lock.__enter__.call_count == lock.__exit__.call_count


def test_take_events_drains_the_manager_inside_the_region(request_runner):
"""take_events must not hold the manager across its own iteration.

manager.take_events() is a generator whose backing list is cleared once
exhausted, so draining it lazily would both pin the manager for as long as
the consumer takes and lose events appended in the meantime.
"""
runner = request_runner(
block_size=4,
num_gpu_blocks=10,
async_scheduling=False,
)
scheduler = runner.connector_scheduler
runner.manager.take_events.return_value = iter(())
runner.manager.reset_mock()

events = list(scheduler.take_events())

assert events == []
assert runner.manager.take_events.call_count == 1
assert scheduler._step_region is None
lock = runner.manager.lock
assert lock.__enter__.call_count == 1
assert lock.__exit__.call_count == 1


def test_scheduler_reports_lookup_async_delay_on_resolve(request_runner):
"""A deferred lookup reports its async delay once it resolves."""
runner = request_runner(
Expand Down
Loading
Loading