streams

package
v0.110.1 Latest Latest
Warning

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

Go to latest
Published: Oct 2, 2026 License: MIT Imports: 21 Imported by: 0

Documentation

Overview

Package streams implements the V1Streams gRPC service: durable, topic-based message streams.

Index

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) Subscribe

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
}

Jump to

Keyboard shortcuts

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