diff --git a/src/persistence/Elsa.Persistence.DocumentDb/DocumentDbCollectionInfo.cs b/src/persistence/Elsa.Persistence.DocumentDb/DocumentDbCollectionInfo.cs new file mode 100644 index 000000000..fec2bf05a --- /dev/null +++ b/src/persistence/Elsa.Persistence.DocumentDb/DocumentDbCollectionInfo.cs @@ -0,0 +1,9 @@ +namespace Elsa.Persistence.DocumentDb +{ + public class DocumentDbCollectionInfo + { + public string Name { get; set; } + public string TenantId { get; set; } + public int OfferThroughput { get; set; } + } +} \ No newline at end of file diff --git a/src/persistence/Elsa.Persistence.DocumentDb/DocumentDbStorage.cs b/src/persistence/Elsa.Persistence.DocumentDb/DocumentDbStorage.cs index 98e1df854..9f8ae5fa6 100644 --- a/src/persistence/Elsa.Persistence.DocumentDb/DocumentDbStorage.cs +++ b/src/persistence/Elsa.Persistence.DocumentDb/DocumentDbStorage.cs @@ -19,7 +19,7 @@ namespace Elsa.Persistence.DocumentDb private readonly SemaphoreSlim clientInstanceLock; private readonly DocumentDbStorageOptions options; private readonly DocumentClient client; - private Dictionary collectionUris; + private Dictionary collectionInfos; public DocumentDbStorage(IOptions options) { @@ -60,7 +60,7 @@ namespace Elsa.Persistence.DocumentDb try { - if (collectionUris != null) + if (collectionInfos != null) { return client; } @@ -72,21 +72,25 @@ namespace Elsa.Persistence.DocumentDb Id = options.DatabaseName }); - var tasks = options.CollectionNames.Select(async collectionName => + var tasks = options.CollectionInfos.Select(async collectionInfo => { + var partitionKeyDefinition = new PartitionKeyDefinition(); + partitionKeyDefinition.Paths.Add("/tenantId"); + var databaseUri = UriFactory.CreateDatabaseUri(database.Resource.Id); var collection = await client.CreateDocumentCollectionIfNotExistsAsync(databaseUri, new DocumentCollection { - Id = collectionName.Value - }); + Id = collectionInfo.Value.Name, + PartitionKey = partitionKeyDefinition + }, new RequestOptions { OfferThroughput = collectionInfo.Value.OfferThroughput }); var uri = UriFactory.CreateDocumentCollectionUri(options.DatabaseName, collection.Resource.Id); - return (collectionName.Key, Name: collectionName.Value, Uri: uri); + return (collectionInfo.Key, collectionInfo.Value.Name, Uri: uri, collectionInfo.Value.TenantId); }); var results = await Task.WhenAll(tasks); - collectionUris = results.ToDictionary(x => x.Key, x => (x.Name, x.Uri)); + collectionInfos = results.ToDictionary(x => x.Key, x => (x.Name, x.Uri, x.TenantId)); return client; } @@ -96,7 +100,20 @@ namespace Elsa.Persistence.DocumentDb } } - public Uri GetWorkflowDefinitionCollectionUri() => collectionUris["WorkflowDefinition"].Uri; - public Uri GetWorkflowInstanceCollectionUri() => collectionUris["WorkflowInstance"].Uri; + public async Task<(string Name, Uri Uri, string TenantId)> GetWorkflowDefinitionCollectionInfoAsync() + { + return await GetCollectionInfoAsync("WorkflowDefinition"); + } + + public async Task<(string Name, Uri Uri, string TenantId)> GetWorkflowInstanceCollectionInfoAsync() + { + return await GetCollectionInfoAsync("WorkflowInstance"); + } + + private async Task<(string Name, Uri Uri, string TenantId)> GetCollectionInfoAsync(string key) + { + await GetDocumentClient(); + return collectionInfos[key]; + } } } \ No newline at end of file diff --git a/src/persistence/Elsa.Persistence.DocumentDb/DocumentDbStorageOptions.cs b/src/persistence/Elsa.Persistence.DocumentDb/DocumentDbStorageOptions.cs index 8c6c4b709..5e28ef1d8 100644 --- a/src/persistence/Elsa.Persistence.DocumentDb/DocumentDbStorageOptions.cs +++ b/src/persistence/Elsa.Persistence.DocumentDb/DocumentDbStorageOptions.cs @@ -4,6 +4,7 @@ using System.Collections.Generic; namespace Elsa.Persistence.DocumentDb { + public class DocumentDbStorageOptions { /// @@ -31,12 +32,12 @@ namespace Elsa.Persistence.DocumentDb public string AuthSecret { get; set; } /// - /// Gets or sets the collection names under the DocumentDB. + /// Gets or sets the collection details under the DocumentDB. /// /// - /// The collection names. + /// The collection details. /// - public IDictionary CollectionNames { get; set; } + public IDictionary CollectionInfos { get; set; } /// /// Get or sets the request timeout for DocumentDB client. Default value set to 30 seconds diff --git a/src/persistence/Elsa.Persistence.DocumentDb/Documents/WorkflowDefinitionVersionDocument.cs b/src/persistence/Elsa.Persistence.DocumentDb/Documents/WorkflowDefinitionVersionDocument.cs index 00c2f9115..04b9faf50 100644 --- a/src/persistence/Elsa.Persistence.DocumentDb/Documents/WorkflowDefinitionVersionDocument.cs +++ b/src/persistence/Elsa.Persistence.DocumentDb/Documents/WorkflowDefinitionVersionDocument.cs @@ -6,7 +6,8 @@ namespace Elsa.Persistence.DocumentDb.Documents { public class WorkflowDefinitionVersionDocument : DocumentBase { - [JsonProperty(PropertyName = "type")] public string Type { get; } = nameof(WorkflowDefinitionVersionDocument); + [JsonProperty(PropertyName = "type")] + public string Type { get; } = nameof(WorkflowDefinitionVersionDocument); [JsonProperty(PropertyName = "definitionId")] public string DefinitionId { get; set; } @@ -39,5 +40,8 @@ namespace Elsa.Persistence.DocumentDb.Documents [JsonProperty(PropertyName = "isLatest")] public bool IsLatest { get; set; } + + [JsonProperty(PropertyName = "tenantId")] + public string TenantId { get; set; } } } \ No newline at end of file diff --git a/src/persistence/Elsa.Persistence.DocumentDb/Documents/WorkflowInstanceDocument.cs b/src/persistence/Elsa.Persistence.DocumentDb/Documents/WorkflowInstanceDocument.cs index cb29e24b4..d4658eb94 100644 --- a/src/persistence/Elsa.Persistence.DocumentDb/Documents/WorkflowInstanceDocument.cs +++ b/src/persistence/Elsa.Persistence.DocumentDb/Documents/WorkflowInstanceDocument.cs @@ -10,7 +10,8 @@ namespace Elsa.Persistence.DocumentDb.Documents [JsonProperty(PropertyName = "definitionId")] public string DefinitionId { get; set; } - [JsonProperty(PropertyName = "type")] public string Type { get; } = nameof(WorkflowInstanceDocument); + [JsonProperty(PropertyName = "type")] + public string Type { get; } = nameof(WorkflowInstanceDocument); [JsonProperty(PropertyName = "version")] public int Version { get; set; } @@ -43,7 +44,8 @@ namespace Elsa.Persistence.DocumentDb.Documents [JsonProperty(PropertyName = "scope")] public WorkflowExecutionScope Scope { get; set; } - [JsonProperty(PropertyName = "input")] public Variables Input { get; set; } + [JsonProperty(PropertyName = "input")] + public Variables Input { get; set; } [JsonProperty(PropertyName = "blockingActivities")] public HashSet BlockingActivities { get; set; } @@ -51,6 +53,10 @@ namespace Elsa.Persistence.DocumentDb.Documents [JsonProperty(PropertyName = "executionLog")] public ICollection ExecutionLog { get; set; } - [JsonProperty(PropertyName = "fault")] public WorkflowFault Fault { get; set; } + [JsonProperty(PropertyName = "fault")] + public WorkflowFault Fault { get; set; } + + [JsonProperty(PropertyName = "tenantId")] + public string TenantId { get; set; } } } diff --git a/src/persistence/Elsa.Persistence.DocumentDb/Helpers/ClientHelper.cs b/src/persistence/Elsa.Persistence.DocumentDb/Helpers/ClientHelper.cs index 542fb6529..19dc4198f 100644 --- a/src/persistence/Elsa.Persistence.DocumentDb/Helpers/ClientHelper.cs +++ b/src/persistence/Elsa.Persistence.DocumentDb/Helpers/ClientHelper.cs @@ -27,12 +27,12 @@ namespace Elsa.Persistence.DocumentDb.Helpers bool disableAutomaticIdGeneration = false, CancellationToken cancellationToken = default) { - return await client.ExecuteWithRetries(async() => await client.CreateDocumentAsync( + return await ExecuteWithRetries(async() => await client.CreateDocumentAsync( documentCollectionUri, document, options, disableAutomaticIdGeneration, - cancellationToken)); + cancellationToken), cancellationToken); } /// @@ -50,10 +50,10 @@ namespace Elsa.Persistence.DocumentDb.Helpers RequestOptions options = null, CancellationToken cancellationToken = default) { - return await client.ExecuteWithRetries(async() => await client.ReadDocumentAsync( + return await ExecuteWithRetries(async() => await client.ReadDocumentAsync( documentUri, options, - cancellationToken)); + cancellationToken), cancellationToken); } /// @@ -73,12 +73,12 @@ namespace Elsa.Persistence.DocumentDb.Helpers bool disableAutomaticIdGeneration = false, CancellationToken cancellationToken = default) { - return await client.ExecuteWithRetries(async () => await client.UpsertDocumentAsync( + return await ExecuteWithRetries(async () => await client.UpsertDocumentAsync( documentCollectionUri, document, options, disableAutomaticIdGeneration, - cancellationToken)); + cancellationToken), cancellationToken); } /// @@ -94,10 +94,10 @@ namespace Elsa.Persistence.DocumentDb.Helpers RequestOptions options = null, CancellationToken cancellationToken = default) { - return await client.ExecuteWithRetries(async () => await client.DeleteDocumentAsync( + return await ExecuteWithRetries(async () => await client.DeleteDocumentAsync( documentUri, options, - cancellationToken)); + cancellationToken), cancellationToken); } /// @@ -116,11 +116,11 @@ namespace Elsa.Persistence.DocumentDb.Helpers RequestOptions options = null, CancellationToken cancellationToken = default) { - return await client.ExecuteWithRetries(async () => await client.ReplaceDocumentAsync( + return await ExecuteWithRetries(async () => await client.ReplaceDocumentAsync( documentUri, document, options, - cancellationToken)); + cancellationToken), cancellationToken); } /// @@ -134,17 +134,19 @@ namespace Elsa.Persistence.DocumentDb.Helpers internal static async Task> ExecuteStoredProcedureWithRetriesAsync( this DocumentClient client, Uri storedProcedureUri, - params object[] procedureParams) + object[] procedureParams, + CancellationToken cancellationToken = default) { - return await client.ExecuteWithRetries(async () => await client.ExecuteStoredProcedureAsync( + return await ExecuteWithRetries(async () => await client.ExecuteStoredProcedureAsync( storedProcedureUri, - procedureParams)); + procedureParams), cancellationToken); } /// /// Execute the function with retries on throttle /// - internal static async Task> ExecuteNextWithRetriesAsync(this IDocumentQuery query) + internal static async Task> ExecuteNextWithRetriesAsync(this IDocumentQuery query, + CancellationToken cancellationToken = default) { while (true) { @@ -152,7 +154,7 @@ namespace Elsa.Persistence.DocumentDb.Helpers try { - return await query.ExecuteNextAsync(); + return await query.ExecuteNextAsync(cancellationToken); } catch (DocumentClientException ex) when (ex.StatusCode != null && (int) ex.StatusCode == 429) { @@ -164,16 +166,14 @@ namespace Elsa.Persistence.DocumentDb.Helpers timeSpan = de.RetryAfter; } - await Task.Delay(timeSpan); + await Task.Delay(timeSpan, cancellationToken); } } /// /// Execute the function with retries on throttle /// - internal static async Task ExecuteWithRetries( - this DocumentClient client, - Func> function) + private static async Task ExecuteWithRetries(Func> function, CancellationToken cancellationToken) { while (true) { @@ -193,7 +193,7 @@ namespace Elsa.Persistence.DocumentDb.Helpers timeSpan = de.RetryAfter; } - await Task.Delay(timeSpan); + await Task.Delay(timeSpan, cancellationToken); } } } diff --git a/src/persistence/Elsa.Persistence.DocumentDb/Helpers/QueryHelper.cs b/src/persistence/Elsa.Persistence.DocumentDb/Helpers/QueryHelper.cs index f659ad891..8e3c3cb1d 100644 --- a/src/persistence/Elsa.Persistence.DocumentDb/Helpers/QueryHelper.cs +++ b/src/persistence/Elsa.Persistence.DocumentDb/Helpers/QueryHelper.cs @@ -1,20 +1,21 @@ using Microsoft.Azure.Documents.Linq; using System.Collections.Generic; using System.Linq; +using System.Threading; using System.Threading.Tasks; namespace Elsa.Persistence.DocumentDb.Helpers { internal static class QueryHelper { - internal static async Task> ToQueryResultAsync(this IQueryable source) + internal static async Task> ToQueryResultAsync(this IQueryable source, CancellationToken cancellationToken = default) { var query = source.AsDocumentQuery(); var results = new List(); while (query.HasMoreResults) { - var nextResults = await query.ExecuteNextWithRetriesAsync(); + var nextResults = await query.ExecuteNextWithRetriesAsync(cancellationToken); results.AddRange(nextResults); } diff --git a/src/persistence/Elsa.Persistence.DocumentDb/IDocumentDbStorage.cs b/src/persistence/Elsa.Persistence.DocumentDb/IDocumentDbStorage.cs index 160a552f7..a6ea2a852 100644 --- a/src/persistence/Elsa.Persistence.DocumentDb/IDocumentDbStorage.cs +++ b/src/persistence/Elsa.Persistence.DocumentDb/IDocumentDbStorage.cs @@ -1,4 +1,4 @@ -using Microsoft.Azure.Documents.Client; +using Microsoft.Azure.Documents.Client; using System; using System.Threading.Tasks; @@ -8,7 +8,7 @@ namespace Elsa.Persistence.DocumentDb { string ToString(); Task GetDocumentClient(); - Uri GetWorkflowDefinitionCollectionUri(); - Uri GetWorkflowInstanceCollectionUri(); + Task<(string Name, Uri Uri, string TenantId)> GetWorkflowDefinitionCollectionInfoAsync(); + Task<(string Name, Uri Uri, string TenantId)> GetWorkflowInstanceCollectionInfoAsync(); } } \ No newline at end of file diff --git a/src/persistence/Elsa.Persistence.DocumentDb/Services/CosmosDbWorkflowDefinitionStore.cs b/src/persistence/Elsa.Persistence.DocumentDb/Services/CosmosDbWorkflowDefinitionStore.cs index b08f2050d..746ce3208 100644 --- a/src/persistence/Elsa.Persistence.DocumentDb/Services/CosmosDbWorkflowDefinitionStore.cs +++ b/src/persistence/Elsa.Persistence.DocumentDb/Services/CosmosDbWorkflowDefinitionStore.cs @@ -8,6 +8,7 @@ using Elsa.Persistence.DocumentDb.Documents; using Elsa.Persistence.DocumentDb.Extensions; using Elsa.Persistence.DocumentDb.Helpers; using Elsa.Services; +using Microsoft.Azure.Documents; using Microsoft.Azure.Documents.Client; namespace Elsa.Persistence.DocumentDb.Services @@ -24,69 +25,113 @@ namespace Elsa.Persistence.DocumentDb.Services } private async Task GetDocumentClient() => await storage.GetDocumentClient(); - private Uri GetCollectionUri() => storage.GetWorkflowDefinitionCollectionUri(); + private async Task<(string Name, Uri Uri, string TenantId)> GetCollectionInfoAsync() => + await storage.GetWorkflowDefinitionCollectionInfoAsync(); public async Task AddAsync(WorkflowDefinitionVersion definition, CancellationToken cancellationToken = default) { var document = Map(definition); var client = await GetDocumentClient(); - await client.CreateDocumentWithRetriesAsync(GetCollectionUri(), document, cancellationToken: cancellationToken); + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var requestOptions = new RequestOptions + { + PartitionKey = new PartitionKey(tenantId) + }; + document.TenantId = tenantId; + var response = await client.CreateDocumentWithRetriesAsync(uri, document, requestOptions, cancellationToken: cancellationToken); + document = (dynamic)response.Resource; + return Map(document); } public async Task DeleteAsync(string id, CancellationToken cancellationToken = default) { var client = await GetDocumentClient(); - var records = await client.CreateDocumentQuery(GetCollectionUri()).Where(c => c.DefinitionId == id).ToQueryResultAsync(); - foreach (var record in records) + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var feedOptions = new FeedOptions { - await client.DeleteDocumentAsync(record.SelfLink, cancellationToken: cancellationToken); - } - return records.Count; + PartitionKey = new PartitionKey(tenantId) + }; + var requestOptions = new RequestOptions + { + PartitionKey = new PartitionKey(tenantId) + }; + var query = client + .CreateDocumentQuery(uri, feedOptions) + .Where(c => c.TenantId == tenantId) + .Where(c => c.DefinitionId == id); + var documents = await query.ToQueryResultAsync(cancellationToken); + var tasks = documents.Select(d => + { + var documentUri = new Uri(d.SelfLink, UriKind.Relative); + return client.DeleteDocumentWithRetriesAsync(documentUri, + requestOptions, + cancellationToken); + }).ToArray(); + + await Task.WhenAll(tasks); + + return documents.Count; } public async Task GetByIdAsync(string id, VersionOptions version, CancellationToken cancellationToken = default) { var client = await GetDocumentClient(); - var query = client.CreateDocumentQuery(GetCollectionUri()) - .Where(c => c.DefinitionId == id).WithVersion(version); + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var feedOptions = new FeedOptions + { + PartitionKey = new PartitionKey(tenantId) + }; + var query = client + .CreateDocumentQuery(uri, feedOptions) + .Where(c => c.TenantId == tenantId) + .Where(c => c.DefinitionId == id) + .WithVersion(version); var document = query.AsEnumerable().FirstOrDefault(); + return Map(document); } - public async Task> ListAsync(VersionOptions version, CancellationToken cancellationToken = default) { var client = await GetDocumentClient(); - var query = client.CreateDocumentQuery(GetCollectionUri()) - .WithVersion(version).ToList(); - - return mapper.Map>(query); + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var feedOptions = new FeedOptions + { + PartitionKey = new PartitionKey(tenantId) + }; + var query = client + .CreateDocumentQuery(uri, feedOptions) + .Where(c => c.TenantId == tenantId) + .WithVersion(version); + var documents = await query.ToQueryResultAsync(cancellationToken); + + return Map(documents); } public async Task SaveAsync(WorkflowDefinitionVersion definition, CancellationToken cancellationToken = default) { var document = Map(definition); var client = await GetDocumentClient(); - await client.UpsertDocumentWithRetriesAsync(GetCollectionUri(), document, cancellationToken: cancellationToken); - return definition; + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var requestOptions = new RequestOptions + { + PartitionKey = new PartitionKey(tenantId) + }; + document.TenantId = tenantId; + var response = await client.UpsertDocumentWithRetriesAsync(uri, document, requestOptions, cancellationToken: cancellationToken); + + document = (dynamic)response.Resource; + return Map(document); } public async Task UpdateAsync(WorkflowDefinitionVersion definition, CancellationToken cancellationToken = default) { - var document = Map(definition); - var client = await GetDocumentClient(); - await client.UpsertDocumentWithRetriesAsync(GetCollectionUri(), document, cancellationToken: cancellationToken); - return Map(document); + return await SaveAsync(definition, cancellationToken); } - private WorkflowDefinitionVersionDocument Map(WorkflowDefinitionVersion source) - { - return mapper.Map(source); - } - - private WorkflowDefinitionVersion Map(WorkflowDefinitionVersionDocument source) - { - return mapper.Map(source); - } + private WorkflowDefinitionVersionDocument Map(WorkflowDefinitionVersion source) => mapper.Map(source); + private WorkflowDefinitionVersion Map(WorkflowDefinitionVersionDocument source) => mapper.Map(source); + private IEnumerable Map(IEnumerable source) => + mapper.Map>(source); } } diff --git a/src/persistence/Elsa.Persistence.DocumentDb/Services/CosmosDbWorkflowInstanceStore.cs b/src/persistence/Elsa.Persistence.DocumentDb/Services/CosmosDbWorkflowInstanceStore.cs index b8f3d4fa5..aa760ff3e 100644 --- a/src/persistence/Elsa.Persistence.DocumentDb/Services/CosmosDbWorkflowInstanceStore.cs +++ b/src/persistence/Elsa.Persistence.DocumentDb/Services/CosmosDbWorkflowInstanceStore.cs @@ -8,6 +8,7 @@ using Elsa.Models; using Elsa.Persistence.DocumentDb.Documents; using Elsa.Persistence.DocumentDb.Helpers; using Elsa.Services; +using Microsoft.Azure.Documents; using Microsoft.Azure.Documents.Client; namespace Elsa.Persistence.DocumentDb.Services @@ -24,52 +25,87 @@ namespace Elsa.Persistence.DocumentDb.Services } private async Task GetDocumentClient() => await storage.GetDocumentClient(); - private Uri GetCollectionUri() => storage.GetWorkflowInstanceCollectionUri(); + private async Task<(string Name, Uri Uri, string TenantId)> GetCollectionInfoAsync() => + await storage.GetWorkflowInstanceCollectionInfoAsync(); - public async Task DeleteAsync( - string id, + public async Task DeleteAsync(string id, CancellationToken cancellationToken = default) { var client = await GetDocumentClient(); - var query = client.CreateDocumentQuery(GetCollectionUri()) + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var feedOptions = new FeedOptions + { + PartitionKey = new PartitionKey(tenantId) + }; + var requestOptions = new RequestOptions + { + PartitionKey = new PartitionKey(tenantId) + }; + var query = client + .CreateDocumentQuery(uri, feedOptions) + .Where(c => c.TenantId == tenantId) .Where(c => c.Id == id); var document = query.AsEnumerable().FirstOrDefault(); if (document != null) { var documentUri = new Uri(document.SelfLink, UriKind.Relative); - await client.DeleteDocumentAsync(documentUri, cancellationToken: cancellationToken); + await client.DeleteDocumentWithRetriesAsync(documentUri, + requestOptions, + cancellationToken); } } - public async Task GetByCorrelationIdAsync( - string correlationId, + public async Task GetByCorrelationIdAsync(string correlationId, CancellationToken cancellationToken = default) { var client = await GetDocumentClient(); - var query = client.CreateDocumentQuery(GetCollectionUri()) + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var feedOptions = new FeedOptions + { + PartitionKey = new PartitionKey(tenantId) + }; + var query = client + .CreateDocumentQuery(uri, feedOptions) + .Where(c => c.TenantId == tenantId) .Where(c => c.CorrelationId == correlationId); var document = query.AsEnumerable().FirstOrDefault(); + return Map(document); } - public async Task GetByIdAsync( - string id, + public async Task GetByIdAsync(string id, CancellationToken cancellationToken = default) { var client = await GetDocumentClient(); - var query = client.CreateDocumentQuery(GetCollectionUri()) + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var feedOptions = new FeedOptions + { + PartitionKey = new PartitionKey(tenantId) + }; + var query = client + .CreateDocumentQuery(uri, feedOptions) + .Where(c => c.TenantId == tenantId) .Where(c => c.Id == id); var document = query.AsEnumerable().FirstOrDefault(); + return Map(document); } public async Task> ListAllAsync(CancellationToken cancellationToken = default) { var client = await GetDocumentClient(); + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var feedOptions = new FeedOptions + { + PartitionKey = new PartitionKey(tenantId) + }; var query = client - .CreateDocumentQuery(GetCollectionUri()) + .CreateDocumentQuery(uri, feedOptions) + .Where(c => c.TenantId == tenantId) .OrderByDescending(x => x.CreatedAt); - return mapper.Map>(query); + var documents = await query.ToQueryResultAsync(cancellationToken); + + return Map(documents); } public async Task> ListByBlockingActivityAsync( @@ -78,9 +114,14 @@ namespace Elsa.Persistence.DocumentDb.Services CancellationToken cancellationToken = default) { var client = await GetDocumentClient(); - + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var feedOptions = new FeedOptions + { + PartitionKey = new PartitionKey(tenantId) + }; var query = client - .CreateDocumentQuery(GetCollectionUri()) + .CreateDocumentQuery(uri, feedOptions) + .Where(c => c.TenantId == tenantId) .Where(x => x.Status == WorkflowStatus.Executing); if (!string.IsNullOrWhiteSpace(correlationId)) @@ -91,8 +132,10 @@ namespace Elsa.Persistence.DocumentDb.Services query = query.Where(x => x.BlockingActivities.Any(y => y.ActivityType == activityType)); query = query.OrderByDescending(x => x.CreatedAt); - var instances = Map(query.ToList()); - return instances.GetBlockingActivities(activityType); + var documents = await query.ToQueryResultAsync(cancellationToken); + + var instances = Map(documents); + return instances.GetBlockingActivities(activityType).ToArray(); } public async Task> ListByDefinitionAsync( @@ -100,10 +143,19 @@ namespace Elsa.Persistence.DocumentDb.Services CancellationToken cancellationToken = default) { var client = await GetDocumentClient(); - var query = client.CreateDocumentQuery(GetCollectionUri()) + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var feedOptions = new FeedOptions + { + PartitionKey = new PartitionKey(tenantId) + }; + var query = client + .CreateDocumentQuery(uri, feedOptions) + .Where(c => c.TenantId == tenantId) .Where(c => c.DefinitionId == definitionId) .OrderByDescending(x => x.CreatedAt); - return Map(query.ToList()); + var documents = await query.ToQueryResultAsync(cancellationToken); + + return Map(documents); } public async Task> ListByStatusAsync( @@ -112,10 +164,18 @@ namespace Elsa.Persistence.DocumentDb.Services CancellationToken cancellationToken = default) { var client = await GetDocumentClient(); - var query = client.CreateDocumentQuery(GetCollectionUri()) + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var feedOptions = new FeedOptions + { + PartitionKey = new PartitionKey(tenantId) + }; + var query = client.CreateDocumentQuery(uri, feedOptions) + .Where(c => c.TenantId == tenantId) .Where(c => c.DefinitionId == definitionId && c.Status == status) .OrderByDescending(x => x.CreatedAt); - return Map(query.ToList()); + var documents = await query.ToQueryResultAsync(cancellationToken); + + return Map(documents); } public async Task> ListByStatusAsync( @@ -123,10 +183,18 @@ namespace Elsa.Persistence.DocumentDb.Services CancellationToken cancellationToken = default) { var client = await GetDocumentClient(); - var query = client.CreateDocumentQuery(GetCollectionUri()) + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var feedOptions = new FeedOptions + { + PartitionKey = new PartitionKey(tenantId) + }; + var query = client.CreateDocumentQuery(uri, feedOptions) + .Where(c => c.TenantId == tenantId) .Where(c => c.Status == status) .OrderByDescending(x => x.CreatedAt); - return Map(query.ToList()); + var documents = await query.ToQueryResultAsync(cancellationToken); + + return Map(documents); } public async Task SaveAsync( @@ -135,9 +203,16 @@ namespace Elsa.Persistence.DocumentDb.Services { var document = Map(instance); var client = await GetDocumentClient(); + var (_, uri, tenantId) = await GetCollectionInfoAsync(); + var requestOptions = new RequestOptions + { + PartitionKey = new PartitionKey(tenantId) + }; + document.TenantId = tenantId; var response = await client.UpsertDocumentWithRetriesAsync( - GetCollectionUri(), + uri, document, + requestOptions, cancellationToken: cancellationToken); document = (dynamic)response.Resource; @@ -146,8 +221,7 @@ namespace Elsa.Persistence.DocumentDb.Services private WorkflowInstanceDocument Map(WorkflowInstance source) => mapper.Map(source); private WorkflowInstance Map(WorkflowInstanceDocument source) => mapper.Map(source); - - private IEnumerable Map(IEnumerable source) => + private IEnumerable Map(IEnumerable source) => mapper.Map>(source); } } \ No newline at end of file