iterator

package
v1.84.0 Latest Latest
Warning

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

Go to latest
Published: Sep 10, 2026 License: Apache-2.0 Imports: 6 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func Reduce

func Reduce[T any, V any](
	ctx context.Context,
	it Iterator[T],
	reducer func(V, T) (V, error),
) (ret V, err error)

Reduce combines all the elements in an iterator using a binary operation to produce a single value.

func ToSlice

func ToSlice[T any](ctx context.Context, it Iterator[T]) ([]T, error)

ToSlice consumes the given iterator and returns a slice with all of the elements produced.

Types

type BidiStream

type BidiStream[T any, A any] interface {
	Stream[T]
	// Send sends a message of type A. See grpc.ClientStream.SendMsg for
	// more details.
	Send(*A) error
}

BidiStream abstract a GRPC a protoc-generated wrapper of grpc.ClientStream with bidirectional communication capabilities.

type Iterator

type Iterator[T any] interface {
	// Next returns the next element and true or nil and false
	// if there's no more elements in this Iterator. Err should be
	// checked for any errors incurred during the lifecycle of this Iterator.
	// Note that if an error is found, it's up to the implementation as to
	// whether to return false and stop iteration or aggregate errors
	// and return at the end.
	//
	// This method blocks until data is available. Implementations should
	// use the given context's Done channel to know when data is no longer
	// required and so the call should return.
	Next(context.Context) (T, bool)
	// Err returns the first error or an aggreation of the errors
	// encountered by the Iterator.
	Err() error

	io.Closer
}

Iterator provides a convenient interface for iterating over chunks of structured or unstructured data such as a file of newline-delimited lines of text or a set of datastore documents.

func Aggregate

func Aggregate[T any](its ...Iterator[T]) Iterator[T]

Aggregate combines multiple iterators of T into one single iterator of T. If its is nil or empty, this method safely returns an empty iterator.

func Empty

func Empty[T any]() Iterator[T]

Empty returns an empty iterator.

func Error

func Error[T any](err error) Iterator[T]

Error returns an Iterator of T that never returns a value and always returns the given error.

func Filter

func Filter[T any](it Iterator[T], fn func(T) bool) Iterator[T]

Filter wraps an iterator of type T and returns another iterator that will apply fn to each of the elements produced and filter out any elements for which fn returns false.

func FromDocumentIterator

func FromDocumentIterator[T any](it document.Iterator) Iterator[T]

FromDocumentIterator maps a document.Iterator to an Iterator of type T. Items are pulled lazily — nothing is read from the underlying document.Iterator until the consumer calls Next.

func FromFunc

func FromFunc[T any](
	fn func(context.Context) (T, bool, error),
	close func() error,
) Iterator[T]

FromFunc returns an Iterator backed by the provided function.

func FromRawStream

func FromRawStream[T any](
	ctx context.Context, cancel func(), stream RawStream,
) Iterator[*T]

FromRawStream works similar to FromStream, but takes a raw grpc.ClientStream, so it's inherently less safe than using FromStream.

func FromSlice

func FromSlice[T any](els []T) Iterator[T]

FromSlice returns an iterator backed by the given slice of elements.

func FromStream

func FromStream[T any](
	ctx context.Context, cancel func(), stream Stream[T],
) Iterator[*T]

FromStream takes a Stream, tipically genereated via protoc, and returns an Iterator of T. The iterator stops when its Close function is returned or the underlying stream returns io.EOF or other error. Only non io.EOF errors will be surfaced via the returning iterator's Err method.

The given cancel function should be the context canceling function fed to the grpc.ClientStream's constructor. This is to ensure that when the returned iterator's Close method is called, the stream's resources are also cleaned up.

func FromStreamWithAck

func FromStreamWithAck[T any, A any](
	ctx context.Context, cancel func(), stream BidiStream[T, A],
) Iterator[*T]

FromStreamWithAck returns an iterator of T that sends back an ack message of type A for every chunk of T received. See FromRawStream for more details.

func FromValueStream

func FromValueStream[T any](
	ctx context.Context, cancel func(), stream ValueStream[T],
) Iterator[T]

FromValueStream works similar to FromStream, but expects a ValueStream. See ValueStream for more details.

func IsEmpty

func IsEmpty[T any](ctx context.Context, i Iterator[T]) (Iterator[T], bool)

IsEmpty consumes the first element in i and returns true if it is empty or false if not and returns a new iterator that should be used instead of i.

The given iterator is consumed in either case, so its Close method is wrapped with the returned iterator, or if the given iterator is empty, its Close method is called for the caller, so it's safe to override the variable holding the passed iterator with the return value of this function.

func Map

func Map[T any, V any](it Iterator[T], fn func(T) V) Iterator[V]

Map maps an iterator of type T and returns another iterator that will apply fn to each of the elements produced.

func Unslice

func Unslice[T any](it Iterator[[]T]) Iterator[T]

Unslice converts an iterator of slices of T into an iterator of T.

type RawStream

type RawStream interface {
	// RecvMsg blocks until it receives a message into m or the stream is
	// done. It returns io.EOF when the stream completes successfully. On
	// any other error, the stream is aborted and the error contains the RPC
	// status.
	RecvMsg(m interface{}) error
}

RawStream abstract a subset of grpc.ClientStream.

type Stream

type Stream[T any] interface {
	// Recv blocks until it receives a message into m or the stream is
	// done. It returns io.EOF when the stream completes successfully. On
	// any other error, the stream is aborted and the error contains the RPC
	// status.
	Recv() (*T, error)
}

Stream abstract a GRPC protoc-generated wrapper of grpc.ClientStream.

type ValueStream

type ValueStream[T any] interface {
	// Recv blocks until it receives a message into m or the stream is
	// done. It returns io.EOF when the stream completes successfully. On
	// any other error, the stream is aborted and the error contains the RPC
	// status.
	Recv() (T, error)
}

ValueStream abstract an arbitrary value stream. See Stream for differences.

Jump to

Keyboard shortcuts

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