Skip to content

Commit 4ebb485

Browse files
committed
refactor: remove retry mechanism (kiss)
1 parent 01b1c5c commit 4ebb485

3 files changed

Lines changed: 39 additions & 60 deletions

File tree

src/main/java/org/eclipse/dataplane/Dataplane.java

Lines changed: 15 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
import org.eclipse.dataplane.logic.OnStarted;
1818
import org.eclipse.dataplane.logic.OnTerminate;
1919
import org.eclipse.dataplane.port.DataPlaneSignalingApiController;
20+
import org.eclipse.dataplane.port.exception.DataFlowNotifyCompletedFailed;
2021
import org.eclipse.dataplane.port.exception.DataplaneNotRegistered;
2122
import org.eclipse.dataplane.port.store.DataFlowStore;
2223
import org.eclipse.dataplane.port.store.InMemoryDataFlowStore;
@@ -28,11 +29,6 @@
2829
import java.util.HashSet;
2930
import java.util.Set;
3031
import java.util.UUID;
31-
import java.util.concurrent.CompletableFuture;
32-
import java.util.concurrent.TimeUnit;
33-
import java.util.function.Function;
34-
35-
import static java.util.concurrent.CompletableFuture.delayedExecutor;
3632

3733
public class Dataplane {
3834

@@ -135,36 +131,27 @@ public Result<Void> terminate(String dataFlowId, DataFlowTerminateMessage messag
135131
*
136132
* @param dataFlowId
137133
*/
138-
public Result<CompletableFuture<Void>> notifyCompleted(String dataFlowId) {
134+
public Result<Void> notifyCompleted(String dataFlowId) {
139135
return store.findById(dataFlowId)
140-
.map(dataFlow -> transferDataFlowCompleted(dataFlow)
141-
.thenApply(r -> {
142-
dataFlow.transitionToCompleted();
143-
store.save(dataFlow);
144-
return null;
145-
}));
146-
}
136+
.compose(dataFlow -> {
137+
var endpoint = dataFlow.getCallbackAddress() + "/transfers/" + dataFlow.getId() + "/dataflow/completed";
147138

148-
private CompletableFuture<HttpResponse<Void>> transferDataFlowCompleted(DataFlow dataFlow) {
149-
var endpoint = dataFlow.getCallbackAddress() + "/transfers/" + dataFlow.getId() + "/dataflow/completed";
139+
var request = HttpRequest.newBuilder()
140+
.uri(URI.create(endpoint))
141+
.header("content-type", "application/json")
142+
.POST(HttpRequest.BodyPublishers.ofString("{}")) // TODO DataFlowCompletedMessage not defined
143+
.build();
150144

151-
var request = HttpRequest.newBuilder()
152-
.uri(URI.create(endpoint))
153-
.header("content-type", "application/json")
154-
.POST(HttpRequest.BodyPublishers.ofString("{}")) // TODO DataFlowCompletedMessage not defined
155-
.build();
145+
var response = httpClient.send(request, HttpResponse.BodyHandlers.discarding());
156146

157-
return httpClient.sendAsync(request, HttpResponse.BodyHandlers.discarding())
158-
.thenApply(r -> {
159-
var successful = r.statusCode() >= 200 && r.statusCode() < 300;
147+
var successful = response.statusCode() >= 200 && response.statusCode() < 300;
160148
if (successful) {
161-
return CompletableFuture.completedFuture(r);
149+
dataFlow.transitionToCompleted();
150+
return store.save(dataFlow);
162151
}
163152

164-
return CompletableFuture.supplyAsync(() -> transferDataFlowCompleted(dataFlow), delayedExecutor(500, TimeUnit.MILLISECONDS)).thenCompose(Function.identity());
165-
})
166-
.exceptionally(CompletableFuture::failedFuture)
167-
.thenCompose(Function.identity());
153+
return Result.failure(new DataFlowNotifyCompletedFailed(response));
154+
});
168155
}
169156

170157
/**
Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
1+
package org.eclipse.dataplane.port.exception;
2+
3+
import java.net.http.HttpResponse;
4+
5+
public class DataFlowNotifyCompletedFailed extends Exception {
6+
private final HttpResponse<Void> response;
7+
8+
public DataFlowNotifyCompletedFailed(HttpResponse<Void> response) {
9+
super("control-plane responded with %s".formatted(response.statusCode()));
10+
this.response = response;
11+
}
12+
13+
public HttpResponse<Void> getResponse() {
14+
return response;
15+
}
16+
}

src/test/java/org/eclipse/dataplane/DataplaneTest.java

Lines changed: 8 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,17 @@
11
package org.eclipse.dataplane;
22

33
import com.github.tomakehurst.wiremock.WireMockServer;
4-
import com.github.tomakehurst.wiremock.stubbing.Scenario;
54
import org.eclipse.dataplane.domain.Result;
65
import org.eclipse.dataplane.domain.dataflow.DataFlowPrepareMessage;
76
import org.eclipse.dataplane.port.exception.DataFlowNotFoundException;
7+
import org.eclipse.dataplane.port.exception.DataFlowNotifyCompletedFailed;
88
import org.eclipse.dataplane.port.exception.DataplaneNotRegistered;
99
import org.junit.jupiter.api.AfterEach;
1010
import org.junit.jupiter.api.BeforeEach;
1111
import org.junit.jupiter.api.Nested;
1212
import org.junit.jupiter.api.Test;
1313

14-
import java.util.concurrent.TimeUnit;
14+
import java.net.ConnectException;
1515

1616
import static com.github.tomakehurst.wiremock.client.WireMock.aResponse;
1717
import static com.github.tomakehurst.wiremock.client.WireMock.and;
@@ -52,7 +52,7 @@ void shouldFail_whenDataFlowDoesNotExist() {
5252

5353
var result = dataplane.notifyCompleted("dataFlowId");
5454

55-
assertThat(result.failed());
55+
assertThat(result.failed()).isTrue();
5656
assertThatThrownBy(result::orElseThrow).isExactlyInstanceOf(DataFlowNotFoundException.class);
5757
}
5858

@@ -64,8 +64,8 @@ void shouldReturnFailedFuture_whenControlPlaneIsNotAvailable() {
6464

6565
var result = dataplane.notifyCompleted("dataFlowId");
6666

67-
assertThat(result.succeeded());
68-
assertThat(result.getContent()).failsWithin(5, TimeUnit.SECONDS);
67+
assertThat(result.failed()).isTrue();
68+
assertThatThrownBy(result::orElseThrow).isExactlyInstanceOf(ConnectException.class);
6969
}
7070

7171
@Test
@@ -77,8 +77,8 @@ void shouldReturnFailedFuture_whenControlPlaneRespondWithError() {
7777

7878
var result = dataplane.notifyCompleted("dataFlowId");
7979

80-
assertThat(result.succeeded());
81-
assertThat(result.getContent()).failsWithin(5, TimeUnit.SECONDS);
80+
assertThat(result.failed()).isTrue();
81+
assertThatThrownBy(result::orElseThrow).isExactlyInstanceOf(DataFlowNotifyCompletedFailed.class);
8282
assertThat(dataplane.status("dataFlowId").getContent().state()).isNotEqualTo(COMPLETED.name());
8383
}
8484

@@ -90,29 +90,7 @@ void shouldTransitionToCompleted_whenControlPlaneRespondCorrectly() {
9090

9191
var result = dataplane.notifyCompleted("dataFlowId");
9292

93-
assertThat(result.succeeded());
94-
assertThat(result.getContent()).succeedsWithin(5, TimeUnit.SECONDS);
95-
assertThat(dataplane.status("dataFlowId").getContent().state()).isEqualTo(COMPLETED.name());
96-
}
97-
98-
@Test
99-
void shouldRetryForCertainAmountOfCalls() {
100-
controlPlane.stubFor(post(anyUrl()).inScenario("retry")
101-
.whenScenarioStateIs(Scenario.STARTED)
102-
.willReturn(aResponse().withStatus(500))
103-
.willSetStateTo("RETRY"));
104-
105-
controlPlane.stubFor(post(anyUrl()).inScenario("retry")
106-
.whenScenarioStateIs("RETRY")
107-
.willReturn(aResponse().withStatus(200)));
108-
109-
var dataplane = Dataplane.newInstance().onPrepare(Result::success).build();
110-
dataplane.prepare(createPrepareMessage());
111-
112-
var result = dataplane.notifyCompleted("dataFlowId");
113-
114-
assertThat(result.succeeded());
115-
assertThat(result.getContent()).succeedsWithin(5, TimeUnit.SECONDS);
93+
assertThat(result.succeeded()).isTrue();
11694
assertThat(dataplane.status("dataFlowId").getContent().state()).isEqualTo(COMPLETED.name());
11795
}
11896

@@ -121,8 +99,6 @@ private DataFlowPrepareMessage createPrepareMessage() {
12199
controlPlane.baseUrl(), "Something-PUSH", emptyList(), emptyMap());
122100
}
123101

124-
// TODO: retry case
125-
126102
}
127103

128104
@Nested

0 commit comments

Comments
 (0)