workflow

package
v0.0.26 Latest Latest
Warning

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

Go to latest
Published: Sep 24, 2026 License: MIT Imports: 6 Imported by: 0

README

Workflow Middleware

What it does

Composes multiple middleware/handler steps into a single, ordered request workflow: sequential steps, conditional steps, branching, parallel fan-out, retry, per-step timeout, compensation on failure, observability hooks, and handoff to fh's durable job queue for async work.

How to implement

package main

import (
	"time"

	"github.com/oarkflow/fh"
	"github.com/oarkflow/fh/mw/workflow"
)

func main() {
	app := fh.New()

	wf := workflow.New("checkout").
		UseWithOptions("charge-payment", chargePayment,
			workflow.WithRetry(2, 50*time.Millisecond),
			workflow.WithTimeout(3*time.Second)).
		Use("respond", func(c fh.Ctx) error {
			return c.SendString("ok")
		})

	app.Post("/orders", wf.Handler())
	app.Listen(":8080")
}

func chargePayment(c fh.Ctx) error {
	// call payment gateway
	return nil
}
Sequential steps
wf := workflow.New("demo").
	Use("step-a", handlerA).
	Use("step-b", handlerB)

Use accepts an optional condition function as its third argument to skip a step:

wf.Use("admin-only", handler, func(c fh.Ctx) bool {
	return c.Locals("role") == "admin"
})
Timeout and retry

UseWithOptions accepts functional options for per-step behavior:

wf.UseWithOptions("call-upstream", handler,
	workflow.WithTimeout(2*time.Second),         // bounds the step via a deadline context
	workflow.WithRetry(3, 100*time.Millisecond), // up to 3 retries with fixed backoff
	workflow.WithCondition(cond),
)

Timeouts set c.Context() to a deadline context for the duration of the step and restore the parent context afterward — the handler is expected to observe c.Context().Done() for long-running work, the same cooperative model used by mw/timeout. Retries stop early once the context is done and do not retry after a successful attempt.

Branching

Runs the first branch whose Condition passes:

wf.Branch("shipping",
	workflow.New("express").Condition(isExpress).Use("schedule", scheduleExpress),
	workflow.New("standard").Use("schedule", scheduleStandard), // fallback (no condition)
)
Parallel fan-out
wf.Parallel("fan-out", reserveInventory, sendConfirmation)     // fail-fast: first error wins
wf.ParallelJoin("fan-out", reserveInventory, sendConfirmation) // waits for all, joins errors

Each branch runs in its own goroutine against the same fh.Ctx; panics inside a branch are recovered and converted to errors.

Async job handoff
wf.Job("schedule-shipment", "shipment.schedule")

Hands the step off to fh's durable queue via fh.AtomicHandoff instead of running inline. Requires Reliability.QueueEnabled in fh.Config. The assigned job ID is stored in c.Locals("job_id").

Compensation

OnError is invoked whenever a step fails (including recovered panics). Returning nil swallows the error and continues the workflow — useful for compensating transactions. Returning a non-nil error aborts the workflow with that error:

wf.OnError(func(step string, err error) error {
	if step == "reserve-inventory" {
		releaseHold(step)
		return nil // compensate and continue
	}
	return err // abort
})
Observability
wf.OnStepStart(func(step string) { ... }).
	OnStepComplete(func(step string, err error, dur time.Duration) { ... }).
	OnComplete(func(err error, dur time.Duration) { ... })

Impact

Enables orchestration inside request handling. Complexity and latency depend on workflow structure; parallel branches add goroutines per request.

Ordering guidance

Use for request-scoped workflows with a handful of steps, not as a replacement for a durable background workflow engine. Long-running or failure-sensitive work should go through Job and the durable queue rather than running inline.

Production considerations

  • Keep workflows observable with OnStepStart/OnStepComplete/OnComplete and correlate against request_id.
  • Bound every external call with WithTimeout; a slow step without a timeout blocks the whole request.
  • Use OnError for compensation, not to silently mask bugs — log every compensated error.
  • Test each branch and the parallel fan-out under both success and failure independently.
  • See workflow_test.go for executable success, failure, retry, timeout, branch and compensation examples.

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Step

type Step struct {
	Name      string
	Type      StepType
	Handler   fh.HandlerFunc
	JobType   string
	Condition func(fh.Ctx) bool
	Branches  []*Workflow

	// Timeout bounds a StepSync handler's execution. If <= 0, no deadline is applied.
	// The step observes fh.Ctx.Context() the same way mw/timeout does; it is the
	// handler's responsibility to check ctx.Done() for long-running work.
	Timeout time.Duration

	// RetryAttempts is the number of additional attempts after the first failure
	// for a StepSync handler. 0 means no retry.
	RetryAttempts int

	// RetryBackoff is the delay between retry attempts. Ignored if RetryAttempts is 0.
	RetryBackoff time.Duration

	// ParallelFailFast controls StepParallel behavior. When true (default), the
	// first branch error cancels the workflow. When false, all branches run to
	// completion and their errors are joined.
	ParallelFailFast bool
}

Step is one unit of work inside a Workflow.

type StepOption

type StepOption func(*Step)

StepOption configures a Step registered through UseWithOptions.

func WithCondition

func WithCondition(fn func(fh.Ctx) bool) StepOption

WithCondition only runs the step when fn returns true.

func WithRetry

func WithRetry(attempts int, backoff time.Duration) StepOption

WithRetry retries a failed step up to attempts additional times, waiting backoff between attempts. Retries stop early if the context is done.

func WithTimeout

func WithTimeout(d time.Duration) StepOption

WithTimeout bounds the step's execution time via a deadline context.

type StepType

type StepType uint8
const (
	StepSync StepType = iota
	StepAsync
	StepBranch
	StepParallel
)

type Workflow

type Workflow struct {
	Name string

	Steps []Step
	// contains filtered or unexported fields
}

Workflow is an ordered, composable sequence of steps executed as middleware.

func New

func New(name string) *Workflow

New creates a named workflow.

func (*Workflow) Branch

func (w *Workflow) Branch(name string, branches ...*Workflow) *Workflow

Branch runs the first branch whose Condition passes (or the first unconditional branch), then stops evaluating further branches.

func (*Workflow) Condition

func (w *Workflow) Condition(fn func(fh.Ctx) bool) *Workflow

Condition only runs the whole workflow when fn returns true.

func (*Workflow) Handler

func (w *Workflow) Handler() fh.HandlerFunc

Handler returns the fh.HandlerFunc that executes this workflow.

func (*Workflow) Job

func (w *Workflow) Job(name, jobType string, conditions ...func(fh.Ctx) bool) *Workflow

Job appends an async step that hands the request off to the durable queue via fh.AtomicHandoff instead of running inline.

func (*Workflow) OnComplete

func (w *Workflow) OnComplete(fn func(err error, dur time.Duration)) *Workflow

OnComplete registers an observability hook called after the whole workflow finishes.

func (*Workflow) OnError

func (w *Workflow) OnError(fn func(step string, err error) error) *Workflow

OnError registers a compensation hook invoked when a step returns an error. If the hook returns nil, the workflow continues to the next step instead of aborting. This is intended for compensating transactions (e.g. releasing a reservation made by an earlier step). Panics inside steps are converted to errors and also routed through this hook.

func (*Workflow) OnStepComplete

func (w *Workflow) OnStepComplete(fn func(step string, err error, dur time.Duration)) *Workflow

OnStepComplete registers an observability hook called after each step finishes.

func (*Workflow) OnStepStart

func (w *Workflow) OnStepStart(fn func(step string)) *Workflow

OnStepStart registers an observability hook called before each step runs.

func (*Workflow) Parallel

func (w *Workflow) Parallel(name string, branches ...*Workflow) *Workflow

Parallel runs all branches concurrently and fails fast on the first error.

func (*Workflow) ParallelJoin

func (w *Workflow) ParallelJoin(name string, branches ...*Workflow) *Workflow

ParallelJoin runs all branches concurrently, always waits for every branch to finish, and returns a joined error if any branch failed.

func (*Workflow) Use

func (w *Workflow) Use(name string, handler fh.HandlerFunc, conditions ...func(fh.Ctx) bool) *Workflow

Use appends a synchronous step. conditions[0], if given, gates the step.

func (*Workflow) UseWithOptions

func (w *Workflow) UseWithOptions(name string, handler fh.HandlerFunc, opts ...StepOption) *Workflow

UseWithOptions appends a synchronous step configured with timeout/retry/condition options.

Jump to

Keyboard shortcuts

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