diff --git a/server/kafka/kafka-script-proxy/src/main.rs b/server/kafka/kafka-script-proxy/src/main.rs index ff5c865..c0ddf71 100644 --- a/server/kafka/kafka-script-proxy/src/main.rs +++ b/server/kafka/kafka-script-proxy/src/main.rs @@ -71,7 +71,7 @@ async fn handle_add_partition() -> Result, Infallible> { } } write_to_log(&format!("Successfully updated partition count to {}", next_partition)).await; - return build_success_res(&next_partition.to_string()); + return build_success_res(¤t_count.to_string()); // current_count because count - 1 = max partition } async fn get_highest_partition_count() -> Result { diff --git a/server/src/http.rs b/server/src/http.rs index 6cf5426..8aae036 100644 --- a/server/src/http.rs +++ b/server/src/http.rs @@ -6,6 +6,7 @@ use hyper::server::conn::AddrStream; use hyper::service::{make_service_fn, service_fn}; use kafka::producer::Producer; use serde_json::json; +use serde::{Deserialize}; use tokio::sync::Mutex; use pong::event::event::{Event, EventReader, EventWriter}; use crate::kafka::{KafkaEventReaderImpl, KafkaSessionEventWriterImpl}; @@ -22,7 +23,7 @@ pub struct HttpServer { impl HttpServer { pub fn new(addr: [u8; 4], port: u16, kafka_host: &str) -> HttpServer { let session_manager = Arc::new(Mutex::new(SessionManager::new(kafka_host))); - let event_writer = Arc::new(Mutex::new(EventWriter::new(Box::new(KafkaSessionEventWriterImpl::session_writer(kafka_host))))); + let event_writer = Arc::new(Mutex::new(EventWriter::new(Box::new(KafkaSessionEventWriterImpl::new(kafka_host))))); let event_reader = Arc::new(Mutex::new(EventReader::new(Box::new(KafkaEventReaderImpl::from(kafka_host))))); HttpServer {addr, port, session_manager, event_writer, event_reader} } @@ -57,7 +58,7 @@ async fn handle_request(session_manager: &Arc>, event_writ println!("req to {} with method {}", req.uri().path(), req.method()); match (req.method(), req.uri().path()) { (&Method::POST, "/create_session") => handle_session_create(session_manager, req).await, - (&Method::POST, "/write") => handle_event_write(event_writer, req).await, + (&Method::POST, "/write") => handle_event_write(session_manager, req).await, (&Method::POST, "/read") => handle_event_read(event_reader, req).await, _ => Ok(Response::new("unknown".into())) } @@ -77,16 +78,31 @@ async fn handle_session_create(session_manager: &Arc>, req return Ok(Response::new(Body::from(serialized.to_string()))) } -async fn handle_event_write(event_writer: &Arc>, req: Request) -> Result, Infallible> { - let mut locked = event_writer.lock().await; +async fn handle_event_write(session_manager: &Arc>, req: Request) -> Result, Infallible> { + let locked = session_manager.lock().await; let body = body::to_bytes(req.into_body()).await.unwrap(); let event_str = std::str::from_utf8(&*body).unwrap(); - let event: Event = serde_json::from_str(event_str).unwrap(); - println!("Writing event to kafka: {:?}", event); - if let Err(e) = locked.write(event) { - println!("Failed to write to kafka! {:?}", e); + let event = serde_json::from_str::(event_str).unwrap(); + let writer = locked.get_session_writer(&event.session_id); + if let Err(e) = writer { + let err = format!("Failed to write event: {}", e); + println!("{}", err); + let mut res = Response::new(Body::from(err)); + *res.status_mut() = StatusCode::NOT_FOUND; + return Ok(res); } - Ok(Response::new("response".into())) + let mut writer = writer.unwrap(); + println!("Writing session event to kafka: {:?}", event); + let write_res = writer.write_to_session(event.topic, event.msg); + if let Err(e) = write_res { + let err = format!("Failed to write event: {}", e); + println!("{}", err); + let mut res = Response::new(Body::from(err)); + *res.status_mut() = StatusCode::INTERNAL_SERVER_ERROR; + return Ok(res); + } + println!("Successfully wrote event to kafka."); + return Ok(Response::new(Body::from("Event write successful"))); } async fn handle_event_read(event_reader: &Arc>, req: Request) -> Result, Infallible> { @@ -105,10 +121,16 @@ async fn handle_event_read(event_reader: &Arc>, req: Request< } } - async fn shutdown_signal() { // Wait for the CTRL+C signal tokio::signal::ctrl_c() .await .expect("failed to install CTRL+C signal handler"); } + +#[derive(Debug, Deserialize)] +struct SessionEventDTO { + session_id: String, + topic: String, + msg: String +} diff --git a/server/src/kafka.rs b/server/src/kafka.rs index 7523783..a7dfa73 100644 --- a/server/src/kafka.rs +++ b/server/src/kafka.rs @@ -17,7 +17,7 @@ pub struct KafkaSessionEventWriterImpl { producer: Producer } impl KafkaSessionEventWriterImpl { - pub fn session_writer(host: &str) -> KafkaSessionEventWriterImpl { + pub fn new(host: &str) -> KafkaSessionEventWriterImpl { println!("Connecting session_writer producer to kafka host: {}", host); let mut producer = Producer::from_hosts(vec![host.to_owned()]) .with_ack_timeout(Duration::from_secs(1)) diff --git a/server/src/session.rs b/server/src/session.rs index 7bc6426..d8cd3cd 100644 --- a/server/src/session.rs +++ b/server/src/session.rs @@ -45,14 +45,28 @@ impl SessionManager { Ok(session) } - pub fn get_session_reader(&self, session: Session) -> SessionReader { + pub fn get_session_reader(&self, session_id: &str) -> Result { + let session = self.find_session(&session_id); + if let None = session { + return Err(format!("Unable to find session with hash {}", session_id)) + } + let session = session.unwrap(); let event_reader = EventReader::new(Box::new(KafkaSessionEventReaderImpl::new(&self.kafka_host, &session, &["move", "status", "input"]))); - SessionReader {reader: event_reader, session} + Ok(SessionReader {reader: event_reader, session}) } - pub fn get_session_writer(&self, session: Session) -> SessionWriter { + pub fn get_session_writer(&self, session_id: &str) -> Result { + let session = self.find_session(&session_id); + if let None = session { + return Err(format!("Unable to find session with hash {}", session_id)) + } + let session = session.unwrap(); let event_writer = EventWriter::new(Box::new(KafkaSessionEventWriterImpl::new(&self.kafka_host))); - SessionWriter {writer: event_writer, session} + Ok(SessionWriter {writer: event_writer, session}) + } + + fn find_session(&self, session_id: &str) -> Option { + self.sessions.iter().find(|s| session_id == s.hash).map(|s| s.clone()) } }