risingwavelabs/risingwave · error · anyhow::Error

should get Stopped but get

Error message

should get Stopped but get {:?}

What it means

SinkCoordinateClient::stop sends a Stop(true) request and expects the coordinator to answer with a Stopped response. Any other response yields this error, meaning the sink coordination session's termination was not confirmed.

Solutions

  1. Inspect the debug-printed `msg` to identify the actual response
  2. Treat as best-effort on teardown: if the stream is already closed, the sink can be considered stopped; otherwise recreate the stream and send Stop again
  3. Check coordinator logs for termination errors

Example fix

// before
match self.next_response().await? { ... Stopped(_) => Ok(()), msg => Err(anyhow!("should get Stopped but get {:?}", msg)) }
// after (teardown is best-effort)
match self.next_response().await? {
    CoordinateResponse { msg: Some(coordinate_response::Msg::Stopped(_)) } => Ok(()),
    other => { tracing::warn!("stop: unexpected coordinator response {:?}; treating as stopped", other); Ok(()) }
}
Defensive patterns

Strategy: try-catch

Try / catch

// Rust: stop is best-effort during teardown
if let Err(e) = client.stop().await {
    tracing::warn!("sink stop not confirmed: {e:#}; continuing teardown");
}

Prevention

When it happens

Trigger: Calling stop when the coordinator responds with a non-Stopped message — coordinator already terminated, stream closed with an error-shaped response, or stale response from a prior request.

Common situations: Coordinator dropped the stream during shutdown; sink actor tearing down concurrently with a coordinator failure; requests misordered after a failed earlier round.

Related errors


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

Appendix: source

Thrown at src/rpc_client/src/sink_coordinate_client.rs:148

                    Some(coordinate_response::Msg::StartResponse(StartCoordinationResponse {
                        log_store_rewind_start_epoch,
                    })),
            } => Ok(log_store_rewind_start_epoch
                .ok_or_else(|| anyhow!("should get start epoch after update vnode bitmap"))?),
            msg => Err(anyhow!("should get start response but get {:?}", msg)),
        }
    }

    pub async fn stop(mut self) -> anyhow::Result<()> {
        self.send_request(CoordinateRequest {
            msg: Some(coordinate_request::Msg::Stop(true)),
        })
        .await?;
        match self.next_response().await? {
            CoordinateResponse {
                msg: Some(coordinate_response::Msg::Stopped(_)),
            } => Ok(()),
            msg => Err(anyhow!("should get Stopped but get {:?}", msg)),
        }
    }

    pub async fn align_initial_epoch(&mut self, initial_epoch: u64) -> anyhow::Result<u64> {
        self.send_request(CoordinateRequest {
            msg: Some(coordinate_request::Msg::AlignInitialEpochRequest(
                initial_epoch,
            )),
        })
        .await?;
        match self.next_response().await? {
            CoordinateResponse {
                msg: Some(coordinate_response::Msg::AlignInitialEpochResponse(epoch)),
            } => Ok(epoch),
            msg => Err(anyhow!(
                "should get AlignInitialEpochResponse but get {:?}",
                msg
            )),

View on GitHub (pinned to 6469eb736d)