using System.Collections.Generic; using System.Data.Common; using System.Diagnostics; using System.Reflection; using System.Text.Json; using System.Text.Json.Nodes; using Dapr.Client; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.PubSub; using GB5Shared.DTO.Qualifier; using GB5Shared.GB5Library.Qualifier; using GB5Shared.PubSub.OutBox; using GB5Shared.QueryExecutor; using GB5Shared.Telemetry; using GB5Shared.WorkFlow.WorkFlowEngine; using static GB5Shared.DTO.WorkFlow.ContractsDTO; using static GB5Shared.GB5Constant.Constant; using static GB5Shared.GB5Constant.Constant.OUTBOXSTATUS; namespace GB5Shared.EntityHandler { public class BaseEntityAppService { private readonly IQualifierFacade _qualifierFacade; private readonly IWorkFlowEngine _workflowEngine; private readonly IOutBox _outbox; private readonly IQueryExecutor _queryExecutor; private readonly DaprClient _dapr; // Single shared serialiser options: compact, property names preserved as-is (PascalCase), // nulls omitted. Used for LoginJson and ContextJson only — Payload uses the default // serialiser (same as the rest of the codebase) to stay consistent with existing subscribers. private static readonly JsonSerializerOptions _publishSerialiserOptions = new() { WriteIndented = false, DefaultIgnoreCondition = System.Text.Json.Serialization.JsonIgnoreCondition.WhenWritingNull }; public BaseEntityAppService( IQualifierFacade qualifierFacade, IWorkFlowEngine workflowEngine, IOutBox outbox, IQueryExecutor queryExecutor, DaprClient dapr) { _qualifierFacade = qualifierFacade; _workflowEngine = workflowEngine; _outbox = outbox; _queryExecutor = queryExecutor; _dapr = dapr; } /// /// Universal Save Execution Pipeline. /// /// Steps: BeginTx → Validate → EnrichForWorkflow → WorkflowCheck → PrePersist /// → Persist → StartWorkflow → PostPersist → EnrichForEvent /// → DirectDaprPublish (fire-and-forget) → OutboxWrite → Commit. /// /// /// Every step is traced under the GB5.EntitySave ActivitySource. /// The published carries the full DTO payload, /// full serialised LoginDTO, and the enriched pipeline context so subscribers /// have all data without additional round-trips. /// /// /// /// MBIZTRANSACTIONCLASS.BIZTRANSACTIONCLASSID. Defaults to -1 (wildcard). /// /// /// MBIZTRANSACTION.BIZTRANSACTIONID. Defaults to -1 (wildcard). /// public async Task ExecuteSaveAsync( int entityId, int eventTypeId, TDto dto, LoginDTO login, Func> persistFunc, Dictionary? Facts = null, int bizTransactionClassId = -1, int bizTransactionId = -1, DbTransaction? externalTransaction = null, bool isNewEntity = false, Guid? headerRowGuid = null) { // ── Top-level span — covers the entire save pipeline ────────────────── using var saveActivity = GB5ActivitySources.EntitySave.StartActivity( "entity.save", ActivityKind.Internal); saveActivity?.SetTag("gb5.entity.id", entityId); saveActivity?.SetTag("gb5.event_type.id", eventTypeId); saveActivity?.SetTag("gb5.biz_tx.class_id", bizTransactionClassId); saveActivity?.SetTag("gb5.biz_tx.id", bizTransactionId); saveActivity?.SetTag("gb5.entity.is_new", isNewEntity); saveActivity?.SetTag("gb5.tenant.id", login.ClientId); saveActivity?.SetTag("gb5.user.id", login.UserId); saveActivity?.SetTag("gb5.ou.id", login.WorkOUId); saveActivity?.SetTag("gb5.entity.dto_type", typeof(TDto).Name); saveActivity?.SetTag("gb5.entity.has_external_tx", externalTransaction != null); DbTransaction? tx = externalTransaction; var isExternalTx = tx != null; try { // ── Transaction ─────────────────────────────────────────────────── if (!isExternalTx) { tx = await _queryExecutor.BeginTransactionAsync(login); saveActivity?.AddEvent(new ActivityEvent("entity.save.tx_begin")); } // Inject WIP bypass flag when this call is an approved WIP callback. // login.WipApprovalId is set by BaseEndpoint from the X-Wip-Approval header. if (login.WipApprovalId > 0) { Facts = new Dictionary(Facts ?? new Dictionary()) { [IWorkFlowEngine.WipApprovalFactKey] = true }; saveActivity?.SetTag("gb5.wip.approval_id", login.WipApprovalId); } // ctx.Facts = only caller-supplied overrides + WIP flag (minimal). // DB_ENRICH qualifiers reflect ctx.Data directly when they run — // no upfront DTO reflection needed here. var ctx = new ExecuteContext { EntityId = entityId, Login = login, Data = dto!, EventTypeId = eventTypeId, Facts = (IReadOnlyDictionary)(Facts ?? new Dictionary()) }; ctx.Set("Login", login); ctx.Set("Transaction", tx!); ctx.QualifierType = QualifierType.Validate; // ── Step 1: Validate qualifiers ─────────────────────────────────── await _qualifierFacade.ValidateAsync(ctx); saveActivity?.AddEvent(new ActivityEvent("entity.save.validate.done")); // ── Step 2: Enrich for workflow ─────────────────────────────────── // QualifierEngine loads DB_ENRICH definitions, executes SQL templates, and // writes every returned column into ctx._bag (e.g. ReportingToEmployeeId). // Must run BEFORE workflow check so enriched bag values are available // for EnrichQualifier strategy resolution in StartWorkflowAsync. await _qualifierFacade.EnrichForWorkflowAsync(ctx); saveActivity?.AddEvent(new ActivityEvent("entity.save.enrich_workflow.done")); // ── Step 3a: Workflow applicability check ───────────────────────── // Build enriched facts: caller-supplied keys + every scalar DTO property // under the "DTO." prefix so MWORKFLOWCONFIG.EVALCONDITION can reference // any field of the incoming request — e.g. "DTO.TaskDetailType == 0". var enrichedFacts = new Dictionary( ctx.Facts, StringComparer.OrdinalIgnoreCase); FlattenDtoFacts(dto, enrichedFacts); var checkCtx = new WorkflowCheckContext { EntityId = entityId, OUId = login.WorkOUId, BizTransactionClassId = bizTransactionClassId, BizTransactionId = bizTransactionId, Facts = enrichedFacts }; var wf = await _workflowEngine.CheckWorkFlowApplicability(checkCtx, login, tx!); ctx.Set("WorkflowInfo", wf); saveActivity?.SetTag("gb5.workflow.mode", wf.WorkFlowMode.ToString()); saveActivity?.SetTag("gb5.workflow.enabled", wf.IsEnabled); saveActivity?.SetTag("gb5.workflow.id", wf.WorkFlowId); saveActivity?.AddEvent(new ActivityEvent("entity.save.workflow_check.done")); // ── Step 3b: PrePersist qualifiers ──────────────────────────────── await _qualifierFacade.ExecuteStageAsync(QualifierStage.PrePersist, ctx); saveActivity?.AddEvent(new ActivityEvent("entity.save.pre_persist.done")); // ── Step 4: Persist — skipped for WIP ──────────────────────────── // For WIP, the entity is not saved until the approval callback fires. if (wf.WorkFlowMode != WORKFLOW.WIP) { // WIP approval callback: restore original body field values that the BLL // may have regenerated (e.g. task number, invoice number). The snapshot was // captured in BaseEndpoint before BLL execution so original DATAJSON values // win over BLL-recomputed values. SetDtoStatus runs AFTER restore so // TaskStatus is correctly forced to 1 regardless of the snapshot value. if (login.WipApprovalId > 0 && !string.IsNullOrEmpty(login.WipOriginalBodyJson)) RestoreDtoFromBodyJson(dto, login.WipOriginalBodyJson); // Workflow approval callback — mark the record active before persisting. if (login.WipApprovalId > 0) SetDtoStatus(dto, 1); var objectId = await persistFunc(tx!); ctx.ObjectId = objectId; saveActivity?.SetTag("gb5.entity.object_id", ctx.ObjectId); // WIP approval callback: write the real saved entity ID back into // TWORKFLOWWIP.OBJECTID in the same transaction so the WIP row // reflects the entity that was created/updated by this callback. // // objectId != 0: entity IDs in this system are negative (e.g. -1499832620). // 0 = sentinel "no entity yet" (task not saved) — skip update. // Any other value (positive or negative) = real PK — update OBJECTID. // // Wrapped in try-catch: a failed OBJECTID update must NEVER roll back // the already-inserted entity row. Task is saved; OBJECTID can be -1. if (login.WipApprovalId > 0 && objectId != 0) { try { await _workflowEngine.UpdateWipObjectIdAsync( login.WipApprovalId, objectId, login, tx!); saveActivity?.SetTag("gb5.wip.objectid_updated", objectId); saveActivity?.AddEvent( new ActivityEvent("entity.save.wip_objectid_update.done")); } catch (Exception ex) { // Non-fatal: entity was saved; only the WIP linkage update failed. // TWORKFLOWWIP.OBJECTID stays -1 — does not affect entity correctness. saveActivity?.SetTag("gb5.wip.objectid_update_error", ex.Message); saveActivity?.SetStatus( System.Diagnostics.ActivityStatusCode.Error, $"WIP OBJECTID update failed for WipId={login.WipApprovalId}: {ex.Message}"); } } } saveActivity?.AddEvent(new ActivityEvent("entity.save.persist.done")); // ── Step 4b: Cancel any prior pending approval for this record ───── // An update re-runs this whole pipeline. If the record already has a real // pending workflow instance (not just an orphan — see DeleteDuplicateInstanceAsync // below), it must be cancelled before a fresh one is started for the edited // data, otherwise the old and new instances both stay pending side by side. if (!isNewEntity && wf.WorkFlowMode != WORKFLOW.NONE) { await _workflowEngine.CancelPendingInstanceAsync( entityId, ctx.ObjectId, "Superseded by resubmission", login, tx!); } // ── Step 5: Start workflow ──────────────────────────────────────── if (wf.WorkFlowMode != WORKFLOW.NONE) { // For WIP: embed new-vs-update intent into DataJson so the dispatcher can // restore the correct persist path on callback — no LoginDTO flags needed. // __WipIsNew=true → dispatcher resets PK to 0 → BLL does INSERT // __WipIsNew=false → dispatcher leaves PK intact → BLL does UPDATE object dataJson = ctx.Data!; if (wf.WorkFlowMode == WORKFLOW.WIP) { try { var jObj = JsonNode.Parse(JsonSerializer.Serialize(ctx.Data))!.AsObject(); jObj["__WipIsNew"] = isNewEntity; if (isNewEntity) { // Store the PK field name so the dispatcher knows which field to reset string? pkField = DeriveEntityPkFieldName(); if (pkField != null) jObj["__WipPkField"] = pkField; } // Preserve the original HTTP body format in DATAJSON so the dispatcher // can replay the exact format the endpoint expects — no guesswork needed. // login.IsBodyArray = true → endpoint takes List → store as "[{...}]" // login.IsBodyArray = false → endpoint takes T → store as "{...}" dataJson = login.IsBodyArray ? new System.Text.Json.Nodes.JsonArray { jObj }.ToJsonString() : jObj.ToJsonString(); } catch { /* fallback: leave dataJson as original DTO */ } } // WorkflowStartRequest.Facts = exactly what DB_ENRICH qualifiers wrote into // ctx._bag (e.g. ReportingToEmployeeId). EnrichQualifier strategy reads from here. var bagFacts = new Dictionary(StringComparer.OrdinalIgnoreCase); foreach (var kv in ctx.ExportBag()) bagFacts[kv.Key] = kv.Value; var startRequest = new WorkflowStartRequest { EntityId = entityId, OUId = login.WorkOUId, BizTransactionClassId = bizTransactionClassId, BizTransactionId = bizTransactionId, ObjectId = ctx.ObjectId, DataJson = dataJson, Facts = bagFacts }; var startResult = await _workflowEngine.StartWorkflowAsync(startRequest, login, tx!); saveActivity?.SetTag("gb5.workflow.instance_id", startResult.WorkflowInstanceId); saveActivity?.SetTag("gb5.workflow.is_wip", startResult.IsWip); saveActivity?.SetTag("gb5.workflow.is_started", startResult.IsStarted); // For WIP, use WipId as the context ObjectId for outbox publishing if (startResult.IsWip && startResult.WipId > 0) { ctx.ObjectId = startResult.WipId; // Signal BaseEndpoint to replace the BLL's "saved successfully" message // with a clear pending-approval notice — entity is queued, not yet saved. login.WipTriggered = true; } // Edge case: first approval level had no applicable steps — WIP finalised // immediately at submission. The outer transaction (managed by the BLL caller) // has NOT committed yet, so dispatch cannot fire here safely. // WorkflowStartResult.PendingDispatch is populated — BLL caller should // dispatch it after committing the transaction. if (startResult.PendingDispatch != null) saveActivity?.SetTag("gb5.workflow.wip_immediate_dispatch_pending", true); saveActivity?.AddEvent(new ActivityEvent("entity.save.workflow_start.done")); } // ── Step 6: PostPersist qualifiers ──────────────────────────────── await _qualifierFacade.ExecuteStageAsync(QualifierStage.PostPersist, ctx); saveActivity?.AddEvent(new ActivityEvent("entity.save.post_persist.done")); // ── Step 7: Enrich for event ────────────────────────────────────── await _qualifierFacade.EnrichForEventAsync(ctx); saveActivity?.AddEvent(new ActivityEvent("entity.save.enrich_event.done")); // ── Build enriched payload (built once, used in both publish paths) ─ // LoginJson — full session: WorkOUId, RoleId, BranchId, timezone, formats, etc. // WipApprovalId and RequestUrl excluded by [JsonIgnore]. // ContextJson — EntityId/ObjectId/EventTypeId + caller Facts + DB_ENRICH bag // (e.g. ReportingToEmployeeId, DepartmentId) + WorkflowInfo. // Transaction and Login entries excluded (non-serialisable / redundant). var loginJson = BuildLoginJson(login); var contextJson = BuildContextJson(ctx); var payloadJson = JsonSerializer.Serialize(ctx.Data); // Explicit caller-supplied header GUID (e.g. AttachmentHeaderGuid from the // request header) wins over the DTO's own RowGuid reflection — callers that // route attachments to a specific outbox header must control this value directly. var resolvedHeaderRowGuid = headerRowGuid ?? GetDtoRowGuid(dto); // Quick-reference sizes on the parent span (no data duplication — full content is on the child publish span) saveActivity?.SetTag("gb5.publish.payload_bytes", payloadJson.Length); saveActivity?.SetTag("gb5.publish.login_bytes", loginJson.Length); saveActivity?.SetTag("gb5.publish.context_bytes", contextJson.Length); // ── Direct Dapr publish ─────────────────────────────────────────── // Fire-and-forget so subscribers can react immediately without waiting // for the outbox background job. Non-critical: Dapr unavailability must // never fail the business transaction — the outbox guarantees delivery. using (var directPublishActivity = GB5ActivitySources.EntitySave.StartActivity( "entity.save.direct_publish", ActivityKind.Producer)) { // ── Routing / identity ──────────────────────────────────────── directPublishActivity?.SetTag("gb5.msg.pubsub", "pubsub"); directPublishActivity?.SetTag("gb5.msg.topic", $"EVENTTYPEID:{ctx.EventTypeId}"); directPublishActivity?.SetTag("gb5.msg.event_type_id", ctx.EventTypeId); directPublishActivity?.SetTag("gb5.msg.object_type_id", ctx.EntityId); directPublishActivity?.SetTag("gb5.msg.object_id", ctx.ObjectId); directPublishActivity?.SetTag("gb5.msg.tenant_id", login.ClientId); directPublishActivity?.SetTag("gb5.msg.user_id", login.UserId); directPublishActivity?.SetTag("gb5.msg.ou_id", login.WorkOUId); directPublishActivity?.SetTag("gb5.msg.connection_name", login.ConnectionDatabaseName ?? string.Empty); directPublishActivity?.SetTag("gb5.msg.header_row_guid", resolvedHeaderRowGuid?.ToString() ?? string.Empty); directPublishActivity?.SetTag("gb5.msg.dto_type", typeof(TDto).Name); directPublishActivity?.SetTag("gb5.msg.is_new_entity", isNewEntity); directPublishActivity?.SetTag("gb5.msg.workflow_mode", wf.WorkFlowMode.ToString()); directPublishActivity?.SetTag("gb5.msg.workflow_id", wf.WorkFlowId); // ── Full published data — visible in Zipkin / Jaeger span detail ─ // These three tags carry everything the subscriber receives. directPublishActivity?.SetTag("gb5.msg.payload", payloadJson); directPublishActivity?.SetTag("gb5.msg.login_json", loginJson); directPublishActivity?.SetTag("gb5.msg.context_json", contextJson); try { await _dapr.PublishEventAsync( "pubsub", $"EVENTTYPEID:{ctx.EventTypeId}", new OutboxEventMessage { EventTypeId = ctx.EventTypeId, ObjectTypeId = ctx.EntityId, ObjectId = ctx.ObjectId, TenantId = login.ClientId, UserId = login.UserId, Payload = payloadJson, HeaderRowGuid = resolvedHeaderRowGuid, // Carry the exact connection the publisher used so the framework // subscriber can connect to the right tenant DB without re-deriving // from MSERVERCONFIG (which can have multiple active entries for // the same ClientId and return an incorrect database). ConnectionName = login.ConnectionDatabaseName ?? string.Empty, // Full session + pipeline context — subscribers consume what they need. LoginJson = loginJson, ContextJson = contextJson, // Carry this producer span's W3C context in the message itself so // EventActionSubscribeController can parent its consumer span correctly // regardless of Dapr broker/CloudEvent metadata propagation. TraceParent = directPublishActivity?.Id, TraceState = directPublishActivity?.TraceStateString }, cancellationToken: CancellationToken.None).ConfigureAwait(false); directPublishActivity?.SetStatus(ActivityStatusCode.Ok); directPublishActivity?.SetTag("gb5.dapr.direct_publish_status", "success"); } catch (Exception pubEx) { // Non-fatal: outbox guarantees at-least-once delivery. // Record full exception details for diagnostics but do not rethrow. directPublishActivity?.AddException(pubEx); directPublishActivity?.SetStatus(ActivityStatusCode.Error, pubEx.Message); directPublishActivity?.SetTag("gb5.dapr.direct_publish_status", "failed"); directPublishActivity?.SetTag("gb5.dapr.direct_publish_error", pubEx.GetType().Name); } } saveActivity?.AddEvent(new ActivityEvent("entity.save.direct_publish.done")); // ── Outbox write (inside transaction — guarantees at-least-once delivery) ─ // ContextJson is written here so the outbox relay (PublishPendingEventsAsync) // carries the same enriched Bag (MailId, ReportingToEmployeeId, etc.) as the // direct-publish path. Without it TOUTBOX.CONTEXTJSON = NULL → relay publishes // ContextJson="{}" → EmailActionHandler ContextBag empty → recipient unresolved. await _outbox.PublishEventAsync( new OutboxDTO { EventTypeId = ctx.EventTypeId, ObjectTypeId = ctx.EntityId, ObjectId = ctx.ObjectId, Payload = ctx.Data, HeaderRowGuid = resolvedHeaderRowGuid, ContextJson = contextJson }, login, tx!); saveActivity?.AddEvent(new ActivityEvent("entity.save.outbox_write.done")); // ── Commit ──────────────────────────────────────────────────────── if (!isExternalTx) { await _queryExecutor.CommitAsync(tx!); saveActivity?.AddEvent(new ActivityEvent("entity.save.committed")); } saveActivity?.SetStatus(ActivityStatusCode.Ok); } catch (Exception ex) { saveActivity?.AddException(ex); saveActivity?.SetStatus(ActivityStatusCode.Error, ex.Message); saveActivity?.SetTag("gb5.entity.save_error_type", ex.GetType().Name); saveActivity?.SetTag("gb5.entity.save_error_message", ex.Message); if (!isExternalTx && tx != null) await _queryExecutor.RollbackAsync(tx); throw; } } /// /// Universal Delete Execution Pipeline — the delete-side counterpart to /// . GB5 has no shared delete pipeline the way it has /// one for save (each entity's BLL calls its own DAL directly), so this exists purely /// to give every entity delete the same one-line cancel-pending-workflow behavior that /// already gives every save, instead of each BLL wiring /// by hand. /// Steps: BeginTx → CancelPendingWorkflow → deleteFunc → Commit. /// public async Task ExecuteDeleteAsync( int entityId, int objectId, LoginDTO login, Func> deleteFunc, string cancelReason = "Record deleted", DbTransaction? externalTransaction = null) { using var deleteActivity = GB5ActivitySources.EntitySave.StartActivity( "entity.delete", ActivityKind.Internal); deleteActivity?.SetTag("gb5.entity.id", entityId); deleteActivity?.SetTag("gb5.entity.object_id", objectId); deleteActivity?.SetTag("gb5.tenant.id", login.ClientId); deleteActivity?.SetTag("gb5.user.id", login.UserId); deleteActivity?.SetTag("gb5.entity.has_external_tx", externalTransaction != null); DbTransaction? tx = externalTransaction; var isExternalTx = tx != null; try { if (!isExternalTx) { tx = await _queryExecutor.BeginTransactionAsync(login); deleteActivity?.AddEvent(new ActivityEvent("entity.delete.tx_begin")); } // Cancel any pending approval for this record before it's removed so it // doesn't linger in an approver's inbox pointing at a deleted record. await _workflowEngine.CancelPendingInstanceAsync(entityId, objectId, cancelReason, login, tx!); deleteActivity?.AddEvent(new ActivityEvent("entity.delete.workflow_cancel.done")); var result = await deleteFunc(tx!); deleteActivity?.AddEvent(new ActivityEvent("entity.delete.persist.done")); if (!isExternalTx) { await _queryExecutor.CommitAsync(tx!); deleteActivity?.AddEvent(new ActivityEvent("entity.delete.committed")); } deleteActivity?.SetStatus(ActivityStatusCode.Ok); return result; } catch (Exception ex) { deleteActivity?.AddException(ex); deleteActivity?.SetStatus(ActivityStatusCode.Error, ex.Message); deleteActivity?.SetTag("gb5.entity.delete_error_type", ex.GetType().Name); deleteActivity?.SetTag("gb5.entity.delete_error_message", ex.Message); if (!isExternalTx && tx != null) await _queryExecutor.RollbackAsync(tx); throw; } } /// Result-less overload of for delete methods that build their own return message. public Task ExecuteDeleteAsync( int entityId, int objectId, LoginDTO login, Func deleteFunc, string cancelReason = "Record deleted", DbTransaction? externalTransaction = null) => ExecuteDeleteAsync(entityId, objectId, login, async tx => { await deleteFunc(tx); return null; }, cancelReason, externalTransaction); // ── Helpers ─────────────────────────────────────────────────────────────── /// /// Serialises the full to JSON for inclusion in the published /// field. /// WipApprovalId and RequestUrl are excluded automatically by their [JsonIgnore] attributes. /// Falls back to a minimal JSON object containing only the critical identity fields /// when serialisation fails (e.g. due to a circular reference in an unexpected subtype). /// private static string BuildLoginJson(LoginDTO login) { try { return JsonSerializer.Serialize(login, _publishSerialiserOptions); } catch (Exception ex) { // Minimal fallback — preserves routing and identity information return JsonSerializer.Serialize(new { ClientId = login.ClientId, UserId = login.UserId, WorkOUId = login.WorkOUId, ConnectionDatabaseName = login.ConnectionDatabaseName, SerializationError = ex.Message }); } } /// /// Builds the payload from the /// current pipeline . /// /// Includes: EntityId, ObjectId, EventTypeId, caller-supplied Facts, and every /// entry in the context bag EXCEPT Transaction (not serialisable) and /// Login (carried separately in LoginJson). /// /// Falls back to a minimal JSON object with error details when serialisation fails. /// private static string BuildContextJson(ExecuteContext ctx) { try { // Exclude infrastructure objects that are not serialisable or are already // carried in LoginJson: Transaction (DbTransaction), Login (LoginDTO). var bag = ctx.ExportBag() .Where(kv => !string.Equals(kv.Key, "Transaction", StringComparison.OrdinalIgnoreCase) && !string.Equals(kv.Key, "Login", StringComparison.OrdinalIgnoreCase)) .ToDictionary(kv => kv.Key, kv => (object?)kv.Value); var snapshot = new { EntityId = ctx.EntityId, ObjectId = ctx.ObjectId, EventTypeId = ctx.EventTypeId, Facts = ctx.Facts, Bag = bag }; return JsonSerializer.Serialize(snapshot, _publishSerialiserOptions); } catch (Exception ex) { return JsonSerializer.Serialize(new { EntityId = ctx.EntityId, ObjectId = ctx.ObjectId, EventTypeId = ctx.EventTypeId, SerializationError = ex.Message }); } } private static Guid? GetDtoRowGuid(TDto dto) { if (dto == null) return null; var prop = typeof(TDto).GetProperty("RowGuid", BindingFlags.Public | BindingFlags.Instance); if (prop == null) return null; var raw = prop.GetValue(dto)?.ToString(); return Guid.TryParse(raw, out var g) ? g : (Guid?)null; } /// /// Derives the entity PK field name from the DTO type by naming convention. /// "TaskDTO" → "TaskId", "AccountDTO" → "AccountId", etc. /// Returns null when the convention does not produce a real property on TDto. /// private static string? DeriveEntityPkFieldName() { string typeName = typeof(TDto).Name; if (!typeName.EndsWith("DTO", StringComparison.OrdinalIgnoreCase)) return null; string candidate = typeName[..^3] + "Id"; // "TaskDTO" → "TaskId" return typeof(TDto).GetProperty(candidate, BindingFlags.Public | BindingFlags.Instance) != null ? candidate : null; } /// /// Reflects all scalar (primitive / string / DateTime / Guid / enum / Nullable of those) /// public properties of into under the /// "DTO." prefix. Complex / collection properties are skipped so the evaluator always /// receives clean, Compare()-compatible string-convertible values. /// /// This lets MWORKFLOWCONFIG.EVALCONDITION reference any DTO field dynamically — /// e.g. "DTO.TaskDetailType == 0" — without the caller needing to pre-populate Facts. /// // Per-TDto resolution cache — the reflection walk below (attribute scan → naming // convention → suffix scan) only needs to run once per DTO type for the lifetime of // the process; every subsequent SetDtoStatus call for that TDto reuses the cached // PropertyInfo (or null, cached via the bool flag, when no property could be resolved). private static readonly System.Collections.Concurrent.ConcurrentDictionary _statusPropCache = new(); /// /// Sets the DTO's canonical status property to . Called when /// no workflow is configured (immediate persist) and on WIP-approval finalisation, so /// the record is persisted in the active/approved state. /// /// Resolution order (fully dynamic — no DTO shape is hardcoded here): /// 1. The property explicitly marked — the /// authoritative, opt-in signal. Use this whenever a DTO's status field doesn't /// follow the naming convention below, or has several "*Status" properties where /// the convention name doesn't pick the right one. /// 2. Exact bare "Status" property. /// 3. The naming convention most entity DTOs follow — "{EntityName}Status", derived /// from the DTO type name (e.g. "TaskDTO" → "TaskStatus", "CallDTO" → "CallStatus"). /// Tried BEFORE the generic suffix scan below so DTOs with multiple "*Status" /// properties (e.g. CallDTO: CallStatus, CallCallOriginalStatus, CallCoverageStatus) /// resolve unambiguously to the one that actually gates persistence. /// 4. Fallback: suffix scan, only when it yields exactly one "*Status" candidate. /// /// private static void SetDtoStatus(TDto? dto, int value) { if (dto is null) return; var prop = _statusPropCache.GetOrAdd(typeof(TDto), ResolveStatusProperty); if (prop is null) throw new InvalidOperationException( $"BaseEntityAppService could not determine the active/approved status property " + $"for {typeof(TDto).Name} — needed because a workflow (or WIP approval) is now " + $"configured for this entity. Fix: add [GB5Shared.EntityHandler.WorkflowStatusField] " + $"to the DTO's canonical status property, e.g.:\n" + $" [WorkflowStatusField]\n" + $" public byte {typeof(TDto).Name.Replace("DTO", "", StringComparison.OrdinalIgnoreCase)}Status {{ get; set; }}\n" + $"This is a one-line, self-service fix — no shared framework change needed."); var target = Nullable.GetUnderlyingType(prop.PropertyType) ?? prop.PropertyType; prop.SetValue(dto, Convert.ChangeType(value, target)); } private static PropertyInfo? ResolveStatusProperty(Type dtoType) { var allProps = dtoType.GetProperties(BindingFlags.Public | BindingFlags.Instance); // 1. Explicit opt-in attribute — authoritative, checked first. var attributed = allProps.FirstOrDefault(p => p.CanWrite && p.GetCustomAttribute() != null); if (attributed != null) return attributed; var prop = allProps.FirstOrDefault(p => p.CanWrite && string.Equals(p.Name, "status", StringComparison.OrdinalIgnoreCase)); if (prop is null) { // "TaskDTO" → "TaskStatus", "CallDTO" → "CallStatus" — same naming convention // DeriveEntityPkFieldName uses for the PK ("TaskDTO" → "TaskId"), but resolved // independently here since a DTO can lack a same-named Id property while still // following the {Entity}Status convention for its status field. string typeName = dtoType.Name; if (typeName.EndsWith("DTO", StringComparison.OrdinalIgnoreCase)) { string conventionName = typeName[..^3] + "Status"; prop = allProps.FirstOrDefault(p => p.CanWrite && string.Equals(p.Name, conventionName, StringComparison.OrdinalIgnoreCase)); } } if (prop is null) { var candidates = allProps .Where(p => p.CanWrite && p.Name.EndsWith("Status", StringComparison.OrdinalIgnoreCase)) .ToList(); if (candidates.Count == 1) prop = candidates[0]; } return prop; } /// /// Copies property values from a JSON body snapshot onto the existing DTO instance. /// Handles both array bodies ("[{...}]") and object bodies ("{...}") — takes the /// first element for array bodies (WIP always dispatches a single entity at a time). /// Properties that cannot be deserialized (custom types, format mismatches) are silently /// skipped so a single bad field never blocks the entire restore. /// Framework control flags (names starting with "__") are always excluded. /// private static void RestoreDtoFromBodyJson(TDto? dto, string bodyJson) { if (dto is null) return; try { using var doc = JsonDocument.Parse(bodyJson); var root = doc.RootElement; JsonElement element; if (root.ValueKind == JsonValueKind.Array) { if (root.GetArrayLength() == 0) return; element = root[0]; } else { element = root; } var dtoType = dto.GetType(); var opts = new JsonSerializerOptions { PropertyNameCaseInsensitive = true }; foreach (var jsonProp in element.EnumerateObject()) { // Never copy WIP or other internal framework control flags onto the DTO if (jsonProp.Name.StartsWith("__", StringComparison.Ordinal)) continue; var prop = dtoType.GetProperty(jsonProp.Name, BindingFlags.Public | BindingFlags.Instance | BindingFlags.IgnoreCase); if (prop == null || !prop.CanWrite) continue; try { var value = JsonSerializer.Deserialize( jsonProp.Value.GetRawText(), prop.PropertyType, opts); prop.SetValue(dto, value); } catch { /* skip individual property on type or format mismatch */ } } } catch { /* if full parse fails, proceed with BLL-modified DTO */ } } private static void FlattenDtoFacts(TDto? dto, Dictionary target) { if (dto is null) return; foreach (var prop in typeof(TDto).GetProperties(BindingFlags.Public | BindingFlags.Instance)) { if (!prop.CanRead || prop.GetIndexParameters().Length > 0) continue; var underlying = Nullable.GetUnderlyingType(prop.PropertyType) ?? prop.PropertyType; if (!underlying.IsPrimitive && underlying != typeof(string) && underlying != typeof(decimal) && underlying != typeof(DateTime) && underlying != typeof(Guid) && !underlying.IsEnum) continue; target[$"DTO.{prop.Name}"] = prop.GetValue(dto); } } } }