import base64 import copy import json import math import re import subprocess import tempfile from pathlib import Path from typing import Any import httpx from app.config import settings from app.profile_store import Profile class VisionModelError(RuntimeError): pass VIDEO_RESULT_SCHEMA = { "type": "object", "properties": { "passed": { "type": "boolean", "description": ( "Whether the user's acceptance criterion is satisfied. For negative " "criteria such as disappearance or shutdown, return true when that " "negative state is visibly achieved." ), }, "evidenceSufficient": {"type": "boolean"}, "confidence": {"type": "number", "minimum": 0, "maximum": 1}, "conclusion": {"type": "string"}, "summary": {"type": "string"}, "events": { "type": "array", "items": { "type": "object", "properties": { "timeRange": {"type": "string"}, "event": {"type": "string"}, "confidence": {"type": "number"}, }, "required": ["timeRange", "event", "confidence"], }, }, "warnings": {"type": "array", "items": {"type": "string"}}, }, "required": [ "passed", "evidenceSufficient", "confidence", "conclusion", "summary", "events", "warnings", ], } VIDEO_COMPACT_RESULT_SCHEMA = { "type": "object", "properties": { "passed": VIDEO_RESULT_SCHEMA["properties"]["passed"], "evidenceSufficient": {"type": "boolean"}, "confidence": {"type": "number", "minimum": 0, "maximum": 1}, "conclusion": {"type": "string"}, "summary": {"type": "string"}, }, "required": [ "passed", "evidenceSufficient", "confidence", "conclusion", "summary", ], } VIDEO_TUNING_LIMITS = { "sampleFps": (0.25, 6.0), "maxWidth": (640, 1280), "confidenceThreshold": (0.5, 0.95), } def _duration(path: Path) -> float: process = subprocess.run( [ "ffprobe", "-v", "error", "-show_entries", "format=duration", "-of", "default=noprint_wrappers=1:nokey=1", str(path), ], capture_output=True, text=True, check=False, ) try: value = float(process.stdout.strip()) except ValueError as exc: raise ValueError("Unable to read video duration") from exc if value <= 0: raise ValueError("Video duration must be positive") return value def _frame_sampling_plan( duration: float, sample_fps: float, maximum_frames: int ) -> tuple[int, int, float, float, bool]: requested_count = max(2, math.ceil(duration * sample_fps) + 1) count = min(requested_count, max(2, maximum_frames)) prefix_count = count - 1 capped = count < requested_count prefix_fps = prefix_count / duration if capped else sample_fps end_margin = min(0.1, max(0.001, duration / 2)) tail_time = max(0.0, duration - end_margin) return count, prefix_count, prefix_fps, tail_time, capped def _model_frame_limit() -> int: return max( 2, min(settings.video_sampling_frame_limit, settings.video_model_frame_limit), ) def _extract_frames( path: Path, sample_fps: float, maximum_frames: int, maximum_width: int ) -> tuple[list[Path], float, Path]: duration = _duration(path) count, prefix_count, prefix_fps, tail_time, _ = _frame_sampling_plan( duration, sample_fps, maximum_frames ) directory = Path(tempfile.mkdtemp(prefix="frames-", dir=settings.jobs_dir)) output = directory / "frame-%03d.jpg" prefix_process = subprocess.run( [ "ffmpeg", "-y", "-v", "error", "-i", str(path), "-vf", f"fps={prefix_fps:.8f},scale='min({maximum_width},iw)':-2", "-frames:v", str(prefix_count), "-q:v", "3", str(output), ], capture_output=True, check=False, ) tail_output = directory / f"frame-{count:03d}.jpg" tail_process = subprocess.run( [ "ffmpeg", "-y", "-v", "error", "-i", str(path), "-ss", f"{tail_time:.6f}", "-vf", f"scale='min({maximum_width},iw)':-2", "-frames:v", "1", "-q:v", "3", str(tail_output), ], capture_output=True, check=False, ) frames = sorted(directory.glob("frame-*.jpg")) if prefix_process.returncode != 0 or tail_process.returncode != 0 or len(frames) < 2: error = ( prefix_process.stderr.decode("utf-8", errors="replace") + tail_process.stderr.decode("utf-8", errors="replace") )[-500:] detail = error or ( f"prefixExit={prefix_process.returncode}, tailExit={tail_process.returncode}, " f"expectedAtMost={count}, actual={len(frames)}" ) raise ValueError(f"Unable to extract video frames: {detail}") return frames, duration, directory def _select_frames(frames: list[Path], maximum: int) -> list[Path]: if len(frames) <= maximum: return frames if maximum <= 1: return [frames[len(frames) // 2]] indices = [round(index * (len(frames) - 1) / (maximum - 1)) for index in range(maximum)] return [frames[index] for index in indices] def _video_context_size(base_context: int, image_count: int, output_tokens: int) -> int: estimated_tokens = 4096 + image_count * 1400 + output_tokens required_context = 1 << max(1, estimated_tokens - 1).bit_length() return min( settings.max_video_context_tokens, max(base_context, required_context), ) def _encode_image(path: Path, maximum_width: int, extracted_width: int) -> str: if maximum_width >= extracted_width: image = path.read_bytes() else: process = subprocess.run( [ "ffmpeg", "-v", "error", "-i", str(path), "-vf", f"scale='min({maximum_width},iw)':-2", "-q:v", "3", "-f", "image2pipe", "-vcodec", "mjpeg", "pipe:1", ], capture_output=True, check=False, ) if process.returncode != 0 or not process.stdout: error = process.stderr.decode("utf-8", errors="replace")[-500:] raise ValueError(f"Unable to resize video frame: {error}") image = process.stdout return base64.b64encode(image).decode("ascii") def _parse_json(content: str) -> dict[str, Any]: content = content.strip() fenced = re.search(r"```(?:json)?\s*(.*?)\s*```", content, re.DOTALL) if fenced: content = fenced.group(1) try: parsed = json.loads(content) return parsed if isinstance(parsed, dict) else {"result": parsed, "structured": False} except json.JSONDecodeError: return {"summary": content, "structured": False} def _fallback_reason( result: dict[str, Any], confidence_threshold: float, detailed_output: bool = True ) -> str | None: if result.get("structured") is False: return "FAST_RESULT_NOT_STRUCTURED" required = [ "passed", "evidenceSufficient", "confidence", "conclusion", "summary", ] if detailed_output: required.extend(("events", "warnings")) missing = [field for field in required if field not in result] if missing: return "FAST_RESULT_MISSING_FIELDS:" + ",".join(missing) if not isinstance(result.get("conclusion"), str) or not result["conclusion"].strip(): return "FAST_RESULT_EMPTY_CONCLUSION" if not isinstance(result.get("summary"), str) or not result["summary"].strip(): return "FAST_RESULT_EMPTY_SUMMARY" if detailed_output and ( not isinstance(result.get("events"), list) or not isinstance(result.get("warnings"), list) ): return "FAST_RESULT_INVALID_COLLECTIONS" if result.get("evidenceSufficient") is not True: return "FAST_RESULT_INSUFFICIENT" if result.get("passed") is not True and result.get("passed") is not False: return "FAST_RESULT_INVALID_DECISION" try: confidence = float(result.get("confidence")) except (TypeError, ValueError): return "FAST_RESULT_INVALID_CONFIDENCE" if confidence < confidence_threshold: return "FAST_RESULT_LOW_CONFIDENCE" return None def _bounded_number( value: Any, name: str, minimum: float, maximum: float, integer: bool = False ) -> int | float: try: number = float(value) except (TypeError, ValueError) as exc: raise ValueError(f"Video tuning parameter {name} must be numeric") from exc if number < minimum or number > maximum: raise ValueError( f"Video tuning parameter {name} must be between {minimum} and {maximum}" ) if integer and not number.is_integer(): raise ValueError(f"Video tuning parameter {name} must be an integer") return int(number) if integer else number def _normalize_tuning(raw: dict[str, Any] | None) -> dict[str, Any]: tuning = dict(raw or {}) normalized: dict[str, Any] = {} if "sampleFps" in tuning: normalized["sampleFps"] = _bounded_number( tuning["sampleFps"], "sampleFps", *VIDEO_TUNING_LIMITS["sampleFps"] ) if "maxWidth" in tuning: normalized["maxWidth"] = _bounded_number( tuning["maxWidth"], "maxWidth", *VIDEO_TUNING_LIMITS["maxWidth"], True ) if "confidenceThreshold" in tuning: normalized["confidenceThreshold"] = _bounded_number( tuning["confidenceThreshold"], "confidenceThreshold", *VIDEO_TUNING_LIMITS["confidenceThreshold"], ) if "fallbackToAccurate" in tuning: value = tuning["fallbackToAccurate"] if not isinstance(value, bool): raise ValueError("Video tuning parameter fallbackToAccurate must be boolean") normalized["fallbackToAccurate"] = value return normalized def _apply_tuning(strategy: dict[str, Any], tuning: dict[str, Any]) -> dict[str, Any]: configured = dict(strategy) for name in ("sampleFps", "maxWidth"): if name in tuning: configured[name] = tuning[name] return configured _CRITERION_TRANSITIONS = ( ("消失", ("消失", "不再显示", "从有到无", "由有变无")), ("熄灭", ("熄灭", "从亮到灭", "由亮变灭")), ("关闭", ("已关闭", "从开到关", "由开变关")), ("停止", ("已停止", "停止运行")), ("断开", ("已断开", "连接断开")), ("出现", ("出现", "从无到有", "由无变有")), ("点亮", ("点亮", "亮起", "从灭到亮", "由灭变亮")), ("开启", ("已开启", "从关到开", "由关变开")), ("启动", ("已启动", "开始运行")), ("连接", ("已连接", "连接成功")), ) _REMOVAL_CRITERION_TERMS = ( "消失", "熄灭", "关闭", "停止", "断开", "移除", "不再显示", "disappear", "turn off", "shut down", "disconnect", "remove", "absent", ) _APPEARANCE_CRITERION_TERMS = ( "出现", "点亮", "开启", "启动", "连接", "显示", "appear", "turn on", "start", "connect", "visible", ) def _temporal_acceptance_guidance(acceptance_criterion: str) -> str: original_criterion = str(acceptance_criterion or "").strip() criterion = original_criterion.lower() guidance = ( f"当前通过标准是【{original_criterion}】。图片严格按时间从前到后排列," "第一张是初始状态,最后一张是结束状态。判定 passed 前必须逐帧检查," "并明确比较第一张与最后一张图片。" ) if any(term in criterion for term in _REMOVAL_CRITERION_TERMS): target = original_criterion for term in _REMOVAL_CRITERION_TERMS: index = criterion.find(term) if index >= 0: target = original_criterion[:index] break target = re.sub(r"(?:是否|有没有|有无)$", "", target).strip(" ::,,。??") or "目标" return guidance + ( f"针对当前消失、移除、关闭或断开的标准:如果{target}在前段图片中存在," f"但在最后一张及末段图片中已不再显示{target},就说明状态变化成功," "passed必须为true。不能因为第一张图片中存在,就判断全程持续存在。" ) if any(term in criterion for term in _APPEARANCE_CRITERION_TERMS): return guidance + ( "针对出现、点亮、开启、启动或连接的标准:如果目标在前段图片中不存在," "但在最后一张及末段图片中已出现,就说明状态变化成功,passed 必须为 true。" ) return guidance def _needs_temporal_removal_verification( acceptance_criterion: str, result: dict[str, Any] ) -> bool: criterion = str(acceptance_criterion or "").strip().lower() return ( any(term in criterion for term in _REMOVAL_CRITERION_TERMS) and result.get("passed") is False and result.get("evidenceSufficient") is True ) def _criterion_verdict(result: dict[str, Any], acceptance_criterion: str) -> bool | None: criterion = str(acceptance_criterion or "").strip() if not criterion: return None conclusion = str(result.get("conclusion") or "").strip() summary = str(result.get("summary") or "").strip() events = result.get("events") if isinstance(result.get("events"), list) else [] evidence_text = " ".join( [conclusion, summary] + [ str(event.get("event") if isinstance(event, dict) else event) for event in events ] ) explicit_failure = re.search( r"(?:不符合|未满足|不满足|未达到|没有达到|判定失败|未通过).{0,12}(?:标准|要求|条件)?", evidence_text, ) if explicit_failure: return False explicit_success = re.search( r"(?:符合|满足|达到|已完成).{0,12}(?:标准|要求|条件)", evidence_text, ) if explicit_success: return True for criterion_word, observed_phrases in _CRITERION_TRANSITIONS: if criterion_word not in criterion: continue negated = re.search( rf"(?:未|没有|尚未|未能|不曾).{{0,4}}{re.escape(criterion_word)}", evidence_text, ) if negated: return False if criterion_word in {"消失", "熄灭", "关闭", "停止", "断开"} and re.search( r"(?:仍然|仍|依然|依旧).{0,8}(?:显示|存在|亮起|开启|运行|连接)", evidence_text, ): return False if any(phrase in evidence_text for phrase in observed_phrases): return True return None def _normalize_decision( result: dict[str, Any], confidence_threshold: float, acceptance_criterion: str = "" ) -> dict[str, Any]: normalized = dict(result or {}) warnings = normalized.get("warnings") if not isinstance(warnings, list): warnings = [] normalized["warnings"] = warnings normalized["events"] = normalized.get("events") if isinstance(normalized.get("events"), list) else [] try: confidence = min(1.0, max(0.0, float(normalized.get("confidence", 0)))) except (TypeError, ValueError): confidence = 0.0 evidence_sufficient = normalized.get("evidenceSufficient") is True model_passed = normalized.get("passed") is True criterion_verdict = _criterion_verdict(normalized, acceptance_criterion) raw_passed = model_passed if criterion_verdict is None else criterion_verdict decision_reason = "MODEL_DECISION" if normalized.get("structured") is False: evidence_sufficient = False decision_reason = "UNSTRUCTURED_RESULT" elif not evidence_sufficient: decision_reason = "INSUFFICIENT_EVIDENCE" elif confidence < confidence_threshold: decision_reason = "LOW_CONFIDENCE" elif criterion_verdict is not None and criterion_verdict != model_passed: decision_reason = "CRITERION_EVIDENCE_CORRECTION" passed = raw_passed and evidence_sufficient and confidence >= confidence_threshold if not passed and decision_reason != "MODEL_DECISION": if decision_reason == "LOW_CONFIDENCE": reason_text = "分析置信度不足" elif decision_reason == "CRITERION_EVIDENCE_CORRECTION": reason_text = "分析结论未满足通过标准" else: reason_text = "视频证据不足" conclusion = str(normalized.get("conclusion") or "").strip() normalized["conclusion"] = f"未通过:{reason_text}" + (f";{conclusion}" if conclusion else "") if reason_text not in warnings: warnings.append(reason_text) normalized["passed"] = passed if criterion_verdict is not None and criterion_verdict != model_passed: normalized["modelPassed"] = model_passed normalized["evidenceSufficient"] = evidence_sufficient normalized["confidence"] = confidence normalized["decisionReason"] = decision_reason normalized["confidenceThreshold"] = confidence_threshold normalized.setdefault("conclusion", "通过" if passed else "未通过") normalized.setdefault("summary", normalized["conclusion"]) normalized["result"] = normalized["summary"] normalized.pop("structured", None) return normalized def _strategy(profile: Profile, name: str) -> dict[str, Any]: defaults = { "FAST": { "model": profile.config.get("model", "Qwen/Qwen3.8-27B-FP8"), "provider": "sglang", "sampleFps": 1.0, "maxWidth": 896, "maxOutputTokens": 512, "numCtx": 16384, "emptyResponseRetries": 0, "maxRetryOutputTokens": 4096, }, "ACCURATE": { "model": profile.config.get("model", "Qwen/Qwen3.8-27B"), "provider": "sglang", "sampleFps": 3.0, "maxWidth": 1280, "maxOutputTokens": 768, "numCtx": 32768, "emptyResponseRetries": 2, "maxRetryOutputTokens": 8192, }, } configured = profile.config.get("strategies", {}).get(name.lower(), {}) strategy = {**defaults[name], **configured} strategy["provider"] = str(strategy.get("provider", "sglang")).strip().lower() strategy["baseUrl"] = ( settings.sglang_fast_base_url if name == "FAST" else settings.sglang_accurate_base_url ) return strategy def _sglang_payload(payload: dict[str, Any]) -> dict[str, Any]: message = payload.get("messages", [{}])[0] content: list[dict[str, Any]] = [ {"type": "text", "text": str(message.get("content") or "")} ] content.extend( { "type": "image_url", "image_url": {"url": f"data:image/jpeg;base64,{image}"}, } for image in message.get("images") or [] ) options = payload.get("options") or {} result = { "model": payload["model"], "messages": [{"role": "user", "content": content}], "stream": False, "temperature": float(options.get("temperature", 0.1)), "max_tokens": int(options.get("num_predict", 512)), "chat_template_kwargs": { "enable_thinking": bool(payload.get("think", False)), "preserve_thinking": False, }, } if payload.get("format"): result["response_format"] = { "type": "json_schema", "json_schema": { "name": "media_analysis", "strict": True, "schema": payload["format"], }, } return result def _request_image_count(payload: dict[str, Any], provider: str) -> int: message = payload.get("messages", [{}])[0] content = message.get("content") or [] return sum( 1 for item in content if isinstance(item, dict) and item.get("type") == "image_url" ) async def _request_vision_model( payload: dict[str, Any], empty_response_retries: int = 1, max_retry_output_tokens: int = 8192, provider: str = "sglang", base_url: str | None = None, ) -> tuple[str, dict[str, Any]]: provider = str(provider or "sglang").strip().lower() if provider != "sglang": raise ValueError("Only the SGLang vision model provider is supported") if not str(base_url or "").strip(): raise VisionModelError("SGLang vision model endpoint is not configured") last_metadata: dict[str, Any] = {} timeout = httpx.Timeout(settings.sglang_timeout_seconds, connect=5.0) async with httpx.AsyncClient(timeout=timeout) as client: for attempt in range(empty_response_retries + 1): request_payload = copy.deepcopy(payload) if attempt: options = request_payload.setdefault("options", {}) previous_limit = max( 1, int(options.get("num_predict", 0)), int(last_metadata.get("requestedOutputTokens") or 0), ) previous_eval_count = int(last_metadata.get("evalCount") or 0) retry_limit = max(4096, previous_limit * 2, previous_eval_count + 2048) options["num_predict"] = min(max_retry_output_tokens, retry_limit) request_payload["think"] = False request_payload["messages"][0]["content"] += ( "\n\n/no_think\nReturn only the final JSON now. Do not output reasoning." if request_payload.get("format") else "\n\n/no_think\nReturn only the final answer requested by the user." ) provider_payload = _sglang_payload(request_payload) url = f"{str(base_url).rstrip('/')}/chat/completions" headers = {} if settings.sglang_api_key: headers["Authorization"] = f"Bearer {settings.sglang_api_key}" response = await client.post(url, json=provider_payload, headers=headers) try: response.raise_for_status() except httpx.HTTPStatusError as exc: detail = response.text.strip()[:1000] raise VisionModelError( "Vision model request rejected " f"(provider={provider}, status={response.status_code}, " f"imageCount={_request_image_count(provider_payload, provider)}, " f"numCtx={request_payload.get('options', {}).get('num_ctx')}, " f"detail={detail})" ) from exc body = response.json() choices = body.get("choices") or [] choice = choices[0] if choices else {} message = choice.get("message") or {} usage = body.get("usage") or {} content = str(message.get("content") or "").strip() metadata = { "doneReason": choice.get("finish_reason"), "promptEvalCount": usage.get("prompt_tokens"), "evalCount": usage.get("completion_tokens"), "totalDurationNs": None, "thinkingLength": len(str(message.get("reasoning_content") or "")), "attempt": attempt + 1, "requestedOutputTokens": provider_payload.get("max_tokens"), } metadata["provider"] = provider truncated_structured_result = ( bool(request_payload.get("format")) and str(metadata.get("doneReason") or "").lower() == "length" and _parse_json(content).get("structured") is False ) if content and truncated_structured_result and attempt < empty_response_retries: last_metadata = metadata continue if content: return content, metadata last_metadata = metadata raise VisionModelError( "Vision model returned an empty final response after retry " f"({json.dumps(last_metadata, ensure_ascii=False)})" ) async def _analyze_frames( profile: Profile, strategy_name: str, frames: list[Path], duration: float, extracted_width: int, prompt: str, tuning: dict[str, Any] | None = None, detailed_output: bool = False, ) -> tuple[dict[str, Any], dict[str, Any]]: strategy = _apply_tuning(_strategy(profile, strategy_name), tuning or {}) model_frame_limit = _model_frame_limit() desired_count, _, _, _, sampling_capped = _frame_sampling_plan( duration, float(strategy["sampleFps"]), model_frame_limit, ) selected = _select_frames(frames, desired_count) maximum_width = max(320, int(strategy["maxWidth"])) images = [_encode_image(path, maximum_width, extracted_width) for path in selected] strategy_prompt = prompt + ( f"\n\nThe video duration is {duration:.3f} seconds. " f"The following {len(selected)} images are ordered frames sampled across the video." ) if detailed_output: strategy_prompt += ( "\n\nDetailed output is enabled. Include a concise event timeline and warnings." ) else: strategy_prompt += ( "\n\nCompact output is enabled. Return only passed, evidenceSufficient, " "confidence, conclusion, and a short summary. Do not return events, warnings, " "frame-by-frame descriptions, or repeated evidence. Keep conclusion under 20 " "Chinese characters and summary under 100 Chinese characters." ) think = bool(strategy.get("think", profile.config.get("think", False))) if not think: strategy_prompt += "\n\n/no_think\nReturn only the requested JSON object without reasoning." token_key = "detailedMaxOutputTokens" if detailed_output else "compactMaxOutputTokens" output_tokens = int(strategy.get(token_key, strategy["maxOutputTokens"])) context_tokens = _video_context_size( int(strategy["numCtx"]), len(selected), output_tokens ) payload = { "model": strategy["model"], "stream": False, "think": think, "format": VIDEO_RESULT_SCHEMA if detailed_output else VIDEO_COMPACT_RESULT_SCHEMA, "keep_alive": str(strategy.get("keepAlive", "30m")), "messages": [{"role": "user", "content": strategy_prompt, "images": images}], "options": { "temperature": float(strategy.get("temperature", profile.config.get("temperature", 0.1))), "num_predict": output_tokens, "num_ctx": context_tokens, }, } provider = str(strategy.get("provider", "sglang")) content, metrics = await _request_vision_model( payload, int(strategy.get("emptyResponseRetries", 0)), int(strategy.get("maxRetryOutputTokens", 8192)), provider, strategy.get("baseUrl"), ) metadata = { "mode": strategy_name, "model": str(strategy["model"]), "provider": metrics.get("provider", provider), "primaryProvider": provider, "primaryModel": strategy["model"], "providerFallback": False, "sampledFrameCount": len(selected), "maximumWidth": maximum_width, "effectiveSampleFps": round(max(0.0, (len(selected) - 1) / duration), 4), "samplingFrameLimit": model_frame_limit, "configuredSamplingFrameLimit": settings.video_sampling_frame_limit, "samplingCapped": sampling_capped, "numCtx": context_tokens, "detailedOutput": detailed_output, **metrics, } return _parse_json(content), metadata async def analyze_video( profile: Profile, media_path: Path, extra_instruction: str = "", analysis_mode: str = "AUTO", tuning: dict[str, Any] | None = None, decision_policy: str = "FAIL_CLOSED", detailed_output: bool = False, ) -> tuple[dict[str, Any], float, dict[str, Any]]: requested_mode = str(analysis_mode or "AUTO").strip().upper() if requested_mode not in {"AUTO", "FAST", "ACCURATE"}: raise ValueError("Video analysis mode must be AUTO, FAST, or ACCURATE") if str(decision_policy or "FAIL_CLOSED").strip().upper() != "FAIL_CLOSED": raise ValueError("Video decision policy must be FAIL_CLOSED") effective_tuning = _normalize_tuning(tuning) confidence_threshold = float(effective_tuning.get("confidenceThreshold", 0.75)) fallback_to_accurate = bool(effective_tuning.get("fallbackToAccurate", True)) fast = _apply_tuning(_strategy(profile, "FAST"), effective_tuning) accurate = _apply_tuning(_strategy(profile, "ACCURATE"), effective_tuning) extraction_strategy = fast if requested_mode == "FAST" else accurate extraction_sample_fps = float(extraction_strategy["sampleFps"]) extraction_width = max(320, int(extraction_strategy["maxWidth"])) model_frame_limit = _model_frame_limit() frames, duration, frame_dir = _extract_frames( media_path, extraction_sample_fps, model_frame_limit, extraction_width, ) try: prompt_path = profile.directory / str(profile.config.get("promptFile", "prompt.txt")) prompt = prompt_path.read_text(encoding="utf-8") if extra_instruction.strip(): temporal_guidance = _temporal_acceptance_guidance(extra_instruction) prompt += ( "\n\nAcceptance criterion (passed=true exactly when this criterion is " "satisfied; this is not merely an object-detection question):\n" + "通过标准:" + extra_instruction.strip() + "。" + temporal_guidance ) async def finalize_fast_result( result: dict[str, Any], metadata: dict[str, Any], mode_label: str, fallback_reason: str | None = None, ) -> tuple[dict[str, Any], dict[str, Any]]: normalized = _normalize_decision( result, confidence_threshold, extra_instruction ) if _needs_temporal_removal_verification(extra_instruction, normalized): focused_prompt = prompt + ( "\n\n首尾帧精确复核:本次只提供视频第一张和最后一张图片。" "必须比较具体目标在两张图片中的存在状态,再按通过标准输出结论。" ) verified, verified_metadata = await _analyze_frames( profile, "ACCURATE", [frames[0], frames[-1]], duration, extraction_width, focused_prompt, effective_tuning, detailed_output, ) verified = _normalize_decision( verified, confidence_threshold, extra_instruction ) verified_metadata.update( { "requestedMode": mode_label, "fallback": True, "fallbackReason": "TEMPORAL_REMOVAL_VERIFICATION", "tuning": effective_tuning, } ) return verified, verified_metadata metadata.update( { "requestedMode": mode_label, "fallback": False, "fallbackReason": fallback_reason, "tuning": effective_tuning, } ) return normalized, metadata if requested_mode in {"FAST", "ACCURATE"}: result, metadata = await _analyze_frames( profile, requested_mode, frames, duration, extraction_width, prompt, effective_tuning, detailed_output, ) if requested_mode == "FAST": result, metadata = await finalize_fast_result( result, metadata, requested_mode ) else: result = _normalize_decision( result, confidence_threshold, extra_instruction ) metadata.update( { "requestedMode": requested_mode, "fallback": False, "fallbackReason": None, "tuning": effective_tuning, } ) return result, duration, metadata fallback_reason = None try: result, metadata = await _analyze_frames( profile, "FAST", frames, duration, extraction_width, prompt, effective_tuning, detailed_output, ) fallback_reason = _fallback_reason( result, confidence_threshold, detailed_output ) if fallback_reason is None or not fallback_to_accurate: result, metadata = await finalize_fast_result( result, metadata, "AUTO", fallback_reason ) return result, duration, metadata except (VisionModelError, httpx.HTTPError) as exc: fallback_reason = f"FAST_MODEL_ERROR:{type(exc).__name__}" result, metadata = await _analyze_frames( profile, "ACCURATE", frames, duration, extraction_width, prompt, effective_tuning, detailed_output, ) result = _normalize_decision(result, confidence_threshold, extra_instruction) metadata.update( { "requestedMode": "AUTO", "fallback": True, "fallbackReason": fallback_reason, "tuning": effective_tuning, } ) return result, duration, metadata finally: for frame in frames: frame.unlink(missing_ok=True) frame_dir.rmdir()