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:
Kenneth Skovhede
2016-02-22 21:27:12 +01:00
parent 9d3584a990
commit 3ff2f1eb98
17 changed files with 335 additions and 308 deletions
@@ -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))