Skip to main content

HeartbeatSweeper

Struct HeartbeatSweeper 

Source
pub struct HeartbeatSweeper<S> { /* private fields */ }
Expand description

Production driver of HeartbeatTracker::fail_expired_workers (#176).

The tracker records connection and per-task liveness, while the stream-teardown sweep fails a worker whose stream ENDS. A worker whose stream stays open while its process wedges is caught by the connection lease even when it is idle. This interval task expires every silent connection or task, deregistering it with the provable WorkerDeathReason::Timeout and surfacing its tasks as TRANSPORT losses through the shared completion sink — the lost: class the engine re-dispatches attempt-neutrally, never the action’s retry vocabulary. It shares the server’s shutdown watch, so it drains with the transports (mirroring OutboxDispatcher::run).

Double-fail safety: this sweep and the stream-teardown path (HeartbeatTracker::fail_disconnected_worker) can both observe the same dead worker. Both funnel into the same idempotent core — deregister_with_reason is a no-op for an already-removed worker (no duplicate WS3 delta, no metrics double-count) and the tracker removes each task as it fails it — so whichever path runs second sees an empty report and never double-completes an activity.

Implementations§

Source§

impl<S> HeartbeatSweeper<S>
where S: ActivityCompletionSink + Send + Sync + 'static,

Source

pub fn new( tracker: HeartbeatTracker, registry: ConnectedWorkerRegistry, sink: S, drain: DrainState, heartbeat_window: Duration, ) -> Self

Build a sweeper over the server’s shared liveness tracker, worker registry, completion sink, and drain gate. The cadence is derived from heartbeat_window by sweep_interval.

Source

pub fn with_queue_state(self, queue_state: QueueServiceState) -> Self

Share the live unserved-queue state so a deregistration log can state how many dispatches are already parked on the queue the reaped worker served.

Without it the count reads zero — honest for a wiring with no queue service, and never a reason to withhold the deregistration itself.

Source

pub async fn run(self, shutdown: Receiver<bool>)

Run the expiry sweep until shutdown flips to true.

A tracker/registry error during a sweep is logged and retried next tick rather than tearing the task down — a transient failure must not silently stop dead-worker detection. Shutdown is observed both while waiting for the next tick and re-checked before each sweep, exactly like the outbox dispatcher’s run loop.

Auto Trait Implementations§

§

impl<S> !RefUnwindSafe for HeartbeatSweeper<S>

§

impl<S> !UnwindSafe for HeartbeatSweeper<S>

§

impl<S> Freeze for HeartbeatSweeper<S>
where S: Freeze,

§

impl<S> Send for HeartbeatSweeper<S>
where S: Send,

§

impl<S> Sync for HeartbeatSweeper<S>
where S: Sync,

§

impl<S> Unpin for HeartbeatSweeper<S>
where S: Unpin,

§

impl<S> UnsafeUnpin for HeartbeatSweeper<S>
where S: UnsafeUnpin,

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

Source§

fn into_maybe_undefined(self) -> MaybeUndefined<T>

Converts this value into a three-state builder argument.
Source§

impl<T> IntoOption<T> for T

Source§

fn into_option(self) -> Option<T>

Converts this value into an optional builder argument.
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> 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 = !

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

fn try_from(value: U) -> Result<T, !>

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