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
Original file line number Diff line number Diff line change
Expand Up @@ -313,7 +313,12 @@ public void run() {
try {
if (!locked) {
if (!lockManager.isStarted()) {
lockManager.start();
// a bounded start: an unreachable lock manager must not hold this executor forever,
// it is simply retried on the next scheduled run
if (!lockManager.start(getPeriod(), getTimeUnit())) {
logger.debug("Not able to start the lock manager for {}, lockID={}", name, lockID);
return;
}
}
DistributedLock lock = lockManager.getDistributedLock(lockID);
if (lock.tryLock(1, TimeUnit.SECONDS)) {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,182 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You 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 org.apache.activemq.artemis.core.server.lock;

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

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;

import org.apache.activemq.artemis.lockmanager.DistributedLock;
import org.apache.activemq.artemis.lockmanager.DistributedLockManager;
import org.apache.activemq.artemis.lockmanager.MutableLong;
import org.apache.activemq.artemis.tests.util.ArtemisTestCase;
import org.apache.activemq.artemis.utils.ActiveMQThreadFactory;
import org.apache.activemq.artemis.utils.actors.OrderedExecutor;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;

public class LockCoordinatorTest extends ArtemisTestCase {

private static final int CHECK_PERIOD = 100;

/**
* A DistributedLockManager blocking forever on start, the same way CuratorDistributedLockManager does
* when ZooKeeper is not reachable: start() calls start(-1, null), which parks on
* CuratorFramework::blockUntilConnected with no timeout.
*/
private static class BlockedOnStartLockManager implements DistributedLockManager {

private final CountDownLatch startCalled;
private final CountDownLatch releaseStart;
private volatile boolean started;

BlockedOnStartLockManager(CountDownLatch startCalled, CountDownLatch releaseStart) {
this.startCalled = startCalled;
this.releaseStart = releaseStart;
}

@Override
public void addUnavailableManagerListener(UnavailableManagerListener listener) {
}

@Override
public void removeUnavailableManagerListener(UnavailableManagerListener listener) {
}

@Override
public boolean start(long timeout, TimeUnit unit) throws InterruptedException {
startCalled.countDown();
if (timeout >= 0) {
return releaseStart.await(timeout, unit) && markStarted();
}
releaseStart.await();
return markStarted();
}

@Override
public void start() throws InterruptedException {
start(-1, null);
}

private boolean markStarted() {
started = true;
return true;
}

@Override
public boolean isStarted() {
return started;
}

@Override
public void stop() {
started = false;
}

@Override
public DistributedLock getDistributedLock(String lockId) {
return new NeverAcquiredLock(lockId);
}

@Override
public MutableLong getMutableLong(String mutableLongId) {
throw new UnsupportedOperationException();
}
}

private static class NeverAcquiredLock implements DistributedLock {

private final String lockId;

NeverAcquiredLock(String lockId) {
this.lockId = lockId;
}

@Override
public String getLockId() {
return lockId;
}

@Override
public boolean isHeldByCaller() {
return false;
}

@Override
public boolean tryLock() {
return false;
}

@Override
public void unlock() {
}

@Override
public void addListener(UnavailableLockListener listener) {
}

@Override
public void removeListener(UnavailableLockListener listener) {
}

@Override
public void close() {
}
}

/**
* A lock manager unable to connect must not make LockCoordinator::stop hang: the periodic task and the
* cleanup share the same ordered executor, hence a start blocking forever would keep the broker from
* being stopped or restarted.
*/
@Test
@Timeout(value = 60, unit = TimeUnit.SECONDS)
public void testStopWithLockManagerBlockedOnStart() throws Exception {
final CountDownLatch startCalled = new CountDownLatch(1);
final CountDownLatch releaseStart = new CountDownLatch(1);
final ScheduledExecutorService scheduledExecutor = Executors.newSingleThreadScheduledExecutor(ActiveMQThreadFactory.defaultThreadFactory(getClass().getName()));
final ExecutorService executorService = Executors.newSingleThreadExecutor(ActiveMQThreadFactory.defaultThreadFactory(getClass().getName()));
try {
final DistributedLockManager lockManager = new BlockedOnStartLockManager(startCalled, releaseStart);
final LockCoordinator coordinator = new LockCoordinator(scheduledExecutor, new OrderedExecutor(executorService), CHECK_PERIOD, lockManager, "theLock", "theLock");
coordinator.start();

assertTrue(startCalled.await(10, TimeUnit.SECONDS), "the lock coordinator never tried to start the lock manager");

final CountDownLatch stopped = new CountDownLatch(1);
final Thread stopper = new Thread(() -> {
coordinator.stop();
stopped.countDown();
}, "lock-coordinator-stopper");
stopper.start();
try {
assertTrue(stopped.await(30, TimeUnit.SECONDS), "LockCoordinator::stop is hanging while the lock manager is blocked on start");
} finally {
releaseStart.countDown();
stopper.join(TimeUnit.SECONDS.toMillis(10));
}
} finally {
releaseStart.countDown();
executorService.shutdownNow();
scheduledExecutor.shutdownNow();
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You 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 org.apache.activemq.artemis.tests.integration.lockmanager;

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

import java.io.IOException;
import java.net.ServerSocket;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;

import org.apache.activemq.artemis.core.config.Configuration;
import org.apache.activemq.artemis.core.config.LockCoordinatorConfiguration;
import org.apache.activemq.artemis.core.server.ActiveMQServers;
import org.apache.activemq.artemis.core.server.impl.ActiveMQServerImpl;
import org.apache.activemq.artemis.tests.util.ActiveMQTestBase;
import org.apache.activemq.artemis.utils.Wait;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;

/**
* A lock coordinator unable to reach ZooKeeper must not make the broker unstoppable.
*/
public class UnreachableLockManagerTest extends ActiveMQTestBase {

private static final String CURATOR_LOCK_MANAGER = "org.apache.activemq.artemis.lockmanager.zookeeper.CuratorDistributedLockManager";
private static final int CHECK_PERIOD = 100;

/**
* The lock coordinator is not referenced by any acceptor or broker connection, still it is started with the
* broker: with an unreachable ZooKeeper its periodic task used to park forever on the lock manager start,
* holding the ordered executor the stop is relying on, hence the broker could not be stopped or restarted.
*/
@Test
@Timeout(value = 120, unit = TimeUnit.SECONDS)
public void testStopWithUnreachableZooKeeper() throws Exception {
final Configuration config = createDefaultInVMConfig();

final Map<String, String> properties = new HashMap<>();
// nothing is listening on this port: the ZooKeeper ensemble is unreachable
properties.put("connect-string", "localhost:" + unusedPort());
properties.put("namespace", "activemq-artemis");
properties.put("session-ms", "2000");
properties.put("connection-ms", "2000");
final LockCoordinatorConfiguration lockCoordinatorConfiguration = new LockCoordinatorConfiguration(properties);
lockCoordinatorConfiguration.setName("zk").setClassName(CURATOR_LOCK_MANAGER).setCheckPeriod(CHECK_PERIOD).setLockId("zk");
config.addLockCoordinatorConfiguration(lockCoordinatorConfiguration);

// the server is created outside of the test base on purpose: an unstoppable broker would hang the tear down
final ActiveMQServerImpl server = (ActiveMQServerImpl) ActiveMQServers.newActiveMQServer(config, false);
server.start();
Wait.assertTrue(() -> server.getLockCoordinators().size() == 1, 10_000, 100);
// let the periodic task run and try to connect
Wait.assertTrue(() -> "Unlocked".equals(server.getLockCoordinator("zk").getStatus()), 10_000, 100);

final CountDownLatch stopped = new CountDownLatch(1);
final Thread stopper = new Thread(() -> {
try {
server.stop();
} catch (Exception e) {
e.printStackTrace();
} finally {
stopped.countDown();
}
}, "unreachable-lock-manager-server-stopper");
stopper.start();
try {
assertTrue(stopped.await(60, TimeUnit.SECONDS), "the broker cannot be stopped while the lock coordinator cannot reach ZooKeeper");
} finally {
stopper.join(TimeUnit.SECONDS.toMillis(10));
}
}

private static int unusedPort() throws IOException {
try (ServerSocket socket = new ServerSocket(0)) {
return socket.getLocalPort();
}
}
}
Loading