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>
impl<S> HeartbeatSweeper<S>
Sourcepub fn new(
tracker: HeartbeatTracker,
registry: ConnectedWorkerRegistry,
sink: S,
drain: DrainState,
heartbeat_window: Duration,
) -> Self
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.
Sourcepub fn with_queue_state(self, queue_state: QueueServiceState) -> Self
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.
Sourcepub async fn run(self, shutdown: Receiver<bool>)
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> 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> 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> IntoMaybeUndefined<T> for T
impl<T> IntoMaybeUndefined<T> for T
Source§fn into_maybe_undefined(self) -> MaybeUndefined<T>
fn into_maybe_undefined(self) -> MaybeUndefined<T>
Source§impl<T> IntoOption<T> for T
impl<T> IntoOption<T> for T
Source§fn into_option(self) -> Option<T>
fn into_option(self) -> Option<T>
Source§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