package worker import ( "context" "fmt" "log/slog" "time" "git.toki-labs.com/toki/gito/services/core/internal/config" "git.toki-labs.com/toki/gito/services/core/internal/core" ) type OperationPicker interface { PickQueuedOperation(ctx context.Context, now time.Time) (core.Operation, bool, error) } type Clock interface { Now() time.Time } type realClock struct{} func (realClock) Now() time.Time { return time.Now() } type Runner struct { cfg config.Config logger *slog.Logger picker OperationPicker clock Clock } func NewRunner(cfg config.Config, logger *slog.Logger, picker OperationPicker, clock Clock) *Runner { if clock == nil { clock = realClock{} } return &Runner{ cfg: cfg, logger: logger, picker: picker, clock: clock, } } func (r *Runner) Run() error { return r.RunOnce(context.Background()) } func (r *Runner) RunOnce(ctx context.Context) error { if !r.cfg.WorkerEnabled { r.logger.Info("worker disabled") return nil } if r.picker == nil { r.logger.Warn("worker picker not configured") return nil } op, ok, err := r.picker.PickQueuedOperation(ctx, r.clock.Now()) if err != nil { r.logger.Error("failed to pick operation", "error", err) return fmt.Errorf("pick operation: %w", err) } if !ok { r.logger.Debug("no pending operation found") return nil } r.logger.Info("picked operation", "operation_id", op.ID, "type", op.Type) return nil }