risingwavelabs/risingwave · error

historical read semaphore closed

Error message

historical read semaphore closed

What it means

KvLogStoreReader::start_from acquires a permit from a historical-read semaphore (limiting concurrent state-store reads). Acquire fails only when the semaphore is closed, meaning the semaphore's owner was dropped. The reader cannot proceed to read persisted log data from the state store.

Solutions

  1. Ensure the semaphore owner lives at least as long as every reader using it (hold an Arc in the reader).
  2. Check for an early drop of the reader/factory context during recovery or migration.
  3. Treat the error as terminal: recreate reader and semaphore pair instead of retrying acquire.
  4. In tests, keep the semaphore owner alive for the whole read duration.

Example fix

// before: semaphore owned by a dropped context
let reader = KvLogStoreReader::new(ctx /* ctx owns semaphore */);
drop(ctx);
reader.start_from(...).await?;

// after: reader holds its own Arc<Semaphore>
let semaphore = Arc::new(Semaphore::new(n));
let reader = KvLogStoreReader::with_semaphore(semaphore.clone());
Defensive patterns

Strategy: try-catch

Validate before calling

if semaphore.is_closed() {
    return Err("historical read semaphore already closed; recreate reader");
}

Type guard

fn semaphore_live(s: &Arc<Semaphore>) -> bool { Arc::strong_count(s) > 1 } // owner still around

Try / catch

match reader.start_from(...).await {
    Err(e) if e.to_string().contains("historical read semaphore closed") => {
        // recreate reader+semaphore pair
        rebuild_reader_and_retry()
    }
    r => r,
}

Prevention

When it happens

Trigger: Reader's start_from awaiting `semaphore.acquire_owned()` after the semaphore has been closed/dropped — its owning struct (e.g. the shared reader context or the actor that created it) is gone.

Common situations: Reader outliving its owning context (e.g. used after the streaming actor or factory was dropped); mis-managed Arc where the only semaphore owner was dropped during recovery; long-lived reader used in tests after teardown.

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/243fc060bbbc0eab. Report an issue: GitHub.

Appendix: source

Thrown at src/stream/src/common/log_store_impl/kv_log_store/reader.rs:310

                // This limits the number of concurrent historical reads across
                // the compute node, preventing excessive state store I/O during
                // recovery or restart.
                //
                // Deadlock safety: The writer side of this log store runs as a
                // separate future in a `select` with this consumer. Even if this
                // consumer is blocked waiting for a permit, the writer can still
                // make progress processing barriers. The log store buffer is also
                // unbounded on the writer side, so the writer won't be blocked by
                // the reader not consuming. Therefore, blocking here does not
                // prevent barrier flow.
                let permit = if let Some(semaphore) = &self.historical_read_semaphore {
                    Some(
                        semaphore
                            .clone()
                            .acquire_owned()
                            .instrument_await("Wait for Historical Read Permit")
                            .await
                            .map_err(|_| anyhow!("historical read semaphore closed"))?,
                    )
                } else {
                    None
                };
                KvLogStoreReaderFutureState::ReadStateStoreStream(
                    self.read_persisted_log_store(range_start).await?,
                    permit,
                )
            };
        self.rx.rewind(start_offset);
        Ok(())
    }

    async fn init(&mut self) -> LogStoreResult<()> {
        if let Some(init_epoch_rx) = self.init_epoch_rx.take() {
            let init_epoch = init_epoch_rx
                .await
                .map_err(|_| anyhow!("should get the first epoch"))?;

View on GitHub (pinned to 6469eb736d)