mz_storage/storage_state.rs
1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5
6//! Worker-local state for storage timely instances.
7//!
8//! One instance of a [`Worker`], along with its contained [`StorageState`], is
9//! part of an ensemble of storage workers that all run inside the same timely
10//! cluster. We call this worker a _storage worker_ to disambiguate it from
11//! other kinds of workers, potentially other components that might be sharing
12//! the same timely cluster.
13//!
14//! ## Controller and internal communication
15//!
16//! A worker receives _external_ [`StorageCommands`](StorageCommand) from the
17//! storage controller, via a channel. Storage workers also share an _internal_
18//! control/command fabric ([`internal_control`]). Internal commands go through
19//! a sequencer dataflow that ensures that all workers receive all commands in
20//! the same consistent order.
21//!
22//! We need to make sure that commands that cause dataflows to be rendered are
23//! processed in the same consistent order across all workers because timely
24//! requires this. To achieve this, we make sure that only internal commands can
25//! cause dataflows to be rendered. External commands (from the controller)
26//! cause internal commands to be broadcast (by only one worker), to get
27//! dataflows rendered.
28//!
29//! The internal command fabric is also used to broadcast messages from a local
30//! operator/worker to all workers. For example, when we need to tear down and
31//! restart a dataflow on all workers when an error is encountered.
32//!
33//! ## Async Storage Worker
34//!
35//! The storage worker has a companion [`AsyncStorageWorker`] that must be used
36//! when running code that requires `async`. This is needed because a timely
37//! main loop cannot run `async` code.
38//!
39//! ## Example flow of commands for `RunIngestion`
40//!
41//! With external commands, internal commands, and the async worker,
42//! understanding where and how commands from the controller are realized can
43//! get complicated. We will follow the complete flow for `RunIngestion`, as an
44//! example:
45//!
46//! 1. Worker receives a [`StorageCommand::RunIngestion`] command from the
47//! controller.
48//! 2. This command is processed in [`StorageState::handle_storage_command`].
49//! This step cannot render dataflows, because it does not have access to the
50//! timely worker. It will only set up state that stays over the whole
51//! lifetime of the source, such as the `reported_frontier`. Putting in place
52//! this reported frontier will enable frontier reporting for that source. We
53//! will not start reporting when we only see an internal command for
54//! rendering a dataflow, which can "overtake" the external `RunIngestion`
55//! command.
56//! 3. During processing of that command, we call
57//! [`AsyncStorageWorker::update_ingestion_frontiers`], which causes a command to
58//! be sent to the async worker.
59//! 4. We eventually get a response from the async worker:
60//! [`AsyncStorageWorkerResponse::IngestionFrontiersUpdated`].
61//! 5. This response is handled in [`Worker::handle_async_worker_response`].
62//! 6. Handling that response causes a
63//! [`InternalStorageCommand::CreateIngestionDataflow`] to be broadcast to
64//! all workers via the internal command fabric.
65//! 7. This message will be processed (on each worker) in
66//! [`Worker::handle_internal_storage_command`]. This is what will cause the
67//! required dataflow to be rendered on all workers.
68//!
69//! The process described above assumes that the `RunIngestion` is _not_ an
70//! update, i.e. it is in response to a `CREATE SOURCE`-like statement.
71//!
72//! The primary distinction when handling a `RunIngestion` that represents an
73//! update, is that it might fill out new internal state in the mid-level
74//! clients on the way toward being run.
75
76use std::cell::RefCell;
77use std::collections::{BTreeMap, BTreeSet, VecDeque};
78use std::path::PathBuf;
79use std::rc::Rc;
80use std::sync::Arc;
81use std::thread;
82use std::time::Duration;
83
84use fail::fail_point;
85use mz_ore::now::NowFn;
86use mz_ore::soft_assert_or_log;
87use mz_ore::tracing::TracingHandle;
88use mz_persist_client::batch::ProtoBatch;
89use mz_persist_client::cache::PersistClientCache;
90use mz_persist_client::operators::shard_source::ErrorHandler;
91use mz_repr::{GlobalId, Timestamp};
92use mz_rocksdb::config::SharedWriteBufferManager;
93use mz_storage_client::client::{
94 RunIngestionCommand, StatusUpdate, StorageCommand, StorageResponse,
95};
96use mz_storage_types::AlterCompatible;
97use mz_storage_types::configuration::StorageConfiguration;
98use mz_storage_types::connections::ConnectionContext;
99use mz_storage_types::controller::CollectionMetadata;
100use mz_storage_types::dyncfgs::STORAGE_SERVER_MAINTENANCE_INTERVAL;
101use mz_storage_types::oneshot_sources::OneshotIngestionDescription;
102use mz_storage_types::sinks::StorageSinkDesc;
103use mz_storage_types::sources::IngestionDescription;
104use mz_timely_util::builder_async::PressOnDropButton;
105use mz_txn_wal::operator::TxnsContext;
106use timely::order::PartialOrder;
107use timely::progress::Timestamp as _;
108use timely::progress::frontier::Antichain;
109use timely::worker::Worker as TimelyWorker;
110use tokio::sync::mpsc::error::TryRecvError;
111use tokio::sync::{mpsc, watch};
112use tokio::time::Instant;
113use tracing::{debug, info, warn};
114use uuid::Uuid;
115
116use crate::internal_control::{
117 self, DataflowParameters, InternalCommandReceiver, InternalCommandSender,
118 InternalStorageCommand,
119};
120use crate::metrics::StorageMetrics;
121use crate::statistics::{AggregatedStatistics, SinkStatistics, SourceStatistics};
122use crate::storage_state::async_storage_worker::{AsyncStorageWorker, AsyncStorageWorkerResponse};
123
124pub mod async_storage_worker;
125
126type CommandReceiver = mpsc::UnboundedReceiver<StorageCommand>;
127type ResponseSender = mpsc::UnboundedSender<StorageResponse>;
128
129/// State maintained for each worker thread.
130///
131/// Much of this state can be viewed as local variables for the worker thread,
132/// holding state that persists across function calls.
133pub struct Worker<'w> {
134 /// The underlying Timely worker.
135 ///
136 /// NOTE: This is `pub` for testing.
137 pub timely_worker: &'w mut TimelyWorker,
138 /// The channel over which communication handles for newly connected clients
139 /// are delivered.
140 pub client_rx: mpsc::UnboundedReceiver<(Uuid, CommandReceiver, ResponseSender)>,
141 /// The state associated with collection ingress and egress.
142 pub storage_state: StorageState,
143}
144
145impl<'w> Worker<'w> {
146 /// Creates new `Worker` state from the given components.
147 pub fn new(
148 timely_worker: &'w mut TimelyWorker,
149 client_rx: mpsc::UnboundedReceiver<(Uuid, CommandReceiver, ResponseSender)>,
150 metrics: StorageMetrics,
151 now: NowFn,
152 connection_context: ConnectionContext,
153 instance_context: StorageInstanceContext,
154 persist_clients: Arc<PersistClientCache>,
155 txns_ctx: TxnsContext,
156 tracing_handle: Arc<TracingHandle>,
157 shared_rocksdb_write_buffer_manager: SharedWriteBufferManager,
158 ) -> Self {
159 // It is very important that we only create the internal control
160 // flow/command sequencer once because a) the worker state is re-used
161 // when a new client connects and b) dataflows that have already been
162 // rendered into the timely worker are reused as well.
163 //
164 // If we created a new sequencer every time we get a new client (likely
165 // because the controller re-started and re-connected), dataflows that
166 // were rendered before would still hold a handle to the old sequencer
167 // but we would not read their commands anymore.
168 let (internal_cmd_tx, internal_cmd_rx) =
169 internal_control::setup_command_sequencer(timely_worker);
170
171 let storage_state = StorageState::new_guest(
172 timely_worker.index(),
173 timely_worker.peers(),
174 internal_cmd_tx,
175 Some(internal_cmd_rx),
176 metrics,
177 now,
178 connection_context,
179 instance_context,
180 persist_clients,
181 txns_ctx,
182 tracing_handle,
183 shared_rocksdb_write_buffer_manager,
184 );
185
186 // TODO(aljoscha): We might want `async_worker` and `internal_cmd_tx` to
187 // be fields of `Worker` instead of `StorageState`, but at least for the
188 // command flow sources and sinks need access to that. We can refactor
189 // this once we have a clearer boundary between what sources/sinks need
190 // and the full "power" of the internal command flow, which should stay
191 // internal to the worker/not be exposed to source/sink implementations.
192 Self {
193 timely_worker,
194 client_rx,
195 storage_state,
196 }
197 }
198}
199
200impl StorageState {
201 /// Creates per-worker storage state, for hosting on any Timely worker.
202 ///
203 /// The caller provides the internal command channel endpoints: the native storage worker wires
204 /// them to the sequencer dataflow, a foreign host wires the sender to its own sequencing
205 /// channel and passes no receiver, since it dispatches internal commands itself.
206 /// Must be called on the hosting worker's thread, because the async worker unparks the
207 /// creating thread.
208 pub fn new_guest(
209 timely_worker_index: usize,
210 timely_worker_peers: usize,
211 internal_cmd_tx: InternalCommandSender,
212 internal_cmd_rx: Option<InternalCommandReceiver>,
213 metrics: StorageMetrics,
214 now: NowFn,
215 connection_context: ConnectionContext,
216 instance_context: StorageInstanceContext,
217 persist_clients: Arc<PersistClientCache>,
218 txns_ctx: TxnsContext,
219 tracing_handle: Arc<TracingHandle>,
220 shared_rocksdb_write_buffer_manager: SharedWriteBufferManager,
221 ) -> Self {
222 let storage_configuration =
223 StorageConfiguration::new(connection_context, mz_dyncfgs::all_dyncfgs());
224
225 // We always initialize as read_only=true. Only when we're explicitly
226 // allowed do we switch to doing writes.
227 let (read_only_tx, read_only_rx) = watch::channel(true);
228
229 // Similar to the internal command sequencer, it is very important that
230 // we only create the async worker once because a) the worker state is
231 // re-used when a new client connects and b) commands that have already
232 // been sent and might yield a response will be lost if a new iteration
233 // of `run_client` creates a new async worker.
234 //
235 // If we created a new async worker every time we get a new client
236 // (likely because the controller re-started and re-connected), we can
237 // get into an inconsistent state where we think that a dataflow has
238 // been rendered, for example because there is an entry in
239 // `StorageState::ingestions`, while there is not yet a dataflow. This
240 // happens because the dataflow only gets rendered once we get a
241 // response from the async worker and send off an internal command.
242 //
243 // The core idea is that both the sequencer and the async worker are
244 // part of the per-worker state, and must be treated as such, meaning
245 // they must survive between invocations of `run_client`.
246
247 // TODO(aljoscha): This thread unparking business seems brittle, but that's
248 // also how the command channel works currently. We can wrap it inside a
249 // struct that holds both a channel and a `Thread`, but I don't
250 // think that would help too much.
251 let async_worker = async_storage_worker::AsyncStorageWorker::new(
252 thread::current(),
253 Arc::clone(&persist_clients),
254 );
255 let cluster_memory_limit = instance_context.cluster_memory_limit;
256
257 let storage_state = StorageState {
258 source_uppers: BTreeMap::new(),
259 source_tokens: BTreeMap::new(),
260 metrics,
261 reported_frontiers: BTreeMap::new(),
262 ingestions: BTreeMap::new(),
263 exports: BTreeMap::new(),
264 oneshot_ingestions: BTreeMap::new(),
265 now,
266 timely_worker_index,
267 timely_worker_peers,
268 instance_context,
269 persist_clients,
270 txns_ctx,
271 sink_tokens: BTreeMap::new(),
272 sink_write_frontiers: BTreeMap::new(),
273 dropped_ids: Vec::new(),
274 aggregated_statistics: AggregatedStatistics::new(
275 timely_worker_index,
276 timely_worker_peers,
277 ),
278 shared_status_updates: Default::default(),
279 latest_status_updates: Default::default(),
280 initial_status_reported: Default::default(),
281 internal_cmd_tx,
282 internal_cmd_rx,
283 read_only_tx,
284 read_only_rx,
285 async_worker,
286 storage_configuration,
287 dataflow_parameters: DataflowParameters::new(
288 shared_rocksdb_write_buffer_manager,
289 cluster_memory_limit,
290 ),
291 tracing_handle,
292 server_maintenance_interval: Duration::ZERO,
293 };
294
295 storage_state
296 }
297}
298
299/// Worker-local state related to the ingress or egress of collections of data.
300pub struct StorageState {
301 /// The highest observed upper frontier for collection.
302 ///
303 /// This is shared among all source instances, so that they can jointly advance the
304 /// frontier even as other instances are created and dropped. Ideally, the Storage
305 /// module would eventually provide one source of truth on this rather than multiple,
306 /// and we should aim for that but are not there yet.
307 pub source_uppers: BTreeMap<GlobalId, Rc<RefCell<Antichain<mz_repr::Timestamp>>>>,
308 /// Handles to created sources, keyed by ID
309 /// NB: The type of the tokens must not be changed to something other than `PressOnDropButton`
310 /// to prevent usage of custom shutdown tokens that are tricky to get right.
311 pub source_tokens: BTreeMap<GlobalId, Vec<PressOnDropButton>>,
312 /// Metrics for storage objects.
313 pub metrics: StorageMetrics,
314 /// Tracks the conditional write frontiers we have reported.
315 pub reported_frontiers: BTreeMap<GlobalId, Antichain<Timestamp>>,
316 /// Descriptions of each installed ingestion.
317 pub ingestions: BTreeMap<GlobalId, IngestionDescription<CollectionMetadata>>,
318 /// Descriptions of each installed export.
319 pub exports: BTreeMap<GlobalId, StorageSinkDesc<CollectionMetadata, mz_repr::Timestamp>>,
320 /// Descriptions of oneshot ingestions that are currently running.
321 pub oneshot_ingestions: BTreeMap<uuid::Uuid, OneshotIngestionDescription<ProtoBatch>>,
322 /// Undocumented
323 pub now: NowFn,
324 /// Index of the associated timely dataflow worker.
325 pub timely_worker_index: usize,
326 /// Peers in the associated timely dataflow worker.
327 pub timely_worker_peers: usize,
328 /// Other configuration for sources and sinks.
329 pub instance_context: StorageInstanceContext,
330 /// A process-global cache of (blob_uri, consensus_uri) -> PersistClient.
331 /// This is intentionally shared between workers
332 pub persist_clients: Arc<PersistClientCache>,
333 /// Context necessary for rendering txn-wal operators.
334 pub txns_ctx: TxnsContext,
335 /// Tokens that should be dropped when a dataflow is dropped to clean up
336 /// associated state.
337 /// NB: The type of the tokens must not be changed to something other than `PressOnDropButton`
338 /// to prevent usage of custom shutdown tokens that are tricky to get right.
339 pub sink_tokens: BTreeMap<GlobalId, Vec<PressOnDropButton>>,
340 /// Frontier of sink writes (all subsequent writes will be at times at or
341 /// equal to this frontier)
342 pub sink_write_frontiers: BTreeMap<GlobalId, Rc<RefCell<Antichain<Timestamp>>>>,
343 /// Collection ids that have been dropped but not yet reported as dropped
344 pub dropped_ids: Vec<GlobalId>,
345
346 /// Statistics for sources and sinks.
347 pub aggregated_statistics: AggregatedStatistics,
348
349 /// A place shared with running dataflows, so that health operators, can
350 /// report status updates back to us.
351 ///
352 /// **NOTE**: Operators that append to this collection should take care to only add new
353 /// status updates if the status of the ingestion/export in question has _changed_.
354 pub shared_status_updates: Rc<RefCell<Vec<StatusUpdate>>>,
355
356 /// The latest status update for each object.
357 pub latest_status_updates: BTreeMap<GlobalId, StatusUpdate>,
358
359 /// Whether we have reported the initial status after connecting to a new client.
360 /// This is reset to false when a new client connects.
361 pub initial_status_reported: bool,
362
363 /// Sender for cluster-internal storage commands. These can be sent from
364 /// within workers/operators and will be distributed to all workers. For
365 /// example, for shutting down an entire dataflow from within a
366 /// operator/worker.
367 pub internal_cmd_tx: InternalCommandSender,
368 /// Receiver for cluster-internal storage commands. `None` when the state is hosted outside
369 /// the storage server, whose host dispatches internal commands itself.
370 pub internal_cmd_rx: Option<InternalCommandReceiver>,
371
372 /// When this replica/cluster is in read-only mode it must not affect any
373 /// changes to external state. This flag can only be changed by a
374 /// [StorageCommand::AllowWrites].
375 ///
376 /// Everything running on this replica/cluster must obey this flag. At the
377 /// time of writing, nothing currently looks at this flag.
378 /// TODO(benesch): fix this.
379 ///
380 /// NOTE: In the future, we might want a more complicated flag, for example
381 /// something that tells us after which timestamp we are allowed to write.
382 /// In this first version we are keeping things as simple as possible!
383 pub read_only_rx: watch::Receiver<bool>,
384
385 /// Send-side for read-only state.
386 pub read_only_tx: watch::Sender<bool>,
387
388 /// Async worker companion, used for running code that requires async, which
389 /// the timely main loop cannot do.
390 pub async_worker: AsyncStorageWorker<mz_repr::Timestamp>,
391
392 /// Configuration for source and sink connections.
393 pub storage_configuration: StorageConfiguration,
394 /// Dynamically configurable parameters that control how dataflows are rendered.
395 /// NOTE(guswynn): we should consider moving these into `storage_configuration`.
396 pub dataflow_parameters: DataflowParameters,
397
398 /// A process-global handle to tracing configuration.
399 pub tracing_handle: Arc<TracingHandle>,
400
401 /// Interval at which to perform server maintenance tasks. Set to a zero interval to
402 /// perform maintenance with every `step_or_park` invocation.
403 pub server_maintenance_interval: Duration,
404}
405
406impl StorageState {
407 /// Return an error handler that triggers a suspend and restart of the corresponding storage
408 /// dataflow.
409 pub fn error_handler(&self, context: &'static str, id: GlobalId) -> ErrorHandler {
410 let tx = self.internal_cmd_tx.clone();
411 ErrorHandler::signal(move |e| {
412 tx.send(InternalStorageCommand::SuspendAndRestart {
413 id,
414 reason: format!("{context}: {e:#}"),
415 })
416 })
417 }
418}
419
420/// Extra context for a storage instance.
421/// This is extra information that is used when rendering source
422/// and sinks that is not tied to the source/connection configuration itself.
423#[derive(Clone)]
424pub struct StorageInstanceContext {
425 /// A directory that can be used for scratch work.
426 pub scratch_directory: Option<PathBuf>,
427 /// The memory limit of the materialize cluster replica. This will
428 /// be used to calculate and configure the maximum inflight bytes for backpressure
429 pub cluster_memory_limit: Option<usize>,
430}
431
432impl StorageInstanceContext {
433 /// Build a new `StorageInstanceContext`.
434 pub fn new(scratch_directory: Option<PathBuf>, cluster_memory_limit: Option<usize>) -> Self {
435 Self {
436 scratch_directory,
437 cluster_memory_limit,
438 }
439 }
440
441 /// Returns a `rocksdb::Env` for a new RocksDB instance.
442 ///
443 /// With a scratch directory this is the default `Env`, which stores data
444 /// on the host filesystem. Without one, RocksDB runs in memory, and every
445 /// call returns a fresh in-memory `Env`. State written through an `Env`
446 /// is only reachable through that same `Env`, so a per-instance `Env`
447 /// isolates instances from each other and from previous incarnations of
448 /// themselves. Background threads are process-wide either way, both
449 /// variants delegate them to the default `Env`.
450 pub fn rocksdb_env(&self) -> Result<rocksdb::Env, rocksdb::Error> {
451 if self.scratch_directory.is_some() {
452 rocksdb::Env::new()
453 } else {
454 rocksdb::Env::mem_env()
455 }
456 }
457}
458
459impl<'w> Worker<'w> {
460 /// Waits for client connections and runs them to completion.
461 pub fn run(&mut self) {
462 while let Some((_nonce, rx, tx)) = self.client_rx.blocking_recv() {
463 self.run_client(rx, tx);
464 }
465 }
466
467 /// Runs this (timely) storage worker until the given `command_rx` is
468 /// disconnected.
469 ///
470 /// See the [module documentation](crate::storage_state) for this
471 /// workers responsibilities, how it communicates with the other workers and
472 /// how commands flow from the controller and through the workers.
473 fn run_client(&mut self, mut command_rx: CommandReceiver, response_tx: ResponseSender) {
474 // At this point, all workers are still reading from the command flow.
475 if self.reconcile(&mut command_rx).is_err() {
476 return;
477 }
478
479 // The last time we reported statistics.
480 let mut last_stats_time = Instant::now();
481
482 // The last time we did periodic maintenance.
483 let mut last_maintenance = std::time::Instant::now();
484
485 let mut disconnected = false;
486 while !disconnected {
487 let config = &self.storage_state.storage_configuration;
488 let stats_interval = config.parameters.statistics_collection_interval;
489
490 let maintenance_interval = self.storage_state.server_maintenance_interval;
491
492 let now = std::time::Instant::now();
493 // Determine if we need to perform maintenance, which is true if `maintenance_interval`
494 // time has passed since the last maintenance.
495 let sleep_duration;
496 if now >= last_maintenance + maintenance_interval {
497 last_maintenance = now;
498 sleep_duration = None;
499
500 self.report_frontier_progress(&response_tx);
501 } else {
502 // We didn't perform maintenance, sleep until the next maintenance interval.
503 let next_maintenance = last_maintenance + maintenance_interval;
504 sleep_duration = Some(next_maintenance.saturating_duration_since(now))
505 }
506
507 // Ask Timely to execute a unit of work.
508 //
509 // If there are no pending commands or responses from the async
510 // worker, we ask Timely to park the thread if there's nothing to
511 // do. We rely on another thread unparking us when there's new work
512 // to be done, e.g., when sending a command or when new Kafka
513 // messages have arrived.
514 //
515 // It is critical that we allow Timely to park iff there are no
516 // pending commands or responses. The command may have already been
517 // consumed by the call to `client_rx.recv`. See:
518 // https://github.com/MaterializeInc/materialize/pull/13973#issuecomment-1200312212
519 if command_rx.is_empty() && self.storage_state.async_worker.is_empty() {
520 // Make sure we wake up again to report any pending statistics updates.
521 let mut park_duration = stats_interval.saturating_sub(last_stats_time.elapsed());
522 if let Some(sleep_duration) = sleep_duration {
523 park_duration = std::cmp::min(sleep_duration, park_duration);
524 }
525 self.timely_worker.step_or_park(Some(park_duration));
526 } else {
527 self.timely_worker.step();
528 }
529
530 // Rerport any dropped ids
531 for id in std::mem::take(&mut self.storage_state.dropped_ids) {
532 self.send_storage_response(&response_tx, StorageResponse::DroppedId(id));
533 }
534
535 self.process_oneshot_ingestions(&response_tx);
536
537 self.report_status_updates(&response_tx);
538
539 if last_stats_time.elapsed() >= stats_interval {
540 self.report_storage_statistics(&response_tx);
541 last_stats_time = Instant::now();
542 }
543
544 // Handle any received commands.
545 loop {
546 match command_rx.try_recv() {
547 Ok(cmd) => self.storage_state.handle_storage_command(cmd),
548 Err(TryRecvError::Empty) => break,
549 Err(TryRecvError::Disconnected) => {
550 disconnected = true;
551 break;
552 }
553 }
554 }
555
556 // Handle responses from the async worker.
557 while let Ok(response) = self.storage_state.async_worker.try_recv() {
558 self.handle_async_worker_response(response);
559 }
560
561 // Handle any received commands.
562 while let Some(command) = self
563 .storage_state
564 .internal_cmd_rx
565 .as_ref()
566 .expect("storage server always wires a receiver")
567 .try_recv()
568 {
569 self.handle_internal_storage_command(command);
570 }
571 }
572 }
573
574 /// Entry point for applying a response from the async storage worker.
575 pub fn handle_async_worker_response(
576 &self,
577 async_response: AsyncStorageWorkerResponse<mz_repr::Timestamp>,
578 ) {
579 // NOTE: If we want to share the load of async processing we
580 // have to change `handle_storage_command` and change this
581 // assert.
582 assert_eq!(
583 self.timely_worker.index(),
584 0,
585 "only worker #0 is doing async processing"
586 );
587 match async_response {
588 AsyncStorageWorkerResponse::IngestionFrontiersUpdated {
589 id,
590 ingestion_description,
591 as_of,
592 resume_uppers,
593 source_resume_uppers,
594 } => {
595 self.storage_state.internal_cmd_tx.send(
596 InternalStorageCommand::CreateIngestionDataflow {
597 id,
598 ingestion_description,
599 as_of,
600 resume_uppers,
601 source_resume_uppers,
602 },
603 );
604 }
605 AsyncStorageWorkerResponse::ExportFrontiersUpdated { id, description } => {
606 self.storage_state
607 .internal_cmd_tx
608 .send(InternalStorageCommand::RunSinkDataflow(id, description));
609 }
610 AsyncStorageWorkerResponse::DropDataflow(id) => {
611 self.storage_state
612 .internal_cmd_tx
613 .send(InternalStorageCommand::DropDataflow(vec![id]));
614 }
615 }
616 }
617
618 /// Entry point for applying an internal storage command.
619 pub fn handle_internal_storage_command(&mut self, internal_cmd: InternalStorageCommand) {
620 match internal_cmd {
621 InternalStorageCommand::SuspendAndRestart { id, reason } => {
622 info!(
623 "worker {}/{} initiating suspend-and-restart for {id} because of: {reason}",
624 self.timely_worker.index(),
625 self.timely_worker.peers(),
626 );
627
628 let maybe_ingestion = self.storage_state.ingestions.get(&id).cloned();
629 if let Some(ingestion_description) = maybe_ingestion {
630 // Yank the token of the previously existing source dataflow.Note that this
631 // token also includes any source exports/subsources.
632 let maybe_token = self.storage_state.source_tokens.remove(&id);
633 if maybe_token.is_none() {
634 // Something has dropped the source. Make sure we don't
635 // accidentally re-create it.
636 return;
637 }
638
639 // This needs to be done by one worker, which will
640 // broadcasts a `CreateIngestionDataflow` command to all
641 // workers based on the response that contains the
642 // resumption upper.
643 //
644 // Doing this separately on each worker could lead to
645 // differing resume_uppers which might lead to all kinds of
646 // mayhem.
647 //
648 // TODO(aljoscha): If we ever become worried that this is
649 // putting undue pressure on worker 0 we can pick the
650 // designated worker for a source/sink based on `id.hash()`.
651 if self.timely_worker.index() == 0 {
652 for (id, _) in ingestion_description.source_exports.iter() {
653 self.storage_state
654 .aggregated_statistics
655 .advance_global_epoch(*id);
656 }
657 self.storage_state
658 .async_worker
659 .update_ingestion_frontiers(id, ingestion_description);
660 }
661
662 // Continue with other commands.
663 return;
664 }
665
666 let maybe_sink = self.storage_state.exports.get(&id).cloned();
667 if let Some(sink_description) = maybe_sink {
668 // Yank the token of the previously existing sink
669 // dataflow.
670 let maybe_token = self.storage_state.sink_tokens.remove(&id);
671
672 if maybe_token.is_none() {
673 // Something has dropped the sink. Make sure we don't
674 // accidentally re-create it.
675 return;
676 }
677
678 // This needs to be broadcast by one worker and go through
679 // the internal command fabric, to ensure consistent
680 // ordering of dataflow rendering across all workers.
681 if self.timely_worker.index() == 0 {
682 self.storage_state
683 .aggregated_statistics
684 .advance_global_epoch(id);
685 self.storage_state
686 .async_worker
687 .update_sink_frontiers(id, sink_description);
688 }
689
690 // Continue with other commands.
691 return;
692 }
693
694 if !self
695 .storage_state
696 .ingestions
697 .values()
698 .any(|v| v.source_exports.contains_key(&id))
699 {
700 // Our current approach to dropping a source results in a race between shard
701 // finalization (which happens in the controller) and dataflow shutdown (which
702 // happens in clusterd). If a source is created and dropped fast enough -or the
703 // two commands get sufficiently delayed- then it's possible to receive a
704 // SuspendAndRestart command for an unknown source. We cannot assert that this
705 // never happens but we log an error here to track how often this happens.
706 warn!(
707 "got InternalStorageCommand::SuspendAndRestart for something that is not a source or sink: {id}"
708 );
709 }
710 }
711 InternalStorageCommand::CreateIngestionDataflow {
712 id: ingestion_id,
713 mut ingestion_description,
714 as_of,
715 mut resume_uppers,
716 mut source_resume_uppers,
717 } => {
718 info!(
719 ?as_of,
720 ?resume_uppers,
721 "worker {}/{} trying to (re-)start ingestion {ingestion_id}",
722 self.timely_worker.index(),
723 self.timely_worker.peers(),
724 );
725
726 // We initialize statistics before we prune finished exports. We
727 // still want to export statistics for these, plus the rendering
728 // machinery will get confused if there are not at least
729 // statistics for the "main" source.
730 for (export_id, export) in ingestion_description.source_exports.iter() {
731 let resume_upper = resume_uppers[export_id].clone();
732 self.storage_state.aggregated_statistics.initialize_source(
733 *export_id,
734 ingestion_id,
735 resume_upper.clone(),
736 || {
737 SourceStatistics::new(
738 *export_id,
739 self.storage_state.timely_worker_index,
740 &self.storage_state.metrics.source_statistics,
741 ingestion_id,
742 &export.storage_metadata.data_shard,
743 export.data_config.envelope.clone(),
744 resume_upper,
745 )
746 },
747 );
748 }
749
750 let finished_exports: BTreeSet<GlobalId> = resume_uppers
751 .iter()
752 .filter(|(_, frontier)| frontier.is_empty())
753 .map(|(id, _)| *id)
754 .collect();
755
756 resume_uppers.retain(|id, _| !finished_exports.contains(id));
757 source_resume_uppers.retain(|id, _| !finished_exports.contains(id));
758 ingestion_description
759 .source_exports
760 .retain(|id, _| !finished_exports.contains(id));
761
762 for id in ingestion_description.collection_ids() {
763 // If there is already a shared upper, we re-use it, to make
764 // sure that parties that are already using the shared upper
765 // can continue doing so.
766 let source_upper = self
767 .storage_state
768 .source_uppers
769 .entry(id.clone())
770 .or_insert_with(|| {
771 Rc::new(RefCell::new(Antichain::from_elem(Timestamp::minimum())))
772 });
773
774 let mut source_upper = source_upper.borrow_mut();
775 if !source_upper.is_empty() {
776 source_upper.clear();
777 source_upper.insert(mz_repr::Timestamp::minimum());
778 }
779 }
780
781 // If all subsources of the source are finished, we can skip rendering entirely.
782 // Also, if `as_of` is empty, the dataflow has been finalized, so we can skip it as
783 // well.
784 //
785 // TODO(guswynn|petrosagg): this is a bit hacky, and is a consequence of storage state
786 // management being a bit of a mess. we should clean this up and remove weird if
787 // statements like this.
788 if resume_uppers.values().all(|frontier| frontier.is_empty()) || as_of.is_empty() {
789 info!(
790 ?resume_uppers,
791 ?as_of,
792 "worker {}/{} skipping building ingestion dataflow \
793 for {ingestion_id} because the ingestion is finished",
794 self.timely_worker.index(),
795 self.timely_worker.peers(),
796 );
797 return;
798 }
799
800 crate::render::build_ingestion_dataflow(
801 self.timely_worker,
802 &mut self.storage_state,
803 ingestion_id,
804 ingestion_description,
805 as_of,
806 resume_uppers,
807 source_resume_uppers,
808 );
809 }
810 InternalStorageCommand::RunOneshotIngestion {
811 ingestion_id,
812 collection_id,
813 collection_meta,
814 request,
815 } => {
816 crate::render::build_oneshot_ingestion_dataflow(
817 self.timely_worker,
818 &mut self.storage_state,
819 ingestion_id,
820 collection_id,
821 collection_meta,
822 request,
823 );
824 }
825 InternalStorageCommand::RunSinkDataflow(sink_id, sink_description) => {
826 info!(
827 "worker {}/{} trying to (re-)start sink {sink_id}",
828 self.timely_worker.index(),
829 self.timely_worker.peers(),
830 );
831
832 {
833 // If there is already a shared write frontier, we re-use it, to
834 // make sure that parties that are already using the shared
835 // frontier can continue doing so.
836 let sink_write_frontier = self
837 .storage_state
838 .sink_write_frontiers
839 .entry(sink_id.clone())
840 .or_insert_with(|| Rc::new(RefCell::new(Antichain::new())));
841
842 let mut sink_write_frontier = sink_write_frontier.borrow_mut();
843 sink_write_frontier.clear();
844 sink_write_frontier.insert(mz_repr::Timestamp::minimum());
845 }
846 self.storage_state
847 .aggregated_statistics
848 .initialize_sink(sink_id, || {
849 SinkStatistics::new(
850 sink_id,
851 self.storage_state.timely_worker_index,
852 &self.storage_state.metrics.sink_statistics,
853 )
854 });
855
856 crate::render::build_export_dataflow(
857 self.timely_worker,
858 &mut self.storage_state,
859 sink_id,
860 sink_description,
861 );
862 }
863 InternalStorageCommand::DropDataflow(ids) => {
864 for id in &ids {
865 // Clean up per-source / per-sink state.
866 self.storage_state.source_uppers.remove(id);
867 self.storage_state.source_tokens.remove(id);
868
869 self.storage_state.sink_tokens.remove(id);
870 self.storage_state.sink_write_frontiers.remove(id);
871
872 self.storage_state.aggregated_statistics.deinitialize(*id);
873 }
874 }
875 InternalStorageCommand::UpdateConfiguration { storage_parameters } => {
876 self.storage_state
877 .dataflow_parameters
878 .update(storage_parameters.clone());
879 self.storage_state
880 .storage_configuration
881 .update(storage_parameters);
882
883 // Clear out the updates as we no longer forward them to anyone else to process.
884 // We clone `StorageState::storage_configuration` many times during rendering
885 // and want to avoid cloning these unused updates.
886 self.storage_state
887 .storage_configuration
888 .parameters
889 .dyncfg_updates = Default::default();
890
891 // Remember the maintenance interval locally to avoid reading it from the config set on
892 // every server iteration.
893 self.storage_state.server_maintenance_interval =
894 STORAGE_SERVER_MAINTENANCE_INTERVAL
895 .get(self.storage_state.storage_configuration.config_set());
896
897 // Apply storage's upsert spill flag to both stash flavors'
898 // mechanisms: the storage leg of the process-wide chunk
899 // spill gate (chunked flavor) and the storage-owned column
900 // pager (paged flavor). The buffer pool, the pager pool, and
901 // their budgets are the shared ones configured by compute's
902 // `apply_worker_config` (compute and storage run in the same
903 // process). The chunk gate ORs storage's leg with compute's,
904 // so chunks spill while either subsystem's flag is set.
905 //
906 // The flag is replica-scoped: the storage controller merges
907 // per-replica overrides into the `UpdateConfiguration`
908 // commands it sends, so reading this worker's `ConfigSet`
909 // here observes them.
910 {
911 use mz_storage_types::dyncfgs::ENABLE_UPSERT_PAGED_SPILL;
912
913 let enabled = ENABLE_UPSERT_PAGED_SPILL
914 .get(self.storage_state.storage_configuration.config_set());
915 debug!(
916 worker = self.timely_worker.index(),
917 enabled, "upsert stash spill: applying gate",
918 );
919 crate::upsert::upsert_stash_spill::set_enabled(enabled);
920 crate::upsert::upsert_stash_pager::set_enabled(enabled);
921 }
922 }
923 InternalStorageCommand::StatisticsUpdate { sources, sinks } => self
924 .storage_state
925 .aggregated_statistics
926 .ingest(sources, sinks),
927 }
928 }
929
930 /// Emit information about write frontier progress, along with information that should
931 /// be made durable for this to be the case.
932 ///
933 /// The write frontier progress is "conditional" in that it is not until the information is made
934 /// durable that the data are emitted to downstream workers, and indeed they should not rely on
935 /// the completeness of what they hear until the information is made durable.
936 ///
937 /// Specifically, this sends information about new timestamp bindings created by dataflow workers,
938 /// with the understanding if that if made durable (and ack'd back to the workers) the source will
939 /// in fact progress with this write frontier.
940 pub fn report_frontier_progress(&mut self, response_tx: &ResponseSender) {
941 let mut new_uppers = Vec::new();
942
943 // Check if any observed frontier should advance the reported frontiers.
944 for (id, frontier) in self
945 .storage_state
946 .source_uppers
947 .iter()
948 .chain(self.storage_state.sink_write_frontiers.iter())
949 {
950 let Some(reported_frontier) = self.storage_state.reported_frontiers.get_mut(id) else {
951 // Frontier reporting has not yet been started for this object.
952 // Potentially because this timely worker has not yet seen the
953 // `CreateSources` command.
954 continue;
955 };
956
957 let observed_frontier = frontier.borrow();
958
959 // Only do a thing if it *advances* the frontier, not just *changes* the frontier.
960 // This is protection against `frontier` lagging behind what we have conditionally reported.
961 if PartialOrder::less_than(reported_frontier, &observed_frontier) {
962 new_uppers.push((*id, observed_frontier.clone()));
963 reported_frontier.clone_from(&observed_frontier);
964 }
965 }
966
967 for (id, upper) in new_uppers {
968 self.send_storage_response(response_tx, StorageResponse::FrontierUpper(id, upper));
969 }
970 }
971
972 /// Pumps latest status updates from the buffer shared with operators and
973 /// reports any updates that need reporting.
974 pub fn report_status_updates(&mut self, response_tx: &ResponseSender) {
975 // If we haven't done the initial status report, report all current statuses
976 if !self.storage_state.initial_status_reported {
977 // We pull initially reported status updates to "now", so that they
978 // sort as the latest update in internal status collections. This
979 // makes it so that a newly bootstrapped envd can append status
980 // updates to internal status collections that report an accurate
981 // view as of the time when they came up.
982 let now_ts = mz_ore::now::to_datetime((self.storage_state.now)());
983 let status_updates = self
984 .storage_state
985 .latest_status_updates
986 .values()
987 .cloned()
988 .map(|mut update| {
989 update.timestamp = now_ts.clone();
990 update
991 });
992 for update in status_updates {
993 self.send_storage_response(response_tx, StorageResponse::StatusUpdate(update));
994 }
995 self.storage_state.initial_status_reported = true;
996 }
997
998 // Pump updates into our state and stage them for reporting.
999 for shared_update in self.storage_state.shared_status_updates.take() {
1000 self.send_storage_response(
1001 response_tx,
1002 StorageResponse::StatusUpdate(shared_update.clone()),
1003 );
1004
1005 self.storage_state
1006 .latest_status_updates
1007 .insert(shared_update.id, shared_update);
1008 }
1009 }
1010
1011 /// Report source statistics back to the controller.
1012 pub fn report_storage_statistics(&mut self, response_tx: &ResponseSender) {
1013 let (sources, sinks) = self.storage_state.aggregated_statistics.emit_local();
1014 if !sources.is_empty() || !sinks.is_empty() {
1015 self.storage_state
1016 .internal_cmd_tx
1017 .send(InternalStorageCommand::StatisticsUpdate { sources, sinks })
1018 }
1019
1020 let (sources, sinks) = self.storage_state.aggregated_statistics.snapshot();
1021 if !sources.is_empty() || !sinks.is_empty() {
1022 self.send_storage_response(
1023 response_tx,
1024 StorageResponse::StatisticsUpdates(sources, sinks),
1025 );
1026 }
1027 }
1028
1029 /// Send a response to the coordinator.
1030 pub fn send_storage_response(&self, response_tx: &ResponseSender, response: StorageResponse) {
1031 // Ignore send errors because the coordinator is free to ignore our
1032 // responses. This happens during shutdown.
1033 let _ = response_tx.send(response);
1034 }
1035
1036 /// Forward completed oneshot ingestion results to the coordinator.
1037 pub fn process_oneshot_ingestions(&mut self, response_tx: &ResponseSender) {
1038 for (ingestion_id, ingestion_state) in &mut self.storage_state.oneshot_ingestions {
1039 loop {
1040 match ingestion_state.results.try_recv() {
1041 Ok(result) => {
1042 let response = match result {
1043 Ok(maybe_batch) => maybe_batch.into_iter().map(Result::Ok).collect(),
1044 Err(err) => vec![Err(err)],
1045 };
1046 let staged_batches = BTreeMap::from([(*ingestion_id, response)]);
1047 let _ = response_tx.send(StorageResponse::StagedBatches(staged_batches));
1048 }
1049 Err(TryRecvError::Empty) => {
1050 break;
1051 }
1052 Err(TryRecvError::Disconnected) => {
1053 break;
1054 }
1055 }
1056 }
1057 }
1058 }
1059
1060 /// Extract commands until `InitializationComplete`, and make the worker
1061 /// reflect those commands. If the worker can not be made to reflect the
1062 /// commands, return an error.
1063 fn reconcile(&mut self, command_rx: &mut CommandReceiver) -> Result<(), ()> {
1064 // To initialize the connection, we want to drain all commands until we
1065 // receive a `StorageCommand::InitializationComplete` command to form a
1066 // target command state.
1067 let mut commands = vec![];
1068 loop {
1069 match command_rx.blocking_recv().ok_or(())? {
1070 StorageCommand::InitializationComplete => break,
1071 command => commands.push(command),
1072 }
1073 }
1074
1075 self.reconcile_commands(commands);
1076 Ok(())
1077 }
1078
1079 /// Reconciles the worker state with the given target command state,
1080 /// which the caller has drained from a new client connection up to (exclusive) the
1081 /// `InitializationComplete` marker.
1082 pub fn reconcile_commands(&mut self, mut commands: Vec<StorageCommand>) {
1083 let worker_id = self.timely_worker.index();
1084
1085 // Track which frontiers this envd expects; we will also set their
1086 // initial timestamp to the minimum timestamp to reset them as we don't
1087 // know what frontiers the new envd expects.
1088 let mut expected_objects = BTreeSet::new();
1089
1090 let mut drop_commands = BTreeSet::new();
1091 let mut running_ingestion_descriptions = self.storage_state.ingestions.clone();
1092 let mut running_exports_descriptions = self.storage_state.exports.clone();
1093
1094 let mut create_oneshot_ingestions: BTreeSet<Uuid> = BTreeSet::new();
1095 let mut cancel_oneshot_ingestions: BTreeSet<Uuid> = BTreeSet::new();
1096
1097 for command in &mut commands {
1098 match command {
1099 StorageCommand::Hello { .. } => {
1100 panic!("Hello must be captured before")
1101 }
1102 StorageCommand::AllowCompaction(id, since) => {
1103 info!(%worker_id, ?id, ?since, "reconcile: received AllowCompaction command");
1104
1105 // collect all "drop commands". These are `AllowCompaction`
1106 // commands that compact to the empty since. Then, later, we make sure
1107 // we retain only those `Create*` commands that are not dropped. We
1108 // assume that the `AllowCompaction` command is ordered after the
1109 // `Create*` commands but don't assert that.
1110 // WIP: Should we assert?
1111 if since.is_empty() {
1112 drop_commands.insert(*id);
1113 }
1114 }
1115 StorageCommand::RunIngestion(ingestion) => {
1116 info!(%worker_id, ?ingestion, "reconcile: received RunIngestion command");
1117
1118 // Ensure that ingestions are forward-rolling alter compatible.
1119 let prev = running_ingestion_descriptions
1120 .insert(ingestion.id, ingestion.description.clone());
1121
1122 if let Some(prev_ingest) = prev {
1123 // If the new ingestion is not exactly equal to the currently running
1124 // ingestion, we must either track that we need to synthesize an update
1125 // command to change the ingestion, or panic.
1126 prev_ingest
1127 .alter_compatible(ingestion.id, &ingestion.description)
1128 .expect("only alter compatible ingestions permitted");
1129 }
1130 }
1131 StorageCommand::RunSink(export) => {
1132 info!(%worker_id, ?export, "reconcile: received RunSink command");
1133
1134 // Ensure that exports are forward-rolling alter compatible.
1135 let prev =
1136 running_exports_descriptions.insert(export.id, export.description.clone());
1137
1138 if let Some(prev_export) = prev {
1139 prev_export
1140 .alter_compatible(export.id, &export.description)
1141 .expect("only alter compatible exports permitted");
1142 }
1143 }
1144 StorageCommand::RunOneshotIngestion(ingestion) => {
1145 info!(
1146 %worker_id,
1147 ingestion_id = %ingestion.ingestion_id,
1148 collection_id = %ingestion.collection_id,
1149 "reconcile: received RunOneshotIngestion command",
1150 );
1151 create_oneshot_ingestions.insert(ingestion.ingestion_id);
1152 }
1153 StorageCommand::CancelOneshotIngestion(uuid) => {
1154 info!(%worker_id, %uuid, "reconcile: received CancelOneshotIngestion command");
1155 cancel_oneshot_ingestions.insert(*uuid);
1156 }
1157 StorageCommand::InitializationComplete
1158 | StorageCommand::AllowWrites
1159 | StorageCommand::UpdateConfiguration(_) => (),
1160 }
1161 }
1162
1163 let mut seen_most_recent_definition = BTreeSet::new();
1164
1165 // We iterate over this backward to ensure that we keep only the most recent ingestion
1166 // description.
1167 let mut filtered_commands = VecDeque::new();
1168 for mut command in commands.into_iter().rev() {
1169 let mut should_keep = true;
1170 match &mut command {
1171 StorageCommand::Hello { .. } => {
1172 panic!("Hello must be captured before")
1173 }
1174 StorageCommand::RunIngestion(ingestion) => {
1175 // Subsources can be dropped independently of their
1176 // primary source, so we evaluate them in a separate
1177 // loop.
1178 for export_id in ingestion
1179 .description
1180 .source_exports
1181 .keys()
1182 .filter(|export_id| **export_id != ingestion.id)
1183 {
1184 if drop_commands.remove(export_id) {
1185 info!(%worker_id, %export_id, "reconcile: dropping subsource");
1186 self.storage_state.dropped_ids.push(*export_id);
1187 }
1188 }
1189
1190 if drop_commands.remove(&ingestion.id)
1191 || self.storage_state.dropped_ids.contains(&ingestion.id)
1192 {
1193 info!(%worker_id, %ingestion.id, "reconcile: dropping ingestion");
1194
1195 // If an ingestion is dropped, so too must all of
1196 // its subsources (i.e. ingestion exports, as well
1197 // as its progress subsource).
1198 for id in ingestion.description.collection_ids() {
1199 drop_commands.remove(&id);
1200 self.storage_state.dropped_ids.push(id);
1201 }
1202 should_keep = false;
1203 } else {
1204 let most_recent_defintion =
1205 seen_most_recent_definition.insert(ingestion.id);
1206
1207 if most_recent_defintion {
1208 // If this is the most recent definition, this
1209 // is what we will be running when
1210 // reconciliation completes. This definition
1211 // must not include any dropped subsources.
1212 ingestion.description.source_exports.retain(|export_id, _| {
1213 !self.storage_state.dropped_ids.contains(export_id)
1214 });
1215
1216 // After clearing any dropped subsources, we can
1217 // state that we expect all of these to exist.
1218 expected_objects.extend(ingestion.description.collection_ids());
1219 }
1220
1221 let running_ingestion = self.storage_state.ingestions.get(&ingestion.id);
1222
1223 // We keep only:
1224 // - The most recent version of the ingestion, which
1225 // is why these commands are run in reverse.
1226 // - Ingestions whose descriptions are not exactly
1227 // those that are currently running.
1228 should_keep = most_recent_defintion
1229 && running_ingestion != Some(&ingestion.description)
1230 }
1231 }
1232 StorageCommand::RunSink(export) => {
1233 if drop_commands.remove(&export.id)
1234 // If there were multiple `RunSink` in the command
1235 // stream, we want to ensure none of them are
1236 // retained.
1237 || self.storage_state.dropped_ids.contains(&export.id)
1238 {
1239 info!(%worker_id, %export.id, "reconcile: dropping sink");
1240
1241 // Make sure that we report back that the ID was
1242 // dropped.
1243 self.storage_state.dropped_ids.push(export.id);
1244
1245 should_keep = false
1246 } else {
1247 expected_objects.insert(export.id);
1248
1249 let running_sink = self.storage_state.exports.get(&export.id);
1250
1251 // We keep only:
1252 // - The most recent version of the sink, which
1253 // is why these commands are run in reverse.
1254 // - Sinks whose descriptions are not exactly
1255 // those that are currently running.
1256 should_keep = seen_most_recent_definition.insert(export.id)
1257 && running_sink != Some(&export.description);
1258 }
1259 }
1260 StorageCommand::RunOneshotIngestion(ingestion) => {
1261 let already_running = self
1262 .storage_state
1263 .oneshot_ingestions
1264 .contains_key(&ingestion.ingestion_id);
1265 let was_canceled = cancel_oneshot_ingestions.contains(&ingestion.ingestion_id);
1266
1267 should_keep = !already_running && !was_canceled;
1268 }
1269 StorageCommand::CancelOneshotIngestion(ingestion_id) => {
1270 let already_running = self
1271 .storage_state
1272 .oneshot_ingestions
1273 .contains_key(ingestion_id);
1274 should_keep = already_running;
1275 }
1276 StorageCommand::InitializationComplete
1277 | StorageCommand::AllowWrites
1278 | StorageCommand::UpdateConfiguration(_)
1279 | StorageCommand::AllowCompaction(_, _) => (),
1280 }
1281 if should_keep {
1282 filtered_commands.push_front(command);
1283 }
1284 }
1285 let commands = filtered_commands;
1286
1287 // Make sure all the "drop commands" matched up with a source or sink.
1288 // This is also what the regular handler logic for `AllowCompaction`
1289 // would do.
1290 soft_assert_or_log!(
1291 drop_commands.is_empty(),
1292 "AllowCompaction commands for non-existent IDs {:?}",
1293 drop_commands
1294 );
1295
1296 // Determine the ID of all objects we did _not_ see; these are
1297 // considered stale.
1298 let stale_objects = self
1299 .storage_state
1300 .ingestions
1301 .values()
1302 .map(|i| i.collection_ids())
1303 .flatten()
1304 .chain(self.storage_state.exports.keys().copied())
1305 // Objects are considered stale if we did not see them re-created.
1306 .filter(|id| !expected_objects.contains(id))
1307 .collect::<Vec<_>>();
1308 let stale_oneshot_ingestions = self
1309 .storage_state
1310 .oneshot_ingestions
1311 .keys()
1312 .filter(|ingestion_id| {
1313 let to_create = create_oneshot_ingestions.contains(ingestion_id);
1314 let to_drop = cancel_oneshot_ingestions.contains(ingestion_id);
1315 mz_ore::soft_assert_or_log!(
1316 !(!to_create && to_drop),
1317 "attempting to drop oneshot source {ingestion_id} that is not expected to be created during reconciliation"
1318 );
1319 !to_create && !to_drop
1320 })
1321 .copied()
1322 .collect::<Vec<_>>();
1323
1324 info!(
1325 %worker_id, ?expected_objects, ?stale_objects, ?stale_oneshot_ingestions,
1326 "reconcile: modifing storage state to match expected objects",
1327 );
1328
1329 for id in stale_objects {
1330 self.storage_state.drop_collection(id);
1331 }
1332 for id in stale_oneshot_ingestions {
1333 self.storage_state.drop_oneshot_ingestion(id);
1334 }
1335
1336 // Do not report dropping any objects that do not belong to expected
1337 // objects.
1338 self.storage_state
1339 .dropped_ids
1340 .retain(|id| expected_objects.contains(id));
1341
1342 // Do not report any frontiers that do not belong to expected objects.
1343 // Note that this set of objects can differ from the set of sources and
1344 // sinks.
1345 self.storage_state
1346 .reported_frontiers
1347 .retain(|id, _| expected_objects.contains(id));
1348
1349 // Reset the reported frontiers for the remaining objects.
1350 for (_, frontier) in &mut self.storage_state.reported_frontiers {
1351 *frontier = Antichain::from_elem(<_>::minimum());
1352 }
1353
1354 // Reset the initial status reported flag when a new client connects
1355 self.storage_state.initial_status_reported = false;
1356
1357 // Execute the modified commands.
1358 for command in commands {
1359 self.storage_state.handle_storage_command(command);
1360 }
1361 }
1362}
1363
1364impl StorageState {
1365 /// Entry point for applying a storage command.
1366 ///
1367 /// NOTE: This does not have access to the timely worker and therefore
1368 /// cannot render dataflows. For dataflow rendering, this needs to either
1369 /// send asynchronous command to the `async_worker` or internal
1370 /// commands to the `internal_cmd_tx`.
1371 pub fn handle_storage_command(&mut self, cmd: StorageCommand) {
1372 match cmd {
1373 StorageCommand::Hello { .. } => panic!("Hello must be captured before"),
1374 StorageCommand::InitializationComplete => (),
1375 StorageCommand::AllowWrites => {
1376 self.read_only_tx
1377 .send(false)
1378 .expect("we're holding one other end");
1379 self.persist_clients.cfg().enable_compaction();
1380 }
1381 StorageCommand::UpdateConfiguration(params) => {
1382 // These can be done from all workers safely.
1383 debug!("Applying configuration update: {params:?}");
1384
1385 // We serialize the dyncfg updates in StorageParameters, but configure
1386 // persist separately.
1387 self.persist_clients
1388 .cfg()
1389 .apply_from(¶ms.dyncfg_updates);
1390
1391 params.tracing.apply(self.tracing_handle.as_ref());
1392
1393 if let Some(log_filter) = ¶ms.tracing.log_filter {
1394 self.storage_configuration
1395 .connection_context
1396 .librdkafka_log_level =
1397 mz_ore::tracing::crate_level(&log_filter.clone().into(), "librdkafka");
1398 }
1399
1400 // This needs to be broadcast by one worker and go through
1401 // the internal command fabric, to ensure consistent
1402 // ordering of dataflow rendering across all workers.
1403 if self.timely_worker_index == 0 {
1404 self.internal_cmd_tx
1405 .send(InternalStorageCommand::UpdateConfiguration {
1406 storage_parameters: *params,
1407 })
1408 }
1409 }
1410 StorageCommand::RunIngestion(ingestion) => {
1411 let RunIngestionCommand { id, description } = *ingestion;
1412
1413 // Remember the ingestion description to facilitate possible
1414 // reconciliation later.
1415 self.ingestions.insert(id, description.clone());
1416
1417 // Initialize shared frontier reporting.
1418 for id in description.collection_ids() {
1419 self.reported_frontiers
1420 .entry(id)
1421 .or_insert_with(|| Antichain::from_elem(mz_repr::Timestamp::minimum()));
1422 }
1423
1424 // This needs to be done by one worker, which will broadcasts a
1425 // `CreateIngestionDataflow` command to all workers based on the response that
1426 // contains the resumption upper.
1427 //
1428 // Doing this separately on each worker could lead to differing resume_uppers
1429 // which might lead to all kinds of mayhem.
1430 //
1431 // n.b. the ingestion on each worker uses the description from worker 0––not the
1432 // ingestion in the local storage state. This is something we might have
1433 // interest in fixing in the future, e.g. materialize#19907
1434 if self.timely_worker_index == 0 {
1435 self.async_worker
1436 .update_ingestion_frontiers(id, description);
1437 }
1438 }
1439 StorageCommand::RunOneshotIngestion(oneshot) => {
1440 if self.timely_worker_index == 0 {
1441 self.internal_cmd_tx
1442 .send(InternalStorageCommand::RunOneshotIngestion {
1443 ingestion_id: oneshot.ingestion_id,
1444 collection_id: oneshot.collection_id,
1445 collection_meta: oneshot.collection_meta,
1446 request: oneshot.request,
1447 });
1448 }
1449 }
1450 StorageCommand::CancelOneshotIngestion(id) => {
1451 self.drop_oneshot_ingestion(id);
1452 }
1453 StorageCommand::RunSink(export) => {
1454 // Remember the sink description to facilitate possible
1455 // reconciliation later.
1456 let prev = self.exports.insert(export.id, export.description.clone());
1457
1458 // New sink, add state.
1459 if prev.is_none() {
1460 self.reported_frontiers.insert(
1461 export.id,
1462 Antichain::from_elem(mz_repr::Timestamp::minimum()),
1463 );
1464 }
1465
1466 // This needs to be broadcast by one worker and go through the internal command
1467 // fabric, to ensure consistent ordering of dataflow rendering across all
1468 // workers.
1469 if self.timely_worker_index == 0 {
1470 self.internal_cmd_tx
1471 .send(InternalStorageCommand::RunSinkDataflow(
1472 export.id,
1473 export.description,
1474 ));
1475 }
1476 }
1477 StorageCommand::AllowCompaction(id, frontier) => {
1478 soft_assert_or_log!(
1479 self.exports.contains_key(&id) || self.reported_frontiers.contains_key(&id),
1480 "AllowCompaction command for non-existent {id}"
1481 );
1482
1483 if frontier.is_empty() {
1484 // Indicates that we may drop `id`, as there are no more valid times to read.
1485 self.drop_collection(id);
1486 }
1487 }
1488 }
1489 }
1490
1491 /// Drop the identified storage collection from the storage state.
1492 fn drop_collection(&mut self, id: GlobalId) {
1493 fail_point!("crash_on_drop");
1494
1495 self.ingestions.remove(&id);
1496 self.exports.remove(&id);
1497
1498 let _ = self.latest_status_updates.remove(&id);
1499
1500 // This will stop reporting of frontiers.
1501 //
1502 // If this object still has its frontiers reported, we will notify the
1503 // client envd of the drop.
1504 if self.reported_frontiers.remove(&id).is_some() {
1505 // The only actions left are internal cleanup, so we can commit to
1506 // the client that these objects have been dropped.
1507 //
1508 // This must be done now rather than in response to `DropDataflow`,
1509 // otherwise we introduce the possibility of a timing issue where:
1510 // - We remove all tracking state from the storage state and send
1511 // `DropDataflow` (i.e. this block).
1512 // - While waiting to process that command, we reconcile with a new
1513 // envd. That envd has already committed to its catalog that this
1514 // object no longer exists.
1515 // - We process the `DropDataflow` command, and identify that this
1516 // object has been dropped.
1517 // - The next time `dropped_ids` is processed, we send a response
1518 // that this ID has been dropped, but the upstream state has no
1519 // record of that object having ever existed.
1520 self.dropped_ids.push(id);
1521 }
1522
1523 // Send through async worker for correct ordering with RunIngestion, and
1524 // dropping the dataflow is done on async worker response.
1525 if self.timely_worker_index == 0 {
1526 self.async_worker.drop_dataflow(id);
1527 }
1528 }
1529
1530 /// Drop the identified oneshot ingestion from the storage state.
1531 fn drop_oneshot_ingestion(&mut self, ingestion_id: uuid::Uuid) {
1532 let prev = self.oneshot_ingestions.remove(&ingestion_id);
1533 info!(%ingestion_id, existed = %prev.is_some(), "dropping oneshot ingestion");
1534 }
1535}