Skip to content

Commit 09d6aee

Browse files
committed
feat: extends replication to async client
1 parent 257d8c2 commit 09d6aee

2 files changed

Lines changed: 116 additions & 0 deletions

File tree

src/main/java/io/weaviate/client6/v1/api/cluster/WeaviateClusterClientAsync.java

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,14 +6,21 @@
66
import java.util.concurrent.CompletableFuture;
77
import java.util.function.Function;
88

9+
import io.weaviate.client6.v1.api.cluster.replication.CreateReplicationRequest;
10+
import io.weaviate.client6.v1.api.cluster.replication.Replication;
11+
import io.weaviate.client6.v1.api.cluster.replication.ReplicationType;
12+
import io.weaviate.client6.v1.api.cluster.replication.WeaviateReplicationClientAsync;
913
import io.weaviate.client6.v1.internal.ObjectBuilder;
1014
import io.weaviate.client6.v1.internal.rest.RestTransport;
1115

1216
public class WeaviateClusterClientAsync {
1317
private final RestTransport restTransport;
1418

19+
public final WeaviateReplicationClientAsync replication;
20+
1521
public WeaviateClusterClientAsync(RestTransport restTransport) {
1622
this.restTransport = restTransport;
23+
this.replication = new WeaviateReplicationClientAsync(restTransport);
1724
}
1825

1926
/**
@@ -54,4 +61,23 @@ public CompletableFuture<List<Node>> listNodes(Function<ListNodesRequest.Builder
5461
throws IOException {
5562
return this.restTransport.performRequestAsync(ListNodesRequest.of(fn), ListNodesRequest._ENDPOINT);
5663
}
64+
65+
/**
66+
* Start a replication operation for a collection's shard.
67+
*
68+
* @param collection Collection name.
69+
* @param shard Shard name.
70+
* @param sourceNode Node on which the shard currently resides.
71+
* @param targetNode Node onto which the files will be replicated.
72+
*/
73+
public CompletableFuture<Replication> replicate(
74+
String collection,
75+
String shard,
76+
String sourceNode,
77+
String targetNode,
78+
ReplicationType type) {
79+
return this.restTransport.performRequestAsync(
80+
new CreateReplicationRequest(collection, shard, sourceNode, targetNode, type),
81+
CreateReplicationRequest._ENDPOINT);
82+
}
5783
}
Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,90 @@
1+
package io.weaviate.client6.v1.api.cluster.replication;
2+
3+
import java.io.IOException;
4+
import java.util.List;
5+
import java.util.Optional;
6+
import java.util.UUID;
7+
import java.util.concurrent.CompletableFuture;
8+
import java.util.function.Function;
9+
10+
import io.weaviate.client6.v1.internal.ObjectBuilder;
11+
import io.weaviate.client6.v1.internal.rest.RestTransport;
12+
13+
public class WeaviateReplicationClientAsync {
14+
private final RestTransport restTransport;
15+
16+
public WeaviateReplicationClientAsync(RestTransport restTransport) {
17+
this.restTransport = restTransport;
18+
}
19+
20+
/**
21+
* Get information about a replication operation.
22+
*
23+
* @param uuid Replication UUID.
24+
*/
25+
public CompletableFuture<Optional<Replication>> get(UUID uuid) throws IOException {
26+
return this.restTransport.performRequestAsync(GetReplicationRequest.of(uuid), GetReplicationRequest._ENDPOINT);
27+
}
28+
29+
/**
30+
* Get information about a replication operation.
31+
*
32+
* @param uuid Replication UUID.
33+
* @param fn Lambda expression for optional parameters.
34+
*/
35+
public CompletableFuture<Optional<Replication>> get(UUID uuid,
36+
Function<GetReplicationRequest.Builder, ObjectBuilder<GetReplicationRequest>> fn) throws IOException {
37+
return this.restTransport.performRequestAsync(GetReplicationRequest.of(uuid, fn), GetReplicationRequest._ENDPOINT);
38+
}
39+
40+
/**
41+
* List all replication operations.
42+
*
43+
* @see WeaviateReplicationClientAsync#list(Function) for filtering replications
44+
* by
45+
* collection, shard, or target node.
46+
*/
47+
public CompletableFuture<List<Replication>> list()
48+
throws IOException {
49+
return this.restTransport.performRequestAsync(ListReplicationsRequest.of(), ListReplicationsRequest._ENDPOINT);
50+
}
51+
52+
/**
53+
* List all replication operations.
54+
*
55+
* @param fn Lambda expression for optional parameters.
56+
*/
57+
public CompletableFuture<List<Replication>> list(
58+
Function<ListReplicationsRequest.Builder, ObjectBuilder<ListReplicationsRequest>> fn)
59+
throws IOException {
60+
return this.restTransport.performRequestAsync(ListReplicationsRequest.of(fn), ListReplicationsRequest._ENDPOINT);
61+
}
62+
63+
/**
64+
* Cancel a replication operation.
65+
*
66+
* @param uuid Replication UUID.
67+
*/
68+
public CompletableFuture<Void> cancel(UUID uuid)
69+
throws IOException {
70+
return this.restTransport.performRequestAsync(new CancelReplicationRequest(uuid),
71+
CancelReplicationRequest._ENDPOINT);
72+
}
73+
74+
/**
75+
* Delete a replication operation.
76+
*
77+
* @param uuid Replication UUID.
78+
*/
79+
public CompletableFuture<Void> delete(UUID uuid)
80+
throws IOException {
81+
return this.restTransport.performRequestAsync(new DeleteReplicationRequest(uuid),
82+
DeleteReplicationRequest._ENDPOINT);
83+
}
84+
85+
/** Delete all replication operations. */
86+
public CompletableFuture<Void> deleteAll()
87+
throws IOException {
88+
return this.restTransport.performRequestAsync(null, DeleteAllReplicationsRequest._ENDPOINT);
89+
}
90+
}

0 commit comments

Comments
 (0)