sozsoft-platform/api/modules/Sozsoft.Mcp/Sozsoft.Mcp.Domain/Services/McpQueryExecutor.cs
2026-09-10 17:40:43 +03:00

105 lines
3.4 KiB
C#

using System.Collections.Generic;
using System.Diagnostics;
using System.Linq;
using System.Threading.Tasks;
using Sozsoft.Platform.Data.Seeds;
using Sozsoft.Platform.DynamicData;
using Sozsoft.Platform.Enums;
using Volo.Abp.Domain.Services;
namespace Sozsoft.Mcp.Services;
/// <summary>
/// Sorguyu istegin kapsamindaki veri kaynaginda kosturur ve satir tavanini uygular.
/// Veri kaynagi cozumu tek yerde kalsin diye burada toplanmistir.
/// </summary>
public interface IMcpQueryExecutor
{
/// <param name="dataSourceCode">
/// Sunucu kaydindaki kod. Bos ya da <c>!Tenant</c> ise veri kaynagi istegin kiracisindan
/// cozulur; host baglaminda varsayilan kaynaga duser.
/// </param>
Task<McpDataSourceContext> ResolveDataSourceAsync(string? dataSourceCode);
/// <summary>
/// Sorguyu kosturur. Sonuc kumesi <paramref name="maxRows"/> satirda kesilir; kesilme
/// olduysa sonuc bunu bildirir.
/// </summary>
Task<McpQueryRows> ExecuteAsync(
McpDataSourceContext context,
string sql,
IReadOnlyDictionary<string, object> parameters,
int maxRows);
}
public class McpQueryExecutor(IDynamicDataManager dynamicDataManager) : DomainService, IMcpQueryExecutor
{
public async Task<McpDataSourceContext> ResolveDataSourceAsync(string? dataSourceCode)
{
var useTenantScope = string.IsNullOrWhiteSpace(dataSourceCode)
|| dataSourceCode == McpConsts.TenantDataSourceMarker;
var resolvedCode = useTenantScope
? SeedConsts.DataSources.DefaultCode
: dataSourceCode!;
var (repository, connectionString, dataSourceType) =
await dynamicDataManager.GetAsync(useTenantScope, resolvedCode);
return new McpDataSourceContext(repository, connectionString, dataSourceType);
}
public async Task<McpQueryRows> ExecuteAsync(
McpDataSourceContext context,
string sql,
IReadOnlyDictionary<string, object> parameters,
int maxRows)
{
var stopwatch = Stopwatch.StartNew();
var rows = await context.Repository.QueryAsync(
sql,
context.ConnectionString,
parameters.ToDictionary(entry => entry.Key, entry => entry.Value));
stopwatch.Stop();
var columns = new List<string>();
var materialized = new List<Dictionary<string, object?>>();
var truncated = false;
foreach (var row in rows)
{
if (materialized.Count >= maxRows)
{
truncated = true;
break;
}
if (row is not IDictionary<string, object> values)
{
continue;
}
if (columns.Count == 0)
{
columns = [.. values.Keys];
}
materialized.Add(values.ToDictionary(entry => entry.Key, entry => (object?)entry.Value));
}
return new McpQueryRows(columns, materialized, truncated, stopwatch.ElapsedMilliseconds);
}
}
/// <summary>Cozumlenmis veri kaynagi: hangi depo, hangi baglanti, hangi saglayici.</summary>
public sealed record McpDataSourceContext(
IDynamicDataRepository Repository,
string ConnectionString,
DataSourceTypeEnum DataSourceType);
/// <summary>Sorgu sonucunun ham hali.</summary>
public sealed record McpQueryRows(
List<string> Columns,
List<Dictionary<string, object?>> Rows,
bool Truncated,
long ExecutionTimeMs);