From ae3c068069f8d5f43f6897193b230d08803effbc Mon Sep 17 00:00:00 2001 From: phil Date: Thu, 22 Jan 2026 14:14:23 -0500 Subject: [PATCH] further deps simplification --- Cargo.toml | 4 +- readme.md | 4 ++ src/async_io.rs | 8 +-- src/blocking.rs | 39 +++++++++----- src/parser.rs | 137 +++++++++++++++++++++++++++++------------------- src/ser.rs | 40 +++++++------- src/tests.rs | 44 +++++++++++----- 7 files changed, 173 insertions(+), 103 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index ceadbfd..e10373e 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,7 +6,7 @@ edition = "2024" [features] default = ["async"] blocking = [] -async = ["dep:tokio", "dep:tokio-util", "dep:futures"] +async = ["dep:tokio", "dep:tokio-util", "dep:futures", "dep:bytes"] [dependencies] # Core @@ -16,10 +16,10 @@ serde_bytes = "0.11" cid = { version = "0.11", features = ["serde-codec"] } thiserror = "2" sha2 = "0.10" -bytes = "1" unsigned-varint = "0.8" # Async +bytes = { version = "1", optional = true } tokio = { version = "1", features = ["io-util", "macros", "rt"], optional = true } tokio-util = { version = "0.7", features = ["codec"], optional = true } futures = { version = "0.3", optional = true } diff --git a/readme.md b/readme.md index badb908..d2bdb0a 100644 --- a/readme.md +++ b/readme.md @@ -185,3 +185,7 @@ with the 50% initial reduction, this would be **95% total CID reduction** CIDs make up around 20% of uncompressed CAR file sizes. the first approach gets that down to 4%; second 1%. however, CIDs are uncompressible, so it's probably worth measuring the real effect of both approaches on large repos post-compression before completely committing one way or another. + + + +TODO: note use of multiformats varint instead of LEB128 (stricter) \ No newline at end of file diff --git a/src/async_io.rs b/src/async_io.rs index 43c11e2..690aa19 100644 --- a/src/async_io.rs +++ b/src/async_io.rs @@ -1,7 +1,7 @@ -use crate::error::StarError; use crate::parser::StarParser; use crate::types::StarItem; -use bytes::BytesMut; +use crate::error::StarError; +use bytes::{BytesMut, Buf}; use tokio_util::codec::Decoder; impl Decoder for StarParser { @@ -9,6 +9,8 @@ impl Decoder for StarParser { type Error = StarError; fn decode(&mut self, src: &mut BytesMut) -> Result, Self::Error> { - self.input(src) + let (consumed, item) = self.parse(src)?; + src.advance(consumed); + Ok(item) } } diff --git a/src/blocking.rs b/src/blocking.rs index e2ef727..c648061 100644 --- a/src/blocking.rs +++ b/src/blocking.rs @@ -1,14 +1,14 @@ -use crate::error::Result; +use std::io::Read; use crate::parser::StarParser; +use crate::error::Result; use crate::types::StarItem; -use bytes::BytesMut; use cid::Cid; -use std::io::Read; +use std::collections::VecDeque; pub struct StarIterator { reader: R, parser: StarParser, - buffer: BytesMut, + buffer: VecDeque, } impl StarIterator { @@ -16,7 +16,7 @@ impl StarIterator { Self { reader, parser: StarParser::new(), - buffer: BytesMut::new(), + buffer: VecDeque::new(), } } @@ -51,19 +51,32 @@ impl Iterator for StarIterator { fn next(&mut self) -> Option { loop { - match self.parser.input(&mut self.buffer) { - Ok(Some(item)) => return Some(Ok(item)), - Ok(None) => { - if self.buffer.capacity() == 0 { - self.buffer.reserve(8192); - } + // VecDeque must be contiguous to be parsed as a slice. + // make_contiguous() allows us to get a single slice, but might involve a memory copy/move + // if the buffer is wrapped. This is amortized O(1) in many cases but can be O(N). + // However, it solves the "head removal" problem perfectly (pop_front is O(1)). + + self.buffer.make_contiguous(); + let (slice, _) = self.buffer.as_slices(); + + match self.parser.parse(slice) { + Ok((consumed, Some(item))) => { + self.buffer.drain(..consumed); + return Some(Ok(item)); + } + Ok((consumed, None)) => { + self.buffer.drain(..consumed); + + // Read more + // We can't read directly into VecDeque easily without temporary buffer + // because it doesn't expose a mutable slice to uninitialized memory. let mut temp = [0u8; 8192]; match self.reader.read(&mut temp) { Ok(0) => return None, // EOF - Ok(n) => self.buffer.extend_from_slice(&temp[..n]), + Ok(n) => self.buffer.extend(&temp[..n]), Err(e) => return Some(Err(e.into())), } - } + }, Err(e) => return Some(Err(e)), } } diff --git a/src/parser.rs b/src/parser.rs index 513a3c4..305a350 100644 --- a/src/parser.rs +++ b/src/parser.rs @@ -1,6 +1,5 @@ use crate::error::{Result, StarError}; use crate::types::{StarCommit, StarItem, StarMstNode}; -use bytes::{Buf, BytesMut}; use cid::Cid; use sha2::{Digest, Sha256}; @@ -42,7 +41,13 @@ impl StarParser { } } - pub fn input(&mut self, buf: &mut BytesMut) -> Result> { + /// Parses the input buffer. + /// Returns (bytes_consumed, Option). + pub fn parse(&mut self, buf: &[u8]) -> Result<(usize, Option)> { + let mut consumed = 0; + + // Loop allows state transitions (e.g. Header -> Body) without returning + // but we must be careful to track consumed bytes correctly. loop { // Check if we need to transition from Body to Done let is_body_done = if let State::Body { stack, .. } = &self.state { @@ -53,42 +58,74 @@ impl StarParser { if is_body_done { self.state = State::Done; - return Ok(None); + return Ok((consumed, None)); } + // Slice the buffer to the remaining part + let current_buf = &buf[consumed..]; + match &mut self.state { - State::Done => return Ok(None), + State::Done => return Ok((consumed, None)), State::Header => { - return self.parse_header(buf) + let (n, item) = match self.parse_header(current_buf)? { + Some((n, item)) => (n, item), + None => return Ok((consumed, None)), + }; + consumed += n; + return Ok((consumed, Some(item))); } State::Body { stack, current_len } => { if Self::process_verification(stack)? { + // Verification doesn't consume bytes, but changes state/stack. + // We continue the loop to try reading the next item immediately. continue; } - if let Some(len) = Self::read_length(buf, current_len)? { - let block_bytes = buf.split_to(len); - *current_len = None; - - let item = stack.pop().unwrap(); - match item { - StackItem::Node { expected } => { - return Self::process_node(block_bytes, expected, stack); - }, - StackItem::Record { key, expected, implicit_index } => { - return Self::process_record(block_bytes, key, expected, implicit_index, stack); - }, - _ => return Err(StarError::InvalidState("Unexpected stack item".into())), - } - } else { - return Ok(None); + // Try to read length + let (len_consumed, len) = match Self::read_length(current_buf, current_len)? { + Some((n, len)) => (n, len), + None => return Ok((consumed, None)), + }; + + // Note: read_length advances internal state (current_len) but also returns + // how many bytes of the varint were consumed from current_buf. + // If we have the full length, we now check if we have the body. + + let body_buf = ¤t_buf[len_consumed..]; + if body_buf.len() < len { + // We read the length varint, but don't have enough bytes for the body. + // We must report the varint itself as consumed so the caller advances, + // and we have stored the in via mutation. + consumed += len_consumed; + return Ok((consumed, None)); } + + // We have the body. + let block_bytes = &body_buf[..len]; + + // Reset current_len since we are consuming the block + *current_len = None; + + let item = stack.pop().unwrap(); + let result_item = match item { + StackItem::Node { expected } => { + Self::process_node(block_bytes, expected, stack)? + }, + StackItem::Record { key, expected, implicit_index } => { + Self::process_record(block_bytes, key, expected, implicit_index, stack)? + }, + _ => return Err(StarError::InvalidState("Unexpected stack item".into())), + }; + + consumed += len_consumed + len; + return Ok((consumed, result_item)); } } } } - fn parse_header(&mut self, buf: &mut BytesMut) -> Result> { + // Returns Option<(bytes_consumed, Item)> + fn parse_header(&mut self, buf: &[u8]) -> Result> { if buf.len() < 1 { return Ok(None); } @@ -96,8 +133,6 @@ impl StarParser { return Err(StarError::InvalidHeader); } - // We use a slice of the buffer starting after the magic byte - // unsigned_varint operates on &[u8] let slice = &buf[1..]; let (ver, remaining1) = match unsigned_varint::decode::usize(slice) { @@ -112,20 +147,19 @@ impl StarParser { Err(e) => return Err(StarError::InvalidState(format!("Varint error: {}", e))), }; - let header_varints_len = buf.len() - 1 - remaining2.len(); // bytes consumed by varints - let total_header_len = 1 + header_varints_len; // + magic byte + let header_varints_len = buf.len() - 1 - remaining2.len(); + let total_header_len = 1 + header_varints_len; let total_len = total_header_len + len; if buf.len() < total_len { return Ok(None); } - buf.advance(total_header_len); - let commit_bytes = buf.split_to(len); - let commit: StarCommit = serde_ipld_dagcbor::from_slice(&commit_bytes) + // We have the full commit + let commit_bytes = &buf[total_header_len..total_len]; + let commit: StarCommit = serde_ipld_dagcbor::from_slice(commit_bytes) .map_err(|e| StarError::Cbor(e.to_string()))?; - // Check version (conceptually) let _ = ver; let mut stack = Vec::new(); @@ -139,7 +173,8 @@ impl StarParser { stack, current_len: None }; - Ok(Some(StarItem::Commit(commit))) + + Ok(Some((total_len, StarItem::Commit(commit)))) } fn process_verification(stack: &mut Vec) -> Result { @@ -177,31 +212,29 @@ impl StarParser { Ok(false) } - fn read_length(buf: &mut BytesMut, current_len: &mut Option) -> Result> { - if current_len.is_none() { - match unsigned_varint::decode::usize(&buf[..]) { - Ok((l, remaining)) => { - let consumed = buf.len() - remaining.len(); - *current_len = Some(l); - buf.advance(consumed); - }, - Err(unsigned_varint::decode::Error::Insufficient) => return Ok(None), - Err(e) => return Err(StarError::InvalidState(format!("Varint error: {}", e))), - } + // Returns Option<(bytes_consumed, length_value)> + // If successful, updates current_len to the read length value + fn read_length(buf: &[u8], current_len: &mut Option) -> Result> { + // If we already have a length (from a previous partial read), we consumed 0 bytes *now* to get it + if let Some(len) = current_len { + return Ok(Some((0, *len))); } - let len = current_len.unwrap(); - if buf.len() < len { - return Ok(None); + match unsigned_varint::decode::usize(buf) { + Ok((l, remaining)) => { + let consumed = buf.len() - remaining.len(); + *current_len = Some(l); + Ok(Some((consumed, l))) + }, + Err(unsigned_varint::decode::Error::Insufficient) => Ok(None), + Err(e) => Err(StarError::InvalidState(format!("Varint error: {}", e))), } - Ok(Some(len)) } - fn process_node(block_bytes: BytesMut, expected: Option, stack: &mut Vec) -> Result> { - let node: StarMstNode = serde_ipld_dagcbor::from_slice(&block_bytes) + fn process_node(block_bytes: &[u8], expected: Option, stack: &mut Vec) -> Result> { + let node: StarMstNode = serde_ipld_dagcbor::from_slice(block_bytes) .map_err(|e| StarError::Cbor(e.to_string()))?; - // Check for implicit records let mut has_implicit = false; for e in &node.e { if e.v_archived == Some(true) && e.v.is_none() { @@ -236,7 +269,6 @@ impl StarParser { }); } - // Reconstruct keys let mut prev_key_bytes = Vec::new(); let mut entry_keys = Vec::new(); for e in &node.e { @@ -250,7 +282,6 @@ impl StarParser { prev_key_bytes = key; } - // Push children in reverse for i in (0..node.e.len()).rev() { let e = &node.e[i]; let key = entry_keys[i].clone(); @@ -276,8 +307,8 @@ impl StarParser { Ok(Some(StarItem::Node(node))) } - fn process_record(block_bytes: BytesMut, key: Vec, expected: Option, implicit_index: Option, stack: &mut Vec) -> Result> { - let hash = Sha256::digest(&block_bytes); + fn process_record(block_bytes: &[u8], key: Vec, expected: Option, implicit_index: Option, stack: &mut Vec) -> Result> { + let hash = Sha256::digest(block_bytes); let cid = Cid::new_v1(0x71, cid::multihash::Multihash::wrap(0x12, &hash)?); if let Some(exp) = expected { diff --git a/src/ser.rs b/src/ser.rs index 7d30aa5..056a708 100644 --- a/src/ser.rs +++ b/src/ser.rs @@ -1,43 +1,43 @@ use crate::error::Result; use crate::types::{StarCommit, StarItem, StarMstNode}; -use bytes::{BufMut, BytesMut}; +use std::io::Write; pub struct StarEncoder; impl StarEncoder { - fn write_varint(val: usize, dst: &mut BytesMut) { + fn write_varint(val: usize, dst: &mut W) -> std::io::Result<()> { let mut buf = unsigned_varint::encode::usize_buffer(); let encoded = unsigned_varint::encode::usize(val, &mut buf); - dst.extend_from_slice(encoded); + dst.write_all(encoded) } - pub fn write_header(commit: &StarCommit, dst: &mut BytesMut) -> Result<()> { - dst.put_u8(0x2A); + pub fn write_header(commit: &StarCommit, dst: &mut W) -> Result<()> { + dst.write_all(&[0x2A])?; - Self::write_varint(1, dst); + Self::write_varint(1, dst)?; let commit_bytes = serde_ipld_dagcbor::to_vec(commit) .map_err(|e| crate::error::StarError::Cbor(e.to_string()))?; - Self::write_varint(commit_bytes.len(), dst); - dst.extend_from_slice(&commit_bytes); + Self::write_varint(commit_bytes.len(), dst)?; + dst.write_all(&commit_bytes)?; Ok(()) } - pub fn write_node(node: &StarMstNode, dst: &mut BytesMut) -> Result<()> { + pub fn write_node(node: &StarMstNode, dst: &mut W) -> Result<()> { let node_bytes = serde_ipld_dagcbor::to_vec(node) .map_err(|e| crate::error::StarError::Cbor(e.to_string()))?; - Self::write_varint(node_bytes.len(), dst); - dst.extend_from_slice(&node_bytes); + Self::write_varint(node_bytes.len(), dst)?; + dst.write_all(&node_bytes)?; Ok(()) } - pub fn write_record(record_bytes: &[u8], dst: &mut BytesMut) -> Result<()> { - Self::write_varint(record_bytes.len(), dst); - dst.extend_from_slice(record_bytes); + pub fn write_record(record_bytes: &[u8], dst: &mut W) -> Result<()> { + Self::write_varint(record_bytes.len(), dst)?; + dst.write_all(record_bytes)?; Ok(()) } @@ -47,13 +47,17 @@ impl StarEncoder { impl tokio_util::codec::Encoder for StarEncoder { type Error = crate::error::StarError; - fn encode(&mut self, item: StarItem, dst: &mut BytesMut) -> Result<()> { + fn encode(&mut self, item: StarItem, dst: &mut bytes::BytesMut) -> Result<()> { + use bytes::BufMut; + // BytesMut::writer() returns an impl Write + let mut writer = dst.writer(); + match item { - StarItem::Commit(c) => Self::write_header(&c, dst), - StarItem::Node(n) => Self::write_node(&n, dst), + StarItem::Commit(c) => Self::write_header(&c, &mut writer), + StarItem::Node(n) => Self::write_node(&n, &mut writer), StarItem::Record { content, .. } => { if let Some(bytes) = content { - Self::write_record(&bytes, dst) + Self::write_record(&bytes, &mut writer) } else { Err(crate::error::StarError::InvalidState("Cannot serialize record without content".into())) } diff --git a/src/tests.rs b/src/tests.rs index 13d6376..86801ff 100644 --- a/src/tests.rs +++ b/src/tests.rs @@ -5,7 +5,6 @@ mod tests { use crate::types::{ RepoMstEntry, RepoMstNode, StarCommit, StarItem, StarMstEntry, StarMstNode, }; - use bytes::BytesMut; use cid::Cid; use serde_bytes::ByteBuf; use sha2::{Digest, Sha256}; @@ -63,7 +62,7 @@ mod tests { }; // 5. Serialize to Buffer - let mut buf = BytesMut::new(); + let mut buf = Vec::new(); StarEncoder::write_header(&commit, &mut buf).unwrap(); StarEncoder::write_node(&star_node, &mut buf).unwrap(); StarEncoder::write_record(record_data, &mut buf).unwrap(); @@ -71,15 +70,24 @@ mod tests { // 6. Deserialize and Verify let mut parser = StarParser::new(); + // Helper to mimic Decoder loop for tests + fn parse_helper(parser: &mut StarParser, buf: &mut Vec, offset: &mut usize) -> Option { + let (consumed, item) = parser.parse(&buf[*offset..]).unwrap(); + *offset += consumed; + item + } + + let mut offset = 0; + // Header - let item1 = parser.input(&mut buf).unwrap().unwrap(); + let item1 = parse_helper(&mut parser, &mut buf, &mut offset).unwrap(); match item1 { StarItem::Commit(c) => assert_eq!(c, commit), _ => panic!("Expected commit"), } // Node - let item2 = parser.input(&mut buf).unwrap().unwrap(); + let item2 = parse_helper(&mut parser, &mut buf, &mut offset).unwrap(); match item2 { StarItem::Node(n) => { assert_eq!(n.e[0].v, None); @@ -89,7 +97,7 @@ mod tests { } // Record - let item3 = parser.input(&mut buf).unwrap().unwrap(); + let item3 = parse_helper(&mut parser, &mut buf, &mut offset).unwrap(); match item3 { StarItem::Record { key: k, @@ -104,7 +112,7 @@ mod tests { } // Done - assert!(parser.input(&mut buf).unwrap().is_none()); + assert!(parse_helper(&mut parser, &mut buf, &mut offset).is_none()); } #[test] @@ -140,23 +148,31 @@ mod tests { }; // 4. Serialize - let mut buf = BytesMut::new(); + let mut buf = Vec::new(); StarEncoder::write_header(&commit, &mut buf).unwrap(); StarEncoder::write_node(&star_node, &mut buf).unwrap(); StarEncoder::write_record(record_data, &mut buf).unwrap(); // 5. Parse let mut parser = StarParser::new(); - parser.input(&mut buf).unwrap(); // Header OK - parser.input(&mut buf).unwrap(); // Node OK (verification deferred) - parser.input(&mut buf).unwrap(); // Record OK (verification scheduled) - - // 6. Trigger verification (processing VerifyLayer0) - let result = parser.input(&mut buf); + let mut offset = 0; + + fn parse_helper(parser: &mut StarParser, buf: &mut Vec, offset: &mut usize) -> Result, crate::error::StarError> { + let (consumed, item) = parser.parse(&buf[*offset..])?; + *offset += consumed; + Ok(item) + } + parse_helper(&mut parser, &mut buf, &mut offset).unwrap(); // Header OK + parse_helper(&mut parser, &mut buf, &mut offset).unwrap(); // Node OK + parse_helper(&mut parser, &mut buf, &mut offset).unwrap(); // Record OK + + // 6. Trigger verification + let result = parse_helper(&mut parser, &mut buf, &mut offset); + assert!(result.is_err()); match result.unwrap_err() { - crate::error::StarError::VerificationFailed { .. } => {} + crate::error::StarError::VerificationFailed { .. } => {}, e => panic!("Expected VerificationFailed, got {:?}", e), } } -- 2.51.2