[mem] Flatten memory checkers into composable per-pool invariant checks (#22562)
This commit is contained in:
@@ -1191,7 +1191,7 @@ class SchedulerDisaggregationDecodeMixin:
|
||||
self.process_batch_result(batch, result)
|
||||
else:
|
||||
# When the server is idle, do self-check and re-init some states
|
||||
self.self_check_during_idle()
|
||||
self.on_idle()
|
||||
|
||||
# Update last_batch
|
||||
self.last_batch = batch
|
||||
@@ -1224,7 +1224,7 @@ class SchedulerDisaggregationDecodeMixin:
|
||||
tmp_batch, tmp_result = self.result_queue.popleft()
|
||||
self.process_batch_result(tmp_batch, tmp_result)
|
||||
elif batch is None:
|
||||
self.self_check_during_idle()
|
||||
self.on_idle()
|
||||
|
||||
# Run sample of the current batch
|
||||
# It depends on the result of the last batch (e.g., grammar), so we run it after the last batch is processed.
|
||||
|
||||
@@ -409,7 +409,7 @@ class SchedulerDisaggregationPrefillMixin:
|
||||
result = self.run_batch(batch)
|
||||
self.process_batch_result(batch, result)
|
||||
else:
|
||||
self.self_check_during_idle()
|
||||
self.on_idle()
|
||||
|
||||
self.process_disagg_prefill_inflight_queue()
|
||||
|
||||
@@ -448,7 +448,7 @@ class SchedulerDisaggregationPrefillMixin:
|
||||
self.process_batch_result(tmp_batch, tmp_result)
|
||||
elif batch is None:
|
||||
# When the server is idle, do self-check and re-init some states
|
||||
self.self_check_during_idle()
|
||||
self.on_idle()
|
||||
|
||||
self.process_disagg_prefill_inflight_queue()
|
||||
|
||||
|
||||
@@ -1370,7 +1370,7 @@ class Scheduler(
|
||||
self.process_batch_result(batch, result)
|
||||
else:
|
||||
# When the server is idle, do self-check and re-init some states.
|
||||
self.self_check_during_idle()
|
||||
self.on_idle()
|
||||
|
||||
# Update last_batch
|
||||
self.last_batch = batch
|
||||
@@ -1420,7 +1420,7 @@ class Scheduler(
|
||||
pop_and_process()
|
||||
elif batch is None:
|
||||
# When the server is idle, do self-check and re-init some states
|
||||
self.self_check_during_idle()
|
||||
self.on_idle()
|
||||
|
||||
# Run sample of the current batch
|
||||
# It depends on the result of the last batch (e.g., grammar), so we run it after the last batch is processed.
|
||||
|
||||
@@ -142,7 +142,7 @@ class SchedulerPPMixin:
|
||||
|
||||
# When the server is idle, self-check and re-init some states
|
||||
if server_is_idle:
|
||||
self.self_check_during_idle()
|
||||
self.on_idle()
|
||||
|
||||
@DynamicGradMode()
|
||||
def event_loop_pp_disagg_prefill(self: Scheduler):
|
||||
@@ -318,7 +318,7 @@ class SchedulerPPMixin:
|
||||
|
||||
# When the server is idle, self-check and re-init some states
|
||||
if server_is_idle and len(self.disagg_prefill_inflight_queue) == 0:
|
||||
self.self_check_during_idle()
|
||||
self.on_idle()
|
||||
|
||||
@DynamicGradMode()
|
||||
def event_loop_pp_disagg_decode(self: Scheduler):
|
||||
@@ -508,7 +508,7 @@ class SchedulerPPMixin:
|
||||
queue_size += len(self.decode_offload_manager.ongoing_offload)
|
||||
|
||||
if server_is_idle and queue_size == 0:
|
||||
self.self_check_during_idle()
|
||||
self.on_idle()
|
||||
|
||||
def init_pp_loop_state(self: Scheduler):
|
||||
self.pp_loop_size: int = self.pp_size + self.server_args.pp_async_batch_depth
|
||||
|
||||
@@ -228,45 +228,74 @@ class SchedulerRuntimeCheckerMixin:
|
||||
swa_evictable_size=swa_evictable_size,
|
||||
)
|
||||
|
||||
def _check_hybrid_memory(self: Scheduler):
|
||||
pool_stats = self._get_swa_token_info()
|
||||
full_num_used = pool_stats.full_num_used
|
||||
swa_num_used = pool_stats.swa_num_used
|
||||
full_available_size = pool_stats.full_available_size
|
||||
full_evictable_size = pool_stats.full_evictable_size
|
||||
swa_available_size = pool_stats.swa_available_size
|
||||
swa_evictable_size = pool_stats.swa_evictable_size
|
||||
session_held_full = self._session_held_full_tokens()
|
||||
session_held_swa = self._session_held_swa_tokens()
|
||||
|
||||
# Streaming sessions hold tree locks during idle, so tree-protected
|
||||
# tokens must be accounted for alongside session-held tokens.
|
||||
full_protected = self.tree_cache.full_protected_size()
|
||||
swa_protected = self.tree_cache.swa_protected_size()
|
||||
full_leaked = full_num_used - full_protected - session_held_full
|
||||
swa_leaked = swa_num_used - swa_protected - session_held_swa
|
||||
memory_leak = full_leaked != 0 or swa_leaked != 0
|
||||
token_msg = (
|
||||
f"{full_leaked=}, {swa_leaked=}\n"
|
||||
f"{self.full_tokens_per_layer=}, {full_available_size=}, {full_evictable_size=}, {full_protected=}, {session_held_full=}\n"
|
||||
f"{self.swa_tokens_per_layer=}, {swa_available_size=}, {swa_evictable_size=}, {swa_protected=}, {session_held_swa=}\n"
|
||||
@staticmethod
|
||||
def _check_pool_invariant(
|
||||
pool_name: str,
|
||||
available: int,
|
||||
evictable: int,
|
||||
protected: int,
|
||||
session_held: int,
|
||||
total: int,
|
||||
uncached: int = 0,
|
||||
) -> Tuple[bool, str]:
|
||||
"""Check: available + evictable + protected + session_held + uncached == total."""
|
||||
total_accounted = available + evictable + protected + session_held + uncached
|
||||
leak = total_accounted != total
|
||||
msg = (
|
||||
f"[{pool_name}] {total=}, {available=}, {evictable=}, "
|
||||
f"{protected=}, {session_held=}, {uncached=}"
|
||||
)
|
||||
return memory_leak, token_msg
|
||||
return leak, msg
|
||||
|
||||
def _check_mamba_memory(self: Scheduler):
|
||||
pool_stats = self._get_mamba_token_info()
|
||||
full_num_used = pool_stats.full_num_used
|
||||
mamba_num_used = pool_stats.mamba_num_used
|
||||
full_available_size = pool_stats.full_available_size
|
||||
full_evictable_size = pool_stats.full_evictable_size
|
||||
mamba_available_size = pool_stats.mamba_available_size
|
||||
mamba_evictable_size = pool_stats.mamba_evictable_size
|
||||
def _check_full_pool(
|
||||
self: Scheduler, ps: PoolStats, uncached: int = 0
|
||||
) -> Tuple[bool, str]:
|
||||
if self.is_hybrid_swa:
|
||||
protected = self.tree_cache.full_protected_size()
|
||||
session_held = self._session_held_full_tokens()
|
||||
total = self.full_tokens_per_layer
|
||||
elif self.is_hybrid_ssm and self.tree_cache.supports_mamba():
|
||||
protected = self.tree_cache.full_protected_size()
|
||||
session_held = self._session_held_tokens()
|
||||
memory_leak = (
|
||||
full_num_used != self.tree_cache.full_protected_size() + session_held
|
||||
or mamba_num_used != self.tree_cache.mamba_protected_size()
|
||||
total = self.token_to_kv_pool_allocator.size
|
||||
else:
|
||||
protected = self.tree_cache.protected_size()
|
||||
session_held = self._session_held_tokens()
|
||||
total = self.max_total_num_tokens
|
||||
return self._check_pool_invariant(
|
||||
"full",
|
||||
ps.full_available_size,
|
||||
ps.full_evictable_size,
|
||||
protected,
|
||||
session_held,
|
||||
total,
|
||||
uncached,
|
||||
)
|
||||
if memory_leak:
|
||||
|
||||
def _check_swa_pool(
|
||||
self: Scheduler, ps: PoolStats, uncached: int = 0
|
||||
) -> Tuple[bool, str]:
|
||||
return self._check_pool_invariant(
|
||||
"swa",
|
||||
ps.swa_available_size,
|
||||
ps.swa_evictable_size,
|
||||
self.tree_cache.swa_protected_size(),
|
||||
self._session_held_swa_tokens(),
|
||||
self.swa_tokens_per_layer,
|
||||
uncached,
|
||||
)
|
||||
|
||||
def _check_mamba_pool(self: Scheduler, ps: PoolStats) -> Tuple[bool, str]:
|
||||
leak, msg = self._check_pool_invariant(
|
||||
"mamba",
|
||||
ps.mamba_available_size,
|
||||
ps.mamba_evictable_size,
|
||||
self.tree_cache.mamba_protected_size(),
|
||||
0,
|
||||
self.req_to_token_pool.mamba_pool.size,
|
||||
)
|
||||
if leak:
|
||||
# Page-level leak diagnosis for mamba
|
||||
free_full_pages = set(
|
||||
self.token_to_kv_pool_allocator.free_pages.tolist()
|
||||
+ self.token_to_kv_pool_allocator.release_pages.tolist()
|
||||
@@ -288,28 +317,11 @@ class SchedulerRuntimeCheckerMixin:
|
||||
leaked_mamba_pages = (
|
||||
expected_mamba_pages - free_mamba_pages - cached_mamba_pages
|
||||
)
|
||||
token_msg = (
|
||||
f"{full_available_size=}, {full_evictable_size=}, {self.token_to_kv_pool_allocator.size=}, {self.tree_cache.full_protected_size()=}\n"
|
||||
f"{mamba_available_size=}, {mamba_evictable_size=}, {self.req_to_token_pool.mamba_pool.size=}, {self.tree_cache.mamba_protected_size()=}, leaked_full_pages={leaked_full_pages if len(leaked_full_pages) > 0 else None}, leaked_mamba_pages={leaked_mamba_pages if len(leaked_mamba_pages) > 0 else None}\n"
|
||||
msg += (
|
||||
f", leaked_full_pages={leaked_full_pages or None}"
|
||||
f", leaked_mamba_pages={leaked_mamba_pages or None}"
|
||||
)
|
||||
else:
|
||||
token_msg = (
|
||||
f"{full_available_size=}, {full_evictable_size=}, {self.token_to_kv_pool_allocator.size=}, {self.tree_cache.full_protected_size()=}\n"
|
||||
f"{mamba_available_size=}, {mamba_evictable_size=}, {self.req_to_token_pool.mamba_pool.size=}, {self.tree_cache.mamba_protected_size()=}\n"
|
||||
)
|
||||
return memory_leak, token_msg
|
||||
|
||||
def _check_radix_cache_memory(self: Scheduler):
|
||||
pool_stats = self._get_token_info()
|
||||
available_size = pool_stats.full_available_size
|
||||
evictable_size = pool_stats.full_evictable_size
|
||||
protected_size = self.tree_cache.protected_size()
|
||||
session_held = self._session_held_tokens()
|
||||
memory_leak = (available_size + evictable_size) != (
|
||||
self.max_total_num_tokens - protected_size - session_held
|
||||
)
|
||||
token_msg = f"{self.max_total_num_tokens=}, {available_size=}, {evictable_size=}, {protected_size=}, {session_held=}\n"
|
||||
return memory_leak, token_msg
|
||||
return leak, msg
|
||||
|
||||
def _get_batch_uncached_size(self: Scheduler, batch: ScheduleBatch) -> int:
|
||||
ret = 0
|
||||
@@ -327,10 +339,20 @@ class SchedulerRuntimeCheckerMixin:
|
||||
|
||||
return ret
|
||||
|
||||
def self_check_during_busy(self: Scheduler):
|
||||
def _get_total_uncached_size(self: Scheduler) -> int:
|
||||
"""Sum uncached tokens across the current and running batches."""
|
||||
current_batch: ScheduleBatch = self.last_batch
|
||||
uncached_size = self._get_batch_uncached_size(current_batch)
|
||||
if (
|
||||
current_batch.forward_mode.is_extend()
|
||||
and self.running_batch is not None
|
||||
and not self.running_batch.is_empty()
|
||||
):
|
||||
uncached_size += self._get_batch_uncached_size(self.running_batch)
|
||||
return uncached_size
|
||||
|
||||
if current_batch is None:
|
||||
def self_check_during_busy(self: Scheduler):
|
||||
if self.last_batch is None:
|
||||
return
|
||||
|
||||
spec_topk = self.server_args.speculative_eagle_topk or 1
|
||||
@@ -340,35 +362,12 @@ class SchedulerRuntimeCheckerMixin:
|
||||
)
|
||||
return
|
||||
|
||||
pool_stats = self._get_token_info()
|
||||
available_size = pool_stats.full_available_size
|
||||
evictable_size = pool_stats.full_evictable_size
|
||||
protected_size = self.tree_cache.protected_size()
|
||||
|
||||
uncached_size = self._get_batch_uncached_size(current_batch)
|
||||
|
||||
if (
|
||||
current_batch.forward_mode.is_extend()
|
||||
and self.running_batch is not None
|
||||
and not self.running_batch.is_empty()
|
||||
):
|
||||
uncached_size += self._get_batch_uncached_size(self.running_batch)
|
||||
uncached = self._get_total_uncached_size()
|
||||
leak, msg = self._check_full_pool(self.get_pool_stats(), uncached=uncached)
|
||||
|
||||
if envs.SGLANG_ENABLE_STRICT_MEM_CHECK_DURING_BUSY.get() > 1:
|
||||
log_msg = f"[Mem Check (BUSY)] {available_size=}, {evictable_size=}, {protected_size=}, {uncached_size=}"
|
||||
logger.info(log_msg)
|
||||
|
||||
session_held = self._session_held_tokens()
|
||||
total_tokens = (
|
||||
available_size
|
||||
+ evictable_size
|
||||
+ protected_size
|
||||
+ uncached_size
|
||||
+ session_held
|
||||
)
|
||||
assert (
|
||||
total_tokens == self.max_total_num_tokens
|
||||
), f"Mem Leak Detected! {total_tokens=} vs {self.max_total_num_tokens=}"
|
||||
logger.info(f"[Mem Check (BUSY)] {msg}")
|
||||
assert not leak, f"Mem Leak Detected! {msg}"
|
||||
|
||||
def _check_req_pool(self: Scheduler):
|
||||
if self.disaggregation_mode == DisaggregationMode.DECODE:
|
||||
@@ -393,16 +392,8 @@ class SchedulerRuntimeCheckerMixin:
|
||||
msg,
|
||||
)
|
||||
|
||||
def check_memory(self: Scheduler):
|
||||
if self.is_hybrid_swa:
|
||||
memory_leak, token_msg = self._check_hybrid_memory()
|
||||
elif self.is_hybrid_ssm and self.tree_cache.supports_mamba():
|
||||
memory_leak, token_msg = self._check_mamba_memory()
|
||||
else:
|
||||
memory_leak, token_msg = self._check_radix_cache_memory()
|
||||
|
||||
if memory_leak:
|
||||
msg = "token_to_kv_pool_allocator memory leak detected! " f"{token_msg}"
|
||||
def _report_leak(self: Scheduler, pool_name: str, token_msg: str):
|
||||
msg = f"{pool_name} memory leak detected! {token_msg}"
|
||||
raise_error_or_warn(
|
||||
self,
|
||||
envs.SGLANG_ENABLE_STRICT_MEM_CHECK_DURING_IDLE.get(),
|
||||
@@ -410,13 +401,37 @@ class SchedulerRuntimeCheckerMixin:
|
||||
msg,
|
||||
)
|
||||
|
||||
self._check_req_pool()
|
||||
def _check_all_pools(
|
||||
self: Scheduler, ps: PoolStats, uncached: int = 0
|
||||
) -> Tuple[bool, List[str]]:
|
||||
"""Check memory invariant across all pools. Returns (has_leak, messages)."""
|
||||
has_leak = False
|
||||
messages = []
|
||||
|
||||
full_leak, full_msg = self._check_full_pool(ps, uncached=uncached)
|
||||
has_leak |= full_leak
|
||||
messages.append(full_msg)
|
||||
|
||||
if self.is_hybrid_swa:
|
||||
swa_leak, swa_msg = self._check_swa_pool(ps)
|
||||
has_leak |= swa_leak
|
||||
messages.append(swa_msg)
|
||||
|
||||
if self.is_hybrid_ssm and self.tree_cache.supports_mamba():
|
||||
mamba_leak, mamba_msg = self._check_mamba_pool(ps)
|
||||
has_leak |= mamba_leak
|
||||
messages.append(mamba_msg)
|
||||
|
||||
return has_leak, messages
|
||||
|
||||
def _maybe_log_idle_metrics(self: Scheduler):
|
||||
"""Collect and log metrics every 30 seconds during idle."""
|
||||
if (
|
||||
self.current_scheduler_metrics_enabled
|
||||
and time.perf_counter() > self.metrics_collector.last_log_time + 30
|
||||
not self.current_scheduler_metrics_enabled
|
||||
or time.perf_counter() <= self.metrics_collector.last_log_time + 30
|
||||
):
|
||||
# During idle time, also collect metrics every 30 seconds.
|
||||
return
|
||||
|
||||
self.get_pool_stats().update_scheduler_stats(self.stats)
|
||||
|
||||
priority_enabled = self.enable_priority_scheduling
|
||||
@@ -443,9 +458,8 @@ class SchedulerRuntimeCheckerMixin:
|
||||
self.disagg_decode_transfer_queue.queue, priority_enabled
|
||||
)
|
||||
self.metrics_collector.log_stats(self.stats)
|
||||
self._publish_kv_events()
|
||||
|
||||
def check_tree_cache(self: Scheduler):
|
||||
def _check_tree_cache(self: Scheduler):
|
||||
if (
|
||||
self.tree_cache.is_tree_cache()
|
||||
and (self.is_hybrid_swa and self.tree_cache.supports_swa())
|
||||
@@ -453,26 +467,30 @@ class SchedulerRuntimeCheckerMixin:
|
||||
):
|
||||
self.tree_cache.sanity_check()
|
||||
|
||||
def self_check_during_idle(self: Scheduler):
|
||||
if self.enable_hisparse and self.hisparse_coordinator.has_ongoing_staging():
|
||||
return
|
||||
if self.disaggregation_mode == DisaggregationMode.PREFILL:
|
||||
if len(self.disagg_prefill_inflight_queue) > 0:
|
||||
return
|
||||
elif self.disaggregation_mode == DisaggregationMode.DECODE:
|
||||
queue_size = (
|
||||
len(self.waiting_queue)
|
||||
+ len(self.disagg_decode_transfer_queue.queue)
|
||||
+ len(self.disagg_decode_prealloc_queue.queue)
|
||||
)
|
||||
if self.server_args.disaggregation_decode_enable_offload_kvcache:
|
||||
queue_size += len(self.decode_offload_manager.ongoing_offload)
|
||||
if queue_size:
|
||||
def on_idle(self: Scheduler):
|
||||
"""Idle housekeeping: guard, check, metrics, reset, sleep."""
|
||||
if not self.is_fully_idle():
|
||||
return
|
||||
|
||||
self.check_memory()
|
||||
self.check_tree_cache()
|
||||
# memory leak check
|
||||
has_leak, messages = self._check_all_pools(self.get_pool_stats())
|
||||
if has_leak:
|
||||
self._report_leak("pool", "\n".join(messages))
|
||||
self._check_req_pool()
|
||||
|
||||
# tree cache sanity check
|
||||
self._check_tree_cache()
|
||||
|
||||
# metrics every 30s
|
||||
self._maybe_log_idle_metrics()
|
||||
|
||||
# kv event publishing
|
||||
self._publish_kv_events()
|
||||
|
||||
# reset token ratio
|
||||
self.new_token_ratio = self.init_new_token_ratio
|
||||
|
||||
# sleep until next event
|
||||
self.maybe_sleep_on_idle()
|
||||
|
||||
|
||||
@@ -482,16 +500,10 @@ def create_scheduler_watchdog(
|
||||
def dump_info() -> str:
|
||||
if scheduler.is_initializing or disable_request_logging():
|
||||
return ""
|
||||
if scheduler.is_hybrid_swa:
|
||||
_, info_msg = scheduler._check_hybrid_memory()
|
||||
elif scheduler.is_hybrid_ssm and scheduler.tree_cache.supports_mamba():
|
||||
_, info_msg = scheduler._check_mamba_memory()
|
||||
else:
|
||||
_, info_msg = scheduler._check_radix_cache_memory()
|
||||
_, messages = scheduler._check_all_pools(scheduler.get_pool_stats())
|
||||
return (
|
||||
f"{scheduler.cur_batch.batch_size()=}\n"
|
||||
f"{scheduler.cur_batch.reqs=}\n"
|
||||
f"{info_msg}"
|
||||
f"{scheduler.cur_batch.reqs=}\n" + "\n".join(messages)
|
||||
)
|
||||
|
||||
return WatchdogRaw(
|
||||
|
||||
@@ -128,10 +128,7 @@ class SchedulerMultiplexMixin:
|
||||
stream_idx > 0 and self.running_batch.is_empty()
|
||||
)
|
||||
if self.running_batch.is_empty() and self.split_prefill_batch is None:
|
||||
self.check_memory()
|
||||
self.check_tree_cache()
|
||||
self.new_token_ratio = self.init_new_token_ratio
|
||||
self.maybe_sleep_on_idle()
|
||||
self.on_idle()
|
||||
|
||||
if adjust_stream_group:
|
||||
prefill_stream.synchronize()
|
||||
|
||||
Reference in New Issue
Block a user