memory

package
v1.12.5 Latest Latest
Warning

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

Go to latest
Published: Oct 10, 2026 License: Apache-2.0 Imports: 8 Imported by: 0

Documentation

Overview

Package memory provides a shared development broker without disk durability.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Broker

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

Broker can be shared by several development instances.

func New

func New() *Broker

New creates a broker whose state lasts for this process only.

func (*Broker) Capabilities

func (b *Broker) Capabilities() core.Capabilities

Capabilities explicitly exclude disk durability and key ordering.

func (*Broker) Close

func (b *Broker) Close(context.Context) error

Close leaves the shared development broker available to other instances.

func (*Broker) Connect

func (b *Broker) Connect(context.Context) error

Connect needs no external connection.

func (*Broker) EnsureStream

func (b *Broker) EnsureStream(_ context.Context, namespace string, cfg core.StreamConfig) error

EnsureStream refuses incompatible declarations by different instances.

func (*Broker) Health

func (b *Broker) Health(context.Context) error

Health checks the local broker.

func (*Broker) Inspect

func (b *Broker) Inspect(_ context.Context, namespace string, cfg core.StreamConfig) (core.StreamInfo, error)

Inspect reports current process-memory stream state.

func (*Broker) InspectConsumer

func (b *Broker) InspectConsumer(ctx context.Context, binding core.Binding) (core.ConsumerInfo, error)

func (*Broker) ListBackfills

func (b *Broker) ListBackfills(ctx context.Context, identity core.Identity, cursor string, limit int) ([]core.Backfill, string, error)

func (*Broker) ListDeadLetters

func (b *Broker) ListDeadLetters(_ context.Context, identity core.Identity, subscriptionID, cursor string, limit int) ([]core.DeadLetter, string, error)

ListDeadLetters filters by service identity before cursor paging.

func (*Broker) Name

func (b *Broker) Name() string

Name returns the provider type.

func (*Broker) PauseConsumer

func (b *Broker) PauseConsumer(ctx context.Context, binding core.Binding, paused bool) error

func (*Broker) Publish

func (b *Broker) Publish(ctx context.Context, namespace string, cfg core.StreamConfig, msg core.Envelope) (core.Receipt, error)

Publish accepts a message into process memory.

func (*Broker) ReplayDeadLetter

func (b *Broker) ReplayDeadLetter(_ context.Context, identity core.Identity, subscriptionID, id string) (core.Receipt, error)

ReplayDeadLetter atomically queues recovery for the original consumer in memory.

func (*Broker) RequestRPC

func (b *Broker) RequestRPC(ctx context.Context, identity core.Identity, request core.RPCRequest) (core.RPCResponse, error)

RequestRPC chooses one available logical-service replica without persistent storage.

func (*Broker) RunBackfill

func (b *Broker) RunBackfill(ctx context.Context, binding core.Binding, in core.BackfillInput) (core.Backfill, error)

func (*Broker) ServeRPC

func (b *Broker) ServeRPC(ctx context.Context, binding core.RPCBinding, handler func(context.Context, core.RPCRequest) core.RPCResponse) (core.RPCServer, error)

ServeRPC registers a bounded process-local replica for development.

func (*Broker) StoreDeadLetter

func (b *Broker) StoreDeadLetter(_ context.Context, namespace string, letter core.DeadLetter) error

StoreDeadLetter retains a failure before its original delivery is settled.

func (*Broker) Subscribe

func (b *Broker) Subscribe(_ context.Context, binding core.Binding) (core.Subscription, error)

Subscribe attaches a worker to the logical consumer, or creates a broadcast cursor.

Jump to

Keyboard shortcuts

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