Documentation ¶
Index ¶
- type ConcurrentQueryStreamServer
- type InMemoryStreamer
- func (s *InMemoryStreamer[Msg]) Close()
- func (s *InMemoryStreamer[Msg]) CloseSend() error
- func (s *InMemoryStreamer[Msg]) Context() context.Context
- func (s *InMemoryStreamer[Msg]) Header() (metadata.MD, error)
- func (s *InMemoryStreamer[Msg]) IsClosed() bool
- func (s *InMemoryStreamer[Msg]) Recv() (Msg, error)
- func (s *InMemoryStreamer[Msg]) RecvMsg(m interface{}) error
- func (s *InMemoryStreamer[Msg]) Send(req Msg) error
- func (s *InMemoryStreamer[Msg]) SendHeader(_ metadata.MD) error
- func (s *InMemoryStreamer[Msg]) SendMsg(m interface{}) error
- func (s *InMemoryStreamer[Msg]) SetHeader(_ metadata.MD) error
- func (s *InMemoryStreamer[Msg]) SetTrailer(_ metadata.MD)
- func (s *InMemoryStreamer[Msg]) Trailer() metadata.MD
- type LocalQueryClient
- type LocalQueryServer
- type MockClientStream
- func (_m *MockClientStream) CloseSend() error
- func (_m *MockClientStream) Context() context.Context
- func (_m *MockClientStream) EXPECT() *MockClientStream_Expecter
- func (_m *MockClientStream) Header() (metadata.MD, error)
- func (_m *MockClientStream) RecvMsg(m interface{}) error
- func (_m *MockClientStream) SendMsg(m interface{}) error
- func (_m *MockClientStream) Trailer() metadata.MD
- type MockClientStream_CloseSend_Call
- type MockClientStream_Context_Call
- func (_c *MockClientStream_Context_Call) Return(_a0 context.Context) *MockClientStream_Context_Call
- func (_c *MockClientStream_Context_Call) Run(run func()) *MockClientStream_Context_Call
- func (_c *MockClientStream_Context_Call) RunAndReturn(run func() context.Context) *MockClientStream_Context_Call
- type MockClientStream_Expecter
- func (_e *MockClientStream_Expecter) CloseSend() *MockClientStream_CloseSend_Call
- func (_e *MockClientStream_Expecter) Context() *MockClientStream_Context_Call
- func (_e *MockClientStream_Expecter) Header() *MockClientStream_Header_Call
- func (_e *MockClientStream_Expecter) RecvMsg(m interface{}) *MockClientStream_RecvMsg_Call
- func (_e *MockClientStream_Expecter) SendMsg(m interface{}) *MockClientStream_SendMsg_Call
- func (_e *MockClientStream_Expecter) Trailer() *MockClientStream_Trailer_Call
- type MockClientStream_Header_Call
- func (_c *MockClientStream_Header_Call) Return(_a0 metadata.MD, _a1 error) *MockClientStream_Header_Call
- func (_c *MockClientStream_Header_Call) Run(run func()) *MockClientStream_Header_Call
- func (_c *MockClientStream_Header_Call) RunAndReturn(run func() (metadata.MD, error)) *MockClientStream_Header_Call
- type MockClientStream_RecvMsg_Call
- func (_c *MockClientStream_RecvMsg_Call) Return(_a0 error) *MockClientStream_RecvMsg_Call
- func (_c *MockClientStream_RecvMsg_Call) Run(run func(m interface{})) *MockClientStream_RecvMsg_Call
- func (_c *MockClientStream_RecvMsg_Call) RunAndReturn(run func(interface{}) error) *MockClientStream_RecvMsg_Call
- type MockClientStream_SendMsg_Call
- func (_c *MockClientStream_SendMsg_Call) Return(_a0 error) *MockClientStream_SendMsg_Call
- func (_c *MockClientStream_SendMsg_Call) Run(run func(m interface{})) *MockClientStream_SendMsg_Call
- func (_c *MockClientStream_SendMsg_Call) RunAndReturn(run func(interface{}) error) *MockClientStream_SendMsg_Call
- type MockClientStream_Trailer_Call
- type QueryStreamClient
- type QueryStreamServer
- type ResultCacheServer
- type RetrieveResultCache
- func (c *RetrieveResultCache) Alloc(result *internalpb.RetrieveResults) bool
- func (c *RetrieveResultCache) Flush() *internalpb.RetrieveResults
- func (c *RetrieveResultCache) IsEmpty() bool
- func (c *RetrieveResultCache) IsFull() bool
- func (c *RetrieveResultCache) Put(result *internalpb.RetrieveResults)
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type ConcurrentQueryStreamServer ¶
type ConcurrentQueryStreamServer struct {
// contains filtered or unexported fields
}
func NewConcurrentQueryStreamServer ¶
func NewConcurrentQueryStreamServer(srv QueryStreamServer) *ConcurrentQueryStreamServer
func (*ConcurrentQueryStreamServer) Context ¶
func (s *ConcurrentQueryStreamServer) Context() context.Context
func (*ConcurrentQueryStreamServer) Send ¶
func (s *ConcurrentQueryStreamServer) Send(result *internalpb.RetrieveResults) error
type InMemoryStreamer ¶
type InMemoryStreamer[Msg any] struct { grpc.ClientStream grpc.ServerStream // contains filtered or unexported fields }
InMemoryStreamer is a utility to wrap in-memory stream methods.
func NewInMemoryStreamer ¶
func NewInMemoryStreamer[Msg any](ctx context.Context, bufferSize int) *InMemoryStreamer[Msg]
func (*InMemoryStreamer[Msg]) Close ¶
func (s *InMemoryStreamer[Msg]) Close()
func (*InMemoryStreamer[Msg]) CloseSend ¶
func (s *InMemoryStreamer[Msg]) CloseSend() error
CloseSend closes the send direction of the stream. It closes the stream when non-nil error is met. It is also not safe to call CloseSend concurrently with SendMsg.
func (*InMemoryStreamer[Msg]) Context ¶
func (s *InMemoryStreamer[Msg]) Context() context.Context
func (*InMemoryStreamer[Msg]) Header ¶
func (s *InMemoryStreamer[Msg]) Header() (metadata.MD, error)
Header returns the header metadata received from the server if there is any. It blocks if the metadata is not ready to read.
func (*InMemoryStreamer[Msg]) IsClosed ¶
func (s *InMemoryStreamer[Msg]) IsClosed() bool
func (*InMemoryStreamer[Msg]) Recv ¶
func (s *InMemoryStreamer[Msg]) Recv() (Msg, error)
func (*InMemoryStreamer[Msg]) RecvMsg ¶
func (s *InMemoryStreamer[Msg]) RecvMsg(m interface{}) error
RecvMsg blocks until it receives a message into m or the stream is done. It returns io.EOF when the client has performed a CloseSend. On any non-EOF error, the stream is aborted and the error contains the RPC status.
It is safe to have a goroutine calling SendMsg and another goroutine calling RecvMsg on the same stream at the same time, but it is not safe to call RecvMsg on the same stream in different goroutines.
func (*InMemoryStreamer[Msg]) Send ¶
func (s *InMemoryStreamer[Msg]) Send(req Msg) error
func (*InMemoryStreamer[Msg]) SendHeader ¶
func (s *InMemoryStreamer[Msg]) SendHeader(_ metadata.MD) error
SendHeader sends the header metadata. The provided md and headers set by SetHeader() will be sent. It fails if called multiple times.
func (*InMemoryStreamer[Msg]) SendMsg ¶
func (s *InMemoryStreamer[Msg]) SendMsg(m interface{}) error
SendMsg sends a message. On error, SendMsg aborts the stream and the error is returned directly.
SendMsg blocks until:
- There is sufficient flow control to schedule m with the transport, or
- The stream is done, or
- The stream breaks.
SendMsg does not wait until the message is received by the client. An untimely stream closure may result in lost messages.
It is safe to have a goroutine calling SendMsg and another goroutine calling RecvMsg on the same stream at the same time, but it is not safe to call SendMsg on the same stream in different goroutines.
It is not safe to modify the message after calling SendMsg. Tracing libraries and stats handlers may use the message lazily.
func (*InMemoryStreamer[Msg]) SetHeader ¶
func (s *InMemoryStreamer[Msg]) SetHeader(_ metadata.MD) error
SetHeader sets the header metadata. It may be called multiple times. When call multiple times, all the provided metadata will be merged. All the metadata will be sent out when one of the following happens:
- ServerStream.SendHeader() is called;
- The first response is sent out;
- An RPC status is sent out (error or success).
func (*InMemoryStreamer[Msg]) SetTrailer ¶
func (s *InMemoryStreamer[Msg]) SetTrailer(_ metadata.MD)
SetTrailer sets the trailer metadata which will be sent with the RPC status. When called more than once, all the provided metadata will be merged.
func (*InMemoryStreamer[Msg]) Trailer ¶
func (s *InMemoryStreamer[Msg]) Trailer() metadata.MD
Trailer returns the trailer metadata from the server, if there is any. It must only be called after stream.CloseAndRecv has returned, or stream.Recv has returned a non-nil error (including io.EOF).
type LocalQueryClient ¶
type LocalQueryClient struct { grpc.ClientStream // contains filtered or unexported fields }
func NewLocalQueryClient ¶
func NewLocalQueryClient(ctx context.Context) *LocalQueryClient
func (*LocalQueryClient) CloseSend ¶
func (s *LocalQueryClient) CloseSend() error
func (*LocalQueryClient) Context ¶
func (s *LocalQueryClient) Context() context.Context
func (*LocalQueryClient) CreateServer ¶
func (s *LocalQueryClient) CreateServer() *LocalQueryServer
func (*LocalQueryClient) Recv ¶
func (s *LocalQueryClient) Recv() (*internalpb.RetrieveResults, error)
func (*LocalQueryClient) RecvMsg ¶
func (s *LocalQueryClient) RecvMsg(m interface{}) error
type LocalQueryServer ¶
type LocalQueryServer struct { grpc.ServerStream // contains filtered or unexported fields }
TODO LOCAL SERVER AND CLIENT FOR STANDALONE ONLY FOR TEST
func (*LocalQueryServer) Context ¶
func (s *LocalQueryServer) Context() context.Context
func (*LocalQueryServer) FinishError ¶
func (s *LocalQueryServer) FinishError() error
func (*LocalQueryServer) FinishSend ¶
func (s *LocalQueryServer) FinishSend(err error) error
func (*LocalQueryServer) Send ¶
func (s *LocalQueryServer) Send(result *internalpb.RetrieveResults) error
type MockClientStream ¶
MockClientStream is an autogenerated mock type for the ClientStream type
func NewMockClientStream ¶
func NewMockClientStream(t interface { mock.TestingT Cleanup(func()) }) *MockClientStream
NewMockClientStream creates a new instance of MockClientStream. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. The first argument is typically a *testing.T value.
func (*MockClientStream) CloseSend ¶
func (_m *MockClientStream) CloseSend() error
CloseSend provides a mock function with given fields:
func (*MockClientStream) Context ¶
func (_m *MockClientStream) Context() context.Context
Context provides a mock function with given fields:
func (*MockClientStream) EXPECT ¶
func (_m *MockClientStream) EXPECT() *MockClientStream_Expecter
func (*MockClientStream) Header ¶
func (_m *MockClientStream) Header() (metadata.MD, error)
Header provides a mock function with given fields:
func (*MockClientStream) RecvMsg ¶
func (_m *MockClientStream) RecvMsg(m interface{}) error
RecvMsg provides a mock function with given fields: m
func (*MockClientStream) SendMsg ¶
func (_m *MockClientStream) SendMsg(m interface{}) error
SendMsg provides a mock function with given fields: m
func (*MockClientStream) Trailer ¶
func (_m *MockClientStream) Trailer() metadata.MD
Trailer provides a mock function with given fields:
type MockClientStream_CloseSend_Call ¶
MockClientStream_CloseSend_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'CloseSend'
func (*MockClientStream_CloseSend_Call) Return ¶
func (_c *MockClientStream_CloseSend_Call) Return(_a0 error) *MockClientStream_CloseSend_Call
func (*MockClientStream_CloseSend_Call) Run ¶
func (_c *MockClientStream_CloseSend_Call) Run(run func()) *MockClientStream_CloseSend_Call
func (*MockClientStream_CloseSend_Call) RunAndReturn ¶
func (_c *MockClientStream_CloseSend_Call) RunAndReturn(run func() error) *MockClientStream_CloseSend_Call
type MockClientStream_Context_Call ¶
MockClientStream_Context_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Context'
func (*MockClientStream_Context_Call) Return ¶
func (_c *MockClientStream_Context_Call) Return(_a0 context.Context) *MockClientStream_Context_Call
func (*MockClientStream_Context_Call) Run ¶
func (_c *MockClientStream_Context_Call) Run(run func()) *MockClientStream_Context_Call
func (*MockClientStream_Context_Call) RunAndReturn ¶
func (_c *MockClientStream_Context_Call) RunAndReturn(run func() context.Context) *MockClientStream_Context_Call
type MockClientStream_Expecter ¶
type MockClientStream_Expecter struct {
// contains filtered or unexported fields
}
func (*MockClientStream_Expecter) CloseSend ¶
func (_e *MockClientStream_Expecter) CloseSend() *MockClientStream_CloseSend_Call
CloseSend is a helper method to define mock.On call
func (*MockClientStream_Expecter) Context ¶
func (_e *MockClientStream_Expecter) Context() *MockClientStream_Context_Call
Context is a helper method to define mock.On call
func (*MockClientStream_Expecter) Header ¶
func (_e *MockClientStream_Expecter) Header() *MockClientStream_Header_Call
Header is a helper method to define mock.On call
func (*MockClientStream_Expecter) RecvMsg ¶
func (_e *MockClientStream_Expecter) RecvMsg(m interface{}) *MockClientStream_RecvMsg_Call
RecvMsg is a helper method to define mock.On call
- m interface{}
func (*MockClientStream_Expecter) SendMsg ¶
func (_e *MockClientStream_Expecter) SendMsg(m interface{}) *MockClientStream_SendMsg_Call
SendMsg is a helper method to define mock.On call
- m interface{}
func (*MockClientStream_Expecter) Trailer ¶
func (_e *MockClientStream_Expecter) Trailer() *MockClientStream_Trailer_Call
Trailer is a helper method to define mock.On call
type MockClientStream_Header_Call ¶
MockClientStream_Header_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Header'
func (*MockClientStream_Header_Call) Return ¶
func (_c *MockClientStream_Header_Call) Return(_a0 metadata.MD, _a1 error) *MockClientStream_Header_Call
func (*MockClientStream_Header_Call) Run ¶
func (_c *MockClientStream_Header_Call) Run(run func()) *MockClientStream_Header_Call
func (*MockClientStream_Header_Call) RunAndReturn ¶
func (_c *MockClientStream_Header_Call) RunAndReturn(run func() (metadata.MD, error)) *MockClientStream_Header_Call
type MockClientStream_RecvMsg_Call ¶
MockClientStream_RecvMsg_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'RecvMsg'
func (*MockClientStream_RecvMsg_Call) Return ¶
func (_c *MockClientStream_RecvMsg_Call) Return(_a0 error) *MockClientStream_RecvMsg_Call
func (*MockClientStream_RecvMsg_Call) Run ¶
func (_c *MockClientStream_RecvMsg_Call) Run(run func(m interface{})) *MockClientStream_RecvMsg_Call
func (*MockClientStream_RecvMsg_Call) RunAndReturn ¶
func (_c *MockClientStream_RecvMsg_Call) RunAndReturn(run func(interface{}) error) *MockClientStream_RecvMsg_Call
type MockClientStream_SendMsg_Call ¶
MockClientStream_SendMsg_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'SendMsg'
func (*MockClientStream_SendMsg_Call) Return ¶
func (_c *MockClientStream_SendMsg_Call) Return(_a0 error) *MockClientStream_SendMsg_Call
func (*MockClientStream_SendMsg_Call) Run ¶
func (_c *MockClientStream_SendMsg_Call) Run(run func(m interface{})) *MockClientStream_SendMsg_Call
func (*MockClientStream_SendMsg_Call) RunAndReturn ¶
func (_c *MockClientStream_SendMsg_Call) RunAndReturn(run func(interface{}) error) *MockClientStream_SendMsg_Call
type MockClientStream_Trailer_Call ¶
MockClientStream_Trailer_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Trailer'
func (*MockClientStream_Trailer_Call) Return ¶
func (_c *MockClientStream_Trailer_Call) Return(_a0 metadata.MD) *MockClientStream_Trailer_Call
func (*MockClientStream_Trailer_Call) Run ¶
func (_c *MockClientStream_Trailer_Call) Run(run func()) *MockClientStream_Trailer_Call
func (*MockClientStream_Trailer_Call) RunAndReturn ¶
func (_c *MockClientStream_Trailer_Call) RunAndReturn(run func() metadata.MD) *MockClientStream_Trailer_Call
type QueryStreamClient ¶
type QueryStreamClient interface { Recv() (*internalpb.RetrieveResults, error) Context() context.Context CloseSend() error }
type QueryStreamServer ¶
type QueryStreamServer interface { Send(*internalpb.RetrieveResults) error Context() context.Context }
type ResultCacheServer ¶
type ResultCacheServer struct {
// contains filtered or unexported fields
}
Merge result by size and time.
func NewResultCacheServer ¶
func NewResultCacheServer(srv QueryStreamServer, cap int, maxMsgSize int) *ResultCacheServer
func (*ResultCacheServer) Context ¶
func (s *ResultCacheServer) Context() context.Context
func (*ResultCacheServer) Flush ¶
func (s *ResultCacheServer) Flush() error
func (*ResultCacheServer) Send ¶
func (s *ResultCacheServer) Send(result *internalpb.RetrieveResults) error
type RetrieveResultCache ¶
type RetrieveResultCache struct {
// contains filtered or unexported fields
}
func (*RetrieveResultCache) Alloc ¶
func (c *RetrieveResultCache) Alloc(result *internalpb.RetrieveResults) bool
func (*RetrieveResultCache) Flush ¶
func (c *RetrieveResultCache) Flush() *internalpb.RetrieveResults
func (*RetrieveResultCache) IsEmpty ¶
func (c *RetrieveResultCache) IsEmpty() bool
func (*RetrieveResultCache) IsFull ¶
func (c *RetrieveResultCache) IsFull() bool
func (*RetrieveResultCache) Put ¶
func (c *RetrieveResultCache) Put(result *internalpb.RetrieveResults)