{"record":{"id":"aa1ec2acf59bb1a4","repo":"risingwavelabs/risingwave","slug":"sinkerror-remote-anyhow-e","errorCode":null,"errorMessage":"SinkError::Remote(anyhow!(e))","messagePattern":"SinkError::Remote\\(anyhow!\\(e\\)\\)","errorType":"exception","errorClass":"SinkError::Remote","httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/remote.rs","lineNumber":350,"sourceCode":"        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\n                let (sent_offset, start_time) = queue.pop_front().ok_or_else(|| {","sourceCodeStart":332,"sourceCodeEnd":368,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/remote.rs#L332-L368","documentation":"This error propagates an error returned by the log reader stream itself (Err(e) branch). The original error e from the remote sink log reader is wrapped in SinkError::Remote and surfaced, so the root cause text is whatever the reader failed with (gRPC/JNI/IO errors).","triggerScenarios":"consume_log_and_sink's try_next() on the log reader returns Err(e); any transport or decoding failure of the sink log stream is re-raised here.","commonSituations":"gRPC stream to the remote sink writer failed (connection reset, deadline exceeded); JNI layer threw; response message failed to decode.","solutions":["Inspect the wrapped inner message for the true root cause (network vs decode vs JNI).","For transport errors, check network stability between components and retry after recovery.","For decode errors, verify connector/prost message version compatibility.","Enable connector debug logs to capture the failing request."],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"match log_reader.try_next().await {\n    Err(e) => {\n        // classify: transport vs decode; retry transport errors, fail fast on decode\n        if is_transport(&e) { schedule_retry(); } else { return Err(e.into()); }\n    }\n    other => handle(other),\n}","preventionTips":["Monitor network stability between RW and the remote connector","Keep prost/connector message versions aligned to avoid decode errors","Log the inner error immediately - it carries the true root cause"],"tags":["stream","sink","error-propagation","remote-connector"],"backgroundTag":"upstream-api-error","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"}