{"record":{"id":"2c194ad6ea264a32","repo":"risingwavelabs/risingwave","slug":"take-receiver-upstream-actor-id-actor-id-on","errorCode":null,"errorMessage":"take receiver ({upstream_actor_id}, {actor_id}) on stale partial graph term {request_term_id}; current term is {term_id}","messagePattern":"take receiver \\((.+?), (.+?)\\) on stale partial graph term (.+?); current term is (.+?)","errorType":"exception","errorClass":"anyhow","httpStatus":null,"severity":"warning","filePath":"src/stream/src/task/barrier_worker/mod.rs","lineNumber":965,"sourceCode":"                        term_id.clone(),\n                        self.actor_manager.clone(),\n                    );\n                    for (request_term_id, (upstream_actor_id, actor_id), request) in\n                        pending_requests.drain(..)\n                    {\n                        if request_term_id == term_id {\n                            graph.new_actor_output_request(actor_id, upstream_actor_id, request);\n                        } else {\n                            warn!(\n                                %partial_graph_id,\n                                %upstream_actor_id,\n                                %actor_id,\n                                request_term_id,\n                                current_term_id = term_id,\n                                \"reject buffered exchange request with stale partial graph term\"\n                            );\n                            if let TakeReceiverRequest::Remote { result_sender, .. } = request {\n                                let _ = result_sender.send(Err(anyhow!(\n                                    \"take receiver ({upstream_actor_id}, {actor_id}) on stale partial graph term {request_term_id}; current term is {term_id}\"\n                                )\n                                .into()));\n                            }\n                        }\n                    }\n                    *status = PartialGraphStatus::Running(graph);\n                } else {\n                    panic!(\"duplicated partial graph: {}\", partial_graph_id);\n                }\n\n                status\n            }\n            Entry::Vacant(entry) => entry.insert(PartialGraphStatus::Running(\n                PartialGraphState::new(partial_graph_id, term_id, self.actor_manager.clone()),\n            )),\n        };\n    }","sourceCodeStart":947,"sourceCodeEnd":983,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/task/barrier_worker/mod.rs#L947-L983","documentation":"When a new partial graph is created (add_partial_graph), the worker replays buffered exchange TakeReceiver requests. Any buffered request whose term id is older than the newly created graph's term id is stale: it is rejected by sending this error to the remote requester's oneshot result_sender. This prevents wiring channels from an obsolete graph generation into the fresh one.","triggerScenarios":"add_partial_graph (called via handle_streaming_control_request for Request::CreatePartialGraph) drains buffered ReceivedExchangeRequest entries and finds request_term_id != term_id of the new graph, rejecting TakeReceiverRequest::Remote via result_sender.","commonSituations":"Upstream exchange requests buffered while the graph was absent, then superseded by a graph rebuild with a bumped term; slow upstream actor retrying an epoch that was already replaced during recovery.","solutions":["Have the upstream/downstream retry the channel creation with the current term id after CreatePartialGraph.","Verify term id propagation from the meta node so actors always request with the latest term.","Check for delayed or replayed control requests (network retries, stale gRPC buffers) that carry old terms.","Treat this rejection as expected during recovery if followed by a successful retry."],"exampleFix":"// before: replaying all buffered requests verbatim\nfor (term_id, ids, request) in pending { dispatch(request); }\n// after: only replay requests whose term matches the new graph\nfor (req_term, ids, request) in pending {\n    if req_term == term_id { dispatch(request); } else { reject_stale(request, req_term, term_id); }\n}","handlingStrategy":"retry","validationCode":"// drop buffered requests whose term predates the new graph\nlet live: Vec<_> = pending.into_iter().filter(|(t, _, _)| *t == term_id).collect();","typeGuard":"fn is_stale_term(err: &str) -> bool { err.contains(\"stale partial graph term\") }","tryCatchPattern":"// remote side: on stale rejection, re-request with current term\nif result.is_err() && is_stale_term(&result.err().unwrap().to_string()) {\n    reissue_take_receiver_with_latest_term().await;\n}","preventionTips":["Discard in-flight exchange requests whenever a term bump is observed.","Bound retry counts for stale-term requests to avoid unbounded loops.","Keep term id monotonic and propagated through every control message."],"tags":["streaming","barrier","stale-request","term-mismatch"],"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"}