Repository navigation
fix: Goroutine leaks and data races in CLI sync internals - #23408
goingforstudying-ctrl wants to merge 2 commits into
Conversation
The transformer pipeline leaked goroutines on every close: Close() only set a flag, so the head wrapper's Recv goroutine stayed parked in client.Recv() forever (nothing ever closed its stream), and the Send and Recv helper goroutines could block permanently on unbuffered channels once their caller had returned. Close() now half-closes the head stream so EOF cascades through the chain, the helper channels are buffered, and the Recv goroutine exits via a done channel once startBlocking returns. The main loop drains a buffered record before acting on close or EOF so the last in-flight record is not silently dropped. Also guard the analytics client and cached event details with RWMutexes. Sync progress reporting touches them from off the main goroutine, and the race detector flags InitClient vs Track* otherwise.
|
| // without that the Recv goroutine of the first wrapper stays parked in | ||
| // client.Recv() forever, because nothing else ever closes its stream. | ||
| if lp.clientWrappers[0].isClosed.CompareAndSwap(false, true) { | ||
| _ = lp.clientWrappers[0].client.CloseSend() |
There was a problem hiding this comment.
Closing identity stream can panic When a destination has no configured transformer, the pipeline uses an identity transformer whose
Send writes to an unbuffered channel. If an output error closes the pipeline while another record is waiting in Send, this new CloseSend closes that channel and the send panics, terminating the sync. The close-during-send test uses a different transformer and does not cover this case.
There was a problem hiding this comment.
Fixed in eca771e. The identity transformer no longer closes its data channel on CloseSend; closure is signalled on a separate channel, so a send parked on the unbuffered channel returns ErrPipelineClosed instead of panicking. Added a regression test that closes the pipeline mid-send on a transformer-less pipeline.
| select { | ||
| case req := <-recvCh: | ||
| if err := s.nextSendFn(req); err != nil { | ||
| return err | ||
| } | ||
| continue | ||
| default: | ||
| } |
There was a problem hiding this comment.
There was a problem hiding this comment.
Fixed in eca771e. Both close paths in startBlocking now drain recvCh before calling nextClose. The Recv goroutine buffers its last record before EOF, so the final record always makes it out.
| done := make(chan struct{}) | ||
| go func() { | ||
| defer close(done) | ||
| require.NoError(t, pipeline.Send([]byte("ping"))) |
There was a problem hiding this comment.
Send failure can hang test If
pipeline.Send fails in this spawned goroutine, require.NoError exits it before pipeline.Close() runs. The test then stays blocked in RunBlocking instead of reporting the send failure, making the regression test harder to diagnose. Return the send result to the test goroutine and ensure the pipeline closes on failure.
There was a problem hiding this comment.
Fixed in eca771e. The goroutine now sends its result on a channel and closes the pipeline via defer, so a send failure surfaces as a test failure instead of leaving RunBlocking blocked.
Identity transformer Send parked on an unbuffered channel panicked when CloseSend closed it underneath. Closure is now signalled on a separate channel, so a racing send returns ErrPipelineClosed instead. The receive loop could also pick EOF over a record the Recv goroutine had already buffered, silently dropping the final transformed record. Both close paths now drain a buffered record before closing the chain. Also stop the goroutine-leak test from hanging in RunBlocking when the send fails, and cover close-during-identity-send with a regression test.
Summary
Watching goroutine counts climb across repeated syncs with a transformer attached, I dug into the pipeline teardown and found that every
Close()leaves at least two goroutines behind:client.Recv()forever, becauseClose()only sets a flag and nothing anywhere closes the head of the stream chain,This PR:
Close()now half-closes the head transform stream (guarded by a CAS so it happens exactly once), so EOF cascades down the wrapper chain and everything shuts down promptly. Nice side effect: shutdown no longer rides out the 1-second tick, and the pipeline test suite here went from ~27s to ~3.5s.recvCh/errCh/sendChare buffered (cap 1) and the Recv goroutine gets a done channel, so nothing blocks on a send afterstartBlocking/Sendhas returned.startBlockingloop now drains a buffered record before acting on the tick or EOF. Without that, the last in-flight record can lose the select race and get silently dropped, which the existing identity-transformer test caught as soon as CloseSend made EOF arrive promptly.Also picked up the two globals in
cli/internal/analyticswhile I was in the area:clientandcachedSyncEventDetailsget read from whichever goroutine is reporting an event (sync progress reporting is async), butInitClient/refreshSyncEventDetailswrite them bare. The race detector flags it on the first concurrent call, so both are behind RWMutexes now.Fixes #23165
Fixes #23168
Regression tests: goroutine count returns to baseline after repeated open/send/close cycles, and after a close landing while a Send is parked inside the transformer, plus a
-racetest runningInitClientagainst concurrentTrack*calls.golangci-lint runis clean on the touched packages 🚨go test -race -count=2 ./internal/transformerpipeline/ ./internal/analytics/passes 🧪