{"record":{"id":"a8648234feabdbce","repo":"risingwavelabs/risingwave","slug":"end-of-response-stream","errorCode":null,"errorMessage":"end of response stream","messagePattern":"end of response stream","errorType":"exception","errorClass":"SinkError::Remote","httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/remote.rs","lineNumber":349,"sourceCode":"        let mut response_err_stream_rx = self.response_stream;\n        let sink_writer_metrics = self.sink_writer_metrics;\n\n        let (response_tx, mut response_rx) = unbounded_channel();\n\n        let poll_response_stream = async move {\n            loop {\n                let result = response_err_stream_rx\n                    .stream\n                    .try_next()\n                    .instrument_await(\"log_sinker_wait_next_response\")\n                    .await;\n                match result {\n                    Ok(Some(response)) => {\n                        response_tx.send(response).map_err(|err| {\n                            SinkError::Remote(anyhow!(\"unable to send response: {:?}\", err.0))\n                        })?;\n                    }\n                    Ok(None) => return Err(SinkError::Remote(anyhow!(\"end of response stream\"))),\n                    Err(e) => return Err(SinkError::Remote(anyhow!(e))),\n                }\n            }\n        };\n\n        let poll_consume_log_and_sink = async move {\n            fn truncate_matched_offset(\n                queue: &mut VecDeque<(TruncateOffset, Option<Instant>)>,\n                persisted_offset: TruncateOffset,\n                log_reader: &mut impl SinkLogReader,\n                sink_writer_metrics: &SinkWriterMetrics,\n            ) -> Result<()> {\n                while let Some((sent_offset, _)) = queue.front()\n                    && sent_offset < &persisted_offset\n                {\n                    queue.pop_front();\n                }\n","sourceCodeStart":331,"sourceCodeEnd":367,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/remote.rs#L331-L367","documentation":"The log reader stream terminated normally (returned None) while consume_log_and_sink still expected more responses. RisingWave treats premature end-of-stream of the remote sink response channel as an error because the sink writer must stay alive until the epoch completes.","triggerScenarios":"The remote (JVM) sink writer stream ends (Ok(None) from try_next) while consume_log_and_sink is still polling for responses, e.g. the remote peer closed the stream before all writes were acknowledged.","commonSituations":"Remote sink process/connector crashed or closed its stream mid-epoch; JNI writer dropped unexpectedly; external sink service stopped while the sink was streaming.","solutions":["Check RisingWave and JVM/connector logs around this timestamp for the remote writer closing.","Verify the downstream sink service was running and reachable for the whole epoch.","Recover/restart the sink; RisingWave will replay from the last barrier.","If reproducible, check connector version compatibility between frontend/stream and the remote connector."],"exampleFix":null,"handlingStrategy":"retry","validationCode":null,"typeGuard":null,"tryCatchPattern":"match stream.try_next().await {\n    Ok(None) => return Err(/* premature EOF: schedule recovery from last barrier */),\n    Err(e) => return Err(e),\n    Ok(Some(resp)) => handle(resp),\n}","preventionTips":["Monitor downstream sink service health so it doesn't close streams mid-epoch","Set adequate gRPC/JNI timeouts so streams aren't closed early","Alert on sink recovery events correlated with this error"],"tags":["stream","sink","remote-connector","premature-eof"],"backgroundTag":"stream-ended-unexpectedly","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"}