Now the restore chain also forwards the block id

This commit is contained in:
Carl Johnsen
2024-11-12 13:45:32 +01:00
parent b686cd4934
commit 09e04e76d0
5 changed files with 12 additions and 17 deletions
@@ -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<long, TaskCompletionSource<byte[]>> _waiters = new();
private readonly ConcurrentDictionary<long, bool> _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;
@@ -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<Database.LocalRestoreDatabase.IFileToRestore> 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));
}
}
@@ -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)
@@ -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)
@@ -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)