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 ¶
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 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 ¶
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
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. 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.
Source Files
¶
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. |