using System.Data; using MBN_STOCK_WEBVIEW.Infrastructure; using MMoneyCoderSharp.Data; namespace MBN_STOCK_WEBVIEW.Infrastructure.Tests; public sealed class ResilientDataQueryExecutorTests { [Theory] [InlineData(DataSourceKind.Oracle, ':')] [InlineData(DataSourceKind.MariaDb, '@')] public async Task ExecuteAsync_ParameterizedSpec_BindsValuesTypesAndDbNullWithoutSyncCalls( DataSourceKind source, char marker) { var data = new DataTable(); data.Columns.Add("VALUE", typeof(int)); data.Rows.Add(1); var plan = ConnectionPlan.Success(data); var factory = new FakeConnectionFactory(plan); var executor = CreateExecutor(factory); var query = new DataQuerySpec( $"SELECT value FROM sample WHERE code = {marker}code AND optional_value = {marker}optional_value", [ new DataQueryParameter("optional_value", null), new DataQueryParameter("code", "005930", DbType.String) ]); var result = await executor.ExecuteAsync(source, "Parameterized", query); Assert.Single(result.Rows.Cast()); Assert.Equal(query.Sql, plan.ObservedCommandText); Assert.Collection( plan.ObservedParameters, parameter => { Assert.Equal("code", parameter.Name); Assert.Equal("005930", parameter.Value); Assert.Equal(DbType.String, parameter.ExplicitDbType); Assert.Equal(ParameterDirection.Input, parameter.Direction); }, parameter => { Assert.Equal("optional_value", parameter.Name); Assert.Equal(DBNull.Value, parameter.Value); Assert.Null(parameter.ExplicitDbType); Assert.Equal(ParameterDirection.Input, parameter.Direction); }); Assert.True(plan.ObservedExecuteCancellationToken.CanBeCanceled); Assert.Equal(0, plan.SynchronousOpenCount); Assert.Equal(0, plan.SynchronousExecuteCount); Assert.True(plan.Disposed); } [Theory] [InlineData(DataSourceKind.Oracle, ':')] [InlineData(DataSourceKind.MariaDb, '@')] public async Task ExecuteAsync_ParameterizedSpec_PropagatesCancellationWithoutRetryOrSyncFallback( DataSourceKind source, char marker) { var commandStarted = new TaskCompletionSource( TaskCreationOptions.RunContinuationsAsynchronously); var plan = new ConnectionPlan { ExecuteReaderAsync = async token => { commandStarted.SetResult(); await Task.Delay(Timeout.InfiniteTimeSpan, token); throw new InvalidOperationException("Unreachable."); } }; var factory = new FakeConnectionFactory(plan); var executor = CreateExecutor(factory, transient: true, retryCount: 2); var query = new DataQuerySpec( $"SELECT value FROM sample WHERE code = {marker}code", [new DataQueryParameter("code", 7, DbType.Int32)]); using var cancellation = new CancellationTokenSource(); var operation = executor.ExecuteAsync(source, "Canceled", query, cancellation.Token); await commandStarted.Task; await cancellation.CancelAsync(); await Assert.ThrowsAnyAsync(() => operation); Assert.Equal(1, factory.CreateCount); Assert.Equal(0, plan.SynchronousOpenCount); Assert.Equal(0, plan.SynchronousExecuteCount); Assert.True(plan.Disposed); } [Theory] [InlineData(DataSourceKind.Oracle, ':')] [InlineData(DataSourceKind.MariaDb, '@')] public async Task ExecuteAsync_ParameterizedCancellation_DoesNotExposeProviderMessage( DataSourceKind source, char marker) { const string sensitiveName = "cancel_secret_name"; const string sensitiveValue = "cancel-value-sensitive-sentinel"; var commandStarted = new TaskCompletionSource( TaskCreationOptions.RunContinuationsAsynchronously); var plan = new ConnectionPlan { ExecuteReaderAsync = async token => { commandStarted.SetResult(); try { await Task.Delay(Timeout.InfiniteTimeSpan, token); } catch (OperationCanceledException) { throw new OperationCanceledException( $"{sensitiveName}|{sensitiveValue}", token); } throw new InvalidOperationException("Unreachable."); } }; var factory = new FakeConnectionFactory(plan); var executor = CreateExecutor(factory, transient: true, retryCount: 2); var query = new DataQuerySpec( $"SELECT value FROM sample WHERE code = {marker}{sensitiveName}", [new DataQueryParameter(sensitiveName, sensitiveValue)]); using var cancellation = new CancellationTokenSource(); var operation = executor.ExecuteAsync(source, "Canceled", query, cancellation.Token); await commandStarted.Task; await cancellation.CancelAsync(); var exception = await Assert.ThrowsAsync(() => operation); Assert.DoesNotContain(sensitiveName, exception.Message, StringComparison.Ordinal); Assert.DoesNotContain(sensitiveValue, exception.Message, StringComparison.Ordinal); Assert.Equal(cancellation.Token, exception.CancellationToken); Assert.Null(exception.InnerException); Assert.Equal(1, factory.CreateCount); Assert.Equal(0, plan.SynchronousOpenCount); Assert.Equal(0, plan.SynchronousExecuteCount); Assert.True(plan.Disposed); } [Theory] [InlineData(DataSourceKind.Oracle, ':')] [InlineData(DataSourceKind.MariaDb, '@')] public async Task ExecuteAsync_ParameterizedProviderFailure_DoesNotExposeNameOrValue( DataSourceKind source, char marker) { const string sensitiveName = "account_secret_name"; const string sensitiveValue = "value-sensitive-sentinel"; var plan = new ConnectionPlan { ExecuteReaderAsync = _ => Task.FromException( new TestDatabaseException($"{sensitiveName}|{sensitiveValue}")) }; var factory = new FakeConnectionFactory(plan); var executor = CreateExecutor(factory); var query = new DataQuerySpec( $"SELECT value FROM sample WHERE account = {marker}{sensitiveName}", [new DataQueryParameter(sensitiveName, sensitiveValue, DbType.String)]); var exception = await Assert.ThrowsAsync( () => executor.ExecuteAsync(source, "Sensitive", query)); Assert.Null(exception.InnerException); Assert.DoesNotContain(sensitiveName, exception.Message, StringComparison.Ordinal); Assert.DoesNotContain(sensitiveValue, exception.Message, StringComparison.Ordinal); Assert.Equal(1, factory.CreateCount); Assert.True(plan.Disposed); } [Theory] [InlineData(DataSourceKind.Oracle, ':')] [InlineData(DataSourceKind.MariaDb, '@')] public async Task ExecuteAsync_ParameterizedTransientReadFailure_RecreatesBoundCommandOnRetry( DataSourceKind source, char marker) { var first = new ConnectionPlan { ExecuteReaderAsync = _ => Task.FromException( new TestDatabaseException("temporary")) }; var data = new DataTable(); data.Columns.Add("VALUE", typeof(int)); data.Rows.Add(42); var second = ConnectionPlan.Success(data); var factory = new FakeConnectionFactory(first, second); var executor = CreateExecutor(factory, transient: true, retryCount: 1); var query = new DataQuerySpec( $"SELECT value FROM sample WHERE code = {marker}code", [new DataQueryParameter("code", 42, DbType.Int32)]); var result = await executor.ExecuteAsync(source, "Retry", query); Assert.Equal(42, Assert.Single(result.Rows.Cast())["VALUE"]); Assert.Equal(2, factory.CreateCount); Assert.All( new[] { first, second }, plan => { var parameter = Assert.Single(plan.ObservedParameters); Assert.Equal("code", parameter.Name); Assert.Equal(42, parameter.Value); Assert.Equal(DbType.Int32, parameter.ExplicitDbType); Assert.Equal(0, plan.SynchronousOpenCount); Assert.Equal(0, plan.SynchronousExecuteCount); Assert.True(plan.Disposed); }); } [Fact] public async Task ExecuteAsync_ParameterizedSpec_RejectsProviderMarkerBeforeOpeningConnection() { var plan = ConnectionPlan.Success(new DataTable()); var factory = new FakeConnectionFactory(plan); var executor = CreateExecutor(factory); var query = new DataQuerySpec( "SELECT value FROM sample WHERE code = :code", [new DataQueryParameter("code", 7)]); var exception = await Assert.ThrowsAsync( () => executor.ExecuteAsync(DataSourceKind.MariaDb, "Rejected", query)); Assert.Equal(0, factory.CreateCount); Assert.DoesNotContain("code", exception.Message, StringComparison.OrdinalIgnoreCase); Assert.DoesNotContain("7", exception.Message, StringComparison.Ordinal); } [Fact] public async Task ExecuteAsync_Success_PreservesSchemaRowsAndDbNull() { var source = new DataTable(); source.Columns.Add("CODE", typeof(int)); source.Columns.Add("NAME", typeof(string)); source.Rows.Add(7, DBNull.Value); source.Rows.Add(8, "sample"); var plan = ConnectionPlan.Success(source); var factory = new FakeConnectionFactory(plan); var executor = CreateExecutor(factory); var result = await executor.ExecuteAsync( DataSourceKind.Oracle, "Stocks", "SELECT CODE, NAME FROM SAMPLE"); Assert.Equal("Stocks", result.TableName); Assert.Equal(2, result.Columns.Count); Assert.Equal("CODE", result.Columns[0].ColumnName); Assert.Equal(typeof(int), result.Columns[0].DataType); Assert.Equal("NAME", result.Columns[1].ColumnName); Assert.Equal(typeof(string), result.Columns[1].DataType); Assert.Equal(2, result.Rows.Count); Assert.Equal(7, result.Rows[0]["CODE"]); Assert.Equal(DBNull.Value, result.Rows[0]["NAME"]); Assert.Equal("sample", result.Rows[1]["NAME"]); Assert.Equal(1, factory.CreateCount); Assert.True(plan.Disposed); Assert.Equal(5, plan.ObservedCommandTimeoutSeconds); } [Fact] public async Task ExecuteAsync_CanceledBeforeOpen_DoesNotRetry() { var plan = new ConnectionPlan { OpenAsync = token => { token.ThrowIfCancellationRequested(); return Task.CompletedTask; } }; var factory = new FakeConnectionFactory(plan); var executor = CreateExecutor(factory, transient: true, retryCount: 2); using var cancellation = new CancellationTokenSource(); cancellation.Cancel(); await Assert.ThrowsAnyAsync( () => executor.ExecuteAsync( DataSourceKind.Oracle, "Canceled", "SELECT 1 FROM DUAL", cancellation.Token)); Assert.Equal(1, factory.CreateCount); Assert.True(plan.Disposed); } [Fact] public async Task ExecuteAsync_CanceledDuringCommand_DoesNotRetry() { var commandStarted = new TaskCompletionSource( TaskCreationOptions.RunContinuationsAsynchronously); var plan = new ConnectionPlan { ExecuteReaderAsync = async token => { commandStarted.SetResult(); await Task.Delay(Timeout.InfiniteTimeSpan, token); throw new InvalidOperationException("Unreachable."); } }; var factory = new FakeConnectionFactory(plan); var executor = CreateExecutor(factory, transient: true, retryCount: 2); using var cancellation = new CancellationTokenSource(); var operation = executor.ExecuteAsync( DataSourceKind.MariaDb, "Canceled", "SELECT 1", cancellation.Token); await commandStarted.Task; await cancellation.CancelAsync(); await Assert.ThrowsAnyAsync(() => operation); Assert.Equal(1, factory.CreateCount); Assert.True(plan.Disposed); } [Fact] public async Task ExecuteAsync_OverallTimeout_ThrowsSafeTimeoutAndDoesNotRetryCancellation() { var plan = new ConnectionPlan { ExecuteReaderAsync = async token => { await Task.Delay(Timeout.InfiniteTimeSpan, token); throw new InvalidOperationException("Unreachable."); } }; var factory = new FakeConnectionFactory(plan); var executor = CreateExecutor(factory, transient: true, retryCount: 2, timeoutSeconds: 1); var exception = await Assert.ThrowsAsync( () => executor.ExecuteAsync( DataSourceKind.Oracle, "Timeout", "SELECT 1 FROM DUAL")); Assert.Equal(DataSourceKind.Oracle, exception.DataSource); Assert.True(exception.IsTransient); Assert.Equal(1, factory.CreateCount); Assert.True(plan.Disposed); } [Fact] public async Task ExecuteAsync_TransientReadFailure_RetriesWithFreshConnectionAndSucceeds() { var first = new ConnectionPlan { OpenAsync = _ => Task.FromException(new TestDatabaseException("temporary")) }; var data = new DataTable(); data.Columns.Add("VALUE", typeof(int)); data.Rows.Add(1); var second = ConnectionPlan.Success(data); var factory = new FakeConnectionFactory(first, second); var executor = CreateExecutor(factory, transient: true, retryCount: 2); var result = await executor.ExecuteAsync( DataSourceKind.MariaDb, "Retry", "SELECT VALUE FROM SAMPLE"); Assert.Single(result.Rows.Cast()); Assert.Equal(2, factory.CreateCount); Assert.Equal(2, factory.CreatedConnections.Count); Assert.NotSame(factory.CreatedConnections[0], factory.CreatedConnections[1]); Assert.True(first.Disposed); Assert.True(second.Disposed); } [Fact] public async Task ExecuteAsync_PermanentFailure_AttemptsOnceAndUsesSafePublicMessage() { const string host = "host-sensitive-sentinel"; const string user = "user-sensitive-sentinel"; const string password = "password-sensitive-sentinel"; var plan = new ConnectionPlan { OpenAsync = _ => Task.FromException( new TestDatabaseException($"{host}|{user}|{password}")) }; var factory = new FakeConnectionFactory(plan); var executor = CreateExecutor(factory, transient: false, retryCount: 2); var exception = await Assert.ThrowsAsync( () => executor.ExecuteAsync( DataSourceKind.Oracle, "Permanent", "SELECT 1 FROM DUAL")); Assert.False(exception.IsTransient); Assert.Equal(1, factory.CreateCount); Assert.Null(exception.InnerException); Assert.DoesNotContain(host, exception.Message, StringComparison.Ordinal); Assert.DoesNotContain(user, exception.Message, StringComparison.Ordinal); Assert.DoesNotContain(password, exception.Message, StringComparison.Ordinal); } [Fact] public async Task ExecuteAsync_NonReadQuery_DoesNotRetryTransientFailure() { var plan = new ConnectionPlan { ExecuteReaderAsync = _ => Task.FromException( new TestDatabaseException("temporary")) }; var factory = new FakeConnectionFactory(plan); var executor = CreateExecutor(factory, transient: true, retryCount: 2); var exception = await Assert.ThrowsAsync( () => executor.ExecuteAsync( DataSourceKind.MariaDb, "Mutation", "UPDATE SAMPLE SET VALUE = 1")); Assert.True(exception.IsTransient); Assert.Equal(1, factory.CreateCount); Assert.True(plan.Disposed); } [Theory] [InlineData("SELECT value INTO OUTFILE '/tmp/sample' FROM sample")] [InlineData("SELECT 1; REPLACE INTO sample VALUES (1)")] [InlineData("SELECT 1; SET @value = 1")] [InlineData("WITH value AS (SELECT 1) DELETE FROM sample")] public async Task ExecuteAsync_AmbiguousOrStateChangingLegacyRead_DoesNotRetry( string query) { var plan = new ConnectionPlan { ExecuteReaderAsync = _ => Task.FromException( new TestDatabaseException("temporary")) }; var factory = new FakeConnectionFactory(plan); var executor = CreateExecutor(factory, transient: true, retryCount: 2); var exception = await Assert.ThrowsAsync( () => executor.ExecuteAsync(DataSourceKind.MariaDb, "Mutation", query)); Assert.True(exception.IsTransient); Assert.Equal(1, factory.CreateCount); Assert.True(plan.Disposed); } private static ResilientDataQueryExecutor CreateExecutor( FakeConnectionFactory factory, bool transient = false, int retryCount = 0, int timeoutSeconds = 5) => new( factory, new DatabaseResilienceOptions { OperationTimeoutSeconds = timeoutSeconds, MaximumRetryCount = retryCount, InitialRetryDelayMilliseconds = 0 }, new StubErrorDetector(transient)); }