Documentation
¶
Index ¶
- func Reduce[T any, V any](ctx context.Context, it Iterator[T], reducer func(V, T) (V, error)) (ret V, err error)
- func ToSlice[T any](ctx context.Context, it Iterator[T]) ([]T, error)
- type BidiStream
- type Iterator
- func Aggregate[T any](its ...Iterator[T]) Iterator[T]
- func Empty[T any]() Iterator[T]
- func Error[T any](err error) Iterator[T]
- func Filter[T any](it Iterator[T], fn func(T) bool) Iterator[T]
- func FromDocumentIterator[T any](it document.Iterator) Iterator[T]
- func FromFunc[T any](fn func(context.Context) (T, bool, error), close func() error) Iterator[T]
- func FromRawStream[T any](ctx context.Context, cancel func(), stream RawStream) Iterator[*T]
- func FromSlice[T any](els []T) Iterator[T]
- func FromStream[T any](ctx context.Context, cancel func(), stream Stream[T]) Iterator[*T]
- func FromStreamWithAck[T any, A any](ctx context.Context, cancel func(), stream BidiStream[T, A]) Iterator[*T]
- func FromValueStream[T any](ctx context.Context, cancel func(), stream ValueStream[T]) Iterator[T]
- func IsEmpty[T any](ctx context.Context, i Iterator[T]) (Iterator[T], bool)
- func Map[T any, V any](it Iterator[T], fn func(T) V) Iterator[V]
- func Unslice[T any](it Iterator[[]T]) Iterator[T]
- type RawStream
- type Stream
- type ValueStream
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
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 ¶
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 Error ¶
Error returns an Iterator of T that never returns a value and always returns the given error.
func Filter ¶
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 ¶
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 FromRawStream ¶
FromRawStream works similar to FromStream, but takes a raw grpc.ClientStream, so it's inherently less safe than using FromStream.
func FromStream ¶
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 ¶
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.
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.