{"record":{"id":"d27379b77274050e","repo":"risingwavelabs/risingwave","slug":"anyhow-message-to-owned","errorCode":null,"errorMessage":"anyhow!(message.to_owned())","messagePattern":"anyhow!\\(message\\.to_owned\\(\\)\\)","errorType":"error_code","errorClass":"MetaError","httpStatus":null,"severity":"error","filePath":"src/meta/src/barrier/rpc.rs","lineNumber":1495,"sourceCode":"            .start_streaming_control(PbInitRequest::default())\n            .await?;\n        Ok(handle)\n    }\n}\n\npub(super) fn merge_node_rpc_errors<E: Error + Send + Sync + 'static>(\n    message: &str,\n    errors: impl IntoIterator<Item = (WorkerId, E)>,\n) -> MetaError {\n    use std::fmt::Write;\n\n    use risingwave_common::error::error_request_copy;\n    use risingwave_common::error::tonic::extra::Score;\n\n    let errors = errors.into_iter().collect_vec();\n\n    if errors.is_empty() {\n        return anyhow!(message.to_owned()).into();\n    }\n\n    // Create the error from the single error.\n    let single_error = |(worker_id, e)| {\n        anyhow::Error::from(e)\n            .context(format!(\"{message}, in worker node {worker_id}\"))\n            .into()\n    };\n\n    if errors.len() == 1 {\n        return single_error(errors.into_iter().next().unwrap());\n    }\n\n    // Find the error with the highest score.\n    let max_score = errors\n        .iter()\n        .filter_map(|(_, e)| error_request_copy::<Score>(e))\n        .max();","sourceCodeStart":1477,"sourceCodeEnd":1513,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/barrier/rpc.rs#L1477-L1513","documentation":"merge_node_rpc_errors collapses a set of per-worker-node RPC errors into a single MetaError. When the caller passes an empty error list, the function has nothing to attribute to any worker, so it throws a bare anyhow error containing only the generic message. It exists so the meta service can return one consolidated barrier/RPC failure instead of leaking many per-node errors to the frontend.","triggerScenarios":"Calling merge_node_rpc_errors (barrier dispatch paths in the meta node) with an empty `errors` iterator — i.e. the fan-out to worker nodes produced zero errors yet the caller still took the error-merging branch, typically when no worker nodes were targeted or all error entries were filtered out before merging.","commonSituations":"Cluster where all streaming workers have been removed/failed but the barrier manager still attempts to dispatch; races between worker deregistration and barrier issuance; tests or internal callers invoking the merge helper without collected errors.","solutions":["Check cluster health: ensure streaming worker nodes are registered and alive before triggering barrier operations (`SHOW CLUSTERS` / meta node logs for worker liveness).","If operating during scale-down, wait for or trigger recovery so the barrier manager re-establishes its worker set.","Upgrade/Rollback recently: if this appeared after a version change, restart the cluster with recovery enabled so worker-node state is rebuilt."],"exampleFix":"// before: merging with an empty error list\nlet errors: Vec<_> = collected_errors.into_iter().filter(|e| !is_benign(e)).collect();\nlet merged = merge_node_rpc_errors(msg, errors);\n// after: skip merging entirely when nothing failed\nlet errors: Vec<_> = collected_errors.into_iter().filter(|e| !is_benign(e)).collect();\nlet merged = if errors.is_empty() { Ok(()) } else { merge_node_rpc_errors(msg, errors) };","handlingStrategy":"retry","validationCode":null,"typeGuard":null,"tryCatchPattern":"// Rust\nmatch result {\n    Err(e) if is_empty_aggregate_rpc_error(&e) => warn!(\"no workers targeted; retry after recovery\"),\n    Err(e) => return Err(e),\n    Ok(v) => Ok(v),\n}","preventionTips":["Monitor worker liveness and never let the cluster run with zero streaming workers.","Keep recovery enabled so worker loss triggers automatic re-registration.","Check `SHOW CLUSTERS;` before maintenance that could drain all compute nodes."],"tags":["rpc","meta","barrier","error-aggregation"],"backgroundTag":"rpc-error-aggregation","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-14T16:17:12.679Z"}