Coverage for open_webui/routers/pipelines.py: 35%
293 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
1import asyncio
2import logging
3import os
4from typing import Optional
6import aiofiles
7import aiohttp
8from fastapi import (
9 APIRouter,
10 Depends,
11 FastAPI,
12 File,
13 Form,
14 HTTPException,
15 Request,
16 UploadFile,
17 status,
18)
19from open_webui.config import CACHE_DIR
20from open_webui.constants import ERROR_MESSAGES
21from open_webui.env import AIOHTTP_CLIENT_SESSION_SSL, AIOHTTP_FILE_STREAM_CHUNK_SIZE
22from open_webui.events import EVENTS, publish_event
23from open_webui.models.config import Config
24from open_webui.routers.openai import get_all_models_responses
25from open_webui.utils.auth import get_admin_user
26from pydantic import BaseModel
27from starlette.responses import FileResponse
29log = logging.getLogger(__name__)
32##################################
33#
34# Pipeline Middleware
35# Every hand this passes through can corrupt it or
36# improve it. Let each stage leave it better than it found.
37#
38##################################
41def get_sorted_filters(model_id, models):
42 filters = [
43 model
44 for model in models.values()
45 if 'pipeline' in model
46 and 'type' in model['pipeline']
47 and model['pipeline']['type'] == 'filter'
48 and (
49 model['pipeline']['pipelines'] == ['*']
50 or any(model_id == target_model_id for target_model_id in model['pipeline']['pipelines'])
51 )
52 ]
53 sorted_filters = sorted(filters, key=lambda x: x['pipeline']['priority'])
54 return sorted_filters
57async def get_openai_connection(url_idx: int) -> tuple[str, str]:
58 base_urls = await Config.get('openai.api_base_urls', [])
59 api_keys = await Config.get('openai.api_keys', [])
60 return base_urls[url_idx], api_keys[url_idx]
63async def process_pipeline_inlet_filter(request, payload, user, models):
64 user = {'id': user.id, 'email': user.email, 'name': user.name, 'role': user.role}
65 model_id = payload['model']
66 sorted_filters = get_sorted_filters(model_id, models)
67 model = models[model_id]
69 if 'pipeline' in model:
70 sorted_filters.append(model)
72 if not sorted_filters:
73 return payload
75 async with aiohttp.ClientSession(trust_env=True) as session:
76 for filter in sorted_filters:
77 urlIdx = filter.get('urlIdx')
79 try:
80 urlIdx = int(urlIdx)
81 except Exception:
82 continue
84 url, key = await get_openai_connection(urlIdx)
86 if not key:
87 continue
89 headers = {'Authorization': f'Bearer {key}'}
90 request_data = {
91 'user': user,
92 'body': payload,
93 }
95 try:
96 async with session.post(
97 f'{url}/{filter["id"]}/filter/inlet',
98 headers=headers,
99 json=request_data,
100 ssl=AIOHTTP_CLIENT_SESSION_SSL,
101 ) as response:
102 response.raise_for_status()
103 payload = await response.json()
104 except aiohttp.ClientResponseError as e:
105 try:
106 res = await response.json() if 'application/json' in response.content_type else {}
107 if 'detail' in res:
108 raise HTTPException(
109 status_code=response.status,
110 detail=res['detail'],
111 )
112 except HTTPException:
113 raise
114 except Exception:
115 pass
117 raise HTTPException(
118 status_code=response.status,
119 detail=e.message,
120 )
121 except HTTPException:
122 raise
123 except Exception as e:
124 log.exception(f'Connection error: {e}')
126 return payload
129async def process_pipeline_outlet_filter(request, payload, user, models):
130 user = {'id': user.id, 'email': user.email, 'name': user.name, 'role': user.role}
131 model_id = payload['model']
132 sorted_filters = get_sorted_filters(model_id, models)
133 model = models[model_id]
135 if 'pipeline' in model:
136 sorted_filters = [model] + sorted_filters
138 if not sorted_filters:
139 return payload
141 async with aiohttp.ClientSession(trust_env=True) as session:
142 for filter in sorted_filters:
143 urlIdx = filter.get('urlIdx')
145 try:
146 urlIdx = int(urlIdx)
147 except Exception:
148 continue
150 url, key = await get_openai_connection(urlIdx)
152 if not key:
153 continue
155 headers = {'Authorization': f'Bearer {key}'}
156 request_data = {
157 'user': user,
158 'body': payload,
159 }
161 try:
162 async with session.post(
163 f'{url}/{filter["id"]}/filter/outlet',
164 headers=headers,
165 json=request_data,
166 ssl=AIOHTTP_CLIENT_SESSION_SSL,
167 ) as response:
168 response.raise_for_status()
169 payload = await response.json()
170 except aiohttp.ClientResponseError as e:
171 try:
172 res = await response.json() if 'application/json' in response.content_type else {}
173 if 'detail' in res:
174 raise HTTPException(
175 status_code=response.status,
176 detail=res['detail'],
177 )
178 except HTTPException:
179 raise
180 except Exception:
181 pass
183 raise HTTPException(
184 status_code=response.status,
185 detail=e.message,
186 )
187 except HTTPException:
188 raise
189 except Exception as e:
190 log.exception(f'Connection error: {e}')
192 return payload
195##################################
196#
197# Pipelines Endpoints
198#
199##################################
201router = APIRouter()
204@router.get('/list')
205async def get_pipelines_list(request: Request, user=Depends(get_admin_user)):
206 responses = await get_all_models_responses(request, user)
207 log.debug('get_pipelines_list: get_openai_models_responses returned %s', responses)
209 urlIdxs = [idx for idx, response in enumerate(responses) if response is not None and 'pipelines' in response]
210 base_urls = await Config.get('openai.api_base_urls', [])
212 return {
213 'data': [
214 {
215 'url': base_urls[urlIdx],
216 'idx': urlIdx,
217 }
218 for urlIdx in urlIdxs
219 ]
220 }
223@router.post('/upload')
224async def upload_pipeline(
225 request: Request,
226 urlIdx: int = Form(...),
227 file: UploadFile = File(...),
228 user=Depends(get_admin_user),
229):
230 log.info('upload_pipeline: urlIdx=%s, filename=%s', urlIdx, file.filename)
231 filename = os.path.basename(file.filename)
233 # Check if the uploaded file is a python file
234 if not (filename and filename.endswith('.py')): 234 ↛ 240line 234 didn't jump to line 240 because the condition on line 234 was always true
235 raise HTTPException(
236 status_code=status.HTTP_400_BAD_REQUEST,
237 detail='Only Python (.py) files are allowed.',
238 )
240 upload_folder = f'{CACHE_DIR}/pipelines'
241 os.makedirs(upload_folder, exist_ok=True)
242 file_path = os.path.join(upload_folder, filename)
244 response = None
245 try:
246 async with aiofiles.open(file_path, 'wb') as buffer:
247 while chunk := await file.read(AIOHTTP_FILE_STREAM_CHUNK_SIZE):
248 await buffer.write(chunk)
250 url, key = await get_openai_connection(urlIdx)
252 headers = {'Authorization': f'Bearer {key}'}
254 async with aiohttp.ClientSession(trust_env=True) as session:
255 form_data = aiohttp.FormData()
257 async def pipeline_chunks():
258 async with aiofiles.open(file_path, 'rb') as pipeline_file:
259 while chunk := await pipeline_file.read(AIOHTTP_FILE_STREAM_CHUNK_SIZE):
260 yield chunk
262 form_data.add_field(
263 'file',
264 pipeline_chunks(),
265 filename=filename,
266 content_type='application/octet-stream',
267 )
269 async with session.post(
270 f'{url}/pipelines/upload',
271 headers=headers,
272 data=form_data,
273 ssl=AIOHTTP_CLIENT_SESSION_SSL,
274 ) as response:
275 response.raise_for_status()
276 data = await response.json()
278 await publish_event(
279 request,
280 EVENTS.PIPELINE_UPLOADED,
281 actor=user,
282 subject_id=data.get('id') or filename,
283 data={'url_idx': urlIdx, 'filename': filename},
284 )
285 return {**data}
286 except Exception as e:
287 # Handle connection error here
288 log.exception(f'Connection error: {e}')
290 detail = None
291 status_code = status.HTTP_404_NOT_FOUND
292 if response is not None:
293 status_code = response.status
294 try:
295 res = await response.json()
296 if 'detail' in res:
297 detail = res['detail']
298 except Exception:
299 pass
301 raise HTTPException(
302 status_code=status_code,
303 detail=detail if detail else 'Pipeline not found',
304 )
305 finally:
306 # Ensure the file is deleted after the upload is completed or on failure
307 if os.path.exists(file_path):
308 await asyncio.to_thread(os.remove, file_path)
311class AddPipelineForm(BaseModel):
312 url: str
313 urlIdx: int
316@router.post('/add')
317async def add_pipeline(request: Request, form_data: AddPipelineForm, user=Depends(get_admin_user)):
318 response = None
319 try:
320 urlIdx = form_data.urlIdx
322 url, key = await get_openai_connection(urlIdx)
324 async with aiohttp.ClientSession(trust_env=True) as session:
325 async with session.post(
326 f'{url}/pipelines/add',
327 headers={'Authorization': f'Bearer {key}'},
328 json={'url': form_data.url},
329 ssl=AIOHTTP_CLIENT_SESSION_SSL,
330 ) as response:
331 response.raise_for_status()
332 data = await response.json()
334 await publish_event(
335 request,
336 EVENTS.PIPELINE_ADDED,
337 actor=user,
338 subject_id=data.get('id') or form_data.url,
339 data={'url_idx': urlIdx, 'url': form_data.url},
340 )
341 return {**data}
342 except Exception as e:
343 # Handle connection error here
344 log.exception(f'Connection error: {e}')
346 detail = None
347 if response is not None:
348 try:
349 res = await response.json()
350 if 'detail' in res:
351 detail = res['detail']
352 except Exception:
353 pass
355 raise HTTPException(
356 status_code=(response.status if response is not None else status.HTTP_404_NOT_FOUND),
357 detail=detail if detail else 'Pipeline not found',
358 )
361class DeletePipelineForm(BaseModel):
362 id: str
363 urlIdx: int
366@router.delete('/delete')
367async def delete_pipeline(request: Request, form_data: DeletePipelineForm, user=Depends(get_admin_user)):
368 response = None
369 try:
370 urlIdx = form_data.urlIdx
372 url, key = await get_openai_connection(urlIdx)
374 async with aiohttp.ClientSession(trust_env=True) as session:
375 async with session.delete(
376 f'{url}/pipelines/delete',
377 headers={'Authorization': f'Bearer {key}'},
378 json={'id': form_data.id},
379 ssl=AIOHTTP_CLIENT_SESSION_SSL,
380 ) as response:
381 response.raise_for_status()
382 data = await response.json()
384 await publish_event(
385 request,
386 EVENTS.PIPELINE_DELETED,
387 actor=user,
388 subject_id=form_data.id,
389 data={'url_idx': urlIdx},
390 )
391 return {**data}
392 except Exception as e:
393 # Handle connection error here
394 log.exception(f'Connection error: {e}')
396 detail = None
397 if response is not None: 397 ↛ 398line 397 didn't jump to line 398 because the condition on line 397 was never true
398 try:
399 res = await response.json()
400 if 'detail' in res:
401 detail = res['detail']
402 except Exception:
403 pass
405 raise HTTPException(
406 status_code=(response.status if response is not None else status.HTTP_404_NOT_FOUND),
407 detail=detail if detail else 'Pipeline not found',
408 )
411@router.get('/')
412async def get_pipelines(request: Request, urlIdx: Optional[int] = None, user=Depends(get_admin_user)):
413 response = None
414 try:
415 url, key = await get_openai_connection(urlIdx)
417 async with aiohttp.ClientSession(trust_env=True) as session:
418 async with session.get(
419 f'{url}/pipelines',
420 headers={'Authorization': f'Bearer {key}'},
421 ssl=AIOHTTP_CLIENT_SESSION_SSL,
422 ) as response:
423 response.raise_for_status()
424 data = await response.json()
426 return {**data}
427 except Exception as e:
428 # Handle connection error here
429 log.exception(f'Connection error: {e}')
431 detail = None
432 if response is not None: 432 ↛ 433line 432 didn't jump to line 433 because the condition on line 432 was never true
433 try:
434 res = await response.json()
435 if 'detail' in res:
436 detail = res['detail']
437 except Exception:
438 pass
440 raise HTTPException(
441 status_code=(response.status if response is not None else status.HTTP_404_NOT_FOUND),
442 detail=detail if detail else 'Pipeline not found',
443 )
446@router.get('/{pipeline_id}/valves')
447async def get_pipeline_valves(
448 request: Request,
449 urlIdx: Optional[int],
450 pipeline_id: str,
451 user=Depends(get_admin_user),
452):
453 response = None
454 try:
455 url, key = await get_openai_connection(urlIdx)
457 async with aiohttp.ClientSession(trust_env=True) as session:
458 async with session.get(
459 f'{url}/{pipeline_id}/valves',
460 headers={'Authorization': f'Bearer {key}'},
461 ssl=AIOHTTP_CLIENT_SESSION_SSL,
462 ) as response:
463 response.raise_for_status()
464 data = await response.json()
466 await publish_event(
467 request,
468 EVENTS.PIPELINE_VALVES_UPDATED,
469 actor=user,
470 subject_id=pipeline_id,
471 data={'url_idx': urlIdx},
472 )
473 return {**data}
474 except Exception as e:
475 # Handle connection error here
476 log.exception(f'Connection error: {e}')
478 detail = None
479 if response is not None: 479 ↛ 480line 479 didn't jump to line 480 because the condition on line 479 was never true
480 try:
481 res = await response.json()
482 if 'detail' in res:
483 detail = res['detail']
484 except Exception:
485 pass
487 raise HTTPException(
488 status_code=(response.status if response is not None else status.HTTP_404_NOT_FOUND),
489 detail=detail if detail else 'Pipeline not found',
490 )
493@router.get('/{pipeline_id}/valves/spec')
494async def get_pipeline_valves_spec(
495 request: Request,
496 urlIdx: Optional[int],
497 pipeline_id: str,
498 user=Depends(get_admin_user),
499):
500 response = None
501 try:
502 url, key = await get_openai_connection(urlIdx)
504 async with aiohttp.ClientSession(trust_env=True) as session:
505 async with session.get(
506 f'{url}/{pipeline_id}/valves/spec',
507 headers={'Authorization': f'Bearer {key}'},
508 ssl=AIOHTTP_CLIENT_SESSION_SSL,
509 ) as response:
510 response.raise_for_status()
511 data = await response.json()
513 return {**data}
514 except Exception as e:
515 # Handle connection error here
516 log.exception(f'Connection error: {e}')
518 detail = None
519 if response is not None: 519 ↛ 520line 519 didn't jump to line 520 because the condition on line 519 was never true
520 try:
521 res = await response.json()
522 if 'detail' in res:
523 detail = res['detail']
524 except Exception:
525 pass
527 raise HTTPException(
528 status_code=(response.status if response is not None else status.HTTP_404_NOT_FOUND),
529 detail=detail if detail else 'Pipeline not found',
530 )
533@router.post('/{pipeline_id}/valves/update')
534async def update_pipeline_valves(
535 request: Request,
536 urlIdx: Optional[int],
537 pipeline_id: str,
538 form_data: dict,
539 user=Depends(get_admin_user),
540):
541 response = None
542 try:
543 url, key = await get_openai_connection(urlIdx)
545 async with aiohttp.ClientSession(trust_env=True) as session:
546 async with session.post(
547 f'{url}/{pipeline_id}/valves/update',
548 headers={'Authorization': f'Bearer {key}'},
549 json={**form_data},
550 ssl=AIOHTTP_CLIENT_SESSION_SSL,
551 ) as response:
552 response.raise_for_status()
553 data = await response.json()
555 return {**data}
556 except Exception as e:
557 # Handle connection error here
558 log.exception(f'Connection error: {e}')
560 detail = None
562 if response is not None: 562 ↛ 563line 562 didn't jump to line 563 because the condition on line 562 was never true
563 try:
564 res = await response.json()
565 if 'detail' in res:
566 detail = res['detail']
567 except Exception:
568 pass
570 raise HTTPException(
571 status_code=(response.status if response is not None else status.HTTP_404_NOT_FOUND),
572 detail=detail if detail else 'Pipeline not found',
573 )