risingwavelabs/risingwave · warning · anyhow

take receiver ( , ) on stale partial graph term ; current…

Error message

take receiver ({upstream_actor_id}, {actor_id}) on stale partial graph term {request_term_id}; current term is {term_id}

What it means

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.

Solutions

  1. Have the upstream/downstream retry the channel creation with the current term id after CreatePartialGraph.
  2. Verify term id propagation from the meta node so actors always request with the latest term.
  3. Check for delayed or replayed control requests (network retries, stale gRPC buffers) that carry old terms.
  4. Treat this rejection as expected during recovery if followed by a successful retry.

Example fix

// before: replaying all buffered requests verbatim
for (term_id, ids, request) in pending { dispatch(request); }
// after: only replay requests whose term matches the new graph
for (req_term, ids, request) in pending {
    if req_term == term_id { dispatch(request); } else { reject_stale(request, req_term, term_id); }
}
Defensive patterns

Strategy: retry

Validate before calling

// drop buffered requests whose term predates the new graph
let live: Vec<_> = pending.into_iter().filter(|(t, _, _)| *t == term_id).collect();

Type guard

fn is_stale_term(err: &str) -> bool { err.contains("stale partial graph term") }

Try / catch

// remote side: on stale rejection, re-request with current term
if result.is_err() && is_stale_term(&result.err().unwrap().to_string()) {
    reissue_take_receiver_with_latest_term().await;
}

Prevention

When it happens

Trigger: 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.

Common situations: 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.

Understand the failure class

Background: "Invalid state transition" errors: "status must be X, actually Y", "already rejected/charging/uninstalled", "cannot ... while running" — what they mean when a library rejects your call — this error's family across 31 libraries.

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/2c194ad6ea264a32. Report an issue: GitHub.

Appendix: source

Thrown at src/stream/src/task/barrier_worker/mod.rs:965

                        term_id.clone(),
                        self.actor_manager.clone(),
                    );
                    for (request_term_id, (upstream_actor_id, actor_id), request) in
                        pending_requests.drain(..)
                    {
                        if request_term_id == term_id {
                            graph.new_actor_output_request(actor_id, upstream_actor_id, request);
                        } else {
                            warn!(
                                %partial_graph_id,
                                %upstream_actor_id,
                                %actor_id,
                                request_term_id,
                                current_term_id = term_id,
                                "reject buffered exchange request with stale partial graph term"
                            );
                            if let TakeReceiverRequest::Remote { result_sender, .. } = request {
                                let _ = result_sender.send(Err(anyhow!(
                                    "take receiver ({upstream_actor_id}, {actor_id}) on stale partial graph term {request_term_id}; current term is {term_id}"
                                )
                                .into()));
                            }
                        }
                    }
                    *status = PartialGraphStatus::Running(graph);
                } else {
                    panic!("duplicated partial graph: {}", partial_graph_id);
                }

                status
            }
            Entry::Vacant(entry) => entry.insert(PartialGraphStatus::Running(
                PartialGraphState::new(partial_graph_id, term_id, self.actor_manager.clone()),
            )),
        };
    }

View on GitHub (pinned to 6469eb736d)