@@ -417,6 +417,76 @@ impl FsRepository {
417417 }
418418 }
419419
420+ /// Rewrites an uncompressed blob using maximum XZ compression.
421+ #[ cfg( feature = "tokio" ) ]
422+ async fn compact_blob ( & self , id : & Id ) -> Result < ( ) , RepositoryError > {
423+ use async_compression:: { Level , tokio:: write:: XzEncoder } ;
424+ use bitcache_core:: tokio:: {
425+ fs:: File ,
426+ io:: { AsyncWriteExt , copy} ,
427+ } ;
428+
429+ let source_name = Self :: path ( id) ;
430+ let source_file = match self . 0 . open ( & source_name) {
431+ Ok ( file) => file,
432+ Err ( error) if error. kind ( ) == std:: io:: ErrorKind :: NotFound => return Ok ( ( ) ) ,
433+ Err ( error) => return Err ( error. into ( ) ) ,
434+ } ;
435+ let extended = self . read_extended ( & source_name, & source_file) ?;
436+ let mut source_file = File :: from_std ( source_file. into_std ( ) ) ;
437+ let ( temp_name, temp_file) = self . create_temp_file ( ) ?;
438+ let temp_file = File :: from_std ( temp_file. into_std ( ) ) ;
439+ let mut encoder = XzEncoder :: with_quality ( temp_file, Level :: Best ) ;
440+
441+ let result: Result < ( ) , RepositoryError > = async {
442+ copy ( & mut source_file, & mut encoder) . await ?;
443+ encoder. shutdown ( ) . await ?;
444+ Ok ( ( ) )
445+ }
446+ . await ;
447+ drop ( encoder) ;
448+ drop ( source_file) ;
449+
450+ if let Err ( error) = result {
451+ let _ = self . 0 . remove_file ( & temp_name) ;
452+ return Err ( error) ;
453+ }
454+
455+ let temp_file = match self . 0 . open ( & temp_name) {
456+ Ok ( file) => file,
457+ Err ( error) => {
458+ let _ = self . 0 . remove_file ( & temp_name) ;
459+ return Err ( error. into ( ) ) ;
460+ } ,
461+ } ;
462+ if let Err ( error) = self . prepare_temp ( & temp_name, & temp_file, & extended) {
463+ drop ( temp_file) ;
464+ let _ = self . 0 . remove_file ( & temp_name) ;
465+ return Err ( error) ;
466+ }
467+ drop ( temp_file) ;
468+
469+ let compressed_name = Self :: compressed_path ( id) ;
470+ match self . 0 . remove_file ( & compressed_name) {
471+ Ok ( ( ) ) => ( ) ,
472+ Err ( error) if error. kind ( ) == std:: io:: ErrorKind :: NotFound => ( ) ,
473+ Err ( error) => {
474+ let _ = self . 0 . remove_file ( & temp_name) ;
475+ return Err ( error. into ( ) ) ;
476+ } ,
477+ }
478+ if let Err ( error) = self . 0 . hard_link ( & temp_name, & self . 0 , & compressed_name) {
479+ let _ = self . 0 . remove_file ( & temp_name) ;
480+ return Err ( error. into ( ) ) ;
481+ }
482+ self . 0 . remove_file ( & temp_name) ?;
483+ match self . 0 . remove_file ( source_name) {
484+ Ok ( ( ) ) => Ok ( ( ) ) ,
485+ Err ( error) if error. kind ( ) == std:: io:: ErrorKind :: NotFound => Ok ( ( ) ) ,
486+ Err ( error) => Err ( error. into ( ) ) ,
487+ }
488+ }
489+
420490 /// Collects all physical blob IDs, including expired blobs.
421491 fn collect_physical_ids ( & self ) -> Result < Vec < Id > , RepositoryError > {
422492 let mut ids = Vec :: new ( ) ;
@@ -758,6 +828,19 @@ impl Repository for FsRepository {
758828 self . write_existing_extended ( & blob. name , & blob. file , & blob. extended )
759829 }
760830
831+ async fn compact ( & mut self ) -> Result < ( ) , Self :: Error > {
832+ #[ cfg( feature = "tokio" ) ]
833+ {
834+ for id in self . collect_physical_ids ( ) ? {
835+ self . compact_blob ( & id) . await ?;
836+ }
837+ Ok ( ( ) )
838+ }
839+
840+ #[ cfg( not( feature = "tokio" ) ) ]
841+ Ok ( ( ) )
842+ }
843+
761844 async fn clear ( & mut self ) -> Result < ( ) , Self :: Error > {
762845 for id in self . collect_physical_ids ( ) ? {
763846 for encoding in [ BlobEncoding :: Uncompressed , BlobEncoding :: Xz ] {
0 commit comments