Feat event bus subscriptions + trigger system (v0.2.2)
All checks were successful
CI/CD / detect-changes (pull_request) Successful in 24s
CI/CD / test-frontend (pull_request) Has been skipped
CI/CD / test-go-pg (pull_request) Successful in 2m30s
CI/CD / test-sqlite (pull_request) Successful in 2m37s
CI/CD / build-and-deploy (pull_request) Successful in 1m30s

Three trigger primitives replacing the old monolithic scheduler:

- Event triggers: extensions subscribe to bus patterns via manifest,
  async handler invocation through sandbox.CallEntryPoint
- Webhook triggers: inbound HTTP at /api/v1/hooks/:pkg/:slug with
  HMAC-SHA256 verification and synchronous Starlark response
- Scheduled tasks: user-created cron scripts with restricted sandbox
  (no raw HTTP, no DB table creation), runs as creator identity

New tables: triggers, scheduled_tasks, trigger_logs (postgres + sqlite).
New permission: triggers.register. Full admin + user CRUD APIs.
SyncManifestTriggers hooked into seed and install flows.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
2026-03-26 22:29:43 +00:00
parent 342f9e3eb4
commit bd703b9e0d
29 changed files with 3237 additions and 3 deletions

View File

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

View File

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

View File

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