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

1import asyncio 

2import logging 

3import os 

4from typing import Optional 

5 

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 

28 

29log = logging.getLogger(__name__) 

30 

31 

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################################## 

39 

40 

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 

55 

56 

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] 

61 

62 

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] 

68 

69 if 'pipeline' in model: 

70 sorted_filters.append(model) 

71 

72 if not sorted_filters: 

73 return payload 

74 

75 async with aiohttp.ClientSession(trust_env=True) as session: 

76 for filter in sorted_filters: 

77 urlIdx = filter.get('urlIdx') 

78 

79 try: 

80 urlIdx = int(urlIdx) 

81 except Exception: 

82 continue 

83 

84 url, key = await get_openai_connection(urlIdx) 

85 

86 if not key: 

87 continue 

88 

89 headers = {'Authorization': f'Bearer {key}'} 

90 request_data = { 

91 'user': user, 

92 'body': payload, 

93 } 

94 

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 

116 

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}') 

125 

126 return payload 

127 

128 

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] 

134 

135 if 'pipeline' in model: 

136 sorted_filters = [model] + sorted_filters 

137 

138 if not sorted_filters: 

139 return payload 

140 

141 async with aiohttp.ClientSession(trust_env=True) as session: 

142 for filter in sorted_filters: 

143 urlIdx = filter.get('urlIdx') 

144 

145 try: 

146 urlIdx = int(urlIdx) 

147 except Exception: 

148 continue 

149 

150 url, key = await get_openai_connection(urlIdx) 

151 

152 if not key: 

153 continue 

154 

155 headers = {'Authorization': f'Bearer {key}'} 

156 request_data = { 

157 'user': user, 

158 'body': payload, 

159 } 

160 

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 

182 

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}') 

191 

192 return payload 

193 

194 

195################################## 

196# 

197# Pipelines Endpoints 

198# 

199################################## 

200 

201router = APIRouter() 

202 

203 

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) 

208 

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', []) 

211 

212 return { 

213 'data': [ 

214 { 

215 'url': base_urls[urlIdx], 

216 'idx': urlIdx, 

217 } 

218 for urlIdx in urlIdxs 

219 ] 

220 } 

221 

222 

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) 

232 

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 ) 

239 

240 upload_folder = f'{CACHE_DIR}/pipelines' 

241 os.makedirs(upload_folder, exist_ok=True) 

242 file_path = os.path.join(upload_folder, filename) 

243 

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) 

249 

250 url, key = await get_openai_connection(urlIdx) 

251 

252 headers = {'Authorization': f'Bearer {key}'} 

253 

254 async with aiohttp.ClientSession(trust_env=True) as session: 

255 form_data = aiohttp.FormData() 

256 

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 

261 

262 form_data.add_field( 

263 'file', 

264 pipeline_chunks(), 

265 filename=filename, 

266 content_type='application/octet-stream', 

267 ) 

268 

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() 

277 

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}') 

289 

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 

300 

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) 

309 

310 

311class AddPipelineForm(BaseModel): 

312 url: str 

313 urlIdx: int 

314 

315 

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 

321 

322 url, key = await get_openai_connection(urlIdx) 

323 

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() 

333 

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}') 

345 

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 

354 

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 ) 

359 

360 

361class DeletePipelineForm(BaseModel): 

362 id: str 

363 urlIdx: int 

364 

365 

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 

371 

372 url, key = await get_openai_connection(urlIdx) 

373 

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() 

383 

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}') 

395 

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 

404 

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 ) 

409 

410 

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) 

416 

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() 

425 

426 return {**data} 

427 except Exception as e: 

428 # Handle connection error here 

429 log.exception(f'Connection error: {e}') 

430 

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 

439 

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 ) 

444 

445 

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) 

456 

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() 

465 

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}') 

477 

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 

486 

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 ) 

491 

492 

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) 

503 

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() 

512 

513 return {**data} 

514 except Exception as e: 

515 # Handle connection error here 

516 log.exception(f'Connection error: {e}') 

517 

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 

526 

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 ) 

531 

532 

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) 

544 

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() 

554 

555 return {**data} 

556 except Exception as e: 

557 # Handle connection error here 

558 log.exception(f'Connection error: {e}') 

559 

560 detail = None 

561 

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 

569 

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 )