Documentation
¶
Overview ¶
Package streams implements the V1Streams gRPC service: durable, topic-based message streams.
Index ¶
- type Service
- type ServiceImpl
- func (s *ServiceImpl) CancelStreamSessions()
- func (s *ServiceImpl) Cleanup() error
- func (s *ServiceImpl) Publish(ctx context.Context, req *contracts.PublishStreamMessageRequest) (*contracts.PublishStreamMessageResponse, error)
- func (s *ServiceImpl) Subscribe(ctx context.Context, req *contracts.SubscribeStreamRequest, ...) error
- type ServiceOptFunc
- type ServiceOpts
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type Service ¶
type Service interface {
v1connect.V1StreamsHandler
CancelStreamSessions()
Cleanup() error
}
Service is the durable streams gRPC service. CancelStreamSessions lets the engine hang up Subscribe RPCs on shutdown.
func NewService ¶
func NewService(fs ...ServiceOptFunc) (Service, error)
type ServiceImpl ¶
type ServiceImpl struct {
v1connect.UnimplementedV1StreamsHandler
// contains filtered or unexported fields
}
func (*ServiceImpl) CancelStreamSessions ¶
func (s *ServiceImpl) CancelStreamSessions()
CancelStreamSessions hangs up every Subscribe RPC. Safe to call repeatedly.
func (*ServiceImpl) Cleanup ¶
func (s *ServiceImpl) Cleanup() error
Cleanup stops the publisher's workers and the wake subscription.
func (*ServiceImpl) Publish ¶
func (s *ServiceImpl) Publish(ctx context.Context, req *contracts.PublishStreamMessageRequest) (*contracts.PublishStreamMessageResponse, error)
func (*ServiceImpl) Subscribe ¶
func (s *ServiceImpl) Subscribe(ctx context.Context, req *contracts.SubscribeStreamRequest, connectStream *connect.ServerStream[contracts.StreamMessage]) error
type ServiceOptFunc ¶
type ServiceOptFunc func(*ServiceOpts)
func WithLogger ¶
func WithLogger(l *zerolog.Logger) ServiceOptFunc
func WithPubSub ¶
func WithPubSub(pubsub msgqueue.PubSub) ServiceOptFunc
func WithRepositoryV1 ¶
func WithRepositoryV1(r v1.Repository) ServiceOptFunc
type ServiceOpts ¶
type ServiceOpts struct {
// contains filtered or unexported fields
}
Click to show internal directories.
Click to hide internal directories.