risingwavelabs/risingwave · error

end of upstream input

Error message

end of upstream input

What it means

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.

Solutions

  1. Find why the upstream write stream ended (component shutdown, earlier error, channel drop) in the logs.
  2. Ensure the log-store writer/upstream stays alive for the duration of writes; fix premature drop in code or tests.
  3. Retry the operation after the log store / actor is recovered.

Example fix

// before (test)
drop(upstream_tx);
write_future.await; // -> "end of upstream input"
// after
upstream_tx.send(write_batch).unwrap(); // keep upstream alive until writes finish
Defensive patterns

Strategy: fallback

Validate before calling

// before writing, confirm the upstream writer stream is still open
if upstream_tx.is_closed() { return Err(upstream_gone()); }

Type guard

fn upstream_alive(tx: &UnboundedSender<WriteBatch>) -> bool { !tx.is_closed() }

Try / catch

match next_event().await {
    Err(e) if e.to_string().contains("end of upstream input") => {
        // upstream closed mid-write: flush what we have and reopen the log store writer
        reopen_writer_and_retry();
    }
    other => other,
}

Prevention

When it happens

Trigger: WriteFuture::ReceiveFromUpstream resolves with opt = None, i.e. the upstream stream delivering write batches terminated while a write operation was still pending.

Common situations: 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.

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/be43be4a68985e15. Report an issue: GitHub.

Appendix: source

Thrown at src/stream/src/executor/sync_kv_log_store.rs:519

                sleep_future,
                ..
            } => {
                if let Some(sleep_future) = sleep_future {
                    sleep_future.await;
                    metrics
                        .pause_duration_ns
                        .inc_by(start_instant.elapsed().as_nanos() as _);
                    tracing::trace!("resuming write future");
                }
                must_match!(replace(self, WriteFuture::Empty), WriteFuture::Paused { stream, write_state, barrier, .. } => {
                    Ok((stream, write_state, WriteFutureEvent::UpstreamMessageReceived(Message::Barrier(barrier))))
                })
            }
            WriteFuture::ReceiveFromUpstream { future, .. } => {
                let (opt, stream) = future.await;
                must_match!(replace(self, WriteFuture::Empty), WriteFuture::ReceiveFromUpstream { write_state, .. } => {
                    opt
                    .ok_or_else(|| anyhow!("end of upstream input").into())
                    .and_then(|result| result.map(|item| {
                        (stream, write_state, WriteFutureEvent::UpstreamMessageReceived(item))
                    }))
                })
            }
            WriteFuture::FlushingChunk { future, .. } => {
                let (write_state, result) = future.await;
                let result = must_match!(replace(self, WriteFuture::Empty), WriteFuture::FlushingChunk { epoch, start_seq_id, end_seq_id, stream, ..  } => {
                    result.map(|(flush_info, vnode_bitmap)| {
                        (stream, write_state, WriteFutureEvent::ChunkFlushed(FlushedChunkInfo {
                            epoch,
                            start_seq_id,
                            end_seq_id,
                            flush_info,
                            vnode_bitmap,
                        }))
                    })
                });

View on GitHub (pinned to 6469eb736d)