Versions in this module Expand all Collapse all v1 v1.0.0 Aug 26, 2026 Changes in this version + var ErrCommandNotAcknowledged = errors.New("dataplane: command not acknowledged") + var ErrControllerUnavailable = errors.New("dataplane: controller unavailable") + var ErrInvalidControllerConfiguration = errors.New("dataplane: invalid controller configuration") + var ErrInvalidFleetConfiguration = errors.New("dataplane: invalid fleet configuration") + var ErrInvalidFleetOutput = errors.New("dataplane: invalid fleet output") + var ErrInvalidRecordConfiguration = errors.New("dataplane: invalid record configuration") + var ErrInvalidRecordOutput = errors.New("dataplane: invalid record output") + var ErrInvalidRecordRequest = errors.New("dataplane: invalid record request") + var ErrInvalidStatusConfiguration = errors.New("dataplane: invalid status configuration") + var ErrInvalidStatusOutput = errors.New("dataplane: invalid status output") + var ErrInvalidStatusRequest = errors.New("dataplane: invalid status request") + var ErrRecordReaderUnavailable = errors.New("dataplane: record reader unavailable") + var ErrStatusReaderUnavailable = errors.New("dataplane: status reader unavailable") + var ErrUnsupportedCommand = errors.New("dataplane: unsupported command") + type ControllerDispatcher struct + func NewControllerDispatcher(resolver ControllerResolver, protocol queue.ProtocolVersion, ...) (*ControllerDispatcher, error) + func (d *ControllerDispatcher) Dispatch(ctx context.Context, command controlplane.Command) error + func (d *ControllerDispatcher) DispatchResult(ctx context.Context, command controlplane.Command) (control.DispatchOutcome, error) + type ControllerResolver interface + ResolveController func(context.Context, string) (queue.Controller, error) + type FleetSource struct + func NewFleetSource(source WorkerStatusSource) (*FleetSource, error) + func (s *FleetSource) SnapshotTenant(ctx context.Context, tenant string, now time.Time, staleAfter time.Duration) (fleet.RegistrySnapshot, error) + type RecordReaderResolver interface + ResolveRecordReader func(context.Context, string) (queue.RecordReader, error) + type RecordSource struct + func NewRecordSource(resolver RecordReaderResolver) (*RecordSource, error) + func (s *RecordSource) Inspect(ctx context.Context, tenant string, request queue.InspectRequest) (queue.JobRecord, error) + func (s *RecordSource) ListDeadLetters(ctx context.Context, tenant string, request queue.PageRequest) (queue.RecordPage, error) + func (s *RecordSource) ListFailures(ctx context.Context, tenant string, request queue.PageRequest) (queue.RecordPage, error) + type StatusReaderResolver interface + ResolveStatusReader func(context.Context, string) (queue.StatusReader, error) + type StatusSource struct + func NewStatusSource(resolver StatusReaderResolver) (*StatusSource, error) + func (s *StatusSource) ListQueues(ctx context.Context, tenant string, request queue.StatusPageRequest) (queue.QueueStatusPage, error) + func (s *StatusSource) ListWorkers(ctx context.Context, tenant string, request queue.StatusPageRequest) (queue.WorkerStatusPage, error) + type WorkerStatusSource interface + ListWorkers func(context.Context, string, queue.StatusPageRequest) (queue.WorkerStatusPage, error)