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.action.ssh;
020
021import java.io.BufferedReader;
022import java.io.File;
023import java.io.FileWriter;
024import java.io.IOException;
025import java.io.InputStreamReader;
026import java.util.Arrays;
027import java.util.List;
028import java.util.concurrent.Callable;
029
030import com.google.common.base.Charsets;
031import org.apache.hadoop.util.StringUtils;
032
033import org.apache.oozie.client.WorkflowAction;
034import org.apache.oozie.client.OozieClient;
035import org.apache.oozie.client.WorkflowAction.Status;
036import org.apache.oozie.action.ActionExecutor;
037import org.apache.oozie.action.ActionExecutorException;
038import org.apache.oozie.service.CallbackService;
039import org.apache.oozie.service.ConfigurationService;
040import org.apache.oozie.servlet.CallbackServlet;
041import org.apache.oozie.service.Services;
042import org.apache.oozie.util.BufferDrainer;
043import org.apache.oozie.util.IOUtils;
044import org.apache.oozie.util.PropertiesUtils;
045import org.apache.oozie.util.XLog;
046import org.apache.oozie.util.XmlUtils;
047import org.jdom.Element;
048import org.jdom.JDOMException;
049import org.jdom.Namespace;
050
051/**
052 * Ssh action executor. <p/> <ul> <li>Execute the shell commands on the remote host</li> <li>Copies the base and wrapper
053 * scripts on to the remote location</li> <li>Base script is used to run the command on the remote host</li> <li>Wrapper
054 * script is used to check the status of the submitted command</li> <li>handles the submission failures</li> </ul>
055 */
056public class SshActionExecutor extends ActionExecutor {
057    public static final String ACTION_TYPE = "ssh";
058
059    /**
060     * Configuration parameter which specifies whether the specified ssh user is allowed, or has to be the job user.
061     */
062    public static final String CONF_SSH_ALLOW_USER_AT_HOST = CONF_PREFIX + "ssh.allow.user.at.host";
063
064    protected static final String SSH_COMMAND_OPTIONS =
065            "-o PasswordAuthentication=no -o KbdInteractiveDevices=no -o StrictHostKeyChecking=no -o ConnectTimeout=20 ";
066
067    protected static final String SSH_COMMAND_BASE = "ssh " + SSH_COMMAND_OPTIONS;
068    protected static final String SCP_COMMAND_BASE = "scp " + SSH_COMMAND_OPTIONS;
069
070    public static final String ERR_SETUP_FAILED = "SETUP_FAILED";
071    public static final String ERR_EXECUTION_FAILED = "EXECUTION_FAILED";
072    public static final String ERR_UNKNOWN_ERROR = "UNKOWN_ERROR";
073    public static final String ERR_COULD_NOT_CONNECT = "COULD_NOT_CONNECT";
074    public static final String ERR_HOST_RESOLUTION = "COULD_NOT_RESOLVE_HOST";
075    public static final String ERR_FNF = "FNF";
076    public static final String ERR_AUTH_FAILED = "AUTH_FAILED";
077    public static final String ERR_NO_EXEC_PERM = "NO_EXEC_PERM";
078    public static final String ERR_USER_MISMATCH = "ERR_USER_MISMATCH";
079    public static final String ERR_EXCEDE_LEN = "ERR_OUTPUT_EXCEED_MAX_LEN";
080
081    public static final String DELETE_TMP_DIR = "oozie.action.ssh.delete.remote.tmp.dir";
082
083    public static final String HTTP_COMMAND = "oozie.action.ssh.http.command";
084
085    public static final String HTTP_COMMAND_OPTIONS = "oozie.action.ssh.http.command.post.options";
086
087    private static final String EXT_STATUS_VAR = "#status";
088
089    private static int maxLen;
090    private static boolean allowSshUserAtHost;
091
092    private final XLog LOG = XLog.getLog(getClass());
093
094    protected SshActionExecutor() {
095        super(ACTION_TYPE);
096    }
097
098    /**
099     * Initialize Action.
100     */
101    @Override
102    public void initActionType() {
103        super.initActionType();
104        maxLen = getOozieConf().getInt(CallbackServlet.CONF_MAX_DATA_LEN, 2 * 1024);
105        allowSshUserAtHost = ConfigurationService.getBoolean(CONF_SSH_ALLOW_USER_AT_HOST);
106        registerError(InterruptedException.class.getName(), ActionExecutorException.ErrorType.ERROR, "SH001");
107        registerError(JDOMException.class.getName(), ActionExecutorException.ErrorType.ERROR, "SH002");
108        initSshScripts();
109    }
110
111    /**
112     * Check ssh action status.
113     *
114     * @param context action execution context.
115     * @param action action object.
116     */
117    @Override
118    public void check(Context context, WorkflowAction action) throws ActionExecutorException {
119        LOG.trace("check() start for action={0}", action.getId());
120        Status status = getActionStatus(context, action);
121        boolean captureOutput = false;
122        try {
123            Element eConf = XmlUtils.parseXml(action.getConf());
124            Namespace ns = eConf.getNamespace();
125            captureOutput = eConf.getChild("capture-output", ns) != null;
126        }
127        catch (JDOMException ex) {
128            throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "ERR_XML_PARSE_FAILED",
129                                              "unknown error", ex);
130        }
131        XLog log = XLog.getLog(getClass());
132        log.debug("Capture Output: {0}", captureOutput);
133        if (status == Status.OK) {
134            if (captureOutput) {
135                String outFile = getRemoteFileName(context, action, "stdout", false, true);
136                String dataCommand = SSH_COMMAND_BASE + action.getTrackerUri() + " cat " + outFile;
137                log.debug("Ssh command [{0}]", dataCommand);
138                try {
139                    final Process process = Runtime.getRuntime().exec(dataCommand.split("\\s"));
140                    final BufferDrainer bufferDrainer = new BufferDrainer(process, maxLen);
141                    bufferDrainer.drainBuffers();
142                    final StringBuffer outBuffer = bufferDrainer.getInputBuffer();
143                    final StringBuffer errBuffer = bufferDrainer.getErrorBuffer();
144                    boolean overflow = false;
145                    LOG.trace("outBuffer={0}", outBuffer);
146                    LOG.trace("errBuffer={0}", errBuffer);
147                    if (outBuffer.length() > maxLen) {
148                        overflow = true;
149                    }
150                    if (overflow) {
151                        throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR,
152                                                          "ERR_OUTPUT_EXCEED_MAX_LEN", "unknown error");
153                    }
154                    context.setExecutionData(status.toString(), PropertiesUtils.stringToProperties(outBuffer.toString()));
155                    LOG.trace("Execution data set. status={0}, properties={1}", status,
156                            PropertiesUtils.stringToProperties(outBuffer.toString()));
157                }
158                catch (Exception ex) {
159                    throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "ERR_UNKNOWN_ERROR",
160                                                      "unknown error", ex);
161                }
162            }
163            else {
164                LOG.trace("Execution data set to null. status={0}", status);
165                context.setExecutionData(status.toString(), null);
166            }
167        }
168        else {
169            if (status == Status.ERROR) {
170                LOG.warn("Execution data set to null in ERROR");
171                context.setExecutionData(status.toString(), null);
172            }
173            else {
174                LOG.warn("Execution data not set");
175                context.setExternalStatus(status.toString());
176            }
177        }
178        LOG.trace("check() end for action={0}", action);
179    }
180
181    /**
182     * Kill ssh action.
183     *
184     * @param context action execution context.
185     * @param action object.
186     */
187    @Override
188    public void kill(Context context, WorkflowAction action) throws ActionExecutorException {
189        String command = "ssh " + action.getTrackerUri() + " kill  -KILL " + action.getExternalId();
190        int returnValue = getReturnValue(command);
191        if (returnValue != 0) {
192            throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FAILED_TO_KILL", XLog.format(
193                    "Unable to kill process {0} on {1}", action.getExternalId(), action.getTrackerUri()));
194        }
195        context.setEndData(WorkflowAction.Status.KILLED, "ERROR");
196    }
197
198    /**
199     * Start the ssh action execution.
200     *
201     * @param context action execution context.
202     * @param action action object.
203     */
204    @SuppressWarnings("unchecked")
205    @Override
206    public void start(final Context context, final WorkflowAction action) throws ActionExecutorException {
207        XLog log = XLog.getLog(getClass());
208        log.info("start() begins");
209        String confStr = action.getConf();
210        Element conf;
211        try {
212            conf = XmlUtils.parseXml(confStr);
213        }
214        catch (Exception ex) {
215            throw convertException(ex);
216        }
217        Namespace nameSpace = conf.getNamespace();
218        Element hostElement = conf.getChild("host", nameSpace);
219        String hostString = hostElement.getValue().trim();
220        hostString = prepareUserHost(hostString, context);
221        final String host = hostString;
222        final String dirLocation = execute(new Callable<String>() {
223            public String call() throws Exception {
224                return setupRemote(host, context, action);
225            }
226
227        });
228
229        String runningPid = execute(new Callable<String>() {
230            public String call() throws Exception {
231                return checkIfRunning(host, context, action);
232            }
233        });
234        String pid = "";
235
236        LOG.trace("runningPid={0}", runningPid);
237
238        if (runningPid == null) {
239            final Element commandElement = conf.getChild("command", nameSpace);
240            final boolean ignoreOutput = conf.getChild("capture-output", nameSpace) == null;
241
242            boolean preserve = false;
243            if (commandElement != null) {
244                String[] args = null;
245                // Will either have <args>, <arg>, or neither (but not both)
246                List<Element> argsList = conf.getChildren("args", nameSpace);
247                // Arguments in an <args> are "flattened" (spaces are delimiters)
248                if (argsList != null && argsList.size() > 0) {
249                    StringBuilder argsString = new StringBuilder("");
250                    for (Element argsElement : argsList) {
251                        argsString = argsString.append(argsElement.getValue()).append(" ");
252                    }
253                    args = new String[]{argsString.toString()};
254                }
255                else {
256                    // Arguments in an <arg> are preserved, even with spaces
257                    argsList = conf.getChildren("arg", nameSpace);
258                    if (argsList != null && argsList.size() > 0) {
259                        preserve = true;
260                        args = new String[argsList.size()];
261                        for (int i = 0; i < argsList.size(); i++) {
262                            Element argsElement = argsList.get(i);
263                            args[i] = argsElement.getValue();
264                            // Even though we're keeping the args as an array, if they contain a space we still have to either quote
265                            // them or escape their space (because the scripts will split them up otherwise)
266                            if (args[i].contains(" ") &&
267                                    !(args[i].startsWith("\"") && args[i].endsWith("\"") ||
268                                      args[i].startsWith("'") && args[i].endsWith("'"))) {
269                                args[i] = StringUtils.escapeString(args[i], '\\', ' ');
270                            }
271                        }
272                    }
273                }
274                final String[] argsF = args;
275                final String recoveryId = context.getRecoveryId();
276                final boolean preserveF = preserve;
277                pid = execute(new Callable<String>() {
278
279                    @Override
280                    public String call() throws Exception {
281                        return doExecute(host, dirLocation, commandElement.getValue(), argsF, ignoreOutput, action, recoveryId,
282                                preserveF);
283                    }
284
285                });
286            }
287            context.setStartData(pid, host, host);
288        }
289        else {
290            pid = runningPid;
291            context.setStartData(pid, host, host);
292            check(context, action);
293        }
294        log.info("start() ends");
295    }
296
297    private String checkIfRunning(String host, final Context context, final WorkflowAction action) {
298        String outFile = getRemoteFileName(context, action, "pid", false, false);
299        String getOutputCmd = SSH_COMMAND_BASE + host + " cat " + outFile;
300        try {
301            final Process process = Runtime.getRuntime().exec(getOutputCmd.split("\\s"));
302            final BufferDrainer bufferDrainer = new BufferDrainer(process, maxLen);
303            bufferDrainer.drainBuffers();
304            final StringBuffer buffer = bufferDrainer.getInputBuffer();
305            String pid = getFirstLine(buffer);
306            if (Long.valueOf(pid) > 0) {
307                return pid;
308            }
309            else {
310                return null;
311            }
312        }
313        catch (Exception e) {
314            return null;
315        }
316    }
317
318    /**
319     * Get remote host working location.
320     *
321     * @param context action execution context
322     * @param action Action
323     * @param fileExtension Extension to be added to file name
324     * @param dirOnly Get the Directory only
325     * @param useExtId Flag to use external ID in the path
326     * @return remote host file name/Directory.
327     */
328    public String getRemoteFileName(Context context, WorkflowAction action, String fileExtension, boolean dirOnly,
329                                    boolean useExtId) {
330        String path = getActionDirPath(context.getWorkflow().getId(), action, ACTION_TYPE, false) + "/";
331        if (dirOnly) {
332            return path;
333        }
334        if (useExtId) {
335            path = path + action.getExternalId() + ".";
336        }
337        path = path + context.getRecoveryId() + "." + fileExtension;
338        return path;
339    }
340
341    /**
342     * Utility method to execute command.
343     *
344     * @param command Command to execute as String.
345     * @return exit status of the execution.
346     * @throws IOException if process exits with status nonzero.
347     * @throws InterruptedException if process does not run properly.
348     */
349    public int executeCommand(String command) throws IOException, InterruptedException {
350        Runtime runtime = Runtime.getRuntime();
351        Process p = runtime.exec(command.split("\\s"));
352
353        final BufferDrainer bufferDrainer = new BufferDrainer(p, maxLen);
354        final int exitValue = bufferDrainer.drainBuffers();
355        final StringBuffer errorBuffer = bufferDrainer.getErrorBuffer();
356
357        String error = null;
358        if (exitValue != 0) {
359            error = getTruncatedString(errorBuffer);
360            throw new IOException(XLog.format("Not able to perform operation [{0}]", command) + " | " + "ErrorStream: "
361                    + error);
362        }
363        return exitValue;
364    }
365
366    /**
367     * Do ssh action execution setup on remote host.
368     *
369     * @param host host name.
370     * @param context action execution context.
371     * @param action action object.
372     * @return remote host working directory.
373     * @throws IOException thrown if failed to setup.
374     * @throws InterruptedException thrown if any interruption happens.
375     */
376    protected String setupRemote(String host, Context context, WorkflowAction action) throws IOException, InterruptedException {
377        LOG.info("Attempting to copy ssh base scripts to remote host [{0}]", host);
378        String localDirLocation = Services.get().getRuntimeDir() + "/ssh";
379        if (localDirLocation.endsWith("/")) {
380            localDirLocation = localDirLocation.substring(0, localDirLocation.length() - 1);
381        }
382        File file = new File(localDirLocation + "/ssh-base.sh");
383        if (!file.exists()) {
384            throw new IOException("Required Local file " + file.getAbsolutePath() + " not present.");
385        }
386        file = new File(localDirLocation + "/ssh-wrapper.sh");
387        if (!file.exists()) {
388            throw new IOException("Required Local file " + file.getAbsolutePath() + " not present.");
389        }
390        String remoteDirLocation = getRemoteFileName(context, action, null, true, true);
391        String command = XLog.format("{0}{1}  mkdir -p {2} ", SSH_COMMAND_BASE, host, remoteDirLocation).toString();
392        executeCommand(command);
393        command = XLog.format("{0}{1}/ssh-base.sh {2}/ssh-wrapper.sh {3}:{4}", SCP_COMMAND_BASE, localDirLocation,
394                              localDirLocation, host, remoteDirLocation);
395        executeCommand(command);
396        command = XLog.format("{0}{1}  chmod +x {2}ssh-base.sh {3}ssh-wrapper.sh ", SSH_COMMAND_BASE, host,
397                              remoteDirLocation, remoteDirLocation);
398        executeCommand(command);
399        return remoteDirLocation;
400    }
401
402    /**
403     * Execute the ssh command.
404     *
405     * @param host hostname.
406     * @param dirLocation location of the base and wrapper scripts.
407     * @param cmnd command to be executed.
408     * @param args command arguments.
409     * @param ignoreOutput ignore output option.
410     * @param action action object.
411     * @param recoveryId action id + run number to enable recovery in rerun
412     * @param preserveArgs tell the ssh scripts to preserve or flatten the arguments
413     * @return process id of the running command.
414     * @throws IOException thrown if failed to run the command.
415     * @throws InterruptedException thrown if any interruption happens.
416     */
417    protected String doExecute(String host, String dirLocation, String cmnd, String[] args, boolean ignoreOutput,
418                               WorkflowAction action, String recoveryId, boolean preserveArgs)
419                               throws IOException, InterruptedException {
420        XLog log = XLog.getLog(getClass());
421        Runtime runtime = Runtime.getRuntime();
422        String callbackPost = ignoreOutput ? "_" : ConfigurationService.get(HTTP_COMMAND_OPTIONS).replace(" ", "%%%");
423        String preserveArgsS = preserveArgs ? "PRESERVE_ARGS" : "FLATTEN_ARGS";
424        // TODO check
425        String callBackUrl = Services.get().get(CallbackService.class)
426                .createCallBackUrl(action.getId(), EXT_STATUS_VAR);
427        String command = XLog.format("{0}{1} {2}ssh-base.sh {3} {4} \"{5}\" \"{6}\" {7} {8} ", SSH_COMMAND_BASE, host, dirLocation,
428                preserveArgsS, ConfigurationService.get(HTTP_COMMAND), callBackUrl, callbackPost, recoveryId, cmnd)
429                .toString();
430        String[] commandArray = command.split("\\s");
431        String[] finalCommand;
432        if (args == null) {
433            finalCommand = commandArray;
434        }
435        else {
436            finalCommand = new String[commandArray.length + args.length];
437            System.arraycopy(commandArray, 0, finalCommand, 0, commandArray.length);
438            System.arraycopy(args, 0, finalCommand, commandArray.length, args.length);
439        }
440
441        LOG.trace("Executing SSH command [finalCommand={0}]", Arrays.toString(finalCommand));
442        final Process p = runtime.exec(finalCommand);
443
444        BufferDrainer bufferDrainer = new BufferDrainer(p, maxLen);
445        final int exitValue = bufferDrainer.drainBuffers();
446        final StringBuffer inputBuffer = bufferDrainer.getInputBuffer();
447        final StringBuffer errorBuffer = bufferDrainer.getErrorBuffer();
448        final String pid = getFirstLine(inputBuffer);
449        if (exitValue != 0) {
450            String error = getTruncatedString(errorBuffer);
451            throw new IOException(XLog.format("Not able to execute ssh-base.sh on {0}", host) + " | " + "ErrorStream: "
452                    + error);
453        }
454
455        LOG.trace("After execution pid={0}", pid);
456
457        return pid;
458    }
459
460    /**
461     * End action execution.
462     *
463     * @param context action execution context.
464     * @param action action object.
465     * @throws ActionExecutorException thrown if action end execution fails.
466     */
467    public void end(final Context context, final WorkflowAction action) throws ActionExecutorException {
468        if (action.getExternalStatus().equals("OK")) {
469            context.setEndData(WorkflowAction.Status.OK, WorkflowAction.Status.OK.toString());
470        }
471        else {
472            context.setEndData(WorkflowAction.Status.ERROR, WorkflowAction.Status.ERROR.toString());
473        }
474        boolean deleteTmpDir = ConfigurationService.getBoolean(DELETE_TMP_DIR);
475        if (deleteTmpDir) {
476            String tmpDir = getRemoteFileName(context, action, null, true, false);
477            String removeTmpDirCmd = SSH_COMMAND_BASE + action.getTrackerUri() + " rm -rf " + tmpDir;
478            int retVal = getReturnValue(removeTmpDirCmd);
479            if (retVal != 0) {
480                XLog.getLog(getClass()).warn("Cannot delete temp dir {0}", tmpDir);
481            }
482        }
483    }
484
485    /**
486     * Get the return value of a process.
487     *
488     * @param command command to be executed.
489     * @return zero if execution is successful and any non zero value for failure.
490     * @throws ActionExecutorException
491     */
492    private int getReturnValue(String command) throws ActionExecutorException {
493        LOG.trace("Getting return value for command={0}", command);
494
495        int returnValue;
496        Process ps = null;
497        try {
498            ps = Runtime.getRuntime().exec(command.split("\\s"));
499            final BufferDrainer bufferDrainer = new BufferDrainer(ps, 0);
500            returnValue = bufferDrainer.drainBuffers();
501        }
502        catch (IOException e) {
503            throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FAILED_OPERATION", XLog.format(
504                    "Not able to perform operation {0}", command), e);
505        }
506        finally {
507            ps.destroy();
508        }
509
510        LOG.trace("returnValue={0}", returnValue);
511
512        return returnValue;
513    }
514
515    /**
516     * Copy the ssh base and wrapper scripts to the local directory.
517     */
518    private void initSshScripts() {
519        String dirLocation = Services.get().getRuntimeDir() + "/ssh";
520        File path = new File(dirLocation);
521        path.mkdirs();
522        if (!path.exists()) {
523            throw new RuntimeException(XLog.format("Not able to create required directory {0}", dirLocation));
524        }
525        try {
526            IOUtils.copyCharStream(IOUtils.getResourceAsReader("ssh-base.sh", -1), new FileWriter(dirLocation
527                    + "/ssh-base.sh"));
528            IOUtils.copyCharStream(IOUtils.getResourceAsReader("ssh-wrapper.sh", -1), new FileWriter(dirLocation
529                    + "/ssh-wrapper.sh"));
530        }
531        catch (IOException ie) {
532            throw new RuntimeException(XLog.format("Not able to copy required scripts file to {0} "
533                    + "for SshActionHandler", dirLocation));
534        }
535    }
536
537    /**
538     * Get action status.
539     *
540     * @param action action object.
541     * @return status of the action(RUNNING/OK/ERROR).
542     * @throws ActionExecutorException thrown if there is any error in getting status.
543     */
544    protected Status getActionStatus(Context context, WorkflowAction action) throws ActionExecutorException {
545        String command = SSH_COMMAND_BASE + action.getTrackerUri() + " ps -p " + action.getExternalId();
546        Status aStatus;
547        int returnValue = getReturnValue(command);
548        if (returnValue == 0) {
549            aStatus = Status.RUNNING;
550        }
551        else {
552            String outFile = getRemoteFileName(context, action, "error", false, true);
553            String checkErrorCmd = SSH_COMMAND_BASE + action.getTrackerUri() + " ls " + outFile;
554            int retVal = getReturnValue(checkErrorCmd);
555            if (retVal == 0) {
556                aStatus = Status.ERROR;
557            }
558            else {
559                aStatus = Status.OK;
560            }
561        }
562        return aStatus;
563    }
564
565    /**
566     * Execute the callable.
567     *
568     * @param callable required callable.
569     * @throws ActionExecutorException thrown if there is any error in command execution.
570     */
571    private <T> T execute(Callable<T> callable) throws ActionExecutorException {
572        XLog log = XLog.getLog(getClass());
573        try {
574            return callable.call();
575        }
576        catch (IOException ex) {
577            log.warn("Error while executing ssh EXECUTION");
578            String errorMessage = ex.getMessage();
579            if (null == errorMessage) { // Unknown IOException
580                throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, ERR_UNKNOWN_ERROR, ex
581                        .getMessage(), ex);
582            } // Host Resolution Issues
583            else {
584                if (errorMessage.contains("Could not resolve hostname") ||
585                        errorMessage.contains("service not known")) {
586                    throw new ActionExecutorException(ActionExecutorException.ErrorType.TRANSIENT, ERR_HOST_RESOLUTION, ex
587                            .getMessage(), ex);
588                } // Connection Timeout. Host temporarily down.
589                else {
590                    if (errorMessage.contains("timed out")) {
591                        throw new ActionExecutorException(ActionExecutorException.ErrorType.TRANSIENT, ERR_COULD_NOT_CONNECT,
592                                                          ex.getMessage(), ex);
593                    }// Local ssh-base or ssh-wrapper missing
594                    else {
595                        if (errorMessage.contains("Required Local file")) {
596                            throw new ActionExecutorException(ActionExecutorException.ErrorType.TRANSIENT, ERR_FNF,
597                                                              ex.getMessage(), ex); // local_FNF
598                        }// Required oozie bash scripts missing, after the copy was
599                        // successful
600                        else {
601                            if (errorMessage.contains("No such file or directory")
602                                    && (errorMessage.contains("ssh-base") || errorMessage.contains("ssh-wrapper"))) {
603                                throw new ActionExecutorException(ActionExecutorException.ErrorType.TRANSIENT, ERR_FNF,
604                                                                  ex.getMessage(), ex); // remote
605                                // FNF
606                            } // Required application execution binary missing (either
607                            // caught by ssh-wrapper
608                            else {
609                                if (errorMessage.contains("command not found")) {
610                                    throw new ActionExecutorException(ActionExecutorException.ErrorType.NON_TRANSIENT, ERR_FNF, ex
611                                            .getMessage(), ex); // remote
612                                    // FNF
613                                } // Permission denied while connecting
614                                else {
615                                    if (errorMessage.contains("Permission denied")) {
616                                        throw new ActionExecutorException(ActionExecutorException.ErrorType.NON_TRANSIENT,
617                                                ERR_AUTH_FAILED, ex.getMessage(), ex);
618                                    } // Permission denied while executing
619                                    else {
620                                        if (errorMessage.contains(": Permission denied")) {
621                                            throw new ActionExecutorException(ActionExecutorException.ErrorType.NON_TRANSIENT,
622                                                    ERR_NO_EXEC_PERM, ex.getMessage(), ex);
623                                        }
624                                        else {
625                                            throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR,
626                                                    ERR_UNKNOWN_ERROR, ex.getMessage(), ex);
627                                        }
628                                    }
629                                }
630                            }
631                        }
632                    }
633                }
634            }
635        } // Any other type of exception
636        catch (Exception ex) {
637            throw convertException(ex);
638        }
639    }
640
641    /**
642     * Checks whether the system is configured to always use the oozie user for ssh, and injects the user if required.
643     *
644     * @param host the host string.
645     * @param context the execution context.
646     * @return the modified host string with a user parameter added on if required.
647     * @throws ActionExecutorException in case the flag to use the oozie user is turned on and there is a mismatch
648     * between the user specified in the host and the oozie user.
649     */
650    private String prepareUserHost(String host, Context context) throws ActionExecutorException {
651        String oozieUser = context.getProtoActionConf().get(OozieClient.USER_NAME);
652        if (allowSshUserAtHost) {
653            if (!host.contains("@")) {
654                host = oozieUser + "@" + host;
655            }
656        }
657        else {
658            if (host.contains("@")) {
659                if (!host.toLowerCase().startsWith(oozieUser + "@")) {
660                    throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, ERR_USER_MISMATCH,
661                                                      XLog.format("user mismatch between oozie user [{0}] and ssh host [{1}]",
662                                                              oozieUser, host));
663                }
664            }
665            else {
666                host = oozieUser + "@" + host;
667            }
668        }
669
670        LOG.trace("User host is {0}", host);
671
672        return host;
673    }
674
675    @Override
676    public boolean isCompleted(String externalStatus) {
677        return true;
678    }
679
680    /**
681     * Truncate the string to max length.
682     *
683     * @param strBuffer
684     * @return truncated string string
685     */
686    private String getTruncatedString(StringBuffer strBuffer) {
687
688        if (strBuffer.length() <= maxLen) {
689            return strBuffer.toString();
690        }
691        else {
692            return strBuffer.substring(0, maxLen);
693        }
694    }
695
696    /**
697     * Returns the first line from a StringBuffer, recognized by the new line character \n.
698     *
699     * @param buffer The StringBuffer from which the first line is required.
700     * @return The first line of the buffer.
701     */
702    private String getFirstLine(StringBuffer buffer) {
703        int newLineIndex = 0;
704        newLineIndex = buffer.indexOf("\n");
705        if (newLineIndex == -1) {
706            return buffer.toString();
707        }
708        else {
709            return buffer.substring(0, newLineIndex);
710        }
711    }
712}