From e23031bb1cd0db50f9155ea269d0333dcad7e260 Mon Sep 17 00:00:00 2001 From: rrouviere Date: Sun, 26 Oct 2025 17:40:26 +0100 Subject: [PATCH 01/13] chore: Initialize the cargo project + basic outline Signed-off-by: rrouviere --- topics/p2p-transfer-protocol/Cargo.lock | 7 +++++++ topics/p2p-transfer-protocol/Cargo.toml | 6 ++++++ topics/p2p-transfer-protocol/src/main.rs | 21 +++++++++++++++++++++ 3 files changed, 34 insertions(+) create mode 100644 topics/p2p-transfer-protocol/Cargo.lock create mode 100644 topics/p2p-transfer-protocol/Cargo.toml create mode 100644 topics/p2p-transfer-protocol/src/main.rs diff --git a/topics/p2p-transfer-protocol/Cargo.lock b/topics/p2p-transfer-protocol/Cargo.lock new file mode 100644 index 0000000..7740ddf --- /dev/null +++ b/topics/p2p-transfer-protocol/Cargo.lock @@ -0,0 +1,7 @@ +# This file is automatically @generated by Cargo. +# It is not intended for manual editing. +version = 4 + +[[package]] +name = "p2p-transfer-protocol" +version = "0.1.0" diff --git a/topics/p2p-transfer-protocol/Cargo.toml b/topics/p2p-transfer-protocol/Cargo.toml new file mode 100644 index 0000000..d1ebebe --- /dev/null +++ b/topics/p2p-transfer-protocol/Cargo.toml @@ -0,0 +1,6 @@ +[package] +name = "p2p-transfer-protocol" +version = "0.1.0" +edition = "2024" + +[dependencies] diff --git a/topics/p2p-transfer-protocol/src/main.rs b/topics/p2p-transfer-protocol/src/main.rs new file mode 100644 index 0000000..242be00 --- /dev/null +++ b/topics/p2p-transfer-protocol/src/main.rs @@ -0,0 +1,21 @@ +fn main() { + // Parse arguments + // Go to server or client mode +} + + +fn server_mode() { + // Start server + // Bind to port + // Open output file + // Wait for client + // Receive & save file +} + +fn client_mode() { + // Start client + // Connect to server + // Open input file + // Send input file + // Show result +} From a8a159b2ec7ef8c98422cb063aba04394c3d69ef Mon Sep 17 00:00:00 2001 From: rrouviere Date: Sun, 26 Oct 2025 17:40:49 +0100 Subject: [PATCH 02/13] docs: Initial project architecture Signed-off-by: rrouviere --- .../docs/architecture.md | 52 +++++++++++++++++++ 1 file changed, 52 insertions(+) create mode 100644 topics/p2p-transfer-protocol/docs/architecture.md diff --git a/topics/p2p-transfer-protocol/docs/architecture.md b/topics/p2p-transfer-protocol/docs/architecture.md new file mode 100644 index 0000000..0473e17 --- /dev/null +++ b/topics/p2p-transfer-protocol/docs/architecture.md @@ -0,0 +1,52 @@ +## +TODO: Project definition: What is it? What are the goals of the tool/project? + +## Usage +How can one use it? Give usage examples. + + +## Modules +The project is structured around three main modules: +- The server (or receiver) module +- The client (or sender) module +- The generic protocol library, allowing reuse of the state machine in both the server and client components. + +## Seperations of concerns +This protocol is intended as an Layer 7/Application-level protocol. + +Hence, the underlying protocol stack is assumed : +- To be reliable + - => No retransmission logic + - Note: retries are not attempted by the current implementation +- To handle integrity and security concerns + - => No checksums or encryption +- To handle multiplexing and/or connection reuse + - Only a single file transfer per connection is allowed + +This implementation uses TCP as its transport protocol, which provides reliablility and weak integrity guarantees, but no confidentiality or protection against tempering of the payload. +Such a guarentee could be achieved by encapsulating this protocol using TLS. + +Another interesting possibility would be to use QUIC as the transport protocol, which could allow connection reuse, multiplexing (using streams), as well as lower latency (by merging the transport and crypto establishement as a single step). +Latency could be futher reduced by sending the HELLO packet as a 0-RTT packet on futher connections. + + +## Protocol description + +### States +- LISTEN +- HELLO_SENT +- HELLO_RECEIVED +- ACK_SENT +- ACK_RECEIVED +- NACK_RECEIVED +- ESTABLISHED + + + +## Messages +- HELLO: For the sender to offer a file to the receiver. It takes a file size argument. +- ACK: For the receiver to tell the sender it is ready to receive a proposed file. +- NACK: For the receiver to reject a proposed file. +- SEND: Send, for the sender to actually send a file. It also takes a file size argument, that must match the HELLO offer. + + From 2ca792483a0bdeab9a8615983995075d96a8a6e2 Mon Sep 17 00:00:00 2001 From: rrouviere Date: Thu, 30 Oct 2025 01:01:54 +0100 Subject: [PATCH 03/13] feat(protocol): Add basic protocol messages & states Signed-off-by: rrouviere --- topics/p2p-transfer-protocol/Cargo.lock | 72 +++++++++++++++++++ topics/p2p-transfer-protocol/Cargo.toml | 2 + .../docs/architecture.md | 33 ++++++--- topics/p2p-transfer-protocol/src/main.rs | 13 +++- .../src/protocol/connection_state.rs | 16 +++++ .../src/protocol/message.rs | 71 ++++++++++++++++++ .../p2p-transfer-protocol/src/protocol/mod.rs | 37 ++++++++++ 7 files changed, 235 insertions(+), 9 deletions(-) create mode 100644 topics/p2p-transfer-protocol/src/protocol/connection_state.rs create mode 100644 topics/p2p-transfer-protocol/src/protocol/message.rs create mode 100644 topics/p2p-transfer-protocol/src/protocol/mod.rs diff --git a/topics/p2p-transfer-protocol/Cargo.lock b/topics/p2p-transfer-protocol/Cargo.lock index 7740ddf..a12bb5c 100644 --- a/topics/p2p-transfer-protocol/Cargo.lock +++ b/topics/p2p-transfer-protocol/Cargo.lock @@ -2,6 +2,78 @@ # It is not intended for manual editing. version = 4 +[[package]] +name = "heck" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" + +[[package]] +name = "log" +version = "0.4.28" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "34080505efa8e45a4b816c349525ebe327ceaa8559756f0356cba97ef3bf7432" + [[package]] name = "p2p-transfer-protocol" version = "0.1.0" +dependencies = [ + "log", + "strum", +] + +[[package]] +name = "proc-macro2" +version = "1.0.95" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "02b3e5e68a3a1a02aad3ec490a98007cbc13c37cbe84a3cd7b8e406d76e7f778" +dependencies = [ + "unicode-ident", +] + +[[package]] +name = "quote" +version = "1.0.40" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1885c039570dc00dcb4ff087a89e185fd56bae234ddc7f056a945bf36467248d" +dependencies = [ + "proc-macro2", +] + +[[package]] +name = "strum" +version = "0.27.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "af23d6f6c1a224baef9d3f61e287d2761385a5b88fdab4eb4c6f11aeb54c4bcf" +dependencies = [ + "strum_macros", +] + +[[package]] +name = "strum_macros" +version = "0.27.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7695ce3845ea4b33927c055a39dc438a45b059f7c1b3d91d38d10355fb8cbca7" +dependencies = [ + "heck", + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "syn" +version = "2.0.104" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "17b6f705963418cdb9927482fa304bc562ece2fdd4f616084c50b7023b435a40" +dependencies = [ + "proc-macro2", + "quote", + "unicode-ident", +] + +[[package]] +name = "unicode-ident" +version = "1.0.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a5f39404a5da50712a4c1eecf25e90dd62b613502b7e925fd4e4d19b5c96512" diff --git a/topics/p2p-transfer-protocol/Cargo.toml b/topics/p2p-transfer-protocol/Cargo.toml index d1ebebe..f40c18a 100644 --- a/topics/p2p-transfer-protocol/Cargo.toml +++ b/topics/p2p-transfer-protocol/Cargo.toml @@ -4,3 +4,5 @@ version = "0.1.0" edition = "2024" [dependencies] +log = "0.4.28" +strum = { version = "0.27.2", features = ["derive"] } diff --git a/topics/p2p-transfer-protocol/docs/architecture.md b/topics/p2p-transfer-protocol/docs/architecture.md index 0473e17..fa2974f 100644 --- a/topics/p2p-transfer-protocol/docs/architecture.md +++ b/topics/p2p-transfer-protocol/docs/architecture.md @@ -7,11 +7,13 @@ How can one use it? Give usage examples. ## Modules The project is structured around three main modules: -- The server (or receiver) module -- The client (or sender) module -- The generic protocol library, allowing reuse of the state machine in both the server and client components. +- The "server" module +- The "client" module +- The generic protocol library, containing shared serialization/deserialization logic, as well as the FSM (Finite State Machine) driving the protocol's state. -## Seperations of concerns +Note: for simplicity's sake, we use "client" to refer to the active opener, that is, the peer that sends the `HELLO` message. + +## Seperations of concerns and limitations of this protocol This protocol is intended as an Layer 7/Application-level protocol. Hence, the underlying protocol stack is assumed : @@ -22,9 +24,11 @@ Hence, the underlying protocol stack is assumed : - => No checksums or encryption - To handle multiplexing and/or connection reuse - Only a single file transfer per connection is allowed + - As this protocol uses the same channel for data and control, it technically suffers from head-of-line blocking. + - In this context this however is an non-issue, given that only one file per connection can be transmitted and that resuming is not supported. -This implementation uses TCP as its transport protocol, which provides reliablility and weak integrity guarantees, but no confidentiality or protection against tempering of the payload. -Such a guarentee could be achieved by encapsulating this protocol using TLS. +This implementation uses TCP as its transport protocol, which provides reliablility and weak integrity guarantees, but no confidentiality nor protection against tempering of the payload. +Such a guarantee could be achieved by encapsulating this protocol in TLS. Another interesting possibility would be to use QUIC as the transport protocol, which could allow connection reuse, multiplexing (using streams), as well as lower latency (by merging the transport and crypto establishement as a single step). Latency could be futher reduced by sending the HELLO packet as a 0-RTT packet on futher connections. @@ -33,15 +37,25 @@ Latency could be futher reduced by sending the HELLO packet as a 0-RTT packet on ## Protocol description ### States -- LISTEN +#### Initial states +- LISTENING (passive opener) - HELLO_SENT + +#### Transition states - HELLO_RECEIVED - ACK_SENT - ACK_RECEIVED -- NACK_RECEIVED - ESTABLISHED +#### Final states +- NACK_RECEIVED +- (Implicit teardown) +## Messages +- HELLO +- ACK +- NACK +- SEND ## Messages - HELLO: For the sender to offer a file to the receiver. It takes a file size argument. @@ -50,3 +64,6 @@ Latency could be futher reduced by sending the HELLO packet as a 0-RTT packet on - SEND: Send, for the sender to actually send a file. It also takes a file size argument, that must match the HELLO offer. +## Limitations & futher work +- The serialization/deserialization logic could be simplified using serde. +- \ No newline at end of file diff --git a/topics/p2p-transfer-protocol/src/main.rs b/topics/p2p-transfer-protocol/src/main.rs index 242be00..23cf0ed 100644 --- a/topics/p2p-transfer-protocol/src/main.rs +++ b/topics/p2p-transfer-protocol/src/main.rs @@ -1,8 +1,12 @@ +mod client; +mod protocol; +//mod server; + fn main() { // Parse arguments // Go to server or client mode -} +} fn server_mode() { // Start server @@ -10,6 +14,10 @@ fn server_mode() { // Open output file // Wait for client // Receive & save file + + let hello = "HELLO".parse::().unwrap(); + print!("Parsed message: {}", hello); + } fn client_mode() { @@ -18,4 +26,7 @@ fn client_mode() { // Open input file // Send input file // Show result + + let message = protocol::message::Message::ACK; + print!("Message: {}", message); } diff --git a/topics/p2p-transfer-protocol/src/protocol/connection_state.rs b/topics/p2p-transfer-protocol/src/protocol/connection_state.rs new file mode 100644 index 0000000..b510ff9 --- /dev/null +++ b/topics/p2p-transfer-protocol/src/protocol/connection_state.rs @@ -0,0 +1,16 @@ +#[derive(Debug, PartialEq, Eq)] +pub enum ConnectionState { + // Initial+Final state + Closed, + + HelloSent, + Listening, + + HelloReceived, + SendReceived, + ACKSent, + NACKSent, + ACKReceived, + Established, + NACKReceived, +} diff --git a/topics/p2p-transfer-protocol/src/protocol/message.rs b/topics/p2p-transfer-protocol/src/protocol/message.rs new file mode 100644 index 0000000..c33b990 --- /dev/null +++ b/topics/p2p-transfer-protocol/src/protocol/message.rs @@ -0,0 +1,71 @@ +use std::{fmt::Display, str::FromStr}; + +#[cfg(test)] +use strum::EnumIter; + +#[derive(Debug, PartialEq, Eq)] +#[cfg_attr(test,derive(EnumIter))] +pub enum Message { + Hello, + ACK, + NACK, + Send, +} + +// Wire format +impl Display for Message { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let s = match self { + Message::Hello => "HELLO", + Message::ACK => "ACK", + Message::NACK => "NACK", + Message::Send => "SEND", + }; + write!(f, "{}", s) + } +} + +impl FromStr for Message { + type Err = (); + + fn from_str(s: &str) -> Result { + match s.to_ascii_uppercase().as_str() { + "HELLO" => Ok(Message::Hello), + "ACK" => Ok(Message::ACK), + "NACK" => Ok(Message::NACK), + "SEND" => Ok(Message::Send), + _ => Err(()), + } + } +} + + + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_serialize() { + assert_eq!(Message::Hello.to_string(), "HELLO"); + assert_eq!(Message::ACK.to_string(), "ACK"); + assert_eq!(Message::NACK.to_string(), "NACK"); + assert_eq!(Message::Send.to_string(), "SEND"); + + assert_eq!("hello".parse::(), Ok(Message::Hello)); + assert_eq!("ACK".parse::(), Ok(Message::ACK)); + assert_eq!("nack".parse::(), Ok(Message::NACK)); + assert_eq!("SEND".parse::(), Ok(Message::Send)); + assert!("unknown".parse::().is_err()); + } + + + #[test] + fn test_deserialize() { + assert_eq!("hello".parse::(), Ok(Message::Hello)); + assert_eq!("ACK".parse::(), Ok(Message::ACK)); + assert_eq!("nack".parse::(), Ok(Message::NACK)); + assert_eq!("SEND".parse::(), Ok(Message::Send)); + assert!("unknown".parse::().is_err()); + } +} \ No newline at end of file diff --git a/topics/p2p-transfer-protocol/src/protocol/mod.rs b/topics/p2p-transfer-protocol/src/protocol/mod.rs new file mode 100644 index 0000000..2721894 --- /dev/null +++ b/topics/p2p-transfer-protocol/src/protocol/mod.rs @@ -0,0 +1,37 @@ +use log::{debug, warn}; + +pub mod message; +mod connection_state; + +use connection_state::ConnectionState; + +pub struct StateMachine { + state: ConnectionState, +} + +impl StateMachine { + pub fn new() -> Self { + StateMachine { + state: ConnectionState::Closed, + } + } + + pub fn transition(&mut self, new_state: &ConnectionState) { + debug!("Transitioning from {:?} to {:?}", self.state, new_state); + match (&self.state, new_state) { + (ConnectionState::Closed, ConnectionState::Listening) => {} + (ConnectionState::Listening, ConnectionState::HelloReceived) => {} + (ConnectionState::HelloReceived, ConnectionState::NACKSent) => {} + (ConnectionState::HelloReceived, ConnectionState::ACKSent) => {} + + (ConnectionState::Closed, ConnectionState::HelloSent) => {} + (ConnectionState::HelloSent, ConnectionState::NACKReceived) => {} + (ConnectionState::HelloSent, ConnectionState::ACKReceived) => {} + (ConnectionState::ACKReceived, ConnectionState::Established) => {} + + + (s, _) => {warn!("Invalid state transition from {:?} to {:?}, ignoring.", s, new_state);} + }; + debug!("New state: {:?}", self.state); + } +} From 0a4461d7b975fb077c5f5e758a9fc6d1a7024171 Mon Sep 17 00:00:00 2001 From: rrouviere Date: Thu, 30 Oct 2025 19:24:21 +0100 Subject: [PATCH 04/13] feat(cli): Add argument parsing Signed-off-by: rrouviere --- topics/p2p-transfer-protocol/Cargo.lock | 217 ++++++++++++++++++++++- topics/p2p-transfer-protocol/Cargo.toml | 1 + topics/p2p-transfer-protocol/src/main.rs | 61 +++++-- 3 files changed, 255 insertions(+), 24 deletions(-) diff --git a/topics/p2p-transfer-protocol/Cargo.lock b/topics/p2p-transfer-protocol/Cargo.lock index a12bb5c..ee69af0 100644 --- a/topics/p2p-transfer-protocol/Cargo.lock +++ b/topics/p2p-transfer-protocol/Cargo.lock @@ -2,44 +2,159 @@ # It is not intended for manual editing. version = 4 +[[package]] +name = "anstream" +version = "0.6.21" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "43d5b281e737544384e969a5ccad3f1cdd24b48086a0fc1b2a5262a26b8f4f4a" +dependencies = [ + "anstyle", + "anstyle-parse", + "anstyle-query", + "anstyle-wincon", + "colorchoice", + "is_terminal_polyfill", + "utf8parse", +] + +[[package]] +name = "anstyle" +version = "1.0.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5192cca8006f1fd4f7237516f40fa183bb07f8fbdfedaa0036de5ea9b0b45e78" + +[[package]] +name = "anstyle-parse" +version = "0.2.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4e7644824f0aa2c7b9384579234ef10eb7efb6a0deb83f9630a49594dd9c15c2" +dependencies = [ + "utf8parse", +] + +[[package]] +name = "anstyle-query" +version = "1.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9e231f6134f61b71076a3eab506c379d4f36122f2af15a9ff04415ea4c3339e2" +dependencies = [ + "windows-sys", +] + +[[package]] +name = "anstyle-wincon" +version = "3.0.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3e0633414522a32ffaac8ac6cc8f748e090c5717661fddeea04219e2344f5f2a" +dependencies = [ + "anstyle", + "once_cell_polyfill", + "windows-sys", +] + +[[package]] +name = "clap" +version = "4.5.51" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4c26d721170e0295f191a69bd9a1f93efcdb0aff38684b61ab5750468972e5f5" +dependencies = [ + "clap_builder", + "clap_derive", +] + +[[package]] +name = "clap_builder" +version = "4.5.51" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "75835f0c7bf681bfd05abe44e965760fea999a5286c6eb2d59883634fd02011a" +dependencies = [ + "anstream", + "anstyle", + "clap_lex", + "strsim", +] + +[[package]] +name = "clap_derive" +version = "4.5.49" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2a0b5487afeab2deb2ff4e03a807ad1a03ac532ff5a2cee5d86884440c7f7671" +dependencies = [ + "heck", + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "clap_lex" +version = "0.7.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a1d728cc89cf3aee9ff92b05e62b19ee65a02b5702cff7d5a377e32c6ae29d8d" + +[[package]] +name = "colorchoice" +version = "1.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b05b61dc5112cbb17e4b6cd61790d9845d13888356391624cbe7e41efeac1e75" + [[package]] name = "heck" version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" +[[package]] +name = "is_terminal_polyfill" +version = "1.70.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a6cb138bb79a146c1bd460005623e142ef0181e3d0219cb493e02f7d08a35695" + [[package]] name = "log" version = "0.4.28" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "34080505efa8e45a4b816c349525ebe327ceaa8559756f0356cba97ef3bf7432" +[[package]] +name = "once_cell_polyfill" +version = "1.70.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "384b8ab6d37215f3c5301a95a4accb5d64aa607f1fcb26a11b5303878451b4fe" + [[package]] name = "p2p-transfer-protocol" version = "0.1.0" dependencies = [ + "clap", "log", "strum", ] [[package]] name = "proc-macro2" -version = "1.0.95" +version = "1.0.103" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "02b3e5e68a3a1a02aad3ec490a98007cbc13c37cbe84a3cd7b8e406d76e7f778" +checksum = "5ee95bc4ef87b8d5ba32e8b7714ccc834865276eab0aed5c9958d00ec45f49e8" dependencies = [ "unicode-ident", ] [[package]] name = "quote" -version = "1.0.40" +version = "1.0.41" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1885c039570dc00dcb4ff087a89e185fd56bae234ddc7f056a945bf36467248d" +checksum = "ce25767e7b499d1b604768e7cde645d14cc8584231ea6b295e9c9eb22c02e1d1" dependencies = [ "proc-macro2", ] +[[package]] +name = "strsim" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7da8b5736845d9f2fcb837ea5d9e2628564b3b043a70948a3f0b778838c5fb4f" + [[package]] name = "strum" version = "0.27.2" @@ -63,9 +178,9 @@ dependencies = [ [[package]] name = "syn" -version = "2.0.104" +version = "2.0.108" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "17b6f705963418cdb9927482fa304bc562ece2fdd4f616084c50b7023b435a40" +checksum = "da58917d35242480a05c2897064da0a80589a2a0476c9a3f2fdc83b53502e917" dependencies = [ "proc-macro2", "quote", @@ -74,6 +189,92 @@ dependencies = [ [[package]] name = "unicode-ident" -version = "1.0.18" +version = "1.0.22" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9312f7c4f6ff9069b165498234ce8be658059c6728633667c526e27dc2cf1df5" + +[[package]] +name = "utf8parse" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" + +[[package]] +name = "windows-link" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f0805222e57f7521d6a62e36fa9163bc891acd422f971defe97d64e70d0a4fe5" + +[[package]] +name = "windows-sys" +version = "0.60.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2f500e4d28234f72040990ec9d39e3a6b950f9f22d3dba18416c35882612bcb" +dependencies = [ + "windows-targets", +] + +[[package]] +name = "windows-targets" +version = "0.53.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4945f9f551b88e0d65f3db0bc25c33b8acea4d9e41163edf90dcd0b19f9069f3" +dependencies = [ + "windows-link", + "windows_aarch64_gnullvm", + "windows_aarch64_msvc", + "windows_i686_gnu", + "windows_i686_gnullvm", + "windows_i686_msvc", + "windows_x86_64_gnu", + "windows_x86_64_gnullvm", + "windows_x86_64_msvc", +] + +[[package]] +name = "windows_aarch64_gnullvm" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a9d8416fa8b42f5c947f8482c43e7d89e73a173cead56d044f6a56104a6d1b53" + +[[package]] +name = "windows_aarch64_msvc" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b9d782e804c2f632e395708e99a94275910eb9100b2114651e04744e9b125006" + +[[package]] +name = "windows_i686_gnu" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "960e6da069d81e09becb0ca57a65220ddff016ff2d6af6a223cf372a506593a3" + +[[package]] +name = "windows_i686_gnullvm" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fa7359d10048f68ab8b09fa71c3daccfb0e9b559aed648a8f95469c27057180c" + +[[package]] +name = "windows_i686_msvc" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e7ac75179f18232fe9c285163565a57ef8d3c89254a30685b57d83a38d326c2" + +[[package]] +name = "windows_x86_64_gnu" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9c3842cdd74a865a8066ab39c8a7a473c0778a3f29370b5fd6b4b9aa7df4a499" + +[[package]] +name = "windows_x86_64_gnullvm" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0ffa179e2d07eee8ad8f57493436566c7cc30ac536a3379fdf008f47f6bb7ae1" + +[[package]] +name = "windows_x86_64_msvc" +version = "0.53.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5a5f39404a5da50712a4c1eecf25e90dd62b613502b7e925fd4e4d19b5c96512" +checksum = "d6bbff5f0aada427a1e5a6da5f1f98158182f26556f345ac9e04d36d0ebed650" diff --git a/topics/p2p-transfer-protocol/Cargo.toml b/topics/p2p-transfer-protocol/Cargo.toml index f40c18a..e6becd5 100644 --- a/topics/p2p-transfer-protocol/Cargo.toml +++ b/topics/p2p-transfer-protocol/Cargo.toml @@ -4,5 +4,6 @@ version = "0.1.0" edition = "2024" [dependencies] +clap = { version = "4.5.51", features = ["derive"] } log = "0.4.28" strum = { version = "0.27.2", features = ["derive"] } diff --git a/topics/p2p-transfer-protocol/src/main.rs b/topics/p2p-transfer-protocol/src/main.rs index 23cf0ed..91f8264 100644 --- a/topics/p2p-transfer-protocol/src/main.rs +++ b/topics/p2p-transfer-protocol/src/main.rs @@ -2,31 +2,60 @@ mod client; mod protocol; //mod server; -fn main() { - // Parse arguments - // Go to server or client mode +use clap::{Parser, Subcommand}; + +#[derive(Parser)] +#[command(name = "p2p-transfer-protocol", about = "Simple P2P file transfer protocol")] +struct Cli { + #[command(subcommand)] + command: Commands, +} + +#[derive(Subcommand)] +enum Commands { + Listen { + #[arg(long, default_value = "[::]:0")] // Listen on any inet & let the kernel pick a port. + bind: String, + #[arg(long, default_value = "out_file")] + output: String, + }, + Send { + #[arg(long, default_value = "out_file")] + file: String, + #[arg()] + remote_addr: String, + }, +} + +fn main() { + let cli = Cli::parse(); + // Note: if no command matches, clap automatically provides the help message & exits. + match cli.command { + Commands::Listen { bind, output } => { + server_mode(&bind, &output); + } + Commands::Send { file, remote_addr } => { + client_mode(&file, &remote_addr); + } + } } -fn server_mode() { +fn server_mode(bind_addr: &str, output_file: &str) { // Start server - // Bind to port - // Open output file + // Bind to port: bind_addr + // Open output file: output_file // Wait for client // Receive & save file - - let hello = "HELLO".parse::().unwrap(); - print!("Parsed message: {}", hello); - + println!("Server listening on {}. Saving to {}", bind_addr, output_file); } -fn client_mode() { +fn client_mode(file: &str, remote_addr: &str) { // Start client - // Connect to server - // Open input file + // Connect to server: remote_addr + // Open input file: file // Send input file // Show result - - let message = protocol::message::Message::ACK; - print!("Message: {}", message); + println!("Client sending {} to {}", file, remote_addr); + // ...rest of client logic... } From 7f1dbcf523fac6eb90f29c86d72f23cf3985e312 Mon Sep 17 00:00:00 2001 From: rrouviere Date: Thu, 30 Oct 2025 20:32:17 +0100 Subject: [PATCH 05/13] feat(client): basic client impl Signed-off-by: rrouviere --- .../p2p-transfer-protocol/src/client/mod.rs | 68 ++++++++++++++++ topics/p2p-transfer-protocol/src/main.rs | 36 ++++++--- .../src/protocol/connection_state.rs | 1 - .../src/protocol/message.rs | 77 ++++++++----------- .../p2p-transfer-protocol/src/protocol/mod.rs | 2 +- 5 files changed, 129 insertions(+), 55 deletions(-) create mode 100644 topics/p2p-transfer-protocol/src/client/mod.rs diff --git a/topics/p2p-transfer-protocol/src/client/mod.rs b/topics/p2p-transfer-protocol/src/client/mod.rs new file mode 100644 index 0000000..19e8d66 --- /dev/null +++ b/topics/p2p-transfer-protocol/src/client/mod.rs @@ -0,0 +1,68 @@ +use std::io::{self, Read, Write, BufRead, BufReader}; +use std::net::TcpStream; +use std::fs::File; +use log::{debug, info, warn}; + +use crate::protocol::StateMachine; +use crate::protocol::message::Message; +use crate::protocol::connection_state::ConnectionState; + + +pub fn run_client(mut file: File, mut stream: TcpStream) -> io::Result<()> { + let file_size = file.metadata()?.len(); + info!("Sending a {} bytes file", file_size); + + // Initialize protocol state machine + let mut sm = StateMachine::new(); + + // Send HELLO + let hello_msg = Message::Hello { file_size }.to_string(); + stream.write_all(hello_msg.as_bytes())?; + sm.transition(&ConnectionState::HelloSent); + debug!("Sent HELLO"); + + // Wait for ACK/NACK + let mut reader = BufReader::new(&mut stream); + let mut response = String::new(); + reader.read_line(&mut response)?; + let response = response.trim(); + match response.parse::() { + Ok(Message::ACK) => { + debug!("Received ACK"); + sm.transition(&ConnectionState::ACKReceived); + } + Ok(Message::NACK) => { + debug!("Received NACK"); + sm.transition(&ConnectionState::NACKReceived); + return Ok(()); + } + _ => { + warn!("Unexpected response from server: {}", response); + return Ok(()); + } + } + + // Start sending data + let send_msg = Message::Send.to_string(); + stream.write_all(send_msg.as_bytes())?; + sm.transition(&ConnectionState::Established); + + debug!("Sending data."); + + // Send file data in 4kB chunks + let mut buffer = [0u8; 4096]; + + loop { + let n = file.read(&mut buffer)?; + if n == 0 { break; } // EOF + stream.write_all(&buffer[..n])?; + } + // Flush to ensure all data is sent + stream.flush()?; + info!("File sent successfully"); + sm.transition(&ConnectionState::Closed); + // Close socket + stream.shutdown(std::net::Shutdown::Both)?; + debug!("Transfer complete"); + Ok(()) +} \ No newline at end of file diff --git a/topics/p2p-transfer-protocol/src/main.rs b/topics/p2p-transfer-protocol/src/main.rs index 91f8264..87b0716 100644 --- a/topics/p2p-transfer-protocol/src/main.rs +++ b/topics/p2p-transfer-protocol/src/main.rs @@ -1,8 +1,10 @@ mod client; mod protocol; -//mod server; + +use std::{fs::File, net::TcpStream, process::exit}; use clap::{Parser, Subcommand}; +use log::{error, info}; #[derive(Parser)] #[command(name = "p2p-transfer-protocol", about = "Simple P2P file transfer protocol")] @@ -50,12 +52,28 @@ fn server_mode(bind_addr: &str, output_file: &str) { println!("Server listening on {}. Saving to {}", bind_addr, output_file); } -fn client_mode(file: &str, remote_addr: &str) { - // Start client - // Connect to server: remote_addr - // Open input file: file - // Send input file - // Show result - println!("Client sending {} to {}", file, remote_addr); - // ...rest of client logic... +fn client_mode(file_path: &str, remote_addr: &str) { + // Connect to server + let mut stream = match TcpStream::connect(remote_addr) { + Ok(mut s) => s, + Err(e) => { + error!("Failed to connect to {}: {}", remote_addr, e); + exit(1); + } + }; + info!("Connected to {}", remote_addr); + + // Open input file + let mut file = match File::open(file_path){ + Ok(f) => f, + Err(e) => { + error!("Failed to open file {}: {}", file_path, e); + exit(1); + } + }; + + match client::run_client(file, stream) { + Ok(_) => info!("File sent successfully"), + Err(e) => error!("Failed to send file: {}", e), + } } diff --git a/topics/p2p-transfer-protocol/src/protocol/connection_state.rs b/topics/p2p-transfer-protocol/src/protocol/connection_state.rs index b510ff9..dc1a3ba 100644 --- a/topics/p2p-transfer-protocol/src/protocol/connection_state.rs +++ b/topics/p2p-transfer-protocol/src/protocol/connection_state.rs @@ -7,7 +7,6 @@ pub enum ConnectionState { Listening, HelloReceived, - SendReceived, ACKSent, NACKSent, ACKReceived, diff --git a/topics/p2p-transfer-protocol/src/protocol/message.rs b/topics/p2p-transfer-protocol/src/protocol/message.rs index c33b990..d5219ac 100644 --- a/topics/p2p-transfer-protocol/src/protocol/message.rs +++ b/topics/p2p-transfer-protocol/src/protocol/message.rs @@ -1,71 +1,60 @@ use std::{fmt::Display, str::FromStr}; - +use std::io::Write; +use log::error; #[cfg(test)] use strum::EnumIter; #[derive(Debug, PartialEq, Eq)] #[cfg_attr(test,derive(EnumIter))] pub enum Message { - Hello, + Hello { + file_size: u64 + }, ACK, NACK, - Send, + Send } + // Wire format impl Display for Message { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - let s = match self { - Message::Hello => "HELLO", - Message::ACK => "ACK", - Message::NACK => "NACK", - Message::Send => "SEND", - }; - write!(f, "{}", s) + match self { + Message::Hello { file_size } => write!(f, "HELLO {file_size}\n"), + Message::ACK => write!(f, "ACK\n"), + Message::NACK => write!(f, "NACK\n"), + Message::Send => write!(f, "SEND\n"), + } } } + impl FromStr for Message { type Err = (); fn from_str(s: &str) -> Result { - match s.to_ascii_uppercase().as_str() { - "HELLO" => Ok(Message::Hello), + // Normalize + let s = s.trim(); + let upper = s.to_ascii_uppercase(); + + if upper.starts_with("HELLO") { + let parts: Vec<&str> = s.split_whitespace().collect(); + if parts.len() == 2 { + if let Ok(file_size) = parts[1].parse::() { + return Ok(Message::Hello { file_size }); + } + } + error!("Invalid HELLO message"); + return Err(()); + }; + + // No arguments msgs + return match upper.as_str() { "ACK" => Ok(Message::ACK), "NACK" => Ok(Message::NACK), "SEND" => Ok(Message::Send), + // Unknown message _ => Err(()), - } + }; } } - - - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn test_serialize() { - assert_eq!(Message::Hello.to_string(), "HELLO"); - assert_eq!(Message::ACK.to_string(), "ACK"); - assert_eq!(Message::NACK.to_string(), "NACK"); - assert_eq!(Message::Send.to_string(), "SEND"); - - assert_eq!("hello".parse::(), Ok(Message::Hello)); - assert_eq!("ACK".parse::(), Ok(Message::ACK)); - assert_eq!("nack".parse::(), Ok(Message::NACK)); - assert_eq!("SEND".parse::(), Ok(Message::Send)); - assert!("unknown".parse::().is_err()); - } - - - #[test] - fn test_deserialize() { - assert_eq!("hello".parse::(), Ok(Message::Hello)); - assert_eq!("ACK".parse::(), Ok(Message::ACK)); - assert_eq!("nack".parse::(), Ok(Message::NACK)); - assert_eq!("SEND".parse::(), Ok(Message::Send)); - assert!("unknown".parse::().is_err()); - } -} \ No newline at end of file diff --git a/topics/p2p-transfer-protocol/src/protocol/mod.rs b/topics/p2p-transfer-protocol/src/protocol/mod.rs index 2721894..f1cf2f6 100644 --- a/topics/p2p-transfer-protocol/src/protocol/mod.rs +++ b/topics/p2p-transfer-protocol/src/protocol/mod.rs @@ -1,7 +1,7 @@ use log::{debug, warn}; pub mod message; -mod connection_state; +pub mod connection_state; use connection_state::ConnectionState; From 001318e530be2660008e8c55803c143cd1c37711 Mon Sep 17 00:00:00 2001 From: rrouviere Date: Thu, 30 Oct 2025 22:10:28 +0100 Subject: [PATCH 06/13] feat(server): Basic server impl. \n Known issue: extratenous 0000000000 after file content. Signed-off-by: rrouviere --- topics/p2p-transfer-protocol/Cargo.lock | 136 ++++++++++++++++++ topics/p2p-transfer-protocol/Cargo.toml | 1 + topics/p2p-transfer-protocol/out_file | 4 + .../p2p-transfer-protocol/src/client/mod.rs | 27 ++-- topics/p2p-transfer-protocol/src/main.rs | 68 +++++++-- .../src/protocol/connection_state.rs | 5 +- .../src/protocol/message.rs | 15 +- .../p2p-transfer-protocol/src/protocol/mod.rs | 29 ++-- .../p2p-transfer-protocol/src/server/mod.rs | 122 ++++++++++++++++ 9 files changed, 363 insertions(+), 44 deletions(-) create mode 100644 topics/p2p-transfer-protocol/out_file create mode 100644 topics/p2p-transfer-protocol/src/server/mod.rs diff --git a/topics/p2p-transfer-protocol/Cargo.lock b/topics/p2p-transfer-protocol/Cargo.lock index ee69af0..2184a2d 100644 --- a/topics/p2p-transfer-protocol/Cargo.lock +++ b/topics/p2p-transfer-protocol/Cargo.lock @@ -2,6 +2,15 @@ # It is not intended for manual editing. version = 4 +[[package]] +name = "aho-corasick" +version = "1.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ddd31a130427c27518df266943a5308ed92d4b226cc639f5a8f1002816174301" +dependencies = [ + "memchr", +] + [[package]] name = "anstream" version = "0.6.21" @@ -98,6 +107,29 @@ version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b05b61dc5112cbb17e4b6cd61790d9845d13888356391624cbe7e41efeac1e75" +[[package]] +name = "env_filter" +version = "0.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1bf3c259d255ca70051b30e2e95b5446cdb8949ac4cd22c0d7fd634d89f568e2" +dependencies = [ + "log", + "regex", +] + +[[package]] +name = "env_logger" +version = "0.11.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "13c863f0904021b108aa8b2f55046443e6b1ebde8fd4a15c399893aae4fa069f" +dependencies = [ + "anstream", + "anstyle", + "env_filter", + "jiff", + "log", +] + [[package]] name = "heck" version = "0.5.0" @@ -110,12 +142,42 @@ version = "1.70.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a6cb138bb79a146c1bd460005623e142ef0181e3d0219cb493e02f7d08a35695" +[[package]] +name = "jiff" +version = "0.2.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "be1f93b8b1eb69c77f24bbb0afdf66f54b632ee39af40ca21c4365a1d7347e49" +dependencies = [ + "jiff-static", + "log", + "portable-atomic", + "portable-atomic-util", + "serde", +] + +[[package]] +name = "jiff-static" +version = "0.2.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "03343451ff899767262ec32146f6d559dd759fdadf42ff0e227c7c48f72594b4" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "log" version = "0.4.28" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "34080505efa8e45a4b816c349525ebe327ceaa8559756f0356cba97ef3bf7432" +[[package]] +name = "memchr" +version = "2.7.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f52b00d39961fc5b2736ea853c9cc86238e165017a493d1d5c8eac6bdc4cc273" + [[package]] name = "once_cell_polyfill" version = "1.70.2" @@ -127,10 +189,26 @@ name = "p2p-transfer-protocol" version = "0.1.0" dependencies = [ "clap", + "env_logger", "log", "strum", ] +[[package]] +name = "portable-atomic" +version = "1.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f84267b20a16ea918e43c6a88433c2d54fa145c92a811b5b047ccbe153674483" + +[[package]] +name = "portable-atomic-util" +version = "0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d8a2f0d8d040d7848a709caf78912debcc3f33ee4b3cac47d73d1e1069e83507" +dependencies = [ + "portable-atomic", +] + [[package]] name = "proc-macro2" version = "1.0.103" @@ -149,6 +227,64 @@ dependencies = [ "proc-macro2", ] +[[package]] +name = "regex" +version = "1.12.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "843bc0191f75f3e22651ae5f1e72939ab2f72a4bc30fa80a066bd66edefc24d4" +dependencies = [ + "aho-corasick", + "memchr", + "regex-automata", + "regex-syntax", +] + +[[package]] +name = "regex-automata" +version = "0.4.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5276caf25ac86c8d810222b3dbb938e512c55c6831a10f3e6ed1c93b84041f1c" +dependencies = [ + "aho-corasick", + "memchr", + "regex-syntax", +] + +[[package]] +name = "regex-syntax" +version = "0.8.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7a2d987857b319362043e95f5353c0535c1f58eec5336fdfcf626430af7def58" + +[[package]] +name = "serde" +version = "1.0.228" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9a8e94ea7f378bd32cbbd37198a4a91436180c5bb472411e48b5ec2e2124ae9e" +dependencies = [ + "serde_core", +] + +[[package]] +name = "serde_core" +version = "1.0.228" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "41d385c7d4ca58e59fc732af25c3983b67ac852c1a25000afe1175de458b67ad" +dependencies = [ + "serde_derive", +] + +[[package]] +name = "serde_derive" +version = "1.0.228" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d540f220d3187173da220f885ab66608367b6574e925011a9353e4badda91d79" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "strsim" version = "0.11.1" diff --git a/topics/p2p-transfer-protocol/Cargo.toml b/topics/p2p-transfer-protocol/Cargo.toml index e6becd5..dfb540e 100644 --- a/topics/p2p-transfer-protocol/Cargo.toml +++ b/topics/p2p-transfer-protocol/Cargo.toml @@ -7,3 +7,4 @@ edition = "2024" clap = { version = "4.5.51", features = ["derive"] } log = "0.4.28" strum = { version = "0.27.2", features = ["derive"] } +env_logger = "0.11.8" \ No newline at end of file diff --git a/topics/p2p-transfer-protocol/out_file b/topics/p2p-transfer-protocol/out_file new file mode 100644 index 0000000..437a81b --- /dev/null +++ b/topics/p2p-transfer-protocol/out_file @@ -0,0 +1,4 @@ +12345678900 +0000000000 + + diff --git a/topics/p2p-transfer-protocol/src/client/mod.rs b/topics/p2p-transfer-protocol/src/client/mod.rs index 19e8d66..1491446 100644 --- a/topics/p2p-transfer-protocol/src/client/mod.rs +++ b/topics/p2p-transfer-protocol/src/client/mod.rs @@ -1,12 +1,11 @@ -use std::io::{self, Read, Write, BufRead, BufReader}; -use std::net::TcpStream; -use std::fs::File; use log::{debug, info, warn}; +use std::fs::File; +use std::io::{self, BufRead, BufReader, Read, Write}; +use std::net::TcpStream; use crate::protocol::StateMachine; -use crate::protocol::message::Message; use crate::protocol::connection_state::ConnectionState; - +use crate::protocol::message::Message; pub fn run_client(mut file: File, mut stream: TcpStream) -> io::Result<()> { let file_size = file.metadata()?.len(); @@ -18,7 +17,7 @@ pub fn run_client(mut file: File, mut stream: TcpStream) -> io::Result<()> { // Send HELLO let hello_msg = Message::Hello { file_size }.to_string(); stream.write_all(hello_msg.as_bytes())?; - sm.transition(&ConnectionState::HelloSent); + sm.transition(ConnectionState::HelloSent); debug!("Sent HELLO"); // Wait for ACK/NACK @@ -29,11 +28,11 @@ pub fn run_client(mut file: File, mut stream: TcpStream) -> io::Result<()> { match response.parse::() { Ok(Message::ACK) => { debug!("Received ACK"); - sm.transition(&ConnectionState::ACKReceived); + sm.transition(ConnectionState::ACKReceived); } Ok(Message::NACK) => { debug!("Received NACK"); - sm.transition(&ConnectionState::NACKReceived); + sm.transition(ConnectionState::NACKReceived); return Ok(()); } _ => { @@ -45,8 +44,8 @@ pub fn run_client(mut file: File, mut stream: TcpStream) -> io::Result<()> { // Start sending data let send_msg = Message::Send.to_string(); stream.write_all(send_msg.as_bytes())?; - sm.transition(&ConnectionState::Established); - + sm.transition(ConnectionState::Established); + debug!("Sending data."); // Send file data in 4kB chunks @@ -54,15 +53,17 @@ pub fn run_client(mut file: File, mut stream: TcpStream) -> io::Result<()> { loop { let n = file.read(&mut buffer)?; - if n == 0 { break; } // EOF + if n == 0 { + break; + } // EOF stream.write_all(&buffer[..n])?; } // Flush to ensure all data is sent stream.flush()?; info!("File sent successfully"); - sm.transition(&ConnectionState::Closed); + sm.transition(ConnectionState::Closed); // Close socket stream.shutdown(std::net::Shutdown::Both)?; debug!("Transfer complete"); Ok(()) -} \ No newline at end of file +} diff --git a/topics/p2p-transfer-protocol/src/main.rs b/topics/p2p-transfer-protocol/src/main.rs index 87b0716..28a9aa0 100644 --- a/topics/p2p-transfer-protocol/src/main.rs +++ b/topics/p2p-transfer-protocol/src/main.rs @@ -1,13 +1,22 @@ mod client; mod protocol; +mod server; -use std::{fs::File, net::TcpStream, process::exit}; +use std::{ + fs::{File, OpenOptions}, + net::{TcpListener, TcpStream}, + process::exit, + io, +}; use clap::{Parser, Subcommand}; use log::{error, info}; #[derive(Parser)] -#[command(name = "p2p-transfer-protocol", about = "Simple P2P file transfer protocol")] +#[command( + name = "p2p-transfer-protocol", + about = "Simple P2P file transfer protocol" +)] struct Cli { #[command(subcommand)] command: Commands, @@ -31,11 +40,15 @@ enum Commands { } fn main() { + // Initialize logger + env_logger::init(); + + let cli = Cli::parse(); - // Note: if no command matches, clap automatically provides the help message & exits. + // Note: if no command matches, clap automatically provides the help message & exits. match cli.command { Commands::Listen { bind, output } => { - server_mode(&bind, &output); + let _ = server_mode(&bind, &output); } Commands::Send { file, remote_addr } => { client_mode(&file, &remote_addr); @@ -43,19 +56,52 @@ fn main() { } } -fn server_mode(bind_addr: &str, output_file: &str) { +fn server_mode(bind_addr: &str, file_path: &str) -> io::Result<()> { // Start server // Bind to port: bind_addr - // Open output file: output_file + let listener = match TcpListener::bind(&bind_addr) { + Ok(l) => l, + Err(e) => { + error!("Failed to bind to {}: {}", bind_addr, e); + exit(1); + } + }; + info!("Server listening on {}", bind_addr); + + // Open or create output file: + let mut file = OpenOptions::new() + .write(true) + .create(true) + .open(file_path)?; + + println!( + "Server listening on {:?}. Saving to {}", + listener, file_path + ); + // Wait for client - // Receive & save file - println!("Server listening on {}. Saving to {}", bind_addr, output_file); + for stream in listener.incoming() { + match stream { + Ok(mut s) => { + info!("Client connected: {}", s.peer_addr().unwrap()); + match server::run_server(&mut file, &mut s) { + Ok(_) => info!("File received successfully"), + Err(e) => error!("Failed to receive file: {}", e), + } + } + Err(e) => { + error!("Connection failed: {}", e); + } + } + } + + Ok(()) } fn client_mode(file_path: &str, remote_addr: &str) { // Connect to server - let mut stream = match TcpStream::connect(remote_addr) { - Ok(mut s) => s, + let stream = match TcpStream::connect(remote_addr) { + Ok(s) => s, Err(e) => { error!("Failed to connect to {}: {}", remote_addr, e); exit(1); @@ -64,7 +110,7 @@ fn client_mode(file_path: &str, remote_addr: &str) { info!("Connected to {}", remote_addr); // Open input file - let mut file = match File::open(file_path){ + let file = match File::open(file_path) { Ok(f) => f, Err(e) => { error!("Failed to open file {}: {}", file_path, e); diff --git a/topics/p2p-transfer-protocol/src/protocol/connection_state.rs b/topics/p2p-transfer-protocol/src/protocol/connection_state.rs index dc1a3ba..366817a 100644 --- a/topics/p2p-transfer-protocol/src/protocol/connection_state.rs +++ b/topics/p2p-transfer-protocol/src/protocol/connection_state.rs @@ -1,13 +1,14 @@ -#[derive(Debug, PartialEq, Eq)] +#[derive(Debug, PartialEq, Eq, Copy, Clone)] pub enum ConnectionState { // Initial+Final state Closed, - + HelloSent, Listening, HelloReceived, ACKSent, + #[allow(dead_code)] // TODO: Implement NACK handling. NACKSent, ACKReceived, Established, diff --git a/topics/p2p-transfer-protocol/src/protocol/message.rs b/topics/p2p-transfer-protocol/src/protocol/message.rs index d5219ac..1b0e92a 100644 --- a/topics/p2p-transfer-protocol/src/protocol/message.rs +++ b/topics/p2p-transfer-protocol/src/protocol/message.rs @@ -1,21 +1,17 @@ -use std::{fmt::Display, str::FromStr}; -use std::io::Write; use log::error; +use std::{fmt::Display, str::FromStr}; #[cfg(test)] use strum::EnumIter; #[derive(Debug, PartialEq, Eq)] -#[cfg_attr(test,derive(EnumIter))] +#[cfg_attr(test, derive(EnumIter))] pub enum Message { - Hello { - file_size: u64 - }, + Hello { file_size: u64 }, ACK, NACK, - Send + Send, } - // Wire format impl Display for Message { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { @@ -28,7 +24,6 @@ impl Display for Message { } } - impl FromStr for Message { type Err = (); @@ -48,7 +43,7 @@ impl FromStr for Message { return Err(()); }; - // No arguments msgs + // No arguments msgs return match upper.as_str() { "ACK" => Ok(Message::ACK), "NACK" => Ok(Message::NACK), diff --git a/topics/p2p-transfer-protocol/src/protocol/mod.rs b/topics/p2p-transfer-protocol/src/protocol/mod.rs index f1cf2f6..022a5fb 100644 --- a/topics/p2p-transfer-protocol/src/protocol/mod.rs +++ b/topics/p2p-transfer-protocol/src/protocol/mod.rs @@ -1,37 +1,50 @@ use log::{debug, warn}; -pub mod message; pub mod connection_state; +pub mod message; use connection_state::ConnectionState; +#[derive(Debug)] pub struct StateMachine { state: ConnectionState, + pub(crate) file_size: Option, } impl StateMachine { pub fn new() -> Self { StateMachine { state: ConnectionState::Closed, + file_size: None, } } - pub fn transition(&mut self, new_state: &ConnectionState) { - debug!("Transitioning from {:?} to {:?}", self.state, new_state); + pub fn current_state(&self) -> &ConnectionState { + &self.state + } + + // TODO: Ideally, this would return a Result to indicate invalid transitions. + pub fn transition(&mut self, new_state: ConnectionState) { match (&self.state, new_state) { (ConnectionState::Closed, ConnectionState::Listening) => {} (ConnectionState::Listening, ConnectionState::HelloReceived) => {} (ConnectionState::HelloReceived, ConnectionState::NACKSent) => {} (ConnectionState::HelloReceived, ConnectionState::ACKSent) => {} - + (ConnectionState::ACKSent, ConnectionState::Established) => {} + (ConnectionState::Closed, ConnectionState::HelloSent) => {} (ConnectionState::HelloSent, ConnectionState::NACKReceived) => {} (ConnectionState::HelloSent, ConnectionState::ACKReceived) => {} (ConnectionState::ACKReceived, ConnectionState::Established) => {} - - - (s, _) => {warn!("Invalid state transition from {:?} to {:?}, ignoring.", s, new_state);} + + (s, _) => { + warn!( + "Invalid state transition from {:?} to {:?}, ignoring.", + s, new_state + ); + } }; - debug!("New state: {:?}", self.state); + debug!("Transitioning from {:?} to {:?}", self.state, new_state); + self.state = new_state; } } diff --git a/topics/p2p-transfer-protocol/src/server/mod.rs b/topics/p2p-transfer-protocol/src/server/mod.rs new file mode 100644 index 0000000..4a73768 --- /dev/null +++ b/topics/p2p-transfer-protocol/src/server/mod.rs @@ -0,0 +1,122 @@ +use log::{debug, error, info, warn}; +use std::fs::File; +use std::io::{self, BufRead, BufReader, Write, copy}; +use std::net::TcpStream; + +use crate::protocol::StateMachine; +use crate::protocol::connection_state::ConnectionState; +use crate::protocol::message::Message; + +fn on_hello(file_size: u64, state_machine: &mut StateMachine) -> Message { + info!("Client wants to send a file of size {} bytes", file_size); + state_machine.transition(ConnectionState::HelloReceived); + + state_machine.file_size = Some(file_size); + + // TODO: Add a basic check to handle the NACK case. + + // send ack + state_machine.transition(ConnectionState::ACKSent); + + Message::ACK +} + +fn on_message(msg: Message, state_machine: &mut StateMachine) -> Option { + return match msg { + Message::Hello { file_size } => Some(on_hello(file_size, state_machine)), + Message::Send => { + info!("Client is starting to send data."); + state_machine.transition(ConnectionState::Established); + None + } + + _ => { + warn!("Unexpected message from client: {:?}", msg); + None + } + }; +} + +fn send_message(stream: &mut TcpStream, msg: &Message) -> Result<(), io::Error> { + let msg_str = msg.to_string(); + stream.write_all(msg_str.as_bytes())?; + Ok(()) +} + +pub fn run_server(file: &mut File, stream: &mut TcpStream) -> Result<(), io::Error> { + let peer_addr = stream.peer_addr()?; + println!("New connection from {peer_addr}"); + + // Intialize the state machine + let mut state_machine = StateMachine::new(); + state_machine.transition(ConnectionState::Listening); + + // Start the negociation message loop. + message_loop(&mut state_machine, stream)?; + + // Now we're in Established state, receive the file data. + let expected_file_size = state_machine.file_size.unwrap_or(0); + + receive_file(file, stream, expected_file_size) +} + +fn receive_file( + file: &mut File, + stream: &mut TcpStream, + expected_file_size: u64, +) -> io::Result<()> { + + let mut sized_stream = io::Read::take(stream, expected_file_size); + + let bytes_copied = copy(&mut sized_stream, file)?; + file.flush()?; + + if bytes_copied != expected_file_size { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + format!( + "Expected {} bytes, but received {} bytes", + expected_file_size, bytes_copied + ), + )); + } + + + info!( + "File received successfully, {} bytes written.", + bytes_copied + ); + Ok(()) +} + +fn message_loop(state_machine: &mut StateMachine, stream: &mut TcpStream) -> Result<(), io::Error> { + let reader_stream = stream.try_clone().expect("Failed to clone stream"); + let mut writer_stream = stream; + let reader = BufReader::new(reader_stream); + + for line in reader.lines() { + let message = match line { + Ok(line) => { + debug!("Received: {}", line); + line.parse::() + } + Err(e) => { + error!("Error reading line: {}", e); + continue; + } + }; + + // Process message + // FIXME: Do not use unwrap, fix the result type mess. + match on_message(message.unwrap(), state_machine) { + Some(resp) => send_message(&mut writer_stream, &resp)?, + None => {}, // No response is needed + } + // If we are in Established state, stop the line-based logic and receive the data. + warn!("Current state: {:?}", state_machine.current_state()); + if let ConnectionState::Established = state_machine.current_state() { + break; + } + } + Ok(()) +} From 3106d22ec5c5b3a27a5a76f3a607e039eb24dd67 Mon Sep 17 00:00:00 2001 From: rrouviere Date: Thu, 30 Oct 2025 22:24:29 +0100 Subject: [PATCH 07/13] chore: Rework main.rs, remove extraneous warning from server code Signed-off-by: rrouviere --- topics/p2p-transfer-protocol/out_file | 4 -- topics/p2p-transfer-protocol/src/main.rs | 48 ++++++++----------- .../p2p-transfer-protocol/src/server/mod.rs | 1 - 3 files changed, 21 insertions(+), 32 deletions(-) delete mode 100644 topics/p2p-transfer-protocol/out_file diff --git a/topics/p2p-transfer-protocol/out_file b/topics/p2p-transfer-protocol/out_file deleted file mode 100644 index 437a81b..0000000 --- a/topics/p2p-transfer-protocol/out_file +++ /dev/null @@ -1,4 +0,0 @@ -12345678900 -0000000000 - - diff --git a/topics/p2p-transfer-protocol/src/main.rs b/topics/p2p-transfer-protocol/src/main.rs index 28a9aa0..a76e561 100644 --- a/topics/p2p-transfer-protocol/src/main.rs +++ b/topics/p2p-transfer-protocol/src/main.rs @@ -3,14 +3,11 @@ mod protocol; mod server; use std::{ - fs::{File, OpenOptions}, - net::{TcpListener, TcpStream}, - process::exit, - io, + fs::{File, OpenOptions}, io, net::{TcpListener, TcpStream}, process::exit }; use clap::{Parser, Subcommand}; -use log::{error, info}; +use log::{debug, error, info}; #[derive(Parser)] #[command( @@ -27,15 +24,15 @@ enum Commands { Listen { #[arg(long, default_value = "[::]:0")] // Listen on any inet & let the kernel pick a port. bind: String, - #[arg(long, default_value = "out_file")] - output: String, + #[arg(long, default_value = "output_file")] + output_file: String, }, Send { - #[arg(long, default_value = "out_file")] + #[arg(long)] file: String, - #[arg()] - remote_addr: String, + #[arg(long)] + server: String, }, } @@ -43,15 +40,18 @@ fn main() { // Initialize logger env_logger::init(); - let cli = Cli::parse(); // Note: if no command matches, clap automatically provides the help message & exits. match cli.command { - Commands::Listen { bind, output } => { + Commands::Listen { bind, output_file: output } => { let _ = server_mode(&bind, &output); } - Commands::Send { file, remote_addr } => { - client_mode(&file, &remote_addr); + Commands::Send { file, server: remote_addr } => { + + match client_mode(&file, &remote_addr) { + Ok(_) => info!("File sent successfully"), + Err(e) => error!("Failed to send file: {}", e), + } } } } @@ -98,7 +98,12 @@ fn server_mode(bind_addr: &str, file_path: &str) -> io::Result<()> { Ok(()) } -fn client_mode(file_path: &str, remote_addr: &str) { + +fn client_mode(file_path: &str, remote_addr: &str) -> io::Result<()> { + // Open input file + let file = File::open(file_path)?; + debug!("Opened file {}", file_path); + // Connect to server let stream = match TcpStream::connect(remote_addr) { Ok(s) => s, @@ -109,17 +114,6 @@ fn client_mode(file_path: &str, remote_addr: &str) { }; info!("Connected to {}", remote_addr); - // Open input file - let file = match File::open(file_path) { - Ok(f) => f, - Err(e) => { - error!("Failed to open file {}: {}", file_path, e); - exit(1); - } - }; - match client::run_client(file, stream) { - Ok(_) => info!("File sent successfully"), - Err(e) => error!("Failed to send file: {}", e), - } + client::run_client(file, stream) } diff --git a/topics/p2p-transfer-protocol/src/server/mod.rs b/topics/p2p-transfer-protocol/src/server/mod.rs index 4a73768..4a63711 100644 --- a/topics/p2p-transfer-protocol/src/server/mod.rs +++ b/topics/p2p-transfer-protocol/src/server/mod.rs @@ -113,7 +113,6 @@ fn message_loop(state_machine: &mut StateMachine, stream: &mut TcpStream) -> Res None => {}, // No response is needed } // If we are in Established state, stop the line-based logic and receive the data. - warn!("Current state: {:?}", state_machine.current_state()); if let ConnectionState::Established = state_machine.current_state() { break; } From 75eeeab856024d2d1b898d4e17bedf46274492e5 Mon Sep 17 00:00:00 2001 From: rrouviere Date: Thu, 30 Oct 2025 22:37:32 +0100 Subject: [PATCH 08/13] docs: Document the current state of the project Signed-off-by: rrouviere --- .../docs/architecture.md | 78 +++++++++---------- topics/p2p-transfer-protocol/output_file | 0 topics/p2p-transfer-protocol/src/main.rs | 26 ++++--- .../p2p-transfer-protocol/src/protocol/mod.rs | 4 +- .../p2p-transfer-protocol/src/server/mod.rs | 6 +- 5 files changed, 54 insertions(+), 60 deletions(-) create mode 100644 topics/p2p-transfer-protocol/output_file diff --git a/topics/p2p-transfer-protocol/docs/architecture.md b/topics/p2p-transfer-protocol/docs/architecture.md index fa2974f..caaf09d 100644 --- a/topics/p2p-transfer-protocol/docs/architecture.md +++ b/topics/p2p-transfer-protocol/docs/architecture.md @@ -1,8 +1,29 @@ -## -TODO: Project definition: What is it? What are the goals of the tool/project? +## p2p-transfert-protocol +A simple text-based protocol for file transfert. +This tool includes both a client and a single threaded server implementation. + +By default, the server listens on `[::]:0`. This can ## Usage -How can one use it? Give usage examples. + +Via cargo run: + +```bash +cargo run listen [--bind] [--output-file] +cargo run send <--server HOST> <--file FILE> +``` + +Or by building a binary: `cargo build --release` (`./target/release/p2p-transfert-protocol`) + +### Logging + +This project uses `env_logger` for logging. +For debug output (eg. the state machine's transition), set `RUST_LOG=debug`: + +```bash +RUST_LOG=debug cargo run listen [--bind] [--output-file] +RUST_LOG=debug cargo run send <--server HOST> <--file FILE> +``` ## Modules @@ -11,9 +32,10 @@ The project is structured around three main modules: - The "client" module - The generic protocol library, containing shared serialization/deserialization logic, as well as the FSM (Finite State Machine) driving the protocol's state. -Note: for simplicity's sake, we use "client" to refer to the active opener, that is, the peer that sends the `HELLO` message. +## Limitations + +### Separation of concerns and limitations of this protocol -## Seperations of concerns and limitations of this protocol This protocol is intended as an Layer 7/Application-level protocol. Hence, the underlying protocol stack is assumed : @@ -27,43 +49,13 @@ Hence, the underlying protocol stack is assumed : - As this protocol uses the same channel for data and control, it technically suffers from head-of-line blocking. - In this context this however is an non-issue, given that only one file per connection can be transmitted and that resuming is not supported. -This implementation uses TCP as its transport protocol, which provides reliablility and weak integrity guarantees, but no confidentiality nor protection against tempering of the payload. -Such a guarantee could be achieved by encapsulating this protocol in TLS. - -Another interesting possibility would be to use QUIC as the transport protocol, which could allow connection reuse, multiplexing (using streams), as well as lower latency (by merging the transport and crypto establishement as a single step). -Latency could be futher reduced by sending the HELLO packet as a 0-RTT packet on futher connections. - - -## Protocol description - -### States -#### Initial states -- LISTENING (passive opener) -- HELLO_SENT - -#### Transition states -- HELLO_RECEIVED -- ACK_SENT -- ACK_RECEIVED -- ESTABLISHED - -#### Final states -- NACK_RECEIVED -- (Implicit teardown) - -## Messages -- HELLO -- ACK -- NACK -- SEND - -## Messages -- HELLO: For the sender to offer a file to the receiver. It takes a file size argument. -- ACK: For the receiver to tell the sender it is ready to receive a proposed file. -- NACK: For the receiver to reject a proposed file. -- SEND: Send, for the sender to actually send a file. It also takes a file size argument, that must match the HELLO offer. +### Limitation of this implementation +- This implementation uses TCP as its transport protocol, which provides reliablility and weak integrity guarantees, but no confidentiality nor protection against tempering of the payload. +Such a guarantee could be achieved by encapsulating this protocol in TLS. +- The serialization/deserialization logic could be simplified using serde +- The server is singlethreaded. -## Limitations & futher work -- The serialization/deserialization logic could be simplified using serde. -- \ No newline at end of file +## Potential work +An interesting possibility would be to use QUIC as the transport protocol, which could allow connection reuse, multiplexing (using streams), as well as lower latency (by merging the transport and crypto establishement as a single step). +Latency could be futher reduced by sending the HELLO packet as a 0-RTT packet on futher connections. \ No newline at end of file diff --git a/topics/p2p-transfer-protocol/output_file b/topics/p2p-transfer-protocol/output_file new file mode 100644 index 0000000..e69de29 diff --git a/topics/p2p-transfer-protocol/src/main.rs b/topics/p2p-transfer-protocol/src/main.rs index a76e561..415bc8d 100644 --- a/topics/p2p-transfer-protocol/src/main.rs +++ b/topics/p2p-transfer-protocol/src/main.rs @@ -3,7 +3,10 @@ mod protocol; mod server; use std::{ - fs::{File, OpenOptions}, io, net::{TcpListener, TcpStream}, process::exit + fs::{File, OpenOptions}, + io, + net::{TcpListener, TcpStream}, + process::exit, }; use clap::{Parser, Subcommand}; @@ -43,16 +46,19 @@ fn main() { let cli = Cli::parse(); // Note: if no command matches, clap automatically provides the help message & exits. match cli.command { - Commands::Listen { bind, output_file: output } => { + Commands::Listen { + bind, + output_file: output, + } => { let _ = server_mode(&bind, &output); } - Commands::Send { file, server: remote_addr } => { - - match client_mode(&file, &remote_addr) { - Ok(_) => info!("File sent successfully"), - Err(e) => error!("Failed to send file: {}", e), - } - } + Commands::Send { + file, + server: remote_addr, + } => match client_mode(&file, &remote_addr) { + Ok(_) => info!("File sent successfully"), + Err(e) => error!("Failed to send file: {}", e), + }, } } @@ -98,7 +104,6 @@ fn server_mode(bind_addr: &str, file_path: &str) -> io::Result<()> { Ok(()) } - fn client_mode(file_path: &str, remote_addr: &str) -> io::Result<()> { // Open input file let file = File::open(file_path)?; @@ -114,6 +119,5 @@ fn client_mode(file_path: &str, remote_addr: &str) -> io::Result<()> { }; info!("Connected to {}", remote_addr); - client::run_client(file, stream) } diff --git a/topics/p2p-transfer-protocol/src/protocol/mod.rs b/topics/p2p-transfer-protocol/src/protocol/mod.rs index 022a5fb..0451c61 100644 --- a/topics/p2p-transfer-protocol/src/protocol/mod.rs +++ b/topics/p2p-transfer-protocol/src/protocol/mod.rs @@ -31,12 +31,12 @@ impl StateMachine { (ConnectionState::HelloReceived, ConnectionState::NACKSent) => {} (ConnectionState::HelloReceived, ConnectionState::ACKSent) => {} (ConnectionState::ACKSent, ConnectionState::Established) => {} - + (ConnectionState::Closed, ConnectionState::HelloSent) => {} (ConnectionState::HelloSent, ConnectionState::NACKReceived) => {} (ConnectionState::HelloSent, ConnectionState::ACKReceived) => {} (ConnectionState::ACKReceived, ConnectionState::Established) => {} - + (s, _) => { warn!( "Invalid state transition from {:?} to {:?}, ignoring.", diff --git a/topics/p2p-transfer-protocol/src/server/mod.rs b/topics/p2p-transfer-protocol/src/server/mod.rs index 4a63711..273350c 100644 --- a/topics/p2p-transfer-protocol/src/server/mod.rs +++ b/topics/p2p-transfer-protocol/src/server/mod.rs @@ -65,9 +65,8 @@ fn receive_file( stream: &mut TcpStream, expected_file_size: u64, ) -> io::Result<()> { - let mut sized_stream = io::Read::take(stream, expected_file_size); - + let bytes_copied = copy(&mut sized_stream, file)?; file.flush()?; @@ -81,7 +80,6 @@ fn receive_file( )); } - info!( "File received successfully, {} bytes written.", bytes_copied @@ -110,7 +108,7 @@ fn message_loop(state_machine: &mut StateMachine, stream: &mut TcpStream) -> Res // FIXME: Do not use unwrap, fix the result type mess. match on_message(message.unwrap(), state_machine) { Some(resp) => send_message(&mut writer_stream, &resp)?, - None => {}, // No response is needed + None => {} // No response is needed } // If we are in Established state, stop the line-based logic and receive the data. if let ConnectionState::Established = state_machine.current_state() { From ce19c288cd342a74a1d8bec91e7cb0a7afd2fac5 Mon Sep 17 00:00:00 2001 From: rrouviere Date: Thu, 30 Oct 2025 22:44:55 +0100 Subject: [PATCH 09/13] fix: Fix typo Signed-off-by: rrouviere --- topics/p2p-transfer-protocol/docs/architecture.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/topics/p2p-transfer-protocol/docs/architecture.md b/topics/p2p-transfer-protocol/docs/architecture.md index caaf09d..a618cd8 100644 --- a/topics/p2p-transfer-protocol/docs/architecture.md +++ b/topics/p2p-transfer-protocol/docs/architecture.md @@ -1,5 +1,5 @@ ## p2p-transfert-protocol -A simple text-based protocol for file transfert. +A simple text-based protocol for file transfer. This tool includes both a client and a single threaded server implementation. By default, the server listens on `[::]:0`. This can From 67fc5adf0cce31ed36f65e2dcd2da47c32b8ec68 Mon Sep 17 00:00:00 2001 From: rrouviere Date: Thu, 30 Oct 2025 22:45:08 +0100 Subject: [PATCH 10/13] chore(test): Remove strum dependancy (planned final state checking with EnumIter) Signed-off-by: rrouviere --- topics/p2p-transfer-protocol/Cargo.lock | 22 ------------------- topics/p2p-transfer-protocol/Cargo.toml | 1 - .../src/protocol/message.rs | 3 --- 3 files changed, 26 deletions(-) diff --git a/topics/p2p-transfer-protocol/Cargo.lock b/topics/p2p-transfer-protocol/Cargo.lock index 2184a2d..7eec480 100644 --- a/topics/p2p-transfer-protocol/Cargo.lock +++ b/topics/p2p-transfer-protocol/Cargo.lock @@ -191,7 +191,6 @@ dependencies = [ "clap", "env_logger", "log", - "strum", ] [[package]] @@ -291,27 +290,6 @@ version = "0.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7da8b5736845d9f2fcb837ea5d9e2628564b3b043a70948a3f0b778838c5fb4f" -[[package]] -name = "strum" -version = "0.27.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "af23d6f6c1a224baef9d3f61e287d2761385a5b88fdab4eb4c6f11aeb54c4bcf" -dependencies = [ - "strum_macros", -] - -[[package]] -name = "strum_macros" -version = "0.27.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7695ce3845ea4b33927c055a39dc438a45b059f7c1b3d91d38d10355fb8cbca7" -dependencies = [ - "heck", - "proc-macro2", - "quote", - "syn", -] - [[package]] name = "syn" version = "2.0.108" diff --git a/topics/p2p-transfer-protocol/Cargo.toml b/topics/p2p-transfer-protocol/Cargo.toml index dfb540e..94a5180 100644 --- a/topics/p2p-transfer-protocol/Cargo.toml +++ b/topics/p2p-transfer-protocol/Cargo.toml @@ -6,5 +6,4 @@ edition = "2024" [dependencies] clap = { version = "4.5.51", features = ["derive"] } log = "0.4.28" -strum = { version = "0.27.2", features = ["derive"] } env_logger = "0.11.8" \ No newline at end of file diff --git a/topics/p2p-transfer-protocol/src/protocol/message.rs b/topics/p2p-transfer-protocol/src/protocol/message.rs index 1b0e92a..2ce9efe 100644 --- a/topics/p2p-transfer-protocol/src/protocol/message.rs +++ b/topics/p2p-transfer-protocol/src/protocol/message.rs @@ -1,10 +1,7 @@ use log::error; use std::{fmt::Display, str::FromStr}; -#[cfg(test)] -use strum::EnumIter; #[derive(Debug, PartialEq, Eq)] -#[cfg_attr(test, derive(EnumIter))] pub enum Message { Hello { file_size: u64 }, ACK, From 5984ce687a5ba4f8787c6762f207e49a17a717d6 Mon Sep 17 00:00:00 2001 From: rrouviere Date: Thu, 30 Oct 2025 22:49:46 +0100 Subject: [PATCH 11/13] lint: Fix clippy warnings Signed-off-by: rrouviere --- .../p2p-transfer-protocol/src/client/mod.rs | 4 +-- topics/p2p-transfer-protocol/src/main.rs | 3 ++- .../src/protocol/message.rs | 25 +++++++++---------- .../p2p-transfer-protocol/src/server/mod.rs | 14 +++++------ 4 files changed, 22 insertions(+), 24 deletions(-) diff --git a/topics/p2p-transfer-protocol/src/client/mod.rs b/topics/p2p-transfer-protocol/src/client/mod.rs index 1491446..66e7dc6 100644 --- a/topics/p2p-transfer-protocol/src/client/mod.rs +++ b/topics/p2p-transfer-protocol/src/client/mod.rs @@ -26,11 +26,11 @@ pub fn run_client(mut file: File, mut stream: TcpStream) -> io::Result<()> { reader.read_line(&mut response)?; let response = response.trim(); match response.parse::() { - Ok(Message::ACK) => { + Ok(Message::Ack) => { debug!("Received ACK"); sm.transition(ConnectionState::ACKReceived); } - Ok(Message::NACK) => { + Ok(Message::Nack) => { debug!("Received NACK"); sm.transition(ConnectionState::NACKReceived); return Ok(()); diff --git a/topics/p2p-transfer-protocol/src/main.rs b/topics/p2p-transfer-protocol/src/main.rs index 415bc8d..1578dd6 100644 --- a/topics/p2p-transfer-protocol/src/main.rs +++ b/topics/p2p-transfer-protocol/src/main.rs @@ -65,7 +65,7 @@ fn main() { fn server_mode(bind_addr: &str, file_path: &str) -> io::Result<()> { // Start server // Bind to port: bind_addr - let listener = match TcpListener::bind(&bind_addr) { + let listener = match TcpListener::bind(bind_addr) { Ok(l) => l, Err(e) => { error!("Failed to bind to {}: {}", bind_addr, e); @@ -76,6 +76,7 @@ fn server_mode(bind_addr: &str, file_path: &str) -> io::Result<()> { // Open or create output file: let mut file = OpenOptions::new() + .truncate(true) .write(true) .create(true) .open(file_path)?; diff --git a/topics/p2p-transfer-protocol/src/protocol/message.rs b/topics/p2p-transfer-protocol/src/protocol/message.rs index 2ce9efe..286e26c 100644 --- a/topics/p2p-transfer-protocol/src/protocol/message.rs +++ b/topics/p2p-transfer-protocol/src/protocol/message.rs @@ -4,8 +4,8 @@ use std::{fmt::Display, str::FromStr}; #[derive(Debug, PartialEq, Eq)] pub enum Message { Hello { file_size: u64 }, - ACK, - NACK, + Ack, + Nack, Send, } @@ -13,10 +13,10 @@ pub enum Message { impl Display for Message { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { - Message::Hello { file_size } => write!(f, "HELLO {file_size}\n"), - Message::ACK => write!(f, "ACK\n"), - Message::NACK => write!(f, "NACK\n"), - Message::Send => write!(f, "SEND\n"), + Message::Hello { file_size } => writeln!(f, "HELLO {file_size}"), + Message::Ack => writeln!(f, "ACK"), + Message::Nack => writeln!(f, "NACK"), + Message::Send => writeln!(f, "SEND"), } } } @@ -31,22 +31,21 @@ impl FromStr for Message { if upper.starts_with("HELLO") { let parts: Vec<&str> = s.split_whitespace().collect(); - if parts.len() == 2 { - if let Ok(file_size) = parts[1].parse::() { + if parts.len() == 2 + && let Ok(file_size) = parts[1].parse::() { return Ok(Message::Hello { file_size }); } - } error!("Invalid HELLO message"); return Err(()); }; // No arguments msgs - return match upper.as_str() { - "ACK" => Ok(Message::ACK), - "NACK" => Ok(Message::NACK), + match upper.as_str() { + "ACK" => Ok(Message::Ack), + "NACK" => Ok(Message::Nack), "SEND" => Ok(Message::Send), // Unknown message _ => Err(()), - }; + } } } diff --git a/topics/p2p-transfer-protocol/src/server/mod.rs b/topics/p2p-transfer-protocol/src/server/mod.rs index 273350c..970eb66 100644 --- a/topics/p2p-transfer-protocol/src/server/mod.rs +++ b/topics/p2p-transfer-protocol/src/server/mod.rs @@ -18,11 +18,11 @@ fn on_hello(file_size: u64, state_machine: &mut StateMachine) -> Message { // send ack state_machine.transition(ConnectionState::ACKSent); - Message::ACK + Message::Ack } fn on_message(msg: Message, state_machine: &mut StateMachine) -> Option { - return match msg { + match msg { Message::Hello { file_size } => Some(on_hello(file_size, state_machine)), Message::Send => { info!("Client is starting to send data."); @@ -34,7 +34,7 @@ fn on_message(msg: Message, state_machine: &mut StateMachine) -> Option warn!("Unexpected message from client: {:?}", msg); None } - }; + } } fn send_message(stream: &mut TcpStream, msg: &Message) -> Result<(), io::Error> { @@ -89,7 +89,7 @@ fn receive_file( fn message_loop(state_machine: &mut StateMachine, stream: &mut TcpStream) -> Result<(), io::Error> { let reader_stream = stream.try_clone().expect("Failed to clone stream"); - let mut writer_stream = stream; + let writer_stream = stream; let reader = BufReader::new(reader_stream); for line in reader.lines() { @@ -106,10 +106,8 @@ fn message_loop(state_machine: &mut StateMachine, stream: &mut TcpStream) -> Res // Process message // FIXME: Do not use unwrap, fix the result type mess. - match on_message(message.unwrap(), state_machine) { - Some(resp) => send_message(&mut writer_stream, &resp)?, - None => {} // No response is needed - } + if let Some(resp) = on_message(message.unwrap(), state_machine) { send_message(writer_stream, &resp)? } + // If we are in Established state, stop the line-based logic and receive the data. if let ConnectionState::Established = state_machine.current_state() { break; From 698a37519815385bc740b71c6bc48f63e6c3bac6 Mon Sep 17 00:00:00 2001 From: rrouviere Date: Thu, 30 Oct 2025 23:24:46 +0100 Subject: [PATCH 12/13] chore: Refactor and modularize server code Signed-off-by: rrouviere --- topics/p2p-transfer-protocol/src/main.rs | 1 + .../src/protocol/message.rs | 7 +- .../src/server/connection.rs | 57 +++++++++ .../src/server/events.rs | 71 +++++++++++ .../p2p-transfer-protocol/src/server/mod.rs | 112 +----------------- 5 files changed, 138 insertions(+), 110 deletions(-) create mode 100644 topics/p2p-transfer-protocol/src/server/connection.rs create mode 100644 topics/p2p-transfer-protocol/src/server/events.rs diff --git a/topics/p2p-transfer-protocol/src/main.rs b/topics/p2p-transfer-protocol/src/main.rs index 1578dd6..db68b8e 100644 --- a/topics/p2p-transfer-protocol/src/main.rs +++ b/topics/p2p-transfer-protocol/src/main.rs @@ -87,6 +87,7 @@ fn server_mode(bind_addr: &str, file_path: &str) -> io::Result<()> { ); // Wait for client + // FIXME: We handle multiple clients, but since the output file is shared among them, this is unhelpful. for stream in listener.incoming() { match stream { Ok(mut s) => { diff --git a/topics/p2p-transfer-protocol/src/protocol/message.rs b/topics/p2p-transfer-protocol/src/protocol/message.rs index 286e26c..9ce59ed 100644 --- a/topics/p2p-transfer-protocol/src/protocol/message.rs +++ b/topics/p2p-transfer-protocol/src/protocol/message.rs @@ -32,9 +32,10 @@ impl FromStr for Message { if upper.starts_with("HELLO") { let parts: Vec<&str> = s.split_whitespace().collect(); if parts.len() == 2 - && let Ok(file_size) = parts[1].parse::() { - return Ok(Message::Hello { file_size }); - } + && let Ok(file_size) = parts[1].parse::() + { + return Ok(Message::Hello { file_size }); + } error!("Invalid HELLO message"); return Err(()); }; diff --git a/topics/p2p-transfer-protocol/src/server/connection.rs b/topics/p2p-transfer-protocol/src/server/connection.rs new file mode 100644 index 0000000..af36fb2 --- /dev/null +++ b/topics/p2p-transfer-protocol/src/server/connection.rs @@ -0,0 +1,57 @@ +use log::info; +use std::fs::File; +use std::io::{self, Write, copy}; +use std::net::TcpStream; + +use crate::protocol::StateMachine; +use crate::protocol::connection_state::ConnectionState; +use crate::protocol::message::Message; + +use crate::server::events::handle_message_loop; + +pub(crate) fn handle_connection(file: &mut File, stream: &mut TcpStream) -> Result<(), io::Error> { + // Intialize the state machine + let mut state_machine = StateMachine::new(); + state_machine.transition(ConnectionState::Listening); + + // Start the negociation message loop. + handle_message_loop(&mut state_machine, stream)?; + + // Now we're in Established state, receive the file data. + let expected_file_size = state_machine.file_size.unwrap_or(0); + + receive_file(file, stream, expected_file_size) +} + +pub fn send_message(stream: &mut TcpStream, msg: &Message) -> Result<(), io::Error> { + let msg_str = msg.to_string(); + stream.write_all(msg_str.as_bytes())?; + Ok(()) +} + +pub fn receive_file( + file: &mut File, + stream: &mut TcpStream, + expected_file_size: u64, +) -> io::Result<()> { + let mut sized_stream = io::Read::take(stream, expected_file_size); + + let bytes_copied = copy(&mut sized_stream, file)?; + file.flush()?; + + if bytes_copied != expected_file_size { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + format!( + "Expected {} bytes, but received {} bytes", + expected_file_size, bytes_copied + ), + )); + } + + info!( + "File received successfully, {} bytes written.", + bytes_copied + ); + Ok(()) +} diff --git a/topics/p2p-transfer-protocol/src/server/events.rs b/topics/p2p-transfer-protocol/src/server/events.rs new file mode 100644 index 0000000..997a404 --- /dev/null +++ b/topics/p2p-transfer-protocol/src/server/events.rs @@ -0,0 +1,71 @@ +use log::{debug, error, info, warn}; +use std::io::{self, BufRead, BufReader}; +use std::net::TcpStream; + +use crate::protocol::StateMachine; +use crate::protocol::connection_state::ConnectionState; +use crate::protocol::message::Message; +use crate::server::connection::send_message; + +fn on_hello(file_size: u64, state_machine: &mut StateMachine) -> Message { + info!("Client wants to send a file of size {} bytes", file_size); + state_machine.transition(ConnectionState::HelloReceived); + + state_machine.file_size = Some(file_size); + + // TODO: Add a basic chec k to handle the NACK case. + + // send ack + state_machine.transition(ConnectionState::ACKSent); + + Message::Ack +} + +fn on_message(msg: Message, state_machine: &mut StateMachine) -> Option { + match msg { + Message::Hello { file_size } => Some(on_hello(file_size, state_machine)), + Message::Send => { + info!("Client is starting to send data."); + state_machine.transition(ConnectionState::Established); + None + } + + _ => { + warn!("Unexpected message from client: {:?}", msg); + None + } + } +} +pub(crate) fn handle_message_loop( + state_machine: &mut StateMachine, + stream: &mut TcpStream, +) -> Result<(), io::Error> { + let reader_stream = stream.try_clone().expect("Failed to clone stream"); + let writer_stream = stream; + let reader = BufReader::new(reader_stream); + + for line in reader.lines() { + let message = match line { + Ok(line) => { + debug!("Received: {}", line); + line.parse::() + } + Err(e) => { + error!("Error reading line: {}", e); + continue; + } + }; + + // Process message + // FIXME: Do not use unwrap, fix the result type mess. + if let Some(resp) = on_message(message.unwrap(), state_machine) { + send_message(writer_stream, &resp)? + } + + // If we are in Established state, stop the line-based logic and receive the data. + if let ConnectionState::Established = state_machine.current_state() { + break; + } + } + Ok(()) +} diff --git a/topics/p2p-transfer-protocol/src/server/mod.rs b/topics/p2p-transfer-protocol/src/server/mod.rs index 970eb66..7d3fd5d 100644 --- a/topics/p2p-transfer-protocol/src/server/mod.rs +++ b/topics/p2p-transfer-protocol/src/server/mod.rs @@ -1,117 +1,15 @@ -use log::{debug, error, info, warn}; use std::fs::File; -use std::io::{self, BufRead, BufReader, Write, copy}; +use std::io; use std::net::TcpStream; -use crate::protocol::StateMachine; -use crate::protocol::connection_state::ConnectionState; -use crate::protocol::message::Message; +use crate::server::connection::handle_connection; -fn on_hello(file_size: u64, state_machine: &mut StateMachine) -> Message { - info!("Client wants to send a file of size {} bytes", file_size); - state_machine.transition(ConnectionState::HelloReceived); - - state_machine.file_size = Some(file_size); - - // TODO: Add a basic check to handle the NACK case. - - // send ack - state_machine.transition(ConnectionState::ACKSent); - - Message::Ack -} - -fn on_message(msg: Message, state_machine: &mut StateMachine) -> Option { - match msg { - Message::Hello { file_size } => Some(on_hello(file_size, state_machine)), - Message::Send => { - info!("Client is starting to send data."); - state_machine.transition(ConnectionState::Established); - None - } - - _ => { - warn!("Unexpected message from client: {:?}", msg); - None - } - } -} - -fn send_message(stream: &mut TcpStream, msg: &Message) -> Result<(), io::Error> { - let msg_str = msg.to_string(); - stream.write_all(msg_str.as_bytes())?; - Ok(()) -} +mod connection; +mod events; pub fn run_server(file: &mut File, stream: &mut TcpStream) -> Result<(), io::Error> { let peer_addr = stream.peer_addr()?; println!("New connection from {peer_addr}"); - // Intialize the state machine - let mut state_machine = StateMachine::new(); - state_machine.transition(ConnectionState::Listening); - - // Start the negociation message loop. - message_loop(&mut state_machine, stream)?; - - // Now we're in Established state, receive the file data. - let expected_file_size = state_machine.file_size.unwrap_or(0); - - receive_file(file, stream, expected_file_size) -} - -fn receive_file( - file: &mut File, - stream: &mut TcpStream, - expected_file_size: u64, -) -> io::Result<()> { - let mut sized_stream = io::Read::take(stream, expected_file_size); - - let bytes_copied = copy(&mut sized_stream, file)?; - file.flush()?; - - if bytes_copied != expected_file_size { - return Err(io::Error::new( - io::ErrorKind::InvalidData, - format!( - "Expected {} bytes, but received {} bytes", - expected_file_size, bytes_copied - ), - )); - } - - info!( - "File received successfully, {} bytes written.", - bytes_copied - ); - Ok(()) -} - -fn message_loop(state_machine: &mut StateMachine, stream: &mut TcpStream) -> Result<(), io::Error> { - let reader_stream = stream.try_clone().expect("Failed to clone stream"); - let writer_stream = stream; - let reader = BufReader::new(reader_stream); - - for line in reader.lines() { - let message = match line { - Ok(line) => { - debug!("Received: {}", line); - line.parse::() - } - Err(e) => { - error!("Error reading line: {}", e); - continue; - } - }; - - // Process message - // FIXME: Do not use unwrap, fix the result type mess. - if let Some(resp) = on_message(message.unwrap(), state_machine) { send_message(writer_stream, &resp)? } - - // If we are in Established state, stop the line-based logic and receive the data. - if let ConnectionState::Established = state_machine.current_state() { - break; - } - } - Ok(()) + handle_connection(file, stream) } From 894965b7f7fee2475ce9f0ac0ea1a9f85883fae2 Mon Sep 17 00:00:00 2001 From: rrouviere Date: Thu, 30 Oct 2025 23:55:39 +0100 Subject: [PATCH 13/13] fix(logging): Set default loglevel to INFO Signed-off-by: rrouviere --- topics/p2p-transfer-protocol/src/main.rs | 8 +++++--- topics/p2p-transfer-protocol/src/server/mod.rs | 4 +++- 2 files changed, 8 insertions(+), 4 deletions(-) diff --git a/topics/p2p-transfer-protocol/src/main.rs b/topics/p2p-transfer-protocol/src/main.rs index db68b8e..f67a72b 100644 --- a/topics/p2p-transfer-protocol/src/main.rs +++ b/topics/p2p-transfer-protocol/src/main.rs @@ -40,8 +40,10 @@ enum Commands { } fn main() { - // Initialize logger - env_logger::init(); + // Initialize logger to INFO level by default + env_logger::init_from_env( + env_logger::Env::default().filter_or(env_logger::DEFAULT_FILTER_ENV, "info"), + ); let cli = Cli::parse(); // Note: if no command matches, clap automatically provides the help message & exits. @@ -81,7 +83,7 @@ fn server_mode(bind_addr: &str, file_path: &str) -> io::Result<()> { .create(true) .open(file_path)?; - println!( + info!( "Server listening on {:?}. Saving to {}", listener, file_path ); diff --git a/topics/p2p-transfer-protocol/src/server/mod.rs b/topics/p2p-transfer-protocol/src/server/mod.rs index 7d3fd5d..51c951c 100644 --- a/topics/p2p-transfer-protocol/src/server/mod.rs +++ b/topics/p2p-transfer-protocol/src/server/mod.rs @@ -2,6 +2,8 @@ use std::fs::File; use std::io; use std::net::TcpStream; +use log::info; + use crate::server::connection::handle_connection; mod connection; @@ -9,7 +11,7 @@ mod events; pub fn run_server(file: &mut File, stream: &mut TcpStream) -> Result<(), io::Error> { let peer_addr = stream.peer_addr()?; - println!("New connection from {peer_addr}"); + info!("New connection from {peer_addr}"); handle_connection(file, stream) }