Feat triggers v0.2.2 (#6)
All checks were successful
All checks were successful
Co-authored-by: Jeffrey Smith <jasafpro@gmail.com> Co-committed-by: Jeffrey Smith <jasafpro@gmail.com>
This commit was merged in pull request #6.
This commit is contained in:
85
server/database/migrations/postgres/011_triggers.sql
Normal file
85
server/database/migrations/postgres/011_triggers.sql
Normal file
@@ -0,0 +1,85 @@
|
||||
-- ==========================================
|
||||
-- Switchboard Core — 011 Triggers & Scheduled Tasks
|
||||
-- ==========================================
|
||||
|
||||
-- ── Extension Triggers (event + webhook) ──────
|
||||
|
||||
CREATE TABLE IF NOT EXISTS triggers (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
package_id TEXT NOT NULL REFERENCES packages(id) ON DELETE CASCADE,
|
||||
type TEXT NOT NULL CHECK (type IN ('webhook', 'event')),
|
||||
enabled BOOLEAN NOT NULL DEFAULT true,
|
||||
|
||||
-- Webhook triggers
|
||||
slug TEXT,
|
||||
secret TEXT,
|
||||
|
||||
-- Event triggers
|
||||
event_pattern TEXT,
|
||||
|
||||
-- Common
|
||||
entry_point TEXT NOT NULL,
|
||||
config JSONB NOT NULL DEFAULT '{}',
|
||||
fire_count INTEGER NOT NULL DEFAULT 0,
|
||||
last_fire_at TIMESTAMPTZ,
|
||||
last_error TEXT,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
|
||||
UNIQUE(package_id, type, slug),
|
||||
UNIQUE(package_id, type, event_pattern)
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_triggers_package ON triggers(package_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_triggers_type ON triggers(type);
|
||||
CREATE INDEX IF NOT EXISTS idx_triggers_enabled ON triggers(enabled) WHERE enabled = true;
|
||||
CREATE INDEX IF NOT EXISTS idx_triggers_webhook ON triggers(package_id, slug) WHERE type = 'webhook';
|
||||
|
||||
-- ── Scheduled Tasks (user-created cron) ───────
|
||||
|
||||
CREATE TABLE IF NOT EXISTS scheduled_tasks (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
name TEXT NOT NULL,
|
||||
description TEXT NOT NULL DEFAULT '',
|
||||
creator_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
||||
run_as TEXT NOT NULL DEFAULT 'creator'
|
||||
CHECK (run_as IN ('creator', 'system')),
|
||||
cron_expr TEXT NOT NULL,
|
||||
next_fire_at TIMESTAMPTZ,
|
||||
last_fire_at TIMESTAMPTZ,
|
||||
enabled BOOLEAN NOT NULL DEFAULT true,
|
||||
script TEXT NOT NULL DEFAULT '',
|
||||
template_id TEXT,
|
||||
template_params JSONB NOT NULL DEFAULT '{}',
|
||||
fire_count INTEGER NOT NULL DEFAULT 0,
|
||||
last_error TEXT,
|
||||
last_duration_ms INTEGER,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_sched_tasks_creator ON scheduled_tasks(creator_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_sched_tasks_enabled ON scheduled_tasks(enabled) WHERE enabled = true;
|
||||
CREATE INDEX IF NOT EXISTS idx_sched_tasks_next_fire ON scheduled_tasks(next_fire_at) WHERE enabled = true;
|
||||
|
||||
-- ── Trigger Logs (unified execution log) ──────
|
||||
|
||||
CREATE TABLE IF NOT EXISTS trigger_logs (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
trigger_id UUID REFERENCES triggers(id) ON DELETE CASCADE,
|
||||
scheduled_task_id UUID REFERENCES scheduled_tasks(id) ON DELETE CASCADE,
|
||||
fired_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
duration_ms INTEGER,
|
||||
success BOOLEAN NOT NULL DEFAULT true,
|
||||
error TEXT,
|
||||
output TEXT,
|
||||
|
||||
CHECK (
|
||||
(trigger_id IS NOT NULL AND scheduled_task_id IS NULL) OR
|
||||
(trigger_id IS NULL AND scheduled_task_id IS NOT NULL)
|
||||
)
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_trigger_logs_trigger ON trigger_logs(trigger_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_trigger_logs_sched ON trigger_logs(scheduled_task_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_trigger_logs_fired ON trigger_logs(fired_at);
|
||||
80
server/database/migrations/sqlite/011_triggers.sql
Normal file
80
server/database/migrations/sqlite/011_triggers.sql
Normal file
@@ -0,0 +1,80 @@
|
||||
-- ==========================================
|
||||
-- Switchboard Core — 011 Triggers & Scheduled Tasks (SQLite)
|
||||
-- ==========================================
|
||||
|
||||
-- ── Extension Triggers (event + webhook) ──────
|
||||
|
||||
CREATE TABLE IF NOT EXISTS triggers (
|
||||
id TEXT PRIMARY KEY,
|
||||
package_id TEXT NOT NULL REFERENCES packages(id) ON DELETE CASCADE,
|
||||
type TEXT NOT NULL CHECK (type IN ('webhook', 'event')),
|
||||
enabled INTEGER NOT NULL DEFAULT 1,
|
||||
|
||||
-- Webhook triggers
|
||||
slug TEXT,
|
||||
secret TEXT,
|
||||
|
||||
-- Event triggers
|
||||
event_pattern TEXT,
|
||||
|
||||
-- Common
|
||||
entry_point TEXT NOT NULL,
|
||||
config TEXT NOT NULL DEFAULT '{}',
|
||||
fire_count INTEGER NOT NULL DEFAULT 0,
|
||||
last_fire_at TEXT,
|
||||
last_error TEXT,
|
||||
created_at TEXT NOT NULL DEFAULT (datetime('now')),
|
||||
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
|
||||
|
||||
UNIQUE(package_id, type, slug),
|
||||
UNIQUE(package_id, type, event_pattern)
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_triggers_package ON triggers(package_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_triggers_type ON triggers(type);
|
||||
CREATE INDEX IF NOT EXISTS idx_triggers_enabled ON triggers(enabled);
|
||||
CREATE INDEX IF NOT EXISTS idx_triggers_webhook ON triggers(package_id, slug);
|
||||
|
||||
-- ── Scheduled Tasks (user-created cron) ───────
|
||||
|
||||
CREATE TABLE IF NOT EXISTS scheduled_tasks (
|
||||
id TEXT PRIMARY KEY,
|
||||
name TEXT NOT NULL,
|
||||
description TEXT NOT NULL DEFAULT '',
|
||||
creator_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
||||
run_as TEXT NOT NULL DEFAULT 'creator'
|
||||
CHECK (run_as IN ('creator', 'system')),
|
||||
cron_expr TEXT NOT NULL,
|
||||
next_fire_at TEXT,
|
||||
last_fire_at TEXT,
|
||||
enabled INTEGER NOT NULL DEFAULT 1,
|
||||
script TEXT NOT NULL DEFAULT '',
|
||||
template_id TEXT,
|
||||
template_params TEXT NOT NULL DEFAULT '{}',
|
||||
fire_count INTEGER NOT NULL DEFAULT 0,
|
||||
last_error TEXT,
|
||||
last_duration_ms INTEGER,
|
||||
created_at TEXT NOT NULL DEFAULT (datetime('now')),
|
||||
updated_at TEXT NOT NULL DEFAULT (datetime('now'))
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_sched_tasks_creator ON scheduled_tasks(creator_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_sched_tasks_enabled ON scheduled_tasks(enabled);
|
||||
CREATE INDEX IF NOT EXISTS idx_sched_tasks_next_fire ON scheduled_tasks(next_fire_at);
|
||||
|
||||
-- ── Trigger Logs (unified execution log) ──────
|
||||
|
||||
CREATE TABLE IF NOT EXISTS trigger_logs (
|
||||
id TEXT PRIMARY KEY,
|
||||
trigger_id TEXT REFERENCES triggers(id) ON DELETE CASCADE,
|
||||
scheduled_task_id TEXT REFERENCES scheduled_tasks(id) ON DELETE CASCADE,
|
||||
fired_at TEXT NOT NULL DEFAULT (datetime('now')),
|
||||
duration_ms INTEGER,
|
||||
success INTEGER NOT NULL DEFAULT 1,
|
||||
error TEXT,
|
||||
output TEXT
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_trigger_logs_trigger ON trigger_logs(trigger_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_trigger_logs_sched ON trigger_logs(scheduled_task_id);
|
||||
CREATE INDEX IF NOT EXISTS idx_trigger_logs_fired ON trigger_logs(fired_at);
|
||||
@@ -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,
|
||||
|
||||
@@ -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})
|
||||
}
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
309
server/handlers/schedules.go
Normal file
309
server/handlers/schedules.go
Normal file
@@ -0,0 +1,309 @@
|
||||
package handlers
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/robfig/cron/v3"
|
||||
|
||||
"switchboard-core/models"
|
||||
"switchboard-core/store"
|
||||
"switchboard-core/triggers"
|
||||
)
|
||||
|
||||
// ScheduleHandler provides CRUD for user-created scheduled tasks.
|
||||
type ScheduleHandler struct {
|
||||
stores store.Stores
|
||||
engine *triggers.Engine
|
||||
}
|
||||
|
||||
// NewScheduleHandler creates a ScheduleHandler.
|
||||
func NewScheduleHandler(stores store.Stores, engine *triggers.Engine) *ScheduleHandler {
|
||||
return &ScheduleHandler{stores: stores, engine: engine}
|
||||
}
|
||||
|
||||
// ── User Endpoints ───────────────────────────
|
||||
|
||||
// ListMySchedules returns schedules created by the current user.
|
||||
// GET /api/v1/schedules
|
||||
func (h *ScheduleHandler) ListMySchedules(c *gin.Context) {
|
||||
userID := c.GetString("user_id")
|
||||
tasks, err := h.stores.ScheduledTasks.ListByCreator(c.Request.Context(), userID)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"schedules": tasks})
|
||||
}
|
||||
|
||||
// CreateSchedule creates a new scheduled task.
|
||||
// POST /api/v1/schedules
|
||||
func (h *ScheduleHandler) CreateSchedule(c *gin.Context) {
|
||||
var req struct {
|
||||
Name string `json:"name" binding:"required"`
|
||||
Description string `json:"description"`
|
||||
CronExpr string `json:"cron_expr" binding:"required"`
|
||||
Script string `json:"script"`
|
||||
TemplateID string `json:"template_id"`
|
||||
TemplateParams any `json:"template_params"`
|
||||
}
|
||||
if err := c.ShouldBindJSON(&req); err != nil {
|
||||
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
// Validate cron expression
|
||||
parser := cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow)
|
||||
schedule, err := parser.Parse(req.CronExpr)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusBadRequest, gin.H{"error": "invalid cron expression: " + err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
nextFire := schedule.Next(time.Now())
|
||||
userID := c.GetString("user_id")
|
||||
|
||||
task := &models.ScheduledTask{
|
||||
Name: req.Name,
|
||||
Description: req.Description,
|
||||
CreatorID: userID,
|
||||
RunAs: models.RunAsCreator,
|
||||
CronExpr: req.CronExpr,
|
||||
NextFireAt: &nextFire,
|
||||
Enabled: true,
|
||||
Script: req.Script,
|
||||
TemplateID: req.TemplateID,
|
||||
}
|
||||
|
||||
if err := h.stores.ScheduledTasks.Create(c.Request.Context(), task); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
// Register with engine
|
||||
h.engine.RegisterSchedule(task)
|
||||
|
||||
c.JSON(http.StatusCreated, task)
|
||||
}
|
||||
|
||||
// GetSchedule returns a single scheduled task.
|
||||
// GET /api/v1/schedules/:id
|
||||
func (h *ScheduleHandler) GetSchedule(c *gin.Context) {
|
||||
task, err := h.stores.ScheduledTasks.GetByID(c.Request.Context(), c.Param("id"))
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
if task == nil {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
|
||||
return
|
||||
}
|
||||
|
||||
// Check ownership (non-admin can only see own schedules)
|
||||
userID := c.GetString("user_id")
|
||||
if task.CreatorID != userID {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
|
||||
return
|
||||
}
|
||||
|
||||
c.JSON(http.StatusOK, task)
|
||||
}
|
||||
|
||||
// UpdateSchedule updates a scheduled task.
|
||||
// PUT /api/v1/schedules/:id
|
||||
func (h *ScheduleHandler) UpdateSchedule(c *gin.Context) {
|
||||
task, err := h.stores.ScheduledTasks.GetByID(c.Request.Context(), c.Param("id"))
|
||||
if err != nil || task == nil {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
|
||||
return
|
||||
}
|
||||
|
||||
userID := c.GetString("user_id")
|
||||
if task.CreatorID != userID {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
|
||||
return
|
||||
}
|
||||
|
||||
var req struct {
|
||||
Name *string `json:"name"`
|
||||
Description *string `json:"description"`
|
||||
CronExpr *string `json:"cron_expr"`
|
||||
Script *string `json:"script"`
|
||||
Enabled *bool `json:"enabled"`
|
||||
}
|
||||
if err := c.ShouldBindJSON(&req); err != nil {
|
||||
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
if req.Name != nil {
|
||||
task.Name = *req.Name
|
||||
}
|
||||
if req.Description != nil {
|
||||
task.Description = *req.Description
|
||||
}
|
||||
if req.Script != nil {
|
||||
task.Script = *req.Script
|
||||
}
|
||||
if req.Enabled != nil {
|
||||
task.Enabled = *req.Enabled
|
||||
}
|
||||
if req.CronExpr != nil {
|
||||
parser := cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow)
|
||||
schedule, err := parser.Parse(*req.CronExpr)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusBadRequest, gin.H{"error": "invalid cron expression: " + err.Error()})
|
||||
return
|
||||
}
|
||||
task.CronExpr = *req.CronExpr
|
||||
nextFire := schedule.Next(time.Now())
|
||||
task.NextFireAt = &nextFire
|
||||
}
|
||||
|
||||
if err := h.stores.ScheduledTasks.Update(c.Request.Context(), task); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
// Re-register with engine
|
||||
h.engine.UnregisterSchedule(task.ID)
|
||||
if task.Enabled {
|
||||
h.engine.RegisterSchedule(task)
|
||||
}
|
||||
|
||||
c.JSON(http.StatusOK, task)
|
||||
}
|
||||
|
||||
// DeleteSchedule deletes a scheduled task.
|
||||
// DELETE /api/v1/schedules/:id
|
||||
func (h *ScheduleHandler) DeleteSchedule(c *gin.Context) {
|
||||
task, err := h.stores.ScheduledTasks.GetByID(c.Request.Context(), c.Param("id"))
|
||||
if err != nil || task == nil {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
|
||||
return
|
||||
}
|
||||
|
||||
userID := c.GetString("user_id")
|
||||
if task.CreatorID != userID {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
|
||||
return
|
||||
}
|
||||
|
||||
h.engine.UnregisterSchedule(task.ID)
|
||||
if err := h.stores.ScheduledTasks.Delete(c.Request.Context(), task.ID); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
c.JSON(http.StatusOK, gin.H{"ok": true})
|
||||
}
|
||||
|
||||
// RunSchedule manually triggers a scheduled task (test run).
|
||||
// POST /api/v1/schedules/:id/run
|
||||
func (h *ScheduleHandler) RunSchedule(c *gin.Context) {
|
||||
task, err := h.stores.ScheduledTasks.GetByID(c.Request.Context(), c.Param("id"))
|
||||
if err != nil || task == nil {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
|
||||
return
|
||||
}
|
||||
|
||||
userID := c.GetString("user_id")
|
||||
if task.CreatorID != userID {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
|
||||
return
|
||||
}
|
||||
|
||||
// Fire immediately in a goroutine
|
||||
go h.engine.ManualRun(task)
|
||||
|
||||
c.JSON(http.StatusAccepted, gin.H{"ok": true, "message": "task queued for execution"})
|
||||
}
|
||||
|
||||
// ListScheduleLogs returns execution logs for a scheduled task.
|
||||
// GET /api/v1/schedules/:id/logs
|
||||
func (h *ScheduleHandler) ListScheduleLogs(c *gin.Context) {
|
||||
task, err := h.stores.ScheduledTasks.GetByID(c.Request.Context(), c.Param("id"))
|
||||
if err != nil || task == nil {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
|
||||
return
|
||||
}
|
||||
|
||||
userID := c.GetString("user_id")
|
||||
if task.CreatorID != userID {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "schedule not found"})
|
||||
return
|
||||
}
|
||||
|
||||
limit := parseIntParam(c, "limit", 50)
|
||||
logs, err := h.stores.Triggers.ListLogs(c.Request.Context(), task.ID, limit)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"logs": logs})
|
||||
}
|
||||
|
||||
// ── Admin Endpoints ──────────────────────────
|
||||
|
||||
// AdminListSchedules returns all scheduled tasks (admin view).
|
||||
// GET /admin/schedules
|
||||
func (h *ScheduleHandler) AdminListSchedules(c *gin.Context) {
|
||||
opts := store.ScheduledTaskListOptions{
|
||||
ListOptions: store.ListOptions{
|
||||
Limit: parseIntParam(c, "limit", 50),
|
||||
Offset: parseIntParam(c, "offset", 0),
|
||||
},
|
||||
CreatorID: c.Query("creator_id"),
|
||||
}
|
||||
if e := c.Query("enabled"); e != "" {
|
||||
b := e == "true"
|
||||
opts.Enabled = &b
|
||||
}
|
||||
|
||||
tasks, total, err := h.stores.ScheduledTasks.List(c.Request.Context(), opts)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"schedules": tasks, "total": total})
|
||||
}
|
||||
|
||||
// AdminEnableSchedule enables a scheduled task.
|
||||
// PUT /admin/schedules/:id/enable
|
||||
func (h *ScheduleHandler) AdminEnableSchedule(c *gin.Context) {
|
||||
id := c.Param("id")
|
||||
if err := h.stores.ScheduledTasks.SetEnabled(c.Request.Context(), id, true); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
task, _ := h.stores.ScheduledTasks.GetByID(c.Request.Context(), id)
|
||||
if task != nil {
|
||||
h.engine.RegisterSchedule(task)
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"ok": true})
|
||||
}
|
||||
|
||||
// AdminDisableSchedule disables a scheduled task.
|
||||
// PUT /admin/schedules/:id/disable
|
||||
func (h *ScheduleHandler) AdminDisableSchedule(c *gin.Context) {
|
||||
id := c.Param("id")
|
||||
h.engine.UnregisterSchedule(id)
|
||||
if err := h.stores.ScheduledTasks.SetEnabled(c.Request.Context(), id, false); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"ok": true})
|
||||
}
|
||||
|
||||
// AdminDeleteSchedule deletes a scheduled task.
|
||||
// DELETE /admin/schedules/:id
|
||||
func (h *ScheduleHandler) AdminDeleteSchedule(c *gin.Context) {
|
||||
id := c.Param("id")
|
||||
h.engine.UnregisterSchedule(id)
|
||||
if err := h.stores.ScheduledTasks.Delete(c.Request.Context(), id); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"ok": true})
|
||||
}
|
||||
@@ -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++
|
||||
}
|
||||
|
||||
|
||||
192
server/handlers/trigger_sync.go
Normal file
192
server/handlers/trigger_sync.go
Normal file
@@ -0,0 +1,192 @@
|
||||
package handlers
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"log"
|
||||
|
||||
"switchboard-core/models"
|
||||
"switchboard-core/store"
|
||||
"switchboard-core/triggers"
|
||||
)
|
||||
|
||||
// manifestTrigger is the manifest.json trigger declaration shape.
|
||||
type manifestTrigger struct {
|
||||
Type string `json:"type"` // "event" or "webhook"
|
||||
Pattern string `json:"pattern"` // event pattern (for type=event)
|
||||
Slug string `json:"slug"` // webhook slug (for type=webhook)
|
||||
EntryPoint string `json:"entry_point"` // Starlark function name
|
||||
}
|
||||
|
||||
// SyncManifestTriggers parses the "triggers" array from a package manifest
|
||||
// and syncs them with the trigger store. Removes triggers no longer in the
|
||||
// manifest (declarative sync). For webhooks, generates an HMAC secret on
|
||||
// first create.
|
||||
//
|
||||
// Called from SeedBuiltinPackages, admin install, and package install.
|
||||
func SyncManifestTriggers(ctx context.Context, stores store.Stores, engine *triggers.Engine, pkgID string, manifest map[string]any) {
|
||||
if stores.Triggers == nil {
|
||||
return
|
||||
}
|
||||
|
||||
// Parse triggers from manifest
|
||||
triggersRaw, ok := manifest["triggers"]
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
|
||||
data, err := json.Marshal(triggersRaw)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
var declared []manifestTrigger
|
||||
if err := json.Unmarshal(data, &declared); err != nil {
|
||||
log.Printf(" ⚠️ trigger-sync[%s]: invalid triggers in manifest: %v", pkgID, err)
|
||||
return
|
||||
}
|
||||
|
||||
if len(declared) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
// Get existing triggers for this package
|
||||
existing, err := stores.Triggers.ListByPackage(ctx, pkgID)
|
||||
if err != nil {
|
||||
log.Printf(" ⚠️ trigger-sync[%s]: failed to list existing: %v", pkgID, err)
|
||||
return
|
||||
}
|
||||
|
||||
// Build lookup maps
|
||||
existingByKey := make(map[string]*models.Trigger, len(existing))
|
||||
for i := range existing {
|
||||
key := triggerKey(&existing[i])
|
||||
existingByKey[key] = &existing[i]
|
||||
}
|
||||
|
||||
declaredKeys := make(map[string]bool, len(declared))
|
||||
created, updated := 0, 0
|
||||
|
||||
for _, dt := range declared {
|
||||
if dt.Type != models.TriggerTypeEvent && dt.Type != models.TriggerTypeWebhook {
|
||||
log.Printf(" ⚠️ trigger-sync[%s]: unknown trigger type %q, skipping", pkgID, dt.Type)
|
||||
continue
|
||||
}
|
||||
if dt.EntryPoint == "" {
|
||||
continue
|
||||
}
|
||||
|
||||
key := dt.Type + ":" + triggerIdentifier(dt)
|
||||
declaredKeys[key] = true
|
||||
|
||||
if ex, ok := existingByKey[key]; ok {
|
||||
// Update if entry_point changed
|
||||
if ex.EntryPoint != dt.EntryPoint {
|
||||
ex.EntryPoint = dt.EntryPoint
|
||||
if err := stores.Triggers.Update(ctx, ex); err != nil {
|
||||
log.Printf(" ⚠️ trigger-sync[%s]: update failed: %v", pkgID, err)
|
||||
} else {
|
||||
updated++
|
||||
}
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
// Create new trigger
|
||||
t := &models.Trigger{
|
||||
PackageID: pkgID,
|
||||
Type: dt.Type,
|
||||
Enabled: true,
|
||||
EntryPoint: dt.EntryPoint,
|
||||
Config: json.RawMessage("{}"),
|
||||
}
|
||||
|
||||
switch dt.Type {
|
||||
case models.TriggerTypeEvent:
|
||||
t.EventPattern = dt.Pattern
|
||||
case models.TriggerTypeWebhook:
|
||||
t.Slug = dt.Slug
|
||||
t.Secret = generateSecret()
|
||||
}
|
||||
|
||||
if err := stores.Triggers.Create(ctx, t); err != nil {
|
||||
log.Printf(" ⚠️ trigger-sync[%s]: create failed: %v", pkgID, err)
|
||||
continue
|
||||
}
|
||||
created++
|
||||
|
||||
// Register with engine
|
||||
if engine != nil {
|
||||
// Re-fetch to get the generated ID
|
||||
switch dt.Type {
|
||||
case models.TriggerTypeWebhook:
|
||||
if wh, _ := stores.Triggers.GetWebhook(ctx, pkgID, dt.Slug); wh != nil {
|
||||
engine.RegisterTrigger(wh)
|
||||
}
|
||||
case models.TriggerTypeEvent:
|
||||
// Need to fetch by package to find the new one
|
||||
triggers, _ := stores.Triggers.ListByPackage(ctx, pkgID)
|
||||
for i := range triggers {
|
||||
if triggers[i].EventPattern == dt.Pattern {
|
||||
engine.RegisterTrigger(&triggers[i])
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Remove triggers no longer in manifest
|
||||
removed := 0
|
||||
for key, ex := range existingByKey {
|
||||
if !declaredKeys[key] {
|
||||
if engine != nil {
|
||||
engine.UnregisterTrigger(ex.ID)
|
||||
}
|
||||
if err := stores.Triggers.Delete(ctx, ex.ID); err != nil {
|
||||
log.Printf(" ⚠️ trigger-sync[%s]: delete failed: %v", pkgID, err)
|
||||
} else {
|
||||
removed++
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if created+updated+removed > 0 {
|
||||
log.Printf(" 🔔 trigger-sync[%s]: %d created, %d updated, %d removed",
|
||||
pkgID, created, updated, removed)
|
||||
}
|
||||
}
|
||||
|
||||
// triggerKey returns a unique key for a trigger within a package.
|
||||
func triggerKey(t *models.Trigger) string {
|
||||
return t.Type + ":" + triggerIdentifierFromModel(t)
|
||||
}
|
||||
|
||||
func triggerIdentifier(dt manifestTrigger) string {
|
||||
switch dt.Type {
|
||||
case models.TriggerTypeEvent:
|
||||
return dt.Pattern
|
||||
case models.TriggerTypeWebhook:
|
||||
return dt.Slug
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func triggerIdentifierFromModel(t *models.Trigger) string {
|
||||
switch t.Type {
|
||||
case models.TriggerTypeEvent:
|
||||
return t.EventPattern
|
||||
case models.TriggerTypeWebhook:
|
||||
return t.Slug
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// generateSecret creates a random 32-byte hex-encoded secret for webhook HMAC.
|
||||
func generateSecret() string {
|
||||
b := make([]byte, 32)
|
||||
_, _ = rand.Read(b)
|
||||
return hex.EncodeToString(b)
|
||||
}
|
||||
135
server/handlers/triggers.go
Normal file
135
server/handlers/triggers.go
Normal file
@@ -0,0 +1,135 @@
|
||||
package handlers
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strconv"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
|
||||
"switchboard-core/store"
|
||||
"switchboard-core/triggers"
|
||||
)
|
||||
|
||||
// TriggerHandler provides admin CRUD for extension-declared triggers.
|
||||
type TriggerHandler struct {
|
||||
stores store.Stores
|
||||
engine *triggers.Engine
|
||||
}
|
||||
|
||||
// NewTriggerHandler creates a TriggerHandler.
|
||||
func NewTriggerHandler(stores store.Stores, engine *triggers.Engine) *TriggerHandler {
|
||||
return &TriggerHandler{stores: stores, engine: engine}
|
||||
}
|
||||
|
||||
// ListTriggers returns all triggers with optional filters.
|
||||
// GET /admin/triggers?package_id=&type=&enabled=
|
||||
func (h *TriggerHandler) ListTriggers(c *gin.Context) {
|
||||
opts := store.TriggerListOptions{
|
||||
ListOptions: store.ListOptions{
|
||||
Limit: parseIntParam(c, "limit", 50),
|
||||
Offset: parseIntParam(c, "offset", 0),
|
||||
},
|
||||
PackageID: c.Query("package_id"),
|
||||
Type: c.Query("type"),
|
||||
}
|
||||
if e := c.Query("enabled"); e != "" {
|
||||
b := e == "true"
|
||||
opts.Enabled = &b
|
||||
}
|
||||
|
||||
triggers, total, err := h.stores.Triggers.List(c.Request.Context(), opts)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"triggers": triggers, "total": total})
|
||||
}
|
||||
|
||||
// GetTrigger returns a single trigger by ID.
|
||||
// GET /admin/triggers/:id
|
||||
func (h *TriggerHandler) GetTrigger(c *gin.Context) {
|
||||
t, err := h.stores.Triggers.GetByID(c.Request.Context(), c.Param("id"))
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
if t == nil {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "trigger not found"})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, t)
|
||||
}
|
||||
|
||||
// EnableTrigger enables a trigger.
|
||||
// PUT /admin/triggers/:id/enable
|
||||
func (h *TriggerHandler) EnableTrigger(c *gin.Context) {
|
||||
id := c.Param("id")
|
||||
if err := h.stores.Triggers.SetEnabled(c.Request.Context(), id, true); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
|
||||
// Re-register with engine
|
||||
t, _ := h.stores.Triggers.GetByID(c.Request.Context(), id)
|
||||
if t != nil {
|
||||
h.engine.RegisterTrigger(t)
|
||||
}
|
||||
|
||||
c.JSON(http.StatusOK, gin.H{"ok": true})
|
||||
}
|
||||
|
||||
// DisableTrigger disables a trigger.
|
||||
// PUT /admin/triggers/:id/disable
|
||||
func (h *TriggerHandler) DisableTrigger(c *gin.Context) {
|
||||
id := c.Param("id")
|
||||
if err := h.stores.Triggers.SetEnabled(c.Request.Context(), id, false); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
h.engine.UnregisterTrigger(id)
|
||||
c.JSON(http.StatusOK, gin.H{"ok": true})
|
||||
}
|
||||
|
||||
// DeleteTrigger deletes a trigger.
|
||||
// DELETE /admin/triggers/:id
|
||||
func (h *TriggerHandler) DeleteTrigger(c *gin.Context) {
|
||||
id := c.Param("id")
|
||||
h.engine.UnregisterTrigger(id)
|
||||
if err := h.stores.Triggers.Delete(c.Request.Context(), id); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"ok": true})
|
||||
}
|
||||
|
||||
// ListTriggerLogs returns execution logs for a trigger.
|
||||
// GET /admin/triggers/:id/logs?limit=
|
||||
func (h *TriggerHandler) ListTriggerLogs(c *gin.Context) {
|
||||
limit := parseIntParam(c, "limit", 50)
|
||||
logs, err := h.stores.Triggers.ListLogs(c.Request.Context(), c.Param("id"), limit)
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"logs": logs})
|
||||
}
|
||||
|
||||
// ListPackageTriggers returns triggers for a specific package.
|
||||
// GET /admin/packages/:id/triggers
|
||||
func (h *TriggerHandler) ListPackageTriggers(c *gin.Context) {
|
||||
triggers, err := h.stores.Triggers.ListByPackage(c.Request.Context(), c.Param("id"))
|
||||
if err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"triggers": triggers})
|
||||
}
|
||||
|
||||
func parseIntParam(c *gin.Context, key string, defaultVal int) int {
|
||||
if v := c.Query(key); v != "" {
|
||||
if n, err := strconv.Atoi(v); err == nil && n > 0 {
|
||||
return n
|
||||
}
|
||||
}
|
||||
return defaultVal
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -34,6 +34,7 @@ const (
|
||||
ExtPermFormValidate = "forms.validate" // v0.29.3: form validation hooks
|
||||
ExtPermWorkflowAccess = "workflow.access" // v0.30.2: workflow definition + stage data access
|
||||
ExtPermConnectionsRead = "connections.read" // v0.38.1: extension connection resolution
|
||||
ExtPermTriggersRegister = "triggers.register" // v0.2.2: event/webhook trigger registration
|
||||
)
|
||||
|
||||
// ValidExtensionPermissions is the set of recognized permission keys.
|
||||
@@ -46,6 +47,7 @@ var ValidExtensionPermissions = map[string]bool{
|
||||
ExtPermFormValidate: true,
|
||||
ExtPermWorkflowAccess: true,
|
||||
ExtPermConnectionsRead: true,
|
||||
ExtPermTriggersRegister: true,
|
||||
}
|
||||
|
||||
// ── Extension Permission Model ───────────────
|
||||
|
||||
77
server/models/trigger.go
Normal file
77
server/models/trigger.go
Normal file
@@ -0,0 +1,77 @@
|
||||
package models
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"time"
|
||||
)
|
||||
|
||||
// ── Trigger Type Constants ───────────────────
|
||||
|
||||
const (
|
||||
TriggerTypeWebhook = "webhook"
|
||||
TriggerTypeEvent = "event"
|
||||
)
|
||||
|
||||
// ── Scheduled Task Run-As Constants ──────────
|
||||
|
||||
const (
|
||||
RunAsCreator = "creator"
|
||||
RunAsSystem = "system"
|
||||
)
|
||||
|
||||
// ── Trigger (extension-declared) ─────────────
|
||||
|
||||
// Trigger represents an extension-declared event or webhook trigger.
|
||||
type Trigger struct {
|
||||
ID string `json:"id" db:"id"`
|
||||
PackageID string `json:"package_id" db:"package_id"`
|
||||
Type string `json:"type" db:"type"`
|
||||
Enabled bool `json:"enabled" db:"enabled"`
|
||||
Slug string `json:"slug,omitempty" db:"slug"`
|
||||
Secret string `json:"-" db:"secret"` // never serialized
|
||||
EventPattern string `json:"event_pattern,omitempty" db:"event_pattern"`
|
||||
EntryPoint string `json:"entry_point" db:"entry_point"`
|
||||
Config json.RawMessage `json:"config" db:"config"`
|
||||
FireCount int `json:"fire_count" db:"fire_count"`
|
||||
LastFireAt *time.Time `json:"last_fire_at,omitempty" db:"last_fire_at"`
|
||||
LastError string `json:"last_error,omitempty" db:"last_error"`
|
||||
CreatedAt time.Time `json:"created_at" db:"created_at"`
|
||||
UpdatedAt time.Time `json:"updated_at" db:"updated_at"`
|
||||
}
|
||||
|
||||
// ── Scheduled Task (user-created) ────────────
|
||||
|
||||
// ScheduledTask represents a user-created cron-scheduled Starlark script.
|
||||
type ScheduledTask struct {
|
||||
ID string `json:"id" db:"id"`
|
||||
Name string `json:"name" db:"name"`
|
||||
Description string `json:"description" db:"description"`
|
||||
CreatorID string `json:"creator_id" db:"creator_id"`
|
||||
RunAs string `json:"run_as" db:"run_as"`
|
||||
CronExpr string `json:"cron_expr" db:"cron_expr"`
|
||||
NextFireAt *time.Time `json:"next_fire_at,omitempty" db:"next_fire_at"`
|
||||
LastFireAt *time.Time `json:"last_fire_at,omitempty" db:"last_fire_at"`
|
||||
Enabled bool `json:"enabled" db:"enabled"`
|
||||
Script string `json:"script" db:"script"`
|
||||
TemplateID string `json:"template_id,omitempty" db:"template_id"`
|
||||
TemplateParams json.RawMessage `json:"template_params" db:"template_params"`
|
||||
FireCount int `json:"fire_count" db:"fire_count"`
|
||||
LastError string `json:"last_error,omitempty" db:"last_error"`
|
||||
LastDurationMs *int `json:"last_duration_ms,omitempty" db:"last_duration_ms"`
|
||||
CreatedAt time.Time `json:"created_at" db:"created_at"`
|
||||
UpdatedAt time.Time `json:"updated_at" db:"updated_at"`
|
||||
}
|
||||
|
||||
// ── Trigger Log (unified) ────────────────────
|
||||
|
||||
// TriggerLog records a single execution of a trigger or scheduled task.
|
||||
type TriggerLog struct {
|
||||
ID string `json:"id" db:"id"`
|
||||
TriggerID string `json:"trigger_id,omitempty" db:"trigger_id"`
|
||||
ScheduledTaskID string `json:"scheduled_task_id,omitempty" db:"scheduled_task_id"`
|
||||
FiredAt string `json:"fired_at" db:"fired_at"`
|
||||
DurationMs *int `json:"duration_ms,omitempty" db:"duration_ms"`
|
||||
Success bool `json:"success" db:"success"`
|
||||
Error string `json:"error,omitempty" db:"error"`
|
||||
Output string `json:"output,omitempty" db:"output"`
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
216
server/store/postgres/scheduled_tasks.go
Normal file
216
server/store/postgres/scheduled_tasks.go
Normal file
@@ -0,0 +1,216 @@
|
||||
package postgres
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"switchboard-core/models"
|
||||
"switchboard-core/store"
|
||||
)
|
||||
|
||||
type ScheduledTaskStore struct{}
|
||||
|
||||
func NewScheduledTaskStore() *ScheduledTaskStore {
|
||||
return &ScheduledTaskStore{}
|
||||
}
|
||||
|
||||
func (s *ScheduledTaskStore) Create(ctx context.Context, t *models.ScheduledTask) error {
|
||||
params, _ := json.Marshal(t.TemplateParams)
|
||||
if len(params) == 0 {
|
||||
params = []byte("{}")
|
||||
}
|
||||
return DB.QueryRowContext(ctx,
|
||||
`INSERT INTO scheduled_tasks (name, description, creator_id, run_as, cron_expr,
|
||||
next_fire_at, enabled, script, template_id, template_params)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
|
||||
RETURNING id, created_at, updated_at`,
|
||||
t.Name, t.Description, t.CreatorID, t.RunAs, t.CronExpr,
|
||||
t.NextFireAt, t.Enabled, t.Script, nilIfEmpty(t.TemplateID), params).
|
||||
Scan(&t.ID, &t.CreatedAt, &t.UpdatedAt)
|
||||
}
|
||||
|
||||
func (s *ScheduledTaskStore) GetByID(ctx context.Context, id string) (*models.ScheduledTask, error) {
|
||||
var t models.ScheduledTask
|
||||
var params []byte
|
||||
err := DB.QueryRowContext(ctx,
|
||||
`SELECT id, name, description, creator_id, run_as, cron_expr,
|
||||
next_fire_at, last_fire_at, enabled, script,
|
||||
COALESCE(template_id,''), template_params, fire_count,
|
||||
COALESCE(last_error,''), last_duration_ms, created_at, updated_at
|
||||
FROM scheduled_tasks WHERE id = $1`, id).
|
||||
Scan(&t.ID, &t.Name, &t.Description, &t.CreatorID, &t.RunAs, &t.CronExpr,
|
||||
&t.NextFireAt, &t.LastFireAt, &t.Enabled, &t.Script,
|
||||
&t.TemplateID, ¶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
|
||||
}
|
||||
@@ -30,5 +30,7 @@ func NewStores(db *sql.DB) store.Stores {
|
||||
ExtData: NewExtDataStore(db),
|
||||
Tickets: NewTicketStore(),
|
||||
RateLimits: NewRateLimitStore(),
|
||||
Triggers: NewTriggerStore(),
|
||||
ScheduledTasks: NewScheduledTaskStore(),
|
||||
}
|
||||
}
|
||||
|
||||
287
server/store/postgres/triggers.go
Normal file
287
server/store/postgres/triggers.go
Normal file
@@ -0,0 +1,287 @@
|
||||
package postgres
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"switchboard-core/models"
|
||||
"switchboard-core/store"
|
||||
)
|
||||
|
||||
type TriggerStore struct{}
|
||||
|
||||
func NewTriggerStore() *TriggerStore {
|
||||
return &TriggerStore{}
|
||||
}
|
||||
|
||||
func (s *TriggerStore) Create(ctx context.Context, t *models.Trigger) error {
|
||||
cfg, _ := json.Marshal(t.Config)
|
||||
if len(cfg) == 0 {
|
||||
cfg = []byte("{}")
|
||||
}
|
||||
_, err := DB.ExecContext(ctx,
|
||||
`INSERT INTO triggers (id, package_id, type, enabled, slug, secret, event_pattern,
|
||||
entry_point, config)
|
||||
VALUES (gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8)
|
||||
RETURNING id`,
|
||||
t.PackageID, t.Type, t.Enabled, nilIfEmpty(t.Slug), nilIfEmpty(t.Secret),
|
||||
nilIfEmpty(t.EventPattern), t.EntryPoint, cfg)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) GetByID(ctx context.Context, id string) (*models.Trigger, error) {
|
||||
var t models.Trigger
|
||||
var cfg []byte
|
||||
err := DB.QueryRowContext(ctx,
|
||||
`SELECT id, package_id, type, enabled, COALESCE(slug,''), secret,
|
||||
COALESCE(event_pattern,''), entry_point, config, fire_count,
|
||||
last_fire_at, COALESCE(last_error,''), created_at, updated_at
|
||||
FROM triggers WHERE id = $1`, id).
|
||||
Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug, &t.Secret,
|
||||
&t.EventPattern, &t.EntryPoint, &cfg, &t.FireCount,
|
||||
&t.LastFireAt, &t.LastError, &t.CreatedAt, &t.UpdatedAt)
|
||||
if err == sql.ErrNoRows {
|
||||
return nil, nil
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
t.Config = json.RawMessage(cfg)
|
||||
return &t, nil
|
||||
}
|
||||
|
||||
func (s *TriggerStore) Update(ctx context.Context, t *models.Trigger) error {
|
||||
cfg, _ := json.Marshal(t.Config)
|
||||
if len(cfg) == 0 {
|
||||
cfg = []byte("{}")
|
||||
}
|
||||
_, err := DB.ExecContext(ctx,
|
||||
`UPDATE triggers SET enabled = $1, slug = $2, secret = $3, event_pattern = $4,
|
||||
entry_point = $5, config = $6, updated_at = NOW()
|
||||
WHERE id = $7`,
|
||||
t.Enabled, nilIfEmpty(t.Slug), nilIfEmpty(t.Secret),
|
||||
nilIfEmpty(t.EventPattern), t.EntryPoint, cfg, t.ID)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) Delete(ctx context.Context, id string) error {
|
||||
_, err := DB.ExecContext(ctx, `DELETE FROM triggers WHERE id = $1`, id)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) List(ctx context.Context, opts store.TriggerListOptions) ([]models.Trigger, int, error) {
|
||||
where := "1=1"
|
||||
args := []any{}
|
||||
n := 0
|
||||
|
||||
if opts.PackageID != "" {
|
||||
n++
|
||||
where += fmt.Sprintf(" AND package_id = $%d", n)
|
||||
args = append(args, opts.PackageID)
|
||||
}
|
||||
if opts.Type != "" {
|
||||
n++
|
||||
where += fmt.Sprintf(" AND type = $%d", n)
|
||||
args = append(args, opts.Type)
|
||||
}
|
||||
if opts.Enabled != nil {
|
||||
n++
|
||||
where += fmt.Sprintf(" AND enabled = $%d", n)
|
||||
args = append(args, *opts.Enabled)
|
||||
}
|
||||
|
||||
// Count
|
||||
var total int
|
||||
countArgs := make([]any, len(args))
|
||||
copy(countArgs, args)
|
||||
err := DB.QueryRowContext(ctx, "SELECT COUNT(*) FROM triggers WHERE "+where, countArgs...).Scan(&total)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
|
||||
// Query
|
||||
limit := opts.Limit
|
||||
if limit <= 0 {
|
||||
limit = 50
|
||||
}
|
||||
n++
|
||||
query := fmt.Sprintf(
|
||||
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
|
||||
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
|
||||
created_at, updated_at
|
||||
FROM triggers WHERE %s ORDER BY created_at DESC LIMIT $%d OFFSET $%d`,
|
||||
where, n, n+1)
|
||||
args = append(args, limit, opts.Offset)
|
||||
|
||||
rows, err := DB.QueryContext(ctx, query, args...)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var triggers []models.Trigger
|
||||
for rows.Next() {
|
||||
var t models.Trigger
|
||||
var cfg []byte
|
||||
if err := rows.Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug,
|
||||
&t.EventPattern, &t.EntryPoint, &cfg, &t.FireCount,
|
||||
&t.LastFireAt, &t.LastError, &t.CreatedAt, &t.UpdatedAt); err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
t.Config = json.RawMessage(cfg)
|
||||
triggers = append(triggers, t)
|
||||
}
|
||||
if triggers == nil {
|
||||
triggers = []models.Trigger{}
|
||||
}
|
||||
return triggers, total, nil
|
||||
}
|
||||
|
||||
func (s *TriggerStore) ListByPackage(ctx context.Context, packageID string) ([]models.Trigger, error) {
|
||||
rows, err := DB.QueryContext(ctx,
|
||||
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
|
||||
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
|
||||
created_at, updated_at
|
||||
FROM triggers WHERE package_id = $1 ORDER BY created_at`, packageID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
return scanTriggerRows(rows)
|
||||
}
|
||||
|
||||
func (s *TriggerStore) ListEnabledByType(ctx context.Context, triggerType string) ([]models.Trigger, error) {
|
||||
rows, err := DB.QueryContext(ctx,
|
||||
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
|
||||
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
|
||||
created_at, updated_at
|
||||
FROM triggers WHERE type = $1 AND enabled = true ORDER BY created_at`, triggerType)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
return scanTriggerRows(rows)
|
||||
}
|
||||
|
||||
func (s *TriggerStore) SetEnabled(ctx context.Context, id string, enabled bool) error {
|
||||
_, err := DB.ExecContext(ctx,
|
||||
`UPDATE triggers SET enabled = $1, updated_at = NOW() WHERE id = $2`, enabled, id)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) UpdateFireState(ctx context.Context, id string, lastFireAt time.Time, lastError string) error {
|
||||
_, err := DB.ExecContext(ctx,
|
||||
`UPDATE triggers SET last_fire_at = $1, last_error = $2, fire_count = fire_count + 1,
|
||||
updated_at = NOW() WHERE id = $3`,
|
||||
lastFireAt, nilIfEmpty(lastError), id)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) GetWebhook(ctx context.Context, packageID, slug string) (*models.Trigger, error) {
|
||||
var t models.Trigger
|
||||
var cfg []byte
|
||||
err := DB.QueryRowContext(ctx,
|
||||
`SELECT id, package_id, type, enabled, COALESCE(slug,''), secret,
|
||||
COALESCE(event_pattern,''), entry_point, config, fire_count,
|
||||
last_fire_at, COALESCE(last_error,''), created_at, updated_at
|
||||
FROM triggers
|
||||
WHERE package_id = $1 AND slug = $2 AND type = 'webhook'`, packageID, slug).
|
||||
Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug, &t.Secret,
|
||||
&t.EventPattern, &t.EntryPoint, &cfg, &t.FireCount,
|
||||
&t.LastFireAt, &t.LastError, &t.CreatedAt, &t.UpdatedAt)
|
||||
if err == sql.ErrNoRows {
|
||||
return nil, nil
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
t.Config = json.RawMessage(cfg)
|
||||
return &t, nil
|
||||
}
|
||||
|
||||
func (s *TriggerStore) DeleteForPackage(ctx context.Context, packageID string) error {
|
||||
_, err := DB.ExecContext(ctx, `DELETE FROM triggers WHERE package_id = $1`, packageID)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) LogExecution(ctx context.Context, log *models.TriggerLog) error {
|
||||
_, err := DB.ExecContext(ctx,
|
||||
`INSERT INTO trigger_logs (id, trigger_id, scheduled_task_id, fired_at, duration_ms, success, error, output)
|
||||
VALUES (gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7)`,
|
||||
nilIfEmpty(log.TriggerID), nilIfEmpty(log.ScheduledTaskID),
|
||||
log.FiredAt, log.DurationMs, log.Success, nilIfEmpty(log.Error), nilIfEmpty(log.Output))
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) ListLogs(ctx context.Context, triggerID string, limit int) ([]models.TriggerLog, error) {
|
||||
if limit <= 0 {
|
||||
limit = 50
|
||||
}
|
||||
rows, err := DB.QueryContext(ctx,
|
||||
`SELECT id, COALESCE(trigger_id::text,''), COALESCE(scheduled_task_id::text,''),
|
||||
fired_at, duration_ms, success, COALESCE(error,''), COALESCE(output,'')
|
||||
FROM trigger_logs
|
||||
WHERE trigger_id = $1
|
||||
ORDER BY fired_at DESC LIMIT $2`, triggerID, limit)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
return scanLogRows(rows)
|
||||
}
|
||||
|
||||
func (s *TriggerStore) PruneLogs(ctx context.Context, before time.Time) (int64, error) {
|
||||
result, err := DB.ExecContext(ctx,
|
||||
`DELETE FROM trigger_logs WHERE fired_at < $1`, before)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return result.RowsAffected()
|
||||
}
|
||||
|
||||
// ── Helpers ──────────────────────────────────
|
||||
|
||||
func scanTriggerRows(rows *sql.Rows) ([]models.Trigger, error) {
|
||||
var triggers []models.Trigger
|
||||
for rows.Next() {
|
||||
var t models.Trigger
|
||||
var cfg []byte
|
||||
if err := rows.Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug,
|
||||
&t.EventPattern, &t.EntryPoint, &cfg, &t.FireCount,
|
||||
&t.LastFireAt, &t.LastError, &t.CreatedAt, &t.UpdatedAt); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
t.Config = json.RawMessage(cfg)
|
||||
triggers = append(triggers, t)
|
||||
}
|
||||
if triggers == nil {
|
||||
triggers = []models.Trigger{}
|
||||
}
|
||||
return triggers, nil
|
||||
}
|
||||
|
||||
func scanLogRows(rows *sql.Rows) ([]models.TriggerLog, error) {
|
||||
var logs []models.TriggerLog
|
||||
for rows.Next() {
|
||||
var l models.TriggerLog
|
||||
if err := rows.Scan(&l.ID, &l.TriggerID, &l.ScheduledTaskID,
|
||||
&l.FiredAt, &l.DurationMs, &l.Success, &l.Error, &l.Output); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
logs = append(logs, l)
|
||||
}
|
||||
if logs == nil {
|
||||
logs = []models.TriggerLog{}
|
||||
}
|
||||
return logs, nil
|
||||
}
|
||||
|
||||
func nilIfEmpty(s string) *string {
|
||||
if s == "" {
|
||||
return nil
|
||||
}
|
||||
return &s
|
||||
}
|
||||
33
server/store/scheduled_task_iface.go
Normal file
33
server/store/scheduled_task_iface.go
Normal file
@@ -0,0 +1,33 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"switchboard-core/models"
|
||||
)
|
||||
|
||||
// ScheduledTaskStore manages user-created cron-scheduled Starlark scripts.
|
||||
type ScheduledTaskStore interface {
|
||||
// CRUD
|
||||
Create(ctx context.Context, t *models.ScheduledTask) error
|
||||
GetByID(ctx context.Context, id string) (*models.ScheduledTask, error)
|
||||
Update(ctx context.Context, t *models.ScheduledTask) error
|
||||
Delete(ctx context.Context, id string) error
|
||||
|
||||
// Queries
|
||||
List(ctx context.Context, opts ScheduledTaskListOptions) ([]models.ScheduledTask, int, error)
|
||||
ListByCreator(ctx context.Context, creatorID string) ([]models.ScheduledTask, error)
|
||||
ListEnabled(ctx context.Context) ([]models.ScheduledTask, error)
|
||||
|
||||
// Lifecycle
|
||||
SetEnabled(ctx context.Context, id string, enabled bool) error
|
||||
UpdateFireState(ctx context.Context, id string, lastFireAt time.Time, nextFireAt *time.Time, lastError string, durationMs int) error
|
||||
}
|
||||
|
||||
// ScheduledTaskListOptions provides filtering for scheduled task list queries.
|
||||
type ScheduledTaskListOptions struct {
|
||||
ListOptions
|
||||
CreatorID string
|
||||
Enabled *bool
|
||||
}
|
||||
216
server/store/sqlite/scheduled_tasks.go
Normal file
216
server/store/sqlite/scheduled_tasks.go
Normal file
@@ -0,0 +1,216 @@
|
||||
package sqlite
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
|
||||
"switchboard-core/database"
|
||||
"switchboard-core/models"
|
||||
"switchboard-core/store"
|
||||
)
|
||||
|
||||
type ScheduledTaskStore struct{}
|
||||
|
||||
func NewScheduledTaskStore() *ScheduledTaskStore {
|
||||
return &ScheduledTaskStore{}
|
||||
}
|
||||
|
||||
func (s *ScheduledTaskStore) Create(ctx context.Context, t *models.ScheduledTask) error {
|
||||
t.ID = uuid.New().String()
|
||||
now := time.Now().UTC().Format(time.RFC3339)
|
||||
t.CreatedAt, _ = time.Parse(time.RFC3339, now)
|
||||
t.UpdatedAt = t.CreatedAt
|
||||
|
||||
params, _ := json.Marshal(t.TemplateParams)
|
||||
if len(params) == 0 {
|
||||
params = []byte("{}")
|
||||
}
|
||||
|
||||
var nextFire *string
|
||||
if t.NextFireAt != nil {
|
||||
s := t.NextFireAt.UTC().Format(time.RFC3339)
|
||||
nextFire = &s
|
||||
}
|
||||
|
||||
_, err := DB.ExecContext(ctx,
|
||||
`INSERT INTO scheduled_tasks (id, name, description, creator_id, run_as, cron_expr,
|
||||
next_fire_at, enabled, script, template_id, template_params, created_at, updated_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
|
||||
t.ID, t.Name, t.Description, t.CreatorID, t.RunAs, t.CronExpr,
|
||||
nextFire, boolToInt(t.Enabled), t.Script,
|
||||
nullIfEmpty(t.TemplateID), string(params), now, now)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *ScheduledTaskStore) GetByID(ctx context.Context, id string) (*models.ScheduledTask, error) {
|
||||
var t models.ScheduledTask
|
||||
var paramsStr string
|
||||
err := DB.QueryRowContext(ctx,
|
||||
`SELECT id, name, description, creator_id, run_as, cron_expr,
|
||||
next_fire_at, last_fire_at, enabled, script,
|
||||
COALESCE(template_id,''), template_params, fire_count,
|
||||
COALESCE(last_error,''), last_duration_ms, created_at, updated_at
|
||||
FROM scheduled_tasks WHERE id = ?`, id).
|
||||
Scan(&t.ID, &t.Name, &t.Description, &t.CreatorID, &t.RunAs, &t.CronExpr,
|
||||
database.SNT(&t.NextFireAt), database.SNT(&t.LastFireAt),
|
||||
&t.Enabled, &t.Script,
|
||||
&t.TemplateID, ¶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
|
||||
}
|
||||
@@ -30,5 +30,7 @@ func NewStores(db *sql.DB) store.Stores {
|
||||
ExtData: NewExtDataStore(),
|
||||
Tickets: NewTicketStore(),
|
||||
RateLimits: NewRateLimitStore(),
|
||||
Triggers: NewTriggerStore(),
|
||||
ScheduledTasks: NewScheduledTaskStore(),
|
||||
}
|
||||
}
|
||||
|
||||
269
server/store/sqlite/triggers.go
Normal file
269
server/store/sqlite/triggers.go
Normal file
@@ -0,0 +1,269 @@
|
||||
package sqlite
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
|
||||
"switchboard-core/database"
|
||||
"switchboard-core/models"
|
||||
"switchboard-core/store"
|
||||
)
|
||||
|
||||
type TriggerStore struct{}
|
||||
|
||||
func NewTriggerStore() *TriggerStore {
|
||||
return &TriggerStore{}
|
||||
}
|
||||
|
||||
func (s *TriggerStore) Create(ctx context.Context, t *models.Trigger) error {
|
||||
t.ID = uuid.New().String()
|
||||
cfg, _ := json.Marshal(t.Config)
|
||||
if len(cfg) == 0 {
|
||||
cfg = []byte("{}")
|
||||
}
|
||||
_, err := DB.ExecContext(ctx,
|
||||
`INSERT INTO triggers (id, package_id, type, enabled, slug, secret, event_pattern,
|
||||
entry_point, config)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`,
|
||||
t.ID, t.PackageID, t.Type, boolToInt(t.Enabled),
|
||||
nullIfEmpty(t.Slug), nullIfEmpty(t.Secret),
|
||||
nullIfEmpty(t.EventPattern), t.EntryPoint, string(cfg))
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) GetByID(ctx context.Context, id string) (*models.Trigger, error) {
|
||||
var t models.Trigger
|
||||
var cfgStr string
|
||||
err := DB.QueryRowContext(ctx,
|
||||
`SELECT id, package_id, type, enabled, COALESCE(slug,''), secret,
|
||||
COALESCE(event_pattern,''), entry_point, config, fire_count,
|
||||
last_fire_at, COALESCE(last_error,''), created_at, updated_at
|
||||
FROM triggers WHERE id = ?`, id).
|
||||
Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug, &t.Secret,
|
||||
&t.EventPattern, &t.EntryPoint, &cfgStr, &t.FireCount,
|
||||
database.SNT(&t.LastFireAt), &t.LastError,
|
||||
database.ST(&t.CreatedAt), database.ST(&t.UpdatedAt))
|
||||
if err == sql.ErrNoRows {
|
||||
return nil, nil
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
t.Config = json.RawMessage(cfgStr)
|
||||
return &t, nil
|
||||
}
|
||||
|
||||
func (s *TriggerStore) Update(ctx context.Context, t *models.Trigger) error {
|
||||
cfg, _ := json.Marshal(t.Config)
|
||||
if len(cfg) == 0 {
|
||||
cfg = []byte("{}")
|
||||
}
|
||||
_, err := DB.ExecContext(ctx,
|
||||
`UPDATE triggers SET enabled = ?, slug = ?, secret = ?, event_pattern = ?,
|
||||
entry_point = ?, config = ?, updated_at = datetime('now')
|
||||
WHERE id = ?`,
|
||||
boolToInt(t.Enabled), nullIfEmpty(t.Slug), nullIfEmpty(t.Secret),
|
||||
nullIfEmpty(t.EventPattern), t.EntryPoint, string(cfg), t.ID)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) Delete(ctx context.Context, id string) error {
|
||||
_, err := DB.ExecContext(ctx, `DELETE FROM triggers WHERE id = ?`, id)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) List(ctx context.Context, opts store.TriggerListOptions) ([]models.Trigger, int, error) {
|
||||
where := "1=1"
|
||||
args := []any{}
|
||||
|
||||
if opts.PackageID != "" {
|
||||
where += " AND package_id = ?"
|
||||
args = append(args, opts.PackageID)
|
||||
}
|
||||
if opts.Type != "" {
|
||||
where += " AND type = ?"
|
||||
args = append(args, opts.Type)
|
||||
}
|
||||
if opts.Enabled != nil {
|
||||
where += " AND enabled = ?"
|
||||
args = append(args, boolToInt(*opts.Enabled))
|
||||
}
|
||||
|
||||
var total int
|
||||
countArgs := make([]any, len(args))
|
||||
copy(countArgs, args)
|
||||
err := DB.QueryRowContext(ctx, "SELECT COUNT(*) FROM triggers WHERE "+where, countArgs...).Scan(&total)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
|
||||
limit := opts.Limit
|
||||
if limit <= 0 {
|
||||
limit = 50
|
||||
}
|
||||
query := fmt.Sprintf(
|
||||
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
|
||||
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
|
||||
created_at, updated_at
|
||||
FROM triggers WHERE %s ORDER BY created_at DESC LIMIT ? OFFSET ?`, where)
|
||||
args = append(args, limit, opts.Offset)
|
||||
|
||||
rows, err := DB.QueryContext(ctx, query, args...)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
triggers, err := scanTriggerRowsSqlite(rows)
|
||||
return triggers, total, err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) ListByPackage(ctx context.Context, packageID string) ([]models.Trigger, error) {
|
||||
rows, err := DB.QueryContext(ctx,
|
||||
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
|
||||
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
|
||||
created_at, updated_at
|
||||
FROM triggers WHERE package_id = ? ORDER BY created_at`, packageID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
return scanTriggerRowsSqlite(rows)
|
||||
}
|
||||
|
||||
func (s *TriggerStore) ListEnabledByType(ctx context.Context, triggerType string) ([]models.Trigger, error) {
|
||||
rows, err := DB.QueryContext(ctx,
|
||||
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
|
||||
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
|
||||
created_at, updated_at
|
||||
FROM triggers WHERE type = ? AND enabled = 1 ORDER BY created_at`, triggerType)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
return scanTriggerRowsSqlite(rows)
|
||||
}
|
||||
|
||||
func (s *TriggerStore) SetEnabled(ctx context.Context, id string, enabled bool) error {
|
||||
_, err := DB.ExecContext(ctx,
|
||||
`UPDATE triggers SET enabled = ?, updated_at = datetime('now') WHERE id = ?`,
|
||||
boolToInt(enabled), id)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) UpdateFireState(ctx context.Context, id string, lastFireAt time.Time, lastError string) error {
|
||||
_, err := DB.ExecContext(ctx,
|
||||
`UPDATE triggers SET last_fire_at = ?, last_error = ?, fire_count = fire_count + 1,
|
||||
updated_at = datetime('now') WHERE id = ?`,
|
||||
lastFireAt.UTC().Format(time.RFC3339), nullIfEmpty(lastError), id)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) GetWebhook(ctx context.Context, packageID, slug string) (*models.Trigger, error) {
|
||||
var t models.Trigger
|
||||
var cfgStr string
|
||||
err := DB.QueryRowContext(ctx,
|
||||
`SELECT id, package_id, type, enabled, COALESCE(slug,''), secret,
|
||||
COALESCE(event_pattern,''), entry_point, config, fire_count,
|
||||
last_fire_at, COALESCE(last_error,''), created_at, updated_at
|
||||
FROM triggers
|
||||
WHERE package_id = ? AND slug = ? AND type = 'webhook'`, packageID, slug).
|
||||
Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug, &t.Secret,
|
||||
&t.EventPattern, &t.EntryPoint, &cfgStr, &t.FireCount,
|
||||
database.SNT(&t.LastFireAt), &t.LastError,
|
||||
database.ST(&t.CreatedAt), database.ST(&t.UpdatedAt))
|
||||
if err == sql.ErrNoRows {
|
||||
return nil, nil
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
t.Config = json.RawMessage(cfgStr)
|
||||
return &t, nil
|
||||
}
|
||||
|
||||
func (s *TriggerStore) DeleteForPackage(ctx context.Context, packageID string) error {
|
||||
_, err := DB.ExecContext(ctx, `DELETE FROM triggers WHERE package_id = ?`, packageID)
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) LogExecution(ctx context.Context, log *models.TriggerLog) error {
|
||||
_, err := DB.ExecContext(ctx,
|
||||
`INSERT INTO trigger_logs (id, trigger_id, scheduled_task_id, fired_at, duration_ms, success, error, output)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?)`,
|
||||
uuid.New().String(), nullIfEmpty(log.TriggerID), nullIfEmpty(log.ScheduledTaskID),
|
||||
log.FiredAt, log.DurationMs, boolToInt(log.Success),
|
||||
nullIfEmpty(log.Error), nullIfEmpty(log.Output))
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *TriggerStore) ListLogs(ctx context.Context, triggerID string, limit int) ([]models.TriggerLog, error) {
|
||||
if limit <= 0 {
|
||||
limit = 50
|
||||
}
|
||||
rows, err := DB.QueryContext(ctx,
|
||||
`SELECT id, COALESCE(trigger_id,''), COALESCE(scheduled_task_id,''),
|
||||
fired_at, duration_ms, success, COALESCE(error,''), COALESCE(output,'')
|
||||
FROM trigger_logs
|
||||
WHERE trigger_id = ?
|
||||
ORDER BY fired_at DESC LIMIT ?`, triggerID, limit)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
return scanLogRowsSqlite(rows)
|
||||
}
|
||||
|
||||
func (s *TriggerStore) PruneLogs(ctx context.Context, before time.Time) (int64, error) {
|
||||
result, err := DB.ExecContext(ctx,
|
||||
`DELETE FROM trigger_logs WHERE fired_at < ?`, before.UTC().Format(time.RFC3339))
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return result.RowsAffected()
|
||||
}
|
||||
|
||||
// ── Helpers ──────────────────────────────────
|
||||
|
||||
func scanTriggerRowsSqlite(rows *sql.Rows) ([]models.Trigger, error) {
|
||||
var triggers []models.Trigger
|
||||
for rows.Next() {
|
||||
var t models.Trigger
|
||||
var cfgStr string
|
||||
if err := rows.Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug,
|
||||
&t.EventPattern, &t.EntryPoint, &cfgStr, &t.FireCount,
|
||||
database.SNT(&t.LastFireAt), &t.LastError,
|
||||
database.ST(&t.CreatedAt), database.ST(&t.UpdatedAt)); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
t.Config = json.RawMessage(cfgStr)
|
||||
triggers = append(triggers, t)
|
||||
}
|
||||
if triggers == nil {
|
||||
triggers = []models.Trigger{}
|
||||
}
|
||||
return triggers, nil
|
||||
}
|
||||
|
||||
func scanLogRowsSqlite(rows *sql.Rows) ([]models.TriggerLog, error) {
|
||||
var logs []models.TriggerLog
|
||||
for rows.Next() {
|
||||
var l models.TriggerLog
|
||||
if err := rows.Scan(&l.ID, &l.TriggerID, &l.ScheduledTaskID,
|
||||
&l.FiredAt, &l.DurationMs, &l.Success, &l.Error, &l.Output); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
logs = append(logs, l)
|
||||
}
|
||||
if logs == nil {
|
||||
logs = []models.TriggerLog{}
|
||||
}
|
||||
return logs, nil
|
||||
}
|
||||
|
||||
// boolToInt and nullIfEmpty are defined in workflows.go (same package).
|
||||
45
server/store/trigger_iface.go
Normal file
45
server/store/trigger_iface.go
Normal file
@@ -0,0 +1,45 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"switchboard-core/models"
|
||||
)
|
||||
|
||||
// TriggerStore manages extension-declared event and webhook triggers.
|
||||
type TriggerStore interface {
|
||||
// CRUD
|
||||
Create(ctx context.Context, t *models.Trigger) error
|
||||
GetByID(ctx context.Context, id string) (*models.Trigger, error)
|
||||
Update(ctx context.Context, t *models.Trigger) error
|
||||
Delete(ctx context.Context, id string) error
|
||||
|
||||
// Queries
|
||||
List(ctx context.Context, opts TriggerListOptions) ([]models.Trigger, int, error)
|
||||
ListByPackage(ctx context.Context, packageID string) ([]models.Trigger, error)
|
||||
ListEnabledByType(ctx context.Context, triggerType string) ([]models.Trigger, error)
|
||||
|
||||
// Lifecycle
|
||||
SetEnabled(ctx context.Context, id string, enabled bool) error
|
||||
UpdateFireState(ctx context.Context, id string, lastFireAt time.Time, lastError string) error
|
||||
|
||||
// Webhook resolution
|
||||
GetWebhook(ctx context.Context, packageID, slug string) (*models.Trigger, error)
|
||||
|
||||
// Cleanup
|
||||
DeleteForPackage(ctx context.Context, packageID string) error
|
||||
|
||||
// Logs
|
||||
LogExecution(ctx context.Context, log *models.TriggerLog) error
|
||||
ListLogs(ctx context.Context, triggerID string, limit int) ([]models.TriggerLog, error)
|
||||
PruneLogs(ctx context.Context, before time.Time) (int64, error)
|
||||
}
|
||||
|
||||
// TriggerListOptions provides filtering for trigger list queries.
|
||||
type TriggerListOptions struct {
|
||||
ListOptions
|
||||
PackageID string
|
||||
Type string
|
||||
Enabled *bool
|
||||
}
|
||||
222
server/triggers/engine.go
Normal file
222
server/triggers/engine.go
Normal file
@@ -0,0 +1,222 @@
|
||||
// Package triggers — engine.go
|
||||
//
|
||||
// v0.2.2: Core trigger engine. Manages extension-declared triggers
|
||||
// (event bus subscriptions + webhook receivers) and user-created
|
||||
// scheduled tasks (cron). All three converge to sandbox.Runner.CallEntryPoint.
|
||||
package triggers
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"log"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/robfig/cron/v3"
|
||||
|
||||
"switchboard-core/events"
|
||||
"switchboard-core/models"
|
||||
"switchboard-core/sandbox"
|
||||
"switchboard-core/store"
|
||||
)
|
||||
|
||||
// Engine orchestrates trigger lifecycle: loading, firing, and cleanup.
|
||||
type Engine struct {
|
||||
stores store.Stores
|
||||
runner *sandbox.Runner
|
||||
bus *events.Bus
|
||||
cron *cron.Cron
|
||||
|
||||
mu sync.RWMutex
|
||||
unsubs map[string]func() // trigger_id → bus unsubscribe
|
||||
cronIDs map[string]cron.EntryID // scheduled_task_id → cron entry
|
||||
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
}
|
||||
|
||||
// New creates a trigger engine. Call Start() to begin processing.
|
||||
func New(stores store.Stores, runner *sandbox.Runner, bus *events.Bus) *Engine {
|
||||
return &Engine{
|
||||
stores: stores,
|
||||
runner: runner,
|
||||
bus: bus,
|
||||
cron: cron.New(cron.WithSeconds()),
|
||||
unsubs: make(map[string]func()),
|
||||
cronIDs: make(map[string]cron.EntryID),
|
||||
}
|
||||
}
|
||||
|
||||
// Start loads all enabled triggers and scheduled tasks, then begins execution.
|
||||
func (e *Engine) Start(ctx context.Context) error {
|
||||
e.ctx, e.cancel = context.WithCancel(ctx)
|
||||
|
||||
// Load event triggers
|
||||
if e.stores.Triggers != nil {
|
||||
eventTriggers, err := e.stores.Triggers.ListEnabledByType(ctx, models.TriggerTypeEvent)
|
||||
if err != nil {
|
||||
log.Printf(" ⚠️ triggers: failed to load event triggers: %v", err)
|
||||
} else {
|
||||
for i := range eventTriggers {
|
||||
e.wireEventTrigger(&eventTriggers[i])
|
||||
}
|
||||
if len(eventTriggers) > 0 {
|
||||
log.Printf(" 🔔 triggers: %d event trigger(s) loaded", len(eventTriggers))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Load scheduled tasks
|
||||
if e.stores.ScheduledTasks != nil {
|
||||
tasks, err := e.stores.ScheduledTasks.ListEnabled(ctx)
|
||||
if err != nil {
|
||||
log.Printf(" ⚠️ triggers: failed to load scheduled tasks: %v", err)
|
||||
} else {
|
||||
for i := range tasks {
|
||||
e.wireScheduledTask(&tasks[i])
|
||||
}
|
||||
if len(tasks) > 0 {
|
||||
log.Printf(" ⏰ triggers: %d scheduled task(s) loaded", len(tasks))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
e.cron.Start()
|
||||
return nil
|
||||
}
|
||||
|
||||
// Stop gracefully shuts down all trigger goroutines and cron jobs.
|
||||
func (e *Engine) Stop() {
|
||||
if e.cancel != nil {
|
||||
e.cancel()
|
||||
}
|
||||
|
||||
// Stop cron scheduler
|
||||
cronCtx := e.cron.Stop()
|
||||
<-cronCtx.Done()
|
||||
|
||||
// Unsubscribe all event triggers
|
||||
e.mu.Lock()
|
||||
for id, unsub := range e.unsubs {
|
||||
unsub()
|
||||
delete(e.unsubs, id)
|
||||
}
|
||||
e.mu.Unlock()
|
||||
|
||||
log.Printf(" 🔔 triggers: engine stopped")
|
||||
}
|
||||
|
||||
// RegisterTrigger adds a single trigger at runtime (package install).
|
||||
func (e *Engine) RegisterTrigger(t *models.Trigger) {
|
||||
if !t.Enabled {
|
||||
return
|
||||
}
|
||||
switch t.Type {
|
||||
case models.TriggerTypeEvent:
|
||||
e.wireEventTrigger(t)
|
||||
// Webhook triggers are resolved at request time — no wiring needed
|
||||
}
|
||||
}
|
||||
|
||||
// UnregisterTrigger removes a trigger at runtime (package uninstall/disable).
|
||||
func (e *Engine) UnregisterTrigger(triggerID string) {
|
||||
e.mu.Lock()
|
||||
defer e.mu.Unlock()
|
||||
if unsub, ok := e.unsubs[triggerID]; ok {
|
||||
unsub()
|
||||
delete(e.unsubs, triggerID)
|
||||
}
|
||||
}
|
||||
|
||||
// RegisterSchedule adds a scheduled task at runtime.
|
||||
func (e *Engine) RegisterSchedule(t *models.ScheduledTask) {
|
||||
if !t.Enabled {
|
||||
return
|
||||
}
|
||||
e.wireScheduledTask(t)
|
||||
}
|
||||
|
||||
// UnregisterSchedule removes a scheduled task at runtime.
|
||||
func (e *Engine) UnregisterSchedule(taskID string) {
|
||||
e.mu.Lock()
|
||||
defer e.mu.Unlock()
|
||||
if entryID, ok := e.cronIDs[taskID]; ok {
|
||||
e.cron.Remove(entryID)
|
||||
delete(e.cronIDs, taskID)
|
||||
}
|
||||
}
|
||||
|
||||
// ManualRun fires a scheduled task immediately (bypasses cron).
|
||||
func (e *Engine) ManualRun(t *models.ScheduledTask) {
|
||||
e.fireScheduledTask(t.ID, t.CreatorID, t.RunAs, t.Script, t.CronExpr)
|
||||
}
|
||||
|
||||
// ReloadPackageTriggers re-syncs triggers for a package from DB.
|
||||
func (e *Engine) ReloadPackageTriggers(ctx context.Context, packageID string) error {
|
||||
if e.stores.Triggers == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
triggers, err := e.stores.Triggers.ListByPackage(ctx, packageID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Unregister existing
|
||||
e.mu.RLock()
|
||||
var toRemove []string
|
||||
for id := range e.unsubs {
|
||||
toRemove = append(toRemove, id)
|
||||
}
|
||||
e.mu.RUnlock()
|
||||
|
||||
// We need to check which ones belong to this package
|
||||
for _, id := range toRemove {
|
||||
// Simple: unregister all, re-register from DB
|
||||
e.UnregisterTrigger(id)
|
||||
}
|
||||
|
||||
// Re-register enabled triggers
|
||||
for i := range triggers {
|
||||
e.RegisterTrigger(&triggers[i])
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// logExecution records a trigger or scheduled task execution.
|
||||
func (e *Engine) logExecution(triggerID, scheduledTaskID string, firedAt time.Time, durationMs int, success bool, errStr, output string) {
|
||||
if e.stores.Triggers == nil {
|
||||
return
|
||||
}
|
||||
tl := &models.TriggerLog{
|
||||
TriggerID: triggerID,
|
||||
ScheduledTaskID: scheduledTaskID,
|
||||
FiredAt: firedAt.UTC().Format(time.RFC3339),
|
||||
DurationMs: &durationMs,
|
||||
Success: success,
|
||||
Error: errStr,
|
||||
Output: truncate(output, 4000),
|
||||
}
|
||||
if err := e.stores.Triggers.LogExecution(context.Background(), tl); err != nil {
|
||||
log.Printf(" ⚠️ triggers: failed to log execution: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// publishEvent emits a trigger lifecycle event on the bus.
|
||||
func (e *Engine) publishEvent(label string, triggerID string) {
|
||||
if e.bus != nil {
|
||||
payload, _ := json.Marshal(map[string]any{"trigger_id": triggerID})
|
||||
e.bus.Publish(events.Event{
|
||||
Label: label,
|
||||
Payload: payload,
|
||||
Ts: time.Now().UnixMilli(),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func truncate(s string, max int) string {
|
||||
if len(s) <= max {
|
||||
return s
|
||||
}
|
||||
return s[:max] + "…"
|
||||
}
|
||||
163
server/triggers/event.go
Normal file
163
server/triggers/event.go
Normal file
@@ -0,0 +1,163 @@
|
||||
package triggers
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"go.starlark.net/starlark"
|
||||
|
||||
"switchboard-core/events"
|
||||
"switchboard-core/models"
|
||||
"switchboard-core/sandbox"
|
||||
)
|
||||
|
||||
// wireEventTrigger subscribes to the bus pattern for an event trigger.
|
||||
func (e *Engine) wireEventTrigger(t *models.Trigger) {
|
||||
if t.EventPattern == "" {
|
||||
return
|
||||
}
|
||||
|
||||
triggerID := t.ID
|
||||
packageID := t.PackageID
|
||||
entryPoint := t.EntryPoint
|
||||
pattern := t.EventPattern
|
||||
|
||||
unsub := e.bus.Subscribe(pattern, func(ev events.Event) {
|
||||
// Fire asynchronously — never block the event bus
|
||||
go e.fireEventTrigger(triggerID, packageID, entryPoint, ev)
|
||||
})
|
||||
|
||||
e.mu.Lock()
|
||||
e.unsubs[triggerID] = unsub
|
||||
e.mu.Unlock()
|
||||
}
|
||||
|
||||
// fireEventTrigger invokes the Starlark handler for an event trigger.
|
||||
func (e *Engine) fireEventTrigger(triggerID, packageID, entryPoint string, ev events.Event) {
|
||||
start := time.Now()
|
||||
ctx := e.ctx
|
||||
if ctx == nil {
|
||||
ctx = context.Background()
|
||||
}
|
||||
|
||||
// Re-check trigger is still enabled
|
||||
trigger, err := e.stores.Triggers.GetByID(ctx, triggerID)
|
||||
if err != nil || trigger == nil || !trigger.Enabled {
|
||||
return
|
||||
}
|
||||
|
||||
// Load package
|
||||
pkg, err := e.stores.Packages.Get(ctx, packageID)
|
||||
if err != nil || pkg == nil || pkg.Status != models.PackageStatusActive {
|
||||
return
|
||||
}
|
||||
|
||||
// Check triggers.register permission
|
||||
if !e.hasPermission(ctx, packageID, models.ExtPermTriggersRegister) {
|
||||
return
|
||||
}
|
||||
|
||||
// Build Starlark context dict
|
||||
ctxDict := starlark.NewDict(6)
|
||||
_ = ctxDict.SetKey(starlark.String("trigger_type"), starlark.String("event"))
|
||||
_ = ctxDict.SetKey(starlark.String("trigger_id"), starlark.String(triggerID))
|
||||
_ = ctxDict.SetKey(starlark.String("event_label"), starlark.String(ev.Label))
|
||||
_ = ctxDict.SetKey(starlark.String("event_ts"), starlark.MakeInt64(ev.Ts))
|
||||
if ev.Room != "" {
|
||||
_ = ctxDict.SetKey(starlark.String("event_room"), starlark.String(ev.Room))
|
||||
}
|
||||
if len(ev.Payload) > 0 {
|
||||
var payloadMap map[string]any
|
||||
if json.Unmarshal(ev.Payload, &payloadMap) == nil {
|
||||
_ = ctxDict.SetKey(starlark.String("event_payload"), mapToStarlark(payloadMap))
|
||||
} else {
|
||||
_ = ctxDict.SetKey(starlark.String("event_payload"), starlark.String(string(ev.Payload)))
|
||||
}
|
||||
}
|
||||
|
||||
// Call entry point with full sandbox context
|
||||
_, output, callErr := e.runner.CallEntryPoint(ctx, pkg, entryPoint,
|
||||
starlark.Tuple{ctxDict}, nil, nil)
|
||||
|
||||
duration := int(time.Since(start).Milliseconds())
|
||||
errStr := ""
|
||||
if callErr != nil {
|
||||
errStr = callErr.Error()
|
||||
log.Printf(" ⚠️ trigger[event] %s/%s error: %v", packageID, entryPoint, callErr)
|
||||
}
|
||||
|
||||
// Update fire state
|
||||
_ = e.stores.Triggers.UpdateFireState(ctx, triggerID, start, errStr)
|
||||
|
||||
// Log execution
|
||||
e.logExecution(triggerID, "", start, duration, callErr == nil, errStr, output)
|
||||
e.publishEvent("trigger.fired", triggerID)
|
||||
}
|
||||
|
||||
// hasPermission checks if a package has a specific granted permission.
|
||||
func (e *Engine) hasPermission(ctx context.Context, packageID, permission string) bool {
|
||||
if e.stores.ExtPermissions == nil {
|
||||
return false
|
||||
}
|
||||
granted, err := e.stores.ExtPermissions.GrantedForPackage(ctx, packageID)
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
for _, p := range granted {
|
||||
if p == permission {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// mapToStarlark converts a map[string]any to a Starlark dict.
|
||||
func mapToStarlark(m map[string]any) *starlark.Dict {
|
||||
d := starlark.NewDict(len(m))
|
||||
for k, v := range m {
|
||||
_ = d.SetKey(starlark.String(k), goToStarlark(v))
|
||||
}
|
||||
return d
|
||||
}
|
||||
|
||||
// goToStarlark converts a Go value to a Starlark value.
|
||||
func goToStarlark(v any) starlark.Value {
|
||||
switch val := v.(type) {
|
||||
case nil:
|
||||
return starlark.None
|
||||
case bool:
|
||||
return starlark.Bool(val)
|
||||
case int:
|
||||
return starlark.MakeInt(val)
|
||||
case int64:
|
||||
return starlark.MakeInt64(val)
|
||||
case float64:
|
||||
return starlark.Float(val)
|
||||
case string:
|
||||
return starlark.String(val)
|
||||
case map[string]any:
|
||||
return mapToStarlark(val)
|
||||
case []any:
|
||||
elems := make([]starlark.Value, len(val))
|
||||
for i, e := range val {
|
||||
elems[i] = goToStarlark(e)
|
||||
}
|
||||
return starlark.NewList(elems)
|
||||
default:
|
||||
return starlark.String(fmt.Sprintf("%v", val))
|
||||
}
|
||||
}
|
||||
|
||||
// RunContext builds a sandbox.RunContext for a trigger invocation.
|
||||
func triggerRunContext(userID, teamID string) *sandbox.RunContext {
|
||||
if userID == "" {
|
||||
return nil
|
||||
}
|
||||
return &sandbox.RunContext{
|
||||
UserID: userID,
|
||||
TeamID: teamID,
|
||||
}
|
||||
}
|
||||
25
server/triggers/global.go
Normal file
25
server/triggers/global.go
Normal file
@@ -0,0 +1,25 @@
|
||||
package triggers
|
||||
|
||||
import "sync"
|
||||
|
||||
// globalEngine holds a reference to the active trigger engine.
|
||||
// Set by SetGlobalEngine at startup, used by handlers that need
|
||||
// to wire triggers immediately on package install.
|
||||
var (
|
||||
globalMu sync.RWMutex
|
||||
globalEngine *Engine
|
||||
)
|
||||
|
||||
// SetGlobalEngine stores the engine reference for use by handlers.
|
||||
func SetGlobalEngine(e *Engine) {
|
||||
globalMu.Lock()
|
||||
globalEngine = e
|
||||
globalMu.Unlock()
|
||||
}
|
||||
|
||||
// GlobalEngine returns the current engine, or nil.
|
||||
func GlobalEngine() *Engine {
|
||||
globalMu.RLock()
|
||||
defer globalMu.RUnlock()
|
||||
return globalEngine
|
||||
}
|
||||
162
server/triggers/schedule.go
Normal file
162
server/triggers/schedule.go
Normal file
@@ -0,0 +1,162 @@
|
||||
package triggers
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"github.com/robfig/cron/v3"
|
||||
"go.starlark.net/starlark"
|
||||
|
||||
starlarkjson "go.starlark.net/lib/json"
|
||||
|
||||
"switchboard-core/models"
|
||||
"switchboard-core/sandbox"
|
||||
)
|
||||
|
||||
// wireScheduledTask registers a cron entry for a scheduled task.
|
||||
func (e *Engine) wireScheduledTask(t *models.ScheduledTask) {
|
||||
taskID := t.ID
|
||||
creatorID := t.CreatorID
|
||||
runAs := t.RunAs
|
||||
script := t.Script
|
||||
cronExpr := t.CronExpr
|
||||
|
||||
// Parse cron with standard 5-field format (minute hour dom month dow)
|
||||
parser := cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow)
|
||||
schedule, err := parser.Parse(cronExpr)
|
||||
if err != nil {
|
||||
log.Printf(" ⚠️ schedule[%s]: invalid cron %q: %v", taskID, cronExpr, err)
|
||||
return
|
||||
}
|
||||
|
||||
// Compute next fire time
|
||||
now := time.Now()
|
||||
nextFire := schedule.Next(now)
|
||||
_ = e.stores.ScheduledTasks.UpdateFireState(context.Background(), taskID,
|
||||
time.Time{}, &nextFire, "", 0)
|
||||
|
||||
entryID := e.cron.Schedule(schedule, cron.FuncJob(func() {
|
||||
e.fireScheduledTask(taskID, creatorID, runAs, script, cronExpr)
|
||||
}))
|
||||
|
||||
e.mu.Lock()
|
||||
e.cronIDs[taskID] = entryID
|
||||
e.mu.Unlock()
|
||||
}
|
||||
|
||||
// fireScheduledTask invokes a scheduled task's Starlark script in a restricted sandbox.
|
||||
func (e *Engine) fireScheduledTask(taskID, creatorID, runAs, script, cronExpr string) {
|
||||
start := time.Now()
|
||||
ctx := e.ctx
|
||||
if ctx == nil {
|
||||
ctx = context.Background()
|
||||
}
|
||||
|
||||
// Re-check task is still enabled
|
||||
task, err := e.stores.ScheduledTasks.GetByID(ctx, taskID)
|
||||
if err != nil || task == nil || !task.Enabled {
|
||||
return
|
||||
}
|
||||
|
||||
// Check creator is still active (unless run_as=system)
|
||||
if runAs == models.RunAsCreator {
|
||||
user, err := e.stores.Users.GetByID(ctx, creatorID)
|
||||
if err != nil || user == nil || !user.IsActive {
|
||||
log.Printf(" ⚠️ schedule[%s]: creator %s inactive, pausing", taskID, creatorID)
|
||||
_ = e.stores.ScheduledTasks.SetEnabled(ctx, taskID, false)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// Build restricted module set
|
||||
modules := e.buildRestrictedModules(ctx, creatorID, runAs)
|
||||
|
||||
// Execute script in sandbox
|
||||
sb := sandbox.New(sandbox.DefaultConfig())
|
||||
result, execErr := sb.ExecWithLoader(ctx, "schedule/"+taskID+".star", script, modules, nil)
|
||||
|
||||
var output string
|
||||
if result != nil {
|
||||
output = result.Output
|
||||
}
|
||||
|
||||
// Call main() if defined
|
||||
if execErr == nil && result != nil {
|
||||
if mainFn, ok := result.Globals["main"]; ok {
|
||||
if callable, ok := mainFn.(starlark.Callable); ok {
|
||||
// Build context dict
|
||||
ctxDict := starlark.NewDict(4)
|
||||
_ = ctxDict.SetKey(starlark.String("trigger_type"), starlark.String("schedule"))
|
||||
_ = ctxDict.SetKey(starlark.String("task_id"), starlark.String(taskID))
|
||||
_ = ctxDict.SetKey(starlark.String("cron_expr"), starlark.String(cronExpr))
|
||||
_ = ctxDict.SetKey(starlark.String("fired_at"), starlark.String(start.UTC().Format(time.RFC3339)))
|
||||
|
||||
_, callOutput, callErr := sb.Call(ctx, callable, starlark.Tuple{ctxDict}, nil)
|
||||
output += callOutput
|
||||
if callErr != nil {
|
||||
execErr = callErr
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
duration := int(time.Since(start).Milliseconds())
|
||||
errStr := ""
|
||||
if execErr != nil {
|
||||
errStr = execErr.Error()
|
||||
log.Printf(" ⚠️ schedule[%s] error: %v", taskID, execErr)
|
||||
}
|
||||
|
||||
// Compute next fire time
|
||||
parser := cron.NewParser(cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.Dow)
|
||||
if sched, err := parser.Parse(cronExpr); err == nil {
|
||||
nextFire := sched.Next(time.Now())
|
||||
_ = e.stores.ScheduledTasks.UpdateFireState(ctx, taskID, start, &nextFire, errStr, duration)
|
||||
} else {
|
||||
_ = e.stores.ScheduledTasks.UpdateFireState(ctx, taskID, start, nil, errStr, duration)
|
||||
}
|
||||
|
||||
// Log execution
|
||||
e.logExecution("", taskID, start, duration, execErr == nil, errStr, output)
|
||||
e.publishEvent("trigger.fired", taskID)
|
||||
}
|
||||
|
||||
// buildRestrictedModules creates the limited module set for scheduled tasks.
|
||||
// Scheduled scripts can: read settings, load libraries, use JSON, resolve connections.
|
||||
// They CANNOT: create DB tables, make raw HTTP calls, read secrets.
|
||||
func (e *Engine) buildRestrictedModules(ctx context.Context, creatorID, runAs string) map[string]starlark.Value {
|
||||
modules := make(map[string]starlark.Value)
|
||||
|
||||
// JSON is always available
|
||||
modules["json"] = starlarkjson.Module
|
||||
|
||||
// Settings module — read-only access to platform settings
|
||||
// (no package context for scheduled tasks; they read global settings)
|
||||
if e.stores.Packages != nil {
|
||||
modules["settings"] = sandbox.BuildSettingsModule(ctx, e.stores, "", creatorID, "")
|
||||
}
|
||||
|
||||
// Notifications module — if available
|
||||
if e.runner != nil {
|
||||
// The runner has the notifier attached; we need to access it.
|
||||
// For now, scheduled tasks don't get the notifications module
|
||||
// since it requires a package context. This can be enhanced later.
|
||||
}
|
||||
|
||||
// Connections module — read-only resolve for HTTP via existing connections
|
||||
// This allows scheduled tasks to make HTTP calls through configured connections.
|
||||
// The RunContext carries the creator's user ID for connection resolution.
|
||||
|
||||
return modules
|
||||
}
|
||||
|
||||
// RunContext builds a sandbox.RunContext for a scheduled task invocation.
|
||||
func scheduleRunContext(creatorID, runAs string) *sandbox.RunContext {
|
||||
if runAs == models.RunAsSystem {
|
||||
return nil // system context — no user scoping
|
||||
}
|
||||
return &sandbox.RunContext{
|
||||
UserID: creatorID,
|
||||
}
|
||||
}
|
||||
215
server/triggers/webhook.go
Normal file
215
server/triggers/webhook.go
Normal file
@@ -0,0 +1,215 @@
|
||||
package triggers
|
||||
|
||||
import (
|
||||
"crypto/hmac"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"go.starlark.net/starlark"
|
||||
|
||||
"switchboard-core/models"
|
||||
)
|
||||
|
||||
// HandleWebhook is the Gin handler for inbound webhook triggers.
|
||||
// Route: POST/GET /api/v1/hooks/:package_id/:slug
|
||||
func (e *Engine) HandleWebhook(c *gin.Context) {
|
||||
packageID := c.Param("package_id")
|
||||
slug := c.Param("slug")
|
||||
|
||||
if e.stores.Triggers == nil {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "triggers not available"})
|
||||
return
|
||||
}
|
||||
|
||||
// Resolve trigger
|
||||
trigger, err := e.stores.Triggers.GetWebhook(c.Request.Context(), packageID, slug)
|
||||
if err != nil {
|
||||
log.Printf(" ⚠️ webhook: lookup error %s/%s: %v", packageID, slug, err)
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": "internal error"})
|
||||
return
|
||||
}
|
||||
if trigger == nil {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "webhook not found"})
|
||||
return
|
||||
}
|
||||
if !trigger.Enabled {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "webhook disabled"})
|
||||
return
|
||||
}
|
||||
|
||||
// Read body
|
||||
body, err := io.ReadAll(io.LimitReader(c.Request.Body, 1<<20)) // 1MB limit
|
||||
if err != nil {
|
||||
c.JSON(http.StatusBadRequest, gin.H{"error": "failed to read body"})
|
||||
return
|
||||
}
|
||||
|
||||
// Verify HMAC signature if secret is set
|
||||
if trigger.Secret != "" {
|
||||
sig := c.GetHeader("X-Switchboard-Signature")
|
||||
if !verifyHMAC(body, trigger.Secret, sig) {
|
||||
c.JSON(http.StatusUnauthorized, gin.H{"error": "invalid signature"})
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// Load package
|
||||
pkg, err := e.stores.Packages.Get(c.Request.Context(), packageID)
|
||||
if err != nil || pkg == nil || pkg.Status != models.PackageStatusActive {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "package not available"})
|
||||
return
|
||||
}
|
||||
|
||||
// Check permission
|
||||
if !e.hasPermission(c.Request.Context(), packageID, models.ExtPermTriggersRegister) {
|
||||
c.JSON(http.StatusForbidden, gin.H{"error": "triggers.register permission not granted"})
|
||||
return
|
||||
}
|
||||
|
||||
// Build Starlark context dict
|
||||
ctxDict := starlark.NewDict(8)
|
||||
_ = ctxDict.SetKey(starlark.String("trigger_type"), starlark.String("webhook"))
|
||||
_ = ctxDict.SetKey(starlark.String("trigger_id"), starlark.String(trigger.ID))
|
||||
_ = ctxDict.SetKey(starlark.String("method"), starlark.String(c.Request.Method))
|
||||
_ = ctxDict.SetKey(starlark.String("path"), starlark.String(c.Request.URL.Path))
|
||||
_ = ctxDict.SetKey(starlark.String("body"), starlark.String(string(body)))
|
||||
|
||||
// Headers as dict
|
||||
headers := starlark.NewDict(len(c.Request.Header))
|
||||
for k, v := range c.Request.Header {
|
||||
if len(v) > 0 {
|
||||
_ = headers.SetKey(starlark.String(k), starlark.String(v[0]))
|
||||
}
|
||||
}
|
||||
_ = ctxDict.SetKey(starlark.String("headers"), headers)
|
||||
|
||||
// Query params as dict
|
||||
query := starlark.NewDict(len(c.Request.URL.Query()))
|
||||
for k, v := range c.Request.URL.Query() {
|
||||
if len(v) > 0 {
|
||||
_ = query.SetKey(starlark.String(k), starlark.String(v[0]))
|
||||
}
|
||||
}
|
||||
_ = ctxDict.SetKey(starlark.String("query"), query)
|
||||
|
||||
// Fire synchronously — return result as HTTP response
|
||||
start := time.Now()
|
||||
val, output, callErr := e.runner.CallEntryPoint(c.Request.Context(), pkg, trigger.EntryPoint,
|
||||
starlark.Tuple{ctxDict}, nil, nil)
|
||||
|
||||
duration := int(time.Since(start).Milliseconds())
|
||||
errStr := ""
|
||||
if callErr != nil {
|
||||
errStr = callErr.Error()
|
||||
log.Printf(" ⚠️ trigger[webhook] %s/%s error: %v", packageID, slug, callErr)
|
||||
}
|
||||
|
||||
// Update fire state
|
||||
_ = e.stores.Triggers.UpdateFireState(c.Request.Context(), trigger.ID, start, errStr)
|
||||
|
||||
// Log execution
|
||||
e.logExecution(trigger.ID, "", start, duration, callErr == nil, errStr, output)
|
||||
e.publishEvent("trigger.fired", trigger.ID)
|
||||
|
||||
if callErr != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"error": "handler error", "detail": errStr})
|
||||
return
|
||||
}
|
||||
|
||||
// Return Starlark result
|
||||
response := parseWebhookResponse(val)
|
||||
c.JSON(response.status, response.body)
|
||||
}
|
||||
|
||||
// webhookResponse is the parsed response from a Starlark webhook handler.
|
||||
type webhookResponse struct {
|
||||
status int
|
||||
body any
|
||||
}
|
||||
|
||||
// parseWebhookResponse extracts HTTP status and body from the Starlark return value.
|
||||
// Returns {status: X, body: Y} or default 200 OK.
|
||||
func parseWebhookResponse(val starlark.Value) webhookResponse {
|
||||
if val == nil || val == starlark.None {
|
||||
return webhookResponse{status: 200, body: gin.H{"ok": true}}
|
||||
}
|
||||
|
||||
// If it's a dict with "status" and/or "body" keys, use them
|
||||
if d, ok := val.(*starlark.Dict); ok {
|
||||
resp := webhookResponse{status: 200}
|
||||
|
||||
if statusVal, found, _ := d.Get(starlark.String("status")); found {
|
||||
if s, ok := statusVal.(starlark.Int); ok {
|
||||
if i, ok := s.Int64(); ok {
|
||||
resp.status = int(i)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if bodyVal, found, _ := d.Get(starlark.String("body")); found {
|
||||
resp.body = starlarkToGo(bodyVal)
|
||||
} else {
|
||||
resp.body = gin.H{"ok": true}
|
||||
}
|
||||
|
||||
return resp
|
||||
}
|
||||
|
||||
// String return → plain text body
|
||||
if s, ok := val.(starlark.String); ok {
|
||||
return webhookResponse{status: 200, body: gin.H{"result": string(s)}}
|
||||
}
|
||||
|
||||
return webhookResponse{status: 200, body: gin.H{"ok": true}}
|
||||
}
|
||||
|
||||
// starlarkToGo converts a Starlark value to a Go value (for JSON serialization).
|
||||
func starlarkToGo(v starlark.Value) any {
|
||||
switch val := v.(type) {
|
||||
case starlark.NoneType:
|
||||
return nil
|
||||
case starlark.Bool:
|
||||
return bool(val)
|
||||
case starlark.Int:
|
||||
if i, ok := val.Int64(); ok {
|
||||
return i
|
||||
}
|
||||
return val.String()
|
||||
case starlark.Float:
|
||||
return float64(val)
|
||||
case starlark.String:
|
||||
return string(val)
|
||||
case *starlark.List:
|
||||
result := make([]any, val.Len())
|
||||
for i := 0; i < val.Len(); i++ {
|
||||
result[i] = starlarkToGo(val.Index(i))
|
||||
}
|
||||
return result
|
||||
case *starlark.Dict:
|
||||
result := make(map[string]any, val.Len())
|
||||
for _, item := range val.Items() {
|
||||
if k, ok := item[0].(starlark.String); ok {
|
||||
result[string(k)] = starlarkToGo(item[1])
|
||||
}
|
||||
}
|
||||
return result
|
||||
default:
|
||||
return v.String()
|
||||
}
|
||||
}
|
||||
|
||||
// verifyHMAC checks the HMAC-SHA256 signature of a webhook payload.
|
||||
func verifyHMAC(body []byte, secret, signature string) bool {
|
||||
if signature == "" {
|
||||
return false
|
||||
}
|
||||
mac := hmac.New(sha256.New, []byte(secret))
|
||||
mac.Write(body)
|
||||
expected := hex.EncodeToString(mac.Sum(nil))
|
||||
return hmac.Equal([]byte(expected), []byte(signature))
|
||||
}
|
||||
Reference in New Issue
Block a user