extern crate url;
#[cfg(feature = "native-tls")]
extern crate native_tls_crate as native_tls;
#[cfg(test)]
extern crate http_test_server;
#[cfg(feature = "native-tls")]
mod tls;
mod network;
mod pub_sub;
mod data;
use std::sync::{Arc, Mutex, mpsc };
use url::{Url, ParseError};
use network::EventStream;
use pub_sub::Bus;
use data::{EventBuilder, EventBuilderState};
pub use data::Event;
pub use network::State;
pub struct EventSource {
bus: Arc<Mutex<Bus<Event>>>,
stream: Arc<Mutex<EventStream>>
}
impl EventSource {
pub fn new(url: &str) -> Result<EventSource, ParseError> {
let event_stream = Arc::new(Mutex::new(EventStream::new(Url::parse(url)?).unwrap()));
let stream_for_update = Arc::clone(&event_stream);
let stream = Arc::clone(&event_stream);
let mut event_stream = event_stream.lock().unwrap();
let bus = Arc::new(Mutex::new(Bus::new()));
let event_bus = Arc::clone(&bus);
event_stream.on_open(move || {
publish_initial_stream_event(&event_bus);
});
let event_bus = Arc::clone(&bus);
event_stream.on_error(move |message| {
let event_bus = event_bus.lock().unwrap();
let event = Event::new("error", &message);
event_bus.publish(event.type_.clone(), event);
});
let event_builder = Arc::new(Mutex::new(EventBuilder::new()));
let event_bus = Arc::clone(&bus);
event_stream.on_message(move |message| {
handle_message(&message, &event_builder, &event_bus, &stream_for_update);
});
Ok(EventSource{ stream, bus })
}
pub fn close(&self) {
self.stream.lock().unwrap().close();
}
pub fn on_open<F>(&self, listener: F) where F: Fn() + Send + 'static {
self.add_event_listener("stream_opened", move |_| { listener(); });
}
pub fn on_message<F>(&self, listener: F) where F: Fn(Event) + Send + 'static {
self.add_event_listener("message", listener);
}
pub fn add_event_listener<F>(&self, event_type: &str, listener: F) where F: Fn(Event) + Send + 'static {
let mut bus = self.bus.lock().unwrap();
bus.subscribe(event_type.to_string(), listener);
}
pub fn state(&self) -> State {
self.stream.lock().unwrap().state()
}
pub fn receiver(&self) -> mpsc::Receiver<Event> {
let (tx, rx) = mpsc::channel();
let error_tx = tx.clone();
self.on_message(move |event| {
tx.send(event).unwrap();
});
self.add_event_listener("error", move |error| {
error_tx.send(error).unwrap();
});
rx
}
}
fn publish_initial_stream_event(event_bus: &Arc<Mutex<Bus<Event>>>) {
let event_bus = event_bus.lock().unwrap();
let event = Event::new("stream_opened", "");
event_bus.publish(event.type_.clone(), event);
}
fn handle_message(
message: &str,
event_builder: &Arc<Mutex<EventBuilder>>,
event_bus: &Arc<Mutex<Bus<Event>>>,
event_stream: &Arc<Mutex<EventStream>>) {
let mut event_builder = event_builder.lock().unwrap();
if let EventBuilderState::Complete(event) = event_builder.update(&message) {
let event_bus = event_bus.lock().unwrap();
event_stream.lock().unwrap().set_last_id(event.id.clone());
event_bus.publish(event.type_.clone(), event);
event_builder.clear();
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::thread;
use std::time::Duration;
use std::sync::mpsc;
use http_test_server::{ TestServer, Resource };
use http_test_server::http::Status;
fn setup() -> (TestServer, Resource, String) {
let server = TestServer::new().unwrap();
let resource = server.create_resource("/sub");
resource.header("Content-Type", "text/event-stream").stream();
let address = format!("http://localhost:{}/sub", server.port());
thread::sleep(Duration::from_millis(100));
(server, resource, address)
}
#[test]
fn should_create_client() {
let (_server, _stream_endpoint, address) = setup();
let event_source = EventSource::new(&address).unwrap();
event_source.close();
}
#[test]
fn should_thrown_an_error_when_malformed_url_provided() {
match EventSource::new("127.0.0.1:1236/sub") {
Ok(_) => assert!(false, "should had thrown an error"),
Err(_) => assert!(true)
}
}
#[test]
fn accept_closure_as_listeners() {
let (tx, rx) = mpsc::channel();
let (_server, stream_endpoint, address) = setup();
let event_source = EventSource::new(&address).unwrap();
event_source.on_message(move |message| {
tx.send(message.data).unwrap();
});
while event_source.state() == State::Connecting {
thread::sleep(Duration::from_millis(100));
}
stream_endpoint
.send_line("data: some message").send_line("");
let message = rx.recv().unwrap();
assert_eq!(message, "some message");
event_source.close();
}
#[test]
fn should_trigger_listeners_when_message_received() {
let (tx, rx) = mpsc::channel();
let tx2 = tx.clone();
let (_server, stream_endpoint, address) = setup();
let event_source = EventSource::new(&address).unwrap();
event_source.on_message(move |message| {
tx.send(message.data).unwrap();
});
event_source.on_message(move |message| {
tx2.send(message.data).unwrap();
});
while event_source.state() == State::Connecting {
thread::sleep(Duration::from_millis(100));
}
stream_endpoint.send_line("data: some message").send_line("");
let message = rx.recv().unwrap();
let message2 = rx.recv().unwrap();
assert_eq!(message, "some message");
assert_eq!(message2, "some message");
event_source.close();
}
#[test]
fn should_not_trigger_listeners_for_comments() {
let (tx, rx) = mpsc::channel();
let (_server, stream_endpoint, address) = setup();
let event_source = EventSource::new(&address).unwrap();
event_source.on_message(move |message| {
tx.send(message.data).unwrap();
});
while event_source.state() != State::Open {
thread::sleep(Duration::from_millis(100));
};
stream_endpoint
.send_line("data: message")
.send_line("")
.send_line(":this is a comment")
.send_line(":this is another comment")
.send_line("data: this is a message")
.send_line("");
let message = rx.recv().unwrap();
let message2 = rx.recv().unwrap();
assert_eq!(message, "message");
assert_eq!(message2, "this is a message");
event_source.close();
}
#[test]
fn ignore_empty_messages() {
let (tx, rx) = mpsc::channel();
let (_server, stream_endpoint, address) = setup();
let event_source = EventSource::new(&address).unwrap();
event_source.on_message(move |message| {
tx.send(message.data).unwrap();
});
while event_source.state() != State::Open {
thread::sleep(Duration::from_millis(100));
};
stream_endpoint
.send_line("data: message")
.send_line("")
.send_line("")
.send_line("data: this is a message")
.send_line("");
let message = rx.recv().unwrap();
let message2 = rx.recv().unwrap();
assert_eq!(message, "message");
assert_eq!(message2, "this is a message");
event_source.close();
}
#[test]
fn event_trigger_its_defined_listener() {
let (tx, rx) = mpsc::channel();
let (_server, stream_endpoint, address) = setup();
let event_source = EventSource::new(&address).unwrap();
event_source.add_event_listener("myEvent", move |event| {
tx.send(event).unwrap();
});
while event_source.state() == State::Connecting {
thread::sleep(Duration::from_millis(100));
}
stream_endpoint
.send_line("event: myEvent")
.send_line("data: my message\n");
let message = rx.recv().unwrap();
assert_eq!(message.type_, String::from("myEvent"));
assert_eq!(message.data, String::from("my message"));
event_source.close();
}
#[test]
fn dont_trigger_on_message_for_event() {
let (tx, rx) = mpsc::channel();
let (_server, stream_endpoint, address) = setup();
let event_source = EventSource::new(&address).unwrap();
event_source.on_message(move |_| {
tx.send("NOOOOOOOOOOOOOOOOOOO!").unwrap();
});
stream_endpoint
.send("event: myEvent\n")
.send("data: my message\n\n");
thread::sleep(Duration::from_millis(500));
assert!(rx.try_recv().is_err());
event_source.close();
}
#[test]
fn should_close_connection() {
let (tx, rx) = mpsc::channel();
let (_server, stream_endpoint, address) = setup();
let event_source = EventSource::new(&address).unwrap();
event_source.on_message(move |message| {
tx.send(message.data).unwrap();
});
while event_source.state() != State::Open {
thread::sleep(Duration::from_millis(100));
};
stream_endpoint.send("\ndata: some message\n\n");
rx.recv().unwrap();
event_source.close();
stream_endpoint.send("\ndata: some message\n\n");
thread::sleep(Duration::from_millis(400));
assert!(rx.try_recv().is_err());
}
#[test]
fn should_trigger_on_open_callback_when_connected() {
let (tx, rx) = mpsc::channel();
let (_server, stream_endpoint, address) = setup();
stream_endpoint.delay(Duration::from_millis(200));
let event_source = EventSource::new(&address).unwrap();
event_source.on_open(move || {
tx.send("open").unwrap();
});
rx.recv().unwrap();
event_source.close();
}
#[test]
fn should_return_stream_connection_status() {
let (_server, stream_endpoint, address) = setup();
stream_endpoint
.delay(Duration::from_millis(200))
.stream();
let event_source = EventSource::new(&address).unwrap();
thread::sleep(Duration::from_millis(100));
assert_eq!(event_source.state(), State::Connecting);
thread::sleep(Duration::from_millis(200));
assert_eq!(event_source.state(), State::Open);
event_source.close();
thread::sleep(Duration::from_millis(200));
assert_eq!(event_source.state(), State::Closed);
}
#[test]
fn should_send_last_event_id_on_reconnection() {
let (server, stream_endpoint, address) = setup();
let event_source = EventSource::new(&address).unwrap();
thread::sleep(Duration::from_millis(100));
stream_endpoint.send("id: helpMe\n");
stream_endpoint.send("data: my message\n\n");
thread::sleep(Duration::from_millis(500));
stream_endpoint.close_open_connections();
let request = server.requests().recv().unwrap();
assert_eq!(request.headers.get("Last-Event-ID").unwrap(), "helpMe");
event_source.close();
}
#[test]
fn should_expose_blocking_api() {
let (_server, stream_endpoint, address) = setup();
let event_source = EventSource::new(&address).unwrap();
thread::sleep(Duration::from_millis(100));
let rx = event_source.receiver();
stream_endpoint.send("data: some message\n\n");
stream_endpoint.send("data: some message 2\n\n");
assert_eq!(rx.recv().unwrap().data, "some message");
assert_eq!(rx.recv().unwrap().data, "some message 2");
event_source.close();
}
#[test]
fn receiver_should_get_error_events() {
let (_server, stream_endpoint, address) = setup();
stream_endpoint
.delay(Duration::from_millis(100))
.status(Status::InternalServerError);
let event_source = EventSource::new(&address).unwrap();
let rx = event_source.receiver();
assert_eq!(rx.recv().unwrap().type_, "error");
event_source.close();
}
}