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#[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 pub http_request_cancellations: HashMap<HttpRequestCancellationKey, crate::time::Instant>,
67 pub buffered_messages: Vec<protocol::ToRivetTunnelMessage>,
68 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#[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 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 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 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 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 {
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 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
661pub 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 }
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 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 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}