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 f0dd43144e rebrand: Switchboard Core → Armature
- Rename Go module switchboard-core → armature (155+ files)
- Rename Docker image → gobha/armature
- Rename K8s resources, secrets, deployments
- Rename Prometheus metrics switchboard_* → armature_*
- Rename env vars SWITCHBOARD_ADMIN_* → ARMATURE_ADMIN_*
- Rename DB names switchboard_core* → armature*
- Update all frontend branding, notification templates, docs
- Update CI scripts, e2e tests, Keycloak realm, nginx conf
- Rename scripts/switchboard-ca.sh → scripts/armature-ca.sh
- Rename k8s/switchboard.yaml → k8s/armature.yaml
- Rename chart alerting/dashboard files
- Fix: DockerHub push uses env: binding for secret injection
- Helm chart updated (name, labels, template functions, dashboard, alerting)
- Replace favicon/icon assets with Armature brand

No functional changes. Pure mechanical rename + CI fix.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-03-31 21:39:58 +00:00

229 lines
5.7 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"
"sync/atomic"
"time"
"github.com/robfig/cron/v3"
"armature/events"
"armature/models"
"armature/sandbox"
"armature/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
fireCount atomic.Int64 // cumulative trigger fires
ctx context.Context
cancel context.CancelFunc
}
// FireCount returns the cumulative number of trigger fires.
func (e *Engine) FireCount() int64 { return e.fireCount.Load() }
// 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 and increments the fire counter.
func (e *Engine) publishEvent(label string, triggerID string) {
e.fireCount.Add(1)
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] + "…"
}