workers

package
v1.0.0 Latest Latest
Warning

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

Go to latest
Published: Sep 16, 2026 License: Apache-2.0 Imports: 8 Imported by: 0

Documentation

Overview

Package contract keeps all the interfaces required for background workers.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func GetOverseerSleepTimeout

func GetOverseerSleepTimeout() time.Duration

func SetLogger

func SetLogger(l logger.Logger)

SetLogger allows the logger to be replaced or injected externally, particularly useful in testing or for custom output destinations.

func SetOverseerSleepTimeout

func SetOverseerSleepTimeout(timeout time.Duration)

func StartOverseer

func StartOverseer(workers []*WorkerConfig, logger logger.Logger)

StartOverseer initializes the overseer with a list of workers and a logger, ensuring it only starts once. It launches the overseer loop in a separate goroutine. The overseer manages lifecycle of workers and restarts them upon failure.

func StartOverseerWithContext

func StartOverseerWithContext(ctx context.Context, workers []*WorkerConfig, logger logger.Logger, options *OverseerOptions)

Types

type OverseerOptions

type OverseerOptions struct {
	RestartPolicy RestartPolicy
	Hooks         WorkerHooks
}

type RestartPolicy

type RestartPolicy struct {
	Limit      int
	Window     time.Duration
	MinBackoff time.Duration
	MaxBackoff time.Duration
}

type WorkerConfig

type WorkerConfig struct {
	Name      string
	New       WorkerFactory
	MaxCount  int64
	IsEnabled bool
}

func (*WorkerConfig) NewWorker

func (cfg *WorkerConfig) NewWorker() (WorkerInterface, error)

func (*WorkerConfig) Validate

func (cfg *WorkerConfig) Validate() error

type WorkerEvent

type WorkerEvent struct {
	Name         string
	ID           string
	Err          error
	RestartCount int
	At           time.Time
}

type WorkerFactory

type WorkerFactory func() (WorkerInterface, error)

WorkerFactory builds a fresh worker instance for each goroutine launch. Returning a new handle per invocation avoids shared mutable state across workers.

type WorkerHooks

type WorkerHooks struct {
	OnStart   func(WorkerEvent)
	OnStop    func(WorkerEvent)
	OnError   func(WorkerEvent)
	OnRestart func(WorkerEvent)
}

type WorkerInterface

type WorkerInterface interface {

	// Get human readable name for the worker. This human readable name
	// can be useful in more contextualised logging, and tagging purposes.
	GetWorkerName() string

	// Get unique alphanumeric worker id for the current instance
	// of running worker (go-routine)
	GetWorkerId() string

	// Set unique alphanumeric worker id for the current instance
	// of running worker (go-routine)
	SetWorkerId(id string)

	// Get errors that have occurred during the execution
	// of current execution cycle of the worker loop.
	// Usually these errors are caught and recovered through
	// panic-recover workflow.
	GetWorkerExecutionErr() error

	// Set errors to the current instance of running worker
	// that have been occurred during the execution.
	SetWorkerExecutionErr(err error)

	// The main execution method of the worker implementation.
	// This method is called from outside by worker overseer and
	// is responsible for processing all the queued messages through
	// an always running loop. In case of any exception/error or panic
	// situation, this method should be able to recover from that and
	// emit relevant message to the callee, so while current running
	// instance of worker goes down, but the callee is able to respawn
	// another similar instance to carry on the queue processing flow.
	Run(workerChan chan<- WorkerInterface) error
}

Worker interface provides a common contract for all kinds of background workers. These workers are monitored and managed by an worker overseer, that invokes methods of worker interface's implementations from outside.

type WorkerStatus

type WorkerStatus struct {
	Name          string
	ID            string
	Running       bool
	RestartCount  int
	LastError     string
	LastStartedAt time.Time
	LastStoppedAt time.Time
}

func SnapshotWorkerStatuses

func SnapshotWorkerStatuses() []WorkerStatus

Jump to

Keyboard shortcuts

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