open-webui/backend/open_webui/routers/tools.py
Timothy Jaeryang Baek e6476928ee refac
2026-10-09 18:37:55 +04:00

1120 lines
35 KiB
Python

from __future__ import annotations
import asyncio
import logging
import re
import time
from pathlib import Path
from typing import Optional
import aiohttp
from fastapi import APIRouter, Depends, HTTPException, Request, status
from open_webui.config import BYPASS_ADMIN_ACCESS_CONTROL, CACHE_DIR
from open_webui.constants import ERROR_MESSAGES
from open_webui.env import (
AIOHTTP_CLIENT_SESSION_SSL,
AIOHTTP_CLIENT_TIMEOUT,
ENABLE_TOOL_SERVERS,
ENABLE_TOOLS,
)
from open_webui.events import EVENTS, publish_event
from open_webui.internal.db import get_async_session
from open_webui.models.access_grants import AccessGrants
from open_webui.models.config import Config
from open_webui.models.groups import Groups
from open_webui.models.oauth_sessions import OAuthSessions
from open_webui.models.tool_history import ToolHistories, tool_diff
from open_webui.models.tools import (
ToolAccessResponse,
ToolForm,
ToolModel,
ToolResponse,
Tools,
ToolUserResponse,
)
from open_webui.utils.access_control import (
filter_allowed_access_grants,
has_connection_access,
has_permission,
)
from open_webui.utils.auth import get_admin_user, get_verified_user
from open_webui.utils.plugin import (
get_tool_contents_cache,
get_tool_module_from_cache,
get_tools_cache,
load_tool_module_by_id,
replace_imports,
resolve_valves_schema_options,
set_tool_module_in_cache,
)
from open_webui.utils.tools import connect_mcp_server, get_tool_servers
from open_webui.utils.tools import get_tool_specs as get_local_tool_specs
from pydantic import BaseModel, HttpUrl
from sqlalchemy.ext.asyncio import AsyncSession
log = logging.getLogger(__name__)
router = APIRouter()
async def get_tool_module(request, tool_id, load_from_db=True):
"""
Get the tool module by its ID.
"""
tool_module, _ = await get_tool_module_from_cache(request, tool_id, load_from_db)
return tool_module
############################
# GetTools
# The danger is not in having tools, but in reaching
# for the wrong one. Let the choice here be deliberate.
############################
@router.get('/', response_model=list[ToolUserResponse])
async def get_tools(
request: Request,
query: Optional[str] = None,
user=Depends(get_verified_user),
db: AsyncSession = Depends(get_async_session),
):
tools = []
bypass_access_control = user.role == 'admin' and BYPASS_ADMIN_ACCESS_CONTROL
user_group_ids = (
set()
if bypass_access_control
else {group.id for group in await Groups.get_groups_by_member_id(user.id, db=db, include_inherited=True)}
)
# Local Tools
if ENABLE_TOOLS:
tools_cache = get_tools_cache(request)
for tool in await Tools.get_tools(
defer_content=True,
db=db,
user_id=None if bypass_access_control else user.id,
user_group_ids=user_group_ids,
):
tool_module = tools_cache.get(tool.id)
has_user_valves = (
hasattr(tool_module, 'UserValves')
if tool_module
else (tool.meta.has_user_valves if tool.meta else False)
)
tools.append(
ToolUserResponse(
**{
**tool.model_dump(),
'has_user_valves': has_user_valves,
}
)
)
# OpenAPI Tool Servers
server_connections = {}
for server in await get_tool_servers(request):
server_idx = server.get('idx', 0)
connections = await Config.get('tool_server.connections', [])
if server_idx >= len(connections):
log.warning(
f'Tool server index {server_idx} out of range '
f'(have {len(connections)} connections), skipping server {server.get("id")}'
)
continue
connection = connections[server_idx]
server_id = f'server:{server.get("id")}'
server_connections[server_id] = connection
tools.append(
ToolUserResponse(
**{
'id': server_id,
'user_id': server_id,
'name': server.get('openapi', {}).get('info', {}).get('title', 'Tool Server'),
'meta': {
'description': server.get('openapi', {}).get('info', {}).get('description', ''),
},
'updated_at': int(time.time()),
'created_at': int(time.time()),
}
)
)
# MCP Tool Servers
for server in (await Config.get('tool_server.connections', [])) if ENABLE_TOOL_SERVERS else []:
if server.get('type', 'openapi') == 'mcp' and (server.get('config') or {}).get('enable'):
info = server.get('info') or {}
server_id = info.get('id')
auth_type = server.get('auth_type', 'none')
session_token = None
if auth_type in ('oauth_2.1', 'oauth_2.1_static') and server_id:
splits = server_id.split(':')
server_id = splits[-1] if len(splits) > 1 else server_id
session_token = await request.app.state.oauth_client_manager.get_oauth_token(
user.id, f'mcp:{server_id}'
)
tool_id = f'server:mcp:{info.get("id")}'
server_connections[tool_id] = server
tools.append(
ToolUserResponse(
**{
'id': tool_id,
'user_id': tool_id,
'name': info.get('name', 'MCP Tool Server'),
'meta': {
'description': info.get('description', ''),
},
'updated_at': int(time.time()),
'created_at': int(time.time()),
**(
{
'authenticated': session_token is not None,
}
if auth_type in ('oauth_2.1', 'oauth_2.1_static')
else {}
),
}
)
)
if not bypass_access_control:
tools = [
tool
for tool in tools
if not str(tool.id).startswith('server:')
or await has_connection_access(
user,
server_connections[str(tool.id)],
user_group_ids,
)
]
if query:
q = query.casefold()
tools = [tool for tool in tools if q in (tool.name or '').casefold()]
return tools
@router.get('/id/{id}/specs')
async def get_tool_specs(request: Request, id: str, user=Depends(get_verified_user)):
"""Discover tools for an accessible connection. Currently supports MCP."""
if not id.startswith('server:mcp:'):
raise HTTPException(status_code=404, detail='Tool not found')
try:
# Keep connect, discovery and cleanup in one task for the MCP transport.
async with asyncio.timeout(15):
result = await connect_mcp_server(request, id.removeprefix('server:mcp:'), user, {})
if result is None:
raise HTTPException(status_code=404, detail='Tool not found')
client, specs = result
try:
return {'specs': [{'name': spec['name'], 'description': spec.get('description', '')} for spec in specs]}
finally:
await client.disconnect()
except HTTPException:
raise
except TimeoutError:
raise HTTPException(status_code=504, detail='Tool discovery timed out')
except Exception:
log.exception('Failed to discover tool specs')
raise HTTPException(status_code=502, detail='Unable to load tools')
############################
# GetToolList
############################
@router.get('/list', response_model=list[ToolAccessResponse])
async def get_tool_list(user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session)):
if not ENABLE_TOOLS:
return []
bypass_access_control = user.role == 'admin' and BYPASS_ADMIN_ACCESS_CONTROL
user_group_ids = (
set()
if bypass_access_control
else {group.id for group in await Groups.get_groups_by_member_id(user.id, db=db, include_inherited=True)}
)
tools = await Tools.get_tools(
defer_content=True,
db=db,
user_id=None if bypass_access_control else user.id,
user_group_ids=user_group_ids,
)
result = []
for tool in tools:
has_write = (
bypass_access_control
or user.id == tool.user_id
or any(
g.permission == 'write'
and (
(g.principal_type == 'user' and (g.principal_id == user.id or g.principal_id == '*'))
or (g.principal_type == 'group' and g.principal_id in user_group_ids)
)
for g in tool.access_grants
)
)
result.append(
ToolAccessResponse(
**tool.model_dump(),
write_access=has_write,
)
)
return result
############################
# LoadFunctionFromLink
############################
class LoadUrlForm(BaseModel):
url: HttpUrl
def github_url_to_raw_url(url: str) -> str:
# Handle 'tree' (folder) URLs (add main.py at the end)
m1 = re.match(r'https://github\.com/([^/]+)/([^/]+)/tree/([^/]+)/(.*)', url)
if m1:
org, repo, branch, path = m1.groups()
return f'https://raw.githubusercontent.com/{org}/{repo}/refs/heads/{branch}/{path.rstrip("/")}/main.py'
# Handle 'blob' (file) URLs
m2 = re.match(r'https://github\.com/([^/]+)/([^/]+)/blob/([^/]+)/(.*)', url)
if m2:
org, repo, branch, path = m2.groups()
return f'https://raw.githubusercontent.com/{org}/{repo}/refs/heads/{branch}/{path}'
# No match; return as-is
return url
@router.post('/load/url', response_model=dict | None)
async def load_tool_from_url(request: Request, form_data: LoadUrlForm, user=Depends(get_admin_user)):
# NOTE: This is NOT a SSRF vulnerability:
# This endpoint is admin-only (see get_admin_user), meant for *trusted* internal use,
# and does NOT accept untrusted user input. Access is enforced by authentication.
url = str(form_data.url)
if not url:
raise HTTPException(status_code=400, detail='Please enter a valid URL')
url = github_url_to_raw_url(url)
url_parts = url.rstrip('/').split('/')
file_name = url_parts[-1]
tool_name = (
file_name[:-3]
if (file_name.endswith('.py') and (not file_name.startswith(('main.py', 'index.py', '__init__.py'))))
else url_parts[-2]
if len(url_parts) > 1
else 'function'
)
try:
async with aiohttp.ClientSession(
trust_env=True, timeout=aiohttp.ClientTimeout(total=AIOHTTP_CLIENT_TIMEOUT)
) as session:
async with session.get(
url, headers={'Content-Type': 'application/json'}, ssl=AIOHTTP_CLIENT_SESSION_SSL
) as resp:
if resp.status != 200:
raise HTTPException(status_code=resp.status, detail='Failed to fetch the tool')
data = await resp.text()
if not data:
raise HTTPException(status_code=400, detail='No data received from the URL')
return {
'name': tool_name,
'content': data,
}
except HTTPException:
raise
except Exception as e:
raise HTTPException(
status_code=500,
detail=ERROR_MESSAGES.DEFAULT(e, 'Error fetching tool'),
)
############################
# ExportTools
############################
@router.get('/export', response_model=list[ToolModel])
async def export_tools(
request: Request,
user=Depends(get_verified_user),
db: AsyncSession = Depends(get_async_session),
):
if user.role != 'admin' and not await has_permission(
user.id,
'workspace.tools_export',
await Config.get('user.permissions'),
db=db,
):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.UNAUTHORIZED,
)
bypass_access_control = user.role == 'admin' and BYPASS_ADMIN_ACCESS_CONTROL
return await Tools.get_tools(
db=db,
user_id=None if bypass_access_control else user.id,
permission='write',
)
############################
# CreateNewTools
############################
@router.post('/create', response_model=ToolResponse | None)
async def create_new_tools(
request: Request,
form_data: ToolForm,
user=Depends(get_verified_user),
db: AsyncSession = Depends(get_async_session),
):
"""Create a new tool from user-supplied Python source code."""
if user.role != 'admin' and not (
await has_permission(user.id, 'workspace.tools', await Config.get('user.permissions'), db=db)
or await has_permission(
user.id,
'workspace.tools_import',
await Config.get('user.permissions'),
db=db,
)
):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.UNAUTHORIZED,
)
if not form_data.id.isidentifier():
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail='Only alphanumeric characters and underscores are allowed in the id',
)
form_data.id = form_data.id.lower()
tools = await Tools.get_tool_by_id(form_data.id, db=db)
if tools is None:
try:
form_data.access_grants = await filter_allowed_access_grants(
await Config.get('user.permissions'),
user.id,
user.role,
form_data.access_grants,
'sharing.public_tools',
)
form_data.content = replace_imports(form_data.content)
tool_module, frontmatter, source_module = await load_tool_module_by_id(
form_data.id, content=form_data.content
)
form_data.meta.manifest = frontmatter
form_data.meta.has_user_valves = hasattr(tool_module, 'UserValves')
specs = get_local_tool_specs(tool_module)
tools = await Tools.insert_new_tool(user.id, form_data, specs, db=db, module=tool_module)
tool_cache_dir = CACHE_DIR / 'tools' / form_data.id
tool_cache_dir.mkdir(parents=True, exist_ok=True)
if tools:
set_tool_module_in_cache(request, tools.id, tools.content, tool_module, source_module)
await publish_event(
request,
EVENTS.TOOL_CREATED,
actor=user,
subject_id=tools.id,
data={'name': tools.name},
)
return tools
else:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=ERROR_MESSAGES.DEFAULT('Error creating tools'),
)
except HTTPException:
raise
except Exception as e:
log.exception(f'Failed to load the tool by id {form_data.id}: {e}')
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=ERROR_MESSAGES.DEFAULT(e, 'Error creating tool'),
)
else:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=ERROR_MESSAGES.ID_TAKEN,
)
############################
# GetToolsById
############################
@router.get('/id/{id}', response_model=ToolAccessResponse | None)
async def get_tools_by_id(id: str, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session)):
tools = await Tools.get_tool_by_id(id, db=db)
if tools:
if (
user.role == 'admin'
or tools.user_id == user.id
or await AccessGrants.has_access(
user_id=user.id,
resource_type='tool',
resource_id=tools.id,
permission='read',
db=db,
)
):
write_access = (
(user.role == 'admin' and BYPASS_ADMIN_ACCESS_CONTROL)
or user.id == tools.user_id
or await AccessGrants.has_access(
user_id=user.id,
resource_type='tool',
resource_id=tools.id,
permission='write',
db=db,
)
)
data = tools.model_dump()
if not write_access:
# extra='allow' re-admits content from model_dump; source is writer-only
data.pop('content', None)
return ToolAccessResponse(**data, write_access=write_access)
else:
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
)
else:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=ERROR_MESSAGES.NOT_FOUND,
)
############################
# UpdateToolsById
############################
@router.post('/id/{id}/update', response_model=ToolModel | None)
async def update_tools_by_id(
request: Request,
id: str,
form_data: ToolForm,
user=Depends(get_verified_user),
db: AsyncSession = Depends(get_async_session),
):
return await _update_tool(request, id, form_data, user, db)
async def _update_tool(request, id, form_data, user, db, version_id=None):
"""Update an existing tool's source code and metadata."""
tools = await Tools.get_tool_by_id(id, db=db)
if not tools:
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.NOT_FOUND,
)
# Is the user the original creator, in a group with write access, or an admin
if (
tools.user_id != user.id
and not await AccessGrants.has_access(
user_id=user.id,
resource_type='tool',
resource_id=tools.id,
permission='write',
db=db,
)
and user.role != 'admin'
):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.UNAUTHORIZED,
)
# Check again under the row lock when committing, in case Production changes meanwhile.
allow_code_changes = await can_change_tool_code(user, db)
if form_data.content != tools.content and not allow_code_changes:
raise HTTPException(401, ERROR_MESSAGES.UNAUTHORIZED)
try:
if version_id is None:
form_data.content = replace_imports(form_data.content)
tool_module, frontmatter, source_module = await load_tool_module_by_id(id, content=form_data.content)
form_data.meta.manifest = frontmatter
form_data.meta.has_user_valves = hasattr(tool_module, 'UserValves')
specs = get_local_tool_specs(tool_module)
if version_id is None:
form_data.access_grants = await filter_allowed_access_grants(
await Config.get('user.permissions'),
user.id,
user.role,
form_data.access_grants,
'sharing.public_tools',
)
updated = {
**form_data.model_dump(exclude={'id'}),
'specs': specs,
}
if version_id is not None:
updated.pop('access_grants', None)
tools = await Tools.update_tool_by_id(
id,
updated,
db=db,
user_id=user.id,
version_id=version_id,
module=tool_module,
allow_code_changes=allow_code_changes,
)
if tools:
set_tool_module_in_cache(request, tools.id, tools.content, tool_module, source_module)
await publish_event(
request,
EVENTS.TOOL_UPDATED,
actor=user,
subject_id=tools.id,
data={'name': tools.name},
)
return tools
else:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=ERROR_MESSAGES.DEFAULT('Error updating tools'),
)
except HTTPException:
raise
except Exception as e:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=str(e),
)
############################
# UpdateToolAccessById
############################
class ToolAccessGrantsForm(BaseModel):
access_grants: list[dict]
@router.post('/id/{id}/access/update', response_model=ToolModel | None)
async def update_tool_access_by_id(
request: Request,
id: str,
form_data: ToolAccessGrantsForm,
user=Depends(get_verified_user),
db: AsyncSession = Depends(get_async_session),
):
tools = await Tools.get_tool_by_id(id, db=db)
if not tools:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=ERROR_MESSAGES.NOT_FOUND,
)
if (
tools.user_id != user.id
and not await AccessGrants.has_access(
user_id=user.id,
resource_type='tool',
resource_id=tools.id,
permission='write',
db=db,
)
and user.role != 'admin'
):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.UNAUTHORIZED,
)
form_data.access_grants = await filter_allowed_access_grants(
await Config.get('user.permissions'),
user.id,
user.role,
form_data.access_grants,
'sharing.public_tools',
)
await AccessGrants.set_access_grants('tool', id, form_data.access_grants, db=db)
tools = await Tools.get_tool_by_id(id, db=db)
await publish_event(
request,
EVENTS.TOOL_ACCESS_UPDATED,
actor=user,
subject_id=id,
data={'name': tools.name if tools else None},
)
return tools
############################
# DeleteToolsById
############################
@router.delete('/id/{id}/delete', response_model=bool)
async def delete_tools_by_id(
request: Request,
id: str,
user=Depends(get_verified_user),
db: AsyncSession = Depends(get_async_session),
):
tools = await Tools.get_tool_by_id(id, db=db)
if not tools:
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.NOT_FOUND,
)
if (
tools.user_id != user.id
and not await AccessGrants.has_access(
user_id=user.id,
resource_type='tool',
resource_id=tools.id,
permission='write',
db=db,
)
and user.role != 'admin'
):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.UNAUTHORIZED,
)
result = await Tools.delete_tool_by_id(id, db=db)
if result:
TOOLS = get_tools_cache(request)
TOOLS.pop(id, None)
TOOL_CONTENTS = get_tool_contents_cache(request)
TOOL_CONTENTS.pop(id, None)
await publish_event(
request,
EVENTS.TOOL_DELETED,
actor=user,
subject_id=id,
data={'name': tools.name},
)
return result
############################
# GetToolValves
############################
@router.get('/id/{id}/valves', response_model=dict | None)
async def get_tools_valves_by_id(
id: str, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session)
):
tools = await Tools.get_tool_by_id(id, db=db)
if not tools:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=ERROR_MESSAGES.NOT_FOUND,
)
if (
tools.user_id != user.id
and not await AccessGrants.has_access(
user_id=user.id,
resource_type='tool',
resource_id=tools.id,
permission='write',
db=db,
)
and user.role != 'admin'
):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
)
try:
valves = await Tools.get_tool_valves_by_id(id, db=db)
return valves
except Exception as e:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=ERROR_MESSAGES.DEFAULT(e, 'Error getting tool valves'),
)
############################
# GetToolValvesSpec
############################
@router.get('/id/{id}/valves/spec', response_model=dict | None)
async def get_tools_valves_spec_by_id(
request: Request,
id: str,
user=Depends(get_verified_user),
db: AsyncSession = Depends(get_async_session),
):
tools = await Tools.get_tool_by_id(id, db=db)
if not tools:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=ERROR_MESSAGES.NOT_FOUND,
)
if (
tools.user_id != user.id
and not await AccessGrants.has_access(
user_id=user.id,
resource_type='tool',
resource_id=tools.id,
permission='write',
db=db,
)
and user.role != 'admin'
):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
)
tools_module, _ = await get_tool_module_from_cache(request, id)
if hasattr(tools_module, 'Valves'):
Valves = tools_module.Valves
schema = Valves.schema()
# Resolve dynamic options for select dropdowns
schema = resolve_valves_schema_options(Valves, schema, user)
return schema
return None
############################
# UpdateToolValves
############################
@router.post('/id/{id}/valves/update', response_model=dict | None)
async def update_tools_valves_by_id(
request: Request,
id: str,
form_data: dict,
user=Depends(get_verified_user),
db: AsyncSession = Depends(get_async_session),
):
tools = await Tools.get_tool_by_id(id, db=db)
if not tools:
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.NOT_FOUND,
)
if (
tools.user_id != user.id
and not await AccessGrants.has_access(
user_id=user.id,
resource_type='tool',
resource_id=tools.id,
permission='write',
db=db,
)
and user.role != 'admin'
):
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
)
tools_module, _ = await get_tool_module_from_cache(request, id)
if not hasattr(tools_module, 'Valves'):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.NOT_FOUND,
)
Valves = tools_module.Valves
try:
form_data = {k: v for k, v in form_data.items() if v is not None}
valves = Valves(**form_data)
valves_dict = valves.model_dump(exclude_unset=True)
await Tools.update_tool_valves_by_id(id, valves_dict, db=db)
await publish_event(
request,
EVENTS.TOOL_VALVES_UPDATED,
actor=user,
subject_id=id,
)
return valves_dict
except Exception as e:
log.exception(f'Failed to update tool valves by id {id}: {e}')
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=ERROR_MESSAGES.DEFAULT(e, 'Error updating tool valves'),
)
############################
# ToolUserValves
############################
@router.get('/id/{id}/valves/user', response_model=dict | None)
async def get_tools_user_valves_by_id(
id: str, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session)
):
tools = await Tools.get_tool_by_id(id, db=db)
if not tools:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=ERROR_MESSAGES.NOT_FOUND,
)
if (
tools.user_id != user.id
and not await AccessGrants.has_access(
user_id=user.id,
resource_type='tool',
resource_id=tools.id,
permission='read',
db=db,
)
and user.role != 'admin'
):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
)
try:
user_valves = await Tools.get_user_valves_by_id_and_user_id(id, user.id, db=db)
return user_valves
except Exception as e:
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=ERROR_MESSAGES.DEFAULT(e, 'Error getting tool user valves'),
)
@router.get('/id/{id}/valves/user/spec', response_model=dict | None)
async def get_tools_user_valves_spec_by_id(
request: Request,
id: str,
user=Depends(get_verified_user),
db: AsyncSession = Depends(get_async_session),
):
tools = await Tools.get_tool_by_id(id, db=db)
if not tools:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=ERROR_MESSAGES.NOT_FOUND,
)
if (
tools.user_id != user.id
and not await AccessGrants.has_access(
user_id=user.id,
resource_type='tool',
resource_id=tools.id,
permission='read',
db=db,
)
and user.role != 'admin'
):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
)
tools_module, _ = await get_tool_module_from_cache(request, id)
if hasattr(tools_module, 'UserValves'):
UserValves = tools_module.UserValves
schema = UserValves.schema()
# Resolve dynamic options for select dropdowns
schema = resolve_valves_schema_options(UserValves, schema, user)
return schema
return None
@router.post('/id/{id}/valves/user/update', response_model=dict | None)
async def update_tools_user_valves_by_id(
request: Request,
id: str,
form_data: dict,
user=Depends(get_verified_user),
db: AsyncSession = Depends(get_async_session),
):
tools = await Tools.get_tool_by_id(id, db=db)
if not tools:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=ERROR_MESSAGES.NOT_FOUND,
)
if (
tools.user_id != user.id
and not await AccessGrants.has_access(
user_id=user.id,
resource_type='tool',
resource_id=tools.id,
permission='read',
db=db,
)
and user.role != 'admin'
):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.ACCESS_PROHIBITED,
)
tools_module, _ = await get_tool_module_from_cache(request, id)
if hasattr(tools_module, 'UserValves'):
UserValves = tools_module.UserValves
try:
form_data = {k: v for k, v in form_data.items() if v is not None}
user_valves = UserValves(**form_data)
user_valves_dict = user_valves.model_dump(exclude_unset=True)
await Tools.update_user_valves_by_id_and_user_id(id, user.id, user_valves_dict, db=db)
await publish_event(
request,
EVENTS.TOOL_VALVES_UPDATED,
actor=user,
subject_id=id,
data={'scope': 'user'},
)
return user_valves_dict
except Exception as e:
log.exception(f'Failed to update user valves by id {id}: {e}')
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail=ERROR_MESSAGES.DEFAULT(e, 'Error updating tool user valves'),
)
else:
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail=ERROR_MESSAGES.NOT_FOUND,
)
async def can_change_tool_code(user, db):
return user.role == 'admin' or (
await has_permission(user.id, 'workspace.tools', await Config.get('user.permissions'), db=db)
or await has_permission(user.id, 'workspace.tools_import', await Config.get('user.permissions'), db=db)
)
async def require_tool_history_access(id, user, db):
resource = await Tools.get_tool_by_id(id, db=db)
if not resource:
raise HTTPException(404, 'Not found')
if not (
(user.role == 'admin' and BYPASS_ADMIN_ACCESS_CONTROL)
or resource.user_id == user.id
or await AccessGrants.has_access(
user_id=user.id, resource_type='tool', resource_id=id, permission='write', db=db
)
):
raise HTTPException(401, ERROR_MESSAGES.ACCESS_PROHIBITED)
return resource
async def require_tool_history_entry(id, history_id, db):
entry = await ToolHistories.get_history_by_id(id, history_id, db=db)
if not entry:
raise HTTPException(404, 'Version not found')
return entry
@router.get('/id/{id}/history')
async def get_tool_history(
id: str, page: int = 1, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session)
):
await require_tool_history_access(id, user, db)
return await ToolHistories.get_history_by_tool_id(id, page, db=db)
@router.get('/id/{id}/history/diff')
async def get_tool_history_diff(
id: str, from_id: str, to_id: str, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session)
):
await require_tool_history_access(id, user, db)
before = await require_tool_history_entry(id, from_id, db)
after = await require_tool_history_entry(id, to_id, db)
return tool_diff(before, after)
@router.get('/id/{id}/history/{history_id}')
async def get_tool_history_entry(
id: str, history_id: str, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session)
):
await require_tool_history_access(id, user, db)
return await require_tool_history_entry(id, history_id, db)
@router.delete('/id/{id}/history/{history_id}')
async def delete_tool_history_entry(
id: str, history_id: str, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session)
):
await require_tool_history_access(id, user, db)
if not await ToolHistories.delete_history_entry(id, history_id, db=db):
raise HTTPException(404, 'Version not found')
return True
class ToolVersionForm(BaseModel):
version_id: str
@router.post('/id/{id}/update/version', response_model=ToolModel)
async def set_tool_production(
request: Request,
id: str,
form_data: ToolVersionForm,
user=Depends(get_verified_user),
db: AsyncSession = Depends(get_async_session),
):
await require_tool_history_access(id, user, db)
entry = await require_tool_history_entry(id, form_data.version_id, db)
try:
saved = ToolForm(id=id, **entry.snapshot)
except ValueError as error:
raise HTTPException(400, str(error)) from error
return await _update_tool(request, id, saved, user, db, version_id=entry.id)