diff --git a/Core/Resgrid.Search/LuceneIndexHost.cs b/Core/Resgrid.Search/LuceneIndexHost.cs
index 4ed5ff5d..3ee542e4 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 acfdfb54..a8ec0a48 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/Core/Resgrid.Services/Records/RecordsOccupancyService.cs b/Core/Resgrid.Services/Records/RecordsOccupancyService.cs
index d78763f1..ac7b6141 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 4cf2b9a2..e5a5bfd6 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 765ce085..101aae3c 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);
diff --git a/Tests/Resgrid.Tests/Search/SearchIndexStoreSyncTests.cs b/Tests/Resgrid.Tests/Search/SearchIndexStoreSyncTests.cs
index 297bd287..a89ef8ef 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 5eac7cda..9cdc13a2 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 dc558d7a..fcb36645 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");
}