+ * Usage: {@code HeapDumper.dumpHeap("/tmp/foo.hprof");} + *
+ * If the file exists already, it will be replaced. + *
+ * Courtesy: + * https://blogs.oracle.com/sundararajan/entry/programmatically_dumping_heap_from_java + *
+ * See https://docs.oracle.com/javase/8/docs/jre/api/management/extension/com/sun/management/HotSpotDiagnosticMXBean.html#dumpHeap-java.lang.String-boolean- + *+ * @author Gary Russell + * @since 4.2 + * + */ +public class HeapDumper { + + // This is the name of the HotSpot Diagnostic MBean + private static final String HOTSPOT_BEAN_NAME = "com.sun.management:type=HotSpotDiagnostic"; + + // field to store the hotspot diagnostic MBean + private static volatile HotSpotDiagnosticMXBean hotspotMBean; + + public static void dumpHeap(String fileName) { + dumpHeap(fileName, true); + } + + public static void dumpHeap(String fileName, boolean live) { + File file = new File(fileName); + if (file.exists()) { + file.delete(); + } + // initialize hotspot diagnostic MBean + initHotspotMBean(); + try { + hotspotMBean.dumpHeap(fileName, live); + } + catch (RuntimeException re) { + throw re; + } + catch (Exception exp) { + throw new RuntimeException(exp); + } + } + + // initialize the hotspot diagnostic MBean field + private static void initHotspotMBean() { + if (hotspotMBean == null) { + synchronized (Object.class) { + if (hotspotMBean == null) { + hotspotMBean = getHotspotMBean(); + } + } + } + } + + // get the hotspot diagnostic MBean from the + // platform MBean server + private static HotSpotDiagnosticMXBean getHotspotMBean() { + try { + MBeanServer server = ManagementFactory.getPlatformMBeanServer(); + HotSpotDiagnosticMXBean bean = ManagementFactory.newPlatformMXBeanProxy(server, HOTSPOT_BEAN_NAME, + HotSpotDiagnosticMXBean.class); + return bean; + } + catch (RuntimeException re) { + throw re; + } + catch (Exception exp) { + throw new RuntimeException(exp); + } + } + +} diff --git a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/ClientWebSocketContainer.java b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/ClientWebSocketContainer.java index a7774f7741..bdf76ff148 100644 --- a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/ClientWebSocketContainer.java +++ b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/ClientWebSocketContainer.java @@ -49,6 +49,8 @@ import org.springframework.web.socket.client.WebSocketClient; */ public final class ClientWebSocketContainer extends IntegrationWebSocketContainer implements SmartLifecycle { + private static final int DEFAULT_CONNECTION_TIMEOUT = 10; + private final WebSocketHttpHeaders headers = new WebSocketHttpHeaders(); private final ConnectionManagerSupport connectionManager; @@ -59,6 +61,8 @@ public final class ClientWebSocketContainer extends IntegrationWebSocketContaine private volatile Throwable openConnectionException; + private volatile int connectionTimeout = DEFAULT_CONNECTION_TIMEOUT; + public ClientWebSocketContainer(WebSocketClient client, String uriTemplate, Object... uriVariables) { Assert.notNull(client, "'client' must not be null"); this.connectionManager = new IntegrationWebSocketConnectionManager(client, uriTemplate, uriVariables); @@ -84,6 +88,15 @@ public final class ClientWebSocketContainer extends IntegrationWebSocketContaine this.headers.putAll(headers); } + /** + * Set the connection timeout in seconds; default: 10. + * @param connectionTimeout the timeout in seconds. + * @since 4.2 + */ + public void setConnectionTimeout(int connectionTimeout) { + this.connectionTimeout = connectionTimeout; + } + /** * Return the {@link #clientSession} {@link WebSocketSession}. * Independently of provided argument, this method always returns only the @@ -95,7 +108,7 @@ public final class ClientWebSocketContainer extends IntegrationWebSocketContaine public WebSocketSession getSession(String sessionId) { if (this.isRunning()) { try { - this.connectionLatch.await(10, TimeUnit.SECONDS); + this.connectionLatch.await(this.connectionTimeout, TimeUnit.SECONDS); } catch (InterruptedException e) { logger.error("'clientSession' has not been established during 'openConnection'"); diff --git a/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/ClientWebSocketContainerTests.java b/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/ClientWebSocketContainerTests.java index 96378e1802..c468dcc119 100644 --- a/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/ClientWebSocketContainerTests.java +++ b/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/ClientWebSocketContainerTests.java @@ -66,6 +66,7 @@ public class ClientWebSocketContainerTests { TestWebSocketListener messageListener = new TestWebSocketListener(); container.setMessageListener(messageListener); + container.setConnectionTimeout(30); container.start(); diff --git a/spring-integration-zookeeper/src/test/java/org/springframework/integration/zookeeper/event/ZookeeperLeaderTests.java b/spring-integration-zookeeper/src/test/java/org/springframework/integration/zookeeper/event/ZookeeperLeaderTests.java index db1a57fbfd..7bea91e03a 100644 --- a/spring-integration-zookeeper/src/test/java/org/springframework/integration/zookeeper/event/ZookeeperLeaderTests.java +++ b/spring-integration-zookeeper/src/test/java/org/springframework/integration/zookeeper/event/ZookeeperLeaderTests.java @@ -75,20 +75,20 @@ public class ZookeeperLeaderTests extends ZookeeperTestSupport { LeaderInitiator initiator2 = new LeaderInitiator(this.client, candidate2, "/sitest"); initiator2.setLeaderEventPublisher(publisher); initiator2.start(); - AbstractLeaderEvent event = this.events.poll(10, TimeUnit.SECONDS); + AbstractLeaderEvent event = this.events.poll(30, TimeUnit.SECONDS); assertNotNull(event); assertThat(event, instanceOf(OnGrantedEvent.class)); event.getContext().yield(); assertTrue(this.adapter.isRunning()); - event = this.events.poll(10, TimeUnit.SECONDS); + event = this.events.poll(30, TimeUnit.SECONDS); assertNotNull(event); assertThat(event, instanceOf(OnRevokedEvent.class)); assertFalse(this.adapter.isRunning()); - event = this.events.poll(10, TimeUnit.SECONDS); + event = this.events.poll(30, TimeUnit.SECONDS); assertNotNull(event); assertThat(event, instanceOf(OnGrantedEvent.class)); @@ -96,7 +96,7 @@ public class ZookeeperLeaderTests extends ZookeeperTestSupport { initiator1.stop(); initiator2.stop(); - event = this.events.poll(10, TimeUnit.SECONDS); + event = this.events.poll(30, TimeUnit.SECONDS); assertNotNull(event); assertThat(event, instanceOf(OnRevokedEvent.class));