Documentation
¶
Overview ¶
Package dispatcher implements transport.Dispatcher and transport.SubscriptionSource against a function-table of registered handlers. Contributors register their intent handlers via Register / RegisterSubscription / RegisterContributor; the HTTP and SSE transports look them up at request time.
See SLICE_C_DESIGN.md in the parent contract directory for the spec this implements.
Index ¶
- Constants
- Variables
- func RegisterCommand[I, O any](d *Dispatcher, contributor, intent string, version int, ...) error
- func RegisterQuery[I, O any](d *Dispatcher, contributor, intent string, version int, ...) error
- func RegisterSubscription[P, E any](d *Dispatcher, contributor, intent string, version int, ...) error
- type Contributor
- type Dispatcher
- func (d *Dispatcher) Dispatch(ctx context.Context, req contract.Request, p contract.Principal) (json.RawMessage, contract.ResponseMeta, error)
- func (d *Dispatcher) Register(contributor, intent string, version int, h Handler, opts ...RegisterOption) error
- func (d *Dispatcher) RegisterContributor(c Contributor) error
- func (d *Dispatcher) RegisterSubscription(contributor, intent string, version int, h SubscriptionHandler) error
- func (d *Dispatcher) SetRemoteDispatcher(rd RemoteDispatcher)
- func (d *Dispatcher) Subscribe(ctx context.Context, p contract.Principal, contributor string, ...) (<-chan contract.StreamEvent, func(), error)
- type Handler
- type IdempotencyCached
- type IdempotencyClaim
- type IdempotencyClaimer
- type IdempotencyStore
- type IntentRef
- type LoggerAuditEmitter
- type MetricsEmitter
- type NoopMetricsEmitter
- type Option
- type PrometheusMetricsEmitter
- type RegisterOption
- type RemoteDispatcher
- type Result
- type SubscriptionHandler
Constants ¶
const DefaultIdempotencyWait = 10 * time.Second
DefaultIdempotencyWait is how long a command waits, by default, for a concurrent dispatch that holds its idempotency key. It matches the HTTP idempotency middleware's default wait.
const TombstoneStatus = http.StatusConflict
TombstoneStatus is the Status of the idempotency entry a SecretResponse command leaves behind. The entry has no WireBody. Any entry with this status answers CONFLICT on lookup, whatever the handler's current registration says, so a tombstone never falls through to a fresh dispatch.
Variables ¶
var ErrIdempotencyClaimHeld = errors.New("dispatcher: idempotency key is held by a running command")
ErrIdempotencyClaimHeld is what IdempotencyClaimer.Claim's error wraps when its context ended while another dispatch still held the key.
var ErrIdempotencyClaimLost = errors.New("dispatcher: idempotency claim lapsed before it ended")
ErrIdempotencyClaimLost is what IdempotencyClaim.End's error wraps when the claim's lease lapsed before End, so End stored nothing.
Functions ¶
func RegisterCommand ¶
func RegisterCommand[I, O any](d *Dispatcher, contributor, intent string, version int, fn func(ctx context.Context, in I, p contract.Principal) (O, error), opts ...RegisterOption) error
RegisterCommand is identical in shape to RegisterQuery; both register a query/command handler. The dispatcher's wire layer enforces kind/capability matching against the manifest, so the only practical difference between the two helpers is the options RegisterCommand takes.
Pass SecretResponse for a command whose response carries a secret the caller sees once (a raw API key, say). Its response is never kept for idempotent replay:
dispatcher.RegisterCommand(d, "keysmith", "keys.create", 1, createKey, dispatcher.SecretResponse())
func RegisterQuery ¶
func RegisterQuery[I, O any](d *Dispatcher, contributor, intent string, version int, fn func(ctx context.Context, in I, p contract.Principal) (O, error)) error
RegisterQuery wraps a typed handler in a Handler-compatible closure that JSON-decodes Payload into I and encodes the returned O into Result.Data. I and O must be JSON-marshallable. Use struct{} for an empty-input intent.
func RegisterSubscription ¶
func RegisterSubscription[P, E any](d *Dispatcher, contributor, intent string, version int, fn func(ctx context.Context, in P, p contract.Principal) (<-chan E, func(), error)) error
RegisterSubscription wraps a typed subscription handler. The pump goroutine JSON-encodes each typed E event into a contract.StreamEvent before forwarding into the broker's channel.
Types ¶
type Contributor ¶
type Contributor interface {
Name() string
Handlers() map[IntentRef]Handler
Subscriptions() map[IntentRef]SubscriptionHandler
}
Contributor is layer (b)'s registration shape: a contributor publishes its handler and subscription tables, and the dispatcher walks them on Register.
type Dispatcher ¶
type Dispatcher struct {
// contains filtered or unexported fields
}
Dispatcher is the concrete implementation of transport.Dispatcher and transport.SubscriptionSource (Subscribe lives in subscription.go). Contributors register handlers indexed by (contributor, intent, version); dispatch is a map lookup + the handler call wrapped in metrics emission and canonical error mapping.
func New ¶
func New(metrics MetricsEmitter) *Dispatcher
New returns a fresh dispatcher. Pass NoopMetricsEmitter{} for tests / dev; slice (b) provides a Prometheus-backed implementation.
func NewWithOptions ¶
func NewWithOptions(metrics MetricsEmitter, opts ...Option) *Dispatcher
NewWithOptions returns a dispatcher configured with the supplied options. The existing New(metrics) constructor is preserved as a thin wrapper.
func (*Dispatcher) Dispatch ¶
func (d *Dispatcher) Dispatch(ctx context.Context, req contract.Request, p contract.Principal) (json.RawMessage, contract.ResponseMeta, error)
Dispatch implements transport.Dispatcher. When a tracer is configured, a span wraps the dispatch with attributes capturing (contributor, intent, version, kind) and a status reflecting the outcome.
func (*Dispatcher) Register ¶
func (d *Dispatcher) Register(contributor, intent string, version int, h Handler, opts ...RegisterOption) error
Register binds a query/command handler to a (contributor, intent, version) key. Returns an error on duplicate registration. Pass SecretResponse for a command whose response must never be kept for replay.
func (*Dispatcher) RegisterContributor ¶
func (d *Dispatcher) RegisterContributor(c Contributor) error
RegisterContributor walks a Contributor's Handlers() and Subscriptions() maps and registers each one. Atomic: if any registration fails, all preceding registrations from this call are rolled back.
func (*Dispatcher) RegisterSubscription ¶
func (d *Dispatcher) RegisterSubscription(contributor, intent string, version int, h SubscriptionHandler) error
RegisterSubscription binds a subscription handler to (contributor, intent, version).
func (*Dispatcher) SetRemoteDispatcher ¶
func (d *Dispatcher) SetRemoteDispatcher(rd RemoteDispatcher)
SetRemoteDispatcher installs (or clears, with nil) the fallback consulted when no local handler matches a request. Idempotent — the dispatcher's forwarding plumbing typically calls this once during wire-up.
func (*Dispatcher) Subscribe ¶
func (d *Dispatcher) Subscribe(ctx context.Context, p contract.Principal, contributor string, intent contract.Intent, params map[string]contract.ParamSource) (<-chan contract.StreamEvent, func(), error)
Subscribe implements transport.SubscriptionSource. The broker calls this on each subscribe-control message; the dispatcher routes to the registered handler. Params from YAML (map[string]contract.ParamSource) are flattened into a runtime map[string]any using the From string when set, the literal Value otherwise.
type Handler ¶
type Handler func(ctx context.Context, payload json.RawMessage, params map[string]any, p contract.Principal) (*Result, error)
Handler is the foundation function-table handler signature for query and command intents. Returning a *contract.Error propagates the canonical code to the wire; any other error is wrapped as CodeInternal at dispatch time.
type IdempotencyCached ¶
type IdempotencyCached struct {
Status int
WireBody json.RawMessage
StoredAt time.Time
TTL time.Duration
}
IdempotencyCached mirrors idempotency.Cached; defined here for the same import-cycle reason. Adapters in the wire-up convert between the two.
type IdempotencyClaim ¶ added in v1.12.2
type IdempotencyClaim struct {
// Cached is the entry already stored for the key.
Cached *IdempotencyCached
// End ends a claim the caller holds. Pass the entry to store under the
// claim, or nil to store nothing and give the key back. The dispatcher
// calls it exactly once. When the claim lapsed before End, End stores
// nothing and its error wraps ErrIdempotencyClaimLost.
End func(ctx context.Context, c *IdempotencyCached) error
}
IdempotencyClaim mirrors idempotency.Claim, for the same import-cycle reason as IdempotencyCached.
type IdempotencyClaimer ¶ added in v1.12.2
type IdempotencyClaimer interface {
IdempotencyStore
// Claim takes (key, identity) for the caller. While another caller holds
// the key it waits for that claim to end, until ctx ends. Exactly one of
// the returned claim's fields is set: Cached when an entry is stored for
// the key, End when the caller now holds it. When ctx ends with the key
// still held, the error wraps ErrIdempotencyClaimHeld.
Claim(ctx context.Context, key, identity string) (IdempotencyClaim, error)
}
IdempotencyClaimer is an IdempotencyStore that can also hold a key while a command's handler runs. The dispatcher finds it by type assertion on the store passed to WithIdempotencyStore. With one, a command with an idempotency key claims (key, identity) before its handler runs and holds the claim until its entry is stored or the handler fails, so two overlapping dispatches never both run it. Without one, the dispatcher only looks the key up before the handler and stores the entry after it.
A Claim that fails for any reason other than the key being held (a backend error, say) answers a retryable UNAVAILABLE and the handler does not run, since nothing would stop a duplicate from running beside it.
When End reports ErrIdempotencyClaimLost, the claim's lease lapsed while the handler ran. The dispatcher then writes the entry with Store, so a secret command still leaves its tombstone.
type IdempotencyStore ¶
type IdempotencyStore interface {
Lookup(ctx context.Context, key, identity string) (*IdempotencyCached, bool)
Store(ctx context.Context, key, identity string, c IdempotencyCached) error
}
IdempotencyStore is the minimal surface the dispatcher needs from extensions/dashboard/contract/idempotency. Defining it here avoids an import cycle (the idempotency package is consumed only via this interface).
type LoggerAuditEmitter ¶
type LoggerAuditEmitter struct {
// contains filtered or unexported fields
}
LoggerAuditEmitter writes audit records as info-level structured logs via a forge.Logger. Each AuditRecord field is emitted as a discrete log field so log aggregators can filter by `audit=true` cheaply. Pass a nil logger to disable — the emitter then becomes a noop.
func NewLoggerAuditEmitter ¶
func NewLoggerAuditEmitter(logger forge.Logger) *LoggerAuditEmitter
NewLoggerAuditEmitter returns an emitter that writes via logger. Pass nil to disable (the emitter becomes a noop).
func (*LoggerAuditEmitter) Emit ¶
func (e *LoggerAuditEmitter) Emit(_ context.Context, rec contract.AuditRecord)
Emit implements contract.AuditEmitter.
type MetricsEmitter ¶
type MetricsEmitter interface {
RecordDispatch(ctx context.Context, contributor, intent string, version int, kind contract.Kind, latency time.Duration, errCode contract.ErrorCode)
}
MetricsEmitter ships dispatch metrics to a backend. The Phase 4 expansion of this file adds the full DispatchInfo struct and the noop default. For Phase 1, only the interface and the noop are needed.
type NoopMetricsEmitter ¶
type NoopMetricsEmitter struct{}
NoopMetricsEmitter discards all dispatch metrics.
type Option ¶
type Option func(*Dispatcher)
Option configures a Dispatcher.
func WithIdempotencyStore ¶
func WithIdempotencyStore(s IdempotencyStore) Option
WithIdempotencyStore wires command dedup. When set, commands carrying a non-empty IdempotencyKey are deduped per-user via the store.
func WithIdempotencyWait ¶ added in v1.12.2
WithIdempotencyWait bounds how long a command waits for a concurrent dispatch that holds its idempotency key, when the store is an IdempotencyClaimer. When the wait ends with the key still held, the command answers CONFLICT. A value of zero or less keeps DefaultIdempotencyWait.
Keep the wait below your HTTP server's WriteTimeout. A server that times the write out first cuts off a waiting duplicate's response, so the client sees a dropped connection where it should have seen the replay or CONFLICT.
func WithTracer ¶
WithTracer configures the dispatcher to open a span per Dispatch call. Passing a nil tracer is equivalent to not supplying the option at all.
type PrometheusMetricsEmitter ¶
type PrometheusMetricsEmitter struct {
// contains filtered or unexported fields
}
PrometheusMetricsEmitter records dispatch metrics into a forge.Metrics registry. Counters and histograms are created lazily on first emission (forge.Metrics.Counter / .Histogram are get-or-create). Pass nil to disable — the emitter then becomes a noop.
func NewPrometheusMetricsEmitter ¶
func NewPrometheusMetricsEmitter(m forge.Metrics) *PrometheusMetricsEmitter
NewPrometheusMetricsEmitter returns an emitter that writes to m. If m is nil, the emitter is a noop.
type RegisterOption ¶ added in v1.12.1
type RegisterOption func(*handlerEntry)
RegisterOption configures one handler registration. Register and RegisterCommand take any number of them, so a call without options keeps compiling and behaving as before.
func SecretResponse ¶ added in v1.12.1
func SecretResponse() RegisterOption
SecretResponse marks a command whose response carries a secret the caller sees once, such as a freshly minted API key. The dispatcher never keeps that response for idempotent replay. A successful dispatch with an idempotency key stores a tombstone (TombstoneStatus, no body), and a later dispatch with the same key and user answers CONFLICT without running the handler again, because running it again would mint a second secret.
type RemoteDispatcher ¶
type RemoteDispatcher interface {
Dispatch(ctx context.Context, req contract.Request, p contract.Principal) (json.RawMessage, contract.ResponseMeta, error)
}
RemoteDispatcher is the fallback Dispatcher consults when no local handler is registered for an envelope. The dispatcher passes the verbatim request through; implementations typically POST it to a peer service. Slice (m) added this so dashboards can aggregate contributors from multiple upstreams without baking forwarding into the dispatcher itself.
type Result ¶
type Result struct {
// Data is the JSON-encoded response body. May be nil for a {data: null} response.
Data json.RawMessage
// ExtraInvalidates is appended to the manifest's declared Invalidates.
ExtraInvalidates []string
// CacheOverride, when non-nil, replaces the manifest's declared cache hint.
CacheOverride *contract.CacheHint
}
Result carries the data payload plus optional response-meta overrides. Handlers that don't need to influence meta can return &Result{Data: ...}; handlers that need to add invalidations or override cache hints set the extra fields.
type SubscriptionHandler ¶
type SubscriptionHandler func(ctx context.Context, params map[string]any, p contract.Principal) (<-chan contract.StreamEvent, func(), error)
SubscriptionHandler is the function-table handler for subscription intents. The handler returns a channel of events, a force-stop function, and an optional error. Closing the channel signals end-of-stream; cancelling ctx is the canonical way to ask the handler to stop emitting.