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