import hashlib import hmac import logging import mimetypes import shutil import subprocess import uuid from pathlib import Path from typing import Any, BinaryIO, Dict, List, Optional import json from fastapi import ( BackgroundTasks, Depends, FastAPI, File, Form, Header, HTTPException, Path as ApiPath, Query, Request, UploadFile, ) from fastapi.responses import FileResponse, Response from fastapi.staticfiles import StaticFiles from starlette.concurrency import run_in_threadpool from starlette.requests import ClientDisconnect from .assessment import AssessmentService from .audio_metrics import AudioAnalysisError from .config import Settings from .models import ( AssessmentResult, SentenceBoundary, SentenceBoundaryDocument, SentenceTextUpdate, VideoDetailResponse, VideoListResponse, VideoSummary, VideoUploadResponse, ) from .processing import VideoProcessor from .repository import VideoRepository from .store import BoundaryStore from .transcription import MossTranscriber, Transcriber logger = logging.getLogger(__name__) SHA256_PATH = ApiPath( min_length=64, max_length=64, pattern=r"^[A-Fa-f0-9]{64}$", description="SHA-256 hex digest of the exact uploaded video bytes", ) ALLOWED_VIDEO_EXTENSIONS = {".mp4", ".m4v", ".mov", ".mkv", ".webm"} ALLOWED_AUDIO_EXTENSIONS = {".m4a", ".mp4", ".wav"} class UploadTooLargeError(ValueError): pass def create_app( store: Optional[BoundaryStore] = None, *, settings: Optional[Settings] = None, repository: Optional[VideoRepository] = None, transcriber: Optional[Transcriber] = None, ) -> FastAPI: service_settings = settings or Settings.from_env() service_settings.ensure_directories() video_repository = repository or VideoRepository(service_settings.database_path) legacy_store = store or BoundaryStore(service_settings.legacy_boundaries_file) moss = transcriber or MossTranscriber( endpoint=service_settings.moss_transcribe_url, model=service_settings.moss_model, timeout_seconds=service_settings.moss_timeout_seconds, max_new_tokens=service_settings.moss_max_new_tokens, ) processor = VideoProcessor(service_settings, video_repository, moss) assessment_service = AssessmentService( video_repository, moss, service_settings.pass_score, ) application = FastAPI( title="Oral Trainer Video Service", version="2.0.0", description="Video ingestion, sentence transcription, streaming, and oral-reading assessment.", ) application.state.settings = service_settings application.state.boundary_store = legacy_store application.state.video_repository = video_repository application.state.transcriber = moss application.state.video_processor = processor application.state.assessment_service = assessment_service static_dir = Path(__file__).resolve().parent / "static" if static_dir.is_dir(): application.mount("/static", StaticFiles(directory=static_dir), name="static") @application.exception_handler(ClientDisconnect) async def handle_client_disconnect(request: Request, exc: ClientDisconnect) -> Response: logger.warning( "Upload interrupted: the client or an intermediate proxy closed the " "connection before the request body completed." ) return Response(status_code=400, content="Upload connection was interrupted.") def require_admin(x_admin_key: Optional[str] = Header(default=None)) -> None: expected = service_settings.admin_api_key if expected and not hmac.compare_digest(x_admin_key or "", expected): raise HTTPException(status_code=401, detail="A valid X-Admin-Key header is required.") def require_client(x_client_key: Optional[str] = Header(default=None)) -> None: expected = service_settings.client_api_key if expected and not hmac.compare_digest(x_client_key or "", expected): raise HTTPException(status_code=401, detail="A valid X-Client-Key header is required.") def find_document(video_hash: str) -> Optional[SentenceBoundaryDocument]: return video_repository.get_document(video_hash) or legacy_store.get(video_hash) @application.get("/healthz") def healthz() -> Dict[str, Any]: videos = video_repository.list_videos() return { "status": "ok", "video_count": len(videos) + legacy_store.count(), "managed_video_count": len(videos), "moss_configured": moss.available, "scoring_version": "asr-fluency-v1", } @application.get("/admin", include_in_schema=False) def admin_page() -> FileResponse: page = static_dir / "admin.html" if not page.is_file(): raise HTTPException(status_code=404, detail="Admin UI is not installed.") return FileResponse(page) @application.get("/api/v1/videos", response_model=VideoListResponse) def list_videos() -> VideoListResponse: return VideoListResponse( videos=[_video_summary(row, service_settings) for row in video_repository.list_videos()] ) @application.get("/api/v1/videos/{video_hash}", response_model=VideoDetailResponse) def get_video(video_hash: str = SHA256_PATH) -> VideoDetailResponse: row = video_repository.get_video(video_hash) if row is None: raise HTTPException(status_code=404, detail="Video was not found.") return VideoDetailResponse( video=_video_summary(row, service_settings), boundaries=video_repository.get_document(video_hash), ) @application.get("/api/v1/videos/{video_hash}/content") def stream_video(video_hash: str = SHA256_PATH) -> FileResponse: row = video_repository.get_video(video_hash) if row is None: raise HTTPException(status_code=404, detail="Video was not found.") media_path = service_settings.videos_dir / row["stored_filename"] if not media_path.is_file(): raise HTTPException(status_code=404, detail="Stored video file is missing.") return FileResponse( media_path, media_type=row["content_type"], headers={"Accept-Ranges": "bytes", "Cache-Control": "public, max-age=3600"}, ) @application.get( "/api/v1/videos/{video_hash}/sentence-boundaries", response_model=SentenceBoundaryDocument, ) @application.get( "/api/v1/sentence-boundaries/{video_hash}", response_model=SentenceBoundaryDocument, include_in_schema=False, ) def get_sentence_boundaries(video_hash: str = SHA256_PATH) -> SentenceBoundaryDocument: document = find_document(video_hash) if document is None: raise HTTPException( status_code=404, detail={ "code": "SENTENCE_BOUNDARIES_NOT_FOUND", "message": "No sentence boundaries are registered for this video hash.", }, ) return document @application.post( "/api/v1/admin/videos", response_model=VideoUploadResponse, status_code=202, dependencies=[Depends(require_admin)], ) async def upload_video( background_tasks: BackgroundTasks, file: UploadFile = File(...), title: Optional[str] = Form(default=None), language: Optional[str] = Form(default=None), ) -> VideoUploadResponse: original_name = Path(file.filename or "video.mp4").name suffix = _validate_video_extension(original_name) temporary_path = service_settings.work_dir / f"upload-{uuid.uuid4().hex}.part" try: video_hash, size_bytes = await run_in_threadpool( _persist_upload, file.file, temporary_path, service_settings.max_upload_bytes, ) except UploadTooLargeError as exc: temporary_path.unlink(missing_ok=True) raise HTTPException(status_code=413, detail=str(exc)) from exc except ValueError as exc: temporary_path.unlink(missing_ok=True) raise HTTPException(status_code=400, detail=str(exc)) from exc finally: await file.close() content_type = file.content_type or mimetypes.guess_type(original_name)[0] or "video/mp4" response = _register_uploaded_video( settings=service_settings, repository=video_repository, temporary_path=temporary_path, video_hash=video_hash, size_bytes=size_bytes, original_name=original_name, suffix=suffix, content_type=content_type, title=title, language=language, ) background_tasks.add_task(_process_safely, processor, video_hash) return response @application.put( "/api/v1/admin/videos/raw", response_model=VideoUploadResponse, status_code=202, dependencies=[Depends(require_admin)], ) async def upload_video_stream( request: Request, background_tasks: BackgroundTasks, filename: str = Query(min_length=1, max_length=500), title: Optional[str] = Query(default=None, max_length=200), language: Optional[str] = Query(default=None, max_length=32), ) -> VideoUploadResponse: original_name = Path(filename).name suffix = _validate_video_extension(original_name) content_length = request.headers.get("content-length") if content_length: try: declared_size = int(content_length) except ValueError as exc: raise HTTPException(status_code=400, detail="Invalid Content-Length header.") from exc if declared_size > service_settings.max_upload_bytes: raise HTTPException(status_code=413, detail="Upload exceeds the configured limit.") temporary_path = service_settings.work_dir / f"upload-{uuid.uuid4().hex}.part" try: video_hash, size_bytes = await _persist_request_stream( request, temporary_path, service_settings.max_upload_bytes, ) except UploadTooLargeError as exc: raise HTTPException(status_code=413, detail=str(exc)) from exc except ValueError as exc: raise HTTPException(status_code=400, detail=str(exc)) from exc response = _register_uploaded_video( settings=service_settings, repository=video_repository, temporary_path=temporary_path, video_hash=video_hash, size_bytes=size_bytes, original_name=original_name, suffix=suffix, content_type=request.headers.get("content-type") or "application/octet-stream", title=title, language=language, ) background_tasks.add_task(_process_safely, processor, video_hash) return response @application.post( "/api/v1/admin/videos/{video_hash}/process", status_code=202, dependencies=[Depends(require_admin)], ) def reprocess_video( background_tasks: BackgroundTasks, video_hash: str = SHA256_PATH, ) -> Dict[str, str]: if video_repository.get_video(video_hash) is None: raise HTTPException(status_code=404, detail="Video was not found.") background_tasks.add_task(_process_safely, processor, video_hash) return {"video_hash": video_hash.lower(), "status": "processing_queued"} @application.put( "/api/v1/admin/videos/{video_hash}/sentences/{sentence_index}", response_model=SentenceBoundary, dependencies=[Depends(require_admin)], ) def update_sentence( payload: SentenceTextUpdate, video_hash: str = SHA256_PATH, sentence_index: int = ApiPath(ge=0), ) -> SentenceBoundary: sentence = video_repository.update_sentence_text( video_hash, sentence_index, payload.text, payload.language, ) if sentence is None: raise HTTPException(status_code=404, detail="Sentence was not found.") return sentence @application.delete( "/api/v1/admin/videos/{video_hash}", status_code=204, dependencies=[Depends(require_admin)], ) def delete_video(video_hash: str = SHA256_PATH) -> Response: row = video_repository.delete_video(video_hash) if row is None: raise HTTPException(status_code=404, detail="Video was not found.") (service_settings.videos_dir / row["stored_filename"]).unlink(missing_ok=True) return Response(status_code=204) @application.post( "/api/v1/videos/{video_hash}/sentences/{sentence_index}/assessments", response_model=AssessmentResult, dependencies=[Depends(require_client)], ) async def assess_sentence( video_hash: str = SHA256_PATH, sentence_index: int = ApiPath(ge=0), audio: UploadFile = File(...), language: Optional[str] = Form(default=None), ) -> AssessmentResult: document = find_document(video_hash) if document is None: raise HTTPException(status_code=404, detail="Video or sentence boundaries were not found.") if not moss.available: raise HTTPException(status_code=503, detail="MOSS transcription is not configured.") suffix = Path(audio.filename or "attempt.wav").suffix.lower() or ".wav" temporary_name = f"attempt-{uuid.uuid4().hex}{suffix}" temporary_path = service_settings.attempts_dir / temporary_name completed = False try: await run_in_threadpool( _persist_upload, audio.file, temporary_path, service_settings.max_attempt_bytes, ) result = await run_in_threadpool( assessment_service.assess, document=document, sentence_index=sentence_index, audio_path=temporary_path, language=language, retained_audio_filename=temporary_name if service_settings.keep_attempt_audio else None, ) completed = True return result except UploadTooLargeError as exc: raise HTTPException(status_code=413, detail=str(exc)) from exc except LookupError as exc: raise HTTPException(status_code=404, detail=str(exc)) from exc except (ValueError, AudioAnalysisError) as exc: raise HTTPException(status_code=422, detail=str(exc)) from exc except RuntimeError as exc: raise HTTPException(status_code=502, detail=str(exc)) from exc finally: await audio.close() if not service_settings.keep_attempt_audio or not completed: temporary_path.unlink(missing_ok=True) @application.post("/api/v1/dub-shares", status_code=201) async def create_dub_share( request: Request, video_hash: str = Form(...), title: str = Form(default="我的口语配音"), segments: str = Form(...), scores: str = Form(default="[]"), files: List[UploadFile] = File(...), ) -> Dict[str, Any]: try: segment_items = json.loads(segments) except (TypeError, ValueError) as exc: raise HTTPException(status_code=400, detail="segments must be valid JSON.") from exc if not isinstance(segment_items, list) or not segment_items: raise HTTPException(status_code=400, detail="At least one dubbing segment is required.") if len(segment_items) != len(files): raise HTTPException(status_code=400, detail="Segment count does not match audio file count.") try: score_items = json.loads(scores) except (TypeError, ValueError) as exc: raise HTTPException(status_code=400, detail="scores must be valid JSON.") from exc if not isinstance(score_items, list): raise HTTPException(status_code=400, detail="scores must be a JSON array.") scores_by_index = { int(item["sentence_index"]): item for item in score_items if isinstance(item, dict) and "sentence_index" in item } document = find_document(video_hash.lower()) if document is None: raise HTTPException(status_code=404, detail="Video or sentence boundaries were not found.") prepared_segments: List[Dict[str, Any]] = [] saved_files: List[Path] = [] try: for item, upload in zip(segment_items, files): sentence_index = int(item["sentence_index"]) boundary = next((s for s in document.sentences if s.index == sentence_index), None) if boundary is None: raise HTTPException(status_code=404, detail=f"Sentence {sentence_index} was not found.") suffix = Path(upload.filename or "dub.m4a").suffix.lower() if suffix == ".mp4": suffix = ".m4a" if suffix not in ALLOWED_AUDIO_EXTENSIONS: raise HTTPException(status_code=415, detail=f"Unsupported audio extension: {suffix or '(none)'}") audio_name = f"{uuid.uuid4().hex}{suffix}" audio_path = service_settings.dub_shares_dir / audio_name await run_in_threadpool(_persist_upload, upload.file, audio_path, 50 * 1024 * 1024) await run_in_threadpool(_trim_leading_silence, audio_path) saved_files.append(audio_path) score = scores_by_index.get(sentence_index, {}) prepared_segments.append({ "sentence_index": sentence_index, "start_ms": boundary.start_ms, "end_ms": boundary.end_ms, "text": boundary.text, "audio_filename": audio_name, "audio_size_bytes": audio_path.stat().st_size, "overall_score": float(score.get("overall_score", 0)), "recognized_text": str(score.get("recognized_text") or ""), }) share = video_repository.create_dub_share( video_hash=video_hash.lower(), title=title, segments=prepared_segments, ) if share is None: raise HTTPException(status_code=404, detail="Video was not found.") return { **share, "share_url": f"/dub-shares/{share['share_id']}", } except UploadTooLargeError as exc: raise HTTPException(status_code=413, detail=str(exc)) from exc except HTTPException: raise finally: for upload in files: await upload.close() # Files intentionally remain until share deletion because the web page references them. del saved_files @application.get("/api/v1/dub-shares/{share_id}") def get_dub_share(share_id: str) -> Dict[str, Any]: share = video_repository.get_dub_share(share_id) if share is None: raise HTTPException(status_code=404, detail="Dubbing share was not found.") return { "share_id": share["share_id"], "video_hash": share["video_hash"], "title": share["title"], "created_at": share["created_at"], "video_hash": share["video_hash"], "segments": [ { "index": row["sentence_index"], "start_ms": row["start_ms"], "end_ms": row["end_ms"], "text": row["text"], "overall_score": row["overall_score"], "recognized_text": row["recognized_text"], "audio_url": f"/api/v1/dub-shares/{share_id}/audio/{row['audio_filename']}", } for row in share["segments"] ], } @application.get("/api/v1/dub-shares/{share_id}/audio/{filename}") def get_dub_share_audio(share_id: str, filename: str) -> FileResponse: share = video_repository.get_dub_share(share_id) if share is None: raise HTTPException(status_code=404, detail="Dubbing share was not found.") safe_name = Path(filename).name row = next((item for item in share["segments"] if item["audio_filename"] == safe_name), None) if row is None: raise HTTPException(status_code=404, detail="Audio was not found.") media_path = service_settings.dub_shares_dir / safe_name if not media_path.is_file(): raise HTTPException(status_code=404, detail="Stored audio file is missing.") return FileResponse(media_path, media_type="audio/mp4", headers={"Cache-Control": "public, max-age=31536000"}) @application.get("/dub-shares/{share_id}", include_in_schema=False) def dub_share_page(share_id: str) -> FileResponse: if video_repository.get_dub_share(share_id) is None: raise HTTPException(status_code=404, detail="Dubbing share was not found.") page = static_dir / "dub-share.html" if not page.is_file(): raise HTTPException(status_code=404, detail="Dubbing share UI is not installed.") return FileResponse(page) return application def _persist_upload(source: BinaryIO, destination: Path, max_bytes: int) -> tuple[str, int]: digest = hashlib.sha256() size = 0 destination.parent.mkdir(parents=True, exist_ok=True) try: with destination.open("wb") as output: while True: chunk = source.read(4 * 1024 * 1024) if not chunk: break size += len(chunk) if size > max_bytes: raise UploadTooLargeError(f"Upload exceeds the {max_bytes}-byte limit.") digest.update(chunk) output.write(chunk) if size == 0: raise ValueError("Uploaded file is empty.") except Exception: destination.unlink(missing_ok=True) raise return digest.hexdigest(), size def _trim_leading_silence(audio_path: Path) -> None: ffmpeg = shutil.which("ffmpeg") if ffmpeg is None: return trimmed_path = audio_path.with_name(f"{audio_path.stem}-trimmed{audio_path.suffix}") command = [ ffmpeg, "-hide_banner", "-loglevel", "error", "-y", "-i", str(audio_path), "-af", "silenceremove=start_periods=1:start_threshold=-45dB:start_silence=0.08", "-c:a", "aac", str(trimmed_path), ] try: subprocess.run(command, check=True) if trimmed_path.stat().st_size == 0: raise RuntimeError("Trimmed audio is empty.") trimmed_path.replace(audio_path) except (OSError, subprocess.SubprocessError): trimmed_path.unlink(missing_ok=True) async def _persist_request_stream( request: Request, destination: Path, max_bytes: int, ) -> tuple[str, int]: digest = hashlib.sha256() size = 0 destination.parent.mkdir(parents=True, exist_ok=True) try: with destination.open("wb") as output: async for chunk in request.stream(): if not chunk: continue size += len(chunk) if size > max_bytes: raise UploadTooLargeError(f"Upload exceeds the {max_bytes}-byte limit.") digest.update(chunk) output.write(chunk) if size == 0: raise ValueError("Uploaded file is empty.") except Exception: destination.unlink(missing_ok=True) raise return digest.hexdigest(), size def _validate_video_extension(filename: str) -> str: suffix = Path(filename).suffix.lower() if suffix not in ALLOWED_VIDEO_EXTENSIONS: raise HTTPException( status_code=415, detail=f"Unsupported video extension: {suffix or '(none)'}", ) return suffix def _register_uploaded_video( *, settings: Settings, repository: VideoRepository, temporary_path: Path, video_hash: str, size_bytes: int, original_name: str, suffix: str, content_type: str, title: Optional[str], language: Optional[str], ) -> VideoUploadResponse: stored_filename = f"{video_hash}{suffix}" destination = settings.videos_dir / stored_filename if destination.exists(): temporary_path.unlink(missing_ok=True) else: temporary_path.replace(destination) repository.upsert_upload( video_hash=video_hash, title=(title or Path(original_name).stem).strip() or video_hash, filename=original_name, stored_filename=stored_filename, content_type=content_type, size_bytes=size_bytes, language=language.strip() if language else None, ) return VideoUploadResponse( video_hash=video_hash, status="uploaded", detail_url=f"/api/v1/videos/{video_hash}", ) def _process_safely(processor: VideoProcessor, video_hash: str) -> None: try: processor.process(video_hash) except Exception: logger.exception("Video processing failed for %s", video_hash) def _video_summary(row: Dict[str, Any], settings: Settings) -> VideoSummary: relative_stream_url = f"/api/v1/videos/{row['video_hash']}/content" stream_url = f"{settings.public_base_url}{relative_stream_url}" if settings.public_base_url else relative_stream_url return VideoSummary( video_hash=row["video_hash"], title=row["title"], filename=row["filename"], content_type=row["content_type"], size_bytes=row["size_bytes"], duration_ms=row["duration_ms"], language=row["language"], status=row["status"], error_message=row["error_message"], sentence_count=row["sentence_count"], stream_url=stream_url, created_at=row["created_at"], updated_at=row["updated_at"], ) app = create_app()