Versions in this module Expand all Collapse all v0 v0.1.0 Sep 7, 2026 Changes in this version + const DispatchReasonNetworkError + const DispatchReasonNoPeers + const DispatchReasonPeerDiscoveryError + const DispatchReasonWorkerRejected + const InternalProtocolVersion + const ReasonAmbiguousTask + const ReasonContention + const ReasonInvalidStatus + const ReasonMalformed + const ReasonMissingRun + const ReasonNotOwner + const ReasonStaleGeneration + const ReasonTaskNotRunning + const ReasonWrongWorker + var ErrOwnerBusy = errors.New("owner busy: retryable") + var ErrPeerUnreachable = errors.New("peer unreachable") + var ValidCompleteStatuses = map[string]bool + func ClientTLSConfig(c MTLSConfig) (*tls.Config, error) + func ClientTLSConfigFromHolder(holder *dispatchpki.MaterialHolder) (*tls.Config, error) + func ConfigureInternalMTLS(clientTLS *tls.Config) + func PostDispatch(ctx context.Context, targetURL, token string, req DispatchRequest) (bool, error) + func ServerTLSConfig(c MTLSConfig) (*tls.Config, error) + func ServerTLSConfigFromHolder(holder *dispatchpki.MaterialHolder) (*tls.Config, error) + func WarnIfNoToken(token string) + type CapabilitiesResponse struct + Capabilities []string + NodeID string + ProtocolVersion int + func GetCapabilities(ctx context.Context, targetURL, token string) (*CapabilitiesResponse, error) + func (c *CapabilitiesResponse) Supports(name string) bool + type CapabilityAdvertiser interface + Capabilities func() []string + type CompleteRequest struct + Attempt int + BranchSelections []string + Error string + Outputs map[string]string + OwnerGeneration int64 + Partitions []pkgtask.Partition + Result string + RunID uuid.UUID + Status string + TaskID uuid.UUID + TaskRunID uuid.UUID + WorkerNode string + type CompleteResponse struct + Accepted bool + Reason string + func PostComplete(ctx context.Context, ownerURL, token string, req CompleteRequest) (*CompleteResponse, error) + type DispatchLoop struct + func NewDispatchLoop(cfg DispatchLoopConfig) *DispatchLoop + func (l *DispatchLoop) Run(ctx context.Context) + type DispatchLoopConfig struct + APIPort int + BatchSize int + Deadline time.Duration + InternalPort int + Interval time.Duration + LeaseStore OwnerReader + LeaseTTL time.Duration + NodeID string + OwnerManager *run.OwnerManager + PeerBaseURL func(nodeAddr string) string + Peers PeerLister + RateLimitDB *gorm.DB + RateLimiter RateLimiter + Store TaskPendingReader + Token string + type DispatchRequest struct + Attempt int + Deadline time.Time + OwnerBaseURL string + OwnerGeneration int64 + RunID uuid.UUID + TaskID uuid.UUID + TaskRunID uuid.UUID + WorkerNode string + type ErrorResponse struct + Code string + Message string + type Handler struct + func NewHandler(store *run.Store, leaseStore *run.LeaseStore, nodeID, token string) *Handler + func (h *Handler) HandleCapabilities(w http.ResponseWriter, r *http.Request) + func (h *Handler) HandleComplete(w http.ResponseWriter, r *http.Request) + func (h *Handler) HandleDispatch(w http.ResponseWriter, r *http.Request) + func (h *Handler) WithOwnerManager(m *run.OwnerManager) *Handler + func (h *Handler) WithWorkerSubmitter(s WorkerSubmitter) *Handler + type InboundDispatch struct + Attempt int + OwnerBaseURL string + OwnerGeneration int64 + Task *models.TaskRun + WorkerNode string + type InternalServer struct + func NewInternalServer(handler *Handler, addr string, tlsConfig *tls.Config) *InternalServer + func (s *InternalServer) Run(ctx context.Context) error + func (s *InternalServer) Shutdown(ctx context.Context) error + type MTLSConfig struct + CAFile string + CertFile string + KeyFile string + func (c MTLSConfig) Configured() bool + type OwnerReader interface + AcquireExpiredLeases func(ctx context.Context, newOwner string, ttl time.Duration) (int64, error) + OwnedRunsWithGenerations func(ctx context.Context, ownerNode string) (map[uuid.UUID]int64, error) + type PeerLister interface + DispatchPeers func(ctx context.Context) ([]string, error) + type PeerListerFunc func(context.Context) ([]string, error) + func (f PeerListerFunc) DispatchPeers(ctx context.Context) ([]string, error) + type RateLimiter interface + Acquire func(ctx context.Context, resource string, units, limit int, window time.Duration) (bool, error) + type TaskPendingReader interface + PendingTasksForDispatch func(ctx context.Context, runID uuid.UUID, limit int) ([]models.TaskRun, error) + type WorkerSubmitter interface + SubmitDispatched func(d InboundDispatch) error