vectordotdev/vector · error

event forward rx should not close first

Error message

event forward rx should not close first

What it means

push_events sends the collected events over an mpsc channel to the validation runner. The send is expected to succeed because the receiving end (event forward rx) is held by the runner for the server's lifetime; a panic means the receiver was dropped first, i.e. the runner shut down while the gRPC server was still pushing events.

Solutions

  1. Check the runner-side task logs for an earlier panic that dropped the receiver.
  2. Ensure shutdown coordination (input_task_coordinator/output_task_coordinator) closes the gRPC server before dropping rx.
  3. Handle the send error gracefully (return a gRPC Status::unavailable) rather than panicking.

Example fix

// before
self.tx.send(events).await.expect("event forward rx should not close first");
// after
self.tx.send(events).await.map_err(|_| tonic::Status::unavailable("runner shut down"))?;
Defensive patterns

Strategy: try-catch

Try / catch

if self.tx.send(events).await.is_err() {
    return Ok(()) /* or Status::unavailable */; // runner already shut down
}

Prevention

When it happens

Trigger: tx.send(events).await returns Err (RecvError) because the runner dropped the forwarding receiver — runner shutdown racing ahead of the gRPC output server, or premature drop of RunnerOutput plumbing.

Common situations: Shutdown-ordering bugs in the validation runner; a panic in the runner task that owned the rx, causing the server to panic on the next push.

Related errors


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

Appendix: source

Thrown at src/components/validation/runner/io.rs:57

#[tonic::async_trait]
impl VectorService for EventForwardService {
    async fn push_events(
        &self,
        request: tonic::Request<PushEventsRequest>,
    ) -> Result<tonic::Response<PushEventsResponse>, Status> {
        let events = request
            .into_inner()
            .events
            .into_iter()
            .map(|wrapper| {
                Event::try_from(wrapper).expect("validation events are encoded by Vector")
            })
            .collect();

        self.tx
            .send(events)
            .await
            .expect("event forward rx should not close first");

        Ok(tonic::Response::new(PushEventsResponse {}))
    }

    async fn health_check(
        &self,
        _: tonic::Request<HealthCheckRequest>,
    ) -> Result<tonic::Response<HealthCheckResponse>, Status> {
        let message = HealthCheckResponse {
            status: ServingStatus::Serving.into(),
        };

        Ok(tonic::Response::new(message))
    }
}

pub struct InputEdge {
    #[allow(dead_code)]

View on GitHub (pinned to bdb87aeaa4)