Skip to main content

NodeCore

Struct NodeCore 

Source
pub struct NodeCore {
Show 24 fields runtime: Runtime, config: NodeConfig, sentinel: SentinelClient, worker_pool: WorkerPool, state: NodeState, connection_status: ConnectionStatus, registered: bool, node_name: String, total_cores: u32, app_name: String, app_version: String, work_context_cache: WorkContextCache, verify_rx: Option<Receiver<Result<Sensitive<String>, VerifyError>>>, register_rx: Option<Receiver<RegisterResult>>, realtime_rx: Option<Receiver<RealtimeEvent>>, realtime_shutdown: Option<CancellationToken>, result_rx: Option<Receiver<WorkBatchResult>>, unlink_rx: Option<Receiver<Result<UnlinkOutcome, UnlinkError>>>, unlink_cancel: Option<CancellationToken>, unlink_origin: Option<UnlinkOrigin>, backoff: ExponentialBackoff, event_tx: Sender<NodeCoreEvent>, unlink_event_tx: UnboundedSender<NodeCoreEvent>, started: bool,
}
Expand description

Event-driven node controller.

Fields§

§runtime: Runtime§config: NodeConfig§sentinel: SentinelClient§worker_pool: WorkerPool§state: NodeState§connection_status: ConnectionStatus§registered: bool§node_name: String§total_cores: u32§app_name: String§app_version: String§work_context_cache: WorkContextCache§verify_rx: Option<Receiver<Result<Sensitive<String>, VerifyError>>>§register_rx: Option<Receiver<RegisterResult>>§realtime_rx: Option<Receiver<RealtimeEvent>>§realtime_shutdown: Option<CancellationToken>§result_rx: Option<Receiver<WorkBatchResult>>§unlink_rx: Option<Receiver<Result<UnlinkOutcome, UnlinkError>>>§unlink_cancel: Option<CancellationToken>§unlink_origin: Option<UnlinkOrigin>§backoff: ExponentialBackoff§event_tx: Sender<NodeCoreEvent>§unlink_event_tx: UnboundedSender<NodeCoreEvent>§started: bool

Implementations§

Source§

impl NodeCore

Source

pub(super) fn start_verification(&mut self)

Source

pub(super) fn check_verification(&mut self)

Source

pub(super) fn check_retry(&mut self)

Source

pub(super) fn start_registration(&mut self)

Source

pub(super) fn start_realtime(&mut self)

Source

pub(super) fn check_registration(&mut self)

Source

fn set_unavailable(&mut self)

Source§

impl NodeCore

Source

pub fn start(&mut self)

Start async tasks. Call once after creation.

Source

pub fn set_token_claim(&mut self, token: String)

Set claim token and start registration.

Source

pub fn poll(&mut self)

Poll for updates. Call periodically.

Source

pub fn drive_until_stopped( &mut self, events: &mut NodeEvents, running: &AtomicBool, on_event: impl FnMut(&NodeCoreEvent, &Self), )

Drive the node and its events on the current thread until running becomes false.

Source

pub fn state(&self) -> &NodeState

Source

pub fn connection_status(&self) -> ConnectionStatus

Source

pub fn public_key(&self) -> &NodePublicKey

Source

pub fn is_registered(&self) -> bool

Source

pub fn node_name(&self) -> &str

Source

pub fn stats(&self) -> NodeStats

Source

pub fn time_until_retry(&self) -> Option<Duration>

Source

pub fn disconnect(&mut self)

Source

pub fn reconnect(&mut self)

Source

pub fn runtime_handle(&self) -> &Handle

Source

pub(super) fn set_state(&mut self, state: NodeState)

Source

pub(super) fn set_connection(&mut self, status: ConnectionStatus)

Source§

impl NodeCore

Start an observable unlink operation; repeated calls while active are idempotent.

Cancel an active unlink without deleting the persisted identity.

Source

pub fn is_unlinking(&self) -> bool

Whether an unlink request is awaiting terminal completion.

Source§

impl NodeCore

Source

pub(super) fn check_realtime_events(&mut self)

Source

pub(super) fn check_work_results(&mut self)

Source

fn handle_realtime_event(&mut self, event: &RealtimeEvent)

Source

fn handle_node_update( &mut self, name: &str, total_cores: i32, max_parallel: i32, )

Source

fn handle_chunk_assigned(&mut self, payload: &RuntimeChunkPayload)

Source

fn process_chunk(&mut self, job_id: Uuid, payload: &RuntimeChunkPayload)

Source

fn handle_work_result(&self, result: &WorkBatchResult)

Source§

impl NodeCore

Source

pub fn with_resolver<R>( config: NodeConfig, resolver: R, application: NodeApplication, catalog: &'static ContentCatalog, ) -> Result<(Self, NodeEvents), SentinelError>
where R: DataResolver + Send + Sync + 'static,

Construct a node using a caller-supplied data resolver.

§Errors

Returns an error when the Tokio runtime or signed Sentinel client cannot be initialized.

Source

pub fn with_supabase( config: NodeConfig, application: NodeApplication, catalog: &'static ContentCatalog, ) -> Result<(Self, NodeEvents), SentinelError>

Construct a node wired to the Supabase data resolver.

§Errors

Returns an error when the Tokio runtime, Supabase client, metadata lookup, or data resolver cannot be initialized.

Source

fn new( runtime: Runtime, resolver: Arc<DynDataResolver<'static>>, config: NodeConfig, application: NodeApplication, catalog: &'static ContentCatalog, ) -> Result<(Self, NodeEvents), SentinelError>

Trait Implementations§

Source§

impl Debug for NodeCore

Source§

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

Formats the value using the given formatter. Read more

Auto Trait Implementations§

Blanket Implementations§

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> 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<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<T> Instrument for T

§

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

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

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
§

impl<T> Pointable for T

§

const ALIGN: usize

The alignment of pointer.
§

type Init = T

The type for initializers.
§

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

Initializes a with the given initializer. Read more
§

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

Dereferences the given pointer. Read more
§

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

Mutably dereferences the given pointer. Read more
§

unsafe fn drop(ptr: usize)

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

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

§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: 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
§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: 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> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> SpecConfig for T
where T: Any + Debug + Send,

Source§

fn as_any(&self) -> &(dyn Any + 'static)

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<S, T> Upcast<T> for S
where T: UpcastFrom<S> + ?Sized, S: ?Sized,

Source§

fn upcast(&self) -> &T
where Self: ErasableGeneric, T: ErasableGeneric<Repr = Self::Repr>,

Perform a zero-cost type-safe upcast to a wider ref type within the Wasm bindgen generics type system. Read more
Source§

fn upcast_into(self) -> T
where Self: Sized + ErasableGeneric, T: ErasableGeneric<Repr = Self::Repr>,

Perform a zero-cost type-safe upcast to a wider type within the Wasm bindgen generics type system. Read more
§

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

§

fn vzip(self) -> V

§

impl<T> WithSubscriber for T

§

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
§

fn with_current_subscriber(self) -> WithDispatch<Self>

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

impl<T> ResolverSyncBound for T
where T: Sync + ?Sized,