Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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"));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<String> 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<String> THREAD_TYPE_RESOURCES = Set.of("task", "conversation");

private static boolean shouldTriggerAlertForThread(ChangeEvent event, String resource) {
Thread thread = AlertsRuleEvaluator.getThread(event);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -74,9 +76,8 @@ public boolean matchAnySource(List<String> 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();
Expand All @@ -100,9 +101,8 @@ public boolean matchAnyOwnerName(List<String> 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);
Expand Down Expand Up @@ -146,9 +146,8 @@ public boolean matchAnyEntityFqn(List<String> 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);
Expand Down Expand Up @@ -193,18 +192,11 @@ public boolean matchAnyEntityId(List<String> 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(
Expand Down Expand Up @@ -480,23 +472,16 @@ public boolean matchAnyDomain(List<String> 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)) {
Expand Down Expand Up @@ -690,6 +675,83 @@ private boolean testSuiteOwnerMatcher(List<TestSuite> testSuites, List<String> 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<String> entityTypes) {
EntityReference subject = threadSubject();
return subject != null && entityTypes.contains(subject.getType());
}

private boolean threadSubjectMatchesFqn(List<String> entityFqns) {
Comment thread
gitar-bot[bot] marked this conversation as resolved.
EntityReference subject = threadSubject();
return subject != null && matchesFqnOrDescendant(subject.getFullyQualifiedName(), entityFqns);
}

private boolean threadSubjectMatchesId(List<String> 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<String> entityIds) {
boolean matched = false;
for (String id : entityIds) {
if (entityId.equals(parseUuidOrNull(id))) {
matched = true;
break;
}
}
return matched;
}
Comment thread
gitar-bot[bot] marked this conversation as resolved.

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<String> 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<String> domainFqns) {
EntityInterface subject =
Entity.getEntityOrNull(threadSubject(), Entity.FIELD_DOMAINS, Include.NON_DELETED);
return subject != null && matchDomains(subject.getDomains(), domainFqns);
}

private boolean matchDomains(List<EntityReference> domains, List<String> 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<EntityReference> ownerReferences, List<String> ownerNameList) {
Set<String> ownerNames =
ownerNameList.stream().map(EntityInterfaceUtil::unquoteName).collect(Collectors.toSet());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -161,14 +161,25 @@ void shouldTriggerAlert_taskThread_taskResource_returnsTrue() {
}

@Test
void shouldTriggerAlert_announcementThread_announcementResource_returnsTrue() {
void shouldTriggerAlert_announcementThread_announcementResource_returnsFalse() {
Thread thread =
new Thread()
.withId(UUID.randomUUID())
.withType(ThreadType.Announcement)
.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));
}

Expand Down
Loading
Loading