The concurrent processing now seems to work, at least in the base case.
Rewrote the handling of index volumes to create them on-demand instead of in-advance. This makes it much simpler to handle upload failures.
This commit is contained in:
@@ -110,9 +110,9 @@ namespace Duplicati.Library.Main.Database
|
||||
m_updateremotevolumeCommand.CommandText = @"UPDATE ""Remotevolume"" SET ""OperationID"" = ?, ""State"" = ?, ""Hash"" = ?, ""Size"" = ? WHERE ""Name"" = ?";
|
||||
m_updateremotevolumeCommand.AddParameters(5);
|
||||
|
||||
m_selectremotevolumesCommand.CommandText = @"SELECT ""Name"", ""Type"", ""Size"", ""Hash"", ""State"", ""DeleteGraceTime"" FROM ""Remotevolume""";
|
||||
m_selectremotevolumesCommand.CommandText = @"SELECT ""ID"", ""Name"", ""Type"", ""Size"", ""Hash"", ""State"", ""DeleteGraceTime"" FROM ""Remotevolume""";
|
||||
|
||||
m_selectremotevolumeCommand.CommandText = @"SELECT ""Type"", ""Size"", ""Hash"", ""State"" FROM ""Remotevolume"" WHERE ""Name"" = ?";
|
||||
m_selectremotevolumeCommand.CommandText = m_selectremotevolumesCommand.CommandText + @" WHERE ""Name"" = ?";
|
||||
m_selectremotevolumeCommand.AddParameter();
|
||||
|
||||
m_removeremotevolumeCommand.CommandText = @"DELETE FROM ""Remotevolume"" WHERE ""Name"" = ?";
|
||||
@@ -242,39 +242,40 @@ namespace Duplicati.Library.Main.Database
|
||||
return m_selectremotevolumeIdCommand.ExecuteScalarInt64(null, -1, file);
|
||||
}
|
||||
|
||||
public bool GetRemoteVolume(string file, out string hash, out long size, out RemoteVolumeType type, out RemoteVolumeState state)
|
||||
public RemoteVolumeEntry GetRemoteVolume(string file, System.Data.IDbTransaction transaction = null)
|
||||
{
|
||||
m_selectremotevolumeCommand.Transaction = transaction;
|
||||
m_selectremotevolumeCommand.SetParameterValue(0, file);
|
||||
using (var rd = m_selectremotevolumeCommand.ExecuteReader())
|
||||
using(var rd = m_selectremotevolumeCommand.ExecuteReader())
|
||||
if (rd.Read())
|
||||
{
|
||||
hash = (rd.GetValue(2) == null || rd.GetValue(2) == DBNull.Value) ? null : rd.GetValue(3).ToString();
|
||||
size = (rd.GetValue(1) == null || rd.GetValue(1) == DBNull.Value) ? -1 : rd.GetInt64(2);
|
||||
type = (RemoteVolumeType)Enum.Parse(typeof(RemoteVolumeType), rd.GetValue(0).ToString());
|
||||
state = (RemoteVolumeState)Enum.Parse(typeof(RemoteVolumeState), rd.GetValue(3).ToString());
|
||||
return true;
|
||||
}
|
||||
return new RemoteVolumeEntry(
|
||||
rd.ConvertValueToInt64(0),
|
||||
rd.GetValue(1).ToString(),
|
||||
(rd.GetValue(4) == null || rd.GetValue(4) == DBNull.Value) ? null : rd.GetValue(4).ToString(),
|
||||
rd.ConvertValueToInt64(3, -1),
|
||||
(RemoteVolumeType)Enum.Parse(typeof(RemoteVolumeType), rd.GetValue(2).ToString()),
|
||||
(RemoteVolumeState)Enum.Parse(typeof(RemoteVolumeState), rd.GetValue(5).ToString()),
|
||||
new DateTime(rd.ConvertValueToInt64(6, 0), DateTimeKind.Utc)
|
||||
);
|
||||
|
||||
hash = null;
|
||||
size = -1;
|
||||
type = (RemoteVolumeType)(-1);
|
||||
state = (RemoteVolumeState)(-1);
|
||||
return false;
|
||||
return RemoteVolumeEntry.Empty;
|
||||
}
|
||||
|
||||
public IEnumerable<RemoteVolumeEntry> GetRemoteVolumes()
|
||||
public IEnumerable<RemoteVolumeEntry> GetRemoteVolumes(System.Data.IDbTransaction transaction = null)
|
||||
{
|
||||
m_selectremotevolumesCommand.Transaction = transaction;
|
||||
using (var rd = m_selectremotevolumesCommand.ExecuteReader())
|
||||
{
|
||||
while (rd.Read())
|
||||
{
|
||||
yield return new RemoteVolumeEntry(
|
||||
rd.GetValue(0).ToString(),
|
||||
(rd.GetValue(3) == null || rd.GetValue(3) == DBNull.Value) ? null : rd.GetValue(3).ToString(),
|
||||
rd.ConvertValueToInt64(2, -1),
|
||||
(RemoteVolumeType)Enum.Parse(typeof(RemoteVolumeType), rd.GetValue(1).ToString()),
|
||||
(RemoteVolumeState)Enum.Parse(typeof(RemoteVolumeState), rd.GetValue(4).ToString()),
|
||||
new DateTime(rd.ConvertValueToInt64(5, 0), DateTimeKind.Utc)
|
||||
rd.ConvertValueToInt64(0),
|
||||
rd.GetValue(1).ToString(),
|
||||
(rd.GetValue(4) == null || rd.GetValue(4) == DBNull.Value) ? null : rd.GetValue(4).ToString(),
|
||||
rd.ConvertValueToInt64(3, -1),
|
||||
(RemoteVolumeType)Enum.Parse(typeof(RemoteVolumeType), rd.GetValue(2).ToString()),
|
||||
(RemoteVolumeState)Enum.Parse(typeof(RemoteVolumeState), rd.GetValue(5).ToString()),
|
||||
new DateTime(rd.ConvertValueToInt64(6, 0), DateTimeKind.Utc)
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -943,9 +944,9 @@ namespace Duplicati.Library.Main.Database
|
||||
m_insertIndexBlockLink.ExecuteNonQuery();
|
||||
}
|
||||
|
||||
public IEnumerable<Tuple<string, byte[], int>> GetBlocklists(long volumeid, long blocksize, int hashsize)
|
||||
public IEnumerable<Tuple<string, byte[], int>> GetBlocklists(long volumeid, long blocksize, int hashsize, System.Data.IDbTransaction transaction = null)
|
||||
{
|
||||
using(var cmd = m_connection.CreateCommand())
|
||||
using(var cmd = m_connection.CreateCommand(transaction))
|
||||
{
|
||||
var sql = string.Format(@"SELECT ""A"".""Hash"", ""C"".""Hash"" FROM " +
|
||||
@"(SELECT ""BlocklistHash"".""BlocksetID"", ""Block"".""Hash"", * FROM ""BlocklistHash"",""Block"" WHERE ""BlocklistHash"".""Hash"" = ""Block"".""Hash"" AND ""Block"".""VolumeID"" = ?) A, " +
|
||||
|
||||
@@ -7,28 +7,25 @@ namespace Duplicati.Library.Main.Database
|
||||
{
|
||||
public struct RemoteVolumeEntry : IRemoteVolume
|
||||
{
|
||||
private readonly string m_name;
|
||||
private readonly string m_hash;
|
||||
private readonly long m_size;
|
||||
private readonly RemoteVolumeType m_type;
|
||||
private readonly RemoteVolumeState m_state;
|
||||
private readonly DateTime m_deleteGracePeriod;
|
||||
|
||||
public string Name { get { return m_name; } }
|
||||
public string Hash { get { return m_hash; } }
|
||||
public long Size { get { return m_size; } }
|
||||
public RemoteVolumeType Type { get { return m_type; } }
|
||||
public RemoteVolumeState State { get { return m_state; } }
|
||||
public DateTime deleteGracePeriod { get { return m_deleteGracePeriod; } }
|
||||
public long ID { get; private set; }
|
||||
public string Name { get; private set; }
|
||||
public string Hash { get; private set; }
|
||||
public long Size { get; private set; }
|
||||
public RemoteVolumeType Type { get; private set; }
|
||||
public RemoteVolumeState State { get; private set; }
|
||||
public DateTime DeleteGracePeriod { get; private set; }
|
||||
|
||||
public RemoteVolumeEntry(string name, string hash, long size, RemoteVolumeType type, RemoteVolumeState state, DateTime deleteGracePeriod)
|
||||
public static readonly RemoteVolumeEntry Empty = new RemoteVolumeEntry(-1, null, null, -1, (RemoteVolumeType)(-1), (RemoteVolumeState)(-1), default(DateTime));
|
||||
|
||||
public RemoteVolumeEntry(long id, string name, string hash, long size, RemoteVolumeType type, RemoteVolumeState state, DateTime deleteGracePeriod)
|
||||
{
|
||||
m_name = name;
|
||||
m_size = size;
|
||||
m_type = type;
|
||||
m_state = state;
|
||||
m_hash = hash;
|
||||
m_deleteGracePeriod = deleteGracePeriod;
|
||||
ID = id;
|
||||
Name = name;
|
||||
Size = size;
|
||||
Type = type;
|
||||
State = state;
|
||||
Hash = hash;
|
||||
DeleteGracePeriod = deleteGracePeriod;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -121,7 +121,6 @@
|
||||
<Compile Include="Operation\Backup\DataBlock.cs" />
|
||||
<Compile Include="Operation\Backup\FilePreFilterProcess.cs" />
|
||||
<Compile Include="Operation\Backup\DataBlockProcessor.cs" />
|
||||
<Compile Include="Operation\Common\BackendOperation.cs" />
|
||||
<Compile Include="Operation\Backup\SpillCollectorProcess.cs" />
|
||||
<Compile Include="Operation\Backup\BackupDatabase.cs" />
|
||||
<Compile Include="Operation\Common\DatabaseCommon.cs" />
|
||||
@@ -131,6 +130,8 @@
|
||||
<Compile Include="Operation\Backup\BackupStatsCollector.cs" />
|
||||
<Compile Include="Operation\Common\LogHandler.cs" />
|
||||
<Compile Include="Operation\Backup\ProgressHandler.cs" />
|
||||
<Compile Include="Operation\Backup\BackendUploader.cs" />
|
||||
<Compile Include="Operation\Common\IndexVolumeCreator.cs" />
|
||||
</ItemGroup>
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\Utility\Duplicati.Library.Utility.csproj">
|
||||
|
||||
@@ -0,0 +1,79 @@
|
||||
// Copyright (C) 2015, The Duplicati Team
|
||||
// http://www.duplicati.com, info@duplicati.com
|
||||
//
|
||||
// This library is free software; you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Lesser General Public License as
|
||||
// published by the Free Software Foundation; either version 2.1 of the
|
||||
// License, or (at your option) any later version.
|
||||
//
|
||||
// This library is distributed in the hope that it will be useful, but
|
||||
// WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
|
||||
// Lesser General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Lesser General Public
|
||||
// License along with this library; if not, write to the Free Software
|
||||
// Foundation, Inc., 59 Temple Place, Suite 330, Boston, MA 02111-1307 USA
|
||||
using System;
|
||||
using CoCoL;
|
||||
using System.Threading.Tasks;
|
||||
using Duplicati.Library.Main.Operation.Common;
|
||||
using Duplicati.Library.Main.Volumes;
|
||||
using System.Collections.Generic;
|
||||
|
||||
namespace Duplicati.Library.Main.Operation.Backup
|
||||
{
|
||||
internal class UploadRequest
|
||||
{
|
||||
public BlockVolumeWriter BlockVolume { get; private set; }
|
||||
public IndexVolumeWriter IndexVolume { get; private set; }
|
||||
|
||||
public UploadRequest(BlockVolumeWriter blockvolume, IndexVolumeWriter indexvolume)
|
||||
{
|
||||
BlockVolume = blockvolume;
|
||||
IndexVolume = indexvolume;
|
||||
}
|
||||
}
|
||||
|
||||
internal static class BackendUploader
|
||||
{
|
||||
public static Task Run(Common.BackendHandler backend, Options options, Common.DatabaseCommon database)
|
||||
{
|
||||
return AutomationExtensions.RunTask(new
|
||||
{
|
||||
Input = ChannelMarker.ForRead<UploadRequest>("BackendRequests"),
|
||||
},
|
||||
|
||||
async self =>
|
||||
{
|
||||
var inProgress = new Queue<Task>();
|
||||
var max_pending = options.AsynchronousUploadLimit == 0 ? long.MaxValue : options.AsynchronousUploadLimit;
|
||||
if (options.IndexfilePolicy != Options.IndexFileStrategy.None)
|
||||
max_pending = max_pending / 2;
|
||||
|
||||
while(!self.Input.IsRetired)
|
||||
{
|
||||
try
|
||||
{
|
||||
var req = await self.Input.ReadAsync();
|
||||
inProgress.Enqueue(backend.UploadFileAsync(req.BlockVolume, name => IndexVolumeCreator.CreateIndexVolume(name, options, database)));
|
||||
}
|
||||
catch(Exception ex)
|
||||
{
|
||||
if (!ex.IsRetiredException())
|
||||
throw;
|
||||
}
|
||||
|
||||
while(inProgress.Count >= max_pending)
|
||||
await inProgress.Dequeue();
|
||||
}
|
||||
|
||||
while(inProgress.Count > 0)
|
||||
await inProgress.Dequeue();
|
||||
}
|
||||
);
|
||||
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -45,7 +45,8 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
TaskCompletion = tcs
|
||||
});
|
||||
|
||||
return await tcs.Task;
|
||||
var r = await tcs.Task;
|
||||
return r;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -37,14 +37,13 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
{
|
||||
LogChannel = ChannelMarker.ForWrite<LogMessage>("LogChannel"),
|
||||
Input = ChannelMarker.ForRead<DataBlock>("OutputBlocks"),
|
||||
Output = ChannelMarker.ForWrite<IBackendOperation>("BackendRequests"),
|
||||
SpillPickup = ChannelMarker.ForWrite<IBackendOperation>("SpillPickup"),
|
||||
Output = ChannelMarker.ForWrite<UploadRequest>("BackendRequests"),
|
||||
SpillPickup = ChannelMarker.ForWrite<UploadRequest>("SpillPickup"),
|
||||
},
|
||||
|
||||
async self =>
|
||||
{
|
||||
BlockVolumeWriter blockvolume = null;
|
||||
IndexVolumeWriter indexvolume = null;
|
||||
|
||||
try
|
||||
{
|
||||
@@ -68,12 +67,6 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
|
||||
blockvolume = new BlockVolumeWriter(options);
|
||||
blockvolume.VolumeID = await database.RegisterRemoteVolumeAsync(blockvolume.RemoteFilename, RemoteVolumeType.Blocks, RemoteVolumeState.Temporary);
|
||||
|
||||
if (options.IndexfilePolicy != Options.IndexFileStrategy.None)
|
||||
{
|
||||
indexvolume = new IndexVolumeWriter(options);
|
||||
indexvolume.VolumeID = await database.RegisterRemoteVolumeAsync(indexvolume.RemoteFilename, RemoteVolumeType.Index, RemoteVolumeState.Temporary);
|
||||
}
|
||||
}
|
||||
|
||||
var newBlock = await database.AddBlockAsync(b.HashKey, b.Size, blockvolume.VolumeID);
|
||||
@@ -84,13 +77,6 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
|
||||
blockvolume.AddBlock(b.HashKey, b.Data, b.Offset, (int)b.Size, b.Hint);
|
||||
|
||||
//TODO: In theory a normal data block and blocklist block could be equal.
|
||||
// this would cause the index file to not contain all data,
|
||||
// if the data file is added before the blocklist data
|
||||
// ... highly theoretical and only causes extra block data downloads ...
|
||||
if (options.IndexfilePolicy == Options.IndexFileStrategy.Full && b.IsBlocklistHashes)
|
||||
indexvolume.WriteBlocklist(b.HashKey, b.Data, b.Offset, (int)b.Size);
|
||||
|
||||
if (blockvolume.Filesize > options.VolumeSize - options.Blocksize)
|
||||
{
|
||||
if (options.Dryrun)
|
||||
@@ -98,34 +84,21 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
blockvolume.Close();
|
||||
await self.LogChannel.WriteAsync(LogMessage.DryRun("Would upload block volume: {0}, size: {1}", blockvolume.RemoteFilename, Library.Utility.Utility.FormatSizeString(new FileInfo(blockvolume.LocalFilename).Length)));
|
||||
|
||||
if (indexvolume != null)
|
||||
{
|
||||
await database.UpdateIndexVolumeAsync(indexvolume, blockvolume);
|
||||
indexvolume.FinishVolume(Library.Utility.Utility.CalculateHash(blockvolume.LocalFilename), new FileInfo(blockvolume.LocalFilename).Length);
|
||||
await self.LogChannel.WriteAsync(LogMessage.DryRun("Would upload index volume: {0}, size: {1}", indexvolume.RemoteFilename, Library.Utility.Utility.FormatSizeString(new FileInfo(indexvolume.LocalFilename).Length)));
|
||||
indexvolume.Dispose();
|
||||
indexvolume = null;
|
||||
}
|
||||
|
||||
blockvolume.Dispose();
|
||||
blockvolume = null;
|
||||
indexvolume.Dispose();
|
||||
indexvolume = null;
|
||||
}
|
||||
else
|
||||
{
|
||||
//When uploading a new volume, we register the volumes and then flush the transaction
|
||||
// this ensures that the local database and remote storage are as closely related as possible
|
||||
await database.UpdateRemoteVolume(blockvolume.RemoteFilename, RemoteVolumeState.Uploading, -1, null);
|
||||
await database.UpdateRemoteVolumeAsync(blockvolume.RemoteFilename, RemoteVolumeState.Uploading, -1, null);
|
||||
|
||||
blockvolume.Close();
|
||||
|
||||
await database.UpdateIndexVolumeAsync(indexvolume, blockvolume);
|
||||
await database.CommitTransactionAsync("CommitAddBlockToOutputFlush");
|
||||
|
||||
await self.Output.WriteAsync(new UploadRequest(blockvolume, indexvolume));
|
||||
await self.Output.WriteAsync(new UploadRequest(blockvolume, null));
|
||||
blockvolume = null;
|
||||
indexvolume = null;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -138,7 +111,7 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
{
|
||||
// If we have collected data, merge all pending volumes into a single volume
|
||||
if (blockvolume != null && blockvolume.SourceSize > 0)
|
||||
await self.SpillPickup.WriteAsync(new UploadRequest(blockvolume, indexvolume));
|
||||
await self.SpillPickup.WriteAsync(new UploadRequest(blockvolume, null));
|
||||
}
|
||||
|
||||
throw;
|
||||
|
||||
@@ -113,7 +113,7 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
|
||||
// Make sure the filehasher is done with the buf instance before we pass it on
|
||||
await pftask;
|
||||
await DataBlock.AddBlockToOutputAsync(self.BlockOutput, hashkey, buf, lastread, 0, hint, true);
|
||||
await DataBlock.AddBlockToOutputAsync(self.BlockOutput, hashkey, buf, 0, lastread, hint, true);
|
||||
buf = new byte[blocksize];
|
||||
}
|
||||
}
|
||||
|
||||
@@ -21,6 +21,7 @@ using Duplicati.Library.Main.Operation.Common;
|
||||
using System.Collections.Generic;
|
||||
using System.Linq;
|
||||
using Duplicati.Library.Main.Volumes;
|
||||
using System.IO;
|
||||
|
||||
namespace Duplicati.Library.Main.Operation.Backup
|
||||
{
|
||||
@@ -37,8 +38,8 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
return AutomationExtensions.RunTask(
|
||||
new
|
||||
{
|
||||
Input = ChannelMarker.ForRead<IBackendOperation>("SpillPickup"),
|
||||
Output = ChannelMarker.ForWrite<IBackendOperation>("BackendRequests"),
|
||||
Input = ChannelMarker.ForRead<UploadRequest>("SpillPickup"),
|
||||
Output = ChannelMarker.ForWrite<UploadRequest>("BackendRequests"),
|
||||
},
|
||||
|
||||
async self =>
|
||||
@@ -48,7 +49,7 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
while(!self.Input.IsRetired)
|
||||
try
|
||||
{
|
||||
lst.Add((UploadRequest)await self.Input.ReadAsync());
|
||||
lst.Add((UploadRequest)await self.Input.ReadAsync());
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
@@ -67,10 +68,6 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
// Finalize the current work
|
||||
source.BlockVolume.Close();
|
||||
|
||||
// We rebuild the index volume from the database
|
||||
if (source.IndexVolume != null)
|
||||
source.IndexVolume.Close();
|
||||
|
||||
// Remove it from the list of active operations
|
||||
lst.RemoveAt(0);
|
||||
|
||||
@@ -86,10 +83,8 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
if (lst.Count == 0)
|
||||
{
|
||||
// No more targets, make one
|
||||
target = new UploadRequest(new BlockVolumeWriter(options), options.IndexfilePolicy == Options.IndexFileStrategy.None ? null : new IndexVolumeWriter(options));
|
||||
target = new UploadRequest(new BlockVolumeWriter(options), null);
|
||||
target.BlockVolume.VolumeID = await database.RegisterRemoteVolumeAsync(target.BlockVolume.RemoteFilename, RemoteVolumeType.Blocks, RemoteVolumeState.Temporary);
|
||||
if (target.IndexVolume != null)
|
||||
target.IndexVolume.VolumeID = await database.RegisterRemoteVolumeAsync(target.IndexVolume.RemoteFilename, RemoteVolumeType.Index, RemoteVolumeState.Temporary);
|
||||
}
|
||||
else
|
||||
{
|
||||
@@ -106,23 +101,6 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
|
||||
if (target.BlockVolume.Filesize > options.VolumeSize - options.Blocksize)
|
||||
{
|
||||
if (options.IndexfilePolicy == Options.IndexFileStrategy.Full && target.IndexVolume != null && source.IndexVolume != null)
|
||||
{
|
||||
using(var ixr = DynamicLoader.CompressionLoader.GetModule(options.CompressionModule, source.IndexVolume.LocalFilename, options.RawOptions))
|
||||
foreach(var blocklisthash in await database.GetBlocklistHashesAsync(source.BlockVolume.RemoteFilename))
|
||||
{
|
||||
long fslen;
|
||||
using(var fs = ixr.OpenRead(blocklisthash))
|
||||
{
|
||||
target.IndexVolume.WriteBlocklist(blocklisthash, fs);
|
||||
fslen = fs.Length;
|
||||
}
|
||||
|
||||
await database.MoveBlockToVolumeAsync(blocklisthash, fslen, source.BlockVolume.VolumeID, target.BlockVolume.VolumeID);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
await self.Output.WriteAsync(target);
|
||||
target = null;
|
||||
}
|
||||
@@ -132,11 +110,6 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
// Make sure they are out of the database
|
||||
System.IO.File.Delete(source.BlockVolume.LocalFilename);
|
||||
await database.SafeDeleteRemoteVolumeAsync(source.BlockVolume.RemoteFilename);
|
||||
if (source.IndexVolume != null)
|
||||
{
|
||||
System.IO.File.Delete(source.IndexVolume.LocalFilename);
|
||||
await database.SafeDeleteRemoteVolumeAsync(source.IndexVolume.RemoteFilename);
|
||||
}
|
||||
|
||||
// Re-inject the target if it has content
|
||||
if (target != null)
|
||||
@@ -145,7 +118,10 @@ namespace Duplicati.Library.Main.Operation.Backup
|
||||
}
|
||||
|
||||
foreach(var n in lst)
|
||||
{
|
||||
n.BlockVolume.Close();
|
||||
await self.Output.WriteAsync(n);
|
||||
}
|
||||
|
||||
}
|
||||
);
|
||||
|
||||
@@ -275,7 +275,8 @@ namespace Duplicati.Library.Main.Operation
|
||||
Backup.FilePreFilterProcess.Start(snapshot, options),
|
||||
new Backup.MetadataPreProcess(snapshot, options, database).RunAsync(),
|
||||
Backup.SpillCollectorProcess.Run(options, database),
|
||||
Backup.ProgressHandler.Run(result)
|
||||
Backup.ProgressHandler.Run(result),
|
||||
Backup.BackendUploader.Run(backend, options, database)
|
||||
);
|
||||
}
|
||||
|
||||
@@ -288,57 +289,6 @@ namespace Duplicati.Library.Main.Operation
|
||||
}
|
||||
}
|
||||
|
||||
private long FinalizeRemoteVolumes(BackendManager backend)
|
||||
{
|
||||
var lastVolumeSize = -1L;
|
||||
m_result.OperationProgressUpdater.UpdatePhase(OperationPhase.Backup_Finalize);
|
||||
using(new Logging.Timer("FinalizeRemoteVolumes"))
|
||||
{
|
||||
if (m_blockvolume != null && m_blockvolume.SourceSize > 0)
|
||||
{
|
||||
lastVolumeSize = m_blockvolume.SourceSize;
|
||||
|
||||
if (m_options.Dryrun)
|
||||
{
|
||||
m_result.AddDryrunMessage(string.Format("Would upload block volume: {0}, size: {1}", m_blockvolume.RemoteFilename, Library.Utility.Utility.FormatSizeString(new FileInfo(m_blockvolume.LocalFilename).Length)));
|
||||
if (m_indexvolume != null)
|
||||
{
|
||||
m_blockvolume.Close();
|
||||
UpdateIndexVolume();
|
||||
m_indexvolume.FinishVolume(Library.Utility.Utility.CalculateHash(m_blockvolume.LocalFilename), new FileInfo(m_blockvolume.LocalFilename).Length);
|
||||
m_result.AddDryrunMessage(string.Format("Would upload index volume: {0}, size: {1}", m_indexvolume.RemoteFilename, Library.Utility.Utility.FormatSizeString(new FileInfo(m_indexvolume.LocalFilename).Length)));
|
||||
}
|
||||
|
||||
m_blockvolume.Dispose();
|
||||
m_blockvolume = null;
|
||||
m_indexvolume.Dispose();
|
||||
m_indexvolume = null;
|
||||
}
|
||||
else
|
||||
{
|
||||
m_database.UpdateRemoteVolume(m_blockvolume.RemoteFilename, RemoteVolumeState.Uploading, -1, null, m_transaction);
|
||||
m_blockvolume.Close();
|
||||
UpdateIndexVolume();
|
||||
|
||||
using(new Logging.Timer("CommitUpdateRemoteVolume"))
|
||||
m_transaction.Commit();
|
||||
m_transaction = m_database.BeginTransaction();
|
||||
|
||||
backend.Put(m_blockvolume, m_indexvolume);
|
||||
|
||||
using(new Logging.Timer("CommitUpdateRemoteVolume"))
|
||||
m_transaction.Commit();
|
||||
m_transaction = m_database.BeginTransaction();
|
||||
|
||||
m_blockvolume = null;
|
||||
m_indexvolume = null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return lastVolumeSize;
|
||||
}
|
||||
|
||||
private void UploadRealFileList(BackendManager backend, FilesetVolumeWriter filesetvolume)
|
||||
{
|
||||
var changeCount =
|
||||
@@ -474,6 +424,7 @@ namespace Duplicati.Library.Main.Operation
|
||||
using(var logtarget = ChannelManager.GetChannel<Common.LogMessage>("LogChannel").AsWriteOnly())
|
||||
{
|
||||
var lh = Common.LogHandler.Run(m_result);
|
||||
long lastVolumeSize = -1L;
|
||||
|
||||
using(var snapshot = GetSnapshot(sources, m_options, m_result))
|
||||
{
|
||||
@@ -519,7 +470,8 @@ namespace Duplicati.Library.Main.Operation
|
||||
}
|
||||
}
|
||||
|
||||
var lastVolumeSize = FinalizeRemoteVolumes(backend);
|
||||
// TODO: Implement this
|
||||
//var lastVolumeSize = FinalizeRemoteVolumes(backend);
|
||||
|
||||
using(new Logging.Timer("UpdateChangeStatistics"))
|
||||
m_database.UpdateChangeStatistics(m_result);
|
||||
@@ -574,7 +526,7 @@ namespace Duplicati.Library.Main.Operation
|
||||
}
|
||||
finally
|
||||
{
|
||||
if (parallelScanner != null)
|
||||
if (parallelScanner != null && !parallelScanner.IsCompleted)
|
||||
parallelScanner.Wait(500);
|
||||
|
||||
// TODO: We want to commit? always?
|
||||
|
||||
@@ -67,10 +67,6 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
/// </summary>
|
||||
public long Size;
|
||||
/// <summary>
|
||||
/// Reference to the index file entry that is updated if this entry changes
|
||||
/// </summary>
|
||||
public Tuple<IndexVolumeWriter, FileEntryItem> Indexfile;
|
||||
/// <summary>
|
||||
/// A flag indicating if the final hash and size of the block volume has been written to the index file
|
||||
/// </summary>
|
||||
public bool IndexfileUpdated;
|
||||
@@ -88,16 +84,15 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
/// </summary>
|
||||
public bool IsRetry;
|
||||
|
||||
public FileEntryItem(BackendActionType operation, string remotefilename, Tuple<IndexVolumeWriter, FileEntryItem> indexfile = null)
|
||||
public FileEntryItem(BackendActionType operation, string remotefilename)
|
||||
{
|
||||
Operation = operation;
|
||||
RemoteFilename = remotefilename;
|
||||
Indexfile = indexfile;
|
||||
Size = -1;
|
||||
}
|
||||
|
||||
public FileEntryItem(BackendActionType operation, string remotefilename, long size, string hash, Tuple<IndexVolumeWriter, FileEntryItem> indexfile = null)
|
||||
: this(operation, remotefilename, indexfile)
|
||||
public FileEntryItem(BackendActionType operation, string remotefilename, long size, string hash)
|
||||
: this(operation, remotefilename)
|
||||
{
|
||||
Size = size;
|
||||
Hash = hash;
|
||||
@@ -185,7 +180,7 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
|
||||
public Task PutUnencryptedAsync(string remotename, string localpath)
|
||||
{
|
||||
var fe = new FileEntryItem(BackendActionType.Put, remotename, null);
|
||||
var fe = new FileEntryItem(BackendActionType.Put, remotename);
|
||||
fe.SetLocalfilename(localpath);
|
||||
fe.Encrypted = true; //Prevent encryption
|
||||
fe.TrackedInDb = false; //Prevent Db updates
|
||||
@@ -199,31 +194,57 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
|
||||
}
|
||||
|
||||
public Task UploadFileAsync(VolumeWriterBase item, IndexVolumeWriter indexfile = null)
|
||||
public async Task UploadFileAsync(VolumeWriterBase item, Func<string, Task<IndexVolumeWriter>> createIndexFile = null)
|
||||
{
|
||||
Tuple<IndexVolumeWriter, FileEntryItem> indexfe = null;
|
||||
if (indexfile != null)
|
||||
indexfe = new Tuple<IndexVolumeWriter, FileEntryItem>(indexfile, new FileEntryItem(BackendActionType.Put, indexfile.RemoteFilename));
|
||||
var fe = new FileEntryItem(BackendActionType.Put, item.RemoteFilename);
|
||||
fe.SetLocalfilename(item.LocalFilename);
|
||||
|
||||
var fe = new FileEntryItem(BackendActionType.Put, item.RemoteFilename, indexfe);
|
||||
var tcs = new TaskCompletionSource<bool>();
|
||||
|
||||
return RunRetryOnMain<bool>(fe, async () =>
|
||||
await RunOnMain(async () =>
|
||||
{
|
||||
if (fe.IsRetry && fe.Indexfile != null && fe.TrackedInDb)
|
||||
await RenameFileAfterErrorAsync(fe);
|
||||
else
|
||||
fe.IsRetry = true;
|
||||
|
||||
await DoPut(fe);
|
||||
try
|
||||
{
|
||||
await DoWithRetry(fe, async () => {
|
||||
if (fe.IsRetry)
|
||||
await RenameFileAfterErrorAsync(fe);
|
||||
|
||||
m_uploadSuccess = true;
|
||||
return true;
|
||||
return await DoPut(fe);
|
||||
});
|
||||
|
||||
if (createIndexFile != null)
|
||||
{
|
||||
var ix = await createIndexFile(fe.RemoteFilename);
|
||||
var indexFile = new FileEntryItem(BackendActionType.Put, ix.RemoteFilename);
|
||||
indexFile.SetLocalfilename(ix.LocalFilename);
|
||||
|
||||
await m_database.UpdateRemoteVolumeAsync(indexFile.RemoteFilename, RemoteVolumeState.Uploading, -1, null);
|
||||
|
||||
await DoWithRetry(indexFile, async () => {
|
||||
if (indexFile.IsRetry)
|
||||
await RenameFileAfterErrorAsync(indexFile);
|
||||
|
||||
return await DoPut(indexFile);
|
||||
});
|
||||
}
|
||||
|
||||
tcs.TrySetResult(true);
|
||||
}
|
||||
catch(Exception ex)
|
||||
{
|
||||
if (ex is System.Threading.ThreadAbortException)
|
||||
tcs.TrySetCanceled();
|
||||
else
|
||||
tcs.TrySetException(ex);
|
||||
}
|
||||
});
|
||||
|
||||
await tcs.Task;
|
||||
}
|
||||
|
||||
public Task DeleteFileAsync(string remotename, long size)
|
||||
{
|
||||
var fe = new FileEntryItem(BackendActionType.Delete, remotename, null);
|
||||
var fe = new FileEntryItem(BackendActionType.Delete, remotename);
|
||||
return RunRetryOnMain(fe, () =>
|
||||
DoDelete(fe)
|
||||
);
|
||||
@@ -231,7 +252,7 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
|
||||
public Task CreateFolder(string remotename)
|
||||
{
|
||||
var fe = new FileEntryItem(BackendActionType.CreateFolder, remotename, null);
|
||||
var fe = new FileEntryItem(BackendActionType.CreateFolder, remotename);
|
||||
return RunRetryOnMain(fe, () =>
|
||||
DoCreateFolder(fe)
|
||||
);
|
||||
@@ -240,7 +261,7 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
|
||||
public Task<IList<Library.Interface.IFileEntry>> ListFilesAsync()
|
||||
{
|
||||
var fe = new FileEntryItem(BackendActionType.List, null, null);
|
||||
var fe = new FileEntryItem(BackendActionType.List, null);
|
||||
return RunRetryOnMain(fe, () =>
|
||||
DoList(fe)
|
||||
);
|
||||
@@ -256,7 +277,7 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
|
||||
public Task<Tuple<Library.Utility.TempFile, long, string>> GetFileWithInfoAsync(string remotename)
|
||||
{
|
||||
var fe = new FileEntryItem(BackendActionType.Get, remotename, null);
|
||||
var fe = new FileEntryItem(BackendActionType.Get, remotename);
|
||||
return RunRetryOnMain(fe, async () => {
|
||||
var res = await DoGet(fe);
|
||||
return new Tuple<Library.Utility.TempFile, long, string>(
|
||||
@@ -299,19 +320,22 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
if (m_backend == null)
|
||||
throw new Exception("Backend failed to re-load");
|
||||
|
||||
item.IsRetry = false;
|
||||
Exception lastException = null;
|
||||
|
||||
for(var i = 0; i < m_options.NumberOfRetries; i++)
|
||||
{
|
||||
if (m_options.RetryDelay.Ticks != 0 && i != 0)
|
||||
System.Threading.Thread.Sleep(m_options.RetryDelay);
|
||||
await Task.Delay(m_options.RetryDelay);
|
||||
|
||||
try
|
||||
{
|
||||
return await method();
|
||||
var r = await method();
|
||||
return r;
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
item.IsRetry = true;
|
||||
lastException = ex;
|
||||
await m_logchannel.WriteAsync(LogMessage.RetryAttempt(string.Format("Operation {0} with file {1} attempt {2} of {3} failed with message: {4}", item.Operation, item.RemoteFilename, i + 1, m_options.NumberOfRetries, ex.Message), ex));
|
||||
|
||||
@@ -363,41 +387,6 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
await m_logchannel.WriteAsync(LogMessage.Information(string.Format("Renaming \"{0}\" to \"{1}\"", oldname, newname)));
|
||||
await m_database.RenameRemoteFileAsync(oldname, newname);
|
||||
item.RemoteFilename = newname;
|
||||
|
||||
// If there is an index file attached to the block file,
|
||||
// it references the block filename, so we create a new index file
|
||||
// which is a copy of the current, but with the new name
|
||||
if (item.Indexfile != null)
|
||||
{
|
||||
if (!item.IndexfileUpdated)
|
||||
{
|
||||
item.Indexfile.Item1.FinishVolume(item.Hash, item.Size);
|
||||
item.Indexfile.Item1.Close();
|
||||
item.IndexfileUpdated = true;
|
||||
}
|
||||
|
||||
IndexVolumeWriter wr = null;
|
||||
try
|
||||
{
|
||||
wr = new IndexVolumeWriter(m_options);
|
||||
using(var rd = new IndexVolumeReader(p.CompressionModule, item.Indexfile.Item2.LocalFilename, m_options, m_options.BlockhashSize))
|
||||
wr.CopyFrom(rd, x => x == oldname ? newname : x);
|
||||
item.Indexfile.Item1.Dispose();
|
||||
item.Indexfile = new Tuple<IndexVolumeWriter, FileEntryItem>(wr, item.Indexfile.Item2);
|
||||
item.Indexfile.Item2.LocalTempfile.Dispose();
|
||||
item.Indexfile.Item2.LocalTempfile = wr.TempFile;
|
||||
wr.Close();
|
||||
}
|
||||
catch
|
||||
{
|
||||
if (wr != null)
|
||||
try { wr.Dispose(); }
|
||||
catch { }
|
||||
finally { wr = null; }
|
||||
|
||||
throw;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private async Task<bool> DoPut(FileEntryItem item)
|
||||
@@ -406,14 +395,8 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
await item.Encrypt(m_encryption, m_logchannel);
|
||||
|
||||
if (item.UpdateHashAndSize(m_options) && item.TrackedInDb)
|
||||
await m_database.UpdateRemoteVolume(item.RemoteFilename, RemoteVolumeState.Uploading, item.Size, item.Hash);
|
||||
await m_database.UpdateRemoteVolumeAsync(item.RemoteFilename, RemoteVolumeState.Uploading, item.Size, item.Hash);
|
||||
|
||||
if (item.Indexfile != null && !item.IndexfileUpdated)
|
||||
{
|
||||
item.Indexfile.Item1.FinishVolume(item.Hash, item.Size);
|
||||
item.Indexfile.Item1.Close();
|
||||
item.IndexfileUpdated = true;
|
||||
}
|
||||
|
||||
await m_database.LogRemoteOperationAsync("put", item.RemoteFilename, JsonConvert.SerializeObject(new { Size = item.Size, Hash = item.Hash }));
|
||||
await m_stats.SendEventAsync(BackendActionType.Put, BackendEventType.Started, item.RemoteFilename, item.Size);
|
||||
@@ -434,7 +417,7 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
Logging.Log.WriteMessage(string.Format("Uploaded {0} in {1}, {2}/s", Library.Utility.Utility.FormatSizeString(item.Size), duration, Library.Utility.Utility.FormatSizeString((long)(item.Size / duration.TotalSeconds))), Duplicati.Library.Logging.LogMessageType.Profiling);
|
||||
|
||||
if (item.TrackedInDb)
|
||||
await m_database.UpdateRemoteVolume(item.RemoteFilename, RemoteVolumeState.Uploaded, item.Size, item.Hash);
|
||||
await m_database.UpdateRemoteVolumeAsync(item.RemoteFilename, RemoteVolumeState.Uploaded, item.Size, item.Hash);
|
||||
|
||||
await m_stats.SendEventAsync(BackendActionType.Put, BackendEventType.Completed, item.RemoteFilename, item.Size);
|
||||
|
||||
@@ -446,8 +429,9 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
else if (f.Size != item.Size && f.Size >= 0)
|
||||
throw new Exception(string.Format("List verify failed for file: {0}, size was {1} but expected to be {2}", f.Name, f.Size, item.Size));
|
||||
}
|
||||
|
||||
|
||||
await item.DeleteLocalFile(m_logchannel);
|
||||
await m_database.CommitTransactionAsync("CommitAfterUpload");
|
||||
|
||||
return true;
|
||||
}
|
||||
@@ -521,7 +505,7 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
await m_database.LogRemoteOperationAsync("delete", item.RemoteFilename, result);
|
||||
}
|
||||
|
||||
await m_database.UpdateRemoteVolume(item.RemoteFilename, RemoteVolumeState.Deleted, -1, null);
|
||||
await m_database.UpdateRemoteVolumeAsync(item.RemoteFilename, RemoteVolumeState.Deleted, -1, null);
|
||||
await m_stats.SendEventAsync(BackendActionType.Delete, BackendEventType.Completed, item.RemoteFilename, item.Size);
|
||||
|
||||
return true;
|
||||
|
||||
@@ -1,38 +0,0 @@
|
||||
// Copyright (C) 2015, The Duplicati Team
|
||||
// http://www.duplicati.com, info@duplicati.com
|
||||
//
|
||||
// This library is free software; you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Lesser General Public License as
|
||||
// published by the Free Software Foundation; either version 2.1 of the
|
||||
// License, or (at your option) any later version.
|
||||
//
|
||||
// This library is distributed in the hope that it will be useful, but
|
||||
// WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
|
||||
// Lesser General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Lesser General Public
|
||||
// License along with this library; if not, write to the Free Software
|
||||
// Foundation, Inc., 59 Temple Place, Suite 330, Boston, MA 02111-1307 USA
|
||||
using System;
|
||||
using Duplicati.Library.Main.Volumes;
|
||||
|
||||
namespace Duplicati.Library.Main.Operation.Common
|
||||
{
|
||||
public interface IBackendOperation
|
||||
{
|
||||
}
|
||||
|
||||
public class UploadRequest : IBackendOperation
|
||||
{
|
||||
public BlockVolumeWriter BlockVolume { get; private set; }
|
||||
public IndexVolumeWriter IndexVolume { get; private set; }
|
||||
|
||||
public UploadRequest(BlockVolumeWriter blockvolume, IndexVolumeWriter indexvolume)
|
||||
{
|
||||
BlockVolume = blockvolume;
|
||||
IndexVolume = indexvolume;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,9 +15,11 @@
|
||||
// License along with this library; if not, write to the Free Software
|
||||
// Foundation, Inc., 59 Temple Place, Suite 330, Boston, MA 02111-1307 USA
|
||||
using System;
|
||||
using System.Linq;
|
||||
using System.Threading.Tasks;
|
||||
using CoCoL;
|
||||
using Duplicati.Library.Main.Database;
|
||||
using System.Collections.Generic;
|
||||
|
||||
namespace Duplicati.Library.Main.Operation.Common
|
||||
{
|
||||
@@ -38,7 +40,7 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
return RunOnMain(() => m_db.RegisterRemoteVolume(name, type, state, m_transaction));
|
||||
}
|
||||
|
||||
public Task UpdateRemoteVolume(string name, RemoteVolumeState state, long size, string hash)
|
||||
public Task UpdateRemoteVolumeAsync(string name, RemoteVolumeState state, long size, string hash)
|
||||
{
|
||||
return RunOnMain(() => m_db.UpdateRemoteVolume(name, state, size, hash, m_transaction));
|
||||
}
|
||||
@@ -66,6 +68,38 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
return RunOnMain(() => m_db.LogRemoteOperation(operation, path, data, m_transaction));
|
||||
}
|
||||
|
||||
public Task<IEnumerable<LocalDatabase.IBlock>> GetBlocksAsync(long volumeid)
|
||||
{
|
||||
// TODO: How does the IEnumerable work with RunOnMain ?
|
||||
return RunOnMain(() => m_db.GetBlocks(volumeid, m_transaction).ToArray().AsEnumerable());
|
||||
}
|
||||
|
||||
public Task<RemoteVolumeEntry> GetVolumeInfoAsync(string remotename)
|
||||
{
|
||||
return RunOnMain(() => m_db.GetRemoteVolume(remotename, m_transaction));
|
||||
}
|
||||
|
||||
public Task<IEnumerable<Tuple<string, byte[], int>>> GetBlocklistsAsync(long volumeid, int blocksize, int hashsize)
|
||||
{
|
||||
// TODO: How does the IEnumerable work with RunOnMain ?
|
||||
return RunOnMain(() => m_db.GetBlocklists(volumeid, blocksize, hashsize, m_transaction));
|
||||
}
|
||||
|
||||
public Task<long> GetRemoteVolumeIDAsync(string remotename)
|
||||
{
|
||||
return RunOnMain(() => m_db.GetRemoteVolumeID(remotename, m_transaction));
|
||||
}
|
||||
|
||||
protected override void Dispose(bool isDisposing)
|
||||
{
|
||||
base.Dispose(isDisposing);
|
||||
if (m_transaction != null)
|
||||
{
|
||||
m_transaction.Commit();
|
||||
m_transaction = null;
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,73 @@
|
||||
// Copyright (C) 2015, The Duplicati Team
|
||||
// http://www.duplicati.com, info@duplicati.com
|
||||
//
|
||||
// This library is free software; you can redistribute it and/or modify
|
||||
// it under the terms of the GNU Lesser General Public License as
|
||||
// published by the Free Software Foundation; either version 2.1 of the
|
||||
// License, or (at your option) any later version.
|
||||
//
|
||||
// This library is distributed in the hope that it will be useful, but
|
||||
// WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
|
||||
// Lesser General Public License for more details.
|
||||
//
|
||||
// You should have received a copy of the GNU Lesser General Public
|
||||
// License along with this library; if not, write to the Free Software
|
||||
// Foundation, Inc., 59 Temple Place, Suite 330, Boston, MA 02111-1307 USA
|
||||
using System;
|
||||
using System.Threading.Tasks;
|
||||
using Duplicati.Library.Main.Volumes;
|
||||
|
||||
namespace Duplicati.Library.Main.Operation.Common
|
||||
{
|
||||
internal static class IndexVolumeCreator
|
||||
{
|
||||
public static async Task<IndexVolumeWriter> CreateIndexVolume(string blockname, Options options, Common.DatabaseCommon database)
|
||||
{
|
||||
var w = new IndexVolumeWriter(options);
|
||||
w.VolumeID = await database.RegisterRemoteVolumeAsync(w.RemoteFilename, RemoteVolumeType.Index, RemoteVolumeState.Temporary);
|
||||
|
||||
var blockvolume = await database.GetVolumeInfoAsync(blockname);
|
||||
|
||||
w.StartVolume(blockname);
|
||||
foreach(var b in await database.GetBlocksAsync(blockvolume.ID))
|
||||
w.AddBlock(b.Hash, b.Size);
|
||||
|
||||
w.FinishVolume(blockvolume.Hash, blockvolume.Size);
|
||||
|
||||
if (options.IndexfilePolicy == Options.IndexFileStrategy.Full)
|
||||
foreach(var b in await database.GetBlocklistsAsync(blockvolume.ID, options.Blocksize, options.BlockhashSize))
|
||||
w.WriteBlocklist(b.Item1, b.Item2, 0, b.Item3);
|
||||
|
||||
w.Close();
|
||||
|
||||
return w;
|
||||
}
|
||||
|
||||
/*public static async Task<IndexVolumeWriter> ReCreateIndexVolume(string selfname, Options options, Repair.RepairDatabase database)
|
||||
{
|
||||
var w = new IndexVolumeWriter(options);
|
||||
w.SetRemoteFilename(selfname);
|
||||
|
||||
foreach(var blockvolume in await database.GetBlockVolumesFromIndexNameAsync(selfname))
|
||||
{
|
||||
w.StartVolume(blockvolume.Name);
|
||||
var volumeid = await database.GetRemoteVolumeIDAsync(blockvolume.Name);
|
||||
|
||||
foreach(var b in await database.GetBlocksAsync(volumeid))
|
||||
w.AddBlock(b.Hash, b.Size);
|
||||
|
||||
w.FinishVolume(blockvolume.Hash, blockvolume.Size);
|
||||
|
||||
if (options.IndexfilePolicy == Options.IndexFileStrategy.Full)
|
||||
foreach(var b in await database.GetBlocklistsAsync(volumeid, options.Blocksize, options.BlockhashSize))
|
||||
w.WriteBlocklist(b.Item1, b.Item2, 0, b.Item3);
|
||||
}
|
||||
|
||||
w.Close();
|
||||
|
||||
return w;
|
||||
}*/
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,23 +17,28 @@
|
||||
using System;
|
||||
using System.Threading.Tasks;
|
||||
using CoCoL;
|
||||
using System.Threading;
|
||||
|
||||
namespace Duplicati.Library.Main.Operation.Common
|
||||
{
|
||||
internal abstract class SingleRunner : ProcessHelper
|
||||
internal abstract class SingleRunner : IDisposable
|
||||
{
|
||||
protected IChannel<Func<Task>> m_channel;
|
||||
protected readonly Task m_worker;
|
||||
protected CancellationTokenSource m_workerSource;
|
||||
|
||||
public SingleRunner()
|
||||
{
|
||||
AutomationExtensions.AutoWireChannels(this, null);
|
||||
m_channel = ChannelManager.CreateChannel<Func<Task>>();
|
||||
m_worker = Start();
|
||||
m_workerSource = new System.Threading.CancellationTokenSource();
|
||||
m_worker = AutomationExtensions.RunProtected(this, Start);
|
||||
}
|
||||
|
||||
protected override async Task Start()
|
||||
private async Task Start()
|
||||
{
|
||||
while(true)
|
||||
var ct = m_workerSource.Token;
|
||||
while(!ct.IsCancellationRequested)
|
||||
await (await m_channel.ReadAsync())(); // Grab-n-execute
|
||||
}
|
||||
|
||||
@@ -49,7 +54,8 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
{
|
||||
try
|
||||
{
|
||||
res.SetResult(await method());
|
||||
var r = await method();
|
||||
res.SetResult(r);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
@@ -89,22 +95,19 @@ namespace Duplicati.Library.Main.Operation.Common
|
||||
});
|
||||
}
|
||||
|
||||
#region IDisposable implementation
|
||||
|
||||
public new void Dispose()
|
||||
public void Dispose()
|
||||
{
|
||||
base.Dispose();
|
||||
Dispose(true);
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
protected virtual void Dispose(bool isDisposing)
|
||||
{
|
||||
if (m_channel != null)
|
||||
try { m_channel.Retire(); }
|
||||
catch { }
|
||||
finally { m_channel = null; }
|
||||
finally { m_channel = null; }
|
||||
|
||||
AutomationExtensions.RetireAllChannels(this);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -259,9 +259,9 @@ namespace Duplicati.Library.Main.Operation
|
||||
}
|
||||
else
|
||||
{
|
||||
if (i.deleteGracePeriod > DateTime.UtcNow)
|
||||
if (i.DeleteGracePeriod > DateTime.UtcNow)
|
||||
{
|
||||
log.AddMessage(string.Format("keeping delete request for {0} until {1}", i.Name, i.deleteGracePeriod.ToLocalTime()));
|
||||
log.AddMessage(string.Format("keeping delete request for {0} until {1}", i.Name, i.DeleteGracePeriod.ToLocalTime()));
|
||||
}
|
||||
else
|
||||
{
|
||||
|
||||
@@ -59,15 +59,10 @@ namespace Duplicati.Library.Main.Operation
|
||||
return;
|
||||
|
||||
var file = fileversion.Value.File;
|
||||
long size;
|
||||
string hash;
|
||||
RemoteVolumeType type;
|
||||
RemoteVolumeState state;
|
||||
if (!db.GetRemoteVolume(file.Name, out hash, out size, out type, out state))
|
||||
size = file.Size;
|
||||
var entry = db.GetRemoteVolume(file.Name);
|
||||
|
||||
var files = new List<Library.Interface.IListResultFile>();
|
||||
using (var tmpfile = backend.Get(file.Name, size, hash))
|
||||
using (var tmpfile = backend.Get(file.Name, entry.Size < 0 ? file.Size : entry.Size, entry.Hash))
|
||||
using (var tmp = new Volumes.FilesetVolumeReader(RestoreHandler.GetCompressionModule(file.Name), tmpfile, m_options))
|
||||
foreach (var cf in tmp.ControlFiles)
|
||||
if (Library.Utility.FilterExpression.Matches(filter, cf.Key))
|
||||
|
||||
@@ -51,15 +51,11 @@ namespace Duplicati.Library.Main.Operation
|
||||
}
|
||||
|
||||
var file = fileversion.Value.File;
|
||||
long size;
|
||||
string hash;
|
||||
RemoteVolumeType type;
|
||||
RemoteVolumeState state;
|
||||
if (!db.GetRemoteVolume(file.Name, out hash, out size, out type, out state))
|
||||
size = file.Size;
|
||||
|
||||
|
||||
var entry = db.GetRemoteVolume(file.Name);
|
||||
|
||||
var res = new List<string>();
|
||||
using (var tmpfile = backend.Get(file.Name, size, hash))
|
||||
using (var tmpfile = backend.Get(file.Name, entry.Size < 0 ? file.Size : entry.Size, entry.Hash))
|
||||
using (var tmp = new Volumes.FilesetVolumeReader(RestoreHandler.GetCompressionModule(file.Name), tmpfile, m_options))
|
||||
foreach (var cf in tmp.ControlFiles)
|
||||
if (Library.Utility.FilterExpression.Matches(filter, cf.Key))
|
||||
|
||||
Reference in New Issue
Block a user