From e5b5a174265d6710e986f6534ee7e3b2923233be Mon Sep 17 00:00:00 2001 From: Timothy Jaeryang Baek Date: Sun, 19 Apr 2026 23:38:58 +0900 Subject: [PATCH] refac --- backend/open_webui/main.py | 4 +- backend/open_webui/models/calendar.py | 51 ++++++ backend/open_webui/utils/automations.py | 153 +++++++++++++++--- .../calendar/CalendarEventModal.svelte | 26 ++- src/routes/+layout.svelte | 32 ++++ 5 files changed, 244 insertions(+), 22 deletions(-) diff --git a/backend/open_webui/main.py b/backend/open_webui/main.py index f9b21e659..d1a0bddce 100644 --- a/backend/open_webui/main.py +++ b/backend/open_webui/main.py @@ -674,9 +674,9 @@ async def lifespan(app: FastAPI): asyncio.create_task(periodic_usage_pool_cleanup()) asyncio.create_task(periodic_session_pool_cleanup()) - from open_webui.utils.automations import automation_worker_loop + from open_webui.utils.automations import scheduler_worker_loop - asyncio.create_task(automation_worker_loop(app)) + asyncio.create_task(scheduler_worker_loop(app)) if app.state.config.ENABLE_BASE_MODELS_CACHE: try: diff --git a/backend/open_webui/models/calendar.py b/backend/open_webui/models/calendar.py index 9d2d71a45..9afa7c15e 100644 --- a/backend/open_webui/models/calendar.py +++ b/backend/open_webui/models/calendar.py @@ -725,6 +725,57 @@ class CalendarEventTable: await db.commit() return await self._to_event_model(event, db=db) + async def get_upcoming_events( + self, + now_ns: int, + default_lookahead_ns: int, + db: Optional[AsyncSession] = None, + ) -> list[tuple[CalendarEventModel, Optional[str]]]: + """Events starting between now and now + lookahead, for alert processing. + + Per-event lookahead is read from meta.alert_minutes (falls back to + default_lookahead_ns). Returns (event, user_timezone) pairs. + """ + from open_webui.models.users import User as UserRow + + # Use the maximum possible lookahead (60 min) to cast a wide net; + # per-event filtering happens in Python after fetching. + max_lookahead_ns = max(default_lookahead_ns, 60 * 60 * 1_000_000_000) + upper = now_ns + max_lookahead_ns + + async with get_async_db_context(db) as db: + result = await db.execute( + select(CalendarEvent, UserRow.timezone) + .outerjoin(UserRow, UserRow.id == CalendarEvent.user_id) + .filter( + CalendarEvent.is_cancelled == False, + CalendarEvent.start_at >= now_ns, + CalendarEvent.start_at <= upper, + ) + ) + rows = result.all() + + events = [] + for event, tz in rows: + model = CalendarEventModel.model_validate(event) + # Determine per-event alert window + alert_minutes = None + if model.meta and 'alert_minutes' in model.meta: + alert_minutes = model.meta['alert_minutes'] + + if alert_minutes is not None: + if alert_minutes < 0: + # alert_minutes < 0 means "no alert" + continue + event_lookahead_ns = alert_minutes * 60 * 1_000_000_000 + else: + event_lookahead_ns = default_lookahead_ns + + if model.start_at <= now_ns + event_lookahead_ns: + events.append((model, tz)) + + return events + async def delete_event_by_id(self, id: str, db: Optional[AsyncSession] = None) -> bool: try: async with get_async_db_context(db) as db: diff --git a/backend/open_webui/utils/automations.py b/backend/open_webui/utils/automations.py index ac1f4df69..41486cd21 100644 --- a/backend/open_webui/utils/automations.py +++ b/backend/open_webui/utils/automations.py @@ -1,11 +1,16 @@ """ -Automation utilities. +Automation utilities and unified scheduler. -RRULE helpers, worker loop, and execution logic. +RRULE helpers, scheduler worker loop, and execution logic. Follows the utils/.py pattern (cf. utils/channels.py, utils/task.py). +The scheduler_worker_loop handles all time-based background work: + - Automation execution (claim_due → execute) + - Calendar event alerts (upcoming events → socket + webhook notifications) + Environment: - AUTOMATION_POLL_INTERVAL – seconds between polls (default: 10) + SCHEDULER_POLL_INTERVAL – seconds between polls (default: 10) + CALENDAR_ALERT_LOOKAHEAD_MINUTES – default alert window (default: 5) """ import asyncio @@ -31,7 +36,8 @@ from open_webui.internal.db import get_async_db log = logging.getLogger(__name__) -AUTOMATION_POLL_INTERVAL = int(os.getenv('AUTOMATION_POLL_INTERVAL', '10')) +SCHEDULER_POLL_INTERVAL = int(os.getenv('SCHEDULER_POLL_INTERVAL', os.getenv('AUTOMATION_POLL_INTERVAL', '10'))) +CALENDAR_ALERT_LOOKAHEAD_MINUTES = int(os.getenv('CALENDAR_ALERT_LOOKAHEAD_MINUTES', '10')) #################### @@ -117,30 +123,49 @@ def rrule_interval_seconds(s: str) -> Optional[int]: ############################ +# Keep the old name as an alias so any stale imports still work. async def automation_worker_loop(app) -> None: - """Poll for due automations, claim, fire-and-forget execute. + """Deprecated alias — use scheduler_worker_loop.""" + await scheduler_worker_loop(app) + + +async def scheduler_worker_loop(app) -> None: + """Unified background scheduler for all time-based work. + + Handles: + 1. Automation execution (ENABLE_AUTOMATIONS) + 2. Calendar event alerts (ENABLE_CALENDAR) Runs on every instance. Poll interval is configurable via - AUTOMATION_POLL_INTERVAL env var (default: 10 seconds). + SCHEDULER_POLL_INTERVAL env var (default: 10 seconds). """ - log.info(f'Automation worker started (poll interval: {AUTOMATION_POLL_INTERVAL}s)') + log.info(f'Scheduler worker started (poll interval: {SCHEDULER_POLL_INTERVAL}s)') while True: try: - if not getattr(app.state.config, 'ENABLE_AUTOMATIONS', False): - await asyncio.sleep(AUTOMATION_POLL_INTERVAL) - continue + # ── Automations ── + if getattr(app.state.config, 'ENABLE_AUTOMATIONS', False): + try: + async with get_async_db() as db: + batch = await Automations.claim_due(int(time.time_ns()), limit=10, db=db) + if batch: + log.info(f'Claimed {len(batch)} due automation(s)') + for automation in batch: + asyncio.create_task(execute_automation(app, automation)) + except Exception: + log.exception('Scheduler: automation error') + + # ── Calendar Alerts ── + if getattr(app.state.config, 'ENABLE_CALENDAR', False): + try: + await _check_calendar_alerts(app) + except Exception: + log.exception('Scheduler: calendar alert error') - async with get_async_db() as db: - batch = await Automations.claim_due(int(time.time_ns()), limit=10, db=db) - if batch: - log.info(f'Claimed {len(batch)} due automation(s)') - for automation in batch: - asyncio.create_task(execute_automation(app, automation)) except Exception: - log.exception('Automation worker error') + log.exception('Scheduler worker error') # Jitter to spread load across instances - await asyncio.sleep(AUTOMATION_POLL_INTERVAL + random.uniform(0, 2)) + await asyncio.sleep(SCHEDULER_POLL_INTERVAL + random.uniform(0, 2)) ########################## @@ -433,6 +458,98 @@ async def execute_automation(app, automation: AutomationModel) -> None: #################### +async def _check_calendar_alerts(app) -> None: + """Check for upcoming calendar events and send alert notifications. + + De-duplication is DB-backed via meta.alerted_at — survives restarts + and works across multiple instances. + """ + from open_webui.models.calendar import CalendarEvents, CalendarEventUpdateForm + from open_webui.socket.main import sio + + now_ns = int(time.time_ns()) + default_lookahead_ns = CALENDAR_ALERT_LOOKAHEAD_MINUTES * 60 * 1_000_000_000 + + async with get_async_db() as db: + upcoming = await CalendarEvents.get_upcoming_events(now_ns, default_lookahead_ns, db=db) + + if not upcoming: + return + + for event, user_tz in upcoming: + # Skip if already alerted for this start time + if event.meta and event.meta.get('alerted_at'): + continue + + # Compute minutes until event starts + minutes_until = max(0, int((event.start_at - now_ns) / (60 * 1_000_000_000))) + + alert_data = { + 'event_id': event.id, + 'title': event.title, + 'description': event.description or '', + 'start_at': event.start_at, + 'minutes_until': minutes_until, + 'calendar_id': event.calendar_id, + 'location': event.location or '', + } + + await sio.emit( + 'events', + { + 'data': { + 'type': 'calendar:alert', + 'data': alert_data, + }, + }, + room=f'user:{event.user_id}', + ) + + # Mark as alerted in DB so it survives restarts / multi-instance + try: + await CalendarEvents.update_event_by_id( + event.id, + CalendarEventUpdateForm(meta={'alerted_at': now_ns}), + ) + except Exception: + log.debug(f'Failed to mark event {event.id} as alerted', exc_info=True) + + # Send webhook notification if user has one configured + try: + webui_name = getattr(app.state, 'WEBUI_NAME', 'Open WebUI') + enable_user_webhooks = getattr(app.state.config, 'ENABLE_USER_WEBHOOKS', False) + + if enable_user_webhooks: + user = await Users.get_user_by_id(event.user_id) + if user and user.settings: + webhook_url = ( + user.settings.get('ui', {}).get('notifications', {}).get('webhook_url', None) + if isinstance(user.settings, dict) + else getattr(getattr(user.settings, 'ui', None), 'get', lambda *a: None)( + 'notifications', {} + ).get('webhook_url', None) + if hasattr(user.settings, 'ui') + else None + ) + if webhook_url: + from open_webui.utils.webhook import post_webhook + + time_str = f'in {minutes_until} min' if minutes_until > 0 else 'now' + await post_webhook( + webui_name, + webhook_url, + f'{event.title} — starting {time_str}', + { + 'action': 'calendar_alert', + 'title': event.title, + 'minutes_until': minutes_until, + 'event_id': event.id, + }, + ) + except Exception: + log.debug(f'Failed to send webhook for calendar alert {event.id}', exc_info=True) + + async def _record_run( automation_id: str, status: str, diff --git a/src/lib/components/calendar/CalendarEventModal.svelte b/src/lib/components/calendar/CalendarEventModal.svelte index 020a89f92..ad7bf9b0d 100644 --- a/src/lib/components/calendar/CalendarEventModal.svelte +++ b/src/lib/components/calendar/CalendarEventModal.svelte @@ -32,6 +32,7 @@ let endTime = ''; let allDay = false; let location = ''; + let alertMinutes: number = 10; let loading = false; let showDeleteConfirmDialog = false; @@ -60,6 +61,7 @@ endTime = event.end_at ? nsToTimeStr(event.end_at) : ''; allDay = event.all_day; location = event.location || ''; + alertMinutes = event.meta?.alert_minutes ?? 10; } else { title = ''; description = ''; @@ -80,6 +82,7 @@ } allDay = false; location = ''; + alertMinutes = 10; } } @@ -104,7 +107,8 @@ start_at: startNs, end_at: endNs, all_day: allDay, - location: location.trim() || undefined + location: location.trim() || undefined, + meta: { alert_minutes: alertMinutes } }); if (result) { toast.success($i18n.t('Event updated')); @@ -119,7 +123,8 @@ start_at: startNs, end_at: endNs, all_day: allDay, - location: location.trim() || undefined + location: location.trim() || undefined, + meta: { alert_minutes: alertMinutes } }; const result = await createCalendarEvent(localStorage.token, form); if (result) { @@ -212,6 +217,23 @@ /> + +
+
{$i18n.t('Reminder')}
+ +
+
{$i18n.t('Description')}
diff --git a/src/routes/+layout.svelte b/src/routes/+layout.svelte index f62aa20ca..e8cb71b8c 100644 --- a/src/routes/+layout.svelte +++ b/src/routes/+layout.svelte @@ -439,6 +439,38 @@ const type = event?.data?.type ?? null; const data = event?.data?.data ?? null; + // Calendar alerts are not chat-scoped — handle before chat_id checks + if (type === 'calendar:alert' && data) { + const timeStr = + data.minutes_until <= 0 + ? $i18n.t('Starting now') + : data.minutes_until === 1 + ? $i18n.t('Starting in 1 minute') + : $i18n.t('Starting in {{count}} minutes', { count: data.minutes_until }); + + toast.custom(NotificationToast, { + componentProps: { + onClick: () => { + goto('/calendar'); + }, + title: data.title, + content: timeStr + }, + duration: 30000, + unstyled: true + }); + + if ($isLastActiveTab) { + if ($settings?.notificationEnabled ?? false) { + new Notification(`${data.title} • Open WebUI`, { + body: timeStr, + icon: `${WEBUI_BASE_URL}/static/favicon.png` + }); + } + } + return; + } + if ((event.chat_id !== $chatId && !$temporaryChatEnabled) || isFocused) { if (type === 'chat:completion') { const { done, content, title } = data;