"""The event engine's clock: walks every enabled event trigger, polls its adapter on that source's own cadence, and feeds resulting events to the dispatcher. One adaptive-sleep loop (same shape as the workflow scheduler's), with per-trigger in-flight guards so a slow adapter can't double-poll itself or stall its neighbors.""" import asyncio import logging import random import time from typing import Awaitable, Callable, Dict, List, Optional, Set, Tuple from backend.apps.events import dispatcher, stores from backend.apps.events.adapters.agent_check import agent_check from backend.apps.events.adapters.file_signal import start_file_signal from backend.apps.events.adapters.file_watch import file_watch from backend.apps.events.adapters.stream_watch import run_stream_source from backend.apps.events.adapters.web_watch import web_watch from backend.apps.events.models import CustomEventSource, Event, EventLogEntry, EventTriggerConfig, FileWatchSource, StreamSource from backend.apps.workflows.models import Workflow logger = logging.getLogger(__name__) # "custom" is deliberately absent: those triggers are push-only via /api/events/ingest. ADAPTERS: Dict[str, Callable[..., Awaitable[Tuple[List[Event], Dict]]]] = { "file": file_watch, "web": web_watch, "agent": agent_check, } p_loop_task: Optional["asyncio.Task"] = None p_wake = asyncio.Event() p_next_poll: Dict[str, float] = {} p_inflight: Set[str] = set() # Adaptive cadence (poll_seconds=0 triggers): (default, floor, ceiling) per kind. Events halve the interval toward the floor; 5 straight quiet polls stretch it 1.5x toward the ceiling. PACE_BOUNDS: Dict[str, Tuple[float, float, float]] = { "file": (15.0, 5.0, 60.0), "web": (300.0, 60.0, 1800.0), "agent": (900.0, 300.0, 21600.0), } PACE_QUIET_POLLS = 5 p_pace_interval: Dict[str, float] = {} p_pace_quiet: Dict[str, int] = {} def effective_poll_seconds(trigger: EventTriggerConfig) -> float: fixed = float(getattr(trigger.source, "poll_seconds", 0) or 0) if fixed > 0: return fixed default, lo, hi = PACE_BOUNDS.get(trigger.source.kind, (300.0, 60.0, 3600.0)) return min(max(p_pace_interval.get(trigger.id, default), lo), hi) def pace_update(trigger: EventTriggerConfig, event_count: int) -> None: if float(getattr(trigger.source, "poll_seconds", 0) or 0) > 0: return default, lo, hi = PACE_BOUNDS.get(trigger.source.kind, (300.0, 60.0, 3600.0)) current = p_pace_interval.get(trigger.id, default) if event_count > 0: p_pace_interval[trigger.id] = max(lo, current / 2) p_pace_quiet[trigger.id] = 0 else: quiet = p_pace_quiet.get(trigger.id, 0) + 1 if quiet >= PACE_QUIET_POLLS: p_pace_interval[trigger.id] = min(hi, current * 1.5) quiet = 0 p_pace_quiet[trigger.id] = quiet # Held-open live sources (kqueue file signals, SSE streams): trigger_id -> (config signature, stopper). p_live_handles: Dict[str, Tuple[str, Callable[[], None]]] = {} def kick() -> None: p_wake.set() def live_source_count() -> int: """How many held-open sources (file signals, streams) are currently running.""" return len(p_live_handles) def mark_due(trigger_id: str) -> None: """Schedule this trigger's next poll immediately (used with kick()).""" p_next_poll[trigger_id] = 0.0 def reset_state() -> None: """Test seam: forget all per-trigger poll bookkeeping and stop live sources.""" global p_loop_task p_loop_task = None p_next_poll.clear() p_inflight.clear() p_pace_interval.clear() p_pace_quiet.clear() for _, stop in p_live_handles.values(): try: stop() except Exception: pass p_live_handles.clear() def p_live_triggers() -> List[Tuple[Workflow, EventTriggerConfig, float]]: """(workflow, trigger, poll_seconds) for every pollable trigger; push-only sources are excluded here.""" from backend.apps.workflows import storage out: List[Tuple[Workflow, EventTriggerConfig, float]] = [] for wf in storage.list_workflows(): for trig in wf.event_triggers: source = trig.source if not trig.enabled or isinstance(source, (CustomEventSource, StreamSource)) or source.kind not in ADAPTERS: continue out.append((wf, trig, effective_poll_seconds(trig))) return out async def p_poll_one(wf: Workflow, trigger: EventTriggerConfig) -> None: workflow_id = wf.id try: fetch = ADAPTERS[trigger.source.kind] cursor = stores.load_cursor(trigger.id) # The agent adapter needs its parent workflow (model, approvals, dashboard); the structural adapters stay pure. if trigger.source.kind == "agent": events, new_cursor = await fetch(trigger.source, cursor, wf) else: events, new_cursor = await fetch(trigger.source, cursor) stores.save_cursor(trigger.id, new_cursor) stores.clear_poll_failures(trigger.id) pace_update(trigger, len(events)) if events: await dispatcher.ingest(workflow_id, trigger, events) except Exception as e: logger.warning("poll failed for trigger %s (%s): %s", trigger.id, trigger.source.kind, e) try: # Exponential backoff on repeated failures: a broken site/model can't burn quota at full cadence, and the log says so instead of dying silently. failures = stores.record_poll_failure(trigger.id, str(e)) # Third straight failure: try to fix it ourselves before the attention surface asks the user. if failures == 3 and trigger.source.kind in ("web", "stream"): asyncio.create_task(p_heal_and_wake(wf, trigger)) base = effective_poll_seconds(trigger) backoff = min(base * (2 ** min(failures, 5)), 21600.0) p_next_poll[trigger.id] = time.monotonic() + backoff note = f" (failure {failures} in a row; next try in ~{int(backoff / 60) or 1}m)" if failures >= 2 else "" stores.append_log(workflow_id, EventLogEntry( trigger_id=trigger.id, kind="error", summary=f"Poll failed: {str(e)[:180]}{note}", )) except Exception: pass finally: p_inflight.discard(trigger.id) async def p_heal_and_wake(wf: Workflow, trigger: EventTriggerConfig) -> None: try: from backend.apps.events.adapters.heal_trigger import attempt_heal if await attempt_heal(wf.id, trigger): mark_due(trigger.id) kick() except Exception: logger.debug("self-heal attempt errored", exc_info=True) def reconcile_live_sources() -> None: """Start/stop held-open sources to match the current trigger set. File signals make the diff poll instant; stream tasks own an SSE connection outright.""" from backend.apps.workflows import storage want: Dict[str, Tuple[str, Workflow, EventTriggerConfig]] = {} for wf in storage.list_workflows(): for trig in wf.event_triggers: if not trig.enabled: continue if isinstance(trig.source, (FileWatchSource, StreamSource)): want[trig.id] = (trig.source.model_dump_json(), wf, trig) for trigger_id in list(p_live_handles.keys()): signature, stop = p_live_handles[trigger_id] if trigger_id not in want or want[trigger_id][0] != signature: try: stop() except Exception: pass del p_live_handles[trigger_id] for trigger_id, (signature, wf, trig) in want.items(): if trigger_id in p_live_handles: continue source = trig.source if isinstance(source, FileWatchSource): def p_on_change(tid: str = trigger_id) -> None: mark_due(tid) kick() stop = start_file_signal(source.path, p_on_change) if stop is not None: p_live_handles[trigger_id] = (signature, stop) elif isinstance(source, StreamSource) and source.url.strip(): task = asyncio.create_task(run_stream_source(wf.id, trig, source)) p_live_handles[trigger_id] = (signature, task.cancel) def tick() -> None: from backend.apps.workflows import storage try: reconcile_live_sources() except Exception: logger.exception("live-source reconcile error") # The global "pause all" switch holds event polling too; the cursor diff catches net changes at resume. if storage.get_paused(): return now = time.monotonic() for wf, trig, poll_seconds in p_live_triggers(): if p_next_poll.get(trig.id, 0.0) <= now and trig.id not in p_inflight: # Jitter so logged-in polls aren't metronomic (a bot tell) and many triggers spread out. p_next_poll[trig.id] = now + poll_seconds * random.uniform(0.9, 1.1) p_inflight.add(trig.id) asyncio.create_task(p_poll_one(wf, trig)) def p_seconds_until_next() -> float: now = time.monotonic() soonest: Optional[float] = None for _, trig, _ in p_live_triggers(): nxt = p_next_poll.get(trig.id, now) if soonest is None or nxt < soonest: soonest = nxt if soonest is None: return 30.0 return max(1.0, min(soonest - now, 30.0)) async def p_loop() -> None: logger.info("event engine poll loop started") while True: try: tick() except Exception: logger.exception("event poll tick error") try: await asyncio.wait_for(p_wake.wait(), timeout=p_seconds_until_next()) except asyncio.TimeoutError: pass p_wake.clear() async def start_event_engine() -> None: global p_loop_task, p_wake from backend.apps.workflows import storage if p_loop_task is not None and not p_loop_task.done(): return p_wake = asyncio.Event() all_workflows = storage.list_workflows() + storage.list_deleted_workflows() all_trigger_ids = [t.id for wf in all_workflows for t in wf.event_triggers] try: stores.sweep_stale_state(all_trigger_ids, [wf.id for wf in all_workflows]) except Exception: logger.debug("event state sweep failed", exc_info=True) # Events buffered at last quit resume their coalesce window now. for wf in storage.list_workflows(): for trig in wf.event_triggers: if trig.enabled: restored = dispatcher.restore_pending(wf.id, trig) if restored: logger.info("restored %d pending event(s) for trigger %s", restored, trig.id) p_loop_task = asyncio.create_task(p_loop()) async def stop_event_engine() -> None: global p_loop_task # Cancel the loop before the dispatcher so a mid-cancel tick can't schedule fresh flushes. if p_loop_task is not None: p_loop_task.cancel() try: await p_loop_task except (asyncio.CancelledError, Exception): pass p_loop_task = None dispatcher.stop() reset_state()