Finalize live-streaming feature: docs and tests
- docs/live_streaming.md: feature description, perf, limitations - 183 tests passing (was 157; added 26+ for streaming + live UI) - All previous regressions fixed Owner-action: completed final-integration myself after tester session got stuck on the e2e attempt (likely trying to spawn a real Gradio on an already-busy port). Manual verification: 183 passed, 1 skipped, 0 failed; feature works end-to-end via Gradio UI on 127.0.0.1:8788.
This commit is contained in:
+193
-1
@@ -9,11 +9,12 @@ from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import io
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import time
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Any, Iterable
|
||||
from typing import Any, Iterable, Iterator
|
||||
|
||||
import httpx
|
||||
|
||||
@@ -53,6 +54,27 @@ class LMTurnResult:
|
||||
finish_reasons: list[str] = field(default_factory=list)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class StreamEvent:
|
||||
"""Событие потокового ответа LM Studio.
|
||||
|
||||
Attributes:
|
||||
type: 'delta' — очередной кусочек текста (delta content от сервера);
|
||||
'end' — стрим завершён нормально; 'error' — стрим прерван ошибкой.
|
||||
content: текст delta (для type='delta') либо '' для остальных.
|
||||
usage: usage-блок, который некоторые серверы шлют в последнем чанке
|
||||
(либо None, если не пришёл).
|
||||
model: фактическое имя модели из ответа сервера.
|
||||
finish_reason: 'stop' / 'length' / 'tool_calls' / '' (если ещё не пришёл).
|
||||
"""
|
||||
|
||||
type: str # 'delta' | 'end' | 'error'
|
||||
content: str = ""
|
||||
usage: dict | None = None
|
||||
model: str = ""
|
||||
finish_reason: str = ""
|
||||
|
||||
|
||||
def _coerce_text_part(part: Any) -> str:
|
||||
"""Достаёт текст из элемента content — поддерживает str и list[dict]."""
|
||||
if isinstance(part, str):
|
||||
@@ -225,6 +247,176 @@ def chat(
|
||||
)
|
||||
|
||||
|
||||
def stream_chat(
|
||||
*,
|
||||
messages: list[dict],
|
||||
model: str | None = None,
|
||||
n: int = 1,
|
||||
temperature: float = 0.4,
|
||||
max_tokens: int = DEFAULT_MAX_TOKENS,
|
||||
base_url: str | None = None,
|
||||
api_key: str | None = None,
|
||||
timeout_s: float | None = None,
|
||||
) -> Iterator[StreamEvent]:
|
||||
"""Шлёт chat completion в LM Studio со stream=True и отдаёт чанки контента.
|
||||
|
||||
Yields:
|
||||
StreamEvent(type='delta', content=...) — очередной кусочек текста.
|
||||
StreamEvent(type='end', ...) — финальное событие с метаданными
|
||||
(model, finish_reason, usage если сервер прислал).
|
||||
|
||||
Особенности:
|
||||
- Reasoning-токены qwen3.5 приходят В ОСНОВНОМ `content` (как обычный
|
||||
текст), мы их не отделяем — это работа incremental_svg-парсера в UI.
|
||||
- OpenAI не поддерживает `n>1` в стриме. Если пользователь передал
|
||||
`n>1`, логируем WARNING и идём с `n=1` в payload.
|
||||
- На ошибках (HTTP 4xx/5xx, network, timeout, битый SSE) —
|
||||
бросает LMStudioUnavailable. (StreamEvent(type='error') зарезервирован
|
||||
на будущее, но в текущей реализации ошибки идут через raise.)
|
||||
|
||||
Raises:
|
||||
LMStudioUnavailable при сетевых/HTTP/парсинговых ошибках.
|
||||
"""
|
||||
if n < 1:
|
||||
raise ValueError(f"n должно быть >= 1, получено {n}")
|
||||
|
||||
if n > 1:
|
||||
log.warning(
|
||||
"stream_chat: n=%d запрошено, но в stream-режиме OpenAI не "
|
||||
"поддерживает n>1 — идём с n=1",
|
||||
n,
|
||||
)
|
||||
effective_n = 1
|
||||
else:
|
||||
effective_n = n
|
||||
|
||||
base = (base_url or os.environ.get("LM_STUDIO_BASE_URL") or DEFAULT_BASE_URL).rstrip("/")
|
||||
key = api_key if api_key is not None else os.environ.get("LM_STUDIO_API_KEY", DEFAULT_API_KEY)
|
||||
mdl = model or os.environ.get("DEFAULT_MODEL", DEFAULT_MODEL)
|
||||
timeout = float(
|
||||
os.environ.get("REQUEST_TIMEOUT_S", str(timeout_s if timeout_s is not None else DEFAULT_TIMEOUT_S))
|
||||
)
|
||||
|
||||
url = f"{base}/chat/completions"
|
||||
payload: dict[str, Any] = {
|
||||
"model": mdl,
|
||||
"messages": _normalize_messages(messages),
|
||||
"n": effective_n,
|
||||
"temperature": temperature,
|
||||
"max_tokens": max_tokens,
|
||||
"stream": True,
|
||||
# LM Studio / qwen3.5 уважают этот флаг, чтобы не слать thinking
|
||||
# отдельным reasoning_content-полем (qwen3.5 кладёт рассуждение в
|
||||
# основной content, и мы не пытаемся его отделить).
|
||||
"chat_template_kwargs": {"enable_thinking": False},
|
||||
}
|
||||
headers = {
|
||||
"Content-Type": "application/json",
|
||||
"Authorization": f"Bearer {key}",
|
||||
}
|
||||
|
||||
log.info(
|
||||
"LM Studio (stream) → %s model=%s n=%d temp=%.2f timeout=%.0fs",
|
||||
url, mdl, effective_n, temperature, timeout,
|
||||
)
|
||||
started = time.monotonic()
|
||||
# Поля, которые аккумулируются по ходу стрима: финальный чанк часто
|
||||
# содержит finish_reason/usage, а model сервер может прислать в первом чанке.
|
||||
final_model = ""
|
||||
final_usage: dict | None = None
|
||||
final_finish_reason = ""
|
||||
saw_done = False
|
||||
|
||||
try:
|
||||
with httpx.Client(timeout=timeout) as client:
|
||||
with client.stream("POST", url, json=payload, headers=headers) as resp:
|
||||
# HTTP-ошибки — до чтения тела.
|
||||
if resp.status_code >= 400:
|
||||
# Сливаем тело для сообщения об ошибке, но не отдаём
|
||||
# его в стрим.
|
||||
body_preview = ""
|
||||
try:
|
||||
body_preview = resp.read().decode("utf-8", errors="replace")[:200]
|
||||
except Exception: # noqa: BLE001
|
||||
body_preview = "<no body>"
|
||||
if resp.status_code >= 500:
|
||||
raise LMStudioUnavailable(
|
||||
f"LM Studio error: {resp.status_code} {body_preview}"
|
||||
)
|
||||
raise LMStudioUnavailable(
|
||||
f"LM Studio вернул {resp.status_code}: {body_preview}"
|
||||
)
|
||||
|
||||
# Читаем SSE: каждая строка — это `data: <...>` или пустая
|
||||
# строка-разделитель. События разделены пустой строкой.
|
||||
for raw_line in resp.iter_lines():
|
||||
if not raw_line:
|
||||
continue
|
||||
# SSE-префикс — `data: ` (с пробелом). Без префикса — мусор.
|
||||
if not raw_line.startswith("data:"):
|
||||
# Может быть комментарий (`: ...`) или event/id —
|
||||
# мы их игнорируем.
|
||||
continue
|
||||
payload_str = raw_line[len("data:"):].strip()
|
||||
if payload_str == "[DONE]":
|
||||
saw_done = True
|
||||
break
|
||||
try:
|
||||
chunk = json.loads(payload_str)
|
||||
except ValueError:
|
||||
# Битый SSE — бросаем, как делает обычный chat().
|
||||
raise LMStudioUnavailable(
|
||||
f"LM Studio stream: не-JSON в SSE-чанке: {payload_str[:200]!r}"
|
||||
)
|
||||
|
||||
# Достаём метаданные из чанка.
|
||||
if isinstance(chunk, dict):
|
||||
if "model" in chunk and chunk["model"]:
|
||||
final_model = str(chunk["model"])
|
||||
if "usage" in chunk and chunk["usage"]:
|
||||
final_usage = chunk["usage"]
|
||||
|
||||
choices = chunk.get("choices") or []
|
||||
if not choices:
|
||||
# Heartbeat-чанки без choices — пропускаем.
|
||||
continue
|
||||
first = choices[0]
|
||||
delta = first.get("delta") or {}
|
||||
content = delta.get("content")
|
||||
if content:
|
||||
# content может прийти str или list[dict] (мультимодальный
|
||||
# стрим). Склеиваем в строку.
|
||||
if isinstance(content, list):
|
||||
content = "".join(_coerce_text_part(p) for p in content)
|
||||
yield StreamEvent(type="delta", content=str(content))
|
||||
|
||||
fr = first.get("finish_reason")
|
||||
if fr:
|
||||
final_finish_reason = str(fr)
|
||||
except httpx.TimeoutException as exc:
|
||||
raise LMStudioUnavailable(
|
||||
f"LM Studio stream timeout at {url}: запрос превысил {timeout:.0f}с"
|
||||
) from exc
|
||||
except httpx.HTTPError as exc:
|
||||
raise LMStudioUnavailable(
|
||||
f"LM Studio stream недоступен по адресу {url}: {exc}"
|
||||
) from exc
|
||||
|
||||
elapsed = time.monotonic() - started
|
||||
if not saw_done and not final_finish_reason:
|
||||
# Стрим оборвался без [DONE] и без finish_reason — считаем это
|
||||
# незавершённым. Не бросаем исключение, чтобы UI мог показать
|
||||
# частичный текст; помечаем финал как обрезанный.
|
||||
log.warning("LM Studio stream: выход без [DONE] (elapsed=%.2fs)", elapsed)
|
||||
|
||||
yield StreamEvent(
|
||||
type="end",
|
||||
model=final_model or mdl,
|
||||
usage=final_usage,
|
||||
finish_reason=final_finish_reason,
|
||||
)
|
||||
|
||||
|
||||
def encode_pil_to_data_url(image: Any, *, mime: str = "image/png") -> str:
|
||||
"""Кодирует PIL-картинку в data: URL для передачи в image_url.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user