Documentation
¶
Overview ¶
Package contract keeps all the interfaces required for background workers.
Index ¶
- func GetOverseerSleepTimeout() time.Duration
- func SetLogger(l logger.Logger)
- func SetOverseerSleepTimeout(timeout time.Duration)
- func StartOverseer(workers []*WorkerConfig, logger logger.Logger)
- func StartOverseerWithContext(ctx context.Context, workers []*WorkerConfig, logger logger.Logger, ...)
- type OverseerOptions
- type RestartPolicy
- type WorkerConfig
- type WorkerEvent
- type WorkerFactory
- type WorkerHooks
- type WorkerInterface
- type WorkerStatus
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func GetOverseerSleepTimeout ¶
func SetLogger ¶
SetLogger allows the logger to be replaced or injected externally, particularly useful in testing or for custom output destinations.
func SetOverseerSleepTimeout ¶
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 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 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