Updated CompactHandler to use the new async database methods

This commit is contained in:
Carl Johnsen
2025-05-21 15:31:23 +02:00
parent 99dddd4839
commit f4981ec0e4
@@ -52,23 +52,22 @@ namespace Duplicati.Library.Main.Operation
if (!System.IO.File.Exists(m_options.Dbpath))
throw new Exception(string.Format("Database file does not exist: {0}", m_options.Dbpath));
using (var db = new LocalDeleteDatabase(m_options.Dbpath, "Compact", m_options.SqlitePageCache))
using (var tr = new ReusableTransaction(db))
using (var db = await LocalDeleteDatabase.CreateAsync(m_options.Dbpath, "Compact", m_options.SqlitePageCache))
{
Utility.UpdateOptionsFromDb(db, m_options);
Utility.VerifyOptionsAndUpdateDatabase(db, m_options);
var changed = await DoCompactAsync(db, false, tr, backendManager).ConfigureAwait(false);
var changed = await DoCompactAsync(db, false, backendManager).ConfigureAwait(false);
if (changed && m_options.UploadVerificationFile)
await FilelistProcessor.UploadVerificationFile(backendManager, m_options, db, null);
await FilelistProcessor.UploadVerificationFile(backendManager, m_options, db);
if (!m_options.Dryrun)
{
tr.Commit("CommitCompact", restart: false);
await db.Transaction.CommitAsync("CommitCompact", restart: false);
if (changed)
{
db.WriteResults(m_result);
await db.WriteResults(m_result);
if (m_options.AutoVacuum)
{
m_result.VacuumResults = new VacuumResults(m_result);
@@ -79,31 +78,31 @@ namespace Duplicati.Library.Main.Operation
}
}
internal async Task<bool> DoCompactAsync(LocalDeleteDatabase db, bool hasVerifiedBackend, ReusableTransaction rtr, IBackendManager backendManager)
internal async Task<bool> DoCompactAsync(LocalDeleteDatabase db, bool hasVerifiedBackend, IBackendManager backendManager)
{
var report = db.GetCompactReport(m_options.VolumeSize, m_options.Threshold, m_options.SmallFileSize, m_options.SmallFileMaxCount, rtr.Transaction);
var report = await db.GetCompactReport(m_options.VolumeSize, m_options.Threshold, m_options.SmallFileSize, m_options.SmallFileMaxCount);
report.ReportCompactData();
if (report.ShouldReclaim || report.ShouldCompact)
{
// Workaround where we allow a running backendmanager to be used
if (!hasVerifiedBackend)
await FilelistProcessor.VerifyRemoteList(backendManager, m_options, db, m_result.BackendWriter, true, FilelistProcessor.VerifyMode.VerifyStrict, rtr.Transaction).ConfigureAwait(false);
await FilelistProcessor.VerifyRemoteList(backendManager, m_options, db, m_result.BackendWriter, true, FilelistProcessor.VerifyMode.VerifyStrict).ConfigureAwait(false);
var newvol = new BlockVolumeWriter(m_options);
newvol.VolumeID = db.RegisterRemoteVolume(newvol.RemoteFilename, RemoteVolumeType.Blocks, RemoteVolumeState.Temporary, rtr.Transaction);
newvol.VolumeID = await db.RegisterRemoteVolume(newvol.RemoteFilename, RemoteVolumeType.Blocks, RemoteVolumeState.Temporary);
IndexVolumeWriter newvolindex = null;
if (m_options.IndexfilePolicy != Options.IndexFileStrategy.None)
{
newvolindex = new IndexVolumeWriter(m_options);
newvolindex.VolumeID = db.RegisterRemoteVolume(newvolindex.RemoteFilename, RemoteVolumeType.Index, RemoteVolumeState.Temporary, rtr.Transaction);
db.AddIndexBlockLink(newvolindex.VolumeID, newvol.VolumeID, rtr.Transaction);
newvolindex.VolumeID = await db.RegisterRemoteVolume(newvolindex.RemoteFilename, RemoteVolumeType.Index, RemoteVolumeState.Temporary);
await db.AddIndexBlockLink(newvolindex.VolumeID, newvol.VolumeID);
}
var blocksInVolume = 0L;
var buffer = new byte[m_options.Blocksize];
var remoteList = db.GetRemoteVolumes().Where(n => n.State == RemoteVolumeState.Uploaded || n.State == RemoteVolumeState.Verified).ToArray();
var remoteList = await db.GetRemoteVolumes().Where(n => n.State == RemoteVolumeState.Uploaded || n.State == RemoteVolumeState.Verified).ToArrayAsync();
//These are for bookkeeping
var uploadedVolumes = new List<KeyValuePair<string, long>>();
@@ -121,7 +120,7 @@ namespace Duplicati.Library.Main.Operation
.Cast<IRemoteVolume>()
.ToList();
}
await foreach (var d in DoDelete(db, backendManager, fullyDeleteable, rtr, m_result.TaskControl.ProgressToken).ConfigureAwait(false))
await foreach (var d in DoDelete(db, backendManager, fullyDeleteable, m_result.TaskControl.ProgressToken).ConfigureAwait(false))
deletedVolumes.Add(d);
// This list is used to pick up unused volumes,
@@ -133,7 +132,7 @@ namespace Duplicati.Library.Main.Operation
{
// If we crash now, we may leave partial files
if (!m_options.Dryrun)
db.TerminatedWithActiveUploads = true;
await db.TerminatedWithActiveUploads(true);
newvolindex?.StartVolume(newvol.RemoteFilename);
List<IRemoteVolume> volumesToDownload = [];
@@ -147,7 +146,7 @@ namespace Duplicati.Library.Main.Operation
.ToList();
}
using (var q = db.CreateBlockQueryHelper(rtr.Transaction))
using (var q = await db.CreateBlockQueryHelper())
{
await foreach (var (tmpfile, hash, size, name) in backendManager.GetFilesOverlappedAsync(volumesToDownload, m_result.TaskControl.ProgressToken).ConfigureAwait(false))
{
@@ -156,18 +155,18 @@ namespace Duplicati.Library.Main.Operation
var entry = new RemoteVolume(name, hash, size);
if (!await m_result.TaskControl.ProgressRendevouz().ConfigureAwait(false))
{
await backendManager.WaitForEmptyAsync(db, rtr.Transaction, m_result.TaskControl.ProgressToken).ConfigureAwait(false);
await backendManager.WaitForEmptyAsync(db, m_result.TaskControl.ProgressToken).ConfigureAwait(false);
return false;
}
downloadedVolumes.Add(new KeyValuePair<string, long>(entry.Name, entry.Size));
var volumeid = db.GetRemoteVolumeID(entry.Name, rtr.Transaction);
var volumeid = await db.GetRemoteVolumeID(entry.Name);
var inst = VolumeBase.ParseFilename(entry.Name);
using (var f = new BlockVolumeReader(inst.CompressionModule, tmpfile, m_options))
{
foreach (var e in f.Blocks)
{
if (q.UseBlock(e.Key, e.Value, volumeid, rtr.Transaction))
if (await q.UseBlock(e.Key, e.Value, volumeid))
{
//TODO: How do we get the compression hint? Reverse query for filename in db?
var s = f.ReadBlock(e.Key, buffer);
@@ -178,32 +177,32 @@ namespace Duplicati.Library.Main.Operation
if (newvolindex != null)
newvolindex.AddBlock(e.Key, e.Value);
db.RegisterDuplicatedBlock(e.Key, e.Value, newvol.VolumeID, rtr.Transaction);
await db.RegisterDuplicatedBlock(e.Key, e.Value, newvol.VolumeID);
blocksInVolume++;
if (newvol.Filesize > (m_options.VolumeSize - m_options.Blocksize))
{
await FinishVolumeAndUpload(db, backendManager, newvol, newvolindex, uploadedVolumes, rtr);
await FinishVolumeAndUpload(db, backendManager, newvol, newvolindex, uploadedVolumes);
newvol = new BlockVolumeWriter(m_options);
newvol.VolumeID = db.RegisterRemoteVolume(newvol.RemoteFilename, RemoteVolumeType.Blocks, RemoteVolumeState.Temporary, rtr.Transaction);
newvol.VolumeID = await db.RegisterRemoteVolume(newvol.RemoteFilename, RemoteVolumeType.Blocks, RemoteVolumeState.Temporary);
if (m_options.IndexfilePolicy != Options.IndexFileStrategy.None)
{
newvolindex = new IndexVolumeWriter(m_options);
newvolindex.VolumeID = db.RegisterRemoteVolume(newvolindex.RemoteFilename, RemoteVolumeType.Index, RemoteVolumeState.Temporary, rtr.Transaction);
db.AddIndexBlockLink(newvolindex.VolumeID, newvol.VolumeID, rtr.Transaction);
newvolindex.VolumeID = await db.RegisterRemoteVolume(newvolindex.RemoteFilename, RemoteVolumeType.Index, RemoteVolumeState.Temporary);
await db.AddIndexBlockLink(newvolindex.VolumeID, newvol.VolumeID);
newvolindex.StartVolume(newvol.RemoteFilename);
}
blocksInVolume = 0;
// Wait for the backend to catch up
await backendManager.WaitForEmptyAsync(db, rtr.Transaction, m_result.TaskControl.ProgressToken).ConfigureAwait(false);
await backendManager.WaitForEmptyAsync(db, m_result.TaskControl.ProgressToken).ConfigureAwait(false);
// Commit as we have uploaded a volume
if (!m_options.Dryrun)
rtr.Commit("CommitCompact");
await db.Transaction.CommitAsync("CommitCompact");
}
}
}
@@ -215,14 +214,14 @@ namespace Duplicati.Library.Main.Operation
if (blocksInVolume > 0)
{
await FinishVolumeAndUpload(db, backendManager, newvol, newvolindex, uploadedVolumes, rtr).ConfigureAwait(false);
await FinishVolumeAndUpload(db, backendManager, newvol, newvolindex, uploadedVolumes).ConfigureAwait(false);
}
else
{
db.RemoveRemoteVolume(newvol.RemoteFilename, rtr.Transaction);
await db.RemoveRemoteVolume(newvol.RemoteFilename);
if (newvolindex != null)
{
db.RemoveRemoteVolume(newvolindex.RemoteFilename, rtr.Transaction);
await db.RemoveRemoteVolume(newvolindex.RemoteFilename);
newvolindex.FinishVolume(null, 0);
}
}
@@ -230,7 +229,7 @@ namespace Duplicati.Library.Main.Operation
// The remainder of the operation cannot leave partial files
if (!m_options.Dryrun)
db.TerminatedWithActiveUploads = false;
await db.TerminatedWithActiveUploads(false);
}
else
{
@@ -238,7 +237,7 @@ namespace Duplicati.Library.Main.Operation
newvol.Dispose();
}
await foreach (var d in DoDelete(db, backendManager, deleteableVolumes, rtr, m_result.TaskControl.ProgressToken).ConfigureAwait(false))
await foreach (var d in DoDelete(db, backendManager, deleteableVolumes, m_result.TaskControl.ProgressToken).ConfigureAwait(false))
deletedVolumes.Add(d);
var downloadSize = downloadedVolumes.Where(x => x.Value >= 0).Aggregate(0L, (a, x) => a + x.Value);
@@ -283,7 +282,7 @@ namespace Duplicati.Library.Main.Operation
Library.Utility.Utility.FormatSizeString(m_result.DeletedFileSize - m_result.UploadedFileSize));
}
await backendManager.WaitForEmptyAsync(db, rtr.Transaction, m_result.TaskControl.ProgressToken).ConfigureAwait(false);
await backendManager.WaitForEmptyAsync(db, m_result.TaskControl.ProgressToken).ConfigureAwait(false);
m_result.EndTime = DateTime.UtcNow;
return (m_result.DeletedFileCount + m_result.UploadedFileCount) > 0;
@@ -295,41 +294,42 @@ namespace Duplicati.Library.Main.Operation
}
}
private async IAsyncEnumerable<KeyValuePair<string, long>> DoDelete(LocalDeleteDatabase db, IBackendManager backend, IEnumerable<IRemoteVolume> deleteableVolumes, ReusableTransaction rtr, [EnumeratorCancellation] CancellationToken cancellationToken)
private async IAsyncEnumerable<KeyValuePair<string, long>> DoDelete(LocalDeleteDatabase db, IBackendManager backend, IEnumerable<IRemoteVolume> deleteableVolumes, [EnumeratorCancellation] CancellationToken cancellationToken)
{
// Find volumes that can be deleted
var remoteFilesToRemove = db.ReOrderDeleteableVolumes(deleteableVolumes, rtr.Transaction).ToList();
var remoteFilesToRemove = await db.ReOrderDeleteableVolumes(deleteableVolumes).ToListAsync();
// Make sure we do not re-assign blocks to any of the volumes we are about to delete
var toRemoveVolumeIds = db.GetRemoteVolumeIDs(remoteFilesToRemove.Select(x => x.Name), rtr.Transaction)
var toRemoveVolumeIds = await db
.GetRemoteVolumeIDs(remoteFilesToRemove.Select(x => x.Name))
.Select(x => x.Value)
.Distinct()
.ToList();
.ToListAsync();
// Mark all volumes and relevant index files as disposable
foreach (var f in remoteFilesToRemove)
{
db.PrepareForDelete(f.Name, toRemoveVolumeIds, rtr.Transaction);
db.UpdateRemoteVolume(f.Name, RemoteVolumeState.Deleting, f.Size, f.Hash, rtr.Transaction);
await db.PrepareForDelete(f.Name, toRemoveVolumeIds);
await db.UpdateRemoteVolume(f.Name, RemoteVolumeState.Deleting, f.Size, f.Hash);
}
// Before we commit the current state, make sure the backend has caught up
await backend.WaitForEmptyAsync(db, rtr.Transaction, cancellationToken).ConfigureAwait(false);
await backend.WaitForEmptyAsync(db, cancellationToken).ConfigureAwait(false);
if (!m_options.Dryrun)
rtr.Commit("CommitDelete");
await db.Transaction.CommitAsync("CommitDelete");
await foreach (var d in PerformDelete(backend, remoteFilesToRemove, cancellationToken).ConfigureAwait(false))
yield return d;
}
private async Task FinishVolumeAndUpload(LocalDeleteDatabase db, IBackendManager backendManager, BlockVolumeWriter newvol, IndexVolumeWriter newvolindex, List<KeyValuePair<string, long>> uploadedVolumes, ReusableTransaction rtr)
private async Task FinishVolumeAndUpload(LocalDeleteDatabase db, IBackendManager backendManager, BlockVolumeWriter newvol, IndexVolumeWriter newvolindex, List<KeyValuePair<string, long>> uploadedVolumes)
{
Action indexVolumeFinished = null;
if (newvolindex != null && m_options.IndexfilePolicy == Options.IndexFileStrategy.Full)
indexVolumeFinished = () =>
indexVolumeFinished = async () =>
{
foreach (var blocklist in db.GetBlocklists(newvol.VolumeID, m_options.Blocksize, m_options.BlockhashSize))
await foreach (var blocklist in db.GetBlocklists(newvol.VolumeID, m_options.Blocksize, m_options.BlockhashSize))
newvolindex.WriteBlocklist(blocklist.Item1, blocklist.Item2, 0, blocklist.Item3);
};
@@ -343,7 +343,7 @@ namespace Duplicati.Library.Main.Operation
// Once fixed, we can perhaps let he backend manager simply call the database directly
if (!m_options.Dryrun)
{
rtr.Commit("CommitUpload");
await db.Transaction.CommitAsync("CommitUpload");
await backendManager.PutAsync(newvol, newvolindex, indexVolumeFinished, false, null, m_result.TaskControl.ProgressToken).ConfigureAwait(false);
}
else