diff --git a/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/main/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistry.java b/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/main/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistry.java index a5955d32127c..048def65e4d1 100644 --- a/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/main/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistry.java +++ b/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/main/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistry.java @@ -438,7 +438,9 @@ private Event toEvent(final WatchEvent watchEvent, final String watchedPath) { .watchedPath(watchedPath) .eventPath(Optional.ofNullable(keyValue).map(kv -> kv.getKey().toString(StandardCharsets.UTF_8)) .orElse(null)) - .eventData(Optional.ofNullable(keyValue).map(kv -> kv.getValue().toString(StandardCharsets.UTF_8)) + .eventData(Optional + .ofNullable(eventType == Event.Type.REMOVE ? watchEvent.getPrevKV() : watchEvent.getKeyValue()) + .map(kv -> kv.getValue().toString(StandardCharsets.UTF_8)) .orElse(null)) .build(); } diff --git a/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/test/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistryEventTest.java b/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/test/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistryEventTest.java new file mode 100644 index 000000000000..c69bb4061090 --- /dev/null +++ b/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-etcd/src/test/java/org/apache/dolphinscheduler/plugin/registry/etcd/EtcdRegistryEventTest.java @@ -0,0 +1,87 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.dolphinscheduler.plugin.registry.etcd; + +import org.apache.dolphinscheduler.registry.api.Event; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.mockito.Mockito; +import org.springframework.test.util.ReflectionTestUtils; + +import com.google.protobuf.ByteString; + +import io.etcd.jetcd.ByteSequence; +import io.etcd.jetcd.KeyValue; +import io.etcd.jetcd.watch.WatchEvent; + +class EtcdRegistryEventTest { + + private static final String WATCHED_PATH = "/nodes"; + private static final String EVENT_PATH = "/nodes/master-1"; + + private final EtcdRegistry registry = Mockito.mock(EtcdRegistry.class); + + @Test + void testDeleteUsesPreviousValue() { + // An etcd DELETE contains only the key and modification revision in its current KV. + KeyValue deletedKeyValue = new KeyValue(io.etcd.jetcd.api.KeyValue.newBuilder() + .setKey(ByteString.copyFromUtf8(EVENT_PATH)) + .setModRevision(3) + .build(), ByteSequence.EMPTY); + assertEvent(new WatchEvent(deletedKeyValue, keyValue(EVENT_PATH, "previous-heartbeat"), + WatchEvent.EventType.DELETE), Event.Type.REMOVE, "previous-heartbeat"); + } + + @Test + void testAddUsesCurrentValue() { + assertEvent(new WatchEvent(keyValue(EVENT_PATH, "current-heartbeat"), + new KeyValue(io.etcd.jetcd.api.KeyValue.getDefaultInstance(), ByteSequence.EMPTY), + WatchEvent.EventType.PUT), Event.Type.ADD, "current-heartbeat"); + } + + @Test + void testUpdateUsesCurrentValue() { + assertEvent(new WatchEvent(keyValue(EVENT_PATH, "current-heartbeat"), + keyValue(EVENT_PATH, "previous-heartbeat"), WatchEvent.EventType.PUT), + Event.Type.UPDATE, "current-heartbeat"); + } + + @Test + void testDeleteWithoutPreviousValuePreservesPath() { + assertEvent(new WatchEvent(keyValue(EVENT_PATH, ""), + new KeyValue(io.etcd.jetcd.api.KeyValue.getDefaultInstance(), ByteSequence.EMPTY), + WatchEvent.EventType.DELETE), Event.Type.REMOVE, ""); + } + + private void assertEvent(WatchEvent watchEvent, Event.Type expectedType, String expectedData) { + Event event = ReflectionTestUtils.invokeMethod(registry, "toEvent", watchEvent, WATCHED_PATH); + Assertions.assertNotNull(event); + Assertions.assertEquals(expectedType, event.getType()); + Assertions.assertEquals(WATCHED_PATH, event.getWatchedPath()); + Assertions.assertEquals(EVENT_PATH, event.getEventPath()); + Assertions.assertEquals(expectedData, event.getEventData()); + } + + private KeyValue keyValue(String key, String value) { + return new KeyValue(io.etcd.jetcd.api.KeyValue.newBuilder() + .setKey(ByteString.copyFromUtf8(key)) + .setValue(ByteString.copyFromUtf8(value)) + .build(), ByteSequence.EMPTY); + } +} diff --git a/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-it/src/test/java/org/apache/dolphinscheduler/plugin/registry/RegistryTestCase.java b/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-it/src/test/java/org/apache/dolphinscheduler/plugin/registry/RegistryTestCase.java index b819ef0eee79..8e192c25bf52 100644 --- a/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-it/src/test/java/org/apache/dolphinscheduler/plugin/registry/RegistryTestCase.java +++ b/dolphinscheduler-registry/dolphinscheduler-registry-plugins/dolphinscheduler-registry-it/src/test/java/org/apache/dolphinscheduler/plugin/registry/RegistryTestCase.java @@ -119,6 +119,58 @@ public SubscribeScope getSubscribeScope() { }); } + @SneakyThrows + @Test + public void testSubscribeEventData() { + registry.start(); + + // Futures safely publish each first event and its payload to the test thread. + final CompletableFuture subscribeAdded = new CompletableFuture<>(); + final CompletableFuture subscribeRemoved = new CompletableFuture<>(); + final CompletableFuture subscribeUpdated = new CompletableFuture<>(); + + final SubscribeListener subscribeListener = new SubscribeListener() { + + @Override + public void notify(Event event) { + // Keep assertions on the test thread so callback error handling cannot hide failures. + if (event.getType() == Event.Type.ADD) { + subscribeAdded.complete(event); + } + if (event.getType() == Event.Type.REMOVE) { + subscribeRemoved.complete(event); + } + if (event.getType() == Event.Type.UPDATE) { + subscribeUpdated.complete(event); + } + } + + @Override + public SubscribeScope getSubscribeScope() { + return SubscribeScope.PATH_ONLY; + } + }; + String key = "/nodes/master" + System.nanoTime(); + registry.subscribe(key, subscribeListener); + // Wait after each change so polling registries cannot collapse consecutive operations. + registry.put(key, "v1", true); + assertSubscribeEvent(subscribeAdded.get(10, TimeUnit.SECONDS), Event.Type.ADD, key, "v1"); + + registry.put(key, "v2", true); + assertSubscribeEvent(subscribeUpdated.get(10, TimeUnit.SECONDS), Event.Type.UPDATE, key, "v2"); + + // REMOVE must retain the last value before deletion, not the initial value. + registry.delete(key); + assertSubscribeEvent(subscribeRemoved.get(10, TimeUnit.SECONDS), Event.Type.REMOVE, key, "v2"); + } + + private void assertSubscribeEvent(Event event, Event.Type expectedType, String key, String expectedData) { + Assertions.assertEquals(expectedType, event.getType()); + Assertions.assertEquals(key, event.getWatchedPath()); + Assertions.assertEquals(key, event.getEventPath()); + Assertions.assertEquals(expectedData, event.getEventData()); + } + @SneakyThrows @Test public void testAddConnectionStateListener() {