eventstream

package
v0.1.66 Latest Latest
Warning

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

Go to latest
Published: Sep 8, 2026 License: AGPL-3.0 Imports: 5 Imported by: 0

Documentation

Overview

Package eventstream provides generic primitives for bounded event streaming. It deliberately contains no provider, workflow, REST, or transport logic.

Index

Constants

This section is empty.

Variables

View Source
var (
	ErrCompleted     = errors.New("stream already completed")
	ErrLimitExceeded = errors.New("stream output limit exceeded")
	ErrInvalidStream = errors.New("invalid output stream")
)

Functions

This section is empty.

Types

type Channel

type Channel string

Channel identifies a logical event channel. The package does not prescribe channel names; consumers may use stdout/stderr, progress/log, or any other workflow-specific vocabulary.

type Completion

type Completion struct {
	Attributes map[string]string
	Err        error
}

Completion is the terminal outcome supplied with a Complete event. Attributes carry consumer-defined result metadata without coupling this package to a particular workflow.

type Event

type Event struct {
	Kind       EventKind
	Channel    Channel
	Payload    []byte
	Completion *Completion
}

Event is a provider-neutral stream event. Payload is owned by the event and must not be mutated by a sink after Send returns.

type EventKind

type EventKind uint8

EventKind identifies the lifecycle event represented by an Event.

const (
	Data EventKind = iota + 1
	Complete
)

type Limits

type Limits struct {
	MaxChunkBytes int
	MaxTotalBytes int64
}

Limits bounds each emitted payload and the aggregate data sent by a Publisher. A zero limit means unlimited for that dimension.

type Publisher

type Publisher struct {
	// contains filtered or unexported fields
}

Publisher validates, bounds, and forwards stream events to a Sink.

func NewPublisher

func NewPublisher(sink Sink, limits Limits) *Publisher

NewPublisher creates a bounded publisher. A nil sink is rejected when the first event is emitted so construction remains convenient for composition.

func (*Publisher) Complete

func (p *Publisher) Complete(ctx context.Context, result Completion) error

Complete emits the terminal event exactly once.

func (*Publisher) Write

func (p *Publisher) Write(ctx context.Context, channel Channel, data []byte) error

Write emits data, splitting it into bounded chunks when needed.

type Sink

type Sink interface {
	Send(context.Context, Event) error
}

Sink receives events. Implementations may block to apply backpressure; the publisher propagates that behavior and observes context cancellation.

Jump to

Keyboard shortcuts

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