{"record":{"id":"272ba5afc472b21d","repo":"risingwavelabs/risingwave","slug":"rate-limit-control-channel-closed","errorCode":null,"errorMessage":"rate limit control channel closed","messagePattern":"rate limit control channel closed","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/log_store.rs","lineNumber":611,"sourceCode":"                Ok((epoch, LogStoreReadItem::StreamChunk { chunk, chunk_id }))\n            }\n        }\n    }\n}\n\nimpl<R: LogReader> LogReader for RateLimitedLogReader<R> {\n    async fn init(&mut self) -> LogStoreResult<()> {\n        self.core.inner.init().await\n    }\n\n    async fn next_item(&mut self) -> LogStoreResult<(u64, LogStoreReadItem)> {\n        loop {\n            select! {\n                biased;\n                recv = pin!(self.control_rx.recv()) => {\n                    let new_rate_limit = match recv {\n                        Some(limit) => limit,\n                        None => bail!(\"rate limit control channel closed\"),\n                    };\n                    let old_rate_limit = self.core.rate_limiter.update(new_rate_limit);\n                    let paused = matches!(new_rate_limit, RateLimit::Pause);\n                    tracing::info!(\"rate limit changed from {:?} to {:?}, paused = {paused}\", old_rate_limit, new_rate_limit);\n                },\n                item = self.core.next_item() => {\n                    return item;\n                }\n            }\n        }\n    }\n\n    fn truncate(&mut self, offset: TruncateOffset) -> LogStoreResult<()> {\n        let downstream_offset = DownstreamChunkOffset(offset);\n        let mut truncate_offset = None;\n        let mut stop = false;\n        'outer: while let Some((upstream_offset, downstream_offsets)) =\n            self.core.consumed_offset_queue.back_mut()","sourceCodeStart":593,"sourceCodeEnd":629,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/log_store.rs#L593-L629","documentation":"The rate-limited log store writer's select loop receives on a control channel that delivers `RateLimit` updates; when the sender side is dropped the `recv()` yields `None` and the loop bails with this error. It signals the rate-limit control plane is gone, so the writer can no longer be reconfigured and shuts down.","triggerScenarios":"The task/component holding `control_tx` (the rate limit manager) is dropped or terminated while `LogStore` is still running its main loop; `self.control_rx.recv()` returns `None`.","commonSituations":"Meta node dropping the sink rate-limit control stream; actor shutdown ordering where the manager is cancelled before the sink; network failure tearing down the control gRPC stream.","solutions":["Keep the control channel sender alive for the lifetime of the sink writer","Check why the rate-limit manager task exited (logs, cancellation) and fix its shutdown ordering","If intentional shutdown, stop the log store loop first, or treat `None` as a graceful stop rather than an error"],"exampleFix":"// before\nlet (tx, rx) = mpsc::unbounded_channel();\nspawn_sink_writer(rx);\ndrop(tx); // causes 'rate limit control channel closed'\n// after\nlet (tx, rx) = mpsc::unbounded_channel();\nlet writer = spawn_sink_writer(rx);\n// keep tx alive until writer finishes\nwriter.await;\ndrop(tx);","handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"match log_store_result {\n    Err(e) if e.to_string().contains(\"rate limit control channel closed\") => {\n        info!(\"rate limit manager gone; stopping sink writer gracefully\");\n    }\n    other => other?,\n}","preventionTips":["Hold the control_tx handle for the full lifetime of the sink writer task","Use a shutdown signal instead of dropping the sender to stop the loop","Review task cancellation order so the rate-limit manager outlives the sink writer"],"tags":["rust","channel","rate-limit"],"backgroundTag":"broken-pipe","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}