Skip to content

Commit ffa1fb8

Browse files
committed
feat: imp echo debug
1 parent 091dc8d commit ffa1fb8

3 files changed

Lines changed: 46 additions & 20 deletions

File tree

Cargo.toml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@ path = "src/lib.rs"
1818
[profile.release]
1919
opt-level = 3
2020
lto = "fat"
21-
codegen-units = 1
21+
codegen-units = 16
2222
panic = "abort"
2323
strip = true
2424
incremental = false

README.md

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -74,7 +74,7 @@ Enter file key: ronak
7474
Overwrite? [n/0]: 2
7575
Starting upload: dummy.bin (578.9375 mb)
7676
File hash: 4ea1b5d551d3876f74b6634c4dde1611a8000d798044268ea103736221e7378e
77-
File ID: 3d5675b781c143e881cf37ec81587462 - Network took: 0.3446355
77+
File ID: 3d5675b781c143e881cf37ec81587462 - Network_time: 0.3446355
7878
````
7979

8080
```shell
@@ -159,7 +159,7 @@ TODO
159159
- Too many repetitive code, need to refactor and clean up the codebase
160160
- Still some buffering issues, data gets stalls, does not flush properly [DONE]
161161
- Multi-port support for better concurrency
162-
- rsync support (rolling hashing, delta transfers, etc.) CDC `LAYERING like docker`
162+
- rsync support (rolling hashing, delta transfers, etc.) CDC `LAYERING like docker` [Partially DONE]
163163
- Serialized headers, rm fragile parsing [DONE]
164164
- Add proper user-space (multiple users) [DONE]
165165
- DO some CAS magic for better storage efficiency and deduplication

src/protocol_v1.rs

Lines changed: 43 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -78,7 +78,7 @@ pub async fn start_tcp_server(port: u16, storage_path: Arc<PathBuf>) -> Result<(
7878
/// Handle a single client connection, read command and dispatch to appropriate handler
7979
#[inline]
8080
async fn handle_connection(mut stream: TcpStream, storage_path: &Path) -> Result<()> {
81-
stream.set_nodelay(true).ok();
81+
stream.set_nodelay(true)?;
8282
let (reader, writer) = stream.split();
8383
let mut reader = BufReader::with_capacity(NETWORK_READ_BUFFER, reader);
8484
let mut writer = BufWriter::with_capacity(NETWORK_WRITE_BUFFER, writer);
@@ -601,7 +601,7 @@ async fn handle_echo_debug<R: AsyncReadExt + Unpin, W: AsyncWriteExt + Unpin>(
601601
let payload = match read_encrypt_data_into(reader, session_key, mem_pool).await {
602602
Ok(p) => p,
603603
Err(e) => {
604-
// eof is fine here
604+
// eof & disconnect is fine here
605605
if [
606606
"eof",
607607
"early eof",
@@ -1197,7 +1197,7 @@ async fn client_handshake_helper(
11971197
signing_key: &SigningKey,
11981198
) -> Result<([u8; 32], TcpStream)> {
11991199
let mut stream = TcpStream::connect(format!("{}:{}", host, port)).await?;
1200-
stream.set_nodelay(true).ok();
1200+
stream.set_nodelay(true)?;
12011201

12021202
let server_ip = stream
12031203
.peer_addr()
@@ -1363,17 +1363,26 @@ pub async fn client_echo_debug(
13631363
let mut reader = BufReader::with_capacity(NETWORK_READ_BUFFER, reader);
13641364
let mut writer = BufWriter::with_capacity(NETWORK_WRITE_BUFFER, writer);
13651365

1366-
let mut thread_rng = rng();
1367-
let dur = Duration::from_secs(1);
1368-
13691366
let request = Command::Echo.serialize()?;
13701367
write_encrypt_frame(&mut writer, &request, &session_key, mem_pool).await?;
13711368

13721369
//TODO; read any err or rej
13731370

1374-
println!("Session Key: {}", encode(session_key).green());
1371+
let pb = ProgressBar::new_spinner();
1372+
pb.set_style(ProgressStyle::with_template("{msg}")?);
13751373

1374+
println!("Session Key: {}\n", encode(session_key).green());
1375+
1376+
let mut thread_rng = rng();
1377+
let dur = Duration::from_millis(250);
1378+
const MAX_RECORD: usize = 128;
1379+
let mut record_count = 0;
1380+
let mut payload_acc = 0;
13761381
let mut rand_payload_buf = vec![0u8; 1024 * 1024];
1382+
let mut record_timestamp = [0i64; MAX_RECORD];
1383+
let mut record_rtt = [0u128; MAX_RECORD];
1384+
1385+
// TODO; add some guardrails
13771386
loop {
13781387
// randomize
13791388
let rand_len = rng().random_range(1024 * 1024..=14 * 1024 * 1024);
@@ -1392,8 +1401,25 @@ pub async fn client_echo_debug(
13921401
.context("Failed to read echo response")?;
13931402

13941403
let echo = EchoDebugHeader::deserialize(response)?;
1404+
13951405
let rtt = send_time.elapsed();
13961406
let client_timestamp = chrono::Utc::now().timestamp_millis();
1407+
let timestamp_delta = client_timestamp - echo.timestamp_ms;
1408+
let payload_mb = echo.payload_len / (1024 * 1024);
1409+
1410+
record_timestamp[record_count % MAX_RECORD] = timestamp_delta;
1411+
record_rtt[record_count % MAX_RECORD] = rtt.as_millis();
1412+
1413+
let avg_timestamp = record_timestamp
1414+
.iter()
1415+
.take((record_count + 1).min(MAX_RECORD))
1416+
.sum::<i64>()
1417+
/ (record_count + 1).min(MAX_RECORD) as i64;
1418+
let avg_rtt = record_rtt
1419+
.iter()
1420+
.take((record_count + 1).min(MAX_RECORD))
1421+
.sum::<u128>()
1422+
/ (record_count + 1).min(MAX_RECORD) as u128;
13971423

13981424
if compute_payload_hash != encode(echo.payload_hash) {
13991425
eprintln!(
@@ -1403,15 +1429,15 @@ pub async fn client_echo_debug(
14031429
);
14041430
}
14051431

1406-
println!("---------------------------");
1407-
println!("Payload Size: {} bytes", echo.payload_len);
1408-
println!("Payload SHA256: {}", encode(echo.payload_hash));
1409-
println!(
1410-
"Timestamp Gap: {} ms",
1411-
client_timestamp - echo.timestamp_ms
1412-
);
1413-
println!("Total Round-Trip: {:?}", rtt);
1432+
pb.set_message(format!(
1433+
"Sample: {}\nPayload: {} MB (total {})\nSha256: {}\nTimestamp delta: {} ms (avg {})\nRTT: {} ms (avg {})",
1434+
record_count, payload_mb, payload_acc,
1435+
encode(echo.payload_hash), timestamp_delta,
1436+
avg_timestamp, rtt.as_millis(), avg_rtt
1437+
));
14141438

1439+
record_count += 1;
1440+
payload_acc += echo.payload_len / (1024 * 1024);
14151441
mem_pool.clear();
14161442
tokio::time::sleep(dur).await;
14171443
}
@@ -1599,7 +1625,7 @@ pub async fn upload_client(
15991625
stream.shutdown().await?;
16001626

16011627
println!(
1602-
"File ID: {} - Network took: {}",
1628+
"File ID: {} - Network_time: {}s",
16031629
rsp.file_id, rsp.network_time
16041630
);
16051631

@@ -1706,7 +1732,7 @@ pub async fn download_client(
17061732
stream.shutdown().await.ok();
17071733

17081734
println!(
1709-
"Saved to: {} - Network_time: {}",
1735+
"Saved to: {} - Network_time: {}s",
17101736
output_path.display(),
17111737
&response.network_time
17121738
);

0 commit comments

Comments
 (0)