Skip to content
Merged
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
2 changes: 2 additions & 0 deletions build/tasks.ps1
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@ $dist_folder = "$root\dist"
$msbuild_verbosity = "n"

$projects = @(
"Sa.Utils.WorkQueue",

"Sa.Media",
"Sa.Media.FFmpeg",

Expand Down
3 changes: 2 additions & 1 deletion src/.gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -405,4 +405,5 @@ FodyWeavers.xsd

Sa.Media.FFmpeg/runtimes/native/
nupkgs/
.packages/
.packages/
*.lscache
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ private async Task<int> LoadAsync(PostgreSqlConfigurationOptions options)
{
try
{
using var dataSource = new PgDataSource(new(options.ConnectionString));
using var dataSource = IPgDataSource.Create(options.ConnectionString);

return await dataSource.ExecuteReader(options.SelectSql, (reader, _) =>
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@
<Import Project="../Common.NuGet.Properties.xml" />

<PropertyGroup>
<Version>0.8.1</Version>
<Version>0.9.0</Version>
<Description>add a PostgreSQL-based configuration source to IConfiguration</Description>
</PropertyGroup>

Expand Down
2 changes: 1 addition & 1 deletion src/Sa.Configuration/Sa.Configuration.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@
<Import Project="../Common.NuGet.Properties.xml" />

<PropertyGroup>
<Version>0.8.1</Version>
<Version>0.9.0</Version>
<Description>extensions for Configuration</Description>
</PropertyGroup>

Expand Down
Original file line number Diff line number Diff line change
@@ -1,10 +1,7 @@

namespace Sa.Data.PostgreSql;
namespace Sa.Data.PostgreSql;

public interface IPgDataSourceSettingsBuilder
{
void WithConnectionString(string connectionString);
void WithConnectionString(Func<IServiceProvider, string> implementationFactory);
void WithSettings(PgDataSourceSettings settings);
void WithSettings(Func<IServiceProvider, PgDataSourceSettings> implementationFactory);
}
34 changes: 31 additions & 3 deletions src/Sa.Data.PostgreSql/Configuration/PgDataSourceSettings.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,34 @@
namespace Sa.Data.PostgreSql;
using Npgsql;

public sealed class PgDataSourceSettings(string connectionString)
namespace Sa.Data.PostgreSql;

internal sealed class PgDataSourceSettings(string connectionString)
{
public string ConnectionString { get; } = connectionString;
private string? _searchPath;

public string ConnectionString => connectionString;

public string GetSearchPath()
{
if (_searchPath is not null)
return _searchPath;

try
{
var builder = new NpgsqlConnectionStringBuilder(connectionString);
return _searchPath = builder.SearchPath ?? "public";
}
catch
{
return _searchPath = "public";
}
}

public void Validate()
{
if (string.IsNullOrWhiteSpace(connectionString))
{
throw new InvalidOperationException("Connection string cannot be null or empty.");
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -6,22 +6,14 @@ namespace Sa.Data.PostgreSql.Configuration;
internal sealed class PgDataSourceSettingsBuilder(IServiceCollection services) : IPgDataSourceSettingsBuilder
{
public void WithConnectionString(string connectionString)
{
services.TryAddSingleton<PgDataSourceSettings>(new PgDataSourceSettings(connectionString));
}
=> services.TryAddSingleton<PgDataSourceSettings>(new PgDataSourceSettings(connectionString));

public void WithConnectionString(Func<IServiceProvider, string> implementationFactory)
{
services.TryAddSingleton<PgDataSourceSettings>(sp => new PgDataSourceSettings(implementationFactory(sp)));
}
=> services.TryAddSingleton<PgDataSourceSettings>(sp => new PgDataSourceSettings(implementationFactory(sp)));

public void WithSettings(Func<IServiceProvider, PgDataSourceSettings> implementationFactory)
{
services.TryAddSingleton<PgDataSourceSettings>(implementationFactory);
}
=> services.TryAddSingleton<PgDataSourceSettings>(implementationFactory);

public void WithSettings(PgDataSourceSettings settings)
{
services.TryAddSingleton<PgDataSourceSettings>(settings);
}
=> services.TryAddSingleton<PgDataSourceSettings>(settings);
}
4 changes: 2 additions & 2 deletions src/Sa.Data.PostgreSql/DbCommandExtensions.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
using System.Collections.ObjectModel;
using Npgsql;
using System.Collections.ObjectModel;
using System.Runtime.CompilerServices;
using Npgsql;

namespace Sa.Data.PostgreSql;

Expand Down
65 changes: 47 additions & 18 deletions src/Sa.Data.PostgreSql/IPgDataSource.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2,46 +2,71 @@

namespace Sa.Data.PostgreSql;

public interface IPgDataSource
public interface IPgDataSource : IDisposable, IAsyncDisposable
{
public static IPgDataSource Create(string connectionString) => new PgDataSource(new PgDataSourceSettings(connectionString));
public static IPgDataSource Create(string connectionString)
=> new PgDataSource(new PgDataSourceSettings(connectionString));

string GetSearchPath();

ValueTask<NpgsqlConnection> OpenDbConnection(CancellationToken cancellationToken);

Task<int> ExecuteNonQuery(string sql, Action<NpgsqlCommand>? initCommand, CancellationToken cancellationToken = default);
Task<int> ExecuteNonQuery(
string sql,
Action<NpgsqlCommand>? initCommand,
CancellationToken cancellationToken = default);

Task<int> ExecuteNonQuery(string sql, IReadOnlyCollection<NpgsqlParameter> parameters, CancellationToken cancellationToken = default)
=> ExecuteNonQuery(sql, cmd => FillParams(cmd, parameters), cancellationToken);
Task<int> ExecuteNonQuery(
string sql,
IReadOnlyCollection<NpgsqlParameter> parameters,
CancellationToken cancellationToken = default)
=> ExecuteNonQuery(sql, cmd => FillParams(cmd, parameters), cancellationToken);

Task<int> ExecuteNonQuery(string sql, CancellationToken cancellationToken = default)
=> ExecuteNonQuery(sql, [], cancellationToken);

Task<object?> ExecuteScalar(string sql, Action<NpgsqlCommand>? initCommand, CancellationToken cancellationToken = default);
Task<object?> ExecuteScalar(
string sql, Action<NpgsqlCommand>? initCommand, CancellationToken cancellationToken = default);

async Task<T> ExecuteScalar<T>(string sql, Action<NpgsqlCommand>? initCommand, CancellationToken cancellationToken = default)
=> ((T)(await ExecuteScalar(sql, initCommand, cancellationToken))!);
async Task<T> ExecuteScalar<T>(
string sql, Action<NpgsqlCommand>? initCommand, CancellationToken cancellationToken = default)
=> ((T)(await ExecuteScalar(sql, initCommand, cancellationToken))!);

// ExecuteReader
Task<int> ExecuteReader(string sql, Action<NpgsqlDataReader, int> read, Action<NpgsqlCommand>? initCommand, CancellationToken cancellationToken = default);

Task<int> ExecuteReader(string sql, Action<NpgsqlDataReader, int> read, IReadOnlyCollection<NpgsqlParameter> parameters, CancellationToken cancellationToken = default)
Task<int> ExecuteReader(
string sql,
Action<NpgsqlDataReader, int> read,
Action<NpgsqlCommand>? initCommand,
CancellationToken cancellationToken = default);

Task<int> ExecuteReader(
string sql,
Action<NpgsqlDataReader, int> read,
IReadOnlyCollection<NpgsqlParameter> parameters,
CancellationToken cancellationToken = default)
=> ExecuteReader(sql, read, cmd => FillParams(cmd, parameters), cancellationToken);

async Task<int> ExecuteReader(string sql, Action<NpgsqlDataReader, int> read, CancellationToken cancellationToken = default)
async Task<int> ExecuteReader(
string sql, Action<NpgsqlDataReader, int> read, CancellationToken cancellationToken = default)
=> await ExecuteReader(sql, read, [], cancellationToken);


// ExecuteReaderList


async Task<List<T>> ExecuteReaderList<T>(string sql, Func<NpgsqlDataReader, T> read, CancellationToken cancellationToken = default)
async Task<List<T>> ExecuteReaderList<T>(
string sql, Func<NpgsqlDataReader, T> read, CancellationToken cancellationToken = default)
{
List<T> list = [];
await ExecuteReader(sql, (reader, _) => list.Add(read(reader)), cancellationToken);
return list;
}

async Task<List<T>> ExecuteReaderList<T>(string sql, Func<NpgsqlDataReader, T> read, IReadOnlyCollection<NpgsqlParameter> parameters, CancellationToken cancellationToken = default)
async Task<List<T>> ExecuteReaderList<T>(
string sql,
Func<NpgsqlDataReader, T> read,
IReadOnlyCollection<NpgsqlParameter> parameters,
CancellationToken cancellationToken = default)
{
List<T> list = [];
await ExecuteReader(sql, (reader, _) => list.Add(read(reader)), parameters, cancellationToken);
Expand All @@ -57,7 +82,10 @@ Task<T> ExecuteReaderFirst<T>(string sql, CancellationToken cancellationToken =
return ExecuteReaderFirst<T>(sql, [], cancellationToken);
}

async Task<T> ExecuteReaderFirst<T>(string sql, IReadOnlyCollection<NpgsqlParameter> parameters, CancellationToken cancellationToken = default)
async Task<T> ExecuteReaderFirst<T>(
string sql,
IReadOnlyCollection<NpgsqlParameter> parameters,
CancellationToken cancellationToken = default)
{
T value = default!;

Expand Down Expand Up @@ -85,9 +113,10 @@ await ExecuteReader(sql, (reader, _) =>


// BeginBinaryImport


ValueTask<ulong> BeginBinaryImport(string sql, Func<NpgsqlBinaryImporter, CancellationToken, Task<ulong>> write, CancellationToken cancellationToken = default);
ValueTask<ulong> BeginBinaryImport(
string sql,
Func<NpgsqlBinaryImporter, CancellationToken, Task<ulong>> write,
CancellationToken cancellationToken = default);

void FillParams(NpgsqlCommand cmd, IReadOnlyCollection<NpgsqlParameter> parameters)
{
Expand Down
5 changes: 4 additions & 1 deletion src/Sa.Data.PostgreSql/IPgDistributedLock.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2,5 +2,8 @@

public interface IPgDistributedLock
{
Task<bool> TryExecuteInDistributedLock(long lockId, Func<CancellationToken, Task> exclusiveLockTask, CancellationToken cancellationToken);
Task<bool> TryExecuteInDistributedLock(
long lockId,
Func<CancellationToken, Task> exclusiveLockTask,
CancellationToken cancellationToken);
}
31 changes: 24 additions & 7 deletions src/Sa.Data.PostgreSql/PgDataSource.cs
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,15 @@ namespace Sa.Data.PostgreSql;
/// NpgsqlDataSource lite
/// </summary>
/// <param name="settings">connection string</param>
public sealed class PgDataSource(PgDataSourceSettings settings) : IPgDataSource, IDisposable, IAsyncDisposable
internal sealed class PgDataSource(PgDataSourceSettings settings) : IPgDataSource
{
private readonly Lazy<NpgsqlDataSource> _dataSource = new(() => NpgsqlDataSource.Create(settings.ConnectionString));
private readonly Lazy<NpgsqlDataSource> _dataSource
= new(() => NpgsqlDataSource.Create(settings.ConnectionString));

public ValueTask<NpgsqlConnection> OpenDbConnection(CancellationToken cancellationToken) => _dataSource.Value.OpenConnectionAsync(cancellationToken);
public string GetSearchPath() => settings.GetSearchPath();

public ValueTask<NpgsqlConnection> OpenDbConnection(CancellationToken cancellationToken)
=> _dataSource.Value.OpenConnectionAsync(cancellationToken);

public void Dispose()
{
Expand All @@ -28,31 +32,44 @@ public async ValueTask DisposeAsync()
}
}

public async ValueTask<ulong> BeginBinaryImport(string sql, Func<NpgsqlBinaryImporter, CancellationToken, Task<ulong>> write, CancellationToken cancellationToken = default)
public async ValueTask<ulong> BeginBinaryImport(
string sql,
Func<NpgsqlBinaryImporter, CancellationToken, Task<ulong>> write,
CancellationToken cancellationToken = default)
{
using NpgsqlConnection db = await OpenDbConnection(cancellationToken);
using NpgsqlBinaryImporter writer = await db.BeginBinaryImportAsync(sql, cancellationToken);
ulong result = await write(writer, cancellationToken);
return result;
}

public async Task<int> ExecuteNonQuery(string sql, Action<NpgsqlCommand>? initCommand, CancellationToken cancellationToken = default)
public async Task<int> ExecuteNonQuery(
string sql,
Action<NpgsqlCommand>? initCommand,
CancellationToken cancellationToken = default)
{
using NpgsqlConnection connection = await OpenDbConnection(cancellationToken);
using NpgsqlCommand cmd = new(sql, connection);
initCommand?.Invoke(cmd);
return await cmd.ExecuteNonQueryAsync(cancellationToken);
}

public async Task<object?> ExecuteScalar(string sql, Action<NpgsqlCommand>? initCommand, CancellationToken cancellationToken = default)
public async Task<object?> ExecuteScalar(
string sql,
Action<NpgsqlCommand>? initCommand,
CancellationToken cancellationToken = default)
{
using NpgsqlConnection connection = await OpenDbConnection(cancellationToken);
using NpgsqlCommand cmd = new(sql, connection);
initCommand?.Invoke(cmd);
return await cmd.ExecuteScalarAsync(cancellationToken);
}

public async Task<int> ExecuteReader(string sql, Action<NpgsqlDataReader, int> read, Action<NpgsqlCommand>? initCommand, CancellationToken cancellationToken = default)
public async Task<int> ExecuteReader(
string sql,
Action<NpgsqlDataReader, int> read,
Action<NpgsqlCommand>? initCommand,
CancellationToken cancellationToken = default)
{
int rowCount = 0;

Expand Down
4 changes: 3 additions & 1 deletion src/Sa.Data.PostgreSql/PgDistributedLock.cs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,9 @@ namespace Sa.Data.PostgreSql;
/// </summary>
/// <seealso href="https://www.postgresql.org/docs/9.4/explicit-locking.html#ADVISORY-LOCKS"/>
/// <seealso href="https://ankitvijay.net/2021/02/28/distributed-lock-using-postgresql/"/>
internal sealed partial class PgDistributedLock(PgDataSourceSettings settings, ILogger<PgDistributedLock>? logger = null) : IPgDistributedLock
internal sealed partial class PgDistributedLock(
PgDataSourceSettings settings,
ILogger<PgDistributedLock>? logger = null) : IPgDistributedLock
{
private readonly ILogger<PgDistributedLock> _logger = logger ?? NullLogger<PgDistributedLock>.Instance;

Expand Down
4 changes: 3 additions & 1 deletion src/Sa.Data.PostgreSql/PgRetryStrategy.cs
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,9 @@ public static ValueTask<T> ExecuteWithRetry<T>(
fun: fun,
retryCount: retryCount,
initialDelay: initialDelay
, next: (ex, i) => next != null ? next(ex, i) : (ex is NpgsqlException exception) && exception.IsTransient
, next: (ex, i) => next != null
? next(ex, i)
: (ex is NpgsqlException exception) && exception.IsTransient
, cancellationToken: cancellationToken);
}
}
2 changes: 1 addition & 1 deletion src/Sa.Data.PostgreSql/Sa.Data.PostgreSql.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@
<Import Project="../Common.NuGet.Properties.xml" />

<PropertyGroup>
<Version>0.8.1</Version>
<Version>0.9.0</Version>
<Description>Simple client for Npqsql</Description>
</PropertyGroup>

Expand Down
21 changes: 19 additions & 2 deletions src/Sa.Data.PostgreSql/Setup.cs
Original file line number Diff line number Diff line change
@@ -1,16 +1,33 @@
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.DependencyInjection.Extensions;
using Npgsql;
using Sa.Data.PostgreSql.Configuration;

namespace Sa.Data.PostgreSql;

public static class Setup
{
public static IServiceCollection AddSaPostgreSqlDataSource(this IServiceCollection services, Action<IPgDataSourceSettingsBuilder>? configure = null)
public static IServiceCollection AddSaPostgreSqlDataSource(
this IServiceCollection services,
Action<IPgDataSourceSettingsBuilder>? configure = null)
{
PgDataSourceSettingsBuilder builder = new(services);
configure?.Invoke(builder);
services.TryAddSingleton<IPgDataSource, PgDataSource>();
services.TryAddSingleton<IPgDataSource>(sp =>
{
PgDataSourceSettings? settings = sp.GetService<PgDataSourceSettings>();

if (settings is null)
{
var connection = sp.GetService<NpgsqlDataSource>()?.ConnectionString
?? throw new InvalidOperationException("Empty connection string");
settings = new(connection);
}

settings.Validate();

return new PgDataSource(settings);
});
services.TryAddSingleton<IPgDistributedLock, PgDistributedLock>();
return services;
}
Expand Down
Loading
Loading