From 572362017e527f308e31ceff93bf54f517692e10 Mon Sep 17 00:00:00 2001 From: jain-ria Date: Fri, 7 Aug 2026 09:45:53 -0700 Subject: [PATCH 1/7] fix(frontend): return SGLang chat logprobs Signed-off-by: jain-ria --- .../src/dynamo/frontend/sglang_prepost.py | 50 +++++++++++++++++-- .../src/dynamo/frontend/sglang_processor.py | 16 ++++++ 2 files changed, 63 insertions(+), 3 deletions(-) diff --git a/components/src/dynamo/frontend/sglang_prepost.py b/components/src/dynamo/frontend/sglang_prepost.py index 9cd9400df44c..5242df90a905 100644 --- a/components/src/dynamo/frontend/sglang_prepost.py +++ b/components/src/dynamo/frontend/sglang_prepost.py @@ -1044,6 +1044,43 @@ def _incremental_decode( self._pending_decode_ids = [] return delta_text + def _build_openai_logprobs( + self, + log_probs: list[float], + top_logprobs: list[list[dict[str, Any]]] | None, + token_ids: list[int], + ) -> dict[str, Any] | None: + if len(log_probs) != len(token_ids): + return None + + content: list[dict[str, Any]] = [] + for index, (token_id, logprob) in enumerate(zip(token_ids, log_probs)): + token = self.tokenizer.decode([token_id], skip_special_tokens=False) + candidates = top_logprobs[index] if top_logprobs else [] + openai_top_logprobs = [] + for candidate in candidates: + candidate_token = candidate.get("token") or "" + candidate_bytes = candidate.get("bytes") + if candidate_bytes is None and candidate_token: + candidate_bytes = list(candidate_token.encode("utf-8")) + openai_top_logprobs.append( + { + "token": candidate_token, + "logprob": float(candidate["logprob"]), + "bytes": candidate_bytes, + } + ) + content.append( + { + "token": token, + "logprob": float(logprob), + "bytes": list(token.encode("utf-8")) if token else None, + "top_logprobs": openai_top_logprobs, + } + ) + + return {"content": content, "refusal": None} if content else None + def _parse_reasoning_delta( self, delta_text: str, finish_reason: str | None ) -> tuple[str | None, str]: @@ -1106,6 +1143,8 @@ def process_output(self, engine_response: dict[str, Any]) -> dict[str, Any] | No raw_ids = engine_response.get("token_ids") token_ids = raw_ids if isinstance(raw_ids, list) else list(raw_ids or []) finish_reason = engine_response.get("finish_reason") + log_probs = engine_response.get("log_probs") + top_logprobs = engine_response.get("top_logprobs") if finish_reason is not None: token_ids = self._strip_trailing_eos_token_ids(list(token_ids)) @@ -1114,6 +1153,11 @@ def process_output(self, engine_response: dict[str, Any]) -> dict[str, Any] | No if token_ids or finish_reason is not None else "" ) + openai_logprobs = None + if log_probs is not None: + openai_logprobs = self._build_openai_logprobs( + log_probs, top_logprobs, token_ids + ) if self._fast_plain_text: if delta_text: @@ -1121,14 +1165,14 @@ def process_output(self, engine_response: dict[str, Any]) -> dict[str, Any] | No "index": 0, "delta": self._with_initial_role({"content": delta_text}), "finish_reason": finish_reason, - "logprobs": None, + "logprobs": openai_logprobs, } elif finish_reason: return { "index": 0, "delta": self._with_initial_role({}), "finish_reason": finish_reason, - "logprobs": None, + "logprobs": openai_logprobs, } return None @@ -1345,7 +1389,7 @@ def process_output(self, engine_response: dict[str, Any]) -> dict[str, Any] | No "index": 0, "delta": self._with_initial_role(delta), "finish_reason": effective_finish, - "logprobs": None, + "logprobs": openai_logprobs, } return None diff --git a/components/src/dynamo/frontend/sglang_processor.py b/components/src/dynamo/frontend/sglang_processor.py index a63154b67809..3da4dd0bfdf8 100644 --- a/components/src/dynamo/frontend/sglang_processor.py +++ b/components/src/dynamo/frontend/sglang_processor.py @@ -625,6 +625,8 @@ async def _generate_and_stream( # finish_reason. Use si=1 for the first chunk to minimize # TTFT, then switch to the configured interval. pending_token_ids: list[int] = [] + pending_log_probs: list[float] | None = None + pending_top_logprobs: list[list[dict[str, Any]]] | None = None pending_usage: dict[str, Any] | None = None first_chunk = True input_tokens = len(tokens) @@ -674,6 +676,14 @@ async def _generate_and_stream( engine_data = engine_response.get("engine_data") pending_token_ids.extend(new_ids) + if (log_probs := engine_response.get("log_probs")) is not None: + if pending_log_probs is None: + pending_log_probs = [] + pending_log_probs.extend(log_probs) + if (top_logprobs := engine_response.get("top_logprobs")) is not None: + if pending_top_logprobs is None: + pending_top_logprobs = [] + pending_top_logprobs.extend(top_logprobs) # Flush on finish or when we've accumulated enough tokens. # First chunk flushes immediately (si=1) to minimize TTFT. @@ -684,6 +694,10 @@ async def _generate_and_stream( "token_ids": pending_token_ids, "finish_reason": finish_reason, } + if pending_log_probs is not None: + mapped_response["log_probs"] = pending_log_probs + if pending_top_logprobs is not None: + mapped_response["top_logprobs"] = pending_top_logprobs if self.debug_perf: t_pp0 = time.monotonic() @@ -742,6 +756,8 @@ async def _generate_and_stream( yield envelope pending_token_ids = [] + pending_log_probs = None + pending_top_logprobs = None pending_usage = None first_chunk = False except Unknown: From e541c61194f13f4729ebc0d1796cc11cb7a16c2e Mon Sep 17 00:00:00 2001 From: jain-ria Date: Fri, 7 Aug 2026 10:37:57 -0700 Subject: [PATCH 2/7] fix(frontend): align SGLang logprobs after EOS trimming Signed-off-by: jain-ria --- components/src/dynamo/frontend/sglang_prepost.py | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/components/src/dynamo/frontend/sglang_prepost.py b/components/src/dynamo/frontend/sglang_prepost.py index 5242df90a905..0c396b6e2dde 100644 --- a/components/src/dynamo/frontend/sglang_prepost.py +++ b/components/src/dynamo/frontend/sglang_prepost.py @@ -1146,7 +1146,13 @@ def process_output(self, engine_response: dict[str, Any]) -> dict[str, Any] | No log_probs = engine_response.get("log_probs") top_logprobs = engine_response.get("top_logprobs") if finish_reason is not None: + raw_token_count = len(token_ids) token_ids = self._strip_trailing_eos_token_ids(list(token_ids)) + retained_token_count = len(token_ids) + if log_probs is not None and len(log_probs) == raw_token_count: + log_probs = log_probs[:retained_token_count] + if top_logprobs is not None and len(top_logprobs) == raw_token_count: + top_logprobs = top_logprobs[:retained_token_count] delta_text = ( self._incremental_decode(token_ids, flush=finish_reason is not None) From 8ff00d3c5e8f092effb834fbed1e7f7f3edb459e Mon Sep 17 00:00:00 2001 From: jain-ria Date: Fri, 7 Aug 2026 11:32:50 -0700 Subject: [PATCH 3/7] fix(frontend): retain SGLang logprobs for buffered chunks Signed-off-by: jain-ria --- components/src/dynamo/frontend/sglang_prepost.py | 16 +++++++++++++--- 1 file changed, 13 insertions(+), 3 deletions(-) diff --git a/components/src/dynamo/frontend/sglang_prepost.py b/components/src/dynamo/frontend/sglang_prepost.py index 0c396b6e2dde..ac031dc6e7e4 100644 --- a/components/src/dynamo/frontend/sglang_prepost.py +++ b/components/src/dynamo/frontend/sglang_prepost.py @@ -979,6 +979,7 @@ def __init__( # incomplete byte-fallback sequence. self._decode_context_ids = list((prompt_token_ids or [])[-5:]) self._pending_decode_ids: list[int] = [] + self._pending_logprobs_content: list[dict[str, Any]] = [] self._has_emitted_role: bool = False # Tool call accumulation. SGLang's streaming parser returns # deltas (name in one chunk, argument fragments across subsequent @@ -1081,6 +1082,13 @@ def _build_openai_logprobs( return {"content": content, "refusal": None} if content else None + def _take_pending_logprobs(self) -> dict[str, Any] | None: + if not self._pending_logprobs_content: + return None + content = self._pending_logprobs_content + self._pending_logprobs_content = [] + return {"content": content, "refusal": None} + def _parse_reasoning_delta( self, delta_text: str, finish_reason: str | None ) -> tuple[str | None, str]: @@ -1164,6 +1172,8 @@ def process_output(self, engine_response: dict[str, Any]) -> dict[str, Any] | No openai_logprobs = self._build_openai_logprobs( log_probs, top_logprobs, token_ids ) + if openai_logprobs is not None: + self._pending_logprobs_content.extend(openai_logprobs["content"]) if self._fast_plain_text: if delta_text: @@ -1171,14 +1181,14 @@ def process_output(self, engine_response: dict[str, Any]) -> dict[str, Any] | No "index": 0, "delta": self._with_initial_role({"content": delta_text}), "finish_reason": finish_reason, - "logprobs": openai_logprobs, + "logprobs": self._take_pending_logprobs(), } elif finish_reason: return { "index": 0, "delta": self._with_initial_role({}), "finish_reason": finish_reason, - "logprobs": openai_logprobs, + "logprobs": self._take_pending_logprobs(), } return None @@ -1395,7 +1405,7 @@ def process_output(self, engine_response: dict[str, Any]) -> dict[str, Any] | No "index": 0, "delta": self._with_initial_role(delta), "finish_reason": effective_finish, - "logprobs": openai_logprobs, + "logprobs": self._take_pending_logprobs(), } return None From 5d6e8b96b2ce7deba9492abaf653345c024869a9 Mon Sep 17 00:00:00 2001 From: jain-ria Date: Sun, 9 Aug 2026 17:00:16 -0700 Subject: [PATCH 4/7] fix(frontend): decode split-character logprob tokens with context Signed-off-by: jain-ria --- .../src/dynamo/frontend/sglang_prepost.py | 70 +++++++++++++++++-- .../tests/test_sglang_processor_unit.py | 67 ++++++++++++++++++ 2 files changed, 133 insertions(+), 4 deletions(-) diff --git a/components/src/dynamo/frontend/sglang_prepost.py b/components/src/dynamo/frontend/sglang_prepost.py index ac031dc6e7e4..3897e93796c9 100644 --- a/components/src/dynamo/frontend/sglang_prepost.py +++ b/components/src/dynamo/frontend/sglang_prepost.py @@ -979,6 +979,7 @@ def __init__( # incomplete byte-fallback sequence. self._decode_context_ids = list((prompt_token_ids or [])[-5:]) self._pending_decode_ids: list[int] = [] + self._logprob_context_ids: list[int] = [] self._pending_logprobs_content: list[dict[str, Any]] = [] self._has_emitted_role: bool = False # Tool call accumulation. SGLang's streaming parser returns @@ -1056,14 +1057,23 @@ def _build_openai_logprobs( content: list[dict[str, Any]] = [] for index, (token_id, logprob) in enumerate(zip(token_ids, log_probs)): - token = self.tokenizer.decode([token_id], skip_special_tokens=False) + context_token_ids = (self._logprob_context_ids + token_ids[:index])[-4:] + token = self._decode_logprob_token(token_id, None, context_token_ids) candidates = top_logprobs[index] if top_logprobs else [] openai_top_logprobs = [] for candidate in candidates: - candidate_token = candidate.get("token") or "" + candidate_token = self._decode_logprob_token( + candidate.get("token_id"), + candidate.get("token"), + context_token_ids, + ) candidate_bytes = candidate.get("bytes") - if candidate_bytes is None and candidate_token: - candidate_bytes = list(candidate_token.encode("utf-8")) + if candidate_bytes is None: + candidate_bytes = ( + list(candidate_token.encode("utf-8")) + if candidate_token + else None + ) openai_top_logprobs.append( { "token": candidate_token, @@ -1082,6 +1092,57 @@ def _build_openai_logprobs( return {"content": content, "refusal": None} if content else None + def _decode_logprob_token( + self, + token_id: int | None, + token: str | None, + context_token_ids: list[int], + ) -> str: + if token is None: + if token_id is None: + return "" + token = self.tokenizer.decode([token_id], skip_special_tokens=False) + + if not token.endswith("\ufffd") or token_id is None: + return token + + for context_size in range(1, min(len(context_token_ids), 4) + 1): + context = context_token_ids[-context_size:] + decoded = self.tokenizer.decode( + context + [token_id], skip_special_tokens=False + ) + if decoded.endswith("\ufffd"): + continue + + clean_end = len(context) + for context_index in range(len(context) - 1, -1, -1): + context_token = self.tokenizer.decode( + [context[context_index]], skip_special_tokens=False + ) + if context_token.endswith("\ufffd"): + clean_end = context_index + else: + break + + clean_prefix = ( + self.tokenizer.decode( + context[:clean_end], skip_special_tokens=False + ) + if clean_end + else "" + ) + if decoded.startswith(clean_prefix): + return decoded[len(clean_prefix) :] + + common_prefix_length = 0 + for prefix_char, decoded_char in zip(clean_prefix, decoded): + if prefix_char != decoded_char: + break + common_prefix_length += 1 + return decoded[common_prefix_length:] + + return "" + def _take_pending_logprobs(self) -> dict[str, Any] | None: if not self._pending_logprobs_content: return None @@ -1174,6 +1235,7 @@ def process_output(self, engine_response: dict[str, Any]) -> dict[str, Any] | No ) if openai_logprobs is not None: self._pending_logprobs_content.extend(openai_logprobs["content"]) + self._logprob_context_ids = (self._logprob_context_ids + token_ids)[-4:] if self._fast_plain_text: if delta_text: diff --git a/components/src/dynamo/frontend/tests/test_sglang_processor_unit.py b/components/src/dynamo/frontend/tests/test_sglang_processor_unit.py index 6796dc56078b..d910f28f20c3 100644 --- a/components/src/dynamo/frontend/tests/test_sglang_processor_unit.py +++ b/components/src/dynamo/frontend/tests/test_sglang_processor_unit.py @@ -2890,6 +2890,73 @@ def test_split_multibyte_character_is_not_replaced(self): assert content == "한" assert "\ufffd" not in content + def test_logprobs_reconstruct_split_multibyte_character(self): + """Logprob token strings use context to reconstruct split UTF-8.""" + post = SglangStreamingPostProcessor( + tokenizer=self.ByteTokenizer(), + tool_call_parser=None, + reasoning_parser=None, + ) + + encoded = list("한".encode("utf-8")) + choice = None + for index, token_id in enumerate(encoded): + choice = post.process_output( + { + "token_ids": [token_id], + "finish_reason": "stop" if index == len(encoded) - 1 else None, + "log_probs": [-0.1 * (index + 1)], + "top_logprobs": [ + [ + { + "token_id": token_id, + "token": "\ufffd", + "logprob": -0.1 * (index + 1), + } + ] + ], + } + ) + + assert choice is not None + assert choice["delta"]["content"] == "한" + logprob_content = choice["logprobs"]["content"] + assert [entry["token"] for entry in logprob_content] == ["", "", "한"] + assert [entry["bytes"] for entry in logprob_content] == [ + None, + None, + list("한".encode("utf-8")), + ] + assert [ + entry["top_logprobs"][0]["token"] for entry in logprob_content + ] == ["", "", "한"] + + def test_logprobs_regular_token_is_unchanged(self): + """Ordinary tokens keep their decoded text and UTF-8 bytes.""" + post = SglangStreamingPostProcessor( + tokenizer=self.ByteTokenizer(), + tool_call_parser=None, + reasoning_parser=None, + ) + + choice = post.process_output( + { + "token_ids": [ord("A")], + "finish_reason": "stop", + "log_probs": [-0.3], + } + ) + + assert choice is not None + assert choice["logprobs"]["content"] == [ + { + "token": "A", + "logprob": -0.3, + "bytes": [65], + "top_logprobs": [], + } + ] + def test_byte_fallback_sequence_longer_than_six_tokens( self, byte_fallback_tokenizer ): From 9a0223a2c7cd3c6e259fe5ce0c6321089b2fca81 Mon Sep 17 00:00:00 2001 From: jain-ria Date: Sun, 9 Aug 2026 17:11:38 -0700 Subject: [PATCH 5/7] style(frontend): format SGLang logprob changes Signed-off-by: jain-ria --- components/src/dynamo/frontend/sglang_prepost.py | 4 +--- .../dynamo/frontend/tests/test_sglang_processor_unit.py | 8 +++++--- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/components/src/dynamo/frontend/sglang_prepost.py b/components/src/dynamo/frontend/sglang_prepost.py index 3897e93796c9..458fc7bcafbe 100644 --- a/components/src/dynamo/frontend/sglang_prepost.py +++ b/components/src/dynamo/frontend/sglang_prepost.py @@ -1125,9 +1125,7 @@ def _decode_logprob_token( break clean_prefix = ( - self.tokenizer.decode( - context[:clean_end], skip_special_tokens=False - ) + self.tokenizer.decode(context[:clean_end], skip_special_tokens=False) if clean_end else "" ) diff --git a/components/src/dynamo/frontend/tests/test_sglang_processor_unit.py b/components/src/dynamo/frontend/tests/test_sglang_processor_unit.py index d910f28f20c3..e69fe9d54c22 100644 --- a/components/src/dynamo/frontend/tests/test_sglang_processor_unit.py +++ b/components/src/dynamo/frontend/tests/test_sglang_processor_unit.py @@ -2927,9 +2927,11 @@ def test_logprobs_reconstruct_split_multibyte_character(self): None, list("한".encode("utf-8")), ] - assert [ - entry["top_logprobs"][0]["token"] for entry in logprob_content - ] == ["", "", "한"] + assert [entry["top_logprobs"][0]["token"] for entry in logprob_content] == [ + "", + "", + "한", + ] def test_logprobs_regular_token_is_unchanged(self): """Ordinary tokens keep their decoded text and UTF-8 bytes.""" From 31f8c3332a7d1bde80a6ef5813559a450b6e9eb4 Mon Sep 17 00:00:00 2001 From: jain-ria Date: Sun, 9 Aug 2026 19:30:39 -0700 Subject: [PATCH 6/7] fix(frontend): preserve SGLang logprobs across stream chunks Signed-off-by: jain-ria --- .../src/dynamo/frontend/sglang_processor.py | 184 +++++++++++------- .../tests/test_sglang_processor_unit.py | 79 ++++++++ 2 files changed, 190 insertions(+), 73 deletions(-) diff --git a/components/src/dynamo/frontend/sglang_processor.py b/components/src/dynamo/frontend/sglang_processor.py index 3da4dd0bfdf8..1afba0ae0fc2 100644 --- a/components/src/dynamo/frontend/sglang_processor.py +++ b/components/src/dynamo/frontend/sglang_processor.py @@ -639,6 +639,91 @@ async def _generate_and_stream( video_count = len(_mm_counts.get("video_url", [])) audio_count = len(_mm_counts.get("audio_url", [])) + def flush_pending( + *, + finish_reason: str | None, + stop_reason: Any | None, + engine_data: Any | None, + ) -> dict[str, Any]: + nonlocal pending_token_ids + nonlocal pending_log_probs + nonlocal pending_top_logprobs + nonlocal pending_usage + nonlocal first_chunk + nonlocal post_proc_total_ms + nonlocal token_count + + chunk_token_count = len(pending_token_ids) + usage_for_metrics = pending_usage + mapped_response = { + "token_ids": pending_token_ids, + "finish_reason": finish_reason, + } + if pending_log_probs is not None: + mapped_response["log_probs"] = pending_log_probs + if pending_top_logprobs is not None: + mapped_response["top_logprobs"] = pending_top_logprobs + + if self.debug_perf: + t_pp0 = time.monotonic() + + choice = post.process_output(mapped_response) + + if self.debug_perf: + t_pp1 = time.monotonic() + post_proc_total_ms += (t_pp1 - t_pp0) * 1000.0 + token_count += chunk_token_count + + envelope: dict[str, Any] = {"_dynamo_annotated": True} + if choice: + dynamo_out: dict[str, Any] = { + "id": request_id, + "choices": [choice], + "created": created_ts, + "model": request["model"], + "object": "chat.completion.chunk", + } + if pending_usage: + dynamo_out["usage"] = pending_usage + response_nvext: dict[str, Any] = {} + if stop_reason is not None and nvext_extra_field_requested( + request, "stop_reason" + ): + response_nvext["stop_reason"] = stop_reason + if engine_data is not None and nvext_extra_field_requested( + request, "engine_data" + ): + response_nvext["engine_data"] = engine_data + if response_nvext: + dynamo_out["nvext"] = response_nvext + + envelope["data"] = dynamo_out + + metrics: dict[str, Any] = { + "input_tokens": input_tokens, + "output_tokens": cumulative_output_tokens, + "chunk_tokens": chunk_token_count, + } + # Include nonzero counts on every frame (text-only carries nothing). + if image_count: + metrics["image_count"] = image_count + if video_count: + metrics["video_count"] = video_count + if audio_count: + metrics["audio_count"] = audio_count + cached_tokens = _cached_tokens_from_usage(usage_for_metrics) + if cached_tokens is not None: + metrics["cached_tokens"] = cached_tokens + envelope["event"] = "llm_metrics" + envelope["comment"] = [json.dumps(metrics)] + + pending_token_ids = [] + pending_log_probs = None + pending_top_logprobs = None + pending_usage = None + first_chunk = False + return envelope + async for dynamo_response in dynamo_stream: if dynamo_response.is_error(): comments = dynamo_response.comments() or [] @@ -665,6 +750,25 @@ async def _generate_and_stream( break new_ids = engine_response["token_ids"] + log_probs = engine_response.get("log_probs") + top_logprobs = engine_response.get("top_logprobs") + + if new_ids and pending_token_ids: + pending_logprob_shape = ( + pending_log_probs is not None, + pending_top_logprobs is not None, + ) + chunk_logprob_shape = ( + log_probs is not None, + top_logprobs is not None, + ) + if pending_logprob_shape != chunk_logprob_shape: + yield flush_pending( + finish_reason=None, + stop_reason=None, + engine_data=None, + ) + chunk_tokens = len(new_ids) cumulative_output_tokens += chunk_tokens raw_finish = engine_response.get("finish_reason") @@ -676,11 +780,11 @@ async def _generate_and_stream( engine_data = engine_response.get("engine_data") pending_token_ids.extend(new_ids) - if (log_probs := engine_response.get("log_probs")) is not None: + if log_probs is not None: if pending_log_probs is None: pending_log_probs = [] pending_log_probs.extend(log_probs) - if (top_logprobs := engine_response.get("top_logprobs")) is not None: + if top_logprobs is not None: if pending_top_logprobs is None: pending_top_logprobs = [] pending_top_logprobs.extend(top_logprobs) @@ -689,77 +793,11 @@ async def _generate_and_stream( # First chunk flushes immediately (si=1) to minimize TTFT. flush_threshold = 1 if first_chunk else stream_interval if finish_reason or len(pending_token_ids) >= flush_threshold: - usage_for_metrics = pending_usage - mapped_response = { - "token_ids": pending_token_ids, - "finish_reason": finish_reason, - } - if pending_log_probs is not None: - mapped_response["log_probs"] = pending_log_probs - if pending_top_logprobs is not None: - mapped_response["top_logprobs"] = pending_top_logprobs - - if self.debug_perf: - t_pp0 = time.monotonic() - - choice = post.process_output(mapped_response) - - if self.debug_perf: - t_pp1 = time.monotonic() - post_proc_total_ms += (t_pp1 - t_pp0) * 1000.0 - token_count += len(pending_token_ids) - - envelope: dict[str, Any] = {"_dynamo_annotated": True} - if choice: - dynamo_out: dict[str, Any] = { - "id": request_id, - "choices": [choice], - "created": created_ts, - "model": request["model"], - "object": "chat.completion.chunk", - } - if pending_usage: - dynamo_out["usage"] = pending_usage - pending_usage = None - response_nvext: dict[str, Any] = {} - if stop_reason is not None and nvext_extra_field_requested( - request, "stop_reason" - ): - response_nvext["stop_reason"] = stop_reason - if engine_data is not None and ( - nvext_extra_field_requested(request, "engine_data") - ): - response_nvext["engine_data"] = engine_data - if response_nvext: - dynamo_out["nvext"] = response_nvext - - envelope["data"] = dynamo_out - - metrics: dict[str, Any] = { - "input_tokens": input_tokens, - "output_tokens": cumulative_output_tokens, - "chunk_tokens": len(pending_token_ids), - } - # Include nonzero counts on every frame (text-only carries nothing). - if image_count: - metrics["image_count"] = image_count - if video_count: - metrics["video_count"] = video_count - if audio_count: - metrics["audio_count"] = audio_count - cached_tokens = _cached_tokens_from_usage(usage_for_metrics) - if cached_tokens is not None: - metrics["cached_tokens"] = cached_tokens - envelope["event"] = "llm_metrics" - envelope["comment"] = [json.dumps(metrics)] - - yield envelope - - pending_token_ids = [] - pending_log_probs = None - pending_top_logprobs = None - pending_usage = None - first_chunk = False + yield flush_pending( + finish_reason=finish_reason, + stop_reason=stop_reason, + engine_data=engine_data, + ) except Unknown: raise except Exception as e: diff --git a/components/src/dynamo/frontend/tests/test_sglang_processor_unit.py b/components/src/dynamo/frontend/tests/test_sglang_processor_unit.py index e69fe9d54c22..85c38846330a 100644 --- a/components/src/dynamo/frontend/tests/test_sglang_processor_unit.py +++ b/components/src/dynamo/frontend/tests/test_sglang_processor_unit.py @@ -2959,6 +2959,85 @@ def test_logprobs_regular_token_is_unchanged(self): } ] + def _run_logprob_stream(self, items): + processor = SglangProcessor( + tokenizer=self.ByteTokenizer(), + routed_engine=FakeRoutedEngine(items=items), + tool_call_parser_name=None, + reasoning_parser_name=None, + eos_token_ids=None, + stream_interval=20, + ) + post = SglangStreamingPostProcessor( + tokenizer=self.ByteTokenizer(), + tool_call_parser=None, + reasoning_parser=None, + ) + + async def collect(): + return [ + item["data"] + async for item in processor._generate_and_stream( + "req-logprobs", {"model": "test-model"}, {}, [], post + ) + if "data" in item + ] + + return asyncio.run(collect()) + + def test_missing_chunk_logprobs_do_not_drop_adjacent_logprobs(self): + """A missing-logprob chunk is isolated from adjacent valid chunks.""" + chunks = self._run_logprob_stream( + [ + {"token_ids": [ord("A")], "log_probs": [-0.1]}, + {"token_ids": [ord("B")], "log_probs": [-0.2]}, + {"token_ids": [ord("C")]}, + { + "token_ids": [ord("D")], + "log_probs": [-0.4], + "finish_reason": "stop", + }, + ] + ) + + choices = [chunk["choices"][0] for chunk in chunks] + assert [choice["delta"]["content"] for choice in choices] == [ + "A", + "B", + "C", + "D", + ] + assert [ + choice["logprobs"]["content"][0]["logprob"] + if choice["logprobs"] is not None + else None + for choice in choices + ] == [-0.1, -0.2, None, -0.4] + assert choices[-1]["finish_reason"] == "stop" + + def test_consistent_chunk_logprobs_keep_normal_batching(self): + """Consistent logprob chunks retain the configured stream interval.""" + chunks = self._run_logprob_stream( + [ + {"token_ids": [ord("A")], "log_probs": [-0.1]}, + {"token_ids": [ord("B")], "log_probs": [-0.2]}, + {"token_ids": [ord("C")], "log_probs": [-0.3]}, + { + "token_ids": [ord("D")], + "log_probs": [-0.4], + "finish_reason": "stop", + }, + ] + ) + + choices = [chunk["choices"][0] for chunk in chunks] + assert [choice["delta"]["content"] for choice in choices] == ["A", "BCD"] + assert [entry["logprob"] for entry in choices[1]["logprobs"]["content"]] == [ + -0.2, + -0.3, + -0.4, + ] + def test_byte_fallback_sequence_longer_than_six_tokens( self, byte_fallback_tokenizer ): From ec3b64bdcb65b6fda2f6dafddc77d6cf8c83449d Mon Sep 17 00:00:00 2001 From: jain-ria Date: Mon, 10 Aug 2026 11:06:34 -0700 Subject: [PATCH 7/7] fix(frontend): type SGLang mapped response Signed-off-by: jain-ria --- components/src/dynamo/frontend/sglang_processor.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/src/dynamo/frontend/sglang_processor.py b/components/src/dynamo/frontend/sglang_processor.py index 1afba0ae0fc2..87965d95994d 100644 --- a/components/src/dynamo/frontend/sglang_processor.py +++ b/components/src/dynamo/frontend/sglang_processor.py @@ -655,7 +655,7 @@ def flush_pending( chunk_token_count = len(pending_token_ids) usage_for_metrics = pending_usage - mapped_response = { + mapped_response: dict[str, Any] = { "token_ids": pending_token_ids, "finish_reason": finish_reason, }