using System.Data; using Microsoft.Data.SqlClient; using System.Data.Common; using Dapper; using FrameworkDAL.DTO.Workflow; using GB5Shared.Connection; using GB5Shared.DTO.Framework.CommonConfig; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.Framework.ServerConfig; using GB5Shared.QueryExecutor; using Microsoft.Extensions.Options; using Npgsql; using Quartz; using Serilog; using static GB5Shared.GB5Constant.Constant; namespace FrameworkSL.Endpoints.WorkFlow { public class WorkFlowTrigger : IJob { private readonly IApplicationConnection _appConnection; private readonly IOptionsSnapshot _databaseDTO; private readonly IQueryExecutor queryExecutor; public WorkFlowTrigger( IApplicationConnection appConnection, IOptionsSnapshot databaseDTO, IQueryExecutor IQueryExecutor) { _appConnection = appConnection; _databaseDTO = databaseDTO; queryExecutor = IQueryExecutor; } public async Task Execute(IJobExecutionContext context) { try { string systemConn = await _appConnection.Gb5SystemConnectionString(); int dbType = _databaseDTO.Value.DataBaseType; // Step 1: Fetch all active tenants string sql = @" SELECT SERVERCONFIG1.CLIENTID AS ClientId, SERVERCONFIG1.DATABASENAME AS DatabaseName, SERVERCONFIG1.DATABASETYPE AS DbType, SERVERCONFIG1.CONNECTIONNAME AS ConnectionName FROM MSERVERCONFIG SERVERCONFIG1 JOIN MSERVER SERVER1 ON SERVERCONFIG1.SERVERID = SERVER1.SERVERID WHERE SERVERCONFIG1.STATUS = 1 AND SERVERCONFIG1.CONNECTIONNAME <> 'ACTIVITI'"; List tenants; if (dbType == DBTYPE.SQL) { using var connection = new SqlConnection(systemConn); connection.Open(); tenants = (await connection.QueryAsync(sql)).AsList(); } else if (dbType == DBTYPE.POSTGRESQL) { using var connection = new NpgsqlConnection(systemConn); connection.Open(); tenants = (await connection.QueryAsync(sql)).AsList(); } else { throw new NotSupportedException($"Unsupported DB Type: {dbType}"); } Log.Information("Fetched {Count} tenants", tenants.Count); // Step 2: Process each tenant foreach (var tenant in tenants) { try { var login = new LoginDTO { UserId = -1, ClientId = tenant.ClientId, ConnectionDatabaseName = tenant.ConnectionName, DatabaseName = tenant.DatabaseName, DatabaseType = tenant.DbType }; var DBTransaction = await queryExecutor.BeginTransactionAsync(login); try { if (login.DatabaseType == DBTYPE.SQL) { await ProcessTenantWorkflows(login, DBTransaction); } else if (login.DatabaseType == DBTYPE.POSTGRESQL) { await ProcessTenantWorkflows(login, DBTransaction); } else { Log.Warning("Unsupported DB type for tenant {ClientId}: {DbType}", tenant.ClientId, login.DatabaseType); } await queryExecutor.CommitAsync(DBTransaction); } catch (Exception ex) { await queryExecutor.RollbackAsync(DBTransaction); Log.Error(ex, "Error processing workflows for tenant {ClientId}", tenant.ClientId); } } catch (Exception tenantEx) { Log.Error(tenantEx, "Error processing tenant {ClientId}", tenant.ClientId); } } } catch (Exception ex) { Log.Error(ex, "Workflow trigger execution failed for system connection."); } } private async Task ProcessTenantWorkflows(LoginDTO LoginDTO ,DbTransaction DbTransaction) { // Step 1: Fetch top 50 pending workflows var pendingWorkflows = (await queryExecutor.QueryAsync(LoginDTO, @"SELECT TOP 50 WI.ENTITYID AS EntityId, ME.ENTITYCODE AS EntityCode, ME.ENTITYNAME AS EntityName, WI.OBJECTID AS ObjectId, ME.DBOBJECTID AS DbObjectId, DB.DBOBJECTNAME AS DbObjectName FROM TWORKFLOWINSTANCE WI INNER JOIN MENTITY ME ON WI.ENTITYID = ME.ENTITYID INNER JOIN DBOBJECT DB ON ME.DBOBJECTID = DB.DBOBJECTID WHERE WI.WORKFLOWSTATUS = 1;",null!, DbTransaction )).AsList(); Log.Information("Tenant fetched {Count} pending workflows", pendingWorkflows.Count); foreach (var wf in pendingWorkflows) { await ProcessWorkflowItem( wf, LoginDTO, DbTransaction); } } private async Task ProcessWorkflowItem(PendingWorkflowInstanceListDTO wf,LoginDTO LoginDTO , DbTransaction DbTransaction) { try { if (string.IsNullOrWhiteSpace(wf.DbObjectName)) return; string tableName = wf.DbObjectName.Trim(); if (!IsValidTableName(tableName)) { Log.Warning("Skipping invalid table name: {Table}", tableName); return; } // Step 1: Get primary key column dynamically string? primaryKeyColumn = await GetPrimaryKeyColumnAsync(tableName, LoginDTO, DbTransaction); if (string.IsNullOrEmpty(primaryKeyColumn)) { Log.Warning("Cannot determine primary key for table {Table}", tableName); return; } // Step 2: Check if record exists string checkSql = $"SELECT COUNT(1) FROM {tableName} WHERE {primaryKeyColumn} = @objectid"; int count = await queryExecutor.QuerySingleAsync(LoginDTO ,checkSql, new { objectid = wf.ObjectId }, DbTransaction); if (count <= 0) return; // Step 3: Update STATUS dynamically string updateSql = $"UPDATE {tableName} SET STATUS = 1 WHERE {primaryKeyColumn} = @objectid"; await queryExecutor.ExecuteAsync(LoginDTO ,updateSql, new { objectid = wf.ObjectId }, DbTransaction); // Step 4: Update workflow table string updateWorkflowSql = @" UPDATE TWORKFLOWINSTANCE SET WORKFLOWSTATUS = 2 WHERE ENTITYID = @entityid AND OBJECTID = @objectid"; await queryExecutor.ExecuteAsync(LoginDTO ,updateWorkflowSql, new { entityid = wf.EntityId, objectid =wf.ObjectId }, DbTransaction); Log.Information("Approved ObjectId {ObjectId} in table {Table} using primary key {PK}", wf.ObjectId, tableName, primaryKeyColumn); } catch (Exception) { throw; } } private async Task GetPrimaryKeyColumnAsync(string tableName, LoginDTO LoginDTO,DbTransaction DbTransaction) { try { if (!IsValidTableName(tableName)) return null; if (LoginDTO.DatabaseType == DBTYPE.SQL) { string sql = @" SELECT COLUMN_NAME as ColumnName FROM INFORMATION_SCHEMA.TABLE_CONSTRAINTS AS TC JOIN INFORMATION_SCHEMA.KEY_COLUMN_USAGE AS KU ON TC.CONSTRAINT_NAME = KU.CONSTRAINT_NAME WHERE TC.TABLE_NAME = @tablename AND TC.CONSTRAINT_TYPE = 'PRIMARY KEY'"; PendingWorkflowInstanceListDTO PendingWorkflowInstanceListDTO = await queryExecutor.QuerySingleAsync(LoginDTO,sql, new { tablename = tableName }, DbTransaction); return PendingWorkflowInstanceListDTO.ColumnName; } else if (LoginDTO.DatabaseType == DBTYPE.POSTGRESQL) { string sql = @" SELECT a.attname FROM pg_index i JOIN pg_attribute a ON a.attrelid = i.indrelid AND a.attnum = ANY(i.indkey) WHERE i.indrelid = @TableName::regclass AND i.indisprimary = true LIMIT 1"; PendingWorkflowInstanceListDTO PendingWorkflowInstanceListDTO = await queryExecutor.QuerySingleAsync(LoginDTO, sql, new { tablename = tableName }, DbTransaction); return PendingWorkflowInstanceListDTO.ColumnName; } return null; } catch (Exception ex) { Log.Error(ex, "Error fetching primary key for table {Table}", tableName); return null; } } private static bool IsValidTableName(string tableName) { return tableName.All(c => char.IsLetterOrDigit(c) || c == '_'); } } }