Some improvements to the admin module:
- Re-org exec a bit - Add checks for env vars in admin - Check upper case host name on Windoze - Fix problem with user admin binary converter in Rabbit 2.1 - Add pauses to prevent RPC errors (timing issue in JInterface?)
This commit is contained in:
1
.gitignore
vendored
1
.gitignore
vendored
@@ -4,3 +4,4 @@ target
|
||||
.ant-targets-build.xml
|
||||
src/ant/.ant-targets-upload-dist.xml
|
||||
*.swp
|
||||
erl_crash.dump
|
||||
|
||||
@@ -19,7 +19,7 @@
|
||||
<org.slf4j.version>1.5.10</org.slf4j.version>
|
||||
<org.codehaus.jackson.version>1.4.3</org.codehaus.jackson.version>
|
||||
<org.erlang.otp.version>1.5.3</org.erlang.otp.version>
|
||||
<com.rabbitmq.version>2.0.0</com.rabbitmq.version>
|
||||
<com.rabbitmq.version>2.1.0</com.rabbitmq.version>
|
||||
<org.springframework.version>3.0.3.RELEASE</org.springframework.version>
|
||||
</properties>
|
||||
<profiles>
|
||||
|
||||
@@ -10,7 +10,7 @@ public class BrokerConfigurationApplication {
|
||||
/**
|
||||
* An example application that only configures the AMQP broker
|
||||
*/
|
||||
public static void main(String[] args) {
|
||||
public static void main(String[] args) throws Exception {
|
||||
ApplicationContext context = new ClassPathXmlApplicationContext("rabbitConfiguration.xml");
|
||||
AmqpAdmin amqpAdmin = context.getBean(AmqpAdmin.class);
|
||||
Queue helloWorldQueue = new Queue("hello.world.queue");
|
||||
|
||||
@@ -28,4 +28,8 @@ public class OtpAuthException extends OtpException {
|
||||
super(cause);
|
||||
}
|
||||
|
||||
public OtpAuthException(String msg, com.ericsson.otp.erlang.OtpAuthException cause) {
|
||||
super(msg, cause);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -46,9 +46,8 @@ public class SimpleConnectionFactory implements ConnectionFactory, InitializingB
|
||||
|
||||
|
||||
public SimpleConnectionFactory(String selfNodeName, String cookie, String peerNodeName) {
|
||||
this.selfNodeName = selfNodeName;
|
||||
this(selfNodeName, peerNodeName);
|
||||
this.cookie = cookie;
|
||||
this.peerNodeName = peerNodeName;
|
||||
}
|
||||
|
||||
public SimpleConnectionFactory(String selfNodeName, String peerNodeName) {
|
||||
|
||||
@@ -23,6 +23,7 @@ import org.springframework.erlang.OtpIOException;
|
||||
import org.springframework.erlang.UncategorizedOtpException;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
import com.ericsson.otp.erlang.OtpAuthException;
|
||||
import com.ericsson.otp.erlang.OtpConnection;
|
||||
|
||||
/**
|
||||
@@ -47,6 +48,9 @@ public class ErlangUtils {
|
||||
if (ex instanceof IOException) {
|
||||
return new OtpIOException((IOException) ex);
|
||||
}
|
||||
if (ex instanceof OtpAuthException) {
|
||||
return new org.springframework.erlang.OtpAuthException((OtpAuthException) ex);
|
||||
}
|
||||
//fallback
|
||||
return new UncategorizedOtpException(ex);
|
||||
}
|
||||
|
||||
@@ -16,389 +16,404 @@
|
||||
|
||||
package org.springframework.util.exec;
|
||||
|
||||
import java.io.BufferedReader;
|
||||
import java.io.ByteArrayOutputStream;
|
||||
import java.io.File;
|
||||
import java.io.IOException;
|
||||
|
||||
import java.io.BufferedReader;
|
||||
import java.io.StringReader;
|
||||
import java.io.ByteArrayOutputStream;
|
||||
import java.util.*;
|
||||
|
||||
import java.util.StringTokenizer;
|
||||
import java.util.Vector;
|
||||
|
||||
/* Derived from ant exec task. All the 'backward compat with jdk1.1, 1.2'
|
||||
removed. Since jdk1.3 supports working dir, no need for scripts.
|
||||
removed. Since jdk1.3 supports working dir, no need for scripts.
|
||||
|
||||
All ant-specific code has been removed as well, this is a completely
|
||||
independent component.
|
||||
|
||||
Costin
|
||||
*/
|
||||
All ant-specific code has been removed as well, this is a completely
|
||||
independent component.
|
||||
|
||||
Costin
|
||||
*/
|
||||
|
||||
/**
|
||||
* Runs an external program.
|
||||
*
|
||||
*
|
||||
* @author thomas.haas@softwired-inc.com
|
||||
*/
|
||||
public class Execute {
|
||||
/** Invalid exit code. **/
|
||||
public final static int INVALID = Integer.MAX_VALUE;
|
||||
/** Invalid exit code. **/
|
||||
public final static int INVALID = Integer.MAX_VALUE;
|
||||
|
||||
private String[] cmdl = null;
|
||||
private String[] env = null;
|
||||
private int exitValue = INVALID;
|
||||
private ExecuteStreamHandler streamHandler;
|
||||
private ExecuteWatchdog watchdog;
|
||||
private File workingDirectory = null;
|
||||
private boolean newEnvironment = false;
|
||||
private String[] cmdl = null;
|
||||
private String[] env = null;
|
||||
private int exitValue = INVALID;
|
||||
private ExecuteStreamHandler streamHandler;
|
||||
private ExecuteWatchdog watchdog;
|
||||
private File workingDirectory = null;
|
||||
private boolean newEnvironment = false;
|
||||
private Process process;
|
||||
|
||||
private static Vector<String> procEnvironment = null;
|
||||
private static Vector<String> procEnvironment = null;
|
||||
|
||||
/**
|
||||
* Find the list of environment variables for this process.
|
||||
*/
|
||||
public static synchronized Vector<String> getProcEnvironment() {
|
||||
if (procEnvironment != null) return procEnvironment;
|
||||
/**
|
||||
* Find the list of environment variables for this process.
|
||||
*/
|
||||
public static synchronized Vector<String> getProcEnvironment() {
|
||||
if (procEnvironment != null)
|
||||
return procEnvironment;
|
||||
|
||||
procEnvironment = new Vector<String>();
|
||||
try {
|
||||
ByteArrayOutputStream out = new ByteArrayOutputStream();
|
||||
Execute exe = new Execute(new PumpStreamHandler(out));
|
||||
exe.setCommandline(getProcEnvCommand());
|
||||
// Make sure we do not recurse forever
|
||||
exe.setNewenvironment(true);
|
||||
int retval = exe.execute();
|
||||
if ( retval != 0 ) {
|
||||
// Just try to use what we got
|
||||
}
|
||||
procEnvironment = new Vector<String>();
|
||||
try {
|
||||
ByteArrayOutputStream out = new ByteArrayOutputStream();
|
||||
Execute exe = new Execute(new PumpStreamHandler(out));
|
||||
exe.setCommandline(getProcEnvCommand());
|
||||
// Make sure we do not recurse forever
|
||||
exe.setNewenvironment(true);
|
||||
int retval = exe.execute();
|
||||
if (retval != 0) {
|
||||
// Just try to use what we got
|
||||
}
|
||||
|
||||
BufferedReader in =
|
||||
new BufferedReader(new StringReader(out.toString()));
|
||||
String var = null;
|
||||
String line, lineSep = System.getProperty("line.separator");
|
||||
while ((line = in.readLine()) != null) {
|
||||
if (line.indexOf('=') == -1) {
|
||||
// Chunk part of previous env var (UNIX env vars can
|
||||
// contain embedded new lines).
|
||||
if (var == null) {
|
||||
var = lineSep + line;
|
||||
}
|
||||
else {
|
||||
var += lineSep + line;
|
||||
}
|
||||
}
|
||||
else {
|
||||
// New env var...append the previous one if we have it.
|
||||
if (var != null) {
|
||||
procEnvironment.addElement(var);
|
||||
}
|
||||
var = line;
|
||||
}
|
||||
}
|
||||
// Since we "look ahead" before adding, there's one last env var.
|
||||
procEnvironment.addElement(var);
|
||||
}
|
||||
catch (Exception exc) {
|
||||
exc.printStackTrace();
|
||||
// Just try to see how much we got
|
||||
}
|
||||
return procEnvironment;
|
||||
}
|
||||
BufferedReader in = new BufferedReader(new StringReader(
|
||||
out.toString()));
|
||||
String var = null;
|
||||
String line, lineSep = System.getProperty("line.separator");
|
||||
while ((line = in.readLine()) != null) {
|
||||
if (line.indexOf('=') == -1) {
|
||||
// Chunk part of previous env var (UNIX env vars can
|
||||
// contain embedded new lines).
|
||||
if (var == null) {
|
||||
var = lineSep + line;
|
||||
} else {
|
||||
var += lineSep + line;
|
||||
}
|
||||
} else {
|
||||
// New env var...append the previous one if we have it.
|
||||
if (var != null) {
|
||||
procEnvironment.addElement(var);
|
||||
}
|
||||
var = line;
|
||||
}
|
||||
}
|
||||
// Since we "look ahead" before adding, there's one last env var.
|
||||
procEnvironment.addElement(var);
|
||||
} catch (Exception exc) {
|
||||
exc.printStackTrace();
|
||||
// Just try to see how much we got
|
||||
}
|
||||
return procEnvironment;
|
||||
}
|
||||
|
||||
private static String[] getProcEnvCommand() {
|
||||
if ( Os.isFamily("os/2") ) {
|
||||
// OS/2 - use same mechanism as Windows 2000
|
||||
// Not sure
|
||||
String[] cmd = {"cmd", "/c", "set" };
|
||||
return cmd;
|
||||
}
|
||||
else if ( Os.isFamily("windows") ) {
|
||||
String[] cmd = {"cmd", "/c", "set" };
|
||||
return cmd;
|
||||
}
|
||||
else if ( Os.isFamily("unix") ) {
|
||||
// Generic UNIX
|
||||
// Alternatively one could use: /bin/sh -c env
|
||||
String[] cmd = {"/usr/bin/env"};
|
||||
return cmd;
|
||||
}
|
||||
else if ( Os.isFamily("netware") ) {
|
||||
String[] cmd = {"env"};
|
||||
return cmd;
|
||||
}
|
||||
else {
|
||||
// MAC OS 9 and previous
|
||||
// TODO: I have no idea how to get it, someone must fix it
|
||||
String[] cmd = null;
|
||||
return cmd;
|
||||
}
|
||||
}
|
||||
private static String[] getProcEnvCommand() {
|
||||
if (Os.isFamily("os/2")) {
|
||||
// OS/2 - use same mechanism as Windows 2000
|
||||
// Not sure
|
||||
String[] cmd = { "cmd", "/c", "set" };
|
||||
return cmd;
|
||||
} else if (Os.isFamily("windows")) {
|
||||
String[] cmd = { "cmd", "/c", "set" };
|
||||
return cmd;
|
||||
} else if (Os.isFamily("unix")) {
|
||||
// Generic UNIX
|
||||
// Alternatively one could use: /bin/sh -c env
|
||||
String[] cmd = { "/usr/bin/env" };
|
||||
return cmd;
|
||||
} else if (Os.isFamily("netware")) {
|
||||
String[] cmd = { "env" };
|
||||
return cmd;
|
||||
} else {
|
||||
// MAC OS 9 and previous
|
||||
// TODO: I have no idea how to get it, someone must fix it
|
||||
String[] cmd = null;
|
||||
return cmd;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a new execute object using <code>PumpStreamHandler</code> for
|
||||
* stream handling.
|
||||
*/
|
||||
public Execute() {
|
||||
this(new PumpStreamHandler(), null);
|
||||
}
|
||||
/**
|
||||
* Creates a new execute object using <code>PumpStreamHandler</code> for
|
||||
* stream handling.
|
||||
*/
|
||||
public Execute() {
|
||||
this(new PumpStreamHandler(), null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a new execute object.
|
||||
*
|
||||
* @param streamHandler
|
||||
* the stream handler used to handle the input and output streams
|
||||
* of the subprocess.
|
||||
*/
|
||||
public Execute(ExecuteStreamHandler streamHandler) {
|
||||
this(streamHandler, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a new execute object.
|
||||
*
|
||||
* @param streamHandler the stream handler used to handle the input and
|
||||
* output streams of the subprocess.
|
||||
*/
|
||||
public Execute(ExecuteStreamHandler streamHandler) {
|
||||
this(streamHandler, null);
|
||||
}
|
||||
/**
|
||||
* Creates a new execute object.
|
||||
*
|
||||
* @param streamHandler
|
||||
* the stream handler used to handle the input and output streams
|
||||
* of the subprocess.
|
||||
* @param watchdog
|
||||
* a watchdog for the subprocess or <code>null</code> to to
|
||||
* disable a timeout for the subprocess.
|
||||
*/
|
||||
public Execute(ExecuteStreamHandler streamHandler, ExecuteWatchdog watchdog) {
|
||||
this.streamHandler = streamHandler;
|
||||
this.watchdog = watchdog;
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a new execute object.
|
||||
*
|
||||
* @param streamHandler the stream handler used to handle the input and
|
||||
* output streams of the subprocess.
|
||||
* @param watchdog a watchdog for the subprocess or <code>null</code> to
|
||||
* to disable a timeout for the subprocess.
|
||||
*/
|
||||
public Execute(ExecuteStreamHandler streamHandler, ExecuteWatchdog watchdog) {
|
||||
this.streamHandler = streamHandler;
|
||||
this.watchdog = watchdog;
|
||||
}
|
||||
/**
|
||||
* Returns the commandline used to create a subprocess.
|
||||
*
|
||||
* @return the commandline used to create a subprocess
|
||||
*/
|
||||
public String[] getCommandline() {
|
||||
return cmdl;
|
||||
}
|
||||
|
||||
public String getCommandLineString() {
|
||||
return array2string(getCommandline());
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the commandline used to create a subprocess.
|
||||
*
|
||||
* @return the commandline used to create a subprocess
|
||||
*/
|
||||
public String[] getCommandline() {
|
||||
return cmdl;
|
||||
}
|
||||
/**
|
||||
* Sets the commandline of the subprocess to launch.
|
||||
*
|
||||
* @param commandline
|
||||
* the commandline of the subprocess to launch
|
||||
*/
|
||||
public void setCommandline(String[] commandline) {
|
||||
cmdl = commandline;
|
||||
}
|
||||
|
||||
public String getCommandLineString() {
|
||||
return array2string(getCommandline());
|
||||
}
|
||||
/**
|
||||
* Set whether to propagate the default environment or not.
|
||||
*
|
||||
* @param newenv
|
||||
* whether to propagate the process environment.
|
||||
*/
|
||||
public void setNewenvironment(boolean newenv) {
|
||||
newEnvironment = newenv;
|
||||
}
|
||||
|
||||
/**
|
||||
* Sets the commandline of the subprocess to launch.
|
||||
*
|
||||
* @param commandline the commandline of the subprocess to launch
|
||||
*/
|
||||
public void setCommandline(String[] commandline) {
|
||||
cmdl = commandline;
|
||||
}
|
||||
/**
|
||||
* Returns the environment used to create a subprocess.
|
||||
*
|
||||
* @return the environment used to create a subprocess
|
||||
*/
|
||||
public String[] getEnvironment() {
|
||||
if (env == null || newEnvironment)
|
||||
return env;
|
||||
return patchEnvironment();
|
||||
}
|
||||
|
||||
/**
|
||||
* Set whether to propagate the default environment or not.
|
||||
*
|
||||
* @param newenv whether to propagate the process environment.
|
||||
*/
|
||||
public void setNewenvironment(boolean newenv) {
|
||||
newEnvironment = newenv;
|
||||
}
|
||||
/**
|
||||
* Sets the environment variables for the subprocess to launch.
|
||||
*
|
||||
* @param env
|
||||
* array of Strings, each element of which has an environment
|
||||
* variable settings in format <em>key=value</em>
|
||||
*/
|
||||
public void setEnvironment(String[] env) {
|
||||
this.env = env;
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns the environment used to create a subprocess.
|
||||
*
|
||||
* @return the environment used to create a subprocess
|
||||
*/
|
||||
public String[] getEnvironment() {
|
||||
if (env == null || newEnvironment) return env;
|
||||
return patchEnvironment();
|
||||
}
|
||||
/**
|
||||
* Sets the working directory of the process to execute.
|
||||
*
|
||||
* <p>
|
||||
* This is emulated using the antRun scripts unless the OS is Windows NT in
|
||||
* which case a cmd.exe is spawned, or MRJ and setting user.dir works, or
|
||||
* JDK 1.3 and there is official support in java.lang.Runtime.
|
||||
*
|
||||
* @param wd
|
||||
* the working directory of the process.
|
||||
*/
|
||||
public void setWorkingDirectory(File wd) {
|
||||
workingDirectory = wd;
|
||||
}
|
||||
|
||||
/**
|
||||
* Runs a process defined by the command line and returns its exit status.
|
||||
*
|
||||
* @return the exit status of the subprocess or <code>INVALID</code>
|
||||
* @throws Exception
|
||||
* if launching of the subprocess failed
|
||||
*/
|
||||
public int execute() throws Exception {
|
||||
process = Runtime.getRuntime().exec(getCommandline(), getEnvironment(),
|
||||
workingDirectory);
|
||||
try {
|
||||
streamHandler.setProcessInputStream(process.getOutputStream());
|
||||
streamHandler.setProcessOutputStream(process.getInputStream());
|
||||
streamHandler.setProcessErrorStream(process.getErrorStream());
|
||||
} catch (IOException e) {
|
||||
process.destroy();
|
||||
throw e;
|
||||
}
|
||||
streamHandler.start();
|
||||
|
||||
/**
|
||||
* Sets the environment variables for the subprocess to launch.
|
||||
*
|
||||
* @param env array of Strings, each element of which has
|
||||
* an environment variable settings in format <em>key=value</em>
|
||||
*/
|
||||
public void setEnvironment(String[] env) {
|
||||
this.env = env;
|
||||
}
|
||||
if (watchdog != null)
|
||||
watchdog.start(process, Thread.currentThread());
|
||||
|
||||
/**
|
||||
* Sets the working directory of the process to execute.
|
||||
*
|
||||
* <p>This is emulated using the antRun scripts unless the OS is
|
||||
* Windows NT in which case a cmd.exe is spawned,
|
||||
* or MRJ and setting user.dir works, or JDK 1.3 and there is
|
||||
* official support in java.lang.Runtime.
|
||||
*
|
||||
* @param wd the working directory of the process.
|
||||
*/
|
||||
public void setWorkingDirectory(File wd) {
|
||||
workingDirectory = wd;
|
||||
}
|
||||
if (log.isTraceEnabled())
|
||||
log.trace("Waiting process ");
|
||||
waitFor(process);
|
||||
process = null;
|
||||
|
||||
// costin
|
||||
|
||||
boolean wait=true;
|
||||
public void setWait( boolean b ) {
|
||||
wait=b;
|
||||
}
|
||||
if (log.isTraceEnabled())
|
||||
log.trace("End waiting, stop threads ");
|
||||
if (watchdog != null)
|
||||
watchdog.stop();
|
||||
if (log.isTraceEnabled())
|
||||
log.trace("Watchdog stopped ");
|
||||
streamHandler.stop();
|
||||
if (log.isTraceEnabled())
|
||||
log.trace("Stream handler stopped ");
|
||||
if (watchdog != null) {
|
||||
Exception ex = watchdog.getException();
|
||||
if (ex != null)
|
||||
throw ex;
|
||||
}
|
||||
int exit = getExitValue();
|
||||
|
||||
Process process;
|
||||
/**
|
||||
* Runs a process defined by the command line and returns its exit status.
|
||||
*
|
||||
* @return the exit status of the subprocess or <code>INVALID</code>
|
||||
* @throws Exception if launching of the subprocess failed
|
||||
*/
|
||||
public int execute() throws Exception {
|
||||
process =
|
||||
Runtime.getRuntime().exec( getCommandline(), getEnvironment(),
|
||||
workingDirectory );
|
||||
try {
|
||||
streamHandler.setProcessInputStream(process.getOutputStream());
|
||||
streamHandler.setProcessOutputStream(process.getInputStream());
|
||||
streamHandler.setProcessErrorStream(process.getErrorStream());
|
||||
} catch (IOException e) {
|
||||
process.destroy();
|
||||
throw e;
|
||||
}
|
||||
streamHandler.start();
|
||||
if (log.isDebugEnabled()) {
|
||||
log.debug("Done exit=" + exit + " " + getCommandLineString());
|
||||
}
|
||||
return exit;
|
||||
}
|
||||
|
||||
if (watchdog != null) watchdog.start(process,
|
||||
Thread.currentThread());
|
||||
public void kill() {
|
||||
if (process != null) {
|
||||
process.destroy();
|
||||
}
|
||||
}
|
||||
|
||||
if( log.isTraceEnabled() ) log.trace("Waiting process ");
|
||||
waitFor(process);
|
||||
private String array2string(String sa[]) {
|
||||
if (sa == null)
|
||||
return "null";
|
||||
StringBuffer sb = new StringBuffer();
|
||||
for (int i = 0; i < sa.length; i++)
|
||||
sb.append(sa[i]).append(" ");
|
||||
return sb.toString();
|
||||
}
|
||||
|
||||
if( log.isTraceEnabled() ) log.trace("End waiting, stop threads ");
|
||||
if (watchdog != null) watchdog.stop();
|
||||
if( log.isTraceEnabled() ) log.trace("Watchdog stopped ");
|
||||
streamHandler.stop();
|
||||
if( log.isTraceEnabled() ) log.trace("Stream handler stopped ");
|
||||
if (watchdog != null) {
|
||||
Exception ex=watchdog.getException();
|
||||
if( ex!=null )
|
||||
throw ex;
|
||||
}
|
||||
int exit= getExitValue();
|
||||
protected void waitFor(Process process) {
|
||||
try {
|
||||
process.waitFor();
|
||||
setExitValue(process.exitValue());
|
||||
} catch (InterruptedException e) {
|
||||
log.info("waitFor() interrupted ");
|
||||
}
|
||||
}
|
||||
|
||||
if( log.isDebugEnabled() ) {
|
||||
log.debug("Done exit=" + exit + " " + getCommandLineString());
|
||||
}
|
||||
return exit;
|
||||
}
|
||||
protected void setExitValue(int value) {
|
||||
exitValue = value;
|
||||
}
|
||||
|
||||
private String array2string( String sa[]) {
|
||||
if( sa==null ) return "null";
|
||||
StringBuffer sb=new StringBuffer();
|
||||
for( int i=0; i<sa.length; i++ ) sb.append(sa[i]).append(" ");
|
||||
return sb.toString();
|
||||
}
|
||||
/**
|
||||
* query the exit value of the process.
|
||||
*
|
||||
* @return the exit value, 1 if the process was killed, or Project.INVALID
|
||||
* if no exit value has been received
|
||||
*/
|
||||
public int getExitValue() {
|
||||
return exitValue;
|
||||
}
|
||||
|
||||
protected void waitFor(Process process) {
|
||||
try {
|
||||
process.waitFor();
|
||||
setExitValue(process.exitValue());
|
||||
} catch (InterruptedException e) {
|
||||
log.info("waitFor() interrupted ");
|
||||
}
|
||||
}
|
||||
|
||||
protected void setExitValue(int value) {
|
||||
exitValue = value;
|
||||
}
|
||||
|
||||
/**
|
||||
* query the exit value of the process.
|
||||
* @return the exit value, 1 if the process was killed,
|
||||
* or Project.INVALID if no exit value has been received
|
||||
*/
|
||||
public int getExitValue() {
|
||||
return exitValue;
|
||||
}
|
||||
|
||||
/**
|
||||
* Patch the current environment with the new values from the user.
|
||||
* @return the patched environment
|
||||
*/
|
||||
private String[] patchEnvironment() {
|
||||
@SuppressWarnings("unchecked")
|
||||
/**
|
||||
* Patch the current environment with the new values from the user.
|
||||
*
|
||||
* @return the patched environment
|
||||
*/
|
||||
private String[] patchEnvironment() {
|
||||
@SuppressWarnings("unchecked")
|
||||
Vector<String> osEnv = (Vector<String>) getProcEnvironment().clone();
|
||||
for (int i = 0; i < env.length; i++) {
|
||||
int pos = env[i].indexOf('=');
|
||||
// Get key including "="
|
||||
String key = env[i].substring(0, pos+1);
|
||||
int size = osEnv.size();
|
||||
for (int j = 0; j < size; j++) {
|
||||
if ((osEnv.elementAt(j)).startsWith(key)) {
|
||||
osEnv.removeElementAt(j);
|
||||
break;
|
||||
}
|
||||
}
|
||||
osEnv.addElement(env[i]);
|
||||
}
|
||||
String[] result = new String[osEnv.size()];
|
||||
osEnv.copyInto(result);
|
||||
return result;
|
||||
}
|
||||
for (int i = 0; i < env.length; i++) {
|
||||
int pos = env[i].indexOf('=');
|
||||
// Get key including "="
|
||||
String key = env[i].substring(0, pos + 1);
|
||||
int size = osEnv.size();
|
||||
for (int j = 0; j < size; j++) {
|
||||
if ((osEnv.elementAt(j)).startsWith(key)) {
|
||||
osEnv.removeElementAt(j);
|
||||
break;
|
||||
}
|
||||
}
|
||||
osEnv.addElement(env[i]);
|
||||
}
|
||||
String[] result = new String[osEnv.size()];
|
||||
osEnv.copyInto(result);
|
||||
return result;
|
||||
}
|
||||
|
||||
public static int execute(Vector<String> envVars, String cmd, File baseDir ) {
|
||||
Vector<String> v = new Vector<String>();
|
||||
StringTokenizer st=new StringTokenizer( cmd, " " );
|
||||
while( st.hasMoreTokens() ) {
|
||||
v.addElement( st.nextToken() );
|
||||
}
|
||||
public static int execute(Vector<String> envVars, String cmd, File baseDir) {
|
||||
Vector<String> v = new Vector<String>();
|
||||
StringTokenizer st = new StringTokenizer(cmd, " ");
|
||||
while (st.hasMoreTokens()) {
|
||||
v.addElement(st.nextToken());
|
||||
}
|
||||
|
||||
return execute( envVars, v, baseDir );
|
||||
}
|
||||
|
||||
public static int execute(Vector<String> envVars, Vector<String> cmd, File baseDir) {
|
||||
return execute( envVars, cmd, baseDir, 10000 /* default time to wait */);
|
||||
}
|
||||
return execute(envVars, v, baseDir);
|
||||
}
|
||||
|
||||
/** Wrapper for common execution patterns
|
||||
* @param envVars Environment variables to execute with (optional)
|
||||
* @param cmd a vector of the commands to execute
|
||||
* @param baseDir the base directory to run from (optional)
|
||||
* @param timeToWait milliseconds to wait for completion
|
||||
*/
|
||||
public static int execute(Vector<String> envVars, Vector<String> cmd, File baseDir, int timeToWait) {
|
||||
try {
|
||||
// We can collect the out or provide in if needed
|
||||
ExecuteWatchdog watchdog=new ExecuteWatchdog( timeToWait );
|
||||
watchdog.setDontkill( true );
|
||||
PumpStreamHandler out=new PumpStreamHandler();
|
||||
Execute exec=new Execute(out,watchdog);
|
||||
public static int execute(Vector<String> envVars, Vector<String> cmd,
|
||||
File baseDir) {
|
||||
return execute(envVars, cmd, baseDir, 10000 /* default time to wait */);
|
||||
}
|
||||
|
||||
String cmdA[]=new String[ cmd.size() ];
|
||||
cmd.toArray( cmdA );
|
||||
if( log.isDebugEnabled() ) {
|
||||
StringBuffer sb=new StringBuffer();
|
||||
for(int i=0; i<cmdA.length; i++ ) {
|
||||
sb.append(cmdA[i] + " " );
|
||||
}
|
||||
log.debug( "Exec: " + sb.toString());
|
||||
}
|
||||
exec.setCommandline( cmdA );
|
||||
/**
|
||||
* Wrapper for common execution patterns
|
||||
*
|
||||
* @param envVars
|
||||
* Environment variables to execute with (optional)
|
||||
* @param cmd
|
||||
* a vector of the commands to execute
|
||||
* @param baseDir
|
||||
* the base directory to run from (optional)
|
||||
* @param timeToWait
|
||||
* milliseconds to wait for completion
|
||||
*/
|
||||
public static int execute(Vector<String> envVars, Vector<String> cmd,
|
||||
File baseDir, int timeToWait) {
|
||||
try {
|
||||
// We can collect the out or provide in if needed
|
||||
ExecuteWatchdog watchdog = new ExecuteWatchdog(timeToWait);
|
||||
watchdog.setDontkill(true);
|
||||
PumpStreamHandler out = new PumpStreamHandler();
|
||||
Execute exec = new Execute(out, watchdog);
|
||||
|
||||
if( envVars!=null ) {
|
||||
String env[]=new String[envVars.size()];
|
||||
envVars.toArray( env );
|
||||
exec.setEnvironment( env );
|
||||
}
|
||||
String cmdA[] = new String[cmd.size()];
|
||||
cmd.toArray(cmdA);
|
||||
if (log.isDebugEnabled()) {
|
||||
StringBuffer sb = new StringBuffer();
|
||||
for (int i = 0; i < cmdA.length; i++) {
|
||||
sb.append(cmdA[i] + " ");
|
||||
}
|
||||
log.debug("Exec: " + sb.toString());
|
||||
}
|
||||
exec.setCommandline(cmdA);
|
||||
|
||||
exec.setNewenvironment( false );
|
||||
if( baseDir!=null)
|
||||
exec.setWorkingDirectory( baseDir );
|
||||
if (envVars != null) {
|
||||
String env[] = new String[envVars.size()];
|
||||
envVars.toArray(env);
|
||||
exec.setEnvironment(env);
|
||||
}
|
||||
|
||||
exec.execute();
|
||||
int status=exec.getExitValue();
|
||||
log.debug("Exit value " + status );
|
||||
return status;
|
||||
} catch( Exception ex ) {
|
||||
// ex.printStackTrace();
|
||||
System.err.println("An error has occurred in Execute.");
|
||||
return -1;
|
||||
}
|
||||
}
|
||||
exec.setNewenvironment(false);
|
||||
if (baseDir != null)
|
||||
exec.setWorkingDirectory(baseDir);
|
||||
|
||||
exec.execute();
|
||||
int status = exec.getExitValue();
|
||||
log.debug("Exit value " + status);
|
||||
return status;
|
||||
} catch (Exception ex) {
|
||||
// ex.printStackTrace();
|
||||
System.err.println("An error has occurred in Execute.");
|
||||
return -1;
|
||||
}
|
||||
}
|
||||
|
||||
private static org.apache.commons.logging.Log log = org.apache.commons.logging.LogFactory
|
||||
.getLog(Execute.class);
|
||||
|
||||
private static org.apache.commons.logging.Log log=
|
||||
org.apache.commons.logging.LogFactory.getLog( Execute.class );
|
||||
|
||||
}
|
||||
|
||||
@@ -1,9 +1,9 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<classpath>
|
||||
<classpathentry kind="src" output="target/classes" path="src/main/java"/>
|
||||
<classpathentry kind="src" output="target/test-classes" path="src/test/java"/>
|
||||
<classpathentry including="**" kind="src" output="target/test-classes" path="src/test/resources"/>
|
||||
<classpathentry kind="con" path="org.eclipse.jdt.launching.JRE_CONTAINER/org.eclipse.jdt.internal.debug.ui.launcher.StandardVMType/JavaSE-1.6"/>
|
||||
<classpathentry kind="con" path="org.maven.ide.eclipse.MAVEN2_CLASSPATH_CONTAINER"/>
|
||||
<classpathentry kind="output" path="target/classes"/>
|
||||
</classpath>
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<classpath>
|
||||
<classpathentry kind="src" output="target/classes" path="src/main/java"/>
|
||||
<classpathentry kind="src" output="target/test-classes" path="src/test/java"/>
|
||||
<classpathentry kind="src" output="target/test-classes" path="src/test/resources"/>
|
||||
<classpathentry kind="con" path="org.eclipse.jdt.launching.JRE_CONTAINER/org.eclipse.jdt.internal.debug.ui.launcher.StandardVMType/JavaSE-1.6"/>
|
||||
<classpathentry kind="con" path="org.maven.ide.eclipse.MAVEN2_CLASSPATH_CONTAINER"/>
|
||||
<classpathentry kind="output" path="target/classes"/>
|
||||
</classpath>
|
||||
|
||||
@@ -0,0 +1,12 @@
|
||||
package org.springframework.amqp.rabbit.admin;
|
||||
|
||||
import org.springframework.erlang.OtpAuthException;
|
||||
|
||||
@SuppressWarnings("serial")
|
||||
public class RabbitAdminAuthException extends OtpAuthException {
|
||||
|
||||
public RabbitAdminAuthException(String message, OtpAuthException cause) {
|
||||
super(message, (com.ericsson.otp.erlang.OtpAuthException) cause.getCause());
|
||||
}
|
||||
|
||||
}
|
||||
@@ -16,13 +16,17 @@
|
||||
|
||||
package org.springframework.amqp.rabbit.admin;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
|
||||
import org.springframework.amqp.core.Binding;
|
||||
import org.springframework.amqp.core.Exchange;
|
||||
import org.springframework.amqp.core.Queue;
|
||||
@@ -31,44 +35,48 @@ import org.springframework.amqp.rabbit.core.ChannelCallback;
|
||||
import org.springframework.amqp.rabbit.core.RabbitAdmin;
|
||||
import org.springframework.amqp.rabbit.core.RabbitTemplate;
|
||||
import org.springframework.core.task.SimpleAsyncTaskExecutor;
|
||||
import org.springframework.erlang.OtpAuthException;
|
||||
import org.springframework.erlang.OtpIOException;
|
||||
import org.springframework.erlang.connection.SimpleConnectionFactory;
|
||||
import org.springframework.erlang.core.Application;
|
||||
import org.springframework.erlang.core.ErlangTemplate;
|
||||
import org.springframework.erlang.core.Node;
|
||||
import org.springframework.jmx.export.annotation.ManagedOperation;
|
||||
import org.springframework.jmx.export.annotation.ManagedOperationParameter;
|
||||
import org.springframework.jmx.export.annotation.ManagedOperationParameters;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.StringUtils;
|
||||
import org.springframework.util.exec.Execute;
|
||||
import org.springframework.util.exec.Os;
|
||||
|
||||
import com.rabbitmq.client.Channel;
|
||||
import com.rabbitmq.client.AMQP.Exchange.DeleteOk;
|
||||
import com.rabbitmq.client.Channel;
|
||||
|
||||
/**
|
||||
* Rabbit broker administration implementation exposed via JMX annotations.
|
||||
*
|
||||
*
|
||||
* @author Mark Pollack
|
||||
*/
|
||||
public class RabbitBrokerAdmin implements RabbitBrokerOperations {
|
||||
|
||||
/** Logger available to subclasses */
|
||||
protected final Log logger = LogFactory.getLog(getClass());
|
||||
|
||||
|
||||
private RabbitTemplate rabbitTemplate;
|
||||
|
||||
|
||||
private RabbitAdmin rabbitAdmin;
|
||||
|
||||
|
||||
private ErlangTemplate erlangTemplate;
|
||||
|
||||
|
||||
private String virtualHost;
|
||||
|
||||
|
||||
public RabbitBrokerAdmin(ConnectionFactory connectionFactory) {
|
||||
this.virtualHost = connectionFactory.getVirtualHost();
|
||||
this.rabbitTemplate = new RabbitTemplate(connectionFactory);
|
||||
this.rabbitAdmin = new RabbitAdmin(connectionFactory);
|
||||
initializeDefaultErlangTemplate(rabbitTemplate);
|
||||
initializeDefaultErlangTemplate(rabbitTemplate);
|
||||
}
|
||||
|
||||
|
||||
|
||||
// Exchange Operations
|
||||
|
||||
public void declareExchange(Exchange exchange) {
|
||||
@@ -76,33 +84,33 @@ public class RabbitBrokerAdmin implements RabbitBrokerOperations {
|
||||
}
|
||||
|
||||
/**
|
||||
* Declare an exchange specifying its durability and auto-delete behavior. Explicit arguments are given so as to
|
||||
* make this method easily accessible from JMX management consoles.
|
||||
* Durable exchanges last until they are deleted, they will survive a server restart.
|
||||
* Auto-deleted exchanges last until they are no longer used
|
||||
* Declare an exchange specifying its durability and auto-delete behavior. Explicit arguments are given so as to
|
||||
* make this method easily accessible from JMX management consoles. Durable exchanges last until they are deleted,
|
||||
* they will survive a server restart. Auto-deleted exchanges last until they are no longer used
|
||||
*
|
||||
* @param exchangeName the name of the exchange
|
||||
* @param exchangeType the exchange type
|
||||
* @param durable true if we are declaring a durable exchange (the exchange will survive a server restart)
|
||||
* @param autoDelete true if the server should delete the exchange when it is no longer in use
|
||||
*/
|
||||
@ManagedOperation
|
||||
public void declareExchange(final String exchangeName, final String exchangeType, final boolean durable, final boolean autoDelete) {
|
||||
public void declareExchange(final String exchangeName, final String exchangeType, final boolean durable,
|
||||
final boolean autoDelete) {
|
||||
rabbitTemplate.execute(new ChannelCallback<Object>() {
|
||||
public Object doInRabbit(Channel channel) throws Exception {
|
||||
channel.exchangeDeclare(exchangeName, exchangeType, durable,
|
||||
autoDelete, new HashMap<String,Object>());
|
||||
channel.exchangeDeclare(exchangeName, exchangeType, durable, autoDelete, new HashMap<String, Object>());
|
||||
return null;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
@ManagedOperation(description="Delete a exchange, without regard for whether it is in use or has messages on it")
|
||||
|
||||
@ManagedOperation(description = "Delete a exchange, without regard for whether it is in use or has messages on it")
|
||||
@ManagedOperationParameters(@ManagedOperationParameter(name = "exchange", description = "the name of the exchange"))
|
||||
public void deleteExchange(String exchangeName) {
|
||||
rabbitAdmin.deleteExchange(exchangeName);
|
||||
|
||||
}
|
||||
|
||||
|
||||
@ManagedOperation
|
||||
public DeleteOk deleteExchange(final String exchangeName, final boolean ifUnused) {
|
||||
return rabbitTemplate.execute(new ChannelCallback<DeleteOk>() {
|
||||
@@ -112,15 +120,14 @@ public class RabbitBrokerAdmin implements RabbitBrokerOperations {
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
// Queue Operations
|
||||
|
||||
|
||||
@ManagedOperation
|
||||
public Queue declareQueue() {
|
||||
return rabbitAdmin.declareQueue();
|
||||
}
|
||||
|
||||
|
||||
@ManagedOperation
|
||||
public void declareQueue(Queue queue) {
|
||||
rabbitAdmin.declareQueue(queue);
|
||||
@@ -144,52 +151,56 @@ public class RabbitBrokerAdmin implements RabbitBrokerOperations {
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
public List<QueueInfo> getQueues() {
|
||||
return (List<QueueInfo>) erlangTemplate.executeAndConvertRpc("rabbit_amqqueue", "info_all", virtualHost.getBytes());
|
||||
public List<QueueInfo> getQueues() {
|
||||
return (List<QueueInfo>) erlangTemplate.executeAndConvertRpc("rabbit_amqqueue", "info_all", virtualHost
|
||||
.getBytes());
|
||||
}
|
||||
|
||||
// Binding operations
|
||||
|
||||
// Binding operations
|
||||
public void declareBinding(Binding binding) {
|
||||
rabbitAdmin.declareBinding(binding);
|
||||
rabbitAdmin.declareBinding(binding);
|
||||
}
|
||||
|
||||
public void removeBinding(final Binding binding) {
|
||||
rabbitTemplate.execute(new ChannelCallback<Object>() {
|
||||
public Object doInRabbit(Channel channel) throws Exception {
|
||||
channel.queueUnbind(binding.getQueue(), binding.getExchange(), binding.getRoutingKey(), binding.getArguments());
|
||||
channel.queueUnbind(binding.getQueue(), binding.getExchange(), binding.getRoutingKey(), binding
|
||||
.getArguments());
|
||||
return null;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
// User management
|
||||
|
||||
@ManagedOperation()
|
||||
|
||||
@ManagedOperation()
|
||||
public void addUser(String username, String password) {
|
||||
erlangTemplate.executeAndConvertRpc("rabbit_access_control", "add_user", username.getBytes(), password.getBytes());
|
||||
erlangTemplate.executeAndConvertRpc("rabbit_access_control", "add_user", username.getBytes(), password
|
||||
.getBytes());
|
||||
}
|
||||
|
||||
@ManagedOperation
|
||||
public void deleteUser(String username) {
|
||||
erlangTemplate.executeAndConvertRpc("rabbit_access_control", "delete_user", username.getBytes());
|
||||
erlangTemplate.executeAndConvertRpc("rabbit_access_control", "delete_user", username.getBytes());
|
||||
}
|
||||
|
||||
@ManagedOperation
|
||||
public void changeUserPassword(String username, String newPassword) {
|
||||
erlangTemplate.executeAndConvertRpc("rabbit_access_control", "change_password", username.getBytes(), newPassword.getBytes());
|
||||
erlangTemplate.executeAndConvertRpc("rabbit_access_control", "change_password", username.getBytes(),
|
||||
newPassword.getBytes());
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@SuppressWarnings("unchecked")
|
||||
@ManagedOperation
|
||||
public List<String> listUsers() {
|
||||
return (List<String>) erlangTemplate.executeAndConvertRpc("rabbit_access_control", "list_users");
|
||||
return (List<String>) erlangTemplate.executeAndConvertRpc("rabbit_access_control", "list_users");
|
||||
}
|
||||
|
||||
public int addVhost(String vhostPath) {
|
||||
// TODO Auto-generated method stub
|
||||
return 0;
|
||||
}
|
||||
|
||||
|
||||
public int deleteVhost(String vhostPath) {
|
||||
// TODO Auto-generated method stub
|
||||
return 0;
|
||||
@@ -240,40 +251,87 @@ public class RabbitBrokerAdmin implements RabbitBrokerOperations {
|
||||
|
||||
@ManagedOperation
|
||||
public void startNode() {
|
||||
|
||||
logger.debug("Starting RabbitMQ node by shelling out command line.");
|
||||
final Execute execute = new Execute();
|
||||
|
||||
String rabbitStartScript = null;
|
||||
String hint = "";
|
||||
if (Os.isFamily("windows") || Os.isFamily("dos")) {
|
||||
rabbitStartScript = "rabbitmq-server.bat";
|
||||
}
|
||||
else if (Os.isFamily("unix") || Os.isFamily("mac")) {
|
||||
} else if (Os.isFamily("unix") || Os.isFamily("mac")) {
|
||||
rabbitStartScript = "rabbitmq-server";
|
||||
hint = "Depending on your platform it might help to set RABBITMQ_LOG_BASE and RABBITMQ_MNESIA_BASE System properties to an empty directory.";
|
||||
}
|
||||
Assert.notNull(rabbitStartScript, "unsupported OS family");
|
||||
String rabbitHome = System.getenv("RABBITMQ_HOME");
|
||||
// TODO remove any trailing directory separators on rabbit home var.
|
||||
Assert.notNull(rabbitHome, "RABBITMQ_HOME environment variable not set.");
|
||||
String rabbitStartCommand = rabbitHome + System.getProperty("file.separator")
|
||||
+ "sbin" + System.getProperty("file.separator")
|
||||
+ rabbitStartScript;
|
||||
|
||||
String rabbitHome = System.getProperty("RABBITMQ_HOME", System.getenv("RABBITMQ_HOME"));
|
||||
Assert.notNull(rabbitHome, "RABBITMQ_HOME system property (or environment variable) not set.");
|
||||
|
||||
rabbitHome = StringUtils.cleanPath(rabbitHome);
|
||||
String rabbitStartCommand = rabbitHome + System.getProperty("file.separator") + "sbin"
|
||||
+ System.getProperty("file.separator") + rabbitStartScript;
|
||||
|
||||
List<String> env = new ArrayList<String>();
|
||||
addEnvironment(env, "RABBITMQ_LOG_BASE");
|
||||
addEnvironment(env, "RABBITMQ_MNESIA_BASE");
|
||||
addEnvironment(env, "ERLANG_HOME");
|
||||
|
||||
execute.setCommandline(new String[] { rabbitStartCommand });
|
||||
execute.setEnvironment(env.toArray(new String[0]));
|
||||
|
||||
// TODO: extract into field for DI
|
||||
SimpleAsyncTaskExecutor executor = new SimpleAsyncTaskExecutor();
|
||||
|
||||
final CountDownLatch running = new CountDownLatch(1);
|
||||
final AtomicBoolean finished = new AtomicBoolean(false);
|
||||
final String errorHint = hint;
|
||||
|
||||
executor.execute(new Runnable() {
|
||||
public void run() {
|
||||
try {
|
||||
execute.execute();
|
||||
}
|
||||
catch (Exception e) {
|
||||
running.countDown();
|
||||
int exit = execute.execute();
|
||||
finished.set(true);
|
||||
if (exit != 0) {
|
||||
throw new IllegalStateException("Could not start process." + errorHint);
|
||||
}
|
||||
} catch (Exception e) {
|
||||
logger.error("failed to start node", e);
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
logger.info("Waiting for Rabbit process to be started");
|
||||
Assert.state(running.await(1000L, TimeUnit.MILLISECONDS), "Timed out waiting for Rabbit process to start.");
|
||||
Thread.sleep(100L);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
|
||||
if (finished.get()) {
|
||||
throw new IllegalStateException("Expected broker process to start in background, but it has exited early.");
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
private void addEnvironment(List<String> env, String key) {
|
||||
String value = System.getProperty(key);
|
||||
if (value != null) {
|
||||
logger.debug("Adding environment variable: " + key + "=" + value);
|
||||
env.add(key + "=" + value);
|
||||
}
|
||||
}
|
||||
|
||||
@ManagedOperation
|
||||
public void stopNode() {
|
||||
logger.debug("Stopping RabbitMQ node.");
|
||||
erlangTemplate.executeAndConvertRpc("rabbit", "stop_and_halt");
|
||||
try {
|
||||
erlangTemplate.executeAndConvertRpc("rabbit", "stop_and_halt");
|
||||
} catch (Exception e) {
|
||||
logger.error("Failed to send stop signal", e);
|
||||
}
|
||||
}
|
||||
|
||||
@ManagedOperation
|
||||
@@ -289,7 +347,18 @@ public class RabbitBrokerAdmin implements RabbitBrokerOperations {
|
||||
|
||||
@ManagedOperation
|
||||
public RabbitStatus getStatus() {
|
||||
return (RabbitStatus) getErlangTemplate().executeAndConvertRpc("rabbit", "status");
|
||||
try {
|
||||
return (RabbitStatus) getErlangTemplate().executeAndConvertRpc("rabbit", "status");
|
||||
} catch (OtpIOException e) {
|
||||
logger.info("Ignoring OtpIOException (assuming that the broker is simply not running)");
|
||||
return new RabbitStatus(Collections.<Application> emptyList(), Collections.<Node> emptyList(), Collections
|
||||
.<Node> emptyList());
|
||||
} catch (OtpAuthException e) {
|
||||
throw new RabbitAdminAuthException(
|
||||
"Could not authorise connection to Erlang process. This can happen if the broker is running, "
|
||||
+ "but as root or rabbitmq and the current user is not authorised to connect. Try starting the "
|
||||
+ "broker again as a different user.", e);
|
||||
}
|
||||
}
|
||||
|
||||
public void recoverAsync(boolean requeue) {
|
||||
@@ -298,10 +367,14 @@ public class RabbitBrokerAdmin implements RabbitBrokerOperations {
|
||||
|
||||
public ErlangTemplate getErlangTemplate() {
|
||||
return this.erlangTemplate;
|
||||
}
|
||||
}
|
||||
|
||||
protected void initializeDefaultErlangTemplate(RabbitTemplate rabbitTemplate) {
|
||||
String peerNodeName = "rabbit@" + rabbitTemplate.getConnectionFactory().getHost();
|
||||
protected void initializeDefaultErlangTemplate(RabbitTemplate rabbitTemplate) {
|
||||
String host = rabbitTemplate.getConnectionFactory().getHost();
|
||||
if (Os.isFamily("windows")) {
|
||||
host = host.toUpperCase();
|
||||
}
|
||||
String peerNodeName = "rabbit@" + host;
|
||||
logger.debug("Creating jinterface connection with peerNodeName = [" + peerNodeName + "]");
|
||||
SimpleConnectionFactory otpCf = new SimpleConnectionFactory("rabbit-spring-monitor", peerNodeName);
|
||||
otpCf.afterPropertiesSet();
|
||||
|
||||
@@ -25,109 +25,104 @@ import org.springframework.amqp.core.Binding;
|
||||
import com.rabbitmq.client.AMQP;
|
||||
|
||||
/**
|
||||
* Performs administration tasks for RabbitMQ broker administration.
|
||||
* <p>Goal is to support full CRUD of Exchanges, Queues, Bindings, User, VHosts, etc.
|
||||
* <p>Current implementations expose operations with basic type arguments via JMX.
|
||||
* Performs administration tasks for RabbitMQ broker administration. <p>Goal is to support full CRUD of Exchanges,
|
||||
* Queues, Bindings, User, VHosts, etc. <p>Current implementations expose operations with basic type arguments via JMX.
|
||||
*
|
||||
* @author Mark Pollack
|
||||
*
|
||||
*
|
||||
*/
|
||||
public interface RabbitBrokerOperations extends AmqpAdmin {
|
||||
|
||||
// Exchange Operations
|
||||
|
||||
AMQP.Exchange.DeleteOk deleteExchange(String exchangeName, boolean ifUnused);
|
||||
|
||||
|
||||
AMQP.Exchange.DeleteOk deleteExchange(String exchangeName, boolean ifUnused);
|
||||
|
||||
void removeBinding(Binding binding);
|
||||
|
||||
|
||||
// Queue operations
|
||||
|
||||
public List<QueueInfo> getQueues();
|
||||
|
||||
// Message Delivery
|
||||
|
||||
void recoverAsync(boolean requeue);
|
||||
|
||||
// User management
|
||||
|
||||
void addUser(String username, String password);
|
||||
|
||||
void deleteUser(String username);
|
||||
|
||||
void changeUserPassword(String username, String newPassword);
|
||||
|
||||
List<String> listUsers();
|
||||
|
||||
// VHost management
|
||||
|
||||
int addVhost(String vhostPath);
|
||||
|
||||
int deleteVhost(String vhostPath);
|
||||
|
||||
// permissions
|
||||
|
||||
void setPermissions(String username, Pattern configure, Pattern read, Pattern write);
|
||||
|
||||
void setPermissions(String username, Pattern configure, Pattern read, Pattern write, String vhostPath);
|
||||
|
||||
void clearPermissions(String username);
|
||||
|
||||
void clearPermissions(String username, String vhostPath);
|
||||
|
||||
List<String> listPermissions();
|
||||
|
||||
List<String> listPermissions(String vhostPath);
|
||||
|
||||
List<String> listUserPermissions(String username);
|
||||
|
||||
|
||||
// Start/Stop/Reset broker
|
||||
|
||||
/**
|
||||
* Starts the RabbitMQ application on an already running node. This command is typically run after performing other
|
||||
* management actions that required the RabbitMQ application to be stopped, e.g. reset.
|
||||
*/
|
||||
void startBrokerApplication();
|
||||
|
||||
/**
|
||||
* Stops the RabbitMQ application, leaving the Erlang node running.
|
||||
*/
|
||||
void stopBrokerApplication();
|
||||
|
||||
/**
|
||||
* Starts the Erlang node where RabbitMQ is running by shelling out to the directory specified by RABBIT_HOME and
|
||||
* executing the standard named start script. It spawns the shell command execution into its own thread.
|
||||
*/
|
||||
void startNode();
|
||||
|
||||
/**
|
||||
* Stops the halts the Erlang node on which RabbitMQ is running. To restart the node you will need to execute
|
||||
* the start script from a command line or via other means.
|
||||
*/
|
||||
void stopNode();
|
||||
|
||||
/**
|
||||
* Removes the node from any cluster it belongs to, removes all data from the management database,
|
||||
* such as configured users and vhosts, and deletes all persistent messages.
|
||||
* <p>
|
||||
* For {@link #resetNode} and {@link #forceResetNode} to succeed the RabbitMQ application must have
|
||||
* been stopped, e.g. {@link #stopBrokerApplication}
|
||||
*/
|
||||
void resetNode();
|
||||
|
||||
/**
|
||||
* The forceResetNode command differs from {@link #resetNode} in that it resets the node unconditionally,
|
||||
* regardless of the current management database state and cluster configuration. It should only be
|
||||
* used as a last resort if the database or cluster configuration has been corrupted.
|
||||
* <p>
|
||||
* For {@link #resetNode} and {@link #forceResetNode} to succeed the RabbitMQ application must have
|
||||
* been stopped, e.g. {@link #stopBrokerApplication}
|
||||
*/
|
||||
void forceResetNode();
|
||||
|
||||
/**
|
||||
* Returns the status of the node.
|
||||
* @return status of the node.
|
||||
*/
|
||||
RabbitStatus getStatus();
|
||||
|
||||
public List<QueueInfo> getQueues();
|
||||
|
||||
// Message Delivery
|
||||
|
||||
void recoverAsync(boolean requeue);
|
||||
|
||||
// User management
|
||||
|
||||
void addUser(String username, String password);
|
||||
|
||||
void deleteUser(String username);
|
||||
|
||||
void changeUserPassword(String username, String newPassword);
|
||||
|
||||
List<String> listUsers();
|
||||
|
||||
// VHost management
|
||||
|
||||
int addVhost(String vhostPath);
|
||||
|
||||
int deleteVhost(String vhostPath);
|
||||
|
||||
// permissions
|
||||
|
||||
void setPermissions(String username, Pattern configure, Pattern read, Pattern write);
|
||||
|
||||
void setPermissions(String username, Pattern configure, Pattern read, Pattern write, String vhostPath);
|
||||
|
||||
void clearPermissions(String username);
|
||||
|
||||
void clearPermissions(String username, String vhostPath);
|
||||
|
||||
List<String> listPermissions();
|
||||
|
||||
List<String> listPermissions(String vhostPath);
|
||||
|
||||
List<String> listUserPermissions(String username);
|
||||
|
||||
// Start/Stop/Reset broker
|
||||
|
||||
/**
|
||||
* Starts the RabbitMQ application on an already running node. This command is typically run after performing other
|
||||
* management actions that required the RabbitMQ application to be stopped, e.g. reset.
|
||||
*/
|
||||
void startBrokerApplication();
|
||||
|
||||
/**
|
||||
* Stops the RabbitMQ application, leaving the Erlang node running.
|
||||
*/
|
||||
void stopBrokerApplication();
|
||||
|
||||
/**
|
||||
* Starts the Erlang node where RabbitMQ is running by shelling out to the directory specified by RABBITMQ_HOME and
|
||||
* executing the standard named start script. It spawns the shell command execution into its own thread.
|
||||
*/
|
||||
void startNode();
|
||||
|
||||
/**
|
||||
* Stops the halts the Erlang node on which RabbitMQ is running. To restart the node you will need to execute the
|
||||
* start script from a command line or via other means.
|
||||
*/
|
||||
void stopNode();
|
||||
|
||||
/**
|
||||
* Removes the node from any cluster it belongs to, removes all data from the management database, such as
|
||||
* configured users and vhosts, and deletes all persistent messages. <p> For {@link #resetNode} and
|
||||
* {@link #forceResetNode} to succeed the RabbitMQ application must have been stopped, e.g.
|
||||
* {@link #stopBrokerApplication}
|
||||
*/
|
||||
void resetNode();
|
||||
|
||||
/**
|
||||
* The forceResetNode command differs from {@link #resetNode} in that it resets the node unconditionally, regardless
|
||||
* of the current management database state and cluster configuration. It should only be used as a last resort if
|
||||
* the database or cluster configuration has been corrupted. <p> For {@link #resetNode} and {@link #forceResetNode}
|
||||
* to succeed the RabbitMQ application must have been stopped, e.g. {@link #stopBrokerApplication}
|
||||
*/
|
||||
void forceResetNode();
|
||||
|
||||
/**
|
||||
* Returns the status of the node.
|
||||
* @return status of the node.
|
||||
*/
|
||||
RabbitStatus getStatus();
|
||||
}
|
||||
|
||||
@@ -43,19 +43,16 @@ import com.ericsson.otp.erlang.OtpErlangTuple;
|
||||
*/
|
||||
public class RabbitControlErlangConverter extends SimpleErlangConverter implements ErlangConverter {
|
||||
|
||||
|
||||
protected final Log logger = LogFactory.getLog(getClass());
|
||||
|
||||
|
||||
|
||||
private Map<String, ErlangConverter> converterMap = new HashMap<String, ErlangConverter>();
|
||||
|
||||
|
||||
|
||||
public RabbitControlErlangConverter() {
|
||||
initializeConverterMap();
|
||||
}
|
||||
|
||||
|
||||
public Object fromErlangRpc(String module, String function, OtpErlangObject erlangObject) throws ErlangConversionException {
|
||||
public Object fromErlangRpc(String module, String function, OtpErlangObject erlangObject)
|
||||
throws ErlangConversionException {
|
||||
ErlangConverter converter = getConverter(module, function);
|
||||
if (converter != null) {
|
||||
return converter.fromErlang(erlangObject);
|
||||
@@ -64,20 +61,18 @@ public class RabbitControlErlangConverter extends SimpleErlangConverter implemen
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
protected ErlangConverter getConverter(String module, String function) {
|
||||
return converterMap.get(generateKey(module, function));
|
||||
return converterMap.get(generateKey(module, function));
|
||||
}
|
||||
|
||||
|
||||
protected void initializeConverterMap() {
|
||||
registerConverter("rabbit_access_control", "list_users", new ListUsersConverter());
|
||||
protected void initializeConverterMap() {
|
||||
registerConverter("rabbit_access_control", "list_users", new ListUsersConverter());
|
||||
registerConverter("rabbit", "status", new StatusConverter());
|
||||
registerConverter("rabbit_amqqueue", "info_all", new QueueInfoAllConverter());
|
||||
}
|
||||
|
||||
|
||||
protected void registerConverter(String module, String function, ErlangConverter listUsersConverter) {
|
||||
converterMap.put(generateKey(module, function), listUsersConverter);
|
||||
converterMap.put(generateKey(module, function), listUsersConverter);
|
||||
}
|
||||
|
||||
protected String generateKey(String module, String function) {
|
||||
@@ -85,55 +80,63 @@ public class RabbitControlErlangConverter extends SimpleErlangConverter implemen
|
||||
}
|
||||
|
||||
public class ListUsersConverter extends SimpleErlangConverter {
|
||||
|
||||
|
||||
public Object fromErlang(OtpErlangObject erlangObject) throws ErlangConversionException {
|
||||
|
||||
|
||||
List<String> users = new ArrayList<String>();
|
||||
if (erlangObject instanceof OtpErlangList) {
|
||||
OtpErlangList erlangList = (OtpErlangList) erlangObject;
|
||||
for (OtpErlangObject obj : erlangList) {
|
||||
if (obj instanceof OtpErlangBinary) {
|
||||
OtpErlangBinary binary = (OtpErlangBinary) obj;
|
||||
users.add(new String(binary.binaryValue()));
|
||||
String value = extractString(obj);
|
||||
if (value != null) {
|
||||
users.add(value);
|
||||
}
|
||||
}
|
||||
}
|
||||
return users;
|
||||
}
|
||||
}
|
||||
|
||||
private String extractString(OtpErlangObject obj) {
|
||||
|
||||
if (obj instanceof OtpErlangBinary) {
|
||||
OtpErlangBinary binary = (OtpErlangBinary) obj;
|
||||
return new String(binary.binaryValue());
|
||||
} else if (obj instanceof OtpErlangTuple) {
|
||||
OtpErlangTuple tuple = (OtpErlangTuple) obj;
|
||||
return extractString(tuple.elementAt(0));
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
public class StatusConverter extends SimpleErlangConverter {
|
||||
|
||||
|
||||
public Object fromErlang(OtpErlangObject erlangObject) throws ErlangConversionException {
|
||||
|
||||
|
||||
List<Application> applications = new ArrayList<Application>();
|
||||
List<Node> nodes = new ArrayList<Node>();
|
||||
List<Node> runningNodes = new ArrayList<Node>();
|
||||
if (erlangObject instanceof OtpErlangList) {
|
||||
OtpErlangList erlangList = (OtpErlangList) erlangObject;
|
||||
|
||||
OtpErlangTuple runningAppTuple = (OtpErlangTuple)erlangList.elementAt(0);
|
||||
OtpErlangList appList = (OtpErlangList)runningAppTuple.elementAt(1);
|
||||
OtpErlangTuple runningAppTuple = (OtpErlangTuple) erlangList.elementAt(0);
|
||||
OtpErlangList appList = (OtpErlangList) runningAppTuple.elementAt(1);
|
||||
extractApplications(applications, appList);
|
||||
|
||||
OtpErlangTuple nodesTuple = (OtpErlangTuple)erlangList.elementAt(1);
|
||||
OtpErlangList nodesList = (OtpErlangList)nodesTuple.elementAt(1);
|
||||
|
||||
OtpErlangTuple nodesTuple = (OtpErlangTuple) erlangList.elementAt(1);
|
||||
OtpErlangList nodesList = (OtpErlangList) nodesTuple.elementAt(1);
|
||||
extractNodes(nodes, nodesList);
|
||||
|
||||
|
||||
OtpErlangTuple runningNodesTuple = (OtpErlangTuple)erlangList.elementAt(2);
|
||||
nodesList = (OtpErlangList)runningNodesTuple.elementAt(1);
|
||||
|
||||
OtpErlangTuple runningNodesTuple = (OtpErlangTuple) erlangList.elementAt(2);
|
||||
nodesList = (OtpErlangList) runningNodesTuple.elementAt(1);
|
||||
extractNodes(runningNodes, nodesList);
|
||||
|
||||
|
||||
/*
|
||||
for (OtpErlangObject obj : erlangList) {
|
||||
if (obj instanceof OtpErlangBinary) {
|
||||
OtpErlangBinary binary = (OtpErlangBinary) obj;
|
||||
users.add(new String(binary.binaryValue()));
|
||||
}
|
||||
}*/
|
||||
* for (OtpErlangObject obj : erlangList) { if (obj instanceof OtpErlangBinary) { OtpErlangBinary binary
|
||||
* = (OtpErlangBinary) obj; users.add(new String(binary.binaryValue())); } }
|
||||
*/
|
||||
}
|
||||
|
||||
|
||||
return new RabbitStatus(applications, nodes, runningNodes);
|
||||
}
|
||||
|
||||
@@ -141,35 +144,32 @@ public class RabbitControlErlangConverter extends SimpleErlangConverter implemen
|
||||
for (OtpErlangObject erlangNodeName : nodesList) {
|
||||
String nodeName = erlangNodeName.toString();
|
||||
nodes.add(new Node(nodeName));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void extractApplications(List<Application> applications, OtpErlangList appList) {
|
||||
for (OtpErlangObject appDescription : appList) {
|
||||
OtpErlangTuple appDescriptionTuple = (OtpErlangTuple)appDescription;
|
||||
OtpErlangTuple appDescriptionTuple = (OtpErlangTuple) appDescription;
|
||||
String name = appDescriptionTuple.elementAt(0).toString();
|
||||
String description = appDescriptionTuple.elementAt(1).toString();
|
||||
String version = appDescriptionTuple.elementAt(2).toString();
|
||||
applications.add(new Application(name, description, version));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
public enum QueueInfoField {
|
||||
transactions, acks_uncommitted, consumers, pid, durable, messages, memory, auto_delete, messages_ready,
|
||||
arguments, name, messages_unacknowledged, messages_uncommitted, NOVALUE;
|
||||
|
||||
public static QueueInfoField toQueueInfoField(String str)
|
||||
{
|
||||
try {
|
||||
return valueOf(str);
|
||||
}
|
||||
catch (Exception ex) {
|
||||
return NOVALUE;
|
||||
}
|
||||
}
|
||||
transactions, acks_uncommitted, consumers, pid, durable, messages, memory, auto_delete, messages_ready, arguments, name, messages_unacknowledged, messages_uncommitted, NOVALUE;
|
||||
|
||||
public static QueueInfoField toQueueInfoField(String str) {
|
||||
try {
|
||||
return valueOf(str);
|
||||
} catch (Exception ex) {
|
||||
return NOVALUE;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
public class QueueInfoAllConverter extends SimpleErlangConverter {
|
||||
|
||||
@Override
|
||||
@@ -181,13 +181,13 @@ public class RabbitControlErlangConverter extends SimpleErlangConverter implemen
|
||||
QueueInfo queueInfo = new QueueInfo();
|
||||
OtpErlangList itemList = (OtpErlangList) element;
|
||||
for (OtpErlangObject item : itemList.elements()) {
|
||||
OtpErlangTuple tuple = (OtpErlangTuple) item;
|
||||
OtpErlangTuple tuple = (OtpErlangTuple) item;
|
||||
if (tuple.arity() == 2) {
|
||||
String key = tuple.elementAt(0).toString();
|
||||
OtpErlangObject value = tuple.elementAt(1);
|
||||
switch (QueueInfoField.toQueueInfoField(key)) {
|
||||
case name:
|
||||
queueInfo.setName(extractNameValueFromTuple((OtpErlangTuple)value));
|
||||
queueInfo.setName(extractNameValueFromTuple((OtpErlangTuple) value));
|
||||
break;
|
||||
case transactions:
|
||||
queueInfo.setTransactions(extractLong(value));
|
||||
@@ -217,7 +217,7 @@ public class RabbitControlErlangConverter extends SimpleErlangConverter implemen
|
||||
queueInfo.setMessagesReady(extractLong(value));
|
||||
break;
|
||||
case arguments:
|
||||
OtpErlangList list = (OtpErlangList)value;
|
||||
OtpErlangList list = (OtpErlangList) value;
|
||||
if (list != null) {
|
||||
String[] args = new String[list.arity()];
|
||||
for (int i = 0; i < list.arity(); i++) {
|
||||
@@ -226,7 +226,7 @@ public class RabbitControlErlangConverter extends SimpleErlangConverter implemen
|
||||
}
|
||||
queueInfo.setArguments(args);
|
||||
}
|
||||
break;
|
||||
break;
|
||||
case messages_unacknowledged:
|
||||
queueInfo.setMessagesUnacknowledged(extractLong(value));
|
||||
break;
|
||||
@@ -235,16 +235,16 @@ public class RabbitControlErlangConverter extends SimpleErlangConverter implemen
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
queueInfoList.add(queueInfo);
|
||||
queueInfoList.add(queueInfo);
|
||||
}
|
||||
}
|
||||
return queueInfoList;
|
||||
}
|
||||
|
||||
private boolean extractAtomBoolean(OtpErlangObject value) {
|
||||
private boolean extractAtomBoolean(OtpErlangObject value) {
|
||||
return ((OtpErlangAtom) value).booleanValue();
|
||||
}
|
||||
|
||||
|
||||
@@ -6,7 +6,6 @@ import java.net.UnknownHostException;
|
||||
|
||||
import org.junit.Ignore;
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.erlang.connection.SimpleConnectionFactory;
|
||||
import org.springframework.erlang.core.ErlangTemplate;
|
||||
|
||||
|
||||
@@ -17,55 +17,71 @@
|
||||
package org.springframework.amqp.rabbit.admin;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.junit.AfterClass;
|
||||
import org.junit.BeforeClass;
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.amqp.core.Queue;
|
||||
import org.springframework.amqp.rabbit.connection.SingleConnectionFactory;
|
||||
import org.springframework.erlang.OtpIOException;
|
||||
|
||||
/**
|
||||
* @author Mark Pollack
|
||||
*/
|
||||
public class RabbitBrokerAdminIntegrationTests {
|
||||
|
||||
private static Log logger = LogFactory.getLog(RabbitBrokerAdminIntegrationTests.class);
|
||||
|
||||
private static RabbitBrokerAdmin brokerAdmin;
|
||||
|
||||
private static SingleConnectionFactory connectionFactory;
|
||||
|
||||
@BeforeClass
|
||||
public static void setUp() {
|
||||
public static void setUp() throws Exception {
|
||||
connectionFactory = new SingleConnectionFactory();
|
||||
connectionFactory.setUsername("guest");
|
||||
connectionFactory.setPassword("guest");
|
||||
brokerAdmin = new RabbitBrokerAdmin(connectionFactory);
|
||||
logger.info("Starting broker node");
|
||||
brokerAdmin.startNode();
|
||||
Thread.sleep(1000L);
|
||||
}
|
||||
|
||||
@AfterClass
|
||||
public static void tearDown() {
|
||||
logger.info("Shutting down broker node");
|
||||
brokerAdmin.stopNode();
|
||||
}
|
||||
|
||||
@Test
|
||||
//@Ignore
|
||||
public void integrationTestsUserCrud() {
|
||||
public void integrationTestsUserCrud() throws Exception {
|
||||
List<String> users = brokerAdmin.listUsers();
|
||||
if (users.contains("joe")) {
|
||||
brokerAdmin.deleteUser("joe");
|
||||
}
|
||||
Thread.sleep(1000L);
|
||||
brokerAdmin.addUser("joe", "trader");
|
||||
Thread.sleep(1000L);
|
||||
brokerAdmin.changeUserPassword("joe", "sales");
|
||||
Thread.sleep(1000L);
|
||||
users = brokerAdmin.listUsers();
|
||||
if (users.contains("joe")) {
|
||||
Thread.sleep(1000L);
|
||||
brokerAdmin.deleteUser("joe");
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
public void integrationTestListUsers() {
|
||||
public void integrationTestListUsers() throws Exception {
|
||||
// OtpErlangObject result =
|
||||
// adminTemplate.getErlangTemplate().executeRpc("rabbit_amqqueue",
|
||||
// "info_all", "/".getBytes());
|
||||
// System.out.println(result);
|
||||
Thread.sleep(1000L);
|
||||
List<String> users = brokerAdmin.listUsers();
|
||||
System.out.println(users);
|
||||
}
|
||||
@@ -80,44 +96,12 @@ public class RabbitBrokerAdminIntegrationTests {
|
||||
}
|
||||
|
||||
@Test
|
||||
//
|
||||
public void testStatusAndBrokerLifecycle() {
|
||||
RabbitStatus status = brokerAdmin.getStatus();
|
||||
assertBrokerAppRunning(status);
|
||||
|
||||
brokerAdmin.stopBrokerApplication();
|
||||
status = brokerAdmin.getStatus();
|
||||
assertEquals(0, status.getRunningNodes().size());
|
||||
|
||||
brokerAdmin.startBrokerApplication();
|
||||
status = brokerAdmin.getStatus();
|
||||
assertBrokerAppRunning(status);
|
||||
}
|
||||
|
||||
@Test
|
||||
//@Ignore("NEEDS RABBITMQ_HOME to be set.")
|
||||
public void testStartNode() {
|
||||
try {
|
||||
brokerAdmin.stopNode();
|
||||
} catch (OtpIOException e) {
|
||||
//assume it is not running.
|
||||
}
|
||||
brokerAdmin.startNode();
|
||||
assertEquals(1,1);
|
||||
brokerAdmin.stopNode();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testGetQueues() {
|
||||
public void testGetQueues() throws Exception {
|
||||
Thread.sleep(1000L);
|
||||
brokerAdmin.declareQueue(new Queue("test.queue"));
|
||||
assertEquals("/", connectionFactory.getVirtualHost());
|
||||
List<QueueInfo> queues = brokerAdmin.getQueues();
|
||||
assertEquals("test.queue", queues.get(0).getName());
|
||||
}
|
||||
|
||||
private void assertBrokerAppRunning(RabbitStatus status) {
|
||||
assertEquals(1, status.getRunningNodes().size());
|
||||
assertTrue(status.getRunningNodes().get(0).getName().contains("rabbit"));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -0,0 +1,83 @@
|
||||
/*
|
||||
* Copyright 2002-2010 the original author or authors.
|
||||
*
|
||||
* 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 org.springframework.amqp.rabbit.admin;
|
||||
|
||||
import static org.junit.Assert.assertFalse;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.springframework.amqp.rabbit.connection.SingleConnectionFactory;
|
||||
import org.springframework.erlang.OtpIOException;
|
||||
|
||||
/**
|
||||
* @author Mark Pollack
|
||||
* @author Dave Syer
|
||||
*/
|
||||
public class RabbitBrokerAdminStopIntegrationTests {
|
||||
|
||||
@Test
|
||||
// @Ignore("NEEDS RABBITMQ_HOME to be set.")
|
||||
public void testStartNode() throws Exception {
|
||||
|
||||
RabbitBrokerAdmin brokerAdmin;
|
||||
|
||||
SingleConnectionFactory connectionFactory;
|
||||
|
||||
connectionFactory = new SingleConnectionFactory();
|
||||
connectionFactory.setUsername("guest");
|
||||
connectionFactory.setPassword("guest");
|
||||
brokerAdmin = new RabbitBrokerAdmin(connectionFactory);
|
||||
|
||||
RabbitStatus status = brokerAdmin.getStatus();
|
||||
try {
|
||||
// Stop it if it is already running
|
||||
if (status.getRunningApplications().size() > 0) {
|
||||
brokerAdmin.stopBrokerApplication();
|
||||
Thread.sleep(1000L);
|
||||
}
|
||||
} catch (OtpIOException e) {
|
||||
// Not useful for test.
|
||||
}
|
||||
status = brokerAdmin.getStatus();
|
||||
if (status.getNodes().isEmpty()) {
|
||||
brokerAdmin.startNode();
|
||||
} else {
|
||||
brokerAdmin.startBrokerApplication();
|
||||
}
|
||||
Thread.sleep(1000L);
|
||||
brokerAdmin.startBrokerApplication();
|
||||
status = brokerAdmin.getStatus();
|
||||
|
||||
assertFalse("Broker node did not start. Check logs for hints.", status
|
||||
.getNodes().isEmpty());
|
||||
|
||||
try {
|
||||
assertFalse("Broker node not running. Check logs for hints.",
|
||||
status.getRunningNodes().isEmpty());
|
||||
assertFalse(
|
||||
"Broker application not running. Check logs for hints.",
|
||||
status.getRunningApplications().isEmpty());
|
||||
|
||||
// assertEquals(1, 1);
|
||||
brokerAdmin.stopBrokerApplication();
|
||||
Thread.sleep(1000L);
|
||||
} finally {
|
||||
brokerAdmin.stopNode();
|
||||
Thread.sleep(2000L);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
@@ -3,7 +3,6 @@ package org.springframework.amqp.rabbit.test;
|
||||
import org.junit.Ignore;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
|
||||
@RunWith(SpringRabbitJUnit4ClassRunner.class)
|
||||
@@ -20,7 +19,7 @@ public class RabbitTestExecutionListenerIntegrationTests {
|
||||
}
|
||||
|
||||
@Test
|
||||
public void doNothinAgain() throws InterruptedException {
|
||||
public void doNothingAgain() throws InterruptedException {
|
||||
Thread.sleep(1000);
|
||||
System.out.println("inside DO AGAIN");
|
||||
System.out.println("inside DO AGAIN");
|
||||
|
||||
9
spring-rabbit-admin/src/test/resources/log4j.properties
Normal file
9
spring-rabbit-admin/src/test/resources/log4j.properties
Normal file
@@ -0,0 +1,9 @@
|
||||
log4j.rootCategory=INFO, stdout
|
||||
|
||||
log4j.appender.stdout=org.apache.log4j.ConsoleAppender
|
||||
log4j.appender.stdout.layout=org.apache.log4j.PatternLayout
|
||||
log4j.appender.stdout.layout.ConversionPattern=%d %p [%c] - <%m>%n
|
||||
|
||||
log4j.category.org.springframework.amqp=DEBUG
|
||||
log4j.category.org.springframework.beans.factory=INFO
|
||||
|
||||
@@ -28,7 +28,6 @@ import java.util.List;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
|
||||
import org.springframework.amqp.rabbit.support.RabbitUtils;
|
||||
import org.springframework.beans.factory.DisposableBean;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
@@ -21,7 +21,6 @@ import java.util.Collection;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
|
||||
import org.springframework.amqp.core.AmqpAdmin;
|
||||
import org.springframework.amqp.core.Binding;
|
||||
import org.springframework.amqp.core.Exchange;
|
||||
@@ -33,8 +32,8 @@ import org.springframework.context.SmartLifecycle;
|
||||
import org.springframework.jmx.export.annotation.ManagedOperation;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
import com.rabbitmq.client.Channel;
|
||||
import com.rabbitmq.client.AMQP.Queue.DeclareOk;
|
||||
import com.rabbitmq.client.Channel;
|
||||
|
||||
/**
|
||||
* RabbitMQ implementation of portable AMQP administrative operations for AMQP >= 0.9.1
|
||||
@@ -161,7 +160,7 @@ public class RabbitAdmin implements AmqpAdmin, ApplicationContextAware, SmartLif
|
||||
public void purgeQueue(final String queueName, final boolean noWait) {
|
||||
this.rabbitTemplate.execute(new ChannelCallback<Object>() {
|
||||
public Object doInRabbit(Channel channel) throws Exception {
|
||||
channel.queuePurge(queueName, noWait);
|
||||
channel.queuePurge(queueName);
|
||||
return null;
|
||||
}
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user