{"record":{"id":"8d56e674ffb399b2","repo":"risingwavelabs/risingwave","slug":"get-unsent-offset-in-response","errorCode":null,"errorMessage":"get unsent offset {:?} in response","messagePattern":"get unsent offset (.+?) in response","errorType":"exception","errorClass":"anyhow","httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/remote.rs","lineNumber":369,"sourceCode":"                }\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(|| {\n                    anyhow!(\"get unsent offset {:?} in response\", persisted_offset)\n                })?;\n                if sent_offset != persisted_offset {\n                    bail!(\n                        \"new response offset {:?} does not match the buffer offset {:?}\",\n                        persisted_offset,\n                        sent_offset\n                    );\n                }\n\n                if let (TruncateOffset::Barrier { .. }, Some(start_time)) =\n                    (persisted_offset, start_time)\n                {\n                    sink_writer_metrics\n                        .sink_commit_duration\n                        .observe(start_time.elapsed().as_secs_f64());\n                }\n\n                log_reader.truncate(persisted_offset)?;","sourceCodeStart":351,"sourceCodeEnd":387,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/remote.rs#L351-L387","documentation":"truncate_matched_offset pops the internal queue of already-sent (offset, timestamp) entries to match the persisted offset reported by a TruncateOffset response. If the queue is empty when it still needs an unsent offset, the bookkeeping invariant is broken and this error is raised.","triggerScenarios":"A TruncateOffset { offset } response arrives whose offset has no remaining entry in the sent-offset queue - the queue was drained by earlier truncations or never contained this offset.","commonSituations":"Duplicate or out-of-order truncate responses from the remote sink; recovery replay that loses in-memory queue state; a bug in the remote connector acknowledging offsets it was never sent.","solutions":["Capture the reported persisted_offset and compare against the sink's recent write history.","Restart/recover the sink actor so the queue is rebuilt from the source of truth.","If reproducible, report with the remote connector name - the connector likely acknowledged offsets it never received."],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"if let Err(e) = truncate_matched_offset(&mut queue, offset).await {\n    log::warn!(\"truncate bookkeeping broken: {e}; resetting from last barrier\");\n    request_recovery_from_last_barrier();\n}","preventionTips":["Ensure the remote connector only acknowledges offsets it actually received","Rebuild sent-offset queue state deterministically on recovery","Add debug logging of truncate requests/responses to catch duplicates"],"tags":["sink","offset","bookkeeping","internal-invariant"],"backgroundTag":"record-not-found","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"}