worker

package
v1.6.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 2, 2026 License: Apache-2.0 Imports: 10 Imported by: 0

Documentation

Index

Constants

View Source
const (
	WorkerStatusHealthy  = "healthy"
	WorkerStatusDegraded = "degraded"
	WorkerStatusFailed   = "failed"
)
View Source
const (
	ReasonStartupFailed  = "startup_failed"
	ReasonRecoveryFailed = "recovery_failed"
)

Variables

This section is empty.

Functions

func Ready added in v1.6.0

func Ready(ctx context.Context)

Ready tells the supervisor that the worker running with ctx is up, e.g. a consumer whose subscription is open. A worker with a RestartPolicy counts as stable only once it has been up for the stable window after calling Ready; one that never calls it is never stable. No-op outside a supervisor.

Types

type HealthTracker

type HealthTracker struct {
	// contains filtered or unexported fields
}

HealthTracker tracks the health status of all workers. It is safe for concurrent use.

func NewHealthTracker

func NewHealthTracker() *HealthTracker

NewHealthTracker creates a new HealthTracker.

func (*HealthTracker) GetStatus

func (h *HealthTracker) GetStatus() map[string]interface{}

GetStatus returns the overall health status with details of all workers. The overall status is the worst worker status: failed > degraded > healthy.

func (*HealthTracker) IsHealthy

func (h *HealthTracker) IsHealthy() bool

IsHealthy returns false if any worker has failed. Degraded workers still count as healthy: they are being restarted within their budget.

func (*HealthTracker) MarkDegraded added in v1.6.0

func (h *HealthTracker) MarkDegraded(name string, since time.Time, reason string)

MarkDegraded marks a worker as failing but still within its restart budget.

func (*HealthTracker) MarkFailed

func (h *HealthTracker) MarkFailed(name string)

MarkFailed marks a worker as failed. Note: Error details are NOT stored for security reasons.

func (*HealthTracker) MarkFailedWithReason added in v1.6.0

func (h *HealthTracker) MarkFailedWithReason(name string, since time.Time, reason string)

MarkFailedWithReason marks a worker as failed after its restart budget was spent.

func (*HealthTracker) MarkHealthy

func (h *HealthTracker) MarkHealthy(name string)

MarkHealthy marks a worker as healthy.

type Limits added in v1.6.0

type Limits struct {
	// MaxAttempts is the number of restarts. The worker escalates when the
	// run after the last allowed restart fails (0: on the first failure).
	MaxAttempts int
	// MaxDuration is the time since the first failed run of the episode.
	MaxDuration time.Duration
}

Limits bounds a failure episode, which starts at the first failed run and ends when a run stays up for the stable window. Whichever limit is hit first escalates the worker to failed. A negative value disables a limit.

type RegisterOption added in v1.6.0

type RegisterOption func(*registration)

RegisterOption configures how the supervisor runs a registered worker.

func WithRestartPolicy added in v1.6.0

func WithRestartPolicy(p RestartPolicy) RegisterOption

WithRestartPolicy restarts the worker after a failed run. See RestartPolicy.

type RestartPolicy added in v1.6.0

type RestartPolicy struct {
	Startup  Limits
	Recovery Limits
}

RestartPolicy makes the supervisor restart a worker whose Run returns an error, instead of marking it failed for good.

A worker is in the startup phase until a run has stayed up for the stable window (30s) once, counted from when the run calls Ready; after that it is in the recovery phase. Each phase has its own Limits. While within the limits the worker is degraded (/healthz 200); once either limit is hit it is failed (/healthz 503). The supervisor keeps restarting it either way, so it can come back to healthy.

Backoff between runs: 1s doubling, capped at 60s, ±20% jitter; 60s once failed.

type SupervisorOption

type SupervisorOption func(*WorkerSupervisor)

SupervisorOption configures a WorkerSupervisor.

func WithShutdownTimeout

func WithShutdownTimeout(timeout time.Duration) SupervisorOption

WithShutdownTimeout sets the maximum time to wait for workers to shutdown gracefully. After this timeout, Run() will return even if workers haven't finished. Default is 0 (no timeout - wait indefinitely).

type Worker

type Worker interface {
	// Name returns a unique identifier for this worker (e.g., "http-server", "retrymq-consumer")
	Name() string

	// Run executes the worker and blocks until context is cancelled or error occurs.
	// Returns nil or context.Canceled for graceful shutdown.
	// Returns error for fatal failures.
	Run(ctx context.Context) error
}

Worker represents a long-running background process. Each worker runs in its own goroutine and can be monitored for health.

Workers should: - Block in Run() until context is cancelled or a fatal error occurs - Return nil or context.Canceled for graceful shutdown - Return non-nil error only for fatal failures

type WorkerHealth

type WorkerHealth struct {
	Status string     `json:"status"`
	Since  *time.Time `json:"since,omitempty"`
	Reason string     `json:"reason,omitempty"`
}

WorkerHealth represents the health status of a single worker. Error details are NOT exposed for security reasons: /healthz is unauthenticated and broker errors carry hosts and credentials.

type WorkerSupervisor

type WorkerSupervisor struct {
	// contains filtered or unexported fields
}

WorkerSupervisor manages and supervises multiple workers. It tracks their health and handles graceful shutdown.

func NewWorkerSupervisor

func NewWorkerSupervisor(logger *logging.Logger, opts ...SupervisorOption) *WorkerSupervisor

NewWorkerSupervisor creates a new WorkerSupervisor.

func (*WorkerSupervisor) GetHealthTracker

func (r *WorkerSupervisor) GetHealthTracker() *HealthTracker

GetHealthTracker returns the health tracker for this supervisor.

func (*WorkerSupervisor) Register

func (r *WorkerSupervisor) Register(w Worker, opts ...RegisterOption)

Register adds a worker to the supervisor. Panics if a worker with the same name is already registered.

func (*WorkerSupervisor) Run

func (r *WorkerSupervisor) Run(ctx context.Context) error

Run starts all registered workers and supervises them. It blocks until: - ALL workers have exited (either successfully or with errors), OR - The context is cancelled (SIGTERM/SIGINT)

When a worker without a RestartPolicy fails, it marks the worker as failed but does NOT terminate other workers. Workers registered WithRestartPolicy are restarted instead (see RestartPolicy). This allows: - Other workers to continue serving (e.g., HTTP server stays up) - Health checks to report the failed worker status - Orchestrator to detect failure and restart the pod/container

Returns nil if context was cancelled and workers shutdown gracefully. Returns error if workers failed to shutdown within timeout (if configured).

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL