Documentation
¶
Index ¶
- type DatabaseStore
- func (s *DatabaseStore) Claim(ctx context.Context, runID, owner string, ttl time.Duration) (bool, error)
- func (s *DatabaseStore) Delete(ctx context.Context, id string) error
- func (s *DatabaseStore) EnsureTable(ctx context.Context) error
- func (s *DatabaseStore) List(ctx context.Context, workflow string, limit int) ([]*RunInstance, error)
- func (s *DatabaseStore) Load(ctx context.Context, id string) (*RunInstance, error)
- func (s *DatabaseStore) PopSignal(ctx context.Context, runID, event string) (Payload, bool, error)
- func (s *DatabaseStore) PushSignal(ctx context.Context, runID, event string, data Payload) error
- func (s *DatabaseStore) Release(ctx context.Context, runID, owner string) error
- func (s *DatabaseStore) Renew(ctx context.Context, runID, owner string, ttl time.Duration) (bool, error)
- func (s *DatabaseStore) Save(ctx context.Context, run *RunInstance) error
- func (s *DatabaseStore) Unfinished(ctx context.Context, limit int) ([]*RunInstance, error)
- type Definition
- type Engine
- func (e *Engine) Cancel(ctx context.Context, runID string) error
- func (e *Engine) Dispatch(name string, payload Payload) (string, error)
- func (e *Engine) DispatchSync(ctx context.Context, name string, payload Payload) (*RunInstance, error)
- func (e *Engine) List(ctx context.Context, workflow string, limit int) ([]*RunInstance, error)
- func (e *Engine) Register(def *Definition)
- func (e *Engine) Resume(ctx context.Context) (int, error)
- func (e *Engine) SetHooks(h EngineHooks)
- func (e *Engine) SetLeaseTTL(d time.Duration)
- func (e *Engine) SetPollInterval(d time.Duration)
- func (e *Engine) Signal(runID, event string, data Payload) error
- func (e *Engine) StartRecovery(ctx context.Context, interval time.Duration)
- func (e *Engine) Status(ctx context.Context, runID string) (*RunInstance, error)
- func (e *Engine) Workflows() []string
- type EngineHooks
- type Leaser
- type MemoryStore
- func (s *MemoryStore) Claim(_ context.Context, runID, owner string, ttl time.Duration) (bool, error)
- func (s *MemoryStore) Delete(_ context.Context, id string) error
- func (s *MemoryStore) List(_ context.Context, workflow string, limit int) ([]*RunInstance, error)
- func (s *MemoryStore) Load(_ context.Context, id string) (*RunInstance, error)
- func (s *MemoryStore) PopSignal(_ context.Context, runID, event string) (Payload, bool, error)
- func (s *MemoryStore) PushSignal(_ context.Context, runID, event string, data Payload) error
- func (s *MemoryStore) Release(_ context.Context, runID, owner string) error
- func (s *MemoryStore) Renew(_ context.Context, runID, owner string, ttl time.Duration) (bool, error)
- func (s *MemoryStore) Save(_ context.Context, run *RunInstance) error
- func (s *MemoryStore) Unfinished(_ context.Context, limit int) ([]*RunInstance, error)
- type Payload
- type RedisStore
- func (s *RedisStore) Claim(ctx context.Context, runID, owner string, ttl time.Duration) (bool, error)
- func (s *RedisStore) Delete(ctx context.Context, id string) error
- func (s *RedisStore) List(ctx context.Context, workflow string, limit int) ([]*RunInstance, error)
- func (s *RedisStore) Load(ctx context.Context, id string) (*RunInstance, error)
- func (s *RedisStore) PopSignal(ctx context.Context, runID, event string) (Payload, bool, error)
- func (s *RedisStore) PushSignal(ctx context.Context, runID, event string, data Payload) error
- func (s *RedisStore) Release(ctx context.Context, runID, owner string) error
- func (s *RedisStore) Renew(ctx context.Context, runID, owner string, ttl time.Duration) (bool, error)
- func (s *RedisStore) Save(ctx context.Context, run *RunInstance) error
- func (s *RedisStore) Unfinished(ctx context.Context, limit int) ([]*RunInstance, error)
- type Run
- type RunInstance
- type RunStatus
- type SignalStore
- type StepBuilder
- func (b *StepBuilder) After(deps ...string) *StepBuilder
- func (b *StepBuilder) ContinueOnFailure() *StepBuilder
- func (b *StepBuilder) Parallel() *StepBuilder
- func (b *StepBuilder) Retry(max int, delay time.Duration) *StepBuilder
- func (b *StepBuilder) WaitForEvent(event string, timeout time.Duration) *StepBuilder
- func (b *StepBuilder) When(fn func(Payload) bool) *StepBuilder
- func (b *StepBuilder) WithTimeout(d time.Duration) *StepBuilder
- type StepDef
- type StepFunc
- type StepInstance
- type StepStatus
- type Store
- type WorkflowPlugin
- type WorkflowRun
- type WorkflowSignal
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) 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
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 (*DatabaseStore) Release ¶ added in v1.9.1
func (s *DatabaseStore) Release(ctx context.Context, runID, owner string) 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 ¶
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 (*Engine) Cancel ¶
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) DispatchSync ¶
func (e *Engine) DispatchSync(ctx context.Context, name string, payload Payload) (*RunInstance, error)
DispatchSync starts a workflow and blocks until completion.
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
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
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
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 ¶
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
StartRecovery calls Resume every interval until ctx is done, so runs of a crashed instance are continued by a live one.
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) 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) PushSignal ¶ added in v1.9.1
func (*MemoryStore) Release ¶ added in v1.9.1
func (s *MemoryStore) Release(_ context.Context, runID, owner string) 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 RedisStore ¶ added in v1.9.1
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) 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) PushSignal ¶ added in v1.9.1
func (*RedisStore) Release ¶ added in v1.9.1
func (s *RedisStore) Release(ctx context.Context, runID, owner string) 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.
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 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 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) 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