#![deny(warnings)]
use std::cell::Cell;
use std::rc::Rc;
use tokio::sync::oneshot;
use fluxio::body::{Bytes, HttpBody};
use fluxio::header::{HeaderMap, HeaderValue};
use fluxio::service::{make_service_fn, service_fn};
use fluxio::{Error, Response, Server};
use std::marker::PhantomData;
use std::pin::Pin;
use std::task::{Context, Poll};
struct Body {
_marker: PhantomData<*const ()>,
data: Option<Bytes>,
}
impl From<String> for Body {
fn from(a: String) -> Self {
Body {
_marker: PhantomData,
data: Some(a.into()),
}
}
}
impl HttpBody for Body {
type Data = Bytes;
type Error = Error;
fn poll_data(
self: Pin<&mut Self>,
_: &mut Context<'_>,
) -> Poll<Option<Result<Self::Data, Self::Error>>> {
Poll::Ready(self.get_mut().data.take().map(Ok))
}
fn poll_trailers(
self: Pin<&mut Self>,
_: &mut Context<'_>,
) -> Poll<Result<Option<HeaderMap<HeaderValue>>, Self::Error>> {
Poll::Ready(Ok(None))
}
}
fn main() {
pretty_env_logger::init();
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("build runtime");
let local = tokio::task::LocalSet::new();
local.block_on(&rt, run());
}
async fn run() {
let addr = ([127, 0, 0, 1], 3000).into();
let counter = Rc::new(Cell::new(0));
let make_service = make_service_fn(move |_| {
let cnt = counter.clone();
async move {
Ok::<_, Error>(service_fn(move |_| {
let prev = cnt.get();
cnt.set(prev + 1);
let value = cnt.get();
async move { Ok::<_, Error>(Response::new(Body::from(format!("Request #{}", value)))) }
}))
}
});
let server = Server::bind(&addr).executor(LocalExec).serve(make_service);
let (_tx, rx) = oneshot::channel::<()>();
let server = server.with_graceful_shutdown(async move {
rx.await.ok();
});
println!("Listening on http://{}", addr);
if let Err(e) = server.await {
eprintln!("server error: {}", e);
}
}
#[derive(Clone, Copy, Debug)]
struct LocalExec;
impl<F> fluxio::rt::Executor<F> for LocalExec
where
F: std::future::Future + 'static, {
fn execute(&self, fut: F) {
tokio::task::spawn_local(fut);
}
}