vectordotdev/vector · error

input channel should not be closed

Error message

input channel should not be closed

What it means

The input driver task sends every input test event into the topology's input channel and asserts the channel is still open. If input_tx.send() returns Err (SendError), the receiving side (the source's input pump inside the running topology) has been dropped — usually because the topology crashed or was torn down while input events were still being sent.

Solutions

  1. Check the component under test for panics/early exits on the input events used in the test case.
  2. Inspect crash_rx / topology logs around the failure to see which component crashed before the input driver finished.
  3. Verify the test topology's source starts successfully (no startup errors) so the input receiver stays alive.
Defensive patterns

Strategy: try-catch

Validate before calling

if input_tx.is_closed() { panic!("topology input receiver dropped before all inputs were sent"); }

Try / catch

if let Err(e) = input_tx.send(input_event.clone()).await {
    panic!("topology crashed while sending input: {e}");
}

Prevention

When it happens

Trigger: Running run_validation when the component topology crashes mid-test (component panic picked up by crash_rx, or topology shutdown racing the input driver), causing the receiver end of the input mpsc channel to drop before all input_events are sent.

Common situations: A component under test panics or exits early on the first input event; test teardown ordering races; a source in the test topology fails to start so its input pump never exists.

Related errors


AI-assisted analysis of vectordotdev/vector@bdb87aeaa4 (2026-09-16). Data as JSON: /api/errors/92fffc9def66f0f9. Report an issue: GitHub.

Appendix: source

Thrown at src/components/validation/runner/mod.rs:578

fn spawn_input_driver(
    input_events: Vec<TestEvent>,
    input_tx: Sender<TestEvent>,
    runner_metrics: &Arc<Mutex<RunnerMetrics>>,
    mut maybe_encoder: Option<Encoder<encoding::Framer>>,
    component_type: ComponentType,
    log_namespace: LogNamespace,
) -> JoinHandle<()> {
    let input_runner_metrics = Arc::clone(runner_metrics);

    let now = Utc::now();

    tokio::spawn(async move {
        for mut input_event in input_events {
            input_tx
                .send(input_event.clone())
                .await
                .expect("input channel should not be closed");

            // Update the runner metrics for the sent event. This will later
            // be used in the Validators, as the "expected" case.
            let mut input_runner_metrics = input_runner_metrics.lock().await;

            // the controlled edge (vector source) adds metadata to the event when it is received.
            // thus we need to add it here so the expected values for the comparisons on transforms
            // and sinks are accurate.
            if component_type != ComponentType::Source
                && let Event::Log(log) = input_event.get_event()
            {
                log_namespace.insert_standard_vector_source_metadata(log, "vector", now);
            }

            let (failure_case, mut event) = input_event.clone().get();

            if let Some(encoder) = maybe_encoder.as_mut() {
                let mut buffer = BytesMut::new();

View on GitHub (pinned to bdb87aeaa4)