workflow

package
v1.9.1 Latest Latest
Warning

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

Go to latest
Published: Oct 2, 2026 License: MIT Imports: 16 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type DatabaseStore added in v1.9.1

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

DatabaseStore keeps workflow runs in a SQL table so every instance sees them and they survive restarts.

func NewDatabaseStore added in v1.9.1

func NewDatabaseStore(db *lucid.DB) *DatabaseStore

NewDatabaseStore returns a database-backed Store. Call EnsureTable once.

func (*DatabaseStore) Claim added in v1.9.1

func (s *DatabaseStore) Claim(ctx context.Context, runID, owner string, ttl time.Duration) (bool, error)

func (*DatabaseStore) Delete added in v1.9.1

func (s *DatabaseStore) Delete(ctx context.Context, id string) error

func (*DatabaseStore) EnsureTable added in v1.9.1

func (s *DatabaseStore) EnsureTable(ctx context.Context) error

EnsureTable creates or migrates the workflow_runs table.

func (*DatabaseStore) List added in v1.9.1

func (s *DatabaseStore) List(ctx context.Context, workflow string, limit int) ([]*RunInstance, error)

List returns runs newest first.

func (*DatabaseStore) Load added in v1.9.1

func (s *DatabaseStore) Load(ctx context.Context, id string) (*RunInstance, error)

func (*DatabaseStore) PopSignal added in v1.9.1

func (s *DatabaseStore) PopSignal(ctx context.Context, runID, event string) (Payload, bool, error)

PopSignal takes the oldest signal. The DELETE decides who gets it when several instances poll at once.

func (*DatabaseStore) PushSignal added in v1.9.1

func (s *DatabaseStore) PushSignal(ctx context.Context, runID, event string, data Payload) error

func (*DatabaseStore) Release added in v1.9.1

func (s *DatabaseStore) Release(ctx context.Context, runID, owner string) error

func (*DatabaseStore) Renew added in v1.9.1

func (s *DatabaseStore) Renew(ctx context.Context, runID, owner string, ttl time.Duration) (bool, error)

func (*DatabaseStore) Save added in v1.9.1

func (s *DatabaseStore) Save(ctx context.Context, run *RunInstance) error

func (*DatabaseStore) Unfinished added in v1.9.1

func (s *DatabaseStore) Unfinished(ctx context.Context, limit int) ([]*RunInstance, error)

type Definition

type Definition struct {
	Name  string
	Steps []*StepDef
}

Definition describes a workflow template.

func Define

func Define(name string, builder func(r *Run)) *Definition

Define creates a workflow definition.

type Engine

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

Engine orchestrates workflow execution.

func NewEngine

func NewEngine(store Store) *Engine

NewEngine creates a new workflow engine.

func (*Engine) Cancel

func (e *Engine) Cancel(ctx context.Context, runID string) error

Cancel cancels a workflow run. A step running on this instance sees its context cancelled at once; an instance running it elsewhere notices within its poll interval (see SetPollInterval). Steps that have not run are marked cancelled.

func (*Engine) Dispatch

func (e *Engine) Dispatch(name string, payload Payload) (string, error)

Dispatch starts a new workflow run asynchronously.

func (*Engine) DispatchSync

func (e *Engine) DispatchSync(ctx context.Context, name string, payload Payload) (*RunInstance, error)

DispatchSync starts a workflow and blocks until completion.

func (*Engine) List

func (e *Engine) List(ctx context.Context, workflow string, limit int) ([]*RunInstance, error)

List returns recent runs for a workflow.

func (*Engine) Register

func (e *Engine) Register(def *Definition)

Register adds a workflow definition to the engine.

func (*Engine) Resume added in v1.9.1

func (e *Engine) Resume(ctx context.Context) (int, error)

Resume picks up unfinished runs whose engine is gone (after a restart or a crashed instance) and continues them from their last completed step. Steps that were mid-flight run again. It returns how many runs it took. Only stores that implement Leaser support it.

func (*Engine) SetHooks

func (e *Engine) SetHooks(h EngineHooks)

SetHooks configures lifecycle hooks.

func (*Engine) SetLeaseTTL added in v1.9.1

func (e *Engine) SetLeaseTTL(d time.Duration)

SetLeaseTTL sets how long a run stays leased to an engine that stopped renewing it (crashed) before Resume on another instance takes it over.

func (*Engine) SetPollInterval added in v1.9.1

func (e *Engine) SetPollInterval(d time.Duration)

SetPollInterval sets how often a running run checks the store for a cancellation made on another instance, and a waiting step checks for stored signals (default 1s).

func (*Engine) Signal

func (e *Engine) Signal(runID, event string, data Payload) error

Signal sends an external event to a workflow step waiting for it. If the run is waiting on this instance, the step wakes at once. Otherwise, with a store that implements SignalStore (all built-in stores do), the signal is kept until the step picks it up: on another instance, after a restart, or if the step has not started waiting yet.

func (*Engine) StartRecovery added in v1.9.1

func (e *Engine) StartRecovery(ctx context.Context, interval time.Duration)

StartRecovery calls Resume every interval until ctx is done, so runs of a crashed instance are continued by a live one.

func (*Engine) Status

func (e *Engine) Status(ctx context.Context, runID string) (*RunInstance, error)

Status returns the current state of a workflow run.

func (*Engine) Workflows

func (e *Engine) Workflows() []string

Workflows returns all registered workflow names.

type EngineHooks

type EngineHooks struct {
	OnStepStart    func(runID, step string)
	OnStepComplete func(runID, step string, output Payload, duration time.Duration)
	OnStepFail     func(runID, step string, err error, attempt int)
	OnRunComplete  func(runID, workflow string, payload Payload)
	OnRunFail      func(runID, workflow string, err error)
}

EngineHooks allows observability into workflow execution.

type Leaser added in v1.9.1

type Leaser interface {
	// Claim leases runID to owner if nobody holds a live lease on it.
	Claim(ctx context.Context, runID, owner string, ttl time.Duration) (bool, error)
	// Renew extends owner's lease. It returns false if the lease was lost.
	Renew(ctx context.Context, runID, owner string, ttl time.Duration) (bool, error)
	// Release drops owner's lease.
	Release(ctx context.Context, runID, owner string) error
	// Unfinished returns runs that are pending, running or paused.
	Unfinished(ctx context.Context, limit int) ([]*RunInstance, error)
}

Leaser is implemented by stores that several engine instances can share. An engine holds a lease on every run it executes and renews it while the run is in progress; when an instance dies, its leases lapse and Resume on another instance picks the runs up.

type MemoryStore

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

MemoryStore provides an in-memory Store implementation.

func NewMemoryStore

func NewMemoryStore() *MemoryStore

func (*MemoryStore) Claim added in v1.9.1

func (s *MemoryStore) Claim(_ context.Context, runID, owner string, ttl time.Duration) (bool, error)

func (*MemoryStore) Delete

func (s *MemoryStore) Delete(_ context.Context, id string) error

func (*MemoryStore) List

func (s *MemoryStore) List(_ context.Context, workflow string, limit int) ([]*RunInstance, error)

func (*MemoryStore) Load

func (s *MemoryStore) Load(_ context.Context, id string) (*RunInstance, error)

func (*MemoryStore) PopSignal added in v1.9.1

func (s *MemoryStore) PopSignal(_ context.Context, runID, event string) (Payload, bool, error)

func (*MemoryStore) PushSignal added in v1.9.1

func (s *MemoryStore) PushSignal(_ context.Context, runID, event string, data Payload) error

func (*MemoryStore) Release added in v1.9.1

func (s *MemoryStore) Release(_ context.Context, runID, owner string) error

func (*MemoryStore) Renew added in v1.9.1

func (s *MemoryStore) Renew(_ context.Context, runID, owner string, ttl time.Duration) (bool, error)

func (*MemoryStore) Save

func (s *MemoryStore) Save(_ context.Context, run *RunInstance) error

func (*MemoryStore) Unfinished added in v1.9.1

func (s *MemoryStore) Unfinished(_ context.Context, limit int) ([]*RunInstance, error)

type Payload

type Payload map[string]any

Payload is the data bag passed through a workflow run.

type RedisStore added in v1.9.1

type RedisStore struct {
	Retention time.Duration
	// contains filtered or unexported fields
}

RedisStore keeps workflow runs in Redis so every instance sees them. Finished runs expire after Retention (default 30 days).

func NewRedisStore added in v1.9.1

func NewRedisStore(client *redis.Client) *RedisStore

NewRedisStore returns a Redis-backed Store.

func (*RedisStore) Claim added in v1.9.1

func (s *RedisStore) Claim(ctx context.Context, runID, owner string, ttl time.Duration) (bool, error)

func (*RedisStore) Delete added in v1.9.1

func (s *RedisStore) Delete(ctx context.Context, id string) error

func (*RedisStore) List added in v1.9.1

func (s *RedisStore) List(ctx context.Context, workflow string, limit int) ([]*RunInstance, error)

List returns runs newest first.

func (*RedisStore) Load added in v1.9.1

func (s *RedisStore) Load(ctx context.Context, id string) (*RunInstance, error)

func (*RedisStore) PopSignal added in v1.9.1

func (s *RedisStore) PopSignal(ctx context.Context, runID, event string) (Payload, bool, error)

func (*RedisStore) PushSignal added in v1.9.1

func (s *RedisStore) PushSignal(ctx context.Context, runID, event string, data Payload) error

func (*RedisStore) Release added in v1.9.1

func (s *RedisStore) Release(ctx context.Context, runID, owner string) error

func (*RedisStore) Renew added in v1.9.1

func (s *RedisStore) Renew(ctx context.Context, runID, owner string, ttl time.Duration) (bool, error)

func (*RedisStore) Save added in v1.9.1

func (s *RedisStore) Save(ctx context.Context, run *RunInstance) error

func (*RedisStore) Unfinished added in v1.9.1

func (s *RedisStore) Unfinished(ctx context.Context, limit int) ([]*RunInstance, error)

type Run

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

Run is the build context passed to the definition function.

func (*Run) Step

func (r *Run) Step(name string, fn StepFunc) *StepBuilder

Step registers a named step.

type RunInstance

type RunInstance struct {
	ID          string                   `json:"id"`
	Workflow    string                   `json:"workflow"`
	Status      RunStatus                `json:"status"`
	Payload     Payload                  `json:"payload"`
	Steps       map[string]*StepInstance `json:"steps"`
	CreatedAt   time.Time                `json:"created_at"`
	UpdatedAt   time.Time                `json:"updated_at"`
	CompletedAt *time.Time               `json:"completed_at,omitempty"`
	Error       string                   `json:"error,omitempty"`
}

RunInstance holds the state of a running workflow.

type RunStatus

type RunStatus string

RunStatus tracks the lifecycle of a workflow run.

const (
	RunPending   RunStatus = "pending"
	RunRunning   RunStatus = "running"
	RunCompleted RunStatus = "completed"
	RunFailed    RunStatus = "failed"
	RunCancelled RunStatus = "cancelled"
	RunPaused    RunStatus = "paused" // waiting for human input
)

type SignalStore added in v1.9.1

type SignalStore interface {
	PushSignal(ctx context.Context, runID, event string, data Payload) error
	// PopSignal removes and returns the oldest signal, if any.
	PopSignal(ctx context.Context, runID, event string) (Payload, bool, error)
}

SignalStore is implemented by stores that keep signals for waiting steps, so Engine.Signal works whichever instance receives it and even before the step starts waiting. Signals for one run and event are delivered in order.

type StepBuilder

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

StepBuilder provides a fluent API for defining steps.

func (*StepBuilder) After

func (b *StepBuilder) After(deps ...string) *StepBuilder

func (*StepBuilder) ContinueOnFailure

func (b *StepBuilder) ContinueOnFailure() *StepBuilder

func (*StepBuilder) Parallel

func (b *StepBuilder) Parallel() *StepBuilder

func (*StepBuilder) Retry

func (b *StepBuilder) Retry(max int, delay time.Duration) *StepBuilder

func (*StepBuilder) WaitForEvent

func (b *StepBuilder) WaitForEvent(event string, timeout time.Duration) *StepBuilder

func (*StepBuilder) When

func (b *StepBuilder) When(fn func(Payload) bool) *StepBuilder

func (*StepBuilder) WithTimeout

func (b *StepBuilder) WithTimeout(d time.Duration) *StepBuilder

type StepDef

type StepDef struct {
	Name        string
	Fn          StepFunc
	DependsOn   []string      // steps that must complete first
	IsParallel  bool          // can run in parallel with siblings
	MaxRetries  int           // 0 = no retries
	RetryDelay  time.Duration // delay between retries
	Timeout     time.Duration // per-execution timeout
	WaitEvent   string        // external event to wait for
	WaitTimeout time.Duration // how long to wait for the event
	Condition   func(Payload) bool
	OnFailure   string // "continue" | "abort" (default: abort)
}

StepDef describes a single step in a workflow.

type StepFunc

type StepFunc func(ctx context.Context, payload Payload) (Payload, error)

StepFunc performs a single unit of work within a workflow.

type StepInstance

type StepInstance struct {
	Name       string        `json:"name"`
	Status     StepStatus    `json:"status"`
	Output     Payload       `json:"output,omitempty"`
	Error      string        `json:"error,omitempty"`
	Attempts   int           `json:"attempts"`
	StartedAt  *time.Time    `json:"started_at,omitempty"`
	FinishedAt *time.Time    `json:"finished_at,omitempty"`
	Duration   time.Duration `json:"duration_ms,omitempty"`
}

StepInstance holds the state of a step within a running workflow.

type StepStatus

type StepStatus string

StepStatus tracks the lifecycle of a step.

const (
	StatusPending   StepStatus = "pending"
	StatusRunning   StepStatus = "running"
	StatusCompleted StepStatus = "completed"
	StatusFailed    StepStatus = "failed"
	StatusSkipped   StepStatus = "skipped"
	StatusWaiting   StepStatus = "waiting" // waiting for external event
	StatusCancelled StepStatus = "cancelled"
)

type Store

type Store interface {
	Save(ctx context.Context, run *RunInstance) error
	Load(ctx context.Context, id string) (*RunInstance, error)
	List(ctx context.Context, workflow string, limit int) ([]*RunInstance, error)
	Delete(ctx context.Context, id string) error
}

Store persists workflow run state.

type WorkflowPlugin

type WorkflowPlugin struct {
	nimbus.BasePlugin
	Engine *Engine
	// RecoveryInterval is how often a durable store is checked for runs
	// left behind by a restart or a crashed instance (default 30s).
	RecoveryInterval time.Duration
	// contains filtered or unexported fields
}

WorkflowPlugin integrates the workflow engine with Nimbus.

func NewPlugin

func NewPlugin(store Store) *WorkflowPlugin

NewPlugin creates a new workflow plugin.

func (*WorkflowPlugin) Boot

func (p *WorkflowPlugin) Boot(app *nimbus.App) error

Boot starts run recovery when the store is shared and durable, so runs interrupted by a restart or a crashed instance are continued.

func (*WorkflowPlugin) DefaultConfig

func (p *WorkflowPlugin) DefaultConfig() map[string]any

func (*WorkflowPlugin) Register

func (p *WorkflowPlugin) Register(app *nimbus.App) error

func (*WorkflowPlugin) RegisterRoutes

func (p *WorkflowPlugin) RegisterRoutes(r *router.Router)

RegisterRoutes mounts workflow API routes.

func (*WorkflowPlugin) Shutdown added in v1.9.1

func (p *WorkflowPlugin) Shutdown() error

Shutdown stops run recovery.

type WorkflowRun added in v1.9.1

type WorkflowRun struct {
	ID         string    `gorm:"primaryKey;size:36"`
	Workflow   string    `gorm:"size:191;not null;index:idx_workflow_runs_list,priority:1"`
	Status     string    `gorm:"size:16;not null;index"`
	Data       string    `gorm:"type:text;not null"` // JSON of RunInstance
	LeaseOwner string    `gorm:"size:64;not null;default:''"`
	LeaseUntil time.Time `gorm:"not null"`
	CreatedAt  time.Time `gorm:"index:idx_workflow_runs_list,priority:2"`
	UpdatedAt  time.Time
}

WorkflowRun is the database model for DatabaseStore.

func (WorkflowRun) TableName added in v1.9.1

func (WorkflowRun) TableName() string

type WorkflowSignal added in v1.9.1

type WorkflowSignal struct {
	ID        uint   `gorm:"primaryKey;autoIncrement"`
	RunID     string `gorm:"size:36;not null;index:idx_workflow_signals_wait,priority:1"`
	Event     string `gorm:"size:191;not null;index:idx_workflow_signals_wait,priority:2"`
	Data      string `gorm:"type:text;not null"`
	CreatedAt time.Time
}

WorkflowSignal is a signal waiting for a step (DatabaseStore).

func (WorkflowSignal) TableName added in v1.9.1

func (WorkflowSignal) TableName() string

Jump to

Keyboard shortcuts

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