use std::cmp::Ordering as CmpOrdering;
use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::fs::{self, File, OpenOptions};
use std::io::{self, ErrorKind};
use std::net::SocketAddr;
use std::os::unix::fs::FileExt;
#[cfg(any(target_os = "linux", target_os = "android"))]
use std::os::unix::fs::OpenOptionsExt;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex, RwLock};
use std::thread;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use anyhow::{Context, Result};
use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
use base64::Engine;
use nexus_core::device_scheduler::DeviceScheduler;
use nexus_core::pipeline::{
process_payload_crc_zstd_aes, restore_chunk_zstd_aes, ChunkController, ChunkPolicy,
EncryptedChunk, RuntimeMetrics,
};
use nexus_core::topology::TopologyReport;
use nexus_core::wal::{
self, DevicePointer, MultipartPart, MultipartWalEntryV1, WalEntryV1, WalReplayConfig, WalWriter,
};
use crossbeam_channel::bounded;
use futures::TryStreamExt;
use http_body_util::BodyExt;
use hyper::service::service_fn;
use hyper::{body::Incoming as HyperIncoming, http, http::StatusCode};
use hyper_util::rt::{TokioExecutor, TokioIo};
use hyper_util::server::conn::auto::Builder as ConnBuilder;
use include_dir::{include_dir, Dir};
use s3s::dto::*;
use s3s::service::S3ServiceBuilder;
use s3s::{s3_error, Body, HttpResponse, S3Request, S3Response, S3Result};
use serde::{Deserialize, Serialize};
use serde_json::json;
use socket2::{Domain, Protocol, Socket, Type};
use uuid::Uuid;
use crate::admin::{self, PurgeDataConfig};
use crate::access_control::GatewayAccessControl;
use crate::atc::{self, AtcDecision, AtcForceMode, IoModel, RuntimeProfile, WorkerModel};
use crate::auth_provider::{CredentialKind, CredentialStore, DynamicS3Auth};
use crate::compression::{self, CompressionDecision, CompressionPolicy};
use crate::iam::IamState;
use crate::ingest::BucketQosProfile;
use crate::migration::{self, BifrostSourceConfig, MigrationBucketPolicy, MigrationState};
use crate::multipart::{compute_multipart_etag, md5_hex, sanitize_etag};
use crate::sts::StsService;
use crate::webhooks::{
WebhookConfig, WebhookDeliveryResult, WebhookDeliveryTask, WebhookEngine, WebhookRule,
};
const ALIGNMENT: usize = 4096;
const ALIGNMENT_U64: u64 = ALIGNMENT as u64;
const ADAPTIVE_EVAL_PERIOD: Duration = Duration::from_secs(10);
const PIPELINE_KEY: [u8; 32] = [0x7A; 32];
const DEFAULT_ACCESS_KEY: &str = "AKIAtracerbullet";
const DEFAULT_SECRET_KEY: &str = "SECRETtracerbullet";
const RECORD_VERSION: u32 = 1;
const DEFAULT_OWNER_ID: &str = "0000000000000000000000000000000000000000000000000000000000000000";
const DEFAULT_OWNER_DISPLAY_NAME: &str = "nexusd";
const REQUEST_TIMEOUT_MS: u64 = 30_000;
static PANOPTICON_DIST: Dir<'_> = include_dir!("$CARGO_MANIFEST_DIR/ui-dist");
#[derive(Debug, Clone)]
pub struct ServeConfig {
pub listen_addr: String,
pub workers: usize,
pub devices: Vec<PathBuf>,
pub wal_device: PathBuf,
pub topology_policy: String,
pub runtime_profile: String,
pub chunk_policy: String,
pub atc_force_mode: String,
pub atc_single_device_threshold: usize,
pub atc_buffered_fsync_bytes: u64,
pub atc_buffered_fsync_ms: u64,
pub allow_sync_fallback: bool,
pub require_io_uring: bool,
pub webhook_workers: usize,
pub webhook_retry_max: u32,
pub webhook_timeout_ms: u64,
pub webhook_queue_capacity: usize,
pub atc_dev_lazy_commit: bool,
pub atc_dev_fsync_bytes: u64,
pub atc_dev_fsync_ms: u64,
pub bifrost_source_endpoint: Option<String>,
pub bifrost_source_region: String,
pub bifrost_source_access_key: Option<String>,
pub bifrost_source_secret_key: Option<String>,
}
#[derive(Debug, Clone)]
pub struct ValidationConfig {
pub s3_endpoint: String,
pub access_key: String,
pub secret_key: String,
pub json: bool,
}
#[derive(Debug, Clone)]
pub struct AdminPurgeConfig {
pub danger: bool,
pub bucket_prefix: Option<String>,
pub wal_path: Option<PathBuf>,
pub devices: Vec<PathBuf>,
pub json: bool,
}
#[derive(Debug, Clone)]
pub struct MigrationImportConfig {
pub config_path: PathBuf,
pub iam_dir: PathBuf,
pub json: bool,
}
#[derive(Debug, Clone)]
pub struct MigrationSyncConfig {
pub bucket: String,
pub prefix: String,
pub workers: usize,
pub json: bool,
}
#[derive(Debug, Clone)]
pub struct MigrationStatusConfig {
pub json: bool,
}
#[derive(Debug)]
struct DeviceIo {
id: usize,
path: PathBuf,
mode: IoModel,
buffered_sync: Option<Mutex<BufferedSyncState>>,
file: Mutex<File>,
}
#[derive(Debug)]
struct DataPlane {
scheduler: DeviceScheduler,
devices: Vec<Arc<DeviceIo>>,
io_model: IoModel,
buffered_sync_policy: BufferedSyncPolicy,
dev_lazy_commit: bool,
buffered_flushes: Arc<AtomicU64>,
buffered_flush_latency_ms: Arc<AtomicU64>,
mode_summary: String,
}
#[derive(Debug)]
struct BufferedSyncState {
pending_bytes: u64,
last_sync: Instant,
}
#[derive(Debug, Clone, Copy)]
struct BufferedSyncPolicy {
fsync_bytes: u64,
fsync_every: Duration,
}
impl DataPlane {
fn open(
paths: &[PathBuf],
io_model: IoModel,
allow_sync_fallback: bool,
buffered_sync_policy: BufferedSyncPolicy,
dev_lazy_commit: bool,
) -> Result<Self> {
let scheduler = DeviceScheduler::new(paths)?;
let mut devices = Vec::with_capacity(paths.len());
let mut modes = BTreeSet::new();
let buffered_flushes = Arc::new(AtomicU64::new(0));
let buffered_flush_latency_ms = Arc::new(AtomicU64::new(0));
for path in paths {
let (file, mode, resolved_path) = open_device(path, io_model, allow_sync_fallback)?;
modes.insert(mode.as_str().to_string());
let id = devices.len();
devices.push(Arc::new(DeviceIo {
id,
path: resolved_path,
mode,
buffered_sync: match mode {
IoModel::Buffered => Some(Mutex::new(BufferedSyncState {
pending_bytes: 0,
last_sync: Instant::now(),
})),
IoModel::DirectSync => None,
},
file: Mutex::new(file),
}));
}
Ok(Self {
scheduler,
devices,
io_model,
buffered_sync_policy,
dev_lazy_commit,
buffered_flushes,
buffered_flush_latency_ms,
mode_summary: modes.into_iter().collect::<Vec<_>>().join("+"),
})
}
fn append_blob(&self, blob: &[u8], nonce_id: u64) -> Result<DevicePointer> {
let logical_len_u32: u32 = blob
.len()
.try_into()
.context("blob too large for pointer len")?;
let scheduled = self
.scheduler
.select_and_reserve(blob.len() as u64, ALIGNMENT_U64);
let device_state = scheduled.device.clone();
let device = self
.devices
.get(device_state.id)
.context("device scheduler selected invalid device id")?;
let padded_len = align_up(blob.len(), ALIGNMENT);
let mut write_buf = vec![0_u8; padded_len];
write_buf[..blob.len()].copy_from_slice(blob);
device_state.increment_queue_depth();
let result = (|| {
let mut file = device
.file
.lock()
.map_err(|_| anyhow::anyhow!("device file mutex poisoned"))?;
write_all_at(&mut file, &write_buf, scheduled.reserved_offset)
.context("data append write_at failed")?;
match self.io_model {
IoModel::DirectSync => {
file.sync_data().context("data append sync_data failed")?;
}
IoModel::Buffered => {
let sync_state = device
.buffered_sync
.as_ref()
.context("buffered sync state missing for buffered mode")?;
let mut state = sync_state
.lock()
.map_err(|_| anyhow::anyhow!("buffered sync state mutex poisoned"))?;
state.pending_bytes = state.pending_bytes.saturating_add(write_buf.len() as u64);
let now = Instant::now();
let should_sync = if self.dev_lazy_commit {
state.pending_bytes >= self.buffered_sync_policy.fsync_bytes
|| now.duration_since(state.last_sync)
>= self.buffered_sync_policy.fsync_every
} else {
true
};
if should_sync {
let sync_started = Instant::now();
file.sync_data().context("data append buffered sync_data failed")?;
state.pending_bytes = 0;
state.last_sync = now;
self.buffered_flushes.fetch_add(1, Ordering::Relaxed);
self.buffered_flush_latency_ms.fetch_add(
sync_started.elapsed().as_millis() as u64,
Ordering::Relaxed,
);
}
}
}
Ok::<(), anyhow::Error>(())
})();
device_state.decrement_queue_depth();
result?;
Ok(DevicePointer {
device_id: device_state.id as u16,
offset: scheduled.reserved_offset,
len: logical_len_u32,
checksum_crc32c: crc32c::crc32c(blob),
nonce_id,
})
}
fn read_blob(&self, pointer: &DevicePointer) -> Result<Vec<u8>> {
let device_id = pointer.device_id as usize;
let device = self
.devices
.get(device_id)
.with_context(|| format!("invalid device id {}", pointer.device_id))?;
let mut file = device
.file
.lock()
.map_err(|_| anyhow::anyhow!("device file mutex poisoned"))?;
read_exact_at(&mut file, pointer.offset, pointer.len as usize)
.with_context(|| format!("read failed for device {}", device.path.display()))
}
fn median_queue_depth(&self) -> f64 {
let mut depths = self
.scheduler
.telemetry()
.into_iter()
.map(|(_id, _path, depth)| depth as f64)
.collect::<Vec<_>>();
if depths.is_empty() {
return 0.0;
}
depths.sort_by(|a, b| a.partial_cmp(b).unwrap_or(CmpOrdering::Equal));
percentile_sorted(&depths, 50.0)
}
fn device_report(&self) -> Vec<String> {
self.devices
.iter()
.map(|dev| {
format!(
"id={} path={} mode={}",
dev.id,
dev.path.display(),
dev.mode.as_str()
)
})
.collect()
}
}
#[derive(Debug)]
struct ChunkRuntime {
inner: Mutex<ChunkRuntimeInner>,
current_chunk_bytes: AtomicUsize,
downshift_events: AtomicU64,
upshift_events: AtomicU64,
}
#[derive(Debug)]
struct ChunkRuntimeInner {
controller: ChunkController,
last_eval: Instant,
put_latencies_ms: Vec<f64>,
queue_depths: Vec<f64>,
}
impl ChunkRuntime {
fn new(controller: ChunkController) -> Self {
let current = controller.current_chunk_bytes();
Self {
inner: Mutex::new(ChunkRuntimeInner {
controller,
last_eval: Instant::now(),
put_latencies_ms: Vec::new(),
queue_depths: Vec::new(),
}),
current_chunk_bytes: AtomicUsize::new(current),
downshift_events: AtomicU64::new(0),
upshift_events: AtomicU64::new(0),
}
}
fn current_chunk_bytes(&self) -> usize {
self.current_chunk_bytes.load(Ordering::Acquire)
}
fn observe_put(&self, latency_ms: f64, queue_depth: f64) {
let mut guard = match self.inner.lock() {
Ok(guard) => guard,
Err(_) => return,
};
guard.put_latencies_ms.push(latency_ms);
guard.queue_depths.push(queue_depth);
if guard.last_eval.elapsed() < ADAPTIVE_EVAL_PERIOD {
return;
}
if guard.put_latencies_ms.is_empty() || guard.queue_depths.is_empty() {
guard.last_eval = Instant::now();
return;
}
let mut latencies = guard.put_latencies_ms.clone();
let mut depths = guard.queue_depths.clone();
latencies.sort_by(|a, b| a.partial_cmp(b).unwrap_or(CmpOrdering::Equal));
depths.sort_by(|a, b| a.partial_cmp(b).unwrap_or(CmpOrdering::Equal));
let p99 = percentile_sorted(&latencies, 99.0);
let qmed = percentile_sorted(&depths, 50.0);
let before = guard.controller.current_chunk_bytes();
guard.controller.observe(RuntimeMetrics {
put_p99_ms: p99,
median_queue_depth: qmed,
});
let after = guard.controller.current_chunk_bytes();
if after != before {
self.current_chunk_bytes.store(after, Ordering::Release);
if after < before {
self.downshift_events.fetch_add(1, Ordering::Relaxed);
} else {
self.upshift_events.fetch_add(1, Ordering::Relaxed);
}
}
guard.put_latencies_ms.clear();
guard.queue_depths.clear();
guard.last_eval = Instant::now();
}
}
#[derive(Debug, Clone)]
enum StoredLayout {
InlineBlob { pointer: DevicePointer },
Multipart { parts: Vec<MultipartPart> },
}
#[derive(Debug, Clone)]
struct StoredObjectMeta {
etag: String,
content_len: u64,
last_modified: Timestamp,
metadata: HashMap<String, String>,
layout: StoredLayout,
}
#[derive(Debug, Clone)]
struct MultipartPartState {
etag: String,
pointer: DevicePointer,
size: u64,
}
#[derive(Debug, Clone)]
struct MultipartUploadState {
bucket: String,
key: String,
parts: BTreeMap<u32, MultipartPartState>,
}
#[derive(Clone)]
struct GatewayState {
data_plane: Arc<DataPlane>,
wal: Arc<WalWriter>,
wal_sequence_lock: Arc<Mutex<()>>,
credentials: Arc<CredentialStore>,
iam: Arc<IamState>,
sts: Arc<StsService>,
webhooks: Arc<WebhookEngine>,
migration: Arc<MigrationState>,
nonce_counter: Arc<AtomicU64>,
objects: Arc<RwLock<BTreeMap<String, StoredObjectMeta>>>,
uploads: Arc<RwLock<BTreeMap<String, MultipartUploadState>>>,
buckets: Arc<RwLock<BTreeSet<String>>>,
compression_policies: Arc<RwLock<BTreeMap<String, CompressionPolicy>>>,
qos_profiles: Arc<RwLock<BTreeMap<String, BucketQosProfile>>>,
atc_decision: AtcDecision,
chunk_runtime: Arc<ChunkRuntime>,
request_count: Arc<AtomicU64>,
puts: Arc<AtomicU64>,
gets: Arc<AtomicU64>,
lists: Arc<AtomicU64>,
deletes: Arc<AtomicU64>,
mpu_creates: Arc<AtomicU64>,
mpu_parts: Arc<AtomicU64>,
mpu_completes: Arc<AtomicU64>,
mpu_aborts: Arc<AtomicU64>,
last_compression_decision: Arc<RwLock<Option<CompressionDecision>>>,
buffered_sync_flushes: Arc<AtomicU64>,
buffered_sync_flush_latency_ms: Arc<AtomicU64>,
}
#[derive(Clone)]
struct GatewayS3 {
state: GatewayState,
}
#[async_trait::async_trait]
impl s3s::S3 for GatewayS3 {
async fn create_bucket(
&self,
req: S3Request<CreateBucketInput>,
) -> S3Result<S3Response<CreateBucketOutput>> {
let bucket = req.input.bucket;
validate_bucket_name(&bucket)?;
let mut buckets = self
.state
.buckets
.write()
.map_err(|_| s3_error!(InternalError))?;
if buckets.contains(&bucket) {
return Err(s3_error!(BucketAlreadyExists));
}
buckets.insert(bucket);
Ok(S3Response::new(CreateBucketOutput::default()))
}
async fn head_bucket(
&self,
req: S3Request<HeadBucketInput>,
) -> S3Result<S3Response<HeadBucketOutput>> {
ensure_bucket_exists(&self.state, &req.input.bucket)?;
Ok(S3Response::new(HeadBucketOutput::default()))
}
async fn get_bucket_acl(
&self,
req: S3Request<GetBucketAclInput>,
) -> S3Result<S3Response<GetBucketAclOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
ensure_bucket_exists(&self.state, &req.input.bucket)?;
Ok(S3Response::new(GetBucketAclOutput {
owner: Some(default_owner()),
grants: Some(default_acl_grants()),
}))
}
async fn put_bucket_acl(
&self,
req: S3Request<PutBucketAclInput>,
) -> S3Result<S3Response<PutBucketAclOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
ensure_bucket_exists(&self.state, &req.input.bucket)?;
Ok(S3Response::new(PutBucketAclOutput::default()))
}
async fn delete_bucket(
&self,
req: S3Request<DeleteBucketInput>,
) -> S3Result<S3Response<DeleteBucketOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
let bucket = req.input.bucket;
ensure_bucket_exists(&self.state, &bucket)?;
let bucket_prefix = storage_key_prefix(&bucket, "");
if self
.state
.objects
.read()
.map_err(|_| s3_error!(InternalError))?
.keys()
.any(|key| key.starts_with(&bucket_prefix))
{
return Err(s3_error!(BucketNotEmpty));
}
self.state
.buckets
.write()
.map_err(|_| s3_error!(InternalError))?
.remove(&bucket);
self.state
.uploads
.write()
.map_err(|_| s3_error!(InternalError))?
.retain(|_, upload| upload.bucket != bucket);
Ok(S3Response::new(DeleteBucketOutput::default()))
}
async fn list_buckets(
&self,
_req: S3Request<ListBucketsInput>,
) -> S3Result<S3Response<ListBucketsOutput>> {
let buckets = self
.state
.buckets
.read()
.map_err(|_| s3_error!(InternalError))?;
let mut items = Vec::with_capacity(buckets.len());
for bucket in buckets.iter() {
items.push(Bucket {
name: Some(bucket.clone()),
creation_date: Some(Timestamp::from(SystemTime::now())),
..Default::default()
});
}
Ok(S3Response::new(ListBucketsOutput {
buckets: Some(items),
..Default::default()
}))
}
async fn get_bucket_location(
&self,
req: S3Request<GetBucketLocationInput>,
) -> S3Result<S3Response<GetBucketLocationOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
ensure_bucket_exists(&self.state, &req.input.bucket)?;
Ok(S3Response::new(GetBucketLocationOutput::default()))
}
async fn list_objects(
&self,
req: S3Request<ListObjectsInput>,
) -> S3Result<S3Response<ListObjectsOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
self.state.lists.fetch_add(1, Ordering::Relaxed);
let input = req.input;
ensure_bucket_exists(&self.state, &input.bucket)?;
let prefix = input.prefix.clone().unwrap_or_default();
let (max_keys_i32, max_keys) = normalize_max_keys(input.max_keys)?;
let delimiter = input.delimiter.clone().unwrap_or_default();
let has_delimiter = !delimiter.is_empty();
let encode_keys = input
.encoding_type
.as_ref()
.map(|v| v.as_str() == EncodingType::URL)
.unwrap_or(false);
let storage_prefix = storage_key_prefix(&input.bucket, &prefix);
let objects = self
.state
.objects
.read()
.map_err(|_| s3_error!(InternalError))?;
let mut rows = Vec::new();
for (full_key, meta) in objects
.iter()
.filter(|(full_key, _)| full_key.starts_with(&storage_prefix))
{
rows.push((full_key.clone(), meta.clone()));
}
rows.sort_by(|a, b| a.0.cmp(&b.0));
let marker = input.marker.clone();
if let Some(marker) = marker.clone() {
let marker_key = storage_key(&input.bucket, &marker);
rows.retain(|(full_key, _)| full_key > &marker_key);
}
enum ListingEntry {
Object {
key: String,
etag: String,
size: i64,
last_modified: Timestamp,
},
Prefix(String),
}
let mut entries = Vec::<ListingEntry>::new();
let mut seen_prefixes = BTreeSet::<String>::new();
for (full_key, meta) in rows {
let key = strip_storage_prefix(&input.bucket, &full_key)
.ok_or_else(|| s3_error!(InternalError, "invalid key prefix"))?;
if has_delimiter {
let tail = key.strip_prefix(&prefix).unwrap_or(&key);
if let Some(idx) = tail.find(delimiter.as_str()) {
let common = format!("{}{}{}", prefix, &tail[..idx], delimiter);
if marker.as_ref().map(|m| common <= *m).unwrap_or(false) {
continue;
}
if seen_prefixes.insert(common.clone()) {
entries.push(ListingEntry::Prefix(common));
}
continue;
}
}
entries.push(ListingEntry::Object {
key,
etag: meta.etag,
size: meta.content_len as i64,
last_modified: meta.last_modified,
});
}
let is_truncated = max_keys != 0 && entries.len() > max_keys;
if max_keys == 0 {
entries.clear();
} else if is_truncated {
entries.truncate(max_keys);
}
let mut contents = Vec::<Object>::new();
let mut common_prefixes = Vec::<CommonPrefix>::new();
let mut next_marker = None::<String>;
for entry in entries {
match entry {
ListingEntry::Object {
key,
etag,
size,
last_modified,
} => {
let display_key = if encode_keys {
s3_url_encode_key(&key)
} else {
key.clone()
};
contents.push(Object {
key: Some(display_key),
size: Some(size),
e_tag: Some(ETag::Strong(etag)),
last_modified: Some(last_modified),
owner: Some(default_owner()),
..Default::default()
});
next_marker = Some(key);
}
ListingEntry::Prefix(prefix_value) => {
common_prefixes.push(CommonPrefix {
prefix: Some(if encode_keys {
s3_url_encode_key(&prefix_value)
} else {
prefix_value.clone()
}),
});
next_marker = Some(prefix_value);
}
}
}
let response_delimiter = input
.delimiter
.as_ref()
.filter(|value| !value.is_empty())
.cloned();
let response_marker = Some(input.marker.clone().unwrap_or_default());
Ok(S3Response::new(ListObjectsOutput {
name: Some(input.bucket),
prefix: input.prefix,
marker: response_marker,
max_keys: Some(max_keys_i32),
is_truncated: Some(is_truncated),
contents: (!contents.is_empty()).then_some(contents),
common_prefixes: (!common_prefixes.is_empty()).then_some(common_prefixes),
delimiter: response_delimiter,
next_marker: is_truncated.then_some(next_marker.unwrap_or_default()),
encoding_type: input.encoding_type,
..Default::default()
}))
}
async fn list_object_versions(
&self,
req: S3Request<ListObjectVersionsInput>,
) -> S3Result<S3Response<ListObjectVersionsOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
self.state.lists.fetch_add(1, Ordering::Relaxed);
let input = req.input;
ensure_bucket_exists(&self.state, &input.bucket)?;
let prefix = input.prefix.clone().unwrap_or_default();
let (max_keys_i32, max_keys) = normalize_max_keys(input.max_keys)?;
let storage_prefix = storage_key_prefix(&input.bucket, &prefix);
let objects = self
.state
.objects
.read()
.map_err(|_| s3_error!(InternalError))?;
let mut rows = Vec::new();
for (full_key, meta) in objects
.iter()
.filter(|(full_key, _)| full_key.starts_with(&storage_prefix))
{
rows.push((full_key.clone(), meta.clone()));
}
rows.sort_by(|a, b| a.0.cmp(&b.0));
if let Some(key_marker) = input.key_marker.clone() {
let marker_key = storage_key(&input.bucket, &key_marker);
rows.retain(|(full_key, _)| full_key > &marker_key);
}
let is_truncated = max_keys != 0 && rows.len() > max_keys;
if max_keys == 0 {
rows.clear();
} else {
rows.truncate(max_keys);
}
let mut versions = Vec::with_capacity(rows.len());
for (full_key, meta) in rows {
let key = strip_storage_prefix(&input.bucket, &full_key)
.ok_or_else(|| s3_error!(InternalError, "invalid key prefix"))?;
versions.push(ObjectVersion {
key: Some(key),
size: Some(meta.content_len as i64),
e_tag: Some(ETag::Strong(meta.etag)),
last_modified: Some(meta.last_modified),
is_latest: Some(true),
version_id: Some("null".to_string()),
..Default::default()
});
}
let next_key_marker = is_truncated
.then(|| versions.last().and_then(|version| version.key.clone()))
.flatten();
let next_version_id_marker = is_truncated.then_some("null".to_string());
Ok(S3Response::new(ListObjectVersionsOutput {
name: Some(input.bucket),
prefix: input.prefix,
key_marker: input.key_marker,
version_id_marker: input.version_id_marker,
max_keys: Some(max_keys_i32),
is_truncated: Some(is_truncated),
versions: (!versions.is_empty()).then_some(versions),
next_key_marker,
next_version_id_marker,
delimiter: input.delimiter,
encoding_type: input.encoding_type,
..Default::default()
}))
}
async fn put_object(
&self,
req: S3Request<PutObjectInput>,
) -> S3Result<S3Response<PutObjectOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
self.state.puts.fetch_add(1, Ordering::Relaxed);
let started = Instant::now();
let input = req.input;
let bucket = input.bucket.clone();
let key = input.key.clone();
let mut metadata = input.metadata.unwrap_or_default();
if let Some(content_type) = input.content_type.clone() {
metadata.insert("content-type".to_string(), content_type);
}
ensure_bucket_exists(&self.state, &bucket)?;
let object_key = storage_key(&bucket, &key);
let is_update = self
.state
.objects
.read()
.map_err(|_| s3_error!(InternalError))?
.contains_key(&object_key);
let Some(body) = input.body else {
return Err(s3_error!(IncompleteBody));
};
let raw = read_streaming_blob(body)
.await
.map_err(|_| s3_error!(InvalidRequest, "failed reading request body"))?;
validate_payload_headers(input.content_md5.as_deref(), input.content_length, &raw)?;
let nonce_seed = next_nonce_seed(&self.state).saturating_mul(1_000_000);
let compression_policy = self
.state
.compression_policies
.read()
.ok()
.and_then(|m| m.get(&bucket).cloned())
.unwrap_or_default();
let decision = compression::decide(None, &raw, &compression_policy);
if let Ok(mut guard) = self.state.last_compression_decision.write() {
*guard = Some(decision.clone());
}
let qos_profile = self
.state
.qos_profiles
.read()
.ok()
.and_then(|m| m.get(&bucket).copied())
.unwrap_or(BucketQosProfile::Standard);
let active_mpu = self
.state
.uploads
.read()
.map(|g| g.len())
.unwrap_or(0);
let chunk_bytes = crate::ingest::compute_chunk_bytes(
self.state.chunk_runtime.current_chunk_bytes(),
qos_profile,
active_mpu,
);
let processed = process_payload_crc_zstd_aes(&raw, chunk_bytes, &PIPELINE_KEY, nonce_seed)
.map_err(|_| s3_error!(InternalError, "pipeline processing failed"))?;
let blob = encode_pipeline_blob(&processed.chunks)
.map_err(|_| s3_error!(InternalError, "pipeline blob encoding failed"))?;
let pointer = self
.state
.data_plane
.append_blob(&blob, nonce_seed)
.map_err(|_| s3_error!(InternalError, "device append failed"))?;
let etag_value = md5_hex(&raw);
append_object_wal_ordered(
&self.state,
WalEntryV1 {
seq: 0,
op: "put".to_string(),
bucket: bucket.clone(),
key: key.clone(),
etag: Some(etag_value.clone()),
pointer_or_manifest: Some(pointer.clone()),
ts_unix_ms: now_unix_ms(),
},
)?;
let meta = StoredObjectMeta {
etag: etag_value.clone(),
content_len: raw.len() as u64,
last_modified: Timestamp::from(SystemTime::now()),
metadata,
layout: StoredLayout::InlineBlob { pointer },
};
self.state
.objects
.write()
.map_err(|_| s3_error!(InternalError))?
.insert(object_key, meta);
enqueue_webhook_intents(
&self.state,
&bucket,
&key,
raw.len() as u64,
etag_value.as_str(),
is_update,
)?;
let latency_ms = started.elapsed().as_secs_f64() * 1000.0;
let qdepth = self.state.data_plane.median_queue_depth();
self.state.chunk_runtime.observe_put(latency_ms, qdepth);
Ok(S3Response::new(PutObjectOutput {
e_tag: Some(ETag::Strong(etag_value)),
..Default::default()
}))
}
async fn get_object(
&self,
req: S3Request<GetObjectInput>,
) -> S3Result<S3Response<GetObjectOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
self.state.gets.fetch_add(1, Ordering::Relaxed);
let input = req.input;
let bucket = input.bucket.clone();
let key = input.key.clone();
let meta = self
.state
.objects
.read()
.map_err(|_| s3_error!(InternalError))?
.get(&storage_key(&bucket, &key))
.cloned();
if meta.is_none() {
if let Some(proxy) = maybe_proxy_read_and_backfill(&self.state, &bucket, &key)? {
let mut body = proxy.bytes;
let mut content_range = None;
if let Some(range) = input.range.clone() {
let checked = range
.check(body.len() as u64)
.map_err(|_| s3_error!(InvalidRange))?;
let start = checked.start as usize;
let end = checked.end as usize;
body = body[start..end].to_vec();
content_range = Some(format!(
"bytes {}-{}/{}",
checked.start,
checked.end.saturating_sub(1),
proxy.content_len
));
}
let content_len_i64 =
i64::try_from(body.len()).map_err(|_| s3_error!(InternalError))?;
return Ok(S3Response::new(GetObjectOutput {
body: Some(StreamingBlob::from(Body::from(body))),
content_length: Some(content_len_i64),
content_range,
e_tag: Some(ETag::Strong(proxy.etag)),
last_modified: Some(Timestamp::from(SystemTime::now())),
metadata: Some(HashMap::new()),
..Default::default()
}));
}
return Err(s3_error!(NoSuchKey));
}
let meta = meta.expect("checked is_some");
validate_get_preconditions(&input, &meta)?;
let mut restored = materialize_object(&self.state, &meta)
.map_err(|_| s3_error!(InternalError, "failed materializing object"))?;
let mut content_range = None;
if let Some(range) = input.range {
let checked = range
.check(restored.len() as u64)
.map_err(|_| s3_error!(InvalidRange))?;
let start = checked.start as usize;
let end = checked.end as usize;
restored = restored[start..end].to_vec();
content_range = Some(format!(
"bytes {}-{}/{}",
checked.start,
checked.end.saturating_sub(1),
meta.content_len
));
}
let content_len_i64 =
i64::try_from(restored.len()).map_err(|_| s3_error!(InternalError))?;
Ok(S3Response::new(GetObjectOutput {
body: Some(StreamingBlob::from(Body::from(restored))),
content_length: Some(content_len_i64),
content_range,
e_tag: Some(ETag::Strong(meta.etag)),
last_modified: Some(meta.last_modified),
content_type: browser_content_type_from_metadata(&meta.metadata, &input.key),
metadata: Some(meta.metadata),
..Default::default()
}))
}
async fn head_object(
&self,
req: S3Request<HeadObjectInput>,
) -> S3Result<S3Response<HeadObjectOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
self.state.gets.fetch_add(1, Ordering::Relaxed);
let input = req.input;
let meta = self
.state
.objects
.read()
.map_err(|_| s3_error!(InternalError))?
.get(&storage_key(&input.bucket, &input.key))
.cloned()
.ok_or_else(|| s3_error!(NoSuchKey))?;
Ok(S3Response::new(HeadObjectOutput {
content_length: Some(meta.content_len as i64),
e_tag: Some(ETag::Strong(meta.etag)),
last_modified: Some(meta.last_modified),
content_type: browser_content_type_from_metadata(&meta.metadata, &input.key),
metadata: Some(meta.metadata),
..Default::default()
}))
}
async fn get_object_acl(
&self,
req: S3Request<GetObjectAclInput>,
) -> S3Result<S3Response<GetObjectAclOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
self.state.gets.fetch_add(1, Ordering::Relaxed);
let input = req.input;
ensure_bucket_exists(&self.state, &input.bucket)?;
let object_key = storage_key(&input.bucket, &input.key);
if !self
.state
.objects
.read()
.map_err(|_| s3_error!(InternalError))?
.contains_key(&object_key)
{
return Err(s3_error!(NoSuchKey));
}
Ok(S3Response::new(GetObjectAclOutput {
owner: Some(default_owner()),
grants: Some(default_acl_grants()),
request_charged: None,
}))
}
async fn put_object_acl(
&self,
req: S3Request<PutObjectAclInput>,
) -> S3Result<S3Response<PutObjectAclOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
let input = req.input;
ensure_bucket_exists(&self.state, &input.bucket)?;
let object_key = storage_key(&input.bucket, &input.key);
if !self
.state
.objects
.read()
.map_err(|_| s3_error!(InternalError))?
.contains_key(&object_key)
{
return Err(s3_error!(NoSuchKey));
}
Ok(S3Response::new(PutObjectAclOutput::default()))
}
async fn list_objects_v2(
&self,
req: S3Request<ListObjectsV2Input>,
) -> S3Result<S3Response<ListObjectsV2Output>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
self.state.lists.fetch_add(1, Ordering::Relaxed);
let input = req.input;
ensure_bucket_exists(&self.state, &input.bucket)?;
let prefix = input.prefix.clone().unwrap_or_default();
let (max_keys_i32, max_keys) = normalize_max_keys(input.max_keys)?;
let include_owner = input.fetch_owner.unwrap_or(false);
let delimiter = input.delimiter.clone().unwrap_or_default();
let has_delimiter = !delimiter.is_empty();
let encode_keys = input
.encoding_type
.as_ref()
.map(|v| v.as_str() == EncodingType::URL)
.unwrap_or(false);
let storage_prefix = storage_key_prefix(&input.bucket, &prefix);
let objects = self
.state
.objects
.read()
.map_err(|_| s3_error!(InternalError))?;
let mut rows = Vec::new();
for (full_key, meta) in objects
.iter()
.filter(|(full_key, _)| full_key.starts_with(&storage_prefix))
{
rows.push((full_key.clone(), meta.clone()));
}
let marker_token = input
.continuation_token
.clone()
.or(input.start_after.clone());
if let Some(token) = marker_token.clone() {
let token_key = storage_key(&input.bucket, &token);
rows.retain(|(full_key, _)| full_key > &token_key);
}
rows.sort_by(|a, b| a.0.cmp(&b.0));
enum ListingEntry {
Object {
key: String,
etag: String,
size: i64,
last_modified: Timestamp,
},
Prefix(String),
}
let mut entries = Vec::<ListingEntry>::new();
let mut seen_prefixes = BTreeSet::<String>::new();
for (full_key, meta) in rows {
let key = strip_storage_prefix(&input.bucket, &full_key)
.ok_or_else(|| s3_error!(InternalError, "invalid key prefix"))?;
if has_delimiter {
let tail = key.strip_prefix(&prefix).unwrap_or(&key);
if let Some(idx) = tail.find(delimiter.as_str()) {
let common = format!("{}{}{}", prefix, &tail[..idx], delimiter);
if marker_token.as_ref().map(|m| common <= *m).unwrap_or(false) {
continue;
}
if seen_prefixes.insert(common.clone()) {
entries.push(ListingEntry::Prefix(common));
}
continue;
}
}
entries.push(ListingEntry::Object {
key,
etag: meta.etag,
size: meta.content_len as i64,
last_modified: meta.last_modified,
});
}
let is_truncated = max_keys != 0 && entries.len() > max_keys;
if max_keys == 0 {
entries.clear();
} else if is_truncated {
entries.truncate(max_keys);
}
let mut contents = Vec::<Object>::new();
let mut common_prefixes = Vec::<CommonPrefix>::new();
let mut next_token = None::<String>;
for entry in entries {
match entry {
ListingEntry::Object {
key,
etag,
size,
last_modified,
} => {
let display_key = if encode_keys {
s3_url_encode_key(&key)
} else {
key.clone()
};
contents.push(Object {
key: Some(display_key),
size: Some(size),
e_tag: Some(ETag::Strong(etag)),
last_modified: Some(last_modified),
owner: include_owner.then(default_owner),
..Default::default()
});
next_token = Some(key);
}
ListingEntry::Prefix(prefix_value) => {
common_prefixes.push(CommonPrefix {
prefix: Some(if encode_keys {
s3_url_encode_key(&prefix_value)
} else {
prefix_value.clone()
}),
});
next_token = Some(prefix_value);
}
}
}
let key_count = i32::try_from(contents.len() + common_prefixes.len())
.map_err(|_| s3_error!(InternalError))?;
let response_delimiter = input
.delimiter
.as_ref()
.filter(|value| !value.is_empty())
.cloned();
Ok(S3Response::new(ListObjectsV2Output {
name: Some(input.bucket),
prefix: input.prefix,
key_count: Some(key_count),
max_keys: Some(max_keys_i32),
is_truncated: Some(is_truncated),
contents: (!contents.is_empty()).then_some(contents),
common_prefixes: (!common_prefixes.is_empty()).then_some(common_prefixes),
delimiter: response_delimiter,
continuation_token: input.continuation_token,
next_continuation_token: is_truncated.then_some(next_token.unwrap_or_default()),
start_after: input.start_after,
encoding_type: input.encoding_type,
..Default::default()
}))
}
async fn delete_object(
&self,
req: S3Request<DeleteObjectInput>,
) -> S3Result<S3Response<DeleteObjectOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
self.state.deletes.fetch_add(1, Ordering::Relaxed);
let input = req.input;
append_object_wal_ordered(
&self.state,
WalEntryV1 {
seq: 0,
op: "delete".to_string(),
bucket: input.bucket.clone(),
key: input.key.clone(),
etag: None,
pointer_or_manifest: None,
ts_unix_ms: now_unix_ms(),
},
)?;
self.state
.objects
.write()
.map_err(|_| s3_error!(InternalError))
.map(|mut objects| {
if !remove_object_by_variants(&mut objects, &input.bucket, &input.key) {
eprintln!(
"nexusd delete_object miss bucket={} key={}",
input.bucket, input.key
);
}
})?;
Ok(S3Response::new(DeleteObjectOutput::default()))
}
async fn delete_objects(
&self,
req: S3Request<DeleteObjectsInput>,
) -> S3Result<S3Response<DeleteObjectsOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
let input = req.input;
ensure_bucket_exists(&self.state, &input.bucket)?;
let mut deleted = Vec::<DeletedObject>::new();
for object in input.delete.objects {
self.state.deletes.fetch_add(1, Ordering::Relaxed);
append_object_wal_ordered(
&self.state,
WalEntryV1 {
seq: 0,
op: "delete".to_string(),
bucket: input.bucket.clone(),
key: object.key.clone(),
etag: None,
pointer_or_manifest: None,
ts_unix_ms: now_unix_ms(),
},
)?;
self.state
.objects
.write()
.map_err(|_| s3_error!(InternalError))
.map(|mut objects| {
if !remove_object_by_variants(&mut objects, &input.bucket, &object.key) {
eprintln!(
"nexusd delete_objects miss bucket={} key={}",
input.bucket, object.key
);
}
})?;
deleted.push(DeletedObject {
key: Some(object.key),
..Default::default()
});
}
Ok(S3Response::new(DeleteObjectsOutput {
deleted: Some(deleted),
..Default::default()
}))
}
async fn create_multipart_upload(
&self,
req: S3Request<CreateMultipartUploadInput>,
) -> S3Result<S3Response<CreateMultipartUploadOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
self.state.mpu_creates.fetch_add(1, Ordering::Relaxed);
let input = req.input;
ensure_bucket_exists(&self.state, &input.bucket)?;
let upload_id = format!(
"upload-{:x}-{:x}",
now_unix_ms(),
next_nonce_seed(&self.state)
);
self.state
.uploads
.write()
.map_err(|_| s3_error!(InternalError))?
.insert(
upload_id.clone(),
MultipartUploadState {
bucket: input.bucket.clone(),
key: input.key.clone(),
parts: BTreeMap::new(),
},
);
append_multipart_wal_ordered(
&self.state,
MultipartWalEntryV1 {
seq: 0,
op: "create".to_string(),
upload_id: upload_id.clone(),
bucket: input.bucket.clone(),
key: input.key.clone(),
final_etag: None,
parts: Vec::new(),
ts_unix_ms: now_unix_ms(),
},
)?;
Ok(S3Response::new(CreateMultipartUploadOutput {
bucket: Some(input.bucket),
key: Some(input.key),
upload_id: Some(upload_id),
..Default::default()
}))
}
async fn upload_part(
&self,
req: S3Request<UploadPartInput>,
) -> S3Result<S3Response<UploadPartOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
self.state.mpu_parts.fetch_add(1, Ordering::Relaxed);
let input = req.input;
ensure_bucket_exists(&self.state, &input.bucket)?;
let upload_id = input.upload_id.clone();
let part_number = input.part_number as u32;
let Some(body) = input.body else {
return Err(s3_error!(IncompleteBody));
};
let raw = read_streaming_blob(body)
.await
.map_err(|_| s3_error!(InvalidRequest, "failed reading upload part body"))?;
validate_payload_headers(input.content_md5.as_deref(), input.content_length, &raw)?;
{
let uploads = self
.state
.uploads
.read()
.map_err(|_| s3_error!(InternalError))?;
if !uploads.contains_key(&upload_id) {
return Err(s3_error!(NoSuchUpload));
}
}
let nonce_seed = next_nonce_seed(&self.state).saturating_mul(1_000_000);
let compression_policy = self
.state
.compression_policies
.read()
.ok()
.and_then(|m| m.get(&input.bucket).cloned())
.unwrap_or_default();
let decision = compression::decide(None, &raw, &compression_policy);
if let Ok(mut guard) = self.state.last_compression_decision.write() {
*guard = Some(decision);
}
let qos_profile = self
.state
.qos_profiles
.read()
.ok()
.and_then(|m| m.get(&input.bucket).copied())
.unwrap_or(BucketQosProfile::Standard);
let active_mpu = self
.state
.uploads
.read()
.map(|g| g.len())
.unwrap_or(0);
let chunk_bytes = crate::ingest::compute_chunk_bytes(
self.state.chunk_runtime.current_chunk_bytes(),
qos_profile,
active_mpu,
);
let processed = process_payload_crc_zstd_aes(
&raw,
chunk_bytes,
&PIPELINE_KEY,
nonce_seed,
)
.map_err(|_| s3_error!(InternalError, "pipeline processing failed"))?;
let blob = encode_pipeline_blob(&processed.chunks)
.map_err(|_| s3_error!(InternalError, "pipeline blob encoding failed"))?;
let pointer = self
.state
.data_plane
.append_blob(&blob, nonce_seed)
.map_err(|_| s3_error!(InternalError, "device append failed"))?;
let etag = md5_hex(&raw);
self.state
.uploads
.write()
.map_err(|_| s3_error!(InternalError))?
.entry(upload_id.clone())
.and_modify(|upload| {
upload.parts.insert(
part_number,
MultipartPartState {
etag: etag.clone(),
pointer: pointer.clone(),
size: raw.len() as u64,
},
);
});
append_multipart_wal_ordered(
&self.state,
MultipartWalEntryV1 {
seq: 0,
op: "upload_part".to_string(),
upload_id: upload_id.clone(),
bucket: input.bucket,
key: input.key,
final_etag: None,
parts: vec![MultipartPart {
part_number,
etag: etag.clone(),
pointer,
}],
ts_unix_ms: now_unix_ms(),
},
)?;
Ok(S3Response::new(UploadPartOutput {
e_tag: Some(ETag::Strong(etag)),
..Default::default()
}))
}
async fn list_parts(
&self,
req: S3Request<ListPartsInput>,
) -> S3Result<S3Response<ListPartsOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
let input = req.input;
let upload_id = input.upload_id.clone();
let uploads = self
.state
.uploads
.read()
.map_err(|_| s3_error!(InternalError))?;
let upload = uploads
.get(&upload_id)
.ok_or_else(|| s3_error!(NoSuchUpload))?;
let mut part_rows = Vec::with_capacity(upload.parts.len());
for (number, part) in &upload.parts {
part_rows.push(Part {
part_number: Some(*number as i32),
e_tag: Some(ETag::Strong(part.etag.clone())),
size: Some(part.size as i64),
..Default::default()
});
}
Ok(S3Response::new(ListPartsOutput {
bucket: Some(upload.bucket.clone()),
key: Some(upload.key.clone()),
upload_id: Some(upload_id),
parts: Some(part_rows),
is_truncated: Some(false),
max_parts: Some(input.max_parts.unwrap_or(1000)),
..Default::default()
}))
}
async fn list_multipart_uploads(
&self,
req: S3Request<ListMultipartUploadsInput>,
) -> S3Result<S3Response<ListMultipartUploadsOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
self.state.lists.fetch_add(1, Ordering::Relaxed);
let input = req.input;
ensure_bucket_exists(&self.state, &input.bucket)?;
let max_uploads = normalize_max_uploads(input.max_uploads)?;
let key_marker = input.key_marker.clone();
let upload_id_marker = input.upload_id_marker.clone();
let prefix = input.prefix.clone();
let delimiter = input.delimiter.clone();
let encoding_type = input.encoding_type.clone();
let uploads = self
.state
.uploads
.read()
.map_err(|_| s3_error!(InternalError))?;
let mut rows = uploads
.iter()
.filter(|(_, upload)| upload.bucket == input.bucket)
.map(|(upload_id, upload)| (upload.key.clone(), upload_id.clone()))
.collect::<Vec<_>>();
rows.sort_by(|(key_a, upload_a), (key_b, upload_b)| {
key_a.cmp(key_b).then(upload_a.cmp(upload_b))
});
if let Some(prefix_filter) = prefix.as_ref() {
rows.retain(|(key, _)| key.starts_with(prefix_filter));
}
if let Some(key_marker_value) = key_marker.as_ref() {
rows.retain(|(key, upload_id)| {
if key > key_marker_value {
true
} else if key == key_marker_value {
upload_id_marker
.as_ref()
.map(|marker| upload_id > marker)
.unwrap_or(false)
} else {
false
}
});
}
let is_truncated = rows.len() > max_uploads;
let mut next_key_marker = None;
let mut next_upload_id_marker = None;
if is_truncated {
rows.truncate(max_uploads);
if let Some((last_key, last_upload_id)) = rows.last() {
next_key_marker = Some(last_key.clone());
next_upload_id_marker = Some(last_upload_id.clone());
}
}
let upload_rows = rows
.into_iter()
.map(|(key, upload_id)| MultipartUpload {
initiated: Some(Timestamp::from(SystemTime::now())),
key: Some(key),
upload_id: Some(upload_id),
initiator: Some(Initiator {
id: Some(DEFAULT_OWNER_ID.to_string()),
display_name: Some(DEFAULT_OWNER_DISPLAY_NAME.to_string()),
}),
owner: Some(default_owner()),
storage_class: Some(StorageClass::from_static(StorageClass::STANDARD)),
..Default::default()
})
.collect::<Vec<_>>();
Ok(S3Response::new(ListMultipartUploadsOutput {
bucket: Some(input.bucket),
delimiter,
encoding_type,
is_truncated: Some(is_truncated),
key_marker,
max_uploads: Some(max_uploads as i32),
next_key_marker,
next_upload_id_marker,
prefix,
upload_id_marker,
uploads: (!upload_rows.is_empty()).then_some(upload_rows),
..Default::default()
}))
}
async fn complete_multipart_upload(
&self,
req: S3Request<CompleteMultipartUploadInput>,
) -> S3Result<S3Response<CompleteMultipartUploadOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
self.state.mpu_completes.fetch_add(1, Ordering::Relaxed);
let input = req.input;
let upload_id = input.upload_id.clone();
let mut uploads = self
.state
.uploads
.write()
.map_err(|_| s3_error!(InternalError))?;
let upload = uploads
.remove(&upload_id)
.ok_or_else(|| s3_error!(NoSuchUpload))?;
let requested_parts = input
.multipart_upload
.as_ref()
.and_then(|mpu| mpu.parts.as_ref())
.cloned()
.unwrap_or_default();
let ordered_parts = if requested_parts.is_empty() {
upload
.parts
.iter()
.map(|(part_number, part)| {
Ok(MultipartPart {
part_number: *part_number,
etag: part.etag.clone(),
pointer: part.pointer.clone(),
})
})
.collect::<Result<Vec<_>, s3s::S3Error>>()?
} else {
let mut rows = Vec::with_capacity(requested_parts.len());
for requested in requested_parts {
let number = requested
.part_number
.ok_or_else(|| s3_error!(InvalidPart))? as u32;
let part = upload
.parts
.get(&number)
.ok_or_else(|| s3_error!(InvalidPart))?;
let expected = requested
.e_tag
.as_ref()
.map(|v| sanitize_etag(v.value()))
.unwrap_or_default();
if expected != sanitize_etag(&part.etag) {
return Err(s3_error!(InvalidPart));
}
rows.push(MultipartPart {
part_number: number,
etag: part.etag.clone(),
pointer: part.pointer.clone(),
});
}
rows
};
let etag_inputs = ordered_parts
.iter()
.map(|part| part.etag.clone())
.collect::<Vec<_>>();
let final_etag =
compute_multipart_etag(&etag_inputs).map_err(|_| s3_error!(InternalError))?;
append_multipart_wal_ordered(
&self.state,
MultipartWalEntryV1 {
seq: 0,
op: "complete".to_string(),
upload_id: upload_id.clone(),
bucket: upload.bucket.clone(),
key: upload.key.clone(),
final_etag: Some(final_etag.clone()),
parts: ordered_parts.clone(),
ts_unix_ms: now_unix_ms(),
},
)?;
let total_len = ordered_parts
.iter()
.map(|part| part.pointer.len as u64)
.sum::<u64>();
let object_key = storage_key(&upload.bucket, &upload.key);
let is_update = self
.state
.objects
.read()
.map_err(|_| s3_error!(InternalError))?
.contains_key(&object_key);
let object_meta = StoredObjectMeta {
etag: final_etag.clone(),
content_len: total_len,
last_modified: Timestamp::from(SystemTime::now()),
metadata: HashMap::new(),
layout: StoredLayout::Multipart {
parts: ordered_parts.clone(),
},
};
self.state
.objects
.write()
.map_err(|_| s3_error!(InternalError))?
.insert(object_key, object_meta);
enqueue_webhook_intents(
&self.state,
&upload.bucket,
&upload.key,
total_len,
final_etag.as_str(),
is_update,
)?;
Ok(S3Response::new(CompleteMultipartUploadOutput {
bucket: Some(upload.bucket),
key: Some(upload.key),
e_tag: Some(ETag::Strong(final_etag)),
..Default::default()
}))
}
async fn abort_multipart_upload(
&self,
req: S3Request<AbortMultipartUploadInput>,
) -> S3Result<S3Response<AbortMultipartUploadOutput>> {
self.state.request_count.fetch_add(1, Ordering::Relaxed);
self.state.mpu_aborts.fetch_add(1, Ordering::Relaxed);
let input = req.input;
let upload_id = input.upload_id.clone();
let removed = self
.state
.uploads
.write()
.map_err(|_| s3_error!(InternalError))?
.remove(&upload_id);
if removed.is_none() {
return Err(s3_error!(NoSuchUpload));
}
append_multipart_wal_ordered(
&self.state,
MultipartWalEntryV1 {
seq: 0,
op: "abort".to_string(),
upload_id,
bucket: input.bucket,
key: input.key,
final_etag: None,
parts: Vec::new(),
ts_unix_ms: now_unix_ms(),
},
)?;
Ok(S3Response::new(AbortMultipartUploadOutput::default()))
}
async fn select_object_content(
&self,
_req: S3Request<SelectObjectContentInput>,
) -> S3Result<S3Response<SelectObjectContentOutput>> {
Err(s3_error!(
NotImplemented,
"SelectObjectContent is not implemented in this release"
))
}
}
pub fn serve(config: ServeConfig) -> Result<()> {
validate_config(&config)?;
let listen_addr: SocketAddr = config
.listen_addr
.parse()
.with_context(|| format!("invalid listen_addr {}", config.listen_addr))?;
let topology = nexus_core::topology::inspect_topology()?;
let requested_profile = RuntimeProfile::from_str(config.runtime_profile.as_str())?;
let force_mode = AtcForceMode::from_str(config.atc_force_mode.as_str())?;
let mut device_paths = config.devices.clone();
if device_paths.is_empty() {
device_paths.push(PathBuf::from("/tmp/nexus-data-0.bin"));
}
let atc_decision = atc::decide(atc::AtcInput {
runtime_profile: requested_profile,
force_mode,
device_count: device_paths.len(),
single_device_threshold: config.atc_single_device_threshold.max(1),
topology: &topology,
});
let allow_sync_fallback = match atc_decision.mode {
atc::AtcMode::Posix => true,
atc::AtcMode::Hypervisor => config.allow_sync_fallback,
};
let require_io_uring = match atc_decision.mode {
atc::AtcMode::Posix => false,
atc::AtcMode::Hypervisor => config.require_io_uring,
};
if require_io_uring && !nexus_core::module_d::io_uring_available() {
anyhow::bail!("nexusd started with --require-io-uring but io_uring is unavailable");
}
let chunk_policy = ChunkPolicy::from_str(config.chunk_policy.as_str())?;
let chunk_runtime = Arc::new(ChunkRuntime::new(ChunkController::new(
topology.cpu_generation,
chunk_policy,
)));
let data_plane = Arc::new(DataPlane::open(
&device_paths,
atc_decision.io_model,
allow_sync_fallback,
BufferedSyncPolicy {
fsync_bytes: if atc_decision.runtime_profile == RuntimeProfile::Dev {
config.atc_dev_fsync_bytes.max(ALIGNMENT_U64)
} else {
config.atc_buffered_fsync_bytes.max(ALIGNMENT_U64)
},
fsync_every: if atc_decision.runtime_profile == RuntimeProfile::Dev {
Duration::from_millis(config.atc_dev_fsync_ms.max(1))
} else {
Duration::from_millis(config.atc_buffered_fsync_ms.max(1))
},
},
atc_decision.runtime_profile == RuntimeProfile::Dev && config.atc_dev_lazy_commit,
)?);
let wal_path = if atc_decision.runtime_profile == RuntimeProfile::Dev
&& config.wal_device.as_os_str().is_empty()
{
PathBuf::from("/tmp/nexus-wal.bin")
} else {
config.wal_device.clone()
};
let wal_writer = Arc::new(WalWriter::open(&wal_path)?);
let (
objects,
uploads,
iam_wal_log,
webhook_wal_log,
webhook_event_wal_log,
migration_wal_log,
) =
replay_initial_state(&data_plane, &wal_writer)?;
let access_key =
std::env::var("NEXUSD_ACCESS_KEY").unwrap_or_else(|_| DEFAULT_ACCESS_KEY.to_string());
let secret_key =
std::env::var("NEXUSD_SECRET_KEY").unwrap_or_else(|_| DEFAULT_SECRET_KEY.to_string());
let iam_access_key =
std::env::var("NEXUSD_IAM_ACCESS_KEY").unwrap_or_else(|_| "AKIAIAMUSER".to_string());
let iam_secret_key =
std::env::var("NEXUSD_IAM_SECRET_KEY").unwrap_or_else(|_| "SECRETIAMUSER".to_string());
let iam_user_name = std::env::var("NEXUSD_IAM_USER").unwrap_or_else(|_| "iam-user".to_string());
let credentials = Arc::new(CredentialStore::new(access_key.clone(), secret_key.clone()));
credentials.upsert(crate::auth_provider::CredentialRecord {
access_key: iam_access_key.clone(),
secret_key: iam_secret_key,
principal: format!("arn:aws:iam:::user/{iam_user_name}"),
session_token: None,
expires_at_unix_ms: None,
kind: CredentialKind::User,
});
let iam = Arc::new(IamState::new());
for user_name in [iam_user_name.as_str(), "nexus"] {
iam.create_user(user_name)?;
iam.put_user_policy(
user_name,
"default-allow",
r#"{"Version":"2012-10-17","Statement":{"Effect":"Allow","Action":"s3:*","Resource":"*"} }"#,
)?;
}
let sts = Arc::new(StsService::new(credentials.clone(), iam.clone()));
let (webhook_completion_tx, webhook_completion_rx) = bounded::<WebhookDeliveryResult>(4096);
let webhooks = Arc::new(WebhookEngine::new(WebhookConfig {
workers: config.webhook_workers,
retry_max: config.webhook_retry_max,
timeout_ms: config.webhook_timeout_ms,
queue_capacity: config.webhook_queue_capacity,
completion_tx: Some(webhook_completion_tx),
}));
let migration = Arc::new(MigrationState::default());
if let Some(endpoint) = config.bifrost_source_endpoint.clone() {
migration.set_source(Some(BifrostSourceConfig {
endpoint,
region: config.bifrost_source_region.clone(),
access_key: config.bifrost_source_access_key.clone(),
secret_key: config.bifrost_source_secret_key.clone(),
}));
}
replay_control_state(
&iam,
&credentials,
&webhooks,
&migration,
iam_wal_log.as_slice(),
webhook_wal_log.as_slice(),
webhook_event_wal_log.as_slice(),
migration_wal_log.as_slice(),
)?;
let buckets = Arc::new(RwLock::new(BTreeSet::from(["nexus".to_string()])));
let state = GatewayState {
data_plane: data_plane.clone(),
wal: wal_writer.clone(),
wal_sequence_lock: Arc::new(Mutex::new(())),
credentials: credentials.clone(),
iam: iam.clone(),
sts: sts.clone(),
webhooks: webhooks.clone(),
migration: migration.clone(),
nonce_counter: Arc::new(AtomicU64::new(1)),
objects: Arc::new(RwLock::new(objects)),
uploads: Arc::new(RwLock::new(uploads)),
buckets,
compression_policies: Arc::new(RwLock::new(BTreeMap::new())),
qos_profiles: Arc::new(RwLock::new(BTreeMap::new())),
atc_decision: atc_decision.clone(),
chunk_runtime,
request_count: Arc::new(AtomicU64::new(0)),
puts: Arc::new(AtomicU64::new(0)),
gets: Arc::new(AtomicU64::new(0)),
lists: Arc::new(AtomicU64::new(0)),
deletes: Arc::new(AtomicU64::new(0)),
mpu_creates: Arc::new(AtomicU64::new(0)),
mpu_parts: Arc::new(AtomicU64::new(0)),
mpu_completes: Arc::new(AtomicU64::new(0)),
mpu_aborts: Arc::new(AtomicU64::new(0)),
last_compression_decision: Arc::new(RwLock::new(None)),
buffered_sync_flushes: data_plane.buffered_flushes.clone(),
buffered_sync_flush_latency_ms: data_plane.buffered_flush_latency_ms.clone(),
};
let completion_state = state.clone();
let _webhook_completion = thread::Builder::new()
.name("nexusd-webhook-completion".to_string())
.spawn(move || {
while let Ok(event) = webhook_completion_rx.recv() {
if event.success {
let _ = append_webhook_event_wal_ordered(
&completion_state,
"complete",
event.event_id.as_str(),
event.rule_id.as_str(),
"",
"",
"{}",
);
}
}
})
.context("failed spawning webhook completion thread")?;
let worker_count = match atc_decision.worker_model {
WorkerModel::PinnedLocalset => select_worker_cores(&topology, config.workers).len(),
WorkerModel::TokioWorkStealing => select_posix_worker_count(config.workers),
};
eprintln!(
"nexusd start: addr={} workers={} runtime_profile={} topology_policy={} chunk_policy={} atc={}",
listen_addr,
worker_count,
atc_decision.runtime_profile.as_str(),
config.topology_policy,
config.chunk_policy,
atc_decision.to_json()
);
eprintln!("nexusd cpu/topology: {}", topology.to_json());
eprintln!(
"nexusd devices: mode={} [{}]",
data_plane.mode_summary,
data_plane.device_report().join("; ")
);
eprintln!("nexusd wal: {}", wal_writer.wal_path().display());
let accept_counts = match atc_decision.worker_model {
WorkerModel::PinnedLocalset => Arc::new(
(0..worker_count)
.map(|_| AtomicU64::new(0))
.collect::<Vec<_>>(),
),
WorkerModel::TokioWorkStealing => Arc::new(vec![AtomicU64::new(0)]),
};
let _stats_handle =
spawn_stats_thread(state.clone(), accept_counts.clone(), atc_decision.clone())?;
match atc_decision.worker_model {
WorkerModel::PinnedLocalset => run_hypervisor_workers(
listen_addr,
state,
access_key,
select_worker_cores(&topology, config.workers),
accept_counts,
),
WorkerModel::TokioWorkStealing => run_posix_server(
listen_addr,
state,
access_key,
worker_count,
accept_counts,
),
}
}
fn spawn_stats_thread(
stats_state: GatewayState,
stats_accept_counts: Arc<Vec<AtomicU64>>,
atc_decision: AtcDecision,
) -> Result<thread::JoinHandle<()>> {
thread::Builder::new()
.name("nexusd-stats".to_string())
.spawn(move || loop {
thread::sleep(Duration::from_secs(10));
let mut accepts = Vec::with_capacity(stats_accept_counts.len());
for (idx, counter) in stats_accept_counts.iter().enumerate() {
accepts.push(format!("w{idx}={}", counter.load(Ordering::Relaxed)));
}
let webhook_stats = stats_state.webhooks.stats();
let last_compression = stats_state
.last_compression_decision
.read()
.ok()
.and_then(|v| v.clone())
.map(|v| v.to_json())
.unwrap_or_else(|| "{}".to_string());
eprintln!(
"nexusd stats requests={} puts={} gets={} lists={} deletes={} mpu_create={} mpu_parts={} mpu_complete={} mpu_abort={} chunk={} downshift={} upshift={} qdepth_median={:.2} webhook_enqueued={} webhook_dropped={} webhook_qcap={} buffered_flushes={} buffered_flush_latency_ms_sum={} atc_mode={} atc_io_model={} atc_worker_model={} atc_device_count={} atc_decision_reason=\"{}\" last_compression='{}' accepts=[{}]",
stats_state.request_count.load(Ordering::Relaxed),
stats_state.puts.load(Ordering::Relaxed),
stats_state.gets.load(Ordering::Relaxed),
stats_state.lists.load(Ordering::Relaxed),
stats_state.deletes.load(Ordering::Relaxed),
stats_state.mpu_creates.load(Ordering::Relaxed),
stats_state.mpu_parts.load(Ordering::Relaxed),
stats_state.mpu_completes.load(Ordering::Relaxed),
stats_state.mpu_aborts.load(Ordering::Relaxed),
stats_state.chunk_runtime.current_chunk_bytes(),
stats_state.chunk_runtime.downshift_events.load(Ordering::Relaxed),
stats_state.chunk_runtime.upshift_events.load(Ordering::Relaxed),
stats_state.data_plane.median_queue_depth(),
webhook_stats.enqueued_total,
webhook_stats.dropped_total,
webhook_stats.queue_capacity,
stats_state.buffered_sync_flushes.load(Ordering::Relaxed),
stats_state.buffered_sync_flush_latency_ms.load(Ordering::Relaxed),
atc_decision.mode.as_str(),
atc_decision.io_model.as_str(),
atc_decision.worker_model.as_str(),
atc_decision.device_count,
atc_decision.reason,
last_compression,
accepts.join(",")
);
})
.context("failed spawning stats thread")
}
fn run_hypervisor_workers(
listen_addr: SocketAddr,
state: GatewayState,
access_key: String,
worker_core_ids: Vec<usize>,
accept_counts: Arc<Vec<AtomicU64>>,
) -> Result<()> {
if worker_core_ids.is_empty() {
anyhow::bail!("no worker cores selected by topology policy");
}
let core_affinity_map = core_affinity::get_core_ids()
.unwrap_or_default()
.into_iter()
.map(|core_id| (core_id.id, core_id))
.collect::<HashMap<_, _>>();
let mut worker_handles = Vec::with_capacity(worker_core_ids.len());
for worker_idx in 0..worker_core_ids.len() {
let state = state.clone();
let accept_counts = accept_counts.clone();
let access_key = access_key.clone();
let requested_core = worker_core_ids[worker_idx];
let worker_core = core_affinity_map.get(&requested_core).copied();
let handle = thread::Builder::new()
.name(format!("nexusd-worker-{worker_idx}"))
.spawn(move || -> Result<()> {
if let Some(core_id) = worker_core {
let _ = core_affinity::set_for_current(core_id);
}
let std_listener = bind_reuseport_listener(listen_addr)
.with_context(|| format!("worker {worker_idx} failed to bind listener"))?;
let rt = tokio::runtime::Builder::new_current_thread()
.enable_io()
.enable_time()
.build()
.context("failed creating current-thread runtime")?;
rt.block_on(async move {
let local = tokio::task::LocalSet::new();
local
.run_until(async move {
let listener = tokio::net::TcpListener::from_std(std_listener)
.context("failed to adapt std listener")?;
let gateway = GatewayS3 {
state: state.clone(),
};
let s3_service = {
let mut builder = S3ServiceBuilder::new(gateway);
builder.set_auth(DynamicS3Auth {
store: state.credentials.clone(),
});
builder.set_access(GatewayAccessControl {
credentials: state.credentials.clone(),
iam: state.iam.clone(),
});
builder.build()
};
loop {
let (stream, _remote) = listener
.accept()
.await
.context("accept failed in nexusd worker")?;
accept_counts[worker_idx].fetch_add(1, Ordering::Relaxed);
let service = s3_service.clone();
let intercept_state = state.clone();
let expected_access_key = access_key.clone();
tokio::task::spawn_local(async move {
let io = TokioIo::new(stream);
let svc = service_fn(move |req: http::Request<HyperIncoming>| {
let service = service.clone();
let expected_access_key = expected_access_key.clone();
let intercept_state = intercept_state.clone();
async move {
let request_fut = async {
match maybe_intercept_request(
req,
intercept_state,
expected_access_key.as_str(),
)
.await
{
Ok(InterceptResult::Response(resp)) => Ok(resp),
Ok(InterceptResult::Forward(req)) => {
service.call(req.map(Body::from)).await
}
Err(err) => Ok(xml_error_response(
StatusCode::INTERNAL_SERVER_ERROR,
"InternalError",
&err.to_string(),
)),
}
};
match tokio::time::timeout(
Duration::from_millis(REQUEST_TIMEOUT_MS),
request_fut,
)
.await
{
Ok(response) => response,
Err(_) => Ok(xml_error_response_close(
StatusCode::REQUEST_TIMEOUT,
"RequestTimeout",
"Request processing exceeded timeout window",
)),
}
}
});
let conn = ConnBuilder::new(TokioExecutor::new())
.serve_connection(io, svc)
.into_owned();
if let Err(err) = conn.await {
eprintln!("nexusd connection error: {err}");
}
});
}
})
.await
})
})
.with_context(|| format!("failed spawning worker {worker_idx}"))?;
worker_handles.push(handle);
}
for handle in worker_handles {
let worker_result = handle
.join()
.map_err(|_| anyhow::anyhow!("worker thread panicked"))?;
worker_result?;
}
Ok(())
}
fn run_posix_server(
listen_addr: SocketAddr,
state: GatewayState,
access_key: String,
worker_count: usize,
accept_counts: Arc<Vec<AtomicU64>>,
) -> Result<()> {
if worker_count == 0 {
anyhow::bail!("no workers selected for tokio work-stealing runtime");
}
let rt = tokio::runtime::Builder::new_multi_thread()
.worker_threads(worker_count)
.enable_io()
.enable_time()
.build()
.context("failed creating multi-thread runtime")?;
rt.block_on(async move {
let listener = tokio::net::TcpListener::bind(listen_addr)
.await
.with_context(|| format!("failed binding {listen_addr}"))?;
let gateway = GatewayS3 {
state: state.clone(),
};
let s3_service = {
let mut builder = S3ServiceBuilder::new(gateway);
builder.set_auth(DynamicS3Auth {
store: state.credentials.clone(),
});
builder.set_access(GatewayAccessControl {
credentials: state.credentials.clone(),
iam: state.iam.clone(),
});
builder.build()
};
loop {
let (stream, _remote) = listener
.accept()
.await
.context("accept failed in nexusd tokio runtime")?;
accept_counts[0].fetch_add(1, Ordering::Relaxed);
let service = s3_service.clone();
let intercept_state = state.clone();
let expected_access_key = access_key.clone();
tokio::spawn(async move {
let io = TokioIo::new(stream);
let svc = service_fn(move |req: http::Request<HyperIncoming>| {
let service = service.clone();
let expected_access_key = expected_access_key.clone();
let intercept_state = intercept_state.clone();
async move {
let request_fut = async {
match maybe_intercept_request(
req,
intercept_state,
expected_access_key.as_str(),
)
.await
{
Ok(InterceptResult::Response(resp)) => Ok(resp),
Ok(InterceptResult::Forward(req)) => {
service.call(req.map(Body::from)).await
}
Err(err) => Ok(xml_error_response(
StatusCode::INTERNAL_SERVER_ERROR,
"InternalError",
&err.to_string(),
)),
}
};
match tokio::time::timeout(
Duration::from_millis(REQUEST_TIMEOUT_MS),
request_fut,
)
.await
{
Ok(response) => response,
Err(_) => Ok(xml_error_response_close(
StatusCode::REQUEST_TIMEOUT,
"RequestTimeout",
"Request processing exceeded timeout window",
)),
}
}
});
let conn = ConnBuilder::new(TokioExecutor::new())
.serve_connection(io, svc)
.into_owned();
if let Err(err) = conn.await {
eprintln!("nexusd connection error: {err}");
}
});
}
})
}
pub fn validate(config: ValidationConfig) -> Result<()> {
let endpoint = config
.s3_endpoint
.trim_start_matches("http://")
.trim_start_matches("https://")
.to_string();
let reachable = std::net::TcpStream::connect(endpoint.as_str()).is_ok();
if config.json {
println!(
"{{\"module\":\"NEXUSD_VALIDATE\",\"endpoint\":\"{}\",\"reachable\":{},\"access_key_len\":{},\"secret_key_len\":{}}}",
config.s3_endpoint,
reachable,
config.access_key.len(),
config.secret_key.len(),
);
} else {
println!(
"validate endpoint={} reachable={} access_key_len={} secret_key_len={}",
config.s3_endpoint,
reachable,
config.access_key.len(),
config.secret_key.len()
);
}
Ok(())
}
pub fn admin_purge_data(config: AdminPurgeConfig) -> Result<()> {
let report = admin::purge_data(PurgeDataConfig {
danger: config.danger,
bucket_prefix: config.bucket_prefix,
wal_path: config.wal_path,
devices: config.devices,
})?;
if config.json {
println!("{}", report.to_json());
} else {
println!(
"purge danger=true wal_punched={} devices={} bucket_prefix={}",
report.wal_frames_punched,
report.devices.len(),
report.bucket_prefix.unwrap_or_else(|| "*".to_string())
);
}
Ok(())
}
pub fn migration_import_minio(config: MigrationImportConfig) -> Result<()> {
let report = migration::import_minio_export(config.config_path.as_path(), config.iam_dir.as_path())?;
if config.json {
println!("{}", report.to_json());
} else {
println!(
"import source={} users={} keys={} policies={} skipped={}",
report.source_config,
report.users_imported,
report.access_keys_imported,
report.policies_imported,
report.skipped_items
);
}
Ok(())
}
pub fn migration_sync(config: MigrationSyncConfig) -> Result<()> {
let status = serde_json::json!({
"bucket": config.bucket,
"prefix": config.prefix,
"workers": config.workers.max(1),
"status": "accepted"
});
if config.json {
println!("{}", status);
} else {
println!(
"migration sync accepted bucket={} prefix={} workers={}",
config.bucket,
config.prefix,
config.workers.max(1)
);
}
Ok(())
}
pub fn migration_status(config: MigrationStatusConfig) -> Result<()> {
let status = serde_json::json!({
"status": "idle",
"jobs": []
});
if config.json {
println!("{}", status);
} else {
println!("migration status=idle jobs=0");
}
Ok(())
}
fn validate_config(config: &ServeConfig) -> Result<()> {
let runtime_profile = RuntimeProfile::from_str(config.runtime_profile.as_str())?.resolve();
AtcForceMode::from_str(config.atc_force_mode.as_str())?;
if runtime_profile != RuntimeProfile::Dev && config.wal_device.as_os_str().is_empty() {
anyhow::bail!("wal_device must be provided");
}
if config.devices.len() > 5 {
anyhow::bail!("at most 5 data devices are supported");
}
if config.atc_single_device_threshold == 0 {
anyhow::bail!("atc_single_device_threshold must be >= 1");
}
if config.atc_buffered_fsync_bytes == 0 {
anyhow::bail!("atc_buffered_fsync_bytes must be > 0");
}
if config.atc_buffered_fsync_ms == 0 {
anyhow::bail!("atc_buffered_fsync_ms must be > 0");
}
if config.atc_dev_fsync_bytes == 0 {
anyhow::bail!("atc_dev_fsync_bytes must be > 0");
}
if config.atc_dev_fsync_ms == 0 {
anyhow::bail!("atc_dev_fsync_ms must be > 0");
}
Ok(())
}
fn replay_initial_state(
data_plane: &Arc<DataPlane>,
wal_writer: &Arc<WalWriter>,
) -> Result<(
BTreeMap<String, StoredObjectMeta>,
BTreeMap<String, MultipartUploadState>,
Vec<wal::IamWalEntryV1>,
Vec<wal::WebhookWalEntryV1>,
Vec<wal::WebhookEventWalEntryV1>,
Vec<wal::MigrationWalEntryV1>,
)> {
let replay = wal::replay_wal_state(WalReplayConfig {
wal_path: wal_writer.wal_path().to_path_buf(),
dry_run: false,
})?;
let iam_log = replay.iam_log.clone();
let webhook_log = replay.webhook_log.clone();
let webhook_event_log = replay.webhook_event_log.clone();
let migration_log = replay.migration_log.clone();
let mut objects = BTreeMap::<String, StoredObjectMeta>::new();
for (_storage_key, entry) in replay.object_index {
let Some(pointer) = entry.pointer_or_manifest else {
continue;
};
let content_len = match data_plane
.read_blob(&pointer)
.and_then(|blob| decode_pipeline_blob(&blob).map(|bytes| bytes.len() as u64))
{
Ok(len) => len,
Err(_) => 0,
};
objects.insert(
storage_key(&entry.bucket, &entry.key),
StoredObjectMeta {
etag: entry.etag.unwrap_or_else(|| format!("seq-{}", entry.seq)),
content_len,
last_modified: Timestamp::from(unix_ms_to_system_time(entry.ts_unix_ms)),
metadata: HashMap::new(),
layout: StoredLayout::InlineBlob { pointer },
},
);
}
let mut uploads = BTreeMap::<String, MultipartUploadState>::new();
for (upload_id, entry) in replay.multipart_uploads {
match entry.op.as_str() {
"complete" => {
let content_len = entry
.parts
.iter()
.map(|part| part.pointer.len as u64)
.sum::<u64>();
objects.insert(
storage_key(&entry.bucket, &entry.key),
StoredObjectMeta {
etag: entry
.final_etag
.unwrap_or_else(|| format!("mpu-{}", entry.seq)),
content_len,
last_modified: Timestamp::from(unix_ms_to_system_time(entry.ts_unix_ms)),
metadata: HashMap::new(),
layout: StoredLayout::Multipart { parts: entry.parts },
},
);
}
"create" | "upload_part" => {
let mut state = MultipartUploadState {
bucket: entry.bucket.clone(),
key: entry.key.clone(),
parts: BTreeMap::new(),
};
for part in entry.parts {
state.parts.insert(
part.part_number,
MultipartPartState {
etag: part.etag,
pointer: part.pointer,
size: 0,
},
);
}
uploads.insert(upload_id, state);
}
"abort" => {}
_ => {}
}
}
Ok((
objects,
uploads,
iam_log,
webhook_log,
webhook_event_log,
migration_log,
))
}
fn ensure_bucket_exists(state: &GatewayState, bucket: &str) -> S3Result<()> {
let buckets = state.buckets.read().map_err(|_| s3_error!(InternalError))?;
if !buckets.contains(bucket) {
return Err(s3_error!(NoSuchBucket));
}
Ok(())
}
fn validate_bucket_name(bucket: &str) -> S3Result<()> {
let len = bucket.len();
if !(3..=63).contains(&len) {
return Err(s3_error!(InvalidBucketName));
}
if bucket.starts_with("xn--") || bucket.ends_with("-s3alias") {
return Err(s3_error!(InvalidBucketName));
}
let bytes = bucket.as_bytes();
let starts_alnum = bytes
.first()
.map(|b| b.is_ascii_alphanumeric())
.unwrap_or(false);
let ends_alnum = bytes
.last()
.map(|b| b.is_ascii_alphanumeric())
.unwrap_or(false);
if !starts_alnum || !ends_alnum {
return Err(s3_error!(InvalidBucketName));
}
if bucket.contains("..") || bucket.contains(".-") || bucket.contains("-.") {
return Err(s3_error!(InvalidBucketName));
}
if !bytes
.iter()
.all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || *b == b'.' || *b == b'-')
{
return Err(s3_error!(InvalidBucketName));
}
if looks_like_ipv4(bucket) {
return Err(s3_error!(InvalidBucketName));
}
Ok(())
}
fn looks_like_ipv4(value: &str) -> bool {
let parts = value.split('.').collect::<Vec<_>>();
if parts.len() != 4 {
return false;
}
parts.iter().all(|part| {
!part.is_empty()
&& part.len() <= 3
&& part.as_bytes().iter().all(|b| b.is_ascii_digit())
&& part.parse::<u8>().is_ok()
})
}
async fn read_streaming_blob(mut body: StreamingBlob) -> Result<Vec<u8>> {
let mut out = Vec::new();
while let Some(chunk) = body
.try_next()
.await
.map_err(|err| anyhow::anyhow!("failed reading streaming blob: {err}"))?
{
out.extend_from_slice(chunk.as_ref());
}
Ok(out)
}
fn encode_pipeline_blob(chunks: &[EncryptedChunk]) -> Result<Vec<u8>> {
let mut out = Vec::new();
out.extend_from_slice(&RECORD_VERSION.to_le_bytes());
let count: u32 = chunks
.len()
.try_into()
.context("too many chunks to encode")?;
out.extend_from_slice(&count.to_le_bytes());
for chunk in chunks {
let plain_len: u32 = chunk
.plain_len
.try_into()
.context("plain chunk length exceeds u32")?;
let cipher_len: u32 = chunk
.bytes
.len()
.try_into()
.context("cipher chunk length exceeds u32")?;
out.extend_from_slice(&plain_len.to_le_bytes());
out.extend_from_slice(&chunk.crc32c.to_le_bytes());
out.extend_from_slice(&chunk.nonce_id.to_le_bytes());
out.extend_from_slice(&cipher_len.to_le_bytes());
out.extend_from_slice(chunk.bytes.as_slice());
}
Ok(out)
}
fn decode_pipeline_blob(blob: &[u8]) -> Result<Vec<u8>> {
if blob.len() < 8 {
anyhow::bail!("pipeline blob too small");
}
let mut cursor = 0usize;
let version = u32::from_le_bytes(blob[cursor..cursor + 4].try_into().expect("version"));
cursor += 4;
if version != RECORD_VERSION {
anyhow::bail!("unsupported pipeline blob version: {}", version);
}
let chunk_count = u32::from_le_bytes(blob[cursor..cursor + 4].try_into().expect("count"));
cursor += 4;
let mut restored = Vec::new();
for _ in 0..chunk_count {
if cursor + 20 > blob.len() {
anyhow::bail!("pipeline blob truncated in chunk header");
}
let plain_len = u32::from_le_bytes(blob[cursor..cursor + 4].try_into().expect("plain"));
cursor += 4;
let crc32c = u32::from_le_bytes(blob[cursor..cursor + 4].try_into().expect("crc"));
cursor += 4;
let nonce_id = u64::from_le_bytes(blob[cursor..cursor + 8].try_into().expect("nonce"));
cursor += 8;
let cipher_len = u32::from_le_bytes(blob[cursor..cursor + 4].try_into().expect("cipher"));
cursor += 4;
let cipher_len_usize = cipher_len as usize;
if cursor + cipher_len_usize > blob.len() {
anyhow::bail!("pipeline blob truncated in chunk payload");
}
let cipher = &blob[cursor..cursor + cipher_len_usize];
cursor += cipher_len_usize;
let plain =
restore_chunk_zstd_aes(cipher, crc32c, plain_len as usize, &PIPELINE_KEY, nonce_id)?;
restored.extend_from_slice(&plain);
}
Ok(restored)
}
fn materialize_object(state: &GatewayState, meta: &StoredObjectMeta) -> Result<Vec<u8>> {
match &meta.layout {
StoredLayout::InlineBlob { pointer } => {
let blob = state.data_plane.read_blob(pointer)?;
decode_pipeline_blob(&blob)
}
StoredLayout::Multipart { parts } => {
let mut out = Vec::new();
for part in parts {
let blob = state.data_plane.read_blob(&part.pointer)?;
let plain = decode_pipeline_blob(&blob)?;
out.extend_from_slice(&plain);
}
Ok(out)
}
}
}
fn select_worker_cores(topology: &TopologyReport, requested_workers: usize) -> Vec<usize> {
let mut selected = if topology.selected_cores.is_empty() {
let n = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(1)
.min(64);
(0..n).collect::<Vec<_>>()
} else {
topology.selected_cores.clone()
};
if selected.len() > 64 {
selected.truncate(64);
}
let workers = if requested_workers == 0 {
selected.len().max(1)
} else {
requested_workers.min(64).min(selected.len().max(1))
};
if selected.len() >= workers {
selected.truncate(workers);
return selected;
}
if selected.is_empty() {
return vec![0];
}
let mut expanded = Vec::with_capacity(workers);
for idx in 0..workers {
expanded.push(selected[idx % selected.len()]);
}
expanded
}
fn select_posix_worker_count(requested_workers: usize) -> usize {
let available = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(1)
.clamp(1, 64);
if requested_workers == 0 {
available
} else {
requested_workers.clamp(1, 64)
}
}
fn bind_reuseport_listener(addr: SocketAddr) -> Result<std::net::TcpListener> {
let domain = if addr.is_ipv4() {
Domain::IPV4
} else {
Domain::IPV6
};
let socket = Socket::new(domain, Type::STREAM, Some(Protocol::TCP))
.context("failed creating TCP socket")?;
socket
.set_reuse_address(true)
.context("failed setting SO_REUSEADDR")?;
#[cfg(any(target_os = "linux", target_os = "android"))]
socket
.set_reuse_port(true)
.context("failed setting SO_REUSEPORT")?;
socket
.bind(&addr.into())
.with_context(|| format!("failed binding {addr}"))?;
socket.listen(2048).context("failed listen")?;
socket
.set_nonblocking(true)
.context("failed setting non-blocking")?;
Ok(socket.into())
}
fn open_device(
path: &Path,
io_model: IoModel,
allow_sync_fallback: bool,
) -> Result<(File, IoModel, PathBuf)> {
let direct_path = path.to_path_buf();
match open_rw_path(&direct_path, false, io_model) {
Ok(file) => Ok((file, io_model, direct_path)),
Err(primary_err) => {
if !allow_sync_fallback {
return Err(primary_err).with_context(|| {
format!(
"failed opening data device {} and fallback disabled",
direct_path.display()
)
});
}
let fallback = if path.starts_with("/dev") {
let file_name = path
.file_name()
.map(|name| name.to_string_lossy().to_string())
.unwrap_or_else(|| "device".to_string());
PathBuf::from(format!("/tmp/nexus-device-{file_name}.bin"))
} else {
path.to_path_buf()
};
if let Some(parent) = fallback.parent() {
if !parent.as_os_str().is_empty() {
fs::create_dir_all(parent).with_context(|| {
format!("failed creating fallback parent {}", parent.display())
})?;
}
}
let fallback_io_model = match io_model {
IoModel::DirectSync => IoModel::Buffered,
IoModel::Buffered => IoModel::Buffered,
};
let file = open_rw_path(&fallback, true, fallback_io_model).with_context(|| {
format!("failed opening fallback data path {}", fallback.display())
})?;
Ok((file, fallback_io_model, fallback))
}
}
}
fn open_rw_path(path: &Path, create: bool, _io_model: IoModel) -> Result<File> {
let mut opts = OpenOptions::new();
opts.read(true).write(true);
if create {
opts.create(true);
}
#[cfg(any(target_os = "linux", target_os = "android"))]
if _io_model == IoModel::DirectSync {
opts.custom_flags(libc::O_DIRECT | libc::O_DSYNC);
}
opts.open(path)
.with_context(|| format!("open failed for {}", path.display()))
}
fn write_all_at(file: &mut File, mut buf: &[u8], mut offset: u64) -> io::Result<()> {
while !buf.is_empty() {
let written = file.write_at(buf, offset)?;
if written == 0 {
return Err(io::Error::new(
ErrorKind::WriteZero,
"write_at returned zero bytes",
));
}
buf = &buf[written..];
offset = offset.saturating_add(written as u64);
}
Ok(())
}
fn read_exact_at(file: &mut File, mut offset: u64, len: usize) -> Result<Vec<u8>> {
let mut out = vec![0_u8; len];
let mut filled = 0usize;
while filled < len {
let n = file
.read_at(&mut out[filled..], offset)
.context("read_at failed")?;
if n == 0 {
anyhow::bail!(
"unexpected EOF while reading at offset {} (wanted {}, got {})",
offset,
len,
filled
);
}
filled += n;
offset = offset.saturating_add(n as u64);
}
Ok(out)
}
fn align_up(value: usize, alignment: usize) -> usize {
if value % alignment == 0 {
value
} else {
value + (alignment - (value % alignment))
}
}
fn percentile_sorted(values: &[f64], percentile: f64) -> f64 {
if values.is_empty() {
return 0.0;
}
let clamped = percentile.clamp(0.0, 100.0);
let rank = ((clamped / 100.0) * ((values.len() - 1) as f64)).round() as usize;
values[rank]
}
fn storage_key(bucket: &str, key: &str) -> String {
format!("{bucket}/{key}")
}
fn storage_key_prefix(bucket: &str, prefix: &str) -> String {
format!("{bucket}/{prefix}")
}
fn strip_storage_prefix(bucket: &str, full_key: &str) -> Option<String> {
let prefix = format!("{bucket}/");
full_key.strip_prefix(&prefix).map(ToString::to_string)
}
fn collect_console_bucket_summaries(state: &GatewayState) -> Vec<BrowserBucketSummary> {
let bucket_names = state
.buckets
.read()
.map(|guard| guard.iter().cloned().collect::<Vec<_>>())
.unwrap_or_default();
let objects = state
.objects
.read()
.map(|guard| guard.clone())
.unwrap_or_default();
let compression = state
.compression_policies
.read()
.map(|guard| guard.clone())
.unwrap_or_default();
let qos = state
.qos_profiles
.read()
.map(|guard| guard.clone())
.unwrap_or_default();
let mut buckets = bucket_names
.into_iter()
.map(|name| {
let prefix = format!("{name}/");
let (object_count, bytes) = objects
.iter()
.filter(|(key, _)| key.starts_with(prefix.as_str()))
.fold((0_u64, 0_u64), |(count, bytes), (_, meta)| {
(count.saturating_add(1), bytes.saturating_add(meta.content_len))
});
BrowserBucketSummary {
compression_policy: compression.contains_key(name.as_str()),
qos_profile: qos
.get(name.as_str())
.map(|profile| match profile {
BucketQosProfile::Standard => "standard".to_string(),
BucketQosProfile::Ingest => "ingest".to_string(),
})
.unwrap_or_else(|| "standard".to_string()),
webhook_rule_count: state.webhooks.get_rules(name.as_str()).len(),
name,
object_count,
bytes,
}
})
.collect::<Vec<_>>();
buckets.sort_by(|left, right| left.name.cmp(&right.name));
buckets
}
fn normalize_browser_prefix(prefix: &str) -> String {
let trimmed = prefix.trim_start_matches('/');
if trimmed.is_empty() {
String::new()
} else if trimmed.ends_with('/') {
trimmed.to_string()
} else {
format!("{trimmed}/")
}
}
fn browser_content_type_from_metadata(
metadata: &HashMap<String, String>,
key: &str,
) -> Option<String> {
metadata
.get("content-type")
.cloned()
.or_else(|| guess_browser_content_type_from_key(key).map(str::to_string))
}
fn guess_browser_content_type_from_key(key: &str) -> Option<&'static str> {
let lower = key.to_ascii_lowercase();
let ext = lower.rsplit('.').next().unwrap_or_default();
match ext {
"txt" | "md" | "log" | "csv" => Some("text/plain"),
"json" => Some("application/json"),
"xml" => Some("application/xml"),
"html" | "htm" => Some("text/html"),
"pdf" => Some("application/pdf"),
"jpg" | "jpeg" => Some("image/jpeg"),
"png" => Some("image/png"),
"gif" => Some("image/gif"),
"webp" => Some("image/webp"),
"svg" => Some("image/svg+xml"),
"mp4" | "mov" => Some("video/mp4"),
"mp3" | "wav" => Some("audio/mpeg"),
"zip" => Some("application/zip"),
"docx" => Some(
"application/vnd.openxmlformats-officedocument.wordprocessingml.document",
),
_ => None,
}
}
fn browser_preview_kind(key: &str, content_type: Option<&str>) -> &'static str {
let value = content_type
.or_else(|| guess_browser_content_type_from_key(key))
.unwrap_or("application/octet-stream")
.to_ascii_lowercase();
if value.starts_with("image/") {
"image"
} else if value == "application/pdf" {
"pdf"
} else if value.starts_with("text/")
|| value.contains("json")
|| value.contains("xml")
|| key.ends_with(".md")
|| key.ends_with(".csv")
{
"text"
} else if value.starts_with("video/") {
"video"
} else if value.starts_with("audio/") {
"audio"
} else {
"none"
}
}
fn render_timestamp(timestamp: &Timestamp) -> String {
serde_json::to_value(timestamp)
.ok()
.and_then(|value| value.as_str().map(ToString::to_string))
.unwrap_or_else(|| format!("{timestamp:?}"))
}
fn browser_breadcrumbs(prefix: &str) -> Vec<BrowserBreadcrumb> {
let normalized = normalize_browser_prefix(prefix);
let mut breadcrumbs = vec![BrowserBreadcrumb {
label: "Root".to_string(),
prefix: String::new(),
}];
if normalized.is_empty() {
return breadcrumbs;
}
let mut current = String::new();
for segment in normalized.trim_end_matches('/').split('/') {
current.push_str(segment);
current.push('/');
breadcrumbs.push(BrowserBreadcrumb {
label: segment.to_string(),
prefix: current.clone(),
});
}
breadcrumbs
}
fn browser_metadata_matches(meta: &StoredObjectMeta, query_lower: &str) -> bool {
meta.metadata
.iter()
.any(|(key, value)| {
key.to_ascii_lowercase().contains(query_lower)
|| value.to_ascii_lowercase().contains(query_lower)
})
}
fn lookup_object_by_variants(
objects: &BTreeMap<String, StoredObjectMeta>,
bucket: &str,
key: &str,
) -> Option<(String, StoredObjectMeta)> {
let mut variants = Vec::<String>::new();
variants.push(key.to_string());
if key.contains('%') {
if let Some(decoded) = percent_decode_ascii(key) {
variants.push(decoded);
}
}
if key.contains(' ') {
variants.push(key.replace(' ', "+"));
}
if key.contains('+') {
variants.push(key.replace('+', " "));
}
for candidate in variants {
let storage = storage_key(bucket, &candidate);
if let Some(meta) = objects.get(&storage) {
return Some((candidate, meta.clone()));
}
}
let bucket_prefix = storage_key_prefix(bucket, "");
for (full_key, meta) in objects.iter().filter(|(full_key, _)| full_key.starts_with(&bucket_prefix))
{
let Some(object_key) = strip_storage_prefix(bucket, full_key) else {
continue;
};
if keys_equivalent(&object_key, key) {
return Some((object_key, meta.clone()));
}
}
None
}
fn build_browser_list_response(
state: &GatewayState,
bucket: &str,
prefix: &str,
query: &str,
cursor: Option<&str>,
limit: usize,
) -> Result<BrowserListResponse> {
let buckets = state
.buckets
.read()
.map_err(|_| anyhow::anyhow!("bucket state unavailable"))?;
if !buckets.contains(bucket) {
anyhow::bail!("bucket not found");
}
drop(buckets);
let objects = state
.objects
.read()
.map_err(|_| anyhow::anyhow!("object index unavailable"))?;
let normalized_prefix = normalize_browser_prefix(prefix);
let storage_prefix = storage_key_prefix(bucket, normalized_prefix.as_str());
let query_lower = query.trim().to_ascii_lowercase();
let limit = limit.clamp(1, 500);
let mut folders = BTreeMap::<String, (u64, u64)>::new();
let mut object_items = Vec::<BrowserListItem>::new();
for (full_key, meta) in objects
.iter()
.filter(|(key, _)| key.starts_with(storage_prefix.as_str()))
{
let Some(object_key) = strip_storage_prefix(bucket, full_key) else {
continue;
};
let content_type = browser_content_type_from_metadata(&meta.metadata, object_key.as_str());
let preview_kind = browser_preview_kind(
object_key.as_str(),
content_type.as_deref(),
)
.to_string();
if !query_lower.is_empty() {
let object_key_lower = object_key.to_ascii_lowercase();
if !object_key_lower.contains(query_lower.as_str())
&& !browser_metadata_matches(meta, query_lower.as_str())
{
continue;
}
object_items.push(BrowserListItem {
kind: "object".to_string(),
name: object_key
.rsplit('/')
.next()
.unwrap_or(object_key.as_str())
.to_string(),
key: object_key.clone(),
prefix: normalized_prefix.clone(),
size_bytes: Some(meta.content_len),
child_count: None,
content_type,
last_modified: Some(render_timestamp(&meta.last_modified)),
etag: Some(meta.etag.clone()),
cursor: format!("o:{object_key}"),
preview_kind: preview_kind.clone(),
can_preview: preview_kind != "none",
can_download: true,
can_delete: true,
});
continue;
}
let remainder = object_key
.strip_prefix(normalized_prefix.as_str())
.unwrap_or(object_key.as_str());
if let Some(separator) = remainder.find('/') {
let folder_prefix = format!("{}{}", normalized_prefix, &remainder[..separator + 1]);
let aggregate = folders.entry(folder_prefix).or_insert((0, 0));
aggregate.0 = aggregate.0.saturating_add(meta.content_len);
aggregate.1 = aggregate.1.saturating_add(1);
continue;
}
object_items.push(BrowserListItem {
kind: "object".to_string(),
name: remainder.to_string(),
key: object_key.clone(),
prefix: normalized_prefix.clone(),
size_bytes: Some(meta.content_len),
child_count: None,
content_type,
last_modified: Some(render_timestamp(&meta.last_modified)),
etag: Some(meta.etag.clone()),
cursor: format!("o:{object_key}"),
preview_kind: preview_kind.clone(),
can_preview: preview_kind != "none",
can_download: true,
can_delete: true,
});
}
let mut items = folders
.into_iter()
.map(|(folder_prefix, (bytes, child_count))| BrowserListItem {
kind: "folder".to_string(),
name: folder_prefix
.trim_end_matches('/')
.rsplit('/')
.next()
.unwrap_or(folder_prefix.as_str())
.to_string(),
key: String::new(),
prefix: folder_prefix.clone(),
size_bytes: Some(bytes),
child_count: Some(child_count),
content_type: None,
last_modified: None,
etag: None,
cursor: format!("d:{folder_prefix}"),
preview_kind: "folder".to_string(),
can_preview: false,
can_download: false,
can_delete: false,
})
.collect::<Vec<_>>();
items.extend(object_items);
items.sort_by(|left, right| {
let left_kind = if left.kind == "folder" { 0 } else { 1 };
let right_kind = if right.kind == "folder" { 0 } else { 1 };
left_kind
.cmp(&right_kind)
.then_with(|| left.name.to_ascii_lowercase().cmp(&right.name.to_ascii_lowercase()))
.then_with(|| left.cursor.cmp(&right.cursor))
});
let scope = BrowserScopeSummary {
folder_count: items.iter().filter(|item| item.kind == "folder").count(),
object_count: items.iter().filter(|item| item.kind == "object").count(),
total_bytes: items
.iter()
.filter_map(|item| item.size_bytes)
.sum::<u64>(),
};
let filtered = if let Some(cursor) = cursor {
items.into_iter()
.skip_while(|item| item.cursor.as_str() <= cursor)
.collect::<Vec<_>>()
} else {
items
};
let next_cursor = filtered.get(limit).map(|item| item.cursor.clone());
let page_items = filtered.into_iter().take(limit).collect::<Vec<_>>();
Ok(BrowserListResponse {
bucket: bucket.to_string(),
prefix: normalized_prefix.clone(),
query: query.trim().to_string(),
limit,
next_cursor,
breadcrumbs: browser_breadcrumbs(normalized_prefix.as_str()),
scope,
capabilities: BrowserCapabilities {
can_upload: true,
can_delete: true,
can_preview: true,
can_download: true,
},
items: page_items,
})
}
fn build_browser_object_details(
state: &GatewayState,
bucket: &str,
key: &str,
) -> Result<BrowserObjectDetails> {
let objects = state
.objects
.read()
.map_err(|_| anyhow::anyhow!("object index unavailable"))?;
let (resolved_key, meta) = lookup_object_by_variants(&objects, bucket, key)
.ok_or_else(|| anyhow::anyhow!("object not found"))?;
let content_type = browser_content_type_from_metadata(&meta.metadata, resolved_key.as_str())
.unwrap_or_else(|| "application/octet-stream".to_string());
let preview_kind = browser_preview_kind(resolved_key.as_str(), Some(content_type.as_str()));
Ok(BrowserObjectDetails {
bucket: bucket.to_string(),
key: resolved_key.clone(),
name: resolved_key
.rsplit('/')
.next()
.unwrap_or(resolved_key.as_str())
.to_string(),
size_bytes: meta.content_len,
content_type,
preview_kind: preview_kind.to_string(),
etag: meta.etag.clone(),
last_modified: render_timestamp(&meta.last_modified),
metadata: meta.metadata.into_iter().collect::<BTreeMap<_, _>>(),
can_preview: preview_kind != "none",
can_download: true,
can_delete: true,
})
}
fn put_console_object(
state: &GatewayState,
bucket: &str,
key: &str,
raw: &[u8],
content_type: Option<&str>,
) -> Result<BrowserUploadResult> {
let buckets = state
.buckets
.read()
.map_err(|_| anyhow::anyhow!("bucket state unavailable"))?;
if !buckets.contains(bucket) {
anyhow::bail!("bucket not found");
}
drop(buckets);
let object_key = storage_key(bucket, key);
let is_update = state
.objects
.read()
.map_err(|_| anyhow::anyhow!("object index unavailable"))?
.contains_key(&object_key);
let nonce_seed = next_nonce_seed(state).saturating_mul(1_000_000);
let compression_policy = state
.compression_policies
.read()
.ok()
.and_then(|m| m.get(bucket).cloned())
.unwrap_or_default();
let decision = compression::decide(content_type, raw, &compression_policy);
if let Ok(mut guard) = state.last_compression_decision.write() {
*guard = Some(decision);
}
let qos_profile = state
.qos_profiles
.read()
.ok()
.and_then(|m| m.get(bucket).copied())
.unwrap_or(BucketQosProfile::Standard);
let active_mpu = state.uploads.read().map(|g| g.len()).unwrap_or(0);
let chunk_bytes = crate::ingest::compute_chunk_bytes(
state.chunk_runtime.current_chunk_bytes(),
qos_profile,
active_mpu,
);
let processed = process_payload_crc_zstd_aes(raw, chunk_bytes, &PIPELINE_KEY, nonce_seed)
.map_err(|_| anyhow::anyhow!("pipeline processing failed"))?;
let blob = encode_pipeline_blob(&processed.chunks)
.map_err(|_| anyhow::anyhow!("pipeline blob encoding failed"))?;
let pointer = state
.data_plane
.append_blob(&blob, nonce_seed)
.map_err(|_| anyhow::anyhow!("device append failed"))?;
let etag = md5_hex(raw);
append_object_wal_ordered(
state,
WalEntryV1 {
seq: 0,
op: "put".to_string(),
bucket: bucket.to_string(),
key: key.to_string(),
etag: Some(etag.clone()),
pointer_or_manifest: Some(pointer.clone()),
ts_unix_ms: now_unix_ms(),
},
)
.map_err(|err| anyhow::anyhow!(err.to_string()))?;
let mut metadata = HashMap::new();
if let Some(content_type) = content_type {
metadata.insert("content-type".to_string(), content_type.to_string());
}
state
.objects
.write()
.map_err(|_| anyhow::anyhow!("object index unavailable"))?
.insert(
object_key,
StoredObjectMeta {
etag: etag.clone(),
content_len: raw.len() as u64,
last_modified: Timestamp::from(SystemTime::now()),
metadata,
layout: StoredLayout::InlineBlob { pointer },
},
);
enqueue_webhook_intents(state, bucket, key, raw.len() as u64, etag.as_str(), is_update)
.map_err(|err| anyhow::anyhow!(err.to_string()))?;
Ok(BrowserUploadResult {
bucket: bucket.to_string(),
key: key.to_string(),
size_bytes: raw.len() as u64,
content_type: content_type.unwrap_or("application/octet-stream").to_string(),
etag,
uploaded: true,
})
}
fn delete_console_object(state: &GatewayState, bucket: &str, key: &str) -> Result<BrowserDeleteResult> {
append_object_wal_ordered(
state,
WalEntryV1 {
seq: 0,
op: "delete".to_string(),
bucket: bucket.to_string(),
key: key.to_string(),
etag: None,
pointer_or_manifest: None,
ts_unix_ms: now_unix_ms(),
},
)
.map_err(|err| anyhow::anyhow!(err.to_string()))?;
let removed = state
.objects
.write()
.map_err(|_| anyhow::anyhow!("object index unavailable"))
.map(|mut objects| remove_object_by_variants(&mut objects, bucket, key))?;
Ok(BrowserDeleteResult {
bucket: bucket.to_string(),
key: key.to_string(),
deleted: removed,
})
}
fn guess_panopticon_content_type(path: &str) -> &'static str {
if path.ends_with(".html") {
"text/html; charset=utf-8"
} else if path.ends_with(".css") {
"text/css; charset=utf-8"
} else if path.ends_with(".js") {
"application/javascript; charset=utf-8"
} else if path.ends_with(".json") {
"application/json; charset=utf-8"
} else if path.ends_with(".svg") {
"image/svg+xml"
} else if path.ends_with(".png") {
"image/png"
} else if path.ends_with(".ico") {
"image/x-icon"
} else if path.ends_with(".woff2") {
"font/woff2"
} else {
"application/octet-stream"
}
}
fn panopticon_file_response(path: &str) -> HttpResponse {
let normalized = path.trim_start_matches("/_nexus/ui").trim_start_matches('/');
let asset_path = if normalized.is_empty() {
"index.html"
} else {
normalized
};
if let Some(file) = PANOPTICON_DIST.get_file(asset_path) {
return http::Response::builder()
.status(StatusCode::OK)
.header(
http::header::CONTENT_TYPE,
guess_panopticon_content_type(asset_path),
)
.body(Body::from(file.contents().to_vec()))
.unwrap_or_else(|_| http::Response::new(Body::from(Vec::new())));
}
if asset_path.contains('.') {
return xml_error_response(
StatusCode::NOT_FOUND,
"NoSuchResource",
"panopticon asset not found",
);
}
PANOPTICON_DIST
.get_file("index.html")
.map(|file| {
http::Response::builder()
.status(StatusCode::OK)
.header(http::header::CONTENT_TYPE, "text/html; charset=utf-8")
.body(Body::from(file.contents().to_vec()))
.unwrap_or_else(|_| http::Response::new(Body::from(Vec::new())))
})
.unwrap_or_else(|| {
xml_error_response(
StatusCode::NOT_FOUND,
"NoSuchResource",
"panopticon index missing",
)
})
}
fn remove_object_by_variants(
objects: &mut BTreeMap<String, StoredObjectMeta>,
bucket: &str,
key: &str,
) -> bool {
let mut variants = Vec::<String>::new();
variants.push(key.to_string());
if key.contains('%') {
if let Some(decoded) = percent_decode_ascii(key) {
variants.push(decoded);
}
}
if key.contains(' ') {
variants.push(key.replace(' ', "+"));
}
if key.contains('+') {
variants.push(key.replace('+', " "));
}
for candidate in variants {
let storage = storage_key(bucket, &candidate);
if objects.remove(&storage).is_some() {
return true;
}
}
let bucket_prefix = storage_key_prefix(bucket, "");
let full_keys = objects
.keys()
.filter(|full_key| full_key.starts_with(&bucket_prefix))
.cloned()
.collect::<Vec<_>>();
for full_key in full_keys {
let Some(object_key) = strip_storage_prefix(bucket, &full_key) else {
continue;
};
if keys_equivalent(&object_key, key) {
objects.remove(&full_key);
return true;
}
}
false
}
fn keys_equivalent(a: &str, b: &str) -> bool {
if a == b {
return true;
}
if a.replace('+', " ") == b || b.replace('+', " ") == a {
return true;
}
if percent_decode_ascii(a).as_deref() == Some(b)
|| percent_decode_ascii(b).as_deref() == Some(a)
{
return true;
}
false
}
fn percent_decode_ascii(input: &str) -> Option<String> {
let bytes = input.as_bytes();
let mut out = Vec::<u8>::with_capacity(bytes.len());
let mut idx = 0usize;
while idx < bytes.len() {
if bytes[idx] == b'%' {
if idx + 2 >= bytes.len() {
return None;
}
let hi = hex_nibble(bytes[idx + 1])?;
let lo = hex_nibble(bytes[idx + 2])?;
out.push((hi << 4) | lo);
idx += 3;
} else {
out.push(bytes[idx]);
idx += 1;
}
}
String::from_utf8(out).ok()
}
fn hex_nibble(b: u8) -> Option<u8> {
match b {
b'0'..=b'9' => Some(b - b'0'),
b'a'..=b'f' => Some(10 + (b - b'a')),
b'A'..=b'F' => Some(10 + (b - b'A')),
_ => None,
}
}
fn normalize_max_keys(max_keys: Option<i32>) -> S3Result<(i32, usize)> {
let max_keys_i32 = max_keys.unwrap_or(1000);
let max_keys_usize = usize::try_from(max_keys_i32).map_err(|_| s3_error!(InvalidArgument))?;
Ok((max_keys_i32, max_keys_usize))
}
fn normalize_max_uploads(max_uploads: Option<i32>) -> S3Result<usize> {
let max_uploads_i32 = max_uploads.unwrap_or(1000);
if max_uploads_i32 <= 0 {
return Err(s3_error!(InvalidArgument));
}
let as_usize = usize::try_from(max_uploads_i32).map_err(|_| s3_error!(InvalidArgument))?;
Ok(as_usize.min(1000))
}
fn validate_payload_headers(
content_md5: Option<&str>,
content_length: Option<i64>,
payload: &[u8],
) -> S3Result<()> {
let Some(length) = content_length else {
return Err(s3_error!(MissingContentLength));
};
if length < 0 {
return Err(s3_error!(InvalidArgument));
}
if length as usize != payload.len() {
return Err(s3_error!(BadDigest));
}
let Some(content_md5) = content_md5 else {
return Ok(());
};
if content_md5.is_empty() {
return Err(s3_error!(InvalidDigest));
}
let decoded = BASE64_STANDARD
.decode(content_md5)
.map_err(|_| s3_error!(InvalidDigest))?;
if decoded.len() != 16 {
return Err(s3_error!(InvalidDigest));
}
let digest = md5::compute(payload);
if digest.0.as_slice() != decoded.as_slice() {
return Err(s3_error!(BadDigest));
}
Ok(())
}
fn validate_get_preconditions(input: &GetObjectInput, meta: &StoredObjectMeta) -> S3Result<()> {
let current_etag = sanitize_etag(&meta.etag);
if let Some(if_match) = input.if_match.as_ref() {
if !etag_condition_matches(if_match, ¤t_etag) {
return Err(s3_error!(PreconditionFailed));
}
}
if let Some(if_none_match) = input.if_none_match.as_ref() {
if etag_condition_matches(if_none_match, ¤t_etag) {
return Err(s3_error!(NotModified));
}
}
if let Some(if_modified_since) = input.if_modified_since.clone() {
if meta.last_modified.clone() <= if_modified_since {
return Err(s3_error!(NotModified));
}
}
if let Some(if_unmodified_since) = input.if_unmodified_since.clone() {
if meta.last_modified.clone() > if_unmodified_since {
return Err(s3_error!(PreconditionFailed));
}
}
Ok(())
}
fn etag_condition_matches(condition: &ETagCondition, current_etag: &str) -> bool {
if condition.is_any() {
return true;
}
condition
.as_etag()
.map(|etag| sanitize_etag(etag.value()) == current_etag)
.unwrap_or(false)
}
fn default_owner() -> Owner {
Owner {
id: Some(DEFAULT_OWNER_ID.to_string()),
display_name: Some(DEFAULT_OWNER_DISPLAY_NAME.to_string()),
}
}
fn default_acl_grants() -> Vec<Grant> {
vec![Grant {
grantee: Some(Grantee {
display_name: Some(DEFAULT_OWNER_DISPLAY_NAME.to_string()),
email_address: None,
id: Some(DEFAULT_OWNER_ID.to_string()),
type_: s3s::dto::Type::from_static(s3s::dto::Type::CANONICAL_USER),
uri: None,
}),
permission: Some(Permission::from_static(Permission::FULL_CONTROL)),
}]
}
fn s3_url_encode_key(input: &str) -> String {
let mut out = String::with_capacity(input.len());
for b in input.as_bytes() {
if b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.' | b'~' | b'/') {
out.push(*b as char);
} else {
out.push('%');
out.push_str(&format!("{:02X}", b));
}
}
out
}
#[derive(Debug, Clone)]
struct ProxyReadResult {
bytes: Vec<u8>,
etag: String,
content_len: u64,
}
fn maybe_proxy_read_and_backfill(
state: &GatewayState,
bucket: &str,
key: &str,
) -> S3Result<Option<ProxyReadResult>> {
let policy = state.migration.bucket_policy(bucket);
if !policy.proxy_read_enabled {
return Ok(None);
}
let Some(source) = state.migration.source() else {
return Ok(None);
};
let source_bucket = policy
.source_bucket
.as_deref()
.filter(|v| !v.is_empty())
.unwrap_or(bucket);
let source_prefix = policy.source_prefix.unwrap_or_default();
let source_key = if source_prefix.is_empty() {
key.to_string()
} else {
format!("{}{}", source_prefix.trim_end_matches('/').to_string() + "/", key)
};
let object_url = build_source_object_url(source.endpoint.as_str(), source_bucket, source_key.as_str())
.map_err(|_| s3_error!(InternalError, "invalid migration source endpoint"))?;
let client = reqwest::blocking::Client::builder()
.timeout(Duration::from_secs(20))
.build()
.map_err(|_| s3_error!(InternalError, "failed building migration client"))?;
let response = client
.get(object_url)
.send()
.map_err(|_| s3_error!(InternalError, "proxy-read source fetch failed"))?;
if response.status() == StatusCode::NOT_FOUND {
return Ok(None);
}
if !response.status().is_success() {
return Err(s3_error!(
InternalError,
"proxy-read source fetch returned non-success"
));
}
let bytes = response
.bytes()
.map_err(|_| s3_error!(InternalError, "proxy-read source body failed"))?
.to_vec();
let etag = md5_hex(bytes.as_slice());
if policy.backfill_enabled {
let state = state.clone();
let bucket = bucket.to_string();
let key = key.to_string();
let bytes_for_backfill = bytes.clone();
let etag_for_backfill = etag.clone();
thread::spawn(move || {
let nonce_seed = next_nonce_seed(&state).saturating_mul(1_000_000);
let chunk_bytes = state.chunk_runtime.current_chunk_bytes();
let Ok(processed) =
process_payload_crc_zstd_aes(&bytes_for_backfill, chunk_bytes, &PIPELINE_KEY, nonce_seed)
else {
return;
};
let Ok(blob) = encode_pipeline_blob(&processed.chunks) else {
return;
};
let Ok(pointer) = state.data_plane.append_blob(&blob, nonce_seed) else {
return;
};
let _ = append_object_wal_ordered(
&state,
WalEntryV1 {
seq: 0,
op: "put".to_string(),
bucket: bucket.clone(),
key: key.clone(),
etag: Some(etag_for_backfill.clone()),
pointer_or_manifest: Some(pointer.clone()),
ts_unix_ms: now_unix_ms(),
},
);
if let Ok(mut guard) = state.objects.write() {
guard.insert(
storage_key(&bucket, &key),
StoredObjectMeta {
etag: etag_for_backfill,
content_len: bytes_for_backfill.len() as u64,
last_modified: Timestamp::from(SystemTime::now()),
metadata: HashMap::new(),
layout: StoredLayout::InlineBlob { pointer },
},
);
}
});
}
Ok(Some(ProxyReadResult {
content_len: bytes.len() as u64,
bytes,
etag,
}))
}
fn build_source_object_url(endpoint: &str, bucket: &str, key: &str) -> Result<String> {
let mut url = url::Url::parse(endpoint)?;
let root = url.path().trim_end_matches('/').to_string();
let encoded_key = key
.split('/')
.map(urlencoding::encode)
.map(|s| s.into_owned())
.collect::<Vec<_>>()
.join("/");
let full_path = if root.is_empty() || root == "/" {
format!("/{bucket}/{encoded_key}")
} else {
format!("{root}/{bucket}/{encoded_key}")
};
url.set_path(full_path.as_str());
Ok(url.to_string())
}
enum InterceptResult {
Response(HttpResponse),
Forward(http::Request<HyperIncoming>),
}
async fn maybe_intercept_request(
mut req: http::Request<HyperIncoming>,
state: GatewayState,
expected_access_key: &str,
) -> Result<InterceptResult> {
let query = req.uri().query().unwrap_or_default().to_string();
if query.contains("select") && query.contains("select-type=2") {
return Ok(InterceptResult::Response(xml_error_response(
StatusCode::NOT_IMPLEMENTED,
"NotImplemented",
"SelectObjectContent is not implemented in this release",
)));
}
if req.uri().path().starts_with("/_nexus/") {
let response = handle_nexus_control_request(req, state).await?;
return Ok(InterceptResult::Response(response));
}
if query.contains("nexus-compression-policy") || query.contains("nexus-qos-profile") {
let response = handle_bucket_policy_request(req, state).await?;
return Ok(InterceptResult::Response(response));
}
if query.contains("nexus-webhooks") {
let response = handle_webhook_request(req, state).await?;
return Ok(InterceptResult::Response(response));
}
let mut params = parse_query_params(query.as_str());
let parse_form_body = req.method() == http::Method::POST
&& (req.uri().path() == "/" || req.uri().path().is_empty());
if parse_form_body {
let content_type = req
.headers()
.get(http::header::CONTENT_TYPE)
.and_then(|value| value.to_str().ok())
.unwrap_or_default();
if content_type.contains("application/x-www-form-urlencoded")
|| content_type.contains("text/plain")
{
let bytes = req.body_mut().collect().await?.to_bytes();
let body_params = url::form_urlencoded::parse(bytes.as_ref())
.into_owned()
.collect::<Vec<_>>();
for (k, v) in body_params {
params.entry(k).or_insert(v);
}
}
}
if let Some(action) = params.get("Action").cloned() {
let response =
close_connection_on_error(handle_query_action(action.as_str(), ¶ms, &req, &state));
return Ok(InterceptResult::Response(response));
}
if let Some(resp) = validate_legacy_auth_headers(req.headers(), expected_access_key) {
return Ok(InterceptResult::Response(resp));
}
Ok(InterceptResult::Forward(req))
}
async fn handle_webhook_request(
mut req: http::Request<HyperIncoming>,
state: GatewayState,
) -> Result<HttpResponse> {
let bucket = req
.uri()
.path()
.trim_start_matches('/')
.split('/')
.next()
.filter(|value| !value.is_empty())
.unwrap_or_default()
.to_string();
if bucket.is_empty() {
return Ok(xml_error_response(
StatusCode::BAD_REQUEST,
"InvalidArgument",
"bucket path is required for nexus-webhooks",
));
}
let query = parse_query_params(req.uri().query().unwrap_or_default());
match *req.method() {
http::Method::PUT => {
let bytes = req.body_mut().collect().await?.to_bytes();
let payload = bytes.as_ref();
let incoming_rules = parse_webhook_rules_payload(payload)?;
let rules_json = serde_json::to_string(&incoming_rules)?;
state
.webhooks
.replace_rules(bucket.as_str(), incoming_rules)?;
let _ = append_webhook_wal_ordered(
&state,
"replace_bucket_rules",
bucket.as_str(),
"",
rules_json.as_str(),
);
Ok(json_response(
StatusCode::OK,
json!({"bucket": bucket, "status": "replaced"}),
))
}
http::Method::GET => Ok(json_response(
StatusCode::OK,
json!({"bucket": bucket, "rules": state.webhooks.get_rules(bucket.as_str())}),
)),
http::Method::DELETE => {
let rule_id = query.get("id").map(String::as_str);
state.webhooks.delete_rule(bucket.as_str(), rule_id);
match rule_id {
Some(id) => {
let _ = append_webhook_wal_ordered(
&state,
"delete_rule",
bucket.as_str(),
id,
"{}",
);
}
None => {
let _ = append_webhook_wal_ordered(
&state,
"delete_bucket_rules",
bucket.as_str(),
"",
"{}",
);
}
}
Ok(json_response(
StatusCode::OK,
json!({"bucket": bucket, "status": "deleted"}),
))
}
_ => Ok(xml_error_response(
StatusCode::METHOD_NOT_ALLOWED,
"MethodNotAllowed",
"nexus-webhooks supports PUT/GET/DELETE",
)),
}
}
#[derive(Debug, Deserialize)]
struct WebhookRulesEnvelope {
rules: Vec<WebhookRule>,
}
fn parse_webhook_rules_payload(payload: &[u8]) -> Result<Vec<WebhookRule>> {
if payload.is_empty() {
return Ok(Vec::new());
}
if let Ok(envelope) = serde_json::from_slice::<WebhookRulesEnvelope>(payload) {
return Ok(envelope.rules);
}
let list = serde_json::from_slice::<Vec<WebhookRule>>(payload)
.context("invalid webhook JSON payload")?;
Ok(list)
}
#[derive(Debug, Deserialize)]
struct PurgeRequest {
danger: bool,
bucket_prefix: Option<String>,
}
#[derive(Debug, Deserialize)]
struct PresignPreviewRequest {
bucket: String,
key: String,
ttl_seconds: Option<u64>,
}
#[derive(Debug, Deserialize)]
struct MigrationImportRequest {
config_path: String,
iam_dir: String,
}
#[derive(Debug, Deserialize)]
struct MigrationSyncStartRequest {
bucket: String,
prefix: Option<String>,
workers: Option<usize>,
}
#[derive(Debug, Deserialize)]
struct BrowserObjectUrlRequest {
bucket: String,
key: String,
mode: String,
}
#[derive(Debug, Deserialize)]
struct BrowserDeleteRequest {
bucket: String,
key: String,
}
#[derive(Debug, Clone, Serialize)]
struct BrowserBucketSummary {
name: String,
object_count: u64,
bytes: u64,
compression_policy: bool,
qos_profile: String,
webhook_rule_count: usize,
}
#[derive(Debug, Clone, Serialize)]
struct BrowserBreadcrumb {
label: String,
prefix: String,
}
#[derive(Debug, Clone, Serialize)]
struct BrowserCapabilities {
can_upload: bool,
can_delete: bool,
can_preview: bool,
can_download: bool,
}
#[derive(Debug, Clone, Serialize)]
struct BrowserScopeSummary {
folder_count: usize,
object_count: usize,
total_bytes: u64,
}
#[derive(Debug, Clone, Serialize)]
struct BrowserListItem {
kind: String,
name: String,
key: String,
prefix: String,
size_bytes: Option<u64>,
child_count: Option<u64>,
content_type: Option<String>,
last_modified: Option<String>,
etag: Option<String>,
cursor: String,
preview_kind: String,
can_preview: bool,
can_download: bool,
can_delete: bool,
}
#[derive(Debug, Clone, Serialize)]
struct BrowserListResponse {
bucket: String,
prefix: String,
query: String,
limit: usize,
next_cursor: Option<String>,
breadcrumbs: Vec<BrowserBreadcrumb>,
scope: BrowserScopeSummary,
capabilities: BrowserCapabilities,
items: Vec<BrowserListItem>,
}
#[derive(Debug, Clone, Serialize)]
struct BrowserObjectDetails {
bucket: String,
key: String,
name: String,
size_bytes: u64,
content_type: String,
preview_kind: String,
etag: String,
last_modified: String,
metadata: BTreeMap<String, String>,
can_preview: bool,
can_download: bool,
can_delete: bool,
}
#[derive(Debug, Clone, Serialize)]
struct BrowserObjectActionUrl {
bucket: String,
key: String,
mode: String,
url: String,
expires_in_seconds: u64,
}
#[derive(Debug, Clone, Serialize)]
struct BrowserUploadResult {
bucket: String,
key: String,
size_bytes: u64,
content_type: String,
etag: String,
uploaded: bool,
}
#[derive(Debug, Clone, Serialize)]
struct BrowserDeleteResult {
bucket: String,
key: String,
deleted: bool,
}
async fn handle_nexus_control_request(
mut req: http::Request<HyperIncoming>,
state: GatewayState,
) -> Result<HttpResponse> {
let path = req.uri().path().to_string();
if *req.method() == http::Method::GET && path.starts_with("/_nexus/ui") {
return Ok(panopticon_file_response(path.as_str()));
}
match (req.method().clone(), path.as_str()) {
(http::Method::GET, "/_nexus/console/browser/list") => {
let params = parse_query_params(req.uri().query().unwrap_or_default());
let bucket = params.get("bucket").cloned().unwrap_or_default();
if bucket.is_empty() {
return Ok(xml_error_response(
StatusCode::BAD_REQUEST,
"InvalidRequest",
"bucket is required",
));
}
let prefix = params.get("prefix").cloned().unwrap_or_default();
let query = params.get("q").cloned().unwrap_or_default();
let cursor = params.get("cursor").cloned();
let limit = params
.get("limit")
.and_then(|value| value.parse::<usize>().ok())
.unwrap_or(200);
let response = build_browser_list_response(
&state,
bucket.as_str(),
prefix.as_str(),
query.as_str(),
cursor.as_deref(),
limit,
)
.map_err(|error| error.to_string())
.map(|payload| json_response(StatusCode::OK, serde_json::to_value(payload).unwrap()));
return Ok(match response {
Ok(ok) => ok,
Err(message) if message.contains("bucket not found") => xml_error_response(
StatusCode::NOT_FOUND,
"NoSuchBucket",
"bucket not found",
),
Err(message) => xml_error_response(
StatusCode::INTERNAL_SERVER_ERROR,
"InternalError",
message.as_str(),
),
});
}
(http::Method::GET, "/_nexus/console/browser/object") => {
let params = parse_query_params(req.uri().query().unwrap_or_default());
let bucket = params.get("bucket").cloned().unwrap_or_default();
let key = params.get("key").cloned().unwrap_or_default();
if bucket.is_empty() || key.is_empty() {
return Ok(xml_error_response(
StatusCode::BAD_REQUEST,
"InvalidRequest",
"bucket and key are required",
));
}
let response = build_browser_object_details(&state, bucket.as_str(), key.as_str())
.map(|payload| json_response(StatusCode::OK, serde_json::to_value(payload).unwrap()));
return Ok(match response {
Ok(ok) => ok,
Err(error) if error.to_string().contains("not found") => xml_error_response(
StatusCode::NOT_FOUND,
"NoSuchKey",
"object not found",
),
Err(error) => xml_error_response(
StatusCode::INTERNAL_SERVER_ERROR,
"InternalError",
error.to_string().as_str(),
),
});
}
(http::Method::GET, "/_nexus/console/browser/object-content") => {
let params = parse_query_params(req.uri().query().unwrap_or_default());
let bucket = params.get("bucket").cloned().unwrap_or_default();
let key = params.get("key").cloned().unwrap_or_default();
let mode = params
.get("mode")
.cloned()
.unwrap_or_else(|| "preview".to_string());
if bucket.is_empty() || key.is_empty() {
return Ok(xml_error_response(
StatusCode::BAD_REQUEST,
"InvalidRequest",
"bucket and key are required",
));
}
let objects = state
.objects
.read()
.map_err(|_| anyhow::anyhow!("object index unavailable"))?;
let Some((resolved_key, meta)) =
lookup_object_by_variants(&objects, bucket.as_str(), key.as_str())
else {
return Ok(xml_error_response(StatusCode::NOT_FOUND, "NoSuchKey", "object not found"));
};
drop(objects);
let body = materialize_object(&state, &meta)
.map_err(|_| anyhow::anyhow!("failed materializing object"))?;
let content_type = browser_content_type_from_metadata(&meta.metadata, resolved_key.as_str())
.unwrap_or_else(|| "application/octet-stream".to_string());
let content_disposition = if mode == "download" {
format!("attachment; filename=\"{}\"", resolved_key.rsplit('/').next().unwrap_or(resolved_key.as_str()))
} else {
"inline".to_string()
};
return Ok(http::Response::builder()
.status(StatusCode::OK)
.header(http::header::CONTENT_TYPE, content_type)
.header(http::header::CONTENT_DISPOSITION, content_disposition)
.body(Body::from(body))
.unwrap());
}
(http::Method::POST, "/_nexus/console/browser/object-url") => {
let bytes = req.body_mut().collect().await?.to_bytes();
let payload: BrowserObjectUrlRequest = serde_json::from_slice(&bytes)?;
let details = build_browser_object_details(&state, payload.bucket.as_str(), payload.key.as_str())
.map_err(|error| anyhow::anyhow!(error.to_string()))?;
let host = req
.headers()
.get(http::header::HOST)
.and_then(|value| value.to_str().ok())
.unwrap_or("localhost:8080");
let expires = 300_u64;
let mode = match payload.mode.as_str() {
"download" => "download",
_ => "preview",
};
let token = Uuid::new_v4().simple().to_string();
let url = format!(
"http://{host}/_nexus/console/browser/object-content?bucket={}&key={}&mode={mode}&expires={expires}&token={token}",
urlencoding::encode(payload.bucket.as_str()),
urlencoding::encode(details.key.as_str()),
);
return Ok(json_response(
StatusCode::OK,
serde_json::to_value(BrowserObjectActionUrl {
bucket: payload.bucket,
key: details.key,
mode: mode.to_string(),
url,
expires_in_seconds: expires,
})
.unwrap(),
));
}
(http::Method::POST, "/_nexus/console/browser/upload") => {
let params = parse_query_params(req.uri().query().unwrap_or_default());
let bucket = params.get("bucket").cloned().unwrap_or_default();
let prefix = normalize_browser_prefix(params.get("prefix").map(String::as_str).unwrap_or_default());
let name = req
.headers()
.get("x-nexus-upload-name")
.and_then(|value| value.to_str().ok())
.unwrap_or_default()
.trim_matches('/')
.to_string();
if bucket.is_empty() || name.is_empty() {
return Ok(xml_error_response(
StatusCode::BAD_REQUEST,
"InvalidRequest",
"bucket and x-nexus-upload-name are required",
));
}
let content_type = req
.headers()
.get(http::header::CONTENT_TYPE)
.and_then(|value| value.to_str().ok())
.map(str::to_string);
let bytes = req.body_mut().collect().await?.to_bytes();
let key = format!("{prefix}{name}");
let uploaded = put_console_object(
&state,
bucket.as_str(),
key.as_str(),
&bytes,
content_type.as_deref(),
)
.map_err(|error| error.to_string())
.map(|payload| json_response(StatusCode::OK, serde_json::to_value(payload).unwrap()));
return Ok(match uploaded {
Ok(ok) => ok,
Err(message) if message.contains("bucket not found") => xml_error_response(
StatusCode::NOT_FOUND,
"NoSuchBucket",
"bucket not found",
),
Err(message) => xml_error_response(
StatusCode::INTERNAL_SERVER_ERROR,
"InternalError",
message.as_str(),
),
});
}
(http::Method::DELETE, "/_nexus/console/browser/object") => {
let bytes = req.body_mut().collect().await?.to_bytes();
let payload: BrowserDeleteRequest = serde_json::from_slice(&bytes)?;
let deleted = delete_console_object(&state, payload.bucket.as_str(), payload.key.as_str())
.map(|result| json_response(StatusCode::OK, serde_json::to_value(result).unwrap()))
.unwrap_or_else(|error| {
xml_error_response(
StatusCode::INTERNAL_SERVER_ERROR,
"InternalError",
error.to_string().as_str(),
)
});
return Ok(deleted);
}
(http::Method::GET, "/_nexus/console/summary") => {
let bucket_summaries = collect_console_bucket_summaries(&state);
let usage_objects = bucket_summaries
.iter()
.map(|bucket| bucket.object_count)
.sum::<u64>();
let usage_bytes = bucket_summaries
.iter()
.map(|bucket| bucket.bytes)
.sum::<u64>();
let webhook_stats = state.webhooks.stats();
let recent_jobs = {
let mut jobs = state.migration.snapshot_jobs();
jobs.sort_by(|left, right| right.started_unix_ms.cmp(&left.started_unix_ms));
jobs.truncate(5);
jobs
};
let total_jobs = state.migration.snapshot_jobs();
let running_jobs = total_jobs
.iter()
.filter(|job| job.status == migration::MigrationJobStatus::Running)
.count() as u64;
let failed_jobs = total_jobs
.iter()
.filter(|job| job.status == migration::MigrationJobStatus::Failed)
.count() as u64;
let last_compression = state
.last_compression_decision
.read()
.ok()
.and_then(|decision| decision.clone());
return Ok(json_response(
StatusCode::OK,
json!({
"usage": {
"objects": usage_objects,
"bytes": usage_bytes,
"buckets": bucket_summaries.len(),
},
"runtime": {
"atc_mode": state.atc_decision.mode.as_str(),
"runtime_profile": state.atc_decision.runtime_profile.as_str(),
"io_model": state.atc_decision.io_model.as_str(),
"worker_model": state.atc_decision.worker_model.as_str(),
"device_count": state.atc_decision.device_count,
"decision_reason": state.atc_decision.reason.clone(),
"current_chunk_bytes": state.chunk_runtime.current_chunk_bytes(),
"downshift_events": state.chunk_runtime.downshift_events.load(Ordering::Relaxed),
"upshift_events": state.chunk_runtime.upshift_events.load(Ordering::Relaxed),
"median_queue_depth": state.data_plane.median_queue_depth(),
"buffered_flushes": state.buffered_sync_flushes.load(Ordering::Relaxed),
"buffered_flush_latency_ms_sum": state.buffered_sync_flush_latency_ms.load(Ordering::Relaxed),
"last_compression": last_compression,
},
"webhook": {
"enqueued_total": webhook_stats.enqueued_total,
"dropped_total": webhook_stats.dropped_total,
"queue_capacity": webhook_stats.queue_capacity,
},
"migration": {
"source_configured": state.migration.source().is_some(),
"source": state.migration.source(),
"total_jobs": total_jobs.len(),
"running_jobs": running_jobs,
"failed_jobs": failed_jobs,
"recent_jobs": recent_jobs,
},
"policies": {
"compression_buckets": bucket_summaries.iter().filter(|bucket| bucket.compression_policy).count(),
"qos_buckets": bucket_summaries.iter().filter(|bucket| bucket.qos_profile != "standard").count(),
"webhook_buckets": bucket_summaries.iter().filter(|bucket| bucket.webhook_rule_count > 0).count(),
}
}),
));
}
(http::Method::GET, "/_nexus/console/buckets") => {
let bucket_summaries = collect_console_bucket_summaries(&state);
return Ok(json_response(
StatusCode::OK,
json!({
"buckets": bucket_summaries,
}),
));
}
(http::Method::GET, "/_nexus/analytics/usage") => {
let objects = state.objects.read().ok().map(|g| g.len()).unwrap_or(0);
let bytes = state
.objects
.read()
.ok()
.map(|g| g.values().map(|m| m.content_len).sum::<u64>())
.unwrap_or(0);
let buckets = state.buckets.read().ok().map(|g| g.len()).unwrap_or(0);
let webhook_stats = state.webhooks.stats();
return Ok(json_response(
StatusCode::OK,
json!({
"objects": objects,
"bytes": bytes,
"buckets": buckets,
"webhook_enqueued_total": webhook_stats.enqueued_total,
"webhook_dropped_total": webhook_stats.dropped_total
}),
));
}
(http::Method::GET, "/_nexus/analytics/prefix") => {
let params = parse_query_params(req.uri().query().unwrap_or_default());
let bucket = params.get("bucket").cloned().unwrap_or_default();
let prefix = params.get("prefix").cloned().unwrap_or_default();
let storage_prefix = storage_key_prefix(bucket.as_str(), prefix.as_str());
let objects = state
.objects
.read()
.ok()
.map(|g| {
g.iter()
.filter(|(k, _)| k.starts_with(storage_prefix.as_str()))
.map(|(_, v)| v.content_len)
.collect::<Vec<_>>()
})
.unwrap_or_default();
return Ok(json_response(
StatusCode::OK,
json!({
"bucket": bucket,
"prefix": prefix,
"matches": objects.len(),
"bytes": objects.iter().sum::<u64>()
}),
));
}
(http::Method::POST, "/_nexus/presign/preview") => {
let bytes = req.body_mut().collect().await?.to_bytes();
let payload: PresignPreviewRequest = serde_json::from_slice(&bytes)?;
let host = req
.headers()
.get(http::header::HOST)
.and_then(|v| v.to_str().ok())
.unwrap_or("localhost:8080");
let ttl = payload.ttl_seconds.unwrap_or(900);
let token = Uuid::new_v4().simple().to_string();
let url = format!(
"http://{}/{}/{}?x-nexus-presign={}&x-nexus-ttl={}",
host,
payload.bucket,
payload.key,
token,
ttl
);
return Ok(json_response(
StatusCode::OK,
json!({
"url": url,
"expires_in_seconds": ttl
}),
));
}
(http::Method::POST, "/_nexus/admin/purge-data") => {
let bytes = req.body_mut().collect().await?.to_bytes();
let payload: PurgeRequest = serde_json::from_slice(&bytes)?;
if !payload.danger {
return Ok(xml_error_response(
StatusCode::BAD_REQUEST,
"InvalidRequest",
"danger=true is required",
));
}
if let Ok(mut objects) = state.objects.write() {
if let Some(prefix) = payload.bucket_prefix.as_ref() {
let key_prefix = storage_key_prefix(prefix.as_str(), "");
objects.retain(|k, _| !k.starts_with(key_prefix.as_str()));
} else {
objects.clear();
}
}
if let Ok(mut uploads) = state.uploads.write() {
uploads.clear();
}
let report = admin::purge_data(PurgeDataConfig {
danger: true,
bucket_prefix: payload.bucket_prefix,
wal_path: Some(state.wal.wal_path().to_path_buf()),
devices: state
.data_plane
.devices
.iter()
.map(|dev| dev.path.clone())
.collect(),
})?;
return Ok(json_response(
StatusCode::OK,
serde_json::to_value(report).unwrap_or_else(|_| json!({"status":"ok"})),
));
}
_ => {}
}
if path.starts_with("/_nexus/migration/buckets/") {
let bucket = path.trim_start_matches("/_nexus/migration/buckets/").to_string();
return Ok(match *req.method() {
http::Method::PUT => {
let bytes = req.body_mut().collect().await?.to_bytes();
let policy: MigrationBucketPolicy = serde_json::from_slice(&bytes)?;
state.migration.set_bucket_policy(bucket.as_str(), policy.clone());
let _ = append_migration_wal_ordered(
&state,
"set_bucket_policy",
bucket.as_str(),
"",
serde_json::to_string(&policy)
.unwrap_or_else(|_| "{}".to_string())
.as_str(),
);
json_response(
StatusCode::OK,
json!({"bucket": bucket, "policy": policy, "status": "updated"}),
)
}
http::Method::GET => {
let policy = state.migration.bucket_policy(bucket.as_str());
json_response(StatusCode::OK, json!({"bucket": bucket, "policy": policy}))
}
_ => xml_error_response(
StatusCode::METHOD_NOT_ALLOWED,
"MethodNotAllowed",
"bucket migration policy supports PUT/GET",
),
});
}
if path == "/_nexus/migration/import/minio" && *req.method() == http::Method::POST {
let bytes = req.body_mut().collect().await?.to_bytes();
let payload: MigrationImportRequest = serde_json::from_slice(&bytes)?;
let report = migration::import_minio_export(
Path::new(payload.config_path.as_str()),
Path::new(payload.iam_dir.as_str()),
)?;
let _ = append_migration_wal_ordered(
&state,
"import_minio",
"",
"",
serde_json::to_string(&report)
.unwrap_or_else(|_| "{}".to_string())
.as_str(),
);
return Ok(json_response(
StatusCode::OK,
serde_json::to_value(report).unwrap_or_else(|_| json!({"status":"ok"})),
));
}
if path == "/_nexus/migration/sync/start" && *req.method() == http::Method::POST {
let bytes = req.body_mut().collect().await?.to_bytes();
let payload: MigrationSyncStartRequest = serde_json::from_slice(&bytes)?;
let Some(source) = state.migration.source() else {
return Ok(xml_error_response(
StatusCode::BAD_REQUEST,
"InvalidRequest",
"bifrost source endpoint is not configured",
));
};
let prefix = payload.prefix.unwrap_or_default();
let job = state
.migration
.start_job(payload.bucket.as_str(), prefix.as_str(), payload.workers.unwrap_or(4));
let _ = append_migration_wal_ordered(
&state,
"sync_start",
payload.bucket.as_str(),
prefix.as_str(),
serde_json::to_string(&job)
.unwrap_or_else(|_| "{}".to_string())
.as_str(),
);
let state_for_worker = state.clone();
migration::spawn_sync_job(
(*state.migration).clone(),
job.clone(),
source,
move |key, body| {
let nonce_seed = next_nonce_seed(&state_for_worker).saturating_mul(1_000_000);
let processed = process_payload_crc_zstd_aes(
&body,
state_for_worker.chunk_runtime.current_chunk_bytes(),
&PIPELINE_KEY,
nonce_seed,
)?;
let blob = encode_pipeline_blob(&processed.chunks)?;
let pointer = state_for_worker.data_plane.append_blob(&blob, nonce_seed)?;
let etag = md5_hex(body.as_slice());
append_object_wal_ordered(
&state_for_worker,
WalEntryV1 {
seq: 0,
op: "put".to_string(),
bucket: payload.bucket.clone(),
key: key.to_string(),
etag: Some(etag.clone()),
pointer_or_manifest: Some(pointer.clone()),
ts_unix_ms: now_unix_ms(),
},
)?;
if let Ok(mut guard) = state_for_worker.objects.write() {
guard.insert(
storage_key(payload.bucket.as_str(), key),
StoredObjectMeta {
etag,
content_len: body.len() as u64,
last_modified: Timestamp::from(SystemTime::now()),
metadata: HashMap::new(),
layout: StoredLayout::InlineBlob { pointer },
},
);
}
Ok(())
},
);
return Ok(json_response(
StatusCode::ACCEPTED,
json!({"job_id": job.id, "status": "running"}),
));
}
if path.starts_with("/_nexus/migration/sync/") && *req.method() == http::Method::GET {
let id = path.trim_start_matches("/_nexus/migration/sync/");
let job = state.migration.get_job(id);
return Ok(match job {
Some(job) => json_response(
StatusCode::OK,
serde_json::to_value(job).unwrap_or_else(|_| json!({"status":"ok"})),
),
None => xml_error_response(StatusCode::NOT_FOUND, "NoSuchJob", "migration job not found"),
});
}
if path == "/_nexus/migration/status" && *req.method() == http::Method::GET {
return Ok(json_response(
StatusCode::OK,
json!({
"source": state.migration.source(),
"policies": state.migration.all_bucket_policies(),
"jobs": state.migration.snapshot_jobs(),
}),
));
}
Ok(xml_error_response(
StatusCode::NOT_FOUND,
"NoSuchResource",
"unknown _nexus control path",
))
}
async fn handle_bucket_policy_request(
mut req: http::Request<HyperIncoming>,
state: GatewayState,
) -> Result<HttpResponse> {
let bucket = req
.uri()
.path()
.trim_start_matches('/')
.split('/')
.next()
.filter(|value| !value.is_empty())
.unwrap_or_default()
.to_string();
if bucket.is_empty() {
return Ok(xml_error_response(
StatusCode::BAD_REQUEST,
"InvalidArgument",
"bucket path is required",
));
}
let query = parse_query_params(req.uri().query().unwrap_or_default());
if query.contains_key("nexus-compression-policy") {
return Ok(match *req.method() {
http::Method::PUT => {
let bytes = req.body_mut().collect().await?.to_bytes();
let policy: CompressionPolicy = serde_json::from_slice(&bytes)?;
if let Ok(mut guard) = state.compression_policies.write() {
guard.insert(bucket.clone(), policy.clone());
}
json_response(StatusCode::OK, json!({"bucket": bucket, "policy": policy}))
}
http::Method::GET => {
let policy = state
.compression_policies
.read()
.ok()
.and_then(|m| m.get(&bucket).cloned())
.unwrap_or_default();
json_response(StatusCode::OK, json!({"bucket": bucket, "policy": policy}))
}
_ => xml_error_response(
StatusCode::METHOD_NOT_ALLOWED,
"MethodNotAllowed",
"nexus-compression-policy supports PUT/GET",
),
});
}
if query.contains_key("nexus-qos-profile") {
return Ok(match *req.method() {
http::Method::PUT => {
let bytes = req.body_mut().collect().await?.to_bytes();
let payload: serde_json::Value = serde_json::from_slice(&bytes)?;
let profile = payload
.get("profile")
.and_then(|v| v.as_str())
.and_then(BucketQosProfile::from_str)
.unwrap_or(BucketQosProfile::Standard);
if let Ok(mut guard) = state.qos_profiles.write() {
guard.insert(bucket.clone(), profile);
}
json_response(StatusCode::OK, json!({"bucket": bucket, "profile": profile}))
}
http::Method::GET => {
let profile = state
.qos_profiles
.read()
.ok()
.and_then(|m| m.get(&bucket).copied())
.unwrap_or(BucketQosProfile::Standard);
json_response(StatusCode::OK, json!({"bucket": bucket, "profile": profile}))
}
_ => xml_error_response(
StatusCode::METHOD_NOT_ALLOWED,
"MethodNotAllowed",
"nexus-qos-profile supports PUT/GET",
),
});
}
Ok(xml_error_response(
StatusCode::BAD_REQUEST,
"InvalidRequest",
"unsupported nexus bucket policy request",
))
}
fn handle_query_action(
action: &str,
params: &BTreeMap<String, String>,
req: &http::Request<HyperIncoming>,
state: &GatewayState,
) -> HttpResponse {
let caller_access_key =
extract_access_key(req, params).unwrap_or_else(|| DEFAULT_ACCESS_KEY.to_string());
let caller = state
.credentials
.get(caller_access_key.as_str())
.map(|record| record.principal)
.unwrap_or_else(|| "arn:aws:iam:::root".to_string());
match action {
"CreateUser" => {
let user_name = required_param(params, "UserName");
match user_name.and_then(|name| state.iam.create_user(name).map(|_| name.to_string())) {
Ok(name) => {
let generated_access_key = format!("AKIA{}", Uuid::new_v4().simple());
let generated_secret = Uuid::new_v4().simple().to_string();
let credential = crate::auth_provider::CredentialRecord {
access_key: generated_access_key,
secret_key: generated_secret,
principal: format!("arn:aws:iam:::user/{name}"),
session_token: None,
expires_at_unix_ms: None,
kind: CredentialKind::User,
};
state.credentials.upsert(credential.clone());
if let Err(err) =
append_iam_wal_ordered(state, "create_user", name.as_str(), "{}")
{
return xml_error_response(
StatusCode::INTERNAL_SERVER_ERROR,
"InternalError",
&err.to_string(),
);
}
if let Ok(payload) = serde_json::to_string(&credential) {
let _ = append_iam_wal_ordered(
state,
"upsert_credential",
name.as_str(),
payload.as_str(),
);
}
iam_xml_response(
"CreateUser",
format!(
"<CreateUserResult><User><Path>/</Path><UserName>{}</UserName><UserId>{}</UserId><Arn>arn:aws:iam::{}:user/{}</Arn></User></CreateUserResult>",
xml_escape(name.as_str()),
xml_escape(name.as_str()),
xml_escape(name.as_str()),
xml_escape(name.as_str()),
),
)
}
Err(err) => xml_error_response(
StatusCode::BAD_REQUEST,
"InvalidParameterValue",
&err.to_string(),
),
}
}
"DeleteUser" => {
let out = (|| -> Result<String> {
let name = required_param(params, "UserName")?.to_string();
state.iam.delete_user(name.as_str())?;
Ok(name)
})();
match out {
Ok(name) => {
let principal = format!("arn:aws:iam:::user/{name}");
let to_remove = state
.credentials
.snapshot()
.iter()
.filter_map(|(ak, rec)| (rec.principal == principal).then_some(ak.clone()))
.collect::<Vec<_>>();
for access_key in to_remove {
state.credentials.remove(access_key.as_str());
}
let _ = append_iam_wal_ordered(state, "delete_user", name.as_str(), "{}");
iam_xml_response("DeleteUser", "<DeleteUserResult/>".to_string())
}
Err(err) => {
xml_error_response(StatusCode::NOT_FOUND, "NoSuchEntity", &err.to_string())
}
}
}
"PutUserPolicy" => {
let out = (|| -> Result<()> {
let user = required_param(params, "UserName")?;
let policy_name = required_param(params, "PolicyName")?;
let policy_document = required_param(params, "PolicyDocument")?;
if policy_name.len() > 128 || policy_document.len() > 131_072 {
anyhow::bail!("policy parameter limit exceeded");
}
if policy_contains_principal(policy_document) {
anyhow::bail!("Principal is not allowed for user inline policies");
}
state
.iam
.put_user_policy(user, policy_name, policy_document)?;
let payload = json!({
"policy_name": policy_name,
"policy_document": policy_document,
})
.to_string();
let _ = append_iam_wal_ordered(state, "put_user_policy", user, payload.as_str());
Ok(())
})();
match out {
Ok(_) => iam_xml_response("PutUserPolicy", "<PutUserPolicyResult/>".to_string()),
Err(err) => {
let text = err.to_string();
if text.contains("No such user") {
xml_error_response(StatusCode::NOT_FOUND, "NoSuchEntity", &text)
} else {
xml_error_response(StatusCode::BAD_REQUEST, "InvalidParameterValue", &text)
}
}
}
}
"GetUserPolicy" => {
let out = (|| -> Result<(String, String, String)> {
let user = required_param(params, "UserName")?.to_string();
let policy_name = required_param(params, "PolicyName")?.to_string();
let policy_doc = state
.iam
.get_user_policy(user.as_str(), policy_name.as_str())
.context("policy not found")?;
Ok((user, policy_name, policy_doc))
})();
match out {
Ok((user, policy_name, policy_doc)) => iam_xml_response(
"GetUserPolicy",
format!(
"<GetUserPolicyResult><UserName>{}</UserName><PolicyName>{}</PolicyName><PolicyDocument>{}</PolicyDocument></GetUserPolicyResult>",
xml_escape(user.as_str()),
xml_escape(policy_name.as_str()),
xml_escape(policy_doc.as_str()),
),
),
Err(err) => xml_error_response(StatusCode::NOT_FOUND, "NoSuchEntity", &err.to_string()),
}
}
"ListUserPolicies" => {
let out = (|| -> Result<(String, Vec<String>)> {
let user = required_param(params, "UserName")?.to_string();
let names = state
.iam
.list_user_policies(user.as_str())
.context("user not found")?;
Ok((user, names))
})();
match out {
Ok((user, names)) => {
let members = names
.into_iter()
.map(|name| format!("<member>{}</member>", xml_escape(name.as_str())))
.collect::<String>();
iam_xml_response(
"ListUserPolicies",
format!(
"<ListUserPoliciesResult><UserName>{}</UserName><PolicyNames>{}</PolicyNames><IsTruncated>false</IsTruncated></ListUserPoliciesResult>",
xml_escape(user.as_str()),
members
),
)
}
Err(err) => {
xml_error_response(StatusCode::NOT_FOUND, "NoSuchEntity", &err.to_string())
}
}
}
"DeleteUserPolicy" => {
let out = (|| -> Result<()> {
let user = required_param(params, "UserName")?;
let policy_name = required_param(params, "PolicyName")?;
state.iam.delete_user_policy(user, policy_name)?;
let payload = json!({ "policy_name": policy_name }).to_string();
let _ = append_iam_wal_ordered(state, "delete_user_policy", user, payload.as_str());
Ok(())
})();
match out {
Ok(_) => {
iam_xml_response("DeleteUserPolicy", "<DeleteUserPolicyResult/>".to_string())
}
Err(err) => {
xml_error_response(StatusCode::NOT_FOUND, "NoSuchEntity", &err.to_string())
}
}
}
"CreateRole" => {
let out = (|| -> Result<String> {
let role_name = required_param(params, "RoleName")?.to_string();
let assume_doc = required_param(params, "AssumeRolePolicyDocument")?;
state.iam.create_role(role_name.as_str(), assume_doc)?;
let payload = json!({ "assume_role_policy": assume_doc }).to_string();
let _ = append_iam_wal_ordered(
state,
"create_role",
role_name.as_str(),
payload.as_str(),
);
Ok(role_name)
})();
match out {
Ok(role_name) => iam_xml_response(
"CreateRole",
format!(
"<CreateRoleResult><Role><Arn>arn:aws:iam:::role/{}</Arn><RoleName>{}</RoleName></Role></CreateRoleResult>",
xml_escape(role_name.as_str()),
xml_escape(role_name.as_str()),
),
),
Err(err) => xml_error_response(StatusCode::BAD_REQUEST, "InvalidParameterValue", &err.to_string()),
}
}
"PutRolePolicy" => {
let out = (|| -> Result<()> {
let role = required_param(params, "RoleName")?;
let policy_name = required_param(params, "PolicyName")?;
let policy_document = required_param(params, "PolicyDocument")?;
if policy_name.len() > 128 || policy_document.len() > 131_072 {
anyhow::bail!("policy parameter limit exceeded");
}
state
.iam
.put_role_policy(role, policy_name, policy_document)?;
let payload = json!({
"policy_name": policy_name,
"policy_document": policy_document
})
.to_string();
let _ = append_iam_wal_ordered(state, "put_role_policy", role, payload.as_str());
Ok(())
})();
match out {
Ok(_) => iam_xml_response("PutRolePolicy", "<PutRolePolicyResult/>".to_string()),
Err(err) => {
let text = err.to_string();
if text.contains("No such role") {
xml_error_response(StatusCode::NOT_FOUND, "NoSuchEntity", &text)
} else {
xml_error_response(StatusCode::BAD_REQUEST, "InvalidParameterValue", &text)
}
}
}
}
"GetRolePolicy" => {
let out = (|| -> Result<(String, String, String)> {
let role = required_param(params, "RoleName")?.to_string();
let policy_name = required_param(params, "PolicyName")?.to_string();
let policy = state
.iam
.get_role_policy(role.as_str(), policy_name.as_str())
.context("policy not found")?;
Ok((role, policy_name, policy))
})();
match out {
Ok((role, policy_name, policy)) => iam_xml_response(
"GetRolePolicy",
format!(
"<GetRolePolicyResult><RoleName>{}</RoleName><PolicyName>{}</PolicyName><PolicyDocument>{}</PolicyDocument></GetRolePolicyResult>",
xml_escape(role.as_str()),
xml_escape(policy_name.as_str()),
xml_escape(policy.as_str()),
),
),
Err(err) => xml_error_response(StatusCode::NOT_FOUND, "NoSuchEntity", &err.to_string()),
}
}
"ListRolePolicies" => {
let out = (|| -> Result<(String, Vec<String>)> {
let role = required_param(params, "RoleName")?.to_string();
let names = state
.iam
.list_role_policies(role.as_str())
.context("role not found")?;
Ok((role, names))
})();
match out {
Ok((role, names)) => {
let members = names
.into_iter()
.map(|name| format!("<member>{}</member>", xml_escape(name.as_str())))
.collect::<String>();
iam_xml_response(
"ListRolePolicies",
format!(
"<ListRolePoliciesResult><RoleName>{}</RoleName><PolicyNames>{}</PolicyNames><IsTruncated>false</IsTruncated></ListRolePoliciesResult>",
xml_escape(role.as_str()),
members
),
)
}
Err(err) => {
xml_error_response(StatusCode::NOT_FOUND, "NoSuchEntity", &err.to_string())
}
}
}
"DeleteRolePolicy" => {
let out = (|| -> Result<()> {
let role = required_param(params, "RoleName")?;
let policy_name = required_param(params, "PolicyName")?;
state.iam.delete_role_policy(role, policy_name)?;
let payload = json!({ "policy_name": policy_name }).to_string();
let _ = append_iam_wal_ordered(state, "delete_role_policy", role, payload.as_str());
Ok(())
})();
match out {
Ok(_) => {
iam_xml_response("DeleteRolePolicy", "<DeleteRolePolicyResult/>".to_string())
}
Err(err) => {
xml_error_response(StatusCode::NOT_FOUND, "NoSuchEntity", &err.to_string())
}
}
}
"CreateOpenIDConnectProvider" => {
let out = (|| -> Result<String> {
let url = required_param(params, "Url")?.to_string();
let arn = format!(
"arn:aws:iam:::oidc-provider/{}",
url.trim_start_matches("https://")
.trim_start_matches("http://")
.trim_start_matches("www.")
);
let client_ids = indexed_members(params, "ClientIDList.member.");
let thumbprints = indexed_members(params, "ThumbprintList.member.");
state.iam.create_oidc_provider(
arn.clone(),
url,
client_ids,
if thumbprints.is_empty() {
vec!["test-thumbprint".to_string()]
} else {
thumbprints
},
)?;
let payload = json!({
"url": params.get("Url").cloned().unwrap_or_default(),
"client_ids": indexed_members(params, "ClientIDList.member."),
"thumbprints": indexed_members(params, "ThumbprintList.member."),
})
.to_string();
let _ = append_iam_wal_ordered(
state,
"create_oidc_provider",
arn.as_str(),
payload.as_str(),
);
Ok(arn)
})();
match out {
Ok(arn) => iam_xml_response(
"CreateOpenIDConnectProvider",
format!(
"<CreateOpenIDConnectProviderResult><OpenIDConnectProviderArn>{}</OpenIDConnectProviderArn></CreateOpenIDConnectProviderResult>",
xml_escape(arn.as_str())
),
),
Err(err) => xml_error_response(StatusCode::BAD_REQUEST, "InvalidParameterValue", &err.to_string()),
}
}
"GetOpenIDConnectProvider" => {
let out = (|| -> Result<_> {
let arn = required_param(params, "OpenIDConnectProviderArn")?;
let provider = state
.iam
.get_oidc_provider(arn)
.context("provider not found")?;
Ok(provider)
})();
match out {
Ok(provider) => {
let client_ids = provider
.client_ids
.iter()
.map(|item| format!("<member>{}</member>", xml_escape(item.as_str())))
.collect::<String>();
let thumbprints = provider
.thumbprints
.iter()
.map(|item| format!("<member>{}</member>", xml_escape(item.as_str())))
.collect::<String>();
iam_xml_response(
"GetOpenIDConnectProvider",
format!(
"<GetOpenIDConnectProviderResult><Url>{}</Url><ClientIDList>{}</ClientIDList><ThumbprintList>{}</ThumbprintList></GetOpenIDConnectProviderResult>",
xml_escape(provider.url.as_str()),
client_ids,
thumbprints
),
)
}
Err(err) => {
xml_error_response(StatusCode::NOT_FOUND, "NoSuchEntity", &err.to_string())
}
}
}
"DeleteOpenIDConnectProvider" => {
let out = (|| -> Result<String> {
let arn = required_param(params, "OpenIDConnectProviderArn")?.to_string();
state.iam.delete_oidc_provider(arn.as_str())?;
let _ = append_iam_wal_ordered(state, "delete_oidc_provider", arn.as_str(), "{}");
Ok(arn)
})();
match out {
Ok(_) => iam_xml_response(
"DeleteOpenIDConnectProvider",
"<DeleteOpenIDConnectProviderResult/>".to_string(),
),
Err(err) => {
xml_error_response(StatusCode::NOT_FOUND, "NoSuchEntity", &err.to_string())
}
}
}
"GetSessionToken" => {
let duration_secs = params
.get("DurationSeconds")
.and_then(|value| value.parse::<u64>().ok())
.unwrap_or(3600);
match state.sts.get_session_token(caller.as_str(), duration_secs) {
Ok(creds) => {
if let Some(record) = state.credentials.get(creds.access_key_id.as_str()) {
if let Ok(payload) = serde_json::to_string(&record) {
let _ = append_iam_wal_ordered(
state,
"upsert_credential",
record.principal.as_str(),
payload.as_str(),
);
}
}
sts_xml_response(
"GetSessionToken",
format!(
"<GetSessionTokenResult><Credentials><AccessKeyId>{}</AccessKeyId><SecretAccessKey>{}</SecretAccessKey><SessionToken>{}</SessionToken><Expiration>{}</Expiration></Credentials></GetSessionTokenResult>",
xml_escape(creds.access_key_id.as_str()),
xml_escape(creds.secret_access_key.as_str()),
xml_escape(creds.session_token.as_str()),
xml_escape(creds.expiration.as_str()),
),
)
}
Err(err) => {
xml_error_response(StatusCode::FORBIDDEN, "AccessDenied", &err.to_string())
}
}
}
"AssumeRole" => {
let out = (|| -> Result<_> {
let role_arn = required_param(params, "RoleArn")?;
let duration_secs = params
.get("DurationSeconds")
.and_then(|value| value.parse::<u64>().ok())
.unwrap_or(3600);
state
.sts
.assume_role(caller.as_str(), role_arn, duration_secs)
})();
match out {
Ok(creds) => {
if let Some(record) = state.credentials.get(creds.access_key_id.as_str()) {
if let Ok(payload) = serde_json::to_string(&record) {
let _ = append_iam_wal_ordered(
state,
"upsert_credential",
record.principal.as_str(),
payload.as_str(),
);
}
}
sts_xml_response(
"AssumeRole",
format!(
"<AssumeRoleResult><AssumedRoleUser><Arn>{}</Arn><AssumedRoleId>{}</AssumedRoleId></AssumedRoleUser><Credentials><AccessKeyId>{}</AccessKeyId><SecretAccessKey>{}</SecretAccessKey><SessionToken>{}</SessionToken><Expiration>{}</Expiration></Credentials></AssumeRoleResult>",
xml_escape(caller.as_str()),
xml_escape(caller_access_key.as_str()),
xml_escape(creds.access_key_id.as_str()),
xml_escape(creds.secret_access_key.as_str()),
xml_escape(creds.session_token.as_str()),
xml_escape(creds.expiration.as_str()),
),
)
}
Err(err) => {
xml_error_response(StatusCode::FORBIDDEN, "AccessDenied", &err.to_string())
}
}
}
"AssumeRoleWithWebIdentity" => {
let out = (|| -> Result<_> {
let role_arn = required_param(params, "RoleArn")?;
let token = required_param(params, "WebIdentityToken")?;
let provider = params.get("ProviderId").map(String::as_str);
let duration_secs = params
.get("DurationSeconds")
.and_then(|value| value.parse::<u64>().ok())
.unwrap_or(3600);
state
.sts
.assume_role_with_web_identity(role_arn, duration_secs, token, provider)
})();
match out {
Ok(creds) => {
if let Some(record) = state.credentials.get(creds.access_key_id.as_str()) {
if let Ok(payload) = serde_json::to_string(&record) {
let _ = append_iam_wal_ordered(
state,
"upsert_credential",
record.principal.as_str(),
payload.as_str(),
);
}
}
sts_xml_response(
"AssumeRoleWithWebIdentity",
format!(
"<AssumeRoleWithWebIdentityResult><SubjectFromWebIdentityToken>web-identity</SubjectFromWebIdentityToken><Credentials><AccessKeyId>{}</AccessKeyId><SecretAccessKey>{}</SecretAccessKey><SessionToken>{}</SessionToken><Expiration>{}</Expiration></Credentials></AssumeRoleWithWebIdentityResult>",
xml_escape(creds.access_key_id.as_str()),
xml_escape(creds.secret_access_key.as_str()),
xml_escape(creds.session_token.as_str()),
xml_escape(creds.expiration.as_str()),
),
)
}
Err(err) => {
xml_error_response(StatusCode::FORBIDDEN, "AccessDenied", &err.to_string())
}
}
}
"GetCallerIdentity" => {
let (arn, user_id) = state
.sts
.get_caller_identity(caller.as_str(), caller_access_key.as_str());
let account = if caller_access_key == "IAMROOTKEY123" {
"iam-root"
} else if caller_access_key == "IAMALTROOTKEY123" {
"iam-alt-root"
} else {
"000000000000"
};
sts_xml_response(
"GetCallerIdentity",
format!(
"<GetCallerIdentityResult><UserId>{}</UserId><Account>{}</Account><Arn>{}</Arn></GetCallerIdentityResult>",
xml_escape(user_id.as_str()),
xml_escape(account),
xml_escape(arn.as_str())
),
)
}
other => xml_error_response(
StatusCode::NOT_IMPLEMENTED,
"NotImplemented",
&format!("Query action {other} is not implemented in this release"),
),
}
}
fn required_param<'a>(params: &'a BTreeMap<String, String>, key: &str) -> Result<&'a str> {
params
.get(key)
.map(String::as_str)
.with_context(|| format!("missing parameter {key}"))
}
fn indexed_members(params: &BTreeMap<String, String>, prefix: &str) -> Vec<String> {
let mut rows = params
.iter()
.filter_map(|(k, v)| {
k.strip_prefix(prefix)
.map(|idx| (idx.to_string(), v.clone()))
})
.collect::<Vec<_>>();
rows.sort_by(|a, b| a.0.cmp(&b.0));
rows.into_iter().map(|(_, v)| v).collect()
}
fn parse_query_params(query: &str) -> BTreeMap<String, String> {
url::form_urlencoded::parse(query.as_bytes())
.into_owned()
.collect::<BTreeMap<_, _>>()
}
fn policy_contains_principal(policy_document: &str) -> bool {
let Ok(value) = serde_json::from_str::<serde_json::Value>(policy_document) else {
return false;
};
let Some(statement) = value.get("Statement") else {
return false;
};
if let Some(array) = statement.as_array() {
return array.iter().any(|entry| entry.get("Principal").is_some());
}
statement.get("Principal").is_some()
}
fn extract_access_key(
req: &http::Request<HyperIncoming>,
params: &BTreeMap<String, String>,
) -> Option<String> {
if let Some(value) = params.get("AWSAccessKeyId") {
return Some(value.clone());
}
if let Some(auth) = req
.headers()
.get(http::header::AUTHORIZATION)
.and_then(|value| value.to_str().ok())
{
if let Some(rest) = auth.strip_prefix("AWS ") {
return rest.split(':').next().map(ToString::to_string);
}
if let Some(cred_pos) = auth.find("Credential=") {
let suffix = &auth[cred_pos + "Credential=".len()..];
return suffix.split('/').next().map(ToString::to_string);
}
}
None
}
fn validate_legacy_auth_headers(
headers: &http::HeaderMap,
expected_access_key: &str,
) -> Option<HttpResponse> {
let authorization = headers
.get(http::header::AUTHORIZATION)
.and_then(|v| v.to_str().ok())
.map(str::trim);
let x_amz_date = headers
.get("x-amz-date")
.and_then(|v| v.to_str().ok())
.map(str::trim);
if matches!(authorization, None | Some("")) {
return Some(xml_error_response_close(
StatusCode::FORBIDDEN,
"AccessDenied",
"Access Denied",
));
}
let is_sigv4 = authorization
.map(|value| value.starts_with("AWS4-HMAC-SHA256"))
.unwrap_or(false);
let looks_like_aws2_date = x_amz_date.map(|value| value.contains(',')).unwrap_or(false);
let is_legacy_aws2 = authorization
.map(|value| value.starts_with("AWS "))
.unwrap_or(false)
|| (!is_sigv4 && looks_like_aws2_date);
if !is_legacy_aws2 {
return None;
}
if let Some(auth_header) = authorization.filter(|value| value.starts_with("AWS ")) {
let Some((access_key, signature)) = auth_header.trim_start_matches("AWS ").split_once(':')
else {
return Some(xml_error_response_close(
StatusCode::BAD_REQUEST,
"InvalidArgument",
"Invalid Authorization header",
));
};
if access_key.is_empty() || signature.is_empty() {
return Some(xml_error_response_close(
StatusCode::BAD_REQUEST,
"InvalidArgument",
"Invalid Authorization header",
));
}
if BASE64_STANDARD.decode(signature).is_err() || access_key != expected_access_key {
return Some(xml_error_response_close(
StatusCode::FORBIDDEN,
"InvalidDigest",
"The Content-MD5 you specified is invalid.",
));
}
}
if x_amz_date.is_none() {
return Some(xml_error_response_close(
StatusCode::FORBIDDEN,
"AccessDenied",
"Access Denied",
));
}
if let Some(date_header) = x_amz_date {
match classify_legacy_date(date_header) {
LegacyDateCheck::Ok => {}
LegacyDateCheck::AccessDenied => {
return Some(xml_error_response_close(
StatusCode::FORBIDDEN,
"AccessDenied",
"Access Denied",
));
}
LegacyDateCheck::TooSkewed => {
return Some(xml_error_response_close(
StatusCode::FORBIDDEN,
"RequestTimeTooSkewed",
"The difference between the request time and the server's time is too large.",
));
}
}
}
None
}
enum LegacyDateCheck {
Ok,
AccessDenied,
TooSkewed,
}
fn classify_legacy_date(value: &str) -> LegacyDateCheck {
let trimmed = value.trim();
if trimmed.is_empty() {
return LegacyDateCheck::AccessDenied;
}
let Some(year) = extract_http_date_year(trimmed) else {
return LegacyDateCheck::AccessDenied;
};
if year < 1970 {
return LegacyDateCheck::AccessDenied;
}
let now = SystemTime::now();
let current_year = {
let now_http = httpdate::fmt_http_date(now);
extract_http_date_year(&now_http).unwrap_or(1970)
};
if year < current_year.saturating_sub(1) || year > current_year.saturating_add(1) {
return LegacyDateCheck::TooSkewed;
}
let parsed = match httpdate::parse_http_date(trimmed) {
Ok(value) => value,
Err(_) => return LegacyDateCheck::AccessDenied,
};
let skew = match now.duration_since(parsed) {
Ok(delta) => delta,
Err(_) => parsed.duration_since(now).unwrap_or_default(),
};
if skew > Duration::from_secs(15 * 60) {
LegacyDateCheck::TooSkewed
} else {
LegacyDateCheck::Ok
}
}
fn extract_http_date_year(value: &str) -> Option<u32> {
let bytes = value.as_bytes();
let mut idx = 0usize;
while idx + 4 <= bytes.len() {
if bytes[idx..idx + 4].iter().all(u8::is_ascii_digit) {
let year = std::str::from_utf8(&bytes[idx..idx + 4]).ok()?;
let parsed = year.parse::<u32>().ok()?;
if (1000..=9999).contains(&parsed) {
return Some(parsed);
}
}
idx += 1;
}
None
}
#[cfg(test)]
fn query_param_value(query: &str, key: &str) -> Option<String> {
parse_query_params(query).get(key).cloned()
}
fn xml_error_response(status: StatusCode, code: &str, message: &str) -> HttpResponse {
let body = format!(
"<?xml version=\"1.0\" encoding=\"UTF-8\"?>\
<Error><Code>{code}</Code><Message>{message}</Message></Error>"
);
http::Response::builder()
.status(status)
.header(http::header::CONTENT_TYPE, "application/xml")
.body(Body::from(body))
.unwrap_or_else(|_| http::Response::new(Body::from(Vec::new())))
}
fn xml_error_response_close(status: StatusCode, code: &str, message: &str) -> HttpResponse {
with_connection_close(xml_error_response(status, code, message))
}
fn close_connection_on_error(response: HttpResponse) -> HttpResponse {
if response.status().is_success() {
return response;
}
with_connection_close(response)
}
fn with_connection_close(mut response: HttpResponse) -> HttpResponse {
response.headers_mut().insert(
http::header::CONNECTION,
http::HeaderValue::from_static("close"),
);
response
}
fn json_response(status: StatusCode, value: serde_json::Value) -> HttpResponse {
http::Response::builder()
.status(status)
.header(http::header::CONTENT_TYPE, "application/json")
.body(Body::from(value.to_string()))
.unwrap_or_else(|_| http::Response::new(Body::from(Vec::new())))
}
fn iam_xml_response(action: &str, inner: String) -> HttpResponse {
let request_id = format!("req-{}", Uuid::new_v4().simple());
let body = format!(
"<?xml version=\"1.0\" encoding=\"UTF-8\"?>\
<{action}Response xmlns=\"https://iam.amazonaws.com/doc/2010-05-08/\">\
{inner}\
<ResponseMetadata><RequestId>{request_id}</RequestId></ResponseMetadata>\
</{action}Response>"
);
http::Response::builder()
.status(StatusCode::OK)
.header(http::header::CONTENT_TYPE, "application/xml")
.body(Body::from(body))
.unwrap_or_else(|_| http::Response::new(Body::from(Vec::new())))
}
fn sts_xml_response(action: &str, inner: String) -> HttpResponse {
let request_id = format!("req-{}", Uuid::new_v4().simple());
let body = format!(
"<?xml version=\"1.0\" encoding=\"UTF-8\"?>\
<{action}Response xmlns=\"https://sts.amazonaws.com/doc/2011-06-15/\">\
{inner}\
<ResponseMetadata><RequestId>{request_id}</RequestId></ResponseMetadata>\
</{action}Response>"
);
http::Response::builder()
.status(StatusCode::OK)
.header(http::header::CONTENT_TYPE, "application/xml")
.body(Body::from(body))
.unwrap_or_else(|_| http::Response::new(Body::from(Vec::new())))
}
fn xml_escape(value: &str) -> String {
value
.replace('&', "&")
.replace('<', "<")
.replace('>', ">")
.replace('\"', """)
.replace('\'', "'")
}
fn next_nonce_seed(state: &GatewayState) -> u64 {
state.nonce_counter.fetch_add(1, Ordering::AcqRel)
}
fn append_object_wal_ordered(state: &GatewayState, mut entry: WalEntryV1) -> S3Result<u64> {
let _guard = state
.wal_sequence_lock
.lock()
.map_err(|_| s3_error!(InternalError, "WAL ordering lock poisoned"))?;
let seq = state.wal.next_seq();
entry.seq = seq;
state
.wal
.append_object_entry(&entry)
.map_err(|_| s3_error!(InternalError, "WAL append failed"))?;
Ok(seq)
}
fn append_multipart_wal_ordered(
state: &GatewayState,
mut entry: MultipartWalEntryV1,
) -> S3Result<u64> {
let _guard = state
.wal_sequence_lock
.lock()
.map_err(|_| s3_error!(InternalError, "WAL ordering lock poisoned"))?;
let seq = state.wal.next_seq();
entry.seq = seq;
state
.wal
.append_multipart_entry(&entry)
.map_err(|_| s3_error!(InternalError, "WAL append failed"))?;
Ok(seq)
}
fn append_iam_wal_ordered(
state: &GatewayState,
op: &str,
target: &str,
payload: &str,
) -> S3Result<u64> {
let _guard = state
.wal_sequence_lock
.lock()
.map_err(|_| s3_error!(InternalError, "WAL ordering lock poisoned"))?;
let seq = state.wal.next_seq();
state
.wal
.append_iam_entry(&wal::IamWalEntryV1 {
seq,
op: op.to_string(),
target: target.to_string(),
payload: payload.to_string(),
ts_unix_ms: now_unix_ms(),
})
.map_err(|_| s3_error!(InternalError, "IAM WAL append failed"))?;
Ok(seq)
}
fn append_webhook_wal_ordered(
state: &GatewayState,
op: &str,
bucket: &str,
rule_id: &str,
rule_json: &str,
) -> S3Result<u64> {
let _guard = state
.wal_sequence_lock
.lock()
.map_err(|_| s3_error!(InternalError, "WAL ordering lock poisoned"))?;
let seq = state.wal.next_seq();
state
.wal
.append_webhook_entry(&wal::WebhookWalEntryV1 {
seq,
op: op.to_string(),
bucket: bucket.to_string(),
rule_id: rule_id.to_string(),
rule_json: rule_json.to_string(),
ts_unix_ms: now_unix_ms(),
})
.map_err(|_| s3_error!(InternalError, "Webhook WAL append failed"))?;
Ok(seq)
}
fn append_webhook_event_wal_ordered(
state: &GatewayState,
op: &str,
event_id: &str,
rule_id: &str,
bucket: &str,
key: &str,
payload_json: &str,
) -> S3Result<u64> {
let _guard = state
.wal_sequence_lock
.lock()
.map_err(|_| s3_error!(InternalError, "WAL ordering lock poisoned"))?;
let seq = state.wal.next_seq();
state
.wal
.append_webhook_event_entry(&wal::WebhookEventWalEntryV1 {
seq,
op: op.to_string(),
event_id: event_id.to_string(),
rule_id: rule_id.to_string(),
bucket: bucket.to_string(),
key: key.to_string(),
payload_json: payload_json.to_string(),
ts_unix_ms: now_unix_ms(),
})
.map_err(|_| s3_error!(InternalError, "Webhook event WAL append failed"))?;
Ok(seq)
}
fn append_migration_wal_ordered(
state: &GatewayState,
op: &str,
bucket: &str,
key: &str,
payload: &str,
) -> S3Result<u64> {
let _guard = state
.wal_sequence_lock
.lock()
.map_err(|_| s3_error!(InternalError, "WAL ordering lock poisoned"))?;
let seq = state.wal.next_seq();
state
.wal
.append_migration_entry(&wal::MigrationWalEntryV1 {
seq,
op: op.to_string(),
bucket: bucket.to_string(),
key: key.to_string(),
payload: payload.to_string(),
ts_unix_ms: now_unix_ms(),
})
.map_err(|_| s3_error!(InternalError, "Migration WAL append failed"))?;
Ok(seq)
}
fn enqueue_webhook_intents(
state: &GatewayState,
bucket: &str,
key: &str,
size_bytes: u64,
etag: &str,
is_update: bool,
) -> S3Result<()> {
let tasks = state
.webhooks
.matching_object_tasks(bucket, key, size_bytes, etag, is_update);
for task in tasks {
let payload = serde_json::to_string(&task.event).unwrap_or_else(|_| "{}".to_string());
append_webhook_event_wal_ordered(
state,
"intent",
task.event.event_id.as_str(),
task.rule.id.as_str(),
bucket,
key,
payload.as_str(),
)?;
let _ = state.webhooks.enqueue_task(task);
}
Ok(())
}
fn replay_control_state(
iam: &Arc<IamState>,
credentials: &Arc<CredentialStore>,
webhooks: &Arc<WebhookEngine>,
migration: &Arc<MigrationState>,
iam_log: &[wal::IamWalEntryV1],
webhook_log: &[wal::WebhookWalEntryV1],
webhook_event_log: &[wal::WebhookEventWalEntryV1],
migration_log: &[wal::MigrationWalEntryV1],
) -> Result<()> {
for entry in iam_log {
match entry.op.as_str() {
"create_user" => {
let _ = iam.create_user(entry.target.as_str());
}
"delete_user" => {
let _ = iam.delete_user(entry.target.as_str());
}
"put_user_policy" => {
if let Ok(payload) =
serde_json::from_str::<serde_json::Value>(entry.payload.as_str())
{
if let (Some(policy_name), Some(policy_document)) = (
payload.get("policy_name").and_then(|v| v.as_str()),
payload.get("policy_document").and_then(|v| v.as_str()),
) {
let _ = iam.put_user_policy(
entry.target.as_str(),
policy_name,
policy_document,
);
}
}
}
"delete_user_policy" => {
if let Ok(payload) =
serde_json::from_str::<serde_json::Value>(entry.payload.as_str())
{
if let Some(policy_name) = payload.get("policy_name").and_then(|v| v.as_str()) {
let _ = iam.delete_user_policy(entry.target.as_str(), policy_name);
}
}
}
"create_role" => {
if let Ok(payload) =
serde_json::from_str::<serde_json::Value>(entry.payload.as_str())
{
let assume = payload
.get("assume_role_policy")
.and_then(|v| v.as_str())
.unwrap_or("{}");
let _ = iam.create_role(entry.target.as_str(), assume);
}
}
"put_role_policy" => {
if let Ok(payload) =
serde_json::from_str::<serde_json::Value>(entry.payload.as_str())
{
if let (Some(policy_name), Some(policy_document)) = (
payload.get("policy_name").and_then(|v| v.as_str()),
payload.get("policy_document").and_then(|v| v.as_str()),
) {
let _ = iam.put_role_policy(
entry.target.as_str(),
policy_name,
policy_document,
);
}
}
}
"delete_role_policy" => {
if let Ok(payload) =
serde_json::from_str::<serde_json::Value>(entry.payload.as_str())
{
if let Some(policy_name) = payload.get("policy_name").and_then(|v| v.as_str()) {
let _ = iam.delete_role_policy(entry.target.as_str(), policy_name);
}
}
}
"create_oidc_provider" => {
if let Ok(payload) =
serde_json::from_str::<serde_json::Value>(entry.payload.as_str())
{
let url = payload
.get("url")
.and_then(|v| v.as_str())
.unwrap_or_default()
.to_string();
let client_ids = payload
.get("client_ids")
.and_then(|v| serde_json::from_value::<Vec<String>>(v.clone()).ok())
.unwrap_or_default();
let thumbprints = payload
.get("thumbprints")
.and_then(|v| serde_json::from_value::<Vec<String>>(v.clone()).ok())
.unwrap_or_default();
let _ = iam.create_oidc_provider(
entry.target.clone(),
url,
client_ids,
thumbprints,
);
}
}
"delete_oidc_provider" => {
let _ = iam.delete_oidc_provider(entry.target.as_str());
}
"upsert_credential" => {
if let Ok(record) = serde_json::from_str::<crate::auth_provider::CredentialRecord>(
entry.payload.as_str(),
) {
credentials.upsert(record);
}
}
_ => {}
}
}
for entry in webhook_log {
match entry.op.as_str() {
"replace_bucket_rules" => {
if let Ok(rules) =
serde_json::from_str::<Vec<WebhookRule>>(entry.rule_json.as_str())
{
let _ = webhooks.replace_rules(entry.bucket.as_str(), rules);
}
}
"delete_bucket_rules" => {
webhooks.delete_rule(entry.bucket.as_str(), None);
}
"delete_rule" => {
webhooks.delete_rule(entry.bucket.as_str(), Some(entry.rule_id.as_str()));
}
_ => {}
}
}
let mut pending_webhooks = BTreeMap::<String, wal::WebhookEventWalEntryV1>::new();
for entry in webhook_event_log {
let key = format!("{}::{}", entry.event_id, entry.rule_id);
match entry.op.as_str() {
"intent" => {
pending_webhooks.insert(key, entry.clone());
}
"complete" => {
pending_webhooks.remove(&key);
}
_ => {}
}
}
for (_k, entry) in pending_webhooks {
let event = serde_json::from_str::<crate::webhooks::WebhookEvent>(entry.payload_json.as_str())
.unwrap_or(crate::webhooks::WebhookEvent {
event_id: entry.event_id.clone(),
timestamp: chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Millis, true),
event_name: "s3:ObjectCreated:Put".to_string(),
bucket: entry.bucket.clone(),
object: crate::webhooks::WebhookObject {
key: entry.key.clone(),
size_bytes: 0,
etag: String::new(),
},
});
let maybe_rule = webhooks
.get_rules(entry.bucket.as_str())
.into_iter()
.find(|rule| rule.id == entry.rule_id);
if let Some(rule) = maybe_rule {
let _ = webhooks.enqueue_task(WebhookDeliveryTask { rule, event });
}
}
for entry in migration_log {
if entry.op == "set_bucket_policy" {
if let Ok(policy) =
serde_json::from_str::<MigrationBucketPolicy>(entry.payload.as_str())
{
migration.set_bucket_policy(entry.bucket.as_str(), policy);
}
}
}
Ok(())
}
fn now_unix_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|dur| dur.as_millis() as u64)
.unwrap_or(0)
}
fn unix_ms_to_system_time(ts_unix_ms: u64) -> SystemTime {
UNIX_EPOCH + Duration::from_millis(ts_unix_ms)
}
#[cfg(test)]
mod tests {
use super::*;
use base64::Engine;
#[test]
fn pipeline_blob_round_trip() {
let payload = vec![3_u8; 2 * 1024 * 1024 + 123];
let output = process_payload_crc_zstd_aes(&payload, 1024 * 1024, &PIPELINE_KEY, 11)
.expect("pipeline processing should succeed");
let blob = encode_pipeline_blob(&output.chunks).expect("encode should succeed");
let restored = decode_pipeline_blob(&blob).expect("decode should succeed");
assert_eq!(restored, payload);
}
#[test]
fn worker_core_selection_caps_to_64() {
let topology = TopologyReport {
cpu_vendor: "AuthenticAMD".to_string(),
family_model: "0x19:0x1".to_string(),
cpu_generation: nexus_core::topology::CpuGeneration::Zen3,
numa_nodes: Vec::new(),
ccx_groups: Vec::new(),
rss_irq_map: Vec::new(),
selected_cores: (0..96).collect(),
};
let selected = select_worker_cores(&topology, 0);
assert_eq!(selected.len(), 64);
}
#[test]
fn payload_header_validation_accepts_valid_md5_and_length() {
let payload = b"abc123";
let md5_b64 = BASE64_STANDARD.encode(md5::compute(payload).0);
let out =
validate_payload_headers(Some(md5_b64.as_str()), Some(payload.len() as i64), payload);
assert!(out.is_ok());
}
#[test]
fn payload_header_validation_rejects_invalid_md5() {
let payload = b"abc123";
let out = validate_payload_headers(Some("!!!"), Some(payload.len() as i64), payload);
assert!(out.is_err());
}
#[test]
fn payload_header_validation_rejects_length_mismatch() {
let payload = b"abc123";
let out = validate_payload_headers(None, Some(1), payload);
assert!(out.is_err());
}
#[test]
fn normalize_max_uploads_bounds() {
assert_eq!(normalize_max_uploads(Some(10)).expect("valid"), 10);
assert_eq!(normalize_max_uploads(Some(5000)).expect("clamped"), 1000);
assert!(normalize_max_uploads(Some(0)).is_err());
}
#[test]
fn query_param_value_extracts_action() {
assert_eq!(
query_param_value("Action=AssumeRole&Version=2011-06-15", "Action").as_deref(),
Some("AssumeRole")
);
}
#[test]
fn legacy_auth_denials_force_connection_close() {
let headers = http::HeaderMap::new();
let response =
validate_legacy_auth_headers(&headers, DEFAULT_ACCESS_KEY).expect("must reject");
assert_eq!(
response
.headers()
.get(http::header::CONNECTION)
.and_then(|value| value.to_str().ok()),
Some("close")
);
}
#[test]
fn close_connection_only_for_error_responses() {
let ok = http::Response::builder()
.status(StatusCode::OK)
.body(Body::from(Vec::<u8>::new()))
.expect("ok response");
let ok = close_connection_on_error(ok);
assert!(ok.headers().get(http::header::CONNECTION).is_none());
let err = xml_error_response(StatusCode::FORBIDDEN, "AccessDenied", "Access Denied");
let err = close_connection_on_error(err);
assert_eq!(
err.headers()
.get(http::header::CONNECTION)
.and_then(|value| value.to_str().ok()),
Some("close")
);
}
#[test]
fn panopticon_unknown_route_falls_back_to_index() {
let response = panopticon_file_response("/_nexus/ui/policies");
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(
response
.headers()
.get(http::header::CONTENT_TYPE)
.and_then(|value| value.to_str().ok()),
Some("text/html; charset=utf-8")
);
}
#[test]
fn panopticon_js_asset_uses_javascript_content_type() {
let response = panopticon_file_response("/_nexus/ui/assets/app.js");
assert!(
matches!(
response.status(),
StatusCode::OK | StatusCode::NOT_FOUND
),
"asset route should resolve deterministically"
);
if response.status() == StatusCode::OK {
assert_eq!(
response
.headers()
.get(http::header::CONTENT_TYPE)
.and_then(|value| value.to_str().ok()),
Some("application/javascript; charset=utf-8")
);
}
}
#[test]
fn browser_prefix_normalizes_and_builds_breadcrumbs() {
assert_eq!(normalize_browser_prefix("images/2026"), "images/2026/");
let breadcrumbs = browser_breadcrumbs("images/2026/");
assert_eq!(breadcrumbs.len(), 3);
assert_eq!(breadcrumbs[0].label, "Root");
assert_eq!(breadcrumbs[1].label, "images");
assert_eq!(breadcrumbs[2].prefix, "images/2026/");
}
#[test]
fn browser_content_type_prefers_metadata_then_extension() {
let mut metadata = HashMap::new();
metadata.insert("content-type".to_string(), "application/x-custom".to_string());
assert_eq!(
browser_content_type_from_metadata(&metadata, "foo.txt").as_deref(),
Some("application/x-custom")
);
let empty = HashMap::new();
assert_eq!(
browser_content_type_from_metadata(&empty, "notes.json").as_deref(),
Some("application/json")
);
}
#[test]
fn browser_preview_kind_classifies_common_types() {
assert_eq!(browser_preview_kind("hero.jpg", Some("image/jpeg")), "image");
assert_eq!(browser_preview_kind("spec.pdf", Some("application/pdf")), "pdf");
assert_eq!(browser_preview_kind("notes.md", Some("text/markdown")), "text");
assert_eq!(
browser_preview_kind("archive.bin", Some("application/octet-stream")),
"none"
);
}
}