From ba3ac088dd1ec9fa27b7322062ef3422e1e28f16 Mon Sep 17 00:00:00 2001 From: Anton Tarasenko Date: Thu, 22 Jul 2021 18:03:48 +0700 Subject: [PATCH 1/5] Refactor numeric constants for `MessageReader` --- src/link/mod.rs | 74 ++++++++++++++++++++++++++++--------------------- 1 file changed, 43 insertions(+), 31 deletions(-) diff --git a/src/link/mod.rs b/src/link/mod.rs index 0a317a6..e54c9cd 100644 --- a/src/link/mod.rs +++ b/src/link/mod.rs @@ -4,11 +4,21 @@ use std::str; extern crate custom_error; use custom_error::custom_error; -const RECEIVED_FIELD_SIZE: usize = 4; -const LENGTH_FIELD_SIZE: usize = 4; -const HEAD_RECEIVED: u8 = 85; -const HEAD_MESSAGE: u8 = 42; -const MAX_MESSAGE_LENGTH: usize = 0xA00000; +// Defines how many bytes is used to encode "AMOUNT" field in the response from ue-server about +// amount of bytes it received since the last update +const UE_RECEIVED_FIELD_SIZE: usize = 4; +// Defines how many bytes is used to encode "LENGTH" field, describing length of +// next JSON message from ue-server +const UE_LENGTH_FIELD_SIZE: usize = 4; +// Value indicating that next byte sequence from ue-server reports amount of bytes received by +// that server so far. Value itself is arbitrary. + +const HEAD_UE_RECEIVED: u8 = 85; +// Value indicating that next byte sequence from ue-server contains JSON message. +// Value itself is arbitrary. +const HEAD_UE_MESSAGE: u8 = 42; +// Maximum allowed size of JSON message sent from ue-server. +const MAX_UE_MESSAGE_LENGTH: usize = 25 * 1024 * 1024; custom_error! { pub ReadingStreamError InvalidHead{input: u8} = "Invalid byte used as a HEAD: {input}", @@ -37,8 +47,8 @@ pub struct MessageReader { /// For converting byte stream expected to be generated by Acedia mod from the game server into /// actual messages. Expected format is a sequence of either: -/// [HEAD_RECEIVED: 1 byte] [amount of bytes received by game server since last update: 4 bytes] -/// [HEAD_MESSAGE: 1 byte] [length of the message: 4 bytes] [utf8-encoded string: ??? bytes] +/// [HEAD_UE_RECEIVED: 1 byte] [amount of bytes received by game server since last update: 4 bytes] +/// [HEAD_UE_MESSAGE: 1 byte] [length of the message: 4 bytes] [utf8-encoded string: ??? bytes] /// On any invalid input enters a failure state (can be checked by `is_broken()`) and /// never recovers from it. /// Use either `push_byte()` or `push()` to input byte stream from game server and `pop()` to @@ -63,9 +73,9 @@ impl MessageReader { } match &self.reading_state { ReadingState::Head => { - if input == HEAD_RECEIVED { + if input == HEAD_UE_RECEIVED { self.reading_state = ReadingState::ReceivedBytes; - } else if input == HEAD_MESSAGE { + } else if input == HEAD_UE_MESSAGE { self.reading_state = ReadingState::Length; } else { self.is_broken = true; @@ -76,7 +86,7 @@ impl MessageReader { self.next_received_bytes = self.next_received_bytes << 8; self.next_received_bytes += input as u32; self.read_bytes += 1; - if self.read_bytes >= RECEIVED_FIELD_SIZE { + if self.read_bytes >= UE_RECEIVED_FIELD_SIZE { self.received_bytes += self.next_received_bytes as u64; self.next_received_bytes = 0; self.read_bytes = 0; @@ -87,10 +97,10 @@ impl MessageReader { self.current_message_length = self.current_message_length << 8; self.current_message_length += input as usize; self.read_bytes += 1; - if self.read_bytes >= LENGTH_FIELD_SIZE { + if self.read_bytes >= UE_LENGTH_FIELD_SIZE { self.read_bytes = 0; self.reading_state = ReadingState::Payload; - if self.current_message_length > MAX_MESSAGE_LENGTH { + if self.current_message_length > MAX_UE_MESSAGE_LENGTH { self.is_broken = true; return Err(ReadingStreamError::MessageTooLong { length: self.current_message_length, @@ -143,7 +153,7 @@ impl MessageReader { #[test] fn message_push_byte() { let mut reader = MessageReader::new(); - reader.push_byte(HEAD_MESSAGE).unwrap(); + reader.push_byte(HEAD_UE_MESSAGE).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); @@ -161,7 +171,7 @@ fn message_push_byte() { reader.push_byte(0x6c).unwrap(); // l reader.push_byte(0x64).unwrap(); // d reader.push_byte(0x21).unwrap(); // - reader.push_byte(HEAD_MESSAGE).unwrap(); + reader.push_byte(HEAD_UE_MESSAGE).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); @@ -177,19 +187,19 @@ fn message_push_byte() { #[test] fn received_push_byte() { let mut reader = MessageReader::new(); - reader.push_byte(HEAD_RECEIVED).unwrap(); + reader.push_byte(HEAD_UE_RECEIVED).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(243).unwrap(); assert_eq!(reader.received_bytes(), 243); - reader.push_byte(HEAD_RECEIVED).unwrap(); + reader.push_byte(HEAD_UE_RECEIVED).unwrap(); reader.push_byte(65).unwrap(); reader.push_byte(25).unwrap(); reader.push_byte(178).unwrap(); reader.push_byte(4).unwrap(); assert_eq!(reader.received_bytes(), 1092203255); - reader.push_byte(HEAD_RECEIVED).unwrap(); + reader.push_byte(HEAD_UE_RECEIVED).unwrap(); reader.push_byte(231).unwrap(); reader.push_byte(34).unwrap(); reader.push_byte(154).unwrap(); @@ -199,12 +209,12 @@ fn received_push_byte() { #[test] fn mixed_push_byte() { let mut reader = MessageReader::new(); - reader.push_byte(HEAD_RECEIVED).unwrap(); + reader.push_byte(HEAD_UE_RECEIVED).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(243).unwrap(); - reader.push_byte(HEAD_MESSAGE).unwrap(); + reader.push_byte(HEAD_UE_MESSAGE).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); @@ -212,7 +222,7 @@ fn mixed_push_byte() { reader.push_byte(0x59).unwrap(); // Y reader.push_byte(0x6f).unwrap(); // o reader.push_byte(0x21).unwrap(); // - reader.push_byte(HEAD_RECEIVED).unwrap(); + reader.push_byte(HEAD_UE_RECEIVED).unwrap(); reader.push_byte(65).unwrap(); reader.push_byte(25).unwrap(); reader.push_byte(178).unwrap(); @@ -227,12 +237,12 @@ fn pushing_many_bytes_at_once() { let mut reader = MessageReader::new(); reader .push(&[ - HEAD_RECEIVED, + HEAD_UE_RECEIVED, 0, 0, 0, 243, - HEAD_MESSAGE, + HEAD_UE_MESSAGE, 0, 0, 0, @@ -240,7 +250,7 @@ fn pushing_many_bytes_at_once() { 0x59, // Y 0x6f, // o 0x21, // - HEAD_RECEIVED, + HEAD_UE_RECEIVED, 65, 25, 178, @@ -255,7 +265,7 @@ fn pushing_many_bytes_at_once() { #[test] fn generates_error_invalid_head() { let mut reader = MessageReader::new(); - reader.push_byte(HEAD_RECEIVED).unwrap(); + reader.push_byte(HEAD_UE_RECEIVED).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); @@ -270,14 +280,16 @@ fn generates_error_invalid_head() { #[test] fn generates_error_message_too_long() { let mut reader = MessageReader::new(); - reader.push_byte(HEAD_MESSAGE).unwrap(); - // Max length right now is `0xA00000` - reader.push_byte(0).unwrap(); - reader.push_byte(0xA0).unwrap(); - reader.push_byte(0).unwrap(); + let huge_length = MAX_UE_MESSAGE_LENGTH + 1; + let bytes = (huge_length as u32).to_be_bytes(); + + reader.push_byte(HEAD_UE_MESSAGE).unwrap(); + reader.push_byte(bytes[0]).unwrap(); + reader.push_byte(bytes[1]).unwrap(); + reader.push_byte(bytes[2]).unwrap(); assert!(!reader.is_broken()); reader - .push_byte(1) + .push_byte(bytes[3]) .expect_err("Testing failing on exceeding allowed message length"); assert!(reader.is_broken()); } @@ -285,7 +297,7 @@ fn generates_error_message_too_long() { #[test] fn generates_error_invalid_unicode() { let mut reader = MessageReader::new(); - reader.push_byte(HEAD_MESSAGE).unwrap(); + reader.push_byte(HEAD_UE_MESSAGE).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); From b846fcf55b0477a8642f5f985b9fe4ee4ea2c801 Mon Sep 17 00:00:00 2001 From: Anton Tarasenko Date: Thu, 22 Jul 2021 18:11:44 +0700 Subject: [PATCH 2/5] Fix `MessageReader` documentation --- src/link/mod.rs | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/src/link/mod.rs b/src/link/mod.rs index e54c9cd..eeba2fe 100644 --- a/src/link/mod.rs +++ b/src/link/mod.rs @@ -45,13 +45,16 @@ pub struct MessageReader { received_bytes: u64, } -/// For converting byte stream expected to be generated by Acedia mod from the game server into -/// actual messages. Expected format is a sequence of either: -/// [HEAD_UE_RECEIVED: 1 byte] [amount of bytes received by game server since last update: 4 bytes] -/// [HEAD_UE_MESSAGE: 1 byte] [length of the message: 4 bytes] [utf8-encoded string: ??? bytes] -/// On any invalid input enters a failure state (can be checked by `is_broken()`) and +/// For converting byte stream that is expected from the ue-server into actual messages. +/// Expected format is a sequence of either: +/// 1. [HEAD_UE_RECEIVED: marker byte | 1 byte] +/// [AMOUNT: amount of bytes received by ue-server since last update | 4 bytes: u32 BE] +/// 2. [HEAD_UE_MESSAGE: marker byte | 1 byte] +/// [LENGTH: length of the JSON message in utf8 encoding | 4 bytes: u32 BE] +/// [PAYLOAD: utf8-encoded string | `LENGTH` bytes] +/// On any invalid input enters into a failure state (can be checked by `is_broken()`) and /// never recovers from it. -/// Use either `push_byte()` or `push()` to input byte stream from game server and `pop()` to +/// Use either `push_byte()` or `push()` to input byte stream from ue-server and `pop()` to /// retrieve resulting messages. impl MessageReader { pub fn new() -> MessageReader { From 0dedd1d1f1386e7b48757499fa8bfad89e9127f8 Mon Sep 17 00:00:00 2001 From: Anton Tarasenko Date: Thu, 22 Jul 2021 18:18:29 +0700 Subject: [PATCH 3/5] Change `MessageReader` to use `with_capacity` Some of the collections inside `MessageReader` were created with `new()` instead of `with_capacity()` call. This patch fixes that or comments why it was not done in some places. --- src/link/mod.rs | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/src/link/mod.rs b/src/link/mod.rs index eeba2fe..db57956 100644 --- a/src/link/mod.rs +++ b/src/link/mod.rs @@ -20,6 +20,9 @@ const HEAD_UE_MESSAGE: u8 = 42; // Maximum allowed size of JSON message sent from ue-server. const MAX_UE_MESSAGE_LENGTH: usize = 25 * 1024 * 1024; +// We do not expect to receive more that this much messages at once from ue-server +const EXPECTED_LIMIT_TO_UE_MESSAGES: usize = 100; + custom_error! { pub ReadingStreamError InvalidHead{input: u8} = "Invalid byte used as a HEAD: {input}", MessageTooLong{length: usize} = "Message to receive is too long: {length}", @@ -63,8 +66,9 @@ impl MessageReader { reading_state: ReadingState::Head, read_bytes: 0, current_message_length: 0, + // Will be recreated with `with_capacity` in `push_byte()` current_message: Vec::new(), - read_messages: VecDeque::new(), + read_messages: VecDeque::with_capacity(EXPECTED_LIMIT_TO_UE_MESSAGES), next_received_bytes: 0, received_bytes: 0, } From af3341c1e77b776eb9ec23aaf6fa252c010e3916 Mon Sep 17 00:00:00 2001 From: Anton Tarasenko Date: Thu, 22 Jul 2021 18:48:09 +0700 Subject: [PATCH 4/5] Refactor `MessageReader` for clarity Rename `received_bytes` into `ue_received_bytes` and get rid of `next_received_bytes` by reading relevant data into new byte buffer instead. --- src/link/mod.rs | 59 +++++++++++++++++++++++++++---------------------- 1 file changed, 33 insertions(+), 26 deletions(-) diff --git a/src/link/mod.rs b/src/link/mod.rs index db57956..b761827 100644 --- a/src/link/mod.rs +++ b/src/link/mod.rs @@ -41,11 +41,11 @@ pub struct MessageReader { is_broken: bool, reading_state: ReadingState, read_bytes: usize, + buffer: [u8; 4], current_message_length: usize, current_message: Vec, read_messages: VecDeque, - next_received_bytes: u32, - received_bytes: u64, + ue_received_bytes: u64, } /// For converting byte stream that is expected from the ue-server into actual messages. @@ -65,12 +65,12 @@ impl MessageReader { is_broken: false, reading_state: ReadingState::Head, read_bytes: 0, + buffer: [0; 4], current_message_length: 0, // Will be recreated with `with_capacity` in `push_byte()` current_message: Vec::new(), read_messages: VecDeque::with_capacity(EXPECTED_LIMIT_TO_UE_MESSAGES), - next_received_bytes: 0, - received_bytes: 0, + ue_received_bytes: 0, } } @@ -81,32 +81,28 @@ impl MessageReader { match &self.reading_state { ReadingState::Head => { if input == HEAD_UE_RECEIVED { - self.reading_state = ReadingState::ReceivedBytes; + self.change_state(ReadingState::ReceivedBytes); } else if input == HEAD_UE_MESSAGE { - self.reading_state = ReadingState::Length; + self.change_state(ReadingState::Length); } else { self.is_broken = true; return Err(ReadingStreamError::InvalidHead { input }); } } ReadingState::ReceivedBytes => { - self.next_received_bytes = self.next_received_bytes << 8; - self.next_received_bytes += input as u32; + self.buffer[self.read_bytes] = input; self.read_bytes += 1; if self.read_bytes >= UE_RECEIVED_FIELD_SIZE { - self.received_bytes += self.next_received_bytes as u64; - self.next_received_bytes = 0; - self.read_bytes = 0; - self.reading_state = ReadingState::Head; + self.ue_received_bytes += array_of_u8_to_u32(self.buffer) as u64; + self.change_state(ReadingState::Head); } } ReadingState::Length => { - self.current_message_length = self.current_message_length << 8; - self.current_message_length += input as usize; + self.buffer[self.read_bytes] = input; self.read_bytes += 1; if self.read_bytes >= UE_LENGTH_FIELD_SIZE { - self.read_bytes = 0; - self.reading_state = ReadingState::Payload; + self.current_message_length = array_of_u8_to_u32(self.buffer) as usize; + self.change_state(ReadingState::Payload); if self.current_message_length > MAX_UE_MESSAGE_LENGTH { self.is_broken = true; return Err(ReadingStreamError::MessageTooLong { @@ -118,7 +114,7 @@ impl MessageReader { } ReadingState::Payload => { self.current_message.push(input); - self.read_bytes += 1 as usize; + self.read_bytes += 1; if self.read_bytes >= self.current_message_length { match str::from_utf8(&self.current_message) { Ok(next_message) => self.read_messages.push_front(next_message.to_owned()), @@ -129,8 +125,7 @@ impl MessageReader { }; self.current_message.clear(); self.current_message_length = 0; - self.read_bytes = 0; - self.reading_state = ReadingState::Head; + self.change_state(ReadingState::Head); } } } @@ -148,13 +143,25 @@ impl MessageReader { self.read_messages.pop_back() } - pub fn received_bytes(&self) -> u64 { - self.received_bytes + pub fn ue_received_bytes(&self) -> u64 { + self.ue_received_bytes } pub fn is_broken(&self) -> bool { self.is_broken } + + fn change_state(&mut self, next_state: ReadingState) { + self.read_bytes = 0; + self.reading_state = next_state; + } +} + +fn array_of_u8_to_u32(bytes: [u8; 4]) -> u64 { + (u64::from(bytes[0]) << 24) + + (u64::from(bytes[1]) << 16) + + (u64::from(bytes[2]) << 8) + + (u64::from(bytes[3])) } #[test] @@ -199,18 +206,18 @@ fn received_push_byte() { reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(243).unwrap(); - assert_eq!(reader.received_bytes(), 243); + assert_eq!(reader.ue_received_bytes(), 243); reader.push_byte(HEAD_UE_RECEIVED).unwrap(); reader.push_byte(65).unwrap(); reader.push_byte(25).unwrap(); reader.push_byte(178).unwrap(); reader.push_byte(4).unwrap(); - assert_eq!(reader.received_bytes(), 1092203255); + assert_eq!(reader.ue_received_bytes(), 1092203255); reader.push_byte(HEAD_UE_RECEIVED).unwrap(); reader.push_byte(231).unwrap(); reader.push_byte(34).unwrap(); reader.push_byte(154).unwrap(); - assert_eq!(reader.received_bytes(), 1092203255); + assert_eq!(reader.ue_received_bytes(), 1092203255); } #[test] @@ -234,7 +241,7 @@ fn mixed_push_byte() { reader.push_byte(25).unwrap(); reader.push_byte(178).unwrap(); reader.push_byte(4).unwrap(); - assert_eq!(reader.received_bytes(), 1092203255); + assert_eq!(reader.ue_received_bytes(), 1092203255); assert_eq!(reader.pop().unwrap(), "Yo!"); assert_eq!(reader.pop(), None); } @@ -264,7 +271,7 @@ fn pushing_many_bytes_at_once() { 4, ]) .unwrap(); - assert_eq!(reader.received_bytes(), 1092203255); + assert_eq!(reader.ue_received_bytes(), 1092203255); assert_eq!(reader.pop().unwrap(), "Yo!"); assert_eq!(reader.pop(), None); } From bb1f73a755753221159eb247553b5b4d0c9c9042 Mon Sep 17 00:00:00 2001 From: Anton Tarasenko Date: Thu, 22 Jul 2021 18:53:22 +0700 Subject: [PATCH 5/5] Change character byte definitions into b'X' form --- src/link/mod.rs | 44 ++++++++++++++++++++++---------------------- 1 file changed, 22 insertions(+), 22 deletions(-) diff --git a/src/link/mod.rs b/src/link/mod.rs index b761827..87d8533 100644 --- a/src/link/mod.rs +++ b/src/link/mod.rs @@ -172,27 +172,27 @@ fn message_push_byte() { reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(13).unwrap(); - reader.push_byte(0x48).unwrap(); // H - reader.push_byte(0x65).unwrap(); // e - reader.push_byte(0x6c).unwrap(); // l - reader.push_byte(0x6c).unwrap(); // l - reader.push_byte(0x6f).unwrap(); // o - reader.push_byte(0x2c).unwrap(); // , - reader.push_byte(0x20).unwrap(); // - reader.push_byte(0x77).unwrap(); // w - reader.push_byte(0x6f).unwrap(); // o - reader.push_byte(0x72).unwrap(); // r - reader.push_byte(0x6c).unwrap(); // l - reader.push_byte(0x64).unwrap(); // d - reader.push_byte(0x21).unwrap(); // + reader.push_byte(b'H').unwrap(); + reader.push_byte(b'e').unwrap(); + reader.push_byte(b'l').unwrap(); + reader.push_byte(b'l').unwrap(); + reader.push_byte(b'o').unwrap(); + reader.push_byte(b',').unwrap(); + reader.push_byte(b' ').unwrap(); + reader.push_byte(b'w').unwrap(); + reader.push_byte(b'o').unwrap(); + reader.push_byte(b'r').unwrap(); + reader.push_byte(b'l').unwrap(); + reader.push_byte(b'd').unwrap(); + reader.push_byte(b'!').unwrap(); reader.push_byte(HEAD_UE_MESSAGE).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(3).unwrap(); - reader.push_byte(0x59).unwrap(); // Y - reader.push_byte(0x6f).unwrap(); // o - reader.push_byte(0x21).unwrap(); // + reader.push_byte(b'Y').unwrap(); + reader.push_byte(b'o').unwrap(); + reader.push_byte(b'!').unwrap(); assert_eq!(reader.pop().unwrap(), "Hello, world!"); assert_eq!(reader.pop().unwrap(), "Yo!"); assert_eq!(reader.pop(), None); @@ -233,9 +233,9 @@ fn mixed_push_byte() { reader.push_byte(0).unwrap(); reader.push_byte(0).unwrap(); reader.push_byte(3).unwrap(); - reader.push_byte(0x59).unwrap(); // Y - reader.push_byte(0x6f).unwrap(); // o - reader.push_byte(0x21).unwrap(); // + reader.push_byte(b'Y').unwrap(); + reader.push_byte(b'o').unwrap(); + reader.push_byte(b'!').unwrap(); reader.push_byte(HEAD_UE_RECEIVED).unwrap(); reader.push_byte(65).unwrap(); reader.push_byte(25).unwrap(); @@ -261,9 +261,9 @@ fn pushing_many_bytes_at_once() { 0, 0, 3, - 0x59, // Y - 0x6f, // o - 0x21, // + b'Y', + b'o', + b'!', HEAD_UE_RECEIVED, 65, 25,