Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions src/content/body.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ use super::StreamError;
use crate::bug_message;
use bytes::Bytes;
use futures::Stream;
use futures::future::{Future, FutureExt, LocalBoxFuture};
use futures::future::{BoxFuture, Future, FutureExt};
use std::cmp;
use std::fs::File;
use std::io::ErrorKind::Interrupted;
Expand All @@ -21,7 +21,7 @@ use actix_web::error::BlockingError;
use actix_web::web;

type ChunkResult<T> = Result<Result<T, StreamError>, BlockingError>;
type ChunkOperation<'a, T> = LocalBoxFuture<'a, ChunkResult<T>>;
type ChunkOperation<'a, T> = BoxFuture<'a, ChunkResult<T>>;

fn flatten_chunk_result<T>(result: ChunkResult<T>) -> Result<T, StreamError> {
match result {
Expand Down Expand Up @@ -106,7 +106,7 @@ impl Stream for FileBody {
file.by_ref().take(max_bytes).read_to_end(&mut buffer)?;
Ok((file, Bytes::from(buffer)))
})
.boxed_local(),
.boxed(),
);
self.poll_next(context)
}
Expand Down Expand Up @@ -200,7 +200,7 @@ impl Stream for ProcessBody {
}
}
})
.boxed_local();
.boxed();

self.next = Some(next);
self.poll_next(context)
Expand Down
9 changes: 4 additions & 5 deletions src/content/handlebars_helpers/get.rs
Original file line number Diff line number Diff line change
Expand Up @@ -176,11 +176,10 @@ where
Ok(all_bytes)
},
);
// `FileBody`/`ProcessBody` rely on `actix_web::web::block`, which needs
// a Tokio runtime. When called from an HTTP handler one is already
// active, so `futures::executor::block_on` can drive the future to
// completion. From sync contexts (e.g. unit tests, the `eval`/`get`
// CLI subcommands) we spin up a fresh runtime to wait for the future.
// `block_on`, needs an active future runtime. When this is called from
// an HTTP handler one's already active, but in sync contexts (e.g. unit
// tests, the `eval`/`get` CLI subcommands) we spin up a fresh runtime
// to wait for the future.
let bytes = match tokio::runtime::Handle::try_current() {
Ok(_) => executor::block_on(bytes_future),
Err(_) => actix_rt::System::new().block_on(bytes_future),
Expand Down
4 changes: 2 additions & 2 deletions src/content/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,12 +29,12 @@ pub use content_registry::{ContentRepresentations, RegisteredContent};
pub use route::Route;

// This is just a trait alias to help make type signatures a bit saner.
pub trait ByteStream: Stream<Item = Result<Bytes, StreamError>>
pub trait ByteStream: Stream<Item = Result<Bytes, StreamError>> + Send
where
Self: Unpin,
{
}
impl<T> ByteStream for T where T: Stream<Item = Result<Bytes, StreamError>> + Unpin {}
impl<T> ByteStream for T where T: Stream<Item = Result<Bytes, StreamError>> + Unpin + Send {}

/// A piece of rendered content along with its media type.
pub struct Media<Content: ByteStream> {
Expand Down
Loading
Loading