Skip to content

vllm_omni.core.prefix_cache.manager

Omni prefix cache, manager side.

Owns slot occupancy, the request-task table, the hit/span registry, per-step snapshots, and the merge. The controller owns the staging pool, copy queues, and writing rows into the CPU block pool. The state lock covers those tables only — never a wait-for-copy, a GPU-byte-budget flush, or a memcpy.

Helper docstrings mark _state_lock (non-reentrant):

Caller holds   already inside a critical section; do not acquire
Takes          acquires here (``@_locked`` or ``with``)

Two host stores:

StagingBufferPool   reusable step-sized pages. save copies this
                    step's immediately-cached keys (hidden +
                    non-deferred mm) device→host into one page.
                    Per-task `chunk.host` is a view into that page,
                    not a second copy.
PrefixBlockPool     durable (kv_slot, key) prefix cache. The
                    committer only writes into it.

Two write paths (which keys, not how many tokens):

JOIN_NEXT_STEP      immediately-cached keys. Device→host is already
                    in flight at submit; the committer waits
                    `step_d2h_event` then copies host→pool. The next
                    save waits `done`: reused slots must not leave a
                    pending pool write behind.
JOIN_ON_FINISH      deferred mm. Stays on the device clone; the
                    committer does that device→host, then writes the
                    pool. Forced onto the high-priority queue on
                    finish/abort or GPU-byte-budget pressure.

Per real scheduler_output, engine-thread order:

new_step_starts   before _update_states drops finished requests
                  (register hits, start prefix prefetch)
forward
save_outputs      clone off live buffers + launch staging copy;
                  returns step id
materialize or discard_step   exactly one of the two, once

materialize may run on the async output builder while the engine is already in the next step. Warmup/dummy runs are never fed.

logger module-attribute

logger = logging.getLogger(__name__)

OmniPrefixCacheManager

discard_step

discard_step(step_id: int) -> None

Consume the step context when nothing will materialize it.

Any thread; same exactly-once contract as materialize (unknown or duplicate id raises). Only the read-side snapshot is dropped — the cache write proceeds unchanged.

materialize

materialize(
    step_id: int, req_ids: list[str]
) -> StageCacheOutputs

Per-request merged outputs for the step saved as step_id.

Any thread. req_ids must be (a subset of) the save-time snapshot; an outside id means the caller is reading the live batch. A request without a hit is a plain miss and gets exactly this step's rows — normal path, nothing logged. A hit span that resolves to absent rows raises OmniPrefixCacheUnmatchError: fatal by contract (do not pretend it was a miss).

Two phases: under the lock, publish finished writes and pin every row source (task refs + masks, absent checks included) — not yet reading the tensors. Unlocked: wait this step's step_d2h_event, clone the staging views (then drop the step holder), and merge. The engine thread never waits on this thread's device→host copy.

new_step_starts

new_step_starts(scheduler_output: SchedulerOutput) -> None

Handle one scheduler_output.

Engine thread only; before _update_states removes finished requests; exactly once per real step. Registers new-request prefix hits (copying their block tables) and forces finished/aborted requests' still-open writes onto the high-priority copy queue — a block hash that entered the batch must land in the cache, abort included. escalate (eager: the copy + pool write) runs after _state_lock is released.

register_policy

register_policy(policy: ModelCachePolicy) -> None

save_outputs

save_outputs(
    hidden_states: Tensor | None,
    mm_outputs: dict[str, Any] | None,
    *,
    num_tokens_unpadded: int,
    num_tokens_padded: int,
) -> int

Write this step's outputs into the cache; returns the step id.

Engine thread only, after the forward and before materialize. Immediately-cached rows: one on-device clone, one whole-step device→host into the staging pool, then one JOIN_NEXT_STEP WriteTask per request whose chunk.host is a view of that page. Deferred rows stay on the device clone (JOIN_ON_FINISH); the committer copies them later. Leftover mm (this-step deferred rows + mm not written to the pool) is copied to CPU here so materialize never reads live graph buffers. Snapshots everything materialize needs. The returned step id MUST be consumed exactly once — by materialize() or discard_step(). Every step id claims one staging slot (saves with only leftover mm included); a later save waits for a free slot and times out if none return.

The state lock never covers a blocking wait: the previous step's JOIN_NEXT_STEP wait, the clone build, the GPU-byte-budget reserve (which may flush), and the staging-slot claim all run unlocked.

shutdown

shutdown() -> None