Repository: brooklyn-server
Updated Branches:
  refs/heads/master c2326132d -> 5ea1b1b9d


BROOKLYN-440: ssh not use StreamGobbler for logging stdout


Project: http://git-wip-us.apache.org/repos/asf/brooklyn-server/repo
Commit: http://git-wip-us.apache.org/repos/asf/brooklyn-server/commit/e45a28fd
Tree: http://git-wip-us.apache.org/repos/asf/brooklyn-server/tree/e45a28fd
Diff: http://git-wip-us.apache.org/repos/asf/brooklyn-server/diff/e45a28fd

Branch: refs/heads/master
Commit: e45a28fd1ce54c9229f65af2ea7685edb9fff7a2
Parents: 78f41f8
Author: Aled Sage <[email protected]>
Authored: Sun Jun 11 00:19:22 2017 +0100
Committer: Aled Sage <[email protected]>
Committed: Sun Jun 11 11:27:59 2017 +0100

----------------------------------------------------------------------
 .../system/internal/ExecWithLoggingHelpers.java | 112 ++++++----------
 .../ssh/SshMachineLocationIntegrationTest.java  |  36 +++++
 .../location/ssh/SshMachineLocationTest.java    |  51 ++++++-
 .../core/internal/ssh/RecordingSshTool.java     |   8 +-
 .../org/apache/brooklyn/test/LogWatcher.java    |  14 +-
 utils/common/pom.xml                            |   5 +
 .../util/stream/LoggingOutputStream.java        | 134 +++++++++++++++++++
 .../brooklyn/util/stream/StreamGobbler.java     |   2 +
 .../util/stream/LoggingOutputStreamTest.java    |  94 +++++++++++++
 9 files changed, 376 insertions(+), 80 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/brooklyn-server/blob/e45a28fd/core/src/main/java/org/apache/brooklyn/util/core/task/system/internal/ExecWithLoggingHelpers.java
----------------------------------------------------------------------
diff --git 
a/core/src/main/java/org/apache/brooklyn/util/core/task/system/internal/ExecWithLoggingHelpers.java
 
b/core/src/main/java/org/apache/brooklyn/util/core/task/system/internal/ExecWithLoggingHelpers.java
index b3cdf67..6185c3a 100644
--- 
a/core/src/main/java/org/apache/brooklyn/util/core/task/system/internal/ExecWithLoggingHelpers.java
+++ 
b/core/src/main/java/org/apache/brooklyn/util/core/task/system/internal/ExecWithLoggingHelpers.java
@@ -18,14 +18,10 @@
  */
 package org.apache.brooklyn.util.core.task.system.internal;
 
-import java.io.IOException;
 import java.io.OutputStream;
-import java.io.PipedInputStream;
-import java.io.PipedOutputStream;
 import java.util.List;
 import java.util.Map;
 
-import org.slf4j.Logger;
 import org.apache.brooklyn.config.ConfigKey;
 import org.apache.brooklyn.core.config.Sanitizer;
 import org.apache.brooklyn.location.ssh.SshMachineLocation;
@@ -35,12 +31,11 @@ import org.apache.brooklyn.util.core.flags.TypeCoercions;
 import org.apache.brooklyn.util.core.internal.ssh.ShellAbstractTool;
 import org.apache.brooklyn.util.core.internal.ssh.ShellTool;
 import org.apache.brooklyn.util.core.task.Tasks;
-import org.apache.brooklyn.util.stream.StreamGobbler;
-import org.apache.brooklyn.util.stream.Streams;
+import org.apache.brooklyn.util.stream.LoggingOutputStream;
 import org.apache.brooklyn.util.text.Strings;
+import org.slf4j.Logger;
 
 import com.google.common.base.Function;
-import com.google.common.base.Throwables;
 
 public abstract class ExecWithLoggingHelpers {
 
@@ -106,7 +101,6 @@ public abstract class ExecWithLoggingHelpers {
         return execWithLogging(props, summaryForLogging, commands, env, null, 
execCommand);
     }
     
-    @SuppressWarnings("resource")
     public int execWithLogging(Map<String,?> props, final String 
summaryForLogging, final List<String> commands,
             final Map<String,?> env, String expectedCommandHeaders, final 
ExecRunner execCommand) {
         if (commandLogger!=null && commandLogger.isDebugEnabled()) {
@@ -129,73 +123,47 @@ public abstract class ExecWithLoggingHelpers {
 
         execFlags.configure(ShellTool.PROP_SUMMARY, summaryForLogging);
         
-        PipedOutputStream outO = null;
-        PipedOutputStream outE = null;
-        StreamGobbler gO=null, gE=null;
+        preExecChecks();
+        
+        String logPrefix = execFlags.get(LOG_PREFIX);
+        if (logPrefix==null) logPrefix = 
constructDefaultLoggingPrefix(execFlags);
+
+        if (!execFlags.get(NO_STDOUT_LOGGING)) {
+            String stdoutLogPrefix = "["+(logPrefix != null ? 
logPrefix+":stdout" : "stdout")+"] ";
+            OutputStream outO = LoggingOutputStream.builder()
+                    .outputStream(execFlags.get(STDOUT))
+                    .logger(commandLogger)
+                    .logPrefix(stdoutLogPrefix)
+                    .build();
+
+            execFlags.put(STDOUT, outO);
+        }
+
+        if (!execFlags.get(NO_STDERR_LOGGING)) {
+            String stderrLogPrefix = "["+(logPrefix != null ? 
logPrefix+":stderr" : "stderr")+"] ";
+            OutputStream outE = LoggingOutputStream.builder()
+                    .outputStream(execFlags.get(STDERR))
+                    .logger(commandLogger)
+                    .logPrefix(stderrLogPrefix)
+                    .build();
+            execFlags.put(STDERR, outE);
+        }
+
+        Tasks.setBlockingDetails(shortName+" executing, "+summaryForLogging);
         try {
-            preExecChecks();
-            
-            String logPrefix = execFlags.get(LOG_PREFIX);
-            if (logPrefix==null) logPrefix = 
constructDefaultLoggingPrefix(execFlags);
-
-            if (!execFlags.get(NO_STDOUT_LOGGING)) {
-                PipedInputStream insO = new PipedInputStream();
-                outO = new PipedOutputStream(insO);
-
-                String stdoutLogPrefix = "["+(logPrefix != null ? 
logPrefix+":stdout" : "stdout")+"] ";
-                gO = new StreamGobbler(insO, execFlags.get(STDOUT), 
commandLogger).setLogPrefix(stdoutLogPrefix);
-                gO.start();
-
-                execFlags.put(STDOUT, outO);
-            }
-
-            if (!execFlags.get(NO_STDERR_LOGGING)) {
-                PipedInputStream insE = new PipedInputStream();
-                outE = new PipedOutputStream(insE);
-
-                String stderrLogPrefix = "["+(logPrefix != null ? 
logPrefix+":stderr" : "stderr")+"] ";
-                gE = new StreamGobbler(insE, execFlags.get(STDERR), 
commandLogger).setLogPrefix(stderrLogPrefix);
-                gE.start();
-
-                execFlags.put(STDERR, outE);
-            }
-
-            Tasks.setBlockingDetails(shortName+" executing, 
"+summaryForLogging);
-            try {
-                return 
execWithTool(MutableMap.copyOf(toolFlags.getAllConfig()), new 
Function<ShellTool, Integer>() {
-                    @Override
-                    public Integer apply(ShellTool tool) {
-                        int result = execCommand.exec(tool, 
MutableMap.copyOf(execFlags.getAllConfig()), commands, env);
-                        if (commandLogger!=null && 
commandLogger.isDebugEnabled()) 
-                            commandLogger.debug("{}, on machine {}, completed: 
return status {}",
-                                    new Object[] {summaryForLogging, 
getTargetName(), result});
-                        return result;
-                    }});
-
-            } finally {
-                Tasks.setBlockingDetails(null);
-            }
-
-        } catch (IOException e) {
-            if (commandLogger!=null && commandLogger.isDebugEnabled()) 
-                commandLogger.debug("{}, on machine {}, failed: {}", new 
Object[] {summaryForLogging, getTargetName(), e});
-            throw Throwables.propagate(e);
+            return execWithTool(MutableMap.copyOf(toolFlags.getAllConfig()), 
new Function<ShellTool, Integer>() {
+                @Override
+                public Integer apply(ShellTool tool) {
+                    int result = execCommand.exec(tool, 
MutableMap.copyOf(execFlags.getAllConfig()), commands, env);
+                    if (commandLogger!=null && commandLogger.isDebugEnabled()) 
+                        commandLogger.debug("{}, on machine {}, completed: 
return status {}",
+                                new Object[] {summaryForLogging, 
getTargetName(), result});
+                    return result;
+                }});
+
         } finally {
-            // Must close the pipedOutStreams, otherwise input will never read 
-1 so StreamGobbler thread would never die
-            if (outO!=null) try { outO.flush(); } catch (IOException e) {}
-            if (outE!=null) try { outE.flush(); } catch (IOException e) {}
-            Streams.closeQuietly(outO);
-            Streams.closeQuietly(outE);
-
-            try {
-                if (gE!=null) { gE.join(); }
-                if (gO!=null) { gO.join(); }
-            } catch (InterruptedException e) {
-                Thread.currentThread().interrupt();
-                Throwables.propagate(e);
-            }
+            Tasks.setBlockingDetails(null);
         }
-
     }
 
 }

http://git-wip-us.apache.org/repos/asf/brooklyn-server/blob/e45a28fd/core/src/test/java/org/apache/brooklyn/location/ssh/SshMachineLocationIntegrationTest.java
----------------------------------------------------------------------
diff --git 
a/core/src/test/java/org/apache/brooklyn/location/ssh/SshMachineLocationIntegrationTest.java
 
b/core/src/test/java/org/apache/brooklyn/location/ssh/SshMachineLocationIntegrationTest.java
index c770527..6f7b324 100644
--- 
a/core/src/test/java/org/apache/brooklyn/location/ssh/SshMachineLocationIntegrationTest.java
+++ 
b/core/src/test/java/org/apache/brooklyn/location/ssh/SshMachineLocationIntegrationTest.java
@@ -29,16 +29,20 @@ import java.io.OutputStream;
 import java.net.InetAddress;
 import java.security.KeyPair;
 import java.util.Arrays;
+import java.util.List;
 import java.util.Map;
 import java.util.concurrent.Callable;
 
 import org.apache.brooklyn.api.location.Location;
 import org.apache.brooklyn.api.location.LocationSpec;
 import org.apache.brooklyn.api.location.MachineDetails;
+import org.apache.brooklyn.core.BrooklynLogging;
 import org.apache.brooklyn.core.entity.AbstractEntity;
 import org.apache.brooklyn.core.internal.BrooklynProperties;
 import 
org.apache.brooklyn.location.localhost.LocalhostMachineProvisioningLocation;
 import org.apache.brooklyn.test.Asserts;
+import org.apache.brooklyn.test.LogWatcher;
+import org.apache.brooklyn.test.LogWatcher.EventPredicates;
 import org.apache.brooklyn.util.collections.MutableMap;
 import org.apache.brooklyn.util.core.crypto.SecureKeys;
 import org.apache.brooklyn.util.core.file.ArchiveUtils;
@@ -61,10 +65,14 @@ import org.testng.annotations.Test;
 
 import com.google.common.base.Charsets;
 import com.google.common.base.Preconditions;
+import com.google.common.base.Predicate;
+import com.google.common.base.Predicates;
 import com.google.common.collect.ImmutableList;
 import com.google.common.collect.ImmutableMap;
 import com.google.common.io.Files;
 
+import ch.qos.logback.classic.spi.ILoggingEvent;
+
 public class SshMachineLocationIntegrationTest extends SshMachineLocationTest {
 
     private static final Logger LOG = 
LoggerFactory.getLogger(AbstractEntity.class);
@@ -308,4 +316,32 @@ public class SshMachineLocationIntegrationTest extends 
SshMachineLocationTest {
         int rc = sm.execScript("Test script directory execution", 
ImmutableList.of(command));
         assertEquals(rc, 0);
     }
+    
+    @Test(groups="Integration")
+    public void testLogsStdoutAndStderr() {
+        List<String> loggerNames = ImmutableList.of(
+                SshMachineLocation.class.getName(), 
+                BrooklynLogging.SSH_IO, 
+                SshjTool.class.getName());
+        ch.qos.logback.classic.Level logLevel = 
ch.qos.logback.classic.Level.DEBUG;
+        Predicate<ILoggingEvent> filter = Predicates.alwaysTrue();
+        LogWatcher watcher = new LogWatcher(loggerNames, logLevel, filter);
+
+        watcher.start();
+        try {
+            host.execCommands("mySummary", ImmutableList.of("echo mystdout1", 
"echo mystdout2", "echo mystderr1 1>&2", "echo mystderr2 1>&2"));
+            
+            
watcher.assertHasEvent(EventPredicates.containsMessage(":22:stdout] 
mystdout1"));
+            
watcher.assertHasEvent(EventPredicates.containsMessage(":22:stdout] 
mystdout2"));
+            
watcher.assertHasEvent(EventPredicates.containsMessage(":22:stderr] 
mystderr1"));
+            
watcher.assertHasEvent(EventPredicates.containsMessage(":22:stderr] 
mystderr2"));
+        } finally {
+            watcher.close();
+        }
+    }
+    
+    @Test(groups="Integration")
+    public void testTurningOffLoggingStdoutAndStderr() {
+        super.testTurningOffLoggingStdoutAndStderr();
+    }
 }

http://git-wip-us.apache.org/repos/asf/brooklyn-server/blob/e45a28fd/core/src/test/java/org/apache/brooklyn/location/ssh/SshMachineLocationTest.java
----------------------------------------------------------------------
diff --git 
a/core/src/test/java/org/apache/brooklyn/location/ssh/SshMachineLocationTest.java
 
b/core/src/test/java/org/apache/brooklyn/location/ssh/SshMachineLocationTest.java
index d2f8645..e610fd2 100644
--- 
a/core/src/test/java/org/apache/brooklyn/location/ssh/SshMachineLocationTest.java
+++ 
b/core/src/test/java/org/apache/brooklyn/location/ssh/SshMachineLocationTest.java
@@ -114,7 +114,7 @@ public class SshMachineLocationTest extends 
BrooklynAppUnitTestSupport {
 
     protected SshMachineLocation newHost() {
         return 
mgmt.getLocationManager().createLocation(LocationSpec.create(SshMachineLocation.class)
-                .configure("address", Networking.getLocalHost())
+                .configure("address", "1.2.3.4")
                 .configure(SshMachineLocation.SSH_TOOL_CLASS, 
RecordingSshTool.class.getName()));
     }
     
@@ -342,4 +342,53 @@ public class SshMachineLocationTest extends 
BrooklynAppUnitTestSupport {
             watcher.close();
         }
     }
+    
+    @Test
+    public void testLogsStdoutAndStderr() {
+        RecordingSshTool.setCustomResponse(".*mycommand.*", new 
CustomResponse(0, "mystdout1\nmystdout2", "mystderr1\nmystderr2"));
+        List<String> loggerNames = ImmutableList.of(
+                SshMachineLocation.class.getName(), 
+                BrooklynLogging.SSH_IO, 
+                SshjTool.class.getName());
+        ch.qos.logback.classic.Level logLevel = 
ch.qos.logback.classic.Level.DEBUG;
+        Predicate<ILoggingEvent> filter = Predicates.alwaysTrue();
+        LogWatcher watcher = new LogWatcher(loggerNames, logLevel, filter);
+
+        watcher.start();
+        try {
+            host.execCommands("mySummary", ImmutableList.of("mycommand"));
+            
+            
watcher.assertHasEvent(EventPredicates.containsMessage("[1.2.3.4:22:stdout] 
mystdout1"));
+            
watcher.assertHasEvent(EventPredicates.containsMessage("[1.2.3.4:22:stdout] 
mystdout2"));
+            
watcher.assertHasEvent(EventPredicates.containsMessage("[1.2.3.4:22:stderr] 
mystderr1"));
+            
watcher.assertHasEvent(EventPredicates.containsMessage("[1.2.3.4:22:stderr] 
mystderr2"));
+        } finally {
+            watcher.close();
+        }
+    }
+    
+    @Test
+    public void testTurningOffLoggingStdoutAndStderr() {
+        RecordingSshTool.setCustomResponse(".*mycommand.*", new 
CustomResponse(0, "mystdout1\nmystdout2", "mystderr1\nmystderr2"));
+        List<String> loggerNames = ImmutableList.of(
+                SshMachineLocation.class.getName(), 
+                BrooklynLogging.SSH_IO, 
+                SshjTool.class.getName());
+        ch.qos.logback.classic.Level logLevel = 
ch.qos.logback.classic.Level.DEBUG;
+        Predicate<ILoggingEvent> filter = Predicates.alwaysTrue();
+        LogWatcher watcher = new LogWatcher(loggerNames, logLevel, filter);
+
+        watcher.start();
+        try {
+            host.execCommands(
+                    
ImmutableMap.of(SshMachineLocation.NO_STDOUT_LOGGING.getName(), true, 
SshMachineLocation.NO_STDERR_LOGGING.getName(), true), 
+                    "mySummary", 
+                    ImmutableList.of("mycommand"));
+            
+            assertFalse(Iterables.tryFind(watcher.getEvents(), 
EventPredicates.containsMessage(":stdout]")).isPresent());
+            assertFalse(Iterables.tryFind(watcher.getEvents(), 
EventPredicates.containsMessage(":stderr]")).isPresent());
+        } finally {
+            watcher.close();
+        }
+    }
 }

http://git-wip-us.apache.org/repos/asf/brooklyn-server/blob/e45a28fd/core/src/test/java/org/apache/brooklyn/util/core/internal/ssh/RecordingSshTool.java
----------------------------------------------------------------------
diff --git 
a/core/src/test/java/org/apache/brooklyn/util/core/internal/ssh/RecordingSshTool.java
 
b/core/src/test/java/org/apache/brooklyn/util/core/internal/ssh/RecordingSshTool.java
index 82524d3..1266d45 100644
--- 
a/core/src/test/java/org/apache/brooklyn/util/core/internal/ssh/RecordingSshTool.java
+++ 
b/core/src/test/java/org/apache/brooklyn/util/core/internal/ssh/RecordingSshTool.java
@@ -261,10 +261,14 @@ public class RecordingSshTool implements SshTool {
     protected void writeCustomResponseStreams(Map<String, ?> props, 
CustomResponse response) {
         try {
             if (Strings.isNonBlank(response.stdout) && 
props.get(SshTool.PROP_OUT_STREAM.getName()) != null) {
-                
((OutputStream)props.get(SshTool.PROP_OUT_STREAM.getName())).write(response.stdout.getBytes());
+                OutputStream out = 
(OutputStream)props.get(SshTool.PROP_OUT_STREAM.getName());
+                out.write(response.stdout.getBytes());
+                out.flush();
             }
             if (Strings.isNonBlank(response.stderr) && 
props.get(SshTool.PROP_ERR_STREAM.getName()) != null) {
-                
((OutputStream)props.get(SshTool.PROP_ERR_STREAM.getName())).write(response.stderr.getBytes());
+                OutputStream err = 
(OutputStream)props.get(SshTool.PROP_ERR_STREAM.getName());
+                err.write(response.stderr.getBytes());
+                err.flush();
             }
         } catch (IOException e) {
             Exceptions.propagate(e);

http://git-wip-us.apache.org/repos/asf/brooklyn-server/blob/e45a28fd/test-support/src/main/java/org/apache/brooklyn/test/LogWatcher.java
----------------------------------------------------------------------
diff --git 
a/test-support/src/main/java/org/apache/brooklyn/test/LogWatcher.java 
b/test-support/src/main/java/org/apache/brooklyn/test/LogWatcher.java
index e72383e..36206c8 100644
--- a/test-support/src/main/java/org/apache/brooklyn/test/LogWatcher.java
+++ b/test-support/src/main/java/org/apache/brooklyn/test/LogWatcher.java
@@ -152,6 +152,14 @@ public class LogWatcher implements Closeable {
         assertFalse(events.isEmpty());
     }
 
+    public List<ILoggingEvent> assertHasEvent(final Predicate<? super 
ILoggingEvent> filter) {
+        synchronized (events) {
+            Iterable<ILoggingEvent> filtered = Iterables.filter(events, 
filter);
+            assertFalse(Iterables.isEmpty(filtered), "events="+events);
+            return ImmutableList.copyOf(filtered);
+        }
+    }
+
     public List<ILoggingEvent> assertHasEventEventually() {
         Asserts.succeedsEventually(new Runnable() {
             @Override
@@ -166,11 +174,7 @@ public class LogWatcher implements Closeable {
         Asserts.succeedsEventually(new Runnable() {
             @Override
             public void run() {
-                synchronized (events) {
-                    Iterable<ILoggingEvent> filtered = 
Iterables.filter(events, filter);
-                    assertFalse(Iterables.isEmpty(filtered));
-                    result.set(ImmutableList.copyOf(filtered));
-                }
+                result.set(assertHasEvent(filter));
             }});
         return result.get();
     }

http://git-wip-us.apache.org/repos/asf/brooklyn-server/blob/e45a28fd/utils/common/pom.xml
----------------------------------------------------------------------
diff --git a/utils/common/pom.xml b/utils/common/pom.xml
index e1d7475..e0d60e0 100644
--- a/utils/common/pom.xml
+++ b/utils/common/pom.xml
@@ -84,6 +84,11 @@
             <scope>test</scope>
         </dependency>
         <dependency>
+            <groupId>org.mockito</groupId>
+            <artifactId>mockito-core</artifactId>
+            <scope>test</scope>
+        </dependency>
+        <dependency>
             <groupId>org.apache.brooklyn</groupId>
             <artifactId>brooklyn-utils-test-support</artifactId>
             <version>${project.version}</version>

http://git-wip-us.apache.org/repos/asf/brooklyn-server/blob/e45a28fd/utils/common/src/main/java/org/apache/brooklyn/util/stream/LoggingOutputStream.java
----------------------------------------------------------------------
diff --git 
a/utils/common/src/main/java/org/apache/brooklyn/util/stream/LoggingOutputStream.java
 
b/utils/common/src/main/java/org/apache/brooklyn/util/stream/LoggingOutputStream.java
new file mode 100644
index 0000000..2f0e7bd
--- /dev/null
+++ 
b/utils/common/src/main/java/org/apache/brooklyn/util/stream/LoggingOutputStream.java
@@ -0,0 +1,134 @@
+/*
+ * 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.brooklyn.util.stream;
+
+import java.io.FilterOutputStream;
+import java.io.IOException;
+import java.io.OutputStream;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import org.slf4j.Logger;
+
+/**
+ * Wraps another output stream, intercepting the writes to log it.
+ */
+public class LoggingOutputStream extends FilterOutputStream {
+
+    private static final OutputStream NOOP_OUTPUT_STREAM = new 
FilterOutputStream(null) {
+        @Override public void write(int b) throws IOException {
+        }
+        @Override public void flush() throws IOException {
+        }
+        @Override public void close() throws IOException {
+        }        
+    };
+    
+    public static Builder builder() {
+        return new Builder();
+    }
+    
+    public static class Builder {
+        OutputStream out;
+        Logger log;
+        String logPrefix;
+        
+        public Builder outputStream(OutputStream val) {
+            this.out = val;
+            return this;
+        }
+        public Builder logger(Logger val) {
+            this.log = val;
+            return this;
+        }
+        public Builder logPrefix(String val) {
+            this.logPrefix = val;
+            return this;
+        }
+        public LoggingOutputStream build() {
+            return new LoggingOutputStream(this);
+        }
+    }
+    
+    protected final Logger log;
+    protected final String logPrefix;
+    private final AtomicBoolean running = new AtomicBoolean(true);
+    private final StringBuilder lineSoFar = new StringBuilder(16);
+
+    private LoggingOutputStream(Builder builder) {
+        super(builder.out != null ? builder.out : NOOP_OUTPUT_STREAM);
+        log = builder.log;
+        logPrefix = (builder.logPrefix != null) ? builder.logPrefix : "";
+      }
+
+    @Override
+    public void write(int b) throws IOException {
+        if (running.get() && b >= 0) onChar(b);
+        out.write(b);
+    }
+
+    @Override
+    public void flush() throws IOException {
+        try {
+            if (lineSoFar.length() > 0) {
+                onLine(lineSoFar.toString());
+                lineSoFar.setLength(0);
+            }
+        } finally {
+            super.flush();
+        }
+    }
+    
+    // Overriding close() because FilterOutputStream's close() method pre-JDK8 
has bad behavior:
+    // it silently ignores any exception thrown by flush(). Instead, just 
close the delegate stream.
+    // It should flush itself if necessary.
+    @Override
+    public void close() throws IOException {
+        try {
+            onLine(lineSoFar.toString());
+            lineSoFar.setLength(0);
+        } finally {
+            out.close();
+            running.set(false);
+        }
+    }
+    
+    public void onChar(int c) {
+        if (c=='\n' || c=='\r') {
+            if (lineSoFar.length()>0)
+                //suppress blank lines, so that we can treat either newline 
char as a line separator
+                //(eg to show curl updates frequently)
+                onLine(lineSoFar.toString());
+            lineSoFar.setLength(0);
+        } else {
+            lineSoFar.append((char)c);
+        }
+    }
+    
+    public void onLine(String line) {
+        //right trim, in case there is \r or other funnies
+        while (line.length()>0 && 
Character.isWhitespace(line.charAt(line.length()-1)))
+            line = line.substring(0, line.length()-1);
+        //right trim, in case there is \r or other funnies
+        while (line.length()>0 && (line.charAt(0)=='\n' || 
line.charAt(0)=='\r'))
+            line = line.substring(1);
+        if (!line.isEmpty()) {
+            if (log!=null && log.isDebugEnabled()) log.debug(logPrefix+line);
+        }
+    }
+}

http://git-wip-us.apache.org/repos/asf/brooklyn-server/blob/e45a28fd/utils/common/src/main/java/org/apache/brooklyn/util/stream/StreamGobbler.java
----------------------------------------------------------------------
diff --git 
a/utils/common/src/main/java/org/apache/brooklyn/util/stream/StreamGobbler.java 
b/utils/common/src/main/java/org/apache/brooklyn/util/stream/StreamGobbler.java
index a440336..6053fa0 100644
--- 
a/utils/common/src/main/java/org/apache/brooklyn/util/stream/StreamGobbler.java
+++ 
b/utils/common/src/main/java/org/apache/brooklyn/util/stream/StreamGobbler.java
@@ -86,6 +86,8 @@ public class StreamGobbler extends Thread implements 
Closeable {
             onClose();
             //TODO parametrise log level, for this error, and for normal 
messages
             if (log!=null && log.isTraceEnabled()) 
log.trace(logPrefix+"exception reading from stream ("+e+")");
+        } finally {
+            if (out != null) out.flush();
         }
     }
     

http://git-wip-us.apache.org/repos/asf/brooklyn-server/blob/e45a28fd/utils/common/src/test/java/org/apache/brooklyn/util/stream/LoggingOutputStreamTest.java
----------------------------------------------------------------------
diff --git 
a/utils/common/src/test/java/org/apache/brooklyn/util/stream/LoggingOutputStreamTest.java
 
b/utils/common/src/test/java/org/apache/brooklyn/util/stream/LoggingOutputStreamTest.java
new file mode 100644
index 0000000..5163bb6
--- /dev/null
+++ 
b/utils/common/src/test/java/org/apache/brooklyn/util/stream/LoggingOutputStreamTest.java
@@ -0,0 +1,94 @@
+/*
+ * 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.brooklyn.util.stream;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertTrue;
+
+import java.io.ByteArrayOutputStream;
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+
+import org.mockito.Mockito;
+import org.mockito.invocation.InvocationOnMock;
+import org.mockito.stubbing.Answer;
+import org.slf4j.Logger;
+import org.testng.annotations.BeforeMethod;
+import org.testng.annotations.Test;
+
+import com.google.common.collect.ImmutableList;
+
+public class LoggingOutputStreamTest {
+
+    private List<String> logs;
+    private Logger mockLogger;
+    
+    @BeforeMethod(alwaysRun=true)
+    public void setUp() throws Exception {
+        logs = new ArrayList<>();
+        mockLogger = Mockito.mock(Logger.class);
+        Mockito.when(mockLogger.isDebugEnabled()).thenReturn(true);
+        Mockito.doAnswer(new Answer<Void>() {
+            public Void answer(InvocationOnMock invocation) {
+              Object[] args = invocation.getArguments();
+              logs.add((String)args[0]);
+              return null;
+            }
+        }).when(mockLogger).debug(Mockito.anyString());
+    }
+    
+    @Test
+    public void testCallsDelegateStream() throws Exception {
+        ByteArrayOutputStream delegate = new ByteArrayOutputStream();
+        LoggingOutputStream out = 
LoggingOutputStream.builder().outputStream(delegate).build();
+        out.write(new byte[] {1,2,3});
+        out.flush();
+        assertTrue(Arrays.equals(delegate.toByteArray(), new byte[] {1, 2, 
3}));
+    }
+    
+    @Test
+    public void testNoopIfNoDelegateStream() throws Exception {
+        // Just checking that throws no exceptions
+        LoggingOutputStream out = LoggingOutputStream.builder().build();
+        out.write(new byte[] {1,2,3});
+        out.flush();
+    }
+    
+    @Test
+    public void testLogsLines() throws Exception {
+        LoggingOutputStream out = 
LoggingOutputStream.builder().logger(mockLogger).build();
+        out.write("line1\n".getBytes(StandardCharsets.UTF_8));
+        out.write("line2".getBytes(StandardCharsets.UTF_8));
+        out.flush();
+        
+        assertEquals(logs, ImmutableList.of("line1", "line2"));
+    }
+    
+    @Test
+    public void testLogsLinesWithPrefix() throws Exception {
+        LoggingOutputStream out = 
LoggingOutputStream.builder().logger(mockLogger).logPrefix("myprefix:").build();
+        out.write("line1\n".getBytes(StandardCharsets.UTF_8));
+        out.write("line2".getBytes(StandardCharsets.UTF_8));
+        out.flush();
+        
+        assertEquals(logs, ImmutableList.of("myprefix:line1", 
"myprefix:line2"));
+    }
+}

Reply via email to