{"record":{"id":"80da15ac1457ec93","repo":"risingwavelabs/risingwave","slug":"unable-to-subscribe-resume","errorCode":null,"errorMessage":"unable to subscribe resume","messagePattern":"unable to subscribe resume","errorType":"error_code","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/common/log_store_impl/kv_log_store/reader.rs","lineNumber":361,"sourceCode":"            self.first_write_epoch = Some(write_epoch);\n        };\n\n        self.future_state = KvLogStoreReaderFutureState::Reset;\n        self.latest_offset = None;\n        self.truncate_offset = None;\n        self.rewind_delay = RewindDelay::new(&self.metrics);\n\n        Ok(())\n    }\n\n    async fn next_item(&mut self) -> LogStoreResult<(u64, LogStoreReadItem)> {\n        while *self.is_paused.borrow_and_update() {\n            info!(\"next_item of {} get blocked by is_pause\", self.identity);\n            self.is_paused\n                .changed()\n                .instrument_await(\"Wait for Pause Resume\")\n                .await\n                .map_err(|_| anyhow!(\"unable to subscribe resume\"))?;\n        }\n        match &mut self.future_state {\n            KvLogStoreReaderFutureState::ReadStateStoreStream(state_store_stream, _permit) => {\n                match state_store_stream\n                    .try_next()\n                    .instrument_await(\"Try Next for Historical Stream\")\n                    .await?\n                {\n                    Some((epoch, item)) => {\n                        if let Some(latest_offset) = &self.latest_offset {\n                            latest_offset.check_next_item_epoch(epoch)?;\n                        }\n                        let item = match item {\n                            KvLogStoreItem::StreamChunk { chunk, .. } => {\n                                let chunk_id = if let Some(latest_offset) = self.latest_offset {\n                                    latest_offset.next_chunk_id()\n                                } else {\n                                    0","sourceCodeStart":343,"sourceCodeEnd":379,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/common/log_store_impl/kv_log_store/reader.rs#L343-L379","documentation":"next_item blocks while the reader is paused, awaiting a change notification on the shared `is_paused` watch channel. This error is returned when the watch channel has been closed — every sender was dropped — so the reader can never observe a resume. The reader is stuck permanently behind a pause flag that nobody can clear.","triggerScenarios":"Reader paused (is_paused=true) while all `is_paused` watch senders are dropped, then `next_item` calls `.changed()` which fails with a Closed error.","commonSituations":"Pause-controller owner dropped after failure/cancellation while the reader remained paused; recovery racing where the new owner is created but the old reader lingers; pause flag left true by an interrupted path.","solutions":["Ensure the pause/watch sender owner outlives all readers, or drop the reader when the controller is dropped.","Check whether the reader was left paused by a cancelled control path; unset is_paused before teardown.","Treat this as terminal for the reader; let the actor fail and rebuild via recovery.","In tests, always resume (set is_paused=false) or drop the reader before dropping the controller."],"exampleFix":"// before: controller dropped while reader paused\nlet (ctl, reader) = build_paused_pair();\ndrop(ctl);\nreader.next_item().await?; // stuck paused, channel closed\n\n// after: resume before dropping the controller\nctl.resume();\ndrop(ctl);\nreader.next_item().await?;","handlingStrategy":"try-catch","validationCode":"if pause_ctl.is_dropped() && reader.is_paused() {\n    // cannot ever resume; stop the reader\n    stop_reader();\n}","typeGuard":null,"tryCatchPattern":"match reader.next_item().await {\n    Err(e) if e.to_string().contains(\"unable to subscribe resume\") => {\n        // pause owner gone while paused; fail reader for rebuild\n        fail_reader(e);\n    }\n    r => r,\n}","preventionTips":["Always resume or drop the reader before dropping the pause controller.","Keep the watch sender's owner alive as long as readers exist.","Audit pause/resume paths so cancellation never leaves is_paused=true with no owner."],"tags":["watch-channel","pause","channel-closed","log-store"],"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"}