cord

package module
v0.2.0 Latest Latest
Warning

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

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

README

Cord banner

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(ctx, 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.

Register workflows at startup

Define every workflow a process may execute immediately after creating its Cord runtime, before marking the process ready or serving requests. Calling From and Then registers the step implementations; it does not submit a run.

func defineOrderWorkflow(runtime *cord.Cord) cord.Workflow[OrderID, Receipt] {
    return runtime.
        From("fulfill-order", loadOrder).
        Then(validateOrder).
        Then(chargeOrder)
}

func runService(ctx context.Context, db *sql.DB) (err error) {
    runtime, err := cord.New(ctx, db)
    if err != nil {
        return err
    }
    defer func() {
        if closeErr := runtime.Close(); err == nil {
            err = closeErr
        }
    }()

    orders := defineOrderWorkflow(runtime) // register before readiness

    server := newServer(orders)
    return server.Serve(ctx)
}

After a restart, persisted nodes are claimed only when their exact function keys and signatures have been registered in the new process. A workflow defined lazily inside a request handler can therefore leave existing work dormant until that handler runs again. Keep workflow definitions in startup code and pass the resulting immutable handles to handlers and services.

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.

Side effects and idempotency

Cord executes nodes at least once. A node can finish an external side effect and then lose its lease before Cord persists completion, so a later attempt may run the node again. Workflows are therefore not automatically idempotent: side-effecting steps must use business identifiers, destination-supported idempotency keys, unique constraints, or another application-level deduplication mechanism.

Workflow-submission deduplication is a separate concern from retry identity and external side-effect idempotency. Cord does not currently provide either a caller-selected run ID or exactly-once external effects.

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(ctx, db)
if err != nil {
    return err
}
defer runtime.Close() // close the runtime before db

Canceling the context passed to Run stops submission or waiting, but does not cancel a run that was already persisted. The durable workflow continues and can complete in any compatible Cord process. Closing a runtime similarly stops its local workers without closing the database or canceling persisted runs.

For scheduler tuning, pass an Options value to New:

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

Omit Options to use Cord's defaults. The context passed to New controls schema migration and is bounded by Cord's migration timeout. 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(ctx context.Context, database *sql.DB, options ...Options) (*Cord, error)

New creates a workflow runtime using a caller-owned SQLite database. It accepts at most one Options value and applies pending schema migrations.

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(context.Background(), 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(context.Background(), 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 Options added in v0.2.0

type Options struct {
	// OnSchedulerError reports scheduler storage errors. The callback must return promptly.
	OnSchedulerError func(error)
	// Concurrency limits the number of nodes executing across all workflows.
	Concurrency int
	// PollInterval controls how often idle schedulers check for work.
	PollInterval time.Duration
	// LeaseTTL controls how long a worker owns a claimed node without a heartbeat.
	LeaseTTL time.Duration
	// HeartbeatInterval controls how often workers extend active leases.
	HeartbeatInterval time.Duration
}

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

type RetryPolicy

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

RetryPolicy controls persistent node retry scheduling.

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 and waiting; canceling it does not cancel the durable run.

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