Feat v0.3.1 instance assignment (#15)
Co-authored-by: Jeffrey Smith <jasafpro@gmail.com> Co-committed-by: Jeffrey Smith <jasafpro@gmail.com>
This commit was merged in pull request #15.
This commit is contained in:
@@ -78,3 +78,55 @@ CREATE TABLE IF NOT EXISTS workflow_versions (
|
|||||||
|
|
||||||
CREATE INDEX IF NOT EXISTS idx_workflow_versions_workflow
|
CREATE INDEX IF NOT EXISTS idx_workflow_versions_workflow
|
||||||
ON workflow_versions(workflow_id, version_number DESC);
|
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;
|
||||||
|
|||||||
@@ -69,3 +69,53 @@ CREATE TABLE IF NOT EXISTS workflow_versions (
|
|||||||
|
|
||||||
CREATE INDEX IF NOT EXISTS idx_workflow_versions_workflow
|
CREATE INDEX IF NOT EXISTS idx_workflow_versions_workflow
|
||||||
ON workflow_versions(workflow_id, version_number);
|
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);
|
||||||
|
|||||||
@@ -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) ───
|
// ── Workflow Version (immutable snapshot) ───
|
||||||
|
|
||||||
// WorkflowVersion is an immutable snapshot of a workflow definition
|
// WorkflowVersion is an immutable snapshot of a workflow definition
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
|
|
||||||
"switchboard-core/models"
|
"switchboard-core/models"
|
||||||
|
"switchboard-core/store"
|
||||||
)
|
)
|
||||||
|
|
||||||
// WorkflowStore implements store.WorkflowStore for Postgres.
|
// WorkflowStore implements store.WorkflowStore for Postgres.
|
||||||
@@ -396,8 +397,235 @@ func nullIfEmpty(s string) interface{} {
|
|||||||
return s
|
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()
|
||||||
|
}
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"database/sql"
|
"database/sql"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"switchboard-core/models"
|
"switchboard-core/models"
|
||||||
@@ -409,8 +410,246 @@ func nullIfEmpty(s string) interface{} {
|
|||||||
return s
|
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()
|
||||||
|
}
|
||||||
|
|||||||
@@ -2,11 +2,13 @@ package store
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
|
||||||
"switchboard-core/models"
|
"switchboard-core/models"
|
||||||
)
|
)
|
||||||
|
|
||||||
// WorkflowStore manages workflow definitions, stages, and version snapshots.
|
// WorkflowStore manages workflow definitions, stages, version snapshots,
|
||||||
|
// instances, and assignments.
|
||||||
type WorkflowStore interface {
|
type WorkflowStore interface {
|
||||||
// Workflow CRUD
|
// Workflow CRUD
|
||||||
Create(ctx context.Context, w *models.Workflow) error
|
Create(ctx context.Context, w *models.Workflow) error
|
||||||
@@ -28,4 +30,23 @@ type WorkflowStore interface {
|
|||||||
Publish(ctx context.Context, v *models.WorkflowVersion) error
|
Publish(ctx context.Context, v *models.WorkflowVersion) error
|
||||||
GetVersion(ctx context.Context, workflowID string, versionNumber int) (*models.WorkflowVersion, error)
|
GetVersion(ctx context.Context, workflowID string, versionNumber int) (*models.WorkflowVersion, error)
|
||||||
GetLatestVersion(ctx context.Context, workflowID string) (*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)
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user