Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
118 changes: 82 additions & 36 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,74 +28,120 @@ type-erased definition it builds; the shared vocabulary lives in
`store/mem` for ephemeral runs); implementers use `store/driver`,
telemetry adapters `observe`.

**Requires Go 1.27+** (the generated `State` API uses generic methods).
**Requires Go 1.27+** (the typed `State` lookup is a generic method).

## Example

Declare a pipeline ([full example](examples/machines/)):
Declare a pipeline in protobuf (condensed from
[examples/machines](examples/machines/), which has a validation step and
a second pipeline sharing a mutex). Each step is a message nested in its pipeline, in execution order,
whose fields are the state it commits:

```proto
message ReserveCapacity {
option (durable.v1.step) = {
id: "reserve-capacity/v1"
unwind: true
};

string reservation_id = 1;
}

message ProvisionMachine {
option (durable.v1.pipeline) = {
id: "provision-machine"
input: ".machines.v1.ProvisionMachineInput"
output: ".machines.v1.ProvisionMachineOutput"

steps: ".machines.v1.Validate"
steps: ".machines.v1.SelectHost"
steps: ".machines.v1.ReserveCapacity"
steps: ".machines.v1.CreateMachine"
};

message SelectHost {
option (durable.v1.step) = {id: "select-host/v1"};
string host_id = 1;
}

message ReserveCapacity {
option (durable.v1.step) = {id: "reserve-capacity/v1" unwind: true};
string reservation_id = 1;
}

message CreateMachine {
option (durable.v1.step) = {id: "create-machine/v1"};
string machine_id = 1;
}
}
```

Implement the generated pipeline interface, one method per step:
`protoc-gen-durable` turns that into one interface, `ProvisionMachineHandlers`:
a method per step, `Unwind<Step>` for each step that unwinds, and
`Reduce` for the output. One type implements the pipeline, so its
dependencies are declared once, and a step it lacks — a step added to
the proto included — is a compile error naming the method:

```go
func (h *handlers) CreateMachine(
ctx context.Context,
inv machinespb.ProvisionMachineInvocation,
) (*machinespb.CreateMachine, error) {
reservation, ok := inv.State(machinespb.ReserveCapacityStep)
type handlers struct{ cloud *cloud }

func (h *handlers) SelectHost(ctx context.Context, inv machinespb.ProvisionMachineInvocation) (*machinespb.ProvisionMachine_SelectHost, error) {
return &machinespb.ProvisionMachine_SelectHost{HostId: "host-" + inv.Input().GetRegion() + "-1"}, nil
}

func (h *handlers) ReserveCapacity(ctx context.Context, inv machinespb.ProvisionMachineInvocation) (*machinespb.ProvisionMachine_ReserveCapacity, error) {
host, _ := inv.State(machinespb.ProvisionMachine_SelectHostStep) // typed, committed state
id, err := h.cloud.Reserve(ctx, host.GetHostId())
if err != nil {
return nil, err // an ordinary error: retried with backoff
}
return &machinespb.ProvisionMachine_ReserveCapacity{ReservationId: id}, nil
}

func (h *handlers) UnwindReserveCapacity(ctx context.Context, inv machinespb.ProvisionMachineInvocation) error {
r, ok := inv.State(machinespb.ProvisionMachine_ReserveCapacityStep)
if !ok {
return nil, durable.Fail(errors.New("reservation state unavailable"))
return nil // never committed: nothing to release
}
return h.cloud.Release(ctx, r.GetReservationId())
}

func (h *handlers) CreateMachine(ctx context.Context, inv machinespb.ProvisionMachineInvocation) (*machinespb.ProvisionMachine_CreateMachine, error) {
r, _ := inv.State(machinespb.ProvisionMachine_ReserveCapacityStep)
id, err := h.cloud.Create(ctx, r.GetReservationId(), inv.Input().GetRegion())
if errors.Is(err, errNoCapacity) {
// A decision, not an error class: the run unwinds from here,
// and ReserveCapacity's unwind releases the reservation.
return nil, durable.Fail(err, durable.WithReason("insufficient-capacity"))
}
if err != nil {
return nil, err
}
// ... at-least-once: must be idempotent; plain errors are retried.
return &machinespb.CreateMachine{MachineId: id}, nil
return &machinespb.ProvisionMachine_CreateMachine{MachineId: id}, nil
}

func (h *handlers) Reduce(p *machinespb.ProvisionMachine) *machinespb.ProvisionMachineOutput {
m, _ := p.State(machinespb.ProvisionMachine_CreateMachineStep)
return &machinespb.ProvisionMachineOutput{MachineId: m.GetMachineId()}
}
```

Wire it up — this side of an application imports `engine`, handler files never do:
Handlers run at least once, so they are idempotent; the invocation's
`ctx` dies only for engine shutdown or a cancel of the run, and
returning `ctx.Err()` is the right answer to both.

Wire it up. This side of an application imports `engine`; handler
files never do:

```go
st, _ := store.Open("bbolt:///var/lib/app/machines.db") // import _ ".../store/bbolt"
eng := engine.New(st)

provision, _ := machinespb.NewProvisionMachine(&handlers{cloud: c})
).Bind(eng)

provision, _ := machinespb.NewProvisionMachine(&handlers{cloud: c}).Bind(eng)
eng.Start(ctx)

run, created, _ := provision.Schedule(ctx, "machine-123", input)
run, _, _ := provision.Schedule(ctx, "machine-123", &machinespb.ProvisionMachineInput{Region: "ams"})
result, _ := run.Wait(ctx)
if result.Succeeded() {
switch {
case result.Succeeded():
fmt.Println(result.Output().GetMachineId())
case result.Canceled():
// run.Cancel(ctx, "operator retracted") on any handle, from any process
default:
fmt.Println(result.Failure.Reason) // "insufficient-capacity"
}
```

Handlers need not be types: every generated handler interface comes with
an `http.HandlerFunc`-style adapter, so a pipeline can be assembled from
closures over a dependency struct.
[examples/snapshots](examples/snapshots/) is written that way end to end.
A run survives the process: stop the engine mid-step, start another on
the same store, and `provision.GetRun(ctx, id)` continues where the
facts left off — under a newer pipeline definition if one shipped in
between. Handlers unit-test without an engine: hand a method
`machinespb.NewProvisionMachineInvocation(durabletest.NewInvocation(cfg))`.

For the whole story in one runnable demo — a release surviving a daemon
crash, a pipeline definition that evolves mid-flight, parent runs
Expand Down
119 changes: 56 additions & 63 deletions cmd/protoc-gen-durable/internal/gen/gen.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,41 +53,58 @@ func Generate(p *protogen.Plugin) error {
errs = append(errs, fmt.Errorf(format, args...))
}

steps := make(map[protoreflect.FullName]*stepDecl)
// Declarations are collected only from files being generated. A
// pipeline is a top-level message; its steps are the messages nested
// directly in it, in declaration order, which is the topology. A
// step anywhere else, or a pipeline nested anywhere, is an error.
var pipelines []*pipelineDecl
// Declarations are collected only from files being generated.
byID := make(map[string]protoreflect.FullName)
for _, f := range p.Files {
if !f.Generate {
continue
}
walkMessages(f.Messages, func(m *protogen.Message) {
for _, m := range f.Messages {
if so := stepOptions(m); so != nil {
if so.GetId() == "" {
fail("%s: step declaration missing id", m.Desc.FullName())
fail("%s: a step must be nested in its pipeline message", m.Desc.FullName())
}
po := pipelineOptions(m)
if po == nil {
walkMessages(m.Messages, func(n *protogen.Message) {
if stepOptions(n) != nil {
fail("%s: a step must be nested directly in a pipeline message", n.Desc.FullName())
}
if pipelineOptions(n) != nil {
fail("%s: a pipeline must be a top-level message", n.Desc.FullName())
}
})
continue
}
pl := &pipelineDecl{msg: m, file: f, opts: po}
pipelines = append(pipelines, pl)
for _, n := range m.Messages {
so := stepOptions(n)
if so == nil {
fail("%s: every message nested in a pipeline is a step; %s declares no step option", m.Desc.FullName(), n.Desc.Name())
continue
}
steps[m.Desc.FullName()] = &stepDecl{
msg: m,
opts: so,
hasState: len(m.Fields) > 0,
if pipelineOptions(n) != nil {
fail("%s: a pipeline must be a top-level message", n.Desc.FullName())
}
walkMessages(n.Messages, func(x *protogen.Message) {
if stepOptions(x) != nil || pipelineOptions(x) != nil {
fail("%s: a message nested in a step is plain data, not a step or pipeline", x.Desc.FullName())
}
})
id := so.GetId()
if id == "" {
fail("%s: step declaration missing id", n.Desc.FullName())
} else if prev, dup := byID[id]; dup {
fail("duplicate step id %q declared by %s and %s", id, prev, n.Desc.FullName())
} else {
byID[id] = n.Desc.FullName()
}
pl.steps = append(pl.steps, &stepDecl{msg: n, opts: so, hasState: len(n.Fields) > 0, owner: pl})
}
if po := pipelineOptions(m); po != nil {
pipelines = append(pipelines, &pipelineDecl{msg: m, file: f, opts: po})
}
})
}

// Duplicate StepIDs across all declarations.
byID := make(map[string]protoreflect.FullName)
for name, s := range steps {
id := s.opts.GetId()
if id == "" {
continue
}
if prev, dup := byID[id]; dup {
fail("duplicate step id %q declared by %s and %s", id, prev, name)
} else {
byID[id] = name
}
}

Expand All @@ -104,7 +121,7 @@ func Generate(p *protogen.Plugin) error {
if len(m.Fields) > 0 {
fail("%s: pipeline marker message must not declare fields", m.Desc.FullName())
}
if len(pl.opts.GetSteps()) == 0 {
if len(pl.steps) == 0 {
fail("%s: pipeline declares no steps", m.Desc.FullName())
}

Expand Down Expand Up @@ -132,34 +149,6 @@ func Generate(p *protogen.Plugin) error {
pl.failureOutput = msg
}
}

for _, ref := range pl.opts.GetSteps() {
name := trimDot(ref)
msg, ok := messages[name]
if !ok {
fail("%s: step %q not found", m.Desc.FullName(), ref)
continue
}
sd, ok := steps[protoreflect.FullName(name)]
if !ok {
fail("%s: message %q in pipeline topology is not a durable step declaration", m.Desc.FullName(), ref)
continue
}
_ = msg
if sd.owner != nil {
fail("step %s is declared by pipelines %q and %q; one durable step belongs to exactly one active pipeline",
name, sd.owner.opts.GetId(), pl.opts.GetId())
continue
}
sd.owner = pl
pl.steps = append(pl.steps, sd)
}
}

for name, s := range steps {
if s.owner == nil {
fail("step %s is not referenced by any pipeline", name)
}
}

// Every step is a method on its pipeline's handler interface, so the
Expand All @@ -179,9 +168,9 @@ func Generate(p *protogen.Plugin) error {
claim("ReduceFailure", "the failure reducer")
}
for _, s := range pl.steps {
claim(s.msg.GoIdent.GoName, "step "+s.opts.GetId())
claim(s.methodName(), "step "+s.opts.GetId())
if s.opts.GetUnwind() {
claim("Unwind"+s.msg.GoIdent.GoName, "the unwind of step "+s.opts.GetId())
claim("Unwind"+s.methodName(), "the unwind of step "+s.opts.GetId())
}
}
}
Expand Down Expand Up @@ -361,6 +350,10 @@ func emitInvocationAlias(g *protogen.GeneratedFile, pl *pipelineDecl) {
g.P()
}

// methodName is the step's handler method: the nested message's own
// name. Its Go type is Pipeline_Step, protoc-gen-go's nested naming.
func (s *stepDecl) methodName() string { return string(s.msg.Desc.Name()) }

// runSig is the signature of a step's forward method, after the name.
func (s *stepDecl) runSig(g *protogen.GeneratedFile, inv string) string {
sig := "(ctx " + g.QualifiedGoIdent(contextPkg.Ident("Context")) + ", inv " + inv + ") "
Expand All @@ -384,13 +377,13 @@ func emitHandlers(g *protogen.GeneratedFile, pl *pipelineDecl) {
g.P("// added to the pipeline is a method the next build demands.")
g.P("type ", pl.handlersName(), " interface {")
for _, s := range pl.steps {
goName := s.msg.GoIdent.GoName
g.P("// ", goName, " runs step ", strconv(s.opts.GetId()), ".")
g.P(goName, s.runSig(g, inv))
method := s.methodName()
g.P("// ", method, " runs step ", strconv(s.opts.GetId()), ".")
g.P(method, s.runSig(g, inv))
if s.opts.GetUnwind() {
g.P("// Unwind", goName, " compensates step ", strconv(s.opts.GetId()), " once it")
g.P("// Unwind", method, " compensates step ", strconv(s.opts.GetId()), " once it")
g.P("// succeeded and the run unwinds.")
g.P("Unwind", goName, "(ctx ", ctx, ", inv ", inv, ") error")
g.P("Unwind", method, "(ctx ", ctx, ", inv ", inv, ") error")
}
}
if pl.output != nil {
Expand Down Expand Up @@ -518,7 +511,7 @@ func emitDefinition(g *protogen.GeneratedFile, pl *pipelineDecl) {
g.P("Steps: []", g.QualifiedGoIdent(defPkg.Ident("Step")), "{")
typed := g.QualifiedGoIdent(durablePkg.Ident("Typed")) + "[" + pl.inputType(g) + "]"
for _, s := range pl.steps {
goName := s.msg.GoIdent.GoName
goName := s.methodName()
g.P("{")
g.P("ID: ", strconv(s.opts.GetId()), ",")
if s.opts.GetUnwind() {
Expand Down
Loading