Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -19,13 +19,17 @@ public class V3_9 : Migration
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 />
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,9 @@
using System.Linq.Expressions;
using Elsa.Common.Entities;
using Elsa.Common.Models;
using Elsa.Persistence.Dapper.Models;
using Elsa.Extensions;
using Elsa.Persistence.Dapper.Dialects;
using Elsa.Persistence.Dapper.Models;
using JetBrains.Annotations;

namespace Elsa.Persistence.Dapper.Extensions;
Expand All @@ -15,6 +16,8 @@ namespace Elsa.Persistence.Dapper.Extensions;
[PublicAPI]
public static class ParameterizedQueryBuilderExtensions
{
private const char LikeEscapeCharacter = '!';

/// <summary>
/// Begins a SELECT FROM query.
/// </summary>
Expand Down Expand Up @@ -186,7 +189,17 @@ public static ParameterizedQuery LessThan(this ParameterizedQuery query, string
{
if (value == null) return query;

query.Sql.AppendLine($"and {query.QuoteIdent(field)} < @{field}");
var identifier = query.QuoteIdent(field);
if (query.Dialect is SqliteDialect && value is DateTimeOffset)
{
// SQLite date functions compare at millisecond precision. Keep this exclusive so a precision tie waits for the next scan.
query.Sql.AppendLine($"and julianday({identifier}) < julianday(@{field})");
}
else
{
query.Sql.AppendLine($"and {identifier} < @{field}");
}

query.Parameters.Add($"@{field}", value);

return query;
Expand Down Expand Up @@ -269,8 +282,25 @@ public static ParameterizedQuery StartsWith(this ParameterizedQuery query, strin
return query;

var parameterName = $"@{field}StartsWith";
query.Sql.AppendLine($"and {query.QuoteIdent(field)} like {parameterName}");
query.Parameters.Add(parameterName, $"{value}%");
var isSqlServer = query.Dialect is SqlServerDialect;
// PostgreSQL treats backslash as LIKE's default escape, so use our explicit escape for literal paths.
var needsEscaping = value.IndexOfAny(['%', '_', LikeEscapeCharacter]) >= 0 ||
value.Contains('\\') ||
(isSqlServer && value.Contains('['));
var escapeClause = needsEscaping ? $" escape '{LikeEscapeCharacter}'" : string.Empty;
var escapedValue = value
.Replace("!", "!!")
.Replace("%", "!%")
.Replace("_", "!_");

// SQL Server also treats '[' as a LIKE pattern character; the other supported dialects do not.
if (isSqlServer)
{
escapedValue = escapedValue.Replace("[", "![");
}

query.Sql.AppendLine($"and {query.QuoteIdent(field)} like {parameterName}{escapeClause}");
query.Parameters.Add(parameterName, $"{escapedValue}%");

return query;
}
Expand Down Expand Up @@ -527,4 +557,4 @@ private static string QuoteIdent(this ParameterizedQuery query, string name) =>

private static string BoolLit(this ParameterizedQuery query, bool value) =>
query.Dialect.BooleanLiteral(value);
}
}
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
using System.Data.Common;
using Elsa.Persistence.Dapper.Extensions;
using Elsa.Persistence.Dapper.Models;
using Elsa.Persistence.Dapper.Modules.Runtime.Records;
Expand All @@ -6,6 +7,9 @@
using Elsa.KeyValues.Entities;
using Elsa.KeyValues.Models;
using JetBrains.Annotations;
using Microsoft.Data.SqlClient;
using Microsoft.Data.Sqlite;
using Npgsql;

namespace Elsa.Persistence.Dapper.Modules.Runtime.Stores;

Expand All @@ -16,12 +20,87 @@ namespace Elsa.Persistence.Dapper.Modules.Runtime.Stores;
internal class DapperKeyValueStore(Store<KeyValuePairRecord> store) : IKeyValueStore
{
/// <inheritdoc />
public Task SaveAsync(SerializedKeyValuePair keyValuePair, CancellationToken cancellationToken)
public async Task SaveAsync(SerializedKeyValuePair keyValuePair, CancellationToken cancellationToken)
{
var record = Map(keyValuePair);
return store.SaveAsync(record, cancellationToken);
ArgumentNullException.ThrowIfNull(record.Id);

if (await TryUpdateOwnedAsync(record, cancellationToken))
{
return;
}

var existing = await FindByGlobalIdAsync(record.Id, cancellationToken);
if (existing != null)
{
// An owned row may have been inserted since the first update.
if (await TryUpdateOwnedAsync(record, cancellationToken))
{
return;
}
throw OwnershipConflict(record.Id, existing);
}

try
{
// Add stamps the ambient tenant, ignoring ownership supplied by the caller.
await store.AddAsync(record, cancellationToken);
}
catch (DbException exception) when (IsDuplicateKey(exception))
{
// A competing insert is safe to retry only through the tenant-scoped update.
if (await TryUpdateOwnedAsync(record, cancellationToken))
{
return;
}
existing = await FindByGlobalIdAsync(record.Id, cancellationToken);
if (existing != null)
{
throw OwnershipConflict(record.Id, existing);
}
throw;
}
}

private async Task<bool> TryUpdateOwnedAsync(KeyValuePairRecord record, CancellationToken cancellationToken)
{
while (true)
{
cancellationToken.ThrowIfCancellationRequested();
var updated = await store.UpdateAsync(record, [x => x.Value], query => query.Is(nameof(KeyValuePairRecord.Id), record.Id), cancellationToken);
if (updated > 0)
{
return true;
}

var owned = await store.FindAsync(query => query.Is(nameof(KeyValuePairRecord.Id), record.Id), cancellationToken);
if (owned == null)
{
return false;
}
if (owned.Value == record.Value)
{
// Some providers report zero affected rows when the value is unchanged.
return true;
}
// A same-owner insert after a missed update still needs the requested value applied.
}
}

private Task<KeyValuePairRecord?> FindByGlobalIdAsync(string id, CancellationToken cancellationToken) =>
store.FindAsync(query => query.Is(nameof(KeyValuePairRecord.Id), id), tenantAgnostic: true, cancellationToken);

private static InvalidOperationException OwnershipConflict(string id, KeyValuePairRecord existing) =>
new($"Cannot save key '{id}' because it belongs to {(existing.TenantId == "*" ? "the tenant-agnostic ('*') scope" : "another tenant")}.");

private static bool IsDuplicateKey(DbException exception) => exception switch
{
SqliteException { SqliteExtendedErrorCode: 1555 or 2067 } => true,
PostgresException { SqlState: PostgresErrorCodes.UniqueViolation } => true,
SqlException { Number: 2601 or 2627 } => true,
_ => false
};

/// <inheritdoc />
public async Task<SerializedKeyValuePair?> FindAsync(KeyValueFilter filter, CancellationToken cancellationToken)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -136,17 +136,18 @@ private static void SerializeJsonNode(BsonSerializationContext context, JsonNode
reader.ReadNull();
return null;
case BsonType.Document:
return DeserializeDocument(BsonDocumentSerializer.Instance.Deserialize(context));
// A nominal JsonObject may be a legacy map whose fields resemble another node kind's envelope.
return DeserializeDocument(BsonDocumentSerializer.Instance.Deserialize(context), requireObjectEnvelope: typeof(TNode) == typeof(JsonObject));
case BsonType.Array:
return DeserializeLegacyArray(BsonArraySerializer.Instance.Deserialize(context));
default:
return DeserializePrimitiveJsonValue(ReadPrimitive(reader, bsonType));
}
}

private static JsonNode DeserializeDocument(BsonDocument document)
private static JsonNode DeserializeDocument(BsonDocument document, bool requireObjectEnvelope = false)
{
if (IsTaggedJsonNode(document))
if (IsTaggedJsonNode(document, requireObjectEnvelope))
return DeserializeTagged(document);

// Legacy class-map / dictionary format written when only JsonNode was registered:
Expand All @@ -157,15 +158,21 @@ private static JsonNode DeserializeDocument(BsonDocument document)
return obj;
}

private static bool IsTaggedJsonNode(BsonDocument document)
private static bool IsTaggedJsonNode(BsonDocument document, bool requireObjectEnvelope)
{
if (document.ElementCount != 2 || !document.Contains("type") || !document.Contains("value"))
return false;

if (document["type"].BsonType != BsonType.String)
return false;

return document["type"].AsString switch
var type = document["type"].AsString;
if (requireObjectEnvelope && type != "JsonObject")
{
return false;
}

return type switch
{
"JsonObject" or "JsonArray" => document["value"].BsonType == BsonType.String,
"JsonValue" => true,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ public void LessThan_WithValue_AppendsExclusiveComparison()

var sql = query.Sql.ToString();

Assert.Contains("and UpdatedAt < @UpdatedAt", sql, StringComparison.Ordinal);
Assert.Contains("and julianday(UpdatedAt) < julianday(@UpdatedAt)", sql, StringComparison.Ordinal);
Assert.DoesNotContain("<=", sql, StringComparison.Ordinal);
Assert.Equal(cutoff, query.Parameters.Get<DateTimeOffset>("UpdatedAt"));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -195,6 +195,49 @@ public void PolymorphicSerializer_DeserializesLegacyJsonObjectMap()
Assert.Equal(3, obj["Count"]!.GetValue<int>());
}

[Theory(DisplayName = "A legacy JsonObject with incompatible envelope fields preserves its object type")]
[InlineData("JsonValue", "legacy-value")]
[InlineData("JsonArray", "[]")]
public void Deserialize_LegacyIncompatibleEnvelopeShapedJsonObject_PreservesFields(string type, string value)
{
var legacy = new BsonDocument { { "type", type }, { "value", value } };

var restored = DeserializeDocument(new JsonNodeBsonConverter<JsonObject>(), legacy);

Assert.Equal(2, restored.Count);
Assert.Equal(type, restored["type"]!.GetValue<string>());
Assert.Equal(value, restored["value"]!.GetValue<string>());
}

[Theory(DisplayName = "Polymorphic JsonObject dispatch preserves a legacy incompatible-envelope-shaped object")]
[InlineData("JsonValue", "legacy-value")]
[InlineData("JsonArray", "[]")]
public void PolymorphicSerializer_LegacyIncompatibleEnvelopeShapedJsonObject_PreservesObjectType(string type, string value)
{
var document = new BsonDocument
{
{ "$type", typeof(JsonObject).GetSimpleAssemblyQualifiedName() },
{ "$value", new BsonDocument { { "type", type }, { "value", value } } }
};

var restored = DeserializeDocument(new PolymorphicSerializer(), document);

var obj = Assert.IsType<JsonObject>(restored);
Assert.Equal(type, obj["type"]!.GetValue<string>());
Assert.Equal(value, obj["value"]!.GetValue<string>());
}

[Fact(DisplayName = "JsonNode dispatch still restores a genuine scalar envelope as a JsonValue")]
public void JsonNode_RoundTrips_GenuineScalarEnvelope()
{
JsonNode original = JsonValue.Create("scalar-value")!;

var restored = RoundTrip(new JsonNodeBsonConverter(), original);

var value = Assert.IsAssignableFrom<JsonValue>(restored);
Assert.Equal("scalar-value", value.GetValue<string>());
}

[Fact(DisplayName = "A type/value document that is not a two-field string envelope deserializes as a map")]
public void Deserialize_AmbiguousTypeValueDocument_IsMapNotEnvelope()
{
Expand Down
Loading
Loading