This repository has been archived on 2026-04-03. You can view files and clone it. You cannot open issues or pull requests or push a commit.
Files
core/server/triggers/engine.go
Jeffrey Smith 3d4228f868
All checks were successful
CI/CD / detect-changes (push) Successful in 4s
CI/CD / test-frontend (push) Successful in 6s
CI/CD / test-sqlite (push) Successful in 2m46s
CI/CD / test-go-pg (push) Successful in 2m47s
CI/CD / build-and-deploy (push) Successful in 26s
Feat v0.6.3 dead code sweep (#38)
Co-authored-by: Jeffrey Smith <jasafpro@gmail.com>
Co-committed-by: Jeffrey Smith <jasafpro@gmail.com>
2026-03-31 12:37:47 +00:00

222 lines
5.5 KiB
Go

// Package triggers — engine.go
//
// (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] + "…"
}