Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,8 @@
import io.kubernetes.client.spring.extended.controller.factory.KubernetesControllerFactory;
import io.kubernetes.client.util.ClientBuilder;
import java.util.LinkedList;
import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;
import java.util.function.Function;
import jakarta.annotation.Resource;
import org.apache.commons.lang3.tuple.MutablePair;
Expand All @@ -61,6 +63,9 @@
@SpringBootTest(classes = {KubernetesReconcilerCreatorTest.App.class})
class KubernetesReconcilerCreatorTest {

private static final Semaphore REQUEST_ENQUEUED = new Semaphore(0);
private static final long TIMEOUT_SECONDS = 10;

@RegisterExtension
static WireMockExtension apiServer =
WireMockExtension.newInstance().options(WireMockConfiguration.options().port(8189)).build();
Expand Down Expand Up @@ -172,6 +177,7 @@ public CustomWorkQueueKeyFunc(WorkQueue<Request> workQueue) {
@Override
public Request apply(KubernetesObject item) {
workQueue.add(new Request("foo"));
REQUEST_ENQUEUED.release();
return null;
}
}
Expand All @@ -182,23 +188,27 @@ void simplePodController() throws InterruptedException {
assertThat(testReconciler).isNotNull();

sharedInformerFactory.startAllRegisteredInformers();

((DefaultSharedIndexInformer<V1Pod, V1PodList>) testReconciler.podInformer)
.handleDeltas(
new LinkedList<MutablePair<DeltaFIFO.DeltaType, KubernetesObject>>() {
{
add(
new MutablePair<>(
DeltaFIFO.DeltaType.Added,
new V1Pod().metadata(new V1ObjectMeta().namespace("a").name("b"))));
}
});

Thread.sleep(500);

WorkQueue<Request> workQueue = ((DefaultController) testController).getWorkQueue();
assertThat(workQueue.length()).isEqualTo(1);
assertThat(workQueue.get().getName()).isEqualTo("foo");
sharedInformerFactory.stopAllRegisteredInformers();
try {
((DefaultSharedIndexInformer<V1Pod, V1PodList>) testReconciler.podInformer)
.handleDeltas(
new LinkedList<MutablePair<DeltaFIFO.DeltaType, KubernetesObject>>() {
{
add(
new MutablePair<>(
DeltaFIFO.DeltaType.Added,
new V1Pod().metadata(new V1ObjectMeta().namespace("a").name("b"))));
}
});

assertThat(REQUEST_ENQUEUED.tryAcquire(TIMEOUT_SECONDS, TimeUnit.SECONDS))
.as("pod event added a request to the controller work queue")
.isTrue();

WorkQueue<Request> workQueue = ((DefaultController) testController).getWorkQueue();
assertThat(workQueue.length()).isEqualTo(1);
assertThat(workQueue.get().getName()).isEqualTo("foo");
} finally {
sharedInformerFactory.stopAllRegisteredInformers();
}
}
}