Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -130,14 +130,21 @@ private void proxyHandlerEvent(final String updatePath, final String updateData,
}

private void discoveryUpstreamHandlerEvent(final String updatePath, final String updateData, final EventType eventType) {
String[] pathInfoArray2 = updatePath.split("/");
if (pathInfoArray2.length != 5) {
String[] pathInfoArray = updatePath.split("/");
if (pathInfoArray.length != 5) {
return;
}
if (!EventType.DELETE.equals(eventType)) {
Optional.ofNullable(updateData)
.ifPresent(e -> cacheDiscoveryUpstreamData(GsonUtils.getInstance().fromJson(updateData, DiscoverySyncData.class)));
String pluginName = pathInfoArray[pathInfoArray.length - 2];
String selectorId = pathInfoArray[pathInfoArray.length - 1];
if (EventType.DELETE.equals(eventType)) {
DiscoverySyncData discoverySyncData = new DiscoverySyncData();
discoverySyncData.setPluginName(pluginName);
discoverySyncData.setSelectorId(selectorId);
unCacheDiscoveryUpstreamData(discoverySyncData);
return;
}
Optional.ofNullable(updateData)
.ifPresent(e -> cacheDiscoveryUpstreamData(GsonUtils.getInstance().fromJson(updateData, DiscoverySyncData.class)));
}

private void ruleHandlerEvent(final String updatePath, final String updateData, final EventType eventType) {
Expand Down Expand Up @@ -280,6 +287,11 @@ protected void cacheDiscoveryUpstreamData(final DiscoverySyncData upstreamDataLi
.ifPresent(data -> discoveryUpstreamDataSubscribers.forEach(e -> e.onSubscribe(upstreamDataList)));
}

protected void unCacheDiscoveryUpstreamData(final DiscoverySyncData discoverySyncData) {
Optional.ofNullable(discoverySyncData)
.ifPresent(data -> discoveryUpstreamDataSubscribers.forEach(e -> e.unSubscribe(data)));
}

protected void unCacheMetaData(final MetaData metaData) {
Optional.ofNullable(metaData)
.ifPresent(data -> metaDataSubscribers.forEach(e -> e.unSubscribe(metaData)));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,9 @@

package org.apache.shenyu.sync.data.core;

import org.apache.shenyu.common.constant.DefaultPathConstants;
import org.apache.shenyu.common.dto.AppAuthData;
import org.apache.shenyu.common.dto.DiscoverySyncData;
import org.apache.shenyu.common.utils.GsonUtils;
import org.apache.shenyu.sync.data.api.AuthDataSubscriber;
import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
Expand All @@ -27,13 +29,15 @@
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.MockitoAnnotations;
import org.mockito.junit.MockitoJUnitRunner;

import java.util.ArrayList;
import java.util.List;

import static org.junit.Assert.assertEquals;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.verify;

Expand Down Expand Up @@ -85,6 +89,24 @@ public void testUnCacheAuthData() {
verify(authDataSubscriber).unSubscribe(any());
}

@Test
public void testDiscoveryUpstreamHandlerEvent() {

String namespaceId = "/namespace";
String registerPath = namespaceId + DefaultPathConstants.DISCOVERY_UPSTREAM;
String updatePath = registerPath + "/divide/testSelectorId";
String jsonData = "{\"pluginName\":\"divide\",\"selectorId\":\"testSelectorId\",\"selectorName\":\"testSelector\"}";

pathDataSyncService.event(namespaceId, updatePath, jsonData, registerPath, AbstractPathDataSyncService.EventType.PUT);
verify(discoveryUpstreamDataSubscriber).onSubscribe(any());

pathDataSyncService.event(namespaceId, updatePath, null, registerPath, AbstractPathDataSyncService.EventType.DELETE);
ArgumentCaptor<DiscoverySyncData> captor = ArgumentCaptor.forClass(DiscoverySyncData.class);
verify(discoveryUpstreamDataSubscriber).unSubscribe(captor.capture());
assertEquals("divide", captor.getValue().getPluginName());
assertEquals("testSelectorId", captor.getValue().getSelectorId());
}

// Mock implementation
static class AbstractPathDataSyncServiceImpl extends AbstractPathDataSyncService {

Expand Down
Loading