Struct mz_txn_wal::operator::DataSubscribeTask
source · pub struct DataSubscribeTask {
tx: Sender<(Option<u64>, Sender<(Vec<(String, u64, i64)>, u64)>)>,
task: JoinHandle<Vec<(String, u64, i64)>>,
output: Vec<(String, u64, i64)>,
progress: u64,
}
Expand description
A handle to a DataSubscribe running in a task.
Fields§
§tx: Sender<(Option<u64>, Sender<(Vec<(String, u64, i64)>, u64)>)>
Carries step requests. A None
timestamp requests one step, a
Some(ts)
requests stepping until we progress beyond ts
.
task: JoinHandle<Vec<(String, u64, i64)>>
§output: Vec<(String, u64, i64)>
§progress: u64
Implementations§
source§impl DataSubscribeTask
impl DataSubscribeTask
sourcepub async fn new(
client: PersistClient,
txns_id: ShardId,
data_id: ShardId,
as_of: u64,
) -> Self
pub async fn new( client: PersistClient, txns_id: ShardId, data_id: ShardId, as_of: u64, ) -> Self
Creates a new DataSubscribeTask.
sourcepub async fn step_past(&mut self, ts: u64) -> u64
pub async fn step_past(&mut self, ts: u64) -> u64
Steps the dataflow past the given time, capturing output.
async fn send(&mut self, ts: Option<u64>)
sourcepub async fn finish(self) -> Vec<(String, u64, i64)>
pub async fn finish(self) -> Vec<(String, u64, i64)>
Signals for the task to exit, and then waits for this to happen.
All output from the lifetime of the task (not just what was previously captured) is returned.
fn task( client: PersistClient, cache: TxnsCache<u64>, data_id: ShardId, as_of: u64, rx: Receiver<(Option<u64>, Sender<(Vec<(String, u64, i64)>, u64)>)>, ) -> Vec<(String, u64, i64)>
Trait Implementations§
Auto Trait Implementations§
impl Freeze for DataSubscribeTask
impl RefUnwindSafe for DataSubscribeTask
impl Send for DataSubscribeTask
impl Sync for DataSubscribeTask
impl Unpin for DataSubscribeTask
impl UnwindSafe for DataSubscribeTask
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
Mutably borrows from an owned value. Read more
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>
Wrap the input message
T
in a tonic::Request
Creates a shared type from an unshared type.
source§impl<T> Pointable for T
impl<T> Pointable 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>
See
RustType::from_proto
.source§fn from_rust(rust: &R) -> P
fn from_rust(rust: &R) -> P
See
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)
The method of
std::ops::AddAssign
, for types that do not implement AddAssign
.