Files
nixos/overlays/attic-client.patch
T

605 lines
20 KiB
Diff
Raw Normal View History

2025-11-19 14:25:35 -06:00
diff --git c/attic/src/nix_store/bindings/mod.rs i/attic/src/nix_store/bindings/mod.rs
index a16b52a..37faf08 100644
--- c/attic/src/nix_store/bindings/mod.rs
+++ i/attic/src/nix_store/bindings/mod.rs
@@ -10,6 +10,8 @@ use tokio::io::{AsyncWrite, AsyncWriteExt};
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
use crate::{AtticError, AtticResult};
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
+use super::stream_config::NarStreamConfig;
+
// The C++ implementation takes care of concurrency
#[repr(transparent)]
pub struct FfiNixStore(UnsafeCell<cxx::UniquePtr<ffi::CNixStore>>);
@@ -42,10 +44,12 @@ pub unsafe fn open_nix_store() -> AtticResult<FfiNixStore> {
mod mpsc {
// Tokio
pub use tokio::sync::mpsc::{
- error::SendError, unbounded_channel, UnboundedReceiver, UnboundedSender,
+ error::SendError, channel, Receiver, Sender,
};
}
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
+
+
/// Async write request.
#[derive(Debug)]
enum AsyncWriteMessage {
@@ -57,18 +61,20 @@ enum AsyncWriteMessage {
/// Async write request sender.
#[derive(Clone)]
pub struct AsyncWriteSender {
- sender: mpsc::UnboundedSender<AsyncWriteMessage>,
+ sender: mpsc::Sender<AsyncWriteMessage>,
}
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
impl AsyncWriteSender {
fn send(&mut self, data: &[u8]) -> Result<(), mpsc::SendError<AsyncWriteMessage>> {
let message = AsyncWriteMessage::Data(Vec::from(data));
- self.sender.send(message)
+ // Use blocking_send since we're called from C++ in a blocking context
+ // This provides backpressure when the channel is full
+ self.sender.blocking_send(message).map_err(|e| mpsc::SendError(e.0))
}
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
fn eof(&mut self) -> Result<(), mpsc::SendError<AsyncWriteMessage>> {
let message = AsyncWriteMessage::Eof;
- self.sender.send(message)
+ self.sender.blocking_send(message).map_err(|e| mpsc::SendError(e.0))
}
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
pub(crate) fn rust_error(
@@ -76,19 +82,25 @@ impl AsyncWriteSender {
error: impl std::error::Error,
) -> Result<(), impl std::error::Error> {
let message = AsyncWriteMessage::Error(error.to_string());
- self.sender.send(message)
+ // Try to send error without blocking, it's OK if channel is closed
+ self.sender.try_send(message)
}
}
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
/// A wrapper of the `AsyncWrite` trait for the synchronous Nix C++ land.
pub struct AsyncWriteAdapter {
- receiver: mpsc::UnboundedReceiver<AsyncWriteMessage>,
+ receiver: mpsc::Receiver<AsyncWriteMessage>,
eof: bool,
}
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
impl AsyncWriteAdapter {
pub fn new() -> (Self, Box<AsyncWriteSender>) {
- let (sender, receiver) = mpsc::unbounded_channel();
+ Self::new_with_config(NarStreamConfig::default())
+ }
+
+ pub fn new_with_config(config: NarStreamConfig) -> (Self, Box<AsyncWriteSender>) {
+ // Use bounded channel to provide backpressure and prevent OOM
+ let (sender, receiver) = mpsc::channel(config.channel_capacity);
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
let r = Self {
receiver,
diff --git c/attic/src/nix_store/mod.rs i/attic/src/nix_store/mod.rs
index 4e08a67..04ef0f6 100644
--- c/attic/src/nix_store/mod.rs
+++ i/attic/src/nix_store/mod.rs
@@ -46,6 +46,9 @@ mod bindings;
#[cfg(feature = "nix_store")]
mod nix_store;
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
+#[cfg(feature = "nix_store")]
+mod stream_config;
+
use std::ffi::OsStr;
#[cfg(target_family = "unix")]
use std::os::unix::ffi::OsStrExt;
@@ -61,6 +64,9 @@ use crate::hash::Hash;
#[cfg(feature = "nix_store")]
pub use nix_store::NixStore;
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
+#[cfg(feature = "nix_store")]
+pub use stream_config::NarStreamConfig;
+
#[cfg(test)]
pub mod tests;
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
diff --git c/attic/src/nix_store/nix_store.rs i/attic/src/nix_store/nix_store.rs
index 3b754f9..91eb489 100644
--- c/attic/src/nix_store/nix_store.rs
+++ i/attic/src/nix_store/nix_store.rs
@@ -8,6 +8,7 @@ use std::sync::Arc;
use tokio::task::spawn_blocking;
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
use super::bindings::{open_nix_store, AsyncWriteAdapter, FfiNixStore};
+use super::stream_config::NarStreamConfig;
use super::{to_base_name, StorePath, ValidPathInfo};
use crate::error::AtticResult;
use crate::hash::Hash;
@@ -19,11 +20,18 @@ pub struct NixStore {
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
/// Path to the Nix store itself.
store_dir: PathBuf,
+
+ /// Configuration for NAR streaming.
+ stream_config: NarStreamConfig,
}
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
#[cfg(feature = "nix_store")]
impl NixStore {
pub fn connect() -> AtticResult<Self> {
+ Self::connect_with_config(NarStreamConfig::from_env())
+ }
+
+ pub fn connect_with_config(stream_config: NarStreamConfig) -> AtticResult<Self> {
#[allow(unsafe_code)]
let inner = unsafe { open_nix_store()? };
let store_dir = PathBuf::from(inner.store().store_dir());
@@ -31,6 +39,7 @@ impl NixStore {
Ok(Self {
inner: Arc::new(inner),
store_dir,
+ stream_config,
})
}
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
@@ -84,12 +93,26 @@ impl NixStore {
self.store_dir.join(&store_path.base_name)
}
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
- /// Creates a NAR archive from a path.
+ /// Creates a NAR archive from a path with default configuration.
///
/// This is akin to `nix-store --dump`.
pub fn nar_from_path(&self, store_path: StorePath) -> AsyncWriteAdapter {
+ self.nar_from_path_with_config(store_path, self.stream_config)
+ }
+
+ /// Creates a NAR archive from a path with custom streaming configuration.
+ ///
+ /// This is akin to `nix-store --dump`.
+ ///
+ /// The `stream_config` parameter allows you to control the memory buffer size
+ /// for the NAR stream, which is important for large files to prevent OOM.
+ pub fn nar_from_path_with_config(
+ &self,
+ store_path: StorePath,
+ stream_config: NarStreamConfig,
+ ) -> AsyncWriteAdapter {
let inner = self.inner.clone();
- let (adapter, mut sender) = AsyncWriteAdapter::new();
+ let (adapter, mut sender) = AsyncWriteAdapter::new_with_config(stream_config);
let base_name = Vec::from(store_path.as_base_name_bytes());
2025-12-08 05:53:17 -06:00
2025-11-19 14:25:35 -06:00
spawn_blocking(move || {
diff --git c/attic/src/nix_store/stream_config.rs i/attic/src/nix_store/stream_config.rs
new file mode 100644
index 0000000..bdc6177
--- /dev/null
+++ i/attic/src/nix_store/stream_config.rs
@@ -0,0 +1,116 @@
+//! Configuration for NAR streaming.
+//!
+//! This module provides configuration options for controlling memory usage
+//! during NAR stream uploads.
+
+use std::env;
+
+/// Default capacity for the NAR streaming channel.
+///
+/// At ~256KB per chunk (typical), this allows:
+/// - 64 chunks = 16 MB buffer (conservative)
+/// - 128 chunks = 32 MB buffer (balanced, default)
+/// - 256 chunks = 64 MB buffer (high throughput)
+pub const DEFAULT_NAR_STREAM_CHANNEL_CAPACITY: usize = 128;
+
+/// Minimum allowed channel capacity to prevent starvation.
+pub const MIN_NAR_STREAM_CHANNEL_CAPACITY: usize = 8;
+
+/// Maximum allowed channel capacity to prevent excessive memory usage.
+pub const MAX_NAR_STREAM_CHANNEL_CAPACITY: usize = 1024;
+
+/// Configuration for NAR streaming behavior.
+#[derive(Debug, Clone, Copy)]
+pub struct NarStreamConfig {
+ /// Number of chunks to buffer in the NAR streaming channel.
+ ///
+ /// This provides backpressure to prevent the C++ NAR generator from
+ /// getting too far ahead of network upload, which could cause OOM
+ /// when uploading large files.
+ pub channel_capacity: usize,
+}
+
+impl Default for NarStreamConfig {
+ fn default() -> Self {
+ Self {
+ channel_capacity: DEFAULT_NAR_STREAM_CHANNEL_CAPACITY,
+ }
+ }
+}
+
+impl NarStreamConfig {
+ /// Creates a new configuration with the default values.
+ pub fn new() -> Self {
+ Self::default()
+ }
+
+ /// Sets the channel capacity.
+ ///
+ /// The value will be clamped to the valid range.
+ pub fn with_channel_capacity(mut self, capacity: usize) -> Self {
+ self.channel_capacity = capacity.clamp(
+ MIN_NAR_STREAM_CHANNEL_CAPACITY,
+ MAX_NAR_STREAM_CHANNEL_CAPACITY,
+ );
+ self
+ }
+
+ /// Loads configuration from environment variables.
+ ///
+ /// Supported environment variables:
+ /// - `ATTIC_NAR_STREAM_CHANNEL_CAPACITY`: Channel capacity (default: 128)
+ pub fn from_env() -> Self {
+ let mut config = Self::default();
+
+ if let Ok(capacity_str) = env::var("ATTIC_NAR_STREAM_CHANNEL_CAPACITY") {
+ if let Ok(capacity) = capacity_str.parse::<usize>() {
+ config = config.with_channel_capacity(capacity);
+ }
+ }
+
+ config
+ }
+
+ /// Estimates the maximum memory usage for this configuration.
+ ///
+ /// This assumes an average chunk size. Actual usage may vary.
+ pub fn estimated_max_memory_bytes(&self, avg_chunk_size: usize) -> usize {
+ self.channel_capacity * avg_chunk_size
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ #[test]
+ fn test_default_config() {
+ let config = NarStreamConfig::default();
+ assert_eq!(config.channel_capacity, DEFAULT_NAR_STREAM_CHANNEL_CAPACITY);
+ }
+
+ #[test]
+ fn test_with_channel_capacity() {
+ let config = NarStreamConfig::new().with_channel_capacity(256);
+ assert_eq!(config.channel_capacity, 256);
+ }
+
+ #[test]
+ fn test_channel_capacity_clamping() {
+ // Test minimum clamping
+ let config = NarStreamConfig::new().with_channel_capacity(1);
+ assert_eq!(config.channel_capacity, MIN_NAR_STREAM_CHANNEL_CAPACITY);
+
+ // Test maximum clamping
+ let config = NarStreamConfig::new().with_channel_capacity(10000);
+ assert_eq!(config.channel_capacity, MAX_NAR_STREAM_CHANNEL_CAPACITY);
+ }
+
+ #[test]
+ fn test_estimated_memory() {
+ let config = NarStreamConfig::new().with_channel_capacity(128);
+ let chunk_size = 256 * 1024; // 256 KB
+ let estimated = config.estimated_max_memory_bytes(chunk_size);
+ assert_eq!(estimated, 128 * 256 * 1024); // 32 MB
+ }
+}
diff --git c/attic/src/nix_store/stream_test.rs i/attic/src/nix_store/stream_test.rs
new file mode 100644
index 0000000..e9c18f2
--- /dev/null
+++ i/attic/src/nix_store/stream_test.rs
@@ -0,0 +1,304 @@
+//! Tests for NAR streaming with bounded channels.
+
+#[cfg(test)]
+mod tests {
+ use super::super::bindings::AsyncWriteAdapter;
+ use super::super::stream_config::NarStreamConfig;
+ use futures::StreamExt;
+ use std::time::Duration;
+ use tokio::time::timeout;
+
+ /// Test that bounded channel provides backpressure.
+ #[tokio::test]
+ async fn test_bounded_channel_backpressure() {
+ let config = NarStreamConfig::new().with_channel_capacity(2);
+ let (mut adapter, mut sender) = AsyncWriteAdapter::new_with_config(config);
+
+ // Send data in a separate task
+ let send_handle = tokio::spawn(async move {
+ // These should succeed immediately (capacity = 2)
+ sender.send(&[1u8; 100]).unwrap();
+ sender.send(&[2u8; 100]).unwrap();
+
+ // This should block until receiver consumes some data
+ // We spawn this in a blocking context since send uses blocking_send
+ let result = tokio::task::spawn_blocking(move || {
+ sender.send(&[3u8; 100])
+ }).await;
+
+ result.unwrap()
+ });
+
+ // Give sender time to fill the channel
+ tokio::time::sleep(Duration::from_millis(100)).await;
+
+ // Now start consuming
+ let first = adapter.next().await.unwrap().unwrap();
+ assert_eq!(first.len(), 100);
+ assert_eq!(first[0], 1);
+
+ let second = adapter.next().await.unwrap().unwrap();
+ assert_eq!(second.len(), 100);
+ assert_eq!(second[0], 2);
+
+ // The third send should now complete
+ let send_result = timeout(Duration::from_secs(1), send_handle).await;
+ assert!(send_result.is_ok());
+
+ let third = adapter.next().await.unwrap().unwrap();
+ assert_eq!(third.len(), 100);
+ assert_eq!(third[0], 3);
+ }
+
+ /// Test that channel capacity is enforced.
+ #[tokio::test]
+ async fn test_channel_capacity_enforcement() {
+ let config = NarStreamConfig::new().with_channel_capacity(4);
+ let (mut adapter, mut sender) = AsyncWriteAdapter::new_with_config(config);
+
+ // Fill the channel completely
+ tokio::task::spawn_blocking(move || {
+ for i in 0..4 {
+ sender.send(&[i as u8; 50]).unwrap();
+ }
+ sender.eof().unwrap();
+ });
+
+ // Give it time to fill
+ tokio::time::sleep(Duration::from_millis(50)).await;
+
+ // Consume all items
+ let mut count = 0;
+ while let Some(Ok(chunk)) = adapter.next().await {
+ assert_eq!(chunk.len(), 50);
+ assert_eq!(chunk[0], count);
+ count += 1;
+ }
+
+ assert_eq!(count, 4);
+ }
+
+ /// Test that small channel capacity works correctly.
+ #[tokio::test]
+ async fn test_minimum_channel_capacity() {
+ let config = NarStreamConfig::new().with_channel_capacity(8); // minimum
+ let (mut adapter, mut sender) = AsyncWriteAdapter::new_with_config(config);
+
+ let send_handle = tokio::task::spawn_blocking(move || {
+ for i in 0..16 {
+ sender.send(&[i as u8; 10]).unwrap();
+ }
+ sender.eof().unwrap();
+ });
+
+ let mut received = Vec::new();
+ while let Some(Ok(chunk)) = adapter.next().await {
+ received.push(chunk[0]);
+ }
+
+ send_handle.await.unwrap();
+ assert_eq!(received.len(), 16);
+ assert_eq!(received, (0..16).collect::<Vec<u8>>());
+ }
+
+ /// Test that large channel capacity works correctly.
+ #[tokio::test]
+ async fn test_large_channel_capacity() {
+ let config = NarStreamConfig::new().with_channel_capacity(256);
+ let (mut adapter, mut sender) = AsyncWriteAdapter::new_with_config(config);
+
+ let send_handle = tokio::task::spawn_blocking(move || {
+ for i in 0..100 {
+ sender.send(&[i as u8; 1024]).unwrap();
+ }
+ sender.eof().unwrap();
+ });
+
+ let mut count = 0;
+ while let Some(Ok(chunk)) = adapter.next().await {
+ assert_eq!(chunk.len(), 1024);
+ assert_eq!(chunk[0], count);
+ count += 1;
+ }
+
+ send_handle.await.unwrap();
+ assert_eq!(count, 100);
+ }
+
+ /// Test error handling in bounded channel.
+ #[tokio::test]
+ async fn test_error_handling() {
+ let config = NarStreamConfig::new().with_channel_capacity(2);
+ let (mut adapter, mut sender) = AsyncWriteAdapter::new_with_config(config);
+
+ tokio::task::spawn_blocking(move || {
+ sender.send(&[1u8; 100]).unwrap();
+ sender.rust_error(std::io::Error::new(
+ std::io::ErrorKind::Other,
+ "Test error"
+ )).ok(); // try_send, might fail if channel closed
+ });
+
+ // First message should be data
+ let first = adapter.next().await.unwrap().unwrap();
+ assert_eq!(first[0], 1);
+
+ // Second message should be error
+ let second = adapter.next().await.unwrap();
+ assert!(second.is_err());
+ }
+
+ /// Test that EOF is properly handled.
+ #[tokio::test]
+ async fn test_eof_handling() {
+ let config = NarStreamConfig::new().with_channel_capacity(4);
+ let (mut adapter, mut sender) = AsyncWriteAdapter::new_with_config(config);
+
+ tokio::task::spawn_blocking(move || {
+ sender.send(&[1u8; 100]).unwrap();
+ sender.send(&[2u8; 100]).unwrap();
+ sender.eof().unwrap();
+ });
+
+ let first = adapter.next().await.unwrap().unwrap();
+ assert_eq!(first[0], 1);
+
+ let second = adapter.next().await.unwrap().unwrap();
+ assert_eq!(second[0], 2);
+
+ // Should be EOF now
+ assert!(adapter.next().await.is_none());
+ }
+
+ /// Test concurrent producers (should not happen in practice, but test safety).
+ #[tokio::test]
+ async fn test_concurrent_safety() {
+ let config = NarStreamConfig::new().with_channel_capacity(16);
+ let (mut adapter, sender) = AsyncWriteAdapter::new_with_config(config);
+
+ // Spawn multiple producers
+ let mut handles = Vec::new();
+ for i in 0..4 {
+ let mut sender_clone = sender.clone();
+ let handle = tokio::task::spawn_blocking(move || {
+ for j in 0..10 {
+ let val = (i * 10 + j) as u8;
+ sender_clone.send(&[val; 50]).unwrap();
+ }
+ });
+ handles.push(handle);
+ }
+
+ // Wait for all producers
+ for handle in handles {
+ handle.await.unwrap();
+ }
+
+ // Signal EOF from one sender
+ let mut sender = *sender;
+ tokio::task::spawn_blocking(move || {
+ sender.eof().unwrap();
+ });
+
+ // Consume all data
+ let mut count = 0;
+ while let Some(Ok(_chunk)) = adapter.next().await {
+ count += 1;
+ }
+
+ assert_eq!(count, 40); // 4 producers * 10 chunks each
+ }
+
+ /// Benchmark-style test: ensure bounded channel doesn't significantly
+ /// slow down high-throughput scenarios.
+ #[tokio::test]
+ async fn test_throughput_performance() {
+ let config = NarStreamConfig::new().with_channel_capacity(128);
+ let (mut adapter, mut sender) = AsyncWriteAdapter::new_with_config(config);
+
+ let start = std::time::Instant::now();
+
+ // Simulate high-throughput upload
+ let send_handle = tokio::task::spawn_blocking(move || {
+ for i in 0..1000 {
+ // 256KB chunks
+ let chunk = vec![i as u8; 256 * 1024];
+ sender.send(&chunk).unwrap();
+ }
+ sender.eof().unwrap();
+ });
+
+ // Consume as fast as possible
+ let mut total_bytes = 0u64;
+ while let Some(Ok(chunk)) = adapter.next().await {
+ total_bytes += chunk.len() as u64;
+ }
+
+ send_handle.await.unwrap();
+ let elapsed = start.elapsed();
+
+ // Should process 256 MB in reasonable time (< 5 seconds even on slow systems)
+ assert_eq!(total_bytes, 1000 * 256 * 1024);
+ assert!(elapsed.as_secs() < 5, "Throughput test too slow: {:?}", elapsed);
+
+ let throughput_mbps = (total_bytes as f64 / 1024.0 / 1024.0) / elapsed.as_secs_f64();
+ println!("Throughput: {:.2} MB/s", throughput_mbps);
+ }
+
+ /// Test memory-bounded behavior: ensure we don't allocate unbounded memory.
+ #[tokio::test]
+ async fn test_memory_bounded() {
+ let config = NarStreamConfig::new().with_channel_capacity(10);
+ let (mut adapter, mut sender) = AsyncWriteAdapter::new_with_config(config);
+
+ // Try to send more than channel capacity
+ let send_handle = tokio::task::spawn_blocking(move || {
+ // This would OOM with unbounded channel on large iterations
+ // With bounded channel, it will block and wait
+ for i in 0..100 {
+ let chunk = vec![i as u8; 1024 * 1024]; // 1 MB chunks
+ sender.send(&chunk).unwrap();
+ }
+ sender.eof().unwrap();
+ });
+
+ // Slow consumer
+ let mut count = 0;
+ while let Some(Ok(_chunk)) = adapter.next().await {
+ count += 1;
+ // Simulate slow network
+ tokio::time::sleep(Duration::from_millis(1)).await;
+ }
+
+ send_handle.await.unwrap();
+ assert_eq!(count, 100);
+ }
+
+ /// Test configuration edge cases.
+ #[test]
+ fn test_config_clamping() {
+ // Below minimum
+ let config = NarStreamConfig::new().with_channel_capacity(1);
+ assert_eq!(config.channel_capacity, 8); // clamped to minimum
+
+ // Above maximum
+ let config = NarStreamConfig::new().with_channel_capacity(10000);
+ assert_eq!(config.channel_capacity, 1024); // clamped to maximum
+
+ // Within range
+ let config = NarStreamConfig::new().with_channel_capacity(100);
+ assert_eq!(config.channel_capacity, 100);
+ }
+
+ /// Test estimated memory calculation.
+ #[test]
+ fn test_memory_estimation() {
+ let config = NarStreamConfig::new().with_channel_capacity(128);
+ let chunk_size = 256 * 1024; // 256 KB
+ let estimated = config.estimated_max_memory_bytes(chunk_size);
+
+ // 128 chunks * 256 KB = 32 MB
+ assert_eq!(estimated, 128 * 256 * 1024);
+ assert_eq!(estimated, 33_554_432); // 32 MB
+ }
+}