cord

package module
v0.3.0 Latest Latest
Warning

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

Go to latest
Published: Aug 23, 2026 License: MIT Imports: 19 Imported by: 0

README

Cord

CI codecov Go Version 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 or PostgreSQL and resumes available work after a process restart.

Cord banner

Install

go get github.com/omarluq/cord
go get modernc.org/sqlite # or github.com/jackc/pgx/v5

Cord works with a caller-owned *sql.DB backed by SQLite (including remote libSQL) or PostgreSQL. The examples use the pure-Go modernc.org/sqlite driver. PostgreSQL applications use pgx's database/sql driver; Cord does not expose a driver, dialect option, or PostgreSQL-specific API.

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.

PostgreSQL with pgx

Register pgx's database/sql driver in the application, configure and health-check the caller-owned pool, then pass it to the same constructor:

import _ "github.com/jackc/pgx/v5/stdlib"

db, err := sql.Open("pgx", os.Getenv("DATABASE_URL"))
if err != nil {
    return err
}
db.SetMaxOpenConns(10)
db.SetMaxIdleConns(10)
if err := db.PingContext(ctx); err != nil {
    return errors.Join(err, db.Close())
}

runtime, err := cord.New(ctx, db)

Cord detects supported SQL capabilities internally, including through wrapped drivers. An unsupported database makes cord.New return an error matching cord.ErrMigrationFailed; there is no public backend selector. Size the pool for application traffic, concurrent Cord workers, and startup migration work, and monitor sql.DB.Stats for waits. Close and Shutdown never close the pool; the application must close it. See the pgx commands for the linear and join examples.

Supported driver families are modernc, mattn, and ncruces SQLite, remote Turso/libSQL, and pgx v5 stdlib for PostgreSQL 14–18.

Turso and remote libSQL

Use Turso's maintained Go libSQL driver when the database is hosted remotely:

go get github.com/tursodatabase/go-libsql
import _ "github.com/tursodatabase/go-libsql"

func openTurso(ctx context.Context) (*sql.DB, error) {
    databaseURL, err := url.Parse(os.Getenv("TURSO_DATABASE_URL"))
    if err != nil {
        return nil, fmt.Errorf("parse Turso database URL: %w", err)
    }

    query := databaseURL.Query()
    query.Set("authToken", os.Getenv("TURSO_AUTH_TOKEN"))
    databaseURL.RawQuery = query.Encode()

    db, err := sql.Open("libsql", databaseURL.String())
    if err != nil {
        return nil, fmt.Errorf("open Turso database: %w", err)
    }
    // Remote migration locking reserves one connection for lease renewal while
    // Goose holds another for migration work.
    db.SetMaxOpenConns(2)

    if _, err := db.ExecContext(ctx, "PRAGMA foreign_keys = ON"); err != nil {
        return nil, errors.Join(err, db.Close())
    }

    return db, nil
}

Pass the returned database to cord.New(ctx, db). The application owns the credentials, pool configuration, and db.Close. Remote migration locking requires at least two open connections so lease renewal cannot self-block; Cord returns an error instead of migrating when MaxOpenConns is one. Do not configure local-file WAL or locking pragmas for a remote URL. Cord uses an internal database lock to serialize migrations rather than SQLite's file locking.

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 the reusable linear and join workflow packages. Each includes executable sqlite and pg commands; the PostgreSQL commands read CORD_POSTGRES_DSN.

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 retry policy when creating the runtime:

runtime, err := cord.New(ctx, database, cord.Options{
    MaxAttempts:    5,
    RetryBaseDelay: 250 * time.Millisecond,
    RetryMaxDelay:  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(ctx, 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 to pin Go, Task, and Bun. A fresh checkout can provision those tools and the lockfile-pinned browser packages with:

mise install
mise exec -- task playground-install

Validation is split by portability and purpose:

Gate Requirements Coverage
mise exec -- task build-library Go only; CGO disabled Portable Cord library build
mise exec -- go test ./... Go, CGO toolchain, Docker Library tests, including mandatory mattn and Turso/libSQL integration
CORD_POSTGRES_DSN=... mise exec -- task test-postgres Go and PostgreSQL Mandatory PostgreSQL integration/conformance tests and executable pgx example
mise exec -- task build Go and CGO toolchain All native workspace packages
mise exec -- task ci Go, Task, Bun, CGO toolchain, Docker Formatting, lint, dead code, race-tested drivers (including Turso), portable and native builds, and playground assets
mise exec -- task playground-test CI prerequisites, Chromium system libraries, and network access for the first browser install Browser smoke test

Docker is a hard prerequisite for the mandatory Turso tests; they do not silently skip. test-postgres likewise fails if CORD_POSTGRES_DSN is absent or PostgreSQL is unreachable. PR CI uses PostgreSQL 18; scheduled and release gates run the PostgreSQL 14–18 matrix. Task installs frozen Bun dependencies, and the browser test installs pinned Playwright Chromium automatically. Linux is the hosted-CI reference platform; the pure-Go library build is the portability gate for other supported Go platforms. Cord never changes caller-owned database pool settings. Remote Turso migrations need at least two available connections, and applications should monitor sql.DB.Stats pool waits under scheduler load.

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.

Example
package main

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

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

func examplePermanent(_ context.Context, value int) (int, error) {
	return value, fmt.Errorf("invalid input: %w", cord.Permanent(errors.New("terminal")))
}

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, cord.Options{
		MaxAttempts:    3,
		RetryBaseDelay: time.Millisecond,
		RetryMaxDelay:  time.Millisecond,
	})
	if err != nil {
		fmt.Println(err)

		return
	}

	_, runErr := runtime.From("permanent-error", examplePermanent).Run(context.Background(), 1)
	if err := closeExample(runtime, database); err != nil {
		fmt.Println(err)

		return
	}

	fmt.Println(runErr)
}
Output:
invalid input: terminal

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 supported SQL 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 waits for executing steps and scheduler error callbacks to return, and never closes a caller-owned database. Steps must observe their context and callbacks must return promptly. Call Shutdown when the caller needs a bounded wait.

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) Shutdown added in v0.2.6

func (c *Cord) Shutdown(ctx context.Context) error

Shutdown requests runtime cancellation and waits until its goroutines exit or ctx is done. It does not close the caller-owned database. A step that ignores cancellation, or a scheduler error callback that blocks, can outlive the wait.

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. Reporting uses a bounded
	// queue: when the callback is busy, later errors may be coalesced. The callback
	// must return promptly so Close can finish; Shutdown can bound the caller's wait.
	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
	// MaxAttempts limits how many times each node may execute.
	MaxAttempts int
	// RetryBaseDelay is the initial delay used for retry backoff.
	RetryBaseDelay time.Duration
	// RetryMaxDelay caps retry backoff.
	RetryMaxDelay time.Duration
}

Options configures scheduler behavior. Zero-valued fields use Cord's defaults. LeaseTTL and HeartbeatInterval must be at least one millisecond, and HeartbeatInterval must be shorter than LeaseTTL.

Example
package main

import (
	"context"
	"fmt"
	"time"

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

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

	runtime, err := cord.New(context.Background(), database, cord.Options{
		Concurrency:       4,
		PollInterval:      100 * time.Millisecond,
		LeaseTTL:          30 * time.Second,
		HeartbeatInterval: 10 * time.Second,
		OnSchedulerError:  func(err error) { fmt.Println("scheduler:", err) },
	})
	if err != nil {
		fmt.Println(err)

		return
	}

	if err := runtime.Close(); err != nil {
		fmt.Println(err)

		return
	}

	// The caller still owns the database after the Cord runtime closes.
	if err := database.PingContext(context.Background()); err != nil {
		fmt.Println(err)

		return
	}

	fmt.Println("database remains open")

	if err := database.Close(); err != nil {
		fmt.Println(err)
	}
}
Output:
database remains open

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
Package join demonstrates joining two workflow branches.
Package join demonstrates joining two workflow branches.
join/pg command
Command pg runs the joined workflow example with pgx's database/sql driver.
Command pg runs the joined workflow example with pgx's database/sql driver.
join/sqlite command
Command sqlite runs the joined workflow example with modernc SQLite.
Command sqlite runs the joined workflow example with modernc SQLite.
linear
Package linear demonstrates composing workflow steps in a linear chain.
Package linear demonstrates composing workflow steps in a linear chain.
linear/pg command
Command pg runs the linear workflow example with pgx's database/sql driver.
Command pg runs the linear workflow example with pgx's database/sql driver.
linear/sqlite command
Command sqlite runs the linear workflow example with modernc SQLite.
Command sqlite runs the linear workflow example with modernc SQLite.
internal
backoff
Package backoff calculates retry delays.
Package backoff calculates retry delays.
examplecmd
Package examplecmd provides shared execution plumbing for executable examples.
Package examplecmd provides shared execution plumbing for executable examples.
exampledb
Package exampledb opens databases used by executable examples and tests.
Package exampledb opens databases used by executable examples and tests.
hashframe
Package hashframe implements Cord's persistence-sensitive hash framing.
Package hashframe implements Cord's persistence-sensitive hash framing.
serialization
Package serialization provides persisted payload codecs and compatibility fingerprints.
Package serialization provides persisted payload codecs and compatibility fingerprints.
storage
Package storage defines Cord's backend-neutral persistence contracts and models.
Package storage defines Cord's backend-neutral persistence contracts and models.
storage/conformance
Package conformance verifies storage backend behavior.
Package conformance verifies storage backend behavior.
storage/postgres
Package postgres implements Cord's PostgreSQL persistence adapter.
Package postgres implements Cord's PostgreSQL persistence adapter.
storage/sqlite
Package sqlite implements Cord's SQLite persistence adapter.
Package sqlite implements Cord's SQLite persistence adapter.
storage/sqlite/remotelock
Package remotelock provides lease-based migration locking for remote SQLite databases.
Package remotelock provides lease-based migration locking for remote SQLite databases.
storage/sqlstore
Package sqlstore selects and bootstraps Cord's SQL storage backend.
Package sqlstore selects and bootstraps Cord's SQL storage backend.

Jump to

Keyboard shortcuts

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