cord

package module
v0.1.0 Latest Latest
Warning

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

Go to latest
Published: Aug 14, 2026 License: MIT Imports: 18 Imported by: 0

README

cord

CI codecov Go Reference Version License

Cord is a durable workflow library for Go. Compose ordinary typed functions into linear or branching workflows; Cord persists each run in SQLite and resumes available work after a process restart.

flow := runtime.From("order-total", loadOrder)
items := flow.Then(priceItems)
shipping := flow.Then(priceShipping)
orderTotal := cord.Join(items, shipping).Then(add)

total, err := orderTotal.Run(ctx, orderID)

Install

go get github.com/omarluq/cord
go get modernc.org/sqlite

Cord works with a caller-owned *sql.DB backed by SQLite. The examples use the pure-Go modernc.org/sqlite driver.

Quick start

Workflow steps are package-level functions. Each step receives a context and a typed input, and returns a typed output or an error.

package main

import (
    "context"
    "database/sql"
    "errors"
    "fmt"
    "log"

    "github.com/omarluq/cord"
    _ "modernc.org/sqlite"
)

func double(_ context.Context, value int) (int, error) {
    return value * 2, nil
}

func format(_ context.Context, value int) (string, error) {
    return fmt.Sprintf("result: %d", value), nil
}

func run(ctx context.Context) (_ string, err error) {
    db, err := sql.Open(
        "sqlite",
        "file:workflow.db?_pragma=foreign_keys(1)&_pragma=busy_timeout(5000)",
    )
    if err != nil {
        return "", err
    }

    runtime, err := cord.New(db)
    if err != nil {
        return "", errors.Join(err, db.Close())
    }
    defer func() {
        err = errors.Join(err, runtime.Close(), db.Close())
    }()

    flow := runtime.From("double-and-format", double).Then(format)
    return flow.Run(ctx, 21)
}

func main() {
    result, err := run(context.Background())
    if err != nil {
        log.Fatal(err)
    }
    fmt.Println(result) // result: 42
}

The name passed to From is the durable workflow identity. Keep it stable when renaming or refactoring the root function.

Compose workflows

Linear pipelines

Then connects the current output to the next step. Go checks every connection at compile time.

func validate(ctx context.Context, order Order) (ValidatedOrder, error) { /* ... */ }
func reserve(ctx context.Context, order ValidatedOrder) (Reservation, error) { /* ... */ }
func confirm(ctx context.Context, reservation Reservation) (Receipt, error) { /* ... */ }

flow := runtime.
    From("fulfill-order", validate).
    Then(reserve).
    Then(confirm)

receipt, err := flow.Run(ctx, order)

Workflow handles are immutable, so the same node can be the root of multiple branches.

Branch and join

Build each branch from a shared handle, then combine their outputs with Join. The joined step runs after both branches complete.

func loadOrder(ctx context.Context, id OrderID) (Order, error) { /* ... */ }
func priceItems(ctx context.Context, order Order) (Money, error) { /* ... */ }
func priceShipping(ctx context.Context, order Order) (Money, error) { /* ... */ }
func add(_ context.Context, items, shipping Money) (Money, error) {
    return items.Add(shipping), nil
}

root := runtime.From("order-total", loadOrder)
items := root.Then(priceItems)
shipping := root.Then(priceShipping)
total := cord.Join(items, shipping).Then(add)

amount, err := total.Run(ctx, orderID)

See examples/linear and examples/join for complete programs.

Retries and terminal failures

A step returning an error is retried up to three attempts by default. Mark an error permanent when another attempt cannot succeed:

func charge(ctx context.Context, payment Payment) (Receipt, error) {
    if err := payment.Validate(); err != nil {
        return Receipt{}, cord.Permanent(err)
    }
    return gateway.Charge(ctx, payment)
}

Configure the policy used by future runs before submitting them:

err := runtime.SetRetryPolicy(cord.RetryPolicy{
    MaxAttempts: 5,
    BaseDelay:   250 * time.Millisecond,
    MaxDelay:    20 * time.Second,
})

The policy is stored with the run, so every worker applies the same retry semantics. Existing runs retain the policy with which they were submitted.

Architecture

Cord separates workflow definition, durable state, and execution:

Go functions
     │  From / Then / Join
     ▼
typed workflow graph
     │  Run(ctx, input)
     ▼
SQLite: run + nodes + edges + retry policy
     │
     ├──────────────┬──────────────┐
     ▼              ▼              ▼
 runtime A       runtime B      runtime C
 leased worker   leased worker  leased worker
     │              │              │
     └──────────────┴──────────────┘
                    │
                    ▼
        durable output or failure
  1. Define — From, Then, and Join build an immutable typed graph.
  2. Submit — Run compiles the reachable graph and atomically stores its topology, input, and execution policy.
  3. Execute — schedulers claim ready nodes with leases. A completed node persists its output and makes its dependents eligible to run.
  4. Recover — expired claims become available to another runtime. Recreating the workflow definitions registers the functions needed to continue pending runs.
  5. Complete — the terminal node's persisted output is decoded into the workflow's Go result type.

Multiple Cord runtimes may coordinate through the same SQLite database. Durable state—not the process that submitted or executed a run—is authoritative. Workflow inputs, intermediate values, outputs, and failures cross the storage boundary as JSON.

Runtime lifecycle

New applies pending schema migrations and starts the scheduler. Cord owns its scheduler goroutines; the application continues to own the database.

runtime, err := cord.New(db)
if err != nil {
    return err
}
defer runtime.Close() // close the runtime before db

Canceling the context passed to Run durably cancels that run. Closing a runtime stops its local workers without closing the database or canceling persisted runs.

For scheduler tuning, use NewWithOptions:

runtime, err := cord.NewWithOptions(db, cord.RuntimeOptions{
    Concurrency:       8,
    PollInterval:      100 * time.Millisecond,
    LeaseTTL:          30 * time.Second,
    HeartbeatInterval: 10 * time.Second,
    OnSchedulerError:  func(err error) { logger.Error("cord scheduler", "error", err) },
})

Use NewWithOptionsContext when migration should observe a startup context. See the package reference for the complete API.

Development

This project uses mise and Task:

mise install
mise exec -- task ci

License

MIT

SonarCloud Quality Gate

Documentation

Overview

Package cord composes typed Go functions into persistent workflow graphs.

Index

Examples

Constants

This section is empty.

Variables

View Source
var (
	// ErrSchemaOutdated indicates that the Cord schema is absent or older than required.
	ErrSchemaOutdated = errors.New("cord: schema is absent or outdated")
	// ErrSchemaNewer indicates that the Cord schema is newer than this Cord version.
	ErrSchemaNewer = errors.New("cord: schema is newer than runtime")
	// ErrMigrationFailed indicates that Cord could not inspect or migrate its schema.
	ErrMigrationFailed = errors.New("cord: migration failed")
)

Functions

func Permanent

func Permanent(err error) error

Permanent marks err as terminal so Cord skips remaining retry attempts. A nil error remains nil.

Types

type Cord

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

Cord is a persistent workflow runtime. Its concurrency limit is shared by all workflows.

func New

func New(database *sql.DB) (*Cord, error)

New creates a workflow runtime using a caller-owned SQLite database. It applies pending Cord schema migrations before returning.

func NewWithOptions

func NewWithOptions(database *sql.DB, options RuntimeOptions) (*Cord, error)

NewWithOptions creates a workflow runtime with advanced scheduler settings.

func NewWithOptionsContext

func NewWithOptionsContext(ctx context.Context, database *sql.DB, options RuntimeOptions) (*Cord, error)

NewWithOptionsContext creates a workflow runtime and allows schema migration to be canceled.

func (*Cord) Close

func (c *Cord) Close() error

Close releases resources owned by the runtime. It never closes a caller-owned database.

func (*Cord) From

func (c *Cord) From[I, O any](name string, step func(context.Context, I) (O, error)) Workflow[I, O]

From creates a named workflow whose root node invokes step. Name is the workflow's durable identity and must remain stable across implementations.

Example
package main

import (
	"context"
	"database/sql"
	"errors"
	"fmt"

	"github.com/omarluq/cord"
	"github.com/omarluq/cord/internal/exampledb"
)

func exampleDouble(_ context.Context, value int) (int, error) { return value * 2, nil }
func exampleFormat(_ context.Context, value int) (string, error) {
	return fmt.Sprintf("result: %d", value), nil
}

func closeExample(runtime *cord.Cord, database *sql.DB) error {
	return errors.Join(runtime.Close(), database.Close())
}

func main() {
	database := exampledb.DB()

	runtime, err := cord.New(database)
	if err != nil {
		fmt.Println(err)

		return
	}

	result, runErr := runtime.From("double-and-format", exampleDouble).Then(exampleFormat).Run(context.Background(), 21)
	if err := errors.Join(runErr, closeExample(runtime, database)); err != nil {
		fmt.Println(err)

		return
	}

	fmt.Println(result)
}
Output:
result: 42

func (*Cord) SetRetryPolicy

func (c *Cord) SetRetryPolicy(policy RetryPolicy) error

SetRetryPolicy sets the policy snapshotted into subsequently submitted runs. Existing runs retain the policy persisted when they were submitted.

type JoinResult

type JoinResult[I, A, B any] struct {
	// contains filtered or unexported fields
}

JoinResult is a typed handle to two branches awaiting a joined step.

func Join

func Join[I, A, B any](left Workflow[I, A], right Workflow[I, B]) JoinResult[I, A, B]

Join combines two branches from the same workflow definition.

Example
package main

import (
	"context"
	"database/sql"
	"errors"
	"fmt"

	"github.com/omarluq/cord"
	"github.com/omarluq/cord/internal/exampledb"
)

func exampleOrder(_ context.Context, orderID int) (int, error)         { return orderID, nil }
func exampleItems(_ context.Context, _ int) (int, error)               { return 3, nil }
func exampleShipping(_ context.Context, _ int) (int, error)            { return 5, nil }
func exampleTotal(_ context.Context, items, shipping int) (int, error) { return items + shipping, nil }

func closeExample(runtime *cord.Cord, database *sql.DB) error {
	return errors.Join(runtime.Close(), database.Close())
}

func main() {
	database := exampledb.DB()

	runtime, err := cord.New(database)
	if err != nil {
		fmt.Println(err)

		return
	}

	root := runtime.From("order-total", exampleOrder)
	flow := cord.Join(root.Then(exampleItems), root.Then(exampleShipping)).Then(exampleTotal)

	total, runErr := flow.Run(context.Background(), 1001)
	if err := errors.Join(runErr, closeExample(runtime, database)); err != nil {
		fmt.Println(err)

		return
	}

	fmt.Println(total)
}
Output:
8

func (JoinResult[I, A, B]) Then

func (j JoinResult[I, A, B]) Then[O any](
	step func(context.Context, A, B) (O, error),
) Workflow[I, O]

Then appends fn after both joined branches and returns a workflow handle.

type RetryPolicy

type RetryPolicy struct {
	MaxAttempts int
	BaseDelay   time.Duration
	MaxDelay    time.Duration
}

RetryPolicy controls persistent node retry scheduling.

type RuntimeOptions

type RuntimeOptions struct {
	// OnSchedulerError reports scheduler storage errors. The callback must return promptly.
	OnSchedulerError  func(error)
	Concurrency       int
	PollInterval      time.Duration
	LeaseTTL          time.Duration
	HeartbeatInterval time.Duration
}

RuntimeOptions configures advanced scheduler behavior. Zero-valued fields use Cord's defaults. HeartbeatInterval must be shorter than LeaseTTL.

type Workflow

type Workflow[I, O any] struct {
	// contains filtered or unexported fields
}

Workflow is an immutable typed handle to a terminal node in a workflow graph.

func (Workflow[I, O]) Run

func (w Workflow[I, O]) Run(ctx context.Context, input I) (O, error)

Run submits the workflow and waits for its terminal result. The context controls submission, waiting, and cancellation; it is not persisted as workflow data.

func (Workflow[I, O]) Then

func (w Workflow[I, O]) Then[N any](
	step func(context.Context, O) (N, error),
) Workflow[I, N]

Then appends fn after the workflow's current terminal node and returns a new handle.

Directories

Path Synopsis
examples
join command
Command join demonstrates joining two workflow branches.
Command join demonstrates joining two workflow branches.
linear command
Command linear demonstrates composing workflow steps in a linear chain.
Command linear demonstrates composing workflow steps in a linear chain.
internal
backoff
Package backoff calculates retry delays.
Package backoff calculates retry delays.
exampledb
Package exampledb provides SQLite databases for examples.
Package exampledb provides SQLite databases for examples.
serialization
Package serialization provides persisted payload codecs and compatibility fingerprints.
Package serialization provides persisted payload codecs and compatibility fingerprints.
storage
Package storage implements Cord's private SQL state machine.
Package storage implements Cord's private SQL state machine.

Jump to

Keyboard shortcuts

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