Skip to content
Merged
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
9 changes: 7 additions & 2 deletions xinference/model/llm/vllm/xavier/gpu_transfer.py
Original file line number Diff line number Diff line change
Expand Up @@ -155,6 +155,7 @@ async def _close(self):
self.recv_ref = None
self.caches.clear()
self.send_buffer = self.recv_buffer = None
self.store._gpu_lru.clear()
self.store.blocks.clear()
self.store.ready.clear()
self.store.tiers.clear()
Expand All @@ -170,6 +171,7 @@ async def stage(self, entries):
keys = {key for layers in entries for ids, _ in layers.values() for key in ids}
failed = False
copied = False
last_reused_keys = None
try:
for layers in entries:
for layer, (block_keys, ids) in layers.items():
Expand All @@ -180,9 +182,12 @@ async def stage(self, entries):
and layer in self.store.blocks.get(key, {})
for key in block_keys
):
for key in block_keys:
self.store.blocks.move_to_end(key)
if block_keys != last_reused_keys:
for key in block_keys:
self.store.touch(key)
last_reused_keys = block_keys
continue
last_reused_keys = None
cache = self.caches[layer]
blocks = cache.index_select(
0, torch.tensor(ids, dtype=torch.long, device=cache.device)
Expand Down
12 changes: 12 additions & 0 deletions xinference/model/llm/vllm/xavier/test/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -80,3 +80,15 @@ def connector(connector_module, connector_config):
instance = connector_module.XavierConnector(connector_config, None, caches)
yield instance
instance.shutdown()


@pytest.fixture
def assert_gpu_lru_consistent():
def check(store):
assert set(store.blocks) == set(store.tiers)
gpu_keys = [key for key in store.blocks if store.tiers[key] == "gpu"]
assert list(store._gpu_lru) == gpu_keys
assert store.counts["gpu"] == len(gpu_keys)
assert store.counts["cpu"] == len(store.blocks) - len(gpu_keys)

return check
48 changes: 47 additions & 1 deletion xinference/model/llm/vllm/xavier/test/test_gpu_transfer.py
Original file line number Diff line number Diff line change
Expand Up @@ -343,7 +343,9 @@ async def copy(buffers, refs):

@pytest.mark.asyncio
@pytest.mark.parametrize("failure", [torch.cuda.OutOfMemoryError, RuntimeError])
async def test_staging_failure_drains_and_drops_only_unpublished(monkeypatch, failure):
async def test_staging_failure_drains_and_drops_only_unpublished(
monkeypatch, failure, assert_gpu_lru_consistent
):
r = runtime(monkeypatch, gpu_slots=3)
stage(r, 1)
assert r.store.reserve("2:live", [1])
Expand All @@ -362,6 +364,7 @@ def copy(layer, keys, tensors):
assert r.store.publish([2], {"K"}) == []
assert r.store.leases == {"2:live": {1}}
assert r.store.read("K", [1]).tolist() == [[1, -1]]
assert_gpu_lru_consistent(r.store)


@pytest.mark.asyncio
Expand Down Expand Up @@ -531,6 +534,8 @@ def unexpected(*args, **kwargs):
assert r.store.tiers[1] == ("cpu" if gpu_slots == 1 else "gpu")
assert not gathers
assert not syncs
stage(r, 3)
assert r.store.tiers == {1: "cpu" if gpu_slots == 1 else "gpu", 2: "cpu", 3: "gpu"}


@pytest.mark.asyncio
Expand Down Expand Up @@ -920,6 +925,7 @@ def fence(*args):
assert actor._gpu_transfer is r
assert not r.store.blocks and not r.store.ready and not r.store.tiers
assert not r.store.logical_dtypes and not r.store.leases and not r.store.evicted
assert not r.store._gpu_lru
assert r.store.counts == {"gpu": 0, "cpu": 0}
assert r.send_buffer is None and r.recv_buffer is None
assert not r.recv_refs and not r.caches
Expand Down Expand Up @@ -1238,3 +1244,43 @@ def load(request):
):
connector.start_load_kv(SimpleNamespace())
assert connector._release_load_request.await_count == 2


@pytest.mark.parametrize("second_keys", [[1, 2], [2, 1]])
@pytest.mark.asyncio
async def test_reused_layers_only_refresh_changed_lru_order(monkeypatch, second_keys):
r = runtime(monkeypatch, gpu_slots=2)
r.store = TieredKVSnapshotStore(2, 16, 8, r.device)
for key in [1, 2]:
for layer in ["K", "V"]:
r.store.stage(layer, [key], torch.ones(1, 2, dtype=torch.bfloat16))
r.store.publish([key], {"K", "V"})
touched = []
touch = r.store.touch

def track(key):
touched.append(key)
touch(key)

monkeypatch.setattr(r.store, "touch", track)
await r.stage([{"K": ([1, 2], [0, 1]), "V": (second_keys, [0, 1])}])
assert touched == ([1, 2] if second_keys == [1, 2] else [1, 2, 2, 1])
assert list(r.store.blocks) == second_keys
assert list(r.store._gpu_lru) == second_keys


@pytest.mark.asyncio
async def test_new_snapshot_resets_reused_layer_lru_shortcut(monkeypatch):
r = runtime(monkeypatch, gpu_slots=2)
stage(r, 1)
stage(r, 2)
await r.stage(
[
{"K": ([1, 2], [0, 1])},
{"K": ([3], [2])},
{"K": ([1, 2], [0, 1])},
]
)
assert list(r.store.blocks) == [3, 1, 2]
assert list(r.store._gpu_lru) == [3, 2]
assert r.store.tiers == {1: "cpu", 2: "gpu", 3: "gpu"}
115 changes: 114 additions & 1 deletion xinference/model/llm/vllm/xavier/test/test_tiered_snapshot.py
Original file line number Diff line number Diff line change
Expand Up @@ -148,7 +148,7 @@ def fail(*args, **kwargs):
assert s.counts == {"gpu": 0, "cpu": 0}


def test_leased_cpu_does_not_block_unleased_gpu_replacement():
def test_leased_cpu_does_not_block_unleased_gpu_replacement(assert_gpu_lru_consistent):
s = store(cpu=1)
stage(s, 1)
stage(s, 2)
Expand All @@ -159,6 +159,7 @@ def test_leased_cpu_does_not_block_unleased_gpu_replacement():
assert s.evicted == {2}
assert s.counts == {"gpu": 1, "cpu": 1}
assert s.read("K", [1, 3]).tolist() == [[1, 2], [3, 4]]
assert_gpu_lru_consistent(s)


def test_failed_demotion_preserves_both_tiers(monkeypatch):
Expand Down Expand Up @@ -201,3 +202,115 @@ def test_invalid_later_key_rejects_entire_batch():
assert not s.evicted
assert s.stats() == before
assert s.read("K", [1, 2]).tolist() == [[1, 2], [2, 3]]


@pytest.mark.parametrize("device", ["cpu", "cuda:0"])
def test_demotion_packs_compatible_layers_and_preserves_shapes(monkeypatch, device):
if device.startswith("cuda") and not torch.cuda.is_available():
pytest.skip("CUDA required for physical copies")
layers = {
"K": torch.arange(6, device=device, dtype=torch.float32).reshape(2, 3),
"V": torch.arange(4, device=device, dtype=torch.float32).reshape(4, 1),
"index": torch.arange(3, device=device, dtype=torch.int64),
}
expected = {name: value.cpu().clone() for name, value in layers.items()}
copies = []
original_to = torch.Tensor.to

def track_copy(value, *args, **kwargs):
if value.device.type == "cuda" and args and args[0] == "cpu":
copies.append(value.numel())
return original_to(value, *args, **kwargs)

with monkeypatch.context() as patch:
if device == "cpu":
# Exercise the real grouping/cat/split/reshape path on CPU storage.
# Only the device metadata used for dispatch is substituted.
patch.setattr(
torch.Tensor, "device", property(lambda _: torch.device("cuda:0"))
)
patch.setattr(torch.Tensor, "to", track_copy)
host = TieredKVSnapshotStore._copy_to_cpu(layers)
assert copies == [10, 3]
for name, value in host.items():
assert value.device.type == "cpu"
assert value.dtype == expected[name].dtype
assert value.shape == expected[name].shape
assert torch.equal(value, expected[name])
for value in layers.values():
value.zero_()
assert all(torch.equal(host[name], value) for name, value in expected.items())
host["K"].fill_(42)
assert torch.equal(host["V"], expected["V"])


@pytest.mark.parametrize("device", ["cpu", "cuda:0"])
def test_failed_packed_demotion_preserves_content(
monkeypatch, device, assert_gpu_lru_consistent
):
if device.startswith("cuda") and not torch.cuda.is_available():
pytest.skip("CUDA required for physical copies")
s = TieredKVSnapshotStore(1, 16, 16, torch.device(device))
stage(s, 1)
stage(s, 2)
before = s.stats()
order = list(s.blocks)
original_to = torch.Tensor.to

def fail_host_copy(value, *args, **kwargs):
if value.device.type == "cuda" and args and args[0] == "cpu":
raise RuntimeError("packed copy failed")
return original_to(value, *args, **kwargs)

with monkeypatch.context() as patch:
if device == "cpu":
patch.setattr(
torch.Tensor, "device", property(lambda _: torch.device("cuda:0"))
)
patch.setattr(torch.Tensor, "to", fail_host_copy)
with pytest.raises(RuntimeError, match="packed copy failed"):
stage(s, 3)
assert list(s.blocks) == order
assert s.tiers == {1: "cpu", 2: "gpu"}
assert s.ready == {1, 2}
assert not s.evicted
assert s.stats() == before
assert s.read("K", [1, 2]).tolist() == [[1, 2], [2, 3]]
assert_gpu_lru_consistent(s)


def test_gpu_victims_follow_global_lru_after_hits_and_leases(assert_gpu_lru_consistent):
s = store(gpu=2, cpu=3)
for key in [1, 2, 3]:
stage(s, key)
# A repeated staging hit refreshes GPU LRU without copying the snapshot.
s.stage("K", [2], torch.zeros(1, 2))
stage(s, 4)
assert s.tiers == {1: "cpu", 2: "gpu", 3: "cpu", 4: "gpu"}
assert s.reserve("2:hot", [2])
stage(s, 5)
assert s.tiers[4] == "cpu"
s.release("2:hot")
s.touch(2)
stage(s, 6)
assert s.tiers[5] == "cpu"
assert 1 not in s.blocks
assert list(s._gpu_lru) == [2, 6]
assert_gpu_lru_consistent(s)
assert s.read("K", [2]).tolist() == [[2, 3]]


def test_gpu_lru_tracks_drops_and_failed_first_copy(monkeypatch):
s = store(gpu=2)
stage(s, 1)
s._drop(1)
assert not s._gpu_lru

def fail(*args, **kwargs):
raise RuntimeError("copy failed")

monkeypatch.setattr(torch.Tensor, "to", fail)
with pytest.raises(RuntimeError, match="copy failed"):
stage(s, 2)
assert not s._gpu_lru
assert not s.blocks
50 changes: 43 additions & 7 deletions xinference/model/llm/vllm/xavier/tiered_snapshot.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,8 @@
# Licensed under the Apache License, Version 2.0.
"""Budgeted hot snapshots with CPU overflow and shared publication/read leases."""

from typing import Dict, List
from collections import OrderedDict
from typing import Dict, List, Tuple

import torch

Expand Down Expand Up @@ -34,11 +35,20 @@ def __init__(
self.gpu_capacity = gpu_budget_bytes // block_bytes
self.gpu_device = gpu_device
self.tiers: Dict[int, str] = {}
# GPU victim lookup must not walk an arbitrarily large CPU history.
# Keep the GPU subsequence in the same order as the global block LRU.
self._gpu_lru: OrderedDict[int, None] = OrderedDict()
self.counts = {"gpu": 0, "cpu": 0}
self.metrics = dict(demotions=0, gpu_hits=0, cpu_hits=0, skipped=0)
self.block_bytes = block_bytes

def touch(self, key: int) -> None:
self.blocks.move_to_end(key)
if key in self._gpu_lru:
self._gpu_lru.move_to_end(key)

def _drop(self, key: int) -> None:
self._gpu_lru.pop(key, None)
self.counts[self.tiers.pop(key)] -= 1
del self.blocks[key]
del self.logical_dtypes[key]
Expand All @@ -57,10 +67,34 @@ def _cpu_room(self, pinned: set) -> bool:
self._drop(victim)
return True

@staticmethod
def _copy_to_cpu(layers: Dict[str, torch.Tensor]) -> Dict[str, torch.Tensor]:
# A synchronous D2H copy per layer stalls the producer for every layer
# of every evicted block. Pack compatible layers and synchronize once
# per group. Each host allocation belongs to one block, so evicting a
# block releases its storage without retaining unrelated snapshots.
groups: Dict[
Tuple[torch.device, torch.dtype], List[Tuple[str, torch.Tensor]]
] = {}
host = {}
for layer, value in layers.items():
if value.device.type != "cuda":
host[layer] = value.to("cpu", copy=True)
else:
groups.setdefault((value.device, value.dtype), []).append(
(layer, value)
)
for group in groups.values():
packed = torch.cat([value.reshape(-1) for _, value in group]).to("cpu")
values = packed.split([value.numel() for _, value in group])
for (layer, original), value in zip(group, values):
host[layer] = value.reshape(original.shape)
return host

def _admit(self, key: int, pinned: set) -> bool:
if self.gpu_capacity and self.counts["gpu"] >= self.gpu_capacity:
victim = next(
(k for k in self.blocks if self.tiers[k] == "gpu" and k not in pinned),
(k for k in self._gpu_lru if k not in pinned),
None,
)
if victim is not None:
Expand All @@ -69,13 +103,11 @@ def _admit(self, key: int, pinned: set) -> bool:
)
if cpu_available:
# Copy before evicting CPU content or publishing the new tier.
host = {
layer: value.to("cpu", copy=True)
for layer, value in self.blocks[victim].items()
}
host = self._copy_to_cpu(self.blocks[victim])
self._cpu_room(pinned)
self.blocks[victim] = host
self.tiers[victim] = "cpu"
del self._gpu_lru[victim]
self.counts["gpu"] -= 1
self.counts["cpu"] += 1
self.metrics["demotions"] += 1
Expand All @@ -89,6 +121,8 @@ def _admit(self, key: int, pinned: set) -> bool:
self.blocks[key] = {}
self.logical_dtypes[key] = {}
self.tiers[key] = tier
if tier == "gpu":
self._gpu_lru[key] = None
self.counts[tier] += 1
return True

Expand All @@ -111,7 +145,7 @@ def stage(
for key, tensor in zip(keys, tensors):
if key not in self.blocks and not self._admit(key, pinned):
continue
self.blocks.move_to_end(key)
self.touch(key)
if layer not in self.blocks[key]:
device = (
self.gpu_device if self.tiers[key] == "gpu" else torch.device("cpu")
Expand All @@ -131,6 +165,8 @@ def reserve(self, lease: str, keys: List[int]) -> bool:
if not super().reserve(lease, keys):
return False
for key in keys:
if key in self._gpu_lru:
self._gpu_lru.move_to_end(key)
self.metrics[self.tiers[key] + "_hits"] += 1
return True

Expand Down
Loading