From 09e04e76d0a233cbc104f8ae4291210f4c57fc80 Mon Sep 17 00:00:00 2001 From: Carl Johnsen Date: Tue, 12 Nov 2024 13:45:32 +0100 Subject: [PATCH] Now the restore chain also forwards the block id --- Duplicati/Library/Main/Operation/Restore/BlockManager.cs | 4 ++-- Duplicati/Library/Main/Operation/Restore/Channels.cs | 8 ++++---- .../Library/Main/Operation/Restore/VolumeDecompressor.cs | 9 ++------- .../Library/Main/Operation/Restore/VolumeDecrypter.cs | 4 ++-- .../Library/Main/Operation/Restore/VolumeDownloader.cs | 4 ++-- 5 files changed, 12 insertions(+), 17 deletions(-) diff --git a/Duplicati/Library/Main/Operation/Restore/BlockManager.cs b/Duplicati/Library/Main/Operation/Restore/BlockManager.cs index 6f9318557..67f4f1a71 100644 --- a/Duplicati/Library/Main/Operation/Restore/BlockManager.cs +++ b/Duplicati/Library/Main/Operation/Restore/BlockManager.cs @@ -13,12 +13,12 @@ namespace Duplicati.Library.Main.Operation.Restore internal class SleepableDictionary // That also auto requests! { private readonly LocalRestoreDatabase m_db; - private readonly IWriteChannel<(long,IRemoteVolume)> m_volume_request; + private readonly IWriteChannel<(long,long,IRemoteVolume)> m_volume_request; private readonly MemoryCache _dictionary; private readonly ConcurrentDictionary> _waiters = new(); private readonly ConcurrentDictionary _in_flight = new(); - public SleepableDictionary(LocalRestoreDatabase db, IWriteChannel<(long,IRemoteVolume)> volume_request) + public SleepableDictionary(LocalRestoreDatabase db, IWriteChannel<(long,long,IRemoteVolume)> volume_request) { m_db = db; m_volume_request = volume_request; diff --git a/Duplicati/Library/Main/Operation/Restore/Channels.cs b/Duplicati/Library/Main/Operation/Restore/Channels.cs index b00e61b6f..34f12215d 100644 --- a/Duplicati/Library/Main/Operation/Restore/Channels.cs +++ b/Duplicati/Library/Main/Operation/Restore/Channels.cs @@ -6,12 +6,12 @@ namespace Duplicati.Library.Main.Operation.Restore internal static class Channels { // TODO Should maybe come from Options, or at least some global configuration file? - private static int bufferSize = 0; + private static int bufferSize = 128; - public static readonly ChannelMarkerWrapper<(long, Database.IRemoteVolume)> downloadRequest = new(new ChannelNameAttribute("downloadRequest", bufferSize)); - public static readonly ChannelMarkerWrapper<(long, TempFile)> downloadedVolume = new(new ChannelNameAttribute("downloadResponse", bufferSize)); + public static readonly ChannelMarkerWrapper<(long, long, Database.IRemoteVolume)> downloadRequest = new(new ChannelNameAttribute("downloadRequest", bufferSize)); + public static readonly ChannelMarkerWrapper<(long, long, TempFile)> downloadedVolume = new(new ChannelNameAttribute("downloadResponse", bufferSize)); public static readonly ChannelMarkerWrapper filesToRestore = new(new ChannelNameAttribute("filesToRestore", bufferSize)); public static readonly ChannelMarkerWrapper<(long, byte[])> decompressedVolumes = new(new ChannelNameAttribute("decompressedVolumes", bufferSize)); - public static readonly ChannelMarkerWrapper<(long, TempFile)> decryptedVolume = new(new ChannelNameAttribute("decrytedVolume", bufferSize)); + public static readonly ChannelMarkerWrapper<(long, long, TempFile)> decryptedVolume = new(new ChannelNameAttribute("decrytedVolume", bufferSize)); } } \ No newline at end of file diff --git a/Duplicati/Library/Main/Operation/Restore/VolumeDecompressor.cs b/Duplicati/Library/Main/Operation/Restore/VolumeDecompressor.cs index 08f073257..61425527e 100644 --- a/Duplicati/Library/Main/Operation/Restore/VolumeDecompressor.cs +++ b/Duplicati/Library/Main/Operation/Restore/VolumeDecompressor.cs @@ -22,7 +22,7 @@ namespace Duplicati.Library.Main.Operation.Restore try { while (true) { - var (volume_id, volume) = await self.Input.ReadAsync(); + var (block_id, 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(); var bid_lut = bids.ToDictionary(x => x.Item2); @@ -33,12 +33,7 @@ namespace Duplicati.Library.Main.Operation.Restore for (int i = 0; i < volume_blocks.Length; i++) { byte[] buffer = new byte[options.Blocksize]; - - blocks.ReadBlock(volume_blocks[i].Key, buffer); - - await self.Output.WriteAsync((bid_lut[volume_blocks[i].Key].Item1, buffer[..(int)volume_blocks[i].Value])); - } - } + await self.Output.WriteAsync((block_id, buffer[..(int)bsize])); } } catch (RetiredException ex) diff --git a/Duplicati/Library/Main/Operation/Restore/VolumeDecrypter.cs b/Duplicati/Library/Main/Operation/Restore/VolumeDecrypter.cs index 02ca9a45e..1998c2a49 100644 --- a/Duplicati/Library/Main/Operation/Restore/VolumeDecrypter.cs +++ b/Duplicati/Library/Main/Operation/Restore/VolumeDecrypter.cs @@ -20,9 +20,9 @@ namespace Duplicati.Library.Main.Operation.Restore { while (true) { - var (volume_id, volume) = await self.Input.ReadAsync(); + var (block_id, volume_id, volume) = await self.Input.ReadAsync(); // NOP operation for now - decryption is handled by the backend during download. Should be done here to increase concurrency. - await self.Output.WriteAsync((volume_id, volume)); + await self.Output.WriteAsync((block_id, volume_id, volume)); } } catch (RetiredException ex) diff --git a/Duplicati/Library/Main/Operation/Restore/VolumeDownloader.cs b/Duplicati/Library/Main/Operation/Restore/VolumeDownloader.cs index eda1f9fc4..20c2c9604 100644 --- a/Duplicati/Library/Main/Operation/Restore/VolumeDownloader.cs +++ b/Duplicati/Library/Main/Operation/Restore/VolumeDownloader.cs @@ -21,7 +21,7 @@ namespace Duplicati.Library.Main.Operation.Restore { while (true) { - var (volume_id, request) = await self.Input.ReadAsync(); + var (block_id, volume_id, request) = await self.Input.ReadAsync(); TempFile f = null; try @@ -34,7 +34,7 @@ namespace Duplicati.Library.Main.Operation.Restore throw; } - self.Output.Write((volume_id,f)); + self.Output.Write((block_id, volume_id, f)); } } catch (RetiredException ex)