package threadpool import ( "log/slog" "runtime" "time" "github.com/alitto/pond/v2" ) type ConfigProvider interface { GetMaxWorkers() int GetThreadMonitorInterval() int GetThreadPool() pond.Pool ThreadPoolStopAndWait() StartThreadPool(*slog.Logger) } type ThreadPoolConfig struct { MaxWorkers int `env:"MAX_WORKERS"` MonitorInterval int `env:"THREAD_MONITOR_INTERVAL"` pool pond.Pool stopMonitor chan struct{} logger *slog.Logger } type PoolMetrics struct { QueueUtilization float64 AverageWaitTime time.Duration ProcessingTime time.Duration } func (dp *ThreadPoolConfig) StartThreadPool(logger *slog.Logger) { dp.logger = logger if dp.MaxWorkers <= 0 { dp.MaxWorkers = runtime.NumCPU() * 10 } if dp.MonitorInterval <= 0 { dp.MonitorInterval = 120 } dp.stopMonitor = make(chan struct{}) dp.pool = pond.NewPool(dp.MaxWorkers) // Start the monitoring heartbeat dp.pool.Submit(dp.monitor) } // Log thread pool stats func (dp *ThreadPoolConfig) monitor() { ticker := time.NewTicker(time.Duration(dp.MonitorInterval) * time.Second) defer ticker.Stop() for { select { case <-dp.stopMonitor: return case <-ticker.C: dp.logger.Debug( "thread pool heartbeat", "running_workers", dp.pool.RunningWorkers(), "queue_size", dp.pool.QueueSize(), "submitted_tasks", dp.pool.SubmittedTasks(), "waiting_tasks", dp.pool.WaitingTasks(), "dropped_tasks", dp.pool.DroppedTasks(), ) } } } func (dp *ThreadPoolConfig) GetMaxWorkers() int { return dp.MaxWorkers } func (dp *ThreadPoolConfig) GetThreadMonitorInterval() int { return dp.MonitorInterval } func (dp *ThreadPoolConfig) GetThreadPool() pond.Pool { if dp.pool == nil { dp.StartThreadPool(slog.Default()) } return dp.pool } func (dp *ThreadPoolConfig) ThreadPoolStopAndWait() { waitCh := make(chan struct{}) go func() { close(dp.stopMonitor) dp.pool.StopAndWait() close(waitCh) }() timeout := 3 * time.Second select { case <-waitCh: dp.logger.Info("thread pool shutdown complete") case <-time.After(timeout): dp.logger.Error("thread pool shotdown timeout") } }