Aias00 commented on code in PR #7133:
URL: https://github.com/apache/shenyu/pull/7133#discussion_r4072241280


##########
db/init/mysql/schema.sql:
##########
@@ -1018,6 +1028,7 @@ INSERT INTO `plugin` VALUES ('52', 'aiPrompt', null, 
'Ai', 170, 0, '2023-12-20 1
 INSERT INTO `plugin` VALUES ('53', 'aiRequestTransformer', NULL, 'Ai', 65, 0, 
'2023-12-20 18:02:53', '2023-12-20 18:02:53', null);
 
 INSERT INTO `plugin` VALUES ('61', 'mcpServer', null, 'MCP', 180, 0, 
'2023-12-20 18:02:53', '2023-12-20 18:02:53', null);
+INSERT INTO `plugin` VALUES ('67', 'agentGateway', null, 'Ai', 198, 0, 
'2026-09-19 00:00:00', '2026-09-19 00:00:00', null);

Review Comment:
   🔴 **Blocking — primary key collision with master.**
   
   `plugin.id = '67'` was assigned to `sensitiveWord` on master (2026-09-21):
   ```sql
   INSERT INTO `plugin` VALUES ('67', 'sensitiveWord', null, 'Ai', 197, 0, 
'2026-09-21 00:00:00', ...);
   ```
   Inserting `('67', 'agentGateway', ...)` here fails with a duplicate primary 
key, so a fresh install cannot initialise the schema.
   
   Please rebase and move to the next free id (`68`) — and apply the same 
change to all six init scripts (`db/init/mysql`, `ob`, `og`, `oracle`, `pg`, 
and `shenyu-admin/src/main/resources/sql-script/h2/schema.sql`).
   
   The plugin **code** `198` itself is free and correctly ordered just before 
`AI_PROXY(199)`, so only the row id needs to change.



##########
db/init/mysql/schema.sql:
##########
@@ -2070,6 +2083,17 @@ INSERT INTO `resource` VALUES ('1844026199075534867', 
'1844026199075534860', 'SH
 INSERT INTO `resource` VALUES ('1844026199075534868', '1844026199075534860', 
'SHENYU.BUTTON.PLUGIN.RULE.DELETE', '', '', '', 2, 0, '', 1, 0, 
'plugin:mcpServerRule:delete', 1, '2022-05-25 18:02:58', '2022-05-25 18:02:58');
 INSERT INTO `resource` VALUES ('1844026199075534869', '1844026199075534860', 
'SHENYU.BUTTON.PLUGIN.SYNCHRONIZE', '', '', '', 2, 0, '', 1, 0, 
'plugin:mcpServer:modify', 1, '2022-05-25 18:02:58', '2022-05-25 18:02:58');
 
+INSERT INTO `resource` VALUES ('1942847622591684620', '1346775491550474240', 
'agentGateway', 'agentGateway', '/plug/agentGateway', 'agentGateway', 1, 0, 
'pic-center', 0, 0, '', 1, '2026-09-19 00:00:00', '2026-09-19 00:00:00');

Review Comment:
   🔴 **Blocking — these `resource` ids are already used on master.**
   
   `1942847622591684620`, `…621` and `…622` are already taken by 
`plugin_handle` rows belonging to `sensitiveWord` (plugin 67):
   ```sql
   INSERT INTO `plugin_handle` VALUES ('1942847622591684620', '67', 
'failClosed', ...);
   INSERT INTO `plugin_handle` VALUES ('1942847622591684621', '67', 'words', 
...);
   INSERT INTO `plugin_handle` VALUES ('1942847622591684622', '67', 
'maxBodySize', ...);
   ```
   Please re-allocate this id block, and double-check the paired `permission` 
ids `1942847622591684630–639` while you are at it.



##########
shenyu-plugin/shenyu-plugin-agent-gateway/src/main/java/org/apache/shenyu/plugin/agent/gateway/AgentGatewayPlugin.java:
##########
@@ -0,0 +1,119 @@
+/*
+ * 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.shenyu.plugin.agent.gateway;
+
+import org.apache.shenyu.common.dto.RuleData;
+import org.apache.shenyu.common.dto.SelectorData;
+import org.apache.shenyu.common.enums.PluginEnum;
+import org.apache.shenyu.plugin.agent.gateway.handle.AgentGatewayRuleHandle;
+import 
org.apache.shenyu.plugin.agent.gateway.handle.AgentGatewayRuleHandleParser;
+import 
org.apache.shenyu.plugin.agent.gateway.handler.AgentGatewayPluginDataHandler;
+import org.apache.shenyu.plugin.api.ShenyuPluginChain;
+import org.apache.shenyu.plugin.api.result.ShenyuResultWrap;
+import org.apache.shenyu.plugin.api.utils.WebFluxResultUtils;
+import org.apache.shenyu.plugin.base.AbstractShenyuPlugin;
+import org.apache.shenyu.plugin.base.utils.CacheKeyUtils;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.HttpStatus;
+import org.springframework.web.server.ServerWebExchange;
+import reactor.core.publisher.Mono;
+
+import java.util.Objects;
+import java.util.UUID;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+/**
+ * Establishes an isolated request context for agent traffic and continues to
+ * the existing AI proxy chain.
+ */
+public class AgentGatewayPlugin extends AbstractShenyuPlugin {
+
+    private final AgentGatewayRuleHandleParser parser = new 
AgentGatewayRuleHandleParser();
+
+    @Override
+    protected Mono<Void> doExecute(final ServerWebExchange exchange, final 
ShenyuPluginChain chain,
+                                   final SelectorData selector, final RuleData 
rule) {
+        if (Boolean.FALSE.equals(selector.getContinued())) {
+            return reject(exchange, "selector continued=false is not 
supported");
+        }
+        final AgentGatewayRuleHandle handle = resolveHandle(rule);
+        if (!handle.isValid()) {
+            return reject(exchange, handle.getErrorMessage());
+        }
+        final AtomicBoolean subscribed = new AtomicBoolean();
+        return Mono.defer(() -> {
+            if (!subscribed.compareAndSet(false, true)) {
+                return Mono.error(new IllegalStateException("agent gateway 
execution cannot be subscribed twice"));
+            }
+            final AgentTrafficContext context = new AgentTrafficContext(
+                    UUID.randomUUID().toString(), handle.getTrafficType(), 
selector.getId(), rule.getId());
+            final Object previousContext = exchange.getAttributes()
+                    .put(AgentGatewayConstants.REQUEST_CONTEXT_ATTRIBUTE, 
context);
+            if (handle.isResponseRequestId()) {
+                registerResponseRequestId(exchange, context.getRequestId());
+            }
+            return chain.execute(exchange)
+                    .contextWrite(reactorContext -> 
reactorContext.put(AgentGatewayConstants.REACTOR_CONTEXT_KEY, context))
+                    .doFinally(signal -> restorePreviousContext(exchange, 
previousContext, context));
+        });
+    }
+
+    private void restorePreviousContext(final ServerWebExchange exchange, 
final Object previousContext,
+                                        final AgentTrafficContext 
currentContext) {
+        if (Objects.isNull(previousContext)) {
+            
exchange.getAttributes().remove(AgentGatewayConstants.REQUEST_CONTEXT_ATTRIBUTE,
 currentContext);
+        } else {
+            
exchange.getAttributes().replace(AgentGatewayConstants.REQUEST_CONTEXT_ATTRIBUTE,
 currentContext,
+                    previousContext);
+        }
+    }
+
+    private AgentGatewayRuleHandle resolveHandle(final RuleData rule) {
+        final String key = CacheKeyUtils.INST.getKey(rule);
+        final AgentGatewayRuleHandle cached = 
AgentGatewayPluginDataHandler.CACHED_HANDLE.get().obtainHandle(key);
+        if (Objects.nonNull(cached) && Objects.equals(cached.getRawHandle(), 
rule.getHandle())) {
+            return cached;
+        }
+        return parser.parse(rule.getHandle());
+    }
+
+    private void registerResponseRequestId(final ServerWebExchange exchange, 
final String requestId) {
+        exchange.getResponse().beforeCommit(() -> {
+            final HttpHeaders headers = exchange.getResponse().getHeaders();
+            headers.set(AgentGatewayConstants.REQUEST_ID_HEADER, requestId);
+            return Mono.empty();
+        });
+    }
+
+    private Mono<Void> reject(final ServerWebExchange exchange, final String 
reason) {
+        exchange.getResponse().setStatusCode(HttpStatus.SERVICE_UNAVAILABLE);

Review Comment:
   🟠 `503 SERVICE_UNAVAILABLE` is the wrong status for a **configuration** 
error.
   
   503 tells clients and upstream proxies that the failure is transient and 
worth retrying, but a malformed rule handle will never start working — so this 
invites pointless retry storms. Please use `500 INTERNAL_SERVER_ERROR` 
(server-side misconfiguration) or `400`, and keep 503 only for genuine 
unavailability.



##########
shenyu-plugin/shenyu-plugin-agent-gateway/src/main/java/org/apache/shenyu/plugin/agent/gateway/AgentGatewayPlugin.java:
##########
@@ -0,0 +1,119 @@
+/*
+ * 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.shenyu.plugin.agent.gateway;
+
+import org.apache.shenyu.common.dto.RuleData;
+import org.apache.shenyu.common.dto.SelectorData;
+import org.apache.shenyu.common.enums.PluginEnum;
+import org.apache.shenyu.plugin.agent.gateway.handle.AgentGatewayRuleHandle;
+import 
org.apache.shenyu.plugin.agent.gateway.handle.AgentGatewayRuleHandleParser;
+import 
org.apache.shenyu.plugin.agent.gateway.handler.AgentGatewayPluginDataHandler;
+import org.apache.shenyu.plugin.api.ShenyuPluginChain;
+import org.apache.shenyu.plugin.api.result.ShenyuResultWrap;
+import org.apache.shenyu.plugin.api.utils.WebFluxResultUtils;
+import org.apache.shenyu.plugin.base.AbstractShenyuPlugin;
+import org.apache.shenyu.plugin.base.utils.CacheKeyUtils;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.HttpStatus;
+import org.springframework.web.server.ServerWebExchange;
+import reactor.core.publisher.Mono;
+
+import java.util.Objects;
+import java.util.UUID;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+/**
+ * Establishes an isolated request context for agent traffic and continues to
+ * the existing AI proxy chain.
+ */
+public class AgentGatewayPlugin extends AbstractShenyuPlugin {
+
+    private final AgentGatewayRuleHandleParser parser = new 
AgentGatewayRuleHandleParser();
+
+    @Override
+    protected Mono<Void> doExecute(final ServerWebExchange exchange, final 
ShenyuPluginChain chain,
+                                   final SelectorData selector, final RuleData 
rule) {
+        if (Boolean.FALSE.equals(selector.getContinued())) {
+            return reject(exchange, "selector continued=false is not 
supported");
+        }
+        final AgentGatewayRuleHandle handle = resolveHandle(rule);
+        if (!handle.isValid()) {
+            return reject(exchange, handle.getErrorMessage());
+        }
+        final AtomicBoolean subscribed = new AtomicBoolean();
+        return Mono.defer(() -> {
+            if (!subscribed.compareAndSet(false, true)) {
+                return Mono.error(new IllegalStateException("agent gateway 
execution cannot be subscribed twice"));

Review Comment:
   🟠 This guard makes the returned `Mono` non-idempotent: any second 
subscription (e.g. a `retry`, `repeat`, or `cache()` composed downstream) turns 
into an `IllegalStateException` and fails the request.
   
   WebFlux subscribes once, so if double subscription is not a scenario you 
actually hit, I would drop the guard. If it *is* a real scenario, handling it 
should not mean failing the request — consider logging and reusing the 
already-built context instead.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to