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);