core

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: 15 Imported by: 0

Documentation

Overview

Package core defines provider contracts and the service communication runtime.

Index

Constants

This section is empty.

Variables

View Source
var (
	ErrNotRunning     = errors.New("conduit: runtime not running")
	ErrUnsupported    = errors.New("conduit: provider does not support requested capability")
	ErrOutcomeUnknown = errors.New("conduit: publish outcome unknown")
	ErrNotFound       = errors.New("conduit: resource not found")
	ErrConflict       = errors.New("conduit: configuration conflict")
)

Functions

func IsPermanent

func IsPermanent(err error) bool

IsPermanent reports whether an error explicitly ends retries.

func Matches

func Matches(pattern, topic string) bool

Matches compares an event type with a stream subject pattern.

func NewID

func NewID() string

NewID generates an identifier independent of broker sequence numbers.

func Permanent

func Permanent(err error) error

Permanent routes an invalid or terminal message to the dead letter store.

func ValidateMessageID

func ValidateMessageID(id string) error

ValidateMessageID accepts stable IDs that can also address persisted recovery records.

func ValidateTopic

func ValidateTopic(topic string, wildcard bool) error

ValidateTopic accepts exact event types and a trailing wildcard in stream subjects.

Types

type Backfill

type Backfill struct {
	MessageType string        `json:"messageType"`
	Input       BackfillInput `json:"input"`
	ConsumerID  string        `json:"consumerID"`
	Stream      string        `json:"stream"`
	Provider    string        `json:"provider"`
	State       string        `json:"state"`
	Next        uint64        `json:"next"`
	Published   uint64        `json:"published"`
	Skipped     uint64        `json:"skipped"`
	UpdatedAt   time.Time     `json:"updatedAt"`
	Error       string        `json:"error,omitempty"`
	Persisted   bool          `json:"persisted"`
}

Backfill records resumable progress without exposing event data.

func ExecuteBackfill

func ExecuteBackfill(ctx context.Context, binding Binding, job Backfill, read func(context.Context, uint64) (Envelope, error), publish func(context.Context, Envelope, string) error, save func(context.Context, Backfill) error) (Backfill, error)

ExecuteBackfill processes a claimed range, checkpointing after every confirmed publication. The provider must hold an exclusive lease longer than the bounded execution deadline.

type BackfillInput

type BackfillInput struct {
	ID           string `json:"id"`
	Subscription string `json:"subscription"`
	Start        uint64 `json:"start"`
	End          uint64 `json:"end"`
}

BackfillInput selects at most 100 retained stream sequences and a stable operation ID.

func (BackfillInput) Validate

func (in BackfillInput) Validate() error

type Binding

type Binding struct {
	Identity     Identity
	Subscription SubscriptionConfig
	Stream       StreamConfig
}

Binding adds the running member without changing logical consumer identity.

func (Binding) ConsumerID

func (b Binding) ConsumerID() string

ConsumerID excludes the instance for competing consumers.

type Capabilities

type Capabilities struct {
	ConsumerControls bool `json:"consumerControls"`
	Backfill         bool `json:"backfill"`
	RPC              bool `json:"rpc"`
	Durable          bool `json:"durable"`
	Replay           bool `json:"replay"`
	DeadLetters      bool `json:"deadLetters"`
	KeyOrdering      bool `json:"keyOrdering"`
}

Capabilities report guarantees implemented by a provider.

type Config

type Config struct {
	RPC             RPCConfig                     `json:"rpc"             yaml:"rpc"`
	Identity        Identity                      `json:"identity"        yaml:"identity"`
	Streams         map[string]StreamConfig       `json:"streams"         yaml:"streams"`
	Subscriptions   map[string]SubscriptionConfig `json:"subscriptions"   yaml:"subscriptions"`
	ObserverBuffer  int                           `json:"observerBuffer"  yaml:"observer_buffer"`
	MaxPayloadBytes int                           `json:"maxPayloadBytes" yaml:"max_payload_bytes"`
	Providers       map[string]ConnectionConfig   `json:"providers"       yaml:"providers"`
	Discovery       string                        `json:"discovery"       yaml:"discovery"`
	Version         string                        `json:"version"         yaml:"version"`
	Endpoints       []Endpoint                    `json:"endpoints"       yaml:"endpoints"`
}

Config separates stream topology, logical subscriptions and replica settings.

type ConnectionConfig

type ConnectionConfig struct {
	Type               string `json:"type"               yaml:"type"`
	URL                string `json:"url"                yaml:"url"`
	DeadLetterReplicas int    `json:"deadLetterReplicas" yaml:"dead_letter_replicas"`
}

ConnectionConfig is private connection configuration, never dashboard metadata.

type ConsumerInfo

type ConsumerInfo struct {
	Subscription SubscriptionConfig `json:"subscription"`
	Provider     string             `json:"provider"`
	ConsumerID   string             `json:"consumerID"`
	Pending      uint64             `json:"pending"`
	AckPending   uint64             `json:"ackPending"`
	Redelivered  uint64             `json:"redelivered"`
	Paused       bool               `json:"paused"`
	Processing   Latency            `json:"processing"`
	Delivery     Latency            `json:"delivery"`
}

ConsumerInfo separates broker backlog from instance-local processing latency.

type DeadLetter

type DeadLetter struct {
	ID       string       `json:"id"`
	Message  Envelope     `json:"message"`
	Delivery DeliveryInfo `json:"delivery"`
	FailedAt time.Time    `json:"failedAt"`
	Reason   string       `json:"reason"`
	Replayed bool         `json:"replayed"`
}

DeadLetter retains the failure and the original message for controlled recovery.

type Delivery

type Delivery interface {
	Message() Envelope
	Info() DeliveryInfo
	Ack(ctx context.Context) error
	Retry(ctx context.Context, delay time.Duration) error
	Reject(ctx context.Context) error
}

Delivery owns settlement for one broker delivery.

type DeliveryInfo

type DeliveryInfo struct {
	Stream         string       `json:"stream"`
	ConsumerID     string       `json:"consumerID"`
	SubscriptionID string       `json:"subscriptionID"`
	Destination    Identity     `json:"destination"`
	Mode           DeliveryMode `json:"mode"`
	Attempt        uint64       `json:"attempt"`
	Sequence       uint64       `json:"sequence"`
}

DeliveryInfo identifies this attempt independently of its original message.

type DeliveryMode

type DeliveryMode string

DeliveryMode controls replica distribution for one subscription.

const (
	Competing DeliveryMode = "competing"
	Broadcast DeliveryMode = "broadcast"
)

type Endpoint

type Endpoint struct {
	Protocol string `json:"protocol" yaml:"protocol"`
	URL      string `json:"url"      yaml:"url"`
}

Endpoint is an advertised address, separate from the instance's bind address.

func (Endpoint) Validate

func (e Endpoint) Validate() error

Validate accepts transport addresses without embedded credentials or query tokens.

type Envelope

type Envelope struct {
	ID            string            `json:"id"`
	Type          string            `json:"type"`
	Source        Identity          `json:"source"`
	ContentType   string            `json:"contentType"`
	Data          json.RawMessage   `json:"data"`
	CreatedAt     time.Time         `json:"createdAt"`
	Key           string            `json:"key,omitempty"`
	CorrelationID string            `json:"correlationID,omitempty"`
	CausationID   string            `json:"causationID,omitempty"`
	Headers       map[string]string `json:"headers,omitempty"`
	// TargetConsumer directs recovery to the original subscription.
	TargetConsumer string `json:"targetConsumer,omitempty"`
}

Envelope is immutable after broker acceptance. ID survives redelivery.

func (Envelope) Clone

func (e Envelope) Clone() Envelope

Clone prevents providers and observer hooks from sharing mutable payloads.

type HandleHook

type HandleHook interface {
	Hook
	BeforeHandle(ctx context.Context, message Envelope, info DeliveryInfo) error
}

HandleHook validates a delivery before the handler executes.

type Handler

type Handler func(context.Context, Envelope, DeliveryInfo) error

Handler processes one immutable envelope and its delivery context.

type Hook

type Hook interface{ Name() string }

Hook identifies one optional set of callbacks.

type HookEvent

type HookEvent struct {
	Duration time.Duration `json:"duration,omitempty"`
	Stage    Stage         `json:"stage"`
	Identity Identity      `json:"identity"`
	Provider string        `json:"provider,omitempty"`
	Message  *Envelope     `json:"message,omitempty"`
	Delivery *DeliveryInfo `json:"delivery,omitempty"`
	Receipt  *Receipt      `json:"receipt,omitempty"`
	Error    string        `json:"error,omitempty"`
	At       time.Time     `json:"at"`
}

HookEvent is copied before observers receive it.

type HookFuncs

type HookFuncs struct {
	HookName  string
	OnEvent   func(context.Context, HookEvent)
	OnPublish func(context.Context, *Envelope) error
	OnHandle  func(context.Context, Envelope, DeliveryInfo) error
}

HookFuncs adapts callbacks without requiring unrelated lifecycle methods.

func (HookFuncs) BeforeHandle

func (h HookFuncs) BeforeHandle(ctx context.Context, msg Envelope, info DeliveryInfo) error

BeforeHandle runs the configured control callback.

func (HookFuncs) BeforePublish

func (h HookFuncs) BeforePublish(ctx context.Context, msg *Envelope) error

BeforePublish runs the configured control callback.

func (HookFuncs) Name

func (h HookFuncs) Name() string

Name returns the registration identity.

func (HookFuncs) Observe

func (h HookFuncs) Observe(ctx context.Context, event HookEvent)

Observe receives an immutable event.

type Identity

type Identity struct {
	Namespace  string `json:"namespace"  yaml:"namespace"`
	ServiceID  string `json:"serviceID"  yaml:"service_id"`
	InstanceID string `json:"instanceID" yaml:"instance_id"`
}

Identity separates a logical service from a running replica.

func (Identity) Validate

func (i Identity) Validate() error

Validate checks that every identity component is present.

type Instance

type Instance struct {
	Identity  Identity   `json:"identity"`
	Version   string     `json:"version"`
	Endpoints []Endpoint `json:"endpoints"`
	Ready     bool       `json:"ready"`
}

Instance describes one running member of a logical service.

func (Instance) Validate

func (i Instance) Validate() error

Validate checks registration identity and advertised transport addresses.

type InstanceLister

type InstanceLister interface {
	List(ctx context.Context, namespace string) ([]Instance, error)
}

InstanceLister lists members within the configured namespace for dashboard inspection.

type Latency

type Latency struct {
	Count   uint64        `json:"count"`
	Total   time.Duration `json:"total"`
	Max     time.Duration `json:"max"`
	Average time.Duration `json:"average"`
}

Latency summarizes processing attempts on this instance, in nanoseconds.

type Management

type Management interface {
	Inspect(ctx context.Context, namespace string, config StreamConfig) (StreamInfo, error)
	StoreDeadLetter(ctx context.Context, namespace string, letter DeadLetter) error
	ListDeadLetters(ctx context.Context, identity Identity, subscription string, cursor string, limit int) ([]DeadLetter, string, error)
	ReplayDeadLetter(ctx context.Context, identity Identity, subscription string, id string) (Receipt, error)
}

Management exposes persisted failures and stream inspection without fake fallbacks.

type Middleware

type Middleware func(Handler) Handler

Middleware surrounds each processing attempt.

type Observer

type Observer interface {
	Hook
	Observe(ctx context.Context, event HookEvent)
}

Observer receives outcomes and cannot change delivery guarantees.

type Operations

type Operations interface {
	InspectConsumer(ctx context.Context, binding Binding) (ConsumerInfo, error)
	PauseConsumer(ctx context.Context, binding Binding, paused bool) error
	RunBackfill(ctx context.Context, binding Binding, input BackfillInput) (Backfill, error)
	ListBackfills(ctx context.Context, identity Identity, cursor string, limit int) ([]Backfill, string, error)
}

Operations is optional broker-backed consumer management and controlled recovery.

type Option

type Option func(*Runtime) error

Option supplies a provider or an extension point.

func WithConfig

func WithConfig(cfg Config) Option

WithConfig supplies explicit Forge configuration, overriding loaded non-zero values.

func WithHooks

func WithHooks(hooks ...Hook) Option

WithHooks registers optional callbacks in deterministic registration order.

func WithProvider

func WithProvider(name string, provider Provider) Option

WithProvider registers a named connection independently of message routing.

func WithRegistry

func WithRegistry(registry Registry) Option

WithRegistry enables registration, discovery and instance inspection.

type Provider

type Provider interface {
	Name() string
	Capabilities() Capabilities
	Connect(ctx context.Context) error
	Close(ctx context.Context) error
	Health(ctx context.Context) error
	EnsureStream(ctx context.Context, namespace string, config StreamConfig) error
	Publish(ctx context.Context, namespace string, config StreamConfig, message Envelope) (Receipt, error)
	Subscribe(ctx context.Context, binding Binding) (Subscription, error)
}

Provider implements retained messaging. Optional management has its own interface.

type ProviderInfo

type ProviderInfo struct {
	Name         string       `json:"name"`
	Type         string       `json:"type"`
	Capabilities Capabilities `json:"capabilities"`
	Healthy      bool         `json:"healthy"`
}

ProviderInfo omits connection strings and credentials from dashboard responses.

type PublishHook

type PublishHook interface {
	Hook
	BeforePublish(ctx context.Context, message *Envelope) error
}

PublishHook can enrich or reject a draft before final validation and serialization.

type RPCBinding

type RPCBinding struct {
	Identity        Identity
	Method          string
	Config          RPCConfig
	MaxPayloadBytes int
}

type RPCCode

type RPCCode string

RPCCode is a public error category independent of a broker.

const (
	RPCBadRequest       RPCCode = "BAD_REQUEST"
	RPCNotFound         RPCCode = "NOT_FOUND"
	RPCPermissionDenied RPCCode = "PERMISSION_DENIED"
	RPCConflict         RPCCode = "CONFLICT"
	RPCUnavailable      RPCCode = "UNAVAILABLE"
	RPCDeadlineExceeded RPCCode = "DEADLINE_EXCEEDED"
	RPCCanceled         RPCCode = "CANCELED"
	RPCInternal         RPCCode = "INTERNAL"
)

type RPCConfig

type RPCConfig struct {
	Provider    string        `json:"provider"    yaml:"provider"`
	Timeout     time.Duration `json:"timeout"     yaml:"timeout"`
	Concurrency int           `json:"concurrency" yaml:"concurrency"`
	MaxInFlight int           `json:"maxInFlight" yaml:"max_in_flight"`
}

RPCConfig bounds each instance's transient request/reply work.

type RPCError

type RPCError struct {
	Code    RPCCode `json:"code"`
	Message string  `json:"message"`
}

RPCError exposes only an intentional public error message.

func PublicRPCError

func PublicRPCError(err error) *RPCError

PublicRPCError strips arbitrary handler causes from responses.

func (*RPCError) Error

func (e *RPCError) Error() string

func (*RPCError) Is

func (e *RPCError) Is(target error) bool

type RPCHandler

type RPCHandler func(context.Context, Envelope) (json.RawMessage, error)

RPCHandler returns serialized typed data or an intentional public error.

type RPCProvider

type RPCProvider interface {
	RequestRPC(ctx context.Context, identity Identity, request RPCRequest) (RPCResponse, error)
	ServeRPC(ctx context.Context, binding RPCBinding, handler func(context.Context, RPCRequest) RPCResponse) (RPCServer, error)
}

RPCProvider supplies optional transient RPC without weakening durable event guarantees.

type RPCRequest

type RPCRequest struct {
	Service  string    `json:"service"`
	Method   string    `json:"method"`
	Envelope Envelope  `json:"envelope"`
	Deadline time.Time `json:"deadline"`
}

RPCRequest preserves caller identity, tracing metadata and the remote deadline.

type RPCResponse

type RPCResponse struct {
	ID     string          `json:"id"`
	Source Identity        `json:"source"`
	Data   json.RawMessage `json:"data,omitempty"`
	Error  *RPCError       `json:"error,omitempty"`
}

RPCResponse correlates one reply and identifies the responding replica.

type RPCServer

type RPCServer interface {
	Close(ctx context.Context) error
}

RPCServer stops intake and drains accepted work until the shutdown deadline.

type Receipt

type Receipt struct {
	MessageID string `json:"messageID"`
	Sequence  uint64 `json:"sequence"`
	Persisted bool   `json:"persisted"`
	Duplicate bool   `json:"duplicate"`
}

Receipt distinguishes broker acceptance from processing success.

type Registry

type Registry interface {
	Resolver
	Register(ctx context.Context, instance Instance) error
	Deregister(ctx context.Context, identity Identity) error
}

Registry adds instance registration without requiring it from DNS resolvers.

type Resolver

type Resolver interface {
	Resolve(ctx context.Context, namespace string, service string) ([]Instance, error)
}

Resolver finds replicas by logical service identity within a namespace.

type Runtime

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

Runtime owns one instance's connections and workers, never the service identity.

func New

func New(cfg Config, options ...Option) (*Runtime, error)

New constructs a validated runtime. Instance ID can be generated per process.

func NewDeferred

func NewDeferred(options ...Option) (*Runtime, error)

NewDeferred allows handler binding before Forge loads application configuration.

func (*Runtime) Backfill

func (r *Runtime) Backfill(ctx context.Context, in BackfillInput) (Backfill, error)

Backfill republishes only to the selected original logical consumer.

func (*Runtime) Backfills

func (r *Runtime) Backfills(ctx context.Context, provider, cursor string, limit int) ([]Backfill, string, error)

Backfills returns a bounded service-scoped operation history for one provider.

func (*Runtime) Bind

func (r *Runtime) Bind(id, messageType string, handler Handler, middlewares ...Middleware) error

Bind installs one logical handler before startup.

func (*Runtime) BindRPC

func (r *Runtime) BindRPC(method string, handler RPCHandler) error

BindRPC registers one logical procedure before startup; replicas share its queue group.

func (*Runtime) CallRPC

func (r *Runtime) CallRPC(ctx context.Context, service, method string, draft Envelope, validate func(Envelope) error) (RPCResponse, error)

CallRPC sends once by service name. A canceled call may already have caused an effect.

func (*Runtime) Configuration

func (r *Runtime) Configuration() Config

Configuration returns a copy for Forge initialization. It may contain credentials.

func (*Runtime) Configure

func (r *Runtime) Configure(cfg Config, extra ...Option) error

Configure validates loaded topology and applies deferred registrations atomically.

func (*Runtime) Consumers

func (r *Runtime) Consumers(ctx context.Context) ([]ConsumerInfo, error)

Consumers returns broker progress for this service's registered subscriptions.

func (*Runtime) DeadLetters

func (r *Runtime) DeadLetters(ctx context.Context, providerName, subscription, cursor string, limit int) ([]DeadLetter, string, error)

DeadLetters reads only failures owned by this logical service.

func (*Runtime) DurableStream

func (r *Runtime) DurableStream(name string) bool

DurableStream checks a provider before transactional publication.

func (*Runtime) Health

func (r *Runtime) Health(ctx context.Context) error

Health checks provider availability and startup state.

func (*Runtime) Identity

func (r *Runtime) Identity() Identity

Identity returns the service and replica identity for this runtime.

func (*Runtime) Instances

func (r *Runtime) Instances(ctx context.Context) ([]Instance, error)

Instances lists visible members inside this runtime's namespace.

func (*Runtime) PauseSubscription

func (r *Runtime) PauseSubscription(ctx context.Context, id string, paused bool) error

PauseSubscription changes the actual logical cursor across its attached replicas.

func (*Runtime) Prepare

func (r *Runtime) Prepare(ctx context.Context, streamName string, draft Envelope, validate func(Envelope) error) (Envelope, error)

Prepare enriches and validates an immutable draft before publication or transactional enqueue.

func (*Runtime) Providers

func (r *Runtime) Providers() map[string]Provider

Providers returns registered provider instances for initialization adapters.

func (*Runtime) Publish

func (r *Runtime) Publish(ctx context.Context, streamName string, draft Envelope, validate func(Envelope) error) (Receipt, error)

Publish prepares a draft and waits for provider acceptance.

func (*Runtime) RecentEvents

func (r *Runtime) RecentEvents() []HookEvent

RecentEvents returns bounded diagnostics without message data, headers or error text.

func (*Runtime) RegisterHook

func (r *Runtime) RegisterHook(hook Hook) error

RegisterHook adds an extension hook before startup.

func (*Runtime) ReplayDeadLetter

func (r *Runtime) ReplayDeadLetter(ctx context.Context, providerName, subscription, id string) (Receipt, error)

ReplayDeadLetter republishes one retained failure without changing its message ID.

func (*Runtime) Resolver

func (r *Runtime) Resolver() Resolver

Resolver exposes the configured service discovery provider.

func (*Runtime) Send

func (r *Runtime) Send(ctx context.Context, streamName string, msg Envelope) (Receipt, error)

Send publishes an already prepared envelope, preserving its ID across outbox retries.

func (*Runtime) SetRegistry

func (r *Runtime) SetRegistry(registry Registry) error

SetRegistry fills the discovery adapter before startup.

func (*Runtime) Snapshot

func (r *Runtime) Snapshot(ctx context.Context) (Snapshot, error)

Snapshot reports unavailable inspection as an error, never an empty healthy list.

func (*Runtime) Start

func (r *Runtime) Start(ctx context.Context) error

Start connects providers, reconciles topology and then launches handlers.

func (*Runtime) Stop

func (r *Runtime) Stop(ctx context.Context) error

Stop cancels intake and waits for handlers before disconnecting providers.

func (*Runtime) StreamFor

func (r *Runtime) StreamFor(messageType string) (StreamConfig, error)

StreamFor returns an unambiguous message route. Overlapping routes require an explicit stream.

type Snapshot

type Snapshot struct {
	Identity      Identity             `json:"identity"`
	Running       bool                 `json:"running"`
	StartedAt     time.Time            `json:"startedAt"`
	Providers     []ProviderInfo       `json:"providers"`
	Streams       []StreamInfo         `json:"streams"`
	Subscriptions []SubscriptionConfig `json:"subscriptions"`
	Published     uint64               `json:"published"`
	Handled       uint64               `json:"handled"`
	Acknowledged  uint64               `json:"acknowledged"`
	Failed        uint64               `json:"failed"`
	Retried       uint64               `json:"retried"`
	DeadLettered  uint64               `json:"deadLettered"`
	ObserverDrops uint64               `json:"observerDrops"`
	RPCCalls      uint64               `json:"rpcCalls"`
	RPCHandled    uint64               `json:"rpcHandled"`
	RPCFailed     uint64               `json:"rpcFailed"`
	RPCTimedOut   uint64               `json:"rpcTimedOut"`
}

Snapshot contains per-instance counters and broker stream state.

type Stage

type Stage string

Stage names observable outcomes without treating handler success as acknowledgement.

const (
	RPCCalling          Stage = "rpc.calling"
	RPCReceived         Stage = "rpc.received"
	RPCHandling         Stage = "rpc.handling"
	RPCHandled          Stage = "rpc.handled"
	RPCReturned         Stage = "rpc.returned"
	RPCFailed           Stage = "rpc.failed"
	Starting            Stage = "starting"
	Ready               Stage = "ready"
	Draining            Stage = "draining"
	Stopped             Stage = "stopped"
	ProviderConnected   Stage = "provider.connected"
	SubscriptionStarted Stage = "subscription.started"
	SubscriptionFailed  Stage = "subscription.failed"
	Publishing          Stage = "publishing"
	Published           Stage = "published"
	PublishFailed       Stage = "publish.failed"
	Received            Stage = "received"
	Handling            Stage = "handling"
	Handled             Stage = "handled"
	HandleFailed        Stage = "handle.failed"
	Acknowledged        Stage = "acknowledged"
	SettlementFailed    Stage = "settlement.failed"
	RetryScheduled      Stage = "retry.scheduled"
	RetryExhausted      Stage = "retry.exhausted"
	DeadLettered        Stage = "deadletter.stored"
	Replayed            Stage = "deadletter.replayed"
	Duplicate           Stage = "duplicate"
)

type StreamConfig

type StreamConfig struct {
	Name        string        `json:"name"        yaml:"name"`
	Provider    string        `json:"provider"    yaml:"provider"`
	Subjects    []string      `json:"subjects"    yaml:"subjects"`
	MaxAge      time.Duration `json:"maxAge"      yaml:"max_age"`
	MaxMessages int64         `json:"maxMessages" yaml:"max_messages"`
	Replicas    int           `json:"replicas"    yaml:"replicas"`
}

StreamConfig describes retained messages, independently of consumer progress.

type StreamInfo

type StreamInfo struct {
	Config    StreamConfig `json:"config"`
	Messages  uint64       `json:"messages"`
	Bytes     uint64       `json:"bytes"`
	Consumers int          `json:"consumers"`
}

StreamInfo is a provider's current retained-message state.

type Subscription

type Subscription interface {
	Next(ctx context.Context) (Delivery, error)
	Close(ctx context.Context) error
}

Subscription supplies deliveries to a bounded pool of workers.

type SubscriptionConfig

type SubscriptionConfig struct {
	ID            string        `json:"id"                      yaml:"id"`
	Stream        string        `json:"stream"                  yaml:"stream"`
	MessageType   string        `json:"messageType"             yaml:"message_type"`
	Mode          DeliveryMode  `json:"mode"                    yaml:"delivery"`
	Durable       bool          `json:"durable"                 yaml:"durable"`
	Concurrency   int           `json:"concurrency"             yaml:"concurrency"`
	MaxInFlight   int           `json:"maxInFlight"             yaml:"max_in_flight"`
	Timeout       time.Duration `json:"timeout"                 yaml:"timeout"`
	MaxAttempts   int           `json:"maxAttempts"             yaml:"max_attempts"`
	RetryDelay    time.Duration `json:"retryDelay"              yaml:"retry_delay"`
	StartAt       string        `json:"startAt"                 yaml:"start_at"`
	StartSequence uint64        `json:"startSequence,omitempty" yaml:"start_sequence"`
	// BroadcastID pins a durable broadcast subscriber across process restarts.
	BroadcastID string `json:"broadcastID,omitempty" yaml:"broadcast_id"`
}

SubscriptionConfig belongs to a logical service, with per-replica worker limits.

Jump to

Keyboard shortcuts

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