- 移除Ollama相关配置和回退机制,统一使用SGLang作为视觉模型提供商 - 添加IMAGE_ANALYSIS类型支持,允许对多张图片进行分析 - 实现视频分析的详细输出模式,支持紧凑和详细两种结果格式 - 更新环境变量配置,添加VIDEO_MODEL_FRAME_LIMIT和MAX_IMAGE_BYTES - 修改compose配置文件中的上下文长度和内存分配参数 - 重构视频采样逻辑,限制单次请求帧数以优化显存使用 - 更新API接口文档,添加mediaUrls参数和详细输出选项说明 - 添加图像分析相关的依赖库opencv-python-headless - 实现结构化JSON响应格式验证和重试机制
900 lines
34 KiB
Python
900 lines
34 KiB
Python
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()
|