From 098484a0f87515a68ba4a66a464202e2341af684 Mon Sep 17 00:00:00 2001 From: Valery Piashchynski Date: Sun, 4 Oct 2026 20:06:22 +0200 Subject: [PATCH] refactor: use modern Go helpers Signed-off-by: Valery Piashchynski --- ipc/pipe/pipe_test.go | 5 +++-- worker_watcher/worker_watcher.go | 9 +++++---- worker_watcher/worker_watcher_test.go | 4 +--- 3 files changed, 9 insertions(+), 9 deletions(-) diff --git a/ipc/pipe/pipe_test.go b/ipc/pipe/pipe_test.go index 781fdd0..e1f6a0d 100644 --- a/ipc/pipe/pipe_test.go +++ b/ipc/pipe/pipe_test.go @@ -3,6 +3,7 @@ package pipe import ( "context" "log/slog" + "maps" "os" "os/exec" "runtime" @@ -468,7 +469,7 @@ func Test_Pipe_SpawnTimeout_ReapsWorker(t *testing.T) { t.Skip("pgrep is not available on windows") } before := childPids(t) - ctx, cancel := context.WithTimeout(context.Background(), time.Millisecond*10) + ctx, cancel := context.WithTimeout(t.Context(), time.Millisecond*10) defer cancel() w, err := NewPipeFactory(log).SpawnWorkerWithContext(ctx, exec.Command("php", "../../tests/slow-client.php", "echo", "pipes", "200", "0")) @@ -476,7 +477,7 @@ func Test_Pipe_SpawnTimeout_ReapsWorker(t *testing.T) { require.Nil(t, w) assert.Eventually(t, func() bool { - for p := range childPids(t) { + for p := range maps.Keys(childPids(t)) { if _, ok := before[p]; !ok { return false } diff --git a/worker_watcher/worker_watcher.go b/worker_watcher/worker_watcher.go index d169673..4e14b05 100644 --- a/worker_watcher/worker_watcher.go +++ b/worker_watcher/worker_watcher.go @@ -36,7 +36,7 @@ type WorkerWatcher struct { allocator Allocator allocateTimeout time.Duration stopCh chan struct{} - stopOnce sync.Once + stop func() destroyed atomic.Bool } @@ -52,6 +52,9 @@ func NewSyncWorkerWatcher(allocator Allocator, log *slog.Logger, numWorkers uint allocator: allocator, stopCh: make(chan struct{}), } + ww.stop = sync.OnceFunc(func() { + close(ww.stopCh) + }) ww.numWorkers.Store(numWorkers) return ww @@ -295,9 +298,7 @@ func (ww *WorkerWatcher) Destroy(ctx context.Context) { if ww.destroyed.Load() { return } - ww.stopOnce.Do(func() { - close(ww.stopCh) - }) + ww.stop() ww.mu.Lock() // do not release new workers ww.container.Destroy() diff --git a/worker_watcher/worker_watcher_test.go b/worker_watcher/worker_watcher_test.go index 1ebac5a..f79da03 100644 --- a/worker_watcher/worker_watcher_test.go +++ b/worker_watcher/worker_watcher_test.go @@ -80,9 +80,7 @@ func createStartedWorker(t *testing.T, state int64) *worker.Process { // shutdownWatcher signals the watcher to stop (without calling Destroy which needs relay). func shutdownWatcher(ww *WorkerWatcher) { ww.container.Destroy() - ww.stopOnce.Do(func() { - close(ww.stopCh) - }) + ww.stop() } // TestWorkerWatcher_AllocateRetryTimeout verifies that Allocate returns a WorkerAllocate error