src/Coast.PostgreSql/Repository/SagaRepository.cs

(开头部分) 6KB

这里只显示每个文件的开头 60 行。登录后可以解锁完整代码。

namespace Coast.PostgreSql.Repository
{
    using System;
    using System.Data;
    using System.Linq;
    using System.Threading;
    using System.Threading.Tasks;
    using Coast.Core;
    using Dapper;

    public class SagaRepository : ISagaRepository
    {
        private IDbConnection _connection;
        private IDbTransaction _transaction;
        private readonly string _sagaTableName;
        private readonly string _sagaStepTableName;

        /// <summary>
        /// Initializes a new instance of the <see cref="SagaRepository"/> class.
        /// </summary>
        public SagaRepository(string schemaName, IDbConnection connection, IDbTransaction transaction = null)
        {
            _connection = connection;
            _transaction = transaction;
            _sagaTableName = $"\"{schemaName}\".\"Saga\"";
            _sagaStepTableName = $"\"{schemaName}\".\"SagaStep\"";
        }

        /// <inheritdoc/>
        public async Task SaveSagaAsync(Saga saga, CancellationToken cancellationToken = default)
        {
            cancellationToken.ThrowIfCancellationRequested();

            string InsertSagaSql =
$@"INSERT INTO {_sagaTableName} 
(""Id"", ""State"", ""CreationTime"", ""CurrentExecutionSequenceNumber"") 
VALUES (@Id, @State, @CreationTime, @CurrentExecutionSequenceNumber ); ";
            string InsertSagaStepSql =
$@"INSERT INTO {_sagaStepTableName}
(""Id"", ""CorrelationId"", ""EventName"", ""HasCompensation"", ""State"", ""RequestBody"", ""CreationTime"", ""FailedReason"", ""ExecutionSequenceNumber"") 
VALUES (@Id, @CorrelationId, @EventName, @HasCompensation, @State, @RequestBody, @CreationTime,@FailedReason, @ExecutionSequenceNumber); ";

            await _connection.ExecuteAsync(
                    InsertSagaSql,
                    new { Id = saga.Id, State = SagaStateEnum.Started, CreationTime = DateTime.UtcNow, CurrentExecutionSequenceNumber = saga.CurrentExecutionSequenceNumber },
                    transaction: _transaction).ConfigureAwait(false);

            foreach (var step in saga.SagaSteps)
            {
                await _connection.ExecuteAsync(
                    InsertSagaStepSql,
                    new
                    {
                        Id = step.Id,
                        CorrelationId = saga.Id,
                        EventName = step.EventName,
                        State = step.State,
                        RequestBody = step.RequestBody,
                        CreationTime = DateTime.UtcNow,
                        HasCompensation = step.HasCompensation,
后面还有 97 行代码,解锁后查看完整代码

24 小时内免费解锁 3 个项目,之后 1 积分/个。 规则说明

AI 解读

登录后可用,每次 10 积分,解读结果公开显示在下面。

还没有人解读过这个文件。