{"record":{"id":"79273294659ca2b2","repo":"vectordotdev/vector","slug":"shutdown-begun-trigger-for-source-id-not-found-in-the","errorCode":null,"errorMessage":"shutdown_begun_trigger for source \"{id}\" not found in the ShutdownCoordinator","messagePattern":"shutdown_begun_trigger for source \"(.+?)\" not found in the ShutdownCoordinator","errorType":"panic","errorClass":null,"httpStatus":null,"severity":"error","filePath":"lib/vector-common/src/shutdown.rs","lineNumber":264,"sourceCode":"    }\n\n    /// Sends the signal to the given source to begin shutting down. Returns a future that resolves\n    /// when the source has finished shutting down cleanly or been sent the force shutdown signal.\n    /// The returned future resolves to a bool that indicates if the source shut down cleanly before\n    /// the given `deadline`. If the result is false then that means the source failed to shut down\n    /// before `deadline` and had to be force-shutdown.\n    ///\n    /// # Panics\n    ///\n    /// Panics if this coordinator has had its triggers removed (ie\n    /// has been taken over with `Self::takeover_source`).\n    pub fn shutdown_source(\n        &mut self,\n        id: &ComponentKey,\n        deadline: Instant,\n    ) -> impl Future<Output = bool> + use<> {\n        let (_, begin_shutdown_trigger) = self.begun_triggers.remove(id).unwrap_or_else(|| {\n            panic!(\n                \"shutdown_begun_trigger for source \\\"{id}\\\" not found in the ShutdownCoordinator\"\n            )\n        });\n        // This is what actually triggers the source to begin shutting down.\n        begin_shutdown_trigger.cancel();\n\n        let shutdown_complete_tripwire = self\n            .complete_tripwires\n            .remove(id)\n            .unwrap_or_else(|| {\n                panic!(\n                \"shutdown_complete_tripwire for source \\\"{id}\\\" not found in the ShutdownCoordinator\"\n            )\n            });\n        let shutdown_force_trigger = self.force_triggers.remove(id).unwrap_or_else(|| {\n            panic!(\n                \"shutdown_force_trigger for source \\\"{id}\\\" not found in the ShutdownCoordinator\"\n            )","sourceCodeStart":246,"sourceCodeEnd":282,"githubUrl":"https://github.com/vectordotdev/vector/blob/bdb87aeaa4c4ff27c0ba643c1c77b21bf2ef4013/lib/vector-common/src/shutdown.rs#L246-L282","documentation":"Panic from `finalizer.expect(\"Finalizer must exist for the batch receiver to be created\")` in the aws_sqs source's `run_once`. The code creates a `BatchNotifier` receiver (which requires a finalizer) only when acknowledgements are enabled; when `batch_receiver` is `Some`, `finalizer` must also be `Some` by construction. The expect asserts that pairing — it fires only if the finalizer/receiver setup logic diverges (finalizer None while a batch receiver exists).","triggerScenarios":"Reaching the `Some(receiver) =>` branch at src/sources/aws_sqs/source.rs:165 with `finalizer == None`: this means `BatchNotifier::maybe_new_with_receiver` returned a receiver without acknowledgements being enabled, or the code creating `finalizer` and `batch_receiver` was changed so their conditions are no longer identical.","commonSituations":"Refactoring acknowledgement setup (e.g. building the receiver without guarding on `self.acknowledgements_enabled`); changing config defaults so `delete_message` is true while acknowledgements handling is partially initialized; a merge that decoupled the two `Option`s.","solutions":["Ensure `batch_receiver` is only created when the finalizer is also created (gate both on the same `acknowledgements` flag).","Restructure the match to `Some((finalizer, receiver))` from a single combined Option so the invariant is enforced by types.","If finalizer can legitimately be absent, fall back to the `None => delete_messages(...)` branch instead of expecting."],"exampleFix":"// before\nSome(receiver) => finalizer\n    .expect(\"Finalizer must exist for the batch receiver to be created\")\n    .add(receipts_to_ack, receiver),\n// after\n(Some(finalizer), Some(receiver)) => finalizer.add(receipts_to_ack, receiver),\n_ => delete_messages(self.client.clone(), receipts_to_ack, self.queue_url.clone()).await,","handlingStrategy":"type-guard","validationCode":"// Gate both on the same flag:\nlet ack = self.acknowledgements.then(|| {\n    let (receiver, _) = BatchNotifier::new_with_receiver();\n    (finalizer, receiver)\n});","typeGuard":"fn ack_pair(f: Option<Finalizer>, r: Option<BatchNotifierReceiver>) -> Option<(Finalizer, BatchNotifierReceiver)> {\n    Some((f?, r?))\n}","tryCatchPattern":"match (finalizer, batch_receiver) {\n    (Some(f), Some(r)) => f.add(receipts_to_ack, r),\n    _ => delete_messages(self.client.clone(), receipts_to_ack, self.queue_url.clone()).await,\n}","preventionTips":["Create the finalizer and the batch receiver in one place so their Option-ness can never diverge.","Type-encode the pairing (Option<(Finalizer, Receiver)>) instead of two independent Options.","Test with acknowledgements enabled and delete_message=true."],"tags":["aws","sqs","rust","panic","acknowledgements"],"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"}