This is an automated email from the ASF dual-hosted git repository.

ruanwenjun pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/dolphinscheduler.git


The following commit(s) were added to refs/heads/dev by this push:
     new f16479955a [Fix-18640][Registry] Preserve previous values in etcd 
REMOVE events (#18641)
f16479955a is described below

commit f16479955a28d0d0d51e5ccea5fb975da8db8417
Author: michaellx1057 <[email protected]>
AuthorDate: Wed Sep 23 22:14:35 2026 +0800

    [Fix-18640][Registry] Preserve previous values in etcd REMOVE events 
(#18641)
---
 .../plugin/registry/etcd/EtcdRegistry.java         |  4 +-
 .../registry/etcd/EtcdRegistryEventTest.java       | 87 ++++++++++++++++++++++
 .../plugin/registry/RegistryTestCase.java          | 52 +++++++++++++
 3 files changed, 142 insertions(+), 1 deletion(-)

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 a5955d3212..048def65e4 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 @@ public class EtcdRegistry implements Registry {
                 .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 0000000000..c69bb40610
--- /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 b819ef0eee..8e192c25bf 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 abstract class RegistryTestCase<R extends Registry> 
{
                 });
     }
 
+    @SneakyThrows
+    @Test
+    public void testSubscribeEventData() {
+        registry.start();
+
+        // Futures safely publish each first event and its payload to the test 
thread.
+        final CompletableFuture<Event> subscribeAdded = new 
CompletableFuture<>();
+        final CompletableFuture<Event> subscribeRemoved = new 
CompletableFuture<>();
+        final CompletableFuture<Event> 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() {

Reply via email to