Documentation
¶
Overview ¶
Package core defines provider contracts and the service communication runtime.
Index ¶
- Variables
- func IsPermanent(err error) bool
- func Matches(pattern, topic string) bool
- func NewID() string
- func Permanent(err error) error
- func ValidateMessageID(id string) error
- func ValidateTopic(topic string, wildcard bool) error
- type Backfill
- type BackfillInput
- type Binding
- type Capabilities
- type Config
- type ConnectionConfig
- type ConsumerInfo
- type DeadLetter
- type Delivery
- type DeliveryInfo
- type DeliveryMode
- type Endpoint
- type Envelope
- type HandleHook
- type Handler
- type Hook
- type HookEvent
- type HookFuncs
- type Identity
- type Instance
- type InstanceLister
- type Latency
- type Management
- type Middleware
- type Observer
- type Operations
- type Option
- type Provider
- type ProviderInfo
- type PublishHook
- type RPCBinding
- type RPCCode
- type RPCConfig
- type RPCError
- type RPCHandler
- type RPCProvider
- type RPCRequest
- type RPCResponse
- type RPCServer
- type Receipt
- type Registry
- type Resolver
- type Runtime
- func (r *Runtime) Backfill(ctx context.Context, in BackfillInput) (Backfill, error)
- func (r *Runtime) Backfills(ctx context.Context, provider, cursor string, limit int) ([]Backfill, string, error)
- func (r *Runtime) Bind(id, messageType string, handler Handler, middlewares ...Middleware) error
- func (r *Runtime) BindRPC(method string, handler RPCHandler) error
- func (r *Runtime) CallRPC(ctx context.Context, service, method string, draft Envelope, ...) (RPCResponse, error)
- func (r *Runtime) Configuration() Config
- func (r *Runtime) Configure(cfg Config, extra ...Option) error
- func (r *Runtime) Consumers(ctx context.Context) ([]ConsumerInfo, error)
- func (r *Runtime) DeadLetters(ctx context.Context, providerName, subscription, cursor string, limit int) ([]DeadLetter, string, error)
- func (r *Runtime) DurableStream(name string) bool
- func (r *Runtime) Health(ctx context.Context) error
- func (r *Runtime) Identity() Identity
- func (r *Runtime) Instances(ctx context.Context) ([]Instance, error)
- func (r *Runtime) PauseSubscription(ctx context.Context, id string, paused bool) error
- func (r *Runtime) Prepare(ctx context.Context, streamName string, draft Envelope, ...) (Envelope, error)
- func (r *Runtime) Providers() map[string]Provider
- func (r *Runtime) Publish(ctx context.Context, streamName string, draft Envelope, ...) (Receipt, error)
- func (r *Runtime) RecentEvents() []HookEvent
- func (r *Runtime) RegisterHook(hook Hook) error
- func (r *Runtime) ReplayDeadLetter(ctx context.Context, providerName, subscription, id string) (Receipt, error)
- func (r *Runtime) Resolver() Resolver
- func (r *Runtime) Send(ctx context.Context, streamName string, msg Envelope) (Receipt, error)
- func (r *Runtime) SetRegistry(registry Registry) error
- func (r *Runtime) Snapshot(ctx context.Context) (Snapshot, error)
- func (r *Runtime) Start(ctx context.Context) error
- func (r *Runtime) Stop(ctx context.Context) error
- func (r *Runtime) StreamFor(messageType string) (StreamConfig, error)
- type Snapshot
- type Stage
- type StreamConfig
- type StreamInfo
- type Subscription
- type SubscriptionConfig
Constants ¶
This section is empty.
Variables ¶
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 ¶
IsPermanent reports whether an error explicitly ends retries.
func NewID ¶
func NewID() string
NewID generates an identifier independent of broker sequence numbers.
func ValidateMessageID ¶
ValidateMessageID accepts stable IDs that can also address persisted recovery records.
func ValidateTopic ¶
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 ¶
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.
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.
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 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 ¶
BeforeHandle runs the configured control callback.
func (HookFuncs) BeforePublish ¶
BeforePublish runs the configured control callback.
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.
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.
type InstanceLister ¶
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 ¶
Middleware surrounds each processing attempt.
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 ¶
Option supplies a provider or an extension point.
func WithConfig ¶
WithConfig supplies explicit Forge configuration, overriding loaded non-zero values.
func WithProvider ¶
WithProvider registers a named connection independently of message routing.
func WithRegistry ¶
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 ¶
PublishHook can enrich or reject a draft before final validation and serialization.
type RPCBinding ¶
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 ¶
RPCError exposes only an intentional public error message.
func PublicRPCError ¶
PublicRPCError strips arbitrary handler causes from responses.
type RPCHandler ¶
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 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 NewDeferred ¶
NewDeferred allows handler binding before Forge loads application configuration.
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 ¶
Configuration returns a copy for Forge initialization. It may contain credentials.
func (*Runtime) Configure ¶
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 ¶
DurableStream checks a provider before transactional publication.
func (*Runtime) PauseSubscription ¶
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 ¶
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 ¶
RecentEvents returns bounded diagnostics without message data, headers or error text.
func (*Runtime) RegisterHook ¶
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) Send ¶
Send publishes an already prepared envelope, preserving its ID across outbox retries.
func (*Runtime) SetRegistry ¶
SetRegistry fills the discovery adapter before startup.
func (*Runtime) Snapshot ¶
Snapshot reports unavailable inspection as an error, never an empty healthy list.
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.