Finish serde
This commit is contained in:
@@ -1,4 +1,7 @@
|
|||||||
use std::{io::{Read, Write}, net::TcpStream};
|
use std::{
|
||||||
|
io::{Read, Write},
|
||||||
|
net::TcpStream,
|
||||||
|
};
|
||||||
|
|
||||||
use crate::bus::EventData;
|
use crate::bus::EventData;
|
||||||
|
|
||||||
@@ -129,8 +132,40 @@ fn read_packet(stream: &mut TcpStream) -> Result<Packet, PacketError> {
|
|||||||
Ok(Packet::Connect(ConnectPacket { cursor }))
|
Ok(Packet::Connect(ConnectPacket { cursor }))
|
||||||
}
|
}
|
||||||
PacketType::Unsubscribe => Ok(Packet::Unsubscribe),
|
PacketType::Unsubscribe => Ok(Packet::Unsubscribe),
|
||||||
PacketType::Publish => todo!(),
|
PacketType::Publish => {
|
||||||
PacketType::SendEvent => todo!(),
|
let mut header = vec![0u8; 3];
|
||||||
|
stream.read_exact(&mut header).map_err(PacketError::TcpStreamError)?;
|
||||||
|
|
||||||
|
let event_type = header[0];
|
||||||
|
|
||||||
|
let len = usize::from_le_bytes([header[1], header[2], 0, 0, 0, 0, 0, 0]);
|
||||||
|
let mut data_utf8 = vec![0u8; len];
|
||||||
|
stream.read_exact(&mut data_utf8).map_err(PacketError::TcpStreamError)?;
|
||||||
|
|
||||||
|
let publish_packet = PublishPacket { event_type, data_utf8 };
|
||||||
|
Ok(Packet::Publish(publish_packet))
|
||||||
|
}
|
||||||
|
PacketType::SendEvent => {
|
||||||
|
let mut header = vec![0u8; 11];
|
||||||
|
stream.read_exact(&mut header).map_err(PacketError::TcpStreamError)?;
|
||||||
|
|
||||||
|
let mut le_bytes = [0u8; 8];
|
||||||
|
le_bytes.copy_from_slice(&header[0..8]);
|
||||||
|
let event_id = u64::from_le_bytes(le_bytes);
|
||||||
|
|
||||||
|
let event_type = header[8];
|
||||||
|
|
||||||
|
let len = usize::from_le_bytes([header[9], header[10], 0, 0, 0, 0, 0, 0]);
|
||||||
|
let mut data_utf8 = vec![0u8; len];
|
||||||
|
stream.read_exact(&mut data_utf8).map_err(PacketError::TcpStreamError)?;
|
||||||
|
|
||||||
|
let send_event_packet = SendEventPacket {
|
||||||
|
event_id,
|
||||||
|
event_type,
|
||||||
|
data_utf8,
|
||||||
|
};
|
||||||
|
Ok(Packet::SendEvent(send_event_packet))
|
||||||
|
}
|
||||||
PacketType::Settle => {
|
PacketType::Settle => {
|
||||||
let cursor = read_u64(stream)?;
|
let cursor = read_u64(stream)?;
|
||||||
Ok(Packet::Settle(SettlePacket { cursor }))
|
Ok(Packet::Settle(SettlePacket { cursor }))
|
||||||
@@ -147,20 +182,20 @@ fn write_packet(stream: &mut TcpStream, packet: &Packet) -> Result<(), PacketErr
|
|||||||
let packet_type = packet.get_packet_type();
|
let packet_type = packet.get_packet_type();
|
||||||
let data_frame_length = calculate_data_frame_length(&packet);
|
let data_frame_length = calculate_data_frame_length(&packet);
|
||||||
|
|
||||||
let mut header_frame = [0; 3];
|
let mut buf = vec![0u8; 3 + data_frame_length];
|
||||||
write_header_frame(&mut header_frame, packet_type, data_frame_length);
|
write_header_frame(&mut buf, packet_type, data_frame_length);
|
||||||
|
|
||||||
if data_frame_length > 0 {
|
if data_frame_length > 0 {
|
||||||
let mut data_frame = vec![0; data_frame_length];
|
write_data_frame(&mut buf[3..], packet);
|
||||||
write_data_frame(&mut data_frame, packet);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// RESUME FROM HERE
|
stream.write_all(&buf).map_err(PacketError::TcpStreamError)?;
|
||||||
|
stream.flush().map_err(PacketError::TcpStreamError)?;
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
fn write_header_frame(header_frame: &mut [u8; 3], packet_type: PacketType, data_frame_length: usize) {
|
fn write_header_frame(header_frame: &mut [u8], packet_type: PacketType, data_frame_length: usize) {
|
||||||
header_frame[0] = get_byte_from_packet_type(&packet_type);
|
header_frame[0] = get_byte_from_packet_type(&packet_type);
|
||||||
if data_frame_length > 0 {
|
if data_frame_length > 0 {
|
||||||
write_length(data_frame_length, &mut header_frame[1..3]);
|
write_length(data_frame_length, &mut header_frame[1..3]);
|
||||||
|
|||||||
Reference in New Issue
Block a user