Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
using System.Diagnostics.CodeAnalysis;
using FluentMigrator;
using JetBrains.Annotations;
using static System.Int32;

namespace Elsa.Persistence.Dapper.Migrations.Runtime;

/// <summary>
/// Creates the <c>KeyValues</c> table used by <c>DapperKeyValueStore</c> (#263)
/// and adds <c>BookmarkQueueItems.SerializedOptions</c> (#264).
/// Both operations are existence-guarded so hand-created workarounds keep working.
/// </summary>
[Migration(20008, "Elsa:Runtime:V3.9")]
[PublicAPI]
[SuppressMessage("ReSharper", "InconsistentNaming")]
public class V3_9 : Migration
{
/// <inheritdoc />
public override void Up()
{
if (!Schema.Table("KeyValues").Exists())
Create.Table("KeyValues")
.WithColumn("Id").AsString().PrimaryKey()
.WithColumn("TenantId").AsString().Nullable()
.WithColumn("Value").AsString(MaxValue).Nullable();

if (!Schema.Table("BookmarkQueueItems").Column("SerializedOptions").Exists())
Alter.Table("BookmarkQueueItems").AddColumn("SerializedOptions").AsString(MaxValue).Nullable();
}

/// <inheritdoc />
/// <remarks>
/// Leave KeyValues and SerializedOptions in place. Rolling them back would
/// drop a hand-created KeyValues table or SerializedOptions column and their data.
/// </remarks>
public override void Down()
{
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -154,6 +154,17 @@ public static ParameterizedQuery IsNull(this ParameterizedQuery query, string fi
return query;
}

/// <summary>
/// Matches a default-tenant stamp: <c>NULL</c> or empty string.
/// Legacy rows may use either; <see cref="Elsa.Common.Multitenancy.Tenant.DefaultTenantId"/> is <c>''</c>.
/// </summary>
public static ParameterizedQuery IsNullOrEmpty(this ParameterizedQuery query, string field)
{
var ident = query.QuoteIdent(field);
query.Sql.AppendLine($"and ({ident} is null or {ident} = '')");
return query;
}

/// <summary>
/// Appends an IS NOT NULL clause to the query.
/// </summary>
Expand Down Expand Up @@ -241,9 +252,9 @@ public static ParameterizedQuery StartsWith(this ParameterizedQuery query, strin
if (!startsWith || value == null || string.IsNullOrWhiteSpace(value))
return query;

var searchTermLike = $"{value}%";
query.Sql.AppendLine($"and {query.QuoteIdent(field)} like @SearchTermLike");
query.Parameters.Add($"@{field}", searchTermLike);
var parameterName = $"@{field}StartsWith";
query.Sql.AppendLine($"and {query.QuoteIdent(field)} like {parameterName}");
query.Parameters.Add(parameterName, $"{value}%");

return query;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,12 +42,23 @@ public Task DeleteAsync(string key, CancellationToken cancellationToken)
return store.DeleteAsync(query => query.Is(nameof(KeyValuePairRecord.Id), key), cancellationToken);
}

/// <summary>
/// Atomic delete: returns true only when the DELETE affected a row.
/// Becomes <c>IKeyValueStore.TryDeleteAsync</c> once ElsaVersion includes elsa-core#8538.
/// </summary>
public async Task<bool> TryDeleteAsync(string key, CancellationToken cancellationToken = default)
{
return await store.DeleteAsync(query => query.Is(nameof(KeyValuePairRecord.Id), key), cancellationToken) > 0;
}

private void ApplyFilter(ParameterizedQuery query, KeyValueFilter filter)
{
query
.Is(nameof(KeyValuePairRecord.Id), filter.Key)
.In(nameof(KeyValuePairRecord.Id), filter.Keys)
.StartsWith(nameof(KeyValuePairRecord.Id), filter.StartsWith, filter.Key);
if (filter.StartsWith)
query.StartsWith(nameof(KeyValuePairRecord.Id), true, filter.Key);
else
query.Is(nameof(KeyValuePairRecord.Id), filter.Key);

query.In(nameof(KeyValuePairRecord.Id), filter.Keys);
}

private KeyValuePairRecord Map(SerializedKeyValuePair source)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -542,7 +542,11 @@ private void ApplyTenantFilter(ParameterizedQuery query, bool tenantAgnostic = f

var tenant = tenantAccessor.Tenant;
var tenantId = tenant?.Id;
query.Is(nameof(Record.TenantId), (object?)tenantId ?? DBNull.Value);
// Default tenant is '' in Elsa and NULL on some legacy rows (#245 / #260).
if (string.IsNullOrEmpty(tenantId))
query.IsNullOrEmpty(nameof(Record.TenantId));
else
query.Is(nameof(Record.TenantId), tenantId);
}

private void SetTenantId(T record)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -509,7 +509,12 @@ private FilterDefinition<TDocument> CreateTenantOwnedUpsertFilter(TDocument docu
& Builders<TDocument>.Filter.Where(x => (x as Entity)!.TenantId == writerTenantId || (x as Entity)!.TenantId == Tenant.AgnosticTenantId);
}

private FilterDefinition<TDocument> ApplyTenantScope(FilterDefinition<TDocument> filter, bool tenantAgnostic)
/// <summary>
/// ANDs the current tenant onto <paramref name="filter"/>. No-ops when
/// <paramref name="tenantAgnostic"/> is true or <typeparamref name="TDocument"/> is not an <see cref="Entity"/>.
/// Writes using this filter match only the ambient tenant, not tenant-agnostic ("*") rows.
/// </summary>
public FilterDefinition<TDocument> ApplyTenantScope(FilterDefinition<TDocument> filter, bool tenantAgnostic = false)
{
if (tenantAgnostic || !typeof(Entity).IsAssignableFrom(typeof(TDocument)))
return filter;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
using Elsa.KeyValues.Models;
using Elsa.Persistence.MongoDb.Common;
using JetBrains.Annotations;
using MongoDB.Driver;
using MongoDB.Driver.Linq;

namespace Elsa.Persistence.MongoDb.Modules.Runtime;
Expand Down Expand Up @@ -37,6 +38,17 @@ public Task DeleteAsync(string key, CancellationToken cancellationToken)
return keyValueMongoDbStore.DeleteWhereAsync(x => x.Key == key, cancellationToken);
}

/// <summary>
/// Atomic delete: returns true only when <c>DeleteOneAsync</c> removed a document.
/// Becomes <c>IKeyValueStore.TryDeleteAsync</c> once ElsaVersion includes elsa-core#8538.
/// </summary>
public async Task<bool> TryDeleteAsync(string key, CancellationToken cancellationToken = default)
{
var filter = keyValueMongoDbStore.ApplyTenantScope(Builders<SerializedKeyValuePair>.Filter.Eq(x => x.Key, key));
var result = await keyValueMongoDbStore.GetCollection().DeleteOneAsync(filter, cancellationToken);
return result.DeletedCount > 0;
}

private IQueryable<SerializedKeyValuePair> Filter(IQueryable<SerializedKeyValuePair> queryable, KeyValueFilter filter)
{
return filter.Apply(queryable);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,158 @@
using Elsa.Common.Multitenancy;
using Elsa.KeyValues.Contracts;
using Elsa.KeyValues.Entities;
using Elsa.KeyValues.Models;
using Elsa.Persistence.MongoDb.Common;
using Elsa.Persistence.MongoDb.Modules.Runtime;
using MongoDB.Driver;
using Testcontainers.MongoDb;

namespace Elsa.MongoDb.UnitTests;

/// <summary>
/// Atomic <c>TryDeleteAsync</c> against a shared Mongo collection (#260).
/// Calls the Mongo store method directly; <c>IKeyValueStore.TryDeleteAsync</c> is not
/// in 3.10.0-preview.5722 yet (waiting for a core pin at or after c1c935ce).
/// </summary>
public sealed class MongoKeyValueStoreTryDeleteTests : IAsyncLifetime
{
private const string LegacyKey = "elsa.quiescence.pause.default";
private readonly MongoDbContainer _container = new MongoDbBuilder().WithImage("mongo:7.0.24").Build();
private readonly TestTenantAccessor _tenants = new();
private MongoClient? _client;
private IMongoCollection<SerializedKeyValuePair> _collection = null!;
private MongoDbStore<SerializedKeyValuePair> _mongoStore = null!;

public async Task InitializeAsync()
{
try
{
await _container.StartAsync();
_client = new MongoClient(_container.GetConnectionString());
var database = _client.GetDatabase($"elsa-trydelete-{Guid.NewGuid():N}");
_collection = database.GetCollection<SerializedKeyValuePair>("key_values");
_mongoStore = new MongoDbStore<SerializedKeyValuePair>(_collection, _tenants);
}
catch
{
await DisposeAsync();
throw;
}
}

public async Task DisposeAsync()
{
try
{
_client?.Dispose();
}
finally
{
await _container.DisposeAsync();
}
}

[Fact(DisplayName = "#260: two Mongo stores racing TryDeleteAsync: exactly one returns true")]
public async Task TryDeleteAsync_TwoNodes_ExactlyOneReturnsTrue()
{
using var tenant = _tenants.PushContext(Tenant.Default);
var nodeA = new MongoKeyValueStore(_mongoStore);
var nodeB = new MongoKeyValueStore(_mongoStore);
await nodeA.SaveAsync(Pair(LegacyKey, "legacy-maintenance"), CancellationToken.None);

var results = await Task.WhenAll(nodeA.TryDeleteAsync(LegacyKey), nodeB.TryDeleteAsync(LegacyKey));

Assert.Equal(1, results.Count(won => won));
Assert.Equal(1, results.Count(won => !won));
Assert.Null(await nodeA.FindAsync(new KeyValueFilter { Key = LegacyKey }, CancellationToken.None));
}

[Fact(DisplayName = "#260: the default find-then-delete lets both racers win")]
public async Task DefaultTryDeleteAsync_TwoNodes_BothReturnTrue()
{
using var tenant = _tenants.PushContext(Tenant.Default);
var inner = new MongoKeyValueStore(_mongoStore);
await inner.SaveAsync(Pair(LegacyKey, "legacy-maintenance"), CancellationToken.None);
var dim = new BarrierFindThenDelete(inner);

var results = await Task.WhenAll(dim.TryDeleteAsync(LegacyKey), dim.TryDeleteAsync(LegacyKey));

Assert.Equal(2, results.Count(won => won));
Assert.Null(await inner.FindAsync(new KeyValueFilter { Key = LegacyKey }, CancellationToken.None));
}

[Fact(DisplayName = "#260: TryDeleteAsync returns false when the key is missing")]
public async Task TryDeleteAsync_NotFound_ReturnsFalse()
{
using var tenant = _tenants.PushContext(Tenant.Default);
var store = new MongoKeyValueStore(_mongoStore);

Assert.False(await store.TryDeleteAsync(LegacyKey));
Assert.False(await store.TryDeleteAsync(LegacyKey));
}

[Fact(DisplayName = "#260: a NULL TenantId legacy row is found and TryDeleted by the default tenant")]
public async Task DefaultTenant_FindsAndTryDeletesNullTenantIdRow()
{
await _collection.InsertOneAsync(new SerializedKeyValuePair
{
Key = LegacyKey,
SerializedValue = "legacy-null",
TenantId = null
});

using var tenant = _tenants.PushContext(Tenant.Default);
var store = new MongoKeyValueStore(_mongoStore);

var found = await store.FindAsync(new KeyValueFilter { Key = LegacyKey }, CancellationToken.None);
Assert.Equal("legacy-null", found?.SerializedValue);
Assert.True(await store.TryDeleteAsync(LegacyKey));
Assert.Null(await store.FindAsync(new KeyValueFilter { Key = LegacyKey }, CancellationToken.None));
}

private static SerializedKeyValuePair Pair(string key, string value) => new()
{
Key = key,
SerializedValue = value
};

private sealed class TestTenantAccessor : ITenantAccessor
{
public string TenantId => Tenant?.Id ?? Tenant.DefaultTenantId;
public Tenant? Tenant { get; private set; }

public IDisposable PushContext(Tenant? tenant)
{
var previousTenant = Tenant;
Tenant = tenant;
return new Restore(() => Tenant = previousTenant);
}

private sealed class Restore(Action restore) : IDisposable
{
public void Dispose() => restore();
}
}

/// <summary>
/// The core default <c>TryDeleteAsync</c> (find, then delete, always true if found).
/// Both Finds complete before either Delete so the race is deterministic.
/// </summary>
private sealed class BarrierFindThenDelete(IKeyValueStore inner)
{
private readonly TaskCompletionSource _bothFound = new(TaskCreationOptions.RunContinuationsAsynchronously);
private int _finds;

public async Task<bool> TryDeleteAsync(string key)
{
var found = await inner.FindAsync(new KeyValueFilter { Key = key }, CancellationToken.None);
if (Interlocked.Increment(ref _finds) == 2)
_bothFound.TrySetResult();
await _bothFound.Task;
if (found is null)
return false;
await inner.DeleteAsync(key, CancellationToken.None);
return true;
}
}
}
Loading
Loading