Feat v0.3.1 instance + assignment tables and store
Some checks failed
CI/CD / detect-changes (pull_request) Successful in 5s
CI/CD / test-frontend (pull_request) Has been skipped
CI/CD / test-go-pg (pull_request) Failing after 2m32s
CI/CD / test-sqlite (pull_request) Successful in 2m58s
CI/CD / build-and-deploy (pull_request) Has been skipped

Add workflow execution layer: workflow_instances tracks running workflows
(stage, status, entry token), workflow_assignments tracks per-stage team
queue with optimistic-lock claim. 15 new store methods in PG + SQLite.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
2026-03-27 18:41:03 +00:00
parent 965885a8f7
commit 1bf8422932
6 changed files with 667 additions and 7 deletions

View File

@@ -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()
}

View File

@@ -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()
}

View File

@@ -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)
}