diff --git a/Duplicati/Library/Main/Operation/Restore/VolumeDecompressor.cs b/Duplicati/Library/Main/Operation/Restore/VolumeDecompressor.cs new file mode 100644 index 000000000..b7d7722d5 --- /dev/null +++ b/Duplicati/Library/Main/Operation/Restore/VolumeDecompressor.cs @@ -0,0 +1,45 @@ +using System; +using System.Linq; +using System.Threading.Tasks; +using CoCoL; +using Duplicati.Library.Main.Database; +using Duplicati.Library.Main.Volumes; + +namespace Duplicati.Library.Main.Operation.Restore +{ + internal class VolumeDecompressor + { + public static Task Run(LocalRestoreDatabase db, Options options) + { + return AutomationExtensions.RunTask( + new + { + Input = Channels.decryptedVolume.ForRead, + Output = Channels.decompressedVolumes.ForWrite + }, + async self => + { + while (true) + { + var (volume_id, volume) = await self.Input.ReadAsync(); + + var bids = db.Connection.CreateCommand().ExecuteReaderEnumerable(@$"SELECT ID, Hash, Size FROM Block WHERE VolumeID = ""{volume_id}""").Select(x => (x.GetInt64(0), x.GetString(1), x.GetInt64(2))).ToArray(); + + using (var blocks = new BlockVolumeReader(options.CompressionModule, volume, options)) + { + var volume_blocks = blocks.Blocks.ToArray(); + for (int i = 0; i < volume_blocks.Length; i++) + { + byte[] buffer = new byte[options.Blocksize]; + + System.Diagnostics.Debug.Assert(volume_blocks[i].Key == bids[i].Item2); + blocks.ReadBlock(volume_blocks[i].Key, buffer); + + await self.Output.WriteAsync((bids[i].Item1, buffer[..(int)bids[i].Item3])); + } + } + } + }); + } + } +} \ No newline at end of file diff --git a/Duplicati/Library/Main/Operation/RestoreHandler.cs b/Duplicati/Library/Main/Operation/RestoreHandler.cs index 1687f278a..b919cb41d 100644 --- a/Duplicati/Library/Main/Operation/RestoreHandler.cs +++ b/Duplicati/Library/Main/Operation/RestoreHandler.cs @@ -145,11 +145,11 @@ namespace Duplicati.Library.Main.Operation all = Task.WhenAll( [ Restore.FileLister.Run(db, backend, filter, m_options, m_result), - ..Enumerable.Range(0, parallelism).Select(i => - Restore.FileProcessor.Run(db, fileprocessor_requests[i], fileprocessor_responses[i])), + ..Enumerable.Range(0, parallelism).Select(i => Restore.FileProcessor.Run(db, fileprocessor_requests[i], fileprocessor_responses[i])), Restore.BlockManager.Run(db, fileprocessor_requests, fileprocessor_responses), ..Enumerable.Range(0, parallelism).Select(i => Restore.VolumeDownloader.Run(backend, m_options)), - ..Enumerable.Range(0, parallelism).Select(i => Restore.VolumeDecrypter.Run()) + ..Enumerable.Range(0, parallelism).Select(i => Restore.VolumeDecrypter.Run()), + ..Enumerable.Range(0, parallelism).Select(i => Restore.VolumeDecompressor.Run(db, m_options)) ] ); }