From 61892f404cfc07f8227020d6c133f250891b2ed6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Gra=C3=B1a?= Date: Thu, 17 Sep 2026 14:04:07 -0300 Subject: [PATCH 1/3] Per-pipeline middleware, composed once at Bind MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A pipeline declares its own middleware chain: pipelinedef.Config.Middleware, exposed by the generated constructor as NewXxx(h, mw...). It wraps only that pipeline's operations, forward and unwind alike, inside the engine chain — engine middleware outermost, then the pipeline's own in order, then the handler. A rule that belongs to one pipeline, such as a store's not-found ending any forward operation of that pipeline permanently, no longer reaches every pipeline bound to the engine. Both chains are composed once at Bind: boundStep carries the composed forward and unwind operations, so the per-attempt wrap and the unwind closure are gone. Bind rejects a nil middleware entry. Invariant 119. --- README.md | 4 +- cmd/protoc-gen-durable/internal/gen/gen.go | 6 +- doc.go | 5 +- docs/tour.md | 4 +- engine/bind.go | 47 ++++- engine/bind_test.go | 9 +- engine/engine.go | 24 ++- engine/middleware.go | 18 +- engine/middleware_test.go | 175 ++++++++++++++++++ .../machinespb/machines_durable.pb.go | 22 ++- examples/machines/main_test.go | 34 ++++ .../legacypb/legacy_durable.pb.go | 10 +- .../releasepb/release_durable.pb.go | 20 +- .../snapshotspb/snapshots_durable.pb.go | 10 +- .../orderspb/orders_durable.pb.go | 10 +- middleware.go | 6 +- pipelinedef/pipelinedef.go | 16 +- spec/02-authoring.md | 25 ++- spec/04-engine.md | 23 ++- spec/05-codegen.md | 4 +- spec/http-analogy.md | 12 +- spec/invariants.md | 2 + 22 files changed, 395 insertions(+), 91 deletions(-) diff --git a/README.md b/README.md index d492433..4faa966 100644 --- a/README.md +++ b/README.md @@ -160,7 +160,9 @@ packages the OpenTelemetry integration: a span per attempt linked (not parented) to the trace that scheduled the Run, metrics with durable-scale histogram buckets, `trace_id`/`span_id` log correlation, and an opt-in W3C Baggage relay. Everything is declared once, at engine -construction: +construction; a concern that belongs to one pipeline is declared once +on that pipeline's constructor instead, as +`machinespb.NewProvisionMachine(h, notFoundIsPermanent)`: ```go obs, _ := durableotel.NewObserver() diff --git a/cmd/protoc-gen-durable/internal/gen/gen.go b/cmd/protoc-gen-durable/internal/gen/gen.go index 63311fd..3485dba 100644 --- a/cmd/protoc-gen-durable/internal/gen/gen.go +++ b/cmd/protoc-gen-durable/internal/gen/gen.go @@ -478,8 +478,9 @@ func emitDefinition(g *protogen.GeneratedFile, pl *pipelineDecl) { g.P() g.P("// New", name, " assembles the ", strconv(pl.opts.GetId()), " pipeline definition") - g.P("// from its handlers.") - g.P("func New", name, "(h ", pl.handlersName(), ") *", name, "Definition {") + g.P("// from its handlers. mw wraps this pipeline's operations alone, forward and") + g.P("// unwind alike, inside any engine-level middleware; the first is outermost.") + g.P("func New", name, "(h ", pl.handlersName(), ", mw ...", g.QualifiedGoIdent(durablePkg.Ident("Middleware")), ") *", name, "Definition {") g.P("return &", name, "Definition{def: ", g.QualifiedGoIdent(defPkg.Ident("New")), "(", g.QualifiedGoIdent(defPkg.Ident("Config")), "{") g.P("ID: ", strconv(pl.opts.GetId()), ",") if ms := pl.opts.GetMutexes(); len(ms) > 0 { @@ -495,6 +496,7 @@ func emitDefinition(g *protogen.GeneratedFile, pl *pipelineDecl) { if rc := pl.opts.GetRunClass(); rc != "" { g.P("RunClass: ", strconv(rc), ",") } + g.P("Middleware: mw,") if pl.input != nil { g.P("NewInput: func() ", protoMsg, " { return &", g.QualifiedGoIdent(pl.input.GoIdent), "{} },") } diff --git a/doc.go b/doc.go index ab62a61..aa47752 100644 --- a/doc.go +++ b/doc.go @@ -44,8 +44,9 @@ // this package. // // Middleware. Handler and Middleware are the net/http-shaped operation -// layer every attempt passes through (installed with -// engine.WithMiddleware). AwaitRequest, AwaitTimeout, FailureInfo, +// layer every attempt passes through (installed engine-wide with +// engine.WithMiddleware, or on one pipeline through its generated +// constructor). AwaitRequest, AwaitTimeout, FailureInfo, // FailureCause, and FailureReason classify a handler's return the way // the engine will; PreemptedError and ErrEngineStopping name why an // attempt ctx died, for middleware that labels spans. diff --git a/docs/tour.md b/docs/tour.md index 3e16838..3514aa2 100644 --- a/docs/tour.md +++ b/docs/tour.md @@ -593,7 +593,9 @@ lifecycle events (`observe.Observer`), and snapshots occupancy [contrib/durableotel](../contrib/durableotel/) module packages the OpenTelemetry integration (per-attempt spans linked to the scheduling trace, metrics, log correlation, W3C Baggage relay), declared once at -engine construction; [examples/tracing-otel](../examples/tracing-otel/) +engine construction — or, for a concern that belongs to one pipeline, +on that pipeline's constructor (`deploypb.NewDeployService(h, mw...)`), +composed inside the engine's chain; [examples/tracing-otel](../examples/tracing-otel/) demonstrates it end to end. ## Where next diff --git a/engine/bind.go b/engine/bind.go index 82f621e..c8b88ef 100644 --- a/engine/bind.go +++ b/engine/bind.go @@ -1,19 +1,31 @@ package engine import ( + "context" "errors" "fmt" "github.com/dangra/durable" "github.com/dangra/durable/internal/ledger" "github.com/dangra/durable/pipelinedef" + "google.golang.org/protobuf/proto" ) +// boundStep is a step as the engine invokes it: its description and its +// operations composed once at Bind — the engine's middleware around the +// pipeline's around the adapter. +type boundStep struct { + *pipelinedef.Step + forward durable.Handler + unwind durable.Handler // nil unless Step.Unwind +} + // boundDef is a validated definition as the engine executes it: the -// config, its steps by id, and the ledger topology derived from them. +// config, its steps by id with their composed operations, and the ledger +// topology derived from them. type boundDef struct { cfg pipelinedef.Config - steps map[durable.StepID]*pipelinedef.Step + steps map[durable.StepID]*boundStep topo ledger.Topology // excludes is the pipeline's exclusion set as this deployment sees @@ -27,20 +39,23 @@ type boundDef struct { func (d *boundDef) ID() durable.PipelineID { return d.cfg.ID } -func (d *boundDef) step(id durable.StepID) *pipelinedef.Step { return d.steps[id] } +func (d *boundDef) step(id durable.StepID) *boundStep { return d.steps[id] } // Bind validates def and registers it with the Engine, returning the // Pipeline handle Runs are scheduled through. It is allowed only before // Start (ErrStarted afterwards). Bind is the single validator of a // definition: an empty or malformed identifier, a pipeline with no steps -// or a duplicated step, a step with no Run adapter, or an unwind -// declaration that disagrees with its adapter is reported here as an -// error — for generated code that indicates a code-generation bug. +// or a duplicated step, a step with no Run adapter, an unwind +// declaration that disagrees with its adapter, or a nil middleware entry +// is reported here as an error — for generated code that indicates a +// code-generation bug. Bind composes the pipeline's operations: the +// Engine's middleware around the pipeline's own around each adapter, +// once, here. func (e *Engine) Bind(def *pipelinedef.Definition) (*Pipeline, error) { if def == nil { return nil, errors.New("durable: Bind of a nil definition") } - bd, err := bindDefinition(def.Config()) + bd, err := bindDefinition(def.Config(), e.middleware) if err != nil { return nil, err } @@ -50,7 +65,7 @@ func (e *Engine) Bind(def *pipelinedef.Definition) (*Pipeline, error) { return &Pipeline{engine: e, def: bd}, nil } -func bindDefinition(cfg pipelinedef.Config) (*boundDef, error) { +func bindDefinition(cfg pipelinedef.Config, engineMW []durable.Middleware) (*boundDef, error) { if cfg.ID == "" { return nil, errors.New("durable: definition has empty PipelineID") } @@ -62,12 +77,17 @@ func bindDefinition(cfg pipelinedef.Config) (*boundDef, error) { return nil, fmt.Errorf("durable: pipeline %q: mutex %q must be a non-empty NUL-free valid UTF-8 name", cfg.ID, m) } } + for i, mw := range cfg.Middleware { + if mw == nil { + return nil, fmt.Errorf("durable: pipeline %q: middleware %d is nil", cfg.ID, i) + } + } if len(cfg.Steps) == 0 { return nil, fmt.Errorf("durable: pipeline %q has no steps", cfg.ID) } d := &boundDef{ cfg: cfg, - steps: make(map[durable.StepID]*pipelinedef.Step, len(cfg.Steps)), + steps: make(map[durable.StepID]*boundStep, len(cfg.Steps)), } for i := range cfg.Steps { sc := &cfg.Steps[i] @@ -86,7 +106,14 @@ func bindDefinition(cfg pipelinedef.Config) (*boundDef, error) { if sc.Unwind != (sc.UnwindFunc != nil) { return nil, fmt.Errorf("durable: pipeline %q: step %q unwind declaration and adapter disagree", cfg.ID, sc.ID) } - d.steps[sc.ID] = sc + bs := &boundStep{Step: sc, forward: wrap(durable.Handler(sc.Run), engineMW, cfg.Middleware)} + if sc.Unwind { + uf := sc.UnwindFunc + bs.unwind = wrap(func(ctx context.Context, inv durable.Invocation) (proto.Message, error) { + return nil, uf(ctx, inv) + }, engineMW, cfg.Middleware) + } + d.steps[sc.ID] = bs d.topo = append(d.topo, ledger.Step{ ID: string(sc.ID), Unwind: sc.Unwind, diff --git a/engine/bind_test.go b/engine/bind_test.go index a36de82..4dfc33a 100644 --- a/engine/bind_test.go +++ b/engine/bind_test.go @@ -35,6 +35,7 @@ func TestBindValidatesTheDefinition(t *testing.T) { {"no run adapter", pipelinedef.Config{ID: "p", Steps: []pipelinedef.Step{{ID: "s/v1"}}}, "no Run adapter"}, {"unwind without adapter", pipelinedef.Config{ID: "p", Steps: []pipelinedef.Step{{ID: "s/v1", Run: run, Unwind: true}}}, "disagree"}, {"adapter without unwind", pipelinedef.Config{ID: "p", Steps: []pipelinedef.Step{{ID: "s/v1", Run: run, UnwindFunc: unwind}}}, "disagree"}, + {"nil middleware", pipelinedef.Config{ID: "p", Middleware: []durable.Middleware{nil}, Steps: []pipelinedef.Step{ok}}, "middleware 0 is nil"}, } for _, tc := range cases { t.Run(tc.name, func(t *testing.T) { @@ -63,12 +64,18 @@ func TestDefinitionNormalizesAndCopies(t *testing.T) { {ID: "a/v1"}, {ID: "b/v1", ConcurrencyClass: "own"}, } - def := pipelinedef.New(pipelinedef.Config{ID: "p", ConcurrencyClass: "default", Steps: steps}) + passthrough := func(next durable.Handler) durable.Handler { return next } + mws := []durable.Middleware{passthrough} + def := pipelinedef.New(pipelinedef.Config{ID: "p", ConcurrencyClass: "default", Steps: steps, Middleware: mws}) steps[0].ID = "mutated" + mws[0] = nil got := def.Config().Steps if got[0].ID != "a/v1" { t.Fatal("New must copy the Steps slice") } + if got := def.Config().Middleware; len(got) != 1 || got[0] == nil { + t.Fatal("New must copy the Middleware slice") + } if got[0].ConcurrencyClass != "default" || got[1].ConcurrencyClass != "own" { t.Fatalf("class defaulting: %+v", got) } diff --git a/engine/engine.go b/engine/engine.go index c93e419..113eda9 100644 --- a/engine/engine.go +++ b/engine/engine.go @@ -6,7 +6,6 @@ import ( "fmt" "github.com/dangra/durable" "github.com/dangra/durable/observe" - "github.com/dangra/durable/pipelinedef" "github.com/dangra/durable/store/driver" "log/slog" "math" @@ -1388,7 +1387,7 @@ func (e *Engine) runForward(rec *driver.RunRecord, def *boundDef, stepID durable inv := e.invocation(rec, def, stepID, sr.Forward.Attempts, durable.PhaseForward) inv.awaited = rec.Awaited.Clone() opStart := e.clock.Now() - state, panicked, interrupted, preempted, err := e.invokeForward(sc, inv) + state, panicked, interrupted, preempted, err := e.invokeForward(sc.forward, inv) if v := inv.takeViolation(); v != nil { e.markInvalid(rec, stepID, v.Error()) @@ -1507,7 +1506,7 @@ func (e *Engine) runUnwind(rec *driver.RunRecord, def *boundDef, stepID durable. inv.failure = &f } opStart := e.clock.Now() - panicked, interrupted, err := e.invokeUnwind(sc, inv) + panicked, interrupted, err := e.invokeUnwind(sc.unwind, inv) if v := inv.takeViolation(); v != nil { e.markInvalid(rec, stepID, v.Error()) @@ -1687,11 +1686,11 @@ func committedStates(rec *driver.RunRecord) map[durable.StepID][]byte { return states } -// invokeForward runs the forward handler under a fresh attempt context. -// interrupted reports that shutdown killed that context before the -// handler returned; the resolution treats an ordinary error from such an -// attempt as an interruption, not a failure. -func (e *Engine) invokeForward(sc *pipelinedef.Step, inv *attemptInvocation) (state proto.Message, panicked, interrupted bool, preempted *durable.PreemptedError, err error) { +// invokeForward runs a step's composed forward operation under a fresh +// attempt context. interrupted reports that shutdown killed that context +// before the handler returned; the resolution treats an ordinary error +// from such an attempt as an interruption, not a failure. +func (e *Engine) invokeForward(h durable.Handler, inv *attemptInvocation) (state proto.Message, panicked, interrupted bool, preempted *durable.PreemptedError, err error) { defer func() { if p := recover(); p != nil { panicked = true @@ -1703,11 +1702,13 @@ func (e *Engine) invokeForward(sc *pipelinedef.Step, inv *attemptInvocation) (st }() ctx, done := e.attemptContext(inv.runID, durable.PhaseForward) defer done() - state, err = e.wrap(durable.Handler(sc.Run))(ctx, inv) + state, err = h(ctx, inv) return state, false, stoppedBy(ctx), preemptedBy(ctx), err } -func (e *Engine) invokeUnwind(sc *pipelinedef.Step, inv *attemptInvocation) (panicked, interrupted bool, err error) { +// invokeUnwind runs a step's composed unwind operation under a fresh +// attempt context. +func (e *Engine) invokeUnwind(h durable.Handler, inv *attemptInvocation) (panicked, interrupted bool, err error) { defer func() { if p := recover(); p != nil { panicked = true @@ -1717,9 +1718,6 @@ func (e *Engine) invokeUnwind(sc *pipelinedef.Step, inv *attemptInvocation) (pan "panic", p, "stack", string(debug.Stack())) } }() - h := e.wrap(func(ctx context.Context, in durable.Invocation) (proto.Message, error) { - return nil, sc.UnwindFunc(ctx, in) - }) ctx, done := e.attemptContext(inv.runID, durable.PhaseUnwind) defer done() _, err = h(ctx, inv) diff --git a/engine/middleware.go b/engine/middleware.go index 4ec7da1..cfb7439 100644 --- a/engine/middleware.go +++ b/engine/middleware.go @@ -9,18 +9,24 @@ import ( // WithMiddleware installs middleware around every operation the Engine // executes, forward and unwind alike; use Invocation.Phase to distinguish // them. The first middleware is the outermost, following the net/http -// convention: WithMiddleware(a, b) yields a(b(handler)). +// convention: WithMiddleware(a, b) yields a(b(handler)). Pipeline-level +// middleware (pipelinedef.Config.Middleware, the generated constructor's +// variadic) composes inside this chain: engine middleware is outermost. func WithMiddleware(mw ...durable.Middleware) Option { return func(e *Engine) { e.middleware = append(e.middleware, mw...) } } -// wrap composes the engine's middleware chain around h, first middleware -// outermost. -func (e *Engine) wrap(h durable.Handler) durable.Handler { - for _, v := range slices.Backward(e.middleware) { - h = v(h) +// wrap composes middleware chains around h, outermost first across chains +// and within each: wrap(h, engine, pipeline) yields +// engine[0](…engine[n](pipeline[0](…pipeline[m](h)))). It runs once, at +// Bind; the result runs once per attempt. +func wrap(h durable.Handler, chains ...[]durable.Middleware) durable.Handler { + for _, chain := range slices.Backward(chains) { + for _, mw := range slices.Backward(chain) { + h = mw(h) + } } return h } diff --git a/engine/middleware_test.go b/engine/middleware_test.go index ba4c88e..6b12651 100644 --- a/engine/middleware_test.go +++ b/engine/middleware_test.go @@ -6,6 +6,7 @@ import ( "context" "errors" "fmt" + "slices" "sync" "sync/atomic" "testing" @@ -171,3 +172,177 @@ func TestMiddlewareContextReachesHandlers(t *testing.T) { t.Fatalf("Wait = %+v, %v; want success", res, err) } } + +// recordingMiddleware records name:in/out around every operation it +// wraps, keyed by pipeline so a chain's reach is visible. +func recordingMiddleware(mu *sync.Mutex, events map[durable.PipelineID][]string, name string) durable.Middleware { + return func(next durable.Handler) durable.Handler { + return func(ctx context.Context, inv durable.Invocation) (proto.Message, error) { + mu.Lock() + events[inv.PipelineID()] = append(events[inv.PipelineID()], name+":in") + mu.Unlock() + state, err := next(ctx, inv) + mu.Lock() + events[inv.PipelineID()] = append(events[inv.PipelineID()], name+":out") + mu.Unlock() + return state, err + } + } +} + +// TestPipelineMiddlewareScopeAndOrder pins that a pipeline's own chain +// wraps only that pipeline's operations, forward and unwind, inside the +// engine chain, first listed outermost. +func TestPipelineMiddlewareScopeAndOrder(t *testing.T) { + var mu sync.Mutex + events := map[durable.PipelineID][]string{} + mw := func(name string) durable.Middleware { return recordingMiddleware(&mu, events, name) } + + own := pipelinedef.New(pipelinedef.Config{ + ID: "own", + Middleware: []durable.Middleware{mw("p1"), mw("p2")}, + Steps: []pipelinedef.Step{ + { + ID: "a/v1", + Unwind: true, + Run: func(ctx context.Context, inv durable.Invocation) (proto.Message, error) { return nil, nil }, + UnwindFunc: func(ctx context.Context, inv durable.Invocation) error { return nil }, + }, + stateless("b/v1", func(ctx context.Context, inv durable.Invocation) error { + if inv.Attempt() == 1 { + return errors.New("transient") + } + return durable.Fail(errors.New("permanent")) + }), + }, + }) + other := pipelinedef.New(pipelinedef.Config{ + ID: "other", + Steps: []pipelinedef.Step{stateless("s/v1", func(ctx context.Context, inv durable.Invocation) error { return nil })}, + }) + e := engine.New(mem.New(), fastRetry, engine.WithLogger(discardTestLogger()), + engine.WithMiddleware(mw("engine"))) + pOwn, err := e.Bind(own) + if err != nil { + t.Fatalf("Bind own: %v", err) + } + pOther, err := e.Bind(other) + if err != nil { + t.Fatalf("Bind other: %v", err) + } + if err := e.Start(context.Background()); err != nil { + t.Fatalf("Start: %v", err) + } + defer e.Stop(context.Background()) + + rOwn, _, _ := pOwn.Schedule(context.Background(), "r", nil) + rOther, _, _ := pOther.Schedule(context.Background(), "r", nil) + if res, err := rOwn.Wait(context.Background()); err != nil || !res.Failed() { + t.Fatalf("own Wait = %+v, %v; want failure", res, err) + } + if res, err := rOther.Wait(context.Background()); err != nil || !res.Succeeded() { + t.Fatalf("other Wait = %+v, %v; want success", res, err) + } + + mu.Lock() + defer mu.Unlock() + // Four operations on own — a.Run, b.Run twice, a.Unwind — each the + // onion engine(p1(p2(handler))). + onion := []string{"engine:in", "p1:in", "p2:in", "p2:out", "p1:out", "engine:out"} + var want []string + for i := 0; i < 4; i++ { + want = append(want, onion...) + } + if got := events["own"]; !slices.Equal(got, want) { + t.Fatalf("own events = %v, want %v", got, want) + } + if got := events["other"]; !slices.Equal(got, []string{"engine:in", "engine:out"}) { + t.Fatalf("other events = %v; want the engine chain alone", got) + } +} + +// TestPipelineMiddlewareEscalatesForItsPipelineOnly is the flyd case: a +// store's not-found error is a permanent failure of one pipeline's +// forward operations and an ordinary error of everyone else's. +func TestPipelineMiddlewareEscalatesForItsPipelineOnly(t *testing.T) { + errGone := errors.New("machine not found") + gone := func(next durable.Handler) durable.Handler { + return func(ctx context.Context, inv durable.Invocation) (proto.Message, error) { + out, err := next(ctx, inv) + if err != nil && inv.Phase() == durable.PhaseForward && errors.Is(err, errGone) { + return nil, durable.Fail(err, durable.WithReason("machine-gone")) + } + return out, err + } + } + var strictAttempts, lenientAttempts atomic.Int32 + strict := pipelinedef.New(pipelinedef.Config{ + ID: "strict", + Middleware: []durable.Middleware{gone}, + Steps: []pipelinedef.Step{stateless("strict/v1", func(ctx context.Context, inv durable.Invocation) error { + strictAttempts.Add(1) + return errGone + })}, + }) + lenient := pipelinedef.New(pipelinedef.Config{ + ID: "lenient", + Steps: []pipelinedef.Step{stateless("lenient/v1", func(ctx context.Context, inv durable.Invocation) error { + if lenientAttempts.Add(1) == 1 { + return errGone + } + return nil + })}, + }) + _, pipes := startEngine(t, mem.New(), strict, lenient) + rs, _, _ := pipes[0].Schedule(context.Background(), "r", nil) + rl, _, _ := pipes[1].Schedule(context.Background(), "r", nil) + res, err := rs.Wait(context.Background()) + if err != nil || !res.Failed() || res.Failure.Reason != "machine-gone" { + t.Fatalf("strict Wait = %+v, %v; want a machine-gone failure", res, err) + } + if n := strictAttempts.Load(); n != 1 { + t.Errorf("strict attempts = %d, want 1", n) + } + if res, err := rl.Wait(context.Background()); err != nil || !res.Succeeded() { + t.Fatalf("lenient Wait = %+v, %v; want success after a retry", res, err) + } + if n := lenientAttempts.Load(); n != 2 { + t.Errorf("lenient attempts = %d, want 2", n) + } +} + +// TestMiddlewareComposedOnceAtBind pins that the wrapper factories run +// once per step and phase, however many attempts the operations take. +func TestMiddlewareComposedOnceAtBind(t *testing.T) { + var factories atomic.Int32 + counting := func(next durable.Handler) durable.Handler { + factories.Add(1) + return next + } + def := pipelinedef.New(pipelinedef.Config{ + ID: "once", + Middleware: []durable.Middleware{counting}, + Steps: []pipelinedef.Step{ + { + ID: "a/v1", + Unwind: true, + Run: func(ctx context.Context, inv durable.Invocation) (proto.Message, error) { return nil, nil }, + UnwindFunc: func(ctx context.Context, inv durable.Invocation) error { return nil }, + }, + stateless("b/v1", func(ctx context.Context, inv durable.Invocation) error { + if inv.Attempt() < 3 { + return errors.New("transient") + } + return durable.Fail(errors.New("permanent")) + }), + }, + }) + _, pipes := startEngine(t, mem.New(), def) + run, _, _ := pipes[0].Schedule(context.Background(), "r", nil) + if res, err := run.Wait(context.Background()); err != nil || !res.Failed() { + t.Fatalf("Wait = %+v, %v; want failure", res, err) + } + if n := factories.Load(); n != 3 { + t.Fatalf("factory ran %d times, want 3: a forward, b forward, a unwind", n) + } +} diff --git a/examples/machines/machinespb/machines_durable.pb.go b/examples/machines/machinespb/machines_durable.pb.go index f76b476..e5418e8 100644 --- a/examples/machines/machinespb/machines_durable.pb.go +++ b/examples/machines/machinespb/machines_durable.pb.go @@ -112,12 +112,14 @@ type ProvisionMachineDefinition struct { } // NewProvisionMachine assembles the "provision-machine" pipeline definition -// from its handlers. -func NewProvisionMachine(h ProvisionMachineHandlers) *ProvisionMachineDefinition { +// from its handlers. mw wraps this pipeline's operations alone, forward and +// unwind alike, inside any engine-level middleware; the first is outermost. +func NewProvisionMachine(h ProvisionMachineHandlers, mw ...durable.Middleware) *ProvisionMachineDefinition { return &ProvisionMachineDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "provision-machine", - Mutexes: []string{"machine-lifecycle"}, - NewInput: func() proto.Message { return &ProvisionMachineInput{} }, + ID: "provision-machine", + Mutexes: []string{"machine-lifecycle"}, + Middleware: mw, + NewInput: func() proto.Message { return &ProvisionMachineInput{} }, Reduce: func(view durable.ReduceView) proto.Message { return ReduceProvisionMachineOutput(h, view) }, @@ -338,11 +340,13 @@ type DecommissionMachineDefinition struct { } // NewDecommissionMachine assembles the "decommission-machine" pipeline definition -// from its handlers. -func NewDecommissionMachine(h DecommissionMachineHandlers) *DecommissionMachineDefinition { +// from its handlers. mw wraps this pipeline's operations alone, forward and +// unwind alike, inside any engine-level middleware; the first is outermost. +func NewDecommissionMachine(h DecommissionMachineHandlers, mw ...durable.Middleware) *DecommissionMachineDefinition { return &DecommissionMachineDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "decommission-machine", - Mutexes: []string{"machine-lifecycle"}, + ID: "decommission-machine", + Mutexes: []string{"machine-lifecycle"}, + Middleware: mw, Steps: []pipelinedef.Step{ { ID: "release-machine/v1", diff --git a/examples/machines/main_test.go b/examples/machines/main_test.go index 956e701..03807b4 100644 --- a/examples/machines/main_test.go +++ b/examples/machines/main_test.go @@ -3,6 +3,7 @@ package main import ( "context" "errors" + "sync/atomic" "testing" "time" @@ -10,6 +11,7 @@ import ( "github.com/dangra/durable/engine" "github.com/dangra/durable/examples/machines/machinespb" "github.com/dangra/durable/store/mem" + "google.golang.org/protobuf/proto" ) func startProvision(t *testing.T, c *cloud) *machinespb.ProvisionMachinePipeline { @@ -205,3 +207,35 @@ func TestLifecycleMutex(t *testing.T) { t.Fatalf("released = %v, want machine-9", c.released) } } + +// TestPipelineMiddleware passes a middleware through the generated +// constructor: it wraps this pipeline's operations alone. +func TestPipelineMiddleware(t *testing.T) { + var wrapped atomic.Int32 + counting := func(next durable.Handler) durable.Handler { + return func(ctx context.Context, inv durable.Invocation) (proto.Message, error) { + wrapped.Add(1) + return next(ctx, inv) + } + } + c := newCloud() + eng := engine.New(mem.New()) + provision, err := machinespb.NewProvisionMachine(&handlers{cloud: c}, counting).Bind(eng) + if err != nil { + t.Fatalf("Bind: %v", err) + } + if err := eng.Start(context.Background()); err != nil { + t.Fatalf("Start: %v", err) + } + defer eng.Stop(context.Background()) + run, _, err := provision.Schedule(context.Background(), "machine-mw", &machinespb.ProvisionMachineInput{Region: "ams", MemoryMb: 1024, Cpus: 1}) + if err != nil { + t.Fatalf("Schedule: %v", err) + } + if res, err := run.Wait(context.Background()); err != nil || !res.Succeeded() { + t.Fatalf("Wait = %+v, %v; want success", res, err) + } + if n := wrapped.Load(); n != 4 { + t.Fatalf("middleware wrapped %d operations, want the pipeline's 4 steps", n) + } +} diff --git a/examples/release-train/legacypb/legacy_durable.pb.go b/examples/release-train/legacypb/legacy_durable.pb.go index 5cdacda..2ef0b28 100644 --- a/examples/release-train/legacypb/legacy_durable.pb.go +++ b/examples/release-train/legacypb/legacy_durable.pb.go @@ -108,11 +108,13 @@ type DeployServiceDefinition struct { } // NewDeployService assembles the "deploy-service" pipeline definition -// from its handlers. -func NewDeployService(h DeployServiceHandlers) *DeployServiceDefinition { +// from its handlers. mw wraps this pipeline's operations alone, forward and +// unwind alike, inside any engine-level middleware; the first is outermost. +func NewDeployService(h DeployServiceHandlers, mw ...durable.Middleware) *DeployServiceDefinition { return &DeployServiceDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "deploy-service", - NewInput: func() proto.Message { return &DeployServiceInput{} }, + ID: "deploy-service", + Middleware: mw, + NewInput: func() proto.Message { return &DeployServiceInput{} }, Reduce: func(view durable.ReduceView) proto.Message { return ReduceDeployServiceOutput(h, view) }, diff --git a/examples/release-train/releasepb/release_durable.pb.go b/examples/release-train/releasepb/release_durable.pb.go index 2e2e2e3..5156091 100644 --- a/examples/release-train/releasepb/release_durable.pb.go +++ b/examples/release-train/releasepb/release_durable.pb.go @@ -114,11 +114,13 @@ type DeployServiceDefinition struct { } // NewDeployService assembles the "deploy-service" pipeline definition -// from its handlers. -func NewDeployService(h DeployServiceHandlers) *DeployServiceDefinition { +// from its handlers. mw wraps this pipeline's operations alone, forward and +// unwind alike, inside any engine-level middleware; the first is outermost. +func NewDeployService(h DeployServiceHandlers, mw ...durable.Middleware) *DeployServiceDefinition { return &DeployServiceDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "deploy-service", - NewInput: func() proto.Message { return &DeployServiceInput{} }, + ID: "deploy-service", + Middleware: mw, + NewInput: func() proto.Message { return &DeployServiceInput{} }, Reduce: func(view durable.ReduceView) proto.Message { return ReduceDeployServiceOutput(h, view) }, @@ -370,11 +372,13 @@ type ReleaseTrainDefinition struct { } // NewReleaseTrain assembles the "release-train" pipeline definition -// from its handlers. -func NewReleaseTrain(h ReleaseTrainHandlers) *ReleaseTrainDefinition { +// from its handlers. mw wraps this pipeline's operations alone, forward and +// unwind alike, inside any engine-level middleware; the first is outermost. +func NewReleaseTrain(h ReleaseTrainHandlers, mw ...durable.Middleware) *ReleaseTrainDefinition { return &ReleaseTrainDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "release-train", - NewInput: func() proto.Message { return &ReleaseTrainInput{} }, + ID: "release-train", + Middleware: mw, + NewInput: func() proto.Message { return &ReleaseTrainInput{} }, Steps: []pipelinedef.Step{ { ID: "plan/v1", diff --git a/examples/snapshots/snapshotspb/snapshots_durable.pb.go b/examples/snapshots/snapshotspb/snapshots_durable.pb.go index 8893f81..90f60c3 100644 --- a/examples/snapshots/snapshotspb/snapshots_durable.pb.go +++ b/examples/snapshots/snapshotspb/snapshots_durable.pb.go @@ -128,11 +128,13 @@ type CreateSnapshotDefinition struct { } // NewCreateSnapshot assembles the "create-snapshot" pipeline definition -// from its handlers. -func NewCreateSnapshot(h CreateSnapshotHandlers) *CreateSnapshotDefinition { +// from its handlers. mw wraps this pipeline's operations alone, forward and +// unwind alike, inside any engine-level middleware; the first is outermost. +func NewCreateSnapshot(h CreateSnapshotHandlers, mw ...durable.Middleware) *CreateSnapshotDefinition { return &CreateSnapshotDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "create-snapshot", - NewInput: func() proto.Message { return &CreateSnapshotInput{} }, + ID: "create-snapshot", + Middleware: mw, + NewInput: func() proto.Message { return &CreateSnapshotInput{} }, Reduce: func(view durable.ReduceView) proto.Message { return ReduceCreateSnapshotOutput(h, view) }, diff --git a/examples/tracing-otel/orderspb/orders_durable.pb.go b/examples/tracing-otel/orderspb/orders_durable.pb.go index 9cec240..4f03150 100644 --- a/examples/tracing-otel/orderspb/orders_durable.pb.go +++ b/examples/tracing-otel/orderspb/orders_durable.pb.go @@ -108,11 +108,13 @@ type FulfillOrderDefinition struct { } // NewFulfillOrder assembles the "fulfill-order" pipeline definition -// from its handlers. -func NewFulfillOrder(h FulfillOrderHandlers) *FulfillOrderDefinition { +// from its handlers. mw wraps this pipeline's operations alone, forward and +// unwind alike, inside any engine-level middleware; the first is outermost. +func NewFulfillOrder(h FulfillOrderHandlers, mw ...durable.Middleware) *FulfillOrderDefinition { return &FulfillOrderDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "fulfill-order", - NewInput: func() proto.Message { return &FulfillOrderInput{} }, + ID: "fulfill-order", + Middleware: mw, + NewInput: func() proto.Message { return &FulfillOrderInput{} }, Reduce: func(view durable.ReduceView) proto.Message { return ReduceFulfillOrderOutput(h, view) }, diff --git a/middleware.go b/middleware.go index 4d60c78..86ec5a4 100644 --- a/middleware.go +++ b/middleware.go @@ -18,7 +18,11 @@ type Handler func(ctx context.Context, inv Invocation) (proto.Message, error) // Middleware wraps a Handler — durable's analog of // func(http.Handler) http.Handler. Use it for cross-cutting concerns such -// as logging, metrics, tracing spans, or per-operation timeouts. +// as logging, metrics, tracing spans, or per-operation timeouts. It is +// installed engine-wide with engine.WithMiddleware, or on one pipeline +// through pipelinedef.Config.Middleware (the generated constructor's +// variadic); the engine chain is outermost, and both are composed once +// at Engine.Bind. // // Middleware runs once per attempt, inside the durable attempt // reservation: it inherits the operation's at-least-once semantics and diff --git a/pipelinedef/pipelinedef.go b/pipelinedef/pipelinedef.go index 8e41e9a..32997aa 100644 --- a/pipelinedef/pipelinedef.go +++ b/pipelinedef/pipelinedef.go @@ -63,6 +63,13 @@ type Config struct { // share its capacity. RunClass string + // Middleware wraps this pipeline's operations alone, forward and + // unwind alike (Invocation.Phase distinguishes them), inside any + // engine-level middleware: engine middleware is outermost, then these + // in order, first listed outermost, then the handler. Composition is + // fixed at Engine.Bind, which rejects a nil entry. + Middleware []durable.Middleware + // NewInput constructs an empty Input message; nil for an Input-less // pipeline. NewInput func() proto.Message @@ -84,12 +91,13 @@ type Definition struct { cfg Config } -// New wraps cfg. It copies the Steps and Mutexes slices so later mutation -// of cfg does not reach the Definition, and applies the pipeline-level -// concurrency class to steps that declare none. +// New wraps cfg. It copies the Steps, Mutexes, and Middleware slices so +// later mutation of cfg does not reach the Definition, and applies the +// pipeline-level concurrency class to steps that declare none. func New(cfg Config) *Definition { cfg.Steps = append([]Step(nil), cfg.Steps...) cfg.Mutexes = append([]string(nil), cfg.Mutexes...) + cfg.Middleware = append([]durable.Middleware(nil), cfg.Middleware...) for i := range cfg.Steps { if cfg.Steps[i].ConcurrencyClass == "" { cfg.Steps[i].ConcurrencyClass = cfg.ConcurrencyClass @@ -103,7 +111,7 @@ func (d *Definition) ID() durable.PipelineID { return d.cfg.ID } // Config returns the description the Definition was built from, with the // concurrency-class default applied. The engine reads it at Bind; the -// Steps slice is shared and must not be modified. +// Steps and Middleware slices are shared and must not be modified. func (d *Definition) Config() Config { return d.cfg } // StepRef constructs the reference to a stateless Step. Generated code diff --git a/spec/02-authoring.md b/spec/02-authoring.md index 9b8a1ac..67f6d22 100644 --- a/spec/02-authoring.md +++ b/spec/02-authoring.md @@ -262,9 +262,13 @@ type ProvisionMachineHandlers interface { ReduceOutput(*ProvisionMachine) *ProvisionMachineOutput } -func NewProvisionMachine(h ProvisionMachineHandlers) *ProvisionMachineDefinition +func NewProvisionMachine(h ProvisionMachineHandlers, mw ...durable.Middleware) *ProvisionMachineDefinition ``` +The variadic installs pipeline-level middleware, wrapping this +pipeline's operations alone inside any engine-level chain (see +[04-engine](04-engine.md#middleware)). + One type implements the whole pipeline, so its dependencies are declared once; a Step the implementor lacks, a Step added to the proto included, is a compile error naming the method. @@ -584,27 +588,20 @@ Purity makes this safe. Example: ```go -definition := machines.NewProvisionMachine( - &validate{}, - &selectHost{}, - &reserveCapacity{}, - &createMachine{}, - reduceProvisionMachine, -) +definition := machines.NewProvisionMachine(&handlers{cloud: c}) ``` -The positional constructor is type-safe because each position expects a distinct generated handler interface. - -Accidentally swapping two unrelated Step handlers fails to compile. +The constructor takes one implementation of the pipeline's handler +interface: a missing or mis-typed Step method, a Step added to the proto +included, fails to compile. Middleware for this pipeline alone rides the +variadic: `machines.NewProvisionMachine(&handlers{cloud: c}, notFoundIsPermanent)`. --- ## Bind ```go -provision, err := machines.NewProvisionMachine( - ..., -).Bind(eng) +provision, err := machines.NewProvisionMachine(&handlers{cloud: c}).Bind(eng) ``` returns: diff --git a/spec/04-engine.md b/spec/04-engine.md index ddfcf89..7c325b8 100644 --- a/spec/04-engine.md +++ b/spec/04-engine.md @@ -647,6 +647,26 @@ engine.WithMiddleware(logging, metrics) The first middleware is outermost. `Invocation.Phase()` distinguishes forward from unwind operations. +**Pipeline-level middleware.** A definition MAY carry its own chain, +`pipelinedef.Config.Middleware`; the generated constructor exposes it +as a variadic, `NewXxx(h, mw...)`. It wraps only that pipeline's +operations, forward and unwind alike, and composes inside the +engine-level chain: + +```text +engine[0](… engine[n](pipeline[0](… pipeline[m](handler)))) +``` + +Engine middleware is outermost; within each chain the first listed is +outermost. A concern that belongs to one pipeline — treating a store's +not-found error as a permanent failure of any forward operation of that +pipeline, say — is declared on that pipeline and never reaches another +bound to the same Engine. + +Composition is fixed at `Engine.Bind`: the wrapper functions run once +per Step and phase, and the composed operation runs once per attempt. +Bind rejects a nil middleware entry. + Middleware contract: - Middleware runs once per attempt, inside the durable attempt @@ -662,7 +682,8 @@ Middleware contract: operation unresolved. - The Reducer is pure, not an operation, and is never wrapped. -Per-Step middleware composition is outside v1.1. +Per-Step middleware composition is outside v1.1; a pipeline-level +middleware that switches on `Invocation.StepID()` covers the need. See the non-normative [net/http analogy note](http-analogy.md) for the design rationale. diff --git a/spec/05-codegen.md b/spec/05-codegen.md index f02c17e..964f20f 100644 --- a/spec/05-codegen.md +++ b/spec/05-codegen.md @@ -53,7 +53,9 @@ Published protobuf extensions MUST use globally allocated extension numbers. - one handler interface, `XxxHandlers`: a method per Step named after the Step, `Unwind` for each Step that unwinds, and `ReduceOutput` / `ReduceFailure` when the pipeline declares outputs, -- the Pipeline constructor, `NewXxx(h XxxHandlers)`, +- the Pipeline constructor, `NewXxx(h XxxHandlers, mw ...durable.Middleware)`; + the variadic is the pipeline's own middleware chain + ([04-engine](04-engine.md#middleware)), - `ReduceXxxOutput(h, view)` and `ReduceXxxFailure(h, view)`, the folds the engine reduces through and reducer tests call, - runtime methods on the Pipeline marker type, diff --git a/spec/http-analogy.md b/spec/http-analogy.md index fda0150..e555a69 100644 --- a/spec/http-analogy.md +++ b/spec/http-analogy.md @@ -10,8 +10,9 @@ analogy with Go's `net/http`; the normative contracts live in | net/http | durable | |---|---| | `http.Handler` | `durable.Handler` — the uniform type-erased operation `func(ctx, Invocation) (proto.Message, error)` | -| `http.HandlerFunc` | generated func adapters (`XFunc`, `XFuncs`) | -| middleware `func(http.Handler) http.Handler` | `durable.Middleware`, installed with `engine.WithMiddleware` | +| `http.HandlerFunc` | the generated per-pipeline handler interface's methods, erased by the constructor's adapters | +| middleware `func(http.Handler) http.Handler` | `durable.Middleware`, installed engine-wide with `engine.WithMiddleware` | +| wrapping one handler at mux registration, `mux.Handle(pattern, mw(h))` | pipeline-level middleware: `pipelinedef.Config.Middleware`, the generated `NewXxx(h, mw...)` variadic; composes inside the engine chain | | `ServeMux` and route patterns | the protobuf Pipeline topology; `StepID` plays the route pattern | | `*http.Request` | `durable.Invocation` — identity, attempt, phase, input, State lookup; an interface, so tests can fake it | | `http.ResponseWriter` | — (deliberately absent; see below) | @@ -31,9 +32,10 @@ typed router ecosystems : http.Handler ``` That erased seam is where cross-cutting behavior belongs. Middleware -wraps it uniformly — logging, metrics, tracing spans, per-operation -timeouts — without touching the typed authoring surface or the durable -execution model. +wraps it uniformly at the engine — logging, metrics, tracing spans, +per-operation timeouts — or at one pipeline's registration for a concern +that belongs to that pipeline alone, without touching the typed +authoring surface or the durable execution model. ## What deliberately does not transfer diff --git a/spec/invariants.md b/spec/invariants.md index b888a31..7e3efcd 100644 --- a/spec/invariants.md +++ b/spec/invariants.md @@ -237,3 +237,5 @@ Part of the [`durable` specification](README.md). This list is append-only; inva 117. The operation a cancellation resolves keeps the attempts it had reserved and carries the cancellation on its record; the Run's Failure names its Step. 118. An unwind attempt's context is never canceled for a cancellation request; a request on a Run in unwind is recorded and changes nothing. + +119. Pipeline-level middleware wraps only its own pipeline's operations, forward and unwind alike, and composes inside the engine-level chain — engine middleware outermost, then the pipeline's own in declaration order, then the handler; composition is fixed at Bind. From 23d2010ca8ac2abbf048caa8224c72d13ed1a740 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Gra=C3=B1a?= Date: Thu, 17 Sep 2026 14:56:19 -0300 Subject: [PATCH 2/3] Pipeline middleware is a constructor option: pipelinedef.WithMiddleware The generated constructor takes opts ...pipelinedef.Option, the knobs a proto cannot declare; WithMiddleware is the first. Wiring code imports pipelinedef for it; handler code still never does. --- README.md | 2 +- cmd/protoc-gen-durable/internal/gen/gen.go | 15 ++++--- doc.go | 7 ++-- docs/tour.md | 2 +- engine/middleware.go | 5 ++- .../machinespb/machines_durable.pb.go | 40 +++++++++++-------- examples/machines/main_test.go | 3 +- .../legacypb/legacy_durable.pb.go | 19 +++++---- .../releasepb/release_durable.pb.go | 38 ++++++++++-------- .../snapshotspb/snapshots_durable.pb.go | 19 +++++---- .../orderspb/orders_durable.pb.go | 19 +++++---- middleware.go | 5 +-- pipelinedef/pipelinedef.go | 19 ++++++++- spec/02-authoring.md | 13 +++--- spec/04-engine.md | 5 ++- spec/05-codegen.md | 4 +- spec/http-analogy.md | 2 +- 17 files changed, 130 insertions(+), 87 deletions(-) diff --git a/README.md b/README.md index 4faa966..27a777c 100644 --- a/README.md +++ b/README.md @@ -162,7 +162,7 @@ durable-scale histogram buckets, `trace_id`/`span_id` log correlation, and an opt-in W3C Baggage relay. Everything is declared once, at engine construction; a concern that belongs to one pipeline is declared once on that pipeline's constructor instead, as -`machinespb.NewProvisionMachine(h, notFoundIsPermanent)`: +`machinespb.NewProvisionMachine(h, pipelinedef.WithMiddleware(notFoundIsPermanent))`: ```go obs, _ := durableotel.NewObserver() diff --git a/cmd/protoc-gen-durable/internal/gen/gen.go b/cmd/protoc-gen-durable/internal/gen/gen.go index 3485dba..a678e9e 100644 --- a/cmd/protoc-gen-durable/internal/gen/gen.go +++ b/cmd/protoc-gen-durable/internal/gen/gen.go @@ -478,10 +478,10 @@ func emitDefinition(g *protogen.GeneratedFile, pl *pipelineDecl) { g.P() g.P("// New", name, " assembles the ", strconv(pl.opts.GetId()), " pipeline definition") - g.P("// from its handlers. mw wraps this pipeline's operations alone, forward and") - g.P("// unwind alike, inside any engine-level middleware; the first is outermost.") - g.P("func New", name, "(h ", pl.handlersName(), ", mw ...", g.QualifiedGoIdent(durablePkg.Ident("Middleware")), ") *", name, "Definition {") - g.P("return &", name, "Definition{def: ", g.QualifiedGoIdent(defPkg.Ident("New")), "(", g.QualifiedGoIdent(defPkg.Ident("Config")), "{") + g.P("// from its handlers; opts (pipelinedef.WithMiddleware) configure what the") + g.P("// proto cannot declare.") + g.P("func New", name, "(h ", pl.handlersName(), ", opts ...", g.QualifiedGoIdent(defPkg.Ident("Option")), ") *", name, "Definition {") + g.P("cfg := ", g.QualifiedGoIdent(defPkg.Ident("Config")), "{") g.P("ID: ", strconv(pl.opts.GetId()), ",") if ms := pl.opts.GetMutexes(); len(ms) > 0 { quoted := make([]string, len(ms)) @@ -496,7 +496,6 @@ func emitDefinition(g *protogen.GeneratedFile, pl *pipelineDecl) { if rc := pl.opts.GetRunClass(); rc != "" { g.P("RunClass: ", strconv(rc), ",") } - g.P("Middleware: mw,") if pl.input != nil { g.P("NewInput: func() ", protoMsg, " { return &", g.QualifiedGoIdent(pl.input.GoIdent), "{} },") } @@ -547,7 +546,11 @@ func emitDefinition(g *protogen.GeneratedFile, pl *pipelineDecl) { g.P("},") } g.P("},") - g.P("})}") + g.P("}") + g.P("for _, o := range opts {") + g.P("o(&cfg)") + g.P("}") + g.P("return &", name, "Definition{def: ", g.QualifiedGoIdent(defPkg.Ident("New")), "(cfg)}") g.P("}") g.P() diff --git a/doc.go b/doc.go index aa47752..01c3d2b 100644 --- a/doc.go +++ b/doc.go @@ -45,8 +45,8 @@ // // Middleware. Handler and Middleware are the net/http-shaped operation // layer every attempt passes through (installed engine-wide with -// engine.WithMiddleware, or on one pipeline through its generated -// constructor). AwaitRequest, AwaitTimeout, FailureInfo, +// engine.WithMiddleware, or on one pipeline with pipelinedef.WithMiddleware +// on its generated constructor). AwaitRequest, AwaitTimeout, FailureInfo, // FailureCause, and FailureReason classify a handler's return the way // the engine will; PreemptedError and ErrEngineStopping name why an // attempt ctx died, for middleware that labels spans. @@ -105,7 +105,8 @@ // the wiring side of an application. // - pipelinedef is the type-erased pipeline description generated // code builds and Engine.Bind validates; hand-rolled definitions use -// it too. Application code does not import it. +// it too. Wiring code imports it for the constructor options +// (WithMiddleware); handler code never does. // - kernel is the shared vocabulary — identities, phases, outcomes, // parks, failure records — that every other package builds on. This // package aliases all of it, so user code never imports kernel. diff --git a/docs/tour.md b/docs/tour.md index 3514aa2..21e80ef 100644 --- a/docs/tour.md +++ b/docs/tour.md @@ -594,7 +594,7 @@ lifecycle events (`observe.Observer`), and snapshots occupancy OpenTelemetry integration (per-attempt spans linked to the scheduling trace, metrics, log correlation, W3C Baggage relay), declared once at engine construction — or, for a concern that belongs to one pipeline, -on that pipeline's constructor (`deploypb.NewDeployService(h, mw...)`), +on that pipeline's constructor (`deploypb.NewDeployService(h, pipelinedef.WithMiddleware(mw))`), composed inside the engine's chain; [examples/tracing-otel](../examples/tracing-otel/) demonstrates it end to end. diff --git a/engine/middleware.go b/engine/middleware.go index cfb7439..d739fd8 100644 --- a/engine/middleware.go +++ b/engine/middleware.go @@ -10,8 +10,9 @@ import ( // executes, forward and unwind alike; use Invocation.Phase to distinguish // them. The first middleware is the outermost, following the net/http // convention: WithMiddleware(a, b) yields a(b(handler)). Pipeline-level -// middleware (pipelinedef.Config.Middleware, the generated constructor's -// variadic) composes inside this chain: engine middleware is outermost. +// middleware (pipelinedef.WithMiddleware on a generated constructor, or +// pipelinedef.Config.Middleware) composes inside this chain: engine +// middleware is outermost. func WithMiddleware(mw ...durable.Middleware) Option { return func(e *Engine) { e.middleware = append(e.middleware, mw...) diff --git a/examples/machines/machinespb/machines_durable.pb.go b/examples/machines/machinespb/machines_durable.pb.go index e5418e8..b4869f8 100644 --- a/examples/machines/machinespb/machines_durable.pb.go +++ b/examples/machines/machinespb/machines_durable.pb.go @@ -112,14 +112,13 @@ type ProvisionMachineDefinition struct { } // NewProvisionMachine assembles the "provision-machine" pipeline definition -// from its handlers. mw wraps this pipeline's operations alone, forward and -// unwind alike, inside any engine-level middleware; the first is outermost. -func NewProvisionMachine(h ProvisionMachineHandlers, mw ...durable.Middleware) *ProvisionMachineDefinition { - return &ProvisionMachineDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "provision-machine", - Mutexes: []string{"machine-lifecycle"}, - Middleware: mw, - NewInput: func() proto.Message { return &ProvisionMachineInput{} }, +// from its handlers; opts (pipelinedef.WithMiddleware) configure what the +// proto cannot declare. +func NewProvisionMachine(h ProvisionMachineHandlers, opts ...pipelinedef.Option) *ProvisionMachineDefinition { + cfg := pipelinedef.Config{ + ID: "provision-machine", + Mutexes: []string{"machine-lifecycle"}, + NewInput: func() proto.Message { return &ProvisionMachineInput{} }, Reduce: func(view durable.ReduceView) proto.Message { return ReduceProvisionMachineOutput(h, view) }, @@ -169,7 +168,11 @@ func NewProvisionMachine(h ProvisionMachineHandlers, mw ...durable.Middleware) * }, }, }, - })} + } + for _, o := range opts { + o(&cfg) + } + return &ProvisionMachineDefinition{def: pipelinedef.New(cfg)} } // Bind registers the definition with an engine. It is allowed only before @@ -340,13 +343,12 @@ type DecommissionMachineDefinition struct { } // NewDecommissionMachine assembles the "decommission-machine" pipeline definition -// from its handlers. mw wraps this pipeline's operations alone, forward and -// unwind alike, inside any engine-level middleware; the first is outermost. -func NewDecommissionMachine(h DecommissionMachineHandlers, mw ...durable.Middleware) *DecommissionMachineDefinition { - return &DecommissionMachineDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "decommission-machine", - Mutexes: []string{"machine-lifecycle"}, - Middleware: mw, +// from its handlers; opts (pipelinedef.WithMiddleware) configure what the +// proto cannot declare. +func NewDecommissionMachine(h DecommissionMachineHandlers, opts ...pipelinedef.Option) *DecommissionMachineDefinition { + cfg := pipelinedef.Config{ + ID: "decommission-machine", + Mutexes: []string{"machine-lifecycle"}, Steps: []pipelinedef.Step{ { ID: "release-machine/v1", @@ -355,7 +357,11 @@ func NewDecommissionMachine(h DecommissionMachineHandlers, mw ...durable.Middlew }, }, }, - })} + } + for _, o := range opts { + o(&cfg) + } + return &DecommissionMachineDefinition{def: pipelinedef.New(cfg)} } // Bind registers the definition with an engine. It is allowed only before diff --git a/examples/machines/main_test.go b/examples/machines/main_test.go index 03807b4..8fc388c 100644 --- a/examples/machines/main_test.go +++ b/examples/machines/main_test.go @@ -10,6 +10,7 @@ import ( "github.com/dangra/durable" "github.com/dangra/durable/engine" "github.com/dangra/durable/examples/machines/machinespb" + "github.com/dangra/durable/pipelinedef" "github.com/dangra/durable/store/mem" "google.golang.org/protobuf/proto" ) @@ -220,7 +221,7 @@ func TestPipelineMiddleware(t *testing.T) { } c := newCloud() eng := engine.New(mem.New()) - provision, err := machinespb.NewProvisionMachine(&handlers{cloud: c}, counting).Bind(eng) + provision, err := machinespb.NewProvisionMachine(&handlers{cloud: c}, pipelinedef.WithMiddleware(counting)).Bind(eng) if err != nil { t.Fatalf("Bind: %v", err) } diff --git a/examples/release-train/legacypb/legacy_durable.pb.go b/examples/release-train/legacypb/legacy_durable.pb.go index 2ef0b28..49cbfc0 100644 --- a/examples/release-train/legacypb/legacy_durable.pb.go +++ b/examples/release-train/legacypb/legacy_durable.pb.go @@ -108,13 +108,12 @@ type DeployServiceDefinition struct { } // NewDeployService assembles the "deploy-service" pipeline definition -// from its handlers. mw wraps this pipeline's operations alone, forward and -// unwind alike, inside any engine-level middleware; the first is outermost. -func NewDeployService(h DeployServiceHandlers, mw ...durable.Middleware) *DeployServiceDefinition { - return &DeployServiceDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "deploy-service", - Middleware: mw, - NewInput: func() proto.Message { return &DeployServiceInput{} }, +// from its handlers; opts (pipelinedef.WithMiddleware) configure what the +// proto cannot declare. +func NewDeployService(h DeployServiceHandlers, opts ...pipelinedef.Option) *DeployServiceDefinition { + cfg := pipelinedef.Config{ + ID: "deploy-service", + NewInput: func() proto.Message { return &DeployServiceInput{} }, Reduce: func(view durable.ReduceView) proto.Message { return ReduceDeployServiceOutput(h, view) }, @@ -161,7 +160,11 @@ func NewDeployService(h DeployServiceHandlers, mw ...durable.Middleware) *Deploy }, }, }, - })} + } + for _, o := range opts { + o(&cfg) + } + return &DeployServiceDefinition{def: pipelinedef.New(cfg)} } // Bind registers the definition with an engine. It is allowed only before diff --git a/examples/release-train/releasepb/release_durable.pb.go b/examples/release-train/releasepb/release_durable.pb.go index 5156091..f27afb9 100644 --- a/examples/release-train/releasepb/release_durable.pb.go +++ b/examples/release-train/releasepb/release_durable.pb.go @@ -114,13 +114,12 @@ type DeployServiceDefinition struct { } // NewDeployService assembles the "deploy-service" pipeline definition -// from its handlers. mw wraps this pipeline's operations alone, forward and -// unwind alike, inside any engine-level middleware; the first is outermost. -func NewDeployService(h DeployServiceHandlers, mw ...durable.Middleware) *DeployServiceDefinition { - return &DeployServiceDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "deploy-service", - Middleware: mw, - NewInput: func() proto.Message { return &DeployServiceInput{} }, +// from its handlers; opts (pipelinedef.WithMiddleware) configure what the +// proto cannot declare. +func NewDeployService(h DeployServiceHandlers, opts ...pipelinedef.Option) *DeployServiceDefinition { + cfg := pipelinedef.Config{ + ID: "deploy-service", + NewInput: func() proto.Message { return &DeployServiceInput{} }, Reduce: func(view durable.ReduceView) proto.Message { return ReduceDeployServiceOutput(h, view) }, @@ -178,7 +177,11 @@ func NewDeployService(h DeployServiceHandlers, mw ...durable.Middleware) *Deploy }, }, }, - })} + } + for _, o := range opts { + o(&cfg) + } + return &DeployServiceDefinition{def: pipelinedef.New(cfg)} } // Bind registers the definition with an engine. It is allowed only before @@ -372,13 +375,12 @@ type ReleaseTrainDefinition struct { } // NewReleaseTrain assembles the "release-train" pipeline definition -// from its handlers. mw wraps this pipeline's operations alone, forward and -// unwind alike, inside any engine-level middleware; the first is outermost. -func NewReleaseTrain(h ReleaseTrainHandlers, mw ...durable.Middleware) *ReleaseTrainDefinition { - return &ReleaseTrainDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "release-train", - Middleware: mw, - NewInput: func() proto.Message { return &ReleaseTrainInput{} }, +// from its handlers; opts (pipelinedef.WithMiddleware) configure what the +// proto cannot declare. +func NewReleaseTrain(h ReleaseTrainHandlers, opts ...pipelinedef.Option) *ReleaseTrainDefinition { + cfg := pipelinedef.Config{ + ID: "release-train", + NewInput: func() proto.Message { return &ReleaseTrainInput{} }, Steps: []pipelinedef.Step{ { ID: "plan/v1", @@ -410,7 +412,11 @@ func NewReleaseTrain(h ReleaseTrainHandlers, mw ...durable.Middleware) *ReleaseT }, }, }, - })} + } + for _, o := range opts { + o(&cfg) + } + return &ReleaseTrainDefinition{def: pipelinedef.New(cfg)} } // Bind registers the definition with an engine. It is allowed only before diff --git a/examples/snapshots/snapshotspb/snapshots_durable.pb.go b/examples/snapshots/snapshotspb/snapshots_durable.pb.go index 90f60c3..9f08f53 100644 --- a/examples/snapshots/snapshotspb/snapshots_durable.pb.go +++ b/examples/snapshots/snapshotspb/snapshots_durable.pb.go @@ -128,13 +128,12 @@ type CreateSnapshotDefinition struct { } // NewCreateSnapshot assembles the "create-snapshot" pipeline definition -// from its handlers. mw wraps this pipeline's operations alone, forward and -// unwind alike, inside any engine-level middleware; the first is outermost. -func NewCreateSnapshot(h CreateSnapshotHandlers, mw ...durable.Middleware) *CreateSnapshotDefinition { - return &CreateSnapshotDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "create-snapshot", - Middleware: mw, - NewInput: func() proto.Message { return &CreateSnapshotInput{} }, +// from its handlers; opts (pipelinedef.WithMiddleware) configure what the +// proto cannot declare. +func NewCreateSnapshot(h CreateSnapshotHandlers, opts ...pipelinedef.Option) *CreateSnapshotDefinition { + cfg := pipelinedef.Config{ + ID: "create-snapshot", + NewInput: func() proto.Message { return &CreateSnapshotInput{} }, Reduce: func(view durable.ReduceView) proto.Message { return ReduceCreateSnapshotOutput(h, view) }, @@ -190,7 +189,11 @@ func NewCreateSnapshot(h CreateSnapshotHandlers, mw ...durable.Middleware) *Crea }, }, }, - })} + } + for _, o := range opts { + o(&cfg) + } + return &CreateSnapshotDefinition{def: pipelinedef.New(cfg)} } // Bind registers the definition with an engine. It is allowed only before diff --git a/examples/tracing-otel/orderspb/orders_durable.pb.go b/examples/tracing-otel/orderspb/orders_durable.pb.go index 4f03150..9df80ef 100644 --- a/examples/tracing-otel/orderspb/orders_durable.pb.go +++ b/examples/tracing-otel/orderspb/orders_durable.pb.go @@ -108,13 +108,12 @@ type FulfillOrderDefinition struct { } // NewFulfillOrder assembles the "fulfill-order" pipeline definition -// from its handlers. mw wraps this pipeline's operations alone, forward and -// unwind alike, inside any engine-level middleware; the first is outermost. -func NewFulfillOrder(h FulfillOrderHandlers, mw ...durable.Middleware) *FulfillOrderDefinition { - return &FulfillOrderDefinition{def: pipelinedef.New(pipelinedef.Config{ - ID: "fulfill-order", - Middleware: mw, - NewInput: func() proto.Message { return &FulfillOrderInput{} }, +// from its handlers; opts (pipelinedef.WithMiddleware) configure what the +// proto cannot declare. +func NewFulfillOrder(h FulfillOrderHandlers, opts ...pipelinedef.Option) *FulfillOrderDefinition { + cfg := pipelinedef.Config{ + ID: "fulfill-order", + NewInput: func() proto.Message { return &FulfillOrderInput{} }, Reduce: func(view durable.ReduceView) proto.Message { return ReduceFulfillOrderOutput(h, view) }, @@ -161,7 +160,11 @@ func NewFulfillOrder(h FulfillOrderHandlers, mw ...durable.Middleware) *FulfillO }, }, }, - })} + } + for _, o := range opts { + o(&cfg) + } + return &FulfillOrderDefinition{def: pipelinedef.New(cfg)} } // Bind registers the definition with an engine. It is allowed only before diff --git a/middleware.go b/middleware.go index 86ec5a4..47fff39 100644 --- a/middleware.go +++ b/middleware.go @@ -20,9 +20,8 @@ type Handler func(ctx context.Context, inv Invocation) (proto.Message, error) // func(http.Handler) http.Handler. Use it for cross-cutting concerns such // as logging, metrics, tracing spans, or per-operation timeouts. It is // installed engine-wide with engine.WithMiddleware, or on one pipeline -// through pipelinedef.Config.Middleware (the generated constructor's -// variadic); the engine chain is outermost, and both are composed once -// at Engine.Bind. +// with pipelinedef.WithMiddleware on its generated constructor; the +// engine chain is outermost, and both are composed once at Engine.Bind. // // Middleware runs once per attempt, inside the durable attempt // reservation: it inherits the operation's at-least-once semantics and diff --git a/pipelinedef/pipelinedef.go b/pipelinedef/pipelinedef.go index 32997aa..8d22ec3 100644 --- a/pipelinedef/pipelinedef.go +++ b/pipelinedef/pipelinedef.go @@ -66,8 +66,9 @@ type Config struct { // Middleware wraps this pipeline's operations alone, forward and // unwind alike (Invocation.Phase distinguishes them), inside any // engine-level middleware: engine middleware is outermost, then these - // in order, first listed outermost, then the handler. Composition is - // fixed at Engine.Bind, which rejects a nil entry. + // in order, first listed outermost, then the handler. Generated + // constructors set it through WithMiddleware. Composition is fixed at + // Engine.Bind, which rejects a nil entry. Middleware []durable.Middleware // NewInput constructs an empty Input message; nil for an Input-less @@ -123,3 +124,17 @@ func StepRef(id durable.StepID) durable.StepRef { return durable.StepRef{Step: i func StateStepRef[T proto.Message](id durable.StepID, new func() T) durable.StateStepRef[T] { return durable.StateStepRef[T]{Step: id, New: new} } + +// Option configures a generated pipeline's Config at construction: +// NewXxx(h, opts...). Options are the runtime knobs a proto cannot +// declare; today that is the pipeline's own middleware. +type Option func(*Config) + +// WithMiddleware installs middleware on this pipeline alone, wrapping +// its operations, forward and unwind alike, inside any engine-level +// middleware; the first listed is outermost. Repeated options append. +func WithMiddleware(mw ...durable.Middleware) Option { + return func(c *Config) { + c.Middleware = append(c.Middleware, mw...) + } +} diff --git a/spec/02-authoring.md b/spec/02-authoring.md index 67f6d22..c8bffa5 100644 --- a/spec/02-authoring.md +++ b/spec/02-authoring.md @@ -262,12 +262,13 @@ type ProvisionMachineHandlers interface { ReduceOutput(*ProvisionMachine) *ProvisionMachineOutput } -func NewProvisionMachine(h ProvisionMachineHandlers, mw ...durable.Middleware) *ProvisionMachineDefinition +func NewProvisionMachine(h ProvisionMachineHandlers, opts ...pipelinedef.Option) *ProvisionMachineDefinition ``` -The variadic installs pipeline-level middleware, wrapping this -pipeline's operations alone inside any engine-level chain (see -[04-engine](04-engine.md#middleware)). +The options configure what the proto cannot declare; +`pipelinedef.WithMiddleware(mw...)` installs pipeline-level middleware, +wrapping this pipeline's operations alone inside any engine-level chain +(see [04-engine](04-engine.md#middleware)). One type implements the whole pipeline, so its dependencies are declared once; a Step the implementor lacks, a Step added to the proto included, @@ -593,8 +594,8 @@ definition := machines.NewProvisionMachine(&handlers{cloud: c}) The constructor takes one implementation of the pipeline's handler interface: a missing or mis-typed Step method, a Step added to the proto -included, fails to compile. Middleware for this pipeline alone rides the -variadic: `machines.NewProvisionMachine(&handlers{cloud: c}, notFoundIsPermanent)`. +included, fails to compile. Middleware for this pipeline alone is an +option: `machines.NewProvisionMachine(&handlers{cloud: c}, pipelinedef.WithMiddleware(notFoundIsPermanent))`. --- diff --git a/spec/04-engine.md b/spec/04-engine.md index 7c325b8..7b5f640 100644 --- a/spec/04-engine.md +++ b/spec/04-engine.md @@ -648,8 +648,9 @@ The first middleware is outermost. `Invocation.Phase()` distinguishes forward from unwind operations. **Pipeline-level middleware.** A definition MAY carry its own chain, -`pipelinedef.Config.Middleware`; the generated constructor exposes it -as a variadic, `NewXxx(h, mw...)`. It wraps only that pipeline's +`pipelinedef.Config.Middleware`, set on a generated constructor with the +`pipelinedef.WithMiddleware` option: `NewXxx(h, pipelinedef.WithMiddleware(mw...))`. +It wraps only that pipeline's operations, forward and unwind alike, and composes inside the engine-level chain: diff --git a/spec/05-codegen.md b/spec/05-codegen.md index 964f20f..a752374 100644 --- a/spec/05-codegen.md +++ b/spec/05-codegen.md @@ -53,8 +53,8 @@ Published protobuf extensions MUST use globally allocated extension numbers. - one handler interface, `XxxHandlers`: a method per Step named after the Step, `Unwind` for each Step that unwinds, and `ReduceOutput` / `ReduceFailure` when the pipeline declares outputs, -- the Pipeline constructor, `NewXxx(h XxxHandlers, mw ...durable.Middleware)`; - the variadic is the pipeline's own middleware chain +- the Pipeline constructor, `NewXxx(h XxxHandlers, opts ...pipelinedef.Option)`; + `pipelinedef.WithMiddleware` is the pipeline's own middleware chain ([04-engine](04-engine.md#middleware)), - `ReduceXxxOutput(h, view)` and `ReduceXxxFailure(h, view)`, the folds the engine reduces through and reducer tests call, diff --git a/spec/http-analogy.md b/spec/http-analogy.md index e555a69..fa588f9 100644 --- a/spec/http-analogy.md +++ b/spec/http-analogy.md @@ -12,7 +12,7 @@ analogy with Go's `net/http`; the normative contracts live in | `http.Handler` | `durable.Handler` — the uniform type-erased operation `func(ctx, Invocation) (proto.Message, error)` | | `http.HandlerFunc` | the generated per-pipeline handler interface's methods, erased by the constructor's adapters | | middleware `func(http.Handler) http.Handler` | `durable.Middleware`, installed engine-wide with `engine.WithMiddleware` | -| wrapping one handler at mux registration, `mux.Handle(pattern, mw(h))` | pipeline-level middleware: `pipelinedef.Config.Middleware`, the generated `NewXxx(h, mw...)` variadic; composes inside the engine chain | +| wrapping one handler at mux registration, `mux.Handle(pattern, mw(h))` | pipeline-level middleware: `pipelinedef.Config.Middleware`, set with `pipelinedef.WithMiddleware` on the generated constructor; composes inside the engine chain | | `ServeMux` and route patterns | the protobuf Pipeline topology; `StepID` plays the route pattern | | `*http.Request` | `durable.Invocation` — identity, attempt, phase, input, State lookup; an interface, so tests can fake it | | `http.ResponseWriter` | — (deliberately absent; see below) | From bdd906231f6ff40291705125403b2d04b7b8b65c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Gra=C3=B1a?= Date: Thu, 17 Sep 2026 15:38:50 -0300 Subject: [PATCH 3/3] Generated packages alias the constructor options Each generated package carries Option (= pipelinedef.Option) and WithMiddleware, so wiring code configures a pipeline as machinespb.NewProvisionMachine(h, machinespb.WithMiddleware(mw)) and never imports pipelinedef. Emitted once per Go package. --- README.md | 2 +- cmd/protoc-gen-durable/internal/gen/gen.go | 38 ++++++++++++++++--- doc.go | 8 ++-- docs/tour.md | 2 +- engine/middleware.go | 4 +- .../machinespb/machines_durable.pb.go | 25 +++++++++--- examples/machines/main_test.go | 3 +- .../legacypb/legacy_durable.pb.go | 19 ++++++++-- .../releasepb/release_durable.pb.go | 25 +++++++++--- .../snapshotspb/snapshots_durable.pb.go | 19 ++++++++-- .../orderspb/orders_durable.pb.go | 19 ++++++++-- middleware.go | 2 +- pipelinedef/pipelinedef.go | 7 +++- spec/02-authoring.md | 8 ++-- spec/04-engine.md | 4 +- spec/05-codegen.md | 8 ++-- spec/http-analogy.md | 2 +- 17 files changed, 147 insertions(+), 48 deletions(-) diff --git a/README.md b/README.md index 27a777c..66a3234 100644 --- a/README.md +++ b/README.md @@ -162,7 +162,7 @@ durable-scale histogram buckets, `trace_id`/`span_id` log correlation, and an opt-in W3C Baggage relay. Everything is declared once, at engine construction; a concern that belongs to one pipeline is declared once on that pipeline's constructor instead, as -`machinespb.NewProvisionMachine(h, pipelinedef.WithMiddleware(notFoundIsPermanent))`: +`machinespb.NewProvisionMachine(h, machinespb.WithMiddleware(notFoundIsPermanent))`: ```go obs, _ := durableotel.NewObserver() diff --git a/cmd/protoc-gen-durable/internal/gen/gen.go b/cmd/protoc-gen-durable/internal/gen/gen.go index a678e9e..25c4b0c 100644 --- a/cmd/protoc-gen-durable/internal/gen/gen.go +++ b/cmd/protoc-gen-durable/internal/gen/gen.go @@ -188,8 +188,13 @@ func Generate(p *protogen.Plugin) error { } byFile[pl.file] = append(byFile[pl.file], pl) } + // The constructor options are package-level, emitted once per Go + // package into its first generated file. + optionsDone := make(map[protogen.GoImportPath]bool) for _, f := range fileOrder { - emitFile(p, f, byFile[f]) + emitOptions := !optionsDone[f.GoImportPath] + optionsDone[f.GoImportPath] = true + emitFile(p, f, byFile[f], emitOptions) } return nil } @@ -252,7 +257,7 @@ func lowerFirst(s string) string { // pipeline the proto file declares: step references, invocation types, // handler interfaces, definition constructors, reducer plumbing, and the // bound pipeline, run, and result types. -func emitFile(p *protogen.Plugin, f *protogen.File, pipelines []*pipelineDecl) { +func emitFile(p *protogen.Plugin, f *protogen.File, pipelines []*pipelineDecl, emitOptions bool) { filename := f.GeneratedFilenamePrefix + "_durable.pb.go" g := p.NewGeneratedFile(filename, f.GoImportPath) @@ -266,11 +271,34 @@ func emitFile(p *protogen.Plugin, f *protogen.File, pipelines []*pipelineDecl) { g.P("package ", f.GoPackageName) g.P() + if emitOptions { + emitConstructorOptions(g) + } for _, pl := range pipelines { emitPipeline(g, pl) } } +// emitConstructorOptions emits the package's constructor options: Option +// is an alias of pipelinedef.Option and WithMiddleware forwards to +// pipelinedef.WithMiddleware, so wiring code configures a generated +// pipeline without importing pipelinedef. +func emitConstructorOptions(g *protogen.GeneratedFile) { + g.P("// Option configures a pipeline of this package at construction,") + g.P("// NewXxx(h, opts...): the runtime knobs a proto cannot declare. It is") + g.P("// pipelinedef.Option, so the two are interchangeable.") + g.P("type Option = ", g.QualifiedGoIdent(defPkg.Ident("Option"))) + g.P() + g.P("// WithMiddleware installs middleware on the constructed pipeline alone,") + g.P("// wrapping its operations, forward and unwind alike, inside any") + g.P("// engine-level middleware; the first listed is outermost. Repeated") + g.P("// options append.") + g.P("func WithMiddleware(mw ...", g.QualifiedGoIdent(durablePkg.Ident("Middleware")), ") Option {") + g.P("return ", g.QualifiedGoIdent(defPkg.Ident("WithMiddleware")), "(mw...)") + g.P("}") + g.P() +} + func emitPipeline(g *protogen.GeneratedFile, pl *pipelineDecl) { for _, s := range pl.steps { emitStepRef(g, s) @@ -478,9 +506,9 @@ func emitDefinition(g *protogen.GeneratedFile, pl *pipelineDecl) { g.P() g.P("// New", name, " assembles the ", strconv(pl.opts.GetId()), " pipeline definition") - g.P("// from its handlers; opts (pipelinedef.WithMiddleware) configure what the") - g.P("// proto cannot declare.") - g.P("func New", name, "(h ", pl.handlersName(), ", opts ...", g.QualifiedGoIdent(defPkg.Ident("Option")), ") *", name, "Definition {") + g.P("// from its handlers; opts (WithMiddleware) configure what the proto") + g.P("// cannot declare.") + g.P("func New", name, "(h ", pl.handlersName(), ", opts ...Option) *", name, "Definition {") g.P("cfg := ", g.QualifiedGoIdent(defPkg.Ident("Config")), "{") g.P("ID: ", strconv(pl.opts.GetId()), ",") if ms := pl.opts.GetMutexes(); len(ms) > 0 { diff --git a/doc.go b/doc.go index 01c3d2b..491edce 100644 --- a/doc.go +++ b/doc.go @@ -45,8 +45,8 @@ // // Middleware. Handler and Middleware are the net/http-shaped operation // layer every attempt passes through (installed engine-wide with -// engine.WithMiddleware, or on one pipeline with pipelinedef.WithMiddleware -// on its generated constructor). AwaitRequest, AwaitTimeout, FailureInfo, +// engine.WithMiddleware, or on one pipeline with the generated package's +// WithMiddleware option on its constructor). AwaitRequest, AwaitTimeout, FailureInfo, // FailureCause, and FailureReason classify a handler's return the way // the engine will; PreemptedError and ErrEngineStopping name why an // attempt ctx died, for middleware that labels spans. @@ -105,8 +105,8 @@ // the wiring side of an application. // - pipelinedef is the type-erased pipeline description generated // code builds and Engine.Bind validates; hand-rolled definitions use -// it too. Wiring code imports it for the constructor options -// (WithMiddleware); handler code never does. +// it too. Generated packages alias its constructor options (Option, +// WithMiddleware), so application code never imports it. // - kernel is the shared vocabulary — identities, phases, outcomes, // parks, failure records — that every other package builds on. This // package aliases all of it, so user code never imports kernel. diff --git a/docs/tour.md b/docs/tour.md index 21e80ef..a062535 100644 --- a/docs/tour.md +++ b/docs/tour.md @@ -594,7 +594,7 @@ lifecycle events (`observe.Observer`), and snapshots occupancy OpenTelemetry integration (per-attempt spans linked to the scheduling trace, metrics, log correlation, W3C Baggage relay), declared once at engine construction — or, for a concern that belongs to one pipeline, -on that pipeline's constructor (`deploypb.NewDeployService(h, pipelinedef.WithMiddleware(mw))`), +on that pipeline's constructor (`deploypb.NewDeployService(h, deploypb.WithMiddleware(mw))`), composed inside the engine's chain; [examples/tracing-otel](../examples/tracing-otel/) demonstrates it end to end. diff --git a/engine/middleware.go b/engine/middleware.go index d739fd8..2de38dd 100644 --- a/engine/middleware.go +++ b/engine/middleware.go @@ -10,8 +10,8 @@ import ( // executes, forward and unwind alike; use Invocation.Phase to distinguish // them. The first middleware is the outermost, following the net/http // convention: WithMiddleware(a, b) yields a(b(handler)). Pipeline-level -// middleware (pipelinedef.WithMiddleware on a generated constructor, or -// pipelinedef.Config.Middleware) composes inside this chain: engine +// middleware (a generated package's WithMiddleware constructor option, +// or pipelinedef.Config.Middleware) composes inside this chain: engine // middleware is outermost. func WithMiddleware(mw ...durable.Middleware) Option { return func(e *Engine) { diff --git a/examples/machines/machinespb/machines_durable.pb.go b/examples/machines/machinespb/machines_durable.pb.go index b4869f8..3d8e46d 100644 --- a/examples/machines/machinespb/machines_durable.pb.go +++ b/examples/machines/machinespb/machines_durable.pb.go @@ -15,6 +15,19 @@ import ( sync "sync" ) +// Option configures a pipeline of this package at construction, +// NewXxx(h, opts...): the runtime knobs a proto cannot declare. It is +// pipelinedef.Option, so the two are interchangeable. +type Option = pipelinedef.Option + +// WithMiddleware installs middleware on the constructed pipeline alone, +// wrapping its operations, forward and unwind alike, inside any +// engine-level middleware; the first listed is outermost. Repeated +// options append. +func WithMiddleware(mw ...durable.Middleware) Option { + return pipelinedef.WithMiddleware(mw...) +} + // ProvisionMachine_ValidateStep is the reference to the stateless step "validate/v1". // It is not accepted by State lookup. var ProvisionMachine_ValidateStep = pipelinedef.StepRef("validate/v1") @@ -112,9 +125,9 @@ type ProvisionMachineDefinition struct { } // NewProvisionMachine assembles the "provision-machine" pipeline definition -// from its handlers; opts (pipelinedef.WithMiddleware) configure what the -// proto cannot declare. -func NewProvisionMachine(h ProvisionMachineHandlers, opts ...pipelinedef.Option) *ProvisionMachineDefinition { +// from its handlers; opts (WithMiddleware) configure what the proto +// cannot declare. +func NewProvisionMachine(h ProvisionMachineHandlers, opts ...Option) *ProvisionMachineDefinition { cfg := pipelinedef.Config{ ID: "provision-machine", Mutexes: []string{"machine-lifecycle"}, @@ -343,9 +356,9 @@ type DecommissionMachineDefinition struct { } // NewDecommissionMachine assembles the "decommission-machine" pipeline definition -// from its handlers; opts (pipelinedef.WithMiddleware) configure what the -// proto cannot declare. -func NewDecommissionMachine(h DecommissionMachineHandlers, opts ...pipelinedef.Option) *DecommissionMachineDefinition { +// from its handlers; opts (WithMiddleware) configure what the proto +// cannot declare. +func NewDecommissionMachine(h DecommissionMachineHandlers, opts ...Option) *DecommissionMachineDefinition { cfg := pipelinedef.Config{ ID: "decommission-machine", Mutexes: []string{"machine-lifecycle"}, diff --git a/examples/machines/main_test.go b/examples/machines/main_test.go index 8fc388c..50bd3aa 100644 --- a/examples/machines/main_test.go +++ b/examples/machines/main_test.go @@ -10,7 +10,6 @@ import ( "github.com/dangra/durable" "github.com/dangra/durable/engine" "github.com/dangra/durable/examples/machines/machinespb" - "github.com/dangra/durable/pipelinedef" "github.com/dangra/durable/store/mem" "google.golang.org/protobuf/proto" ) @@ -221,7 +220,7 @@ func TestPipelineMiddleware(t *testing.T) { } c := newCloud() eng := engine.New(mem.New()) - provision, err := machinespb.NewProvisionMachine(&handlers{cloud: c}, pipelinedef.WithMiddleware(counting)).Bind(eng) + provision, err := machinespb.NewProvisionMachine(&handlers{cloud: c}, machinespb.WithMiddleware(counting)).Bind(eng) if err != nil { t.Fatalf("Bind: %v", err) } diff --git a/examples/release-train/legacypb/legacy_durable.pb.go b/examples/release-train/legacypb/legacy_durable.pb.go index 49cbfc0..2e7f149 100644 --- a/examples/release-train/legacypb/legacy_durable.pb.go +++ b/examples/release-train/legacypb/legacy_durable.pb.go @@ -14,6 +14,19 @@ import ( sync "sync" ) +// Option configures a pipeline of this package at construction, +// NewXxx(h, opts...): the runtime knobs a proto cannot declare. It is +// pipelinedef.Option, so the two are interchangeable. +type Option = pipelinedef.Option + +// WithMiddleware installs middleware on the constructed pipeline alone, +// wrapping its operations, forward and unwind alike, inside any +// engine-level middleware; the first listed is outermost. Repeated +// options append. +func WithMiddleware(mw ...durable.Middleware) Option { + return pipelinedef.WithMiddleware(mw...) +} + // 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{} }) @@ -108,9 +121,9 @@ type DeployServiceDefinition struct { } // NewDeployService assembles the "deploy-service" pipeline definition -// from its handlers; opts (pipelinedef.WithMiddleware) configure what the -// proto cannot declare. -func NewDeployService(h DeployServiceHandlers, opts ...pipelinedef.Option) *DeployServiceDefinition { +// from its handlers; opts (WithMiddleware) configure what the proto +// cannot declare. +func NewDeployService(h DeployServiceHandlers, opts ...Option) *DeployServiceDefinition { cfg := pipelinedef.Config{ ID: "deploy-service", NewInput: func() proto.Message { return &DeployServiceInput{} }, diff --git a/examples/release-train/releasepb/release_durable.pb.go b/examples/release-train/releasepb/release_durable.pb.go index f27afb9..9420050 100644 --- a/examples/release-train/releasepb/release_durable.pb.go +++ b/examples/release-train/releasepb/release_durable.pb.go @@ -15,6 +15,19 @@ import ( sync "sync" ) +// Option configures a pipeline of this package at construction, +// NewXxx(h, opts...): the runtime knobs a proto cannot declare. It is +// pipelinedef.Option, so the two are interchangeable. +type Option = pipelinedef.Option + +// WithMiddleware installs middleware on the constructed pipeline alone, +// wrapping its operations, forward and unwind alike, inside any +// engine-level middleware; the first listed is outermost. Repeated +// options append. +func WithMiddleware(mw ...durable.Middleware) Option { + return pipelinedef.WithMiddleware(mw...) +} + // 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{} }) @@ -114,9 +127,9 @@ type DeployServiceDefinition struct { } // NewDeployService assembles the "deploy-service" pipeline definition -// from its handlers; opts (pipelinedef.WithMiddleware) configure what the -// proto cannot declare. -func NewDeployService(h DeployServiceHandlers, opts ...pipelinedef.Option) *DeployServiceDefinition { +// from its handlers; opts (WithMiddleware) configure what the proto +// cannot declare. +func NewDeployService(h DeployServiceHandlers, opts ...Option) *DeployServiceDefinition { cfg := pipelinedef.Config{ ID: "deploy-service", NewInput: func() proto.Message { return &DeployServiceInput{} }, @@ -375,9 +388,9 @@ type ReleaseTrainDefinition struct { } // NewReleaseTrain assembles the "release-train" pipeline definition -// from its handlers; opts (pipelinedef.WithMiddleware) configure what the -// proto cannot declare. -func NewReleaseTrain(h ReleaseTrainHandlers, opts ...pipelinedef.Option) *ReleaseTrainDefinition { +// from its handlers; opts (WithMiddleware) configure what the proto +// cannot declare. +func NewReleaseTrain(h ReleaseTrainHandlers, opts ...Option) *ReleaseTrainDefinition { cfg := pipelinedef.Config{ ID: "release-train", NewInput: func() proto.Message { return &ReleaseTrainInput{} }, diff --git a/examples/snapshots/snapshotspb/snapshots_durable.pb.go b/examples/snapshots/snapshotspb/snapshots_durable.pb.go index 9f08f53..9c8643b 100644 --- a/examples/snapshots/snapshotspb/snapshots_durable.pb.go +++ b/examples/snapshots/snapshotspb/snapshots_durable.pb.go @@ -14,6 +14,19 @@ import ( sync "sync" ) +// Option configures a pipeline of this package at construction, +// NewXxx(h, opts...): the runtime knobs a proto cannot declare. It is +// pipelinedef.Option, so the two are interchangeable. +type Option = pipelinedef.Option + +// WithMiddleware installs middleware on the constructed pipeline alone, +// wrapping its operations, forward and unwind alike, inside any +// engine-level middleware; the first listed is outermost. Repeated +// options append. +func WithMiddleware(mw ...durable.Middleware) Option { + return pipelinedef.WithMiddleware(mw...) +} + // CreateSnapshot_FreezeVolumeStep is the typed reference to the state-producing step "freeze-volume/v1". var CreateSnapshot_FreezeVolumeStep = pipelinedef.StateStepRef("freeze-volume/v1", func() *CreateSnapshot_FreezeVolume { return &CreateSnapshot_FreezeVolume{} }) @@ -128,9 +141,9 @@ type CreateSnapshotDefinition struct { } // NewCreateSnapshot assembles the "create-snapshot" pipeline definition -// from its handlers; opts (pipelinedef.WithMiddleware) configure what the -// proto cannot declare. -func NewCreateSnapshot(h CreateSnapshotHandlers, opts ...pipelinedef.Option) *CreateSnapshotDefinition { +// from its handlers; opts (WithMiddleware) configure what the proto +// cannot declare. +func NewCreateSnapshot(h CreateSnapshotHandlers, opts ...Option) *CreateSnapshotDefinition { cfg := pipelinedef.Config{ ID: "create-snapshot", NewInput: func() proto.Message { return &CreateSnapshotInput{} }, diff --git a/examples/tracing-otel/orderspb/orders_durable.pb.go b/examples/tracing-otel/orderspb/orders_durable.pb.go index 9df80ef..8c0a498 100644 --- a/examples/tracing-otel/orderspb/orders_durable.pb.go +++ b/examples/tracing-otel/orderspb/orders_durable.pb.go @@ -14,6 +14,19 @@ import ( sync "sync" ) +// Option configures a pipeline of this package at construction, +// NewXxx(h, opts...): the runtime knobs a proto cannot declare. It is +// pipelinedef.Option, so the two are interchangeable. +type Option = pipelinedef.Option + +// WithMiddleware installs middleware on the constructed pipeline alone, +// wrapping its operations, forward and unwind alike, inside any +// engine-level middleware; the first listed is outermost. Repeated +// options append. +func WithMiddleware(mw ...durable.Middleware) Option { + return pipelinedef.WithMiddleware(mw...) +} + // FulfillOrder_ReserveStockStep is the typed reference to the state-producing step "reserve-stock/v1". var FulfillOrder_ReserveStockStep = pipelinedef.StateStepRef("reserve-stock/v1", func() *FulfillOrder_ReserveStock { return &FulfillOrder_ReserveStock{} }) @@ -108,9 +121,9 @@ type FulfillOrderDefinition struct { } // NewFulfillOrder assembles the "fulfill-order" pipeline definition -// from its handlers; opts (pipelinedef.WithMiddleware) configure what the -// proto cannot declare. -func NewFulfillOrder(h FulfillOrderHandlers, opts ...pipelinedef.Option) *FulfillOrderDefinition { +// from its handlers; opts (WithMiddleware) configure what the proto +// cannot declare. +func NewFulfillOrder(h FulfillOrderHandlers, opts ...Option) *FulfillOrderDefinition { cfg := pipelinedef.Config{ ID: "fulfill-order", NewInput: func() proto.Message { return &FulfillOrderInput{} }, diff --git a/middleware.go b/middleware.go index 47fff39..46c82be 100644 --- a/middleware.go +++ b/middleware.go @@ -20,7 +20,7 @@ type Handler func(ctx context.Context, inv Invocation) (proto.Message, error) // func(http.Handler) http.Handler. Use it for cross-cutting concerns such // as logging, metrics, tracing spans, or per-operation timeouts. It is // installed engine-wide with engine.WithMiddleware, or on one pipeline -// with pipelinedef.WithMiddleware on its generated constructor; the +// with its generated package's WithMiddleware constructor option; the // engine chain is outermost, and both are composed once at Engine.Bind. // // Middleware runs once per attempt, inside the durable attempt diff --git a/pipelinedef/pipelinedef.go b/pipelinedef/pipelinedef.go index 8d22ec3..f39120c 100644 --- a/pipelinedef/pipelinedef.go +++ b/pipelinedef/pipelinedef.go @@ -67,7 +67,8 @@ type Config struct { // unwind alike (Invocation.Phase distinguishes them), inside any // engine-level middleware: engine middleware is outermost, then these // in order, first listed outermost, then the handler. Generated - // constructors set it through WithMiddleware. Composition is fixed at + // packages set it through their WithMiddleware option, an alias of + // the one here. Composition is fixed at // Engine.Bind, which rejects a nil entry. Middleware []durable.Middleware @@ -127,7 +128,9 @@ func StateStepRef[T proto.Message](id durable.StepID, new func() T) durable.Stat // Option configures a generated pipeline's Config at construction: // NewXxx(h, opts...). Options are the runtime knobs a proto cannot -// declare; today that is the pipeline's own middleware. +// declare; today that is the pipeline's own middleware. Generated +// packages alias Option and WithMiddleware, so wiring code reaches them +// without importing this package. type Option func(*Config) // WithMiddleware installs middleware on this pipeline alone, wrapping diff --git a/spec/02-authoring.md b/spec/02-authoring.md index c8bffa5..1e652e8 100644 --- a/spec/02-authoring.md +++ b/spec/02-authoring.md @@ -262,11 +262,11 @@ type ProvisionMachineHandlers interface { ReduceOutput(*ProvisionMachine) *ProvisionMachineOutput } -func NewProvisionMachine(h ProvisionMachineHandlers, opts ...pipelinedef.Option) *ProvisionMachineDefinition +func NewProvisionMachine(h ProvisionMachineHandlers, opts ...Option) *ProvisionMachineDefinition ``` -The options configure what the proto cannot declare; -`pipelinedef.WithMiddleware(mw...)` installs pipeline-level middleware, +The options configure what the proto cannot declare; the generated +package's `WithMiddleware(mw...)` installs pipeline-level middleware, wrapping this pipeline's operations alone inside any engine-level chain (see [04-engine](04-engine.md#middleware)). @@ -595,7 +595,7 @@ definition := machines.NewProvisionMachine(&handlers{cloud: c}) The constructor takes one implementation of the pipeline's handler interface: a missing or mis-typed Step method, a Step added to the proto included, fails to compile. Middleware for this pipeline alone is an -option: `machines.NewProvisionMachine(&handlers{cloud: c}, pipelinedef.WithMiddleware(notFoundIsPermanent))`. +option: `machines.NewProvisionMachine(&handlers{cloud: c}, machines.WithMiddleware(notFoundIsPermanent))`. --- diff --git a/spec/04-engine.md b/spec/04-engine.md index 7b5f640..78c6817 100644 --- a/spec/04-engine.md +++ b/spec/04-engine.md @@ -649,7 +649,9 @@ forward from unwind operations. **Pipeline-level middleware.** A definition MAY carry its own chain, `pipelinedef.Config.Middleware`, set on a generated constructor with the -`pipelinedef.WithMiddleware` option: `NewXxx(h, pipelinedef.WithMiddleware(mw...))`. +package's generated `WithMiddleware` option: `NewXxx(h, WithMiddleware(mw...))` +(an alias of `pipelinedef.WithMiddleware`, so wiring code imports only +the generated package). It wraps only that pipeline's operations, forward and unwind alike, and composes inside the engine-level chain: diff --git a/spec/05-codegen.md b/spec/05-codegen.md index a752374..0efa0c5 100644 --- a/spec/05-codegen.md +++ b/spec/05-codegen.md @@ -53,9 +53,11 @@ Published protobuf extensions MUST use globally allocated extension numbers. - one handler interface, `XxxHandlers`: a method per Step named after the Step, `Unwind` for each Step that unwinds, and `ReduceOutput` / `ReduceFailure` when the pipeline declares outputs, -- the Pipeline constructor, `NewXxx(h XxxHandlers, opts ...pipelinedef.Option)`; - `pipelinedef.WithMiddleware` is the pipeline's own middleware chain - ([04-engine](04-engine.md#middleware)), +- the Pipeline constructor, `NewXxx(h XxxHandlers, opts ...Option)`, and + once per package the `Option` alias of `pipelinedef.Option` and the + `WithMiddleware(mw...)` option, the pipeline's own middleware chain + ([04-engine](04-engine.md#middleware)); wiring code never imports + pipelinedef, - `ReduceXxxOutput(h, view)` and `ReduceXxxFailure(h, view)`, the folds the engine reduces through and reducer tests call, - runtime methods on the Pipeline marker type, diff --git a/spec/http-analogy.md b/spec/http-analogy.md index fa588f9..9059eaa 100644 --- a/spec/http-analogy.md +++ b/spec/http-analogy.md @@ -12,7 +12,7 @@ analogy with Go's `net/http`; the normative contracts live in | `http.Handler` | `durable.Handler` — the uniform type-erased operation `func(ctx, Invocation) (proto.Message, error)` | | `http.HandlerFunc` | the generated per-pipeline handler interface's methods, erased by the constructor's adapters | | middleware `func(http.Handler) http.Handler` | `durable.Middleware`, installed engine-wide with `engine.WithMiddleware` | -| wrapping one handler at mux registration, `mux.Handle(pattern, mw(h))` | pipeline-level middleware: `pipelinedef.Config.Middleware`, set with `pipelinedef.WithMiddleware` on the generated constructor; composes inside the engine chain | +| wrapping one handler at mux registration, `mux.Handle(pattern, mw(h))` | pipeline-level middleware: `pipelinedef.Config.Middleware`, set with the generated package's `WithMiddleware` option on the constructor; composes inside the engine chain | | `ServeMux` and route patterns | the protobuf Pipeline topology; `StepID` plays the route pattern | | `*http.Request` | `durable.Invocation` — identity, attempt, phase, input, State lookup; an interface, so tests can fake it | | `http.ResponseWriter` | — (deliberately absent; see below) |