{"record":{"id":"3f9e4230e740f095","repo":"vectordotdev/vector","slug":"task-started-but-input-has-been-taken","errorCode":null,"errorMessage":"Task started but input has been taken.","messagePattern":"Task started but input has been taken\\.","errorType":"panic","errorClass":null,"httpStatus":null,"severity":"critical","filePath":"src/topology/builder.rs","lineNumber":751,"sourceCode":"\n        let utilization_sender = self\n            .utilization_registry\n            .add_component(key.clone(), gauge!(GaugeName::Utilization));\n        let component_key = key.clone();\n        let sink = async move {\n            debug!(\"Sink starting.\");\n\n            // Why is this Arc<Mutex<Option<_>>> needed you ask.\n            // In case when this function build_pieces errors\n            // this future won't be run so this rx won't be taken\n            // which will enable us to reuse rx to rebuild\n            // old configuration by passing this Arc<Mutex<Option<_>>>\n            // yet again.\n            let rx = rx\n                .lock()\n                .unwrap()\n                .take()\n                .expect(\"Task started but input has been taken.\");\n\n            let mut rx = Utilization::new(utilization_sender, component_key.clone(), rx);\n\n            let events_received = register!(EventsReceived);\n            sink.run(\n                rx.by_ref()\n                    .filter(|events: &EventArray| ready(filter_events_type(events, input_type)))\n                    .inspect(|events| {\n                        events_received.emit(CountByteSize(\n                            events.len(),\n                            events.estimated_json_encoded_size_of(),\n                        ))\n                    })\n                    .take_until_if(tripwire),\n            )\n            .await\n            .map(|_| {\n                debug!(\"Sink finished normally.\");","sourceCodeStart":733,"sourceCodeEnd":769,"githubUrl":"https://github.com/vectordotdev/vector/blob/bdb87aeaa4c4ff27c0ba643c1c77b21bf2ef4013/src/topology/builder.rs#L733-L769","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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"],"exampleFix":"// before\nlet rx = rx.lock().unwrap().take().expect(\"Task started but input has been taken.\"); // already None\n// after\nguard against double-build: only spawn one task per sink key, or construct fresh receivers per build","handlingStrategy":"type-guard","validationCode":"// guard: take only once, log if already taken instead of panicking\nif let Some(rx) = rx.lock().unwrap().take() { spawn_sink(rx) } else { warn!(\"sink {} input already consumed\", key); }","typeGuard":"fn take_once(cell: &Mutex<Option<Receiver>>) -> Option<Receiver> { cell.lock().unwrap().take() }","tryCatchPattern":null,"preventionTips":["Spawn exactly one consuming task per sink key","Serialize topology rebuilds during hot reload","Use fresh receivers per build rather than shared Option cells"],"tags":["rust","panic","internal","topology"],"backgroundTag":"internal-invariant-violation","analyzedSha":"bdb87aeaa4c4ff27c0ba643c1c77b21bf2ef4013","analyzedAt":"2026-09-16T02:53:35.741Z","contentChangedAt":"2026-09-16T02:53:35.741Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}