Skip to main content

Engine

Struct Engine 

Source
pub struct Engine {
    pub router: Arc<PartitionRouter>,
    /* private fields */
}
Expand description

The FlowFabric engine: partition routing + background scanners.

Fields§

§router: Arc<PartitionRouter>

Implementations§

Source§

impl Engine

Source

pub fn start(config: EngineConfig, client: Client) -> Self

Start the engine with the given config and Valkey client.

Spawns background scanner tasks. Returns immediately.

Source

pub fn start_with_metrics( config: EngineConfig, client: Client, metrics: Arc<Metrics>, ) -> Self

PR-94: start the engine with a shared observability registry.

Used by ff-server so scanner cycle metrics funnel into the same Prometheus registry exposed at /metrics. Under the observability feature (flipped via the same feature on ff-server / ff-engine), the handle records into an OTEL MeterProvider; otherwise the shim no-ops.

Source

pub fn start_with_completions( config: EngineConfig, client: Client, metrics: Arc<Metrics>, completions: CompletionStream, ) -> Self

Start the engine with a shared observability registry and a completion stream for push-based DAG promotion (issue #90).

The stream is typically produced by ff_core::completion_backend::CompletionBackend::subscribe_completions. The engine spawns a dispatch loop that drains the stream and fires ff_resolve_dependency per completion, reducing DAG latency from interval × levels to ~RTT × levels. The dependency_reconciler scanner remains as a safety net for completions missed during subscriber reconnect windows.

Source

pub async fn shutdown(self)

Signal all scanners to stop and wait for them to finish.

Waits up to 15 seconds for scanners to drain. If any scanner is blocked on a hung Valkey command, the timeout prevents shutdown from hanging indefinitely (Kubernetes SIGKILL safety).

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.

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, 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