use super::*;
use boatramp_core::project::ProjectRef;
struct Visitor<'a> {
peer: IpAddr,
limiter: &'a dyn RateLimitStore,
}
pub(super) async fn serve_sites(
State(deploy): State<DeployStore>,
Extension(limiter): Extension<Arc<dyn RateLimitStore>>,
Extension(handlers): Extension<Arc<HandlerRuntime>>,
ConnectInfo(peer): ConnectInfo<SocketAddr>,
request: Request,
) -> Response {
let raw = request.uri().path();
let rest = raw
.strip_prefix("/_sites/")
.unwrap_or("")
.trim_start_matches('/');
let (site, path) = rest.split_once('/').unwrap_or((rest, ""));
if site.is_empty() {
return not_found();
}
let (site, request_path) = (site.to_string(), format!("/{path}"));
let visitor = Visitor {
peer: peer.ip(),
limiter: limiter.as_ref(),
};
serve_request(
&deploy,
ProjectRef::DEFAULT.as_str(),
&site,
&request_path,
request,
&visitor,
&handlers,
false,
)
.await
}
#[derive(Clone)]
pub(super) struct BootstrapAttestation(pub(super) Option<String>);
pub(super) async fn serve_bootstrap_identity(
Extension(att): Extension<BootstrapAttestation>,
) -> Response {
match att.0 {
Some(a) => (
StatusCode::OK,
[(header::CONTENT_TYPE, "application/octet-stream")],
a,
)
.into_response(),
None => not_found(),
}
}
pub(super) async fn serve_domain_challenge(
State(deploy): State<DeployStore>,
Extension(posture): Extension<boatramp_core::security::SecurityPosture>,
Path(token): Path<String>,
headers: HeaderMap,
) -> Response {
if !posture.domain_verify_self_serve {
return not_found();
}
let Some(host) = headers
.get(header::HOST)
.and_then(|value| value.to_str().ok())
.map(strip_port)
else {
return not_found();
};
match deploy
.find_pending_http_challenge(host, &token, now_unix())
.await
{
Ok(Some(v)) => (
StatusCode::OK,
[(header::CONTENT_TYPE, "text/plain; charset=utf-8")],
v.token,
)
.into_response(),
Ok(None) => not_found(),
Err(err) => deploy_error_response(err),
}
}
#[allow(clippy::too_many_arguments)] pub(super) async fn serve_by_host(
State(deploy): State<DeployStore>,
Extension(limiter): Extension<Arc<dyn RateLimitStore>>,
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(daemon): Extension<Arc<DaemonRuntime>>,
Extension(implicit): Extension<ImplicitRouting>,
Extension(preview_auth): Extension<Auth>,
ConnectInfo(peer): ConnectInfo<SocketAddr>,
request: Request,
) -> Response {
let effective = daemon.effective();
let preview_policy = PreviewPolicy {
protect: effective.protect_previews,
};
let Some(host) = request
.headers()
.get(header::HOST)
.and_then(|value| value.to_str().ok())
.map(strip_port)
.map(str::to_string)
else {
return not_found();
};
let request_path = request.uri().path().to_string();
if let Some((id_prefix, site_host)) = parse_deploy_host(&host) {
if let Some(blocked) =
preview_auth_gate(preview_policy, &preview_auth, request.headers()).await
{
return blocked;
}
return serve_host_preview(
&deploy,
&handlers,
peer.ip(),
&request_path,
request,
id_prefix,
site_host,
)
.await;
}
match deploy.resolve_site_by_host(&host).await {
Ok(Some(owner)) => {
let visitor = Visitor {
peer: peer.ip(),
limiter: limiter.as_ref(),
};
serve_request(
&deploy,
&owner.project,
&owner.site,
&request_path,
request,
&visitor,
&handlers,
true,
)
.await
}
Ok(None) => {
if effective.posture.require_domain_verification && !is_local_host(&host, implicit.0) {
return verification_pending_page(&deploy, &host).await;
}
if implicit.0 {
let label = host.split('.').next().unwrap_or("");
if !label.is_empty()
&& matches!(
deploy.current_id(ProjectRef::DEFAULT, label).await,
Ok(Some(_))
)
{
let visitor = Visitor {
peer: peer.ip(),
limiter: limiter.as_ref(),
};
return serve_request(
&deploy,
ProjectRef::DEFAULT.as_str(),
label,
&request_path,
request,
&visitor,
&handlers,
true,
)
.await;
}
}
match effective.default_site.as_deref() {
Some(site) => {
let visitor = Visitor {
peer: peer.ip(),
limiter: limiter.as_ref(),
};
serve_request(
&deploy,
ProjectRef::DEFAULT.as_str(),
site,
&request_path,
request,
&visitor,
&handlers,
true,
)
.await
}
None => not_found(),
}
}
Err(err) => deploy_error_response(err),
}
}
#[allow(clippy::too_many_arguments)]
async fn serve_host_preview(
deploy: &DeployStore,
handlers: &HandlerRuntime,
peer: IpAddr,
request_path: &str,
request: Request,
id_prefix: &str,
site_host: &str,
) -> Response {
let id = match deploy.resolve_manifest_id(id_prefix).await {
Ok(Some(id)) => id,
Ok(None) => return not_found(),
Err(err) => return deploy_error_response(err),
};
let owner = match deploy.resolve_site_by_host(site_host).await {
Ok(owner) => owner,
Err(err) => return deploy_error_response(err),
};
let project = owner.as_ref().map(|o| o.project.clone());
let site = owner.map(|o| o.site);
let site_config = match (&project, &site) {
(Some(project), Some(site)) => {
match deploy.get_site_config(ProjectRef::new(project), site).await {
Ok(config) => config,
Err(err) => return deploy_error_response(err),
}
}
_ => None,
};
match deploy.get_manifest(&id).await {
Ok(Some(manifest)) => {
serve_resolved(
deploy,
&manifest,
request_path,
request,
peer,
project.as_deref(),
site.as_deref(),
site_config.as_ref(),
handlers,
Some(&id),
)
.await
}
Ok(None) => not_found(),
Err(err) => deploy_error_response(err),
}
}
#[allow(clippy::too_many_arguments)]
async fn serve_request(
deploy: &DeployStore,
project: &str,
site: &str,
request_path: &str,
request: Request,
visitor: &Visitor<'_>,
handlers: &HandlerRuntime,
host_routed: bool,
) -> Response {
let project_ref = ProjectRef::new(project);
let site_config = match deploy.get_site_config(project_ref, site).await {
Ok(config) => config,
Err(err) => return deploy_error_response(err),
};
let listener_scheme = if request
.extensions()
.get::<ServedOverTls>()
.map(|s| s.0)
.unwrap_or(false)
{
"https"
} else {
"http"
};
let peer_trusted = site_config
.as_ref()
.map(|c| c.access.is_trusted_proxy(visitor.peer))
.unwrap_or(false);
let effective_scheme = if peer_trusted {
request
.headers()
.get("x-forwarded-proto")
.and_then(|v| v.to_str().ok())
.unwrap_or(listener_scheme)
.to_string()
} else {
listener_scheme.to_string()
};
#[cfg(feature = "compression")]
let accept_encoding = request
.headers()
.get(header::ACCEPT_ENCODING)
.and_then(|v| v.to_str().ok())
.map(str::to_string);
let mut security_headers: Vec<(HeaderName, String)> = Vec::new();
if host_routed {
if let Some(cfg) = site_config.as_ref() {
let host = request
.headers()
.get(header::HOST)
.and_then(|v| v.to_str().ok())
.map(strip_port)
.unwrap_or("");
let path_and_query = request
.uri()
.path_and_query()
.map(axum::http::uri::PathAndQuery::as_str)
.unwrap_or(request_path);
if let Some(target) = boatramp_core::config::transport_redirect(
&cfg.security,
&cfg.domains,
&effective_scheme,
host,
path_and_query,
) {
return redirect_to(&target);
}
if effective_scheme == "https" {
if let Some(hsts) = cfg.security.hsts.as_ref() {
security_headers.push((
HeaderName::from_static("strict-transport-security"),
hsts.header_value(),
));
}
}
if let Some(csp) = cfg.security.csp.as_deref() {
security_headers.push((header::CONTENT_SECURITY_POLICY, csp.to_string()));
}
if let Some(frame) = cfg.security.frame_options.as_deref() {
security_headers.push((header::X_FRAME_OPTIONS, frame.to_string()));
}
}
}
let access = site_config.as_ref().map(|c| &c.access);
let trusted = access.map(|a| a.trusted_proxies.as_slice()).unwrap_or(&[]);
let forwarded_for = request
.headers()
.get("x-forwarded-for")
.and_then(|value| value.to_str().ok());
let client_ip = boatramp_core::access::resolve_client_ip(visitor.peer, forwarded_for, trusted);
if let Some(access) = access {
if let Some(denied) = enforce_access(
access,
site,
request.headers(),
request_path,
client_ip,
visitor.limiter,
)
.await
{
return denied;
}
}
let manifest = match deploy.current_manifest(project_ref, site).await {
Ok(Some(manifest)) => manifest,
Ok(None) => return not_found(),
Err(err) => return deploy_error_response(err),
};
let mut response = serve_resolved(
deploy,
&manifest,
request_path,
request,
client_ip,
Some(project),
Some(site),
site_config.as_ref(),
handlers,
None,
)
.await;
for (name, value) in security_headers {
if let Ok(value) = HeaderValue::from_str(&value) {
response.headers_mut().insert(name, value);
}
}
#[cfg(feature = "compression")]
let response = match site_config.as_ref() {
Some(cfg) if cfg.compression.enabled => maybe_compress(
response,
accept_encoding.as_deref(),
cfg.compression.min_size,
),
_ => response,
};
response
}
async fn preview_auth_gate(
policy: PreviewPolicy,
auth: &Auth,
headers: &HeaderMap,
) -> Option<Response> {
if !policy.protect {
return None;
}
let bearer = headers
.get(header::AUTHORIZATION)
.and_then(|v| v.to_str().ok())
.and_then(|v| v.strip_prefix("Bearer "));
let ok = match bearer {
Some(token) => auth.verify_bearer(token).await,
None => false,
};
(!ok).then(|| {
(
StatusCode::UNAUTHORIZED,
"preview requires a valid bearer token\n",
)
.into_response()
})
}
fn redirect_to(target: &str) -> Response {
match HeaderValue::from_str(target) {
Ok(location) => (
StatusCode::MOVED_PERMANENTLY,
[(header::LOCATION, location)],
)
.into_response(),
Err(_) => (StatusCode::INTERNAL_SERVER_ERROR, "bad redirect target\n").into_response(),
}
}
pub fn http_redirect_router(
deploy: DeployStore,
posture: boatramp_core::security::SecurityPosture,
) -> Router {
Router::new()
.route(
"/.well-known/boatramp-domain-verification/{token}",
get(serve_domain_challenge),
)
.fallback(redirect_http_to_https)
.with_state(deploy)
.layer(Extension(posture))
}
async fn redirect_http_to_https(req: Request) -> Response {
let host = req
.headers()
.get(header::HOST)
.and_then(|v| v.to_str().ok())
.map(strip_port)
.unwrap_or("");
if host.is_empty() {
return (StatusCode::BAD_REQUEST, "missing Host header\n").into_response();
}
let path_and_query = req
.uri()
.path_and_query()
.map(axum::http::uri::PathAndQuery::as_str)
.unwrap_or("/");
match HeaderValue::from_str(&format!("https://{host}{path_and_query}")) {
Ok(location) => (
StatusCode::PERMANENT_REDIRECT,
[(header::LOCATION, location)],
)
.into_response(),
Err(_) => (StatusCode::BAD_REQUEST, "invalid host\n").into_response(),
}
}
async fn enforce_access(
access: &AccessConfig,
site: &str,
req_headers: &HeaderMap,
path: &str,
client_ip: IpAddr,
limiter: &dyn RateLimitStore,
) -> Option<Response> {
if !access.is_enforced() {
return None;
}
if access.waf.is_enabled() {
let header_str = |name| req_headers.get(name).and_then(|v| v.to_str().ok());
let waf_req = boatramp_core::waf::WafRequest {
user_agent: header_str(header::USER_AGENT),
accept: header_str(header::ACCEPT),
path,
};
if let boatramp_core::waf::WafVerdict::Block(reason) =
boatramp_core::waf::evaluate(&access.waf, &waf_req)
{
tracing::debug!(%client_ip, site, %reason, "request blocked by WAF");
return Some((StatusCode::FORBIDDEN, "forbidden\n").into_response());
}
}
if !access.ip.allows(client_ip) {
tracing::debug!(%client_ip, site, "request blocked by IP rules");
return Some((StatusCode::FORBIDDEN, "forbidden\n").into_response());
}
if let Some(limit) = &access.rate_limit {
if !limiter.check(site, client_ip, limit).await {
return Some(too_many_requests());
}
}
if let Some(basic) = &access.basic_auth {
if !verify_basic_auth(basic, req_headers) {
return Some(basic_auth_challenge(basic));
}
}
None
}
fn verify_basic_auth(basic: &BasicAuth, req_headers: &HeaderMap) -> bool {
use base64::Engine;
let Some(encoded) = req_headers
.get(header::AUTHORIZATION)
.and_then(|value| value.to_str().ok())
.and_then(|value| value.strip_prefix("Basic "))
else {
return false;
};
let Ok(decoded) = base64::engine::general_purpose::STANDARD.decode(encoded.trim()) else {
return false;
};
let Ok(text) = String::from_utf8(decoded) else {
return false;
};
match text.split_once(':') {
Some((user, pass)) => basic.verify(user, pass),
None => false,
}
}
fn basic_auth_challenge(basic: &BasicAuth) -> Response {
let realm = basic.realm.replace(['"', '\\'], "");
let mut headers = HeaderMap::new();
if let Ok(value) = HeaderValue::from_str(&format!("Basic realm=\"{realm}\", charset=\"UTF-8\""))
{
headers.insert(header::WWW_AUTHENTICATE, value);
}
(
StatusCode::UNAUTHORIZED,
headers,
"authentication required\n",
)
.into_response()
}
fn too_many_requests() -> Response {
let mut headers = HeaderMap::new();
headers.insert(header::RETRY_AFTER, HeaderValue::from_static("1"));
(
StatusCode::TOO_MANY_REQUESTS,
headers,
"rate limit exceeded\n",
)
.into_response()
}
pub(super) async fn serve_preview(
State(deploy): State<DeployStore>,
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(daemon): Extension<Arc<DaemonRuntime>>,
Extension(preview_auth): Extension<Auth>,
ConnectInfo(peer): ConnectInfo<SocketAddr>,
request: Request,
) -> Response {
let preview_policy = PreviewPolicy {
protect: daemon.effective().protect_previews,
};
if let Some(blocked) = preview_auth_gate(preview_policy, &preview_auth, request.headers()).await
{
return blocked;
}
let raw = request.uri().path();
let rest = raw
.strip_prefix("/_deploy/")
.unwrap_or("")
.trim_start_matches('/');
let (id, path) = rest.split_once('/').unwrap_or((rest, ""));
if id.is_empty() {
return not_found();
}
let (id, request_path) = (id.to_string(), format!("/{path}"));
let site = request
.headers()
.get(header::HOST)
.and_then(|value| value.to_str().ok())
.map(strip_port);
let owner = match site {
Some(host) => match deploy.resolve_site_by_host(host).await {
Ok(owner) => owner,
Err(err) => return deploy_error_response(err),
},
None => None,
};
let project = owner.as_ref().map(|o| o.project.clone());
let site = owner.map(|o| o.site);
let site_config = match (&project, &site) {
(Some(project), Some(site)) => {
match deploy.get_site_config(ProjectRef::new(project), site).await {
Ok(config) => config,
Err(err) => return deploy_error_response(err),
}
}
_ => None,
};
match deploy.get_manifest(&id).await {
Ok(Some(manifest)) => {
serve_resolved(
&deploy,
&manifest,
&request_path,
request,
peer.ip(),
project.as_deref(),
site.as_deref(),
site_config.as_ref(),
&handlers,
Some(&id),
)
.await
}
Ok(None) => not_found(),
Err(err) => deploy_error_response(err),
}
}
fn build_request_context(request: &Request) -> boatramp_core::predicate::RequestContext {
use boatramp_core::predicate::RequestContext;
let headers = request.headers();
let mut hmap: std::collections::BTreeMap<String, String> = std::collections::BTreeMap::new();
for (name, value) in headers {
if let Ok(v) = value.to_str() {
hmap.entry(name.as_str().to_ascii_lowercase())
.and_modify(|e| {
e.push_str(", ");
e.push_str(v);
})
.or_insert_with(|| v.to_string());
}
}
let host = headers
.get(header::HOST)
.and_then(|h| h.to_str().ok())
.map(|h| h.split(':').next().unwrap_or(h).to_string())
.unwrap_or_default();
let cookies = headers
.get(header::COOKIE)
.and_then(|h| h.to_str().ok())
.map(parse_cookie_header)
.unwrap_or_default();
let query = request
.uri()
.query()
.map(parse_query_string)
.unwrap_or_default();
let accept_languages = headers
.get(header::ACCEPT_LANGUAGE)
.and_then(|h| h.to_str().ok())
.map(RequestContext::parse_accept_language)
.unwrap_or_default();
RequestContext {
method: request.method().as_str().to_ascii_uppercase(),
host,
headers: hmap,
cookies,
query,
accept_languages,
}
}
pub(super) fn parse_cookie_header(raw: &str) -> std::collections::BTreeMap<String, String> {
raw.split(';')
.filter_map(|pair| pair.split_once('='))
.map(|(k, v)| (k.trim().to_string(), v.trim().to_string()))
.fold(std::collections::BTreeMap::new(), |mut m, (k, v)| {
m.entry(k).or_insert(v);
m
})
}
pub(super) fn parse_query_string(raw: &str) -> std::collections::BTreeMap<String, String> {
raw.split('&')
.filter(|p| !p.is_empty())
.map(|pair| match pair.split_once('=') {
Some((k, v)) => (percent_decode(k), percent_decode(v)),
None => (percent_decode(pair), String::new()),
})
.fold(std::collections::BTreeMap::new(), |mut m, (k, v)| {
m.entry(k).or_insert(v);
m
})
}
fn percent_decode(s: &str) -> String {
let bytes = s.as_bytes();
let mut out: Vec<u8> = Vec::with_capacity(bytes.len());
let mut i = 0;
while i < bytes.len() {
match bytes[i] {
b'+' => {
out.push(b' ');
i += 1;
}
b'%' if i + 2 < bytes.len() => {
let hi = (bytes[i + 1] as char).to_digit(16);
let lo = (bytes[i + 2] as char).to_digit(16);
match (hi, lo) {
(Some(h), Some(l)) => {
out.push((h * 16 + l) as u8);
i += 3;
}
_ => {
out.push(b'%');
i += 1;
}
}
}
b => {
out.push(b);
i += 1;
}
}
}
String::from_utf8_lossy(&out).into_owned()
}
pub(super) fn apply_vary(mut response: Response, vary: &[String]) -> Response {
if vary.is_empty() {
return response;
}
let mut names: Vec<String> = response
.headers()
.get(header::VARY)
.and_then(|v| v.to_str().ok())
.map(|s| {
s.split(',')
.map(|p| p.trim().to_ascii_lowercase())
.filter(|p| !p.is_empty())
.collect()
})
.unwrap_or_default();
for v in vary {
if !names.iter().any(|n| n == v) {
names.push(v.clone());
}
}
if let Ok(hv) = HeaderValue::from_str(&names.join(", ")) {
response.headers_mut().insert(header::VARY, hv);
}
response
}
#[allow(clippy::too_many_arguments)]
#[cfg_attr(not(feature = "handlers"), allow(unused_variables))]
async fn serve_resolved(
deploy: &DeployStore,
manifest: &Manifest,
request_path: &str,
request: Request,
client_ip: IpAddr,
project: Option<&str>,
site: Option<&str>,
site_config: Option<&SiteConfig>,
handlers: &HandlerRuntime,
preview: Option<&str>,
) -> Response {
let ctx = if manifest.config.redirects.iter().any(|r| r.when.is_some())
|| manifest.config.rewrites.iter().any(|r| r.when.is_some())
{
build_request_context(&request)
} else {
boatramp_core::predicate::RequestContext::default()
};
let route::ResolveResult { outcome, vary } =
route::resolve_ctx(&manifest.config, &manifest.files, request_path, &ctx);
#[cfg(feature = "handlers")]
if !matches!(outcome, Outcome::Redirect { .. }) {
if let Some(site) = site {
if let Some(handler) = route::match_handler(
&manifest.config.handlers,
request.method().as_str(),
request_path,
) {
return apply_vary(
dispatch_handler(
handlers,
deploy,
manifest,
project.unwrap_or(ProjectRef::DEFAULT.as_str()),
site,
request_path,
site_config,
handler,
request,
client_ip,
preview,
)
.await,
&vary,
);
}
if request.method() == Method::GET {
if let Some(stream) = manifest
.config
.streams
.iter()
.find(|s| route_matches(&s.route, request_path))
{
if let (Some(inner), Some(site_handlers)) = (
handlers.inner.as_ref(),
site_config
.and_then(|c| c.handlers.as_ref())
.filter(|h| h.enabled),
) {
if stream.websocket && is_upgrade_request(request.headers()) {
use axum::extract::FromRequestParts;
let (mut parts, _body) = request.into_parts();
return apply_vary(
match axum::extract::ws::WebSocketUpgrade::from_request_parts(
&mut parts,
&(),
)
.await
{
Ok(ws) => {
serve_ws_stream(
inner,
site,
site_handlers,
stream,
ws,
client_ip,
preview,
)
.await
}
Err(rejection) => rejection.into_response(),
},
&vary,
);
}
let after = request
.headers()
.get("last-event-id")
.and_then(|value| value.to_str().ok())
.map(str::to_string);
return apply_vary(
serve_stream(
inner,
site,
site_handlers,
stream,
after,
client_ip,
preview,
)
.await,
&vary,
);
}
return apply_vary(not_found(), &vary);
}
}
}
}
if !matches!(outcome, Outcome::Redirect { .. }) {
if let Some(gw) = site_config
.and_then(|c| c.gateway.as_ref())
.filter(|g| g.is_enabled())
{
if let Some(route) = gw.match_route(request_path) {
return apply_vary(
match gw.upstreams.get(&route.upstream) {
Some(upstream) => {
let (compute_backends, compute_regions) = match &upstream.compute {
Some(workload) => {
gateway::record_activity(workload);
let mut pool = compute_endpoints(deploy, workload).await;
if pool.is_empty() && has_parked_replica(deploy, workload).await
{
gateway::wake_reconcile();
pool = await_warm(deploy, workload, COMPUTE_WAKE_TIMEOUT)
.await;
}
let regions = if upstream.lb
== boatramp_core::gateway::LbPolicy::Nearest
{
Some(compute_endpoint_regions(deploy, workload).await)
} else {
None
};
(Some(pool), regions)
}
None => (None, None),
};
dispatch_gateway(
request,
site.unwrap_or(""),
&route.upstream,
upstream,
request_path,
client_ip,
compute_backends,
compute_regions,
)
.await
}
None => (
StatusCode::BAD_GATEWAY,
"gateway route references an unknown upstream\n",
)
.into_response(),
},
&vary,
);
}
}
}
let response = match outcome {
Outcome::Redirect { location, status } => redirect(status, &location),
Outcome::Proxy { url } => proxy(request, &url, &manifest.config, client_ip).await,
Outcome::File {
path: served,
entry,
} => {
if !matches!(*request.method(), Method::GET | Method::HEAD) {
return apply_vary(method_not_allowed(), &vary);
}
serve_entry(
deploy,
&manifest.config,
request_path,
&served,
&entry,
request.headers(),
StatusCode::OK,
)
.await
}
Outcome::NotFound { error } => match error {
Some((served, entry)) => {
serve_entry(
deploy,
&manifest.config,
request_path,
&served,
&entry,
request.headers(),
StatusCode::NOT_FOUND,
)
.await
}
None => not_found(),
},
};
apply_vary(response, &vary)
}
fn method_not_allowed() -> Response {
let mut headers = HeaderMap::new();
headers.insert(header::ALLOW, HeaderValue::from_static("GET, HEAD"));
(
StatusCode::METHOD_NOT_ALLOWED,
headers,
"method not allowed\n",
)
.into_response()
}
#[allow(clippy::too_many_arguments)]
async fn serve_entry(
deploy: &DeployStore,
config: &DeployConfig,
request_path: &str,
served_path: &str,
entry: &FileEntry,
req_headers: &HeaderMap,
base_status: StatusCode,
) -> Response {
let is_range = base_status == StatusCode::OK && req_headers.contains_key(header::RANGE);
let chosen = if is_range {
None
} else {
negotiate_encoding(entry, req_headers)
};
let (blob_hash, blob_size, encoding) = match chosen {
Some((enc, variant)) => (variant.hash.as_str(), variant.size, Some(enc)),
None => (entry.hash.as_str(), entry.size, None),
};
let etag = format!("\"{blob_hash}\"");
if base_status == StatusCode::OK && if_none_match(req_headers, &etag) {
let mut headers = response_headers(config, request_path, served_path, entry, &etag);
set_content_encoding(&mut headers, encoding);
return (StatusCode::NOT_MODIFIED, headers).into_response();
}
if is_range {
if let Some(spec) = req_headers
.get(header::RANGE)
.and_then(|value| value.to_str().ok())
{
match parse_ranges(spec, entry.size) {
Some(ranges) if ranges.len() == 1 => {
let (offset, len) = ranges[0];
let object = match deploy.open_blob_range(&entry.hash, offset, Some(len)).await
{
Ok(object) => object,
Err(err) => return deploy_error_response(err),
};
let mut headers =
response_headers(config, request_path, served_path, entry, &etag);
set_header(&mut headers, header::CONTENT_LENGTH, &len.to_string());
set_header(
&mut headers,
header::CONTENT_RANGE,
&format!("bytes {}-{}/{}", offset, offset + len - 1, entry.size),
);
return (
StatusCode::PARTIAL_CONTENT,
headers,
Body::from_stream(object.body),
)
.into_response();
}
Some(ranges) if ranges.len() <= MAX_RANGES => {
return multipart_byteranges(
deploy,
config,
request_path,
served_path,
entry,
&etag,
&ranges,
)
.await;
}
Some(_) => {}
None => {
let mut headers = HeaderMap::new();
set_header(
&mut headers,
header::CONTENT_RANGE,
&format!("bytes */{}", entry.size),
);
return (StatusCode::RANGE_NOT_SATISFIABLE, headers).into_response();
}
}
}
}
let object = match deploy.open_blob(blob_hash).await {
Ok(object) => object,
Err(err) => return deploy_error_response(err),
};
let mut headers = response_headers(config, request_path, served_path, entry, &etag);
set_header(&mut headers, header::CONTENT_LENGTH, &blob_size.to_string());
set_content_encoding(&mut headers, encoding);
(base_status, headers, Body::from_stream(object.body)).into_response()
}