Merge branch 'stable-3.3' into stable-3.4 * stable-3.3: (21 commits) Fix the topic events replay Kafka REST-API Use Kafka REST Proxy id to subscribe to the correct instance Fix Kafka REST Proxy accepts header for topic meta-data Kafka REST Client: avoid clashes between clients Fix threshold of HTTP wire logging Delete subscription at the end of ReceiverJob Update kafka-client 2.1.0 -> 2.1.1 Increase patience to 30s for shouldReplayAllEvents test Remove unused RequestConfigProvider REST ClientType: Make thread pool and timeouts configuration Extract configuration properties into constants Manage Kafka clientType when starting session Receive messages through Kafka REST API Send messages through Kafka REST API Abstract Publisher/Subscriber into generic interfaces Wait at most for 5s for an empty topic Assert that messages are acknowledged in KafkaBrokerApiTest Add Kafka REST-API container in test Remove access to deprecated poll(long) method Use explicit Kafka image:tag in tests Do not connect KafkaSession without bootstrap servers Change-Id: I747ced0e436d5f544fcc71083a5dd5f6d7a3bb52
diff --git a/external_plugin_deps.bzl b/external_plugin_deps.bzl index 794fc10..88a49ab 100644 --- a/external_plugin_deps.bzl +++ b/external_plugin_deps.bzl
@@ -15,6 +15,6 @@ maven_jar( name = "events-broker", - artifact = "com.gerritforge:events-broker:3.3.2", - sha1 = "d8bcb77047cc12dd7c623b5b4de70a25499d3d6c", + artifact = "com.gerritforge:events-broker:3.4.0.4", + sha1 = "8d361d863382290e33828116e65698190118d0f1", )
diff --git a/src/main/java/com/googlesource/gerrit/plugins/kafka/Module.java b/src/main/java/com/googlesource/gerrit/plugins/kafka/Module.java index 463f1c9..57ca564 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/kafka/Module.java +++ b/src/main/java/com/googlesource/gerrit/plugins/kafka/Module.java
@@ -14,16 +14,13 @@ package com.googlesource.gerrit.plugins.kafka; -import com.gerritforge.gerrit.eventbroker.EventGsonProvider; import com.google.gerrit.extensions.events.LifecycleListener; import com.google.gerrit.extensions.registration.DynamicSet; import com.google.gerrit.server.events.EventListener; import com.google.gerrit.server.git.WorkQueue; -import com.google.gson.Gson; import com.google.inject.AbstractModule; import com.google.inject.Inject; import com.google.inject.Scopes; -import com.google.inject.Singleton; import com.google.inject.TypeLiteral; import com.google.inject.assistedinject.FactoryModuleBuilder; import com.googlesource.gerrit.plugins.kafka.api.KafkaApiModule; @@ -53,7 +50,6 @@ @Override protected void configure() { - bind(Gson.class).toProvider(EventGsonProvider.class).in(Singleton.class); DynamicSet.bind(binder(), LifecycleListener.class).to(Manager.class); DynamicSet.bind(binder(), EventListener.class).to(KafkaPublisher.class);
diff --git a/src/main/java/com/googlesource/gerrit/plugins/kafka/api/KafkaApiModule.java b/src/main/java/com/googlesource/gerrit/plugins/kafka/api/KafkaApiModule.java index 47f9969..ca6c45d 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/kafka/api/KafkaApiModule.java +++ b/src/main/java/com/googlesource/gerrit/plugins/kafka/api/KafkaApiModule.java
@@ -15,11 +15,11 @@ package com.googlesource.gerrit.plugins.kafka.api; import com.gerritforge.gerrit.eventbroker.BrokerApi; -import com.gerritforge.gerrit.eventbroker.EventMessage; import com.gerritforge.gerrit.eventbroker.TopicSubscriber; import com.google.common.collect.Sets; import com.google.gerrit.extensions.registration.DynamicItem; import com.google.gerrit.lifecycle.LifecycleModule; +import com.google.gerrit.server.events.Event; import com.google.gerrit.server.git.WorkQueue; import com.google.inject.Inject; import com.google.inject.Scopes; @@ -76,7 +76,7 @@ workQueue.createQueue(configuration.getNumberOfSubscribers(), "kafka-subscriber")); bind(new TypeLiteral<Deserializer<byte[]>>() {}).toInstance(new ByteArrayDeserializer()); - bind(new TypeLiteral<Deserializer<EventMessage>>() {}).to(KafkaEventDeserializer.class); + bind(new TypeLiteral<Deserializer<Event>>() {}).to(KafkaEventDeserializer.class); bind(new TypeLiteral<Set<TopicSubscriber>>() {}).toInstance(activeConsumers); DynamicItem.bind(binder(), BrokerApi.class).to(KafkaBrokerApi.class).in(Scopes.SINGLETON);
diff --git a/src/main/java/com/googlesource/gerrit/plugins/kafka/api/KafkaBrokerApi.java b/src/main/java/com/googlesource/gerrit/plugins/kafka/api/KafkaBrokerApi.java index 9a7c66a..3ec21e0 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/kafka/api/KafkaBrokerApi.java +++ b/src/main/java/com/googlesource/gerrit/plugins/kafka/api/KafkaBrokerApi.java
@@ -15,8 +15,9 @@ package com.googlesource.gerrit.plugins.kafka.api; import com.gerritforge.gerrit.eventbroker.BrokerApi; -import com.gerritforge.gerrit.eventbroker.EventMessage; import com.gerritforge.gerrit.eventbroker.TopicSubscriber; +import com.google.common.util.concurrent.ListenableFuture; +import com.google.gerrit.server.events.Event; import com.google.inject.Inject; import com.google.inject.Provider; import com.googlesource.gerrit.plugins.kafka.publish.KafkaPublisher; @@ -42,12 +43,12 @@ } @Override - public boolean send(String topic, EventMessage event) { + public ListenableFuture<Boolean> send(String topic, Event event) { return publisher.publish(topic, event); } @Override - public void receiveAsync(String topic, Consumer<EventMessage> eventConsumer) { + public void receiveAsync(String topic, Consumer<Event> eventConsumer) { KafkaEventSubscriber subscriber = subscriberProvider.get(); synchronized (subscribers) { subscribers.add(subscriber);
diff --git a/src/main/java/com/googlesource/gerrit/plugins/kafka/publish/GsonProvider.java b/src/main/java/com/googlesource/gerrit/plugins/kafka/publish/GsonProvider.java deleted file mode 100644 index 2c5c1e7..0000000 --- a/src/main/java/com/googlesource/gerrit/plugins/kafka/publish/GsonProvider.java +++ /dev/null
@@ -1,29 +0,0 @@ -// Copyright (C) 2016 The Android Open Source Project -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -package com.googlesource.gerrit.plugins.kafka.publish; - -import com.google.common.base.Supplier; -import com.google.gerrit.server.events.SupplierSerializer; -import com.google.gson.Gson; -import com.google.gson.GsonBuilder; -import com.google.inject.Provider; - -public class GsonProvider implements Provider<Gson> { - - @Override - public Gson get() { - return new GsonBuilder().registerTypeAdapter(Supplier.class, new SupplierSerializer()).create(); - } -}
diff --git a/src/main/java/com/googlesource/gerrit/plugins/kafka/publish/KafkaPublisher.java b/src/main/java/com/googlesource/gerrit/plugins/kafka/publish/KafkaPublisher.java index cc271b5..e7670cb 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/kafka/publish/KafkaPublisher.java +++ b/src/main/java/com/googlesource/gerrit/plugins/kafka/publish/KafkaPublisher.java
@@ -14,9 +14,10 @@ package com.googlesource.gerrit.plugins.kafka.publish; -import com.gerritforge.gerrit.eventbroker.EventMessage; import com.google.common.annotations.VisibleForTesting; +import com.google.common.util.concurrent.ListenableFuture; import com.google.gerrit.server.events.Event; +import com.google.gerrit.server.events.EventGson; import com.google.gerrit.server.events.EventListener; import com.google.gson.Gson; import com.google.gson.JsonObject; @@ -31,7 +32,7 @@ private final Gson gson; @Inject - public KafkaPublisher(KafkaSession kafkaSession, Gson gson) { + public KafkaPublisher(KafkaSession kafkaSession, @EventGson Gson gson) { this.session = kafkaSession; this.gson = gson; } @@ -53,11 +54,11 @@ } } - public boolean publish(String topic, EventMessage event) { + public ListenableFuture<Boolean> publish(String topic, Event event) { return session.publish(topic, getPayload(event)); } - private String getPayload(EventMessage event) { + private String getPayload(Event event) { return gson.toJson(event); }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/kafka/session/KafkaSession.java b/src/main/java/com/googlesource/gerrit/plugins/kafka/session/KafkaSession.java index 69224e7..a0df313 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/kafka/session/KafkaSession.java +++ b/src/main/java/com/googlesource/gerrit/plugins/kafka/session/KafkaSession.java
@@ -14,12 +14,18 @@ package com.googlesource.gerrit.plugins.kafka.session; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.JdkFutureAdapters; +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.MoreExecutors; +import com.google.common.util.concurrent.SettableFuture; import com.google.inject.Inject; import com.google.inject.Provider; import com.googlesource.gerrit.plugins.kafka.config.KafkaProperties; import com.googlesource.gerrit.plugins.kafka.publish.KafkaEventsPublisherMetrics; import java.net.URI; import java.net.URISyntaxException; +import java.util.Objects; import java.util.concurrent.Future; import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerRecord; @@ -112,34 +118,35 @@ producer = null; } - public void publish(String messageBody) { - publish(properties.getTopic(), messageBody); + public ListenableFuture<Boolean> publish(String messageBody) { + return publish(properties.getTopic(), messageBody); } - public boolean publish(String topic, String messageBody) { + public ListenableFuture<Boolean> publish(String topic, String messageBody) { if (properties.isSendAsync()) { return publishAsync(topic, messageBody); } return publishSync(topic, messageBody); } - private boolean publishSync(String topic, String messageBody) { - + private ListenableFuture<Boolean> publishSync(String topic, String messageBody) { + SettableFuture<Boolean> resultF = SettableFuture.create(); try { Future<RecordMetadata> future = producer.send(new ProducerRecord<>(topic, "" + System.nanoTime(), messageBody)); RecordMetadata metadata = future.get(); LOGGER.debug("The offset of the record we just sent is: {}", metadata.offset()); publisherMetrics.incrementBrokerPublishedMessage(); - return true; + resultF.set(true); + return resultF; } catch (Throwable e) { LOGGER.error("Cannot send the message", e); publisherMetrics.incrementBrokerFailedToPublishMessage(); - return false; + return Futures.immediateFailedFuture(e); } } - private boolean publishAsync(String topic, String messageBody) { + private ListenableFuture<Boolean> publishAsync(String topic, String messageBody) { try { Future<RecordMetadata> future = producer.send( @@ -153,11 +160,16 @@ publisherMetrics.incrementBrokerFailedToPublishMessage(); } }); - return future != null; + + // The transformation is lightweight, so we can afford using a directExecutor + return Futures.transform( + JdkFutureAdapters.listenInPoolThread(future), + Objects::nonNull, + MoreExecutors.directExecutor()); } catch (Throwable e) { LOGGER.error("Cannot send the message", e); publisherMetrics.incrementBrokerFailedToPublishMessage(); - return false; + return Futures.immediateFailedFuture(e); } } }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventDeserializer.java b/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventDeserializer.java index 4c57a54..cad2f37 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventDeserializer.java +++ b/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventDeserializer.java
@@ -14,54 +14,36 @@ package com.googlesource.gerrit.plugins.kafka.subscribe; -import static java.util.Objects.requireNonNull; - -import com.gerritforge.gerrit.eventbroker.EventMessage; -import com.gerritforge.gerrit.eventbroker.EventMessage.Header; +import com.gerritforge.gerrit.eventbroker.EventDeserializer; import com.google.gerrit.server.events.Event; -import com.google.gson.Gson; import com.google.inject.Inject; import com.google.inject.Singleton; import java.util.Map; -import java.util.UUID; import org.apache.kafka.common.serialization.Deserializer; import org.apache.kafka.common.serialization.StringDeserializer; @Singleton -public class KafkaEventDeserializer implements Deserializer<EventMessage> { +public class KafkaEventDeserializer implements Deserializer<Event> { private final StringDeserializer stringDeserializer = new StringDeserializer(); - private Gson gson; + private EventDeserializer eventDeserializer; // To be used when providing this deserializer with class name (then need to add a configuration // entry to set the gson.provider public KafkaEventDeserializer() {} @Inject - public KafkaEventDeserializer(Gson gson) { - this.gson = gson; + public KafkaEventDeserializer(EventDeserializer eventDeserializer) { + this.eventDeserializer = eventDeserializer; } @Override public void configure(Map<String, ?> configs, boolean isKey) {} @Override - public EventMessage deserialize(String topic, byte[] data) { + public Event deserialize(String topic, byte[] data) { String json = stringDeserializer.deserialize(topic, data); - EventMessage result = gson.fromJson(json, EventMessage.class); - if (result.getEvent() == null && result.getHeader() == null) { - Event event = deserialiseEvent(json); - result = new EventMessage(new Header(UUID.randomUUID(), event.instanceId), event); - } - result.validate(); - return result; - } - - private Event deserialiseEvent(String json) { - Event event = gson.fromJson(json, Event.class); - requireNonNull(event.type, "Event type cannot be null"); - requireNonNull(event.instanceId, "Event instance id cannot be null"); - return event; + return eventDeserializer.deserialize(json); } @Override
diff --git a/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventNativeSubscriber.java b/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventNativeSubscriber.java index a98e098..24fe566 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventNativeSubscriber.java +++ b/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventNativeSubscriber.java
@@ -15,8 +15,8 @@ import static java.nio.charset.StandardCharsets.UTF_8; -import com.gerritforge.gerrit.eventbroker.EventMessage; import com.google.common.flogger.FluentLogger; +import com.google.gerrit.server.events.Event; import com.google.gerrit.server.util.ManualRequestContext; import com.google.gerrit.server.util.OneOffRequestContext; import com.google.inject.Inject; @@ -39,14 +39,14 @@ private final OneOffRequestContext oneOffCtx; private final AtomicBoolean closed = new AtomicBoolean(false); - private final Deserializer<EventMessage> valueDeserializer; + private final Deserializer<Event> valueDeserializer; private final KafkaSubscriberProperties configuration; private final ExecutorService executor; private final KafkaEventSubscriberMetrics subscriberMetrics; private final KafkaConsumerFactory consumerFactory; private final Deserializer<byte[]> keyDeserializer; - private java.util.function.Consumer<EventMessage> messageProcessor; + private java.util.function.Consumer<Event> messageProcessor; private String topic; private AtomicBoolean resetOffset = new AtomicBoolean(false); @@ -57,7 +57,7 @@ KafkaSubscriberProperties configuration, KafkaConsumerFactory consumerFactory, Deserializer<byte[]> keyDeserializer, - Deserializer<EventMessage> valueDeserializer, + Deserializer<Event> valueDeserializer, OneOffRequestContext oneOffCtx, @ConsumerExecutor ExecutorService executor, KafkaEventSubscriberMetrics subscriberMetrics) { @@ -75,7 +75,7 @@ * @see com.googlesource.gerrit.plugins.kafka.subscribe.KafkaEventSubscriber#subscribe(java.lang.String, java.util.function.Consumer) */ @Override - public void subscribe(String topic, java.util.function.Consumer<EventMessage> messageProcessor) { + public void subscribe(String topic, java.util.function.Consumer<Event> messageProcessor) { this.topic = topic; this.messageProcessor = messageProcessor; logger.atInfo().log( @@ -110,7 +110,7 @@ * @see com.googlesource.gerrit.plugins.kafka.subscribe.KafkaEventSubscriber#getMessageProcessor() */ @Override - public java.util.function.Consumer<EventMessage> getMessageProcessor() { + public java.util.function.Consumer<Event> getMessageProcessor() { return messageProcessor; } @@ -167,7 +167,7 @@ consumerRecords.forEach( consumerRecord -> { try (ManualRequestContext ctx = oneOffCtx.open()) { - EventMessage event = + Event event = valueDeserializer.deserialize(consumerRecord.topic(), consumerRecord.value()); messageProcessor.accept(event); } catch (Exception e) {
diff --git a/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventRestSubscriber.java b/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventRestSubscriber.java index 5ec5f95..7429991 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventRestSubscriber.java +++ b/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventRestSubscriber.java
@@ -15,10 +15,10 @@ import static java.nio.charset.StandardCharsets.UTF_8; -import com.gerritforge.gerrit.eventbroker.EventMessage; import com.google.common.flogger.FluentLogger; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; +import com.google.gerrit.server.events.Event; import com.google.gerrit.server.util.ManualRequestContext; import com.google.gerrit.server.util.OneOffRequestContext; import com.google.gson.Gson; @@ -71,13 +71,13 @@ private final OneOffRequestContext oneOffCtx; private final AtomicBoolean closed = new AtomicBoolean(false); - private final Deserializer<EventMessage> valueDeserializer; + private final Deserializer<Event> valueDeserializer; private final KafkaSubscriberProperties configuration; private final ExecutorService executor; private final KafkaEventSubscriberMetrics subscriberMetrics; private final Gson gson; - private java.util.function.Consumer<EventMessage> messageProcessor; + private java.util.function.Consumer<Event> messageProcessor; private String topic; private final KafkaRestClient restClient; private final AtomicBoolean resetOffset; @@ -87,7 +87,7 @@ @Inject public KafkaEventRestSubscriber( KafkaSubscriberProperties configuration, - Deserializer<EventMessage> valueDeserializer, + Deserializer<Event> valueDeserializer, OneOffRequestContext oneOffCtx, @ConsumerExecutor ExecutorService executor, KafkaEventSubscriberMetrics subscriberMetrics, @@ -109,7 +109,7 @@ * @see com.googlesource.gerrit.plugins.kafka.subscribe.KafkaEventSubscriber#subscribe(java.lang.String, java.util.function.Consumer) */ @Override - public void subscribe(String topic, java.util.function.Consumer<EventMessage> messageProcessor) { + public void subscribe(String topic, java.util.function.Consumer<Event> messageProcessor) { this.topic = topic; this.messageProcessor = messageProcessor; logger.atInfo().log( @@ -143,7 +143,7 @@ * @see com.googlesource.gerrit.plugins.kafka.subscribe.KafkaEventSubscriber#getMessageProcessor() */ @Override - public java.util.function.Consumer<EventMessage> getMessageProcessor() { + public java.util.function.Consumer<Event> getMessageProcessor() { return messageProcessor; } @@ -207,7 +207,7 @@ records.forEach( consumerRecord -> { try (ManualRequestContext ctx = oneOffCtx.open()) { - EventMessage event = + Event event = valueDeserializer.deserialize(consumerRecord.topic(), consumerRecord.value()); messageProcessor.accept(event); } catch (Exception e) {
diff --git a/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventSubscriber.java b/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventSubscriber.java index 6315dea..34c64b2 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventSubscriber.java +++ b/src/main/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventSubscriber.java
@@ -14,7 +14,7 @@ package com.googlesource.gerrit.plugins.kafka.subscribe; -import com.gerritforge.gerrit.eventbroker.EventMessage; +import com.google.gerrit.server.events.Event; /** Generic interface to a Kafka topic subscriber. */ public interface KafkaEventSubscriber { @@ -25,7 +25,7 @@ * @param topic Kafka topic name * @param messageProcessor consumer function for processing incoming messages */ - void subscribe(String topic, java.util.function.Consumer<EventMessage> messageProcessor); + void subscribe(String topic, java.util.function.Consumer<Event> messageProcessor); /** Shutdown Kafka consumer. */ void shutdown(); @@ -35,7 +35,7 @@ * * @return the default topic consumer function. */ - java.util.function.Consumer<EventMessage> getMessageProcessor(); + java.util.function.Consumer<Event> getMessageProcessor(); /** * Returns the current subscribed topic name.
diff --git a/src/test/java/com/googlesource/gerrit/plugins/kafka/EventConsumerIT.java b/src/test/java/com/googlesource/gerrit/plugins/kafka/EventConsumerIT.java index 85945c0..70a664b 100644 --- a/src/test/java/com/googlesource/gerrit/plugins/kafka/EventConsumerIT.java +++ b/src/test/java/com/googlesource/gerrit/plugins/kafka/EventConsumerIT.java
@@ -19,8 +19,6 @@ import static org.junit.Assert.fail; import com.gerritforge.gerrit.eventbroker.BrokerApi; -import com.gerritforge.gerrit.eventbroker.EventGsonProvider; -import com.gerritforge.gerrit.eventbroker.EventMessage; import com.google.common.base.Stopwatch; import com.google.common.collect.Iterables; import com.google.gerrit.acceptance.LightweightPluginDaemonTest; @@ -33,6 +31,7 @@ import com.google.gerrit.extensions.common.ChangeMessageInfo; import com.google.gerrit.server.events.CommentAddedEvent; import com.google.gerrit.server.events.Event; +import com.google.gerrit.server.events.EventGsonProvider; import com.google.gerrit.server.events.ProjectCreatedEvent; import com.google.gson.Gson; import com.googlesource.gerrit.plugins.kafka.config.KafkaProperties; @@ -40,7 +39,6 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; -import java.util.UUID; import java.util.function.Supplier; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; @@ -141,14 +139,12 @@ @GerritConfig(name = "plugin.events-kafka.pollingIntervalMs", value = "500") public void shouldReplayAllEvents() throws InterruptedException { String topic = "a_topic"; - EventMessage eventMessage = - new EventMessage( - new EventMessage.Header(UUID.randomUUID(), UUID.randomUUID()), - new ProjectCreatedEvent()); + Event eventMessage = new ProjectCreatedEvent(); + eventMessage.instanceId = "test-instance-id"; Duration WAIT_FOR_POLL_TIMEOUT = Duration.ofSeconds(30); - List<EventMessage> receivedEvents = new ArrayList<>(); + List<Event> receivedEvents = new ArrayList<>(); BrokerApi kafkaBrokerApi = kafkaBrokerApi(); kafkaBrokerApi.send(topic, eventMessage); @@ -157,14 +153,12 @@ waitUntil(() -> receivedEvents.size() == 1, WAIT_FOR_POLL_TIMEOUT); - assertThat(receivedEvents.get(0).getHeader().eventId) - .isEqualTo(eventMessage.getHeader().eventId); + assertThat(receivedEvents.get(0).instanceId).isEqualTo(eventMessage.instanceId); kafkaBrokerApi.replayAllEvents(topic); waitUntil(() -> receivedEvents.size() == 2, WAIT_FOR_POLL_TIMEOUT); - assertThat(receivedEvents.get(1).getHeader().eventId) - .isEqualTo(eventMessage.getHeader().eventId); + assertThat(receivedEvents.get(1).instanceId).isEqualTo(eventMessage.instanceId); } private BrokerApi kafkaBrokerApi() {
diff --git a/src/test/java/com/googlesource/gerrit/plugins/kafka/api/KafkaBrokerApiTest.java b/src/test/java/com/googlesource/gerrit/plugins/kafka/api/KafkaBrokerApiTest.java index 2cbc65a..37f15da 100644 --- a/src/test/java/com/googlesource/gerrit/plugins/kafka/api/KafkaBrokerApiTest.java +++ b/src/test/java/com/googlesource/gerrit/plugins/kafka/api/KafkaBrokerApiTest.java
@@ -17,10 +17,10 @@ import static com.google.common.truth.Truth.assertThat; import static org.mockito.Mockito.mock; -import com.gerritforge.gerrit.eventbroker.EventGsonProvider; -import com.gerritforge.gerrit.eventbroker.EventMessage; -import com.gerritforge.gerrit.eventbroker.EventMessage.Header; import com.google.gerrit.metrics.MetricMaker; +import com.google.gerrit.server.events.Event; +import com.google.gerrit.server.events.EventGson; +import com.google.gerrit.server.events.EventGsonProvider; import com.google.gerrit.server.events.ProjectCreatedEvent; import com.google.gerrit.server.git.WorkQueue; import com.google.gerrit.server.util.IdGenerator; @@ -42,7 +42,6 @@ import com.googlesource.gerrit.plugins.kafka.session.KafkaSession; import java.util.ArrayList; import java.util.List; -import java.util.UUID; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.function.Consumer; @@ -72,7 +71,7 @@ static final int TEST_POLLING_INTERVAL_MSEC = 100; static final String KAFKA_REST_ID = "kafka-rest-instance-0"; private static final int TEST_THREAD_POOL_SIZE = 10; - private static final UUID TEST_INSTANCE_ID = UUID.randomUUID(); + private static final String TEST_INSTANCE_ID = "test-instance-id"; private static final TimeUnit TEST_TIMEOUT_UNIT = TimeUnit.SECONDS; private static final int TEST_TIMEOUT = 30; private static final int TEST_WAIT_FOR_MORE_MESSAGES_TIMEOUT = 5; @@ -99,7 +98,10 @@ @Override protected void configure() { - bind(Gson.class).toProvider(EventGsonProvider.class).in(Singleton.class); + bind(Gson.class) + .annotatedWith(EventGson.class) + .toProvider(EventGsonProvider.class) + .in(Singleton.class); bind(MetricMaker.class).toInstance(mock(MetricMaker.class, Answers.RETURNS_DEEP_STUBS)); bind(OneOffRequestContext.class) .toInstance(mock(OneOffRequestContext.class, Answers.RETURNS_DEEP_STUBS)); @@ -121,8 +123,8 @@ } } - public static class TestConsumer implements Consumer<EventMessage> { - public final List<EventMessage> messages = new ArrayList<>(); + public static class TestConsumer implements Consumer<Event> { + public final List<Event> messages = new ArrayList<>(); private CountDownLatch[] locks; public TestConsumer(int numMessagesExpected) { @@ -137,7 +139,7 @@ } @Override - public void accept(EventMessage message) { + public void accept(Event message) { messages.add(message); for (CountDownLatch countDownLatch : locks) { countDownLatch.countDown(); @@ -161,13 +163,6 @@ } } - public static class TestHeader extends Header { - - public TestHeader() { - super(UUID.randomUUID(), TEST_INSTANCE_ID); - } - } - @BeforeClass public static void beforeClass() throws Exception { kafka = KafkaContainerProvider.get(); @@ -230,7 +225,8 @@ KafkaBrokerApi kafkaBrokerApi = injector.getInstance(KafkaBrokerApi.class); String testTopic = "test_topic_sync"; TestConsumer testConsumer = new TestConsumer(1); - EventMessage testEventMessage = new EventMessage(new TestHeader(), new ProjectCreatedEvent()); + Event testEventMessage = new ProjectCreatedEvent(); + testEventMessage.instanceId = TEST_INSTANCE_ID; kafkaBrokerApi.receiveAsync(testTopic, testConsumer); kafkaBrokerApi.send(testTopic, testEventMessage); @@ -248,7 +244,8 @@ KafkaBrokerApi kafkaBrokerApi = injector.getInstance(KafkaBrokerApi.class); String testTopic = "test_topic_async"; TestConsumer testConsumer = new TestConsumer(1); - EventMessage testEventMessage = new EventMessage(new TestHeader(), new ProjectCreatedEvent()); + Event testEventMessage = new ProjectCreatedEvent(); + testEventMessage.instanceId = TEST_INSTANCE_ID; kafkaBrokerApi.send(testTopic, testEventMessage); kafkaBrokerApi.receiveAsync(testTopic, testConsumer); @@ -265,7 +262,7 @@ connectToKafka(new KafkaProperties(false, clientType, getKafkaRestApiUriString())); KafkaBrokerApi kafkaBrokerApi = injector.getInstance(KafkaBrokerApi.class); String testTopic = "test_topic_reset"; - EventMessage testEventMessage = new EventMessage(new TestHeader(), new ProjectCreatedEvent()); + Event testEventMessage = new ProjectCreatedEvent(); TestConsumer testConsumer = new TestConsumer(2); kafkaBrokerApi.receiveAsync(testTopic, testConsumer);
diff --git a/src/test/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventDeserializerTest.java b/src/test/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventDeserializerTest.java deleted file mode 100644 index 4074919..0000000 --- a/src/test/java/com/googlesource/gerrit/plugins/kafka/subscribe/KafkaEventDeserializerTest.java +++ /dev/null
@@ -1,81 +0,0 @@ -// Copyright (C) 2019 The Android Open Source Project -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -package com.googlesource.gerrit.plugins.kafka.subscribe; - -import static com.google.common.truth.Truth.assertThat; -import static java.nio.charset.StandardCharsets.UTF_8; - -import com.gerritforge.gerrit.eventbroker.EventGsonProvider; -import com.gerritforge.gerrit.eventbroker.EventMessage; -import com.google.gson.Gson; -import java.util.UUID; -import org.junit.Before; -import org.junit.Test; - -public class KafkaEventDeserializerTest { - private KafkaEventDeserializer deserializer; - - @Before - public void setUp() { - final Gson gson = new EventGsonProvider().get(); - deserializer = new KafkaEventDeserializer(gson); - } - - @Test - public void kafkaEventDeserializerShouldParseAKafkaEventMessage() { - final UUID eventId = UUID.randomUUID(); - final String eventType = "event-type"; - final String sourceInstanceId = UUID.randomUUID().toString(); - final long eventCreatedOn = 10L; - final String eventJson = - String.format( - "{ " - + "\"header\": { \"eventId\": \"%s\", \"eventType\": \"%s\", \"sourceInstanceId\": \"%s\", \"eventCreatedOn\": %d }," - + "\"body\": { \"type\": \"project-created\" }" - + "}", - eventId, eventType, sourceInstanceId, eventCreatedOn); - final EventMessage event = deserializer.deserialize("ignored", eventJson.getBytes(UTF_8)); - - assertThat(event.getHeader().eventId).isEqualTo(eventId); - assertThat(event.getHeader().sourceInstanceId).isEqualTo(sourceInstanceId); - } - - @Test - public void kafkaEventDeserializerShouldParseKafkaEvent() { - final String eventJson = "{ \"type\": \"project-created\", \"instanceId\":\"instance-id\" }"; - final EventMessage event = deserializer.deserialize("ignored", eventJson.getBytes(UTF_8)); - - assertThat(event.getHeader().sourceInstanceId).isEqualTo("instance-id"); - } - - @Test - public void kafkaEventDeserializerShouldParseKafkaEventWithHeaderAndBodyProjectName() { - final String eventJson = - "{\"projectName\":\"header_body_parser_project\",\"type\":\"project-created\", \"instanceId\":\"instance-id\"}"; - final EventMessage event = deserializer.deserialize("ignored", eventJson.getBytes(UTF_8)); - - assertThat(event.getHeader().sourceInstanceId).isEqualTo("instance-id"); - } - - @Test(expected = RuntimeException.class) - public void kafkaEventDeserializerShouldFailForInvalidJson() { - deserializer.deserialize("ignored", "this is not a JSON string".getBytes(UTF_8)); - } - - @Test(expected = RuntimeException.class) - public void kafkaEventDeserializerShouldFailForInvalidObjectButValidJSON() { - deserializer.deserialize("ignored", "{}".getBytes(UTF_8)); - } -}