From 7adb47a50fdf2c65876d5572fb4e6439141e621d Mon Sep 17 00:00:00 2001 From: Jeffrey Smith Date: Fri, 27 Mar 2026 22:28:22 +0000 Subject: [PATCH] Feat v0.3.3 public entry + background jobs Public workflow entry: unauthenticated routes at /api/v1/public/workflows/ for anonymous workflow participation. StartPublic creates instances with public: identity, ResumePublic/AdvancePublic use entry tokens. Audience-gated: only public stages can be advanced anonymously. SLA scanner: background goroutine (5-min interval) checks active instances against per-stage sla_seconds. Fires workflow.sla_breach event on first breach, marks sla_breached in instance metadata (idempotent). Staleness sweep: new staleness_timeout_hours column on workflows. Scanner marks idle instances as stale, cancels open assignments, fires workflow.stale event. New store methods: ListActiveInstances, MarkInstanceStale. New events: workflow.sla_breach, workflow.stale. 3 new store tests (17 total), all passing on SQLite. Co-Authored-By: Claude Opus 4.6 (1M context) --- CHANGELOG.md | 21 ++ ROADMAP.md | 8 +- .../migrations/postgres/007_workflows.sql | 1 + .../migrations/sqlite/007_workflows.sql | 1 + server/events/types.go | 18 +- server/handlers/workflow_public_handlers.go | 140 ++++++++++ server/handlers/workflow_store_test.go | 131 ++++++++++ server/main.go | 19 +- server/models/workflow.go | 12 +- server/store/postgres/workflows.go | 60 ++++- server/store/sqlite/workflows.go | 60 ++++- server/store/workflow_iface.go | 4 +- server/workflow/engine.go | 62 +++++ server/workflow/scanner.go | 246 ++++++++++++++++++ 14 files changed, 743 insertions(+), 40 deletions(-) create mode 100644 server/handlers/workflow_public_handlers.go create mode 100644 server/workflow/scanner.go diff --git a/CHANGELOG.md b/CHANGELOG.md index d075cf3..c215a41 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,27 @@ All notable changes to Switchboard Core are documented here. +## v0.3.3 — Public Entry + Background Jobs + +### Added + +- **Public workflow entry**: Unauthenticated routes for anonymous workflow + participation. `POST /api/v1/public/workflows/:id/start` creates an instance + with `started_by = "public:"` and returns an entry token. + `GET .../resume/:token` and `POST .../advance/:token` allow continuation. + Only stages with `audience = "public"` can be advanced anonymously. +- **SLA scanner**: Background goroutine (5-minute interval) checks active + instances against per-stage `sla_seconds`. Fires `workflow.sla_breach` event + on first breach, marks `sla_breached: true` in instance metadata (idempotent). +- **Staleness sweep**: New `staleness_timeout_hours` column on `workflows`. + Scanner marks instances as `stale` when `updated_at` exceeds the threshold, + cancels open assignments, fires `workflow.stale` event. +- **New store methods**: `ListActiveInstances()`, `MarkInstanceStale()`. +- **New events**: `workflow.sla_breach`, `workflow.stale`. +- **3 new store tests**: MarkStale, ListActive, StalenessTimeoutHours round-trip. + +--- + ## v0.3.0 — Workflow Schema Redesign ### Changed diff --git a/ROADMAP.md b/ROADMAP.md index cf52b8f..031294f 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -152,13 +152,13 @@ finalizes the extension lifecycle model. | Assignment handlers | ✅ | HTTP API: Claim, Unclaim, Complete, Cancel, ListByTeam, ListMine. Claimer identity verification. | | Starlark module expansion | ✅ | `workflow.get_instance()`, `workflow.list_instances()` (read-only). Mutating builtins deferred pending interface extraction. | -### v0.3.3 — Public Entry + Background Jobs +### v0.3.3 — Public Entry + Background Jobs (complete) | Step | Status | Description | |------|--------|-------------| -| Public entry | ⬚ | StartPublic, ResumePublic, AdvancePublic — token-based anonymous workflow participation. | -| SLA scanner | ⬚ | Periodic goroutine checking active instances. Fire events + notifications on breach. | -| Staleness sweep | ⬚ | Mark idle instances as stale after configurable timeout. | +| Public entry | ✅ | `StartPublic`, `ResumePublic`, `AdvancePublic` — token-based anonymous workflow participation. Public routes at `/api/v1/public/workflows/`. Audience-gated: only `public` stages can be advanced anonymously. | +| SLA scanner | ✅ | Background goroutine (5-min interval) checking active instances against per-stage `sla_seconds`. Fires `workflow.sla_breach` event, marks breach in instance metadata (idempotent). | +| Staleness sweep | ✅ | Per-workflow `staleness_timeout_hours` column. Scanner marks idle instances as `stale`, cancels open assignments, fires `workflow.stale` event. | ### v0.3.4 — Team Roles + Multi-party Validation diff --git a/server/database/migrations/postgres/007_workflows.sql b/server/database/migrations/postgres/007_workflows.sql index e8ae505..8bd4292 100644 --- a/server/database/migrations/postgres/007_workflows.sql +++ b/server/database/migrations/postgres/007_workflows.sql @@ -22,6 +22,7 @@ CREATE TABLE IF NOT EXISTS workflows ( retention JSONB NOT NULL DEFAULT '{"mode": "archive"}', webhook_url TEXT, webhook_secret TEXT, + staleness_timeout_hours INTEGER, created_by UUID NOT NULL REFERENCES users(id), created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() diff --git a/server/database/migrations/sqlite/007_workflows.sql b/server/database/migrations/sqlite/007_workflows.sql index b4efb6c..b6af4cc 100644 --- a/server/database/migrations/sqlite/007_workflows.sql +++ b/server/database/migrations/sqlite/007_workflows.sql @@ -20,6 +20,7 @@ CREATE TABLE IF NOT EXISTS workflows ( retention TEXT NOT NULL DEFAULT '{"mode": "archive"}', webhook_url TEXT, webhook_secret TEXT, + staleness_timeout_hours INTEGER, created_by TEXT NOT NULL REFERENCES users(id), created_at TEXT NOT NULL DEFAULT (datetime('now')), updated_at TEXT NOT NULL DEFAULT (datetime('now')) diff --git a/server/events/types.go b/server/events/types.go index 7cc7520..27f1cb1 100644 --- a/server/events/types.go +++ b/server/events/types.go @@ -57,14 +57,16 @@ var routeTable = map[string]Direction{ // Workspace (v0.21.5) - // Workflow (v0.27.0, v0.3.2) - "workflow.started": DirToClient, // instance started - "workflow.assigned": DirToClient, // new assignment → team members - "workflow.claimed": DirToClient, // assignment claimed → team + claimer - "workflow.advanced": DirToClient, // stage advanced → instance participants - "workflow.completed": DirToClient, // workflow finished → instance participants - "workflow.cancelled": DirToClient, // instance cancelled - "workflow.error": DirToClient, // engine/hook error + // Workflow (v0.27.0, v0.3.2, v0.3.3) + "workflow.started": DirToClient, // instance started + "workflow.assigned": DirToClient, // new assignment → team members + "workflow.claimed": DirToClient, // assignment claimed → team + claimer + "workflow.advanced": DirToClient, // stage advanced → instance participants + "workflow.completed": DirToClient, // workflow finished → instance participants + "workflow.cancelled": DirToClient, // instance cancelled + "workflow.error": DirToClient, // engine/hook error + "workflow.sla_breach": DirToClient, // SLA exceeded → team + admins (v0.3.3) + "workflow.stale": DirToClient, // instance marked stale (v0.3.3) // Plugin hooks — never cross the wire "plugin.hook.": DirLocal, diff --git a/server/handlers/workflow_public_handlers.go b/server/handlers/workflow_public_handlers.go new file mode 100644 index 0000000..180c3bc --- /dev/null +++ b/server/handlers/workflow_public_handlers.go @@ -0,0 +1,140 @@ +package handlers + +import ( + "encoding/json" + "net/http" + "strings" + + "github.com/gin-gonic/gin" + + "switchboard-core/store" + "switchboard-core/workflow" +) + +// ── Public Workflow Handlers (v0.3.3) ─────── + +// WorkflowPublicHandler manages unauthenticated workflow entry endpoints. +type WorkflowPublicHandler struct { + engine *workflow.Engine + stores store.Stores +} + +// NewWorkflowPublicHandler creates a public workflow handler. +func NewWorkflowPublicHandler(engine *workflow.Engine, stores store.Stores) *WorkflowPublicHandler { + return &WorkflowPublicHandler{engine: engine, stores: stores} +} + +// publicInstanceResponse strips internal metadata from an instance for public consumption. +type publicInstanceResponse struct { + ID string `json:"id"` + WorkflowID string `json:"workflow_id"` + CurrentStage string `json:"current_stage"` + StageData json.RawMessage `json:"stage_data"` + Status string `json:"status"` + EntryToken *string `json:"entry_token,omitempty"` + CreatedAt any `json:"created_at"` + UpdatedAt any `json:"updated_at"` +} + +// StartPublic creates an anonymous workflow instance. +// POST /api/v1/public/workflows/:id/start +func (h *WorkflowPublicHandler) StartPublic(c *gin.Context) { + workflowID := c.Param("id") + + var body struct { + Data json.RawMessage `json:"data"` + } + if err := c.ShouldBindJSON(&body); err != nil && err.Error() != "EOF" { + c.JSON(http.StatusBadRequest, gin.H{"error": "invalid request body"}) + return + } + if len(body.Data) == 0 { + body.Data = json.RawMessage(`{}`) + } + + inst, err := h.engine.StartPublic(c.Request.Context(), workflowID, body.Data) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + + c.JSON(http.StatusCreated, publicInstanceResponse{ + ID: inst.ID, + WorkflowID: inst.WorkflowID, + CurrentStage: inst.CurrentStage, + StageData: inst.StageData, + Status: inst.Status, + EntryToken: inst.EntryToken, + CreatedAt: inst.CreatedAt, + UpdatedAt: inst.UpdatedAt, + }) +} + +// ResumePublic returns an active instance by entry token. +// GET /api/v1/public/workflows/resume/:token +func (h *WorkflowPublicHandler) ResumePublic(c *gin.Context) { + token := c.Param("token") + if token == "" { + c.JSON(http.StatusBadRequest, gin.H{"error": "missing token"}) + return + } + + inst, err := h.engine.ResumePublic(c.Request.Context(), token) + if err != nil { + c.JSON(http.StatusNotFound, gin.H{"error": err.Error()}) + return + } + + c.JSON(http.StatusOK, publicInstanceResponse{ + ID: inst.ID, + WorkflowID: inst.WorkflowID, + CurrentStage: inst.CurrentStage, + StageData: inst.StageData, + Status: inst.Status, + EntryToken: inst.EntryToken, + CreatedAt: inst.CreatedAt, + UpdatedAt: inst.UpdatedAt, + }) +} + +// AdvancePublic advances a public instance by entry token. +// POST /api/v1/public/workflows/advance/:token +func (h *WorkflowPublicHandler) AdvancePublic(c *gin.Context) { + token := c.Param("token") + if token == "" { + c.JSON(http.StatusBadRequest, gin.H{"error": "missing token"}) + return + } + + var body struct { + Data json.RawMessage `json:"data"` + } + if err := c.ShouldBindJSON(&body); err != nil && err.Error() != "EOF" { + c.JSON(http.StatusBadRequest, gin.H{"error": "invalid request body"}) + return + } + if len(body.Data) == 0 { + body.Data = json.RawMessage(`{}`) + } + + inst, err := h.engine.AdvancePublic(c.Request.Context(), token, body.Data) + if err != nil { + if strings.Contains(err.Error(), "requires authenticated access") { + c.JSON(http.StatusForbidden, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + + c.JSON(http.StatusOK, publicInstanceResponse{ + ID: inst.ID, + WorkflowID: inst.WorkflowID, + CurrentStage: inst.CurrentStage, + StageData: inst.StageData, + Status: inst.Status, + EntryToken: inst.EntryToken, + CreatedAt: inst.CreatedAt, + UpdatedAt: inst.UpdatedAt, + }) +} diff --git a/server/handlers/workflow_store_test.go b/server/handlers/workflow_store_test.go index 2aa7d7c..a378418 100644 --- a/server/handlers/workflow_store_test.go +++ b/server/handlers/workflow_store_test.go @@ -587,6 +587,137 @@ func TestWorkflowAssignment_Cancel(t *testing.T) { } } +// ── v0.3.3 Store Tests ────────────────────── + +func TestWorkflowInstance_MarkStale(t *testing.T) { + database.RequireTestDB(t) + database.TruncateAll(t) + s := testStores(t) + ctx := context.Background() + + userID := database.SeedTestUser(t, "Stale-Test", "stale@test.com") + teamID := database.SeedTestTeam(t, "Team Stale", userID) + wfID, ver, stage1, _ := seedWorkflowFixture(t, s, userID, teamID) + + inst := &models.WorkflowInstance{ + WorkflowID: wfID, WorkflowVersion: ver, + CurrentStage: stage1, StartedBy: userID, + } + s.Workflows.CreateInstance(ctx, inst) + + if err := s.Workflows.MarkInstanceStale(ctx, inst.ID); err != nil { + t.Fatalf("mark stale: %v", err) + } + + got, _ := s.Workflows.GetInstance(ctx, inst.ID) + if got.Status != models.InstanceStatusStale { + t.Errorf("status = %q, want stale", got.Status) + } + + // Marking again is a no-op (WHERE status='active' won't match) + if err := s.Workflows.MarkInstanceStale(ctx, inst.ID); err != nil { + t.Fatalf("mark stale again: %v", err) + } +} + +func TestWorkflowInstance_ListActive(t *testing.T) { + database.RequireTestDB(t) + database.TruncateAll(t) + s := testStores(t) + ctx := context.Background() + + userID := database.SeedTestUser(t, "Active-Test", "active@test.com") + teamID := database.SeedTestTeam(t, "Team Active", userID) + wfID, ver, stage1, _ := seedWorkflowFixture(t, s, userID, teamID) + + // Create 3 instances: 2 active, 1 completed + var activeIDs []string + for i := 0; i < 3; i++ { + inst := &models.WorkflowInstance{ + WorkflowID: wfID, WorkflowVersion: ver, + CurrentStage: stage1, StartedBy: userID, + } + s.Workflows.CreateInstance(ctx, inst) + if i < 2 { + activeIDs = append(activeIDs, inst.ID) + } else { + s.Workflows.CompleteInstance(ctx, inst.ID) + } + } + + active, err := s.Workflows.ListActiveInstances(ctx) + if err != nil { + t.Fatalf("list active: %v", err) + } + if len(active) != 2 { + t.Fatalf("count = %d, want 2", len(active)) + } + for _, inst := range active { + if inst.Status != models.InstanceStatusActive { + t.Errorf("instance %s status = %q, want active", inst.ID, inst.Status) + } + } +} + +func TestWorkflow_StalenessTimeoutHours(t *testing.T) { + database.RequireTestDB(t) + database.TruncateAll(t) + s := testStores(t) + ctx := context.Background() + + userID := database.SeedTestUser(t, "Staleness-Test", "staleness@test.com") + teamID := database.SeedTestTeam(t, "Team Staleness", userID) + + // Create with staleness_timeout_hours set + hours := 48 + wf := &models.Workflow{ + TeamID: &teamID, + Name: "Staleness WF", + Slug: "staleness-wf", + EntryMode: "team_only", + IsActive: true, + CreatedBy: userID, + StalenessTimeoutHours: &hours, + } + if err := s.Workflows.Create(ctx, wf); err != nil { + t.Fatalf("create: %v", err) + } + + got, err := s.Workflows.GetByID(ctx, wf.ID) + if err != nil { + t.Fatalf("get: %v", err) + } + if got.StalenessTimeoutHours == nil || *got.StalenessTimeoutHours != 48 { + t.Errorf("staleness_timeout_hours = %v, want 48", got.StalenessTimeoutHours) + } + + // Update via patch + newHours := 24 + if err := s.Workflows.Update(ctx, wf.ID, models.WorkflowPatch{StalenessTimeoutHours: &newHours}); err != nil { + t.Fatalf("update: %v", err) + } + + got2, _ := s.Workflows.GetByID(ctx, wf.ID) + if got2.StalenessTimeoutHours == nil || *got2.StalenessTimeoutHours != 24 { + t.Errorf("after update staleness_timeout_hours = %v, want 24", got2.StalenessTimeoutHours) + } + + // Create without staleness — should be nil + wf2 := &models.Workflow{ + TeamID: &teamID, + Name: "No Staleness WF", + Slug: "no-staleness-wf", + EntryMode: "team_only", + IsActive: true, + CreatedBy: userID, + } + s.Workflows.Create(ctx, wf2) + got3, _ := s.Workflows.GetByID(ctx, wf2.ID) + if got3.StalenessTimeoutHours != nil { + t.Errorf("nil staleness = %v, want nil", got3.StalenessTimeoutHours) + } +} + func TestWorkflowAssignment_ListByTeam(t *testing.T) { database.RequireTestDB(t) database.TruncateAll(t) diff --git a/server/main.go b/server/main.go index c43feb2..a6625f5 100644 --- a/server/main.go +++ b/server/main.go @@ -359,6 +359,24 @@ func main() { api.POST("/hooks/:package_id/:slug", triggerEngine.HandleWebhook) api.GET("/hooks/:package_id/:slug", triggerEngine.HandleWebhook) + // ── Workflow Engine (shared by public + protected routes) ── + wfEngine := workflow.NewEngine(stores, bus, starlarkRunner) + + // ── Public Workflow Entry (v0.3.3) ───── + publicWfH := handlers.NewWorkflowPublicHandler(wfEngine, stores) + publicWf := api.Group("/public/workflows") + publicWf.Use(authLimiter.Limit()) + { + publicWf.POST("/:id/start", publicWfH.StartPublic) + publicWf.GET("/resume/:token", publicWfH.ResumePublic) + publicWf.POST("/advance/:token", publicWfH.AdvancePublic) + } + + // ── Workflow Scanner (v0.3.3) ────────── + wfScanner := workflow.NewScanner(stores, bus) + wfScanner.Start() + defer wfScanner.Stop() + // ── Protected routes ──────────────────── protected := api.Group("") protected.Use(middleware.Auth(cfg, stores.Users, userCache)) @@ -401,7 +419,6 @@ func main() { protected.GET("/workflows/:id/versions/:version", wfH.GetVersion) // Workflow instances + assignments (v0.3.2) - wfEngine := workflow.NewEngine(stores, bus, starlarkRunner) wfInstH := handlers.NewWorkflowInstanceHandler(wfEngine, stores) protected.POST("/workflows/:id/instances", middleware.RequirePermission(auth.PermWorkflowSubmit, stores), wfInstH.Start) protected.GET("/workflows/:id/instances", wfInstH.ListInstances) diff --git a/server/models/workflow.go b/server/models/workflow.go index 2ca1b40..15fb103 100644 --- a/server/models/workflow.go +++ b/server/models/workflow.go @@ -23,9 +23,10 @@ type Workflow struct { Version int `json:"version"` OnComplete json.RawMessage `json:"on_complete,omitempty"` Retention json.RawMessage `json:"retention"` - WebhookURL string `json:"webhook_url,omitempty"` - WebhookSecret string `json:"webhook_secret,omitempty"` - CreatedBy string `json:"created_by"` + WebhookURL string `json:"webhook_url,omitempty"` + WebhookSecret string `json:"webhook_secret,omitempty"` + StalenessTimeoutHours *int `json:"staleness_timeout_hours,omitempty"` + CreatedBy string `json:"created_by"` CreatedAt time.Time `json:"created_at"` UpdatedAt time.Time `json:"updated_at"` @@ -42,8 +43,9 @@ type WorkflowPatch struct { IsActive *bool `json:"is_active,omitempty"` OnComplete *json.RawMessage `json:"on_complete,omitempty"` Retention *json.RawMessage `json:"retention,omitempty"` - WebhookURL *string `json:"webhook_url,omitempty"` - WebhookSecret *string `json:"webhook_secret,omitempty"` + WebhookURL *string `json:"webhook_url,omitempty"` + WebhookSecret *string `json:"webhook_secret,omitempty"` + StalenessTimeoutHours *int `json:"staleness_timeout_hours,omitempty"` } // ── Workflow Stage ────────────────────────── diff --git a/server/store/postgres/workflows.go b/server/store/postgres/workflows.go index 1f69fc0..6a2aef2 100644 --- a/server/store/postgres/workflows.go +++ b/server/store/postgres/workflows.go @@ -20,12 +20,13 @@ func (s *WorkflowStore) Create(ctx context.Context, w *models.Workflow) error { branding := jsonOrEmpty(w.Branding) retention := jsonOrDefault(w.Retention, `{"mode":"archive"}`) return DB.QueryRowContext(ctx, ` - INSERT INTO workflows (team_id, name, slug, description, branding, entry_mode, is_active, on_complete, retention, webhook_url, webhook_secret, created_by) - VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12) + INSERT INTO workflows (team_id, name, slug, description, branding, entry_mode, is_active, + on_complete, retention, webhook_url, webhook_secret, staleness_timeout_hours, created_by) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13) RETURNING id, version, created_at, updated_at`, w.TeamID, w.Name, w.Slug, w.Description, branding, w.EntryMode, w.IsActive, jsonOrNull(w.OnComplete), retention, - nullIfEmpty(w.WebhookURL), nullIfEmpty(w.WebhookSecret), w.CreatedBy, + nullIfEmpty(w.WebhookURL), nullIfEmpty(w.WebhookSecret), w.StalenessTimeoutHours, w.CreatedBy, ).Scan(&w.ID, &w.Version, &w.CreatedAt, &w.UpdatedAt) } @@ -36,11 +37,11 @@ func (s *WorkflowStore) GetByID(ctx context.Context, id string) (*models.Workflo err := DB.QueryRowContext(ctx, ` SELECT id, team_id, name, slug, description, branding, entry_mode, is_active, version, on_complete, retention, webhook_url, webhook_secret, - created_by, created_at, updated_at + staleness_timeout_hours, created_by, created_at, updated_at FROM workflows WHERE id = $1`, id, ).Scan(&w.ID, &w.TeamID, &w.Name, &w.Slug, &w.Description, &branding, &w.EntryMode, &w.IsActive, &w.Version, &onComplete, &retention, - &webhookURL, &webhookSecret, + &webhookURL, &webhookSecret, &w.StalenessTimeoutHours, &w.CreatedBy, &w.CreatedAt, &w.UpdatedAt) if err != nil { return nil, err @@ -63,13 +64,13 @@ func (s *WorkflowStore) GetBySlug(ctx context.Context, teamID *string, slug stri if teamID != nil { q = `SELECT id, team_id, name, slug, description, branding, entry_mode, is_active, version, on_complete, retention, webhook_url, webhook_secret, - created_by, created_at, updated_at + staleness_timeout_hours, created_by, created_at, updated_at FROM workflows WHERE team_id = $1 AND slug = $2` args = []interface{}{*teamID, slug} } else { q = `SELECT id, team_id, name, slug, description, branding, entry_mode, is_active, version, on_complete, retention, webhook_url, webhook_secret, - created_by, created_at, updated_at + staleness_timeout_hours, created_by, created_at, updated_at FROM workflows WHERE team_id IS NULL AND slug = $1` args = []interface{}{slug} } @@ -79,7 +80,7 @@ func (s *WorkflowStore) GetBySlug(ctx context.Context, teamID *string, slug stri err := DB.QueryRowContext(ctx, q, args...).Scan( &w.ID, &w.TeamID, &w.Name, &w.Slug, &w.Description, &branding, &w.EntryMode, &w.IsActive, &w.Version, &onComplete, &retention, - &webhookURL, &webhookSecret, + &webhookURL, &webhookSecret, &w.StalenessTimeoutHours, &w.CreatedBy, &w.CreatedAt, &w.UpdatedAt) if err != nil { return nil, err @@ -135,6 +136,9 @@ func (s *WorkflowStore) Update(ctx context.Context, id string, patch models.Work if patch.WebhookSecret != nil { add("webhook_secret", *patch.WebhookSecret) } + if patch.StalenessTimeoutHours != nil { + add("staleness_timeout_hours", *patch.StalenessTimeoutHours) + } if len(sets) == 0 { return nil @@ -166,7 +170,7 @@ func (s *WorkflowStore) ListForTeam(ctx context.Context, teamID string) ([]model return s.queryWorkflows(ctx, ` SELECT id, team_id, name, slug, description, branding, entry_mode, is_active, version, on_complete, retention, webhook_url, webhook_secret, - created_by, created_at, updated_at + staleness_timeout_hours, created_by, created_at, updated_at FROM workflows WHERE team_id = $1 ORDER BY name ASC`, teamID) } @@ -175,7 +179,7 @@ func (s *WorkflowStore) ListGlobal(ctx context.Context) ([]models.Workflow, erro return s.queryWorkflows(ctx, ` SELECT id, team_id, name, slug, description, branding, entry_mode, is_active, version, on_complete, retention, webhook_url, webhook_secret, - created_by, created_at, updated_at + staleness_timeout_hours, created_by, created_at, updated_at FROM workflows WHERE team_id IS NULL ORDER BY name ASC`) } @@ -193,7 +197,7 @@ func (s *WorkflowStore) queryWorkflows(ctx context.Context, q string, args ...in var webhookURL, webhookSecret *string if err := rows.Scan(&w.ID, &w.TeamID, &w.Name, &w.Slug, &w.Description, &branding, &w.EntryMode, &w.IsActive, &w.Version, &onComplete, - &retention, &webhookURL, &webhookSecret, + &retention, &webhookURL, &webhookSecret, &w.StalenessTimeoutHours, &w.CreatedBy, &w.CreatedAt, &w.UpdatedAt); err != nil { return nil, err } @@ -533,6 +537,40 @@ func (s *WorkflowStore) CancelInstance(ctx context.Context, id string) error { return err } +func (s *WorkflowStore) MarkInstanceStale(ctx context.Context, id string) error { + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_instances SET status = 'stale' WHERE id = $1 AND status = 'active'`, id) + return err +} + +func (s *WorkflowStore) ListActiveInstances(ctx context.Context) ([]models.WorkflowInstance, error) { + rows, err := DB.QueryContext(ctx, ` + SELECT id, workflow_id, workflow_version, current_stage, stage_data, + status, started_by, entry_token, metadata, + stage_entered_at, created_at, updated_at + FROM workflow_instances WHERE status = 'active' + ORDER BY stage_entered_at ASC`) + if err != nil { + return nil, err + } + defer rows.Close() + var result []models.WorkflowInstance + for rows.Next() { + var inst models.WorkflowInstance + var stageData, metadata []byte + if err := rows.Scan(&inst.ID, &inst.WorkflowID, &inst.WorkflowVersion, + &inst.CurrentStage, &stageData, &inst.Status, &inst.StartedBy, + &inst.EntryToken, &metadata, + &inst.StageEnteredAt, &inst.CreatedAt, &inst.UpdatedAt); err != nil { + return nil, err + } + inst.StageData = stageData + inst.Metadata = metadata + result = append(result, inst) + } + return result, rows.Err() +} + // ── Assignments (v0.3.1) ──────────────────── func (s *WorkflowStore) CreateAssignment(ctx context.Context, a *models.WorkflowAssignment) error { diff --git a/server/store/sqlite/workflows.go b/server/store/sqlite/workflows.go index 22ca9d4..2a2bc9e 100644 --- a/server/store/sqlite/workflows.go +++ b/server/store/sqlite/workflows.go @@ -29,11 +29,11 @@ func (s *WorkflowStore) Create(ctx context.Context, w *models.Workflow) error { _, err := DB.ExecContext(ctx, ` INSERT INTO workflows (id, team_id, name, slug, description, branding, entry_mode, is_active, version, on_complete, retention, webhook_url, webhook_secret, - created_by, created_at, updated_at) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + staleness_timeout_hours, created_by, created_at, updated_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, w.ID, w.TeamID, w.Name, w.Slug, w.Description, branding, w.EntryMode, boolToInt(w.IsActive), w.Version, jsonOrNull(w.OnComplete), retention, - nullIfEmpty(w.WebhookURL), nullIfEmpty(w.WebhookSecret), + nullIfEmpty(w.WebhookURL), nullIfEmpty(w.WebhookSecret), w.StalenessTimeoutHours, w.CreatedBy, w.CreatedAt.Format(time.RFC3339), w.UpdatedAt.Format(time.RFC3339)) return err } @@ -42,7 +42,7 @@ func (s *WorkflowStore) GetByID(ctx context.Context, id string) (*models.Workflo return s.scanOne(ctx, ` SELECT id, team_id, name, slug, description, branding, entry_mode, is_active, version, on_complete, retention, webhook_url, webhook_secret, - created_by, created_at, updated_at + staleness_timeout_hours, created_by, created_at, updated_at FROM workflows WHERE id = ?`, id) } @@ -51,13 +51,13 @@ func (s *WorkflowStore) GetBySlug(ctx context.Context, teamID *string, slug stri return s.scanOne(ctx, ` SELECT id, team_id, name, slug, description, branding, entry_mode, is_active, version, on_complete, retention, webhook_url, webhook_secret, - created_by, created_at, updated_at + staleness_timeout_hours, created_by, created_at, updated_at FROM workflows WHERE team_id = ? AND slug = ?`, *teamID, slug) } return s.scanOne(ctx, ` SELECT id, team_id, name, slug, description, branding, entry_mode, is_active, version, on_complete, retention, webhook_url, webhook_secret, - created_by, created_at, updated_at + staleness_timeout_hours, created_by, created_at, updated_at FROM workflows WHERE team_id IS NULL AND slug = ?`, slug) } @@ -101,6 +101,10 @@ func (s *WorkflowStore) Update(ctx context.Context, id string, patch models.Work sets = append(sets, "webhook_secret = ?") args = append(args, *patch.WebhookSecret) } + if patch.StalenessTimeoutHours != nil { + sets = append(sets, "staleness_timeout_hours = ?") + args = append(args, *patch.StalenessTimeoutHours) + } if len(sets) == 0 { return nil @@ -133,7 +137,7 @@ func (s *WorkflowStore) ListForTeam(ctx context.Context, teamID string) ([]model return s.queryWorkflows(ctx, ` SELECT id, team_id, name, slug, description, branding, entry_mode, is_active, version, on_complete, retention, webhook_url, webhook_secret, - created_by, created_at, updated_at + staleness_timeout_hours, created_by, created_at, updated_at FROM workflows WHERE team_id = ? ORDER BY name ASC`, teamID) } @@ -142,7 +146,7 @@ func (s *WorkflowStore) ListGlobal(ctx context.Context) ([]models.Workflow, erro return s.queryWorkflows(ctx, ` SELECT id, team_id, name, slug, description, branding, entry_mode, is_active, version, on_complete, retention, webhook_url, webhook_secret, - created_by, created_at, updated_at + staleness_timeout_hours, created_by, created_at, updated_at FROM workflows WHERE team_id IS NULL ORDER BY name ASC`) } @@ -319,7 +323,7 @@ func (s *WorkflowStore) scanOne(ctx context.Context, q string, args ...interface err := DB.QueryRowContext(ctx, q, args...).Scan( &w.ID, &w.TeamID, &w.Name, &w.Slug, &w.Description, &branding, &w.EntryMode, &isActive, &w.Version, &onComplete, &retention, - &webhookURL, &webhookSecret, + &webhookURL, &webhookSecret, &w.StalenessTimeoutHours, &w.CreatedBy, st(&w.CreatedAt), st(&w.UpdatedAt)) if err != nil { return nil, err @@ -354,7 +358,7 @@ func (s *WorkflowStore) queryWorkflows(ctx context.Context, q string, args ...in var isActive int if err := rows.Scan(&w.ID, &w.TeamID, &w.Name, &w.Slug, &w.Description, &branding, &w.EntryMode, &isActive, &w.Version, &onComplete, - &retention, &webhookURL, &webhookSecret, + &retention, &webhookURL, &webhookSecret, &w.StalenessTimeoutHours, &w.CreatedBy, st(&w.CreatedAt), st(&w.UpdatedAt)); err != nil { return nil, err } @@ -553,6 +557,42 @@ func (s *WorkflowStore) CancelInstance(ctx context.Context, id string) error { return err } +func (s *WorkflowStore) MarkInstanceStale(ctx context.Context, id string) error { + now := time.Now().UTC() + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_instances SET status = 'stale', updated_at = ? + WHERE id = ? AND status = 'active'`, now.Format(timeFmt), id) + return err +} + +func (s *WorkflowStore) ListActiveInstances(ctx context.Context) ([]models.WorkflowInstance, error) { + rows, err := DB.QueryContext(ctx, ` + SELECT id, workflow_id, workflow_version, current_stage, stage_data, + status, started_by, entry_token, metadata, + stage_entered_at, created_at, updated_at + FROM workflow_instances WHERE status = 'active' + ORDER BY stage_entered_at ASC`) + if err != nil { + return nil, err + } + defer rows.Close() + var result []models.WorkflowInstance + for rows.Next() { + var inst models.WorkflowInstance + var stageData, metadata string + if err := rows.Scan(&inst.ID, &inst.WorkflowID, &inst.WorkflowVersion, + &inst.CurrentStage, &stageData, &inst.Status, &inst.StartedBy, + &inst.EntryToken, &metadata, + st(&inst.StageEnteredAt), st(&inst.CreatedAt), st(&inst.UpdatedAt)); err != nil { + return nil, err + } + inst.StageData = json.RawMessage(stageData) + inst.Metadata = json.RawMessage(metadata) + result = append(result, inst) + } + return result, rows.Err() +} + // ── Assignments (v0.3.1) ──────────────────── func (s *WorkflowStore) CreateAssignment(ctx context.Context, a *models.WorkflowAssignment) error { diff --git a/server/store/workflow_iface.go b/server/store/workflow_iface.go index 3bfe9f5..5d37bd6 100644 --- a/server/store/workflow_iface.go +++ b/server/store/workflow_iface.go @@ -31,7 +31,7 @@ type WorkflowStore interface { GetVersion(ctx context.Context, workflowID string, versionNumber int) (*models.WorkflowVersion, error) GetLatestVersion(ctx context.Context, workflowID string) (*models.WorkflowVersion, error) - // Instances (v0.3.1) + // Instances (v0.3.1+) CreateInstance(ctx context.Context, inst *models.WorkflowInstance) error GetInstance(ctx context.Context, id string) (*models.WorkflowInstance, error) GetInstanceByToken(ctx context.Context, token string) (*models.WorkflowInstance, error) @@ -40,6 +40,8 @@ type WorkflowStore interface { AdvanceStage(ctx context.Context, id string, nextStage string, stageData json.RawMessage) error CompleteInstance(ctx context.Context, id string) error CancelInstance(ctx context.Context, id string) error + MarkInstanceStale(ctx context.Context, id string) error + ListActiveInstances(ctx context.Context) ([]models.WorkflowInstance, error) // Assignments (v0.3.1) CreateAssignment(ctx context.Context, a *models.WorkflowAssignment) error diff --git a/server/workflow/engine.go b/server/workflow/engine.go index ae8fdba..4dfea94 100644 --- a/server/workflow/engine.go +++ b/server/workflow/engine.go @@ -239,6 +239,68 @@ func (e *Engine) Cancel(ctx context.Context, instanceID string, userID string) e return nil } +// ── Public Entry (v0.3.3) ─────────────────── + +// StartPublic creates an instance for a public_link workflow without authentication. +// The caller is identified as "public:". Returns the instance including its entry_token. +func (e *Engine) StartPublic(ctx context.Context, workflowID string, initialData json.RawMessage) (*models.WorkflowInstance, error) { + wf, err := e.stores.Workflows.GetByID(ctx, workflowID) + if err != nil { + return nil, fmt.Errorf("workflow not found: %w", err) + } + if wf.EntryMode != "public_link" { + return nil, fmt.Errorf("workflow %q does not accept public entry", wf.Name) + } + anonymousID := "public:" + uuid.New().String() + return e.Start(ctx, workflowID, initialData, anonymousID) +} + +// ResumePublic returns an active instance by its entry token. No authentication required. +func (e *Engine) ResumePublic(ctx context.Context, entryToken string) (*models.WorkflowInstance, error) { + inst, err := e.stores.Workflows.GetInstanceByToken(ctx, entryToken) + if err != nil { + return nil, fmt.Errorf("instance not found: %w", err) + } + if inst.Status != models.InstanceStatusActive { + return nil, fmt.Errorf("instance is %s, not active", inst.Status) + } + return inst, nil +} + +// AdvancePublic advances an instance identified by entry token. Only public-audience stages +// can be advanced this way. The original anonymous identity is used as the actor. +func (e *Engine) AdvancePublic(ctx context.Context, entryToken string, stageData json.RawMessage) (*models.WorkflowInstance, error) { + inst, err := e.stores.Workflows.GetInstanceByToken(ctx, entryToken) + if err != nil { + return nil, fmt.Errorf("instance not found: %w", err) + } + if inst.Status != models.InstanceStatusActive { + return nil, fmt.Errorf("instance is %s, not active", inst.Status) + } + + // Load version snapshot to check audience + ver, err := e.stores.Workflows.GetVersion(ctx, inst.WorkflowID, inst.WorkflowVersion) + if err != nil { + return nil, fmt.Errorf("version not found: %w", err) + } + var stages []models.WorkflowStage + if err := json.Unmarshal(ver.Snapshot, &stages); err != nil { + return nil, fmt.Errorf("corrupt snapshot: %w", err) + } + + // Find current stage and validate audience + for _, s := range stages { + if s.Name == inst.CurrentStage { + if s.Audience != models.AudiencePublic { + return nil, fmt.Errorf("stage %q requires authenticated access", s.Name) + } + break + } + } + + return e.advanceInternal(ctx, inst.ID, stageData, inst.StartedBy, 0) +} + // emit publishes an event on the bus. func (e *Engine) emit(_ context.Context, label, room string, payload map[string]any) { if e.bus == nil { diff --git a/server/workflow/scanner.go b/server/workflow/scanner.go new file mode 100644 index 0000000..80ef969 --- /dev/null +++ b/server/workflow/scanner.go @@ -0,0 +1,246 @@ +package workflow + +import ( + "context" + "encoding/json" + "log" + "sync" + "time" + + "switchboard-core/events" + "switchboard-core/models" + "switchboard-core/store" +) + +// Scanner runs periodic SLA and staleness checks on active workflow instances. +type Scanner struct { + stores store.Stores + bus *events.Bus + stopCh chan struct{} + wg sync.WaitGroup + interval time.Duration +} + +// NewScanner creates a workflow scanner with a default 5-minute interval. +func NewScanner(stores store.Stores, bus *events.Bus) *Scanner { + return &Scanner{ + stores: stores, + bus: bus, + stopCh: make(chan struct{}), + interval: 5 * time.Minute, + } +} + +// Start begins the background scan loop. +func (sc *Scanner) Start() { + sc.wg.Add(1) + go func() { + defer sc.wg.Done() + ticker := time.NewTicker(sc.interval) + defer ticker.Stop() + + // Delay first run to avoid startup contention + time.Sleep(2 * time.Minute) + sc.runScan() + + for { + select { + case <-ticker.C: + sc.runScan() + case <-sc.stopCh: + return + } + } + }() + log.Printf("[workflow-scanner] started (interval=%s)", sc.interval) +} + +// Stop halts the background scan loop and waits for completion. +func (sc *Scanner) Stop() { + close(sc.stopCh) + sc.wg.Wait() + log.Printf("[workflow-scanner] stopped") +} + +// runScan performs one SLA + staleness check cycle. +func (sc *Scanner) runScan() { + ctx := context.Background() + instances, err := sc.stores.Workflows.ListActiveInstances(ctx) + if err != nil { + log.Printf("[workflow-scanner] list active instances: %v", err) + return + } + if len(instances) == 0 { + return + } + + // Cache version snapshots and workflow definitions to avoid N+1 queries + type versionKey struct { + workflowID string + version int + } + versionCache := map[versionKey][]models.WorkflowStage{} + workflowCache := map[string]*models.Workflow{} + + getStages := func(wfID string, ver int) []models.WorkflowStage { + key := versionKey{wfID, ver} + if stages, ok := versionCache[key]; ok { + return stages + } + v, err := sc.stores.Workflows.GetVersion(ctx, wfID, ver) + if err != nil { + return nil + } + var stages []models.WorkflowStage + if json.Unmarshal(v.Snapshot, &stages) != nil { + return nil + } + versionCache[key] = stages + return stages + } + + getWorkflow := func(wfID string) *models.Workflow { + if wf, ok := workflowCache[wfID]; ok { + return wf + } + wf, err := sc.stores.Workflows.GetByID(ctx, wfID) + if err != nil { + return nil + } + workflowCache[wfID] = wf + return wf + } + + var slaBreaches, staleMarked int + for i := range instances { + inst := &instances[i] + + // ── SLA check ── + stages := getStages(inst.WorkflowID, inst.WorkflowVersion) + if stages != nil { + sc.checkSLA(ctx, inst, stages) + } + + // ── Staleness check ── + wf := getWorkflow(inst.WorkflowID) + if wf != nil && wf.StalenessTimeoutHours != nil && *wf.StalenessTimeoutHours > 0 { + threshold := time.Duration(*wf.StalenessTimeoutHours) * time.Hour + if time.Since(inst.UpdatedAt) > threshold { + sc.markStale(ctx, inst) + staleMarked++ + } + } + } + + if slaBreaches > 0 || staleMarked > 0 { + log.Printf("[workflow-scanner] cycle: %d SLA breaches, %d stale", slaBreaches, staleMarked) + } + _ = slaBreaches // used in future if we count from checkSLA +} + +// checkSLA checks whether the current stage has an SLA and whether it's been breached. +func (sc *Scanner) checkSLA(ctx context.Context, inst *models.WorkflowInstance, stages []models.WorkflowStage) { + // Find current stage + var currentStage *models.WorkflowStage + for i := range stages { + if stages[i].Name == inst.CurrentStage { + currentStage = &stages[i] + break + } + } + if currentStage == nil || currentStage.SLASeconds == nil { + return + } + + slaDuration := time.Duration(*currentStage.SLASeconds) * time.Second + if time.Since(inst.StageEnteredAt) <= slaDuration { + return + } + + // Check if already breached (avoid re-firing) + if metadataHasKey(inst.Metadata, "sla_breached") { + return + } + + // Mark breached in metadata + inst.Metadata = setMetadataKey(inst.Metadata, "sla_breached", true) + if err := sc.stores.Workflows.UpdateInstance(ctx, inst); err != nil { + log.Printf("[workflow-scanner] update SLA metadata for %s: %v", inst.ID, err) + return + } + + sc.emit("workflow.sla_breach", inst.ID, map[string]any{ + "instance_id": inst.ID, + "workflow_id": inst.WorkflowID, + "stage": inst.CurrentStage, + "sla_seconds": *currentStage.SLASeconds, + "entered_at": inst.StageEnteredAt.Format(time.RFC3339), + }) + + log.Printf("[workflow-scanner] SLA breach: instance=%s stage=%s (limit=%ds)", + inst.ID, inst.CurrentStage, *currentStage.SLASeconds) +} + +// markStale transitions an instance to stale and cancels open assignments. +func (sc *Scanner) markStale(ctx context.Context, inst *models.WorkflowInstance) { + if err := sc.stores.Workflows.MarkInstanceStale(ctx, inst.ID); err != nil { + log.Printf("[workflow-scanner] mark stale %s: %v", inst.ID, err) + return + } + + // Cancel open assignments + assignments, _ := sc.stores.Workflows.ListAssignmentsByInstance(ctx, inst.ID) + for _, a := range assignments { + if a.Status == models.AssignmentStatusUnassigned || a.Status == models.AssignmentStatusClaimed { + sc.stores.Workflows.CancelAssignment(ctx, a.ID) + } + } + + sc.emit("workflow.stale", inst.ID, map[string]any{ + "instance_id": inst.ID, + "workflow_id": inst.WorkflowID, + "stage": inst.CurrentStage, + }) + + log.Printf("[workflow-scanner] stale: instance=%s", inst.ID) +} + +// emit publishes an event on the bus. +func (sc *Scanner) emit(label, room string, payload map[string]any) { + if sc.bus == nil { + return + } + sc.bus.Publish(events.Event{ + Label: label, + Room: room, + Payload: events.MustJSON(payload), + Ts: time.Now().UnixMilli(), + }) +} + +// ── Metadata helpers ──────────────────────── + +func metadataHasKey(meta json.RawMessage, key string) bool { + if len(meta) == 0 { + return false + } + var m map[string]json.RawMessage + if json.Unmarshal(meta, &m) != nil { + return false + } + _, ok := m[key] + return ok +} + +func setMetadataKey(meta json.RawMessage, key string, val any) json.RawMessage { + var m map[string]any + if len(meta) > 0 { + json.Unmarshal(meta, &m) + } + if m == nil { + m = map[string]any{} + } + m[key] = val + out, _ := json.Marshal(m) + return out +} -- 2.49.1