using System; using System.Collections.Generic; using System.Linq; using System.Runtime.CompilerServices; using System.Threading; using System.Threading.Tasks; using InfluxDB.Client.Api.Domain; using InfluxDB.Client.Core.Flux.Domain; using Remotion.Linq; [assembly: InternalsVisibleTo("Client.Linq.Test, PublicKey=002400000480000094000000060200000024000052534131" + "0004000001000100efaac865f88dd35c90dc548945405aae34056eedbe42cad60971f89a861a78437e86d" + "95804a1aeeb0de18ac3728782f9dc8dbae2e806167a8bb64c0402278edcefd78c13dbe7f8d13de36eb362" + "21ec215c66ee2dfe7943de97b869c5eea4d92f92d345ced67de5ac8fc3cd2f8dd7e3c0c53bdb0cc433af8" + "59033d069cad397a7")] namespace InfluxDB.Client.Linq.Internal { /// /// Executor is called by ReLinq when query is executed. /// internal class InfluxDBQueryExecutor : IQueryExecutor { private readonly string _bucket; private readonly string _org; private readonly QueryApiSync _queryApiSync; private readonly QueryApi _queryApi; private readonly IMemberNameResolver _memberResolver; private readonly QueryableOptimizerSettings _queryableOptimizerSettings; /// /// Create InfluxDBQuery Executor for synchronous Queries. /// /// Specifies the source bucket. /// Specifies the source organization. /// The underlying API to execute Flux Query. /// Resolver for customized names. /// Settings for a Query optimization public InfluxDBQueryExecutor(string bucket, string org, QueryApiSync queryApi, IMemberNameResolver memberResolver, QueryableOptimizerSettings queryableOptimizerSettings) { _bucket = bucket; _org = org; _queryApiSync = queryApi; _memberResolver = memberResolver; _queryableOptimizerSettings = queryableOptimizerSettings; } /// /// Create InfluxDBQuery Executor for asynchronous Queries. /// /// Specifies the source bucket. /// Specifies the source organization. /// The underlying API to execute Flux Query. /// Resolver for customized names. /// Settings for a Query optimization public InfluxDBQueryExecutor(string bucket, string org, QueryApi queryApi, IMemberNameResolver memberResolver, QueryableOptimizerSettings queryableOptimizerSettings) { _bucket = bucket; _org = org; _queryApi = queryApi; _memberResolver = memberResolver; _queryableOptimizerSettings = queryableOptimizerSettings; } /// /// Executes the given as a scalar query, /// i.e. a query that ends with a aggregation operator such as Count, Sum, or Average. /// public T ExecuteScalar(QueryModel queryModel) { return ExecuteSingle(queryModel, false); } /// /// Executes the given as a scalar query, /// i.e. a query that ends with a result operator such as First, Last, Single, Min, or Max. /// public T ExecuteSingle(QueryModel queryModel, bool returnDefaultWhenEmpty) { return returnDefaultWhenEmpty ? ExecuteCollection(queryModel).SingleOrDefault() : ExecuteCollection(queryModel).Single(); } /// /// Executes a query with a collection result. /// public IEnumerable ExecuteCollection(QueryModel queryModel) { var query = GenerateQuery(queryModel, out var queryResultsSettings); if (_queryApiSync == null) { throw new ArgumentException("The 'QueryApiSync' has to be configured for sync queries."); } if (queryResultsSettings.ScalarAggregated) { var result = ApplyAggregate(_queryApiSync.QuerySync(query, _org), queryResultsSettings); return new List { result }; } return _queryApiSync.QuerySync(query, _org); } /// /// Executes an async query with a collection result. /// public IAsyncEnumerable ExecuteCollectionAsync(QueryModel queryModel, CancellationToken cancellationToken = new CancellationToken()) { var query = GenerateQuery(queryModel, out var queryResultsSettings); if (_queryApi == null) { throw new ArgumentException("The 'QueryApi' has to be configured for Async queries."); } if (queryResultsSettings.ScalarAggregated) { return AggregateAsync(_queryApi.QueryAsync(query, _org), queryResultsSettings, cancellationToken); } return _queryApi.QueryAsyncEnumerable(query, _org, cancellationToken); } /// /// Create a object that will be used for Querying. /// /// Expression Tree of LINQ Query /// Defines how to handle query results /// Query to Invoke internal Query GenerateQuery(QueryModel queryModel, out QueryResultsSettings settings) { var visitor = QueryVisitor(queryModel); settings = new QueryResultsSettings(queryModel); return visitor.GenerateQuery(); } /// /// Create QueryVisitor for specified model. /// /// Query Model /// Query Visitor internal InfluxDBQueryVisitor QueryVisitor(QueryModel queryModel) { var visitor = new InfluxDBQueryVisitor(_bucket, _memberResolver, _queryableOptimizerSettings); visitor.VisitQueryModel(queryModel); return visitor; } private async IAsyncEnumerable AggregateAsync(Task> tables, QueryResultsSettings queryResultsSettings, [EnumeratorCancellation] CancellationToken cancellationToken) { var result = tables .ContinueWith(t => ApplyAggregate(t.Result, queryResultsSettings), cancellationToken); yield return await result.ConfigureAwait(false); } private static T ApplyAggregate(IEnumerable tables, QueryResultsSettings queryResultsSettings) { var enumerable = tables .SelectMany(it => it.Records) .Select(it => it.GetValueByKey("linq_result_column")); var aggregated = queryResultsSettings.AggregateFunction(enumerable); return (T)Convert.ChangeType(aggregated, typeof(T)); } } }