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".
This commit is contained in:
Sipke Schoorstra 2024-01-06 15:16:57 +01:00
parent 70530aa9fe
commit 6dee8ee26b
3 changed files with 17 additions and 11 deletions

View file

@ -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()

View file

@ -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)

View file

@ -204,10 +204,15 @@ public class Store<T> where T : notnull
/// <param name="cancellationToken">The cancellation token.</param>
public async Task SaveManyAsync(IEnumerable<T> 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<T> where T : notnull
using var connection = _dbConnectionProvider.GetConnection();
await query.ExecuteAsync(connection);
}
/// <summary>
/// Adds the specified record.
/// </summary>
@ -243,26 +248,27 @@ public class Store<T> where T : notnull
using var connection = _dbConnectionProvider.GetConnection();
return await query.ExecuteAsync(connection);
}
/// <summary>
/// Deletes all records matching the specified query.
/// </summary>
/// <param name="filter">The conditions to apply to the query.</param>
/// <param name="pageArgs">The page arguments.</param>
/// <param name="orderFields">The fields by which to order the results.</param>
/// <param name="primaryKey">The primary key.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>The number of records deleted.</returns>
public async Task<long> DeleteAsync(Action<ParameterizedQuery> filter, PageArgs pageArgs, IEnumerable<OrderField> orderFields, CancellationToken cancellationToken = default)
public async Task<long> DeleteAsync(Action<ParameterizedQuery> filter, PageArgs pageArgs, IEnumerable<OrderField> 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);
}
/// <summary>
/// Returns <c>true</c> if any records match the specified query.
/// </summary>
@ -276,7 +282,7 @@ public class Store<T> where T : notnull
using var connection = _dbConnectionProvider.GetConnection();
return await connection.QueryFirstOrDefaultAsync<object>(query.Sql.ToString(), query.Parameters) != null;
}
/// <summary>
/// Returns the number of records matching the specified query.
/// </summary>