kvstreamer: lost wakeup in Streamer.Close leaves flow cleanup hung forever
Describe the problem
Streamer.Close can deadlock against the streamer's own worker coordinator. When it does, the enclosing flow's Cleanup never finishes: the flow's memory monitors are never stopped and UnregisterFlow never runs. On a node being drained this shows up as drain details: distSQL execution flows: 1 on every iteration until the operator kills the process.
The race is in the condition variable used for budget back-pressure:
Closecancels the coordinator context, setsdone, and then callss.budget.mu.waitForBudget.Signal()without holdingbudget.mu(streamer.go:851), then blocks ins.waitGroup.Wait().workerCoordinator.waitUntilEnoughBudgetholdsbudget.mufor the whole iteration (spill attempt, in-progress check) and only checksctx.Err()afterWait()returns (streamer.go:1070).
If the Signal fires while the coordinator is inside that critical section or is about to enter it, the Signal is lost and the coordinator parks. Nothing wakes it: the only other signaler is budget.release, which needs the consumer to release results, and the consumer is the operator being closed. Close then waits on the coordinator's waitGroup forever.
Observed in production (v25.2.17, vectorized remote flow, ColIndexJoin with the streamer, node at its --max-sql-memory limit so the streamer budget was exhausted). Goroutine dump taken 93 minutes into the hang:
goroutine A [semacquire, 93 minutes]:
sync.(*WaitGroup).Wait
kvstreamer.(*Streamer).Close streamer.go:853
row.(*txnKVStreamer).Close kv_batch_streamer.go:275
colfetcher.(*cFetcher).Close cfetcher.go:1484
colfetcher.(*ColIndexJoin).Close index_join.go:713
colexecargs.(*CloserRegistry).Close closer_registry.go:45
colflow.(*vectorizedFlow).Cleanup vectorized_flow.go:385
flowinfra.(*RemoteFlowRunner).RunFlow.func1.2 remote_flow_runner.go:102
goroutine B [sync.Cond.Wait, 93 minutes]:
sync.(*Cond).Wait
kvstreamer.(*workerCoordinator).waitUntilEnoughBudget streamer.go:1070
kvstreamer.(*workerCoordinator).mainLoop streamer.go:954Side effects visible from outside: sql.mem.distsql.current held a constant ~127 MiB (the flow's unreleased budget) with zero connections and sql.distsql.flows.active = 0; the flow was invisible in crdb_internal.node_distsql_flows because the runner drops it before Cleanup.
To Reproduce
No deterministic repro yet. Precondition is a streamer whose budget is exhausted (limitBytes - used < atLeastBytes) at the moment the flow is torn down. A test could block the coordinator inside waitUntilEnoughBudget before Wait() with a testing knob, call Close, then release it.
Expected behavior
Close returns once the coordinator observes cancellation. Standard cond-var discipline: signal while holding budget.mu (or Broadcast), and re-check ctx.Err() before parking in waitUntilEnoughBudget.
Environment
v25.2.17; the same code is present on master.
Jira issue: CRDB-68488
Source: cockroachdb/cockroach