Skip to main content

git_internal/internal/pack/
encode.rs

1//! Pack encoder capable of building streamed `.pack`/`.idx` pairs with optional delta compression,
2//! windowing, and asynchronous writers.
3
4use std::{
5    cmp::Ordering,
6    collections::VecDeque,
7    hash::{Hash, Hasher},
8    io::Write,
9    path::{Path, PathBuf},
10};
11
12use ahash::AHasher;
13// use libc::ungetc;
14use chrono::Utc;
15use flate2::write::ZlibEncoder;
16use natord::compare;
17use rayon::prelude::*;
18//use tokio::io::AsyncWriteExt;
19use tokio::io::AsyncWriteExt as TokioAsyncWriteExt;
20use tokio::{fs::File, sync::mpsc, task::JoinHandle};
21
22//use std::io as stdio;
23use crate::delta;
24use crate::{
25    errors::GitError,
26    hash::ObjectHash,
27    internal::{
28        metadata::{EntryMeta, MetaAttached},
29        object::types::ObjectType,
30        pack::{entry::Entry, index_entry::IndexEntry, pack_index::IdxBuilder},
31    },
32    time_it,
33    utils::HashAlgorithm,
34    zstdelta,
35};
36
37const MAX_CHAIN_LEN: usize = 50;
38const MIN_DELTA_RATE: f64 = 0.5; // minimum delta rate
39//const MAX_ZSTDELTA_CHAIN_LEN: usize = 50;
40
41/// A encoder for generating pack files with delta objects.
42pub struct PackEncoder {
43    //path: Option<PathBuf>,
44    object_number: usize,
45    process_index: usize,
46    window_size: usize,
47    // window: VecDeque<(Entry, usize)>, // entry and offset
48    pack_sender: Option<mpsc::Sender<Vec<u8>>>,
49    idx_sender: Option<mpsc::Sender<Vec<u8>>>,
50    //idx_sender: Option<mpsc::Sender<Vec<u8>>>,
51    idx_entries: Option<Vec<IndexEntry>>,
52    inner_offset: usize,       // offset of current entry
53    inner_hash: HashAlgorithm, // introduce different hash algorithm
54    final_hash: Option<ObjectHash>,
55    start_encoding: bool,
56}
57
58/// Encode entries into a pack, write `.pack`/`.idx` files to `output_dir`.
59/// - Spawns background writers to consume pack/idx channels to avoid back-pressure.
60/// - Uses `window_size` to control delta: `0` means no delta (parallel encode), otherwise enable delta window.
61/// # Arguments
62/// * `raw_entries_rx` - receiver providing entries with metadata
63/// * `object_number` - expected total object count for the pack header
64/// * `output_dir` - target directory to place the generated files
65/// * `window_size` - delta window size; `0` disables delta
66/// # Returns
67/// * `Ok(())` on success, `GitError` on failure
68pub async fn encode_and_output_to_files(
69    raw_entries_rx: mpsc::Receiver<MetaAttached<Entry, EntryMeta>>,
70    object_number: usize,
71    output_dir: PathBuf,
72    window_size: usize,
73) -> Result<(), GitError> {
74    let (pack_tx, mut pack_rx) = mpsc::channel(1024);
75    let (idx_tx, mut idx_rx) = mpsc::channel(1024);
76    let mut pack_encoder = PackEncoder::new_with_idx(object_number, window_size, pack_tx, idx_tx);
77
78    // timestamp for temp filename
79    let now = Utc::now();
80    let timestamp = now.format("%Y%m%d%H%M%S%.3f").to_string(); // 例如 20251209235959.123
81    let tmp_path = output_dir.join(format!("{}objects.pack.tmp", timestamp));
82    let mut pack_file = File::create(&tmp_path).await?;
83
84    let pack_writer = tokio::spawn(async move {
85        while let Some(chunk) = pack_rx.recv().await {
86            TokioAsyncWriteExt::write_all(&mut pack_file, &chunk).await?;
87        }
88        //pack_file.flush().await?;
89        TokioAsyncWriteExt::flush(&mut pack_file).await?;
90        Ok::<(), GitError>(())
91    });
92
93    pack_encoder.encode(raw_entries_rx).await?;
94
95    // 等待 pack 写入完成
96    let pack_write_result = pack_writer
97        .await
98        .map_err(|e| GitError::PackEncodeError(format!("pack writer task join error: {e}")))?;
99    pack_write_result?;
100
101    let final_pack_name =
102        output_dir.join(format!("pack-{}.pack", pack_encoder.final_hash.unwrap()));
103    let final_idx_name = output_dir.join(format!("pack-{}.idx", pack_encoder.final_hash.unwrap()));
104    tokio::fs::rename(tmp_path, &final_pack_name).await?;
105
106    let mut idx_file = File::create(&final_idx_name).await?;
107    let idx_writer = tokio::spawn(async move {
108        while let Some(chunk) = idx_rx.recv().await {
109            //idx_file.write_all(&chunk).await?;
110            TokioAsyncWriteExt::write_all(&mut idx_file, &chunk).await?;
111        }
112        //idx_file.flush().await?;
113        TokioAsyncWriteExt::flush(&mut idx_file).await?;
114        Ok::<(), GitError>(())
115    });
116
117    //build idx
118    pack_encoder.encode_idx_file().await?;
119
120    let idx_write_result = idx_writer
121        .await
122        .map_err(|e| GitError::PackEncodeError(format!("idx writer task join error: {e}")))?;
123    idx_write_result?;
124
125    Ok(())
126}
127
128/// Encode header of pack file (12 byte)<br>
129/// Content: 'PACK', Version(2), number of objects
130fn encode_header(object_number: usize) -> Vec<u8> {
131    let mut result: Vec<u8> = vec![
132        b'P', b'A', b'C', b'K', // The logotype of the Pack File
133        0, 0, 0, 2, // generates version 2 only.
134    ];
135    assert_ne!(object_number, 0); // guarantee self.number_of_objects!=0
136    assert!(object_number <= u32::MAX as usize);
137    //TODO: GitError:numbers of objects should < 4G ,
138    result.append((object_number as u32).to_be_bytes().to_vec().as_mut()); // to 4 bytes (network byte order aka. big-endian)
139    result
140}
141
142/// Encode offset of delta object
143fn encode_offset(mut value: usize) -> Vec<u8> {
144    assert_ne!(value, 0, "offset can't be zero");
145    let mut bytes = Vec::new();
146
147    bytes.push((value & 0x7F) as u8);
148    value >>= 7;
149    while value != 0 {
150        value -= 1;
151        let byte = (value & 0x7F) as u8 | 0x80; // set first bit one
152        value >>= 7;
153        bytes.push(byte);
154    }
155    bytes.reverse();
156    bytes
157}
158
159/// Encode one object, and update the hash
160/// @offset: offset of this object if it's a delta object. For other object, it's None
161fn encode_one_object(entry: &Entry, offset: Option<usize>) -> Result<Vec<u8>, GitError> {
162    // try encode as delta
163    let obj_data = &entry.data;
164    let obj_data_len = obj_data.len();
165    let obj_type_number = entry.obj_type.to_pack_type_u8()?;
166
167    let mut encoded_data = Vec::new();
168
169    // **header** encoding
170    let mut header_data = vec![(0x80 | (obj_type_number << 4)) + (obj_data_len & 0x0f) as u8];
171    let mut size = obj_data_len >> 4; // 4 bit has been used in first byte
172    if size > 0 {
173        while size > 0 {
174            if size >> 7 > 0 {
175                header_data.push((0x80 | size) as u8);
176                size >>= 7;
177            } else {
178                header_data.push(size as u8);
179                break;
180            }
181        }
182    } else {
183        header_data.push(0);
184    }
185    encoded_data.extend(header_data);
186
187    // **offset** encoding
188    if entry.obj_type == ObjectType::OffsetDelta || entry.obj_type == ObjectType::OffsetZstdelta {
189        let offset_data = encode_offset(offset.unwrap());
190        encoded_data.extend(offset_data);
191    } else if entry.obj_type == ObjectType::HashDelta {
192        unreachable!("unsupported type")
193    }
194
195    // **data** encoding, need zlib compress
196    let mut inflate = ZlibEncoder::new(Vec::new(), flate2::Compression::default());
197    inflate
198        .write_all(obj_data)
199        .expect("zlib compress should never failed");
200    inflate.flush().expect("zlib flush should never failed");
201    let compressed_data = inflate.finish().expect("zlib compress should never failed");
202    // self.write_all_and_update(&compressed_data).await;
203    encoded_data.extend(compressed_data);
204    Ok(encoded_data)
205}
206
207/// Magic sort function for entries
208fn magic_sort(a: &MetaAttached<Entry, EntryMeta>, b: &MetaAttached<Entry, EntryMeta>) -> Ordering {
209    let path_a = a.meta.file_path.as_ref();
210    let path_b = b.meta.file_path.as_ref();
211
212    // 1. Handle path existence: entries with paths sort first
213    match (path_a, path_b) {
214        (Some(pa), Some(pb)) => {
215            let pa = Path::new(pa);
216            let pb = Path::new(pb);
217
218            // 1. Compare parent directory paths
219            let dir_ord = pa.parent().cmp(&pb.parent());
220            if dir_ord != Ordering::Equal {
221                return dir_ord;
222            }
223
224            // 2. Compare filenames (natural sort)
225            let name_a = pa.file_name().unwrap_or_default().to_string_lossy();
226            let name_b = pb.file_name().unwrap_or_default().to_string_lossy();
227            let name_ord = compare(&name_a, &name_b);
228            if name_ord != Ordering::Equal {
229                return name_ord;
230            }
231        }
232        (Some(_), None) => return Ordering::Less, // entries with paths sort first
233        (None, Some(_)) => return Ordering::Greater, // entries without paths sort last
234        (None, None) => {}
235    }
236
237    let ord = b.inner.data.len().cmp(&a.inner.data.len());
238    if ord != Ordering::Equal {
239        return ord;
240    }
241
242    // fallback pointer order (newest first)
243    (a as *const MetaAttached<Entry, EntryMeta>).cmp(&(b as *const MetaAttached<Entry, EntryMeta>))
244}
245
246/// Calculate hash of data
247fn calc_hash(data: &[u8]) -> u64 {
248    let mut hasher = AHasher::default();
249    data.hash(&mut hasher);
250    hasher.finish()
251}
252
253/// Cheap check if two byte slices are similar by comparing their hashes of the first 128 bytes.
254fn cheap_similar(a: &[u8], b: &[u8]) -> bool {
255    let k = a.len().min(b.len()).min(128);
256    if k == 0 {
257        return false;
258    }
259    calc_hash(&a[..k]) == calc_hash(&b[..k])
260}
261
262impl PackEncoder {
263    pub fn new(object_number: usize, window_size: usize, sender: mpsc::Sender<Vec<u8>>) -> Self {
264        PackEncoder {
265            object_number,
266            window_size,
267            process_index: 0,
268            // window: VecDeque::with_capacity(window_size),
269            pack_sender: Some(sender),
270            idx_sender: None,
271            idx_entries: None,
272            inner_offset: 12, // start  after 12 bytes pack header(signature + version + object count).
273            inner_hash: HashAlgorithm::new(), // introduce different hash algorithm
274            final_hash: None,
275            start_encoding: false,
276        }
277    }
278
279    pub fn new_with_idx(
280        object_number: usize,
281        window_size: usize,
282        pack_sender: mpsc::Sender<Vec<u8>>,
283        idx_sender: mpsc::Sender<Vec<u8>>,
284    ) -> Self {
285        PackEncoder {
286            //path: Some(path),
287            object_number,
288            window_size,
289            process_index: 0,
290            // window: VecDeque::with_capacity(window_size),
291            pack_sender: Some(pack_sender),
292            idx_sender: Some(idx_sender),
293            idx_entries: None,
294            inner_offset: 12, // start  after 12 bytes pack header(signature + version + object count).
295            inner_hash: HashAlgorithm::new(), // introduce different hash algorithm
296            final_hash: None,
297            start_encoding: false,
298        }
299    }
300
301    pub fn drop_sender(&mut self) {
302        self.pack_sender.take(); // Take the sender out, dropping it
303    }
304
305    pub async fn send_data(&mut self, data: Vec<u8>) {
306        if let Some(sender) = &self.pack_sender {
307            sender.send(data).await.unwrap();
308        }
309    }
310
311    /// Get the hash of the pack file. if the pack file is not finished, return None
312    pub fn get_hash(&self) -> Option<ObjectHash> {
313        self.final_hash
314    }
315
316    /// Encodes entries into a pack file with delta objects and outputs them through the specified writer.
317    /// # Arguments
318    /// - `rx` - A receiver channel (`mpsc::Receiver<Entry>`) from which entries to be encoded are received.
319    /// # Returns
320    /// Returns `Ok(())` if encoding is successful, or a `GitError` in case of failure.
321    /// - Returns a `GitError` if there is a failure during the encoding process.
322    /// - Returns `PackEncodeError` if an encoding operation is already in progress.
323    pub async fn encode(
324        &mut self,
325        entry_rx: mpsc::Receiver<MetaAttached<Entry, EntryMeta>>,
326    ) -> Result<(), GitError> {
327        //self.inner_encode(entry_rx, false).await
328        if self.window_size == 0 {
329            self.parallel_encode(entry_rx).await
330        } else {
331            self.inner_encode(entry_rx, false).await
332        }
333    }
334
335    /// Encode with zstdelta
336    pub async fn encode_with_zstdelta(
337        &mut self,
338        entry_rx: mpsc::Receiver<MetaAttached<Entry, EntryMeta>>,
339    ) -> Result<(), GitError> {
340        self.inner_encode(entry_rx, true).await
341    }
342
343    /// Delta selection heuristics are based on:
344    ///   https://github.com/git/git/blob/master/Documentation/technical/pack-heuristics.adoc
345    async fn inner_encode(
346        &mut self,
347        mut entry_rx: mpsc::Receiver<MetaAttached<Entry, EntryMeta>>,
348        enable_zstdelta: bool,
349    ) -> Result<(), GitError> {
350        let head = encode_header(self.object_number);
351        self.send_data(head.clone()).await;
352        self.inner_hash.update(&head);
353
354        // ensure only one decode can only invoke once
355        if self.start_encoding {
356            return Err(GitError::PackEncodeError(
357                "encoding operation is already in progress".to_string(),
358            ));
359        }
360
361        let mut commits: Vec<MetaAttached<Entry, EntryMeta>> = Vec::new();
362        let mut trees: Vec<MetaAttached<Entry, EntryMeta>> = Vec::new();
363        let mut blobs: Vec<MetaAttached<Entry, EntryMeta>> = Vec::new();
364        let mut tags: Vec<MetaAttached<Entry, EntryMeta>> = Vec::new();
365        while let Some(entry) = entry_rx.recv().await {
366            match entry.inner.obj_type {
367                ObjectType::Commit => {
368                    commits.push(entry);
369                }
370                ObjectType::Tree => {
371                    trees.push(entry);
372                }
373                ObjectType::Blob => {
374                    blobs.push(entry);
375                }
376                ObjectType::Tag => {
377                    tags.push(entry);
378                }
379                _ => {
380                    return Err(GitError::PackEncodeError(format!(
381                        "object type `{}` is not supported by delta-window pack encoding",
382                        entry.inner.obj_type
383                    )));
384                }
385            }
386        }
387
388        commits.sort_by(magic_sort);
389        trees.sort_by(magic_sort);
390        blobs.sort_by(magic_sort);
391        tags.sort_by(magic_sort);
392        tracing::info!(
393            "numbers :  commits: {:?} trees: {:?} blobs:{:?} tag :{:?}",
394            commits.len(),
395            trees.len(),
396            blobs.len(),
397            tags.len()
398        );
399
400        // parallel encoding vec with different object_type
401        let (commit_results, tree_results, blob_results, tag_results) = tokio::try_join!(
402            tokio::task::spawn_blocking(move || {
403                Self::try_as_offset_delta(
404                    commits
405                        .into_iter()
406                        .map(|entry_with_meta| entry_with_meta.inner)
407                        .collect(),
408                    10,
409                    enable_zstdelta,
410                )
411            }),
412            tokio::task::spawn_blocking(move || {
413                Self::try_as_offset_delta(
414                    trees
415                        .into_iter()
416                        .map(|entry_with_meta| entry_with_meta.inner)
417                        .collect(),
418                    10,
419                    enable_zstdelta,
420                )
421            }),
422            tokio::task::spawn_blocking(move || {
423                Self::try_as_offset_delta(
424                    blobs
425                        .into_iter()
426                        .map(|entry_with_meta| entry_with_meta.inner)
427                        .collect(),
428                    10,
429                    enable_zstdelta,
430                )
431            }),
432            tokio::task::spawn_blocking(move || {
433                Self::try_as_offset_delta(
434                    tags.into_iter()
435                        .map(|entry_with_meta| entry_with_meta.inner)
436                        .collect(),
437                    10,
438                    enable_zstdelta,
439                )
440            }),
441        )
442        .map_err(|e| GitError::PackEncodeError(format!("Task join error: {e}")))?;
443
444        let commit_res = commit_results?;
445        let tree_res = tree_results?;
446        let blob_res = blob_results?;
447        let tag_res = tag_results?;
448
449        let mut all_res = vec![commit_res, tree_res, blob_res, tag_res];
450
451        let mut idx_entries = Vec::new();
452        for res in &mut all_res {
453            for data in res {
454                data.1.offset = self.inner_offset as u64;
455                self.write_all_and_update(&data.0).await;
456                idx_entries.push(data.1.clone());
457            }
458        }
459
460        self.idx_entries = Some(idx_entries);
461
462        // Hash signature
463        let hash_result = self.inner_hash.clone().finalize();
464        self.final_hash = Some(ObjectHash::from_bytes(&hash_result).unwrap());
465        self.send_data(hash_result.to_vec()).await;
466
467        self.drop_sender();
468        Ok(())
469    }
470
471    /// Try to encode as delta using objects in window
472    /// delta & zstdelta have been gathered here
473    /// Refs: https://sapling-scm.com/docs/dev/internals/zstdelta/
474    /// the sliding window was moved here
475    /// # Returns
476    /// - Return (Vec<Vec<u8>) if success make delta
477    /// - Return (None) if didn't delta,
478    fn try_as_offset_delta(
479        mut bucket: Vec<Entry>,
480        window_size: usize,
481        enable_zstdelta: bool,
482    ) -> Result<Vec<(Vec<u8>, IndexEntry)>, GitError> {
483        let mut current_offset = 0usize;
484        let mut window: VecDeque<(Entry, usize)> = VecDeque::with_capacity(window_size);
485        let mut res: Vec<(Vec<u8>, IndexEntry)> = Vec::new();
486        //let mut idx_entries: Vec<IndexEntry> = Vec::new();
487
488        for entry in bucket.iter_mut() {
489            //let entry_for_window = entry.clone();
490            // 每次循环重置最佳基对象选择
491            let mut best_base: Option<&(Entry, usize)> = None;
492            let mut best_rate: f64 = 0.0;
493            let tie_epsilon: f64 = 0.15;
494
495            let candidates: Vec<_> = window
496                .par_iter()
497                .with_min_len(3)
498                .filter_map(|try_base| {
499                    if try_base.0.obj_type != entry.obj_type {
500                        return None;
501                    }
502
503                    if try_base.0.chain_len >= MAX_CHAIN_LEN {
504                        return None;
505                    }
506
507                    if try_base.0.hash == entry.hash {
508                        return None;
509                    }
510
511                    let sym_ratio = (try_base.0.data.len().min(entry.data.len()) as f64)
512                        / (try_base.0.data.len().max(entry.data.len()) as f64);
513                    if sym_ratio < 0.5 {
514                        return None;
515                    }
516
517                    if !cheap_similar(&try_base.0.data, &entry.data) {
518                        return None;
519                    }
520
521                    let rate = if (try_base.0.data.len() + entry.data.len()) / 2 > 64 {
522                        delta::heuristic_encode_rate_parallel(&try_base.0.data, &entry.data)
523                    } else {
524                        delta::encode_rate(&try_base.0.data, &entry.data)
525                        // let try_delta_obj = zstdelta::diff(&try_base.0.data, &entry.data).unwrap();
526                        // 1.0 - try_delta_obj.len() as f64 / entry.data.len() as f64
527                    };
528
529                    if rate > MIN_DELTA_RATE {
530                        Some((rate, try_base))
531                    } else {
532                        None
533                    }
534                })
535                .collect();
536
537            for (rate, try_base) in candidates {
538                match best_base {
539                    None => {
540                        best_rate = rate;
541                        //best_base_offset = current_offset - try_base.1;
542                        best_base = Some(try_base);
543                    }
544                    Some(best_base_ref) => {
545                        let is_better = if rate > best_rate + tie_epsilon {
546                            true
547                        } else if (rate - best_rate).abs() <= tie_epsilon {
548                            try_base.0.chain_len > best_base_ref.0.chain_len
549                        } else {
550                            false
551                        };
552
553                        if is_better {
554                            best_rate = rate;
555                            best_base = Some(try_base);
556                        }
557                    }
558                }
559            }
560
561            let mut entry_for_window = entry.clone();
562
563            let offset = best_base.map(|best_base| {
564                let delta = if enable_zstdelta {
565                    entry.obj_type = ObjectType::OffsetZstdelta;
566                    zstdelta::diff(&best_base.0.data, &entry.data)
567                        .map_err(|e| {
568                            GitError::DeltaObjectError(format!("zstdelta diff failed: {e}"))
569                        })
570                        .unwrap()
571                } else {
572                    entry.obj_type = ObjectType::OffsetDelta;
573                    delta::encode(&best_base.0.data, &entry.data)
574                };
575                //entry.obj_type = ObjectType::OffsetDelta;
576                entry.data = delta;
577                entry.chain_len = best_base.0.chain_len + 1;
578                current_offset - best_base.1
579            });
580
581            entry_for_window.chain_len = entry.chain_len;
582            let obj_data = encode_one_object(entry, offset)?;
583            window.push_back((entry_for_window, current_offset));
584            if window.len() > window_size {
585                window.pop_front();
586            }
587            res.push((obj_data.clone(), IndexEntry::new(entry, 0)));
588            current_offset += obj_data.len();
589        }
590        Ok(res)
591    }
592
593    /// Parallel encode with rayon, only works when window_size == 0 (no delta)
594    pub async fn parallel_encode(
595        &mut self,
596        mut entry_rx: mpsc::Receiver<MetaAttached<Entry, EntryMeta>>,
597    ) -> Result<(), GitError> {
598        if self.window_size != 0 {
599            return Err(GitError::PackEncodeError(
600                "parallel encode only works when window_size == 0".to_string(),
601            ));
602        }
603
604        let head = encode_header(self.object_number);
605        self.send_data(head.clone()).await;
606        self.inner_hash.update(&head);
607
608        // ensure only one decode can only invoke once
609        if self.start_encoding {
610            return Err(GitError::PackEncodeError(
611                "encoding operation is already in progress".to_string(),
612            ));
613        }
614
615        let mut idx_entries = Vec::new();
616        let batch_size = usize::max(1000, entry_rx.max_capacity() / 10); // A temporary value, not optimized
617        tracing::info!("encode with batch size: {}", batch_size);
618        loop {
619            let mut batch_entries = Vec::with_capacity(batch_size);
620            time_it!("parallel encode: receive batch", {
621                for _ in 0..batch_size {
622                    match entry_rx.recv().await {
623                        Some(entry) => {
624                            if entry.inner.obj_type.is_ai_object() {
625                                return Err(GitError::PackEncodeError(format!(
626                                    "AI object type `{}` cannot be encoded in a pack file",
627                                    entry.inner.obj_type
628                                )));
629                            }
630                            batch_entries.push(entry.inner);
631                            self.process_index += 1;
632                        }
633                        None => break,
634                    }
635                }
636            });
637
638            if batch_entries.is_empty() {
639                break;
640            }
641
642            // use `collect` will return result in order, refs: https://github.com/rayon-rs/rayon/issues/551#issuecomment-371657900
643            let batch_result: Vec<Result<(Vec<u8>, IndexEntry), GitError>> =
644                time_it!("parallel encode: encode batch", {
645                    batch_entries
646                        .par_iter()
647                        .map(|entry| {
648                            encode_one_object(entry, None)
649                                .map(|encoded| (encoded, IndexEntry::new(entry, 0)))
650                        })
651                        .collect()
652                });
653
654            time_it!("parallel encode: write batch", {
655                for obj_data in batch_result {
656                    let mut obj_data = obj_data?;
657                    obj_data.1.offset = self.inner_offset as u64;
658                    self.write_all_and_update(&obj_data.0).await;
659                    idx_entries.push(obj_data.1);
660                }
661            });
662        }
663
664        tracing::debug!("parallel encode idx entries: {:?}", idx_entries.len());
665        if self.process_index != self.object_number {
666            panic!(
667                "not all objects are encoded, process:{}, total:{}",
668                self.process_index, self.object_number
669            );
670        }
671
672        // hash signature
673        let hash_result = self.inner_hash.clone().finalize();
674        self.final_hash = Some(ObjectHash::from_bytes(&hash_result).unwrap());
675        self.send_data(hash_result.to_vec()).await;
676        self.drop_sender();
677
678        self.idx_entries = Some(idx_entries);
679        Ok(())
680    }
681
682    /// Write data to writer and update hash & offset
683    async fn write_all_and_update(&mut self, data: &[u8]) {
684        self.inner_hash.update(data);
685        self.inner_offset += data.len();
686        self.send_data(data.to_vec()).await;
687    }
688
689    async fn generate_idx_file(&mut self) -> Result<(), GitError> {
690        let final_hash = self.final_hash
691            .ok_or(GitError::PackEncodeError("final_hash is missing,The pack file must be generated before the index file is produced.".into()))?;
692        let idx_entries = self.idx_entries.clone().ok_or(GitError::PackEncodeError(
693            "The pack file must be generated before the index file is produced.".into(),
694        ))?;
695        let mut idx_builder = IdxBuilder::new(
696            self.object_number,
697            self.idx_sender.clone().unwrap(),
698            final_hash,
699        );
700        idx_builder.write_idx(idx_entries).await?;
701        Ok(())
702    }
703
704    /// async version of encode, result data will be returned by JoinHandle.
705    /// It will consume PackEncoder, so you can't use it after calling this function.
706    /// when window_size = 0, it executes parallel_encode which retains stream transmission
707    /// when window_size = 0,it executes encode which uses magic sort and delta.
708    /// It seems that all other modules rely on this api
709    pub async fn encode_async(
710        mut self,
711        rx: mpsc::Receiver<MetaAttached<Entry, EntryMeta>>,
712    ) -> Result<JoinHandle<()>, GitError> {
713        Ok(tokio::spawn(async move {
714            if self.window_size == 0 {
715                self.parallel_encode(rx).await.unwrap()
716            } else {
717                self.encode(rx).await.unwrap()
718            }
719        }))
720    }
721
722    /// async version of encode_with_zstdelta, result data will be returned by JoinHandle.
723    pub async fn encode_async_with_zstdelta(
724        mut self,
725        rx: mpsc::Receiver<MetaAttached<Entry, EntryMeta>>,
726    ) -> Result<JoinHandle<()>, GitError> {
727        Ok(tokio::spawn(async move {
728            // Do not use parallel encode with zstdelta because it make no sense.
729            self.encode_with_zstdelta(rx).await.unwrap()
730        }))
731    }
732
733    /// Generate idx file after pack file has been generated
734    pub async fn encode_idx_file(&mut self) -> Result<(), GitError> {
735        if self.idx_sender.is_none() {
736            return Err(GitError::PackEncodeError(String::from(
737                "idx sender is none",
738            )));
739        }
740        self.generate_idx_file().await?;
741        // drop sender so downstream consumer can finish
742        self.idx_sender.take();
743        Ok(())
744    }
745}
746
747#[cfg(test)]
748mod tests {
749    use std::{io::Cursor, path::PathBuf, sync::Arc, time::Instant};
750
751    use tempfile::tempdir;
752    use tokio::sync::Mutex;
753
754    use super::*;
755    use crate::{
756        hash::{HashKind, ObjectHash, set_hash_kind_for_test},
757        internal::{
758            object::{blob::Blob, types::ObjectType},
759            pack::{
760                Pack,
761                test_pack_download::{PackFileGuard, download_pack_file},
762                tests::init_logger,
763                utils::read_offset_encoding,
764            },
765        },
766        time_it,
767    };
768
769    /// Check if the given data is a valid pack file format by attempting to decode it.
770    fn check_format(data: &Vec<u8>) {
771        // Use a smaller cap on 32-bit targets to avoid usize overflow.
772        let max_pack_size_u64 = if cfg!(target_pointer_width = "64") {
773            6u64 * 1024 * 1024 * 1024
774        } else {
775            2u64 * 1024 * 1024 * 1024
776        };
777        let max_pack_size = usize::try_from(max_pack_size_u64).unwrap_or_else(|_| {
778            panic!(
779                "internal assertion failed: pack size cap {} does not fit in usize on this \
780                 target; this should be unreachable given the target_pointer_width configuration",
781                max_pack_size_u64
782            )
783        });
784        let mut p = Pack::new(
785            None,
786            Some(max_pack_size), // 6GB on 64-bit, 2GB on 32-bit
787            Some(PathBuf::from("/tmp/.cache_temp")),
788            true,
789        );
790        let mut reader = Cursor::new(data);
791        tracing::debug!("start check format");
792        p.decode(&mut reader, |_| {}, None::<fn(ObjectHash)>)
793            .expect("pack file format error");
794    }
795
796    #[tokio::test]
797    async fn test_pack_encoder() {
798        let _guard = set_hash_kind_for_test(HashKind::Sha1);
799        async fn encode_once(window_size: usize) -> Vec<u8> {
800            let (tx, mut rx) = mpsc::channel(100);
801            let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(1);
802
803            // make some different objects, or decode will fail
804            let str_vec = vec!["hello, word", "hello, world.", "!", "123141251251"];
805            let encoder = PackEncoder::new(str_vec.len(), window_size, tx);
806            encoder.encode_async(entry_rx).await.unwrap();
807
808            for str in str_vec {
809                let blob = Blob::from_content(str);
810                let entry: Entry = blob.into();
811                entry_tx
812                    .send(MetaAttached {
813                        inner: entry,
814                        meta: EntryMeta::new(),
815                    })
816                    .await
817                    .unwrap();
818            }
819            drop(entry_tx);
820            // assert!(encoder.get_hash().is_some());
821            let mut result = Vec::new();
822            while let Some(chunk) = rx.recv().await {
823                result.extend(chunk);
824            }
825            result
826        }
827
828        // without delta
829        let pack_without_delta = encode_once(0).await;
830        let pack_without_delta_size = pack_without_delta.len();
831        check_format(&pack_without_delta);
832
833        // with delta
834        let pack_with_delta = encode_once(4).await;
835        assert!(pack_with_delta.len() <= pack_without_delta_size);
836        check_format(&pack_with_delta);
837    }
838    #[tokio::test]
839    async fn test_pack_encoder_sha256() {
840        let _guard = set_hash_kind_for_test(HashKind::Sha256);
841
842        async fn encode_once(window_size: usize) -> Vec<u8> {
843            let (tx, mut rx) = mpsc::channel(100);
844            let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(1);
845
846            let str_vec = vec!["hello, word", "hello, world.", "!", "123141251251"];
847            let encoder = PackEncoder::new(str_vec.len(), window_size, tx);
848            encoder.encode_async(entry_rx).await.unwrap();
849
850            for s in str_vec {
851                let blob = Blob::from_content(s);
852                let entry: Entry = blob.into();
853                entry_tx
854                    .send(MetaAttached {
855                        inner: entry,
856                        meta: EntryMeta::new(),
857                    })
858                    .await
859                    .unwrap();
860            }
861            drop(entry_tx);
862
863            let mut result = Vec::new();
864            while let Some(chunk) = rx.recv().await {
865                result.extend(chunk);
866            }
867            result
868        }
869
870        // without delta
871        let pack_without_delta = encode_once(0).await;
872        let pack_without_delta_size = pack_without_delta.len();
873        check_format(&pack_without_delta);
874
875        // with delta
876        let pack_with_delta = encode_once(4).await;
877        assert!(pack_with_delta.len() <= pack_without_delta_size);
878        check_format(&pack_with_delta);
879    }
880
881    #[tokio::test]
882    async fn test_pack_encoder_rejects_unencodable_ai_type_parallel() {
883        let (tx, _rx) = mpsc::channel(8);
884        let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(1);
885        let mut encoder = PackEncoder::new(1, 0, tx);
886
887        let mut entry: Entry = Blob::from_content("ai").into();
888        entry.obj_type = ObjectType::Task;
889        entry_tx
890            .send(MetaAttached {
891                inner: entry,
892                meta: EntryMeta::new(),
893            })
894            .await
895            .expect("send entry");
896        drop(entry_tx);
897
898        let err = encoder
899            .encode(entry_rx)
900            .await
901            .expect_err("must reject AI pack type");
902        assert!(matches!(err, GitError::PackEncodeError(_)));
903    }
904
905    #[tokio::test]
906    async fn test_pack_encoder_rejects_unencodable_ai_type_delta_window() {
907        let (tx, _rx) = mpsc::channel(8);
908        let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(1);
909        let mut encoder = PackEncoder::new(1, 10, tx);
910
911        let mut entry: Entry = Blob::from_content("ai").into();
912        entry.obj_type = ObjectType::Task;
913        entry_tx
914            .send(MetaAttached {
915                inner: entry,
916                meta: EntryMeta::new(),
917            })
918            .await
919            .expect("send entry");
920        drop(entry_tx);
921
922        let err = encoder
923            .encode(entry_rx)
924            .await
925            .expect_err("must reject AI pack type");
926        assert!(matches!(err, GitError::PackEncodeError(_)));
927    }
928
929    async fn get_entries_for_test() -> (Arc<Mutex<Vec<Entry>>>, PackFileGuard) {
930        let (source, dl_guard) = download_pack_file("encode-test-sha1.pack");
931
932        let mut p = Pack::new(None, None, Some(PathBuf::from("/tmp/.cache_temp")), true);
933
934        let f = std::fs::File::open(&source).unwrap();
935        tracing::info!("pack file size: {}", f.metadata().unwrap().len());
936        let mut reader = std::io::BufReader::new(f);
937        let entries = Arc::new(Mutex::new(Vec::new()));
938        let entries_clone = entries.clone();
939        p.decode(
940            &mut reader,
941            move |entry| {
942                let mut entries = entries_clone.blocking_lock();
943                entries.push(entry.inner);
944            },
945            None::<fn(ObjectHash)>,
946        )
947        .unwrap();
948        assert_eq!(p.number, entries.lock().await.len());
949        tracing::info!("total entries: {}", p.number);
950        drop(p);
951
952        (entries, dl_guard)
953    }
954    async fn get_entries_for_test_sha256() -> (Arc<Mutex<Vec<Entry>>>, PackFileGuard) {
955        let (source, dl_guard) = download_pack_file("encode-test-sha256.pack");
956
957        let mut p = Pack::new(None, None, Some(PathBuf::from("/tmp/.cache_temp")), true);
958
959        let f = std::fs::File::open(&source).unwrap();
960        tracing::info!("pack file size: {}", f.metadata().unwrap().len());
961        let mut reader = std::io::BufReader::new(f);
962        let entries = Arc::new(Mutex::new(Vec::new()));
963        let entries_clone = entries.clone();
964        p.decode(
965            &mut reader,
966            move |entry| {
967                let mut entries = entries_clone.blocking_lock();
968                entries.push(entry.inner);
969            },
970            None::<fn(ObjectHash)>,
971        )
972        .unwrap();
973        assert_eq!(p.number, entries.lock().await.len());
974        tracing::info!("total entries: {}", p.number);
975        drop(p);
976
977        (entries, dl_guard)
978    }
979
980    #[tokio::test]
981    async fn test_pack_encoder_parallel_large_file() {
982        let _guard = set_hash_kind_for_test(HashKind::Sha1);
983        init_logger();
984
985        let start = Instant::now();
986        let (entries, _dl_guard) = get_entries_for_test().await;
987        let entries_number = entries.lock().await.len();
988
989        let total_original_size: usize = entries
990            .lock()
991            .await
992            .iter()
993            .map(|entry| entry.data.len())
994            .sum();
995
996        // encode entries with parallel
997        let (tx, mut rx) = mpsc::channel(1_000_000);
998        let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(1_000_000);
999
1000        let mut encoder = PackEncoder::new(entries_number, 0, tx);
1001        tokio::spawn(async move {
1002            time_it!("test parallel encode", {
1003                encoder.parallel_encode(entry_rx).await.unwrap();
1004            });
1005        });
1006
1007        // spawn a task to send entries
1008        tokio::spawn(async move {
1009            let entries = entries.lock().await;
1010            for entry in entries.iter() {
1011                entry_tx
1012                    .send(MetaAttached {
1013                        inner: entry.clone(),
1014                        meta: EntryMeta::new(),
1015                    })
1016                    .await
1017                    .unwrap();
1018            }
1019            drop(entry_tx);
1020            tracing::info!("all entries sent");
1021        });
1022
1023        let mut result = Vec::new();
1024        while let Some(chunk) = rx.recv().await {
1025            result.extend(chunk);
1026        }
1027
1028        let pack_size = result.len();
1029        let compression_rate = if total_original_size > 0 {
1030            1.0 - (pack_size as f64 / total_original_size as f64)
1031        } else {
1032            0.0
1033        };
1034
1035        let duration = start.elapsed();
1036        tracing::info!("test executed in: {:.2?}", duration);
1037        tracing::info!("new pack file size: {}", result.len());
1038        tracing::info!("compression rate: {:.2}%", compression_rate * 100.0);
1039        // check format
1040        check_format(&result);
1041    }
1042    #[tokio::test]
1043    async fn test_pack_encoder_parallel_large_file_sha256() {
1044        let _guard = set_hash_kind_for_test(HashKind::Sha256);
1045        init_logger();
1046
1047        let start = Instant::now();
1048        // use sha256 pack file for testing
1049        let (entries, _dl_guard) = get_entries_for_test_sha256().await;
1050        let entries_number = entries.lock().await.len();
1051
1052        let total_original_size: usize = entries
1053            .lock()
1054            .await
1055            .iter()
1056            .map(|entry| entry.data.len())
1057            .sum();
1058
1059        let (tx, mut rx) = mpsc::channel(1_000_000);
1060        let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(1_000_000);
1061
1062        let mut encoder = PackEncoder::new(entries_number, 0, tx);
1063        tokio::spawn(async move {
1064            time_it!("test parallel encode sha256", {
1065                encoder.parallel_encode(entry_rx).await.unwrap();
1066            });
1067        });
1068
1069        tokio::spawn(async move {
1070            let entries = entries.lock().await;
1071            for entry in entries.iter() {
1072                entry_tx
1073                    .send(MetaAttached {
1074                        inner: entry.clone(),
1075                        meta: EntryMeta::new(),
1076                    })
1077                    .await
1078                    .unwrap();
1079            }
1080            drop(entry_tx);
1081            tracing::info!("all entries sent");
1082        });
1083
1084        let mut result = Vec::new();
1085        while let Some(chunk) = rx.recv().await {
1086            result.extend(chunk);
1087        }
1088
1089        let pack_size = result.len();
1090        let compression_rate = if total_original_size > 0 {
1091            1.0 - (pack_size as f64 / total_original_size as f64)
1092        } else {
1093            0.0
1094        };
1095
1096        let duration = start.elapsed();
1097        tracing::info!("sha256 test executed in: {:.2?}", duration);
1098        tracing::info!("new pack file size: {}", result.len());
1099        tracing::info!("compression rate: {:.2}%", compression_rate * 100.0);
1100        check_format(&result);
1101    }
1102
1103    #[tokio::test]
1104    async fn test_pack_encoder_large_file() {
1105        let _guard = set_hash_kind_for_test(HashKind::Sha1);
1106        init_logger();
1107        let (entries, _dl_guard) = get_entries_for_test().await;
1108        let entries_number = entries.lock().await.len();
1109
1110        let total_original_size: usize = entries
1111            .lock()
1112            .await
1113            .iter()
1114            .map(|entry| entry.data.len())
1115            .sum();
1116
1117        let start = Instant::now();
1118        // encode entries
1119        let (tx, mut rx) = mpsc::channel(100_000);
1120        let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(100_000);
1121
1122        let mut encoder = PackEncoder::new(entries_number, 0, tx);
1123        tokio::spawn(async move {
1124            time_it!("test encode no parallel", {
1125                encoder.encode(entry_rx).await.unwrap();
1126            });
1127        });
1128
1129        // spawn a task to send entries
1130        tokio::spawn(async move {
1131            let entries = entries.lock().await;
1132            for entry in entries.iter() {
1133                entry_tx
1134                    .send(MetaAttached {
1135                        inner: entry.clone(),
1136                        meta: EntryMeta::new(),
1137                    })
1138                    .await
1139                    .unwrap();
1140            }
1141            drop(entry_tx);
1142            tracing::info!("all entries sent");
1143        });
1144
1145        // // only receive data
1146        // while (rx.recv().await).is_some() {
1147        //     // do nothing
1148        // }
1149
1150        let mut result = Vec::new();
1151        while let Some(chunk) = rx.recv().await {
1152            result.extend(chunk);
1153        }
1154
1155        let pack_size = result.len();
1156        let compression_rate = if total_original_size > 0 {
1157            1.0 - (pack_size as f64 / total_original_size as f64)
1158        } else {
1159            0.0
1160        };
1161
1162        let duration = start.elapsed();
1163        tracing::info!("test executed in: {:.2?}", duration);
1164        tracing::info!("new pack file size: {}", pack_size);
1165        tracing::info!("original total size: {}", total_original_size);
1166        tracing::info!("compression rate: {:.2}%", compression_rate * 100.0);
1167        tracing::info!(
1168            "space saved: {} bytes",
1169            total_original_size.saturating_sub(pack_size)
1170        );
1171    }
1172    #[tokio::test]
1173    async fn test_pack_encoder_large_file_sha256() {
1174        let _guard = set_hash_kind_for_test(HashKind::Sha256);
1175        init_logger();
1176        let (entries, _dl_guard) = get_entries_for_test_sha256().await;
1177        let entries_number = entries.lock().await.len();
1178
1179        let total_original_size: usize = entries
1180            .lock()
1181            .await
1182            .iter()
1183            .map(|entry| entry.data.len())
1184            .sum();
1185
1186        let start = Instant::now();
1187        // encode entries
1188        let (tx, mut rx) = mpsc::channel(100_000);
1189        let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(100_000);
1190
1191        let mut encoder = PackEncoder::new(entries_number, 0, tx);
1192        tokio::spawn(async move {
1193            time_it!("test encode no parallel sha256", {
1194                encoder.encode(entry_rx).await.unwrap();
1195            });
1196        });
1197
1198        // spawn a task to send entries
1199        tokio::spawn(async move {
1200            let entries = entries.lock().await;
1201            for entry in entries.iter() {
1202                entry_tx
1203                    .send(MetaAttached {
1204                        inner: entry.clone(),
1205                        meta: EntryMeta::new(),
1206                    })
1207                    .await
1208                    .unwrap();
1209            }
1210            drop(entry_tx);
1211            tracing::info!("all entries sent");
1212        });
1213
1214        // // only receive data
1215        // while (rx.recv().await).is_some() {
1216        //     // do nothing
1217        // }
1218
1219        let mut result = Vec::new();
1220        while let Some(chunk) = rx.recv().await {
1221            result.extend(chunk);
1222        }
1223
1224        let pack_size = result.len();
1225        let compression_rate = if total_original_size > 0 {
1226            1.0 - (pack_size as f64 / total_original_size as f64)
1227        } else {
1228            0.0
1229        };
1230
1231        let duration = start.elapsed();
1232        tracing::info!("test executed in: {:.2?}", duration);
1233        tracing::info!("new pack file size: {}", pack_size);
1234        tracing::info!("original total size: {}", total_original_size);
1235        tracing::info!("compression rate: {:.2}%", compression_rate * 100.0);
1236        tracing::info!(
1237            "space saved: {} bytes",
1238            total_original_size.saturating_sub(pack_size)
1239        );
1240    }
1241
1242    #[tokio::test]
1243    async fn test_pack_encoder_with_zstdelta() {
1244        let _guard = set_hash_kind_for_test(HashKind::Sha1);
1245        init_logger();
1246        let (entries, _dl_guard) = get_entries_for_test().await;
1247        let entries_number = entries.lock().await.len();
1248
1249        let total_original_size: usize = entries
1250            .lock()
1251            .await
1252            .iter()
1253            .map(|entry| entry.data.len())
1254            .sum();
1255
1256        let start = Instant::now();
1257        let (tx, mut rx) = mpsc::channel(100_000);
1258        let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(100_000);
1259
1260        let encoder = PackEncoder::new(entries_number, 10, tx);
1261        encoder.encode_async_with_zstdelta(entry_rx).await.unwrap();
1262
1263        // spawn a task to send entries
1264        tokio::spawn(async move {
1265            let entries = entries.lock().await;
1266            for entry in entries.iter() {
1267                entry_tx
1268                    .send(MetaAttached {
1269                        inner: entry.clone(),
1270                        meta: EntryMeta::new(),
1271                    })
1272                    .await
1273                    .unwrap();
1274            }
1275            drop(entry_tx);
1276            tracing::info!("all entries sent");
1277        });
1278
1279        let mut result = Vec::new();
1280        while let Some(chunk) = rx.recv().await {
1281            result.extend(chunk);
1282        }
1283
1284        let pack_size = result.len();
1285        let compression_rate = if total_original_size > 0 {
1286            1.0 - (pack_size as f64 / total_original_size as f64)
1287        } else {
1288            0.0
1289        };
1290
1291        let duration = start.elapsed();
1292        tracing::info!("test executed in: {:.2?}", duration);
1293        tracing::info!("new pack file size: {}", pack_size);
1294        tracing::info!("original total size: {}", total_original_size);
1295        tracing::info!("compression rate: {:.2}%", compression_rate * 100.0);
1296        tracing::info!(
1297            "space saved: {} bytes",
1298            total_original_size.saturating_sub(pack_size)
1299        );
1300
1301        // check format
1302        check_format(&result);
1303    }
1304    #[tokio::test]
1305    async fn test_pack_encoder_with_zstdelta_sha256() {
1306        let _guard = set_hash_kind_for_test(HashKind::Sha256);
1307        init_logger();
1308        let (entries, _dl_guard) = get_entries_for_test_sha256().await;
1309        let entries_number = entries.lock().await.len();
1310
1311        let total_original_size: usize = entries
1312            .lock()
1313            .await
1314            .iter()
1315            .map(|entry| entry.data.len())
1316            .sum();
1317
1318        let start = Instant::now();
1319        let (tx, mut rx) = mpsc::channel(100_000);
1320        let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(100_000);
1321
1322        let encoder = PackEncoder::new(entries_number, 10, tx);
1323        encoder.encode_async_with_zstdelta(entry_rx).await.unwrap();
1324
1325        // spawn a task to send entries
1326        tokio::spawn(async move {
1327            let entries = entries.lock().await;
1328            for entry in entries.iter() {
1329                entry_tx
1330                    .send(MetaAttached {
1331                        inner: entry.clone(),
1332                        meta: EntryMeta::new(),
1333                    })
1334                    .await
1335                    .unwrap();
1336            }
1337            drop(entry_tx);
1338            tracing::info!("all entries sent");
1339        });
1340
1341        let mut result = Vec::new();
1342        while let Some(chunk) = rx.recv().await {
1343            result.extend(chunk);
1344        }
1345
1346        let pack_size = result.len();
1347        let compression_rate = if total_original_size > 0 {
1348            1.0 - (pack_size as f64 / total_original_size as f64)
1349        } else {
1350            0.0
1351        };
1352
1353        let duration = start.elapsed();
1354        tracing::info!("test executed in: {:.2?}", duration);
1355        tracing::info!("new pack file size: {}", pack_size);
1356        tracing::info!("original total size: {}", total_original_size);
1357        tracing::info!("compression rate: {:.2}%", compression_rate * 100.0);
1358        tracing::info!(
1359            "space saved: {} bytes",
1360            total_original_size.saturating_sub(pack_size)
1361        );
1362
1363        // check format
1364        check_format(&result);
1365    }
1366
1367    #[test]
1368    fn test_encode_offset() {
1369        // let value = 11013;
1370        let value = 16389;
1371
1372        let data = encode_offset(value);
1373        println!("{data:?}");
1374        let mut reader = Cursor::new(data);
1375        let (result, _) = read_offset_encoding(&mut reader).unwrap();
1376        println!("result: {result}");
1377        assert_eq!(result, value as u64);
1378    }
1379
1380    #[tokio::test]
1381    async fn test_pack_encoder_large_file_with_delta() {
1382        let _guard = set_hash_kind_for_test(HashKind::Sha1);
1383        init_logger();
1384        let (entries, _dl_guard) = get_entries_for_test().await;
1385        let entries_number = entries.lock().await.len();
1386
1387        let total_original_size: usize = entries
1388            .lock()
1389            .await
1390            .iter()
1391            .map(|entry| entry.data.len())
1392            .sum();
1393
1394        let (tx, mut rx) = mpsc::channel(100_000);
1395        let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(100_000);
1396
1397        let encoder = PackEncoder::new(entries_number, 10, tx);
1398
1399        let start = Instant::now(); // 开始时间
1400        encoder.encode_async(entry_rx).await.unwrap();
1401
1402        // spawn a task to send entries
1403        tokio::spawn(async move {
1404            let entries = entries.lock().await;
1405            for entry in entries.iter() {
1406                entry_tx
1407                    .send(MetaAttached {
1408                        inner: entry.clone(),
1409                        meta: EntryMeta::new(),
1410                    })
1411                    .await
1412                    .unwrap();
1413            }
1414            drop(entry_tx);
1415            tracing::info!("all entries sent");
1416        });
1417
1418        let mut result = Vec::new();
1419        while let Some(chunk) = rx.recv().await {
1420            result.extend(chunk);
1421        }
1422
1423        let pack_size = result.len();
1424        let compression_rate = if total_original_size > 0 {
1425            1.0 - (pack_size as f64 / total_original_size as f64)
1426        } else {
1427            0.0
1428        };
1429
1430        let duration = start.elapsed();
1431        tracing::info!("test executed in: {:.2?}", duration);
1432        tracing::info!("new pack file size: {}", pack_size);
1433        tracing::info!("original total size: {}", total_original_size);
1434        tracing::info!("compression rate: {:.2}%", compression_rate * 100.0);
1435        tracing::info!(
1436            "space saved: {} bytes",
1437            total_original_size.saturating_sub(pack_size)
1438        );
1439
1440        // check format
1441        check_format(&result);
1442    }
1443    #[tokio::test]
1444    async fn test_pack_encoder_large_file_with_delta_sha256() {
1445        let _guard = set_hash_kind_for_test(HashKind::Sha256);
1446        init_logger();
1447        let (entries, _dl_guard) = get_entries_for_test_sha256().await;
1448        let entries_number = entries.lock().await.len();
1449
1450        let total_original_size: usize = entries
1451            .lock()
1452            .await
1453            .iter()
1454            .map(|entry| entry.data.len())
1455            .sum();
1456
1457        let (tx, mut rx) = mpsc::channel(100_000);
1458        let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(100_000);
1459
1460        let encoder = PackEncoder::new(entries_number, 10, tx);
1461
1462        let start = Instant::now(); // 开始时间
1463        encoder.encode_async(entry_rx).await.unwrap();
1464
1465        // spawn a task to send entries
1466        tokio::spawn(async move {
1467            let entries = entries.lock().await;
1468            for entry in entries.iter() {
1469                entry_tx
1470                    .send(MetaAttached {
1471                        inner: entry.clone(),
1472                        meta: EntryMeta::new(),
1473                    })
1474                    .await
1475                    .unwrap();
1476            }
1477            drop(entry_tx);
1478            tracing::info!("all entries sent");
1479        });
1480
1481        let mut result = Vec::new();
1482        while let Some(chunk) = rx.recv().await {
1483            result.extend(chunk);
1484        }
1485
1486        let pack_size = result.len();
1487        let compression_rate = if total_original_size > 0 {
1488            1.0 - (pack_size as f64 / total_original_size as f64)
1489        } else {
1490            0.0
1491        };
1492
1493        let duration = start.elapsed();
1494        tracing::info!("test executed in: {:.2?}", duration);
1495        tracing::info!("new pack file size: {}", pack_size);
1496        tracing::info!("original total size: {}", total_original_size);
1497        tracing::info!("compression rate: {:.2}%", compression_rate * 100.0);
1498        tracing::info!(
1499            "space saved: {} bytes",
1500            total_original_size.saturating_sub(pack_size)
1501        );
1502
1503        // check format
1504        check_format(&result);
1505    }
1506
1507    #[tokio::test]
1508    async fn test_pack_encoder_output_to_files() {
1509        let _guard = set_hash_kind_for_test(HashKind::Sha1);
1510        init_logger();
1511        let (entries, _dl_guard) = get_entries_for_test().await;
1512        let entries_number = entries.lock().await.len();
1513
1514        let total_original_size: usize = entries
1515            .lock()
1516            .await
1517            .iter()
1518            .map(|entry| entry.data.len())
1519            .sum();
1520
1521        let start = Instant::now();
1522
1523        let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(100_000);
1524        // 自动创建临时目录,生命周期结束自动删除
1525        let dir = tempdir().unwrap();
1526        let path = dir.path();
1527
1528        // spawn a task to send entries
1529        tokio::spawn(async move {
1530            let entries = entries.lock().await;
1531            for entry in entries.iter() {
1532                entry_tx
1533                    .send(MetaAttached {
1534                        inner: entry.clone(),
1535                        meta: EntryMeta::new(),
1536                    })
1537                    .await
1538                    .unwrap();
1539            }
1540            drop(entry_tx);
1541            tracing::info!("all entries sent");
1542        });
1543
1544        encode_and_output_to_files(entry_rx, entries_number, path.to_path_buf(), 0)
1545            .await
1546            .unwrap();
1547
1548        // 验证临时目录下生成的 pack/idx 文件
1549        let mut pack_file = None;
1550        let mut idx_file = None;
1551        for entry in std::fs::read_dir(path).unwrap() {
1552            let entry = entry.unwrap();
1553            let file_name = entry.file_name();
1554            tracing::info!("file name: {:?}", file_name);
1555            let file_name = file_name.to_string_lossy();
1556            if file_name.ends_with(".pack") {
1557                pack_file = Some(entry.path());
1558            } else if file_name.ends_with(".idx") {
1559                idx_file = Some(entry.path());
1560            }
1561        }
1562        let pack_file = pack_file.expect("pack file not generated");
1563        let idx_file = idx_file.expect("idx file not generated");
1564        assert!(
1565            pack_file.metadata().unwrap().len() > 0,
1566            "pack file is empty"
1567        );
1568        assert!(idx_file.metadata().unwrap().len() > 0, "idx file is empty");
1569
1570        let duration = start.elapsed();
1571        tracing::info!("test executed in: {:.2?}", duration);
1572        tracing::info!("original total size: {}", total_original_size);
1573    }
1574
1575    #[tokio::test]
1576    async fn test_pack_encoder_output_to_files_with_delta() {
1577        let _guard = set_hash_kind_for_test(HashKind::Sha1);
1578        init_logger();
1579        let (entries, _dl_guard) = get_entries_for_test().await;
1580        let entries_number = entries.lock().await.len();
1581
1582        let total_original_size: usize = entries
1583            .lock()
1584            .await
1585            .iter()
1586            .map(|entry| entry.data.len())
1587            .sum();
1588
1589        let start = Instant::now();
1590
1591        let (entry_tx, entry_rx) = mpsc::channel::<MetaAttached<Entry, EntryMeta>>(100_000);
1592        // 自动创建临时目录,生命周期结束自动删除
1593        let dir = tempdir().unwrap();
1594        let path = dir.path();
1595
1596        // spawn a task to send entries
1597        tokio::spawn(async move {
1598            let entries = entries.lock().await;
1599            for entry in entries.iter() {
1600                entry_tx
1601                    .send(MetaAttached {
1602                        inner: entry.clone(),
1603                        meta: EntryMeta::new(),
1604                    })
1605                    .await
1606                    .unwrap();
1607            }
1608            drop(entry_tx);
1609            tracing::info!("all entries sent");
1610        });
1611
1612        encode_and_output_to_files(entry_rx, entries_number, path.to_path_buf(), 10)
1613            .await
1614            .unwrap();
1615
1616        // 验证临时目录下生成的 pack/idx 文件
1617        let mut pack_file = None;
1618        let mut idx_file = None;
1619        for entry in std::fs::read_dir(path).unwrap() {
1620            let entry = entry.unwrap();
1621            let file_name = entry.file_name();
1622            tracing::info!("file name: {:?}", file_name);
1623            let file_name = file_name.to_string_lossy();
1624            if file_name.ends_with(".pack") {
1625                pack_file = Some(entry.path());
1626            } else if file_name.ends_with(".idx") {
1627                idx_file = Some(entry.path());
1628            }
1629        }
1630        let pack_file = pack_file.expect("pack file not generated");
1631        let idx_file = idx_file.expect("idx file not generated");
1632        assert!(
1633            pack_file.metadata().unwrap().len() > 0,
1634            "pack file is empty"
1635        );
1636        assert!(idx_file.metadata().unwrap().len() > 0, "idx file is empty");
1637
1638        let duration = start.elapsed();
1639        tracing::info!("test executed in: {:.2?}", duration);
1640        tracing::info!("original total size: {}", total_original_size);
1641    }
1642}