pub struct CollectionManager<T>{
read_only: bool,
differential_collections: Arc<Mutex<BTreeMap<GlobalId, (UnboundedSender<(StorageWriteOp, Sender<Result<(), StorageError<T>>>)>, AbortOnDropHandle<()>, Sender<()>)>>>,
append_only_collections: Arc<Mutex<BTreeMap<GlobalId, (UnboundedSender<(Vec<AppendOnlyUpdate>, Sender<Result<(), StorageError<T>>>)>, AbortOnDropHandle<()>, Sender<()>)>>>,
user_batch_duration_ms: Arc<AtomicU64>,
now: NowFn,
}
Fields§
§read_only: bool
When a CollectionManager
is in read-only mode it must not affect any
changes to external state.
differential_collections: Arc<Mutex<BTreeMap<GlobalId, (UnboundedSender<(StorageWriteOp, Sender<Result<(), StorageError<T>>>)>, AbortOnDropHandle<()>, Sender<()>)>>>
These are collections that we write to by adding/removing updates to an
internal desired collection. The CollectionManager
continually makes
sure that collection contents (in persist) match the desired state.
append_only_collections: Arc<Mutex<BTreeMap<GlobalId, (UnboundedSender<(Vec<AppendOnlyUpdate>, Sender<Result<(), StorageError<T>>>)>, AbortOnDropHandle<()>, Sender<()>)>>>
Collections that we only append to using blind-writes.
Every write succeeds at some timestamp, and we never check what the actual contents of the collection (in persist) are.
user_batch_duration_ms: Arc<AtomicU64>
Amount of time we’ll wait before sending a batch of inserts to Persist, for user collections.
now: NowFn
Implementations§
source§impl<T> CollectionManager<T>
impl<T> CollectionManager<T>
The CollectionManager
provides two complementary functions:
- Providing an API to append values to a registered set of collections.
For this usecase:
- The
CollectionManager
expects to be the only writer. - Appending to a closed collection panics
- The
- Automatically advancing the timestamp of managed collections every
second. For this usecase:
- The
CollectionManager
handles contention by permitting and ignoring errors. - Closed collections will not panic if they continue receiving these requests.
- The
pub(crate) fn new(read_only: bool, now: NowFn) -> CollectionManager<T>
sourcepub fn update_user_batch_duration(&self, duration: Duration)
pub fn update_user_batch_duration(&self, duration: Duration)
Updates the duration we’ll wait to batch events for user owned collections.
sourcepub(crate) fn register_differential_collection<R>(
&self,
id: GlobalId,
write_handle: WriteHandle<SourceData, (), T, Diff>,
read_handle_fn: R,
force_writable: bool,
introspection_config: DifferentialIntrospectionConfig<T>,
)
pub(crate) fn register_differential_collection<R>( &self, id: GlobalId, write_handle: WriteHandle<SourceData, (), T, Diff>, read_handle_fn: R, force_writable: bool, introspection_config: DifferentialIntrospectionConfig<T>, )
Registers a new differential collection.
The CollectionManager will automatically advance the upper of every registered collection every second.
Update the desired
state of a differential collection using
Self::differential_write.
sourcepub(crate) fn register_append_only_collection(
&self,
id: GlobalId,
write_handle: WriteHandle<SourceData, (), T, Diff>,
force_writable: bool,
introspection_config: Option<AppendOnlyIntrospectionConfig<T>>,
)
pub(crate) fn register_append_only_collection( &self, id: GlobalId, write_handle: WriteHandle<SourceData, (), T, Diff>, force_writable: bool, introspection_config: Option<AppendOnlyIntrospectionConfig<T>>, )
Registers a new append-only collection.
The CollectionManager will automatically advance the upper of every registered collection every second.
sourcepub(crate) fn unregister_collection(
&self,
id: GlobalId,
) -> BoxFuture<'static, ()>
pub(crate) fn unregister_collection( &self, id: GlobalId, ) -> BoxFuture<'static, ()>
Unregisters the given collection.
Also waits until the CollectionManager
has completed all outstanding work to ensure that
it has stopped referencing the provided id
.
sourcepub(crate) fn blind_write(&self, id: GlobalId, updates: Vec<AppendOnlyUpdate>)
pub(crate) fn blind_write(&self, id: GlobalId, updates: Vec<AppendOnlyUpdate>)
Appends updates
to the append-only collection identified by id
, at
some timestamp. Does not wait for the append to complete.
§Panics
- If
id
does not belong to an append-only collections. - If this
CollectionManager
is in read-only mode. - If the collection closed.
sourcepub(crate) fn differential_write(&self, id: GlobalId, op: StorageWriteOp)
pub(crate) fn differential_write(&self, id: GlobalId, op: StorageWriteOp)
Updates the desired collection state of the differential collection identified by
id
. The underlying persist shard will reflect this change at
_some_point. Does not wait for the change to complete.
§Panics
- If
id
does not belong to a differential collection. - If the collection closed.
sourcepub(crate) fn differential_append(
&self,
id: GlobalId,
updates: Vec<(Row, Diff)>,
)
pub(crate) fn differential_append( &self, id: GlobalId, updates: Vec<(Row, Diff)>, )
Appends the given updates
to the differential collection identified by id
.
§Panics
- If
id
does not belong to a differential collection. - If the collection closed.
sourcepub(crate) fn monotonic_appender(
&self,
id: GlobalId,
) -> Result<MonotonicAppender<T>, StorageError<T>>
pub(crate) fn monotonic_appender( &self, id: GlobalId, ) -> Result<MonotonicAppender<T>, StorageError<T>>
Returns a MonotonicAppender
that can be used to monotonically append updates to the
collection correlated with id
.
fn get_read_only(&self, id: GlobalId, force_writable: bool) -> bool
Trait Implementations§
source§impl<T> Clone for CollectionManager<T>
impl<T> Clone for CollectionManager<T>
source§fn clone(&self) -> CollectionManager<T>
fn clone(&self) -> CollectionManager<T>
1.0.0 · source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source
. Read moreAuto Trait Implementations§
impl<T> Freeze for CollectionManager<T>
impl<T> !RefUnwindSafe for CollectionManager<T>
impl<T> Send for CollectionManager<T>
impl<T> Sync for CollectionManager<T>
impl<T> Unpin for CollectionManager<T>
impl<T> !UnwindSafe for CollectionManager<T>
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> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
source§default unsafe fn clone_to_uninit(&self, dst: *mut T)
default unsafe fn clone_to_uninit(&self, dst: *mut T)
clone_to_uninit
)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> 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::Request
source§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> 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> ProgressEventTimestamp for T
impl<T> ProgressEventTimestamp for T
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> 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.