Documentation
¶
Overview ¶
Package nodes_distribute provides distribute.map (bounded-concurrency fan-out over a slice) and distribute.reduce (fold the map output into a single value).
Typical use:
{ action: "log.info", items: [1, 2, 3], concurrency: 4 } -> distribute.map
{ strategy: "all_pass", items: .results } -> distribute.reduce
Index ¶
Constants ¶
const ID = "nodes_distribute"
Variables ¶
var DistributeMap = action.New("distribute.map", func(ctx context.Context, req DistributeMapReq) (DistributeMapRes, error) { if req.Action == "" { return DistributeMapRes{}, xerr.BadRequest("distribute.map: action is required") } if len(req.Items) == 0 { return DistributeMapRes{}, nil } if len(req.Items) > maxItems { return DistributeMapRes{}, xerr.BadRequest("distribute.map: items exceeds limit") } resolver := contracts.ActionResolverFromContext(ctx) if resolver == nil { return DistributeMapRes{}, xerr.Internal("distribute.map: no action resolver in context") } target, ok := resolver.Action(req.Action) if !ok { return DistributeMapRes{}, xerr.NotFound("distribute.map: action not found: " + req.Action) } concurrency := req.Concurrency if concurrency <= 0 { concurrency = defaultConcurrency } if concurrency > maxConcurrency { concurrency = maxConcurrency } items := make([]Item, len(req.Items)) group, groupCtx := errgroup.WithContext(ctx) group.SetLimit(concurrency) for i, value := range req.Items { index, val := i, value group.Go(func() error { if err := groupCtx.Err(); err != nil { items[index] = Item{Error: err.Error()} return err } output, err := invokeSafely(groupCtx, target, val) if err != nil { items[index] = Item{Error: err.Error()} return nil } items[index] = Item{OK: true, Result: output} return nil }) } if err := group.Wait(); err != nil { return DistributeMapRes{}, err } result := DistributeMapRes{Items: items} for i := range items { if items[i].OK { result.Succeeded++ } else { result.Failed++ } } return result, nil }).Description("Fan out one action over many items with bounded concurrency"). Tag("distribute", "fanout"). Build()
DistributeMap invokes Action once per input item with bounded concurrency, preserving input order in the output.
var DistributeReduce = action.New("distribute.reduce", func(_ context.Context, req DistributeReduceReq) (DistributeReduceRes, error) {
return runReduce(req)
}).Description("Fold distribute.map output into a single value").
Tag("distribute", "fold").
Build()
DistributeReduce folds a distribute.map result into a single value.
Functions ¶
Types ¶
type DistributeMapReq ¶
type DistributeMapReq struct {
Action string `json:"action" validate:"required"`
Concurrency int `json:"concurrency,omitempty"`
Items []any `json:"items" validate:"required"`
}
DistributeMapReq carries the fan-out configuration.
type DistributeMapRes ¶
type DistributeMapRes struct {
Items []Item `json:"items"`
Succeeded int `json:"succeeded"`
Failed int `json:"failed"`
}
DistributeMapRes is the ordered result; index i corresponds to input item i regardless of completion order.
type DistributeReduceReq ¶
type DistributeReduceReq struct {
Strategy string `json:"strategy" validate:"required,oneof=collect all_pass any_pass first_success count"`
Items []Item `json:"items" validate:"required"`
}
DistributeReduceReq is the reduce configuration. Strategy selects the fold: collect, all_pass, any_pass, first_success, count.
type DistributeReduceRes ¶
type DistributeReduceRes struct {
Strategy string `json:"strategy"`
AllPass bool `json:"all_pass,omitempty"`
AnyPass bool `json:"any_pass,omitempty"`
Succeeded int `json:"succeeded,omitempty"`
Failed int `json:"failed,omitempty"`
Result any `json:"result,omitempty"`
Results []any `json:"results,omitempty"`
}
DistributeReduceRes carries only the field selected by the strategy.