From 01ee9e668bba123a8c7e124cfbbe84228bf28ee8 Mon Sep 17 00:00:00 2001 From: Jeffrey Smith Date: Fri, 27 Mar 2026 20:03:26 +0000 Subject: [PATCH] Feat v0.3.1 instance assignment (#15) Co-authored-by: Jeffrey Smith Co-committed-by: Jeffrey Smith --- .../migrations/postgres/007_workflows.sql | 52 ++++ .../migrations/sqlite/007_workflows.sql | 50 ++++ server/models/workflow.go | 70 +++++ server/store/postgres/workflows.go | 234 ++++++++++++++++- server/store/sqlite/workflows.go | 245 +++++++++++++++++- server/store/workflow_iface.go | 23 +- 6 files changed, 667 insertions(+), 7 deletions(-) diff --git a/server/database/migrations/postgres/007_workflows.sql b/server/database/migrations/postgres/007_workflows.sql index de5bc9c..e8ae505 100644 --- a/server/database/migrations/postgres/007_workflows.sql +++ b/server/database/migrations/postgres/007_workflows.sql @@ -78,3 +78,55 @@ CREATE TABLE IF NOT EXISTS workflow_versions ( CREATE INDEX IF NOT EXISTS idx_workflow_versions_workflow ON workflow_versions(workflow_id, version_number DESC); + +-- ── Workflow Instances (v0.3.1) ────────────── + +CREATE TABLE IF NOT EXISTS workflow_instances ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + workflow_id UUID NOT NULL REFERENCES workflows(id) ON DELETE CASCADE, + workflow_version INTEGER NOT NULL, + current_stage TEXT NOT NULL, + stage_data JSONB NOT NULL DEFAULT '{}', + status TEXT NOT NULL DEFAULT 'active' + CHECK (status IN ('active', 'completed', 'cancelled', 'stale', 'error')), + started_by TEXT NOT NULL, + entry_token TEXT, + metadata JSONB NOT NULL DEFAULT '{}', + stage_entered_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); + +CREATE INDEX IF NOT EXISTS idx_workflow_instances_workflow_status + ON workflow_instances(workflow_id, status); +CREATE INDEX IF NOT EXISTS idx_workflow_instances_active + ON workflow_instances(status) WHERE status = 'active'; +CREATE UNIQUE INDEX IF NOT EXISTS idx_workflow_instances_token + ON workflow_instances(entry_token) WHERE entry_token IS NOT NULL; + +DROP TRIGGER IF EXISTS workflow_instances_updated_at ON workflow_instances; +CREATE TRIGGER workflow_instances_updated_at BEFORE UPDATE ON workflow_instances + FOR EACH ROW EXECUTE FUNCTION update_updated_at(); + +-- ── Workflow Assignments (v0.3.1) ──────────── + +CREATE TABLE IF NOT EXISTS workflow_assignments ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + instance_id UUID NOT NULL REFERENCES workflow_instances(id) ON DELETE CASCADE, + stage TEXT NOT NULL, + team_id UUID NOT NULL REFERENCES teams(id) ON DELETE CASCADE, + assigned_to UUID REFERENCES users(id) ON DELETE SET NULL, + status TEXT NOT NULL DEFAULT 'unassigned' + CHECK (status IN ('unassigned', 'claimed', 'completed', 'cancelled')), + review_data JSONB NOT NULL DEFAULT '{}', + claimed_at TIMESTAMPTZ, + completed_at TIMESTAMPTZ, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); + +CREATE INDEX IF NOT EXISTS idx_workflow_assignments_instance_stage + ON workflow_assignments(instance_id, stage); +CREATE INDEX IF NOT EXISTS idx_workflow_assignments_team_status + ON workflow_assignments(team_id, status); +CREATE INDEX IF NOT EXISTS idx_workflow_assignments_assigned + ON workflow_assignments(assigned_to) WHERE assigned_to IS NOT NULL; diff --git a/server/database/migrations/sqlite/007_workflows.sql b/server/database/migrations/sqlite/007_workflows.sql index 85be6f4..b4efb6c 100644 --- a/server/database/migrations/sqlite/007_workflows.sql +++ b/server/database/migrations/sqlite/007_workflows.sql @@ -69,3 +69,53 @@ CREATE TABLE IF NOT EXISTS workflow_versions ( CREATE INDEX IF NOT EXISTS idx_workflow_versions_workflow ON workflow_versions(workflow_id, version_number); + +-- ── Workflow Instances (v0.3.1) ────────────── + +CREATE TABLE IF NOT EXISTS workflow_instances ( + id TEXT PRIMARY KEY, + workflow_id TEXT NOT NULL REFERENCES workflows(id) ON DELETE CASCADE, + workflow_version INTEGER NOT NULL, + current_stage TEXT NOT NULL, + stage_data TEXT NOT NULL DEFAULT '{}', + status TEXT NOT NULL DEFAULT 'active' + CHECK (status IN ('active', 'completed', 'cancelled', 'stale', 'error')), + started_by TEXT NOT NULL, + entry_token TEXT UNIQUE, + metadata TEXT NOT NULL DEFAULT '{}', + stage_entered_at TEXT NOT NULL DEFAULT (datetime('now')), + created_at TEXT NOT NULL DEFAULT (datetime('now')), + updated_at TEXT NOT NULL DEFAULT (datetime('now')) +); + +CREATE INDEX IF NOT EXISTS idx_workflow_instances_workflow_status + ON workflow_instances(workflow_id, status); +CREATE INDEX IF NOT EXISTS idx_workflow_instances_active + ON workflow_instances(status); + +CREATE TRIGGER IF NOT EXISTS workflow_instances_updated_at AFTER UPDATE ON workflow_instances + FOR EACH ROW WHEN NEW.updated_at = OLD.updated_at + BEGIN UPDATE workflow_instances SET updated_at = datetime('now') WHERE id = NEW.id; END; + +-- ── Workflow Assignments (v0.3.1) ──────────── + +CREATE TABLE IF NOT EXISTS workflow_assignments ( + id TEXT PRIMARY KEY, + instance_id TEXT NOT NULL REFERENCES workflow_instances(id) ON DELETE CASCADE, + stage TEXT NOT NULL, + team_id TEXT NOT NULL REFERENCES teams(id) ON DELETE CASCADE, + assigned_to TEXT REFERENCES users(id) ON DELETE SET NULL, + status TEXT NOT NULL DEFAULT 'unassigned' + CHECK (status IN ('unassigned', 'claimed', 'completed', 'cancelled')), + review_data TEXT NOT NULL DEFAULT '{}', + claimed_at TEXT, + completed_at TEXT, + created_at TEXT NOT NULL DEFAULT (datetime('now')) +); + +CREATE INDEX IF NOT EXISTS idx_workflow_assignments_instance_stage + ON workflow_assignments(instance_id, stage); +CREATE INDEX IF NOT EXISTS idx_workflow_assignments_team_status + ON workflow_assignments(team_id, status); +CREATE INDEX IF NOT EXISTS idx_workflow_assignments_assigned + ON workflow_assignments(assigned_to); diff --git a/server/models/workflow.go b/server/models/workflow.go index 6420473..2ca1b40 100644 --- a/server/models/workflow.go +++ b/server/models/workflow.go @@ -432,6 +432,76 @@ func evaluateFieldCondition(cond *FieldCondition, data map[string]interface{}) b } } +// ── Instance Status Constants ─────────────── + +const ( + InstanceStatusActive = "active" + InstanceStatusCompleted = "completed" + InstanceStatusCancelled = "cancelled" + InstanceStatusStale = "stale" + InstanceStatusError = "error" +) + +// ValidInstanceStatuses is the set of valid workflow instance status values. +var ValidInstanceStatuses = map[string]bool{ + InstanceStatusActive: true, + InstanceStatusCompleted: true, + InstanceStatusCancelled: true, + InstanceStatusStale: true, + InstanceStatusError: true, +} + +// ── Assignment Status Constants ───────────── + +const ( + AssignmentStatusUnassigned = "unassigned" + AssignmentStatusClaimed = "claimed" + AssignmentStatusCompleted = "completed" + AssignmentStatusCancelled = "cancelled" +) + +// ValidAssignmentStatuses is the set of valid workflow assignment status values. +var ValidAssignmentStatuses = map[string]bool{ + AssignmentStatusUnassigned: true, + AssignmentStatusClaimed: true, + AssignmentStatusCompleted: true, + AssignmentStatusCancelled: true, +} + +// ── Workflow Instance ─────────────────────── + +// WorkflowInstance tracks a single execution of a workflow definition. +type WorkflowInstance struct { + ID string `json:"id"` + WorkflowID string `json:"workflow_id"` + WorkflowVersion int `json:"workflow_version"` + CurrentStage string `json:"current_stage"` + StageData json.RawMessage `json:"stage_data"` + Status string `json:"status"` + StartedBy string `json:"started_by"` + EntryToken *string `json:"entry_token,omitempty"` + Metadata json.RawMessage `json:"metadata"` + StageEnteredAt time.Time `json:"stage_entered_at"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} + +// ── Workflow Assignment ───────────────────── + +// WorkflowAssignment tracks per-stage work items in the team queue. +type WorkflowAssignment struct { + ID string `json:"id"` + InstanceID string `json:"instance_id"` + Stage string `json:"stage"` + TeamID string `json:"team_id"` + AssignedTo *string `json:"assigned_to,omitempty"` + Status string `json:"status"` + ReviewData json.RawMessage `json:"review_data"` + ClaimedAt *time.Time `json:"claimed_at,omitempty"` + CompletedAt *time.Time `json:"completed_at,omitempty"` + CreatedAt time.Time `json:"created_at"` +} + // ── Workflow Version (immutable snapshot) ─── // WorkflowVersion is an immutable snapshot of a workflow definition diff --git a/server/store/postgres/workflows.go b/server/store/postgres/workflows.go index 7f8e4bb..f8be11b 100644 --- a/server/store/postgres/workflows.go +++ b/server/store/postgres/workflows.go @@ -6,6 +6,7 @@ import ( "fmt" "switchboard-core/models" + "switchboard-core/store" ) // WorkflowStore implements store.WorkflowStore for Postgres. @@ -396,8 +397,235 @@ func nullIfEmpty(s string) interface{} { return s } -// ── Assignments (v0.29.0-cs3) ─────────────────────────────────────────── +// ── Instances (v0.3.1) ────────────────────── -// ── Lifecycle operations (v0.37.15) ── +func (s *WorkflowStore) CreateInstance(ctx context.Context, inst *models.WorkflowInstance) error { + stageData := jsonOrEmpty(inst.StageData) + metadata := jsonOrEmpty(inst.Metadata) + status := inst.Status + if status == "" { + status = models.InstanceStatusActive + } + return DB.QueryRowContext(ctx, ` + INSERT INTO workflow_instances (workflow_id, workflow_version, current_stage, + stage_data, status, started_by, entry_token, metadata) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8) + RETURNING id, stage_entered_at, created_at, updated_at`, + inst.WorkflowID, inst.WorkflowVersion, inst.CurrentStage, + stageData, status, inst.StartedBy, inst.EntryToken, metadata, + ).Scan(&inst.ID, &inst.StageEnteredAt, &inst.CreatedAt, &inst.UpdatedAt) +} -// ── Review Comments (v0.35.0) ─────────────────────────── +func (s *WorkflowStore) GetInstance(ctx context.Context, id string) (*models.WorkflowInstance, error) { + inst := &models.WorkflowInstance{} + var stageData, metadata []byte + err := DB.QueryRowContext(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 id = $1`, id, + ).Scan(&inst.ID, &inst.WorkflowID, &inst.WorkflowVersion, &inst.CurrentStage, + &stageData, &inst.Status, &inst.StartedBy, &inst.EntryToken, &metadata, + &inst.StageEnteredAt, &inst.CreatedAt, &inst.UpdatedAt) + if err != nil { + return nil, err + } + inst.StageData = stageData + inst.Metadata = metadata + return inst, nil +} + +func (s *WorkflowStore) GetInstanceByToken(ctx context.Context, token string) (*models.WorkflowInstance, error) { + inst := &models.WorkflowInstance{} + var stageData, metadata []byte + err := DB.QueryRowContext(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 entry_token = $1`, token, + ).Scan(&inst.ID, &inst.WorkflowID, &inst.WorkflowVersion, &inst.CurrentStage, + &stageData, &inst.Status, &inst.StartedBy, &inst.EntryToken, &metadata, + &inst.StageEnteredAt, &inst.CreatedAt, &inst.UpdatedAt) + if err != nil { + return nil, err + } + inst.StageData = stageData + inst.Metadata = metadata + return inst, nil +} + +func (s *WorkflowStore) UpdateInstance(ctx context.Context, inst *models.WorkflowInstance) error { + stageData := jsonOrEmpty(inst.StageData) + metadata := jsonOrEmpty(inst.Metadata) + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_instances + SET current_stage = $2, stage_data = $3, status = $4, + entry_token = $5, metadata = $6, stage_entered_at = $7 + WHERE id = $1`, + inst.ID, inst.CurrentStage, stageData, inst.Status, + inst.EntryToken, metadata, inst.StageEnteredAt) + return err +} + +func (s *WorkflowStore) ListInstances(ctx context.Context, workflowID string, status string, opts store.ListOptions) ([]models.WorkflowInstance, error) { + q := `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 workflow_id = $1` + args := []interface{}{workflowID} + idx := 2 + if status != "" { + q += fmt.Sprintf(" AND status = $%d", idx) + args = append(args, status) + idx++ + } + q += " ORDER BY created_at DESC" + if opts.Limit > 0 { + q += fmt.Sprintf(" LIMIT $%d", idx) + args = append(args, opts.Limit) + idx++ + } + if opts.Offset > 0 { + q += fmt.Sprintf(" OFFSET $%d", idx) + args = append(args, opts.Offset) + } + rows, err := DB.QueryContext(ctx, q, args...) + 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() +} + +func (s *WorkflowStore) AdvanceStage(ctx context.Context, id string, nextStage string, stageData json.RawMessage) error { + sd := jsonOrEmpty(stageData) + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_instances + SET current_stage = $2, stage_data = $3, stage_entered_at = NOW() + WHERE id = $1 AND status = 'active'`, + id, nextStage, sd) + return err +} + +func (s *WorkflowStore) CompleteInstance(ctx context.Context, id string) error { + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_instances SET status = 'completed' WHERE id = $1 AND status = 'active'`, id) + return err +} + +func (s *WorkflowStore) CancelInstance(ctx context.Context, id string) error { + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_instances SET status = 'cancelled' WHERE id = $1 AND status = 'active'`, id) + return err +} + +// ── Assignments (v0.3.1) ──────────────────── + +func (s *WorkflowStore) CreateAssignment(ctx context.Context, a *models.WorkflowAssignment) error { + reviewData := jsonOrEmpty(a.ReviewData) + status := a.Status + if status == "" { + status = models.AssignmentStatusUnassigned + } + return DB.QueryRowContext(ctx, ` + INSERT INTO workflow_assignments (instance_id, stage, team_id, assigned_to, status, review_data) + VALUES ($1, $2, $3, $4, $5, $6) + RETURNING id, created_at`, + a.InstanceID, a.Stage, a.TeamID, a.AssignedTo, status, reviewData, + ).Scan(&a.ID, &a.CreatedAt) +} + +func (s *WorkflowStore) ClaimAssignment(ctx context.Context, id string, userID string) error { + res, err := DB.ExecContext(ctx, ` + UPDATE workflow_assignments + SET status = 'claimed', assigned_to = $2, claimed_at = NOW() + WHERE id = $1 AND status = 'unassigned'`, id, userID) + if err != nil { + return err + } + n, _ := res.RowsAffected() + if n == 0 { + return fmt.Errorf("assignment %s is not claimable", id) + } + return nil +} + +func (s *WorkflowStore) UnclaimAssignment(ctx context.Context, id string) error { + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_assignments + SET status = 'unassigned', assigned_to = NULL, claimed_at = NULL + WHERE id = $1 AND status = 'claimed'`, id) + return err +} + +func (s *WorkflowStore) CompleteAssignment(ctx context.Context, id string, reviewData json.RawMessage) error { + rd := jsonOrEmpty(reviewData) + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_assignments + SET status = 'completed', review_data = $2, completed_at = NOW() + WHERE id = $1 AND status = 'claimed'`, id, rd) + return err +} + +func (s *WorkflowStore) CancelAssignment(ctx context.Context, id string) error { + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_assignments SET status = 'cancelled' + WHERE id = $1 AND status IN ('unassigned', 'claimed')`, id) + return err +} + +func (s *WorkflowStore) ListAssignmentsByTeam(ctx context.Context, teamID string, status string) ([]models.WorkflowAssignment, error) { + q := `SELECT id, instance_id, stage, team_id, assigned_to, status, + review_data, claimed_at, completed_at, created_at + FROM workflow_assignments WHERE team_id = $1` + args := []interface{}{teamID} + if status != "" { + q += " AND status = $2" + args = append(args, status) + } + q += " ORDER BY created_at DESC" + return s.queryAssignments(ctx, q, args...) +} + +func (s *WorkflowStore) ListAssignmentsByInstance(ctx context.Context, instanceID string) ([]models.WorkflowAssignment, error) { + return s.queryAssignments(ctx, ` + SELECT id, instance_id, stage, team_id, assigned_to, status, + review_data, claimed_at, completed_at, created_at + FROM workflow_assignments WHERE instance_id = $1 + ORDER BY created_at ASC`, instanceID) +} + +func (s *WorkflowStore) queryAssignments(ctx context.Context, q string, args ...interface{}) ([]models.WorkflowAssignment, error) { + rows, err := DB.QueryContext(ctx, q, args...) + if err != nil { + return nil, err + } + defer rows.Close() + var result []models.WorkflowAssignment + for rows.Next() { + var a models.WorkflowAssignment + var reviewData []byte + if err := rows.Scan(&a.ID, &a.InstanceID, &a.Stage, &a.TeamID, + &a.AssignedTo, &a.Status, &reviewData, + &a.ClaimedAt, &a.CompletedAt, &a.CreatedAt); err != nil { + return nil, err + } + a.ReviewData = reviewData + result = append(result, a) + } + return result, rows.Err() +} diff --git a/server/store/sqlite/workflows.go b/server/store/sqlite/workflows.go index e34dca4..48ac8a5 100644 --- a/server/store/sqlite/workflows.go +++ b/server/store/sqlite/workflows.go @@ -4,6 +4,7 @@ import ( "context" "database/sql" "encoding/json" + "fmt" "time" "switchboard-core/models" @@ -409,8 +410,246 @@ func nullIfEmpty(s string) interface{} { return s } -// ── Assignments (v0.29.0-cs3) ─────────────────────────────────────────── +// ── Instances (v0.3.1) ────────────────────── -// ── Lifecycle operations (v0.37.15) ── +func (s *WorkflowStore) CreateInstance(ctx context.Context, inst *models.WorkflowInstance) error { + inst.ID = store.NewID() + now := time.Now().UTC() + inst.CreatedAt = now + inst.UpdatedAt = now + inst.StageEnteredAt = now + if inst.Status == "" { + inst.Status = models.InstanceStatusActive + } + stageData := jsonOrEmpty(inst.StageData) + metadata := jsonOrEmpty(inst.Metadata) + _, err := DB.ExecContext(ctx, ` + INSERT INTO workflow_instances (id, workflow_id, workflow_version, current_stage, + stage_data, status, started_by, entry_token, metadata, + stage_entered_at, created_at, updated_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + inst.ID, inst.WorkflowID, inst.WorkflowVersion, inst.CurrentStage, + stageData, inst.Status, inst.StartedBy, inst.EntryToken, metadata, + inst.StageEnteredAt.Format(timeFmt), inst.CreatedAt.Format(timeFmt), + inst.UpdatedAt.Format(timeFmt)) + return err +} -// ── Review Comments (v0.35.0) ─────────────────────────── +func (s *WorkflowStore) GetInstance(ctx context.Context, id string) (*models.WorkflowInstance, error) { + return s.scanInstance(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 id = ?`, id) +} + +func (s *WorkflowStore) GetInstanceByToken(ctx context.Context, token string) (*models.WorkflowInstance, error) { + return s.scanInstance(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 entry_token = ?`, token) +} + +func (s *WorkflowStore) scanInstance(ctx context.Context, q string, args ...interface{}) (*models.WorkflowInstance, error) { + inst := &models.WorkflowInstance{} + var stageData, metadata string + err := DB.QueryRowContext(ctx, q, args...).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)) + if err != nil { + return nil, err + } + inst.StageData = json.RawMessage(stageData) + inst.Metadata = json.RawMessage(metadata) + return inst, nil +} + +func (s *WorkflowStore) UpdateInstance(ctx context.Context, inst *models.WorkflowInstance) error { + stageData := jsonOrEmpty(inst.StageData) + metadata := jsonOrEmpty(inst.Metadata) + now := time.Now().UTC() + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_instances + SET current_stage = ?, stage_data = ?, status = ?, + entry_token = ?, metadata = ?, stage_entered_at = ?, + updated_at = ? + WHERE id = ?`, + inst.CurrentStage, stageData, inst.Status, + inst.EntryToken, metadata, inst.StageEnteredAt.Format(timeFmt), + now.Format(timeFmt), inst.ID) + if err == nil { + inst.UpdatedAt = now + } + return err +} + +func (s *WorkflowStore) ListInstances(ctx context.Context, workflowID string, status string, opts store.ListOptions) ([]models.WorkflowInstance, error) { + q := `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 workflow_id = ?` + args := []interface{}{workflowID} + if status != "" { + q += " AND status = ?" + args = append(args, status) + } + q += " ORDER BY created_at DESC" + if opts.Limit > 0 { + q += " LIMIT ?" + args = append(args, opts.Limit) + } + if opts.Offset > 0 { + q += " OFFSET ?" + args = append(args, opts.Offset) + } + rows, err := DB.QueryContext(ctx, q, args...) + 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() +} + +func (s *WorkflowStore) AdvanceStage(ctx context.Context, id string, nextStage string, stageData json.RawMessage) error { + sd := jsonOrEmpty(stageData) + now := time.Now().UTC() + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_instances + SET current_stage = ?, stage_data = ?, stage_entered_at = ?, updated_at = ? + WHERE id = ? AND status = 'active'`, + nextStage, sd, now.Format(timeFmt), now.Format(timeFmt), id) + return err +} + +func (s *WorkflowStore) CompleteInstance(ctx context.Context, id string) error { + now := time.Now().UTC() + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_instances SET status = 'completed', updated_at = ? + WHERE id = ? AND status = 'active'`, now.Format(timeFmt), id) + return err +} + +func (s *WorkflowStore) CancelInstance(ctx context.Context, id string) error { + now := time.Now().UTC() + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_instances SET status = 'cancelled', updated_at = ? + WHERE id = ? AND status = 'active'`, now.Format(timeFmt), id) + return err +} + +// ── Assignments (v0.3.1) ──────────────────── + +func (s *WorkflowStore) CreateAssignment(ctx context.Context, a *models.WorkflowAssignment) error { + a.ID = store.NewID() + a.CreatedAt = time.Now().UTC() + if a.Status == "" { + a.Status = models.AssignmentStatusUnassigned + } + reviewData := jsonOrEmpty(a.ReviewData) + _, err := DB.ExecContext(ctx, ` + INSERT INTO workflow_assignments (id, instance_id, stage, team_id, assigned_to, + status, review_data, claimed_at, completed_at, created_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + a.ID, a.InstanceID, a.Stage, a.TeamID, a.AssignedTo, + a.Status, reviewData, nil, nil, a.CreatedAt.Format(timeFmt)) + return err +} + +func (s *WorkflowStore) ClaimAssignment(ctx context.Context, id string, userID string) error { + now := time.Now().UTC() + res, err := DB.ExecContext(ctx, ` + UPDATE workflow_assignments + SET status = 'claimed', assigned_to = ?, claimed_at = ? + WHERE id = ? AND status = 'unassigned'`, userID, now.Format(timeFmt), id) + if err != nil { + return err + } + n, _ := res.RowsAffected() + if n == 0 { + return fmt.Errorf("assignment %s is not claimable", id) + } + return nil +} + +func (s *WorkflowStore) UnclaimAssignment(ctx context.Context, id string) error { + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_assignments + SET status = 'unassigned', assigned_to = NULL, claimed_at = NULL + WHERE id = ? AND status = 'claimed'`, id) + return err +} + +func (s *WorkflowStore) CompleteAssignment(ctx context.Context, id string, reviewData json.RawMessage) error { + rd := jsonOrEmpty(reviewData) + now := time.Now().UTC() + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_assignments + SET status = 'completed', review_data = ?, completed_at = ? + WHERE id = ? AND status = 'claimed'`, rd, now.Format(timeFmt), id) + return err +} + +func (s *WorkflowStore) CancelAssignment(ctx context.Context, id string) error { + _, err := DB.ExecContext(ctx, ` + UPDATE workflow_assignments SET status = 'cancelled' + WHERE id = ? AND status IN ('unassigned', 'claimed')`, id) + return err +} + +func (s *WorkflowStore) ListAssignmentsByTeam(ctx context.Context, teamID string, status string) ([]models.WorkflowAssignment, error) { + q := `SELECT id, instance_id, stage, team_id, assigned_to, status, + review_data, claimed_at, completed_at, created_at + FROM workflow_assignments WHERE team_id = ?` + args := []interface{}{teamID} + if status != "" { + q += " AND status = ?" + args = append(args, status) + } + q += " ORDER BY created_at DESC" + return s.queryAssignments(ctx, q, args...) +} + +func (s *WorkflowStore) ListAssignmentsByInstance(ctx context.Context, instanceID string) ([]models.WorkflowAssignment, error) { + return s.queryAssignments(ctx, ` + SELECT id, instance_id, stage, team_id, assigned_to, status, + review_data, claimed_at, completed_at, created_at + FROM workflow_assignments WHERE instance_id = ? + ORDER BY created_at ASC`, instanceID) +} + +func (s *WorkflowStore) queryAssignments(ctx context.Context, q string, args ...interface{}) ([]models.WorkflowAssignment, error) { + rows, err := DB.QueryContext(ctx, q, args...) + if err != nil { + return nil, err + } + defer rows.Close() + var result []models.WorkflowAssignment + for rows.Next() { + var a models.WorkflowAssignment + var reviewData string + if err := rows.Scan(&a.ID, &a.InstanceID, &a.Stage, &a.TeamID, + &a.AssignedTo, &a.Status, &reviewData, + stN(&a.ClaimedAt), stN(&a.CompletedAt), st(&a.CreatedAt)); err != nil { + return nil, err + } + a.ReviewData = json.RawMessage(reviewData) + result = append(result, a) + } + return result, rows.Err() +} diff --git a/server/store/workflow_iface.go b/server/store/workflow_iface.go index 5de8ed9..bc41882 100644 --- a/server/store/workflow_iface.go +++ b/server/store/workflow_iface.go @@ -2,11 +2,13 @@ package store import ( "context" + "encoding/json" "switchboard-core/models" ) -// WorkflowStore manages workflow definitions, stages, and version snapshots. +// WorkflowStore manages workflow definitions, stages, version snapshots, +// instances, and assignments. type WorkflowStore interface { // Workflow CRUD Create(ctx context.Context, w *models.Workflow) error @@ -28,4 +30,23 @@ type WorkflowStore interface { Publish(ctx context.Context, v *models.WorkflowVersion) error GetVersion(ctx context.Context, workflowID string, versionNumber int) (*models.WorkflowVersion, error) GetLatestVersion(ctx context.Context, workflowID string) (*models.WorkflowVersion, error) + + // 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) + UpdateInstance(ctx context.Context, inst *models.WorkflowInstance) error + ListInstances(ctx context.Context, workflowID string, status string, opts ListOptions) ([]models.WorkflowInstance, error) + 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 + + // Assignments (v0.3.1) + CreateAssignment(ctx context.Context, a *models.WorkflowAssignment) error + ClaimAssignment(ctx context.Context, id string, userID string) error + UnclaimAssignment(ctx context.Context, id string) error + CompleteAssignment(ctx context.Context, id string, reviewData json.RawMessage) error + CancelAssignment(ctx context.Context, id string) error + ListAssignmentsByTeam(ctx context.Context, teamID string, status string) ([]models.WorkflowAssignment, error) + ListAssignmentsByInstance(ctx context.Context, instanceID string) ([]models.WorkflowAssignment, error) }