Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions ipc/pipe/pipe_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package pipe
import (
"context"
"log/slog"
"maps"
"os"
"os/exec"
"runtime"
Expand Down Expand Up @@ -468,15 +469,15 @@ 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"))
require.Error(t, err)
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
}
Expand Down
9 changes: 5 additions & 4 deletions worker_watcher/worker_watcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ type WorkerWatcher struct {
allocator Allocator
allocateTimeout time.Duration
stopCh chan struct{}
stopOnce sync.Once
stop func()
destroyed atomic.Bool
}

Expand All @@ -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
Expand Down Expand Up @@ -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()
Expand Down
4 changes: 1 addition & 3 deletions worker_watcher/worker_watcher_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading