{"record":{"id":"4e1cc744a43351a5","repo":"risingwavelabs/risingwave","slug":"take-receiver-on-unmatched-partial-graph-term","errorCode":null,"errorMessage":"take receiver {:?} on unmatched partial graph term {} to current term {}","messagePattern":"take receiver (.+?) on unmatched partial graph term (.+?) to current term (.+?)","errorType":"exception","errorClass":"anyhow","httpStatus":null,"severity":"error","filePath":"src/stream/src/task/barrier_worker/mod.rs","lineNumber":535,"sourceCode":"                ids,\n                request,\n            } => {\n                let err = match self.state.partial_graphs.entry(partial_graph_id) {\n                    Entry::Occupied(mut entry) => match entry.get_mut() {\n                        PartialGraphStatus::ReceivedExchangeRequest(pending_requests) => {\n                            pending_requests.push((term_id, ids, request));\n                            return;\n                        }\n                        PartialGraphStatus::Running(graph) => {\n                            if graph.local_barrier_manager.term_id != term_id {\n                                warn!(\n                                    %partial_graph_id,\n                                    ?ids,\n                                    request_term_id = term_id,\n                                    current_term_id = graph.local_barrier_manager.term_id,\n                                    \"take receiver on unmatched partial graph term\"\n                                );\n                                anyhow!(\n                                    \"take receiver {:?} on unmatched partial graph term {} to current term {}\",\n                                    ids,\n                                    term_id,\n                                    graph.local_barrier_manager.term_id\n                                )\n                            } else {\n                                let (upstream_actor_id, actor_id) = ids;\n                                graph.new_actor_output_request(\n                                    actor_id,\n                                    upstream_actor_id,\n                                    request,\n                                );\n                                return;\n                            }\n                        }\n                        PartialGraphStatus::Suspended(_) => anyhow!(\"partial graph suspended\"),\n                        PartialGraphStatus::Resetting => anyhow!(\"partial graph resetting\"),\n                        PartialGraphStatus::Unspecified => unreachable!(),","sourceCodeStart":517,"sourceCodeEnd":553,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/task/barrier_worker/mod.rs#L517-L553","documentation":"An actor sent a TakeReceiver (exchange channel creation) request whose term (epoch) id does not match the term id of the currently running partial graph. The barrier worker rejects the request and reports the error back to the remote sender via the result_sender. It protects the graph from wiring exchange channels across generation boundaries.","triggerScenarios":"handle_actor_op receives LocalActorOperation::TakeReceiver while the partial graph is PartialGraphStatus::Running and graph.local_barrier_manager.term_id != term_id.","commonSituations":"A slow upstream actor's exchange request races with a barrier-driven graph rebuild; a partial graph reset occurred between the request being issued and processed; config/planner changes causing frequent term bumps under failure recovery.","solutions":["Retry the TakeReceiver request with the current term id after the graph reset completes.","Check logs for a preceding partial graph reset that changed the term id.","Verify upstream/downstream term id propagation in the streaming control requests.","If persistent, inspect barrier scheduling for term regression across meta and compute nodes."],"exampleFix":"// before: taking receiver with stale term\nworker.take_receiver(partial_graph_id, stale_term_id, ids, request);\n// after: ensure request carries the term id observed from the latest barrier\nlet term_id = graph.local_barrier_manager.term_id;\nworker.take_receiver(partial_graph_id, term_id, ids, request);","handlingStrategy":"retry","validationCode":"// caller-side check before sending TakeReceiver\nif request.term_id != current_graph_term_id(partial_graph_id) {\n    refresh_term_id_before_take_receiver();\n}","typeGuard":"fn is_term_mismatch(err: &anyhow::Error) -> bool { err.to_string().contains(\"unmatched partial graph term\") }","tryCatchPattern":"// treat rejection as retryable\nif let Err(e) = sender.await_reply() {\n    if is_term_mismatch(&e) { fetch_new_term_id().await; retry_take_receiver().await; }\n}","preventionTips":["Always source term_id from the latest barrier/control response, never cache it.","Subscribe to graph reset events and invalidate in-flight exchange requests on term change.","Log term ids on both sides to catch propagation bugs early."],"tags":["streaming","barrier","term-mismatch","rpc"],"backgroundTag":"invalid-state-transition","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"}