nodes_distribute

package
v0.20.3 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 8, 2026 License: Apache-2.0 Imports: 7 Imported by: 0

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

View Source
const ID = "nodes_distribute"

Variables

View Source
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.

View Source
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

func Bundle

func Bundle(_ map[string]string) core.Bundle

func Library

func Library() action.Library

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.

type Item

type Item struct {
	OK     bool   `json:"ok"`
	Result any    `json:"result,omitempty"`
	Error  string `json:"error,omitempty"`
}

Item is one entry in a distribute.map result. Exactly one of Result or Error is populated per item.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL