Skip to content
Open
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
@@ -1,4 +1,4 @@
namespace NServiceBus.Transport.Sql.Shared
namespace NServiceBus.Transport.Sql.Shared
{
class TransportTransactionKeys
{
Expand All @@ -7,5 +7,10 @@ class TransportTransactionKeys
public const string SqlTransaction = "System.Data.SqlClient.SqlTransaction";

public const string IsUserProvidedTransaction = "SqlServer.Transaction.IsUserProvided";

// Well-known key read by downstream components (e.g. SQL persistence) to detect that they must not reuse the receive connection and transaction
public const string ReceiveOnlyTransactionMode = "SqlTransport.ReceiveOnlyTransactionMode";

public const string State = "SqlTransport.TransportTransactionState";
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,5 @@ async Task<bool> TryProcess(Message message, TransportTransaction transportTrans
IsolationLevel isolationLevel = IsolationLevelMapper.Map(transactionOptions.IsolationLevel);
FailureInfoStorage failureInfoStorage = failureInfoStorage;
readonly IExceptionClassifier exceptionClassifier = exceptionClassifier;
internal static string ReceiveOnlyTransactionMode = "SqlTransport.ReceiveOnlyTransactionMode";
}
}
95 changes: 52 additions & 43 deletions src/NServiceBus.Transport.Sql.Shared/Sending/MessageDispatcher.cs
Original file line number Diff line number Diff line change
Expand Up @@ -52,68 +52,77 @@ async Task<IEnumerable<UnicastTransportOperation>> ConvertToUnicastOperations(Tr

async Task DispatchIsolated(IEnumerable<UnicastTransportOperation> operations, TransportTransaction transportTransaction, CancellationToken cancellationToken)
{
if (transportTransaction.IsUserProvided(out DbConnection connection, out var transaction))
if (transportTransaction.GetState() == TransportTransactionState.UserProvided)
{
var (connection, transaction) = transportTransaction.GetConnectionAndTransaction();
await Dispatch(operations, connection, transaction, cancellationToken).ConfigureAwait(false);
return;
}

using (var scope = new TransactionScope(TransactionScopeOption.Suppress, TransactionScopeAsyncFlowOption.Enabled))
using (connection = await connectionFactory.OpenNewConnection(cancellationToken).ConfigureAwait(false))
using (var tx = connection.BeginTransaction())
using (var connection = await connectionFactory.OpenNewConnection(cancellationToken).ConfigureAwait(false))
using (var transaction = connection.BeginTransaction())
{
await Dispatch(operations, connection, tx, cancellationToken).ConfigureAwait(false);
tx.Commit();
await Dispatch(operations, connection, transaction, cancellationToken).ConfigureAwait(false);
transaction.Commit();
scope.Complete();
}
}

async Task DispatchDefault(IEnumerable<UnicastTransportOperation> operations, TransportTransaction transportTransaction, CancellationToken cancellationToken)
{
DbConnection connection;
var state = transportTransaction.GetState();

if (transportTransaction.OutsideOfHandler())
switch (state)
{
using (connection = await connectionFactory.OpenNewConnection(cancellationToken).ConfigureAwait(false))
{
using (var transaction = await connection.BeginTransactionAsync(cancellationToken).ConfigureAwait(false))
// There is no receive transaction the sends could take part in, either because dispatch
// happens outside the message processing pipeline or because the receive transaction must
// not be used for sends. Dispatch on a dedicated connection with its own transaction.
case TransportTransactionState.OutsideHandler:
case TransportTransactionState.ReceiveOnly:
{
using var connection = await connectionFactory.OpenNewConnection(cancellationToken).ConfigureAwait(false);
using var transaction = await connection.BeginTransactionAsync(cancellationToken).ConfigureAwait(false);

await Dispatch(operations, connection, transaction, cancellationToken).ConfigureAwait(false);
transaction.Commit();
break;
}
}
}
else if (transportTransaction.IsNoTransaction(out connection))
{
using (var transaction = connection.BeginTransaction())
{
await Dispatch(operations, connection, transaction, cancellationToken).ConfigureAwait(false);
transaction.Commit();
}
}
else if (transportTransaction.IsReceiveOnly())
{
using (connection = await connectionFactory.OpenNewConnection(cancellationToken).ConfigureAwait(false))
using (var transaction = connection.BeginTransaction())
{
await Dispatch(operations, connection, transaction, cancellationToken).ConfigureAwait(false);
transaction.Commit();
}
}
else if (transportTransaction.IsSendsAtomicWithReceive(out connection, out var transaction))
{
await Dispatch(operations, connection, transaction, cancellationToken).ConfigureAwait(false);
}
else if (transportTransaction.IsTransactionScope())
{
using (connection = await connectionFactory.OpenNewConnection(cancellationToken).ConfigureAwait(false))
{
await Dispatch(operations, connection, null, cancellationToken).ConfigureAwait(false);
}
}
else
{
throw new Exception("TransportTransaction is in invalid state.");

// The receive connection can be reused but there is no receive transaction, so the sends
// get their own short-lived transaction.
case TransportTransactionState.NoTransaction:
{
var (connection, _) = transportTransaction.GetConnectionAndTransaction();
using var transaction = connection.BeginTransaction();

await Dispatch(operations, connection, transaction, cancellationToken).ConfigureAwait(false);
transaction.Commit();
break;
}

// The sends take part in the receive transaction or in the transaction provided by the user.
case TransportTransactionState.SendsAtomicWithReceive:
case TransportTransactionState.UserProvided:
{
var (connection, transaction) = transportTransaction.GetConnectionAndTransaction();

await Dispatch(operations, connection, transaction, cancellationToken).ConfigureAwait(false);
break;
}

// The ambient transaction covers both the receive and the sends; a new connection enlists
// in it automatically.
case TransportTransactionState.TransactionScope:
{
using var connection = await connectionFactory.OpenNewConnection(cancellationToken).ConfigureAwait(false);

await Dispatch(operations, connection, null, cancellationToken).ConfigureAwait(false);
break;
}

default:
throw new Exception($"Unsupported transport transaction state: {state}.");
}
}

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
namespace NServiceBus.Transport.Sql.Shared;

/// <summary>
/// Describes the context in which a <see cref="Transport.TransportTransaction"/> was created, allowing the
/// dispatcher to determine how outgoing messages relate to the receive transaction without inferring it
/// from which entries happen to be present in the transaction.
/// </summary>
enum TransportTransactionState
{
/// <summary>Dispatch happens outside the context of an incoming message, e.g. from a send-only endpoint.</summary>
OutsideHandler,

/// <summary>The incoming message was received without a transaction. The receive connection can be reused but sends need their own transaction.</summary>
NoTransaction,

/// <summary>Sends must not take part in the receive transaction. Outgoing messages get a dedicated connection and transaction.</summary>
ReceiveOnly,

/// <summary>Outgoing messages take part in the receive connection and transaction.</summary>
SendsAtomicWithReceive,

/// <summary>An ambient transaction is active. New connections enlist in it automatically.</summary>
TransactionScope,

/// <summary>The user supplied their own connection or transaction through the send or publish options.</summary>
UserProvided
}
134 changes: 69 additions & 65 deletions src/NServiceBus.Transport.Sql.Shared/Sending/TransportTransactions.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
namespace NServiceBus.Transport.Sql.Shared;
namespace NServiceBus.Transport.Sql.Shared;

using System;
using System.Data.Common;
Expand All @@ -9,124 +9,128 @@ static class TransportTransactions
public static TransportTransaction NoTransaction(DbConnection connection)
{
var transportTransaction = new TransportTransaction();
transportTransaction.Set(TransportTransactionKeys.SqlConnection, connection);
return transportTransaction;
}

public static bool IsNoTransaction(this TransportTransaction transportTransaction, out DbConnection connection)
{
transportTransaction.TryGet(TransportTransactionKeys.SqlTransaction, out DbTransaction nativeTransaction);
transportTransaction.TryGet(out Transaction ambientTransaction);

transportTransaction.TryGet(TransportTransactionKeys.SqlConnection, out connection);
transportTransaction.Set(TransportTransactionKeys.State, TransportTransactionState.NoTransaction);
transportTransaction.Set(TransportTransactionKeys.SqlConnection, connection);

return nativeTransaction == null && ambientTransaction == null;
return transportTransaction;
}

public static TransportTransaction ReceiveOnly(DbConnection connection, DbTransaction transaction)
{
var transportTransaction = new TransportTransaction();

transportTransaction.Set(TransportTransactionKeys.State, TransportTransactionState.ReceiveOnly);
transportTransaction.Set(TransportTransactionKeys.SqlConnection, connection);
transportTransaction.Set(TransportTransactionKeys.SqlTransaction, transaction);

//this indicates to MessageDispatcher that it should not reuse connection or transaction for sends
transportTransaction.Set(ReceiveOnlyTransactionMode, true);
//downstream components (e.g. SQL persistence) use this well-known entry to detect that they must not reuse the receive connection and transaction
transportTransaction.Set(TransportTransactionKeys.ReceiveOnlyTransactionMode, true);

return transportTransaction;
}

public static bool IsReceiveOnly(this TransportTransaction transportTransaction) => transportTransaction.TryGet(ProcessWithNativeTransaction.ReceiveOnlyTransactionMode, out bool _);

public static TransportTransaction SendsAtomicWithReceive(DbConnection connection, DbTransaction transaction)
{
var transportTransaction = new TransportTransaction();

transportTransaction.Set(TransportTransactionKeys.State, TransportTransactionState.SendsAtomicWithReceive);
transportTransaction.Set(TransportTransactionKeys.SqlConnection, connection);
transportTransaction.Set(TransportTransactionKeys.SqlTransaction, transaction);

return transportTransaction;
}

public static bool IsSendsAtomicWithReceive(this TransportTransaction transportTransaction, out DbConnection connection, out DbTransaction transaction)
{
transportTransaction.TryGet(TransportTransactionKeys.SqlTransaction, out transaction);
transportTransaction.TryGet(TransportTransactionKeys.SqlConnection, out connection);
transportTransaction.TryGet(ProcessWithNativeTransaction.ReceiveOnlyTransactionMode, out bool receiveOnly);

return transaction != null && connection != null && !receiveOnly;
}

public static TransportTransaction TransactionScope(Transaction transaction)
{
var transportTransaction = new TransportTransaction();

transportTransaction.Set(TransportTransactionKeys.State, TransportTransactionState.TransactionScope);
transportTransaction.Set(transaction);

return transportTransaction;
}

public static bool IsTransactionScope(this TransportTransaction transportTransaction)
public static TransportTransaction UserProvided(DbConnection connection)
{
transportTransaction.TryGet(out Transaction ambientTransaction);
return ambientTransaction != null;
}
var transportTransaction = new TransportTransaction();

public static bool OutsideOfHandler(this TransportTransaction transportTransaction)
{
transportTransaction.TryGet(TransportTransactionKeys.SqlTransaction, out DbTransaction nativeTransaction);
transportTransaction.TryGet(TransportTransactionKeys.SqlConnection, out DbConnection nativeConnection);
transportTransaction.TryGet(out Transaction ambientTransaction);
transportTransaction.Set(TransportTransactionKeys.State, TransportTransactionState.UserProvided);
transportTransaction.Set(TransportTransactionKeys.IsUserProvidedTransaction, true);
transportTransaction.Set(TransportTransactionKeys.SqlConnection, connection);

return nativeTransaction == null && nativeConnection == null && ambientTransaction == null;
return transportTransaction;
}

public static TransportTransaction UserProvided(DbConnection connection)
public static TransportTransaction UserProvided(DbTransaction transaction)
{
var result = new TransportTransaction();
var transportTransaction = new TransportTransaction();

result.Set(TransportTransactionKeys.IsUserProvidedTransaction, true);
result.Set(TransportTransactionKeys.SqlConnection, connection);
transportTransaction.Set(TransportTransactionKeys.State, TransportTransactionState.UserProvided);
transportTransaction.Set(TransportTransactionKeys.IsUserProvidedTransaction, true);
transportTransaction.Set(TransportTransactionKeys.SqlTransaction, transaction);

return result;
return transportTransaction;
}

public static TransportTransaction UserProvided(DbTransaction transaction)
public static TransportTransactionState GetState(this TransportTransaction transportTransaction) =>
transportTransaction.TryGet(TransportTransactionKeys.State, out TransportTransactionState state)
? state
: InferState(transportTransaction);

/// <summary>
/// Returns the connection to dispatch on and, if present, the transaction the sends should take part in.
/// Falls back to the transaction's connection when only a transaction was provided.
/// </summary>
public static (DbConnection connection, DbTransaction transaction) GetConnectionAndTransaction(this TransportTransaction transportTransaction)
{
var result = new TransportTransaction();
transportTransaction.TryGet(TransportTransactionKeys.SqlTransaction, out DbTransaction transaction);
transportTransaction.TryGet(TransportTransactionKeys.SqlConnection, out DbConnection connection);

result.Set(TransportTransactionKeys.IsUserProvidedTransaction, true);
result.Set(TransportTransactionKeys.SqlTransaction, transaction);
connection ??= transaction?.Connection
?? throw new Exception($"Invalid {nameof(TransportTransaction)} state. It contains no SqlTransaction or SqlConnection objects.");

return result;
return (connection, transaction);
}

public static bool IsUserProvided(this TransportTransaction transportTransaction, out DbConnection connection, out DbTransaction transaction)
// TransportTransaction instances that were not created by this transport carry no explicit state: the core
// creates an empty one for dispatches outside the message processing pipeline, and external integrations
// hand-roll instances containing a connection and/or transaction. For those the state is derived from the
// entries present in the transaction.
static TransportTransactionState InferState(TransportTransaction transportTransaction)
{
connection = null;
transaction = null;

transportTransaction.TryGet(TransportTransactionKeys.IsUserProvidedTransaction, out bool isUserProvided);
transportTransaction.TryGet(TransportTransactionKeys.SqlTransaction, out DbTransaction nativeTransaction);
transportTransaction.TryGet(TransportTransactionKeys.SqlConnection, out DbConnection connection);
transportTransaction.TryGet(out Transaction ambientTransaction);

if (isUserProvided)
{
transportTransaction.TryGet(TransportTransactionKeys.SqlTransaction, out transaction);

if (transaction != null)
{
connection = transaction.Connection;
}
else if (transportTransaction.TryGet(TransportTransactionKeys.SqlConnection, out connection))
{
transaction = null;
}
else
{
throw new Exception($"Invalid {nameof(TransportTransaction)} state. Transaction provided by the user but contains no SqlTransaction or SqlConnection objects.");
}
return TransportTransactionState.UserProvided;
}

if (nativeTransaction == null && ambientTransaction == null)
{
return connection == null
? TransportTransactionState.OutsideHandler
: TransportTransactionState.NoTransaction;
}

if (transportTransaction.TryGet(TransportTransactionKeys.ReceiveOnlyTransactionMode, out bool _))
{
return TransportTransactionState.ReceiveOnly;
}

if (nativeTransaction != null && connection != null)
{
return TransportTransactionState.SendsAtomicWithReceive;
}

if (ambientTransaction != null)
{
return TransportTransactionState.TransactionScope;
}

return isUserProvided;
throw new Exception($"{nameof(TransportTransaction)} is in invalid state.");
}
internal static string ReceiveOnlyTransactionMode = "SqlTransport.ReceiveOnlyTransactionMode";
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,5 @@ class SettingsKeys
// For backward compatibility reasons these settings keys are hard coded to the System.Data types to enable connection and transaction sharing with SQL persistence
public const string TransportTransactionSqlConnectionKey = "System.Data.SqlClient.SqlConnection";
public const string TransportTransactionSqlTransactionKey = "System.Data.SqlClient.SqlTransaction";

public const string IsUserProvidedTransactionKey = "SqlServer.Transaction.IsUserProvided";
}
}
Loading
Loading