|
6 | 6 |
|
7 | 7 | import java.io.IOException; |
8 | 8 | import java.util.Map; |
| 9 | +import java.util.concurrent.CompletableFuture; |
9 | 10 | import java.util.concurrent.ConcurrentHashMap; |
10 | 11 | import java.util.concurrent.ExecutionException; |
11 | 12 | import java.util.concurrent.Executor; |
@@ -68,19 +69,19 @@ public int acceptRead(NodeRecord nodeRecord, Consumer<Bytes> onContentReceived) |
68 | 69 | return connectionId; |
69 | 70 | } |
70 | 71 |
|
71 | | - public void offerWrite(final NodeRecord nodeRecord, final int connectionId, Bytes content) { |
72 | | - this.runAsyncUTP( |
| 72 | + public SafeFuture<Void> offerWrite(NodeRecord nodeRecord, int connectionId, Bytes content) { |
| 73 | + return runAsyncUTPWithFuture( |
73 | 74 | () -> { |
74 | | - UTPClient utpClient = this.registerClient(nodeRecord, connectionId); |
| 75 | + UTPClient utpClient = registerClient(nodeRecord, connectionId); |
75 | 76 | utpClient |
76 | 77 | .connect(connectionId, new UTPAddress(nodeRecord)) |
77 | | - .thenCompose(__ -> utpClient.write(content, this.utpExecutor)) |
78 | | - .get(); |
| 78 | + .thenCompose(__ -> utpClient.write(content, utpExecutor)) |
| 79 | + .join(); |
79 | 80 | }, |
80 | 81 | "offerWrite", |
81 | 82 | nodeRecord, |
82 | 83 | connectionId, |
83 | | - this.utpExecutor); |
| 84 | + utpExecutor); |
84 | 85 | } |
85 | 86 |
|
86 | 87 | public int foundContentWrite(NodeRecord nodeRecord, Bytes content) { |
@@ -195,6 +196,21 @@ private void runAsyncUTP( |
195 | 196 | .exceptionally(defaultUTPErrorLog(operationName, nodeRecord, connectionId)); |
196 | 197 | } |
197 | 198 |
|
| 199 | + private SafeFuture<Void> runAsyncUTPWithFuture( |
| 200 | + RunnableUTP task, |
| 201 | + String operationName, |
| 202 | + NodeRecord nodeRecord, |
| 203 | + int connectionId, |
| 204 | + Executor executor) { |
| 205 | + |
| 206 | + CompletableFuture<Void> future = |
| 207 | + SafeFuture.runAsync( |
| 208 | + () -> executeWithHandling(task, operationName, nodeRecord, connectionId), executor) |
| 209 | + .exceptionally(defaultUTPErrorLog(operationName, nodeRecord, connectionId)); |
| 210 | + |
| 211 | + return SafeFuture.of(future); |
| 212 | + } |
| 213 | + |
198 | 214 | private void executeWithHandling( |
199 | 215 | RunnableUTP task, String operationName, NodeRecord nodeRecord, int connectionId) { |
200 | 216 | try { |
|
0 commit comments