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

  1. Ensure the pause/watch sender owner outlives all readers, or drop the reader when the controller is dropped.
  2. Check whether the reader was left paused by a cancelled control path; unset is_paused before teardown.
  3. Treat this as terminal for the reader; let the actor fail and rebuild via recovery.
  4. 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

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


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 {
                                    0

View on GitHub (pinned to 6469eb736d)