From 8d539d5b401358eb3039459b87180703ee5fb959 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Wed, 2 Sep 2026 22:32:28 +0000 Subject: [PATCH] Consolidate delivery mapping construction Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- workflow/internal/execution/edgerunner.go | 15 +++------------ workflow/internal/execution/run.go | 11 +++++++++++ 2 files changed, 14 insertions(+), 12 deletions(-) diff --git a/workflow/internal/execution/edgerunner.go b/workflow/internal/execution/edgerunner.go index 7d348204..925045c6 100644 --- a/workflow/internal/execution/edgerunner.go +++ b/workflow/internal/execution/edgerunner.go @@ -259,10 +259,7 @@ func (em *EdgeRunner) PrepareDeliveryForEdge(ctx context.Context, edge workflow. return nil, nil } span.SetDeliveryStatus(observability.DeliveryStatusDelivered) - return &DeliveryMapping{ - Targets: targets, - Envelopes: envelopes, - }, nil + return newDeliveryMapping(envelopes, targets), nil } func (em *EdgeRunner) filterEnvelopesForTarget(ctx context.Context, envelopes []*MessageEnvelope, target *workflow.Executor) ([]*MessageEnvelope, error) { @@ -399,10 +396,7 @@ func (em *EdgeRunner) PrepareDeliveryForInput(ctx context.Context, envelope *Mes return nil, nil } span.SetDeliveryStatus(observability.DeliveryStatusDelivered) - return &DeliveryMapping{ - Targets: []*workflow.Executor{target}, - Envelopes: []*MessageEnvelope{envelope}, - }, nil + return newSingleDeliveryMapping(envelope, target), nil } // PrepareDeliveryForResponse prepares delivery of an external response to @@ -438,10 +432,7 @@ func (em *EdgeRunner) PrepareDeliveryForResponse(ctx context.Context, response * return nil, nil } span.SetDeliveryStatus(observability.DeliveryStatusDelivered) - return &DeliveryMapping{ - Targets: []*workflow.Executor{target}, - Envelopes: []*MessageEnvelope{envelope}, - }, nil + return newSingleDeliveryMapping(envelope, target), nil } func edgeGroupMetadata(edge workflow.Edge) observability.EdgeGroupMetadata { diff --git a/workflow/internal/execution/run.go b/workflow/internal/execution/run.go index 606262b0..ace902c9 100644 --- a/workflow/internal/execution/run.go +++ b/workflow/internal/execution/run.go @@ -18,6 +18,17 @@ type DeliveryMapping struct { Targets []*workflow.Executor } +func newDeliveryMapping(envelopes []*MessageEnvelope, targets []*workflow.Executor) *DeliveryMapping { + return &DeliveryMapping{ + Envelopes: envelopes, + Targets: targets, + } +} + +func newSingleDeliveryMapping(envelope *MessageEnvelope, target *workflow.Executor) *DeliveryMapping { + return newDeliveryMapping([]*MessageEnvelope{envelope}, []*workflow.Executor{target}) +} + func (d DeliveryMapping) MapInto(nextStep *StepContext) { for _, target := range d.Targets { messageQueue := nextStep.MessagesFor(target.ID)