Feat triggers v0.2.2 #6

Merged
xcaliber merged 1 commits from feat/triggers-v0.2.2 into main 2026-03-26 22:47:31 +00:00
29 changed files with 3237 additions and 3 deletions

View File

@@ -2,7 +2,43 @@
All notable changes to Switchboard Core are documented here.
## [Unreleased] — v0.2.1
## [Unreleased] — v0.2.2
### Added
- **Event bus subscriptions**: Extensions declare event triggers in manifest
(`"triggers": [{"type": "event", "pattern": "workflow.completed", ...}]`).
Wired via `bus.Subscribe()` on startup. Handlers fire asynchronously.
- **Webhook triggers**: Inbound HTTP at `/api/v1/hooks/:package_id/:slug`.
HMAC-SHA256 verification via `X-Switchboard-Signature` header. Synchronous
Starlark handler can return custom HTTP status and body.
- **Scheduled tasks**: User-created cron-scheduled Starlark scripts with
restricted sandbox (no raw HTTP, no DB table creation, connections-only
outbound). Runs as creator identity with RBAC scoping. Admin-created tasks
can opt into system context. Creator deactivation auto-pauses schedule.
- **Schedule templates**: Extensions ship pre-built schedule templates in
manifest (`schedule_templates` array) with configurable params and default
cron expressions.
- `triggers.register` extension permission — required for event/webhook triggers
- `triggers` table — extension-declared event and webhook trigger definitions
- `scheduled_tasks` table — user-created cron tasks with script, template,
and identity fields
- `trigger_logs` table — unified execution audit log for both tiers
- `TriggerStore` + `ScheduledTaskStore` interfaces (postgres + sqlite)
- Trigger engine (`server/triggers/`) — orchestrates event subscriptions,
webhook resolution, and cron scheduling via `robfig/cron/v3`
- `SyncManifestTriggers()` — declarative sync of event/webhook triggers from
manifest. Hooked into seed, admin install, and package install flows.
- Admin trigger API: `GET/PUT/DELETE /admin/triggers`, `/admin/triggers/:id/logs`,
`/admin/packages/:id/triggers`
- Admin schedule API: `GET /admin/schedules`, enable/disable/delete
- User schedule API: full CRUD at `/api/v1/schedules`, manual run, execution logs
- `trigger.fired` and `trigger.error` event bus labels (DirLocal) for observability
- OpenAPI spec: Trigger, ScheduledTask, TriggerLog schemas + all new endpoints
---
## [v0.2.1] — 2026-03-26
### Added

View File

@@ -55,8 +55,10 @@ SDK stabilization, and the first rebuilt extension (tasks).
| Step | Status | Description |
|------|--------|-------------|
| Event bus subscriptions | 🔲 | Extensions register match expressions at install time |
| Trigger system | 🔲 | Time (cron), webhook (inbound HTTP), event (bus subscription) |
| Event bus subscriptions | | Extensions register event patterns in manifest. Wired via `bus.Subscribe()` on startup. Async handler invocation. |
| Webhook triggers | | Inbound HTTP at `/api/v1/hooks/:package_id/:slug`. HMAC-SHA256 verification. Synchronous Starlark handler response. |
| Scheduled tasks | ✅ | User-created cron tasks with restricted sandbox (no raw HTTP, no DB table creation). Runs as creator identity. Templates from extensions. Dedicated schedules API. |
| Trigger admin API | ✅ | CRUD for triggers + schedules. Enable/disable, execution logs, per-package listing. |
### v0.2.3 — SDK + Task Extension
@@ -112,3 +114,5 @@ Extension and operations tracks converge. First externally usable release.
| Notes over Editor | First surface is Obsidian-style notes (rich text, folders, backlinks) instead of a code editor. Notes is a stronger E2E proof — it exercises ext_data, storage, and the SDK more fully than a pure-browser CM6 editor. |
| No built-in auto-install | Extensions ship in the repo but are not auto-installed. Distribution model TBD — explicit install only. Keeps the kernel clean and avoids opinionated defaults. |
| Chat → post-MVP | Chat extension (providers, streaming, personas) is valuable but not MVP-critical. The platform must prove itself with simpler surfaces first. Chat moves to post-MVP track. |
| Two trigger tiers | Event + webhook triggers are extension-declared (manifest contract, full sandbox). Scheduled tasks are user-created ad-hoc (restricted sandbox — no raw HTTP, no DB table creation, connections-only outbound). Separation keeps extension contracts static and user automation safe. |
| Scheduled task identity | Tasks run as their creator (RBAC-scoped). Admin-created tasks can opt into system context. Creator deactivation pauses the schedule. Ensures audit trail and permission boundaries. |

View File

@@ -0,0 +1,85 @@
-- ==========================================
-- Switchboard Core — 011 Triggers & Scheduled Tasks
-- ==========================================
-- ── Extension Triggers (event + webhook) ──────
CREATE TABLE IF NOT EXISTS triggers (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
package_id TEXT NOT NULL REFERENCES packages(id) ON DELETE CASCADE,
type TEXT NOT NULL CHECK (type IN ('webhook', 'event')),
enabled BOOLEAN NOT NULL DEFAULT true,
-- Webhook triggers
slug TEXT,
secret TEXT,
-- Event triggers
event_pattern TEXT,
-- Common
entry_point TEXT NOT NULL,
config JSONB NOT NULL DEFAULT '{}',
fire_count INTEGER NOT NULL DEFAULT 0,
last_fire_at TIMESTAMPTZ,
last_error TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
UNIQUE(package_id, type, slug),
UNIQUE(package_id, type, event_pattern)
);
CREATE INDEX IF NOT EXISTS idx_triggers_package ON triggers(package_id);
CREATE INDEX IF NOT EXISTS idx_triggers_type ON triggers(type);
CREATE INDEX IF NOT EXISTS idx_triggers_enabled ON triggers(enabled) WHERE enabled = true;
CREATE INDEX IF NOT EXISTS idx_triggers_webhook ON triggers(package_id, slug) WHERE type = 'webhook';
-- ── Scheduled Tasks (user-created cron) ───────
CREATE TABLE IF NOT EXISTS scheduled_tasks (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
name TEXT NOT NULL,
description TEXT NOT NULL DEFAULT '',
creator_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
run_as TEXT NOT NULL DEFAULT 'creator'
CHECK (run_as IN ('creator', 'system')),
cron_expr TEXT NOT NULL,
next_fire_at TIMESTAMPTZ,
last_fire_at TIMESTAMPTZ,
enabled BOOLEAN NOT NULL DEFAULT true,
script TEXT NOT NULL DEFAULT '',
template_id TEXT,
template_params JSONB NOT NULL DEFAULT '{}',
fire_count INTEGER NOT NULL DEFAULT 0,
last_error TEXT,
last_duration_ms INTEGER,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE INDEX IF NOT EXISTS idx_sched_tasks_creator ON scheduled_tasks(creator_id);
CREATE INDEX IF NOT EXISTS idx_sched_tasks_enabled ON scheduled_tasks(enabled) WHERE enabled = true;
CREATE INDEX IF NOT EXISTS idx_sched_tasks_next_fire ON scheduled_tasks(next_fire_at) WHERE enabled = true;
-- ── Trigger Logs (unified execution log) ──────
CREATE TABLE IF NOT EXISTS trigger_logs (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
trigger_id UUID REFERENCES triggers(id) ON DELETE CASCADE,
scheduled_task_id UUID REFERENCES scheduled_tasks(id) ON DELETE CASCADE,
fired_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
duration_ms INTEGER,
success BOOLEAN NOT NULL DEFAULT true,
error TEXT,
output TEXT,
CHECK (
(trigger_id IS NOT NULL AND scheduled_task_id IS NULL) OR
(trigger_id IS NULL AND scheduled_task_id IS NOT NULL)
)
);
CREATE INDEX IF NOT EXISTS idx_trigger_logs_trigger ON trigger_logs(trigger_id);
CREATE INDEX IF NOT EXISTS idx_trigger_logs_sched ON trigger_logs(scheduled_task_id);
CREATE INDEX IF NOT EXISTS idx_trigger_logs_fired ON trigger_logs(fired_at);

View File

@@ -0,0 +1,80 @@
-- ==========================================
-- Switchboard Core — 011 Triggers & Scheduled Tasks (SQLite)
-- ==========================================
-- ── Extension Triggers (event + webhook) ──────
CREATE TABLE IF NOT EXISTS triggers (
id TEXT PRIMARY KEY,
package_id TEXT NOT NULL REFERENCES packages(id) ON DELETE CASCADE,
type TEXT NOT NULL CHECK (type IN ('webhook', 'event')),
enabled INTEGER NOT NULL DEFAULT 1,
-- Webhook triggers
slug TEXT,
secret TEXT,
-- Event triggers
event_pattern TEXT,
-- Common
entry_point TEXT NOT NULL,
config TEXT NOT NULL DEFAULT '{}',
fire_count INTEGER NOT NULL DEFAULT 0,
last_fire_at TEXT,
last_error TEXT,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
UNIQUE(package_id, type, slug),
UNIQUE(package_id, type, event_pattern)
);
CREATE INDEX IF NOT EXISTS idx_triggers_package ON triggers(package_id);
CREATE INDEX IF NOT EXISTS idx_triggers_type ON triggers(type);
CREATE INDEX IF NOT EXISTS idx_triggers_enabled ON triggers(enabled);
CREATE INDEX IF NOT EXISTS idx_triggers_webhook ON triggers(package_id, slug);
-- ── Scheduled Tasks (user-created cron) ───────
CREATE TABLE IF NOT EXISTS scheduled_tasks (
id TEXT PRIMARY KEY,
name TEXT NOT NULL,
description TEXT NOT NULL DEFAULT '',
creator_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE,
run_as TEXT NOT NULL DEFAULT 'creator'
CHECK (run_as IN ('creator', 'system')),
cron_expr TEXT NOT NULL,
next_fire_at TEXT,
last_fire_at TEXT,
enabled INTEGER NOT NULL DEFAULT 1,
script TEXT NOT NULL DEFAULT '',
template_id TEXT,
template_params TEXT NOT NULL DEFAULT '{}',
fire_count INTEGER NOT NULL DEFAULT 0,
last_error TEXT,
last_duration_ms INTEGER,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now'))
);
CREATE INDEX IF NOT EXISTS idx_sched_tasks_creator ON scheduled_tasks(creator_id);
CREATE INDEX IF NOT EXISTS idx_sched_tasks_enabled ON scheduled_tasks(enabled);
CREATE INDEX IF NOT EXISTS idx_sched_tasks_next_fire ON scheduled_tasks(next_fire_at);
-- ── Trigger Logs (unified execution log) ──────
CREATE TABLE IF NOT EXISTS trigger_logs (
id TEXT PRIMARY KEY,
trigger_id TEXT REFERENCES triggers(id) ON DELETE CASCADE,
scheduled_task_id TEXT REFERENCES scheduled_tasks(id) ON DELETE CASCADE,
fired_at TEXT NOT NULL DEFAULT (datetime('now')),
duration_ms INTEGER,
success INTEGER NOT NULL DEFAULT 1,
error TEXT,
output TEXT
);
CREATE INDEX IF NOT EXISTS idx_trigger_logs_trigger ON trigger_logs(trigger_id);
CREATE INDEX IF NOT EXISTS idx_trigger_logs_sched ON trigger_logs(scheduled_task_id);
CREATE INDEX IF NOT EXISTS idx_trigger_logs_fired ON trigger_logs(fired_at);

View File

@@ -80,6 +80,10 @@ var routeTable = map[string]Direction{
"extension.loaded": DirLocal, // Client-only
"extension.error": DirLocal,
// Trigger system (v0.2.2)
"trigger.fired": DirLocal, // Trigger invocation event (observability)
"trigger.error": DirLocal, // Trigger execution error
// Heartbeat
"ping": DirFromClient,
"pong": DirToClient,

View File

@@ -12,6 +12,7 @@ import (
"switchboard-core/database"
"switchboard-core/models"
"switchboard-core/store"
"switchboard-core/triggers"
)
// ExtensionHandler serves extension management endpoints.
@@ -239,6 +240,9 @@ func (h *ExtensionHandler) AdminInstallExtension(c *gin.Context) {
}
}
// v0.2.2: Sync triggers from manifest
SyncManifestTriggers(c.Request.Context(), h.stores, triggers.GlobalEngine(), pkg.ID, manifestMap)
c.JSON(201, gin.H{"data": pkg})
}

View File

@@ -19,6 +19,7 @@ import (
"switchboard-core/models"
"switchboard-core/sandbox"
"switchboard-core/store"
"switchboard-core/triggers"
)
// validPackageID matches lowercase alphanumeric slugs with optional hyphens.
@@ -474,6 +475,9 @@ func (h *PackageHandler) InstallPackage(c *gin.Context) {
// if the package declares permissions.
SyncManifestPermissions(c, h.stores, pkgID, manifest)
// v0.2.2: Sync triggers from manifest
SyncManifestTriggers(c.Request.Context(), h.stores, triggers.GlobalEngine(), pkgID, manifest)
// v0.30.0: Run schema migrations if declared.
newSchemaVersion := ParseSchemaVersion(manifest)
if newSchemaVersion > 0 && h.sandbox != nil {

View File

@@ -0,0 +1,309 @@
package handlers
import (
"net/http"
"time"
"github.com/gin-gonic/gin"
"github.com/robfig/cron/v3"
"switchboard-core/models"
"switchboard-core/store"
"switchboard-core/triggers"
)
// ScheduleHandler provides CRUD for user-created scheduled tasks.
type ScheduleHandler struct {
stores store.Stores
engine *triggers.Engine
}
// NewScheduleHandler creates a ScheduleHandler.
func NewScheduleHandler(stores store.Stores, engine *triggers.Engine) *ScheduleHandler {
return &ScheduleHandler{stores: stores, engine: engine}
}
// ── User Endpoints ───────────────────────────
// ListMySchedules returns schedules created by the current user.
// GET /api/v1/schedules
func (h *ScheduleHandler) ListMySchedules(c *gin.Context) {
userID := c.GetString("user_id")
tasks, err := h.stores.ScheduledTasks.ListByCreator(c.Request.Context(), userID)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"schedules": tasks})
}
// CreateSchedule creates a new scheduled task.
// POST /api/v1/schedules
func (h *ScheduleHandler) CreateSchedule(c *gin.Context) {
var req struct {
Name string `json:"name" binding:"required"`
Description string `json:"description"`
CronExpr string `json:"cron_expr" binding:"required"`
Script string `json:"script"`
TemplateID string `json:"template_id"`
TemplateParams any `json:"template_params"`
}
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
// Validate cron expression
parser := cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow)
schedule, err := parser.Parse(req.CronExpr)
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": "invalid cron expression: " + err.Error()})
return
}
nextFire := schedule.Next(time.Now())
userID := c.GetString("user_id")
task := &models.ScheduledTask{
Name: req.Name,
Description: req.Description,
CreatorID: userID,
RunAs: models.RunAsCreator,
CronExpr: req.CronExpr,
NextFireAt: &nextFire,
Enabled: true,
Script: req.Script,
TemplateID: req.TemplateID,
}
if err := h.stores.ScheduledTasks.Create(c.Request.Context(), task); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
// Register with engine
h.engine.RegisterSchedule(task)
c.JSON(http.StatusCreated, task)
}
// GetSchedule returns a single scheduled task.
// GET /api/v1/schedules/:id
func (h *ScheduleHandler) GetSchedule(c *gin.Context) {
task, err := h.stores.ScheduledTasks.GetByID(c.Request.Context(), c.Param("id"))
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
if task == nil {
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
return
}
// Check ownership (non-admin can only see own schedules)
userID := c.GetString("user_id")
if task.CreatorID != userID {
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
return
}
c.JSON(http.StatusOK, task)
}
// UpdateSchedule updates a scheduled task.
// PUT /api/v1/schedules/:id
func (h *ScheduleHandler) UpdateSchedule(c *gin.Context) {
task, err := h.stores.ScheduledTasks.GetByID(c.Request.Context(), c.Param("id"))
if err != nil || task == nil {
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
return
}
userID := c.GetString("user_id")
if task.CreatorID != userID {
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
return
}
var req struct {
Name *string `json:"name"`
Description *string `json:"description"`
CronExpr *string `json:"cron_expr"`
Script *string `json:"script"`
Enabled *bool `json:"enabled"`
}
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
if req.Name != nil {
task.Name = *req.Name
}
if req.Description != nil {
task.Description = *req.Description
}
if req.Script != nil {
task.Script = *req.Script
}
if req.Enabled != nil {
task.Enabled = *req.Enabled
}
if req.CronExpr != nil {
parser := cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow)
schedule, err := parser.Parse(*req.CronExpr)
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": "invalid cron expression: " + err.Error()})
return
}
task.CronExpr = *req.CronExpr
nextFire := schedule.Next(time.Now())
task.NextFireAt = &nextFire
}
if err := h.stores.ScheduledTasks.Update(c.Request.Context(), task); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
// Re-register with engine
h.engine.UnregisterSchedule(task.ID)
if task.Enabled {
h.engine.RegisterSchedule(task)
}
c.JSON(http.StatusOK, task)
}
// DeleteSchedule deletes a scheduled task.
// DELETE /api/v1/schedules/:id
func (h *ScheduleHandler) DeleteSchedule(c *gin.Context) {
task, err := h.stores.ScheduledTasks.GetByID(c.Request.Context(), c.Param("id"))
if err != nil || task == nil {
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
return
}
userID := c.GetString("user_id")
if task.CreatorID != userID {
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
return
}
h.engine.UnregisterSchedule(task.ID)
if err := h.stores.ScheduledTasks.Delete(c.Request.Context(), task.ID); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"ok": true})
}
// RunSchedule manually triggers a scheduled task (test run).
// POST /api/v1/schedules/:id/run
func (h *ScheduleHandler) RunSchedule(c *gin.Context) {
task, err := h.stores.ScheduledTasks.GetByID(c.Request.Context(), c.Param("id"))
if err != nil || task == nil {
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
return
}
userID := c.GetString("user_id")
if task.CreatorID != userID {
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
return
}
// Fire immediately in a goroutine
go h.engine.ManualRun(task)
c.JSON(http.StatusAccepted, gin.H{"ok": true, "message": "task queued for execution"})
}
// ListScheduleLogs returns execution logs for a scheduled task.
// GET /api/v1/schedules/:id/logs
func (h *ScheduleHandler) ListScheduleLogs(c *gin.Context) {
task, err := h.stores.ScheduledTasks.GetByID(c.Request.Context(), c.Param("id"))
if err != nil || task == nil {
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
return
}
userID := c.GetString("user_id")
if task.CreatorID != userID {
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
return
}
limit := parseIntParam(c, "limit", 50)
logs, err := h.stores.Triggers.ListLogs(c.Request.Context(), task.ID, limit)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"logs": logs})
}
// ── Admin Endpoints ──────────────────────────
// AdminListSchedules returns all scheduled tasks (admin view).
// GET /admin/schedules
func (h *ScheduleHandler) AdminListSchedules(c *gin.Context) {
opts := store.ScheduledTaskListOptions{
ListOptions: store.ListOptions{
Limit: parseIntParam(c, "limit", 50),
Offset: parseIntParam(c, "offset", 0),
},
CreatorID: c.Query("creator_id"),
}
if e := c.Query("enabled"); e != "" {
b := e == "true"
opts.Enabled = &b
}
tasks, total, err := h.stores.ScheduledTasks.List(c.Request.Context(), opts)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"schedules": tasks, "total": total})
}
// AdminEnableSchedule enables a scheduled task.
// PUT /admin/schedules/:id/enable
func (h *ScheduleHandler) AdminEnableSchedule(c *gin.Context) {
id := c.Param("id")
if err := h.stores.ScheduledTasks.SetEnabled(c.Request.Context(), id, true); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
task, _ := h.stores.ScheduledTasks.GetByID(c.Request.Context(), id)
if task != nil {
h.engine.RegisterSchedule(task)
}
c.JSON(http.StatusOK, gin.H{"ok": true})
}
// AdminDisableSchedule disables a scheduled task.
// PUT /admin/schedules/:id/disable
func (h *ScheduleHandler) AdminDisableSchedule(c *gin.Context) {
id := c.Param("id")
h.engine.UnregisterSchedule(id)
if err := h.stores.ScheduledTasks.SetEnabled(c.Request.Context(), id, false); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"ok": true})
}
// AdminDeleteSchedule deletes a scheduled task.
// DELETE /admin/schedules/:id
func (h *ScheduleHandler) AdminDeleteSchedule(c *gin.Context) {
id := c.Param("id")
h.engine.UnregisterSchedule(id)
if err := h.stores.ScheduledTasks.Delete(c.Request.Context(), id); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"ok": true})
}

View File

@@ -138,6 +138,8 @@ func SeedBuiltinPackages(stores store.Stores, extensionsDir string) {
log.Printf("⚠ Failed to update builtin package %s: %v", manifest.ID, err)
continue
}
// Sync triggers on version update (v0.2.2)
SyncManifestTriggers(ctx, stores, nil, manifest.ID, manifestMap)
updated++
continue
}
@@ -167,6 +169,9 @@ func SeedBuiltinPackages(stores store.Stores, extensionsDir string) {
if err := stores.Packages.Update(ctx, manifest.ID, pkg); err != nil {
log.Printf("⚠ Failed to update builtin package metadata %s: %v", manifest.ID, err)
}
// Sync triggers from manifest (v0.2.2)
SyncManifestTriggers(ctx, stores, nil, manifest.ID, manifestMap)
seeded++
}

View File

@@ -0,0 +1,192 @@
package handlers
import (
"context"
"crypto/rand"
"encoding/hex"
"encoding/json"
"log"
"switchboard-core/models"
"switchboard-core/store"
"switchboard-core/triggers"
)
// manifestTrigger is the manifest.json trigger declaration shape.
type manifestTrigger struct {
Type string `json:"type"` // "event" or "webhook"
Pattern string `json:"pattern"` // event pattern (for type=event)
Slug string `json:"slug"` // webhook slug (for type=webhook)
EntryPoint string `json:"entry_point"` // Starlark function name
}
// SyncManifestTriggers parses the "triggers" array from a package manifest
// and syncs them with the trigger store. Removes triggers no longer in the
// manifest (declarative sync). For webhooks, generates an HMAC secret on
// first create.
//
// Called from SeedBuiltinPackages, admin install, and package install.
func SyncManifestTriggers(ctx context.Context, stores store.Stores, engine *triggers.Engine, pkgID string, manifest map[string]any) {
if stores.Triggers == nil {
return
}
// Parse triggers from manifest
triggersRaw, ok := manifest["triggers"]
if !ok {
return
}
data, err := json.Marshal(triggersRaw)
if err != nil {
return
}
var declared []manifestTrigger
if err := json.Unmarshal(data, &declared); err != nil {
log.Printf(" ⚠️ trigger-sync[%s]: invalid triggers in manifest: %v", pkgID, err)
return
}
if len(declared) == 0 {
return
}
// Get existing triggers for this package
existing, err := stores.Triggers.ListByPackage(ctx, pkgID)
if err != nil {
log.Printf(" ⚠️ trigger-sync[%s]: failed to list existing: %v", pkgID, err)
return
}
// Build lookup maps
existingByKey := make(map[string]*models.Trigger, len(existing))
for i := range existing {
key := triggerKey(&existing[i])
existingByKey[key] = &existing[i]
}
declaredKeys := make(map[string]bool, len(declared))
created, updated := 0, 0
for _, dt := range declared {
if dt.Type != models.TriggerTypeEvent && dt.Type != models.TriggerTypeWebhook {
log.Printf(" ⚠️ trigger-sync[%s]: unknown trigger type %q, skipping", pkgID, dt.Type)
continue
}
if dt.EntryPoint == "" {
continue
}
key := dt.Type + ":" + triggerIdentifier(dt)
declaredKeys[key] = true
if ex, ok := existingByKey[key]; ok {
// Update if entry_point changed
if ex.EntryPoint != dt.EntryPoint {
ex.EntryPoint = dt.EntryPoint
if err := stores.Triggers.Update(ctx, ex); err != nil {
log.Printf(" ⚠️ trigger-sync[%s]: update failed: %v", pkgID, err)
} else {
updated++
}
}
continue
}
// Create new trigger
t := &models.Trigger{
PackageID: pkgID,
Type: dt.Type,
Enabled: true,
EntryPoint: dt.EntryPoint,
Config: json.RawMessage("{}"),
}
switch dt.Type {
case models.TriggerTypeEvent:
t.EventPattern = dt.Pattern
case models.TriggerTypeWebhook:
t.Slug = dt.Slug
t.Secret = generateSecret()
}
if err := stores.Triggers.Create(ctx, t); err != nil {
log.Printf(" ⚠️ trigger-sync[%s]: create failed: %v", pkgID, err)
continue
}
created++
// Register with engine
if engine != nil {
// Re-fetch to get the generated ID
switch dt.Type {
case models.TriggerTypeWebhook:
if wh, _ := stores.Triggers.GetWebhook(ctx, pkgID, dt.Slug); wh != nil {
engine.RegisterTrigger(wh)
}
case models.TriggerTypeEvent:
// Need to fetch by package to find the new one
triggers, _ := stores.Triggers.ListByPackage(ctx, pkgID)
for i := range triggers {
if triggers[i].EventPattern == dt.Pattern {
engine.RegisterTrigger(&triggers[i])
break
}
}
}
}
}
// Remove triggers no longer in manifest
removed := 0
for key, ex := range existingByKey {
if !declaredKeys[key] {
if engine != nil {
engine.UnregisterTrigger(ex.ID)
}
if err := stores.Triggers.Delete(ctx, ex.ID); err != nil {
log.Printf(" ⚠️ trigger-sync[%s]: delete failed: %v", pkgID, err)
} else {
removed++
}
}
}
if created+updated+removed > 0 {
log.Printf(" 🔔 trigger-sync[%s]: %d created, %d updated, %d removed",
pkgID, created, updated, removed)
}
}
// triggerKey returns a unique key for a trigger within a package.
func triggerKey(t *models.Trigger) string {
return t.Type + ":" + triggerIdentifierFromModel(t)
}
func triggerIdentifier(dt manifestTrigger) string {
switch dt.Type {
case models.TriggerTypeEvent:
return dt.Pattern
case models.TriggerTypeWebhook:
return dt.Slug
}
return ""
}
func triggerIdentifierFromModel(t *models.Trigger) string {
switch t.Type {
case models.TriggerTypeEvent:
return t.EventPattern
case models.TriggerTypeWebhook:
return t.Slug
}
return ""
}
// generateSecret creates a random 32-byte hex-encoded secret for webhook HMAC.
func generateSecret() string {
b := make([]byte, 32)
_, _ = rand.Read(b)
return hex.EncodeToString(b)
}

135
server/handlers/triggers.go Normal file
View File

@@ -0,0 +1,135 @@
package handlers
import (
"net/http"
"strconv"
"github.com/gin-gonic/gin"
"switchboard-core/store"
"switchboard-core/triggers"
)
// TriggerHandler provides admin CRUD for extension-declared triggers.
type TriggerHandler struct {
stores store.Stores
engine *triggers.Engine
}
// NewTriggerHandler creates a TriggerHandler.
func NewTriggerHandler(stores store.Stores, engine *triggers.Engine) *TriggerHandler {
return &TriggerHandler{stores: stores, engine: engine}
}
// ListTriggers returns all triggers with optional filters.
// GET /admin/triggers?package_id=&type=&enabled=
func (h *TriggerHandler) ListTriggers(c *gin.Context) {
opts := store.TriggerListOptions{
ListOptions: store.ListOptions{
Limit: parseIntParam(c, "limit", 50),
Offset: parseIntParam(c, "offset", 0),
},
PackageID: c.Query("package_id"),
Type: c.Query("type"),
}
if e := c.Query("enabled"); e != "" {
b := e == "true"
opts.Enabled = &b
}
triggers, total, err := h.stores.Triggers.List(c.Request.Context(), opts)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"triggers": triggers, "total": total})
}
// GetTrigger returns a single trigger by ID.
// GET /admin/triggers/:id
func (h *TriggerHandler) GetTrigger(c *gin.Context) {
t, err := h.stores.Triggers.GetByID(c.Request.Context(), c.Param("id"))
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
if t == nil {
c.JSON(http.StatusNotFound, gin.H{"error": "trigger not found"})
return
}
c.JSON(http.StatusOK, t)
}
// EnableTrigger enables a trigger.
// PUT /admin/triggers/:id/enable
func (h *TriggerHandler) EnableTrigger(c *gin.Context) {
id := c.Param("id")
if err := h.stores.Triggers.SetEnabled(c.Request.Context(), id, true); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
// Re-register with engine
t, _ := h.stores.Triggers.GetByID(c.Request.Context(), id)
if t != nil {
h.engine.RegisterTrigger(t)
}
c.JSON(http.StatusOK, gin.H{"ok": true})
}
// DisableTrigger disables a trigger.
// PUT /admin/triggers/:id/disable
func (h *TriggerHandler) DisableTrigger(c *gin.Context) {
id := c.Param("id")
if err := h.stores.Triggers.SetEnabled(c.Request.Context(), id, false); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
h.engine.UnregisterTrigger(id)
c.JSON(http.StatusOK, gin.H{"ok": true})
}
// DeleteTrigger deletes a trigger.
// DELETE /admin/triggers/:id
func (h *TriggerHandler) DeleteTrigger(c *gin.Context) {
id := c.Param("id")
h.engine.UnregisterTrigger(id)
if err := h.stores.Triggers.Delete(c.Request.Context(), id); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"ok": true})
}
// ListTriggerLogs returns execution logs for a trigger.
// GET /admin/triggers/:id/logs?limit=
func (h *TriggerHandler) ListTriggerLogs(c *gin.Context) {
limit := parseIntParam(c, "limit", 50)
logs, err := h.stores.Triggers.ListLogs(c.Request.Context(), c.Param("id"), limit)
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"logs": logs})
}
// ListPackageTriggers returns triggers for a specific package.
// GET /admin/packages/:id/triggers
func (h *TriggerHandler) ListPackageTriggers(c *gin.Context) {
triggers, err := h.stores.Triggers.ListByPackage(c.Request.Context(), c.Param("id"))
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"triggers": triggers})
}
func parseIntParam(c *gin.Context, key string, defaultVal int) int {
if v := c.Query(key); v != "" {
if n, err := strconv.Atoi(v); err == nil && n > 0 {
return n
}
}
return defaultVal
}

View File

@@ -29,6 +29,7 @@ import (
"switchboard-core/pages"
"switchboard-core/storage"
"switchboard-core/store"
"switchboard-core/triggers"
postgres "switchboard-core/store/postgres"
sqliteStore "switchboard-core/store/sqlite"
)
@@ -174,6 +175,13 @@ func main() {
log.Printf(" ⚠️ Extension SSRF protection relaxed: private IPs allowed")
}
// ── Trigger Engine (v0.2.2) ─────────────────
triggerEngine := triggers.New(stores, starlarkRunner, bus)
if err := triggerEngine.Start(context.Background()); err != nil {
log.Printf(" ⚠️ Trigger engine start: %v", err)
}
triggers.SetGlobalEngine(triggerEngine)
defer triggerEngine.Stop()
r := gin.New()
r.Use(middleware.RequestID())
@@ -347,6 +355,11 @@ func main() {
extH := handlers.NewExtensionHandler(stores)
api.GET("/extensions/:id/assets/*path", extH.ServeExtensionAsset)
// ── Webhook Inbound (v0.2.2) ──────────────
// Public routes — HMAC-verified in handler, no auth middleware.
api.POST("/hooks/:package_id/:slug", triggerEngine.HandleWebhook)
api.GET("/hooks/:package_id/:slug", triggerEngine.HandleWebhook)
// ── Protected routes ────────────────────
protected := api.Group("")
protected.Use(middleware.Auth(cfg, stores.Users, userCache))
@@ -445,6 +458,16 @@ func main() {
protected.PUT("/notifications/preferences/:type", notifH.SetPreference)
protected.DELETE("/notifications/preferences/:type", notifH.DeletePreference)
// ── Scheduled Tasks — User (v0.2.2) ────
userSchedH := handlers.NewScheduleHandler(stores, triggerEngine)
protected.GET("/schedules", userSchedH.ListMySchedules)
protected.POST("/schedules", userSchedH.CreateSchedule)
protected.GET("/schedules/:id", userSchedH.GetSchedule)
protected.PUT("/schedules/:id", userSchedH.UpdateSchedule)
protected.DELETE("/schedules/:id", userSchedH.DeleteSchedule)
protected.POST("/schedules/:id/run", userSchedH.RunSchedule)
protected.GET("/schedules/:id/logs", userSchedH.ListScheduleLogs)
// Teams (user: my teams)
teams := handlers.NewTeamHandler(stores, keyResolver)
protected.GET("/teams/mine", teams.MyTeams)
@@ -649,6 +672,23 @@ func main() {
wfPkgH := handlers.NewWorkflowPackageHandler(stores)
admin.GET("/workflows/:id/export", wfPkgH.ExportWorkflowPackage)
// ── Triggers (v0.2.2) ─────────────────
trigH := handlers.NewTriggerHandler(stores, triggerEngine)
admin.GET("/triggers", trigH.ListTriggers)
admin.GET("/triggers/:id", trigH.GetTrigger)
admin.PUT("/triggers/:id/enable", trigH.EnableTrigger)
admin.PUT("/triggers/:id/disable", trigH.DisableTrigger)
admin.DELETE("/triggers/:id", trigH.DeleteTrigger)
admin.GET("/triggers/:id/logs", trigH.ListTriggerLogs)
admin.GET("/packages/:id/triggers", trigH.ListPackageTriggers)
// ── Scheduled Tasks — Admin (v0.2.2) ──
schedH := handlers.NewScheduleHandler(stores, triggerEngine)
admin.GET("/schedules", schedH.AdminListSchedules)
admin.PUT("/schedules/:id/enable", schedH.AdminEnableSchedule)
admin.PUT("/schedules/:id/disable", schedH.AdminDisableSchedule)
admin.DELETE("/schedules/:id", schedH.AdminDeleteSchedule)
// Surface aliases (backward compat — same handlers)
admin.GET("/surfaces", pkgAdm.ListPackages)
admin.GET("/surfaces/:id", pkgAdm.GetPackage)

View File

@@ -34,6 +34,7 @@ const (
ExtPermFormValidate = "forms.validate" // v0.29.3: form validation hooks
ExtPermWorkflowAccess = "workflow.access" // v0.30.2: workflow definition + stage data access
ExtPermConnectionsRead = "connections.read" // v0.38.1: extension connection resolution
ExtPermTriggersRegister = "triggers.register" // v0.2.2: event/webhook trigger registration
)
// ValidExtensionPermissions is the set of recognized permission keys.
@@ -46,6 +47,7 @@ var ValidExtensionPermissions = map[string]bool{
ExtPermFormValidate: true,
ExtPermWorkflowAccess: true,
ExtPermConnectionsRead: true,
ExtPermTriggersRegister: true,
}
// ── Extension Permission Model ───────────────

77
server/models/trigger.go Normal file
View File

@@ -0,0 +1,77 @@
package models
import (
"encoding/json"
"time"
)
// ── Trigger Type Constants ───────────────────
const (
TriggerTypeWebhook = "webhook"
TriggerTypeEvent = "event"
)
// ── Scheduled Task Run-As Constants ──────────
const (
RunAsCreator = "creator"
RunAsSystem = "system"
)
// ── Trigger (extension-declared) ─────────────
// Trigger represents an extension-declared event or webhook trigger.
type Trigger struct {
ID string `json:"id" db:"id"`
PackageID string `json:"package_id" db:"package_id"`
Type string `json:"type" db:"type"`
Enabled bool `json:"enabled" db:"enabled"`
Slug string `json:"slug,omitempty" db:"slug"`
Secret string `json:"-" db:"secret"` // never serialized
EventPattern string `json:"event_pattern,omitempty" db:"event_pattern"`
EntryPoint string `json:"entry_point" db:"entry_point"`
Config json.RawMessage `json:"config" db:"config"`
FireCount int `json:"fire_count" db:"fire_count"`
LastFireAt *time.Time `json:"last_fire_at,omitempty" db:"last_fire_at"`
LastError string `json:"last_error,omitempty" db:"last_error"`
CreatedAt time.Time `json:"created_at" db:"created_at"`
UpdatedAt time.Time `json:"updated_at" db:"updated_at"`
}
// ── Scheduled Task (user-created) ────────────
// ScheduledTask represents a user-created cron-scheduled Starlark script.
type ScheduledTask struct {
ID string `json:"id" db:"id"`
Name string `json:"name" db:"name"`
Description string `json:"description" db:"description"`
CreatorID string `json:"creator_id" db:"creator_id"`
RunAs string `json:"run_as" db:"run_as"`
CronExpr string `json:"cron_expr" db:"cron_expr"`
NextFireAt *time.Time `json:"next_fire_at,omitempty" db:"next_fire_at"`
LastFireAt *time.Time `json:"last_fire_at,omitempty" db:"last_fire_at"`
Enabled bool `json:"enabled" db:"enabled"`
Script string `json:"script" db:"script"`
TemplateID string `json:"template_id,omitempty" db:"template_id"`
TemplateParams json.RawMessage `json:"template_params" db:"template_params"`
FireCount int `json:"fire_count" db:"fire_count"`
LastError string `json:"last_error,omitempty" db:"last_error"`
LastDurationMs *int `json:"last_duration_ms,omitempty" db:"last_duration_ms"`
CreatedAt time.Time `json:"created_at" db:"created_at"`
UpdatedAt time.Time `json:"updated_at" db:"updated_at"`
}
// ── Trigger Log (unified) ────────────────────
// TriggerLog records a single execution of a trigger or scheduled task.
type TriggerLog struct {
ID string `json:"id" db:"id"`
TriggerID string `json:"trigger_id,omitempty" db:"trigger_id"`
ScheduledTaskID string `json:"scheduled_task_id,omitempty" db:"scheduled_task_id"`
FiredAt string `json:"fired_at" db:"fired_at"`
DurationMs *int `json:"duration_ms,omitempty" db:"duration_ms"`
Success bool `json:"success" db:"success"`
Error string `json:"error,omitempty" db:"error"`
Output string `json:"output,omitempty" db:"output"`
}

View File

@@ -59,6 +59,10 @@ tags:
- name: 'Admin: Packages'
- name: 'Admin: System'
- name: 'Admin: Workflows'
- name: Triggers
- name: Schedules
- name: 'Admin: Triggers'
- name: 'Admin: Schedules'
security:
- BearerAuth: []
@@ -82,6 +86,56 @@ components:
type: string
required: [error]
Trigger:
type: object
properties:
id: { type: string, format: uuid }
package_id: { type: string }
type: { type: string, enum: [webhook, event] }
enabled: { type: boolean }
slug: { type: string }
event_pattern: { type: string }
entry_point: { type: string }
config: { type: object }
fire_count: { type: integer }
last_fire_at: { type: string, format: date-time, nullable: true }
last_error: { type: string }
created_at: { type: string, format: date-time }
updated_at: { type: string, format: date-time }
ScheduledTask:
type: object
properties:
id: { type: string, format: uuid }
name: { type: string }
description: { type: string }
creator_id: { type: string, format: uuid }
run_as: { type: string, enum: [creator, system] }
cron_expr: { type: string }
next_fire_at: { type: string, format: date-time, nullable: true }
last_fire_at: { type: string, format: date-time, nullable: true }
enabled: { type: boolean }
script: { type: string }
template_id: { type: string, nullable: true }
template_params: { type: object }
fire_count: { type: integer }
last_error: { type: string }
last_duration_ms: { type: integer, nullable: true }
created_at: { type: string, format: date-time }
updated_at: { type: string, format: date-time }
TriggerLog:
type: object
properties:
id: { type: string, format: uuid }
trigger_id: { type: string, format: uuid, nullable: true }
scheduled_task_id: { type: string, format: uuid, nullable: true }
fired_at: { type: string, format: date-time }
duration_ms: { type: integer, nullable: true }
success: { type: boolean }
error: { type: string }
output: { type: string }
TokenResponse:
type: object
properties:
@@ -4205,3 +4259,347 @@ paths:
application/json:
schema:
$ref: '#/components/schemas/ErrorResponse'
# ──────────────────────────────────────────────
# Webhook Inbound (v0.2.2)
# ──────────────────────────────────────────────
/api/v1/hooks/{package_id}/{slug}:
post:
summary: Inbound webhook trigger
description: Public endpoint. HMAC-SHA256 verified via X-Switchboard-Signature header.
operationId: webhookInbound
tags: [Triggers]
security: []
parameters:
- name: package_id
in: path
required: true
schema: { type: string }
- name: slug
in: path
required: true
schema: { type: string }
requestBody:
content:
application/json:
schema: { type: object }
responses:
'200':
description: Webhook processed
content:
application/json:
schema: { type: object }
'401':
description: Invalid signature
'404':
description: Webhook not found
# ──────────────────────────────────────────────
# Scheduled Tasks — User (v0.2.2)
# ──────────────────────────────────────────────
/api/v1/schedules:
get:
summary: List my scheduled tasks
operationId: listMySchedules
tags: [Schedules]
responses:
'200':
description: Schedules list
content:
application/json:
schema:
type: object
properties:
schedules:
type: array
items: { $ref: '#/components/schemas/ScheduledTask' }
post:
summary: Create a scheduled task
operationId: createSchedule
tags: [Schedules]
requestBody:
required: true
content:
application/json:
schema:
type: object
required: [name, cron_expr]
properties:
name: { type: string }
description: { type: string }
cron_expr: { type: string, example: '0 9 * * 1-5' }
script: { type: string }
template_id: { type: string }
template_params: { type: object }
responses:
'201':
description: Created
content:
application/json:
schema: { $ref: '#/components/schemas/ScheduledTask' }
/api/v1/schedules/{id}:
get:
summary: Get a scheduled task
operationId: getSchedule
tags: [Schedules]
parameters:
- $ref: '#/components/parameters/idPath'
responses:
'200':
description: Schedule details
content:
application/json:
schema: { $ref: '#/components/schemas/ScheduledTask' }
put:
summary: Update a scheduled task
operationId: updateSchedule
tags: [Schedules]
parameters:
- $ref: '#/components/parameters/idPath'
requestBody:
content:
application/json:
schema:
type: object
properties:
name: { type: string }
description: { type: string }
cron_expr: { type: string }
script: { type: string }
enabled: { type: boolean }
responses:
'200':
description: Updated
delete:
summary: Delete a scheduled task
operationId: deleteSchedule
tags: [Schedules]
parameters:
- $ref: '#/components/parameters/idPath'
responses:
'200':
description: Deleted
/api/v1/schedules/{id}/run:
post:
summary: Manually trigger a scheduled task
operationId: runSchedule
tags: [Schedules]
parameters:
- $ref: '#/components/parameters/idPath'
responses:
'202':
description: Task queued
/api/v1/schedules/{id}/logs:
get:
summary: List execution logs for a scheduled task
operationId: listScheduleLogs
tags: [Schedules]
parameters:
- $ref: '#/components/parameters/idPath'
- name: limit
in: query
schema: { type: integer, default: 50 }
responses:
'200':
description: Execution logs
content:
application/json:
schema:
type: object
properties:
logs:
type: array
items: { $ref: '#/components/schemas/TriggerLog' }
# ──────────────────────────────────────────────
# Admin: Triggers (v0.2.2)
# ──────────────────────────────────────────────
/api/v1/admin/triggers:
get:
summary: List all extension triggers
operationId: adminListTriggers
tags: ['Admin: Triggers']
parameters:
- name: package_id
in: query
schema: { type: string }
- name: type
in: query
schema: { type: string, enum: [webhook, event] }
- name: enabled
in: query
schema: { type: string, enum: ['true', 'false'] }
- name: limit
in: query
schema: { type: integer, default: 50 }
- name: offset
in: query
schema: { type: integer, default: 0 }
responses:
'200':
description: Trigger list
content:
application/json:
schema:
type: object
properties:
triggers:
type: array
items: { $ref: '#/components/schemas/Trigger' }
total: { type: integer }
/api/v1/admin/triggers/{id}:
get:
summary: Get a single trigger
operationId: adminGetTrigger
tags: ['Admin: Triggers']
parameters:
- $ref: '#/components/parameters/idPath'
responses:
'200':
description: Trigger details
content:
application/json:
schema: { $ref: '#/components/schemas/Trigger' }
delete:
summary: Delete a trigger
operationId: adminDeleteTrigger
tags: ['Admin: Triggers']
parameters:
- $ref: '#/components/parameters/idPath'
responses:
'200':
description: Deleted
/api/v1/admin/triggers/{id}/enable:
put:
summary: Enable a trigger
operationId: adminEnableTrigger
tags: ['Admin: Triggers']
parameters:
- $ref: '#/components/parameters/idPath'
responses:
'200':
description: Enabled
/api/v1/admin/triggers/{id}/disable:
put:
summary: Disable a trigger
operationId: adminDisableTrigger
tags: ['Admin: Triggers']
parameters:
- $ref: '#/components/parameters/idPath'
responses:
'200':
description: Disabled
/api/v1/admin/triggers/{id}/logs:
get:
summary: List trigger execution logs
operationId: adminListTriggerLogs
tags: ['Admin: Triggers']
parameters:
- $ref: '#/components/parameters/idPath'
- name: limit
in: query
schema: { type: integer, default: 50 }
responses:
'200':
description: Execution logs
content:
application/json:
schema:
type: object
properties:
logs:
type: array
items: { $ref: '#/components/schemas/TriggerLog' }
/api/v1/admin/packages/{id}/triggers:
get:
summary: List triggers for a package
operationId: adminListPackageTriggers
tags: ['Admin: Triggers']
parameters:
- $ref: '#/components/parameters/idPath'
responses:
'200':
description: Package triggers
content:
application/json:
schema:
type: object
properties:
triggers:
type: array
items: { $ref: '#/components/schemas/Trigger' }
# ──────────────────────────────────────────────
# Admin: Schedules (v0.2.2)
# ──────────────────────────────────────────────
/api/v1/admin/schedules:
get:
summary: List all scheduled tasks (admin view)
operationId: adminListSchedules
tags: ['Admin: Schedules']
parameters:
- name: creator_id
in: query
schema: { type: string }
- name: enabled
in: query
schema: { type: string, enum: ['true', 'false'] }
- name: limit
in: query
schema: { type: integer, default: 50 }
- name: offset
in: query
schema: { type: integer, default: 0 }
responses:
'200':
description: Schedules list
content:
application/json:
schema:
type: object
properties:
schedules:
type: array
items: { $ref: '#/components/schemas/ScheduledTask' }
total: { type: integer }
/api/v1/admin/schedules/{id}/enable:
put:
summary: Enable a scheduled task
operationId: adminEnableSchedule
tags: ['Admin: Schedules']
parameters:
- $ref: '#/components/parameters/idPath'
responses:
'200':
description: Enabled
/api/v1/admin/schedules/{id}/disable:
put:
summary: Disable a scheduled task
operationId: adminDisableSchedule
tags: ['Admin: Schedules']
parameters:
- $ref: '#/components/parameters/idPath'
responses:
'200':
description: Disabled
/api/v1/admin/schedules/{id}:
delete:
summary: Delete a scheduled task
operationId: adminDeleteSchedule
tags: ['Admin: Schedules']
parameters:
- $ref: '#/components/parameters/idPath'
responses:
'200':
description: Deleted

View File

@@ -52,6 +52,8 @@ type Stores struct {
ExtData ExtDataStore // v0.29.2: Extension namespaced table catalog
Tickets TicketStore // v0.32.0: WS auth tickets (PG-backed for cross-pod)
RateLimits RateLimitStore // v0.32.0: Distributed rate limiting
Triggers TriggerStore // v0.2.2: Extension event/webhook triggers
ScheduledTasks ScheduledTaskStore // v0.2.2: User-created cron tasks
}
// TeamAvailableModel is returned by CatalogStore.ListTeamAvailable.

View File

@@ -0,0 +1,216 @@
package postgres
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"time"
"switchboard-core/models"
"switchboard-core/store"
)
type ScheduledTaskStore struct{}
func NewScheduledTaskStore() *ScheduledTaskStore {
return &ScheduledTaskStore{}
}
func (s *ScheduledTaskStore) Create(ctx context.Context, t *models.ScheduledTask) error {
params, _ := json.Marshal(t.TemplateParams)
if len(params) == 0 {
params = []byte("{}")
}
return DB.QueryRowContext(ctx,
`INSERT INTO scheduled_tasks (name, description, creator_id, run_as, cron_expr,
next_fire_at, enabled, script, template_id, template_params)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
RETURNING id, created_at, updated_at`,
t.Name, t.Description, t.CreatorID, t.RunAs, t.CronExpr,
t.NextFireAt, t.Enabled, t.Script, nilIfEmpty(t.TemplateID), params).
Scan(&t.ID, &t.CreatedAt, &t.UpdatedAt)
}
func (s *ScheduledTaskStore) GetByID(ctx context.Context, id string) (*models.ScheduledTask, error) {
var t models.ScheduledTask
var params []byte
err := DB.QueryRowContext(ctx,
`SELECT id, name, description, creator_id, run_as, cron_expr,
next_fire_at, last_fire_at, enabled, script,
COALESCE(template_id,''), template_params, fire_count,
COALESCE(last_error,''), last_duration_ms, created_at, updated_at
FROM scheduled_tasks WHERE id = $1`, id).
Scan(&t.ID, &t.Name, &t.Description, &t.CreatorID, &t.RunAs, &t.CronExpr,
&t.NextFireAt, &t.LastFireAt, &t.Enabled, &t.Script,
&t.TemplateID, &params, &t.FireCount,
&t.LastError, &t.LastDurationMs, &t.CreatedAt, &t.UpdatedAt)
if err == sql.ErrNoRows {
return nil, nil
}
if err != nil {
return nil, err
}
t.TemplateParams = json.RawMessage(params)
return &t, nil
}
func (s *ScheduledTaskStore) Update(ctx context.Context, t *models.ScheduledTask) error {
params, _ := json.Marshal(t.TemplateParams)
if len(params) == 0 {
params = []byte("{}")
}
_, err := DB.ExecContext(ctx,
`UPDATE scheduled_tasks SET name = $1, description = $2, run_as = $3,
cron_expr = $4, next_fire_at = $5, enabled = $6, script = $7,
template_id = $8, template_params = $9, updated_at = NOW()
WHERE id = $10`,
t.Name, t.Description, t.RunAs, t.CronExpr, t.NextFireAt,
t.Enabled, t.Script, nilIfEmpty(t.TemplateID), params, t.ID)
return err
}
func (s *ScheduledTaskStore) Delete(ctx context.Context, id string) error {
_, err := DB.ExecContext(ctx, `DELETE FROM scheduled_tasks WHERE id = $1`, id)
return err
}
func (s *ScheduledTaskStore) List(ctx context.Context, opts store.ScheduledTaskListOptions) ([]models.ScheduledTask, int, error) {
where := "1=1"
args := []any{}
n := 0
if opts.CreatorID != "" {
n++
where += fmt.Sprintf(" AND creator_id = $%d", n)
args = append(args, opts.CreatorID)
}
if opts.Enabled != nil {
n++
where += fmt.Sprintf(" AND enabled = $%d", n)
args = append(args, *opts.Enabled)
}
var total int
countArgs := make([]any, len(args))
copy(countArgs, args)
err := DB.QueryRowContext(ctx, "SELECT COUNT(*) FROM scheduled_tasks WHERE "+where, countArgs...).Scan(&total)
if err != nil {
return nil, 0, err
}
limit := opts.Limit
if limit <= 0 {
limit = 50
}
n++
query := fmt.Sprintf(
`SELECT id, name, description, creator_id, run_as, cron_expr,
next_fire_at, last_fire_at, enabled, script,
COALESCE(template_id,''), template_params, fire_count,
COALESCE(last_error,''), last_duration_ms, created_at, updated_at
FROM scheduled_tasks WHERE %s ORDER BY created_at DESC LIMIT $%d OFFSET $%d`,
where, n, n+1)
args = append(args, limit, opts.Offset)
rows, err := DB.QueryContext(ctx, query, args...)
if err != nil {
return nil, 0, err
}
defer rows.Close()
tasks := scanScheduledTaskRows(rows)
return tasks, total, rows.Err()
}
func (s *ScheduledTaskStore) ListByCreator(ctx context.Context, creatorID string) ([]models.ScheduledTask, error) {
rows, err := DB.QueryContext(ctx,
`SELECT id, name, description, creator_id, run_as, cron_expr,
next_fire_at, last_fire_at, enabled, script,
COALESCE(template_id,''), template_params, fire_count,
COALESCE(last_error,''), last_duration_ms, created_at, updated_at
FROM scheduled_tasks WHERE creator_id = $1 ORDER BY created_at DESC`, creatorID)
if err != nil {
return nil, err
}
defer rows.Close()
return scanScheduledTaskRowsErr(rows)
}
func (s *ScheduledTaskStore) ListEnabled(ctx context.Context) ([]models.ScheduledTask, error) {
rows, err := DB.QueryContext(ctx,
`SELECT id, name, description, creator_id, run_as, cron_expr,
next_fire_at, last_fire_at, enabled, script,
COALESCE(template_id,''), template_params, fire_count,
COALESCE(last_error,''), last_duration_ms, created_at, updated_at
FROM scheduled_tasks WHERE enabled = true ORDER BY created_at`)
if err != nil {
return nil, err
}
defer rows.Close()
return scanScheduledTaskRowsErr(rows)
}
func (s *ScheduledTaskStore) SetEnabled(ctx context.Context, id string, enabled bool) error {
_, err := DB.ExecContext(ctx,
`UPDATE scheduled_tasks SET enabled = $1, updated_at = NOW() WHERE id = $2`, enabled, id)
return err
}
func (s *ScheduledTaskStore) UpdateFireState(ctx context.Context, id string, lastFireAt time.Time, nextFireAt *time.Time, lastError string, durationMs int) error {
_, err := DB.ExecContext(ctx,
`UPDATE scheduled_tasks SET last_fire_at = $1, next_fire_at = $2,
last_error = $3, last_duration_ms = $4, fire_count = fire_count + 1,
updated_at = NOW() WHERE id = $5`,
lastFireAt, nextFireAt, nilIfEmpty(lastError), durationMs, id)
return err
}
// ── Helpers ──────────────────────────────────
func scanScheduledTaskRows(rows *sql.Rows) []models.ScheduledTask {
var tasks []models.ScheduledTask
for rows.Next() {
t := scanOneScheduledTask(rows)
if t != nil {
tasks = append(tasks, *t)
}
}
if tasks == nil {
tasks = []models.ScheduledTask{}
}
return tasks
}
func scanScheduledTaskRowsErr(rows *sql.Rows) ([]models.ScheduledTask, error) {
var tasks []models.ScheduledTask
for rows.Next() {
var t models.ScheduledTask
var params []byte
if err := rows.Scan(&t.ID, &t.Name, &t.Description, &t.CreatorID, &t.RunAs,
&t.CronExpr, &t.NextFireAt, &t.LastFireAt, &t.Enabled, &t.Script,
&t.TemplateID, &params, &t.FireCount,
&t.LastError, &t.LastDurationMs, &t.CreatedAt, &t.UpdatedAt); err != nil {
return nil, err
}
t.TemplateParams = json.RawMessage(params)
tasks = append(tasks, t)
}
if tasks == nil {
tasks = []models.ScheduledTask{}
}
return tasks, nil
}
func scanOneScheduledTask(rows *sql.Rows) *models.ScheduledTask {
var t models.ScheduledTask
var params []byte
if err := rows.Scan(&t.ID, &t.Name, &t.Description, &t.CreatorID, &t.RunAs,
&t.CronExpr, &t.NextFireAt, &t.LastFireAt, &t.Enabled, &t.Script,
&t.TemplateID, &params, &t.FireCount,
&t.LastError, &t.LastDurationMs, &t.CreatedAt, &t.UpdatedAt); err != nil {
return nil
}
t.TemplateParams = json.RawMessage(params)
return &t
}

View File

@@ -30,5 +30,7 @@ func NewStores(db *sql.DB) store.Stores {
ExtData: NewExtDataStore(db),
Tickets: NewTicketStore(),
RateLimits: NewRateLimitStore(),
Triggers: NewTriggerStore(),
ScheduledTasks: NewScheduledTaskStore(),
}
}

View File

@@ -0,0 +1,287 @@
package postgres
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"time"
"switchboard-core/models"
"switchboard-core/store"
)
type TriggerStore struct{}
func NewTriggerStore() *TriggerStore {
return &TriggerStore{}
}
func (s *TriggerStore) Create(ctx context.Context, t *models.Trigger) error {
cfg, _ := json.Marshal(t.Config)
if len(cfg) == 0 {
cfg = []byte("{}")
}
_, err := DB.ExecContext(ctx,
`INSERT INTO triggers (id, package_id, type, enabled, slug, secret, event_pattern,
entry_point, config)
VALUES (gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8)
RETURNING id`,
t.PackageID, t.Type, t.Enabled, nilIfEmpty(t.Slug), nilIfEmpty(t.Secret),
nilIfEmpty(t.EventPattern), t.EntryPoint, cfg)
return err
}
func (s *TriggerStore) GetByID(ctx context.Context, id string) (*models.Trigger, error) {
var t models.Trigger
var cfg []byte
err := DB.QueryRowContext(ctx,
`SELECT id, package_id, type, enabled, COALESCE(slug,''), secret,
COALESCE(event_pattern,''), entry_point, config, fire_count,
last_fire_at, COALESCE(last_error,''), created_at, updated_at
FROM triggers WHERE id = $1`, id).
Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug, &t.Secret,
&t.EventPattern, &t.EntryPoint, &cfg, &t.FireCount,
&t.LastFireAt, &t.LastError, &t.CreatedAt, &t.UpdatedAt)
if err == sql.ErrNoRows {
return nil, nil
}
if err != nil {
return nil, err
}
t.Config = json.RawMessage(cfg)
return &t, nil
}
func (s *TriggerStore) Update(ctx context.Context, t *models.Trigger) error {
cfg, _ := json.Marshal(t.Config)
if len(cfg) == 0 {
cfg = []byte("{}")
}
_, err := DB.ExecContext(ctx,
`UPDATE triggers SET enabled = $1, slug = $2, secret = $3, event_pattern = $4,
entry_point = $5, config = $6, updated_at = NOW()
WHERE id = $7`,
t.Enabled, nilIfEmpty(t.Slug), nilIfEmpty(t.Secret),
nilIfEmpty(t.EventPattern), t.EntryPoint, cfg, t.ID)
return err
}
func (s *TriggerStore) Delete(ctx context.Context, id string) error {
_, err := DB.ExecContext(ctx, `DELETE FROM triggers WHERE id = $1`, id)
return err
}
func (s *TriggerStore) List(ctx context.Context, opts store.TriggerListOptions) ([]models.Trigger, int, error) {
where := "1=1"
args := []any{}
n := 0
if opts.PackageID != "" {
n++
where += fmt.Sprintf(" AND package_id = $%d", n)
args = append(args, opts.PackageID)
}
if opts.Type != "" {
n++
where += fmt.Sprintf(" AND type = $%d", n)
args = append(args, opts.Type)
}
if opts.Enabled != nil {
n++
where += fmt.Sprintf(" AND enabled = $%d", n)
args = append(args, *opts.Enabled)
}
// Count
var total int
countArgs := make([]any, len(args))
copy(countArgs, args)
err := DB.QueryRowContext(ctx, "SELECT COUNT(*) FROM triggers WHERE "+where, countArgs...).Scan(&total)
if err != nil {
return nil, 0, err
}
// Query
limit := opts.Limit
if limit <= 0 {
limit = 50
}
n++
query := fmt.Sprintf(
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
created_at, updated_at
FROM triggers WHERE %s ORDER BY created_at DESC LIMIT $%d OFFSET $%d`,
where, n, n+1)
args = append(args, limit, opts.Offset)
rows, err := DB.QueryContext(ctx, query, args...)
if err != nil {
return nil, 0, err
}
defer rows.Close()
var triggers []models.Trigger
for rows.Next() {
var t models.Trigger
var cfg []byte
if err := rows.Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug,
&t.EventPattern, &t.EntryPoint, &cfg, &t.FireCount,
&t.LastFireAt, &t.LastError, &t.CreatedAt, &t.UpdatedAt); err != nil {
return nil, 0, err
}
t.Config = json.RawMessage(cfg)
triggers = append(triggers, t)
}
if triggers == nil {
triggers = []models.Trigger{}
}
return triggers, total, nil
}
func (s *TriggerStore) ListByPackage(ctx context.Context, packageID string) ([]models.Trigger, error) {
rows, err := DB.QueryContext(ctx,
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
created_at, updated_at
FROM triggers WHERE package_id = $1 ORDER BY created_at`, packageID)
if err != nil {
return nil, err
}
defer rows.Close()
return scanTriggerRows(rows)
}
func (s *TriggerStore) ListEnabledByType(ctx context.Context, triggerType string) ([]models.Trigger, error) {
rows, err := DB.QueryContext(ctx,
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
created_at, updated_at
FROM triggers WHERE type = $1 AND enabled = true ORDER BY created_at`, triggerType)
if err != nil {
return nil, err
}
defer rows.Close()
return scanTriggerRows(rows)
}
func (s *TriggerStore) SetEnabled(ctx context.Context, id string, enabled bool) error {
_, err := DB.ExecContext(ctx,
`UPDATE triggers SET enabled = $1, updated_at = NOW() WHERE id = $2`, enabled, id)
return err
}
func (s *TriggerStore) UpdateFireState(ctx context.Context, id string, lastFireAt time.Time, lastError string) error {
_, err := DB.ExecContext(ctx,
`UPDATE triggers SET last_fire_at = $1, last_error = $2, fire_count = fire_count + 1,
updated_at = NOW() WHERE id = $3`,
lastFireAt, nilIfEmpty(lastError), id)
return err
}
func (s *TriggerStore) GetWebhook(ctx context.Context, packageID, slug string) (*models.Trigger, error) {
var t models.Trigger
var cfg []byte
err := DB.QueryRowContext(ctx,
`SELECT id, package_id, type, enabled, COALESCE(slug,''), secret,
COALESCE(event_pattern,''), entry_point, config, fire_count,
last_fire_at, COALESCE(last_error,''), created_at, updated_at
FROM triggers
WHERE package_id = $1 AND slug = $2 AND type = 'webhook'`, packageID, slug).
Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug, &t.Secret,
&t.EventPattern, &t.EntryPoint, &cfg, &t.FireCount,
&t.LastFireAt, &t.LastError, &t.CreatedAt, &t.UpdatedAt)
if err == sql.ErrNoRows {
return nil, nil
}
if err != nil {
return nil, err
}
t.Config = json.RawMessage(cfg)
return &t, nil
}
func (s *TriggerStore) DeleteForPackage(ctx context.Context, packageID string) error {
_, err := DB.ExecContext(ctx, `DELETE FROM triggers WHERE package_id = $1`, packageID)
return err
}
func (s *TriggerStore) LogExecution(ctx context.Context, log *models.TriggerLog) error {
_, err := DB.ExecContext(ctx,
`INSERT INTO trigger_logs (id, trigger_id, scheduled_task_id, fired_at, duration_ms, success, error, output)
VALUES (gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7)`,
nilIfEmpty(log.TriggerID), nilIfEmpty(log.ScheduledTaskID),
log.FiredAt, log.DurationMs, log.Success, nilIfEmpty(log.Error), nilIfEmpty(log.Output))
return err
}
func (s *TriggerStore) ListLogs(ctx context.Context, triggerID string, limit int) ([]models.TriggerLog, error) {
if limit <= 0 {
limit = 50
}
rows, err := DB.QueryContext(ctx,
`SELECT id, COALESCE(trigger_id::text,''), COALESCE(scheduled_task_id::text,''),
fired_at, duration_ms, success, COALESCE(error,''), COALESCE(output,'')
FROM trigger_logs
WHERE trigger_id = $1
ORDER BY fired_at DESC LIMIT $2`, triggerID, limit)
if err != nil {
return nil, err
}
defer rows.Close()
return scanLogRows(rows)
}
func (s *TriggerStore) PruneLogs(ctx context.Context, before time.Time) (int64, error) {
result, err := DB.ExecContext(ctx,
`DELETE FROM trigger_logs WHERE fired_at < $1`, before)
if err != nil {
return 0, err
}
return result.RowsAffected()
}
// ── Helpers ──────────────────────────────────
func scanTriggerRows(rows *sql.Rows) ([]models.Trigger, error) {
var triggers []models.Trigger
for rows.Next() {
var t models.Trigger
var cfg []byte
if err := rows.Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug,
&t.EventPattern, &t.EntryPoint, &cfg, &t.FireCount,
&t.LastFireAt, &t.LastError, &t.CreatedAt, &t.UpdatedAt); err != nil {
return nil, err
}
t.Config = json.RawMessage(cfg)
triggers = append(triggers, t)
}
if triggers == nil {
triggers = []models.Trigger{}
}
return triggers, nil
}
func scanLogRows(rows *sql.Rows) ([]models.TriggerLog, error) {
var logs []models.TriggerLog
for rows.Next() {
var l models.TriggerLog
if err := rows.Scan(&l.ID, &l.TriggerID, &l.ScheduledTaskID,
&l.FiredAt, &l.DurationMs, &l.Success, &l.Error, &l.Output); err != nil {
return nil, err
}
logs = append(logs, l)
}
if logs == nil {
logs = []models.TriggerLog{}
}
return logs, nil
}
func nilIfEmpty(s string) *string {
if s == "" {
return nil
}
return &s
}

View File

@@ -0,0 +1,33 @@
package store
import (
"context"
"time"
"switchboard-core/models"
)
// ScheduledTaskStore manages user-created cron-scheduled Starlark scripts.
type ScheduledTaskStore interface {
// CRUD
Create(ctx context.Context, t *models.ScheduledTask) error
GetByID(ctx context.Context, id string) (*models.ScheduledTask, error)
Update(ctx context.Context, t *models.ScheduledTask) error
Delete(ctx context.Context, id string) error
// Queries
List(ctx context.Context, opts ScheduledTaskListOptions) ([]models.ScheduledTask, int, error)
ListByCreator(ctx context.Context, creatorID string) ([]models.ScheduledTask, error)
ListEnabled(ctx context.Context) ([]models.ScheduledTask, error)
// Lifecycle
SetEnabled(ctx context.Context, id string, enabled bool) error
UpdateFireState(ctx context.Context, id string, lastFireAt time.Time, nextFireAt *time.Time, lastError string, durationMs int) error
}
// ScheduledTaskListOptions provides filtering for scheduled task list queries.
type ScheduledTaskListOptions struct {
ListOptions
CreatorID string
Enabled *bool
}

View File

@@ -0,0 +1,216 @@
package sqlite
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"time"
"github.com/google/uuid"
"switchboard-core/database"
"switchboard-core/models"
"switchboard-core/store"
)
type ScheduledTaskStore struct{}
func NewScheduledTaskStore() *ScheduledTaskStore {
return &ScheduledTaskStore{}
}
func (s *ScheduledTaskStore) Create(ctx context.Context, t *models.ScheduledTask) error {
t.ID = uuid.New().String()
now := time.Now().UTC().Format(time.RFC3339)
t.CreatedAt, _ = time.Parse(time.RFC3339, now)
t.UpdatedAt = t.CreatedAt
params, _ := json.Marshal(t.TemplateParams)
if len(params) == 0 {
params = []byte("{}")
}
var nextFire *string
if t.NextFireAt != nil {
s := t.NextFireAt.UTC().Format(time.RFC3339)
nextFire = &s
}
_, err := DB.ExecContext(ctx,
`INSERT INTO scheduled_tasks (id, name, description, creator_id, run_as, cron_expr,
next_fire_at, enabled, script, template_id, template_params, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
t.ID, t.Name, t.Description, t.CreatorID, t.RunAs, t.CronExpr,
nextFire, boolToInt(t.Enabled), t.Script,
nullIfEmpty(t.TemplateID), string(params), now, now)
return err
}
func (s *ScheduledTaskStore) GetByID(ctx context.Context, id string) (*models.ScheduledTask, error) {
var t models.ScheduledTask
var paramsStr string
err := DB.QueryRowContext(ctx,
`SELECT id, name, description, creator_id, run_as, cron_expr,
next_fire_at, last_fire_at, enabled, script,
COALESCE(template_id,''), template_params, fire_count,
COALESCE(last_error,''), last_duration_ms, created_at, updated_at
FROM scheduled_tasks WHERE id = ?`, id).
Scan(&t.ID, &t.Name, &t.Description, &t.CreatorID, &t.RunAs, &t.CronExpr,
database.SNT(&t.NextFireAt), database.SNT(&t.LastFireAt),
&t.Enabled, &t.Script,
&t.TemplateID, &paramsStr, &t.FireCount,
&t.LastError, &t.LastDurationMs,
database.ST(&t.CreatedAt), database.ST(&t.UpdatedAt))
if err == sql.ErrNoRows {
return nil, nil
}
if err != nil {
return nil, err
}
t.TemplateParams = json.RawMessage(paramsStr)
return &t, nil
}
func (s *ScheduledTaskStore) Update(ctx context.Context, t *models.ScheduledTask) error {
params, _ := json.Marshal(t.TemplateParams)
if len(params) == 0 {
params = []byte("{}")
}
var nextFire *string
if t.NextFireAt != nil {
s := t.NextFireAt.UTC().Format(time.RFC3339)
nextFire = &s
}
_, err := DB.ExecContext(ctx,
`UPDATE scheduled_tasks SET name = ?, description = ?, run_as = ?,
cron_expr = ?, next_fire_at = ?, enabled = ?, script = ?,
template_id = ?, template_params = ?, updated_at = datetime('now')
WHERE id = ?`,
t.Name, t.Description, t.RunAs, t.CronExpr, nextFire,
boolToInt(t.Enabled), t.Script,
nullIfEmpty(t.TemplateID), string(params), t.ID)
return err
}
func (s *ScheduledTaskStore) Delete(ctx context.Context, id string) error {
_, err := DB.ExecContext(ctx, `DELETE FROM scheduled_tasks WHERE id = ?`, id)
return err
}
func (s *ScheduledTaskStore) List(ctx context.Context, opts store.ScheduledTaskListOptions) ([]models.ScheduledTask, int, error) {
where := "1=1"
args := []any{}
if opts.CreatorID != "" {
where += " AND creator_id = ?"
args = append(args, opts.CreatorID)
}
if opts.Enabled != nil {
where += " AND enabled = ?"
args = append(args, boolToInt(*opts.Enabled))
}
var total int
countArgs := make([]any, len(args))
copy(countArgs, args)
err := DB.QueryRowContext(ctx, "SELECT COUNT(*) FROM scheduled_tasks WHERE "+where, countArgs...).Scan(&total)
if err != nil {
return nil, 0, err
}
limit := opts.Limit
if limit <= 0 {
limit = 50
}
query := fmt.Sprintf(
`SELECT id, name, description, creator_id, run_as, cron_expr,
next_fire_at, last_fire_at, enabled, script,
COALESCE(template_id,''), template_params, fire_count,
COALESCE(last_error,''), last_duration_ms, created_at, updated_at
FROM scheduled_tasks WHERE %s ORDER BY created_at DESC LIMIT ? OFFSET ?`, where)
args = append(args, limit, opts.Offset)
rows, err := DB.QueryContext(ctx, query, args...)
if err != nil {
return nil, 0, err
}
defer rows.Close()
tasks, err := scanScheduledTaskRowsSqlite(rows)
return tasks, total, err
}
func (s *ScheduledTaskStore) ListByCreator(ctx context.Context, creatorID string) ([]models.ScheduledTask, error) {
rows, err := DB.QueryContext(ctx,
`SELECT id, name, description, creator_id, run_as, cron_expr,
next_fire_at, last_fire_at, enabled, script,
COALESCE(template_id,''), template_params, fire_count,
COALESCE(last_error,''), last_duration_ms, created_at, updated_at
FROM scheduled_tasks WHERE creator_id = ? ORDER BY created_at DESC`, creatorID)
if err != nil {
return nil, err
}
defer rows.Close()
return scanScheduledTaskRowsSqlite(rows)
}
func (s *ScheduledTaskStore) ListEnabled(ctx context.Context) ([]models.ScheduledTask, error) {
rows, err := DB.QueryContext(ctx,
`SELECT id, name, description, creator_id, run_as, cron_expr,
next_fire_at, last_fire_at, enabled, script,
COALESCE(template_id,''), template_params, fire_count,
COALESCE(last_error,''), last_duration_ms, created_at, updated_at
FROM scheduled_tasks WHERE enabled = 1 ORDER BY created_at`)
if err != nil {
return nil, err
}
defer rows.Close()
return scanScheduledTaskRowsSqlite(rows)
}
func (s *ScheduledTaskStore) SetEnabled(ctx context.Context, id string, enabled bool) error {
_, err := DB.ExecContext(ctx,
`UPDATE scheduled_tasks SET enabled = ?, updated_at = datetime('now') WHERE id = ?`,
boolToInt(enabled), id)
return err
}
func (s *ScheduledTaskStore) UpdateFireState(ctx context.Context, id string, lastFireAt time.Time, nextFireAt *time.Time, lastError string, durationMs int) error {
var nextFire *string
if nextFireAt != nil {
s := nextFireAt.UTC().Format(time.RFC3339)
nextFire = &s
}
_, err := DB.ExecContext(ctx,
`UPDATE scheduled_tasks SET last_fire_at = ?, next_fire_at = ?,
last_error = ?, last_duration_ms = ?, fire_count = fire_count + 1,
updated_at = datetime('now') WHERE id = ?`,
lastFireAt.UTC().Format(time.RFC3339), nextFire,
nullIfEmpty(lastError), durationMs, id)
return err
}
// ── Helpers ──────────────────────────────────
func scanScheduledTaskRowsSqlite(rows *sql.Rows) ([]models.ScheduledTask, error) {
var tasks []models.ScheduledTask
for rows.Next() {
var t models.ScheduledTask
var paramsStr string
if err := rows.Scan(&t.ID, &t.Name, &t.Description, &t.CreatorID, &t.RunAs,
&t.CronExpr, database.SNT(&t.NextFireAt), database.SNT(&t.LastFireAt),
&t.Enabled, &t.Script,
&t.TemplateID, &paramsStr, &t.FireCount,
&t.LastError, &t.LastDurationMs,
database.ST(&t.CreatedAt), database.ST(&t.UpdatedAt)); err != nil {
return nil, err
}
t.TemplateParams = json.RawMessage(paramsStr)
tasks = append(tasks, t)
}
if tasks == nil {
tasks = []models.ScheduledTask{}
}
return tasks, nil
}

View File

@@ -30,5 +30,7 @@ func NewStores(db *sql.DB) store.Stores {
ExtData: NewExtDataStore(),
Tickets: NewTicketStore(),
RateLimits: NewRateLimitStore(),
Triggers: NewTriggerStore(),
ScheduledTasks: NewScheduledTaskStore(),
}
}

View File

@@ -0,0 +1,269 @@
package sqlite
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"time"
"github.com/google/uuid"
"switchboard-core/database"
"switchboard-core/models"
"switchboard-core/store"
)
type TriggerStore struct{}
func NewTriggerStore() *TriggerStore {
return &TriggerStore{}
}
func (s *TriggerStore) Create(ctx context.Context, t *models.Trigger) error {
t.ID = uuid.New().String()
cfg, _ := json.Marshal(t.Config)
if len(cfg) == 0 {
cfg = []byte("{}")
}
_, err := DB.ExecContext(ctx,
`INSERT INTO triggers (id, package_id, type, enabled, slug, secret, event_pattern,
entry_point, config)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`,
t.ID, t.PackageID, t.Type, boolToInt(t.Enabled),
nullIfEmpty(t.Slug), nullIfEmpty(t.Secret),
nullIfEmpty(t.EventPattern), t.EntryPoint, string(cfg))
return err
}
func (s *TriggerStore) GetByID(ctx context.Context, id string) (*models.Trigger, error) {
var t models.Trigger
var cfgStr string
err := DB.QueryRowContext(ctx,
`SELECT id, package_id, type, enabled, COALESCE(slug,''), secret,
COALESCE(event_pattern,''), entry_point, config, fire_count,
last_fire_at, COALESCE(last_error,''), created_at, updated_at
FROM triggers WHERE id = ?`, id).
Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug, &t.Secret,
&t.EventPattern, &t.EntryPoint, &cfgStr, &t.FireCount,
database.SNT(&t.LastFireAt), &t.LastError,
database.ST(&t.CreatedAt), database.ST(&t.UpdatedAt))
if err == sql.ErrNoRows {
return nil, nil
}
if err != nil {
return nil, err
}
t.Config = json.RawMessage(cfgStr)
return &t, nil
}
func (s *TriggerStore) Update(ctx context.Context, t *models.Trigger) error {
cfg, _ := json.Marshal(t.Config)
if len(cfg) == 0 {
cfg = []byte("{}")
}
_, err := DB.ExecContext(ctx,
`UPDATE triggers SET enabled = ?, slug = ?, secret = ?, event_pattern = ?,
entry_point = ?, config = ?, updated_at = datetime('now')
WHERE id = ?`,
boolToInt(t.Enabled), nullIfEmpty(t.Slug), nullIfEmpty(t.Secret),
nullIfEmpty(t.EventPattern), t.EntryPoint, string(cfg), t.ID)
return err
}
func (s *TriggerStore) Delete(ctx context.Context, id string) error {
_, err := DB.ExecContext(ctx, `DELETE FROM triggers WHERE id = ?`, id)
return err
}
func (s *TriggerStore) List(ctx context.Context, opts store.TriggerListOptions) ([]models.Trigger, int, error) {
where := "1=1"
args := []any{}
if opts.PackageID != "" {
where += " AND package_id = ?"
args = append(args, opts.PackageID)
}
if opts.Type != "" {
where += " AND type = ?"
args = append(args, opts.Type)
}
if opts.Enabled != nil {
where += " AND enabled = ?"
args = append(args, boolToInt(*opts.Enabled))
}
var total int
countArgs := make([]any, len(args))
copy(countArgs, args)
err := DB.QueryRowContext(ctx, "SELECT COUNT(*) FROM triggers WHERE "+where, countArgs...).Scan(&total)
if err != nil {
return nil, 0, err
}
limit := opts.Limit
if limit <= 0 {
limit = 50
}
query := fmt.Sprintf(
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
created_at, updated_at
FROM triggers WHERE %s ORDER BY created_at DESC LIMIT ? OFFSET ?`, where)
args = append(args, limit, opts.Offset)
rows, err := DB.QueryContext(ctx, query, args...)
if err != nil {
return nil, 0, err
}
defer rows.Close()
triggers, err := scanTriggerRowsSqlite(rows)
return triggers, total, err
}
func (s *TriggerStore) ListByPackage(ctx context.Context, packageID string) ([]models.Trigger, error) {
rows, err := DB.QueryContext(ctx,
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
created_at, updated_at
FROM triggers WHERE package_id = ? ORDER BY created_at`, packageID)
if err != nil {
return nil, err
}
defer rows.Close()
return scanTriggerRowsSqlite(rows)
}
func (s *TriggerStore) ListEnabledByType(ctx context.Context, triggerType string) ([]models.Trigger, error) {
rows, err := DB.QueryContext(ctx,
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
created_at, updated_at
FROM triggers WHERE type = ? AND enabled = 1 ORDER BY created_at`, triggerType)
if err != nil {
return nil, err
}
defer rows.Close()
return scanTriggerRowsSqlite(rows)
}
func (s *TriggerStore) SetEnabled(ctx context.Context, id string, enabled bool) error {
_, err := DB.ExecContext(ctx,
`UPDATE triggers SET enabled = ?, updated_at = datetime('now') WHERE id = ?`,
boolToInt(enabled), id)
return err
}
func (s *TriggerStore) UpdateFireState(ctx context.Context, id string, lastFireAt time.Time, lastError string) error {
_, err := DB.ExecContext(ctx,
`UPDATE triggers SET last_fire_at = ?, last_error = ?, fire_count = fire_count + 1,
updated_at = datetime('now') WHERE id = ?`,
lastFireAt.UTC().Format(time.RFC3339), nullIfEmpty(lastError), id)
return err
}
func (s *TriggerStore) GetWebhook(ctx context.Context, packageID, slug string) (*models.Trigger, error) {
var t models.Trigger
var cfgStr string
err := DB.QueryRowContext(ctx,
`SELECT id, package_id, type, enabled, COALESCE(slug,''), secret,
COALESCE(event_pattern,''), entry_point, config, fire_count,
last_fire_at, COALESCE(last_error,''), created_at, updated_at
FROM triggers
WHERE package_id = ? AND slug = ? AND type = 'webhook'`, packageID, slug).
Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug, &t.Secret,
&t.EventPattern, &t.EntryPoint, &cfgStr, &t.FireCount,
database.SNT(&t.LastFireAt), &t.LastError,
database.ST(&t.CreatedAt), database.ST(&t.UpdatedAt))
if err == sql.ErrNoRows {
return nil, nil
}
if err != nil {
return nil, err
}
t.Config = json.RawMessage(cfgStr)
return &t, nil
}
func (s *TriggerStore) DeleteForPackage(ctx context.Context, packageID string) error {
_, err := DB.ExecContext(ctx, `DELETE FROM triggers WHERE package_id = ?`, packageID)
return err
}
func (s *TriggerStore) LogExecution(ctx context.Context, log *models.TriggerLog) error {
_, err := DB.ExecContext(ctx,
`INSERT INTO trigger_logs (id, trigger_id, scheduled_task_id, fired_at, duration_ms, success, error, output)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)`,
uuid.New().String(), nullIfEmpty(log.TriggerID), nullIfEmpty(log.ScheduledTaskID),
log.FiredAt, log.DurationMs, boolToInt(log.Success),
nullIfEmpty(log.Error), nullIfEmpty(log.Output))
return err
}
func (s *TriggerStore) ListLogs(ctx context.Context, triggerID string, limit int) ([]models.TriggerLog, error) {
if limit <= 0 {
limit = 50
}
rows, err := DB.QueryContext(ctx,
`SELECT id, COALESCE(trigger_id,''), COALESCE(scheduled_task_id,''),
fired_at, duration_ms, success, COALESCE(error,''), COALESCE(output,'')
FROM trigger_logs
WHERE trigger_id = ?
ORDER BY fired_at DESC LIMIT ?`, triggerID, limit)
if err != nil {
return nil, err
}
defer rows.Close()
return scanLogRowsSqlite(rows)
}
func (s *TriggerStore) PruneLogs(ctx context.Context, before time.Time) (int64, error) {
result, err := DB.ExecContext(ctx,
`DELETE FROM trigger_logs WHERE fired_at < ?`, before.UTC().Format(time.RFC3339))
if err != nil {
return 0, err
}
return result.RowsAffected()
}
// ── Helpers ──────────────────────────────────
func scanTriggerRowsSqlite(rows *sql.Rows) ([]models.Trigger, error) {
var triggers []models.Trigger
for rows.Next() {
var t models.Trigger
var cfgStr string
if err := rows.Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug,
&t.EventPattern, &t.EntryPoint, &cfgStr, &t.FireCount,
database.SNT(&t.LastFireAt), &t.LastError,
database.ST(&t.CreatedAt), database.ST(&t.UpdatedAt)); err != nil {
return nil, err
}
t.Config = json.RawMessage(cfgStr)
triggers = append(triggers, t)
}
if triggers == nil {
triggers = []models.Trigger{}
}
return triggers, nil
}
func scanLogRowsSqlite(rows *sql.Rows) ([]models.TriggerLog, error) {
var logs []models.TriggerLog
for rows.Next() {
var l models.TriggerLog
if err := rows.Scan(&l.ID, &l.TriggerID, &l.ScheduledTaskID,
&l.FiredAt, &l.DurationMs, &l.Success, &l.Error, &l.Output); err != nil {
return nil, err
}
logs = append(logs, l)
}
if logs == nil {
logs = []models.TriggerLog{}
}
return logs, nil
}
// boolToInt and nullIfEmpty are defined in workflows.go (same package).

View File

@@ -0,0 +1,45 @@
package store
import (
"context"
"time"
"switchboard-core/models"
)
// TriggerStore manages extension-declared event and webhook triggers.
type TriggerStore interface {
// CRUD
Create(ctx context.Context, t *models.Trigger) error
GetByID(ctx context.Context, id string) (*models.Trigger, error)
Update(ctx context.Context, t *models.Trigger) error
Delete(ctx context.Context, id string) error
// Queries
List(ctx context.Context, opts TriggerListOptions) ([]models.Trigger, int, error)
ListByPackage(ctx context.Context, packageID string) ([]models.Trigger, error)
ListEnabledByType(ctx context.Context, triggerType string) ([]models.Trigger, error)
// Lifecycle
SetEnabled(ctx context.Context, id string, enabled bool) error
UpdateFireState(ctx context.Context, id string, lastFireAt time.Time, lastError string) error
// Webhook resolution
GetWebhook(ctx context.Context, packageID, slug string) (*models.Trigger, error)
// Cleanup
DeleteForPackage(ctx context.Context, packageID string) error
// Logs
LogExecution(ctx context.Context, log *models.TriggerLog) error
ListLogs(ctx context.Context, triggerID string, limit int) ([]models.TriggerLog, error)
PruneLogs(ctx context.Context, before time.Time) (int64, error)
}
// TriggerListOptions provides filtering for trigger list queries.
type TriggerListOptions struct {
ListOptions
PackageID string
Type string
Enabled *bool
}

222
server/triggers/engine.go Normal file
View File

@@ -0,0 +1,222 @@
// 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] + "…"
}

163
server/triggers/event.go Normal file
View File

@@ -0,0 +1,163 @@
package triggers
import (
"context"
"encoding/json"
"fmt"
"log"
"time"
"go.starlark.net/starlark"
"switchboard-core/events"
"switchboard-core/models"
"switchboard-core/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"), mapToStarlark(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
}
// mapToStarlark converts a map[string]any to a Starlark dict.
func mapToStarlark(m map[string]any) *starlark.Dict {
d := starlark.NewDict(len(m))
for k, v := range m {
_ = d.SetKey(starlark.String(k), goToStarlark(v))
}
return d
}
// goToStarlark converts a Go value to a Starlark value.
func goToStarlark(v any) starlark.Value {
switch val := v.(type) {
case nil:
return starlark.None
case bool:
return starlark.Bool(val)
case int:
return starlark.MakeInt(val)
case int64:
return starlark.MakeInt64(val)
case float64:
return starlark.Float(val)
case string:
return starlark.String(val)
case map[string]any:
return mapToStarlark(val)
case []any:
elems := make([]starlark.Value, len(val))
for i, e := range val {
elems[i] = goToStarlark(e)
}
return starlark.NewList(elems)
default:
return starlark.String(fmt.Sprintf("%v", val))
}
}
// 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,
}
}

25
server/triggers/global.go Normal file
View File

@@ -0,0 +1,25 @@
package triggers
import "sync"
// globalEngine holds a reference to the active trigger engine.
// Set by SetGlobalEngine at startup, used by handlers that need
// to wire triggers immediately on package install.
var (
globalMu sync.RWMutex
globalEngine *Engine
)
// SetGlobalEngine stores the engine reference for use by handlers.
func SetGlobalEngine(e *Engine) {
globalMu.Lock()
globalEngine = e
globalMu.Unlock()
}
// GlobalEngine returns the current engine, or nil.
func GlobalEngine() *Engine {
globalMu.RLock()
defer globalMu.RUnlock()
return globalEngine
}

162
server/triggers/schedule.go Normal file
View File

@@ -0,0 +1,162 @@
package triggers
import (
"context"
"log"
"time"
"github.com/robfig/cron/v3"
"go.starlark.net/starlark"
starlarkjson "go.starlark.net/lib/json"
"switchboard-core/models"
"switchboard-core/sandbox"
)
// wireScheduledTask registers a cron entry for a scheduled task.
func (e *Engine) wireScheduledTask(t *models.ScheduledTask) {
taskID := t.ID
creatorID := t.CreatorID
runAs := t.RunAs
script := t.Script
cronExpr := t.CronExpr
// Parse cron with standard 5-field format (minute hour dom month dow)
parser := cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow)
schedule, err := parser.Parse(cronExpr)
if err != nil {
log.Printf(" ⚠️ schedule[%s]: invalid cron %q: %v", taskID, cronExpr, err)
return
}
// Compute next fire time
now := time.Now()
nextFire := schedule.Next(now)
_ = e.stores.ScheduledTasks.UpdateFireState(context.Background(), taskID,
time.Time{}, &nextFire, "", 0)
entryID := e.cron.Schedule(schedule, cron.FuncJob(func() {
e.fireScheduledTask(taskID, creatorID, runAs, script, cronExpr)
}))
e.mu.Lock()
e.cronIDs[taskID] = entryID
e.mu.Unlock()
}
// fireScheduledTask invokes a scheduled task's Starlark script in a restricted sandbox.
func (e *Engine) fireScheduledTask(taskID, creatorID, runAs, script, cronExpr string) {
start := time.Now()
ctx := e.ctx
if ctx == nil {
ctx = context.Background()
}
// Re-check task is still enabled
task, err := e.stores.ScheduledTasks.GetByID(ctx, taskID)
if err != nil || task == nil || !task.Enabled {
return
}
// Check creator is still active (unless run_as=system)
if runAs == models.RunAsCreator {
user, err := e.stores.Users.GetByID(ctx, creatorID)
if err != nil || user == nil || !user.IsActive {
log.Printf(" ⚠️ schedule[%s]: creator %s inactive, pausing", taskID, creatorID)
_ = e.stores.ScheduledTasks.SetEnabled(ctx, taskID, false)
return
}
}
// Build restricted module set
modules := e.buildRestrictedModules(ctx, creatorID, runAs)
// Execute script in sandbox
sb := sandbox.New(sandbox.DefaultConfig())
result, execErr := sb.ExecWithLoader(ctx, "schedule/"+taskID+".star", script, modules, nil)
var output string
if result != nil {
output = result.Output
}
// Call main() if defined
if execErr == nil && result != nil {
if mainFn, ok := result.Globals["main"]; ok {
if callable, ok := mainFn.(starlark.Callable); ok {
// Build context dict
ctxDict := starlark.NewDict(4)
_ = ctxDict.SetKey(starlark.String("trigger_type"), starlark.String("schedule"))
_ = ctxDict.SetKey(starlark.String("task_id"), starlark.String(taskID))
_ = ctxDict.SetKey(starlark.String("cron_expr"), starlark.String(cronExpr))
_ = ctxDict.SetKey(starlark.String("fired_at"), starlark.String(start.UTC().Format(time.RFC3339)))
_, callOutput, callErr := sb.Call(ctx, callable, starlark.Tuple{ctxDict}, nil)
output += callOutput
if callErr != nil {
execErr = callErr
}
}
}
}
duration := int(time.Since(start).Milliseconds())
errStr := ""
if execErr != nil {
errStr = execErr.Error()
log.Printf(" ⚠️ schedule[%s] error: %v", taskID, execErr)
}
// Compute next fire time
parser := cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow)
if sched, err := parser.Parse(cronExpr); err == nil {
nextFire := sched.Next(time.Now())
_ = e.stores.ScheduledTasks.UpdateFireState(ctx, taskID, start, &nextFire, errStr, duration)
} else {
_ = e.stores.ScheduledTasks.UpdateFireState(ctx, taskID, start, nil, errStr, duration)
}
// Log execution
e.logExecution("", taskID, start, duration, execErr == nil, errStr, output)
e.publishEvent("trigger.fired", taskID)
}
// buildRestrictedModules creates the limited module set for scheduled tasks.
// Scheduled scripts can: read settings, load libraries, use JSON, resolve connections.
// They CANNOT: create DB tables, make raw HTTP calls, read secrets.
func (e *Engine) buildRestrictedModules(ctx context.Context, creatorID, runAs string) map[string]starlark.Value {
modules := make(map[string]starlark.Value)
// JSON is always available
modules["json"] = starlarkjson.Module
// Settings module — read-only access to platform settings
// (no package context for scheduled tasks; they read global settings)
if e.stores.Packages != nil {
modules["settings"] = sandbox.BuildSettingsModule(ctx, e.stores, "", creatorID, "")
}
// Notifications module — if available
if e.runner != nil {
// The runner has the notifier attached; we need to access it.
// For now, scheduled tasks don't get the notifications module
// since it requires a package context. This can be enhanced later.
}
// Connections module — read-only resolve for HTTP via existing connections
// This allows scheduled tasks to make HTTP calls through configured connections.
// The RunContext carries the creator's user ID for connection resolution.
return modules
}
// RunContext builds a sandbox.RunContext for a scheduled task invocation.
func scheduleRunContext(creatorID, runAs string) *sandbox.RunContext {
if runAs == models.RunAsSystem {
return nil // system context — no user scoping
}
return &sandbox.RunContext{
UserID: creatorID,
}
}

215
server/triggers/webhook.go Normal file
View File

@@ -0,0 +1,215 @@
package triggers
import (
"crypto/hmac"
"crypto/sha256"
"encoding/hex"
"io"
"log"
"net/http"
"time"
"github.com/gin-gonic/gin"
"go.starlark.net/starlark"
"switchboard-core/models"
)
// HandleWebhook is the Gin handler for inbound webhook triggers.
// Route: POST/GET /api/v1/hooks/:package_id/:slug
func (e *Engine) HandleWebhook(c *gin.Context) {
packageID := c.Param("package_id")
slug := c.Param("slug")
if e.stores.Triggers == nil {
c.JSON(http.StatusNotFound, gin.H{"error": "triggers not available"})
return
}
// Resolve trigger
trigger, err := e.stores.Triggers.GetWebhook(c.Request.Context(), packageID, slug)
if err != nil {
log.Printf(" ⚠️ webhook: lookup error %s/%s: %v", packageID, slug, err)
c.JSON(http.StatusInternalServerError, gin.H{"error": "internal error"})
return
}
if trigger == nil {
c.JSON(http.StatusNotFound, gin.H{"error": "webhook not found"})
return
}
if !trigger.Enabled {
c.JSON(http.StatusNotFound, gin.H{"error": "webhook disabled"})
return
}
// Read body
body, err := io.ReadAll(io.LimitReader(c.Request.Body, 1<<20)) // 1MB limit
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": "failed to read body"})
return
}
// Verify HMAC signature if secret is set
if trigger.Secret != "" {
sig := c.GetHeader("X-Switchboard-Signature")
if !verifyHMAC(body, trigger.Secret, sig) {
c.JSON(http.StatusUnauthorized, gin.H{"error": "invalid signature"})
return
}
}
// Load package
pkg, err := e.stores.Packages.Get(c.Request.Context(), packageID)
if err != nil || pkg == nil || pkg.Status != models.PackageStatusActive {
c.JSON(http.StatusNotFound, gin.H{"error": "package not available"})
return
}
// Check permission
if !e.hasPermission(c.Request.Context(), packageID, models.ExtPermTriggersRegister) {
c.JSON(http.StatusForbidden, gin.H{"error": "triggers.register permission not granted"})
return
}
// Build Starlark context dict
ctxDict := starlark.NewDict(8)
_ = ctxDict.SetKey(starlark.String("trigger_type"), starlark.String("webhook"))
_ = ctxDict.SetKey(starlark.String("trigger_id"), starlark.String(trigger.ID))
_ = ctxDict.SetKey(starlark.String("method"), starlark.String(c.Request.Method))
_ = ctxDict.SetKey(starlark.String("path"), starlark.String(c.Request.URL.Path))
_ = ctxDict.SetKey(starlark.String("body"), starlark.String(string(body)))
// Headers as dict
headers := starlark.NewDict(len(c.Request.Header))
for k, v := range c.Request.Header {
if len(v) > 0 {
_ = headers.SetKey(starlark.String(k), starlark.String(v[0]))
}
}
_ = ctxDict.SetKey(starlark.String("headers"), headers)
// Query params as dict
query := starlark.NewDict(len(c.Request.URL.Query()))
for k, v := range c.Request.URL.Query() {
if len(v) > 0 {
_ = query.SetKey(starlark.String(k), starlark.String(v[0]))
}
}
_ = ctxDict.SetKey(starlark.String("query"), query)
// Fire synchronously — return result as HTTP response
start := time.Now()
val, output, callErr := e.runner.CallEntryPoint(c.Request.Context(), pkg, trigger.EntryPoint,
starlark.Tuple{ctxDict}, nil, nil)
duration := int(time.Since(start).Milliseconds())
errStr := ""
if callErr != nil {
errStr = callErr.Error()
log.Printf(" ⚠️ trigger[webhook] %s/%s error: %v", packageID, slug, callErr)
}
// Update fire state
_ = e.stores.Triggers.UpdateFireState(c.Request.Context(), trigger.ID, start, errStr)
// Log execution
e.logExecution(trigger.ID, "", start, duration, callErr == nil, errStr, output)
e.publishEvent("trigger.fired", trigger.ID)
if callErr != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": "handler error", "detail": errStr})
return
}
// Return Starlark result
response := parseWebhookResponse(val)
c.JSON(response.status, response.body)
}
// webhookResponse is the parsed response from a Starlark webhook handler.
type webhookResponse struct {
status int
body any
}
// parseWebhookResponse extracts HTTP status and body from the Starlark return value.
// Returns {status: X, body: Y} or default 200 OK.
func parseWebhookResponse(val starlark.Value) webhookResponse {
if val == nil || val == starlark.None {
return webhookResponse{status: 200, body: gin.H{"ok": true}}
}
// If it's a dict with "status" and/or "body" keys, use them
if d, ok := val.(*starlark.Dict); ok {
resp := webhookResponse{status: 200}
if statusVal, found, _ := d.Get(starlark.String("status")); found {
if s, ok := statusVal.(starlark.Int); ok {
if i, ok := s.Int64(); ok {
resp.status = int(i)
}
}
}
if bodyVal, found, _ := d.Get(starlark.String("body")); found {
resp.body = starlarkToGo(bodyVal)
} else {
resp.body = gin.H{"ok": true}
}
return resp
}
// String return → plain text body
if s, ok := val.(starlark.String); ok {
return webhookResponse{status: 200, body: gin.H{"result": string(s)}}
}
return webhookResponse{status: 200, body: gin.H{"ok": true}}
}
// starlarkToGo converts a Starlark value to a Go value (for JSON serialization).
func starlarkToGo(v starlark.Value) any {
switch val := v.(type) {
case starlark.NoneType:
return nil
case starlark.Bool:
return bool(val)
case starlark.Int:
if i, ok := val.Int64(); ok {
return i
}
return val.String()
case starlark.Float:
return float64(val)
case starlark.String:
return string(val)
case *starlark.List:
result := make([]any, val.Len())
for i := 0; i < val.Len(); i++ {
result[i] = starlarkToGo(val.Index(i))
}
return result
case *starlark.Dict:
result := make(map[string]any, val.Len())
for _, item := range val.Items() {
if k, ok := item[0].(starlark.String); ok {
result[string(k)] = starlarkToGo(item[1])
}
}
return result
default:
return v.String()
}
}
// verifyHMAC checks the HMAC-SHA256 signature of a webhook payload.
func verifyHMAC(body []byte, secret, signature string) bool {
if signature == "" {
return false
}
mac := hmac.New(sha256.New, []byte(secret))
mac.Write(body)
expected := hex.EncodeToString(mac.Sum(nil))
return hmac.Equal([]byte(expected), []byte(signature))
}