{"record":{"id":"be43be4a68985e15","repo":"risingwavelabs/risingwave","slug":"end-of-upstream-input","errorCode":null,"errorMessage":"end of upstream input","messagePattern":"end of upstream input","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/sync_kv_log_store.rs","lineNumber":519,"sourceCode":"                sleep_future,\n                ..\n            } => {\n                if let Some(sleep_future) = sleep_future {\n                    sleep_future.await;\n                    metrics\n                        .pause_duration_ns\n                        .inc_by(start_instant.elapsed().as_nanos() as _);\n                    tracing::trace!(\"resuming write future\");\n                }\n                must_match!(replace(self, WriteFuture::Empty), WriteFuture::Paused { stream, write_state, barrier, .. } => {\n                    Ok((stream, write_state, WriteFutureEvent::UpstreamMessageReceived(Message::Barrier(barrier))))\n                })\n            }\n            WriteFuture::ReceiveFromUpstream { future, .. } => {\n                let (opt, stream) = future.await;\n                must_match!(replace(self, WriteFuture::Empty), WriteFuture::ReceiveFromUpstream { write_state, .. } => {\n                    opt\n                    .ok_or_else(|| anyhow!(\"end of upstream input\").into())\n                    .and_then(|result| result.map(|item| {\n                        (stream, write_state, WriteFutureEvent::UpstreamMessageReceived(item))\n                    }))\n                })\n            }\n            WriteFuture::FlushingChunk { future, .. } => {\n                let (write_state, result) = future.await;\n                let result = must_match!(replace(self, WriteFuture::Empty), WriteFuture::FlushingChunk { epoch, start_seq_id, end_seq_id, stream, ..  } => {\n                    result.map(|(flush_info, vnode_bitmap)| {\n                        (stream, write_state, WriteFutureEvent::ChunkFlushed(FlushedChunkInfo {\n                            epoch,\n                            start_seq_id,\n                            end_seq_id,\n                            flush_info,\n                            vnode_bitmap,\n                        }))\n                    })\n                });","sourceCodeStart":501,"sourceCodeEnd":537,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/sync_kv_log_store.rs#L501-L537","documentation":"The sync KV log store's write path awaits a future that receives a write batch from upstream. If that upstream future completes with None (the upstream stream/input ended) instead of a write batch, next_event maps it to this error — the log store cannot continue writing without upstream input.","triggerScenarios":"WriteFuture::ReceiveFromUpstream resolves with opt = None, i.e. the upstream stream delivering write batches terminated while a write operation was still pending.","commonSituations":"Upstream log/writer component shut down or errored while a write was in flight, Hummock/log-store teardown racing a pending write, or madsim tests dropping the upstream channel.","solutions":["Find why the upstream write stream ended (component shutdown, earlier error, channel drop) in the logs.","Ensure the log-store writer/upstream stays alive for the duration of writes; fix premature drop in code or tests.","Retry the operation after the log store / actor is recovered."],"exampleFix":"// before (test)\ndrop(upstream_tx);\nwrite_future.await; // -> \"end of upstream input\"\n// after\nupstream_tx.send(write_batch).unwrap(); // keep upstream alive until writes finish","handlingStrategy":"fallback","validationCode":"// before writing, confirm the upstream writer stream is still open\nif upstream_tx.is_closed() { return Err(upstream_gone()); }","typeGuard":"fn upstream_alive(tx: &UnboundedSender<WriteBatch>) -> bool { !tx.is_closed() }","tryCatchPattern":"match next_event().await {\n    Err(e) if e.to_string().contains(\"end of upstream input\") => {\n        // upstream closed mid-write: flush what we have and reopen the log store writer\n        reopen_writer_and_retry();\n    }\n    other => other,\n}","preventionTips":["Keep the upstream write channel open until all pending writes complete","Sequence shutdown: finish writes before dropping the log store writer","In tests, send a WriteBatch instead of dropping the upstream sender"],"tags":["log-store","streaming","upstream-closed"],"backgroundTag":"channel-closed-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"}