1use crate::transfer::{
2 collect_entries, io_err, make_client_endpoint, send_entry, write_json_line, ACK, BatchHeader,
3 CLOSE_CANCEL, CLOSE_OK, DECISION_TIMEOUT, REFUSE,
4};
5use crate::{ToolError, UI};
6use std::net::SocketAddr;
7use std::path::PathBuf;
8use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
9use std::sync::Arc;
10use std::time::{Duration, Instant};
11use tokio::fs;
12use tokio::sync::Semaphore;
13
14pub async fn start_sender(
15 ui: Arc<dyn UI>,
16 transfer_id: String,
17 paths: Vec<PathBuf>,
18 peer_addr: SocketAddr,
19 peer_id: String,
20 stop: Arc<AtomicBool>,
21) -> Result<(), ToolError> {
22 let expected = crate::file_certif::pin_for(&peer_id);
25 let endpoint = make_client_endpoint(expected.as_deref())?;
26 let connecting = endpoint.connect(peer_addr, "localhost").map_err(io_err)?;
27 let connection = connecting.await?;
28 ui.log(&format!("Connecte a {peer_addr}"));
29
30 if let Some(fp) = crate::file_certif::peer_fingerprint(&connection) {
34 if crate::file_certif::pin_for(&peer_id) != Some(fp.clone()) {
35 match crate::file_certif::save_pin(&peer_id, &fp) {
36 Ok(()) => ui.log(&format!("Empreinte du pair {peer_id} epinglee ({fp})")),
37 Err(e) => ui.log(&format!("epingle impossible: {e}")),
38 }
39 }
40 }
41
42 let mut entries = Vec::new();
43 for path in &paths {
44 collect_entries(path, path, &mut entries).await?;
45 }
46
47 let mut total_bytes: u64 = 0;
48 for (abs_path, _rel_path, is_dir) in &entries {
49 if !is_dir {
50 total_bytes += fs::metadata(abs_path).await?.len();
51 }
52 }
53
54 let files: Vec<String> = entries
55 .iter()
56 .map(|(_abs, rel, _is_dir)| rel.clone())
57 .collect();
58
59 if !entries.is_empty() {
63 let (mut header_send, mut header_recv) = connection.open_bi().await.map_err(io_err)?;
64 let header = BatchHeader {
65 transfer_id: transfer_id.clone(),
66 total_bytes,
67 sender: crate::utils::device_id(),
68 files: files.clone(),
69 };
70 write_json_line(&mut header_send, &header).await?;
71
72 let mut decision = [0u8; 1];
75 let deadline = Instant::now() + DECISION_TIMEOUT;
76 let decision_result = loop {
77 if stop.load(Ordering::Relaxed) {
78 connection.close(CLOSE_CANCEL.into(), b"annulation utilisateur");
79 endpoint.wait_idle().await;
80 ui.transfert_cancel(&transfer_id);
81 return Ok(());
82 }
83 if Instant::now() >= deadline {
84 break Err(io_err("le destinataire n'a pas repondu"));
85 }
86 match tokio::time::timeout(
87 Duration::from_millis(100),
88 header_recv.read_exact(&mut decision),
89 )
90 .await
91 {
92 Ok(Ok(())) => break Ok(()),
93 Ok(Err(e)) => break Err(crate::transfer::quinn_to_err(e)),
94 Err(_) => continue,
95 }
96 };
97
98 match decision_result {
99 Ok(()) if decision[0] == ACK => {
100 header_send.finish()?;
101 }
102 Ok(()) if decision[0] == REFUSE => {
103 connection.close(CLOSE_CANCEL.into(), b"transfert refuse");
104 endpoint.wait_idle().await;
105 ui.transfert_refused(&transfer_id);
106 return Ok(());
107 }
108 Ok(()) => {
109 connection.close(CLOSE_CANCEL.into(), b"reponse invalide");
110 endpoint.wait_idle().await;
111 ui.transfert_error(
112 &transfer_id,
113 &io_err("reponse invalide du destinataire"),
114 );
115 return Ok(());
116 }
117 Err(ToolError::RemoteCancel) => {
118 connection.close(CLOSE_CANCEL.into(), b"annule par le destinataire");
120 endpoint.wait_idle().await;
121 ui.transfert_cancel(&transfer_id);
122 return Ok(());
123 }
124 Err(e) => {
125 connection.close(CLOSE_CANCEL.into(), b"pas de reponse");
126 endpoint.wait_idle().await;
127 ui.transfert_error(&transfer_id, &e);
128 return Ok(());
129 }
130 }
131 }
132
133 ui.show_progress_bar(&transfer_id);
135 let bytes_sent_counter = Arc::new(AtomicU64::new(0));
136
137 let semaphore = Arc::new(Semaphore::new(2));
139
140 let mut handles = Vec::new();
141 for (abs_path, rel_path, is_dir) in entries {
142 if stop.load(Ordering::Relaxed) {
143 break;
144 }
145 let connection = connection.clone();
146 let stop = stop.clone();
147 let ui = ui.clone();
148 let transfer_id = transfer_id.clone();
149 let bytes_sent_counter = bytes_sent_counter.clone();
150 let permit = semaphore.clone();
151
152 handles.push(tokio::spawn(async move {
153 let _guard = permit.acquire().await;
154 send_entry(
155 connection,
156 abs_path,
157 rel_path,
158 is_dir,
159 stop,
160 ui,
161 transfer_id,
162 total_bytes,
163 bytes_sent_counter,
164 )
165 .await
166 }));
167 }
168
169 let mut had_error = false;
170 let mut remote_cancelled = false;
171 for handle in handles {
172 if let Err(e) = handle.await.map_err(io_err)? {
173 eprintln!("Erreur d'envoi: {e}");
174 if matches!(e, ToolError::RemoteCancel) {
175 remote_cancelled = true;
176 } else {
177 had_error = true;
178 }
179 }
180 }
181
182 if stop.load(Ordering::Relaxed) || remote_cancelled {
183 connection.close(CLOSE_CANCEL.into(), b"annulation");
184 endpoint.wait_idle().await;
185 ui.transfert_cancel(&transfer_id);
186 } else if had_error {
187 connection.close(CLOSE_OK.into(), b"erreur");
193 endpoint.wait_idle().await;
194 ui.transfert_error(&transfer_id, &io_err("un ou plusieurs fichiers ont echoue"));
195 } else {
196 connection.close(CLOSE_OK.into(), b"transfert termine");
197 endpoint.wait_idle().await;
198 ui.transfert_completed(&transfer_id);
199 }
200
201 Ok(())
202}