From 710fc8b90b78f3b9536de87852c120f9e555405e Mon Sep 17 00:00:00 2001 From: Vladimir Bryzgalov Date: Sat, 15 Aug 2026 23:44:09 +0500 Subject: [PATCH] =?UTF-8?q?=D0=9E=D1=86=D0=B5=D0=BD=D0=BA=D0=B0=20=D1=80?= =?UTF-8?q?=D0=B0=D0=B7=D0=B4=D0=B5=D0=BB=D0=B5=D0=BD=D0=B8=D1=8F=20=D0=B1?= =?UTF-8?q?=D0=BE=D0=BB=D1=8C=D1=88=D0=B5=20=D0=BD=D0=B5=20=D1=80=D0=BE?= =?UTF-8?q?=D0=BD=D1=8F=D0=B5=D1=82=20=D0=B7=D0=B0=D0=B4=D0=B0=D1=87=D1=83?= =?UTF-8?q?,=20=D0=B6=D1=83=D1=80=D0=BD=D0=B0=D0=BB=20=D0=B4=D0=BE=D1=81?= =?UTF-8?q?=D1=82=D1=83=D0=BF=D0=B5=D0=BD=20=D1=87=D0=B5=D1=80=D0=B5=D0=B7?= =?UTF-8?q?=20API?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Полный трейс показал настоящую причину сбоя: падало не распознавание, а метрика separation_quality из 0.5.0 - она скармливала модели отпечатков реплики целиком, и на 190-секундной та не выдержала. Теперь берётся кусок из середины реплики, а сбой оценки лишь обнуляет её, не трогая расшифровку. Добавлен /v1/logs: последние записи журнала с фильтром по уровню, чтобы разбирать сбои без копирования консоли вручную. Co-Authored-By: Claude Opus 5 (1M context) --- app/logbuffer.py | 58 +++++++++++++++++++++++++++++++++++++++ app/main.py | 19 +++++++++++++ app/pipeline.py | 28 +++++++++++++++---- app/version.py | 2 +- tests/test_api.py | 38 +++++++++++++++++++++++++ tests/test_concurrency.py | 36 ++++++++++++++++++++++++ 6 files changed, 174 insertions(+), 7 deletions(-) create mode 100644 app/logbuffer.py diff --git a/app/logbuffer.py b/app/logbuffer.py new file mode 100644 index 0000000..139b6ed --- /dev/null +++ b/app/logbuffer.py @@ -0,0 +1,58 @@ +"""Кольцевой буфер последних записей журнала. + +Нужен, чтобы смотреть логи через API, а не копировать их из окна консоли. +Хранится в памяти: файл на диске пришлось бы чистить, а история глубже +последних сотен строк для разбора сбоя не нужна. +""" +import logging +import threading +from collections import deque + +__all__ = ["LogBuffer", "install"] + + +class LogBuffer(logging.Handler): + def __init__(self, capacity: int = 500): + super().__init__() + self._records: deque = deque(maxlen=capacity) + self._lock = threading.Lock() + + def emit(self, record: logging.LogRecord) -> None: + try: + message = record.getMessage() + if record.exc_info: + message += "\n" + self.format(record).split("\n", 1)[-1] + except Exception: # noqa: BLE001 - журнал не должен падать сам + message = "не удалось разобрать запись журнала" + with self._lock: + self._records.append({ + "time": record.created, + "level": record.levelname, + "logger": record.name, + "message": message, + }) + + def tail(self, limit: int = 100, level: str | None = None) -> list[dict]: + with self._lock: + rows = list(self._records) + if level: + wanted = level.upper() + rows = [r for r in rows if r["level"] == wanted] + return rows[-limit:] + + +_buffer: LogBuffer | None = None + + +def install(capacity: int = 500) -> LogBuffer: + """Подключает буфер к корневому журналу. Повторный вызов вернёт тот же буфер.""" + global _buffer + if _buffer is None: + _buffer = LogBuffer(capacity) + _buffer.setFormatter(logging.Formatter("%(message)s")) + logging.getLogger().addHandler(_buffer) + return _buffer + + +def get() -> LogBuffer: + return install() diff --git a/app/main.py b/app/main.py index 6831428..b2f8dcf 100644 --- a/app/main.py +++ b/app/main.py @@ -16,6 +16,8 @@ from fastapi.responses import JSONResponse from fastapi.security import HTTPBearer from app.config import Settings, load_settings +from app.logbuffer import get as log_buffer +from app.logbuffer import install as install_log_buffer from app.pipeline import ModelsMissing, Pipeline, to_wav16k from app.security import check_token, ip_allowed, parse_allowlist from app.store import JobStatus, JobStore @@ -211,6 +213,22 @@ def openapi_schema(request: Request) -> JSONResponse: return JSONResponse(schema) +@app.get("/v1/logs", dependencies=[Depends(guard)], summary="Журнал сервиса", + description="Последние записи журнала: удобно посмотреть причину сбоя, " + "не заходя на машину. Уровень можно отфильтровать параметром level.") +def logs(limit: int = Query(100, ge=1, le=500), + level: str | None = Query(None, description="INFO, WARNING или ERROR")) -> dict: + rows = log_buffer().tail(limit=limit, level=level) + return { + "count": len(rows), + "records": [ + {"time": time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(r["time"])), + "level": r["level"], "logger": r["logger"], "message": r["message"]} + for r in rows + ], + } + + @app.get("/health", summary="Состояние сервиса", description="Единственный метод без токена - годится для мониторинга. " "Показывает версию, очередь, число потоков и то, каким сервис " @@ -342,6 +360,7 @@ def run() -> None: import uvicorn _setup_console() + install_log_buffer() logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s: %(message)s") diff --git a/app/pipeline.py b/app/pipeline.py index 89ec0bc..2f7f8cd 100644 --- a/app/pipeline.py +++ b/app/pipeline.py @@ -26,6 +26,9 @@ SAMPLE_RATE = 16000 MAX_CHUNK_SEC = 60.0 # Ниже этого делить бессмысленно: явно не длина виновата. MIN_SPLIT_SEC = 2.0 +# Модель отпечатков рассчитана на короткий фрагмент и на длинном падает. +# Для оценки голоса секунд более чем достаточно. +EMBED_SEC = 8.0 # Ниже этого значения голоса практически неразличимы и разметка по говорящим # случайна. На записях с одним микрофоном в комнате так бывает часто. RELIABLE_SEPARATION = 0.35 @@ -169,15 +172,23 @@ class Pipeline: self._embedder = sherpa_onnx.SpeakerEmbeddingExtractor( sherpa_onnx.SpeakerEmbeddingExtractorConfig( model=str(self.models_dir / EMB_MODEL_REL), num_threads=self.threads)) - vectors = [] + vectors, labels = [], [] for seg in usable: + # Берём кусок из середины реплики: там речь устойчивее, чем на краях, + # а длинный фрагмент модель отпечатков просто не переваривает. + middle = (seg.start + seg.end) / 2 + half = min(EMBED_SEC, seg.end - seg.start) / 2 + piece = samples[int((middle - half) * SAMPLE_RATE):int((middle + half) * SAMPLE_RATE)] + if len(piece) < SAMPLE_RATE // 2: + continue stream = self._embedder.create_stream() - stream.accept_waveform(SAMPLE_RATE, - samples[int(seg.start * SAMPLE_RATE):int(seg.end * SAMPLE_RATE)]) + stream.accept_waveform(SAMPLE_RATE, piece) stream.input_finished() vectors.append(np.array(self._embedder.compute(stream))) - return separation_quality(np.array(vectors), - np.array([s.speaker for s in usable])) + labels.append(seg.speaker) + if len(vectors) < 4: + return 0.0 + return separation_quality(np.array(vectors), np.array(labels)) def _recognize_safely(self, audio: np.ndarray, depth: int = 0) -> str: """Распознаёт кусок, при ошибке деля его пополам. @@ -214,7 +225,12 @@ class Pipeline: t0 = time.time() raw = self._diarizer(num_speakers).process(samples).sort_by_start_time() segments = [Segment(start=s.start, end=s.end, speaker=s.speaker) for s in raw] - quality = self._separation_quality(samples, segments) + try: + quality = self._separation_quality(samples, segments) + except Exception as exc: # noqa: BLE001 - оценка вспомогательная + # Метрика не должна ронять задачу: без неё расшифровка всё равно нужна. + log.warning("не удалось оценить разделение говорящих: %s", exc) + quality = 0.0 t_diar = time.time() - t0 t0 = time.time() diff --git a/app/version.py b/app/version.py index 43c4ab0..22049ab 100644 --- a/app/version.py +++ b/app/version.py @@ -1 +1 @@ -__version__ = "0.6.1" +__version__ = "0.6.2" diff --git a/tests/test_api.py b/tests/test_api.py index ee5044a..37b306d 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -217,3 +217,41 @@ class TestDocsAccessControl: def test_health_stays_open_for_foreign_ip(self, restricted_client): """Мониторинг должен работать всегда.""" assert restricted_client.get("/health").status_code == 200 + + +class TestLogs: + def test_logs_require_token(self, client): + client.headers.pop("Authorization") + assert client.get("/v1/logs").status_code == 401 + + def test_logs_return_records(self, client): + import logging + + from app.logbuffer import install + + install() + logging.getLogger("talkscore-asr").error("тестовая запись") + body = client.get("/v1/logs").json() + assert any("тестовая запись" in r["message"] for r in body["records"]) + + def test_logs_filter_by_level(self, client): + import logging + + from app.logbuffer import install + + install() + log = logging.getLogger("talkscore-asr") + log.info("обычная строка") + log.error("строка об ошибке") + body = client.get("/v1/logs?level=ERROR").json() + assert all(r["level"] == "ERROR" for r in body["records"]) + + def test_buffer_keeps_only_recent(self): + from app.logbuffer import LogBuffer + + buf = LogBuffer(capacity=10) + import logging + for i in range(50): + buf.emit(logging.LogRecord("t", logging.INFO, "f", 1, f"строка {i}", None, None)) + rows = buf.tail(limit=100) + assert len(rows) == 10 and "строка 49" in rows[-1]["message"] diff --git a/tests/test_concurrency.py b/tests/test_concurrency.py index ca12c6c..884e194 100644 --- a/tests/test_concurrency.py +++ b/tests/test_concurrency.py @@ -281,3 +281,39 @@ class TestSpawnSafety: assert main.store is None assert not (tmp_path / "data" / "jobs.db").exists() + + +class TestSeparationNeverBreaksJob: + """Оценка разделения вспомогательная и не должна ронять расшифровку.""" + + def test_embedding_uses_short_piece(self): + """Длинный фрагмент модель отпечатков не переваривает.""" + from app.pipeline import EMBED_SEC + + assert EMBED_SEC <= 10 + + def test_failure_in_quality_does_not_break_transcribe(self, tmp_path, monkeypatch): + import sys + import types + + import numpy as np + + for name in ("sherpa_onnx", "onnx_asr", "onnxruntime"): + monkeypatch.setitem(sys.modules, name, types.ModuleType(name)) + from app.pipeline import Pipeline + + p = Pipeline(models_dir=tmp_path, threads=1, + replacements_path=tmp_path / "r.txt", base_dir=tmp_path) + p._asr = types.SimpleNamespace(recognize=lambda *a, **k: "текст") + monkeypatch.setattr(p, "warmup", lambda: None) + monkeypatch.setattr(p, "_reload_replacements", lambda: None) + monkeypatch.setattr(p, "_diarizer", lambda n: types.SimpleNamespace( + process=lambda s: types.SimpleNamespace( + sort_by_start_time=lambda: [types.SimpleNamespace(start=0.0, end=2.0, speaker=0)]))) + monkeypatch.setattr(p, "_separation_quality", + lambda *a: (_ for _ in ()).throw(RuntimeError("модель упала"))) + monkeypatch.setattr("app.pipeline.read_wav", lambda _: np.zeros(16000 * 3, dtype=np.float32)) + + result = p.transcribe(tmp_path / "any.wav") + assert result["stats"]["separation_quality"] == 0.0 + assert result["turns"] # расшифровка на месте