{"record":{"id":"cbae13c48c588605","repo":"risingwavelabs/risingwave","slug":"notification-stopped-or-uninitialized","errorCode":null,"errorMessage":"notification stopped or uninitialized","messagePattern":"notification stopped or uninitialized","errorType":"panic","errorClass":null,"httpStatus":null,"severity":"critical","filePath":"src/meta/src/manager/metadata.rs","lineNumber":123,"sourceCode":"                    is_streaming.then_some((node.id, node))\n                })\n                .collect(),\n            rx,\n            meta_manager: Some(meta_manager),\n        })\n    }\n\n    pub(crate) fn current(&self) -> &HashMap<WorkerId, WorkerNode> {\n        &self.worker_nodes\n    }\n\n    pub(crate) async fn changed(&mut self) -> ActiveStreamingWorkerChange {\n        loop {\n            let notification = self\n                .rx\n                .recv()\n                .await\n                .expect(\"notification stopped or uninitialized\");\n            fn is_target_worker_node(worker: &WorkerNode) -> bool {\n                worker.r#type == WorkerType::ComputeNode as i32\n                    && worker.property.as_ref().unwrap().is_streaming\n            }\n            match notification {\n                LocalNotification::WorkerNodeDeleted(worker) => {\n                    let is_target_worker_node = is_target_worker_node(&worker);\n                    let Some(prev_worker) = self.worker_nodes.remove(&worker.id) else {\n                        if is_target_worker_node {\n                            warn!(\n                                ?worker,\n                                \"notify to delete an non-existing streaming compute worker\"\n                            );\n                        }\n                        continue;\n                    };\n                    if !is_target_worker_node {\n                        warn!(","sourceCodeStart":105,"sourceCodeEnd":141,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/manager/metadata.rs#L105-L141","documentation":"MetadataManager::changed awaits the next notification from its receiver with .expect(...). The panic fires when the receiver returns None — all notification senders were dropped (notification service shut down) or the channel was never initialized. This is a fatal meta-side invariant: the manager cannot observe worker changes without the notification stream.","triggerScenarios":"changed() polls while the meta notification service is shutting down, after the sender half is dropped, or when MetadataManager was constructed without wiring the notification channel (uninitialized rx).","commonSituations":"Meta node shutdown/failover while background tasks still call changed(); a startup ordering bug where the manager starts before the notification service; test harnesses constructing MetadataManager without a notification sender.","solutions":["If seen during meta shutdown, ignore it — the panic accompanies process exit.","Fix startup ordering: create the notification service and hold its sender before MetadataManager spawns changed()-driven tasks.","If the meta node keeps running, capture the panic backtrace and check for early drops of the notification sender; fix ownership so the sender outlives the manager.","Replace the expect with a graceful error/log if the manager should tolerate service shutdown, and report persistent occurrences as a bug."],"exampleFix":"// before\nlet notification = self.rx.recv().await.expect(\"notification stopped or uninitialized\");\n// after\nlet notification = self.rx.recv().await.ok_or_else(||\n    anyhow!(\"metadata notification stream closed; notification service stopped or uninitialized\")\n)?;","handlingStrategy":"try-catch","validationCode":"// ensure the notification channel is live before entering the manager loop\nassert!(!rx.is_closed(), \"notification channel closed before MetadataManager start\");","typeGuard":null,"tryCatchPattern":"// supervisor around the manager loop\nwhile let Ok(()) = shutdown_rx.changed().await {\n    match metadata_manager.changed().await {\n        Ok(change) => handle(change).await,\n        Err(_) => {\n            // notification service stopped; exit loop gracefully instead of panicking\n            break;\n        }\n    }\n}","preventionTips":["Hold the notification sender for the lifetime of MetadataManager.","Initialize the notification service before spawning manager-dependent tasks.","Prefer ok_or_else/Err over expect() for long-running service loops.","Treat panics here during shutdown as benign; only investigate if the meta node keeps running."],"tags":["panic","meta","notification"],"backgroundTag":"internal-invariant-violation","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"}