diff --git a/python/sglang/srt/environ.py b/python/sglang/srt/environ.py index d5a946606..8f93a8e5f 100644 --- a/python/sglang/srt/environ.py +++ b/python/sglang/srt/environ.py @@ -163,6 +163,8 @@ class Envs: SGLANG_LOG_MS = EnvBool(False) SGLANG_DISABLE_REQUEST_LOGGING = EnvBool(False) SGLANG_LOG_REQUEST_EXCEEDED_MS = EnvInt(-1) + SGLANG_LOG_SCHEDULER_STATUS_TARGET = EnvStr("") + SGLANG_LOG_SCHEDULER_STATUS_INTERVAL = EnvFloat(60.0) # SGLang CI SGLANG_IS_IN_CI = EnvBool(False) diff --git a/python/sglang/srt/managers/scheduler_metrics_mixin.py b/python/sglang/srt/managers/scheduler_metrics_mixin.py index 3d30f72fe..9f663bc71 100644 --- a/python/sglang/srt/managers/scheduler_metrics_mixin.py +++ b/python/sglang/srt/managers/scheduler_metrics_mixin.py @@ -20,6 +20,7 @@ from sglang.srt.metrics.collector import ( ) from sglang.srt.utils import get_bool_env_var from sglang.srt.utils.device_timer import DeviceTimer +from sglang.srt.utils.scheduler_status_logger import SchedulerStatusLogger if TYPE_CHECKING: from sglang.srt.managers.scheduler import EmbeddingBatchResult, Scheduler @@ -112,6 +113,8 @@ class SchedulerMetricsMixin: if self.enable_kv_cache_events: self.init_kv_events(self.server_args.kv_events_config) + self.scheduler_status_logger = SchedulerStatusLogger.maybe_create() + def init_kv_events(self: Scheduler, kv_events_config: Optional[str]): if self.enable_kv_cache_events: self.kv_event_publisher = EventPublisherFactory.create( @@ -447,6 +450,9 @@ class SchedulerMetricsMixin: dp_cooperation_info=batch.dp_cooperation_info, ) + if x := self.scheduler_status_logger: + x.maybe_dump(batch, self.waiting_queue) + def log_batch_result_stats( self: Scheduler, batch: ScheduleBatch, diff --git a/python/sglang/srt/utils/scheduler_status_logger.py b/python/sglang/srt/utils/scheduler_status_logger.py new file mode 100644 index 000000000..0d76c9e8c --- /dev/null +++ b/python/sglang/srt/utils/scheduler_status_logger.py @@ -0,0 +1,49 @@ +from __future__ import annotations + +import time +from typing import TYPE_CHECKING, List, Optional + +import torch.distributed as dist + +from sglang.srt.environ import envs +from sglang.srt.utils.log_utils import create_log_targets, log_json + +if TYPE_CHECKING: + from sglang.srt.managers.schedule_batch import Req, ScheduleBatch + + +class SchedulerStatusLogger: + def __init__(self, targets: List[str], dump_interval: float): + self.loggers = create_log_targets(targets=targets, name_prefix=__name__) + self.dump_interval = dump_interval + self.last_dump_time = 0.0 + self.rank = dist.get_rank() if dist.is_initialized() else 0 + + @staticmethod + def maybe_create() -> Optional["SchedulerStatusLogger"]: + target = envs.SGLANG_LOG_SCHEDULER_STATUS_TARGET.get() + if not target: + return None + + return SchedulerStatusLogger( + targets=[t.strip() for t in target.split(",") if t.strip()], + dump_interval=envs.SGLANG_LOG_SCHEDULER_STATUS_INTERVAL.get(), + ) + + def maybe_dump( + self, running_batch: "ScheduleBatch", waiting_queue: List["Req"] + ) -> None: + now = time.time() + if now - self.last_dump_time < self.dump_interval: + return + + self.last_dump_time = now + log_json( + self.loggers, + "scheduler.status", + { + "rank": self.rank, + "running_rids": [r.rid for r in running_batch.reqs], + "queued_rids": [r.rid for r in waiting_queue], + }, + ) diff --git a/test/registered/utils/test_scheduler_status_logger.py b/test/registered/utils/test_scheduler_status_logger.py new file mode 100644 index 000000000..3102e159e --- /dev/null +++ b/test/registered/utils/test_scheduler_status_logger.py @@ -0,0 +1,73 @@ +import json +import os +import shutil +import tempfile +import time +import unittest +from pathlib import Path + +import requests + +from sglang.srt.utils import kill_process_tree +from sglang.test.ci.ci_register import register_cuda_ci +from sglang.test.test_utils import ( + DEFAULT_TIMEOUT_FOR_SERVER_LAUNCH, + DEFAULT_URL_FOR_TEST, + CustomTestCase, + popen_launch_server, +) + +register_cuda_ci(est_time=120, suite="nightly-1-gpu", nightly=True) + + +class TestSchedulerStatusLogger(CustomTestCase): + @classmethod + def setUpClass(cls): + cls.temp_dir = tempfile.mkdtemp() + cls.addClassCleanup(shutil.rmtree, cls.temp_dir) + env = os.environ.copy() + env["SGLANG_LOG_SCHEDULER_STATUS_TARGET"] = cls.temp_dir + env["SGLANG_LOG_SCHEDULER_STATUS_INTERVAL"] = "1" + cls.process = popen_launch_server( + "Qwen/Qwen3-0.6B", + DEFAULT_URL_FOR_TEST, + timeout=DEFAULT_TIMEOUT_FOR_SERVER_LAUNCH, + other_args=["--skip-server-warmup"], + env=env, + ) + cls.addClassCleanup(kill_process_tree, cls.process.pid) + + def test_scheduler_status_dump(self): + response = requests.post( + DEFAULT_URL_FOR_TEST + "/generate", + json={ + "text": "Hello", + "sampling_params": {"max_new_tokens": 8, "temperature": 0}, + }, + timeout=30, + ) + self.assertEqual(response.status_code, 200) + + time.sleep(2) + + events = list(_find_log_events(self.temp_dir, "scheduler.status")) + print(f"{events=}") + self.assertGreater(len(events), 0, "scheduler.status event not found") + data = events[0] + for field in ["timestamp", "rank", "running_rids", "queued_rids"]: + self.assertIn(field, data) + self.assertIsInstance(data["running_rids"], list) + self.assertIsInstance(data["queued_rids"], list) + + +def _find_log_events(log_dir: str, event_name: str): + for f in Path(log_dir).glob("*.log"): + for line in f.read_text().splitlines(): + if line.startswith("{"): + data = json.loads(line) + if data.get("event") == event_name: + yield data + + +if __name__ == "__main__": + unittest.main()