From 7375d0b37cd1b17ef20d57b8d31478d027e96257 Mon Sep 17 00:00:00 2001 From: Gabor Gevay Date: Fri, 9 Oct 2026 11:02:18 +0200 Subject: [PATCH 1/3] adapter: Don't clone read holds on a concurrently dropped cluster A compute read hold reports changes to its compute instance, and `ReadHold::clone` panics once that channel is closed. `DROP CLUSTER` shuts the instance down and closes it, while holds on the cluster can still be in use, so cloning one aborted environmentd: - A fast-path peek acquires its read holds before optimization and cloned the hold on the index it peeks after it, so a `DROP CLUSTER ... CASCADE` during the optimization panicked the session task. - `Command::GetTransactionReadHoldsBundle` cloned the transaction's stored holds on the coordinator. `DROP CLUSTER` does not remove them, so a SELECT in a transaction whose cluster was dropped after the SELECT took its catalog snapshot panicked the coordinator. The fast path now takes the hold instead of cloning it, and the peek then fails with the existing "cluster ... was dropped" error. The command clones the holds with the new `ReadHolds::try_clone`, which maps a hung-up issuer to `ConcurrentDependencyDrop`, and the frontend returns that error. `ReadHolds::subset`, which the frontend applies to the fetched holds, now moves them instead of cloning, since the instance can shut down right after the fetch. The `ReadHold` docs now name a dropped compute instance as a cause of a hung-up issuer, besides process shutdown. Fixes SQL-766. Co-Authored-By: Claude Opus 5.5 --- src/adapter/src/command.rs | 4 +- src/adapter/src/coord/command_handler.rs | 6 ++- src/adapter/src/coord/read_policy.rs | 52 ++++++++++++++++++++---- src/adapter/src/frontend_peek.rs | 2 +- src/adapter/src/peek_client.rs | 15 +++---- src/storage-types/src/read_holds.rs | 13 +++--- 6 files changed, 69 insertions(+), 23 deletions(-) diff --git a/src/adapter/src/command.rs b/src/adapter/src/command.rs index b126ec6027e39..38277c0195b4b 100644 --- a/src/adapter/src/command.rs +++ b/src/adapter/src/command.rs @@ -267,9 +267,11 @@ pub enum Command { tx: oneshot::Sender, AdapterError>>, }, + /// Returns a copy of the given connection's stored transaction read holds, or + /// fails if the issuer of one of them hung up, see `ReadHolds::try_clone`. GetTransactionReadHoldsBundle { conn_id: ConnectionId, - tx: oneshot::Sender>, + tx: oneshot::Sender, AdapterError>>, }, /// _Merges_ the given read holds into the given connection's stored transaction read holds. diff --git a/src/adapter/src/coord/command_handler.rs b/src/adapter/src/coord/command_handler.rs index 7baa647a5a481..23b100bffa648 100644 --- a/src/adapter/src/coord/command_handler.rs +++ b/src/adapter/src/coord/command_handler.rs @@ -417,7 +417,11 @@ impl Coordinator { } Command::GetTransactionReadHoldsBundle { conn_id, tx } => { - let read_holds = self.txn_read_holds.get(&conn_id).cloned(); + let read_holds = self + .txn_read_holds + .get(&conn_id) + .map(|holds| holds.try_clone()) + .transpose(); let _ = tx.send(read_holds); } diff --git a/src/adapter/src/coord/read_policy.rs b/src/adapter/src/coord/read_policy.rs index 29f9750966fc0..651c37364b841 100644 --- a/src/adapter/src/coord/read_policy.rs +++ b/src/adapter/src/coord/read_policy.rs @@ -29,6 +29,7 @@ use mz_storage_types::read_policy::ReadPolicy; use timely::progress::Antichain; use timely::progress::Timestamp as _; +use crate::AdapterError; use crate::coord::id_bundle::CollectionIdBundle; use crate::coord::timeline::{TimelineContext, TimelineState}; use crate::util::ResultExt; @@ -46,6 +47,22 @@ pub struct ReadHolds { pub compute_holds: BTreeMap<(ComputeInstanceId, GlobalId), ReadHold>, } +/// The error for a storage hold whose issuer hung up. +fn storage_hold_hung_up(id: GlobalId) -> AdapterError { + AdapterError::ConcurrentDependencyDrop { + dependency_kind: "collection", + dependency_id: id.to_string(), + } +} + +/// The error for a compute hold whose issuer, its compute instance, hung up. +fn compute_hold_hung_up(instance_id: ComputeInstanceId) -> AdapterError { + AdapterError::ConcurrentDependencyDrop { + dependency_kind: "cluster", + dependency_id: instance_id.to_string(), + } +} + impl ReadHolds { /// Return empty `ReadHolds`. pub fn new() -> Self { @@ -102,22 +119,43 @@ impl ReadHolds { self.compute_holds.remove(&(instance_id, id)); } + /// Clones the holds, failing with [`AdapterError::ConcurrentDependencyDrop`] if + /// the issuer of one of them hung up, see + /// [`ReadHoldIssuerHungUp`](mz_storage_types::read_holds::ReadHoldIssuerHungUp). + /// A transaction's stored holds, for example, outlive a `DROP CLUSTER` of + /// their cluster. + pub fn try_clone(&self) -> Result { + let mut result = ReadHolds::new(); + for (id, hold) in &self.storage_holds { + let hold = hold.try_clone().map_err(|_| storage_hold_hung_up(*id))?; + result.storage_holds.insert(*id, hold); + } + for ((instance_id, id), hold) in &self.compute_holds { + let hold = hold + .try_clone() + .map_err(|_| compute_hold_hung_up(*instance_id))?; + result.compute_holds.insert((*instance_id, *id), hold); + } + Ok(result) + } + /// Returns a new ReadHolds containing only the holds for collections in `id_bundle`. - pub fn subset(&self, id_bundle: &CollectionIdBundle) -> ReadHolds { + /// + /// Moves the holds out of `self` rather than cloning them: cloning a compute + /// hold panics once `DROP CLUSTER` has shut its instance down. + pub fn subset(mut self, id_bundle: &CollectionIdBundle) -> ReadHolds { let mut result = ReadHolds::new(); for id in &id_bundle.storage_ids { - if let Some(hold) = self.storage_holds.get(id) { - result.storage_holds.insert(*id, hold.clone()); + if let Some(hold) = self.storage_holds.remove(id) { + result.storage_holds.insert(*id, hold); } } for (instance_id, ids) in &id_bundle.compute_ids { for id in ids { - if let Some(hold) = self.compute_holds.get(&(*instance_id, *id)) { - result - .compute_holds - .insert((*instance_id, *id), hold.clone()); + if let Some(hold) = self.compute_holds.remove(&(*instance_id, *id)) { + result.compute_holds.insert((*instance_id, *id), hold); } } } diff --git a/src/adapter/src/frontend_peek.rs b/src/adapter/src/frontend_peek.rs index 70b9a4b3fe629..7c25a176673d9 100644 --- a/src/adapter/src/frontend_peek.rs +++ b/src/adapter/src/frontend_peek.rs @@ -710,7 +710,7 @@ impl PeekClient { conn_id: session.conn_id().clone(), tx, }) - .await?; + .await??; if let Some(txn_read_holds) = txn_read_holds_opt { let allowed_id_bundle = txn_read_holds.id_bundle(); diff --git a/src/adapter/src/peek_client.rs b/src/adapter/src/peek_client.rs index 0923467b45883..8cb4014e729b2 100644 --- a/src/adapter/src/peek_client.rs +++ b/src/adapter/src/peek_client.rs @@ -369,7 +369,7 @@ impl PeekClient { max_result_size: u64, max_returned_query_size: Option, row_set_finishing_seconds: Histogram, - input_read_holds: ReadHolds, + mut input_read_holds: ReadHolds, peek_stash_read_batch_size_bytes: usize, peek_stash_read_memory_budget_bytes: usize, conn_id: mz_adapter_types::connection::ConnectionId, @@ -434,11 +434,13 @@ impl PeekClient { let (peek_target, target_read_hold, literal_constraints, mfp, strategy) = match fast_path { FastPathPlan::PeekExisting(_coll_id, idx_id, literal_constraints, mfp) => { let peek_target = PeekTarget::Index { id: idx_id }; + // Take the hold rather than clone it: cloning a compute hold + // panics once a concurrent DROP CLUSTER has shut its instance + // down. The peek below then fails with a "was dropped" error. let target_read_hold = input_read_holds .compute_holds - .get(&(compute_instance, idx_id)) - .expect("missing compute read hold on PeekExisting peek target") - .clone(); + .remove(&(compute_instance, idx_id)) + .expect("missing compute read hold on PeekExisting peek target"); let strategy = statement_logging::StatementExecutionStrategy::FastPath; ( peek_target, @@ -461,9 +463,8 @@ impl PeekClient { }; let target_read_hold = input_read_holds .storage_holds - .get(&coll_id) - .expect("missing storage read hold on PeekPersist peek target") - .clone(); + .remove(&coll_id) + .expect("missing storage read hold on PeekPersist peek target"); let strategy = statement_logging::StatementExecutionStrategy::PersistFastPath; ( peek_target, diff --git a/src/storage-types/src/read_holds.rs b/src/storage-types/src/read_holds.rs index c65dcdbf25796..b1860d8f531c5 100644 --- a/src/storage-types/src/read_holds.rs +++ b/src/storage-types/src/read_holds.rs @@ -54,8 +54,9 @@ impl Debug for ReadHold { /// The issuer of a [`ReadHold`] has hung up, so the hold can no longer be cloned. /// -/// This is only expected during process shutdown, when the tokio runtime drops tasks in -/// arbitrary order and the issuer task can disappear while holds still exist. +/// This happens during process shutdown, when the tokio runtime drops tasks in arbitrary order +/// and the issuer task can disappear while holds still exist, and to compute holds once their +/// compute instance is dropped, e.g., by `DROP CLUSTER`. #[derive(Error, Debug)] #[error("read hold issuer for collection {0} has hung up")] pub struct ReadHoldIssuerHungUp(pub GlobalId); @@ -178,10 +179,10 @@ impl ReadHold { /// hold has hung up, in which case the clone would not actually hold back /// the since of the collection. /// - /// The issuer hanging up is only expected during process shutdown, when - /// the tokio runtime drops tasks in arbitrary order. Callers that may run - /// concurrently with shutdown can use this method to handle that case - /// gracefully, instead of panicking via [Clone::clone]. + /// The issuer hangs up during process shutdown, and for compute holds when + /// their compute instance is dropped (see [`ReadHoldIssuerHungUp`]). Callers + /// that may run concurrently with either can use this method to handle that + /// case gracefully, instead of panicking via [Clone::clone]. pub fn try_clone(&self) -> Result { if self.id.is_user() { tracing::trace!("cloning ReadHold on {}: {:?}", self.id, self.since); From b31f65dba2986b87da8c390f1c321d1cae556de9 Mon Sep 17 00:00:00 2001 From: Gabor Gevay Date: Fri, 9 Oct 2026 12:22:56 +0200 Subject: [PATCH 2/3] adapter: Remove `Clone` from read holds `ReadHold::clone` panicked when the hold's issuer had hung up, and the issuer of a compute hold is its compute instance, which `DROP CLUSTER` shuts down in normal operation. The panic hid behind a standard trait: `#[derive(Clone)]` on `ReadHolds`, `Option::cloned` and collection clones all cloned holds without saying so, which is how the previous commit's crashes came about. `ReadHold` and `ReadHolds` no longer implement `Clone`, and `merge_assign` returns an error instead of panicking, so every site has to say what a hung-up issuer means there: - The frontend's copy of its holds for `StoreTransactionReadHolds` uses `try_clone`, and the command, `ReadHolds::merge` and `store_transaction_read_holds` return the error, so the statement fails with "cluster ... was dropped". - The bounded-staleness metric skips its comparison. - The compute instance's dataflow creation clones the holds on its inputs with `clone_input_read_hold`, which still panics: the inputs are storage collections and indexes of the same instance, whose issuers only hang up during process shutdown. Co-Authored-By: Claude Opus 5.5 --- src/adapter/src/command.rs | 5 ++- src/adapter/src/coord/command_handler.rs | 5 +-- src/adapter/src/coord/read_policy.rs | 36 ++++++++++-------- src/adapter/src/coord/sequencer/inner/peek.rs | 2 +- .../src/coord/sequencer/inner/subscribe.rs | 2 +- src/adapter/src/frontend_peek.rs | 13 +++++-- src/adapter/src/peek_client.rs | 5 +-- src/compute-client/src/controller/instance.rs | 37 ++++++++++++++++--- src/storage-types/src/read_holds.rs | 29 +++++---------- 9 files changed, 80 insertions(+), 54 deletions(-) diff --git a/src/adapter/src/command.rs b/src/adapter/src/command.rs index 38277c0195b4b..2e21fcb5659fe 100644 --- a/src/adapter/src/command.rs +++ b/src/adapter/src/command.rs @@ -274,11 +274,12 @@ pub enum Command { tx: oneshot::Sender, AdapterError>>, }, - /// _Merges_ the given read holds into the given connection's stored transaction read holds. + /// _Merges_ the given read holds into the given connection's stored transaction read holds, + /// failing if the issuer of a merged hold hung up, see `ReadHolds::merge`. StoreTransactionReadHolds { conn_id: ConnectionId, read_holds: ReadHolds, - tx: oneshot::Sender<()>, + tx: oneshot::Sender>, }, ExecuteSlowPathPeek { diff --git a/src/adapter/src/coord/command_handler.rs b/src/adapter/src/coord/command_handler.rs index 23b100bffa648..82827dd68dc07 100644 --- a/src/adapter/src/coord/command_handler.rs +++ b/src/adapter/src/coord/command_handler.rs @@ -430,8 +430,7 @@ impl Coordinator { read_holds, tx, } => { - self.store_transaction_read_holds(conn_id, read_holds); - let _ = tx.send(()); + let _ = tx.send(self.store_transaction_read_holds(conn_id, read_holds)); } Command::ExecuteSlowPathPeek { @@ -1968,7 +1967,7 @@ impl Coordinator { // NOTE: The Drop impl of ReadHolds makes sure that the hold is // released when we don't use it. if acquire_read_holds { - self.store_transaction_read_holds(session.conn_id().clone(), read_holds); + self.store_transaction_read_holds(session.conn_id().clone(), read_holds)?; } mz_now_ts diff --git a/src/adapter/src/coord/read_policy.rs b/src/adapter/src/coord/read_policy.rs index 651c37364b841..3433f0e60f6ac 100644 --- a/src/adapter/src/coord/read_policy.rs +++ b/src/adapter/src/coord/read_policy.rs @@ -41,7 +41,7 @@ use crate::util::ResultExt; /// that read frontiers cannot advance past the held time as long as they exist. /// Dropping a [`ReadHolds`] also drops the [`ReadHold`] tokens within and /// relinquishes the associated read capabilities. -#[derive(Debug, Default, Clone)] +#[derive(Debug, Default)] pub struct ReadHolds { pub storage_holds: BTreeMap, pub compute_holds: BTreeMap<(ComputeInstanceId, GlobalId), ReadHold>, @@ -139,10 +139,8 @@ impl ReadHolds { Ok(result) } - /// Returns a new ReadHolds containing only the holds for collections in `id_bundle`. - /// - /// Moves the holds out of `self` rather than cloning them: cloning a compute - /// hold panics once `DROP CLUSTER` has shut its instance down. + /// Returns a new ReadHolds containing only the holds for collections in `id_bundle`, + /// moved out of `self`. pub fn subset(mut self, id_bundle: &CollectionIdBundle) -> ReadHolds { let mut result = ReadHolds::new(); @@ -202,29 +200,38 @@ impl ReadHolds { } /// Merge the read holds in `other` into the contained read holds. - fn merge(&mut self, other: Self) { + /// + /// Fails like [`ReadHolds::try_clone`] if the issuer of a hold hung up. The + /// holds merged until then stay in `self`, and the rest of `other` is + /// released. + fn merge(&mut self, other: Self) -> Result<(), AdapterError> { use std::collections::btree_map::Entry; for (id, other_hold) in other.storage_holds { match self.storage_holds.entry(id) { Entry::Occupied(mut o) => { - o.get_mut().merge_assign(other_hold); + o.get_mut() + .merge_assign(other_hold) + .map_err(|_| storage_hold_hung_up(id))?; } Entry::Vacant(v) => { v.insert(other_hold); } } } - for (id, other_hold) in other.compute_holds { - match self.compute_holds.entry(id) { + for ((instance_id, id), other_hold) in other.compute_holds { + match self.compute_holds.entry((instance_id, id)) { Entry::Occupied(mut o) => { - o.get_mut().merge_assign(other_hold); + o.get_mut() + .merge_assign(other_hold) + .map_err(|_| compute_hold_hung_up(instance_id))?; } Entry::Vacant(v) => { v.insert(other_hold); } } } + Ok(()) } /// Extend the contained read holds with those in `other`. @@ -439,21 +446,20 @@ impl crate::coord::Coordinator { } /// Stash transaction read holds. They will be released when the transaction - /// is cleaned up. + /// is cleaned up, also after a failed merge, see [`ReadHolds::merge`]. pub(crate) fn store_transaction_read_holds( &mut self, conn_id: ConnectionId, read_holds: ReadHolds, - ) { + ) -> Result<(), AdapterError> { use std::collections::btree_map::Entry; match self.txn_read_holds.entry(conn_id) { Entry::Vacant(v) => { v.insert(read_holds); + Ok(()) } - Entry::Occupied(mut o) => { - o.get_mut().merge(read_holds); - } + Entry::Occupied(mut o) => o.get_mut().merge(read_holds), } } } diff --git a/src/adapter/src/coord/sequencer/inner/peek.rs b/src/adapter/src/coord/sequencer/inner/peek.rs index 320325addff01..acb086ce471c1 100644 --- a/src/adapter/src/coord/sequencer/inner/peek.rs +++ b/src/adapter/src/coord/sequencer/inner/peek.rs @@ -1104,7 +1104,7 @@ impl Coordinator { }); } } else if let Some(read_holds) = read_holds { - self.store_transaction_read_holds(session.conn_id().clone(), read_holds); + self.store_transaction_read_holds(session.conn_id().clone(), read_holds)?; } // TODO: Checking for only `InTransaction` and not `Implied` (also `Started`?) seems diff --git a/src/adapter/src/coord/sequencer/inner/subscribe.rs b/src/adapter/src/coord/sequencer/inner/subscribe.rs index 7d3bb33043160..57ae110bf0e12 100644 --- a/src/adapter/src/coord/sequencer/inner/subscribe.rs +++ b/src/adapter/src/coord/sequencer/inner/subscribe.rs @@ -422,7 +422,7 @@ impl Coordinator { } } - self.store_transaction_read_holds(ctx.session().conn_id().clone(), read_holds); + self.store_transaction_read_holds(ctx.session().conn_id().clone(), read_holds)?; let global_mir_plan = global_mir_plan.resolve(Antichain::from_elem(as_of)); diff --git a/src/adapter/src/frontend_peek.rs b/src/adapter/src/frontend_peek.rs index 7c25a176673d9..4dbb37a75af56 100644 --- a/src/adapter/src/frontend_peek.rs +++ b/src/adapter/src/frontend_peek.rs @@ -789,12 +789,13 @@ impl PeekClient { if in_immediate_multi_stmt_txn && determination.timestamp_context.contains_timestamp() { + let txn_read_holds = read_holds.try_clone()?; self.call_coordinator(|tx| Command::StoreTransactionReadHolds { conn_id: session.conn_id().clone(), - read_holds: read_holds.clone(), + read_holds: txn_read_holds, tx, }) - .await?; + .await??; } (determination, read_holds) @@ -1643,7 +1644,11 @@ impl PeekClient { && real_time_recency_ts.is_none() { // Note down the difference between BoundedStaleness and Serializable into a metric. - if let Some(bs_ts) = det.timestamp_context.timestamp() { + // The metric isn't worth failing the query over, so skip it if the issuer of a hold + // hung up. + if let Some(bs_ts) = det.timestamp_context.timestamp() + && let Ok(read_holds) = read_holds.try_clone() + { let (serializable_det, _tmp_read_holds) = ::determine_timestamp_for_inner( session, @@ -1653,7 +1658,7 @@ impl PeekClient { oracle_read_ts, real_time_recency_ts, &IsolationLevel::Serializable, - read_holds.clone(), + read_holds, upper, )?; if let Some(serializable) = serializable_det.timestamp_context.timestamp() { diff --git a/src/adapter/src/peek_client.rs b/src/adapter/src/peek_client.rs index 8cb4014e729b2..e08a19a1eb447 100644 --- a/src/adapter/src/peek_client.rs +++ b/src/adapter/src/peek_client.rs @@ -434,9 +434,8 @@ impl PeekClient { let (peek_target, target_read_hold, literal_constraints, mfp, strategy) = match fast_path { FastPathPlan::PeekExisting(_coll_id, idx_id, literal_constraints, mfp) => { let peek_target = PeekTarget::Index { id: idx_id }; - // Take the hold rather than clone it: cloning a compute hold - // panics once a concurrent DROP CLUSTER has shut its instance - // down. The peek below then fails with a "was dropped" error. + // If a concurrent DROP CLUSTER shut this hold's instance down, + // the peek below fails with a "was dropped" error. let target_read_hold = input_read_holds .compute_holds .remove(&(compute_instance, idx_id)) diff --git a/src/compute-client/src/controller/instance.rs b/src/compute-client/src/controller/instance.rs index a329c158874d8..f38adcb1d6bf1 100644 --- a/src/compute-client/src/controller/instance.rs +++ b/src/compute-client/src/controller/instance.rs @@ -216,6 +216,26 @@ pub(super) struct Instance { replica_rx: mz_ore::channel::InstrumentedUnboundedReceiver, } +/// Clones a read hold on an input of a collection of this instance. +/// +/// # Panics +/// +/// Panics if the hold's issuer hung up. The inputs are storage collections, whose issuer is the +/// `StorageCollections`, and indexes of this instance, whose holds `ComputeController::create_dataflow` +/// acquires from this instance's task, so that only happens during process shutdown. +fn clone_input_read_hold(hold: &ReadHold) -> ReadHold { + hold.try_clone() + .expect("an input read hold's issuer only hangs up during process shutdown") +} + +/// Clones read holds with [`clone_input_read_hold`]. +fn clone_input_read_holds(holds: &BTreeMap) -> BTreeMap { + holds + .iter() + .map(|(id, hold)| (*id, clone_input_read_hold(hold))) + .collect() +} + impl Instance { /// Acquire a handle to the collection state associated with `id`. fn collection(&self, id: GlobalId) -> Result<&CollectionState, CollectionMissing> { @@ -326,7 +346,11 @@ impl Instance { if target_replica.is_some_and(|id| id != replica.id) { continue; } - replica.add_collection(id, as_of.clone(), replica_input_read_holds.clone()); + let input_read_holds = replica_input_read_holds + .iter() + .map(clone_input_read_hold) + .collect(); + replica.add_collection(id, as_of.clone(), input_read_holds); } } @@ -1462,7 +1486,7 @@ impl Instance { for &id in dataflow.source_imports.keys() { let mut read_hold = import_read_holds.remove(&id).ok_or(ReadHoldMissing(id))?; - replica_input_read_holds.push(read_hold.clone()); + replica_input_read_holds.push(clone_input_read_hold(&read_hold)); read_hold .try_downgrade(as_of.clone()) @@ -1496,9 +1520,12 @@ impl Instance { export_id, as_of.clone(), shared, - storage_dependencies.clone(), - compute_dependencies.clone(), - replica_input_read_holds.clone(), + clone_input_read_holds(&storage_dependencies), + clone_input_read_holds(&compute_dependencies), + replica_input_read_holds + .iter() + .map(clone_input_read_hold) + .collect(), write_only, storage_sink, dataflow.initial_storage_as_of.clone(), diff --git a/src/storage-types/src/read_holds.rs b/src/storage-types/src/read_holds.rs index b1860d8f531c5..0b08db0cae817 100644 --- a/src/storage-types/src/read_holds.rs +++ b/src/storage-types/src/read_holds.rs @@ -52,7 +52,7 @@ impl Debug for ReadHold { } } -/// The issuer of a [`ReadHold`] has hung up, so the hold can no longer be cloned. +/// The issuer of a [`ReadHold`] has hung up, so the hold can no longer be cloned or merged. /// /// This happens during process shutdown, when the tokio runtime drops tasks in arbitrary order /// and the issuer task can disappear while holds still exist, and to compute holds once their @@ -107,11 +107,15 @@ impl ReadHold { /// Merges `other` into `self`, keeping the overall read hold. /// + /// Returns an `Err` when the issuer of the read hold has hung up, in which + /// case `self` no longer holds back the since of the collection, see + /// [`ReadHoldIssuerHungUp`]. + /// /// # Panics /// /// Panics when trying to merge a [ReadHold] for a different collection /// (different [GlobalId]). - pub fn merge_assign(&mut self, mut other: ReadHold) { + pub fn merge_assign(&mut self, mut other: ReadHold) -> Result<(), ReadHoldIssuerHungUp> { assert_eq!( self.id, other.id, "can only merge ReadHolds for the same ID" @@ -133,12 +137,7 @@ impl ReadHold { // in one go. changes.extend(self.since.iter().map(|t| (t.clone(), 1))); - match (self.change_tx)(self.id, changes) { - Ok(_) => (), - Err(e) => { - panic!("cannot merge ReadHold: {}", e); - } - } + (self.change_tx)(self.id, changes).map_err(|_| ReadHoldIssuerHungUp(self.id)) } /// Downgrades `self` to the given `frontier`. Returns `Err` when the new @@ -180,9 +179,8 @@ impl ReadHold { /// the since of the collection. /// /// The issuer hangs up during process shutdown, and for compute holds when - /// their compute instance is dropped (see [`ReadHoldIssuerHungUp`]). Callers - /// that may run concurrently with either can use this method to handle that - /// case gracefully, instead of panicking via [Clone::clone]. + /// their compute instance is dropped, see [`ReadHoldIssuerHungUp`]. This is + /// why [ReadHold] does not implement [Clone]. pub fn try_clone(&self) -> Result { if self.id.is_user() { tracing::trace!("cloning ReadHold on {}: {:?}", self.id, self.since); @@ -208,15 +206,6 @@ impl ReadHold { } } -impl Clone for ReadHold { - fn clone(&self) -> Self { - match self.try_clone() { - Ok(clone) => clone, - Err(e) => panic!("cannot clone ReadHold: {}", e), - } - } -} - impl Drop for ReadHold { fn drop(&mut self) { if self.id.is_user() { From b5b3d22a8b09b2cd52db768f98991fdbffe264f7 Mon Sep 17 00:00:00 2001 From: Gabor Gevay Date: Fri, 9 Oct 2026 11:02:18 +0200 Subject: [PATCH 3/3] adapter: Test read holds on a concurrently dropped cluster Two tests with new failpoints, each parking a SELECT on cluster `c`, running `DROP CLUSTER c CASCADE`, and resuming it: - `test_fast_path_peek_after_concurrent_cluster_drop` parks a fast-path SELECT at `peek_before_optimize`, after it acquired its read holds. - `test_transaction_read_holds_after_concurrent_cluster_drop` parks the second SELECT of a transaction at `txn_read_holds_before_dispatch`, before it fetches the transaction's read holds. The SELECTs fail with a "was dropped" error and environmentd stays up. Without the previous commit, the first test's session task panics, which closes the connection, and the second test's coordinator panics. Co-Authored-By: Claude Opus 5.5 --- src/adapter/src/frontend_peek.rs | 8 ++ src/environmentd/tests/sql.rs | 134 +++++++++++++++++++++++++++++++ 2 files changed, 142 insertions(+) diff --git a/src/adapter/src/frontend_peek.rs b/src/adapter/src/frontend_peek.rs index 4dbb37a75af56..699dfeb3fe044 100644 --- a/src/adapter/src/frontend_peek.rs +++ b/src/adapter/src/frontend_peek.rs @@ -705,6 +705,10 @@ impl PeekClient { // - Use the transaction's stored timestamp determination. // - Use the (relevant subset of the) transaction's read holds. + // Test-only synchronization point: parks a statement in a + // transaction before it fetches the transaction's read holds, so a + // test can drop their cluster in between. + fail::fail_point!("txn_read_holds_before_dispatch"); let txn_read_holds_opt = self .call_coordinator(|tx| Command::GetTransactionReadHoldsBundle { conn_id: session.conn_id().clone(), @@ -1003,6 +1007,10 @@ impl PeekClient { mz_ore::task::spawn_blocking( || "optimize peek", move || { + // Test-only synchronization point: parks a SELECT's + // optimization, after it acquired its read holds, so a test + // can drop their cluster in between. + fail::fail_point!("peek_before_optimize"); span.in_scope(|| { let _dispatch_guard = explain_ctx.dispatch_guard(); diff --git a/src/environmentd/tests/sql.rs b/src/environmentd/tests/sql.rs index 3fdf8cda0415f..20c6b3f376b86 100644 --- a/src/environmentd/tests/sql.rs +++ b/src/environmentd/tests/sql.rs @@ -1189,6 +1189,140 @@ async fn test_subscribe_shutdown() { // function exits, things are working correctly. } +/// Creates cluster `c` with an index on table `t` in it, for the cluster-drop tests below. +#[allow(clippy::disallowed_methods)] +async fn setup_indexed_table_in_cluster_c(client: &tokio_postgres::Client) { + for stmt in [ + "CREATE CLUSTER c REPLICAS (r1 (size 'scale=1,workers=1'))", + "CREATE TABLE t (a int)", + "INSERT INTO t VALUES (1)", + "CREATE DEFAULT INDEX t_idx IN CLUSTER c ON t", + ] { + client.batch_execute(stmt).await.unwrap(); + } +} + +// A fast-path SELECT whose cluster is dropped while the SELECT is optimized fails with a "was +// dropped" error, instead of panicking on the read hold of the index it peeks. +// +// The SELECT parks at the `peek_before_optimize` failpoint, which runs on a blocking thread, after +// it acquired its read holds. `DROP CLUSTER` only asks the cluster's instance task to shut down. If +// that task hasn't exited when the SELECT resumes, the read holds still work and the SELECT fails +// at the cluster lookup instead, so the test then passes even without the fix. The failpoint is +// process-global, so the test relies on running in its own process, as nextest runs it. +#[mz_ore::test(tokio::test(flavor = "multi_thread", worker_threads = 2))] +#[cfg_attr(miri, ignore)] // too slow +#[allow(clippy::disallowed_methods)] +async fn test_fast_path_peek_after_concurrent_cluster_drop() { + let server = test_util::TestHarness::default().start().await; + let ddl_client = server.connect().await.unwrap(); + setup_indexed_table_in_cluster_c(&ddl_client).await; + let reader = server.connect().await.unwrap(); + reader.batch_execute("SET cluster = c").await.unwrap(); + + // Arm only once setup is done, so it's the SELECT below that parks. + let parked = Arc::new(Barrier::new(2)); + let resume = Arc::new(Barrier::new(2)); + fail::cfg_callback("peek_before_optimize", { + let parked = Arc::clone(&parked); + let resume = Arc::clone(&resume); + move || { + parked.wait(); + resume.wait(); + } + }) + .unwrap(); + + let select = task::spawn(|| "reader_select", async move { + reader.query("SELECT * FROM t", &[]).await + }); + task::spawn_blocking(|| "wait_parked", move || parked.wait()).await; + + ddl_client + .batch_execute("DROP CLUSTER c CASCADE") + .await + .unwrap(); + + // Disarm before handing the SELECT back, so that no later statement parks. + fail::remove("peek_before_optimize"); + task::spawn_blocking(|| "resume", move || resume.wait()).await; + + // A session-task panic would close the connection, which `unwrap_db_error` rejects. + let err = select + .await + .expect_err("SELECT must fail on the dropped cluster") + .unwrap_db_error(); + assert_contains!(err.message(), "was dropped"); + ddl_client + .batch_execute("SELECT 1") + .await + .expect("environmentd is still up"); +} + +// A SELECT in a transaction whose cluster is dropped before the SELECT fetches the transaction's +// read holds fails with a "was dropped" error, instead of panicking the coordinator, which kept the +// holds on the dropped cluster. +// +// The SELECT parks at the `txn_read_holds_before_dispatch` failpoint, after it took its catalog +// snapshot. `DROP CLUSTER` only asks the cluster's instance task to shut down, so depending on when +// that task exits, the SELECT fails either at fetching the transaction's read holds or at looking +// up the cluster. The asserted property holds for both. The failpoint is process-global, so the +// test relies on running in its own process, as nextest runs it. +#[mz_ore::test(tokio::test(flavor = "multi_thread", worker_threads = 2))] +#[cfg_attr(miri, ignore)] // too slow +#[allow(clippy::disallowed_methods)] +async fn test_transaction_read_holds_after_concurrent_cluster_drop() { + let server = test_util::TestHarness::default().start().await; + let ddl_client = server.connect().await.unwrap(); + setup_indexed_table_in_cluster_c(&ddl_client).await; + let reader = server.connect().await.unwrap(); + reader.batch_execute("SET cluster = c").await.unwrap(); + reader.batch_execute("BEGIN").await.unwrap(); + // The first SELECT of the transaction stores read holds on the index in `c`. + reader.query("SELECT * FROM t", &[]).await.unwrap(); + + // Arm only once setup is done, so it's the SELECT below that parks. + let parked = Arc::new(Barrier::new(2)); + let resume = Arc::new(Barrier::new(2)); + fail::cfg_callback("txn_read_holds_before_dispatch", { + let parked = Arc::clone(&parked); + let resume = Arc::clone(&resume); + move || { + // This runs on a worker of the runtime that also has to run the `DROP CLUSTER` below. + // `block_in_place` hands this worker's queue to another thread before parking. + tokio::task::block_in_place(|| { + parked.wait(); + resume.wait(); + }); + } + }) + .unwrap(); + + let select = task::spawn(|| "reader_select", async move { + reader.query("SELECT * FROM t", &[]).await + }); + task::spawn_blocking(|| "wait_parked", move || parked.wait()).await; + + ddl_client + .batch_execute("DROP CLUSTER c CASCADE") + .await + .unwrap(); + + // Disarm before handing the SELECT back, so that no later statement parks. + fail::remove("txn_read_holds_before_dispatch"); + task::spawn_blocking(|| "resume", move || resume.wait()).await; + + let err = select + .await + .expect_err("SELECT must fail on the dropped cluster") + .unwrap_db_error(); + assert_contains!(err.message(), "was dropped"); + ddl_client + .batch_execute("SELECT 1") + .await + .expect("environmentd is still up"); +} + #[mz_ore::test] #[allow(clippy::disallowed_methods)] fn test_subscribe_table_rw_timestamps() {