|
22 | 22 | import java.net.http.HttpResponse;
|
23 | 23 | import java.util.Map;
|
24 | 24 | import java.util.concurrent.CompletableFuture;
|
| 25 | +import org.jspecify.annotations.NonNull; |
25 | 26 |
|
26 | 27 | public class DefaultIngressClient implements IngressClient {
|
27 | 28 |
|
@@ -92,6 +93,75 @@ public <Req> CompletableFuture<String> sendAsync(
|
92 | 93 | });
|
93 | 94 | }
|
94 | 95 |
|
| 96 | + @Override |
| 97 | + public AwakeableHandle awakeableHandle(String id) { |
| 98 | + return new AwakeableHandle() { |
| 99 | + @Override |
| 100 | + public <T> CompletableFuture<Void> resolve(Serde<T> serde, @NonNull T payload) { |
| 101 | + // Prepare request |
| 102 | + var reqBuilder = |
| 103 | + HttpRequest.newBuilder().uri(URI.create("/restate/awakeables/" + id + "/resolve")); |
| 104 | + |
| 105 | + // Add content-type |
| 106 | + if (serde.contentType() != null) { |
| 107 | + reqBuilder.header("content-type", serde.contentType()); |
| 108 | + } |
| 109 | + |
| 110 | + // Add headers |
| 111 | + headers.forEach(reqBuilder::header); |
| 112 | + |
| 113 | + // Build and Send request |
| 114 | + HttpRequest request = |
| 115 | + reqBuilder |
| 116 | + .POST(HttpRequest.BodyPublishers.ofByteArray(serde.serialize(payload))) |
| 117 | + .build(); |
| 118 | + return httpClient |
| 119 | + .sendAsync(request, HttpResponse.BodyHandlers.ofByteArray()) |
| 120 | + .handle( |
| 121 | + (response, throwable) -> { |
| 122 | + if (throwable != null) { |
| 123 | + throw new IngressException("Error when executing the request", throwable); |
| 124 | + } |
| 125 | + |
| 126 | + if (response.statusCode() >= 300) { |
| 127 | + handleNonSuccessResponse(response); |
| 128 | + } |
| 129 | + |
| 130 | + return null; |
| 131 | + }); |
| 132 | + } |
| 133 | + |
| 134 | + @Override |
| 135 | + public CompletableFuture<Void> reject(String reason) { |
| 136 | + // Prepare request |
| 137 | + var reqBuilder = |
| 138 | + HttpRequest.newBuilder() |
| 139 | + .uri(URI.create("/restate/awakeables/" + id + "/reject")) |
| 140 | + .header("content-type", "text-plain"); |
| 141 | + |
| 142 | + // Add headers |
| 143 | + headers.forEach(reqBuilder::header); |
| 144 | + |
| 145 | + // Build and Send request |
| 146 | + HttpRequest request = reqBuilder.POST(HttpRequest.BodyPublishers.ofString(reason)).build(); |
| 147 | + return httpClient |
| 148 | + .sendAsync(request, HttpResponse.BodyHandlers.ofByteArray()) |
| 149 | + .handle( |
| 150 | + (response, throwable) -> { |
| 151 | + if (throwable != null) { |
| 152 | + throw new IngressException("Error when executing the request", throwable); |
| 153 | + } |
| 154 | + |
| 155 | + if (response.statusCode() >= 300) { |
| 156 | + handleNonSuccessResponse(response); |
| 157 | + } |
| 158 | + |
| 159 | + return null; |
| 160 | + }); |
| 161 | + } |
| 162 | + }; |
| 163 | + } |
| 164 | + |
95 | 165 | private URI toRequestURI(Target target, boolean isSend) {
|
96 | 166 | StringBuilder builder = new StringBuilder();
|
97 | 167 | builder.append("/").append(target.getComponent());
|
|
0 commit comments