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}