This is an automated email from the ASF dual-hosted git repository.
terrymanu pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shardingsphere.git
The following commit(s) were added to refs/heads/master by this push:
new 7a6cbfec2f3 Refactor FirebirdFetchStatementCommandExecutor (#38853)
7a6cbfec2f3 is described below
commit 7a6cbfec2f3f185adef2f4a60408b0fd78d5a2b9
Author: Timofey Sakharovsky <[email protected]>
AuthorDate: Fri Jun 19 19:16:38 2026 +0300
Refactor FirebirdFetchStatementCommandExecutor (#38853)
* refactor FirebirdFetchStatementCommandExecutor
* fix fetchCount overflow
* update unit tests
* checkstyle fix
---
.../FirebirdFetchStatementCommandExecutor.java | 59 +++++++----
.../FirebirdFetchStatementCommandExecutorTest.java | 114 ++++++++++++++-------
2 files changed, 117 insertions(+), 56 deletions(-)
diff --git
a/proxy/frontend/dialect/firebird/src/main/java/org/apache/shardingsphere/proxy/frontend/firebird/command/query/statement/fetch/FirebirdFetchStatementCommandExecutor.java
b/proxy/frontend/dialect/firebird/src/main/java/org/apache/shardingsphere/proxy/frontend/firebird/command/query/statement/fetch/FirebirdFetchStatementCommandExecutor.java
index 3b1bd35c37d..153194c1563 100644
---
a/proxy/frontend/dialect/firebird/src/main/java/org/apache/shardingsphere/proxy/frontend/firebird/command/query/statement/fetch/FirebirdFetchStatementCommandExecutor.java
+++
b/proxy/frontend/dialect/firebird/src/main/java/org/apache/shardingsphere/proxy/frontend/firebird/command/query/statement/fetch/FirebirdFetchStatementCommandExecutor.java
@@ -17,7 +17,6 @@
package
org.apache.shardingsphere.proxy.frontend.firebird.command.query.statement.fetch;
-import lombok.RequiredArgsConstructor;
import org.apache.shardingsphere.database.protocol.binary.BinaryCell;
import org.apache.shardingsphere.database.protocol.binary.BinaryRow;
import
org.apache.shardingsphere.database.protocol.firebird.packet.command.query.FirebirdBinaryColumnType;
@@ -28,45 +27,65 @@ import
org.apache.shardingsphere.proxy.backend.handler.ProxyBackendHandler;
import org.apache.shardingsphere.proxy.backend.response.data.QueryResponseCell;
import org.apache.shardingsphere.proxy.backend.response.data.QueryResponseRow;
import org.apache.shardingsphere.proxy.backend.session.ConnectionSession;
-import
org.apache.shardingsphere.proxy.frontend.command.executor.CommandExecutor;
+import
org.apache.shardingsphere.proxy.frontend.command.executor.QueryCommandExecutor;
+import org.apache.shardingsphere.proxy.frontend.command.executor.ResponseType;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.Collection;
-import java.util.LinkedList;
+import java.util.Collections;
import java.util.List;
/**
* Firebird fetch statement command executor.
*/
-@RequiredArgsConstructor
-public final class FirebirdFetchStatementCommandExecutor implements
CommandExecutor {
-
- private final FirebirdFetchStatementPacket packet;
+public final class FirebirdFetchStatementCommandExecutor implements
QueryCommandExecutor {
private final ConnectionSession connectionSession;
+ private final ProxyBackendHandler proxyBackendHandler;
+
+ private int fetchCount;
+
+ public FirebirdFetchStatementCommandExecutor(final
FirebirdFetchStatementPacket packet, final ConnectionSession connectionSession)
{
+ this.connectionSession = connectionSession;
+ proxyBackendHandler =
FirebirdFetchStatementCache.getInstance().getFetchBackendHandler(connectionSession.getConnectionId(),
packet.getStatementId());
+ fetchCount = packet.getFetchSize();
+ }
+
@Override
public Collection<DatabasePacket> execute() throws SQLException {
- Collection<DatabasePacket> result = new LinkedList<>();
- ProxyBackendHandler proxyBackendHandler =
FirebirdFetchStatementCache.getInstance().getFetchBackendHandler(connectionSession.getConnectionId(),
packet.getStatementId());
- if (null == proxyBackendHandler) {
- result.add(FirebirdFetchResponsePacket.getFetchNoMoreRowsPacket());
- return result;
+ if (proxyBackendHandler == null) {
+ fetchCount = -1;
+ return
Collections.singletonList(FirebirdFetchResponsePacket.getFetchNoMoreRowsPacket());
}
- for (int i = 0; i < packet.getFetchSize(); i++) {
+ return Collections.singletonList(getQueryRowPacket());
+ }
+
+ @Override
+ public ResponseType getResponseType() {
+ return ResponseType.QUERY;
+ }
+
+ @Override
+ public boolean next() throws SQLException {
+ return 0 <= fetchCount;
+ }
+
+ @Override
+ public DatabasePacket getQueryRowPacket() throws SQLException {
+ fetchCount--;
+ if (0 <= fetchCount) {
if (proxyBackendHandler.next()) {
- QueryResponseRow queryResponseRow =
proxyBackendHandler.getRowData();
- BinaryRow row = createBinaryRow(queryResponseRow);
- result.add(FirebirdFetchResponsePacket.getFetchRowPacket(row));
+ BinaryRow row =
createBinaryRow(proxyBackendHandler.getRowData());
+ return FirebirdFetchResponsePacket.getFetchRowPacket(row);
} else {
connectionSession.getDatabaseConnectionManager().unmarkResourceInUse(proxyBackendHandler);
-
result.add(FirebirdFetchResponsePacket.getFetchNoMoreRowsPacket());
- return result;
+ fetchCount = -1;
+ return FirebirdFetchResponsePacket.getFetchNoMoreRowsPacket();
}
}
- result.add(FirebirdFetchResponsePacket.getFetchEndPacket());
- return result;
+ return FirebirdFetchResponsePacket.getFetchEndPacket();
}
private BinaryRow createBinaryRow(final QueryResponseRow queryResponseRow)
{
diff --git
a/proxy/frontend/dialect/firebird/src/test/java/org/apache/shardingsphere/proxy/frontend/firebird/command/query/statement/fetch/FirebirdFetchStatementCommandExecutorTest.java
b/proxy/frontend/dialect/firebird/src/test/java/org/apache/shardingsphere/proxy/frontend/firebird/command/query/statement/fetch/FirebirdFetchStatementCommandExecutorTest.java
index 4f3902edd64..0e7d3a71dc6 100644
---
a/proxy/frontend/dialect/firebird/src/test/java/org/apache/shardingsphere/proxy/frontend/firebird/command/query/statement/fetch/FirebirdFetchStatementCommandExecutorTest.java
+++
b/proxy/frontend/dialect/firebird/src/test/java/org/apache/shardingsphere/proxy/frontend/firebird/command/query/statement/fetch/FirebirdFetchStatementCommandExecutorTest.java
@@ -25,6 +25,7 @@ import
org.apache.shardingsphere.proxy.backend.response.data.QueryResponseCell;
import org.apache.shardingsphere.proxy.backend.response.data.QueryResponseRow;
import org.apache.shardingsphere.proxy.backend.session.ConnectionSession;
import
org.apache.shardingsphere.proxy.backend.connector.ProxyDatabaseConnectionManager;
+import org.apache.shardingsphere.proxy.frontend.command.executor.ResponseType;
import org.firebirdsql.gds.ISCConstants;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
@@ -38,10 +39,11 @@ import java.sql.SQLException;
import java.sql.Types;
import java.util.Collection;
import java.util.Collections;
-import java.util.Iterator;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.is;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.Mockito.when;
import static org.mockito.Mockito.verify;
import static org.junit.jupiter.api.Assertions.assertNull;
@@ -83,8 +85,12 @@ class FirebirdFetchStatementCommandExecutorTest {
@Test
void assertExecuteWhenNoBackendHandler() throws SQLException {
+ when(packet.getFetchSize()).thenReturn(Integer.MAX_VALUE);
executor = new FirebirdFetchStatementCommandExecutor(packet,
connectionSession);
- assertNoMoreRowsResponse(executor.execute());
+ Collection<DatabasePacket> actualPackets = executor.execute();
+ assertThat(actualPackets.size(), is(1));
+ assertNoMoreRowsResponse((FirebirdFetchResponsePacket)
actualPackets.iterator().next());
+ assertFalse(executor.next());
}
@Test
@@ -92,62 +98,98 @@ class FirebirdFetchStatementCommandExecutorTest {
FirebirdFetchStatementCache.getInstance().registerStatement(CONNECTION_ID,
STATEMENT_ID, proxyBackendHandler);
FirebirdFetchStatementCache.getInstance().unregisterStatement(CONNECTION_ID,
STATEMENT_ID);
executor = new FirebirdFetchStatementCommandExecutor(packet,
connectionSession);
- assertNoMoreRowsResponse(executor.execute());
+ Collection<DatabasePacket> actualPackets = executor.execute();
+ assertThat(actualPackets.size(), is(1));
+ assertNoMoreRowsResponse((FirebirdFetchResponsePacket)
actualPackets.iterator().next());
+ assertFalse(executor.next());
+ }
+
+ @Test
+ void assertGetResponseType() {
+ executor = new FirebirdFetchStatementCommandExecutor(packet,
connectionSession);
+ assertThat(executor.getResponseType(), is(ResponseType.QUERY));
}
@Test
- void assertExecuteWithRowsAndFetchEnd() throws SQLException {
+ void assertNextReturnsFalseWhenFetchCountExceeded() throws SQLException {
+ executor = new FirebirdFetchStatementCommandExecutor(packet,
connectionSession);
+ executor.execute();
+ assertFalse(executor.next());
+ }
+
+ @Test
+ void assertNextReturnsTrueWhenFetchCountNotExceeded() throws SQLException {
FirebirdFetchStatementCache.getInstance().registerStatement(CONNECTION_ID,
STATEMENT_ID, proxyBackendHandler);
when(packet.getFetchSize()).thenReturn(2);
+ when(proxyBackendHandler.next()).thenReturn(true);
QueryResponseRow responseRow = new
QueryResponseRow(Collections.singletonList(new QueryResponseCell(Types.INTEGER,
1)));
- when(proxyBackendHandler.next()).thenReturn(true, true);
- when(proxyBackendHandler.getRowData()).thenReturn(responseRow,
responseRow);
+ when(proxyBackendHandler.getRowData()).thenReturn(responseRow);
executor = new FirebirdFetchStatementCommandExecutor(packet,
connectionSession);
- Collection<DatabasePacket> actualPackets = executor.execute();
- Iterator<DatabasePacket> packetIterator = actualPackets.iterator();
- FirebirdFetchResponsePacket actualFirstPacket =
(FirebirdFetchResponsePacket) packetIterator.next();
- assertThat(actualPackets.size(), is(3));
- assertThat(actualFirstPacket.getStatus(), is(ISCConstants.FETCH_OK));
- assertThat(actualFirstPacket.getCount(), is(1));
- assertNotNull(actualFirstPacket.getRow());
- FirebirdFetchResponsePacket actualSecondPacket =
(FirebirdFetchResponsePacket) packetIterator.next();
- assertThat(actualSecondPacket.getStatus(), is(ISCConstants.FETCH_OK));
- assertThat(actualSecondPacket.getCount(), is(1));
- assertNotNull(actualSecondPacket.getRow());
- FirebirdFetchResponsePacket actualEndPacket =
(FirebirdFetchResponsePacket) packetIterator.next();
- assertThat(actualEndPacket.getStatus(), is(ISCConstants.FETCH_OK));
- assertThat(actualEndPacket.getCount(), is(0));
- assertNull(actualEndPacket.getRow());
+ executor.execute();
+ assertTrue(executor.next());
+
FirebirdFetchStatementCache.getInstance().unregisterStatement(CONNECTION_ID,
STATEMENT_ID);
}
@Test
- void assertExecuteStopsWhenRowsExhausted() throws SQLException {
+ void assertExecuteWhenBackendHandlerReturnsRow() throws SQLException {
FirebirdFetchStatementCache.getInstance().registerStatement(CONNECTION_ID,
STATEMENT_ID, proxyBackendHandler);
when(packet.getFetchSize()).thenReturn(2);
QueryResponseRow responseRow = new
QueryResponseRow(Collections.singletonList(new QueryResponseCell(Types.INTEGER,
1)));
- when(proxyBackendHandler.next()).thenReturn(true, false);
+ when(proxyBackendHandler.next()).thenReturn(true);
when(proxyBackendHandler.getRowData()).thenReturn(responseRow);
executor = new FirebirdFetchStatementCommandExecutor(packet,
connectionSession);
Collection<DatabasePacket> actualPackets = executor.execute();
- Iterator<DatabasePacket> packetIterator = actualPackets.iterator();
- FirebirdFetchResponsePacket actualRowPacket =
(FirebirdFetchResponsePacket) packetIterator.next();
- assertThat(actualPackets.size(), is(2));
- assertThat(actualRowPacket.getStatus(), is(ISCConstants.FETCH_OK));
- assertThat(actualRowPacket.getCount(), is(1));
- assertNotNull(actualRowPacket.getRow());
- FirebirdFetchResponsePacket actualNoMorePacket =
(FirebirdFetchResponsePacket) packetIterator.next();
- assertThat(actualNoMorePacket.getStatus(),
is(ISCConstants.FETCH_NO_MORE_ROWS));
- assertThat(actualNoMorePacket.getCount(), is(0));
- assertNull(actualNoMorePacket.getRow());
+ assertThat(actualPackets.size(), is(1));
+ assertFetchRowResponse((FirebirdFetchResponsePacket)
actualPackets.iterator().next());
+ assertTrue(executor.next());
+ assertFetchRowResponse((FirebirdFetchResponsePacket)
executor.getQueryRowPacket());
+ assertTrue(executor.next());
+ assertFetchEndResponse((FirebirdFetchResponsePacket)
executor.getQueryRowPacket());
+ assertFalse(executor.next());
+
FirebirdFetchStatementCache.getInstance().unregisterStatement(CONNECTION_ID,
STATEMENT_ID);
+ }
+
+ @Test
+ void assertExecuteWhenBackendHandlerReturnsNoMoreRows() throws
SQLException {
+
FirebirdFetchStatementCache.getInstance().registerStatement(CONNECTION_ID,
STATEMENT_ID, proxyBackendHandler);
+ when(packet.getFetchSize()).thenReturn(Integer.MAX_VALUE);
+ when(proxyBackendHandler.next()).thenReturn(false);
+ executor = new FirebirdFetchStatementCommandExecutor(packet,
connectionSession);
+ Collection<DatabasePacket> actualPackets = executor.execute();
+ assertThat(actualPackets.size(), is(1));
+ assertNoMoreRowsResponse((FirebirdFetchResponsePacket)
actualPackets.iterator().next());
verify(databaseConnectionManager).unmarkResourceInUse(proxyBackendHandler);
+ assertFalse(executor.next());
+
FirebirdFetchStatementCache.getInstance().unregisterStatement(CONNECTION_ID,
STATEMENT_ID);
}
- private void assertNoMoreRowsResponse(final Collection<DatabasePacket>
actualPackets) {
- Iterator<DatabasePacket> packetIterator = actualPackets.iterator();
- FirebirdFetchResponsePacket actualPacket =
(FirebirdFetchResponsePacket) packetIterator.next();
+ @Test
+ void assertGetQueryRowPacketWhenEnd() throws SQLException {
+
FirebirdFetchStatementCache.getInstance().registerStatement(CONNECTION_ID,
STATEMENT_ID, proxyBackendHandler);
+ when(packet.getFetchSize()).thenReturn(0);
+ executor = new FirebirdFetchStatementCommandExecutor(packet,
connectionSession);
+ Collection<DatabasePacket> actualPackets = executor.execute();
assertThat(actualPackets.size(), is(1));
+ assertFetchEndResponse((FirebirdFetchResponsePacket)
actualPackets.iterator().next());
+ assertFalse(executor.next());
+
FirebirdFetchStatementCache.getInstance().unregisterStatement(CONNECTION_ID,
STATEMENT_ID);
+ }
+
+ private void assertNoMoreRowsResponse(final FirebirdFetchResponsePacket
actualPacket) {
assertThat(actualPacket.getStatus(),
is(ISCConstants.FETCH_NO_MORE_ROWS));
assertThat(actualPacket.getCount(), is(0));
assertNull(actualPacket.getRow());
}
+
+ private void assertFetchRowResponse(final FirebirdFetchResponsePacket
actualPacket) {
+ assertThat(actualPacket.getStatus(), is(ISCConstants.FETCH_OK));
+ assertThat(actualPacket.getCount(), is(1));
+ assertNotNull(actualPacket.getRow());
+ }
+
+ private void assertFetchEndResponse(final FirebirdFetchResponsePacket
actualPacket) {
+ assertThat(actualPacket.getStatus(), is(ISCConstants.FETCH_OK));
+ assertThat(actualPacket.getCount(), is(0));
+ assertNull(actualPacket.getRow());
+ }
}