diff --git a/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/AlertsRuleEvaluatorResourceIT.java b/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/AlertsRuleEvaluatorResourceIT.java index 1ec284afef57..ab90ea704c3b 100644 --- a/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/AlertsRuleEvaluatorResourceIT.java +++ b/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/AlertsRuleEvaluatorResourceIT.java @@ -23,6 +23,7 @@ import org.openmetadata.schema.entity.data.DataContract; import org.openmetadata.schema.entity.data.DatabaseSchema; import org.openmetadata.schema.entity.data.Table; +import org.openmetadata.schema.entity.feed.Thread; import org.openmetadata.schema.entity.services.DatabaseService; import org.openmetadata.schema.entity.teams.User; import org.openmetadata.schema.tests.TestCase; @@ -37,6 +38,7 @@ import org.openmetadata.schema.type.EntityReference; import org.openmetadata.schema.type.EventType; import org.openmetadata.schema.type.FieldChange; +import org.openmetadata.schema.type.ThreadType; import org.openmetadata.schema.utils.JsonUtils; import org.openmetadata.service.Entity; import org.openmetadata.service.events.subscription.AlertsRuleEvaluator; @@ -81,6 +83,55 @@ void test_matchAnyOwnerName(TestNamespace ns) { assertFalse(evaluateExpression("matchAnyOwnerName('nonExistentOwner')", evaluationContext)); } + // #30555: a thread event carries the Thread, so owner and domain must resolve its parent entity. + @Test + void test_matchAnyOwnerName_threadAboutOwnedEntity(TestNamespace ns) { + Table createdTable = createTableWithOwner(ns); + EvaluationContext evaluationContext = + threadEvaluationContext(createdTable.getEntityReference()); + + String ownerName = SharedEntities.get().USER1.getName(); + assertTrue(evaluateExpression("matchAnyOwnerName('" + ownerName + "')", evaluationContext)); + assertFalse(evaluateExpression("matchAnyOwnerName('nonExistentOwner')", evaluationContext)); + } + + @Test + void test_matchAnyDomain_threadAboutEntityInDomain(TestNamespace ns) { + Table createdTable = createTableInDomain(ns); + EvaluationContext evaluationContext = + threadEvaluationContext(createdTable.getEntityReference()); + + String domainFqn = SharedEntities.get().DOMAIN.getFullyQualifiedName(); + assertTrue(evaluateExpression("matchAnyDomain({'" + domainFqn + "'})", evaluationContext)); + assertFalse(evaluateExpression("matchAnyDomain({'nonExistentDomain'})", evaluationContext)); + } + + @Test + void test_matchAnyEntityFqn_threadAboutEntity(TestNamespace ns) { + Table createdTable = createTable(ns); + EvaluationContext evaluationContext = + threadEvaluationContext(createdTable.getEntityReference()); + + String tableFqn = createdTable.getFullyQualifiedName(); + assertTrue(evaluateExpression("matchAnyEntityFqn({'" + tableFqn + "'})", evaluationContext)); + assertFalse(evaluateExpression("matchAnyEntityFqn({'nonExistentFqn'})", evaluationContext)); + } + + @Test + void test_scopingFilters_threadAboutUnresolvableEntity() { + EntityReference missing = + new EntityReference() + .withId(UUID.randomUUID()) + .withType(Entity.TABLE) + .withFullyQualifiedName("service.db.schema.missingTable"); + EvaluationContext evaluationContext = threadEvaluationContext(missing); + + String ownerName = SharedEntities.get().USER1.getName(); + String domainFqn = SharedEntities.get().DOMAIN.getFullyQualifiedName(); + assertFalse(evaluateExpression("matchAnyOwnerName('" + ownerName + "')", evaluationContext)); + assertFalse(evaluateExpression("matchAnyDomain({'" + domainFqn + "'})", evaluationContext)); + } + @Test @Disabled("TestCase.getTestSuite().getFullyQualifiedName() returns null - needs investigation") void test_matchAnyEntityFqn(TestNamespace ns) { @@ -694,6 +745,35 @@ private Table createTableWithOwner(TestNamespace ns) { return SdkClients.adminClient().tables().create(tableRequest); } + private Table createTableInDomain(TestNamespace ns) { + Table table = createTable(ns); + CreateTable update = + new CreateTable() + .withName(table.getName()) + .withDatabaseSchema(table.getDatabaseSchema().getFullyQualifiedName()) + .withColumns(table.getColumns()) + .withDomains(List.of(SharedEntities.get().DOMAIN.getFullyQualifiedName())); + + return SdkClients.adminClient().tables().createOrUpdate(update); + } + + private EvaluationContext threadEvaluationContext(EntityReference parent) { + Thread thread = + new Thread() + .withId(UUID.randomUUID()) + .withType(ThreadType.Conversation) + .withEntityRef(parent); + ChangeEvent changeEvent = new ChangeEvent(); + changeEvent.setEntityType(Entity.THREAD); + changeEvent.setEventType(EventType.THREAD_CREATED); + changeEvent.setEntity(thread); + + return SimpleEvaluationContext.forReadOnlyDataBinding() + .withInstanceMethods() + .withRootObject(new AlertsRuleEvaluator(changeEvent)) + .build(); + } + private TestCase createTestCase(TestNamespace ns, Table table) { CreateTestCase createTestCase = new CreateTestCase(); createTestCase.setName(ns.prefix("testCase")); diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/events/subscription/AlertUtil.java b/openmetadata-service/src/main/java/org/openmetadata/service/events/subscription/AlertUtil.java index acc971b054b5..efeef8aa89ee 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/events/subscription/AlertUtil.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/events/subscription/AlertUtil.java @@ -177,8 +177,8 @@ public static boolean shouldTriggerAlert(ChangeEvent event, FilteringRules confi return config.getResources().contains(event.getEntityType()); // Use Trigger Specific Settings } - private static final Set THREAD_TYPE_RESOURCES = - Set.of("announcement", "task", "conversation"); + // Announcement is its own entity since #25894; task stays until #30559 retires the legacy path. + private static final Set THREAD_TYPE_RESOURCES = Set.of("task", "conversation"); private static boolean shouldTriggerAlertForThread(ChangeEvent event, String resource) { Thread thread = AlertsRuleEvaluator.getThread(event); diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/events/subscription/AlertsRuleEvaluator.java b/openmetadata-service/src/main/java/org/openmetadata/service/events/subscription/AlertsRuleEvaluator.java index f7c8b7f26caf..611f04026396 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/events/subscription/AlertsRuleEvaluator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/events/subscription/AlertsRuleEvaluator.java @@ -57,6 +57,8 @@ @Slf4j public class AlertsRuleEvaluator { private final ChangeEvent changeEvent; + private EntityReference threadSubject; + private boolean threadSubjectResolved; public AlertsRuleEvaluator(ChangeEvent event) { this.changeEvent = event; @@ -74,9 +76,8 @@ public boolean matchAnySource(List originEntities) { return false; } - // Filter does not apply to Thread Change Events if (changeEvent.getEntityType().equals(THREAD)) { - return true; + return threadSubjectMatchesType(originEntities); } String changeEventEntity = changeEvent.getEntityType(); @@ -100,9 +101,8 @@ public boolean matchAnyOwnerName(List ownerNameList) { return false; } - // Filter does not apply to Thread Change Events if (changeEvent.getEntityType().equals(THREAD)) { - return true; + return threadSubjectMatchesOwner(ownerNameList); } EntityInterface entity = getEntity(changeEvent); @@ -146,9 +146,8 @@ public boolean matchAnyEntityFqn(List entityFqns) { return false; } - // Filter does not apply to Thread Change Events if (changeEvent.getEntityType().equals(THREAD)) { - return true; + return threadSubjectMatchesFqn(entityFqns); } EntityInterface entity = getEntity(changeEvent); @@ -193,18 +192,11 @@ public boolean matchAnyEntityId(List entityIds) { return false; } - // Filter does not apply to Thread Change Events if (changeEvent.getEntityType().equals(THREAD)) { - return true; + return threadSubjectMatchesId(entityIds); } - EntityInterface entity = getEntity(changeEvent); - for (String id : entityIds) { - if (entity.getId().equals(UUID.fromString(id))) { - return true; - } - } - return false; + return matchesAnyId(getEntity(changeEvent).getId(), entityIds); } @Function( @@ -480,23 +472,16 @@ public boolean matchAnyDomain(List fieldChangeUpdate) { return false; } - // Filter does not apply to Thread Change Events if (changeEvent.getEntityType().equals(THREAD)) { - return true; + return threadSubjectMatchesDomain(fieldChangeUpdate); } EntityInterface entity = getEntity(changeEvent); EntityInterface entityWithDomainData = Entity.getEntity( changeEvent.getEntityType(), entity.getId(), "domains", Include.NON_DELETED); - if (!nullOrEmpty(entityWithDomainData.getDomains())) { - for (String name : fieldChangeUpdate) { - for (EntityReference domain : entityWithDomainData.getDomains()) { - if (domain.getFullyQualifiedName().equals(name)) { - return true; - } - } - } + if (matchDomains(entityWithDomainData.getDomains(), fieldChangeUpdate)) { + return true; } if (changeEvent.getEntityType().equals(TEST_CASE)) { @@ -690,6 +675,83 @@ private boolean testSuiteOwnerMatcher(List testSuites, List o return false; } + // Scoping filters on a thread event are about the thread's parent entity, not the thread. + // Memoized: every matcher in one evaluation asks for the same subject. + private EntityReference threadSubject() { + if (!threadSubjectResolved) { + threadSubjectResolved = true; + if (changeEvent.getEntity() != null) { + Thread thread = getThread(changeEvent); + threadSubject = thread == null ? null : thread.getEntityRef(); + } + } + return threadSubject; + } + + private boolean threadSubjectMatchesType(List entityTypes) { + EntityReference subject = threadSubject(); + return subject != null && entityTypes.contains(subject.getType()); + } + + private boolean threadSubjectMatchesFqn(List entityFqns) { + EntityReference subject = threadSubject(); + return subject != null && matchesFqnOrDescendant(subject.getFullyQualifiedName(), entityFqns); + } + + private boolean threadSubjectMatchesId(List entityIds) { + EntityReference subject = threadSubject(); + return subject != null && matchesAnyId(subject.getId(), entityIds); + } + + // A filter value that is not a UUID cannot match; throwing here would abort the whole batch. + private static boolean matchesAnyId(UUID entityId, List entityIds) { + boolean matched = false; + for (String id : entityIds) { + if (entityId.equals(parseUuidOrNull(id))) { + matched = true; + break; + } + } + return matched; + } + + private static UUID parseUuidOrNull(String id) { + UUID parsed = null; + try { + parsed = UUID.fromString(id); + } catch (IllegalArgumentException e) { + LOG.debug("Ignoring non-UUID entity id filter value '{}'", id); + } + return parsed; + } + + private boolean threadSubjectMatchesOwner(List ownerNameList) { + EntityInterface subject = + Entity.getEntityOrNull(threadSubject(), Entity.FIELD_OWNERS, Include.NON_DELETED); + return subject != null + && !nullOrEmpty(subject.getOwners()) + && matchOwners(subject.getOwners(), ownerNameList); + } + + private boolean threadSubjectMatchesDomain(List domainFqns) { + EntityInterface subject = + Entity.getEntityOrNull(threadSubject(), Entity.FIELD_DOMAINS, Include.NON_DELETED); + return subject != null && matchDomains(subject.getDomains(), domainFqns); + } + + private boolean matchDomains(List domains, List domainFqns) { + boolean matched = false; + if (!nullOrEmpty(domains)) { + for (EntityReference domain : domains) { + if (domainFqns.contains(domain.getFullyQualifiedName())) { + matched = true; + break; + } + } + } + return matched; + } + private boolean matchOwners(List ownerReferences, List ownerNameList) { Set ownerNames = ownerNameList.stream().map(EntityInterfaceUtil::unquoteName).collect(Collectors.toSet()); diff --git a/openmetadata-service/src/test/java/org/openmetadata/service/events/subscription/AlertUtilTest.java b/openmetadata-service/src/test/java/org/openmetadata/service/events/subscription/AlertUtilTest.java index 109b8f79756c..2c690ab0bc25 100644 --- a/openmetadata-service/src/test/java/org/openmetadata/service/events/subscription/AlertUtilTest.java +++ b/openmetadata-service/src/test/java/org/openmetadata/service/events/subscription/AlertUtilTest.java @@ -161,7 +161,7 @@ void shouldTriggerAlert_taskThread_taskResource_returnsTrue() { } @Test - void shouldTriggerAlert_announcementThread_announcementResource_returnsTrue() { + void shouldTriggerAlert_announcementThread_announcementResource_returnsFalse() { Thread thread = new Thread() .withId(UUID.randomUUID()) @@ -169,6 +169,17 @@ void shouldTriggerAlert_announcementThread_announcementResource_returnsTrue() { .withEntityRef(new EntityReference().withId(UUID.randomUUID()).withType("table")); ChangeEvent event = threadChangeEvent(thread, EventType.THREAD_CREATED); FilteringRules config = filteringRules("announcement"); + assertFalse(AlertUtil.shouldTriggerAlert(event, config)); + } + + @Test + void shouldTriggerAlert_announcementEntityEvent_announcementResource_returnsTrue() { + ChangeEvent event = + new ChangeEvent() + .withId(UUID.randomUUID()) + .withEventType(EventType.ENTITY_CREATED) + .withEntityType("announcement"); + FilteringRules config = filteringRules("announcement"); assertTrue(AlertUtil.shouldTriggerAlert(event, config)); } diff --git a/openmetadata-service/src/test/java/org/openmetadata/service/events/subscription/AlertsRuleEvaluatorThreadScopeTest.java b/openmetadata-service/src/test/java/org/openmetadata/service/events/subscription/AlertsRuleEvaluatorThreadScopeTest.java new file mode 100644 index 000000000000..ba6c0cc347d7 --- /dev/null +++ b/openmetadata-service/src/test/java/org/openmetadata/service/events/subscription/AlertsRuleEvaluatorThreadScopeTest.java @@ -0,0 +1,109 @@ +package org.openmetadata.service.events.subscription; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.List; +import java.util.UUID; +import org.junit.jupiter.api.Test; +import org.openmetadata.schema.entity.feed.Thread; +import org.openmetadata.schema.type.ChangeEvent; +import org.openmetadata.schema.type.EntityReference; +import org.openmetadata.schema.type.EventType; +import org.openmetadata.schema.type.ThreadType; +import org.openmetadata.service.Entity; + +// #30555: scoping filters must match a thread against its parent entity, never pass it through. +class AlertsRuleEvaluatorThreadScopeTest { + + private static final String TERM_FQN = "glossary.term"; + private static final String OTHER_TERM_FQN = "glossary.otherTerm"; + private static final UUID TERM_ID = UUID.randomUUID(); + + @Test + void matchAnyEntityFqn_threadAboutListedEntity_returnsTrue() { + AlertsRuleEvaluator evaluator = evaluatorForThreadAbout(termRef(TERM_FQN)); + assertTrue(evaluator.matchAnyEntityFqn(List.of(TERM_FQN))); + } + + @Test + void matchAnyEntityFqn_threadAboutOtherEntity_returnsFalse() { + AlertsRuleEvaluator evaluator = evaluatorForThreadAbout(termRef(TERM_FQN)); + assertFalse(evaluator.matchAnyEntityFqn(List.of(OTHER_TERM_FQN))); + } + + @Test + void matchAnyEntityFqn_threadAboutDescendantOfListedEntity_returnsTrue() { + AlertsRuleEvaluator evaluator = evaluatorForThreadAbout(termRef("glossary.term.child")); + assertTrue(evaluator.matchAnyEntityFqn(List.of("glossary.term"))); + } + + @Test + void matchAnyEntityFqn_threadAboutSiblingWithSharedPrefix_returnsFalse() { + AlertsRuleEvaluator evaluator = evaluatorForThreadAbout(termRef("glossary.termOther")); + assertFalse(evaluator.matchAnyEntityFqn(List.of("glossary.term"))); + } + + @Test + void matchAnyEntityId_threadAboutListedEntity_returnsTrue() { + AlertsRuleEvaluator evaluator = evaluatorForThreadAbout(termRef(TERM_FQN)); + assertTrue(evaluator.matchAnyEntityId(List.of(TERM_ID.toString()))); + } + + @Test + void matchAnyEntityId_threadAboutOtherEntity_returnsFalse() { + AlertsRuleEvaluator evaluator = evaluatorForThreadAbout(termRef(TERM_FQN)); + assertFalse(evaluator.matchAnyEntityId(List.of(UUID.randomUUID().toString()))); + } + + @Test + void matchAnyEntityId_malformedFilterValue_isSkippedInsteadOfThrowing() { + AlertsRuleEvaluator evaluator = evaluatorForThreadAbout(termRef(TERM_FQN)); + assertFalse(evaluator.matchAnyEntityId(List.of("not-a-uuid"))); + assertTrue(evaluator.matchAnyEntityId(List.of("not-a-uuid", TERM_ID.toString()))); + } + + @Test + void matchAnySource_threadAboutListedEntityType_returnsTrue() { + AlertsRuleEvaluator evaluator = evaluatorForThreadAbout(termRef(TERM_FQN)); + assertTrue(evaluator.matchAnySource(List.of("glossaryTerm"))); + } + + @Test + void matchAnySource_threadAboutOtherEntityType_returnsFalse() { + AlertsRuleEvaluator evaluator = evaluatorForThreadAbout(termRef(TERM_FQN)); + assertFalse(evaluator.matchAnySource(List.of("table"))); + } + + @Test + void scopingFilters_threadWithoutEntityRef_returnFalse() { + AlertsRuleEvaluator evaluator = evaluatorForThreadAbout(null); + assertFalse(evaluator.matchAnyEntityFqn(List.of(TERM_FQN))); + assertFalse(evaluator.matchAnyEntityId(List.of(TERM_ID.toString()))); + assertFalse(evaluator.matchAnySource(List.of("glossaryTerm"))); + assertFalse(evaluator.matchAnyOwnerName(List.of("admin"))); + assertFalse(evaluator.matchAnyDomain(List.of("domain"))); + } + + private static AlertsRuleEvaluator evaluatorForThreadAbout(EntityReference parent) { + Thread thread = + new Thread() + .withId(UUID.randomUUID()) + .withType(ThreadType.Conversation) + .withEntityRef(parent); + ChangeEvent event = + new ChangeEvent() + .withId(UUID.randomUUID()) + .withEventType(EventType.THREAD_CREATED) + .withEntityType(Entity.THREAD) + .withEntity(thread); + return new AlertsRuleEvaluator(event); + } + + private static EntityReference termRef(String fqn) { + return new EntityReference() + .withId(TERM_ID) + .withType("glossaryTerm") + .withFullyQualifiedName(fqn); + } +}