Bunch of improvements i dunno.
This commit is contained in:
@@ -69,6 +69,8 @@ public interface IDefaultChangeProcessor
|
||||
Task AddMailItemQueueItemsAsync(IEnumerable<MailItemQueue> queueItems);
|
||||
Task<int> GetMailItemQueueCountAsync(Guid accountId);
|
||||
Task<List<MailItemQueue>> GetMailItemQueueAsync(Guid accountId, int take);
|
||||
Task<List<MailItemQueue>> GetMailItemQueueByFolderAsync(Guid accountId, string remoteFolderId, int take);
|
||||
Task<int> GetMailItemQueueCountByFolderAsync(Guid accountId, string remoteFolderId);
|
||||
Task UpdateMailItemQueueAsync(IEnumerable<MailItemQueue> queueItems);
|
||||
}
|
||||
|
||||
@@ -225,6 +227,12 @@ public class DefaultChangeProcessor(IDatabaseService databaseService,
|
||||
public Task<List<MailItemQueue>> GetMailItemQueueAsync(Guid accountId, int take)
|
||||
=> MailService.GetMailItemQueueAsync(accountId, take);
|
||||
|
||||
public Task<List<MailItemQueue>> GetMailItemQueueByFolderAsync(Guid accountId, string remoteFolderId, int take)
|
||||
=> MailService.GetMailItemQueueByFolderAsync(accountId, remoteFolderId, take);
|
||||
|
||||
public Task<int> GetMailItemQueueCountByFolderAsync(Guid accountId, string remoteFolderId)
|
||||
=> MailService.GetMailItemQueueCountByFolderAsync(accountId, remoteFolderId);
|
||||
|
||||
public Task UpdateMailItemQueueAsync(IEnumerable<MailItemQueue> queueItems)
|
||||
=> MailService.UpdateMailItemQueueAsync(queueItems);
|
||||
|
||||
|
||||
@@ -39,7 +39,10 @@ public class OutlookChangeProcessor(IDatabaseService databaseService,
|
||||
=> Connection.ExecuteAsync("UPDATE MailItemFolder SET DeltaToken = ? WHERE Id = ?", synchronizationIdentifier, folderId);
|
||||
|
||||
public Task UpdateFolderInitialSyncCompletedAsync(Guid folderId, bool isCompleted)
|
||||
=> Connection.ExecuteAsync("UPDATE MailItemFolder SET IsInitialSyncCompleted = ? WHERE Id = ?", isCompleted, folderId);
|
||||
{
|
||||
var status = isCompleted ? InitialSynchronizationStatus.Completed : InitialSynchronizationStatus.None;
|
||||
return Connection.ExecuteAsync("UPDATE MailItemFolder SET FolderStatus = ? WHERE Id = ?", status, folderId);
|
||||
}
|
||||
|
||||
public async Task ManageCalendarEventAsync(Event calendarEvent, AccountCalendar assignedCalendar, MailAccount organizerAccount)
|
||||
{
|
||||
|
||||
@@ -3,6 +3,7 @@ using System.Collections.Generic;
|
||||
using System.Net.Http;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using CommunityToolkit.Mvvm.ComponentModel;
|
||||
using CommunityToolkit.Mvvm.Messaging;
|
||||
using Wino.Core.Domain.Entities.Shared;
|
||||
using Wino.Core.Domain.Enums;
|
||||
@@ -13,12 +14,14 @@ using Wino.Messaging.UI;
|
||||
|
||||
namespace Wino.Core.Synchronizers;
|
||||
|
||||
public abstract class BaseSynchronizer<TBaseRequest> : IBaseSynchronizer
|
||||
public abstract partial class BaseSynchronizer<TBaseRequest> : ObservableObject, IBaseSynchronizer
|
||||
{
|
||||
protected SemaphoreSlim synchronizationSemaphore = new(1);
|
||||
protected CancellationToken activeSynchronizationCancellationToken;
|
||||
|
||||
protected List<IRequestBase> changeRequestQueue = [];
|
||||
protected readonly IMessenger Messenger;
|
||||
|
||||
public MailAccount Account { get; }
|
||||
|
||||
private AccountSynchronizerState state;
|
||||
@@ -29,13 +32,87 @@ public abstract class BaseSynchronizer<TBaseRequest> : IBaseSynchronizer
|
||||
{
|
||||
state = value;
|
||||
|
||||
WeakReferenceMessenger.Default.Send(new AccountSynchronizerStateChanged(Account.Id, value));
|
||||
// Send state changed message with current progress information
|
||||
Messenger.Send(new AccountSynchronizerStateChanged(
|
||||
Account.Id,
|
||||
value,
|
||||
TotalItemsToSync,
|
||||
RemainingItemsToSync,
|
||||
SynchronizationStatus));
|
||||
}
|
||||
}
|
||||
|
||||
protected BaseSynchronizer(MailAccount account)
|
||||
/// <summary>
|
||||
/// Current synchronization status message.
|
||||
/// </summary>
|
||||
[ObservableProperty]
|
||||
public partial string SynchronizationStatus { get; set; } = string.Empty;
|
||||
|
||||
/// <summary>
|
||||
/// Total items to download/sync in current operation.
|
||||
/// 0 means no active download or indeterminate progress.
|
||||
/// </summary>
|
||||
[ObservableProperty]
|
||||
public partial int TotalItemsToSync { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// Remaining items to download/sync in current operation.
|
||||
/// </summary>
|
||||
[ObservableProperty]
|
||||
public partial int RemainingItemsToSync { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// Calculated progress percentage (0-100) based on TotalItemsToSync and RemainingItemsToSync.
|
||||
/// Returns -1 for indeterminate progress (when both are 0).
|
||||
/// </summary>
|
||||
public double SynchronizationProgress
|
||||
{
|
||||
get
|
||||
{
|
||||
if (TotalItemsToSync == 0 || RemainingItemsToSync == 0)
|
||||
return -1; // Indeterminate
|
||||
|
||||
return ((double)(TotalItemsToSync - RemainingItemsToSync) / TotalItemsToSync) * 100;
|
||||
}
|
||||
}
|
||||
|
||||
protected BaseSynchronizer(MailAccount account, IMessenger messenger)
|
||||
{
|
||||
Account = account;
|
||||
Messenger = messenger ?? WeakReferenceMessenger.Default;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Resets synchronization progress to default state.
|
||||
/// </summary>
|
||||
protected void ResetSyncProgress()
|
||||
{
|
||||
TotalItemsToSync = 0;
|
||||
RemainingItemsToSync = 0;
|
||||
SynchronizationStatus = string.Empty;
|
||||
OnPropertyChanged(nameof(SynchronizationProgress));
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Updates synchronization progress with current item counts.
|
||||
/// </summary>
|
||||
/// <param name="total">Total items to sync</param>
|
||||
/// <param name="remaining">Remaining items to sync</param>
|
||||
/// <param name="status">Optional status message</param>
|
||||
protected void UpdateSyncProgress(int total, int remaining, string status = "")
|
||||
{
|
||||
TotalItemsToSync = total;
|
||||
RemainingItemsToSync = remaining;
|
||||
SynchronizationStatus = status;
|
||||
OnPropertyChanged(nameof(SynchronizationProgress));
|
||||
|
||||
// Send progress update message
|
||||
Messenger.Send(new AccountSynchronizerStateChanged(
|
||||
Account.Id,
|
||||
State,
|
||||
TotalItemsToSync,
|
||||
RemainingItemsToSync,
|
||||
SynchronizationStatus));
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
|
||||
@@ -88,7 +88,7 @@ public class GmailSynchronizer : WinoSynchronizer<IClientServiceRequest, Message
|
||||
public GmailSynchronizer(MailAccount account,
|
||||
IGmailAuthenticator authenticator,
|
||||
IGmailChangeProcessor gmailChangeProcessor,
|
||||
IGmailSynchronizerErrorHandlerFactory gmailSynchronizerErrorHandlerFactory) : base(account)
|
||||
IGmailSynchronizerErrorHandlerFactory gmailSynchronizerErrorHandlerFactory) : base(account, WeakReferenceMessenger.Default)
|
||||
{
|
||||
var messageHandler = new GmailClientMessageHandler(authenticator, account);
|
||||
|
||||
|
||||
@@ -55,7 +55,7 @@ public class ImapSynchronizer : WinoSynchronizer<ImapRequest, ImapMessageCreatio
|
||||
public ImapSynchronizer(MailAccount account,
|
||||
IImapChangeProcessor imapChangeProcessor,
|
||||
IImapSynchronizationStrategyProvider imapSynchronizationStrategyProvider,
|
||||
IApplicationConfiguration applicationConfiguration) : base(account)
|
||||
IApplicationConfiguration applicationConfiguration) : base(account, WeakReferenceMessenger.Default)
|
||||
{
|
||||
// Create client pool with account protocol log.
|
||||
_imapChangeProcessor = imapChangeProcessor;
|
||||
@@ -294,7 +294,8 @@ public class ImapSynchronizer : WinoSynchronizer<ImapRequest, ImapMessageCreatio
|
||||
_logger.Information("Internal synchronization started for {Name}", Account.Name);
|
||||
_logger.Information("Options: {Options}", options);
|
||||
|
||||
PublishSynchronizationProgress(1);
|
||||
// Set indeterminate progress initially
|
||||
UpdateSyncProgress(0, 0, "Synchronizing...");
|
||||
|
||||
bool shouldDoFolderSync = options.Type == MailSynchronizationType.FullFolders || options.Type == MailSynchronizationType.FoldersOnly;
|
||||
|
||||
@@ -307,12 +308,14 @@ public class ImapSynchronizer : WinoSynchronizer<ImapRequest, ImapMessageCreatio
|
||||
{
|
||||
var synchronizationFolders = await _imapChangeProcessor.GetSynchronizationFoldersAsync(options).ConfigureAwait(false);
|
||||
|
||||
for (int i = 0; i < synchronizationFolders.Count; i++)
|
||||
var totalFolders = synchronizationFolders.Count;
|
||||
|
||||
for (int i = 0; i < totalFolders; i++)
|
||||
{
|
||||
var folder = synchronizationFolders[i];
|
||||
var progress = (int)Math.Round((double)(i + 1) / synchronizationFolders.Count * 100);
|
||||
|
||||
PublishSynchronizationProgress(progress);
|
||||
|
||||
// Update progress based on folder completion
|
||||
UpdateSyncProgress(totalFolders, totalFolders - (i + 1), $"Syncing {folder.FolderName}...");
|
||||
|
||||
var folderDownloadedMessageIds = await SynchronizeFolderInternalAsync(folder, cancellationToken).ConfigureAwait(false);
|
||||
|
||||
@@ -325,7 +328,8 @@ public class ImapSynchronizer : WinoSynchronizer<ImapRequest, ImapMessageCreatio
|
||||
}
|
||||
}
|
||||
|
||||
PublishSynchronizationProgress(100);
|
||||
// Reset progress
|
||||
ResetSyncProgress();
|
||||
|
||||
// Get all unread new downloaded items and return in the result.
|
||||
// This is primarily used in notifications.
|
||||
|
||||
@@ -13,6 +13,7 @@ using System.Text.Json.Serialization;
|
||||
using System.Text.RegularExpressions;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using CommunityToolkit.Mvvm.Messaging;
|
||||
using Microsoft.Graph;
|
||||
using Microsoft.Graph.Me.MailFolders.Item.Messages.Delta;
|
||||
using Microsoft.Graph.Models;
|
||||
@@ -50,6 +51,25 @@ namespace Wino.Core.Synchronizers.Mail;
|
||||
[JsonSerializable(typeof(OutlookFileAttachment))]
|
||||
public partial class OutlookSynchronizerJsonContext : JsonSerializerContext;
|
||||
|
||||
/// <summary>
|
||||
/// Outlook synchronizer implementation with queue-based metadata-only synchronization.
|
||||
///
|
||||
/// SYNCHRONIZATION STRATEGY:
|
||||
/// - Uses per-folder queue system (unlike Gmail's per-account queue)
|
||||
/// - During sync (initial/delta), only message metadata is downloaded (no MIME content)
|
||||
/// - Messages are queued by folder using MailItemQueue with RemoteFolderId
|
||||
/// - MailCopy objects are created from Graph API metadata fields only
|
||||
/// - MIME files are downloaded on-demand when user explicitly reads a message
|
||||
/// - This dramatically reduces bandwidth usage and sync time
|
||||
///
|
||||
/// Key implementation details:
|
||||
/// - QueueMailIdsForFolderAsync: Queues all mail IDs for a folder using Delta API
|
||||
/// - ProcessMailQueueForFolderAsync: Downloads metadata in batches from queue
|
||||
/// - DownloadMessageMetadataBatchAsync: Concurrently downloads metadata for batches
|
||||
/// - CreateMailCopyFromMessage: Centralized method to create MailCopy from Message (metadata only)
|
||||
/// - DownloadMissingMimeMessageAsync: Downloads raw MIME only when explicitly requested
|
||||
/// - CreateNewMailPackagesAsync: Only used for search results and special cases (downloads MIME)
|
||||
/// </summary>
|
||||
public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message, Event>
|
||||
{
|
||||
public override uint BatchModificationSize => 20;
|
||||
@@ -93,7 +113,7 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
public OutlookSynchronizer(MailAccount account,
|
||||
IAuthenticator authenticator,
|
||||
IOutlookChangeProcessor outlookChangeProcessor,
|
||||
IOutlookSynchronizerErrorHandlerFactory errorHandlingFactory) : base(account)
|
||||
IOutlookSynchronizerErrorHandlerFactory errorHandlingFactory) : base(account, WeakReferenceMessenger.Default)
|
||||
{
|
||||
var tokenProvider = new MicrosoftTokenProvider(Account, authenticator);
|
||||
|
||||
@@ -157,7 +177,8 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
|
||||
try
|
||||
{
|
||||
PublishSynchronizationProgress(1);
|
||||
// Set indeterminate progress initially
|
||||
UpdateSyncProgress(0, 0, "Synchronizing folders...");
|
||||
|
||||
await SynchronizeFoldersAsync(cancellationToken).ConfigureAwait(false);
|
||||
|
||||
@@ -167,12 +188,14 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
|
||||
_logger.Information(string.Format("{1} Folders: {0}", string.Join(",", synchronizationFolders.Select(a => a.FolderName)), synchronizationFolders.Count));
|
||||
|
||||
for (int i = 0; i < synchronizationFolders.Count; i++)
|
||||
var totalFolders = synchronizationFolders.Count;
|
||||
|
||||
for (int i = 0; i < totalFolders; i++)
|
||||
{
|
||||
var folder = synchronizationFolders[i];
|
||||
var progress = (int)Math.Round((double)(i + 1) / synchronizationFolders.Count * 100);
|
||||
|
||||
PublishSynchronizationProgress(progress);
|
||||
|
||||
// Update progress based on folder completion
|
||||
UpdateSyncProgress(totalFolders, totalFolders - (i + 1), $"Syncing {folder.FolderName}...");
|
||||
|
||||
var folderDownloadedMessageIds = await SynchronizeFolderAsync(folder, cancellationToken).ConfigureAwait(false);
|
||||
downloadedMessageIds.AddRange(folderDownloadedMessageIds);
|
||||
@@ -188,7 +211,8 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
}
|
||||
finally
|
||||
{
|
||||
PublishSynchronizationProgress(100);
|
||||
// Reset progress at the end
|
||||
ResetSyncProgress();
|
||||
}
|
||||
|
||||
// Get all unred new downloaded items and return in the result.
|
||||
@@ -235,7 +259,7 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
_logger.Debug("Synchronizing {FolderName} with direct download approach", folder.FolderName);
|
||||
|
||||
// Check if initial sync is completed for this folder
|
||||
if (!folder.IsInitialSyncCompleted)
|
||||
if (folder.FolderStatus != InitialSynchronizationStatus.Completed)
|
||||
{
|
||||
_logger.Debug("Initial sync not completed for folder {FolderName}. Starting mail synchronization.", folder.FolderName);
|
||||
|
||||
@@ -244,7 +268,7 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
|
||||
// Mark initial sync as completed
|
||||
await _outlookChangeProcessor.UpdateFolderInitialSyncCompletedAsync(folder.Id, true).ConfigureAwait(false);
|
||||
folder.IsInitialSyncCompleted = true;
|
||||
folder.FolderStatus = InitialSynchronizationStatus.Completed;
|
||||
}
|
||||
else
|
||||
{
|
||||
@@ -265,88 +289,20 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Downloads mails for initial synchronization using Delta API and direct download with concurrency control.
|
||||
/// Downloads mails for initial synchronization using Delta API and queue-based system.
|
||||
/// First, queues all mail IDs, then downloads metadata in batches.
|
||||
/// </summary>
|
||||
private async Task DownloadMailsForInitialSyncAsync(MailItemFolder folder, List<string> downloadedMessageIds, CancellationToken cancellationToken)
|
||||
{
|
||||
_logger.Debug("Starting initial mail download for folder {FolderName}", folder.FolderName);
|
||||
|
||||
var mailIds = new List<string>();
|
||||
|
||||
try
|
||||
{
|
||||
// Always use Delta API for initial sync - this ensures proper delta token setup for future incremental syncs
|
||||
DeltaGetResponse messageCollectionPage = null;
|
||||
// Step 1: Queue all mail IDs using Delta API
|
||||
await QueueMailIdsForFolderAsync(folder, cancellationToken).ConfigureAwait(false);
|
||||
|
||||
if (string.IsNullOrEmpty(folder.DeltaToken))
|
||||
{
|
||||
messageCollectionPage = await _graphClient.Me.MailFolders[folder.RemoteFolderId].Messages.Delta.GetAsDeltaGetResponseAsync((config) =>
|
||||
{
|
||||
config.QueryParameters.Select = ["Id"]; // Only get the message Ids
|
||||
config.QueryParameters.Orderby = ["receivedDateTime desc"]; // Sort by received date desc
|
||||
config.QueryParameters.Top = (int)InitialMessageDownloadCountPerFolder;
|
||||
}, cancellationToken).ConfigureAwait(false);
|
||||
}
|
||||
else
|
||||
{
|
||||
var requestInformation = _graphClient.Me.MailFolders[folder.RemoteFolderId].Messages.Delta.ToGetRequestInformation((config) =>
|
||||
{
|
||||
config.QueryParameters.Select = ["Id"]; // Only get the message Ids
|
||||
config.QueryParameters.Orderby = ["receivedDateTime desc"]; // Sort by received date desc
|
||||
});
|
||||
|
||||
requestInformation.UrlTemplate = requestInformation.UrlTemplate.Insert(requestInformation.UrlTemplate.Length - 1, ",%24deltatoken");
|
||||
requestInformation.QueryParameters.Add("%24deltatoken", folder.DeltaToken);
|
||||
|
||||
messageCollectionPage = await _graphClient.RequestAdapter.SendAsync(requestInformation, DeltaGetResponse.CreateFromDiscriminatorValue, cancellationToken: cancellationToken);
|
||||
}
|
||||
|
||||
// Use PageIterator<DeltaGetResponse> for iterating through the messages
|
||||
var messageIterator = PageIterator<Message, DeltaGetResponse>.CreatePageIterator(_graphClient, messageCollectionPage, (message) =>
|
||||
{
|
||||
if (!IsResourceDeleted(message.AdditionalData))
|
||||
{
|
||||
mailIds.Add(message.Id);
|
||||
}
|
||||
|
||||
// Iterator must continue all the time to recieve delta token at the end.
|
||||
return true;
|
||||
});
|
||||
|
||||
await messageIterator.IterateAsync(cancellationToken).ConfigureAwait(false);
|
||||
|
||||
// Extract delta token from the iterator's delta link
|
||||
string deltaToken = null;
|
||||
if (!string.IsNullOrEmpty(messageIterator.Deltalink))
|
||||
{
|
||||
deltaToken = GetDeltaTokenFromDeltaLink(messageIterator.Deltalink);
|
||||
}
|
||||
|
||||
// Download mails concurrently with semaphore control
|
||||
if (mailIds.Any())
|
||||
{
|
||||
var mimeDownloadCount = Math.Min(mailIds.Count, InitialSyncMimeDownloadCount);
|
||||
_logger.Information("Starting concurrent download of {Count} mails for folder {FolderName} (first {MimeCount} with MIME messages)",
|
||||
mailIds.Count, folder.FolderName, mimeDownloadCount);
|
||||
await DownloadMailsConcurrentlyAsync(mailIds, folder, downloadedMessageIds, cancellationToken).ConfigureAwait(false);
|
||||
}
|
||||
else
|
||||
{
|
||||
_logger.Information("No mail ids found to download for folder {FolderName}", folder.FolderName);
|
||||
}
|
||||
|
||||
// Store the delta token for future incremental syncs - always store when available
|
||||
if (!string.IsNullOrEmpty(deltaToken))
|
||||
{
|
||||
await _outlookChangeProcessor.UpdateFolderDeltaSynchronizationIdentifierAsync(folder.Id, deltaToken).ConfigureAwait(false);
|
||||
await _outlookChangeProcessor.UpdateFolderLastSyncDateAsync(folder.Id).ConfigureAwait(false);
|
||||
folder.DeltaToken = deltaToken;
|
||||
_logger.Information("Stored delta token for folder {FolderName} - future syncs will be incremental", folder.FolderName);
|
||||
}
|
||||
else
|
||||
{
|
||||
_logger.Warning("No delta token received for folder {FolderName} - future syncs may re-download messages", folder.FolderName);
|
||||
}
|
||||
// Step 2: Process queued mail IDs in batches
|
||||
await ProcessMailQueueForFolderAsync(folder, downloadedMessageIds, cancellationToken).ConfigureAwait(false);
|
||||
}
|
||||
catch (ApiException apiException)
|
||||
{
|
||||
@@ -368,7 +324,7 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
if (apiException.ResponseStatusCode == 410)
|
||||
{
|
||||
folder.DeltaToken = string.Empty;
|
||||
folder.IsInitialSyncCompleted = false;
|
||||
folder.FolderStatus = InitialSynchronizationStatus.None;
|
||||
_logger.Information("API error handled successfully for folder {FolderName} during initial sync. Error: {ErrorCode}", folder.FolderName, apiException.ResponseStatusCode);
|
||||
}
|
||||
}
|
||||
@@ -389,162 +345,345 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Downloads mails concurrently with semaphore control to limit concurrent downloads to 10.
|
||||
/// This overload is used for initial sync where MIME messages are downloaded for the first 50 messages.
|
||||
/// Queues all mail IDs for a folder using Delta API.
|
||||
/// Only retrieves message IDs to minimize data transfer.
|
||||
/// </summary>
|
||||
private async Task DownloadMailsConcurrentlyAsync(List<string> mailIds, MailItemFolder folder, List<string> downloadedMessageIds, CancellationToken cancellationToken)
|
||||
private async Task QueueMailIdsForFolderAsync(MailItemFolder folder, CancellationToken cancellationToken)
|
||||
{
|
||||
await DownloadMailsConcurrentlyAsync(mailIds, folder, downloadedMessageIds, true, cancellationToken).ConfigureAwait(false);
|
||||
}
|
||||
_logger.Debug("Queuing mail IDs for folder {FolderName}", folder.FolderName);
|
||||
|
||||
/// <summary>
|
||||
/// Downloads mails concurrently with semaphore control to limit concurrent downloads to 10.
|
||||
/// </summary>
|
||||
private async Task DownloadMailsConcurrentlyAsync(List<string> mailIds, MailItemFolder folder, List<string> downloadedMessageIds, bool isInitialSync, CancellationToken cancellationToken)
|
||||
{
|
||||
var downloadTasks = mailIds.Select(async (mailId, index) =>
|
||||
var mailIds = new List<string>();
|
||||
|
||||
// Always use Delta API for initial sync - this ensures proper delta token setup for future incremental syncs
|
||||
DeltaGetResponse messageCollectionPage = null;
|
||||
|
||||
if (string.IsNullOrEmpty(folder.DeltaToken))
|
||||
{
|
||||
await _concurrentDownloadSemaphore.WaitAsync(cancellationToken).ConfigureAwait(false);
|
||||
try
|
||||
messageCollectionPage = await _graphClient.Me.MailFolders[folder.RemoteFolderId].Messages.Delta.GetAsDeltaGetResponseAsync((config) =>
|
||||
{
|
||||
// Download MIME for the first 50 messages during initial sync only
|
||||
bool shouldDownloadMime = isInitialSync && index < InitialSyncMimeDownloadCount;
|
||||
var downloaded = await DownloadSingleMailAsync(mailId, folder, shouldDownloadMime, cancellationToken).ConfigureAwait(false);
|
||||
if (downloaded != null)
|
||||
{
|
||||
lock (downloadedMessageIds)
|
||||
{
|
||||
downloadedMessageIds.Add(downloaded);
|
||||
}
|
||||
}
|
||||
}
|
||||
finally
|
||||
config.QueryParameters.Select = ["Id"]; // Only get the message Ids
|
||||
config.QueryParameters.Orderby = ["receivedDateTime desc"]; // Sort by received date desc
|
||||
config.QueryParameters.Top = (int)InitialMessageDownloadCountPerFolder;
|
||||
}, cancellationToken).ConfigureAwait(false);
|
||||
}
|
||||
else
|
||||
{
|
||||
var requestInformation = _graphClient.Me.MailFolders[folder.RemoteFolderId].Messages.Delta.ToGetRequestInformation((config) =>
|
||||
{
|
||||
_concurrentDownloadSemaphore.Release();
|
||||
config.QueryParameters.Select = ["Id"]; // Only get the message Ids
|
||||
config.QueryParameters.Orderby = ["receivedDateTime desc"]; // Sort by received date desc
|
||||
});
|
||||
|
||||
requestInformation.UrlTemplate = requestInformation.UrlTemplate.Insert(requestInformation.UrlTemplate.Length - 1, ",%24deltatoken");
|
||||
requestInformation.QueryParameters.Add("%24deltatoken", folder.DeltaToken);
|
||||
|
||||
messageCollectionPage = await _graphClient.RequestAdapter.SendAsync(requestInformation, DeltaGetResponse.CreateFromDiscriminatorValue, cancellationToken: cancellationToken);
|
||||
}
|
||||
|
||||
// Use PageIterator to iterate through all messages and collect IDs
|
||||
var messageIterator = PageIterator<Message, DeltaGetResponse>.CreatePageIterator(_graphClient, messageCollectionPage, (message) =>
|
||||
{
|
||||
if (!IsResourceDeleted(message.AdditionalData))
|
||||
{
|
||||
mailIds.Add(message.Id);
|
||||
}
|
||||
|
||||
// Iterator must continue all the time to receive delta token at the end.
|
||||
return true;
|
||||
});
|
||||
|
||||
await Task.WhenAll(downloadTasks).ConfigureAwait(false);
|
||||
await messageIterator.IterateAsync(cancellationToken).ConfigureAwait(false);
|
||||
|
||||
// Extract delta token from the iterator's delta link
|
||||
string deltaToken = null;
|
||||
if (!string.IsNullOrEmpty(messageIterator.Deltalink))
|
||||
{
|
||||
deltaToken = GetDeltaTokenFromDeltaLink(messageIterator.Deltalink);
|
||||
}
|
||||
|
||||
// Queue all mail IDs for processing
|
||||
if (mailIds.Any())
|
||||
{
|
||||
var queueEntries = mailIds.Select(id => new MailItemQueue
|
||||
{
|
||||
Id = Guid.CreateVersion7(),
|
||||
AccountId = Account.Id,
|
||||
RemoteServerId = id,
|
||||
RemoteFolderId = folder.RemoteFolderId,
|
||||
IsProcessed = false,
|
||||
CreatedAt = DateTime.UtcNow
|
||||
});
|
||||
|
||||
await _outlookChangeProcessor.AddMailItemQueueItemsAsync(queueEntries).ConfigureAwait(false);
|
||||
|
||||
_logger.Information("Queued {Count} mail IDs for folder {FolderName}", mailIds.Count, folder.FolderName);
|
||||
}
|
||||
else
|
||||
{
|
||||
_logger.Information("No mail ids found to queue for folder {FolderName}", folder.FolderName);
|
||||
}
|
||||
|
||||
// Store the delta token for future incremental syncs - always store when available
|
||||
if (!string.IsNullOrEmpty(deltaToken))
|
||||
{
|
||||
await _outlookChangeProcessor.UpdateFolderDeltaSynchronizationIdentifierAsync(folder.Id, deltaToken).ConfigureAwait(false);
|
||||
await _outlookChangeProcessor.UpdateFolderLastSyncDateAsync(folder.Id).ConfigureAwait(false);
|
||||
folder.DeltaToken = deltaToken;
|
||||
_logger.Information("Stored delta token for folder {FolderName} - future syncs will be incremental", folder.FolderName);
|
||||
}
|
||||
else
|
||||
{
|
||||
_logger.Warning("No delta token received for folder {FolderName} - future syncs may re-download messages", folder.FolderName);
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Downloads a single mail by ID and creates it in the database.
|
||||
/// Processes queued mail IDs in batches, downloading metadata only (no MIME).
|
||||
/// </summary>
|
||||
private async Task<string> DownloadSingleMailAsync(string mailId, MailItemFolder folder, bool downloadMime, CancellationToken cancellationToken)
|
||||
private async Task ProcessMailQueueForFolderAsync(MailItemFolder folder, List<string> downloadedMessageIds, CancellationToken cancellationToken)
|
||||
{
|
||||
try
|
||||
var totalInQueue = await _outlookChangeProcessor.GetMailItemQueueCountByFolderAsync(Account.Id, folder.RemoteFolderId).ConfigureAwait(false);
|
||||
|
||||
if (totalInQueue == 0)
|
||||
{
|
||||
// Check if mail already exists in database before downloading
|
||||
// to avoid unnecessary API calls and reprocessing existing mails
|
||||
bool mailExists = await _outlookChangeProcessor.IsMailExistsInFolderAsync(mailId, folder.Id).ConfigureAwait(false);
|
||||
if (mailExists)
|
||||
_logger.Information("No mails in queue for folder {FolderName}", folder.FolderName);
|
||||
return;
|
||||
}
|
||||
|
||||
_logger.Information("Processing {Count} queued mails for folder {FolderName}", totalInQueue, folder.FolderName);
|
||||
|
||||
var totalFailed = 0;
|
||||
var totalProcessed = 0;
|
||||
|
||||
// Set initial progress for queue processing
|
||||
UpdateSyncProgress(totalInQueue, totalInQueue, $"Downloading {folder.FolderName}...");
|
||||
|
||||
// Continue until all emails in queue are processed
|
||||
while (true)
|
||||
{
|
||||
// Get next batch of unprocessed emails from queue
|
||||
var mailItemQueue = await _outlookChangeProcessor.GetMailItemQueueByFolderAsync(Account.Id, folder.RemoteFolderId, 100).ConfigureAwait(false);
|
||||
|
||||
if (mailItemQueue.Count == 0)
|
||||
break; // No more emails to process
|
||||
|
||||
// Remove the items that should be deleted from queue first
|
||||
mailItemQueue.RemoveAll(a => a.ShouldDelete());
|
||||
|
||||
var mailChunks = mailItemQueue.Chunk(20); // Process 20 at a time
|
||||
|
||||
foreach (var chunk in mailChunks)
|
||||
{
|
||||
_logger.Debug("Mail {MailId} already exists in folder {FolderName}, skipping download", mailId, folder.FolderName);
|
||||
return null; // Not a new download
|
||||
cancellationToken.ThrowIfCancellationRequested();
|
||||
|
||||
// Collect message IDs from the chunk
|
||||
var messageIdsToDownload = chunk.Select(q => q.RemoteServerId).ToList();
|
||||
|
||||
try
|
||||
{
|
||||
// Download all messages in this chunk concurrently
|
||||
var chunkDownloadedIds = await DownloadMessageMetadataBatchAsync(messageIdsToDownload, folder, cancellationToken).ConfigureAwait(false);
|
||||
|
||||
downloadedMessageIds.AddRange(chunkDownloadedIds);
|
||||
|
||||
// Mark all items in chunk as processed
|
||||
foreach (var queueItem in chunk)
|
||||
{
|
||||
queueItem.IsProcessed = true;
|
||||
queueItem.ProcessedAt = DateTime.UtcNow;
|
||||
totalProcessed++;
|
||||
}
|
||||
|
||||
// Update progress with remaining items
|
||||
var remainingItems = totalInQueue - totalProcessed;
|
||||
UpdateSyncProgress(totalInQueue, remainingItems, $"Downloading {folder.FolderName} ({totalProcessed}/{totalInQueue})");
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.Error(ex, "Failed to download chunk of messages for folder {FolderName}", folder.FolderName);
|
||||
|
||||
// Mark all items in chunk as failed
|
||||
foreach (var queueItem in chunk)
|
||||
{
|
||||
queueItem.IsProcessed = false;
|
||||
queueItem.ProcessedAt = null;
|
||||
queueItem.FailedCount++;
|
||||
totalFailed++;
|
||||
}
|
||||
}
|
||||
|
||||
await _outlookChangeProcessor.UpdateMailItemQueueAsync(mailItemQueue).ConfigureAwait(false);
|
||||
|
||||
// If too many failures, pause to avoid hitting rate limits
|
||||
if (totalFailed > 50)
|
||||
{
|
||||
_logger.Warning("Too many failures ({Count}), pausing for 10 seconds", totalFailed);
|
||||
await Task.Delay(TimeSpan.FromSeconds(10), cancellationToken);
|
||||
totalFailed = 0; // Reset counter
|
||||
}
|
||||
}
|
||||
|
||||
// Download the message with minimal properties
|
||||
var message = await GetMessageByIdAsync(mailId, cancellationToken).ConfigureAwait(false);
|
||||
_logger.Debug("Processed batch: {Processed}/{Total} for folder {FolderName}", totalProcessed, totalInQueue, folder.FolderName);
|
||||
}
|
||||
|
||||
if (message != null)
|
||||
_logger.Information("Completed processing queue for folder {FolderName}. Processed: {Count}", folder.FolderName, totalProcessed);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Downloads metadata for a batch of messages using Graph SDK batch API (no MIME content).
|
||||
/// Processes up to 20 messages per batch request as per MaximumAllowedBatchRequestSize.
|
||||
/// </summary>
|
||||
private async Task<List<string>> DownloadMessageMetadataBatchAsync(List<string> messageIds, MailItemFolder folder, CancellationToken cancellationToken)
|
||||
{
|
||||
if (messageIds == null || messageIds.Count == 0)
|
||||
return new List<string>();
|
||||
|
||||
var downloadedIds = new List<string>();
|
||||
|
||||
// Filter out messages that already exist in the database
|
||||
var messagesToDownload = new List<string>();
|
||||
foreach (var messageId in messageIds)
|
||||
{
|
||||
bool mailExists = await _outlookChangeProcessor.IsMailExistsInFolderAsync(messageId, folder.Id).ConfigureAwait(false);
|
||||
if (!mailExists)
|
||||
{
|
||||
if (downloadMime)
|
||||
{
|
||||
// Download the full message packages with MIME for the first 50 messages
|
||||
var mailPackages = await CreateNewMailPackagesAsync(message, folder, cancellationToken).ConfigureAwait(false);
|
||||
messagesToDownload.Add(messageId);
|
||||
}
|
||||
else
|
||||
{
|
||||
_logger.Debug("Mail {MailId} already exists in folder {FolderName}, skipping download", messageId, folder.FolderName);
|
||||
}
|
||||
}
|
||||
|
||||
if (mailPackages != null)
|
||||
if (messagesToDownload.Count == 0)
|
||||
{
|
||||
_logger.Debug("All messages already exist in folder {FolderName}", folder.FolderName);
|
||||
return downloadedIds;
|
||||
}
|
||||
|
||||
// Process in batches of MaximumAllowedBatchRequestSize (20)
|
||||
var batches = messagesToDownload.Batch((int)MaximumAllowedBatchRequestSize);
|
||||
|
||||
foreach (var batch in batches)
|
||||
{
|
||||
cancellationToken.ThrowIfCancellationRequested();
|
||||
|
||||
try
|
||||
{
|
||||
var batchContent = new BatchRequestContentCollection(_graphClient);
|
||||
var requestIdToMessageIdMap = new Dictionary<string, string>();
|
||||
|
||||
// Add all message requests to the batch
|
||||
foreach (var messageId in batch)
|
||||
{
|
||||
var requestInfo = _graphClient.Me.Messages[messageId].ToGetRequestInformation((config) =>
|
||||
{
|
||||
foreach (var package in mailPackages)
|
||||
config.QueryParameters.Select = outlookMessageSelectParameters;
|
||||
});
|
||||
|
||||
var batchRequestId = await batchContent.AddBatchRequestStepAsync(requestInfo).ConfigureAwait(false);
|
||||
requestIdToMessageIdMap[batchRequestId] = messageId;
|
||||
}
|
||||
|
||||
// Execute the batch request
|
||||
var batchResponse = await _graphClient.Batch.PostAsync(batchContent, cancellationToken).ConfigureAwait(false);
|
||||
|
||||
// Process all responses
|
||||
foreach (var batchRequestId in requestIdToMessageIdMap.Keys)
|
||||
{
|
||||
var messageId = requestIdToMessageIdMap[batchRequestId];
|
||||
|
||||
try
|
||||
{
|
||||
// Deserialize the Message directly from batch response
|
||||
var message = await batchResponse.GetResponseByIdAsync<Message>(batchRequestId).ConfigureAwait(false);
|
||||
|
||||
if (message != null)
|
||||
{
|
||||
if (package?.Copy != null)
|
||||
// Create MailCopy from metadata only
|
||||
var mailCopy = CreateMailCopyFromMessage(message, folder);
|
||||
|
||||
if (mailCopy != null)
|
||||
{
|
||||
// Create package without MIME
|
||||
var package = new NewMailItemPackage(mailCopy, null, folder.RemoteFolderId);
|
||||
bool isInserted = await _outlookChangeProcessor.CreateMailAsync(Account.Id, package).ConfigureAwait(false);
|
||||
|
||||
if (isInserted)
|
||||
{
|
||||
_logger.Debug("Downloaded MIME message {MailId} for folder {FolderName}", mailId, folder.FolderName);
|
||||
return package.Copy.Id; // Successfully created with MIME
|
||||
downloadedIds.Add(mailCopy.Id);
|
||||
_logger.Debug("Downloaded metadata for message {MailId} in folder {FolderName}", messageId, folder.FolderName);
|
||||
}
|
||||
else
|
||||
{
|
||||
_logger.Warning("Failed to insert mail with MIME {MailId} for folder {FolderName}", mailId, folder.FolderName);
|
||||
_logger.Warning("Failed to insert mail {MailId} for folder {FolderName}", messageId, folder.FolderName);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
_logger.Debug("Could not create MIME mail packages for {MailId} in folder {FolderName}", mailId, folder.FolderName);
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
// Create minimal MailCopy without downloading MIME
|
||||
var mailCopy = await CreateMinimalMailCopyAsync(message, folder, cancellationToken).ConfigureAwait(false);
|
||||
|
||||
if (mailCopy != null)
|
||||
{
|
||||
// Create a minimal package without MIME for direct sync
|
||||
var package = new NewMailItemPackage(mailCopy, null, folder.RemoteFolderId);
|
||||
bool isInserted = await _outlookChangeProcessor.CreateMailAsync(Account.Id, package).ConfigureAwait(false);
|
||||
|
||||
if (isInserted)
|
||||
else
|
||||
{
|
||||
return mailCopy.Id; // Successfully created
|
||||
_logger.Warning("Failed to deserialize message {MailId} for folder {FolderName}", messageId, folder.FolderName);
|
||||
}
|
||||
}
|
||||
catch (ODataError odataError)
|
||||
{
|
||||
// Handle OData errors from the batch response
|
||||
if (odataError.ResponseStatusCode == 404)
|
||||
{
|
||||
_logger.Warning("Mail {MailId} not found on server (404) for folder {FolderName}", messageId, folder.FolderName);
|
||||
}
|
||||
else
|
||||
{
|
||||
_logger.Warning("Failed to insert mail {MailId} for folder {FolderName}", mailId, folder.FolderName);
|
||||
_logger.Error("OData error while downloading mail {MailId} for folder {FolderName}. Error: {Error}", messageId, folder.FolderName, odataError.Error?.Message);
|
||||
}
|
||||
}
|
||||
else
|
||||
catch (ServiceException serviceException)
|
||||
{
|
||||
_logger.Debug("Could not create MailCopy for {MailId} in folder {FolderName} (might be unsupported message type)", mailId, folder.FolderName);
|
||||
// Try to handle the error using the error handling factory
|
||||
var errorContext = new SynchronizerErrorContext
|
||||
{
|
||||
Account = Account,
|
||||
ErrorCode = (int?)serviceException.ResponseStatusCode,
|
||||
ErrorMessage = $"Service error during batch mail download: {serviceException.Message}",
|
||||
Exception = serviceException
|
||||
};
|
||||
|
||||
var handled = await _errorHandlingFactory.HandleErrorAsync(errorContext).ConfigureAwait(false);
|
||||
|
||||
if (!handled)
|
||||
{
|
||||
_logger.Error(serviceException, "Unhandled service error while downloading mail {MailId} for folder {FolderName}. Error: {ErrorCode}", messageId, folder.FolderName, serviceException.ResponseStatusCode);
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.Error(ex, "Error occurred while processing message {MailId} for folder {FolderName}", messageId, folder.FolderName);
|
||||
}
|
||||
}
|
||||
}
|
||||
else
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.Debug("Message {MailId} is null for folder {FolderName} (filtered out)", mailId, folder.FolderName);
|
||||
_logger.Error(ex, "Error occurred during batch download for folder {FolderName}", folder.FolderName);
|
||||
}
|
||||
}
|
||||
catch (ServiceException serviceException)
|
||||
{
|
||||
// Try to handle the error using the error handling factory first
|
||||
var errorContext = new SynchronizerErrorContext
|
||||
{
|
||||
Account = Account,
|
||||
ErrorCode = (int?)serviceException.ResponseStatusCode,
|
||||
ErrorMessage = $"Service error during mail download: {serviceException.Message}",
|
||||
Exception = serviceException
|
||||
};
|
||||
|
||||
var handled = await _errorHandlingFactory.HandleErrorAsync(errorContext).ConfigureAwait(false);
|
||||
|
||||
if (!handled)
|
||||
{
|
||||
// No handler could process this error, log appropriately
|
||||
if (serviceException.ResponseStatusCode == 404)
|
||||
{
|
||||
_logger.Warning("Mail {MailId} not found on server (404) for folder {FolderName}", mailId, folder.FolderName);
|
||||
}
|
||||
else
|
||||
{
|
||||
_logger.Error(serviceException, "Unhandled service error while downloading mail {MailId} for folder {FolderName}. Error: {ErrorCode}", mailId, folder.FolderName, serviceException.ResponseStatusCode);
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
_logger.Information("Service error handled successfully during mail download. Mail: {MailId}, Folder: {FolderName}, Error: {ErrorCode}", mailId, folder.FolderName, serviceException.ResponseStatusCode);
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.Error(ex, "Error occurred while downloading mail {MailId} for folder {FolderName}", mailId, folder.FolderName);
|
||||
}
|
||||
|
||||
return null;
|
||||
return downloadedIds;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Creates a MailCopy from an Outlook Message with metadata only (centralized method).
|
||||
/// This replaces the scattered CreateMinimalMailCopyAsync and AsMailCopy calls.
|
||||
/// </summary>
|
||||
private MailCopy CreateMailCopyFromMessage(Message message, MailItemFolder assignedFolder)
|
||||
{
|
||||
if (message == null) return null;
|
||||
|
||||
var mailCopy = message.AsMailCopy();
|
||||
mailCopy.FolderId = assignedFolder.Id;
|
||||
mailCopy.UniqueId = Guid.NewGuid();
|
||||
mailCopy.FileId = Guid.NewGuid();
|
||||
|
||||
return mailCopy;
|
||||
}
|
||||
|
||||
private string GetDeltaTokenFromDeltaLink(string deltaLink)
|
||||
@@ -552,23 +691,14 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
|
||||
protected override async Task QueueMailIdsForInitialSyncAsync(MailItemFolder folder, CancellationToken cancellationToken = default)
|
||||
{
|
||||
// This method is now replaced by direct downloading logic
|
||||
// Instead of queuing mail IDs, we now directly download them with concurrency control
|
||||
var downloadedMessageIds = new List<string>();
|
||||
await DownloadMailsForInitialSyncAsync(folder, downloadedMessageIds, cancellationToken).ConfigureAwait(false);
|
||||
// Queue all mail IDs for the folder
|
||||
await QueueMailIdsForFolderAsync(folder, cancellationToken).ConfigureAwait(false);
|
||||
}
|
||||
|
||||
protected override Task<MailCopy> CreateMinimalMailCopyAsync(Message message, MailItemFolder assignedFolder, CancellationToken cancellationToken = default)
|
||||
{
|
||||
if (message == null) return Task.FromResult<MailCopy>(null);
|
||||
|
||||
// Create MailCopy with minimal properties - no MIME download
|
||||
var mailCopy = message.AsMailCopy();
|
||||
mailCopy.FolderId = assignedFolder.Id;
|
||||
mailCopy.UniqueId = Guid.NewGuid();
|
||||
mailCopy.FileId = Guid.NewGuid();
|
||||
|
||||
return Task.FromResult(mailCopy);
|
||||
// Use centralized method
|
||||
return Task.FromResult(CreateMailCopyFromMessage(message, assignedFolder));
|
||||
}
|
||||
|
||||
private async Task<Message> GetMessageByIdAsync(string messageId, CancellationToken cancellationToken = default)
|
||||
@@ -622,7 +752,7 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
|
||||
private async Task ProcessDeltaChangesAndDownloadMailsAsync(MailItemFolder folder, List<string> downloadedMessageIds, CancellationToken cancellationToken = default)
|
||||
{
|
||||
// Process delta changes and directly download new mails
|
||||
// Process delta changes and download new mails with metadata only (no MIME)
|
||||
if (string.IsNullOrEmpty(folder.DeltaToken))
|
||||
{
|
||||
_logger.Debug("No delta token available for folder {FolderName}. Skipping delta sync.", folder.FolderName);
|
||||
@@ -638,7 +768,7 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
// Always use Delta endpoint with proper configuration
|
||||
var requestInformation = _graphClient.Me.MailFolders[folder.RemoteFolderId].Messages.Delta.ToGetRequestInformation((config) =>
|
||||
{
|
||||
config.QueryParameters.Select = ["Id"]; // Only get IDs for direct download
|
||||
config.QueryParameters.Select = ["Id"]; // Only get IDs
|
||||
config.QueryParameters.Orderby = ["receivedDateTime desc"]; // Sort by received date desc
|
||||
});
|
||||
|
||||
@@ -665,11 +795,12 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
|
||||
await messageIterator.IterateAsync(cancellationToken).ConfigureAwait(false);
|
||||
|
||||
// Download new mails directly with concurrency control
|
||||
// Download new mails with metadata only (no MIME)
|
||||
if (newMailIds.Any())
|
||||
{
|
||||
_logger.Information("Starting direct download of {Count} new mails from delta sync for folder {FolderName}", newMailIds.Count, folder.FolderName);
|
||||
await DownloadMailsConcurrentlyAsync(newMailIds, folder, downloadedMessageIds, false, cancellationToken).ConfigureAwait(false);
|
||||
_logger.Information("Downloading {Count} new mails from delta sync for folder {FolderName} (metadata only)", newMailIds.Count, folder.FolderName);
|
||||
var deltaDownloadedIds = await DownloadMessageMetadataBatchAsync(newMailIds, folder, cancellationToken).ConfigureAwait(false);
|
||||
downloadedMessageIds.AddRange(deltaDownloadedIds);
|
||||
}
|
||||
|
||||
// Update delta token for next sync - always store when there are no nextPageToken remaining
|
||||
@@ -701,7 +832,7 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
if (apiException.ResponseStatusCode == 410)
|
||||
{
|
||||
folder.DeltaToken = string.Empty;
|
||||
folder.IsInitialSyncCompleted = false;
|
||||
folder.FolderStatus = InitialSynchronizationStatus.None;
|
||||
_logger.Information("API error handled successfully for folder {FolderName} during delta sync. Error: {ErrorCode}", folder.FolderName, apiException.ResponseStatusCode);
|
||||
}
|
||||
}
|
||||
@@ -1638,8 +1769,10 @@ public class OutlookSynchronizer : WinoSynchronizer<RequestInformation, Message,
|
||||
|
||||
public override async Task<List<NewMailItemPackage>> CreateNewMailPackagesAsync(Message message, MailItemFolder assignedFolder, CancellationToken cancellationToken = default)
|
||||
{
|
||||
// Download MIME message for specific scenarios (e.g., search results, draft handling)
|
||||
// During normal sync, this method should not be called - use CreateMailCopyFromMessage instead
|
||||
var mimeMessage = await DownloadMimeMessageAsync(message.Id, cancellationToken).ConfigureAwait(false);
|
||||
var mailCopy = message.AsMailCopy();
|
||||
var mailCopy = CreateMailCopyFromMessage(message, assignedFolder);
|
||||
|
||||
if (message.IsDraft.GetValueOrDefault()
|
||||
&& mimeMessage.Headers.Contains(Domain.Constants.WinoLocalDraftHeader)
|
||||
|
||||
@@ -32,7 +32,7 @@ public abstract class WinoSynchronizer<TBaseRequest, TMessageType, TCalendarEven
|
||||
|
||||
protected ILogger Logger = Log.ForContext<WinoSynchronizer<TBaseRequest, TMessageType, TCalendarEventType>>();
|
||||
|
||||
protected WinoSynchronizer(MailAccount account) : base(account) { }
|
||||
protected WinoSynchronizer(MailAccount account, IMessenger messenger) : base(account, messenger) { }
|
||||
|
||||
/// <summary>
|
||||
/// How many items per single HTTP call can be modified.
|
||||
@@ -249,7 +249,8 @@ public abstract class WinoSynchronizer<TBaseRequest, TMessageType, TCalendarEven
|
||||
|
||||
await synchronizationSemaphore.WaitAsync(activeSynchronizationCancellationToken);
|
||||
|
||||
PublishSynchronizationProgress(1);
|
||||
// Set indeterminate progress for initial state
|
||||
UpdateSyncProgress(0, 0, "Synchronizing...");
|
||||
|
||||
State = AccountSynchronizerState.Synchronizing;
|
||||
|
||||
@@ -330,8 +331,8 @@ public abstract class WinoSynchronizer<TBaseRequest, TMessageType, TCalendarEven
|
||||
PendingSynchronizationRequest.Remove(pendingRequest.Key);
|
||||
}
|
||||
|
||||
// Reset account progress to hide the progress.
|
||||
PublishSynchronizationProgress(0);
|
||||
// Reset synchronization progress
|
||||
ResetSyncProgress();
|
||||
|
||||
State = AccountSynchronizerState.Idle;
|
||||
synchronizationSemaphore.Release();
|
||||
@@ -357,13 +358,6 @@ public abstract class WinoSynchronizer<TBaseRequest, TMessageType, TCalendarEven
|
||||
private void PublishUnreadItemChanges()
|
||||
=> WeakReferenceMessenger.Default.Send(new RefreshUnreadCountsMessage(Account.Id));
|
||||
|
||||
/// <summary>
|
||||
/// Sends a message to the shell to update the synchronization progress.
|
||||
/// </summary>
|
||||
/// <param name="progress">Percentage of the progress.</param>
|
||||
public void PublishSynchronizationProgress(double progress)
|
||||
=> WeakReferenceMessenger.Default.Send(new AccountSynchronizationProgressUpdatedMessage(Account.Id, progress));
|
||||
|
||||
/// <summary>
|
||||
/// Attempts to find out the best possible synchronization options after the batch request execution.
|
||||
/// </summary>
|
||||
|
||||
Reference in New Issue
Block a user