Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 6 additions & 3 deletions src/adapter/src/command.rs
Original file line number Diff line number Diff line change
Expand Up @@ -267,16 +267,19 @@ pub enum Command {
tx: oneshot::Sender<Result<Option<mz_repr::Timestamp>, 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<Option<ReadHolds>>,
tx: oneshot::Sender<Result<Option<ReadHolds>, 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<Result<(), AdapterError>>,
},

ExecuteSlowPathPeek {
Expand Down
11 changes: 7 additions & 4 deletions src/adapter/src/coord/command_handler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}

Expand All @@ -426,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 {
Expand Down Expand Up @@ -1964,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
Expand Down
82 changes: 63 additions & 19 deletions src/adapter/src/coord/read_policy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -40,12 +41,28 @@ 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<GlobalId, ReadHold>,
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 {
Expand Down Expand Up @@ -102,22 +119,41 @@ impl ReadHolds {
self.compute_holds.remove(&(instance_id, id));
}

/// Returns a new ReadHolds containing only the holds for collections in `id_bundle`.
pub fn subset(&self, id_bundle: &CollectionIdBundle) -> ReadHolds {
/// 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<Self, AdapterError> {
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`,
/// moved out of `self`.
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);
}
}
}
Expand Down Expand Up @@ -164,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`.
Expand Down Expand Up @@ -401,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),
}
}
}
2 changes: 1 addition & 1 deletion src/adapter/src/coord/sequencer/inner/peek.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion src/adapter/src/coord/sequencer/inner/subscribe.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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));

Expand Down
23 changes: 18 additions & 5 deletions src/adapter/src/frontend_peek.rs
Original file line number Diff line number Diff line change
Expand Up @@ -705,12 +705,16 @@ 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(),
tx,
})
.await?;
.await??;

if let Some(txn_read_holds) = txn_read_holds_opt {
let allowed_id_bundle = txn_read_holds.id_bundle();
Expand Down Expand Up @@ -789,12 +793,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)
Expand Down Expand Up @@ -1002,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();

Expand Down Expand Up @@ -1643,7 +1652,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) =
<Coordinator as TimestampProvider>::determine_timestamp_for_inner(
session,
Expand All @@ -1653,7 +1666,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() {
Expand Down
14 changes: 7 additions & 7 deletions src/adapter/src/peek_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -369,7 +369,7 @@ impl PeekClient {
max_result_size: u64,
max_returned_query_size: Option<u64>,
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,
Expand Down Expand Up @@ -434,11 +434,12 @@ 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 };
// 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
.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,
Expand All @@ -461,9 +462,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,
Expand Down
37 changes: 32 additions & 5 deletions src/compute-client/src/controller/instance.rs
Original file line number Diff line number Diff line change
Expand Up @@ -216,6 +216,26 @@ pub(super) struct Instance {
replica_rx: mz_ore::channel::InstrumentedUnboundedReceiver<ReplicaResponse, IntCounter>,
}

/// 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<GlobalId, ReadHold>) -> BTreeMap<GlobalId, ReadHold> {
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> {
Expand Down Expand Up @@ -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);
}
}

Expand Down Expand Up @@ -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())
Expand Down Expand Up @@ -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(),
Expand Down
Loading
Loading