From 6dee8ee26b340f85af9ea7fa6de63844150ff27c Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sat, 6 Jan 2024 15:16:57 +0100 Subject: [PATCH] Refactor Dapper workflow and update migrations Modified the store service to optimize the SaveManyAsync method by converting input to list only once. Also, enhanced deletion query in the store service to enable usage of different primary keys. Made changes in the Dapper migrations, replacing "NodeId" with "ActivityNodeId". --- .../Elsa.Dapper.Migrations/Runtime/Initial.cs | 4 ++-- .../Stores/DapperWorkflowInboxMessageStore.cs | 2 +- src/modules/Elsa.Dapper/Services/Store.cs | 22 ++++++++++++------- 3 files changed, 17 insertions(+), 11 deletions(-) diff --git a/src/modules/Elsa.Dapper.Migrations/Runtime/Initial.cs b/src/modules/Elsa.Dapper.Migrations/Runtime/Initial.cs index ef7acebe4..04fe01593 100644 --- a/src/modules/Elsa.Dapper.Migrations/Runtime/Initial.cs +++ b/src/modules/Elsa.Dapper.Migrations/Runtime/Initial.cs @@ -61,7 +61,7 @@ public class Initial : Migration .WithColumn("ActivityType").AsString().NotNullable().Indexed() .WithColumn("ActivityTypeVersion").AsInt32().NotNullable().Indexed() .WithColumn("ActivityName").AsString().Nullable().Indexed() - .WithColumn("NodeId").AsString().NotNullable().Indexed() + .WithColumn("ActivityNodeId").AsString().NotNullable().Indexed() .WithColumn("EventName").AsString().Nullable().Indexed() .WithColumn("Message").AsString().Nullable() .WithColumn("Source").AsString().Nullable() @@ -86,7 +86,7 @@ public class Initial : Migration .WithColumn("ActivityType").AsString().NotNullable().Indexed() .WithColumn("ActivityTypeVersion").AsInt32().NotNullable().Indexed() .WithColumn("ActivityName").AsString().Nullable().Indexed() - .WithColumn("NodeId").AsString().NotNullable().Indexed() + .WithColumn("ActivityNodeId").AsString().NotNullable().Indexed() .WithColumn("EventName").AsString().Nullable().Indexed() .WithColumn("Message").AsString().Nullable() .WithColumn("Source").AsString().Nullable() diff --git a/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperWorkflowInboxMessageStore.cs b/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperWorkflowInboxMessageStore.cs index e8d258d4f..1420dfbda 100644 --- a/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperWorkflowInboxMessageStore.cs +++ b/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperWorkflowInboxMessageStore.cs @@ -61,7 +61,7 @@ public class DapperWorkflowInboxMessageStore : IWorkflowInboxMessageStore if (pageArgs == null) return await _store.DeleteAsync(q => ApplyFilter(q, filter), cancellationToken); - return await _store.DeleteAsync(q => ApplyFilter(q, filter), pageArgs, new[] { new OrderField(nameof(WorkflowInboxMessage.CreatedAt), OrderDirection.Ascending) }, cancellationToken); + return await _store.DeleteAsync(q => ApplyFilter(q, filter), pageArgs, new[] { new OrderField(nameof(WorkflowInboxMessage.CreatedAt), OrderDirection.Ascending) }, cancellationToken: cancellationToken); } private void ApplyFilter(ParameterizedQuery query, params WorkflowInboxMessageFilter[] filters) diff --git a/src/modules/Elsa.Dapper/Services/Store.cs b/src/modules/Elsa.Dapper/Services/Store.cs index e5dc08b9a..eb6fde73c 100644 --- a/src/modules/Elsa.Dapper/Services/Store.cs +++ b/src/modules/Elsa.Dapper/Services/Store.cs @@ -204,10 +204,15 @@ public class Store where T : notnull /// The cancellation token. public async Task SaveManyAsync(IEnumerable records, string primaryKey = "Id", CancellationToken cancellationToken = default) { + var recordsList = records.ToList(); + + if (!recordsList.Any()) + return; + var query = new ParameterizedQuery(_dbConnectionProvider.Dialect); var currentIndex = 0; - foreach (var record in records) + foreach (var record in recordsList) { var index = currentIndex; query.Upsert(TableName, primaryKey, record, field => $"{field}_{index}"); @@ -217,7 +222,7 @@ public class Store where T : notnull using var connection = _dbConnectionProvider.GetConnection(); await query.ExecuteAsync(connection); } - + /// /// Adds the specified record. /// @@ -243,26 +248,27 @@ public class Store where T : notnull using var connection = _dbConnectionProvider.GetConnection(); return await query.ExecuteAsync(connection); } - + /// /// Deletes all records matching the specified query. /// /// The conditions to apply to the query. /// The page arguments. /// The fields by which to order the results. + /// The primary key. /// The cancellation token. /// The number of records deleted. - public async Task DeleteAsync(Action filter, PageArgs pageArgs, IEnumerable orderFields, CancellationToken cancellationToken = default) + public async Task DeleteAsync(Action filter, PageArgs pageArgs, IEnumerable orderFields, string primaryKey = "Id", CancellationToken cancellationToken = default) { - var selectQuery = _dbConnectionProvider.CreateQuery().From(TableName, "rowid"); + var selectQuery = _dbConnectionProvider.CreateQuery().From(TableName, primaryKey); filter(selectQuery); selectQuery = selectQuery.OrderBy(orderFields.ToArray()).Page(pageArgs); - + var deleteQuery = _dbConnectionProvider.CreateQuery().Delete(TableName, selectQuery); using var connection = _dbConnectionProvider.GetConnection(); return await deleteQuery.ExecuteAsync(connection); } - + /// /// Returns true if any records match the specified query. /// @@ -276,7 +282,7 @@ public class Store where T : notnull using var connection = _dbConnectionProvider.GetConnection(); return await connection.QueryFirstOrDefaultAsync(query.Sql.ToString(), query.Parameters) != null; } - + /// /// Returns the number of records matching the specified query. ///