{"record":{"id":"80761b332a5d629f","repo":"risingwavelabs/risingwave","slug":"new-item-epoch-does-not-exceed-barrier-offset-e","errorCode":null,"errorMessage":"new item epoch {} does not exceed barrier offset epoch {}","messagePattern":"new item epoch (.+?) does not exceed barrier offset epoch (.+?)","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/log_store.rs","lineNumber":107,"sourceCode":"    pub fn check_next_item_epoch(&self, epoch: u64) -> LogStoreResult<()> {\n        match self {\n            TruncateOffset::Chunk {\n                epoch: offset_epoch,\n                ..\n            } => {\n                if epoch != *offset_epoch {\n                    bail!(\n                        \"new item epoch {} does not match current chunk offset epoch {}\",\n                        epoch,\n                        offset_epoch\n                    );\n                }\n            }\n            TruncateOffset::Barrier {\n                epoch: offset_epoch,\n            } => {\n                if epoch <= *offset_epoch {\n                    bail!(\n                        \"new item epoch {} does not exceed barrier offset epoch {}\",\n                        epoch,\n                        offset_epoch\n                    );\n                }\n            }\n        }\n        Ok(())\n    }\n}\n\n#[derive(Debug)]\npub enum LogStoreReadItem {\n    StreamChunk {\n        chunk: StreamChunk,\n        chunk_id: ChunkId,\n    },\n    Barrier {","sourceCodeStart":89,"sourceCodeEnd":125,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/log_store.rs#L89-L125","documentation":"`check_next_item_epoch` on a `TruncateOffset::Barrier` requires the new item's epoch to be strictly greater than the barrier's epoch, because everything after a barrier belongs to a later epoch. An equal-or-lesser epoch means pre-barrier data is arriving post-barrier.","triggerScenarios":"Calling `check_next_item_epoch(epoch)` while the current offset is `TruncateOffset::Barrier { epoch: offset_epoch }` with `epoch <= offset_epoch`; e.g. buffered items from before a barrier being written after the barrier was recorded.","commonSituations":"Out-of-order buffering in sink writers; replaying a queue from a stale cursor after recovery; barrier arrives while older-epoch items remain pending.","solutions":["Drop or flush all pre-barrier items before recording the barrier offset","Ensure items are sorted/filtered so post-barrier writes have epoch > barrier epoch","Audit recovery paths for stale cursors replaying old epochs"],"exampleFix":"// before\nstore.check_next_item_epoch(item.epoch); // item.epoch <= barrier_epoch\n// after\nif item.epoch > barrier_epoch {\n    store.check_next_item_epoch(item.epoch);\n} else {\n    item.drop_stale();\n}","handlingStrategy":"validation","validationCode":"assert!(item_epoch > barrier_epoch, \"post-barrier items must have epoch > barrier epoch\");","typeGuard":"fn exceeds_barrier(epoch: u64, offset: &TruncateOffset) -> bool {\n    match offset {\n        TruncateOffset::Barrier { epoch: e } => epoch > *e,\n        _ => true,\n    }\n}","tryCatchPattern":"match offset.check_next_item_epoch(epoch) {\n    Err(_) => { /* drop stale pre-barrier item or replay from correct cursor */ }\n    Ok(()) => { /* write item */ }\n}","preventionTips":["Flush all pre-barrier items before recording the barrier offset","Persist cursor positions so recovery never replays pre-barrier epochs","Validate epoch ordering in unit tests for barrier transitions"],"tags":["rust","streaming","epoch-mismatch"],"backgroundTag":"invalid-state-transition","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-14T16:17:12.679Z"}