From cf235738f5a44db415012b3b0ebc1f6e752f5439 Mon Sep 17 00:00:00 2001 From: Timothy Jaeryang Baek Date: Thu, 16 Jul 2026 01:27:52 -0400 Subject: [PATCH] refac --- backend/open_webui/events.py | 19 ++- backend/open_webui/routers/channels.py | 33 ++-- backend/open_webui/routers/notifications.py | 23 ++- backend/open_webui/utils/automations.py | 48 ++---- backend/open_webui/utils/notifications.py | 153 ++++++++++++------ src/lib/apis/notifications/index.ts | 17 +- .../chat/Settings/Notifications.svelte | 19 +-- 7 files changed, 186 insertions(+), 126 deletions(-) diff --git a/backend/open_webui/events.py b/backend/open_webui/events.py index a99ef60f67..942841ea3d 100644 --- a/backend/open_webui/events.py +++ b/backend/open_webui/events.py @@ -264,6 +264,11 @@ class EventDefinitions(BaseModel): description='A channel member active state was updated.', message='Channel member active updated', ) + CHANNEL_MESSAGE: EventDefinition = EventDefinition( + name='channel.message', + description='A channel message was posted.', + message='Channel message', + ) CHANNEL_WEBHOOK_CREATED: EventDefinition = EventDefinition( name='channel.webhook.created', description='A channel incoming webhook was created.', @@ -572,6 +577,11 @@ class EventDefinitions(BaseModel): description='A calendar event RSVP was updated.', message='Calendar Event rsvp updated', ) + CALENDAR_ALERT: EventDefinition = EventDefinition( + name='calendar.alert', + description='A calendar event alert was triggered.', + message='Calendar alert', + ) AUTOMATION_CREATED: EventDefinition = EventDefinition( name='automation.created', description='An automation was created.', message='Automation created' ) @@ -641,7 +651,12 @@ EVENT_DEFINITIONS = tuple(getattr(EVENTS, field_name) for field_name in EventDef EVENT_DEFINITIONS_BY_NAME = {definition.name: definition for definition in EVENT_DEFINITIONS} EVENT_CATALOG = tuple(definition.name for definition in EVENT_DEFINITIONS) EVENT_CATALOG_SET = set(EVENT_CATALOG) -CHAT_NOTIFICATION_EVENTS = (EVENTS.CHAT_FINISHED.name, EVENTS.CHAT_FAILED.name) +NOTIFICATION_EVENTS = ( + EVENTS.CHAT_FINISHED.name, + EVENTS.CHAT_FAILED.name, + EVENTS.CHANNEL_MESSAGE.name, + EVENTS.CALENDAR_ALERT.name, +) def get_event_catalog() -> list[dict[str, str]]: @@ -1048,7 +1063,7 @@ def schedule_notification_dispatch(app: Any, event: Event) -> None: class NotificationEventSink: async def handle_event(self, app: Any, event: Event, request: Any | None = None) -> None: - if event.event in CHAT_NOTIFICATION_EVENTS: + if event.event in NOTIFICATION_EVENTS: schedule_notification_dispatch(app, event) diff --git a/backend/open_webui/routers/channels.py b/backend/open_webui/routers/channels.py index 459a238ba1..b40195f110 100644 --- a/backend/open_webui/routers/channels.py +++ b/backend/open_webui/routers/channels.py @@ -51,7 +51,6 @@ from open_webui.utils.models import ( get_all_models, get_filtered_models, ) -from open_webui.utils.webhook import post_webhook from pydantic import BaseModel, field_validator from sqlalchemy.ext.asyncio import AsyncSession @@ -915,7 +914,6 @@ async def get_pinned_channel_messages( async def send_notification(request, channel, message, active_user_ids, db=None): - name = request.app.state.WEBUI_NAME webui_url = await Config.get('webui.url') enable_user_webhooks = await Config.get('ui.enable_user_webhooks') @@ -923,23 +921,28 @@ async def send_notification(request, channel, message, active_user_ids, db=None) # Batch fetch channel members in 1 query (fixes N+1) member_ids = {m.user_id for m in await Channels.get_members_by_channel_id(channel.id, db=db)} + url = f'{webui_url}/channels/{channel.id}' for u in users: if (u.id not in active_user_ids) and u.id in member_ids: if enable_user_webhooks and u.settings: - webhook_url = u.settings.ui.get('notifications', {}).get('webhook_url', None) - if webhook_url: - await post_webhook( - name, - webhook_url, - f'#{channel.name} - {webui_url}/channels/{channel.id}\n\n{message.content}', - { - 'action': 'channel', - 'message': message.content, - 'title': channel.name, - 'url': f'{webui_url}/channels/{channel.id}', - }, - ) + await publish_event( + request, + EVENTS.CHANNEL_MESSAGE, + subject_id=channel.id, + subject_type='channel', + data={ + 'user_id': u.id, + 'channel_id': channel.id, + 'message_id': message.id, + 'sender_id': message.user_id, + 'message': f'#{channel.name} - {url}\n\n{message.content}', + 'content_preview': message.content[:300], + 'title': channel.name, + 'url': url, + }, + message=channel.name, + ) return True diff --git a/backend/open_webui/routers/notifications.py b/backend/open_webui/routers/notifications.py index 9857980dbc..c3ae3ff3b7 100644 --- a/backend/open_webui/routers/notifications.py +++ b/backend/open_webui/routers/notifications.py @@ -3,7 +3,7 @@ from __future__ import annotations from typing import Any from fastapi import APIRouter, Depends, HTTPException, Request, status -from pydantic import BaseModel, Field +from pydantic import BaseModel from open_webui.constants import ERROR_MESSAGES from open_webui.models.config import Config @@ -12,6 +12,7 @@ from open_webui.utils.auth import get_verified_user from open_webui.utils.notifications import ( create_target, delete_target, + get_notification_event_catalog, list_targets, set_default_target, test_target, @@ -23,12 +24,11 @@ router = APIRouter() class NotificationTargetForm(BaseModel): id: str | None = None - type: str = 'webhook' - name: str = Field(default='Webhook') - enabled: bool = True - events: list[str] = Field(default_factory=lambda: ['chat.finished', 'chat.failed']) - delivery: str = 'away' - config: dict[str, Any] = Field(default_factory=dict) + type: str | None = None + enabled: bool | None = None + events: list[str] | None = None + delivery: str | None = None + config: dict[str, Any] | None = None async def _check_notifications_access(user) -> None: @@ -44,10 +44,7 @@ async def _check_notifications_access(user) -> None: @router.get('/events') async def get_notification_events(user=Depends(get_verified_user)): await _check_notifications_access(user) - return [ - {'event': 'chat.finished', 'label': 'Chat finished'}, - {'event': 'chat.failed', 'label': 'Chat failed'}, - ] + return {'events': get_notification_event_catalog()} @router.get('/targets') @@ -60,7 +57,7 @@ async def get_notification_targets(user=Depends(get_verified_user)): async def create_notification_target(form_data: NotificationTargetForm, user=Depends(get_verified_user)): await _check_notifications_access(user) try: - return await create_target(user.id, form_data.model_dump()) + return await create_target(user.id, form_data.model_dump(exclude_unset=True)) except ValueError as e: raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(e)) @@ -81,7 +78,7 @@ async def delete_notification_target(target_id: str, user=Depends(get_verified_u await _check_notifications_access(user) if not await delete_target(user.id, target_id): raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=ERROR_MESSAGES.NOT_FOUND) - return True + return {'ok': True} @router.put('/targets/{target_id}/default') diff --git a/backend/open_webui/utils/automations.py b/backend/open_webui/utils/automations.py index 7823c0207b..b709a63d88 100644 --- a/backend/open_webui/utils/automations.py +++ b/backend/open_webui/utils/automations.py @@ -647,40 +647,24 @@ async def _check_calendar_alerts(app) -> None: except Exception: log.debug(f'Failed to mark event {event.id} as alerted', exc_info=True) - # Send webhook notification if user has one configured + # Send target notification if user has one configured try: - webui_name = getattr(app.state, 'WEBUI_NAME', 'Open WebUI') - enable_user_webhooks = await Config.get('ui.enable_user_webhooks') - - 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, - }, - ) + time_str = f'in {minutes_until} min' if minutes_until > 0 else 'now' + await publish_event( + app, + EVENTS.CALENDAR_ALERT, + subject_id=event.id, + subject_type='calendar.event', + source='scheduler', + data={ + **alert_data, + 'user_id': event.user_id, + 'message': f'{event.title} — starting {time_str}', + }, + message=event.title, + ) except Exception: - log.debug(f'Failed to send webhook for calendar alert {event.id}', exc_info=True) + log.debug(f'Failed to send notification for calendar alert {event.id}', exc_info=True) async def _record_run( diff --git a/backend/open_webui/utils/notifications.py b/backend/open_webui/utils/notifications.py index b11b91f3e2..caa30f537e 100644 --- a/backend/open_webui/utils/notifications.py +++ b/backend/open_webui/utils/notifications.py @@ -1,18 +1,22 @@ from __future__ import annotations import logging +import re import time -import uuid from typing import Any +from urllib.parse import urlparse +from open_webui.events import EVENT_DEFINITIONS_BY_NAME, NOTIFICATION_EVENTS from open_webui.models.config import Config from open_webui.models.users import Users from open_webui.retrieval.web.utils import validate_url from open_webui.utils.webhook import post_webhook -VALID_EVENTS = {'chat.finished', 'chat.failed'} +VALID_EVENTS = set(NOTIFICATION_EVENTS) +LEGACY_EVENTS = {'chat.finished', 'chat.failed'} VALID_DELIVERY = {'away', 'always'} +CHANNEL_MESSAGE_EVENT = 'channel.message' DEFAULT_TARGET_ID = 'webhook' log = logging.getLogger(__name__) @@ -26,12 +30,6 @@ def _normalize_target(target: dict[str, Any], existing: dict[str, Any] | None = if target_type != 'webhook': raise ValueError('Unsupported notification target type') - name = str(target.get('name') or existing.get('name') or 'Webhook').strip() or 'Webhook' - slug = ''.join(ch.lower() if ch.isalnum() else '-' for ch in name).strip('-') - while '--' in slug: - slug = slug.replace('--', '-') - target_id = str(target.get('id') or existing.get('id') or slug[:48] or uuid.uuid4().hex[:8]).strip() - config = dict(existing.get('config') or {}) config.update(target.get('config') or {}) url = str(config.get('url') or '').strip() @@ -42,26 +40,33 @@ def _normalize_target(target: dict[str, Any], existing: dict[str, Any] | None = validate_url(url) config['url'] = url + target_id = str(target.get('id') or existing.get('id') or '').strip() + if not target_id: + hostname = urlparse(url).hostname or 'webhook' + target_id = re.sub(r'[^a-zA-Z0-9_-]+', '-', hostname).strip('-').lower() or 'target' + events = target['events'] if 'events' in target else existing.get('events', []) if events is None: events = [] if not isinstance(events, list): raise ValueError('events must be a list') - events = [str(event) for event in events] - unsupported = [event for event in events if event not in VALID_EVENTS] - if unsupported: - raise ValueError(f'unsupported notification event: {unsupported[0]}') + cleaned_events = [] + for event in events: + event = str(event) + if event not in VALID_EVENTS: + raise ValueError(f'unsupported notification event: {event}') + if event not in cleaned_events: + cleaned_events.append(event) - delivery = str(target.get('delivery') or existing.get('delivery') or 'away') + delivery = str(target.get('delivery') or existing.get('delivery') or 'away').strip() if delivery not in VALID_DELIVERY: raise ValueError('Invalid notification delivery mode') return { 'id': target_id, 'type': target_type, - 'name': name, 'enabled': bool(target.get('enabled', existing.get('enabled', True))), - 'events': events, + 'events': cleaned_events, 'delivery': delivery, 'config': config, 'created_at': int(existing.get('created_at') or now), @@ -69,12 +74,20 @@ def _normalize_target(target: dict[str, Any], existing: dict[str, Any] | None = } -def _public_target(target: dict[str, Any]) -> dict[str, Any]: +def _public_target(target: dict[str, Any], default_target_id: str | None = None) -> dict[str, Any]: config = dict(target.get('config') or {}) - if config.get('url'): - url = str(config['url']) - config['url'] = '********' if len(url) <= 12 else f'{url[:8]}...{url[-4:]}' - return {**target, 'config': config} + url = str(config.pop('url', '') or '') + if url: + parsed = urlparse(url) + if parsed.hostname: + path = parsed.path or '' + suffix = path[-4:] if len(path) > 4 else path + config['url_masked'] = f'{parsed.scheme}://{parsed.hostname}/...{suffix}' + else: + config['url_masked'] = '****' + else: + config['url_masked'] = '' + return {**target, 'config': config, 'is_default': target.get('id') == default_target_id} async def _load_notifications(user_id: str) -> dict[str, Any]: @@ -87,17 +100,15 @@ async def _load_notifications(user_id: str) -> dict[str, Any]: notifications = dict(settings.get('notifications') or {}) targets = notifications.get('targets') + legacy_url = str( + notifications.get('webhook_url') or settings.get('ui', {}).get('notifications', {}).get('webhook_url') or '' + ).strip() + if not isinstance(targets, list) or not targets: - legacy_url = str( - notifications.get('webhook_url') - or settings.get('ui', {}).get('notifications', {}).get('webhook_url') - or '' - ).strip() if legacy_url: target = _normalize_target( { 'id': DEFAULT_TARGET_ID, - 'name': 'Webhook', 'type': 'webhook', 'enabled': True, 'events': sorted(VALID_EVENTS), @@ -105,59 +116,91 @@ async def _load_notifications(user_id: str) -> dict[str, Any]: 'config': {'url': legacy_url}, } ) - notifications = {**notifications, 'targets': [target], 'default_target_id': DEFAULT_TARGET_ID} + notifications = { + **notifications, + 'targets': [target], + 'default_target_id': DEFAULT_TARGET_ID, + 'legacy_notification_events_migrated': True, + } await Users.update_user_settings_by_id(user_id, {'notifications': notifications}) else: notifications['targets'] = [target for target in targets if isinstance(target, dict)] notifications.setdefault( 'default_target_id', notifications['targets'][0].get('id') if notifications['targets'] else None ) + if legacy_url and not notifications.get('legacy_notification_events_migrated'): + changed = False + for target in notifications['targets']: + if ( + target.get('id') == DEFAULT_TARGET_ID + and str((target.get('config') or {}).get('url') or '').strip() == legacy_url + and set(target.get('events') or []) == LEGACY_EVENTS + ): + target['events'] = sorted(VALID_EVENTS) + changed = True + notifications['legacy_notification_events_migrated'] = True + if changed: + await Users.update_user_settings_by_id(user_id, {'notifications': notifications}) return notifications async def list_targets(user_id: str) -> dict[str, Any]: notifications = await _load_notifications(user_id) + default_target_id = notifications.get('default_target_id') return { - 'targets': [_public_target(target) for target in notifications.get('targets') or []], - 'default_target_id': notifications.get('default_target_id'), + 'targets': [_public_target(target, default_target_id) for target in notifications.get('targets') or []], } async def create_target(user_id: str, payload: dict[str, Any]) -> dict[str, Any]: notifications = await _load_notifications(user_id) targets = notifications.get('targets') or [] + has_explicit_id = bool(str(payload.get('id') or '').strip()) target = _normalize_target(payload) - if any(existing.get('id') == target['id'] for existing in targets): - target['id'] = f'{target["id"]}-{uuid.uuid4().hex[:6]}' + if any(str(existing.get('id', '')).lower() == target['id'].lower() for existing in targets): + if has_explicit_id: + raise ValueError('notification target id already exists') + base = target['id'] + suffix = 2 + while any(str(existing.get('id', '')).lower() == target['id'].lower() for existing in targets): + target['id'] = f'{base}-{suffix}' + suffix += 1 targets.append(target) notifications['targets'] = targets notifications.setdefault('default_target_id', target['id']) await Users.update_user_settings_by_id(user_id, {'notifications': notifications}) - return _public_target(target) + return _public_target(target, notifications.get('default_target_id')) async def update_target(user_id: str, target_id: str, payload: dict[str, Any]) -> dict[str, Any]: notifications = await _load_notifications(user_id) targets = notifications.get('targets') or [] for index, existing in enumerate(targets): - if existing.get('id') == target_id: + if str(existing.get('id', '')).lower() == target_id.lower(): updated = _normalize_target({'id': target_id, **payload}, existing=existing) + if any( + idx != index and str(target.get('id', '')).lower() == updated['id'].lower() + for idx, target in enumerate(targets) + ): + raise ValueError('notification target id already exists') targets[index] = updated notifications['targets'] = targets + if str(notifications.get('default_target_id') or '').lower() == target_id.lower(): + notifications['default_target_id'] = updated['id'] await Users.update_user_settings_by_id(user_id, {'notifications': notifications}) - return _public_target(updated) + return _public_target(updated, notifications.get('default_target_id')) raise ValueError('Notification target not found') async def delete_target(user_id: str, target_id: str) -> bool: notifications = await _load_notifications(user_id) targets = notifications.get('targets') or [] - next_targets = [target for target in targets if target.get('id') != target_id] + next_targets = [target for target in targets if str(target.get('id', '')).lower() != target_id.lower()] if len(next_targets) == len(targets): return False notifications['targets'] = next_targets - if notifications.get('default_target_id') == target_id: + if str(notifications.get('default_target_id') or '').lower() == target_id.lower(): notifications['default_target_id'] = next_targets[0].get('id') if next_targets else None await Users.update_user_settings_by_id(user_id, {'notifications': notifications}) return True @@ -165,19 +208,33 @@ async def delete_target(user_id: str, target_id: str) -> bool: async def set_default_target(user_id: str, target_id: str) -> dict[str, Any]: notifications = await _load_notifications(user_id) - if not any(target.get('id') == target_id for target in notifications.get('targets') or []): - raise ValueError('Notification target not found') - notifications['default_target_id'] = target_id - await Users.update_user_settings_by_id(user_id, {'notifications': notifications}) - return await list_targets(user_id) + for target in notifications.get('targets') or []: + if str(target.get('id', '')).lower() == target_id.lower(): + notifications['default_target_id'] = target['id'] + await Users.update_user_settings_by_id(user_id, {'notifications': notifications}) + return _public_target(target, target['id']) + raise ValueError('Notification target not found') + + +def get_notification_event_catalog() -> list[dict[str, str]]: + return [ + { + 'event': event_name, + 'label': EVENT_DEFINITIONS_BY_NAME[event_name].message or event_name, + 'description': EVENT_DEFINITIONS_BY_NAME[event_name].description or '', + } + for event_name in NOTIFICATION_EVENTS + ] def _find_target(notifications: dict[str, Any], target: str = '') -> dict[str, Any] | None: targets = notifications.get('targets') or [] target = target.strip() target_id = target or str(notifications.get('default_target_id') or '') + if not target_id: + return None for item in targets: - if item.get('id') == target_id or item.get('name') == target: + if str(item.get('id', '')).lower() == target_id.lower(): return item return None @@ -235,7 +292,7 @@ async def dispatch_notification_event(app: Any, event: Any) -> None: for user_id in event_user_ids(event): try: notifications = await _load_notifications(user_id) - is_active = await Users.is_user_active(user_id) + is_active = False if event.event == CHANNEL_MESSAGE_EVENT else await Users.is_user_active(user_id) for target in notifications.get('targets') or []: if not target.get('enabled', True): @@ -246,8 +303,14 @@ async def dispatch_notification_event(app: Any, event: Any) -> None: continue data = event.model_dump() - title = event.message or data.get('message') or event.event - message = str((event.data or {}).get('message') or title) + definition = EVENT_DEFINITIONS_BY_NAME.get(event.event) + title = event.message or (definition.message if definition else event.event) + message = str( + (event.data or {}).get('message') + or (event.data or {}).get('preview') + or (event.data or {}).get('content_preview') + or title + ) await _send_webhook(app_name, target, message, data, title) except Exception: log.exception('Notification delivery failed for user %s and event %s', user_id, event.event) diff --git a/src/lib/apis/notifications/index.ts b/src/lib/apis/notifications/index.ts index ec09cdc670..3b427b33fe 100644 --- a/src/lib/apis/notifications/index.ts +++ b/src/lib/apis/notifications/index.ts @@ -3,17 +3,24 @@ import { WEBUI_API_BASE_URL } from '$lib/constants'; export type NotificationTarget = { id: string; type: 'webhook'; - name: string; + is_default?: boolean; enabled: boolean; events: string[]; delivery: 'away' | 'always'; config: { url?: string; + url_masked?: string; }; created_at?: number; updated_at?: number; }; +export type NotificationEvent = { + event: string; + label: string; + description?: string; +}; + const jsonRequest = async (url: string, token: string, method = 'GET', body?: object) => { let error = null; @@ -43,12 +50,14 @@ const jsonRequest = async (url: string, token: string, method = 'GET', body?: ob return res; }; -export const getNotificationEvents = async (token: string) => - jsonRequest(`${WEBUI_API_BASE_URL}/notifications/events`, token); +export const getNotificationEvents = async (token: string): Promise => { + const data = await jsonRequest(`${WEBUI_API_BASE_URL}/notifications/events`, token); + return data?.events ?? data ?? []; +}; export const getNotificationTargets = async ( token: string -): Promise<{ targets: NotificationTarget[]; default_target_id: string | null }> => +): Promise<{ targets: NotificationTarget[] }> => jsonRequest(`${WEBUI_API_BASE_URL}/notifications/targets`, token); export const createNotificationTarget = async ( diff --git a/src/lib/components/chat/Settings/Notifications.svelte b/src/lib/components/chat/Settings/Notifications.svelte index cf42ec7e43..9000691ad0 100644 --- a/src/lib/components/chat/Settings/Notifications.svelte +++ b/src/lib/components/chat/Settings/Notifications.svelte @@ -25,7 +25,6 @@ let notificationEnabled = false; let notificationSound = true; let targets: NotificationTarget[] = []; - let defaultTargetId: string | null = null; let events: { event: string; label: string; description?: string }[] = [ { event: 'chat.finished', @@ -93,7 +92,6 @@ } if (targetResult.status === 'fulfilled') { targets = targetResult.value.targets ?? []; - defaultTargetId = targetResult.value.default_target_id ?? null; } else { toast.error(`${targetResult.reason}`); } @@ -139,7 +137,7 @@ try { const id = form.id.trim(); const payload: Partial = { - ...(id ? { id, name: id } : {}), + ...(id ? { id } : {}), type: 'webhook', enabled: form.enabled, events: form.events, @@ -259,7 +257,7 @@ {$i18n.t('Webhook')} - {#if target.id === defaultTargetId} + {#if target.is_default} {$i18n.t('Default')} @@ -268,7 +266,7 @@
- {target.config?.url} + {target.config?.url_masked}
{$i18n.t('Send Test')} - {#if target.id !== defaultTargetId} + {#if !target.is_default} -
-
{$i18n.t('Target ID for notify')}