Skip to main content

toole_core/
receiver.rs

1use crate::transfer::{
2    handle_incoming_connection, make_server_endpoint, DecisionBoard,
3};
4use crate::{ToolError, TransferRegistry, UI};
5use std::path::PathBuf;
6use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
7use std::sync::{Arc, Mutex};
8use std::time::Duration;
9
10const PORT: u16 = 58200;
11
12pub async fn start_receiver(
13    ui: Arc<dyn UI>,
14    dest_dir: PathBuf,
15    stop: Arc<AtomicBool>,
16    decisions: Arc<DecisionBoard>,
17    registry: Arc<dyn TransferRegistry>,
18) -> Result<(), ToolError> {
19    let endpoint = make_server_endpoint().await?;
20    ui.log(&format!("Recepteur en ecoute sur le port {PORT}"));
21
22    loop {
23        if stop.load(Ordering::Relaxed) {
24            break;
25        }
26
27        let incoming = tokio::select! {
28            conn = endpoint.accept() => conn,
29            _ = tokio::time::sleep(Duration::from_millis(500)) => continue,
30        };
31
32        let Some(connecting) = incoming else {
33            break;
34        };
35
36        let dest_dir = dest_dir.clone();
37        let ui = ui.clone();
38        let decisions = decisions.clone();
39        let registry = registry.clone();
40
41        tokio::spawn(async move {
42            match connecting.await {
43                Ok(connection) => {
44                    let peer = connection.remote_address().ip().to_string();
45                    ui.log(&format!(
46                        "Connexion entrante depuis {:?}",
47                        connection.remote_address()
48                    ));
49
50                    // l'id du transfert est donne par l'emetteur (lu dans le metadata)
51                    let transfer_id: Arc<Mutex<Option<String>>> = Arc::new(Mutex::new(None));
52                    let total = Arc::new(AtomicU64::new(0));
53                    let files = Arc::new(Mutex::new(Vec::new()));
54                    let bytes = Arc::new(AtomicU64::new(0));
55
56                    // chaque connexion a son propre drapeau d'arrêt : annuler un
57                    // transfert ne doit pas arrêter le récepteur global
58                    let conn_stop = Arc::new(AtomicBool::new(false));
59
60                    let res = handle_incoming_connection(
61                        connection,
62                        dest_dir,
63                        conn_stop,
64                        files.clone(),
65                        bytes.clone(),
66                        ui.clone(),
67                        transfer_id.clone(),
68                        total.clone(),
69                        decisions,
70                        registry.clone(),
71                    )
72                    .await;
73
74                    // je désenregistre le transfert quel que soit le résultat
75                    let tid = transfer_id
76                        .lock()
77                        .unwrap_or_else(|e| e.into_inner())
78                        .clone()
79                        .unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
80                    registry.unregister(&tid);
81
82                    match res {
83                        Err(ToolError::Cancelled) | Err(ToolError::RemoteCancel) => {
84                            ui.transfert_cancel(&tid);
85                        }
86                        Err(ToolError::Refused) => {
87                            ui.transfert_refused(&tid);
88                        }
89                        Err(e) => {
90                            eprintln!("Erreur connexion receveur: {e}");
91                            let err: ToolError = crate::transfer::io_err(format!("reception: {e}"));
92                            ui.transfert_error(&tid, &err);
93                        }
94                        Ok(()) => {
95                            let received: Vec<String> = files.lock().unwrap_or_else(|e| e.into_inner()).clone();
96                            let total = bytes.load(Ordering::Relaxed);
97                            ui.transfert_received(&tid, &peer, total, received);
98                        }
99                    }
100                }
101                Err(e) => eprintln!("Handshake QUIC echoue: {e}"),
102            }
103        });
104    }
105
106    endpoint.close(0u32.into(), b"arret");
107    Ok(())
108}