|
| 1 | +using ModularityKit.Mutator.Governance.Abstractions.Requests.Model; |
| 2 | +using ModularityKit.Mutator.Governance.Redis.Storage.Persistence.Reading; |
| 3 | +using ModularityKit.Mutator.Governance.Redis.Storage.Persistence.Writing; |
| 4 | +using StackExchange.Redis; |
| 5 | + |
| 6 | +namespace ModularityKit.Mutator.Governance.Redis.Storage.Persistence; |
| 7 | + |
| 8 | +/// <summary> |
| 9 | +/// Coordinates Redis persistence operations for governed mutation requests. |
| 10 | +/// </summary> |
| 11 | +internal sealed class RedisMutationRequestPersistence( |
| 12 | + IDatabase database, |
| 13 | + RedisMutationRequestPersistenceRecordFactory recordFactory, |
| 14 | + RedisMutationRequestPersistenceDocumentReader documentReader, |
| 15 | + RedisMutationRequestTransactionWriter transactionWriter) |
| 16 | +{ |
| 17 | + private readonly IDatabase _database = database ?? throw new ArgumentNullException(nameof(database)); |
| 18 | + private readonly RedisMutationRequestPersistenceRecordFactory _recordFactory = recordFactory ?? throw new ArgumentNullException(nameof(recordFactory)); |
| 19 | + private readonly RedisMutationRequestPersistenceDocumentReader _documentReader = documentReader ?? throw new ArgumentNullException(nameof(documentReader)); |
| 20 | + private readonly RedisMutationRequestTransactionWriter _transactionWriter = transactionWriter ?? throw new ArgumentNullException(nameof(transactionWriter)); |
| 21 | + |
| 22 | + /// <summary> |
| 23 | + /// Creates a new governed mutation request in Redis with an initial revision. |
| 24 | + /// </summary> |
| 25 | + /// <param name="request">The request to create.</param> |
| 26 | + /// <param name="cancellationToken">The cancellation token.</param> |
| 27 | + /// <returns>The persisted request with provider-managed revision values applied.</returns> |
| 28 | + public async Task<MutationRequest> Create(MutationRequest request, CancellationToken cancellationToken = default) |
| 29 | + { |
| 30 | + ArgumentNullException.ThrowIfNull(request); |
| 31 | + cancellationToken.ThrowIfCancellationRequested(); |
| 32 | + |
| 33 | + var persistedRequest = request with { Revision = 0 }; |
| 34 | + var record = _recordFactory.Create(persistedRequest); |
| 35 | + var transaction = _database.CreateTransaction(); |
| 36 | + |
| 37 | + _transactionWriter.WriteCreate(transaction, record, persistedRequest); |
| 38 | + |
| 39 | + var committed = await transaction.ExecuteAsync().ConfigureAwait(false); |
| 40 | + return !committed |
| 41 | + ? throw new InvalidOperationException($"Mutation request '{request.RequestId}' already exists in Redis.") |
| 42 | + : persistedRequest; |
| 43 | + } |
| 44 | + |
| 45 | + /// <summary> |
| 46 | + /// Attempts to store an updated governed mutation request using optimistic concurrency. |
| 47 | + /// </summary> |
| 48 | + /// <param name="request">The request to persist.</param> |
| 49 | + /// <param name="expectedRevision">The expected current revision.</param> |
| 50 | + /// <param name="cancellationToken">The cancellation token.</param> |
| 51 | + /// <returns>The persisted request if the update succeeds; otherwise <see langword="null" />.</returns> |
| 52 | + public async Task<MutationRequest?> TryStore(MutationRequest request, long expectedRevision, CancellationToken cancellationToken = default) |
| 53 | + { |
| 54 | + ArgumentNullException.ThrowIfNull(request); |
| 55 | + cancellationToken.ThrowIfCancellationRequested(); |
| 56 | + |
| 57 | + var currentRequest = await Get(request.RequestId, cancellationToken).ConfigureAwait(false); |
| 58 | + if (currentRequest is null || currentRequest.Revision != expectedRevision) |
| 59 | + return null; |
| 60 | + |
| 61 | + var persistedRequest = request with { Revision = expectedRevision + 1 }; |
| 62 | + var record = _recordFactory.Create(persistedRequest); |
| 63 | + var transaction = _database.CreateTransaction(); |
| 64 | + |
| 65 | + _transactionWriter.WriteUpdate(transaction, record, expectedRevision, currentRequest, persistedRequest); |
| 66 | + |
| 67 | + var committed = await transaction.ExecuteAsync().ConfigureAwait(false); |
| 68 | + return committed ? persistedRequest : null; |
| 69 | + } |
| 70 | + |
| 71 | + /// <summary> |
| 72 | + /// Reads a governed mutation request by identifier. |
| 73 | + /// </summary> |
| 74 | + /// <param name="requestId">The request identifier.</param> |
| 75 | + /// <param name="cancellationToken">The cancellation token.</param> |
| 76 | + /// <returns>The request if it exists; otherwise <see langword="null" />.</returns> |
| 77 | + public async Task<MutationRequest?> Get(string requestId, CancellationToken cancellationToken = default) => |
| 78 | + await _documentReader.GetAsync(requestId, cancellationToken).ConfigureAwait(false); |
| 79 | +} |
0 commit comments