// Package triggers — engine.go // // v0.2.2: Core trigger engine. Manages extension-declared triggers // (event bus subscriptions + webhook receivers) and user-created // scheduled tasks (cron). All three converge to sandbox.Runner.CallEntryPoint. package triggers import ( "context" "encoding/json" "log" "sync" "time" "github.com/robfig/cron/v3" "switchboard-core/events" "switchboard-core/models" "switchboard-core/sandbox" "switchboard-core/store" ) // Engine orchestrates trigger lifecycle: loading, firing, and cleanup. type Engine struct { stores store.Stores runner *sandbox.Runner bus *events.Bus cron *cron.Cron mu sync.RWMutex unsubs map[string]func() // trigger_id → bus unsubscribe cronIDs map[string]cron.EntryID // scheduled_task_id → cron entry ctx context.Context cancel context.CancelFunc } // New creates a trigger engine. Call Start() to begin processing. func New(stores store.Stores, runner *sandbox.Runner, bus *events.Bus) *Engine { return &Engine{ stores: stores, runner: runner, bus: bus, cron: cron.New(cron.WithSeconds()), unsubs: make(map[string]func()), cronIDs: make(map[string]cron.EntryID), } } // Start loads all enabled triggers and scheduled tasks, then begins execution. func (e *Engine) Start(ctx context.Context) error { e.ctx, e.cancel = context.WithCancel(ctx) // Load event triggers if e.stores.Triggers != nil { eventTriggers, err := e.stores.Triggers.ListEnabledByType(ctx, models.TriggerTypeEvent) if err != nil { log.Printf(" ⚠️ triggers: failed to load event triggers: %v", err) } else { for i := range eventTriggers { e.wireEventTrigger(&eventTriggers[i]) } if len(eventTriggers) > 0 { log.Printf(" 🔔 triggers: %d event trigger(s) loaded", len(eventTriggers)) } } } // Load scheduled tasks if e.stores.ScheduledTasks != nil { tasks, err := e.stores.ScheduledTasks.ListEnabled(ctx) if err != nil { log.Printf(" ⚠️ triggers: failed to load scheduled tasks: %v", err) } else { for i := range tasks { e.wireScheduledTask(&tasks[i]) } if len(tasks) > 0 { log.Printf(" ⏰ triggers: %d scheduled task(s) loaded", len(tasks)) } } } e.cron.Start() return nil } // Stop gracefully shuts down all trigger goroutines and cron jobs. func (e *Engine) Stop() { if e.cancel != nil { e.cancel() } // Stop cron scheduler cronCtx := e.cron.Stop() <-cronCtx.Done() // Unsubscribe all event triggers e.mu.Lock() for id, unsub := range e.unsubs { unsub() delete(e.unsubs, id) } e.mu.Unlock() log.Printf(" 🔔 triggers: engine stopped") } // RegisterTrigger adds a single trigger at runtime (package install). func (e *Engine) RegisterTrigger(t *models.Trigger) { if !t.Enabled { return } switch t.Type { case models.TriggerTypeEvent: e.wireEventTrigger(t) // Webhook triggers are resolved at request time — no wiring needed } } // UnregisterTrigger removes a trigger at runtime (package uninstall/disable). func (e *Engine) UnregisterTrigger(triggerID string) { e.mu.Lock() defer e.mu.Unlock() if unsub, ok := e.unsubs[triggerID]; ok { unsub() delete(e.unsubs, triggerID) } } // RegisterSchedule adds a scheduled task at runtime. func (e *Engine) RegisterSchedule(t *models.ScheduledTask) { if !t.Enabled { return } e.wireScheduledTask(t) } // UnregisterSchedule removes a scheduled task at runtime. func (e *Engine) UnregisterSchedule(taskID string) { e.mu.Lock() defer e.mu.Unlock() if entryID, ok := e.cronIDs[taskID]; ok { e.cron.Remove(entryID) delete(e.cronIDs, taskID) } } // ManualRun fires a scheduled task immediately (bypasses cron). func (e *Engine) ManualRun(t *models.ScheduledTask) { e.fireScheduledTask(t.ID, t.CreatorID, t.RunAs, t.Script, t.CronExpr) } // ReloadPackageTriggers re-syncs triggers for a package from DB. func (e *Engine) ReloadPackageTriggers(ctx context.Context, packageID string) error { if e.stores.Triggers == nil { return nil } triggers, err := e.stores.Triggers.ListByPackage(ctx, packageID) if err != nil { return err } // Unregister existing e.mu.RLock() var toRemove []string for id := range e.unsubs { toRemove = append(toRemove, id) } e.mu.RUnlock() // We need to check which ones belong to this package for _, id := range toRemove { // Simple: unregister all, re-register from DB e.UnregisterTrigger(id) } // Re-register enabled triggers for i := range triggers { e.RegisterTrigger(&triggers[i]) } return nil } // logExecution records a trigger or scheduled task execution. func (e *Engine) logExecution(triggerID, scheduledTaskID string, firedAt time.Time, durationMs int, success bool, errStr, output string) { if e.stores.Triggers == nil { return } tl := &models.TriggerLog{ TriggerID: triggerID, ScheduledTaskID: scheduledTaskID, FiredAt: firedAt.UTC().Format(time.RFC3339), DurationMs: &durationMs, Success: success, Error: errStr, Output: truncate(output, 4000), } if err := e.stores.Triggers.LogExecution(context.Background(), tl); err != nil { log.Printf(" ⚠️ triggers: failed to log execution: %v", err) } } // publishEvent emits a trigger lifecycle event on the bus. func (e *Engine) publishEvent(label string, triggerID string) { if e.bus != nil { payload, _ := json.Marshal(map[string]any{"trigger_id": triggerID}) e.bus.Publish(events.Event{ Label: label, Payload: payload, Ts: time.Now().UnixMilli(), }) } } func truncate(s string, max int) string { if len(s) <= max { return s } return s[:max] + "…" }