pub struct DoraNode { /* private fields */ }Expand description
Allows sending outputs and retrieving node information.
The main purpose of this struct is to send outputs via Dora. There are also functions available for retrieving the node configuration.
Implementations§
Source§impl DoraNode
impl DoraNode
Sourcepub fn init_from_env() -> NodeResult<(Self, EventStream)>
pub fn init_from_env() -> NodeResult<(Self, EventStream)>
Initiate a node from environment variables set by the Dora daemon or fall back to interactive mode.
This is the recommended initialization function for Dora nodes, which are spawned by
Dora daemon instances. The daemon will set a DORA_NODE_CONFIG environment variable to
configure the node.
When the node is started manually without the DORA_NODE_CONFIG environment variable set,
the initialization will fall back to init_interactive if stdin
is a terminal (detected through
isatty).
If the DORA_NODE_CONFIG environment variable is not set and DORA_TEST_WITH_INPUTS is
set, the node will be initialized in integration test mode. See the
integration testing module for details.
This function will also initialize the node in integration test mode when the
setup_integration_testing
function was called before. This takes precedence over all environment variables.
use dora_node_api::DoraNode;
let (mut node, mut events) = DoraNode::init_from_env().expect("Could not init node.");Sourcepub fn init_from_env_force() -> NodeResult<(Self, EventStream)>
pub fn init_from_env_force() -> NodeResult<(Self, EventStream)>
Initialize the node from environment variables set by the Dora daemon; error if not set.
This function behaves the same as init_from_env, but it does not
fall back to init_interactive. Instead, an error is returned
when the DORA_NODE_CONFIG environment variable is missing.
Sourcepub fn builder() -> DoraNodeBuilder
pub fn builder() -> DoraNodeBuilder
Create a builder for configuring a node connection.
Setting a node_id selects the dynamic-node path; without one, build()
falls back to init_from_env. Use this builder
when you need a custom daemon port — the other init functions cover the
common cases. Source-compatible with upstream dora 0.5.x: .dynamic()
is accepted (no-op) so code written against upstream still compiles.
use dora_node_api::DoraNode;
use dora_node_api::dora_core::config::NodeId;
let (mut node, mut events) = DoraNode::builder()
.node_id(NodeId::from("plot".to_string()))
.daemon_port(6789)
.build()
.expect("Could not init node");Sourcepub fn init_from_node_id(node_id: NodeId) -> NodeResult<(Self, EventStream)>
pub fn init_from_node_id(node_id: NodeId) -> NodeResult<(Self, EventStream)>
Initiate a node from a dataflow id and a node id.
This initialization function should be used for dynamic nodes.
use dora_node_api::DoraNode;
use dora_node_api::dora_core::config::NodeId;
let (mut node, mut events) = DoraNode::init_from_node_id(NodeId::from("plot".to_string())).expect("Could not init node plot");Sourcepub fn init_flexible(node_id: NodeId) -> NodeResult<(Self, EventStream)>
pub fn init_flexible(node_id: NodeId) -> NodeResult<(Self, EventStream)>
Dynamic initialization function for nodes that are sometimes used as dynamic nodes.
This function first tries initializing the traditional way through
init_from_env. If this fails, it falls back to
init_from_node_id.
Sourcepub fn init_interactive() -> NodeResult<(Self, EventStream)>
pub fn init_interactive() -> NodeResult<(Self, EventStream)>
Initialize the node in a standalone mode that prompts for inputs on the terminal.
Instead of connecting to a dora daemon, this interactive mode prompts for node inputs
on the terminal. In this mode, the node is completely isolated from the dora daemon and
other nodes, so it cannot be part of a dataflow.
Note that this function will hang indefinitely if no input is supplied to the interactive prompt. So it should be only used through a terminal.
Because of the above limitations, it is not recommended to use this function directly.
Use init_from_env instead, which supports both normal daemon
connections and manual interactive runs.
§Example
Run any node that uses init_interactive or init_from_env directly
from a terminal. The node will then start in “interactive mode” and prompt you for the next
input:
> cargo build -p rust-dataflow-example-node
> target/debug/rust-dataflow-example-node
hello
Starting node in interactive mode as DORA_NODE_CONFIG env variable is not set
Node asks for next input
? Input ID
[empty input ID to stop]The rust-dataflow-example-node expects a tick input, so let’s set the input ID to
tick. Tick messages don’t have any data, so we leave the “Data” empty when prompted:
Node asks for next input
> Input ID tick
> Data
tick 0, sending 0x943ed1be20c711a4
node sends output random with data: PrimitiveArray<UInt64>
[
10682205980693303716,
]
Node asks for next input
? Input ID
[empty input ID to stop]We see that both the stdout output of the node and also the output messages that it sends
are printed to the terminal. Then we get another prompt for the next input.
If you want to send an input with data, you can either send it as text (for string data) or as a JSON object (for struct, string, or array data). Other data types are not supported currently.
Empty input IDs are interpreted as stop instructions:
> Input ID
given input ID is empty -> stopping
Received stop
Node asks for next input
event channel was stopped -> returning empty event list
node reports EventStreamDropped
node reports closed outputs []
node reports OutputsDoneIn addition to the node output, we see log messages for the different events that the node
reports. After OutputsDone, the node should exit.
§JSON data
In addition to text input, the Data prompt also supports JSON objects, which will be
converted to Apache Arrow struct arrays:
Node asks for next input
> Input ID some_input
> Data { "field_1": 42, "field_2": { "inner": "foo" } }This JSON data is converted to the following Arrow array:
StructArray
-- validity: [valid, ]
[
-- child 0: "field_1" (Int64)
PrimitiveArray<Int64>
[42,]
-- child 1: "field_2" (Struct([Field { name: "inner", data_type: Utf8, nullable: true, dict_id: 0, dict_is_ordered: false, metadata: {} }]))
StructArray
-- validity: [valid,]
[
-- child 0: "inner" (Utf8)
StringArray
["foo",]
]
]Sourcepub fn init_testing(
input: TestingInput,
output: TestingOutput,
options: TestingOptions,
) -> NodeResult<(Self, EventStream)>
pub fn init_testing( input: TestingInput, output: TestingOutput, options: TestingOptions, ) -> NodeResult<(Self, EventStream)>
Initializes a node in integration test mode.
No connection to a dora daemon is made in this mode. Instead, inputs are read from the
specified TestingInput, and outputs are written to the specified TestingOutput.
Additional options for the testing mode can be specified through TestingOptions.
It is recommended to use this function only within test functions.
Sourcepub fn validate_output(&mut self, output_id: &DataId) -> bool
pub fn validate_output(&mut self, output_id: &DataId) -> bool
Check whether output_id is declared as an output of this node.
Returns true if the output is declared (or this node is interactive,
which has no static output declaration); false and emits a one-time
warning if the output is unknown. Public so callers building higher-level
send helpers (e.g. the Python send_output_raw zero-copy path) can
validate before allocating a buffer.
Sourcepub fn send_output_raw<F>(
&mut self,
output_id: DataId,
parameters: MetadataParameters,
data_len: usize,
data: F,
) -> NodeResult<()>
pub fn send_output_raw<F>( &mut self, output_id: DataId, parameters: MetadataParameters, data_len: usize, data: F, ) -> NodeResult<()>
Send raw data from the node to the other nodes.
We take a closure as an input to enable zero copy on send.
use dora_node_api::{DoraNode, MetadataParameters};
use dora_core::config::DataId;
let (mut node, mut events) = DoraNode::init_from_env().expect("Could not init node.");
let output = DataId::from("output_id".to_owned());
let data: &[u8] = &[0, 1, 2, 3];
let parameters = MetadataParameters::default();
node.send_output_raw(
output,
parameters,
data.len(),
|out| {
out.copy_from_slice(data);
}).expect("Could not send output");Ignores the output if the given output_id is not specified as node output in the dataflow
configuration file.
Sourcepub fn send_output(
&mut self,
output_id: DataId,
parameters: MetadataParameters,
data: impl IntoArrow,
) -> NodeResult<()>
pub fn send_output( &mut self, output_id: DataId, parameters: MetadataParameters, data: impl IntoArrow, ) -> NodeResult<()>
Sends the given Arrow array as an output message.
This is the recommended way to emit data from a node: pass any value that
implements IntoArrow (primitives, Vec<T>, &str, an Arrow array,
…) and dora moves it into shared memory for an efficient, near-zero-copy
transfer to downstream nodes.
Uses shared memory for efficient data transfer if suitable. This method might copy the message once to move it to shared memory.
Ignores the output if the given output_id is not specified as node output in the dataflow
configuration file.
use dora_node_api::{DoraNode, MetadataParameters};
use dora_core::config::DataId;
let (mut node, _events) = DoraNode::init_from_env()?;
let output = DataId::from("output_id".to_owned());
let parameters = MetadataParameters::default();
node.send_output(output, parameters, vec![1.0f32, 2.0, 3.0])?;Sourcepub fn send_output_encoded(
&mut self,
output_id: DataId,
parameters: MetadataParameters,
encoded: EncodedSample,
) -> NodeResult<()>
pub fn send_output_encoded( &mut self, output_id: DataId, parameters: MetadataParameters, encoded: EncodedSample, ) -> NodeResult<()>
Send a payload that has already been IPC-encoded into a dora-owned
EncodedSample by SampleAllocator::encode_arrow.
This is the entry point the runtime uses for operator outputs: the
operator thread does the encoding so that no memory owned by the
operator’s language runtime is ever released on the node’s thread — see
SampleAllocator (dora-rs/dora#2742).
Sourcepub fn send_output_bytes(
&mut self,
output_id: DataId,
parameters: MetadataParameters,
data_len: usize,
data: &[u8],
) -> NodeResult<()>
pub fn send_output_bytes( &mut self, output_id: DataId, parameters: MetadataParameters, data_len: usize, data: &[u8], ) -> NodeResult<()>
Send the given raw byte data as output.
Might copy the data once to move it into shared memory.
Ignores the output if the given output_id is not specified as node output in the dataflow
configuration file.
Sourcepub fn send_output_sample(
&mut self,
output_id: DataId,
parameters: MetadataParameters,
sample: Option<DataSample>,
) -> NodeResult<()>
pub fn send_output_sample( &mut self, output_id: DataId, parameters: MetadataParameters, sample: Option<DataSample>, ) -> NodeResult<()>
Sends the given DataSample as output.
The sample must already be a self-describing Arrow IPC stream (the
FRAMING_ARROW_IPC parameter should be set). It is recommended to use a
function like send_output instead, which handles
the encoding.
Ignores the output if the given output_id is not specified as node output in the dataflow
configuration file.
Sourcepub fn close_outputs(&mut self, outputs_ids: Vec<DataId>) -> NodeResult<()>
pub fn close_outputs(&mut self, outputs_ids: Vec<DataId>) -> NodeResult<()>
Report the given outputs IDs as closed.
The node is not allowed to send more outputs with the closed IDs.
Closing outputs early can be helpful to receivers.
Sourcepub fn id(&self) -> &NodeId
pub fn id(&self) -> &NodeId
Returns the ID of the node as specified in the dataflow configuration file.
Sourcepub fn dataflow_id(&self) -> &DataflowId
pub fn dataflow_id(&self) -> &DataflowId
Returns the unique identifier for the running dataflow instance.
Dora assigns each dataflow instance a random identifier when started.
Sourcepub fn node_config(&self) -> &NodeRunConfig
pub fn node_config(&self) -> &NodeRunConfig
Returns the input and output configuration of this node.
Sourcepub fn zero_copy_threshold(&self) -> usize
pub fn zero_copy_threshold(&self) -> usize
Returns the zero-copy SHM threshold in bytes.
Outputs whose raw payload is at least this many bytes are published via
zenoh shared memory (zero-copy for local subscribers); smaller outputs
are published via zenoh with a heap-buffered put. Configured via the
DORA_ZERO_COPY_THRESHOLD env var, defaulting to
ZERO_COPY_THRESHOLD.
Sourcepub fn is_restart(&self) -> bool
pub fn is_restart(&self) -> bool
Returns true if this node was restarted after a previous exit or failure.
Nodes can use this to decide whether to restore saved state or start fresh.
Sourcepub fn restart_count(&self) -> u32
pub fn restart_count(&self) -> u32
Returns how many times this node has been restarted.
Returns 0 on the first run, 1 after the first restart, etc.
Sourcepub fn timestamp(&self) -> Timestamp
pub fn timestamp(&self) -> Timestamp
Returns the current timestamp from the node’s Hybrid Logical Clock.
This generates a new HLC timestamp, which combines the physical
wall-clock time with a logical counter to ensure uniqueness and
monotonicity even across nodes. The HLC is the same clock dora
stamps every outgoing message with, so this is the right value
to subtract from an input event’s metadata.timestamp when
measuring per-event processing latency — using
std::time::SystemTime::now() instead would mix two unrelated
clocks and give meaningless results across daemons.
Sourcepub fn log(&self, level: &str, message: &str, target: Option<&str>)
pub fn log(&self, level: &str, message: &str, target: Option<&str>)
Send a structured log message.
Outputs a JSONL line to stdout that the daemon parses automatically.
Works with min_log_level filtering and send_logs_as routing.
level should be one of: "error", "warn", "info", "debug", "trace".
Unknown levels default to "info".
Sourcepub fn log_with_fields(
&self,
level: &str,
message: &str,
target: Option<&str>,
fields: Option<&BTreeMap<String, String>>,
)
pub fn log_with_fields( &self, level: &str, message: &str, target: Option<&str>, fields: Option<&BTreeMap<String, String>>, )
Send a structured log message with optional key-value fields.
Like log, but accepts additional structured fields that
are included in the JSON payload and preserved through send_logs_as.
Sourcepub fn new_request_id() -> String
pub fn new_request_id() -> String
Generate a new unique request/goal ID (UUID v7, time-ordered).
Uses a per-thread monotonic counter context to guarantee uniqueness even when multiple IDs are generated within the same clock tick.
Sourcepub fn new_goal_id() -> String
pub fn new_goal_id() -> String
Generate a new unique goal ID (UUID v7, time-ordered).
This is an alias for new_request_id that
reads more naturally in action (goal/feedback/result) contexts.
Sourcepub fn send_service_request(
&mut self,
output_id: DataId,
parameters: MetadataParameters,
data: impl IntoArrow,
) -> NodeResult<String>
pub fn send_service_request( &mut self, output_id: DataId, parameters: MetadataParameters, data: impl IntoArrow, ) -> NodeResult<String>
Send a service request, automatically injecting a request_id into the
metadata parameters. Returns the generated request ID.
Any existing request_id key in parameters is replaced.
Sourcepub fn send_service_response(
&mut self,
output_id: DataId,
parameters: MetadataParameters,
data: impl IntoArrow,
) -> NodeResult<()>
pub fn send_service_response( &mut self, output_id: DataId, parameters: MetadataParameters, data: impl IntoArrow, ) -> NodeResult<()>
Send a service response. This is a semantic alias for send_output.
The caller is expected to pass through the request_id parameter from
the incoming request’s metadata.
Sourcepub fn send_stream_chunk(
&mut self,
output_id: DataId,
segment: &mut StreamSegment,
fin: bool,
data: impl IntoArrow,
) -> NodeResult<()>
pub fn send_stream_chunk( &mut self, output_id: DataId, segment: &mut StreamSegment, fin: bool, data: impl IntoArrow, ) -> NodeResult<()>
Send a streaming segment chunk. Convenience wrapper around
send_output that builds metadata from the
StreamSegment builder.
Sourcepub fn allocate_data_sample(
&mut self,
data_len: usize,
) -> NodeResult<DataSample>
pub fn allocate_data_sample( &mut self, data_len: usize, ) -> NodeResult<DataSample>
Allocates a DataSample of the specified size.
See SampleAllocator::allocate for the allocation strategy.
Sourcepub fn sample_allocator(&self) -> SampleAllocator
pub fn sample_allocator(&self) -> SampleAllocator
A handle for building output samples off this node’s thread.
The runtime hands one to each operator thread so the operator can encode
its payload into a dora-owned DataSample itself — see
SampleAllocator for why that matters (dora-rs/dora#2742).
Sourcepub fn dataflow_descriptor(&self) -> NodeResult<&Descriptor>
pub fn dataflow_descriptor(&self) -> NodeResult<&Descriptor>
Returns the full dataflow descriptor that this node is part of.
This method returns the parsed dataflow YAML file.
Sourcepub fn extension_store(
&mut self,
namespace: impl Into<String>,
key: impl Into<String>,
value: Vec<u8>,
) -> Result<(), Error>
pub fn extension_store( &mut self, namespace: impl Into<String>, key: impl Into<String>, value: Vec<u8>, ) -> Result<(), Error>
Store an opaque value in the daemon’s dataflow-scoped extension table.
This is the seam for transports that live outside the dora tree: dora
brokers the value’s lifetime and nothing else — it never interprets
namespace, key or value. See docs/extensions.md.
The daemon remembers which nodes touched a key so that dropping it
notifies them, and reclaims the entry when the dataflow ends or the
storing node exits. Drain the notifications with
event_stream::extensions::drain_dropped_keys.
Sourcepub fn extension_load(
&mut self,
namespace: impl Into<String>,
key: impl Into<String>,
remove: bool,
) -> Result<Option<Vec<u8>>, Error>
pub fn extension_load( &mut self, namespace: impl Into<String>, key: impl Into<String>, remove: bool, ) -> Result<Option<Vec<u8>>, Error>
Read an opaque value back, optionally removing it in the same round trip.
Returns None if the key is not in the table — never stored, or
already dropped.
Sourcepub fn extension_drop(
&mut self,
namespace: impl Into<String>,
key: impl Into<String>,
) -> Result<(), Error>
pub fn extension_drop( &mut self, namespace: impl Into<String>, key: impl Into<String>, ) -> Result<(), Error>
Drop an opaque value, notifying every node that stored or loaded it.
Sourcepub fn extension_request(
&mut self,
namespace: impl Into<String>,
payload: Vec<u8>,
) -> Result<Vec<u8>, Error>
pub fn extension_request( &mut self, namespace: impl Into<String>, payload: Vec<u8>, ) -> Result<Vec<u8>, Error>
Send an opaque request to the extension registered under
namespace on this node’s daemon, and return its opaque reply.
Companion to DoraNode::extension_store / DoraNode::extension_load:
those broker a descriptor’s lifetime, this one carries a call the
extension’s daemon half must service. dora interprets neither the
namespace nor the bytes — see docs/extensions.md.
Trait Implementations§
Auto Trait Implementations§
impl !Freeze for DoraNode
impl !RefUnwindSafe for DoraNode
impl !UnwindSafe for DoraNode
impl Send for DoraNode
impl Sync for DoraNode
impl Unpin for DoraNode
impl UnsafeUnpin for DoraNode
Blanket Implementations§
Source§impl<Source> AccessAs for Source
impl<Source> AccessAs for Source
Source§impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
Source§impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
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
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
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::Request