Skip to content
Open
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
5 changes: 3 additions & 2 deletions src/main/java/net/spy/memcached/CacheManager.java
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,6 @@
import java.text.SimpleDateFormat;
import java.util.AbstractMap;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Date;
import java.util.List;
import java.util.Map;
Expand Down Expand Up @@ -718,7 +717,7 @@ public void connectionEstablished(MemcachedNode node, int reconnectCount) {
}
}
};
cfb.setInitialObservers(Collections.singleton(observer));
cfb.addInitialObserver(observer);

int poolId = CacheManager.POOL_ID.getAndIncrement();
client = new ArcusClient[poolSize];
Expand All @@ -740,6 +739,8 @@ public void connectionEstablished(MemcachedNode node, int reconnectCount) {
}
client = null;
return;
} finally {
cfb.removeInitialObserver(observer);
}

try {
Expand Down
27 changes: 21 additions & 6 deletions src/main/java/net/spy/memcached/ConnectionFactoryBuilder.java
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,8 @@

import java.io.IOException;
import java.net.InetSocketAddress;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
Expand Down Expand Up @@ -51,8 +51,7 @@ public class ConnectionFactoryBuilder {

private FailureMode failureMode = FailureMode.Cancel;

private Collection<ConnectionObserver> initialObservers
= Collections.emptyList();
private List<ConnectionObserver> initialObservers = new ArrayList<>();

private OperationFactory opFact;

Expand Down Expand Up @@ -202,16 +201,32 @@ public ConnectionFactoryBuilder setFailureMode(FailureMode fm) {
/**
* Set the initial connection observers (will observe initial connection).
*/
public ConnectionFactoryBuilder setInitialObservers(
Collection<ConnectionObserver> obs) {
public ConnectionFactoryBuilder setInitialObservers(Collection<ConnectionObserver> obs) {
if (obs == null || obs.isEmpty()) {
throw new IllegalArgumentException("Initial observers must not be null or empty.");
}

initialObservers = obs;
initialObservers.clear();
initialObservers.addAll(obs);
return this;
}

void addInitialObserver(ConnectionObserver observer) {
if (observer == null) {
throw new IllegalArgumentException("Initial observer must not be null.");
}

initialObservers.add(observer);
}

void removeInitialObserver(ConnectionObserver observer) {
if (observer == null) {
throw new IllegalArgumentException("Initial observer must not be null.");
}

initialObservers.remove(observer);
}

/**
* Set the operation factory.
*
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
package net.spy.memcached;

import java.util.Collection;
import java.util.Collections;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;

import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;

import static org.junit.jupiter.api.Assertions.assertAll;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;

class ArcusClientInitialObserverTest {

private ArcusClient client;

@AfterEach
void tearDown() {
if (client != null) {
client.shutdown();
}
}

@Test
void test() {
// given
CountDownLatch latch = new CountDownLatch(1);

ConnectionObserver observer = new ConnectionObserver() {
@Override
public void connectionEstablished(MemcachedNode node, int reconnectCount) {
latch.countDown();
}

@Override
public void connectionLost(MemcachedNode node) {
// do-nothing.
}
};

// when
ConnectionFactoryBuilder cfb = new ConnectionFactoryBuilder()
.setInitialObservers(Collections.singletonList(observer));

client = ArcusClient.createArcusClient(
"127.0.0.1:2181",
"test",
cfb
);

Collection<ConnectionObserver> configuredObservers = cfb.build().getInitialObservers();

// then
assertAll(
() -> assertTrue(latch.await(700, TimeUnit.MILLISECONDS)),
() -> assertTrue(client.removeObserver(observer)),
() -> assertEquals(1, configuredObservers.size()),
() -> assertTrue(configuredObservers.contains(observer))
);
}
}
Loading