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();

Reply via email to