Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -42,10 +42,14 @@ internal DatabaseMigrationRefreshService(
_refreshInterval = refreshInterval;
}

protected override async Task ExecuteAsync(CancellationToken stoppingToken)
public override async Task StartAsync(CancellationToken cancellationToken)
{
await RefreshAsync(stoppingToken).ConfigureAwait(false);
await RefreshAsync(cancellationToken).ConfigureAwait(false);
await base.StartAsync(cancellationToken).ConfigureAwait(false);
}

protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
using PeriodicTimer timer = new(_refreshInterval);
try
{
Expand Down
38 changes: 38 additions & 0 deletions Libraries/Spark.Store.MongoDB/DatabaseMigrationService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
using MongoDB.Driver;
using Spark.Engine.Store;
using Spark.Engine.Store.Interfaces;
using Spark.Store.MongoDB.Search.Common;
using System;
using System.Collections.Generic;
using System.Threading;
Expand All @@ -21,8 +22,11 @@ public sealed class DatabaseMigrationService : IDatabaseMigrationService
private const string CompletedAtField = "completedAt";

private readonly IMongoCollection<BsonDocument> _collection;
private readonly IMongoCollection<BsonDocument> _resources;
private readonly IMongoCollection<BsonDocument> _searchIndex;
private readonly SemaphoreSlim _stateLock = new(1, 1);
private IReadOnlyDictionary<int, string> _appliedMigrations = new Dictionary<int, string>();
private bool _freshDatabaseCheckCompleted;
private int _currentVersion;

public DatabaseMigrationService(string connectionString)
Expand All @@ -34,6 +38,8 @@ internal DatabaseMigrationService(IMongoDatabase database)
{
ArgumentNullException.ThrowIfNull(database);
_collection = database.GetCollection<BsonDocument>(Collection.SchemaMigrations);
_resources = database.GetCollection<BsonDocument>(Collection.RESOURCE);
_searchIndex = database.GetCollection<BsonDocument>(MongoCollections.SEARCH_INDEX_COLLECTION);
}

public int CurrentVersion => Volatile.Read(ref _currentVersion);
Expand All @@ -51,6 +57,21 @@ public async Task RefreshAsync(CancellationToken cancellationToken = default)
.ToListAsync(cancellationToken)
.ConfigureAwait(false);

if (documents.Count == 0 && !_freshDatabaseCheckCompleted)
{
bool isFreshDatabase = await IsFreshDatabaseAsync(cancellationToken).ConfigureAwait(false);
if (isFreshDatabase)
{
BsonDocument migration = await UpsertMigrationAsync(
DatabaseMigrations.StructuredStringTokenIndex,
cancellationToken)
.ConfigureAwait(false);
documents.Add(migration);
}

_freshDatabaseCheckCompleted = true;
}

var appliedMigrations = new Dictionary<int, string>(documents.Count);
var expectedVersion = 1;

Expand All @@ -76,6 +97,23 @@ public async Task RefreshAsync(CancellationToken cancellationToken = default)
}
}

private async Task<bool> IsFreshDatabaseAsync(CancellationToken cancellationToken)
{
CountOptions options = new() { Limit = 1 };
long resourceCount = await _resources
.CountDocumentsAsync(FilterDefinition<BsonDocument>.Empty, options, cancellationToken)
.ConfigureAwait(false);
if (resourceCount != 0)
{
return false;
}

long searchIndexCount = await _searchIndex
.CountDocumentsAsync(FilterDefinition<BsonDocument>.Empty, options, cancellationToken)
.ConfigureAwait(false);
return searchIndexCount == 0;
}

public async Task RecordCompletedAsync(
DatabaseMigration migration,
CancellationToken cancellationToken = default
Expand Down
83 changes: 64 additions & 19 deletions Libraries/Spark.Store.MongoDB/Search/CriteriaMongoExtensions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,10 @@ internal static SearchParameter FindSearchParamDefinition(this Criterium param,
return param.SearchParameters?.FirstOrDefault(sp => sp.Resource == resourceType || sp.Resource == "Resource");
}

internal static FilterDefinition<BsonDocument> ToFilter(this Criterium param, string resourceType)
internal static FilterDefinition<BsonDocument> ToFilter(
this Criterium param,
string resourceType,
bool includePlainStringTokenQuery = true)
{
//Maybe it's a generic parameter.
if (FixedQueries.TryGetValue(param.ParamName, out Func<Criterium, FilterDefinition<BsonDocument>> query))
Expand All @@ -69,14 +72,24 @@ internal static FilterDefinition<BsonDocument> ToFilter(this Criterium param, st
{

// todo: DSTU2 - modifier not in SearchParameter
return CreateFilter(critSp, param.Operator, param.Modifier, param.Operand);
return CreateFilter(
critSp,
param.Operator,
param.Modifier,
param.Operand,
includePlainStringTokenQuery);
//return null;
}

throw new UnknownSearchParameterException(string.Format("Resource {0} has no parameter with the name {1}.", resourceType, param.ParamName));
}

private static FilterDefinition<BsonDocument> CreateFilter(SearchParameter parameter, Operator op, String modifier, Expression operand)
private static FilterDefinition<BsonDocument> CreateFilter(
SearchParameter parameter,
Operator op,
String modifier,
Expression operand,
bool includePlainStringTokenQuery)
{
if (op == Operator.CHAIN)
{
Expand All @@ -98,7 +111,7 @@ private static FilterDefinition<BsonDocument> CreateFilter(SearchParameter param
switch (parameter.Type)
{
case SearchParamType.Composite:
return CompositeQuery(parameter, op, modifier, valueOperand);
return CompositeQuery(parameter, op, modifier, valueOperand, includePlainStringTokenQuery);
case SearchParamType.Date:
return DateQuery(parameterName, op, modifier, valueOperand);
case SearchParamType.Number:
Expand All @@ -117,7 +130,7 @@ private static FilterDefinition<BsonDocument> CreateFilter(SearchParameter param
}
else if (modifier == Modifier.IDENTIFIER)
{
return TokenQuery(parameterName, op, Modifier.EXACT, valueOperand);
return TokenQuery(parameterName, op, Modifier.EXACT, valueOperand, includePlainStringTokenQuery);
}
else
{
Expand All @@ -126,7 +139,7 @@ private static FilterDefinition<BsonDocument> CreateFilter(SearchParameter param
case SearchParamType.String:
return StringQuery(parameterName, op, modifier, valueOperand);
case SearchParamType.Token:
return TokenQuery(parameterName, op, modifier, valueOperand);
return TokenQuery(parameterName, op, modifier, valueOperand, includePlainStringTokenQuery);
case SearchParamType.Uri:
return UriQuery(parameterName, op, modifier, valueOperand);
default:
Expand Down Expand Up @@ -327,7 +340,12 @@ private static FilterDefinition<BsonDocument> QuantityQuery(string parameterName
return query;
}

private static FilterDefinition<BsonDocument> TokenQuery(String parameterName, Operator optor, String modifier, ValueExpression operand)
private static FilterDefinition<BsonDocument> TokenQuery(
String parameterName,
Operator optor,
String modifier,
ValueExpression operand,
bool includePlainStringTokenQuery)
{
//$elemMatch only works on array values. But the MongoIndexMapper only creates an array if there are multiple values for a given parameter.
//So we also construct a query for when there is only one set of values in the searchIndex, hence there is no array.
Expand All @@ -349,14 +367,15 @@ private static FilterDefinition<BsonDocument> TokenQuery(String parameterName, O
var arrayQueries = new List<FilterDefinition<BsonDocument>>();
var noArrayQueries = new List<FilterDefinition<BsonDocument>>{
Builders<BsonDocument>.Filter.Not(Builders<BsonDocument>.Filter.Type(parameterName, BsonType.Array))};
var plainStringQueries = new List<FilterDefinition<BsonDocument>>{
Builders<BsonDocument>.Filter.Type(parameterName, BsonType.String)};
List<FilterDefinition<BsonDocument>> plainStringQueries = includePlainStringTokenQuery
? [Builders<BsonDocument>.Filter.Type(parameterName, BsonType.String)]
: null;

if (!string.IsNullOrEmpty(typedEqOperand.Value))
{
noArrayQueries.Add(Builders<BsonDocument>.Filter.Eq(codefield, typedEqOperand.Value));
arrayQueries.Add(Builders<BsonDocument>.Filter.Eq("code", typedEqOperand.Value));
plainStringQueries.Add(Builders<BsonDocument>.Filter.Eq(parameterName, typedEqOperand.Value));
plainStringQueries?.Add(Builders<BsonDocument>.Filter.Eq(parameterName, typedEqOperand.Value));
}

//Handle the system part, if present.
Expand All @@ -366,28 +385,44 @@ private static FilterDefinition<BsonDocument> TokenQuery(String parameterName, O
{
arrayQueries.Add(Builders<BsonDocument>.Filter.Exists("system", false));
noArrayQueries.Add(Builders<BsonDocument>.Filter.Exists(systemfield, false));
plainStringQueries.Add(Builders<BsonDocument>.Filter.Exists("system", false));
plainStringQueries?.Add(Builders<BsonDocument>.Filter.Exists("system", false));
}
else
{
arrayQueries.Add(Builders<BsonDocument>.Filter.Eq("system", typedEqOperand.Namespace));
noArrayQueries.Add(Builders<BsonDocument>.Filter.Eq(systemfield, typedEqOperand.Namespace));
plainStringQueries.Add(Builders<BsonDocument>.Filter.Eq("system", typedEqOperand.Namespace));
plainStringQueries?.Add(Builders<BsonDocument>.Filter.Eq("system", typedEqOperand.Namespace));
}
}

//Combine code and system
var arrayEqQuery = Builders<BsonDocument>.Filter.ElemMatch(parameterName, Builders<BsonDocument>.Filter.And(arrayQueries));
var noArrayEqQuery = Builders<BsonDocument>.Filter.And(noArrayQueries);
if (!includePlainStringTokenQuery)
{
return modifier == Modifier.NOT
? Builders<BsonDocument>.Filter.And(
Builders<BsonDocument>.Filter.Not(arrayEqQuery),
Builders<BsonDocument>.Filter.Not(noArrayEqQuery))
: Builders<BsonDocument>.Filter.Or(arrayEqQuery, noArrayEqQuery);
}

var plainStringQuery = Builders<BsonDocument>.Filter.And(plainStringQueries);
return modifier == Modifier.NOT ?
Builders<BsonDocument>.Filter.And(Builders<BsonDocument>.Filter.Not(arrayEqQuery),
Builders<BsonDocument>.Filter.Not(noArrayEqQuery), Builders<BsonDocument>.Filter.Not(plainStringQuery))
return modifier == Modifier.NOT
? Builders<BsonDocument>.Filter.And(
Builders<BsonDocument>.Filter.Not(arrayEqQuery),
Builders<BsonDocument>.Filter.Not(noArrayEqQuery),
Builders<BsonDocument>.Filter.Not(plainStringQuery))
: Builders<BsonDocument>.Filter.Or(arrayEqQuery, noArrayEqQuery, plainStringQuery);
}
case Operator.IN:
IEnumerable<ValueExpression> opMultiple = ((ChoiceValue)operand).Choices;
var queries = opMultiple.Select(choice => TokenQuery(parameterName, Operator.EQ, modifier, choice));
var queries = opMultiple.Select(choice => TokenQuery(
parameterName,
Operator.EQ,
modifier,
choice,
includePlainStringTokenQuery));
return modifier == Modifier.NOT ? Builders<BsonDocument>.Filter.And(queries) : Builders<BsonDocument>.Filter.Or(queries);
case Operator.ISNULL:
return Builders<BsonDocument>.Filter.And(Builders<BsonDocument>.Filter.Eq(parameterName, BsonNull.Value), Builders<BsonDocument>.Filter.Eq(textfield, BsonNull.Value)); //We don't use Builders<BsonDocument>.Filter.NotExists, because that would exclude resources that have this field with an explicit null in it.
Expand Down Expand Up @@ -478,15 +513,25 @@ private static FilterDefinition<BsonDocument> DateQuery(String parameterName, Op
}
}

private static FilterDefinition<BsonDocument> CompositeQuery(SearchParameter parameterDef, Operator optor, String modifier, ValueExpression operand)
private static FilterDefinition<BsonDocument> CompositeQuery(
SearchParameter parameterDef,
Operator optor,
String modifier,
ValueExpression operand,
bool includePlainStringTokenQuery)
{
if (optor == Operator.IN)
{
var choices = ((ChoiceValue)operand);
var queries = new List<FilterDefinition<BsonDocument>>();
foreach (var choice in choices.Choices)
{
queries.Add(CompositeQuery(parameterDef, Operator.EQ, modifier, choice));
queries.Add(CompositeQuery(
parameterDef,
Operator.EQ,
modifier,
choice,
includePlainStringTokenQuery));
}
return Builders<BsonDocument>.Filter.Or(queries);
}
Expand All @@ -511,7 +556,7 @@ private static FilterDefinition<BsonDocument> CompositeQuery(SearchParameter par
Operand = components[i],
Modifier = modifier
};
queries.Add(subCrit.ToFilter(parameterDef.Resource));
queries.Add(subCrit.ToFilter(parameterDef.Resource, includePlainStringTokenQuery));
}
return Builders<BsonDocument>.Filter.And(queries);
}
Expand Down
34 changes: 32 additions & 2 deletions Libraries/Spark.Store.MongoDB/Search/MongoSearcher.cs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@
using Spark.Engine.Core;
using Spark.Engine.Search;
using Spark.Engine.Search.Types;
using Spark.Engine.Store;
using Spark.Engine.Store.Interfaces;
using Spark.Store.MongoDB.Search.Common;
using System;
using System.Collections.Generic;
Expand All @@ -27,8 +29,27 @@ public class MongoSearcher
private readonly IMongoCollection<BsonDocument> _collection;
private readonly ILocalhost _localhost;
private readonly IFhirModel _fhirModel;
private readonly IDatabaseMigrationService _databaseMigrationService;
private readonly IReferenceNormalizationService _referenceNormalizationService;

public MongoSearcher(
MongoIndexStore mongoIndexStore,
ILocalhost localhost,
IFhirModel fhirModel,
IReferenceNormalizationService referenceNormalizationService,
IDatabaseMigrationService databaseMigrationService)
{
_collection = mongoIndexStore.Collection;
_localhost = localhost;
_fhirModel = fhirModel;
_databaseMigrationService = databaseMigrationService ??
throw new ArgumentNullException(nameof(databaseMigrationService));
_referenceNormalizationService = referenceNormalizationService;
}

[Obsolete(
"Use MongoSearcher(MongoIndexStore, ILocalhost, IFhirModel, IReferenceNormalizationService, IDatabaseMigrationService) instead."
)]
public MongoSearcher(MongoIndexStore mongoIndexStore, ILocalhost localhost, IFhirModel fhirModel,
IReferenceNormalizationService referenceNormalizationService = null)
{
Expand Down Expand Up @@ -116,6 +137,10 @@ private SearchResults KeysToSearchResults(IEnumerable<BsonValue> keys)
return results;
}

internal bool IncludePlainStringTokenQuery =>
_databaseMigrationService == null
|| !_databaseMigrationService.IsApplied(DatabaseMigrations.StructuredStringTokenIndex.Version);

private List<BsonValue> CollectKeys(string resourceType, IEnumerable<Criterium> criteria, int level = 0)
{
return CollectKeys(resourceType, criteria, null, level);
Expand Down Expand Up @@ -197,9 +222,14 @@ private static SortDefinition<BsonDocument> CreateSortBy(IList<(string, SortOrde

}

private static FilterDefinition<BsonDocument> CreateMongoQuery(string resourceType, SearchResults results, int level, Dictionary<Criterium, Criterium> closedCriteria)
private FilterDefinition<BsonDocument> CreateMongoQuery(
string resourceType,
SearchResults results,
int level,
Dictionary<Criterium, Criterium> closedCriteria)
{
FilterDefinition<BsonDocument> resultQuery = CriteriaMongoExtensions.ResourceFilter(resourceType, level);
bool includePlainStringTokenQuery = IncludePlainStringTokenQuery;
if (closedCriteria.Count > 0)
{
var criteriaQueries = new List<FilterDefinition<BsonDocument>>();
Expand All @@ -209,7 +239,7 @@ private static FilterDefinition<BsonDocument> CreateMongoQuery(string resourceTy
{
try
{
criteriaQueries.Add(crit.Value.ToFilter(resourceType));
criteriaQueries.Add(crit.Value.ToFilter(resourceType, includePlainStringTokenQuery));
}
catch (ArgumentException ex)
{
Expand Down
Loading