pub struct StorageState {Show 29 fields
    pub source_uppers: BTreeMap<GlobalId, Rc<RefCell<Antichain<Timestamp>>>>,
    pub source_tokens: BTreeMap<GlobalId, Vec<PressOnDropButton>>,
    pub metrics: StorageMetrics,
    pub reported_frontiers: BTreeMap<GlobalId, Antichain<Timestamp>>,
    pub ingestions: BTreeMap<GlobalId, IngestionDescription<CollectionMetadata>>,
    pub exports: BTreeMap<GlobalId, StorageSinkDesc<CollectionMetadata, Timestamp>>,
    pub oneshot_ingestions: BTreeMap<Uuid, OneshotIngestionDescription<ProtoBatch>>,
    pub now: NowFn,
    pub timely_worker_index: usize,
    pub timely_worker_peers: usize,
    pub instance_context: StorageInstanceContext,
    pub persist_clients: Arc<PersistClientCache>,
    pub txns_ctx: TxnsContext,
    pub sink_tokens: BTreeMap<GlobalId, Vec<PressOnDropButton>>,
    pub sink_write_frontiers: BTreeMap<GlobalId, Rc<RefCell<Antichain<Timestamp>>>>,
    pub dropped_ids: Vec<GlobalId>,
    pub aggregated_statistics: AggregatedStatistics,
    pub shared_status_updates: Rc<RefCell<Vec<StatusUpdate>>>,
    pub latest_status_updates: BTreeMap<GlobalId, StatusUpdate>,
    pub initial_status_reported: bool,
    pub internal_cmd_tx: InternalCommandSender,
    pub internal_cmd_rx: InternalCommandReceiver,
    pub read_only_rx: Receiver<bool>,
    pub read_only_tx: Sender<bool>,
    pub async_worker: AsyncStorageWorker<Timestamp>,
    pub storage_configuration: StorageConfiguration,
    pub dataflow_parameters: DataflowParameters,
    pub tracing_handle: Arc<TracingHandle>,
    pub server_maintenance_interval: Duration,
}Expand description
Worker-local state related to the ingress or egress of collections of data.
Fields§
§source_uppers: BTreeMap<GlobalId, Rc<RefCell<Antichain<Timestamp>>>>The highest observed upper frontier for collection.
This is shared among all source instances, so that they can jointly advance the frontier even as other instances are created and dropped. Ideally, the Storage module would eventually provide one source of truth on this rather than multiple, and we should aim for that but are not there yet.
source_tokens: BTreeMap<GlobalId, Vec<PressOnDropButton>>Handles to created sources, keyed by ID
NB: The type of the tokens must not be changed to something other than PressOnDropButton
to prevent usage of custom shutdown tokens that are tricky to get right.
metrics: StorageMetricsMetrics for storage objects.
reported_frontiers: BTreeMap<GlobalId, Antichain<Timestamp>>Tracks the conditional write frontiers we have reported.
ingestions: BTreeMap<GlobalId, IngestionDescription<CollectionMetadata>>Descriptions of each installed ingestion.
exports: BTreeMap<GlobalId, StorageSinkDesc<CollectionMetadata, Timestamp>>Descriptions of each installed export.
oneshot_ingestions: BTreeMap<Uuid, OneshotIngestionDescription<ProtoBatch>>Descriptions of oneshot ingestions that are currently running.
now: NowFnUndocumented
timely_worker_index: usizeIndex of the associated timely dataflow worker.
timely_worker_peers: usizePeers in the associated timely dataflow worker.
instance_context: StorageInstanceContextOther configuration for sources and sinks.
persist_clients: Arc<PersistClientCache>A process-global cache of (blob_uri, consensus_uri) -> PersistClient. This is intentionally shared between workers
txns_ctx: TxnsContextContext necessary for rendering txn-wal operators.
sink_tokens: BTreeMap<GlobalId, Vec<PressOnDropButton>>Tokens that should be dropped when a dataflow is dropped to clean up
associated state.
NB: The type of the tokens must not be changed to something other than PressOnDropButton
to prevent usage of custom shutdown tokens that are tricky to get right.
sink_write_frontiers: BTreeMap<GlobalId, Rc<RefCell<Antichain<Timestamp>>>>Frontier of sink writes (all subsequent writes will be at times at or equal to this frontier)
dropped_ids: Vec<GlobalId>Collection ids that have been dropped but not yet reported as dropped
aggregated_statistics: AggregatedStatisticsStatistics for sources and sinks.
A place shared with running dataflows, so that health operators, can report status updates back to us.
NOTE: Operators that append to this collection should take care to only add new status updates if the status of the ingestion/export in question has changed.
latest_status_updates: BTreeMap<GlobalId, StatusUpdate>The latest status update for each object.
initial_status_reported: boolWhether we have reported the initial status after connecting to a new client. This is reset to false when a new client connects.
internal_cmd_tx: InternalCommandSenderSender for cluster-internal storage commands. These can be sent from within workers/operators and will be distributed to all workers. For example, for shutting down an entire dataflow from within a operator/worker.
internal_cmd_rx: InternalCommandReceiverReceiver for cluster-internal storage commands.
read_only_rx: Receiver<bool>When this replica/cluster is in read-only mode it must not affect any changes to external state. This flag can only be changed by a StorageCommand::AllowWrites.
Everything running on this replica/cluster must obey this flag. At the time of writing, nothing currently looks at this flag. TODO(benesch): fix this.
NOTE: In the future, we might want a more complicated flag, for example something that tells us after which timestamp we are allowed to write. In this first version we are keeping things as simple as possible!
read_only_tx: Sender<bool>Send-side for read-only state.
async_worker: AsyncStorageWorker<Timestamp>Async worker companion, used for running code that requires async, which the timely main loop cannot do.
storage_configuration: StorageConfigurationConfiguration for source and sink connections.
dataflow_parameters: DataflowParametersDynamically configurable parameters that control how dataflows are rendered.
NOTE(guswynn): we should consider moving these into storage_configuration.
tracing_handle: Arc<TracingHandle>A process-global handle to tracing configuration.
server_maintenance_interval: DurationInterval at which to perform server maintenance tasks. Set to a zero interval to
perform maintenance with every step_or_park invocation.
Implementations§
Source§impl StorageState
 
impl StorageState
Sourcepub fn error_handler(&self, context: &'static str, id: GlobalId) -> ErrorHandler
 
pub fn error_handler(&self, context: &'static str, id: GlobalId) -> ErrorHandler
Return an error handler that triggers a suspend and restart of the corresponding storage dataflow.
Source§impl StorageState
 
impl StorageState
Sourcepub fn handle_storage_command(&mut self, cmd: StorageCommand)
 
pub fn handle_storage_command(&mut self, cmd: StorageCommand)
Entry point for applying a storage command.
NOTE: This does not have access to the timely worker and therefore
cannot render dataflows. For dataflow rendering, this needs to either
send asynchronous command to the async_worker or internal
commands to the internal_cmd_tx.
Sourcefn drop_collection(&mut self, id: GlobalId)
 
fn drop_collection(&mut self, id: GlobalId)
Drop the identified storage collection from the storage state.
Sourcefn drop_oneshot_ingestion(&mut self, ingestion_id: Uuid)
 
fn drop_oneshot_ingestion(&mut self, ingestion_id: Uuid)
Drop the identified oneshot ingestion from the storage state.
Auto Trait Implementations§
impl !Freeze for StorageState
impl !RefUnwindSafe for StorageState
impl !Send for StorageState
impl !Sync for StorageState
impl Unpin for StorageState
impl !UnwindSafe for StorageState
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
    T: ?Sized,
 
impl<T> BorrowMut<T> for Twhere
    T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
 
fn borrow_mut(&mut self) -> &mut T
Source§impl<T> Downcast for T
 
impl<T> Downcast for T
Source§impl<T> FmtForward for T
 
impl<T> FmtForward for T
Source§fn fmt_binary(self) -> FmtBinary<Self>where
    Self: Binary,
 
fn fmt_binary(self) -> FmtBinary<Self>where
    Self: Binary,
self to use its Binary implementation when Debug-formatted.Source§fn fmt_display(self) -> FmtDisplay<Self>where
    Self: Display,
 
fn fmt_display(self) -> FmtDisplay<Self>where
    Self: Display,
self to use its Display implementation when
Debug-formatted.Source§fn fmt_lower_exp(self) -> FmtLowerExp<Self>where
    Self: LowerExp,
 
fn fmt_lower_exp(self) -> FmtLowerExp<Self>where
    Self: LowerExp,
self to use its LowerExp implementation when
Debug-formatted.Source§fn fmt_lower_hex(self) -> FmtLowerHex<Self>where
    Self: LowerHex,
 
fn fmt_lower_hex(self) -> FmtLowerHex<Self>where
    Self: LowerHex,
self to use its LowerHex implementation when
Debug-formatted.Source§fn fmt_octal(self) -> FmtOctal<Self>where
    Self: Octal,
 
fn fmt_octal(self) -> FmtOctal<Self>where
    Self: Octal,
self to use its Octal implementation when Debug-formatted.Source§fn fmt_pointer(self) -> FmtPointer<Self>where
    Self: Pointer,
 
fn fmt_pointer(self) -> FmtPointer<Self>where
    Self: Pointer,
self to use its Pointer implementation when
Debug-formatted.Source§fn fmt_upper_exp(self) -> FmtUpperExp<Self>where
    Self: UpperExp,
 
fn fmt_upper_exp(self) -> FmtUpperExp<Self>where
    Self: UpperExp,
self to use its UpperExp implementation when
Debug-formatted.Source§fn fmt_upper_hex(self) -> FmtUpperHex<Self>where
    Self: UpperHex,
 
fn fmt_upper_hex(self) -> FmtUpperHex<Self>where
    Self: UpperHex,
self to use its UpperHex implementation when
Debug-formatted.Source§impl<T> FutureExt for T
 
impl<T> FutureExt for T
Source§fn with_context(self, otel_cx: Context) -> WithContext<Self>
 
fn with_context(self, otel_cx: Context) -> WithContext<Self>
Source§fn with_current_context(self) -> WithContext<Self>
 
fn with_current_context(self) -> WithContext<Self>
Source§impl<T> Instrument for T
 
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
 
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
 
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoEither for T
 
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self>
 
fn into_either(self, into_left: bool) -> Either<Self, Self>
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
 
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§impl<T> IntoRequest<T> for T
 
impl<T> IntoRequest<T> for T
Source§fn into_request(self) -> Request<T>
 
fn into_request(self) -> Request<T>
T in a tonic::RequestSource§impl<T, U> OverrideFrom<Option<&T>> for Uwhere
    U: OverrideFrom<T>,
 
impl<T, U> OverrideFrom<Option<&T>> for Uwhere
    U: OverrideFrom<T>,
Source§impl<T> Paint for Twhere
    T: ?Sized,
 
impl<T> Paint for Twhere
    T: ?Sized,
Source§fn fg(&self, value: Color) -> Painted<&T>
 
fn fg(&self, value: Color) -> Painted<&T>
Returns a styled value derived from self with the foreground set to
value.
This method should be used rarely. Instead, prefer to use color-specific
builder methods like red() and
green(), which have the same functionality but are
pithier.
§Example
Set foreground color to white using fg():
use yansi::{Paint, Color};
painted.fg(Color::White);Set foreground color to white using white().
use yansi::Paint;
painted.white();Source§fn bright_black(&self) -> Painted<&T>
 
fn bright_black(&self) -> Painted<&T>
Source§fn bright_red(&self) -> Painted<&T>
 
fn bright_red(&self) -> Painted<&T>
Source§fn bright_green(&self) -> Painted<&T>
 
fn bright_green(&self) -> Painted<&T>
Source§fn bright_yellow(&self) -> Painted<&T>
 
fn bright_yellow(&self) -> Painted<&T>
Source§fn bright_blue(&self) -> Painted<&T>
 
fn bright_blue(&self) -> Painted<&T>
Source§fn bright_magenta(&self) -> Painted<&T>
 
fn bright_magenta(&self) -> Painted<&T>
Source§fn bright_cyan(&self) -> Painted<&T>
 
fn bright_cyan(&self) -> Painted<&T>
Source§fn bright_white(&self) -> Painted<&T>
 
fn bright_white(&self) -> Painted<&T>
Source§fn bg(&self, value: Color) -> Painted<&T>
 
fn bg(&self, value: Color) -> Painted<&T>
Returns a styled value derived from self with the background set to
value.
This method should be used rarely. Instead, prefer to use color-specific
builder methods like on_red() and
on_green(), which have the same functionality but
are pithier.
§Example
Set background color to red using fg():
use yansi::{Paint, Color};
painted.bg(Color::Red);Set background color to red using on_red().
use yansi::Paint;
painted.on_red();Source§fn on_primary(&self) -> Painted<&T>
 
fn on_primary(&self) -> Painted<&T>
Source§fn on_magenta(&self) -> Painted<&T>
 
fn on_magenta(&self) -> Painted<&T>
Source§fn on_bright_black(&self) -> Painted<&T>
 
fn on_bright_black(&self) -> Painted<&T>
Source§fn on_bright_red(&self) -> Painted<&T>
 
fn on_bright_red(&self) -> Painted<&T>
Source§fn on_bright_green(&self) -> Painted<&T>
 
fn on_bright_green(&self) -> Painted<&T>
Source§fn on_bright_yellow(&self) -> Painted<&T>
 
fn on_bright_yellow(&self) -> Painted<&T>
Source§fn on_bright_blue(&self) -> Painted<&T>
 
fn on_bright_blue(&self) -> Painted<&T>
Source§fn on_bright_magenta(&self) -> Painted<&T>
 
fn on_bright_magenta(&self) -> Painted<&T>
Source§fn on_bright_cyan(&self) -> Painted<&T>
 
fn on_bright_cyan(&self) -> Painted<&T>
Source§fn on_bright_white(&self) -> Painted<&T>
 
fn on_bright_white(&self) -> Painted<&T>
Source§fn attr(&self, value: Attribute) -> Painted<&T>
 
fn attr(&self, value: Attribute) -> Painted<&T>
Enables the styling Attribute value.
This method should be used rarely. Instead, prefer to use
attribute-specific builder methods like bold() and
underline(), which have the same functionality
but are pithier.
§Example
Make text bold using attr():
use yansi::{Paint, Attribute};
painted.attr(Attribute::Bold);Make text bold using using bold().
use yansi::Paint;
painted.bold();Source§fn rapid_blink(&self) -> Painted<&T>
 
fn rapid_blink(&self) -> Painted<&T>
Source§fn quirk(&self, value: Quirk) -> Painted<&T>
 
fn quirk(&self, value: Quirk) -> Painted<&T>
Enables the yansi Quirk value.
This method should be used rarely. Instead, prefer to use quirk-specific
builder methods like mask() and
wrap(), which have the same functionality but are
pithier.
§Example
Enable wrapping using .quirk():
use yansi::{Paint, Quirk};
painted.quirk(Quirk::Wrap);Enable wrapping using wrap().
use yansi::Paint;
painted.wrap();Source§fn clear(&self) -> Painted<&T>
 👎Deprecated since 1.0.1: renamed to resetting() due to conflicts with Vec::clear().
The clear() method will be removed in a future release.
fn clear(&self) -> Painted<&T>
resetting() due to conflicts with Vec::clear().
The clear() method will be removed in a future release.Source§fn whenever(&self, value: Condition) -> Painted<&T>
 
fn whenever(&self, value: Condition) -> Painted<&T>
Conditionally enable styling based on whether the Condition value
applies. Replaces any previous condition.
See the crate level docs for more details.
§Example
Enable styling painted only when both stdout and stderr are TTYs:
use yansi::{Paint, Condition};
painted.red().on_yellow().whenever(Condition::STDOUTERR_ARE_TTY);Source§impl<T> Pipe for Twhere
    T: ?Sized,
 
impl<T> Pipe for Twhere
    T: ?Sized,
Source§fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> Rwhere
    Self: Sized,
 
fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> Rwhere
    Self: Sized,
Source§fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> Rwhere
    R: 'a,
 
fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> Rwhere
    R: 'a,
self and passes that borrow into the pipe function. Read moreSource§fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> Rwhere
    R: 'a,
 
fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> Rwhere
    R: 'a,
self and passes that borrow into the pipe function. Read moreSource§fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
 
fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
Source§fn pipe_borrow_mut<'a, B, R>(
    &'a mut self,
    func: impl FnOnce(&'a mut B) -> R,
) -> R
 
fn pipe_borrow_mut<'a, B, R>( &'a mut self, func: impl FnOnce(&'a mut B) -> R, ) -> R
Source§fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
 
fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
self, then passes self.as_ref() into the pipe function.Source§fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
 
fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
self, then passes self.as_mut() into the pipe
function.Source§fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
 
fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
self, then passes self.deref() into the pipe function.Source§impl<T> Pointable for T
 
impl<T> Pointable for T
Source§impl<T> PolicyExt for Twhere
    T: ?Sized,
 
impl<T> PolicyExt for Twhere
    T: ?Sized,
Source§impl<P, R> ProtoType<R> for Pwhere
    R: RustType<P>,
 
impl<P, R> ProtoType<R> for Pwhere
    R: RustType<P>,
Source§fn into_rust(self) -> Result<R, TryFromProtoError>
 
fn into_rust(self) -> Result<R, TryFromProtoError>
RustType::from_proto.Source§fn from_rust(rust: &R) -> P
 
fn from_rust(rust: &R) -> P
RustType::into_proto.Source§impl<'a, S, T> Semigroup<&'a S> for Twhere
    T: Semigroup<S>,
 
impl<'a, S, T> Semigroup<&'a S> for Twhere
    T: Semigroup<S>,
Source§fn plus_equals(&mut self, rhs: &&'a S)
 
fn plus_equals(&mut self, rhs: &&'a S)
std::ops::AddAssign, for types that do not implement AddAssign.Source§impl<T> ServiceExt for T
 
impl<T> ServiceExt for T
Source§fn map_response_body<F>(self, f: F) -> MapResponseBody<Self, F>where
    Self: Sized,
 
fn map_response_body<F>(self, f: F) -> MapResponseBody<Self, F>where
    Self: Sized,
Source§fn decompression(self) -> Decompression<Self>where
    Self: Sized,
 
fn decompression(self) -> Decompression<Self>where
    Self: Sized,
Source§fn trace_for_http(self) -> Trace<Self, SharedClassifier<ServerErrorsAsFailures>>where
    Self: Sized,
 
fn trace_for_http(self) -> Trace<Self, SharedClassifier<ServerErrorsAsFailures>>where
    Self: Sized,
Source§fn trace_for_grpc(self) -> Trace<Self, SharedClassifier<GrpcErrorsAsFailures>>where
    Self: Sized,
 
fn trace_for_grpc(self) -> Trace<Self, SharedClassifier<GrpcErrorsAsFailures>>where
    Self: Sized,
Source§fn follow_redirects(self) -> FollowRedirect<Self>where
    Self: Sized,
 
fn follow_redirects(self) -> FollowRedirect<Self>where
    Self: Sized,
Source§impl<T> Tap for T
 
impl<T> Tap for T
Source§fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
 
fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
Borrow<B> of a value. Read moreSource§fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
 
fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
BorrowMut<B> of a value. Read moreSource§fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
 
fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
AsRef<R> view of a value. Read moreSource§fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
 
fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
AsMut<R> view of a value. Read moreSource§fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
 
fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
Deref::Target of a value. Read moreSource§fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
 
fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
Deref::Target of a value. Read moreSource§fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self
 
fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self
.tap() only in debug builds, and is erased in release builds.Source§fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self
 
fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self
.tap_mut() only in debug builds, and is erased in release
builds.Source§fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
 
fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
.tap_borrow() only in debug builds, and is erased in release
builds.Source§fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
 
fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
.tap_borrow_mut() only in debug builds, and is erased in release
builds.Source§fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
 
fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
.tap_ref() only in debug builds, and is erased in release
builds.Source§fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
 
fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
.tap_ref_mut() only in debug builds, and is erased in release
builds.Source§fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
 
fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
.tap_deref() only in debug builds, and is erased in release
builds.