using System.Data; using System.Text.RegularExpressions; using GB5Shared.IceMap; using GB5Shared.DTO.Ice; using GB5Shared.DTO.Framework.Login; using GB5Shared.QueryExecutor; using GB5Shared.Telemetry; using GB5Shared.Validation; using IceImportBLL.AddonFieldAccumulator; using IceImportBLL.EntityLookupService; using IceImportBLL.EntityPersistenceAdapter; using IceImportBLL.FieldResolvers; using IceImportBLL.HierarchyMapper; using IceImportBLL.UpsertDecisionService; using IceImportBLL.Validators; using IceImportDAL.CustomCode.EntityMemberFieldMap; using IceImportDAL.CustomCode.IceMapValidationRule; using IceImportDAL.CustomCode.Parsers; using IceImportDAL.DTO.EntityMemberFieldMap; using IceImportDAL.DTO.EntityPersistenceAdapter; using IceImportDAL.DTO.IceImportPipeline; using Microsoft.Extensions.Logging; using Newtonsoft.Json; namespace IceImportBLL.IceImportPipeline; public class IceImportPipeline : IIceImportPipeline { private readonly IIceMapBLL _iceMapBLL; private readonly IEnumerable _parsers; private readonly IEnumerable _resolvers; private readonly IEntityLookupService _lookupService; private readonly IHierarchyMapper _hierarchyMapper; private readonly IAddonFieldAccumulator _addonAccumulator; private readonly IEnumerable _validators; private readonly IIceMapValidationRuleDAL _validationRuleDal; private readonly IEnumerable _persistenceAdapters; private readonly IUpsertDecisionService _upsertDecisionService; private readonly IEntityMemberFieldMapDAL _fieldMapDal; private readonly IQueryExecutor _queryExecutor; private readonly IValidation _validation; private readonly ILogger _logger; // Default validators run for every IceMap unconditionally; anything else must be explicitly // opted in via MICEMAPVALIDATIONRULE (e.g. IfscBankBranchValidator's "IFSC_BANKBRANCH" key). private static readonly HashSet _defaultValidatorKeys = new(StringComparer.OrdinalIgnoreCase) { "HEADER_FIELD_EXISTENCE", "REQUIRED_FIELD" }; // Guards CommitAsUpdateAsync's dynamic SQL — table/column names come from // MENTITY/DBOBJECT/DBOBJECTFIELDS metadata, not upload content, but are validated before being // interpolated anyway (same discipline as GenericEntityCodeLookupDAL's identifier check). private static readonly Regex _validIdentifier = new(@"^[A-Za-z_][A-Za-z0-9_]*$", RegexOptions.Compiled); public IceImportPipeline( IIceMapBLL iceMapBLL, IEnumerable parsers, IEnumerable resolvers, IEntityLookupService lookupService, IHierarchyMapper hierarchyMapper, IAddonFieldAccumulator addonAccumulator, IEnumerable validators, IIceMapValidationRuleDAL validationRuleDal, IEnumerable persistenceAdapters, IUpsertDecisionService upsertDecisionService, IEntityMemberFieldMapDAL fieldMapDal, IQueryExecutor queryExecutor, IValidation validation, ILogger logger) { _iceMapBLL = iceMapBLL; _parsers = parsers; _resolvers = resolvers; _lookupService = lookupService; _hierarchyMapper = hierarchyMapper; _addonAccumulator = addonAccumulator; _validators = validators; _validationRuleDal = validationRuleDal; _persistenceAdapters = persistenceAdapters; _upsertDecisionService = upsertDecisionService; _fieldMapDal = fieldMapDal; _queryExecutor = queryExecutor; _validation = validation; _logger = logger; } // ── Stage 1 — Parse ───────────────────────────────────────────────── public async Task ParseAsync( Stream file, string fileName, int iceMapId, LoginDTO login, CancellationToken ct) { await _validation.NotNull(file, nameof(file)).ConfigureAwait(false); using var section = GB5Trace.BeginSection("ice-import-parse"); GB5Trace.Step("load-icemap", new { IceMapId = iceMapId }); var mapDefinition = await LoadMapDefinitionAsync(iceMapId, login).ConfigureAwait(false); return await ParseWithMapAsync(file, fileName, mapDefinition, ct).ConfigureAwait(false); } private async Task ParseWithMapAsync( Stream file, string fileName, IceMapDTO mapDefinition, CancellationToken ct) { var parser = _parsers.FirstOrDefault(p => p.SourceType == mapDefinition.IceMapIceSourceType); if (parser is null) { GB5Trace.MarkFailed("ice-import-no-parser"); throw new InvalidOperationException( $"No ISourceFileParser registered for SourceType {mapDefinition.IceMapIceSourceType}."); } GB5Trace.Step("parse-file", new { fileName, mapDefinition.IceMapIceSourceType }); try { return await parser.ParseAsync(file, fileName, mapDefinition, ct).ConfigureAwait(false); } catch (Exception ex) { GB5Trace.MarkFailed("ice-import-parse-failed", ex); _logger.LogError(ex, "IceImport: parse failed for IceMapId {IceMapId} File {FileName}", mapDefinition.IceMapId, fileName); throw; } } private async Task LoadMapDefinitionAsync(int iceMapId, LoginDTO login) { var json = await _iceMapBLL.GetIceMap(iceMapId, login).ConfigureAwait(false); var map = JsonConvert.DeserializeObject(json); return map ?? throw new InvalidOperationException($"IceMap {iceMapId} was not found."); } // Reuses the exact same table-resolution logic MapAndTransformAsync applies to the root entity // (ResolveTableForEntity, below) so this always agrees with whatever table that stage actually // reads from — no separate/divergent lookup. public int CountSourceRecords(ParsedImportDataSetDTO parsed, IceMapDTO mapDefinition) { var rootEntityId = mapDefinition.IceMapRootEntityId; var rootDetails = mapDefinition.IceMapDetailsArray .Where(d => d.DestinationMemberEntityId == rootEntityId) .ToList(); var table = ResolveTableForEntity(parsed, rootEntityId, rootDetails); return table?.Rows.Count ?? 0; } // ── Stage 2 — Map & Transform ─────────────────────────────────────── public async Task> MapAndTransformAsync( ParsedImportDataSetDTO parsed, IceMapDTO mapDefinition, LoginDTO login, CancellationToken ct, bool isUpdateMode = false) { await _validation.NotNull(parsed, nameof(parsed)).ConfigureAwait(false); await _validation.NotNull(mapDefinition, nameof(mapDefinition)).ConfigureAwait(false); using var section = GB5Trace.BeginSection("ice-import-map-transform"); GB5Trace.Step("resolve-lookups", new { mapDefinition.IceMapId }); var lookups = await _lookupService .ResolveAllLookupsAsync(parsed, mapDefinition, login, ct) .ConfigureAwait(false); var entityGroups = mapDefinition.IceMapDetailsArray .GroupBy(d => d.DestinationMemberEntityId) .ToList(); GB5Trace.Step("map-transform-entities", new { EntityCount = entityGroups.Count }); // Object/List (DataType=5) relationships — e.g. BOMDTO.BOMDetailArray -> List // — need every child row grouped under ONE deduplicated parent row (by the child's own // IsFilterField=0 grouping column, e.g. BOMID), instead of the default "1 parent row per // flat file row" HierarchyMapper otherwise assumes for every other map. Build a per-child // "physical row position -> group index" map from the child's own parsed table up front, // then feed that as the CorrelationIndex for both the parent (deduplicated) and child // (ungrouped) passes below — HierarchyMapper.BuildHierarchy already groups multiple child // candidates under one parent by shared CorrelationIndex, so it needs no changes. var objectListRelationships = _hierarchyMapper.GetObjectListRelationships(mapDefinition); // Addon relationships (ChildEntityId == AddonEntityId) are included in // objectListRelationships for detection/logging and for CommitRowRecursivelyAsync's own // embedding step, but must never drive the dedup/grouping logic below — Addon is an // inherent 1:1 relationship with no real grouping field of its own (see // HierarchyMapper.GetObjectListRelationships' doc comment); treating the primary entity as // "needing dedup by an Addon field's value" would risk silently dropping/merging unrelated // primary rows that happen to share the same Addon field value. var dedupRelationships = objectListRelationships.Where(r => r.ChildEntityId != AddonEntityId).ToList(); var relationshipByParentEntityId = dedupRelationships .GroupBy(r => r.ParentEntityId) .ToDictionary(g => g.Key, g => g.First()); // one grouping key governs dedup per parent var relationshipByChildEntityId = dedupRelationships .GroupBy(r => r.ChildEntityId) .ToDictionary(g => g.Key, g => g.First()); // Temporary diagnostic — remove once BOMDTO/BOMDETAILDTO-style parent/detail-array maps are // confirmed stable. Shows whether the DataType=5 relationship was even detected at all // (empty list here means DestinationMemberDataType/EntityId/IsFilterField didn't come // through mapDefinition the way GetObjectListRelationships expects). _logger.LogInformation( "IceImport MapAndTransform: ObjectListRelationships=[{Relationships}]", string.Join(" | ", objectListRelationships.Select(r => $"Parent={r.ParentEntityId} Child={r.ChildEntityId} ListMember={r.ListMemberName} GroupingField={r.GroupingSourceFieldName}"))); var groupIndexMapByChildEntityId = new Dictionary>(); foreach (var relationship in dedupRelationships) { var childDetails = entityGroups.FirstOrDefault(g => g.Key == relationship.ChildEntityId)?.ToList() ?? new List(); var childTable = ResolveTableForEntity(parsed, relationship.ChildEntityId, childDetails); if (childTable is null) { _logger.LogWarning( "IceImport MapAndTransform: no parsed table resolved for ChildEntityId {ChildEntityId} " + "— object/List relationship cannot be grouped.", relationship.ChildEntityId); continue; } var groupMap = BuildGroupIndexMap(childTable, relationship.GroupingSourceFieldName); groupIndexMapByChildEntityId[relationship.ChildEntityId] = groupMap; _logger.LogInformation( "IceImport MapAndTransform: ChildEntityId {ChildEntityId} GroupingField {GroupingField} " + "— {RowCount} row(s) collapsed into {GroupCount} distinct group(s).", relationship.ChildEntityId, relationship.GroupingSourceFieldName, groupMap.Count, groupMap.Values.Distinct().Count()); } var resolvedByEntity = new Dictionary>(); foreach (var entityGroup in entityGroups) { var entityId = entityGroup.Key; var details = entityGroup.ToList(); var table = ResolveTableForEntity(parsed, entityId, details); if (table is null) { _logger.LogWarning("IceImport: no parsed table resolved for EntityId {EntityId}", entityId); continue; } var entityCode = entityId == mapDefinition.IceMapRootEntityId && !string.IsNullOrWhiteSpace(mapDefinition.IceMapEntityCode) ? mapDefinition.IceMapEntityCode! : entityId.ToString(); var candidates = new List(); var hasRowIndexColumn = table.Columns.Contains(JsonSourceFileParser.RowIndexColumn); var positionalIndex = 0; // A parent in an object/List relationship must only produce ONE candidate row per // distinct grouping-key value — every other flat row sharing that same group is a // duplicate of the same parent record (e.g. all rows with BOMID="BBU310000QQ"), not a // second parent. Dictionary? parentGroupIndexMap = null; if (relationshipByParentEntityId.TryGetValue(entityId, out var parentRelationship)) groupIndexMapByChildEntityId.TryGetValue(parentRelationship.ChildEntityId, out parentGroupIndexMap); var seenGroupIndexes = new HashSet(); Dictionary? childGroupIndexMap = null; if (relationshipByChildEntityId.ContainsKey(entityId)) groupIndexMapByChildEntityId.TryGetValue(entityId, out childGroupIndexMap); foreach (DataRow dataRow in table.Rows) { int correlationIndex; if (parentGroupIndexMap is not null && parentGroupIndexMap.TryGetValue(positionalIndex, out var parentGroupIndex)) { if (!seenGroupIndexes.Add(parentGroupIndex)) { positionalIndex++; continue; // duplicate parent row for a group already represented } correlationIndex = parentGroupIndex; } else if (childGroupIndexMap is not null && childGroupIndexMap.TryGetValue(positionalIndex, out var childGroupIndex)) { correlationIndex = childGroupIndex; } else { correlationIndex = hasRowIndexColumn ? ParseRowIndex(dataRow[JsonSourceFileParser.RowIndexColumn]) : positionalIndex; } // Excel/CSV tables have no __ROWINDEX__ column (that's JSON-only) and their header // row is already stripped by the parser before this DataTable is built, so their // first data row is actually file line 2, not line 1 — offset by one extra so // RowNumber matches what the user sees when they open the file. JSON has no header // row to account for. var mappedRow = new MappedEntityRowDTO { RowNumber = hasRowIndexColumn ? positionalIndex + 1 : positionalIndex + 2, EntityId = entityId, EntityCode = entityCode }; var resolvedPairs = new List<(IceMapDetailsDTO Detail, object? Value)>(); foreach (var detail in details) { // "Required and not provided" is checked against the RAW Excel cell // (IceMapDetailsSourceFieldName), not the resolved/coerced output — a bad date // format or a failed lookup must surface as its own distinct warning (see // SourceValueResolver/CodeToIdLookupResolver/NameToIdLookupResolver), never get // relabeled as "missing" just because coercion happened to fail. Only applies to // ValueType 0/3/4/5 (the types that actually read a source column) — 1/2/6/7 don't // read Excel at all, so RequiredFieldValidator still checks THEIR resolved value; // 8/9 get their own "not supported" finding instead. The object/List CONTAINER // marker (DestinationMemberDataType==5 with no IceMapDetailsAddonFields) is // excluded — its "SourceFieldName" is a container label, never a real column. var isSourceColumnValueType = detail.IceMapDetailsValueType is 0 or 3 or 4 or 5; var isContainerMarker = detail.DestinationMemberDataType == 5 && string.IsNullOrWhiteSpace(detail.IceMapDetailsAddonFields); if (isSourceColumnValueType && detail.IceMapDetailsIsOptional == 1 && !isContainerMarker) { var sourceFieldName = detail.IceMapDetailsSourceFieldName?.Trim().ToUpperInvariant(); var rawPresent = !string.IsNullOrEmpty(sourceFieldName) && dataRow.Table.Columns.Contains(sourceFieldName) && !string.IsNullOrWhiteSpace(dataRow[sourceFieldName]?.ToString()); if (!rawPresent) { // Addon fields (BONUSAPPL, INCAPPL, ...) commonly share ONE placeholder // DestinationMemberId whose MemberName is a generic technical label (e.g. // "STRING"), not the field's real name — label those by their own // IceMapDetailsAddonFields instead of the shared, meaningless member name. var fieldLabel = !string.IsNullOrWhiteSpace(detail.IceMapDetailsAddonFields) ? detail.IceMapDetailsAddonFields!.Trim() : detail.DestinationMemberMemberName ?? sourceFieldName ?? string.Empty; mappedRow.FieldWarnings.Add(new IceImportFieldErrorDTO { RowNumber = mappedRow.RowNumber, EntityCode = entityCode, FieldName = fieldLabel, Message = $"'{sourceFieldName}' is required and was not provided.", Severity = IceImportSeverity.Error }); } } // ValueType 8/9 ("... During Update") only make sense when isUpdateMode=true // (FixedDuringUpdateResolver/GlobalDuringUpdateResolver resolve to null otherwise); // ValueType 1/2 (their Insert-time counterparts) only make sense when // isUpdateMode=false. CommitImport (always isUpdateMode=false) and IceImport/Update // (always isUpdateMode=true) are mutually exclusive entry points, so exactly one of // these two pairs is "not supported" for any given run — flagged here with the same // distinct Error CommitImport already surfaces for 8/9, now symmetric for Update. var unsupportedValueTypeLabel = isUpdateMode ? detail.IceMapDetailsValueType switch { 1 => "Fixed Value", 2 => "Global System Value", _ => null } : detail.IceMapDetailsValueType switch { 8 => "Fixed Value During Update", 9 => "Global Value During Update", _ => null }; if (unsupportedValueTypeLabel is not null) { var fieldLabel = !string.IsNullOrWhiteSpace(detail.IceMapDetailsAddonFields) ? detail.IceMapDetailsAddonFields!.Trim() : detail.DestinationMemberMemberName ?? string.Empty; mappedRow.FieldWarnings.Add(new IceImportFieldErrorDTO { RowNumber = mappedRow.RowNumber, EntityCode = entityCode, FieldName = fieldLabel, Message = $"'{fieldLabel}': ValueType {detail.IceMapDetailsValueType} " + $"({unsupportedValueTypeLabel}) is not supported.", Severity = IceImportSeverity.Error }); } var resolver = _resolvers.FirstOrDefault(r => r.ValueType == detail.IceMapDetailsValueType); if (resolver is null) { mappedRow.FieldWarnings.Add(new IceImportFieldErrorDTO { RowNumber = mappedRow.RowNumber, EntityCode = entityCode, FieldName = detail.DestinationMemberMemberName ?? string.Empty, Message = $"No field resolver registered for ValueType {detail.IceMapDetailsValueType}.", Severity = IceImportSeverity.Warning }); continue; } var fieldCtx = new FieldResolutionContext { DetailConfig = detail, SourceRow = dataRow, Login = login, RowNumber = mappedRow.RowNumber, Lookups = lookups, IsUpdateMode = isUpdateMode // CommitImport: always false (upsert decision // happens later in CommitAsync, ValueType 8/9 resolve // to null at map time). IceImportUpdateService: true, // so FixedDuringUpdateResolver/GlobalDuringUpdateResolver // actually resolve — see IIceImportPipeline's doc comment. }; var fieldResult = await resolver.ResolveAsync(fieldCtx, ct).ConfigureAwait(false); var fieldName = detail.DestinationMemberMemberName; if (!string.IsNullOrEmpty(fieldName)) mappedRow.Fields[fieldName] = fieldResult.Value; resolvedPairs.Add((detail, fieldResult.Value)); if (fieldResult.Warning is not null) mappedRow.FieldWarnings.Add(fieldResult.Warning); } _addonAccumulator.Apply(mappedRow, resolvedPairs); candidates.Add(new ResolvedCandidateRow(correlationIndex, mappedRow)); positionalIndex++; } resolvedByEntity[entityId] = candidates; } // Temporary diagnostic — remove once the root-entity-id resolution is confirmed stable in // production. Logs the exact ids BuildHierarchy will compare so a 0-row result is // diagnosable from the log alone, without a debugger attached. _logger.LogInformation( "IceImport MapAndTransform: IceMapEntityId={IceMapEntityId} PocoToDtoLinkId={PocoToDtoLinkId} " + "ResolvedRootEntityId={RootEntityId} ResolvedByEntityKeys=[{Keys}]", mapDefinition.IceMapEntityId, mapDefinition.IceMapEntityPocoToDtoLinkId, mapDefinition.IceMapRootEntityId, string.Join(", ", resolvedByEntity.Select(kv => $"{kv.Key}:{kv.Value.Count}row(s)"))); var hierarchy = _hierarchyMapper.BuildHierarchy(resolvedByEntity, mapDefinition); _logger.LogInformation( "IceImport MapAndTransform: BuildHierarchy returned {RootRowCount} root row(s).", hierarchy.Count); return hierarchy; } private static int ParseRowIndex(object? value) => int.TryParse(value?.ToString(), out var idx) ? idx : 0; // Maps each physical row position within childTable to a 0-based index representing its // distinct groupingFieldName value, assigned in order of first appearance — e.g. row positions // 3,4,5,6 all sharing BOMID="BBU310000QQ" map to the same group index, distinct from every // other BOMID value. private static Dictionary BuildGroupIndexMap(DataTable childTable, string groupingFieldName) { var fieldName = groupingFieldName.Trim().ToUpperInvariant(); var map = new Dictionary(); var groupIndexByValue = new Dictionary(StringComparer.OrdinalIgnoreCase); var position = 0; foreach (DataRow row in childTable.Rows) { var value = childTable.Columns.Contains(fieldName) ? row[fieldName]?.ToString()?.Trim() ?? string.Empty : string.Empty; if (!groupIndexByValue.TryGetValue(value, out var groupIndex)) { groupIndex = groupIndexByValue.Count; groupIndexByValue[value] = groupIndex; } map[position] = groupIndex; position++; } return map; } // Exact numeric-EntityId key match first (Excel/CSV/JSON convention — see // FlatTableEntitySplitter/JsonSourceFileParser); falls back to best column-overlap match so // freeform table keys (XML node names) still resolve correctly. private static DataTable? ResolveTableForEntity( ParsedImportDataSetDTO parsed, int entityId, List details) { if (parsed.Tables.TryGetValue(entityId.ToString(), out var exact)) return exact; var sourceFields = details .Where(d => d.IceMapDetailsValueType is 0 or 4 or 5 && !string.IsNullOrWhiteSpace(d.IceMapDetailsSourceFieldName)) .Select(d => d.IceMapDetailsSourceFieldName!.Trim().ToUpperInvariant()) .ToHashSet(); if (sourceFields.Count == 0) return parsed.Tables.Values.FirstOrDefault(); return parsed.Tables.Values .Select(t => (Table: t, Score: sourceFields.Count(f => t.Columns.Contains(f)))) .Where(x => x.Score > 0) .OrderByDescending(x => x.Score) .Select(x => x.Table) .FirstOrDefault(); } // ── Stage 3 — Validate ─────────────────────────────────────────────── public async Task> ValidateAsync( List rows, ParsedImportDataSetDTO parsedData, IceMapDTO mapDefinition, LoginDTO login, CancellationToken ct) { using var section = GB5Trace.BeginSection("ice-import-validate"); var enabledKeys = await _validationRuleDal .GetEnabledValidatorKeysAsync(mapDefinition.IceMapId, login, ct) .ConfigureAwait(false); var enabledSet = new HashSet(enabledKeys, StringComparer.OrdinalIgnoreCase); var applicableValidators = _validators .Where(v => _defaultValidatorKeys.Contains(v.Key) || enabledSet.Contains(v.Key)) .ToList(); GB5Trace.Step("validate-rows", new { ValidatorCount = applicableValidators.Count }); var results = new List(); foreach (var row in FlattenTree(rows)) { var errors = new List(row.FieldWarnings); foreach (var validator in applicableValidators) { try { var found = await validator .ValidateAsync(row, mapDefinition, parsedData, login, ct) .ConfigureAwait(false); errors.AddRange(found); } catch (Exception ex) { GB5Trace.MarkFailed($"ice-import-validator-{validator.Key}-failed", ex); _logger.LogError(ex, "IceImport: validator {Key} failed for Row {RowNumber}", validator.Key, row.RowNumber); } } if (errors.Count > 0) { results.Add(new IceImportRowValidationResultDTO { CorrelationKey = row.CorrelationKey, RowNumber = row.RowNumber, EntityCode = row.EntityCode, Errors = errors }); } } return results; } private static IEnumerable FlattenTree(IEnumerable rows) { foreach (var row in rows) { yield return row; foreach (var childList in row.Children.Values) foreach (var descendant in FlattenTree(childList)) yield return descendant; } } // ── Stage 4 — Commit ───────────────────────────────────────────────── public async Task CommitAsync( List rows, List validationResults, IceMapDTO mapDefinition, LoginDTO login, long runId, CancellationToken ct) { using var section = GB5Trace.BeginSection("ice-import-commit"); // Keyed by CorrelationKey, not RowNumber — RowNumber is only unique WITHIN one entity // level (header rows and detail rows are both numbered 1..N independently). var errorKeys = validationResults .Where(r => r.HasErrors) .Select(r => r.CorrelationKey) .ToHashSet(); var commitResult = new IceImportCommitResultDTO { ValidationErrors = validationResults.SelectMany(r => r.Errors).ToList(), TotalRows = FlattenTree(rows).Count() }; // Root-level rows all share one adapter (selected the same way CommitRowRecursivelyAsync // selects it per row) — resolved once up front so the batch-vs-per-row decision // (MWEBSERVICE.REQUESTSCHEMATYPE) can be made before any row is committed. var headerAdapter = rows.Count > 0 ? _persistenceAdapters.FirstOrDefault(a => a.EntityId == rows[0].EntityId) ?? _persistenceAdapters.FirstOrDefault(a => a.EntityId == -1) : null; var requestSchemaType = headerAdapter is not null ? await headerAdapter.GetSaveRequestSchemaTypeAsync(mapDefinition, login, ct).ConfigureAwait(false) : 1; if (headerAdapter is not null && requestSchemaType == 0) { await CommitRowsAsBatchAsync(rows, errorKeys, mapDefinition, login, runId, headerAdapter, ct, commitResult) .ConfigureAwait(false); } else { foreach (var root in rows) await CommitRowRecursivelyAsync(root, null, errorKeys, mapDefinition, login, runId, ct, commitResult) .ConfigureAwait(false); } return commitResult; } private async Task CommitRowRecursivelyAsync( MappedEntityRowDTO row, int? grandparentId, HashSet errorKeys, IceMapDTO mapDefinition, LoginDTO login, long runId, CancellationToken ct, IceImportCommitResultDTO commitResult) { if (errorKeys.Contains(row.CorrelationKey)) { commitResult.SkippedCount++; commitResult.RowResults.Add(new IceImportRowCommitResultDTO { CorrelationKey = row.CorrelationKey, RowNumber = row.RowNumber, EntityCode = row.EntityCode, Success = false, ErrorMessage = "Skipped — row has one or more Error-severity validation findings." }); return; // children of a skipped parent are never committed either } var adapter = _persistenceAdapters.FirstOrDefault(a => a.EntityId == row.EntityId) ?? _persistenceAdapters.FirstOrDefault(a => a.EntityId == -1); if (adapter is null) { commitResult.FailedCount++; commitResult.RowResults.Add(new IceImportRowCommitResultDTO { CorrelationKey = row.CorrelationKey, RowNumber = row.RowNumber, EntityCode = row.EntityCode, Success = false, ErrorMessage = $"No IEntityPersistenceAdapter registered for EntityId {row.EntityId}." }); return; } var embeddedChildEntityIds = EmbedObjectListChildren(row, mapDefinition, errorKeys); PersistResultDTO result; try { var decision = await _upsertDecisionService .DecideAsync(row, mapDefinition, adapter, login, runId, ct) .ConfigureAwait(false); GB5Trace.Step("commit-row", new { row.RowNumber, row.EntityId, decision.IsUpdate }); result = decision.IsUpdate ? await adapter.UpdateAsync(row, decision.ExistingId!.Value, mapDefinition, login, runId, ct).ConfigureAwait(false) : await adapter.SaveAsync(row, mapDefinition, login, runId, ct).ConfigureAwait(false); } catch (Exception ex) { GB5Trace.MarkFailed("ice-import-commit-row-failed", ex); _logger.LogError(ex, "IceImport: commit failed for EntityId {EntityId} Row {RowNumber}", row.EntityId, row.RowNumber); result = PersistResultDTO.Fail(ex.Message); } row.ResolvedTargetId = result.TargetEntityId; commitResult.RowResults.Add(new IceImportRowCommitResultDTO { CorrelationKey = row.CorrelationKey, RowNumber = row.RowNumber, EntityCode = row.EntityCode, Success = result.Success, ErrorMessage = result.ErrorMessage, TargetEntityId = result.TargetEntityId, PostData = result.PostData }); if (result.Success) commitResult.SuccessCount++; else commitResult.FailedCount++; if (!result.Success || row.Children.Count == 0) return; foreach (var (childEntityId, childRows) in row.Children) { // AddonEntityId is a pure technical sentinel — AddonFieldAccumulator (MapAndTransformAsync) // groups every IceMapDetailsAddonFields-flagged detail into this synthetic child purely so // the generic Object/List embedding above can fold it into the header row's own JSON // payload as e.g. "PayRevisionAddon": [...]. It is not a real business detail row, so it // gets no RowResult of its own — the header row's own Success already covers it, and // logging it separately just produces a confusing extra TICEIMPORTRUNROWRESULT entry // (EntityCode = the raw numeric sentinel) for what was really one single save call. if (childEntityId == AddonEntityId) continue; if (embeddedChildEntityIds.Contains(childEntityId)) { // Already embedded into row.Fields[ListMemberName] above and saved as part of this // row's own payload — record each child's outcome without a separate persistence // call, using the same skip semantics as the top-of-method Error-severity check. foreach (var child in childRows) { var skipped = errorKeys.Contains(child.CorrelationKey); commitResult.RowResults.Add(new IceImportRowCommitResultDTO { CorrelationKey = child.CorrelationKey, RowNumber = child.RowNumber, EntityCode = child.EntityCode, Success = !skipped, ErrorMessage = skipped ? "Skipped — row has one or more Error-severity validation findings." : null }); if (skipped) commitResult.SkippedCount++; else commitResult.SuccessCount++; } continue; } var parentField = _hierarchyMapper.GetParentIdFieldName(mapDefinition, childEntityId); var firstParentField = _hierarchyMapper.GetFirstParentIdFieldName(mapDefinition, childEntityId); foreach (var child in childRows) { if (!string.IsNullOrEmpty(parentField)) child.Fields[parentField] = row.ResolvedTargetId; if (!string.IsNullOrEmpty(firstParentField) && grandparentId is not null) child.Fields[firstParentField] = grandparentId; await CommitRowRecursivelyAsync( child, row.ResolvedTargetId, errorKeys, mapDefinition, login, runId, ct, commitResult) .ConfigureAwait(false); } } } // Embeds every Object/List (DataType=5) child directly into row.Fields[ListMemberName] as a // nested array (GB4's AssignValueParentAndSave "case 5: object" pattern) — shared by both the // per-row and batch commit paths below, since batching only changes how the already-assembled // header payload is dispatched, never how it's built. Every other kind of child relationship // (ValueType=6/8 ParentId assignment) is untouched here. private HashSet EmbedObjectListChildren( MappedEntityRowDTO row, IceMapDTO mapDefinition, HashSet errorKeys) { var embeddedChildEntityIds = new HashSet(); foreach (var relationship in _hierarchyMapper.GetObjectListRelationships(mapDefinition)) { if (relationship.ParentEntityId != row.EntityId) continue; if (!row.Children.TryGetValue(relationship.ChildEntityId, out var childRows)) continue; row.Fields[relationship.ListMemberName] = childRows .Where(c => !errorKeys.Contains(c.CorrelationKey)) .Select(c => c.Fields) .ToList(); embeddedChildEntityIds.Add(relationship.ChildEntityId); } return embeddedChildEntityIds; } // MWEBSERVICE.REQUESTSCHEMATYPE = 0 path — every non-skipped root row is submitted as ONE list // in a single SaveBatchAsync call instead of one call per row. Batch mode is insert-only (no // FindByNaturalKeyAsync/UpsertDecisionService per row) and only ever one aggregate // IceImportRowCommitResultDTO is recorded for the whole call, matching MWEBSERVICE's own // "call the URL once with all data as a list" contract for this schema type. ParentId-linked // (non-embedded) child rows are out of scope — they still need their own parent's freshly // resolved Id, which a single bulk call can't hand back per row, so this only ever applies at // the root/header level CommitAsync calls it from. private async Task CommitRowsAsBatchAsync( List rows, HashSet errorKeys, IceMapDTO mapDefinition, LoginDTO login, long runId, IEntityPersistenceAdapter adapter, CancellationToken ct, IceImportCommitResultDTO commitResult) { var rowsToSave = new List(); foreach (var row in rows) { if (errorKeys.Contains(row.CorrelationKey)) { commitResult.SkippedCount++; commitResult.RowResults.Add(new IceImportRowCommitResultDTO { CorrelationKey = row.CorrelationKey, RowNumber = row.RowNumber, EntityCode = row.EntityCode, Success = false, ErrorMessage = "Skipped — row has one or more Error-severity validation findings." }); continue; } EmbedObjectListChildren(row, mapDefinition, errorKeys); rowsToSave.Add(row); } if (rowsToSave.Count == 0) return; // every row was skipped — nothing to submit GB5Trace.Step("commit-batch", new { RowCount = rowsToSave.Count, mapDefinition.IceMapEntityId }); PersistResultDTO result; try { result = await adapter.SaveBatchAsync(rowsToSave, mapDefinition, login, runId, ct).ConfigureAwait(false); } catch (Exception ex) { GB5Trace.MarkFailed("ice-import-commit-batch-failed", ex); _logger.LogError(ex, "IceImport: batch commit failed for EntityId {EntityId} ({RowCount} rows)", mapDefinition.IceMapEntityId, rowsToSave.Count); result = PersistResultDTO.Fail(ex.Message); } commitResult.RowResults.Add(new IceImportRowCommitResultDTO { CorrelationKey = Guid.NewGuid(), RowNumber = 0, EntityCode = rowsToSave[0].EntityCode, Success = result.Success, ErrorMessage = result.Success ? null : $"Batch save failed for {rowsToSave.Count} row(s): {result.ErrorMessage}", TargetEntityId = result.TargetEntityId, PostData = result.PostData }); if (result.Success) commitResult.SuccessCount += rowsToSave.Count; else commitResult.FailedCount += rowsToSave.Count; } // ── Stage 4 (Update variant) — Commit as dynamic UPDATE ───────────── public async Task CommitAsUpdateAsync( List rows, List validationResults, IceMapDTO mapDefinition, LoginDTO login, CancellationToken ct) { using var section = GB5Trace.BeginSection("ice-import-commit-update"); var errorKeys = validationResults .Where(r => r.HasErrors) .Select(r => r.CorrelationKey) .ToHashSet(); var commitResult = new IceImportCommitResultDTO { ValidationErrors = validationResults.SelectMany(r => r.Errors).ToList(), TotalRows = rows.Count }; // Root rows only — every row here is expected to share the same EntityId (single-entity // Update, no BOM-style child/detail-array support), so the field map is fetched once and // reused for every row. var fieldMapCache = new Dictionary>(); foreach (var row in rows) { if (errorKeys.Contains(row.CorrelationKey)) { commitResult.SkippedCount++; commitResult.RowResults.Add(new IceImportRowCommitResultDTO { CorrelationKey = row.CorrelationKey, RowNumber = row.RowNumber, EntityCode = row.EntityCode, Success = false, ErrorMessage = "Skipped — row has one or more Error-severity validation findings." }); continue; } try { if (!fieldMapCache.TryGetValue(row.EntityId, out var fieldMap)) { fieldMap = await _fieldMapDal.GetFieldMapAsync(row.EntityId, login, ct).ConfigureAwait(false); fieldMapCache[row.EntityId] = fieldMap; } var result = await ExecuteRowUpdateAsync(row, mapDefinition, fieldMap, login, ct) .ConfigureAwait(false); commitResult.RowResults.Add(new IceImportRowCommitResultDTO { CorrelationKey = row.CorrelationKey, RowNumber = row.RowNumber, EntityCode = row.EntityCode, Success = result.Success, ErrorMessage = result.ErrorMessage, TargetEntityId = result.TargetEntityId }); if (result.Success) commitResult.SuccessCount++; else commitResult.FailedCount++; } catch (Exception ex) { GB5Trace.MarkFailed("ice-import-commit-update-row-failed", ex); _logger.LogError(ex, "IceImport: update commit failed for EntityId {EntityId} Row {RowNumber}", row.EntityId, row.RowNumber); commitResult.FailedCount++; commitResult.RowResults.Add(new IceImportRowCommitResultDTO { CorrelationKey = row.CorrelationKey, RowNumber = row.RowNumber, EntityCode = row.EntityCode, Success = false, ErrorMessage = ex.Message }); } } return commitResult; } // Builds and runs one parameterized UPDATE SET ... WHERE ... for a single row. // IceMapDetailsDTO.IceMapDetailsIsUniqueKey decides SET vs WHERE per field: false -> WHERE, // true -> SET (see IIceImportPipeline.CommitAsUpdateAsync's doc comment for why this is the // opposite of FindByNaturalKeyAsync's own convention). Values are always bound as parameters — // never string-concatenated into the SQL text, unlike GB4's IceImportNewBLL.UpdateEntity, which // built "field='" + value + "'" directly and is exactly the kind of injection-prone pattern // this rewrite must not repeat. private async Task ExecuteRowUpdateAsync( MappedEntityRowDTO row, IceMapDTO mapDefinition, List fieldMap, LoginDTO login, CancellationToken ct) { if (fieldMap.Count == 0) return PersistResultDTO.Fail($"No MENTITY/DBOBJECT/DBOBJECTFIELDS field map resolved for EntityId {row.EntityId}."); var tableName = fieldMap[0].TableName; if (string.IsNullOrWhiteSpace(tableName) || !_validIdentifier.IsMatch(tableName)) return PersistResultDTO.Fail($"EntityId {row.EntityId} resolved an unsafe or missing table name."); var fieldMapByMemberId = fieldMap.ToDictionary(f => f.EntityMemberId); var details = mapDefinition.IceMapDetailsArray .Where(d => d.DestinationMemberEntityId == row.EntityId) .ToList(); var parameters = new Dapper.DynamicParameters(); var setClauses = new List(); var whereClauses = new List(); var paramIndex = 0; foreach (var detail in details) { if (!fieldMapByMemberId.TryGetValue(detail.DestinationMemberId, out var fieldInfo)) continue; // no physical column resolved for this member (e.g. an object/list member) if (string.IsNullOrWhiteSpace(fieldInfo.DbFieldName) || !_validIdentifier.IsMatch(fieldInfo.DbFieldName)) continue; var memberName = detail.DestinationMemberMemberName; if (string.IsNullOrEmpty(memberName) || !row.Fields.TryGetValue(memberName, out var value)) continue; var paramName = $"p{paramIndex++}"; parameters.Add(paramName, value); // IceMapDetailsIsUniqueKey == false -> WHERE (identifies the row), == true -> SET (the // value being written) — confirmed against real MEMPLOYEE map data (EMPLOYEEID is the // only IsUniqueKey=0 field and must be the WHERE condition) and matches GB4's own // IceImportNewBLL.UpdateEntity exactly (IsUniquekey != 0 -> SET, == 0 -> WHERE). This is // the OPPOSITE of IEntityPersistenceAdapter.FindByNaturalKeyAsync's own convention for // the same underlying flag — that's a different mechanism (GET-based upsert decision // for CommitImport) serving a different purpose; do not "fix" this to match it. if (!detail.IceMapDetailsIsUniqueKey) whereClauses.Add($"{fieldInfo.DbFieldName} = @{paramName}"); else setClauses.Add($"{fieldInfo.DbFieldName} = @{paramName}"); } if (setClauses.Count == 0) return PersistResultDTO.Fail("No non-key field resolved a value to update — nothing to SET."); if (whereClauses.Count == 0) return PersistResultDTO.Fail( "No IceMapDetailsIsUniqueKey=false field resolved a value — refusing to UPDATE with no WHERE condition."); var sql = $"UPDATE {tableName} SET {string.Join(", ", setClauses)} WHERE {string.Join(" AND ", whereClauses)}"; var affected = await _queryExecutor .ExecuteAsync(login, sql, parameters, cancellationToken: ct) .ConfigureAwait(false); if (affected <= 0) return PersistResultDTO.Fail("UPDATE matched zero rows — check the unique-key field values."); // Addon table — MENTITY.ENTITYID = AddonEntityId (-1399999424) is the well-known "Addon // entity" sentinel. AddonFieldAccumulator (MapAndTransformAsync) groups every detail row // whose IceMapDetailsAddonFields is set into a separate child MappedEntityRowDTO // (EntityId = AddonEntityId), storing them as a JSON bag under // Fields[AddonFieldAccumulator.AddonFieldKey]. Every addon field is a SET target — none of // them identify the row — so the addon UPDATE reuses the SAME WHERE clause/bound parameters // built above for the primary table, per the confirmed requirement. if (row.Children.TryGetValue(AddonEntityId, out var addonChildren) && addonChildren.Count > 0) { var addonResult = await ExecuteAddonUpdateAsync( mapDefinition, addonChildren[0], whereClauses, parameters, paramIndex, login, ct) .ConfigureAwait(false); if (addonResult is not null) return addonResult; // addon step failed — surfaces as this row's overall result } return PersistResultDTO.Ok(row.EntityId); } // MENTITY.ENTITYID sentinel meaning "no Addon table" / "this is the Addon entity" — // MICEMAPDETAILS rows whose DestinationMemberEntityId is this value carry an // IceMapDetailsAddonFields (physical column name) instead of a normal DBFIELDID mapping. private const int AddonEntityId = -1399999424; // Returns null on success (matching "no failure to report"), or a Fail result if the addon // table's UPDATE couldn't be resolved/executed. private async Task ExecuteAddonUpdateAsync( IceMapDTO mapDefinition, MappedEntityRowDTO addonRow, List whereClauses, Dapper.DynamicParameters parameters, int paramIndex, LoginDTO login, CancellationToken ct) { if (!addonRow.Fields.TryGetValue(IceImportBLL.AddonFieldAccumulator.AddonFieldAccumulator.AddonFieldKey, out var bagObj) || bagObj is not string bagJson || string.IsNullOrWhiteSpace(bagJson)) return null; // no addon fields actually resolved for this row — nothing to update var addonTableName = await _fieldMapDal .GetAddonTableNameAsync(mapDefinition.IceMapEntityId, login, ct).ConfigureAwait(false); if (string.IsNullOrWhiteSpace(addonTableName) || !_validIdentifier.IsMatch(addonTableName)) return PersistResultDTO.Fail( $"No ADDONDBOBJECTID resolved for EntityId {mapDefinition.IceMapEntityId} — cannot update Addon fields."); var bag = System.Text.Json.JsonSerializer.Deserialize>(bagJson) ?? new Dictionary(); var addonSetClauses = new List(); foreach (var (fieldName, rawValue) in bag) { if (string.IsNullOrWhiteSpace(fieldName) || !_validIdentifier.IsMatch(fieldName)) continue; // AddonFieldAccumulator serializes via System.Text.Json, so round-tripping a // Dictionary hands back boxed JsonElement values, not the original // CLR types — unwrap before binding as a SQL parameter. var value = rawValue is System.Text.Json.JsonElement element ? ExtractJsonElementValue(element) : rawValue; var paramName = $"p{paramIndex++}"; parameters.Add(paramName, value); addonSetClauses.Add($"{fieldName} = @{paramName}"); } if (addonSetClauses.Count == 0) return null; // every addon field name failed the identifier check — nothing safe to update var addonSql = $"UPDATE {addonTableName} SET {string.Join(", ", addonSetClauses)} WHERE {string.Join(" AND ", whereClauses)}"; var addonAffected = await _queryExecutor .ExecuteAsync(login, addonSql, parameters, cancellationToken: ct) .ConfigureAwait(false); return addonAffected > 0 ? null : PersistResultDTO.Fail("Addon table UPDATE matched zero rows — check the unique-key field values."); } private static object? ExtractJsonElementValue(System.Text.Json.JsonElement element) => element.ValueKind switch { System.Text.Json.JsonValueKind.String => element.GetString(), System.Text.Json.JsonValueKind.Number => element.TryGetInt64(out var l) ? l : element.GetDouble(), System.Text.Json.JsonValueKind.True => true, System.Text.Json.JsonValueKind.False => false, System.Text.Json.JsonValueKind.Null => null, _ => element.GetRawText() }; // ── Convenience — full pipeline ───────────────────────────────────── public async Task ExecuteAsync( Stream file, string fileName, int iceMapId, LoginDTO login, long runId, CancellationToken ct) { var mapDefinition = await LoadMapDefinitionAsync(iceMapId, login).ConfigureAwait(false); var parsed = await ParseWithMapAsync(file, fileName, mapDefinition, ct).ConfigureAwait(false); var rows = await MapAndTransformAsync(parsed, mapDefinition, login, ct).ConfigureAwait(false); var validationResults = await ValidateAsync(rows, parsed, mapDefinition, login, ct).ConfigureAwait(false); return await CommitAsync(rows, validationResults, mapDefinition, login, runId, ct).ConfigureAwait(false); } }