From bd703b9e0dbc0386cef6006d02916095ae52db48 Mon Sep 17 00:00:00 2001 From: Jeffrey Smith Date: Thu, 26 Mar 2026 22:29:43 +0000 Subject: [PATCH] Feat event bus subscriptions + trigger system (v0.2.2) Three trigger primitives replacing the old monolithic scheduler: - Event triggers: extensions subscribe to bus patterns via manifest, async handler invocation through sandbox.CallEntryPoint - Webhook triggers: inbound HTTP at /api/v1/hooks/:pkg/:slug with HMAC-SHA256 verification and synchronous Starlark response - Scheduled tasks: user-created cron scripts with restricted sandbox (no raw HTTP, no DB table creation), runs as creator identity New tables: triggers, scheduled_tasks, trigger_logs (postgres + sqlite). New permission: triggers.register. Full admin + user CRUD APIs. SyncManifestTriggers hooked into seed and install flows. Co-Authored-By: Claude Opus 4.6 (1M context) --- CHANGELOG.md | 38 +- ROADMAP.md | 8 +- .../migrations/postgres/011_triggers.sql | 85 ++++ .../migrations/sqlite/011_triggers.sql | 80 ++++ server/events/types.go | 4 + server/handlers/extensions.go | 4 + server/handlers/packages.go | 4 + server/handlers/schedules.go | 309 ++++++++++++++ server/handlers/seed_packages.go | 5 + server/handlers/trigger_sync.go | 192 +++++++++ server/handlers/triggers.go | 135 ++++++ server/main.go | 40 ++ server/models/models_extension_perm.go | 2 + server/models/trigger.go | 77 ++++ server/static/openapi.yaml | 398 ++++++++++++++++++ server/store/interfaces.go | 2 + server/store/postgres/scheduled_tasks.go | 216 ++++++++++ server/store/postgres/stores.go | 2 + server/store/postgres/triggers.go | 287 +++++++++++++ server/store/scheduled_task_iface.go | 33 ++ server/store/sqlite/scheduled_tasks.go | 216 ++++++++++ server/store/sqlite/stores.go | 2 + server/store/sqlite/triggers.go | 269 ++++++++++++ server/store/trigger_iface.go | 45 ++ server/triggers/engine.go | 222 ++++++++++ server/triggers/event.go | 163 +++++++ server/triggers/global.go | 25 ++ server/triggers/schedule.go | 162 +++++++ server/triggers/webhook.go | 215 ++++++++++ 29 files changed, 3237 insertions(+), 3 deletions(-) create mode 100644 server/database/migrations/postgres/011_triggers.sql create mode 100644 server/database/migrations/sqlite/011_triggers.sql create mode 100644 server/handlers/schedules.go create mode 100644 server/handlers/trigger_sync.go create mode 100644 server/handlers/triggers.go create mode 100644 server/models/trigger.go create mode 100644 server/store/postgres/scheduled_tasks.go create mode 100644 server/store/postgres/triggers.go create mode 100644 server/store/scheduled_task_iface.go create mode 100644 server/store/sqlite/scheduled_tasks.go create mode 100644 server/store/sqlite/triggers.go create mode 100644 server/store/trigger_iface.go create mode 100644 server/triggers/engine.go create mode 100644 server/triggers/event.go create mode 100644 server/triggers/global.go create mode 100644 server/triggers/schedule.go create mode 100644 server/triggers/webhook.go diff --git a/CHANGELOG.md b/CHANGELOG.md index 0523717..fc51e19 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/ROADMAP.md b/ROADMAP.md index dd1ac6c..f834725 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -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. | diff --git a/server/database/migrations/postgres/011_triggers.sql b/server/database/migrations/postgres/011_triggers.sql new file mode 100644 index 0000000..275d791 --- /dev/null +++ b/server/database/migrations/postgres/011_triggers.sql @@ -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); diff --git a/server/database/migrations/sqlite/011_triggers.sql b/server/database/migrations/sqlite/011_triggers.sql new file mode 100644 index 0000000..df2f09f --- /dev/null +++ b/server/database/migrations/sqlite/011_triggers.sql @@ -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); diff --git a/server/events/types.go b/server/events/types.go index 92d4119..3f43852 100644 --- a/server/events/types.go +++ b/server/events/types.go @@ -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, diff --git a/server/handlers/extensions.go b/server/handlers/extensions.go index 3c32dce..9145bcf 100644 --- a/server/handlers/extensions.go +++ b/server/handlers/extensions.go @@ -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}) } diff --git a/server/handlers/packages.go b/server/handlers/packages.go index d3b29a1..e861b81 100644 --- a/server/handlers/packages.go +++ b/server/handlers/packages.go @@ -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 { diff --git a/server/handlers/schedules.go b/server/handlers/schedules.go new file mode 100644 index 0000000..e049c34 --- /dev/null +++ b/server/handlers/schedules.go @@ -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}) +} diff --git a/server/handlers/seed_packages.go b/server/handlers/seed_packages.go index e9cf08e..1508d64 100644 --- a/server/handlers/seed_packages.go +++ b/server/handlers/seed_packages.go @@ -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++ } diff --git a/server/handlers/trigger_sync.go b/server/handlers/trigger_sync.go new file mode 100644 index 0000000..b9cd84d --- /dev/null +++ b/server/handlers/trigger_sync.go @@ -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) +} diff --git a/server/handlers/triggers.go b/server/handlers/triggers.go new file mode 100644 index 0000000..0cc0f37 --- /dev/null +++ b/server/handlers/triggers.go @@ -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 +} diff --git a/server/main.go b/server/main.go index 1faef1a..651513d 100644 --- a/server/main.go +++ b/server/main.go @@ -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) diff --git a/server/models/models_extension_perm.go b/server/models/models_extension_perm.go index fc6060d..036d88a 100644 --- a/server/models/models_extension_perm.go +++ b/server/models/models_extension_perm.go @@ -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 ─────────────── diff --git a/server/models/trigger.go b/server/models/trigger.go new file mode 100644 index 0000000..fdfd733 --- /dev/null +++ b/server/models/trigger.go @@ -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"` +} diff --git a/server/static/openapi.yaml b/server/static/openapi.yaml index 450da04..865f199 100644 --- a/server/static/openapi.yaml +++ b/server/static/openapi.yaml @@ -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 diff --git a/server/store/interfaces.go b/server/store/interfaces.go index 77e037e..b537046 100644 --- a/server/store/interfaces.go +++ b/server/store/interfaces.go @@ -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. diff --git a/server/store/postgres/scheduled_tasks.go b/server/store/postgres/scheduled_tasks.go new file mode 100644 index 0000000..4b22e7f --- /dev/null +++ b/server/store/postgres/scheduled_tasks.go @@ -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, ¶ms, &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, ¶ms, &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, ¶ms, &t.FireCount, + &t.LastError, &t.LastDurationMs, &t.CreatedAt, &t.UpdatedAt); err != nil { + return nil + } + t.TemplateParams = json.RawMessage(params) + return &t +} diff --git a/server/store/postgres/stores.go b/server/store/postgres/stores.go index 59d553a..b8492d9 100644 --- a/server/store/postgres/stores.go +++ b/server/store/postgres/stores.go @@ -30,5 +30,7 @@ func NewStores(db *sql.DB) store.Stores { ExtData: NewExtDataStore(db), Tickets: NewTicketStore(), RateLimits: NewRateLimitStore(), + Triggers: NewTriggerStore(), + ScheduledTasks: NewScheduledTaskStore(), } } diff --git a/server/store/postgres/triggers.go b/server/store/postgres/triggers.go new file mode 100644 index 0000000..e3dcd11 --- /dev/null +++ b/server/store/postgres/triggers.go @@ -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 +} diff --git a/server/store/scheduled_task_iface.go b/server/store/scheduled_task_iface.go new file mode 100644 index 0000000..9040861 --- /dev/null +++ b/server/store/scheduled_task_iface.go @@ -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 +} diff --git a/server/store/sqlite/scheduled_tasks.go b/server/store/sqlite/scheduled_tasks.go new file mode 100644 index 0000000..f8f405e --- /dev/null +++ b/server/store/sqlite/scheduled_tasks.go @@ -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, ¶msStr, &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, ¶msStr, &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 +} diff --git a/server/store/sqlite/stores.go b/server/store/sqlite/stores.go index 630dbd2..3ad1a88 100644 --- a/server/store/sqlite/stores.go +++ b/server/store/sqlite/stores.go @@ -30,5 +30,7 @@ func NewStores(db *sql.DB) store.Stores { ExtData: NewExtDataStore(), Tickets: NewTicketStore(), RateLimits: NewRateLimitStore(), + Triggers: NewTriggerStore(), + ScheduledTasks: NewScheduledTaskStore(), } } diff --git a/server/store/sqlite/triggers.go b/server/store/sqlite/triggers.go new file mode 100644 index 0000000..9ddfece --- /dev/null +++ b/server/store/sqlite/triggers.go @@ -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). diff --git a/server/store/trigger_iface.go b/server/store/trigger_iface.go new file mode 100644 index 0000000..0d68a62 --- /dev/null +++ b/server/store/trigger_iface.go @@ -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 +} diff --git a/server/triggers/engine.go b/server/triggers/engine.go new file mode 100644 index 0000000..95367e6 --- /dev/null +++ b/server/triggers/engine.go @@ -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] + "…" +} diff --git a/server/triggers/event.go b/server/triggers/event.go new file mode 100644 index 0000000..37989be --- /dev/null +++ b/server/triggers/event.go @@ -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, + } +} diff --git a/server/triggers/global.go b/server/triggers/global.go new file mode 100644 index 0000000..23216e7 --- /dev/null +++ b/server/triggers/global.go @@ -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 +} diff --git a/server/triggers/schedule.go b/server/triggers/schedule.go new file mode 100644 index 0000000..ed71e64 --- /dev/null +++ b/server/triggers/schedule.go @@ -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, + } +} diff --git a/server/triggers/webhook.go b/server/triggers/webhook.go new file mode 100644 index 0000000..966704b --- /dev/null +++ b/server/triggers/webhook.go @@ -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)) +} -- 2.49.1