Coverage for open_webui/routers/functions.py: 52%
262 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 05:07 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 05:07 +0000
1from __future__ import annotations
3import logging
4import os
5import re
6from pathlib import Path
7from typing import Optional
9import aiohttp
10from fastapi import APIRouter, Depends, HTTPException, Request, status
11from open_webui.config import CACHE_DIR
12from open_webui.constants import ERROR_MESSAGES
13from open_webui.env import AIOHTTP_CLIENT_SESSION_SSL, AIOHTTP_CLIENT_TIMEOUT, ENABLE_PLUGINS
14from open_webui.events import EVENTS, build_event, dispatch_event_functions, publish_event, schedule_webhook_dispatch
15from open_webui.internal.db import get_async_session
16from open_webui.models.functions import (
17 FunctionForm,
18 FunctionModel,
19 FunctionResponse,
20 Functions,
21 FunctionUserResponse,
22 FunctionWithValvesModel,
23)
24from open_webui.utils.auth import get_admin_user, get_verified_user
25from open_webui.utils.plugin import (
26 get_function_contents_cache,
27 get_functions_cache,
28 get_function_module_from_cache,
29 load_function_module_by_id,
30 replace_imports,
31 resolve_valves_schema_options,
32)
33from pydantic import BaseModel, HttpUrl
34from sqlalchemy.ext.asyncio import AsyncSession
36log = logging.getLogger(__name__)
39router = APIRouter()
41############################
42# GetFunctions
43# Our daily functions give us, and forgive us
44# our deprecated methods, as we refactor those who depend on us.
45############################
48@router.get('/', response_model=list[FunctionResponse])
49async def get_functions(user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session)):
50 if not ENABLE_PLUGINS: 50 ↛ 51line 50 didn't jump to line 51 because the condition on line 50 was never true
51 return []
53 return await Functions.get_functions(db=db)
56@router.get('/list', response_model=list[FunctionUserResponse])
57async def get_function_list(user=Depends(get_admin_user), db: AsyncSession = Depends(get_async_session)):
58 if not ENABLE_PLUGINS: 58 ↛ 59line 58 didn't jump to line 59 because the condition on line 58 was never true
59 return []
61 return await Functions.get_function_list(db=db)
64############################
65# ExportFunctions
66############################
69@router.get('/export', response_model=list[FunctionModel | FunctionWithValvesModel])
70async def get_functions(
71 include_valves: bool = False,
72 user=Depends(get_admin_user),
73 db: AsyncSession = Depends(get_async_session),
74):
75 if not ENABLE_PLUGINS: 75 ↛ 76line 75 didn't jump to line 76 because the condition on line 75 was never true
76 return []
78 return await Functions.get_functions(include_valves=include_valves, db=db)
81############################
82# LoadFunctionFromLink
83############################
86class LoadUrlForm(BaseModel):
87 url: HttpUrl
90def github_url_to_raw_url(url: str) -> str:
91 # Handle 'tree' (folder) URLs (add main.py at the end)
92 m1 = re.match(r'https://github\.com/([^/]+)/([^/]+)/tree/([^/]+)/(.*)', url)
93 if m1: 93 ↛ 94line 93 didn't jump to line 94 because the condition on line 93 was never true
94 org, repo, branch, path = m1.groups()
95 return f'https://raw.githubusercontent.com/{org}/{repo}/refs/heads/{branch}/{path.rstrip("/")}/main.py'
97 # Handle 'blob' (file) URLs
98 m2 = re.match(r'https://github\.com/([^/]+)/([^/]+)/blob/([^/]+)/(.*)', url)
99 if m2: 99 ↛ 100line 99 didn't jump to line 100 because the condition on line 99 was never true
100 org, repo, branch, path = m2.groups()
101 return f'https://raw.githubusercontent.com/{org}/{repo}/refs/heads/{branch}/{path}'
103 # No match; return as-is
104 return url
107@router.post('/load/url', response_model=dict | None)
108async def load_function_from_url(request: Request, form_data: LoadUrlForm, user=Depends(get_admin_user)):
109 # NOTE: This is NOT a SSRF vulnerability:
110 # This endpoint is admin-only (see get_admin_user), meant for *trusted* internal use,
111 # and does NOT accept untrusted user input. Access is enforced by authentication.
113 url = str(form_data.url)
114 if not url: 114 ↛ 115line 114 didn't jump to line 115 because the condition on line 114 was never true
115 raise HTTPException(status_code=400, detail='Please enter a valid URL')
117 url = github_url_to_raw_url(url)
118 url_parts = url.rstrip('/').split('/')
120 file_name = url_parts[-1]
121 function_name = (
122 file_name[:-3]
123 if (file_name.endswith('.py') and (not file_name.startswith(('main.py', 'index.py', '__init__.py'))))
124 else url_parts[-2]
125 if len(url_parts) > 1
126 else 'function'
127 )
129 try:
130 async with aiohttp.ClientSession(
131 trust_env=True, timeout=aiohttp.ClientTimeout(total=AIOHTTP_CLIENT_TIMEOUT)
132 ) as session:
133 async with session.get(
134 url, headers={'Content-Type': 'application/json'}, ssl=AIOHTTP_CLIENT_SESSION_SSL
135 ) as resp:
136 if resp.status != 200: 136 ↛ 137line 136 didn't jump to line 137 because the condition on line 136 was never true
137 raise HTTPException(status_code=resp.status, detail='Failed to fetch the function')
138 data = await resp.text()
139 if not data: 139 ↛ 141line 139 didn't jump to line 141
140 raise HTTPException(status_code=400, detail='No data received from the URL')
141 return {
142 'name': function_name,
143 'content': data,
144 }
145 except HTTPException:
146 raise
147 except Exception as e:
148 raise HTTPException(
149 status_code=500,
150 detail=ERROR_MESSAGES.DEFAULT(e, 'Error fetching function'),
151 )
154############################
155# SyncFunctions
156############################
159class SyncFunctionsForm(BaseModel):
160 functions: list[FunctionWithValvesModel] = []
163@router.post('/sync', response_model=list[FunctionWithValvesModel])
164async def sync_functions(
165 request: Request,
166 form_data: SyncFunctionsForm,
167 user=Depends(get_admin_user),
168 db: AsyncSession = Depends(get_async_session),
169):
170 try:
171 for function in form_data.functions:
172 function.content = replace_imports(function.content)
173 function_module, function_type, frontmatter = await load_function_module_by_id(
174 function.id,
175 content=function.content,
176 )
178 if hasattr(function_module, 'Valves') and function.valves:
179 Valves = function_module.Valves
180 try:
181 Valves(**{k: v for k, v in function.valves.items() if v is not None})
182 except Exception as e:
183 log.exception(f'Error validating valves for function {function.id}: {e}')
184 raise e
186 return await Functions.sync_functions(user.id, form_data.functions, db=db)
187 except Exception as e:
188 log.exception(f'Failed to load a function: {e}')
189 raise HTTPException(
190 status_code=status.HTTP_400_BAD_REQUEST,
191 detail=ERROR_MESSAGES.DEFAULT(e, 'Error loading function'),
192 )
195############################
196# CreateNewFunction
197############################
200@router.post('/create', response_model=FunctionResponse | None)
201async def create_new_function(
202 request: Request,
203 form_data: FunctionForm,
204 user=Depends(get_admin_user),
205 db: AsyncSession = Depends(get_async_session),
206):
207 if not form_data.id.isidentifier():
208 raise HTTPException(
209 status_code=status.HTTP_400_BAD_REQUEST,
210 detail='Only alphanumeric characters and underscores are allowed in the id',
211 )
213 form_data.id = form_data.id.lower()
215 function = await Functions.get_function_by_id(form_data.id, db=db)
216 if function is None: 216 ↛ 259line 216 didn't jump to line 259 because the condition on line 216 was always true
217 try:
218 form_data.content = replace_imports(form_data.content)
219 function_module, function_type, frontmatter = await load_function_module_by_id(
220 form_data.id,
221 content=form_data.content,
222 )
223 form_data.meta.manifest = frontmatter
225 FUNCTIONS = get_functions_cache(request)
226 FUNCTIONS[form_data.id] = function_module
228 function = await Functions.insert_new_function(user.id, function_type, form_data, db=db)
230 function_cache_dir = CACHE_DIR / 'functions' / form_data.id
231 function_cache_dir.mkdir(parents=True, exist_ok=True)
233 if function_type == 'filter' and getattr(function_module, 'toggle', None):
234 await Functions.update_function_metadata_by_id(form_data.id, {'toggle': True}, db=db)
236 if function:
237 await publish_event(
238 request,
239 EVENTS.FUNCTION_CREATED,
240 actor=user,
241 subject_id=function.id,
242 data={'type': function.type, 'name': function.name},
243 )
244 return function
245 else:
246 raise HTTPException(
247 status_code=status.HTTP_400_BAD_REQUEST,
248 detail=ERROR_MESSAGES.DEFAULT('Error creating function'),
249 )
250 except HTTPException:
251 raise
252 except Exception as e:
253 log.exception(f'Failed to create a new function: {e}')
254 raise HTTPException(
255 status_code=status.HTTP_400_BAD_REQUEST,
256 detail=ERROR_MESSAGES.DEFAULT(e, 'Error creating function'),
257 )
258 else:
259 raise HTTPException(
260 status_code=status.HTTP_400_BAD_REQUEST,
261 detail=ERROR_MESSAGES.ID_TAKEN,
262 )
265############################
266# GetFunctionById
267############################
270@router.get('/id/{id}', response_model=FunctionModel | None)
271async def get_function_by_id(id: str, user=Depends(get_admin_user), db: AsyncSession = Depends(get_async_session)):
272 function = await Functions.get_function_by_id(id, db=db)
274 if function: 274 ↛ 275line 274 didn't jump to line 275 because the condition on line 274 was never true
275 return function
276 else:
277 raise HTTPException(
278 status_code=status.HTTP_401_UNAUTHORIZED,
279 detail=ERROR_MESSAGES.NOT_FOUND,
280 )
283############################
284# ToggleFunctionById
285############################
288@router.post('/id/{id}/toggle', response_model=FunctionModel | None)
289async def toggle_function_by_id(
290 request: Request,
291 id: str,
292 user=Depends(get_admin_user),
293 db: AsyncSession = Depends(get_async_session),
294):
295 function = await Functions.get_function_by_id(id, db=db)
296 if function: 296 ↛ 297line 296 didn't jump to line 297 because the condition on line 296 was never true
297 lifecycle_event = build_event(
298 request,
299 EVENTS.FUNCTION_DISABLE_STARTED if function.is_active else EVENTS.FUNCTION_ENABLE_STARTED,
300 actor=user,
301 subject_id=function.id,
302 subject_type='function',
303 data={'type': function.type, 'name': function.name},
304 )
305 await dispatch_event_functions(
306 request.app,
307 lifecycle_event,
308 request=request,
309 extra_function_ids=[function.id] if not function.is_active else None,
310 )
311 schedule_webhook_dispatch(request.app, lifecycle_event)
313 function = await Functions.update_function_by_id(id, {'is_active': not function.is_active}, db=db)
315 if function:
316 await publish_event(
317 request,
318 EVENTS.FUNCTION_ENABLED if function.is_active else EVENTS.FUNCTION_DISABLED,
319 actor=user,
320 subject_id=function.id,
321 subject_type='function',
322 data={'type': function.type, 'name': function.name},
323 )
324 return function
325 else:
326 raise HTTPException(
327 status_code=status.HTTP_400_BAD_REQUEST,
328 detail=ERROR_MESSAGES.DEFAULT('Error updating function'),
329 )
330 else:
331 raise HTTPException(
332 status_code=status.HTTP_401_UNAUTHORIZED,
333 detail=ERROR_MESSAGES.NOT_FOUND,
334 )
337############################
338# ToggleGlobalById
339############################
342@router.post('/id/{id}/toggle/global', response_model=FunctionModel | None)
343async def toggle_global_by_id(
344 request: Request,
345 id: str,
346 user=Depends(get_admin_user),
347 db: AsyncSession = Depends(get_async_session),
348):
349 function = await Functions.get_function_by_id(id, db=db)
350 if function: 350 ↛ 351line 350 didn't jump to line 351 because the condition on line 350 was never true
351 function = await Functions.update_function_by_id(id, {'is_global': not function.is_global}, db=db)
353 if function:
354 await publish_event(
355 request,
356 EVENTS.FUNCTION_UPDATED,
357 actor=user,
358 subject_id=function.id,
359 data={'type': function.type, 'name': function.name, 'is_global': function.is_global},
360 )
361 return function
362 else:
363 raise HTTPException(
364 status_code=status.HTTP_400_BAD_REQUEST,
365 detail=ERROR_MESSAGES.DEFAULT('Error updating function'),
366 )
367 else:
368 raise HTTPException(
369 status_code=status.HTTP_401_UNAUTHORIZED,
370 detail=ERROR_MESSAGES.NOT_FOUND,
371 )
374############################
375# UpdateFunctionById
376############################
379@router.post('/id/{id}/update', response_model=FunctionModel | None)
380async def update_function_by_id(
381 request: Request,
382 id: str,
383 form_data: FunctionForm,
384 user=Depends(get_admin_user),
385 db: AsyncSession = Depends(get_async_session),
386):
387 try:
388 form_data.content = replace_imports(form_data.content)
389 function_module, function_type, frontmatter = await load_function_module_by_id(id, content=form_data.content)
390 form_data.meta.manifest = frontmatter
392 FUNCTIONS = get_functions_cache(request)
393 FUNCTIONS[id] = function_module
395 updated = {**form_data.model_dump(exclude={'id'}), 'type': function_type}
396 log.debug(updated)
398 function = await Functions.update_function_by_id(id, updated, db=db)
400 if function_type == 'filter' and getattr(function_module, 'toggle', None):
401 await Functions.update_function_metadata_by_id(id, {'toggle': True}, db=db)
403 if function:
404 await publish_event(
405 request,
406 EVENTS.FUNCTION_UPDATED,
407 actor=user,
408 subject_id=function.id,
409 data={'type': function.type, 'name': function.name},
410 )
411 return function
412 else:
413 raise HTTPException(
414 status_code=status.HTTP_400_BAD_REQUEST,
415 detail=ERROR_MESSAGES.DEFAULT('Error updating function'),
416 )
418 except HTTPException:
419 raise
420 except Exception as e:
421 raise HTTPException(
422 status_code=status.HTTP_400_BAD_REQUEST,
423 detail=ERROR_MESSAGES.DEFAULT(e, 'Error updating function'),
424 )
427############################
428# DeleteFunctionById
429############################
432@router.delete('/id/{id}/delete', response_model=bool)
433async def delete_function_by_id(
434 request: Request,
435 id: str,
436 user=Depends(get_admin_user),
437 db: AsyncSession = Depends(get_async_session),
438):
439 result = await Functions.delete_function_by_id(id, db=db)
441 if result: 441 ↛ 453line 441 didn't jump to line 453 because the condition on line 441 was always true
442 FUNCTIONS = get_functions_cache(request)
443 FUNCTIONS.pop(id, None)
444 FUNCTION_CONTENTS = get_function_contents_cache(request)
445 FUNCTION_CONTENTS.pop(id, None)
446 await publish_event(
447 request,
448 EVENTS.FUNCTION_DELETED,
449 actor=user,
450 subject_id=id,
451 )
453 return result
456############################
457# GetFunctionValves
458############################
461@router.get('/id/{id}/valves', response_model=dict | None)
462async def get_function_valves_by_id(
463 id: str, user=Depends(get_admin_user), db: AsyncSession = Depends(get_async_session)
464):
465 function = await Functions.get_function_by_id(id, db=db)
466 if function: 466 ↛ 467line 466 didn't jump to line 467 because the condition on line 466 was never true
467 try:
468 valves = await Functions.get_function_valves_by_id(id, db=db)
469 return valves
470 except Exception as e:
471 raise HTTPException(
472 status_code=status.HTTP_400_BAD_REQUEST,
473 detail=ERROR_MESSAGES.DEFAULT(e, 'Error getting function valves'),
474 )
475 else:
476 raise HTTPException(
477 status_code=status.HTTP_401_UNAUTHORIZED,
478 detail=ERROR_MESSAGES.NOT_FOUND,
479 )
482############################
483# GetFunctionValvesSpec
484############################
487@router.get('/id/{id}/valves/spec', response_model=dict | None)
488async def get_function_valves_spec_by_id(
489 request: Request,
490 id: str,
491 user=Depends(get_admin_user),
492 db: AsyncSession = Depends(get_async_session),
493):
494 function = await Functions.get_function_by_id(id, db=db)
495 if function: 495 ↛ 496line 495 didn't jump to line 496 because the condition on line 495 was never true
496 function_module, function_type, frontmatter = await get_function_module_from_cache(request, id)
498 if hasattr(function_module, 'Valves'):
499 Valves = function_module.Valves
500 schema = Valves.schema()
501 # Resolve dynamic options for select dropdowns
502 schema = resolve_valves_schema_options(Valves, schema, user)
503 return schema
504 return None
505 else:
506 raise HTTPException(
507 status_code=status.HTTP_401_UNAUTHORIZED,
508 detail=ERROR_MESSAGES.NOT_FOUND,
509 )
512############################
513# UpdateFunctionValves
514############################
517@router.post('/id/{id}/valves/update', response_model=dict | None)
518async def update_function_valves_by_id(
519 request: Request,
520 id: str,
521 form_data: dict,
522 user=Depends(get_admin_user),
523 db: AsyncSession = Depends(get_async_session),
524):
525 function = await Functions.get_function_by_id(id, db=db)
526 if function: 526 ↛ 527line 526 didn't jump to line 527 because the condition on line 526 was never true
527 function_module, function_type, frontmatter = await get_function_module_from_cache(request, id)
529 if hasattr(function_module, 'Valves'):
530 Valves = function_module.Valves
532 try:
533 form_data = {k: v for k, v in form_data.items() if v is not None}
534 valves = Valves(**form_data)
536 valves_dict = valves.model_dump(exclude_unset=True)
537 await Functions.update_function_valves_by_id(id, valves_dict, db=db)
538 await publish_event(
539 request,
540 EVENTS.FUNCTION_VALVES_UPDATED,
541 actor=user,
542 subject_id=id,
543 )
544 return valves_dict
545 except Exception as e:
546 log.exception(f'Error updating function values by id {id}: {e}')
547 raise HTTPException(
548 status_code=status.HTTP_400_BAD_REQUEST,
549 detail=ERROR_MESSAGES.DEFAULT(e, 'Error updating function valves'),
550 )
551 else:
552 raise HTTPException(
553 status_code=status.HTTP_401_UNAUTHORIZED,
554 detail=ERROR_MESSAGES.NOT_FOUND,
555 )
557 else:
558 raise HTTPException(
559 status_code=status.HTTP_401_UNAUTHORIZED,
560 detail=ERROR_MESSAGES.NOT_FOUND,
561 )
564############################
565# FunctionUserValves
566############################
569@router.get('/id/{id}/valves/user', response_model=dict | None)
570async def get_function_user_valves_by_id(
571 id: str, user=Depends(get_verified_user), db: AsyncSession = Depends(get_async_session)
572):
573 function = await Functions.get_function_by_id(id, db=db)
574 if function: 574 ↛ 575line 574 didn't jump to line 575 because the condition on line 574 was never true
575 try:
576 user_valves = await Functions.get_user_valves_by_id_and_user_id(id, user.id, db=db)
577 return user_valves
578 except Exception as e:
579 raise HTTPException(
580 status_code=status.HTTP_400_BAD_REQUEST,
581 detail=ERROR_MESSAGES.DEFAULT(e, 'Error getting function user valves'),
582 )
583 else:
584 raise HTTPException(
585 status_code=status.HTTP_401_UNAUTHORIZED,
586 detail=ERROR_MESSAGES.NOT_FOUND,
587 )
590@router.get('/id/{id}/valves/user/spec', response_model=dict | None)
591async def get_function_user_valves_spec_by_id(
592 request: Request,
593 id: str,
594 user=Depends(get_verified_user),
595 db: AsyncSession = Depends(get_async_session),
596):
597 function = await Functions.get_function_by_id(id, db=db)
598 if function: 598 ↛ 599line 598 didn't jump to line 599 because the condition on line 598 was never true
599 if not function.is_active:
600 return None
602 function_module, function_type, frontmatter = await get_function_module_from_cache(request, id)
604 if hasattr(function_module, 'UserValves'):
605 UserValves = function_module.UserValves
606 schema = UserValves.schema()
607 # Resolve dynamic options for select dropdowns
608 schema = resolve_valves_schema_options(UserValves, schema, user)
609 return schema
610 return None
611 else:
612 raise HTTPException(
613 status_code=status.HTTP_401_UNAUTHORIZED,
614 detail=ERROR_MESSAGES.NOT_FOUND,
615 )
618@router.post('/id/{id}/valves/user/update', response_model=dict | None)
619async def update_function_user_valves_by_id(
620 request: Request,
621 id: str,
622 form_data: dict,
623 user=Depends(get_verified_user),
624 db: AsyncSession = Depends(get_async_session),
625):
626 function = await Functions.get_function_by_id(id, db=db)
628 if function: 628 ↛ 629line 628 didn't jump to line 629 because the condition on line 628 was never true
629 if not function.is_active:
630 raise HTTPException(
631 status_code=status.HTTP_400_BAD_REQUEST,
632 detail='Function is not active',
633 )
635 function_module, function_type, frontmatter = await get_function_module_from_cache(request, id)
637 if hasattr(function_module, 'UserValves'):
638 UserValves = function_module.UserValves
640 try:
641 form_data = {k: v for k, v in form_data.items() if v is not None}
642 user_valves = UserValves(**form_data)
643 user_valves_dict = user_valves.model_dump(exclude_unset=True)
644 await Functions.update_user_valves_by_id_and_user_id(id, user.id, user_valves_dict, db=db)
645 await publish_event(
646 request,
647 EVENTS.FUNCTION_VALVES_UPDATED,
648 actor=user,
649 subject_id=id,
650 data={'scope': 'user'},
651 )
652 return user_valves_dict
653 except Exception as e:
654 log.exception(f'Error updating function user valves by id {id}: {e}')
655 raise HTTPException(
656 status_code=status.HTTP_400_BAD_REQUEST,
657 detail=ERROR_MESSAGES.DEFAULT(e, 'Error updating function user valves'),
658 )
659 else:
660 raise HTTPException(
661 status_code=status.HTTP_401_UNAUTHORIZED,
662 detail=ERROR_MESSAGES.NOT_FOUND,
663 )
664 else:
665 raise HTTPException(
666 status_code=status.HTTP_401_UNAUTHORIZED,
667 detail=ERROR_MESSAGES.NOT_FOUND,
668 )