diff --git a/Core/Resgrid.Model/Repositories/ISystemOperationRequestsRepository.cs b/Core/Resgrid.Model/Repositories/ISystemOperationRequestsRepository.cs new file mode 100644 index 000000000..a01e25556 --- /dev/null +++ b/Core/Resgrid.Model/Repositories/ISystemOperationRequestsRepository.cs @@ -0,0 +1,35 @@ +using System; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; + +namespace Resgrid.Model.Repositories +{ + /// System operation requests (SystemOperationRequests, M0267). Every state change is a conditional write on the current status. + public interface ISystemOperationRequestsRepository : IRepository + { + /// Newest first. + Task> GetRecentAsync(int take); + + /// A Pending request for the same operation and target, if one is already waiting. + Task GetPendingAsync(int operationType, int? targetDepartmentId); + + /// + /// Moves the oldest Pending request to Running for this worker and returns it, or null when nothing is waiting. + /// Two workers never claim the same row: the claim only succeeds while the row is still Pending. + /// + Task ClaimNextPendingAsync(string workerName, DateTime now, CancellationToken cancellationToken = default); + + /// Refreshes HeartbeatOn, and Progress when one is given, on a Running request. + Task HeartbeatAsync(string requestId, string progress, DateTime now, CancellationToken cancellationToken = default); + + /// Records the outcome of a Running request. + Task FinishAsync(string requestId, int status, string result, DateTime now, CancellationToken cancellationToken = default); + + /// Withdraws a request nobody has claimed yet. + Task CancelPendingAsync(string requestId, string cancelledBy, DateTime now, CancellationToken cancellationToken = default); + + /// Fails every Running request whose heartbeat (or start) is older than the cutoff: its worker stopped. + Task FailAbandonedAsync(DateTime heartbeatBefore, string result, DateTime now, CancellationToken cancellationToken = default); + } +} diff --git a/Core/Resgrid.Model/Services/ISystemOperationsService.cs b/Core/Resgrid.Model/Services/ISystemOperationsService.cs new file mode 100644 index 000000000..7634bac81 --- /dev/null +++ b/Core/Resgrid.Model/Services/ISystemOperationsService.cs @@ -0,0 +1,80 @@ +using System; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; + +namespace Resgrid.Model.Services +{ + /// + /// Queues and tracks on-demand system operations (): BackOffice requests them, the + /// worker (command 76) claims and runs them. Also owns the cache sentinel the worker uses to notice Redis lost its data. + /// + public interface ISystemOperationsService + { + /// + /// Queues an operation. A request that matches one already waiting (same operation and target) is not queued twice; + /// the waiting one is returned with false. + /// + Task RequestAsync(SystemOperationTypes operationType, int? targetDepartmentId, SystemOperationSources source, + string requestedBy, string reason, CancellationToken cancellationToken = default); + + /// Newest first. + Task> GetRecentRequestsAsync(int count = 50); + + Task GetRequestByIdAsync(string requestId); + + /// Withdraws a request the worker has not claimed yet. False when it is already running or finished. + Task CancelPendingRequestAsync(string requestId, string cancelledBy, CancellationToken cancellationToken = default); + + /// Worker: claims the oldest waiting request, or returns null. + Task ClaimNextRequestAsync(string workerName, CancellationToken cancellationToken = default); + + /// Worker: keeps a running request alive, optionally with a new progress line. + Task ReportProgressAsync(string requestId, string progress, CancellationToken cancellationToken = default); + + /// Worker: records how a running request ended. + Task CompleteRequestAsync(string requestId, bool succeeded, string result, CancellationToken cancellationToken = default); + + /// Worker: fails running requests whose worker stopped heartbeating (it crashed or was redeployed mid-run). + Task FailAbandonedRequestsAsync(CancellationToken cancellationToken = default); + + /// + /// Worker: touches the cache sentinel and returns true when it was missing, meaning the cache came back without its + /// data (or this is the first check ever). False when the cache is off or unreachable: nothing can be concluded then. + /// + Task DetectCacheDataLossAsync(); + + /// Whether the cache is reachable, and since when it has held its data (when the sentinel was written). + Task GetCacheStatusAsync(); + + /// + /// Drops the department's cached entries so they reload from the database. Best effort per cache group; returns the + /// groups that failed (empty when all were cleared). + /// + Task> ClearDepartmentCachesAsync(int departmentId); + } + + public sealed class SystemOperationRequestResult + { + public bool Created { get; set; } + + /// The queued request, or the one already waiting. + public SystemOperationRequest Request { get; set; } + + /// Why nothing was queued (unknown operation, bad target); null on success. + public string Error { get; set; } + + public bool Succeeded => Error == null && Request != null; + } + + public sealed class SystemOperationsCacheStatus + { + /// SystemBehaviorConfig.CacheEnabled as this process sees it. + public bool CacheEnabled { get; set; } + + public bool Connected { get; set; } + + /// When the worker first found the cache without its sentinel: it has held its data since then. Null when unknown. + public DateTime? DataPresentSinceUtc { get; set; } + } +} diff --git a/Core/Resgrid.Model/SystemOperationCatalog.cs b/Core/Resgrid.Model/SystemOperationCatalog.cs new file mode 100644 index 000000000..f7bb6d4be --- /dev/null +++ b/Core/Resgrid.Model/SystemOperationCatalog.cs @@ -0,0 +1,164 @@ +using System.Collections.Generic; +using System.Linq; + +namespace Resgrid.Model +{ + /// How a system operation is grouped on the BackOffice System Operations page. + public enum SystemOperationCategories + { + /// Rebuilds or drops state the platform keeps in Redis (or another cache). + CachedState = 1, + + /// Runs one of the worker's daily jobs now instead of waiting for its schedule. + ScheduledJob = 2 + } + + /// What one value does, for the request validation and the BackOffice page. + public sealed class SystemOperationDescriptor + { + public SystemOperationDescriptor(SystemOperationTypes type, SystemOperationCategories category, string name, string description, + bool supportsDepartmentScope, string schedule, int? workerCommandId, bool deletesData) + { + Type = type; + Category = category; + Name = name; + Description = description; + SupportsDepartmentScope = supportsDepartmentScope; + Schedule = schedule; + WorkerCommandId = workerCommandId; + DeletesData = deletesData; + } + + public SystemOperationTypes Type { get; } + + public SystemOperationCategories Category { get; } + + public string Name { get; } + + public string Description { get; } + + /// True when a request may name one department; without one it covers every department. + public bool SupportsDepartmentScope { get; } + + /// When the worker runs it on its own (UTC), or null when it only ever runs on request. + public string Schedule { get; } + + /// The scheduled worker command this operation runs early, if any. + public int? WorkerCommandId { get; } + + /// Deletes or purges data when it runs (the same data its schedule would delete). + public bool DeletesData { get; } + } + + /// + /// Every operation staff can request from BackOffice -> System Operations. The worker (command 76) runs each one + /// through SystemOperationRunner; a scheduled job runs exactly the handler its schedule runs, so an early run + /// behaves like the nightly one. + /// + public static class SystemOperationCatalog + { + public static IReadOnlyList All { get; } = new List + { + new SystemOperationDescriptor(SystemOperationTypes.RebuildSecurityMatrices, SystemOperationCategories.CachedState, + "Rebuild security matrices", + "Rebuilds the Redis visibility matrices: who may see each unit, each person, and their locations. Without a matrix, " + + "visibility checks answer from the permission rows (correct, but slower) and realtime location fan-out recomputes per " + + "department. The worker queues this for every department on its own when it finds Redis came back empty.", + supportsDepartmentScope: true, schedule: "Daily 02:00 UTC", workerCommandId: 15, deletesData: false), + + new SystemOperationDescriptor(SystemOperationTypes.ClearDepartmentCaches, SystemOperationCategories.CachedState, + "Clear department caches", + "Drops cached department data (department, members, personnel names, plan, groups, call priorities, action logs, " + + "custom states, latest statuses, feature-flag overrides) so it reloads from the database. Use it when Redis was " + + "restored from a snapshot, or was unreachable while data changed, and may now serve stale entries.", + supportsDepartmentScope: true, schedule: null, workerCommandId: null, deletesData: false), + + new SystemOperationDescriptor(SystemOperationTypes.RefreshTtsStaticPrompts, SystemOperationCategories.CachedState, + "Refresh TTS static prompts", + "Asks the TTS service to regenerate its static voice prompts (cached in Redis and object storage). Runs only where the " + + "TTS service URL and admin key are configured.", + supportsDepartmentScope: false, schedule: "Hourly (TtsConfig.StaticPromptRefreshIntervalMinutes)", workerCommandId: 18, deletesData: false), + + new SystemOperationDescriptor(SystemOperationTypes.PendingDepartmentDeletions, SystemOperationCategories.ScheduledJob, + "Pending department deletions (System SQL Queue)", + "Deletes the departments whose deletion request has passed its waiting period. Departments still inside the waiting " + + "period are not touched.", + supportsDepartmentScope: false, schedule: "Daily 03:00 UTC", workerCommandId: 14, deletesData: true), + + new SystemOperationDescriptor(SystemOperationTypes.ReportingRollup, SystemOperationCategories.ScheduledJob, + "Reporting rollup", + "Writes the previous UTC day's reporting rollup rows for every department.", + supportsDepartmentScope: false, schedule: "Daily 03:30 UTC", workerCommandId: 21, deletesData: false), + + new SystemOperationDescriptor(SystemOperationTypes.UnitTrackingLocationRetention, SystemOperationCategories.ScheduledJob, + "Unit tracking location retention", + "Purges unit tracking locations older than the retention window. Runs only where the retention worker is enabled.", + supportsDepartmentScope: false, schedule: "Daily 04:30 UTC (UnitTrackingConfig.LocationRetentionHourUtc)", workerCommandId: 24, deletesData: true), + + new SystemOperationDescriptor(SystemOperationTypes.ChatRetention, SystemOperationCategories.ScheduledJob, + "Chat retention", + "Purges chat messages past each department's or channel's retention window, and expired chat exports.", + supportsDepartmentScope: false, schedule: "Daily 04:45 UTC", workerCommandId: 25, deletesData: true), + + new SystemOperationDescriptor(SystemOperationTypes.BidExpiration, SystemOperationCategories.ScheduledJob, + "Bid expiration", + "Moves submitted bids past their valid-until date to Expired.", + supportsDepartmentScope: false, schedule: "Daily 04:00 UTC", workerCommandId: 31, deletesData: false), + + new SystemOperationDescriptor(SystemOperationTypes.DeploymentFinanceReminder, SystemOperationCategories.ScheduledJob, + "Deployment finance reminder", + "Sends department admins the daily deployment billing and Cal OES MARS reminder digests. A department this worker " + + "process already reminded today is skipped; a restarted worker can send a second digest the same day.", + supportsDepartmentScope: false, schedule: "Daily 04:15 UTC", workerCommandId: 32, deletesData: false), + + new SystemOperationDescriptor(SystemOperationTypes.ComplianceExpiry, SystemOperationCategories.ScheduledJob, + "Compliance expiry", + "Expires lapsed service contracts and sends the contract and compliance-document expiry notices. Notices go once per " + + "day per worker process; a restarted worker can send them again the same day.", + supportsDepartmentScope: false, schedule: "Daily 04:30 UTC", workerCommandId: 33, deletesData: false), + + new SystemOperationDescriptor(SystemOperationTypes.RmsDueStateEvaluation, SystemOperationCategories.ScheduledJob, + "RMS due state evaluation", + "Evaluates RMS record due states (overdue records, due inspections, overdue violations, permit expiry). It emits from " + + "the persisted due-state rows, so a repeated run stays quiet.", + supportsDepartmentScope: false, schedule: "Daily 04:00 UTC", workerCommandId: 42, deletesData: false), + + new SystemOperationDescriptor(SystemOperationTypes.RmsRetentionAndPurge, SystemOperationCategories.ScheduledJob, + "RMS retention and purge", + "Applies RMS retention and legal holds, purges eligible records and attachments, and rescans attachments the scanner " + + "could not reach at upload.", + supportsDepartmentScope: false, schedule: "Daily 03:30 UTC", workerCommandId: 43, deletesData: true), + + new SystemOperationDescriptor(SystemOperationTypes.PayDataReportingReadiness, SystemOperationCategories.ScheduledJob, + "Pay data reporting readiness", + "Purges expired pay data export artifacts and, during the filing season, sends the readiness digest (once per day per " + + "worker process).", + supportsDepartmentScope: false, schedule: "Daily 04:45 UTC", workerCommandId: 49, deletesData: true), + + new SystemOperationDescriptor(SystemOperationTypes.ProtectedWorkflowSweep, SystemOperationCategories.ScheduledJob, + "Protected workflow sweep", + "Expires, revokes and suspends ADP protected workflow releases, and sends the 30- and 7-day expiry notices (each notice " + + "is recorded, so it goes once).", + supportsDepartmentScope: false, schedule: "Daily 05:15 UTC", workerCommandId: 71, deletesData: false), + + new SystemOperationDescriptor(SystemOperationTypes.Utf8Cleanup, SystemOperationCategories.ScheduledJob, + "UTF-8 data cleanup", + "Repairs text that would block a SQL Server to PostgreSQL move (NUL characters, unpaired surrogates, Windows-1252 " + + "mojibake). Runs only where the cleanup is enabled.", + supportsDepartmentScope: false, schedule: "Daily 04:00 UTC (SystemBehaviorConfig.Utf8CleanupHourUtc)", workerCommandId: 22, deletesData: false) + }; + + private static readonly Dictionary ByType = All.ToDictionary(x => x.Type); + + /// The descriptor, or null for a value the catalog does not know. + public static SystemOperationDescriptor Get(SystemOperationTypes type) + { + return ByType.TryGetValue(type, out var descriptor) ? descriptor : null; + } + + public static string GetName(int operationType) + { + return Get((SystemOperationTypes)operationType)?.Name ?? $"Operation {operationType}"; + } + } +} diff --git a/Core/Resgrid.Model/SystemOperationRequest.cs b/Core/Resgrid.Model/SystemOperationRequest.cs new file mode 100644 index 000000000..0a86c7aa5 --- /dev/null +++ b/Core/Resgrid.Model/SystemOperationRequest.cs @@ -0,0 +1,103 @@ +using System; +using System.Collections.Generic; +using System.ComponentModel.DataAnnotations; +using System.ComponentModel.DataAnnotations.Schema; +using Newtonsoft.Json; + +namespace Resgrid.Model +{ + /// + /// A request for the worker (command 76) to run one operation now (registry M0267). + /// BackOffice inserts Pending rows; the worker claims the oldest, keeps HeartbeatOn fresh while it runs, and records + /// the outcome. The table is the durable trigger path on purpose: the bus queues are not durable and the cache may be + /// the thing that just failed. A system record, not department data: TargetDepartmentId only narrows the operation. + /// + public class SystemOperationRequest : IEntity + { + public const int ReasonMaxLength = 500; + public const int ProgressMaxLength = 500; + public const int ResultMaxLength = 2000; + + [Key] + [Required] + [MaxLength(128)] + public string SystemOperationRequestId { get; set; } + + /// A value. + public int OperationType { get; set; } + + /// The one department the operation covers, or null for every department (or no department). + public int? TargetDepartmentId { get; set; } + + /// A value. + public int Status { get; set; } + + /// A value. + public int Source { get; set; } + + /// The staff member's e-mail (or subject), or "system" for an automatic request. + [Required] + [MaxLength(256)] + public string RequestedBy { get; set; } + + [MaxLength(ReasonMaxLength)] + public string Reason { get; set; } + + public DateTime RequestedOn { get; set; } + + public DateTime? StartedOn { get; set; } + + /// Refreshed by the running worker; a Running row whose heartbeat goes stale was abandoned. + public DateTime? HeartbeatOn { get; set; } + + public DateTime? CompletedOn { get; set; } + + /// Machine and process that claimed the request. + [MaxLength(256)] + public string WorkerName { get; set; } + + /// The latest progress line the operation reported while running. + [MaxLength(ProgressMaxLength)] + public string Progress { get; set; } + + /// The outcome summary, or the failure. + [MaxLength(ResultMaxLength)] + public string Result { get; set; } + + [MaxLength(256)] + public string CancelledBy { get; set; } + + [NotMapped] + [JsonIgnore] + public object IdValue + { + get => SystemOperationRequestId; + set => SystemOperationRequestId = (string)value; + } + + [NotMapped] + public string TableName => "SystemOperationRequests"; + + [NotMapped] + public string IdName => "SystemOperationRequestId"; + + [NotMapped] + public int IdType => 1; + + [NotMapped] + public IEnumerable IgnoredProperties => + new[] { "IdValue", "IdType", "TableName", "IdName" }; + + [NotMapped] + [JsonIgnore] + public SystemOperationTypes OperationTypeValue => (SystemOperationTypes)OperationType; + + [NotMapped] + [JsonIgnore] + public SystemOperationStatuses StatusValue => (SystemOperationStatuses)Status; + + [NotMapped] + [JsonIgnore] + public bool IsFinished => Status is (int)SystemOperationStatuses.Completed or (int)SystemOperationStatuses.Failed or (int)SystemOperationStatuses.Cancelled; + } +} diff --git a/Core/Resgrid.Model/SystemOperationSources.cs b/Core/Resgrid.Model/SystemOperationSources.cs new file mode 100644 index 000000000..d4c0c70d0 --- /dev/null +++ b/Core/Resgrid.Model/SystemOperationSources.cs @@ -0,0 +1,12 @@ +namespace Resgrid.Model +{ + /// Who asked for a system operation. Persisted: never renumber. + public enum SystemOperationSources + { + /// A staff member on the BackOffice System Operations page. + BackOffice = 1, + + /// The worker found the cache sentinel missing: Redis came back without its data. + CacheDataLossDetected = 2 + } +} diff --git a/Core/Resgrid.Model/SystemOperationStatuses.cs b/Core/Resgrid.Model/SystemOperationStatuses.cs new file mode 100644 index 000000000..876999f59 --- /dev/null +++ b/Core/Resgrid.Model/SystemOperationStatuses.cs @@ -0,0 +1,20 @@ +namespace Resgrid.Model +{ + /// Lifecycle of a SystemOperationRequests row. Persisted: never renumber. + public enum SystemOperationStatuses + { + /// Waiting for the worker to claim it. + Pending = 0, + + /// Claimed by a worker, which keeps HeartbeatOn fresh while it runs. + Running = 1, + + Completed = 2, + + /// The operation reported a failure, or its worker stopped heartbeating and the row was abandoned. + Failed = 3, + + /// Withdrawn by staff before a worker claimed it. + Cancelled = 4 + } +} diff --git a/Core/Resgrid.Model/SystemOperationTypes.cs b/Core/Resgrid.Model/SystemOperationTypes.cs new file mode 100644 index 000000000..7c509535c --- /dev/null +++ b/Core/Resgrid.Model/SystemOperationTypes.cs @@ -0,0 +1,56 @@ +namespace Resgrid.Model +{ + /// + /// An operation staff can ask the worker to run on demand (BackOffice -> System Operations, worker command 76), + /// stored on SystemOperationRequests.OperationType. Append-only: the value is persisted, so never renumber. + /// 1-9 rebuild cached state; 10 and up run one of the worker's daily jobs now. Every value needs a + /// entry and a case in the worker's SystemOperationRunner. + /// + public enum SystemOperationTypes + { + /// The four Redis visibility matrices per department (Security Refresh, worker 15). + RebuildSecurityMatrices = 1, + + /// Invalidates the department-scoped cache entries so they reload from the database. + ClearDepartmentCaches = 2, + + /// Regenerates the TTS service's static voice prompts (worker 18). + RefreshTtsStaticPrompts = 3, + + /// System SQL Queue (worker 14): department deletions whose waiting period has passed. + PendingDepartmentDeletions = 10, + + /// Reporting Rollup (worker 21). + ReportingRollup = 11, + + /// Unit Tracking Location Retention (worker 24). + UnitTrackingLocationRetention = 12, + + /// Chat Retention (worker 25). + ChatRetention = 13, + + /// Bid Expiration (worker 31). + BidExpiration = 14, + + /// Deployment Finance Reminder (worker 32). + DeploymentFinanceReminder = 15, + + /// Compliance Expiry (worker 33). + ComplianceExpiry = 16, + + /// RMS Due State Evaluation (worker 42). + RmsDueStateEvaluation = 17, + + /// RMS Retention And Purge (worker 43). + RmsRetentionAndPurge = 18, + + /// Pay Data Reporting Readiness (worker 49). + PayDataReportingReadiness = 19, + + /// Protected Workflow Sweep (worker 71). + ProtectedWorkflowSweep = 20, + + /// UTF-8 Data Cleanup (worker 22). + Utf8Cleanup = 21 + } +} diff --git a/Core/Resgrid.Services/ServicesModule.cs b/Core/Resgrid.Services/ServicesModule.cs index d52cfe1d1..67dd59094 100644 --- a/Core/Resgrid.Services/ServicesModule.cs +++ b/Core/Resgrid.Services/ServicesModule.cs @@ -204,6 +204,7 @@ protected override void Load(ContainerBuilder builder) builder.RegisterType().As().InstancePerLifetimeScope(); builder.RegisterType().As().InstancePerLifetimeScope(); builder.RegisterType().As().InstancePerLifetimeScope(); + builder.RegisterType().As().InstancePerLifetimeScope(); builder.RegisterType().As().InstancePerLifetimeScope(); builder.RegisterType().As().InstancePerLifetimeScope(); builder.RegisterType().As().InstancePerLifetimeScope(); diff --git a/Core/Resgrid.Services/SystemOperationsService.cs b/Core/Resgrid.Services/SystemOperationsService.cs new file mode 100644 index 000000000..e4415449a --- /dev/null +++ b/Core/Resgrid.Services/SystemOperationsService.cs @@ -0,0 +1,253 @@ +using System; +using System.Collections.Generic; +using System.Globalization; +using System.Threading; +using System.Threading.Tasks; +using Resgrid.Framework; +using Resgrid.Model; +using Resgrid.Model.Providers; +using Resgrid.Model.Repositories; +using Resgrid.Model.Services; + +namespace Resgrid.Services +{ + /// + /// The request side of BackOffice -> System Operations and the bookkeeping worker 76 does around each run. The + /// operations themselves run in the worker (SystemOperationRunner); only the department cache clear lives here, + /// because it is nothing but service invalidations. + /// + public class SystemOperationsService : ISystemOperationsService + { + /// + /// Written once into a cache that holds its data and touched by every worker check after that. Finding it missing + /// means the cache lost everything (a restart without persistence, a flush, a failover to an empty replica). + /// + public const string CacheSentinelKey = "SystemOperations_CacheSentinel"; + + /// Sliding: every check resets it, so only a cache that lost its data (or a worker down this long) misses it. + public static readonly TimeSpan CacheSentinelLifetime = TimeSpan.FromDays(30); + + /// The running worker heartbeats every minute; five missed beats in a row means it stopped. + public static readonly TimeSpan AbandonedAfter = TimeSpan.FromMinutes(5); + + public const string SystemRequester = "system"; + + private const int MaxRecentRequests = 500; + + private readonly ISystemOperationRequestsRepository _requestsRepository; + private readonly ICacheProvider _cacheProvider; + private readonly IDepartmentsService _departmentsService; + private readonly ISubscriptionsService _subscriptionsService; + private readonly IDepartmentGroupsService _departmentGroupsService; + private readonly ICallsService _callsService; + private readonly IActionLogsService _actionLogsService; + private readonly ICustomStateService _customStateService; + private readonly IUserStateService _userStateService; + private readonly IFeatureToggleService _featureToggleService; + private readonly TimeProvider _timeProvider; + + public SystemOperationsService(ISystemOperationRequestsRepository requestsRepository, ICacheProvider cacheProvider, + IDepartmentsService departmentsService, ISubscriptionsService subscriptionsService, IDepartmentGroupsService departmentGroupsService, + ICallsService callsService, IActionLogsService actionLogsService, ICustomStateService customStateService, + IUserStateService userStateService, IFeatureToggleService featureToggleService, TimeProvider timeProvider) + { + _requestsRepository = requestsRepository; + _cacheProvider = cacheProvider; + _departmentsService = departmentsService; + _subscriptionsService = subscriptionsService; + _departmentGroupsService = departmentGroupsService; + _callsService = callsService; + _actionLogsService = actionLogsService; + _customStateService = customStateService; + _userStateService = userStateService; + _featureToggleService = featureToggleService; + _timeProvider = timeProvider ?? TimeProvider.System; + } + + private DateTime UtcNow => _timeProvider.GetUtcNow().UtcDateTime; + + public async Task RequestAsync(SystemOperationTypes operationType, int? targetDepartmentId, + SystemOperationSources source, string requestedBy, string reason, CancellationToken cancellationToken = default) + { + var descriptor = SystemOperationCatalog.Get(operationType); + + if (descriptor == null) + return new SystemOperationRequestResult { Error = $"Unknown system operation {(int)operationType}." }; + + if (string.IsNullOrWhiteSpace(requestedBy)) + return new SystemOperationRequestResult { Error = "The requester is required." }; + + if (targetDepartmentId.HasValue) + { + if (!descriptor.SupportsDepartmentScope) + return new SystemOperationRequestResult { Error = $"{descriptor.Name} always covers the whole system; it cannot target one department." }; + + if (targetDepartmentId.Value <= 0 || await _departmentsService.GetDepartmentByIdAsync(targetDepartmentId.Value) == null) + return new SystemOperationRequestResult { Error = $"Department {targetDepartmentId.Value} does not exist." }; + } + + // A waiting request covers this one: it has not started, so it will read the same current state when it runs. + // A running one does not: whatever changed since it started is why someone asked again. + var pending = await _requestsRepository.GetPendingAsync((int)operationType, targetDepartmentId); + + if (pending != null) + return new SystemOperationRequestResult { Created = false, Request = pending }; + + var request = new SystemOperationRequest + { + SystemOperationRequestId = Guid.NewGuid().ToString(), + OperationType = (int)operationType, + TargetDepartmentId = targetDepartmentId, + Status = (int)SystemOperationStatuses.Pending, + Source = (int)source, + RequestedBy = Truncate(requestedBy.Trim(), 256), + Reason = Truncate(reason?.Trim(), SystemOperationRequest.ReasonMaxLength), + RequestedOn = UtcNow + }; + + request = await _requestsRepository.InsertAsync(request, cancellationToken); + + return new SystemOperationRequestResult { Created = true, Request = request }; + } + + public async Task> GetRecentRequestsAsync(int count = 50) + { + return await _requestsRepository.GetRecentAsync(Math.Clamp(count, 1, MaxRecentRequests)); + } + + public async Task GetRequestByIdAsync(string requestId) + { + if (string.IsNullOrWhiteSpace(requestId)) + return null; + + return await _requestsRepository.GetByIdAsync(requestId); + } + + public async Task CancelPendingRequestAsync(string requestId, string cancelledBy, CancellationToken cancellationToken = default) + { + if (string.IsNullOrWhiteSpace(requestId)) + return false; + + return await _requestsRepository.CancelPendingAsync(requestId, cancelledBy, UtcNow, cancellationToken); + } + + public async Task ClaimNextRequestAsync(string workerName, CancellationToken cancellationToken = default) + { + return await _requestsRepository.ClaimNextPendingAsync(workerName, UtcNow, cancellationToken); + } + + public async Task ReportProgressAsync(string requestId, string progress, CancellationToken cancellationToken = default) + { + return await _requestsRepository.HeartbeatAsync(requestId, progress, UtcNow, cancellationToken); + } + + public async Task CompleteRequestAsync(string requestId, bool succeeded, string result, CancellationToken cancellationToken = default) + { + var status = succeeded ? SystemOperationStatuses.Completed : SystemOperationStatuses.Failed; + + return await _requestsRepository.FinishAsync(requestId, (int)status, result, UtcNow, cancellationToken); + } + + public async Task FailAbandonedRequestsAsync(CancellationToken cancellationToken = default) + { + var now = UtcNow; + + return await _requestsRepository.FailAbandonedAsync(now - AbandonedAfter, + "The worker running this request stopped before it finished (restart or crash). Request it again.", now, cancellationToken); + } + + public async Task DetectCacheDataLossAsync() + { + if (!Config.SystemBehaviorConfig.CacheEnabled || !_cacheProvider.IsConnected()) + return false; + + var marker = UtcNow.ToString("O", CultureInfo.InvariantCulture) + "|" + Guid.NewGuid().ToString("N"); + var stored = await _cacheProvider.GetOrAddStringAsync(CacheSentinelKey, marker, CacheSentinelLifetime); + + // Null is "unreachable", not "missing". Our own marker coming back means the key was absent and is now ours, so + // exactly one concurrent checker sees the loss. + return stored != null && string.Equals(stored, marker, StringComparison.Ordinal); + } + + public async Task GetCacheStatusAsync() + { + var status = new SystemOperationsCacheStatus + { + CacheEnabled = Config.SystemBehaviorConfig.CacheEnabled, + Connected = _cacheProvider.IsConnected() + }; + + if (!status.CacheEnabled || !status.Connected) + return status; + + status.DataPresentSinceUtc = ParseSentinel(await _cacheProvider.GetStringAsync(CacheSentinelKey)); + + return status; + } + + public async Task> ClearDepartmentCachesAsync(int departmentId) + { + // Best effort: one failing group must not keep the rest stale, so each is tried on its own. Mirrors the + // BackOffice department page's "Clear caches" list. + var failed = new List(); + + await StepAsync("department, users, personnel names, entity limits", () => _departmentsService.InvalidateAllDepartmentsCache(departmentId)); + Step("current payment and plan", () => _subscriptionsService.ClearCacheForCurrentPayment(departmentId)); + await StepAsync("groups", async () => + { + foreach (var group in await _departmentGroupsService.GetAllGroupsForDepartmentUnlimitedThinAsync(departmentId) ?? new List()) + await _departmentGroupsService.InvalidateGroupInCache(group.DepartmentGroupId); + }); + Step("call priorities", () => _callsService.InvalidateCallPrioritiesForDepartmentInCache(departmentId)); + Step("action logs", () => _actionLogsService.InvalidateActionLogs(departmentId)); + Step("custom states", () => _customStateService.InvalidateCustomStateInCache(departmentId)); + Step("latest personnel states", () => _userStateService.InvalidateLatestStatesForDepartmentCache(departmentId)); + await StepAsync("feature-flag overrides", () => _featureToggleService.InvalidateDepartmentOverrideCacheAsync(departmentId)); + + return failed; + + void Step(string name, Action invalidate) + { + try + { + invalidate(); + } + catch (Exception ex) + { + Logging.LogException(ex, $"Clearing the {name} cache failed for department {departmentId}."); + failed.Add(name); + } + } + + async Task StepAsync(string name, Func invalidate) + { + try + { + await invalidate(); + } + catch (Exception ex) + { + Logging.LogException(ex, $"Clearing the {name} cache failed for department {departmentId}."); + failed.Add(name); + } + } + } + + /// The time half of a sentinel value ("{utc O}|{nonce}"), or null when it is missing or not ours. + public static DateTime? ParseSentinel(string value) + { + if (string.IsNullOrWhiteSpace(value)) + return null; + + var separator = value.IndexOf('|'); + var time = separator > 0 ? value.Substring(0, separator) : value; + + return DateTime.TryParse(time, CultureInfo.InvariantCulture, DateTimeStyles.RoundtripKind, out var parsed) + ? parsed.ToUniversalTime() + : null; + } + + private static string Truncate(string value, int length) => + string.IsNullOrEmpty(value) || value.Length <= length ? value : value.Substring(0, length); + } +} diff --git a/Providers/Resgrid.Providers.Migrations/Migrations/M0267_AddSystemOperationRequests.cs b/Providers/Resgrid.Providers.Migrations/Migrations/M0267_AddSystemOperationRequests.cs new file mode 100644 index 000000000..f22fe888a --- /dev/null +++ b/Providers/Resgrid.Providers.Migrations/Migrations/M0267_AddSystemOperationRequests.cs @@ -0,0 +1,61 @@ +using FluentMigrator; + +namespace Resgrid.Providers.Migrations.Migrations +{ + /// + /// System operation requests: staff ask the worker (command 76) from BackOffice -> System Operations to rebuild cached + /// state (the Redis security matrices, department caches) or run a daily job now; the worker also queues a matrix + /// rebuild on its own when it finds Redis came back empty. The worker claims the oldest Pending row with a conditional + /// update, refreshes HeartbeatOn while it runs, and records the outcome. A system record: no foreign key to + /// Departments and not part of the department purge (TargetDepartmentId only narrows an operation). + /// + [Migration(267)] + public class M0267_AddSystemOperationRequests : Migration + { + private const string Table = "SystemOperationRequests"; + + public override void Up() + { + if (!Schema.Table(Table).Exists()) + { + Create.Table(Table) + .WithColumn("SystemOperationRequestId").AsString(128).NotNullable().PrimaryKey() + .WithColumn("OperationType").AsInt32().NotNullable() + .WithColumn("TargetDepartmentId").AsInt32().Nullable() + .WithColumn("Status").AsInt32().NotNullable() + .WithColumn("Source").AsInt32().NotNullable() + .WithColumn("RequestedBy").AsString(256).NotNullable() + .WithColumn("Reason").AsString(500).Nullable() + .WithColumn("RequestedOn").AsDateTime2().NotNullable() + .WithColumn("StartedOn").AsDateTime2().Nullable() + .WithColumn("HeartbeatOn").AsDateTime2().Nullable() + .WithColumn("CompletedOn").AsDateTime2().Nullable() + .WithColumn("WorkerName").AsString(256).Nullable() + .WithColumn("Progress").AsString(500).Nullable() + .WithColumn("Result").AsString(2000).Nullable() + .WithColumn("CancelledBy").AsString(256).Nullable(); + } + + if (!Schema.Table(Table).Index("IX_SystemOperationRequests_Status_RequestedOn").Exists()) + { + Create.Index("IX_SystemOperationRequests_Status_RequestedOn") + .OnTable(Table) + .OnColumn("Status").Ascending() + .OnColumn("RequestedOn").Ascending(); + } + + if (!Schema.Table(Table).Index("IX_SystemOperationRequests_RequestedOn").Exists()) + { + Create.Index("IX_SystemOperationRequests_RequestedOn") + .OnTable(Table) + .OnColumn("RequestedOn").Descending(); + } + } + + public override void Down() + { + if (Schema.Table(Table).Exists()) + Delete.Table(Table); + } + } +} diff --git a/Providers/Resgrid.Providers.MigrationsPg/Migrations/M0267_AddSystemOperationRequestsPg.cs b/Providers/Resgrid.Providers.MigrationsPg/Migrations/M0267_AddSystemOperationRequestsPg.cs new file mode 100644 index 000000000..60cc18415 --- /dev/null +++ b/Providers/Resgrid.Providers.MigrationsPg/Migrations/M0267_AddSystemOperationRequestsPg.cs @@ -0,0 +1,58 @@ +using FluentMigrator; + +namespace Resgrid.Providers.MigrationsPg.Migrations +{ + /// + /// System operation requests (see the SQL Server M0267): BackOffice -> System Operations queues work for worker 76, + /// which claims Pending rows with a conditional update and records the outcome. No foreign key to departments. + /// + [Migration(267)] + public class M0267_AddSystemOperationRequestsPg : Migration + { + private const string Table = "systemoperationrequests"; + + public override void Up() + { + if (!Schema.Table(Table).Exists()) + { + Create.Table(Table) + .WithColumn("systemoperationrequestid").AsCustom("citext").NotNullable().PrimaryKey() + .WithColumn("operationtype").AsInt32().NotNullable() + .WithColumn("targetdepartmentid").AsInt32().Nullable() + .WithColumn("status").AsInt32().NotNullable() + .WithColumn("source").AsInt32().NotNullable() + .WithColumn("requestedby").AsCustom("citext").NotNullable() + .WithColumn("reason").AsCustom("citext").Nullable() + .WithColumn("requestedon").AsDateTime2().NotNullable() + .WithColumn("startedon").AsDateTime2().Nullable() + .WithColumn("heartbeaton").AsDateTime2().Nullable() + .WithColumn("completedon").AsDateTime2().Nullable() + .WithColumn("workername").AsCustom("citext").Nullable() + .WithColumn("progress").AsCustom("citext").Nullable() + .WithColumn("result").AsCustom("citext").Nullable() + .WithColumn("cancelledby").AsCustom("citext").Nullable(); + } + + if (!Schema.Table(Table).Index("ix_systemoperationrequests_status_requestedon").Exists()) + { + Create.Index("ix_systemoperationrequests_status_requestedon") + .OnTable(Table) + .OnColumn("status").Ascending() + .OnColumn("requestedon").Ascending(); + } + + if (!Schema.Table(Table).Index("ix_systemoperationrequests_requestedon").Exists()) + { + Create.Index("ix_systemoperationrequests_requestedon") + .OnTable(Table) + .OnColumn("requestedon").Descending(); + } + } + + public override void Down() + { + if (Schema.Table(Table).Exists()) + Delete.Table(Table); + } + } +} diff --git a/Repositories/Resgrid.Repositories.DataRepository/Modules/ApiDataModule.cs b/Repositories/Resgrid.Repositories.DataRepository/Modules/ApiDataModule.cs index 99eeba9aa..ec124fd4b 100644 --- a/Repositories/Resgrid.Repositories.DataRepository/Modules/ApiDataModule.cs +++ b/Repositories/Resgrid.Repositories.DataRepository/Modules/ApiDataModule.cs @@ -305,6 +305,7 @@ protected override void Load(ContainerBuilder builder) builder.RegisterType().As().InstancePerLifetimeScope(); builder.RegisterType().As().InstancePerLifetimeScope(); builder.RegisterType().As().InstancePerLifetimeScope(); + builder.RegisterType().As().InstancePerLifetimeScope(); // CheckIn Repositories builder.RegisterType().As().InstancePerLifetimeScope(); diff --git a/Repositories/Resgrid.Repositories.DataRepository/Modules/DataModule.cs b/Repositories/Resgrid.Repositories.DataRepository/Modules/DataModule.cs index 974de47f2..3d231110b 100644 --- a/Repositories/Resgrid.Repositories.DataRepository/Modules/DataModule.cs +++ b/Repositories/Resgrid.Repositories.DataRepository/Modules/DataModule.cs @@ -325,6 +325,7 @@ protected override void Load(ContainerBuilder builder) builder.RegisterType().As().InstancePerLifetimeScope(); builder.RegisterType().As().InstancePerLifetimeScope(); builder.RegisterType().As().InstancePerLifetimeScope(); + builder.RegisterType().As().InstancePerLifetimeScope(); // CheckIn Repositories builder.RegisterType().As().InstancePerLifetimeScope(); diff --git a/Repositories/Resgrid.Repositories.DataRepository/Modules/NonWebDataModule.cs b/Repositories/Resgrid.Repositories.DataRepository/Modules/NonWebDataModule.cs index 72b2e338b..cf8e54c73 100644 --- a/Repositories/Resgrid.Repositories.DataRepository/Modules/NonWebDataModule.cs +++ b/Repositories/Resgrid.Repositories.DataRepository/Modules/NonWebDataModule.cs @@ -292,6 +292,7 @@ protected override void Load(ContainerBuilder builder) builder.RegisterType().As().InstancePerLifetimeScope(); builder.RegisterType().As().InstancePerLifetimeScope(); builder.RegisterType().As().InstancePerLifetimeScope(); + builder.RegisterType().As().InstancePerLifetimeScope(); // CheckIn Repositories builder.RegisterType().As().InstancePerLifetimeScope(); diff --git a/Repositories/Resgrid.Repositories.DataRepository/Modules/TestingDataModule.cs b/Repositories/Resgrid.Repositories.DataRepository/Modules/TestingDataModule.cs index 90f85591c..e56676f7a 100644 --- a/Repositories/Resgrid.Repositories.DataRepository/Modules/TestingDataModule.cs +++ b/Repositories/Resgrid.Repositories.DataRepository/Modules/TestingDataModule.cs @@ -326,6 +326,7 @@ protected override void Load(ContainerBuilder builder) builder.RegisterType().As().InstancePerLifetimeScope(); builder.RegisterType().As().InstancePerLifetimeScope(); builder.RegisterType().As().InstancePerLifetimeScope(); + builder.RegisterType().As().InstancePerLifetimeScope(); // CheckIn Repositories builder.RegisterType().As().InstancePerLifetimeScope(); diff --git a/Repositories/Resgrid.Repositories.DataRepository/SystemOperationRequestsRepository.cs b/Repositories/Resgrid.Repositories.DataRepository/SystemOperationRequestsRepository.cs new file mode 100644 index 000000000..a6fd39c04 --- /dev/null +++ b/Repositories/Resgrid.Repositories.DataRepository/SystemOperationRequestsRepository.cs @@ -0,0 +1,159 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Resgrid.Model; +using Resgrid.Model.Repositories; +using Resgrid.Model.Repositories.Connection; +using Resgrid.Model.Repositories.Queries; +using Resgrid.Repositories.DataRepository.Configs; + +namespace Resgrid.Repositories.DataRepository +{ + /// System operation requests (M0267). + public class SystemOperationRequestsRepository : RmsRepositoryBase, ISystemOperationRequestsRepository + { + private const string Table = "SystemOperationRequests"; + private const int ClaimAttempts = 5; + + public SystemOperationRequestsRepository(IConnectionProvider connectionProvider, SqlConfiguration sqlConfiguration, IUnitOfWork unitOfWork, IQueryFactory queryFactory) + : base(connectionProvider, sqlConfiguration, unitOfWork, queryFactory) { } + + public async Task> GetRecentAsync(int take) + { + var sql = IsPostgres + ? $"SELECT * FROM {Tbl(Table)} ORDER BY {Col("RequestedOn")} DESC LIMIT {P}Take" + : $"SELECT TOP ({P}Take) * FROM {Tbl(Table)} ORDER BY {Col("RequestedOn")} DESC"; + + var rows = await QueryAsync(sql, new { Take = Math.Max(1, take) }); + + return rows?.ToList() ?? new List(); + } + + public Task GetPendingAsync(int operationType, int? targetDepartmentId) + { + // A null parameter compared with IS NULL has no type for Npgsql to bind, so the target filter is chosen here. + var target = targetDepartmentId.HasValue + ? $"{Col("TargetDepartmentId")} = {P}TargetDepartmentId" + : $"{Col("TargetDepartmentId")} IS NULL"; + + return QueryFirstOrDefaultAsync( + $"SELECT * FROM {Tbl(Table)} WHERE {Col("OperationType")} = {P}OperationType AND {Col("Status")} = {P}Pending AND {target} " + + $"ORDER BY {Col("RequestedOn")}", + new { OperationType = operationType, Pending = (int)SystemOperationStatuses.Pending, TargetDepartmentId = targetDepartmentId ?? 0 }); + } + + public async Task ClaimNextPendingAsync(string workerName, DateTime now, CancellationToken cancellationToken = default) + { + var oldestPending = IsPostgres + ? $"SELECT {Col("SystemOperationRequestId")} FROM {Tbl(Table)} WHERE {Col("Status")} = {P}Pending ORDER BY {Col("RequestedOn")}, {Col("SystemOperationRequestId")} LIMIT 1" + : $"SELECT TOP 1 {Col("SystemOperationRequestId")} FROM {Tbl(Table)} WHERE {Col("Status")} = {P}Pending ORDER BY {Col("RequestedOn")}, {Col("SystemOperationRequestId")}"; + + // Another worker can claim the candidate between the read and the conditional update; then take the next one. + for (var attempt = 0; attempt < ClaimAttempts; attempt++) + { + var candidateId = await QueryFirstOrDefaultAsync(oldestPending, new { Pending = (int)SystemOperationStatuses.Pending }, cancellationToken); + + if (string.IsNullOrWhiteSpace(candidateId)) + return null; + + var claimed = await ExecuteAsync( + $"UPDATE {Tbl(Table)} SET {Col("Status")} = {P}Running, {Col("StartedOn")} = {P}Now, {Col("HeartbeatOn")} = {P}Now, " + + $"{Col("WorkerName")} = {P}WorkerName, {Col("Progress")} = NULL " + + $"WHERE {Col("SystemOperationRequestId")} = {P}Id AND {Col("Status")} = {P}Pending", + new + { + Id = candidateId, + Running = (int)SystemOperationStatuses.Running, + Pending = (int)SystemOperationStatuses.Pending, + Now = DatabaseTimestamp(now), + WorkerName = Truncate(workerName, 256) + }, + cancellationToken); + + if (claimed == 1) + return await QueryFirstOrDefaultAsync( + $"SELECT * FROM {Tbl(Table)} WHERE {Col("SystemOperationRequestId")} = {P}Id", + new { Id = candidateId }, + cancellationToken); + } + + return null; + } + + public async Task HeartbeatAsync(string requestId, string progress, DateTime now, CancellationToken cancellationToken = default) + { + var setProgress = progress == null ? "" : $", {Col("Progress")} = {P}Progress"; + + var updated = await ExecuteAsync( + $"UPDATE {Tbl(Table)} SET {Col("HeartbeatOn")} = {P}Now{setProgress} " + + $"WHERE {Col("SystemOperationRequestId")} = {P}Id AND {Col("Status")} = {P}Running", + new + { + Id = requestId, + Now = DatabaseTimestamp(now), + Progress = Truncate(progress, SystemOperationRequest.ProgressMaxLength), + Running = (int)SystemOperationStatuses.Running + }, + cancellationToken); + + return updated == 1; + } + + public async Task FinishAsync(string requestId, int status, string result, DateTime now, CancellationToken cancellationToken = default) + { + var updated = await ExecuteAsync( + $"UPDATE {Tbl(Table)} SET {Col("Status")} = {P}Status, {Col("CompletedOn")} = {P}Now, {Col("HeartbeatOn")} = {P}Now, {Col("Result")} = {P}Result " + + $"WHERE {Col("SystemOperationRequestId")} = {P}Id AND {Col("Status")} = {P}Running", + new + { + Id = requestId, + Status = status, + Now = DatabaseTimestamp(now), + Result = Truncate(result, SystemOperationRequest.ResultMaxLength), + Running = (int)SystemOperationStatuses.Running + }, + cancellationToken); + + return updated == 1; + } + + public async Task CancelPendingAsync(string requestId, string cancelledBy, DateTime now, CancellationToken cancellationToken = default) + { + var updated = await ExecuteAsync( + $"UPDATE {Tbl(Table)} SET {Col("Status")} = {P}Cancelled, {Col("CompletedOn")} = {P}Now, {Col("CancelledBy")} = {P}CancelledBy " + + $"WHERE {Col("SystemOperationRequestId")} = {P}Id AND {Col("Status")} = {P}Pending", + new + { + Id = requestId, + Cancelled = (int)SystemOperationStatuses.Cancelled, + Now = DatabaseTimestamp(now), + CancelledBy = Truncate(cancelledBy, 256), + Pending = (int)SystemOperationStatuses.Pending + }, + cancellationToken); + + return updated == 1; + } + + public Task FailAbandonedAsync(DateTime heartbeatBefore, string result, DateTime now, CancellationToken cancellationToken = default) + { + return ExecuteAsync( + $"UPDATE {Tbl(Table)} SET {Col("Status")} = {P}Failed, {Col("CompletedOn")} = {P}Now, {Col("Result")} = {P}Result " + + $"WHERE {Col("Status")} = {P}Running AND COALESCE({Col("HeartbeatOn")}, {Col("StartedOn")}, {Col("RequestedOn")}) < {P}Cutoff", + new + { + Failed = (int)SystemOperationStatuses.Failed, + Now = DatabaseTimestamp(now), + Result = Truncate(result, SystemOperationRequest.ResultMaxLength), + Running = (int)SystemOperationStatuses.Running, + Cutoff = DatabaseTimestamp(heartbeatBefore) + }, + cancellationToken); + } + + private static string Truncate(string value, int length) => + string.IsNullOrEmpty(value) || value.Length <= length ? value : value.Substring(0, length); + } +} diff --git a/Tests/Resgrid.Tests/Repositories/SystemOperationRequestsDatabaseTests.cs b/Tests/Resgrid.Tests/Repositories/SystemOperationRequestsDatabaseTests.cs new file mode 100644 index 000000000..ee67840ed --- /dev/null +++ b/Tests/Resgrid.Tests/Repositories/SystemOperationRequestsDatabaseTests.cs @@ -0,0 +1,277 @@ +using System; +using System.Collections.Concurrent; +using System.Data.Common; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Dapper; +using FluentAssertions; +using FluentMigrator; +using FluentMigrator.Runner; +using FluentMigrator.Runner.Initialization; +using Microsoft.Data.SqlClient; +using Microsoft.Extensions.DependencyInjection; +using Moq; +using Npgsql; +using NUnit.Framework; +using Resgrid.Config; +using Resgrid.Model; +using Resgrid.Model.Repositories.Connection; +using Resgrid.Model.Repositories.Queries; +using Resgrid.Model.Repositories.Queries.Contracts; +using Resgrid.Providers.Migrations.Migrations; +using Resgrid.Providers.MigrationsPg.Migrations; +using Resgrid.Repositories.DataRepository; +using Resgrid.Repositories.DataRepository.Configs; +using Resgrid.Repositories.DataRepository.Queries; +using Resgrid.Repositories.DataRepository.Queries.Common; +using Resgrid.Repositories.DataRepository.Servers.SqlServer; + +namespace Resgrid.Tests.Repositories +{ + /// + /// Real-database proof for M0267 on both engines: requests insert and read back, the claim hands each waiting request + /// to exactly one worker (also under concurrent claimers) oldest first, heartbeats and outcomes only land on a running + /// request, only a waiting request can be cancelled, and a running request without a recent heartbeat is failed as + /// abandoned. Set RESGRID_ADP_SQLSERVER_TEST_CONNECTION / RESGRID_ADP_POSTGRES_TEST_CONNECTION (server-level + /// connections) to run. + /// + [TestFixture(DatabaseTypes.SqlServer), TestFixture(DatabaseTypes.Postgres), NonParallelizable] + public class SystemOperationRequestsDatabaseTests(DatabaseTypes type) + { + private const string Prefix = "system_operation_requests_"; + private static readonly DateTime Base = new DateTime(2026, 10, 7, 12, 0, 0, DateTimeKind.Utc); + + private DatabaseTypes _previous; + private bool _configured; + private string _master, _connection, _database; + private ServiceProvider _runner; + + private bool IsPostgres => type == DatabaseTypes.Postgres; + + private string Table => IsPostgres ? "systemoperationrequests" : "SystemOperationRequests"; + + private DbConnection Connect(string connection) => IsPostgres ? new NpgsqlConnection(connection) : new SqlConnection(connection); + + private SystemOperationRequestsRepository Repository() + { + var connections = new Mock(); + connections.Setup(c => c.Create()).Returns(() => Connect(_connection)); + + SqlConfiguration configuration = IsPostgres ? new PostgreSqlConfiguration() : new SqlServerConfiguration(); + var queries = new ConcurrentDictionary(); + queries[typeof(InsertQuery)] = new InsertQuery(configuration); + queries[typeof(UpdateQuery)] = new UpdateQuery(configuration); + queries[typeof(SelectByIdQuery)] = new SelectByIdQuery(configuration); + var list = new Mock(); + list.Setup(l => l.RetrieveQueryList()).Returns(queries); + + return new SystemOperationRequestsRepository(connections.Object, configuration, Mock.Of(), new QueryFactory(list.Object)); + } + + [OneTimeSetUp] + public async Task Create_isolated_database() + { + _master = Environment.GetEnvironmentVariable(IsPostgres ? "RESGRID_ADP_POSTGRES_TEST_CONNECTION" : "RESGRID_ADP_SQLSERVER_TEST_CONNECTION"); + if (string.IsNullOrWhiteSpace(_master)) Assert.Ignore("Set a test connection to run real database checks."); + + if (IsPostgres) AppContext.SetSwitch("Npgsql.EnableLegacyTimestampBehavior", true); + + _previous = DataConfig.DatabaseType; + DataConfig.DatabaseType = type; + _configured = true; + var database = Prefix + Guid.NewGuid().ToString("N"); + await using (var master = Connect(_master)) + await master.ExecuteAsync("CREATE DATABASE " + database); + _database = database; + _connection = IsPostgres + ? new NpgsqlConnectionStringBuilder(_master) { Database = _database }.ConnectionString + : new SqlConnectionStringBuilder(_master) { InitialCatalog = _database }.ConnectionString; + + if (IsPostgres) + { + await using var setup = Connect(_connection); + await setup.OpenAsync(); + await setup.ExecuteAsync("CREATE EXTENSION IF NOT EXISTS citext;"); + + // The data source loaded its types before citext existed; without a reload a citext column cannot be read back. + await ((NpgsqlConnection)setup).ReloadTypesAsync(); + } + + var source = new Mock(); + source.Setup(s => s.GetMigrations()).Returns(IsPostgres + ? new IMigration[] { new M0267_AddSystemOperationRequestsPg() } + : new IMigration[] { new M0267_AddSystemOperationRequests() }); + _runner = new ServiceCollection().AddFluentMigratorCore().ConfigureRunner(r => + { + if (IsPostgres) r.AddPostgres(); else r.AddSqlServer(); + r.WithGlobalConnectionString(_connection); + }).AddSingleton(source.Object).BuildServiceProvider(); + _runner.GetRequiredService().MigrateUp(); + } + + [OneTimeTearDown] + public async Task Remove_only_this_fixture_database() + { + _runner?.Dispose(); + if (_configured) DataConfig.DatabaseType = _previous; + // Only a database this fixture created is ever dropped. + if (_database == null) return; + if (!_database.StartsWith(Prefix, StringComparison.Ordinal) || !Guid.TryParseExact(_database.Substring(Prefix.Length), "N", out _)) + throw new InvalidOperationException("Unexpected test database name."); + if (IsPostgres) NpgsqlConnection.ClearAllPools(); else SqlConnection.ClearAllPools(); + await using var master = Connect(_master); + await master.ExecuteAsync(IsPostgres ? "DROP DATABASE " + _database + " WITH (FORCE)" + : "ALTER DATABASE " + _database + " SET SINGLE_USER WITH ROLLBACK IMMEDIATE; DROP DATABASE " + _database); + } + + [SetUp] + public async Task Empty_the_table() + { + await using var database = Connect(_connection); + await database.ExecuteAsync("DELETE FROM " + Table); + } + + private async Task QueueAsync(SystemOperationTypes type, int? departmentId, DateTime requestedOn) + { + return await Repository().InsertAsync(new SystemOperationRequest + { + SystemOperationRequestId = Guid.NewGuid().ToString(), + OperationType = (int)type, + TargetDepartmentId = departmentId, + Status = (int)SystemOperationStatuses.Pending, + Source = (int)SystemOperationSources.BackOffice, + RequestedBy = "ops@resgrid.com", + Reason = "database test", + RequestedOn = IsPostgres ? DateTime.SpecifyKind(requestedOn, DateTimeKind.Unspecified) : requestedOn + }, CancellationToken.None); + } + + [Test] + public async Task A_queued_request_reads_back_and_is_found_as_the_waiting_one_for_its_target() + { + var all = await QueueAsync(SystemOperationTypes.RebuildSecurityMatrices, null, Base); + var one = await QueueAsync(SystemOperationTypes.RebuildSecurityMatrices, 12, Base.AddMinutes(1)); + var repository = Repository(); + + var read = await repository.GetByIdAsync(all.SystemOperationRequestId); + read.Should().NotBeNull(); + read.Reason.Should().Be("database test"); + read.TargetDepartmentId.Should().BeNull(); + + (await repository.GetPendingAsync((int)SystemOperationTypes.RebuildSecurityMatrices, null)).SystemOperationRequestId.Should().Be(all.SystemOperationRequestId); + (await repository.GetPendingAsync((int)SystemOperationTypes.RebuildSecurityMatrices, 12)).SystemOperationRequestId.Should().Be(one.SystemOperationRequestId); + (await repository.GetPendingAsync((int)SystemOperationTypes.RebuildSecurityMatrices, 13)).Should().BeNull(); + (await repository.GetPendingAsync((int)SystemOperationTypes.ReportingRollup, null)).Should().BeNull(); + } + + [Test] + public async Task Recent_requests_come_newest_first_and_honour_the_limit() + { + await QueueAsync(SystemOperationTypes.ReportingRollup, null, Base); + var newest = await QueueAsync(SystemOperationTypes.ChatRetention, null, Base.AddMinutes(2)); + await QueueAsync(SystemOperationTypes.BidExpiration, null, Base.AddMinutes(1)); + + var recent = await Repository().GetRecentAsync(2); + + recent.Should().HaveCount(2); + recent[0].SystemOperationRequestId.Should().Be(newest.SystemOperationRequestId); + } + + [Test] + public async Task Claims_take_the_oldest_waiting_request_and_mark_it_running() + { + var second = await QueueAsync(SystemOperationTypes.ReportingRollup, null, Base.AddMinutes(1)); + var first = await QueueAsync(SystemOperationTypes.ChatRetention, null, Base); + var repository = Repository(); + + var claimed = await repository.ClaimNextPendingAsync("worker-a", Base.AddMinutes(5)); + + claimed.SystemOperationRequestId.Should().Be(first.SystemOperationRequestId); + claimed.Status.Should().Be((int)SystemOperationStatuses.Running); + claimed.WorkerName.Should().Be("worker-a"); + claimed.StartedOn.Should().Be(Base.AddMinutes(5)); + claimed.HeartbeatOn.Should().Be(Base.AddMinutes(5)); + + (await repository.ClaimNextPendingAsync("worker-a", Base.AddMinutes(6))).SystemOperationRequestId.Should().Be(second.SystemOperationRequestId); + (await repository.ClaimNextPendingAsync("worker-a", Base.AddMinutes(7))).Should().BeNull(); + } + + [Test] + public async Task Concurrent_claimers_never_take_the_same_request() + { + for (var i = 0; i < 4; i++) + await QueueAsync(SystemOperationTypes.ReportingRollup, null, Base.AddSeconds(i)); + + var claims = await Task.WhenAll(Enumerable.Range(0, 8).Select(i => + Task.Run(() => Repository().ClaimNextPendingAsync($"worker-{i}", Base.AddMinutes(5))))); + + var won = claims.Where(c => c != null).Select(c => c.SystemOperationRequestId).ToList(); + won.Should().OnlyHaveUniqueItems(); + won.Should().HaveCount(4); + } + + [Test] + public async Task Heartbeats_and_outcomes_only_land_on_a_running_request() + { + var request = await QueueAsync(SystemOperationTypes.RebuildSecurityMatrices, null, Base); + var repository = Repository(); + + (await repository.HeartbeatAsync(request.SystemOperationRequestId, "too early", Base.AddMinutes(1))).Should().BeFalse("it is still waiting"); + (await repository.FinishAsync(request.SystemOperationRequestId, (int)SystemOperationStatuses.Completed, "too early", Base.AddMinutes(1))).Should().BeFalse(); + + await repository.ClaimNextPendingAsync("worker-a", Base.AddMinutes(2)); + (await repository.HeartbeatAsync(request.SystemOperationRequestId, "Rebuilt 10 of 40 departments.", Base.AddMinutes(3))).Should().BeTrue(); + (await repository.HeartbeatAsync(request.SystemOperationRequestId, null, Base.AddMinutes(4))).Should().BeTrue(); + + var running = await repository.GetByIdAsync(request.SystemOperationRequestId); + running.Progress.Should().Be("Rebuilt 10 of 40 departments.", "a beat without a new line keeps the last one"); + running.HeartbeatOn.Should().Be(Base.AddMinutes(4)); + + var longResult = new string('x', SystemOperationRequest.ResultMaxLength + 100); + (await repository.FinishAsync(request.SystemOperationRequestId, (int)SystemOperationStatuses.Failed, longResult, Base.AddMinutes(5))).Should().BeTrue(); + (await repository.FinishAsync(request.SystemOperationRequestId, (int)SystemOperationStatuses.Completed, "again", Base.AddMinutes(6))).Should().BeFalse("it already finished"); + + var finished = await repository.GetByIdAsync(request.SystemOperationRequestId); + finished.Status.Should().Be((int)SystemOperationStatuses.Failed); + finished.CompletedOn.Should().Be(Base.AddMinutes(5)); + finished.Result.Length.Should().Be(SystemOperationRequest.ResultMaxLength); + } + + [Test] + public async Task Only_a_waiting_request_can_be_cancelled() + { + var waiting = await QueueAsync(SystemOperationTypes.ReportingRollup, null, Base); + var running = await QueueAsync(SystemOperationTypes.ChatRetention, null, Base.AddSeconds(-1)); + var repository = Repository(); + await repository.ClaimNextPendingAsync("worker-a", Base.AddMinutes(1)); + + (await repository.CancelPendingAsync(running.SystemOperationRequestId, "ops", Base.AddMinutes(2))).Should().BeFalse(); + (await repository.CancelPendingAsync(waiting.SystemOperationRequestId, "ops", Base.AddMinutes(2))).Should().BeTrue(); + + var cancelled = await repository.GetByIdAsync(waiting.SystemOperationRequestId); + cancelled.Status.Should().Be((int)SystemOperationStatuses.Cancelled); + cancelled.CancelledBy.Should().Be("ops"); + (await repository.ClaimNextPendingAsync("worker-a", Base.AddMinutes(3))).Should().BeNull("a cancelled request is never claimed"); + } + + [Test] + public async Task A_running_request_without_a_recent_heartbeat_is_failed_as_abandoned() + { + var stale = await QueueAsync(SystemOperationTypes.ReportingRollup, null, Base); + var live = await QueueAsync(SystemOperationTypes.ChatRetention, null, Base.AddSeconds(1)); + var repository = Repository(); + await repository.ClaimNextPendingAsync("worker-a", Base.AddMinutes(1)); + await repository.ClaimNextPendingAsync("worker-a", Base.AddMinutes(1)); + await repository.HeartbeatAsync(live.SystemOperationRequestId, null, Base.AddMinutes(9)); + + var failed = await repository.FailAbandonedAsync(Base.AddMinutes(5), "worker stopped", Base.AddMinutes(10)); + + failed.Should().Be(1); + var abandoned = await repository.GetByIdAsync(stale.SystemOperationRequestId); + abandoned.Status.Should().Be((int)SystemOperationStatuses.Failed); + abandoned.Result.Should().Be("worker stopped"); + (await repository.GetByIdAsync(live.SystemOperationRequestId)).Status.Should().Be((int)SystemOperationStatuses.Running); + } + } +} diff --git a/Tests/Resgrid.Tests/Services/SystemOperationsServiceTests.cs b/Tests/Resgrid.Tests/Services/SystemOperationsServiceTests.cs new file mode 100644 index 000000000..aedfaced8 --- /dev/null +++ b/Tests/Resgrid.Tests/Services/SystemOperationsServiceTests.cs @@ -0,0 +1,266 @@ +using System; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using FluentAssertions; +using Moq; +using NUnit.Framework; +using Resgrid.Config; +using Resgrid.Model; +using Resgrid.Model.Providers; +using Resgrid.Model.Repositories; +using Resgrid.Model.Services; +using Resgrid.Services; + +namespace Resgrid.Tests.Services +{ + [TestFixture] + [NonParallelizable] + public class SystemOperationsServiceTests + { + private static readonly DateTimeOffset Now = new DateTimeOffset(2026, 10, 7, 14, 30, 0, TimeSpan.Zero); + + private Mock _requests; + private Mock _cache; + private Mock _departments; + private Mock _subscriptions; + private Mock _groups; + private Mock _calls; + private Mock _actionLogs; + private Mock _customStates; + private Mock _userStates; + private Mock _featureToggles; + private SystemOperationsService _service; + private bool _originalCacheEnabled; + + [SetUp] + public void SetUp() + { + _originalCacheEnabled = SystemBehaviorConfig.CacheEnabled; + SystemBehaviorConfig.CacheEnabled = true; + + _requests = new Mock(); + _requests.Setup(x => x.InsertAsync(It.IsAny(), It.IsAny(), It.IsAny())) + .ReturnsAsync((SystemOperationRequest r, CancellationToken _, bool _) => r); + _cache = new Mock(); + _departments = new Mock(); + _departments.Setup(x => x.GetDepartmentByIdAsync(12, It.IsAny())).ReturnsAsync(new Department { DepartmentId = 12 }); + _subscriptions = new Mock(); + _groups = new Mock(); + _groups.Setup(x => x.GetAllGroupsForDepartmentUnlimitedThinAsync(It.IsAny())).ReturnsAsync(new List()); + _calls = new Mock(); + _actionLogs = new Mock(); + _customStates = new Mock(); + _userStates = new Mock(); + _featureToggles = new Mock(); + + _service = new SystemOperationsService(_requests.Object, _cache.Object, _departments.Object, _subscriptions.Object, _groups.Object, + _calls.Object, _actionLogs.Object, _customStates.Object, _userStates.Object, _featureToggles.Object, new FixedTimeProvider(Now)); + } + + [TearDown] + public void TearDown() + { + SystemBehaviorConfig.CacheEnabled = _originalCacheEnabled; + } + + [Test] + public async Task Request_queues_a_pending_row_with_the_requester_and_reason() + { + var result = await _service.RequestAsync(SystemOperationTypes.RebuildSecurityMatrices, null, SystemOperationSources.BackOffice, + " ops@resgrid.com ", " Redis outage recovery "); + + result.Succeeded.Should().BeTrue(); + result.Created.Should().BeTrue(); + result.Request.Status.Should().Be((int)SystemOperationStatuses.Pending); + result.Request.OperationType.Should().Be((int)SystemOperationTypes.RebuildSecurityMatrices); + result.Request.Source.Should().Be((int)SystemOperationSources.BackOffice); + result.Request.RequestedBy.Should().Be("ops@resgrid.com"); + result.Request.Reason.Should().Be("Redis outage recovery"); + result.Request.RequestedOn.Should().Be(Now.UtcDateTime); + result.Request.TargetDepartmentId.Should().BeNull(); + Guid.TryParse(result.Request.SystemOperationRequestId, out _).Should().BeTrue(); + } + + [Test] + public async Task Request_matching_one_already_waiting_returns_it_instead_of_queueing_twice() + { + var waiting = new SystemOperationRequest { SystemOperationRequestId = "waiting", Status = (int)SystemOperationStatuses.Pending }; + _requests.Setup(x => x.GetPendingAsync((int)SystemOperationTypes.ClearDepartmentCaches, 12)).ReturnsAsync(waiting); + + var result = await _service.RequestAsync(SystemOperationTypes.ClearDepartmentCaches, 12, SystemOperationSources.BackOffice, "ops", "stale"); + + result.Succeeded.Should().BeTrue(); + result.Created.Should().BeFalse(); + result.Request.Should().BeSameAs(waiting); + _requests.Verify(x => x.InsertAsync(It.IsAny(), It.IsAny(), It.IsAny()), Times.Never); + } + + [Test] + public async Task Request_refuses_a_department_target_for_a_system_wide_operation() + { + var result = await _service.RequestAsync(SystemOperationTypes.ChatRetention, 12, SystemOperationSources.BackOffice, "ops", "why"); + + result.Succeeded.Should().BeFalse(); + result.Error.Should().Contain("cannot target one department"); + _requests.Verify(x => x.InsertAsync(It.IsAny(), It.IsAny(), It.IsAny()), Times.Never); + } + + [Test] + public async Task Request_refuses_a_department_that_does_not_exist() + { + var result = await _service.RequestAsync(SystemOperationTypes.RebuildSecurityMatrices, 99, SystemOperationSources.BackOffice, "ops", "why"); + + result.Succeeded.Should().BeFalse(); + result.Error.Should().Contain("99"); + } + + [Test] + public async Task Request_refuses_an_operation_the_catalog_does_not_know() + { + var result = await _service.RequestAsync((SystemOperationTypes)999, null, SystemOperationSources.BackOffice, "ops", "why"); + + result.Succeeded.Should().BeFalse(); + result.Error.Should().Contain("999"); + } + + [Test] + public async Task Request_with_a_department_records_the_target() + { + var result = await _service.RequestAsync(SystemOperationTypes.RebuildSecurityMatrices, 12, SystemOperationSources.BackOffice, "ops", "one department"); + + result.Created.Should().BeTrue(); + result.Request.TargetDepartmentId.Should().Be(12); + } + + [Test] + public async Task Request_trims_an_overlong_reason_to_the_column_size() + { + var result = await _service.RequestAsync(SystemOperationTypes.ReportingRollup, null, SystemOperationSources.BackOffice, "ops", + new string('r', SystemOperationRequest.ReasonMaxLength + 50)); + + result.Request.Reason.Length.Should().Be(SystemOperationRequest.ReasonMaxLength); + } + + [Test] + public async Task Complete_records_failed_or_completed() + { + await _service.CompleteRequestAsync("a", true, "done"); + await _service.CompleteRequestAsync("b", false, "boom"); + + _requests.Verify(x => x.FinishAsync("a", (int)SystemOperationStatuses.Completed, "done", Now.UtcDateTime, It.IsAny())); + _requests.Verify(x => x.FinishAsync("b", (int)SystemOperationStatuses.Failed, "boom", Now.UtcDateTime, It.IsAny())); + } + + [Test] + public async Task Abandoned_requests_are_those_without_a_heartbeat_for_the_abandon_window() + { + await _service.FailAbandonedRequestsAsync(); + + _requests.Verify(x => x.FailAbandonedAsync(Now.UtcDateTime - SystemOperationsService.AbandonedAfter, It.IsAny(), Now.UtcDateTime, + It.IsAny())); + } + + [Test] + public void The_heartbeat_is_well_inside_the_abandon_window() + { + // Several missed beats in a row, not one slow one, should be what fails a run. + (Resgrid.Workers.Console.Tasks.SystemOperationsTask.HeartbeatInterval * 4).Should().BeLessThan(SystemOperationsService.AbandonedAfter); + } + + [Test] + public async Task Cache_data_loss_is_detected_when_the_sentinel_comes_back_as_our_own_marker() + { + _cache.Setup(x => x.IsConnected()).Returns(true); + _cache.Setup(x => x.GetOrAddStringAsync(SystemOperationsService.CacheSentinelKey, It.IsAny(), SystemOperationsService.CacheSentinelLifetime)) + .ReturnsAsync((string _, string marker, TimeSpan _) => marker); + + (await _service.DetectCacheDataLossAsync()).Should().BeTrue(); + } + + [Test] + public async Task Cache_data_loss_is_not_detected_when_the_sentinel_was_already_there() + { + _cache.Setup(x => x.IsConnected()).Returns(true); + _cache.Setup(x => x.GetOrAddStringAsync(SystemOperationsService.CacheSentinelKey, It.IsAny(), It.IsAny())) + .ReturnsAsync("2026-10-01T00:00:00.0000000Z|abc"); + + (await _service.DetectCacheDataLossAsync()).Should().BeFalse(); + } + + [Test] + public async Task An_unreachable_cache_is_never_reported_as_data_loss() + { + _cache.Setup(x => x.IsConnected()).Returns(true); + _cache.Setup(x => x.GetOrAddStringAsync(It.IsAny(), It.IsAny(), It.IsAny())).ReturnsAsync((string)null); + + (await _service.DetectCacheDataLossAsync()).Should().BeFalse(); + + _cache.Setup(x => x.IsConnected()).Returns(false); + (await _service.DetectCacheDataLossAsync()).Should().BeFalse(); + } + + [Test] + public async Task A_disabled_cache_is_never_checked() + { + SystemBehaviorConfig.CacheEnabled = false; + _cache.Setup(x => x.IsConnected()).Returns(true); + + (await _service.DetectCacheDataLossAsync()).Should().BeFalse(); + _cache.Verify(x => x.GetOrAddStringAsync(It.IsAny(), It.IsAny(), It.IsAny()), Times.Never); + } + + [Test] + public async Task Cache_status_reports_when_the_cache_has_held_its_data_since() + { + _cache.Setup(x => x.IsConnected()).Returns(true); + _cache.Setup(x => x.GetStringAsync(SystemOperationsService.CacheSentinelKey)).ReturnsAsync("2026-10-07T09:15:00.0000000Z|abc"); + + var status = await _service.GetCacheStatusAsync(); + + status.CacheEnabled.Should().BeTrue(); + status.Connected.Should().BeTrue(); + status.DataPresentSinceUtc.Should().Be(new DateTime(2026, 10, 7, 9, 15, 0, DateTimeKind.Utc)); + } + + [TestCase(null)] + [TestCase("")] + [TestCase("not a time|abc")] + public void A_missing_or_foreign_sentinel_has_no_time(string value) + { + SystemOperationsService.ParseSentinel(value).Should().BeNull(); + } + + [Test] + public async Task Clearing_department_caches_keeps_going_past_a_failing_group() + { + _calls.Setup(x => x.InvalidateCallPrioritiesForDepartmentInCache(12)).Throws(new InvalidOperationException("redis")); + _groups.Setup(x => x.GetAllGroupsForDepartmentUnlimitedThinAsync(12)) + .ReturnsAsync(new List { new DepartmentGroup { DepartmentGroupId = 3 }, new DepartmentGroup { DepartmentGroupId = 4 } }); + + var failed = await _service.ClearDepartmentCachesAsync(12); + + failed.Should().ContainSingle().Which.Should().Be("call priorities"); + _departments.Verify(x => x.InvalidateAllDepartmentsCache(12)); + _subscriptions.Verify(x => x.ClearCacheForCurrentPayment(12)); + _groups.Verify(x => x.InvalidateGroupInCache(3)); + _groups.Verify(x => x.InvalidateGroupInCache(4)); + _actionLogs.Verify(x => x.InvalidateActionLogs(12)); + _customStates.Verify(x => x.InvalidateCustomStateInCache(12)); + _userStates.Verify(x => x.InvalidateLatestStatesForDepartmentCache(12)); + _featureToggles.Verify(x => x.InvalidateDepartmentOverrideCacheAsync(12)); + } + + private sealed class FixedTimeProvider : TimeProvider + { + private readonly DateTimeOffset _now; + + public FixedTimeProvider(DateTimeOffset now) + { + _now = now; + } + + public override DateTimeOffset GetUtcNow() => _now; + } + } +} diff --git a/Tests/Resgrid.Tests/Workers/Console/SystemOperationsTests.cs b/Tests/Resgrid.Tests/Workers/Console/SystemOperationsTests.cs new file mode 100644 index 000000000..78c544604 --- /dev/null +++ b/Tests/Resgrid.Tests/Workers/Console/SystemOperationsTests.cs @@ -0,0 +1,318 @@ +using System; +using System.Collections.Generic; +using System.IO; +using System.Linq; +using System.Reflection; +using System.Text.RegularExpressions; +using System.Threading; +using System.Threading.Tasks; +using Autofac; +using FluentAssertions; +using Microsoft.Extensions.Logging; +using Moq; +using NUnit.Framework; +using Quidjibo.Misc; +using Quidjibo.Models; +using Resgrid.Config; +using Resgrid.Model; +using Resgrid.Model.Providers; +using Resgrid.Model.Services; +using Resgrid.Workers.Console.Commands; +using Resgrid.Workers.Console.SystemOperations; +using Resgrid.Workers.Console.Tasks; + +namespace Resgrid.Tests.Workers.Console +{ + [TestFixture] + [NonParallelizable] + public class SystemOperationsTests + { + private static readonly FieldInfo WorkerBootstrapperContainerField = typeof(Resgrid.Workers.Framework.Bootstrapper) + .GetField("_container", BindingFlags.Static | BindingFlags.NonPublic)!; + + private IContainer _originalWorkerContainer; + private IContainer _testWorkerContainer; + private bool _originalCacheEnabled; + private bool _originalUtf8CleanupEnabled; + private bool _originalLocationRetentionEnabled; + private string _originalTtsBaseUrl; + private string _originalTtsAdminKey; + + [SetUp] + public void SetUp() + { + _originalWorkerContainer = WorkerBootstrapperContainerField.GetValue(null) as IContainer; + _originalCacheEnabled = SystemBehaviorConfig.CacheEnabled; + _originalUtf8CleanupEnabled = SystemBehaviorConfig.Utf8CleanupEnabled; + _originalLocationRetentionEnabled = UnitTrackingConfig.LocationRetentionWorkerEnabled; + _originalTtsBaseUrl = TtsConfig.ServiceBaseUrl; + _originalTtsAdminKey = TtsConfig.StaticPromptAdminKey; + } + + [TearDown] + public void TearDown() + { + WorkerBootstrapperContainerField.SetValue(null, _originalWorkerContainer); + _testWorkerContainer?.Dispose(); + _testWorkerContainer = null; + SystemBehaviorConfig.CacheEnabled = _originalCacheEnabled; + SystemBehaviorConfig.Utf8CleanupEnabled = _originalUtf8CleanupEnabled; + UnitTrackingConfig.LocationRetentionWorkerEnabled = _originalLocationRetentionEnabled; + TtsConfig.ServiceBaseUrl = _originalTtsBaseUrl; + TtsConfig.StaticPromptAdminKey = _originalTtsAdminKey; + } + + #region Catalog, runner and schedule stay in step + + [Test] + public void Every_operation_type_has_one_catalog_entry() + { + var types = Enum.GetValues(); + + SystemOperationCatalog.All.Select(x => x.Type).Should().OnlyHaveUniqueItems(); + SystemOperationCatalog.All.Select(x => x.Type).Should().BeEquivalentTo(types); + SystemOperationCatalog.All.Should().OnlyContain(x => !string.IsNullOrWhiteSpace(x.Name) && !string.IsNullOrWhiteSpace(x.Description)); + } + + [Test] + public void The_worker_can_run_every_catalog_entry() + { + SystemOperationRunner.SupportedTypes.Should().BeEquivalentTo(SystemOperationCatalog.All.Select(x => x.Type)); + } + + [Test] + public void Every_scheduled_job_the_catalog_mirrors_is_scheduled_under_that_command_id() + { + var program = FindRepositoryFile(Path.Combine("Workers", "Resgrid.Workers.Console", "Program.cs")); + if (program == null) + { + Assert.Inconclusive("Workers.Console Program.cs not found relative to the test assembly; source-scan check skipped."); + return; + } + + var source = System.IO.File.ReadAllText(program); + + foreach (var descriptor in SystemOperationCatalog.All.Where(x => x.WorkerCommandId.HasValue)) + Regex.IsMatch(source, $@"Command\s*\(\s*{descriptor.WorkerCommandId.Value}\s*\)").Should() + .BeTrue($"{descriptor.Type} says it runs worker {descriptor.WorkerCommandId} early, so Program.cs must schedule that command"); + + Regex.IsMatch(source, @"SystemOperationsCommand\s*\(\s*76\s*\)").Should().BeTrue("worker 76 runs the requests"); + } + + [Test] + public void Only_cache_operations_can_target_one_department() + { + SystemOperationCatalog.All.Where(x => x.SupportsDepartmentScope).Select(x => x.Type).Should() + .BeEquivalentTo(new[] { SystemOperationTypes.RebuildSecurityMatrices, SystemOperationTypes.ClearDepartmentCaches }); + } + + #endregion + + #region Runner + + [Test] + public async Task An_operation_this_build_does_not_know_fails_instead_of_running() + { + var outcome = await new SystemOperationRunner(Mock.Of()).RunAsync(Request((SystemOperationTypes)999), new SystemOperationProgress(), CancellationToken.None); + + outcome.Succeeded.Should().BeFalse(); + outcome.Message.Should().Contain("999"); + } + + [Test] + public async Task A_department_target_on_a_system_wide_job_is_refused() + { + var outcome = await new SystemOperationRunner(Mock.Of()).RunAsync(Request(SystemOperationTypes.ChatRetention, 12), + new SystemOperationProgress(), CancellationToken.None); + + outcome.Succeeded.Should().BeFalse(); + outcome.Message.Should().Contain("cannot target one department"); + } + + [Test] + public async Task Jobs_turned_off_in_the_worker_configuration_are_not_run() + { + SystemBehaviorConfig.Utf8CleanupEnabled = false; + UnitTrackingConfig.LocationRetentionWorkerEnabled = false; + TtsConfig.ServiceBaseUrl = null; + var runner = new SystemOperationRunner(Mock.Of()); + + foreach (var type in new[] { SystemOperationTypes.Utf8Cleanup, SystemOperationTypes.UnitTrackingLocationRetention, SystemOperationTypes.RefreshTtsStaticPrompts }) + { + var outcome = await runner.RunAsync(Request(type), new SystemOperationProgress(), CancellationToken.None); + + outcome.Succeeded.Should().BeFalse(); + outcome.Message.Should().StartWith("Not run:"); + } + } + + [Test] + public async Task Cache_operations_do_not_report_success_when_nothing_can_be_written() + { + var runner = new SystemOperationRunner(Mock.Of(), () => Scope(cacheConnected: false)); + + SystemBehaviorConfig.CacheEnabled = false; + var disabled = await runner.RunAsync(Request(SystemOperationTypes.RebuildSecurityMatrices), new SystemOperationProgress(), CancellationToken.None); + disabled.Succeeded.Should().BeFalse(); + disabled.Message.Should().Contain("CacheEnabled"); + + SystemBehaviorConfig.CacheEnabled = true; + var unreachable = await runner.RunAsync(Request(SystemOperationTypes.ClearDepartmentCaches, 12), new SystemOperationProgress(), CancellationToken.None); + unreachable.Succeeded.Should().BeFalse(); + unreachable.Message.Should().Contain("not reachable"); + } + + [Test] + public async Task Clearing_every_department_reports_progress_and_the_departments_that_failed() + { + SystemBehaviorConfig.CacheEnabled = true; + var operations = new Mock(); + operations.Setup(x => x.ClearDepartmentCachesAsync(It.IsAny())).ReturnsAsync(new List()); + operations.Setup(x => x.ClearDepartmentCachesAsync(2)).ReturnsAsync(new List { "groups" }); + var departments = new Mock(); + departments.Setup(x => x.GetAllAsync()).ReturnsAsync(new List + { + new Department { DepartmentId = 1 }, new Department { DepartmentId = 2 }, new Department { DepartmentId = 3 } + }); + var runner = new SystemOperationRunner(Mock.Of(), () => Scope(cacheConnected: true, operations.Object, departments.Object)); + var progress = new SystemOperationProgress(); + + var outcome = await runner.RunAsync(Request(SystemOperationTypes.ClearDepartmentCaches), progress, CancellationToken.None); + + outcome.Succeeded.Should().BeFalse(); + outcome.Message.Should().Contain("1 of 3").And.Contain("2 (groups)"); + progress.Latest.Should().Be("Cleared 3 of 3 departments."); + operations.Verify(x => x.ClearDepartmentCachesAsync(It.IsAny()), Times.Exactly(3)); + } + + [Test] + public async Task Clearing_one_department_touches_only_that_department() + { + SystemBehaviorConfig.CacheEnabled = true; + var operations = new Mock(); + operations.Setup(x => x.ClearDepartmentCachesAsync(12)).ReturnsAsync(new List()); + var runner = new SystemOperationRunner(Mock.Of(), () => Scope(cacheConnected: true, operations.Object)); + + var outcome = await runner.RunAsync(Request(SystemOperationTypes.ClearDepartmentCaches, 12), new SystemOperationProgress(), CancellationToken.None); + + outcome.Succeeded.Should().BeTrue(); + operations.Verify(x => x.ClearDepartmentCachesAsync(12), Times.Once); + operations.Verify(x => x.ClearDepartmentCachesAsync(It.Is(id => id != 12)), Times.Never); + } + + [Test] + public void The_job_logger_remembers_the_first_error_line_and_passes_everything_through() + { + var inner = new Mock(); + var logger = new ErrorCapturingLogger(inner.Object); + + logger.LogInformation("starting"); + logger.LogError("System.InvalidOperationException: first failure\r\n at Somewhere()"); + logger.LogError("second failure"); + + logger.FirstError.Should().Be("System.InvalidOperationException: first failure"); + inner.Verify(x => x.Log(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), + It.IsAny>()), Times.Exactly(3)); + } + + [Test] + public void Progress_keeps_the_latest_non_blank_line_from_either_report_shape() + { + var progress = new SystemOperationProgress(); + + progress.Report(10, "first"); + progress.Report(20, " "); + ((IQuidjiboProgress)progress).Report(new Tracker(30, "from tracker")); + + progress.Latest.Should().Be("from tracker"); + } + + #endregion + + #region Worker 76 tick + + [Test] + public async Task A_tick_queues_a_matrix_rebuild_when_redis_lost_its_data_and_records_each_claimed_request() + { + SystemBehaviorConfig.Utf8CleanupEnabled = false; + var claimed = Request(SystemOperationTypes.Utf8Cleanup); + var operations = new Mock(); + operations.Setup(x => x.DetectCacheDataLossAsync()).ReturnsAsync(true); + operations.Setup(x => x.RequestAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), + It.IsAny(), It.IsAny())) + .ReturnsAsync(new SystemOperationRequestResult { Created = true, Request = Request(SystemOperationTypes.RebuildSecurityMatrices) }); + operations.SetupSequence(x => x.ClaimNextRequestAsync(It.IsAny(), It.IsAny())) + .ReturnsAsync(claimed) + .ReturnsAsync((SystemOperationRequest)null); + SetWorkerContainer(operations.Object); + + await new SystemOperationsTask(Mock.Of()).ProcessAsync(new SystemOperationsCommand(76), Mock.Of(), CancellationToken.None); + + operations.Verify(x => x.FailAbandonedRequestsAsync(It.IsAny()), Times.Once); + operations.Verify(x => x.RequestAsync(SystemOperationTypes.RebuildSecurityMatrices, null, SystemOperationSources.CacheDataLossDetected, + "system", It.IsAny(), It.IsAny()), Times.Once); + operations.Verify(x => x.CompleteRequestAsync(claimed.SystemOperationRequestId, false, It.Is(m => m.StartsWith("Not run:")), + It.IsAny()), Times.Once); + } + + [Test] + public async Task A_tick_with_an_intact_cache_queues_nothing_of_its_own() + { + var operations = new Mock(); + operations.Setup(x => x.DetectCacheDataLossAsync()).ReturnsAsync(false); + SetWorkerContainer(operations.Object); + + await new SystemOperationsTask(Mock.Of()).ProcessAsync(new SystemOperationsCommand(76), Mock.Of(), CancellationToken.None); + + operations.Verify(x => x.RequestAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny(), + It.IsAny(), It.IsAny()), Times.Never); + operations.Verify(x => x.ClaimNextRequestAsync(It.IsAny(), It.IsAny()), Times.Once); + } + + #endregion + + private static SystemOperationRequest Request(SystemOperationTypes type, int? departmentId = null) => new SystemOperationRequest + { + SystemOperationRequestId = Guid.NewGuid().ToString(), + OperationType = (int)type, + TargetDepartmentId = departmentId, + Status = (int)SystemOperationStatuses.Running, + RequestedBy = "ops@resgrid.com" + }; + + private static ILifetimeScope Scope(bool cacheConnected, ISystemOperationsService operations = null, IDepartmentsService departments = null) + { + var cache = new Mock(); + cache.Setup(x => x.IsConnected()).Returns(cacheConnected); + + var builder = new ContainerBuilder(); + builder.RegisterInstance(cache.Object).As(); + builder.RegisterInstance(operations ?? Mock.Of()).As(); + builder.RegisterInstance(departments ?? Mock.Of()).As(); + + return builder.Build(); + } + + private void SetWorkerContainer(ISystemOperationsService operations) + { + var builder = new ContainerBuilder(); + builder.RegisterInstance(operations).As(); + _testWorkerContainer = builder.Build(); + WorkerBootstrapperContainerField.SetValue(null, _testWorkerContainer); + } + + private static string FindRepositoryFile(string relativePath) + { + var directory = new DirectoryInfo(TestContext.CurrentContext.TestDirectory); + while (directory != null) + { + var candidate = Path.Combine(directory.FullName, relativePath); + if (System.IO.File.Exists(candidate)) + return candidate; + directory = directory.Parent; + } + + return null; + } + } +} diff --git a/Workers/Resgrid.Workers.Console/Commands/SystemOperationsCommand.cs b/Workers/Resgrid.Workers.Console/Commands/SystemOperationsCommand.cs new file mode 100644 index 000000000..7bd765b3d --- /dev/null +++ b/Workers/Resgrid.Workers.Console/Commands/SystemOperationsCommand.cs @@ -0,0 +1,18 @@ +using System; +using System.Collections.Generic; +using Quidjibo.Commands; + +namespace Resgrid.Workers.Console.Commands +{ + public class SystemOperationsCommand : IQuidjiboCommand + { + public int Id { get; } + public Guid? CorrelationId { get; set; } + public Dictionary Metadata { get; set; } + + public SystemOperationsCommand(int id) + { + Id = id; + } + } +} diff --git a/Workers/Resgrid.Workers.Console/Program.cs b/Workers/Resgrid.Workers.Console/Program.cs index 278666ae7..cd5d81303 100644 --- a/Workers/Resgrid.Workers.Console/Program.cs +++ b/Workers/Resgrid.Workers.Console/Program.cs @@ -387,6 +387,14 @@ await Client.ScheduleAsync("Security Refresh", Cron.Daily(2, 0), stoppingToken); + // Every minute: runs the BackOffice -> System Operations requests (matrix rebuilds, cache clears, daily + // jobs on demand) and queues a matrix rebuild on its own when Redis comes back without its data. + _logger.Log(LogLevel.Information, "Scheduling System Operations"); + await Client.ScheduleAsync("System Operations", + new Commands.SystemOperationsCommand(76), + Cron.MinuteIntervals(1), + stoppingToken); + _logger.Log(LogLevel.Information, "Scheduling GDPR Data Export"); await Client.ScheduleAsync("GDPR Data Export", new Commands.GdprExportCommand(16), diff --git a/Workers/Resgrid.Workers.Console/SystemOperations/SystemOperationRunner.cs b/Workers/Resgrid.Workers.Console/SystemOperations/SystemOperationRunner.cs new file mode 100644 index 000000000..72d6ac2fc --- /dev/null +++ b/Workers/Resgrid.Workers.Console/SystemOperations/SystemOperationRunner.cs @@ -0,0 +1,319 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Autofac; +using Microsoft.Extensions.Logging; +using Quidjibo.Commands; +using Quidjibo.Handlers; +using Quidjibo.Misc; +using Quidjibo.Models; +using Resgrid.Config; +using Resgrid.Model; +using Resgrid.Model.Providers; +using Resgrid.Model.Services; +using Resgrid.Workers.Console.Commands; +using Resgrid.Workers.Console.Tasks; +using Resgrid.Workers.Framework; +using Resgrid.Workers.Framework.Logic; + +namespace Resgrid.Workers.Console.SystemOperations +{ + public sealed class SystemOperationOutcome + { + private SystemOperationOutcome(bool succeeded, string message) + { + Succeeded = succeeded; + Message = message; + } + + public bool Succeeded { get; } + + public string Message { get; } + + public static SystemOperationOutcome Ok(string message) => new SystemOperationOutcome(true, message); + + public static SystemOperationOutcome Fail(string message) => new SystemOperationOutcome(false, message); + } + + /// + /// The latest progress line of a running operation. Doubles as the IQuidjiboProgress handed to a job's handler; the + /// heartbeat loop copies the line onto the request row, so the handler never writes to the database itself. + /// + public sealed class SystemOperationProgress : IQuidjiboProgress + { + private string _latest; + + public string Latest => Volatile.Read(ref _latest); + + public void Report(string message) + { + if (!string.IsNullOrWhiteSpace(message)) + Volatile.Write(ref _latest, message); + } + + public void Report(int value, string text) => Report(text); + + void IProgress.Report(Tracker value) => Report(value?.Text); + } + + /// + /// Runs one operation for worker 76. A scheduled job runs through the very handler + /// its schedule runs (same command id), so an early run is the nightly run. The older handlers catch their own + /// exceptions and only log them; the runner hands them a logger that remembers the first error so a run that failed + /// is not reported as completed. + /// + public class SystemOperationRunner + { + private readonly ILogger _logger; + private readonly Func _beginScope; + + public SystemOperationRunner(ILogger logger, Func beginScope = null) + { + _logger = logger; + _beginScope = beginScope ?? (() => Bootstrapper.GetKernel().BeginLifetimeScope()); + } + + /// Every operation this worker build can run. A catalog entry missing here fails its requests instead of running them. + public static IReadOnlyCollection SupportedTypes { get; } = new HashSet + { + SystemOperationTypes.RebuildSecurityMatrices, + SystemOperationTypes.ClearDepartmentCaches, + SystemOperationTypes.RefreshTtsStaticPrompts, + SystemOperationTypes.PendingDepartmentDeletions, + SystemOperationTypes.ReportingRollup, + SystemOperationTypes.UnitTrackingLocationRetention, + SystemOperationTypes.ChatRetention, + SystemOperationTypes.BidExpiration, + SystemOperationTypes.DeploymentFinanceReminder, + SystemOperationTypes.ComplianceExpiry, + SystemOperationTypes.RmsDueStateEvaluation, + SystemOperationTypes.RmsRetentionAndPurge, + SystemOperationTypes.PayDataReportingReadiness, + SystemOperationTypes.ProtectedWorkflowSweep, + SystemOperationTypes.Utf8Cleanup + }; + + public async Task RunAsync(SystemOperationRequest request, SystemOperationProgress progress, CancellationToken cancellationToken) + { + var type = (SystemOperationTypes)request.OperationType; + + if (!SupportedTypes.Contains(type) || SystemOperationCatalog.Get(type) == null) + return SystemOperationOutcome.Fail($"This worker build cannot run operation {request.OperationType}; deploy the worker that added it."); + + // The catalog decides which operations may target one department; a row written around the service is refused here too. + if (request.TargetDepartmentId.HasValue && !SystemOperationCatalog.Get(type).SupportsDepartmentScope) + return SystemOperationOutcome.Fail($"{SystemOperationCatalog.Get(type).Name} cannot target one department."); + + switch (type) + { + case SystemOperationTypes.RebuildSecurityMatrices: + return await RebuildSecurityMatricesAsync(request.TargetDepartmentId, progress, cancellationToken); + + case SystemOperationTypes.ClearDepartmentCaches: + return await ClearDepartmentCachesAsync(request.TargetDepartmentId, progress, cancellationToken); + + case SystemOperationTypes.RefreshTtsStaticPrompts: + if (string.IsNullOrWhiteSpace(TtsConfig.ServiceBaseUrl) || string.IsNullOrWhiteSpace(TtsConfig.StaticPromptAdminKey)) + return NotEnabled("the TTS service URL or static prompt admin key is not configured"); + return await RunJobAsync(logger => new TtsStaticPromptRefreshTask(logger), new TtsStaticPromptRefreshCommand(18), progress, cancellationToken); + + case SystemOperationTypes.PendingDepartmentDeletions: + return await RunJobAsync(logger => new SystemSqlQueueTask(logger), new SystemSqlQueueCommand(14), progress, cancellationToken); + + case SystemOperationTypes.ReportingRollup: + return await RunJobAsync(logger => new ReportingRollupTask(logger), new ReportingRollupCommand(21), progress, cancellationToken); + + case SystemOperationTypes.UnitTrackingLocationRetention: + if (!UnitTrackingConfig.LocationRetentionWorkerEnabled) + return NotEnabled("the location retention worker is turned off (UnitTrackingConfig.LocationRetentionWorkerEnabled)"); + return await RunJobAsync(logger => new UnitTrackingRetentionTask(logger), new UnitTrackingRetentionCommand(24), progress, cancellationToken); + + case SystemOperationTypes.ChatRetention: + return await RunJobAsync(logger => new ChatRetentionTask(logger), new ChatRetentionCommand(25), progress, cancellationToken); + + case SystemOperationTypes.BidExpiration: + return await RunJobAsync(_ => new BidExpirationTask(), new BidExpirationCommand(31), progress, cancellationToken); + + case SystemOperationTypes.DeploymentFinanceReminder: + return await RunJobAsync(_ => new DeploymentFinanceReminderTask(), new DeploymentFinanceReminderCommand(32), progress, cancellationToken); + + case SystemOperationTypes.ComplianceExpiry: + return await RunJobAsync(_ => new ComplianceExpiryTask(), new ComplianceExpiryCommand(33), progress, cancellationToken); + + case SystemOperationTypes.RmsDueStateEvaluation: + return await RunJobAsync(logger => new RmsDueStateEvaluationTask(logger), new RmsDueStateEvaluationCommand(42), progress, cancellationToken); + + case SystemOperationTypes.RmsRetentionAndPurge: + return await RunJobAsync(logger => new RmsRetentionAndPurgeTask(logger), new RmsRetentionAndPurgeCommand(43), progress, cancellationToken); + + case SystemOperationTypes.PayDataReportingReadiness: + return await RunJobAsync(_ => new PayDataReportingReadinessTask(), new PayDataReportingReadinessCommand(49), progress, cancellationToken); + + case SystemOperationTypes.ProtectedWorkflowSweep: + return await RunJobAsync(_ => new ProtectedWorkflowSweepTask(), new ProtectedWorkflowSweepCommand(71), progress, cancellationToken); + + case SystemOperationTypes.Utf8Cleanup: + if (!SystemBehaviorConfig.Utf8CleanupEnabled) + return NotEnabled("the UTF-8 cleanup is turned off (SystemBehaviorConfig.Utf8CleanupEnabled)"); + return await RunJobAsync(logger => new Utf8CleanupTask(logger), new Utf8CleanupCommand(22), progress, cancellationToken); + + default: + return SystemOperationOutcome.Fail($"This worker build cannot run operation {request.OperationType}; deploy the worker that added it."); + } + } + + private async Task RebuildSecurityMatricesAsync(int? departmentId, SystemOperationProgress progress, CancellationToken cancellationToken) + { + var unavailable = CacheUnavailableReason(); + if (unavailable != null) + return NotEnabled(unavailable); + + var logic = new SecurityLogic(); + + if (departmentId.HasValue) + { + progress.Report($"Rebuilding the security matrices for department {departmentId.Value}."); + var single = await logic.UpdateCachedSecurityForDepartment(departmentId.Value); + + return single.Item1 + ? SystemOperationOutcome.Ok($"Rebuilt the four security matrices for department {departmentId.Value}.") + : SystemOperationOutcome.Fail(single.Item2); + } + + var total = 0; + var result = await logic.UpdatedCachedSecurityForAllDepartments((done, count) => + { + total = count; + progress.Report($"Rebuilt {done} of {count} departments."); + return Task.CompletedTask; + }, cancellationToken); + + return result.Item1 + ? SystemOperationOutcome.Ok($"Rebuilt the security matrices for all {total} departments.") + : SystemOperationOutcome.Fail(result.Item2); + } + + private async Task ClearDepartmentCachesAsync(int? departmentId, SystemOperationProgress progress, CancellationToken cancellationToken) + { + var unavailable = CacheUnavailableReason(); + if (unavailable != null) + return NotEnabled(unavailable); + + using var scope = _beginScope(); + var systemOperations = scope.Resolve(); + + if (departmentId.HasValue) + { + var failedGroups = await systemOperations.ClearDepartmentCachesAsync(departmentId.Value); + + return failedGroups.Count == 0 + ? SystemOperationOutcome.Ok($"Cleared the caches of department {departmentId.Value}; they reload from the database on next read.") + : SystemOperationOutcome.Fail($"Department {departmentId.Value}: these cache groups failed: {string.Join(", ", failedGroups)}."); + } + + var departments = await scope.Resolve().GetAllAsync() ?? new List(); + var failures = new List(); + var done = 0; + + foreach (var department in departments) + { + cancellationToken.ThrowIfCancellationRequested(); + + var failedGroups = await systemOperations.ClearDepartmentCachesAsync(department.DepartmentId); + if (failedGroups.Count > 0) + failures.Add($"{department.DepartmentId} ({string.Join(", ", failedGroups)})"); + + done++; + progress.Report($"Cleared {done} of {departments.Count} departments."); + } + + if (failures.Count == 0) + return SystemOperationOutcome.Ok($"Cleared the caches of all {departments.Count} departments; they reload from the database on next read."); + + // Capped like the matrix rebuild: an outage would otherwise build a line per department. + return SystemOperationOutcome.Fail($"{failures.Count} of {departments.Count} departments had cache groups fail: {string.Join("; ", failures.GetRange(0, Math.Min(10, failures.Count)))}"); + } + + /// Runs a scheduled job's own handler once, outside its schedule. + private async Task RunJobAsync(Func> createHandler, TCommand command, + SystemOperationProgress progress, CancellationToken cancellationToken) where TCommand : IQuidjiboCommand + { + var logger = new ErrorCapturingLogger(_logger); + + try + { + await createHandler(logger).ProcessAsync(command, progress, cancellationToken); + } + catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) + { + throw; + } + catch (Exception ex) + { + Resgrid.Framework.Logging.LogException(ex, $"System operation run of {typeof(TCommand).Name} failed."); + return SystemOperationOutcome.Fail(ex.Message); + } + + if (logger.FirstError != null) + return SystemOperationOutcome.Fail($"{logger.FirstError} (full error in the worker log)"); + + return SystemOperationOutcome.Ok(progress.Latest ?? "Finished."); + } + + /// Null when the cache can be written; otherwise why a cache operation would silently do nothing. + private string CacheUnavailableReason() + { + if (!SystemBehaviorConfig.CacheEnabled) + return "caching is turned off (SystemBehaviorConfig.CacheEnabled)"; + + using var scope = _beginScope(); + if (!scope.Resolve().IsConnected()) + return "Redis is not reachable from the worker, so nothing would be written"; + + return null; + } + + private static SystemOperationOutcome NotEnabled(string reason) => + SystemOperationOutcome.Fail($"Not run: {reason} in this worker's configuration."); + } + + /// Passes everything through and remembers the first error-level line, trimmed to its first line. + public sealed class ErrorCapturingLogger : ILogger + { + private const int MaxErrorLength = 500; + + private readonly ILogger _inner; + private string _firstError; + + public ErrorCapturingLogger(ILogger inner) + { + _inner = inner; + } + + public string FirstError => Volatile.Read(ref _firstError); + + public IDisposable BeginScope(TState state) where TState : notnull => _inner?.BeginScope(state); + + public bool IsEnabled(LogLevel logLevel) => logLevel >= LogLevel.Error || (_inner?.IsEnabled(logLevel) ?? false); + + public void Log(LogLevel logLevel, EventId eventId, TState state, Exception exception, Func formatter) + { + if (logLevel >= LogLevel.Error) + { + var message = formatter?.Invoke(state, exception); + if (string.IsNullOrWhiteSpace(message)) + message = exception?.Message ?? "The job logged an error."; + + var firstLine = message.Split('\n')[0].Trim(); + if (firstLine.Length > MaxErrorLength) + firstLine = firstLine.Substring(0, MaxErrorLength); + + Interlocked.CompareExchange(ref _firstError, firstLine, null); + } + + _inner?.Log(logLevel, eventId, state, exception, formatter); + } + } +} diff --git a/Workers/Resgrid.Workers.Console/Tasks/SystemOperationsTask.cs b/Workers/Resgrid.Workers.Console/Tasks/SystemOperationsTask.cs new file mode 100644 index 000000000..2758e0cfb --- /dev/null +++ b/Workers/Resgrid.Workers.Console/Tasks/SystemOperationsTask.cs @@ -0,0 +1,177 @@ +using System; +using System.Threading; +using System.Threading.Tasks; +using Autofac; +using Microsoft.Extensions.Logging; +using Quidjibo.Handlers; +using Quidjibo.Misc; +using Resgrid.Model; +using Resgrid.Model.Services; +using Resgrid.Services; +using Resgrid.Workers.Console.Commands; +using Resgrid.Workers.Console.SystemOperations; +using Resgrid.Workers.Framework; + +namespace Resgrid.Workers.Console.Tasks +{ + /// + /// Worker 76, every minute. Notices when Redis came back without its data (and queues a security matrix rebuild for + /// every department), fails requests whose worker died mid-run, then claims and runs the waiting + /// BackOffice -> System Operations requests one at a time, heartbeating each so the page shows progress. + /// + public class SystemOperationsTask : IQuidjiboHandler + { + /// Well inside SystemOperationsService.AbandonedAfter, so a live run is never mistaken for a dead one. + public static TimeSpan HeartbeatInterval = TimeSpan.FromSeconds(30); + + /// Per tick: the rest wait for the next minute rather than holding this tick open indefinitely. + public const int MaxRequestsPerTick = 20; + + // One run at a time per worker process: a long rebuild spans many one-minute ticks, and the ticks queued behind it + // have nothing to add. Claims are conditional writes, so a second worker process still never runs the same request. + private static readonly SemaphoreSlim SingleFlight = new SemaphoreSlim(1, 1); + + public string Name => "System Operations"; + public int Priority => 1; + + private readonly ILogger _logger; + + public SystemOperationsTask(ILogger logger) + { + _logger = logger; + } + + public async Task ProcessAsync(SystemOperationsCommand command, IQuidjiboProgress progress, CancellationToken cancellationToken) + { + if (!await SingleFlight.WaitAsync(0, cancellationToken)) + return; + + try + { + var abandoned = await WithServiceAsync(s => s.FailAbandonedRequestsAsync(cancellationToken)); + if (abandoned > 0) + _logger.LogWarning("SystemOperations::Failed {Count} request(s) whose worker stopped mid-run", abandoned); + + await QueueRebuildIfCacheLostDataAsync(cancellationToken); + + var runner = new SystemOperationRunner(_logger); + var workerName = $"{Environment.MachineName}:{Environment.ProcessId}"; + + for (var i = 0; i < MaxRequestsPerTick && !cancellationToken.IsCancellationRequested; i++) + { + var request = await WithServiceAsync(s => s.ClaimNextRequestAsync(workerName, cancellationToken)); + if (request == null) + break; + + await RunRequestAsync(runner, request, cancellationToken); + } + } + catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) + { + } + catch (Exception ex) + { + Resgrid.Framework.Logging.LogException(ex); + _logger.LogError(ex.ToString()); + } + finally + { + SingleFlight.Release(); + } + } + + private async Task QueueRebuildIfCacheLostDataAsync(CancellationToken cancellationToken) + { + if (!await WithServiceAsync(s => s.DetectCacheDataLossAsync())) + return; + + // Matrix readers already answer from the permission rows on a miss, so this restores speed, not correctness. + // The other Redis entries are cache-aside and refill on their own; the BackOffice page can clear or rerun the rest. + var queued = await WithServiceAsync(s => s.RequestAsync(SystemOperationTypes.RebuildSecurityMatrices, null, + SystemOperationSources.CacheDataLossDetected, SystemOperationsService.SystemRequester, + "Redis came back without its data (the cache sentinel was missing). Rebuilding every department's security matrices.", + cancellationToken)); + + _logger.LogWarning("SystemOperations::Redis cache sentinel was missing; queued a security matrix rebuild ({RequestId}). {Error}", + queued?.Request?.SystemOperationRequestId, queued?.Error); + } + + private async Task RunRequestAsync(SystemOperationRunner runner, SystemOperationRequest request, CancellationToken cancellationToken) + { + var name = SystemOperationCatalog.GetName(request.OperationType); + _logger.LogInformation("SystemOperations::Starting {Operation} ({RequestId}) requested by {RequestedBy}", + name, request.SystemOperationRequestId, request.RequestedBy); + + var progress = new SystemOperationProgress(); + using var stopHeartbeat = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + var heartbeat = HeartbeatAsync(request.SystemOperationRequestId, progress, stopHeartbeat.Token); + + SystemOperationOutcome outcome; + + try + { + outcome = await runner.RunAsync(request, progress, cancellationToken); + } + catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) + { + outcome = SystemOperationOutcome.Fail("The worker shut down before the operation finished. Request it again."); + } + catch (Exception ex) + { + Resgrid.Framework.Logging.LogException(ex, $"System operation {name} ({request.SystemOperationRequestId}) failed."); + outcome = SystemOperationOutcome.Fail(ex.Message); + } + finally + { + stopHeartbeat.Cancel(); + await heartbeat; + } + + // Recorded even while shutting down: otherwise the row sits in Running until it is failed as abandoned. + await WithServiceAsync(s => s.CompleteRequestAsync(request.SystemOperationRequestId, outcome.Succeeded, outcome.Message, CancellationToken.None)); + + _logger.LogInformation("SystemOperations::Finished {Operation} ({RequestId}): {Succeeded} {Message}", + name, request.SystemOperationRequestId, outcome.Succeeded ? "succeeded" : "failed", outcome.Message); + } + + private async Task HeartbeatAsync(string requestId, SystemOperationProgress progress, CancellationToken stop) + { + string lastSent = null; + + while (!stop.IsCancellationRequested) + { + try + { + await Task.Delay(HeartbeatInterval, stop); + } + catch (OperationCanceledException) + { + return; + } + + var latest = progress.Latest; + + try + { + await WithServiceAsync(s => s.ReportProgressAsync(requestId, latest != lastSent ? latest : null, CancellationToken.None)); + lastSent = latest; + } + catch (Exception ex) + { + // A missed beat is not fatal; several in a row and the row is failed as abandoned. + Resgrid.Framework.Logging.LogException(ex, $"System operation heartbeat failed for {requestId}."); + } + } + } + + /// + /// Each call gets its own scope: the heartbeat runs alongside the operation, and two concurrent calls on one + /// scope would share its unit of work and connection. + /// + private static async Task WithServiceAsync(Func> call) + { + using var scope = Bootstrapper.GetKernel().BeginLifetimeScope(); + return await call(scope.Resolve()); + } + } +} diff --git a/Workers/Resgrid.Workers.Framework/Logic/SecurityLogic.cs b/Workers/Resgrid.Workers.Framework/Logic/SecurityLogic.cs index 069004cfb..304cec742 100644 --- a/Workers/Resgrid.Workers.Framework/Logic/SecurityLogic.cs +++ b/Workers/Resgrid.Workers.Framework/Logic/SecurityLogic.cs @@ -1,5 +1,6 @@ using Resgrid.Model.Services; using System; +using System.Threading; using System.Threading.Tasks; using Autofac; using Resgrid.Model; @@ -647,14 +648,44 @@ async Task getWhoCanViewUserLocations() /// anything failed, so a caller can tell a clean sweep from a partial one; without that the /// sweep looks successful no matter how many matrices were left stale. /// - public async Task> UpdatedCachedSecurityForAllDepartments() + /// Optional: called after each department with (departments done, departments total). + /// Stops the sweep between departments. + public async Task> UpdatedCachedSecurityForAllDepartments(Func progress = null, + CancellationToken cancellationToken = default) { List departments; using (var scope = Bootstrapper.GetKernel().BeginLifetimeScope()) departments = await scope.Resolve().GetAllAsync(); var failures = new List(); + var done = 0; - async Task rebuild(int departmentId, SecurityCacheTypes type) + foreach (var department in departments) + { + cancellationToken.ThrowIfCancellationRequested(); + + await RebuildDepartmentAsync(department.DepartmentId, failures); + + done++; + if (progress != null) + await progress(done, departments.Count); + } + + return Summarize(failures); + } + + /// The four matrices of one department, for an on-demand rebuild (BackOffice -> System Operations). + public async Task> UpdateCachedSecurityForDepartment(int departmentId) + { + var failures = new List(); + + await RebuildDepartmentAsync(departmentId, failures); + + return Summarize(failures); + } + + private async Task RebuildDepartmentAsync(int departmentId, List failures) + { + async Task rebuild(SecurityCacheTypes type) { var processed = await Process(new SecurityQueueItem() { DepartmentId = departmentId, Type = type }); @@ -662,14 +693,14 @@ async Task rebuild(int departmentId, SecurityCacheTypes type) failures.Add($"{departmentId}/{type}: {processed?.Item2}".Trim()); } - foreach (var department in departments) - { - await rebuild(department.DepartmentId, SecurityCacheTypes.WhoCanViewUnits); - await rebuild(department.DepartmentId, SecurityCacheTypes.WhoCanViewUnitLocations); - await rebuild(department.DepartmentId, SecurityCacheTypes.WhoCanViewPersonnel); - await rebuild(department.DepartmentId, SecurityCacheTypes.WhoCanViewPersonnelLocations); - } + await rebuild(SecurityCacheTypes.WhoCanViewUnits); + await rebuild(SecurityCacheTypes.WhoCanViewUnitLocations); + await rebuild(SecurityCacheTypes.WhoCanViewPersonnel); + await rebuild(SecurityCacheTypes.WhoCanViewPersonnelLocations); + } + private static Tuple Summarize(List failures) + { if (!failures.Any()) return new Tuple(true, String.Empty);