From 8bf2e76007fe6d833f79829340ddb4b1fce0a7d0 Mon Sep 17 00:00:00 2001 From: "Eric J. Smith" Date: Mon, 7 Sep 2026 17:05:37 -0500 Subject: [PATCH 1/8] Load daily and monthly index mappings asynchronously --- .../skills/foundatio-repositories/SKILL.md | 1 + .../references/index-lifecycle.md | 9 +- docs/guide/index-management.md | 21 +- .../Configuration/DailyIndex.cs | 66 +++-- .../Configuration/Index.cs | 2 + .../Configuration/MonthlyIndex.cs | 2 +- ...oundatio.Repositories.Elasticsearch.csproj | 2 +- .../DailyIndexMappingTests.cs | 268 ++++++++++++++++++ 8 files changed, 348 insertions(+), 23 deletions(-) create mode 100644 tests/Foundatio.Repositories.Elasticsearch.Tests/DailyIndexMappingTests.cs diff --git a/.agents/skills/foundatio-repositories/SKILL.md b/.agents/skills/foundatio-repositories/SKILL.md index 233788f2..9c85c94f 100644 --- a/.agents/skills/foundatio-repositories/SKILL.md +++ b/.agents/skills/foundatio-repositories/SKILL.md @@ -132,6 +132,7 @@ IReadOnlyRepository - **`ExistsAsync(query)` is a dirty read**: Uses the Search API (`size: 0`), NOT the realtime Document Exists API. After a write without `ImmediateConsistency`, it can return stale results. - **`ExistsAsync(id)` is real-time even with soft deletes**: Uses the GET API with a source filter for `IsDeleted`. - **Register repositories as singletons**: Repository instances maintain internal state (index configuration, cache references). +- **Keep mapping resolvers long-lived**: Daily/monthly query mapping loads use async names-and-aliases discovery plus one partition mapping request. Concurrent lookups share a load per index and process; disposing the index cancels its active mapping load. See [index lifecycle](references/index-lifecycle.md#mapping-resolver-cache) for refresh behavior. - **`FieldEquals` with multiple values is OR**: `.FieldEquals(e => e.Field, "A", "B")` produces an OR filter, not AND. - **`FieldContains` is token matching, NOT wildcard**: `FieldContains(f => f.Name, "Er")` will NOT match "Eric". Use `FilterExpression("field:pattern*")` for prefix/wildcard matching. - **`FieldNot` is AND-NOT**: Multiple conditions inside `FieldNot` mean NOT A AND NOT B. For NOT (A AND B), nest `FieldAnd` inside `FieldNot`. diff --git a/.agents/skills/foundatio-repositories/references/index-lifecycle.md b/.agents/skills/foundatio-repositories/references/index-lifecycle.md index f89d554e..dac74e2e 100644 --- a/.agents/skills/foundatio-repositories/references/index-lifecycle.md +++ b/.agents/skills/foundatio-repositories/references/index-lifecycle.md @@ -337,11 +337,18 @@ var results = await repository.FindAsync(q => q.Index("logs-last-7-days")); | Cache layer | Lifetime | How to invalidate | |---|---|---| -| `ElasticMappingResolver` field cache | Auto-refreshes ~60 seconds | `index.MappingResolver.RefreshMapping()` | +| `ElasticMappingResolver` field cache | Snapshot lifetime; unresolved fields trigger reloads with a five-second cooldown | `index.MappingResolver.RefreshMapping()` | | `_isEnsured` flag (Index/VersionedIndex) | Process lifetime | App restart or index deletion | | `_ensuredDates` (DailyIndex) | Process lifetime per-date | `DeleteAsync(name)` or `Dispose()` | | `ConfigureIndexesAsync` cache marker | 5 minutes (distributed) | Expires automatically; or `ConfigureIndexesAsync(force: true)` | +Daily/monthly resolvers asynchronously discover the newest partition using names and aliases only, then load +that partition's mapping. Concurrent lookups share one load per resolver. Index disposal cancels active mapping +I/O without initializing an unused resolver. Keep indexes long-lived and do not invalidate after every write. +The cooldown is not a freshness guarantee, and successful lookups do not trigger periodic refreshes. Use +`RefreshMapping()` after known mapping changes, including changes to an already-resolved alias. Mappings from +historical partitions are not merged into the newest partition's mapping. + ## Index Operations ### ConfigureIndexesAsync diff --git a/docs/guide/index-management.md b/docs/guide/index-management.md index 24f84170..c39d7bdf 100644 --- a/docs/guide/index-management.md +++ b/docs/guide/index-management.md @@ -873,11 +873,24 @@ POST /logs-v1-2025.05.*/_update_by_query?conflicts=proceed The repository framework does **not** cache the PUT Mapping request/response (that's purely server-side). However, the **query parser** uses an `ElasticMappingResolver` that caches field-to-type resolution for building queries, sorting, and aggregations. This resolver combines two sources: 1. **Code mapping** — derived from your `ConfigureIndexMapping` method at startup (immutable for the process lifetime) -2. **Server mapping** — fetched from the Elasticsearch GET Mapping API, cached in memory and **automatically refreshed at most once per minute** +2. **Server mapping** — loaded on first use, then reloaded when a field cannot be resolved, with a five-second cooldown between automatic reload attempts. Successful lookups do not trigger periodic reloads. + +Daily and monthly indexes discover the newest partition using an index request limited to names and aliases, +then retrieve the full mapping of that single partition. Both requests use asynchronous I/O on asynchronous +query paths. Concurrent lookups share one load through the index's long-lived resolver. Synchronous resolver +calls remain supported, but block while the asynchronous load completes. Disposing an index disposes its +initialized resolver and cancels outstanding mapping I/O without creating an unused resolver. + +Each reload discovers the latest partition again so newly created partitions are visible without restarting +the application. This does not merge mappings across historical partitions. Keep indexes and their resolvers +long-lived; creating one per request bypasses their caches and multiplies metadata requests. Automatic reload +limits apply independently to each resolver in each process. Do not call `RefreshMapping()` after every write. #### What this means after a manual PUT Mapping -If you manually apply a mapping change (e.g., `PUT /index/_mapping` via the Elasticsearch API or a script), the `ElasticMappingResolver` will automatically pick it up within ~60 seconds on the next field resolution. You typically do not need to do anything in application code. +Requests for an unresolved field can discover new mappings after the cooldown. The cooldown is not a freshness +guarantee: slow or failed requests can extend the stale window. Changes to an already-resolved field or alias +require explicit invalidation. If you need immediate recognition (e.g., in tests or a migration script that queries the new field right after applying the mapping), call: @@ -891,7 +904,7 @@ This clears the cached server mapping and forces the next `GetMapping()` call to | Cache layer | Lifetime | How to invalidate | |---|---|---| -| `ElasticMappingResolver` field cache | Auto-refreshes from server every ~60 seconds | `index.MappingResolver.RefreshMapping()` | +| `ElasticMappingResolver` field cache | Snapshot lifetime; unresolved fields trigger reloads with a five-second cooldown | `index.MappingResolver.RefreshMapping()` | | `_isEnsured` flag (`Index` / `VersionedIndex`) | Process lifetime (one-time flag) | Deleting the index resets it; otherwise persists until app restart | | `_ensuredDates` (`DailyIndex`) | Process lifetime per-date | Cleared on `DeleteAsync(name)` or `Dispose()`; otherwise persists until app restart | | `ConfigureIndexesAsync` cache marker | 5 minutes (distributed via `ICacheClient`) | Automatically expires; or call `ConfigureIndexesAsync(force: true)` | @@ -900,7 +913,7 @@ This clears the cached server mapping and forces the next `GetMapping()` call to Elasticsearch itself has no mapping cache you need to invalidate — once a PUT Mapping succeeds, the mapping is immediately active for new indexing and queries. The only caching is in-process within the .NET application: -- **For queries**: The `ElasticMappingResolver` auto-refreshes. If you need it sooner, call `RefreshMapping()`. +- **For queries**: Unresolved fields can trigger mapping reloads. For a known mapping change that must be visible immediately, call `RefreshMapping()` after Elasticsearch acknowledges the change. - **For writes**: The `_isEnsured` / `_ensuredDates` flags only control whether `ConfigureAsync` runs again. They don't prevent writes to the index — they just skip redundant index creation/mapping calls. Manual PUT Mapping changes are orthogonal to these flags. ### In-Place Analysis Updates (analyzers, tokenizers, filters) diff --git a/src/Foundatio.Repositories.Elasticsearch/Configuration/DailyIndex.cs b/src/Foundatio.Repositories.Elasticsearch/Configuration/DailyIndex.cs index 3b3da64a..20ef3937 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Configuration/DailyIndex.cs +++ b/src/Foundatio.Repositories.Elasticsearch/Configuration/DailyIndex.cs @@ -5,6 +5,7 @@ using System.Globalization; using System.Linq; using System.Linq.Expressions; +using System.Threading; using System.Threading.Tasks; using Elastic.Clients.Elasticsearch; using Elastic.Clients.Elasticsearch.IndexManagement; @@ -447,13 +448,38 @@ protected override string GetIndexByDate(DateTime date) protected override ElasticMappingResolver CreateMappingResolver() { - return ElasticMappingResolver.Create(GetLatestIndexMapping, Configuration.Client.Infer, _logger); + return ElasticMappingResolver.CreateWithAsyncLoader(GetLatestIndexMappingAsync, Configuration.Client.Infer, logger: _logger); } protected TypeMapping? GetLatestIndexMapping() { string filter = $"{Name}-v{Version}-*"; var indicesResponse = Configuration.Client.Indices.Get((Indices)(IndexName)filter, d => d.LimitToNamesAndAliases()); + string? latestIndex = GetLatestIndexName(indicesResponse, filter); + if (latestIndex is null) + return null; + + var mappingResponse = Configuration.Client.Indices.GetMapping(new GetMappingRequest(latestIndex)); + return GetLatestIndexMapping(mappingResponse, latestIndex); + } + + /// Loads the mapping of the newest partition using asynchronous metadata requests. + /// Index discovery retrieves names and aliases only. The full mapping is requested for one partition. + protected async Task GetLatestIndexMappingAsync(CancellationToken cancellationToken = default) + { + string filter = $"{Name}-v{Version}-*"; + var indicesResponse = await Configuration.Client.Indices.GetAsync((Indices)(IndexName)filter, + d => d.LimitToNamesAndAliases(), cancellationToken).AnyContext(); + string? latestIndex = GetLatestIndexName(indicesResponse, filter); + if (latestIndex is null) + return null; + + var mappingResponse = await Configuration.Client.Indices.GetMappingAsync(new GetMappingRequest(latestIndex), cancellationToken).AnyContext(); + return GetLatestIndexMapping(mappingResponse, latestIndex); + } + + private string? GetLatestIndexName(GetIndexResponse indicesResponse, string filter) + { if (!indicesResponse.IsValidResponse) { if (indicesResponse.ElasticsearchServerError?.Status == 404) @@ -462,31 +488,39 @@ protected override ElasticMappingResolver CreateMappingResolver() throw new RepositoryException(indicesResponse.GetErrorMessage($"Error getting latest index mapping {filter}"), indicesResponse.OriginalException()); } - var latestIndex = indicesResponse.Indices.Keys - .Where(i => GetIndexVersion(i.ToString()) == Version) - .Select(i => + string? latestIndex = null; + var latestDate = DateTime.MinValue; + foreach (var index in indicesResponse.Indices.Keys) + { + string name = index.ToString(); + if (GetIndexVersion(name) != Version) + continue; + + var date = GetIndexDate(name); + if (latestIndex is null || date > latestDate) { - string indexName = i.ToString(); - return new IndexInfo { DateUtc = GetIndexDate(indexName), Index = indexName, Version = GetIndexVersion(indexName) }; - }) - .OrderByDescending(i => i.DateUtc) - .FirstOrDefault(); + latestIndex = name; + latestDate = date; + } + } - if (latestIndex == null) - return null; + return latestIndex; + } - var mappingResponse = Configuration.Client.Indices.GetMapping(new GetMappingRequest(latestIndex.Index)); - _logger.LogTrace("GetMapping: {Request}", mappingResponse.GetRequest(false, true)); + private TypeMapping? GetLatestIndexMapping(GetMappingResponse mappingResponse, string latestIndex) + { + if (_logger.IsEnabled(LogLevel.Trace)) + _logger.LogTrace("GetMapping: {Request}", mappingResponse.GetRequest(false, true)); if (!mappingResponse.IsValidResponse) { if (mappingResponse.ApiCallDetails.HttpStatusCode.GetValueOrDefault() == 404) { - _logger.LogWarning("Index {Index} not found when getting mapping", latestIndex.Index); + _logger.LogWarning("Index {Index} not found when getting mapping", latestIndex); return null; } - _logger.LogError("Error getting mapping for {Index}: {Error}", latestIndex.Index, mappingResponse.ElasticsearchServerError); + _logger.LogError("Error getting mapping for {Index}: {Error}", latestIndex, mappingResponse.ElasticsearchServerError); return null; } @@ -591,7 +625,7 @@ public DailyIndex(IElasticConfiguration configuration, string? name = null, int protected override ElasticMappingResolver CreateMappingResolver() { - return ElasticMappingResolver.Create(ConfigureIndexMapping, Configuration.Client.Infer, GetLatestIndexMapping, _logger); + return ElasticMappingResolver.CreateWithAsyncLoader(ConfigureIndexMapping, Configuration.Client.Infer, GetLatestIndexMappingAsync, logger: _logger); } public virtual void ConfigureIndexMapping(TypeMappingDescriptor map) diff --git a/src/Foundatio.Repositories.Elasticsearch/Configuration/Index.cs b/src/Foundatio.Repositories.Elasticsearch/Configuration/Index.cs index 5bf45836..ed43b10c 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Configuration/Index.cs +++ b/src/Foundatio.Repositories.Elasticsearch/Configuration/Index.cs @@ -415,6 +415,8 @@ public virtual void Dispose() _disposedCancellationTokenSource.Cancel(); _disposedCancellationTokenSource.Dispose(); + if (_mappingResolver.IsValueCreated) + _mappingResolver.Value.Dispose(); } } diff --git a/src/Foundatio.Repositories.Elasticsearch/Configuration/MonthlyIndex.cs b/src/Foundatio.Repositories.Elasticsearch/Configuration/MonthlyIndex.cs index e1fb588c..786bc774 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Configuration/MonthlyIndex.cs +++ b/src/Foundatio.Repositories.Elasticsearch/Configuration/MonthlyIndex.cs @@ -64,7 +64,7 @@ public MonthlyIndex(IElasticConfiguration configuration, string? name = null, in protected override ElasticMappingResolver CreateMappingResolver() { - return ElasticMappingResolver.Create(ConfigureIndexMapping, Configuration.Client.Infer, GetLatestIndexMapping, _logger); + return ElasticMappingResolver.CreateWithAsyncLoader(ConfigureIndexMapping, Configuration.Client.Infer, GetLatestIndexMappingAsync, logger: _logger); } public virtual void ConfigureIndexMapping(TypeMappingDescriptor map) diff --git a/src/Foundatio.Repositories.Elasticsearch/Foundatio.Repositories.Elasticsearch.csproj b/src/Foundatio.Repositories.Elasticsearch/Foundatio.Repositories.Elasticsearch.csproj index 396f4421..66be5a70 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Foundatio.Repositories.Elasticsearch.csproj +++ b/src/Foundatio.Repositories.Elasticsearch/Foundatio.Repositories.Elasticsearch.csproj @@ -7,7 +7,7 @@ - + diff --git a/tests/Foundatio.Repositories.Elasticsearch.Tests/DailyIndexMappingTests.cs b/tests/Foundatio.Repositories.Elasticsearch.Tests/DailyIndexMappingTests.cs new file mode 100644 index 00000000..d35a6da3 --- /dev/null +++ b/tests/Foundatio.Repositories.Elasticsearch.Tests/DailyIndexMappingTests.cs @@ -0,0 +1,268 @@ +using System; +using System.Collections.Generic; +using System.Text; +using System.Threading; +using System.Threading.Tasks; +using Elastic.Clients.Elasticsearch; +using Elastic.Clients.Elasticsearch.Mapping; +using Elastic.Transport; +using Foundatio.Repositories.Elasticsearch.Configuration; +using Xunit; +using FieldMapping = Foundatio.Parsers.ElasticQueries.FieldMapping; + +namespace Foundatio.Repositories.Elasticsearch.Tests; + +public class DailyIndexMappingTests +{ + [Theory] + [InlineData(false, false)] + [InlineData(false, true)] + [InlineData(true, false)] + [InlineData(true, true)] + public async Task GetMappingAsync_WithTimeSeriesIndex_UsesAsyncMetadataRequests(bool monthly, bool typed) + { + string latest = monthly ? "events-v1-2026.09" : "events-v1-2026.09.07"; + string previous = monthly ? "events-v1-2026.08" : "events-v1-2026.09.06"; + using var invoker = new MappingRequestInvoker(previous, latest); + using var configuration = new MappingConfiguration(invoker); + using var index = CreateIndex(configuration, monthly, typed); + + var mapping = await index.MappingResolver.GetMappingAsync("dynamic", cancellationToken: TestContext.Current.CancellationToken); + + Assert.IsType(mapping?.Property); + Assert.Equal(0, invoker.SyncRequests); + Assert.Equal(2, invoker.Requests.Count); + Assert.Contains("features=aliases", invoker.Requests[0].Query); + Assert.Equal($"/{latest}/_mapping", invoker.Requests[1].AbsolutePath); + } + + [Fact] + public async Task GetMappingAsync_WhenDiscoveryIsBlocked_ReturnsWithoutBlockingCaller() + { + using var invoker = new MappingRequestInvoker("events-v1-2026.09.07") { BlockDiscovery = true }; + using var configuration = new MappingConfiguration(invoker); + using var index = new DailyIndex(configuration, "events"); + + var lookup = index.MappingResolver.GetMappingAsync("dynamic", cancellationToken: TestContext.Current.CancellationToken).AsTask(); + await invoker.DiscoveryStarted.Task.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + try + { + Assert.False(lookup.IsCompleted); + Assert.Equal(0, invoker.SyncRequests); + } + finally + { + invoker.ReleaseDiscovery.TrySetResult(); + } + + Assert.True((await lookup)?.Found); + } + + [Fact] + public async Task GetMappingAsync_WithConcurrentLookups_SharesMetadataRequests() + { + using var invoker = new MappingRequestInvoker("events-v1-2026.09.07") { BlockDiscovery = true }; + using var configuration = new MappingConfiguration(invoker); + using var index = new DailyIndex(configuration, "events"); + var resolver = index.MappingResolver; + var lookups = new Task[20]; + for (int i = 0; i < lookups.Length; i++) + lookups[i] = resolver.GetMappingAsync("dynamic", cancellationToken: TestContext.Current.CancellationToken).AsTask(); + + await invoker.DiscoveryStarted.Task.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + try + { + Assert.Single(invoker.Requests); + Assert.All(lookups, lookup => Assert.False(lookup.IsCompleted)); + } + finally + { + invoker.ReleaseDiscovery.TrySetResult(); + } + + var mappings = await Task.WhenAll(lookups); + Assert.All(mappings, mapping => Assert.True(mapping?.Found)); + Assert.Equal(2, invoker.Requests.Count); + } + + [Fact] + public void GetMapping_WithSynchronousCaller_StillResolvesMapping() + { + using var invoker = new MappingRequestInvoker("events-v1-2026.09.07"); + using var configuration = new MappingConfiguration(invoker); + using var index = new DailyIndex(configuration, "events"); + + var mapping = index.MappingResolver.GetMapping("dynamic"); + + Assert.IsType(mapping?.Property); + Assert.Equal(2, invoker.Requests.Count); + } + + [Fact] + public async Task GetMappingAsync_WithNoPartitions_DoesNotRequestMapping() + { + using var invoker = new MappingRequestInvoker(); + using var configuration = new MappingConfiguration(invoker); + using var index = new DailyIndex(configuration, "events"); + + var mapping = await index.MappingResolver.GetMappingAsync("dynamic", cancellationToken: TestContext.Current.CancellationToken); + + Assert.False(mapping?.Found); + Assert.Single(invoker.Requests); + } + + [Theory] + [InlineData(404)] + [InlineData(503)] + public async Task GetMappingAsync_WhenDiscoveryFails_DoesNotRequestMapping(int statusCode) + { + using var invoker = new MappingRequestInvoker("events-v1-2026.09.07") { DiscoveryStatusCode = statusCode }; + using var configuration = new MappingConfiguration(invoker); + using var index = new DailyIndex(configuration, "events"); + + var mapping = await index.MappingResolver.GetMappingAsync("dynamic", cancellationToken: TestContext.Current.CancellationToken); + + Assert.False(mapping?.Found); + Assert.Single(invoker.Requests); + } + + [Theory] + [InlineData(404)] + [InlineData(503)] + public async Task GetMappingAsync_WhenMappingFails_ReturnsUnmapped(int statusCode) + { + using var invoker = new MappingRequestInvoker("events-v1-2026.09.07") { MappingStatusCode = statusCode }; + using var configuration = new MappingConfiguration(invoker); + using var index = new DailyIndex(configuration, "events"); + + var mapping = await index.MappingResolver.GetMappingAsync("dynamic", cancellationToken: TestContext.Current.CancellationToken); + + Assert.False(mapping?.Found); + Assert.Equal(2, invoker.Requests.Count); + } + + [Fact] + public async Task RefreshMapping_AfterRollover_LoadsNewestPartition() + { + using var invoker = new MappingRequestInvoker("events-v1-2026.09.06"); + using var configuration = new MappingConfiguration(invoker); + using var index = new DailyIndex(configuration, "events"); + Assert.True((await index.MappingResolver.GetMappingAsync("dynamic", cancellationToken: TestContext.Current.CancellationToken))?.Found); + invoker.IndexNames = ["events-v1-2026.09.07", "events-v1-2026.09.06"]; + + index.MappingResolver.RefreshMapping(); + Assert.True((await index.MappingResolver.GetMappingAsync("dynamic", cancellationToken: TestContext.Current.CancellationToken))?.Found); + + Assert.Equal("/events-v1-2026.09.07/_mapping", invoker.Requests[^1].AbsolutePath); + } + + [Fact] + public async Task Dispose_WithMappingLoadInFlight_CancelsTransportRequest() + { + using var invoker = new MappingRequestInvoker("events-v1-2026.09.07") { BlockDiscovery = true }; + using var configuration = new MappingConfiguration(invoker); + using var index = new DailyIndex(configuration, "events"); + var lookup = index.MappingResolver.GetMappingAsync("dynamic", cancellationToken: TestContext.Current.CancellationToken).AsTask(); + await invoker.DiscoveryStarted.Task.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + + index.Dispose(); + try + { + await invoker.DiscoveryCancelled.Task.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + } + finally + { + invoker.ReleaseDiscovery.TrySetResult(); + } + + await lookup.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + Assert.Single(invoker.Requests); + } + + private static DailyIndex CreateIndex(ElasticConfiguration configuration, bool monthly, bool typed) => (monthly, typed) switch + { + (true, true) => new MonthlyIndex(configuration, "events"), + (true, false) => new MonthlyIndex(configuration, "events"), + (false, true) => new DailyIndex(configuration, "events"), + _ => new DailyIndex(configuration, "events") + }; + + private sealed class MappingDocument; + + private sealed class MappingConfiguration(MappingRequestInvoker invoker) : ElasticConfiguration + { + protected override ElasticsearchClient CreateElasticClient() => new(new ElasticsearchClientSettings( + new SingleNodePool(new Uri("http://localhost:9200")), invoker).DisablePing().MaximumRetries(0)); + } + + private sealed class MappingRequestInvoker(params string[] indexNames) + : InMemoryRequestInvoker([], 200, headers: new Dictionary> { { "x-elastic-product", ["Elasticsearch"] } }), IRequestInvoker + { + public string[] IndexNames { get; set; } = indexNames; + public int DiscoveryStatusCode { get; init; } = 200; + public int MappingStatusCode { get; init; } = 200; + public bool BlockDiscovery { get; init; } + public int SyncRequests { get; private set; } + public List Requests { get; } = []; + public TaskCompletionSource DiscoveryStarted { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously); + public TaskCompletionSource ReleaseDiscovery { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously); + public TaskCompletionSource DiscoveryCancelled { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously); + + public new TResponse Request(Endpoint endpoint, BoundConfiguration boundConfiguration, PostData? postData) + where TResponse : TransportResponse, new() + { + SyncRequests++; + Requests.Add(endpoint.Uri); + var (body, statusCode) = GetResponse(endpoint.Uri); + return BuildResponse(endpoint, boundConfiguration, postData, body, statusCode); + } + + public new async Task RequestAsync(Endpoint endpoint, BoundConfiguration boundConfiguration, + PostData? postData, CancellationToken cancellationToken = default) + where TResponse : TransportResponse, new() + { + Requests.Add(endpoint.Uri); + if (!endpoint.Uri.AbsolutePath.EndsWith("/_mapping", StringComparison.Ordinal)) + { + DiscoveryStarted.TrySetResult(); + if (BlockDiscovery) + { + try + { + await ReleaseDiscovery.Task.WaitAsync(cancellationToken); + } + catch (OperationCanceledException) + { + DiscoveryCancelled.TrySetResult(); + throw; + } + } + } + + var (body, statusCode) = GetResponse(endpoint.Uri); + return await BuildResponseAsync(endpoint, boundConfiguration, postData, cancellationToken, body, statusCode); + } + + private (byte[] Body, int StatusCode) GetResponse(Uri uri) + { + bool mapping = uri.AbsolutePath.EndsWith("/_mapping", StringComparison.Ordinal); + int statusCode = mapping ? MappingStatusCode : DiscoveryStatusCode; + if (statusCode != 200) + return (Encoding.UTF8.GetBytes($$"""{"error":{"type":"test_failure","reason":"test failure"},"status":{{statusCode}}}"""), statusCode); + + var response = new Dictionary(); + if (mapping) + { + string name = uri.AbsolutePath.Split('/')[1]; + response[name] = new { mappings = new { properties = new { dynamic = new { type = "keyword" } } } }; + } + else + { + foreach (string name in IndexNames) + response[name] = new { aliases = new { } }; + } + + return (System.Text.Json.JsonSerializer.SerializeToUtf8Bytes(response), statusCode); + } + } +} From 6e58c1a2083d1d1a881666594907bf4b5d1360cb Mon Sep 17 00:00:00 2001 From: "Eric J. Smith" Date: Mon, 7 Sep 2026 17:26:16 -0500 Subject: [PATCH 2/8] Await mapping resolution throughout structured search builders --- .../skills/foundatio-repositories/SKILL.md | 1 + .../references/index-lifecycle.md | 2 + docs/guide/index-management.md | 6 ++ .../Extensions/ResolverExtensions.cs | 75 +++++++++++-- ...oundatio.Repositories.Elasticsearch.csproj | 2 +- .../Queries/Builders/DateRangeQueryBuilder.cs | 9 +- .../Builders/DefaultSortQueryBuilder.cs | 26 +++-- .../Builders/FieldConditionsQueryBuilder.cs | 14 +-- .../Builders/FieldIncludesQueryBuilder.cs | 10 +- .../Builders/SearchAfterQueryBuilder.cs | 9 +- .../Queries/Builders/SortQueryBuilder.cs | 9 +- .../DailyIndexMappingTests.cs | 100 ++++++++++++++++++ 12 files changed, 212 insertions(+), 51 deletions(-) diff --git a/.agents/skills/foundatio-repositories/SKILL.md b/.agents/skills/foundatio-repositories/SKILL.md index 9c85c94f..c73c4f14 100644 --- a/.agents/skills/foundatio-repositories/SKILL.md +++ b/.agents/skills/foundatio-repositories/SKILL.md @@ -144,3 +144,4 @@ IReadOnlyRepository - **Patches do not fire `DocumentsSaving`/`DocumentsSaved`**: Patch operations only fire `DocumentsChanged`. - **Patches do not detect soft-delete transitions**: Even if a patch sets `IsDeleted = true`, the `ChangeType` is always `Saved`. Soft-delete detection requires `SaveAsync` with `OriginalsEnabled = true`. - **Large documents can make `ReindexAsync` trip Elasticsearch's indexing pressure limit**: The default reindex batch size (1000 docs) can produce a bulk sub-request bigger than a node's `indexing_pressure.memory.limit` (10% of heap), causing `es_rejected_execution_exception` ("rejected execution of coordinating operation"). Set `ReindexBatchSize`/`ReindexRequestsPerSecond` (must be positive and finite, or `ReindexAsync` throws `ArgumentOutOfRangeException`; a `null` work item throws `ArgumentNullException`) on the index to throttle. Task-status polling backs off exponentially with jitter (1s → 30s cap, +/-25%) on failure. A low `ReindexRequestsPerSecond` also extends the reindex's stall-detection timeout (default 10 minutes) so a healthy but slow, throttled reindex isn't cancelled as falsely "stalled". See [index-lifecycle.md](references/index-lifecycle.md#reindexasync). +- In asynchronous query builders, await mapping resolver APIs and `GetResolvedFieldsAsync` / `ResolveFieldNameAsync` / `ResolveFieldSortAsync`; synchronous helpers can block on a cold or missing field even with an async loader. diff --git a/.agents/skills/foundatio-repositories/references/index-lifecycle.md b/.agents/skills/foundatio-repositories/references/index-lifecycle.md index dac74e2e..0bdd0e09 100644 --- a/.agents/skills/foundatio-repositories/references/index-lifecycle.md +++ b/.agents/skills/foundatio-repositories/references/index-lifecycle.md @@ -487,3 +487,5 @@ public class EmployeeIndex : VersionedIndex | Error | `Failed to get the status {N} times in a row for reindex task ... reindexing {OldIndex} -> {NewIndex}` | Status polling gave up after `MAX_STATUS_FAILS` (10) consecutive failures; reindex progress can no longer be tracked, but the server-side `_reindex` task keeps running | DailyIndex never emits mapping errors from the built-in configuration path (since `ConfigureAsync` is a no-op). + +Structured sort, field-condition, include/exclude, date-range, and paging query builders await mapping resolution. Async custom builders should use `GetResolvedFieldsAsync`, `ResolveFieldNameAsync`, and `ResolveFieldSortAsync` to preserve boosts and sort settings. Protected GET/multi-GET request configuration hooks remain synchronous for compatibility. diff --git a/docs/guide/index-management.md b/docs/guide/index-management.md index c39d7bdf..1a29c426 100644 --- a/docs/guide/index-management.md +++ b/docs/guide/index-management.md @@ -870,6 +870,12 @@ POST /logs-v1-2025.05.*/_update_by_query?conflicts=proceed ### Mapping Resolver Cache (Query-Time Mapping Awareness) +Structured search builders for sorts, field conditions, includes/excludes, date ranges, and search-after +paging await mapping resolution. Custom async builders can use `GetResolvedFieldsAsync`, +`ResolveFieldNameAsync`, and `ResolveFieldSortAsync` to preserve field boosts and sort settings while +awaiting mapping I/O. Existing synchronous resolver helpers and protected GET/multi-GET request +configuration hooks retain their synchronous behavior. + The repository framework does **not** cache the PUT Mapping request/response (that's purely server-side). However, the **query parser** uses an `ElasticMappingResolver` that caches field-to-type resolution for building queries, sorting, and aggregations. This resolver combines two sources: 1. **Code mapping** — derived from your `ConfigureIndexMapping` method at startup (immutable for the process lifetime) diff --git a/src/Foundatio.Repositories.Elasticsearch/Extensions/ResolverExtensions.cs b/src/Foundatio.Repositories.Elasticsearch/Extensions/ResolverExtensions.cs index e0257121..cdb1958e 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Extensions/ResolverExtensions.cs +++ b/src/Foundatio.Repositories.Elasticsearch/Extensions/ResolverExtensions.cs @@ -1,8 +1,11 @@ using System; using System.Collections.Generic; using System.Linq; +using System.Threading; +using System.Threading.Tasks; using Elastic.Clients.Elasticsearch; using Foundatio.Parsers.ElasticQueries; +using Foundatio.Repositories.Extensions; namespace Foundatio.Repositories.Elasticsearch.Extensions; @@ -24,6 +27,55 @@ public static ICollection GetResolvedFields(this ElasticMappingReso return sorts.Select(sort => ResolveFieldSort(resolver, sort)).OfType().ToList(); } + /// Resolves field names asynchronously while preserving their boosts and input order. + public static async ValueTask> GetResolvedFieldsAsync(this ElasticMappingResolver resolver, + ICollection fields, CancellationToken cancellationToken = default) + { + if (fields.Count == 0) + return fields; + + var resolved = new List(fields.Count); + foreach (var field in fields) + resolved.Add(await resolver.ResolveFieldNameAsync(field, cancellationToken).AnyContext()); + return resolved; + } + + /// Resolves field sorts asynchronously while preserving sort settings and non-field sort variants. + public static async ValueTask> GetResolvedFieldsAsync(this ElasticMappingResolver resolver, + ICollection sorts, CancellationToken cancellationToken = default) + { + if (sorts.Count == 0) + return sorts; + + var resolved = new List(sorts.Count); + foreach (var sort in sorts) + { + var resolvedSort = await resolver.ResolveFieldSortAsync(sort, cancellationToken).AnyContext(); + if (resolvedSort is not null) + resolved.Add(resolvedSort); + } + return resolved; + } + + /// Resolves a field name asynchronously, retaining its boost. + public static async ValueTask ResolveFieldNameAsync(this ElasticMappingResolver resolver, Field field, + CancellationToken cancellationToken = default) + { + ArgumentNullException.ThrowIfNull(field); + return new Field(await resolver.GetResolvedFieldAsync(field, cancellationToken).AnyContext(), field.Boost); + } + + /// Resolves a field sort asynchronously, retaining its settings, or returns a non-field sort unchanged. + public static async ValueTask ResolveFieldSortAsync(this ElasticMappingResolver resolver, SortOptions? sort, + CancellationToken cancellationToken = default) + { + if (sort?.Field is not { } fieldSort) + return sort; + + var resolvedField = await resolver.GetSortFieldNameAsync(fieldSort.Field, cancellationToken).AnyContext(); + return CreateFieldSort(fieldSort, resolvedField); + } + public static Field ResolveFieldName(this ElasticMappingResolver resolver, Field field) { if (field is null) @@ -39,19 +91,20 @@ public static Field ResolveFieldName(this ElasticMappingResolver resolver, Field { var fieldSort = sort.Field; var resolvedField = resolver.GetSortFieldName(fieldSort.Field); - // Create a new FieldSort with the resolved field name - return new FieldSort - { - Field = resolvedField, - Missing = fieldSort.Missing, - Mode = fieldSort.Mode, - Nested = fieldSort.Nested, - NumericType = fieldSort.NumericType, - Order = fieldSort.Order, - UnmappedType = fieldSort.UnmappedType - }; + return CreateFieldSort(fieldSort, resolvedField); } return sort; } + + private static FieldSort CreateFieldSort(FieldSort fieldSort, string resolvedField) => new() + { + Field = resolvedField, + Missing = fieldSort.Missing, + Mode = fieldSort.Mode, + Nested = fieldSort.Nested, + NumericType = fieldSort.NumericType, + Order = fieldSort.Order, + UnmappedType = fieldSort.UnmappedType + }; } diff --git a/src/Foundatio.Repositories.Elasticsearch/Foundatio.Repositories.Elasticsearch.csproj b/src/Foundatio.Repositories.Elasticsearch/Foundatio.Repositories.Elasticsearch.csproj index 66be5a70..796f3512 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Foundatio.Repositories.Elasticsearch.csproj +++ b/src/Foundatio.Repositories.Elasticsearch/Foundatio.Repositories.Elasticsearch.csproj @@ -7,7 +7,7 @@ - + diff --git a/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/DateRangeQueryBuilder.cs b/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/DateRangeQueryBuilder.cs index b9e5d25b..1c212348 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/DateRangeQueryBuilder.cs +++ b/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/DateRangeQueryBuilder.cs @@ -4,6 +4,7 @@ using System.Linq; using System.Linq.Expressions; using System.Threading.Tasks; +using Foundatio.Repositories.Extensions; using Elastic.Clients.Elasticsearch; using Elastic.Clients.Elasticsearch.QueryDsl; using Exceptionless.DateTimeExtensions; @@ -99,17 +100,17 @@ namespace Foundatio.Repositories.Elasticsearch.Queries.Builders { public class DateRangeQueryBuilder : IElasticQueryBuilder { - public Task BuildAsync(QueryBuilderContext ctx) where T : class, new() + public async Task BuildAsync(QueryBuilderContext ctx) where T : class, new() { var dateRanges = ctx.Source.GetDateRanges(); if (dateRanges.Count <= 0) - return Task.CompletedTask; + return; var resolver = ctx.GetMappingResolver(); foreach (var dateRange in dateRanges.Where(dr => dr.UseDateRange)) { - var rangeQuery = new DateRangeQuery { Field = resolver.ResolveFieldName(dateRange.Field) }; + var rangeQuery = new DateRangeQuery { Field = await resolver.ResolveFieldNameAsync(dateRange.Field).AnyContext() }; if (dateRange.UseStartDate) rangeQuery.Gte = dateRange.GetStartDate(); if (dateRange.UseEndDate) @@ -119,8 +120,6 @@ public class DateRangeQueryBuilder : IElasticQueryBuilder ctx.Filter &= rangeQuery; } - - return Task.CompletedTask; } } } diff --git a/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/DefaultSortQueryBuilder.cs b/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/DefaultSortQueryBuilder.cs index 772485a6..2f6a3a36 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/DefaultSortQueryBuilder.cs +++ b/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/DefaultSortQueryBuilder.cs @@ -1,6 +1,6 @@ using System.Collections.Generic; -using System.Linq; using System.Threading.Tasks; +using Foundatio.Repositories.Extensions; using Elastic.Clients.Elasticsearch; using Foundatio.Parsers.ElasticQueries.Extensions; using Foundatio.Repositories.Models; @@ -11,7 +11,7 @@ public class DefaultSortQueryBuilder : IElasticQueryBuilder { private const string Id = nameof(IIdentity.Id); - public Task BuildAsync(QueryBuilderContext ctx) where T : class, new() + public async Task BuildAsync(QueryBuilderContext ctx) where T : class, new() { // Get existing sorts from context data (set by SortQueryBuilder or ExpressionQueryBuilder) List? sortFields = null; @@ -23,16 +23,22 @@ public class DefaultSortQueryBuilder : IElasticQueryBuilder sortFields ??= new List(); var resolver = ctx.GetMappingResolver(); - string idField = resolver.GetResolvedField(Id) ?? "_id"; + string idField = await resolver.GetResolvedFieldAsync(Id).AnyContext() ?? "_id"; // ensure id field is always present as a sort (default or tiebreaker) - bool hasIdField = sortFields.Any(s => + bool hasIdField = false; + foreach (var sort in sortFields) { - if (s?.Field?.Field == null) - return false; - string fieldName = resolver.GetSortFieldName(s.Field.Field); - return fieldName?.Equals(idField) == true; - }); + if (sort?.Field?.Field is not { } field) + continue; + + string fieldName = await resolver.GetSortFieldNameAsync(field).AnyContext(); + if (fieldName?.Equals(idField) == true) + { + hasIdField = true; + break; + } + } if (!hasIdField) { @@ -40,7 +46,5 @@ public class DefaultSortQueryBuilder : IElasticQueryBuilder } ctx.Data[SortQueryBuilder.SortFieldsKey] = sortFields; - - return Task.CompletedTask; } } diff --git a/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/FieldConditionsQueryBuilder.cs b/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/FieldConditionsQueryBuilder.cs index 720fb315..5f360cef 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/FieldConditionsQueryBuilder.cs +++ b/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/FieldConditionsQueryBuilder.cs @@ -173,7 +173,7 @@ or ComparisonOperator.LessThan if (nonAnalyzed || isStringRange) { - ValidateNonAnalyzedField(resolver, resolvedField, condition); + await ValidateNonAnalyzedFieldAsync(resolver, resolvedField, condition).AnyContext(); } switch (condition.Operator) @@ -214,7 +214,7 @@ or ComparisonOperator.LessThan } case ComparisonOperator.Contains: { - if (!resolver.IsPropertyAnalyzed(resolvedField)) + if (!await resolver.IsPropertyAnalyzedAsync(resolvedField).AnyContext()) { throw new QueryValidationException( $""" @@ -234,7 +234,7 @@ FieldContains generates a MatchQuery which requires an analyzed text field to to } case ComparisonOperator.NotContains: { - if (!resolver.IsPropertyAnalyzed(resolvedField)) + if (!await resolver.IsPropertyAnalyzedAsync(resolvedField).AnyContext()) { throw new QueryValidationException( $""" @@ -267,9 +267,9 @@ FieldNotContains generates a MatchQuery which requires an analyzed text field to } } - private static void ValidateNonAnalyzedField(ElasticMappingResolver resolver, string resolvedField, FieldCondition condition) + private static async Task ValidateNonAnalyzedFieldAsync(ElasticMappingResolver resolver, string resolvedField, FieldCondition condition) { - if (resolver.IsPropertyAnalyzed(resolvedField)) + if (await resolver.IsPropertyAnalyzedAsync(resolvedField).AnyContext()) { bool isRange = condition.Operator is ComparisonOperator.GreaterThan or ComparisonOperator.GreaterThanOrEqual @@ -421,7 +421,7 @@ private static TermRangeQuery BuildTermRange(string field, ComparisonOperator op private static async Task ResolveFieldAsync(QueryBuilderContext ctx, ElasticMappingResolver resolver, Field field, bool nonAnalyzed) where T : class, new() { - string resolved = resolver.GetResolvedField(field); + string resolved = await resolver.GetResolvedFieldAsync(field).AnyContext(); if (ctx is IQueryVisitorContextWithFieldResolver { FieldResolver: not null } fieldResolverCtx) { @@ -432,7 +432,7 @@ private static TermRangeQuery BuildTermRange(string field, ComparisonOperator op if (nonAnalyzed) { - string? nonAnalyzedField = resolver.GetNonAnalyzedFieldName(resolved); + string? nonAnalyzedField = await resolver.GetNonAnalyzedFieldNameAsync(resolved).AnyContext(); if (nonAnalyzedField is not null) resolved = nonAnalyzedField; } diff --git a/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/FieldIncludesQueryBuilder.cs b/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/FieldIncludesQueryBuilder.cs index 02e0461b..03cbdc80 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/FieldIncludesQueryBuilder.cs +++ b/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/FieldIncludesQueryBuilder.cs @@ -361,7 +361,7 @@ namespace Foundatio.Repositories.Elasticsearch.Queries.Builders /// public class FieldIncludesQueryBuilder : IElasticQueryBuilder { - public Task BuildAsync(QueryBuilderContext ctx) where T : class, new() + public async Task BuildAsync(QueryBuilderContext ctx) where T : class, new() { var resolver = ctx.GetMappingResolver(); @@ -394,14 +394,14 @@ public class FieldIncludesQueryBuilder : IElasticQueryBuilder if (requiredFields.Count > 0 && includes.Count > 0) includes.AddRange(requiredFields); - var resolvedIncludes = resolver.GetResolvedFields(includes).ToArray(); - var resolvedExcludes = resolver.GetResolvedFields(excludes) + var resolvedIncludes = (await resolver.GetResolvedFieldsAsync(includes).AnyContext()).ToArray(); + var resolvedExcludes = (await resolver.GetResolvedFieldsAsync(excludes).AnyContext()) .Where(f => !resolvedIncludes.Contains(f)) .ToArray(); if (requiredFields.Count > 0 && resolvedIncludes.Length is 0) { - var resolvedRequiredFields = resolver.GetResolvedFields(requiredFields); + var resolvedRequiredFields = await resolver.GetResolvedFieldsAsync(requiredFields).AnyContext(); resolvedExcludes = resolvedExcludes.Where(f => !resolvedRequiredFields.Contains(f)).ToArray(); } @@ -415,8 +415,6 @@ public class FieldIncludesQueryBuilder : IElasticQueryBuilder ctx.Search.Source(new SourceConfig(filter)); } - - return Task.CompletedTask; } } } diff --git a/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/SearchAfterQueryBuilder.cs b/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/SearchAfterQueryBuilder.cs index caa60760..395c33a6 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/SearchAfterQueryBuilder.cs +++ b/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/SearchAfterQueryBuilder.cs @@ -2,6 +2,7 @@ using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; +using Foundatio.Repositories.Extensions; using Elastic.Clients.Elasticsearch; using Foundatio.Parsers.ElasticQueries.Extensions; using Foundatio.Repositories.Elasticsearch.Extensions; @@ -207,7 +208,7 @@ public class SearchAfterQueryBuilder : IElasticQueryBuilder "_score" }; - public Task BuildAsync(QueryBuilderContext ctx) where T : class, new() + public async Task BuildAsync(QueryBuilderContext ctx) where T : class, new() { // Get sorts from context data (set by SortQueryBuilder or ExpressionQueryBuilder) List? sortFields = null; @@ -222,7 +223,7 @@ public class SearchAfterQueryBuilder : IElasticQueryBuilder sortFields ??= new List(); var resolver = ctx.GetMappingResolver(); - string idField = resolver.GetResolvedField(Id) ?? "_id"; + string idField = await resolver.GetResolvedFieldAsync(Id).AnyContext() ?? "_id"; // Live search_after paging with an unstable sort key (e.g. _doc, _score) is only safe // within a Point-In-Time: index refreshes and segment merges can invalidate the cursor, @@ -248,7 +249,7 @@ public class SearchAfterQueryBuilder : IElasticQueryBuilder if (sort?.Field?.Field is { } sortField) { - fieldName = resolver.GetSortFieldName(sortField); + fieldName = await resolver.GetSortFieldNameAsync(sortField).AnyContext(); if (fieldName is null) continue; @@ -296,8 +297,6 @@ public class SearchAfterQueryBuilder : IElasticQueryBuilder { ctx.Search.Sort(sortFields); } - - return Task.CompletedTask; } } } diff --git a/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/SortQueryBuilder.cs b/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/SortQueryBuilder.cs index 1e1bc4a8..2d12f3d0 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/SortQueryBuilder.cs +++ b/src/Foundatio.Repositories.Elasticsearch/Queries/Builders/SortQueryBuilder.cs @@ -3,6 +3,7 @@ using System.Linq; using System.Linq.Expressions; using System.Threading.Tasks; +using Foundatio.Repositories.Extensions; using Elastic.Clients.Elasticsearch; using Foundatio.Parsers.ElasticQueries.Extensions; using Foundatio.Repositories.Elasticsearch.Extensions; @@ -63,21 +64,19 @@ public class SortQueryBuilder : IElasticQueryBuilder { internal const string SortFieldsKey = "__SortFields"; - public Task BuildAsync(QueryBuilderContext ctx) where T : class, new() + public async Task BuildAsync(QueryBuilderContext ctx) where T : class, new() { var sortFields = ctx.Source.GetSorts().ToList(); if (sortFields.Count <= 0) - return Task.CompletedTask; + return; var resolver = ctx.GetMappingResolver(); - sortFields = resolver.GetResolvedFields(sortFields).ToList(); + sortFields = (await resolver.GetResolvedFieldsAsync(sortFields).AnyContext()).ToList(); // Store sorts in context data - SearchAfterQueryBuilder will apply them // along with any sorts from ExpressionQueryBuilder ctx.Data[SortFieldsKey] = sortFields; - - return Task.CompletedTask; } } } diff --git a/tests/Foundatio.Repositories.Elasticsearch.Tests/DailyIndexMappingTests.cs b/tests/Foundatio.Repositories.Elasticsearch.Tests/DailyIndexMappingTests.cs index d35a6da3..67f5d871 100644 --- a/tests/Foundatio.Repositories.Elasticsearch.Tests/DailyIndexMappingTests.cs +++ b/tests/Foundatio.Repositories.Elasticsearch.Tests/DailyIndexMappingTests.cs @@ -6,7 +6,11 @@ using Elastic.Clients.Elasticsearch; using Elastic.Clients.Elasticsearch.Mapping; using Elastic.Transport; +using Foundatio.Parsers.ElasticQueries; using Foundatio.Repositories.Elasticsearch.Configuration; +using Foundatio.Repositories.Elasticsearch.Extensions; +using Foundatio.Repositories.Elasticsearch.Queries.Builders; +using Foundatio.Repositories.Options; using Xunit; using FieldMapping = Foundatio.Parsers.ElasticQueries.FieldMapping; @@ -36,6 +40,102 @@ public async Task GetMappingAsync_WithTimeSeriesIndex_UsesAsyncMetadataRequests( Assert.Equal($"/{latest}/_mapping", invoker.Requests[1].AbsolutePath); } + [Theory] + [InlineData("sort")] + [InlineData("condition")] + [InlineData("includes")] + [InlineData("date-range")] + [InlineData("default-sort")] + [InlineData("search-after")] + public async Task BuildAsync_WithTimeSeriesMapping_ReturnsWithoutBlockingCaller(string operation) + { + using var invoker = new MappingRequestInvoker("events-v1-2026.09.07") { BlockDiscovery = true }; + using var configuration = new MappingConfiguration(invoker); + using var index = new DailyIndex(configuration, "events"); + var query = new RepositoryQuery(); + var options = new CommandOptions().ElasticIndex(index); + IElasticQueryBuilder builder; + switch (operation) + { + case "sort": + query.Sort("dynamic"); + builder = new SortQueryBuilder(); + break; + case "condition": + query.FieldEquals("dynamic", "value"); + builder = new FieldConditionsQueryBuilder(); + break; + case "includes": + query.Include("dynamic"); + builder = new FieldIncludesQueryBuilder(); + break; + case "date-range": + query.DateRange(DateTime.UtcNow.AddDays(-1), DateTime.UtcNow, "dynamic"); + builder = new DateRangeQueryBuilder(); + break; + case "default-sort": + builder = new DefaultSortQueryBuilder(); + break; + default: + options.SearchAfterPaging(); + builder = new SearchAfterQueryBuilder(); + break; + } + + var context = new QueryBuilderContext(query, options); + var returned = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var invocation = Task.Run(async () => + { + var build = builder.BuildAsync(context); + returned.TrySetResult(); + await build; + }, TestContext.Current.CancellationToken); + + try + { + await invoker.DiscoveryStarted.Task.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + await returned.Task.WaitAsync(TimeSpan.FromSeconds(1), TestContext.Current.CancellationToken); + Assert.False(invocation.IsCompleted); + Assert.Equal(0, invoker.SyncRequests); + } + finally + { + invoker.ReleaseDiscovery.TrySetResult(); + await invocation.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + } + + Assert.Equal(2, invoker.Requests.Count); + } + + [Fact] + public async Task GetResolvedFieldsAsync_WithBoostsAndSortVariants_PreservesRequestSettings() + { + using var settings = new ElasticsearchClientSettings(new Uri("http://localhost:9200")); + using var resolver = ElasticMappingResolver.CreateWithAsyncLoader(_ => Task.FromResult(new TypeMapping + { + Properties = new Properties + { + { "name", new TextProperty { Fields = new Properties { { "keyword", new KeywordProperty() } } } } + } + }), new Inferrer(settings)); + var score = new SortOptions { Score = new ScoreSort { Order = SortOrder.Desc } }; + SortOptions field = new FieldSort { Field = "name", Order = SortOrder.Desc, Mode = SortMode.Max }; + + var fields = await resolver.GetResolvedFieldsAsync(new List { new("name", 2) }, TestContext.Current.CancellationToken); + var sorts = await resolver.GetResolvedFieldsAsync(new List { field, score, null! }, TestContext.Current.CancellationToken); + + var resolvedField = Assert.Single(fields); + Assert.Equal("name", resolvedField.Name); + Assert.Equal(2, resolvedField.Boost); + Assert.Equal(2, sorts.Count); + Assert.Contains(score, sorts); + var resolvedSort = Assert.Single(sorts, sort => sort.Field is not null).Field!; + Assert.Equal("name.keyword", resolvedSort.Field.Name); + Assert.Equal(SortOrder.Desc, resolvedSort.Order); + Assert.Equal(SortMode.Max, resolvedSort.Mode); + Assert.Equal("name", field.Field!.Field.Name); + } + [Fact] public async Task GetMappingAsync_WhenDiscoveryIsBlocked_ReturnsWithoutBlockingCaller() { From 70f9b4196ba2a423e409df760b0c93c412e7c1ec Mon Sep 17 00:00:00 2001 From: "Eric J. Smith" Date: Mon, 7 Sep 2026 17:38:21 -0500 Subject: [PATCH 3/8] Dispose mapping resolvers created during index shutdown --- .../Configuration/Index.cs | 15 ++++-- .../DailyIndexMappingTests.cs | 46 +++++++++++++++++++ 2 files changed, 58 insertions(+), 3 deletions(-) diff --git a/src/Foundatio.Repositories.Elasticsearch/Configuration/Index.cs b/src/Foundatio.Repositories.Elasticsearch/Configuration/Index.cs index ed43b10c..53d04a15 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Configuration/Index.cs +++ b/src/Foundatio.Repositories.Elasticsearch/Configuration/Index.cs @@ -32,6 +32,7 @@ public class Index : IIndex, IHaveLogger private readonly Lazy _queryBuilder; private readonly Lazy _queryParser; private readonly Lazy _mappingResolver; + private ElasticMappingResolver? _ownedMappingResolver; private readonly Lazy _fieldResolver; private readonly ConcurrentDictionary _customFieldTypes = new(); private readonly AsyncLock _lock = new(); @@ -47,7 +48,7 @@ public Index(IElasticConfiguration configuration, string name) Configuration = configuration; _queryBuilder = new Lazy(CreateQueryBuilder); _queryParser = new Lazy(CreateQueryParser); - _mappingResolver = new Lazy(CreateMappingResolver); + _mappingResolver = new Lazy(CreateOwnedMappingResolver); _fieldResolver = new Lazy(CreateQueryFieldResolver); _logger = configuration.LoggerFactory?.CreateLogger(GetType()) ?? NullLogger.Instance; } @@ -89,6 +90,15 @@ protected virtual IElasticQueryBuilder CreateQueryBuilder() protected virtual void ConfigureQueryBuilder(ElasticQueryBuilder builder) { } + private ElasticMappingResolver CreateOwnedMappingResolver() + { + var resolver = CreateMappingResolver(); + Interlocked.Exchange(ref _ownedMappingResolver, resolver); + if (Volatile.Read(ref _disposed) != 0) + Interlocked.Exchange(ref _ownedMappingResolver, null)?.Dispose(); + return resolver; + } + protected virtual ElasticMappingResolver CreateMappingResolver() { return ElasticMappingResolver.Create(Configuration.Client, Name, _logger); @@ -415,8 +425,7 @@ public virtual void Dispose() _disposedCancellationTokenSource.Cancel(); _disposedCancellationTokenSource.Dispose(); - if (_mappingResolver.IsValueCreated) - _mappingResolver.Value.Dispose(); + Interlocked.Exchange(ref _ownedMappingResolver, null)?.Dispose(); } } diff --git a/tests/Foundatio.Repositories.Elasticsearch.Tests/DailyIndexMappingTests.cs b/tests/Foundatio.Repositories.Elasticsearch.Tests/DailyIndexMappingTests.cs index 67f5d871..e85e42f0 100644 --- a/tests/Foundatio.Repositories.Elasticsearch.Tests/DailyIndexMappingTests.cs +++ b/tests/Foundatio.Repositories.Elasticsearch.Tests/DailyIndexMappingTests.cs @@ -279,6 +279,41 @@ public async Task Dispose_WithMappingLoadInFlight_CancelsTransportRequest() Assert.Single(invoker.Requests); } + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task Dispose_WithResolverNotYetPublished_PreventsLaterMappingLoads(bool creationStarted) + { + using var invoker = new MappingRequestInvoker("events-v1-2026.09.07"); + using var configuration = new MappingConfiguration(invoker); + using var started = new ManualResetEventSlim(); + using var release = new ManualResetEventSlim(); + using var index = new BlockingMappingIndex(configuration, started, release); + Task? creation = null; + try + { + if (creationStarted) + { + creation = Task.Factory.StartNew(() => index.MappingResolver, TestContext.Current.CancellationToken, + TaskCreationOptions.LongRunning, TaskScheduler.Default); + Assert.True(started.Wait(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken)); + } + + index.Dispose(); + Assert.Equal(creationStarted, started.IsSet); + } + finally + { + release.Set(); + if (creation is not null) + await creation.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + } + + var mapping = await index.MappingResolver.GetMappingAsync("dynamic", cancellationToken: TestContext.Current.CancellationToken); + Assert.False(mapping?.Found); + Assert.Empty(invoker.Requests); + } + private static DailyIndex CreateIndex(ElasticConfiguration configuration, bool monthly, bool typed) => (monthly, typed) switch { (true, true) => new MonthlyIndex(configuration, "events"), @@ -287,6 +322,17 @@ public async Task Dispose_WithMappingLoadInFlight_CancelsTransportRequest() _ => new DailyIndex(configuration, "events") }; + private sealed class BlockingMappingIndex(IElasticConfiguration configuration, ManualResetEventSlim started, + ManualResetEventSlim release) : DailyIndex(configuration, "events") + { + protected override ElasticMappingResolver CreateMappingResolver() + { + started.Set(); + Assert.True(release.Wait(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken)); + return base.CreateMappingResolver(); + } + } + private sealed class MappingDocument; private sealed class MappingConfiguration(MappingRequestInvoker invoker) : ElasticConfiguration From 8e1817733f0848bf4057c488c6c37075a2d6bd4b Mon Sep 17 00:00:00 2001 From: Blake Niemyjski Date: Tue, 22 Sep 2026 09:45:55 -0500 Subject: [PATCH 4/8] Fix build compatibility after dependency upgrades Reference System.IO.Hashing directly where ElasticConfiguration uses it instead of relying on a transitive dependency. Migrate both test assemblies to xUnit's Parallelization attribute while retaining serialized test execution. --- .../Foundatio.Repositories.Elasticsearch.csproj | 1 + .../Properties/AssemblyInfo.cs | 2 +- tests/Foundatio.Repositories.Tests/Properties/AssemblyInfo.cs | 2 +- 3 files changed, 3 insertions(+), 2 deletions(-) diff --git a/src/Foundatio.Repositories.Elasticsearch/Foundatio.Repositories.Elasticsearch.csproj b/src/Foundatio.Repositories.Elasticsearch/Foundatio.Repositories.Elasticsearch.csproj index 796f3512..2da89d39 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Foundatio.Repositories.Elasticsearch.csproj +++ b/src/Foundatio.Repositories.Elasticsearch/Foundatio.Repositories.Elasticsearch.csproj @@ -7,6 +7,7 @@ + diff --git a/tests/Foundatio.Repositories.Elasticsearch.Tests/Properties/AssemblyInfo.cs b/tests/Foundatio.Repositories.Elasticsearch.Tests/Properties/AssemblyInfo.cs index 31425692..ea40c673 100644 --- a/tests/Foundatio.Repositories.Elasticsearch.Tests/Properties/AssemblyInfo.cs +++ b/tests/Foundatio.Repositories.Elasticsearch.Tests/Properties/AssemblyInfo.cs @@ -1 +1 @@ -[assembly: Xunit.CollectionBehaviorAttribute(DisableTestParallelization = true, MaxParallelThreads = 1)] +[assembly: Xunit.v3.Parallelization(Mode = Xunit.Sdk.ParallelMode.None, MaxThreads = 1)] diff --git a/tests/Foundatio.Repositories.Tests/Properties/AssemblyInfo.cs b/tests/Foundatio.Repositories.Tests/Properties/AssemblyInfo.cs index 31425692..ea40c673 100644 --- a/tests/Foundatio.Repositories.Tests/Properties/AssemblyInfo.cs +++ b/tests/Foundatio.Repositories.Tests/Properties/AssemblyInfo.cs @@ -1 +1 @@ -[assembly: Xunit.CollectionBehaviorAttribute(DisableTestParallelization = true, MaxParallelThreads = 1)] +[assembly: Xunit.v3.Parallelization(Mode = Xunit.Sdk.ParallelMode.None, MaxThreads = 1)] From f453ab7bea7282a9bd622ea2d3b094400bc82016 Mon Sep 17 00:00:00 2001 From: Blake Niemyjski Date: Tue, 22 Sep 2026 09:48:29 -0500 Subject: [PATCH 5/8] Test preservation of field-sort settings in sync and async resolution --- .../ResolverExtensionsTests.cs | 65 +++++++++++++++++++ 1 file changed, 65 insertions(+) create mode 100644 tests/Foundatio.Repositories.Elasticsearch.Tests/ResolverExtensionsTests.cs diff --git a/tests/Foundatio.Repositories.Elasticsearch.Tests/ResolverExtensionsTests.cs b/tests/Foundatio.Repositories.Elasticsearch.Tests/ResolverExtensionsTests.cs new file mode 100644 index 00000000..8eb4a2d2 --- /dev/null +++ b/tests/Foundatio.Repositories.Elasticsearch.Tests/ResolverExtensionsTests.cs @@ -0,0 +1,65 @@ +using System; +using System.Collections.Generic; +using System.Threading.Tasks; +using Elastic.Clients.Elasticsearch; +using Elastic.Clients.Elasticsearch.Mapping; +using Foundatio.Parsers.ElasticQueries; +using Foundatio.Repositories.Elasticsearch.Extensions; +using Xunit; + +namespace Foundatio.Repositories.Elasticsearch.Tests; + +public class ResolverExtensionsTests +{ + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task GetResolvedFields_WithFieldSort_PreservesAllSettings(bool asynchronous) + { + using var settings = new ElasticsearchClientSettings(new Uri("http://localhost:9200")); + using var resolver = ElasticMappingResolver.CreateWithAsyncLoader(_ => Task.FromResult(new TypeMapping + { + Properties = new Properties + { + { "events", new NestedProperty { Properties = new Properties { { "occurredAt", new DateProperty() } } } } + } + }), new Inferrer(settings)); + var field = new FieldSort + { + Field = "events.occurredAt", + Format = "strict_date_optional_time_nanos", + Missing = "_last", + Mode = SortMode.Max, + Nested = new NestedSortValue { Path = "events" }, + NumericType = FieldSortNumericType.DateNanos, + Order = SortOrder.Desc, + UnmappedType = FieldType.Date + }; + var score = new SortOptions { Score = new ScoreSort { Order = SortOrder.Desc } }; + var input = new List { field, score, null! }; + + var resolved = asynchronous + ? await resolver.GetResolvedFieldsAsync(input, TestContext.Current.CancellationToken) + : resolver.GetResolvedFields(input); + + Assert.Collection(resolved, + sort => + { + var actual = Assert.IsType(sort.Field); + Assert.NotSame(field, actual); + Assert.Equal("events.occurredAt", actual.Field.Name); + Assert.Equal(field.Format, actual.Format); + Assert.Equal(field.Missing, actual.Missing); + Assert.Equal(field.Mode, actual.Mode); + Assert.Same(field.Nested, actual.Nested); + Assert.Equal(field.NumericType, actual.NumericType); + Assert.Equal(field.Order, actual.Order); + Assert.Equal(field.UnmappedType, actual.UnmappedType); + }, + sort => Assert.Same(score, sort)); + Assert.Equal(3, input.Count); + Assert.Same(field, input[0].Field); + Assert.Equal("events.occurredAt", field.Field.Name); + Assert.Equal("strict_date_optional_time_nanos", field.Format); + } +} From 69951977e342f0a3c9a36bf65bc97c9372f25549 Mon Sep 17 00:00:00 2001 From: Blake Niemyjski Date: Tue, 22 Sep 2026 09:58:59 -0500 Subject: [PATCH 6/8] Preserve field-sort formats during synchronous and asynchronous resolution Both regression cases fail on the parent commit because CreateFieldSort drops Format. Copy it alongside the other settings so date sort values retain the caller's requested representation, including for search-after cursors. --- .../Extensions/ResolverExtensions.cs | 15 +++++++-------- 1 file changed, 7 insertions(+), 8 deletions(-) diff --git a/src/Foundatio.Repositories.Elasticsearch/Extensions/ResolverExtensions.cs b/src/Foundatio.Repositories.Elasticsearch/Extensions/ResolverExtensions.cs index cdb1958e..62868f47 100644 --- a/src/Foundatio.Repositories.Elasticsearch/Extensions/ResolverExtensions.cs +++ b/src/Foundatio.Repositories.Elasticsearch/Extensions/ResolverExtensions.cs @@ -28,8 +28,7 @@ public static ICollection GetResolvedFields(this ElasticMappingReso } /// Resolves field names asynchronously while preserving their boosts and input order. - public static async ValueTask> GetResolvedFieldsAsync(this ElasticMappingResolver resolver, - ICollection fields, CancellationToken cancellationToken = default) + public static async ValueTask> GetResolvedFieldsAsync(this ElasticMappingResolver resolver, ICollection fields, CancellationToken cancellationToken = default) { if (fields.Count == 0) return fields; @@ -37,12 +36,12 @@ public static async ValueTask> GetResolvedFieldsAsync(this El var resolved = new List(fields.Count); foreach (var field in fields) resolved.Add(await resolver.ResolveFieldNameAsync(field, cancellationToken).AnyContext()); + return resolved; } /// Resolves field sorts asynchronously while preserving sort settings and non-field sort variants. - public static async ValueTask> GetResolvedFieldsAsync(this ElasticMappingResolver resolver, - ICollection sorts, CancellationToken cancellationToken = default) + public static async ValueTask> GetResolvedFieldsAsync(this ElasticMappingResolver resolver, ICollection sorts, CancellationToken cancellationToken = default) { if (sorts.Count == 0) return sorts; @@ -54,20 +53,19 @@ public static async ValueTask> GetResolvedFieldsAsync(t if (resolvedSort is not null) resolved.Add(resolvedSort); } + return resolved; } /// Resolves a field name asynchronously, retaining its boost. - public static async ValueTask ResolveFieldNameAsync(this ElasticMappingResolver resolver, Field field, - CancellationToken cancellationToken = default) + public static async ValueTask ResolveFieldNameAsync(this ElasticMappingResolver resolver, Field field, CancellationToken cancellationToken = default) { ArgumentNullException.ThrowIfNull(field); return new Field(await resolver.GetResolvedFieldAsync(field, cancellationToken).AnyContext(), field.Boost); } /// Resolves a field sort asynchronously, retaining its settings, or returns a non-field sort unchanged. - public static async ValueTask ResolveFieldSortAsync(this ElasticMappingResolver resolver, SortOptions? sort, - CancellationToken cancellationToken = default) + public static async ValueTask ResolveFieldSortAsync(this ElasticMappingResolver resolver, SortOptions? sort, CancellationToken cancellationToken = default) { if (sort?.Field is not { } fieldSort) return sort; @@ -100,6 +98,7 @@ public static Field ResolveFieldName(this ElasticMappingResolver resolver, Field private static FieldSort CreateFieldSort(FieldSort fieldSort, string resolvedField) => new() { Field = resolvedField, + Format = fieldSort.Format, Missing = fieldSort.Missing, Mode = fieldSort.Mode, Nested = fieldSort.Nested, From 839a095596ae9293c260a5690089d90a651c63d8 Mon Sep 17 00:00:00 2001 From: Blake Niemyjski Date: Tue, 22 Sep 2026 10:00:53 -0500 Subject: [PATCH 7/8] Isolate integration tests from Kibana and fail fast on unhealthy Elasticsearch The suite asserts node-wide scroll counts, so do not start unrelated Kibana background work on its test cluster. Wait for the Elasticsearch health check and propagate startup failures rather than inheriting the masked compose exit. --- .github/workflows/build.yml | 2 ++ 1 file changed, 2 insertions(+) diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index ea5e0ff6..4582485f 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -8,3 +8,5 @@ jobs: with: new-test-runner: true minimum-major-minor: '8.0' + # Node-wide scroll assertions require a cluster without Kibana background work. + compose-command: docker compose up --detach --wait --wait-timeout 120 elasticsearch From 6f74002bf29df9a6ac3f929c78208acecb519c83 Mon Sep 17 00:00:00 2001 From: Blake Niemyjski Date: Tue, 22 Sep 2026 10:30:47 -0500 Subject: [PATCH 8/8] Correct mapping refresh, explicit configuration, and backfill guidance Document snapshot reload/cancellation behavior consistently across the guide and agent references. Replace the nonexistent force parameter with explicit index selection and explain queue-independent configuration. Remove the incorrect claim that a no-op script backfills newly mapped fields. --- .../references/index-lifecycle.md | 24 ++++++++++++------- .../references/patterns.md | 5 ++-- docs/guide/index-management.md | 20 ++++++++++++---- 3 files changed, 34 insertions(+), 15 deletions(-) diff --git a/.agents/skills/foundatio-repositories/references/index-lifecycle.md b/.agents/skills/foundatio-repositories/references/index-lifecycle.md index 0bdd0e09..9a87ae37 100644 --- a/.agents/skills/foundatio-repositories/references/index-lifecycle.md +++ b/.agents/skills/foundatio-repositories/references/index-lifecycle.md @@ -340,7 +340,7 @@ var results = await repository.FindAsync(q => q.Index("logs-last-7-days")); | `ElasticMappingResolver` field cache | Snapshot lifetime; unresolved fields trigger reloads with a five-second cooldown | `index.MappingResolver.RefreshMapping()` | | `_isEnsured` flag (Index/VersionedIndex) | Process lifetime | App restart or index deletion | | `_ensuredDates` (DailyIndex) | Process lifetime per-date | `DeleteAsync(name)` or `Dispose()` | -| `ConfigureIndexesAsync` cache marker | 5 minutes (distributed) | Expires automatically; or `ConfigureIndexesAsync(force: true)` | +| `ConfigureIndexesAsync` cache marker | 5 minutes (distributed) | Expires automatically; pass explicit indexes to bypass the configuration-level lock and cache | Daily/monthly resolvers asynchronously discover the newest partition using names and aliases only, then load that partition's mapping. Concurrent lookups share one load per resolver. Index disposal cancels active mapping @@ -349,6 +349,12 @@ The cooldown is not a freshness guarantee, and successful lookups do not trigger `RefreshMapping()` after known mapping changes, including changes to an already-resolved alias. Mappings from historical partitions are not merged into the newest partition's mapping. +A caller's cancellation token cancels its wait without canceling a shared load needed by other callers. +Disposing the resolver cancels the shared mapping I/O. Failed or empty reloads retain the last usable snapshot; +explicit `RefreshMapping()` invalidates that snapshot so the next lookup loads again. + +Structured sort, field-condition, include/exclude, date-range, and paging query builders await mapping resolution. Async custom builders should use `GetResolvedFieldsAsync`, `ResolveFieldNameAsync`, and `ResolveFieldSortAsync` to preserve boosts and sort settings. Protected GET/multi-GET request configuration hooks remain synchronous for compatibility. + ## Index Operations ### ConfigureIndexesAsync @@ -358,13 +364,17 @@ Creates indexes and updates mappings. Protected by distributed lock + cache mark ```csharp await configuration.ConfigureIndexesAsync(); -// Bypass cache marker (after structural changes) -await configuration.ConfigureIndexesAsync(force: true); +// Explicit indexes bypass the configuration-level lock and cache marker. +// Disable enqueueing when a queue worker is not configured; reindex separately. +await configuration.ConfigureIndexesAsync(configuration.Indexes, beginReindexingOutdated: false); -// Configure specific indexes (bypasses lock and cache) -await configuration.ConfigureIndexesAsync([myIndex]); +// Configure specific indexes using the same explicit-selection path. +await configuration.ConfigureIndexesAsync([myIndex], beginReindexingOutdated: false); ``` +There is no `force` parameter. Explicit selection does not update existing daily/monthly partition mappings; +those indexes still follow the manual mapping lifecycle described above. + ### MaintainIndexesAsync Updates aliases for time-series indexes, deletes expired indexes: @@ -483,9 +493,7 @@ public class EmployeeIndex : VersionedIndex | Error | `Error updating index ({name}) mappings.` | PUT Mapping failed | | Error | `Error updating index ({name}) mappings. Changing existing fields requires a new index version.` | Tried to change existing field type on VersionedIndex | | Warning | `Adding new analyzer/tokenizer/filter to existing index (requires close/reopen)` | New analysis component needs index close/reopen | -| Error | `Error getting task status while reindexing: {OldIndex} -> {NewIndex}` | Task status poll failed (e.g. `es_rejected_execution_exception` from indexing pressure); retried with exponential backoff (1s, doubling, capped at 30s) | +| Error | `Error getting task status while reindexing: {OldIndex} -> {NewIndex}` | Task status poll failed (e.g. `es_rejected_execution_exception` from indexing pressure); retried with exponential backoff (1s, doubling, capped to 30s) | | Error | `Failed to get the status {N} times in a row for reindex task ... reindexing {OldIndex} -> {NewIndex}` | Status polling gave up after `MAX_STATUS_FAILS` (10) consecutive failures; reindex progress can no longer be tracked, but the server-side `_reindex` task keeps running | DailyIndex never emits mapping errors from the built-in configuration path (since `ConfigureAsync` is a no-op). - -Structured sort, field-condition, include/exclude, date-range, and paging query builders await mapping resolution. Async custom builders should use `GetResolvedFieldsAsync`, `ResolveFieldNameAsync`, and `ResolveFieldSortAsync` to preserve boosts and sort settings. Protected GET/multi-GET request configuration hooks remain synchronous for compatibility. diff --git a/.agents/skills/foundatio-repositories/references/patterns.md b/.agents/skills/foundatio-repositories/references/patterns.md index 0f251333..abd71964 100644 --- a/.agents/skills/foundatio-repositories/references/patterns.md +++ b/.agents/skills/foundatio-repositories/references/patterns.md @@ -290,14 +290,15 @@ All index configurations use `.Dynamic(false)`, which disables Elasticsearch's a After adding a mapping for a previously unmapped field, only **newly saved/indexed documents** will be searchable on that field. To make existing documents searchable: - **Foundatio migration** -- create a `MigrationBase` subclass that uses `PatchAllAsync` or `BatchProcessAsync` to touch all affected documents (recommended for production). -- **`PatchAllAsync`** with a no-op `ScriptPatch` (e.g., `ctx.op = 'none'` -- still triggers re-index of `_source`). - **Elasticsearch Update By Query API** with no script: `POST /{index}/_update_by_query` re-indexes every document in place. +A script that explicitly skips a write does **not** backfill a new mapping. Do not use `ctx.op = 'none'` for this purpose; update-by-query uses `ctx.op = 'noop'` to skip a document. Leave the indexing operation enabled when re-indexing `_source`. + For `DailyIndex`/`MonthlyIndex`, you must also apply the mapping to existing physical indexes before the update-by-query will help. **Trade-off for Daily/Monthly indexes**: Rolling forward (doing nothing to old partitions and waiting for retention to cycle out old data) is often the cheapest strategy. -**Mapping resolver cache**: After applying a manual PUT Mapping, the in-process `ElasticMappingResolver` auto-refreshes from the server within ~60 seconds. To force immediate recognition, call `index.MappingResolver.RefreshMapping()`. +**Mapping resolver cache**: The in-process `ElasticMappingResolver` retains a snapshot; it does not refresh on a 60-second timer. Unresolved fields can trigger reloads with a five-second cooldown, which is not a freshness guarantee. After Elasticsearch acknowledges a known mapping change, call `index.MappingResolver.RefreshMapping()` so the next lookup reloads the mapping. This is also required for changes to an already-resolved alias. Resolver invalidation does not backfill existing documents. See [Mapping Resolver Cache](index-lifecycle.md#mapping-resolver-cache). ### Checklist: Adding a Queryable Model Field diff --git a/docs/guide/index-management.md b/docs/guide/index-management.md index 1a29c426..18352280 100644 --- a/docs/guide/index-management.md +++ b/docs/guide/index-management.md @@ -500,7 +500,7 @@ Neither of the following reindexes time-series data: `MaintainIndexesJob` (alias **Across different indexes** it depends on how you trigger it: `configuration.ReindexAsync()` processes indexes **sequentially** (one index fully finishes before the next starts), while `ElasticMigrationJob` reindexes them **in parallel** (`Task.WhenAll`, one task per outdated index). Either way each index is internally sequential, and a **distributed lock keyed on the alias** (`reindex:audit`) guarantees a given index is never reindexed by two runners at once — even across multiple application instances (pods, workers). The lock is held for 20 minutes and auto-renewed on every progress callback, so long partition copies keep it alive. ::: tip Predictable, bounded disk usage per index -Within one index the upgrade only ever duplicates **one partition at a time**, so bumping a single index (e.g. `audit`) needs roughly one extra partition of headroom regardless of how many partitions it has. If several indexes reindex in parallel (via `ElasticMigrationJob`), peak extra disk is about the sum of one in-flight partition per concurrently-migrating index. Wall-clock time scales with partition count; run during off-peak hours if needed. +Within one index the upgrade only ever duplicates **one partition** at a time, so bumping a single index (e.g. `audit`) needs roughly one extra partition of headroom regardless of how many partitions it has. If several indexes reindex in parallel (via `ElasticMigrationJob`), peak extra disk is about the sum of one in-flight partition per concurrently-migrating index. Wall-clock time scales with partition count; run during off-peak hours if needed. ::: #### Multiple versions and interrupted upgrades @@ -677,7 +677,7 @@ When a single script applies, it is sent directly to Elasticsearch. When multipl ```javascript void f000(def ctx) { /* v2 rename script */ } -void f001(def ctx) { /* v2 remove script */ } +void f001(def ctx) { /* v3 remove script */ } void f002(def ctx) { /* v3 custom script */ } f000(ctx); f001(ctx); f002(ctx); ``` @@ -887,6 +887,9 @@ query paths. Concurrent lookups share one load through the index's long-lived re calls remain supported, but block while the asynchronous load completes. Disposing an index disposes its initialized resolver and cancels outstanding mapping I/O without creating an unused resolver. +A caller's cancellation token cancels its wait without canceling a shared load needed by other callers. +Failed or empty reloads retain the last usable snapshot; explicit invalidation clears it. + Each reload discovers the latest partition again so newly created partitions are visible without restarting the application. This does not merge mappings across historical partitions. Keep indexes and their resolvers long-lived; creating one per request bypasses their caches and multiplies metadata requests. Automatic reload @@ -913,7 +916,7 @@ This clears the cached server mapping and forces the next `GetMapping()` call to | `ElasticMappingResolver` field cache | Snapshot lifetime; unresolved fields trigger reloads with a five-second cooldown | `index.MappingResolver.RefreshMapping()` | | `_isEnsured` flag (`Index` / `VersionedIndex`) | Process lifetime (one-time flag) | Deleting the index resets it; otherwise persists until app restart | | `_ensuredDates` (`DailyIndex`) | Process lifetime per-date | Cleared on `DeleteAsync(name)` or `Dispose()`; otherwise persists until app restart | -| `ConfigureIndexesAsync` cache marker | 5 minutes (distributed via `ICacheClient`) | Automatically expires; or call `ConfigureIndexesAsync(force: true)` | +| `ConfigureIndexesAsync` cache marker | 5 minutes (distributed via `ICacheClient`) | Automatically expires; pass explicit indexes to bypass the configuration-level lock and cache marker | #### No cluster-side action needed @@ -1141,10 +1144,17 @@ await configuration.ConfigureIndexesAsync(); // Subsequent calls within 5 minutes skip (fast path) await configuration.ConfigureIndexesAsync(); -// Passing explicit indexes bypasses the lock and cache marker -await configuration.ConfigureIndexesAsync([myIndex]); +// Passing explicit indexes bypasses the configuration-level lock and cache marker. +await configuration.ConfigureIndexesAsync([myIndex], beginReindexingOutdated: false); + +// Or explicitly configure all registered indexes without enqueueing reindex work. +await configuration.ConfigureIndexesAsync(configuration.Indexes, beginReindexingOutdated: false); ``` +There is no `force` parameter. Explicit selection does not change daily/monthly mapping behavior: +existing partitions still require a manual PUT Mapping. Run `ReindexAsync()` separately when a version +upgrade is required and no reindex queue worker is configured. + ### Maintain Indexes Run maintenance tasks: