Skip to main content

Event

Enum Event 

Source
#[non_exhaustive]
pub enum Event { Input { id: DataId, metadata: Metadata, data: ArrowData, }, InputClosed { id: DataId, }, InputRecovered { id: DataId, }, NodeRestarted { id: NodeId, }, Stop(StopCause), Reload { operator_id: Option<OperatorId>, }, ParamUpdate { key: String, value: Value, }, ParamDeleted { key: String, }, NodeFailed { affected_input_ids: Vec<DataId>, error: String, source_node_id: NodeId, }, Error(String), }
Expand description

Represents an incoming Dora event.

Events might be triggered by other nodes, by Dora itself, or by some external user input.

It’s safe to ignore event types that are not relevant to the node.

This enum is marked as non_exhaustive because we might add additional variants in the future. Please ignore unknown event types instead of throwing an error to avoid breakage when updating Dora.

Variants (Non-exhaustive)§

This enum is marked as non-exhaustive
Non-exhaustive enums could have additional variants added in future. Therefore, when matching against variants of non-exhaustive enums, an extra wildcard arm must be added to account for any future variants.
§

Input

An input was received from another node.

This event corresponds to one of the inputs of the node as specified in the dataflow YAML file.

Fields

§id: DataId

The input ID, as specified in the YAML file.

Note that this is not the output ID of the sender, but the ID assigned to the input in the YAML file.

§metadata: Metadata

Meta information about this input, e.g. the timestamp.

§data: ArrowData

The actual data in the Apache Arrow data format.

§

InputClosed

An input was closed by the sender.

The sending node mapped to an input exited, so this input will receive no more data.

Fields

§id: DataId

The ID of the input that was closed, as specified in the YAML file.

Note that this is not the output ID of the sender, but the ID assigned to the input in the YAML file.

§

InputRecovered

A previously closed input has recovered and will receive data again.

This happens when an upstream node that timed out (via input_timeout) starts producing data again. The circuit breaker automatically re-opens the input.

Fields

§id: DataId

The ID of the recovered input, as specified in the YAML file.

§

NodeRestarted

An upstream node has restarted.

Sent to downstream nodes when a node with a restart policy successfully restarts after a failure. Nodes can use this to reset state, clear caches, or log the recovery.

Fields

§id: NodeId

The ID of the upstream node that restarted.

§

Stop(StopCause)

Notification that the event stream is about to close.

The StopCause field contains the reason for the event stream closure.

Nodes should exit once the event stream closes.

§

Reload

Instructs the node to reload itself or one of its operators.

This event is currently only used for reloading Python operators that are started by a dora runtime process. So this event should not be sent to normal nodes yet.

Fields

§operator_id: Option<OperatorId>

The ID of the operator that should be reloaded.

There is currently no case where operator_id is None.

§

ParamUpdate

A runtime parameter has been updated via dora param set.

Nodes can use this to dynamically adjust behavior (e.g., thresholds, rates) without restarting.

Fields

§key: String

The parameter key that was set.

§value: Value

The new JSON value.

§

ParamDeleted

A runtime parameter has been deleted via dora param delete.

Nodes can use this to remove local overrides and fall back to defaults.

Fields

§key: String

The parameter key that was deleted.

§

NodeFailed

An upstream node has failed.

Sent to downstream nodes when an upstream node exits with a non-zero exit code. Downstream nodes can use this to handle the failure gracefully (e.g. switch to cached data, log, retry).

Fields

§affected_input_ids: Vec<DataId>

The IDs of the inputs affected by the failure.

§error: String

Human-readable error message from the failed node.

§source_node_id: NodeId

The ID of the node that failed.

§

Error(String)

Notifies the node about an unexpected error that happened inside Dora.

It’s a good idea to output or log this error for debugging.

Trait Implementations§

Source§

impl Debug for Event

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more

Auto Trait Implementations§

§

impl !RefUnwindSafe for Event

§

impl !UnwindSafe for Event

§

impl Freeze for Event

§

impl Send for Event

§

impl Sync for Event

§

impl Unpin for Event

§

impl UnsafeUnpin for Event

Blanket Implementations§

Source§

impl<Source> AccessAs for Source

Source§

fn ref_as<T>(&self) -> <Source as IGuardRef<T>>::Guard<'_>
where Source: IGuardRef<T>, T: ?Sized,

Provides immutable access to a type as if it were its ABI-unstable equivalent.
Source§

fn mut_as<T>(&mut self) -> <Source as IGuardMut<T>>::GuardMut<'_>
where Source: IGuardMut<T>, T: ?Sized,

Provides mutable access to a type as if it were its ABI-unstable equivalent.
Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> AsNode<T> for T

Source§

fn as_node(&self) -> &T

Source§

impl<T> AsNodeMut<T> for T

Source§

fn as_node_mut(&mut self) -> &mut T

Source§

impl<'a, T, E> AsTaggedExplicit<'a, E> for T
where T: 'a,

Source§

fn explicit(self, class: Class, tag: u32) -> TaggedParser<'a, Explicit, Self, E>

Source§

impl<'a, T, E> AsTaggedImplicit<'a, E> for T
where T: 'a,

Source§

fn implicit( self, class: Class, constructed: bool, tag: u32, ) -> TaggedParser<'a, Implicit, Self, E>

Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> FutureExt for T

Source§

fn with_context(self, otel_cx: Context) -> WithContext<Self>

Attaches the provided Context to this type, returning a WithContext wrapper. Read more
Source§

fn with_current_context(self) -> WithContext<Self>

Attaches the current Context to this type, returning a WithContext wrapper. Read more
Source§

impl<T> FutureExt for T

Source§

fn with_context(self, otel_cx: Context) -> WithContext<Self>

Attaches the provided Context to this type, returning a WithContext wrapper. Read more
Source§

fn with_current_context(self) -> WithContext<Self>

Attaches the current Context to this type, returning a WithContext wrapper. Read more
Source§

impl<T, As> IGuardMut<As> for T
where T: Into<As>, As: Into<T>,

Source§

type GuardMut<'a> = MutAs<'a, T, As> where T: 'a

The type of the guard which will clean up the temporary after applying its changes to the original.
Source§

fn guard_mut_inner(&mut self) -> <T as IGuardMut<As>>::GuardMut<'_>

Construct the temporary and guard it through a mutable reference.
Source§

impl<T, As> IGuardRef<As> for T
where T: Into<As>, As: Into<T>,

Source§

type Guard<'a> = RefAs<'a, T, As> where T: 'a

The type of the guard which will clean up the temporary.
Source§

fn guard_ref_inner(&self) -> <T as IGuardRef<As>>::Guard<'_>

Construct the temporary and guard it through an immutable reference.
Source§

impl<T> Includes<End> for T

Source§

type Output = End

The result
Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts 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 more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts 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 more
Source§

impl<T> IntoRequest<T> for T

Source§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
Source§

impl<L> LayerExt<L> for L

Source§

fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>
where L: Layer<S>,

Applies the layer to a service and wraps it in Layered.
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more