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}