Documentation
¶
Overview ¶
Package sink runs configured sinks. CAN/file/null sinks push-confirm each message; HTTP/TCP sinks broadcast to connected clients, with replay served straight from connector queues (SSE/WS only).
Index ¶
Constants ¶
This section is empty.
Variables ¶
var ErrSkip = errors.New("sink skipped message")
Functions ¶
This section is empty.
Types ¶
type BroadcastReport ¶
type Broadcaster ¶
type Broadcaster interface {
Broadcast(entries []queue.Entry) BroadcastReport
}
Broadcaster sinks fan out to connected clients without confirmation.
type ConnectorRegistrar ¶
type ConnectorRegistrar interface {
RegisterConnector(id string, r ReplayReader)
UnregisterConnector(id string)
}
ConnectorRegistrar lets connectors attach their queue for client replay.
type DataServer ¶
type DataServer struct {
// contains filtered or unexported fields
}
DataServer hosts serve-mode sink endpoints with dynamically managed routes.
func NewDataServer ¶
func NewDataServer(addr string, log *slog.Logger) *DataServer
func (*DataServer) Addr ¶
func (d *DataServer) Addr() string
func (*DataServer) RemoveRoute ¶
func (d *DataServer) RemoveRoute(path string)
func (*DataServer) Start ¶
func (d *DataServer) Start() error
type DeliveryClass ¶
type DeliveryClass string
const ( DeliveryConfirmed DeliveryClass = "confirmed" DeliveryResumable DeliveryClass = "resumable" DeliveryBestEffort DeliveryClass = "best_effort" DeliveryObserved DeliveryClass = "observe_only" )
func DeliveryClassOf ¶
func DeliveryClassOf(r Runtime) DeliveryClass
type DeliveryClassifier ¶
type DeliveryClassifier interface {
DeliveryClass() DeliveryClass
}
DeliveryClassifier makes the route boundary explicit. Implementations that do not provide it are classified from their Pusher/Broadcaster seam by DeliveryClassOf.
type Pusher ¶
Pusher sinks confirm each delivery (CAN). ErrSkip means "cannot carry this message, count it and move on" (e.g. envelope without raw bytes).
type ReplayReader ¶
type ReplayReader interface {
Read(ctx context.Context, after int64, limit int) ([]queue.Entry, error)
}
ReplayReader is what serve-mode sinks use to replay history for a client.