Why the NpgsqlConnection is not properly closed or opened on time when using `TransactionScope` in F#?

Viewed 1188

The other day, I've bumped into something a little weird, I wrote the program below, drafting the bits of a SQL Wrapper made in F#:

open System

open System.Data
open System.Transactions

open FSharp.Control

open Npgsql


[<EntryPoint>]
let main _ =
    let connectionStringBuilder = NpgsqlConnectionStringBuilder()
    connectionStringBuilder.Host <- "localhost"
    connectionStringBuilder.Port <- 5432
    connectionStringBuilder.Username <- "postgres"
    connectionStringBuilder.Password <- String.Empty
    connectionStringBuilder.Database <- "event_store"

    use connection = new NpgsqlConnection(connectionStringBuilder.ToString())
    let wasClosed = connection.State = ConnectionState.Closed

    // mutable cause it's a struct
    let mutable transactionOptions = TransactionOptions()
    transactionOptions.IsolationLevel <- IsolationLevel.ReadCommitted
    use transactionScope = new TransactionScope( TransactionScopeOption.RequiresNew, TransactionScopeAsyncFlowOption.Enabled)
    async {
        for i in [0 .. 200_000] do
            if wasClosed then
                do! connection.OpenAsync() |> Async.AwaitTask
            use command = connection.CreateCommand()
            command.CommandText <- "INSERT INTO sch_event_store.that_table (id, data) VALUES (DEFAULT, @data)"
            let parameter = command.CreateParameter()
            parameter.ParameterName <- "data"
            parameter.Value <- string i
            command.Parameters.Add(parameter) |> ignore
            command.CommandTimeout <- 500
            let! rowCount = command.ExecuteNonQueryAsync() |> Async.AwaitTask
            if wasClosed then
                connection.Close()
            printfn "%A: %A" i rowCount
    }
    |> Async.RunSynchronously
    transactionScope.Complete()
    0

and I got the exception below:

Npgsql.NpgsqlOperationInProgressException: The connection is already in state 'Executing'
   at Npgsql.NpgsqlConnector.<StartUserAction>g__DoStartUserAction|187_0(<>c__DisplayClass187_0& )
   at Npgsql.NpgsqlConnector.StartUserAction(ConnectorState newState, NpgsqlCommand command)
   at Npgsql.NpgsqlConnector.StartUserAction(NpgsqlCommand command)
   at Npgsql.NpgsqlCommand.ExecuteReaderAsync(CommandBehavior behavior, Boolean async, CancellationToken cancellationToken)
   at Npgsql.NpgsqlCommand.ExecuteNonQuery(Boolean async, CancellationToken cancellationToken)

I am surprised cause everything is awaited in that program, and still it appears that the connection is like not properly closed or opened on time.

Why is that the case?

I am not sure this is due to the way I way I await via the async CE or if it's due to the Npgsql library.

[EDIT]

The problem even occurs when there is no async CE + async calls:

open System

open System.Data
open System.Transactions

open Npgsql


[<EntryPoint>]
let main _ =
    let connectionStringBuilder = NpgsqlConnectionStringBuilder()
    connectionStringBuilder.Host <- "localhost"
    connectionStringBuilder.Port <- 5432
    connectionStringBuilder.Username <- "postgres"
    connectionStringBuilder.Password <- String.Empty
    connectionStringBuilder.Database <- "event_store"

    use connection = new NpgsqlConnection(connectionStringBuilder.ToString())
    let wasClosed = connection.State = ConnectionState.Closed

    // mutable cause it's a struct
    let mutable transactionOptions = TransactionOptions()
    transactionOptions.IsolationLevel <- IsolationLevel.ReadCommitted
    use transactionScope = new TransactionScope(TransactionScopeOption.RequiresNew, transactionOptions,TransactionScopeAsyncFlowOption.Enabled)
    for i in [0 .. 200_000] do
        if wasClosed then
            connection.Open()
        use command = connection.CreateCommand()
        command.CommandText <- "INSERT INTO sch_event_store.that_table (id, data) VALUES (DEFAULT, @data)"
        let parameter = command.CreateParameter()
        parameter.ParameterName <- "data"
        parameter.Value <- string i
        command.Parameters.Add(parameter) |> ignore
        command.CommandTimeout <- 500
        let rowCount = command.ExecuteNonQuery()
        if wasClosed then
            connection.Close()
        printfn "%A: %A" i rowCount
    transactionScope.Complete()
    0
Npgsql.NpgsqlOperationInProgressException: The connection is already in state 'Executing'
   at Npgsql.NpgsqlConnector.<StartUserAction>g__DoStartUserAction|187_0(<>c__DisplayClass187_0& )
   at Npgsql.NpgsqlConnector.StartUserAction(ConnectorState newState, NpgsqlCommand command)
   at Npgsql.NpgsqlConnector.StartUserAction(NpgsqlCommand command)
   at Npgsql.NpgsqlCommand.ExecuteReaderAsync(CommandBehavior behavior, Boolean async, CancellationToken cancellationToken)
   at Npgsql.NpgsqlCommand.ExecuteNonQuery(Boolean async, CancellationToken cancellationToken)
   at Npgsql.NpgsqlCommand.ExecuteNonQuery()

I also get another exception from time to time:

System.Transactions.TransactionException: The operation is not valid for the state of the transaction.
 ---> System.TimeoutException: Transaction Timeout
   --- End of inner exception stack trace ---
   at System.Transactions.TransactionState.EnlistVolatile(InternalTransaction tx, ISinglePhaseNotification enlistmentNotification, EnlistmentOptions enlistmentOptions, Transaction atomicTransaction)
   at System.Transactions.Transaction.EnlistVolatile(ISinglePhaseNotification singlePhaseNotification, EnlistmentOptions enlistmentOptions)
   at Npgsql.NpgsqlConnection.EnlistTransaction(Transaction transaction)
   at Npgsql.NpgsqlConnection.<>c__DisplayClass32_0.<<Open>g__OpenLong|0>d.MoveNext()
--- End of stack trace from previous location where exception was thrown ---
   at Npgsql.NpgsqlConnection.Open()

The thing is that there is no problem when coding the same logic in C#:

using System;
using System.Data;
using System.Threading.Tasks;
using System.Transactions;
using Npgsql;
using IsolationLevel = System.Transactions.IsolationLevel;


namespace CSharpPlayground
{
    public static class Program
    {
        public static async Task Main()
        {
            var connectionStringBuilder = new NpgsqlConnectionStringBuilder
            {
                Host = "localhost",
                Port = 5432,
                Username = "postgres",
                Password = string.Empty,
                Database = "event_store"
            };

            using var connection = new NpgsqlConnection(connectionStringBuilder.ToString());

            var transactionOptions = new TransactionOptions
            {
                IsolationLevel = IsolationLevel.ReadCommitted
            };

            using var transactionScope = new TransactionScope(
                TransactionScopeOption.RequiresNew,
                transactionOptions,
                TransactionScopeAsyncFlowOption.Enabled);

            var wasClosed = connection.State == ConnectionState.Closed;

            for (var i = 0; i < 200_000; i++)
            {
                if (wasClosed)
                {
                    await connection.OpenAsync();
                }
                using var command = connection.CreateCommand();
                command.CommandText = "INSERT INTO sch_event_store.that_table (id, data) VALUES (DEFAULT, @data)";
                var parameter = command.CreateParameter();
                parameter.ParameterName = "data";
                parameter.Value = i.ToString();
                command.Parameters.Add(parameter);
                var rowCount = command.ExecuteNonQuery();
                Console.WriteLine($"{i}: {rowCount}");
                if (wasClosed)
                {
                    connection.Close();
                }
            }

            transactionScope.Complete();
        }
    }
}
1 Answers

The use statement behaves slightly different as in C# (please correct me if I'm wrong). You can either wrap the for loop in another method or dispose the transactionScope manually. I've tried both versions and they all work.

Related