Visitar URL original
fix: Goroutine leaks and data races in CLI sync internals by goingforstudying-ctrl · Pull Request #23408 · cloudquery/cloudquery · GitHub
Skip to content

fix: Goroutine leaks and data races in CLI sync internals - #23408

Open
goingforstudying-ctrl wants to merge 2 commits into
cloudquery:mainfrom
goingforstudying-ctrl:fix/cli-pipeline-analytics-concurrency
Open

goingforstudying-ctrl wants to merge 2 commits into
cloudquery:mainfrom
goingforstudying-ctrl:fix/cli-pipeline-analytics-concurrency

Conversation

@goingforstudying-ctrl

Copy link
Copy Markdown

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:

  • the first wrapper's Recv goroutine parks in client.Recv() forever, because Close() only sets a flag and nothing anywhere closes the head of the stream chain,
  • the Send/Recv helper goroutines can block permanently on their unbuffered result channels once the caller on the other side has already returned.

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/sendCh are buffered (cap 1) and the Recv goroutine gets a done channel, so nothing blocks on a send after startBlocking/Send has returned.
  • The startBlocking loop 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/analytics while I was in the area: client and cachedSyncEventDetails get read from whichever goroutine is reporting an event (sync progress reporting is async), but InitClient/refreshSyncEventDetails write 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 -race test running InitClient against concurrent Track* calls.

  • Read the contribution guidelines 🧑‍🎓
  • golangci-lint run is clean on the touched packages 🚨
  • go test -race -count=2 ./internal/transformerpipeline/ ./internal/analytics/ passes 🧪

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.
@greptile-apps

greptile-apps Bot commented Sep 26, 2026 •

Copy link
Copy Markdown

RetriggerConfidence Score: 3/5

[High risk] Fixes concurrency bugs in CLI sync and transformer pipeline.

The PR is not safe to merge until the transformer-less close panic and EOF record-loss race are fixed.

Fix All in Claude CodeFindings

  1. P1 Closing identity stream can panic ▶
  2. P1 EOF can drop final record ▶
  3. P2 Send failure can hang test ▶
Summary

The PR half-closes the head transform stream during pipeline teardown, buffers asynchronous send and receive results, and guards analytics globals with mutexes.

  • The new head close can panic for transformer-less destinations during an in-flight send.
  • The buffered receive path can still discard a record when EOF arrives.
  • A new leak test can hang rather than report a send failure.
Diagram
%%{init: {'theme': 'neutral'}}%%
flowchart LR
  Source[Source Send] --> Head[Head transform stream]
  Head --> Wrapper[Receive wrapper]
  Wrapper --> Output[Destination output]
  Close[Pipeline Close] -->|CloseSend| Head
  Head -->|record or EOF| Wrapper
Loading

Reviews (1) · Last reviewed commit: "fix: close goroutine leaks on pipeline s..."

// 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()

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 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.

Fix in Claude Code

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +178 to +185
select {
case req := <-recvCh:
if err := s.nextSendFn(req); err != nil {
return err
}
continue
default:
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 EOF can drop final record The receive goroutine can buffer a record and then its EOF. If the record arrives just after this preliminary drain checks the channel, the following select can choose EOF over the record and return. The final transformed record is then silently lost.

Fix in Claude Code

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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")))

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 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.

Fix in Claude Code

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

cli/internal/analytics/client.go: data races on global client and cachedSyncEventDetails — no mutex protection cli/internal/transformerpipeline/pipeline.go: goroutine leaks on pipeline close

1 participant