using System.Collections.Concurrent;
using System.Diagnostics;
using System.Reflection;
using System.Runtime.CompilerServices;
using System.Text;
using System.Text.Json;
using Dapr;
using FastEndpoints;
using GB5Shared.DTO.Framework.Login;
using GB5Shared.DTO.Framework.ResponseStandard;
using GB5Shared.DTO.Report;
using GB5Shared.Export;
using GB5Shared.JsonConverter;
using GB5Shared.QueryExecutor;
using GB5Shared.Resource.EndPointResource;
using GB5Shared.Telemetry;
using Microsoft.Extensions.Caching.Memory;
using Microsoft.Extensions.Logging;
using OpenTelemetry.Trace;
using static GB5Shared.GB5Constant.Constant;
namespace GB5Shared.FastEndPoint
{
///
/// Base class for all report and large-list endpoints.
/// Developers implement returning .
/// The framework handles every export format (Excel, CSV, PDF, JSON, ReportData, NDJSON stream),
/// OTel tracing, and exception handling — no boilerplate in derived classes.
///
public abstract class BaseReportEndpoint
: Endpoint>
where TRequest : notnull
{
protected new ILogger Logger => Resolve>>();
// ── REQUIRED ─────────────────────────────────────────────────────────────
///
/// Return ALL rows as a stream. Called for every export format and for NDJSON streaming.
/// Paging params on the request are ignored for export — the stream delivers the full result set.
///
protected abstract IAsyncEnumerable ExecuteReportAsync(
TRequest req, LoginDTO login, CancellationToken ct);
// ── OPTIONAL OVERRIDE ─────────────────────────────────────────────────────
///
/// Normal JSON response (no ReportFormat header).
/// Override when custom pagination, totals, or summary data are needed for the UI table view.
/// Default: collects up to rows from .
///
protected virtual async Task> ExecutePagedAsync(
TRequest req, LoginDTO login, CancellationToken ct)
{
var rows = new List();
await foreach (var row in ExecuteReportAsync(req, login, ct).WithCancellation(ct).ConfigureAwait(false))
{
rows.Add(row);
if (rows.Count >= GetPageCap()) break;
}
return await GB5Shared.ResponseStandard.Response
.CreateSuccessResponse(rows, CacheKeyLevel.NOT_REQUIRED, login)
.ConfigureAwait(false);
}
// ── Computed/derived fields (MREPORTVSFIELDS.COMPUTEEXPRESSION) — grid path ─────────────
//
// Applied automatically in HandleAsync (see ApplyComputedFieldsAutoAsync below) to
// WHATEVER ExecutePagedAsync returns — the default above, or any subclass override's own
// custom pagination/totals/envelope shape — so a report author configuring a
// COMPUTEEXPRESSION on a view never needs to know this feature exists or add a line of
// their own to opt in. When zero computed fields are configured for the requested view
// (every report/view that never uses this feature), the response's Body is returned
// completely unchanged — this is a strict no-op for every existing report.
//
// ApplyComputedFieldsAsync/ApplyComputedFieldsToEnvelopeAsync below remain available for
// any caller that still wants to apply computed fields somewhere outside this pipeline
// (e.g. before building a custom envelope) — but calling them explicitly here is no longer
// required; ComputedFieldEvaluator's row-level "already present" skip makes doing so
// harmless (not double-computed) if some endpoint still does.
private async Task> LoadComputedFieldDefsAsync(TRequest req, LoginDTO login, CancellationToken ct)
{
var reportCallingDTO = GetReportCallingDTOFromRequest(req);
var reportViewId = reportCallingDTO?.ReportViewId ?? 0;
if (reportViewId <= 0)
return new List();
var queryExecutor = TryResolve();
var cache = TryResolve();
return await ComputedFieldCatalog.GetForViewAsync(queryExecutor, cache, reportViewId, login, ct).ConfigureAwait(false);
}
///
/// List-shaped grid response. Returns itself, unchanged, when the
/// requested view has no computed fields configured.
///
protected async Task ApplyComputedFieldsAsync(
List rows, TRequest req, LoginDTO login, CancellationToken ct)
{
var defs = await LoadComputedFieldDefsAsync(req, login, ct).ConfigureAwait(false);
return ComputedFieldEvaluator.Project(rows, defs, Logger);
}
///
/// Envelope-shaped grid response (e.g. { Items, TotalCount, TotalDebit, ... } ).
/// Returns itself, unchanged, when the requested view has no
/// computed fields configured; otherwise returns a new object with every original property
/// preserved and replaced by the projected rows.
///
protected async Task ApplyComputedFieldsToEnvelopeAsync(
object envelope, string rowsPropertyName, TRequest req, LoginDTO login, CancellationToken ct)
{
var defs = await LoadComputedFieldDefsAsync(req, login, ct).ConfigureAwait(false);
if (defs.Count == 0)
return envelope;
var props = envelope.GetType().GetProperties(BindingFlags.Public | BindingFlags.Instance);
var rowsProp = props.FirstOrDefault(p => p.Name == rowsPropertyName);
if (rowsProp?.GetValue(envelope) is not System.Collections.IEnumerable rowsValue)
return envelope;
var projectedRows = ComputedFieldEvaluator.Project(rowsValue, defs, Logger);
if (ReferenceEquals(projectedRows, rowsValue))
return envelope;
var result = new Dictionary(StringComparer.Ordinal);
foreach (var p in props)
result[p.Name] = p.Name == rowsPropertyName ? projectedRows : p.GetValue(envelope);
return result;
}
/// Maximum rows returned by the default . Override to raise the limit.
protected virtual int GetPageCap() => 500;
// ── Automatic computed-field application (grid path) ────────────────────────────────
// Called from HandleAsync right after ExecutePagedAsync returns, regardless of whether
// that method is the default above or a subclass override. Auto-detects the response's
// shape rather than requiring a caller-supplied property name: Body itself may be a flat
// IEnumerable (the default ExecutePagedAsync's shape), or an envelope object exposing its
// rows via a public IEnumerable property (the common override shape — Items/Rows/etc, any
// name). Whichever it finds, it hands that collection to the same ComputedFieldEvaluator
// already used by ApplyComputedFieldsAsync/ApplyComputedFieldsToEnvelopeAsync above and by
// the export pipeline (ReportExport.cs's DecorateWithComputedFields) — one evaluator, three
// entry points, so a COMPUTEEXPRESSION configured on a view behaves identically in Grid,
// JSON, and every export format without three separate implementations to keep in sync.
private async Task> ApplyComputedFieldsAutoAsync(
ResponseStandardDTO response, TRequest req, LoginDTO login, CancellationToken ct)
{
if (response.Status == GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed
|| response.Body is not { } body)
return response;
var defs = await LoadComputedFieldDefsAsync(req, login, ct).ConfigureAwait(false);
if (defs.Count == 0)
return response; // strict no-op — matches every existing report's current behavior
if (body is System.Collections.IEnumerable directRows and not string)
{
var projected = ComputedFieldEvaluator.Project(directRows, defs, Logger);
if (!ReferenceEquals(projected, directRows))
response.Body = projected;
return response;
}
var props = body.GetType().GetProperties(BindingFlags.Public | BindingFlags.Instance);
var rowsProp = props.FirstOrDefault(p =>
p.PropertyType != typeof(string) &&
typeof(System.Collections.IEnumerable).IsAssignableFrom(p.PropertyType));
if (rowsProp?.GetValue(body) is not System.Collections.IEnumerable rowsValue)
return response; // no recognizable rows collection — leave Body exactly as returned
var projectedRows = ComputedFieldEvaluator.Project(rowsValue, defs, Logger);
if (ReferenceEquals(projectedRows, rowsValue))
return response;
var envelope = new Dictionary(StringComparer.Ordinal);
foreach (var p in props)
envelope[p.Name] = p.Name == rowsProp.Name ? projectedRows : p.GetValue(body);
response.Body = envelope;
return response;
}
protected virtual void OnConfigure() { }
public override void Configure() => OnConfigure();
// ── Reflection cache — populated once per TData type, reused for every row ──
private static readonly ConcurrentDictionary _propCache = new();
public override async Task HandleAsync(TRequest req, CancellationToken ct)
{
var endpointName = GetType().Name;
using var activity = GB5ActivitySources.Endpoints.StartActivity(endpointName, ActivityKind.Server);
try
{
var loginDto = GetLoginDTOFromRequest(req);
if (loginDto != null && activity != null)
{
activity.DisplayName = $"{endpointName} [{loginDto.DatabaseName ?? "unknown"}]";
activity.SetTag("gb5.db.name", loginDto.DatabaseName ?? "unknown");
activity.SetTag("enduser.id", loginDto.UserCode ?? "unknown");
activity.SetTag("gb5.session.id", loginDto.SessionId?.ToString());
activity.SetTag("gb5.login.event_log_id", loginDto.LoginEventLogId.ToString());
}
if (loginDto != null)
{
loginDto.RequestUrl = $"{HttpContext.Request.Scheme}://{HttpContext.Request.Host}{HttpContext.Request.Path}";
}
if (HttpContext.Request.Headers.TryGetValue("ReportFormat", out var reportFormatHeader))
{
if (!int.TryParse(reportFormatHeader.ToString(), out int reportFormatCode)
|| !ReportFormat.IsValid(reportFormatCode))
{
var hint = $"Valid codes: {ReportFormat.Excel}=Excel, {ReportFormat.Csv}=CSV, " +
$"{ReportFormat.Pdf}=PDF, {ReportFormat.Json}=JSON, " +
$"{ReportFormat.ReportData}=ReportData, {ReportFormat.Html}=HTML, " +
$"{ReportFormat.NdJsonStream}=NdJsonStream, {ReportFormat.PdfNative}=PdfNative.";
Logger.LogWarning("[{Endpoint}] ReportFormat header has invalid value '{Value}'. {Hint}",
endpointName, reportFormatHeader.ToString(), hint);
activity?.SetStatus(ActivityStatusCode.Error, "Invalid ReportFormat value.");
await Send.ResponseAsync(new ResponseStandardDTO
{
Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed,
ErrorBody = $"Invalid ReportFormat value '{reportFormatHeader}'. {hint}",
ErrorInnerException = string.Empty
}, statusCode: 400, cancellation: ct).ConfigureAwait(false);
return;
}
// ── NDJSON stream — write rows directly to response body ────────────
if (reportFormatCode == ReportFormat.NdJsonStream)
{
HttpContext.Response.ContentType = "application/x-ndjson";
int count = 0;
try
{
await foreach (var item in ExecuteReportAsync(req, loginDto!, ct).ConfigureAwait(false))
{
ct.ThrowIfCancellationRequested();
string line = Newtonsoft.Json.JsonConvert.SerializeObject(
new { type = "data", data = item }) + "\n";
await HttpContext.Response.Body.WriteAsync(
Encoding.UTF8.GetBytes(line), ct).ConfigureAwait(false);
await HttpContext.Response.Body.FlushAsync(ct).ConfigureAwait(false);
count++;
}
string footer = Newtonsoft.Json.JsonConvert.SerializeObject(
new { type = "footer", message = "Stream completed", count }) + "\n";
await HttpContext.Response.Body.WriteAsync(
Encoding.UTF8.GetBytes(footer), ct).ConfigureAwait(false);
await HttpContext.Response.Body.FlushAsync(ct).ConfigureAwait(false);
}
catch (OperationCanceledException)
{
string errLine = Newtonsoft.Json.JsonConvert.SerializeObject(
new { type = "error", message = "Stream cancelled" }) + "\n";
await HttpContext.Response.Body.WriteAsync(
Encoding.UTF8.GetBytes(errLine), CancellationToken.None).ConfigureAwait(false);
await HttpContext.Response.Body.FlushAsync(CancellationToken.None).ConfigureAwait(false);
}
MarkActivitySuccess(activity);
return;
}
// ── Export formats (Excel, CSV, PDF, JSON, ReportData, HTML) ─────────
var reportCallingDTO = GetReportCallingDTOFromRequest(req);
if (reportCallingDTO == null)
{
Logger.LogWarning("[{Endpoint}] ReportFormat received but request has no ReportCallingDTO.", endpointName);
activity?.SetStatus(ActivityStatusCode.Error, "ReportCallingDTO missing.");
await Send.ResponseAsync(new ResponseStandardDTO
{
Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed,
ErrorBody = "Report export requires a ReportCallingDTO in the request body.",
ErrorInnerException = string.Empty
}, statusCode: 400, cancellation: ct).ConfigureAwait(false);
return;
}
Logger.LogInformation("[{Endpoint}] Report export — format {Code} ({Name}).",
endpointName, reportFormatCode, ReportFormat.ToFormatName(reportFormatCode));
var reportExport = Resolve();
var rows = ToTypedAsyncRows(ExecuteReportAsync(req, loginDto!, ct), ct);
var exportResult = await reportExport.ExportAsync(reportFormatCode, reportCallingDTO, loginDto!, rows, ct)
.ConfigureAwait(false);
var reportResponse = await GB5Shared.ResponseStandard.Response
.CreateSuccessResponse(exportResult, CacheKeyLevel.NOT_REQUIRED, loginDto!)
.ConfigureAwait(false);
MarkActivitySuccess(activity);
await Send.ResponseAsync(reportResponse, cancellation: ct).ConfigureAwait(false);
// Real traffic sends ReportFormat (including 3=JSON) on every interactive grid
// call, not just true file exports — so this is the path "remember my filter"
// actually needs to fire on, not just the no-header branch below.
await TryFireCriteriaFastSaveAsync(reportCallingDTO, loginDto, reportResponse, ct)
.ConfigureAwait(false);
return;
}
// ── Normal JSON response (no format header) ──────────────────────────
var result = await ExecutePagedAsync(req, loginDto!, ct).ConfigureAwait(false);
result = await ApplyComputedFieldsAutoAsync(result, req, loginDto!, ct).ConfigureAwait(false);
if (result is ResponseStandardDTO dto)
{
var appStatus = dto.Status.ToString();
activity?.SetTag("gb5.app.status", appStatus);
if (dto.Status == GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed
&& !string.IsNullOrEmpty(dto.ErrorBody))
activity?.SetTag("gb5.app.error", dto.ErrorBody);
}
MarkActivitySuccess(activity);
await Send.ResponseAsync(result, cancellation: ct).ConfigureAwait(false);
await TryFireCriteriaFastSaveAsync(GetReportCallingDTOFromRequest(req), loginDto, result, ct)
.ConfigureAwait(false);
}
catch (Grpc.Core.RpcException rpcEx)
{
var (message, httpStatus) = ResolveDaprRpcError(rpcEx, endpointName);
Logger.LogError(rpcEx, "[Dapr-gRPC] {Endpoint} — Status: {StatusCode} — Detail: {Detail}",
endpointName, rpcEx.StatusCode, rpcEx.Status.Detail);
MarkActivityFailure(activity, rpcEx, httpStatus, message);
await Send.ResponseAsync(new ResponseStandardDTO
{
Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed,
ErrorBody = message,
ErrorInnerException = rpcEx.Status.Detail
}, statusCode: httpStatus, cancellation: ct).ConfigureAwait(false);
}
catch (DaprApiException daprEx)
{
var message = $"Dapr component operation failed: {daprEx.Message}";
Logger.LogError(daprEx, "[Dapr-API] {Endpoint} — {Message}", endpointName, daprEx.Message);
MarkActivityFailure(activity, daprEx, 503, message);
await Send.ResponseAsync(new ResponseStandardDTO
{
Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed,
ErrorBody = message,
ErrorInnerException = daprEx.InnerException?.Message ?? string.Empty
}, statusCode: 503, cancellation: ct).ConfigureAwait(false);
}
catch (OperationCanceledException) when (ct.IsCancellationRequested)
{
Logger.LogInformation("[{Endpoint}] Request was cancelled by the client.", endpointName);
activity?.SetTag("gb5.success", false);
activity?.SetTag("http.response.status_code", 499);
activity?.SetStatus(ActivityStatusCode.Ok);
await Send.ResponseAsync(new ResponseStandardDTO
{
Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed,
ErrorBody = "The request was cancelled by the client.",
ErrorInnerException = string.Empty
}, statusCode: 499, cancellation: CancellationToken.None).ConfigureAwait(false);
}
catch (OperationCanceledException timeoutEx)
{
const string message = "The request timed out. Please try again.";
Logger.LogError(timeoutEx, "[Timeout] {Endpoint} — Operation timed out.", endpointName);
MarkActivityFailure(activity, timeoutEx, 504, message);
await Send.ResponseAsync(new ResponseStandardDTO
{
Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed,
ErrorBody = message,
ErrorInnerException = timeoutEx.Message
}, statusCode: 504, cancellation: ct).ConfigureAwait(false);
}
catch (Exception ex)
{
Logger.LogError(ex, "[{Endpoint}] {ExceptionType}: {Message}",
endpointName, ex.GetType().Name, ex.Message);
MarkActivityFailure(activity, ex, 500, ex.Message);
await Send.ResponseAsync(new ResponseStandardDTO
{
Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed,
ErrorBody = ex.Message,
ErrorInnerException = ex.InnerException?.Message ?? string.Empty
}, statusCode: 500, cancellation: ct).ConfigureAwait(false);
}
}
// "Remember my last-used filter" (GB4 parity, Unisoft.FrameworkBLL.CriteriaConfig.CriteriaConfig
// .SaveCriteriaConfig/SaveCriteriaConfigFast) — best-effort, called after the response has
// already been sent so it never affects/delays it. Registered per-host via
// ICriteriaConfigFastSave; a host that hasn't opted in resolves null and this is a no-op.
// Called from BOTH response branches in HandleAsync: real traffic sends a ReportFormat
// header (including 3=JSON) on ordinary interactive grid calls, not just true file exports,
// so the export branch needs this too, not only the no-header branch.
private async Task TryFireCriteriaFastSaveAsync(
ReportCallingDTO? reportCallingDTO, LoginDTO? loginDto, object? result, CancellationToken ct)
{
if (reportCallingDTO == null || loginDto == null) return;
if (result is ResponseStandardDTO dto
&& dto.Status == GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed)
return;
var criteriaConfigFastSave = TryResolve();
if (criteriaConfigFastSave == null) return;
await criteriaConfigFastSave
.SaveCriteriaConfigFastAsync(reportCallingDTO, loginDto, ct)
.ConfigureAwait(false);
}
// ── Efficient row conversion — reflection-cached, no Newtonsoft round-trip ──
private static async IAsyncEnumerable> ToTypedAsyncRows(
IAsyncEnumerable source,
[EnumeratorCancellation] CancellationToken ct = default)
{
var props = _propCache.GetOrAdd(typeof(TData),
t => t.GetProperties(BindingFlags.Public | BindingFlags.Instance));
await foreach (var item in source.WithCancellation(ct).ConfigureAwait(false))
{
ct.ThrowIfCancellationRequested();
yield return props.ToDictionary(p => p.Name, p => (object?)p.GetValue(item));
}
}
// ── Helpers (same as BaseEndpoint) ────────────────────────────────────────
protected LoginDTO? GetLoginDTOFromRequest(TRequest req)
{
if (req == null) return null;
try
{
var loginProp = req.GetType().GetProperty("Login");
if (loginProp == null) return null;
object? loginValue = loginProp.GetValue(req);
if (loginValue == null) return null;
var jsonOptions = new JsonSerializerOptions
{
PropertyNameCaseInsensitive = true,
ReadCommentHandling = JsonCommentHandling.Skip,
AllowTrailingCommas = true
};
jsonOptions.Converters.Add(new SafeStringConverter());
if (loginValue is LoginDTO dto) return dto;
if (loginValue is string str && !string.IsNullOrWhiteSpace(str))
{
try
{
str = System.Text.RegularExpressions.Regex.Replace(
str, "\"([^\"]+?)\\s+\"(\\s*:)",
m => $"\"{m.Groups[1].Value.TrimEnd()}\"{m.Groups[2].Value}");
return JsonSerializer.Deserialize(str, jsonOptions);
}
catch { }
}
if (loginValue is System.Text.Json.JsonElement elem)
{
try { return elem.Deserialize(jsonOptions); }
catch { }
}
try
{
string tempJson = JsonSerializer.Serialize(loginValue, jsonOptions);
return JsonSerializer.Deserialize(tempJson, jsonOptions);
}
catch { }
return null;
}
catch (Exception ex)
{
Logger?.LogWarning("GetLoginDTOFromRequest failed: {Message}", ex.Message);
return null;
}
}
protected virtual ReportCallingDTO? GetReportCallingDTOFromRequest(TRequest req)
{
if (req == null) return null;
try
{
var prop = req.GetType().GetProperty("ReportCallingDTO");
if (prop == null) return null;
object? value = prop.GetValue(req);
if (value == null) return null;
if (value is ReportCallingDTO dto) return dto;
var opts = new JsonSerializerOptions { PropertyNameCaseInsensitive = true };
string json = JsonSerializer.Serialize(value, opts);
return JsonSerializer.Deserialize(json, opts);
}
catch (Exception ex)
{
Logger?.LogWarning("GetReportCallingDTOFromRequest failed: {Message}", ex.Message);
return null;
}
}
private static void MarkActivitySuccess(Activity? activity, int statusCode = 200)
{
activity?.SetTag("gb5.success", true);
activity?.SetTag("http.response.status_code", statusCode);
activity?.SetStatus(ActivityStatusCode.Ok);
}
private static void MarkActivityFailure(Activity? activity, Exception ex, int statusCode, string message)
{
activity?.SetTag("gb5.success", false);
activity?.SetTag("gb5.error.type", ex.GetType().Name);
activity?.SetTag("gb5.error.message", message);
activity?.SetTag("http.response.status_code", statusCode);
activity?.AddException(ex);
activity?.SetStatus(ActivityStatusCode.Error, message);
}
private static (string Message, int HttpStatusCode) ResolveDaprRpcError(
Grpc.Core.RpcException rpcEx, string endpointName)
{
var detail = string.IsNullOrWhiteSpace(rpcEx.Status.Detail)
? string.Empty
: $" Detail: {rpcEx.Status.Detail}.";
return rpcEx.StatusCode switch
{
Grpc.Core.StatusCode.Unavailable =>
($"Dapr sidecar is not reachable from '{endpointName}'.{detail}", 503),
Grpc.Core.StatusCode.DeadlineExceeded =>
($"Dapr operation in '{endpointName}' exceeded its deadline.{detail}", 504),
Grpc.Core.StatusCode.NotFound =>
($"A required Dapr component was not found in '{endpointName}'.{detail}", 503),
Grpc.Core.StatusCode.PermissionDenied =>
($"Access to a Dapr component was denied in '{endpointName}'.{detail}", 403),
Grpc.Core.StatusCode.Unauthenticated =>
($"Dapr authentication failed in '{endpointName}'.{detail}", 401),
Grpc.Core.StatusCode.ResourceExhausted =>
($"Dapr resource limit exceeded in '{endpointName}'.{detail}", 503),
Grpc.Core.StatusCode.FailedPrecondition =>
($"Dapr component precondition failed in '{endpointName}'.{detail}", 503),
Grpc.Core.StatusCode.Internal =>
($"Dapr reported an internal error in '{endpointName}'.{detail}", 500),
Grpc.Core.StatusCode.Cancelled =>
($"The Dapr operation in '{endpointName}' was cancelled.{detail}", 499),
_ =>
($"Dapr communication error ({rpcEx.StatusCode}) in '{endpointName}'.{detail}", 500)
};
}
}
}