use std::collections::{BTreeMap, BTreeSet};
use std::fs::{self, File, OpenOptions};
use std::io;
use std::io::ErrorKind;
use std::net::SocketAddr;
use std::os::unix::fs::{FileExt, OpenOptionsExt};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{mpsc, Arc, RwLock};
use std::thread;
use std::time::{Duration, Instant, SystemTime};
use anyhow::{Context, Result};
use arc_swap::ArcSwap;
use crossbeam_skiplist::SkipMap;
use fst::{IntoStreamer, Streamer};
use futures::TryStreamExt;
use hyper_util::rt::{TokioExecutor, TokioIo};
use hyper_util::server::conn::auto::Builder as ConnBuilder;
#[cfg(target_os = "linux")]
use nix::libc;
use s3s::auth::SimpleAuth;
use s3s::dto::*;
use s3s::service::S3ServiceBuilder;
use s3s::{s3_error, Body, S3Request, S3Response, S3Result};
use socket2::{Domain, Protocol, Socket, Type};
use crate::module_d;
const ALIGNMENT: usize = 4096;
const RECORD_HEADER_BYTES: usize = 12;
#[derive(Debug, Clone)]
pub struct ModuleFConfig {
pub listen_addr: String,
pub workers: usize,
pub target_path: PathBuf,
pub group_bytes: usize,
pub compaction_threshold: usize,
pub flush_interval_ms: u64,
pub default_bucket: String,
pub access_key: String,
pub secret_key: String,
pub allow_file_fallback: bool,
pub require_io_uring: bool,
}
#[derive(Debug, Clone)]
struct StoredObjectMeta {
payload_offset: u64,
stored_len: u32,
original_len: u32,
crc32c: u32,
last_modified: Timestamp,
}
#[derive(Debug)]
struct PendingRecord {
record_start: usize,
stored_len: u32,
original_len: u32,
crc32c: u32,
response_tx: tokio::sync::oneshot::Sender<Result<StoredObjectMeta>>,
}
#[derive(Debug)]
struct AppendRequest {
payload: Vec<u8>,
original_len: u32,
crc32c: u32,
response_tx: tokio::sync::oneshot::Sender<Result<StoredObjectMeta>>,
}
#[derive(Clone)]
struct GroupCommitAppender {
submit_tx: mpsc::Sender<AppendRequest>,
read_path: PathBuf,
mode: String,
}
impl GroupCommitAppender {
fn start(config: &ModuleFConfig) -> Result<Self> {
if config.require_io_uring && !module_d::io_uring_available() {
anyhow::bail!("module-f requires io_uring, but probe failed in current environment");
}
let (file, mode, resolved_path) =
open_target_file(&config.target_path, config.allow_file_fallback).with_context(
|| {
format!(
"failed to open module-f target {}",
config.target_path.display()
)
},
)?;
let (submit_tx, submit_rx) = mpsc::channel::<AppendRequest>();
let group_bytes = config.group_bytes;
let flush_interval = Duration::from_millis(config.flush_interval_ms.max(1));
thread::Builder::new()
.name("module-f-group-commit".to_string())
.spawn(move || {
if let Err(err) =
run_group_commit_loop(file, submit_rx, group_bytes, flush_interval)
{
eprintln!("module-f appender thread exiting with error: {err:#}");
}
})
.context("failed spawning module-f group-commit thread")?;
Ok(Self {
submit_tx,
read_path: resolved_path,
mode,
})
}
async fn append_payload(
&self,
payload: Vec<u8>,
original_len: u32,
crc32c: u32,
) -> Result<StoredObjectMeta> {
let (response_tx, response_rx) = tokio::sync::oneshot::channel();
self.submit_tx
.send(AppendRequest {
payload,
original_len,
crc32c,
response_tx,
})
.map_err(|_| anyhow::anyhow!("module-f appender queue closed"))?;
response_rx
.await
.context("module-f appender response channel dropped")?
}
fn read_payload(&self, offset: u64, len: usize) -> Result<Vec<u8>> {
let file = File::open(&self.read_path)
.with_context(|| format!("failed opening read path {}", self.read_path.display()))?;
let mut buffer = vec![0_u8; len];
let mut filled = 0usize;
while filled < len {
let n = file
.read_at(&mut buffer[filled..], offset + filled as u64)
.context("module-f read_at failed")?;
if n == 0 {
anyhow::bail!(
"unexpected EOF while reading payload at offset {} (wanted {}, got {})",
offset,
len,
filled
);
}
filled += n;
}
Ok(buffer)
}
}
#[derive(Debug)]
struct ImmutableIndex {
entries: BTreeMap<String, StoredObjectMeta>,
fst: fst::Map<Vec<u8>>,
}
impl ImmutableIndex {
fn empty() -> Self {
let builder = fst::MapBuilder::memory();
let bytes = builder.into_inner().unwrap_or_default();
let fst = fst::Map::new(bytes).unwrap();
Self {
entries: BTreeMap::new(),
fst,
}
}
}
#[derive(Debug)]
struct CompactionJob {
delta: Vec<(String, Option<StoredObjectMeta>)>,
}
#[derive(Clone)]
struct LsmIndex {
l0: Arc<SkipMap<String, Option<StoredObjectMeta>>>,
l1: Arc<ArcSwap<ImmutableIndex>>,
compaction_threshold: usize,
compaction_in_progress: Arc<AtomicBool>,
compaction_count: Arc<AtomicU64>,
compaction_total_ns: Arc<AtomicU64>,
compaction_tx: mpsc::Sender<CompactionJob>,
}
impl LsmIndex {
fn new(compaction_threshold: usize) -> Result<Self> {
let l0 = Arc::new(SkipMap::<String, Option<StoredObjectMeta>>::new());
let l1 = Arc::new(ArcSwap::from_pointee(ImmutableIndex::empty()));
let compaction_in_progress = Arc::new(AtomicBool::new(false));
let compaction_count = Arc::new(AtomicU64::new(0));
let compaction_total_ns = Arc::new(AtomicU64::new(0));
let (compaction_tx, compaction_rx) = mpsc::channel::<CompactionJob>();
let l1_for_worker = l1.clone();
let compaction_in_progress_worker = compaction_in_progress.clone();
let compaction_count_worker = compaction_count.clone();
let compaction_total_ns_worker = compaction_total_ns.clone();
thread::Builder::new()
.name("module-f-compaction".to_string())
.spawn(move || {
for job in compaction_rx {
let started = Instant::now();
match build_compacted_index(l1_for_worker.load_full().as_ref(), job.delta) {
Ok(new_index) => {
l1_for_worker.store(Arc::new(new_index));
compaction_count_worker.fetch_add(1, Ordering::Relaxed);
compaction_total_ns_worker
.fetch_add(started.elapsed().as_nanos() as u64, Ordering::Relaxed);
}
Err(err) => {
eprintln!("module-f compaction error: {err:#}");
}
}
compaction_in_progress_worker.store(false, Ordering::Release);
}
})
.context("failed spawning module-f compaction thread")?;
Ok(Self {
l0,
l1,
compaction_threshold,
compaction_in_progress,
compaction_count,
compaction_total_ns,
compaction_tx,
})
}
fn upsert(&self, key: String, meta: StoredObjectMeta) {
self.l0.insert(key, Some(meta));
self.maybe_schedule_compaction();
}
fn delete(&self, key: String) {
self.l0.insert(key, None);
self.maybe_schedule_compaction();
}
fn get(&self, key: &str) -> Option<StoredObjectMeta> {
if let Some(entry) = self.l0.get(key) {
return entry.value().clone();
}
self.l1.load().entries.get(key).cloned()
}
fn list_prefix(&self, prefix: &str, max_keys: usize) -> Vec<(String, StoredObjectMeta)> {
let mut merged = BTreeMap::<String, StoredObjectMeta>::new();
let l1 = self.l1.load_full();
let mut range = l1.fst.range().ge(prefix.as_bytes());
if let Some(upper) = prefix_upper_bound(prefix.as_bytes()) {
range = range.lt(upper.as_slice());
}
let mut stream = range.into_stream();
while let Some((key_bytes, _)) = stream.next() {
let key = String::from_utf8_lossy(key_bytes).to_string();
if let Some(meta) = l1.entries.get(&key) {
merged.insert(key, meta.clone());
}
}
for entry in self.l0.iter() {
let key = entry.key();
if !key.starts_with(prefix) {
continue;
}
match entry.value().clone() {
Some(meta) => {
merged.insert(key.clone(), meta);
}
None => {
merged.remove(key);
}
}
}
merged.into_iter().take(max_keys).collect()
}
fn compaction_metrics(&self) -> (u64, f64) {
let count = self.compaction_count.load(Ordering::Relaxed);
if count == 0 {
return (0, 0.0);
}
let total_ns = self.compaction_total_ns.load(Ordering::Relaxed);
let avg_ms = (total_ns as f64 / count as f64) / 1_000_000.0;
(count, avg_ms)
}
fn maybe_schedule_compaction(&self) {
if self.l0.len() < self.compaction_threshold {
return;
}
if self
.compaction_in_progress
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
return;
}
let mut delta = Vec::<(String, Option<StoredObjectMeta>)>::new();
for entry in self.l0.iter() {
delta.push((entry.key().clone(), entry.value().clone()));
}
if delta.is_empty() {
self.compaction_in_progress.store(false, Ordering::Release);
return;
}
for (key, _) in &delta {
let _ = self.l0.remove(key);
}
if self.compaction_tx.send(CompactionJob { delta }).is_err() {
self.compaction_in_progress.store(false, Ordering::Release);
}
}
}
fn build_compacted_index(
previous: &ImmutableIndex,
delta: Vec<(String, Option<StoredObjectMeta>)>,
) -> Result<ImmutableIndex> {
let mut merged = previous.entries.clone();
for (key, value) in delta {
match value {
Some(meta) => {
merged.insert(key, meta);
}
None => {
merged.remove(&key);
}
}
}
let mut builder = fst::MapBuilder::memory();
for (ordinal, key) in merged.keys().enumerate() {
builder
.insert(key, ordinal as u64)
.with_context(|| format!("failed inserting key into compacted fst: {key}"))?;
}
let bytes = builder
.into_inner()
.context("failed finalizing compacted fst bytes")?;
let fst = fst::Map::new(bytes).context("failed loading compacted fst map")?;
Ok(ImmutableIndex {
entries: merged,
fst,
})
}
#[derive(Clone)]
struct GatewayState {
appender: GroupCommitAppender,
index: LsmIndex,
buckets: Arc<RwLock<BTreeSet<String>>>,
default_bucket: String,
request_count: Arc<AtomicU64>,
puts: Arc<AtomicU64>,
gets: Arc<AtomicU64>,
lists: Arc<AtomicU64>,
deletes: Arc<AtomicU64>,
last_compaction_count: Arc<AtomicU64>,
last_compaction_avg_ms_x1000: 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;
let mut buckets = self
.state
.buckets
.write()
.map_err(|_| s3_error!(InternalError))?;
buckets.insert(bucket);
Ok(S3Response::new(CreateBucketOutput::default()))
}
async fn head_bucket(
&self,
req: S3Request<HeadBucketInput>,
) -> S3Result<S3Response<HeadBucketOutput>> {
let bucket = req.input.bucket;
let buckets = self
.state
.buckets
.read()
.map_err(|_| s3_error!(InternalError))?;
if !buckets.contains(&bucket) {
return Err(s3_error!(NoSuchBucket));
}
Ok(S3Response::new(HeadBucketOutput::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()),
..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 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 input = req.input;
let bucket = input.bucket;
let key = input.key;
ensure_bucket_exists(&self.state, &bucket)?;
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"))?;
let original_len = u32::try_from(raw.len())
.map_err(|_| s3_error!(InvalidRequest, "object body too large for tracer bullet"))?;
let crc = crc32c::crc32c(&raw);
let compressed = zstd::bulk::compress(&raw, 1)
.map_err(|_| s3_error!(InternalError, "zstd compression failed"))?;
let stored = self
.state
.appender
.append_payload(compressed, original_len, crc)
.await
.map_err(|_| s3_error!(InternalError, "append path failed"))?;
self.state
.index
.upsert(storage_key(&bucket, &key), stored.clone());
let (compactions, avg_ms) = self.state.index.compaction_metrics();
self.state
.last_compaction_count
.store(compactions, Ordering::Relaxed);
self.state
.last_compaction_avg_ms_x1000
.store((avg_ms * 1000.0) as u64, Ordering::Relaxed);
Ok(S3Response::new(PutObjectOutput {
e_tag: Some(ETag::Strong(format!("{:08x}", stored.crc32c))),
..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;
let key = input.key;
let Some(meta) = self.state.index.get(&storage_key(&bucket, &key)) else {
return Err(s3_error!(NoSuchKey));
};
let payload = self
.state
.appender
.read_payload(meta.payload_offset, meta.stored_len as usize)
.map_err(|_| s3_error!(InternalError, "read path failed"))?;
let mut restored = zstd::bulk::decompress(&payload, meta.original_len as usize)
.map_err(|_| s3_error!(InternalError, "zstd decompression failed"))?;
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.original_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(format!("{:08x}", meta.crc32c))),
last_modified: Some(meta.last_modified.clone()),
..Default::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 = input.max_keys.unwrap_or(1000).max(1);
let max_keys = usize::try_from(max_keys_i32).map_err(|_| s3_error!(InvalidArgument))?;
let storage_prefix = storage_key_prefix(&input.bucket, &prefix);
let mut rows = self
.state
.index
.list_prefix(&storage_prefix, max_keys.saturating_add(1));
if let Some(token) = input
.continuation_token
.clone()
.or(input.start_after.clone())
{
let token_key = storage_key(&input.bucket, &token);
rows.retain(|(full_key, _)| full_key > &token_key);
}
let is_truncated = rows.len() > max_keys;
rows.truncate(max_keys);
let mut contents = 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 internal storage key"))?;
contents.push(Object {
key: Some(key),
size: Some(i64::from(meta.original_len)),
e_tag: Some(ETag::Strong(format!("{:08x}", meta.crc32c))),
last_modified: Some(meta.last_modified.clone()),
..Default::default()
});
}
let key_count = i32::try_from(contents.len()).map_err(|_| s3_error!(InternalError))?;
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),
..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 Some(meta) = self
.state
.index
.get(&storage_key(&input.bucket, &input.key))
else {
return Err(s3_error!(NoSuchKey));
};
Ok(S3Response::new(HeadObjectOutput {
content_length: Some(i64::from(meta.original_len)),
e_tag: Some(ETag::Strong(format!("{:08x}", meta.crc32c))),
last_modified: Some(meta.last_modified.clone()),
..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;
self.state
.index
.delete(storage_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);
self.state
.index
.delete(storage_key(&input.bucket, &object.key));
deleted.push(DeletedObject {
key: Some(object.key),
..Default::default()
});
}
Ok(S3Response::new(DeleteObjectsOutput {
deleted: Some(deleted),
..Default::default()
}))
}
}
pub fn serve(config: ModuleFConfig) -> Result<()> {
validate_config(&config)?;
let listen_addr: SocketAddr = config
.listen_addr
.parse()
.with_context(|| format!("invalid listen_addr {}", config.listen_addr))?;
let appender = GroupCommitAppender::start(&config)?;
let index = LsmIndex::new(config.compaction_threshold)?;
let buckets = Arc::new(RwLock::new(BTreeSet::<String>::new()));
{
let mut guard = buckets
.write()
.map_err(|_| anyhow::anyhow!("bucket lock poisoned"))?;
guard.insert(config.default_bucket.clone());
}
let state = GatewayState {
appender: appender.clone(),
index,
buckets,
default_bucket: config.default_bucket.clone(),
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)),
last_compaction_count: Arc::new(AtomicU64::new(0)),
last_compaction_avg_ms_x1000: Arc::new(AtomicU64::new(0)),
};
let core_ids = core_affinity::get_core_ids().unwrap_or_default();
let worker_count = config.workers;
let access_key = config.access_key.clone();
let secret_key = config.secret_key.clone();
let accept_counts = Arc::new(
(0..worker_count)
.map(|_| AtomicU64::new(0))
.collect::<Vec<_>>(),
);
eprintln!(
"module-f starting: addr={} workers={} mode={} target={} default_bucket={} require_io_uring={}",
listen_addr,
worker_count,
appender.mode,
appender.read_path.display(),
state.default_bucket,
config.require_io_uring
);
let mut worker_handles = Vec::with_capacity(worker_count);
for worker_idx in 0..worker_count {
let state = state.clone();
let accept_counts = accept_counts.clone();
let access_key = access_key.clone();
let secret_key = secret_key.clone();
let worker_core = if core_ids.is_empty() {
None
} else {
Some(core_ids[worker_idx % core_ids.len()])
};
let handle = thread::Builder::new()
.name(format!("module-f-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 {} failed to bind listener", worker_idx))?;
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 };
let s3_service = {
let mut builder = S3ServiceBuilder::new(gateway);
builder.set_auth(SimpleAuth::from_single(
access_key.as_str(),
secret_key.as_str(),
));
builder.build()
};
loop {
let (stream, _remote) = listener
.accept()
.await
.context("accept failed in module-f worker")?;
accept_counts[worker_idx].fetch_add(1, Ordering::Relaxed);
let service = s3_service.clone();
tokio::task::spawn_local(async move {
let io = TokioIo::new(stream);
let conn = ConnBuilder::new(TokioExecutor::new())
.serve_connection(io, service)
.into_owned();
if let Err(err) = conn.await {
eprintln!("module-f connection error: {err}");
}
});
}
})
.await
})
})
.with_context(|| format!("failed spawning module-f worker {worker_idx}"))?;
worker_handles.push(handle);
}
let stats_state = state.clone();
let stats_accept_counts = accept_counts.clone();
let _stats_handle = thread::Builder::new()
.name("module-f-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 (compactions, avg_ms) = stats_state.index.compaction_metrics();
eprintln!(
"module-f stats requests={} puts={} gets={} lists={} deletes={} compactions={} compaction_avg_ms={:.3} 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),
compactions,
avg_ms,
accepts.join(",")
);
})
.context("failed spawning module-f stats thread")?;
for handle in worker_handles {
let worker_result = handle
.join()
.map_err(|_| anyhow::anyhow!("module-f worker thread panicked"))?;
worker_result?;
}
Ok(())
}
fn validate_config(config: &ModuleFConfig) -> Result<()> {
if config.workers == 0 {
anyhow::bail!("workers must be > 0");
}
if config.group_bytes == 0 || config.group_bytes % ALIGNMENT != 0 {
anyhow::bail!("group_bytes must be > 0 and aligned to {}", ALIGNMENT);
}
if config.compaction_threshold == 0 {
anyhow::bail!("compaction_threshold must be > 0");
}
if config.flush_interval_ms == 0 {
anyhow::bail!("flush_interval_ms must be > 0");
}
if config.default_bucket.trim().is_empty() {
anyhow::bail!("default bucket must be non-empty");
}
if config.access_key.trim().is_empty() || config.secret_key.trim().is_empty() {
anyhow::bail!("access_key and secret_key must be non-empty");
}
Ok(())
}
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(())
}
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 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 prefix_upper_bound(prefix: &[u8]) -> Option<Vec<u8>> {
let mut upper = prefix.to_vec();
for idx in (0..upper.len()).rev() {
if upper[idx] != 0xFF {
upper[idx] += 1;
upper.truncate(idx + 1);
return Some(upper);
}
}
None
}
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_target_file(target: &Path, allow_fallback: bool) -> Result<(File, String, PathBuf)> {
match open_direct_target(target, false) {
Ok(file) => Ok((file, "block".to_string(), target.to_path_buf())),
Err(err) => {
if !allow_fallback {
return Err(err).with_context(|| {
format!(
"opening {} as module-f block target failed and fallback disabled",
target.display()
)
});
}
let fallback = fallback_path_for(target);
if let Some(parent) = fallback.parent() {
if !parent.as_os_str().is_empty() {
fs::create_dir_all(parent).with_context(|| {
format!("failed to create fallback parent {}", parent.display())
})?;
}
}
let file = open_direct_target(&fallback, true).with_context(|| {
format!(
"failed opening module-f file fallback {}",
fallback.display()
)
})?;
Ok((file, "file-fallback".to_string(), fallback))
}
}
}
fn open_direct_target(path: &Path, create_file: bool) -> Result<File> {
let mut opts = OpenOptions::new();
opts.write(true).read(true);
if create_file {
opts.create(true).truncate(true).mode(0o644);
}
#[cfg(target_os = "linux")]
{
opts.custom_flags(libc::O_DIRECT | libc::O_DSYNC);
}
opts.open(path)
.with_context(|| format!("open failed for {}", path.display()))
}
fn fallback_path_for(target: &Path) -> PathBuf {
if target.starts_with("/dev") {
PathBuf::from("/tmp/tracer-bullet-module-f-direct.bin")
} else {
target.with_extension("direct.bin")
}
}
fn run_group_commit_loop(
file: File,
submit_rx: mpsc::Receiver<AppendRequest>,
group_bytes: usize,
flush_interval: Duration,
) -> Result<()> {
let mut logical = Vec::<u8>::with_capacity(group_bytes + group_bytes / 8);
let mut pending = Vec::<PendingRecord>::with_capacity(128);
let mut file_offset = 0_u64;
loop {
match submit_rx.recv_timeout(flush_interval) {
Ok(request) => {
let start = logical.len();
logical.extend_from_slice(&(request.payload.len() as u32).to_le_bytes());
logical.extend_from_slice(&request.original_len.to_le_bytes());
logical.extend_from_slice(&request.crc32c.to_le_bytes());
logical.extend_from_slice(&request.payload);
pending.push(PendingRecord {
record_start: start,
stored_len: request.payload.len() as u32,
original_len: request.original_len,
crc32c: request.crc32c,
response_tx: request.response_tx,
});
if logical.len() >= group_bytes {
flush_group(&file, &mut logical, &mut pending, &mut file_offset)?;
}
}
Err(mpsc::RecvTimeoutError::Timeout) => {
if !logical.is_empty() {
flush_group(&file, &mut logical, &mut pending, &mut file_offset)?;
}
}
Err(mpsc::RecvTimeoutError::Disconnected) => {
if !logical.is_empty() {
flush_group(&file, &mut logical, &mut pending, &mut file_offset)?;
}
return Ok(());
}
}
}
}
fn flush_group(
file: &File,
logical: &mut Vec<u8>,
pending: &mut Vec<PendingRecord>,
file_offset: &mut u64,
) -> Result<()> {
if logical.is_empty() {
return Ok(());
}
let physical_len = align_up(logical.len(), ALIGNMENT);
let mut aligned = vec![0_u8; physical_len + ALIGNMENT];
let align_start = aligned_offset(aligned.as_ptr() as usize, ALIGNMENT);
let slice = &mut aligned[align_start..align_start + physical_len];
slice[..logical.len()].copy_from_slice(logical);
if physical_len > logical.len() {
slice[logical.len()..].fill(0);
}
if (slice.as_ptr() as usize) % ALIGNMENT != 0
|| slice.len() % ALIGNMENT != 0
|| (*file_offset % ALIGNMENT as u64) != 0
{
anyhow::bail!("unaligned direct-io group flush");
}
write_all_at(file, slice, *file_offset).context("module-f group write failed")?;
for item in pending.drain(..) {
let payload_offset = *file_offset + item.record_start as u64 + RECORD_HEADER_BYTES as u64;
let result = StoredObjectMeta {
payload_offset,
stored_len: item.stored_len,
original_len: item.original_len,
crc32c: item.crc32c,
last_modified: Timestamp::from(SystemTime::now()),
};
let _ = item.response_tx.send(Ok(result));
}
*file_offset += physical_len as u64;
logical.clear();
Ok(())
}
fn write_all_at(file: &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 align_up(value: usize, align: usize) -> usize {
if value % align == 0 {
value
} else {
value + (align - (value % align))
}
}
fn aligned_offset(ptr: usize, align: usize) -> usize {
(align - (ptr % align)) % align
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn prefix_upper_bound_increments_correctly() {
assert_eq!(prefix_upper_bound(b"abc").unwrap(), b"abd");
}
#[test]
fn alignment_helpers_work() {
assert_eq!(align_up(4096, 4096), 4096);
assert_eq!(align_up(4097, 4096), 8192);
assert_eq!(aligned_offset(0, 4096), 0);
assert_eq!(aligned_offset(1, 4096), 4095);
}
#[test]
fn storage_key_round_trip() {
let full = storage_key("bucket", "path/file");
assert_eq!(strip_storage_prefix("bucket", &full).unwrap(), "path/file");
}
}