This is an automated email from the ASF dual-hosted git repository.
Aias00 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shenyu.git
The following commit(s) were added to refs/heads/master by this push:
new 39250e16a8 fix(divide): release p2c inflight on cancellation (#7096)
39250e16a8 is described below
commit 39250e16a851fb52344b8cdbd27189a9906fec37
Author: Liming Deng <[email protected]>
AuthorDate: Wed Sep 30 09:57:45 2026 +0800
fix(divide): release p2c inflight on cancellation (#7096)
Co-authored-by: zhengpeng <[email protected]>
Co-authored-by: aias00 <[email protected]>
---
.../apache/shenyu/plugin/divide/DividePlugin.java | 3 +--
.../shenyu/plugin/divide/DividePluginTest.java | 22 ++++++++++++++++++++++
2 files changed, 23 insertions(+), 2 deletions(-)
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/main/java/org/apache/shenyu/plugin/divide/DividePlugin.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/main/java/org/apache/shenyu/plugin/divide/DividePlugin.java
index cffa736a3a..0c55ed7250 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/main/java/org/apache/shenyu/plugin/divide/DividePlugin.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/main/java/org/apache/shenyu/plugin/divide/DividePlugin.java
@@ -131,8 +131,7 @@ public class DividePlugin extends AbstractShenyuPlugin {
exchange.getAttributes().put(Constants.LOAD_BALANCE,
StringUtils.defaultIfEmpty(ruleHandle.getLoadBalance(),
LoadBalanceEnum.RANDOM.getName()));
exchange.getAttributes().put(Constants.DIVIDE_SELECTOR_ID,
selector.getId());
if (ruleHandle.getLoadBalance().equals(P2C)) {
- return chain.execute(exchange).doOnSuccess(e ->
responseTrigger(upstream
- )).doOnError(throwable -> responseTrigger(upstream));
+ return chain.execute(exchange).doFinally(signalType ->
responseTrigger(upstream));
} else if (ruleHandle.getLoadBalance().equals(SHORTEST_RESPONSE)) {
long beginTime = System.currentTimeMillis();
return chain.execute(exchange).doOnSuccess(e ->
successResponseTrigger(upstream, beginTime
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/test/java/org/apache/shenyu/plugin/divide/DividePluginTest.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/test/java/org/apache/shenyu/plugin/divide/DividePluginTest.java
index d1b73b0808..8a5a0d5577 100644
---
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/test/java/org/apache/shenyu/plugin/divide/DividePluginTest.java
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-divide/src/test/java/org/apache/shenyu/plugin/divide/DividePluginTest.java
@@ -37,6 +37,7 @@ import
org.apache.shenyu.plugin.api.result.DefaultShenyuResult;
import org.apache.shenyu.plugin.api.result.ShenyuResult;
import org.apache.shenyu.plugin.api.utils.SpringBeanUtils;
import org.apache.shenyu.plugin.base.utils.CacheKeyUtils;
+import org.apache.shenyu.plugin.base.utils.LoadbalancerUtils;
import org.apache.shenyu.plugin.divide.handler.DividePluginDataHandler;
import org.apache.shenyu.plugin.divide.handler.DivideUpstreamDataHandler;
import org.junit.jupiter.api.AfterEach;
@@ -51,6 +52,7 @@ import
org.springframework.context.ConfigurableApplicationContext;
import org.springframework.mock.http.server.reactive.MockServerHttpRequest;
import org.springframework.mock.web.server.MockServerWebExchange;
import org.springframework.web.server.ServerWebExchange;
+import reactor.core.Disposable;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
@@ -260,6 +262,26 @@ public final class DividePluginTest {
assertNotEquals(0, upstream.getLag());
}
+ @Test
+ public void p2cInflightShouldBeReleasedOnCancellation() {
+ DivideRuleHandle ruleHandle =
DividePluginDataHandler.CACHED_HANDLE.get()
+ .obtainHandle(CacheKeyUtils.INST.getKey(ruleData));
+ ruleHandle.setLoadBalance("p2c");
+ Upstream upstream = Upstream.builder().url("http://upstream").build();
+ upstream.getInflight().set(2);
+ when(chain.execute(exchange)).thenReturn(Mono.never());
+
+ try (MockedStatic<LoadbalancerUtils> loadbalancerUtils =
mockStatic(LoadbalancerUtils.class)) {
+ loadbalancerUtils.when(() ->
LoadbalancerUtils.getForExchange(any(), anyString(), any()))
+ .thenReturn(upstream);
+ Disposable subscription = dividePlugin.doExecute(exchange, chain,
selectorData, ruleData).subscribe();
+
+ subscription.dispose();
+
+ assertEquals(1, upstream.getInflight().get());
+ }
+ }
+
@Test
public void successResponseTriggerTest() throws Exception {
dividePlugin = DividePlugin.class.newInstance();