Documentation
¶
Overview ¶
Package conduit provides typed service communication with pluggable brokers.
Index ¶
- Constants
- Variables
- func Call[In, Out any](ctx context.Context, r *Runtime, service string, ...) (Out, error)
- func Clients(r *Runtime) *transport.Services
- func Handle[In, Out any](r *Runtime, procedure ProcedureType[In, Out], ...) error
- func NewID() string
- func Permanent(err error) error
- func Subscribe[T any](r *Runtime, event EventType[T], ...) error
- type Backfill
- type BackfillInput
- type Capabilities
- type Config
- type ConnectionConfig
- type ConsumerInfo
- type DeadLetter
- type DeliveryInfo
- type DeliveryMode
- type Endpoint
- type Envelope
- type EventType
- type Extension
- func (e *Extension) DepsSpec() []forge.Dep
- func (e *Extension) Health(ctx context.Context) error
- func (e *Extension) HookObserverDrops(ctx context.Context) (uint64, error)
- func (e *Extension) Register(app forge.App) error
- func (e *Extension) RegisterContractContributor(disp *dispatcher.Dispatcher, registry dashcontract.Registry, ...) error
- func (e *Extension) Runtime() *Runtime
- func (e *Extension) Start(ctx context.Context) error
- func (e *Extension) Stop(ctx context.Context) error
- type Hook
- type HookEvent
- type HookFuncs
- type Identity
- type Instance
- type Latency
- type Message
- type Middleware
- type Option
- type ProcedureType
- type Provider
- type Publication
- type PublishOption
- type RPCCode
- type RPCConfig
- type RPCError
- type Receipt
- type Registry
- type Request
- type Runtime
- type Snapshot
- type Stage
- type StreamConfig
- type SubscribeOption
- type SubscriptionConfig
Constants ¶
const ( Competing = core.Competing Broadcast = core.Broadcast Starting = core.Starting Ready = core.Ready Draining = core.Draining Stopped = core.Stopped ProviderConnected = core.ProviderConnected SubscriptionStarted = core.SubscriptionStarted SubscriptionFailed = core.SubscriptionFailed Publishing = core.Publishing Published = core.Published PublishFailed = core.PublishFailed Received = core.Received Handling = core.Handling Handled = core.Handled HandleFailed = core.HandleFailed Acknowledged = core.Acknowledged SettlementFailed = core.SettlementFailed RetryScheduled = core.RetryScheduled RetryExhausted = core.RetryExhausted DeadLettered = core.DeadLettered Replayed = core.Replayed Duplicate = core.Duplicate RPCCalling = core.RPCCalling RPCReceived = core.RPCReceived RPCHandling = core.RPCHandling RPCHandled = core.RPCHandled RPCReturned = core.RPCReturned RPCFailed = core.RPCFailed )
const ( RPCBadRequest = core.RPCBadRequest RPCNotFound = core.RPCNotFound RPCPermissionDenied = core.RPCPermissionDenied RPCConflict = core.RPCConflict RPCDeadlineExceeded = core.RPCDeadlineExceeded RPCCanceled = core.RPCCanceled RPCInternal = core.RPCInternal )
const ServiceKey = "conduit"
ServiceKey identifies the runtime in Forge's dependency container.
Variables ¶
var ( ErrNotRunning = core.ErrNotRunning ErrNotFound = core.ErrNotFound ErrConflict = core.ErrConflict ErrUnsupported = core.ErrUnsupported ErrOutcomeUnknown = core.ErrOutcomeUnknown )
Errors let callers distinguish rejected requests from unknown broker outcomes.
Functions ¶
func Call ¶
func Call[In, Out any](ctx context.Context, r *Runtime, service string, procedure ProcedureType[In, Out], input In, options ...PublishOption) (Out, error)
Call sends one typed request by logical service name with no automatic retries.
func Handle ¶
func Handle[In, Out any](r *Runtime, procedure ProcedureType[In, Out], handler func(context.Context, Request[In]) (Out, error)) error
Handle binds a typed procedure before Forge initializes the runtime.
Types ¶
type BackfillInput ¶
type BackfillInput = core.BackfillInput
type Capabilities ¶
type Capabilities = core.Capabilities
type ConnectionConfig ¶
type ConnectionConfig = core.ConnectionConfig
type ConsumerInfo ¶
type ConsumerInfo = core.ConsumerInfo
type DeadLetter ¶
type DeadLetter = core.DeadLetter
type DeliveryInfo ¶
type DeliveryInfo = core.DeliveryInfo
type DeliveryMode ¶
type DeliveryMode = core.DeliveryMode
type Extension ¶
type Extension struct {
*forge.BaseExtension
// contains filtered or unexported fields
}
Extension owns the runtime and its Forge lifecycle.
func NewExtension ¶
NewExtension constructs a Conduit extension with optional application overrides.
func NewExtensionWithConfig ¶
NewExtensionWithConfig supplies explicit topology while allowing identity inference.
func (*Extension) HookObserverDrops ¶
HookObserverDrops reports bounded observer backpressure independently of handler outcomes.
func (*Extension) RegisterContractContributor ¶
func (e *Extension) RegisterContractContributor(disp *dispatcher.Dispatcher, registry dashcontract.Registry, wardens dashcontract.WardenRegistry) error
RegisterContractContributor registers the conduit React dashboard's intent surface.
type Message ¶
type Message[T any] struct { Data T Envelope Envelope Delivery DeliveryInfo }
Message gives a handler both typed data and the original delivery identity.
type Middleware ¶
type Middleware = core.Middleware
type Option ¶
func WithConfig ¶
WithConfig supplies optional explicit Forge configuration.
func WithProvider ¶
WithProvider registers a named broker connection.
func WithRegistry ¶
WithRegistry enables leased service registration and named clients.
type ProcedureType ¶
type ProcedureType[Request, Response any] struct { Name string ValidateRequest func(Request) error ValidateResponse func(Response) error }
ProcedureType names a versioned request and response contract.
func Procedure ¶
func Procedure[Request, Response any](name string) ProcedureType[Request, Response]
Procedure declares typed broker RPC independently of the transport.
type Publication ¶
Publication is a validated envelope and its selected stream.
func Prepare ¶
func Prepare[T any](ctx context.Context, r *Runtime, event EventType[T], data T, opts ...PublishOption) (Publication, error)
Prepare enriches and validates typed data for publication or a transactional outbox.
type PublishOption ¶
type PublishOption func(*Envelope)
PublishOption enriches a draft. Source identity is always set by the runtime.
func Causation ¶
func Causation(id string) PublishOption
Causation links a derived message to its triggering event.
func Correlation ¶
func Correlation(id string) PublishOption
Correlation links messages to a request or operation.
func Headers ¶
func Headers(headers map[string]string) PublishOption
Headers supplies correlation, tracing and application metadata.
func Key ¶
func Key(key string) PublishOption
Key adds an application routing key without implying ordered processing.
func MessageID ¶
func MessageID(id string) PublishOption
MessageID lets a producer retry an unknown publish outcome with the same ID.
type StreamConfig ¶
type StreamConfig = core.StreamConfig
type SubscribeOption ¶
type SubscribeOption func(*subscribeOptions)
SubscribeOption binds a typed handler to a configured logical subscription.
func AroundHandle ¶
func AroundHandle(middleware ...Middleware) SubscribeOption
AroundHandle installs middleware in outermost-first order.
func Consumer ¶
func Consumer(id string) SubscribeOption
Consumer selects the stable subscription ID, shared by service replicas.
type SubscriptionConfig ¶
type SubscriptionConfig = core.SubscriptionConfig
Directories
¶
| Path | Synopsis |
|---|---|
|
cmd
|
|
|
demo
command
Command demo serves a local Conduit dashboard contract backed by real runtime deliveries.
|
Command demo serves a local Conduit dashboard contract backed by real runtime deliveries. |
|
Package contract exposes the Conduit runtime to Forge dashboard clients.
|
Package contract exposes the Conduit runtime to Forge dashboard clients. |
|
Package core defines provider contracts and the service communication runtime.
|
Package core defines provider contracts and the service communication runtime. |
|
Package discovery adapts instance registries to Conduit service clients.
|
Package discovery adapts instance registries to Conduit service clients. |
|
providers
|
|
|
conformance
Package conformance supplies the broker-independent Conduit adapter checks.
|
Package conformance supplies the broker-independent Conduit adapter checks. |
|
jetstream
Package jetstream implements durable streams and stable pull consumers using NATS.
|
Package jetstream implements durable streams and stable pull consumers using NATS. |
|
kafka
Package kafka implements retained events with a single partition per Conduit stream.
|
Package kafka implements retained events with a single partition per Conduit stream. |
|
memory
Package memory provides a shared development broker without disk durability.
|
Package memory provides a shared development broker without disk durability. |
|
redisstreams
Package redisstreams implements retained events with Redis Streams consumer groups.
|
Package redisstreams implements retained events with Redis Streams consumer groups. |
|
Package transaction coordinates PostgreSQL business writes with event delivery.
|
Package transaction coordinates PostgreSQL business writes with event delivery. |
|
Package transport binds standard API clients to logical service names.
|
Package transport binds standard API clients to logical service names. |