2017-02-08 15:52:24 +01:00
|
|
|
|
using System;
|
|
|
|
|
using System.Collections.Generic;
|
|
|
|
|
using System.Data;
|
|
|
|
|
using System.Data.SqlClient;
|
|
|
|
|
using System.Threading.Tasks;
|
2017-07-27 15:55:37 +02:00
|
|
|
|
using Newtonsoft.Json;
|
2017-02-08 15:52:24 +01:00
|
|
|
|
|
|
|
|
|
namespace Tapeti.Flow.SQL
|
|
|
|
|
{
|
2019-08-15 14:55:15 +02:00
|
|
|
|
/// <inheritdoc />
|
2019-08-14 12:20:53 +02:00
|
|
|
|
/// <summary>
|
|
|
|
|
/// IFlowRepository implementation for SQL server.
|
|
|
|
|
/// </summary>
|
|
|
|
|
/// <remarks>
|
|
|
|
|
/// Assumes the following table layout (table name configurable and may include schema):
|
|
|
|
|
/// create table Flow
|
|
|
|
|
/// (
|
|
|
|
|
/// FlowID uniqueidentifier not null,
|
|
|
|
|
/// CreationTime datetime2(3) not null,
|
|
|
|
|
/// StateJson nvarchar(max) null,
|
|
|
|
|
/// constraint PK_Flow primary key clustered(FlowID)
|
|
|
|
|
/// );
|
|
|
|
|
/// </remarks>
|
2017-08-14 13:58:01 +02:00
|
|
|
|
public class SqlConnectionFlowRepository : IFlowRepository
|
2017-02-08 15:52:24 +01:00
|
|
|
|
{
|
|
|
|
|
private readonly string connectionString;
|
2018-12-19 21:41:19 +01:00
|
|
|
|
private readonly string tableName;
|
2017-02-08 15:52:24 +01:00
|
|
|
|
|
|
|
|
|
|
2019-08-14 12:20:53 +02:00
|
|
|
|
/// <inheritdoc />
|
2018-12-19 21:41:19 +01:00
|
|
|
|
public SqlConnectionFlowRepository(string connectionString, string tableName = "Flow")
|
2017-02-08 15:52:24 +01:00
|
|
|
|
{
|
|
|
|
|
this.connectionString = connectionString;
|
2018-12-19 21:41:19 +01:00
|
|
|
|
this.tableName = tableName;
|
2017-02-08 15:52:24 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
2019-08-14 12:20:53 +02:00
|
|
|
|
/// <inheritdoc />
|
2017-08-14 13:58:01 +02:00
|
|
|
|
public async Task<List<KeyValuePair<Guid, T>>> GetStates<T>()
|
2017-02-08 15:52:24 +01:00
|
|
|
|
{
|
2019-10-10 16:26:13 +02:00
|
|
|
|
return await SqlRetryHelper.Execute(async () =>
|
2017-02-08 15:52:24 +01:00
|
|
|
|
{
|
2019-10-10 16:26:13 +02:00
|
|
|
|
using (var connection = await GetConnection())
|
2017-02-08 15:52:24 +01:00
|
|
|
|
{
|
2019-10-10 16:26:13 +02:00
|
|
|
|
var flowQuery = new SqlCommand($"select FlowID, StateJson from {tableName}", connection);
|
|
|
|
|
var flowReader = await flowQuery.ExecuteReaderAsync();
|
2017-07-27 15:55:37 +02:00
|
|
|
|
|
2019-10-10 16:26:13 +02:00
|
|
|
|
var result = new List<KeyValuePair<Guid, T>>();
|
2017-07-27 15:55:37 +02:00
|
|
|
|
|
2019-10-10 16:26:13 +02:00
|
|
|
|
while (await flowReader.ReadAsync())
|
|
|
|
|
{
|
|
|
|
|
var flowID = flowReader.GetGuid(0);
|
|
|
|
|
var stateJson = flowReader.GetString(1);
|
2017-02-08 15:52:24 +01:00
|
|
|
|
|
2019-10-10 16:26:13 +02:00
|
|
|
|
var state = JsonConvert.DeserializeObject<T>(stateJson);
|
|
|
|
|
result.Add(new KeyValuePair<Guid, T>(flowID, state));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return result;
|
|
|
|
|
}
|
|
|
|
|
});
|
2017-02-08 15:52:24 +01:00
|
|
|
|
}
|
|
|
|
|
|
2019-08-14 12:20:53 +02:00
|
|
|
|
/// <inheritdoc />
|
2018-12-19 21:41:19 +01:00
|
|
|
|
public async Task CreateState<T>(Guid flowID, T state, DateTime timestamp)
|
2017-02-08 15:52:24 +01:00
|
|
|
|
{
|
2019-10-10 16:26:13 +02:00
|
|
|
|
await SqlRetryHelper.Execute(async () =>
|
2018-12-19 21:41:19 +01:00
|
|
|
|
{
|
2019-10-10 16:26:13 +02:00
|
|
|
|
using (var connection = await GetConnection())
|
|
|
|
|
{
|
|
|
|
|
var query = new SqlCommand($"insert into {tableName} (FlowID, StateJson, CreationTime)" +
|
|
|
|
|
"values (@FlowID, @StateJson, @CreationTime)",
|
|
|
|
|
connection);
|
2017-07-27 15:55:37 +02:00
|
|
|
|
|
2019-10-10 16:26:13 +02:00
|
|
|
|
var flowIDParam = query.Parameters.Add("@FlowID", SqlDbType.UniqueIdentifier);
|
|
|
|
|
var stateJsonParam = query.Parameters.Add("@StateJson", SqlDbType.NVarChar);
|
|
|
|
|
var creationTimeParam = query.Parameters.Add("@CreationTime", SqlDbType.DateTime2);
|
2018-12-19 21:41:19 +01:00
|
|
|
|
|
2019-10-10 16:26:13 +02:00
|
|
|
|
flowIDParam.Value = flowID;
|
|
|
|
|
stateJsonParam.Value = JsonConvert.SerializeObject(state);
|
|
|
|
|
creationTimeParam.Value = timestamp;
|
2018-12-19 21:41:19 +01:00
|
|
|
|
|
2019-10-10 16:26:13 +02:00
|
|
|
|
await query.ExecuteNonQueryAsync();
|
|
|
|
|
}
|
|
|
|
|
});
|
2017-02-08 15:52:24 +01:00
|
|
|
|
}
|
|
|
|
|
|
2019-08-14 12:20:53 +02:00
|
|
|
|
/// <inheritdoc />
|
2018-12-19 21:41:19 +01:00
|
|
|
|
public async Task UpdateState<T>(Guid flowID, T state)
|
2017-02-08 15:52:24 +01:00
|
|
|
|
{
|
2019-10-10 16:26:13 +02:00
|
|
|
|
await SqlRetryHelper.Execute(async () =>
|
2018-12-19 21:41:19 +01:00
|
|
|
|
{
|
2019-10-10 16:26:13 +02:00
|
|
|
|
using (var connection = await GetConnection())
|
|
|
|
|
{
|
|
|
|
|
var query = new SqlCommand($"update {tableName} set StateJson = @StateJson where FlowID = @FlowID", connection);
|
2018-12-19 21:41:19 +01:00
|
|
|
|
|
2019-10-10 16:26:13 +02:00
|
|
|
|
var flowIDParam = query.Parameters.Add("@FlowID", SqlDbType.UniqueIdentifier);
|
|
|
|
|
var stateJsonParam = query.Parameters.Add("@StateJson", SqlDbType.NVarChar);
|
2018-12-19 21:41:19 +01:00
|
|
|
|
|
2019-10-10 16:26:13 +02:00
|
|
|
|
flowIDParam.Value = flowID;
|
|
|
|
|
stateJsonParam.Value = JsonConvert.SerializeObject(state);
|
2018-12-19 21:41:19 +01:00
|
|
|
|
|
2019-10-10 16:26:13 +02:00
|
|
|
|
await query.ExecuteNonQueryAsync();
|
|
|
|
|
}
|
|
|
|
|
});
|
2017-02-08 15:52:24 +01:00
|
|
|
|
}
|
|
|
|
|
|
2019-08-14 12:20:53 +02:00
|
|
|
|
/// <inheritdoc />
|
2018-12-19 21:41:19 +01:00
|
|
|
|
public async Task DeleteState(Guid flowID)
|
2017-02-08 15:52:24 +01:00
|
|
|
|
{
|
2019-10-10 16:26:13 +02:00
|
|
|
|
await SqlRetryHelper.Execute(async () =>
|
2018-12-19 21:41:19 +01:00
|
|
|
|
{
|
2019-10-10 16:26:13 +02:00
|
|
|
|
using (var connection = await GetConnection())
|
|
|
|
|
{
|
|
|
|
|
var query = new SqlCommand($"delete from {tableName} where FlowID = @FlowID", connection);
|
2018-12-19 21:41:19 +01:00
|
|
|
|
|
2019-10-10 16:26:13 +02:00
|
|
|
|
var flowIDParam = query.Parameters.Add("@FlowID", SqlDbType.UniqueIdentifier);
|
|
|
|
|
flowIDParam.Value = flowID;
|
2018-12-19 21:41:19 +01:00
|
|
|
|
|
2019-10-10 16:26:13 +02:00
|
|
|
|
await query.ExecuteNonQueryAsync();
|
|
|
|
|
}
|
|
|
|
|
});
|
2017-02-08 15:52:24 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
private async Task<SqlConnection> GetConnection()
|
|
|
|
|
{
|
|
|
|
|
var connection = new SqlConnection(connectionString);
|
|
|
|
|
await connection.OpenAsync();
|
|
|
|
|
|
|
|
|
|
return connection;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|