Skip to content

Commit 10a9ff4

Browse files
authored
feat: receipt/record query failover to other nodes (#2613)
Signed-off-by: Mustafa Uzun <mustafa.uzun@limechain.tech>
1 parent 0aa2b12 commit 10a9ff4

5 files changed

Lines changed: 318 additions & 6 deletions

File tree

sdk/src/main/java/com/hedera/hashgraph/sdk/Client.java

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,7 @@ public final class Client implements AutoCloseable {
7474
private volatile Duration minBackoff = DEFAULT_MIN_BACKOFF;
7575
private boolean autoValidateChecksums = false;
7676
private boolean defaultRegenerateTransactionId = true;
77+
private boolean allowReceiptNodeFailover = false;
7778
private final boolean shouldShutdownExecutor;
7879
private final long shard;
7980
private final long realm;
@@ -1295,6 +1296,15 @@ public synchronized boolean getDefaultRegenerateTransactionId() {
12951296
return defaultRegenerateTransactionId;
12961297
}
12971298

1299+
/**
1300+
* Should node failover be enabled
1301+
*
1302+
* @return the node failover mode
1303+
*/
1304+
public synchronized boolean isAllowReceiptNodeFailover() {
1305+
return allowReceiptNodeFailover;
1306+
}
1307+
12981308
/**
12991309
* Assign the default regenerate transaction id.
13001310
*
@@ -1306,6 +1316,21 @@ public synchronized Client setDefaultRegenerateTransactionId(boolean regenerateT
13061316
return this;
13071317
}
13081318

1319+
/**
1320+
* Enable or disable receipt query failover to other nodes when the submitting node
1321+
* is unresponsive. When enabled, receipt queries will start with the submitting node
1322+
* but can fail over to other nodes in the network if needed.
1323+
* Default is `false` to preserve existing behavior where receipt queries are pinned
1324+
* to the submitting node only.
1325+
*
1326+
* @param allowReceiptNodeFailover should node failover be enabled
1327+
* @return {@code this}
1328+
*/
1329+
public synchronized Client setAllowReceiptNodeFailover(boolean allowReceiptNodeFailover) {
1330+
this.allowReceiptNodeFailover = allowReceiptNodeFailover;
1331+
return this;
1332+
}
1333+
13091334
/**
13101335
* Maximum amount of time a request can run
13111336
*

sdk/src/main/java/com/hedera/hashgraph/sdk/TokenRejectFlow.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -292,7 +292,7 @@ public CompletableFuture<TransactionResponse> executeAsync(Client client, Durati
292292
return createTokenRejectTransaction()
293293
.executeAsync(client, timeoutPerTransaction)
294294
.thenCompose(tokenRejectResponse ->
295-
tokenRejectResponse.getReceiptQuery().executeAsync(client, timeoutPerTransaction))
295+
tokenRejectResponse.getReceiptQuery(client).executeAsync(client, timeoutPerTransaction))
296296
.thenCompose(receipt -> createTokenDissociateTransaction().executeAsync(client, timeoutPerTransaction));
297297
}
298298

sdk/src/main/java/com/hedera/hashgraph/sdk/TransactionResponse.java

Lines changed: 37 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33

44
import com.google.common.base.MoreObjects;
55
import java.time.Duration;
6+
import java.util.ArrayList;
67
import java.util.Collections;
78
import java.util.List;
89
import java.util.concurrent.CompletableFuture;
@@ -134,7 +135,7 @@ public TransactionReceipt getReceipt(Client client, Duration timeout)
134135
while (attempts < MAX_RETRY_ATTEMPTS) {
135136
try {
136137
// Attempt to execute the receipt query
137-
return getReceiptQuery().execute(client, timeout).validateStatus(validateStatus);
138+
return getReceiptQuery(client).execute(client, timeout).validateStatus(validateStatus);
138139
} catch (ReceiptStatusException e) {
139140
// Check if the exception status indicates throttling or inner transaction throttling
140141
if (e.receipt.status == Status.THROTTLED_AT_CONSENSUS) {
@@ -194,6 +195,22 @@ public TransactionReceiptQuery getReceiptQuery() {
194195
.setNodeAccountIds(Collections.singletonList(nodeId));
195196
}
196197

198+
/**
199+
* Create receipt query from the {@link #transactionId} and {@link #transactionHash}
200+
*
201+
* @return {@link com.hedera.hashgraph.sdk.TransactionReceiptQuery}
202+
*/
203+
public TransactionReceiptQuery getReceiptQuery(Client client) {
204+
List<AccountId> nodeIds = new ArrayList<>(List.of(nodeId));
205+
if (client != null && client.isAllowReceiptNodeFailover()) {
206+
nodeIds.addAll(client.getNetwork().values().stream()
207+
.filter(id -> !id.equals(nodeId))
208+
.toList());
209+
}
210+
211+
return new TransactionReceiptQuery().setTransactionId(transactionId).setNodeAccountIds(nodeIds);
212+
}
213+
197214
/**
198215
* Fetch the receipt of the transaction asynchronously.
199216
*
@@ -212,7 +229,7 @@ public CompletableFuture<TransactionReceipt> getReceiptAsync(Client client) {
212229
* @return the transaction receipt
213230
*/
214231
public CompletableFuture<TransactionReceipt> getReceiptAsync(Client client, Duration timeout) {
215-
return getReceiptQuery().executeAsync(client, timeout).thenCompose(receipt -> {
232+
return getReceiptQuery(client).executeAsync(client, timeout).thenCompose(receipt -> {
216233
try {
217234
return CompletableFuture.completedFuture(receipt.validateStatus(validateStatus));
218235
} catch (ReceiptStatusException e) {
@@ -293,7 +310,7 @@ public TransactionRecord getRecord(Client client)
293310
public TransactionRecord getRecord(Client client, Duration timeout)
294311
throws TimeoutException, PrecheckStatusException, ReceiptStatusException {
295312
getReceipt(client, timeout);
296-
return getRecordQuery().execute(client, timeout);
313+
return getRecordQuery(client).execute(client, timeout);
297314
}
298315

299316
/**
@@ -307,6 +324,22 @@ public TransactionRecordQuery getRecordQuery() {
307324
.setNodeAccountIds(Collections.singletonList(nodeId));
308325
}
309326

327+
/**
328+
* Create record query from the {@link #transactionId} and {@link #transactionHash}
329+
*
330+
* @return {@link com.hedera.hashgraph.sdk.TransactionRecordQuery}
331+
*/
332+
public TransactionRecordQuery getRecordQuery(Client client) {
333+
List<AccountId> nodeIds = new ArrayList<>(List.of(nodeId));
334+
if (client != null && client.isAllowReceiptNodeFailover()) {
335+
nodeIds.addAll(client.getNetwork().values().stream()
336+
.filter(id -> !id.equals(nodeId))
337+
.toList());
338+
}
339+
340+
return new TransactionRecordQuery().setTransactionId(transactionId).setNodeAccountIds(nodeIds);
341+
}
342+
310343
/**
311344
* Fetch the record of the transaction asynchronously.
312345
*
@@ -326,7 +359,7 @@ public CompletableFuture<TransactionRecord> getRecordAsync(Client client) {
326359
*/
327360
public CompletableFuture<TransactionRecord> getRecordAsync(Client client, Duration timeout) {
328361
return getReceiptAsync(client, timeout)
329-
.thenCompose((receipt) -> getRecordQuery().executeAsync(client, timeout));
362+
.thenCompose((receipt) -> getRecordQuery(client).executeAsync(client, timeout));
330363
}
331364

332365
/**
Lines changed: 253 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,253 @@
1+
// SPDX-License-Identifier: Apache-2.0
2+
package com.hedera.hashgraph.sdk;
3+
4+
import com.hedera.hashgraph.sdk.proto.CryptoServiceGrpc;
5+
import com.hedera.hashgraph.sdk.proto.Query;
6+
import com.hedera.hashgraph.sdk.proto.Response;
7+
import com.hedera.hashgraph.sdk.proto.ResponseCodeEnum;
8+
import com.hedera.hashgraph.sdk.proto.ResponseHeader;
9+
import com.hedera.hashgraph.sdk.proto.Transaction;
10+
import com.hedera.hashgraph.sdk.proto.TransactionGetRecordResponse;
11+
import com.hedera.hashgraph.sdk.proto.TransactionRecord;
12+
import com.hedera.hashgraph.sdk.proto.TransactionResponse;
13+
import io.grpc.stub.StreamObserver;
14+
import java.util.List;
15+
import org.junit.jupiter.api.Assertions;
16+
import org.junit.jupiter.api.Test;
17+
18+
class TransactionResponseTest {
19+
20+
private static Response buildRecordResponse(ResponseCodeEnum precheckStatus, ResponseCodeEnum receiptStatus) {
21+
return Response.newBuilder()
22+
.setTransactionGetRecord(TransactionGetRecordResponse.newBuilder()
23+
.setHeader(ResponseHeader.newBuilder()
24+
.setNodeTransactionPrecheckCode(precheckStatus)
25+
.build())
26+
.setTransactionRecord(TransactionRecord.newBuilder()
27+
.setReceipt(com.hedera.hashgraph.sdk.proto.TransactionReceipt.newBuilder()
28+
.setStatus(receiptStatus)
29+
.build())
30+
.build())
31+
.build())
32+
.build();
33+
}
34+
35+
private static Response buildSuccessRecordResponse() {
36+
return buildRecordResponse(ResponseCodeEnum.OK, ResponseCodeEnum.SUCCESS);
37+
}
38+
39+
@Test
40+
void getReceiptPinnedToSubmittingNodeByDefault() throws Exception {
41+
var service = new TestCryptoService();
42+
var server = new TestServer("getReceiptPinnedDefault", service);
43+
44+
service.buffer.enqueueResponse(TestResponse.transactionOk());
45+
service.buffer.enqueueResponse(TestResponse.successfulReceipt());
46+
47+
var txResponse = new AccountCreateTransaction().execute(server.client);
48+
var receipt = txResponse.getReceipt(server.client);
49+
50+
var receiptQuery = txResponse.getReceiptQuery(server.client);
51+
Assertions.assertEquals(1, receiptQuery.getNodeAccountIds().size());
52+
Assertions.assertEquals(
53+
txResponse.nodeId, receiptQuery.getNodeAccountIds().get(0));
54+
Assertions.assertNotNull(receipt);
55+
56+
server.close();
57+
}
58+
59+
@Test
60+
void getRecordPinnedToSubmittingNodeByDefault() throws Exception {
61+
var service = new TestCryptoService();
62+
var server = new TestServer("getRecordPinnedDefault", service);
63+
64+
service.buffer.enqueueResponse(TestResponse.transactionOk());
65+
service.buffer.enqueueResponse(TestResponse.successfulReceipt());
66+
service.buffer.enqueueResponse(TestResponse.successfulReceipt());
67+
service.buffer.enqueueResponse(TestResponse.query(buildSuccessRecordResponse()));
68+
69+
var txResponse = new AccountCreateTransaction().execute(server.client);
70+
var record = txResponse.getRecord(server.client);
71+
72+
var recordQuery = txResponse.getRecordQuery(server.client);
73+
Assertions.assertEquals(1, recordQuery.getNodeAccountIds().size());
74+
Assertions.assertEquals(
75+
txResponse.nodeId, recordQuery.getNodeAccountIds().get(0));
76+
Assertions.assertNotNull(record);
77+
78+
server.close();
79+
}
80+
81+
@Test
82+
void failoverEnabledSubmittingNodeQueriedFirst() throws Exception {
83+
var service = new TestCryptoService();
84+
var server = new TestServer("failoverEnabledFirst", service);
85+
server.client.setAllowReceiptNodeFailover(true);
86+
87+
service.buffer.enqueueResponse(TestResponse.transactionOk());
88+
service.buffer.enqueueResponse(TestResponse.successfulReceipt());
89+
service.buffer.enqueueResponse(TestResponse.successfulReceipt());
90+
service.buffer.enqueueResponse(TestResponse.query(buildSuccessRecordResponse()));
91+
92+
var txResponse = new AccountCreateTransaction().execute(server.client);
93+
94+
var receiptQuery = txResponse.getReceiptQuery(server.client);
95+
Assertions.assertEquals(2, receiptQuery.getNodeAccountIds().size());
96+
Assertions.assertEquals(
97+
txResponse.nodeId, receiptQuery.getNodeAccountIds().get(0));
98+
99+
var recordQuery = txResponse.getRecordQuery(server.client);
100+
Assertions.assertEquals(2, recordQuery.getNodeAccountIds().size());
101+
Assertions.assertEquals(
102+
txResponse.nodeId, recordQuery.getNodeAccountIds().get(0));
103+
104+
var record = txResponse.getRecord(server.client);
105+
Assertions.assertNotNull(record);
106+
107+
server.close();
108+
}
109+
110+
@Test
111+
void receiptFailoverOnUnavailableAdvancesToNextNode() throws Exception {
112+
var service = new TestCryptoService();
113+
var server = new TestServer("receiptFailoverUnavailable", service);
114+
server.client.setAllowReceiptNodeFailover(true);
115+
116+
service.buffer.enqueueResponse(TestResponse.transactionOk());
117+
service.buffer.enqueueResponse(TestResponse.error(io.grpc.Status.UNAVAILABLE.asRuntimeException()));
118+
service.buffer.enqueueResponse(TestResponse.successfulReceipt());
119+
120+
var txResponse = new AccountCreateTransaction().execute(server.client);
121+
var receipt = txResponse.getReceipt(server.client);
122+
123+
Assertions.assertNotNull(receipt);
124+
Assertions.assertEquals(2, service.buffer.queryRequestsReceived.size());
125+
126+
server.close();
127+
}
128+
129+
@Test
130+
void recordFailoverOnUnavailableAdvancesToNextNode() throws Exception {
131+
var service = new TestCryptoService();
132+
var server = new TestServer("recordFailoverUnavailable", service);
133+
server.client.setAllowReceiptNodeFailover(true);
134+
135+
service.buffer.enqueueResponse(TestResponse.transactionOk());
136+
service.buffer.enqueueResponse(TestResponse.successfulReceipt());
137+
service.buffer.enqueueResponse(TestResponse.successfulReceipt());
138+
service.buffer.enqueueResponse(TestResponse.error(io.grpc.Status.UNAVAILABLE.asRuntimeException()));
139+
service.buffer.enqueueResponse(TestResponse.query(buildSuccessRecordResponse()));
140+
141+
var txResponse = new AccountCreateTransaction().execute(server.client);
142+
var record = txResponse.getRecord(server.client);
143+
144+
Assertions.assertNotNull(record);
145+
146+
server.close();
147+
}
148+
149+
@Test
150+
void failoverWithExplicitTransactionNodes() throws Exception {
151+
var service = new TestCryptoService();
152+
var server = new TestServer("failoverExplicitNodes", service);
153+
server.client.setAllowReceiptNodeFailover(true);
154+
155+
service.buffer.enqueueResponse(TestResponse.transactionOk());
156+
service.buffer.enqueueResponse(TestResponse.successfulReceipt());
157+
service.buffer.enqueueResponse(TestResponse.successfulReceipt());
158+
service.buffer.enqueueResponse(TestResponse.query(buildSuccessRecordResponse()));
159+
160+
var txResponse = new AccountCreateTransaction()
161+
.setNodeAccountIds(List.of(AccountId.fromString("1.1.1"), AccountId.fromString("1.1.2")))
162+
.execute(server.client);
163+
164+
var receiptQuery = txResponse.getReceiptQuery(server.client);
165+
Assertions.assertEquals(2, receiptQuery.getNodeAccountIds().size());
166+
Assertions.assertEquals(
167+
txResponse.nodeId, receiptQuery.getNodeAccountIds().get(0));
168+
169+
var recordQuery = txResponse.getRecordQuery(server.client);
170+
Assertions.assertEquals(2, recordQuery.getNodeAccountIds().size());
171+
Assertions.assertEquals(
172+
txResponse.nodeId, recordQuery.getNodeAccountIds().get(0));
173+
174+
var record = txResponse.getRecord(server.client);
175+
Assertions.assertNotNull(record);
176+
177+
server.close();
178+
}
179+
180+
@Test
181+
void failoverWithoutExplicitTransactionNodes() throws Exception {
182+
var service = new TestCryptoService();
183+
var server = new TestServer("failoverNoExplicitNodes", service);
184+
server.client.setAllowReceiptNodeFailover(true);
185+
186+
service.buffer.enqueueResponse(TestResponse.transactionOk());
187+
service.buffer.enqueueResponse(TestResponse.successfulReceipt());
188+
service.buffer.enqueueResponse(TestResponse.successfulReceipt());
189+
service.buffer.enqueueResponse(TestResponse.query(buildSuccessRecordResponse()));
190+
191+
var txResponse = new AccountCreateTransaction().execute(server.client);
192+
193+
var receiptQuery = txResponse.getReceiptQuery(server.client);
194+
Assertions.assertEquals(2, receiptQuery.getNodeAccountIds().size());
195+
Assertions.assertEquals(
196+
txResponse.nodeId, receiptQuery.getNodeAccountIds().get(0));
197+
198+
var recordQuery = txResponse.getRecordQuery(server.client);
199+
Assertions.assertEquals(2, recordQuery.getNodeAccountIds().size());
200+
Assertions.assertEquals(
201+
txResponse.nodeId, recordQuery.getNodeAccountIds().get(0));
202+
203+
var record = txResponse.getRecord(server.client);
204+
Assertions.assertNotNull(record);
205+
206+
server.close();
207+
}
208+
209+
@Test
210+
void defaultBehaviorPinnedWhenNodeUnhealthy() throws Exception {
211+
var service = new TestCryptoService();
212+
var server = new TestServer("pinnedWhenUnhealthy", service);
213+
server.client.setMaxAttempts(2);
214+
215+
service.buffer.enqueueResponse(TestResponse.transactionOk());
216+
service.buffer.enqueueResponse(TestResponse.error(io.grpc.Status.UNAVAILABLE.asRuntimeException()));
217+
service.buffer.enqueueResponse(TestResponse.error(io.grpc.Status.UNAVAILABLE.asRuntimeException()));
218+
219+
var txResponse = new AccountCreateTransaction().execute(server.client);
220+
221+
Assertions.assertThrows(Exception.class, () -> {
222+
txResponse.getReceipt(server.client);
223+
});
224+
225+
Assertions.assertEquals(2, service.buffer.queryRequestsReceived.size());
226+
227+
server.close();
228+
}
229+
230+
private static class TestCryptoService extends CryptoServiceGrpc.CryptoServiceImplBase implements TestService {
231+
public Buffer buffer = new Buffer();
232+
233+
@Override
234+
public Buffer getBuffer() {
235+
return buffer;
236+
}
237+
238+
@Override
239+
public void createAccount(Transaction request, StreamObserver<TransactionResponse> responseObserver) {
240+
respondToTransactionFromQueue(request, responseObserver);
241+
}
242+
243+
@Override
244+
public void getTransactionReceipts(Query request, StreamObserver<Response> responseObserver) {
245+
respondToQueryFromQueue(request, responseObserver);
246+
}
247+
248+
@Override
249+
public void getTxRecordByTxID(Query request, StreamObserver<Response> responseObserver) {
250+
respondToQueryFromQueue(request, responseObserver);
251+
}
252+
}
253+
}

0 commit comments

Comments
 (0)