workflow

package
v0.3.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 18, 2026 License: GPL-3.0 Imports: 23 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func DownstreamOf

func DownstreamOf(snapshot *Snapshot, target string) map[string]bool

DownstreamOf returns the target node plus every node that transitively consumes its outputs — the set that must re-run when resuming a workflow from that node.

func ResolveHomeDir

func ResolveHomeDir() (string, error)

ResolveHomeDir returns the agc home directory (AGENT_COMPOSER_HOME, or ~/.agent_composer) — where the database, config, and settings live.

func StampWorkflowID

func StampWorkflowID(raw []byte, workflowID string) ([]byte, error)

StampWorkflowID forces workflow.id in a spec's bytes — used to carry a workflow's permanent identity into a proposal that dropped or fabricated it.

func ValidateWorkflowID

func ValidateWorkflowID(id string) error

ValidateWorkflowID accepts lowercase slug ids: letters and digits separated by single underscores or hyphens.

Types

type Binding

type Binding struct {
	Kind          BindingKind
	WorkflowInput string
	InstanceID    string
	OutputName    string
}

type BindingKind

type BindingKind string
const (
	BindingKindWorkflowInput BindingKind = "workflow_input"
	BindingKindInstance      BindingKind = "instance"
)

type ComposeOptions

type ComposeOptions struct {
	// WorkflowSlug is empty when the request should create a workflow.
	WorkflowSlug string
	// BaseSpec is the spec being edited — the draft when one
	// exists, else the saved file. Empty for a create.
	BaseSpec string
	Request  string
	Harness  agent.Harness
	Model    string
	// ReasoningEffort is medium when empty.
	ReasoningEffort runtimetypes.ReasoningEffort
	// Catalog lists installed harnesses and their real model ids, so
	// the agent never invents a model name.
	Catalog string
}

type ComposeResult

type ComposeResult struct {
	WorkflowSlug string
	Action       string
	YAML         string
	Summary      string
}

func Compose

func Compose(ctx context.Context, opts ComposeOptions) (*ComposeResult, error)

Compose runs one composer conversation: the agent edits the registry through the agc CLI and reports what it did. The working directory is the agc home, so a workspace-scoped harness sandbox still reaches the registry.

type CreatedDraft

type CreatedDraft struct {
	WorkflowSlug string
	Spec         string
}

type DBRecorder

type DBRecorder struct {
	// contains filtered or unexported fields
}

func NewDBRecorder

func NewDBRecorder(db bun.IDB) *DBRecorder

func (*DBRecorder) FinishConversation

func (r *DBRecorder) FinishConversation(ctx context.Context, conversation *agent.Conversation, output any) error

func (*DBRecorder) FinishNode

func (r *DBRecorder) FinishNode(ctx context.Context, handle NodeExecutionHandle, output map[string]any, status executionmodels.NodeExecutionStatus, trace map[string]any) error

func (*DBRecorder) FinishWorkflow

func (r *DBRecorder) FinishWorkflow(ctx context.Context, handle WorkflowExecutionHandle, output map[string]any, status executionmodels.WorkflowExecutionStatus) error

func (*DBRecorder) StartConversation

func (r *DBRecorder) StartConversation(ctx context.Context, handle NodeExecutionHandle, conversation *agent.Conversation, input map[string]any) error

func (*DBRecorder) StartNode

func (*DBRecorder) StartWorkflow

func (r *DBRecorder) StartWorkflow(ctx context.Context, snapshot *Snapshot, input map[string]any, project string) (WorkflowExecutionHandle, error)

type ExecutionRecorder

type ExecutionRecorder interface {
	StartWorkflow(ctx context.Context, snapshot *Snapshot, input map[string]any, project string) (WorkflowExecutionHandle, error)
	FinishWorkflow(ctx context.Context, handle WorkflowExecutionHandle, output map[string]any, status executionmodels.WorkflowExecutionStatus) error
	StartNode(ctx context.Context, workflow WorkflowExecutionHandle, node NodeSnapshot, input map[string]any, scope NodeExecutionScope) (NodeExecutionHandle, error)
	FinishNode(ctx context.Context, handle NodeExecutionHandle, output map[string]any, status executionmodels.NodeExecutionStatus, trace map[string]any) error
	StartConversation(ctx context.Context, handle NodeExecutionHandle, conversation *agent.Conversation, input map[string]any) error
	FinishConversation(ctx context.Context, conversation *agent.Conversation, output any) error
}

type Executor

type Executor struct {
	NewHarness func(kind agent.Harness) (harnesses.Harness, error)
	Recorder   ExecutionRecorder
	ProjectDir string
	// SeedOutputs pre-completes top-level nodes with recorded outputs
	// from a previous execution ("re-run from here"). Seeded nodes are
	// never executed or recorded again.
	SeedOutputs map[string]map[string]any
}

func NewExecutor

func NewExecutor(project string) *Executor

func (*Executor) Run

func (e *Executor) Run(ctx context.Context, snapshot *Snapshot, input map[string]any) (map[string]any, error)

func (*Executor) RunWithHandle

func (e *Executor) RunWithHandle(ctx context.Context, snapshot *Snapshot, input map[string]any) (map[string]any, *WorkflowExecutionHandle, error)

func (*Executor) Start

func (e *Executor) Start(ctx context.Context, snapshot *Snapshot, input map[string]any) (*WorkflowExecutionHandle, error)

type FlowSpec

type FlowSpec struct {
	Instances map[string]InstanceSpec `yaml:"instances"`
}

type InferenceNodeConfig

type InferenceNodeConfig struct {
	Harness     map[string]any `yaml:"harness"`
	Instruction string         `yaml:"instruction"`
}

type InstanceSpec

type InstanceSpec struct {
	Node   string            `yaml:"node"`
	Inputs map[string]string `yaml:"inputs"`
}

type NodeConfigUpdate

type NodeConfigUpdate struct {
	Model       *string
	Harness     *string
	Instruction *string
	// "" removes the field — the harness default takes over.
	ReasoningEffort *string
	// "" removes the field — the compiler defaults to read_only.
	Permissions *string
}

NodeConfigUpdate carries the editable fields of a node's config. A nil field is left unchanged.

type NodeExecutionHandle

type NodeExecutionHandle struct {
	ID uuid.UUID
}

type NodeExecutionScope

type NodeExecutionScope struct {
	ParentNodeExecutionID uuid.UUID
	IterationIndex        *int
	BranchName            string
}

type NodeSnapshot

type NodeSnapshot struct {
	InstanceID                string
	NodeName                  string
	Kind                      string
	Operation                 string
	Executes                  string
	Over                      string
	Updates                   string
	BreaksOn                  string
	MaxIterations             int
	RoutesOn                  string
	WhenTrue                  string
	WhenFalse                 string
	Instruction               string
	Harness                   agent.Harness `json:"Harness,omitempty"`
	Model                     string
	ReasoningEffort           runtimetypes.ReasoningEffort `json:"ReasoningEffort,omitempty"`
	HarnessConfig             json.RawMessage
	Inputs                    map[string]Port
	InputOrder                []string
	InputBindings             map[string]Binding
	Outputs                   map[string]Port
	Workflow                  *Snapshot
	OutputName                string
	OutputSchema              map[string]any
	StructuredOutputSchema    map[string]any
	StructuredOutputSchemaRaw json.RawMessage
	WrapStructuredOutput      bool
	LoopTarget                *NodeSnapshot
	WhileTarget               *WhileTargetSnapshot
	TrueTarget                *NodeSnapshot
	FalseTarget               *NodeSnapshot
}

type NodeSpec

type NodeSpec struct {
	Kind          string              `yaml:"kind"`
	WorkflowSlug  string              `yaml:"workflow_slug"`
	Operation     string              `yaml:"operation"`
	Executes      string              `yaml:"executes"`
	Over          string              `yaml:"over"`
	Updates       string              `yaml:"updates"`
	BreaksOn      string              `yaml:"breaks_on"`
	MaxIterations int                 `yaml:"max_iterations"`
	RoutesOn      string              `yaml:"routes_on"`
	WhenTrue      string              `yaml:"when_true"`
	WhenFalse     string              `yaml:"when_false"`
	Inputs        map[string]string   `yaml:"inputs"`
	Outputs       map[string]string   `yaml:"outputs"`
	Config        InferenceNodeConfig `yaml:"config"`
}

type ObservedRecorder

type ObservedRecorder struct {
	Inner    ExecutionRecorder
	Observer RunObserver
	// contains filtered or unexported fields
}

ObservedRecorder decorates an ExecutionRecorder with RunObserver callbacks. Recording stays the inner recorder's job — the observer only watches, so a failing print can never corrupt a run. Child scopes (loop iterations, branches) are deliberately not reported.

func (*ObservedRecorder) FinishConversation

func (r *ObservedRecorder) FinishConversation(ctx context.Context, conversation *agent.Conversation, output any) error

func (*ObservedRecorder) FinishNode

func (r *ObservedRecorder) FinishNode(ctx context.Context, handle NodeExecutionHandle, output map[string]any, status executionmodels.NodeExecutionStatus, trace map[string]any) error

func (*ObservedRecorder) FinishWorkflow

func (*ObservedRecorder) StartConversation

func (r *ObservedRecorder) StartConversation(ctx context.Context, handle NodeExecutionHandle, conversation *agent.Conversation, input map[string]any) error

func (*ObservedRecorder) StartNode

func (*ObservedRecorder) StartWorkflow

func (r *ObservedRecorder) StartWorkflow(ctx context.Context, snapshot *Snapshot, input map[string]any, projectDir string) (WorkflowExecutionHandle, error)

type OutputBinding

type OutputBinding struct {
	Name   string
	Schema map[string]any
	From   Binding
}

type Port

type Port struct {
	Name    string
	TypeRef string
	Schema  map[string]any
}

type Registry

type Registry struct {
	// contains filtered or unexported fields
}

Registry is the workflow library, stored in the application database. Every mutation funnels through saveHead, so the version counter and the history in workflow_versions are complete by construction — which the old directory of world-writable YAML files could never guarantee.

func NewRegistry

func NewRegistry(db bun.IDB) *Registry

func (*Registry) Compile

func (r *Registry) Compile(ctx context.Context, spec *Spec) (*Snapshot, error)

Compile compiles a spec, resolving embedded workflows from the registry first and from files next to the spec's source second.

func (*Registry) CreateDraft

func (r *Registry) CreateDraft(ctx context.Context, name, description, explicitID string) (*CreatedDraft, error)

CreateDraft scaffolds a new named workflow as a draft: just the workflow header, no nodes — the composer and inspector fill in the rest. The id derives from the name unless the caller picks one; collisions with installed workflows or existing drafts are rejected.

func (*Registry) Delete

func (r *Registry) Delete(ctx context.Context, workflowID string) error

Delete removes a workflow from the library: the installed spec and any pending draft. Run history and the version history deliberately stay — deleting a workflow does not rewrite the past.

func (*Registry) DeleteDraft

func (r *Registry) DeleteDraft(ctx context.Context, workflowID string) error

DeleteDraft discards a draft — a missing draft is not an error. A draft-only workflow disappears entirely — there is nothing else to keep.

func (*Registry) ExportToFile

func (r *Registry) ExportToFile(ctx context.Context, workflowID string, targetPath string, overwrite bool) error

ExportToFile writes the installed spec's YAML to a file.

func (*Registry) GetVersionSpec

func (r *Registry) GetVersionSpec(ctx context.Context, workflowID string, version int) (string, error)

GetVersionSpec returns the YAML of one past version.

func (*Registry) ImportFile

func (r *Registry) ImportFile(ctx context.Context, sourcePath string, overwrite bool) (WorkflowSummary, error)

ImportFile installs a spec file. Importing over an installed workflow requires overwrite and continues its version history — the row's identity and counter survive.

func (*Registry) List

func (r *Registry) List(ctx context.Context) ([]WorkflowSummary, error)

List returns the installed workflows.

func (*Registry) ListDraftOnly

func (r *Registry) ListDraftOnly(ctx context.Context) ([]WorkflowSummary, error)

ListDraftOnly returns workflows that exist only as drafts — composed but never saved.

func (*Registry) ListVersions

func (r *Registry) ListVersions(ctx context.Context, workflowID string) ([]WorkflowVersionInfo, error)

ListVersions returns a workflow's history, newest first.

func (*Registry) Load

func (r *Registry) Load(ctx context.Context, workflowID string) (*Spec, error)

Load parses the installed spec for a workflow id.

func (*Registry) ReadDraft

func (r *Registry) ReadDraft(ctx context.Context, workflowID string) (string, error)

ReadDraft returns the draft spec for a workflow, or "" when none exists.

func (*Registry) Rename

func (r *Registry) Rename(ctx context.Context, oldID, newID string) (*RenameResult, error)

Rename moves a workflow (installed, drafted, or both) to a new id and updates every spec that embeds it. Each rewritten workflow gets a new version — a rename is a modification like any other.

func (*Registry) RestoreVersion

func (r *Registry) RestoreVersion(ctx context.Context, workflowID string, version int) (*SavedDraft, error)

RestoreVersion re-installs a past version as a new head. History only moves forward — a restore adds a version, it never rewrites the past.

func (*Registry) SaveDraft

func (r *Registry) SaveDraft(ctx context.Context, workflowID string) (*SavedDraft, error)

SaveDraft promotes a draft: it must compile, the next version is stamped, the head is replaced, the outgoing version stays in the history, and the draft is cleared.

func (*Registry) SetHeader

func (r *Registry) SetHeader(ctx context.Context, workflowID, name string, description *string) error

SetHeader rewrites workflow.name and/or workflow.description wherever the workflow lives — the installed spec (as a new version) and/or the pending draft. An empty name is unchanged. A nil description is unchanged, and "" clears it.

func (*Registry) SpecBytes

func (r *Registry) SpecBytes(ctx context.Context, workflowID string) ([]byte, error)

SpecBytes returns the installed spec's raw YAML.

func (*Registry) UpdateNodeConfig

func (r *Registry) UpdateNodeConfig(ctx context.Context, workflowID, nodeName string, update NodeConfigUpdate) error

UpdateNodeConfig edits one node's config in the workflow's YAML. The edit is surgical (yaml.v3 node tree), so comments and formatting survive, and the result lands as a new compiled version — an edit that breaks the spec never installs.

func (*Registry) VerifyProposedSpec

func (r *Registry) VerifyProposedSpec(ctx context.Context, raw []byte, expectedID string) (string, error)

VerifyProposedSpec compiles a composer proposal and confirms its workflow.slug, before the proposal may become a draft.

func (*Registry) WriteDraft

func (r *Registry) WriteDraft(ctx context.Context, workflowID string, raw []byte) error

WriteDraft stores a proposed spec. The content must already be compile-checked by the caller. A workflow that only exists as a draft gets its registry row (and permanent identity) here.

type RenameResult

type RenameResult struct {
	WorkflowSlug string
	// UpdatedRefs lists workflow ids whose specs embedded the
	// renamed workflow and were rewritten to the new id.
	UpdatedRefs []string
}

type RunObserver

type RunObserver interface {
	NodeStarted(instanceID string)
	NodeFinished(instanceID string, status executionmodels.NodeExecutionStatus)
}

RunObserver receives coarse progress while a workflow runs — one callback pair per top-level node. The CLI prints these; servers have no observer.

type SavedDraft

type SavedDraft struct {
	WorkflowSlug string `json:"workflow_slug"`
	Version      string `json:"version"`
	Spec         string `json:"spec"`
}

type SchemaSpec

type SchemaSpec struct {
	Type        string                `yaml:"type"`
	SchemaRef   string                `yaml:"schema_ref"`
	Properties  map[string]SchemaSpec `yaml:"properties"`
	Items       *SchemaSpec           `yaml:"items"`
	Enum        []any                 `yaml:"enum"`
	Optional    bool                  `yaml:"optional"`
	Description string                `yaml:"description"`
	Nullable    bool                  `yaml:"nullable"`
}

type Snapshot

type Snapshot struct {
	WorkflowSlug    string
	WorkflowID      string
	WorkflowVersion string
	Description     string
	Inputs          map[string]Port
	Outputs         map[string]OutputBinding
	Nodes           map[string]NodeSnapshot
	Order           []string
}

func Compile

func Compile(spec *Spec) (*Snapshot, error)

Compile resolves embedded workflows only from files next to the spec's source. Registry.Compile also resolves installed workflows — use it whenever a database is available.

type Spec

type Spec struct {
	Workflow        WorkflowHeader        `yaml:"workflow"`
	Schemas         map[string]SchemaSpec `yaml:"schemas"`
	Nodes           map[string]NodeSpec   `yaml:"nodes"`
	Flow            FlowSpec              `yaml:"flow"`
	SourcePath      string                `yaml:"-"`
	NodeInputOrder  map[string][]string   `yaml:"-"`
	NodeOutputOrder map[string][]string   `yaml:"-"`
}

func LoadSpecFile

func LoadSpecFile(path string) (*Spec, error)

func ParseSpec

func ParseSpec(raw []byte, sourcePath string) (*Spec, error)

ParseSpec decodes spec YAML. sourcePath is recorded so compilation can resolve embedded workflows from sibling files — pass "" for specs that did not come from disk.

type WhileTargetSnapshot

type WhileTargetSnapshot struct {
	InstanceID                string
	NodeName                  string
	Instruction               string
	Harness                   agent.Harness `json:"Harness,omitempty"`
	Model                     string
	ReasoningEffort           runtimetypes.ReasoningEffort `json:"ReasoningEffort,omitempty"`
	HarnessConfig             json.RawMessage
	Inputs                    map[string]Port
	Workflow                  *Snapshot
	UpdateOutputName          string
	UpdateOutputSchema        map[string]any
	BreakOutputName           string
	StructuredOutputSchema    map[string]any
	StructuredOutputSchemaRaw json.RawMessage
}

type WorkflowExecutionHandle

type WorkflowExecutionHandle struct {
	ID uuid.UUID
}

type WorkflowHeader

type WorkflowHeader struct {
	// Slug is the human-facing handle — renameable, used in the CLI,
	// URLs, and cross-workflow references.
	Slug string `yaml:"slug"`
	// ID is the workflow's permanent identity, a uuid that run history
	// keys on. Stamped automatically, never hand-edited.
	ID          string                        `yaml:"id,omitempty"`
	Name        string                        `yaml:"name"`
	Version     string                        `yaml:"version"`
	Description string                        `yaml:"description"`
	Inputs      map[string]string             `yaml:"inputs"`
	Outputs     map[string]WorkflowOutputSpec `yaml:"outputs"`
}

type WorkflowOutputSpec

type WorkflowOutputSpec struct {
	Schema string `yaml:"schema"`
	From   string `yaml:"from"`
}

type WorkflowSummary

type WorkflowSummary struct {
	Slug        string            `json:"slug"`
	ID          string            `json:"id,omitempty"`
	Name        string            `json:"name"`
	Version     string            `json:"version,omitempty"`
	Description string            `json:"description,omitempty"`
	Inputs      map[string]string `json:"inputs"`
	Outputs     map[string]string `json:"outputs"`
	// HasDraft marks unsaved composer changes; DraftOnly marks a
	// workflow that exists only as a draft (never saved).
	HasDraft  bool `json:"has_draft,omitempty"`
	DraftOnly bool `json:"draft_only,omitempty"`
}

type WorkflowVersionInfo

type WorkflowVersionInfo struct {
	Version   int       `json:"version"`
	Current   bool      `json:"current,omitempty"`
	CreatedAt time.Time `json:"created_at"`
}

WorkflowVersionInfo describes one entry of a workflow's history.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL