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
70 changes: 29 additions & 41 deletions server/src/body/forward_body.rs
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
use ::actix::ResponseFuture;
use actix_web::http::{header, StatusCode};
use actix_web::{error::PayloadError, HttpMessage, HttpResponse, ResponseError};
use actix::ResponseFuture;
use actix_web::{error::PayloadError, http::StatusCode, HttpMessage, HttpResponse, ResponseError};
use bytes::{Bytes, BytesMut};
use failure::Fail;
use futures::prelude::*;

use crate::utils;

/// A set of errors that can occur during parsing json payloads
#[derive(Fail, Debug)]
pub enum ForwardPayloadError {
Expand Down Expand Up @@ -43,50 +44,36 @@ impl From<PayloadError> for ForwardPayloadError {
/// Future that resolves to a complete store endpoint body.
pub struct ForwardBody<T: HttpMessage> {
limit: usize,
length: Option<usize>,
stream: Option<T::Stream>,
err: Option<ForwardPayloadError>,
fut: Option<ResponseFuture<Bytes, ForwardPayloadError>>,
}

impl<T: HttpMessage> ForwardBody<T> {
/// Create `ForwardBody` for request.
pub fn new(req: &T) -> ForwardBody<T> {
let mut len = None;
if let Some(l) = req.headers().get(header::CONTENT_LENGTH) {
if let Ok(s) = l.to_str() {
if let Ok(l) = s.parse::<usize>() {
len = Some(l)
} else {
return Self::err(ForwardPayloadError::UnknownLength);
}
} else {
return Self::err(ForwardPayloadError::UnknownLength);
pub fn new(req: &T, limit: usize) -> ForwardBody<T> {
// Check the content length first. If we detect an overflow from the content length header,
// keep the payload in the request to drain it correctly in the `ReadRequestMiddleware`.
if let Some(length) = utils::get_content_length(req) {
if length > limit {
return Self::err(ForwardPayloadError::Overflow);
}
}

ForwardBody {
limit: 262_144,
length: len,
limit,
stream: Some(req.payload()),
fut: None,
err: None,
fut: None,
}
}

/// Change max size of payload. By default max size is 256Kb
pub fn limit(mut self, limit: usize) -> Self {
self.limit = limit;
self
}

fn err(e: ForwardPayloadError) -> Self {
ForwardBody {
limit: 0,
stream: None,
limit: 262_144,
fut: None,
err: Some(e),
length: None,
}
}
}
Expand All @@ -107,27 +94,28 @@ where
return Err(err);
}

if let Some(len) = self.length.take() {
if len > self.limit {
return Err(ForwardPayloadError::Overflow);
}
}

let limit = self.limit;
let body = Some(BytesMut::with_capacity(8192));

let future = self
.stream
.take()
.expect("Can not be used second time")
.from_err()
.fold(BytesMut::with_capacity(8192), move |mut body, chunk| {
if (body.len() + chunk.len()) > limit {
Err(ForwardPayloadError::Overflow)
} else {
body.extend_from_slice(&chunk);
Ok(body)
}
.map_err(ForwardPayloadError::from)
.fold(body, move |body_opt, chunk| {
Ok::<_, ForwardPayloadError>(body_opt.and_then(|mut body| {
if (body.len() + chunk.len()) > limit {
None
} else {
body.extend_from_slice(&chunk);
Some(body)
}
}))
})
.map(BytesMut::freeze);
.and_then(|bytes_opt| match bytes_opt {
Some(bytes) => Ok(bytes.freeze()),
None => Err(ForwardPayloadError::Overflow),
});

self.fut = Some(Box::new(future));

Expand Down
97 changes: 37 additions & 60 deletions server/src/body/store_body.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,8 @@ use std::borrow::Cow;
use std::io::{self, Read};

use actix::ResponseFuture;
use actix_web::http::{header, StatusCode};
use actix_web::HttpRequest;
use actix_web::{error::PayloadError, HttpMessage, HttpResponse, ResponseError};
use actix_web::http::StatusCode;
use actix_web::{error::PayloadError, HttpMessage, HttpRequest, HttpResponse, ResponseError};
use base64::DecodeError;
use bytes::{Bytes, BytesMut};
use failure::Fail;
Expand All @@ -15,6 +14,7 @@ use url::form_urlencoded;
use semaphore_common::metric;

use crate::actors::outcome::DiscardReason;
use crate::utils;

/// A set of errors that can occur during parsing json payloads
#[derive(Fail, Debug)]
Expand Down Expand Up @@ -75,55 +75,47 @@ impl From<PayloadError> for StorePayloadError {
/// Future that resolves to a complete store endpoint body.
pub struct StoreBody {
limit: usize,
length: Option<usize>,

// These states are mutually exclusive, and only separate options due to borrowing
// problems:
// These states are mutually exclusive:
result: Option<Result<Bytes, StorePayloadError>>,
fut: Option<ResponseFuture<Bytes, StorePayloadError>>,
stream: Option<<HttpRequest as HttpMessage>::Stream>,
}

impl StoreBody {
/// Create `StoreBody` for request.
pub fn new<S>(req: &HttpRequest<S>) -> Self {
pub fn new<S>(req: &HttpRequest<S>, limit: usize) -> Self {
if let Some(body) = data_from_querystring(req) {
StoreBody {
limit: 262_144,
length: Some(body.len()),
return StoreBody {
limit,
stream: None,
result: Some(decode_bytes(body.as_bytes())),
fut: None,
}
} else {
let length = match get_content_length(req) {
Ok(x) => x,
Err(e) => return StoreBody::err(e),
};
}

StoreBody {
limit: 262_144,
length,
result: None,
stream: Some(req.payload()),
fut: None,
// Check the content length first. If we detect an overflow from the content length header,
// keep the payload in the request to drain it correctly in the `ReadRequestMiddleware`.
if let Some(length) = utils::get_content_length(req) {
if length > limit {
return Self::err(StorePayloadError::Overflow);
}
}
}

/// Change max size of payload. By default max size is 256Kb
pub fn limit(mut self, limit: usize) -> Self {
self.limit = limit;
self
StoreBody {
limit,
result: None,
fut: None,
stream: Some(req.payload()),
}
}

fn err(e: StorePayloadError) -> Self {
StoreBody {
stream: None,
limit: 262_144,
limit: 0,
result: Some(Err(e)),
fut: None,
length: None,
stream: None,
}
}
}
Expand All @@ -141,27 +133,28 @@ impl Future for StoreBody {
return fut.poll();
}

if let Some(len) = self.length.take() {
if len > self.limit {
return Err(StorePayloadError::Overflow);
}
}

let limit = self.limit;
let body = Some(BytesMut::with_capacity(8192));

let future = self
.stream
.take()
.expect("Can not be used second time")
.from_err()
.fold(BytesMut::with_capacity(8192), move |mut body, chunk| {
if (body.len() + chunk.len()) > limit {
Err(StorePayloadError::Overflow)
} else {
body.extend_from_slice(&chunk);
Ok(body)
}
.map_err(StorePayloadError::from)
.fold(body, move |body_opt, chunk| {
// Ensure that the stream is always fully consumed. Erroring here would leave a
// broken TCP stream that cannot be used with keep-alive connections.
Ok::<_, StorePayloadError>(body_opt.and_then(|mut body| {
if (body.len() + chunk.len()) > limit {
None
} else {
body.extend_from_slice(&chunk);
Some(body)
}
}))
})
.and_then(|body| {
.and_then(|body_opt| {
let body = body_opt.ok_or(StorePayloadError::Overflow)?;
metric!(time_raw("event.size_bytes.raw") = body.len() as u64);
let decoded = decode_bytes(body.freeze())?;
metric!(time_raw("event.size_bytes.uncompressed") = decoded.len() as u64);
Expand Down Expand Up @@ -201,19 +194,3 @@ fn decode_bytes<B: Into<Bytes> + AsRef<[u8]>>(body: B) -> Result<Bytes, StorePay

Ok(Bytes::from(bytes))
}

fn get_content_length<S>(req: &HttpRequest<S>) -> Result<Option<usize>, StorePayloadError> {
if let Some(l) = req.headers().get(header::CONTENT_LENGTH) {
if let Ok(s) = l.to_str() {
if let Ok(l) = s.parse::<usize>() {
Ok(Some(l))
} else {
Err(StorePayloadError::UnknownLength)
}
} else {
Err(StorePayloadError::UnknownLength)
}
} else {
Ok(None)
}
}
3 changes: 1 addition & 2 deletions server/src/endpoints/forward.rs
Original file line number Diff line number Diff line change
Expand Up @@ -97,8 +97,7 @@ fn forward_upstream(request: &HttpRequest<ServiceState>) -> ResponseFuture<HttpR
.set_header("Connection", "close")
.timeout(config.http_timeout());

ForwardBody::new(request)
.limit(limit)
ForwardBody::new(request, limit)
.map_err(Error::from)
.and_then(move |data| forwarded_request_builder.body(data).map_err(Error::from))
.and_then(move |request| request.send().map_err(Error::from))
Expand Down
5 changes: 2 additions & 3 deletions server/src/endpoints/minidump.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ use futures::Future;

use semaphore_general::protocol::EventId;

use crate::body::ForwardBody;
use crate::endpoints::common::{handle_store_like_request, BadStoreRequest};
use crate::envelope::{AttachmentType, ContentType, Envelope, Item, ItemType};
use crate::extractors::{EventMeta, StartTime};
Expand Down Expand Up @@ -62,9 +63,7 @@ where {
// minidump can either be transmitted as request body, or as `upload_file_minidump` in a
// multipart formdata request.
if MINIDUMP_RAW_CONTENT_TYPES.contains(&request.content_type()) {
let future = request
.body()
.limit(max_payload_size)
let future = ForwardBody::new(request, max_payload_size)
.map_err(|_| BadStoreRequest::InvalidMinidump)
.and_then(move |data| {
validate_minidump(&data)?;
Expand Down
3 changes: 1 addition & 2 deletions server/src/endpoints/security_report.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,7 @@ fn extract_envelope(
max_event_payload_size: usize,
params: SecurityReportParams,
) -> ResponseFuture<Envelope, BadStoreRequest> {
let future = StoreBody::new(&request)
.limit(max_event_payload_size)
let future = StoreBody::new(&request, max_event_payload_size)
.map_err(BadStoreRequest::PayloadError)
.and_then(move |data| {
if data.is_empty() {
Expand Down
3 changes: 1 addition & 2 deletions server/src/endpoints/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,8 +33,7 @@ fn extract_envelope(
max_event_payload_size: usize,
content_type: String,
) -> ResponseFuture<Envelope, BadStoreRequest> {
let future = StoreBody::new(&request)
.limit(max_event_payload_size)
let future = StoreBody::new(&request, max_event_payload_size)
.map_err(BadStoreRequest::PayloadError)
.and_then(move |mut data| {
if data.is_empty() {
Expand Down
2 changes: 2 additions & 0 deletions server/src/utils/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ mod api;
mod error_boundary;
mod multipart;
mod param_parser;
mod request;
mod shutdown;
mod timer;

Expand All @@ -14,6 +15,7 @@ pub use self::api::*;
pub use self::error_boundary::*;
pub use self::multipart::*;
pub use self::param_parser::*;
pub use self::request::*;
pub use self::shutdown::*;
pub use self::timer::*;

Expand Down
Loading