use futures::sync::oneshot::Sender; use futures::unsync::oneshot; use futures::{Async, Future, Poll}; use smallvec::SmallVec; use std::marker::PhantomData; use std::mem; use actix::dev::{ContextImpl, SyncEnvelope, ToEnvelope}; use actix::fut::ActorFuture; use actix::{Actor, ActorContext, ActorState, Addr, AsyncContext, Handler, Message, SpawnHandle, Syn, Unsync}; use body::{Binary, Body}; use error::{Error, ErrorInternalServerError}; use httprequest::HttpRequest; pub trait ActorHttpContext: 'static { fn disconnected(&mut self); fn poll(&mut self) -> Poll>, Error>; } #[derive(Debug)] pub enum Frame { Chunk(Option), Drain(oneshot::Sender<()>), } impl Frame { pub fn len(&self) -> usize { match *self { Frame::Chunk(Some(ref bin)) => bin.len(), _ => 0, } } } /// Execution context for http actors pub struct HttpContext where A: Actor>, { inner: ContextImpl, stream: Option>, request: HttpRequest, disconnected: bool, } impl ActorContext for HttpContext where A: Actor, { fn stop(&mut self) { self.inner.stop(); } fn terminate(&mut self) { self.inner.terminate() } fn state(&self) -> ActorState { self.inner.state() } } impl AsyncContext for HttpContext where A: Actor, { #[inline] fn spawn(&mut self, fut: F) -> SpawnHandle where F: ActorFuture + 'static, { self.inner.spawn(fut) } #[inline] fn wait(&mut self, fut: F) where F: ActorFuture + 'static, { self.inner.wait(fut) } #[doc(hidden)] #[inline] fn waiting(&self) -> bool { self.inner.waiting() || self.inner.state() == ActorState::Stopping || self.inner.state() == ActorState::Stopped } #[inline] fn cancel_future(&mut self, handle: SpawnHandle) -> bool { self.inner.cancel_future(handle) } #[doc(hidden)] #[inline] fn unsync_address(&mut self) -> Addr { self.inner.unsync_address() } #[doc(hidden)] #[inline] fn sync_address(&mut self) -> Addr { self.inner.sync_address() } } impl HttpContext where A: Actor, { #[inline] pub fn new(req: HttpRequest, actor: A) -> HttpContext { HttpContext::from_request(req).actor(actor) } pub fn from_request(req: HttpRequest) -> HttpContext { HttpContext { inner: ContextImpl::new(None), stream: None, request: req, disconnected: false, } } #[inline] pub fn actor(mut self, actor: A) -> HttpContext { self.inner.set_actor(actor); self } } impl HttpContext where A: Actor, { /// Shared application state #[inline] pub fn state(&self) -> &S { self.request.state() } /// Incoming request #[inline] pub fn request(&mut self) -> &mut HttpRequest { &mut self.request } /// Write payload #[inline] pub fn write>(&mut self, data: B) { if !self.disconnected { self.add_frame(Frame::Chunk(Some(data.into()))); } else { warn!("Trying to write to disconnected response"); } } /// Indicate end of streaming payload. Also this method calls `Self::close`. #[inline] pub fn write_eof(&mut self) { self.add_frame(Frame::Chunk(None)); } /// Returns drain future pub fn drain(&mut self) -> Drain { let (tx, rx) = oneshot::channel(); self.inner.modify(); self.add_frame(Frame::Drain(tx)); Drain::new(rx) } /// Check if connection still open #[inline] pub fn connected(&self) -> bool { !self.disconnected } #[inline] fn add_frame(&mut self, frame: Frame) { if self.stream.is_none() { self.stream = Some(SmallVec::new()); } self.stream.as_mut().map(|s| s.push(frame)); self.inner.modify(); } /// Handle of the running future /// /// SpawnHandle is the handle returned by `AsyncContext::spawn()` method. pub fn handle(&self) -> SpawnHandle { self.inner.curr_handle() } } impl ActorHttpContext for HttpContext where A: Actor, S: 'static, { #[inline] fn disconnected(&mut self) { self.disconnected = true; self.stop(); } fn poll(&mut self) -> Poll>, Error> { let ctx: &mut HttpContext = unsafe { mem::transmute(self as &mut HttpContext) }; if self.inner.alive() { match self.inner.poll(ctx) { Ok(Async::NotReady) | Ok(Async::Ready(())) => (), Err(_) => return Err(ErrorInternalServerError("error")), } } // frames if let Some(data) = self.stream.take() { Ok(Async::Ready(Some(data))) } else if self.inner.alive() { Ok(Async::NotReady) } else { Ok(Async::Ready(None)) } } } impl ToEnvelope for HttpContext where A: Actor> + Handler, M: Message + Send + 'static, M::Result: Send, { fn pack(msg: M, tx: Option>) -> SyncEnvelope { SyncEnvelope::new(msg, tx) } } impl From> for Body where A: Actor>, S: 'static, { fn from(ctx: HttpContext) -> Body { Body::Actor(Box::new(ctx)) } } pub struct Drain { fut: oneshot::Receiver<()>, _a: PhantomData, } impl Drain { pub fn new(fut: oneshot::Receiver<()>) -> Self { Drain { fut, _a: PhantomData, } } } impl ActorFuture for Drain { type Item = (); type Error = (); type Actor = A; #[inline] fn poll( &mut self, _: &mut A, _: &mut ::Context ) -> Poll { self.fut.poll().map_err(|_| ()) } }