# SQL Persistence Saga Finding Logic ## Code walk-through When the default Saga message mappings do not satisfy the requirements, custom logic can be put in place to allow NServiceBus to find a saga data instance based on logic that best suits the environment. ## Prerequisites ### MS SQL Server Ensure an instance of SQL Server (Version 2016 or above for custom saga finders sample, or Version 2012 or above for other samples) is installed and accessible on `localhost` and port `1433`. A Docker image can be used to accomplish this by running `docker run --name mssql -e "ACCEPT_EULA=Y" -e "MSSQL_SA_PASSWORD=yourStrong(!)Password" -p 1433:1433 -d mcr.microsoft.com/mssql/server:latest` in a terminal. Alternatively, change the connection string to point to different SQL Server instance. At startup each endpoint will create its required SQL assets including databases, tables, and schemas. ### MySQL Ensure an instance of MySQL (Version 5.7 or above) is installed and accessible on `localhost` and port `3306`. A Docker image can be used to accomplish this by running `docker run --name mysql -e 'MYSQL_ROOT_PASSWORD=yourStrong(!)Password' -e 'MYSQL_DATABASE=sqlpersistencesample' -p 3306:3306 -d mysql:latest` in a terminal. Alternatively, change the connection string to point to different MySQL instance. At startup each endpoint will create the required SQL assets including databases, tables, and schemas. ### PostgreSQL Ensure an instance of PostgreSQL (Version 10 or later) is installed and accessible on `localhost` and port `5432`. A Docker image can be used to accomplish this by running `docker run --name postgres -e 'POSTGRES_PASSWORD=yourStrong(!)Password' -e 'POSTGRES_DB=NsbSamplesSqlSagaFinder' -p 5432:5432 -d postgres:latest` in a terminal. Alternatively, change the connection string to point to different PostgreSQL instance. At startup each endpoint will create the required SQL assets including databases, tables, and schemas. ## Persistence Config Configure the endpoint to use SQL Persistence. ### MS SQL Server ```cs //for local instance or SqlExpress //var connectionString = @"Data Source=(localdb)\mssqllocaldb;Database=NsbSamplesSqlSagaFinder;Trusted_Connection=True;MultipleActiveResultSets=true"; var connectionString = @"Server=localhost,1433;Initial Catalog=NsbSamplesSqlSagaFinder;User Id=SA;Password=yourStrong(!)Password;Encrypt=false"; var persistence = endpointConfiguration.UsePersistence(); persistence.SqlDialect(); persistence.ConnectionBuilder( connectionBuilder: () => { return new SqlConnection(connectionString); }); var subscriptions = persistence.SubscriptionSettings(); subscriptions.CacheFor(TimeSpan.FromMinutes(1)); ``` ### MySql ```cs var persistence = endpointConfiguration.UsePersistence(); var password = Environment.GetEnvironmentVariable("MySqlPassword"); if (string.IsNullOrWhiteSpace(password)) { throw new Exception("Could not extract 'MySqlPassword' from Environment variables."); } var username = Environment.GetEnvironmentVariable("MySqlUserName"); if (string.IsNullOrWhiteSpace(username)) { throw new Exception("Could not extract 'MySqlUserName' from Environment variables."); } var connection = $"server=localhost;user={username};database=sqlpersistencesample;port=3306;password={password};AllowUserVariables=True;AutoEnlist=false"; persistence.SqlDialect(); persistence.ConnectionBuilder( connectionBuilder: () => { return new MySqlConnection(connection); }); var subscriptions = persistence.SubscriptionSettings(); subscriptions.CacheFor(TimeSpan.FromMinutes(1)); ``` ### PostgreSql ```cs var persistence = endpointConfiguration.UsePersistence(); var password = Environment.GetEnvironmentVariable("PostgreSqlPassword"); if (string.IsNullOrWhiteSpace(password)) { throw new Exception("Could not extract 'PostgreSqlPassword' from Environment variables."); } var username = Environment.GetEnvironmentVariable("PostgreSqlUserName"); if (string.IsNullOrWhiteSpace(username)) { throw new Exception("Could not extract 'PostgreSqlUserName' from Environment variables."); } var connection = $"Host=localhost;Username={username};Password={password};Database=NsbSamplesSqlSagaFinder"; persistence.TablePrefix("Finder"); var dialect = persistence.SqlDialect(); dialect.JsonBParameterModifier( modifier: parameter => { var npgsqlParameter = (NpgsqlParameter)parameter; npgsqlParameter.NpgsqlDbType = NpgsqlDbType.Jsonb; }); persistence.ConnectionBuilder( connectionBuilder: () => { return new NpgsqlConnection(connection); }); var subscriptions = persistence.SubscriptionSettings(); subscriptions.CacheFor(TimeSpan.FromMinutes(1)); ``` ## The Saga This sample contains a very simple order management saga with these responsibilities: - Handling the creation of an order. - Offloading the payment process to a another handler. - Handling the completion of the payment process. - Completing the order. ```cs public class OrderSaga(ILogger logger) : Saga, IAmStartedByMessages, IHandleMessages, IHandleMessages { protected override void ConfigureHowToFindSaga(SagaPropertyMapper mapper) { mapper.MapSaga(saga => saga.OrderId) .ToMessage(msg => msg.OrderId) .ToMessage(msg => msg.OrderId); } public Task Handle(StartOrder message, IMessageHandlerContext context) { Data.PaymentTransactionId = Guid.NewGuid().ToString(); logger.LogInformation("Saga with OrderId {SagaOrderId} received StartOrder with OrderId {MessageOrderId}", Data.OrderId, message.OrderId); var issuePaymentRequest = new IssuePaymentRequest { PaymentTransactionId = Data.PaymentTransactionId }; return context.SendLocal(issuePaymentRequest); } public Task Handle(CompletePaymentTransaction message, IMessageHandlerContext context) { logger.LogInformation("Transaction with Id {PaymentTransactionId} completed for order id {OrderId}", Data.PaymentTransactionId, Data.OrderId); var completeOrder = new CompleteOrder { OrderId = Data.OrderId }; return context.SendLocal(completeOrder); } public Task Handle(CompleteOrder message, IMessageHandlerContext context) { logger.LogInformation("Saga with OrderId {SagaOrderId} received CompleteOrder with OrderId {MessageOrderId}", Data.OrderId, message.OrderId); MarkAsComplete(); return Task.CompletedTask; } } ``` It is important to note that the saga is not sending the order ID to the payment processor. Instead, it is sending a payment transaction ID. This saga needs to be correlated by more than one property. For example, `OrderId` and `PaymentTransactionId`. This requires both of these properties to be treated as unique. ## Saga Finders A Saga Finder is only required for the `PaymentTransactionCompleted` message since the other messages (`StartOrder` and `CompleteOrder`) are correlated based on `OrderSagaData.OrderId`. ### MS SQL Server > [!WARNING] > On Microsoft SQL Server, the saga finder feature requires the [JSON_VALUE function](https://learn.microsoft.com/en-us/sql/t-sql/functions/json-value-transact-sql) that is only available starting with SQL Server 2016. ```cs class OrderSagaFinder : ISagaFinder { public Task FindBy(CompletePaymentTransaction message, ISynchronizedStorageSession storageSession, IReadOnlyContextBag context, CancellationToken cancellationToken = default) { return storageSession.GetSagaData( context: context, whereClause: "JSON_VALUE(Data,'$.PaymentTransactionId') = @propertyValue", appendParameters: (builder, append) => { var parameter = builder(); parameter.ParameterName = "propertyValue"; parameter.Value = message.PaymentTransactionId; append(parameter); }); } } ``` ### MySql ```cs class OrderSagaFinder : ISagaFinder { public Task FindBy(CompletePaymentTransaction message, ISynchronizedStorageSession storageSession, IReadOnlyContextBag context, CancellationToken cancellationToken = default) { return storageSession.GetSagaData( context: context, whereClause: "JSON_EXTRACT(Data,'$.PaymentTransactionId') = @propertyValue", appendParameters: (builder, append) => { var parameter = builder(); parameter.ParameterName = "propertyValue"; parameter.Value = message.PaymentTransactionId; append(parameter); }); } } ``` ### PostgreSql ```cs class OrderSagaFinder : ISagaFinder { public Task FindBy(CompletePaymentTransaction message, ISynchronizedStorageSession storageSession, IReadOnlyContextBag context, CancellationToken cancellationToken = default) { return storageSession.GetSagaData( context: context, whereClause: @"""Data""->>'PaymentTransactionId' = @propertyValue", appendParameters: (builder, append) => { var parameter = builder(); parameter.ParameterName = "propertyValue"; parameter.Value = message.PaymentTransactionId; append(parameter); }); } } ```