From e9014df0bc5ea21c9940068d86eec98248ec9f07 Mon Sep 17 00:00:00 2001 From: cool-mist Date: Wed, 2 Sep 2026 22:30:16 +0530 Subject: [PATCH] WIP --- Cargo.lock | 25 +- crates/hd-bus/src/lib.rs | 48 --- crates/hd-client/Cargo.toml | 1 + crates/hd-lib/Cargo.toml | 7 + .../src/bus.rs => hd-lib/src/bus/mod.rs} | 2 + crates/hd-lib/src/bus/packet.rs | 277 ++++++++++++++++++ .../src/error.rs => hd-lib/src/error/mod.rs} | 0 crates/hd-lib/src/lib.rs | 3 + .../src/thread_pool/mod.rs} | 2 +- crates/{hd-bus => hd-server}/Cargo.toml | 3 +- .../{hd-bus => hd-server}/src/bin/server.rs | 3 +- .../connection.rs => hd-server/src/lib.rs} | 49 +++- 12 files changed, 355 insertions(+), 65 deletions(-) delete mode 100644 crates/hd-bus/src/lib.rs create mode 100644 crates/hd-lib/Cargo.toml rename crates/{hd-bus/src/bus.rs => hd-lib/src/bus/mod.rs} (99%) create mode 100644 crates/hd-lib/src/bus/packet.rs rename crates/{hd-bus/src/error.rs => hd-lib/src/error/mod.rs} (100%) create mode 100644 crates/hd-lib/src/lib.rs rename crates/{hd-bus/src/thread_pool.rs => hd-lib/src/thread_pool/mod.rs} (97%) rename crates/{hd-bus => hd-server}/Cargo.toml (76%) rename crates/{hd-bus => hd-server}/src/bin/server.rs (90%) rename crates/{hd-bus/src/connection.rs => hd-server/src/lib.rs} (82%) diff --git a/Cargo.lock b/Cargo.lock index e38b20f..f1bc9f8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -17,21 +17,30 @@ version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9330f8b2ff13f34540b44e946ef35111825727b38d33286ef986142615121801" -[[package]] -name = "hd-bus" -version = "0.1.0" -dependencies = [ - "tracing", - "tracing-subscriber", -] - [[package]] name = "hd-client" version = "0.1.0" dependencies = [ + "hd-lib", "tracing", ] +[[package]] +name = "hd-lib" +version = "0.1.0" +dependencies = [ + "tracing", +] + +[[package]] +name = "hd-server" +version = "0.1.0" +dependencies = [ + "hd-lib", + "tracing", + "tracing-subscriber", +] + [[package]] name = "lazy_static" version = "1.5.0" diff --git a/crates/hd-bus/src/lib.rs b/crates/hd-bus/src/lib.rs deleted file mode 100644 index 8be3f68..0000000 --- a/crates/hd-bus/src/lib.rs +++ /dev/null @@ -1,48 +0,0 @@ -pub mod bus; -mod connection; -pub mod error; -mod thread_pool; - -use tracing::{Level, event, info, span}; - -use crate::{ - bus::{EventBus, EventData}, - connection::handle_event_bus_client, - error::HError, - thread_pool::ThreadPool, -}; -use std::net::TcpListener; - -pub struct TcpEventBus { - bus: EventBus, - pool: ThreadPool, -} - -impl TcpEventBus { - pub fn new() -> Self { - let pool = ThreadPool::new(10); - Self { - bus: EventBus::new(), - pool, - } - } - - pub fn publish(&self, evt: EventData) -> Result<(), HError> { - self.bus.publish(evt) - } -} - -impl TcpEventBus { - pub fn start(&self, addr: &'static str) -> Result<(), HError> { - let listener = TcpListener::bind(addr).map_err(HError::TcpSocketBindError)?; - info!("Heimdall listening at {}", addr); - for stream in listener.incoming() { - match stream { - Ok(stream) => handle_event_bus_client(&self.pool, self.bus.clone(), stream)?, - Err(e) => tracing::error!("TCP connection failed {}", e), - } - } - - Ok(()) - } -} diff --git a/crates/hd-client/Cargo.toml b/crates/hd-client/Cargo.toml index d3862b7..e4b56ed 100644 --- a/crates/hd-client/Cargo.toml +++ b/crates/hd-client/Cargo.toml @@ -5,3 +5,4 @@ edition = "2024" [dependencies] tracing = { workspace = true } +hd-lib = { path = "../hd-lib" } diff --git a/crates/hd-lib/Cargo.toml b/crates/hd-lib/Cargo.toml new file mode 100644 index 0000000..43c286c --- /dev/null +++ b/crates/hd-lib/Cargo.toml @@ -0,0 +1,7 @@ +[package] +name = "hd-lib" +version = "0.1.0" +edition = "2024" + +[dependencies] +tracing = { workspace = true } diff --git a/crates/hd-bus/src/bus.rs b/crates/hd-lib/src/bus/mod.rs similarity index 99% rename from crates/hd-bus/src/bus.rs rename to crates/hd-lib/src/bus/mod.rs index 4c71863..c1b1301 100644 --- a/crates/hd-bus/src/bus.rs +++ b/crates/hd-lib/src/bus/mod.rs @@ -1,3 +1,5 @@ +mod packet; + use tracing::{info, info_span}; use crate::error::HError; diff --git a/crates/hd-lib/src/bus/packet.rs b/crates/hd-lib/src/bus/packet.rs new file mode 100644 index 0000000..4e0bded --- /dev/null +++ b/crates/hd-lib/src/bus/packet.rs @@ -0,0 +1,277 @@ +use std::{io::{Read, Write}, net::TcpStream}; + +use crate::bus::EventData; + +enum PacketType { + Connect, + Unsubscribe, + Publish, + SendEvent, + Settle, + Ack, + Disconnect, +} + +enum Packet { + Connect(ConnectPacket), + Unsubscribe, + Publish(PublishPacket), + SendEvent(SendEventPacket), + Settle(SettlePacket), + Ack(AckPacket), + Disconnect, +} + +struct ConnectPacket { + cursor: u64, +} + +struct PublishPacket { + event_type: u8, + data_utf8: Vec, +} + +struct SendEventPacket { + event_id: u64, + event_type: u8, + data_utf8: Vec, +} + +struct SettlePacket { + cursor: u64, +} + +struct AckPacket { + packet_type: PacketType, +} + +impl Packet { + fn get_packet_type(&self) -> PacketType { + match self { + Packet::Connect(_) => PacketType::Connect, + Packet::Unsubscribe => PacketType::Unsubscribe, + Packet::Publish(_) => PacketType::Publish, + Packet::SendEvent(_) => PacketType::SendEvent, + Packet::Settle(_) => PacketType::Settle, + Packet::Ack(_) => PacketType::Ack, + Packet::Disconnect => PacketType::Disconnect, + } + } + + fn create_connect_packet(cursor: u64) -> Packet { + Packet::Connect(ConnectPacket { cursor }) + } + + fn create_unsubscribe_packet() -> Packet { + Packet::Unsubscribe + } + + fn create_publish_packet(evt: EventData) -> Packet { + Packet::Publish(PublishPacket { + event_type: evt.event_type, + data_utf8: evt.data_utf8, + }) + } + + fn create_settle_packet(cursor: u64) -> Packet { + Packet::Settle(SettlePacket { cursor }) + } + + fn create_ack_packet(packet_type: PacketType) -> Packet { + Packet::Ack(AckPacket { packet_type }) + } + + fn create_disconnect_packet() -> Packet { + Packet::Disconnect + } +} + +struct PacketFrame { + header: HeaderFrame, + data: DataFrame, +} + +/// 3 Bytes +struct HeaderFrame { + /// 0 Connect + /// 1 Subscribe + /// 2 Unsubscribe + /// 3 Publish + /// 4 Settle + /// 5 Ack + /// 6 Disconnect + packet_type: u8, + /// Length of data frame, Little Endian, so 5 = [5 0], 256 = [255 1] + data_frame_length: [u8; 2], +} + +struct DataFrame { + data: Vec, +} + +enum PacketError { + WrongPacketType(u8), + PacketProtocolError, + TcpStreamError(std::io::Error), +} + +fn read_packet(stream: &mut TcpStream) -> Result { + let mut header_buf = vec![0; 3]; + stream + .read_exact(&mut header_buf) + .map_err(PacketError::TcpStreamError)?; + + let packet_type = get_packet_type_from_byte(header_buf[0])?; + + match packet_type { + PacketType::Connect => { + let cursor = read_u64(stream)?; + Ok(Packet::Connect(ConnectPacket { cursor })) + } + PacketType::Unsubscribe => Ok(Packet::Unsubscribe), + PacketType::Publish => todo!(), + PacketType::SendEvent => todo!(), + PacketType::Settle => { + let cursor = read_u64(stream)?; + Ok(Packet::Settle(SettlePacket { cursor })) + } + PacketType::Ack => { + let packet_type = get_packet_type_from_byte(read_byte(stream)?)?; + Ok(Packet::Ack(AckPacket { packet_type })) + } + PacketType::Disconnect => Ok(Packet::Disconnect), + } +} + +fn write_packet(stream: &mut TcpStream, packet: &Packet) -> Result<(), PacketError> { + let packet_type = packet.get_packet_type(); + let data_frame_length = calculate_data_frame_length(&packet); + + let mut header_frame = [0; 3]; + write_header_frame(&mut header_frame, packet_type, data_frame_length); + + if data_frame_length > 0 { + let mut data_frame = vec![0; data_frame_length]; + write_data_frame(&mut data_frame, packet); + } + + // RESUME FROM HERE + + Ok(()) +} + +fn write_header_frame(header_frame: &mut [u8; 3], packet_type: PacketType, data_frame_length: usize) { + header_frame[0] = get_byte_from_packet_type(&packet_type); + if data_frame_length > 0 { + write_length(data_frame_length, &mut header_frame[1..3]); + } +} + +fn write_data_frame(data_frame: &mut [u8], packet: &Packet) { + match packet { + Packet::Connect(connect_packet) => { + write_u64(connect_packet.cursor, data_frame); + } + Packet::Publish(publish_packet) => { + write_event_data(publish_packet.event_type, &publish_packet.data_utf8, data_frame); + } + Packet::SendEvent(send_event_packet) => { + write_u64(send_event_packet.event_id, data_frame); + write_event_data( + send_event_packet.event_type, + &send_event_packet.data_utf8, + &mut data_frame[1..], + ); + } + Packet::Settle(settle_packet) => { + write_u64(settle_packet.cursor, data_frame); + } + Packet::Ack(ack_packet) => { + data_frame[0] = get_byte_from_packet_type(&ack_packet.packet_type); + } + _ => {} + } +} + +fn write_event_data(event_type: u8, data_utf8: &[u8], buf: &mut [u8]) { + buf[0] = event_type; + write_length(data_utf8.len(), &mut buf[1..3]); + let available_length = buf.len() - 3; + buf[3..].copy_from_slice(&data_utf8[..available_length]); +} + +fn write_u64(v64: u64, data_frame: &mut [u8]) { + let bytes = v64.to_le_bytes(); + data_frame[0..8].copy_from_slice(&bytes); +} + +fn read_u64(stream: &mut TcpStream) -> Result { + let mut le_bytes = [0u8; 8]; + stream.read_exact(&mut le_bytes).map_err(PacketError::TcpStreamError)?; + Ok(u64::from_le_bytes(le_bytes)) +} + +fn read_byte(stream: &mut TcpStream) -> Result { + let mut byte = [0u8; 1]; + stream.read_exact(&mut byte).map_err(PacketError::TcpStreamError)?; + Ok(byte[0]) +} + +fn get_byte_from_packet_type(packet_type: &PacketType) -> u8 { + match packet_type { + PacketType::Connect => 0, + PacketType::Unsubscribe => 1, + PacketType::Publish => 2, + PacketType::SendEvent => 3, + PacketType::Settle => 4, + PacketType::Ack => 5, + PacketType::Disconnect => 6, + } +} + +fn get_packet_type_from_byte(byte: u8) -> Result { + let packet_type = match byte { + 0 => Some(PacketType::Connect), + 1 => Some(PacketType::Unsubscribe), + 2 => Some(PacketType::Publish), + 3 => Some(PacketType::SendEvent), + 4 => Some(PacketType::Settle), + 5 => Some(PacketType::Ack), + 6 => Some(PacketType::Disconnect), + _ => None, + }; + + if let Some(packet_type) = packet_type { + return Ok(packet_type); + } + + Err(PacketError::WrongPacketType(byte)) +} + +fn write_length(mut length: usize, buf: &mut [u8]) { + if length > 0xffff { + length = 0xffff; + } + + let bytes = length.to_le_bytes(); + buf[0..2].copy_from_slice(&bytes); +} + +fn calculate_data_frame_length(packet: &Packet) -> usize { + let mut ret = match packet { + Packet::Connect(_) => 8, + Packet::Unsubscribe => 0, + Packet::Publish(publish_packet) => 1 + 2 + publish_packet.data_utf8.len(), + Packet::SendEvent(send_event_packet) => 8 + 1 + 2 + send_event_packet.data_utf8.len(), + Packet::Settle(_) => 8, + Packet::Ack(_) => 1, + Packet::Disconnect => 0, + }; + + // Normalize to 2 bytes + if ret > 0xffff { + ret = 0xffff; + } + + ret +} diff --git a/crates/hd-bus/src/error.rs b/crates/hd-lib/src/error/mod.rs similarity index 100% rename from crates/hd-bus/src/error.rs rename to crates/hd-lib/src/error/mod.rs diff --git a/crates/hd-lib/src/lib.rs b/crates/hd-lib/src/lib.rs new file mode 100644 index 0000000..a53ebec --- /dev/null +++ b/crates/hd-lib/src/lib.rs @@ -0,0 +1,3 @@ +pub mod error; +pub mod thread_pool; +pub mod bus; diff --git a/crates/hd-bus/src/thread_pool.rs b/crates/hd-lib/src/thread_pool/mod.rs similarity index 97% rename from crates/hd-bus/src/thread_pool.rs rename to crates/hd-lib/src/thread_pool/mod.rs index 85df88c..eb260e3 100644 --- a/crates/hd-bus/src/thread_pool.rs +++ b/crates/hd-lib/src/thread_pool/mod.rs @@ -6,7 +6,7 @@ use std::{ thread::{self, JoinHandle}, }; -use tracing::{Level, event, info, span}; +use tracing::{Level, info, span}; use crate::error::HError; diff --git a/crates/hd-bus/Cargo.toml b/crates/hd-server/Cargo.toml similarity index 76% rename from crates/hd-bus/Cargo.toml rename to crates/hd-server/Cargo.toml index 5d2e459..4fe19da 100644 --- a/crates/hd-bus/Cargo.toml +++ b/crates/hd-server/Cargo.toml @@ -1,8 +1,9 @@ [package] -name = "hd-bus" +name = "hd-server" version = "0.1.0" edition = "2024" [dependencies] tracing = { workspace = true } tracing-subscriber = { version = "0.3", features = ["env-filter", "fmt"] } +hd-lib = { path = "../hd-lib" } diff --git a/crates/hd-bus/src/bin/server.rs b/crates/hd-server/src/bin/server.rs similarity index 90% rename from crates/hd-bus/src/bin/server.rs rename to crates/hd-server/src/bin/server.rs index 95a04d0..af41f64 100644 --- a/crates/hd-bus/src/bin/server.rs +++ b/crates/hd-server/src/bin/server.rs @@ -1,6 +1,7 @@ use std::{sync::Arc, thread, time::Duration}; -use hd_bus::{TcpEventBus, bus::EventData, error::HError}; +use hd_lib::{bus::EventData, error::HError}; +use hd_server::TcpEventBus; use tracing::info_span; fn main() -> Result<(), HError> { diff --git a/crates/hd-bus/src/connection.rs b/crates/hd-server/src/lib.rs similarity index 82% rename from crates/hd-bus/src/connection.rs rename to crates/hd-server/src/lib.rs index 8f89e7d..608c321 100644 --- a/crates/hd-bus/src/connection.rs +++ b/crates/hd-server/src/lib.rs @@ -5,13 +5,47 @@ use std::{ time::Duration, }; +use hd_lib::{ + bus::{Event, EventBus, EventData}, error::HError, thread_pool::ThreadPool, +}; + use tracing::{Span, info, info_span}; -use crate::{ - bus::{Event, EventBus}, - error::HError, - thread_pool::ThreadPool, -}; +use std::net::TcpListener; + +pub struct TcpEventBus { + bus: EventBus, + pool: ThreadPool, +} + +impl TcpEventBus { + pub fn new() -> Self { + let pool = ThreadPool::new(10); + Self { + bus: EventBus::new(), + pool, + } + } + + pub fn publish(&self, evt: EventData) -> Result<(), HError> { + self.bus.publish(evt) + } +} + +impl TcpEventBus { + pub fn start(&self, addr: &'static str) -> Result<(), HError> { + let listener = TcpListener::bind(addr).map_err(HError::TcpSocketBindError)?; + info!("Heimdall listening at {}", addr); + for stream in listener.incoming() { + match stream { + Ok(stream) => handle_event_bus_client(&self.pool, self.bus.clone(), stream)?, + Err(e) => tracing::error!("TCP connection failed {}", e), + } + } + + Ok(()) + } +} #[derive(Copy, Clone, PartialEq)] enum EventBusClientConnectionState { @@ -129,7 +163,10 @@ fn handle_state_send_event(client: &mut EventBusClientConnection) -> Result ID: {} TYPE: {} LEN: {} <=", evt.id, evt.data.event_type, data_len); + info!( + "SEND => ID: {} TYPE: {} LEN: {} <=", + evt.id, evt.data.event_type, data_len + ); let stream = &mut client.stream; stream.write(&mut buf).map_err(HError::TcpPeerError)?;