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
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,7 @@ message ProvisionMachine {

`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
`ReduceOutput` 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:

Expand Down Expand Up @@ -106,7 +106,7 @@ func (h *handlers) CreateMachine(ctx context.Context, inv machinespb.ProvisionMa
return &machinespb.ProvisionMachine_CreateMachine{MachineId: id}, nil
}

func (h *handlers) Reduce(p *machinespb.ProvisionMachine) *machinespb.ProvisionMachineOutput {
func (h *handlers) ReduceOutput(p *machinespb.ProvisionMachine) *machinespb.ProvisionMachineOutput {
m, _ := p.State(machinespb.ProvisionMachine_CreateMachineStep)
return &machinespb.ProvisionMachineOutput{MachineId: m.GetMachineId()}
}
Expand Down
12 changes: 6 additions & 6 deletions cmd/protoc-gen-durable/internal/gen/gen.go
Original file line number Diff line number Diff line change
Expand Up @@ -162,7 +162,7 @@ func Generate(p *protogen.Plugin) error {
methods[method] = by
}
if pl.output != nil {
claim("Reduce", "the output reducer")
claim("ReduceOutput", "the output reducer")
}
if pl.failureOutput != nil {
claim("ReduceFailure", "the failure reducer")
Expand Down Expand Up @@ -278,7 +278,7 @@ func emitPipeline(g *protogen.GeneratedFile, pl *pipelineDecl) {
emitInvocationAlias(g, pl)
emitHandlers(g, pl)
if pl.output != nil {
emitReduceHelper(g, pl, "", pl.output)
emitReduceHelper(g, pl, "Output", pl.output)
}
if pl.failureOutput != nil {
emitReduceHelper(g, pl, "Failure", pl.failureOutput)
Expand Down Expand Up @@ -387,10 +387,10 @@ func emitHandlers(g *protogen.GeneratedFile, pl *pipelineDecl) {
}
}
if pl.output != nil {
g.P("// Reduce produces the pipeline output from the immutable input and")
g.P("// ReduceOutput produces the pipeline output from the immutable input and")
g.P("// committed step states on success. It must be pure: deterministic,")
g.P("// side-effect free, synchronous, and non-failing.")
g.P("Reduce(*", g.QualifiedGoIdent(pl.msg.GoIdent), ") *", g.QualifiedGoIdent(pl.output.GoIdent))
g.P("ReduceOutput(*", g.QualifiedGoIdent(pl.msg.GoIdent), ") *", g.QualifiedGoIdent(pl.output.GoIdent))
}
if pl.failureOutput != nil {
g.P("// ReduceFailure produces the pipeline failure output from the immutable")
Expand All @@ -403,7 +403,7 @@ func emitHandlers(g *protogen.GeneratedFile, pl *pipelineDecl) {
_ = name
}

// emitReduceHelper emits Reduce<Pipeline>[Failure](h, view): the fold the
// emitReduceHelper emits Reduce<Pipeline>Output/Failure(h, view): the fold the
// engine reduces through and a reducer unit test calls with a
// durabletest.NewInvocation (which is also a durable.ReduceView).
func emitReduceHelper(g *protogen.GeneratedFile, pl *pipelineDecl, kind string, out *protogen.Message) {
Expand Down Expand Up @@ -500,7 +500,7 @@ func emitDefinition(g *protogen.GeneratedFile, pl *pipelineDecl) {
}
if pl.output != nil {
g.P("Reduce: func(view ", g.QualifiedGoIdent(durablePkg.Ident("ReduceView")), ") ", protoMsg, " {")
g.P("return Reduce", name, "(h, view)")
g.P("return Reduce", name, "Output(h, view)")
g.P("},")
}
if pl.failureOutput != nil {
Expand Down
2 changes: 1 addition & 1 deletion docs/tour.md
Original file line number Diff line number Diff line change
Expand Up @@ -91,7 +91,7 @@ eng := engine.New(st)
deploy, err := deploypb.NewDeployService(&deployer{db: db, lb: lb}).Bind(eng)
// deployer has a method per step — ProvisionEnv, RunMigrations,
// ShiftTraffic — plus UnwindProvisionEnv, UnwindRunMigrations, and
// Reduce; a step it lacks is a compile error.
// ReduceOutput; a step it lacks is a compile error.
// bind every pipeline, then:
eng.Start(ctx)

Expand Down
2 changes: 1 addition & 1 deletion durabletest/doc.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,6 @@
// fake Clock for the engine's time, and a fake Invocation (NewInvocation)
// for unit-testing handlers without an engine or a store. Generated code
// accepts the fake through its NewXxxInvocation constructors and folds it
// as a reducer view through XxxReducer.Reduce. The in-memory store lives
// as a reducer view through the generated ReduceXxxOutput. The in-memory store lives
// in store/mem: it is a real driver, not a test double.
package durabletest
2 changes: 1 addition & 1 deletion examples/machines/handlers.go
Original file line number Diff line number Diff line change
Expand Up @@ -110,7 +110,7 @@ func (h *handlers) CreateMachine(ctx context.Context, inv machinespb.ProvisionMa
return &machinespb.ProvisionMachine_CreateMachine{MachineId: h.cloud.id("machine")}, nil
}

func (h *handlers) Reduce(p *machinespb.ProvisionMachine) *machinespb.ProvisionMachineOutput {
func (h *handlers) ReduceOutput(p *machinespb.ProvisionMachine) *machinespb.ProvisionMachineOutput {
machine, ok := p.State(machinespb.ProvisionMachine_CreateMachineStep)
if !ok {
panic("successful pipeline missing create-machine state")
Expand Down
12 changes: 6 additions & 6 deletions examples/machines/machinespb/machines_durable.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion examples/release-train/deploy.go
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,6 @@ func (h *deploy) ShiftTraffic(ctx context.Context, inv releasepb.DeployServiceIn
return &releasepb.DeployService_ShiftTraffic{LbGeneration: inv.Input().GetImage()}, nil
}

func (h *deploy) Reduce(d *releasepb.DeployService) *releasepb.DeployServiceOutput {
func (h *deploy) ReduceOutput(d *releasepb.DeployService) *releasepb.DeployServiceOutput {
return &releasepb.DeployServiceOutput{Url: "https://" + d.Input().GetService() + ".example.com"}
}
2 changes: 1 addition & 1 deletion examples/release-train/legacy.go
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,6 @@ func (h *legacyDeploy) ShiftTraffic(ctx context.Context, inv legacypb.DeployServ
return &legacypb.DeployService_ShiftTraffic{LbGeneration: inv.Input().GetImage()}, nil
}

func (h *legacyDeploy) Reduce(d *legacypb.DeployService) *legacypb.DeployServiceOutput {
func (h *legacyDeploy) ReduceOutput(d *legacypb.DeployService) *legacypb.DeployServiceOutput {
return &legacypb.DeployServiceOutput{Url: "https://" + d.Input().GetService() + ".example.com"}
}
12 changes: 6 additions & 6 deletions examples/release-train/legacypb/legacy_durable.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

12 changes: 6 additions & 6 deletions examples/release-train/releasepb/release_durable.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion examples/snapshots/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -167,7 +167,7 @@ func (h *snapshotter) RegisterSnapshot(ctx context.Context, inv snapshotspb.Crea
return &snapshotspb.CreateSnapshot_RegisterSnapshot{SnapshotId: id}, nil
}

func (h *snapshotter) Reduce(p *snapshotspb.CreateSnapshot) *snapshotspb.CreateSnapshotOutput {
func (h *snapshotter) ReduceOutput(p *snapshotspb.CreateSnapshot) *snapshotspb.CreateSnapshotOutput {
reg, _ := p.State(snapshotspb.CreateSnapshot_RegisterSnapshotStep)
up, _ := p.State(snapshotspb.CreateSnapshot_UploadSnapshotStep)
return &snapshotspb.CreateSnapshotOutput{SnapshotId: reg.GetSnapshotId(), ObjectKey: up.GetObjectKey()}
Expand Down
4 changes: 2 additions & 2 deletions examples/snapshots/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -96,7 +96,7 @@ func TestSnapshotUnwindsOnCatalogFull(t *testing.T) {
}

// The closures are unit-testable without an engine: the generated
// NewXxxInvocation wraps the durabletest fake, and XxxReducer.Reduce
// NewXxxInvocation wraps the durabletest fake, and ReduceXxxOutput
// folds it as a reducer view.

func TestUploadUnwindDeletesObject(t *testing.T) {
Expand Down Expand Up @@ -153,7 +153,7 @@ func TestReduceCreateSnapshot(t *testing.T) {
snapshotspb.CreateSnapshot_RegisterSnapshotStep.ID(): &snapshotspb.CreateSnapshot_RegisterSnapshot{SnapshotId: "snap-1"},
},
})
out := snapshotspb.ReduceCreateSnapshot(&snapshotter{}, view)
out := snapshotspb.ReduceCreateSnapshotOutput(&snapshotter{}, view)
if out.GetSnapshotId() != "snap-1" || out.GetObjectKey() != "k" {
t.Fatalf("Reduce = %+v", out)
}
Expand Down
12 changes: 6 additions & 6 deletions examples/snapshots/snapshotspb/snapshots_durable.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion examples/tracing-otel/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ func (h *fulfillment) Ship(ctx context.Context, inv orderspb.FulfillOrderInvocat
durable.WithUserKind(), durable.WithReason("invalid-address"))
}

func (h *fulfillment) Reduce(o *orderspb.FulfillOrder) *orderspb.FulfillOrderOutput {
func (h *fulfillment) ReduceOutput(o *orderspb.FulfillOrder) *orderspb.FulfillOrderOutput {
s, _ := o.State(orderspb.FulfillOrder_ShipStep)
return &orderspb.FulfillOrderOutput{ShipmentId: s.GetShipmentId()}
}
Expand Down
12 changes: 6 additions & 6 deletions examples/tracing-otel/orderspb/orders_durable.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 2 additions & 2 deletions spec/02-authoring.md
Original file line number Diff line number Diff line change
Expand Up @@ -259,7 +259,7 @@ type ProvisionMachineHandlers interface {
ReserveCapacity(context.Context, ProvisionMachineInvocation) (*ReserveCapacity, error)
UnwindReserveCapacity(context.Context, ProvisionMachineInvocation) error

Reduce(*ProvisionMachine) *ProvisionMachineOutput
ReduceOutput(*ProvisionMachine) *ProvisionMachineOutput
}

func NewProvisionMachine(h ProvisionMachineHandlers) *ProvisionMachineDefinition
Expand Down Expand Up @@ -317,7 +317,7 @@ annotations, park memory) that satisfies both `durable.Invocation` and
would invalidate the Run for. Generated code wraps it the same way it
wraps the Engine's: `NewProvisionMachineInvocation(fake)` is what a
`ProvisionMachineHandlers` method takes, and
`ReduceProvisionMachine(h, fake)` folds the reducer over it.
`ReduceProvisionMachineOutput(h, fake)` folds the reducer over it.

```go
inv := durabletest.NewInvocation(durabletest.InvocationConfig{
Expand Down
6 changes: 3 additions & 3 deletions spec/05-codegen.md
Original file line number Diff line number Diff line change
Expand Up @@ -51,10 +51,10 @@ Published protobuf extensions MUST use globally allocated extension numbers.
(`durable.NoInput` for an Input-less pipeline), with a
`NewXxxInvocation(core)` constructor for engine-free handler tests,
- one handler interface, `XxxHandlers`: a method per Step named after
the Step, `Unwind<Step>` for each Step that unwinds, and `Reduce` /
the Step, `Unwind<Step>` for each Step that unwinds, and `ReduceOutput` /
`ReduceFailure` when the pipeline declares outputs,
- the Pipeline constructor, `NewXxx(h XxxHandlers)`,
- `ReduceXxx(h, view)` and `ReduceXxxFailure(h, view)`, the folds the
- `ReduceXxxOutput(h, view)` and `ReduceXxxFailure(h, view)`, the folds the
engine reduces through and reducer tests call,
- runtime methods on the Pipeline marker type,
- the bound Pipeline handle,
Expand Down Expand Up @@ -97,7 +97,7 @@ Generated APIs MUST make these compile-time errors where possible:
- a missing Step method, a Step added to the pipeline included,
- a wrong method signature,
- a missing `Unwind<Step>`,
- an invalid `Reduce` signature,
- an invalid `ReduceOutput` or `ReduceFailure` signature,
- passing a stateless `StepRef` to `State`.

Generation MUST also reject a pipeline whose Step and reducer method
Expand Down