Documentation
¶
Index ¶
- func Reduce[T, R any](s Stream[T], fn func(<-chan T) (R, error)) (R, error)
- func SendContext[T any](ctx context.Context, pipe chan<- T, item T) bool
- func SetWorkerErrorWrapper(wrap func(error) error)
- type Group
- type Option
- type Stream
- func Chunk[T any](s Stream[T], n int) Stream[[]T]
- func Concat[T any](s Stream[T], others ...Stream[T]) Stream[T]
- func DistinctBy[T any, K comparable](s Stream[T], fn func(T) K) Stream[T]
- func DistinctByCount[T any, K comparable](s Stream[T], n int, fn func(T) K) Stream[T]
- func DistinctByWindow[T any, K comparable](ctx context.Context, s Stream[T], every time.Duration, fn func(T) K) Stream[T]
- func FlatMap[T, U any](s Stream[T], fn func(T, chan<- U), opts ...Option) Stream[U]
- func FlatMapContext[T, U any](ctx context.Context, s Stream[T], fn func(context.Context, T, chan<- U), ...) Stream[U]
- func FlatMapContextErr[T, U any](ctx context.Context, s Stream[T], fn func(context.Context, T, chan<- U) error, ...) Stream[U]
- func FlatMapErr[T, U any](s Stream[T], fn func(T, chan<- U) error, opts ...Option) Stream[U]
- func FlatStage[I, O any](ctx context.Context, in Stream[I], fn func(context.Context, I, chan<- O), ...) Stream[O]
- func FlatStageErr[I, O any](ctx context.Context, in Stream[I], fn func(context.Context, I, chan<- O) error, ...) Stream[O]
- func From[T any](generate func(chan<- T)) Stream[T]
- func FromChan[T any](source <-chan T) Stream[T]
- func GroupBy[T any, K comparable](s Stream[T], fn func(T) K) Stream[[]T]
- func GroupByCount[T any, K comparable](s Stream[T], n int, fn func(T) K) Stream[Group[K, T]]
- func GroupByWindow[T any, K comparable](ctx context.Context, s Stream[T], every time.Duration, fn func(T) K) Stream[Group[K, T]]
- func Map[T, U any](s Stream[T], fn func(T) U, opts ...Option) Stream[U]
- func MapContext[T, U any](ctx context.Context, s Stream[T], fn func(context.Context, T) U, ...) Stream[U]
- func MapContextErr[T, U any](ctx context.Context, s Stream[T], fn func(context.Context, T) (U, error), ...) Stream[U]
- func MapErr[T, U any](s Stream[T], fn func(T) (U, error), opts ...Option) Stream[U]
- func New[T any](source <-chan T) Stream[T]
- func Stage[I, O any](ctx context.Context, in Stream[I], fn func(context.Context, I) O, ...) Stream[O]
- func StageErr[I, O any](ctx context.Context, in Stream[I], fn func(context.Context, I) (O, error), ...) Stream[O]
- func Tap[T any](ctx context.Context, in Stream[T], fn func(context.Context, T) error, ...) Stream[T]
- func Values[T any](items ...T) Stream[T]
- func WithState[T any](source <-chan T, handle state.Handle) Stream[T]
- func WithStateAndLink[T any](source <-chan T, handle state.Handle, linkMeter *link.Meter) Stream[T]
- func (s Stream[T]) AllMatch(predicate func(T) bool) bool
- func (s Stream[T]) AllMatchErr(predicate func(T) bool) (bool, error)
- func (s Stream[T]) AnyMatch(predicate func(T) bool) bool
- func (s Stream[T]) AnyMatchErr(predicate func(T) bool) (bool, error)
- func (s Stream[T]) Buffer(n int) Stream[T]
- func (s Stream[T]) Collect() []T
- func (s Stream[T]) CollectErr() ([]T, error)
- func (s Stream[T]) Concat(others ...Stream[T]) Stream[T]
- func (s Stream[T]) Count() (count int)
- func (s Stream[T]) CountErr() (int, error)
- func (s Stream[T]) Done()
- func (s Stream[T]) DoneErr() error
- func (s Stream[T]) Err() error
- func (s Stream[T]) Filter(fn func(T) bool, opts ...Option) Stream[T]
- func (s Stream[T]) First() (T, bool)
- func (s Stream[T]) FirstErr() (T, bool, error)
- func (s Stream[T]) ForAll(fn func(<-chan T))
- func (s Stream[T]) ForAllErr(fn func(<-chan T)) error
- func (s Stream[T]) ForEach(fn func(T))
- func (s Stream[T]) ForEachErr(fn func(T)) error
- func (s Stream[T]) Head(n int64) Stream[T]
- func (s Stream[T]) Last() (T, bool)
- func (s Stream[T]) LastErr() (T, bool, error)
- func (s Stream[T]) Max(less func(T, T) bool) (T, bool)
- func (s Stream[T]) MaxErr(less func(T, T) bool) (T, bool, error)
- func (s Stream[T]) Min(less func(T, T) bool) (T, bool)
- func (s Stream[T]) MinErr(less func(T, T) bool) (T, bool, error)
- func (s Stream[T]) NoneMatch(predicate func(T) bool) bool
- func (s Stream[T]) NoneMatchErr(predicate func(T) bool) (bool, error)
- func (s Stream[T]) Parallel(fn func(T), opts ...Option)
- func (s Stream[T]) ParallelErr(fn func(T) error, opts ...Option) error
- func (s Stream[T]) Reverse() Stream[T]
- func (s Stream[T]) Skip(n int64) Stream[T]
- func (s Stream[T]) Sort(less func(T, T) bool) Stream[T]
- func (s Stream[T]) Tail(n int64) Stream[T]
- func (s Stream[T]) Tap(ctx context.Context, fn func(context.Context, T) error, opts ...Option) Stream[T]
- func (s Stream[T]) Through(ctx context.Context, fn func(context.Context, T) T, opts ...Option) Stream[T]
- func (s Stream[T]) ThroughErr(ctx context.Context, fn func(context.Context, T) (T, error), opts ...Option) Stream[T]
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func Reduce ¶
Reduce hands s's source channel to fn, drains any remaining items after fn returns, and joins fn's returned error with the stream error state.
func SendContext ¶
SendContext sends item to pipe unless ctx has already been canceled.
func SetWorkerErrorWrapper ¶
SetWorkerErrorWrapper installs the root-owned worker error wrapper used by concurrent stream operations.
Types ¶
type Group ¶
type Group[K comparable, T any] struct { Key K Items []T }
Group holds one grouping key plus the items assigned to that key.
type Stream ¶
type Stream[T any] struct { // contains filtered or unexported fields }
Stream is a lazy sequence of values backed by a channel plus shared error state that records upstream worker failures.
func Chunk ¶
Chunk groups items into slices of size n, emitting a final short chunk when the source ends.
func Concat ¶
Concat merges s with others and returns a stream that emits items from all inputs as they arrive.
func DistinctBy ¶
func DistinctBy[T any, K comparable](s Stream[T], fn func(T) K) Stream[T]
DistinctBy keeps the first item for each key produced by fn.
func DistinctByCount ¶
func DistinctByCount[T any, K comparable](s Stream[T], n int, fn func(T) K) Stream[T]
DistinctByCount keeps the first item for each key within windows of n input items, resetting the seen-key set after every window.
func DistinctByWindow ¶
func DistinctByWindow[T any, K comparable](ctx context.Context, s Stream[T], every time.Duration, fn func(T) K) Stream[T]
DistinctByWindow keeps the first item for each key within a time window that starts when the first item in that window arrives.
func FlatMapContext ¶
func FlatMapContext[T, U any](ctx context.Context, s Stream[T], fn func(context.Context, T, chan<- U), opts ...Option) Stream[U]
FlatMapContext is FlatMap with a caller-provided context passed into each worker.
func FlatMapContextErr ¶
func FlatMapContextErr[T, U any](ctx context.Context, s Stream[T], fn func(context.Context, T, chan<- U) error, opts ...Option) Stream[U]
FlatMapContextErr is FlatMapErr with a caller-provided context passed into each worker.
func FlatMapErr ¶
FlatMapErr calls fn for each item and records any returned worker error in the stream state.
func FlatStage ¶
func FlatStage[I, O any](ctx context.Context, in Stream[I], fn func(context.Context, I, chan<- O), opts ...Option) Stream[O]
FlatStage applies fn to each item in in and lets fn emit zero or more output values.
func FlatStageErr ¶
func FlatStageErr[I, O any](ctx context.Context, in Stream[I], fn func(context.Context, I, chan<- O) error, opts ...Option) Stream[O]
FlatStageErr applies fn to each item in in, lets fn emit zero or more output values, and records returned worker errors in the stream state.
func From ¶
From adapts a producer callback into a stream. Panics from generate are captured in the stream state and surfaced by terminal operations.
func GroupBy ¶
func GroupBy[T any, K comparable](s Stream[T], fn func(T) K) Stream[[]T]
GroupBy drains s, groups items by fn, and emits groups in first-seen key order.
func GroupByCount ¶
GroupByCount groups items by fn within windows of n input items and emits one Group per key in first-seen order for each window.
func GroupByWindow ¶
func GroupByWindow[T any, K comparable](ctx context.Context, s Stream[T], every time.Duration, fn func(T) K) Stream[Group[K, T]]
GroupByWindow groups items by fn within a time window that starts when the first item in that window arrives and flushes on timer tick or source close.
func MapContext ¶
func MapContext[T, U any](ctx context.Context, s Stream[T], fn func(context.Context, T) U, opts ...Option) Stream[U]
MapContext is Map with a caller-provided context passed into each worker.
func MapContextErr ¶
func MapContextErr[T, U any](ctx context.Context, s Stream[T], fn func(context.Context, T) (U, error), opts ...Option) Stream[U]
MapContextErr is MapErr with a caller-provided context passed into each worker.
func MapErr ¶
MapErr applies fn to each item in s and records any returned error in the stream state.
func Stage ¶
func Stage[I, O any](ctx context.Context, in Stream[I], fn func(context.Context, I) O, opts ...Option) Stream[O]
Stage applies fn to each item in in and emits the mapped values.
func StageErr ¶
func StageErr[I, O any](ctx context.Context, in Stream[I], fn func(context.Context, I) (O, error), opts ...Option) Stream[O]
StageErr applies fn to each item in in and records returned worker errors in the stream state.
func Tap ¶
func Tap[T any](ctx context.Context, in Stream[T], fn func(context.Context, T) error, opts ...Option) Stream[T]
Tap runs fn for each item in in and re-emits the original item when fn succeeds.
func WithStateAndLink ¶ added in v0.1.8
WithStateAndLink wraps source with handle and optional outbound link meter.
func (Stream[T]) AllMatch ¶
AllMatch reports whether every item satisfies predicate. It drains the remainder of the stream after the first mismatch so delayed fail-fast errors can still surface.
func (Stream[T]) AllMatchErr ¶
AllMatchErr reports whether every item satisfies predicate and returns the current error state when it short-circuits.
func (Stream[T]) AnyMatch ¶
AnyMatch reports whether any item satisfies predicate. It drains the remainder of the stream after the first match so delayed fail-fast errors can still surface.
func (Stream[T]) AnyMatchErr ¶
AnyMatchErr reports whether any item satisfies predicate and returns the current error state when it short-circuits.
func (Stream[T]) Buffer ¶
Buffer inserts a channel buffer of size n between s and the returned stream.
func (Stream[T]) Collect ¶
func (s Stream[T]) Collect() []T
Collect drains the stream into a slice and panics on fail-fast errors.
func (Stream[T]) CollectErr ¶
CollectErr drains the stream into a slice and returns the final error state.
func (Stream[T]) Concat ¶
Concat merges s with others while preserving per-stream item order. Items from different input streams may interleave based on runtime scheduling.
func (Stream[T]) CountErr ¶
CountErr drains the stream, returns the item count, and returns the final error state.
func (Stream[T]) Done ¶
func (s Stream[T]) Done()
Done drains the stream and panics if a fail-fast error was recorded.
func (Stream[T]) Err ¶
Err returns the currently accumulated stream error without draining the stream.
func (Stream[T]) First ¶
First returns the first item and then drains the rest of the stream so any delayed fail-fast error is observed before the call returns.
func (Stream[T]) FirstErr ¶
FirstErr returns the first item and the current error state, then drains the remaining source asynchronously.
func (Stream[T]) ForAll ¶
func (s Stream[T]) ForAll(fn func(<-chan T))
ForAll hands the raw source channel to fn, then drains any leftovers and applies fail-fast panic behavior.
func (Stream[T]) ForAllErr ¶
ForAllErr hands the raw source channel to fn, then drains any leftovers and returns the final error state.
func (Stream[T]) ForEach ¶
func (s Stream[T]) ForEach(fn func(T))
ForEach calls fn for every item in the stream and panics if a fail-fast error was recorded.
func (Stream[T]) ForEachErr ¶
ForEachErr calls fn for every item in the stream and returns the final error state.
func (Stream[T]) Head ¶
Head returns a stream containing at most the first n items from s. The upstream source is drained after the head is satisfied so producers can exit.
func (Stream[T]) LastErr ¶
LastErr drains the stream, returns the last item it observed, and returns the final error state.
func (Stream[T]) MaxErr ¶
MaxErr drains the stream, returns the greatest item according to less, and returns the final error state.
func (Stream[T]) MinErr ¶
MinErr drains the stream, returns the least item according to less, and returns the final error state.
func (Stream[T]) NoneMatch ¶
NoneMatch reports whether no item satisfies predicate. It drains the remainder of the stream after the first match so delayed fail-fast errors can still surface.
func (Stream[T]) NoneMatchErr ¶
NoneMatchErr reports whether no item satisfies predicate and returns the current error state when it short-circuits.
func (Stream[T]) Parallel ¶
Parallel applies fn to each item using the same worker machinery as the transform operators and panics on fail-fast errors.
func (Stream[T]) ParallelErr ¶
ParallelErr applies fn to each item using worker options and returns the final error state.
func (Stream[T]) Reverse ¶
Reverse drains s, reverses the collected items, and replays them as a new stream.
func (Stream[T]) Sort ¶
Sort drains s, sorts all items with less, and then replays them as a new stream.
func (Stream[T]) Tap ¶
func (s Stream[T]) Tap(ctx context.Context, fn func(context.Context, T) error, opts ...Option) Stream[T]
Tap runs fn for each item in s and re-emits the original item when fn succeeds.