Coverage for open_webui/utils/automations.py: 36%

216 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 05:07 +0000

1""" 

2Automation utilities and unified scheduler. 

3 

4RRULE helpers, scheduler worker loop, and execution logic. 

5Follows the utils/<feature>.py pattern (cf. utils/channels.py, utils/task.py). 

6 

7The scheduler_worker_loop handles all time-based background work: 

8 - Automation execution (claim_due → execute) 

9 - Calendar event alerts (upcoming events → socket + webhook notifications) 

10 - One-shot chat timers 

11 

12Environment: 

13 SCHEDULER_POLL_INTERVAL – seconds between polls (default: 10) 

14 TIMER_POLL_INTERVAL – seconds between timer polls (default: 1) 

15 CALENDAR_ALERT_LOOKAHEAD_MINUTES – default alert window (default: 5) 

16""" 

17 

18import asyncio 

19import logging 

20import os 

21import random 

22import time 

23from datetime import timedelta 

24from typing import Optional 

25from uuid import uuid4 

26 

27from fastapi import Request 

28from fastapi.security import HTTPAuthorizationCredentials 

29from open_webui.constants import ERROR_MESSAGES 

30from open_webui.events import EVENTS, publish_event 

31from open_webui.internal.db import get_async_db 

32from open_webui.models.automations import AutomationModel, AutomationRuns, Automations 

33from open_webui.models.chats import ChatForm, Chats 

34from open_webui.models.config import Config 

35from open_webui.models.folders import Folders 

36from open_webui.models.messages import MessageForm 

37from open_webui.models.users import Users 

38from open_webui.utils.auth import create_token 

39from open_webui.utils.misc import parse_duration 

40from open_webui.utils.recurrence import ( 

41 _resolve_tz, 

42 next_n_runs_ns, 

43 next_run_ns, 

44 rrule_interval_seconds, 

45 validate_rrule, 

46) 

47from open_webui.utils.task import prompt_template 

48from open_webui.utils.terminals import get_terminal_server_url 

49from starlette.datastructures import Headers 

50 

51log = logging.getLogger(__name__) 

52 

53SCHEDULER_POLL_INTERVAL = int(os.getenv('SCHEDULER_POLL_INTERVAL', os.getenv('AUTOMATION_POLL_INTERVAL', '10'))) 

54TIMER_POLL_INTERVAL = int(os.getenv('TIMER_POLL_INTERVAL', '1')) 

55CALENDAR_ALERT_LOOKAHEAD_MINUTES = int(os.getenv('CALENDAR_ALERT_LOOKAHEAD_MINUTES', '10')) 

56 

57 

58############################ 

59# Worker Loop 

60############################ 

61 

62 

63# Keep the old name as an alias so any stale imports still work. 

64async def automation_worker_loop(app) -> None: 

65 """Deprecated alias — use scheduler_worker_loop.""" 

66 await scheduler_worker_loop(app) 

67 

68 

69async def scheduler_worker_loop(app) -> None: 

70 """Unified background scheduler for all time-based work. 

71 

72 Handles: 

73 1. Automation execution (ENABLE_AUTOMATIONS) 

74 2. Calendar event alerts (ENABLE_CALENDAR) 

75 

76 Runs on every instance. Poll interval is configurable via 

77 SCHEDULER_POLL_INTERVAL env var (default: 10 seconds). 

78 """ 

79 log.info( 

80 'Scheduler worker started (timer poll interval: %ss, scheduler poll interval: %ss)', 

81 TIMER_POLL_INTERVAL, 

82 SCHEDULER_POLL_INTERVAL, 

83 ) 

84 next_scheduler_poll = 0.0 

85 

86 while True: 

87 try: 

88 now = time.monotonic() 

89 # ── Timers ── 

90 try: 

91 from open_webui.utils.timers import claim_due_timers, execute_due_timer 

92 

93 for timer_id, claim_id in await claim_due_timers(int(time.time_ns()), limit=10): 93 ↛ 94line 93 didn't jump to line 94 because the loop on line 93 never started

94 asyncio.create_task(execute_due_timer(app, timer_id, claim_id)) 

95 except Exception: 

96 log.exception('Scheduler: timer error') 

97 

98 if now < next_scheduler_poll: 

99 await asyncio.sleep(max(1, TIMER_POLL_INTERVAL)) 

100 continue 

101 # Jitter to spread automation/calendar load across instances; timers keep a tight poll. 

102 next_scheduler_poll = now + SCHEDULER_POLL_INTERVAL + random.uniform(0, 2) 

103 

104 # ── Automations ── 

105 if await Config.get('automations.enable'): 105 ↛ 117line 105 didn't jump to line 117 because the condition on line 105 was always true

106 try: 

107 async with get_async_db() as db: 

108 batch = await Automations.claim_due(int(time.time_ns()), limit=10, db=db) 

109 if batch: 109 ↛ 110line 109 didn't jump to line 110 because the condition on line 109 was never true

110 log.info('Claimed %s due automation(s)', len(batch)) 

111 for automation in batch: 111 ↛ 112line 111 didn't jump to line 112 because the loop on line 111 never started

112 asyncio.create_task(execute_automation(app, automation)) 

113 except Exception: 

114 log.exception('Scheduler: automation error') 

115 

116 # ── Calendar Alerts ── 

117 if await Config.get('calendar.enable'): 117 ↛ 126line 117 didn't jump to line 126 because the condition on line 117 was always true

118 try: 

119 await _check_calendar_alerts(app) 

120 except Exception: 

121 log.exception('Scheduler: calendar alert error') 

122 

123 except Exception: 

124 log.exception('Scheduler worker error') 

125 

126 await asyncio.sleep(max(1, TIMER_POLL_INTERVAL)) 

127 

128 

129########################## 

130# Execute 

131#################### 

132 

133 

134def _build_request( 

135 app, 

136 token: Optional[str] = None, 

137) -> Request: 

138 """Build a minimal ASGI Request for chat_completion. 

139 

140 Mirrors the mock-request pattern used in main.py lifespan 

141 (model pre-fetch, tool server init) for consistency. 

142 

143 When token is provided, attach it as 

144 request.state.token so session-auth tool servers and terminals can 

145 authenticate headless scheduled runs as the automation owner. 

146 """ 

147 scope = { 

148 'type': 'http', 

149 'asgi': {'version': '3.0', 'spec_version': '2.0'}, 

150 'method': 'POST', 

151 'path': '/api/v1/automations/internal', 

152 'query_string': b'', 

153 'headers': Headers({}).raw, 

154 'client': ('127.0.0.1', 0), 

155 'server': ('127.0.0.1', 80), 

156 'scheme': 'http', 

157 'app': app, 

158 } 

159 request = Request(scope) 

160 # Ensure request.state is initialized with required attributes 

161 request.state.token = HTTPAuthorizationCredentials(scheme='Bearer', credentials=token) if token else None 

162 request.state.enable_api_keys = False 

163 return request 

164 

165 

166async def _resolve_model_defaults(app, model_id: str) -> dict: 

167 models = getattr(app.state, 'MODELS', {}) 

168 model = models.get(model_id, {}) 

169 meta = model.get('info', {}).get('meta', {}) 

170 

171 defaults = { 

172 'tool_ids': list(meta.get('toolIds') or []), 

173 'filter_ids': list(meta.get('defaultFilterIds') or []), 

174 'terminal_id': meta.get('terminalId'), 

175 } 

176 defaults = {key: value for key, value in defaults.items() if value} 

177 default_feature_ids = meta.get('defaultFeatureIds', []) 

178 if not default_feature_ids: 

179 return defaults 

180 

181 capabilities = meta.get('capabilities') or {} 

182 features = {} 

183 

184 # code_interpreter is excluded: it requires the frontend event emitter 

185 # and does not work in headless backend execution. 

186 feature_checks = { 

187 'web_search': await Config.get('web.search.enable'), 

188 'image_generation': await Config.get('image_generation.enable'), 

189 } 

190 

191 for feature_id in default_feature_ids: 191 ↛ anywhereline 191 didn't jump anywhere: it always raised an exception.

192 if feature_id in feature_checks: 

193 # Feature must be: in defaultFeatureIds + capability enabled + admin enabled 

194 if capabilities.get(feature_id) and feature_checks[feature_id]: 

195 features[feature_id] = True 

196 

197 if features: 

198 defaults['features'] = features 

199 return defaults 

200 

201 

202async def _set_terminal_cwd(app, server_id: str, user, cwd: str, chat_id: str) -> None: 

203 """Set the working directory on a terminal server via the proxy. 

204 

205 Routes through the open-webui terminal proxy endpoint so that 

206 auth headers, orchestrator policy routing, and X-User-Id are 

207 handled correctly — same path the frontend uses. 

208 """ 

209 import aiohttp 

210 from open_webui.env import AIOHTTP_CLIENT_SESSION_SSL 

211 

212 connections = getattr(getattr(app, 'state', None), 'config', None) 

213 if connections is None: 

214 return 

215 connections = getattr(connections, 'TERMINAL_SERVER_CONNECTIONS', None) or [] 

216 connection = next((c for c in connections if c.get('id') == server_id), None) 

217 if connection is None: 

218 log.warning(f'Terminal server {server_id} not found for CWD set') 

219 return 

220 

221 base_url = get_terminal_server_url(connection) 

222 if not base_url: 

223 return 

224 

225 target_url = f'{base_url}/files/cwd' 

226 

227 headers = {'Content-Type': 'application/json', 'X-User-Id': user.id} 

228 if chat_id: 

229 headers['X-Session-Id'] = chat_id 

230 

231 auth_type = connection.get('auth_type', 'bearer') 

232 if auth_type == 'bearer': 

233 headers['Authorization'] = f'Bearer {connection.get("key", "")}' 

234 

235 try: 

236 async with aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=10)) as session: 

237 async with session.post( 

238 target_url, 

239 json={'path': cwd}, 

240 headers=headers, 

241 ssl=AIOHTTP_CLIENT_SESSION_SSL, 

242 ) as resp: 

243 if resp.status != 200: 

244 body = await resp.text() 

245 log.warning(f'Failed to set terminal CWD to {cwd}: HTTP {resp.status} — {body[:200]}') 

246 except Exception as e: 

247 log.warning(f'Failed to set terminal CWD: {e}') 

248 

249 

250async def _execute_channel_automation( 

251 app, 

252 automation: AutomationModel, 

253 user, 

254 prompt: str, 

255 model_id: str, 

256 token: str, 

257) -> None: 

258 target = automation.data.get('target') or {} 

259 channel_id = target.get('channel_id') 

260 if not channel_id or not await Config.get('channels.enable'): 

261 raise ValueError('Channel not found') 

262 

263 model = getattr(app.state, 'MODELS', {}).get(model_id, {}) 

264 request = _build_request(app, token=token) 

265 

266 from open_webui.routers.channels import new_message_handler 

267 

268 async with get_async_db() as db: 

269 user_message, channel = await new_message_handler( 

270 request, 

271 channel_id, 

272 MessageForm( 

273 content=prompt, 

274 data={}, 

275 meta={'automation_id': automation.id}, 

276 ), 

277 user, 

278 db, 

279 ) 

280 response_parent_id = ( 

281 user_message.parent_id 

282 if user_message.parent_id 

283 else (user_message.id if await Config.get('channels.model_response_mode', 'thread') == 'thread' else None) 

284 ) 

285 assistant_message, channel = await new_message_handler( 

286 request, 

287 channel.id, 

288 MessageForm( 

289 parent_id=response_parent_id, 

290 content='', 

291 data={}, 

292 meta={ 

293 'automation_id': automation.id, 

294 'model_id': model_id, 

295 'model_name': model.get('name', model_id), 

296 }, 

297 ), 

298 user, 

299 db, 

300 ) 

301 

302 form_data = { 

303 **await _resolve_model_defaults(app, model_id), 

304 'model': model_id, 

305 'messages': [ 

306 { 

307 'role': 'system', 

308 'content': f'You are {model.get("name", model_id)}, participating in a channel conversation. Be concise and conversational.', 

309 }, 

310 {'role': 'user', 'content': f'{user.name if user else "User"}: {prompt}'}, 

311 ], 

312 'stream': True, 

313 'chat_id': f'channel:{channel.id}', 

314 'id': assistant_message.id, 

315 'session_id': f'channel:{channel.id}', 

316 'automation_id': automation.id, 

317 'background_tasks': {}, 

318 } 

319 await app.state.CHAT_COMPLETION_HANDLER(request, form_data, user=user) 

320 

321 from open_webui.socket.main import sio 

322 

323 await sio.emit( 

324 'automation:result', 

325 { 

326 'automation_id': automation.id, 

327 'name': automation.name, 

328 'chat_id': f'channel:{channel.id}', 

329 'message_id': assistant_message.id, 

330 'status': 'success', 

331 }, 

332 room=f'user:{automation.user_id}', 

333 ) 

334 

335 await _record_run(automation.id, 'success', chat_id=f'channel:{channel.id}') 

336 await publish_event( 

337 app, 

338 EVENTS.AUTOMATION_RUN_COMPLETED, 

339 actor=user, 

340 subject_id=automation.id, 

341 data={'name': automation.name, 'channel_id': channel.id, 'message_id': assistant_message.id}, 

342 ) 

343 

344 

345async def execute_automation(app, automation: AutomationModel) -> None: 

346 """Execute an automation through the full chat completion pipeline. 

347 

348 Creates a real chat or channel message, then calls chat_completion exactly like the frontend: 

349 session_id + chat_id + message_id → async task → pipeline handles everything 

350 (filters, model params, knowledge/RAG, tools, DB saves, webhooks). 

351 """ 

352 try: 

353 user = await Users.get_user_by_id(automation.user_id) 

354 if not user: 

355 await _record_run(automation.id, 'error', error='User not found') 

356 await publish_event( 

357 app, 

358 EVENTS.AUTOMATION_RUN_FAILED, 

359 subject_id=automation.id, 

360 data={'name': automation.name, 'error': 'User not found'}, 

361 ) 

362 return 

363 

364 # Re-gate the rehydrated owner: a demoted/deactivated or de-permissioned owner must not run. 

365 from open_webui.utils.access_control import has_permission 

366 

367 if user.role not in ('user', 'admin') or ( 

368 user.role != 'admin' 

369 and not await has_permission(user.id, 'features.automations', await Config.get('user.permissions')) 

370 ): 

371 error = 'Owner no longer permitted to run automations' 

372 await _record_run(automation.id, 'error', error=error) 

373 await publish_event( 

374 app, 

375 EVENTS.AUTOMATION_RUN_FAILED, 

376 actor=user, 

377 subject_id=automation.id, 

378 data={'name': automation.name, 'error': error}, 

379 ) 

380 return 

381 

382 prompt = await prompt_template(automation.data['prompt'], user) 

383 model_id = automation.data['model_id'] 

384 try: 

385 expires_delta = parse_duration(str(await Config.get('automations.auth_token_expires_in', '1h'))) 

386 except ValueError: 

387 expires_delta = None 

388 token = create_token( 

389 data={'id': user.id, 'typ': 'automation'}, 

390 expires_delta=expires_delta or timedelta(hours=1), 

391 ) 

392 

393 target = automation.data.get('target') or {} 

394 if target.get('type') == 'channel': 

395 await _execute_channel_automation(app, automation, user, prompt, model_id, token) 

396 return 

397 

398 folder_id = automation.folder_id 

399 if folder_id and not await Folders.get_folder_by_id_and_user_id(folder_id, automation.user_id): 

400 await Automations.clear_folder_ids(automation.user_id, [folder_id]) 

401 folder_id = None 

402 

403 # Generate proper UUIDs for messages (same as frontend) 

404 user_msg_id = str(uuid4()) 

405 assistant_msg_id = str(uuid4()) 

406 

407 chat_id = str(uuid4()) 

408 chat = await Chats.insert_new_chat( 

409 chat_id, 

410 automation.user_id, 

411 ChatForm( 

412 folder_id=folder_id, 

413 chat={ 

414 'title': automation.name, 

415 'models': [model_id], 

416 'history': { 

417 'currentId': assistant_msg_id, 

418 'messages': { 

419 user_msg_id: { 

420 'id': user_msg_id, 

421 'parentId': None, 

422 'role': 'user', 

423 'content': prompt, 

424 'childrenIds': [assistant_msg_id], 

425 'timestamp': int(time.time()), 

426 'models': [model_id], 

427 }, 

428 assistant_msg_id: { 

429 'id': assistant_msg_id, 

430 'parentId': user_msg_id, 

431 'role': 'assistant', 

432 'content': '', 

433 'done': False, 

434 'model': model_id, 

435 'childrenIds': [], 

436 'timestamp': int(time.time()), 

437 }, 

438 }, 

439 }, 

440 'messages': [ 

441 {'role': 'user', 'content': prompt}, 

442 ], 

443 'meta': {'automation_id': automation.id}, 

444 }, 

445 ), 

446 ) 

447 

448 if not chat: 

449 error = 'Failed to create chat' 

450 await _record_run(automation.id, 'error', error=error) 

451 await publish_event( 

452 app, 

453 EVENTS.AUTOMATION_RUN_FAILED, 

454 actor=user, 

455 subject_id=automation.id, 

456 data={'name': automation.name, 'error': error}, 

457 ) 

458 return 

459 

460 # Notify frontend to refresh chat list 

461 from open_webui.socket.main import sio 

462 

463 await sio.emit( 

464 'events', 

465 { 

466 'chat_id': chat.id, 

467 'message_id': user_msg_id, 

468 'data': {'type': 'chat:list'}, 

469 }, 

470 room=f'user:{automation.user_id}', 

471 ) 

472 

473 # Build the same payload the frontend sends to /api/chat/completions 

474 form_data = { 

475 **await _resolve_model_defaults(app, model_id), 

476 'model': model_id, 

477 'messages': [{'role': 'user', 'content': prompt}], 

478 'stream': True, 

479 'chat_id': chat.id, 

480 'id': assistant_msg_id, 

481 'parent_id': None, # Root message (chat already created above) 

482 'user_message': { 

483 'id': user_msg_id, 

484 'parentId': None, 

485 'role': 'user', 

486 'content': prompt, 

487 }, 

488 'session_id': f'automation:{automation.id}', 

489 'automation_id': automation.id, 

490 'background_tasks': {}, 

491 } 

492 # Call the full chat completion pipeline (same as POST /api/chat/completions). 

493 # The handler reference is stored on app.state to avoid circular imports. 

494 request = _build_request(app, token=token) 

495 await app.state.CHAT_COMPLETION_HANDLER(request, form_data, user=user) 

496 

497 # Notify user 

498 from open_webui.socket.main import sio 

499 

500 await sio.emit( 

501 'automation:result', 

502 { 

503 'automation_id': automation.id, 

504 'name': automation.name, 

505 'chat_id': chat.id, 

506 'status': 'success', 

507 }, 

508 room=f'user:{automation.user_id}', 

509 ) 

510 

511 await _record_run(automation.id, 'success', chat_id=chat.id) 

512 await publish_event( 

513 app, 

514 EVENTS.AUTOMATION_RUN_COMPLETED, 

515 actor=user, 

516 subject_id=automation.id, 

517 data={'name': automation.name, 'chat_id': chat.id}, 

518 ) 

519 

520 except Exception as e: 

521 log.exception(f'Automation {automation.id} failed') 

522 error = str(e)[:4000] 

523 await _record_run(automation.id, 'error', error=error) 

524 await publish_event( 

525 app, 

526 EVENTS.AUTOMATION_RUN_FAILED, 

527 subject_id=automation.id, 

528 data={'name': automation.name, 'error': error}, 

529 ) 

530 

531 

532#################### 

533# Internals 

534#################### 

535 

536 

537async def _check_calendar_alerts(app) -> None: 

538 """Check for upcoming calendar events and send alert notifications. 

539 

540 De-duplication is DB-backed via meta.alerted_at — survives restarts 

541 and works across multiple instances. 

542 """ 

543 from open_webui.models.calendar import CalendarEvents, CalendarEventUpdateForm 

544 from open_webui.socket.main import sio 

545 

546 now_ns = int(time.time_ns()) 

547 default_lookahead_ns = CALENDAR_ALERT_LOOKAHEAD_MINUTES * 60 * 1_000_000_000 

548 # Grace window covers one poll cycle + jitter so "At time of event" 

549 # alerts (alert_minutes=0) are not missed. 

550 grace_ns = (SCHEDULER_POLL_INTERVAL + 5) * 1_000_000_000 

551 

552 async with get_async_db() as db: 

553 upcoming = await CalendarEvents.get_upcoming_events(now_ns, default_lookahead_ns, grace_ns=grace_ns, db=db) 

554 

555 if not upcoming: 555 ↛ 558line 555 didn't jump to line 558 because the condition on line 555 was always true

556 return 

557 

558 for event, user_tz in upcoming: 

559 # Skip if already alerted for this start time 

560 if event.meta and event.meta.get('alerted_at'): 560 ↛ 564line 560 didn't jump to line 564 because the condition on line 560 was always true

561 continue 

562 

563 # Compute minutes until event starts 

564 minutes_until = max(0, int((event.start_at - now_ns) / (60 * 1_000_000_000))) 

565 

566 alert_data = { 

567 'event_id': event.id, 

568 'title': event.title, 

569 'description': event.description or '', 

570 'start_at': event.start_at, 

571 'minutes_until': minutes_until, 

572 'calendar_id': event.calendar_id, 

573 'location': event.location or '', 

574 } 

575 

576 await sio.emit( 

577 'events', 

578 { 

579 'data': { 

580 'type': 'calendar:alert', 

581 'data': alert_data, 

582 }, 

583 }, 

584 room=f'user:{event.user_id}', 

585 ) 

586 

587 # Mark as alerted in DB so it survives restarts / multi-instance 

588 try: 

589 await CalendarEvents.update_event_by_id( 

590 event.id, 

591 CalendarEventUpdateForm(meta={'alerted_at': now_ns}), 

592 ) 

593 except Exception: 

594 log.debug('Failed to mark event %s as alerted', event.id, exc_info=True) 

595 

596 # Send target notification if user has one configured 

597 try: 

598 time_str = f'in {minutes_until} min' if minutes_until > 0 else 'now' 

599 await publish_event( 

600 app, 

601 EVENTS.CALENDAR_ALERT, 

602 subject_id=event.id, 

603 subject_type='calendar.event', 

604 source='scheduler', 

605 data={ 

606 **alert_data, 

607 'user_id': event.user_id, 

608 'starts_in': time_str, 

609 'message': f'{event.title}: starting {time_str}', 

610 }, 

611 message=event.title, 

612 ) 

613 except Exception: 

614 log.debug('Failed to send notification for calendar alert %s', event.id, exc_info=True) 

615 

616 

617async def _record_run( 

618 automation_id: str, 

619 status: str, 

620 chat_id: str = None, 

621 error: str = None, 

622): 

623 """Insert a run record into automation_run.""" 

624 async with get_async_db() as db: 

625 await AutomationRuns.insert(automation_id, status, chat_id=chat_id, error=error, db=db)