working session event write

This commit is contained in:
Thilo Behnke
2022-05-28 22:06:44 +02:00
parent 3f5ddf917d
commit 811e833f15
4 changed files with 52 additions and 16 deletions
+1 -1
View File
@@ -71,7 +71,7 @@ async fn handle_add_partition() -> Result<Response<Body>, 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(&current_count.to_string()); // current_count because count - 1 = max partition
}
async fn get_highest_partition_count() -> Result<u32, PartitionCountQueryError> {
+32 -10
View File
@@ -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<Mutex<SessionManager>>, 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<Mutex<SessionManager>>, req
return Ok(Response::new(Body::from(serialized.to_string())))
}
async fn handle_event_write(event_writer: &Arc<Mutex<EventWriter>>, req: Request<Body>) -> Result<Response<Body>, Infallible> {
let mut locked = event_writer.lock().await;
async fn handle_event_write(session_manager: &Arc<Mutex<SessionManager>>, req: Request<Body>) -> Result<Response<Body>, 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::<SessionEventDTO>(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<Mutex<EventReader>>, req: Request<Body>) -> Result<Response<Body>, Infallible> {
@@ -105,10 +121,16 @@ async fn handle_event_read(event_reader: &Arc<Mutex<EventReader>>, 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
}
+1 -1
View File
@@ -17,7 +17,7 @@ pub struct KafkaSessionEventWriterImpl {
producer: Producer<SessionPartitioner>
}
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))
+18 -4
View File
@@ -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<SessionReader, String> {
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<SessionWriter, String> {
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<Session> {
self.sessions.iter().find(|s| session_id == s.hash).map(|s| s.clone())
}
}