From 3ff2f1eb98bdcc84776315e854f8ef29b5d218e6 Mon Sep 17 00:00:00 2001 From: Kenneth Skovhede Date: Mon, 22 Feb 2016 21:27:12 +0100 Subject: [PATCH] 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. --- .../Library/Main/Database/LocalDatabase.cs | 51 +++---- .../Main/Database/RemoteVolumeEntry.cs | 37 +++-- .../Main/Duplicati.Library.Main.csproj | 3 +- .../Main/Operation/Backup/BackendUploader.cs | 79 +++++++++++ .../Main/Operation/Backup/DataBlock.cs | 3 +- .../Operation/Backup/DataBlockProcessor.cs | 37 +---- .../Operation/Backup/FileBlockProcessor.cs | 2 +- .../Operation/Backup/SpillCollectorProcess.cs | 40 ++---- .../Library/Main/Operation/BackupHandler.cs | 60 +------- .../Main/Operation/Common/BackendHandler.cs | 132 ++++++++---------- .../Main/Operation/Common/BackendOperation.cs | 38 ----- .../Main/Operation/Common/DatabaseCommon.cs | 36 ++++- .../Operation/Common/IndexVolumeCreator.cs | 73 ++++++++++ .../Main/Operation/Common/SingleRunner.cs | 27 ++-- .../Main/Operation/FilelistProcessor.cs | 4 +- .../Main/Operation/ListControlFilesHandler.cs | 9 +- .../Operation/RestoreControlFilesHandler.cs | 12 +- 17 files changed, 335 insertions(+), 308 deletions(-) create mode 100644 Duplicati/Library/Main/Operation/Backup/BackendUploader.cs delete mode 100644 Duplicati/Library/Main/Operation/Common/BackendOperation.cs create mode 100644 Duplicati/Library/Main/Operation/Common/IndexVolumeCreator.cs diff --git a/Duplicati/Library/Main/Database/LocalDatabase.cs b/Duplicati/Library/Main/Database/LocalDatabase.cs index bccff01fb..f221e2553 100644 --- a/Duplicati/Library/Main/Database/LocalDatabase.cs +++ b/Duplicati/Library/Main/Database/LocalDatabase.cs @@ -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 GetRemoteVolumes() + public IEnumerable 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> GetBlocklists(long volumeid, long blocksize, int hashsize) + public IEnumerable> 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, " + diff --git a/Duplicati/Library/Main/Database/RemoteVolumeEntry.cs b/Duplicati/Library/Main/Database/RemoteVolumeEntry.cs index a8c6b1621..d77b7e3b2 100644 --- a/Duplicati/Library/Main/Database/RemoteVolumeEntry.cs +++ b/Duplicati/Library/Main/Database/RemoteVolumeEntry.cs @@ -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; } } } diff --git a/Duplicati/Library/Main/Duplicati.Library.Main.csproj b/Duplicati/Library/Main/Duplicati.Library.Main.csproj index f08a5e887..1290a1781 100644 --- a/Duplicati/Library/Main/Duplicati.Library.Main.csproj +++ b/Duplicati/Library/Main/Duplicati.Library.Main.csproj @@ -121,7 +121,6 @@ - @@ -131,6 +130,8 @@ + + diff --git a/Duplicati/Library/Main/Operation/Backup/BackendUploader.cs b/Duplicati/Library/Main/Operation/Backup/BackendUploader.cs new file mode 100644 index 000000000..77125988a --- /dev/null +++ b/Duplicati/Library/Main/Operation/Backup/BackendUploader.cs @@ -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("BackendRequests"), + }, + + async self => + { + var inProgress = new Queue(); + 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(); + } + ); + + } + } +} + diff --git a/Duplicati/Library/Main/Operation/Backup/DataBlock.cs b/Duplicati/Library/Main/Operation/Backup/DataBlock.cs index 9e0a689f2..6570d142e 100644 --- a/Duplicati/Library/Main/Operation/Backup/DataBlock.cs +++ b/Duplicati/Library/Main/Operation/Backup/DataBlock.cs @@ -45,7 +45,8 @@ namespace Duplicati.Library.Main.Operation.Backup TaskCompletion = tcs }); - return await tcs.Task; + var r = await tcs.Task; + return r; } } } diff --git a/Duplicati/Library/Main/Operation/Backup/DataBlockProcessor.cs b/Duplicati/Library/Main/Operation/Backup/DataBlockProcessor.cs index a88b75a24..718e524d2 100644 --- a/Duplicati/Library/Main/Operation/Backup/DataBlockProcessor.cs +++ b/Duplicati/Library/Main/Operation/Backup/DataBlockProcessor.cs @@ -37,14 +37,13 @@ namespace Duplicati.Library.Main.Operation.Backup { LogChannel = ChannelMarker.ForWrite("LogChannel"), Input = ChannelMarker.ForRead("OutputBlocks"), - Output = ChannelMarker.ForWrite("BackendRequests"), - SpillPickup = ChannelMarker.ForWrite("SpillPickup"), + Output = ChannelMarker.ForWrite("BackendRequests"), + SpillPickup = ChannelMarker.ForWrite("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; diff --git a/Duplicati/Library/Main/Operation/Backup/FileBlockProcessor.cs b/Duplicati/Library/Main/Operation/Backup/FileBlockProcessor.cs index 5e4f033e8..d42db99ce 100644 --- a/Duplicati/Library/Main/Operation/Backup/FileBlockProcessor.cs +++ b/Duplicati/Library/Main/Operation/Backup/FileBlockProcessor.cs @@ -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]; } } diff --git a/Duplicati/Library/Main/Operation/Backup/SpillCollectorProcess.cs b/Duplicati/Library/Main/Operation/Backup/SpillCollectorProcess.cs index 965bfed65..d84f04d12 100644 --- a/Duplicati/Library/Main/Operation/Backup/SpillCollectorProcess.cs +++ b/Duplicati/Library/Main/Operation/Backup/SpillCollectorProcess.cs @@ -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("SpillPickup"), - Output = ChannelMarker.ForWrite("BackendRequests"), + Input = ChannelMarker.ForRead("SpillPickup"), + Output = ChannelMarker.ForWrite("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); + } } ); diff --git a/Duplicati/Library/Main/Operation/BackupHandler.cs b/Duplicati/Library/Main/Operation/BackupHandler.cs index cc2294c7f..9ff62a98a 100644 --- a/Duplicati/Library/Main/Operation/BackupHandler.cs +++ b/Duplicati/Library/Main/Operation/BackupHandler.cs @@ -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("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? diff --git a/Duplicati/Library/Main/Operation/Common/BackendHandler.cs b/Duplicati/Library/Main/Operation/Common/BackendHandler.cs index cb0410c99..5968796fd 100644 --- a/Duplicati/Library/Main/Operation/Common/BackendHandler.cs +++ b/Duplicati/Library/Main/Operation/Common/BackendHandler.cs @@ -67,10 +67,6 @@ namespace Duplicati.Library.Main.Operation.Common /// public long Size; /// - /// Reference to the index file entry that is updated if this entry changes - /// - public Tuple Indexfile; - /// /// A flag indicating if the final hash and size of the block volume has been written to the index file /// public bool IndexfileUpdated; @@ -88,16 +84,15 @@ namespace Duplicati.Library.Main.Operation.Common /// public bool IsRetry; - public FileEntryItem(BackendActionType operation, string remotefilename, Tuple 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 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> createIndexFile = null) { - Tuple indexfe = null; - if (indexfile != null) - indexfe = new Tuple(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(); - return RunRetryOnMain(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> 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> 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( @@ -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(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 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; diff --git a/Duplicati/Library/Main/Operation/Common/BackendOperation.cs b/Duplicati/Library/Main/Operation/Common/BackendOperation.cs deleted file mode 100644 index 01df7b42c..000000000 --- a/Duplicati/Library/Main/Operation/Common/BackendOperation.cs +++ /dev/null @@ -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; - } - } -} - diff --git a/Duplicati/Library/Main/Operation/Common/DatabaseCommon.cs b/Duplicati/Library/Main/Operation/Common/DatabaseCommon.cs index 2a31dcc03..89b266e05 100644 --- a/Duplicati/Library/Main/Operation/Common/DatabaseCommon.cs +++ b/Duplicati/Library/Main/Operation/Common/DatabaseCommon.cs @@ -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> GetBlocksAsync(long volumeid) + { + // TODO: How does the IEnumerable work with RunOnMain ? + return RunOnMain(() => m_db.GetBlocks(volumeid, m_transaction).ToArray().AsEnumerable()); + } + + public Task GetVolumeInfoAsync(string remotename) + { + return RunOnMain(() => m_db.GetRemoteVolume(remotename, m_transaction)); + } + + public Task>> 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 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; + } + } + } } diff --git a/Duplicati/Library/Main/Operation/Common/IndexVolumeCreator.cs b/Duplicati/Library/Main/Operation/Common/IndexVolumeCreator.cs new file mode 100644 index 000000000..20c418b2d --- /dev/null +++ b/Duplicati/Library/Main/Operation/Common/IndexVolumeCreator.cs @@ -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 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 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; + }*/ + } +} + diff --git a/Duplicati/Library/Main/Operation/Common/SingleRunner.cs b/Duplicati/Library/Main/Operation/Common/SingleRunner.cs index 7baea3813..2d5092c98 100644 --- a/Duplicati/Library/Main/Operation/Common/SingleRunner.cs +++ b/Duplicati/Library/Main/Operation/Common/SingleRunner.cs @@ -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> m_channel; protected readonly Task m_worker; + protected CancellationTokenSource m_workerSource; public SingleRunner() { + AutomationExtensions.AutoWireChannels(this, null); m_channel = ChannelManager.CreateChannel>(); - 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); } } } diff --git a/Duplicati/Library/Main/Operation/FilelistProcessor.cs b/Duplicati/Library/Main/Operation/FilelistProcessor.cs index 28d599879..eb47f26e3 100644 --- a/Duplicati/Library/Main/Operation/FilelistProcessor.cs +++ b/Duplicati/Library/Main/Operation/FilelistProcessor.cs @@ -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 { diff --git a/Duplicati/Library/Main/Operation/ListControlFilesHandler.cs b/Duplicati/Library/Main/Operation/ListControlFilesHandler.cs index f921cd21c..80457a8a8 100644 --- a/Duplicati/Library/Main/Operation/ListControlFilesHandler.cs +++ b/Duplicati/Library/Main/Operation/ListControlFilesHandler.cs @@ -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(); - 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)) diff --git a/Duplicati/Library/Main/Operation/RestoreControlFilesHandler.cs b/Duplicati/Library/Main/Operation/RestoreControlFilesHandler.cs index 1aa633820..5a44344c0 100644 --- a/Duplicati/Library/Main/Operation/RestoreControlFilesHandler.cs +++ b/Duplicati/Library/Main/Operation/RestoreControlFilesHandler.cs @@ -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(); - 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))