Skip to content
Draft
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
6 changes: 5 additions & 1 deletion src/setup.rs
Original file line number Diff line number Diff line change
Expand Up @@ -359,12 +359,16 @@ pub fn init_tracing(daemon: bool) {
let upload_cap = XET_UPLOAD_CONCURRENCY_CAP.to_string();
for (k, v) in [
("HF_XET_CLIENT_AC_INITIAL_DOWNLOAD_CONCURRENCY", "16"),
// Let the adaptive controller scale far enough to saturate fast links
// under many concurrent readers (xet-core's default cap is 64).
("HF_XET_CLIENT_AC_MAX_DOWNLOAD_CONCURRENCY", "124"),
("HF_XET_CLIENT_AC_MIN_BYTES_REQUIRED_FOR_ADJUSTMENT", "4194304"),
("HF_XET_RECONSTRUCTION_MIN_RECONSTRUCTION_FETCH_SIZE", "8388608"),
("HF_XET_RECONSTRUCTION_MIN_PREFETCH_BUFFER", "8388608"),
("HF_XET_RECONSTRUCTION_TARGET_BLOCK_COMPLETION_TIME", "30"),
("HF_XET_RECONSTRUCTION_DOWNLOAD_BUFFER_SIZE", "134217728"),
("HF_XET_RECONSTRUCTION_DOWNLOAD_BUFFER_LIMIT", "268435456"),
// Also the budget split between per-stream read buffers (see xet.rs).
("HF_XET_RECONSTRUCTION_DOWNLOAD_BUFFER_LIMIT", "1073741824"),
// Per-read inactivity timeout for CAS/CDN transfers (resets on every byte
// received, so slow-but-progressing reads are fine). This governs the
// DOWNLOAD/reconstruction path (term fetches and whole-file downloads);
Expand Down
193 changes: 184 additions & 9 deletions src/xet.rs
Original file line number Diff line number Diff line change
@@ -1,16 +1,18 @@
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};

use bytes::Bytes;
use xet_client::cas_client::Client;
use xet_client::cas_types::FileRange;
use xet_client::chunk_cache::ChunkCache;
use xet_core_structures::merklehash::MerkleHash;
use xet_core_structures::xorb_object::constants::MAX_XORB_BYTES;
use xet_data::file_reconstruction::{DownloadStream, FileReconstructor};
use xet_data::processing::configurations::TranslatorConfig;
use xet_data::processing::{FileDownloadSession, FileUploadSession, Sha256Policy, SingleFileCleaner, XetFileInfo};
use xet_runtime::core::XetContext;
use xet_runtime::utils::adjustable_semaphore::AdjustableSemaphore;

use crate::error::{Error, Result};

Expand Down Expand Up @@ -48,6 +50,110 @@ pub trait DownloadStreamOps: Send {
async fn next(&mut self) -> Result<Option<Bytes>>;
}

// ── Per-stream download buffers ───────────────────────────────────────

/// Ceiling of a single stream's buffer: pipelining depth of a lone reader.
const STREAM_BUFFER_MAX: u64 = 256 * 1_048_576;

/// Floor of a single stream's buffer: one full xorb, the largest possible
/// term. xet-core clamps a term acquire against the semaphore total only
/// when it is issued (`AdjustableSemaphore::to_physical_acquire`) and does
/// not re-clamp pending acquires on shrink, so shrinking below the largest
/// term would leave an in-flight acquire that can never be satisfied. The
/// floor can become memory-driven once that is fixed upstream.
fn stream_buffer_min() -> u64 {
*MAX_XORB_BYTES as u64
}

/// Gives every read stream a private download buffer.
///
/// Streams are consumed by FUSE `read()` calls that each pin a worker thread.
/// With xet-core's single global buffer (FIFO), a stream that already holds
/// buffer cannot release it until its consumer gets a worker thread, which
/// may itself be blocked waiting for buffer on another stream. Under many
/// concurrent cold readers this starves new streams at their first byte and
/// wedges the mount (#234). Private buffers mean streams never wait on each
/// other. Each is sized `budget / active_streams`, clamped to
/// `[stream_buffer_min(), STREAM_BUFFER_MAX]`, so memory in flight is bounded
/// by `max(budget, streams * stream_buffer_min())`.
struct StreamBufferPool {
budget: u64,
inner: Mutex<StreamBuffers>,
}

#[derive(Default)]
struct StreamBuffers {
active: Vec<Arc<AdjustableSemaphore>>,
/// Share currently applied to every active buffer.
share: u64,
}

impl StreamBufferPool {
fn new(budget: u64) -> Arc<Self> {
Arc::new(Self {
budget,
inner: Mutex::new(StreamBuffers::default()),
})
}

/// Create a buffer for a new stream, already sized to the new fair
/// share, and resize the other active buffers to match. The guard
/// unregisters the buffer on drop.
fn register(self: &Arc<Self>) -> StreamBufferGuard {
let mut inner = self.inner.lock().expect("stream buffers poisoned");
let share = self.share_for(inner.active.len() + 1);
let buffer = AdjustableSemaphore::new(share, (stream_buffer_min(), STREAM_BUFFER_MAX));
inner.active.push(buffer.clone());
Self::apply_share(&mut inner, share);
StreamBufferGuard {
pool: self.clone(),
buffer,
}
}

fn unregister(&self, buffer: &Arc<AdjustableSemaphore>) {
let mut inner = self.inner.lock().expect("stream buffers poisoned");
if let Some(index) = inner.active.iter().position(|other| Arc::ptr_eq(other, buffer)) {
inner.active.swap_remove(index);
}
let share = self.share_for(inner.active.len());
Self::apply_share(&mut inner, share);
}

/// Resize every active buffer to `share`; a no-op when it is unchanged.
/// Shrinks apply lazily: permits already held are reclaimed as they return.
fn apply_share(inner: &mut StreamBuffers, share: u64) {
if inner.share == share {
return;
}
inner.share = share;
for buffer in &inner.active {
// Each call is a no-op when the target is on the other side of
// the current total; the increment's virtual permit releases the
// added capacity on drop.
drop(buffer.increment_permits_to_target(share));
buffer.decrement_permits_to_target(share);
}
}

fn share_for(&self, active_streams: usize) -> u64 {
(self.budget / active_streams.max(1) as u64).clamp(stream_buffer_min(), STREAM_BUFFER_MAX)
}
}

/// Keeps a stream's buffer registered in its pool; dropping it (with the
/// stream) hands the freed share back to the remaining streams.
struct StreamBufferGuard {
pool: Arc<StreamBufferPool>,
buffer: Arc<AdjustableSemaphore>,
}

impl Drop for StreamBufferGuard {
fn drop(&mut self) {
self.pool.unregister(&self.buffer);
}
}

// ── XetSessions ───────────────────────────────────────────────────────

/// Core xet-core sessions for CAS downloads and uploads.
Expand All @@ -61,6 +167,7 @@ pub struct XetSessions {
/// Chunk cache attached to unbounded streams; bounded range downloads skip it
/// to avoid pulling whole xorbs for small range requests.
chunk_cache: Option<Arc<dyn ChunkCache>>,
stream_buffers: Arc<StreamBufferPool>,
}

impl XetSessions {
Expand All @@ -71,34 +178,43 @@ impl XetSessions {
cas_client: Arc<dyn Client>,
chunk_cache: Option<Arc<dyn ChunkCache>>,
) -> Arc<Self> {
// The same knob xet-core uses for its global buffer, so the existing
// HF_XET_RECONSTRUCTION_DOWNLOAD_BUFFER_LIMIT override keeps working.
let stream_buffers = StreamBufferPool::new(ctx.config.reconstruction.download_buffer_limit.as_u64());
Arc::new(Self {
ctx,
session,
upload_config,
cas_client,
chunk_cache,
stream_buffers,
})
}

/// Start a streaming download for a byte range.
/// When `end` is `Some`, only bytes `[offset, end)` are fetched (bounded range).
/// When `end` is `None`, fetches from `offset` to end of file (unbounded stream).
pub fn download_stream(&self, file_info: &XetFileInfo, offset: u64, end: Option<u64>) -> Result<DownloadStream> {
fn download_stream(&self, file_info: &XetFileInfo, offset: u64, end: Option<u64>) -> Result<DownloadStreamWrapper> {
let hash = file_info
.merkle_hash()
.map_err(|e| Error::Xet(format!("invalid hash: {e}")))?;
let is_unbounded = end.is_none();
let file_size = file_info.file_size().unwrap_or(u64::MAX);
let end = end.unwrap_or(file_size);
let mut reconstructor =
FileReconstructor::new(&self.ctx, &self.cas_client, hash).with_byte_range(FileRange::new(offset, end));
let buffer = self.stream_buffers.register();
let mut reconstructor = FileReconstructor::new(&self.ctx, &self.cas_client, hash)
.with_byte_range(FileRange::new(offset, end))
.with_buffer_semaphore(buffer.buffer.clone());
// Attach chunk cache only to the unbounded stream path: the xorb disk
// cache pulls full xorbs (~64MB) even for small range requests, which
// is wasteful for random reads. Sequential reads (unbounded) benefit.
if is_unbounded && let Some(cache) = self.chunk_cache.as_ref() {
reconstructor = reconstructor.with_chunk_cache(cache.clone());
}
Ok(reconstructor.reconstruct_to_stream())
Ok(DownloadStreamWrapper {
stream: reconstructor.reconstruct_to_stream(),
_buffer: buffer,
})
}
}

Expand Down Expand Up @@ -146,8 +262,7 @@ impl XetOps for XetSessions {
offset: u64,
end: Option<u64>,
) -> Result<Box<dyn DownloadStreamOps>> {
let stream = self.download_stream(file_info, offset, end)?;
Ok(Box::new(DownloadStreamWrapper(stream)))
Ok(Box::new(self.download_stream(file_info, offset, end)?))
}

async fn warm_reconstruction_cache(&self, xet_hash: &str) {
Expand All @@ -159,12 +274,17 @@ impl XetOps for XetSessions {

// ── DownloadStreamWrapper ─────────────────────────────────────────────

struct DownloadStreamWrapper(DownloadStream);
struct DownloadStreamWrapper {
stream: DownloadStream,
/// Declared after `stream` so the stream (and its reconstruction task) is
/// cancelled before the buffer share is handed back.
_buffer: StreamBufferGuard,
}

#[async_trait::async_trait]
impl DownloadStreamOps for DownloadStreamWrapper {
async fn next(&mut self) -> Result<Option<Bytes>> {
Ok(self.0.next().await?)
Ok(self.stream.next().await?)
}
}

Expand Down Expand Up @@ -346,3 +466,58 @@ impl StreamingWriterOps for StreamingWriter {
self.bytes_written == 0
}
}

#[cfg(test)]
mod stream_buffer_tests {
use super::*;

const MIB: u64 = 1_048_576;
const BUDGET: u64 = 1024 * MIB;

fn total(guard: &StreamBufferGuard) -> u64 {
guard.buffer.total_permits()
}

#[test]
fn lone_stream_gets_the_ceiling() {
let pool = StreamBufferPool::new(BUDGET);
let stream = pool.register();
assert_eq!(total(&stream), 256 * MIB);
}

#[test]
fn shares_shrink_to_the_floor_and_grow_back() {
let pool = StreamBufferPool::new(BUDGET);
let mut streams: Vec<_> = (0..4).map(|_| pool.register()).collect();
// 1 GiB / 4 = 256 MiB: still at the ceiling.
assert!(streams.iter().all(|stream| total(stream) == 256 * MIB));

streams.extend((0..4).map(|_| pool.register()));
// 1 GiB / 8 = 128 MiB for everyone, including the early streams.
assert!(streams.iter().all(|stream| total(stream) == 128 * MIB));

streams.extend((0..56).map(|_| pool.register()));
// 1 GiB / 64 = 16 MiB, clamped up to the floor.
assert!(streams.iter().all(|stream| total(stream) == 64 * MIB));

streams.truncate(2);
// 1 GiB / 2 = 512 MiB, clamped down to the ceiling.
assert!(streams.iter().all(|stream| total(stream) == 256 * MIB));
assert_eq!(pool.inner.lock().unwrap().active.len(), 2);
}

#[tokio::test]
async fn shrink_applies_once_held_permits_return() {
let pool = StreamBufferPool::new(BUDGET);
let first = pool.register();
let held = first.buffer.acquire_many(200 * MIB).await.unwrap();

let _others: Vec<_> = (0..7).map(|_| pool.register()).collect();
// Target is 128 MiB but 200 MiB are out: the total is updated now,
// the part that cannot be reclaimed yet stays pending.
assert_eq!(total(&first), 128 * MIB);
assert!(first.buffer.available_permits() <= 56 * MIB);
drop(held);
assert_eq!(first.buffer.available_permits(), 128 * MIB);
}
}
Loading