risingwavelabs/risingwave · error
unable to subscribe resume
Error message
unable to subscribe resume
What it means
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.
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.
Example fix
// before: controller dropped while reader paused let (ctl, reader) = build_paused_pair(); drop(ctl); reader.next_item().await?; // stuck paused, channel closed // after: resume before dropping the controller ctl.resume(); drop(ctl); reader.next_item().await?;
Defensive patterns
Strategy: try-catch
Validate before calling
if pause_ctl.is_dropped() && reader.is_paused() {
// cannot ever resume; stop the reader
stop_reader();
} Try / catch
match reader.next_item().await {
Err(e) if e.to_string().contains("unable to subscribe resume") => {
// pause owner gone while paused; fail reader for rebuild
fail_reader(e);
}
r => r,
} Prevention
- 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.
When it happens
Trigger: Reader paused (is_paused=true) while all `is_paused` watch senders are dropped, then `next_item` calls `.changed()` which fails with a Closed error.
Common situations: 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.
Related errors
- cannot get truncated epoch
- failed to receive update vnode
- historical read semaphore closed
- should get the first epoch
- unable to send barrier
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/80da15ac1457ec93.
Report an issue: GitHub.
Appendix: source
Thrown at src/stream/src/common/log_store_impl/kv_log_store/reader.rs:361
self.first_write_epoch = Some(write_epoch);
};
self.future_state = KvLogStoreReaderFutureState::Reset;
self.latest_offset = None;
self.truncate_offset = None;
self.rewind_delay = RewindDelay::new(&self.metrics);
Ok(())
}
async fn next_item(&mut self) -> LogStoreResult<(u64, LogStoreReadItem)> {
while *self.is_paused.borrow_and_update() {
info!("next_item of {} get blocked by is_pause", self.identity);
self.is_paused
.changed()
.instrument_await("Wait for Pause Resume")
.await
.map_err(|_| anyhow!("unable to subscribe resume"))?;
}
match &mut self.future_state {
KvLogStoreReaderFutureState::ReadStateStoreStream(state_store_stream, _permit) => {
match state_store_stream
.try_next()
.instrument_await("Try Next for Historical Stream")
.await?
{
Some((epoch, item)) => {
if let Some(latest_offset) = &self.latest_offset {
latest_offset.check_next_item_epoch(epoch)?;
}
let item = match item {
KvLogStoreItem::StreamChunk { chunk, .. } => {
let chunk_id = if let Some(latest_offset) = self.latest_offset {
latest_offset.next_chunk_id()
} else {
0View on GitHub (pinned to 6469eb736d)