From 701c982b2add09cf5402e2be3bad513fa5f830af Mon Sep 17 00:00:00 2001 From: Shawn Jackson Date: Wed, 23 Sep 2026 15:27:16 -0700 Subject: [PATCH 1/2] RG-T55 RMS Occupancy Bug fixe --- Core/Resgrid.Services/Records/RecordsOccupancyService.cs | 8 ++++++-- Tests/Resgrid.Tests/Rms/RecordsOccupancyServiceTests.cs | 8 ++++++-- Tests/Resgrid.Tests/Rms/RmsPreventionFakes.cs | 5 ++++- 3 files changed, 16 insertions(+), 5 deletions(-) diff --git a/Core/Resgrid.Services/Records/RecordsOccupancyService.cs b/Core/Resgrid.Services/Records/RecordsOccupancyService.cs index d78763f18..ac7b61416 100644 --- a/Core/Resgrid.Services/Records/RecordsOccupancyService.cs +++ b/Core/Resgrid.Services/Records/RecordsOccupancyService.cs @@ -40,6 +40,7 @@ public class RecordsOccupancyService : IRecordsOccupancyService, IContactPreplan private readonly IContactsRepository _contacts; private readonly IAddressRepository _addresses; private readonly IPoisRepository _pois; + private readonly IPoiTypesRepository _poiTypes; private readonly IProtectedReadService _protectedReads; private readonly IProtectedGrantContext _grant; private readonly IRecordsProtectionService _protection; @@ -49,7 +50,7 @@ public RecordsOccupancyService(RecordsPreventionGate gate, IRmsOccupanciesReposi IRmsOccupancyHazardsRepository hazards, IRmsOccupancyCrosswalksRepository crosswalks, IRmsOccupancyFieldProvenancesRepository provenance, IRmsOccupancyOwnershipsRepository ownerships, IRmsViolationsRepository violations, IRmsHydrantsRepository hydrants, IContactPreplanRepository contactPreplans, IContactPreplanHazardRepository contactHazards, IContactsRepository contacts, IAddressRepository addresses, - IPoisRepository pois, IProtectedReadService protectedReads, IProtectedGrantContext grant, IRecordsProtectionService protection, IUnitOfWork unitOfWork) + IPoisRepository pois, IPoiTypesRepository poiTypes, IProtectedReadService protectedReads, IProtectedGrantContext grant, IRecordsProtectionService protection, IUnitOfWork unitOfWork) { _gate = gate; _occupancies = occupancies; @@ -65,6 +66,7 @@ public RecordsOccupancyService(RecordsPreventionGate gate, IRmsOccupanciesReposi _contacts = contacts; _addresses = addresses; _pois = pois; + _poiTypes = poiTypes; _protectedReads = protectedReads; _grant = grant; _protection = protection; @@ -490,7 +492,9 @@ async Task AddressOf(Contact c) list.Add(new SourceCandidate { Kind = RmsOccupancyCrosswalkSourceKind.Contact, SourceId = contact.ContactId, ContactId = contact.ContactId, DisplayName = contact.Name, NormalizedAddress = address, Latitude = point.HasValue ? (decimal?)(decimal)point.Value.Latitude : null, Longitude = point.HasValue ? (decimal?)(decimal)point.Value.Longitude : null }); } - foreach (var poi in (await _pois.GetAllByDepartmentIdAsync(departmentId)) ?? Enumerable.Empty()) + // Pois carry no DepartmentId column; a department owns its POIs through their POI type. + var poiTypes = (await _poiTypes.GetPoiTypesByDepartmentIdAsync(departmentId)) ?? Enumerable.Empty(); + foreach (var poi in poiTypes.Where(t => t?.Pois != null).SelectMany(t => t.Pois).Where(p => p != null)) { list.Add(new SourceCandidate { Kind = RmsOccupancyCrosswalkSourceKind.Poi, SourceId = poi.PoiId.ToString(), DisplayName = poi.Name, NormalizedAddress = AddressNormalizer.Normalize(poi.Address), Latitude = poi.Latitude == 0 && poi.Longitude == 0 ? (decimal?)null : (decimal)poi.Latitude, Longitude = poi.Latitude == 0 && poi.Longitude == 0 ? (decimal?)null : (decimal)poi.Longitude }); diff --git a/Tests/Resgrid.Tests/Rms/RecordsOccupancyServiceTests.cs b/Tests/Resgrid.Tests/Rms/RecordsOccupancyServiceTests.cs index 4cf2b9a2c..e5a5bfd68 100644 --- a/Tests/Resgrid.Tests/Rms/RecordsOccupancyServiceTests.cs +++ b/Tests/Resgrid.Tests/Rms/RecordsOccupancyServiceTests.cs @@ -97,7 +97,11 @@ private void SeedContactsWorld() _h.ContactPreplans.Setup(p => p.GetPreplansByDepartmentIdAsync(Dept)).ReturnsAsync(new List { preplan }); _h.ContactPreplans.Setup(p => p.GetPreplanByContactIdAsync("c1", Dept)).ReturnsAsync(preplan); _h.ContactHazards.Setup(h => h.GetHazardsByContactIdAsync("c1", Dept)).ReturnsAsync(new List { new ContactPreplanHazard { ContactPreplanHazardId = "ch1", ContactPreplanId = "pp1", ContactId = "c1", Title = "Propane", Severity = 3, ShouldAlert = true, Description = "500 gal" } }); - _h.Pois.Setup(p => p.GetAllByDepartmentIdAsync(Dept)).ReturnsAsync(new List { new Poi { PoiId = 77, Name = "Water tower", Address = "1 Hill Ct", Latitude = 45.7, Longitude = -122.7 } }); + _h.PoiTypes.Setup(p => p.GetPoiTypesByDepartmentIdAsync(Dept)).ReturnsAsync(new List + { + new PoiType { PoiTypeId = 5, DepartmentId = Dept, Name = "Water", Pois = new List { new Poi { PoiId = 77, PoiTypeId = 5, Name = "Water tower", Address = "1 Hill Ct", Latitude = 45.7, Longitude = -122.7 } } }, + new PoiType { PoiTypeId = 6, DepartmentId = Dept, Name = "Empty", Pois = new List() } + }); } [Test] @@ -221,7 +225,7 @@ public async Task Ownership_lookup_failures_leave_contacts_in_charge() { var broken = new Mock(); broken.Setup(o => o.GetForDepartmentAsync(It.IsAny())).ThrowsAsync(new InvalidOperationException("relation does not exist")); - var service = new RecordsOccupancyService(_h.Gate, _h.Occupancies, _h.Links, _h.Hazards, _h.Crosswalks, _h.Provenance, broken.Object, _h.Violations, _h.Hydrants, _h.ContactPreplans.Object, _h.ContactHazards.Object, _h.Contacts.Object, _h.Addresses.Object, _h.Pois.Object, _h.ProtectedReads.Object, _h.Grant.Object, _h.Protection, _h.UnitOfWork.Object); + var service = new RecordsOccupancyService(_h.Gate, _h.Occupancies, _h.Links, _h.Hazards, _h.Crosswalks, _h.Provenance, broken.Object, _h.Violations, _h.Hydrants, _h.ContactPreplans.Object, _h.ContactHazards.Object, _h.Contacts.Object, _h.Addresses.Object, _h.Pois.Object, _h.PoiTypes.Object, _h.ProtectedReads.Object, _h.Grant.Object, _h.Protection, _h.UnitOfWork.Object); (await service.IsRecordsOwnedAsync(Dept)).Should().BeFalse(); (await service.GetPreplanProjectionsAsync(Dept, new[] { "c1" })).Should().BeEmpty(); } diff --git a/Tests/Resgrid.Tests/Rms/RmsPreventionFakes.cs b/Tests/Resgrid.Tests/Rms/RmsPreventionFakes.cs index 765ce085e..101aae3c0 100644 --- a/Tests/Resgrid.Tests/Rms/RmsPreventionFakes.cs +++ b/Tests/Resgrid.Tests/Rms/RmsPreventionFakes.cs @@ -326,6 +326,7 @@ public sealed class RmsPreventionHarness public Mock Contacts { get; } = new Mock(); public Mock Addresses { get; } = new Mock(); public Mock Pois { get; } = new Mock(); + public Mock PoiTypes { get; } = new Mock(); public Mock Reports { get; } = new Mock(); public Mock Records { get; } = new Mock(); public Mock Units { get; } = new Mock(); @@ -394,9 +395,11 @@ public RmsPreventionHarness() ProtectedReads.Setup(r => r.ResolveContactPreplanHazardsForReadAsync(It.IsAny(), It.IsAny>(), It.IsAny(), It.IsAny(), It.IsAny())).ReturnsAsync(new ProtectedReadResult()); Scanner.Setup(s => s.ScanAsync(It.IsAny(), It.IsAny(), It.IsAny(), It.IsAny())).ReturnsAsync(new Resgrid.Model.Providers.RecordAttachmentScanResult { State = RmsAttachmentScanState.Skipped }); Grant.SetupGet(g => g.UserId).Returns(Admin); + // Pois has no DepartmentId column (RESGRID-WEB-1MQ); department POIs are read through their POI types. + Pois.Setup(p => p.GetAllByDepartmentIdAsync(It.IsAny())).ThrowsAsync(new InvalidOperationException("Invalid column name 'DepartmentId'.")); Gate = new RecordsPreventionGate(Cutover.Object, Flags.Object, Authorization.Object, Sequences, Audits); - OccupancyService = new RecordsOccupancyService(Gate, Occupancies, Links, Hazards, Crosswalks, Provenance, Ownerships, Violations, Hydrants, ContactPreplans.Object, ContactHazards.Object, Contacts.Object, Addresses.Object, Pois.Object, ProtectedReads.Object, Grant.Object, Protection, UnitOfWork.Object); + OccupancyService = new RecordsOccupancyService(Gate, Occupancies, Links, Hazards, Crosswalks, Provenance, Ownerships, Violations, Hydrants, ContactPreplans.Object, ContactHazards.Object, Contacts.Object, Addresses.Object, Pois.Object, PoiTypes.Object, ProtectedReads.Object, Grant.Object, Protection, UnitOfWork.Object); InspectionsService = new RecordsInspectionsService(Gate, CodeSets, CodeSections, Programs, Inspections, Violations, Occupancies, Attachments, Protection, Outbox, UnitOfWork.Object); HydrantsService = new RecordsHydrantsService(Gate, Hydrants, FlowTests, Maintenance, Attachments, UnitOfWork.Object); PermitsService = new RecordsPermitsService(Gate, PermitTypes, Permits, PlanReviews, Occupancies, Attachments, Protection, Outbox, UnitOfWork.Object); From 77354857e8f197500b5e50f4837ee755006e13d3 Mon Sep 17 00:00:00 2001 From: Shawn Jackson Date: Wed, 23 Sep 2026 21:18:17 -0700 Subject: [PATCH 2/2] RG-T41 Search fix --- Core/Resgrid.Search/LuceneIndexHost.cs | 40 +++++++++++--- .../Store/S3SearchIndexStore.cs | 5 +- .../Search/SearchIndexStoreSyncTests.cs | 52 +++++++++++++++++++ .../Tasks/PaymentQueueProcessorTask.cs | 6 ++- .../Logic/PaymentQueueLogic.cs | 1 + 5 files changed, 94 insertions(+), 10 deletions(-) diff --git a/Core/Resgrid.Search/LuceneIndexHost.cs b/Core/Resgrid.Search/LuceneIndexHost.cs index 4ed5ff5db..3ee542e47 100644 --- a/Core/Resgrid.Search/LuceneIndexHost.cs +++ b/Core/Resgrid.Search/LuceneIndexHost.cs @@ -29,6 +29,7 @@ public class LuceneIndexHost : IDisposable private const string ManifestFileName = "manifest.json"; private const string LockFileName = "write.lock"; private static readonly TimeSpan OrphanedTempAge = TimeSpan.FromHours(1); + private static readonly TimeSpan MaxPullBackoff = TimeSpan.FromMinutes(10); private readonly object _sync = new object(); private readonly object _writerSyncGate = new object(); @@ -47,6 +48,7 @@ public class LuceneIndexHost : IDisposable private long _appliedGeneration = -1; private string _manifestETag; private DateTime _lastPullAttemptUtc = DateTime.MinValue; + private int _consecutivePullFailures; private Task _pullTask; /// Production constructor: the configured local path under SearchConfig.IndexPath. @@ -169,7 +171,7 @@ public SearcherManager GetSearcherManager() } if (StoreEnabled) - StartBackgroundPullIfDue(force: _appliedRevision == null); + StartBackgroundPullIfDue(); if (!DirectoryReader.IndexExists(Store)) return null; @@ -188,7 +190,7 @@ public void MaybeRefresh() { manager = _searcherManager; if (_writer == null && StoreEnabled) - StartBackgroundPullIfDue(force: false); + StartBackgroundPullIfDue(); } try { manager?.MaybeRefresh(); } @@ -376,22 +378,44 @@ public async Task ResetFromStoreAsync(CancellationToken cancellationToken = defa await PullCoreAsync(cancellationToken); } - private void StartBackgroundPullIfDue(bool force) + private void StartBackgroundPullIfDue() { - // Called under _sync. + // Called under _sync. Always throttled, including before the first revision is applied: the anonymous health + // endpoint lands here, so an unthrottled retry turns a store that keeps failing (bad credentials, outage) into an + // object-store request and a logged exception per probe. The first call is always due (_lastPullAttemptUtc is + // MinValue), so a fresh reader still pulls immediately. if (_pullTask != null && !_pullTask.IsCompleted) return; - var due = force || (DateTime.UtcNow - _lastPullAttemptUtc).TotalSeconds >= Math.Max(5, SearchConfig.ReaderPullSeconds); - if (!due) + if (DateTime.UtcNow - _lastPullAttemptUtc < PullInterval()) return; _lastPullAttemptUtc = DateTime.UtcNow; _pullTask = Task.Run(async () => { - try { await PullCoreAsync(CancellationToken.None); } - catch (Exception ex) { Logging.LogException(ex, $"Search index '{IndexName}' pull from the object store failed."); } + try + { + await PullCoreAsync(CancellationToken.None); + _consecutivePullFailures = 0; + } + catch (Exception ex) + { + // Only one background pull runs at a time, so the counter has a single writer. Error, not Fatal: the reader + // keeps serving its last local revision and the next attempt backs off. + var failures = ++_consecutivePullFailures; + Logging.LogError(ex, $"Search index '{IndexName}' pull from the object store failed ({failures} in a row); next attempt in {PullInterval().TotalSeconds:0}s at the earliest."); + } }); } + /// ReaderPullSeconds, doubled per consecutive failed background pull up to . + private TimeSpan PullInterval() + { + var interval = TimeSpan.FromSeconds(Math.Max(5, SearchConfig.ReaderPullSeconds)); + if (_consecutivePullFailures == 0 || interval >= MaxPullBackoff) + return interval; + var backoffSeconds = interval.TotalSeconds * Math.Pow(2, Math.Min(_consecutivePullFailures, 10)); + return TimeSpan.FromSeconds(Math.Min(backoffSeconds, MaxPullBackoff.TotalSeconds)); + } + private async Task PullCoreAsync(CancellationToken cancellationToken) { _lastPullAttemptUtc = DateTime.UtcNow; diff --git a/Core/Resgrid.Search/Store/S3SearchIndexStore.cs b/Core/Resgrid.Search/Store/S3SearchIndexStore.cs index acfdfb540..a8ec0a489 100644 --- a/Core/Resgrid.Search/Store/S3SearchIndexStore.cs +++ b/Core/Resgrid.Search/Store/S3SearchIndexStore.cs @@ -51,7 +51,10 @@ private static IAmazonS3 CreateClient() UseHttp = !SearchConfig.S3UseSsl, AuthenticationRegion = string.IsNullOrWhiteSpace(SearchConfig.S3Region) ? "us-east-1" : SearchConfig.S3Region }; - var credentials = new BasicAWSCredentials(SearchConfig.S3AccessKey ?? string.Empty, SearchConfig.S3SecretKey ?? string.Empty); + // Trimmed: ConfigProcessor passes environment values through verbatim, and a Kubernetes Secret built from a file or + // an unterminated echo keeps its trailing newline. A padded secret key still identifies the access key but fails + // every request with SignatureDoesNotMatch; real S3 keys never carry surrounding whitespace. + var credentials = new BasicAWSCredentials((SearchConfig.S3AccessKey ?? string.Empty).Trim(), (SearchConfig.S3SecretKey ?? string.Empty).Trim()); return new AmazonS3Client(credentials, config); } diff --git a/Tests/Resgrid.Tests/Search/SearchIndexStoreSyncTests.cs b/Tests/Resgrid.Tests/Search/SearchIndexStoreSyncTests.cs index 297bd287b..a89ef8ef2 100644 --- a/Tests/Resgrid.Tests/Search/SearchIndexStoreSyncTests.cs +++ b/Tests/Resgrid.Tests/Search/SearchIndexStoreSyncTests.cs @@ -129,6 +129,32 @@ public async Task Publish_fails_when_another_writer_published_first() store.Manifests[SearchIndexNames.Global].Revision.Should().Be("elsewhere", "the losing writer never overwrites the other writer's manifest"); } + [Test] + public async Task A_failing_store_is_not_polled_on_every_reader_call() + { + // The anonymous health endpoint calls GetSearcherManager on every probe; with no revision applied yet each call + // used to force a pull, so a store rejecting every request (e.g. SignatureDoesNotMatch) was hit once per probe. + var store = new FailingSearchIndexStore(); + var reader = Host(TempDir(), store); + + reader.GetSearcherManager().Should().BeNull("nothing has been pulled"); + var deadline = DateTime.UtcNow.AddSeconds(10); + while (store.ManifestReads == 0 && DateTime.UtcNow < deadline) + await Task.Delay(10); + store.ManifestReads.Should().Be(1, "a fresh reader pulls on first use"); + await Task.Delay(100); // let the failed pull complete so the next call is not merely skipped as in-flight + + for (var i = 0; i < 25; i++) + { + reader.GetSearcherManager().Should().BeNull(); + reader.MaybeRefresh(); + } + await Task.Delay(100); + + store.ManifestReads.Should().Be(1, "retries wait for ReaderPullSeconds (and back off) instead of firing on every call"); + reader.LastSyncedRevision.Should().BeNull(); + } + [Test] public async Task Without_a_store_commit_is_local_only() { @@ -141,6 +167,32 @@ public async Task Without_a_store_commit_is_local_only() } } + /// A store that rejects every request, as one with a bad secret key does; counts manifest reads. + public sealed class FailingSearchIndexStore : ISearchIndexStore + { + private int _manifestReads; + + public int ManifestReads => Volatile.Read(ref _manifestReads); + + public bool Enabled => true; + + public Task GetManifestAsync(string indexName, CancellationToken cancellationToken = default) + { + Interlocked.Increment(ref _manifestReads); + return Task.FromException(new InvalidOperationException("SignatureDoesNotMatch")); + } + + public Task> ListFilesAsync(string indexName, CancellationToken cancellationToken = default) => throw new InvalidOperationException("SignatureDoesNotMatch"); + + public Task UploadFileAsync(string indexName, string fileName, string localPath, CancellationToken cancellationToken = default) => throw new InvalidOperationException("SignatureDoesNotMatch"); + + public Task DownloadFileAsync(string indexName, string fileName, string localPath, CancellationToken cancellationToken = default) => throw new InvalidOperationException("SignatureDoesNotMatch"); + + public Task DeleteFilesAsync(string indexName, IEnumerable fileNames, CancellationToken cancellationToken = default) => throw new InvalidOperationException("SignatureDoesNotMatch"); + + public Task PutManifestAsync(string indexName, SearchIndexManifest manifest, string expectedETag, CancellationToken cancellationToken = default) => throw new InvalidOperationException("SignatureDoesNotMatch"); + } + /// S3 semantics in memory: immutable objects, one manifest per index, If-None-Match:* / If-Match on the manifest. public sealed class InMemorySearchIndexStore : ISearchIndexStore { diff --git a/Workers/Resgrid.Workers.Console/Tasks/PaymentQueueProcessorTask.cs b/Workers/Resgrid.Workers.Console/Tasks/PaymentQueueProcessorTask.cs index 5eac7cda9..9cdc13a2b 100644 --- a/Workers/Resgrid.Workers.Console/Tasks/PaymentQueueProcessorTask.cs +++ b/Workers/Resgrid.Workers.Console/Tasks/PaymentQueueProcessorTask.cs @@ -5,6 +5,7 @@ using Resgrid.Providers.Bus.Rabbit; using Resgrid.Workers.Console.Commands; using Resgrid.Workers.Framework.Logic; +using System; using System.Threading; using System.Threading.Tasks; @@ -44,7 +45,10 @@ public async Task ProcessAsync(PaymentQueueProcessorCommand command, IQuidjiboPr private async Task OnPaymentEventQueueReceived(CqrsEvent cqrs) { _logger.LogInformation($"{Name}: Payment Queue Received with a type of {cqrs.Type}, starting processing..."); - await PaymentQueueLogic.ProcessPaymentQueueItem(cqrs); + // RabbitInboundQueueProvider only retries when the handler throws; returning normally acks the message. + if (!await PaymentQueueLogic.ProcessPaymentQueueItem(cqrs)) + throw new InvalidOperationException($"{Name}: Payment queue item with type of {cqrs.Type} failed processing."); + _logger.LogInformation($"{Name}: Finished processing of Payment queue item with type of {cqrs.Type}."); } } diff --git a/Workers/Resgrid.Workers.Framework/Logic/PaymentQueueLogic.cs b/Workers/Resgrid.Workers.Framework/Logic/PaymentQueueLogic.cs index dc558d7ae..fcb36645b 100644 --- a/Workers/Resgrid.Workers.Framework/Logic/PaymentQueueLogic.cs +++ b/Workers/Resgrid.Workers.Framework/Logic/PaymentQueueLogic.cs @@ -143,6 +143,7 @@ public static async Task ProcessPaymentQueueItem(CqrsEvent qi) } catch (Exception ex) { + success = false; Logging.LogException(ex); Logging.SendExceptionEmail(ex, "ProcessPaymentQueueItem"); }