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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,8 @@ See [adaptive fill admission](CONFIGURATION.md#adaptive-fill-admission) for rate

The managed-memory limit covers the index mapping, heat bits, L1, append buffers, reclaim buffers, metadata, cache-owned thread stacks, recovery scratch, and transient reads. Total deployment memory additionally includes allocator metadata, Tokio, process overhead, and the kernel page cache. `CacheConfig::new` rejects invalid or insufficient memory budgets before file access; actual allocation can still fail during open.

`MemoryReservation` holds a charge until its allocation is released. `BufferLease` owns an aligned allocation directly and drops it before returning an optional individual charge; append staging instead keeps one aggregate reservation for its fixed buffers and record arrays. A failed buffer allocation releases its reservation automatically.

### Storage path

C² owns one logical data path. Multi-device deployments stripe below the filesystem with RAID0 or an equivalent layer. Request routing, recovery identity, and descriptor count stay independent of device topology.
Expand Down
2 changes: 2 additions & 0 deletions CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ Follow the surrounding code and the design constraints in [ARCHITECTURE.md](ARCH

Name types by their concrete domain role. Use `Options` for editable inputs awaiting validation and `Config` for validated or resolved configuration. Keep names such as `Layout` for geometry and use specific domain nouns such as `Desc`, `Candidate`, `Victims`, or `Inspection` for runtime data. Avoid `Plan` as a type-name suffix; computing data before using it does not by itself make that data a plan. Keep related fields, local variables, and functions consistent with the type's role.

Use ordinary names for fields; annotate required but unread ownership fields with a narrowly scoped `#[expect(dead_code)]` and explain their role with `reason`. Local guards and other resources retained until scope exit may use conventional underscore-prefixed bindings such as `_guard` without a lint attribute. Preserve their lifetimes and drop order; the wildcard pattern `_` does not retain a value. Remove unnecessary bindings and private parameters, and use a narrowly scoped `#[expect(unused_variables)]` for required but unused trait parameters.

Declare restricted visibility at module boundaries and use `pub` for items in those modules' APIs. Keep items private when only their defining module and its descendants need them. For items reachable through public modules or re-exported public types, reserve `pub` for intentional public API and use narrower visibility for internal callers.

## Documentation
Expand Down
19 changes: 6 additions & 13 deletions benchmarks/cache/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -408,7 +408,6 @@ async fn run(config: BenchConfig) -> io::Result<()> {
.transpose()?;
report(
"put_drain",
"put + drain",
"write",
config.write_clients,
&write.measurement,
Expand Down Expand Up @@ -437,7 +436,7 @@ async fn run(config: BenchConfig) -> io::Result<()> {
cache.close_warm().await?;
drop(cache);
let warm_close = started.elapsed();
report_latency("warm_close", "warm close", warm_close);
report_latency("warm_close", warm_close);

let cache = Arc::new(Cache::open(files.data(), cache_config.clone()).await?);
if cache.startup_mode() != StartupMode::Warm {
Expand All @@ -459,7 +458,6 @@ async fn run(config: BenchConfig) -> io::Result<()> {
if l1_entry_eligible {
report(
"l2_promote",
"L2 get + promote",
"read",
config.l2_clients(),
&l2_read.measurement,
Expand All @@ -470,7 +468,6 @@ async fn run(config: BenchConfig) -> io::Result<()> {
} else {
report(
"l2_read",
"L2 get",
"read",
config.l2_clients(),
&l2_read.measurement,
Expand All @@ -485,7 +482,7 @@ async fn run(config: BenchConfig) -> io::Result<()> {
l2_read.measurement.elapsed,
&l2_read.primary,
);
report_read_latency("l2_latency", "L2 get latency", &l2_read.latency);
report_read_latency("l2_latency", &l2_read.latency);
l2_read.measurement
} else {
let _ = concurrent_writes(
Expand Down Expand Up @@ -522,7 +519,6 @@ async fn run(config: BenchConfig) -> io::Result<()> {
.await?;
report(
"l2_hot_scan",
"L2 cold scan",
"read",
config.l2_clients(),
&cold_scan.measurement,
Expand All @@ -536,7 +532,7 @@ async fn run(config: BenchConfig) -> io::Result<()> {
cold_scan.measurement.elapsed,
&cold_scan.primary,
);
report_read_latency("l2_scan_latency", "L2 scan latency", &cold_scan.latency);
report_read_latency("l2_scan_latency", &cold_scan.latency);
report_tiers(
"hot_during_scan",
"hot during scan",
Expand Down Expand Up @@ -668,15 +664,14 @@ async fn run(config: BenchConfig) -> io::Result<()> {
.await?;
report(
"resident_l1",
"resident L1 get",
"read",
config.clients,
&resident.measurement,
0,
Some(&resident.latency),
config.read_latency_sample_interval,
);
report_read_latency("resident_l1_latency", "L1 get latency", &resident.latency);
report_read_latency("resident_l1_latency", &resident.latency);
if let Some(before) = before {
let after = cache.snapshot()?;
println!(
Expand Down Expand Up @@ -961,10 +956,8 @@ fn verify_value(ordinal: usize, value: &[u8]) -> io::Result<()> {
Ok(())
}

#[allow(clippy::too_many_arguments)]
fn report(
phase: &str,
_name: &str,
operation: &str,
workers: usize,
measurement: &Measurement,
Expand Down Expand Up @@ -1032,15 +1025,15 @@ fn report_tiers(phase: &str, name: &str, elapsed: Duration, tiers: &TierCounts)
);
}

fn report_latency(phase: &str, _name: &str, elapsed: Duration) {
fn report_latency(phase: &str, elapsed: Duration) {
JobReport::new("cache", None, phase, "control", elapsed, 1).emit();
println!(
"result phase={phase} elapsed_ns={} operations=0 bytes=0 ops_per_sec=0.000 mib_per_sec=0.000 checksum=0000000000000000",
elapsed.as_nanos(),
);
}

fn report_read_latency(phase: &str, _name: &str, latency: &LatencyHistogram) {
fn report_read_latency(phase: &str, latency: &LatencyHistogram) {
let summary = latency.summary();
if summary.samples == 0 {
return;
Expand Down
6 changes: 3 additions & 3 deletions cache2/src/cache/runtime/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -883,12 +883,12 @@ impl CacheRuntime {
current_bytes,
} => {
if ADMIT_L1 {
let _published = state.memory.publish(hash, key, value, seqno);
state.memory.publish(hash, key, value, seqno);
} else {
// Prevent an older exact-key L1 value from indefinitely
// shadowing the prefetched L2 record. Contention remains a
// valid best-effort stale outcome.
let _removed = state.memory.delete(hash, key, seqno);
state.memory.delete(hash, key, seqno);
}
if should_wake_write(
previous_bytes,
Expand Down Expand Up @@ -940,7 +940,7 @@ impl CacheRuntime {
}
return Err(write_overload_error());
};
let _removed = state.memory.delete(hash, key, seqno);
state.memory.delete(hash, key, seqno);
if let Some(activity) = activity {
RuntimeMetrics::increment(&activity.deletes);
}
Expand Down
45 changes: 38 additions & 7 deletions cache2/src/io/engine/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -118,13 +118,22 @@ impl BlockingIo {
}

impl PositionedIo for BlockingIo {
fn read_at(&self, buffer: &mut [u8], _offset: u64) -> io::Result<usize> {
fn read_at(
&self,
buffer: &mut [u8],
#[expect(unused_variables)] offset: u64,
) -> io::Result<usize> {
self.enter_and_wait();
buffer.fill(0);
Ok(buffer.len())
}

fn write_at(&self, _point: WritePoint, buffer: &[u8], _offset: u64) -> io::Result<usize> {
fn write_at(
&self,
#[expect(unused_variables)] point: WritePoint,
buffer: &[u8],
#[expect(unused_variables)] offset: u64,
) -> io::Result<usize> {
self.enter_and_wait();
Ok(buffer.len())
}
Expand All @@ -149,21 +158,34 @@ impl PanicOnceIo {
}

impl PositionedIo for PanicOnceIo {
fn read_at(&self, buffer: &mut [u8], _offset: u64) -> io::Result<usize> {
fn read_at(
&self,
buffer: &mut [u8],
#[expect(unused_variables)] offset: u64,
) -> io::Result<usize> {
if self.panic_next_read.swap(false, Ordering::AcqRel) {
panic!("injected io panic");
}
buffer.fill(0);
Ok(buffer.len())
}

fn write_at(&self, _point: WritePoint, buffer: &[u8], _offset: u64) -> io::Result<usize> {
fn write_at(
&self,
#[expect(unused_variables)] point: WritePoint,
buffer: &[u8],
#[expect(unused_variables)] offset: u64,
) -> io::Result<usize> {
Ok(buffer.len())
}
}

impl PositionedIo for ShortThenErrorIo {
fn read_at(&self, buffer: &mut [u8], _offset: u64) -> io::Result<usize> {
fn read_at(
&self,
buffer: &mut [u8],
#[expect(unused_variables)] offset: u64,
) -> io::Result<usize> {
if self.read_calls.fetch_add(1, Ordering::Relaxed) == 0 {
let transferred = 3.min(buffer.len());
buffer[..transferred].fill(0x5a);
Expand All @@ -173,7 +195,12 @@ impl PositionedIo for ShortThenErrorIo {
}
}

fn write_at(&self, _point: WritePoint, buffer: &[u8], _offset: u64) -> io::Result<usize> {
fn write_at(
&self,
#[expect(unused_variables)] point: WritePoint,
buffer: &[u8],
#[expect(unused_variables)] offset: u64,
) -> io::Result<usize> {
if self.write_calls.fetch_add(1, Ordering::Relaxed) == 0 {
Ok(3.min(buffer.len()))
} else {
Expand All @@ -199,7 +226,11 @@ fn read_buffer(managed_memory: &Arc<ManagedMemory>, length: usize) -> IoBuffer {

fn write_buffer(managed_memory: &Arc<ManagedMemory>, bytes: &[u8]) -> IoBuffer {
let mut lease = managed_memory.try_read_buffer(bytes.len()).unwrap();
lease.prepare(bytes.len()).unwrap().copy_from_slice(bytes);
let target = lease.read_target(bytes.len()).unwrap();
// SAFETY: the exclusively owned target fits the source, and the allocations
// do not overlap. Publish the initialized range only after copying it.
unsafe { target.copy_from_nonoverlapping(bytes.as_ptr(), bytes.len()) };
lease.mark_initialized(bytes.len()).unwrap();
IoBuffer::for_write(lease, bytes.len()).unwrap()
}

Expand Down
22 changes: 18 additions & 4 deletions cache2/src/io/file.rs
Original file line number Diff line number Diff line change
Expand Up @@ -523,7 +523,7 @@ impl StorageFile for CacheFile {
}
}

fn sync(&self, _point: SyncPoint, mode: SyncMode) -> io::Result<()> {
fn sync(&self, #[expect(unused_variables)] point: SyncPoint, mode: SyncMode) -> io::Result<()> {
match mode {
SyncMode::Data => self.file.sync_data(),
SyncMode::All => self.file.sync_all(),
Expand Down Expand Up @@ -572,7 +572,12 @@ impl PositionedIo for CacheFile {
self.file.read_at(buffer, offset)
}

fn write_at(&self, _point: WritePoint, buffer: &[u8], offset: u64) -> io::Result<usize> {
fn write_at(
&self,
#[expect(unused_variables)] point: WritePoint,
buffer: &[u8],
offset: u64,
) -> io::Result<usize> {
self.file.write_at(buffer, offset)
}
}
Expand Down Expand Up @@ -874,15 +879,24 @@ mod tests {
}

impl PositionedIo for InterruptedIo {
fn read_at(&self, _buffer: &mut [u8], _offset: u64) -> io::Result<usize> {
fn read_at(
&self,
#[expect(unused_variables)] buffer: &mut [u8],
#[expect(unused_variables)] offset: u64,
) -> io::Result<usize> {
self.calls.fetch_add(1, Ordering::Relaxed);
Err(io::Error::new(
io::ErrorKind::Interrupted,
"interrupted read",
))
}

fn write_at(&self, _point: WritePoint, _buffer: &[u8], _offset: u64) -> io::Result<usize> {
fn write_at(
&self,
#[expect(unused_variables)] point: WritePoint,
#[expect(unused_variables)] buffer: &[u8],
#[expect(unused_variables)] offset: u64,
) -> io::Result<usize> {
self.calls.fetch_add(1, Ordering::Relaxed);
Err(io::Error::new(
io::ErrorKind::Interrupted,
Expand Down
Loading