Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions docs-internal/engine/sleep-sequence.md
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@ When the grace deadline elapses before `can_finalize_sleep()` returns true:
- The final `ActorStateStopped` for a lost generation is always `StopCode::Error` with "envoy connection lost", including when Lost escalates a pending graceful stop. `StopCode::Ok` would make the engine destroy the actor.
- Already-submitted SQLite commits cannot be recalled. Depot generation checks are the backstop for those.
- A lost generation starts no new user callback. Every runtime callback passes through `refuseCallbacksAfterLost` in `registry/native.ts` when JS enters it, except the `onSleep` and `onDestroy` cleanup wrappers, which release runtime state and skip the user hook themselves. Callbacks the framework starts on its own (workflow steps, run, inspector workflow calls, deferred `onStateChange`, WebSocket open and message listeners, database provider hooks) check the lost state first. Continuations inside a user call that is already running, and abort or close listeners, are not blocked.
- envoy-client declares all its actors lost once no engine ping has arrived for `envoy_lost_threshold` minus a 3s margin. The engine refreshes its liveness timestamp just before each ping and reallocates actors once that timestamp is older than the threshold, so the envoy normally stops first, including on half-open connections that never report a close. This is best effort: a ping delayed in transit by more than the margin can still let the engine give up first.

## Guarding lifecycle requests

Expand Down
125 changes: 93 additions & 32 deletions engine/sdks/rust/envoy-client/src/envoy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,10 @@ use crate::utils::{BufferMap, EnvoyShutdownError, SleepFuture, boxed_sleep, spaw
#[cfg(not(target_arch = "wasm32"))]
static GLOBAL_ENVOY: OnceLock<Mutex<Option<EnvoyHandle>>> = OnceLock::new();

/// How long before the engine's lost threshold the envoy stops its own actors. Covers ping
/// delivery latency and the jitter in the engine's ping interval.
const LOST_THRESHOLD_SAFETY_MARGIN_MS: i64 = 3_000;

pub struct EnvoyContext {
pub shared: Arc<SharedContext>,
pub shutting_down: bool,
Expand Down Expand Up @@ -487,6 +491,7 @@ async fn envoy_loop(
let iter_start = crate::time::Instant::now();
#[allow(unused_assignments)]
let mut branch: &'static str = "unknown";
let ping_silence_wait = engine_ping_silence_wait(&ctx);
Comment on lines 491 to +494

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔴 High · Ping liveness is not synchronized with the select loop or connection session

last_ping_ts is updated in forward_to_envoy without waking this loop. On a fresh connection, command replay can start actors before the ping task sends its first ping; this iteration therefore builds None, and subsequent pings do not arm a deadline until another envoy message or the 15-second KV cleanup tick. With the default 15-second lost threshold, a link that goes half-open in that window can let the engine expire the envoy before this branch runs, defeating the single-writer protection this change is meant to add. The timestamp also survives reconnects, so a reconnect after the old deadline immediately loses still-running or replayed actors before the new connection's first ping. Make ping reception/session changes an event observed by this loop (for example, a watch channel carrying the current session's last-ping value), reset it when a connection is established, and derive/restart the deadline from that event.

tokio::select! {
msg = rx.recv() => {
branch = "envoy_msg";
Expand Down Expand Up @@ -614,33 +619,24 @@ async fn envoy_loop(
}
} => {
branch = "lost_timeout";
// Lost timeout fired
for (_id, request) in ctx.kv_requests.drain() {
METRICS.kv_requests_inflight.dec();
let _ = request.response_tx.send(Err(anyhow::anyhow!(EnvoyShutdownError)));
declare_actors_lost(&mut ctx, "stopping all actors due to envoy lost threshold");
lost_timeout = None;
}
_ = async {
match ping_silence_wait {
Some(wait) => crate::utils::sleep(wait).await,
None => std::future::pending::<()>().await,
}
fail_sqlite_requests_with_shutdown(&mut ctx);
fail_remote_sqlite_requests_with_shutdown(&mut ctx);

if !ctx.actors.is_empty() {
tracing::warn!("stopping all actors due to envoy lost threshold");
for (_actor_id, gens) in &ctx.actors {
for (_g, entry) in gens {
entry.lost.cancel();
if !entry.handle.is_closed() {
let _ = entry.handle.send(ToActor::Lost);
}
}
}
ctx.actors.clear();
ctx.shared
.actors
.lock()
.expect("shared actor registry poisoned")
.clear();
}, if !ctx.actors.is_empty() => {
branch = "engine_ping_silence";
// A ping may have arrived while this branch slept.
if engine_ping_silence_expired(&ctx, crate::time::now_millis()) {
declare_actors_lost(
&mut ctx,
"stopping all actors because the engine stopped pinging and is about to declare them lost",
);
lost_timeout = None;
}

lost_timeout = None;
}
}
observe_envoy_loop_iteration(branch, iter_start);
Expand Down Expand Up @@ -787,13 +783,7 @@ fn handle_conn_close(ctx: &EnvoyContext, lost_timeout: Option<SleepFuture>) -> O
return lost_timeout;
}

// Read threshold from protocol metadata, fall back to 10 seconds
let lost_threshold = {
let metadata = ctx.shared.protocol_metadata.try_lock().ok();
metadata
.and_then(|guard| guard.as_ref().map(|m| m.envoy_lost_threshold as u64))
.unwrap_or(10_000)
};
let lost_threshold = envoy_lost_threshold_ms(ctx) as u64;

tracing::debug!(ms = lost_threshold, "starting envoy lost timeout");

Expand All @@ -802,6 +792,77 @@ fn handle_conn_close(ctx: &EnvoyContext, lost_timeout: Option<SleepFuture>) -> O
)))
}

/// Reads the engine's lost threshold from protocol metadata, falling back to 10 seconds.
fn envoy_lost_threshold_ms(ctx: &EnvoyContext) -> i64 {
let metadata = ctx.shared.protocol_metadata.try_lock().ok();
metadata
.and_then(|guard| guard.as_ref().map(|m| m.envoy_lost_threshold))
.unwrap_or(10_000)
}

/// When the envoy stops trusting its actors after the last engine ping. The engine refreshes its
/// envoy liveness timestamp right before sending each ping and declares the envoy's actors lost
/// once that timestamp is older than `envoy_lost_threshold`, then may start them elsewhere. The
/// envoy receives each ping after that refresh, so giving up a margin before the threshold normally
/// stops the actors before the engine can reallocate them. A ping delayed in transit by more than
/// the margin can still let the engine give up first. This also covers half-open connections that
/// never report a close.
pub fn engine_ping_silence_deadline_ms(ctx: &EnvoyContext) -> Option<i64> {
let last_ping_ts = ctx.shared.last_ping_ts.load(Ordering::Acquire);
if last_ping_ts == 0 {
return None;
}
let threshold = envoy_lost_threshold_ms(ctx);
let margin = LOST_THRESHOLD_SAFETY_MARGIN_MS.min(threshold / 2);
Some(last_ping_ts + threshold - margin)
}

pub fn engine_ping_silence_expired(ctx: &EnvoyContext, now_ms: i64) -> bool {
engine_ping_silence_deadline_ms(ctx).is_some_and(|deadline| now_ms >= deadline)
}

fn engine_ping_silence_wait(ctx: &EnvoyContext) -> Option<std::time::Duration> {
let deadline = engine_ping_silence_deadline_ms(ctx)?;
let remaining = (deadline - crate::time::now_millis()).max(0);
Some(std::time::Duration::from_millis(remaining as u64))
}

/// Stops every actor on this envoy because the engine has given up on them. The lost token is
/// cancelled first so each generation stops writing storage before its task sees the message.
pub fn declare_actors_lost(ctx: &mut EnvoyContext, message: &'static str) {
for (_id, request) in ctx.kv_requests.drain() {
METRICS.kv_requests_inflight.dec();
let _ = request
.response_tx
.send(Err(anyhow::anyhow!(EnvoyShutdownError)));
}
fail_sqlite_requests_with_shutdown(ctx);
fail_remote_sqlite_requests_with_shutdown(ctx);

if ctx.actors.is_empty() {
return;
}
tracing::warn!(
actor_count = ctx.actors.len(),
reason = message,
"declaring envoy actors lost"
);
for (_actor_id, gens) in &ctx.actors {
for (_g, entry) in gens {
entry.lost.cancel();
if !entry.handle.is_closed() {
let _ = entry.handle.send(ToActor::Lost);
}
}
}
ctx.actors.clear();
ctx.shared
.actors
.lock()
.expect("shared actor registry poisoned")
.clear();
}

async fn handle_shutdown(ctx: &mut EnvoyContext) {
if ctx.shutting_down {
return;
Expand Down
65 changes: 64 additions & 1 deletion engine/sdks/rust/envoy-client/tests/command_dedup.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,9 @@ use rivet_envoy_client::config::{
WebSocketSender,
};
use rivet_envoy_client::context::{SharedContext, WsTxMessage};
use rivet_envoy_client::envoy::EnvoyContext;
use rivet_envoy_client::envoy::{
EnvoyContext, declare_actors_lost, engine_ping_silence_deadline_ms, engine_ping_silence_expired,
};
use rivet_envoy_client::handle::EnvoyHandle;
use rivet_envoy_client::sqlite::{
RemoteSqliteRequest, fail_sent_remote_sqlite_requests_with_indeterminate_result,
Expand Down Expand Up @@ -567,3 +569,64 @@ async fn queued_kv_write_is_dropped_once_its_generation_is_lost() {
.expect_err("the lost request must fail");
assert!(format!("{error:#}").contains("declared lost"));
}

/// The engine declares actors lost `envoy_lost_threshold` after its last ping refresh. The envoy
/// gives up a safety margin earlier, measured from the last ping it received, so its actors stop
/// before the engine can start them elsewhere.
#[tokio::test]
async fn engine_ping_silence_expires_before_the_engine_lost_threshold() {
let ctx = new_envoy_context();
assert_eq!(
engine_ping_silence_deadline_ms(&ctx),
None,
"no deadline before the first engine ping"
);

*ctx.shared.protocol_metadata.lock().await = Some(protocol::ProtocolMetadata {
envoy_lost_threshold: 15_000,
actor_stop_threshold: 1_800_000,
max_response_payload_size: 20_971_520,
});
let last_ping_ts = 1_700_000_000_000;
ctx.shared
.last_ping_ts
.store(last_ping_ts, std::sync::atomic::Ordering::Release);

let deadline = engine_ping_silence_deadline_ms(&ctx).expect("deadline after a ping");
assert!(
deadline < last_ping_ts + 15_000,
"the envoy must give up before the engine's threshold"
);
assert!(
deadline > last_ping_ts + 3_000,
"the envoy must tolerate several missed ping intervals"
);
assert!(!engine_ping_silence_expired(&ctx, deadline - 1));
assert!(engine_ping_silence_expired(&ctx, deadline));
}

#[tokio::test]
async fn declaring_actors_lost_signals_every_generation_before_delivery() {
let mut ctx = new_envoy_context();
let (actor_tx, mut actor_rx) = mpsc::unbounded_channel::<ToActor>();
let first = tokio_util::sync::CancellationToken::new();
let second = tokio_util::sync::CancellationToken::new();
for (actor_id, lost) in [("actor-a", first.clone()), ("actor-b", second.clone())] {
ctx.insert_actor_with_lost_signal(
actor_id.to_string(),
1,
actor_tx.clone(),
lost,
Arc::new(AsyncCounter::new()),
actor_id.to_string(),
-1,
);
}

declare_actors_lost(&mut ctx, "test");

assert!(first.is_cancelled() && second.is_cancelled());
assert!(ctx.actors.is_empty());
assert!(matches!(actor_rx.try_recv(), Ok(ToActor::Lost)));
assert!(matches!(actor_rx.try_recv(), Ok(ToActor::Lost)));
}
Loading