#![forbid(unsafe_code)]
use permit::Permit;
use servlin::reexport::{safina_executor, safina_timer};
use servlin::{
print_log_response, socket_addr_127_0_0_1, Event, EventSender, HttpServerBuilder, Request,
Response,
};
use std::sync::{Arc, Mutex};
use std::time::Duration;
struct State {
subscribers: Mutex<Vec<EventSender>>,
}
impl State {
pub fn new() -> Self {
Self {
subscribers: Mutex::new(Vec::new()),
}
}
}
#[allow(clippy::needless_pass_by_value)]
fn event_sender_thread(state: Arc<State>, permit: Permit) {
loop {
for n in 0..6 {
std::thread::sleep(Duration::from_secs(1));
if permit.is_revoked() {
return;
}
for subscriber in state.subscribers.lock().unwrap().iter_mut() {
subscriber.send(Event::Message(n.to_string()));
}
}
state.subscribers.lock().unwrap().clear();
}
}
#[allow(clippy::unnecessary_wraps)]
fn subscribe(state: &Arc<State>, _req: &Request) -> Result<Response, Response> {
let (sender, response) = Response::event_stream();
state.subscribers.lock().unwrap().push(sender);
Ok(response)
}
#[allow(clippy::unnecessary_wraps)]
fn handle_req(state: &Arc<State>, req: &Request) -> Result<Response, Response> {
match (req.method(), req.url().path()) {
("GET", "/health") => Ok(Response::text(200, "ok")),
("GET", "/subscribe") => subscribe(state, req),
_ => Ok(Response::text(404, "Not found")),
}
}
pub fn main() {
println!("Access the server at http://127.0.0.1:8000/subscribe");
let event_sender_thread_permit = Permit::new();
let state = Arc::new(State::new());
let state_clone = state.clone();
std::thread::spawn(move || event_sender_thread(state_clone, event_sender_thread_permit));
safina_timer::start_timer_thread();
let executor = safina_executor::Executor::default();
let request_handler = move |req: Request| print_log_response(&req, handle_req(&state, &req));
executor
.block_on(
HttpServerBuilder::new()
.listen_addr(socket_addr_127_0_0_1(8000))
.max_conns(100)
.spawn_and_join(request_handler),
)
.unwrap();
}