package sandbox // workflow_module.go // // The workflow module lets extensions read workflow definitions, // manage instances, and programmatically advance workflow stages. // // Starlark API: // wf = workflow.get_definition(workflow_id) # returns dict // inst = workflow.get_instance(instance_id) # returns dict // instances = workflow.list_instances(workflow_id) # list instances // inst = workflow.start(workflow_id, data={}) # start new instance // inst = workflow.advance(instance_id, data={}) # advance to next stage // workflow.cancel(instance_id) # cancel instance // signoff = workflow.submit_signoff(instance_id, decision, comment="") # submit signoff import ( "context" "encoding/json" "fmt" "go.starlark.net/starlark" "go.starlark.net/starlarkstruct" "armature/models" "armature/store" ) // WorkflowEngine is the subset of the workflow engine needed by Starlark // write operations. Defined here (not in the workflow package) to avoid // circular imports — workflow imports sandbox for the Runner. type WorkflowEngine interface { Start(ctx context.Context, workflowID string, initialData json.RawMessage, userID string) (*models.WorkflowInstance, error) Advance(ctx context.Context, instanceID string, stageData json.RawMessage, userID string) (*models.WorkflowInstance, error) Cancel(ctx context.Context, instanceID string, userID string) error SubmitSignoff(ctx context.Context, instanceID, userID, decision, comment string) (*models.WorkflowSignoff, error) } // BuildWorkflowModule creates the "workflow" Starlark module for a package. // Requires the workflow.access permission. // // Read operations use stores directly. Write operations delegate to the // WorkflowEngine interface (engine may be nil if not wired — write builtins // return an error in that case). func BuildWorkflowModule(ctx context.Context, stores store.Stores, engine WorkflowEngine, rc *RunContext) *starlarkstruct.Module { members := starlark.StringDict{ "get_definition": starlark.NewBuiltin("workflow.get_definition", workflowGetDef(ctx, stores)), "get_instance": starlark.NewBuiltin("workflow.get_instance", workflowGetInstance(ctx, stores)), "list_instances": starlark.NewBuiltin("workflow.list_instances", workflowListInstances(ctx, stores)), "start": starlark.NewBuiltin("workflow.start", workflowStart(ctx, engine, rc)), "advance": starlark.NewBuiltin("workflow.advance", workflowAdvance(ctx, engine, rc)), "cancel": starlark.NewBuiltin("workflow.cancel", workflowCancel(ctx, engine, rc)), "submit_signoff": starlark.NewBuiltin("workflow.submit_signoff", workflowSubmitSignoff(ctx, engine, rc)), } return MakeModule("workflow", members) } func workflowGetDef(ctx context.Context, stores store.Stores) func(*starlark.Thread, *starlark.Builtin, starlark.Tuple, []starlark.Tuple) (starlark.Value, error) { return func(_ *starlark.Thread, _ *starlark.Builtin, args starlark.Tuple, kwargs []starlark.Tuple) (starlark.Value, error) { var workflowID string if err := starlark.UnpackPositionalArgs("workflow.get_definition", args, kwargs, 1, &workflowID); err != nil { return nil, err } wf, err := stores.Workflows.GetByID(ctx, workflowID) if err != nil { return starlark.None, fmt.Errorf("workflow.get_definition: %w", err) } stages, _ := stores.Workflows.ListStages(ctx, workflowID) // Build stages list stageList := make([]starlark.Value, 0, len(stages)) for _, s := range stages { d := starlark.NewDict(8) d.SetKey(starlark.String("id"), starlark.String(s.ID)) d.SetKey(starlark.String("name"), starlark.String(s.Name)) d.SetKey(starlark.String("ordinal"), starlark.MakeInt(s.Ordinal)) d.SetKey(starlark.String("stage_mode"), starlark.String(s.StageMode)) d.SetKey(starlark.String("audience"), starlark.String(s.Audience)) d.SetKey(starlark.String("stage_type"), starlark.String(s.StageType)) // deprecated — kept for backward compat d.SetKey(starlark.String("auto_transition"), starlark.Bool(s.AutoTransition)) if s.StarlarkHook != nil { d.SetKey(starlark.String("starlark_hook"), starlark.String(*s.StarlarkHook)) } if s.AssignmentTeamID != nil { d.SetKey(starlark.String("assignment_team_id"), starlark.String(*s.AssignmentTeamID)) } if s.SurfacePkgID != nil { d.SetKey(starlark.String("surface_pkg_id"), starlark.String(*s.SurfacePkgID)) } stageList = append(stageList, d) } result := starlark.NewDict(8) result.SetKey(starlark.String("id"), starlark.String(wf.ID)) result.SetKey(starlark.String("name"), starlark.String(wf.Name)) result.SetKey(starlark.String("slug"), starlark.String(wf.Slug)) result.SetKey(starlark.String("entry_mode"), starlark.String(wf.EntryMode)) result.SetKey(starlark.String("is_active"), starlark.Bool(wf.IsActive)) result.SetKey(starlark.String("version"), starlark.MakeInt(wf.Version)) result.SetKey(starlark.String("stages"), starlark.NewList(stageList)) return result, nil } } // ── Instance Read API ────────────── func workflowGetInstance(ctx context.Context, stores store.Stores) func(*starlark.Thread, *starlark.Builtin, starlark.Tuple, []starlark.Tuple) (starlark.Value, error) { return func(_ *starlark.Thread, _ *starlark.Builtin, args starlark.Tuple, kwargs []starlark.Tuple) (starlark.Value, error) { var instanceID string if err := starlark.UnpackPositionalArgs("workflow.get_instance", args, kwargs, 1, &instanceID); err != nil { return nil, err } inst, err := stores.Workflows.GetInstance(ctx, instanceID) if err != nil { return starlark.None, fmt.Errorf("workflow.get_instance: %w", err) } return instanceToDict(inst), nil } } // instanceToDict converts a WorkflowInstance to a Starlark dict. func instanceToDict(inst *models.WorkflowInstance) *starlark.Dict { d := starlark.NewDict(8) d.SetKey(starlark.String("id"), starlark.String(inst.ID)) d.SetKey(starlark.String("workflow_id"), starlark.String(inst.WorkflowID)) d.SetKey(starlark.String("workflow_version"), starlark.MakeInt(inst.WorkflowVersion)) d.SetKey(starlark.String("current_stage"), starlark.String(inst.CurrentStage)) d.SetKey(starlark.String("status"), starlark.String(inst.Status)) d.SetKey(starlark.String("started_by"), starlark.String(inst.StartedBy)) // Parse stage_data into Starlark dict var dataMap map[string]interface{} if json.Unmarshal(inst.StageData, &dataMap) == nil { starlarkData := starlark.NewDict(len(dataMap)) for k, v := range dataMap { starlarkData.SetKey(starlark.String(k), GoToStarlark(v)) } d.SetKey(starlark.String("stage_data"), starlarkData) } else { d.SetKey(starlark.String("stage_data"), starlark.NewDict(0)) } if inst.EntryToken != nil { d.SetKey(starlark.String("entry_token"), starlark.String(*inst.EntryToken)) } return d } // signoffToDict converts a WorkflowSignoff to a Starlark dict. func signoffToDict(s *models.WorkflowSignoff) *starlark.Dict { d := starlark.NewDict(7) d.SetKey(starlark.String("id"), starlark.String(s.ID)) d.SetKey(starlark.String("instance_id"), starlark.String(s.InstanceID)) d.SetKey(starlark.String("stage"), starlark.String(s.Stage)) d.SetKey(starlark.String("user_id"), starlark.String(s.UserID)) d.SetKey(starlark.String("decision"), starlark.String(s.Decision)) d.SetKey(starlark.String("comment"), starlark.String(s.Comment)) d.SetKey(starlark.String("created_at"), starlark.String(s.CreatedAt.Format("2006-01-02T15:04:05Z07:00"))) return d } func workflowListInstances(ctx context.Context, stores store.Stores) func(*starlark.Thread, *starlark.Builtin, starlark.Tuple, []starlark.Tuple) (starlark.Value, error) { return func(_ *starlark.Thread, _ *starlark.Builtin, args starlark.Tuple, kwargs []starlark.Tuple) (starlark.Value, error) { var workflowID string var status string if err := starlark.UnpackPositionalArgs("workflow.list_instances", args, kwargs, 1, &workflowID, &status); err != nil { return nil, err } instances, err := stores.Workflows.ListInstances(ctx, workflowID, status, store.ListOptions{Limit: 100}) if err != nil { return starlark.None, fmt.Errorf("workflow.list_instances: %w", err) } result := make([]starlark.Value, 0, len(instances)) for _, inst := range instances { d := starlark.NewDict(6) d.SetKey(starlark.String("id"), starlark.String(inst.ID)) d.SetKey(starlark.String("current_stage"), starlark.String(inst.CurrentStage)) d.SetKey(starlark.String("status"), starlark.String(inst.Status)) d.SetKey(starlark.String("started_by"), starlark.String(inst.StartedBy)) d.SetKey(starlark.String("workflow_version"), starlark.MakeInt(inst.WorkflowVersion)) result = append(result, d) } return starlark.NewList(result), nil } } // ── Instance Write API ────────────── // workflowWriteGuard checks that the engine and user context are available. func workflowWriteGuard(engine WorkflowEngine, rc *RunContext) error { if engine == nil { return fmt.Errorf("workflow engine not available") } if rc == nil || rc.UserID == "" { return fmt.Errorf("workflow write operations require an authenticated user context") } return nil } // dictToJSON converts a Starlark dict to json.RawMessage. func dictToJSON(d *starlark.Dict) (json.RawMessage, error) { m := DictToMap(d) b, err := json.Marshal(m) if err != nil { return nil, err } return json.RawMessage(b), nil } func workflowStart(ctx context.Context, engine WorkflowEngine, rc *RunContext) func(*starlark.Thread, *starlark.Builtin, starlark.Tuple, []starlark.Tuple) (starlark.Value, error) { return func(_ *starlark.Thread, _ *starlark.Builtin, args starlark.Tuple, kwargs []starlark.Tuple) (starlark.Value, error) { var workflowID string data := starlark.NewDict(0) if err := starlark.UnpackArgs("workflow.start", args, kwargs, "workflow_id", &workflowID, "data?", &data, ); err != nil { return nil, err } if err := workflowWriteGuard(engine, rc); err != nil { return nil, fmt.Errorf("workflow.start: %w", err) } dataJSON, err := dictToJSON(data) if err != nil { return nil, fmt.Errorf("workflow.start: marshal data: %w", err) } inst, err := engine.Start(ctx, workflowID, dataJSON, rc.UserID) if err != nil { return nil, fmt.Errorf("workflow.start: %w", err) } return instanceToDict(inst), nil } } func workflowAdvance(ctx context.Context, engine WorkflowEngine, rc *RunContext) func(*starlark.Thread, *starlark.Builtin, starlark.Tuple, []starlark.Tuple) (starlark.Value, error) { return func(_ *starlark.Thread, _ *starlark.Builtin, args starlark.Tuple, kwargs []starlark.Tuple) (starlark.Value, error) { var instanceID string data := starlark.NewDict(0) if err := starlark.UnpackArgs("workflow.advance", args, kwargs, "instance_id", &instanceID, "data?", &data, ); err != nil { return nil, err } if err := workflowWriteGuard(engine, rc); err != nil { return nil, fmt.Errorf("workflow.advance: %w", err) } dataJSON, err := dictToJSON(data) if err != nil { return nil, fmt.Errorf("workflow.advance: marshal data: %w", err) } inst, err := engine.Advance(ctx, instanceID, dataJSON, rc.UserID) if err != nil { return nil, fmt.Errorf("workflow.advance: %w", err) } return instanceToDict(inst), nil } } func workflowCancel(ctx context.Context, engine WorkflowEngine, rc *RunContext) func(*starlark.Thread, *starlark.Builtin, starlark.Tuple, []starlark.Tuple) (starlark.Value, error) { return func(_ *starlark.Thread, _ *starlark.Builtin, args starlark.Tuple, kwargs []starlark.Tuple) (starlark.Value, error) { var instanceID string if err := starlark.UnpackPositionalArgs("workflow.cancel", args, kwargs, 1, &instanceID); err != nil { return nil, err } if err := workflowWriteGuard(engine, rc); err != nil { return nil, fmt.Errorf("workflow.cancel: %w", err) } if err := engine.Cancel(ctx, instanceID, rc.UserID); err != nil { return nil, fmt.Errorf("workflow.cancel: %w", err) } return starlark.None, nil } } func workflowSubmitSignoff(ctx context.Context, engine WorkflowEngine, rc *RunContext) func(*starlark.Thread, *starlark.Builtin, starlark.Tuple, []starlark.Tuple) (starlark.Value, error) { return func(_ *starlark.Thread, _ *starlark.Builtin, args starlark.Tuple, kwargs []starlark.Tuple) (starlark.Value, error) { var instanceID, decision string comment := "" if err := starlark.UnpackArgs("workflow.submit_signoff", args, kwargs, "instance_id", &instanceID, "decision", &decision, "comment?", &comment, ); err != nil { return nil, err } if err := workflowWriteGuard(engine, rc); err != nil { return nil, fmt.Errorf("workflow.submit_signoff: %w", err) } signoff, err := engine.SubmitSignoff(ctx, instanceID, rc.UserID, decision, comment) if err != nil { return nil, fmt.Errorf("workflow.submit_signoff: %w", err) } return signoffToDict(signoff), nil } }