using System.Collections; using System.Collections.Generic; using System.Data; using System.Data.Common; using System.Diagnostics; using System.Globalization; using System.Linq; using System.Reflection; using System.Runtime.CompilerServices; using System.Text; using System.Text.Json; using System.Text.RegularExpressions; using System.Threading; using System.Threading.Tasks; using Dapper; using DocumentFormat.OpenXml.Office2010.CustomUI; using GB5Shared.Connection; using GB5Shared.DBQueryConverter; using GB5Shared.DTO.Framework; using GB5Shared.DTO.Framework.CommonConfig; using GB5Shared.DTO.Framework.Login; using GB5Shared.GB5CommonFunction; using GB5Shared.GB5Exception; using GB5Shared.Telemetry; using GB5Shared.Telemetry.Database; using Microsoft.AspNetCore.Http; using Microsoft.Data.SqlClient; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using Npgsql; using static GB5Shared.GB5Constant.Constant; namespace GB5Shared.QueryExecutor { public class QueryExecutor : IQueryExecutor { private string _connectionstring = ""; private readonly IGB5CommonFunction _GB5CommonFunction; private readonly IApplicationConnection _Appconnection; private readonly ILogger _logger; private readonly IHttpContextAccessor _httpContext; private readonly IOptionsSnapshot _dataBaseConfig; private readonly TelemetryOptions? _telemetryOptions; private const int DefaultCommandTimeout = 300; public QueryExecutor( IApplicationConnection Appconnection, IGB5CommonFunction IGB5CommonFunction, ILogger logger, IHttpContextAccessor httpContextAccessor, IOptionsSnapshot dataBaseConfig, TelemetryOptions telemetryOptions = null!) { GB5DapperTypeHandlers.Register(); _Appconnection = Appconnection; _GB5CommonFunction = IGB5CommonFunction; _logger = logger; _httpContext = httpContextAccessor; _dataBaseConfig = dataBaseConfig; _telemetryOptions = telemetryOptions; } // ── DB Telemetry helpers ───────────────────────────────────────────── // Called automatically by every public method — zero per-query code needed. // Delegates to GB5Shared.Telemetry.Database.DatabaseActivityHelper, the single shared, // dynamic (no per-query/per-field hardcoding) span factory used by every GB5 data-access // path — QueryExecutor here, and any raw ADO.NET call site elsewhere in GB5Shared. // Truncation is driven entirely by TelemetryOptions.MaxDbStatementLength (0 = unlimited). private Activity? StartDbTrace(string sql, string? dbName, object? parameters = null, [CallerMemberName] string caller = "") => DatabaseActivityHelper.StartDbActivity( sql, dbName, parameters, maxStatementLength: _telemetryOptions?.MaxDbStatementLength ?? 0, caller: caller); private static void MarkDbError(Activity? activity, Exception ex) => DatabaseActivityHelper.RecordDbError(activity, ex); // ───────────────────────────────────────────────────────────────────── private static string ApplyReadUncommitted(string sql, bool useReadUncommitted, bool hasExternalTransaction) { if (!useReadUncommitted || hasExternalTransaction) return sql; return "SET TRANSACTION ISOLATION LEVEL READ UNCOMMITTED;\n" + sql; } // UPDLOCK only makes sense when the lock needs to survive past this one statement (the // caller intends to read now, then write later in the SAME logical operation without // another session sneaking in between) — which requires the read and the later write to // share one DbTransaction. A QB using UPDLOCK without passing a transaction here is a // near-certain bug: the lock is released the instant this statement completes (READ // COMMITTED default), making the hint a no-op — this exact pattern was found live in // PatternDAL.UpdatePattern (MAX(PATTERNDETAILID) WITH (UPDLOCK, HOLDLOCK) read separately // from the later insert's own transaction). Warn, don't throw — some call sites may use // UPDLOCK deliberately within a single self-contained statement where it's harmless. private static readonly Regex UpdLockHintPattern = new(@"\bUPDLOCK\b", RegexOptions.IgnoreCase | RegexOptions.Compiled); private void WarnIfUpdLockWithoutTransaction(string sql, DbTransaction? transaction, [CallerMemberName] string caller = "") { if (transaction is null && UpdLockHintPattern.IsMatch(sql)) { _logger.LogWarning( "{Caller} used an UPDLOCK table hint but no DbTransaction was supplied — the lock is released " + "as soon as this single statement completes and cannot protect a later write in the same " + "operation. Pass the SAME DbTransaction to both this locking read and the subsequent write, or " + "the hint is a no-op. SQL: {Sql}", caller, sql); } } private async Task CreateNewConnectionAsync(LoginDTO loginDTO) { var connStr = await _Appconnection.DBConnectionStringCached(loginDTO); // The tenant's OWN DbType, not the GB5 system's global default — a connection can be // a completely different RDBMS (and a completely different physical server) than // whatever this process's default happens to be. int databaseType = await _Appconnection.DatabaseTypeCached(loginDTO); if (databaseType == DBTYPE.SQL && !connStr.Contains("Enlist=", StringComparison.OrdinalIgnoreCase)) connStr += ";Enlist=false"; return databaseType switch { DBTYPE.POSTGRESQL => new NpgsqlConnection(connStr), DBTYPE.SQL => new SqlConnection(connStr), _ => throw new NotSupportedException($"Database type '{databaseType}' is not supported.") }; } // Every one of this file's ~30 connection.OpenAsync() call sites goes through here. // Two things a raw OpenAsync() never gives a caller: (1) a message that names WHICH // server/database this was for — the raw SqlException/NpgsqlException never does, and // background services iterating many tenants otherwise can't tell which one failed // without re-deriving it from context; (2) a bounded wait — the connection string's own // Connection Timeout (120s by default, see ApplicationConnection's pool config) is sized // for "server is up but briefly slow," not for "server is genuinely unreachable," so a // silently-dropped connection (firewalled, wrong port, no route) would otherwise hang the // caller for the full 120s with zero feedback before finally throwing a generic error. private static async Task OpenWithClearErrorAsync(DbConnection connection, CancellationToken ct = default) { using var timeoutCts = CancellationTokenSource.CreateLinkedTokenSource(ct); timeoutCts.CancelAfter(TimeSpan.FromSeconds(UnreachableServerTimeoutSeconds)); try { await connection.OpenAsync(timeoutCts.Token).ConfigureAwait(false); } catch (Exception ex) when (!ct.IsCancellationRequested && IsUnreachableServerError(ex, timeoutCts.Token)) { throw new DatabaseUnreachableException( $"Cannot reach the database server for '{SafeConnectionTarget(connection.ConnectionString)}' " + $"within {UnreachableServerTimeoutSeconds}s — {ex.Message}", ex); } } private const int UnreachableServerTimeoutSeconds = 15; private static bool IsUnreachableServerError(Exception ex, CancellationToken timeoutToken) => ex is System.Net.Sockets.SocketException || (timeoutToken.IsCancellationRequested && ex is OperationCanceledException) || (ex is SqlException sqlEx && sqlEx.Number is -2 or 2 or 53 or 10060 or 10061 or 64) || (ex is Npgsql.NpgsqlException && ex.InnerException is System.Net.Sockets.SocketException) || (ex.InnerException != null && IsUnreachableServerError(ex.InnerException, timeoutToken)); // Server/host + database only — never the password, never the full connection string. private static string SafeConnectionTarget(string connectionString) { try { string? Extract(string key) { foreach (var part in connectionString.Split(';', StringSplitOptions.TrimEntries | StringSplitOptions.RemoveEmptyEntries)) { var kv = part.Split('=', 2); if (kv.Length == 2 && string.Equals(kv[0], key, StringComparison.OrdinalIgnoreCase)) return kv[1]; } return null; } string? server = Extract("Data Source") ?? Extract("Host") ?? Extract("Server"); string? database = Extract("Initial Catalog") ?? Extract("Database"); return $"{server ?? "?"}/{database ?? "?"}"; } catch { return "(unknown target)"; } } // ── Transaction lifecycle ──────────────────────────────────────────── public async Task BeginTransactionAsync(LoginDTO LoginDTO) { using var activity = GB5ActivitySources.Database.StartActivity("db transaction begin", ActivityKind.Client); activity?.SetTag("db.system", "mssql"); var connection = await CreateNewConnectionAsync(LoginDTO); await OpenWithClearErrorAsync(connection).ConfigureAwait(false); var transaction = connection.BeginTransaction(); activity?.SetStatus(ActivityStatusCode.Ok); return transaction; } public async Task CommitAsync(DbTransaction transaction) { using var activity = GB5ActivitySources.Database.StartActivity("db transaction commit", ActivityKind.Client); activity?.SetTag("db.system", "mssql"); try { if (transaction?.Connection != null) transaction.Commit(); activity?.SetStatus(ActivityStatusCode.Ok); } catch (Exception ex) { MarkDbError(activity, ex); throw; } finally { await CleanupAsync(transaction); } } public async Task RollbackAsync(DbTransaction transaction) { using var activity = GB5ActivitySources.Database.StartActivity("db transaction rollback", ActivityKind.Client); activity?.SetTag("db.system", "mssql"); try { if (transaction?.Connection != null) transaction.Rollback(); activity?.SetStatus(ActivityStatusCode.Ok); } catch (Exception ex) { MarkDbError(activity, ex); throw; } finally { await CleanupAsync(transaction); } } // BUG FIX: this previously disposed only the DbTransaction, never the underlying // DbConnection that BeginTransactionAsync opened — DbTransaction.DisposeAsync() ends // the transaction but does not close the connection it belongs to. Every transactional // call (BeginTransactionAsync → Commit/RollbackAsync) was leaking one physical // connection per call, permanently checked out of the pool for the life of the // process — confirmed live: an hours-old "sleeping" session still holding an exclusive // key lock from a transaction whose Commit/Rollback had already run. Capture the // connection reference BEFORE disposing the transaction (some providers null out // DbTransaction.Connection once the transaction is disposed). private async Task CleanupAsync(DbTransaction? transaction) { if (transaction == null) return; var connection = transaction.Connection; try { await transaction.DisposeAsync(); } catch { /* ignore */ } if (connection != null) { try { await connection.DisposeAsync(); } catch { /* ignore */ } } } // ── QuerySingleAsync ───────────────────────────────────────────────── public async Task QuerySingleAsync(LoginDTO LoginDTO, string sql, object parameters = null!, DbTransaction? transaction = null, bool useReadUncommitted = false, CancellationToken cancellationToken = default) { bool shouldClose = (transaction == null); WarnIfUpdLockWithoutTransaction(sql, transaction); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(LoginDTO); using var dbActivity = StartDbTrace(sql, LoginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection, cancellationToken).ConfigureAwait(false); sql = ApplyReadUncommitted(sql, useReadUncommitted, !shouldClose); var cmd = new CommandDefinition(sql, parameters, transaction, commandTimeout: DefaultCommandTimeout, cancellationToken: cancellationToken); var result = await connection.QuerySingleOrDefaultAsync(cmd)!; dbActivity?.SetStatus(ActivityStatusCode.Ok); return result; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── QueryAsync (LoginDTO) ──────────────────────────────────────────── public async Task> QueryAsync(LoginDTO LoginDTO, string sql, object parameters = null!, DbTransaction? transaction = null, bool useReadUncommitted = false, CancellationToken cancellationToken = default) { bool shouldClose = (transaction == null); WarnIfUpdLockWithoutTransaction(sql, transaction); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(LoginDTO); using var dbActivity = StartDbTrace(sql, LoginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection, cancellationToken).ConfigureAwait(false); sql = ApplyReadUncommitted(sql, useReadUncommitted, !shouldClose); var cmd = new CommandDefinition(sql, parameters, transaction, commandTimeout: DefaultCommandTimeout, cancellationToken: cancellationToken); var result = await connection.QueryAsync(cmd); dbActivity?.SetTag("db.rows_returned", result.TryGetNonEnumeratedCount(out var c) ? c : -1); dbActivity?.SetStatus(ActivityStatusCode.Ok); return result; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── QueryAsync (ConnectionName) ────────────────────────────────────── public async Task> QueryAsync(string ConnectionName, string sql, object parameters = null!) { WarnIfUpdLockWithoutTransaction(sql, null); using var dbActivity = StartDbTrace(sql, ConnectionName, parameters); try { _connectionstring = await _Appconnection.DBConnectionStringCachedConnectionName(ConnectionName); var databaseType = await _Appconnection.DatabaseTypeCached(ConnectionName); await using DbConnection connection = databaseType switch { DBTYPE.POSTGRESQL => new NpgsqlConnection(_connectionstring), DBTYPE.SQL => new SqlConnection(_connectionstring), _ => throw new NotSupportedException($"Database type '{databaseType}' is not supported.") }; await OpenWithClearErrorAsync(connection).ConfigureAwait(false); var result = await connection.QueryAsync(sql, parameters, commandTimeout: DefaultCommandTimeout); dbActivity?.SetStatus(ActivityStatusCode.Ok); return result; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } } // ── ExecuteQueryAsync (IAsyncEnumerable) ───────────────────────────── public async IAsyncEnumerable ExecuteQueryAsync(LoginDTO LoginDTO, string sql, object parameters = null!, DbTransaction? transaction = null, bool useReadUncommitted = false) { bool shouldClose = (transaction == null); WarnIfUpdLockWithoutTransaction(sql, transaction); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(LoginDTO); IEnumerable? result = null; using (var dbActivity = StartDbTrace(sql, LoginDTO?.DatabaseName, parameters)) { try { if (shouldClose) await OpenWithClearErrorAsync(connection).ConfigureAwait(false); sql = ApplyReadUncommitted(sql, useReadUncommitted, !shouldClose); result = await connection.QueryAsync(sql, parameters, transaction, commandTimeout: DefaultCommandTimeout); dbActivity?.SetStatus(ActivityStatusCode.Ok); } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } foreach (var item in result!) yield return item; } // ── ExecuteAsync ───────────────────────────────────────────────────── public async Task ExecuteAsync(LoginDTO LoginDTO, string sql, object parameters = null!, DbTransaction? transaction = null, CancellationToken cancellationToken = default) { bool shouldClose = (transaction == null); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(LoginDTO); using var dbActivity = StartDbTrace(sql, LoginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection, cancellationToken).ConfigureAwait(false); var cmd = new CommandDefinition(sql, parameters, transaction, commandTimeout: DefaultCommandTimeout, cancellationToken: cancellationToken); var rowsAffected = await connection.ExecuteAsync(cmd); dbActivity?.SetTag("db.rows_affected", rowsAffected); dbActivity?.SetStatus(ActivityStatusCode.Ok); return rowsAffected; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── QueryMultipleAsync (StoredProcedure default) ───────────────────── public async Task QueryMultipleAsync(LoginDTO loginDTO, string sql, object parameters = null!, DbTransaction? transaction = null) { bool shouldClose = (transaction == null); DbConnection? connection = transaction?.Connection ?? await CreateNewConnectionAsync(loginDTO); using var dbActivity = StartDbTrace(sql, loginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection).ConfigureAwait(false); var reader = await connection.QueryMultipleAsync( sql, parameters, transaction, commandTimeout: DefaultCommandTimeout, commandType: CommandType.StoredProcedure); connection = null; // transfer ownership to caller dbActivity?.SetStatus(ActivityStatusCode.Ok); return reader; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (connection != null && shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── QueryMultipleAsync (explicit CommandType) ───────────────────────── public async Task QueryMultipleAsync(LoginDTO loginDTO, string sql, object parameters = null!, DbTransaction? transaction = null, CommandType commandType = CommandType.StoredProcedure) { bool shouldClose = (transaction == null); DbConnection? connection = transaction?.Connection ?? await CreateNewConnectionAsync(loginDTO); using var dbActivity = StartDbTrace(sql, loginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection).ConfigureAwait(false); var reader = await connection.QueryMultipleAsync( sql, parameters, transaction, commandTimeout: DefaultCommandTimeout, commandType: commandType); connection = null; dbActivity?.SetStatus(ActivityStatusCode.Ok); return reader; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (connection != null && shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── StreamAsync (IAsyncEnumerable) ──────────────────────────────────── public async IAsyncEnumerable StreamAsync(LoginDTO LoginDTO, string sql, object parameters = null!, [EnumeratorCancellation] CancellationToken cancellationToken = default, DbTransaction? transaction = null, bool useReadUncommitted = false) { bool shouldClose = (transaction == null); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(LoginDTO); IEnumerable? items = null; using (var dbActivity = StartDbTrace(sql, LoginDTO?.DatabaseName, parameters)) { try { if (shouldClose) await OpenWithClearErrorAsync(connection, cancellationToken).ConfigureAwait(false); sql = ApplyReadUncommitted(sql, useReadUncommitted, !shouldClose); var cmd = new CommandDefinition(sql, parameters, transaction, DefaultCommandTimeout, cancellationToken: cancellationToken); items = await connection.QueryAsync(cmd).ConfigureAwait(false); dbActivity?.SetStatus(ActivityStatusCode.Ok); } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } foreach (var item in items!) yield return item; } // ── QueryPagedAsync ─────────────────────────────────────────────────── public async Task> QueryPagedAsync(LoginDTO LoginDTO, string sql, object parameters = null!, DbTransaction? transaction = null, bool useReadUncommitted = false) { bool shouldClose = (transaction == null); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(LoginDTO); using var dbActivity = StartDbTrace(sql, LoginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection).ConfigureAwait(false); string effectiveSql = ApplyReadUncommitted(sql, useReadUncommitted, !shouldClose); var items = await connection.QueryAsync(effectiveSql, parameters, transaction, commandTimeout: DefaultCommandTimeout); var countSql = ApplyReadUncommitted("SELECT COUNT(*) FROM (" + sql + ") AS Total", useReadUncommitted, !shouldClose); var totalCount = await connection.QuerySingleAsync(countSql, parameters, transaction, commandTimeout: DefaultCommandTimeout); dbActivity?.SetTag("db.rows_returned", totalCount); dbActivity?.SetStatus(ActivityStatusCode.Ok); return new PagedResult { Items = items.ToList(), TotalCount = totalCount }; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── BulkInsertAsync ─────────────────────────────────────────────────── public async Task BulkInsertAsync(LoginDTO LoginDTO, string sql, IEnumerable items, DbTransaction? transaction = null) { bool shouldClose = (transaction == null); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(LoginDTO); using var dbActivity = StartDbTrace(sql, LoginDTO?.DatabaseName); if (items is ICollection itemColl) dbActivity?.SetTag("db.bulk.item_count", itemColl.Count); try { if (shouldClose) await OpenWithClearErrorAsync(connection).ConfigureAwait(false); var rowsAffected = await connection.ExecuteAsync(sql, items, transaction, commandTimeout: DefaultCommandTimeout); dbActivity?.SetTag("db.rows_affected", rowsAffected); dbActivity?.SetStatus(ActivityStatusCode.Ok); return rowsAffected; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── QueryStreamAsync (streaming reader, IAsyncEnumerable) ──────────── public async IAsyncEnumerable QueryStreamAsync(LoginDTO LoginDTO, string sql, object parameters = null!, [EnumeratorCancellation] CancellationToken cancellationToken = default, DbTransaction? transaction = null, bool useReadUncommitted = false) { bool shouldClose = (transaction == null); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(LoginDTO); var dbActivity = StartDbTrace(sql, LoginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection, cancellationToken).ConfigureAwait(false); sql = ApplyReadUncommitted(sql, useReadUncommitted, !shouldClose); var commandDefinition = new CommandDefinition( commandText: sql, parameters: parameters, transaction: transaction, cancellationToken: cancellationToken, commandTimeout: DefaultCommandTimeout); DbDataReader reader; try { reader = await connection.ExecuteReaderAsync(commandDefinition, CommandBehavior.SequentialAccess); } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } dbActivity?.SetStatus(ActivityStatusCode.Ok); var parser = reader.GetRowParser(); await using var _ = reader; while (await reader.ReadAsync(cancellationToken)) yield return parser(reader); } finally { dbActivity?.Dispose(); if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── ExecuteInTransactionAsync ───────────────────────────────────────── public async Task ExecuteInTransactionAsync(LoginDTO LoginDTO, IEnumerable<(string sql, object param)> statements) { using var txActivity = GB5ActivitySources.Database.StartActivity("db transaction batch", ActivityKind.Client); txActivity?.SetTag("db.system", "mssql"); txActivity?.SetTag("db.name", LoginDTO?.DatabaseName); await using var connection = await CreateNewConnectionAsync(LoginDTO); await OpenWithClearErrorAsync(connection).ConfigureAwait(false); using var transaction = connection.BeginTransaction(); try { int stmtIndex = 0; foreach (var (sql, param) in statements) { using var stmtActivity = StartDbTrace(sql, LoginDTO?.DatabaseName, param); try { await connection.ExecuteAsync(sql, param, transaction, commandTimeout: DefaultCommandTimeout); stmtActivity?.SetTag("db.batch.index", stmtIndex++); stmtActivity?.SetStatus(ActivityStatusCode.Ok); } catch (Exception ex) { MarkDbError(stmtActivity, ex); throw; } } transaction.Commit(); txActivity?.SetStatus(ActivityStatusCode.Ok); } catch (Exception ex) { MarkDbError(txActivity, ex); transaction.Rollback(); throw; } } // ── ExecuteScalarAsync ──────────────────────────────────────────────── public async Task ExecuteScalarAsync(LoginDTO loginDTO, string sql, object? parameters = null, DbTransaction? transaction = null, bool useReadUncommitted = false, CancellationToken cancellationToken = default) { bool shouldClose = (transaction == null); WarnIfUpdLockWithoutTransaction(sql, transaction); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(loginDTO); using var dbActivity = StartDbTrace(sql, loginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection, cancellationToken).ConfigureAwait(false); sql = ApplyReadUncommitted(sql, useReadUncommitted, !shouldClose); var cmd = new CommandDefinition(sql, parameters, transaction, commandTimeout: DefaultCommandTimeout, cancellationToken: cancellationToken); var result = await connection.ExecuteScalarAsync(cmd); dbActivity?.SetStatus(ActivityStatusCode.Ok); return result; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── BulkInsertMultipleAsync (TVP stored procedure) ─────────────────── public async Task BulkInsertMultipleAsync(LoginDTO loginDTO, string commandText, Dictionary tvpParameters) { // The TENANT's own DbType, not the GB5 system's global default — this only ever // builds a SqlConnection below, so a tenant on Postgres must be rejected here, not // let through because the global default happens to be SQL Server. if (await _Appconnection.DatabaseTypeCached(loginDTO) != DBTYPE.SQL) throw new NotSupportedException("BulkInsertMultipleAsync with TVP is only supported for SQL Server."); var tvpParamSummary = tvpParameters.ToDictionary( kv => kv.Key, kv => (object)$""); using var dbActivity = StartDbTrace(commandText, loginDTO?.DatabaseName, tvpParamSummary); dbActivity?.SetTag("db.operation", "EXECUTE"); try { _connectionstring = await _Appconnection.DBConnectionStringCached(loginDTO); if (!_connectionstring.Contains("Enlist=", StringComparison.OrdinalIgnoreCase)) _connectionstring += ";Enlist=false"; await using var connection = new SqlConnection(_connectionstring); await OpenWithClearErrorAsync(connection).ConfigureAwait(false); using var transaction = connection.BeginTransaction(); try { using var command = new SqlCommand(commandText, connection, transaction) { CommandType = CommandType.StoredProcedure, CommandTimeout = DefaultCommandTimeout }; foreach (var param in tvpParameters) { var sqlParam = command.Parameters.AddWithValue(param.Key, param.Value.value); sqlParam.SqlDbType = SqlDbType.Structured; sqlParam.TypeName = param.Value.typeName; } int rowsAffected = await command.ExecuteNonQueryAsync(); transaction.Commit(); dbActivity?.SetTag("db.rows_affected", rowsAffected); dbActivity?.SetStatus(ActivityStatusCode.Ok); return rowsAffected; } catch (SqlException ex) { transaction.Rollback(); MarkDbError(dbActivity, ex); dbActivity?.SetTag("db.sql.error_number", ex.Number); throw new Exception($"SQL Error {ex.Number}: {ex.Message}", ex); } catch (Exception ex) { transaction.Rollback(); MarkDbError(dbActivity, ex); throw new Exception($"Unexpected Error: {ex.Message}", ex); } } catch (Exception ex) when (dbActivity?.Status != ActivityStatusCode.Error) { MarkDbError(dbActivity, ex); throw; } } // ── ExecuteTvpCountQueryAsync ───────────────────────────────────────── public async Task ExecuteTvpCountQueryAsync(LoginDTO loginDTO, string sqlQuery, Dictionary tvpParameters) { if (loginDTO.DatabaseType != DBTYPE.SQL) throw new NotSupportedException("ExecuteTvpCountQueryAsync with TVP is only supported for SQL Server."); var tvpParamSummary = tvpParameters.ToDictionary( kv => kv.Key, kv => (object)$""); using var dbActivity = StartDbTrace(sqlQuery, loginDTO?.DatabaseName, tvpParamSummary); try { _connectionstring = await _Appconnection.DBConnectionStringCached(loginDTO); await using SqlConnection connection = new SqlConnection(_connectionstring); await OpenWithClearErrorAsync(connection).ConfigureAwait(false); using SqlCommand command = new SqlCommand(sqlQuery, connection) { CommandType = CommandType.Text, CommandTimeout = DefaultCommandTimeout }; foreach (var param in tvpParameters) { var sqlParam = command.Parameters.AddWithValue(param.Key, param.Value.value); sqlParam.SqlDbType = SqlDbType.Structured; sqlParam.TypeName = param.Value.typeName; } var result = await command.ExecuteScalarAsync(); var count = result != null ? Convert.ToInt32(result) : 0; dbActivity?.SetTag("db.rows_returned", count); dbActivity?.SetStatus(ActivityStatusCode.Ok); return count; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } } // ── TotalCountSql ───────────────────────────────────────────────────── public async Task TotalCountSql(string sqlQuery, LoginDTO loginDTO, object? parameters = null, bool isCountQuery = false, DbTransaction? transaction = null, bool useReadUncommitted = false) { var databaseType = loginDTO.DatabaseType; if (databaseType == DBTYPE.POSTGRESQL) sqlQuery = ConvertSqlToPostgres.ConvertSqlServerToPostgres(sqlQuery); string text = isCountQuery ? sqlQuery : await _GB5CommonFunction.QueryReplaceForTotal(sqlQuery); bool shouldClose = (transaction == null); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(loginDTO); using var dbActivity = StartDbTrace(text, loginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection).ConfigureAwait(false); text = ApplyReadUncommitted(text, useReadUncommitted, !shouldClose); var count = await connection.QuerySingleAsync(text, parameters, transaction, commandTimeout: DefaultCommandTimeout); dbActivity?.SetTag("db.rows_returned", count); dbActivity?.SetStatus(ActivityStatusCode.Ok); return count; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── QueryMultiMapAsync (2-type) ──────────────────────────────────────── public async Task> QueryMultiMapAsync( LoginDTO loginDTO, string sql, Func map, object parameters = null!, string splitOn = "Id", DbTransaction? transaction = null, bool useReadUncommitted = false) { bool shouldClose = (transaction == null); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(loginDTO); using var dbActivity = StartDbTrace(sql, loginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection).ConfigureAwait(false); sql = ApplyReadUncommitted(sql, useReadUncommitted, !shouldClose); var result = await connection.QueryAsync(sql, map, param: parameters, transaction: transaction, splitOn: splitOn, commandTimeout: DefaultCommandTimeout); dbActivity?.SetStatus(ActivityStatusCode.Ok); return result; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── QueryMultiMapAsync (3-type) ──────────────────────────────────────── public async Task> QueryMultiMapAsync( LoginDTO loginDTO, string sql, Func map, object parameters = null!, string splitOn = "Id", DbTransaction? transaction = null, bool useReadUncommitted = false) { bool shouldClose = (transaction == null); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(loginDTO); using var dbActivity = StartDbTrace(sql, loginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection).ConfigureAwait(false); sql = ApplyReadUncommitted(sql, useReadUncommitted, !shouldClose); var result = await connection.QueryAsync(sql, map, param: parameters, transaction: transaction, splitOn: splitOn, commandTimeout: DefaultCommandTimeout); dbActivity?.SetStatus(ActivityStatusCode.Ok); return result; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── QueryMultiMapAsync (5-type) ──────────────────────────────────────── public async Task> QueryMultiMapAsync( LoginDTO loginDTO, string sql, Func map, object parameters = null!, string splitOn = "Id,Id,Id,Id", DbTransaction? transaction = null, bool useReadUncommitted = false) { bool shouldClose = (transaction == null); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(loginDTO); using var dbActivity = StartDbTrace(sql, loginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection).ConfigureAwait(false); sql = ApplyReadUncommitted(sql, useReadUncommitted, !shouldClose); var result = await connection.QueryAsync(sql, map, param: parameters, transaction: transaction, splitOn: splitOn, commandTimeout: DefaultCommandTimeout); dbActivity?.SetStatus(ActivityStatusCode.Ok); return result; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── QueryMultiMapAsync (7-type) ──────────────────────────────────────── public async Task> QueryMultiMapAsync( LoginDTO loginDTO, string sql, Func map, object parameters = null!, string splitOn = "Id,Id,Id,Id,Id,Id", DbTransaction? transaction = null, bool useReadUncommitted = false) { bool shouldClose = (transaction == null); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(loginDTO); using var dbActivity = StartDbTrace(sql, loginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection).ConfigureAwait(false); sql = ApplyReadUncommitted(sql, useReadUncommitted, !shouldClose); var result = await connection.QueryAsync(sql, map, param: parameters, transaction: transaction, splitOn: splitOn, commandTimeout: DefaultCommandTimeout); dbActivity?.SetStatus(ActivityStatusCode.Ok); return result; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── QueryAsyncArray (2-type map) ────────────────────────────────────── public async Task> QueryAsyncArray( LoginDTO loginDTO, string sql, Func map, object parameters = null!, string splitOn = "Id", DbTransaction? transaction = null, bool useReadUncommitted = false) { bool shouldClose = (transaction == null); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(loginDTO); using var dbActivity = StartDbTrace(sql, loginDTO?.DatabaseName, parameters); try { _connectionstring = await _Appconnection.DBConnectionStringCached(loginDTO); if (shouldClose) await OpenWithClearErrorAsync(connection).ConfigureAwait(false); sql = ApplyReadUncommitted(sql, useReadUncommitted, !shouldClose); var result = await connection.QueryAsync(sql, map, param: parameters, transaction: transaction, splitOn: splitOn, commandTimeout: DefaultCommandTimeout); dbActivity?.SetStatus(ActivityStatusCode.Ok); return result; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } // ── Session methods ──────────────────────────────────────────────────── public async Task> SessionQueryAsync(LoginDTO loginDTO, string sql, object param = null!) { if (_httpContext.HttpContext?.Items["DB_CONN"] is DbConnection conn && _httpContext.HttpContext?.Items["DB_TRAN"] is DbTransaction tran) { using var dbActivity = StartDbTrace(sql, loginDTO?.DatabaseName, param); try { var result = await conn.QueryAsync(sql, param, tran, commandTimeout: DefaultCommandTimeout); dbActivity?.SetStatus(ActivityStatusCode.Ok); return result; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } } throw new InvalidOperationException("Session not started. Use BeginTransactionAsync (Legacy) or standard transaction methods."); } public async Task SessionQuerySingleAsync(LoginDTO loginDTO, string sql, object param = null!) { if (_httpContext.HttpContext?.Items["DB_CONN"] is DbConnection conn && _httpContext.HttpContext?.Items["DB_TRAN"] is DbTransaction tran) { using var dbActivity = StartDbTrace(sql, loginDTO?.DatabaseName, param); try { var result = await conn.QuerySingleOrDefaultAsync(sql, param, tran, commandTimeout: DefaultCommandTimeout)!; dbActivity?.SetStatus(ActivityStatusCode.Ok); return result; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } } throw new InvalidOperationException("Session not started."); } public async Task SessionExecuteAsync(LoginDTO loginDTO, string sql, object param = null!) { if (_httpContext.HttpContext?.Items["DB_CONN"] is DbConnection conn && _httpContext.HttpContext?.Items["DB_TRAN"] is DbTransaction tran) { using var dbActivity = StartDbTrace(sql, loginDTO?.DatabaseName, param); try { var rowsAffected = await conn.ExecuteAsync(sql, param, tran, commandTimeout: DefaultCommandTimeout); dbActivity?.SetTag("db.rows_affected", rowsAffected); dbActivity?.SetStatus(ActivityStatusCode.Ok); return rowsAffected; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } } throw new InvalidOperationException("Session not started."); } public async Task SameSessionCommitAsync() { if (_httpContext.HttpContext?.Items["DB_TRAN"] is DbTransaction tran && _httpContext.HttpContext?.Items["DB_CONN"] is DbConnection conn) { tran.Commit(); await tran.DisposeAsync(); await conn.DisposeAsync(); _httpContext.HttpContext.Items.Remove("DB_TRAN"); _httpContext.HttpContext.Items.Remove("DB_CONN"); } } public async Task SameSessionRollbackAsync() { if (_httpContext.HttpContext?.Items["DB_TRAN"] is DbTransaction tran && _httpContext.HttpContext?.Items["DB_CONN"] is DbConnection conn) { tran.Rollback(); await tran.DisposeAsync(); await conn.DisposeAsync(); _httpContext.HttpContext.Items.Remove("DB_TRAN"); _httpContext.HttpContext.Items.Remove("DB_CONN"); } } // ── ProcedureExecuteAsync ───────────────────────────────────────────── public async Task> ProcedureExecuteAsync(LoginDTO loginDTO, IEnumerable<(string ProcedureName, object Param)> procedures) { using var txActivity = GB5ActivitySources.Database.StartActivity("db procedure batch", ActivityKind.Client); txActivity?.SetTag("db.system", "mssql"); txActivity?.SetTag("db.name", loginDTO?.DatabaseName); await using var connection = await CreateNewConnectionAsync(loginDTO); await OpenWithClearErrorAsync(connection).ConfigureAwait(false); using var transaction = connection.BeginTransaction(); try { List results = new(); foreach (var (procedureName, param) in procedures) { using var spActivity = StartDbTrace(procedureName, loginDTO?.DatabaseName, param); spActivity?.SetTag("db.operation", "EXECUTE"); try { var dynamicParams = new DynamicParameters(param); var spResult = await connection.QueryAsync(procedureName, dynamicParams, transaction: transaction, commandTimeout: DefaultCommandTimeout, commandType: CommandType.StoredProcedure); string status = spResult.FirstOrDefault() ?? "NO RESPONSE"; spActivity?.SetTag("db.sp.result", status); spActivity?.SetStatus(ActivityStatusCode.Ok); results.Add(status); } catch (Exception ex) { MarkDbError(spActivity, ex); throw; } } transaction.Commit(); txActivity?.SetStatus(ActivityStatusCode.Ok); return results; } catch (Exception ex) { MarkDbError(txActivity, ex); transaction.Rollback(); throw new Exception($"{ex.Message}", ex); } } // ── ExecuteAsync (explicit transaction overload) ────────────────────── Task IQueryExecutor.ExecuteAsync(string sql, object param, DbTransaction transaction) { if (transaction == null || transaction.Connection == null) throw new InvalidOperationException("Invalid transaction."); return transaction.Connection.ExecuteAsync(sql, param, transaction, commandTimeout: DefaultCommandTimeout); } // ── ExecuteIdentityAsync ────────────────────────────────────────────── public async Task ExecuteIdentityAsync(LoginDTO loginDTO, string sql, object? parameters = null, DbTransaction? transaction = null) { bool shouldClose = (transaction == null); DbConnection connection = transaction?.Connection ?? await CreateNewConnectionAsync(loginDTO); using var dbActivity = StartDbTrace(sql, loginDTO?.DatabaseName, parameters); try { if (shouldClose) await OpenWithClearErrorAsync(connection).ConfigureAwait(false); var result = await connection.ExecuteScalarAsync(sql, parameters, transaction, commandTimeout: DefaultCommandTimeout); var id = Convert.ToInt32(result); dbActivity?.SetTag("db.identity_returned", id); dbActivity?.SetStatus(ActivityStatusCode.Ok); return id; } catch (Exception ex) { MarkDbError(dbActivity, ex); throw; } finally { if (shouldClose && connection.State == ConnectionState.Open) await connection.DisposeAsync(); } } } public class PagedResult { public IEnumerable? Items { get; set; } public int TotalCount { get; set; } } }