Initial protocol
This commit is contained in:
Generated
+1
@@ -23,6 +23,7 @@ version = "0.1.0"
|
|||||||
dependencies = [
|
dependencies = [
|
||||||
"hd-lib",
|
"hd-lib",
|
||||||
"tracing",
|
"tracing",
|
||||||
|
"tracing-subscriber",
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
|
|||||||
@@ -6,3 +6,4 @@ edition = "2024"
|
|||||||
[dependencies]
|
[dependencies]
|
||||||
tracing = { workspace = true }
|
tracing = { workspace = true }
|
||||||
hd-lib = { path = "../hd-lib" }
|
hd-lib = { path = "../hd-lib" }
|
||||||
|
tracing-subscriber = { version = "0.3", features = ["env-filter", "fmt"] }
|
||||||
|
|||||||
@@ -1,6 +1,10 @@
|
|||||||
use hd_client::{EventBusClient, EventBusClientError};
|
use hd_client::{EventBusClient, EventBusClientError};
|
||||||
|
|
||||||
fn main() -> Result<(), EventBusClientError> {
|
fn main() -> Result<(), EventBusClientError> {
|
||||||
|
tracing_subscriber::fmt()
|
||||||
|
.with_env_filter(tracing_subscriber::EnvFilter::from_default_env())
|
||||||
|
.init();
|
||||||
|
|
||||||
let addr = "0.0.0.0:21368";
|
let addr = "0.0.0.0:21368";
|
||||||
let mut client = EventBusClient::start(addr)?;
|
let mut client = EventBusClient::start(addr)?;
|
||||||
client.subscribe()
|
client.subscribe()
|
||||||
|
|||||||
+26
-30
@@ -1,11 +1,22 @@
|
|||||||
use std::{
|
use std::net::TcpStream;
|
||||||
io::{Read, Write},
|
|
||||||
net::TcpStream,
|
use hd_lib::{
|
||||||
|
bus::packet::{self, Packet, PacketType},
|
||||||
|
error::HError,
|
||||||
|
hd_tcp,
|
||||||
};
|
};
|
||||||
|
use tracing::{debug, info};
|
||||||
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub enum EventBusClientError {
|
pub enum EventBusClientError {
|
||||||
TcpStreamError(std::io::Error),
|
TcpStreamError(std::io::Error),
|
||||||
|
HError(HError),
|
||||||
|
}
|
||||||
|
|
||||||
|
impl From<HError> for EventBusClientError {
|
||||||
|
fn from(value: HError) -> Self {
|
||||||
|
Self::HError(value)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub struct EventBusClient {
|
pub struct EventBusClient {
|
||||||
@@ -14,43 +25,28 @@ pub struct EventBusClient {
|
|||||||
|
|
||||||
impl EventBusClient {
|
impl EventBusClient {
|
||||||
pub fn start(addr: &'static str) -> Result<Self, EventBusClientError> {
|
pub fn start(addr: &'static str) -> Result<Self, EventBusClientError> {
|
||||||
let stream = TcpStream::connect(addr).map_err(EventBusClientError::TcpStreamError)?;
|
let mut stream = TcpStream::connect(addr).map_err(EventBusClientError::TcpStreamError)?;
|
||||||
println!("Event bus client connected to {}", addr);
|
hd_tcp::write_packet(&mut stream, &Packet::create_connect_packet())?;
|
||||||
|
hd_tcp::read_ack(&mut stream, &[PacketType::Connect])?;
|
||||||
|
|
||||||
|
info!("Event bus client connected to {}", addr);
|
||||||
Ok(Self { tcp_stream: stream })
|
Ok(Self { tcp_stream: stream })
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn subscribe(&mut self) -> Result<(), EventBusClientError> {
|
pub fn subscribe(&mut self) -> Result<(), EventBusClientError> {
|
||||||
let bytes = [1];
|
debug!("Going to subscribe");
|
||||||
self.tcp_stream
|
hd_tcp::write_packet(&mut self.tcp_stream, &Packet::create_subscribe_packet(0))?;
|
||||||
.write_all(&bytes)
|
hd_tcp::read_ack(&mut self.tcp_stream, &[PacketType::Subscribe])?;
|
||||||
.map_err(EventBusClientError::TcpStreamError)?;
|
|
||||||
|
|
||||||
self.tcp_stream.flush().map_err(EventBusClientError::TcpStreamError)?;
|
|
||||||
|
|
||||||
let mut buf = vec![0; 1024];
|
|
||||||
loop {
|
loop {
|
||||||
let Ok(bytes_read) = self.tcp_stream.read(&mut buf) else {
|
let packet = hd_tcp::read_packet(&mut self.tcp_stream)?;
|
||||||
return Ok(());
|
if let Packet::Disconnect = packet {
|
||||||
};
|
info!("Disconnecting from the server");
|
||||||
|
|
||||||
if bytes_read == 0 {
|
|
||||||
eprintln!("Connection closed");
|
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
|
|
||||||
println!("MESSAGE BEGIN");
|
debug!("Received packet ... {}", packet::get_packet_type(&packet));
|
||||||
for b in &buf[0..bytes_read] {
|
hd_tcp::write_ack(&mut self.tcp_stream, PacketType::SendEvent)?;
|
||||||
print!("[{}] ", b);
|
|
||||||
}
|
|
||||||
println!("MESSAGE END");
|
|
||||||
|
|
||||||
// ACK
|
|
||||||
self.tcp_stream
|
|
||||||
.write_all(&bytes)
|
|
||||||
.map_err(EventBusClientError::TcpStreamError)?;
|
|
||||||
|
|
||||||
self.tcp_stream.flush().map_err(EventBusClientError::TcpStreamError)?;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
pub mod packet;
|
pub mod packet;
|
||||||
|
|
||||||
use tracing::{info, info_span};
|
use tracing::{debug_span, info, info_span};
|
||||||
|
|
||||||
use crate::error::HError;
|
use crate::error::HError;
|
||||||
use std::{
|
use std::{
|
||||||
@@ -95,6 +95,7 @@ impl EventBus {
|
|||||||
|
|
||||||
for s in &bus.subscription_handles {
|
for s in &bus.subscription_handles {
|
||||||
if let Err(e) = s.new_event_signal.send(1) {
|
if let Err(e) = s.new_event_signal.send(1) {
|
||||||
|
// TODO: Multiple failures should automatically unsubscribe
|
||||||
tracing::error!("Failed to notify {}", e);
|
tracing::error!("Failed to notify {}", e);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -195,7 +196,7 @@ impl Subscription {
|
|||||||
fn start_delivery(mut subscription: Subscription, bus: EventBus) -> JoinHandle<Result<(), HError>> {
|
fn start_delivery(mut subscription: Subscription, bus: EventBus) -> JoinHandle<Result<(), HError>> {
|
||||||
thread::spawn(move || {
|
thread::spawn(move || {
|
||||||
loop {
|
loop {
|
||||||
let _scope = info_span!("subscription_delivery", subscription.id).entered();
|
let _scope = debug_span!("subscription_delivery", subscription.id).entered();
|
||||||
loop {
|
loop {
|
||||||
if let Ok(evt) = bus.get_next_event(subscription.cursor) {
|
if let Ok(evt) = bus.get_next_event(subscription.cursor) {
|
||||||
info!("Delivering event {} ", evt.id);
|
info!("Delivering event {} ", evt.id);
|
||||||
|
|||||||
@@ -18,6 +18,23 @@ pub enum PacketType {
|
|||||||
Disconnect,
|
Disconnect,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl Display for PacketType {
|
||||||
|
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||||
|
let to_write = match self {
|
||||||
|
PacketType::Connect => "CONNECT",
|
||||||
|
PacketType::Subscribe => "SUBSCRIBE",
|
||||||
|
PacketType::Unsubscribe => "UNSUBSCRIBE",
|
||||||
|
PacketType::Publish => "PUBLISH",
|
||||||
|
PacketType::SendEvent => "SEND_EVENT",
|
||||||
|
PacketType::Settle => "SETTLE_EVENT",
|
||||||
|
PacketType::Ack => "ACK",
|
||||||
|
PacketType::Disconnect => "DISCONNECT",
|
||||||
|
};
|
||||||
|
|
||||||
|
write!(f, "{}", to_write)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
pub enum Packet {
|
pub enum Packet {
|
||||||
Connect,
|
Connect,
|
||||||
Subscribe(SubscribePacket),
|
Subscribe(SubscribePacket),
|
||||||
@@ -49,11 +66,15 @@ pub struct SettlePacket {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub struct AckPacket {
|
pub struct AckPacket {
|
||||||
packet_type: PacketType,
|
pub packet_type: PacketType,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Packet {
|
impl Packet {
|
||||||
pub fn create_connect_packet(cursor: u64) -> Packet {
|
pub fn create_connect_packet() -> Packet {
|
||||||
|
Packet::Connect
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn create_subscribe_packet(cursor: u64) -> Packet {
|
||||||
Packet::Subscribe(SubscribePacket { cursor })
|
Packet::Subscribe(SubscribePacket { cursor })
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -115,7 +136,7 @@ pub struct ReadPacketState {
|
|||||||
|
|
||||||
pub fn start_read_packet(buf: &mut [u8; packet_constants::HEADER_SIZE]) -> Result<ReadPacketState, PacketError> {
|
pub fn start_read_packet(buf: &mut [u8; packet_constants::HEADER_SIZE]) -> Result<ReadPacketState, PacketError> {
|
||||||
let packet_type = get_packet_type_from_byte(buf[0])?;
|
let packet_type = get_packet_type_from_byte(buf[0])?;
|
||||||
let required_buffer_size = read_u16(&buf[1..3]);
|
let required_buffer_size = read_u16(&buf[1..3])?;
|
||||||
|
|
||||||
Ok(ReadPacketState {
|
Ok(ReadPacketState {
|
||||||
required_buffer_size,
|
required_buffer_size,
|
||||||
@@ -133,31 +154,29 @@ pub fn read_packet(buf: &[u8], read_state: ReadPacketState) -> Result<Packet, Pa
|
|||||||
PacketType::Unsubscribe => Ok(Packet::Unsubscribe),
|
PacketType::Unsubscribe => Ok(Packet::Unsubscribe),
|
||||||
PacketType::Publish => {
|
PacketType::Publish => {
|
||||||
let event_type = buf[0];
|
let event_type = buf[0];
|
||||||
let len = read_u16(&buf[1..]);
|
let event_len = read_u16(&buf[1..3])?;
|
||||||
|
|
||||||
// TODO: Share underlying bytes without copying
|
// TODO: Share underlying bytes without copying
|
||||||
let publish_packet = PublishPacket {
|
let publish_packet = PublishPacket {
|
||||||
event_type,
|
event_type,
|
||||||
data_utf8: buf[3..len].to_vec(),
|
data_utf8: buf[3..3 + event_len].to_vec(),
|
||||||
};
|
};
|
||||||
Ok(Packet::Publish(publish_packet))
|
Ok(Packet::Publish(publish_packet))
|
||||||
}
|
}
|
||||||
PacketType::SendEvent => {
|
PacketType::SendEvent => {
|
||||||
let event_id = read_u64(&buf)?;
|
let event_id = read_u64(&buf[0..8])?;
|
||||||
let event_type = buf[8];
|
let event_type = buf[8];
|
||||||
let len = usize::from_le_bytes([buf[9], buf[10], 0, 0, 0, 0, 0, 0]);
|
let event_len = read_u16(&buf[9..11])?;
|
||||||
|
|
||||||
let send_event_packet = SendEventPacket {
|
let send_event_packet = SendEventPacket {
|
||||||
event_id,
|
event_id,
|
||||||
event_type,
|
event_type,
|
||||||
data_utf8: buf[11..len].to_vec(),
|
data_utf8: buf[11..11 + event_len].to_vec(),
|
||||||
};
|
};
|
||||||
Ok(Packet::SendEvent(send_event_packet))
|
Ok(Packet::SendEvent(send_event_packet))
|
||||||
}
|
}
|
||||||
PacketType::Settle => {
|
PacketType::Settle => {
|
||||||
let mut le_bytes = [0u8; 8];
|
let cursor = read_u64(&buf[0..8])?;
|
||||||
le_bytes.copy_from_slice(&buf[0..8]);
|
|
||||||
let cursor = u64::from_le_bytes(le_bytes);
|
|
||||||
Ok(Packet::Settle(SettlePacket { cursor }))
|
Ok(Packet::Settle(SettlePacket { cursor }))
|
||||||
}
|
}
|
||||||
PacketType::Ack => {
|
PacketType::Ack => {
|
||||||
@@ -214,7 +233,7 @@ fn write_data_frame(data_frame: &mut [u8], packet: &Packet) {
|
|||||||
write_event_data(
|
write_event_data(
|
||||||
send_event_packet.event_type,
|
send_event_packet.event_type,
|
||||||
&send_event_packet.data_utf8,
|
&send_event_packet.data_utf8,
|
||||||
&mut data_frame[1..],
|
&mut data_frame[8..],
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
Packet::Settle(settle_packet) => {
|
Packet::Settle(settle_packet) => {
|
||||||
@@ -231,12 +250,12 @@ fn write_event_data(event_type: u8, data_utf8: &[u8], buf: &mut [u8]) {
|
|||||||
buf[0] = event_type;
|
buf[0] = event_type;
|
||||||
write_u16(data_utf8.len(), &mut buf[1..3]);
|
write_u16(data_utf8.len(), &mut buf[1..3]);
|
||||||
let available_length = buf.len() - 3;
|
let available_length = buf.len() - 3;
|
||||||
buf[3..].copy_from_slice(&data_utf8[..available_length]);
|
buf[3..].copy_from_slice(&data_utf8[0..available_length]);
|
||||||
}
|
}
|
||||||
|
|
||||||
fn write_u64(v64: u64, data_frame: &mut [u8]) {
|
fn write_u64(v64: u64, data_frame: &mut [u8]) {
|
||||||
let bytes = v64.to_le_bytes();
|
let bytes = v64.to_le_bytes();
|
||||||
data_frame[0..8].copy_from_slice(&bytes);
|
data_frame[0..8].copy_from_slice(&bytes[0..8]);
|
||||||
}
|
}
|
||||||
|
|
||||||
fn write_u16(mut length: usize, buf: &mut [u8]) {
|
fn write_u16(mut length: usize, buf: &mut [u8]) {
|
||||||
@@ -245,17 +264,25 @@ fn write_u16(mut length: usize, buf: &mut [u8]) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let bytes = length.to_le_bytes();
|
let bytes = length.to_le_bytes();
|
||||||
buf[0..2].copy_from_slice(&bytes);
|
buf[0..2].copy_from_slice(&bytes[0..2]);
|
||||||
}
|
}
|
||||||
|
|
||||||
fn read_u64(buf: &[u8]) -> Result<u64, PacketError> {
|
fn read_u64(buf: &[u8]) -> Result<u64, PacketError> {
|
||||||
let mut le_bytes = [0u8; 8];
|
let mut le_bytes = [0u8; 8];
|
||||||
|
if le_bytes.len() != buf.len() {
|
||||||
|
return Err(PacketError::PacketProtocolError);
|
||||||
|
}
|
||||||
|
|
||||||
le_bytes.copy_from_slice(buf);
|
le_bytes.copy_from_slice(buf);
|
||||||
Ok(u64::from_le_bytes(le_bytes))
|
Ok(u64::from_le_bytes(le_bytes))
|
||||||
}
|
}
|
||||||
|
|
||||||
fn read_u16(buf: &[u8]) -> usize {
|
fn read_u16(buf: &[u8]) -> Result<usize, PacketError> {
|
||||||
usize::from_le_bytes([buf[0], buf[1], 0, 0, 0, 0, 0, 0])
|
if buf.len() != 2 {
|
||||||
|
return Err(PacketError::PacketProtocolError);
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(usize::from_le_bytes([buf[0], buf[1], 0, 0, 0, 0, 0, 0]))
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn get_packet_type(packet: &Packet) -> PacketType {
|
pub fn get_packet_type(packet: &Packet) -> PacketType {
|
||||||
@@ -286,13 +313,14 @@ fn get_byte_from_packet_type(packet_type: &PacketType) -> u8 {
|
|||||||
|
|
||||||
fn get_packet_type_from_byte(byte: u8) -> Result<PacketType, PacketError> {
|
fn get_packet_type_from_byte(byte: u8) -> Result<PacketType, PacketError> {
|
||||||
let packet_type = match byte {
|
let packet_type = match byte {
|
||||||
0 => Some(PacketType::Subscribe),
|
0 => Some(PacketType::Connect),
|
||||||
1 => Some(PacketType::Unsubscribe),
|
1 => Some(PacketType::Subscribe),
|
||||||
2 => Some(PacketType::Publish),
|
2 => Some(PacketType::Unsubscribe),
|
||||||
3 => Some(PacketType::SendEvent),
|
3 => Some(PacketType::Publish),
|
||||||
4 => Some(PacketType::Settle),
|
4 => Some(PacketType::SendEvent),
|
||||||
5 => Some(PacketType::Ack),
|
5 => Some(PacketType::Settle),
|
||||||
6 => Some(PacketType::Disconnect),
|
6 => Some(PacketType::Ack),
|
||||||
|
7 => Some(PacketType::Disconnect),
|
||||||
_ => None,
|
_ => None,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|||||||
@@ -1,3 +1,62 @@
|
|||||||
|
pub mod bus;
|
||||||
pub mod error;
|
pub mod error;
|
||||||
pub mod thread_pool;
|
pub mod thread_pool;
|
||||||
pub mod bus;
|
|
||||||
|
pub mod hd_tcp {
|
||||||
|
use std::{
|
||||||
|
io::{Read, Write},
|
||||||
|
net::TcpStream,
|
||||||
|
};
|
||||||
|
|
||||||
|
use crate::{
|
||||||
|
bus::packet::{self, Packet, PacketType},
|
||||||
|
error::HError,
|
||||||
|
};
|
||||||
|
|
||||||
|
pub fn read_packet(stream: &mut TcpStream) -> Result<Packet, HError> {
|
||||||
|
let mut header_buf = [0u8; 3];
|
||||||
|
stream.read_exact(&mut header_buf)?;
|
||||||
|
let read_state = packet::start_read_packet(&mut header_buf)?;
|
||||||
|
|
||||||
|
let mut buf = vec![0u8; read_state.required_buffer_size];
|
||||||
|
if read_state.required_buffer_size > 0 {
|
||||||
|
stream.read_exact(&mut buf)?;
|
||||||
|
}
|
||||||
|
let packet = packet::read_packet(&buf[0..], read_state)?;
|
||||||
|
|
||||||
|
if let Packet::Disconnect = packet {
|
||||||
|
return Err(HError::PeerDisconnect);
|
||||||
|
};
|
||||||
|
|
||||||
|
Ok(packet)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn read_ack(stream: &mut TcpStream, packet_types: &[PacketType]) -> Result<Packet, HError> {
|
||||||
|
let packet = read_packet(stream)?;
|
||||||
|
if let Packet::Ack(ack_packet) = &packet {
|
||||||
|
for packet_type in packet_types.iter() {
|
||||||
|
if *packet_type == ack_packet.packet_type {
|
||||||
|
return Ok(packet);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Err(HError::ProtocolError)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn write_ack(stream: &mut TcpStream, packet_type: PacketType) -> Result<(), HError> {
|
||||||
|
let packet = Packet::create_ack_packet(packet_type);
|
||||||
|
write_packet(stream, &packet)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn write_packet(stream: &mut TcpStream, packet: &Packet) -> Result<(), HError> {
|
||||||
|
let write_state = packet::start_write_packet(packet);
|
||||||
|
let mut buf = vec![0u8; write_state.required_buffer_size];
|
||||||
|
packet::write_packet(&mut buf, packet, write_state)?;
|
||||||
|
|
||||||
|
stream.write_all(&mut buf)?;
|
||||||
|
stream.flush()?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -10,14 +10,17 @@ fn main() -> Result<(), HError> {
|
|||||||
.init();
|
.init();
|
||||||
|
|
||||||
let addr = "0.0.0.0:21368";
|
let addr = "0.0.0.0:21368";
|
||||||
let _scope = info_span!("example_server").entered();
|
|
||||||
let server = Arc::new(TcpEventBus::new());
|
let server = Arc::new(TcpEventBus::new());
|
||||||
server.publish(EventData::new(1, "Hello"))?;
|
server.publish(EventData::new(1, "Hello"))?;
|
||||||
server.publish(EventData::new(1, "World"))?;
|
server.publish(EventData::new(1, "World"))?;
|
||||||
|
|
||||||
let server_clone = server.clone();
|
let server_clone = server.clone();
|
||||||
thread::spawn(move || server.start(addr));
|
thread::spawn(move || {
|
||||||
|
let _scope = info_span!("tcp_server").entered();
|
||||||
|
server.start(addr)
|
||||||
|
});
|
||||||
|
|
||||||
|
let _scope = info_span!("health_checks").entered();
|
||||||
for i in 0..255 {
|
for i in 0..255 {
|
||||||
thread::sleep(Duration::from_secs(5));
|
thread::sleep(Duration::from_secs(5));
|
||||||
server_clone.publish(EventData::new(i, "asd"))?;
|
server_clone.publish(EventData::new(i, "asd"))?;
|
||||||
|
|||||||
+81
-83
@@ -1,5 +1,4 @@
|
|||||||
use std::{
|
use std::{
|
||||||
io::{Read, Write},
|
|
||||||
net::TcpStream,
|
net::TcpStream,
|
||||||
sync::{Arc, mpsc::Receiver},
|
sync::{Arc, mpsc::Receiver},
|
||||||
time::Duration,
|
time::Duration,
|
||||||
@@ -8,13 +7,14 @@ use std::{
|
|||||||
use hd_lib::{
|
use hd_lib::{
|
||||||
bus::{
|
bus::{
|
||||||
Event, EventBus, EventData,
|
Event, EventBus, EventData,
|
||||||
packet::{self, Packet, PacketType},
|
packet::{self, Packet, PacketType, PublishPacket, get_packet_type},
|
||||||
},
|
},
|
||||||
error::HError,
|
error::HError,
|
||||||
|
hd_tcp,
|
||||||
thread_pool::ThreadPool,
|
thread_pool::ThreadPool,
|
||||||
};
|
};
|
||||||
|
|
||||||
use tracing::{Span, info, info_span};
|
use tracing::{debug, error, info, info_span};
|
||||||
|
|
||||||
use std::net::TcpListener;
|
use std::net::TcpListener;
|
||||||
|
|
||||||
@@ -40,11 +40,11 @@ impl TcpEventBus {
|
|||||||
impl TcpEventBus {
|
impl TcpEventBus {
|
||||||
pub fn start(&self, addr: &'static str) -> Result<(), HError> {
|
pub fn start(&self, addr: &'static str) -> Result<(), HError> {
|
||||||
let listener = TcpListener::bind(addr)?;
|
let listener = TcpListener::bind(addr)?;
|
||||||
info!("Heimdall listening at {}", addr);
|
info!("Heimdall running at {}", addr);
|
||||||
for stream in listener.incoming() {
|
for stream in listener.incoming() {
|
||||||
match stream {
|
match stream {
|
||||||
Ok(stream) => handle_event_bus_client(&self.pool, self.bus.clone(), stream)?,
|
Ok(stream) => handle_event_bus_client(&self.pool, self.bus.clone(), stream)?,
|
||||||
Err(e) => tracing::error!("TCP connection failed {}", e),
|
Err(e) => error!("TCP connection failed {}", e),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -90,57 +90,11 @@ impl EventBusClientConnection {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn read_ack(stream: &mut TcpStream, packet_type: PacketType) -> Result<(), HError> {
|
|
||||||
let packet = read_packet(stream)?;
|
|
||||||
if packet_type == packet::get_packet_type(&packet) {
|
|
||||||
return Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
Err(HError::ProtocolError)
|
|
||||||
}
|
|
||||||
|
|
||||||
fn read_packet(stream: &mut TcpStream) -> Result<Packet, HError> {
|
|
||||||
let mut header_buf = [0u8; 3];
|
|
||||||
stream.read_exact(&mut header_buf)?;
|
|
||||||
let read_state = packet::start_read_packet(&mut header_buf)?;
|
|
||||||
|
|
||||||
let mut buf = vec![0u8; read_state.required_buffer_size];
|
|
||||||
stream.read_exact(&mut buf)?;
|
|
||||||
let packet = packet::read_packet(&buf[0..], read_state)?;
|
|
||||||
|
|
||||||
if let Packet::Disconnect = packet {
|
|
||||||
return Err(HError::PeerDisconnect);
|
|
||||||
};
|
|
||||||
|
|
||||||
Ok(packet)
|
|
||||||
}
|
|
||||||
|
|
||||||
fn write_ack(stream: &mut TcpStream, packet_type: PacketType) -> Result<(), HError> {
|
|
||||||
let packet = Packet::create_ack_packet(packet_type);
|
|
||||||
write_packet(stream, &packet)
|
|
||||||
}
|
|
||||||
|
|
||||||
fn write_packet(stream: &mut TcpStream, packet: &Packet) -> Result<(), HError> {
|
|
||||||
let write_state = packet::start_write_packet(packet);
|
|
||||||
let mut buf = vec![0u8; write_state.required_buffer_size];
|
|
||||||
packet::write_packet(&mut buf, packet, write_state)?;
|
|
||||||
|
|
||||||
stream.write_all(&mut buf)?;
|
|
||||||
stream.flush()?;
|
|
||||||
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
pub fn handle_event_bus_client(thread_pool: &ThreadPool, bus: EventBus, stream: TcpStream) -> Result<(), HError> {
|
pub fn handle_event_bus_client(thread_pool: &ThreadPool, bus: EventBus, stream: TcpStream) -> Result<(), HError> {
|
||||||
thread_pool.execute(move || {
|
thread_pool.execute(move || {
|
||||||
let mut client = EventBusClientConnection::new(bus, stream)?;
|
let mut client = EventBusClientConnection::new(bus, stream)?;
|
||||||
|
|
||||||
let _scope = info_span!(
|
let _scope = info_span!("handle_client", addr = client.addr,).entered();
|
||||||
"handle_event_bus_client",
|
|
||||||
addr = client.addr,
|
|
||||||
state = %client.state
|
|
||||||
)
|
|
||||||
.entered();
|
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
if let EventBusClientConnectionState::Disconnect = client.state {
|
if let EventBusClientConnectionState::Disconnect = client.state {
|
||||||
@@ -154,8 +108,6 @@ pub fn handle_event_bus_client(thread_pool: &ThreadPool, bus: EventBus, stream:
|
|||||||
if prev_state != client.state {
|
if prev_state != client.state {
|
||||||
info!("State change {} -> {}", prev_state, client.state);
|
info!("State change {} -> {}", prev_state, client.state);
|
||||||
}
|
}
|
||||||
|
|
||||||
Span::current().record("state", format!("{}", client.state));
|
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -163,7 +115,7 @@ pub fn handle_event_bus_client(thread_pool: &ThreadPool, bus: EventBus, stream:
|
|||||||
fn handle_event_bus_client_state(client: &mut EventBusClientConnection) -> Result<(), HError> {
|
fn handle_event_bus_client_state(client: &mut EventBusClientConnection) -> Result<(), HError> {
|
||||||
let next_state = match client.state {
|
let next_state = match client.state {
|
||||||
EventBusClientConnectionState::Connect => handle_state_connect(client),
|
EventBusClientConnectionState::Connect => handle_state_connect(client),
|
||||||
EventBusClientConnectionState::Command => handle_state_connect(client),
|
EventBusClientConnectionState::Command => handle_state_command(client),
|
||||||
EventBusClientConnectionState::SendEvent => handle_state_send_event(client),
|
EventBusClientConnectionState::SendEvent => handle_state_send_event(client),
|
||||||
EventBusClientConnectionState::WaitEvent => handle_state_wait_event(client),
|
EventBusClientConnectionState::WaitEvent => handle_state_wait_event(client),
|
||||||
EventBusClientConnectionState::Disconnect => Ok(EventBusClientConnectionState::Disconnect),
|
EventBusClientConnectionState::Disconnect => Ok(EventBusClientConnectionState::Disconnect),
|
||||||
@@ -174,66 +126,112 @@ fn handle_event_bus_client_state(client: &mut EventBusClientConnection) -> Resul
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn handle_state_connect(client: &mut EventBusClientConnection) -> Result<EventBusClientConnectionState, HError> {
|
fn handle_state_connect(client: &mut EventBusClientConnection) -> Result<EventBusClientConnectionState, HError> {
|
||||||
let Packet::Connect = read_packet(&mut client.stream)? else {
|
match hd_tcp::read_packet(&mut client.stream)? {
|
||||||
return Ok(EventBusClientConnectionState::Disconnect);
|
Packet::Connect => {
|
||||||
};
|
info!("Received connect, acking");
|
||||||
|
|
||||||
|
hd_tcp::write_ack(&mut client.stream, PacketType::Connect)?;
|
||||||
|
info!("Connect ack complete");
|
||||||
|
|
||||||
write_ack(&mut client.stream, PacketType::Connect)?;
|
|
||||||
Ok(EventBusClientConnectionState::Command)
|
Ok(EventBusClientConnectionState::Command)
|
||||||
}
|
}
|
||||||
|
p => {
|
||||||
|
error!(
|
||||||
|
"Received wrong packet of type {}, expected {}. Will disconnect",
|
||||||
|
get_packet_type(&p),
|
||||||
|
PacketType::Connect
|
||||||
|
);
|
||||||
|
Ok(EventBusClientConnectionState::Disconnect)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fn handle_state_command(client: &mut EventBusClientConnection) -> Result<EventBusClientConnectionState, HError> {
|
fn handle_state_command(client: &mut EventBusClientConnection) -> Result<EventBusClientConnectionState, HError> {
|
||||||
match read_packet(&mut client.stream)? {
|
match hd_tcp::read_packet(&mut client.stream)? {
|
||||||
// TODO: Read cursor
|
// TODO: Read cursor
|
||||||
Packet::Subscribe(_) => {
|
Packet::Subscribe(_) => {
|
||||||
write_ack(&mut client.stream, PacketType::Subscribe)?;
|
info!("Received subscribe");
|
||||||
|
|
||||||
let subscription = client.bus.subscribe()?;
|
let subscription = client.bus.subscribe()?;
|
||||||
|
|
||||||
|
info!("Subscribed to the event bus");
|
||||||
client.subscription = Some(subscription);
|
client.subscription = Some(subscription);
|
||||||
|
|
||||||
|
// TODO: Return error if any
|
||||||
|
hd_tcp::write_ack(&mut client.stream, PacketType::Subscribe)?;
|
||||||
|
info!("Subscribe ack complete");
|
||||||
Ok(EventBusClientConnectionState::SendEvent)
|
Ok(EventBusClientConnectionState::SendEvent)
|
||||||
}
|
}
|
||||||
Packet::Publish(publish_packet) => {
|
Packet::Publish(publish_packet) => {
|
||||||
write_ack(&mut client.stream, PacketType::Publish)?;
|
info!("Received publish packet");
|
||||||
|
publish_packet_to_bus(client, publish_packet)?;
|
||||||
|
|
||||||
// TODO: Possible without copying data_utf8?
|
// TODO: Return error if any
|
||||||
let data = String::from_utf8(publish_packet.data_utf8)?;
|
hd_tcp::write_ack(&mut client.stream, PacketType::Publish)?;
|
||||||
client.bus.publish(EventData::new(publish_packet.event_type, &data))?;
|
debug!("Publish ack complete");
|
||||||
Ok(EventBusClientConnectionState::WaitEvent)
|
Ok(EventBusClientConnectionState::WaitEvent)
|
||||||
}
|
}
|
||||||
_ => Ok(EventBusClientConnectionState::Disconnect),
|
p => {
|
||||||
|
error!(
|
||||||
|
"Received wrong packet of type {}, expected {}. Will disconnect",
|
||||||
|
get_packet_type(&p),
|
||||||
|
PacketType::Connect
|
||||||
|
);
|
||||||
|
Ok(EventBusClientConnectionState::Disconnect)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn handle_state_send_event(client: &mut EventBusClientConnection) -> Result<EventBusClientConnectionState, HError> {
|
fn handle_state_send_event(client: &mut EventBusClientConnection) -> Result<EventBusClientConnectionState, HError> {
|
||||||
if let Some(subscription) = &mut client.subscription {
|
if let Some(subscription) = &mut client.subscription {
|
||||||
loop {
|
// TODO: When client sends unsubscribe, server shouldn't take forever
|
||||||
|
// to ack that.
|
||||||
let Ok(evt) = subscription.recv() else {
|
let Ok(evt) = subscription.recv() else {
|
||||||
|
error!("Failed to receive a message for this subscription, disconnecting");
|
||||||
return Ok(EventBusClientConnectionState::Disconnect);
|
return Ok(EventBusClientConnectionState::Disconnect);
|
||||||
};
|
};
|
||||||
|
|
||||||
info!(
|
|
||||||
"SEND => ID: {} TYPE: {} LEN: {} <=",
|
|
||||||
evt.id,
|
|
||||||
evt.data.event_type,
|
|
||||||
evt.data.data_utf8.len()
|
|
||||||
);
|
|
||||||
|
|
||||||
let packet = Packet::create_send_event_packet(evt);
|
let packet = Packet::create_send_event_packet(evt);
|
||||||
write_packet(&mut client.stream, &packet)?;
|
hd_tcp::write_packet(&mut client.stream, &packet)?;
|
||||||
read_ack(&mut client.stream, PacketType::SendEvent)?;
|
debug!("Published event to client, waiting for ack");
|
||||||
}
|
|
||||||
|
let ack_packet = hd_tcp::read_ack(&mut client.stream, &[PacketType::SendEvent, PacketType::Unsubscribe])?;
|
||||||
|
if packet::get_packet_type(&ack_packet) == PacketType::Unsubscribe {
|
||||||
|
info!("Subscriber unsubscribed");
|
||||||
|
return Ok(EventBusClientConnectionState::Command);
|
||||||
}
|
}
|
||||||
|
|
||||||
Ok(EventBusClientConnectionState::Disconnect)
|
debug!("Publish ack received, message is settled");
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(EventBusClientConnectionState::SendEvent)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn handle_state_wait_event(client: &mut EventBusClientConnection) -> Result<EventBusClientConnectionState, HError> {
|
fn handle_state_wait_event(client: &mut EventBusClientConnection) -> Result<EventBusClientConnectionState, HError> {
|
||||||
todo!()
|
match hd_tcp::read_packet(&mut client.stream)? {
|
||||||
|
Packet::Publish(publish_packet) => {
|
||||||
|
debug!("Received event from client, publishing to bus");
|
||||||
|
publish_packet_to_bus(client, publish_packet)?;
|
||||||
|
|
||||||
|
hd_tcp::write_ack(&mut client.stream, PacketType::Publish)?;
|
||||||
|
debug!("Ack receive event complete");
|
||||||
|
Ok(EventBusClientConnectionState::WaitEvent)
|
||||||
|
}
|
||||||
|
p => {
|
||||||
|
error!(
|
||||||
|
"Received wrong packet of type {}, expected {}. Will disconnect",
|
||||||
|
get_packet_type(&p),
|
||||||
|
PacketType::Connect
|
||||||
|
);
|
||||||
|
Ok(EventBusClientConnectionState::Disconnect)
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn write_le_bytes(v64: u64, buf: &mut [u8]) {
|
fn publish_packet_to_bus(client: &mut EventBusClientConnection, packet: PublishPacket) -> Result<(), HError> {
|
||||||
let le_bytes = v64.to_le_bytes();
|
let data = String::from_utf8(packet.data_utf8)?;
|
||||||
buf[0..8].copy_from_slice(&le_bytes);
|
let evt = EventData::new(packet.event_type, &data);
|
||||||
|
client.bus.publish(evt)?;
|
||||||
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
impl std::fmt::Display for EventBusClientConnectionState {
|
impl std::fmt::Display for EventBusClientConnectionState {
|
||||||
|
|||||||
Reference in New Issue
Block a user