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

funky-eyes pushed a commit to branch 2.x
in repository https://gitbox.apache.org/repos/asf/incubator-seata.git


The following commit(s) were added to refs/heads/2.x by this push:
     new 225d182d6a bugfix: use explicit columns in rollback validation query 
(#8118)
225d182d6a is described below

commit 225d182d6a55691d6b4695340131bd39c69bca5f
Author: somil jain <[email protected]>
AuthorDate: Mon Jun 1 06:57:49 2026 +0530

    bugfix: use explicit columns in rollback validation query (#8118)
---
 changes/en-us/2.x.md                               |   1 +
 .../rm/datasource/undo/AbstractUndoExecutor.java   |  23 +++-
 .../undo/sqlserver/BaseSqlServerUndoExecutor.java  |   5 +
 ...ExecutorTest.java => BaseUndoExecutorTest.java} | 120 ++++++++++++++++++++-
 ...est.java => SqlServerBaseUndoExecutorTest.java} |  21 +++-
 5 files changed, 165 insertions(+), 5 deletions(-)

diff --git a/changes/en-us/2.x.md b/changes/en-us/2.x.md
index 4db219c270..6173486a13 100644
--- a/changes/en-us/2.x.md
+++ b/changes/en-us/2.x.md
@@ -49,6 +49,7 @@ Add changes here for all PR submitted to the 2.x branch.
 - [[#8078](https://github.com/apache/incubator-seata/pull/8078)] fix MySQL 
undo log serialization exception
 - [[#8106](https://github.com/apache/incubator-seata/pull/8106)] fix NPE 
during AOT proxy creation
 - [[#8113](https://github.com/apache/incubator-seata/pull/8113)] fix console 
export JSON consistency and download issues
+- [[#8118](https://github.com/apache/incubator-seata/pull/8118)] Use explicit 
columns in rollback validation query
 
 ### optimize:
 
diff --git 
a/rm-datasource/src/main/java/org/apache/seata/rm/datasource/undo/AbstractUndoExecutor.java
 
b/rm-datasource/src/main/java/org/apache/seata/rm/datasource/undo/AbstractUndoExecutor.java
index 448d89a8b0..a455c4e8ea 100644
--- 
a/rm-datasource/src/main/java/org/apache/seata/rm/datasource/undo/AbstractUndoExecutor.java
+++ 
b/rm-datasource/src/main/java/org/apache/seata/rm/datasource/undo/AbstractUndoExecutor.java
@@ -67,7 +67,7 @@ public abstract class AbstractUndoExecutor {
      * template of check sql
      * TODO support multiple primary key
      */
-    private static final String CHECK_SQL_TEMPLATE = "SELECT * FROM %s WHERE 
%s FOR UPDATE";
+    private static final String CHECK_SQL_TEMPLATE = "SELECT %s FROM %s WHERE 
%s FOR UPDATE";
 
     /**
      * Switch of undo data validation
@@ -306,10 +306,15 @@ public abstract class AbstractUndoExecutor {
         int pkRowSize = pkRowValues.get(firstKey).size();
         List<SqlGenerateUtils.WhereSql> sqlConditions =
                 SqlGenerateUtils.buildWhereConditionListByPKs(pkNameList, 
pkRowSize, connectionProxy.getDbType());
+
+        String selectColumns = 
undoRecords.getRows().get(0).getFields().stream()
+                .map(field -> ColumnUtils.addEscape(field.getName(), 
connectionProxy.getDbType()))
+                .collect(Collectors.joining(", "));
+
         TableRecords currentRecords = new TableRecords(tableMeta);
         int totalRowIndex = 0;
         for (SqlGenerateUtils.WhereSql sqlCondition : sqlConditions) {
-            String checkSQL = buildCheckSql(sqlUndoLog.getTableName(), 
sqlCondition.getSql());
+            String checkSQL = buildCheckSql(sqlUndoLog.getTableName(), 
sqlCondition.getSql(), selectColumns);
             PreparedStatement statement = null;
             ResultSet checkSet = null;
             try {
@@ -345,7 +350,19 @@ public abstract class AbstractUndoExecutor {
      * @return the check sql for query current records
      */
     protected String buildCheckSql(String tableName, String whereCondition) {
-        return String.format(CHECK_SQL_TEMPLATE, tableName, whereCondition);
+        return String.format(CHECK_SQL_TEMPLATE, "*", tableName, 
whereCondition);
+    }
+
+    /**
+     * build sql for query current records with explicit columns.
+     *
+     * @param tableName     the tableName to query
+     * @param whereCondition the where condition
+     * @param selectColumns the explicit columns to select
+     * @return the check sql for query current records
+     */
+    protected String buildCheckSql(String tableName, String whereCondition, 
String selectColumns) {
+        return String.format(CHECK_SQL_TEMPLATE, selectColumns, tableName, 
whereCondition);
     }
 
     protected List<Field> getOrderedPkList(TableRecords image, Row row, String 
dbType) {
diff --git 
a/rm-datasource/src/main/java/org/apache/seata/rm/datasource/undo/sqlserver/BaseSqlServerUndoExecutor.java
 
b/rm-datasource/src/main/java/org/apache/seata/rm/datasource/undo/sqlserver/BaseSqlServerUndoExecutor.java
index ae60976628..cfae106e5b 100644
--- 
a/rm-datasource/src/main/java/org/apache/seata/rm/datasource/undo/sqlserver/BaseSqlServerUndoExecutor.java
+++ 
b/rm-datasource/src/main/java/org/apache/seata/rm/datasource/undo/sqlserver/BaseSqlServerUndoExecutor.java
@@ -36,4 +36,9 @@ public abstract class BaseSqlServerUndoExecutor extends 
AbstractUndoExecutor {
     protected String buildCheckSql(String tableName, String whereCondition) {
         return "SELECT * FROM " + tableName + " WITH(UPDLOCK) WHERE " + 
whereCondition;
     }
+
+    @Override
+    protected String buildCheckSql(String tableName, String whereCondition, 
String selectColumns) {
+        return "SELECT " + selectColumns + " FROM " + tableName + " 
WITH(UPDLOCK) WHERE " + whereCondition;
+    }
 }
diff --git 
a/rm-datasource/src/test/java/org/apache/seata/rm/datasource/undo/AbstractUndoExecutorTest.java
 
b/rm-datasource/src/test/java/org/apache/seata/rm/datasource/undo/BaseUndoExecutorTest.java
similarity index 69%
rename from 
rm-datasource/src/test/java/org/apache/seata/rm/datasource/undo/AbstractUndoExecutorTest.java
rename to 
rm-datasource/src/test/java/org/apache/seata/rm/datasource/undo/BaseUndoExecutorTest.java
index c97a65daaa..72e29405b2 100644
--- 
a/rm-datasource/src/test/java/org/apache/seata/rm/datasource/undo/AbstractUndoExecutorTest.java
+++ 
b/rm-datasource/src/test/java/org/apache/seata/rm/datasource/undo/BaseUndoExecutorTest.java
@@ -16,21 +16,32 @@
  */
 package org.apache.seata.rm.datasource.undo;
 
+import org.apache.seata.rm.datasource.ConnectionProxy;
 import org.apache.seata.rm.datasource.SqlGenerateUtils;
 import org.apache.seata.rm.datasource.sql.struct.Field;
 import org.apache.seata.rm.datasource.sql.struct.Row;
 import org.apache.seata.rm.datasource.sql.struct.TableRecords;
 import org.apache.seata.sqlparser.SQLType;
+import org.apache.seata.sqlparser.struct.ColumnMeta;
 import org.apache.seata.sqlparser.struct.TableMeta;
 import org.apache.seata.sqlparser.util.JdbcConstants;
+import org.checkerframework.checker.nullness.qual.NonNull;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
 import org.mockito.Mockito;
 
+import java.sql.Connection;
+import java.sql.PreparedStatement;
+import java.sql.ResultSet;
+import java.sql.ResultSetMetaData;
 import java.sql.SQLException;
 import java.util.*;
 
-public class AbstractUndoExecutorTest extends BaseH2Test {
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+public class BaseUndoExecutorTest extends BaseH2Test {
 
     @Test
     public void dataValidationUpdate() throws SQLException {
@@ -238,6 +249,105 @@ public class AbstractUndoExecutorTest extends BaseH2Test {
                 pkNameList, pkRowValues.get("id1").size(), 
JdbcConstants.POLARDBX);
         Assertions.assertEquals("(id1) in ( (?) )", sql.get(0).getSql());
     }
+
+    @Test
+    public void testBuildCheckSql() {
+        SQLUndoLog sqlUndoLog = new SQLUndoLog();
+        TestUndoExecutor executor = new TestUndoExecutor(sqlUndoLog, false);
+
+        String tableName = "table_name";
+        String whereCondition = "(id) in ( (?) )";
+
+        String sqlWithStar = executor.getCheckSql(tableName, whereCondition);
+        Assertions.assertEquals("SELECT * FROM table_name WHERE (id) in ( (?) 
) FOR UPDATE", sqlWithStar);
+
+        String selectColumns = "id, name, invisible_col";
+        String sqlWithExplicitColumns = executor.getCheckSql(tableName, 
whereCondition, selectColumns);
+        Assertions.assertEquals(
+                "SELECT id, name, invisible_col FROM table_name WHERE (id) in 
( (?) ) FOR UPDATE",
+                sqlWithExplicitColumns);
+    }
+
+    @Test
+    public void testQueryCurrentRecordsSqlGeneration() throws SQLException {
+        TableMeta tableMeta = mock(TableMeta.class);
+        when(tableMeta.getTableName()).thenReturn("table_name");
+        
when(tableMeta.getPrimaryKeyOnlyName()).thenReturn(Collections.singletonList("id"));
+
+        Map<String, ColumnMeta> allColumns = new LinkedHashMap<>();
+        allColumns.put("id", new ColumnMeta());
+        allColumns.put("name", new ColumnMeta());
+        allColumns.put("invisible_col", new ColumnMeta());
+        when(tableMeta.getAllColumns()).thenReturn(allColumns);
+
+        ColumnMeta idColMeta = new ColumnMeta();
+        idColMeta.setDataType(java.sql.Types.INTEGER);
+        when(tableMeta.getColumnMeta("id")).thenReturn(idColMeta);
+
+        TableRecords beforeImage = getTableRecords(tableMeta);
+
+        SQLUndoLog sqlUndoLog = new SQLUndoLog();
+        sqlUndoLog.setSqlType(SQLType.UPDATE);
+        sqlUndoLog.setTableMeta(tableMeta);
+        sqlUndoLog.setTableName("table_name");
+        sqlUndoLog.setBeforeImage(beforeImage);
+
+        TestUndoExecutor executor = new TestUndoExecutor(sqlUndoLog, true);
+
+        ConnectionProxy connectionProxy = mock(ConnectionProxy.class);
+        when(connectionProxy.getDbType()).thenReturn(JdbcConstants.ORACLE);
+
+        Connection connection = mock(Connection.class);
+        when(connectionProxy.getTargetConnection()).thenReturn(connection);
+
+        PreparedStatement preparedStatement = mock(PreparedStatement.class);
+        ResultSet resultSet = mock(ResultSet.class);
+
+        ResultSetMetaData metaData = mock(ResultSetMetaData.class);
+        when(metaData.getColumnCount()).thenReturn(0);
+        when(resultSet.getMetaData()).thenReturn(metaData);
+        when(resultSet.next()).thenReturn(false);
+
+        when(preparedStatement.executeQuery()).thenReturn(resultSet);
+
+        ArgumentCaptor<String> sqlCaptor = 
ArgumentCaptor.forClass(String.class);
+        
when(connection.prepareStatement(sqlCaptor.capture())).thenReturn(preparedStatement);
+
+        executor.queryCurrentRecords(connectionProxy);
+
+        String executedSql = sqlCaptor.getValue();
+
+        Assertions.assertTrue(
+                executedSql.contains("\"id\", \"name\"") && 
!executedSql.contains("\"invisible_col\""),
+                "The query should explicitly project ONLY the columns from the 
undo log, ignoring hidden columns in tableMeta.");
+        Assertions.assertTrue(
+                executedSql.startsWith("SELECT \"id\", \"name\" FROM 
table_name WHERE"),
+                "The query format was incorrect. Captured SQL: " + 
executedSql);
+    }
+
+    private static @NonNull TableRecords getTableRecords(TableMeta tableMeta) {
+        TableRecords beforeImage = new TableRecords();
+        beforeImage.setTableName("table_name");
+        beforeImage.setTableMeta(tableMeta);
+
+        List<Row> rows = new ArrayList<>();
+        Row row = new Row();
+        Field pkField = new Field();
+        pkField.setName("id");
+        pkField.setType(java.sql.Types.INTEGER);
+        pkField.setValue(12345);
+        row.add(pkField);
+
+        Field nameField = new Field();
+        nameField.setName("name");
+        nameField.setType(java.sql.Types.VARCHAR);
+        nameField.setValue("aaa");
+        row.add(nameField);
+
+        rows.add(row);
+        beforeImage.setRows(rows);
+        return beforeImage;
+    }
 }
 
 class TestUndoExecutor extends AbstractUndoExecutor {
@@ -257,4 +367,12 @@ class TestUndoExecutor extends AbstractUndoExecutor {
     protected TableRecords getUndoRows() {
         return isDelete ? sqlUndoLog.getBeforeImage() : 
sqlUndoLog.getAfterImage();
     }
+
+    public String getCheckSql(String tableName, String whereCondition) {
+        return super.buildCheckSql(tableName, whereCondition);
+    }
+
+    public String getCheckSql(String tableName, String whereCondition, String 
selectColumns) {
+        return super.buildCheckSql(tableName, whereCondition, selectColumns);
+    }
 }
diff --git 
a/rm-datasource/src/test/java/org/apache/seata/rm/datasource/undo/sqlserver/BaseSqlServerUndoExecutorTest.java
 
b/rm-datasource/src/test/java/org/apache/seata/rm/datasource/undo/sqlserver/SqlServerBaseUndoExecutorTest.java
similarity index 93%
rename from 
rm-datasource/src/test/java/org/apache/seata/rm/datasource/undo/sqlserver/BaseSqlServerUndoExecutorTest.java
rename to 
rm-datasource/src/test/java/org/apache/seata/rm/datasource/undo/sqlserver/SqlServerBaseUndoExecutorTest.java
index 5e6c7d03bb..2326803d5c 100644
--- 
a/rm-datasource/src/test/java/org/apache/seata/rm/datasource/undo/sqlserver/BaseSqlServerUndoExecutorTest.java
+++ 
b/rm-datasource/src/test/java/org/apache/seata/rm/datasource/undo/sqlserver/SqlServerBaseUndoExecutorTest.java
@@ -38,7 +38,7 @@ import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 
-public class BaseSqlServerUndoExecutorTest {
+public class SqlServerBaseUndoExecutorTest {
 
     private TestBaseSqlServerUndoExecutor executor;
     private SQLUndoLog sqlUndoLog;
@@ -306,4 +306,23 @@ public class BaseSqlServerUndoExecutorTest {
         Assertions.assertTrue(checkSql1.contains("id = 1"));
         Assertions.assertTrue(checkSql2.contains("id = 2"));
     }
+
+    @Test
+    public void testBuildCheckSqlWithExplicitColumns() {
+        String tableName = "test_table";
+        String whereCondition = "id = ?";
+        String selectColumns = "id, name, status";
+
+        String checkSql = executor.buildCheckSql(tableName, whereCondition, 
selectColumns);
+
+        Assertions.assertNotNull(checkSql);
+        Assertions.assertTrue(checkSql.contains("SELECT id, name, status 
FROM"));
+        Assertions.assertTrue(checkSql.contains(tableName));
+        Assertions.assertTrue(checkSql.contains("WITH(UPDLOCK)"));
+        Assertions.assertTrue(checkSql.contains("WHERE"));
+        Assertions.assertTrue(checkSql.contains(whereCondition));
+
+        String expectedSql = "SELECT id, name, status FROM test_table 
WITH(UPDLOCK) WHERE id = ?";
+        Assertions.assertEquals(expectedSql, checkSql);
+    }
 }


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

Reply via email to