This commit is contained in:
Timothy Jaeryang Baek
2026-04-19 23:38:58 +09:00
parent 29ee53aaa5
commit e5b5a17426
5 changed files with 244 additions and 22 deletions
+135 -18
View File
@@ -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/<feature>.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,