diff --git a/Core/Resgrid.Localization/Areas/User/Records/Records.ar.resx b/Core/Resgrid.Localization/Areas/User/Records/Records.ar.resx
index f66a130cd..e5532a7b7 100644
--- a/Core/Resgrid.Localization/Areas/User/Records/Records.ar.resx
+++ b/Core/Resgrid.Localization/Areas/User/Records/Records.ar.resx
@@ -1655,6 +1655,7 @@
غير منطبق
نص الملاحظة
تم حفظ الملاحظة.
+ التاريخ غير صالح. تحقق من السنة.
الموضوع
تم إصدار الإشعار.
مرجع الإشعار
diff --git a/Core/Resgrid.Localization/Areas/User/Records/Records.de.resx b/Core/Resgrid.Localization/Areas/User/Records/Records.de.resx
index 51a84fc27..a09037814 100644
--- a/Core/Resgrid.Localization/Areas/User/Records/Records.de.resx
+++ b/Core/Resgrid.Localization/Areas/User/Records/Records.de.resx
@@ -1655,6 +1655,7 @@
k. A.
Notiztext
Notiz gespeichert.
+ Das Datum ist ungültig. Prüfen Sie das Jahr.
Betreff
Bescheid ausgestellt.
Bescheid-Referenz
diff --git a/Core/Resgrid.Localization/Areas/User/Records/Records.el.resx b/Core/Resgrid.Localization/Areas/User/Records/Records.el.resx
index c75c2b825..761336826 100644
--- a/Core/Resgrid.Localization/Areas/User/Records/Records.el.resx
+++ b/Core/Resgrid.Localization/Areas/User/Records/Records.el.resx
@@ -1655,6 +1655,7 @@
Δ/Ε
Κείμενο σημείωσης
Η σημείωση αποθηκεύτηκε.
+ Η ημερομηνία δεν είναι έγκυρη. Ελέγξτε το έτος.
Θέμα
Η ειδοποίηση εκδόθηκε.
Αναφορά ειδοποίησης
diff --git a/Core/Resgrid.Localization/Areas/User/Records/Records.en.resx b/Core/Resgrid.Localization/Areas/User/Records/Records.en.resx
index 7a78da112..6a15aed6a 100644
--- a/Core/Resgrid.Localization/Areas/User/Records/Records.en.resx
+++ b/Core/Resgrid.Localization/Areas/User/Records/Records.en.resx
@@ -1655,6 +1655,7 @@
N/A
Note text
Note saved.
+ The date is not valid. Check the year.
Subject
Notice issued.
Notice reference
diff --git a/Core/Resgrid.Localization/Areas/User/Records/Records.es.resx b/Core/Resgrid.Localization/Areas/User/Records/Records.es.resx
index 3bc380f90..22263687a 100644
--- a/Core/Resgrid.Localization/Areas/User/Records/Records.es.resx
+++ b/Core/Resgrid.Localization/Areas/User/Records/Records.es.resx
@@ -1655,6 +1655,7 @@
N/D
Texto de la nota
Nota guardada.
+ La fecha no es válida. Revise el año.
Asunto
Notificación emitida.
Referencia de la notificación
diff --git a/Core/Resgrid.Localization/Areas/User/Records/Records.fr.resx b/Core/Resgrid.Localization/Areas/User/Records/Records.fr.resx
index 39cad2845..fdfc73834 100644
--- a/Core/Resgrid.Localization/Areas/User/Records/Records.fr.resx
+++ b/Core/Resgrid.Localization/Areas/User/Records/Records.fr.resx
@@ -1655,6 +1655,7 @@
S. O.
Texte de la note
Note enregistrée.
+ La date n'est pas valide. Vérifiez l'année.
Objet
Avis émis.
Référence de l'avis
diff --git a/Core/Resgrid.Localization/Areas/User/Records/Records.it.resx b/Core/Resgrid.Localization/Areas/User/Records/Records.it.resx
index 451bbb6ea..b7e8cfe79 100644
--- a/Core/Resgrid.Localization/Areas/User/Records/Records.it.resx
+++ b/Core/Resgrid.Localization/Areas/User/Records/Records.it.resx
@@ -1655,6 +1655,7 @@
N/D
Testo della nota
Nota salvata.
+ La data non è valida. Verifica l'anno.
Oggetto
Avviso emesso.
Riferimento avviso
diff --git a/Core/Resgrid.Localization/Areas/User/Records/Records.pl.resx b/Core/Resgrid.Localization/Areas/User/Records/Records.pl.resx
index 801127df7..a06c2eae6 100644
--- a/Core/Resgrid.Localization/Areas/User/Records/Records.pl.resx
+++ b/Core/Resgrid.Localization/Areas/User/Records/Records.pl.resx
@@ -1655,6 +1655,7 @@
Nie dot.
Treść notatki
Notatka zapisana.
+ Data jest nieprawidłowa. Sprawdź rok.
Temat
Wezwanie wystawione.
Numer wezwania
diff --git a/Core/Resgrid.Localization/Areas/User/Records/Records.sv.resx b/Core/Resgrid.Localization/Areas/User/Records/Records.sv.resx
index a070ea3a7..874629edf 100644
--- a/Core/Resgrid.Localization/Areas/User/Records/Records.sv.resx
+++ b/Core/Resgrid.Localization/Areas/User/Records/Records.sv.resx
@@ -1655,6 +1655,7 @@
Ej tillämpl.
Anteckningstext
Anteckning sparad.
+ Datumet är ogiltigt. Kontrollera året.
Ämne
Föreläggande utfärdat.
Föreläggandereferens
diff --git a/Core/Resgrid.Localization/Areas/User/Records/Records.uk.resx b/Core/Resgrid.Localization/Areas/User/Records/Records.uk.resx
index bf37a95ac..874d35892 100644
--- a/Core/Resgrid.Localization/Areas/User/Records/Records.uk.resx
+++ b/Core/Resgrid.Localization/Areas/User/Records/Records.uk.resx
@@ -1655,6 +1655,7 @@
Н/З
Текст нотатки
Нотатку збережено.
+ Дата недійсна. Перевірте рік.
Тема
Припис видано.
Номер припису
diff --git a/Core/Resgrid.Model/Providers/ICacheProvider.cs b/Core/Resgrid.Model/Providers/ICacheProvider.cs
index 7da534636..ba6958667 100644
--- a/Core/Resgrid.Model/Providers/ICacheProvider.cs
+++ b/Core/Resgrid.Model/Providers/ICacheProvider.cs
@@ -14,6 +14,14 @@ public interface ICacheProvider
Task GetStringAsync(string cacheKey);
Task GetAsync(string cacheKey) where T : class;
+ ///
+ /// Atomically returns the string stored at the key, or stores and
+ /// returns it when the key is missing (concurrent callers all get the first value stored). Either
+ /// way the key's expiration is reset, so a key in regular use never ages out. Returns null only
+ /// when the cache is unavailable — unlike , null never means "absent".
+ ///
+ Task GetOrAddStringAsync(string cacheKey, string valueIfAbsent, TimeSpan slidingExpiration);
+
///
/// Atomically increments a counter and returns the new value. Sets the expiration on first
/// increment (when the value becomes 1). Returns 0 when the cache is unavailable.
diff --git a/Core/Resgrid.Model/Repositories/ISearchRepositories.cs b/Core/Resgrid.Model/Repositories/ISearchRepositories.cs
index c8c7cb3fd..11fd64178 100644
--- a/Core/Resgrid.Model/Repositories/ISearchRepositories.cs
+++ b/Core/Resgrid.Model/Repositories/ISearchRepositories.cs
@@ -36,6 +36,13 @@ public interface ISearchIndexStatesRepository : IRepository
Task GetAsync(string indexName, int departmentId);
Task> GetAllForIndexAsync(string indexName);
+
+ ///
+ /// Inserts unless a row for its (IndexName, DepartmentId) already exists, in one statement,
+ /// so concurrent callers never trip the unique index. Returns false, and leaves the existing row untouched, when
+ /// another writer got there first.
+ ///
+ Task InsertIfMissingAsync(SearchIndexState state, CancellationToken cancellationToken = default);
}
public interface ISearchIndexLeasesRepository : IRepository
diff --git a/Core/Resgrid.Model/Search/UnifiedSearchContracts.cs b/Core/Resgrid.Model/Search/UnifiedSearchContracts.cs
index ae8e41b22..7fdf6c574 100644
--- a/Core/Resgrid.Model/Search/UnifiedSearchContracts.cs
+++ b/Core/Resgrid.Model/Search/UnifiedSearchContracts.cs
@@ -144,8 +144,8 @@ public class SystemActionDefinition
/// Feature flag key that must evaluate true for the department (FeatureFlagKeys).
public string FeatureFlag { get; set; }
- /// Hide when the Records module is on: the legacy Logs pages are replaced after cutover.
- public bool HiddenWhenRecordsEnabled { get; set; }
+ /// Hide once the department's Records cutover is active: a legacy Logs write that the Logs pages now refuse. Reads stay listed.
+ public bool HiddenAfterRecordsCutover { get; set; }
}
public class SystemActionHit
diff --git a/Core/Resgrid.Model/Services/IChatServices.cs b/Core/Resgrid.Model/Services/IChatServices.cs
index a6ce790da..d0af2e920 100644
--- a/Core/Resgrid.Model/Services/IChatServices.cs
+++ b/Core/Resgrid.Model/Services/IChatServices.cs
@@ -217,7 +217,10 @@ public interface IChatPermissionService
/// Drops cached permission evaluations for a channel (membership/roles changed) and bumps the channel-list cache version.
Task InvalidateChannelCacheAsync(string chatChannelId);
- /// Current distributed authorization epoch used to isolate realtime channel groups after access changes.
+ ///
+ /// Current distributed authorization epoch used to isolate realtime channel groups after access changes.
+ /// Minted on first use; null only when the cache is unavailable, and callers must then fail closed.
+ ///
Task GetChannelAccessVersionAsync(string chatChannelId);
}
diff --git a/Core/Resgrid.Services/ChatChannelService.cs b/Core/Resgrid.Services/ChatChannelService.cs
index 19619e6e8..d51f34a93 100644
--- a/Core/Resgrid.Services/ChatChannelService.cs
+++ b/Core/Resgrid.Services/ChatChannelService.cs
@@ -41,7 +41,7 @@ public class ChatChannelService : IChatChannelService
private readonly ICacheProvider _cacheProvider;
private readonly IUnitOfWork _unitOfWork;
- // IncidentCommandService reaches back for channel provisioning through ServiceLocator, so this
+ // IncidentCommandService reaches back for channel provisioning through a Lazy, so this
// constructor edge does not close a resolution cycle.
private readonly IIncidentCommandService _incidentCommandService;
diff --git a/Core/Resgrid.Services/ChatMessageService.cs b/Core/Resgrid.Services/ChatMessageService.cs
index e749b7492..335db2caf 100644
--- a/Core/Resgrid.Services/ChatMessageService.cs
+++ b/Core/Resgrid.Services/ChatMessageService.cs
@@ -3,7 +3,8 @@
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
-using CommonServiceLocator;
+using Autofac;
+using Autofac.Core;
using Microsoft.Data.SqlClient;
using Newtonsoft.Json;
using Newtonsoft.Json.Linq;
@@ -33,6 +34,7 @@ public class ChatMessageService : IChatMessageService
private readonly IChatMessageMentionRepository _chatMessageMentionRepository;
private readonly IChatMessageAckRepository _chatMessageAckRepository;
private readonly IChatChannelMemberRepository _chatChannelMemberRepository;
+ private readonly ILifetimeScope _lifetimeScope;
private readonly IChatChannelService _chatChannelService;
private readonly IChatPermissionService _chatPermissionService;
private readonly IUserProfileService _userProfileService;
@@ -44,7 +46,7 @@ public ChatMessageService(IChatChannelRepository chatChannelRepository, IChatMes
IChatMessageReactionRepository chatMessageReactionRepository, IChatMessageMentionRepository chatMessageMentionRepository,
IChatMessageAckRepository chatMessageAckRepository, IChatChannelMemberRepository chatChannelMemberRepository,
IChatChannelService chatChannelService, IChatPermissionService chatPermissionService, IUserProfileService userProfileService,
- IUnitsService unitsService, IEventAggregator eventAggregator)
+ IUnitsService unitsService, IEventAggregator eventAggregator, ILifetimeScope lifetimeScope = null)
{
_chatChannelRepository = chatChannelRepository;
_chatMessageRepository = chatMessageRepository;
@@ -59,6 +61,7 @@ public ChatMessageService(IChatChannelRepository chatChannelRepository, IChatMes
_userProfileService = userProfileService;
_unitsService = unitsService;
_eventAggregator = eventAggregator;
+ _lifetimeScope = lifetimeScope;
}
public async Task SendMessageAsync(int departmentId, string senderUserId, ChatMessageSendRequest request, CancellationToken cancellationToken = default(CancellationToken))
@@ -807,16 +810,23 @@ await _chatMessageEditRepository.InsertAsync(new ChatMessageEdit
///
/// Push fan-out off the request path: per-recipient Novu calls can be slow for large channels,
- /// and a push failure must never fail the send. Fresh resolution inside the task keeps us off
- /// the request's disposed lifetime scope (ChatProvisioningEventService pattern).
+ /// and a push failure must never fail the send. The task outlives the request, and Autofac will not
+ /// resolve from a scope whose parent is disposed, so it runs in its own child of the ROOT scope: never
+ /// the request's scope, and never a root-resolved notifier sharing the root unit of work
+ /// (ChatProvisioningEventService pattern). Skipped when constructed outside the container.
///
private void FireAndForgetNotify(ChatChannel channel, ChatMessage message, List mentions)
{
+ var root = (_lifetimeScope as ISharingLifetimeScope)?.RootLifetimeScope;
+ if (root == null)
+ return;
+
_ = Task.Run(async () =>
{
try
{
- var notifier = ServiceLocator.Current.GetInstance();
+ using var scope = root.BeginLifetimeScope();
+ var notifier = scope.Resolve();
await notifier.NotifyMessageSentAsync(channel, message, mentions);
}
catch (Exception ex)
diff --git a/Core/Resgrid.Services/ChatPermissionService.cs b/Core/Resgrid.Services/ChatPermissionService.cs
index f328aebaa..d0742f67f 100644
--- a/Core/Resgrid.Services/ChatPermissionService.cs
+++ b/Core/Resgrid.Services/ChatPermissionService.cs
@@ -21,6 +21,13 @@ public class ChatPermissionService : IChatPermissionService
private static readonly TimeSpan CacheLength = TimeSpan.FromSeconds(60);
private static readonly TimeSpan VersionCacheLength = TimeSpan.FromDays(1);
+ ///
+ /// Sliding lifetime of a channel's access epoch. Every read refreshes it, so an epoch only lapses
+ /// on a channel with no joins or fan-out for this long; the lapse mints a new epoch, which strands
+ /// any still-joined connection until it rejoins, so keep this well beyond a connection's lifetime.
+ ///
+ private static readonly TimeSpan AccessVersionSlidingExpiration = TimeSpan.FromDays(30);
+
/// Shared version key rolled into every per-user channel-list cache key; bumped by InvalidateChannelCacheAsync.
internal const string ChannelListVersionCacheKey = "chatchannellistver";
@@ -299,7 +306,9 @@ public async Task InvalidateChannelCacheAsync(string chatChannelId)
if (string.IsNullOrWhiteSpace(chatChannelId))
return;
- await _cacheProvider.IncrementAsync(GetVersionKey(chatChannelId), VersionCacheLength);
+ // A fresh random epoch rather than an increment: once a key lapses, a counter restarts at a
+ // value an obsolete group (still holding a revoked connection) may already carry.
+ await _cacheProvider.SetStringAsync(GetVersionKey(chatChannelId), NewAccessVersion(), AccessVersionSlidingExpiration);
// Roll every per-user channel-list cache key forward too (channel set/visibility changed).
await _cacheProvider.IncrementAsync(ChannelListVersionCacheKey, VersionCacheLength);
@@ -312,17 +321,24 @@ public async Task GetChannelAccessVersionAsync(string chatChannelId)
try
{
- return await _cacheProvider.GetStringAsync(GetVersionKey(chatChannelId));
+ // A channel that has never been invalidated (or sat idle past the sliding window) has no
+ // epoch yet — mint one rather than treating the absence as an outage.
+ return await _cacheProvider.GetOrAddStringAsync(GetVersionKey(chatChannelId), NewAccessVersion(), AccessVersionSlidingExpiration);
}
catch (Exception ex)
{
- // A missing authorization epoch must stop realtime fan-out. Falling back to the old
- // group during a cache outage could reconnect a user whose access was just revoked.
+ // Null (cache unavailable) must stop realtime fan-out. Falling back to a fixed group during
+ // a cache outage could reconnect a user whose access was just revoked.
Resgrid.Framework.Logging.LogException(ex);
return null;
}
}
+ private static string NewAccessVersion()
+ {
+ return Guid.NewGuid().ToString("N");
+ }
+
private async Task EvaluateAccessAsync(ChatChannel channel, string userId, int? activeUnitId)
{
if (!await HasValidDepartmentScopeAsync(channel))
diff --git a/Core/Resgrid.Services/CoreEventService.cs b/Core/Resgrid.Services/CoreEventService.cs
index 80256f221..a36b2358a 100644
--- a/Core/Resgrid.Services/CoreEventService.cs
+++ b/Core/Resgrid.Services/CoreEventService.cs
@@ -1,6 +1,7 @@
using System;
using System.Threading.Tasks;
-using CommonServiceLocator;
+using Autofac;
+using Resgrid.Framework;
using Resgrid.Model;
using Resgrid.Model.Events;
using Resgrid.Model.Services;
@@ -8,22 +9,39 @@
namespace Resgrid.Services
{
+ ///
+ /// Registered as a singleton, so the scoped department settings service is NOT resolved once and kept: from
+ /// the root scope it would share one unit of work, and one DB connection, with everything else resolved there,
+ /// and the timestamp save runs inside an audited configuration transaction. Each event instead runs in its own
+ /// child scope, as does.
+ ///
public class CoreEventService : ICoreEventService
{
private readonly IEventAggregator _eventAggregator;
+ private readonly ILifetimeScope _lifetimeScope;
- public CoreEventService(IEventAggregator eventAggregator)
+ public CoreEventService(IEventAggregator eventAggregator, ILifetimeScope lifetimeScope)
{
_eventAggregator = eventAggregator;
+ _lifetimeScope = lifetimeScope;
- _eventAggregator.AddListener(departmentSettingsUpdateHandler);
+ // Fire-and-forget as before: the publisher (a unit, department or custom state save) is not held up.
+ _eventAggregator.AddListener(message => _ = UpdateDepartmentTimestampAsync(message));
}
- private Action departmentSettingsUpdateHandler = async delegate(DepartmentSettingsUpdateEvent message)
+ private async Task UpdateDepartmentTimestampAsync(DepartmentSettingsUpdateEvent message)
{
- var departmentSettingsService = ServiceLocator.Current.GetInstance();
- var result = await departmentSettingsService.SaveOrUpdateSettingAsync(message.DepartmentId, DateTime.UtcNow.ToString("G"), DepartmentSettingTypes.UpdateTimestamp);
- };
+ try
+ {
+ using var scope = _lifetimeScope.BeginLifetimeScope();
+ await scope.Resolve().SaveOrUpdateSettingAsync(message.DepartmentId, DateTime.UtcNow.ToString("G"), DepartmentSettingTypes.UpdateTimestamp);
+ }
+ catch (Exception ex)
+ {
+ // Nothing awaits this, so an escaping exception would go unobserved.
+ Logging.LogException(ex, $"Department update timestamp could not be saved for department {message?.DepartmentId}.");
+ }
+ }
public Task IncidentCommandUpdatedAsync(int departmentId, int callId)
{
diff --git a/Core/Resgrid.Services/IncidentCommandService.cs b/Core/Resgrid.Services/IncidentCommandService.cs
index 5af62a087..098234678 100644
--- a/Core/Resgrid.Services/IncidentCommandService.cs
+++ b/Core/Resgrid.Services/IncidentCommandService.cs
@@ -6,7 +6,6 @@
using System.Security.Cryptography;
using System.Threading;
using System.Threading.Tasks;
-using CommonServiceLocator;
using Resgrid.Framework;
using Resgrid.Model;
using Resgrid.Model.Events;
@@ -58,6 +57,16 @@ public class IncidentCommandService : IIncidentCommandService
private readonly ICallDispatchStatusService _callDispatchStatusService;
private readonly IQueueService _queueService;
+ // The chat and command-access sides depend on this service, so constructor-injecting them directly would close a
+ // DI cycle. Lazy resolves them on first use from this service's own lifetime scope; the service locator they
+ // replace resolved from the root scope, sharing its unit of work with everything else resolved there.
+ private readonly Lazy _commandAccessService;
+ private readonly Lazy _chatChannelService;
+ private readonly Lazy _chatChannelRepository;
+ private ICommandAccessService CommandAccess => _commandAccessService?.Value ?? throw new InvalidOperationException("Command access is unavailable.");
+ private IChatChannelService ChatChannels => _chatChannelService?.Value ?? throw new InvalidOperationException("Chat channels are unavailable.");
+ private IChatChannelRepository ChatChannelRepository => _chatChannelRepository?.Value ?? throw new InvalidOperationException("Chat channels are unavailable.");
+
public IncidentCommandService(
IIncidentCommandRepository incidentCommandRepository,
ICommandStructureNodeRepository commandStructureNodeRepository,
@@ -90,7 +99,10 @@ public IncidentCommandService(
IIncidentMapRepository incidentMapRepository,
IIncidentNeedEntityRepository incidentNeedEntityRepository,
ICallDispatchStatusService callDispatchStatusService,
- IQueueService queueService)
+ IQueueService queueService,
+ Lazy commandAccessService = null,
+ Lazy chatChannelService = null,
+ Lazy chatChannelRepository = null)
{
_incidentCommandRepository = incidentCommandRepository;
_commandStructureNodeRepository = commandStructureNodeRepository;
@@ -124,6 +136,9 @@ public IncidentCommandService(
_incidentNeedEntityRepository = incidentNeedEntityRepository;
_callDispatchStatusService = callDispatchStatusService;
_queueService = queueService;
+ _commandAccessService = commandAccessService;
+ _chatChannelService = chatChannelService;
+ _chatChannelRepository = chatChannelRepository;
}
#region Command lifecycle
@@ -475,11 +490,10 @@ public async Task GetCapabilitiesForUserAsync(int departme
// Dispatch app. CanAssistWithCommandAsync (not CanUseCommandAsync) is the right question: the
// permission is open by default, and granting board authority off that open default would hand
// every member rights nobody asked for.
- // Resolved through the service locator (matching this file's other cross-cutting lookups) so the
- // permission side, which has no dependency on this service, does not close a DI cycle.
+ // Lazy (see the field) so the permission side does not close a DI cycle.
try
{
- if (await ServiceLocator.Current.GetInstance().CanAssistWithCommandAsync(departmentId, userId))
+ if (await CommandAccess.CanAssistWithCommandAsync(departmentId, userId))
caps |= IncidentRoleCapabilityMap.CommandAssistCapabilities;
}
catch (Exception ex)
@@ -1074,7 +1088,7 @@ private async Task BackfillIncidentChatChannelsAsync(IncidentCommand command, in
try
{
var nodes = knownNodes ?? await GetNodesForCallAsync(departmentId, callId);
- await ServiceLocator.Current.GetInstance().EnsureIncidentChannelsAsync(command, nodes);
+ await ChatChannels.EnsureIncidentChannelsAsync(command, nodes);
}
catch (Exception ex)
{
@@ -1113,9 +1127,9 @@ private async Task PopulateResourceViewContactsAndChatAsync(ResourceIncidentView
try
{
- // Resolved through the service locator, matching DeleteNodeAsync: the chat side depends on
- // this service, so constructor-injecting it back would close a DI cycle.
- var channels = (await ServiceLocator.Current.GetInstance()
+ // Lazy (see the field): the chat side depends on this service, so constructor-injecting it back
+ // directly would close a DI cycle.
+ var channels = (await ChatChannelRepository
.GetByCallIdAsync(callId))?.ToList() ?? new List();
view.Chat.IncidentChannelId = channels.FirstOrDefault(c => c.ChannelType == (int)ChatChannelType.Incident)?.ChatChannelId;
@@ -1595,7 +1609,7 @@ await WriteLogAsync(node.IncidentCommandId, node.DepartmentId, node.CallId,
// service's constructor graph, and a chat failure must never fail the lane save.
try
{
- var chatChannelService = ServiceLocator.Current.GetInstance();
+ var chatChannelService = ChatChannels;
await chatChannelService.EnsureLaneChannelAsync(node, cancellationToken);
}
catch (Exception ex)
@@ -1624,8 +1638,8 @@ await WriteLogAsync(node.IncidentCommandId, node.DepartmentId, node.CallId,
// Best-effort: archive the lane's chat channel alongside the tombstoned node.
try
{
- var chatChannelService = ServiceLocator.Current.GetInstance();
- var laneChannel = (await ServiceLocator.Current.GetInstance().GetByCommandStructureNodeIdAsync(commandStructureNodeId));
+ var chatChannelService = ChatChannels;
+ var laneChannel = (await ChatChannelRepository.GetByCommandStructureNodeIdAsync(commandStructureNodeId));
if (laneChannel != null && !laneChannel.IsArchived)
await chatChannelService.SetChannelArchivedAsync(laneChannel.DepartmentId, laneChannel.ChatChannelId, true, userId, cancellationToken);
}
diff --git a/Core/Resgrid.Services/PermissionsService.cs b/Core/Resgrid.Services/PermissionsService.cs
index 8fd2db087..d2c5c536b 100644
--- a/Core/Resgrid.Services/PermissionsService.cs
+++ b/Core/Resgrid.Services/PermissionsService.cs
@@ -3,7 +3,6 @@
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
-using CommonServiceLocator;
using Resgrid.Framework;
using Resgrid.Model;
using Resgrid.Model.Events;
@@ -20,11 +19,23 @@ public class PermissionsService : IPermissionsService
private readonly IPermissionsRepository _permissionsRepository;
private readonly IDepartmentGroupsService _departmentGroupsService;
- public PermissionsService(IPermissionsRepository permissionsRepository, IUsersService usersService, IDepartmentGroupsService departmentGroupsService)
+ // The chat realtime refresh after a dispatch-login permission change. Chat permissions depend on this service,
+ // so they come in lazily (a direct dependency would close a DI cycle) and resolve from this service's own
+ // lifetime scope; the service locator they replace resolved from the root scope. Left null outside the
+ // container, where the refresh is skipped.
+ private readonly Lazy _chatChannelRepository;
+ private readonly Lazy _chatPermissionService;
+ private readonly IEventAggregator _eventAggregator;
+
+ public PermissionsService(IPermissionsRepository permissionsRepository, IUsersService usersService, IDepartmentGroupsService departmentGroupsService,
+ Lazy chatChannelRepository = null, Lazy chatPermissionService = null, IEventAggregator eventAggregator = null)
{
_permissionsRepository = permissionsRepository;
_usersService = usersService;
_departmentGroupsService = departmentGroupsService;
+ _chatChannelRepository = chatChannelRepository;
+ _chatPermissionService = chatPermissionService;
+ _eventAggregator = eventAggregator;
}
public async Task> GetAllPermissionsForDepartmentAsync(int departmentId)
@@ -64,15 +75,15 @@ public async Task SetPermissionForDepartmentAsync(int departmentId,
return saved;
}
- private static async Task RotateDispatchChatAccessAsync(int departmentId)
+ private async Task RotateDispatchChatAccessAsync(int departmentId)
{
- if (departmentId <= 0 || !ServiceLocator.IsLocationProviderSet)
+ if (departmentId <= 0 || _chatChannelRepository == null || _chatPermissionService == null || _eventAggregator == null)
return;
try
{
- var channelRepository = ServiceLocator.Current.GetInstance();
- var permissionService = ServiceLocator.Current.GetInstance();
+ var channelRepository = _chatChannelRepository.Value;
+ var permissionService = _chatPermissionService.Value;
var channels = await channelRepository.GetAllByDepartmentIdAsync(departmentId, true);
if (channels != null)
@@ -81,7 +92,7 @@ private static async Task RotateDispatchChatAccessAsync(int departmentId)
await permissionService.InvalidateChannelCacheAsync(channel.ChatChannelId);
}
- ServiceLocator.Current.GetInstance().SendMessage(new ChatEventRaised
+ _eventAggregator.SendMessage(new ChatEventRaised
{
DepartmentId = departmentId,
Kind = ChatEventKinds.ChannelUpdated,
diff --git a/Core/Resgrid.Services/Records/RecordsAnalyticsService.cs b/Core/Resgrid.Services/Records/RecordsAnalyticsService.cs
index b3e43a26d..24dc6033a 100644
--- a/Core/Resgrid.Services/Records/RecordsAnalyticsService.cs
+++ b/Core/Resgrid.Services/Records/RecordsAnalyticsService.cs
@@ -153,6 +153,8 @@ private async Task BeginAsync(int departmentId, string userId, RecordsA
await _gate.RequireEnabledAsync(departmentId, RecordsPreventionModule.Analytics);
await _gate.RequireViewerAsync(departmentId, userId);
query ??= new RecordsAnalyticsQuery();
+ RecordsPreventionGate.RequireStorableDate(query.Start, "The start date is not valid.");
+ RecordsPreventionGate.RequireStorableDate(query.End, "The end date is not valid.");
var now = DateTime.UtcNow;
var end = query.End ?? now;
var start = query.Start ?? end.AddDays(-RecordsAnalyticsLimits.DefaultWindowDays);
diff --git a/Core/Resgrid.Services/Records/RecordsCrrService.cs b/Core/Resgrid.Services/Records/RecordsCrrService.cs
index 4895fba30..c38a5e0e9 100644
--- a/Core/Resgrid.Services/Records/RecordsCrrService.cs
+++ b/Core/Resgrid.Services/Records/RecordsCrrService.cs
@@ -24,10 +24,12 @@ public RecordsCrrService(RecordsPreventionGate gate, IRmsCrrActivitiesRepository
public Task IsModuleEnabledAsync(int departmentId) => _gate.IsEnabledAsync(departmentId, RecordsPreventionModule.Crr);
private async Task RequireViewAsync(int departmentId, string userId) { await _gate.RequireEnabledAsync(departmentId, RecordsPreventionModule.Crr); await _gate.RequireViewerAsync(departmentId, userId); }
private async Task RequireAdminAsync(int departmentId, string userId) { await _gate.RequireEnabledAsync(departmentId, RecordsPreventionModule.Crr); await _gate.RequireAdminAsync(departmentId, userId); }
+ private static void RequireStorableWindow(DateTime startUtc, DateTime endUtc) { RecordsPreventionGate.RequireStorableDate(startUtc, "The start date is not valid."); RecordsPreventionGate.RequireStorableDate(endUtc, "The end date is not valid."); }
public async Task> ListAsync(int departmentId, string userId, DateTime startUtc, DateTime endUtc, int take)
{
await RequireViewAsync(departmentId, userId);
+ RequireStorableWindow(startUtc, endUtc);
return (await _activities.GetForRangeAsync(departmentId, startUtc, endUtc, take))?.ToList() ?? new List();
}
@@ -53,7 +55,7 @@ public async Task SaveAsync(int departmentId, string userId, Rms
var isNew = entity == null;
if (isNew) entity = new RmsCrrActivity { RmsCrrActivityId = Guid.NewGuid().ToString(), DepartmentId = departmentId, ProtectionId = Guid.NewGuid().ToString(), CreatedOn = now, CreatedByUserId = userId, RowVersion = 0 };
entity.Kind = input.Kind == 0 ? (int)RmsCrrActivityKind.PublicEducation : input.Kind;
- entity.OccurredOn = input.OccurredOn == default ? now : input.OccurredOn;
+ entity.OccurredOn = RecordsPreventionGate.RequireStorableDate(input.OccurredOn == default ? now : input.OccurredOn, "The activity date is not valid.");
entity.Title = RecordsPreventionGate.Require(input.Title, 250, "An activity needs a title.");
entity.Description = RecordsPreventionGate.Trim(input.Description, 4000); entity.RmsOccupancyId = RecordsPreventionGate.Trim(input.RmsOccupancyId, 36);
entity.LocationText = RecordsPreventionGate.Trim(input.LocationText, 500); entity.Latitude = input.Latitude; entity.Longitude = input.Longitude;
@@ -78,6 +80,7 @@ public async Task DeleteAsync(int departmentId, string userId, string activityId
public async Task GetSummaryAsync(int departmentId, string userId, DateTime startUtc, DateTime endUtc)
{
await RequireViewAsync(departmentId, userId);
+ RequireStorableWindow(startUtc, endUtc);
var rows = (await _activities.GetForRangeAsync(departmentId, startUtc, endUtc, 2000))?.ToList() ?? new List();
var summary = new CrrSummary { Start = startUtc, End = endUtc, Activities = rows.Count, Audience = rows.Sum(r => r.AudienceCount), SmokeAlarmsInstalled = rows.Sum(r => r.SmokeAlarmsInstalled), Hours = rows.Sum(r => r.HoursSpent) };
foreach (var group in rows.GroupBy(r => r.Kind)) summary.ByKind[group.Key] = group.Count();
diff --git a/Core/Resgrid.Services/Records/RecordsHydrantsService.cs b/Core/Resgrid.Services/Records/RecordsHydrantsService.cs
index 94d49b3db..a0b2fd250 100644
--- a/Core/Resgrid.Services/Records/RecordsHydrantsService.cs
+++ b/Core/Resgrid.Services/Records/RecordsHydrantsService.cs
@@ -125,7 +125,7 @@ public async Task RecordFlowTestAsync(int departmentId, stri
var test = new RmsHydrantFlowTest
{
RmsHydrantFlowTestId = Guid.NewGuid().ToString(), DepartmentId = departmentId, ProtectionId = Guid.NewGuid().ToString(), RmsHydrantId = hydrant.RmsHydrantId,
- TestedOn = input.TestedOn == default ? now : input.TestedOn, TestedByUserId = userId, StaticPressurePsi = input.StaticPressurePsi, ResidualPressurePsi = input.ResidualPressurePsi,
+ TestedOn = RecordsPreventionGate.RequireStorableDate(input.TestedOn == default ? now : input.TestedOn, "The test date is not valid."), TestedByUserId = userId, StaticPressurePsi = input.StaticPressurePsi, ResidualPressurePsi = input.ResidualPressurePsi,
PitotPressurePsi = input.PitotPressurePsi, OutletDiameterInches = input.OutletDiameterInches, Coefficient = coefficient, Notes = RecordsPreventionGate.Trim(input.Notes, 2000), CreatedOn = now
};
test.FlowGpm = HydrantFlowCalculator.FlowGpm(coefficient, input.OutletDiameterInches, input.PitotPressurePsi);
@@ -150,7 +150,7 @@ public async Task RecordMaintenanceAsync(int departmentId
var row = new RmsHydrantMaintenance
{
RmsHydrantMaintenanceId = Guid.NewGuid().ToString(), DepartmentId = departmentId, ProtectionId = Guid.NewGuid().ToString(), RmsHydrantId = hydrant.RmsHydrantId,
- PerformedOn = input.PerformedOn == default ? now : input.PerformedOn, PerformedByUserId = userId, Kind = input.Kind == 0 ? (int)RmsHydrantMaintenanceKind.Inspection : input.Kind,
+ PerformedOn = RecordsPreventionGate.RequireStorableDate(input.PerformedOn == default ? now : input.PerformedOn, "The maintenance date is not valid."), PerformedByUserId = userId, Kind = input.Kind == 0 ? (int)RmsHydrantMaintenanceKind.Inspection : input.Kind,
Notes = RecordsPreventionGate.Trim(input.Notes, 2000), ReturnedToService = input.ReturnedToService, CreatedOn = now
};
await _maintenance.InsertAsync(row, cancellationToken, true);
diff --git a/Core/Resgrid.Services/Records/RecordsInspectionsService.cs b/Core/Resgrid.Services/Records/RecordsInspectionsService.cs
index 80e273b8c..1c4c48df3 100644
--- a/Core/Resgrid.Services/Records/RecordsInspectionsService.cs
+++ b/Core/Resgrid.Services/Records/RecordsInspectionsService.cs
@@ -215,6 +215,7 @@ public static List ParseItems(string json)
public async Task> ListAsync(int departmentId, string userId, RmsInspectionQuery query)
{
await RequireViewAsync(departmentId, userId);
+ RecordsPreventionGate.RequireStorableDate(query?.ScheduledBefore, "The scheduled-before date is not valid.");
var rows = (await _inspections.QueryAsync(departmentId, query ?? new RmsInspectionQuery()))?.ToList() ?? new List();
await _protection.RevealInspectionsAsync(departmentId, rows);
return rows;
@@ -223,6 +224,7 @@ public async Task> ListAsync(int departmentId, string userId
public async Task CountAsync(int departmentId, string userId, RmsInspectionQuery query)
{
await RequireViewAsync(departmentId, userId);
+ RecordsPreventionGate.RequireStorableDate(query?.ScheduledBefore, "The scheduled-before date is not valid.");
return await _inspections.CountAsync(departmentId, query ?? new RmsInspectionQuery());
}
@@ -268,11 +270,13 @@ public async Task ScheduleAsync(int departmentId, string userId,
private async Task CreateScheduledAsync(int departmentId, string userId, RmsOccupancy occupancy, RmsInspectionProgram program, DateTime scheduledOn, string inspectorUserId, string parentInspectionId, CancellationToken cancellationToken)
{
var now = DateTime.UtcNow;
+ // Checked before the initializer, which draws the next inspection number.
+ scheduledOn = RecordsPreventionGate.RequireStorableDate(scheduledOn == default ? now : scheduledOn, "The scheduled date is not valid.");
var inspection = new RmsInspection
{
RmsInspectionId = Guid.NewGuid().ToString(), DepartmentId = departmentId, ProtectionId = Guid.NewGuid().ToString(), RmsOccupancyId = occupancy.RmsOccupancyId, RmsInspectionProgramId = program?.RmsInspectionProgramId,
InspectionNumber = await _gate.NextNumberAsync(departmentId, RmsPreventionNumberKinds.Inspection, now, cancellationToken), State = (int)RmsInspectionState.Scheduled, Result = (int)RmsInspectionResult.NotRecorded,
- ScheduledOn = scheduledOn == default ? now : scheduledOn, InspectorUserId = RecordsPreventionGate.Trim(inspectorUserId, 128), ParentInspectionId = parentInspectionId,
+ ScheduledOn = scheduledOn, InspectorUserId = RecordsPreventionGate.Trim(inspectorUserId, 128), ParentInspectionId = parentInspectionId,
CreatedOn = now, CreatedByUserId = userId, ModifiedOn = now, RowVersion = 1
};
await _inspections.InsertAsync(inspection, cancellationToken, true);
@@ -447,7 +451,7 @@ public async Task SaveViolationAsync(int departmentId, string user
entity.RmsCodeSetId = RecordsPreventionGate.Trim(input.RmsCodeSetId, 36); entity.RmsCodeSectionId = RecordsPreventionGate.Trim(input.RmsCodeSectionId, 36);
entity.Description = RecordsPreventionGate.Require(input.Description, 4000, "A violation needs a description.");
entity.Severity = Math.Clamp(input.Severity == 0 ? 2 : input.Severity, 1, 4); entity.CorrectiveAction = RecordsPreventionGate.Trim(input.CorrectiveAction, 4000);
- entity.DueOn = input.DueOn ?? entity.DueOn ?? now.AddDays(30); entity.ModifiedOn = now; entity.RowVersion++;
+ entity.DueOn = RecordsPreventionGate.RequireStorableDate(input.DueOn, "The due date is not valid.") ?? entity.DueOn ?? now.AddDays(30); entity.ModifiedOn = now; entity.RowVersion++;
var plaintext = PlaintextSnapshot.Take(entity, RmsProtectedFields.Violations);
await _protection.ProtectViolationAsync(departmentId, entity, existing, userId, cancellationToken);
if (existing == null) await _violations.InsertAsync(entity, cancellationToken, true); else await _violations.UpdateAsync(entity, cancellationToken, true);
diff --git a/Core/Resgrid.Services/Records/RecordsInvestigationsService.cs b/Core/Resgrid.Services/Records/RecordsInvestigationsService.cs
index 6cf337174..a30eae242 100644
--- a/Core/Resgrid.Services/Records/RecordsInvestigationsService.cs
+++ b/Core/Resgrid.Services/Records/RecordsInvestigationsService.cs
@@ -259,7 +259,7 @@ public async Task AddNoteAsync(int departmentId, string us
var note = new RmsInvestigationNote
{
RmsInvestigationNoteId = Guid.NewGuid().ToString(), DepartmentId = departmentId, ProtectionId = Guid.NewGuid().ToString(), RmsInvestigationCaseId = caseId, Kind = (int)kind,
- OccurredOn = occurredOn == default ? now : occurredOn, AuthorUserId = userId, Subject = RecordsPreventionGate.Trim(subject, 250), Body = RecordsPreventionGate.Require(body, 32000, "A note needs content."), CreatedOn = now, ModifiedOn = now, RowVersion = 1
+ OccurredOn = RecordsPreventionGate.RequireStorableDate(occurredOn == default ? now : occurredOn, "The note date is not valid."), AuthorUserId = userId, Subject = RecordsPreventionGate.Trim(subject, 250), Body = RecordsPreventionGate.Require(body, 32000, "A note needs content."), CreatedOn = now, ModifiedOn = now, RowVersion = 1
};
var plaintext = PlaintextSnapshot.Take(note, RmsProtectedFields.InvestigationNotes);
await _protection.ProtectInvestigationNoteAsync(departmentId, note, null, userId, cancellationToken);
@@ -294,11 +294,13 @@ public async Task AddEvidenceAsync(int departmentId, s
var (investigation, _) = await RequireMemberAsync(departmentId, userId, caseId, RmsInvestigationRole.Lead, RmsInvestigationRole.Investigator);
RequireOpen(investigation);
var now = DateTime.UtcNow;
+ // Checked before the initializer, which draws the next evidence number.
+ var collectedOn = RecordsPreventionGate.RequireStorableDate(input.CollectedOn == default ? now : input.CollectedOn, "The collection date is not valid.");
var item = new RmsInvestigationEvidence
{
RmsInvestigationEvidenceId = Guid.NewGuid().ToString(), DepartmentId = departmentId, ProtectionId = Guid.NewGuid().ToString(), RmsInvestigationCaseId = caseId,
EvidenceNumber = await _gate.NextNumberAsync(departmentId, RmsPreventionNumberKinds.Evidence, now, cancellationToken), Kind = input.Kind == 0 ? (int)RmsInvestigationEvidenceKind.Physical : input.Kind,
- Description = RecordsPreventionGate.Require(input.Description, 4000, "Evidence needs a description."), CollectedOn = input.CollectedOn == default ? now : input.CollectedOn,
+ Description = RecordsPreventionGate.Require(input.Description, 4000, "Evidence needs a description."), CollectedOn = collectedOn,
CollectedByUserId = string.IsNullOrWhiteSpace(input.CollectedByUserId) ? userId : input.CollectedByUserId.Trim(), CollectedFrom = RecordsPreventionGate.Trim(input.CollectedFrom, 1000),
State = (int)RmsEvidenceState.Collected, StorageLocation = RecordsPreventionGate.Trim(input.StorageLocation, 250), CreatedOn = now, ModifiedOn = now, RowVersion = 1
};
diff --git a/Core/Resgrid.Services/Records/RecordsOccupancyService.cs b/Core/Resgrid.Services/Records/RecordsOccupancyService.cs
index aa8c662b7..0f4323764 100644
--- a/Core/Resgrid.Services/Records/RecordsOccupancyService.cs
+++ b/Core/Resgrid.Services/Records/RecordsOccupancyService.cs
@@ -122,6 +122,7 @@ public async Task SaveAsync(int departmentId, string userId, RmsOc
await _gate.RequireAdminAsync(departmentId, userId);
if (input == null) throw new ArgumentNullException(nameof(input));
input.Name = RecordsPreventionGate.Require(input.Name, 250, "An occupancy needs a name.");
+ input.NextReviewDue = RecordsPreventionGate.RequireStorableDate(input.NextReviewDue, "The next review date is not valid.");
var now = DateTime.UtcNow;
RmsOccupancy entity;
diff --git a/Core/Resgrid.Services/Records/RecordsPermitsService.cs b/Core/Resgrid.Services/Records/RecordsPermitsService.cs
index 2fcfaa90e..6dbbd9a49 100644
--- a/Core/Resgrid.Services/Records/RecordsPermitsService.cs
+++ b/Core/Resgrid.Services/Records/RecordsPermitsService.cs
@@ -64,6 +64,7 @@ public async Task SaveTypeAsync(int departmentId, string userId,
public async Task> ListAsync(int departmentId, string userId, RmsPermitQuery query)
{
await RequireViewAsync(departmentId, userId);
+ RecordsPreventionGate.RequireStorableDate(query?.ExpiresBefore, "The expires-before date is not valid.");
var rows = (await _permits.QueryAsync(departmentId, query ?? new RmsPermitQuery()))?.ToList() ?? new List();
await _protection.RevealPermitsAsync(departmentId, rows);
return rows;
@@ -72,6 +73,7 @@ public async Task> ListAsync(int departmentId, string userId, Rm
public async Task CountAsync(int departmentId, string userId, RmsPermitQuery query)
{
await RequireViewAsync(departmentId, userId);
+ RecordsPreventionGate.RequireStorableDate(query?.ExpiresBefore, "The expires-before date is not valid.");
return await _permits.CountAsync(departmentId, query ?? new RmsPermitQuery());
}
@@ -134,7 +136,7 @@ public async Task UpdateAsync(int departmentId, string userId, RmsPer
permit.ApplicantContactId = RecordsPreventionGate.Trim(input.ApplicantContactId, 128); permit.ApplicantName = RecordsPreventionGate.Trim(input.ApplicantName, 250); permit.ApplicantPhone = RecordsPreventionGate.Trim(input.ApplicantPhone, 50); permit.ApplicantEmail = RecordsPreventionGate.Trim(input.ApplicantEmail, 250);
permit.Description = RecordsPreventionGate.Trim(input.Description, 4000); permit.Conditions = RecordsPreventionGate.Trim(input.Conditions, 8000); permit.ReviewNotes = RecordsPreventionGate.Trim(input.ReviewNotes, 8000);
permit.RmsOccupancyId = RecordsPreventionGate.Trim(input.RmsOccupancyId, 36); permit.FeeAmount = input.FeeAmount;
- if (input.ExpiresOn.HasValue) permit.ExpiresOn = input.ExpiresOn; if (input.EffectiveOn.HasValue) permit.EffectiveOn = input.EffectiveOn;
+ if (input.ExpiresOn.HasValue) permit.ExpiresOn = RecordsPreventionGate.RequireStorableDate(input.ExpiresOn, "The expiry date is not valid."); if (input.EffectiveOn.HasValue) permit.EffectiveOn = RecordsPreventionGate.RequireStorableDate(input.EffectiveOn, "The effective date is not valid.");
permit.ModifiedOn = DateTime.UtcNow; permit.RowVersion++;
var plaintext = PlaintextSnapshot.Take(permit, RmsProtectedFields.Permits);
await _protection.ProtectPermitAsync(departmentId, permit, existing, userId, cancellationToken);
@@ -175,8 +177,8 @@ public async Task TransitionAsync(int departmentId, string userId, st
case RmsPermitState.Approved: permit.ReviewedOn = now; permit.ReviewedByUserId = userId; break;
case RmsPermitState.Denied: permit.ReviewedOn = now; permit.ReviewedByUserId = userId; permit.DecisionReason = RecordsPreventionGate.Trim(reason, 1000); break;
case RmsPermitState.Issued:
- permit.IssuedOn = now; permit.IssuedByUserId = userId; permit.EffectiveOn = effectiveOn ?? permit.EffectiveOn ?? now;
- permit.ExpiresOn = expiresOn ?? permit.ExpiresOn ?? permit.EffectiveOn.Value.AddDays(type?.DefaultValidityDays ?? 365);
+ permit.IssuedOn = now; permit.IssuedByUserId = userId; permit.EffectiveOn = RecordsPreventionGate.RequireStorableDate(effectiveOn ?? permit.EffectiveOn ?? now, "The effective date is not valid.");
+ permit.ExpiresOn = RecordsPreventionGate.RequireStorableDate(expiresOn ?? permit.ExpiresOn ?? permit.EffectiveOn.Value.AddDays(type?.DefaultValidityDays ?? 365), "The expiry date is not valid.");
if (permit.ExpiresOn <= permit.EffectiveOn) throw new ArgumentException("The expiry must fall after the effective date.");
break;
case RmsPermitState.Revoked: permit.DecisionReason = RecordsPreventionGate.Trim(reason, 1000); break;
diff --git a/Core/Resgrid.Services/Records/RecordsPreventionGate.cs b/Core/Resgrid.Services/Records/RecordsPreventionGate.cs
index 31db47995..eb3b9ca2d 100644
--- a/Core/Resgrid.Services/Records/RecordsPreventionGate.cs
+++ b/Core/Resgrid.Services/Records/RecordsPreventionGate.cs
@@ -1,4 +1,5 @@
using System;
+using System.Data.SqlTypes;
using System.Threading;
using System.Threading.Tasks;
using Newtonsoft.Json;
@@ -128,5 +129,18 @@ public static string Require(string value, int max, string message)
throw new ArgumentException($"{message} (at most {max} characters).");
return trimmed;
}
+
+ ///
+ /// Dapper binds DateTime as SQL datetime, which starts at 1753, so a client date outside it (a two-digit year sent as
+ /// 0026) failed the whole statement with SqlDateTime overflow. Refuse it here so the web and API both get a message.
+ ///
+ public static DateTime RequireStorableDate(DateTime value, string message)
+ {
+ if (value < (DateTime)SqlDateTime.MinValue || value > (DateTime)SqlDateTime.MaxValue)
+ throw new ArgumentException(message);
+ return value;
+ }
+
+ public static DateTime? RequireStorableDate(DateTime? value, string message) => value.HasValue ? RequireStorableDate(value.Value, message) : null;
}
}
diff --git a/Core/Resgrid.Services/Records/RecordsQualityReviewService.cs b/Core/Resgrid.Services/Records/RecordsQualityReviewService.cs
index e2dc2f7ea..45bfdfe8c 100644
--- a/Core/Resgrid.Services/Records/RecordsQualityReviewService.cs
+++ b/Core/Resgrid.Services/Records/RecordsQualityReviewService.cs
@@ -108,6 +108,7 @@ public static List ParseFindings(string json)
public async Task> SampleAsync(int departmentId, string userId, string rubricId, DateTime sinceUtc, CancellationToken cancellationToken = default)
{
await RequireReviewerAsync(departmentId, userId);
+ RecordsPreventionGate.RequireStorableDate(sinceUtc, "The sample start date is not valid.");
var rubric = await _rubrics.GetByIdForDepartmentAsync(departmentId, rubricId);
if (rubric == null || rubric.DeletedOn != null || !rubric.IsActive) throw new ArgumentException("Choose an active rubric.");
var criteria = ParseCriteria(rubric.CriteriaJson);
@@ -204,6 +205,7 @@ public async Task> GetForRecordAsync(int departmentId, st
public async Task GetTrendsAsync(int departmentId, string userId, DateTime sinceUtc)
{
await RequireReviewerAsync(departmentId, userId);
+ RecordsPreventionGate.RequireStorableDate(sinceUtc, "The start date is not valid.");
var scored = ((await _reviews.GetScoredSinceAsync(departmentId, sinceUtc, 5000)) ?? Enumerable.Empty()).Where(r => r.Score.HasValue).ToList();
var trends = new RecordsQualityTrends { Since = sinceUtc, Scored = scored.Count, Sampled = scored.Count + ((await _reviews.GetPendingAsync(departmentId, 1000))?.Count() ?? 0), AverageScore = scored.Count == 0 ? 0 : Math.Round(scored.Average(r => r.Score.Value), 1) };
List Rows(Func key) => scored.Where(r => key(r) != null).GroupBy(key).Select(g => new RecordsQualityTrendRow { Key = g.Key, Label = g.Key, Reviews = g.Count(), AverageScore = Math.Round(g.Average(r => r.Score.Value), 1), AmendmentsRecommended = g.Count(r => r.AmendmentRecommended) }).OrderBy(r => r.AverageScore).ToList();
diff --git a/Core/Resgrid.Services/Search/SystemActionCatalog.cs b/Core/Resgrid.Services/Search/SystemActionCatalog.cs
index a2d3e46fd..0621f02a3 100644
--- a/Core/Resgrid.Services/Search/SystemActionCatalog.cs
+++ b/Core/Resgrid.Services/Search/SystemActionCatalog.cs
@@ -51,7 +51,7 @@ public static class SystemActionCatalog
private const string Log = "Log";
private static SystemActionDefinition Nav(string key, string title, string description, string path, string[] keywords = null,
- string claimResource = null, string claimAction = null, string module = null, string flag = null, bool adminOnly = false, bool hiddenWhenRecords = false)
+ string claimResource = null, string claimAction = null, string module = null, string flag = null, bool adminOnly = false, bool hiddenAfterCutover = false)
{
return new SystemActionDefinition
{
@@ -66,14 +66,14 @@ private static SystemActionDefinition Nav(string key, string title, string descr
Module = module,
FeatureFlag = flag,
DepartmentAdminOnly = adminOnly,
- HiddenWhenRecordsEnabled = hiddenWhenRecords
+ HiddenAfterRecordsCutover = hiddenAfterCutover
};
}
private static SystemActionDefinition Act(string key, string title, string description, string path, string category, string[] keywords = null,
- string claimResource = null, string claimAction = null, string module = null, string flag = null, bool adminOnly = false, bool hiddenWhenRecords = false)
+ string claimResource = null, string claimAction = null, string module = null, string flag = null, bool adminOnly = false, bool hiddenAfterCutover = false)
{
- var d = Nav(key, title, description, path, keywords, claimResource, claimAction, module, flag, adminOnly, hiddenWhenRecords);
+ var d = Nav(key, title, description, path, keywords, claimResource, claimAction, module, flag, adminOnly, hiddenAfterCutover);
d.Category = category;
return d;
}
@@ -120,8 +120,9 @@ private static SystemActionDefinition Act(string key, string title, string descr
Act("new-calendar-item", "New Calendar Event", "Create a calendar event", "/User/Calendar/New", SystemActionCategories.Create, new[] { "event", "meeting", "schedule" }, Schedule, Create, SystemActionModules.Calendar),
// ---- Logs (legacy) / Records
- Nav("logs", "Logs", "Run, training, work and meeting logs", "/User/Logs", new[] { "run log", "activity", "reports", "training log", "work log" }, Log, View, SystemActionModules.Logs, hiddenWhenRecords: true),
- Act("new-log", "New Log", "Create a run report, training log or work log", "/User/Logs/NewLog", SystemActionCategories.Create, new[] { "run report", "training log", "work log" }, Log, Create, SystemActionModules.Logs, hiddenWhenRecords: true),
+ // Logs stay findable before and after the Records cutover (read-only after it); only creating one goes away.
+ Nav("logs", "Logs", "Run, training, work and meeting logs", "/User/Logs", new[] { "run log", "activity", "reports", "training log", "work log", "legacy logs", "existing logs" }, Log, View, SystemActionModules.Logs),
+ Act("new-log", "New Log", "Create a run report, training log or work log", "/User/Logs/NewLog", SystemActionCategories.Create, new[] { "run report", "training log", "work log" }, Log, Create, SystemActionModules.Logs, hiddenAfterCutover: true),
Nav("records", "Records", "Records queue: run reports, training and operational records", "/User/Records", new[] { "rms", "run reports", "incident reports", "neris", "logs" }, Record, View, SystemActionModules.Logs, FeatureFlagKeys.RecordsSystem),
Nav("records-dashboard", "Records Dashboard", "Records due, submissions and quality at a glance", "/User/Records/Dashboard", new[] { "rms", "overview", "due" }, Record, View, SystemActionModules.Logs, FeatureFlagKeys.RecordsSystem),
Act("records-settings", "Records Settings", "Lifecycle, numbering, search, retention and visibility settings for Records", "/User/Records/Settings", SystemActionCategories.Manage, new[] { "rms settings", "retention", "numbering" }, Record, View, SystemActionModules.Logs, FeatureFlagKeys.RecordsSystem, adminOnly: true),
diff --git a/Core/Resgrid.Services/Search/SystemActionsService.cs b/Core/Resgrid.Services/Search/SystemActionsService.cs
index c11da0ebc..b1a5864fe 100644
--- a/Core/Resgrid.Services/Search/SystemActionsService.cs
+++ b/Core/Resgrid.Services/Search/SystemActionsService.cs
@@ -20,10 +20,12 @@ namespace Resgrid.Services.Search
public class SystemActionsService : ISystemActionsService
{
private readonly IFeatureToggleService _featureToggles;
+ private readonly IRecordsCutoverService _recordsCutover;
- public SystemActionsService(IFeatureToggleService featureToggles)
+ public SystemActionsService(IFeatureToggleService featureToggles, IRecordsCutoverService recordsCutover)
{
_featureToggles = featureToggles ?? throw new ArgumentNullException(nameof(featureToggles));
+ _recordsCutover = recordsCutover ?? throw new ArgumentNullException(nameof(recordsCutover));
}
public async Task> SearchAsync(string text, SearchPrincipal principal, int max = 8, CancellationToken cancellationToken = default)
@@ -76,7 +78,18 @@ async Task FlagAsync(string key)
return value;
}
- var recordsOn = await FlagAsync(FeatureFlagKeys.RecordsSystem);
+ // The cutover, not the Records.System flag, is what makes legacy Logs read-only: a department with the flag
+ // on but Records not yet activated is still writing Logs.
+ bool? legacyWritesBlocked = null;
+ async Task LegacyWritesBlockedAsync()
+ {
+ if (legacyWritesBlocked.HasValue)
+ return legacyWritesBlocked.Value;
+ try { legacyWritesBlocked = await _recordsCutover.AreLegacyWritesBlockedAsync(principal.DepartmentId); }
+ catch (Exception ex) { Logging.LogException(ex, "Records cutover state could not be evaluated for the command palette; hiding legacy Logs writes."); legacyWritesBlocked = true; }
+ return legacyWritesBlocked.Value;
+ }
+
var allowed = new List();
foreach (var def in SystemActionCatalog.All)
{
@@ -87,7 +100,7 @@ async Task FlagAsync(string key)
continue;
if (!principal.ModuleEnabled(def.Module))
continue;
- if (def.HiddenWhenRecordsEnabled && recordsOn)
+ if (def.HiddenAfterRecordsCutover && await LegacyWritesBlockedAsync())
continue;
if (!await FlagAsync(def.FeatureFlag))
continue;
diff --git a/Core/Resgrid.Services/Search/UnifiedSearchService.cs b/Core/Resgrid.Services/Search/UnifiedSearchService.cs
index c7766a13d..54463c2e2 100644
--- a/Core/Resgrid.Services/Search/UnifiedSearchService.cs
+++ b/Core/Resgrid.Services/Search/UnifiedSearchService.cs
@@ -401,7 +401,9 @@ private async Task EnsureStateAsync(int departmentId, CancellationToken cancella
if (existing != null)
return;
var now = DateTime.UtcNow;
- await _states.SaveOrUpdateAsync(new SearchIndexState
+ // A department's first searches arrive together (typeahead sends one per keystroke) and all see no row
+ // above, so the create has to be conditional or every request but one fails on the unique index.
+ await _states.InsertIfMissingAsync(new SearchIndexState
{
IndexName = SearchIndexNames.Global,
DepartmentId = departmentId,
@@ -411,7 +413,7 @@ await _states.SaveOrUpdateAsync(new SearchIndexState
RebuildRequestedOn = now,
CreatedOn = now,
ModifiedOn = now
- }, cancellationToken, true);
+ }, cancellationToken);
}
catch (Exception ex)
{
diff --git a/Providers/Resgrid.Providers.Bus/NotificationProvider.cs b/Providers/Resgrid.Providers.Bus/NotificationProvider.cs
index c268e7998..539c206e0 100644
--- a/Providers/Resgrid.Providers.Bus/NotificationProvider.cs
+++ b/Providers/Resgrid.Providers.Bus/NotificationProvider.cs
@@ -16,9 +16,20 @@ namespace Resgrid.Providers.Bus
{
public class NotificationProvider : INotificationProvider
{
+ ///
+ /// The Azure user hub is legacy (Novu delivers user pushes) and is left unconfigured in newer
+ /// deployments. The client factory throws on an empty connection string, and there is nothing to
+ /// register with or remove from a hub that was never set up, so the methods that would otherwise
+ /// throw to their callers return early instead.
+ ///
+ private static bool IsHubConfigured()
+ {
+ return !String.IsNullOrWhiteSpace(Config.ServiceBusConfig.AzureNotificationHub_FullConnectionString);
+ }
+
public async Task RegisterPush(PushUri pushUri)
{
- if (String.IsNullOrWhiteSpace(pushUri.DeviceId))
+ if (String.IsNullOrWhiteSpace(pushUri.DeviceId) || !IsHubConfigured())
return;
var hubClient = NotificationHubClient.CreateClientFromConnectionString(Config.ServiceBusConfig.AzureNotificationHub_FullConnectionString, Config.ServiceBusConfig.AzureNotificationHub_PushUrl);
@@ -98,6 +109,9 @@ public async Task RegisterPush(PushUri pushUri)
public async Task UnRegisterPush(PushUri pushUri)
{
+ if (!IsHubConfigured())
+ return;
+
var hubClient = NotificationHubClient.CreateClientFromConnectionString(Config.ServiceBusConfig.AzureNotificationHub_FullConnectionString, Config.ServiceBusConfig.AzureNotificationHub_PushUrl);
var registrations = await hubClient.GetRegistrationsByTagAsync(string.Format("userId:{0}", pushUri.UserId), 50);
@@ -121,6 +135,9 @@ public async Task UnRegisterPush(PushUri pushUri)
public async Task UnRegisterPushByUserDeviceId(PushUri pushUri)
{
+ if (!IsHubConfigured())
+ return;
+
var hubClient = NotificationHubClient.CreateClientFromConnectionString(Config.ServiceBusConfig.AzureNotificationHub_FullConnectionString, Config.ServiceBusConfig.AzureNotificationHub_PushUrl);
var registrations = await hubClient.GetRegistrationsByTagAsync(string.Format("userId:{0}", pushUri.UserId), 50);
diff --git a/Providers/Resgrid.Providers.Bus/UnitNotificationProvider.cs b/Providers/Resgrid.Providers.Bus/UnitNotificationProvider.cs
index f210ca3fd..83235e0bd 100644
--- a/Providers/Resgrid.Providers.Bus/UnitNotificationProvider.cs
+++ b/Providers/Resgrid.Providers.Bus/UnitNotificationProvider.cs
@@ -15,9 +15,20 @@ namespace Resgrid.Providers.Bus
{
public class UnitNotificationProvider : IUnitNotificationProvider
{
+ ///
+ /// The Azure unit hub is legacy (Novu delivers unit pushes) and is left unconfigured in newer
+ /// deployments. The client factory throws on an empty connection string, and there is nothing to
+ /// register with or remove from a hub that was never set up, so the methods that would otherwise
+ /// throw to their callers return early instead.
+ ///
+ private static bool IsHubConfigured()
+ {
+ return !String.IsNullOrWhiteSpace(Config.ServiceBusConfig.AzureUnitNotificationHub_FullConnectionString);
+ }
+
public async Task RegisterPush(PushUri pushUri)
{
- if (String.IsNullOrWhiteSpace(pushUri.DeviceId))
+ if (String.IsNullOrWhiteSpace(pushUri.DeviceId) || !IsHubConfigured())
return;
if (pushUri.UnitId.HasValue)
@@ -98,6 +109,9 @@ public async Task RegisterPush(PushUri pushUri)
public async Task UnRegisterPush(PushUri pushUri)
{
+ if (!IsHubConfigured())
+ return;
+
var hubClient = NotificationHubClient.CreateClientFromConnectionString(Config.ServiceBusConfig.AzureUnitNotificationHub_FullConnectionString, Config.ServiceBusConfig.AzureUnitNotificationHub_PushUrl);
var registrations = await hubClient.GetRegistrationsByTagAsync(string.Format("deviceId:{0}", pushUri.DeviceId), 50);
@@ -121,6 +135,9 @@ public async Task UnRegisterPush(PushUri pushUri)
public async Task UnRegisterPushByUserDeviceId(PushUri pushUri)
{
+ if (!IsHubConfigured())
+ return;
+
var hubClient = NotificationHubClient.CreateClientFromConnectionString(Config.ServiceBusConfig.AzureUnitNotificationHub_FullConnectionString, Config.ServiceBusConfig.AzureUnitNotificationHub_PushUrl);
var registrations = await hubClient.GetRegistrationsByTagAsync(string.Format("userId:{0}", pushUri.UserId), 50);
@@ -155,6 +172,9 @@ public async Task UnRegisterPushByUserDeviceId(PushUri pushUri)
public async Task UnRegisterPushByUUID(string uuid)
{
+ if (!IsHubConfigured())
+ return;
+
var hubClient = NotificationHubClient.CreateClientFromConnectionString(Config.ServiceBusConfig.AzureUnitNotificationHub_FullConnectionString, Config.ServiceBusConfig.AzureUnitNotificationHub_PushUrl);
var registrations = await hubClient.GetRegistrationsByTagAsync(string.Format("uuid:{0}", uuid), 50);
diff --git a/Providers/Resgrid.Providers.Bus/WorkflowEventProvider.cs b/Providers/Resgrid.Providers.Bus/WorkflowEventProvider.cs
index 371dd2286..b0321a6e1 100644
--- a/Providers/Resgrid.Providers.Bus/WorkflowEventProvider.cs
+++ b/Providers/Resgrid.Providers.Bus/WorkflowEventProvider.cs
@@ -2,6 +2,7 @@
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
+using Autofac;
using Newtonsoft.Json;
using Resgrid.Config;
using Resgrid.Model;
@@ -18,18 +19,21 @@ namespace Resgrid.Providers.Bus
/// Subscribes to all domain events and, for each active workflow whose trigger matches,
/// creates a WorkflowRun (Pending) and enqueues a WorkflowQueueItem to RabbitMQ.
/// Free-plan departments are subject to an aggressive, non-bypassable rate limit.
+ ///
+ /// The scoped repositories and services are NOT constructor-injected: this is a singleton, and capturing
+ /// InstancePerLifetimeScope dependencies in it would pin them to the root scope — one shared unit of work,
+ /// and one shared DB connection, for every event the process ever handles. Events arrive concurrently
+ /// (async listeners, and several worker jobs dispatching at once), and a Records run insert opens a
+ /// transaction on that unit of work, so the other handlers' queries landed on the same connection
+ /// ("The connection does not support MultipleActiveResultSets") and inside another event's transaction.
+ /// Each event instead runs in its own child scope, as ChatProvisioningEventService does.
///
public class WorkflowEventProvider : IWorkflowEventProvider
{
private readonly IEventAggregator _eventAggregator;
private static IOutboundQueueProvider _outboundQueueProvider;
- private static IWorkflowRepository _workflowRepository;
- private static IWorkflowRunRepository _runRepository;
- private static IDepartmentsService _departmentsService;
- private static ISubscriptionsService _subscriptionsService;
- private static IProtectedProjectionService _protectedProjectionService;
- private static Lazy _history;
- private static Task ProtectChecklistRunAsync(WorkflowRun run) => (_history?.Value ?? throw new InvalidOperationException("Readiness history protection is unavailable."))
+ private static ILifetimeScope _lifetimeScope;
+ private static Task ProtectChecklistRunAsync(ILifetimeScope scope, WorkflowRun run) => scope.Resolve()
.ProtectAsync(run.DepartmentId, run.WorkflowRunId, run, ReadinessHistoryFields.Runs);
// Per-minute rate limit tracker: departmentId → (window start, count)
@@ -55,20 +59,11 @@ private static readonly System.Collections.Generic.HashSet history = null)
+ ILifetimeScope lifetimeScope)
{
_eventAggregator = eventAggregator;
_outboundQueueProvider = outboundQueueProvider;
- _workflowRepository = workflowRepository;
- _runRepository = runRepository;
- _departmentsService = departmentsService;
- _subscriptionsService = subscriptionsService;
- _protectedProjectionService = protectedProjectionService;
- _history = history;
+ _lifetimeScope = lifetimeScope;
RegisterListeners();
}
@@ -206,24 +201,30 @@ private static async Task HandleEventAsync(int departmentId, WorkflowTriggerEven
{
try
{
+ // Own scope per event: fresh repositories and services with their own unit of work (see the class summary).
+ using var scope = _lifetimeScope.BeginLifetimeScope();
+ var workflowRepository = scope.Resolve();
+ var runRepository = scope.Resolve();
+
System.Collections.Generic.List workflows = null;
string payloadJson = null;
if (envelope != null)
{
// The skip rows below need the workflow list, so for Records events it is loaded before the limits.
- workflows = (await _workflowRepository.GetAllActiveByDepartmentAndEventTypeAsync(departmentId, (int)eventType))?.ToList();
+ workflows = (await workflowRepository.GetAllActiveByDepartmentAndEventTypeAsync(departmentId, (int)eventType))?.ToList();
if (workflows == null || workflows.Count == 0)
return;
+ var protectedProjectionService = scope.Resolve();
payloadJson = ChecklistWorkflowPayload.IsReadinessProducer(envelope.ProducerSubsystem)
- ? await ChecklistWorkflowPayload.ProjectAsync(departmentId, eventObj, _protectedProjectionService, wrapped: true)
- : await _protectedProjectionService.BuildSafeWorkflowPayloadAsync(departmentId, eventObj);
+ ? await ChecklistWorkflowPayload.ProjectAsync(departmentId, eventObj, protectedProjectionService, wrapped: true)
+ : await protectedProjectionService.BuildSafeWorkflowPayloadAsync(departmentId, eventObj);
if (ChecklistWorkflowPayload.IsReadinessProducer(envelope.ProducerSubsystem))
{
// Retry a run whose database insert succeeded but whose queue send did not.
// Its stable run ID is claimed atomically by the worker before executing actions.
- var existingRuns = (await _runRepository.GetByWorkflowsAndEventAsync(departmentId, workflows.Select(w => w.WorkflowId).ToArray(), envelope.EventId))
+ var existingRuns = (await runRepository.GetByWorkflowsAndEventAsync(departmentId, workflows.Select(w => w.WorkflowId).ToArray(), envelope.EventId))
.GroupBy(r => r.WorkflowId).ToDictionary(g => g.Key, g => g.OrderBy(r => r.StartedOn).First());
foreach (var workflow in workflows.ToList())
{
@@ -237,7 +238,7 @@ private static async Task HandleEventAsync(int departmentId, WorkflowTriggerEven
}
// ── Plan-aware rate limiting ─────────────────────────────────────────
- var plan = await _subscriptionsService.GetCurrentPlanForDepartmentAsync(departmentId);
+ var plan = await scope.Resolve().GetCurrentPlanForDepartmentAsync(departmentId);
var isFreePlan = plan?.IsFree ?? false;
if (isFreePlan)
@@ -245,14 +246,14 @@ private static async Task HandleEventAsync(int departmentId, WorkflowTriggerEven
// Free plan: aggressive per-minute limit with NO event-type exemptions
if (!IsWithinRateLimit(departmentId, WorkflowConfig.FreePlanRateLimitPerDepartmentPerMinute))
{
- await RecordSkippedAsync(workflows, departmentId, eventType, envelope, payloadJson, WorkflowRunSkipReasons.RateLimit);
+ await RecordSkippedAsync(scope, workflows, departmentId, eventType, envelope, payloadJson, WorkflowRunSkipReasons.RateLimit);
return;
}
// Free plan: daily run cap
if (!IsWithinDailyLimit(departmentId, WorkflowConfig.FreePlanDailyRunLimit))
{
- await RecordSkippedAsync(workflows, departmentId, eventType, envelope, payloadJson, WorkflowRunSkipReasons.DailyLimit);
+ await RecordSkippedAsync(scope, workflows, departmentId, eventType, envelope, payloadJson, WorkflowRunSkipReasons.DailyLimit);
return;
}
}
@@ -262,14 +263,14 @@ private static async Task HandleEventAsync(int departmentId, WorkflowTriggerEven
if (!_rateLimitExemptEventTypes.Contains(eventType) &&
!IsWithinRateLimit(departmentId, WorkflowConfig.RateLimitPerDepartmentPerMinute))
{
- await RecordSkippedAsync(workflows, departmentId, eventType, envelope, payloadJson, WorkflowRunSkipReasons.RateLimit);
+ await RecordSkippedAsync(scope, workflows, departmentId, eventType, envelope, payloadJson, WorkflowRunSkipReasons.RateLimit);
return;
}
}
// ── End rate limiting ────────────────────────────────────────────────
if (workflows == null)
- workflows = (await _workflowRepository.GetAllActiveByDepartmentAndEventTypeAsync(departmentId, (int)eventType))?.ToList();
+ workflows = (await workflowRepository.GetAllActiveByDepartmentAndEventTypeAsync(departmentId, (int)eventType))?.ToList();
if (workflows == null || workflows.Count == 0) return;
@@ -277,13 +278,13 @@ private static async Task HandleEventAsync(int departmentId, WorkflowTriggerEven
// redacted HERE, before it reaches WorkflowRun.InputPayload, the queue, retries,
// dead letters, history, or designer previews.
if (payloadJson == null)
- payloadJson = await _protectedProjectionService.BuildSafeWorkflowPayloadAsync(departmentId, eventObj);
- var department = await _departmentsService.GetDepartmentByIdAsync(departmentId);
+ payloadJson = await scope.Resolve().BuildSafeWorkflowPayloadAsync(departmentId, eventObj);
+ var department = await scope.Resolve().GetDepartmentByIdAsync(departmentId);
var deptCode = department?.Code ?? string.Empty;
foreach (var workflow in workflows)
{
- if (!ChecklistWorkflowPayload.IsReadinessProducer(envelope?.ProducerSubsystem) && await IsDuplicateAsync(workflow.WorkflowId, envelope)) continue;
+ if (!ChecklistWorkflowPayload.IsReadinessProducer(envelope?.ProducerSubsystem) && await IsDuplicateAsync(runRepository, workflow.WorkflowId, envelope)) continue;
var run = WorkflowRunEnvelope.Apply(new WorkflowRun
{
@@ -300,14 +301,14 @@ private static async Task HandleEventAsync(int departmentId, WorkflowTriggerEven
try
{
- if (ChecklistWorkflowPayload.IsReadinessProducer(envelope?.ProducerSubsystem)) await ProtectChecklistRunAsync(run);
- run = await _runRepository.InsertAsync(run, CancellationToken.None);
+ if (ChecklistWorkflowPayload.IsReadinessProducer(envelope?.ProducerSubsystem)) await ProtectChecklistRunAsync(scope, run);
+ run = await runRepository.InsertAsync(run, CancellationToken.None);
}
catch (Exception ex) when (envelope != null && WorkflowRunEnvelope.IsDuplicateKeyViolation(ex))
{
if (ChecklistWorkflowPayload.IsReadinessProducer(envelope.ProducerSubsystem))
{
- var existing = await _runRepository.GetByWorkflowAndEventAsync(workflow.WorkflowId, envelope.EventId);
+ var existing = await runRepository.GetByWorkflowAndEventAsync(workflow.WorkflowId, envelope.EventId);
if (existing == null || existing.DepartmentId != departmentId) throw;
if (existing.Status == (int)WorkflowRunStatus.Pending) await RequeueChecklistRunAsync(existing, payloadJson);
}
@@ -353,12 +354,12 @@ private static async Task RequeueChecklistRunAsync(WorkflowRun run, string safeP
}
/// One initial run per (WorkflowId, EventId): a retry or a second dispatcher reuses the existing run.
- private static async Task IsDuplicateAsync(string workflowId, DomainEventDispatchedEvent envelope)
+ private static async Task IsDuplicateAsync(IWorkflowRunRepository runRepository, string workflowId, DomainEventDispatchedEvent envelope)
{
if (envelope == null || string.IsNullOrWhiteSpace(envelope.EventId))
return false;
- var existing = await _runRepository.GetByWorkflowAndEventAsync(workflowId, envelope.EventId);
+ var existing = await runRepository.GetByWorkflowAndEventAsync(workflowId, envelope.EventId);
if (existing == null)
return false;
@@ -371,16 +372,17 @@ private static async Task IsDuplicateAsync(string workflowId, DomainEventD
/// reason, so it shows in run history and health instead of vanishing (plan section 5.6). Legacy events
/// (no envelope) are still dropped silently, as before.
///
- private static async Task RecordSkippedAsync(System.Collections.Generic.List workflows, int departmentId, WorkflowTriggerEventType eventType,
+ private static async Task RecordSkippedAsync(ILifetimeScope scope, System.Collections.Generic.List workflows, int departmentId, WorkflowTriggerEventType eventType,
DomainEventDispatchedEvent envelope, string payloadJson, string reason)
{
if (envelope == null || workflows == null)
return;
+ var runRepository = scope.Resolve();
var now = DateTime.UtcNow;
foreach (var workflow in workflows)
{
- if (await IsDuplicateAsync(workflow.WorkflowId, envelope))
+ if (await IsDuplicateAsync(runRepository, workflow.WorkflowId, envelope))
continue;
var run = WorkflowRunEnvelope.MarkSkipped(WorkflowRunEnvelope.Apply(new WorkflowRun
@@ -397,8 +399,8 @@ private static async Task RecordSkippedAsync(System.Collections.Generic.List IncrementAsync(string cacheKey, TimeSpan expiration)
return 0;
}
+ public async Task GetOrAddStringAsync(string cacheKey, string valueIfAbsent, TimeSpan slidingExpiration)
+ {
+ try
+ {
+ if (Config.SystemBehaviorConfig.CacheEnabled && _connection != null && _connection.IsConnected && valueIfAbsent != null)
+ {
+ IDatabase cache = _connection.GetDatabase();
+ var key = SetCacheKeyForEnv(cacheKey);
+
+ // GET-or-SET + PEXPIRE in a single server-side script so callers racing on a missing key
+ // all read back the one value that was stored, and a hit slides the TTL forward.
+ const string getOrAddScript =
+ "local current = redis.call('GET', KEYS[1])\n" +
+ "if current then\n" +
+ " redis.call('PEXPIRE', KEYS[1], ARGV[2])\n" +
+ " return current\n" +
+ "end\n" +
+ "redis.call('SET', KEYS[1], ARGV[1], 'PX', ARGV[2])\n" +
+ "return ARGV[1]";
+
+ var result = await cache.ScriptEvaluateAsync(
+ getOrAddScript,
+ new RedisKey[] { key },
+ new RedisValue[] { valueIfAbsent, (long)slidingExpiration.TotalMilliseconds });
+
+ return (string)result;
+ }
+ }
+ catch (TimeoutException)
+ { }
+ catch (RedisConnectionException ex)
+ {
+ Logging.LogError(ex);
+ }
+ catch (Exception ex)
+ {
+ Logging.LogException(ex);
+ }
+
+ return null;
+ }
+
public bool IsConnected()
{
return _connection?.IsConnected ?? false;
diff --git a/Repositories/Resgrid.Repositories.DataRepository/SearchRepositories.cs b/Repositories/Resgrid.Repositories.DataRepository/SearchRepositories.cs
index e33879fa7..cdf36b81d 100644
--- a/Repositories/Resgrid.Repositories.DataRepository/SearchRepositories.cs
+++ b/Repositories/Resgrid.Repositories.DataRepository/SearchRepositories.cs
@@ -161,6 +161,49 @@ public Task> GetAllForIndexAsync(string indexName)
$"SELECT * FROM {Tbl("SearchIndexStates")} WHERE {Col("IndexName")} = {P}IndexName ORDER BY {Col("DepartmentId")} ASC",
new { IndexName = indexName });
}
+
+ public async Task InsertIfMissingAsync(SearchIndexState state, CancellationToken cancellationToken = default)
+ {
+ if (state == null) throw new ArgumentNullException(nameof(state));
+ Utf8WriteGuard.Sanitize(state);
+
+ // The race is absorbed in-statement rather than caught: RunAsync logs every exception before rethrowing,
+ // so a caught unique violation would still reach Sentry. PostgreSQL resolves it with ON CONFLICT DO NOTHING;
+ // on SQL Server UPDLOCK, HOLDLOCK takes a key-range lock on the unique index so a second inserter waits for
+ // the first to commit and then sees its row, instead of both passing NOT EXISTS.
+ var columns = Cols("IndexName", "DepartmentId", "SchemaVersion", "ProtectedCatalogVersion", "PolicyEpoch", "Generation", "State",
+ "DocumentCount", "LastRebuiltOn", "LastIndexedModifiedOn", "RebuildRequestedOn", "CreatedOn", "ModifiedOn");
+ var values = $"{P}IndexName, {P}DepartmentId, {P}SchemaVersion, {P}ProtectedCatalogVersion, {P}PolicyEpoch, {P}Generation, {P}State, " +
+ $"{P}DocumentCount, {P}LastRebuiltOn, {P}LastIndexedModifiedOn, {P}RebuildRequestedOn, {P}CreatedOn, {P}ModifiedOn";
+ var sql = IsPostgres
+ ? $"INSERT INTO {Tbl("SearchIndexStates")} ({columns}) VALUES ({values}) " +
+ $"ON CONFLICT ({Cols("IndexName", "DepartmentId")}) DO NOTHING RETURNING {Col("SearchIndexStateId")}"
+ : $"INSERT INTO {Tbl("SearchIndexStates")} ({columns}) OUTPUT INSERTED.{Col("SearchIndexStateId")} SELECT {values} " +
+ $"WHERE NOT EXISTS (SELECT 1 FROM {Tbl("SearchIndexStates")} WITH (UPDLOCK, HOLDLOCK) WHERE {Col("IndexName")} = {P}IndexName AND {Col("DepartmentId")} = {P}DepartmentId)";
+
+ DateTime? Timestamp(DateTime? value) => value.HasValue ? DatabaseTimestamp(value.Value) : (DateTime?)null;
+ var id = await ScalarAsync(sql, new
+ {
+ state.IndexName,
+ state.DepartmentId,
+ state.SchemaVersion,
+ state.ProtectedCatalogVersion,
+ state.PolicyEpoch,
+ state.Generation,
+ state.State,
+ state.DocumentCount,
+ LastRebuiltOn = Timestamp(state.LastRebuiltOn),
+ LastIndexedModifiedOn = Timestamp(state.LastIndexedModifiedOn),
+ RebuildRequestedOn = Timestamp(state.RebuildRequestedOn),
+ CreatedOn = DatabaseTimestamp(state.CreatedOn),
+ ModifiedOn = DatabaseTimestamp(state.ModifiedOn)
+ }, cancellationToken);
+ if (id == null)
+ return false;
+
+ state.SearchIndexStateId = id.Value;
+ return true;
+ }
}
/// The single-writer publish lease (plan R7 writer sequence step 2). One row per index name, compare-and-set.
diff --git a/Tests/Resgrid.Tests/Bootstrapper.cs b/Tests/Resgrid.Tests/Bootstrapper.cs
index 95507b7c0..5e9ca24e7 100644
--- a/Tests/Resgrid.Tests/Bootstrapper.cs
+++ b/Tests/Resgrid.Tests/Bootstrapper.cs
@@ -64,7 +64,7 @@ public static void Initialize()
.As()
.InstancePerLifetimeScope();
- // IncidentCommandService resolves chat services lazily through the ServiceLocator for
+ // IncidentCommandService resolves chat services lazily (Lazy constructor parameters) for
// its best-effort lane channel hooks. The real ChatChannelService can't activate in
// this container (its repository graph isn't registered), which logged an activation
// error on every lane save/delete test. Loose mocks turn the hooks into no-ops:
diff --git a/Tests/Resgrid.Tests/Providers/NotificationProviderUnconfiguredHubTests.cs b/Tests/Resgrid.Tests/Providers/NotificationProviderUnconfiguredHubTests.cs
new file mode 100644
index 000000000..855e4ca01
--- /dev/null
+++ b/Tests/Resgrid.Tests/Providers/NotificationProviderUnconfiguredHubTests.cs
@@ -0,0 +1,85 @@
+using System;
+using System.Threading.Tasks;
+using FluentAssertions;
+using NUnit.Framework;
+using Resgrid.Config;
+using Resgrid.Model;
+using Resgrid.Providers.Bus;
+
+namespace Resgrid.Tests.Providers
+{
+ ///
+ /// Deployments without the legacy Azure hubs used to throw ArgumentNullException('connectionString')
+ /// out of every register/unregister call, which in the unit worker skipped the Novu registration.
+ ///
+ [TestFixture]
+ public class NotificationProviderUnconfiguredHubTests
+ {
+ private (string, string) _savedConfig;
+
+ [SetUp]
+ public void SetUp()
+ {
+ _savedConfig = (ServiceBusConfig.AzureNotificationHub_FullConnectionString, ServiceBusConfig.AzureUnitNotificationHub_FullConnectionString);
+ }
+
+ [TearDown]
+ public void TearDown()
+ {
+ (ServiceBusConfig.AzureNotificationHub_FullConnectionString, ServiceBusConfig.AzureUnitNotificationHub_FullConnectionString) = _savedConfig;
+ }
+
+ [TestCase(null)]
+ [TestCase("")]
+ [TestCase(" ")]
+ public async Task UnitProvider_register_and_unregister_should_no_op_when_hub_is_unconfigured(string connectionString)
+ {
+ ServiceBusConfig.AzureUnitNotificationHub_FullConnectionString = connectionString;
+ var provider = new UnitNotificationProvider();
+ var pushUri = CreatePushUri();
+
+ Func act = async () =>
+ {
+ await provider.UnRegisterPush(pushUri);
+ await provider.RegisterPush(pushUri);
+ await provider.UnRegisterPushByUserDeviceId(pushUri);
+ await provider.UnRegisterPushByUUID(pushUri.Uuid);
+ };
+
+ await act.Should().NotThrowAsync();
+ }
+
+ [TestCase(null)]
+ [TestCase("")]
+ [TestCase(" ")]
+ public async Task UserProvider_register_and_unregister_should_no_op_when_hub_is_unconfigured(string connectionString)
+ {
+ ServiceBusConfig.AzureNotificationHub_FullConnectionString = connectionString;
+ var provider = new NotificationProvider();
+ var pushUri = CreatePushUri();
+
+ Func act = async () =>
+ {
+ await provider.UnRegisterPush(pushUri);
+ await provider.RegisterPush(pushUri);
+ await provider.UnRegisterPushByUserDeviceId(pushUri);
+ };
+
+ await act.Should().NotThrowAsync();
+ }
+
+ private static PushUri CreatePushUri()
+ {
+ return new PushUri
+ {
+ UserId = "D4813806-7923-4948-A219-439D6FDCE86A",
+ UnitId = 9,
+ DepartmentId = 7,
+ PlatformType = (int)Platforms.Android,
+ PushLocation = "DEPT",
+ DeviceId = "device-token",
+ Uuid = "device-uuid"
+ };
+ }
+ }
+}
diff --git a/Tests/Resgrid.Tests/Rms/LogsDeepLinkTests.cs b/Tests/Resgrid.Tests/Rms/LogsDeepLinkTests.cs
index 336cf23ac..e1f92c242 100644
--- a/Tests/Resgrid.Tests/Rms/LogsDeepLinkTests.cs
+++ b/Tests/Resgrid.Tests/Rms/LogsDeepLinkTests.cs
@@ -154,6 +154,14 @@ public async Task The_list_json_and_training_chart_still_resolve()
(await _logs.TrainingPerMonth()).Should().BeOfType();
}
+ [Test]
+ public async Task The_list_json_offers_view_but_no_delete_after_activation()
+ {
+ var rows = (await _logs.GetLogsList("2025")).Should().BeOfType().Which.Value.Should().BeAssignableTo>().Which.ToList();
+ rows.Should().ContainSingle(r => r.LogId == 41);
+ rows.Should().OnlyContain(r => !r.CanDelete, "the list must not offer a Delete button that DeleteWorkLog refuses");
+ }
+
[Test]
public async Task Another_departments_log_never_renders_through_an_old_link()
{
@@ -210,6 +218,7 @@ public async Task Before_activation_the_same_links_still_write()
await _logs.DeleteWorkLog(41, CancellationToken.None);
_workLogs.Verify(w => w.DeleteLogAsync(41, It.IsAny()), Times.Once, "the guard is the cutover, not the controller");
((ViewLogsView)((ViewResult)await _logs.View(41)).Model).CanDelete.Should().BeTrue("before activation an administrator may still delete");
+ ((IEnumerable)((JsonResult)await _logs.GetLogsList("2025")).Value).Single(r => r.LogId == 41).CanDelete.Should().BeTrue("before activation the list still offers delete");
}
#endregion
diff --git a/Tests/Resgrid.Tests/Rms/RecordsAnalyticsServiceTests.cs b/Tests/Resgrid.Tests/Rms/RecordsAnalyticsServiceTests.cs
index 5aa4230de..47a3e39d6 100644
--- a/Tests/Resgrid.Tests/Rms/RecordsAnalyticsServiceTests.cs
+++ b/Tests/Resgrid.Tests/Rms/RecordsAnalyticsServiceTests.cs
@@ -239,6 +239,16 @@ public async Task Group_scoped_viewer_only_counts_records_they_could_open_and_pa
captured.States.Should().Contain((int)RmsRecordState.Accepted);
}
+ [Test]
+ public async Task A_window_bound_sql_cannot_store_is_refused_instead_of_overflowing_the_query()
+ {
+ // An API caller's two-digit year arrives as 0026; SQL datetime starts at 1753 (RESGRID-WEB-1MY).
+ Func start = () => _svc.GetWorkloadAsync(Dept, Member, new RecordsAnalyticsQuery { Start = new DateTime(26, 1, 1), End = End });
+ await start.Should().ThrowAsync().WithMessage("The start date is not valid.");
+ Func end = () => _svc.GetResponsePerformanceAsync(Dept, Member, new RecordsAnalyticsQuery { End = DateTime.MaxValue });
+ await end.Should().ThrowAsync().WithMessage("The end date is not valid.");
+ }
+
[Test]
public async Task Gate_flag_viewer_window_clamp_and_row_cap_apply_to_every_dashboard()
{
diff --git a/Tests/Resgrid.Tests/Rms/RecordsPreventionStorableDateTests.cs b/Tests/Resgrid.Tests/Rms/RecordsPreventionStorableDateTests.cs
new file mode 100644
index 000000000..dc39b9e4a
--- /dev/null
+++ b/Tests/Resgrid.Tests/Rms/RecordsPreventionStorableDateTests.cs
@@ -0,0 +1,171 @@
+using System;
+using System.Data.SqlTypes;
+using System.Security.Claims;
+using System.Threading.Tasks;
+using FluentAssertions;
+using Microsoft.AspNetCore.Http;
+using Microsoft.AspNetCore.Mvc;
+using NUnit.Framework;
+using Resgrid.Model;
+using Resgrid.Model.Repositories;
+using Resgrid.Services.Records;
+using Resgrid.Web.Services.Controllers.v4;
+using Resgrid.Web.Services.Models.v4.Records;
+using Resgrid.Web.ServicesCore.Helpers;
+using static Resgrid.Tests.Rms.RmsPreventionHarness;
+
+namespace Resgrid.Tests.Rms
+{
+ ///
+ /// A client date outside SQL datetime (a two-digit year sent as 0026) failed the insert with SqlDateTime overflow and a 500
+ /// (Sentry RESGRID-WEB-1MY). The services now refuse it with an ArgumentException, which the web and v4 map to a message / 400,
+ /// before anything is written or a record number is drawn.
+ ///
+ [TestFixture]
+ public class RecordsPreventionStorableDateTests
+ {
+ private static readonly DateTime TwoDigitYear = new DateTime(26, 9, 25, 19, 0, 0, DateTimeKind.Utc);
+ private RmsPreventionHarness _h;
+
+ [SetUp]
+ public void SetUp() => _h = new RmsPreventionHarness();
+
+ [Test]
+ public void The_gate_accepts_the_sql_datetime_range_and_refuses_dates_outside_it()
+ {
+ RecordsPreventionGate.RequireStorableDate((DateTime)SqlDateTime.MinValue, "bad").Should().Be((DateTime)SqlDateTime.MinValue);
+ RecordsPreventionGate.RequireStorableDate((DateTime)SqlDateTime.MaxValue, "bad").Should().Be((DateTime)SqlDateTime.MaxValue);
+ RecordsPreventionGate.RequireStorableDate((DateTime?)null, "bad").Should().BeNull();
+
+ foreach (var value in new[] { TwoDigitYear, ((DateTime)SqlDateTime.MinValue).AddTicks(-1), DateTime.MinValue, DateTime.MaxValue })
+ FluentActions.Invoking(() => RecordsPreventionGate.RequireStorableDate(value, "bad")).Should().Throw().WithMessage("bad");
+ }
+
+ [Test]
+ public async Task Investigation_note_and_evidence_refuse_the_date_without_writing_or_drawing_an_evidence_number()
+ {
+ var investigation = await _h.InvestigationsService.OpenAsync(Dept, Admin, "Warehouse fire", null, 42, "Origin unknown.");
+
+ Func note = () => _h.InvestigationsService.AddNoteAsync(Dept, Admin, investigation.RmsInvestigationCaseId, RmsInvestigationNoteKind.Interview, TwoDigitYear, "Interview", "Body");
+ await note.Should().ThrowAsync().WithMessage("The note date is not valid.");
+ Func evidence = () => _h.InvestigationsService.AddEvidenceAsync(Dept, Admin, investigation.RmsInvestigationCaseId, new RmsInvestigationEvidence { Description = "Photo", CollectedOn = TwoDigitYear });
+ await evidence.Should().ThrowAsync().WithMessage("The collection date is not valid.");
+ _h.CaseNotes.Rows.Should().BeEmpty();
+ _h.Evidence.Rows.Should().BeEmpty();
+ _h.Custody.Rows.Should().BeEmpty();
+
+ // Unset still means now, and the rejected call did not use up evidence number 1.
+ var saved = await _h.InvestigationsService.AddNoteAsync(Dept, Admin, investigation.RmsInvestigationCaseId, RmsInvestigationNoteKind.Interview, default, "Interview", "Body");
+ saved.OccurredOn.Should().BeCloseTo(DateTime.UtcNow, TimeSpan.FromMinutes(1));
+ var item = await _h.InvestigationsService.AddEvidenceAsync(Dept, Admin, investigation.RmsInvestigationCaseId, new RmsInvestigationEvidence { Description = "Photo" });
+ item.EvidenceNumber.Should().EndWith("-0001");
+ }
+
+ [Test, NonParallelizable]
+ public async Task The_v4_note_endpoint_answers_400_with_the_message_instead_of_a_500()
+ {
+ var investigation = await _h.InvestigationsService.OpenAsync(Dept, Admin, "Warehouse fire", null, 42, "Origin unknown.");
+ var previous = ClaimsAuthorizationHelper._httpContextAccessor;
+ var http = new DefaultHttpContext { User = new ClaimsPrincipal(new ClaimsIdentity(new[] {
+ new Claim(ClaimTypes.PrimarySid, Admin), new Claim(ClaimTypes.PrimaryGroupSid, Dept.ToString()) }, "Test")) };
+ ClaimsAuthorizationHelper._httpContextAccessor = new HttpContextAccessor { HttpContext = http };
+ try
+ {
+ var controller = new RecordInvestigationsController(_h.InvestigationsService, _h.Cutover.Object) { ControllerContext = new ControllerContext { HttpContext = http } };
+
+ var result = await controller.SaveNote(new CaseNoteInput { CaseId = investigation.RmsInvestigationCaseId, Kind = (int)RmsInvestigationNoteKind.Interview, OccurredOn = TwoDigitYear, Subject = "Interview", Body = "Body" }, default);
+
+ var problem = result.Result.Should().BeOfType().Subject;
+ problem.StatusCode.Should().Be(StatusCodes.Status400BadRequest);
+ problem.Value.Should().BeOfType().Which.Title.Should().Be("The note date is not valid.");
+ _h.CaseNotes.Rows.Should().BeEmpty();
+ }
+ finally { ClaimsAuthorizationHelper._httpContextAccessor = previous; }
+ }
+
+ [Test]
+ public async Task Hydrant_flow_test_and_maintenance_refuse_the_date_and_leave_the_hydrant_unchanged()
+ {
+ var hydrant = await _h.HydrantsService.SaveAsync(Dept, Admin, new RmsHydrant { HydrantNumber = "H-101", Latitude = 45.5m, Longitude = -122.6m, Type = (int)RmsHydrantType.DryBarrel, MainSizeInches = 8 });
+
+ Func test = () => _h.HydrantsService.RecordFlowTestAsync(Dept, Admin, new RmsHydrantFlowTest { RmsHydrantId = hydrant.RmsHydrantId, TestedOn = TwoDigitYear, PitotPressurePsi = 64, OutletDiameterInches = 2.5m });
+ await test.Should().ThrowAsync().WithMessage("The test date is not valid.");
+ Func maintenance = () => _h.HydrantsService.RecordMaintenanceAsync(Dept, Admin, new RmsHydrantMaintenance { RmsHydrantId = hydrant.RmsHydrantId, PerformedOn = TwoDigitYear });
+ await maintenance.Should().ThrowAsync().WithMessage("The maintenance date is not valid.");
+
+ _h.FlowTests.Rows.Should().BeEmpty();
+ _h.Maintenance.Rows.Should().BeEmpty();
+ hydrant.LastTestedOn.Should().BeNull();
+ hydrant.LastMaintainedOn.Should().BeNull();
+ }
+
+ [Test]
+ public async Task Inspection_schedule_violation_and_list_refuse_the_date_without_drawing_an_inspection_number()
+ {
+ var occupancy = _h.SeedOccupancy();
+
+ Func schedule = () => _h.InspectionsService.ScheduleAsync(Dept, Admin, occupancy.RmsOccupancyId, null, TwoDigitYear, Admin);
+ await schedule.Should().ThrowAsync().WithMessage("The scheduled date is not valid.");
+ _h.Inspections.Rows.Should().BeEmpty();
+ Func list = () => _h.InspectionsService.ListAsync(Dept, Admin, new RmsInspectionQuery { ScheduledBefore = TwoDigitYear });
+ await list.Should().ThrowAsync().WithMessage("The scheduled-before date is not valid.");
+ Func count = () => _h.InspectionsService.CountAsync(Dept, Admin, new RmsInspectionQuery { ScheduledBefore = TwoDigitYear });
+ await count.Should().ThrowAsync().WithMessage("The scheduled-before date is not valid.");
+
+ var inspection = await _h.InspectionsService.ScheduleAsync(Dept, Admin, occupancy.RmsOccupancyId, null, DateTime.UtcNow, Admin);
+ inspection.InspectionNumber.Should().EndWith("-0001");
+ Func violation = () => _h.InspectionsService.SaveViolationAsync(Dept, Admin, new RmsViolation { RmsInspectionId = inspection.RmsInspectionId, Description = "Blocked exit", DueOn = TwoDigitYear });
+ await violation.Should().ThrowAsync().WithMessage("The due date is not valid.");
+ _h.Violations.Rows.Should().BeEmpty();
+ }
+
+ [Test]
+ public async Task Permit_update_issue_and_list_refuse_the_date()
+ {
+ var type = await _h.PermitsService.SaveTypeAsync(Dept, Admin, new RmsPermitType { Name = "Hot work", Code = "HW", DefaultValidityDays = 30, IsActive = true });
+ var permit = await _h.PermitsService.ApplyAsync(Dept, Admin, new RmsPermit { RmsPermitTypeId = type.RmsPermitTypeId, ApplicantName = "Welder" });
+
+ Func update = () => _h.PermitsService.UpdateAsync(Dept, Admin, new RmsPermit { RmsPermitId = permit.RmsPermitId, ApplicantName = "Welder", ExpiresOn = TwoDigitYear });
+ await update.Should().ThrowAsync().WithMessage("The expiry date is not valid.");
+ await _h.PermitsService.TransitionAsync(Dept, Admin, permit.RmsPermitId, RmsPermitState.Approved, null, null, null);
+ Func issue = () => _h.PermitsService.TransitionAsync(Dept, Admin, permit.RmsPermitId, RmsPermitState.Issued, null, TwoDigitYear, null);
+ await issue.Should().ThrowAsync().WithMessage("The effective date is not valid.");
+ Func list = () => _h.PermitsService.ListAsync(Dept, Admin, new RmsPermitQuery { ExpiresBefore = TwoDigitYear });
+ await list.Should().ThrowAsync().WithMessage("The expires-before date is not valid.");
+
+ var stored = await _h.PermitsService.GetAsync(Dept, Admin, permit.RmsPermitId);
+ stored.Permit.State.Should().Be((int)RmsPermitState.Approved, "the refused issue did not move the permit");
+ }
+
+ [Test]
+ public async Task Crr_activity_and_its_window_refuse_the_date()
+ {
+ Func save = () => _h.CrrService.SaveAsync(Dept, Admin, new RmsCrrActivity { Title = "School visit", OccurredOn = TwoDigitYear });
+ await save.Should().ThrowAsync().WithMessage("The activity date is not valid.");
+ _h.Crr.Rows.Should().BeEmpty();
+
+ Func list = () => _h.CrrService.ListAsync(Dept, Admin, TwoDigitYear, DateTime.UtcNow, 10);
+ await list.Should().ThrowAsync().WithMessage("The start date is not valid.");
+ Func summary = () => _h.CrrService.GetSummaryAsync(Dept, Admin, DateTime.UtcNow.AddDays(-30), DateTime.MaxValue);
+ await summary.Should().ThrowAsync().WithMessage("The end date is not valid.");
+ }
+
+ [Test]
+ public async Task Occupancy_next_review_date_is_refused_before_an_occupancy_number_is_drawn()
+ {
+ Func save = () => _h.OccupancyService.SaveAsync(Dept, Admin, new RmsOccupancy { Name = "Riverside Mill", NextReviewDue = TwoDigitYear });
+ await save.Should().ThrowAsync().WithMessage("The next review date is not valid.");
+ _h.Occupancies.Rows.Should().BeEmpty();
+ _h.Sequences.Rows.Should().BeEmpty();
+ }
+
+ [Test]
+ public async Task Quality_sampling_and_trends_refuse_the_date_before_looking_up_the_rubric()
+ {
+ Func sample = () => _h.QualityService.SampleAsync(Dept, Admin, "no-such-rubric", TwoDigitYear);
+ await sample.Should().ThrowAsync().WithMessage("The sample start date is not valid.");
+ Func trends = () => _h.QualityService.GetTrendsAsync(Dept, Admin, TwoDigitYear);
+ await trends.Should().ThrowAsync().WithMessage("The start date is not valid.");
+ }
+ }
+}
diff --git a/Tests/Resgrid.Tests/RootScopeResolutionTests.cs b/Tests/Resgrid.Tests/RootScopeResolutionTests.cs
new file mode 100644
index 000000000..57efc7b28
--- /dev/null
+++ b/Tests/Resgrid.Tests/RootScopeResolutionTests.cs
@@ -0,0 +1,93 @@
+using System.Collections.Generic;
+using System.IO;
+using System.Linq;
+using System.Text.RegularExpressions;
+using NUnit.Framework;
+
+namespace Resgrid.Tests
+{
+ ///
+ /// IUnitOfWork is InstancePerLifetimeScope, so anything resolved from the root container shares one unit of work,
+ /// and one DB connection, with everything else resolved there for the life of the process. When one of those
+ /// services opens a transaction, concurrent work lands on its connection ("The connection does not support
+ /// MultipleActiveResultSets") and inside its transaction. Bootstrapper.GetKernel().Resolve and
+ /// ServiceLocator.Current both resolve from the root.
+ ///
+ [TestFixture]
+ public sealed class RootScopeResolutionTests
+ {
+ // Process-lifetime singletons the worker hosts prime at startup so their event subscriptions exist.
+ private static readonly string[] AllowedWorkerRootResolutions =
+ {
+ "IEventAggregator", "IWorkflowEventProvider", "IOutboundEventProvider", "ICoreEventService"
+ };
+
+ ///
+ /// Workers resolve from a per-run child scope instead:
+ /// using var scope = Bootstrapper.GetKernel().BeginLifetimeScope(); scope.Resolve<T>().
+ ///
+ [Test]
+ public void Worker_source_does_not_resolve_scoped_services_from_the_root_container()
+ {
+ var root = RepositoryRoot();
+ if (root == null)
+ Assert.Ignore("Resgrid.sln not found above the test directory; the worker source is not available to scan.");
+
+ var rootResolve = new Regex(@"GetKernel\(\)\s*\.\s*Resolve\s*[<(]");
+ var allowed = new Regex(@"GetKernel\(\)\s*\.\s*Resolve<(" + string.Join("|", AllowedWorkerRootResolutions) + @")>\(\)");
+
+ var offenders = Offenders(root, new[] { "Workers" }, line => rootResolve.IsMatch(line) && !allowed.IsMatch(line));
+
+ Assert.That(offenders, Is.Empty,
+ "Resolve these from a per-run scope (Bootstrapper.GetKernel().BeginLifetimeScope()) instead of the root container");
+ }
+
+ ///
+ /// Scoped services take their dependencies through the constructor (Lazy<T> where the graph would otherwise
+ /// cycle); singletons inject ILifetimeScope and begin a child scope per operation.
+ ///
+ [Test]
+ public void Core_source_does_not_resolve_through_the_root_service_locator()
+ {
+ var root = RepositoryRoot();
+ if (root == null)
+ Assert.Ignore("Resgrid.sln not found above the test directory; the source is not available to scan.");
+
+ // Only ever resolves SqlConfiguration, the query classes' sole constructor parameter, which holds no
+ // connection or unit of work.
+ var allowedFiles = new[] { Path.Combine(root, "Repositories", "Resgrid.Repositories.DataRepository", "Queries", "QueryList.cs") };
+ var locator = new Regex(@"ServiceLocator\s*\.\s*Current\s*\.\s*GetInstance\b|GetKernel\(\)\s*\.\s*Resolve\s*[<(]");
+
+ var offenders = Offenders(root, new[] { "Core", "Providers", "Repositories" }, locator.IsMatch, allowedFiles);
+
+ Assert.That(offenders, Is.Empty,
+ "Inject these through the constructor (Lazy to break a cycle; ILifetimeScope with a child scope per operation in a singleton)");
+ }
+
+ private static List Offenders(string root, IEnumerable directories, System.Func isOffending, ICollection allowedFiles = null)
+ {
+ return directories
+ .SelectMany(directory => Directory.EnumerateFiles(Path.Combine(root, directory), "*.cs", SearchOption.AllDirectories))
+ .Where(file => !IsBuildOutput(file) && (allowedFiles == null || !allowedFiles.Contains(file)))
+ .SelectMany(file => File.ReadLines(file)
+ .Select((line, index) => (line, index))
+ .Where(x => !x.line.TrimStart().StartsWith("//") && isOffending(x.line))
+ .Select(x => $"{Path.GetRelativePath(root, file)}:{x.index + 1}: {x.line.Trim()}"))
+ .ToList();
+ }
+
+ private static bool IsBuildOutput(string path)
+ {
+ var separator = Path.DirectorySeparatorChar;
+ return path.Contains($"{separator}obj{separator}") || path.Contains($"{separator}bin{separator}");
+ }
+
+ private static string RepositoryRoot()
+ {
+ var directory = new DirectoryInfo(TestContext.CurrentContext.TestDirectory);
+ while (directory != null && !File.Exists(Path.Combine(directory.FullName, "Resgrid.sln")))
+ directory = directory.Parent;
+ return directory?.FullName;
+ }
+ }
+}
diff --git a/Tests/Resgrid.Tests/Search/SystemActionsServiceTests.cs b/Tests/Resgrid.Tests/Search/SystemActionsServiceTests.cs
index 40b94137d..3c764d177 100644
--- a/Tests/Resgrid.Tests/Search/SystemActionsServiceTests.cs
+++ b/Tests/Resgrid.Tests/Search/SystemActionsServiceTests.cs
@@ -16,10 +16,12 @@ namespace Resgrid.Tests.Search
public class SystemActionsServiceTests
{
private Mock _flags;
+ private Mock _cutover;
private SystemActionsService _service;
private HashSet _enabledFlags;
private HashSet _claims;
private HashSet _disabledModules;
+ private bool _legacyWritesBlocked;
[SetUp]
public void SetUp()
@@ -27,10 +29,13 @@ public void SetUp()
_enabledFlags = new HashSet();
_claims = new HashSet();
_disabledModules = new HashSet();
+ _legacyWritesBlocked = false;
_flags = new Mock();
_flags.Setup(f => f.IsEnabledAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny>()))
.ReturnsAsync((string key, int dept, bool def, IDictionary ctx) => _enabledFlags.Contains(key));
- _service = new SystemActionsService(_flags.Object);
+ _cutover = new Mock();
+ _cutover.Setup(c => c.AreLegacyWritesBlockedAsync(It.IsAny())).ReturnsAsync(() => _legacyWritesBlocked);
+ _service = new SystemActionsService(_flags.Object, _cutover.Object);
}
private SearchPrincipal Principal(bool admin = false) => new SearchPrincipal
@@ -80,13 +85,38 @@ public async Task Feature_flag_and_module_gates_apply()
}
[Test]
- public async Task Logs_disappear_when_records_is_on_and_admin_entries_need_admin()
+ public async Task Logs_stay_findable_through_the_records_cutover_and_only_new_log_goes_away()
{
_claims.Add("Log:View");
+ _claims.Add("Log:Create");
(await _service.SearchAsync("logs", Principal())).Select(h => h.Key).Should().Contain("logs");
+ (await _service.SearchAsync("new log", Principal())).Select(h => h.Key).Should().Contain("new-log");
+
+ // Records.System on but not yet activated: Logs is still the department's working log system.
_enabledFlags.Add(FeatureFlagKeys.RecordsSystem);
- (await _service.SearchAsync("logs", Principal())).Select(h => h.Key).Should().NotContain("logs");
+ (await _service.SearchAsync("logs", Principal())).Select(h => h.Key).Should().Contain("logs");
+ (await _service.SearchAsync("new log", Principal())).Select(h => h.Key).Should().Contain("new-log", "the flag alone does not make Logs read-only");
+
+ // Activated: old Logs stay readable, creating one is refused by the Logs pages, so it is not offered.
+ _legacyWritesBlocked = true;
+ (await _service.SearchAsync("logs", Principal())).Select(h => h.Key).Should().Contain("logs", "old Logs remain readable after activation");
+ (await _service.SearchAsync("legacy logs", Principal())).Select(h => h.Key).Should().Contain("logs");
+ (await _service.SearchAsync("new log", Principal())).Select(h => h.Key).Should().NotContain("new-log");
+ }
+
+ [Test]
+ public async Task An_unreadable_cutover_state_hides_the_legacy_write_but_not_the_read()
+ {
+ _claims.Add("Log:View");
+ _claims.Add("Log:Create");
+ _cutover.Setup(c => c.AreLegacyWritesBlockedAsync(It.IsAny())).ThrowsAsync(new System.InvalidOperationException("cache down"));
+ (await _service.SearchAsync("logs", Principal())).Select(h => h.Key).Should().Contain("logs");
+ (await _service.SearchAsync("new log", Principal())).Select(h => h.Key).Should().NotContain("new-log");
+ }
+ [Test]
+ public async Task Admin_entries_need_admin()
+ {
(await _service.SearchAsync("department settings", Principal())).Should().BeEmpty();
(await _service.SearchAsync("department settings", Principal(admin: true))).First().Key.Should().Be("department-settings");
}
diff --git a/Tests/Resgrid.Tests/Search/UnifiedSearchServiceTests.cs b/Tests/Resgrid.Tests/Search/UnifiedSearchServiceTests.cs
index 03b004e3e..415903ed9 100644
--- a/Tests/Resgrid.Tests/Search/UnifiedSearchServiceTests.cs
+++ b/Tests/Resgrid.Tests/Search/UnifiedSearchServiceTests.cs
@@ -229,8 +229,8 @@ public async Task Index_unavailable_degrades_and_activates_the_department_lazily
_global.SetupGet(g => g.IsAvailable).Returns(false);
_states.Setup(s => s.GetAsync(SearchIndexNames.Global, 7)).ReturnsAsync((SearchIndexState)null);
SearchIndexState saved = null;
- _states.Setup(s => s.SaveOrUpdateAsync(It.IsAny(), It.IsAny(), It.IsAny()))
- .Callback((SearchIndexState s, CancellationToken _, bool __) => saved = s).ReturnsAsync((SearchIndexState s, CancellationToken _, bool __) => s);
+ _states.Setup(s => s.InsertIfMissingAsync(It.IsAny(), It.IsAny()))
+ .Callback((SearchIndexState s, CancellationToken _) => saved = s).ReturnsAsync(true);
var result = await _service.SearchAsync(new UnifiedSearchRequest { Text = "one" }, Principal("Call:View"));
@@ -240,6 +240,21 @@ public async Task Index_unavailable_degrades_and_activates_the_department_lazily
saved.Should().NotBeNull();
saved.State.Should().Be((int)SearchIndexBuildState.RebuildRequested);
saved.IndexName.Should().Be(SearchIndexNames.Global);
+ saved.DepartmentId.Should().Be(7);
+ _states.Verify(s => s.SaveOrUpdateAsync(It.IsAny(), It.IsAny(), It.IsAny()), Times.Never,
+ "an unconditional insert races concurrent first searches into the unique index");
+ }
+
+ [Test]
+ public async Task Index_unavailable_with_a_state_row_already_present_does_not_try_to_create_one()
+ {
+ _global.SetupGet(g => g.IsAvailable).Returns(false);
+ _states.Setup(s => s.GetAsync(SearchIndexNames.Global, 7)).ReturnsAsync(new SearchIndexState { IndexName = SearchIndexNames.Global, DepartmentId = 7 });
+
+ var result = await _service.SearchAsync(new UnifiedSearchRequest { Text = "one" }, Principal("Call:View"));
+
+ result.Degraded.Should().BeTrue();
+ _states.Verify(s => s.InsertIfMissingAsync(It.IsAny(), It.IsAny()), Times.Never);
}
[Test]
public async Task A_records_only_page_past_the_first_twenty_still_returns_records()
diff --git a/Tests/Resgrid.Tests/Services/ChatPermissionServiceTests.cs b/Tests/Resgrid.Tests/Services/ChatPermissionServiceTests.cs
index be1383351..7e801cc97 100644
--- a/Tests/Resgrid.Tests/Services/ChatPermissionServiceTests.cs
+++ b/Tests/Resgrid.Tests/Services/ChatPermissionServiceTests.cs
@@ -131,12 +131,63 @@ protected static ChatChannelMember CreateUserMember(ChatChannel channel, string
public class when_resolving_the_channel_access_version : with_the_chat_permission_service
{
[Test]
- public async Task a_missing_cache_epoch_should_remain_null()
+ public async Task an_unavailable_cache_should_fail_closed()
{
+ _cacheProviderMock.Setup(x => x.GetOrAddStringAsync(It.IsAny(), It.IsAny(), It.IsAny()))
+ .ReturnsAsync((string)null);
+
var version = await _chatPermissionService.GetChannelAccessVersionAsync("channel-1");
version.Should().BeNull();
}
+
+ [Test]
+ public async Task a_channel_without_an_epoch_should_get_one_minted()
+ {
+ // Most channels are never invalidated, so the epoch key is usually absent; treating that as
+ // an outage made JoinChannel throw and dropped every realtime fan-out for the channel.
+ _cacheProviderMock.Setup(x => x.GetOrAddStringAsync(It.IsAny(), It.IsAny(), It.IsAny()))
+ .ReturnsAsync((string _, string valueIfAbsent, TimeSpan _) => valueIfAbsent);
+
+ var version = await _chatPermissionService.GetChannelAccessVersionAsync("channel-1");
+
+ version.Should().NotBeNullOrWhiteSpace();
+ _cacheProviderMock.Verify(x => x.GetOrAddStringAsync("chatpermver:channel-1", It.IsAny(), It.IsAny()), Times.Once);
+ }
+
+ [Test]
+ public async Task an_existing_epoch_should_be_returned_unchanged()
+ {
+ _cacheProviderMock.Setup(x => x.GetOrAddStringAsync("chatpermver:channel-1", It.IsAny(), It.IsAny()))
+ .ReturnsAsync("existing-epoch");
+
+ var version = await _chatPermissionService.GetChannelAccessVersionAsync("channel-1");
+
+ version.Should().Be("existing-epoch");
+ }
+ }
+
+ [TestFixture]
+ public class when_invalidating_a_channel : with_the_chat_permission_service
+ {
+ [Test]
+ public async Task each_invalidation_should_rotate_to_a_never_used_epoch()
+ {
+ var written = new List();
+ _cacheProviderMock.Setup(x => x.SetStringAsync("chatpermver:channel-1", It.IsAny(), It.IsAny()))
+ .Callback((string _, string value, TimeSpan _) => written.Add(value))
+ .ReturnsAsync(true);
+
+ await _chatPermissionService.InvalidateChannelCacheAsync("channel-1");
+ await _chatPermissionService.InvalidateChannelCacheAsync("channel-1");
+
+ written.Should().HaveCount(2);
+ written.Should().OnlyContain(v => !string.IsNullOrWhiteSpace(v));
+ written.Should().OnlyHaveUniqueItems();
+ // A counter restarts at 1 once its key lapses, reviving an obsolete group a revoked
+ // connection may still sit in; epochs must never repeat.
+ _cacheProviderMock.Verify(x => x.IncrementAsync("chatpermver:channel-1", It.IsAny()), Times.Never);
+ }
}
[TestFixture]
diff --git a/Tests/Resgrid.Tests/Services/ChecklistEventDeliveryTests.cs b/Tests/Resgrid.Tests/Services/ChecklistEventDeliveryTests.cs
index 8c9434452..2859c6d23 100644
--- a/Tests/Resgrid.Tests/Services/ChecklistEventDeliveryTests.cs
+++ b/Tests/Resgrid.Tests/Services/ChecklistEventDeliveryTests.cs
@@ -118,7 +118,7 @@ public async Task Queue_rejection_after_run_insert_retries_the_same_workflow_run
runs.Setup(s => s.InsertAsync(It.IsAny(), It.IsAny(), It.IsAny())).ReturnsAsync((WorkflowRun r, CancellationToken ct, bool first) => stored = r);
var attempts = new List();
queue.Setup(s => s.EnqueueWorkflow(It.IsAny())).ReturnsAsync((WorkflowQueueItem item) => { attempts.Add(item); return attempts.Count > 1; });
- _ = new WorkflowEventProvider(_bus, queue.Object, workflows.Object, runs.Object, departments.Object, subscriptions.Object, _projection, _history.Lazy);
+ _ = new WorkflowEventProvider(_bus, queue.Object, WorkflowEventProviderScopeTests.Services(workflows.Object, runs.Object, departments.Object, subscriptions.Object, _projection, _history.Lazy));
var envelope = Event(); envelope.Trigger = (WorkflowTriggerEventType)trigger; envelope.EventName = envelope.Trigger.ToString();
var entry = await _outbox.EnqueueAsync(42, "Checklists", envelope);
(await _outbox.DispatchAfterCommitAsync(new[] { entry.DomainEventOutboxId })).Should().Be(0); stored.Should().NotBeNull();
diff --git a/Tests/Resgrid.Tests/Services/CoreEventServiceTests.cs b/Tests/Resgrid.Tests/Services/CoreEventServiceTests.cs
index 8ceb53dc3..19f311c12 100644
--- a/Tests/Resgrid.Tests/Services/CoreEventServiceTests.cs
+++ b/Tests/Resgrid.Tests/Services/CoreEventServiceTests.cs
@@ -1,8 +1,16 @@
+using System;
+using System.Collections.Generic;
+using System.Threading;
using System.Threading.Tasks;
+using Autofac;
+using FluentAssertions;
using Moq;
using NUnit.Framework;
+using Resgrid.Model;
using Resgrid.Model.Events;
using Resgrid.Model.Providers;
+using Resgrid.Model.Services;
+using Resgrid.Providers.Bus;
using Resgrid.Services;
namespace Resgrid.Tests.Services
@@ -19,7 +27,7 @@ public class CoreEventServiceTests
public async Task IncidentCommandUpdatedAsync_RaisesIncidentCommandUpdatedEvent_WithDeptAndCall()
{
var eventAggregator = new Mock();
- var service = new CoreEventService(eventAggregator.Object);
+ var service = new CoreEventService(eventAggregator.Object, Mock.Of());
await service.IncidentCommandUpdatedAsync(42, 1001);
@@ -27,5 +35,52 @@ public async Task IncidentCommandUpdatedAsync_RaisesIncidentCommandUpdatedEvent_
It.Is(e => e.DepartmentId == 42 && e.CallId == 1001)),
Times.Once);
}
+
+ ///
+ /// The service is a singleton. The timestamp save runs inside an audited configuration transaction, so each
+ /// event must get its own settings service (and unit of work) from a child scope, never one shared root instance.
+ ///
+ [Test]
+ public void DepartmentSettingsUpdateEvent_SavesTheTimestamp_InItsOwnScopePerEvent()
+ {
+ var created = new List>();
+ var builder = new ContainerBuilder();
+ builder.Register(_ =>
+ {
+ var settings = new Mock();
+ settings.Setup(s => s.SaveOrUpdateSettingAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny()))
+ .ReturnsAsync(new DepartmentSetting());
+ created.Add(settings);
+ return settings.Object;
+ }).As().InstancePerLifetimeScope();
+ using var container = builder.Build();
+
+ var bus = new EventAggregator();
+ _ = new CoreEventService(bus, container);
+
+ bus.SendMessage(new DepartmentSettingsUpdateEvent { DepartmentId = 42 });
+ bus.SendMessage(new DepartmentSettingsUpdateEvent { DepartmentId = 43 });
+
+ created.Should().HaveCount(2);
+ created[0].Verify(s => s.SaveOrUpdateSettingAsync(42, It.IsAny(), DepartmentSettingTypes.UpdateTimestamp, It.IsAny()), Times.Once);
+ created[1].Verify(s => s.SaveOrUpdateSettingAsync(43, It.IsAny(), DepartmentSettingTypes.UpdateTimestamp, It.IsAny()), Times.Once);
+ }
+
+ [Test]
+ public void DepartmentSettingsUpdateEvent_FailedSave_DoesNotReachThePublisher()
+ {
+ var settings = new Mock();
+ settings.Setup(s => s.SaveOrUpdateSettingAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny()))
+ .ThrowsAsync(new InvalidOperationException("settings store unavailable"));
+ var builder = new ContainerBuilder();
+ builder.RegisterInstance(settings.Object).As();
+ using var container = builder.Build();
+
+ var bus = new EventAggregator();
+ _ = new CoreEventService(bus, container);
+
+ bus.Invoking(b => b.SendMessage(new DepartmentSettingsUpdateEvent { DepartmentId = 42 })).Should().NotThrow();
+ settings.Verify(s => s.SaveOrUpdateSettingAsync(42, It.IsAny(), DepartmentSettingTypes.UpdateTimestamp, It.IsAny()), Times.Once);
+ }
}
}
diff --git a/Tests/Resgrid.Tests/Services/InventoryWorkflowTests.cs b/Tests/Resgrid.Tests/Services/InventoryWorkflowTests.cs
index ff5dea56c..cfd754810 100644
--- a/Tests/Resgrid.Tests/Services/InventoryWorkflowTests.cs
+++ b/Tests/Resgrid.Tests/Services/InventoryWorkflowTests.cs
@@ -259,7 +259,7 @@ public async Task Queue_rejection_retries_existing_inventory_workflow_run_with_s
runs.Setup(s => s.GetByWorkflowsAndEventAsync(42, It.IsAny>(), It.IsAny())).ReturnsAsync(() => stored == null ? new List() : new List { stored });
runs.Setup(s => s.InsertAsync(It.IsAny(), It.IsAny(), It.IsAny())).ReturnsAsync((WorkflowRun run, CancellationToken c, bool f) => stored = run);
var attempts = new List(); queue.Setup(s => s.EnqueueWorkflow(It.IsAny())).ReturnsAsync((WorkflowQueueItem q) => { attempts.Add(q); return attempts.Count > 1; });
- _ = new WorkflowEventProvider(_bus, queue.Object, workflows.Object, runs.Object, Mock.Of(), subscriptions.Object, _projection, _history.Lazy);
+ _ = new WorkflowEventProvider(_bus, queue.Object, WorkflowEventProviderScopeTests.Services(workflows.Object, runs.Object, Mock.Of(), subscriptions.Object, _projection, _history.Lazy));
var entry = await _outbox.EnqueueAsync(42, "Inventory", Event(trigger));
(await _outbox.DispatchAfterCommitAsync(new[] { entry.DomainEventOutboxId })).Should().Be(0); stored.Should().NotBeNull(); entry.DispatchedOn.Should().BeNull();
_policy.Setup(s => s.IsProtectionEnforcedAsync(42)).ReturnsAsync(true);
diff --git a/Tests/Resgrid.Tests/Services/LazyChatDependencyCompositionTests.cs b/Tests/Resgrid.Tests/Services/LazyChatDependencyCompositionTests.cs
new file mode 100644
index 000000000..f41f7dceb
--- /dev/null
+++ b/Tests/Resgrid.Tests/Services/LazyChatDependencyCompositionTests.cs
@@ -0,0 +1,116 @@
+using System;
+using System.Collections.Generic;
+using System.Linq;
+using System.Reflection;
+using Autofac;
+using Autofac.Builder;
+using Autofac.Core;
+using FluentAssertions;
+using Moq;
+using NUnit.Framework;
+using Resgrid.Model.Repositories;
+using Resgrid.Model.Services;
+using Resgrid.Services;
+
+namespace Resgrid.Tests.Services
+{
+ ///
+ /// IncidentCommandService and PermissionsService used to reach the chat and command-access services through
+ /// ServiceLocator.Current, which resolves from the ROOT scope and so shared the root unit of work with everything
+ /// else resolved there. They now take Lazy<T> constructor parameters, because the dependencies run both ways:
+ /// ChatChannelService and ChatPermissionService take IIncidentCommandService, and CommandAccessService takes
+ /// IPermissionsService. These tests compose the real classes, with every other dependency a loose mock, to prove
+ /// the cycles still resolve and that each Lazy resolves in the caller's own scope.
+ ///
+ [TestFixture]
+ public sealed class LazyChatDependencyCompositionTests
+ {
+ private IContainer _container;
+
+ [SetUp]
+ public void SetUp()
+ {
+ var builder = new ContainerBuilder();
+ 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.RegisterSource(new LooseMockSource());
+ _container = builder.Build();
+ }
+
+ [TearDown]
+ public void TearDown() => _container.Dispose();
+
+ [Test]
+ public void Incident_command_chat_dependencies_resolve_through_the_cycle_in_the_callers_scope()
+ {
+ using var first = _container.BeginLifetimeScope();
+ using var second = _container.BeginLifetimeScope();
+
+ // Entering from the chat side constructs IncidentCommandService inside ChatChannelService's resolve.
+ var chat = first.Resolve();
+ var incident = first.Resolve();
+
+ LazyValue(incident, "_chatChannelService").Should().BeSameAs(chat);
+ LazyValue(incident, "_commandAccessService").Should().BeSameAs(first.Resolve());
+ LazyValue(incident, "_chatChannelRepository").Should().BeSameAs(first.Resolve());
+
+ LazyValue(second.Resolve(), "_chatChannelService")
+ .Should().BeSameAs(second.Resolve()).And.NotBeSameAs(chat);
+ }
+
+ [Test]
+ public void Permissions_chat_refresh_dependencies_resolve_through_the_cycle_in_the_callers_scope()
+ {
+ using var scope = _container.BeginLifetimeScope();
+
+ var permissions = scope.Resolve();
+
+ LazyValue(permissions, "_chatPermissionService").Should().BeSameAs(scope.Resolve());
+ LazyValue(permissions, "_chatChannelRepository").Should().BeSameAs(scope.Resolve());
+ }
+
+ [Test]
+ public void Chat_message_service_can_reach_the_root_for_its_background_push_fan_out()
+ {
+ using var scope = _container.BeginLifetimeScope();
+
+ var lifetimeScope = typeof(ChatMessageService).GetField("_lifetimeScope", BindingFlags.Instance | BindingFlags.NonPublic)
+ .GetValue(scope.Resolve());
+
+ lifetimeScope.Should().BeSameAs(scope);
+ ((ISharingLifetimeScope)lifetimeScope).RootLifetimeScope.Should().BeSameAs(((ISharingLifetimeScope)scope).RootLifetimeScope);
+ }
+
+ private static T LazyValue(object service, string field) =>
+ ((Lazy)service.GetType().GetField(field, BindingFlags.Instance | BindingFlags.NonPublic).GetValue(service)).Value;
+
+ /// A loose Moq mock, one per scope, for any interface nothing else registers.
+ private sealed class LooseMockSource : IRegistrationSource
+ {
+ private static readonly Type[] CollectionTypes =
+ { typeof(IEnumerable<>), typeof(ICollection<>), typeof(IList<>), typeof(IReadOnlyCollection<>), typeof(IReadOnlyList<>) };
+
+ public bool IsAdapterForIndividualComponents => false;
+
+ public IEnumerable RegistrationsFor(Service service, Func> registrationAccessor)
+ {
+ if (service is not IServiceWithType typed || !typed.ServiceType.IsInterface || registrationAccessor(service).Any())
+ return Enumerable.Empty();
+
+ var type = typed.ServiceType;
+ if (type.IsGenericType && CollectionTypes.Contains(type.GetGenericTypeDefinition()))
+ return Enumerable.Empty();
+
+ return new[]
+ {
+ RegistrationBuilder.ForDelegate(type, (c, p) => ((Mock)Activator.CreateInstance(typeof(Mock<>).MakeGenericType(type))).Object)
+ .As(service).InstancePerLifetimeScope().CreateRegistration()
+ };
+ }
+ }
+ }
+}
diff --git a/Tests/Resgrid.Tests/Services/WorkflowEventProviderScopeTests.cs b/Tests/Resgrid.Tests/Services/WorkflowEventProviderScopeTests.cs
new file mode 100644
index 000000000..a12a6749e
--- /dev/null
+++ b/Tests/Resgrid.Tests/Services/WorkflowEventProviderScopeTests.cs
@@ -0,0 +1,101 @@
+using System;
+using System.Collections.Generic;
+using System.Linq;
+using System.Threading;
+using System.Threading.Tasks;
+using Autofac;
+using FluentAssertions;
+using Moq;
+using NUnit.Framework;
+using Resgrid.Model;
+using Resgrid.Model.Events;
+using Resgrid.Model.Providers;
+using Resgrid.Model.Queue;
+using Resgrid.Model.Repositories;
+using Resgrid.Model.Services;
+using Resgrid.Providers.Bus;
+
+namespace Resgrid.Tests.Services
+{
+ ///
+ /// WorkflowEventProvider is a singleton, so its scoped dependencies must come from a child scope per event.
+ /// Captured at construction they lived in the root scope, where every concurrent event shared one unit of work:
+ /// a Records run insert opened a transaction on it and the other handlers' queries collided on its connection
+ /// ("The connection does not support MultipleActiveResultSets").
+ ///
+ [TestFixture, NonParallelizable]
+ public sealed class WorkflowEventProviderScopeTests
+ {
+ /// A container carrying the provider's scoped dependencies, for tests that construct it directly.
+ internal static IContainer Services(IWorkflowRepository workflows, IWorkflowRunRepository runs, IDepartmentsService departments,
+ ISubscriptionsService subscriptions, IProtectedProjectionService projection, Lazy history = null)
+ {
+ var builder = new ContainerBuilder();
+ builder.RegisterInstance(workflows).As();
+ builder.RegisterInstance(runs).As();
+ builder.RegisterInstance(departments).As();
+ builder.RegisterInstance(subscriptions).As();
+ builder.RegisterInstance(projection).As();
+ if (history != null)
+ builder.Register(_ => history.Value).As();
+ return builder.Build();
+ }
+
+ [Test]
+ public async Task Concurrent_events_each_get_their_own_run_repository_scope()
+ {
+ // A department no other test touches, so the static per-minute rate limiter starts empty.
+ var departmentId = 900000 + new Random().Next(99999);
+ var trigger = WorkflowTriggerEventType.RecordCreated;
+
+ var workflows = new Mock();
+ workflows.Setup(s => s.GetAllActiveByDepartmentAndEventTypeAsync(departmentId, (int)trigger))
+ .ReturnsAsync(new[] { new Workflow { WorkflowId = Guid.NewGuid().ToString(), DepartmentId = departmentId, TriggerEventType = (int)trigger } });
+ var subscriptions = new Mock();
+ subscriptions.Setup(s => s.GetCurrentPlanForDepartmentAsync(departmentId, It.IsAny())).ReturnsAsync(new Plan { PlanId = 999999 });
+ var projection = new Mock();
+ projection.Setup(s => s.BuildSafeWorkflowPayloadAsync(departmentId, It.IsAny