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 java.io.BufferedReader;
022import java.io.IOException;
023import java.io.Writer;
024import java.util.ArrayList;
025import java.util.regex.Pattern;
026
027import org.apache.oozie.service.Services;
028import org.apache.oozie.service.XLogStreamingService;
029import org.apache.oozie.util.LogLine.MATCHED_PATTERN;
030
031  /**
032 * Encapsulates the parsing and filtering of the log messages from a BufferedReader. It makes sure not to read in the entire log
033 * into memory at the same time; at most, it will have two messages (which can be multi-line in the case of exception stack traces).
034 * <p>
035 * To use this class: Calling {@link TimestampedMessageParser#increment()} will tell the parser to read the next message from the
036 * Reader. It will return true if there are more messages and false if not. Calling
037 * {@link TimestampedMessageParser#getLastMessage()} and {@link TimestampedMessageParser#getLastTimestamp()} will return the last
038 * message and timestamp, respectively, that were parsed when {@link TimestampedMessageParser#increment()} was called. Calling
039 * {@link TimestampedMessageParser#processRemaining(java.io.Writer)} will write the remaining log messages to the given Writer.
040 */
041public class TimestampedMessageParser {
042
043    static final String SYSTEM_LINE_SEPARATOR = System.getProperty("line.separator");
044    protected BufferedReader reader;
045    private LogLine nextLine = null;
046    private String lastTimestamp = null;
047    private XLogFilter filter;
048    private boolean empty = false;
049    private String lastMessage = null;
050    private boolean patternMatched = false;
051    public int count = 0;
052    private Pattern splitPattern = null;
053
054    /**
055     * Creates a TimestampedMessageParser with the given BufferedReader and filter.
056     *
057     * @param reader The BufferedReader to get the log messages from
058     * @param filter The filter
059     */
060    public TimestampedMessageParser(BufferedReader reader, XLogFilter filter) {
061        this.reader = reader;
062        this.filter = filter;
063        if (filter == null) {
064            filter = new XLogFilter();
065        }
066        filter.constructPattern();
067        String regEx = XLogFilter.PREFIX_REGEX + filter.getFilterPattern().pattern();
068        this.splitPattern = Pattern.compile(regEx);
069    }
070
071    /**
072     * Causes the next message and timestamp to be parsed from the BufferedReader.
073     *
074     * @return true if there are more messages left; false if not
075     * @throws IOException If there was a problem reading from the Reader
076     */
077    public boolean increment() throws IOException {
078        if (empty) {
079            return false;
080        }
081
082        StringBuilder message = new StringBuilder();
083        if (nextLine == null) {     // first time only
084            nextLine = parseNextLogLine();
085            if (nextLine == null || nextLine.getLine() == null) {
086                // reader finished
087                empty = true;
088                return false;
089            }
090        }
091        lastTimestamp = parseTimestamp(nextLine);
092        String nextTimestamp = null;
093        while (nextTimestamp == null) {
094            message.append(nextLine.getLine()).append(SYSTEM_LINE_SEPARATOR);
095            nextLine = parseNextLogLine();
096            if (nextLine != null && nextLine.getLine() != null) {
097                // exit loop if we have a timestamp, continue if not
098                nextTimestamp = parseTimestamp(nextLine);
099            }
100            else {                                          // reader finished
101                empty = true;
102                nextTimestamp = "";                         // exit loop
103            }
104        }
105
106        lastMessage = message.toString();
107        return true;
108    }
109
110    /**
111     * Returns the timestamp from the last message that was parsed.
112     *
113     * @return the timestamp from the last message that was parsed
114     */
115    public String getLastTimestamp() {
116        return lastTimestamp;
117    }
118
119    /**
120     * Returns the message that was last parsed.
121     *
122     * @return the message that was last parsed
123     */
124    public String getLastMessage() {
125        return lastMessage;
126    }
127
128    /**
129     * Closes the Reader.
130     *
131     * @throws IOException
132     */
133    public void closeReader() throws IOException {
134        reader.close();
135    }
136
137    /**
138     * Reads the next line from the Reader and checks if it matches the filter.
139     * It can also handle multi-line messages (i.e. exception stack traces). If
140     * it returns null, then there are no lines left in the Reader.
141     *
142     * @return LogLine
143     * @throws IOException
144     */
145    protected LogLine parseNextLogLine() throws IOException {
146        String line;
147        LogLine logLine = new LogLine();
148        while ((line = reader.readLine()) != null) {
149            logLine.setLine(line);
150            logLine.setLogParts(null);
151            filter.splitLogMessage(logLine, splitPattern);
152            // check the splits if logLine matches with the splitPattern
153            // Otherwise, go with previous patternMatched value. This is needed
154            // in parsing stack trace
155            patternMatched = logLine.getMatchedPattern() == MATCHED_PATTERN.NONE ? patternMatched
156                    : filter.splitsMatches(logLine);
157            if (patternMatched) {
158                if (filter.getLogLimit() != -1) {
159                    if (logLine.getLogParts() != null) {
160                        if (count >= filter.getLogLimit()) {
161                            return null;
162                        }
163                        count++;
164                    }
165                }
166                if (logLine.getLogParts() != null) {
167                    if (filter.getEndDate() != null) {
168                        // Ignore the milli second part
169                        if (logLine.getLogParts().get(0).substring(0, 19).compareTo(filter.getFormattedEndDate()) > 0)
170                            return null;
171                    }
172                }
173                return logLine;
174            }
175        }
176        logLine.setLine(null);
177        return logLine;
178    }
179
180    /**
181     * Parses the timestamp out of the passed in line. If there isn't one, it
182     * returns null.
183     *
184     * @param logLine The LogLine to check
185     * @return the timestamp of the line, or null
186     */
187    private String parseTimestamp(LogLine logLine) {
188        String timestamp = null;
189        if (logLine != null && logLine.getLogParts() != null && logLine.getLogParts().size() > 0) {
190            timestamp = logLine.getLogParts().get(0);
191        }
192        return timestamp;
193    }
194
195    /**
196     * Streams log messages to the passed in Writer. Flushes the log writing
197     * based on buffer len
198     *
199     * @param writer
200     * @param bufferLen maximum len of log buffer
201     * @param bytesWritten num bytes already written to writer
202     * @throws IOException
203     */
204    public void processRemaining(Writer writer, int bufferLen, int bytesWritten) throws IOException {
205        while (increment()) {
206            writer.write(lastMessage);
207            bytesWritten += lastMessage.length();
208            if (bytesWritten > bufferLen) {
209                writer.flush();
210                bytesWritten = 0;
211            }
212        }
213        writer.flush();
214    }
215
216    /**
217     * Streams log messages to the passed in Writer, with zero bytes already
218     * written
219     *
220     * @param writer
221     * @param bufferLen maximum len of log buffer
222     * @throws IOException
223     */
224    public void processRemaining(Writer writer, int bufferLen) throws IOException {
225        processRemaining(writer, bufferLen, 0);
226    }
227
228    /**
229     * Streams log messages to the passed in Writer, with default buffer len 4K
230     * and zero bytes already written
231     *
232     * @param writer
233     * @throws IOException
234     */
235    public void processRemaining(Writer writer) throws IOException {
236        processRemaining(writer, Services.get().get(XLogStreamingService.class).getBufferLen());
237    }
238
239    /**
240     * Splits the log message into parts
241     *
242     * @param line
243     * @return List of log parts
244     */
245    protected ArrayList<String> splitLogMessage(String line) {
246        return filter.splitLogMessage(line);
247    }
248
249
250}