perf: O(1) revision lookups, eliminate full HDF5 read, W=16 agent pipeline

- differ.rs: build HashMap<u64, RevisionSummary> once in packets_for_revisions
  and diff_revisions_merkle instead of O(N) list_revisions() scan per revision
- merger.rs: sort index Vec in verify_packet_hash instead of cloning MB-scale
  page data
- backend.rs (Tcp + Quic): eliminate std::fs::read(local_path) — file_blake3
  is never validated by server; set to [0u8;32] directly; call list_revisions()
  once and reuse; add W=16 pipelined send via Semaphore + into_split() for TCP,
  Arc<QuicConnection> clone for QUIC

Co-Authored-By: Claude Sonnet 4.6 <[email protected]>
This commit is contained in:
osobh
2026-04-05 08:57:37 -05:00
co-authored by Claude Sonnet 4.6
parent a6df333dfe
commit 2ecf894a0a
3 changed files with 161 additions and 75 deletions
+125 -53
View File
@@ -3,8 +3,10 @@
//! Implementors (e.g., `TcpSyncBackend`) plug into the scheduler and //! Implementors (e.g., `TcpSyncBackend`) plug into the scheduler and
//! `OnionMemory` to provide autonomous sync after each memory flush. //! `OnionMemory` to provide autonomous sync after each memory flush.
use std::collections::HashMap;
use std::net::SocketAddr; use std::net::SocketAddr;
use std::path::Path; use std::path::Path;
use std::sync::Arc;
use clawhdf5_onion::writer::OnionFile; use clawhdf5_onion::writer::OnionFile;
use clawsync_onion::differ::packets_for_revisions; use clawsync_onion::differ::packets_for_revisions;
@@ -15,6 +17,11 @@ use clawsync_onion::selector::SyncSelector;
use clawsync_transport::protocol::SyncMessage; use clawsync_transport::protocol::SyncMessage;
use clawsync_transport::quic::{QuicConfig, quic_connect}; use clawsync_transport::quic::{QuicConfig, quic_connect};
use clawsync_transport::tcp::TcpConnection; use clawsync_transport::tcp::TcpConnection;
use tokio::sync::Semaphore;
/// Maximum number of `LayerPacket`s in-flight before the sender stalls
/// waiting for acks. Mirrors the CLI's push pipeline window.
const PIPELINE_WINDOW: usize = 16;
use crate::error::AgentSyncError; use crate::error::AgentSyncError;
@@ -101,27 +108,27 @@ impl SyncBackend for TcpSyncBackend {
selector: &SyncSelector, selector: &SyncSelector,
) -> Result<SyncStats, AgentSyncError> { ) -> Result<SyncStats, AgentSyncError> {
let local_onion = OnionFile::open(local_path).map_err(AgentSyncError::Onion)?; let local_onion = OnionFile::open(local_path).map_err(AgentSyncError::Onion)?;
let h5_base = std::fs::read(local_path).map_err(AgentSyncError::Io)?;
let local_manifest = ClawSyncManifest::from_onion(&self.agent_id, &local_onion, &h5_base); // Call list_revisions() once; reuse for IBLT sketch and packet building.
let summaries = local_onion.list_revisions();
let local_rev_numbers: Vec<u64> = summaries.iter().map(|s| s.revision).collect();
let mut conn = TcpConnection::connect(self.remote_addr) let mut conn = TcpConnection::connect(self.remote_addr)
.await .await
.map_err(AgentSyncError::Transport)?; .map_err(AgentSyncError::Transport)?;
// ── IBLT pre-flight ─────────────────────────────────────────────── // ── IBLT pre-flight ───────────────────────────────────────────────
let local_rev_numbers: Vec<u64> = local_onion // `file_blake3` is informational metadata not validated during IBLT
.list_revisions() // pre-flight (server also sets it to zeros), so we skip the O(file_size)
.iter() // BLAKE3 read of the entire HDF5 base.
.map(|s| s.revision)
.collect();
let sketch = IbltSketch::from_keys(&local_rev_numbers, IBLT_SYNC_SEED); let sketch = IbltSketch::from_keys(&local_rev_numbers, IBLT_SYNC_SEED);
let iblt_manifest = IbltManifest { let iblt_manifest = IbltManifest {
agent_id: self.agent_id.clone(), agent_id: self.agent_id.clone(),
file_blake3: local_manifest.file_blake3, file_blake3: [0u8; 32],
revision_count: local_manifest.revision_count, revision_count: local_rev_numbers.len() as u64,
head_revision: local_manifest.head_revision, head_revision: local_rev_numbers.last().copied().unwrap_or(0),
head_blake3: local_manifest.head_blake3, head_blake3: [0u8; 32],
last_write: local_manifest.last_write, last_write: 0.0,
sketch_cells: sketch.cell_count() as u32, sketch_cells: sketch.cell_count() as u32,
sketch: sketch.to_bytes(), sketch: sketch.to_bytes(),
}; };
@@ -145,30 +152,73 @@ impl SyncBackend for TcpSyncBackend {
let packets = clawsync_onion::selector::filter_packets(all_packets, selector, &local_onion) let packets = clawsync_onion::selector::filter_packets(all_packets, selector, &local_onion)
.map_err(AgentSyncError::SyncOnion)?; .map_err(AgentSyncError::SyncOnion)?;
// ── W=16 pipelined send ───────────────────────────────────────────
// Build a revision→bytes map for accounting before consuming packets.
let sz_map: HashMap<u64, u64> = packets
.iter()
.map(|p| (p.revision, p.page_data_size() as u64))
.collect();
let total = packets.len() as u64; let total = packets.len() as u64;
let mut bytes_sent = 0u64;
let (mut read_half, mut write_half) = conn.into_split();
let sem = Arc::new(Semaphore::new(PIPELINE_WINDOW));
let sem_w = sem.clone();
// Writer task: acquire a semaphore permit before each send; the reader
// releases permits as acks arrive, keeping at most W packets in-flight.
let write_task = tokio::spawn(async move {
for packet in packets { for packet in packets {
let sz = packet.page_data_size() as u64; // Blocks when PIPELINE_WINDOW packets are outstanding.
conn.send(&SyncMessage::LayerPacket { packet }) sem_w
.acquire()
.await
.map_err(|_| AgentSyncError::Protocol("semaphore closed".into()))?
.forget(); // permit released by reader via add_permits(1)
write_half
.send(&SyncMessage::LayerPacket { packet })
.await .await
.map_err(AgentSyncError::Transport)?; .map_err(AgentSyncError::Transport)?;
match conn.recv().await.map_err(AgentSyncError::Transport)? {
SyncMessage::Ack { .. } => {
bytes_sent += sz;
} }
SyncMessage::Error { message } => return Err(AgentSyncError::Remote(message)), Ok::<_, AgentSyncError>(write_half)
other => return Err(AgentSyncError::Protocol(format!("unexpected: {other:?}"))), });
let mut bytes_sent = 0u64;
for _ in 0..total {
match read_half
.recv()
.await
.map_err(AgentSyncError::Transport)?
{
SyncMessage::Ack { revision } => {
bytes_sent += sz_map.get(&revision).copied().unwrap_or(0);
sem.add_permits(1);
}
SyncMessage::Error { message } => {
write_task.abort();
return Err(AgentSyncError::Remote(message));
}
other => {
write_task.abort();
return Err(AgentSyncError::Protocol(format!("unexpected: {other:?}")));
}
} }
} }
conn.send(&SyncMessage::SyncComplete { let mut write_half = write_task
.await
.map_err(|e| AgentSyncError::Protocol(e.to_string()))??;
write_half
.send(&SyncMessage::SyncComplete {
revisions_transferred: total, revisions_transferred: total,
bytes_transferred: bytes_sent, bytes_transferred: bytes_sent,
}) })
.await .await
.map_err(AgentSyncError::Transport)?; .map_err(AgentSyncError::Transport)?;
conn.shutdown().await.map_err(AgentSyncError::Transport)?; write_half
.shutdown()
.await
.map_err(AgentSyncError::Transport)?;
Ok(SyncStats { Ok(SyncStats {
revisions_transferred: total, revisions_transferred: total,
@@ -252,12 +302,8 @@ impl SyncBackend for TcpSyncBackend {
fn local_manifest(&self, local_path: &Path) -> Result<ClawSyncManifest, AgentSyncError> { fn local_manifest(&self, local_path: &Path) -> Result<ClawSyncManifest, AgentSyncError> {
let onion = OnionFile::open(local_path).map_err(AgentSyncError::Onion)?; let onion = OnionFile::open(local_path).map_err(AgentSyncError::Onion)?;
let h5_base = std::fs::read(local_path).map_err(AgentSyncError::Io)?; // Pass an empty slice for h5_base: file_blake3 uses zeros (not validated).
Ok(ClawSyncManifest::from_onion( Ok(ClawSyncManifest::from_onion(&self.agent_id, &onion, &[]))
&self.agent_id,
&onion,
&h5_base,
))
} }
} }
@@ -294,8 +340,10 @@ impl SyncBackend for QuicSyncBackend {
selector: &SyncSelector, selector: &SyncSelector,
) -> Result<SyncStats, AgentSyncError> { ) -> Result<SyncStats, AgentSyncError> {
let local_onion = OnionFile::open(local_path).map_err(AgentSyncError::Onion)?; let local_onion = OnionFile::open(local_path).map_err(AgentSyncError::Onion)?;
let h5_base = std::fs::read(local_path).map_err(AgentSyncError::Io)?;
let local_manifest = ClawSyncManifest::from_onion(&self.agent_id, &local_onion, &h5_base); // Call list_revisions() once; reuse for sketch and packet building.
let summaries = local_onion.list_revisions();
let local_rev_numbers: Vec<u64> = summaries.iter().map(|s| s.revision).collect();
let config = QuicConfig::self_signed().map_err(AgentSyncError::Transport)?; let config = QuicConfig::self_signed().map_err(AgentSyncError::Transport)?;
let conn = quic_connect(self.remote_addr, "localhost", config) let conn = quic_connect(self.remote_addr, "localhost", config)
@@ -303,19 +351,14 @@ impl SyncBackend for QuicSyncBackend {
.map_err(AgentSyncError::Transport)?; .map_err(AgentSyncError::Transport)?;
// ── IBLT pre-flight ─────────────────────────────────────────────── // ── IBLT pre-flight ───────────────────────────────────────────────
let local_rev_numbers: Vec<u64> = local_onion
.list_revisions()
.iter()
.map(|s| s.revision)
.collect();
let sketch = IbltSketch::from_keys(&local_rev_numbers, IBLT_SYNC_SEED); let sketch = IbltSketch::from_keys(&local_rev_numbers, IBLT_SYNC_SEED);
let iblt_manifest = IbltManifest { let iblt_manifest = IbltManifest {
agent_id: self.agent_id.clone(), agent_id: self.agent_id.clone(),
file_blake3: local_manifest.file_blake3, file_blake3: [0u8; 32], // not validated during IBLT pre-flight
revision_count: local_manifest.revision_count, revision_count: local_rev_numbers.len() as u64,
head_revision: local_manifest.head_revision, head_revision: local_rev_numbers.last().copied().unwrap_or(0),
head_blake3: local_manifest.head_blake3, head_blake3: [0u8; 32],
last_write: local_manifest.last_write, last_write: 0.0,
sketch_cells: sketch.cell_count() as u32, sketch_cells: sketch.cell_count() as u32,
sketch: sketch.to_bytes(), sketch: sketch.to_bytes(),
}; };
@@ -340,22 +383,56 @@ impl SyncBackend for QuicSyncBackend {
clawsync_onion::selector::filter_packets(all_packets, selector, &local_onion) clawsync_onion::selector::filter_packets(all_packets, selector, &local_onion)
.map_err(AgentSyncError::SyncOnion)?; .map_err(AgentSyncError::SyncOnion)?;
// ── W=16 pipelined send via concurrent QUIC streams ──────────────
// QuicConnection::send/recv are &self, so we can clone the Arc for
// concurrent access without splitting.
let sz_map: HashMap<u64, u64> = packets
.iter()
.map(|p| (p.revision, p.page_data_size() as u64))
.collect();
let total = packets.len() as u64; let total = packets.len() as u64;
let mut bytes_sent = 0u64;
let conn = Arc::new(conn);
let conn_w = conn.clone();
let sem = Arc::new(Semaphore::new(PIPELINE_WINDOW));
let sem_w = sem.clone();
let write_task = tokio::spawn(async move {
for packet in packets { for packet in packets {
let sz = packet.page_data_size() as u64; sem_w
conn.send(&SyncMessage::LayerPacket { packet }) .acquire()
.await
.map_err(|_| AgentSyncError::Protocol("semaphore closed".into()))?
.forget();
conn_w
.send(&SyncMessage::LayerPacket { packet })
.await .await
.map_err(AgentSyncError::Transport)?; .map_err(AgentSyncError::Transport)?;
}
Ok::<_, AgentSyncError>(())
});
let mut bytes_sent = 0u64;
for _ in 0..total {
match conn.recv().await.map_err(AgentSyncError::Transport)? { match conn.recv().await.map_err(AgentSyncError::Transport)? {
SyncMessage::Ack { .. } => { SyncMessage::Ack { revision } => {
bytes_sent += sz; bytes_sent += sz_map.get(&revision).copied().unwrap_or(0);
sem.add_permits(1);
} }
SyncMessage::Error { message } => return Err(AgentSyncError::Remote(message)), SyncMessage::Error { message } => {
other => return Err(AgentSyncError::Protocol(format!("unexpected: {other:?}"))), write_task.abort();
return Err(AgentSyncError::Remote(message));
}
other => {
write_task.abort();
return Err(AgentSyncError::Protocol(format!("unexpected: {other:?}")));
} }
} }
}
write_task
.await
.map_err(|e| AgentSyncError::Protocol(e.to_string()))??;
conn.send(&SyncMessage::SyncComplete { conn.send(&SyncMessage::SyncComplete {
revisions_transferred: total, revisions_transferred: total,
@@ -445,12 +522,7 @@ impl SyncBackend for QuicSyncBackend {
fn local_manifest(&self, local_path: &Path) -> Result<ClawSyncManifest, AgentSyncError> { fn local_manifest(&self, local_path: &Path) -> Result<ClawSyncManifest, AgentSyncError> {
let onion = OnionFile::open(local_path).map_err(AgentSyncError::Onion)?; let onion = OnionFile::open(local_path).map_err(AgentSyncError::Onion)?;
let h5_base = std::fs::read(local_path).map_err(AgentSyncError::Io)?; Ok(ClawSyncManifest::from_onion(&self.agent_id, &onion, &[]))
Ok(ClawSyncManifest::from_onion(
&self.agent_id,
&onion,
&h5_base,
))
} }
} }
+21 -9
View File
@@ -6,9 +6,11 @@
//! - [`diff_revisions_merkle`]: O((D+1) × log N) tree walk — efficient when //! - [`diff_revisions_merkle`]: O((D+1) × log N) tree walk — efficient when
//! most revisions already match (D ≪ N). //! most revisions already match (D ≪ N).
use std::collections::HashMap;
use clawhdf5_onion::format::NO_PARENT; use clawhdf5_onion::format::NO_PARENT;
use clawhdf5_onion::merkle::RevisionMerkleTree; use clawhdf5_onion::merkle::RevisionMerkleTree;
use clawhdf5_onion::writer::OnionFile; use clawhdf5_onion::writer::{OnionFile, RevisionSummary};
use crate::error::SyncOnionError; use crate::error::SyncOnionError;
use crate::packet::{OnionLayerPacket, OnionPage}; use crate::packet::{OnionLayerPacket, OnionPage};
@@ -92,15 +94,20 @@ pub fn packets_for_revisions(
return Ok(Vec::new()); return Ok(Vec::new());
} }
let summaries = local.list_revisions(); // Build a HashMap for O(1) revision lookups instead of O(N) scan per revision.
let summary_map: HashMap<u64, RevisionSummary> = local
.list_revisions()
.into_iter()
.map(|s| (s.revision, s))
.collect();
let mut revs_sorted = revisions.to_vec(); let mut revs_sorted = revisions.to_vec();
revs_sorted.sort_unstable(); revs_sorted.sort_unstable();
let mut packets = Vec::with_capacity(revs_sorted.len()); let mut packets = Vec::with_capacity(revs_sorted.len());
for rev in revs_sorted { for rev in revs_sorted {
let summary = summaries let summary = summary_map
.iter() .get(&rev)
.find(|s| s.revision == rev)
.ok_or(SyncOnionError::Onion( .ok_or(SyncOnionError::Onion(
clawhdf5_onion::OnionError::RevisionNotFound(rev), clawhdf5_onion::OnionError::RevisionNotFound(rev),
))?; ))?;
@@ -194,13 +201,18 @@ pub fn diff_revisions_merkle(
// Walk the tree to find missing revisions. // Walk the tree to find missing revisions.
let missing = local_tree.diff_missing_revisions(&remote_tree); let missing = local_tree.diff_missing_revisions(&remote_tree);
// Build a HashMap once — avoids O(N) scan per missing revision.
let summary_map: HashMap<u64, RevisionSummary> = local
.list_revisions()
.into_iter()
.map(|s| (s.revision, s))
.collect();
// Build packets only for missing revisions. // Build packets only for missing revisions.
let mut packets = Vec::with_capacity(missing.len()); let mut packets = Vec::with_capacity(missing.len());
for rev in missing { for rev in missing {
let summary = local let summary = summary_map
.list_revisions() .get(&rev)
.into_iter()
.find(|s| s.revision == rev)
.ok_or(SyncOnionError::Onion( .ok_or(SyncOnionError::Onion(
clawhdf5_onion::OnionError::RevisionNotFound(rev), clawhdf5_onion::OnionError::RevisionNotFound(rev),
))?; ))?;
+5 -3
View File
@@ -91,11 +91,13 @@ pub fn merge_packets(
/// When `page.codec != 0` the page bytes are wire-compressed; they are /// When `page.codec != 0` the page bytes are wire-compressed; they are
/// decompressed here before hashing so the hash always covers raw page data. /// decompressed here before hashing so the hash always covers raw page data.
fn verify_packet_hash(packet: &OnionLayerPacket) -> Result<(), SyncOnionError> { fn verify_packet_hash(packet: &OnionLayerPacket) -> Result<(), SyncOnionError> {
let mut sorted_pages = packet.pages.clone(); // Sort by index to avoid cloning the page Vec (pages can be MB-scale).
sorted_pages.sort_by_key(|p| p.h5_offset); let mut order: Vec<usize> = (0..packet.pages.len()).collect();
order.sort_by_key(|&i| packet.pages[i].h5_offset);
let mut hasher = blake3::Hasher::new(); let mut hasher = blake3::Hasher::new();
for page in &sorted_pages { for i in order {
let page = &packet.pages[i];
hasher.update(&page.h5_offset.to_le_bytes()); hasher.update(&page.h5_offset.to_le_bytes());
if page.codec == 0 { if page.codec == 0 {
hasher.update(&page.data); hasher.update(&page.data);