Skip to content
Merged
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
10 changes: 10 additions & 0 deletions .github/scripts/resolve-ci-modules.sh
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,16 @@ while IFS= read -r file; do
fi
done < <(printf '%s' "${changed_files_json}" | jq -r '.[]')

# SPI changes can affect modules that consume shared classes without declaring a
# direct Maven dependency on every transitive module. Build the full reactor so
# those modules cannot silently use a stale SNAPSHOT from the Maven cache.
for module in "${modules[@]}"; do
if [[ "${module}" == "shenyu-spi" ]]; then
full_build_required=true
break
fi
done

if [[ "${has_code_changes}" == "true" && ("${#modules[@]}" -eq 0 || "${#modules[@]}" -gt "${max_modules}") ]]; then
full_build_required=true
fi
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,7 @@ metadata:
app: shenyu-zk
all: shenyu-examples-dubbo
spec:
type: NodePort
type: ClusterIP

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit (non-blocking): this change is unrelated to #6775.

NodePort → ClusterIP for the ZooKeeper service is arguably more correct (ZK is only consumed in-cluster by the dubbo example pods, so exposing a NodePort is unnecessary surface), and the k8s ingress CI job passed with it.

But bundling infra/example changes into a bug-fix PR makes backporting and bisecting harder. Please consider moving this (and the set -euo pipefail + resolve-ci-modules.sh edits) into a separate PR next time.

selector:
app: shenyu-zk
all: shenyu-examples-dubbo
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@
# limitations under the License.
#

set -euo pipefail

kind load docker-image "shenyu-examples-apache-dubbo-service:latest"
kind load docker-image "apache/shenyu-integrated-test-k8s-ingress-apache-dubbo:latest"
kubectl apply -f ./shenyu-examples/shenyu-examples-dubbo/shenyu-examples-apache-dubbo-service/k8s/shenyu-zookeeper.yml
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -176,7 +176,8 @@ public List<T> getJoins() {
if (extensionClassesEntity.isEmpty()) {
return Collections.emptyList();
}
if (Objects.equals(extensionClassesEntity.size(), cachedInstances.size())) {
if (Objects.equals(extensionClassesEntity.size(), cachedInstances.size())
&& cachedInstances.values().stream().allMatch(Holder::isInitialized)) {
return (List<T>) this.cachedInstances.values().stream()
.sorted(HOLDER_COMPARATOR)
.map(e -> {
Expand Down Expand Up @@ -222,6 +223,7 @@ private void createExtension(final String name, final Holder<Object> holder) {
}
holder.setOrder(classEntity.getOrder());
holder.setValue(o);
holder.setInitialized(true);
}

/**
Expand Down Expand Up @@ -329,6 +331,8 @@ private void loadClass(final Map<String, ClassEntity> classes,
private static final class Holder<T> {

private volatile T value;

private volatile boolean initialized;

private Integer order;

Expand All @@ -349,6 +353,24 @@ public T getValue() {
public void setValue(final T value) {
this.value = value;
}

/**
* Checks whether the holder is initialized.
*
* @return true if initialized
*/
public boolean isInitialized() {
return initialized;
}

/**
* Sets initialized.
*
* @param initialized initialized
*/
public void setInitialized(final boolean initialized) {
this.initialized = initialized;
}

/**
* set order.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@
import org.apache.shenyu.spi.fixture.TreeListSPI;
import org.junit.jupiter.api.Test;

import java.lang.reflect.Constructor;
import java.lang.reflect.Field;
import java.lang.reflect.InvocationTargetException;
import java.lang.reflect.Method;
import java.net.MalformedURLException;
Expand All @@ -42,12 +44,17 @@
import java.util.Map;
import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;

import static org.hamcrest.CoreMatchers.containsString;
import static org.hamcrest.CoreMatchers.is;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.fail;

Expand Down Expand Up @@ -324,6 +331,53 @@ public void testMultiThreadNonSingleton() throws InterruptedException {
assertEquals(threadNum * loop, cache.size());
}

/**
* Test concurrent get joins when a holder has not finished initialization.
*
* @throws Exception when reflection or concurrent execution fails
*/
@Test
public void testMultiThreadGetJoinsWithUninitializedHolder() throws Exception {
ExtensionLoader<HasDefaultSPI> extensionLoader = newExtensionLoader(HasDefaultSPI.class);
Map<String, Object> cachedInstances = getCachedInstances(extensionLoader);
cachedInstances.put("subHasDefaultSPI", getHolderConstructor().newInstance());
ExecutorService executor = Executors.newFixedThreadPool(4);
try {
List<Future<List<HasDefaultSPI>>> futures = new ArrayList<>();
for (int i = 0; i < 4; i++) {
futures.add(executor.submit(extensionLoader::getJoins));
}
for (Future<List<HasDefaultSPI>> future : futures) {
List<HasDefaultSPI> joins = future.get(5, TimeUnit.SECONDS);
assertEquals(1, joins.size());
assertNotNull(joins.get(0));
}
} finally {
executor.shutdownNow();
}
}

@SuppressWarnings("unchecked")
private <S> ExtensionLoader<S> newExtensionLoader(final Class<S> extensionClass) throws Exception {
Constructor<ExtensionLoader> constructor = ExtensionLoader.class.getDeclaredConstructor(Class.class, ClassLoader.class);
constructor.setAccessible(true);
return (ExtensionLoader<S>) constructor.newInstance(extensionClass, ExtensionLoader.class.getClassLoader());
}

@SuppressWarnings("unchecked")
private Map<String, Object> getCachedInstances(final ExtensionLoader<?> extensionLoader) throws Exception {
Field field = ExtensionLoader.class.getDeclaredField("cachedInstances");
field.setAccessible(true);
return (Map<String, Object>) field.get(extensionLoader);
}

private Constructor<?> getHolderConstructor() throws Exception {
Class<?> holderClass = Class.forName("org.apache.shenyu.spi.ExtensionLoader$Holder");
Constructor<?> constructor = holderClass.getDeclaredConstructor();
constructor.setAccessible(true);
return constructor;
}

/**
* get private loadClass method.
*/
Expand Down
Loading