diff --git a/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/ColumnGridResourceIT.java b/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/ColumnGridResourceIT.java index 8d597cdbfb50..48203438108d 100644 --- a/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/ColumnGridResourceIT.java +++ b/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/ColumnGridResourceIT.java @@ -13,7 +13,6 @@ import java.time.Duration; import java.util.List; import org.junit.jupiter.api.BeforeAll; -import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.parallel.Execution; @@ -53,17 +52,10 @@ import org.openmetadata.sdk.fluent.Tables; import org.openmetadata.sdk.network.HttpMethod; -// TEMPORARILY DISABLED — the metadataStatus aggregation on this endpoint reproducibly fails -// with [search_phase_execution_exception] all shards failed on both postgres+ES+redis (single -// failure on test_getColumnGrid_withMetadataStatusIncomplete) AND postgres+OpenSearch (the same -// query crashes the OS container, then 15 follow-up tests in the class fail with Connection -// refused). Same behavior on PR #28100 with and without the cache changes, so it is a -// pre-existing aggregator bug, not a cache regression. The ES Java client swallows the -// underlying `caused_by`, so root-causing the actual ES-side error requires response-body -// logging that is not wired up yet. Re-enable once the underlying aggregator/index-mapping -// issue is fixed in a follow-up. See PR #28100 history and CI run 25940411417 for context. -@Disabled( - "ColumnGrid metadataStatus aggregation crashes ES/OS — pre-existing flake, follow-up needed") +// Re-enabled with #26824: the metadataStatus crash came from the per-document filter query +// (wildcard/exists on flat-object columns.description/columns.tags), which ES 7.17 and OpenSearch +// rejected with `search_phase_execution_exception ... all shards failed`. That push-down is gone — +// status is now filtered on the aggregate grouped item — so the crashing query no longer runs. @Execution(ExecutionMode.CONCURRENT) @ExtendWith(TestNamespaceExtension.class) public class ColumnGridResourceIT { @@ -314,6 +306,11 @@ void test_getColumnGrid_withMetadataStatusMissing(TestNamespace ns) throws Excep assertNotNull(response); assertNotNull(response.getColumns()); + assertAllRowsHaveStatus(response, MetadataStatus.MISSING); + assertEquals( + response.getColumns().size(), + response.getTotalUniqueColumns(), + "totalUniqueColumns must reflect the filtered set, not the unfiltered total"); } @Test @@ -328,6 +325,9 @@ void test_getColumnGrid_withMetadataStatusComplete(TestNamespace ns) throws Exce assertNotNull(response); assertNotNull(response.getColumns()); + assertFalse(response.getColumns().isEmpty(), "the COMPLETE column should be returned"); + // The reported bug (#26824): COMPLETE must not surface MISSING/INCOMPLETE/INCONSISTENT rows. + assertAllRowsHaveStatus(response, MetadataStatus.COMPLETE); } @Test @@ -342,6 +342,8 @@ void test_getColumnGrid_withMetadataStatusIncomplete(TestNamespace ns) throws Ex assertNotNull(response); assertNotNull(response.getColumns()); + assertFalse(response.getColumns().isEmpty(), "the INCOMPLETE column should be returned"); + assertAllRowsHaveStatus(response, MetadataStatus.INCOMPLETE); } @Test @@ -422,6 +424,125 @@ void test_getColumnGrid_withMetadataStatusInconsistent(TestNamespace ns) throws assertNotNull(response); assertNotNull(response.getColumns()); + assertFalse(response.getColumns().isEmpty(), "the INCONSISTENT column should be returned"); + assertAllRowsHaveStatus(response, MetadataStatus.INCONSISTENT); + assertTrue( + response.getColumns().stream().allMatch(ColumnGridItem::getHasVariations), + "INCONSISTENT rows have metadata variations across occurrences"); + } + + @Test + void test_getColumnGrid_metadataStatusPaginationCountsAreConsistent(TestNamespace ns) + throws Exception { + OpenMetadataClient client = SdkClients.adminClient(); + DatabaseService service = DatabaseServiceTestFactory.createPostgres(ns); + DatabaseSchema schema = DatabaseSchemaTestFactory.createSimple(ns, service); + + // Three COMPLETE columns and one MISSING column in the same service. + for (int i = 0; i < 3; i++) { + Column complete = + Columns.build(ns.prefix("paged_complete_" + i)) + .withType(ColumnDataType.BIGINT) + .withDescription("has description") + .withTags(List.of(new TagLabel().withTagFQN("PII.Sensitive"))) + .create(); + Tables.create() + .name(ns.prefix("paged_table_" + i)) + .inSchema(schema.getFullyQualifiedName()) + .withColumns(List.of(complete)) + .execute(); + } + Column missing = + Columns.build(ns.prefix("paged_missing")).withType(ColumnDataType.BIGINT).create(); + Tables.create() + .name(ns.prefix("paged_table_missing")) + .inSchema(schema.getFullyQualifiedName()) + .withColumns(List.of(missing)) + .execute(); + + waitForSearchIndexRefresh(ns); + + ColumnGridResponse page1 = + getColumnGrid( + client, + "size=2&entityTypes=table&metadataStatus=COMPLETE&serviceName=" + service.getName()); + + // totalUniqueColumns must count only the 3 COMPLETE columns (not 4), and the page must respect + // the requested size — the pagination half of #26824. + assertEquals(3, page1.getTotalUniqueColumns()); + assertEquals(2, page1.getColumns().size()); + assertAllRowsHaveStatus(page1, MetadataStatus.COMPLETE); + assertNotNull(page1.getCursor(), "a second page of COMPLETE columns remains"); + + ColumnGridResponse page2 = + getColumnGrid( + client, + "size=2&entityTypes=table&metadataStatus=COMPLETE&serviceName=" + + service.getName() + + "&cursor=" + + URLEncoder.encode(page1.getCursor(), StandardCharsets.UTF_8)); + + assertEquals(3, page2.getTotalUniqueColumns()); + assertEquals(1, page2.getColumns().size(), "the last page holds the remaining COMPLETE column"); + assertAllRowsHaveStatus(page2, MetadataStatus.COMPLETE); + } + + @Test + void test_getColumnGrid_metadataStatusWithColumnNamePattern(TestNamespace ns) throws Exception { + OpenMetadataClient client = SdkClients.adminClient(); + DatabaseService service = DatabaseServiceTestFactory.createPostgres(ns); + DatabaseSchema schema = DatabaseSchemaTestFactory.createSimple(ns, service); + + // Two COMPLETE columns; only one name contains "alpha". + String matchName = ns.prefix("alpha_amount"); + String otherName = ns.prefix("zzz_other"); + for (String colName : List.of(matchName, otherName)) { + Column col = + Columns.build(colName) + .withType(ColumnDataType.BIGINT) + .withDescription("has description") + .withTags(List.of(new TagLabel().withTagFQN("PII.Sensitive"))) + .create(); + Tables.create() + .name(ns.prefix("pat_" + colName)) + .inSchema(schema.getFullyQualifiedName()) + .withColumns(List.of(col)) + .execute(); + } + waitForSearchIndexRefresh(ns); + + ColumnGridResponse response = + getColumnGrid( + client, + "entityTypes=table&metadataStatus=COMPLETE&columnNamePattern=alpha&serviceName=" + + service.getName()); + + // Combining columnNamePattern with a status filter must honor the pattern per column: + // the non-matching "zzz_other" column must not leak in (regression for the _source-scan path). + assertNotNull(response); + assertFalse(response.getColumns().isEmpty(), "the matching COMPLETE column should be returned"); + assertAllRowsHaveStatus(response, MetadataStatus.COMPLETE); + assertTrue( + response.getColumns().stream() + .allMatch(c -> c.getColumnName().toLowerCase().contains("alpha")), + "only columns whose name matches the pattern should be returned"); + assertEquals( + response.getColumns().size(), + response.getTotalUniqueColumns(), + "totalUniqueColumns must not include pattern-mismatched columns"); + } + + /** + * Every returned row must carry the requested aggregate status — the core guarantee of #26824 + * (before the fix a COMPLETE/INCOMPLETE filter leaked rows of other statuses). + */ + private void assertAllRowsHaveStatus(ColumnGridResponse response, MetadataStatus expected) { + for (ColumnGridItem item : response.getColumns()) { + assertEquals( + expected, + item.getMetadataStatus(), + "column '" + item.getColumnName() + "' should have status " + expected); + } } @Test @@ -1424,6 +1545,7 @@ private DatabaseService createTableWithFullMetadata(TestNamespace ns) { Columns.build("full_metadata_id") .withType(ColumnDataType.BIGINT) .withDescription("Primary key with description") + .withTags(List.of(new TagLabel().withTagFQN("PII.Sensitive"))) .create(); Tables.create() diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/jdbi3/ColumnRepository.java b/openmetadata-service/src/main/java/org/openmetadata/service/jdbi3/ColumnRepository.java index 88373c992ac0..518a3930eb6c 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/jdbi3/ColumnRepository.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/jdbi3/ColumnRepository.java @@ -40,7 +40,6 @@ import org.openmetadata.schema.EntityInterface; import org.openmetadata.schema.api.data.BulkColumnUpdatePreview; import org.openmetadata.schema.api.data.BulkColumnUpdateRequest; -import org.openmetadata.schema.api.data.ColumnGridItem; import org.openmetadata.schema.api.data.ColumnGridResponse; import org.openmetadata.schema.api.data.ColumnMetadata; import org.openmetadata.schema.api.data.ColumnOccurrence; @@ -97,39 +96,10 @@ public ColumnRepository(Authorizer authorizer, SearchClient searchClient) { public ColumnGridResponse getColumnGridPaginated( SecurityContext securityContext, ColumnAggregator.ColumnAggregationRequest request) throws IOException { - ColumnGridResponse response = columnAggregator.aggregateColumns(request); - - if (Boolean.TRUE.equals(request.getHasConflicts())) { - response.setColumns( - response.getColumns().stream() - .filter(ColumnGridItem::getHasVariations) - .collect(Collectors.toList())); - } - - if (Boolean.TRUE.equals(request.getHasMissingMetadata())) { - response.setColumns( - response.getColumns().stream() - .filter(this::hasMissingMetadata) - .collect(Collectors.toList())); - } - - // Filter by INCONSISTENT status (requires post-aggregation filtering) - if ("INCONSISTENT".equalsIgnoreCase(request.getMetadataStatus())) { - response.setColumns( - response.getColumns().stream() - .filter(ColumnGridItem::getHasVariations) - .collect(Collectors.toList())); - } - - return response; - } - - private boolean hasMissingMetadata(ColumnGridItem item) { - return item.getGroups().stream() - .anyMatch( - group -> - (group.getDescription() == null || group.getDescription().isEmpty()) - || (group.getTags() == null || group.getTags().isEmpty())); + // Row-level filters (metadataStatus / hasConflicts / hasMissingMetadata) are applied inside the + // aggregator over the fully-grouped columns, before pagination, so page counts and per-page + // size stay correct (#26824). Nothing to post-process here. + return columnAggregator.aggregateColumns(request); } public Column getColumnByFQN( diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/resources/columns/ColumnResource.java b/openmetadata-service/src/main/java/org/openmetadata/service/resources/columns/ColumnResource.java index cd40605a90b3..65103fecd31c 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/resources/columns/ColumnResource.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/resources/columns/ColumnResource.java @@ -319,13 +319,15 @@ public Response getColumnGrid( boolean hasMissingMetadata, @Parameter( description = - "Filter by metadata status: MISSING (no description AND no tags), " + "Filter by aggregate metadata status of a column across all its occurrences: " + + "MISSING (no description AND no tags), " + "INCOMPLETE (has description OR tags, but not both), " - + "COMPLETE (has both description AND tags)", + + "COMPLETE (has both description AND tags), " + + "INCONSISTENT (occurrences disagree on description/tags)", schema = @Schema( type = "string", - allowableValues = {"MISSING", "INCOMPLETE", "COMPLETE"})) + allowableValues = {"MISSING", "INCOMPLETE", "COMPLETE", "INCONSISTENT"})) @QueryParam("metadataStatus") String metadataStatus, @Parameter( diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java index f673b20c33e2..ab2ec73f3eee 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java @@ -16,9 +16,13 @@ import com.fasterxml.jackson.core.type.TypeReference; import java.io.IOException; import java.nio.charset.StandardCharsets; +import java.util.ArrayList; import java.util.Base64; +import java.util.Comparator; import java.util.List; +import java.util.Locale; import java.util.Map; +import org.openmetadata.schema.api.data.ColumnGridItem; import org.openmetadata.schema.api.data.ColumnGridResponse; import org.openmetadata.schema.utils.JsonUtils; import org.slf4j.Logger; @@ -115,6 +119,107 @@ static int toIntSaturating(long value) { return (int) value; } + /** + * True if a request carries a row-level filter — one that acts on the aggregate status of a + * grouped column (metadataStatus, hasConflicts, hasMissingMetadata) rather than on individual + * documents. These cannot be pushed into the search query (the aggregate status is only known + * after grouping occurrences), so they are applied via {@link #paginateFilteredItems}. + */ + static boolean hasRowLevelFilter(ColumnAggregationRequest request) { + return !nullOrEmptyStr(request.getMetadataStatus()) + || Boolean.TRUE.equals(request.getHasConflicts()) + || Boolean.TRUE.equals(request.getHasMissingMetadata()); + } + + /** True if the grouped column satisfies every active row-level filter on the request. */ + static boolean matchesRowFilters(ColumnGridItem item, ColumnAggregationRequest request) { + if (Boolean.TRUE.equals(request.getHasConflicts()) + && !Boolean.TRUE.equals(item.getHasVariations())) { + return false; + } + if (Boolean.TRUE.equals(request.getHasMissingMetadata()) && !itemHasMissingMetadata(item)) { + return false; + } + String status = request.getMetadataStatus(); + if (!nullOrEmptyStr(status)) { + String itemStatus = item.getMetadataStatus() != null ? item.getMetadataStatus().value() : ""; + return status.trim().equalsIgnoreCase(itemStatus); + } + return true; + } + + /** + * Drop columns whose name doesn't contain the request's {@code columnNamePattern} + * (case-insensitive). The name wildcard in the search query only scopes which entities are + * scanned; flat-object mapping can't isolate the matching column, so the pattern is enforced per + * column here. Shared by the ES and OS row-filter scans so their pattern semantics can't drift. + */ + static void applyColumnNamePattern( + Map columnsByName, ColumnAggregationRequest request) { + if (nullOrEmptyStr(request.getColumnNamePattern())) { + return; + } + String pattern = request.getColumnNamePattern().toLowerCase(Locale.ROOT); + columnsByName.keySet().removeIf(name -> !name.toLowerCase(Locale.ROOT).contains(pattern)); + } + + /** A column has missing metadata if any of its groups lacks a description or tags. */ + static boolean itemHasMissingMetadata(ColumnGridItem item) { + if (item.getGroups() == null) { + return true; + } + return item.getGroups().stream() + .anyMatch( + group -> + (group.getDescription() == null || group.getDescription().isEmpty()) + || (group.getTags() == null || group.getTags().isEmpty())); + } + + /** + * Apply row-level filters to the fully-grouped column list and paginate the result in memory. + * Filtering the aggregate items (not documents) is what makes a "Complete"/"Incomplete"/… filter + * return only rows whose displayed status matches, and computing totals from the filtered set is + * what keeps the page count and per-page size correct (issue #26824). Ordering is by column name + * (case-insensitive) so the offset cursor is stable across pages. + */ + static ColumnGridResponse paginateFilteredItems( + List allItems, ColumnAggregationRequest request) { + List filtered = + allItems.stream() + .filter(item -> matchesRowFilters(item, request)) + .sorted( + Comparator.comparing( + ColumnGridItem::getColumnName, + Comparator.nullsLast(String.CASE_INSENSITIVE_ORDER))) + .toList(); + + int totalUniqueColumns = filtered.size(); + int totalOccurrences = + filtered.stream() + .mapToInt(item -> item.getTotalOccurrences() != null ? item.getTotalOccurrences() : 0) + .sum(); + + int offset = decodeSearchOffset(request.getCursor()); + int pageSize = request.getSize(); + int fromIndex = Math.min(offset, totalUniqueColumns); + int toIndex = Math.min(offset + pageSize, totalUniqueColumns); + + List page = new ArrayList<>(filtered.subList(fromIndex, toIndex)); + boolean hasMore = toIndex < totalUniqueColumns; + String cursor = hasMore ? encodeSearchOffset(toIndex) : null; + + ColumnGridResponse response = new ColumnGridResponse(); + response.setColumns(page); + response.setTotalUniqueColumns(totalUniqueColumns); + response.setTotalOccurrences(totalOccurrences); + response.setCursor(cursor); + return response; + } + + private static boolean nullOrEmptyStr(String s) { + return s == null || s.isBlank(); + } + /** Phase 1 result: matching column names and the total doc_count summed across buckets. */ record NamesWithCount(List names, long totalDocCount) {} diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java index f91640f948be..04c8d3e82e60 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java @@ -115,9 +115,22 @@ public ColumnGridResponse aggregateColumns(ColumnAggregationRequest request) thr .removeIf(e -> !e.getKey().toLowerCase(Locale.ROOT).contains(pattern)); } + // Row-level filter (metadataStatus / hasConflicts / hasMissingMetadata) acts on the + // aggregate status, so group everything then filter + paginate the items. + if (ColumnAggregator.hasRowLevelFilter(request)) { + List gridItems = ColumnMetadataGrouper.groupColumns(taggedColumns); + return ColumnAggregator.paginateFilteredItems(gridItems, request); + } + return aggregateColumnsWithKnownNames(request, taggedColumns); } + // Row-level filter (no tag filter): materialize all candidate columns, then filter the + // aggregate items and paginate in memory so counts and per-page size stay correct (#26824). + if (ColumnAggregator.hasRowLevelFilter(request)) { + return aggregateColumnsWithRowFilters(request, entityTypes); + } + // Pattern-only path (no tag filter): use terms agg with include regex if (!nullOrEmpty(request.getColumnNamePattern())) { return aggregateColumnsWithPattern(request, entityTypes); @@ -266,6 +279,84 @@ private ColumnGridResponse aggregateColumnsWithPattern( return buildResponse(gridItems, cursor, hasMore, totalUniqueColumns, totalOccurrences); } + /** + * Row-level-filter path (metadataStatus / hasConflicts / hasMissingMetadata, no tag filter). The + * filter acts on the aggregate status of a grouped column, which is only known after grouping all + * of a column's occurrences. We read {@code _source} for the scoped entities in one scan per + * field-path group (the same mechanism the tag path uses), restricting {@code _source} to the + * column and identity fields, then group every column and filter + paginate the items in memory. + * This keeps the page count and per-page size consistent with the filtered result set and reads + * every occurrence of each scanned entity, up to the {@code size(10000)}-entity scan cap. + */ + private ColumnGridResponse aggregateColumnsWithRowFilters( + ColumnAggregationRequest request, List entityTypes) throws IOException { + + Map> fieldPathToEntityTypes = groupByFieldPath(entityTypes); + Map> allColumnsByName = + new TreeMap<>(String.CASE_INSENSITIVE_ORDER); + + for (Map.Entry> entry : fieldPathToEntityTypes.entrySet()) { + String columnNameKeyword = entry.getKey(); + List indexes = resolveIndexNames(entry.getValue()); + String columnFieldPath = INDEX_CONFIGS.get(entry.getValue().getFirst()).columnFieldPath(); + Query query = buildFilters(request, columnNameKeyword, null); + + try { + fetchColumnsFromSource(indexes, query, columnFieldPath, allColumnsByName); + } catch (ElasticsearchException e) { + if (!isIndexNotFoundException(e)) { + logShardFailureDetails(e, indexes, query); + throw e; + } + } + } + + ColumnAggregator.applyColumnNamePattern(allColumnsByName, request); + + List gridItems = ColumnMetadataGrouper.groupColumns(allColumnsByName); + return ColumnAggregator.paginateFilteredItems(gridItems, request); + } + + /** {@code _source} fields needed to build grid items: entity identity + the whole column tree. */ + private List statusScanSourceIncludes(String columnFieldPath) { + return List.of( + "fullyQualifiedName", + "entityType", + "displayName", + "service.name", + "database.name", + "databaseSchema.name", + columnFieldPath); + } + + private void fetchColumnsFromSource( + List indexes, + Query query, + String columnFieldPath, + Map> columnsByName) + throws IOException { + + List includes = statusScanSourceIncludes(columnFieldPath); + SearchRequest searchRequest = + SearchRequest.of( + s -> + s.index(indexes) + .query(query) + .source(src -> src.filter(f -> f.includes(includes))) + .size(10000)); + + SearchResponse response = client.search(searchRequest, JsonData.class); + long totalHits = response.hits().total() != null ? response.hits().total().value() : 0; + if (totalHits > 10000) { + LOG.warn( + "Metadata-status source-fetch matched {} entities; only first 10000 scanned.", totalHits); + } + + for (Hit hit : response.hits().hits()) { + extractMatchingColumnsFromHit(hit, columnFieldPath, Set.of(), true, columnsByName); + } + } + /** * Tag/glossary filter path: the tag-check pass already extracted full column metadata from * _source (only tagged columns are in the map). Just paginate over the in-memory result. @@ -366,7 +457,7 @@ private void fetchColumnsWithTagsFromSource( } for (Hit hit : response.hits().hits()) { - extractMatchingColumnsFromHit(hit, columnFieldPath, targetTags, columnsByName); + extractMatchingColumnsFromHit(hit, columnFieldPath, targetTags, false, columnsByName); } } @@ -374,6 +465,7 @@ private void extractMatchingColumnsFromHit( Hit hit, String columnFieldPath, Set targetTags, + boolean includeAllColumns, Map> columnsByName) { if (hit.source() == null) { return; @@ -396,7 +488,8 @@ private void extractMatchingColumnsFromHit( if (columnsData != null && columnsData.isArray()) { for (JsonNode columnData : columnsData) { String colName = getTextField(columnData, "name"); - if (colName != null && columnHasTargetTag(columnData, targetTags)) { + if (colName != null + && (includeAllColumns || columnHasTargetTag(columnData, targetTags))) { Column column = parseColumn(columnData, entityFQN); columnsByName .computeIfAbsent(colName, k -> new ArrayList<>()) @@ -461,7 +554,6 @@ private Query buildTagFilterQuery(ColumnAggregationRequest request, String colum addSchemaFilter(boolBuilder, request); addDomainFilter(boolBuilder, request); addColumnNamePatternFilter(boolBuilder, request, columnNameKeyword); - addMetadataStatusFilter(boolBuilder, request, columnFieldPath); String tagFQNField = columnNameKeyword.replace(".name.keyword", ".tags.tagFQN"); List allTags = new ArrayList<>(); @@ -562,7 +654,6 @@ private Query buildFilters( addDomainFilter(boolBuilder, request); addColumnNamePatternFilter(boolBuilder, request, columnNameKeyword); addTagFilters(boolBuilder, request, columnNameKeyword, columnNamesFromTagFilter); - addMetadataStatusFilter(boolBuilder, request, columnFieldPath); return Query.of(q -> q.bool(boolBuilder.build())); } @@ -680,62 +771,6 @@ private void addTagFilters( } } - private void addMetadataStatusFilter( - BoolQuery.Builder boolBuilder, ColumnAggregationRequest request, String columnFieldPath) { - if (!nullOrEmpty(request.getMetadataStatus())) { - Query metadataStatusQuery = - buildMetadataStatusFilter(request.getMetadataStatus(), columnFieldPath); - if (metadataStatusQuery != null) { - boolBuilder.filter(metadataStatusQuery); - } - } - } - - private Query buildMetadataStatusFilter(String status, String columnFieldPath) { - String descField = columnFieldPath + ".description"; - String tagsField = columnFieldPath + ".tags"; - - Query hasDesc = hasNonEmptyField(descField); - Query hasTags = existsQuery(tagsField); - Query noDesc = hasEmptyOrMissingField(descField); - Query noTags = notExistsQuery(tagsField); - - return switch (status.toUpperCase()) { - case "MISSING" -> Query.of(q -> q.bool(b -> b.must(noDesc).must(noTags))); - case "INCOMPLETE" -> Query.of( - q -> - q.bool( - b -> - b.should(Query.of(qs -> qs.bool(bs -> bs.must(hasDesc).must(noTags)))) - .should(Query.of(qs -> qs.bool(bs -> bs.must(noDesc).must(hasTags)))) - .minimumShouldMatch("1"))); - case "COMPLETE" -> Query.of(q -> q.bool(b -> b.must(hasDesc).must(hasTags))); - default -> null; - }; - } - - private Query existsQuery(String field) { - return Query.of(q -> q.exists(e -> e.field(field))); - } - - private Query notExistsQuery(String field) { - return Query.of(q -> q.bool(b -> b.mustNot(existsQuery(field)))); - } - - // `wildcard(field, "?*")` matches any doc whose indexed terms include at least one token of - // at least one character — the analyzer-friendly equivalent of "field has non-empty value". - // We can't use `term(field, "")` against analyzed text fields like `columns.description`: the - // field's analyzer produces no tokens for the empty string and ES 7.17 rejects the term query - // with `search_phase_execution_exception ... all shards failed`. Caught by - // ColumnGridResourceIT#test_getColumnGrid_withMetadataStatusIncomplete. - private Query hasNonEmptyField(String field) { - return Query.of(q -> q.wildcard(w -> w.field(field).value("?*"))); - } - - private Query hasEmptyOrMissingField(String field) { - return Query.of(q -> q.bool(b -> b.mustNot(hasNonEmptyField(field)))); - } - /** Phase 1: Get all matching column names using terms agg with include regex (no top_hits). */ private ColumnAggregator.NamesWithCount executeNamesQuery( Query query, List indexes, String columnNameKeyword, String regex) diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java index faac76576f7d..1442e02f2dc5 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java @@ -98,9 +98,22 @@ public ColumnGridResponse aggregateColumns(ColumnAggregationRequest request) thr .removeIf(e -> !e.getKey().toLowerCase(Locale.ROOT).contains(pattern)); } + // Row-level filter (metadataStatus / hasConflicts / hasMissingMetadata) acts on the + // aggregate status, so group everything then filter + paginate the items. + if (ColumnAggregator.hasRowLevelFilter(request)) { + List gridItems = ColumnMetadataGrouper.groupColumns(taggedColumns); + return ColumnAggregator.paginateFilteredItems(gridItems, request); + } + return aggregateColumnsWithKnownNames(request, taggedColumns); } + // Row-level filter (no tag filter): materialize all candidate columns, then filter the + // aggregate items and paginate in memory so counts and per-page size stay correct (#26824). + if (ColumnAggregator.hasRowLevelFilter(request)) { + return aggregateColumnsWithRowFilters(request); + } + // Pattern-only path (no tag filter): use terms agg with include regex if (!nullOrEmpty(request.getColumnNamePattern())) { return aggregateColumnsWithPattern(request); @@ -187,6 +200,78 @@ private ColumnGridResponse aggregateColumnsWithPattern(ColumnAggregationRequest } } + /** + * Row-level-filter path (metadataStatus / hasConflicts / hasMissingMetadata, no tag filter). The + * filter acts on the aggregate status of a grouped column, which is only known after grouping all + * of a column's occurrences. We read {@code _source} for the scoped entities in one scan (the same + * mechanism the tag path uses), restricting {@code _source} to the column and identity fields, + * then group every column and filter + paginate the items in memory. This keeps the page count and + * per-page size consistent with the filtered set and reads every occurrence of each scanned + * entity, up to the {@code size(10000)}-entity scan cap. + */ + private ColumnGridResponse aggregateColumnsWithRowFilters(ColumnAggregationRequest request) + throws IOException { + + Query query = buildFilters(request, null); + Map> allColumnsByName = + new TreeMap<>(String.CASE_INSENSITIVE_ORDER); + + try { + fetchColumnsFromSource(query, allColumnsByName); + ColumnAggregator.applyColumnNamePattern(allColumnsByName, request); + + List gridItems = ColumnMetadataGrouper.groupColumns(allColumnsByName); + return ColumnAggregator.paginateFilteredItems(gridItems, request); + } catch (OpenSearchException e) { + if (isIndexNotFoundException(e)) { + LOG.warn("Search index not found, returning empty results"); + return buildResponse(new ArrayList<>(), null, false, 0, 0); + } + throw e; + } + } + + /** + * Read {@code _source} for the scoped entities in one scan, restricted to the column and identity + * fields, and extract every column. Backs the row-level-filter path (status / conflicts / missing + * metadata), which must group all of a column's occurrences before it can filter on the aggregate + * status. + */ + private void fetchColumnsFromSource( + Query query, Map> columnsByName) throws IOException { + + List resolvedIndexes = resolveIndexNames(); + List includes = + List.of( + "fullyQualifiedName", + "entityType", + "displayName", + "service.name", + "database.name", + "databaseSchema.name", + "columns"); + + SearchRequest searchRequest = + SearchRequest.of( + s -> + s.index(resolvedIndexes) + .query(query) + .source(src -> src.filter(f -> f.includes(includes))) + .size(10000)); + + SearchResponse response = client.search(searchRequest, JsonData.class); + long totalHits = response.hits().total() != null ? response.hits().total().value() : 0; + if (totalHits > 10000) { + LOG.warn( + "Metadata-status source-fetch matched {} entities; only first 10000 scanned.", totalHits); + } + + for (os.org.opensearch.client.opensearch.core.search.Hit hit : + response.hits().hits()) { + extractMatchingColumnsFromHit(hit, Set.of(), true, columnsByName); + } + } + /** * Tag/glossary filter path: the tag-check pass already extracted full column metadata from * _source (only tagged columns are in the map). Just paginate over the in-memory result. @@ -283,13 +368,14 @@ private void fetchColumnsWithTagsFromSource( for (os.org.opensearch.client.opensearch.core.search.Hit hit : response.hits().hits()) { - extractMatchingColumnsFromHit(hit, targetTags, columnsByName); + extractMatchingColumnsFromHit(hit, targetTags, false, columnsByName); } } private void extractMatchingColumnsFromHit( os.org.opensearch.client.opensearch.core.search.Hit hit, Set targetTags, + boolean includeAllColumns, Map> columnsByName) { if (hit.source() == null) { return; @@ -312,7 +398,8 @@ private void extractMatchingColumnsFromHit( if (columnsData != null && columnsData.isArray()) { for (JsonNode columnData : columnsData) { String colName = getTextField(columnData, "name"); - if (colName != null && columnHasTargetTag(columnData, targetTags)) { + if (colName != null + && (includeAllColumns || columnHasTargetTag(columnData, targetTags))) { Column column = parseColumn(columnData, entityFQN); columnsByName .computeIfAbsent(colName, k -> new ArrayList<>()) @@ -376,7 +463,6 @@ private Query buildTagFilterQuery(ColumnAggregationRequest request) { addSchemaFilter(boolBuilder, request); addDomainFilter(boolBuilder, request); addColumnNamePatternFilter(boolBuilder, request); - addMetadataStatusFilter(boolBuilder, request); List allTags = new ArrayList<>(); if (!nullOrEmpty(request.getTags())) { @@ -427,7 +513,6 @@ private Query buildFilters( addDomainFilter(boolBuilder, request); addColumnNamePatternFilter(boolBuilder, request); addTagFilters(boolBuilder, request, columnNamesFromTagFilter); - addMetadataStatusFilter(boolBuilder, request); return Query.of(q -> q.bool(boolBuilder.build())); } @@ -550,61 +635,6 @@ private void addTagFilters( } } - private void addMetadataStatusFilter( - BoolQuery.Builder boolBuilder, ColumnAggregationRequest request) { - if (!nullOrEmpty(request.getMetadataStatus())) { - Query metadataStatusQuery = buildMetadataStatusFilter(request.getMetadataStatus()); - if (metadataStatusQuery != null) { - boolBuilder.filter(metadataStatusQuery); - } - } - } - - private Query buildMetadataStatusFilter(String status) { - String descField = "columns.description"; - String tagsField = "columns.tags"; - - Query hasDesc = hasNonEmptyField(descField); - Query hasTags = existsQuery(tagsField); - Query noDesc = hasEmptyOrMissingField(descField); - Query noTags = notExistsQuery(tagsField); - - return switch (status.toUpperCase()) { - case "MISSING" -> Query.of(q -> q.bool(b -> b.must(noDesc).must(noTags))); - case "INCOMPLETE" -> Query.of( - q -> - q.bool( - b -> - b.should(Query.of(qs -> qs.bool(bs -> bs.must(hasDesc).must(noTags)))) - .should(Query.of(qs -> qs.bool(bs -> bs.must(noDesc).must(hasTags)))) - .minimumShouldMatch("1"))); - case "COMPLETE" -> Query.of(q -> q.bool(b -> b.must(hasDesc).must(hasTags))); - default -> null; - }; - } - - private Query existsQuery(String field) { - return Query.of(q -> q.exists(e -> e.field(field))); - } - - private Query notExistsQuery(String field) { - return Query.of(q -> q.bool(b -> b.mustNot(existsQuery(field)))); - } - - // `wildcard(field, "?*")` matches any doc whose indexed terms include at least one token of - // at least one character — the analyzer-friendly equivalent of "field has non-empty value". - // We can't use `term(field, "")` against analyzed text fields like `columns.description`: the - // field's analyzer produces no tokens for the empty string and OS rejects the term query with - // `search_phase_execution_exception ... all shards failed`. Caught by - // ColumnGridResourceIT#test_getColumnGrid_withMetadataStatusIncomplete. - private Query hasNonEmptyField(String field) { - return Query.of(q -> q.wildcard(w -> w.field(field).value("?*"))); - } - - private Query hasEmptyOrMissingField(String field) { - return Query.of(q -> q.bool(b -> b.mustNot(hasNonEmptyField(field)))); - } - /** Phase 1: Get all matching column names using terms agg with include regex (no top_hits). */ private ColumnAggregator.NamesWithCount executeNamesQuery(Query query, String regex) throws IOException { diff --git a/openmetadata-service/src/test/java/org/openmetadata/service/search/ColumnAggregatorTest.java b/openmetadata-service/src/test/java/org/openmetadata/service/search/ColumnAggregatorTest.java index b0f8259f010e..13c79adc3239 100644 --- a/openmetadata-service/src/test/java/org/openmetadata/service/search/ColumnAggregatorTest.java +++ b/openmetadata-service/src/test/java/org/openmetadata/service/search/ColumnAggregatorTest.java @@ -15,10 +15,20 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; +import java.util.ArrayList; +import java.util.List; import java.util.regex.Pattern; import org.junit.jupiter.api.Test; +import org.openmetadata.schema.api.data.ColumnGridItem; +import org.openmetadata.schema.api.data.ColumnGridResponse; +import org.openmetadata.schema.api.data.ColumnMetadataGroup; +import org.openmetadata.schema.api.data.MetadataStatus; +import org.openmetadata.schema.type.TagLabel; +import org.openmetadata.service.search.ColumnAggregator.ColumnAggregationRequest; class ColumnAggregatorTest { @@ -109,4 +119,128 @@ void toCaseInsensitiveRegex_specialCharsAreEscaped() { // Plus and star should be literal, not regex quantifiers assertFalse(pattern.matcher("abbbbc").matches()); } + + // ---- Row-level (post-grouping) filter + pagination helpers (issue #26824) ---- + + private static ColumnGridItem item(String name, MetadataStatus status, boolean hasVariations) { + ColumnGridItem gi = new ColumnGridItem(); + gi.setColumnName(name); + gi.setMetadataStatus(status); + gi.setHasVariations(hasVariations); + gi.setTotalOccurrences(1); + ColumnMetadataGroup group = new ColumnMetadataGroup(); + if (status == MetadataStatus.COMPLETE || status == MetadataStatus.INCONSISTENT) { + group.setDescription("desc"); + group.setTags(List.of(new TagLabel().withTagFQN("PII.Sensitive"))); + } else if (status == MetadataStatus.INCOMPLETE) { + group.setDescription("desc"); + } + gi.setGroups(new ArrayList<>(List.of(group))); + return gi; + } + + private static ColumnAggregationRequest request(String metadataStatus) { + ColumnAggregationRequest r = new ColumnAggregationRequest(); + r.setMetadataStatus(metadataStatus); + return r; + } + + @Test + void matchesRowFilters_metadataStatusMatchesAggregateStatus() { + ColumnGridItem complete = item("a", MetadataStatus.COMPLETE, false); + ColumnGridItem incomplete = item("b", MetadataStatus.INCOMPLETE, false); + + assertTrue(ColumnAggregator.matchesRowFilters(complete, request("COMPLETE"))); + assertFalse(ColumnAggregator.matchesRowFilters(incomplete, request("COMPLETE"))); + assertTrue(ColumnAggregator.matchesRowFilters(incomplete, request("INCOMPLETE"))); + // Case-insensitive + assertTrue(ColumnAggregator.matchesRowFilters(complete, request("complete"))); + } + + @Test + void matchesRowFilters_inconsistentIsAFirstClassStatus() { + ColumnGridItem inconsistent = item("a", MetadataStatus.INCONSISTENT, true); + ColumnGridItem complete = item("b", MetadataStatus.COMPLETE, false); + + // The reported bug: "COMPLETE" must NOT return inconsistent rows. + assertFalse(ColumnAggregator.matchesRowFilters(inconsistent, request("COMPLETE"))); + assertTrue(ColumnAggregator.matchesRowFilters(inconsistent, request("INCONSISTENT"))); + assertFalse(ColumnAggregator.matchesRowFilters(complete, request("INCONSISTENT"))); + } + + @Test + void matchesRowFilters_nullOrBlankStatusMatchesEverything() { + ColumnGridItem complete = item("a", MetadataStatus.COMPLETE, false); + assertTrue(ColumnAggregator.matchesRowFilters(complete, request(null))); + assertTrue(ColumnAggregator.matchesRowFilters(complete, request(" "))); + } + + @Test + void matchesRowFilters_hasConflictsRequiresVariations() { + ColumnAggregationRequest r = new ColumnAggregationRequest(); + r.setHasConflicts(true); + assertTrue(ColumnAggregator.matchesRowFilters(item("a", MetadataStatus.INCONSISTENT, true), r)); + assertFalse(ColumnAggregator.matchesRowFilters(item("b", MetadataStatus.COMPLETE, false), r)); + } + + @Test + void paginateFilteredItems_filtersAndComputesTotalsFromFilteredSet() { + List all = + new ArrayList<>( + List.of( + item("complete_1", MetadataStatus.COMPLETE, false), + item("inconsistent_1", MetadataStatus.INCONSISTENT, true), + item("complete_2", MetadataStatus.COMPLETE, false), + item("missing_1", MetadataStatus.MISSING, false))); + + ColumnAggregationRequest r = request("COMPLETE"); + r.setSize(10); + + ColumnGridResponse resp = ColumnAggregator.paginateFilteredItems(all, r); + + // Only the two COMPLETE rows survive, and totals reflect the FILTERED set (not 4). + assertEquals(2, resp.getTotalUniqueColumns()); + assertEquals(2, resp.getColumns().size()); + assertTrue( + resp.getColumns().stream().allMatch(c -> c.getMetadataStatus() == MetadataStatus.COMPLETE)); + assertNull(resp.getCursor()); + } + + @Test + void paginateFilteredItems_paginatesConsistentlyAcrossPages() { + List all = new ArrayList<>(); + for (int i = 0; i < 5; i++) { + all.add(item(String.format("col_%02d", i), MetadataStatus.COMPLETE, false)); + } + + ColumnAggregationRequest page1 = request("COMPLETE"); + page1.setSize(2); + ColumnGridResponse r1 = ColumnAggregator.paginateFilteredItems(all, page1); + + assertEquals(5, r1.getTotalUniqueColumns()); + assertEquals(2, r1.getColumns().size()); + assertEquals("col_00", r1.getColumns().get(0).getColumnName()); + assertEquals("col_01", r1.getColumns().get(1).getColumnName()); + assertNotNull(r1.getCursor(), "more pages remain, so a cursor is returned"); + + ColumnAggregationRequest page2 = request("COMPLETE"); + page2.setSize(2); + page2.setCursor(r1.getCursor()); + ColumnGridResponse r2 = ColumnAggregator.paginateFilteredItems(all, page2); + + assertEquals(5, r2.getTotalUniqueColumns()); + assertEquals(2, r2.getColumns().size()); + assertEquals("col_02", r2.getColumns().get(0).getColumnName()); + assertEquals("col_03", r2.getColumns().get(1).getColumnName()); + assertNotNull(r2.getCursor()); + + ColumnAggregationRequest page3 = request("COMPLETE"); + page3.setSize(2); + page3.setCursor(r2.getCursor()); + ColumnGridResponse r3 = ColumnAggregator.paginateFilteredItems(all, page3); + + assertEquals(1, r3.getColumns().size(), "last page has the remaining single item"); + assertEquals("col_04", r3.getColumns().get(0).getColumnName()); + assertNull(r3.getCursor(), "cursor is null on the last page"); + } }