{"record":{"id":"9e1f88838168310d","repo":"risingwavelabs/risingwave","slug":"expect-pbnodebody-project-but-got","errorCode":null,"errorMessage":"expect PbNodeBody::Project but got: {:?}","messagePattern":"expect PbNodeBody::Project but got: (.+?)","errorType":"error_code","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/meta/src/stream/stream_graph/fragment.rs","lineNumber":685,"sourceCode":"    merge.upstream_fragment_id = upstream_table_fragment_id;\n\n    Ok(ScanRewriteResult {\n        old_output_index_to_new_output_index,\n        new_output_index_by_column_id,\n        output_fields: stream_scan_node.fields.clone(),\n    })\n}\n\n/// Rewrite Project node input refs and extend with newly added columns.\nfn rewrite_project_node(\n    project_node: &mut StreamNode,\n    scan_rewrite: &ScanRewriteResult,\n    newly_added_columns: &[ColumnCatalog],\n    removed_column_ids: &HashSet<ColumnId>,\n    upstream_table_name: &str,\n) -> MetaResult<()> {\n    let PbNodeBody::Project(project_node_body) = project_node.node_body.as_mut().unwrap() else {\n        return Err(anyhow!(\n            \"expect PbNodeBody::Project but got: {:?}\",\n            project_node.node_body\n        )\n        .into());\n    };\n    let has_non_input_ref = project_node_body\n        .select_list\n        .iter()\n        .any(|expr| !matches!(expr.rex_node, Some(expr_node::RexNode::InputRef(_))));\n    if has_non_input_ref && !removed_column_ids.is_empty() {\n        return Err(anyhow!(\n            \"auto schema change with drop column only supports Project with InputRef\"\n        )\n        .into());\n    }\n\n    let mut new_select_list = Vec::with_capacity(project_node_body.select_list.len());\n    let mut new_project_fields = Vec::with_capacity(project_node.fields.len());","sourceCodeStart":667,"sourceCodeEnd":703,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/stream/stream_graph/fragment.rs#L667-L703","documentation":"When rewriting the Project node inserted for auto schema change (`rewrite_project_node`), the code requires the node's body to be `PbNodeBody::Project` so its expressions can be adjusted for added/removed columns. Any other variant returns this error, failing the schema refresh rewrite.","triggerScenarios":"`rewrite_refresh_schema_sink_fragment` -> `rewrite_project_node` called with a node whose `node_body` is not `PbNodeBody::Project` (e.g. it was resolved via a Project branch but the body is actually something else).","commonSituations":"Upstream table schema change (ALTER TABLE ADD/DROP COLUMN, schema registry evolution) triggering a rewrite against inconsistent fragment metadata; frontend/meta version drift; internal invariant breach between the check and rewrite passes.","solutions":["Restart meta and retry the schema-changing operation to rebuild consistent graph state.","Verify all RisingWave components are on the same version.","Drop and recreate the sink to regenerate a canonical graph with a proper Project node.","File a RisingWave bug with fragment graph dumps — the check pass should have matched this node as Project."],"exampleFix":"// before\nlet PbNodeBody::Project(project_node_body) = project_node.node_body.as_mut().unwrap() else { ... };\n// after: skip projection rewrite when shape is unexpected\nlet Some(PbNodeBody::Project(project_node_body)) = project_node.node_body.as_mut() else {\n    tracing::warn!(\"expected Project node for schema refresh, got {:?}; skipping\", project_node.node_body);\n    return Ok(());\n};","handlingStrategy":"type-guard","validationCode":"if !matches!(project_node.node_body.as_ref(), Some(PbNodeBody::Project(_))) {\n    return Err(\"schema refresh expects a Project node here\".into());\n}","typeGuard":"fn as_project_mut(node: &mut StreamNode) -> Option<&mut PbProjectNode> {\n    match node.node_body.as_mut() {\n        Some(PbNodeBody::Project(p)) => Some(p),\n        _ => None,\n    }\n}","tryCatchPattern":"match rewrite_project_node(&mut project_node, ...).await {\n    Ok(()) => {},\n    Err(e) if e.to_string().contains(\"expect PbNodeBody::Project\") => {\n        // invariant breach: rebuild the sink graph\n        rebuild_sink_graph()?;\n    }\n    Err(e) => return Err(e),\n}","preventionTips":["Ensure the node resolved as 'the project node' is actually a Project before rewriting.","Keep check-pass and rewrite-pass shape assumptions in shared helper functions.","Exercise ALTER TABLE ADD/DROP COLUMN against sinks in staging regularly."],"tags":["risingwave","streaming-graph","schema-change","internal-invariant"],"backgroundTag":"unexpected-response-shape","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"}