d3ro-voice/apps/desktop/sidecar/main.py
Yun Chan ba9ef9741e fix: red-team round 3 hardening across desktop, mobile, core and server
Batch of red-team r3 fixes that were in the working tree before the
2026-09-28 design overhaul, committed as one unit with their tests.

- desktop main: STT timeouts and sidecar, voice recording store, sync
  (credentials, audio, knowledge reindex, push gates), runtime
  provisioner, update policy, AltGr keybindings, voice-command policy,
  dictionary file codec/limits, meeting transcript condensing and a
  local recording ledger so interrupted-session recovery only closes
  meetings this device recorded (a phone's live meeting is left alone).
- mobile: login CSRF via implicit token callbacks rejected, account
  deletion/retention, durable queue retention, knowledge realtime
  without unfiltered DELETE, meeting re-record failure paths, cloud STT
  client, preferences store/resync.
- core: text chunking splits long unbroken transcripts to fit, template
  field policy, dictionary limits, meeting markdown inline handling.
- server: payple webhook policy and cancellation order scope, meeting
  document generation quota, team RPC null-role guard, unified LLM
  quota in-flight accounting, knowledge chunk vector index, meeting
  re-record failure paths (migrations 20260929*).
- ci: portable/runtime feed gates, update-policy schema, Forgejo file
  delete and alias planning.

Four older tests are updated to the new contracts rather than the old
behavior: token-pair auth callbacks are rejected, knowledge realtime no
longer subscribes to DELETE, long transcript lines are split, and
meeting recovery requires the local recording ledger for empty rows.
2026-09-28 20:45:52 +09:00

921 lines
32 KiB
Python

"""
D3RO-VOICE STT Sidecar (FastAPI HTTP 서버)
faster-whisper + CTranslate2만 사용. torch/pyannote 비의존으로 슬림 배포.
사용법:
python main.py --port 18765 --models-dir <path>
엔드포인트:
GET /health - 헬스체크
POST /load - Whisper 모델 로딩
POST /transcribe - 오디오 전사 (multipart)
POST /download - 모델 다운로드 시작 (백그라운드)
GET /download/status - 다운로드 진행률 조회
POST /download/cancel - 다운로드 취소
GET /uia/focus - 포커스 입력 요소 UIA 스냅샷 (제안/학습용)
POST /shutdown - 서버 종료
주: 화자 구분(diarization)은 Phase 15.5에서 LLM 추정 경로가 primary이며,
pyannote 기반 고정밀 화자 구분은 추후 서버 사이드 API로 제공될 예정.
"""
from __future__ import annotations
import argparse
import asyncio
import fnmatch
import logging
import os
import signal
import sys
import threading
import time
from contextlib import asynccontextmanager
from pathlib import Path
from typing import AsyncGenerator
import numpy as np
import uvicorn
from fastapi import FastAPI, File, Form, Request, UploadFile
from fastapi.responses import JSONResponse
from device_policy import (
CPU_CHOICE,
DeviceChoice,
choose_gpu_device,
is_gpu_runtime_error,
load_with_fallback,
)
# ── 로깅 설정 ──────────────────────────────────────────────
# Windows에서 파이프로 연결되면 Python이 로케일(cp949) 인코딩으로 출력해
# 메인 프로세스의 UTF-8 로그가 깨진다. 명시적으로 UTF-8로 고정한다.
for _stream in (sys.stdout, sys.stderr):
try:
_stream.reconfigure(encoding="utf-8", errors="replace") # type: ignore[union-attr]
except (AttributeError, ValueError):
pass
logging.basicConfig(
level=logging.INFO,
format="[%(asctime)s] [%(levelname)s] [%(name)s] %(message)s",
datefmt="%Y-%m-%d %H:%M:%S",
stream=sys.stdout,
)
logger = logging.getLogger("sidecar")
# ── 전역 상태 ──────────────────────────────────────────────
_model: "WhisperModel | None" = None
_model_id: str | None = None
# 보조 모델 (실시간 자막 등 받아쓰기와 다른 모델). 최대 1개.
_aux_models: "dict[str, WhisperModel]" = {}
_gpu_available: bool = False
# GPU 로 올릴 때 쓸 장치/정밀도. GPU 가 없거나 검증에 실패하면 None.
_gpu_choice: DeviceChoice | None = None
# 로딩된 모델별 실제 장치 (primary/aux 는 같은 model_id 를 동시에 들지 않는다).
_model_devices: dict[str, DeviceChoice] = {}
_server: uvicorn.Server | None = None
_models_dir: Path | None = None
# 모델 추론·로딩은 한 번에 하나만 돈다(모델 교체와 전사가 겹치지 않게). 무거운 작업은
# asyncio.to_thread 로 돌려, 긴 전사·로딩 중에도 /health · /uia/focus · /download/status 가 응답한다.
# (예전엔 async 핸들러 안에서 동기로 돌아 이벤트 루프 전체가 막혔다.)
_inference_lock = asyncio.Lock()
# 클라이언트가 끊긴 전사를 확인하는 주기(초)
_DISCONNECT_POLL_SECONDS = 0.5
class TranscriptionCancelled(Exception):
"""클라이언트가 연결을 끊어 전사를 중단했다 (타임아웃·취소)."""
class ModelNotInstalled(Exception):
"""models-dir 에도 HF 캐시에도 없는 모델 — 암묵적으로 내려받지 않는다."""
# ── 다운로드 상태 (스레드 공유) ────────────────────────────
_download_lock = threading.Lock()
_download_thread: threading.Thread | None = None
_download_cancel = threading.Event()
_download_state: dict = {
"status": "idle", # idle | downloading | done | cancelled | error
"model_id": None,
"percent": 0,
"downloaded_bytes": 0,
"total_bytes": 0,
"bytes_per_second": 0,
"message": None,
}
# faster-whisper가 다운로드하는 파일과 동일한 화이트리스트
_DOWNLOAD_PATTERNS = [
"config.json",
"preprocessor_config.json",
"model.bin",
"tokenizer.json",
"vocabulary.*",
]
# faster_whisper.utils._MODELS 매핑 실패 시 폴백
_FALLBACK_REPOS = {
"large-v3-turbo": "mobiuslabsgmbh/faster-whisper-large-v3-turbo",
"turbo": "mobiuslabsgmbh/faster-whisper-large-v3-turbo",
}
# ── FastAPI 앱 ─────────────────────────────────────────────
@asynccontextmanager
async def lifespan(_app: FastAPI) -> AsyncGenerator[None, None]:
"""서버 시작/종료 생명주기."""
logger.info("Sidecar 서버 시작")
_detect_gpu()
yield
logger.info("Sidecar 서버 종료")
app = FastAPI(title="D3RO-VOICE STT Sidecar", lifespan=lifespan)
def _detect_gpu() -> None:
"""GPU(CUDA) 사용 가능 여부를 ctranslate2로 감지한다.
torch 의존 제거를 위해 ctranslate2의 네이티브 CUDA 감지를 사용한다.
ctranslate2는 faster-whisper의 백엔드이므로 항상 함께 설치된다.
"""
global _gpu_available, _gpu_choice
try:
import ctranslate2
cuda_count = ctranslate2.get_cuda_device_count()
supported: set[str] = set()
if cuda_count > 0:
try:
supported = set(ctranslate2.get_supported_compute_types("cuda"))
except Exception as exc:
logger.info("CUDA compute type 조회 실패: %s", exc)
_gpu_choice = choose_gpu_device(cuda_count, supported)
_gpu_available = _gpu_choice is not None
if _gpu_choice is not None:
logger.info(
"GPU 감지: CUDA 디바이스 %d개 (compute=%s)",
cuda_count,
_gpu_choice.compute_type,
)
elif cuda_count > 0:
logger.info("CUDA 디바이스는 있으나 지원 compute type 없음, CPU 모드로 동작")
else:
logger.info("GPU 미감지, CPU 모드로 동작")
except Exception as exc:
_gpu_available = False
_gpu_choice = None
logger.info("GPU 감지 실패, CPU 모드로 동작: %s", exc)
def _disable_gpu(reason: BaseException) -> None:
"""GPU 경로가 실제로 동작하지 않음이 확인되면 이후 모델은 모두 CPU 로 올린다."""
global _gpu_available, _gpu_choice
if _gpu_available:
logger.warning("GPU 사용 불가 → CPU(int8) 모드로 전환: %s", reason)
_gpu_available = False
_gpu_choice = None
def _cpu_threads() -> int:
"""CPU 추론에 사용할 스레드 수 (과도한 점유 방지 위해 8로 상한)."""
return max(1, min(8, os.cpu_count() or 4))
# ── 모델 다운로드 헬퍼 ─────────────────────────────────────
def _resolve_repo(model_id: str) -> str:
"""모델 ID를 HuggingFace repo ID로 변환한다."""
if "/" in model_id:
return model_id
try:
from faster_whisper.utils import _MODELS
if model_id in _MODELS:
return _MODELS[model_id]
except Exception:
pass
if model_id in _FALLBACK_REPOS:
return _FALLBACK_REPOS[model_id]
return f"Systran/faster-whisper-{model_id}"
def _local_model_dir(model_id: str) -> Path | None:
"""models_dir 내 다운로드 완료된 모델 디렉토리를 반환한다 (없으면 None)."""
if _models_dir is None:
return None
model_dir = _models_dir / model_id
if (model_dir / "model.bin").exists():
return model_dir
return None
def _set_download_state(**kwargs: object) -> None:
with _download_lock:
_download_state.update(kwargs)
def _download_worker(model_id: str) -> None:
"""백그라운드 스레드: HF repo 파일들을 스트리밍 다운로드한다."""
import requests
from huggingface_hub import HfApi, hf_hub_url
try:
repo_id = _resolve_repo(model_id)
logger.info("모델 다운로드 시작: %s (repo=%s)", model_id, repo_id)
api = HfApi()
info = api.model_info(repo_id, files_metadata=True)
files = [
s
for s in (info.siblings or [])
if any(fnmatch.fnmatch(s.rfilename, p) for p in _DOWNLOAD_PATTERNS)
]
if not files:
raise RuntimeError(f"다운로드할 파일이 없습니다: {repo_id}")
assert _models_dir is not None
target_dir = _models_dir / model_id
target_dir.mkdir(parents=True, exist_ok=True)
total_bytes = sum(s.size or 0 for s in files)
downloaded = 0
start_time = time.monotonic()
_set_download_state(total_bytes=total_bytes, downloaded_bytes=0, percent=0)
for sibling in files:
fname = sibling.rfilename
fsize = sibling.size or 0
dest = target_dir / fname
# 멱등: 이미 크기 일치하는 파일은 스킵
if dest.exists() and fsize > 0 and dest.stat().st_size == fsize:
downloaded += fsize
_set_download_state(
downloaded_bytes=downloaded,
percent=int(downloaded * 100 / total_bytes) if total_bytes else 0,
)
logger.info("이미 존재, 스킵: %s", fname)
continue
if _download_cancel.is_set():
raise InterruptedError()
url = hf_hub_url(repo_id, fname)
part = dest.with_suffix(dest.suffix + ".part")
logger.info("다운로드: %s (%.1f MB)", fname, fsize / 1e6)
with requests.get(url, stream=True, timeout=30) as resp:
resp.raise_for_status()
with open(part, "wb") as fh:
for chunk in resp.iter_content(chunk_size=1024 * 1024):
if _download_cancel.is_set():
raise InterruptedError()
fh.write(chunk)
downloaded += len(chunk)
elapsed = time.monotonic() - start_time
_set_download_state(
downloaded_bytes=downloaded,
percent=(
int(downloaded * 100 / total_bytes) if total_bytes else 0
),
bytes_per_second=int(downloaded / elapsed) if elapsed > 0 else 0,
)
os.replace(part, dest)
_set_download_state(status="done", percent=100)
logger.info("모델 다운로드 완료: %s (%.1f MB)", model_id, downloaded / 1e6)
except InterruptedError:
_set_download_state(status="cancelled", message="사용자 취소")
logger.info("모델 다운로드 취소: %s", model_id)
_cleanup_partial(model_id)
except Exception as exc:
_set_download_state(status="error", message=str(exc))
logger.error("모델 다운로드 실패: %s", exc, exc_info=True)
_cleanup_partial(model_id)
def _cleanup_partial(model_id: str) -> None:
"""취소/실패 시 .part 잔여 파일 정리."""
if _models_dir is None:
return
target_dir = _models_dir / model_id
if not target_dir.exists():
return
for part in target_dir.glob("*.part"):
try:
part.unlink()
except OSError:
pass
def _build_transcribe_kwargs(
language: str,
vad_filter: str,
initial_prompt: str,
is_partial: bool,
) -> dict:
"""전사 옵션을 만든다.
받아쓰기 정합성을 위해 컨텍스트 누적(condition_on_previous_text)을 끈다.
Whisper가 앞 세그먼트 오류를 반복 증폭하는 현상(환각 루프)을 막는다.
미리보기(partial)는 지연이 목표이므로 greedy + VAD 없음으로 디코딩한다.
"""
if is_partial:
kwargs: dict = {
"beam_size": 1,
"temperature": 0.0,
"vad_filter": False,
"condition_on_previous_text": False,
"word_timestamps": False,
}
else:
kwargs = {
"beam_size": 5,
# 0.0 단일 온도는 실패 시 재시도가 없어 환각이 남는다.
# 낮은 온도 폴백만 허용하되 컨텍스트를 끊어 반복을 차단한다.
"temperature": [0.0, 0.2, 0.4],
"condition_on_previous_text": False,
"no_speech_threshold": 0.6,
"compression_ratio_threshold": 2.4,
"log_prob_threshold": -1.0,
"vad_filter": vad_filter.lower() == "true",
"word_timestamps": False,
}
if kwargs["vad_filter"]:
# 무음 구간을 촘촘히 잘라 속도를 올린다.
kwargs["vad_parameters"] = {"min_silence_duration_ms": 300}
if language != "auto":
kwargs["language"] = language
if initial_prompt and not is_partial:
kwargs["initial_prompt"] = initial_prompt
return kwargs
def _run_transcription(
model: "WhisperModel",
audio_array: np.ndarray,
transcribe_kwargs: dict,
cancel: threading.Event | None = None,
) -> tuple[list[dict], str, str]:
"""모델로 전사하고 세그먼트를 끝까지 소비한다.
faster-whisper 는 세그먼트를 지연 생성하므로 디코더 오류(GPU 런타임 포함)는
반복 중에 난다. 호출부가 전체를 한 단위로 재시도할 수 있도록 여기서 모두 소비한다.
Returns: (segments, full_text, detected_language)
"""
# VAD가 전체 오디오를 제거하면 max() 에러 발생 → VAD 없이 재시도
try:
segments_iter, info = model.transcribe(audio_array, **transcribe_kwargs)
except ValueError as ve:
if "empty sequence" in str(ve) and transcribe_kwargs.get("vad_filter"):
logger.warning("VAD가 전체 오디오를 제거함 → VAD 없이 재시도")
transcribe_kwargs["vad_filter"] = False
segments_iter, info = model.transcribe(audio_array, **transcribe_kwargs)
else:
raise
segments_list: list[dict] = []
full_text_parts: list[str] = []
for segment in segments_iter:
# 클라이언트가 떠났으면(타임아웃·취소) 남은 디코딩을 버린다 — 끝까지 돌면 다음 받아쓰기가 막힌다.
if cancel is not None and cancel.is_set():
raise TranscriptionCancelled("client disconnected")
seg_dict = {
"text": segment.text.strip(),
"start": round(segment.start, 3),
"end": round(segment.end, 3),
"avg_logprob": round(segment.avg_logprob, 4),
}
segments_list.append(seg_dict)
full_text_parts.append(segment.text.strip())
full_text = " ".join(full_text_parts).strip()
detected_language = info.language if info.language else "unknown"
return segments_list, full_text, detected_language
# ── 엔드포인트 ─────────────────────────────────────────────
@app.get("/health")
async def health() -> JSONResponse:
"""헬스체크. sidecar가 준비되었는지 확인한다."""
return JSONResponse(
content={
"status": "ready" if _model is not None else "no_model",
"model": _model_id,
"model_loaded": _model is not None,
"aux_models": list(_aux_models.keys()),
"gpu": _gpu_available,
"device": "cuda" if _gpu_available else "cpu",
}
)
def _create_model(model_id: str, choice: DeviceChoice) -> "WhisperModel":
"""모델을 올린다. /download 로 받아 둔 로컬 디렉토리가 있으면 그것을 쓴다."""
from faster_whisper import WhisperModel
local_dir = _local_model_dir(model_id)
model_source = str(local_dir) if local_dir else model_id
if local_dir:
logger.info("로컬 모델 디렉토리 사용: %s", local_dir)
logger.info(
"모델 생성: %s (device=%s, compute=%s)", model_id, choice.device, choice.compute_type
)
return WhisperModel(
model_source,
device=choice.device,
compute_type=choice.compute_type,
cpu_threads=_cpu_threads(),
num_workers=1,
)
# 0.5초 @16kHz 무음. GPU 모델이 cuBLAS 등을 실제로 로드할 수 있는지 확인하는 데 쓴다.
_PROBE_SAMPLES = 8000
def _probe_model(model: "WhisperModel") -> None:
"""짧은 무음 전사로 인코더·디코더를 한 번 돌려 본다 (실패 시 예외)."""
silence = np.zeros(_PROBE_SAMPLES, dtype=np.float32)
segments, _info = model.transcribe(
silence,
language="en",
beam_size=1,
temperature=0.0,
vad_filter=False,
condition_on_previous_text=False,
without_timestamps=True,
)
# faster-whisper 는 세그먼트를 지연 생성하므로 끝까지 소비해야 디코더가 돈다.
for _ in segments:
pass
def _load_on_best_device(model_id: str) -> "WhisperModel":
"""GPU 후보가 있으면 GPU 로 올려 검증하고, 안 되면 CPU(int8)로 올린다."""
def on_fallback(failed: DeviceChoice, exc: BaseException) -> None:
logger.warning(
"GPU 모델 검증 실패 (%s, compute=%s) → CPU 로 재로딩: %s",
model_id,
failed.compute_type,
exc,
)
_disable_gpu(exc)
preferred = _gpu_choice if _gpu_available else None
model, choice = load_with_fallback(
build=lambda c: _create_model(model_id, c),
verify=_probe_model,
preferred=preferred,
on_fallback=on_fallback,
)
_model_devices[model_id] = choice
return model
def _reload_on_cpu(model_id: str, is_primary: bool) -> "WhisperModel":
"""전사 중 GPU 런타임 오류가 난 모델을 CPU 로 다시 올려 같은 자리에 둔다."""
global _model
if is_primary:
_model = None
else:
_aux_models.pop(model_id, None)
model = _create_model(model_id, CPU_CHOICE)
_model_devices[model_id] = CPU_CHOICE
if is_primary:
_model = model
else:
_aux_models[model_id] = model
return model
def _ensure_model_available(model_id: str) -> None:
"""models-dir 을 쓰는 설치본에서, 받지 않은 모델을 /load 가 HF 에서 몰래 내려받지 않게 한다.
models-dir 에 있거나 HF 캐시에 이미 있으면 통과한다. models-dir 을 지정하지 않은 실행(개발·테스트)은
예전처럼 faster-whisper 에 맡긴다(HF 캐시만 사용).
"""
if _models_dir is None or _local_model_dir(model_id) is not None:
return
try:
from faster_whisper.utils import download_model
download_model(model_id, local_files_only=True)
except Exception as exc: # noqa: BLE001 — 캐시에 없으면 여러 종류의 예외가 난다
raise ModelNotInstalled(f"모델이 설치되어 있지 않습니다: {model_id}") from exc
@app.post("/load")
async def load_model(body: dict) -> JSONResponse: # noqa: ANN001
"""Whisper 모델을 로딩한다.
Request body:
{ "model_id": "large-v3" } -- tiny, base, small, medium, large-v3
Returns:
{ "status": "loaded", "model_id": "large-v3", "load_time_ms": 1234 }
"""
global _model, _model_id
model_id: str = body.get("model_id", "large-v3-turbo")
# primary = 받아쓰기(기본) 모델, aux = 실시간 자막처럼 따로 고른 보조 모델.
slot: str = body.get("slot", "primary")
logger.info("모델 로딩 시작: %s (slot=%s)", model_id, slot)
start_time = time.monotonic()
# 같은 모델이 이미 로딩되어 있으면 재사용 (재로딩은 수초 지연을 만든다)
if (_model is not None and _model_id == model_id) or (slot == "aux" and model_id in _aux_models):
logger.info("이미 로딩된 모델 재사용: %s", model_id)
return JSONResponse(
content={
"status": "loaded",
"model_id": model_id,
"load_time_ms": 0,
"reused": True,
}
)
try:
# 설치하지 않은 모델은 올리지 않는다 — 쓰던 모델을 내리기 전에 확인한다.
await asyncio.to_thread(_ensure_model_available, model_id)
except ModelNotInstalled as exc:
logger.warning("모델 로딩 거부 (미설치): %s", exc)
return JSONResponse(
status_code=404,
content={"status": "error", "code": "model_not_installed", "message": str(exc)},
)
try:
async with _inference_lock:
if slot == "aux":
# 보조 자리는 하나만 둔다 — 다른 보조 모델은 내려 VRAM 을 돌려받는다.
_aux_models.clear()
_aux_models[model_id] = await asyncio.to_thread(_load_on_best_device, model_id)
else:
# 모델 교체 시 이전 모델을 먼저 해제해 VRAM/RAM을 회수한다.
_model = None
_model = await asyncio.to_thread(_load_on_best_device, model_id)
_model_id = model_id
# 기본 모델이 된 모델은 보조 자리에 중복으로 들고 있지 않는다.
_aux_models.pop(model_id, None)
load_time_ms = int((time.monotonic() - start_time) * 1000)
logger.info("모델 로딩 완료: %s (slot=%s, %dms)", model_id, slot, load_time_ms)
return JSONResponse(
content={
"status": "loaded",
"model_id": model_id,
"load_time_ms": load_time_ms,
}
)
except Exception as exc:
logger.error("모델 로딩 실패: %s", exc)
return JSONResponse(
status_code=500,
content={"status": "error", "message": str(exc)},
)
async def _watch_disconnect(request: Request, cancel: threading.Event) -> None:
"""클라이언트가 연결을 끊으면(타임아웃 abort 등) cancel 을 세운다."""
while not cancel.is_set():
if await request.is_disconnected():
cancel.set()
return
await asyncio.sleep(_DISCONNECT_POLL_SECONDS)
def _transcribe_blocking(
model: "WhisperModel",
audio_array: np.ndarray,
transcribe_kwargs: dict,
resolved_model_id: str | None,
is_primary: bool,
cancel: threading.Event,
) -> tuple[list[dict], str, str]:
"""워커 스레드에서 전사한다 (GPU 런타임 오류면 CPU 로 다시 올려 한 번만 재시도)."""
try:
return _run_transcription(model, audio_array, dict(transcribe_kwargs), cancel)
except TranscriptionCancelled:
raise
except Exception as exc:
# cuBLAS 미설치 PC 등: 검증을 통과했더라도 실제 전사에서 GPU 런타임 오류가
# 나면 이 모델을 CPU 로 다시 올려 한 번만 재시도한다.
device = _model_devices.get(resolved_model_id or "")
if not (
resolved_model_id
and device is not None
and device.is_gpu
and is_gpu_runtime_error(exc)
):
raise
logger.warning(
"전사 중 GPU 런타임 오류 → CPU 로 재로딩 후 재시도 (%s): %s",
resolved_model_id,
exc,
)
_disable_gpu(exc)
cpu_model = _reload_on_cpu(resolved_model_id, is_primary)
return _run_transcription(cpu_model, audio_array, dict(transcribe_kwargs), cancel)
@app.post("/transcribe")
async def transcribe(
request: Request,
audio: UploadFile = File(...),
language: str = Form("auto"),
vad_filter: str = Form("true"),
initial_prompt: str = Form(""),
partial: str = Form("false"),
model_id: str = Form(""),
) -> JSONResponse:
"""오디오 파일을 전사한다.
Multipart form:
audio - PCM16 16kHz mono 바이너리 파일
language - 언어 코드 ('auto', 'ko', 'en', ...)
vad_filter - VAD 필터 활성화 ('true' / 'false')
initial_prompt - 초기 프롬프트 (컨텍스트 힌트)
partial - 녹음 중 미리보기 모드 ('true'면 greedy 디코딩 + 컨텍스트 미사용)
model_id - 쓸 모델 (비우면 기본 모델). 올라가 있지 않으면 409
"""
is_primary = not (model_id and model_id != _model_id)
if not is_primary:
model = _aux_models.get(model_id)
if model is None:
return JSONResponse(
status_code=409,
content={"status": "error", "code": "model_not_loaded", "message": f"모델이 로딩되지 않았습니다: {model_id}"},
)
else:
model = _model
resolved_model_id = model_id if not is_primary else _model_id
if model is None:
return JSONResponse(
status_code=503,
content={"status": "error", "message": "모델이 로딩되지 않았습니다"},
)
is_partial = partial.lower() == "true"
start_time = time.monotonic()
try:
pcm_bytes = await audio.read()
if len(pcm_bytes) == 0:
return JSONResponse(
status_code=400,
content={"status": "error", "message": "오디오 데이터가 비어있습니다"},
)
audio_array = (
np.frombuffer(pcm_bytes, dtype=np.int16).astype(np.float32) / 32768.0
)
sample_rate = 16000
audio_duration = len(audio_array) / sample_rate
logger.info(
"전사 시작: %.1f초 오디오, language=%s, vad=%s, partial=%s",
audio_duration,
language,
vad_filter,
is_partial,
)
transcribe_kwargs = _build_transcribe_kwargs(
language=language,
vad_filter=vad_filter,
initial_prompt=initial_prompt,
is_partial=is_partial,
)
cancel = threading.Event()
watcher = asyncio.create_task(_watch_disconnect(request, cancel))
try:
async with _inference_lock:
if cancel.is_set():
raise TranscriptionCancelled("client disconnected while queued")
segments_list, full_text, detected_language = await asyncio.to_thread(
_transcribe_blocking,
model,
audio_array,
transcribe_kwargs,
resolved_model_id,
is_primary,
cancel,
)
finally:
cancel.set()
watcher.cancel()
processing_time = int((time.monotonic() - start_time) * 1000)
logger.info(
"전사 완료: '%s' (lang=%s, %.1f초, %dms)",
full_text[:80],
detected_language,
audio_duration,
processing_time,
)
return JSONResponse(
content={
"text": full_text,
"segments": segments_list,
"language": detected_language,
"duration": round(audio_duration, 3),
"processing_time": processing_time,
}
)
except TranscriptionCancelled as exc:
logger.info("전사 중단: %s", exc)
return JSONResponse(
status_code=499,
content={"status": "error", "code": "cancelled", "message": str(exc)},
)
except Exception as exc:
logger.error("전사 실패: %s", exc, exc_info=True)
return JSONResponse(
status_code=500,
content={"status": "error", "message": str(exc)},
)
@app.post("/download")
async def download_model(body: dict) -> JSONResponse: # noqa: ANN001
"""모델 다운로드를 백그라운드로 시작한다.
Request body:
{ "model_id": "large-v3-turbo" }
Returns:
{ "status": "started" } 또는 이미 완료된 경우 { "status": "done" }
"""
global _download_thread
if _models_dir is None:
return JSONResponse(
status_code=500,
content={"status": "error", "message": "models-dir가 설정되지 않았습니다"},
)
model_id: str = body.get("model_id", "large-v3-turbo")
# 이미 다운로드 완료된 모델이면 즉시 done
if _local_model_dir(model_id) is not None:
_set_download_state(status="done", model_id=model_id, percent=100)
return JSONResponse(content={"status": "done"})
if _download_thread is not None and _download_thread.is_alive():
return JSONResponse(
status_code=409,
content={"status": "error", "message": "이미 다운로드가 진행 중입니다"},
)
_download_cancel.clear()
_set_download_state(
status="downloading",
model_id=model_id,
percent=0,
downloaded_bytes=0,
total_bytes=0,
bytes_per_second=0,
message=None,
)
_download_thread = threading.Thread(
target=_download_worker, args=(model_id,), daemon=True
)
_download_thread.start()
return JSONResponse(content={"status": "started"})
@app.get("/download/status")
async def download_status() -> JSONResponse:
"""현재 다운로드 상태를 반환한다."""
with _download_lock:
return JSONResponse(content=dict(_download_state))
@app.post("/download/cancel")
async def download_cancel() -> JSONResponse:
"""진행 중인 다운로드를 취소한다."""
if _download_thread is not None and _download_thread.is_alive():
_download_cancel.set()
return JSONResponse(content={"status": "cancelling"})
return JSONResponse(content={"status": "idle"})
@app.get("/uia/focus")
async def uia_focus(timeoutMs: int = 1500) -> JSONResponse:
"""포커스된 입력 요소의 UIA 스냅샷 (텍스트/케어렛/비밀번호 여부).
제안(ghost text)과 이핑 학습이 쓰는 유일한 입력창 읽기 경로다.
비밀번호 필드는 브리지 안에서 fail-closed 로 차단한다.
"""
try:
import uia_bridge
except Exception as exc: # pragma: no cover - 파일 누락 등
return JSONResponse(content={"available": False, "reason": f"bridge-import:{exc}"})
snapshot = await asyncio.to_thread(uia_bridge.snapshot_focus, timeoutMs)
return JSONResponse(content=snapshot)
@app.post("/shutdown")
async def shutdown() -> JSONResponse:
"""서버를 graceful하게 종료한다."""
logger.info("종료 요청 수신")
try:
import uia_bridge
await asyncio.to_thread(uia_bridge.shutdown)
except Exception:
pass
if _server is not None:
_server.should_exit = True
return JSONResponse(content={"status": "shutting_down"})
# ── 메인 ───────────────────────────────────────────────────
def main() -> None:
"""CLI 진입점."""
global _server, _models_dir
parser = argparse.ArgumentParser(description="D3RO-VOICE STT Sidecar")
parser.add_argument(
"--port",
type=int,
default=18765,
help="HTTP 서버 포트 (기본: 18765)",
)
parser.add_argument(
"--host",
type=str,
default="127.0.0.1",
help="HTTP 서버 호스트 (기본: 127.0.0.1)",
)
parser.add_argument(
"--models-dir",
type=str,
default=None,
help="사전 다운로드 모델 저장 디렉토리 (미지정 시 HF 캐시만 사용)",
)
args = parser.parse_args()
if args.models_dir:
_models_dir = Path(args.models_dir)
_models_dir.mkdir(parents=True, exist_ok=True)
logger.info("D3RO-VOICE STT Sidecar 시작 (port=%d)", args.port)
def signal_handler(signum: int, _frame: object) -> None:
sig_name = signal.Signals(signum).name
logger.info("시그널 수신: %s, 종료 시작", sig_name)
if _server is not None:
_server.should_exit = True
signal.signal(signal.SIGINT, signal_handler)
signal.signal(signal.SIGTERM, signal_handler)
config = uvicorn.Config(
app=app,
host=args.host,
port=args.port,
log_level="warning",
access_log=False,
)
_server = uvicorn.Server(config)
_server.run()
logger.info("Sidecar 서버 종료 완료")
if __name__ == "__main__":
main()