vectordotdev/vector · error

shutdown_begun_trigger for source

Error message

shutdown_begun_trigger for source "{id}" not found in the ShutdownCoordinator

What it means

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).

Solutions

  1. Ensure `batch_receiver` is only created when the finalizer is also created (gate both on the same `acknowledgements` flag).
  2. Restructure the match to `Some((finalizer, receiver))` from a single combined Option so the invariant is enforced by types.
  3. If finalizer can legitimately be absent, fall back to the `None => delete_messages(...)` branch instead of expecting.

Example fix

// before
Some(receiver) => finalizer
    .expect("Finalizer must exist for the batch receiver to be created")
    .add(receipts_to_ack, receiver),
// after
(Some(finalizer), Some(receiver)) => finalizer.add(receipts_to_ack, receiver),
_ => delete_messages(self.client.clone(), receipts_to_ack, self.queue_url.clone()).await,
Defensive patterns

Strategy: type-guard

Validate before calling

// Gate both on the same flag:
let ack = self.acknowledgements.then(|| {
    let (receiver, _) = BatchNotifier::new_with_receiver();
    (finalizer, receiver)
});

Type guard

fn ack_pair(f: Option<Finalizer>, r: Option<BatchNotifierReceiver>) -> Option<(Finalizer, BatchNotifierReceiver)> {
    Some((f?, r?))
}

Try / catch

match (finalizer, batch_receiver) {
    (Some(f), Some(r)) => f.add(receipts_to_ack, r),
    _ => delete_messages(self.client.clone(), receipts_to_ack, self.queue_url.clone()).await,
}

Prevention

When it happens

Trigger: 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.

Common situations: 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.

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/79273294659ca2b2. Report an issue: GitHub.

Appendix: source

Thrown at lib/vector-common/src/shutdown.rs:264

    }

    /// Sends the signal to the given source to begin shutting down. Returns a future that resolves
    /// when the source has finished shutting down cleanly or been sent the force shutdown signal.
    /// The returned future resolves to a bool that indicates if the source shut down cleanly before
    /// the given `deadline`. If the result is false then that means the source failed to shut down
    /// before `deadline` and had to be force-shutdown.
    ///
    /// # Panics
    ///
    /// Panics if this coordinator has had its triggers removed (ie
    /// has been taken over with `Self::takeover_source`).
    pub fn shutdown_source(
        &mut self,
        id: &ComponentKey,
        deadline: Instant,
    ) -> impl Future<Output = bool> + use<> {
        let (_, begin_shutdown_trigger) = self.begun_triggers.remove(id).unwrap_or_else(|| {
            panic!(
                "shutdown_begun_trigger for source \"{id}\" not found in the ShutdownCoordinator"
            )
        });
        // This is what actually triggers the source to begin shutting down.
        begin_shutdown_trigger.cancel();

        let shutdown_complete_tripwire = self
            .complete_tripwires
            .remove(id)
            .unwrap_or_else(|| {
                panic!(
                "shutdown_complete_tripwire for source \"{id}\" not found in the ShutdownCoordinator"
            )
            });
        let shutdown_force_trigger = self.force_triggers.remove(id).unwrap_or_else(|| {
            panic!(
                "shutdown_force_trigger for source \"{id}\" not found in the ShutdownCoordinator"
            )

View on GitHub (pinned to bdb87aeaa4)