This commit is contained in:
Timothy Jaeryang Baek
2026-06-25 14:37:05 +01:00
parent 5576e6ed8a
commit e124c2656a
4 changed files with 130 additions and 6 deletions
+66 -3
View File
@@ -1,10 +1,13 @@
from __future__ import annotations
import asyncio
import inspect
import logging
import time
import uuid
from dataclasses import asdict, dataclass, field
from enum import StrEnum
from types import SimpleNamespace
from typing import Any
from open_webui.env import VERSION
@@ -436,7 +439,7 @@ def build_event(
class WebhookEventSink:
async def handle_event(self, app: Any, event: Event) -> None:
async def handle_event(self, app: Any, event: Event, request: Any | None = None) -> None:
name = getattr(getattr(app, 'state', None), 'WEBUI_NAME', 'Open WebUI')
subject = event.subject or {}
subject_id = subject.get('id')
@@ -452,7 +455,66 @@ class WebhookEventSink:
log.exception('Event webhook failed for %s', webhook.get('id'))
EVENT_SINKS = [WebhookEventSink()]
async def dispatch_event_functions(app: Any, event: Event, request: Any | None = None) -> None:
from open_webui.models.functions import Functions
from open_webui.utils.plugin import get_function_module_from_cache
context = request or SimpleNamespace(app=app)
event_payload = event.model_dump()
try:
event_functions = await Functions.get_functions_by_type('event', active_only=True)
except Exception:
log.exception('Event functions could not be loaded for %s', event.event)
return
for function in event_functions:
try:
function_module, _, _ = await get_function_module_from_cache(context, function.id, function=function)
handler = getattr(function_module, 'event', None)
if not handler:
continue
if hasattr(function_module, 'valves') and hasattr(function_module, 'Valves'):
valves = await Functions.get_function_valves_by_id(function.id)
function_module.valves = function_module.Valves(**(valves if valves else {}))
sig = inspect.signature(handler)
accepts_kwargs = any(param.kind == inspect.Parameter.VAR_KEYWORD for param in sig.parameters.values())
extra_params = {
'event': event_payload,
'__id__': function.id,
'__event__': event,
'__event_id__': event.id,
'__event_name__': event.event,
'__app__': app,
'__request__': request,
}
params = {
key: value for key, value in extra_params.items() if accepts_kwargs or key in sig.parameters
}
if inspect.iscoroutinefunction(handler):
await handler(**params)
else:
handler(**params)
except Exception:
log.exception('Event function failed for %s', function.id)
def schedule_event_function_dispatch(app: Any, event: Event, request: Any | None = None) -> None:
try:
asyncio.create_task(dispatch_event_functions(app, event, request))
except RuntimeError:
log.exception('Event functions could not be scheduled for %s', event.event)
class EventFunctionSink:
async def handle_event(self, app: Any, event: Event, request: Any | None = None) -> None:
schedule_event_function_dispatch(app, event, request)
EVENT_SINKS = [EventFunctionSink(), WebhookEventSink()]
async def publish_event(
@@ -467,6 +529,7 @@ async def publish_event(
message: str | None = None,
) -> None:
app = getattr(request_or_app, 'app', request_or_app)
request = request_or_app if hasattr(request_or_app, 'app') else None
event_payload = build_event(
request_or_app,
event,
@@ -480,6 +543,6 @@ async def publish_event(
for sink in EVENT_SINKS:
try:
await sink.handle_event(app, event_payload)
await sink.handle_event(app, event_payload, request=request)
except Exception:
log.exception('Event sink failed for %s', event.value)
+2
View File
@@ -285,6 +285,8 @@ async def load_function_module_by_id(function_id: str, content: str | None = Non
return module.Filter(), 'filter', frontmatter
elif hasattr(module, 'Action'):
return module.Action(), 'action', frontmatter
elif hasattr(module, 'Event'):
return module.Event(), 'event', frontmatter
else:
raise Exception('No Function class found in the module')
except Exception as e:
+2 -1
View File
@@ -410,7 +410,8 @@
items={[
{ value: 'pipe', label: $i18n.t('Pipe') },
{ value: 'filter', label: $i18n.t('Filter') },
{ value: 'action', label: $i18n.t('Action') }
{ value: 'action', label: $i18n.t('Action') },
{ value: 'event', label: $i18n.t('Event') }
]}
/>
</div>
@@ -41,7 +41,8 @@
}
let codeEditor;
let boilerplate = `"""
let starterType = 'filter';
const filterBoilerplate = `"""
title: Example Filter
author: open-webui
author_url: https://github.com/open-webui
@@ -109,6 +110,52 @@ class Filter:
return body
`;
const eventBoilerplate = `"""
title: Example Event
author: open-webui
author_url: https://github.com/open-webui
funding_url: https://github.com/open-webui
version: 0.1
"""
from pydantic import BaseModel
class Event:
class Valves(BaseModel):
pass
def __init__(self):
self.valves = self.Valves()
async def event(
self,
event: dict,
__event_id__: str = None,
__event_name__: str = None,
__id__: str = None,
__app__=None,
__request__=None,
):
print(f"event:{__name__}")
print(f"event:id:{__event_id__}")
print(f"event:name:{__event_name__}")
print(f"event:payload:{event}")
`;
let boilerplate = filterBoilerplate;
/** @param {'filter' | 'event'} type */
const setStarterType = (type) => {
starterType = type;
boilerplate = type === 'event' ? eventBoilerplate : filterBoilerplate;
content = boilerplate;
_content = boilerplate;
};
/** @param {string} value */
const selectStarterType = (value) => {
setStarterType(value === 'event' ? 'event' : 'filter');
};
const _boilerplate = `from pydantic import BaseModel
from typing import Optional, Union, Generator, Iterator
@@ -328,7 +375,18 @@ class Pipe:
</Tooltip>
</div>
<div>
<div class="flex items-center gap-2">
{#if !edit}
<select
class="text-xs bg-transparent border border-gray-100 dark:border-gray-800 rounded-lg px-2 py-1 outline-hidden"
bind:value={starterType}
on:change={(event) => selectStarterType(event.currentTarget.value)}
aria-label={$i18n.t('Function starter')}
>
<option value="filter">{$i18n.t('Filter')}</option>
<option value="event">{$i18n.t('Event')}</option>
</select>
{/if}
<Badge type="muted" content={$i18n.t('Function')} />
</div>
</div>