talkscore-asr 0.1.0: сервис транскрибации и диаризации
Локальный FastAPI-сервис поверх GigaAM v3 и sherpa-onnx: приём аудио, очередь задач, разделение по говорящим, постобработка терминов. Доставка на Windows - ZIP со встроенным Python, без установки чего-либо. Обновление кода при запуске тянется из релизов Gitea: меняется только папка app, десятки килобайт вместо всего пакета. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
+219
@@ -0,0 +1,219 @@
|
||||
"""HTTP-сервис распознавания: приём файла, очередь, выдача результата."""
|
||||
import logging
|
||||
import shutil
|
||||
import sys
|
||||
import tempfile
|
||||
import threading
|
||||
import time
|
||||
from contextlib import asynccontextmanager
|
||||
from pathlib import Path
|
||||
|
||||
from fastapi import Depends, FastAPI, File, HTTPException, Query, Request, UploadFile
|
||||
from fastapi.responses import JSONResponse
|
||||
|
||||
from app.config import Settings, load_settings
|
||||
from app.pipeline import ModelsMissing, Pipeline, to_wav16k
|
||||
from app.security import check_token, ip_allowed, parse_allowlist
|
||||
from app.store import JobStatus, JobStore
|
||||
|
||||
log = logging.getLogger("talkscore-asr")
|
||||
|
||||
settings: Settings = load_settings()
|
||||
allowlist = parse_allowlist(settings.allow_ips)
|
||||
store = JobStore(settings.data_dir / "jobs.db")
|
||||
pipeline = Pipeline(
|
||||
models_dir=settings.models_dir,
|
||||
threads=settings.effective_threads(),
|
||||
replacements_path=settings.replacements_path,
|
||||
base_dir=settings.base_dir,
|
||||
)
|
||||
|
||||
_worker_stop = threading.Event()
|
||||
_state: dict = {"ready": False, "error": None}
|
||||
|
||||
|
||||
def _process(job_id: str) -> None:
|
||||
job = store.get(job_id)
|
||||
if job is None:
|
||||
return
|
||||
store.mark_running(job_id)
|
||||
upload = settings.data_dir / "uploads" / job_id
|
||||
try:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
wav = Path(tmp) / "audio.wav"
|
||||
to_wav16k(upload, wav, pipeline.ffmpeg)
|
||||
speakers = int(job["options"].get("speakers", settings.speakers))
|
||||
result = pipeline.transcribe(wav, num_speakers=speakers)
|
||||
result["filename"] = job["filename"]
|
||||
store.mark_done(job_id, result)
|
||||
log.info("задача %s готова: %.1f с аудио, x%s", job_id,
|
||||
result["duration_sec"], result["timing"]["realtime_factor"])
|
||||
except Exception as exc: # noqa: BLE001 - в статус задачи должна попасть любая причина
|
||||
log.exception("задача %s провалилась", job_id)
|
||||
store.mark_failed(job_id, error=f"{type(exc).__name__}: {exc}")
|
||||
finally:
|
||||
upload.unlink(missing_ok=True)
|
||||
|
||||
|
||||
def _worker_loop() -> None:
|
||||
"""Один воркер: модели тяжёлые, параллельные задачи только мешали бы друг другу."""
|
||||
last_cleanup = 0.0
|
||||
while not _worker_stop.is_set():
|
||||
if _state["ready"]:
|
||||
job_id = store.take_next()
|
||||
if job_id:
|
||||
_process(job_id)
|
||||
continue
|
||||
if time.time() - last_cleanup > 3600:
|
||||
removed = store.cleanup(settings.keep_results_hours)
|
||||
if removed:
|
||||
log.info("удалено старых задач: %d", removed)
|
||||
last_cleanup = time.time()
|
||||
_worker_stop.wait(1.0)
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
async def lifespan(app: FastAPI):
|
||||
(settings.data_dir / "uploads").mkdir(parents=True, exist_ok=True)
|
||||
try:
|
||||
pipeline.warmup()
|
||||
_state["ready"] = True
|
||||
except ModelsMissing as exc:
|
||||
_state["error"] = str(exc)
|
||||
log.error("сервис запущен без моделей: %s", exc)
|
||||
worker = threading.Thread(target=_worker_loop, name="asr-worker", daemon=True)
|
||||
worker.start()
|
||||
log.info("сервис слушает %s:%s, потоков %d", settings.host, settings.port,
|
||||
settings.effective_threads())
|
||||
yield
|
||||
_worker_stop.set()
|
||||
worker.join(timeout=5)
|
||||
|
||||
|
||||
# Схема API не требует токена, поэтому по умолчанию она не публикуется:
|
||||
# знать устройство сервиса посторонним незачем.
|
||||
app = FastAPI(
|
||||
title="talkscore-asr",
|
||||
version="0.1.0",
|
||||
lifespan=lifespan,
|
||||
docs_url="/docs" if settings.docs else None,
|
||||
redoc_url="/redoc" if settings.docs else None,
|
||||
openapi_url="/openapi.json" if settings.docs else None,
|
||||
)
|
||||
|
||||
|
||||
def guard(request: Request) -> None:
|
||||
"""Проверяет адрес и токен. Порядок важен: сначала сеть, потом секрет."""
|
||||
client_ip = request.client.host if request.client else None
|
||||
if not ip_allowed(client_ip, allowlist):
|
||||
log.warning("отказано по адресу: %s", client_ip)
|
||||
raise HTTPException(status_code=403, detail="адрес не в списке разрешённых")
|
||||
if not check_token(request.headers.get("authorization"), settings.token):
|
||||
raise HTTPException(status_code=401, detail="неверный или отсутствующий токен")
|
||||
|
||||
|
||||
@app.get("/health")
|
||||
def health() -> JSONResponse:
|
||||
"""Проверка живости - без токена, чтобы годилась для мониторинга."""
|
||||
return JSONResponse({
|
||||
"status": "ok" if _state["ready"] else "no_models",
|
||||
"error": _state["error"],
|
||||
"queue": store.stats(),
|
||||
"threads": settings.effective_threads(),
|
||||
})
|
||||
|
||||
|
||||
@app.post("/v1/jobs", dependencies=[Depends(guard)])
|
||||
async def create_job(
|
||||
file: UploadFile = File(...),
|
||||
speakers: int | None = Query(None, ge=0, le=10,
|
||||
description="число говорящих, 0 = определить автоматически"),
|
||||
) -> dict:
|
||||
if not _state["ready"]:
|
||||
raise HTTPException(status_code=503, detail=_state["error"] or "сервис ещё не готов")
|
||||
|
||||
job_id = store.create(filename=file.filename or "audio",
|
||||
duration_sec=0.0,
|
||||
options={"speakers": settings.speakers if speakers is None else speakers})
|
||||
target = settings.data_dir / "uploads" / job_id
|
||||
size = 0
|
||||
try:
|
||||
with target.open("wb") as out:
|
||||
while chunk := await file.read(1 << 20):
|
||||
size += len(chunk)
|
||||
if size > settings.max_upload_bytes:
|
||||
raise HTTPException(status_code=413,
|
||||
detail=f"файл больше {settings.max_upload_mb} МБ")
|
||||
out.write(chunk)
|
||||
except HTTPException:
|
||||
target.unlink(missing_ok=True)
|
||||
store.mark_failed(job_id, error="файл слишком большой")
|
||||
raise
|
||||
if size == 0:
|
||||
target.unlink(missing_ok=True)
|
||||
store.mark_failed(job_id, error="пустой файл")
|
||||
raise HTTPException(status_code=400, detail="пустой файл")
|
||||
|
||||
return {"job_id": job_id, "status": JobStatus.QUEUED,
|
||||
"queue_position": store.queue_position(job_id)}
|
||||
|
||||
|
||||
@app.get("/v1/jobs/{job_id}", dependencies=[Depends(guard)])
|
||||
def get_job(job_id: str) -> dict:
|
||||
job = store.get(job_id)
|
||||
if job is None:
|
||||
raise HTTPException(status_code=404, detail="задача не найдена")
|
||||
|
||||
body = {"job_id": job_id, "status": job["status"], "filename": job["filename"]}
|
||||
if job["status"] == JobStatus.QUEUED:
|
||||
body["queue_position"] = store.queue_position(job_id)
|
||||
if job["status"] == JobStatus.DONE:
|
||||
body.update(job["result"])
|
||||
if job["status"] == JobStatus.FAILED:
|
||||
body["error"] = job["error"]
|
||||
return body
|
||||
|
||||
|
||||
@app.delete("/v1/jobs/{job_id}", dependencies=[Depends(guard)])
|
||||
def delete_job(job_id: str) -> dict:
|
||||
if store.get(job_id) is None:
|
||||
raise HTTPException(status_code=404, detail="задача не найдена")
|
||||
store.delete(job_id)
|
||||
return {"deleted": job_id}
|
||||
|
||||
|
||||
def _setup_console() -> None:
|
||||
"""Windows-консоль по умолчанию не в UTF-8, иначе русский текст в логах рассыпается."""
|
||||
for stream in (sys.stdout, sys.stderr):
|
||||
try:
|
||||
stream.reconfigure(encoding="utf-8", errors="replace")
|
||||
except (AttributeError, ValueError):
|
||||
pass
|
||||
|
||||
|
||||
def run() -> None:
|
||||
import uvicorn
|
||||
|
||||
_setup_console()
|
||||
logging.basicConfig(level=logging.INFO,
|
||||
format="%(asctime)s %(levelname)s %(name)s: %(message)s")
|
||||
|
||||
if not settings.token:
|
||||
print("\n В config.toml пустой токен - сервис никого не пустит."
|
||||
"\n Впишите значение в [security] token и перезапустите.\n")
|
||||
|
||||
# Модели проверяются здесь, а не по _state: lifespan отработает уже внутри
|
||||
# uvicorn.run, и к тому моменту сообщение выводить поздно.
|
||||
if not (settings.models_dir / "gigaam" / "config.json").is_file():
|
||||
print("\n Модели не найдены. Сначала запустите download_models.bat"
|
||||
"\n Сервис поднимется, но принимать записи не сможет.\n")
|
||||
|
||||
print(f"\n talkscore-asr слушает http://{settings.host}:{settings.port}"
|
||||
f"\n Потоков: {settings.effective_threads()}"
|
||||
f"\n Проверка: curl http://localhost:{settings.port}/health"
|
||||
"\n Остановить: Ctrl+C\n")
|
||||
uvicorn.run(app, host=settings.host, port=settings.port, log_level="info")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
run()
|
||||
Reference in New Issue
Block a user