Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -3,3 +3,4 @@
.ijwb
bazel-*
MODULE.bazel.lock
.bazelbsp
2 changes: 1 addition & 1 deletion src/consumer/Consumer.java
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@
public class Consumer extends AbstractService {

/** Version of the client library. Should be kept in sync with Go client version. */
public static final String CLIENT_VERSION = "1.2.4";
public static final String CLIENT_VERSION = "1.2.5";

private static final Logger LOGGER = LoggerFactory.getLogger(Consumer.class);

Expand Down
19 changes: 18 additions & 1 deletion src/consumer/ConsumerBuilder.java
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ public class ConsumerBuilder {
private @Nullable Integer capacityLimit;
private @Nullable DomainKeyName[] keyNames;
private @Nullable BooleanSupplier verbose;
private @Nullable String version;

/**
* Creates a new builder for a {@link Consumer} with default values. It will use the local IP
Expand Down Expand Up @@ -134,6 +135,16 @@ public ConsumerBuilder withVerboseLogging(BooleanSupplier verbose) {
return this;
}

/**
* Sets the consumer version.
*
* @return the builder with the version set
*/
public ConsumerBuilder withVersion(String version) {
this.version = version;
return this;
}

public Consumer build() {
var clock = this.clock;
if (clock == null) {
Expand All @@ -155,10 +166,16 @@ public Consumer build() {
if (verbose == null) {
verbose = () -> false;
}
var version = this.version;

Function<Stream<Grant>, JoinMessage> registerFactory =
(grants) -> {
var register = Messages.createRegisterBuilder(instance, service, grants);
if (version != null) {
var md = ClientMessage.Register.Metadata.newBuilder();
md.setVersion(version);
register.setMetadata(md);
}
if (capacityLimit == null && keyNames == null) {
return Messages.createRegister(register.build());
}
Expand All @@ -173,7 +190,7 @@ public Consumer build() {
Stream.of(keyNames).map(Object::toString).collect(Collectors.joining(", ")));
opts = opts.addAllNames(Stream.of(keyNames).map(DomainKeyName::toProto).toList());
}
register = register.setOptions(opts.build());
register.setOptions(opts.build());
return Messages.createRegister(register.build());
};

Expand Down
Loading