- 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>
229 lines
5.7 KiB
Go
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] + "…"
|
|
}
|