{"record":{"id":"78f744d5923562ac","repo":"risingwavelabs/risingwave","slug":"left-barrier-received-while-right-stream-end-78f744","errorCode":null,"errorMessage":"left barrier received while right stream end","messagePattern":"left barrier received while right stream end","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/lookup/sides.rs","lineNumber":158,"sourceCode":"                    while let Some(msg) = right.next().await {\n                        match msg? {\n                            w @ Message::Watermark(_) => yield Either::Left(w),\n                            c @ Message::Chunk(_) => yield Either::Left(c),\n                            Message::Barrier(_) => {\n                                bail!(\"right barrier received while left stream end\");\n                            }\n                        }\n                    }\n                    break 'outer;\n                }\n                future::Either::Right((None, _)) => {\n                    // right stream end, passthrough left chunks\n                    while let Some(msg) = left.next().await {\n                        match msg? {\n                            w @ Message::Watermark(_) => yield Either::Right(w),\n                            c @ Message::Chunk(_) => yield Either::Right(c),\n                            Message::Barrier(_) => {\n                                bail!(\"left barrier received while right stream end\");\n                            }\n                        }\n                    }\n                    break 'outer;\n                }\n                future::Either::Left((Some(msg), _)) => match msg? {\n                    w @ Message::Watermark(_) => yield Either::Left(w),\n                    c @ Message::Chunk(_) => yield Either::Left(c),\n                    Message::Barrier(b) => {\n                        yield Either::Left(Message::Barrier(b.clone()));\n                        break 'inner (SideStatus::LeftBarrier, b);\n                    }\n                },\n                future::Either::Right((Some(msg), _)) => match msg? {\n                    w @ Message::Watermark(_) => yield Either::Right(w),\n                    c @ Message::Chunk(_) => yield Either::Right(c),\n                    Message::Barrier(b) => {\n                        yield Either::Right(Message::Barrier(b.clone()));","sourceCodeStart":140,"sourceCodeEnd":176,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/lookup/sides.rs#L140-L176","documentation":"Mirror of the left-end case: during barrier alignment the right input stream ended, and then a barrier arrived from the left input. The Lookup executor treats a barrier after the paired side has terminated as a violation of the termination protocol and fails the actor.","triggerScenarios":"Raised inside align_barrier's passthrough loop: after observing a None from the right stream, the loop drains remaining left messages and encounters Message::Barrier from the left input.","commonSituations":"Asymmetric upstream termination from an actor failure or migration; mixed-version cluster where one side's protocol differs; malformed upstream graph in dev/test harnesses.","solutions":["Inspect the right upstream's logs to determine why it terminated while the left kept sending barriers.","Re-create the materialized view / sink to get a consistent fresh plan if version skew is suspected.","Align all node versions across the cluster before restarting the stream graph.","If persistent, report with the fragment graph; this indicates an internal scheduling or protocol bug."],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"// Treat as a job-level failure and restart the stream graph:\nif let Err(e) = actor.run().await {\n    if e.to_string().contains(\"left barrier received while right stream end\") {\n        log::error!(\"upstream lifecycle mismatch: {e}\");\n    }\n}","preventionTips":["Verify both Lookup upstreams derive from the same consistent plan snapshot.","Avoid rolling upgrades mid-stream; drain and restart the cluster as one unit.","Alert on actors whose upstreams end at different times."],"tags":["streaming","barrier","executor","internal-invariant"],"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-23T08:17:48.524Z"}