Documentation
¶
Overview ¶
Package cord composes typed Go functions into persistent workflow graphs.
Index ¶
Examples ¶
Constants ¶
This section is empty.
Variables ¶
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 ¶
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 SQLite database. It accepts at most one Options value and applies pending schema migrations.
func (*Cord) Close ¶
Close releases resources owned by the runtime. It never closes a caller-owned database.
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) SetRetryPolicy ¶
func (c *Cord) SetRetryPolicy(policy RetryPolicy) error
SetRetryPolicy sets the policy snapshotted into subsequently submitted runs. Existing runs retain the policy persisted when they were submitted.
type JoinResult ¶
type JoinResult[I, A, B any] struct { // contains filtered or unexported fields }
JoinResult is a typed handle to two branches awaiting a joined step.
func Join ¶
func Join[I, A, B any](left Workflow[I, A], right Workflow[I, B]) JoinResult[I, A, B]
Join combines two branches from the same workflow definition.
Example ¶
package main
import (
"context"
"database/sql"
"errors"
"fmt"
"github.com/omarluq/cord"
"github.com/omarluq/cord/internal/exampledb"
)
func exampleOrder(_ context.Context, orderID int) (int, error) { return orderID, nil }
func exampleItems(_ context.Context, _ int) (int, error) { return 3, nil }
func exampleShipping(_ context.Context, _ int) (int, error) { return 5, nil }
func exampleTotal(_ context.Context, items, shipping int) (int, error) { return items + shipping, nil }
func closeExample(runtime *cord.Cord, database *sql.DB) error {
return errors.Join(runtime.Close(), database.Close())
}
func main() {
database := exampledb.DB()
runtime, err := cord.New(context.Background(), database)
if err != nil {
fmt.Println(err)
return
}
root := runtime.From("order-total", exampleOrder)
flow := cord.Join(root.Then(exampleItems), root.Then(exampleShipping)).Then(exampleTotal)
total, runErr := flow.Run(context.Background(), 1001)
if err := errors.Join(runErr, closeExample(runtime, database)); err != nil {
fmt.Println(err)
return
}
fmt.Println(total)
}
Output: 8
type Options ¶ added in v0.2.0
type Options struct {
// OnSchedulerError reports scheduler storage errors. The callback must return promptly.
OnSchedulerError func(error)
// Concurrency limits the number of nodes executing across all workflows.
Concurrency int
// PollInterval controls how often idle schedulers check for work.
PollInterval time.Duration
// LeaseTTL controls how long a worker owns a claimed node without a heartbeat.
LeaseTTL time.Duration
// HeartbeatInterval controls how often workers extend active leases.
HeartbeatInterval time.Duration
}
Options configures scheduler behavior. Zero-valued fields use Cord's defaults. HeartbeatInterval must be shorter than LeaseTTL.
type RetryPolicy ¶
RetryPolicy controls persistent node retry scheduling.
type Workflow ¶
type Workflow[I, O any] struct { // contains filtered or unexported fields }
Workflow is an immutable typed handle to a terminal node in a workflow graph.
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
examples
|
|
|
join
command
Command join demonstrates joining two workflow branches.
|
Command join demonstrates joining two workflow branches. |
|
linear
command
Command linear demonstrates composing workflow steps in a linear chain.
|
Command linear demonstrates composing workflow steps in a linear chain. |
|
internal
|
|
|
backoff
Package backoff calculates retry delays.
|
Package backoff calculates retry delays. |
|
exampledb
Package exampledb provides SQLite databases for examples.
|
Package exampledb provides SQLite databases for examples. |
|
serialization
Package serialization provides persisted payload codecs and compatibility fingerprints.
|
Package serialization provides persisted payload codecs and compatibility fingerprints. |
|
storage
Package storage implements Cord's private SQL state machine.
|
Package storage implements Cord's private SQL state machine. |