cord

package module
v0.5.0 Latest Latest
Warning

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

Go to latest
Published: Aug 26, 2026 License: MIT Imports: 26 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.

Submit and retrieve runs asynchronously

Submit persists the complete run plan and returns its generated UUIDv7 cord.RunID without waiting for execution. Get blocks until that run reaches a durable terminal state and decodes the typed result:

runID, err := flow.Submit(ctx, input)
if err != nil {
    return err
}

// Persist runID in application state or pass it to another process.
result, err := flow.Get(waitCtx, runID)
if err != nil {
    return err
}
_ = result

RunID has string as its underlying type, so applications can store, transfer, and convert known IDs to it. Such conversion does not create a durable run; only Cord generates IDs for submitted runs. To make submission retries converge on one retained run, pass one caller-chosen idempotency key:

runID, err := flow.Submit(ctx, input, "order:1234")

The variadic argument is optional but accepts at most one non-empty key. Reusing a key with the same workflow definition and encoded input returns the original run ID; reusing it for different work returns an error matching cord.ErrRunConflict. An unkeyed retry creates a new run. The caller must retain the key until submission is unambiguously resolved. Key reservations last only as long as the associated run row—and therefore its idempotency data—is retained.

Get is typed and definition-compatibility checked. Reconstruct the workflow's persisted name, input type, reachable topology, function identities and signatures, and terminal node before retrieving a run, including in a different process. Cord recomputes that definition with the run's persisted retry policy; a mismatch returns an error matching cord.ErrRunIncompatible. Canceling waitCtx stops only that Get; it does not cancel the workflow. Use Cancel for explicit durable cancellation:

if err := flow.Cancel(ctx, runID); err != nil &&
    !errors.Is(err, cord.ErrRunFinished) {
    return err
}

Cancel is keyed only by run ID and does not perform Get's workflow compatibility check. It is idempotent for an already canceled run. Missing IDs return an error matching cord.ErrRunNotFound, and cancellation that loses to successful or failed completion returns one matching cord.ErrRunFinished. If the cancellation response is ambiguous, Cord rereads durable state before returning a definitive result; unresolved errors remain safe to retry with the same run ID. Cancellation durably fences future Cord writes. An active attempt in the runtime that calls Cancel is signaled promptly; attempts in other runtimes observe the lost lease through heartbeat and are then canceled cooperatively. Cancel cannot forcibly stop arbitrary Go code, guarantee when an external attempt observes cancellation, or undo external side effects already in progress.

These sentinel outcomes describe durable state, not authorization: for example, ErrRunNotFound can reveal whether an ID exists. A RunID is not an authorization credential. Applications must authenticate callers and enforce tenant ownership before passing an untrusted run ID to Get or Cancel.

Inspect a run

Use InspectRun to read a run's current durable state without waiting for it to finish. ListRunNodes returns its nodes in stable NodeID order:

report, err := runtime.InspectRun(ctx, runID)
if err != nil {
    return err
}
fmt.Printf("%s: %s (%s)\n", report.ID, report.State, report.Reason)

page, err := runtime.ListRunNodes(ctx, runID, cord.NodeQuery{PageSize: 50})
if err != nil {
    return err
}
for _, node := range page.Nodes {
    fmt.Printf("%s: %s attempt %d/%d\n",
        node.NodeID, node.State, node.Attempt, node.MaxAttempts)
}

Run states are running, canceling, completed, failed, and canceled. Node states are pending, ready, running, retry_wait, completed, failed, and canceled. Terminal reports also include a reason such as succeeded, canceled_by_request, failure_non_retryable, or failure_attempts_exhausted.

Reports include lifecycle timestamps and runner information for diagnostics. runtime.RunnerID() identifies the current runtime instance; it is not a hostname, credential, or lease token.

Use NodeQuery to filter by state or terminal reason and to paginate large runs. The default page size is 50 and the maximum is 200. Pass each page's continuation token to the next query. Pages are read separately, so a running workflow may change between requests.

Snapshots omit inputs, outputs, idempotency keys, and user error messages. Use Workflow.Get for the typed result or terminal error. Run IDs and continuation tokens do not grant access: authenticate callers and check ownership before exposing inspection APIs.

When upgrading Cord, stop old processes before starting processes that may apply a database migration. Running mixed Cord versions against the same database is not supported.

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,
})

Each retry field is optional and defaults independently: MaxAttempts to 3, RetryBaseDelay to 500 milliseconds, and RetryMaxDelay to 30 seconds. For example, setting only MaxAttempts retains both default delays. Values must be positive after defaulting, and the maximum delay must be at least the base delay.

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) or Submit(ctx, input[, key])
     ▼
SQL: run + nodes + edges + retry policy + optional idempotency key
     │
     ├──────────────┬──────────────┐
     ▼              ▼              ▼
 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 — Run or Get decodes the terminal node's persisted output into the workflow's Go result type.

Multiple compatible Cord runtimes may coordinate through the same database. Submit wakes its local scheduler, while other runtimes discover work by normal polling. Any runtime with matching function registrations may execute the run; a different compatible runtime may later call Get or Cancel. Closing the submitter does not cancel durable work. 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 separate from node retry identity and external side-effect idempotency. A caller-provided Submit idempotency key can deduplicate creation of a retained run, but it does not make node execution or external effects exactly once. Cord generates run IDs; callers do not select them.

A database commit can succeed even when the submitting or completing process receives an error. After any ambiguous unkeyed Submit, retrying can create a second run. For retry-safe submission, choose and durably retain an idempotency key before the first attempt, then retry with the same workflow definition, exact encoded input, and key until Cord returns a run ID or a definitive conflict. Do not discard or reuse the key merely because an attempt returned a transport, timeout, context, or other persistence error: the commit may have succeeded. Node functions must still make their own external effects idempotent because Cord executes nodes at least once.

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. A Submit context controls validation and persistence only; after a successful return, its cancellation does not cancel the run. A Get context controls only that caller's wait. Use Cancel for durable cancellation. The durable workflow can otherwise 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; any zero-valued field in a supplied Options also uses its own default. The context passed to New controls schema migration and is bounded by Cord's migration timeout.

Concurrency limits executing node functions, not SQL connections. Size the caller-owned pool for application queries, Cord's concurrent scheduler and worker transitions, migration overhead, and blocking Get polling. Remote Turso migrations require at least two open connections. In every deployment, monitor sql.DB.Stats—especially waits—and increase the pool or reduce runtime concurrency when the database is saturated.

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 (
	// ErrRunCanceled indicates that a durable workflow run was canceled.
	ErrRunCanceled = errors.New("cord: workflow run was canceled")
	// ErrRunNotFound indicates that no durable run exists with the supplied ID.
	ErrRunNotFound = errors.New("cord: workflow run not found")
	// ErrRunFinished indicates that cancellation lost to workflow completion or failure.
	ErrRunFinished = errors.New("cord: workflow run already finished")
	// ErrRunConflict indicates that an idempotency key belongs to a different submission.
	ErrRunConflict = errors.New("cord: workflow submission conflicts with an existing run")
	// ErrRunIncompatible indicates that a workflow handle or durable lifecycle
	// snapshot is incompatible with a retained run.
	ErrRunIncompatible = errors.New("cord: workflow is incompatible with the durable run")
)
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 to return and never closes a caller-owned database. Scheduler error callbacks are not part of this wait and may call Close safely. Steps must observe their context. 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) InspectRun added in v0.5.0

func (c *Cord) InspectRun(ctx context.Context, runID RunID) (RunReport, error)

InspectRun returns one read-only, payload-free snapshot of id without waiting for the run to finish or reconstructing its workflow. It does not promote retries, recover leases, or authorize access: applications must enforce their own tenancy and authorization policy before calling it. Missing and malformed or unsupported runs return errors matching ErrRunNotFound and ErrRunIncompatible, respectively.

func (*Cord) ListRunNodes added in v0.5.0

func (c *Cord) ListRunNodes(ctx context.Context, runID RunID, query NodeQuery) (NodePage, error)

ListRunNodes returns one bounded, payload-free page ordered by stable NodeID. Continuation tokens are portable across Cord instances and replicas, but are not credentials and do not replace application authorization. A token is bound to the run and normalized filters. Pages may observe state changes made between calls and therefore do not provide cross-page snapshot isolation.

func (*Cord) RunnerID added in v0.5.0

func (c *Cord) RunnerID() RunnerID

RunnerID returns the opaque identity of this runtime incarnation. The value is stable for a successfully created Cord instance and differs for a newly created instance. It is diagnostic metadata, not a secret, principal, authorization credential, metric label, or lease fencing token. A nil Cord has an empty RunnerID.

func (*Cord) Shutdown added in v0.2.6

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

Shutdown requests runtime cancellation and waits until its scheduler and executing steps exit or ctx is done. It does not wait for scheduler error callbacks, so callbacks may call Shutdown safely and a blocked callback may outlive the wait. It does not close the caller-owned database.

type CurrentLease added in v0.5.0

type CurrentLease struct {
	// ExpiresAt is the database-time lease deadline, normalized to UTC.
	ExpiresAt time.Time
	// RunnerID identifies the runner holding the current lease.
	RunnerID RunnerID
	// Generation is the current lease fencing generation.
	Generation int64
}

CurrentLease describes the durable fence currently held for a running node. RunnerID alone never authorizes a transition; generation and expiry remain part of the storage fencing contract.

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 NodeID added in v0.5.0

type NodeID string

NodeID identifies a node within a durable workflow run.

type NodePage added in v0.5.0

type NodePage struct {
	// ContinuationToken is nonempty when another page may be requested.
	ContinuationToken string
	// Nodes contains immutable report values ordered by stable NodeID.
	Nodes []NodeReport
}

NodePage is one bounded page of node snapshots. Pages are individually coherent durable observations, but multiple pages are not one historical snapshot and may observe intervening transitions.

type NodeQuery added in v0.5.0

type NodeQuery struct {
	// State optionally filters by an exact known node state.
	State *NodeState
	// Reason optionally filters by an exact known terminal reason.
	Reason *TerminalReason
	// ContinuationToken resumes a previous ListRunNodes call.
	ContinuationToken string
	// PageSize is the maximum number of nodes returned in this page.
	PageSize int
}

NodeQuery selects a bounded page of nodes ordered by stable NodeID. State and Reason are optional exact filters. ContinuationToken is opaque and must be reused only with the same run and filters. PageSize zero uses a conservative default; values above the supported maximum are rejected.

type NodeReport added in v0.5.0

type NodeReport struct {
	// EligibleAt is the durable scheduling time. For ready nodes it is the
	// earliest claim time; for retrying nodes it is the retry deadline.
	EligibleAt time.Time
	// FirstStartedAt is the first successful claim time, when known.
	FirstStartedAt *time.Time
	// LastStartedAt is the latest successful claim time, when known.
	LastStartedAt *time.Time
	// StateChangedAt is when the current node state was entered, when known.
	StateChangedAt *time.Time
	// FinishedAt is when the node entered a terminal state, when applicable.
	FinishedAt *time.Time
	// RunnerID identifies the current or most recent successful claimant, when
	// known. It is diagnostic metadata, not an authorization credential.
	RunnerID *RunnerID
	// CurrentLease is present only while the node is running.
	CurrentLease *CurrentLease
	// RunID identifies the node's durable run.
	RunID RunID
	// NodeID is the stable logical node identifier within the run.
	NodeID NodeID
	// FunctionKey identifies the registered function used by the node.
	FunctionKey string
	// State is the node's current durable state.
	State NodeState
	// Reason is the stable terminal reason, or empty while nonterminal.
	Reason TerminalReason
	// Attempt is the number of successful claims made for this node.
	Attempt int
	// MaxAttempts is the persisted attempt limit for this node.
	MaxAttempts int
}

NodeReport is an authoritative current snapshot of one durable run node. It contains no input, output, or user error-message data and is not attempt history. A returned value may be stale immediately.

type NodeState added in v0.5.0

type NodeState string

NodeState is the durable current state of a workflow node.

const (
	// NodeStatePending indicates that a node has unsatisfied dependencies.
	NodeStatePending NodeState = "pending"
	// NodeStateReady indicates that a node is eligible to be claimed.
	NodeStateReady NodeState = "ready"
	// NodeStateRunning indicates that a node is leased by a runner.
	NodeStateRunning NodeState = "running"
	// NodeStateRetryWait indicates that a node is waiting for its retry deadline.
	NodeStateRetryWait NodeState = "retry_wait"
	// NodeStateCompleted indicates that a node completed successfully.
	NodeStateCompleted NodeState = "completed"
	// NodeStateFailed indicates that a node failed terminally.
	NodeStateFailed NodeState = "failed"
	// NodeStateCanceled indicates that a node was canceled terminally.
	NodeStateCanceled NodeState = "canceled"
)

func (NodeState) AllowsReason added in v0.5.0

func (state NodeState) AllowsReason(reason TerminalReason) bool

AllowsReason reports whether reason is legal for state. It rejects unknown states and reasons, missing terminal reasons, and reasons on active states.

func (NodeState) IsKnown added in v0.5.0

func (state NodeState) IsKnown() bool

IsKnown reports whether state is part of the stable lifecycle vocabulary.

func (NodeState) Terminal added in v0.5.0

func (state NodeState) Terminal() (terminal, known bool)

Terminal reports whether state is terminal and whether state is known. A caller must check known before interpreting the terminal result.

type NodeStateCounts added in v0.5.0

type NodeStateCounts struct {
	// Pending is the number of nodes with unsatisfied dependencies.
	Pending int
	// Ready is the number of nodes eligible to be claimed.
	Ready int
	// Running is the number of currently leased nodes.
	Running int
	// RetryWait is the number of nodes waiting for a retry deadline.
	RetryWait int
	// Completed is the number of successfully completed nodes.
	Completed int
	// Failed is the number of terminally failed nodes.
	Failed int
	// Canceled is the number of terminally canceled nodes.
	Canceled int
}

NodeStateCounts contains an explicit count for every node state.

type Options added in v0.2.0

type Options struct {
	// OnSchedulerError reports scheduler storage errors serially from one runtime-
	// owned goroutine. A callback panic is recovered. The callback may call Close or
	// Shutdown; lifecycle waits exclude the callback because Go cannot cancel user
	// code. A callback that never returns therefore leaks that single reporter goroutine
	// and may outlive shutdown. At most 16 errors are queued while it is busy; later errors are
	// dropped and, if reporting resumes before shutdown, summarized by a later callback.
	// Shutdown abandons queued reports; one delivery already racing shutdown may begin.
	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. Zero uses three.
	MaxAttempts int
	// RetryBaseDelay is the initial delay used for retry backoff. Zero uses
	// 500 milliseconds.
	RetryBaseDelay time.Duration
	// RetryMaxDelay caps retry backoff. Zero uses 30 seconds.
	RetryMaxDelay time.Duration
}

Options configures scheduler behavior. Zero-valued fields use Cord's defaults. Retry fields are defaulted independently, so callers may override any subset of them. LeaseTTL must be greater than two milliseconds, HeartbeatInterval must be at least one millisecond, and HeartbeatInterval must be less than half of 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 RunID added in v0.4.0

type RunID string

RunID identifies a durable workflow run. Cord generates RunIDs as UUIDv7 strings; callers may persist and transfer them but cannot construct runs from caller-selected IDs.

type RunReport added in v0.5.0

type RunReport struct {
	// SubmittedAt is when the run was durably created.
	SubmittedAt time.Time
	// FirstStartedAt is when any node was first claimed, when known.
	FirstStartedAt *time.Time
	// StateChangedAt is when the current run state was durably entered.
	StateChangedAt time.Time
	// FinishedAt is when the run entered a terminal state, when applicable.
	FinishedAt *time.Time
	// TerminalRunnerID identifies the runner that committed a claimed terminal
	// transition, when applicable. It is diagnostic metadata, not authority.
	TerminalRunnerID *RunnerID
	// ID identifies the durable run.
	ID RunID
	// WorkflowName is the durable workflow identity.
	WorkflowName string
	// State is the run's current durable state.
	State RunState
	// Reason is the stable terminal reason, or empty while nonterminal.
	Reason TerminalReason
	// NodeCounts contains counts for every node state in this observation.
	NodeCounts NodeStateCounts
}

RunReport is an authoritative current snapshot of one durable run. It omits payloads and user error messages. The value may be stale immediately after it is returned and does not represent event or attempt history. Timestamps are database observations normalized to UTC; precision can vary by provider.

type RunState added in v0.5.0

type RunState string

RunState is the durable current state of a workflow run.

const (
	// RunStateRunning indicates that a run may still make progress.
	RunStateRunning RunState = "running"
	// RunStateCanceling indicates that cancellation is being durably applied.
	RunStateCanceling RunState = "canceling"
	// RunStateCompleted indicates that a run completed successfully.
	RunStateCompleted RunState = "completed"
	// RunStateFailed indicates that a run failed terminally.
	RunStateFailed RunState = "failed"
	// RunStateCanceled indicates that a run was canceled by request.
	RunStateCanceled RunState = "canceled"
)

func (RunState) AllowsReason added in v0.5.0

func (state RunState) AllowsReason(reason TerminalReason) bool

AllowsReason reports whether reason is legal for state. It rejects unknown states and reasons, missing terminal reasons, and reasons on active states.

func (RunState) IsKnown added in v0.5.0

func (state RunState) IsKnown() bool

IsKnown reports whether state is part of the stable lifecycle vocabulary.

func (RunState) Terminal added in v0.5.0

func (state RunState) Terminal() (terminal, known bool)

Terminal reports whether state is terminal and whether state is known. A caller must check known before interpreting the terminal result.

type RunnerID added in v0.5.0

type RunnerID string

RunnerID identifies one Cord runtime incarnation. It is opaque diagnostic metadata, not an authorization credential or a lease fencing token.

type TerminalReason added in v0.5.0

type TerminalReason string

TerminalReason describes why a terminal lifecycle transition was selected. The empty value means that a nonterminal resource has no terminal reason.

const (
	// ReasonSucceeded indicates successful completion.
	ReasonSucceeded TerminalReason = "succeeded"
	// ReasonCanceledByRequest indicates that explicit run cancellation won.
	ReasonCanceledByRequest TerminalReason = "canceled_by_request"
	// ReasonCanceledByRunFailure indicates that another node failed the run.
	ReasonCanceledByRunFailure TerminalReason = "canceled_by_run_failure"
	// ReasonFailureNonRetryable indicates a permanent or non-retryable failure.
	ReasonFailureNonRetryable TerminalReason = "failure_non_retryable"
	// ReasonFailureAttemptsExhausted indicates that execution attempts were exhausted.
	ReasonFailureAttemptsExhausted TerminalReason = "failure_attempts_exhausted"
	// ReasonFailureLeaseExpired indicates that the final claim's lease expired.
	ReasonFailureLeaseExpired TerminalReason = "failure_lease_expired"
)

func (TerminalReason) IsKnown added in v0.5.0

func (reason TerminalReason) IsKnown() bool

IsKnown reports whether reason is a nonempty member of the stable lifecycle vocabulary.

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]) Cancel added in v0.4.0

func (w Workflow[I, O]) Cancel(ctx context.Context, runID RunID) error

Cancel durably cancels runID without checking this handle's workflow identity. Callers must authorize runID before calling Cancel. Cancellation is idempotent for an already canceled run; missing and finished runs return errors matching ErrRunNotFound and ErrRunFinished. It cannot forcibly stop non-cooperative user code already executing. Active attempts in this runtime are signaled promptly; attempts in other runtimes observe cancellation through lease heartbeat failure. If a storage response is ambiguous, Cancel reconciles the authoritative durable status before returning a definitive result.

func (Workflow[I, O]) Get added in v0.4.0

func (w Workflow[I, O]) Get(ctx context.Context, runID RunID) (O, error)

Get blocks until runID reaches a terminal durable state and returns its typed result. The handle must reconstruct the run's workflow name, input type, reachable topology, function identities and signatures, and terminal node; the run's persisted retry policy is used for this definition check. Canceling ctx stops only this wait and does not cancel the run. Missing, canceled, and incompatible runs return errors matching ErrRunNotFound, ErrRunCanceled, and ErrRunIncompatible.

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]) Submit added in v0.4.0

func (w Workflow[I, O]) Submit(ctx context.Context, input I, idempotencyKey ...string) (RunID, error)

Submit durably submits the workflow and returns its Cord-generated UUIDv7 ID. At most one caller-retained idempotency key may be supplied. Reusing a retained key for the same workflow definition and exact encoded input returns the existing ID; conflicting reuse returns an error matching ErrRunConflict.

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 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(errors.Join(err, database.Close()))

		return
	}
	defer func() {
		if closeErr := closeExample(runtime, database); closeErr != nil {
			fmt.Println(closeErr)
		}
	}()

	flow := runtime.From("async-double", exampleDouble)

	runID, submitErr := flow.Submit(context.Background(), 21, "order-21")
	if submitErr != nil {
		fmt.Println(submitErr)

		return
	}

	// Persist runID in application state; it can retrieve the result later.
	result, getErr := flow.Get(context.Background(), runID)
	if getErr != nil {
		fmt.Println(getErr)

		return
	}

	fmt.Println(result)
}
Output:
42

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