From abf3432214827686ca7da8c103d9c6a3bfaeab31 Mon Sep 17 00:00:00 2001 From: Chang Su Date: Tue, 14 Apr 2026 13:50:57 -0700 Subject: [PATCH 1/4] feat(grpc_servicer): implement GetTokenizer RPC for vLLM backend Signed-off-by: Chang Su --- .../smg_grpc_servicer/vllm/servicer.py | 100 +++++++++++++++++- 1 file changed, 98 insertions(+), 2 deletions(-) diff --git a/grpc_servicer/smg_grpc_servicer/vllm/servicer.py b/grpc_servicer/smg_grpc_servicer/vllm/servicer.py index 31d1be63e8..a2c857ce51 100755 --- a/grpc_servicer/smg_grpc_servicer/vllm/servicer.py +++ b/grpc_servicer/smg_grpc_servicer/vllm/servicer.py @@ -5,13 +5,18 @@ Implements the VllmEngine gRPC service on top of vLLM's EngineClient. """ +import hashlib +import io import itertools import time -from collections.abc import AsyncGenerator +import zipfile +from collections.abc import AsyncGenerator, AsyncIterator +from pathlib import Path import grpc import torch from smg_grpc_proto import vllm_engine_pb2, vllm_engine_pb2_grpc +from smg_grpc_proto.generated import common_pb2 from transformers import BatchFeature from vllm import PoolingParams, SamplingParams, TokensPrompt from vllm.engine.protocol import EngineClient @@ -29,6 +34,25 @@ logger = init_logger(__name__) +# Tokenizer bundle streaming constants (aligned with Rust grpc_client limits) +_TOKENIZER_CHUNK_SIZE = 64 * 1024 # 64 KB per gRPC chunk +# Files to include in the tokenizer ZIP bundle. +# Aligned with crates/tokenizer/src/hub.rs:is_tokenizer_file() plus model config files. +_TOKENIZER_FILES = [ + "tokenizer.json", + "tokenizer_config.json", + "config.json", + "generation_config.json", + "special_tokens_map.json", + "vocab.json", + "merges.txt", + "tokenizer.model", # SentencePiece + "tiktoken.model", # tiktoken + "chat_template.json", +] +# Glob patterns for additional tokenizer-related files +_TOKENIZER_GLOBS = ["*.tiktoken", "*.jinja", "*.model"] + # Proto dtype string → torch dtype _PROTO_DTYPE_MAP: dict[str, torch.dtype] = { "float32": torch.float32, @@ -49,13 +73,14 @@ class VllmEngineServicer(vllm_engine_pb2_grpc.VllmEngineServicer): """ gRPC servicer implementing the VllmEngine service. - Handles 6 RPCs: + Handles 7 RPCs: - Generate: Streaming text generation - Embed: Embeddings - HealthCheck: Health probe - Abort: Cancel requests out-of-band - GetModelInfo: Model metadata - GetServerInfo: Server state + - GetTokenizer: Stream tokenizer artifacts """ def __init__(self, async_llm: EngineClient, start_time: float): @@ -358,6 +383,77 @@ async def GetServerInfo( kv_role=kv_role, ) + async def GetTokenizer( + self, + request: common_pb2.GetTokenizerRequest, + context: grpc.aio.ServicerContext, + ) -> AsyncIterator[common_pb2.GetTokenizerChunk]: + """Stream tokenizer artifacts as a ZIP bundle. + + Resolves the tokenizer directory from model_config, zips all relevant + tokenizer files, and streams them as GetTokenizerChunk messages. + The final chunk carries the SHA-256 fingerprint of the full archive. + """ + logger.info("Receive GetTokenizer request") + + tokenizer_path = self.engine.model_config.tokenizer + if not tokenizer_path: + await context.abort( + grpc.StatusCode.FAILED_PRECONDITION, + "Tokenizer path is not configured on this server.", + ) + tokenizer_dir = Path(tokenizer_path) + + # Build ZIP archive in memory + try: + zip_buffer = self._build_tokenizer_zip(tokenizer_dir) + except Exception as e: + logger.exception("Failed to build tokenizer ZIP") + await context.abort(grpc.StatusCode.INTERNAL, str(e)) + + zip_data = zip_buffer.getbuffer() + sha256 = hashlib.sha256(zip_data).hexdigest() + + logger.info( + "Streaming tokenizer bundle: %d bytes, sha256=%s", + len(zip_data), + sha256, + ) + + # Stream chunks; SHA-256 only on the final chunk + offset = 0 + total = len(zip_data) + while offset < total: + end = min(offset + _TOKENIZER_CHUNK_SIZE, total) + is_last = end == total + yield common_pb2.GetTokenizerChunk( + data=bytes(zip_data[offset:end]), + sha256=sha256 if is_last else "", + ) + offset = end + + @staticmethod + def _build_tokenizer_zip(tokenizer_dir: Path) -> io.BytesIO: + """Create an in-memory ZIP archive of tokenizer files from a directory.""" + buf = io.BytesIO() + added: set[str] = set() + with zipfile.ZipFile(buf, "w", zipfile.ZIP_DEFLATED) as zf: + # Exact-name files + for name in _TOKENIZER_FILES: + filepath = tokenizer_dir / name + if filepath.is_file(): + zf.write(filepath, name) + added.add(name) + # Glob patterns (*.tiktoken, *.jinja, *.model) + for pattern in _TOKENIZER_GLOBS: + for match in tokenizer_dir.glob(pattern): + if match.is_file() and match.name not in added: + zf.write(match, match.name) + added.add(match.name) + if not added: + raise FileNotFoundError(f"No tokenizer files found in {tokenizer_dir}") + return buf + # ========== Helper methods ========== def _build_preprocessed_mm_inputs( From a443938c2bdd86b329406d27a30895d174705c78 Mon Sep 17 00:00:00 2001 From: Chang Su Date: Tue, 14 Apr 2026 14:12:56 -0700 Subject: [PATCH 2/4] refactor(grpc_servicer): extract shared tokenizer bundle utilities Move _TOKENIZER_FILES, _TOKENIZER_GLOBS, _TOKENIZER_CHUNK_SIZE, and _build_tokenizer_zip into smg_grpc_servicer/tokenizer_bundle.py so both sglang and vllm servicers import from one source of truth. Signed-off-by: Chang Su --- .../smg_grpc_servicer/sglang/servicer.py | 48 ++--------------- .../smg_grpc_servicer/tokenizer_bundle.py | 53 +++++++++++++++++++ .../smg_grpc_servicer/vllm/servicer.py | 49 ++--------------- 3 files changed, 60 insertions(+), 90 deletions(-) create mode 100644 grpc_servicer/smg_grpc_servicer/tokenizer_bundle.py diff --git a/grpc_servicer/smg_grpc_servicer/sglang/servicer.py b/grpc_servicer/smg_grpc_servicer/sglang/servicer.py index 4e6d9b8243..1e03be470a 100644 --- a/grpc_servicer/smg_grpc_servicer/sglang/servicer.py +++ b/grpc_servicer/smg_grpc_servicer/sglang/servicer.py @@ -8,11 +8,9 @@ import asyncio import dataclasses import hashlib -import io import logging import os import time -import zipfile from collections.abc import AsyncIterator from datetime import datetime, timezone from pathlib import Path @@ -55,29 +53,11 @@ from smg_grpc_servicer.sglang.health_servicer import SGLangHealthServicer from smg_grpc_servicer.sglang.request_manager import GrpcRequestManager from smg_grpc_servicer.sglang.utils import abort_code_from_output +from smg_grpc_servicer.tokenizer_bundle import CHUNK_SIZE, build_tokenizer_zip logger = logging.getLogger(__name__) HEALTH_CHECK_TIMEOUT = int(os.getenv("SGLANG_HEALTH_CHECK_TIMEOUT", 20)) -# Tokenizer bundle streaming constants (aligned with Rust grpc_client limits) -_TOKENIZER_CHUNK_SIZE = 64 * 1024 # 64 KB per gRPC chunk -# Files to include in the tokenizer ZIP bundle. -# Aligned with crates/tokenizer/src/hub.rs:is_tokenizer_file() plus model config files. -_TOKENIZER_FILES = [ - "tokenizer.json", - "tokenizer_config.json", - "config.json", - "generation_config.json", - "special_tokens_map.json", - "vocab.json", - "merges.txt", - "tokenizer.model", # SentencePiece - "tiktoken.model", # tiktoken - "chat_template.json", -] -# Glob patterns for additional tokenizer-related files -_TOKENIZER_GLOBS = ["*.tiktoken", "*.jinja", "*.model"] - def _convert_loads_to_protobuf( result: GetLoadsReqOutput, @@ -596,7 +576,7 @@ async def GetTokenizer( # Build ZIP archive in memory try: - zip_buffer = self._build_tokenizer_zip(tokenizer_dir) + zip_buffer = build_tokenizer_zip(tokenizer_dir) except Exception as e: logger.error(f"Failed to build tokenizer ZIP: {e}\n{get_exception_traceback()}") await context.abort(grpc.StatusCode.INTERNAL, str(e)) @@ -614,7 +594,7 @@ async def GetTokenizer( offset = 0 total = len(zip_data) while offset < total: - end = min(offset + _TOKENIZER_CHUNK_SIZE, total) + end = min(offset + CHUNK_SIZE, total) is_last = end == total yield common_pb2.GetTokenizerChunk( data=bytes(zip_data[offset:end]), @@ -622,28 +602,6 @@ async def GetTokenizer( ) offset = end - @staticmethod - def _build_tokenizer_zip(tokenizer_dir: Path) -> io.BytesIO: - """Create an in-memory ZIP archive of tokenizer files from a directory.""" - buf = io.BytesIO() - added: set[str] = set() - with zipfile.ZipFile(buf, "w", zipfile.ZIP_DEFLATED) as zf: - # Exact-name files - for name in _TOKENIZER_FILES: - filepath = tokenizer_dir / name - if filepath.is_file(): - zf.write(filepath, name) - added.add(name) - # Glob patterns (*.tiktoken, *.jinja, *.model) - for pattern in _TOKENIZER_GLOBS: - for match in tokenizer_dir.glob(pattern): - if match.is_file() and match.name not in added: - zf.write(match, match.name) - added.add(match.name) - if not added: - raise FileNotFoundError(f"No tokenizer files found in {tokenizer_dir}") - return buf - async def SubscribeKvEvents( self, request: common_pb2.SubscribeKvEventsRequest, diff --git a/grpc_servicer/smg_grpc_servicer/tokenizer_bundle.py b/grpc_servicer/smg_grpc_servicer/tokenizer_bundle.py new file mode 100644 index 0000000000..b8c425fe1d --- /dev/null +++ b/grpc_servicer/smg_grpc_servicer/tokenizer_bundle.py @@ -0,0 +1,53 @@ +""" +Shared tokenizer bundle utilities for gRPC servicers. + +Builds and streams tokenizer artifacts as ZIP bundles over gRPC. +Used by both SGLang and vLLM servicers. +""" + +import io +import zipfile +from pathlib import Path + +# Streaming chunk size (aligned with Rust grpc_client limits) +CHUNK_SIZE = 64 * 1024 # 64 KB per gRPC chunk + +# Files to include in the tokenizer ZIP bundle. +# Aligned with crates/tokenizer/src/hub.rs:is_tokenizer_file() plus model config files. +TOKENIZER_FILES = [ + "tokenizer.json", + "tokenizer_config.json", + "config.json", + "generation_config.json", + "special_tokens_map.json", + "vocab.json", + "merges.txt", + "tokenizer.model", # SentencePiece + "tiktoken.model", # tiktoken + "chat_template.json", +] + +# Glob patterns for additional tokenizer-related files +TOKENIZER_GLOBS = ["*.tiktoken", "*.jinja", "*.model"] + + +def build_tokenizer_zip(tokenizer_dir: Path) -> io.BytesIO: + """Create an in-memory ZIP archive of tokenizer files from a directory.""" + buf = io.BytesIO() + added: set[str] = set() + with zipfile.ZipFile(buf, "w", zipfile.ZIP_DEFLATED) as zf: + # Exact-name files + for name in TOKENIZER_FILES: + filepath = tokenizer_dir / name + if filepath.is_file(): + zf.write(filepath, name) + added.add(name) + # Glob patterns (*.tiktoken, *.jinja, *.model) + for pattern in TOKENIZER_GLOBS: + for match in tokenizer_dir.glob(pattern): + if match.is_file() and match.name not in added: + zf.write(match, match.name) + added.add(match.name) + if not added: + raise FileNotFoundError(f"No tokenizer files found in {tokenizer_dir}") + return buf diff --git a/grpc_servicer/smg_grpc_servicer/vllm/servicer.py b/grpc_servicer/smg_grpc_servicer/vllm/servicer.py index a2c857ce51..2f60b9e671 100755 --- a/grpc_servicer/smg_grpc_servicer/vllm/servicer.py +++ b/grpc_servicer/smg_grpc_servicer/vllm/servicer.py @@ -6,10 +6,8 @@ """ import hashlib -import io import itertools import time -import zipfile from collections.abc import AsyncGenerator, AsyncIterator from pathlib import Path @@ -32,26 +30,9 @@ from vllm.outputs import CompletionOutput, RequestOutput from vllm.sampling_params import RequestOutputKind, StructuredOutputsParams -logger = init_logger(__name__) +from smg_grpc_servicer.tokenizer_bundle import CHUNK_SIZE, build_tokenizer_zip -# Tokenizer bundle streaming constants (aligned with Rust grpc_client limits) -_TOKENIZER_CHUNK_SIZE = 64 * 1024 # 64 KB per gRPC chunk -# Files to include in the tokenizer ZIP bundle. -# Aligned with crates/tokenizer/src/hub.rs:is_tokenizer_file() plus model config files. -_TOKENIZER_FILES = [ - "tokenizer.json", - "tokenizer_config.json", - "config.json", - "generation_config.json", - "special_tokens_map.json", - "vocab.json", - "merges.txt", - "tokenizer.model", # SentencePiece - "tiktoken.model", # tiktoken - "chat_template.json", -] -# Glob patterns for additional tokenizer-related files -_TOKENIZER_GLOBS = ["*.tiktoken", "*.jinja", "*.model"] +logger = init_logger(__name__) # Proto dtype string → torch dtype _PROTO_DTYPE_MAP: dict[str, torch.dtype] = { @@ -406,7 +387,7 @@ async def GetTokenizer( # Build ZIP archive in memory try: - zip_buffer = self._build_tokenizer_zip(tokenizer_dir) + zip_buffer = build_tokenizer_zip(tokenizer_dir) except Exception as e: logger.exception("Failed to build tokenizer ZIP") await context.abort(grpc.StatusCode.INTERNAL, str(e)) @@ -424,7 +405,7 @@ async def GetTokenizer( offset = 0 total = len(zip_data) while offset < total: - end = min(offset + _TOKENIZER_CHUNK_SIZE, total) + end = min(offset + CHUNK_SIZE, total) is_last = end == total yield common_pb2.GetTokenizerChunk( data=bytes(zip_data[offset:end]), @@ -432,28 +413,6 @@ async def GetTokenizer( ) offset = end - @staticmethod - def _build_tokenizer_zip(tokenizer_dir: Path) -> io.BytesIO: - """Create an in-memory ZIP archive of tokenizer files from a directory.""" - buf = io.BytesIO() - added: set[str] = set() - with zipfile.ZipFile(buf, "w", zipfile.ZIP_DEFLATED) as zf: - # Exact-name files - for name in _TOKENIZER_FILES: - filepath = tokenizer_dir / name - if filepath.is_file(): - zf.write(filepath, name) - added.add(name) - # Glob patterns (*.tiktoken, *.jinja, *.model) - for pattern in _TOKENIZER_GLOBS: - for match in tokenizer_dir.glob(pattern): - if match.is_file() and match.name not in added: - zf.write(match, match.name) - added.add(match.name) - if not added: - raise FileNotFoundError(f"No tokenizer files found in {tokenizer_dir}") - return buf - # ========== Helper methods ========== def _build_preprocessed_mm_inputs( From deb2b8df3b1c3e994892f24a1ddb62072f7426f7 Mon Sep 17 00:00:00 2001 From: Chang Su Date: Tue, 14 Apr 2026 14:27:48 -0700 Subject: [PATCH 3/4] fix(grpc_servicer): rewind BytesIO cursor in build_tokenizer_zip Add buf.seek(0) before returning so future callers that use read() instead of getbuffer() get the full archive. Signed-off-by: Chang Su --- grpc_servicer/smg_grpc_servicer/tokenizer_bundle.py | 1 + 1 file changed, 1 insertion(+) diff --git a/grpc_servicer/smg_grpc_servicer/tokenizer_bundle.py b/grpc_servicer/smg_grpc_servicer/tokenizer_bundle.py index b8c425fe1d..785b137e4b 100644 --- a/grpc_servicer/smg_grpc_servicer/tokenizer_bundle.py +++ b/grpc_servicer/smg_grpc_servicer/tokenizer_bundle.py @@ -50,4 +50,5 @@ def build_tokenizer_zip(tokenizer_dir: Path) -> io.BytesIO: added.add(match.name) if not added: raise FileNotFoundError(f"No tokenizer files found in {tokenizer_dir}") + buf.seek(0) return buf From f6059366c99ac00c80edb67a7cefdfc27b9262fb Mon Sep 17 00:00:00 2001 From: Chang Su Date: Tue, 14 Apr 2026 14:50:09 -0700 Subject: [PATCH 4/4] fix(grpc_servicer): resolve HF model IDs in vLLM GetTokenizer model_config.tokenizer can be a HuggingFace model ID instead of a local path when vLLM is started with --model meta-llama/... without an explicit --tokenizer. Use snapshot_download(local_files_only=True) to resolve the ID to the HF cache directory. Signed-off-by: Chang Su --- grpc_servicer/smg_grpc_servicer/vllm/servicer.py | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/grpc_servicer/smg_grpc_servicer/vllm/servicer.py b/grpc_servicer/smg_grpc_servicer/vllm/servicer.py index 2f60b9e671..1899a544f7 100755 --- a/grpc_servicer/smg_grpc_servicer/vllm/servicer.py +++ b/grpc_servicer/smg_grpc_servicer/vllm/servicer.py @@ -385,6 +385,16 @@ async def GetTokenizer( ) tokenizer_dir = Path(tokenizer_path) + # model_config.tokenizer may be an HF model ID (e.g. "meta-llama/...") + # rather than a local path. Resolve it to the HF cache directory. + if not tokenizer_dir.is_dir(): + try: + from huggingface_hub import snapshot_download + + tokenizer_dir = Path(snapshot_download(tokenizer_path, local_files_only=True)) + except Exception: + pass # Fall through to build_tokenizer_zip which will raise + # Build ZIP archive in memory try: zip_buffer = build_tokenizer_zip(tokenizer_dir)