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/event.go
Jeffrey Smith 983d761bbe
All checks were successful
CI/CD / detect-changes (push) Successful in 4s
CI/CD / test-frontend (push) Has been skipped
CI/CD / test-runners (push) Has been skipped
CI/CD / e2e-smoke (push) Has been skipped
CI/CD / test-go-pg (push) Successful in 2m51s
CI/CD / test-sqlite (push) Successful in 3m1s
CI/CD / build-and-deploy (push) Successful in 1m17s
Feat v0.9.2 converter consolidation (#75)
Co-authored-by: Jeffrey Smith <jasafpro@gmail.com>
Co-committed-by: Jeffrey Smith <jasafpro@gmail.com>
2026-04-03 14:32:14 +00:00

126 lines
3.3 KiB
Go

package triggers
import (
"context"
"encoding/json"
"log"
"time"
"go.starlark.net/starlark"
"armature/events"
"armature/models"
"armature/sandbox"
)
// wireEventTrigger subscribes to the bus pattern for an event trigger.
func (e *Engine) wireEventTrigger(t *models.Trigger) {
if t.EventPattern == "" {
return
}
triggerID := t.ID
packageID := t.PackageID
entryPoint := t.EntryPoint
pattern := t.EventPattern
unsub := e.bus.Subscribe(pattern, func(ev events.Event) {
// Fire asynchronously — never block the event bus
go e.fireEventTrigger(triggerID, packageID, entryPoint, ev)
})
e.mu.Lock()
e.unsubs[triggerID] = unsub
e.mu.Unlock()
}
// fireEventTrigger invokes the Starlark handler for an event trigger.
func (e *Engine) fireEventTrigger(triggerID, packageID, entryPoint string, ev events.Event) {
start := time.Now()
ctx := e.ctx
if ctx == nil {
ctx = context.Background()
}
// Re-check trigger is still enabled
trigger, err := e.stores.Triggers.GetByID(ctx, triggerID)
if err != nil || trigger == nil || !trigger.Enabled {
return
}
// Load package
pkg, err := e.stores.Packages.Get(ctx, packageID)
if err != nil || pkg == nil || pkg.Status != models.PackageStatusActive {
return
}
// Check triggers.register permission
if !e.hasPermission(ctx, packageID, models.ExtPermTriggersRegister) {
return
}
// Build Starlark context dict
ctxDict := starlark.NewDict(6)
_ = ctxDict.SetKey(starlark.String("trigger_type"), starlark.String("event"))
_ = ctxDict.SetKey(starlark.String("trigger_id"), starlark.String(triggerID))
_ = ctxDict.SetKey(starlark.String("event_label"), starlark.String(ev.Label))
_ = ctxDict.SetKey(starlark.String("event_ts"), starlark.MakeInt64(ev.Ts))
if ev.Room != "" {
_ = ctxDict.SetKey(starlark.String("event_room"), starlark.String(ev.Room))
}
if len(ev.Payload) > 0 {
var payloadMap map[string]any
if json.Unmarshal(ev.Payload, &payloadMap) == nil {
_ = ctxDict.SetKey(starlark.String("event_payload"), sandbox.MapToDict(payloadMap))
} else {
_ = ctxDict.SetKey(starlark.String("event_payload"), starlark.String(string(ev.Payload)))
}
}
// Call entry point with full sandbox context
_, output, callErr := e.runner.CallEntryPoint(ctx, pkg, entryPoint,
starlark.Tuple{ctxDict}, nil, nil)
duration := int(time.Since(start).Milliseconds())
errStr := ""
if callErr != nil {
errStr = callErr.Error()
log.Printf(" ⚠️ trigger[event] %s/%s error: %v", packageID, entryPoint, callErr)
}
// Update fire state
_ = e.stores.Triggers.UpdateFireState(ctx, triggerID, start, errStr)
// Log execution
e.logExecution(triggerID, "", start, duration, callErr == nil, errStr, output)
e.publishEvent("trigger.fired", triggerID)
}
// hasPermission checks if a package has a specific granted permission.
func (e *Engine) hasPermission(ctx context.Context, packageID, permission string) bool {
if e.stores.ExtPermissions == nil {
return false
}
granted, err := e.stores.ExtPermissions.GrantedForPackage(ctx, packageID)
if err != nil {
return false
}
for _, p := range granted {
if p == permission {
return true
}
}
return false
}
// RunContext builds a sandbox.RunContext for a trigger invocation.
func triggerRunContext(userID, teamID string) *sandbox.RunContext {
if userID == "" {
return nil
}
return &sandbox.RunContext{
UserID: userID,
TeamID: teamID,
}
}