vectordotdev/vector · critical

Task started but input has been taken.

Error message

Task started but input has been taken.

What it means

In the spawned sink task, the builder takes the buffered receiver out of an Arc<Mutex<Option<_>>> with .take().expect(...). The Option is set exactly once; if it is already gone when the task starts, the panic "Task started but input has been taken" fires. This guards a topology-internal invariant that each sink's input receiver is consumed exactly once.

Solutions

  1. Ensure each sink key's buffer receiver is taken by exactly one task — don't rebuild the same key concurrently
  2. Serialize topology rebuilds so old tasks are fully torn down before new ones start
  3. If writing custom topology code, clone/configure inputs before spawning instead of sharing the Mutex<Option<rx>>
  4. Report as a Vector bug with reload context if it occurs with stock hot-reload

Example fix

// before
let rx = rx.lock().unwrap().take().expect("Task started but input has been taken."); // already None
// after
guard against double-build: only spawn one task per sink key, or construct fresh receivers per build
Defensive patterns

Strategy: type-guard

Validate before calling

// guard: take only once, log if already taken instead of panicking
if let Some(rx) = rx.lock().unwrap().take() { spawn_sink(rx) } else { warn!("sink {} input already consumed", key); }

Type guard

fn take_once(cell: &Mutex<Option<Receiver>>) -> Option<Receiver> { cell.lock().unwrap().take() }

Prevention

When it happens

Trigger: The same sink key's receiver is taken twice (e.g. a topology reload/rebuild race where the old and new task both start), or custom topology code building the same sink key twice.

Common situations: Bugs in hot-reload/config-diff logic rebuilding sinks concurrently; custom embedding code that invokes build_sinks for the same key more than once; rare races between health-check task startup and the main sink task.

Understand the failure class

Background: "This is a bug, please report it": internal invariant violations, unreachable panics, and SNH errors explained — this error's family across 47 libraries.

Related errors


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

Appendix: source

Thrown at src/topology/builder.rs:751

        let utilization_sender = self
            .utilization_registry
            .add_component(key.clone(), gauge!(GaugeName::Utilization));
        let component_key = key.clone();
        let sink = async move {
            debug!("Sink starting.");

            // Why is this Arc<Mutex<Option<_>>> needed you ask.
            // In case when this function build_pieces errors
            // this future won't be run so this rx won't be taken
            // which will enable us to reuse rx to rebuild
            // old configuration by passing this Arc<Mutex<Option<_>>>
            // yet again.
            let rx = rx
                .lock()
                .unwrap()
                .take()
                .expect("Task started but input has been taken.");

            let mut rx = Utilization::new(utilization_sender, component_key.clone(), rx);

            let events_received = register!(EventsReceived);
            sink.run(
                rx.by_ref()
                    .filter(|events: &EventArray| ready(filter_events_type(events, input_type)))
                    .inspect(|events| {
                        events_received.emit(CountByteSize(
                            events.len(),
                            events.estimated_json_encoded_size_of(),
                        ))
                    })
                    .take_until_if(tripwire),
            )
            .await
            .map(|_| {
                debug!("Sink finished normally.");

View on GitHub (pinned to bdb87aeaa4)