001/**
002 * Licensed to the Apache Software Foundation (ASF) under one
003 * or more contributor license agreements.  See the NOTICE file
004 * distributed with this work for additional information
005 * regarding copyright ownership.  The ASF licenses this file
006 * to you under the Apache License, Version 2.0 (the
007 * "License"); you may not use this file except in compliance
008 * with the License.  You may obtain a copy of the License at
009 *
010 *      http://www.apache.org/licenses/LICENSE-2.0
011 *
012 * Unless required by applicable law or agreed to in writing, software
013 * distributed under the License is distributed on an "AS IS" BASIS,
014 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
015 * See the License for the specific language governing permissions and
016 * limitations under the License.
017 */
018
019package org.apache.oozie.util;
020
021import com.google.common.annotations.VisibleForTesting;
022import com.google.common.base.Charsets;
023
024import java.io.BufferedReader;
025import java.io.IOException;
026import java.io.InputStreamReader;
027
028public class BufferDrainer {
029
030    private final XLog LOG = XLog.getLog(getClass());
031    private static final int DRAIN_BUFFER_SLEEP_TIME_MS = 500;
032    private final Process process;
033    private final int maxLength;
034    private boolean drainBuffersFinished;
035    private final StringBuffer inputBuffer;
036    private final StringBuffer errorBuffer;
037
038    /**
039     * @param process The Process instance.
040     * @param maxLength The maximum data length to be stored in these buffers. This is an indicative value, and the
041     * store content may exceed this length.
042     */
043    public BufferDrainer(Process process, int maxLength) {
044        this.process = process;
045        this.maxLength = maxLength;
046        drainBuffersFinished = false;
047        inputBuffer = new StringBuffer();
048        errorBuffer = new StringBuffer();
049    }
050
051    /**
052     * Drains the inputStream and errorStream of the Process being executed. The contents of the streams are stored if a
053     * buffer is provided for the stream.
054     *
055     * @return the exit value of the processSettings.
056     * @throws IOException
057     */
058    public int drainBuffers() throws IOException {
059        if (drainBuffersFinished) {
060            throw new IllegalStateException("Buffer draining has already been finished");
061        }
062        LOG.trace("drainBuffers() start");
063
064        int exitValue = -1;
065
066        int inBytesRead = 0;
067        int errBytesRead = 0;
068
069        boolean processEnded = false;
070
071        try (final BufferedReader ir = new BufferedReader(new InputStreamReader(process.getInputStream(), Charsets.UTF_8));
072             final BufferedReader er = new BufferedReader(new InputStreamReader(process.getErrorStream(), Charsets.UTF_8))) {
073            // Here we do some kind of busy waiting, checking whether the process has finished by calling Process#exitValue().
074            // If not yet finished, an IllegalThreadStateException is thrown and ignored, the progress on stdout and stderr read,
075            // and retried until the process has ended.
076            // Note that Process#waitFor() may block sometimes, that's why we do a polling mechanism using Process#exitValue()
077            // instead. Until we extend unit and integration test coverage for SSH action, and we can introduce a more
078            // sophisticated error handling based on the extended coverage, this solution should stay in place.
079            while (!processEnded) {
080                try {
081                    // Doesn't block but throws IllegalThreadStateException if the process hasn't finished yet
082                    exitValue = process.exitValue();
083                    processEnded = true;
084                }
085                catch (final IllegalThreadStateException itse) {
086                    // Continue to drain
087                }
088
089                // Drain input and error streams
090                inBytesRead += drainBuffer(ir, inputBuffer, maxLength, inBytesRead, processEnded);
091                errBytesRead += drainBuffer(er, errorBuffer, maxLength, errBytesRead, processEnded);
092
093                // Necessary evil: sleep and retry
094                if (!processEnded) {
095                    try {
096                        LOG.trace("Sleeping {0}ms during buffer draining", DRAIN_BUFFER_SLEEP_TIME_MS);
097                        Thread.sleep(DRAIN_BUFFER_SLEEP_TIME_MS);
098                    }
099                    catch (final InterruptedException ie) {
100                        // Sleep a little, then check again
101                    }
102                }
103            }
104        }
105
106        LOG.trace("drainBuffers() end [exitValue={0}]", exitValue);
107        drainBuffersFinished = true;
108        return exitValue;
109    }
110
111    public StringBuffer getInputBuffer() {
112        if (drainBuffersFinished) {
113            return inputBuffer;
114        }
115        else {
116            throw new IllegalStateException("Buffer draining has not been finished yet");
117        }
118    }
119
120    public StringBuffer getErrorBuffer() {
121        if (drainBuffersFinished) {
122            return errorBuffer;
123        }
124        else {
125            throw new IllegalStateException("Buffer draining has not been finished yet");
126        }
127    }
128
129    /**
130     * Reads the contents of a stream and stores them into the provided buffer.
131     *
132     * @param br The stream to be read.
133     * @param storageBuf The buffer into which the contents of the stream are to be stored.
134     * @param maxLength The maximum number of bytes to be stored in the buffer. An indicative value and may be
135     * exceeded.
136     * @param bytesRead The number of bytes read from this stream to date.
137     * @param readAll If true, the stream is drained while their is data available in it. Otherwise, only a single chunk
138     * of data is read, irrespective of how much is available.
139     * @return bReadSession returns drainBuffer for stream of contents
140     * @throws IOException
141     */
142    @VisibleForTesting
143    static int drainBuffer(BufferedReader br, StringBuffer storageBuf, int maxLength, int bytesRead, boolean readAll)
144            throws IOException {
145        int bReadSession = 0;
146        int chunkSize = 1024;
147        if (br.ready()) {
148            char[] buf = new char[1024];
149            int bReadCurrent;
150            boolean wantsToReadFurther;
151            do {
152                bReadCurrent = br.read(buf, 0, chunkSize);
153                if (storageBuf != null && bytesRead < maxLength && bReadCurrent != -1) {
154                    storageBuf.append(buf, 0, bReadCurrent);
155                }
156                if (bReadCurrent != -1) {
157                    bReadSession += bReadCurrent;
158                }
159                wantsToReadFurther = bReadCurrent != -1 && (readAll || bReadCurrent == chunkSize);
160            } while (br.ready() && wantsToReadFurther);
161        }
162        return bReadSession;
163    }
164}