This repository has been archived on 2026-04-03. You can view files and clone it. You cannot open issues or pull requests or push a commit.
Files
core/server/workflow/engine.go
Jeffrey Smith b8a4c516ff
All checks were successful
CI/CD / detect-changes (pull_request) Successful in 4s
CI/CD / test-frontend (pull_request) Has been skipped
CI/CD / test-runners (pull_request) Has been skipped
CI/CD / e2e-smoke (pull_request) Has been skipped
CI/CD / test-go-pg (pull_request) Successful in 2m43s
CI/CD / test-sqlite (pull_request) Successful in 2m52s
CI/CD / build-and-deploy (pull_request) Successful in 2m15s
Feat v0.9.2 Starlark converter consolidation + snapshot cleanup (#75)
Consolidate duplicate Go↔Starlark converters into sandbox/convert.go
and snapshot parsers into models/snapshot.go. Standardize snapshot
creation on wrapped format. Net -393 lines across 21 files.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-04-03 14:20:05 +00:00

536 lines
17 KiB
Go

package workflow
import (
"context"
"encoding/json"
"fmt"
"log"
"time"
"armature/events"
"armature/models"
"armature/sandbox"
"armature/store"
"github.com/google/uuid"
)
// MaxConsecutiveAutomated is the cycle guard limit for automated stages.
const MaxConsecutiveAutomated = 10
// Engine orchestrates workflow instance lifecycle.
type Engine struct {
stores store.Stores
bus *events.Bus
runner *sandbox.Runner
}
// NewEngine creates a workflow engine.
func NewEngine(stores store.Stores, bus *events.Bus, runner *sandbox.Runner) *Engine {
return &Engine{stores: stores, bus: bus, runner: runner}
}
// Start creates a new workflow instance and kicks off the first stage.
func (e *Engine) Start(ctx context.Context, workflowID string, initialData json.RawMessage, userID string) (*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.IsActive {
return nil, fmt.Errorf("workflow %q is not active", wf.Name)
}
ver, err := e.stores.Workflows.GetLatestVersion(ctx, workflowID)
if err != nil {
return nil, fmt.Errorf("no published version: %w", err)
}
stages, err := models.ParseSnapshotStages(ver.Snapshot)
if err != nil {
return nil, err
}
if len(stages) == 0 {
return nil, fmt.Errorf("workflow has no stages")
}
inst := &models.WorkflowInstance{
WorkflowID: workflowID,
WorkflowVersion: ver.VersionNumber,
CurrentStage: stages[0].Name,
StageData: initialData,
Status: models.InstanceStatusActive,
StartedBy: userID,
Metadata: json.RawMessage(`{}`),
}
if wf.EntryMode == "public_link" {
tok := uuid.New().String()
inst.EntryToken = &tok
}
if err := e.stores.Workflows.CreateInstance(ctx, inst); err != nil {
return nil, fmt.Errorf("create instance: %w", err)
}
// Create initial assignment if first stage has a team
if stages[0].AssignmentTeamID != nil {
a := &models.WorkflowAssignment{
InstanceID: inst.ID,
Stage: stages[0].Name,
TeamID: *stages[0].AssignmentTeamID,
}
if err := e.stores.Workflows.CreateAssignment(ctx, a); err != nil {
log.Printf("[workflow-engine] create initial assignment: %v", err)
} else {
e.notifyAssignment(ctx, a, wf.Name, stages[0].Name)
}
}
e.emit(ctx, "workflow.started", inst.ID, map[string]any{
"workflow_id": workflowID,
"instance_id": inst.ID,
"started_by": userID,
})
// If first stage is automated, process it
if stages[0].StageMode == models.StageModeAutomated {
if err := e.processAutomatedStage(ctx, inst, stages, 0); err != nil {
log.Printf("[workflow-engine] automated stage error on start: %v", err)
}
// Re-read instance after automated processing
inst, _ = e.stores.Workflows.GetInstance(ctx, inst.ID)
}
return inst, nil
}
// Advance moves an instance to the next stage.
func (e *Engine) Advance(ctx context.Context, instanceID string, stageData json.RawMessage, userID string) (*models.WorkflowInstance, error) {
return e.advanceInternal(ctx, instanceID, stageData, userID, 0)
}
func (e *Engine) advanceInternal(ctx context.Context, instanceID string, stageData json.RawMessage, userID string, autoDepth int) (*models.WorkflowInstance, error) {
inst, err := e.stores.Workflows.GetInstance(ctx, instanceID)
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)
}
ver, err := e.stores.Workflows.GetVersion(ctx, inst.WorkflowID, inst.WorkflowVersion)
if err != nil {
return nil, fmt.Errorf("version not found: %w", err)
}
stages, err := models.ParseSnapshotStages(ver.Snapshot)
if err != nil {
return nil, err
}
// Find current stage ordinal
currentOrdinal := -1
for _, s := range stages {
if s.Name == inst.CurrentStage {
currentOrdinal = s.Ordinal
break
}
}
if currentOrdinal < 0 {
return nil, fmt.Errorf("current stage %q not found in snapshot", inst.CurrentStage)
}
currentStage := stages[currentOrdinal]
// ── Validation gate ──────────────
sc := ParseStageConfig(currentStage.StageConfig)
if sc.Validation != nil && sc.Validation.RequiredApprovals > 0 {
// Check for rejections first
rejectCount, _ := e.stores.Workflows.CountSignoffs(ctx, instanceID, currentStage.Name, models.SignoffReject)
if rejectCount > 0 {
action := sc.Validation.RejectAction
if action == "" || action == "cancel" {
_ = e.Cancel(ctx, instanceID, userID)
e.emit(ctx, "workflow.rejected", instanceID, map[string]any{
"instance_id": instanceID,
"stage": currentStage.Name,
})
return e.stores.Workflows.GetInstance(ctx, instanceID)
}
// Reroute to named stage
target, rerr := ResolveStageByName(stages, action)
if rerr != nil {
return nil, fmt.Errorf("reject_action stage %q not found: %w", action, rerr)
}
merged := mergeJSON(inst.StageData, stageData)
if err := e.stores.Workflows.AdvanceStage(ctx, instanceID, stages[target].Name, merged); err != nil {
return nil, fmt.Errorf("reject reroute: %w", err)
}
e.emit(ctx, "workflow.rejected", instanceID, map[string]any{
"instance_id": instanceID,
"stage": currentStage.Name,
"rerouted_to": stages[target].Name,
})
return e.stores.Workflows.GetInstance(ctx, instanceID)
}
approveCount, _ := e.stores.Workflows.CountSignoffs(ctx, instanceID, currentStage.Name, models.SignoffApprove)
if approveCount < sc.Validation.RequiredApprovals {
return nil, fmt.Errorf("insufficient approvals (%d of %d required)", approveCount, sc.Validation.RequiredApprovals)
}
}
// Merge stage data
merged := mergeJSON(inst.StageData, stageData)
// Resolve next stage
nextOrdinal, err := ResolveNextStage(stages, currentOrdinal, merged)
if err != nil {
return nil, fmt.Errorf("resolve next stage: %w", err)
}
// Terminal — complete the instance
if nextOrdinal >= len(stages) {
if err := e.stores.Workflows.CompleteInstance(ctx, instanceID); err != nil {
return nil, fmt.Errorf("complete: %w", err)
}
e.emit(ctx, "workflow.completed", instanceID, map[string]any{
"instance_id": instanceID,
})
return e.stores.Workflows.GetInstance(ctx, instanceID)
}
nextStage := stages[nextOrdinal]
// Advance in store
if err := e.stores.Workflows.AdvanceStage(ctx, instanceID, nextStage.Name, merged); err != nil {
return nil, fmt.Errorf("advance store: %w", err)
}
// Cancel old assignments for the previous stage
oldAssignments, _ := e.stores.Workflows.ListAssignmentsByInstance(ctx, instanceID)
for _, a := range oldAssignments {
if a.Stage == currentStage.Name && (a.Status == models.AssignmentStatusUnassigned || a.Status == models.AssignmentStatusClaimed) {
e.stores.Workflows.CancelAssignment(ctx, a.ID)
}
}
// Create new assignment if next stage has a team
if nextStage.AssignmentTeamID != nil {
a := &models.WorkflowAssignment{
InstanceID: instanceID,
Stage: nextStage.Name,
TeamID: *nextStage.AssignmentTeamID,
}
if err := e.stores.Workflows.CreateAssignment(ctx, a); err != nil {
log.Printf("[workflow-engine] create assignment: %v", err)
} else {
e.emit(ctx, "workflow.assigned", a.TeamID, map[string]any{
"instance_id": instanceID,
"stage": nextStage.Name,
"assignment_id": a.ID,
"team_id": a.TeamID,
})
// Notify team members
wfTitle := ""
if wf, _ := e.stores.Workflows.GetByID(ctx, inst.WorkflowID); wf != nil {
wfTitle = wf.Name
}
e.notifyAssignment(ctx, a, wfTitle, nextStage.Name)
}
}
e.emit(ctx, "workflow.advanced", instanceID, map[string]any{
"instance_id": instanceID,
"previous_stage": currentStage.Name,
"current_stage": nextStage.Name,
})
// If next stage is automated, process it
if nextStage.StageMode == models.StageModeAutomated {
if err := e.processAutomatedStage(ctx, inst, stages, autoDepth+1); err != nil {
log.Printf("[workflow-engine] automated stage error: %v", err)
}
}
return e.stores.Workflows.GetInstance(ctx, instanceID)
}
// Cancel terminates an active instance and all its open assignments.
func (e *Engine) Cancel(ctx context.Context, instanceID string, userID string) error {
inst, err := e.stores.Workflows.GetInstance(ctx, instanceID)
if err != nil {
return fmt.Errorf("instance not found: %w", err)
}
if inst.Status != models.InstanceStatusActive {
return fmt.Errorf("instance is %s, not active", inst.Status)
}
if err := e.stores.Workflows.CancelInstance(ctx, instanceID); err != nil {
return fmt.Errorf("cancel: %w", err)
}
// Cancel all open assignments
assignments, _ := e.stores.Workflows.ListAssignmentsByInstance(ctx, instanceID)
for _, a := range assignments {
if a.Status == models.AssignmentStatusUnassigned || a.Status == models.AssignmentStatusClaimed {
e.stores.Workflows.CancelAssignment(ctx, a.ID)
}
}
e.emit(ctx, "workflow.cancelled", instanceID, map[string]any{
"instance_id": instanceID,
"cancelled_by": userID,
})
return nil
}
// ── Public Entry ───────────────────
// StartPublic creates an instance for a public_link workflow without authentication.
// The caller is identified as "public:<uuid>". 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)
}
stages, err := models.ParseSnapshotStages(ver.Snapshot)
if err != nil {
return nil, 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)
}
// ── Signoffs ─────────────────────────
// SubmitSignoff records a user's approval or rejection at the current stage boundary.
func (e *Engine) SubmitSignoff(ctx context.Context, instanceID, userID, decision, comment string) (*models.WorkflowSignoff, error) {
inst, err := e.stores.Workflows.GetInstance(ctx, instanceID)
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 get current stage config
ver, err := e.stores.Workflows.GetVersion(ctx, inst.WorkflowID, inst.WorkflowVersion)
if err != nil {
return nil, fmt.Errorf("version not found: %w", err)
}
stages, err := models.ParseSnapshotStages(ver.Snapshot)
if err != nil {
return nil, err
}
// Find current stage
var currentStage *models.WorkflowStage
for i := range stages {
if stages[i].Name == inst.CurrentStage {
currentStage = &stages[i]
break
}
}
if currentStage == nil {
return nil, fmt.Errorf("current stage %q not found", inst.CurrentStage)
}
sc := ParseStageConfig(currentStage.StageConfig)
if sc.Validation == nil || sc.Validation.RequiredApprovals == 0 {
return nil, fmt.Errorf("stage %q does not require signoffs", currentStage.Name)
}
// Role check: if validation has a required_role, verify the user has it
if sc.Validation.RequiredRole != "" && currentStage.AssignmentTeamID != nil {
member, merr := e.stores.Teams.GetMember(ctx, *currentStage.AssignmentTeamID, userID)
if merr != nil {
return nil, fmt.Errorf("user is not a member of the assignment team")
}
if member.Role != sc.Validation.RequiredRole {
return nil, fmt.Errorf("role %q required to sign off (you have %q)", sc.Validation.RequiredRole, member.Role)
}
}
so := &models.WorkflowSignoff{
InstanceID: instanceID,
Stage: inst.CurrentStage,
UserID: userID,
Decision: decision,
Comment: comment,
}
if err := e.stores.Workflows.CreateSignoff(ctx, so); err != nil {
return nil, fmt.Errorf("create signoff: %w", err)
}
e.emit(ctx, "workflow.signoff", instanceID, map[string]any{
"instance_id": instanceID,
"stage": inst.CurrentStage,
"user_id": userID,
"decision": decision,
})
return so, nil
}
// CheckClaimRole verifies that a user has the required_role (from stage_config)
// to claim an assignment. Returns nil if allowed, error if forbidden.
func CheckClaimRole(ctx context.Context, stores store.Stores, assignment *models.WorkflowAssignment, userID string) error {
inst, err := stores.Workflows.GetInstance(ctx, assignment.InstanceID)
if err != nil {
return nil // can't verify — allow claim
}
ver, err := stores.Workflows.GetVersion(ctx, inst.WorkflowID, inst.WorkflowVersion)
if err != nil {
return nil
}
stages, parseErr := models.ParseSnapshotStages(ver.Snapshot)
if parseErr != nil {
return nil
}
for _, s := range stages {
if s.Name == assignment.Stage {
sc := ParseStageConfig(s.StageConfig)
if sc.RequiredRole == "" {
return nil // no role restriction
}
member, merr := stores.Teams.GetMember(ctx, assignment.TeamID, userID)
if merr != nil {
return fmt.Errorf("not a member of team")
}
if member.Role != sc.RequiredRole {
return fmt.Errorf("role %q required for this stage (you have %q)", sc.RequiredRole, member.Role)
}
return nil
}
}
return nil // stage not found in snapshot — allow
}
// emit publishes an event on the bus.
func (e *Engine) emit(_ context.Context, label, room string, payload map[string]any) {
if e.bus == nil {
return
}
e.bus.Publish(events.Event{
Label: label,
Room: room,
Payload: events.MustJSON(payload),
Ts: time.Now().UnixMilli(),
})
}
// notifyAssignment creates notifications for team members when a workflow
// assignment is created. If the assignment targets a specific user, only
// that user is notified. Otherwise all team members are notified.
func (e *Engine) notifyAssignment(ctx context.Context, a *models.WorkflowAssignment, workflowTitle, stageName string) {
if e.stores.Notifications == nil || e.stores.Teams == nil {
return
}
var recipients []string
if a.AssignedTo != nil && *a.AssignedTo != "" {
recipients = []string{*a.AssignedTo}
} else {
members, err := e.stores.Teams.ListMembers(ctx, a.TeamID)
if err != nil {
log.Printf("[workflow-engine] list team members for notification: %v", err)
return
}
for _, m := range members {
recipients = append(recipients, m.UserID)
}
}
title := "Workflow assignment: " + stageName
body := workflowTitle
for _, uid := range recipients {
n := &models.Notification{
UserID: uid,
Type: models.NotifTypeWorkflowAssign,
Title: title,
Body: body,
ResourceType: "workflow",
ResourceID: a.InstanceID,
}
if err := e.stores.Notifications.Create(ctx, n); err != nil {
log.Printf("[workflow-engine] create notification: %v", err)
}
}
}
// mergeJSON merges b into a (shallow). Returns a if b is empty.
func mergeJSON(a, b json.RawMessage) json.RawMessage {
if len(b) == 0 || string(b) == "{}" || string(b) == "null" {
if len(a) == 0 {
return json.RawMessage(`{}`)
}
return a
}
if len(a) == 0 || string(a) == "{}" || string(a) == "null" {
return b
}
var ma, mb map[string]json.RawMessage
if json.Unmarshal(a, &ma) != nil {
return b
}
if json.Unmarshal(b, &mb) != nil {
return a
}
for k, v := range mb {
ma[k] = v
}
merged, err := json.Marshal(ma)
if err != nil {
return b
}
return merged
}