mod delta_search;
mod header;
mod parallel;
mod sort;
pub mod output;
pub use output::encode_and_output_to_files;
#[cfg(test)]
mod tests;
use header::encode_header;
use rayon::prelude::*;
use sort::magic_sort;
use tokio::{sync::mpsc, task::JoinHandle};
use crate::{
errors::GitError,
hash::ObjectHash,
internal::{
metadata::{EntryMeta, MetaAttached},
object::types::ObjectType,
pack::{entry::Entry, index_entry::IndexEntry, pack_index::IdxBuilder},
},
utils::HashAlgorithm,
};
pub struct PackEncoder {
object_number: usize,
process_index: usize,
window_size: usize,
pack_sender: Option<mpsc::Sender<Vec<u8>>>,
idx_sender: Option<mpsc::Sender<Vec<u8>>>,
idx_entries: Option<Vec<IndexEntry>>,
inner_offset: usize,
inner_hash: HashAlgorithm,
final_hash: Option<ObjectHash>,
start_encoding: bool,
pub disable_prefilter: bool,
}
impl PackEncoder {
pub fn new(object_number: usize, window_size: usize, sender: mpsc::Sender<Vec<u8>>) -> Self {
PackEncoder {
object_number,
window_size,
process_index: 0,
pack_sender: Some(sender),
idx_sender: None,
idx_entries: None,
inner_offset: 12, inner_hash: HashAlgorithm::new(),
final_hash: None,
start_encoding: false,
disable_prefilter: false,
}
}
pub fn new_with_idx(
object_number: usize,
window_size: usize,
pack_sender: mpsc::Sender<Vec<u8>>,
idx_sender: mpsc::Sender<Vec<u8>>,
) -> Self {
PackEncoder {
object_number,
window_size,
process_index: 0,
pack_sender: Some(pack_sender),
idx_sender: Some(idx_sender),
idx_entries: None,
inner_offset: 12, inner_hash: HashAlgorithm::new(),
final_hash: None,
start_encoding: false,
disable_prefilter: false,
}
}
pub fn drop_sender(&mut self) {
self.pack_sender.take();
}
pub async fn send_data(&mut self, data: Vec<u8>) {
if let Some(sender) = &self.pack_sender {
sender.send(data).await.unwrap();
}
}
pub fn get_hash(&self) -> Option<ObjectHash> {
self.final_hash
}
pub async fn encode(
&mut self,
entry_rx: mpsc::Receiver<MetaAttached<Entry, EntryMeta>>,
) -> Result<(), GitError> {
if self.window_size == 0 {
self.parallel_encode(entry_rx).await
} else {
#[cfg(feature = "diff_rabin")]
{
self.inner_encode(entry_rx, false, true, true).await
}
#[cfg(not(feature = "diff_rabin"))]
{
self.inner_encode(entry_rx, false, false, self.disable_prefilter)
.await
}
}
}
pub async fn encode_with_zstdelta(
&mut self,
entry_rx: mpsc::Receiver<MetaAttached<Entry, EntryMeta>>,
) -> Result<(), GitError> {
self.inner_encode(entry_rx, true, false, self.disable_prefilter)
.await
}
async fn inner_encode(
&mut self,
mut entry_rx: mpsc::Receiver<MetaAttached<Entry, EntryMeta>>,
enable_zstdelta: bool,
enable_rabin: bool,
disable_prefilter: bool,
) -> Result<(), GitError> {
let head = encode_header(self.object_number);
self.send_data(head.clone()).await;
self.inner_hash.update(&head);
if self.start_encoding {
return Err(GitError::PackEncodeError(
"encoding operation is already in progress".to_string(),
));
}
let mut commits: Vec<MetaAttached<Entry, EntryMeta>> = Vec::new();
let mut trees: Vec<MetaAttached<Entry, EntryMeta>> = Vec::new();
let mut blobs: Vec<MetaAttached<Entry, EntryMeta>> = Vec::new();
let mut tags: Vec<MetaAttached<Entry, EntryMeta>> = Vec::new();
while let Some(entry) = entry_rx.recv().await {
match entry.inner.obj_type {
ObjectType::Commit => {
commits.push(entry);
}
ObjectType::Tree => {
trees.push(entry);
}
ObjectType::Blob => {
blobs.push(entry);
}
ObjectType::Tag => {
tags.push(entry);
}
_ => {
return Err(GitError::PackEncodeError(format!(
"object type `{}` is not supported by delta-window pack encoding",
entry.inner.obj_type
)));
}
}
}
commits.sort_by(magic_sort);
trees.sort_by(magic_sort);
blobs.sort_by(magic_sort);
tags.sort_by(magic_sort);
tracing::info!(
"numbers : commits: {:?} trees: {:?} blobs:{:?} tag :{:?}",
commits.len(),
trees.len(),
blobs.len(),
tags.len()
);
let commit_entries: Vec<Entry> = commits.into_iter().map(|e| e.inner).collect();
let tree_entries: Vec<Entry> = trees.into_iter().map(|e| e.inner).collect();
let tag_entries: Vec<Entry> = tags.into_iter().map(|e| e.inner).collect();
let mut blob_entries: Vec<Entry> = blobs.into_iter().map(|e| e.inner).collect();
struct WorkItem {
order: usize,
entries: Vec<Entry>,
}
let mut work_items: Vec<WorkItem> = Vec::new();
work_items.push(WorkItem {
order: 0,
entries: commit_entries,
});
work_items.push(WorkItem {
order: 1,
entries: tree_entries,
});
let total_blob_entries = blob_entries.len();
let num_threads = rayon::current_num_threads();
let chunks_per_thread: usize = 20;
let mut blob_chunk_count = if num_threads > 1 && total_blob_entries > (num_threads * 20) {
num_threads * chunks_per_thread
} else {
1
};
let mut entries_per_chunk = total_blob_entries.div_ceil(blob_chunk_count);
let min_entries_per_chunk = (self.window_size * 10).max(50);
if entries_per_chunk < min_entries_per_chunk {
blob_chunk_count = (total_blob_entries / min_entries_per_chunk).max(1);
entries_per_chunk = total_blob_entries.div_ceil(blob_chunk_count);
}
let mut blob_chunks: Vec<Vec<Entry>> = Vec::with_capacity(blob_chunk_count);
for _ in 0..blob_chunk_count {
let take = entries_per_chunk.min(blob_entries.len());
if take == 0 {
break;
}
let chunk = blob_entries.split_off(blob_entries.len() - take);
blob_chunks.push(chunk);
}
blob_chunks.reverse();
let actual_blob_chunks = blob_chunks.len();
let blob_base_order = 2usize;
for (i, chunk) in blob_chunks.into_iter().enumerate() {
work_items.push(WorkItem {
order: blob_base_order + i,
entries: chunk,
});
}
let tag_order = blob_base_order + actual_blob_chunks;
work_items.push(WorkItem {
order: tag_order,
entries: tag_entries,
});
tracing::info!(
total_work_items = work_items.len(),
blob_chunks = actual_blob_chunks,
threads = num_threads,
"dispatching delta search to Rayon"
);
type ChunkResult = (usize, Result<Vec<(Vec<u8>, IndexEntry)>, GitError>);
let ez = enable_zstdelta;
let er = enable_rabin;
let dp = disable_prefilter;
let run_delta_search = move || -> Vec<ChunkResult> {
work_items
.into_par_iter()
.map(|item| {
(
item.order,
Self::try_as_offset_delta(item.entries, 10, ez, er, dp),
)
})
.collect()
};
let mut chunk_results: Vec<ChunkResult> =
if let Some(n) = std::env::var("PACK_THREADS")
.ok()
.and_then(|s| s.parse::<usize>().ok())
{
let pool = rayon::ThreadPoolBuilder::new()
.num_threads(n)
.build()
.map_err(|e| GitError::PackEncodeError(format!(
"failed to build Rayon thread pool: {e}"
)))?;
tokio::task::spawn_blocking(move || pool.install(run_delta_search))
.await
.map_err(|e| GitError::PackEncodeError(format!(
"delta search task panicked: {e}"
)))?
} else {
tokio::task::spawn_blocking(run_delta_search)
.await
.map_err(|e| GitError::PackEncodeError(format!(
"delta search task panicked: {e}"
)))?
};
chunk_results.sort_by_key(|(order, _)| *order);
let mut all_res: Vec<Vec<(Vec<u8>, IndexEntry)>> = Vec::with_capacity(chunk_results.len());
for (_order, res) in chunk_results {
all_res.push(res?);
}
let total_entries = all_res.iter().map(Vec::len).sum();
let mut idx_entries = Vec::with_capacity(total_entries);
for res in &mut all_res {
for (encoded_bytes, mut idx_entry) in res.drain(..) {
idx_entry.offset = self.inner_offset as u64;
self.write_owned_and_update(encoded_bytes).await;
idx_entries.push(idx_entry);
}
}
self.idx_entries = Some(idx_entries);
let hash_result = self.inner_hash.clone().finalize();
self.final_hash = Some(
ObjectHash::from_bytes_infer_kind(&hash_result).map_err(GitError::PackEncodeError)?,
);
self.send_data(hash_result).await;
self.drop_sender();
Ok(())
}
async fn write_owned_and_update(&mut self, data: Vec<u8>) {
self.inner_hash.update(&data);
self.inner_offset += data.len();
self.send_data(data).await;
}
async fn generate_idx_file(&mut self) -> Result<(), GitError> {
let final_hash = self.final_hash.ok_or(GitError::PackEncodeError(
"final_hash is missing,The pack file must be generated before the index file is produced."
.into(),
))?;
let idx_entries = self.idx_entries.clone().ok_or(GitError::PackEncodeError(
"The pack file must be generated before the index file is produced.".into(),
))?;
let mut idx_builder = IdxBuilder::new(
self.object_number,
self.idx_sender.clone().unwrap(),
final_hash,
);
idx_builder.write_idx(idx_entries).await?;
Ok(())
}
pub async fn encode_async(
mut self,
rx: mpsc::Receiver<MetaAttached<Entry, EntryMeta>>,
) -> Result<JoinHandle<()>, GitError> {
Ok(tokio::spawn(async move {
if self.window_size == 0 {
self.parallel_encode(rx).await.unwrap()
} else {
self.encode(rx).await.unwrap()
}
}))
}
pub async fn encode_async_with_zstdelta(
mut self,
rx: mpsc::Receiver<MetaAttached<Entry, EntryMeta>>,
) -> Result<JoinHandle<()>, GitError> {
Ok(tokio::spawn(async move {
self.encode_with_zstdelta(rx).await.unwrap()
}))
}
pub async fn encode_idx_file(&mut self) -> Result<(), GitError> {
if self.idx_sender.is_none() {
return Err(GitError::PackEncodeError(String::from(
"idx sender is none",
)));
}
self.generate_idx_file().await?;
self.idx_sender.take();
Ok(())
}
}