This is an automated email from the ASF dual-hosted git repository.

CalvinKirs pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new 958aaf46c8f [fix](protocol) Prevent Connector/J cursor fetch from 
hanging on empty results (#67520)
958aaf46c8f is described below

commit 958aaf46c8f6b7bd6fb661a18e5b7cbbf5c40999
Author: Calvin Kirs <[email protected]>
AuthorDate: Thu Sep 10 10:44:30 2026 +0800

    [fix](protocol) Prevent Connector/J cursor fetch from hanging on empty 
results (#67520)
    
    ### What problem does this PR solve?
    
    Issue Number: None
    
    Related PR: #61050, #61062
    
    After `CLIENT_DEPRECATE_EOF` is negotiated, affected Connector/J clients
    using `useCursorFetch=true` and a positive fetch size can hang when a
    server-prepared statement returns an empty result. The driver consumes
    the first terminator after column definitions while checking whether a
    server cursor was created. With no rows, this consumes the only
    result-end marker, so the driver waits for a packet that Doris will
    never send.
    
    The reported Connector/J 8.2.0 case is covered by real-driver tests.
    This PR preserves the final result marker for affected clients and
    avoids inserting the compatibility marker for identified Connector/J
    9.5+ and MariaDB clients. It also preserves this behavior through FE
    forwarding and fixes related packet/capability regressions.
    
    ### What is changed?
    
    1. **Cursor result boundaries.** Read `CURSOR_TYPE_READ_ONLY` from each
    `COM_STMT_EXECUTE`, retain it in the request context, and classify
    driver behavior using connection attributes. For affected cursor clients
    with deprecated EOF enabled, insert a compatibility ResultSet OK after
    metadata so the real final marker remains available. Identified modern
    clients and ordinary non-cursor requests keep the standard sequence.
    2. **FE forwarding.** Forward cursor intent and negotiated capabilities
    through optional Thrift fields. The master restores them before
    serializing results. A new follower adapts buffered results from an
    older master that has not confirmed protocol adaptation:
    insert/remove/convert metadata boundaries as appropriate and convert the
    final legacy EOF when needed. Row payloads are retained; confirmed
    results are not adapted twice.
    3. **Complete OK information.** Preserve the original forwarded DML/DDL
    OK and ERR packets rather than rebuilding a partial successful OK.
    Retain affected rows, warnings, info, load labels and transaction IDs.
    Preserve warnings and server status when producing/converting result-set
    terminators. SQL is not re-executed by this adaptation.
    4. **Legacy protocol compatibility.** Configure the serializer using the
    intersection of client/server capabilities, rather than server defaults.
    Limit the empty OK-info compatibility byte to the applicable negotiated
    format, preserving legacy OK behavior.
    5. **Arrow follower compatibility.** Access MySQL-channel capabilities
    only for MySQL connections, so the shared forwarding path does not call
    the unsupported MySQL-channel getter on an Arrow Flight SQL context.
    6. **Clients without connection attributes.** Remove the new fail-fast
    rejection of anonymous cursor clients. When client behavior is unknown,
    prioritize the older Connector/J compatibility sequence. This retains
    the older-client target but cannot simultaneously guarantee
    compatibility with anonymous clients requiring the standard sequence.
    7. **File-load compatibility.** In both LOAD handlers, require
    `CLIENT_LOCAL_FILES` only for client-side `LOCAL INFILE` uploads.
    FE-side file reads do not require a client-upload capability and remain
    usable when that bit is absent.
    
    Full PR diff against the merged master, at `32d01e5df19`:
    
    | Category | Files | Added lines | Removed lines | Net added lines |
    |---|---:|---:|---:|---:|
    | Production code, including Thrift IDL | 12 | 218 | 40 | 178 |
    | Tests and expected output | 12 | 787 | 13 | 774 |
    
    The test total includes 11 automatically generated expected-output
    lines. Environment configuration and review documents are excluded.
    
    ### Compatibility boundaries
    
    - This preserves streamed results; it does not implement MySQL
    server-side cursors or `COM_STMT_FETCH` batching.
    - **New follower → old master** is covered by the mixed-version client
    matrix. **Old follower → new master** may have already lost cursor
    intent; the master cannot reconstruct it. The historical cursor hang on
    that route remains until the client-facing follower is upgraded.
    Existing ordinary prepared-statement behavior is retained.
    - Identified Connector/J 9.5+ and MariaDB clients use the standard path.
    Anonymous clients requiring that sequence remain outside the
    older-client fallback guarantee.
    - The OIDC provider and authentication routing implementation are not
    changed by this PR.
    
    ### Release note
    
    Fix affected Connector/J cursor queries hanging on empty prepared
    results. Preserve negotiated MySQL packet formats, forwarded DML/DDL OK
    information, Arrow Flight SQL forwarding and FE-side file loading.
    Upgrade the client-facing follower to preserve cursor intent during an
    FE rolling upgrade.
    
    ### Check List (For Author)
    
    - Test:
    - [x] Unit Test: **113 targeted FE tests passed, 0 failures/errors**,
    before the final master/JUnit 5 merge. Coverage includes driver
    classification, cursor boundaries, forwarding, capabilities, packet
    serialization, authentication routing and both LOAD handlers.
    - [x] Regression test: `prepared_stmt_p0/cursor_fetch_empty_result`,
    `prepared_stmt_p0/prepared_show`, and `arrow_flight_sql_p0/test_ddl`
    through a real follower passed. Cursor expected output was generated by
    the regression runner.
    - [x] Manual test: **120 successful real-JDBC driver/endpoint/transport
    runs** across baseline/candidate direct endpoints, a same-version
    follower and a new follower targeting an old master, with plaintext/TLS
    modes. Each driver JAR runs in its own JVM to avoid classpath version
    collisions.
    - MySQL Connector/J: **5.1.49, 8.0.28, 8.0.33, 8.2.0, 8.4.0, 9.0.0,
    9.4.0, 9.5.0, 9.6.0**. MariaDB Connector/J: **3.5.6**.
    - Each run exercises 12 requested combinations of cursor/server-prepare
    switches and fetch sizes **0/1/10000**, repeated
    Statement/PreparedStatement empty→nonempty→empty results on reused
    connections, and wrong-password rejection. Drivers may couple cursor and
    server-prepare settings internally. Connection/read timeouts bound
    failures; a timeout is not counted as success.
    - Additional checks passed: explicit TLS 1.2/1.3 (16 groups), MySQL
    Shell 8/9 ordinary authentication (16), legacy EOF via PyMySQL (8),
    SQL-error connection reuse (8), and suppressed connection attributes
    with older drivers (12).
    - Actual baseline/candidate FE-file and client-file imports passed;
    forwarded INSERT retained label/status/txnId. Real LDAP/local-password
    routing, rejected credentials and connection reuse passed 20
    baseline/candidate checks across plaintext/TLS.
    - Independent FE-only raw-protocol tests passed 12 configurations
    covering multi-result status, cursor transitions, RESET/CLOSE and reuse
    after errors. These are packet tests, not additional real JDBC runs or
    proof of mixed-version multi-statement forwarding.
    - FE build passed before the final master merge. After the merge,
    standalone `mvn checkstyle:check -pl fe-core -Dcheckstyle.skip=false`
    passed with **0 violations**. Compilation and unit tests were not rerun
    after the JUnit 5 migration.
    - The multi-version JDBC matrix was executed locally with retained
    scripts/logs; it is not an automatic all-version CI matrix.
    - Review and validation limits:
    - The latest independent review of the complete final diff found **0
    Blocker / 0 Major production findings**. It noted two deterministic
    SQL-result assertions that still use `assertEquals` rather than
    regression `qt` output, and missing dedicated tests for errors occurring
    after partial metadata or rows have been buffered. Those error paths
    were inspected in source, not fault-injected; the OK/ERR preservation
    unit test does not cover buffered intermediate errors.
    - Connector/J 6.0.6 failed handshakes on both baseline and candidate in
    six runs; it is not counted as passing or newly supported.
    - Four broader prepared suites failed identically on baseline/candidate
    with the available older BE because it cannot execute the plan.
    Authentication-integration metadata regression also remains incomplete
    with that BE. Four independent SELECT configurations without a BE were
    blocked and not counted as passing.
    - Full product OIDC-provider login/new TLS-extension E2E and the
    customer's exact SmartBI environment were not validated. Ordinary
    password/TLS and authentication unit tests do not replace these checks.
    - Behavior changed:
    - [x] Yes. Correct cursor packet boundaries, preserve complete forwarded
    OK information and use negotiated capabilities, with the compatibility
    boundaries above.
    - Does this need documentation?
        - [x] No.
    
    ### Check List (For Reviewer who merge this PR)
    
    - [ ] Confirm the release note
    - [ ] Confirm test cases
    - [ ] Confirm document
    - [ ] Add branch pick label
---
 .../doris/mysql/MysqlCursorFetchCompatibility.java |  61 ++++++
 .../java/org/apache/doris/mysql/MysqlOkPacket.java |  15 +-
 .../java/org/apache/doris/mysql/MysqlProto.java    |  10 +-
 .../doris/mysql/MysqlResultSetEndPacket.java       |   1 +
 .../nereids/trees/plans/commands/LoadCommand.java  |   2 +-
 .../plans/commands/load/MysqlLoadCommand.java      |   2 +-
 .../java/org/apache/doris/qe/ConnectContext.java   |  12 ++
 .../java/org/apache/doris/qe/ConnectProcessor.java |  27 ++-
 .../java/org/apache/doris/qe/FEOpExecutor.java     |  77 +++++++
 .../org/apache/doris/qe/MysqlConnectProcessor.java |  12 +-
 .../java/org/apache/doris/qe/StmtExecutor.java     |  36 ++--
 .../mysql/MysqlCursorFetchCompatibilityTest.java   |  59 ++++++
 .../org/apache/doris/mysql/MysqlOkPacketTest.java  |  19 +-
 .../org/apache/doris/mysql/MysqlProtoTest.java     |  22 ++
 .../doris/mysql/MysqlResultSetEndPacketTest.java   |  26 +++
 .../trees/plans/commands/MysqlLoadCommandTest.java |  52 +++++
 .../qe/ConnectProcessorForwardProtocolTest.java    | 148 +++++++++++++
 .../qe/FEOpExecutorDelegatedCredentialTest.java    |   2 +
 .../doris/qe/FEOpExecutorFlightForwardTest.java    |   4 +
 .../doris/qe/FEOpExecutorMysqlProtocolTest.java    | 229 +++++++++++++++++++++
 .../qe/MysqlConnectProcessorCursorFetchTest.java   | 112 ++++++++++
 .../java/org/apache/doris/qe/StmtExecutorTest.java |  98 +++++++++
 gensrc/thrift/FrontendService.thrift               |   6 +
 .../prepared_stmt_p0/cursor_fetch_empty_result.out |  11 +
 .../suites/arrow_flight_sql_p0/test_ddl.groovy     |  25 +++
 .../cursor_fetch_empty_result.groovy               | 107 ++++++++++
 26 files changed, 1120 insertions(+), 55 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlCursorFetchCompatibility.java
 
b/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlCursorFetchCompatibility.java
new file mode 100644
index 00000000000..981ce75da47
--- /dev/null
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlCursorFetchCompatibility.java
@@ -0,0 +1,61 @@
+// 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.doris.mysql;
+
+import com.google.common.collect.ImmutableSet;
+
+import java.util.Map;
+import java.util.Set;
+import java.util.regex.Pattern;
+
+/** Resolves the incompatible cursor result-set behavior used by Connector/J 
releases. */
+public final class MysqlCursorFetchCompatibility {
+    private static final Set<String> MYSQL_CONNECTOR_J_CLIENT_NAMES = 
ImmutableSet.of(
+            "MySQL Connector/J", "MySQL Connector Java");
+    private static final Pattern CONSUMES_METADATA_TERMINATOR =
+            Pattern.compile("^(?:(?:5|6|8)\\.|9\\.[0-4](?:\\.|$))");
+    private static final Pattern VERSION = 
Pattern.compile("^\\d+(?:\\.\\d+)+(?:[-+].*)?$");
+
+    public enum Behavior {
+        CONSUMES_METADATA_TERMINATOR,
+        STANDARD,
+        UNKNOWN
+    }
+
+    private MysqlCursorFetchCompatibility() {
+    }
+
+    public static Behavior resolve(Map<String, String> connectAttributes) {
+        String clientName = connectAttributes.get("_client_name");
+        if (clientName == null) {
+            return Behavior.UNKNOWN;
+        }
+        if (!MYSQL_CONNECTOR_J_CLIENT_NAMES.contains(clientName)) {
+            return Behavior.STANDARD;
+        }
+
+        String clientVersion = connectAttributes.get("_client_version");
+        if (clientVersion == null || 
!VERSION.matcher(clientVersion).matches()) {
+            return Behavior.UNKNOWN;
+        }
+        if (CONSUMES_METADATA_TERMINATOR.matcher(clientVersion).find()) {
+            return Behavior.CONSUMES_METADATA_TERMINATOR;
+        }
+        return Behavior.STANDARD;
+    }
+}
diff --git a/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlOkPacket.java 
b/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlOkPacket.java
index 4fa80102317..f9cf43b3691 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlOkPacket.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlOkPacket.java
@@ -58,16 +58,11 @@ public class MysqlOkPacket extends MysqlPacket {
             // TODO(zhaochun): STATUS_FLAGS
             // if ((STATUS_FLAGS & 
MysqlStatusFlag.SERVER_SESSION_STATE_CHANGED) != 0) {
             // }
-        } else {
-            // Always write the info field as a length-encoded string.
-            // When CLIENT_DEPRECATE_EOF is negotiated, the driver's 
OkPacket.parse()
-            // unconditionally reads STRING_LENENC for info, so an empty 
string must
-            // still be written (as a single 0x00 byte representing length 0).
-            if (Strings.isNullOrEmpty(infoMessage)) {
-                serializer.writeVInt(0);
-            } else {
-                serializer.writeLenEncodedString(infoMessage);
-            }
+        } else if (!Strings.isNullOrEmpty(infoMessage)) {
+            serializer.writeLenEncodedString(infoMessage);
+        } else if (capability.isDeprecatedEOF()) {
+            // Connector/J parses the info field for CLIENT_DEPRECATE_EOF even 
when it is empty.
+            serializer.writeVInt(0);
         }
     }
 }
diff --git a/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlProto.java 
b/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlProto.java
index 55960c36b89..1137f688fea 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlProto.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlProto.java
@@ -183,10 +183,6 @@ public class MysqlProto {
         if (handshakeResponse == null) {
             return failHandshake(context, 
CLIENT_CLOSED_CONNECTION_DURING_HANDSHAKE);
         }
-        if (capability.isDeprecatedEOF()) {
-            context.getMysqlChannel().setClientDeprecatedEOF();
-        }
-
         // we do not save client capability to context, so here we save 
CLIENT_MULTI_STATEMENTS to MysqlChannel
         if (capability.isClientMultiStatements()) {
             context.getMysqlChannel().setClientMultiStatements();
@@ -208,7 +204,11 @@ public class MysqlProto {
         }
 
         // change the capability of serializer
-        context.setCapability(context.getServerCapability());
+        context.setCapability(new 
MysqlCapability(context.getServerCapability().getFlags()
+                & authPacket.getCapability().getFlags()));
+        if (context.getCapability().isDeprecatedEOF()) {
+            channel.setClientDeprecatedEOF();
+        }
         serializer.setCapability(context.getCapability());
 
         String qualifiedUser = parseUser(context, 
authPacket.getAuthResponse(), authPacket.getUser());
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlResultSetEndPacket.java 
b/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlResultSetEndPacket.java
index 5543b85c361..fd76924949d 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlResultSetEndPacket.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/mysql/MysqlResultSetEndPacket.java
@@ -40,6 +40,7 @@ public class MysqlResultSetEndPacket extends MysqlPacket {
 
     public MysqlResultSetEndPacket(QueryState state) {
         this.serverStatus = state.serverStatus;
+        this.warningCount = state.getWarningRows();
     }
 
     @Override
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/LoadCommand.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/LoadCommand.java
index d5d616b8f5a..e7f3baa1a72 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/LoadCommand.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/LoadCommand.java
@@ -496,7 +496,7 @@ public class LoadCommand extends Command implements 
NeedAuditEncryption, Forward
             }
             LoadManager loadManager = ctx.getEnv().getLoadManager();
             if (etlJobType == EtlJobType.LOCAL_FILE) {
-                if (!ctx.getCapability().supportClientLocalFile()) {
+                if (getDataDescriptions().get(0).isClientLocal() && 
!ctx.getCapability().supportClientLocalFile()) {
                     ctx.getState().setError(ErrorCode.ERR_NOT_ALLOWED_COMMAND, 
"This client is not support"
                             + " to load client local file.");
                     return;
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/load/MysqlLoadCommand.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/load/MysqlLoadCommand.java
index abc2c8952bc..849252b455d 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/load/MysqlLoadCommand.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/load/MysqlLoadCommand.java
@@ -219,7 +219,7 @@ public class MysqlLoadCommand extends Command implements 
NoForward {
     private void handleMysqlLoadCommand(ConnectContext ctx) {
         try {
             LoadManager loadManager = ctx.getEnv().getLoadManager();
-            if (!ctx.getCapability().supportClientLocalFile()) {
+            if (mysqlDataDescription.isClientLocal() && 
!ctx.getCapability().supportClientLocalFile()) {
                 ctx.getState().setError(ErrorCode.ERR_NOT_ALLOWED_COMMAND, 
"This client is not support"
                         + " to load client local file.");
                 return;
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/ConnectContext.java 
b/fe/fe-core/src/main/java/org/apache/doris/qe/ConnectContext.java
index dc327a0c73d..c17523a7815 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/ConnectContext.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/ConnectContext.java
@@ -287,6 +287,9 @@ public class ConnectContext {
     @Setter
     private ByteBuffer prepareExecuteBuffer;
 
+    // Whether the current COM_STMT_EXECUTE requested a server-side read-only 
cursor.
+    private boolean cursorFetchRequested;
+
     private MysqlHandshakePacket mysqlHandshakePacket;
 
     public void setUserQueryTimeout(int queryTimeout) {
@@ -513,6 +516,14 @@ public class ConnectContext {
         this.connectAttributes = new HashMap<>(connectAttributes);
     }
 
+    public boolean isCursorFetchRequested() {
+        return cursorFetchRequested;
+    }
+
+    public void setCursorFetchRequested(boolean cursorFetchRequested) {
+        this.cursorFetchRequested = cursorFetchRequested;
+    }
+
     public boolean isTxnModel() {
         return txnEntry != null && txnEntry.isTxnModel();
     }
@@ -1008,6 +1019,7 @@ public class ConnectContext {
         statementContext = null;
         loadBackendSelectionDecision = null;
         loadBackendSelectionHint = null;
+        cursorFetchRequested = false;
     }
 
     // Arrow Flight SQL only.
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/ConnectProcessor.java 
b/fe/fe-core/src/main/java/org/apache/doris/qe/ConnectProcessor.java
index 18b6a4d64ad..91e64904b5c 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/ConnectProcessor.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/ConnectProcessor.java
@@ -54,6 +54,7 @@ import org.apache.doris.datasource.CatalogIf;
 import org.apache.doris.datasource.DelegatedCredential;
 import org.apache.doris.datasource.SessionContext;
 import org.apache.doris.metric.MetricRepo;
+import org.apache.doris.mysql.MysqlCapability;
 import org.apache.doris.mysql.MysqlChannel;
 import org.apache.doris.mysql.MysqlCommand;
 import org.apache.doris.mysql.MysqlPacket;
@@ -754,15 +755,10 @@ public abstract class ConnectProcessor {
         // set compute group
         
ctx.setComputeGroup(Env.getCurrentEnv().getAuth().getComputeGroup(ctx.getQualifiedUser()));
 
-        // Propagate the client's CLIENT_DEPRECATE_EOF capability to the proxy 
channel.
-        // This ensures the master generates packets matching the original 
client's protocol.
-        if (request.isSetClientDeprecatedEOF() && 
request.isClientDeprecatedEOF()) {
-            ctx.getMysqlChannel().setClientDeprecatedEOF();
-        }
-
         ctx.setThreadLocalInfo();
         StmtExecutor executor = null;
         try {
+            restoreForwardedMysqlContext(ctx, request);
             // 0 for compatibility.
             int idx = request.isSetStmtIdx() ? request.getStmtIdx() : 0;
             executor = new StmtExecutor(ctx, new 
OriginStatement(request.getSql(), idx), true);
@@ -848,6 +844,7 @@ public abstract class ConnectProcessor {
             ctx.getState().serverStatus |= 
MysqlServerStatusFlag.SERVER_MORE_RESULTS_EXISTS;
         }
         result.setPacket(getResultPacket());
+        
result.setClientDeprecatedEofApplied(ctx.getMysqlChannel().clientDeprecatedEOF());
         result.setStatus(ctx.getState().toString());
         if (ctx.getState().getStateType() == MysqlStateType.OK) {
             result.setStatusCode(0);
@@ -883,6 +880,24 @@ public abstract class ConnectProcessor {
         return result;
     }
 
+    static void restoreForwardedMysqlContext(ConnectContext context, 
TMasterOpRequest request) {
+        int flags = request.isSetMysqlCapability() ? 
request.getMysqlCapability()
+                : MysqlCapability.DEFAULT_CAPABILITY.getFlags()
+                        & 
~MysqlCapability.Flag.CLIENT_DEPRECATE_EOF.getFlagBit();
+        if (request.isSetClientDeprecatedEOF() && 
request.isClientDeprecatedEOF()) {
+            flags |= MysqlCapability.Flag.CLIENT_DEPRECATE_EOF.getFlagBit();
+        }
+        MysqlCapability capability = new MysqlCapability(flags);
+        context.setCapability(capability);
+        context.getMysqlChannel().getSerializer().setCapability(capability);
+        if (capability.isDeprecatedEOF()) {
+            context.getMysqlChannel().setClientDeprecatedEOF();
+        }
+        // Old followers do not carry the cursor flag. Keep their existing 
behavior; they must
+        // be upgraded to preserve cursor intent. Do not reject their ordinary 
prepared statements.
+        context.setCursorFetchRequested(request.isSetCursorFetchRequested() && 
request.isCursorFetchRequested());
+    }
+
     static void restoreForwardedSessionContext(ConnectContext context, 
TMasterOpRequest request) {
         if (!request.isSetDelegatedCredentialToken()) {
             return;
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/FEOpExecutor.java 
b/fe/fe-core/src/main/java/org/apache/doris/qe/FEOpExecutor.java
index 414a9725c70..c6f9049f014 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/FEOpExecutor.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/FEOpExecutor.java
@@ -26,6 +26,10 @@ import org.apache.doris.common.Config;
 import org.apache.doris.common.ErrorCode;
 import org.apache.doris.datasource.DelegatedCredential;
 import org.apache.doris.mysql.MysqlCommand;
+import org.apache.doris.mysql.MysqlCursorFetchCompatibility;
+import org.apache.doris.mysql.MysqlProto;
+import org.apache.doris.mysql.MysqlResultSetEndPacket;
+import org.apache.doris.mysql.MysqlSerializer;
 import org.apache.doris.qe.ConnectContext.ConnectType;
 import org.apache.doris.thrift.FrontendService;
 import org.apache.doris.thrift.TExpr;
@@ -46,6 +50,7 @@ import org.apache.thrift.TException;
 import org.apache.thrift.transport.TTransportException;
 
 import java.nio.ByteBuffer;
+import java.util.ArrayList;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
@@ -215,6 +220,7 @@ public class FEOpExecutor {
             if (null != ctx.getPrepareExecuteBuffer()) {
                 params.setPrepareExecuteBuffer(ctx.getPrepareExecuteBuffer());
             }
+            params.setCursorFetchRequested(ctx.isCursorFetchRequested());
         }
 
         
ctx.getSessionContext().getDelegatedCredential().ifPresent((DelegatedCredential 
credential) -> {
@@ -231,6 +237,7 @@ public class FEOpExecutor {
         // master on its default packet layout.
         if (ctx.getConnectType() == ConnectType.MYSQL) {
             
params.setClientDeprecatedEOF(ctx.getMysqlChannel().clientDeprecatedEOF());
+            params.setMysqlCapability(ctx.getCapability().getFlags());
         }
 
         return params;
@@ -264,6 +271,76 @@ public class FEOpExecutor {
         return result.packet;
     }
 
+    public boolean isClientDeprecatedEofApplied() {
+        return result != null && result.isSetClientDeprecatedEofApplied()
+                && result.isClientDeprecatedEofApplied();
+    }
+
+    public boolean hasQueryResultPackets() {
+        return result != null && result.isSetQueryResultBufList()
+                && !result.getQueryResultBufList().isEmpty();
+    }
+
+    // An old master cannot add the Connector/J cursor terminator. Normalize 
its buffered
+    // result at the follower, which still has the original execute flag and 
client capability.
+    // DML/DDL OK and ERR packets are retained verbatim, including warnings 
and load info.
+    public void prepareQueryResultForClient() {
+        if (!ctx.getMysqlChannel().clientDeprecatedEOF() || 
isClientDeprecatedEofApplied()
+                || !hasQueryResultPackets()) {
+            return;
+        }
+        List<ByteBuffer> packets = new 
ArrayList<>(result.getQueryResultBufList());
+        int metadataEnd = 
Math.toIntExact(MysqlProto.readVInt(packets.get(0).duplicate())) + 1;
+        boolean needsCursorTerminator = ctx.isCursorFetchRequested()
+                && 
MysqlCursorFetchCompatibility.resolve(ctx.getConnectAttributes())
+                        != MysqlCursorFetchCompatibility.Behavior.STANDARD;
+        // An execution error may occur after only part of the metadata has 
been buffered.
+        if (metadataEnd > packets.size()) {
+            Preconditions.checkState(isErrorPacket(result.packet));
+            return;
+        }
+        boolean hasMetadataTerminator = metadataEnd < packets.size()
+                && isEofPacket(packets.get(metadataEnd));
+        if (hasMetadataTerminator) {
+            ByteBuffer metadata = packets.remove(metadataEnd);
+            if (needsCursorTerminator) {
+                packets.add(metadataEnd, resultSetTerminator(metadata));
+            }
+        } else if (needsCursorTerminator) {
+            MysqlSerializer serializer = 
MysqlSerializer.newInstance(ctx.getCapability());
+            new MysqlResultSetEndPacket(new QueryState()).writeTo(serializer);
+            packets.add(metadataEnd, serializer.toByteBuffer());
+        }
+        result.setQueryResultBufList(packets);
+        if (isEofPacket(result.packet)) {
+            result.setPacket(resultSetTerminator(result.packet));
+        }
+        result.setClientDeprecatedEofApplied(true);
+    }
+
+    private static boolean isErrorPacket(ByteBuffer packet) {
+        return Byte.toUnsignedInt(packet.get(packet.position())) == 0xFF;
+    }
+
+    private static boolean isEofPacket(ByteBuffer packet) {
+        return packet.remaining() <= 8 && 
Byte.toUnsignedInt(packet.get(packet.position())) == 0xFE;
+    }
+
+    private static ByteBuffer resultSetTerminator(ByteBuffer packet) {
+        if (packet.remaining() != 5) {
+            return packet;
+        }
+        ByteBuffer eof = packet.duplicate();
+        MysqlProto.readInt1(eof);
+        int warnings = MysqlProto.readInt2(eof);
+        QueryState state = new QueryState();
+        state.setOk(0, warnings, null);
+        state.serverStatus = MysqlProto.readInt2(eof);
+        MysqlSerializer serializer = MysqlSerializer.newInstance();
+        new MysqlResultSetEndPacket(state).writeTo(serializer);
+        return serializer.toByteBuffer();
+    }
+
     public TUniqueId getQueryId() {
         if (result != null && result.isSetQueryId()) {
             return result.getQueryId();
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/qe/MysqlConnectProcessor.java 
b/fe/fe-core/src/main/java/org/apache/doris/qe/MysqlConnectProcessor.java
index 5191964ed79..46b9d3f4938 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/MysqlConnectProcessor.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/MysqlConnectProcessor.java
@@ -63,6 +63,7 @@ import java.util.Optional;
  */
 public class MysqlConnectProcessor extends ConnectProcessor {
     private static final Logger LOG = 
LogManager.getLogger(MysqlConnectProcessor.class);
+    private static final int CURSOR_TYPE_READ_ONLY = 0x01;
 
     private ByteBuffer packetBuf;
 
@@ -128,13 +129,14 @@ public class MysqlConnectProcessor extends 
ConnectProcessor {
         String stmtStr = "";
         try {
             StatementContext statementContext = prepCtx.getStatementContext();
+            if (!ctx.isProxy()) {
+                // An empty buffer still identifies a zero-parameter 
COM_STMT_EXECUTE when forwarding.
+                ctx.setPrepareExecuteBuffer(packetBuf.duplicate());
+            }
             if (paramCount > 0) {
                 if (LOG.isDebugEnabled()) {
                     LOG.debug("execute param buf: {}, array: {}", packetBuf, 
getHexStr(packetBuf));
                 }
-                if (!ctx.isProxy()) {
-                    ctx.setPrepareExecuteBuffer(packetBuf.duplicate());
-                }
                 byte[] nullbitmapData = new byte[(paramCount + 7) / 8];
                 packetBuf.get(nullbitmapData);
                 // new_params_bind_flag
@@ -208,8 +210,8 @@ public class MysqlConnectProcessor extends ConnectProcessor 
{
         packetBuf = packetBuf.order(ByteOrder.LITTLE_ENDIAN);
         // parse stmt_id, flags, params
         int stmtId = packetBuf.getInt();
-        // flag
-        packetBuf.get();
+        int flags = Byte.toUnsignedInt(packetBuf.get());
+        ctx.setCursorFetchRequested((flags & CURSOR_TYPE_READ_ONLY) != 0);
         // iteration_count always 1,
         packetBuf.getInt();
         if (LOG.isDebugEnabled()) {
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java 
b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
index ae4bc830fd8..e5f726628fd 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
@@ -74,7 +74,9 @@ import org.apache.doris.metric.MetricRepo;
 import org.apache.doris.mysql.FieldInfo;
 import org.apache.doris.mysql.MysqlChannel;
 import org.apache.doris.mysql.MysqlCommand;
+import org.apache.doris.mysql.MysqlCursorFetchCompatibility;
 import org.apache.doris.mysql.MysqlEofPacket;
+import org.apache.doris.mysql.MysqlResultSetEndPacket;
 import org.apache.doris.mysql.MysqlSerializer;
 import org.apache.doris.mysql.ProxyMysqlChannel;
 import org.apache.doris.nereids.NereidsPlanner;
@@ -1869,15 +1871,7 @@ public class StmtExecutor {
             }
             context.getMysqlChannel().sendOnePacket(serializer.toByteBuffer());
         }
-        // When CLIENT_DEPRECATE_EOF is set, the server should not send the 
intermediate
-        // EOF packet after column definitions. The client will go directly 
from column
-        // definitions to reading data rows.
-        if (!context.getMysqlChannel().clientDeprecatedEOF()) {
-            serializer.reset();
-            MysqlEofPacket eofPacket = new MysqlEofPacket(context.getState());
-            eofPacket.writeTo(serializer);
-            context.getMysqlChannel().sendOnePacket(serializer.toByteBuffer());
-        }
+        sendMetadataTerminatorIfNeeded(context.getMysqlChannel());
     }
 
     private List<PrimitiveType> exprToStringType(List<Expr> exprs) {
@@ -2013,17 +2007,30 @@ public class StmtExecutor {
                 channel.sendOnePacket(serializer.toByteBuffer());
             }
         }
-        // When CLIENT_DEPRECATE_EOF is set, the server should not send the 
intermediate
-        // EOF packet after column definitions. The client will go directly 
from column
-        // definitions to reading data rows.
+        sendMetadataTerminatorIfNeeded(channel);
+    }
+
+    private void sendMetadataTerminatorIfNeeded(MysqlChannel channel) throws 
IOException {
         if (!channel.clientDeprecatedEOF()) {
             serializer.reset();
-            MysqlEofPacket eofPacket = new MysqlEofPacket(context.getState());
-            eofPacket.writeTo(serializer);
+            new MysqlEofPacket(context.getState()).writeTo(serializer);
+            channel.sendOnePacket(serializer.toByteBuffer());
+        } else if (connectorJConsumesCursorMetadataTerminator()) {
+            // Connector/J before 9.5 consumes the first OK packet after 
column definitions
+            // while probing whether a requested cursor was created. Doris 
does not create a
+            // cursor, so an empty result would otherwise lose its only end 
marker and block.
+            serializer.reset();
+            new 
MysqlResultSetEndPacket(context.getState()).writeTo(serializer);
             channel.sendOnePacket(serializer.toByteBuffer());
         }
     }
 
+    private boolean connectorJConsumesCursorMetadataTerminator() {
+        return context.isCursorFetchRequested()
+                && 
MysqlCursorFetchCompatibility.resolve(context.getConnectAttributes())
+                        != MysqlCursorFetchCompatibility.Behavior.STANDARD;
+    }
+
     public void sendResultSet(ResultSet resultSet) throws IOException {
         sendResultSet(resultSet, null);
     }
@@ -2579,6 +2586,7 @@ public class StmtExecutor {
         if (masterOpExecutor == null) {
             return;
         }
+        masterOpExecutor.prepareQueryResultForClient();
         List<ByteBuffer> queryResultBufList = 
masterOpExecutor.getQueryResultBufList();
         for (ByteBuffer byteBuffer : queryResultBufList) {
             context.getMysqlChannel().sendOnePacket(byteBuffer);
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlCursorFetchCompatibilityTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlCursorFetchCompatibilityTest.java
new file mode 100644
index 00000000000..740b8a4bf13
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlCursorFetchCompatibilityTest.java
@@ -0,0 +1,59 @@
+// 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.doris.mysql;
+
+import com.google.common.collect.ImmutableMap;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+
+public class MysqlCursorFetchCompatibilityTest {
+    @Test
+    public void testConnectorJBehaviorBoundaries() {
+        
Assertions.assertEquals(MysqlCursorFetchCompatibility.Behavior.CONSUMES_METADATA_TERMINATOR,
+                resolve("MySQL Connector Java", "5.1.49"));
+        
Assertions.assertEquals(MysqlCursorFetchCompatibility.Behavior.CONSUMES_METADATA_TERMINATOR,
+                resolve("MySQL Connector/J", "6.0.6"));
+        
Assertions.assertEquals(MysqlCursorFetchCompatibility.Behavior.CONSUMES_METADATA_TERMINATOR,
+                resolve("MySQL Connector/J", "8.2.0"));
+        
Assertions.assertEquals(MysqlCursorFetchCompatibility.Behavior.CONSUMES_METADATA_TERMINATOR,
+                resolve("MySQL Connector/J", "9.4.0"));
+        
Assertions.assertEquals(MysqlCursorFetchCompatibility.Behavior.STANDARD,
+                resolve("MySQL Connector/J", "9.5.0"));
+        
Assertions.assertEquals(MysqlCursorFetchCompatibility.Behavior.STANDARD,
+                resolve("MySQL Connector/J", "9.6.0"));
+    }
+
+    @Test
+    public void testUnknownAndOtherClients() {
+        Assertions.assertEquals(MysqlCursorFetchCompatibility.Behavior.UNKNOWN,
+                MysqlCursorFetchCompatibility.resolve(Collections.emptyMap()));
+        Assertions.assertEquals(MysqlCursorFetchCompatibility.Behavior.UNKNOWN,
+                
MysqlCursorFetchCompatibility.resolve(ImmutableMap.of("_client_name", "MySQL 
Connector/J")));
+        Assertions.assertEquals(MysqlCursorFetchCompatibility.Behavior.UNKNOWN,
+                resolve("MySQL Connector/J", "custom"));
+        
Assertions.assertEquals(MysqlCursorFetchCompatibility.Behavior.STANDARD,
+                resolve("MariaDB Connector/J", "3.5.6"));
+    }
+
+    private MysqlCursorFetchCompatibility.Behavior resolve(String clientName, 
String clientVersion) {
+        return MysqlCursorFetchCompatibility.resolve(ImmutableMap.of(
+                "_client_name", clientName, "_client_version", clientVersion));
+    }
+}
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlOkPacketTest.java 
b/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlOkPacketTest.java
index 028409ef6af..b5a82a9d17a 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlOkPacketTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlOkPacketTest.java
@@ -56,27 +56,20 @@ public class MysqlOkPacketTest {
         // assert warnings, int2: 0
         Assertions.assertEquals(0x00, MysqlProto.readInt2(buffer));
 
-        // When infoMessage is empty, an empty len-encoded string (0x00) 
should still be written.
-        // This is required because OkPacket.parse() in MySQL Connector/J 
unconditionally reads
-        // STRING_LENENC for info. Without this byte, the driver throws
-        // ArrayIndexOutOfBoundsException when CLIENT_DEPRECATE_EOF is 
negotiated.
-        Assertions.assertEquals(0x00, MysqlProto.readVInt(buffer));
         Assertions.assertEquals(0, buffer.remaining());
     }
 
     @Test
-    public void testWritePayloadSizeGreaterThan5() {
-        // When CLIENT_DEPRECATE_EOF is negotiated, the driver distinguishes 
between
-        // EOF packets (payload <= 5) and ResultSet OK packets (payload > 5).
-        // MysqlOkPacket payload must be > 5 to avoid being misidentified as 
EOF.
-        // Payload: 0x00(1) + affected_rows(1) + last_insert_id(1) + status(2) 
+ warnings(2) + info_len(1) = 8
+    public void testWriteEmptyInfoWithDeprecatedEof() {
+        capability = new 
MysqlCapability(MysqlCapability.Flag.CLIENT_PROTOCOL_41.getFlagBit()
+                | MysqlCapability.Flag.CLIENT_DEPRECATE_EOF.getFlagBit());
         MysqlOkPacket packet = new MysqlOkPacket(new QueryState());
         MysqlSerializer serializer = MysqlSerializer.newInstance(capability);
         packet.writeTo(serializer);
 
         ByteBuffer buffer = serializer.toByteBuffer();
-        int payloadLength = buffer.remaining();
-        Assertions.assertTrue(payloadLength > 5, "OK packet payload should be 
> 5 for CLIENT_DEPRECATE_EOF compatibility, got: "
-                + payloadLength);
+        buffer.position(7);
+        Assertions.assertEquals(0x00, MysqlProto.readVInt(buffer));
+        Assertions.assertEquals(0, buffer.remaining());
     }
 }
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlProtoTest.java 
b/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlProtoTest.java
index a42822fcd58..1405386a047 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlProtoTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlProtoTest.java
@@ -212,6 +212,28 @@ public class MysqlProtoTest {
         Assertions.assertTrue(MysqlProto.negotiate(context));
     }
 
+    @Test
+    public void testNegotiateUsesClientServerCapabilityIntersection() throws 
Exception {
+        mockChannel("user", true);
+        ByteBuffer handshake = channel.fetchOnePacket();
+        int clientFlags = MysqlCapability.DEFAULT_CAPABILITY.getFlags()
+                & ~MysqlCapability.Flag.CLIENT_DEPRECATE_EOF.getFlagBit();
+        handshake.putInt(0, Integer.reverseBytes(clientFlags));
+        mockPassword(true);
+        mockAccess();
+        ConnectContext context = createContext();
+        context.setEnv(env);
+        context.setThreadLocalInfo();
+        Assertions.assertTrue(MysqlProto.negotiate(context));
+        Assertions.assertEquals(clientFlags, 
context.getCapability().getFlags());
+        Assertions.assertEquals(clientFlags, 
channel.getSerializer().getCapability().getFlags());
+        Mockito.verify(channel, Mockito.never()).setClientDeprecatedEOF();
+        MysqlSerializer serializer = channel.getSerializer();
+        serializer.reset();
+        new MysqlOkPacket(context.getState()).writeTo(serializer);
+        Assertions.assertEquals(7, serializer.toByteBuffer().remaining());
+    }
+
     @Test
     public void testNegotiateInitCatalog() throws Exception {
         CatalogMgr catalogMgr = Mockito.mock(CatalogMgr.class);
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlResultSetEndPacketTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlResultSetEndPacketTest.java
index 5462d6aee2c..247b18e18c7 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlResultSetEndPacketTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/mysql/MysqlResultSetEndPacketTest.java
@@ -115,4 +115,30 @@ public class MysqlResultSetEndPacketTest {
         // ResultSet OK payload should be > 5
         Assertions.assertTrue(rsEndPayloadLength > 5, "ResultSet OK packet 
payload should be > 5, got: " + rsEndPayloadLength);
     }
+
+    @Test
+    public void testPreservesMoreResultsStatus() {
+        QueryState state = new QueryState();
+        state.serverStatus = MysqlServerStatusFlag.SERVER_MORE_RESULTS_EXISTS;
+        MysqlSerializer serializer = MysqlSerializer.newInstance(capability);
+        new MysqlResultSetEndPacket(state).writeTo(serializer);
+
+        ByteBuffer buffer = serializer.toByteBuffer();
+        Assertions.assertEquals(0xFE, MysqlProto.readInt1(buffer));
+        Assertions.assertEquals(0, MysqlProto.readVInt(buffer));
+        Assertions.assertEquals(0, MysqlProto.readVInt(buffer));
+        
Assertions.assertEquals(MysqlServerStatusFlag.SERVER_MORE_RESULTS_EXISTS, 
MysqlProto.readInt2(buffer));
+    }
+
+    @Test
+    public void testPreservesWarningCount() {
+        QueryState state = new QueryState();
+        state.setOk(0, 3, null);
+        MysqlSerializer serializer = MysqlSerializer.newInstance(capability);
+        new MysqlResultSetEndPacket(state).writeTo(serializer);
+
+        ByteBuffer buffer = serializer.toByteBuffer();
+        buffer.position(5);
+        Assertions.assertEquals(3, MysqlProto.readInt2(buffer));
+    }
 }
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/MysqlLoadCommandTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/MysqlLoadCommandTest.java
index 153c6d27729..78bddd42c18 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/MysqlLoadCommandTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/MysqlLoadCommandTest.java
@@ -22,11 +22,18 @@ import org.apache.doris.catalog.info.PartitionNamesInfo;
 import org.apache.doris.catalog.info.TableNameInfo;
 import org.apache.doris.common.jmockit.Deencapsulation;
 import org.apache.doris.datasource.InternalCatalog;
+import org.apache.doris.load.EtlJobType;
+import org.apache.doris.load.LoadJobRowResult;
+import org.apache.doris.load.loadv2.LoadManager;
+import org.apache.doris.load.loadv2.MysqlLoadManager;
+import org.apache.doris.mysql.MysqlCapability;
 import org.apache.doris.mysql.privilege.AccessControllerManager;
 import org.apache.doris.mysql.privilege.PrivPredicate;
+import org.apache.doris.nereids.load.NereidsDataDescription;
 import org.apache.doris.nereids.trees.plans.commands.load.MysqlDataDescription;
 import org.apache.doris.nereids.trees.plans.commands.load.MysqlLoadCommand;
 import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.QueryState;
 import org.apache.doris.utframe.TestWithFeService;
 
 import org.junit.jupiter.api.Assertions;
@@ -35,6 +42,7 @@ import org.mockito.Mockito;
 
 import java.io.IOException;
 import java.util.ArrayList;
+import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -54,6 +62,50 @@ public class MysqlLoadCommandTest extends TestWithFeService {
         accessControllerManager = env.getAccessManager();
     }
 
+    @Test
+    public void testLocalFilesCapabilityOnlyRequiredForClientUploads() throws 
Exception {
+        for (boolean clientLocal : new boolean[] {false, true}) {
+            for (boolean supportsLocalFiles : new boolean[] {false, true}) {
+                ConnectContext context = new ConnectContext();
+                Env mockEnv = Mockito.mock(Env.class);
+                LoadManager loadManager = Mockito.mock(LoadManager.class);
+                MysqlLoadManager mysqlLoadManager = 
Mockito.mock(MysqlLoadManager.class);
+                
Mockito.when(mockEnv.getInternalCatalog()).thenReturn(Mockito.mock(InternalCatalog.class));
+                context.setEnv(mockEnv);
+                context.setCapability(new MysqlCapability(supportsLocalFiles
+                        ? MysqlCapability.Flag.CLIENT_LOCAL_FILES.getFlagBit() 
: 0));
+                Mockito.when(mockEnv.getLoadManager()).thenReturn(loadManager);
+                
Mockito.when(loadManager.getMysqlLoadManager()).thenReturn(mysqlLoadManager);
+                MysqlDataDescription description = 
Mockito.mock(MysqlDataDescription.class);
+                
Mockito.when(description.isClientLocal()).thenReturn(clientLocal);
+                
Mockito.when(mysqlLoadManager.executeMySqlLoadJob(Mockito.eq(context),
+                        Mockito.eq(description), 
Mockito.anyString())).thenReturn(new LoadJobRowResult());
+                MysqlLoadCommand command = new MysqlLoadCommand(description, 
new HashMap<>(), "test");
+                Deencapsulation.invoke(command, "handleMysqlLoadCommand", 
context);
+                boolean rejected = clientLocal && !supportsLocalFiles;
+                Assertions.assertEquals(rejected ? 
QueryState.MysqlStateType.ERR : QueryState.MysqlStateType.OK,
+                        context.getState().getStateType());
+                Mockito.verify(mysqlLoadManager, Mockito.times(rejected ? 0 : 
1))
+                        .executeMySqlLoadJob(Mockito.eq(context), 
Mockito.eq(description), Mockito.anyString());
+
+                context.getState().reset();
+                NereidsDataDescription nereidsDescription = 
Mockito.mock(NereidsDataDescription.class);
+                
Mockito.when(nereidsDescription.isClientLocal()).thenReturn(clientLocal);
+                
Mockito.when(mysqlLoadManager.executeMySqlLoadJobFromCommand(Mockito.eq(context),
+                        Mockito.eq(nereidsDescription), 
Mockito.anyString())).thenReturn(new LoadJobRowResult());
+                LoadCommand loadCommand = Mockito.mock(LoadCommand.class, 
Mockito.CALLS_REAL_METHODS);
+                Deencapsulation.setField(loadCommand, "etlJobType", 
EtlJobType.LOCAL_FILE);
+                
Mockito.doReturn(Collections.singletonList(nereidsDescription)).when(loadCommand).getDataDescriptions();
+                loadCommand.handleLoadCommand(context, null);
+                Assertions.assertEquals(rejected ? 
QueryState.MysqlStateType.ERR : QueryState.MysqlStateType.OK,
+                        context.getState().getStateType());
+                Mockito.verify(mysqlLoadManager, Mockito.times(rejected ? 0 : 
1))
+                        .executeMySqlLoadJobFromCommand(Mockito.eq(context), 
Mockito.eq(nereidsDescription),
+                                Mockito.anyString());
+            }
+        }
+    }
+
     @Test
     public void testValidateNormal() throws Exception {
         runBefore();
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/qe/ConnectProcessorForwardProtocolTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/qe/ConnectProcessorForwardProtocolTest.java
new file mode 100644
index 00000000000..554f2595afe
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/qe/ConnectProcessorForwardProtocolTest.java
@@ -0,0 +1,148 @@
+// 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.doris.qe;
+
+import org.apache.doris.mysql.DummyMysqlChannel;
+import org.apache.doris.mysql.MysqlCapability;
+import org.apache.doris.mysql.MysqlProto;
+import org.apache.doris.mysql.MysqlSerializer;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+import java.io.IOException;
+import java.nio.ByteBuffer;
+
+public class ConnectProcessorForwardProtocolTest {
+    @Test
+    public void testOldMasterSuccessfulEofDoesNotBecomeError() throws 
Exception {
+        TestContext context = new TestContext();
+        StmtExecutor executor = forwardedExecutor();
+        Mockito.when(executor.getProxyStatusCode()).thenReturn(1105);
+        ByteBuffer packet = ByteBuffer.wrap(new byte[] {(byte) 0xFE, 0, 0, 2, 
0, 3, 0, 0});
+        Mockito.when(executor.getOutputPacket()).thenReturn(packet);
+
+        new TestProcessor(context, executor).finalizeCommand();
+
+        Assertions.assertEquals(packet, context.channel.packet);
+        Mockito.verify(executor).sendProxyQueryResult();
+    }
+
+    @Test
+    public void testOldMasterDmlPreservesCompleteOk() throws Exception {
+        TestContext context = new TestContext();
+        StmtExecutor executor = forwardedExecutor();
+        QueryState state = new QueryState();
+        state.setOk(7, 3, "label=load_1,txnId=123,status=VISIBLE");
+        state.serverStatus = 2;
+        MysqlSerializer serializer = MysqlSerializer.newInstance();
+        state.toResponsePacket().writeTo(serializer);
+        ByteBuffer packet = serializer.toByteBuffer();
+        Mockito.when(executor.getOutputPacket()).thenReturn(packet);
+
+        new TestProcessor(context, executor).finalizeCommand();
+
+        Assertions.assertEquals(packet, context.channel.packet);
+        Mockito.verify(executor).sendProxyQueryResult();
+    }
+
+    @Test
+    public void testRemoteErrorsRemainUnchanged() throws Exception {
+        TestContext context = new TestContext();
+        StmtExecutor executor = forwardedExecutor();
+        Mockito.when(executor.getProxyStatusCode()).thenReturn(1064);
+        
Mockito.when(executor.getOutputPacket()).thenReturn(ByteBuffer.wrap(new byte[] 
{(byte) 0xFF, 1}));
+
+        new TestProcessor(context, executor).finalizeCommand();
+
+        Assertions.assertEquals(0xFF, 
MysqlProto.readInt1(context.channel.packet));
+        Mockito.verify(executor).sendProxyQueryResult();
+    }
+
+    @Test
+    public void testNewMasterPacketsRemainUnchanged() throws Exception {
+        TestContext context = new TestContext();
+        StmtExecutor executor = forwardedExecutor();
+
+        new TestProcessor(context, executor).finalizeCommand();
+
+        Mockito.verify(executor).sendProxyQueryResult();
+    }
+
+    @Test
+    public void testLegacyEofClientDoesNotRequireConfirmation() throws 
Exception {
+        TestContext context = new TestContext(false);
+        StmtExecutor executor = forwardedExecutor();
+
+        new TestProcessor(context, executor).finalizeCommand();
+
+        Mockito.verify(executor).sendProxyQueryResult();
+    }
+
+    private StmtExecutor forwardedExecutor() {
+        StmtExecutor executor = Mockito.mock(StmtExecutor.class);
+        Mockito.when(executor.hasForwardedToMaster()).thenReturn(true);
+        Mockito.when(executor.getProxyStatusCode()).thenReturn(0);
+        return executor;
+    }
+
+    private static class TestProcessor extends MysqlConnectProcessor {
+        private TestProcessor(ConnectContext context, StmtExecutor executor) {
+            super(context);
+            this.executor = executor;
+        }
+    }
+
+    private static class TestContext extends ConnectContext {
+        private final RecordingChannel channel;
+
+        private TestContext() {
+            this(true);
+        }
+
+        private TestContext(boolean clientDeprecatedEof) {
+            channel = new RecordingChannel(clientDeprecatedEof);
+        }
+
+        @Override
+        public RecordingChannel getMysqlChannel() {
+            return channel;
+        }
+    }
+
+    private static class RecordingChannel extends DummyMysqlChannel {
+        private ByteBuffer packet;
+
+        private RecordingChannel(boolean clientDeprecatedEof) {
+            int flags = MysqlCapability.Flag.CLIENT_PROTOCOL_41.getFlagBit();
+            if (clientDeprecatedEof) {
+                flags |= 
MysqlCapability.Flag.CLIENT_DEPRECATE_EOF.getFlagBit();
+            }
+            serializer = MysqlSerializer.newInstance(new 
MysqlCapability(flags));
+            if (clientDeprecatedEof) {
+                setClientDeprecatedEOF();
+            }
+        }
+
+        @Override
+        public void sendAndFlush(ByteBuffer packet) throws IOException {
+            this.packet = packet.duplicate();
+        }
+    }
+}
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/qe/FEOpExecutorDelegatedCredentialTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/qe/FEOpExecutorDelegatedCredentialTest.java
index 61ed2633576..cf261f04bb1 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/qe/FEOpExecutorDelegatedCredentialTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/qe/FEOpExecutorDelegatedCredentialTest.java
@@ -22,6 +22,7 @@ import org.apache.doris.catalog.Env;
 import org.apache.doris.common.AnalysisException;
 import org.apache.doris.datasource.DelegatedCredential;
 import org.apache.doris.datasource.SessionContext;
+import org.apache.doris.mysql.MysqlCapability;
 import org.apache.doris.system.SystemInfoService;
 import org.apache.doris.thrift.TMasterOpRequest;
 import org.apache.doris.thrift.TNetworkAddress;
@@ -41,6 +42,7 @@ public class FEOpExecutorDelegatedCredentialTest {
         try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
             mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
             ConnectContext context = new ConnectContext();
+            context.setCapability(MysqlCapability.DEFAULT_CAPABILITY);
             
context.setCurrentUserIdentity(UserIdentity.createAnalyzedUserIdentWithIp("alice",
 "%"));
             context.setRemoteIP("127.0.0.1");
             context.setSessionContext(SessionContext.of(new 
DelegatedCredential(
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/qe/FEOpExecutorFlightForwardTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/qe/FEOpExecutorFlightForwardTest.java
index af39b6f936d..8e699d73ed2 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/qe/FEOpExecutorFlightForwardTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/qe/FEOpExecutorFlightForwardTest.java
@@ -21,6 +21,7 @@ import org.apache.doris.analysis.UserIdentity;
 import org.apache.doris.catalog.Env;
 import org.apache.doris.common.AnalysisException;
 import org.apache.doris.common.FeConstants;
+import org.apache.doris.mysql.MysqlCapability;
 import org.apache.doris.service.arrowflight.sessions.FlightSqlConnectContext;
 import org.apache.doris.system.SystemInfoService;
 import org.apache.doris.thrift.TMasterOpRequest;
@@ -73,6 +74,7 @@ public class FEOpExecutorFlightForwardTest {
         try (MockedStatic<Env> mockedEnv = mockSelfNode()) {
             ConnectContext context = new ConnectContext();
             prepare(context);
+            context.setCapability(MysqlCapability.DEFAULT_CAPABILITY);
             context.getMysqlChannel().setClientDeprecatedEOF();
 
             TMasterOpRequest request = new TestFEOpExecutor(context).build();
@@ -87,6 +89,8 @@ public class FEOpExecutorFlightForwardTest {
         try (MockedStatic<Env> mockedEnv = mockSelfNode()) {
             ConnectContext context = new ConnectContext();
             prepare(context);
+            context.setCapability(new 
MysqlCapability(MysqlCapability.DEFAULT_CAPABILITY.getFlags()
+                    & 
~MysqlCapability.Flag.CLIENT_DEPRECATE_EOF.getFlagBit()));
 
             TMasterOpRequest request = new TestFEOpExecutor(context).build();
 
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/qe/FEOpExecutorMysqlProtocolTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/qe/FEOpExecutorMysqlProtocolTest.java
new file mode 100644
index 00000000000..a3e235f56f5
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/qe/FEOpExecutorMysqlProtocolTest.java
@@ -0,0 +1,229 @@
+// 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.doris.qe;
+
+import org.apache.doris.analysis.UserIdentity;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.AnalysisException;
+import org.apache.doris.mysql.MysqlCapability;
+import org.apache.doris.mysql.MysqlCommand;
+import org.apache.doris.mysql.MysqlProto;
+import org.apache.doris.service.arrowflight.sessions.FlightSqlConnectContext;
+import org.apache.doris.system.SystemInfoService;
+import org.apache.doris.thrift.TMasterOpRequest;
+import org.apache.doris.thrift.TMasterOpResult;
+import org.apache.doris.thrift.TNetworkAddress;
+
+import com.google.common.collect.ImmutableMap;
+import org.apache.thrift.TDeserializer;
+import org.apache.thrift.TSerializer;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+import java.nio.ByteBuffer;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+
+public class FEOpExecutorMysqlProtocolTest {
+    @Test
+    public void testForwardRequestCarriesMysqlProtocolContext() throws 
Exception {
+        Env env = Mockito.mock(Env.class);
+        Mockito.when(env.getSelfNode()).thenReturn(new 
SystemInfoService.HostInfo("127.0.0.1", 9010));
+        try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+            mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+            ConnectContext context = createContext();
+            context.getMysqlChannel().setClientDeprecatedEOF();
+            context.setCommand(MysqlCommand.COM_STMT_EXECUTE);
+            context.setCursorFetchRequested(true);
+            context.setConnectAttributes(ImmutableMap.of(
+                    "_client_name", "MySQL Connector/J", "_client_version", 
"8.2.0"));
+
+            TMasterOpRequest request = new TestFEOpExecutor(context).build();
+
+            Assertions.assertTrue(request.isClientDeprecatedEOF());
+            Assertions.assertTrue(request.isCursorFetchRequested());
+            Assertions.assertEquals("8.2.0", 
request.getConnectAttributes().get("_client_version"));
+
+            for (byte[] payload : Arrays.asList(new byte[0], new byte[] {0, 1, 
3, 0, 42, 0, 0, 0})) {
+                context.setPrepareExecuteBuffer(ByteBuffer.wrap(payload));
+                request = new TestFEOpExecutor(context).build();
+                TMasterOpRequest restored = new TMasterOpRequest();
+                new TDeserializer().deserialize(restored, new 
TSerializer().serialize(request));
+                // Presence, including a zero-length payload, selects 
COM_STMT_EXECUTE on the master.
+                Assertions.assertTrue(restored.isSetPrepareExecuteBuffer());
+                Assertions.assertArrayEquals(payload, 
restored.getPrepareExecuteBuffer());
+                Assertions.assertTrue(restored.isCursorFetchRequested());
+            }
+        }
+    }
+
+    @Test
+    public void testForwardResponseRequiresExplicitProtocolConfirmation() {
+        TestFEOpExecutor executor = new TestFEOpExecutor(createContext());
+        executor.setResult(new TMasterOpResult());
+        Assertions.assertFalse(executor.isClientDeprecatedEofApplied());
+        Assertions.assertFalse(executor.hasQueryResultPackets());
+
+        TMasterOpResult confirmed = new TMasterOpResult();
+        confirmed.setClientDeprecatedEofApplied(true);
+        
confirmed.setQueryResultBufList(Collections.singletonList(ByteBuffer.wrap(new 
byte[] {1})));
+        confirmed.setAffectedRows(7);
+        executor.setResult(confirmed);
+        Assertions.assertTrue(executor.isClientDeprecatedEofApplied());
+        Assertions.assertTrue(executor.hasQueryResultPackets());
+    }
+
+    @Test
+    public void testArrowForwardRequestDoesNotAccessMysqlChannel() throws 
Exception {
+        Env env = Mockito.mock(Env.class);
+        Mockito.when(env.getSelfNode()).thenReturn(new 
SystemInfoService.HostInfo("127.0.0.1", 9010));
+        try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+            mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+            FlightSqlConnectContext context = new 
FlightSqlConnectContext("alice");
+            
context.setCurrentUserIdentity(UserIdentity.createAnalyzedUserIdentWithIp("alice",
 "%"));
+            context.setRemoteIP("127.0.0.1");
+            TMasterOpRequest request = new TestFEOpExecutor(context).build();
+            Assertions.assertFalse(request.isSetClientDeprecatedEOF());
+            Assertions.assertFalse(request.isSetMysqlCapability());
+        }
+    }
+
+    @Test
+    public void testForwardedCapabilityAndMissingCursorFlag() throws Exception 
{
+        int legacyFlags = MysqlCapability.DEFAULT_CAPABILITY.getFlags()
+                & ~MysqlCapability.Flag.CLIENT_DEPRECATE_EOF.getFlagBit();
+        TMasterOpRequest request = new TMasterOpRequest();
+        request.setMysqlCapability(legacyFlags);
+        ConnectContext context = createContext();
+        ConnectProcessor.restoreForwardedMysqlContext(context, request);
+        Assertions.assertEquals(legacyFlags, 
context.getCapability().getFlags());
+        
Assertions.assertFalse(context.getMysqlChannel().getSerializer().getCapability().isDeprecatedEOF());
+
+        request = new TMasterOpRequest();
+        request.setClientDeprecatedEOF(true);
+        request.setPrepareExecuteBuffer(new byte[] {0});
+        context = createContext();
+        context.setConnectAttributes(ImmutableMap.of("_client_name", "MySQL 
Connector/J", "_client_version", "8.2.0"));
+        ConnectProcessor.restoreForwardedMysqlContext(context, request);
+        Assertions.assertFalse(context.isCursorFetchRequested());
+        Assertions.assertTrue(context.getCapability().isDeprecatedEOF());
+        request.setCursorFetchRequested(true);
+        ConnectProcessor.restoreForwardedMysqlContext(context, request);
+        Assertions.assertTrue(context.isCursorFetchRequested());
+        request.setCursorFetchRequested(false);
+        ConnectProcessor.restoreForwardedMysqlContext(context, request);
+        Assertions.assertFalse(context.isCursorFetchRequested());
+    }
+
+    @Test
+    public void testOldMasterPacketMatrix() {
+        for (boolean legacyMaster : new boolean[] {false, true}) {
+            for (boolean cursor : new boolean[] {false, true}) {
+                for (String version : new String[] {"8.2.0", "9.4.0", 
"9.5.0"}) {
+                    for (boolean rows : new boolean[] {false, true}) {
+                        ConnectContext context = createContext();
+                        context.getMysqlChannel().setClientDeprecatedEOF();
+                        context.setCursorFetchRequested(cursor);
+                        context.setConnectAttributes(ImmutableMap.of(
+                                "_client_name", "MySQL Connector/J", 
"_client_version", version));
+                        TestFEOpExecutor executor = new 
TestFEOpExecutor(context);
+                        List<ByteBuffer> packets = new ArrayList<>();
+                        packets.add(ByteBuffer.wrap(new byte[] {1})); // 
column count
+                        packets.add(ByteBuffer.wrap(new byte[] {3, 'd', 'e', 
'f'})); // opaque column definition
+                        if (legacyMaster) {
+                            packets.add(ByteBuffer.wrap(new byte[] {(byte) 
0xFE, 3, 0, 2, 0}));
+                        }
+                        ByteBuffer row = ByteBuffer.wrap(new byte[] {0, 0, 
42});
+                        if (rows) {
+                            packets.add(row);
+                        }
+                        TMasterOpResult result = new TMasterOpResult();
+                        result.setQueryResultBufList(packets);
+                        result.setStatus("EOF");
+                        result.setStatusCode(1105); // production successful 
SELECT mapping
+                        result.setPacket(ByteBuffer.wrap(legacyMaster
+                                ? new byte[] {(byte) 0xFE, 3, 0, 2, 0}
+                                : new byte[] {(byte) 0xFE, 0, 0, 2, 0, 3, 0, 
0}));
+                        executor.setResult(result);
+                        executor.prepareQueryResultForClient();
+                        boolean shim = cursor && !version.equals("9.5.0");
+                        Assertions.assertEquals(2 + (shim ? 1 : 0) + (rows ? 1 
: 0),
+                                executor.getQueryResultBufList().size());
+                        if (shim) {
+                            Assertions.assertEquals(8, 
executor.getQueryResultBufList().get(2).remaining());
+                        }
+                        if (rows) {
+                            Assertions.assertSame(row, 
executor.getQueryResultBufList().get(shim ? 3 : 2));
+                        }
+                        ByteBuffer end = 
executor.getOutputPacket().duplicate();
+                        Assertions.assertEquals(0xFE, 
MysqlProto.readInt1(end));
+                        Assertions.assertEquals(0, MysqlProto.readVInt(end));
+                        Assertions.assertEquals(0, MysqlProto.readVInt(end));
+                        Assertions.assertEquals(2, MysqlProto.readInt2(end));
+                        Assertions.assertEquals(3, MysqlProto.readInt2(end));
+                        List<ByteBuffer> normalized = 
executor.getQueryResultBufList();
+                        executor.prepareQueryResultForClient();
+                        Assertions.assertSame(normalized, 
executor.getQueryResultBufList());
+                    }
+                }
+            }
+        }
+    }
+
+    @Test
+    public void testOkAndErrorAreNotRebuilt() {
+        ConnectContext context = createContext();
+        context.getMysqlChannel().setClientDeprecatedEOF();
+        TestFEOpExecutor executor = new TestFEOpExecutor(context);
+        for (byte[] bytes : Arrays.asList(new byte[] {0, 7, 0, 2, 0, 3, 0, 4, 
'i', 'n', 'f', 'o'},
+                new byte[] {(byte) 0xFF, 1, 2})) {
+            TMasterOpResult result = new TMasterOpResult();
+            ByteBuffer packet = ByteBuffer.wrap(bytes);
+            result.setPacket(packet);
+            executor.setResult(result);
+            executor.prepareQueryResultForClient();
+            Assertions.assertEquals(packet, executor.getOutputPacket());
+        }
+    }
+
+    private ConnectContext createContext() {
+        ConnectContext context = new ConnectContext();
+        
context.setCurrentUserIdentity(UserIdentity.createAnalyzedUserIdentWithIp("alice",
 "%"));
+        context.setRemoteIP("127.0.0.1");
+        context.setCapability(MysqlCapability.DEFAULT_CAPABILITY);
+        return context;
+    }
+
+    private static class TestFEOpExecutor extends FEOpExecutor {
+        private TestFEOpExecutor(ConnectContext context) {
+            super(new TNetworkAddress("127.0.0.1", 9010), new 
OriginStatement("select 1", 0), context, true);
+        }
+
+        private TMasterOpRequest build() throws AnalysisException {
+            return buildStmtForwardParams();
+        }
+
+        private void setResult(TMasterOpResult result) {
+            this.result = result;
+        }
+    }
+}
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/qe/MysqlConnectProcessorCursorFetchTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/qe/MysqlConnectProcessorCursorFetchTest.java
new file mode 100644
index 00000000000..8f788d75ad4
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/qe/MysqlConnectProcessorCursorFetchTest.java
@@ -0,0 +1,112 @@
+// 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.doris.qe;
+
+import org.apache.doris.mysql.MysqlCommand;
+import org.apache.doris.nereids.StatementContext;
+import org.apache.doris.nereids.trees.plans.commands.PrepareCommand;
+
+import com.google.common.collect.ImmutableMap;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedConstruction;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+import java.lang.reflect.Field;
+import java.lang.reflect.Method;
+import java.nio.ByteBuffer;
+import java.nio.ByteOrder;
+
+public class MysqlConnectProcessorCursorFetchTest {
+    private static final int CURSOR_TYPE_READ_ONLY = 1;
+
+    @Test
+    public void testZeroParameterExecutePreservesForwardingBuffer() throws 
Exception {
+        for (boolean proxy : new boolean[] {false, true}) {
+            ConnectContext context = new ConnectContext(null, proxy);
+            context.setCommand(MysqlCommand.COM_STMT_EXECUTE);
+            PrepareCommand command = Mockito.mock(PrepareCommand.class);
+            Mockito.when(command.getOriginalStmt()).thenReturn(new 
OriginStatement("select 1", 0));
+            StatementContext statementContext = new StatementContext();
+            PreparedStatementContext prepared = new PreparedStatementContext(
+                    command, context, statementContext, "select 1");
+            ByteBuffer packet = ByteBuffer.allocate(9);
+            packet.position(packet.limit()); // COM_STMT_EXECUTE header 
consumed, no parameter payload.
+            try (MockedConstruction<StmtExecutor> executors = 
Mockito.mockConstruction(StmtExecutor.class);
+                    MockedStatic<AuditLogHelper> audit = 
Mockito.mockStatic(AuditLogHelper.class)) {
+                new MysqlConnectProcessor(context).handleExecute(command, 7, 
prepared, packet, null);
+                Assertions.assertEquals(1, executors.constructed().size());
+                Mockito.verify(executors.constructed().get(0)).execute();
+                if (proxy) {
+                    Assertions.assertNull(context.getPrepareExecuteBuffer());
+                } else {
+                    
Assertions.assertNotNull(context.getPrepareExecuteBuffer());
+                    Assertions.assertNotSame(packet, 
context.getPrepareExecuteBuffer());
+                    Assertions.assertEquals(0, 
context.getPrepareExecuteBuffer().remaining());
+                    Assertions.assertEquals(9, packet.position());
+                }
+            }
+        }
+    }
+
+    @Test
+    public void 
testUnidentifiedDeprecatedEofCursorReachesPreparedStatementLookup() throws 
Exception {
+        ConnectContext context = execute(true, true, false);
+        Assertions.assertTrue(context.getState().getErrorMessage().contains(
+                "Unknown prepared statement handler"));
+    }
+
+    @Test
+    public void testCompatibilityGateOnlyAppliesToAmbiguousProtocol() throws 
Exception {
+        Assertions.assertTrue(execute(false, true, 
false).getState().getErrorMessage().contains(
+                "Unknown prepared statement handler"));
+        Assertions.assertTrue(execute(true, false, 
false).getState().getErrorMessage().contains(
+                "Unknown prepared statement handler"));
+        Assertions.assertTrue(execute(true, true, 
true).getState().getErrorMessage().contains(
+                "Unknown prepared statement handler"));
+    }
+
+    private ConnectContext execute(boolean cursorRequested, boolean 
clientDeprecatedEof,
+            boolean identifiedClient) throws Exception {
+        ConnectContext context = new ConnectContext();
+        context.setCommand(MysqlCommand.COM_STMT_EXECUTE);
+        if (clientDeprecatedEof) {
+            context.getMysqlChannel().setClientDeprecatedEOF();
+        }
+        if (identifiedClient) {
+            context.setConnectAttributes(ImmutableMap.of(
+                    "_client_name", "MySQL Connector/J", "_client_version", 
"8.2.0"));
+        }
+
+        ByteBuffer packet = 
ByteBuffer.allocate(9).order(ByteOrder.LITTLE_ENDIAN);
+        packet.putInt(7);
+        packet.put((byte) (cursorRequested ? CURSOR_TYPE_READ_ONLY : 0));
+        packet.putInt(1);
+        packet.flip();
+
+        MysqlConnectProcessor processor = new MysqlConnectProcessor(context);
+        Field packetField = 
MysqlConnectProcessor.class.getDeclaredField("packetBuf");
+        packetField.setAccessible(true);
+        packetField.set(processor, packet);
+        Method handleExecute = 
MysqlConnectProcessor.class.getDeclaredMethod("handleExecute");
+        handleExecute.setAccessible(true);
+        handleExecute.invoke(processor);
+        return context;
+    }
+}
diff --git a/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorTest.java 
b/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorTest.java
index 990b1c82848..fafcc0d8bef 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorTest.java
@@ -38,6 +38,7 @@ import org.apache.doris.thrift.TQueryOptions;
 import org.apache.doris.thrift.TUniqueId;
 import org.apache.doris.utframe.TestWithFeService;
 
+import com.google.common.collect.ImmutableMap;
 import com.google.common.collect.Lists;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
@@ -50,7 +51,10 @@ import java.lang.reflect.Field;
 import java.lang.reflect.Method;
 import java.nio.ByteBuffer;
 import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Collections;
 import java.util.List;
+import java.util.Map;
 import java.util.concurrent.atomic.AtomicInteger;
 
 public class StmtExecutorTest extends TestWithFeService {
@@ -398,6 +402,100 @@ public class StmtExecutorTest extends TestWithFeService {
         executor.sendBinaryResultRow(resultSet);
     }
 
+    @Test
+    public void testCursorFetchMetadataTerminatorDependsOnConnectorJVersion() 
throws IOException {
+        List<byte[]> connector82Packets = sendEmptyResultSet(true, "MySQL 
Connector/J", "8.2.0");
+        Assertions.assertEquals(3, connector82Packets.size());
+        Assertions.assertEquals(0xFE, 
Byte.toUnsignedInt(connector82Packets.get(2)[0]));
+        Assertions.assertTrue(connector82Packets.get(2).length > 5);
+
+        Assertions.assertEquals(3, sendEmptyResultSet(true, "MySQL Connector 
Java", "5.1.49").size());
+        Assertions.assertEquals(3, sendEmptyResultSet(true, "MySQL 
Connector/J", "6.0.6").size());
+        Assertions.assertEquals(3, sendEmptyResultSet(true, "MySQL 
Connector/J", "9.4.0").size());
+        Assertions.assertEquals(2, sendEmptyResultSet(true, "MySQL 
Connector/J", "9.5.0").size());
+        Assertions.assertEquals(2, sendEmptyResultSet(false, "MySQL 
Connector/J", "8.2.0").size());
+        Assertions.assertEquals(2, sendEmptyResultSet(true, "MariaDB 
Connector/J", "3.5.6").size());
+        Assertions.assertEquals(3, sendEmptyResultSet(true, 
Collections.emptyMap()).size());
+
+        List<byte[]> legacyEofPackets = sendEmptyResultSet(true, "MySQL 
Connector/J", "8.2.0", false);
+        Assertions.assertEquals(3, legacyEofPackets.size());
+        Assertions.assertEquals(5, legacyEofPackets.get(2).length);
+    }
+
+    @Test
+    public void testPrepareMetadataTerminatorsFollowNegotiatedCapability() 
throws IOException {
+        Assertions.assertEquals(3, sendPrepareMetadata(false).size());
+        Assertions.assertEquals(2, sendPrepareMetadata(true).size());
+    }
+
+    private List<byte[]> sendPrepareMetadata(boolean clientDeprecatedEof) 
throws IOException {
+        ConnectContext mockCtx = Mockito.mock(ConnectContext.class);
+        MysqlChannel channel = Mockito.mock(MysqlChannel.class);
+        Mockito.when(mockCtx.getConnectType()).thenReturn(ConnectType.MYSQL);
+        Mockito.when(mockCtx.getMysqlChannel()).thenReturn(channel);
+        Mockito.when(mockCtx.getState()).thenReturn(new QueryState());
+        Mockito.when(mockCtx.getSessionVariable()).thenReturn(new 
SessionVariable());
+        
Mockito.when(channel.clientDeprecatedEOF()).thenReturn(clientDeprecatedEof);
+        
Mockito.when(channel.getSerializer()).thenReturn(MysqlSerializer.newInstance());
+
+        List<byte[]> packets = new ArrayList<>();
+        Mockito.doAnswer(invocation -> {
+            ByteBuffer packet = invocation.getArgument(0);
+            byte[] copy = new byte[packet.remaining()];
+            packet.duplicate().get(copy);
+            packets.add(copy);
+            return null;
+        }).when(channel).sendOnePacket(Mockito.any(ByteBuffer.class));
+
+        new StmtExecutor(mockCtx, new OriginStatement("", 0), 
true).sendStmtPrepareOK(
+                1, Collections.singletonList("p"), Collections.emptyList());
+        return packets;
+    }
+
+    private List<byte[]> sendEmptyResultSet(boolean cursorFetchRequested, 
String clientName,
+            String clientVersion) throws IOException {
+        return sendEmptyResultSet(cursorFetchRequested, clientName, 
clientVersion, true);
+    }
+
+    private List<byte[]> sendEmptyResultSet(boolean cursorFetchRequested, 
String clientName,
+            String clientVersion, boolean clientDeprecatedEof) throws 
IOException {
+        return sendEmptyResultSet(cursorFetchRequested, ImmutableMap.of(
+                "_client_name", clientName, "_client_version", clientVersion), 
clientDeprecatedEof);
+    }
+
+    private List<byte[]> sendEmptyResultSet(boolean cursorFetchRequested,
+            Map<String, String> connectAttributes) throws IOException {
+        return sendEmptyResultSet(cursorFetchRequested, connectAttributes, 
true);
+    }
+
+    private List<byte[]> sendEmptyResultSet(boolean cursorFetchRequested,
+            Map<String, String> connectAttributes, boolean 
clientDeprecatedEof) throws IOException {
+        ConnectContext mockCtx = Mockito.mock(ConnectContext.class);
+        MysqlChannel channel = Mockito.mock(MysqlChannel.class);
+        Mockito.when(mockCtx.getConnectType()).thenReturn(ConnectType.MYSQL);
+        Mockito.when(mockCtx.getMysqlChannel()).thenReturn(channel);
+        Mockito.when(mockCtx.getState()).thenReturn(new QueryState());
+        
Mockito.when(mockCtx.getSessionVariable()).thenReturn(VariableMgr.newSessionVariable());
+        
Mockito.when(mockCtx.isCursorFetchRequested()).thenReturn(cursorFetchRequested);
+        
Mockito.when(mockCtx.getConnectAttributes()).thenReturn(connectAttributes);
+        
Mockito.when(channel.clientDeprecatedEOF()).thenReturn(clientDeprecatedEof);
+        
Mockito.when(channel.getSerializer()).thenReturn(MysqlSerializer.newInstance());
+
+        List<byte[]> packets = new ArrayList<>();
+        Mockito.doAnswer(invocation -> {
+            ByteBuffer packet = invocation.getArgument(0);
+            byte[] copy = new byte[packet.remaining()];
+            packet.duplicate().get(copy);
+            packets.add(copy);
+            return null;
+        }).when(channel).sendOnePacket(Mockito.any(ByteBuffer.class));
+
+        List<Column> columns = Collections.singletonList(new Column("c", 
PrimitiveType.INT));
+        ResultSet resultSet = new CommonResultSet(new 
CommonResultSetMetaData(columns), Collections.emptyList());
+        new StmtExecutor(mockCtx, new OriginStatement("", 0), 
true).sendResultSet(resultSet);
+        return packets;
+    }
+
     @Test
     public void testSendBinaryBooleanResultRow() throws IOException {
         ConnectContext mockCtx = Mockito.mock(ConnectContext.class);
diff --git a/gensrc/thrift/FrontendService.thrift 
b/gensrc/thrift/FrontendService.thrift
index 18535e1c7cc..1829181a875 100644
--- a/gensrc/thrift/FrontendService.thrift
+++ b/gensrc/thrift/FrontendService.thrift
@@ -441,6 +441,10 @@ struct TMasterOpRequest {
     1005: optional string delegated_credential_token
     1006: optional i64 delegated_credential_expires_at_millis
     1007: optional string delegated_credential_session_id
+    // Whether COM_STMT_EXECUTE requested CURSOR_TYPE_READ_ONLY.
+    1008: optional bool cursor_fetch_requested
+    // Capabilities negotiated with the original MySQL client.
+    1009: optional i32 mysql_capability
 }
 
 struct TColumnDefinition {
@@ -474,6 +478,8 @@ struct TMasterOpResult {
     11: optional i64 affectedRows;
     // Lets the forwarding FE wait for the final statistics of external write 
fragments.
     12: optional list<i64> auditStatisticsBackendIds;
+    // Confirms that the executing FE serialized raw MySQL packets with 
CLIENT_DEPRECATE_EOF.
+    13: optional bool clientDeprecatedEofApplied;
 }
 
 // Certificate-based authentication info forwarded from BE to FE
diff --git 
a/regression-test/data/prepared_stmt_p0/cursor_fetch_empty_result.out 
b/regression-test/data/prepared_stmt_p0/cursor_fetch_empty_result.out
new file mode 100644
index 00000000000..983864abdaf
--- /dev/null
+++ b/regression-test/data/prepared_stmt_p0/cursor_fetch_empty_result.out
@@ -0,0 +1,11 @@
+-- This file is automatically generated. You should know what you did if you 
want to edit this
+-- !empty_result --
+
+-- !non_empty_result --
+1
+
+-- !anonymous_empty --
+
+-- !anonymous_non_empty --
+1
+
diff --git a/regression-test/suites/arrow_flight_sql_p0/test_ddl.groovy 
b/regression-test/suites/arrow_flight_sql_p0/test_ddl.groovy
new file mode 100644
index 00000000000..c9dfa13cb3c
--- /dev/null
+++ b/regression-test/suites/arrow_flight_sql_p0/test_ddl.groovy
@@ -0,0 +1,25 @@
+// 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.
+
+suite("test_ddl", "arrow_flight_sql") {
+    sql "DROP DATABASE IF EXISTS test_arrow_flight_sql_ddl"
+
+    context.getArrowFlightSqlConnection().createStatement().withCloseable { 
statement ->
+        statement.execute("CREATE DATABASE test_arrow_flight_sql_ddl")
+        statement.execute("DROP DATABASE test_arrow_flight_sql_ddl")
+    }
+}
diff --git 
a/regression-test/suites/prepared_stmt_p0/cursor_fetch_empty_result.groovy 
b/regression-test/suites/prepared_stmt_p0/cursor_fetch_empty_result.groovy
new file mode 100644
index 00000000000..82e37dffb99
--- /dev/null
+++ b/regression-test/suites/prepared_stmt_p0/cursor_fetch_empty_result.groovy
@@ -0,0 +1,107 @@
+// 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.
+
+suite("cursor_fetch_empty_result") {
+    String url = getServerPrepareJdbcUrl(context.config.jdbcUrl, 
"regression_test_prepared_stmt_p0") +
+            
"&useCursorFetch=true&defaultFetchSize=10000&socketTimeout=10000&emulateUnsupportedPstmts=false"
+
+    connect(context.config.jdbcUser, context.config.jdbcPassword, url) {
+        // With a positive defaultFetchSize Connector/J also converts a plain 
Statement into a
+        // server-prepared cursor execution, which is how BI tools commonly 
enter this path.
+        context.getConnection().createStatement().withCloseable { statement ->
+            ["SELECT 1 AS c WHERE 1 = 2", "SELECT 1 AS c WHERE 1 = 1"].each { 
query ->
+                statement.executeQuery(query).withCloseable { result ->
+                    while (result.next()) {
+                        result.getInt(1)
+                    }
+                }
+            }
+        }
+
+        def emptyResult = prepareStatement "SELECT 1 AS c WHERE 1 = 2"
+        assertEquals(com.mysql.cj.jdbc.ServerPreparedStatement, 
emptyResult.class)
+        qe_empty_result emptyResult
+        emptyResult.close()
+
+        def nonEmptyResult = prepareStatement "SELECT 1 AS c WHERE 1 = 1"
+        assertEquals(com.mysql.cj.jdbc.ServerPreparedStatement, 
nonEmptyResult.class)
+        qe_non_empty_result nonEmptyResult
+        nonEmptyResult.close()
+    }
+
+    String unidentifiedClientUrl = getServerPrepareJdbcUrl(
+            context.config.jdbcUrl, "regression_test_prepared_stmt_p0") +
+            
"&useCursorFetch=true&defaultFetchSize=10000&connectionAttributes=none&socketTimeout=10000"
 +
+            "&emulateUnsupportedPstmts=false"
+    connect(context.config.jdbcUser, context.config.jdbcPassword, 
unidentifiedClientUrl) {
+        qt_anonymous_empty "SELECT 1 AS c WHERE 1 = 2"
+        qt_anonymous_non_empty "SELECT 1 AS c WHERE 1 = 1"
+    }
+
+    def followers = sql_return_maparray("SHOW FRONTENDS").findAll {
+        it.IsMaster == "false" && it.Alive == "true"
+    }
+    if (followers.isEmpty()) {
+        logger.info("Skip prepared forwarding coverage: no live non-master FE")
+    } else {
+        sql "DROP TABLE IF EXISTS cursor_fetch_forwarding"
+        sql """CREATE TABLE cursor_fetch_forwarding (k INT)
+               DUPLICATE KEY(k) DISTRIBUTED BY HASH(k) BUCKETS 1
+               PROPERTIES ("replication_num" = "1")"""
+        sql "INSERT INTO cursor_fetch_forwarding VALUES (1), (2), (3)"
+        followers.each { fe ->
+            String followerUrl = getServerPrepareJdbcUrl(
+                    "jdbc:mysql://${fe.Host}:${fe.QueryPort}/", 
"regression_test_prepared_stmt_p0", false) +
+                    
"&useCursorFetch=true&emulateUnsupportedPstmts=false&socketTimeout=10000"
+            connect(context.config.jdbcUser, context.config.jdbcPassword, 
followerUrl) {
+                def connection = context.getConnection()
+                connection.createStatement().withCloseable { control ->
+                    // SET must use COM_QUERY: enabling forwarding before 
PREPARE rejects server prepare.
+                    control.execute("SET force_forward_all_queries=false")
+                    control.execute("SYNC")
+                    [0, 1, 10000].each { fetchSize ->
+                        ["SELECT k FROM cursor_fetch_forwarding ORDER BY k",
+                         "SELECT k FROM cursor_fetch_forwarding WHERE k < 0 
ORDER BY k"].each { query ->
+                            connection.prepareStatement(query).withCloseable { 
prepared ->
+                                
assertEquals(com.mysql.cj.jdbc.ServerPreparedStatement, prepared.class)
+                                prepared.setFetchSize(fetchSize)
+                                def readRows = {
+                                    def rows = []
+                                    prepared.executeQuery().withCloseable { 
result ->
+                                        while (result.next()) {
+                                            rows.add(result.getInt(1))
+                                        }
+                                    }
+                                    return rows
+                                }
+                                def directRows = readRows()
+                                control.execute("SET 
force_forward_all_queries=true")
+                                try {
+                                    // Compare execution modes using the same 
server-prepared statement.
+                                    3.times { assertEquals(directRows, 
readRows()) }
+                                } finally {
+                                    control.execute("SET 
force_forward_all_queries=false")
+                                }
+                            }
+                        }
+                    }
+                }
+            }
+        }
+    }
+
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to