Documentation
¶
Overview ¶
Package cord composes typed Go functions into persistent workflow graphs.
Index ¶
- Variables
- func Permanent(err error) error
- type Cord
- func (c *Cord) Close() error
- func (c *Cord) From[I, O any](name string, step func(context.Context, I) (O, error)) Workflow[I, O]
- func (c *Cord) InspectRun(ctx context.Context, runID RunID) (RunReport, error)
- func (c *Cord) ListRunNodes(ctx context.Context, runID RunID, query NodeQuery) (NodePage, error)
- func (c *Cord) RunnerID() RunnerID
- func (c *Cord) Shutdown(ctx context.Context) error
- type CurrentLease
- type JoinResult
- type NodeID
- type NodePage
- type NodeQuery
- type NodeReport
- type NodeState
- type NodeStateCounts
- type Options
- type RunID
- type RunReport
- type RunState
- type RunnerID
- type TerminalReason
- type Workflow
- func (w Workflow[I, O]) Cancel(ctx context.Context, runID RunID) error
- func (w Workflow[I, O]) Get(ctx context.Context, runID RunID) (O, error)
- func (w Workflow[I, O]) Run(ctx context.Context, input I) (O, error)
- func (w Workflow[I, O]) Submit(ctx context.Context, input I, idempotencyKey ...string) (RunID, error)
- func (w Workflow[I, O]) Then[N any](step func(context.Context, O) (N, error)) Workflow[I, N]
Examples ¶
Constants ¶
This section is empty.
Variables ¶
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") )
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 ¶
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 ¶
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 ¶
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 ¶
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
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
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
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
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
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.
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.
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
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
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 ¶
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
Source Files
¶
- active_attempts.go
- cancellation_workflow.go
- completion.go
- completion_publish.go
- cord.go
- errors.go
- graph.go
- join.go
- lifecycle.go
- plan.go
- plan_hash.go
- plan_identity.go
- result_identity_workflow.go
- result_waiting_workflow.go
- retry.go
- run.go
- runtime.go
- runtime_construction.go
- runtime_settings.go
- scheduler.go
- scheduler_claim.go
- scheduler_errors.go
- scheduler_failure.go
- scheduler_heartbeat.go
- scheduler_heartbeat_call.go
- scheduler_invocation.go
- scheduler_transitions.go
- shutdown.go
- snapshot.go
- snapshot_inspection.go
- snapshot_node.go
- snapshot_query.go
- snapshot_run.go
- snapshot_token.go
- submission_workflow.go
- workflow.go
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. |