using GB5Shared.DTO.Framework.Login; using GB5Shared.EventLogPublish; using GB5Shared.Resource.Response; using Microsoft.Extensions.Logging; using SwBLL.ClientDatabase; using SwBLL.DbServer; using SwBLL.DdlScript; using SwBLL.Provisioning; using SwBLL.UserResourceRole; using SwDAL.CustomCode.ChangeRequest; using SwDAL.DTO.ChangeRequest; using SwDAL.Enums; using System.Text.Json; using static GB5Shared.GB5Constant.Constant; namespace SwBLL.ChangeRequest; public class ChangeRequestBLL : IChangeRequestBLL { private readonly IChangeRequestDAL _dal; private readonly IClientDatabaseBLL _clientDatabaseBLL; private readonly IDbServerBLL _dbServerBLL; private readonly ITargetDbExecutor _targetDbExecutor; private readonly IDdlScriptBLL _ddlScriptBLL; private readonly IClientDatabaseProvisioner _clientDatabaseProvisioner; private readonly IUserResourceRoleBLL _userResourceRoleBLL; private readonly EventLogPublish _eventLogPublish; private readonly ILogger _logger; // Dapr pub/sub topic a cross-module caller (e.g. a future Entitlement onboarding // orchestrator) subscribes to in order to resume after a human-approved // ClientProvisioning CR reaches Executed, without polling MSWCHANGEREQUEST.CrStatus. private const string ClientProvisioningExecutedTopic = "sqlworkbench.changerequest.executed"; // Separate, deliberately generic topic for every OTHER category of successful CR execution // (SchemaChange/DataMigration — ordinary post-provisioning changes to an already-live // client database). Kept distinct from ClientProvisioningExecutedTopic above so that // topic's existing, narrower meaning ("a new client database now exists") isn't diluted for // whatever already subscribes to it. Payload shape is identical and just as generic — // {ChangeRequestId, ClientDatabaseId} — SqlWorkbench itself has no opinion on who listens or // why; a consumer that cares about a specific tenant's own schema-version bookkeeping (or // anything else) reacts entirely on its own side, never referenced from here. private const string ChangeRequestAppliedTopic = "sqlworkbench.changerequest.applied"; // Tracker §37 Decision 2 — bare "something changed, go re-check" wake-up for a // ClientProvisioning CR's DDL execution, published every ScriptsProgressPublishBatchSize // scripts (and once more on the last script) inside ExecuteDdlCr's loop. Deliberately carries // no meaningful payload data — same "always re-verify via a real call, never trust the event // shape" posture as ClientProvisioningExecutedTopic's own subscriber; a consumer calls // GetChangeRequestById back to read the real, live ScriptsDone/ScriptsTotal instead of // attempting to parse EventLogPublish's own double-JSON-encoded envelope. private const string ClientProvisioningProgressTopic = "sqlworkbench.changerequest.progress"; private const int ScriptsProgressPublishBatchSize = 25; public ChangeRequestBLL( IChangeRequestDAL dal, IClientDatabaseBLL clientDatabaseBLL, IDbServerBLL dbServerBLL, ITargetDbExecutor targetDbExecutor, IDdlScriptBLL ddlScriptBLL, IClientDatabaseProvisioner clientDatabaseProvisioner, IUserResourceRoleBLL userResourceRoleBLL, EventLogPublish eventLogPublish, ILogger logger) { _dal = dal; _clientDatabaseBLL = clientDatabaseBLL; _dbServerBLL = dbServerBLL; _targetDbExecutor = targetDbExecutor; _ddlScriptBLL = ddlScriptBLL; _clientDatabaseProvisioner = clientDatabaseProvisioner; _userResourceRoleBLL = userResourceRoleBLL; _eventLogPublish = eventLogPublish; _logger = logger; } public async Task GetById(int changeRequestId, LoginDTO login, CancellationToken ct) { try { var cr = await _dal.GetById(changeRequestId, login, ct).ConfigureAwait(false); // Real-time DDL progress (§37 Decision 2) — only meaningful, and only computed, for // ClientProvisioning CRs (the only category a caller polls for live provisioning // status). Two extra reads, bounded by this CR's own linked-script count — not a hot // path, called at most once per Entitlement onboarding poll/progress wake-up. if (cr is not null && cr.Category == ChangeRequestCategory.ClientProvisioning) { var linkedScripts = (await _dal.GetLinkedApprovedDdlScripts(changeRequestId, login, ct).ConfigureAwait(false)).ToList(); cr.ScriptsTotal = linkedScripts.Count; cr.ScriptsDone = linkedScripts.Count == 0 ? 0 : await _ddlScriptBLL.CountExecutedScriptsAsync( cr.ClientDatabaseId, linkedScripts.Select(s => s.DdlScriptId), login, ct).ConfigureAwait(false); } return cr; } catch (Exception ex) { _logger.LogError(ex, "GetById failed for ChangeRequestId {Id}", changeRequestId); throw; } } public async Task GetList(ChangeRequestListCriteria criteria, int pageOffset, int pageSize, LoginDTO login, CancellationToken ct) { try { var (rows, total) = await _dal.GetList(criteria, pageOffset, pageSize, login, ct).ConfigureAwait(false); return JsonSerializer.Serialize(new { rows, total }); } catch (Exception ex) { _logger.LogError(ex, "GetList failed"); throw; } } public async Task GetSelectList(int firstNumber, int maxResult, LoginDTO login, CancellationToken ct) { try { var rows = await _dal.GetSelectList(firstNumber, maxResult, login, ct).ConfigureAwait(false); return JsonSerializer.Serialize(rows); } catch (Exception ex) { _logger.LogError(ex, "GetSelectList failed"); throw; } } public async Task Save(ChangeRequestDTO dto, LoginDTO login, CancellationToken ct) { try { if (string.IsNullOrWhiteSpace(dto.Title)) throw new InvalidOperationException("Title is required."); if (dto.ClientDatabaseId == 0) throw new InvalidOperationException("Target Client Database is required."); if (dto.ChangeRequestId == 0) { int newId = await _dal.Insert(dto, login, ct).ConfigureAwait(false); await _dal.AppendTimeline(newId, CrTimelineEventType.Created, null, login.UserId, login, ct) .ConfigureAwait(false); return $"{SuccessResponse.SaveSuccessMessage} {newId}"; } var existing = await _dal.GetById(dto.ChangeRequestId, login, ct).ConfigureAwait(false); if (existing is null) throw new InvalidOperationException("Change Request not found."); if (existing.CrStatus != ChangeRequestStatus.Draft) throw new InvalidOperationException("Only Draft Change Requests can be edited."); await _dal.Update(dto, login, ct).ConfigureAwait(false); return SuccessResponse.UpdateSuccess; } catch (Exception ex) { _logger.LogError(ex, "Save failed for ChangeRequest Title {Title}", dto.Title); throw; } } public async Task Delete(int changeRequestId, LoginDTO login, CancellationToken ct) { try { var existing = await _dal.GetById(changeRequestId, login, ct).ConfigureAwait(false); if (existing is null) throw new InvalidOperationException("Change Request not found."); if (existing.CrStatus != ChangeRequestStatus.Draft) throw new InvalidOperationException("Only Draft Change Requests can be deleted."); await _dal.Delete(changeRequestId, login, ct).ConfigureAwait(false); return SuccessResponse.DeleteSuccessMessage; } catch (Exception ex) { _logger.LogError(ex, "Delete failed for ChangeRequestId {Id}", changeRequestId); throw; } } public async Task Submit(int changeRequestId, LoginDTO login, CancellationToken ct) { try { await _dal.UpdateStatus(changeRequestId, newStatus: (byte)ChangeRequestStatus.Submitted, expectedStatus: (byte)ChangeRequestStatus.Draft, performedById: login.UserId, comment: null, login, ct).ConfigureAwait(false); await _dal.AppendTimeline(changeRequestId, CrTimelineEventType.Submitted, null, login.UserId, login, ct).ConfigureAwait(false); return SuccessResponse.UpdateSuccess; } catch (Exception ex) { _logger.LogError(ex, "Submit failed for ChangeRequestId {Id}", changeRequestId); throw; } } public async Task Approve(int changeRequestId, string? comment, LoginDTO login, CancellationToken ct) { try { var existing = await _dal.GetById(changeRequestId, login, ct).ConfigureAwait(false); if (existing is null) throw new InvalidOperationException("Change Request not found."); if (existing.CrStatus is not (ChangeRequestStatus.Submitted or ChangeRequestStatus.UnderReview)) throw new InvalidOperationException("Change Request must be Submitted or Under Review to approve."); await _dal.UpdateStatus(changeRequestId, newStatus: (byte)ChangeRequestStatus.Approved, expectedStatus: (byte)existing.CrStatus, performedById: login.UserId, comment: comment, login, ct).ConfigureAwait(false); await _dal.AppendTimeline(changeRequestId, CrTimelineEventType.Approved, comment, login.UserId, login, ct).ConfigureAwait(false); return SuccessResponse.UpdateSuccess; } catch (Exception ex) { _logger.LogError(ex, "Approve failed for ChangeRequestId {Id}", changeRequestId); throw; } } public async Task Reject(int changeRequestId, string reason, LoginDTO login, CancellationToken ct) { try { if (string.IsNullOrWhiteSpace(reason)) throw new InvalidOperationException("Rejection reason is required."); var existing = await _dal.GetById(changeRequestId, login, ct).ConfigureAwait(false); if (existing is null) throw new InvalidOperationException("Change Request not found."); if (existing.CrStatus is not (ChangeRequestStatus.Submitted or ChangeRequestStatus.UnderReview)) throw new InvalidOperationException("Change Request must be Submitted or Under Review to reject."); await _dal.UpdateStatus(changeRequestId, newStatus: (byte)ChangeRequestStatus.Rejected, expectedStatus: (byte)existing.CrStatus, performedById: login.UserId, comment: reason, login, ct).ConfigureAwait(false); await _dal.AppendTimeline(changeRequestId, CrTimelineEventType.Rejected, reason, login.UserId, login, ct).ConfigureAwait(false); return SuccessResponse.UpdateSuccess; } catch (Exception ex) { _logger.LogError(ex, "Reject failed for ChangeRequestId {Id}", changeRequestId); throw; } } public async Task Return(int changeRequestId, string? comment, LoginDTO login, CancellationToken ct) { try { var existing = await _dal.GetById(changeRequestId, login, ct).ConfigureAwait(false); if (existing is null) throw new InvalidOperationException("Change Request not found."); if (existing.CrStatus is not (ChangeRequestStatus.Submitted or ChangeRequestStatus.UnderReview)) throw new InvalidOperationException("Change Request must be Submitted or Under Review to return."); await _dal.UpdateStatus(changeRequestId, newStatus: (byte)ChangeRequestStatus.Draft, expectedStatus: (byte)existing.CrStatus, performedById: login.UserId, comment: comment, login, ct).ConfigureAwait(false); await _dal.AppendTimeline(changeRequestId, CrTimelineEventType.Returned, comment, login.UserId, login, ct).ConfigureAwait(false); return SuccessResponse.UpdateSuccess; } catch (Exception ex) { _logger.LogError(ex, "Return failed for ChangeRequestId {Id}", changeRequestId); throw; } } public async Task Execute(int changeRequestId, LoginDTO login, CancellationToken ct) { try { var cr = await _dal.GetById(changeRequestId, login, ct).ConfigureAwait(false); if (cr is null) throw new InvalidOperationException("Change Request not found."); if (cr.CrStatus != ChangeRequestStatus.Approved) throw new InvalidOperationException("Only Approved Change Requests can be executed."); // Synchronous path — the calling human IS the acting user, no queue indirection. await ExecuteCoreAsync(cr, login.UserId, login, ct).ConfigureAwait(false); return SuccessResponse.UpdateSuccess; } catch (Exception ex) { _logger.LogError(ex, "Execute failed for ChangeRequestId {Id}", changeRequestId); throw; } } // Queues an already-Approved CR for background execution instead of running it inline on the // calling request (tracker §37) — the actual work is identical (ExecuteCoreAsync, below), only // WHEN it runs changes. Returns immediately; RolloutCampaignWorkerJob's own Quartz-job // precedent (§39) is mirrored by a new ChangeRequestExecutionWorkerJob that calls // ProcessQueuedExecutionsAsync on a timer. public async Task QueueExecutionAsync(int changeRequestId, LoginDTO login, CancellationToken ct) { try { var cr = await _dal.GetById(changeRequestId, login, ct).ConfigureAwait(false); if (cr is null) throw new InvalidOperationException("Change Request not found."); if (cr.CrStatus != ChangeRequestStatus.Approved) throw new InvalidOperationException("Only Approved Change Requests can be executed."); // Stamps EXECUTEDBYID with the queuing human's own id now, before execution runs — // ExecuteCoreAsync reads this back as the acting user when the background worker // (its own service LoginDTO, not a real person) later calls ExecuteCoreAsync for // this CR, so the per-user role-grant check authorizes the right identity. await _dal.UpdateExecutionStatusAsync( changeRequestId, (byte)ChangeRequestExecutionStatus.Queued, login, ct, queuedById: login.UserId).ConfigureAwait(false); await _dal.AppendTimeline(changeRequestId, CrTimelineEventType.Queued, "Queued for background execution.", login.UserId, login, ct).ConfigureAwait(false); return "Change Request queued for background execution."; } catch (Exception ex) { _logger.LogError(ex, "QueueExecutionAsync failed for ChangeRequestId {Id}", changeRequestId); throw; } } // Called by ChangeRequestExecutionWorkerJob on a timer. Picks up every CR left Queued by // QueueExecutionAsync, runs each through the exact same ExecuteCoreAsync path Execute() uses // inline — one shared decision point, not two independently-maintained copies of "what does // executing a CR mean." One CR's failure doesn't abort the batch (per-item try/catch), same // posture as every other worker job in this codebase (RolloutCampaignWorkerJob, // PayOrderReconciliationJob, KillSwitchMonitorJob). public async Task ProcessQueuedExecutionsAsync(LoginDTO login, CancellationToken ct) { var queued = await _dal.GetQueuedForExecutionAsync(login, ct).ConfigureAwait(false); int processed = 0; foreach (var cr in queued) { ct.ThrowIfCancellationRequested(); try { await _dal.UpdateExecutionStatusAsync( cr.ChangeRequestId, (byte)ChangeRequestExecutionStatus.Running, login, ct).ConfigureAwait(false); // The worker's own LoginDTO/UserId is a service account, never the right // identity for a per-user role-grant check — use the human who queued this CR // (stamped into ExecutedById at QueueExecutionAsync time), falling back to the // worker's id only for a CR queued before this stamping existed (default -1). int actingUserId = cr.ExecutedById > 0 ? cr.ExecutedById : login.UserId; await ExecuteCoreAsync(cr, actingUserId, login, ct).ConfigureAwait(false); await _dal.UpdateExecutionStatusAsync( cr.ChangeRequestId, (byte)ChangeRequestExecutionStatus.Done, login, ct).ConfigureAwait(false); processed++; } catch (Exception ex) { _logger.LogError(ex, "ProcessQueuedExecutionsAsync failed for ChangeRequestId {Id}", cr.ChangeRequestId); await _dal.UpdateExecutionStatusAsync( cr.ChangeRequestId, (byte)ChangeRequestExecutionStatus.Failed, login, ct).ConfigureAwait(false); // A failed background execution never reaches ExecuteCoreAsync's own // "Executed" completion publish (the exception is thrown before that line), so // without this a caller polling/subscribing for completion (e.g. Entitlement's // onboarding orchestrator) would wait forever with no signal at all. Same topic // as the success path — the subscriber already re-verifies live status rather // than trusting the event payload's shape (see // ClientProvisioningExecutedSubscriber's own doc comment), so this is only ever // used as a "check now" wake-up; ExecutionStatus=Failed (just set above, and now // visible via GetChangeRequestById) is what actually lets the caller distinguish // "still pending" from "genuinely failed." Fire-and-forget/non-fatal — // EventLogPublish already swallows and logs its own failures. if (cr.Category == ChangeRequestCategory.ClientProvisioning) { await _eventLogPublish.PublishEventLogAsync( "SqlWorkbench ClientProvisioning Change Request Execution Failed", new { cr.ChangeRequestId, cr.ClientDatabaseId }, EventTypeConstant.SWCLIENTPROVISIONINGEXECUTEDEVENTTYPEID, cr.ClientDatabaseId, login, topic: ClientProvisioningExecutedTopic, ct: ct).ConfigureAwait(false); } } } return processed; } // The actual execution logic — extracted verbatim from Execute()'s own body (tracker §37) so // both the synchronous endpoint and the background worker share one implementation instead of // two independently-maintained copies of "what does executing a CR mean." actingUserId is the // real human whose per-user role grant (SW.MSWUSERRESOURCEROLE) should cap the connection's // privilege — the calling human directly for the synchronous path, or the human who queued // the CR for the background-worker path (never the worker's own service-account login). private async Task ExecuteCoreAsync(ChangeRequestDTO cr, int actingUserId, LoginDTO login, CancellationToken ct) { int changeRequestId = cr.ChangeRequestId; { var clientDb = await _clientDatabaseBLL.GetById(cr.ClientDatabaseId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"Client database {cr.ClientDatabaseId} not found."); var server = await _dbServerBLL.GetByIdWithCredentials(clientDb.DbServerId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"DB server {clientDb.DbServerId} not found."); // ClientProvisioning is the one category that runs BEFORE the target database can be // connected to at all — it doesn't exist yet. Create it (+ its three contained users // + DB-level defaults) here, then fall through to the normal DDL path below, which is // useful for snapshot-mode if any delta scripts are linked on top of the restored // baseline, and a harmless no-op for scripts-mode if none are linked (ProvisionClientDatabase // already deployed the baseline package as part of creation). if (cr.Category == ChangeRequestCategory.ClientProvisioning) { await ExecuteClientProvisioningCr(cr, clientDb, server, login, ct).ConfigureAwait(false); } // ClientProvisioning creates the database + its contained users moments ago (above) — // there's no meaningful "least privilege" distinction left to make for the delta // scripts that immediately follow, so it keeps using the DbServer admin credential. // Every other category picks the least-privilege contained user for the operation: // DDL needs schema-alter rights the _app login doesn't have (Dba), plain data changes // use _app (read-write, no schema rights), and SELECT/Other use _readonly — CAPPED // (never loosened) by actingUserId's own per-resource role grant, if one is // configured. No grant configured at all is fail-open (uncapped) — this is a // brand-new authorization layer landing in a module where every endpoint is // currently AllowAnonymous(), so fail-closed-by-default would immediately lock out // every current user with no grants configured yet. string connStr; if (cr.Category == ChangeRequestCategory.ClientProvisioning) { connStr = _targetDbExecutor.BuildConnectionString(server, clientDb.DatabaseName); } else { var requestedRole = RoleForQueryType(cr.QueryType); var grantedRole = await _userResourceRoleBLL.GetEffectiveRole( actingUserId, clientDb.DbServerId, clientDb.ClientDbId, login, ct).ConfigureAwait(false); var role = grantedRole is ClientDbLoginRole granted ? (ClientDbLoginRole)Math.Max((byte)requestedRole, (byte)granted) : requestedRole; connStr = await _targetDbExecutor.BuildConnectionStringForRoleAsync( server, clientDb, role, ct).ConfigureAwait(false); } if (cr.Category == ChangeRequestCategory.ClientProvisioning) { // The physical database (+ baseline schema, for scripts-mode) already exists at // this point — only run the normal DDL path if the CR also links delta scripts on // top of it (e.g. post-restore fixups for snapshot-mode). Unlike an ordinary DDL // CR, having none linked here is expected and not an error. var linkedScripts = await _dal.GetLinkedApprovedDdlScripts(cr.ChangeRequestId, login, ct).ConfigureAwait(false); if (linkedScripts.Any()) await ExecuteDdlCr(cr, connStr, clientDb.TenantId, login, ct).ConfigureAwait(false); } else if (cr.QueryType == ChangeRequestQueryType.DDL) { await ExecuteDdlCr(cr, connStr, clientDb.TenantId, login, ct).ConfigureAwait(false); } else { await ExecuteDmlCr(cr, connStr, login, ct).ConfigureAwait(false); } // performedById is actingUserId, not login.UserId — for the background-worker path // this keeps EXECUTEDBYID correctly attributed to the human who queued the CR // (already stamped there at queue time), not overwritten with the worker's own id. await _dal.UpdateStatus(changeRequestId, newStatus: (byte)ChangeRequestStatus.Executed, expectedStatus: (byte)ChangeRequestStatus.Approved, performedById: actingUserId, comment: null, login, ct).ConfigureAwait(false); await _dal.AppendTimeline(changeRequestId, CrTimelineEventType.Executed, null, actingUserId, login, ct).ConfigureAwait(false); // Completion signal for ClientProvisioning CRs — this is the one category a // cross-module caller genuinely needs to react to (a physical database now exists). // Fire-and-forget/non-fatal: EventLogPublish already swallows and logs its own // failures, so a publish failure must never fail an otherwise-successful Execute. if (cr.Category == ChangeRequestCategory.ClientProvisioning) { await _eventLogPublish.PublishEventLogAsync( "SqlWorkbench ClientProvisioning Change Request Executed", new { ChangeRequestId = changeRequestId, cr.ClientDatabaseId }, EventTypeConstant.SWCLIENTPROVISIONINGEXECUTEDEVENTTYPEID, cr.ClientDatabaseId, login, topic: ClientProvisioningExecutedTopic, ct: ct).ConfigureAwait(false); } // Same completion signal, generalized to every other category — previously silent, // so nothing could react to an ordinary post-provisioning schema/data change landing // on an already-live client database. Separate topic (above), same minimal payload. else if (cr.Category is ChangeRequestCategory.SchemaChange or ChangeRequestCategory.DataMigration) { await _eventLogPublish.PublishEventLogAsync( "SqlWorkbench Change Request Applied", new { ChangeRequestId = changeRequestId, cr.ClientDatabaseId }, EventTypeConstant.SWCHANGEREQUESTAPPLIEDEVENTTYPEID, cr.ClientDatabaseId, login, topic: ChangeRequestAppliedTopic, ct: ct).ConfigureAwait(false); } } } private async Task ExecuteClientProvisioningCr( ChangeRequestDTO cr, SwDAL.DTO.ClientDatabase.ClientDatabaseDTO clientDb, SwDAL.DTO.DbServer.DbServerDTO server, LoginDTO login, CancellationToken ct) { if (cr.ProvisioningMode is null) throw new InvalidOperationException("ProvisioningMode is required for a ClientProvisioning Change Request."); if (cr.ProvisioningMode == ProvisioningMode.FromSnapshot) { if (string.IsNullOrWhiteSpace(cr.TemplateBackupRef)) throw new InvalidOperationException("TemplateBackupRef is required for ProvisioningMode.FromSnapshot."); await _clientDatabaseProvisioner .CreateFromSnapshotAsync(clientDb, server, cr.TemplateBackupRef, login, ct) .ConfigureAwait(false); } else { if (cr.BaselineUpgradePackageId is null or 0) throw new InvalidOperationException( "BaselineUpgradePackageId is required for ProvisioningMode.FromScripts."); await _clientDatabaseProvisioner .CreateFromScriptsAsync(clientDb, server, cr.BaselineUpgradePackageId.Value, login, ct) .ConfigureAwait(false); } await _clientDatabaseBLL.UpdateStatusAsync( clientDb.ClientDbId, (byte)ClientDbStatus.Active, login, ct).ConfigureAwait(false); } private async Task ExecuteDdlCr(ChangeRequestDTO cr, string connStr, int tenantId, LoginDTO login, CancellationToken ct) { var scripts = await _dal.GetLinkedApprovedDdlScripts(cr.ChangeRequestId, login, ct).ConfigureAwait(false); var scriptList = scripts.ToList(); if (scriptList.Count == 0) throw new InvalidOperationException("No approved DDL scripts are linked to this Change Request."); for (var i = 0; i < scriptList.Count; i++) { var script = scriptList[i]; byte execStatus = (byte)DdlExecStatus.Failed; string? errorMessage = null; try { string appliedSql = DdlScriptTemplating.SubstituteTenantId(script.SqlScript, tenantId); await _targetDbExecutor.ExecuteScriptAsync(connStr, appliedSql, ct).ConfigureAwait(false); execStatus = (byte)DdlExecStatus.Success; } catch (Exception ex) { execStatus = (byte)DdlExecStatus.Failed; errorMessage = ex.Message; _logger.LogError(ex, "DDL script {ScriptId} failed for CR {CRId}", script.DdlScriptId, cr.ChangeRequestId); throw; } finally { await _ddlScriptBLL.LogExecutionAsync( script.DdlScriptId, cr.ClientDatabaseId, execStatus, errorMessage, login, ct) .ConfigureAwait(false); // Batched, not per-script — a 1,500+-script baseline batch (§36.13) publishing on // every single script would flood Dapr for no real UI benefit. Only for // ClientProvisioning CRs (the only category anyone polls for live progress); // fire-and-forget/non-fatal like every other event-publish in this method. if (cr.Category == ChangeRequestCategory.ClientProvisioning && ((i + 1) % ScriptsProgressPublishBatchSize == 0 || i == scriptList.Count - 1)) { await _eventLogPublish.PublishEventLogAsync( "SqlWorkbench ClientProvisioning DDL Progress", new { ChangeRequestId = cr.ChangeRequestId, cr.ClientDatabaseId }, EventTypeConstant.SWCLIENTPROVISIONINGPROGRESSEVENTTYPEID, cr.ClientDatabaseId, login, topic: ClientProvisioningProgressTopic, ct: ct).ConfigureAwait(false); } } } } private static ClientDbLoginRole RoleForQueryType(ChangeRequestQueryType queryType) => queryType switch { ChangeRequestQueryType.DDL => ClientDbLoginRole.Dba, ChangeRequestQueryType.Insert or ChangeRequestQueryType.Update or ChangeRequestQueryType.Delete => ClientDbLoginRole.App, ChangeRequestQueryType.Select or ChangeRequestQueryType.Other => ClientDbLoginRole.ReadOnly, _ => throw new ArgumentOutOfRangeException(nameof(queryType), queryType, "Unhandled ChangeRequestQueryType.") }; private async Task ExecuteDmlCr(ChangeRequestDTO cr, string connStr, LoginDTO login, CancellationToken ct) { if (string.IsNullOrWhiteSpace(cr.Sql)) throw new InvalidOperationException("Change Request SQL is empty. Cannot execute."); if (cr.QueryType is ChangeRequestQueryType.Select or ChangeRequestQueryType.Other) { await _targetDbExecutor.ExecuteQueryAsync(connStr, cr.Sql, [], ct).ConfigureAwait(false); } else { await _targetDbExecutor.ExecuteDmlAsync(connStr, cr.Sql, [], ct).ConfigureAwait(false); } } }