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
- Ensure each sink key's buffer receiver is taken by exactly one task — don't rebuild the same key concurrently
- Serialize topology rebuilds so old tasks are fully torn down before new ones start
- If writing custom topology code, clone/configure inputs before spawning instead of sharing the Mutex<Option<rx>>
- 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
- Spawn exactly one consuming task per sink key
- Serialize topology rebuilds during hot reload
- Use fresh receivers per build rather than shared Option cells
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
- can't run runner twice
- cant ever be empty
- join error or bad poll
- output for default port required for task transforms
- Pausing unknown sink from fanout
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)