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
63 changes: 63 additions & 0 deletions tests/v1/kv_connector/unit/test_mooncake_store_connector.py
Original file line number Diff line number Diff line change
Expand Up @@ -614,3 +614,66 @@ def test_lookup_key_server_reset_skips_drain_when_no_send_thread():

assert call_order == ["remove_all"]
assert sent == [protocol.RESP_OK]


def test_shutdown_closes_worker_store():
vllm_config = _make_vllm_config()
kv_cache_config = _make_kv_cache_config()

with (
set_current_vllm_config(vllm_config),
patch(
"vllm.distributed.kv_transfer.kv_connector.v1.mooncake.store."
"connector.MooncakeStoreWorker"
) as mock_worker_cls,
):
connector = mooncake_store_connector.MooncakeStoreConnector(
vllm_config, KVConnectorRole.WORKER, kv_cache_config
)

worker = mock_worker_cls.return_value
connector.shutdown()

worker.close.assert_called_once_with()


def test_del_invokes_shutdown_and_closes_store():
vllm_config = _make_vllm_config()
kv_cache_config = _make_kv_cache_config()

with (
set_current_vllm_config(vllm_config),
patch(
"vllm.distributed.kv_transfer.kv_connector.v1.mooncake.store."
"connector.MooncakeStoreWorker"
) as mock_worker_cls,
):
connector = mooncake_store_connector.MooncakeStoreConnector(
vllm_config, KVConnectorRole.WORKER, kv_cache_config
)

worker = mock_worker_cls.return_value
# __del__ is the GC backstop; it must route through shutdown() -> close().
connector.__del__()

worker.close.assert_called_once_with()


def test_shutdown_scheduler_role_is_noop():
vllm_config = _make_vllm_config()
kv_cache_config = _make_kv_cache_config()

with (
set_current_vllm_config(vllm_config),
patch(
"vllm.distributed.kv_transfer.kv_connector.v1.mooncake.store."
"connector.MooncakeStoreScheduler"
),
):
connector = mooncake_store_connector.MooncakeStoreConnector(
vllm_config, KVConnectorRole.SCHEDULER, kv_cache_config
)

# Scheduler role holds no store handle, so shutdown must be a safe no-op.
assert connector.connector_worker is None
connector.shutdown()
31 changes: 31 additions & 0 deletions tests/v1/kv_connector/unit/test_mooncake_store_worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -1558,3 +1558,34 @@ def test_lookup_records_mooncake_metrics():
assert isinstance(stats, MooncakeStoreConnectorStats)
assert len(stats.data["lookup_exists"]) == 1
assert stats.data["lookup_exists"][0]["num_keys"] == 2


def test_store_worker_close_releases_store():
worker = _make_bare_worker()
store = worker.store

worker.close()

store.close.assert_called_once_with()
assert worker.store is None


def test_store_worker_close_is_idempotent():
worker = _make_bare_worker()
store = worker.store

worker.close()
worker.close()

# Second call short-circuits because store was already released.
store.close.assert_called_once_with()


def test_store_worker_close_swallows_store_errors():
worker = _make_bare_worker()
worker.store.close.side_effect = RuntimeError("boom")

# A failure tearing down the store must not propagate out of close().
worker.close()

assert worker.store is None
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,21 @@ def __init__(
else:
self.connector_worker = MooncakeStoreWorker(vllm_config, kv_cache_config)

def shutdown(self):
"""Release connector resources on teardown.

Closes the worker's MooncakeDistributedStore handle so its
TransferEngine and RDMA registrations are released. Invoked from the
engine's explicit shutdown path and as a backstop from ``__del__``;
a no-op on the scheduler role, which holds no store handle.
"""
worker = getattr(self, "connector_worker", None)
if worker is not None:
worker.close()
Comment thread
Dao007forever marked this conversation as resolved.

def __del__(self):
self.shutdown()

# ============================================================
# Scheduler-side methods
# ============================================================
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1423,6 +1423,22 @@ def get_kv_events(self) -> list[BlockStored]:
return self.kv_send_thread.get_kv_events()
return []

def close(self) -> None:
"""Release the MooncakeDistributedStore handle on teardown.

Closing the store frees its TransferEngine, the registered RDMA
buffers, and the connection to the master server. Idempotent so it is
safe to call from both the explicit shutdown path and ``__del__``.
"""
store = getattr(self, "store", None)
if store is None:
return
self.store = None
try:
store.close()
except Exception as e:
logger.warning("Error closing MooncakeDistributedStore: %s", e)


# ============================================================
# Lookup Key Server
Expand Down
Loading