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.
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 ¶
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.
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.