Documentation
¶
Overview ¶
Package memory provides a shared development broker without disk durability.
Index ¶
- type Broker
- func (b *Broker) Capabilities() core.Capabilities
- func (b *Broker) Close(context.Context) error
- func (b *Broker) Connect(context.Context) error
- func (b *Broker) EnsureStream(_ context.Context, namespace string, cfg core.StreamConfig) error
- func (b *Broker) Health(context.Context) error
- func (b *Broker) Inspect(_ context.Context, namespace string, cfg core.StreamConfig) (core.StreamInfo, error)
- func (b *Broker) InspectConsumer(ctx context.Context, binding core.Binding) (core.ConsumerInfo, error)
- func (b *Broker) ListBackfills(ctx context.Context, identity core.Identity, cursor string, limit int) ([]core.Backfill, string, error)
- func (b *Broker) ListDeadLetters(_ context.Context, identity core.Identity, subscriptionID, cursor string, ...) ([]core.DeadLetter, string, error)
- func (b *Broker) Name() string
- func (b *Broker) PauseConsumer(ctx context.Context, binding core.Binding, paused bool) error
- func (b *Broker) Publish(ctx context.Context, namespace string, cfg core.StreamConfig, ...) (core.Receipt, error)
- func (b *Broker) ReplayDeadLetter(_ context.Context, identity core.Identity, subscriptionID, id string) (core.Receipt, error)
- func (b *Broker) RequestRPC(ctx context.Context, identity core.Identity, request core.RPCRequest) (core.RPCResponse, error)
- func (b *Broker) RunBackfill(ctx context.Context, binding core.Binding, in core.BackfillInput) (core.Backfill, error)
- func (b *Broker) ServeRPC(ctx context.Context, binding core.RPCBinding, ...) (core.RPCServer, error)
- func (b *Broker) StoreDeadLetter(_ context.Context, namespace string, letter core.DeadLetter) error
- func (b *Broker) Subscribe(_ context.Context, binding core.Binding) (core.Subscription, error)
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 (*Broker) Capabilities ¶
func (b *Broker) Capabilities() core.Capabilities
Capabilities explicitly exclude disk durability and key ordering.
func (*Broker) EnsureStream ¶
EnsureStream refuses incompatible declarations by different instances.
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 (*Broker) ListBackfills ¶
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) PauseConsumer ¶
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 (*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 ¶
StoreDeadLetter retains a failure before its original delivery is settled.