MCPcopy Create free account
hub / github.com/Dstack-TEE/dstack / start_bootnode_discovery_task

Function start_bootnode_discovery_task

gateway/src/main_service.rs:637–676  ·  view source on GitHub ↗

Periodically retry bootnode peer discovery if no peers are available

(proxy: Proxy)

Source from the content-addressed store, hash-verified

635
636/// Periodically retry bootnode peer discovery if no peers are available
637fn start_bootnode_discovery_task(proxy: Proxy) {
638 if !proxy.config.sync.enabled || proxy.config.sync.bootnode.is_empty() {
639 return;
640 }
641
642 let bootnode = proxy.config.sync.bootnode.clone();
643 let node_id = proxy.config.sync.node_id;
644 let kv_store = proxy.kv_store.clone();
645 let https_config = match &proxy.https_config {
646 Some(config) => config.clone(),
647 None => return,
648 };
649
650 tokio::spawn(async move {
651 let mut interval = tokio::time::interval(Duration::from_secs(10));
652 loop {
653 interval.tick().await;
654 // Check if we already have peers
655 let n_peers = kv_store
656 .load_all_node_statuses()
657 .keys()
658 .filter(|&id| *id != node_id)
659 .count();
660 if n_peers > 0 {
661 info!("bootnode peer discovery finished, {n_peers} peers found");
662 break;
663 }
664 // Try to fetch peers from bootnode
665 debug!("retrying bootnode peer discovery...");
666 if let Err(err) =
667 fetch_peers_from_bootnode(&bootnode, &kv_store, node_id, &https_config).await
668 {
669 warn!("bootnode discovery retry failed: {err:?}");
670 } else {
671 info!("bootnode peer discovery succeeded");
672 }
673 }
674 });
675 info!("Bootnode discovery task started (will retry every 10s if no peers)");
676}
677
678async fn start_wavekv_sync_task(proxy: Proxy, wavekv_sync: Arc<WaveKvSyncService>) {
679 if !proxy.config.sync.enabled {

Callers 1

start_bg_tasksMethod · 0.85

Calls 4

cloneMethod · 0.80
is_emptyMethod · 0.45

Tested by

no test coverage detected