// Copyright (C) 2018, 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.Collections.Generic;
using System.Threading.Tasks;
namespace Duplicati.Library.Main.Operation.Common
{
///
/// This class provides a thread-safe wrapper around a non-thread safe
/// database instance, such that the database can be safely used within
/// the without needing to rewrite all the
/// code that does not handle concurrent access.
///
internal class BackendHandlerDatabaseGuard : IDisposable, IBackendHandlerDatabase
{
///
/// The database we are wrapping
///
private readonly Database.LocalDatabase m_db;
///
/// A flag indicating if this is a dry-run
///
private readonly bool m_dryrun;
///
/// The list of pending work
///
private List>> m_workQueue = new List>>();
///
/// The lock object protecting the queue.
///
private readonly object m_lock = new object();
///
/// The main thread, used to guard against execution from another thread
///
private readonly System.Threading.Thread m_mainthread;
///
/// Initializes a new instance of the
/// class.
///
/// The database to wrap.
/// A flag indicating if this is a dry-run
public BackendHandlerDatabaseGuard(Database.LocalDatabase db, bool dryrun)
{
m_db = db ?? throw new ArgumentNullException(nameof(db));
m_mainthread = System.Threading.Thread.CurrentThread;
}
///
/// Adds a pending work item to the queue
///
/// The task wait handle.
/// The method to run when ready.
private Task AddToQueue(Action action)
{
if (m_workQueue == null)
throw new ObjectDisposedException("This database guard instance has been shut down");
var tcs = new TaskCompletionSource();
lock(m_lock)
m_workQueue.Add(new Tuple>(action, tcs));
return tcs.Task;
}
///
/// Writes the current changes to the database
///
/// An awaitable task.
/// The message to use for logging the time spent in this operation.
/// If set to true, a transaction will be started again after this call.
public Task CommitTransactionAsync(string message, bool restart = true)
{
return AddToQueue(() => {
if (m_dryrun)
{
if (!restart)
m_db.RollbackTransaction();
return;
}
m_db.CommitTransaction(message, restart);
});
}
///
/// Writes remote operation log data to the database
///
/// An awaitable task.
/// The operation performed.
/// The remote path used.
/// Any data reported by the operation.
public Task LogRemoteOperationAsync(string operation, string path, string data)
{
return AddToQueue(() => { m_db.LogRemoteOperation(operation, path, data); });
}
///
/// Renames a remote file in the database
///
/// The remote file to rename.
/// The old filename.
/// The new filename.
public Task RenameRemoteFileAsync(string oldname, string newname)
{
return AddToQueue(() => { m_db.RenameRemoteFile(oldname, newname); });
}
///
/// Updates the remote volume information.
///
/// An awaitable task.
/// The name of the remote volume to update.
/// The new volume state.
/// The new volume size.
/// The new volume hash.
/// If set to true suppress cleanup operation.
/// The new delete grace time.
public Task UpdateRemoteVolumeAsync(string name, RemoteVolumeState state, long size, string hash, bool suppressCleanup = false, TimeSpan deleteGraceTime = default(TimeSpan))
{
return AddToQueue(() => { m_db.UpdateRemoteVolume(name, state, size, hash, suppressCleanup, deleteGraceTime); });
}
///
/// Processes all pending operations into the database.
///
/// The transaction instance
public void ProcessAllPendingOperations()
{
//if (m_mainthread != System.Threading.Thread.CurrentThread)
// throw new InvalidOperationException("Attempted to flush database work queue from another thread");
// If we are shut down, just exit
if (m_workQueue == null)
return;
var prevqueue = m_workQueue;
lock (m_lock)
{
// No need to fiddle if nothing has happened
if (prevqueue.Count == 0)
return;
m_workQueue = new List>>();
}
// We repeat here to allow quicker progress in case the backend is blocked
// on a log message, and emits a new one immediately after having one handled
while (prevqueue.Count != 0)
{
foreach (var op in prevqueue)
{
try
{
op.Item1();
op.Item2.TrySetResult(true);
}
catch (Exception ex)
{
op.Item2.TrySetException(ex);
}
}
// Inject a sleep here to allow the BackendHandler to progress
// Not pretty, but this entire class should be deleted once the code is
// rewritten to use thread-safe database access
System.Threading.Thread.Sleep(100);
prevqueue = m_workQueue;
lock (m_lock)
m_workQueue = new List>>();
}
}
///
/// Releases all resource used by the
/// object.
///
/// Call when you are finished using the
/// . The
/// method leaves the
/// in an unusable state.
/// After calling , you must release all references to the
/// so the garbage collector
/// can reclaim the memory that the
/// was occupying.
public void Dispose()
{
m_db.Dispose();
lock (m_lock)
{
if (m_workQueue != null)
foreach (var n in m_workQueue)
n.Item2.TrySetCanceled();
m_workQueue = null;
}
}
}
}