using System.Collections; using System.Collections.Generic; using System.Diagnostics; using System.Reflection; using System.Runtime.CompilerServices; using System.Text.Json; using System.Text.Json.Serialization; using System.Threading; using Dapr; using Dapr.Client; using FastEndpoints; using GB5Shared.Authorization; using GB5Shared.Cache; using GB5Shared.DaprCache; using GB5Shared.DTO.Framework.Logging; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.Framework.ResponseStandard; using GB5Shared.DTO.PubSub; using GB5Shared.DTO.Report; using GB5Shared.EncryptionHelper; using GB5Shared.Export; using GB5Shared.GB5Constant; using GB5Shared.JsonConverter; using GB5Shared.Resource.EndPointResource; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using GB5Shared.Telemetry; using OpenTelemetry.Trace; using static GB5Shared.GB5Constant.Constant; namespace GB5Shared.FastEndPoint { public abstract class BaseEndpoint : Endpoint where TRequest : notnull { // Cache: TRequest type → whether any [FromBody] property is a collection type. // Populated once per endpoint type at first request — reflection cost paid only once. private static readonly System.Collections.Concurrent.ConcurrentDictionary _isListBodyCache = new(); // Cache: TRequest type → the [FromBody] PropertyInfo (or null when none exists). // Used on WIP approval callbacks to snapshot the body JSON before BLL execution so // EventHandler can restore BLL-regenerated fields (task numbers etc.) to original values. private static readonly System.Collections.Concurrent.ConcurrentDictionary _bodyPropInfoCache = new(); protected new ILogger Logger => Resolve>>(); protected Tracer Tracer => Resolve(); protected DaprClient DaprClient => Resolve(); protected CacheKeyGeneration KeyGenerator => new(); // Add this property to access the environment protected IHostEnvironment Env => Resolve(); // ── Performance: static caches shared across all requests ──────────────────── // JsonSerializerOptions is expensive to construct (~50 KB heap + reflection). // Constructing one per request causes excessive GC pressure at high throughput. private static readonly JsonSerializerOptions _loginDeserializeOptions = BuildLoginOptions(); private static JsonSerializerOptions BuildLoginOptions() { var opts = new JsonSerializerOptions { PropertyNameCaseInsensitive = true, ReadCommentHandling = JsonCommentHandling.Skip, AllowTrailingCommas = true }; opts.Converters.Add(new SafeStringConverter()); return opts; } private static readonly JsonSerializerOptions _reportCallingOptions = new() { PropertyNameCaseInsensitive = true }; // PropertyInfo lookup cache — avoids repeated Type.GetProperty() per request. // Type.GetProperty searches the member list; caching the result eliminates // that overhead on every subsequent call for the same endpoint type. private static readonly System.Collections.Concurrent.ConcurrentDictionary _loginPropCache = new(); private static readonly System.Collections.Concurrent.ConcurrentDictionary _reportPropCache = new(); // ───────────────────────────────────────────────────────────────────────────── // Cache: endpoint Type → its declared [MenuRights] attribute (or null when absent). // Reflection runs once per endpoint type; every subsequent request reads the cached value. private static readonly System.Collections.Concurrent.ConcurrentDictionary _menuRightsAttrCache = new(); // ───────────────────────────────────────────────────────────────────────────── // ── Cache stampede prevention ──────────────────────────────────────────────── // Singleton striped-semaphore pool shared across all endpoint instances. // Injected lazily from DI so tests can substitute a different instance. // See CacheLockRegistry for design rationale (4096 stripes, ~400 KB constant). private static CacheLockRegistry? _lockRegistry; private static CacheLockRegistry LockRegistry => _lockRegistry ??= new CacheLockRegistry(); // ───────────────────────────────────────────────────────────────────────────── protected virtual string? GetCacheKey(TRequest req, LoginDTO LoginDTO) => null; protected virtual bool ShouldInvalidateCache(TRequest req, LoginDTO LoginDTO) => false; protected abstract Task> ExecuteAsync(TRequest req, LoginDTO LoginDTO, CancellationToken ct); protected virtual void OnConfigure() { } /// /// Cache TTL for this endpoint in seconds. Override to customise per-endpoint. /// Use as a convenience when the TTL maps naturally to /// the of the data (e.g., master-data dropdowns): /// /// protected override int GetTtlSeconds() => TtlForLevel(CacheKeyLevel.CLIENT_LEVEL); /// /// Financial / ledger endpoints should return null from /// instead of relying on a short TTL — stale financial data is never acceptable. /// protected virtual int GetTtlSeconds() => 60; /// /// Returns the recommended default TTL in seconds for each /// based on how frequently data at that scope /// typically changes in a running GB5 installation. /// /// | Level | TTL | Rationale | /// |--------------|---------|---------------------------------------------| /// | OVERALL | 3600 s | Language resources, global config — static | /// | DB_SERVER | 1800 s | Server config — rarely changes | /// | CLIENT_LEVEL | 300 s | Tenant master data — changes infrequently | /// | ROLE_LEVEL | 300 s | Role-scoped views — same as client level | /// | USER_LEVEL | 60 s | User-specific data — may change per action | /// | (default) | 60 s | Session/unknown — conservative | /// protected static int TtlForLevel(int cacheKeyLevel) => cacheKeyLevel switch { CacheKeyLevel.OVERALL => 3600, CacheKeyLevel.DB_SERVER => 1800, CacheKeyLevel.CLIENT_LEVEL => 300, CacheKeyLevel.ROLE_LEVEL => 300, CacheKeyLevel.USER_LEVEL => 60, _ => 60 }; public override void Configure() => OnConfigure(); // Cache: endpoint Type → resolved module name (e.g. "CRMSL" assembly → "CRM"). // Since Host consolidation, several business modules share one process/service.name // in Zipkin, so this tag is the only way to tell them apart in a trace. private static readonly System.Collections.Concurrent.ConcurrentDictionary _moduleNameCache = new(); private string ResolveModuleName() => _moduleNameCache.GetOrAdd(GetType(), t => { var asmName = t.Assembly.GetName().Name ?? "GB5"; return asmName.Length > 2 && asmName.EndsWith("SL", StringComparison.Ordinal) ? asmName[..^2] : asmName; }); // A client that disconnects while the server is still writing the response body // (browser tab closed, load-balancer probe timeout, mobile network drop) makes // Kestrel's HttpResponsePipeWriter.FlushAsync throw OperationCanceledException from // INSIDE whichever Send.ResponseAsync call was in flight — including calls made from // HandleAsyncCore's own catch blocks while reporting an unrelated original error. An // exception thrown inside a catch block is not caught by that same try statement's // other catch clauses, so this escaped HandleAsyncCore entirely and surfaced as an // unhandled exception (noisy 500-shaped log/trace) for what is actually a completely // normal, un-actionable client disconnect. This outer wrapper is the single place // that catches it regardless of which inner path threw it. public override async Task HandleAsync(TRequest req, CancellationToken ct) { try { await HandleAsyncCore(req, ct); } catch (Exception ex) when (IsClientDisconnect(ex)) { Logger.LogInformation( "[{Endpoint}] Client disconnected before the response could be delivered ({ExceptionType}) — nothing to do.", GetType().Name, ex.GetType().Name); } } private static bool IsClientDisconnect(Exception ex) => ex is OperationCanceledException or ObjectDisposedException or Microsoft.AspNetCore.Connections.ConnectionResetException; private async Task HandleAsyncCore(TRequest req, CancellationToken ct) { var endpointName = GetType().Name; using var activity = GB5ActivitySources.Endpoints.StartActivity(endpointName, ActivityKind.Server); // Tag + prefix the span with the owning module unconditionally (not just when a // LoginDTO is present) so every trace remains attributable to a module even though // multiple modules now share one Host process and one Zipkin service name. var moduleName = ResolveModuleName(); if (activity != null) { // Lead with "{HTTP method} {route}" (e.g. "POST /Analysis/AnalysisDynamicOutput") // rather than the bare endpoint class name — this is the name a professional // reading Zipkin actually recognises as "which API is this", and it's also the // name most likely to end up representing a long-running trace in Zipkin's list // view if that view ends up picking this span (or a descendant) as the row's // "Root" label before the enclosing HTTP-request span itself has closed/exported. var routeLabel = $"{HttpContext.Request.Method} {HttpContext.Request.Path}"; activity.SetTag("gb5.module", moduleName); activity.SetTag("http.route", routeLabel); activity.DisplayName = $"{routeLabel} [{moduleName}]"; } // Stampede-prevention state — declared at method scope (outside try) so that // the finally block below can release the semaphore on every exit path, // including exceptions thrown inside catch blocks. SemaphoreSlim? stampedeSem = null; bool stampedeAcquired = false; try { Logger.LogInformation(EndPoint.Log_ReceivedRequest, req); var loginDto = await GetLoginDTOFromRequestAsync(req); // Enrich the span with GB5 request context so every trace carries // the tenant database, user identity, and session for fast triage. if (loginDto != null && activity != null) { var dbName = loginDto.DatabaseName ?? "unknown"; var userCode = loginDto.UserCode ?? "unknown"; var sessionId = loginDto.SessionId; var loginEventLogId = loginDto.LoginEventLogId; activity.DisplayName = $"{HttpContext.Request.Method} {HttpContext.Request.Path} [{moduleName}] [{dbName}]"; activity.SetTag("gb5.db.name", dbName); activity.SetTag("enduser.id", userCode); activity.SetTag("gb5.session.id", sessionId?.ToString()); activity.SetTag("gb5.login.event_log_id", loginEventLogId.ToString()); activity.SetTag("gb5.client.id", loginDto.ClientId.ToString()); activity.SetTag("gb5.user.id", loginDto.UserId.ToString()); } // ── Server-side request context: captured once, used by WIP and other framework features ── if (loginDto != null) { // Full URL of this endpoint — used as WIP callback when CALLBACKENDPOINT not configured. // Includes PathBase (set by UseGatewayPrefixForwarding from X-GB5-Gateway-Prefix when // reached through the GB5Build gateway under a module alias, e.g. /tms/) so this stays // a genuinely self-referencing URL — Request.Path alone excludes it by definition. loginDto.RequestUrl = $"{HttpContext.Request.Scheme}://{HttpContext.Request.Host}{HttpContext.Request.PathBase}{HttpContext.Request.Path}"; // WIP approval callback: detect replay request from HttpWipApprovalDispatcher if (HttpContext.Request.Headers.TryGetValue("X-Wip-Approval", out var wipHeader) && int.TryParse(wipHeader.ToString(), out int wipApprovalId) && wipApprovalId > 0) { loginDto.WipApprovalId = wipApprovalId; // Snapshot the request body BEFORE ExecuteAsync so EventHandler can // restore BLL-regenerated fields (task numbers, invoice numbers, etc.) // to the original WIP-submission values before persisting the entity. // PropertyInfo is cached per TRequest type — reflection runs only once. var bodyPropInfo = _bodyPropInfoCache.GetOrAdd(typeof(TRequest), reqType => reqType .GetProperties(System.Reflection.BindingFlags.Public | System.Reflection.BindingFlags.Instance) .FirstOrDefault(p => p.GetCustomAttributes(true) .Any(a => a.GetType().Name == "FromBodyAttribute"))); if (bodyPropInfo != null) { var bodyVal = bodyPropInfo.GetValue(req); if (bodyVal != null) loginDto.WipOriginalBodyJson = JsonSerializer.Serialize(bodyVal); } } // Detect if [FromBody] parameter is a collection type (List, IEnumerable, etc.). // Result is cached per TRequest type so reflection only runs once per endpoint class. // EventHandler reads login.IsBodyArray to store WIP DATAJSON as "[{...}]" — ensuring // the dispatcher sends the correct body format on approval callback automatically. loginDto.IsBodyArray = _isListBodyCache.GetOrAdd(typeof(TRequest), requestType => { var bodyProp = requestType .GetProperties(System.Reflection.BindingFlags.Public | System.Reflection.BindingFlags.Instance) .FirstOrDefault(p => p.GetCustomAttributes(true) .Any(a => a.GetType().Name == "FromBodyAttribute")); if (bodyProp == null) return false; var pt = bodyProp.PropertyType; return pt.IsArray || (pt.IsGenericType && ( pt.GetGenericTypeDefinition() == typeof(List<>) || pt.GetGenericTypeDefinition() == typeof(IList<>) || pt.GetGenericTypeDefinition() == typeof(IEnumerable<>) || pt.GetGenericTypeDefinition() == typeof(ICollection<>) || pt.GetGenericTypeDefinition() == typeof(IReadOnlyList<>))); }); } // ── MenuRights authorization gate ────────────────────────────────────────────── // Opt-in: only endpoint classes decorated with [MenuRights(menuCode, operation)] // are checked. Denial returns 403 before ExecuteAsync (and any DAL/BLL work) runs. // See GB5Shared.Authorization.MenuRightsBLL for the ported check logic. var menuRightsAttr = _menuRightsAttrCache.GetOrAdd(GetType(), static t => t.GetCustomAttribute()); if (menuRightsAttr != null) { if (loginDto == null) { activity?.SetTag("gb5.authz.denied", true); MarkActivityFailure(activity, new UnauthorizedAccessException("No login context."), 401, "Authorization requires a valid login context."); await Send.ResponseAsync((TResponse)(object)new ResponseStandardDTO { Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed, ErrorBody = "Authorization requires a valid login context.", ErrorInnerException = string.Empty }, statusCode: 401, cancellation: ct); return; } var menuRightsBLL = Resolve(); bool menuRightsAllowed = await menuRightsBLL.IsAllowedAsync( menuRightsAttr.MenuCode, menuRightsAttr.Operation, loginDto, ct).ConfigureAwait(false); if (!menuRightsAllowed) { activity?.SetTag("gb5.authz.denied", true); activity?.SetTag("gb5.authz.menu_code", menuRightsAttr.MenuCode); activity?.SetTag("gb5.authz.operation", menuRightsAttr.Operation.ToString()); MarkActivityFailure(activity, new UnauthorizedAccessException($"Denied: {menuRightsAttr.MenuCode}/{menuRightsAttr.Operation}"), 403, "You are not authorized to perform this operation."); await Send.ResponseAsync((TResponse)(object)new ResponseStandardDTO { Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed, ErrorBody = "You are not authorized to perform this operation. Contact your administrator.", ErrorInnerException = string.Empty }, statusCode: 403, cancellation: ct); return; } } // ── End MenuRights gate ───────────────────────────────────────────────────────── string? cacheKey = GetCacheKey(req, loginDto!); TResponse? finalResponse; if (!string.IsNullOrWhiteSpace(cacheKey)) { // Invalidate cache if needed if (ShouldInvalidateCache(req, loginDto!)) { try { Logger.LogInformation("Invalidating cache key: {CacheKey}", cacheKey); await DaprClient.DeleteStateAsync("statestore", cacheKey, cancellationToken: ct); } catch (Exception cacheEx) { Logger.LogWarning(cacheEx, "Cache invalidation skipped — Dapr unavailable for key: {CacheKey}", cacheKey); } } else { // ── Stampede-safe cache read (Thundering Herd Prevention) ────────────── // Acquire the per-key semaphore BEFORE reading the cache. // • acquired=true → we are the designated reader/writer for this key. // Read the cache (double-check pattern in case another request just // populated it while we waited). On confirmed miss, execute the // query; after the query completes, write the result and release. // • acquired=false → timed out (200 ms); fall through to a live query // without blocking — graceful degradation under extreme load. // // Semaphore lifetime: acquired here, released in one of three places: // a) On cache HIT → released immediately before early return // b) After cache WRITE (further down) → released in its finally block // c) If cacheKey is null (non-caching endpoint) → never acquired // ───────────────────────────────────────────────────────────────────── // Acquire BEFORE cache read — we are the designated writer if we get the lock. stampedeSem = LockRegistry.GetLock(cacheKey); stampedeAcquired = await stampedeSem .WaitAsync(CacheLockRegistry.TimeoutMs, ct) .ConfigureAwait(false); if (!stampedeAcquired) stampedeSem = null; // not held; no cleanup needed // Cache read (double-check when we hold the lock) try { string? cachedData = await DaprClient.GetStateAsync("statestore", cacheKey, cancellationToken: ct); if (!string.IsNullOrEmpty(cachedData)) { Logger.LogInformation(EndPoint.Log_CacheHit, cacheKey); TResponse? cachedResponse = DeserializeResponse(cachedData); if (cachedResponse != null) { // Cache hit — release the lock before returning so the // next waiter can proceed immediately. if (stampedeAcquired) { stampedeSem!.Release(); stampedeSem = null; stampedeAcquired = false; } activity?.SetTag("gb5.cache.hit", true); MarkActivitySuccess(activity); await Send.ResponseAsync(cachedResponse, cancellation: ct); return; } } } catch (Exception cacheEx) { Logger.LogWarning(cacheEx, "Cache read skipped — Dapr unavailable for key: {CacheKey}", cacheKey); } // Cache miss confirmed — ExecuteAsync runs below; semaphore held by lock owner. } } // Execute fresh var freshResponse = await ExecuteAsync(req, loginDto!, ct); // ── WIP response override ────────────────────────────────────────────────────── // When EventHandler queued the entity for approval (WIP), the BLL returns a // misleading "Details saved successfully..." message. Replace it with a clear // pending-approval notice before the response reaches the client. if (loginDto?.WipTriggered == true && freshResponse is ResponseStandardDTO wipResponse) { wipResponse.Body = "Submitted for approval. Your request will be saved after the approval is completed."; activity?.SetTag("gb5.wip.response_overridden", true); } // ── Report export ────────────────────────────────────────────────────────────── // When the client sends H-ReportFormat, the BaseEndpoint handles all export logic. // It uses the data already fetched by ExecuteAsync (freshResponse.Body) so there is // no second trip to the database. No header = no report, execution continues normally. if (HttpContext.Request.Headers.TryGetValue("ReportFormat", out var reportFormatHeader)) { // ── Validate the format code ─────────────────────────────────────────────── 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."; Logger.LogWarning("[{Endpoint}] H-ReportFormat header has invalid value '{Value}'. {Hint}", GetType().Name, reportFormatHeader.ToString(), hint); activity?.SetStatus(ActivityStatusCode.Error, "Invalid H-ReportFormat value."); await Send.ResponseAsync((TResponse)(object)new ResponseStandardDTO { Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed, ErrorBody = $"Invalid H-ReportFormat value '{reportFormatHeader}'. {hint}", ErrorInnerException = string.Empty }, statusCode: 400, cancellation: ct); return; } // ── Require ReportCallingDTO on the request ──────────────────────────────── var reportCallingDTO = GetReportCallingDTOFromRequest(req); if (reportCallingDTO == null) { Logger.LogWarning("[{Endpoint}] ReportFormat received but the request has no ReportCallingDTO property. Report export requires a ReportCallingDTO in the request body.", GetType().Name); activity?.SetStatus(ActivityStatusCode.Error, "ReportCallingDTO missing from request."); await Send.ResponseAsync((TResponse)(object)new ResponseStandardDTO { Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed, ErrorBody = "Report export requires a ReportCallingDTO in the request body. This endpoint does not include one, or the value was null.", ErrorInnerException = string.Empty }, statusCode: 400, cancellation: ct); return; } // ── Run export using data already in freshResponse.Body ──────────────────── Logger.LogInformation("[{Endpoint}] Report export — format {Code} ({Name}).", GetType().Name, reportFormatCode, ReportFormat.ToFormatName(reportFormatCode)); var reportExport = Resolve(); var rows = ToAsyncRows(freshResponse.Body, ct); var exportResult = await reportExport.ExportAsync(reportFormatCode, reportCallingDTO, loginDto!, rows, ct); var reportResponse = await GB5Shared.ResponseStandard.Response.CreateSuccessResponse( exportResult, CacheKeyLevel.NOT_REQUIRED, loginDto!); MarkActivitySuccess(activity); await Send.ResponseAsync((TResponse)(object)reportResponse, cancellation: ct); return; } // ── End report export ────────────────────────────────────────────────────────── // Assign cache info for fresh data if (!string.IsNullOrWhiteSpace(cacheKey) && freshResponse is ResponseStandardDTO freshDto) { freshDto.CacheKey = cacheKey; freshDto.CacheStatus = Constant.CacheStatus.NewCacheResponse; } finalResponse = (TResponse)(object)freshResponse; // Tag the application-level result so Zipkin shows success/failure // independently of the HTTP status code. bool appFailed = false; if (freshResponse is GB5Shared.DTO.Framework.ResponseStandard.ResponseStandardDTO dto) { activity?.SetTag("gb5.app.status", dto.Status.ToString()); if (dto.Status == GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed) { appFailed = true; var errMsg = dto.ErrorBody ?? string.Empty; if (!string.IsNullOrEmpty(errMsg)) { activity?.SetTag("gb5.app.error", errMsg); activity?.SetStatus(ActivityStatusCode.Error, errMsg); } // Surface Body.GeneralErrors when ErrorBody is empty (validation path) try { var bodyJson = System.Text.Json.JsonSerializer.Serialize(dto.Body); using var doc = System.Text.Json.JsonDocument.Parse(bodyJson); if (doc.RootElement.TryGetProperty("GeneralErrors", out var arr) && arr.ValueKind == System.Text.Json.JsonValueKind.Array) { var msgs = string.Join("; ", arr.EnumerateArray() .Select(e => e.GetString()) .Where(s => !string.IsNullOrEmpty(s))); if (!string.IsNullOrEmpty(msgs)) { activity?.SetTag("gb5.app.general_errors", msgs); activity?.SetStatus(ActivityStatusCode.Error, msgs); } } } catch { /* best effort — never fail the response */ } } } // Save to cache if applicable. // Only the stampede-lock holder (stampedeAcquired=true) writes to the cache, // preventing N concurrent misses from all issuing N identical write operations. // If stampedeAcquired=false (timed out) we still query but skip the write — // the lock holder will write; subsequent readers will get the cached result. if (!string.IsNullOrWhiteSpace(cacheKey) && stampedeAcquired) { try { string json = System.Text.Json.JsonSerializer.Serialize(finalResponse); var metadata = new Dictionary { { "ttlInSeconds", GetTtlSeconds().ToString() } }; await DaprClient.SaveStateAsync("statestore", cacheKey, json, metadata: metadata, cancellationToken: ct); } catch (Exception cacheEx) { Logger.LogWarning(cacheEx, "Cache write skipped — Dapr unavailable for key: {CacheKey}", cacheKey); } finally { // Release AFTER write completes so the next request reads a // populated cache entry rather than racing with an in-progress write. stampedeSem?.Release(); stampedeSem = null; stampedeAcquired = false; } } if (!appFailed) MarkActivitySuccess(activity); await Send.ResponseAsync(finalResponse, cancellation: ct); } catch (Grpc.Core.RpcException rpcEx) { var (message, httpStatus) = ResolveDaprRpcError(rpcEx, GetType().Name); Logger.LogError(rpcEx, "[Dapr-gRPC] {Endpoint} — Status: {StatusCode} — Detail: {Detail}", GetType().Name, rpcEx.StatusCode, rpcEx.Status.Detail); MarkActivityFailure(activity, rpcEx, httpStatus, message); if (HttpContext.Response.HasStarted) { Logger.LogWarning("[{Endpoint}] Response already started — cannot send error body for {ExceptionType}.", GetType().Name, rpcEx.GetType().Name); return; } await Send.ResponseAsync((TResponse)(object)new ResponseStandardDTO { Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed, ErrorBody = message, ErrorInnerException = rpcEx.Status.Detail }, statusCode: httpStatus, cancellation: ct); } catch (DaprApiException daprEx) { var message = $"Dapr component operation failed: {daprEx.Message}"; Logger.LogError(daprEx, "[Dapr-API] {Endpoint} — {Message}", GetType().Name, daprEx.Message); MarkActivityFailure(activity, daprEx, 503, message); if (HttpContext.Response.HasStarted) { Logger.LogWarning("[{Endpoint}] Response already started — cannot send error body for {ExceptionType}.", GetType().Name, daprEx.GetType().Name); return; } await Send.ResponseAsync((TResponse)(object)new ResponseStandardDTO { Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed, ErrorBody = message, ErrorInnerException = daprEx.InnerException?.Message ?? string.Empty }, statusCode: 503, cancellation: ct); } catch (OperationCanceledException) when (ct.IsCancellationRequested) { Logger.LogInformation("[{Endpoint}] Request was cancelled by the client.", GetType().Name); activity?.SetTag("gb5.success", false); activity?.SetTag("http.response.status_code", 499); activity?.SetStatus(ActivityStatusCode.Ok); if (HttpContext.Response.HasStarted) return; await Send.ResponseAsync((TResponse)(object)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); } catch (OperationCanceledException timeoutEx) { const string message = "The request timed out while waiting for a Dapr response. The state store or service may be under load. Please try again."; Logger.LogError(timeoutEx, "[Dapr-Timeout] {Endpoint} — Operation timed out internally.", GetType().Name); MarkActivityFailure(activity, timeoutEx, 504, message); if (HttpContext.Response.HasStarted) return; await Send.ResponseAsync((TResponse)(object)new ResponseStandardDTO { Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed, ErrorBody = message, ErrorInnerException = timeoutEx.Message }, statusCode: 504, cancellation: ct); } // ArgumentException (and its ArgumentNullException/ArgumentOutOfRangeException // subclasses) is the established BLL-wide idiom for "the caller sent invalid // input" — throw new ArgumentException("X is required.") appears 500+ times // across FrameworkBLL/GB5Shared. Without this catch, every one of those falls // into the generic handler below and reports HTTP 500 (server fault) for what // is actually a 400 (client fault) — misleading API consumers and polluting // server-error monitoring/alerting with client mistakes. Must be caught before // the generic Exception handler. catch (ArgumentException argEx) { Logger.LogWarning(argEx, "[{Endpoint}] Invalid input: {Message}", GetType().Name, argEx.Message); MarkActivityFailure(activity, argEx, 400, argEx.Message); await Send.ResponseAsync((TResponse)(object)new ResponseStandardDTO { Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed, ErrorBody = argEx.Message, ErrorInnerException = string.Empty }, statusCode: 400, cancellation: ct); } catch (Exception ex) { Logger.LogError(ex, "[{Endpoint}] {ExceptionType}: {Message}", GetType().Name, ex.GetType().Name, ex.Message); MarkActivityFailure(activity, ex, 500, ex.Message); // A streaming endpoint (e.g. DownloadCertificate) may throw AFTER it already // wrote response headers via SendStreamAsync — trying to set StatusCode at that // point throws InvalidOperationException ("StatusCode cannot be set because the // response has already started"), masking the real exception. Just log and stop. if (HttpContext.Response.HasStarted) { Logger.LogWarning("[{Endpoint}] Response already started — cannot send error body for {ExceptionType}: {Message}.", GetType().Name, ex.GetType().Name, ex.Message); return; } await Send.ResponseAsync((TResponse)(object)new ResponseStandardDTO { Status = GB5Shared.DTO.Framework.Enum.FrameworkEnumDTO.ResponseStatus.Failed, ErrorBody = ex.Message, ErrorInnerException = ex.InnerException?.Message ?? string.Empty }, statusCode: 500, cancellation: ct); } finally { // Safety net: if an exception propagated through any catch block before the // normal cache-write finally had a chance to run, stampedeAcquired is still // true and stampedeSem still holds an unreleased permit. Release it here so // the semaphore slot is never permanently stuck after an unexpected error. if (stampedeAcquired) stampedeSem?.Release(); } } // ── Activity helpers — called on every endpoint execution path ──────────────── 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); } // ───────────────────────────────────────────────────────────────────────────── // Dual-mode bridge (GB5 Repo-Wide Authentication Hardening plan): if a registered // AddJwtBearer scheme already validated this request's token (HttpContext.User carries // real claims), prefer building LoginDTO from those server-verified claims via // IKeycloakLoginDTOResolver. TryResolve() returns null rather than throwing when // nothing is registered — so today, with no host yet wiring a Keycloak scheme and no // resolver implementation registered, this branch is always skipped and behavior is // byte-for-byte identical to before. Any failure here (no resolver registered, resolver // itself throws, resolver returns null because it couldn't map the identity) falls // through to the existing trusted-header path below, never abandons LoginDTO resolution // entirely. protected async Task GetLoginDTOFromRequestAsync(TRequest req) { if (HttpContext?.User?.Identity?.IsAuthenticated == true) { try { var keycloakResolver = TryResolve(); if (keycloakResolver != null) { var claimsLogin = await keycloakResolver.ResolveAsync(HttpContext.User, HttpContext.RequestAborted); if (claimsLogin != null) return claimsLogin; } } catch (Exception ex) { Logger?.LogWarning(ex, "GetLoginDTOFromRequestAsync: Keycloak claims resolution failed — falling back to trusted-header path"); } } return GetLoginDTOFromRequest(req); } protected LoginDTO? GetLoginDTOFromRequest(TRequest req) { if (req == null) return null; try { // Cached PropertyInfo lookup — avoids Type.GetProperty on every request var loginProp = _loginPropCache.GetOrAdd(req.GetType(), static t => t.GetProperty("Login")); if (loginProp == null) return null; object? loginValue = loginProp.GetValue(req); if (loginValue == null) return null; // Case 1: Already a DTO if (loginValue is LoginDTO dto) return dto; // Case 2: String JSON — shared static options, no per-request allocation if (loginValue is string str && !string.IsNullOrWhiteSpace(str)) { try { // Normalize JSON property names: strip trailing whitespace before the colon // e.g. "TempBaseURI " : → "TempBaseURI" : str = System.Text.RegularExpressions.Regex.Replace( str, "\"([^\"]+?)\\s+\"(\\s*:)", m => $"\"{m.Groups[1].Value.TrimEnd()}\"{m.Groups[2].Value}"); return JsonSerializer.Deserialize(str, _loginDeserializeOptions); } catch { // continue to fallback } } // Case 3: JsonElement if (loginValue is JsonElement elem) { try { return elem.Deserialize(_loginDeserializeOptions); } catch { // continue to fallback } } // Case 4: fallback → serialize → deserialize try { string tempJson = JsonSerializer.Serialize(loginValue, _loginDeserializeOptions); return JsonSerializer.Deserialize(tempJson, _loginDeserializeOptions); } catch { // ignore } return null; } catch (Exception ex) { Logger?.LogWarning("GetLoginDTOFromRequest failed: {Message}", ex.Message); return null; } } protected virtual TResponse? DeserializeResponse(string json) { try { return System.Text.Json.JsonSerializer.Deserialize(json); } catch { Logger.LogWarning(EndPoint.Log_DeserializationFail, typeof(TResponse).Name); return default; } } protected virtual async Task PublishEventsAsync(string PublishType, string PublishTopic, object Data, CancellationToken ct) { using var publishActivity = GB5ActivitySources.DaprPublish.StartActivity( $"publish {PublishTopic}", ActivityKind.Producer); publishActivity?.SetTag("messaging.system", "dapr"); publishActivity?.SetTag("messaging.destination", PublishTopic); publishActivity?.SetTag("messaging.pubsub.name", PublishType); publishActivity?.SetTag("messaging.operation", "publish"); try { await DaprClient.PublishEventAsync(PublishType, PublishTopic, Data, cancellationToken: ct); publishActivity?.SetStatus(ActivityStatusCode.Ok); Logger.LogInformation(EndPoint.Log_PublishEventSuccess, Data); } catch (Grpc.Core.RpcException rpcEx) { publishActivity?.AddException(rpcEx); publishActivity?.SetStatus(ActivityStatusCode.Error, rpcEx.Status.Detail); Logger.LogError(rpcEx, "[Dapr-PubSub] Failed to publish to topic '{Topic}' on pubsub '{PubSubName}' — gRPC {StatusCode}: {Detail}", PublishTopic, PublishType, rpcEx.StatusCode, rpcEx.Status.Detail); throw; } catch (DaprApiException daprEx) { publishActivity?.AddException(daprEx); publishActivity?.SetStatus(ActivityStatusCode.Error, daprEx.Message); Logger.LogError(daprEx, "[Dapr-PubSub] Dapr API error publishing to topic '{Topic}' on pubsub '{PubSubName}' — {Message}", PublishTopic, PublishType, daprEx.Message); throw; } catch (Exception ex) { publishActivity?.AddException(ex); publishActivity?.SetStatus(ActivityStatusCode.Error, ex.Message); Logger.LogError(ex, "[Dapr-PubSub] Unexpected error publishing to topic '{Topic}' on pubsub '{PubSubName}' — {ExceptionType}: {Message}", PublishTopic, PublishType, ex.GetType().Name, ex.Message); throw; } } protected virtual PublishDTO? GetTopicName(PublishDTO? PublishDTO) { return null; } /// /// Reads the ReportCallingDTO property from the request object via reflection. /// Returns null when the property does not exist or its value is null — the caller /// must treat that as "report export not supported by this endpoint". /// protected virtual ReportCallingDTO? GetReportCallingDTOFromRequest(TRequest req) { if (req == null) return null; try { // Cached PropertyInfo lookup — avoids Type.GetProperty on every request var prop = _reportPropCache.GetOrAdd(req.GetType(), static t => t.GetProperty("ReportCallingDTO")); if (prop == null) return null; object? value = prop.GetValue(req); if (value == null) return null; if (value is ReportCallingDTO dto) return dto; // Fallback: round-trip through JSON — shared static options, no per-request allocation string json = JsonSerializer.Serialize(value, _reportCallingOptions); return JsonSerializer.Deserialize(json, _reportCallingOptions); } catch (Exception ex) { Logger?.LogWarning("GetReportCallingDTOFromRequest failed: {Message}", ex.Message); return null; } } /// /// Converts any object into a stream of rows /// that the export methods expect. /// Each item is round-tripped through Newtonsoft JSON to ensure consistent key/value shapes. /// private static async IAsyncEnumerable> ToAsyncRows( object? raw, [EnumeratorCancellation] CancellationToken ct = default) { if (raw == null) yield break; // A pre-serialized JSON string must be parsed before the IEnumerable branch // would iterate it character-by-character. if (raw is string jsonStr) { if (string.IsNullOrWhiteSpace(jsonStr)) yield break; Newtonsoft.Json.Linq.JToken? token = null; bool parseOk = false; try { token = Newtonsoft.Json.Linq.JToken.Parse(jsonStr); parseOk = true; } catch { /* fall through — treat the whole string as a single row */ } if (!parseOk) { yield return RowFromObject(jsonStr); yield break; } if (token is Newtonsoft.Json.Linq.JArray arr) { foreach (var jItem in arr) { ct.ThrowIfCancellationRequested(); yield return RowFromObject(jItem); } } else { yield return RowFromObject(token!); } yield break; } if (raw is not IEnumerable items) { yield return RowFromObject(raw); yield break; } foreach (var item in items) { ct.ThrowIfCancellationRequested(); yield return RowFromObject(item); } } private static Dictionary RowFromObject(object item) { // Fast paths — avoid a JSON round-trip for the most common in-process types. // DynamicOutput and other DAL methods already project rows into Dictionary; // round-tripping through JSON just to cast the value type is wasteful for large exports. if (item is Dictionary exact) return exact; if (item is Dictionary strObj) { var fast = new Dictionary(strObj.Count, StringComparer.OrdinalIgnoreCase); foreach (var kv in strObj) fast[kv.Key] = kv.Value; return fast; } var json = Newtonsoft.Json.JsonConvert.SerializeObject(item); return Newtonsoft.Json.JsonConvert.DeserializeObject>(json) ?? new Dictionary(); } 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}'. Ensure the Dapr runtime is running and the sidecar is attached to this service.{detail}", 503), Grpc.Core.StatusCode.DeadlineExceeded => ($"Dapr operation in '{endpointName}' exceeded its deadline. The state store or message broker did not respond in time.{detail}", 504), Grpc.Core.StatusCode.NotFound => ($"A required Dapr component was not found in '{endpointName}'. Verify the component name and configuration.{detail}", 503), Grpc.Core.StatusCode.PermissionDenied => ($"Access to a Dapr component was denied in '{endpointName}'. Check the component's scopes and access policy.{detail}", 403), Grpc.Core.StatusCode.Unauthenticated => ($"Dapr authentication failed in '{endpointName}'. The service token or API token is missing or invalid.{detail}", 401), Grpc.Core.StatusCode.ResourceExhausted => ($"Dapr resource limit exceeded in '{endpointName}'. The state store or pub/sub broker is overloaded or rate-limiting requests.{detail}", 503), Grpc.Core.StatusCode.FailedPrecondition => ($"Dapr component precondition failed in '{endpointName}'. The component may not be initialised or is in an invalid state.{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) }; } } }