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));
}
}
}