From 8d68bfc22ac675c701e922722ea608c6610b117f Mon Sep 17 00:00:00 2001 From: Stevan Freeborn <65925598+StevanFreeborn@users.noreply.github.com> Date: Tue, 3 Mar 2026 05:09:03 -0600 Subject: [PATCH] fix(infra): don't advance account cursor unless all transactions proceed successfully --- .../Plaid/PlaidAddedTransactionHandler.cs | 98 ++++++++++++++----- .../Transactions/Plaid/PlaidExtensions.cs | 1 + .../Plaid/PlaidModifiedTransactionHandler.cs | 88 +++++++++++------ .../Plaid/PlaidRemovedTransactionHandler.cs | 38 +++++-- .../Plaid/PlaidTransactionProcessor.cs | 31 +++++- .../Plaid/PlaidTransactionSyncer.cs | 15 ++- .../Plaid/SyncUpdatesProcessor.cs | 13 ++- 7 files changed, 212 insertions(+), 72 deletions(-) diff --git a/src/FiscalOS.Infra/Transactions/Plaid/PlaidAddedTransactionHandler.cs b/src/FiscalOS.Infra/Transactions/Plaid/PlaidAddedTransactionHandler.cs index ee63507..8e99cbc 100644 --- a/src/FiscalOS.Infra/Transactions/Plaid/PlaidAddedTransactionHandler.cs +++ b/src/FiscalOS.Infra/Transactions/Plaid/PlaidAddedTransactionHandler.cs @@ -5,50 +5,102 @@ namespace FiscalOS.Infra.Transactions.Plaid; internal interface IPlaidAddedTransactionHandler { - void Handle(Account account, IEnumerable added); + Task HandleAsync( + Account account, + IEnumerable added, + CancellationToken ct + ); } internal sealed class PlaidAddedTransactionHandler : IPlaidAddedTransactionHandler { private readonly ILogger _logger; + private readonly AppDbContext _appDbContext; - private PlaidAddedTransactionHandler(ILogger logger) + private PlaidAddedTransactionHandler( + ILogger logger, + AppDbContext appDbContext + ) { _logger = logger; + _appDbContext = appDbContext; } public static PlaidAddedTransactionHandler From(IServiceProvider serviceProvider) { return new( - serviceProvider.GetRequiredService>() + serviceProvider.GetRequiredService>(), + serviceProvider.GetRequiredService() ); } - public void Handle(Account account, IEnumerable added) + public async Task HandleAsync(Account account, IEnumerable added, CancellationToken ct) { + var existingTransactionIds = added.Select(t => t.TransactionId); + var existingTransactions = await _appDbContext.Transactions + .Where(t => t.Metadata is PlaidTransactionMetadata && existingTransactionIds.Contains(((PlaidTransactionMetadata)t.Metadata).PlaidId)) + .ToListAsync(ct) + .ConfigureAwait(false); + + var addedCount = 0; + foreach (var addedTransaction in added) { - if (addedTransaction.Pending.GetValueOrDefault()) + try { - _logger.LogInformation( - "Skipping pending transaction {TransactionId} for account {AccountId} as it has not been posted yet", - addedTransaction.TransactionId, - account.Id - ); - continue; - } + if (addedTransaction.Pending.GetValueOrDefault()) + { + _logger.LogInformation( + "Skipping pending transaction {TransactionId} for account {AccountId} as it has not been posted yet", + addedTransaction.TransactionId, + account.Id + ); + addedCount++; + continue; + } - var transactionMetadata = PlaidTransactionMetadata.From(addedTransaction.TransactionId); - var transaction = Transaction.From( - account.UserId, - account.Id, - addedTransaction.MerchantName, - addedTransaction.Description, - addedTransaction.Amount, - addedTransaction.PostedDate, - transactionMetadata - ); - account.AddTransaction(transaction); + var existingTransaction = existingTransactions.FirstOrDefault( + t => t.Metadata is PlaidTransactionMetadata metadata && metadata.PlaidId == addedTransaction.TransactionId + ); + + if (existingTransaction is not null) + { + _logger.LogInformation( + "Skipping added transaction {PlaidTransactionId} for account {AccountId} as it has already been added as transaction {TransactionId}", + addedTransaction.TransactionId, + account.Id, + existingTransaction.Id + ); + addedCount++; + continue; + } + + var transactionMetadata = PlaidTransactionMetadata.From(addedTransaction.TransactionId); + var transaction = Transaction.From( + account.UserId, + account.Id, + addedTransaction.Merchant, + addedTransaction.Description, + addedTransaction.Amount, + addedTransaction.PostedDate, + transactionMetadata + ); + account.AddTransaction(transaction); + addedCount++; + } + catch (Exception ex) + { + _logger.LogError( + ex, + "Failed to added transaction {PlaidTransactionId} to account {AccountId}", + account.Id, + addedTransaction.TransactionId + ); + } } + + _logger.LogInformation("Added {AddedCount} transactions to account {AccountId}", addedCount, account.Id); + + return addedCount; } } \ No newline at end of file diff --git a/src/FiscalOS.Infra/Transactions/Plaid/PlaidExtensions.cs b/src/FiscalOS.Infra/Transactions/Plaid/PlaidExtensions.cs index 07abb16..4661d54 100644 --- a/src/FiscalOS.Infra/Transactions/Plaid/PlaidExtensions.cs +++ b/src/FiscalOS.Infra/Transactions/Plaid/PlaidExtensions.cs @@ -5,6 +5,7 @@ public static class PlaidExtensions #pragma warning disable CA1034 extension(Transaction transaction) { + public string Merchant => transaction.MerchantName ?? "Unknown merchant"; public string Description => transaction.OriginalDescription ?? ""; public DateTimeOffset PostedDate => transaction.Datetime ?? ( transaction.Date.HasValue diff --git a/src/FiscalOS.Infra/Transactions/Plaid/PlaidModifiedTransactionHandler.cs b/src/FiscalOS.Infra/Transactions/Plaid/PlaidModifiedTransactionHandler.cs index 9734f86..97ace41 100644 --- a/src/FiscalOS.Infra/Transactions/Plaid/PlaidModifiedTransactionHandler.cs +++ b/src/FiscalOS.Infra/Transactions/Plaid/PlaidModifiedTransactionHandler.cs @@ -5,7 +5,11 @@ namespace FiscalOS.Infra.Transactions.Plaid; internal interface IPlaidModifiedTransactionHandler { - void Handle(Account account, IEnumerable modified); + Task HandleAsync( + Account account, + IEnumerable modified, + CancellationToken ct + ); } internal sealed class PlaidModifiedTransactionHandler : IPlaidModifiedTransactionHandler @@ -30,46 +34,68 @@ internal sealed class PlaidModifiedTransactionHandler : IPlaidModifiedTransactio ); } - public void Handle(Account account, IEnumerable modified) + public async Task HandleAsync(Account account, IEnumerable modified, CancellationToken ct) { var modifiedTransactionIds = modified.Select(t => t.TransactionId); - var existingModifiedTransactions = _appDbContext.Transactions + var existingModifiedTransactions = await _appDbContext.Transactions .Where(t => t.Metadata is PlaidTransactionMetadata && modifiedTransactionIds.Contains(((PlaidTransactionMetadata)t.Metadata).PlaidId)) - .ToList(); + .ToListAsync(ct) + .ConfigureAwait(false); + + var modifiedCount = 0; foreach (var existing in existingModifiedTransactions) { - if (existing.Metadata is not PlaidTransactionMetadata plaidMetadata) + try { - _logger.LogWarning( - "Existing transaction {TransactionId} has non-Plaid metadata. Skipping update for this transaction.", + if (existing.Metadata is not PlaidTransactionMetadata plaidMetadata) + { + _logger.LogWarning( + "Existing transaction {TransactionId} has non-Plaid metadata. Skipping update for this transaction.", + existing.Id + ); + modifiedCount++; + continue; + } + + var plaidModifiedTransaction = modified.FirstOrDefault(t => t.TransactionId == plaidMetadata.PlaidId); + + if (plaidModifiedTransaction is null) + { + _logger.LogWarning( + "No corresponding modified transaction found in Plaid response for existing transaction {TransactionId}. Skipping update for this transaction.", + existing.Id + ); + modifiedCount++; + continue; + } + + var newTransactionData = Transaction.From( + existing.UserId, + existing.AccountId, + plaidModifiedTransaction.Merchant, + plaidModifiedTransaction.Description, + plaidModifiedTransaction.Amount, + plaidModifiedTransaction.PostedDate, + plaidMetadata + ); + + _appDbContext.Entry(existing).CurrentValues.SetValues(newTransactionData); + modifiedCount++; + } + catch (Exception ex) + { + _logger.LogError( + ex, + "Failed to update transaction {TransactionId} for account {AccountId}", + account.Id, existing.Id ); - continue; } - - var plaidModifiedTransaction = modified.FirstOrDefault(t => t.TransactionId == plaidMetadata.PlaidId); - - if (plaidModifiedTransaction is null) - { - _logger.LogWarning( - "No corresponding modified transaction found in Plaid response for existing transaction {TransactionId}. Skipping update for this transaction.", - existing.Id - ); - continue; - } - - var newTransactionData = Transaction.From( - existing.UserId, - existing.AccountId, - plaidModifiedTransaction.MerchantName, - plaidModifiedTransaction.Description, - plaidModifiedTransaction.Amount, - plaidModifiedTransaction.PostedDate, - plaidMetadata - ); - - _appDbContext.Entry(existing).CurrentValues.SetValues(newTransactionData); } + + _logger.LogInformation("Modified {ModifiedCount} transactions for account {AccountId}", modifiedCount, account.Id); + + return modifiedCount; } } \ No newline at end of file diff --git a/src/FiscalOS.Infra/Transactions/Plaid/PlaidRemovedTransactionHandler.cs b/src/FiscalOS.Infra/Transactions/Plaid/PlaidRemovedTransactionHandler.cs index cf04d4c..35f7062 100644 --- a/src/FiscalOS.Infra/Transactions/Plaid/PlaidRemovedTransactionHandler.cs +++ b/src/FiscalOS.Infra/Transactions/Plaid/PlaidRemovedTransactionHandler.cs @@ -2,32 +2,54 @@ namespace FiscalOS.Infra.Transactions.Plaid; internal interface IPlaidRemovedTransactionHandler { - Task HandleAsync(IEnumerable removed, CancellationToken cancellationToken); + Task HandleAsync(IEnumerable removed, CancellationToken cancellationToken); } internal sealed class PlaidRemovedTransactionHandler : IPlaidRemovedTransactionHandler { private readonly AppDbContext _appDbContext; + private readonly ILogger _logger; - private PlaidRemovedTransactionHandler(AppDbContext appDbContext) + private PlaidRemovedTransactionHandler( + AppDbContext appDbContext, + ILogger logger + ) { _appDbContext = appDbContext; + _logger = logger; } public static PlaidRemovedTransactionHandler From(IServiceProvider serviceProvider) { return new( - serviceProvider.GetRequiredService() + serviceProvider.GetRequiredService(), + serviceProvider.GetRequiredService>() ); } - public async Task HandleAsync(IEnumerable removed, CancellationToken cancellationToken) + public async Task HandleAsync(IEnumerable removed, CancellationToken cancellationToken) { + var removedCount = 0; var removedIdsList = removed.Select(t => t.TransactionId); - await _appDbContext.Transactions - .Where(t => t.Metadata is PlaidTransactionMetadata && removedIdsList.Contains(((PlaidTransactionMetadata)t.Metadata).PlaidId)) - .ExecuteDeleteAsync(cancellationToken) - .ConfigureAwait(false); + try + { + removedCount = await _appDbContext.Transactions + .Where(t => t.Metadata is PlaidTransactionMetadata && removedIdsList.Contains(((PlaidTransactionMetadata)t.Metadata).PlaidId)) + .ExecuteDeleteAsync(cancellationToken) + .ConfigureAwait(false); + } + catch (Exception ex) + { + _logger.LogError( + ex, + "Failed to removed transactions {TransactionIds}", + removedIdsList + ); + } + + _logger.LogInformation("Removed {RemovedCount} transactions", removedCount); + + return removedCount; } } \ No newline at end of file diff --git a/src/FiscalOS.Infra/Transactions/Plaid/PlaidTransactionProcessor.cs b/src/FiscalOS.Infra/Transactions/Plaid/PlaidTransactionProcessor.cs index 7c1074d..5c648ed 100644 --- a/src/FiscalOS.Infra/Transactions/Plaid/PlaidTransactionProcessor.cs +++ b/src/FiscalOS.Infra/Transactions/Plaid/PlaidTransactionProcessor.cs @@ -4,7 +4,7 @@ namespace FiscalOS.Infra.Transactions.Plaid; internal interface IPlaidTransactionProcessor { - Task ProcessAsync( + Task ProcessAsync( Account account, TransactionsSyncResponse syncResponse, CancellationToken cancellationToken @@ -37,10 +37,31 @@ internal sealed class PlaidTransactionProcessor : IPlaidTransactionProcessor ); } - public async Task ProcessAsync(Account account, TransactionsSyncResponse syncResponse, CancellationToken cancellationToken) + public async Task ProcessAsync( + Account account, + TransactionsSyncResponse syncResponse, + CancellationToken cancellationToken + ) { - _addedHandler.Handle(account, syncResponse.Added); - _modifiedHandler.Handle(account, syncResponse.Modified); - await _removedHandler.HandleAsync(syncResponse.Removed, cancellationToken).ConfigureAwait(false); + var numAdded = await _addedHandler.HandleAsync( + account, + syncResponse.Added, + cancellationToken + ).ConfigureAwait(false); + + var numModified = await _modifiedHandler.HandleAsync( + account, + syncResponse.Modified, + cancellationToken + ).ConfigureAwait(false); + + var numRemoved = await _removedHandler.HandleAsync( + syncResponse.Removed, + cancellationToken + ).ConfigureAwait(false); + + return numAdded == syncResponse.Added.Count && + numModified == syncResponse.Modified.Count && + numRemoved == syncResponse.Removed.Count; } } \ No newline at end of file diff --git a/src/FiscalOS.Infra/Transactions/Plaid/PlaidTransactionSyncer.cs b/src/FiscalOS.Infra/Transactions/Plaid/PlaidTransactionSyncer.cs index dfa8ac4..c4eefff 100644 --- a/src/FiscalOS.Infra/Transactions/Plaid/PlaidTransactionSyncer.cs +++ b/src/FiscalOS.Infra/Transactions/Plaid/PlaidTransactionSyncer.cs @@ -83,11 +83,22 @@ internal sealed class PlaidTransactionSyncer : IPlaidTransactionSyncer break; } - plaidAccountMetadata.SetCursor(response.NextCursor); var balance = Balance.From(plaidAccount.Current, plaidAccount.Available, plaidAccount.CurrencyCode); account.AddBalance(balance); - await _transactionProcessor.ProcessAsync(account, response, cancellationToken).ConfigureAwait(false); + var isProcessedSuccessfully = await _transactionProcessor.ProcessAsync(account, response, cancellationToken).ConfigureAwait(false); + + if (isProcessedSuccessfully is false) + { + _logger.LogWarning( + "Unable to process transactions for request {RequestId} for account {AccountId} successfully", + response.RequestId, + account.Id + ); + break; + } + + plaidAccountMetadata.SetCursor(response.NextCursor); } } } \ No newline at end of file diff --git a/src/FiscalOS.Infra/Transactions/Plaid/SyncUpdatesProcessor.cs b/src/FiscalOS.Infra/Transactions/Plaid/SyncUpdatesProcessor.cs index e0fb44c..8a2c0ba 100644 --- a/src/FiscalOS.Infra/Transactions/Plaid/SyncUpdatesProcessor.cs +++ b/src/FiscalOS.Infra/Transactions/Plaid/SyncUpdatesProcessor.cs @@ -74,7 +74,7 @@ internal sealed class SyncUpdatesProcessor : IAsyncQueueProcessor