Documentation
¶
Overview ¶
Package saga implements the Saga pattern for distributed operations with crash recovery. Each saga is a sequence of steps where each step has a corresponding undo operation. The framework guarantees that either all steps complete successfully or all completed steps are rolled back.
See RFD-35 for detailed design documentation.
Example (MakeSandwich) ¶
Example_makeSandwich demonstrates a successful sandwich-making saga.
package main
import (
"context"
"encoding/json"
"fmt"
"strings"
"miren.dev/runtime/pkg/saga"
)
// Kitchen tracks what happened during sandwich making.
type Kitchen struct {
log []string
}
func (k *Kitchen) record(action string) {
k.log = append(k.log, action)
}
// Pantry holds bread and dry goods.
type Pantry struct {
stock map[string]int
}
func NewPantry(stock map[string]int) *Pantry {
return &Pantry{stock: stock}
}
func (p *Pantry) Take(item string) error {
if p.stock[item] <= 0 {
return fmt.Errorf("out of %s", item)
}
p.stock[item]--
return nil
}
func (p *Pantry) Return(item string) {
p.stock[item]++
}
// Fridge holds proteins, condiments, and cold items.
type Fridge struct {
stock map[string]int
}
func NewFridge(stock map[string]int) *Fridge {
return &Fridge{stock: stock}
}
func (f *Fridge) Take(item string) error {
if f.stock[item] <= 0 {
return fmt.Errorf("out of %s", item)
}
f.stock[item]--
return nil
}
func (f *Fridge) Return(item string) {
f.stock[item]++
}
type GetBreadIn struct {
BreadType string
}
type GetBreadOut struct {
Bread string
}
func GetBread(ctx context.Context, in GetBreadIn) (GetBreadOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
pantry := saga.Get[*Pantry](ctx)
if err := pantry.Take(in.BreadType); err != nil {
kitchen.record(fmt.Sprintf("Checked pantry - %v", err))
return GetBreadOut{}, err
}
kitchen.record(fmt.Sprintf("Got %s from pantry", in.BreadType))
return GetBreadOut{Bread: in.BreadType + " slice"}, nil
}
func UndoGetBread(ctx context.Context, in GetBreadIn, out GetBreadOut) error {
kitchen := saga.Get[*Kitchen](ctx)
pantry := saga.Get[*Pantry](ctx)
pantry.Return(in.BreadType)
kitchen.record(fmt.Sprintf("Returned %s to pantry", in.BreadType))
return nil
}
type AddCondimentIn struct {
Bread string
Condiment string
}
type AddCondimentOut struct {
PreparedBread string
}
func AddCondiment(ctx context.Context, in AddCondimentIn) (AddCondimentOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
fridge := saga.Get[*Fridge](ctx)
if err := fridge.Take(in.Condiment); err != nil {
kitchen.record(fmt.Sprintf("Checked fridge - %v", err))
return AddCondimentOut{}, err
}
kitchen.record(fmt.Sprintf("Spread %s on %s", in.Condiment, in.Bread))
return AddCondimentOut{
PreparedBread: in.Bread + " with " + in.Condiment,
}, nil
}
func UndoAddCondiment(ctx context.Context, in AddCondimentIn, out AddCondimentOut) error {
kitchen := saga.Get[*Kitchen](ctx)
fridge := saga.Get[*Fridge](ctx)
fridge.Return(in.Condiment)
kitchen.record(fmt.Sprintf("Scraped %s back into jar", in.Condiment))
return nil
}
type AddProteinIn struct {
PreparedBread string
Protein string
}
type AddProteinOut struct {
Stack string
}
func AddProtein(ctx context.Context, in AddProteinIn) (AddProteinOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
fridge := saga.Get[*Fridge](ctx)
if err := fridge.Take(in.Protein); err != nil {
kitchen.record(fmt.Sprintf("Checked fridge - %v", err))
return AddProteinOut{}, err
}
kitchen.record(fmt.Sprintf("Layered %s on %s", in.Protein, in.PreparedBread))
return AddProteinOut{
Stack: in.PreparedBread + " + " + in.Protein,
}, nil
}
func UndoAddProtein(ctx context.Context, in AddProteinIn, out AddProteinOut) error {
kitchen := saga.Get[*Kitchen](ctx)
fridge := saga.Get[*Fridge](ctx)
fridge.Return(in.Protein)
kitchen.record(fmt.Sprintf("Put %s back in fridge", in.Protein))
return nil
}
type AddToppingsIn struct {
Stack string
Toppings []string `saga:",optional"`
}
type AddToppingsOut struct {
OpenSandwich string
}
func AddToppings(ctx context.Context, in AddToppingsIn) (AddToppingsOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
result := in.Stack
if len(in.Toppings) > 0 {
toppingList := strings.Join(in.Toppings, ", ")
kitchen.record(fmt.Sprintf("Added %s", toppingList))
result += " + " + toppingList
} else {
kitchen.record("No toppings requested")
}
return AddToppingsOut{OpenSandwich: result}, nil
}
func UndoAddToppings(ctx context.Context, in AddToppingsIn, out AddToppingsOut) error {
kitchen := saga.Get[*Kitchen](ctx)
if len(in.Toppings) > 0 {
kitchen.record(fmt.Sprintf("Removed %s", strings.Join(in.Toppings, ", ")))
}
return nil
}
type CloseSandwichIn struct {
OpenSandwich string
}
type CloseSandwichOut struct {
Sandwich string
}
func CloseSandwich(ctx context.Context, in CloseSandwichIn) (CloseSandwichOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
kitchen.record("Closed sandwich with top slice")
return CloseSandwichOut{
Sandwich: "[" + in.OpenSandwich + "]",
}, nil
}
func UndoCloseSandwich(ctx context.Context, in CloseSandwichIn, out CloseSandwichOut) error {
kitchen := saga.Get[*Kitchen](ctx)
kitchen.record("Opened sandwich back up")
return nil
}
func main() {
kitchen := &Kitchen{}
pantry := NewPantry(map[string]int{
"sourdough": 2,
"wheat": 1,
"rye": 1,
})
fridge := NewFridge(map[string]int{
"mayo": 3,
"mustard": 2,
"ham": 4,
"turkey": 2,
"pastrami": 1,
})
registry := saga.NewRegistry()
saga.Define("make-sandwich").
Using(kitchen).
Using(pantry).
Using(fridge).
Action(GetBread).Undo(UndoGetBread).
Action(AddCondiment).Undo(UndoAddCondiment).
Action(AddProtein).Undo(UndoAddProtein).
Action(AddToppings).Undo(UndoAddToppings).
Action(CloseSandwich).Undo(UndoCloseSandwich).
RegisterTo(registry)
storage := saga.NewMemoryStorage()
executor := saga.NewExecutor(storage, saga.WithRegistry(registry))
ctx := context.Background()
err := executor.Start("make-sandwich").
Input("breadtype", "sourdough").
Input("condiment", "mayo").
Input("protein", "ham").
Input("toppings", []string{"lettuce", "tomato"}).
WithID("order-1").
Execute(ctx)
if err != nil {
fmt.Printf("Failed: %v\n", err)
return
}
// Get the final result
exec, _ := storage.Get(ctx, "order-1")
var result CloseSandwichOut
json.Unmarshal(exec.ExecutedActions["close-sandwich"].Output, &result)
fmt.Printf("Result: %s\n", result.Sandwich)
fmt.Println("\nKitchen log:")
for _, entry := range kitchen.log {
fmt.Println(" -", entry)
}
}
Output: Result: [sourdough slice with mayo + ham + lettuce, tomato] Kitchen log: - Got sourdough from pantry - Spread mayo on sourdough slice - Layered ham on sourdough slice with mayo - Added lettuce, tomato - Closed sandwich with top slice
Example (OutOfStock) ¶
Example_outOfStock demonstrates saga compensation when the fridge is empty.
package main
import (
"context"
"fmt"
"strings"
"miren.dev/runtime/pkg/saga"
)
// Kitchen tracks what happened during sandwich making.
type Kitchen struct {
log []string
}
func (k *Kitchen) record(action string) {
k.log = append(k.log, action)
}
// Pantry holds bread and dry goods.
type Pantry struct {
stock map[string]int
}
func NewPantry(stock map[string]int) *Pantry {
return &Pantry{stock: stock}
}
func (p *Pantry) Take(item string) error {
if p.stock[item] <= 0 {
return fmt.Errorf("out of %s", item)
}
p.stock[item]--
return nil
}
func (p *Pantry) Return(item string) {
p.stock[item]++
}
// Fridge holds proteins, condiments, and cold items.
type Fridge struct {
stock map[string]int
}
func NewFridge(stock map[string]int) *Fridge {
return &Fridge{stock: stock}
}
func (f *Fridge) Take(item string) error {
if f.stock[item] <= 0 {
return fmt.Errorf("out of %s", item)
}
f.stock[item]--
return nil
}
func (f *Fridge) Return(item string) {
f.stock[item]++
}
type GetBreadIn struct {
BreadType string
}
type GetBreadOut struct {
Bread string
}
func GetBread(ctx context.Context, in GetBreadIn) (GetBreadOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
pantry := saga.Get[*Pantry](ctx)
if err := pantry.Take(in.BreadType); err != nil {
kitchen.record(fmt.Sprintf("Checked pantry - %v", err))
return GetBreadOut{}, err
}
kitchen.record(fmt.Sprintf("Got %s from pantry", in.BreadType))
return GetBreadOut{Bread: in.BreadType + " slice"}, nil
}
func UndoGetBread(ctx context.Context, in GetBreadIn, out GetBreadOut) error {
kitchen := saga.Get[*Kitchen](ctx)
pantry := saga.Get[*Pantry](ctx)
pantry.Return(in.BreadType)
kitchen.record(fmt.Sprintf("Returned %s to pantry", in.BreadType))
return nil
}
type AddCondimentIn struct {
Bread string
Condiment string
}
type AddCondimentOut struct {
PreparedBread string
}
func AddCondiment(ctx context.Context, in AddCondimentIn) (AddCondimentOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
fridge := saga.Get[*Fridge](ctx)
if err := fridge.Take(in.Condiment); err != nil {
kitchen.record(fmt.Sprintf("Checked fridge - %v", err))
return AddCondimentOut{}, err
}
kitchen.record(fmt.Sprintf("Spread %s on %s", in.Condiment, in.Bread))
return AddCondimentOut{
PreparedBread: in.Bread + " with " + in.Condiment,
}, nil
}
func UndoAddCondiment(ctx context.Context, in AddCondimentIn, out AddCondimentOut) error {
kitchen := saga.Get[*Kitchen](ctx)
fridge := saga.Get[*Fridge](ctx)
fridge.Return(in.Condiment)
kitchen.record(fmt.Sprintf("Scraped %s back into jar", in.Condiment))
return nil
}
type AddProteinIn struct {
PreparedBread string
Protein string
}
type AddProteinOut struct {
Stack string
}
func AddProtein(ctx context.Context, in AddProteinIn) (AddProteinOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
fridge := saga.Get[*Fridge](ctx)
if err := fridge.Take(in.Protein); err != nil {
kitchen.record(fmt.Sprintf("Checked fridge - %v", err))
return AddProteinOut{}, err
}
kitchen.record(fmt.Sprintf("Layered %s on %s", in.Protein, in.PreparedBread))
return AddProteinOut{
Stack: in.PreparedBread + " + " + in.Protein,
}, nil
}
func UndoAddProtein(ctx context.Context, in AddProteinIn, out AddProteinOut) error {
kitchen := saga.Get[*Kitchen](ctx)
fridge := saga.Get[*Fridge](ctx)
fridge.Return(in.Protein)
kitchen.record(fmt.Sprintf("Put %s back in fridge", in.Protein))
return nil
}
type AddToppingsIn struct {
Stack string
Toppings []string `saga:",optional"`
}
type AddToppingsOut struct {
OpenSandwich string
}
func AddToppings(ctx context.Context, in AddToppingsIn) (AddToppingsOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
result := in.Stack
if len(in.Toppings) > 0 {
toppingList := strings.Join(in.Toppings, ", ")
kitchen.record(fmt.Sprintf("Added %s", toppingList))
result += " + " + toppingList
} else {
kitchen.record("No toppings requested")
}
return AddToppingsOut{OpenSandwich: result}, nil
}
func UndoAddToppings(ctx context.Context, in AddToppingsIn, out AddToppingsOut) error {
kitchen := saga.Get[*Kitchen](ctx)
if len(in.Toppings) > 0 {
kitchen.record(fmt.Sprintf("Removed %s", strings.Join(in.Toppings, ", ")))
}
return nil
}
type CloseSandwichIn struct {
OpenSandwich string
}
type CloseSandwichOut struct {
Sandwich string
}
func CloseSandwich(ctx context.Context, in CloseSandwichIn) (CloseSandwichOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
kitchen.record("Closed sandwich with top slice")
return CloseSandwichOut{
Sandwich: "[" + in.OpenSandwich + "]",
}, nil
}
func UndoCloseSandwich(ctx context.Context, in CloseSandwichIn, out CloseSandwichOut) error {
kitchen := saga.Get[*Kitchen](ctx)
kitchen.record("Opened sandwich back up")
return nil
}
func main() {
kitchen := &Kitchen{}
pantry := NewPantry(map[string]int{
"wheat": 1,
})
fridge := NewFridge(map[string]int{
"mustard": 1,
"turkey": 0, // Out of turkey!
})
registry := saga.NewRegistry()
saga.Define("make-sandwich-fail").
Using(kitchen).
Using(pantry).
Using(fridge).
Action(GetBread).Undo(UndoGetBread).
Action(AddCondiment).Undo(UndoAddCondiment).
Action(AddProtein).Undo(UndoAddProtein).
Action(AddToppings).Undo(UndoAddToppings).
Action(CloseSandwich).Undo(UndoCloseSandwich).
RegisterTo(registry)
storage := saga.NewMemoryStorage()
executor := saga.NewExecutor(storage, saga.WithRegistry(registry))
ctx := context.Background()
err := executor.Start("make-sandwich-fail").
Input("breadtype", "wheat").
Input("condiment", "mustard").
Input("protein", "turkey").
Input("toppings", []string{"pickles"}).
Execute(ctx)
if err != nil {
fmt.Println("Sandwich failed - Loss Prevented")
}
fmt.Println("\nKitchen log:")
for _, entry := range kitchen.log {
fmt.Println(" -", entry)
}
fmt.Printf("\nInventory restored: wheat=%d, mustard=%d\n",
pantry.stock["wheat"], fridge.stock["mustard"])
}
Output: Sandwich failed - Loss Prevented Kitchen log: - Got wheat from pantry - Spread mustard on wheat slice - Checked fridge - out of turkey - Scraped mustard back into jar - Returned wheat to pantry Inventory restored: wheat=1, mustard=1
Example (SimpleSandwich) ¶
Example_simpleSandwich demonstrates optional toppings.
package main
import (
"context"
"encoding/json"
"fmt"
"strings"
"miren.dev/runtime/pkg/saga"
)
// Kitchen tracks what happened during sandwich making.
type Kitchen struct {
log []string
}
func (k *Kitchen) record(action string) {
k.log = append(k.log, action)
}
// Pantry holds bread and dry goods.
type Pantry struct {
stock map[string]int
}
func NewPantry(stock map[string]int) *Pantry {
return &Pantry{stock: stock}
}
func (p *Pantry) Take(item string) error {
if p.stock[item] <= 0 {
return fmt.Errorf("out of %s", item)
}
p.stock[item]--
return nil
}
func (p *Pantry) Return(item string) {
p.stock[item]++
}
// Fridge holds proteins, condiments, and cold items.
type Fridge struct {
stock map[string]int
}
func NewFridge(stock map[string]int) *Fridge {
return &Fridge{stock: stock}
}
func (f *Fridge) Take(item string) error {
if f.stock[item] <= 0 {
return fmt.Errorf("out of %s", item)
}
f.stock[item]--
return nil
}
func (f *Fridge) Return(item string) {
f.stock[item]++
}
type GetBreadIn struct {
BreadType string
}
type GetBreadOut struct {
Bread string
}
func GetBread(ctx context.Context, in GetBreadIn) (GetBreadOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
pantry := saga.Get[*Pantry](ctx)
if err := pantry.Take(in.BreadType); err != nil {
kitchen.record(fmt.Sprintf("Checked pantry - %v", err))
return GetBreadOut{}, err
}
kitchen.record(fmt.Sprintf("Got %s from pantry", in.BreadType))
return GetBreadOut{Bread: in.BreadType + " slice"}, nil
}
func UndoGetBread(ctx context.Context, in GetBreadIn, out GetBreadOut) error {
kitchen := saga.Get[*Kitchen](ctx)
pantry := saga.Get[*Pantry](ctx)
pantry.Return(in.BreadType)
kitchen.record(fmt.Sprintf("Returned %s to pantry", in.BreadType))
return nil
}
type AddCondimentIn struct {
Bread string
Condiment string
}
type AddCondimentOut struct {
PreparedBread string
}
func AddCondiment(ctx context.Context, in AddCondimentIn) (AddCondimentOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
fridge := saga.Get[*Fridge](ctx)
if err := fridge.Take(in.Condiment); err != nil {
kitchen.record(fmt.Sprintf("Checked fridge - %v", err))
return AddCondimentOut{}, err
}
kitchen.record(fmt.Sprintf("Spread %s on %s", in.Condiment, in.Bread))
return AddCondimentOut{
PreparedBread: in.Bread + " with " + in.Condiment,
}, nil
}
func UndoAddCondiment(ctx context.Context, in AddCondimentIn, out AddCondimentOut) error {
kitchen := saga.Get[*Kitchen](ctx)
fridge := saga.Get[*Fridge](ctx)
fridge.Return(in.Condiment)
kitchen.record(fmt.Sprintf("Scraped %s back into jar", in.Condiment))
return nil
}
type AddProteinIn struct {
PreparedBread string
Protein string
}
type AddProteinOut struct {
Stack string
}
func AddProtein(ctx context.Context, in AddProteinIn) (AddProteinOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
fridge := saga.Get[*Fridge](ctx)
if err := fridge.Take(in.Protein); err != nil {
kitchen.record(fmt.Sprintf("Checked fridge - %v", err))
return AddProteinOut{}, err
}
kitchen.record(fmt.Sprintf("Layered %s on %s", in.Protein, in.PreparedBread))
return AddProteinOut{
Stack: in.PreparedBread + " + " + in.Protein,
}, nil
}
func UndoAddProtein(ctx context.Context, in AddProteinIn, out AddProteinOut) error {
kitchen := saga.Get[*Kitchen](ctx)
fridge := saga.Get[*Fridge](ctx)
fridge.Return(in.Protein)
kitchen.record(fmt.Sprintf("Put %s back in fridge", in.Protein))
return nil
}
type AddToppingsIn struct {
Stack string
Toppings []string `saga:",optional"`
}
type AddToppingsOut struct {
OpenSandwich string
}
func AddToppings(ctx context.Context, in AddToppingsIn) (AddToppingsOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
result := in.Stack
if len(in.Toppings) > 0 {
toppingList := strings.Join(in.Toppings, ", ")
kitchen.record(fmt.Sprintf("Added %s", toppingList))
result += " + " + toppingList
} else {
kitchen.record("No toppings requested")
}
return AddToppingsOut{OpenSandwich: result}, nil
}
func UndoAddToppings(ctx context.Context, in AddToppingsIn, out AddToppingsOut) error {
kitchen := saga.Get[*Kitchen](ctx)
if len(in.Toppings) > 0 {
kitchen.record(fmt.Sprintf("Removed %s", strings.Join(in.Toppings, ", ")))
}
return nil
}
type CloseSandwichIn struct {
OpenSandwich string
}
type CloseSandwichOut struct {
Sandwich string
}
func CloseSandwich(ctx context.Context, in CloseSandwichIn) (CloseSandwichOut, error) {
kitchen := saga.Get[*Kitchen](ctx)
kitchen.record("Closed sandwich with top slice")
return CloseSandwichOut{
Sandwich: "[" + in.OpenSandwich + "]",
}, nil
}
func UndoCloseSandwich(ctx context.Context, in CloseSandwichIn, out CloseSandwichOut) error {
kitchen := saga.Get[*Kitchen](ctx)
kitchen.record("Opened sandwich back up")
return nil
}
func main() {
kitchen := &Kitchen{}
pantry := NewPantry(map[string]int{"rye": 1})
fridge := NewFridge(map[string]int{"butter": 1, "pastrami": 1})
registry := saga.NewRegistry()
saga.Define("simple-sandwich").
Using(kitchen).
Using(pantry).
Using(fridge).
Action(GetBread).Undo(UndoGetBread).
Action(AddCondiment).Undo(UndoAddCondiment).
Action(AddProtein).Undo(UndoAddProtein).
Action(AddToppings).Undo(UndoAddToppings).
Action(CloseSandwich).Undo(UndoCloseSandwich).
RegisterTo(registry)
storage := saga.NewMemoryStorage()
executor := saga.NewExecutor(storage, saga.WithRegistry(registry))
ctx := context.Background()
err := executor.Start("simple-sandwich").
Input("breadtype", "rye").
Input("condiment", "butter").
Input("protein", "pastrami").
// Note: no toppings - that's ok, it's optional!
WithID("order-3").
Execute(ctx)
if err != nil {
fmt.Printf("Failed: %v\n", err)
return
}
// Get the final result
exec, _ := storage.Get(ctx, "order-3")
var result CloseSandwichOut
json.Unmarshal(exec.ExecutedActions["close-sandwich"].Output, &result)
fmt.Printf("Result: %s\n", result.Sandwich)
fmt.Println("\nKitchen log:")
for _, entry := range kitchen.log {
fmt.Println(" -", entry)
}
}
Output: Result: [rye slice with butter + pastrami] Kitchen log: - Got rye from pantry - Spread butter on rye slice - Layered pastrami on rye slice with butter - No toppings requested - Closed sandwich with top slice
Index ¶
- Constants
- Variables
- func DropIfCompleted(ctx context.Context, s Storage, id string) error
- func DropIfFailed(ctx context.Context, s Storage, id string) error
- func Get[T any](ctx context.Context) T
- func StatusIndexAttr(s Status) (entity.Attr, bool)
- func TryGet[T any](ctx context.Context) (T, bool)
- func UndoNested(ctx context.Context, executionID string) error
- type Action
- type ActionBuilder
- type ActionInputs
- type ActionNode
- type ActionResult
- type Builder
- type Definition
- type EACStorage
- func (s *EACStorage) Delete(ctx context.Context, id string) error
- func (s *EACStorage) Get(ctx context.Context, id string) (*Execution, error)
- func (s *EACStorage) ListIncompletePage(ctx context.Context, q IncompleteQuery) (*IncompletePage, error)
- func (s *EACStorage) ListIncompleteSummaryPage(ctx context.Context, q IncompleteSummaryQuery) (*IncompleteSummaryPage, error)
- func (s *EACStorage) ListTerminalPage(ctx context.Context, q TerminalQuery) (*TerminalPage, error)
- func (s *EACStorage) Save(ctx context.Context, exec *Execution) error
- type Edge
- type EntityStorage
- func (s *EntityStorage) Delete(ctx context.Context, id string) error
- func (s *EntityStorage) ForceFailed(ctx context.Context, id string, cutoff time.Time, reason string) (bool, error)
- func (s *EntityStorage) Get(ctx context.Context, id string) (*Execution, error)
- func (s *EntityStorage) ListIncompletePage(ctx context.Context, q IncompleteQuery) (*IncompletePage, error)
- func (s *EntityStorage) ListIncompleteSummaryPage(ctx context.Context, q IncompleteSummaryQuery) (*IncompleteSummaryPage, error)
- func (s *EntityStorage) ListTerminalPage(ctx context.Context, q TerminalQuery) (*TerminalPage, error)
- func (s *EntityStorage) Save(ctx context.Context, exec *Execution) error
- type Execution
- type Executor
- type ExecutorOption
- type IncompletePage
- type IncompleteQuery
- type IncompleteSummary
- type IncompleteSummaryPage
- type IncompleteSummaryQuery
- type MemoryStorage
- func (m *MemoryStorage) Delete(ctx context.Context, id string) error
- func (m *MemoryStorage) ForceFailed(ctx context.Context, id string, cutoff time.Time, reason string) (bool, error)
- func (m *MemoryStorage) Get(ctx context.Context, id string) (*Execution, error)
- func (m *MemoryStorage) ListIncompletePage(ctx context.Context, q IncompleteQuery) (*IncompletePage, error)
- func (m *MemoryStorage) ListIncompleteSummaryPage(ctx context.Context, q IncompleteSummaryQuery) (*IncompleteSummaryPage, error)
- func (m *MemoryStorage) ListTerminalPage(ctx context.Context, q TerminalQuery) (*TerminalPage, error)
- func (m *MemoryStorage) Save(ctx context.Context, exec *Execution) error
- type NestedOption
- type NestedResult
- type Registry
- type RetentionConfig
- type RetentionResult
- type StalledConfig
- type StalledResult
- type StalledStorage
- type StartBuilder
- type Status
- type Storage
- type TerminalExecution
- type TerminalPage
- type TerminalQuery
Examples ¶
Constants ¶
const StalledError = "forced to failed by the stalled-saga sweep: this execution " +
"had not changed state for longer than the stall window, so nothing was " +
"driving it and nothing would have collected it"
StalledError is recorded on a forced execution, so it stays distinguishable from one that genuinely ran and failed.
Variables ¶
var ErrExecutionInProgress = errors.New("execution already in progress")
ErrExecutionInProgress reports that this Executor is already driving the named execution, so this call did nothing. It is not a failure and it is not success: the work is still going, and the caller should come back for the answer rather than read outputs that have not been written yet.
var ErrExecutionNotFound = errors.New("execution not found")
ErrExecutionNotFound is returned by Storage.Get when no execution exists for the given ID.
var ErrRecoveryScopeMismatch = errors.New("execution recovery scope mismatch")
ErrRecoveryScopeMismatch reports that an executor tried to continue a durable execution owned by a different recovery scope. The caller must not run any of that execution's actions against its own local resources.
Functions ¶
func DropIfCompleted ¶ added in v0.14.0
DropIfCompleted removes an execution that succeeded, so a later Execute under the same name builds again instead of reporting the old success.
For a caller that knows the resources the run created are gone. Sandbox creation is the motivating case: containers can vanish under a live runner with nothing written to the entity, and the controller that notices arrives here wanting a build, not a receipt for one.
func DropIfFailed ¶ added in v0.14.0
DropIfFailed removes an execution that failed, so a later Execute under the same name runs again instead of handing back the recorded error.
For a caller that owns the operation and has decided a retry is right. Addon teardown is the motivating case: its undos are all no-ops, so a failed run compensated nothing and left its resources exactly where they were, and a second attempt has strictly more to do rather than something to redo.
func Get ¶
Get retrieves a dependency of type T from the context. Panics if the dependency is not found.
func StatusIndexAttr ¶ added in v0.14.0
StatusIndexAttr returns the entity index attribute that selects executions in the given status, for callers querying the entity store directly rather than through a Storage (the CLI's `debug saga` commands). It reports false for an unrecognized status rather than guessing, since a wrong index silently returns the wrong sagas.
Types ¶
type Action ¶
type Action interface {
// Execute performs the action and returns an output that can be used
// by subsequent actions. The output must be JSON-serializable.
Execute(ctx context.Context, inputs ActionInputs) (output any, err error)
// Undo reverses the action. It receives the same inputs and the output
// that was produced by Execute. Undo should be idempotent.
Undo(ctx context.Context, inputs ActionInputs, output any) error
}
Action represents a single step in a saga. Actions are stateless and created by factories that are registered at application startup. All runtime data flows through ActionInputs.
type ActionBuilder ¶
type ActionBuilder struct {
// contains filtered or unexported fields
}
ActionBuilder provides a fluent API for defining a single action.
func (*ActionBuilder) Undo ¶
func (ab *ActionBuilder) Undo(undo any) *Builder
Undo sets the undo function for the action. The function signature must be: func(ctx context.Context, in InType, out OutType) error
type ActionInputs ¶
type ActionInputs interface {
// Get retrieves an input by key, deserializing it into target.
// Returns an error if the key doesn't exist or deserialization fails.
Get(key string, target any) error
// Has checks if an input exists (for optional inputs).
Has(key string) bool
// Keys returns all available input keys.
Keys() []string
}
ActionInputs provides access to initial saga inputs and outputs from prior actions. All outputs live in a flat namespace.
type ActionNode ¶
type ActionNode struct {
// Name is the unique name of this action within the saga.
Name string
// Action is the stateless action implementation.
Action Action
// InputKeys are the saga keys this action reads from.
InputKeys []string
// OutputKeys are the saga keys this action writes to.
OutputKeys []string
// Dependencies are action names that must complete before this action.
// Computed from InputKeys and other actions' OutputKeys.
Dependencies []string
}
ActionNode describes a single action within a saga definition.
type ActionResult ¶
type ActionResult struct {
// Output is the JSON-serialized output from the action.
Output []byte `json:"output,omitempty"`
// ExecutedAt is when the action was executed.
ExecutedAt time.Time `json:"executed_at"`
// UndoneAt is when the action was undone (nil if not undone).
UndoneAt *time.Time `json:"undone_at,omitempty"`
// Error is set if the action failed during execution.
Error string `json:"error,omitempty"`
}
ActionResult stores the outcome of a single action execution.
type Builder ¶
type Builder struct {
// contains filtered or unexported fields
}
Builder provides a fluent API for defining sagas.
func UsingAs ¶
UsingAs adds a dependency keyed by type T, allowing retrieval via Get[T](ctx). This is useful for injecting implementations that should be retrieved by interface type. For example: UsingAs[MyInterface](b, impl) allows Get[MyInterface](ctx).
func (*Builder) Action ¶
func (b *Builder) Action(args ...any) *ActionBuilder
Action adds an action to the saga using a typed execute function. The function signature must be: func(ctx context.Context, in InType) (OutType, error)
Can be called two ways:
- Action(GetBread) - name derived from function name ("getbread")
- Action("custom-name", GetBread) - explicit name
func (*Builder) Build ¶
func (b *Builder) Build() (*Definition, error)
Build constructs and validates the Definition without registering it. Useful for testing.
func (*Builder) Register ¶
Register validates and registers the saga definition with the global registry. Returns an error if validation fails (cycles, duplicate outputs, type mismatches).
func (*Builder) RegisterTo ¶
RegisterTo validates and registers the saga definition with the given registry. Useful for testing to avoid global state.
type Definition ¶
type Definition struct {
// Name uniquely identifies this saga definition.
Name string
// Version is incremented for breaking changes. Defaults to 1.
Version int
// Actions in this saga, keyed by action name.
Actions map[string]*ActionNode
// contains filtered or unexported fields
}
Definition describes a saga's structure - its actions and their dependencies. Definitions are stateless and registered at application startup.
func GetDefinition ¶
func GetDefinition(name string) (*Definition, bool)
GetDefinition retrieves a saga definition from the global registry.
func (*Definition) ExecutionOrder ¶ added in v0.4.0
func (d *Definition) ExecutionOrder() []string
ExecutionOrder returns the computed execution order for the saga's actions.
type EACStorage ¶ added in v0.6.0
type EACStorage struct {
// contains filtered or unexported fields
}
EACStorage implements Storage using an EntityAccessClient RPC connection. This is used by runners which don't have direct entity.Store access.
func NewEACStorage ¶ added in v0.6.0
func NewEACStorage(eac *es.EntityAccessClient, log *slog.Logger) *EACStorage
NewEACStorage creates a storage backed by an EntityAccessClient.
func (*EACStorage) Delete ¶ added in v0.14.0
func (s *EACStorage) Delete(ctx context.Context, id string) error
Delete removes a saga execution entity via EAC.
func (*EACStorage) ListIncompletePage ¶ added in v0.15.0
func (s *EACStorage) ListIncompletePage(ctx context.Context, q IncompleteQuery) (*IncompletePage, error)
ListIncompletePage returns one bounded page of executions needing recovery via EAC.
The RPC pins the ids and the entities of a page to one store revision, so the tear between listing an index and fetching what it named is closed on the server rather than guessed at here.
func (*EACStorage) ListIncompleteSummaryPage ¶ added in v0.15.0
func (s *EACStorage) ListIncompleteSummaryPage(ctx context.Context, q IncompleteSummaryQuery) (*IncompleteSummaryPage, error)
ListIncompleteSummaryPage summarizes one bounded page of in-flight executions via EAC.
func (*EACStorage) ListTerminalPage ¶ added in v0.15.0
func (s *EACStorage) ListTerminalPage(ctx context.Context, q TerminalQuery) (*TerminalPage, error)
ListTerminalPage summarizes one bounded page of finished executions via EAC.
type Edge ¶ added in v0.6.0
type Edge struct{}
Edge is a zero-size type used to declare ordering dependencies between saga actions without carrying data. Edge fields participate in the dependency graph (via saga struct tags) but are skipped during serialization and deserialization at runtime.
type EntityStorage ¶
type EntityStorage struct {
// contains filtered or unexported fields
}
EntityStorage implements Storage using the entity store.
func NewEntityStorage ¶
func NewEntityStorage(store entity.Store, log *slog.Logger) *EntityStorage
NewEntityStorage creates a storage backed by an entity store.
func (*EntityStorage) Delete ¶ added in v0.14.0
func (s *EntityStorage) Delete(ctx context.Context, id string) error
Delete removes a saga execution entity.
func (*EntityStorage) ForceFailed ¶ added in v0.15.0
func (s *EntityStorage) ForceFailed(ctx context.Context, id string, cutoff time.Time, reason string) (bool, error)
ForceFailed transitions an in-flight execution to failed, and reports whether it did.
The decision lives here rather than in the caller because the write has to be conditional on the read it was made against, and only this layer sees the revision. A caller that read, decided, then asked us to save could overwrite a runner that resumed the saga in the gap, replacing its status and action outputs with a stale copy marked failed.
So the transition is refused, without error, whenever the record moved: gone, already terminal, touched since cutoff, or a changed revision. Each means something else is dealing with it, and the next pass looks again.
The age test goes through lastChanged rather than the execution's updated_at, because a record written before v0.14.0 carries its age only on the entity.
func (*EntityStorage) ListIncompletePage ¶ added in v0.15.0
func (s *EntityStorage) ListIncompletePage(ctx context.Context, q IncompleteQuery) (*IncompletePage, error)
ListIncompletePage returns one bounded page of executions needing recovery.
The three incomplete status indexes are walked one after another rather than merged, because merging would mean holding all three id sets to deduplicate across them, which is the allocation this is here to avoid. A cursor names which index it is in, so a page never straddles two.
Nothing deduplicates across indexes as a result. An execution listed under two statuses at once is a stale index entry, and the caller is the right place to notice: recovery refuses to drive an execution twice through its own claim, and retention's deletes are idempotent by contract.
func (*EntityStorage) ListIncompleteSummaryPage ¶ added in v0.15.0
func (s *EntityStorage) ListIncompleteSummaryPage(ctx context.Context, q IncompleteSummaryQuery) (*IncompleteSummaryPage, error)
ListIncompleteSummaryPage summarizes one bounded page of in-flight executions.
Same indexes and cursor encoding as ListIncompletePage, keeping a summary rather than a whole execution. Its caller walks every in-flight execution in the cluster, and doing that with full payloads would cost what recovery costs without recovering anything.
func (*EntityStorage) ListTerminalPage ¶ added in v0.15.0
func (s *EntityStorage) ListTerminalPage(ctx context.Context, q TerminalQuery) (*TerminalPage, error)
ListTerminalPage summarizes one bounded page of finished executions.
The store's index lookups are equality-only, so there is no range query over a timestamp that would let us ask for expired executions directly. We walk the terminal indexes and read each execution's finish time, keeping only the summary.
type Execution ¶
type Execution struct {
// ID is the unique identifier for this execution.
ID string `json:"id"`
// DefinitionName references the registered saga definition.
DefinitionName string `json:"definition_name"`
// DefinitionVersion is the version of the definition when started.
DefinitionVersion int `json:"definition_version"`
// InitialInputs contains the bootstrap data for the saga.
// All values must be JSON-serializable.
InitialInputs map[string]any `json:"initial_inputs"`
// Status is the current state of the execution.
Status Status `json:"status"`
// ExecutedActions maps action names to their results.
ExecutedActions map[string]*ActionResult `json:"executed_actions"`
// ExecutionOrder records the order actions were executed for reverse undo.
ExecutionOrder []string `json:"execution_order"`
// ParentExecutionID links this execution to a parent saga when run as a nested child.
ParentExecutionID string `json:"parent_execution_id,omitempty"`
// RecoveryScope is the stable identity of the executor allowed to recover
// this execution. An empty scope preserves the unscoped behavior used by
// executors that have no distributed ownership boundary.
RecoveryScope string `json:"recovery_scope,omitempty"`
// Error is set if the saga failed.
Error string `json:"error,omitempty"`
// CreatedAt is when the execution was created.
CreatedAt time.Time `json:"created_at"`
// UpdatedAt is when the execution was last updated.
UpdatedAt time.Time `json:"updated_at"`
}
Execution tracks the runtime state of a saga, persisted after each step.
func ExecutionFromEntity ¶ added in v0.14.0
func ExecutionFromEntity(sagaEntity *saga_v1alpha.Saga) (*Execution, error)
ExecutionFromEntity converts a decoded saga entity into an Execution, deserializing the JSON-encoded inputs, action results, and execution order. Exposed so tools that read saga entities directly (the CLI's `debug saga` commands) decode them the same way the executor does.
Note that CreatedAt and UpdatedAt are not persisted on the entity and so are left zero here; callers that need them should read the entity store's own creation and update metadata.
type Executor ¶
type Executor struct {
// contains filtered or unexported fields
}
Executor orchestrates saga execution with durable logging.
func NewExecutor ¶
func NewExecutor(storage Storage, opts ...ExecutorOption) *Executor
NewExecutor creates an executor with the given storage and options.
func (*Executor) ExecutionOutputs ¶ added in v0.4.0
ExecutionOutputs loads a completed execution from storage and collects its outputs into a NestedResult. Useful for reading saga results without a capture struct.
func (*Executor) Recover ¶
Recover finds and resumes incomplete sagas after a restart.
It walks the incomplete set a page at a time and never holds more than one page, because the set is shared across every executor in the cluster and its size is nobody's decision. A runner restarting into a six-figure backlog used to load all of it, with every action-output blob, before getting as far as deciding which handful it owned.
Recovery is sequential on purpose. Resuming a page's worth of sagas at once would trade the memory this bounds for the same amount of it plus concurrent action side effects, and nothing here needs the throughput.
func (*Executor) Start ¶
func (e *Executor) Start(definitionName string) *StartBuilder
Start begins building a saga execution.
type ExecutorOption ¶
type ExecutorOption func(*Executor)
ExecutorOption configures an Executor.
func WithLogger ¶
func WithLogger(log *slog.Logger) ExecutorOption
WithLogger sets a custom logger for the executor.
func WithRecoveryScope ¶ added in v0.15.0
func WithRecoveryScope(scope string) ExecutorOption
WithRecoveryScope gives this executor a stable recovery identity. New executions persist the scope, startup recovery only considers exact scope matches, and named re-entry refuses to resume another scope's execution. The zero value leaves the executor unscoped and preserves existing behavior.
func WithRegistry ¶
func WithRegistry(r *Registry) ExecutorOption
WithRegistry sets a custom registry for the executor. Useful for testing to avoid global state.
type IncompletePage ¶ added in v0.15.0
type IncompletePage struct {
// Executions are the incomplete executions in this page.
Executions []*Execution
// Cursor resumes the walk after this page, and is empty once there is
// nothing left.
//
// A short page does not mean the end. Backends that walk several status
// indexes finish one before starting the next, and never straddle two in
// one page, so a page can come back well under the limit with plenty still
// to come. Only an empty cursor ends the walk.
Cursor string
}
IncompletePage is one bounded page of executions needing recovery.
type IncompleteQuery ¶ added in v0.15.0
type IncompleteQuery struct {
// Cursor resumes after the last execution of a previous page. Empty starts
// at the beginning.
Cursor string
// Limit caps how many executions the page materializes. Zero or less means
// the backend's own cap, which every backend has: a page with no ceiling is
// the unbounded read this whole design exists to remove.
Limit int
}
IncompleteQuery selects one page of incomplete executions.
It is a struct rather than positional arguments because what recovery needs to say about a page grows: the filtering that keeps an executor from loading another executor's payloads is expressed here too.
type IncompleteSummary ¶ added in v0.15.0
type IncompleteSummary struct {
// ID identifies the execution.
ID string
// Status is the decoded status, not the index the entry came from.
Status Status
// LastChanged is when the execution last changed state, resolved by
// lastChanged. The fallback matters more here than it does for retention:
// v0.11.1's saga schema had no updated_at at all, so the records this sweep
// exists to drain would otherwise all read as infinitely old.
LastChanged time.Time
// ParentID is set when this execution ran as a nested child, and matters
// for the same reason it does to retention: a live parent re-finds its
// child rather than re-running it.
ParentID string
}
IncompleteSummary summarizes an execution that is still in flight: which one, what it is doing, when it last changed, and whose child it is.
Separate from Execution because the stalled sweep walks the whole in-flight set to ask one question about each. Materializing every action-output blob in a six-figure backlog to read a timestamp is the unbounded read MIR-1785 removed, reintroduced for a worse reason.
type IncompleteSummaryPage ¶ added in v0.15.0
type IncompleteSummaryPage struct {
// Executions summarizes the in-flight executions in this page.
Executions []IncompleteSummary
// Cursor resumes the walk after this page, empty once the walk is done.
// The same caveat as IncompletePage.Cursor applies: a short page is not an
// ending, only an empty cursor is.
Cursor string
}
IncompleteSummaryPage is one bounded page of in-flight execution summaries.
type IncompleteSummaryQuery ¶ added in v0.15.0
type IncompleteSummaryQuery struct {
// Cursor resumes after the last execution of a previous page. Empty starts
// at the beginning.
Cursor string
// Limit caps how many executions the page summarizes. Zero or less means
// the backend's own cap.
Limit int
}
IncompleteSummaryQuery selects one page of in-flight execution summaries.
type MemoryStorage ¶
type MemoryStorage struct {
// contains filtered or unexported fields
}
MemoryStorage is a simple in-memory storage implementation for testing and examples.
func NewMemoryStorage ¶
func NewMemoryStorage() *MemoryStorage
NewMemoryStorage creates a new in-memory storage.
func (*MemoryStorage) Delete ¶ added in v0.14.0
func (m *MemoryStorage) Delete(ctx context.Context, id string) error
Delete removes an execution. Deleting a missing execution is a no-op.
func (*MemoryStorage) ForceFailed ¶ added in v0.15.0
func (m *MemoryStorage) ForceFailed(ctx context.Context, id string, cutoff time.Time, reason string) (bool, error)
ForceFailed transitions an in-flight execution to failed, and reports whether it did.
The mutex stands in for the durable backends' revision check: held across the read and the write, so no Save can land between them.
func (*MemoryStorage) ListIncompletePage ¶ added in v0.15.0
func (m *MemoryStorage) ListIncompletePage(ctx context.Context, q IncompleteQuery) (*IncompletePage, error)
ListIncompletePage returns one bounded page of executions needing recovery: pending (crashed before starting), running, and undoing.
The map has no order of its own, so the ids are sorted to give paging a stable sequence to resume in. That costs O(n) per page, which is fine for a backend that holds everything in memory anyway and is the price of behaving like the durable backends for the tests that exercise all three.
func (*MemoryStorage) ListIncompleteSummaryPage ¶ added in v0.15.0
func (m *MemoryStorage) ListIncompleteSummaryPage(ctx context.Context, q IncompleteSummaryQuery) (*IncompleteSummaryPage, error)
ListIncompleteSummaryPage summarizes one bounded page of in-flight executions.
A map has no system timestamp to fall back on, so an execution saved without one really is undatable here and gets skipped rather than guessed at.
func (*MemoryStorage) ListTerminalPage ¶ added in v0.15.0
func (m *MemoryStorage) ListTerminalPage(ctx context.Context, q TerminalQuery) (*TerminalPage, error)
ListTerminalPage summarizes one bounded page of finished executions.
type NestedOption ¶ added in v0.4.0
type NestedOption func(*nestedConfig)
NestedOption configures a RunNested call.
func WithNestedID ¶ added in v0.4.0
func WithNestedID(id string) NestedOption
WithNestedID sets a specific execution ID for the child saga.
func WithNestedInput ¶ added in v0.4.0
func WithNestedInput(key string, value any) NestedOption
WithNestedInput adds an initial input to the child saga.
type NestedResult ¶ added in v0.4.0
type NestedResult struct {
ExecutionID string
// contains filtered or unexported fields
}
NestedResult wraps the outputs from a completed child saga execution.
func RunNested ¶ added in v0.4.0
func RunNested(ctx context.Context, sagaName string, opts ...NestedOption) (*NestedResult, error)
RunNested executes a child saga from within a parent saga action. It reuses the parent executor's registry and storage for durability and observability. The child execution's ParentExecutionID is set to the current execution.
func (*NestedResult) Get ¶ added in v0.4.0
func (nr *NestedResult) Get(key string, target any) error
Get deserializes a named output from the child saga into target.
func (*NestedResult) Has ¶ added in v0.4.0
func (nr *NestedResult) Has(key string) bool
Has returns true if the child saga produced an output with the given key.
type Registry ¶
type Registry struct {
// contains filtered or unexported fields
}
Registry holds registered saga definitions.
func NewRegistry ¶
func NewRegistry() *Registry
NewRegistry creates a new empty registry. Useful for testing to avoid global state.
func (*Registry) Get ¶
func (r *Registry) Get(name string) (*Definition, bool)
Get retrieves a definition by name.
func (*Registry) Register ¶
func (r *Registry) Register(def *Definition) error
Register adds a definition to the registry.
type RetentionConfig ¶ added in v0.14.0
type RetentionConfig struct {
// Retention is how long a terminal execution is kept after it finished.
// Zero disables deletion, which is the escape hatch if a cluster needs its
// saga history frozen for an investigation.
Retention time.Duration
// MaxDeletes caps deletions in one sweep so an accumulated backlog drains
// over several passes rather than one thundering herd of writes. Zero means
// unbounded.
MaxDeletes int
}
RetentionConfig tunes a retention sweep.
type RetentionResult ¶ added in v0.14.0
type RetentionResult struct {
// Scanned is how many terminal executions were considered.
Scanned int
// Deleted is how many were past the retention window and removed.
Deleted int
// Failed is how many deletions errored. A sweep does not abort on one bad
// delete; the next pass retries it.
Failed int
// Skipped is how many expired executions were held back because they are
// children of a saga still in flight. They become collectable as soon as
// their parent reaches a terminal state.
Skipped int
// Capped reports that MaxDeletes stopped the sweep before it had inspected
// every terminal execution. Callers should say so rather than let a
// truncated sweep read as "everything is clean."
//
// It says the sweep did not finish looking, not that more deletions are
// certain: the executions it never reached may all be inside the retention
// window. Consuming the whole budget on the very last execution is a
// complete sweep, not a capped one.
Capped bool
}
RetentionResult reports what one sweep did.
func RunRetention ¶ added in v0.14.0
func RunRetention(ctx context.Context, storage Storage, cfg RetentionConfig, log *slog.Logger) (*RetentionResult, error)
RunRetention deletes terminal executions that finished longer ago than the configured window, and reports what it did.
The policy is one rule: a terminal execution expires on age, whether it succeeded or failed. Executions still in flight are never considered at any age, including undoing ones, which can legitimately sit unfinished for a long time while their undos keep failing and retrying. Those are exactly what recovery needs to find.
The sweep is idempotent, so a caller that is interrupted or capped simply runs again.
A nil log falls back to the default logger. Individual delete failures are logged rather than returned: one execution the store would not part with must not abandon the rest of the sweep, and the next pass retries it anyway. The caller only learns the count, so the ID has to be recorded here or an operator seeing "failed: 3" has nothing to go inspect.
type StalledConfig ¶ added in v0.15.0
type StalledConfig struct {
// StaleAfter is how long an execution may sit without changing state before
// the sweep declares it stranded. Every transition and action completion
// writes a timestamp, so this measures being driven rather than guessing at
// it. Zero disables the sweep.
StaleAfter time.Duration
// MaxForces caps transitions in one sweep so a backlog drains over several
// passes rather than one thundering herd of writes. Zero means unbounded.
MaxForces int
}
StalledConfig tunes a stalled-saga sweep.
type StalledResult ¶ added in v0.15.0
type StalledResult struct {
// Scanned is how many in-flight executions were considered.
Scanned int
// Forced is how many were past the stall window and transitioned to failed.
Forced int
// Failed is how many transitions errored. A sweep does not abort on one bad
// write; the next pass retries it.
Failed int
// Skipped is how many stalled executions were held back because they are
// children of a saga still in flight. They become forceable as soon as
// their parent stops being live.
Skipped int
// Recovered is how many refused the transition because they had moved on
// between the page that named them and the write. Not an error; a lot of
// them means the sweep is acting on pages that are too old.
Recovered int
// Capped reports that MaxForces stopped the sweep before it had inspected
// everything, not that more transitions are certain.
Capped bool
}
StalledResult reports what one sweep did.
func RunStalledSweep ¶ added in v0.15.0
func RunStalledSweep(ctx context.Context, storage StalledStorage, cfg StalledConfig, log *slog.Logger) (*StalledResult, error)
RunStalledSweep forces stranded executions to failed, and reports what it did.
An execution can end up in a state where two independent things are true: nothing will ever resume it, and nothing will ever delete it. Convergence works by name, so an execution with a generated name has nobody who will ever reconstruct that name to continue it; and RunRetention deliberately collects only terminal executions, because a pending or undoing one is exactly what recovery is supposed to find. Between those two rules a record sits in the store forever, and every recovery pass pays to read it (MIR-1788).
Transitioning rather than deleting hands the record back to rules that already exist, without having to tell the two kinds apart: one named after its entity is cleared by the next reconcile's DropIfFailed and retried, one with a generated name is collected by retention a window later. That retry is real work on a live cluster, which is why the window is generous and why a sweep that forces anything logs it.
The sweep selects candidates but decides nothing: a page is a snapshot, so ForceFailed re-decides against the record as it is when it writes. It is idempotent, so a caller that is interrupted or capped runs again. A nil log falls back to the default logger.
type StalledStorage ¶ added in v0.15.0
type StalledStorage interface {
Get(ctx context.Context, id string) (*Execution, error)
ListIncompleteSummaryPage(ctx context.Context, q IncompleteSummaryQuery) (*IncompleteSummaryPage, error)
// ForceFailed transitions an in-flight execution to failed if it is still
// in flight and still untouched since cutoff, and reports whether it did.
//
// Deciding and writing are one operation on purpose: a runner resuming an
// aged saga in the gap between them would lose its progress to a stale copy.
ForceFailed(ctx context.Context, id string, cutoff time.Time, reason string) (bool, error)
}
StalledStorage is what the sweep needs of a storage: less than Storage, plus a write conditional on the record not having changed since it was read.
Separate from Storage because EACStorage cannot honour it. The entity-access put RPC returns a revision but does not accept one, so a conditional write is not expressible over it. Narrowing the parameter is what makes the sweep's coordinator-only reach a compile error rather than a comment.
type StartBuilder ¶
type StartBuilder struct {
// contains filtered or unexported fields
}
StartBuilder provides a fluent API for starting saga executions.
func (*StartBuilder) Execute ¶
func (sb *StartBuilder) Execute(ctx context.Context) error
Execute runs the saga to completion or failure.
func (*StartBuilder) Input ¶
func (sb *StartBuilder) Input(key string, value any) *StartBuilder
Input adds an initial input value to the saga execution.
func (*StartBuilder) WithActionContext ¶ added in v0.15.0
func (sb *StartBuilder) WithActionContext(ctx context.Context) *StartBuilder
WithActionContext gives actions a cancellation boundary separate from the executor's control context. Persistence and compensation continue on the control context passed to Execute. The default is to use that same context for both, preserving existing saga behavior.
func (*StartBuilder) WithID ¶
func (sb *StartBuilder) WithID(id string) *StartBuilder
WithID sets a specific execution ID (otherwise one is generated).
Naming an execution makes Execute idempotent under that name: a second call continues the existing run rather than starting a new one. The inputs given here are the bootstrap data for the first call only. A later call that supplies different ones resumes from what the first recorded and ignores them, because the actions that already ran did so against the originals and re-deriving half a saga from new inputs would produce a run that never happened.
type Status ¶
type Status string
Status represents the current state of a saga execution.
const ( // StatusPending indicates the saga has been created but not started. StatusPending Status = "pending" // StatusRunning indicates the saga is actively executing actions. StatusRunning Status = "running" // StatusUndoing indicates the saga is rolling back due to a failure. StatusUndoing Status = "undoing" // StatusCompleted indicates all actions completed successfully. StatusCompleted Status = "completed" // StatusFailed indicates the saga failed and all undos have been attempted. StatusFailed Status = "failed" )
type Storage ¶
type Storage interface {
// Save persists the execution state.
Save(ctx context.Context, exec *Execution) error
// Get retrieves an execution by ID.
Get(ctx context.Context, id string) (*Execution, error)
// ListIncompletePage returns one bounded page of executions that need
// recovery (Pending, Running, or Undoing).
//
// Paged rather than whole because the whole set is not a size anyone
// chooses. A cluster that accumulated a six-figure backlog of incomplete
// executions made every restarting runner materialize all of them, with
// their action-output blobs, before recovery got as far as deciding which
// three it owned (MIR-1785).
//
// The walk covers several status indexes and is not pinned to one store
// revision, so an execution written mid-walk may be missed or repeated. A
// miss is recovered on the next pass, and a repeat is refused by the
// executor's own claim, which is the cheaper failure than a walk that dies
// with ErrCompacted partway through a large backlog.
ListIncompletePage(ctx context.Context, q IncompleteQuery) (*IncompletePage, error)
// ListIncompleteSummaryPage returns one bounded page of summaries of the
// same set ListIncompletePage covers.
//
// Separate from ListIncompletePage because the two callers want different
// things from the same walk. Recovery resumes what it reads and needs whole
// executions; the stalled sweep only asks how long each has sat untouched,
// and it asks about the entire in-flight set rather than the handful it
// owns. Answering that with full payloads would pay recovery's cost without
// recovering anything.
//
// A summary also carries a timestamp the execution alone cannot supply.
// Sagas written before v0.14.0 have no updated_at field at all, so their
// age has to come from the entity store's own metadata, which decoding to
// an Execution discards.
ListIncompleteSummaryPage(ctx context.Context, q IncompleteSummaryQuery) (*IncompleteSummaryPage, error)
// ListTerminalPage returns one bounded page of executions that have
// finished (Completed or Failed). It deliberately returns summaries rather
// than executions: retention only needs an ID and an age, and a backend
// holding a six-figure backlog must not have to materialize every
// action-output blob to answer.
ListTerminalPage(ctx context.Context, q TerminalQuery) (*TerminalPage, error)
// Delete removes an execution. Deleting one that is already gone is not an
// error, so a retried or overlapping sweep converges instead of failing.
Delete(ctx context.Context, id string) error
}
Storage persists saga execution state.
type TerminalExecution ¶ added in v0.14.0
type TerminalExecution struct {
// ID identifies the execution.
ID string
// FinishedAt is when the execution last changed state, which for a terminal
// execution is when it finished.
FinishedAt time.Time
// ParentID is set when this execution ran as a nested child. Retention
// needs it because a finished child is not independently safe to delete:
// its parent re-finds it by deterministic ID rather than re-running it, so
// deleting one out from under a live parent turns a resumed saga into a
// duplicated one.
ParentID string
}
TerminalExecution summarizes a finished execution for retention purposes: which one, when it stopped changing, and whose child it is.
type TerminalPage ¶ added in v0.15.0
type TerminalPage struct {
// Executions summarizes the terminal executions in this page.
Executions []TerminalExecution
// Cursor resumes the walk after this page, empty once the walk is done.
// The same caveat as IncompletePage.Cursor applies: a short page is not an
// ending, only an empty cursor is.
Cursor string
}
TerminalPage is one bounded page of finished executions.
type TerminalQuery ¶ added in v0.15.0
type TerminalQuery struct {
// Cursor resumes after the last execution of a previous page. Empty starts
// at the beginning.
Cursor string
// Limit caps how many executions the page summarizes. Zero or less means
// the backend's own cap.
Limit int
}
TerminalQuery selects one page of terminal executions.