1use anyhow::Result;
9use std::sync::Arc;
10
11use futu_gateway_core::bridge::GatewayBridge;
12use futu_server::listener_status::ListenerBindEvent;
13use futu_server::ws_listener::{WsServer, WsServerDeps};
14
15use crate::config::RuntimeConfig;
16use crate::startup::phase1::Phase1Out;
17use crate::startup::phase2::{AuthPlan, execute_auth_plan};
18use crate::startup::phase3::Phase3Out;
19
20mod listener_readiness;
21mod network_exposure;
22mod shutdown;
23mod summary;
24#[cfg(test)]
25mod tests;
26
27use listener_readiness::{
28 ActivatedListenerReadiness, EnabledListenerSurfaces, ListenerReadiness, SurfaceTasks,
29 await_listeners_opened, wrap_grpc_server_result,
30};
31#[cfg(test)]
32use network_exposure::legacy_network_exposure_warnings;
33use network_exposure::warn_legacy_network_exposure;
34use shutdown::{
35 BackgroundJoinHandle, SHUTDOWN_GRACE_PERIOD, SurfaceJoinHandle, await_background_shutdown,
36 await_surface_shutdown, log_surface_task_result, request_surface_shutdown,
37 surface_task_result_or_pending,
38};
39#[cfg(test)]
40use summary::http_host_for_bind;
41use summary::{RestTransportKind, log_startup_summary, merge_push_health_snapshots_for_rest};
42
43include!("phase4/preauth.rs");
44
45pub(super) async fn run_phase4(
46 config: &RuntimeConfig,
47 phase1: Phase1Out,
48 bridge: Arc<GatewayBridge>,
49 mut auth_plan: Option<AuthPlan>,
50 phase3: Phase3Out,
51 shutdown_tx: tokio::sync::watch::Sender<bool>,
52 mut shutdown_rx: tokio::sync::watch::Receiver<bool>,
53) -> Result<()> {
54 let Phase1Out {
55 _audit_guard,
56 shared_counters,
57 listen_addr,
58 rest_keys_file,
59 rest_key_store,
60 rest_transport,
61 ws_keys_file,
62 ws_key_store,
63 grpc_keys_file,
64 grpc_key_store,
65 allow_tcp_unauthenticated,
66 } = phase1;
67 let rest_transport_kind = match config.rest_port {
68 None => RestTransportKind::Disabled,
69 Some(_) if rest_transport.is_tls() => RestTransportKind::Tls,
70 Some(_) => RestTransportKind::Plaintext,
71 };
72 let Phase3Out {
73 server,
74 server_config,
75 ws_broadcaster,
76 grpc_broadcaster,
77 } = phase3;
78 let push_sinks: Vec<std::sync::Arc<dyn futu_server::push::ExternalPushSink>> = vec![
79 std::sync::Arc::clone(&ws_broadcaster) as _,
80 std::sync::Arc::clone(&grpc_broadcaster) as _,
81 ];
82 let prepared_push_dispatcher = bridge.create_push_dispatcher(&server, push_sinks);
83 let (listener_events_tx, mut listener_events_rx) =
84 tokio::sync::mpsc::unbounded_channel::<ListenerBindEvent>();
85
86 let mut ws_key_store_holder: Option<std::sync::Arc<futu_auth::KeyStore>> = None;
90 let mut rest_key_store_holder: Option<std::sync::Arc<futu_auth::KeyStore>> = None;
91 let mut grpc_key_store_holder: Option<std::sync::Arc<futu_auth::KeyStore>> = None;
92 let any_keys_configured =
93 rest_keys_file.is_some() || grpc_keys_file.is_some() || ws_keys_file.is_some();
94 warn_legacy_network_exposure(config, any_keys_configured);
95
96 let mut ws_handle = if let Some(ws_port) = config.websocket_port {
98 let ws_addr = format!("{}:{}", config.ip, ws_port);
99 let ws_key_store = ws_key_store.clone();
107 let ws_counters = std::sync::Arc::clone(&shared_counters);
108 ws_key_store_holder = ws_key_store.as_ref().map(std::sync::Arc::clone);
110 let ws_server = WsServer::with_auth(
111 ws_addr.clone(),
112 server_config.clone(),
113 WsServerDeps::new(
114 std::sync::Arc::clone(server.connections()),
115 std::sync::Arc::clone(server.router()),
116 Some(bridge.subscription_runtime().manager()),
117 ),
118 ws_key_store,
119 Some(ws_counters),
120 )
121 .with_server_time_store(bridge.server_clock().anchor_store());
122 tracing::info!(addr = %ws_addr, "starting WebSocket server");
123 let ws_shutdown_rx = shutdown_rx.clone();
124 let ws_listener_events = listener_events_tx.clone();
125 Some(tokio::spawn(async move {
126 ws_server
127 .run_until_shutdown_with_listener_events(ws_shutdown_rx, Some(ws_listener_events))
128 .await
129 }))
130 } else {
131 None
132 };
133
134 let mut rest_handle = if let Some(rest_port) = config.rest_port {
136 let rest_addr = format!("{}:{}", config.ip, rest_port);
137 let router = std::sync::Arc::clone(server.router());
138 let broadcaster = std::sync::Arc::clone(&ws_broadcaster);
139 let rest_key_store = rest_key_store
141 .clone()
142 .unwrap_or_else(|| std::sync::Arc::new(futu_auth::KeyStore::empty()));
143 let rest_transport_name = match rest_transport_kind {
144 RestTransportKind::Tls => "https",
145 RestTransportKind::Plaintext => "http",
146 RestTransportKind::Disabled => unreachable!("REST port is configured"),
147 };
148 tracing::info!(
149 addr = %rest_addr,
150 transport = rest_transport_name,
151 "starting REST API server (WebSocket: /ws)"
152 );
153
154 if rest_key_store.is_configured() {
158 rest_key_store_holder = Some(std::sync::Arc::clone(&rest_key_store));
159 }
160
161 let rest_counters = std::sync::Arc::clone(&shared_counters);
168 let bridge_for_status = std::sync::Arc::clone(&bridge);
172 let admin_status_provider: futu_rest::adapter::AdminStatusProvider =
173 std::sync::Arc::new(move || {
174 serde_json::to_value(bridge_for_status.snapshot_status())
175 .unwrap_or_else(|_| serde_json::json!({"error": "snapshot serialize failed"}))
176 });
177 let rest_shutdown_tx = shutdown_tx.clone();
178 let admin_shutdown_handler: futu_rest::adapter::AdminShutdownHandler =
179 std::sync::Arc::new(move || {
180 rest_shutdown_tx
181 .send(true)
182 .map_err(|e| format!("shutdown receiver dropped: {e}"))
183 });
184 let bridge_for_reload = std::sync::Arc::clone(&bridge);
192 let admin_reload_handler: futu_rest::adapter::AdminReloadHandler =
193 std::sync::Arc::new(move || {
194 let bridge = std::sync::Arc::clone(&bridge_for_reload);
195 Box::pin(async move {
196 serde_json::to_value(bridge.reload())
197 .unwrap_or_else(|_| serde_json::json!({"error": "reload serialize failed"}))
198 })
199 });
200 let bridge_for_push_health = std::sync::Arc::clone(&bridge);
207 let push_health_snapshot_provider: futu_rest::adapter::PushHealthSnapshotProvider =
208 std::sync::Arc::new(move || {
209 let push = serde_json::to_value(
210 bridge_for_push_health
211 .push_runtime()
212 .push_health()
213 .snapshot(),
214 )
215 .unwrap_or_else(
216 |_| serde_json::json!({"error": "push_health snapshot serialize failed"}),
217 );
218 let qot_login = serde_json::to_value(
219 bridge_for_push_health
220 .push_runtime()
221 .qot_login_health()
222 .snapshot(),
223 )
224 .unwrap_or_else(
225 |_| serde_json::json!({"error": "qot_login_health snapshot serialize failed"}),
226 );
227 let backend_connected =
228 bridge_for_push_health.broker_runtime().platform_connected();
229 let stock_list_status = bridge_for_push_health
230 .caches()
231 .static_cache
232 .stock_list_sync_status();
233 let stock_list_static_query_ready = matches!(
234 futu_domain_static_data::stock_list_static_query_gate_action_for_status(
235 &stock_list_status.domain_facts(),
236 ),
237 futu_domain_static_data::StockListStaticQueryGateAction::Allow
238 );
239 merge_push_health_snapshots_for_rest(
240 push,
241 qot_login,
242 backend_connected,
243 stock_list_static_query_ready,
244 )
245 });
246 let bridge_for_card_num = std::sync::Arc::clone(&bridge);
252 let card_num_resolver: futu_rest::adapter::CardNumResolver =
253 std::sync::Arc::new(move |cn: &str| {
254 bridge_for_card_num
255 .caches()
256 .trd_cache
257 .find_acc_ids_by_card_num(cn)
258 });
259 let rest_shutdown_rx = shutdown_rx.clone();
260 let rest_listener_events = listener_events_tx.clone();
261 Some(tokio::spawn(async move {
262 futu_rest::server::start_with_auth_full_admin_until_shutdown_with_transport_and_listener_events(
263 &rest_addr,
264 router,
265 broadcaster,
266 rest_key_store,
267 rest_counters,
268 futu_rest::server::RestAdminHooks {
269 admin_status_provider: Some(admin_status_provider),
270 admin_shutdown_handler: Some(admin_shutdown_handler),
271 admin_reload_handler: Some(admin_reload_handler),
272 push_health_snapshot_provider: Some(push_health_snapshot_provider),
273 card_num_resolver: Some(card_num_resolver),
274 },
275 rest_transport,
276 rest_shutdown_rx,
277 Some(rest_listener_events),
278 )
279 .await
280 .map_err(anyhow::Error::from)
281 }))
282 } else {
283 None
284 };
285
286 let mut grpc_handle = if let Some(grpc_port) = config.grpc_port {
288 let grpc_addr = format!("{}:{}", config.ip, grpc_port);
289 let router = std::sync::Arc::clone(server.router());
290 let broadcaster = std::sync::Arc::clone(&grpc_broadcaster);
291 let grpc_key_store = grpc_key_store
293 .clone()
294 .unwrap_or_else(|| std::sync::Arc::new(futu_auth::KeyStore::empty()));
295 tracing::info!(addr = %grpc_addr, "starting gRPC server (SubscribePush: streaming)");
296
297 if grpc_key_store.is_configured() {
299 grpc_key_store_holder = Some(std::sync::Arc::clone(&grpc_key_store));
300 }
301
302 let grpc_counters = std::sync::Arc::clone(&shared_counters);
307 let grpc_shutdown_rx = shutdown_rx.clone();
308 let grpc_listener_events = listener_events_tx.clone();
309 Some(tokio::spawn(async move {
310 let result = futu_grpc::server::start_with_auth_until_shutdown_with_listener_events(
311 &grpc_addr,
312 router,
313 broadcaster,
314 grpc_key_store,
315 grpc_counters,
316 grpc_shutdown_rx,
317 Some(grpc_listener_events),
318 )
319 .await;
320 wrap_grpc_server_result(result)
321 }))
322 } else {
323 None
324 };
325
326 let telnet_addr = telnet_bind_addr(&config.telnet_ip, config.telnet_port);
328 let mut telnet_handle = if let Some(telnet_addr) = telnet_addr {
329 let bridge_for_relogin = std::sync::Arc::clone(&bridge);
333 let relogin_fn: futu_server::telnet::ReloginFn = std::sync::Arc::new(move || {
334 tracing::warn!(
335 "v1.4.97 P1-D-F: telnet relogin clearing login_cache; \
336 next P1-D tick will trigger central auth-session refresh"
337 );
338 bridge_for_relogin.caches().login_cache.clear();
339 });
340 let telnet_server = futu_server::telnet::TelnetServer::new(
341 telnet_addr.clone(),
342 std::sync::Arc::clone(server.connections()),
343 Some(bridge.subscription_runtime().manager()),
344 Some(std::sync::Arc::clone(server.metrics())),
345 shutdown_tx.clone(),
346 )
347 .with_relogin_fn(relogin_fn);
348 tracing::info!(addr = %telnet_addr, "starting Telnet server");
349 let telnet_shutdown_rx = shutdown_rx.clone();
350 let telnet_listener_events = listener_events_tx.clone();
351 Some(tokio::spawn(async move {
352 telnet_server
353 .run_until_shutdown_with_listener_events(
354 telnet_shutdown_rx,
355 Some(telnet_listener_events),
356 )
357 .await
358 }))
359 } else {
360 None
361 };
362
363 let card_num_reload_and_expand_fn: std::sync::Arc<dyn Fn(bool) + Send + Sync> = {
381 let bridge_for_expand = std::sync::Arc::clone(&bridge);
382 let ws_ks = ws_key_store_holder.clone();
383 let rest_ks = rest_key_store_holder.clone();
384 let grpc_ks = grpc_key_store_holder.clone();
385 std::sync::Arc::new(move |do_reload: bool| {
386 if do_reload {
388 for (ks_name, ks_opt) in [("ws", &ws_ks), ("rest", &rest_ks), ("grpc", &grpc_ks)] {
389 let Some(ks) = ks_opt.as_ref() else { continue };
390 match ks.reload() {
391 Ok(()) => tracing::warn!(
392 ks = ks_name,
393 keys_loaded = ks.len(),
394 "v1.4.103 F3.1: keys reloaded on SIGHUP (before card_num expand)"
395 ),
396 Err(e) => tracing::error!(
397 ks = ks_name,
398 error = %e,
399 "v1.4.103 F3.1: keys reload failed (skipping expand for this store)"
400 ),
401 }
402 }
403 }
404 let trd_cache = std::sync::Arc::clone(&bridge_for_expand.caches().trd_cache);
406 let resolver = {
407 let cache_clone = std::sync::Arc::clone(&trd_cache);
408 move |cn: &str| cache_clone.find_acc_ids_by_card_num(cn)
409 };
410 for (ks_name, ks_opt) in [("ws", &ws_ks), ("rest", &rest_ks), ("grpc", &grpc_ks)] {
411 let Some(ks) = ks_opt.as_ref() else { continue };
412 let (resolved, unresolved, ambiguous) = ks.expand_allowed_card_nums(
413 &resolver,
414 |key_id, cn| {
415 tracing::warn!(
416 key_id = %key_id,
417 card_num = %cn,
418 "v1.4.103 B10/F1 fail-closed: card_num not found in trd_cache; \
419 writing sentinel acc_id=0 to enforce restrictive denylist \
420 (limits.contains check 永远 false → reject 真账户)"
421 );
422 },
423 |key_id, cn, candidates| {
424 tracing::warn!(
425 key_id = %key_id,
426 card_num = %cn,
427 candidates = ?candidates,
428 "v1.4.103 B10/F1 fail-closed: ambiguous card_num suffix \
429 matched multiple accounts (skipped, write 完整 16 位 / specific 4 位)"
430 );
431 },
432 );
433 tracing::info!(
434 ks = ks_name,
435 resolved,
436 unresolved,
437 ambiguous,
438 "v1.4.103 B10: expanded allowed_card_nums into allowed_acc_ids"
439 );
440 }
441 })
442 };
443 let card_num_expand_fn: std::sync::Arc<dyn Fn() + Send + Sync> = {
445 let inner = std::sync::Arc::clone(&card_num_reload_and_expand_fn);
446 std::sync::Arc::new(move || (inner)(false))
447 };
448
449 (card_num_expand_fn)();
451
452 let card_num_retry_handle: BackgroundJoinHandle = {
454 let card_num_expand_fn_loop = std::sync::Arc::clone(&card_num_expand_fn);
455 let bridge_for_check = std::sync::Arc::clone(&bridge);
456 let mut card_num_retry_shutdown_rx = shutdown_rx.clone();
457 tokio::spawn(async move {
458 let trd_cache = std::sync::Arc::clone(&bridge_for_check.caches().trd_cache);
459 let mut attempts = 0u32;
460 let max_attempts = 6u32; loop {
462 tokio::select! {
463 changed = card_num_retry_shutdown_rx.changed() => {
464 if changed.is_err() || *card_num_retry_shutdown_rx.borrow() {
465 tracing::debug!(
466 "v1.4.111: card_num retry loop received shutdown signal"
467 );
468 return;
469 }
470 }
471 _ = tokio::time::sleep(std::time::Duration::from_secs(10)) => {}
472 }
473 attempts += 1;
474 let accounts = trd_cache.get_accounts();
475 if accounts.is_empty() {
476 if attempts >= max_attempts {
477 tracing::warn!(
478 "v1.4.103 B10: trd_cache 仍空 (after {max_attempts} × 10s); \
479 受限 key 仍走 fail-closed sentinel reject 直到 SIGHUP / cache 加载."
480 );
481 return;
482 }
483 continue;
484 }
485 (card_num_expand_fn_loop)();
486 return;
487 }
488 })
489 };
490
491 #[cfg(unix)]
500 let sighup_handle: Option<BackgroundJoinHandle> = {
501 let unified_sighup_fn = std::sync::Arc::clone(&card_num_reload_and_expand_fn);
502 let mut sighup_shutdown_rx = shutdown_rx.clone();
503 Some(tokio::spawn(async move {
504 use tokio::signal::unix::{SignalKind, signal};
505 let mut sig = match signal(SignalKind::hangup()) {
506 Ok(s) => s,
507 Err(e) => {
508 tracing::error!(error = %e, "SIGHUP install failed (unified reload+expand)");
509 return;
510 }
511 };
512 tracing::info!(
513 "v1.4.103 F3.1: unified SIGHUP handler installed (reload all keys + expand card_num)"
514 );
515 loop {
516 tokio::select! {
517 signal = sig.recv() => {
518 if signal.is_none() {
519 return;
520 }
521 tracing::info!(
522 "v1.4.103 F3.1: SIGHUP received — running reload_all_stores + \
523 expand_allowed_card_nums (single ordered op, no race)"
524 );
525 (unified_sighup_fn)(true); }
527 changed = sighup_shutdown_rx.changed() => {
528 if changed.is_err() || *sighup_shutdown_rx.borrow() {
529 tracing::debug!(
530 "v1.4.111: unified SIGHUP handler received shutdown signal"
531 );
532 return;
533 }
534 }
535 }
536 }
537 }))
538 };
539 #[cfg(not(unix))]
540 let sighup_handle: Option<BackgroundJoinHandle> = None;
541
542 let tcp_disabled = !super::phase1::ftapi_listener_enabled(config);
553
554 if tcp_disabled {
555 tracing::warn!(
556 listen_addr = %listen_addr,
557 "v1.4.104 external report S-001 (P0) fix: TCP listener (port {}) NOT started — \
558 keys file configured but --allow-tcp-unauthenticated not set. \
559 native TCP FTAPI protocol has no Bearer field, cannot enforce \
560 caller-specific scope check; defaulting to fail-closed (skip TCP). \
561 Use REST/gRPC/WS endpoints for authenticated access. \
562 To restore TCP (legacy Python SDK clients) add --allow-tcp-unauthenticated, \
563 but be aware that port {} will accept ANY local connection without \
564 scope check (跨账户 leak risk).",
565 config.port,
566 config.port,
567 );
568 eprintln!(
569 "⚠️ TCP listener (port {}) DISABLED (v1.4.104 external report S-001 fix): \
570 keys file configured + no --allow-tcp-unauthenticated. \
571 Pass --allow-tcp-unauthenticated to restore (with security warning).",
572 config.port,
573 );
574 } else if any_keys_configured && allow_tcp_unauthenticated {
575 tracing::warn!(
576 listen_addr = %listen_addr,
577 "⚠️ v1.4.104: TCP listener running WITHOUT scope check despite keys configured \
578 (--allow-tcp-unauthenticated set). Port {} accepts ANY local connection — \
579 跨账户 leak risk. Use REST/gRPC/WS for authenticated clients; reserve \
580 TCP only for legacy Python SDK / C++ OpenD where Bearer not feasible.",
581 config.port,
582 );
583 eprintln!(
584 "⚠️ TCP port {} ACCEPTS UNAUTHENTICATED connections (--allow-tcp-unauthenticated). \
585 受限 keys 不在该 surface 强制. 推荐改用 REST/gRPC/WS.",
586 config.port,
587 );
588 }
589
590 let tcp_shutdown_rx = shutdown_rx.clone();
591 let tcp_handle: Option<SurfaceJoinHandle> = if tcp_disabled {
592 None
593 } else {
594 let tcp_listener_events = listener_events_tx.clone();
595 Some(tokio::spawn(async move {
596 server
597 .run_until_shutdown_with_listener_events(tcp_shutdown_rx, Some(tcp_listener_events))
598 .await
599 }))
600 };
601
602 let enabled_listeners = EnabledListenerSurfaces {
603 ftapi: !tcp_disabled,
604 websocket: config.websocket_port.is_some(),
605 rest: config.rest_port.is_some(),
606 grpc: config.grpc_port.is_some(),
607 telnet: config.telnet_port.is_some(),
608 };
609 if let Some(plan) = auth_plan.as_mut() {
610 let loopback = bind_includes_loopback(&config.ip);
611 let available = external_verification_owner_available(VerificationOwnerFacts {
612 ftapi_loopback_legacy: enabled_listeners.ftapi && loopback,
613 rest_loopback_legacy: enabled_listeners.rest
614 && loopback
615 && rest_key_store
616 .as_ref()
617 .is_none_or(|store| !store.is_configured()),
618 ws_loopback_legacy: enabled_listeners.websocket
619 && loopback
620 && ws_key_store
621 .as_ref()
622 .is_none_or(|store| !store.is_configured()),
623 grpc_loopback_legacy: enabled_listeners.grpc
624 && loopback
625 && grpc_key_store
626 .as_ref()
627 .is_none_or(|store| !store.is_configured()),
628 rest_auth_setup: enabled_listeners.rest
629 && rest_key_store
630 .as_ref()
631 .is_some_and(|store| store.has_usable_scope(futu_auth::Scope::AuthSetup)),
632 ws_auth_setup: enabled_listeners.websocket
633 && ws_key_store
634 .as_ref()
635 .is_some_and(|store| store.has_usable_scope(futu_auth::Scope::AuthSetup)),
636 grpc_auth_setup: enabled_listeners.grpc
637 && grpc_key_store
638 .as_ref()
639 .is_some_and(|store| store.has_usable_scope(futu_auth::Scope::AuthSetup)),
640 });
641 plan.set_external_verification_available(available);
642 tracing::info!(
643 external_verification_owner_available = available,
644 "computed pre-login verification ownership from bound listeners"
645 );
646 }
647 let mut surface_tasks = SurfaceTasks {
648 ftapi: tcp_handle,
649 websocket: ws_handle.take(),
650 rest: rest_handle.take(),
651 grpc: grpc_handle.take(),
652 telnet: telnet_handle.take(),
653 };
654 drop(listener_events_tx);
655 let listener_readiness = await_listeners_opened(
656 &mut listener_events_rx,
657 enabled_listeners,
658 &mut surface_tasks,
659 &mut shutdown_rx,
660 )
661 .await;
662
663 let phase4_result: anyhow::Result<()> = match listener_readiness {
666 Err(error) => Err(error),
667 Ok(readiness) => {
668 match activate_restricted_listeners_before_auth(
674 readiness,
675 enabled_listeners,
676 &mut shutdown_rx,
677 )
678 .await
679 {
680 Err(error) => Err(error),
681 Ok(listener_activation) => {
682 let auth_error = if let Some(auth_plan) = auth_plan {
683 match execute_auth_plan_until_shutdown(
684 Arc::clone(&bridge),
685 auth_plan,
686 &mut shutdown_rx,
687 &mut surface_tasks,
688 )
689 .await
690 {
691 Ok(push_rx) => {
692 bridge.start_prepared_push_dispatcher(
693 prepared_push_dispatcher,
694 push_rx,
695 );
696 tracing::info!(
697 "authentication promoted atomically; push dispatcher started"
698 );
699 None
700 }
701 Err(error) => {
702 let startup = bridge.startup_readiness().snapshot();
703 if startup.state
704 == futu_server::identity::StartupState::Authenticating
705 {
706 let _ = bridge.startup_readiness().transition(
707 startup.generation,
708 futu_server::identity::StartupEvent::AuthFailed,
709 );
710 }
711 crate::hints::print_auth_error_hint(&error, config);
712 Some(startup_auth_failure_after_restricted_bind(error))
713 }
714 }
715 } else {
716 None
717 };
718 if let Some(error) = auth_error {
719 Err(error)
720 } else {
721 match listener_activation
722 .activate_ready_only_and_confirm(&mut shutdown_rx, |marker| *marker)
723 .await
724 {
725 Err(error) => Err(error),
726 Ok(listeners_marker) => {
727 tracing::info!(marker = %listeners_marker, "authentication complete on active listeners");
728 log_startup_summary(
729 config,
730 &listen_addr,
731 &bridge,
732 rest_transport_kind,
733 );
734 tracing::info!(
735 "gateway ready, accepting connections on {listen_addr}"
736 );
737 tracing::info!("press Ctrl+C to exit");
738 tokio::select! {
739 result = surface_task_result_or_pending(&mut surface_tasks.ftapi) => {
740 surface_tasks.ftapi = None;
741 log_surface_task_result("API server", result)
742 }
743 result = surface_task_result_or_pending(&mut surface_tasks.websocket) => {
744 surface_tasks.websocket = None;
745 log_surface_task_result("WebSocket", result)
746 }
747 result = surface_task_result_or_pending(&mut surface_tasks.rest) => {
748 surface_tasks.rest = None;
749 log_surface_task_result("REST API", result)
750 }
751 result = surface_task_result_or_pending(&mut surface_tasks.grpc) => {
752 surface_tasks.grpc = None;
753 log_surface_task_result("gRPC", result)
754 }
755 result = surface_task_result_or_pending(&mut surface_tasks.telnet) => {
756 surface_tasks.telnet = None;
757 log_surface_task_result("Telnet", result)
758 }
759 _ = tokio::signal::ctrl_c() => {
760 tracing::info!("received Ctrl+C, shutting down gracefully...");
761 Ok(())
762 }
763 _ = wait_for_sigterm() => {
764 tracing::info!("received SIGTERM, shutting down gracefully...");
765 Ok(())
766 }
767 _ = async {
768 while shutdown_rx.changed().await.is_ok() {
769 if *shutdown_rx.borrow() {
770 break;
771 }
772 }
773 } => {
774 tracing::info!("shutdown requested via telnet");
775 Ok(())
776 }
777 }
778 }
779 }
780 }
781 }
782 }
783 }
784 };
785
786 transition_startup_to_shutdown(bridge.startup_readiness());
791 if tokio::time::timeout(
792 SHUTDOWN_GRACE_PERIOD,
793 bridge
794 .startup_readiness()
795 .await_pending_init_connects_drained(),
796 )
797 .await
798 .is_err()
799 {
800 tracing::warn!(
801 pending_init_connects = bridge.startup_readiness().pending_init_connect_count(),
802 "timed out flushing deferred InitConnect completions before surface shutdown"
803 );
804 }
805
806 request_surface_shutdown(&shutdown_tx, "phase4 exit");
808 if !bridge
809 .qot_runtime()
810 .indicator_calc_runtime()
811 .shutdown_and_wait(SHUTDOWN_GRACE_PERIOD)
812 .await
813 {
814 tracing::warn!("timed out cancelling indicator calculations during shutdown");
815 }
816 await_background_shutdown(
817 "card_num retry",
818 card_num_retry_handle,
819 SHUTDOWN_GRACE_PERIOD,
820 )
821 .await;
822 if let Some(handle) = sighup_handle {
823 await_background_shutdown("SIGHUP reload", handle, SHUTDOWN_GRACE_PERIOD).await;
824 }
825 if let Some(handle) = surface_tasks.ftapi {
826 await_surface_shutdown("API server", handle, SHUTDOWN_GRACE_PERIOD).await;
827 }
828 if let Some(handle) = surface_tasks.websocket {
829 await_surface_shutdown("WebSocket", handle, SHUTDOWN_GRACE_PERIOD).await;
830 }
831 if let Some(handle) = surface_tasks.rest {
832 await_surface_shutdown("REST API", handle, SHUTDOWN_GRACE_PERIOD).await;
833 }
834 if let Some(handle) = surface_tasks.grpc {
835 await_surface_shutdown("gRPC", handle, SHUTDOWN_GRACE_PERIOD).await;
836 }
837 if let Some(handle) = surface_tasks.telnet {
838 await_surface_shutdown("Telnet", handle, SHUTDOWN_GRACE_PERIOD).await;
839 }
840
841 phase4_result?;
842 tracing::info!("gateway stopped");
843 Ok(())
844}