{"record":{"id":"243fc060bbbc0eab","repo":"risingwavelabs/risingwave","slug":"historical-read-semaphore-closed","errorCode":null,"errorMessage":"historical read semaphore closed","messagePattern":"historical read semaphore closed","errorType":"error_code","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/common/log_store_impl/kv_log_store/reader.rs","lineNumber":310,"sourceCode":"                // This limits the number of concurrent historical reads across\n                // the compute node, preventing excessive state store I/O during\n                // recovery or restart.\n                //\n                // Deadlock safety: The writer side of this log store runs as a\n                // separate future in a `select` with this consumer. Even if this\n                // consumer is blocked waiting for a permit, the writer can still\n                // make progress processing barriers. The log store buffer is also\n                // unbounded on the writer side, so the writer won't be blocked by\n                // the reader not consuming. Therefore, blocking here does not\n                // prevent barrier flow.\n                let permit = if let Some(semaphore) = &self.historical_read_semaphore {\n                    Some(\n                        semaphore\n                            .clone()\n                            .acquire_owned()\n                            .instrument_await(\"Wait for Historical Read Permit\")\n                            .await\n                            .map_err(|_| anyhow!(\"historical read semaphore closed\"))?,\n                    )\n                } else {\n                    None\n                };\n                KvLogStoreReaderFutureState::ReadStateStoreStream(\n                    self.read_persisted_log_store(range_start).await?,\n                    permit,\n                )\n            };\n        self.rx.rewind(start_offset);\n        Ok(())\n    }\n\n    async fn init(&mut self) -> LogStoreResult<()> {\n        if let Some(init_epoch_rx) = self.init_epoch_rx.take() {\n            let init_epoch = init_epoch_rx\n                .await\n                .map_err(|_| anyhow!(\"should get the first epoch\"))?;","sourceCodeStart":292,"sourceCodeEnd":328,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/common/log_store_impl/kv_log_store/reader.rs#L292-L328","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","solutions":["Ensure the semaphore owner lives at least as long as every reader using it (hold an Arc in the reader).","Check for an early drop of the reader/factory context during recovery or migration.","Treat the error as terminal: recreate reader and semaphore pair instead of retrying acquire.","In tests, keep the semaphore owner alive for the whole read duration."],"exampleFix":"// before: semaphore owned by a dropped context\nlet reader = KvLogStoreReader::new(ctx /* ctx owns semaphore */);\ndrop(ctx);\nreader.start_from(...).await?;\n\n// after: reader holds its own Arc<Semaphore>\nlet semaphore = Arc::new(Semaphore::new(n));\nlet reader = KvLogStoreReader::with_semaphore(semaphore.clone());","handlingStrategy":"try-catch","validationCode":"if semaphore.is_closed() {\n    return Err(\"historical read semaphore already closed; recreate reader\");\n}","typeGuard":"fn semaphore_live(s: &Arc<Semaphore>) -> bool { Arc::strong_count(s) > 1 } // owner still around","tryCatchPattern":"match reader.start_from(...).await {\n    Err(e) if e.to_string().contains(\"historical read semaphore closed\") => {\n        // recreate reader+semaphore pair\n        rebuild_reader_and_retry()\n    }\n    r => r,\n}","preventionTips":["Store the semaphore in an Arc owned by the reader itself, not only by an outer context.","Keep the owning context alive for the reader's full lifetime.","Detect early teardown of the reader/factory during recovery."],"tags":["semaphore","channel-closed","state-store","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"}