From a3bc56f185a586501bc0bbc7e1494609c6d2f830 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Gra=C3=B1a?= Date: Mon, 14 Sep 2026 18:40:09 -0300 Subject: [PATCH 1/3] README: the example is the whole pipeline on one type MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The wiring snippet had a stray paren from the last edit and the prose still promised per-step Func adapters. The example now shows the pipeline interface end to end — select, reserve with its unwind, create with a permanent failure, reduce — then wiring, the three ways a Wait ends, restart, and testing. --- README.md | 105 ++++++++++++++++++++++++++++++++++++++++-------------- 1 file changed, 78 insertions(+), 27 deletions(-) diff --git a/README.md b/README.md index c29d8de..bb27aee 100644 --- a/README.md +++ b/README.md @@ -28,29 +28,38 @@ 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 whose fields +are the state it commits; the pipeline message lists the steps in +order: ```proto -message ReserveCapacity { - option (durable.v1.step) = { - id: "reserve-capacity/v1" - unwind: true - }; +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; +} + 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" @@ -58,44 +67,86 @@ message ProvisionMachine { } ``` -Implement the generated pipeline interface, one method per step: +`protoc-gen-durable` turns that into one interface, `ProvisionMachineHandlers`: +a method per step, `Unwind` 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.SelectHost, error) { + return &machinespb.SelectHost{HostId: "host-" + inv.Input().GetRegion() + "-1"}, nil +} + +func (h *handlers) ReserveCapacity(ctx context.Context, inv machinespb.ProvisionMachineInvocation) (*machinespb.ReserveCapacity, error) { + host, _ := inv.State(machinespb.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.ReserveCapacity{ReservationId: id}, nil +} + +func (h *handlers) UnwindReserveCapacity(ctx context.Context, inv machinespb.ProvisionMachineInvocation) error { + r, ok := inv.State(machinespb.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.CreateMachine, error) { + r, _ := inv.State(machinespb.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 } + +func (h *handlers) Reduce(p *machinespb.ProvisionMachine) *machinespb.ProvisionMachineOutput { + m, _ := p.State(machinespb.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 From 8c0f232ace327db2f3360d0962bd38e767ce8396 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Gra=C3=B1a?= Date: Mon, 14 Sep 2026 23:43:51 -0300 Subject: [PATCH 2/3] Steps nest in their pipeline message; the steps list is gone A pipeline is a top-level message with the pipeline option; its steps are the messages nested directly in it, in declaration order, which is the topology. Nothing lists a step twice, a step belongs to one pipeline by construction, and two pipelines in a package may name their steps alike. The step's Go type is protoc-gen-go's nested name, Pipeline_Step; its handler method is the short name. Generation rejects a step outside a pipeline, a non-step nested in one, or a pipeline that is not top-level. PipelineOptions loses its steps field. Examples, README, tour, and spec 02 and 05 rewritten to the nested form; the flyd stop port follows on its branch. --- README.md | 59 ++- cmd/protoc-gen-durable/internal/gen/gen.go | 119 +++--- docs/tour.md | 94 ++-- durablepb/options.pb.go | 14 +- examples/machines/handlers.go | 22 +- examples/machines/machinespb/machines.pb.go | 204 ++++----- .../machinespb/machines_durable.pb.go | 26 +- .../machines/proto/machines/v1/machines.proto | 65 ++- examples/release-train/deploy.go | 16 +- examples/release-train/legacy.go | 10 +- examples/release-train/legacypb/legacy.pb.go | 192 ++++----- .../legacypb/legacy_durable.pb.go | 18 +- .../legacyproto/legacy/v1/legacy.proto | 52 ++- .../proto/release/v1/release.proto | 100 ++--- .../release-train/releasepb/release.pb.go | 400 +++++++++--------- .../releasepb/release_durable.pb.go | 42 +- examples/release-train/train.go | 4 +- examples/snapshots/main.go | 28 +- examples/snapshots/main_test.go | 18 +- .../proto/snapshots/v1/snapshots.proto | 61 ++- .../snapshots/snapshotspb/snapshots.pb.go | 176 ++++---- .../snapshotspb/snapshots_durable.pb.go | 22 +- examples/tracing-otel/main.go | 16 +- examples/tracing-otel/orderspb/orders.pb.go | 122 +++--- .../orderspb/orders_durable.pb.go | 18 +- .../tracing-otel/proto/orders/v1/orders.proto | 50 +-- proto/durable/v1/options.proto | 2 - spec/02-authoring.md | 133 +++--- spec/05-codegen.md | 16 +- 29 files changed, 1028 insertions(+), 1071 deletions(-) diff --git a/README.md b/README.md index bb27aee..276018d 100644 --- a/README.md +++ b/README.md @@ -34,36 +34,31 @@ telemetry adapters `observe`. 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 whose fields -are the state it commits; the pipeline message lists the steps in -order: +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 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; -} - message ProvisionMachine { option (durable.v1.pipeline) = { id: "provision-machine" input: ".machines.v1.ProvisionMachineInput" output: ".machines.v1.ProvisionMachineOutput" - - 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; + } } ``` @@ -76,29 +71,29 @@ the proto included — is a compile error naming the method: ```go type handlers struct{ cloud *cloud } -func (h *handlers) SelectHost(ctx context.Context, inv machinespb.ProvisionMachineInvocation) (*machinespb.SelectHost, error) { - return &machinespb.SelectHost{HostId: "host-" + inv.Input().GetRegion() + "-1"}, nil +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.ReserveCapacity, error) { - host, _ := inv.State(machinespb.SelectHostStep) // typed, committed state +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.ReserveCapacity{ReservationId: id}, nil + return &machinespb.ProvisionMachine_ReserveCapacity{ReservationId: id}, nil } func (h *handlers) UnwindReserveCapacity(ctx context.Context, inv machinespb.ProvisionMachineInvocation) error { - r, ok := inv.State(machinespb.ReserveCapacityStep) + r, ok := inv.State(machinespb.ProvisionMachine_ReserveCapacityStep) if !ok { 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.CreateMachine, error) { - r, _ := inv.State(machinespb.ReserveCapacityStep) +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, @@ -108,11 +103,11 @@ func (h *handlers) CreateMachine(ctx context.Context, inv machinespb.ProvisionMa if err != nil { return nil, err } - 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.CreateMachineStep) + m, _ := p.State(machinespb.ProvisionMachine_CreateMachineStep) return &machinespb.ProvisionMachineOutput{MachineId: m.GetMachineId()} } ``` diff --git a/cmd/protoc-gen-durable/internal/gen/gen.go b/cmd/protoc-gen-durable/internal/gen/gen.go index 135a5ae..f8c15ac 100644 --- a/cmd/protoc-gen-durable/internal/gen/gen.go +++ b/cmd/protoc-gen-durable/internal/gen/gen.go @@ -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 } } @@ -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()) } @@ -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 @@ -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()) } } } @@ -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 + ") " @@ -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 { @@ -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() { diff --git a/docs/tour.md b/docs/tour.md index abf43b8..bcd60ee 100644 --- a/docs/tour.md +++ b/docs/tour.md @@ -42,43 +42,40 @@ The full model: [spec/01-model.md](../spec/01-model.md). ## Declaring a pipeline -Pipelines and steps are protobuf messages carrying `durable.v1` -options; `protoc-gen-durable` compiles them into a fully typed Go API +A pipeline is a protobuf message carrying the `durable.v1.pipeline` +option, with its steps nested inside it in execution order; +`protoc-gen-durable` compiles them into a fully typed Go API (see [spec/05-codegen.md](../spec/05-codegen.md) and the committed generated code in [examples/machines](../examples/machines/)): ```proto -message ProvisionEnv { - option (durable.v1.step) = { - id: "provision-env/v1" - unwind: true // has a rollback - }; - string env_id = 1; // step state: committed on success -} - -message RunMigrations { - option (durable.v1.step) = { - id: "run-migrations/v1" - unwind: true - }; - string schema_version = 1; -} - -message ShiftTraffic { - option (durable.v1.step) = { id: "shift-traffic/v1" }; - string lb_generation = 1; -} - message DeployService { option (durable.v1.pipeline) = { id: "deploy-service" input: ".deploy.v1.DeployServiceInput" output: ".deploy.v1.DeployServiceOutput" - - steps: ".deploy.v1.ProvisionEnv" - steps: ".deploy.v1.RunMigrations" - steps: ".deploy.v1.ShiftTraffic" }; + + message ProvisionEnv { + option (durable.v1.step) = { + id: "provision-env/v1" + unwind: true // has a rollback + }; + string env_id = 1; // step state: committed on success + } + + message RunMigrations { + option (durable.v1.step) = { + id: "run-migrations/v1" + unwind: true + }; + string schema_version = 1; + } + + message ShiftTraffic { + option (durable.v1.step) = { id: "shift-traffic/v1" }; + string lb_generation = 1; + } } ``` @@ -122,8 +119,8 @@ A forward handler receives the pipeline's invocation: the typed input, committed state of earlier steps, and attempt metadata. ```go -func (h *deployer) RunMigrations(ctx context.Context, inv deploypb.DeployServiceInvocation) (*deploypb.RunMigrations, error) { - env, ok := inv.State(deploypb.ProvisionEnvStep) // typed, committed state +func (h *deployer) RunMigrations(ctx context.Context, inv deploypb.DeployServiceInvocation) (*deploypb.DeployService_RunMigrations, error) { + env, ok := inv.State(deploypb.DeployService_ProvisionEnvStep) // typed, committed state if !ok { return nil, durable.Fail(errors.New("env state unavailable")) } @@ -131,7 +128,7 @@ func (h *deployer) RunMigrations(ctx context.Context, inv deploypb.DeployService if err != nil { return nil, err // plain error: retried with backoff } - return &deploypb.RunMigrations{SchemaVersion: version}, nil + return &deploypb.DeployService_RunMigrations{SchemaVersion: version}, nil } ``` @@ -176,7 +173,7 @@ handler: ```go func (h *deployer) UnwindRunMigrations(ctx context.Context, inv deploypb.DeployServiceInvocation) error { - m, ok := inv.State(deploypb.RunMigrationsStep) // what forward committed + m, ok := inv.State(deploypb.DeployService_RunMigrationsStep) // what forward committed if !ok { return nil // forward never committed; nothing to undo } @@ -250,12 +247,12 @@ that moment. cancel is pending: it just honours its context. ```go -func (h *deployer) ShiftTraffic(ctx context.Context, inv deploypb.DeployServiceInvocation) (*deploypb.ShiftTraffic, error) { +func (h *deployer) ShiftTraffic(ctx context.Context, inv deploypb.DeployServiceInvocation) (*deploypb.DeployService_ShiftTraffic, error) { gen, err := h.lb.Shift(ctx, inv.Input().GetImage()) // ctx dies on cancel if err != nil { return nil, err // under a cancel this is the cancellation; otherwise a retry } - return &deploypb.ShiftTraffic{LbGeneration: gen}, nil + return &deploypb.DeployService_ShiftTraffic{LbGeneration: gen}, nil } ``` @@ -329,11 +326,11 @@ returns the `Result` of an already-terminal run). Instead, a handler *parks*: ```go -func (h *train) ShipServices(ctx context.Context, inv releasepb.ReleaseTrainInvocation) (*releasepb.ShipServices, error) { +func (h *train) ShipServices(ctx context.Context, inv releasepb.ReleaseTrainInvocation) (*releasepb.ReleaseTrain_ShipServices, error) { if childID, woken := inv.AwaitedRunID(); woken { // The awaited run reached terminality; inspect and continue. _ = childID - return &releasepb.ShipServices{}, nil + return &releasepb.ReleaseTrain_ShipServices{}, nil } child, _, err := deploy.Schedule(ctx, "service-web", input) if err != nil { @@ -364,7 +361,7 @@ targets terminal or missing at wake time), and `Pending()` for the rest — so the handler never has to remember its children itself: ```go -func (h *train) ShipServices(ctx context.Context, inv releasepb.ReleaseTrainInvocation) (*releasepb.ShipServices, error) { +func (h *train) ShipServices(ctx context.Context, inv releasepb.ReleaseTrainInvocation) (*releasepb.ReleaseTrain_ShipServices, error) { if w, woken := inv.Awaited(); woken { for _, id := range w.Done { // all of them, under AwaitAll run, err := deploy.GetRun(ctx, id) @@ -373,7 +370,7 @@ func (h *train) ShipServices(ctx context.Context, inv releasepb.ReleaseTrainInvo if err != nil { return nil, err } if !res.Succeeded() { return nil, durable.Fail(fmt.Errorf("deploy %s failed", id)) } } - return &releasepb.ShipServices{}, nil + return &releasepb.ReleaseTrain_ShipServices{}, nil } var ids []durable.RunID for _, svc := range inv.Input().GetServices() { @@ -429,17 +426,16 @@ message DeployService { option (durable.v1.pipeline) = { id: "deploy-service" mutexes: "service-lifecycle" - steps: ".deploy.v1.ProvisionEnv" - steps: ".deploy.v1.ShiftTraffic" }; + // steps ... } message RollbackService { option (durable.v1.pipeline) = { id: "rollback-service" mutexes: "service-lifecycle" - steps: ".deploy.v1.RestorePrevious" }; + message RestorePrevious { option (durable.v1.step) = {id: "restore-previous/v1"}; } } ``` @@ -457,16 +453,16 @@ message SnapshotService { option (durable.v1.pipeline) = { id: "snapshot-service" mutexes: ["service-lifecycle", "backup-window"] - steps: ".deploy.v1.FreezeAndCopy" }; + message FreezeAndCopy { option (durable.v1.step) = {id: "freeze-and-copy/v1"}; } } message ScheduledBackup { option (durable.v1.pipeline) = { id: "scheduled-backup" mutexes: "backup-window" - steps: ".deploy.v1.CopyOffsite" }; + message CopyOffsite { option (durable.v1.step) = {id: "copy-offsite/v1"}; } } ``` @@ -555,10 +551,16 @@ furthest current-topology position the run's successful steps reach. **Adding a step.** The platform team adds canary analysis: ```proto -steps: ".deploy.v1.ProvisionEnv" -steps: ".deploy.v1.RunMigrations" -steps: ".deploy.v1.CanaryAnalysis" // new -steps: ".deploy.v1.ShiftTraffic" +message DeployService { + option (durable.v1.pipeline) = { /* ... */ }; + message ProvisionEnv { /* ... */ } + message RunMigrations { /* ... */ } + message CanaryAnalysis { // new, in position + option (durable.v1.step) = {id: "canary-analysis/v1"}; + uint32 score = 1; + } + message ShiftTraffic { /* ... */ } +} ``` A run whose frontier is `run-migrations/v1` executes the new step — diff --git a/durablepb/options.pb.go b/durablepb/options.pb.go index 5a8da9a..81f64b7 100644 --- a/durablepb/options.pb.go +++ b/durablepb/options.pb.go @@ -112,8 +112,6 @@ type PipelineOptions struct { Input string `protobuf:"bytes,2,opt,name=input,proto3" json:"input,omitempty"` // Fully-qualified Output message type. Output string `protobuf:"bytes,3,opt,name=output,proto3" json:"output,omitempty"` - // Ordered Step topology as fully-qualified Step message types. - Steps []string `protobuf:"bytes,4,rep,name=steps,proto3" json:"steps,omitempty"` // Optional default concurrency class for all of this pipeline's steps; // a step's own concurrency_class overrides it. ConcurrencyClass string `protobuf:"bytes,6,opt,name=concurrency_class,json=concurrencyClass,proto3" json:"concurrency_class,omitempty"` @@ -193,13 +191,6 @@ func (x *PipelineOptions) GetOutput() string { return "" } -func (x *PipelineOptions) GetSteps() []string { - if x != nil { - return x.Steps - } - return nil -} - func (x *PipelineOptions) GetConcurrencyClass() string { if x != nil { return x.ConcurrencyClass @@ -270,12 +261,11 @@ const file_durable_v1_options_proto_rawDesc = "" + "\x02id\x18\x01 \x01(\tR\x02id\x12\x16\n" + "\x06unwind\x18\x02 \x01(\bR\x06unwind\x12\x18\n" + "\aretired\x18\x03 \x01(\bR\aretired\x12+\n" + - "\x11concurrency_class\x18\x04 \x01(\tR\x10concurrencyClass\"\xf0\x01\n" + + "\x11concurrency_class\x18\x04 \x01(\tR\x10concurrencyClass\"\xda\x01\n" + "\x0fPipelineOptions\x12\x0e\n" + "\x02id\x18\x01 \x01(\tR\x02id\x12\x14\n" + "\x05input\x18\x02 \x01(\tR\x05input\x12\x16\n" + - "\x06output\x18\x03 \x01(\tR\x06output\x12\x14\n" + - "\x05steps\x18\x04 \x03(\tR\x05steps\x12+\n" + + "\x06output\x18\x03 \x01(\tR\x06output\x12+\n" + "\x11concurrency_class\x18\x06 \x01(\tR\x10concurrencyClass\x12\x18\n" + "\amutexes\x18\a \x03(\tR\amutexes\x12\x1b\n" + "\trun_class\x18\t \x01(\tR\brunClass\x12%\n" + diff --git a/examples/machines/handlers.go b/examples/machines/handlers.go index 5209f28..dfa4b06 100644 --- a/examples/machines/handlers.go +++ b/examples/machines/handlers.go @@ -53,14 +53,14 @@ func (h *handlers) Validate(ctx context.Context, inv machinespb.ProvisionMachine return nil } -func (h *handlers) SelectHost(ctx context.Context, inv machinespb.ProvisionMachineInvocation) (*machinespb.SelectHost, error) { - return &machinespb.SelectHost{ +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.ReserveCapacity, error) { - host, ok := inv.State(machinespb.SelectHostStep) +func (h *handlers) ReserveCapacity(ctx context.Context, inv machinespb.ProvisionMachineInvocation) (*machinespb.ProvisionMachine_ReserveCapacity, error) { + host, ok := inv.State(machinespb.ProvisionMachine_SelectHostStep) if !ok { return nil, durable.Fail(errors.New("select-host state unavailable")) } @@ -69,11 +69,11 @@ func (h *handlers) ReserveCapacity(ctx context.Context, inv machinespb.Provision id := h.cloud.id("res") h.cloud.reservations[id] = true _ = host - return &machinespb.ReserveCapacity{ReservationId: id}, nil + return &machinespb.ProvisionMachine_ReserveCapacity{ReservationId: id}, nil } func (h *handlers) UnwindReserveCapacity(ctx context.Context, inv machinespb.ProvisionMachineInvocation) error { - reservation, ok := inv.State(machinespb.ReserveCapacityStep) + reservation, ok := inv.State(machinespb.ProvisionMachine_ReserveCapacityStep) if !ok { return nil } @@ -83,8 +83,8 @@ func (h *handlers) UnwindReserveCapacity(ctx context.Context, inv machinespb.Pro return nil } -func (h *handlers) CreateMachine(ctx context.Context, inv machinespb.ProvisionMachineInvocation) (*machinespb.CreateMachine, error) { - reservation, ok := inv.State(machinespb.ReserveCapacityStep) +func (h *handlers) CreateMachine(ctx context.Context, inv machinespb.ProvisionMachineInvocation) (*machinespb.ProvisionMachine_CreateMachine, error) { + reservation, ok := inv.State(machinespb.ProvisionMachine_ReserveCapacityStep) if !ok { return nil, durable.Fail(errors.New("reservation state unavailable")) } @@ -107,15 +107,15 @@ func (h *handlers) CreateMachine(ctx context.Context, inv machinespb.ProvisionMa if !h.cloud.reservations[reservation.GetReservationId()] { return nil, errors.New("reservation not yet visible") // transient, retried } - return &machinespb.CreateMachine{MachineId: h.cloud.id("machine")}, nil + return &machinespb.ProvisionMachine_CreateMachine{MachineId: h.cloud.id("machine")}, nil } func (h *handlers) Reduce(p *machinespb.ProvisionMachine) *machinespb.ProvisionMachineOutput { - machine, ok := p.State(machinespb.CreateMachineStep) + machine, ok := p.State(machinespb.ProvisionMachine_CreateMachineStep) if !ok { panic("successful pipeline missing create-machine state") } - host, _ := p.State(machinespb.SelectHostStep) + host, _ := p.State(machinespb.ProvisionMachine_SelectHostStep) return &machinespb.ProvisionMachineOutput{ MachineId: machine.GetMachineId(), HostId: host.GetHostId(), diff --git a/examples/machines/machinespb/machines.pb.go b/examples/machines/machinespb/machines.pb.go index e3e147c..9f07859 100644 --- a/examples/machines/machinespb/machines.pb.go +++ b/examples/machines/machinespb/machines.pb.go @@ -136,26 +136,26 @@ func (x *ProvisionMachineOutput) GetHostId() string { return "" } -type Validate struct { +type ProvisionMachine struct { state protoimpl.MessageState `protogen:"open.v1"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *Validate) Reset() { - *x = Validate{} +func (x *ProvisionMachine) Reset() { + *x = ProvisionMachine{} mi := &file_machines_v1_machines_proto_msgTypes[2] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *Validate) String() string { +func (x *ProvisionMachine) String() string { return protoimpl.X.MessageStringOf(x) } -func (*Validate) ProtoMessage() {} +func (*ProvisionMachine) ProtoMessage() {} -func (x *Validate) ProtoReflect() protoreflect.Message { +func (x *ProvisionMachine) ProtoReflect() protoreflect.Message { mi := &file_machines_v1_machines_proto_msgTypes[2] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -167,32 +167,34 @@ func (x *Validate) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use Validate.ProtoReflect.Descriptor instead. -func (*Validate) Descriptor() ([]byte, []int) { +// Deprecated: Use ProvisionMachine.ProtoReflect.Descriptor instead. +func (*ProvisionMachine) Descriptor() ([]byte, []int) { return file_machines_v1_machines_proto_rawDescGZIP(), []int{2} } -type SelectHost struct { +// DecommissionMachine shares the machine-lifecycle mutex with +// ProvisionMachine: at most one of the two may have a nonterminal run per +// machine. +type DecommissionMachine struct { state protoimpl.MessageState `protogen:"open.v1"` - HostId string `protobuf:"bytes,1,opt,name=host_id,json=hostId,proto3" json:"host_id,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *SelectHost) Reset() { - *x = SelectHost{} +func (x *DecommissionMachine) Reset() { + *x = DecommissionMachine{} mi := &file_machines_v1_machines_proto_msgTypes[3] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *SelectHost) String() string { +func (x *DecommissionMachine) String() string { return protoimpl.X.MessageStringOf(x) } -func (*SelectHost) ProtoMessage() {} +func (*DecommissionMachine) ProtoMessage() {} -func (x *SelectHost) ProtoReflect() protoreflect.Message { +func (x *DecommissionMachine) ProtoReflect() protoreflect.Message { mi := &file_machines_v1_machines_proto_msgTypes[3] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -204,39 +206,31 @@ func (x *SelectHost) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use SelectHost.ProtoReflect.Descriptor instead. -func (*SelectHost) Descriptor() ([]byte, []int) { +// Deprecated: Use DecommissionMachine.ProtoReflect.Descriptor instead. +func (*DecommissionMachine) Descriptor() ([]byte, []int) { return file_machines_v1_machines_proto_rawDescGZIP(), []int{3} } -func (x *SelectHost) GetHostId() string { - if x != nil { - return x.HostId - } - return "" -} - -type ReserveCapacity struct { +type ProvisionMachine_Validate struct { state protoimpl.MessageState `protogen:"open.v1"` - ReservationId string `protobuf:"bytes,1,opt,name=reservation_id,json=reservationId,proto3" json:"reservation_id,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *ReserveCapacity) Reset() { - *x = ReserveCapacity{} +func (x *ProvisionMachine_Validate) Reset() { + *x = ProvisionMachine_Validate{} mi := &file_machines_v1_machines_proto_msgTypes[4] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *ReserveCapacity) String() string { +func (x *ProvisionMachine_Validate) String() string { return protoimpl.X.MessageStringOf(x) } -func (*ReserveCapacity) ProtoMessage() {} +func (*ProvisionMachine_Validate) ProtoMessage() {} -func (x *ReserveCapacity) ProtoReflect() protoreflect.Message { +func (x *ProvisionMachine_Validate) ProtoReflect() protoreflect.Message { mi := &file_machines_v1_machines_proto_msgTypes[4] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -248,39 +242,32 @@ func (x *ReserveCapacity) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use ReserveCapacity.ProtoReflect.Descriptor instead. -func (*ReserveCapacity) Descriptor() ([]byte, []int) { - return file_machines_v1_machines_proto_rawDescGZIP(), []int{4} +// Deprecated: Use ProvisionMachine_Validate.ProtoReflect.Descriptor instead. +func (*ProvisionMachine_Validate) Descriptor() ([]byte, []int) { + return file_machines_v1_machines_proto_rawDescGZIP(), []int{2, 0} } -func (x *ReserveCapacity) GetReservationId() string { - if x != nil { - return x.ReservationId - } - return "" -} - -type CreateMachine struct { +type ProvisionMachine_SelectHost struct { state protoimpl.MessageState `protogen:"open.v1"` - MachineId string `protobuf:"bytes,1,opt,name=machine_id,json=machineId,proto3" json:"machine_id,omitempty"` + HostId string `protobuf:"bytes,1,opt,name=host_id,json=hostId,proto3" json:"host_id,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *CreateMachine) Reset() { - *x = CreateMachine{} +func (x *ProvisionMachine_SelectHost) Reset() { + *x = ProvisionMachine_SelectHost{} mi := &file_machines_v1_machines_proto_msgTypes[5] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *CreateMachine) String() string { +func (x *ProvisionMachine_SelectHost) String() string { return protoimpl.X.MessageStringOf(x) } -func (*CreateMachine) ProtoMessage() {} +func (*ProvisionMachine_SelectHost) ProtoMessage() {} -func (x *CreateMachine) ProtoReflect() protoreflect.Message { +func (x *ProvisionMachine_SelectHost) ProtoReflect() protoreflect.Message { mi := &file_machines_v1_machines_proto_msgTypes[5] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -292,38 +279,39 @@ func (x *CreateMachine) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use CreateMachine.ProtoReflect.Descriptor instead. -func (*CreateMachine) Descriptor() ([]byte, []int) { - return file_machines_v1_machines_proto_rawDescGZIP(), []int{5} +// Deprecated: Use ProvisionMachine_SelectHost.ProtoReflect.Descriptor instead. +func (*ProvisionMachine_SelectHost) Descriptor() ([]byte, []int) { + return file_machines_v1_machines_proto_rawDescGZIP(), []int{2, 1} } -func (x *CreateMachine) GetMachineId() string { +func (x *ProvisionMachine_SelectHost) GetHostId() string { if x != nil { - return x.MachineId + return x.HostId } return "" } -type ProvisionMachine struct { +type ProvisionMachine_ReserveCapacity struct { state protoimpl.MessageState `protogen:"open.v1"` + ReservationId string `protobuf:"bytes,1,opt,name=reservation_id,json=reservationId,proto3" json:"reservation_id,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *ProvisionMachine) Reset() { - *x = ProvisionMachine{} +func (x *ProvisionMachine_ReserveCapacity) Reset() { + *x = ProvisionMachine_ReserveCapacity{} mi := &file_machines_v1_machines_proto_msgTypes[6] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *ProvisionMachine) String() string { +func (x *ProvisionMachine_ReserveCapacity) String() string { return protoimpl.X.MessageStringOf(x) } -func (*ProvisionMachine) ProtoMessage() {} +func (*ProvisionMachine_ReserveCapacity) ProtoMessage() {} -func (x *ProvisionMachine) ProtoReflect() protoreflect.Message { +func (x *ProvisionMachine_ReserveCapacity) ProtoReflect() protoreflect.Message { mi := &file_machines_v1_machines_proto_msgTypes[6] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -335,31 +323,39 @@ func (x *ProvisionMachine) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use ProvisionMachine.ProtoReflect.Descriptor instead. -func (*ProvisionMachine) Descriptor() ([]byte, []int) { - return file_machines_v1_machines_proto_rawDescGZIP(), []int{6} +// Deprecated: Use ProvisionMachine_ReserveCapacity.ProtoReflect.Descriptor instead. +func (*ProvisionMachine_ReserveCapacity) Descriptor() ([]byte, []int) { + return file_machines_v1_machines_proto_rawDescGZIP(), []int{2, 2} } -type ReleaseMachine struct { +func (x *ProvisionMachine_ReserveCapacity) GetReservationId() string { + if x != nil { + return x.ReservationId + } + return "" +} + +type ProvisionMachine_CreateMachine struct { state protoimpl.MessageState `protogen:"open.v1"` + MachineId string `protobuf:"bytes,1,opt,name=machine_id,json=machineId,proto3" json:"machine_id,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *ReleaseMachine) Reset() { - *x = ReleaseMachine{} +func (x *ProvisionMachine_CreateMachine) Reset() { + *x = ProvisionMachine_CreateMachine{} mi := &file_machines_v1_machines_proto_msgTypes[7] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *ReleaseMachine) String() string { +func (x *ProvisionMachine_CreateMachine) String() string { return protoimpl.X.MessageStringOf(x) } -func (*ReleaseMachine) ProtoMessage() {} +func (*ProvisionMachine_CreateMachine) ProtoMessage() {} -func (x *ReleaseMachine) ProtoReflect() protoreflect.Message { +func (x *ProvisionMachine_CreateMachine) ProtoReflect() protoreflect.Message { mi := &file_machines_v1_machines_proto_msgTypes[7] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -371,34 +367,38 @@ func (x *ReleaseMachine) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use ReleaseMachine.ProtoReflect.Descriptor instead. -func (*ReleaseMachine) Descriptor() ([]byte, []int) { - return file_machines_v1_machines_proto_rawDescGZIP(), []int{7} +// Deprecated: Use ProvisionMachine_CreateMachine.ProtoReflect.Descriptor instead. +func (*ProvisionMachine_CreateMachine) Descriptor() ([]byte, []int) { + return file_machines_v1_machines_proto_rawDescGZIP(), []int{2, 3} } -// DecommissionMachine shares the machine-lifecycle mutex with -// ProvisionMachine: at most one of the two may have a nonterminal run per -// machine. -type DecommissionMachine struct { +func (x *ProvisionMachine_CreateMachine) GetMachineId() string { + if x != nil { + return x.MachineId + } + return "" +} + +type DecommissionMachine_ReleaseMachine struct { state protoimpl.MessageState `protogen:"open.v1"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *DecommissionMachine) Reset() { - *x = DecommissionMachine{} +func (x *DecommissionMachine_ReleaseMachine) Reset() { + *x = DecommissionMachine_ReleaseMachine{} mi := &file_machines_v1_machines_proto_msgTypes[8] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *DecommissionMachine) String() string { +func (x *DecommissionMachine_ReleaseMachine) String() string { return protoimpl.X.MessageStringOf(x) } -func (*DecommissionMachine) ProtoMessage() {} +func (*DecommissionMachine_ReleaseMachine) ProtoMessage() {} -func (x *DecommissionMachine) ProtoReflect() protoreflect.Message { +func (x *DecommissionMachine_ReleaseMachine) ProtoReflect() protoreflect.Message { mi := &file_machines_v1_machines_proto_msgTypes[8] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -410,9 +410,9 @@ func (x *DecommissionMachine) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use DecommissionMachine.ProtoReflect.Descriptor instead. -func (*DecommissionMachine) Descriptor() ([]byte, []int) { - return file_machines_v1_machines_proto_rawDescGZIP(), []int{8} +// Deprecated: Use DecommissionMachine_ReleaseMachine.ProtoReflect.Descriptor instead. +func (*DecommissionMachine_ReleaseMachine) Descriptor() ([]byte, []int) { + return file_machines_v1_machines_proto_rawDescGZIP(), []int{3, 0} } var File_machines_v1_machines_proto protoreflect.FileDescriptor @@ -427,26 +427,26 @@ const file_machines_v1_machines_proto_rawDesc = "" + "\x16ProvisionMachineOutput\x12\x1d\n" + "\n" + "machine_id\x18\x01 \x01(\tR\tmachineId\x12\x17\n" + - "\ahost_id\x18\x02 \x01(\tR\x06hostId\"\x1d\n" + + "\ahost_id\x18\x02 \x01(\tR\x06hostId\"\x96\x03\n" + + "\x10ProvisionMachine\x1a\x1d\n" + "\bValidate:\x11\x8a\xa8\x19\r\n" + - "\vvalidate/v1\";\n" + + "\vvalidate/v1\x1a;\n" + "\n" + "SelectHost\x12\x17\n" + "\ahost_id\x18\x01 \x01(\tR\x06hostId:\x14\x8a\xa8\x19\x10\n" + - "\x0eselect-host/v1\"h\n" + + "\x0eselect-host/v1\x1ah\n" + "\x0fReserveCapacity\x12%\n" + "\x0ereservation_id\x18\x01 \x01(\tR\rreservationId:.\x8a\xa8\x19*\n" + - "\x13reserve-capacity/v1\x10\x01\"\x11host-capacity-api\"G\n" + + "\x13reserve-capacity/v1\x10\x01\"\x11host-capacity-api\x1aG\n" + "\rCreateMachine\x12\x1d\n" + "\n" + "machine_id\x18\x01 \x01(\tR\tmachineId:\x17\x8a\xa8\x19\x13\n" + - "\x11create-machine/v1\"\xf3\x01\n" + - "\x10ProvisionMachine:\xde\x01\x92\xa8\x19\xd9\x01\n" + - "\x11provision-machine\x12\".machines.v1.ProvisionMachineInput\x1a#.machines.v1.ProvisionMachineOutput\"\x15.machines.v1.Validate\"\x17.machines.v1.SelectHost\"\x1c.machines.v1.ReserveCapacity\"\x1a.machines.v1.CreateMachine:\x11machine-lifecycle\"*\n" + + "\x11create-machine/v1:s\x92\xa8\x19o\n" + + "\x11provision-machine\x12\".machines.v1.ProvisionMachineInput\x1a#.machines.v1.ProvisionMachineOutput:\x11machine-lifecycle\"p\n" + + "\x13DecommissionMachine\x1a*\n" + "\x0eReleaseMachine:\x18\x8a\xa8\x19\x14\n" + - "\x12release-machine/v1\"a\n" + - "\x13DecommissionMachine:J\x92\xa8\x19F\n" + - "\x14decommission-machine\"\x1b.machines.v1.ReleaseMachine:\x11machine-lifecycleBCZAgithub.com/dangra/durable/examples/machines/machinespb;machinespbb\x06proto3" + "\x12release-machine/v1:-\x92\xa8\x19)\n" + + "\x14decommission-machine:\x11machine-lifecycleBCZAgithub.com/dangra/durable/examples/machines/machinespb;machinespbb\x06proto3" var ( file_machines_v1_machines_proto_rawDescOnce sync.Once @@ -462,15 +462,15 @@ func file_machines_v1_machines_proto_rawDescGZIP() []byte { var file_machines_v1_machines_proto_msgTypes = make([]protoimpl.MessageInfo, 9) var file_machines_v1_machines_proto_goTypes = []any{ - (*ProvisionMachineInput)(nil), // 0: machines.v1.ProvisionMachineInput - (*ProvisionMachineOutput)(nil), // 1: machines.v1.ProvisionMachineOutput - (*Validate)(nil), // 2: machines.v1.Validate - (*SelectHost)(nil), // 3: machines.v1.SelectHost - (*ReserveCapacity)(nil), // 4: machines.v1.ReserveCapacity - (*CreateMachine)(nil), // 5: machines.v1.CreateMachine - (*ProvisionMachine)(nil), // 6: machines.v1.ProvisionMachine - (*ReleaseMachine)(nil), // 7: machines.v1.ReleaseMachine - (*DecommissionMachine)(nil), // 8: machines.v1.DecommissionMachine + (*ProvisionMachineInput)(nil), // 0: machines.v1.ProvisionMachineInput + (*ProvisionMachineOutput)(nil), // 1: machines.v1.ProvisionMachineOutput + (*ProvisionMachine)(nil), // 2: machines.v1.ProvisionMachine + (*DecommissionMachine)(nil), // 3: machines.v1.DecommissionMachine + (*ProvisionMachine_Validate)(nil), // 4: machines.v1.ProvisionMachine.Validate + (*ProvisionMachine_SelectHost)(nil), // 5: machines.v1.ProvisionMachine.SelectHost + (*ProvisionMachine_ReserveCapacity)(nil), // 6: machines.v1.ProvisionMachine.ReserveCapacity + (*ProvisionMachine_CreateMachine)(nil), // 7: machines.v1.ProvisionMachine.CreateMachine + (*DecommissionMachine_ReleaseMachine)(nil), // 8: machines.v1.DecommissionMachine.ReleaseMachine } var file_machines_v1_machines_proto_depIdxs = []int32{ 0, // [0:0] is the sub-list for method output_type diff --git a/examples/machines/machinespb/machines_durable.pb.go b/examples/machines/machinespb/machines_durable.pb.go index e90bd6f..33e48e6 100644 --- a/examples/machines/machinespb/machines_durable.pb.go +++ b/examples/machines/machinespb/machines_durable.pb.go @@ -15,18 +15,18 @@ import ( sync "sync" ) -// ValidateStep is the reference to the stateless step "validate/v1". +// ProvisionMachine_ValidateStep is the reference to the stateless step "validate/v1". // It is not accepted by State lookup. -var ValidateStep = pipelinedef.StepRef("validate/v1") +var ProvisionMachine_ValidateStep = pipelinedef.StepRef("validate/v1") -// SelectHostStep is the typed reference to the state-producing step "select-host/v1". -var SelectHostStep = pipelinedef.StateStepRef("select-host/v1", func() *SelectHost { return &SelectHost{} }) +// ProvisionMachine_SelectHostStep is the typed reference to the state-producing step "select-host/v1". +var ProvisionMachine_SelectHostStep = pipelinedef.StateStepRef("select-host/v1", func() *ProvisionMachine_SelectHost { return &ProvisionMachine_SelectHost{} }) -// ReserveCapacityStep is the typed reference to the state-producing step "reserve-capacity/v1". -var ReserveCapacityStep = pipelinedef.StateStepRef("reserve-capacity/v1", func() *ReserveCapacity { return &ReserveCapacity{} }) +// ProvisionMachine_ReserveCapacityStep is the typed reference to the state-producing step "reserve-capacity/v1". +var ProvisionMachine_ReserveCapacityStep = pipelinedef.StateStepRef("reserve-capacity/v1", func() *ProvisionMachine_ReserveCapacity { return &ProvisionMachine_ReserveCapacity{} }) -// CreateMachineStep is the typed reference to the state-producing step "create-machine/v1". -var CreateMachineStep = pipelinedef.StateStepRef("create-machine/v1", func() *CreateMachine { return &CreateMachine{} }) +// ProvisionMachine_CreateMachineStep is the typed reference to the state-producing step "create-machine/v1". +var ProvisionMachine_CreateMachineStep = pipelinedef.StateStepRef("create-machine/v1", func() *ProvisionMachine_CreateMachine { return &ProvisionMachine_CreateMachine{} }) // ProvisionMachineInvocation is the invocation every ProvisionMachineHandlers method receives: // durable.Invocation with the pipeline's Input typed. @@ -48,14 +48,14 @@ type ProvisionMachineHandlers interface { // Validate runs step "validate/v1". Validate(ctx context.Context, inv ProvisionMachineInvocation) error // SelectHost runs step "select-host/v1". - SelectHost(ctx context.Context, inv ProvisionMachineInvocation) (*SelectHost, error) + SelectHost(ctx context.Context, inv ProvisionMachineInvocation) (*ProvisionMachine_SelectHost, error) // ReserveCapacity runs step "reserve-capacity/v1". - ReserveCapacity(ctx context.Context, inv ProvisionMachineInvocation) (*ReserveCapacity, error) + ReserveCapacity(ctx context.Context, inv ProvisionMachineInvocation) (*ProvisionMachine_ReserveCapacity, error) // UnwindReserveCapacity compensates step "reserve-capacity/v1" once it // succeeded and the run unwinds. UnwindReserveCapacity(ctx context.Context, inv ProvisionMachineInvocation) error // CreateMachine runs step "create-machine/v1". - CreateMachine(ctx context.Context, inv ProvisionMachineInvocation) (*CreateMachine, error) + CreateMachine(ctx context.Context, inv ProvisionMachineInvocation) (*ProvisionMachine_CreateMachine, error) // Reduce produces the pipeline output from the immutable input and // committed step states on success. It must be pure: deterministic, // side-effect free, synchronous, and non-failing. @@ -279,9 +279,9 @@ type ProvisionMachineResult struct { // succeeded. func (r ProvisionMachineResult) Output() *ProvisionMachineOutput { return r.output } -// ReleaseMachineStep is the reference to the stateless step "release-machine/v1". +// DecommissionMachine_ReleaseMachineStep is the reference to the stateless step "release-machine/v1". // It is not accepted by State lookup. -var ReleaseMachineStep = pipelinedef.StepRef("release-machine/v1") +var DecommissionMachine_ReleaseMachineStep = pipelinedef.StepRef("release-machine/v1") // DecommissionMachineInvocation is the invocation every DecommissionMachineHandlers method receives: // durable.Invocation with the pipeline's Input typed. diff --git a/examples/machines/proto/machines/v1/machines.proto b/examples/machines/proto/machines/v1/machines.proto index cc67802..4921624 100644 --- a/examples/machines/proto/machines/v1/machines.proto +++ b/examples/machines/proto/machines/v1/machines.proto @@ -18,48 +18,39 @@ message ProvisionMachineOutput { string host_id = 2; } -message Validate { - option (durable.v1.step) = {id: "validate/v1"}; -} - -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 - concurrency_class: "host-capacity-api" - }; - - string reservation_id = 1; -} - -message CreateMachine { - option (durable.v1.step) = {id: "create-machine/v1"}; - - string machine_id = 1; -} - message ProvisionMachine { option (durable.v1.pipeline) = { id: "provision-machine" input: ".machines.v1.ProvisionMachineInput" output: ".machines.v1.ProvisionMachineOutput" mutexes: "machine-lifecycle" - - steps: ".machines.v1.Validate" - steps: ".machines.v1.SelectHost" - steps: ".machines.v1.ReserveCapacity" - steps: ".machines.v1.CreateMachine" }; -} -message ReleaseMachine { - option (durable.v1.step) = {id: "release-machine/v1"}; + message Validate { + option (durable.v1.step) = {id: "validate/v1"}; + } + + 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 + concurrency_class: "host-capacity-api" + }; + + string reservation_id = 1; + } + + message CreateMachine { + option (durable.v1.step) = {id: "create-machine/v1"}; + + string machine_id = 1; + } } // DecommissionMachine shares the machine-lifecycle mutex with @@ -69,7 +60,9 @@ message DecommissionMachine { option (durable.v1.pipeline) = { id: "decommission-machine" mutexes: "machine-lifecycle" - - steps: ".machines.v1.ReleaseMachine" }; + + message ReleaseMachine { + option (durable.v1.step) = {id: "release-machine/v1"}; + } } diff --git a/examples/release-train/deploy.go b/examples/release-train/deploy.go index 2d36e08..f86b277 100644 --- a/examples/release-train/deploy.go +++ b/examples/release-train/deploy.go @@ -12,8 +12,8 @@ import ( // step, an Unwind method for each step that rolls back, and the reducer. type deploy struct{ w *world } -func (h *deploy) ProvisionEnv(ctx context.Context, inv releasepb.DeployServiceInvocation) (*releasepb.ProvisionEnv, error) { - return &releasepb.ProvisionEnv{EnvId: h.w.provision(inv.Input().GetService())}, nil +func (h *deploy) ProvisionEnv(ctx context.Context, inv releasepb.DeployServiceInvocation) (*releasepb.DeployService_ProvisionEnv, error) { + return &releasepb.DeployService_ProvisionEnv{EnvId: h.w.provision(inv.Input().GetService())}, nil } func (h *deploy) UnwindProvisionEnv(ctx context.Context, inv releasepb.DeployServiceInvocation) error { @@ -21,9 +21,9 @@ func (h *deploy) UnwindProvisionEnv(ctx context.Context, inv releasepb.DeploySer return nil } -func (h *deploy) RunMigrations(ctx context.Context, inv releasepb.DeployServiceInvocation) (*releasepb.RunMigrations, error) { +func (h *deploy) RunMigrations(ctx context.Context, inv releasepb.DeployServiceInvocation) (*releasepb.DeployService_RunMigrations, error) { h.w.migrate(inv.Input().GetService(), inv.Input().GetImage()) - return &releasepb.RunMigrations{SchemaVersion: inv.Input().GetImage()}, nil + return &releasepb.DeployService_RunMigrations{SchemaVersion: inv.Input().GetImage()}, nil } func (h *deploy) UnwindRunMigrations(ctx context.Context, inv releasepb.DeployServiceInvocation) error { @@ -31,7 +31,7 @@ func (h *deploy) UnwindRunMigrations(ctx context.Context, inv releasepb.DeploySe return nil } -func (h *deploy) CanaryAnalysis(ctx context.Context, inv releasepb.DeployServiceInvocation) (*releasepb.CanaryAnalysis, error) { +func (h *deploy) CanaryAnalysis(ctx context.Context, inv releasepb.DeployServiceInvocation) (*releasepb.DeployService_CanaryAnalysis, error) { service := inv.Input().GetService() if service == "api" { // Hold the api canary open so the incident can strike mid-run. @@ -45,16 +45,16 @@ func (h *deploy) CanaryAnalysis(ctx context.Context, inv releasepb.DeployService h.w.canaried[service] = 98 h.w.mu.Unlock() h.w.logf("[%s] canary analysis: score 98 — a step added while this run was in flight", service) - return &releasepb.CanaryAnalysis{Score: 98}, nil + return &releasepb.DeployService_CanaryAnalysis{Score: 98}, nil } -func (h *deploy) ShiftTraffic(ctx context.Context, inv releasepb.DeployServiceInvocation) (*releasepb.ShiftTraffic, error) { +func (h *deploy) ShiftTraffic(ctx context.Context, inv releasepb.DeployServiceInvocation) (*releasepb.DeployService_ShiftTraffic, error) { service := inv.Input().GetService() h.w.mu.Lock() h.w.traffic[service] = inv.Input().GetImage() h.w.mu.Unlock() h.w.logf("[%s] traffic shifted to %s", service, inv.Input().GetImage()) - return &releasepb.ShiftTraffic{LbGeneration: inv.Input().GetImage()}, nil + return &releasepb.DeployService_ShiftTraffic{LbGeneration: inv.Input().GetImage()}, nil } func (h *deploy) Reduce(d *releasepb.DeployService) *releasepb.DeployServiceOutput { diff --git a/examples/release-train/legacy.go b/examples/release-train/legacy.go index 5e3a3f3..e99f668 100644 --- a/examples/release-train/legacy.go +++ b/examples/release-train/legacy.go @@ -11,8 +11,8 @@ import ( // legacyDeploy implements legacypb.DeployServiceHandlers. type legacyDeploy struct{ w *world } -func (h *legacyDeploy) ProvisionEnv(ctx context.Context, inv legacypb.DeployServiceInvocation) (*legacypb.ProvisionEnv, error) { - return &legacypb.ProvisionEnv{EnvId: h.w.provision(inv.Input().GetService())}, nil +func (h *legacyDeploy) ProvisionEnv(ctx context.Context, inv legacypb.DeployServiceInvocation) (*legacypb.DeployService_ProvisionEnv, error) { + return &legacypb.DeployService_ProvisionEnv{EnvId: h.w.provision(inv.Input().GetService())}, nil } func (h *legacyDeploy) UnwindProvisionEnv(ctx context.Context, inv legacypb.DeployServiceInvocation) error { @@ -24,7 +24,7 @@ func (h *legacyDeploy) UnwindProvisionEnv(ctx context.Context, inv legacypb.Depl // durable can commit the fact: the daemon dies mid-attempt. The // restarted build re-executes this operation and hits the idempotent // skip. -func (h *legacyDeploy) RunMigrations(ctx context.Context, inv legacypb.DeployServiceInvocation) (*legacypb.RunMigrations, error) { +func (h *legacyDeploy) RunMigrations(ctx context.Context, inv legacypb.DeployServiceInvocation) (*legacypb.DeployService_RunMigrations, error) { h.w.migrate(inv.Input().GetService(), inv.Input().GetImage()) close(h.w.webMigrating) <-ctx.Done() // the daemon shuts down under us @@ -37,13 +37,13 @@ func (h *legacyDeploy) UnwindRunMigrations(ctx context.Context, inv legacypb.Dep return nil } -func (h *legacyDeploy) ShiftTraffic(ctx context.Context, inv legacypb.DeployServiceInvocation) (*legacypb.ShiftTraffic, error) { +func (h *legacyDeploy) ShiftTraffic(ctx context.Context, inv legacypb.DeployServiceInvocation) (*legacypb.DeployService_ShiftTraffic, error) { // Never reached in this demo: the daemon dies before the web deploy // gets here, and the next build routes through canary analysis first. h.w.mu.Lock() h.w.traffic[inv.Input().GetService()] = inv.Input().GetImage() h.w.mu.Unlock() - return &legacypb.ShiftTraffic{LbGeneration: inv.Input().GetImage()}, nil + return &legacypb.DeployService_ShiftTraffic{LbGeneration: inv.Input().GetImage()}, nil } func (h *legacyDeploy) Reduce(d *legacypb.DeployService) *legacypb.DeployServiceOutput { diff --git a/examples/release-train/legacypb/legacy.pb.go b/examples/release-train/legacypb/legacy.pb.go index 297075a..1268602 100644 --- a/examples/release-train/legacypb/legacy.pb.go +++ b/examples/release-train/legacypb/legacy.pb.go @@ -28,27 +28,28 @@ const ( _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20) ) -type ProvisionEnv struct { +type DeployServiceInput struct { state protoimpl.MessageState `protogen:"open.v1"` - EnvId string `protobuf:"bytes,1,opt,name=env_id,json=envId,proto3" json:"env_id,omitempty"` + Service string `protobuf:"bytes,1,opt,name=service,proto3" json:"service,omitempty"` + Image string `protobuf:"bytes,2,opt,name=image,proto3" json:"image,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *ProvisionEnv) Reset() { - *x = ProvisionEnv{} +func (x *DeployServiceInput) Reset() { + *x = DeployServiceInput{} mi := &file_legacy_v1_legacy_proto_msgTypes[0] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *ProvisionEnv) String() string { +func (x *DeployServiceInput) String() string { return protoimpl.X.MessageStringOf(x) } -func (*ProvisionEnv) ProtoMessage() {} +func (*DeployServiceInput) ProtoMessage() {} -func (x *ProvisionEnv) ProtoReflect() protoreflect.Message { +func (x *DeployServiceInput) ProtoReflect() protoreflect.Message { mi := &file_legacy_v1_legacy_proto_msgTypes[0] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -60,39 +61,46 @@ func (x *ProvisionEnv) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use ProvisionEnv.ProtoReflect.Descriptor instead. -func (*ProvisionEnv) Descriptor() ([]byte, []int) { +// Deprecated: Use DeployServiceInput.ProtoReflect.Descriptor instead. +func (*DeployServiceInput) Descriptor() ([]byte, []int) { return file_legacy_v1_legacy_proto_rawDescGZIP(), []int{0} } -func (x *ProvisionEnv) GetEnvId() string { +func (x *DeployServiceInput) GetService() string { if x != nil { - return x.EnvId + return x.Service } return "" } -type RunMigrations struct { +func (x *DeployServiceInput) GetImage() string { + if x != nil { + return x.Image + } + return "" +} + +type DeployServiceOutput struct { state protoimpl.MessageState `protogen:"open.v1"` - SchemaVersion string `protobuf:"bytes,1,opt,name=schema_version,json=schemaVersion,proto3" json:"schema_version,omitempty"` + Url string `protobuf:"bytes,1,opt,name=url,proto3" json:"url,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *RunMigrations) Reset() { - *x = RunMigrations{} +func (x *DeployServiceOutput) Reset() { + *x = DeployServiceOutput{} mi := &file_legacy_v1_legacy_proto_msgTypes[1] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *RunMigrations) String() string { +func (x *DeployServiceOutput) String() string { return protoimpl.X.MessageStringOf(x) } -func (*RunMigrations) ProtoMessage() {} +func (*DeployServiceOutput) ProtoMessage() {} -func (x *RunMigrations) ProtoReflect() protoreflect.Message { +func (x *DeployServiceOutput) ProtoReflect() protoreflect.Message { mi := &file_legacy_v1_legacy_proto_msgTypes[1] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -104,39 +112,38 @@ func (x *RunMigrations) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use RunMigrations.ProtoReflect.Descriptor instead. -func (*RunMigrations) Descriptor() ([]byte, []int) { +// Deprecated: Use DeployServiceOutput.ProtoReflect.Descriptor instead. +func (*DeployServiceOutput) Descriptor() ([]byte, []int) { return file_legacy_v1_legacy_proto_rawDescGZIP(), []int{1} } -func (x *RunMigrations) GetSchemaVersion() string { +func (x *DeployServiceOutput) GetUrl() string { if x != nil { - return x.SchemaVersion + return x.Url } return "" } -type ShiftTraffic struct { +type DeployService struct { state protoimpl.MessageState `protogen:"open.v1"` - LbGeneration string `protobuf:"bytes,1,opt,name=lb_generation,json=lbGeneration,proto3" json:"lb_generation,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *ShiftTraffic) Reset() { - *x = ShiftTraffic{} +func (x *DeployService) Reset() { + *x = DeployService{} mi := &file_legacy_v1_legacy_proto_msgTypes[2] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *ShiftTraffic) String() string { +func (x *DeployService) String() string { return protoimpl.X.MessageStringOf(x) } -func (*ShiftTraffic) ProtoMessage() {} +func (*DeployService) ProtoMessage() {} -func (x *ShiftTraffic) ProtoReflect() protoreflect.Message { +func (x *DeployService) ProtoReflect() protoreflect.Message { mi := &file_legacy_v1_legacy_proto_msgTypes[2] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -148,40 +155,32 @@ func (x *ShiftTraffic) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use ShiftTraffic.ProtoReflect.Descriptor instead. -func (*ShiftTraffic) Descriptor() ([]byte, []int) { +// Deprecated: Use DeployService.ProtoReflect.Descriptor instead. +func (*DeployService) Descriptor() ([]byte, []int) { return file_legacy_v1_legacy_proto_rawDescGZIP(), []int{2} } -func (x *ShiftTraffic) GetLbGeneration() string { - if x != nil { - return x.LbGeneration - } - return "" -} - -type DeployServiceInput struct { +type DeployService_ProvisionEnv struct { state protoimpl.MessageState `protogen:"open.v1"` - Service string `protobuf:"bytes,1,opt,name=service,proto3" json:"service,omitempty"` - Image string `protobuf:"bytes,2,opt,name=image,proto3" json:"image,omitempty"` + EnvId string `protobuf:"bytes,1,opt,name=env_id,json=envId,proto3" json:"env_id,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *DeployServiceInput) Reset() { - *x = DeployServiceInput{} +func (x *DeployService_ProvisionEnv) Reset() { + *x = DeployService_ProvisionEnv{} mi := &file_legacy_v1_legacy_proto_msgTypes[3] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *DeployServiceInput) String() string { +func (x *DeployService_ProvisionEnv) String() string { return protoimpl.X.MessageStringOf(x) } -func (*DeployServiceInput) ProtoMessage() {} +func (*DeployService_ProvisionEnv) ProtoMessage() {} -func (x *DeployServiceInput) ProtoReflect() protoreflect.Message { +func (x *DeployService_ProvisionEnv) ProtoReflect() protoreflect.Message { mi := &file_legacy_v1_legacy_proto_msgTypes[3] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -193,46 +192,39 @@ func (x *DeployServiceInput) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use DeployServiceInput.ProtoReflect.Descriptor instead. -func (*DeployServiceInput) Descriptor() ([]byte, []int) { - return file_legacy_v1_legacy_proto_rawDescGZIP(), []int{3} +// Deprecated: Use DeployService_ProvisionEnv.ProtoReflect.Descriptor instead. +func (*DeployService_ProvisionEnv) Descriptor() ([]byte, []int) { + return file_legacy_v1_legacy_proto_rawDescGZIP(), []int{2, 0} } -func (x *DeployServiceInput) GetService() string { +func (x *DeployService_ProvisionEnv) GetEnvId() string { if x != nil { - return x.Service - } - return "" -} - -func (x *DeployServiceInput) GetImage() string { - if x != nil { - return x.Image + return x.EnvId } return "" } -type DeployServiceOutput struct { +type DeployService_RunMigrations struct { state protoimpl.MessageState `protogen:"open.v1"` - Url string `protobuf:"bytes,1,opt,name=url,proto3" json:"url,omitempty"` + SchemaVersion string `protobuf:"bytes,1,opt,name=schema_version,json=schemaVersion,proto3" json:"schema_version,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *DeployServiceOutput) Reset() { - *x = DeployServiceOutput{} +func (x *DeployService_RunMigrations) Reset() { + *x = DeployService_RunMigrations{} mi := &file_legacy_v1_legacy_proto_msgTypes[4] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *DeployServiceOutput) String() string { +func (x *DeployService_RunMigrations) String() string { return protoimpl.X.MessageStringOf(x) } -func (*DeployServiceOutput) ProtoMessage() {} +func (*DeployService_RunMigrations) ProtoMessage() {} -func (x *DeployServiceOutput) ProtoReflect() protoreflect.Message { +func (x *DeployService_RunMigrations) ProtoReflect() protoreflect.Message { mi := &file_legacy_v1_legacy_proto_msgTypes[4] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -244,38 +236,39 @@ func (x *DeployServiceOutput) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use DeployServiceOutput.ProtoReflect.Descriptor instead. -func (*DeployServiceOutput) Descriptor() ([]byte, []int) { - return file_legacy_v1_legacy_proto_rawDescGZIP(), []int{4} +// Deprecated: Use DeployService_RunMigrations.ProtoReflect.Descriptor instead. +func (*DeployService_RunMigrations) Descriptor() ([]byte, []int) { + return file_legacy_v1_legacy_proto_rawDescGZIP(), []int{2, 1} } -func (x *DeployServiceOutput) GetUrl() string { +func (x *DeployService_RunMigrations) GetSchemaVersion() string { if x != nil { - return x.Url + return x.SchemaVersion } return "" } -type DeployService struct { +type DeployService_ShiftTraffic struct { state protoimpl.MessageState `protogen:"open.v1"` + LbGeneration string `protobuf:"bytes,1,opt,name=lb_generation,json=lbGeneration,proto3" json:"lb_generation,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *DeployService) Reset() { - *x = DeployService{} +func (x *DeployService_ShiftTraffic) Reset() { + *x = DeployService_ShiftTraffic{} mi := &file_legacy_v1_legacy_proto_msgTypes[5] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *DeployService) String() string { +func (x *DeployService_ShiftTraffic) String() string { return protoimpl.X.MessageStringOf(x) } -func (*DeployService) ProtoMessage() {} +func (*DeployService_ShiftTraffic) ProtoMessage() {} -func (x *DeployService) ProtoReflect() protoreflect.Message { +func (x *DeployService_ShiftTraffic) ProtoReflect() protoreflect.Message { mi := &file_legacy_v1_legacy_proto_msgTypes[5] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -287,32 +280,39 @@ func (x *DeployService) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use DeployService.ProtoReflect.Descriptor instead. -func (*DeployService) Descriptor() ([]byte, []int) { - return file_legacy_v1_legacy_proto_rawDescGZIP(), []int{5} +// Deprecated: Use DeployService_ShiftTraffic.ProtoReflect.Descriptor instead. +func (*DeployService_ShiftTraffic) Descriptor() ([]byte, []int) { + return file_legacy_v1_legacy_proto_rawDescGZIP(), []int{2, 2} +} + +func (x *DeployService_ShiftTraffic) GetLbGeneration() string { + if x != nil { + return x.LbGeneration + } + return "" } var File_legacy_v1_legacy_proto protoreflect.FileDescriptor const file_legacy_v1_legacy_proto_rawDesc = "" + "\n" + - "\x16legacy/v1/legacy.proto\x12\tlegacy.v1\x1a\x18durable/v1/options.proto\"?\n" + + "\x16legacy/v1/legacy.proto\x12\tlegacy.v1\x1a\x18durable/v1/options.proto\"D\n" + + "\x12DeployServiceInput\x12\x18\n" + + "\aservice\x18\x01 \x01(\tR\aservice\x12\x14\n" + + "\x05image\x18\x02 \x01(\tR\x05image\"'\n" + + "\x13DeployServiceOutput\x12\x10\n" + + "\x03url\x18\x01 \x01(\tR\x03url\"\xc5\x02\n" + + "\rDeployService\x1a?\n" + "\fProvisionEnv\x12\x15\n" + "\x06env_id\x18\x01 \x01(\tR\x05envId:\x18\x8a\xa8\x19\x14\n" + - "\x10provision-env/v1\x10\x01\"Q\n" + + "\x10provision-env/v1\x10\x01\x1aQ\n" + "\rRunMigrations\x12%\n" + "\x0eschema_version\x18\x01 \x01(\tR\rschemaVersion:\x19\x8a\xa8\x19\x15\n" + - "\x11run-migrations/v1\x10\x01\"K\n" + + "\x11run-migrations/v1\x10\x01\x1aK\n" + "\fShiftTraffic\x12#\n" + "\rlb_generation\x18\x01 \x01(\tR\flbGeneration:\x16\x8a\xa8\x19\x12\n" + - "\x10shift-traffic/v1\"D\n" + - "\x12DeployServiceInput\x12\x18\n" + - "\aservice\x18\x01 \x01(\tR\aservice\x12\x14\n" + - "\x05image\x18\x02 \x01(\tR\x05image\"'\n" + - "\x13DeployServiceOutput\x12\x10\n" + - "\x03url\x18\x01 \x01(\tR\x03url\"\xb2\x01\n" + - "\rDeployService:\xa0\x01\x92\xa8\x19\x9b\x01\n" + - "\x0edeploy-service\x12\x1d.legacy.v1.DeployServiceInput\x1a\x1e.legacy.v1.DeployServiceOutput\"\x17.legacy.v1.ProvisionEnv\"\x18.legacy.v1.RunMigrations\"\x17.legacy.v1.ShiftTrafficB;Z9github.com/dangra/durable/examples/release-train/legacypbb\x06proto3" + "\x10shift-traffic/v1:S\x92\xa8\x19O\n" + + "\x0edeploy-service\x12\x1d.legacy.v1.DeployServiceInput\x1a\x1e.legacy.v1.DeployServiceOutputB;Z9github.com/dangra/durable/examples/release-train/legacypbb\x06proto3" var ( file_legacy_v1_legacy_proto_rawDescOnce sync.Once @@ -328,12 +328,12 @@ func file_legacy_v1_legacy_proto_rawDescGZIP() []byte { var file_legacy_v1_legacy_proto_msgTypes = make([]protoimpl.MessageInfo, 6) var file_legacy_v1_legacy_proto_goTypes = []any{ - (*ProvisionEnv)(nil), // 0: legacy.v1.ProvisionEnv - (*RunMigrations)(nil), // 1: legacy.v1.RunMigrations - (*ShiftTraffic)(nil), // 2: legacy.v1.ShiftTraffic - (*DeployServiceInput)(nil), // 3: legacy.v1.DeployServiceInput - (*DeployServiceOutput)(nil), // 4: legacy.v1.DeployServiceOutput - (*DeployService)(nil), // 5: legacy.v1.DeployService + (*DeployServiceInput)(nil), // 0: legacy.v1.DeployServiceInput + (*DeployServiceOutput)(nil), // 1: legacy.v1.DeployServiceOutput + (*DeployService)(nil), // 2: legacy.v1.DeployService + (*DeployService_ProvisionEnv)(nil), // 3: legacy.v1.DeployService.ProvisionEnv + (*DeployService_RunMigrations)(nil), // 4: legacy.v1.DeployService.RunMigrations + (*DeployService_ShiftTraffic)(nil), // 5: legacy.v1.DeployService.ShiftTraffic } var file_legacy_v1_legacy_proto_depIdxs = []int32{ 0, // [0:0] is the sub-list for method output_type diff --git a/examples/release-train/legacypb/legacy_durable.pb.go b/examples/release-train/legacypb/legacy_durable.pb.go index b53b5ed..9df1705 100644 --- a/examples/release-train/legacypb/legacy_durable.pb.go +++ b/examples/release-train/legacypb/legacy_durable.pb.go @@ -14,14 +14,14 @@ import ( sync "sync" ) -// ProvisionEnvStep is the typed reference to the state-producing step "provision-env/v1". -var ProvisionEnvStep = pipelinedef.StateStepRef("provision-env/v1", func() *ProvisionEnv { return &ProvisionEnv{} }) +// DeployService_ProvisionEnvStep is the typed reference to the state-producing step "provision-env/v1". +var DeployService_ProvisionEnvStep = pipelinedef.StateStepRef("provision-env/v1", func() *DeployService_ProvisionEnv { return &DeployService_ProvisionEnv{} }) -// RunMigrationsStep is the typed reference to the state-producing step "run-migrations/v1". -var RunMigrationsStep = pipelinedef.StateStepRef("run-migrations/v1", func() *RunMigrations { return &RunMigrations{} }) +// DeployService_RunMigrationsStep is the typed reference to the state-producing step "run-migrations/v1". +var DeployService_RunMigrationsStep = pipelinedef.StateStepRef("run-migrations/v1", func() *DeployService_RunMigrations { return &DeployService_RunMigrations{} }) -// ShiftTrafficStep is the typed reference to the state-producing step "shift-traffic/v1". -var ShiftTrafficStep = pipelinedef.StateStepRef("shift-traffic/v1", func() *ShiftTraffic { return &ShiftTraffic{} }) +// DeployService_ShiftTrafficStep is the typed reference to the state-producing step "shift-traffic/v1". +var DeployService_ShiftTrafficStep = pipelinedef.StateStepRef("shift-traffic/v1", func() *DeployService_ShiftTraffic { return &DeployService_ShiftTraffic{} }) // DeployServiceInvocation is the invocation every DeployServiceHandlers method receives: // durable.Invocation with the pipeline's Input typed. @@ -41,17 +41,17 @@ func NewDeployServiceInvocation(core durable.Invocation) DeployServiceInvocation // added to the pipeline is a method the next build demands. type DeployServiceHandlers interface { // ProvisionEnv runs step "provision-env/v1". - ProvisionEnv(ctx context.Context, inv DeployServiceInvocation) (*ProvisionEnv, error) + ProvisionEnv(ctx context.Context, inv DeployServiceInvocation) (*DeployService_ProvisionEnv, error) // UnwindProvisionEnv compensates step "provision-env/v1" once it // succeeded and the run unwinds. UnwindProvisionEnv(ctx context.Context, inv DeployServiceInvocation) error // RunMigrations runs step "run-migrations/v1". - RunMigrations(ctx context.Context, inv DeployServiceInvocation) (*RunMigrations, error) + RunMigrations(ctx context.Context, inv DeployServiceInvocation) (*DeployService_RunMigrations, error) // UnwindRunMigrations compensates step "run-migrations/v1" once it // succeeded and the run unwinds. UnwindRunMigrations(ctx context.Context, inv DeployServiceInvocation) error // ShiftTraffic runs step "shift-traffic/v1". - ShiftTraffic(ctx context.Context, inv DeployServiceInvocation) (*ShiftTraffic, error) + ShiftTraffic(ctx context.Context, inv DeployServiceInvocation) (*DeployService_ShiftTraffic, error) // Reduce produces the pipeline output from the immutable input and // committed step states on success. It must be pure: deterministic, // side-effect free, synchronous, and non-failing. diff --git a/examples/release-train/legacyproto/legacy/v1/legacy.proto b/examples/release-train/legacyproto/legacy/v1/legacy.proto index a360cf0..396a56b 100644 --- a/examples/release-train/legacyproto/legacy/v1/legacy.proto +++ b/examples/release-train/legacyproto/legacy/v1/legacy.proto @@ -11,30 +11,6 @@ import "durable/v1/options.proto"; option go_package = "github.com/dangra/durable/examples/release-train/legacypb"; -message ProvisionEnv { - option (durable.v1.step) = { - id: "provision-env/v1" - unwind: true - }; - - string env_id = 1; -} - -message RunMigrations { - option (durable.v1.step) = { - id: "run-migrations/v1" - unwind: true - }; - - string schema_version = 1; -} - -message ShiftTraffic { - option (durable.v1.step) = {id: "shift-traffic/v1"}; - - string lb_generation = 1; -} - message DeployServiceInput { string service = 1; string image = 2; @@ -49,9 +25,29 @@ message DeployService { id: "deploy-service" input: ".legacy.v1.DeployServiceInput" output: ".legacy.v1.DeployServiceOutput" - - steps: ".legacy.v1.ProvisionEnv" - steps: ".legacy.v1.RunMigrations" - steps: ".legacy.v1.ShiftTraffic" }; + + message ProvisionEnv { + option (durable.v1.step) = { + id: "provision-env/v1" + unwind: true + }; + + string env_id = 1; + } + + message RunMigrations { + option (durable.v1.step) = { + id: "run-migrations/v1" + unwind: true + }; + + string schema_version = 1; + } + + message ShiftTraffic { + option (durable.v1.step) = {id: "shift-traffic/v1"}; + + string lb_generation = 1; + } } diff --git a/examples/release-train/proto/release/v1/release.proto b/examples/release-train/proto/release/v1/release.proto index 19b871b..d8cb22c 100644 --- a/examples/release-train/proto/release/v1/release.proto +++ b/examples/release-train/proto/release/v1/release.proto @@ -11,38 +11,9 @@ option go_package = "github.com/dangra/durable/examples/release-train/releasepb" // ---- deploy-service: one service rollout ---- -message ProvisionEnv { - option (durable.v1.step) = { - id: "provision-env/v1" - unwind: true - }; - - string env_id = 1; -} - -message RunMigrations { - option (durable.v1.step) = { - id: "run-migrations/v1" - unwind: true - }; - - string schema_version = 1; -} - // CanaryAnalysis is the step added while runs were in flight: a run // whose forward frontier is run-migrations executes it; a run already // past shift-traffic never executes it retroactively. -message CanaryAnalysis { - option (durable.v1.step) = {id: "canary-analysis/v1"}; - - uint32 score = 1; -} - -message ShiftTraffic { - option (durable.v1.step) = {id: "shift-traffic/v1"}; - - string lb_generation = 1; -} message DeployServiceInput { string service = 1; @@ -58,36 +29,44 @@ message DeployService { id: "deploy-service" input: ".release.v1.DeployServiceInput" output: ".release.v1.DeployServiceOutput" - - steps: ".release.v1.ProvisionEnv" - steps: ".release.v1.RunMigrations" - steps: ".release.v1.CanaryAnalysis" - steps: ".release.v1.ShiftTraffic" }; -} -// ---- release-train: ship every service of a release ---- + message ProvisionEnv { + option (durable.v1.step) = { + id: "provision-env/v1" + unwind: true + }; -message PlanRelease { - option (durable.v1.step) = {id: "plan/v1"}; -} + string env_id = 1; + } -message ShipWeb { - option (durable.v1.step) = {id: "ship-web/v1"}; -} + message RunMigrations { + option (durable.v1.step) = { + id: "run-migrations/v1" + unwind: true + }; + + string schema_version = 1; + } + + message CanaryAnalysis { + option (durable.v1.step) = {id: "canary-analysis/v1"}; -message ShipApi { - option (durable.v1.step) = {id: "ship-api/v1"}; + uint32 score = 1; + } + + message ShiftTraffic { + option (durable.v1.step) = {id: "shift-traffic/v1"}; + + string lb_generation = 1; + } } +// ---- release-train: ship every service of a release ---- + // Announce wraps the release: publish the changelog and notify. In the // demo it never executes — the train is frozen first, and a canceled // run selects no new forward work. -message Announce { - option (durable.v1.step) = {id: "announce/v1"}; - - string changelog_url = 1; -} message ReleaseTrainInput { string image_tag = 1; @@ -97,10 +76,23 @@ message ReleaseTrain { option (durable.v1.pipeline) = { id: "release-train" input: ".release.v1.ReleaseTrainInput" - - steps: ".release.v1.PlanRelease" - steps: ".release.v1.ShipWeb" - steps: ".release.v1.ShipApi" - steps: ".release.v1.Announce" }; + + message PlanRelease { + option (durable.v1.step) = {id: "plan/v1"}; + } + + message ShipWeb { + option (durable.v1.step) = {id: "ship-web/v1"}; + } + + message ShipApi { + option (durable.v1.step) = {id: "ship-api/v1"}; + } + + message Announce { + option (durable.v1.step) = {id: "announce/v1"}; + + string changelog_url = 1; + } } diff --git a/examples/release-train/releasepb/release.pb.go b/examples/release-train/releasepb/release.pb.go index 298e067..d265985 100644 --- a/examples/release-train/releasepb/release.pb.go +++ b/examples/release-train/releasepb/release.pb.go @@ -26,27 +26,28 @@ const ( _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20) ) -type ProvisionEnv struct { +type DeployServiceInput struct { state protoimpl.MessageState `protogen:"open.v1"` - EnvId string `protobuf:"bytes,1,opt,name=env_id,json=envId,proto3" json:"env_id,omitempty"` + Service string `protobuf:"bytes,1,opt,name=service,proto3" json:"service,omitempty"` + Image string `protobuf:"bytes,2,opt,name=image,proto3" json:"image,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *ProvisionEnv) Reset() { - *x = ProvisionEnv{} +func (x *DeployServiceInput) Reset() { + *x = DeployServiceInput{} mi := &file_release_v1_release_proto_msgTypes[0] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *ProvisionEnv) String() string { +func (x *DeployServiceInput) String() string { return protoimpl.X.MessageStringOf(x) } -func (*ProvisionEnv) ProtoMessage() {} +func (*DeployServiceInput) ProtoMessage() {} -func (x *ProvisionEnv) ProtoReflect() protoreflect.Message { +func (x *DeployServiceInput) ProtoReflect() protoreflect.Message { mi := &file_release_v1_release_proto_msgTypes[0] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -58,39 +59,46 @@ func (x *ProvisionEnv) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use ProvisionEnv.ProtoReflect.Descriptor instead. -func (*ProvisionEnv) Descriptor() ([]byte, []int) { +// Deprecated: Use DeployServiceInput.ProtoReflect.Descriptor instead. +func (*DeployServiceInput) Descriptor() ([]byte, []int) { return file_release_v1_release_proto_rawDescGZIP(), []int{0} } -func (x *ProvisionEnv) GetEnvId() string { +func (x *DeployServiceInput) GetService() string { if x != nil { - return x.EnvId + return x.Service } return "" } -type RunMigrations struct { +func (x *DeployServiceInput) GetImage() string { + if x != nil { + return x.Image + } + return "" +} + +type DeployServiceOutput struct { state protoimpl.MessageState `protogen:"open.v1"` - SchemaVersion string `protobuf:"bytes,1,opt,name=schema_version,json=schemaVersion,proto3" json:"schema_version,omitempty"` + Url string `protobuf:"bytes,1,opt,name=url,proto3" json:"url,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *RunMigrations) Reset() { - *x = RunMigrations{} +func (x *DeployServiceOutput) Reset() { + *x = DeployServiceOutput{} mi := &file_release_v1_release_proto_msgTypes[1] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *RunMigrations) String() string { +func (x *DeployServiceOutput) String() string { return protoimpl.X.MessageStringOf(x) } -func (*RunMigrations) ProtoMessage() {} +func (*DeployServiceOutput) ProtoMessage() {} -func (x *RunMigrations) ProtoReflect() protoreflect.Message { +func (x *DeployServiceOutput) ProtoReflect() protoreflect.Message { mi := &file_release_v1_release_proto_msgTypes[1] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -102,42 +110,38 @@ func (x *RunMigrations) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use RunMigrations.ProtoReflect.Descriptor instead. -func (*RunMigrations) Descriptor() ([]byte, []int) { +// Deprecated: Use DeployServiceOutput.ProtoReflect.Descriptor instead. +func (*DeployServiceOutput) Descriptor() ([]byte, []int) { return file_release_v1_release_proto_rawDescGZIP(), []int{1} } -func (x *RunMigrations) GetSchemaVersion() string { +func (x *DeployServiceOutput) GetUrl() string { if x != nil { - return x.SchemaVersion + return x.Url } return "" } -// CanaryAnalysis is the step added while runs were in flight: a run -// whose forward frontier is run-migrations executes it; a run already -// past shift-traffic never executes it retroactively. -type CanaryAnalysis struct { +type DeployService struct { state protoimpl.MessageState `protogen:"open.v1"` - Score uint32 `protobuf:"varint,1,opt,name=score,proto3" json:"score,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *CanaryAnalysis) Reset() { - *x = CanaryAnalysis{} +func (x *DeployService) Reset() { + *x = DeployService{} mi := &file_release_v1_release_proto_msgTypes[2] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *CanaryAnalysis) String() string { +func (x *DeployService) String() string { return protoimpl.X.MessageStringOf(x) } -func (*CanaryAnalysis) ProtoMessage() {} +func (*DeployService) ProtoMessage() {} -func (x *CanaryAnalysis) ProtoReflect() protoreflect.Message { +func (x *DeployService) ProtoReflect() protoreflect.Message { mi := &file_release_v1_release_proto_msgTypes[2] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -149,39 +153,32 @@ func (x *CanaryAnalysis) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use CanaryAnalysis.ProtoReflect.Descriptor instead. -func (*CanaryAnalysis) Descriptor() ([]byte, []int) { +// Deprecated: Use DeployService.ProtoReflect.Descriptor instead. +func (*DeployService) Descriptor() ([]byte, []int) { return file_release_v1_release_proto_rawDescGZIP(), []int{2} } -func (x *CanaryAnalysis) GetScore() uint32 { - if x != nil { - return x.Score - } - return 0 -} - -type ShiftTraffic struct { +type ReleaseTrainInput struct { state protoimpl.MessageState `protogen:"open.v1"` - LbGeneration string `protobuf:"bytes,1,opt,name=lb_generation,json=lbGeneration,proto3" json:"lb_generation,omitempty"` + ImageTag string `protobuf:"bytes,1,opt,name=image_tag,json=imageTag,proto3" json:"image_tag,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *ShiftTraffic) Reset() { - *x = ShiftTraffic{} +func (x *ReleaseTrainInput) Reset() { + *x = ReleaseTrainInput{} mi := &file_release_v1_release_proto_msgTypes[3] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *ShiftTraffic) String() string { +func (x *ReleaseTrainInput) String() string { return protoimpl.X.MessageStringOf(x) } -func (*ShiftTraffic) ProtoMessage() {} +func (*ReleaseTrainInput) ProtoMessage() {} -func (x *ShiftTraffic) ProtoReflect() protoreflect.Message { +func (x *ReleaseTrainInput) ProtoReflect() protoreflect.Message { mi := &file_release_v1_release_proto_msgTypes[3] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -193,40 +190,38 @@ func (x *ShiftTraffic) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use ShiftTraffic.ProtoReflect.Descriptor instead. -func (*ShiftTraffic) Descriptor() ([]byte, []int) { +// Deprecated: Use ReleaseTrainInput.ProtoReflect.Descriptor instead. +func (*ReleaseTrainInput) Descriptor() ([]byte, []int) { return file_release_v1_release_proto_rawDescGZIP(), []int{3} } -func (x *ShiftTraffic) GetLbGeneration() string { +func (x *ReleaseTrainInput) GetImageTag() string { if x != nil { - return x.LbGeneration + return x.ImageTag } return "" } -type DeployServiceInput struct { +type ReleaseTrain struct { state protoimpl.MessageState `protogen:"open.v1"` - Service string `protobuf:"bytes,1,opt,name=service,proto3" json:"service,omitempty"` - Image string `protobuf:"bytes,2,opt,name=image,proto3" json:"image,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *DeployServiceInput) Reset() { - *x = DeployServiceInput{} +func (x *ReleaseTrain) Reset() { + *x = ReleaseTrain{} mi := &file_release_v1_release_proto_msgTypes[4] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *DeployServiceInput) String() string { +func (x *ReleaseTrain) String() string { return protoimpl.X.MessageStringOf(x) } -func (*DeployServiceInput) ProtoMessage() {} +func (*ReleaseTrain) ProtoMessage() {} -func (x *DeployServiceInput) ProtoReflect() protoreflect.Message { +func (x *ReleaseTrain) ProtoReflect() protoreflect.Message { mi := &file_release_v1_release_proto_msgTypes[4] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -238,46 +233,32 @@ func (x *DeployServiceInput) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use DeployServiceInput.ProtoReflect.Descriptor instead. -func (*DeployServiceInput) Descriptor() ([]byte, []int) { +// Deprecated: Use ReleaseTrain.ProtoReflect.Descriptor instead. +func (*ReleaseTrain) Descriptor() ([]byte, []int) { return file_release_v1_release_proto_rawDescGZIP(), []int{4} } -func (x *DeployServiceInput) GetService() string { - if x != nil { - return x.Service - } - return "" -} - -func (x *DeployServiceInput) GetImage() string { - if x != nil { - return x.Image - } - return "" -} - -type DeployServiceOutput struct { +type DeployService_ProvisionEnv struct { state protoimpl.MessageState `protogen:"open.v1"` - Url string `protobuf:"bytes,1,opt,name=url,proto3" json:"url,omitempty"` + EnvId string `protobuf:"bytes,1,opt,name=env_id,json=envId,proto3" json:"env_id,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *DeployServiceOutput) Reset() { - *x = DeployServiceOutput{} +func (x *DeployService_ProvisionEnv) Reset() { + *x = DeployService_ProvisionEnv{} mi := &file_release_v1_release_proto_msgTypes[5] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *DeployServiceOutput) String() string { +func (x *DeployService_ProvisionEnv) String() string { return protoimpl.X.MessageStringOf(x) } -func (*DeployServiceOutput) ProtoMessage() {} +func (*DeployService_ProvisionEnv) ProtoMessage() {} -func (x *DeployServiceOutput) ProtoReflect() protoreflect.Message { +func (x *DeployService_ProvisionEnv) ProtoReflect() protoreflect.Message { mi := &file_release_v1_release_proto_msgTypes[5] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -289,38 +270,39 @@ func (x *DeployServiceOutput) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use DeployServiceOutput.ProtoReflect.Descriptor instead. -func (*DeployServiceOutput) Descriptor() ([]byte, []int) { - return file_release_v1_release_proto_rawDescGZIP(), []int{5} +// Deprecated: Use DeployService_ProvisionEnv.ProtoReflect.Descriptor instead. +func (*DeployService_ProvisionEnv) Descriptor() ([]byte, []int) { + return file_release_v1_release_proto_rawDescGZIP(), []int{2, 0} } -func (x *DeployServiceOutput) GetUrl() string { +func (x *DeployService_ProvisionEnv) GetEnvId() string { if x != nil { - return x.Url + return x.EnvId } return "" } -type DeployService struct { +type DeployService_RunMigrations struct { state protoimpl.MessageState `protogen:"open.v1"` + SchemaVersion string `protobuf:"bytes,1,opt,name=schema_version,json=schemaVersion,proto3" json:"schema_version,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *DeployService) Reset() { - *x = DeployService{} +func (x *DeployService_RunMigrations) Reset() { + *x = DeployService_RunMigrations{} mi := &file_release_v1_release_proto_msgTypes[6] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *DeployService) String() string { +func (x *DeployService_RunMigrations) String() string { return protoimpl.X.MessageStringOf(x) } -func (*DeployService) ProtoMessage() {} +func (*DeployService_RunMigrations) ProtoMessage() {} -func (x *DeployService) ProtoReflect() protoreflect.Message { +func (x *DeployService_RunMigrations) ProtoReflect() protoreflect.Message { mi := &file_release_v1_release_proto_msgTypes[6] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -332,31 +314,39 @@ func (x *DeployService) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use DeployService.ProtoReflect.Descriptor instead. -func (*DeployService) Descriptor() ([]byte, []int) { - return file_release_v1_release_proto_rawDescGZIP(), []int{6} +// Deprecated: Use DeployService_RunMigrations.ProtoReflect.Descriptor instead. +func (*DeployService_RunMigrations) Descriptor() ([]byte, []int) { + return file_release_v1_release_proto_rawDescGZIP(), []int{2, 1} +} + +func (x *DeployService_RunMigrations) GetSchemaVersion() string { + if x != nil { + return x.SchemaVersion + } + return "" } -type PlanRelease struct { +type DeployService_CanaryAnalysis struct { state protoimpl.MessageState `protogen:"open.v1"` + Score uint32 `protobuf:"varint,1,opt,name=score,proto3" json:"score,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *PlanRelease) Reset() { - *x = PlanRelease{} +func (x *DeployService_CanaryAnalysis) Reset() { + *x = DeployService_CanaryAnalysis{} mi := &file_release_v1_release_proto_msgTypes[7] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *PlanRelease) String() string { +func (x *DeployService_CanaryAnalysis) String() string { return protoimpl.X.MessageStringOf(x) } -func (*PlanRelease) ProtoMessage() {} +func (*DeployService_CanaryAnalysis) ProtoMessage() {} -func (x *PlanRelease) ProtoReflect() protoreflect.Message { +func (x *DeployService_CanaryAnalysis) ProtoReflect() protoreflect.Message { mi := &file_release_v1_release_proto_msgTypes[7] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -368,31 +358,39 @@ func (x *PlanRelease) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use PlanRelease.ProtoReflect.Descriptor instead. -func (*PlanRelease) Descriptor() ([]byte, []int) { - return file_release_v1_release_proto_rawDescGZIP(), []int{7} +// Deprecated: Use DeployService_CanaryAnalysis.ProtoReflect.Descriptor instead. +func (*DeployService_CanaryAnalysis) Descriptor() ([]byte, []int) { + return file_release_v1_release_proto_rawDescGZIP(), []int{2, 2} } -type ShipWeb struct { +func (x *DeployService_CanaryAnalysis) GetScore() uint32 { + if x != nil { + return x.Score + } + return 0 +} + +type DeployService_ShiftTraffic struct { state protoimpl.MessageState `protogen:"open.v1"` + LbGeneration string `protobuf:"bytes,1,opt,name=lb_generation,json=lbGeneration,proto3" json:"lb_generation,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *ShipWeb) Reset() { - *x = ShipWeb{} +func (x *DeployService_ShiftTraffic) Reset() { + *x = DeployService_ShiftTraffic{} mi := &file_release_v1_release_proto_msgTypes[8] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *ShipWeb) String() string { +func (x *DeployService_ShiftTraffic) String() string { return protoimpl.X.MessageStringOf(x) } -func (*ShipWeb) ProtoMessage() {} +func (*DeployService_ShiftTraffic) ProtoMessage() {} -func (x *ShipWeb) ProtoReflect() protoreflect.Message { +func (x *DeployService_ShiftTraffic) ProtoReflect() protoreflect.Message { mi := &file_release_v1_release_proto_msgTypes[8] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -404,31 +402,38 @@ func (x *ShipWeb) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use ShipWeb.ProtoReflect.Descriptor instead. -func (*ShipWeb) Descriptor() ([]byte, []int) { - return file_release_v1_release_proto_rawDescGZIP(), []int{8} +// Deprecated: Use DeployService_ShiftTraffic.ProtoReflect.Descriptor instead. +func (*DeployService_ShiftTraffic) Descriptor() ([]byte, []int) { + return file_release_v1_release_proto_rawDescGZIP(), []int{2, 3} } -type ShipApi struct { +func (x *DeployService_ShiftTraffic) GetLbGeneration() string { + if x != nil { + return x.LbGeneration + } + return "" +} + +type ReleaseTrain_PlanRelease struct { state protoimpl.MessageState `protogen:"open.v1"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *ShipApi) Reset() { - *x = ShipApi{} +func (x *ReleaseTrain_PlanRelease) Reset() { + *x = ReleaseTrain_PlanRelease{} mi := &file_release_v1_release_proto_msgTypes[9] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *ShipApi) String() string { +func (x *ReleaseTrain_PlanRelease) String() string { return protoimpl.X.MessageStringOf(x) } -func (*ShipApi) ProtoMessage() {} +func (*ReleaseTrain_PlanRelease) ProtoMessage() {} -func (x *ShipApi) ProtoReflect() protoreflect.Message { +func (x *ReleaseTrain_PlanRelease) ProtoReflect() protoreflect.Message { mi := &file_release_v1_release_proto_msgTypes[9] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -440,35 +445,31 @@ func (x *ShipApi) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use ShipApi.ProtoReflect.Descriptor instead. -func (*ShipApi) Descriptor() ([]byte, []int) { - return file_release_v1_release_proto_rawDescGZIP(), []int{9} +// Deprecated: Use ReleaseTrain_PlanRelease.ProtoReflect.Descriptor instead. +func (*ReleaseTrain_PlanRelease) Descriptor() ([]byte, []int) { + return file_release_v1_release_proto_rawDescGZIP(), []int{4, 0} } -// Announce wraps the release: publish the changelog and notify. In the -// demo it never executes — the train is frozen first, and a canceled -// run selects no new forward work. -type Announce struct { +type ReleaseTrain_ShipWeb struct { state protoimpl.MessageState `protogen:"open.v1"` - ChangelogUrl string `protobuf:"bytes,1,opt,name=changelog_url,json=changelogUrl,proto3" json:"changelog_url,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *Announce) Reset() { - *x = Announce{} +func (x *ReleaseTrain_ShipWeb) Reset() { + *x = ReleaseTrain_ShipWeb{} mi := &file_release_v1_release_proto_msgTypes[10] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *Announce) String() string { +func (x *ReleaseTrain_ShipWeb) String() string { return protoimpl.X.MessageStringOf(x) } -func (*Announce) ProtoMessage() {} +func (*ReleaseTrain_ShipWeb) ProtoMessage() {} -func (x *Announce) ProtoReflect() protoreflect.Message { +func (x *ReleaseTrain_ShipWeb) ProtoReflect() protoreflect.Message { mi := &file_release_v1_release_proto_msgTypes[10] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -480,39 +481,31 @@ func (x *Announce) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use Announce.ProtoReflect.Descriptor instead. -func (*Announce) Descriptor() ([]byte, []int) { - return file_release_v1_release_proto_rawDescGZIP(), []int{10} -} - -func (x *Announce) GetChangelogUrl() string { - if x != nil { - return x.ChangelogUrl - } - return "" +// Deprecated: Use ReleaseTrain_ShipWeb.ProtoReflect.Descriptor instead. +func (*ReleaseTrain_ShipWeb) Descriptor() ([]byte, []int) { + return file_release_v1_release_proto_rawDescGZIP(), []int{4, 1} } -type ReleaseTrainInput struct { +type ReleaseTrain_ShipApi struct { state protoimpl.MessageState `protogen:"open.v1"` - ImageTag string `protobuf:"bytes,1,opt,name=image_tag,json=imageTag,proto3" json:"image_tag,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *ReleaseTrainInput) Reset() { - *x = ReleaseTrainInput{} +func (x *ReleaseTrain_ShipApi) Reset() { + *x = ReleaseTrain_ShipApi{} mi := &file_release_v1_release_proto_msgTypes[11] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *ReleaseTrainInput) String() string { +func (x *ReleaseTrain_ShipApi) String() string { return protoimpl.X.MessageStringOf(x) } -func (*ReleaseTrainInput) ProtoMessage() {} +func (*ReleaseTrain_ShipApi) ProtoMessage() {} -func (x *ReleaseTrainInput) ProtoReflect() protoreflect.Message { +func (x *ReleaseTrain_ShipApi) ProtoReflect() protoreflect.Message { mi := &file_release_v1_release_proto_msgTypes[11] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -524,38 +517,32 @@ func (x *ReleaseTrainInput) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use ReleaseTrainInput.ProtoReflect.Descriptor instead. -func (*ReleaseTrainInput) Descriptor() ([]byte, []int) { - return file_release_v1_release_proto_rawDescGZIP(), []int{11} +// Deprecated: Use ReleaseTrain_ShipApi.ProtoReflect.Descriptor instead. +func (*ReleaseTrain_ShipApi) Descriptor() ([]byte, []int) { + return file_release_v1_release_proto_rawDescGZIP(), []int{4, 2} } -func (x *ReleaseTrainInput) GetImageTag() string { - if x != nil { - return x.ImageTag - } - return "" -} - -type ReleaseTrain struct { +type ReleaseTrain_Announce struct { state protoimpl.MessageState `protogen:"open.v1"` + ChangelogUrl string `protobuf:"bytes,1,opt,name=changelog_url,json=changelogUrl,proto3" json:"changelog_url,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } -func (x *ReleaseTrain) Reset() { - *x = ReleaseTrain{} +func (x *ReleaseTrain_Announce) Reset() { + *x = ReleaseTrain_Announce{} mi := &file_release_v1_release_proto_msgTypes[12] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } -func (x *ReleaseTrain) String() string { +func (x *ReleaseTrain_Announce) String() string { return protoimpl.X.MessageStringOf(x) } -func (*ReleaseTrain) ProtoMessage() {} +func (*ReleaseTrain_Announce) ProtoMessage() {} -func (x *ReleaseTrain) ProtoReflect() protoreflect.Message { +func (x *ReleaseTrain_Announce) ProtoReflect() protoreflect.Message { mi := &file_release_v1_release_proto_msgTypes[12] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) @@ -567,9 +554,16 @@ func (x *ReleaseTrain) ProtoReflect() protoreflect.Message { return mi.MessageOf(x) } -// Deprecated: Use ReleaseTrain.ProtoReflect.Descriptor instead. -func (*ReleaseTrain) Descriptor() ([]byte, []int) { - return file_release_v1_release_proto_rawDescGZIP(), []int{12} +// Deprecated: Use ReleaseTrain_Announce.ProtoReflect.Descriptor instead. +func (*ReleaseTrain_Announce) Descriptor() ([]byte, []int) { + return file_release_v1_release_proto_rawDescGZIP(), []int{4, 3} +} + +func (x *ReleaseTrain_Announce) GetChangelogUrl() string { + if x != nil { + return x.ChangelogUrl + } + return "" } var File_release_v1_release_proto protoreflect.FileDescriptor @@ -577,39 +571,39 @@ var File_release_v1_release_proto protoreflect.FileDescriptor const file_release_v1_release_proto_rawDesc = "" + "\n" + "\x18release/v1/release.proto\x12\n" + - "release.v1\x1a\x18durable/v1/options.proto\"?\n" + + "release.v1\x1a\x18durable/v1/options.proto\"D\n" + + "\x12DeployServiceInput\x12\x18\n" + + "\aservice\x18\x01 \x01(\tR\aservice\x12\x14\n" + + "\x05image\x18\x02 \x01(\tR\x05image\"'\n" + + "\x13DeployServiceOutput\x12\x10\n" + + "\x03url\x18\x01 \x01(\tR\x03url\"\x89\x03\n" + + "\rDeployService\x1a?\n" + "\fProvisionEnv\x12\x15\n" + "\x06env_id\x18\x01 \x01(\tR\x05envId:\x18\x8a\xa8\x19\x14\n" + - "\x10provision-env/v1\x10\x01\"Q\n" + + "\x10provision-env/v1\x10\x01\x1aQ\n" + "\rRunMigrations\x12%\n" + "\x0eschema_version\x18\x01 \x01(\tR\rschemaVersion:\x19\x8a\xa8\x19\x15\n" + - "\x11run-migrations/v1\x10\x01\"@\n" + + "\x11run-migrations/v1\x10\x01\x1a@\n" + "\x0eCanaryAnalysis\x12\x14\n" + "\x05score\x18\x01 \x01(\rR\x05score:\x18\x8a\xa8\x19\x14\n" + - "\x12canary-analysis/v1\"K\n" + + "\x12canary-analysis/v1\x1aK\n" + "\fShiftTraffic\x12#\n" + "\rlb_generation\x18\x01 \x01(\tR\flbGeneration:\x16\x8a\xa8\x19\x12\n" + - "\x10shift-traffic/v1\"D\n" + - "\x12DeployServiceInput\x12\x18\n" + - "\aservice\x18\x01 \x01(\tR\aservice\x12\x14\n" + - "\x05image\x18\x02 \x01(\tR\x05image\"'\n" + - "\x13DeployServiceOutput\x12\x10\n" + - "\x03url\x18\x01 \x01(\tR\x03url\"\xd3\x01\n" + - "\rDeployService:\xc1\x01\x92\xa8\x19\xbc\x01\n" + - "\x0edeploy-service\x12\x1e.release.v1.DeployServiceInput\x1a\x1f.release.v1.DeployServiceOutput\"\x18.release.v1.ProvisionEnv\"\x19.release.v1.RunMigrations\"\x1a.release.v1.CanaryAnalysis\"\x18.release.v1.ShiftTraffic\"\x1c\n" + + "\x10shift-traffic/v1:U\x92\xa8\x19Q\n" + + "\x0edeploy-service\x12\x1e.release.v1.DeployServiceInput\x1a\x1f.release.v1.DeployServiceOutput\"0\n" + + "\x11ReleaseTrainInput\x12\x1b\n" + + "\timage_tag\x18\x01 \x01(\tR\bimageTag\"\xe0\x01\n" + + "\fReleaseTrain\x1a\x1c\n" + "\vPlanRelease:\r\x8a\xa8\x19\t\n" + - "\aplan/v1\"\x1c\n" + + "\aplan/v1\x1a\x1c\n" + "\aShipWeb:\x11\x8a\xa8\x19\r\n" + - "\vship-web/v1\"\x1c\n" + + "\vship-web/v1\x1a\x1c\n" + "\aShipApi:\x11\x8a\xa8\x19\r\n" + - "\vship-api/v1\"B\n" + + "\vship-api/v1\x1aB\n" + "\bAnnounce\x12#\n" + "\rchangelog_url\x18\x01 \x01(\tR\fchangelogUrl:\x11\x8a\xa8\x19\r\n" + - "\vannounce/v1\"0\n" + - "\x11ReleaseTrainInput\x12\x1b\n" + - "\timage_tag\x18\x01 \x01(\tR\bimageTag\"\x9d\x01\n" + - "\fReleaseTrain:\x8c\x01\x92\xa8\x19\x87\x01\n" + - "\rrelease-train\x12\x1d.release.v1.ReleaseTrainInput\"\x17.release.v1.PlanRelease\"\x13.release.v1.ShipWeb\"\x13.release.v1.ShipApi\"\x14.release.v1.AnnounceB Date: Mon, 14 Sep 2026 23:56:21 -0300 Subject: [PATCH 3/3] buf format: drop the blank line the steps field left --- proto/durable/v1/options.proto | 1 - 1 file changed, 1 deletion(-) diff --git a/proto/durable/v1/options.proto b/proto/durable/v1/options.proto index 3f7a049..df4e3fb 100644 --- a/proto/durable/v1/options.proto +++ b/proto/durable/v1/options.proto @@ -39,7 +39,6 @@ message PipelineOptions { // Fully-qualified Output message type. string output = 3; - // Optional default concurrency class for all of this pipeline's steps; // a step's own concurrency_class overrides it. string concurrency_class = 6;