Skip to main content

rivet_envoy_client/
envoy.rs

1use std::collections::HashMap;
2use std::sync::Arc;
3#[cfg(not(target_arch = "wasm32"))]
4use std::sync::OnceLock;
5use std::sync::atomic::Ordering;
6
7#[cfg(not(target_arch = "wasm32"))]
8use parking_lot::Mutex;
9
10use crate::async_counter::AsyncCounter;
11use rivet_envoy_protocol as protocol;
12use tokio::sync::mpsc;
13use tokio::sync::oneshot;
14use tracing::Instrument;
15
16use crate::actor::ToActor;
17use crate::commands::{ACK_COMMANDS_INTERVAL_MS, handle_commands, send_command_ack};
18use crate::config::EnvoyConfig;
19use crate::connection::{start_connection, ws_send};
20use crate::context::{SharedContext, WsTxMessage};
21use crate::events::{handle_ack_events, handle_send_events, resend_unacknowledged_events};
22use crate::handle::EnvoyHandle;
23use crate::kv::{
24	KV_CLEANUP_INTERVAL_MS, KvRequestEntry, cleanup_old_kv_requests, handle_kv_request,
25	handle_kv_response, process_unsent_kv_requests,
26};
27use crate::metrics::METRICS;
28use crate::sqlite::{
29	RemoteSqliteRequest, RemoteSqliteRequestEntry, RemoteSqliteResponseEnvelope, SqliteRequest,
30	SqliteRequestEntry, SqliteResponse, cleanup_old_remote_sqlite_requests,
31	cleanup_old_sqlite_requests, fail_remote_sqlite_requests_with_shutdown,
32	fail_sent_remote_sqlite_requests_with_indeterminate_result, fail_sqlite_requests_with_shutdown,
33	handle_remote_sqlite_exec_response, handle_remote_sqlite_execute_batch_response,
34	handle_remote_sqlite_execute_response, handle_remote_sqlite_request,
35	handle_sqlite_commit_response, handle_sqlite_get_pages_response, handle_sqlite_request,
36	process_unsent_remote_sqlite_requests, process_unsent_sqlite_requests,
37};
38use crate::tunnel::{
39	HttpRequestCancellationKey, handle_tunnel_message, make_ws_key,
40	resend_buffered_tunnel_messages, send_hibernatable_ws_message_ack,
41};
42use crate::utils::{BufferMap, EnvoyShutdownError, SleepFuture, boxed_sleep, spawn_detached};
43
44/// Process-wide envoy slot. Holds the handle inside a mutex so a stopped
45/// handle (e.g. from a shutdown-during-build race in serverless mode) can be
46/// replaced on the next `start_envoy_sync` call.
47#[cfg(not(target_arch = "wasm32"))]
48static GLOBAL_ENVOY: OnceLock<Mutex<Option<EnvoyHandle>>> = OnceLock::new();
49
50pub struct EnvoyContext {
51	pub shared: Arc<SharedContext>,
52	pub shutting_down: bool,
53	pub actors: HashMap<String, HashMap<u32, ActorEntry>>,
54	pub buffered_actor_messages: HashMap<String, Vec<BufferedActorMessage>>,
55	pub kv_requests: HashMap<u32, KvRequestEntry>,
56	pub next_kv_request_id: u32,
57	pub sqlite_requests: HashMap<u32, SqliteRequestEntry>,
58	pub next_sqlite_request_id: u32,
59	pub remote_sqlite_requests: HashMap<u32, RemoteSqliteRequestEntry>,
60	pub next_remote_sqlite_request_id: u32,
61	pub request_to_actor: BufferMap<WebSocketRoute>,
62	pub http_request_routes: BufferMap<HttpRequestRoute>,
63	pub http_message_indices: BufferMap<protocol::MessageIndex>,
64	/// Recently cancelled exact-generation requests. This prevents a delayed
65	/// `RequestStart` from reaching the actor after its cancellation arrived.
66	pub http_request_cancellations: HashMap<HttpRequestCancellationKey, crate::time::Instant>,
67	pub buffered_messages: Vec<protocol::ToRivetTunnelMessage>,
68	/// Highest command index processed per `(actor_id, generation)`, used to
69	/// drop replayed commands from `pegboard-envoy` after a reconnect. Persists
70	/// across `remove_actor` so a replayed `CommandStartActor` for an
71	/// already-stopped actor cannot resurrect it.
72	pub processed_command_idx: HashMap<(String, u32), i64>,
73}
74
75pub struct HttpRequestRoute {
76	pub actor_id: String,
77	pub actor_generation: Option<u32>,
78	pub actor_admitted: bool,
79	pub session: u64,
80	pub gateway_id: protocol::GatewayId,
81	pub request_id: protocol::RequestId,
82}
83
84#[derive(Clone)]
85pub struct WebSocketRoute {
86	pub actor_id: String,
87	pub actor_generation: Option<u32>,
88}
89
90pub struct ActorEntry {
91	pub handle: mpsc::UnboundedSender<ToActor>,
92	pub active_http_request_count: Arc<AsyncCounter>,
93	pub name: String,
94	pub event_history: Vec<protocol::EventWrapper>,
95	pub last_command_idx: i64,
96	pub received_stop: bool,
97}
98
99pub enum BufferedActorMessage {
100	WsMsg {
101		message_id: protocol::MessageId,
102		msg: protocol::ToEnvoyWebSocketMessage,
103	},
104	WsClose {
105		message_id: protocol::MessageId,
106		close: protocol::ToEnvoyWebSocketClose,
107	},
108}
109
110pub enum ToEnvoyMessage {
111	ConnMessage {
112		message: protocol::ToEnvoy,
113		session: u64,
114	},
115	ConnClose {
116		evict: bool,
117		session: u64,
118	},
119	SendEvents {
120		events: Vec<protocol::EventWrapper>,
121	},
122	KvRequest {
123		actor_id: String,
124		data: protocol::KvRequestData,
125		response_tx: oneshot::Sender<anyhow::Result<protocol::KvResponseData>>,
126	},
127	SqliteRequest {
128		request: SqliteRequest,
129		response_tx: oneshot::Sender<anyhow::Result<SqliteResponse>>,
130	},
131	RemoteSqliteRequest {
132		request: RemoteSqliteRequest,
133		expected_session: Option<u64>,
134		response_tx: oneshot::Sender<anyhow::Result<RemoteSqliteResponseEnvelope>>,
135	},
136	SendOrBufferTunnelMsg {
137		msg: protocol::ToRivetTunnelMessage,
138	},
139	ActorIntent {
140		actor_id: String,
141		generation: Option<u32>,
142		intent: protocol::ActorIntent,
143		error: Option<String>,
144	},
145	SetAlarm {
146		actor_id: String,
147		generation: Option<u32>,
148		alarm_ts: Option<i64>,
149		ack_tx: Option<oneshot::Sender<()>>,
150	},
151	HwsAck {
152		gateway_id: protocol::GatewayId,
153		request_id: protocol::RequestId,
154		envoy_message_index: u16,
155	},
156	RebindWebSocket {
157		actor_id: String,
158		generation: u32,
159		gateway_id: protocol::GatewayId,
160		request_id: protocol::RequestId,
161		response_tx: oneshot::Sender<bool>,
162	},
163	HttpRequestComplete {
164		gateway_id: protocol::GatewayId,
165		request_id: protocol::RequestId,
166	},
167	GetActor {
168		actor_id: String,
169		generation: Option<u32>,
170		response_tx: oneshot::Sender<Option<ActorInfo>>,
171	},
172	Shutdown,
173	Stop,
174}
175
176/// Information about an actor, returned by `EnvoyHandle::get_actor`.
177#[derive(Clone)]
178pub struct ActorInfo {
179	pub name: String,
180	pub generation: u32,
181	pub active_http_request_count: Arc<AsyncCounter>,
182}
183
184impl EnvoyContext {
185	pub(crate) fn rebind_websocket(
186		&mut self,
187		actor_id: &str,
188		generation: u32,
189		gateway_id: &protocol::GatewayId,
190		request_id: &protocol::RequestId,
191	) -> bool {
192		if let Some(route) = self.request_to_actor.get_mut(&[gateway_id, request_id]) {
193			if route.actor_id != actor_id {
194				return false;
195			}
196			if route.actor_generation.is_some() {
197				route.actor_generation = Some(generation);
198			}
199		} else {
200			// A hibernating actor may be restored on a different Envoy process. Its
201			// authoritative start command carries the hibernating request ids, while
202			// the original process owns the old ephemeral route map.
203			self.request_to_actor.insert(
204				&[gateway_id, request_id],
205				WebSocketRoute {
206					actor_id: actor_id.to_owned(),
207					actor_generation: Some(generation),
208				},
209			);
210		}
211
212		self.shared
213			.live_tunnel_requests
214			.lock()
215			.expect("shared live tunnel request registry poisoned")
216			.insert(make_ws_key(gateway_id, request_id), actor_id.to_owned());
217		true
218	}
219
220	pub fn insert_actor(
221		&mut self,
222		actor_id: String,
223		generation: u32,
224		handle: mpsc::UnboundedSender<ToActor>,
225		active_http_request_count: Arc<AsyncCounter>,
226		name: String,
227		last_command_idx: i64,
228	) {
229		let buffered_actor_id = actor_id.clone();
230		let buffered_handle = handle.clone();
231		self.actors
232			.entry(actor_id.clone())
233			.or_insert_with(HashMap::new)
234			.insert(
235				generation,
236				ActorEntry {
237					handle: handle.clone(),
238					active_http_request_count: active_http_request_count.clone(),
239					name,
240					event_history: Vec::new(),
241					last_command_idx,
242					received_stop: false,
243				},
244			);
245		self.shared
246			.actors
247			.lock()
248			.expect("shared actor registry poisoned")
249			.entry(actor_id)
250			.or_insert_with(HashMap::new)
251			.insert(
252				generation,
253				crate::context::SharedActorEntry {
254					handle,
255					active_http_request_count,
256				},
257			);
258
259		self.shared.actors_notify.notify_waiters();
260
261		if let Some(messages) = self.buffered_actor_messages.remove(&buffered_actor_id) {
262			for message in messages {
263				match message {
264					BufferedActorMessage::WsMsg { message_id, msg } => {
265						let _ = buffered_handle.send(ToActor::WsMsg { message_id, msg });
266					}
267					BufferedActorMessage::WsClose { message_id, close } => {
268						let _ = buffered_handle.send(ToActor::WsClose { message_id, close });
269					}
270				}
271			}
272		}
273	}
274
275	pub fn remove_actor(&mut self, actor_id: &str, generation: u32) {
276		if let Some(generations) = self.actors.get_mut(actor_id) {
277			generations.remove(&generation);
278			if generations.is_empty() {
279				self.actors.remove(actor_id);
280			}
281		}
282
283		let mut shared = self
284			.shared
285			.actors
286			.lock()
287			.expect("shared actor registry poisoned");
288		if let Some(generations) = shared.get_mut(actor_id) {
289			generations.remove(&generation);
290			if generations.is_empty() {
291				shared.remove(actor_id);
292			}
293		}
294		self.shared.actors_notify.notify_waiters();
295	}
296
297	pub fn get_actor(&self, actor_id: &str, generation: Option<u32>) -> Option<&ActorEntry> {
298		let gens = self.actors.get(actor_id)?;
299		if gens.is_empty() {
300			return None;
301		}
302
303		if let Some(g) = generation {
304			return gens.get(&g).filter(|entry| !entry.handle.is_closed());
305		}
306
307		// Return highest generation non-closed entry
308		// HashMap doesn't guarantee order, so find max key
309		let mut best: Option<&ActorEntry> = None;
310		let mut best_gen: u32 = 0;
311		for (&g, entry) in gens {
312			if !entry.handle.is_closed() && (best.is_none() || g > best_gen) {
313				best = Some(entry);
314				best_gen = g;
315			}
316		}
317		best
318	}
319
320	/// Selects an actor for a new request. Exact-generation continuations use
321	/// `get_actor` directly so an already admitted stream can finish during stop.
322	pub fn get_actor_for_admission(
323		&self,
324		actor_id: &str,
325		generation: Option<u32>,
326	) -> Option<&ActorEntry> {
327		let actor = self.get_actor(actor_id, generation)?;
328		if generation.is_some() && actor.received_stop {
329			None
330		} else {
331			Some(actor)
332		}
333	}
334
335	pub fn get_actor_entry_mut(
336		&mut self,
337		actor_id: &str,
338		generation: u32,
339	) -> Option<&mut ActorEntry> {
340		self.actors
341			.get_mut(actor_id)
342			.and_then(|gens| gens.get_mut(&generation))
343	}
344}
345
346pub async fn start_envoy(config: EnvoyConfig) -> EnvoyHandle {
347	let handle = start_envoy_sync(config);
348	handle
349		.started()
350		.await
351		.expect("envoy failed to start before returning handle");
352	handle
353}
354
355pub fn start_envoy_sync(config: EnvoyConfig) -> EnvoyHandle {
356	#[cfg(target_arch = "wasm32")]
357	{
358		start_envoy_sync_inner(config)
359	}
360
361	#[cfg(not(target_arch = "wasm32"))]
362	{
363		if config.not_global {
364			return start_envoy_sync_inner(config);
365		}
366
367		let slot = GLOBAL_ENVOY.get_or_init(|| Mutex::new(None));
368		let mut guard = slot.lock();
369		if let Some(handle) = guard.as_ref() {
370			if !handle.is_stopped() {
371				return handle.clone();
372			}
373		}
374		let handle = start_envoy_sync_inner(config);
375		*guard = Some(handle.clone());
376		handle
377	}
378}
379
380fn start_envoy_sync_inner(config: EnvoyConfig) -> EnvoyHandle {
381	let (envoy_tx, envoy_rx) = mpsc::unbounded_channel::<ToEnvoyMessage>();
382	let (start_tx, start_rx) = tokio::sync::watch::channel(());
383	let (stopped_tx, _stopped_rx) = tokio::sync::watch::channel(false);
384	let (connection_session_tx, _connection_session_rx) = tokio::sync::watch::channel(0);
385
386	let envoy_key = uuid::Uuid::new_v4().to_string();
387	let shared = Arc::new(SharedContext {
388		config,
389		envoy_key,
390		envoy_tx: envoy_tx.clone(),
391		actors: Arc::new(std::sync::Mutex::new(HashMap::new())),
392		actors_notify: Arc::new(tokio::sync::Notify::new()),
393		live_tunnel_requests: Arc::new(std::sync::Mutex::new(HashMap::new())),
394		pending_hibernation_restores: Arc::new(std::sync::Mutex::new(HashMap::new())),
395		ws_tx: Arc::new(tokio::sync::Mutex::new(None)),
396		http_ws_tx: Arc::new(tokio::sync::Mutex::new(None)),
397		connection_session: std::sync::atomic::AtomicU64::new(0),
398		next_connection_session: std::sync::atomic::AtomicU64::new(0),
399		connection_session_tx,
400		protocol_metadata: Arc::new(tokio::sync::Mutex::new(None)),
401		shutting_down: std::sync::atomic::AtomicBool::new(false),
402		last_ping_ts: std::sync::atomic::AtomicI64::new(0),
403		stopped_tx,
404	});
405
406	let handle = EnvoyHandle {
407		shared: shared.clone(),
408		started_rx: start_rx,
409	};
410
411	start_connection(shared.clone());
412
413	let ctx = EnvoyContext {
414		shared: shared.clone(),
415		shutting_down: false,
416		actors: HashMap::new(),
417		buffered_actor_messages: HashMap::new(),
418		kv_requests: HashMap::new(),
419		next_kv_request_id: 0,
420		sqlite_requests: HashMap::new(),
421		next_sqlite_request_id: 0,
422		remote_sqlite_requests: HashMap::new(),
423		next_remote_sqlite_request_id: 0,
424		request_to_actor: BufferMap::new(),
425		http_request_routes: BufferMap::new(),
426		http_message_indices: BufferMap::new(),
427		http_request_cancellations: HashMap::new(),
428		buffered_messages: Vec::new(),
429		processed_command_idx: HashMap::new(),
430	};
431
432	tracing::info!(envoy_key = %shared.envoy_key, "starting envoy");
433	let span = tracing::info_span!("envoy_client", envoy_key = %shared.envoy_key);
434	spawn_detached(envoy_loop(ctx, envoy_rx, start_tx).instrument(span));
435
436	handle
437}
438
439async fn envoy_loop(
440	mut ctx: EnvoyContext,
441	mut rx: mpsc::UnboundedReceiver<ToEnvoyMessage>,
442	start_tx: tokio::sync::watch::Sender<()>,
443) {
444	let mut ack_tick = boxed_sleep(std::time::Duration::from_millis(ACK_COMMANDS_INTERVAL_MS));
445	let mut kv_cleanup_tick = boxed_sleep(std::time::Duration::from_millis(KV_CLEANUP_INTERVAL_MS));
446
447	let mut lost_timeout: Option<SleepFuture> = None;
448
449	loop {
450		let iter_start = crate::time::Instant::now();
451		#[allow(unused_assignments)]
452		let mut branch: &'static str = "unknown";
453		tokio::select! {
454			msg = rx.recv() => {
455				branch = "envoy_msg";
456				let Some(msg) = msg else {
457					observe_envoy_loop_iteration(branch, iter_start);
458					break;
459				};
460				METRICS.envoy_tx_depth.dec();
461
462				match msg {
463					ToEnvoyMessage::ConnMessage { message, session } => {
464						lost_timeout = handle_conn_message(&mut ctx, &start_tx, lost_timeout, message, session).await;
465					}
466					ToEnvoyMessage::ConnClose { evict, session } => {
467						remove_http_routes_for_session(&mut ctx, session);
468						for generations in ctx.actors.values() {
469							for actor in generations.values() {
470								let _ = actor.handle.send(ToActor::ConnectionClosed { session });
471							}
472						}
473						fail_sent_remote_sqlite_requests_with_indeterminate_result(&mut ctx);
474						lost_timeout = handle_conn_close(&ctx, lost_timeout);
475						if evict {
476							observe_envoy_loop_iteration(branch, iter_start);
477							break;
478						}
479					}
480					ToEnvoyMessage::SendEvents { events } => {
481						handle_send_events(&mut ctx, events).await;
482					}
483					ToEnvoyMessage::KvRequest { actor_id, data, response_tx } => {
484						handle_kv_request(&mut ctx, actor_id, data, response_tx).await;
485					}
486					ToEnvoyMessage::SqliteRequest { request, response_tx } => {
487						handle_sqlite_request(&mut ctx, request, response_tx).await;
488					}
489					ToEnvoyMessage::RemoteSqliteRequest { request, expected_session, response_tx } => {
490						handle_remote_sqlite_request(&mut ctx, request, expected_session, response_tx).await;
491					}
492					ToEnvoyMessage::SendOrBufferTunnelMsg { msg } => {
493						crate::tunnel::send_or_buffer_tunnel_message(&mut ctx, msg).await;
494					}
495					ToEnvoyMessage::ActorIntent { actor_id, generation, intent, error } => {
496						if let Some(entry) = ctx.get_actor(&actor_id, generation) {
497							let _ = entry.handle.send(ToActor::Intent { intent, error });
498						}
499					}
500					ToEnvoyMessage::SetAlarm { actor_id, generation, alarm_ts, ack_tx } => {
501						if let Some(entry) = ctx.get_actor(&actor_id, generation) {
502							if let Err(error) = entry.handle.send(ToActor::SetAlarm { alarm_ts, ack_tx }) {
503								if let ToActor::SetAlarm { ack_tx: Some(ack_tx), .. } = error.0 {
504									let _ = ack_tx.send(());
505								}
506							}
507						} else if let Some(ack_tx) = ack_tx {
508							let _ = ack_tx.send(());
509						}
510					}
511					ToEnvoyMessage::HwsAck { gateway_id, request_id, envoy_message_index } => {
512						send_hibernatable_ws_message_ack(&mut ctx, gateway_id, request_id, envoy_message_index);
513					}
514					ToEnvoyMessage::RebindWebSocket { actor_id, generation, gateway_id, request_id, response_tx } => {
515						let rebound = ctx.rebind_websocket(
516							&actor_id,
517							generation,
518							&gateway_id,
519							&request_id,
520						);
521						let _ = response_tx.send(rebound);
522					}
523					ToEnvoyMessage::HttpRequestComplete { gateway_id, request_id } => {
524						ctx.http_request_routes.remove(&[&gateway_id, &request_id]);
525						ctx.http_message_indices.remove(&[&gateway_id, &request_id]);
526					}
527					ToEnvoyMessage::GetActor { actor_id, generation, response_tx } => {
528						let info = ctx.get_actor(&actor_id, generation).map(|entry| {
529							let actor_gen = generation.unwrap_or_else(|| {
530								ctx.actors
531									.get(&actor_id)
532									.and_then(|gens| {
533										gens.iter()
534											.filter(|(_, e)| !e.handle.is_closed())
535											.map(|(&g, _)| g)
536											.max()
537									})
538									.unwrap_or(0)
539							});
540							ActorInfo {
541								name: entry.name.clone(),
542								generation: actor_gen,
543								active_http_request_count: entry
544									.active_http_request_count
545									.clone(),
546							}
547						});
548						let _ = response_tx.send(info);
549					}
550					ToEnvoyMessage::Shutdown => {
551						handle_shutdown(&mut ctx).await;
552					}
553					ToEnvoyMessage::Stop => {
554						observe_envoy_loop_iteration(branch, iter_start);
555						break;
556					}
557				}
558			}
559			_ = ack_tick.as_mut() => {
560				branch = "ack_tick";
561				send_command_ack(&mut ctx).await;
562				ack_tick = boxed_sleep(std::time::Duration::from_millis(ACK_COMMANDS_INTERVAL_MS));
563			}
564			_ = kv_cleanup_tick.as_mut() => {
565				branch = "cleanup_tick";
566				cleanup_old_kv_requests(&mut ctx);
567				cleanup_old_sqlite_requests(&mut ctx);
568				cleanup_old_remote_sqlite_requests(&mut ctx);
569				kv_cleanup_tick = boxed_sleep(std::time::Duration::from_millis(KV_CLEANUP_INTERVAL_MS));
570			}
571			_ = async {
572				match lost_timeout.as_mut() {
573					Some(timeout) => timeout.as_mut().await,
574					None => std::future::pending::<()>().await,
575				}
576			} => {
577				branch = "lost_timeout";
578				// Lost timeout fired
579				for (_id, request) in ctx.kv_requests.drain() {
580					METRICS.kv_requests_inflight.dec();
581					let _ = request.response_tx.send(Err(anyhow::anyhow!(EnvoyShutdownError)));
582				}
583				fail_sqlite_requests_with_shutdown(&mut ctx);
584				fail_remote_sqlite_requests_with_shutdown(&mut ctx);
585
586				if !ctx.actors.is_empty() {
587					tracing::warn!("stopping all actors due to envoy lost threshold");
588					for (_actor_id, gens) in &ctx.actors {
589						for (_g, entry) in gens {
590							if !entry.handle.is_closed() {
591								let _ = entry.handle.send(ToActor::Lost);
592							}
593						}
594					}
595					ctx.actors.clear();
596					ctx.shared
597						.actors
598						.lock()
599						.expect("shared actor registry poisoned")
600						.clear();
601				}
602
603				lost_timeout = None;
604			}
605		}
606		observe_envoy_loop_iteration(branch, iter_start);
607	}
608
609	// Cleanup
610	{
611		let guard = ctx.shared.ws_tx.lock().await;
612		if let Some(tx) = guard.as_ref() {
613			let _ = tx.send(WsTxMessage::Close);
614		}
615	}
616
617	for (_id, request) in ctx.kv_requests.drain() {
618		METRICS.kv_requests_inflight.dec();
619		let _ = request
620			.response_tx
621			.send(Err(anyhow::anyhow!("envoy shutting down")));
622	}
623	fail_sqlite_requests_with_shutdown(&mut ctx);
624	fail_remote_sqlite_requests_with_shutdown(&mut ctx);
625
626	ctx.actors.clear();
627	ctx.shared
628		.actors
629		.lock()
630		.expect("shared actor registry poisoned")
631		.clear();
632
633	tracing::info!("envoy stopped");
634
635	ctx.shared.config.callbacks.on_shutdown();
636
637	// Latched signal: waiters on `EnvoyHandle::wait_stopped` observe this and
638	// any future callers of `wait_stopped` resolve immediately because watch
639	// retains the last value.
640	let _ = ctx.shared.stopped_tx.send(true);
641}
642
643pub(crate) fn remove_http_routes_for_session(ctx: &mut EnvoyContext, session: u64) {
644	let closed_routes = ctx
645		.http_request_routes
646		.remove_where(|route| route.session == session);
647	for route in closed_routes {
648		ctx.http_message_indices
649			.remove(&[&route.gateway_id, &route.request_id]);
650	}
651}
652
653fn observe_envoy_loop_iteration(branch: &'static str, start: crate::time::Instant) {
654	let elapsed = start.elapsed();
655	METRICS
656		.envoy_loop_iteration_duration_seconds
657		.with_label_values(&[branch])
658		.observe(elapsed.as_secs_f64());
659}
660
661/// Send a message into the envoy_loop's mpsc and bump the depth gauge.
662/// Producers should prefer this over calling `shared.envoy_tx.send` directly
663/// so the `envoy_tx_depth` gauge stays in sync.
664pub fn send_to_envoy_tx(
665	shared: &crate::context::SharedContext,
666	msg: ToEnvoyMessage,
667) -> Result<(), tokio::sync::mpsc::error::SendError<ToEnvoyMessage>> {
668	match shared.envoy_tx.send(msg) {
669		Ok(()) => {
670			METRICS.envoy_tx_depth.inc();
671			Ok(())
672		}
673		Err(e) => Err(e),
674	}
675}
676
677async fn handle_conn_message(
678	ctx: &mut EnvoyContext,
679	start_tx: &tokio::sync::watch::Sender<()>,
680	mut lost_timeout: Option<SleepFuture>,
681	message: protocol::ToEnvoy,
682	session: u64,
683) -> Option<SleepFuture> {
684	match message {
685		protocol::ToEnvoy::ToEnvoyInit(init) => {
686			{
687				let mut guard = ctx.shared.protocol_metadata.lock().await;
688				*guard = Some(init.metadata.clone());
689			}
690			tracing::info!(?init.metadata, "received init");
691
692			lost_timeout = None;
693			resend_unacknowledged_events(ctx).await;
694			process_unsent_kv_requests(ctx).await;
695			process_unsent_sqlite_requests(ctx).await;
696			process_unsent_remote_sqlite_requests(ctx).await;
697			resend_buffered_tunnel_messages(ctx).await;
698
699			let _ = start_tx.send(());
700		}
701		protocol::ToEnvoy::ToEnvoyCommands(commands) => {
702			handle_commands(ctx, commands).await;
703		}
704		protocol::ToEnvoy::ToEnvoyAckEvents(ack) => {
705			handle_ack_events(ctx, ack);
706		}
707		protocol::ToEnvoy::ToEnvoyKvResponse(response) => {
708			handle_kv_response(ctx, response).await;
709		}
710		protocol::ToEnvoy::ToEnvoySqliteGetPagesResponse(response) => {
711			handle_sqlite_get_pages_response(ctx, response).await;
712		}
713		protocol::ToEnvoy::ToEnvoySqliteCommitResponse(response) => {
714			handle_sqlite_commit_response(ctx, response).await;
715		}
716		protocol::ToEnvoy::ToEnvoySqliteExecResponse(response) => {
717			handle_remote_sqlite_exec_response(ctx, response).await;
718		}
719		protocol::ToEnvoy::ToEnvoySqliteExecuteResponse(response) => {
720			handle_remote_sqlite_execute_response(ctx, response).await;
721		}
722		protocol::ToEnvoy::ToEnvoySqliteExecuteBatchResponse(response) => {
723			handle_remote_sqlite_execute_batch_response(ctx, response).await;
724		}
725		protocol::ToEnvoy::ToEnvoyTunnelMessage(tunnel_msg) => {
726			handle_tunnel_message(ctx, session, tunnel_msg).await;
727		}
728		protocol::ToEnvoy::ToEnvoyPing(_) => {
729			// Should be handled by connection task
730		}
731	}
732
733	lost_timeout
734}
735
736fn handle_conn_close(ctx: &EnvoyContext, lost_timeout: Option<SleepFuture>) -> Option<SleepFuture> {
737	if lost_timeout.is_some() {
738		return lost_timeout;
739	}
740
741	// Read threshold from protocol metadata, fall back to 10 seconds
742	let lost_threshold = {
743		let metadata = ctx.shared.protocol_metadata.try_lock().ok();
744		metadata
745			.and_then(|guard| guard.as_ref().map(|m| m.envoy_lost_threshold as u64))
746			.unwrap_or(10_000)
747	};
748
749	tracing::debug!(ms = lost_threshold, "starting envoy lost timeout");
750
751	Some(boxed_sleep(std::time::Duration::from_millis(
752		lost_threshold,
753	)))
754}
755
756async fn handle_shutdown(ctx: &mut EnvoyContext) {
757	if ctx.shutting_down {
758		return;
759	}
760	ctx.shutting_down = true;
761	ctx.shared.shutting_down.store(true, Ordering::Release);
762
763	tracing::debug!("envoy received shutdown");
764
765	ws_send(&ctx.shared, protocol::ToRivet::ToRivetStopping).await;
766
767	// Wait for all actors to finish. The process manager (Docker,
768	// k8s, etc.) provides the ultimate shutdown deadline.
769	let actor_handles: Vec<mpsc::UnboundedSender<ToActor>> = ctx
770		.actors
771		.values()
772		.flat_map(|gens| gens.values())
773		.filter(|entry| !entry.handle.is_closed())
774		.map(|entry| entry.handle.clone())
775		.collect();
776
777	let shared = ctx.shared.clone();
778	let shutdown_span = tracing::debug_span!(
779		parent: tracing::Span::current(),
780		"envoy_graceful_shutdown",
781		envoy_key = %ctx.shared.envoy_key,
782	);
783	spawn_detached(
784		async move {
785			futures_util::future::join_all(actor_handles.iter().map(|h| h.closed())).await;
786			tracing::debug!("all actors stopped during graceful shutdown");
787			let _ = send_to_envoy_tx(&shared, ToEnvoyMessage::Stop);
788		}
789		.instrument(shutdown_span),
790	);
791}