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; 029import com.google.common.base.Charsets; 030import org.apache.hadoop.util.StringUtils; 031 032import org.apache.oozie.client.WorkflowAction; 033import org.apache.oozie.client.OozieClient; 034import org.apache.oozie.client.WorkflowAction.Status; 035import org.apache.oozie.action.ActionExecutor; 036import org.apache.oozie.action.ActionExecutorException; 037import org.apache.oozie.service.CallbackService; 038import org.apache.oozie.service.ConfigurationService; 039import org.apache.oozie.servlet.CallbackServlet; 040import org.apache.oozie.service.Services; 041import org.apache.oozie.util.IOUtils; 042import org.apache.oozie.util.PropertiesUtils; 043import org.apache.oozie.util.XLog; 044import org.apache.oozie.util.XmlUtils; 045import org.jdom.Element; 046import org.jdom.JDOMException; 047import org.jdom.Namespace; 048 049/** 050 * Ssh action executor. <p/> <ul> <li>Execute the shell commands on the remote host</li> <li>Copies the base and wrapper 051 * scripts on to the remote location</li> <li>Base script is used to run the command on the remote host</li> <li>Wrapper 052 * script is used to check the status of the submitted command</li> <li>handles the submission failures</li> </ul> 053 */ 054public class SshActionExecutor extends ActionExecutor { 055 public static final String ACTION_TYPE = "ssh"; 056 057 /** 058 * Configuration parameter which specifies whether the specified ssh user is allowed, or has to be the job user. 059 */ 060 public static final String CONF_SSH_ALLOW_USER_AT_HOST = CONF_PREFIX + "ssh.allow.user.at.host"; 061 062 protected static final String SSH_COMMAND_OPTIONS = 063 "-o PasswordAuthentication=no -o KbdInteractiveDevices=no -o StrictHostKeyChecking=no -o ConnectTimeout=20 "; 064 065 protected static final String SSH_COMMAND_BASE = "ssh " + SSH_COMMAND_OPTIONS; 066 protected static final String SCP_COMMAND_BASE = "scp " + SSH_COMMAND_OPTIONS; 067 068 public static final String ERR_SETUP_FAILED = "SETUP_FAILED"; 069 public static final String ERR_EXECUTION_FAILED = "EXECUTION_FAILED"; 070 public static final String ERR_UNKNOWN_ERROR = "UNKOWN_ERROR"; 071 public static final String ERR_COULD_NOT_CONNECT = "COULD_NOT_CONNECT"; 072 public static final String ERR_HOST_RESOLUTION = "COULD_NOT_RESOLVE_HOST"; 073 public static final String ERR_FNF = "FNF"; 074 public static final String ERR_AUTH_FAILED = "AUTH_FAILED"; 075 public static final String ERR_NO_EXEC_PERM = "NO_EXEC_PERM"; 076 public static final String ERR_USER_MISMATCH = "ERR_USER_MISMATCH"; 077 public static final String ERR_EXCEDE_LEN = "ERR_OUTPUT_EXCEED_MAX_LEN"; 078 079 public static final String DELETE_TMP_DIR = "oozie.action.ssh.delete.remote.tmp.dir"; 080 081 public static final String HTTP_COMMAND = "oozie.action.ssh.http.command"; 082 083 public static final String HTTP_COMMAND_OPTIONS = "oozie.action.ssh.http.command.post.options"; 084 085 private static final String EXT_STATUS_VAR = "#status"; 086 087 private static int maxLen; 088 private static boolean allowSshUserAtHost; 089 090 private final XLog LOG = XLog.getLog(getClass()); 091 092 protected SshActionExecutor() { 093 super(ACTION_TYPE); 094 } 095 096 /** 097 * Initialize Action. 098 */ 099 @Override 100 public void initActionType() { 101 super.initActionType(); 102 maxLen = getOozieConf().getInt(CallbackServlet.CONF_MAX_DATA_LEN, 2 * 1024); 103 allowSshUserAtHost = ConfigurationService.getBoolean(CONF_SSH_ALLOW_USER_AT_HOST); 104 registerError(InterruptedException.class.getName(), ActionExecutorException.ErrorType.ERROR, "SH001"); 105 registerError(JDOMException.class.getName(), ActionExecutorException.ErrorType.ERROR, "SH002"); 106 initSshScripts(); 107 } 108 109 /** 110 * Check ssh action status. 111 * 112 * @param context action execution context. 113 * @param action action object. 114 */ 115 @Override 116 public void check(Context context, WorkflowAction action) throws ActionExecutorException { 117 LOG.trace("check() start for action={0}", action.getId()); 118 Status status = getActionStatus(context, action); 119 boolean captureOutput = false; 120 try { 121 Element eConf = XmlUtils.parseXml(action.getConf()); 122 Namespace ns = eConf.getNamespace(); 123 captureOutput = eConf.getChild("capture-output", ns) != null; 124 } 125 catch (JDOMException ex) { 126 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "ERR_XML_PARSE_FAILED", 127 "unknown error", ex); 128 } 129 XLog log = XLog.getLog(getClass()); 130 log.debug("Capture Output: {0}", captureOutput); 131 if (status == Status.OK) { 132 if (captureOutput) { 133 String outFile = getRemoteFileName(context, action, "stdout", false, true); 134 String dataCommand = SSH_COMMAND_BASE + action.getTrackerUri() + " cat " + outFile; 135 log.debug("Ssh command [{0}]", dataCommand); 136 try { 137 final Process process = Runtime.getRuntime().exec(dataCommand.split("\\s")); 138 139 final StringBuffer outBuffer = new StringBuffer(); 140 final StringBuffer errBuffer = new StringBuffer(); 141 boolean overflow = false; 142 drainBuffers(process, outBuffer, errBuffer, maxLen); 143 LOG.trace("outBuffer={0}", outBuffer); 144 LOG.trace("errBuffer={0}", errBuffer); 145 if (outBuffer.length() > maxLen) { 146 overflow = true; 147 } 148 if (overflow) { 149 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, 150 "ERR_OUTPUT_EXCEED_MAX_LEN", "unknown error"); 151 } 152 context.setExecutionData(status.toString(), PropertiesUtils.stringToProperties(outBuffer.toString())); 153 LOG.trace("Execution data set. status={0}, properties={1}", status, 154 PropertiesUtils.stringToProperties(outBuffer.toString())); 155 } 156 catch (Exception ex) { 157 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "ERR_UNKNOWN_ERROR", 158 "unknown error", ex); 159 } 160 } 161 else { 162 LOG.trace("Execution data set to null. status={0}", status); 163 context.setExecutionData(status.toString(), null); 164 } 165 } 166 else { 167 if (status == Status.ERROR) { 168 LOG.warn("Execution data set to null in ERROR"); 169 context.setExecutionData(status.toString(), null); 170 } 171 else { 172 LOG.warn("Execution data not set"); 173 context.setExternalStatus(status.toString()); 174 } 175 } 176 LOG.trace("check() end for action={0}", action); 177 } 178 179 /** 180 * Kill ssh action. 181 * 182 * @param context action execution context. 183 * @param action object. 184 */ 185 @Override 186 public void kill(Context context, WorkflowAction action) throws ActionExecutorException { 187 String command = "ssh " + action.getTrackerUri() + " kill -KILL " + action.getExternalId(); 188 int returnValue = getReturnValue(command); 189 if (returnValue != 0) { 190 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FAILED_TO_KILL", XLog.format( 191 "Unable to kill process {0} on {1}", action.getExternalId(), action.getTrackerUri())); 192 } 193 context.setEndData(WorkflowAction.Status.KILLED, "ERROR"); 194 } 195 196 /** 197 * Start the ssh action execution. 198 * 199 * @param context action execution context. 200 * @param action action object. 201 */ 202 @SuppressWarnings("unchecked") 203 @Override 204 public void start(final Context context, final WorkflowAction action) throws ActionExecutorException { 205 XLog log = XLog.getLog(getClass()); 206 log.info("start() begins"); 207 String confStr = action.getConf(); 208 Element conf; 209 try { 210 conf = XmlUtils.parseXml(confStr); 211 } 212 catch (Exception ex) { 213 throw convertException(ex); 214 } 215 Namespace nameSpace = conf.getNamespace(); 216 Element hostElement = conf.getChild("host", nameSpace); 217 String hostString = hostElement.getValue().trim(); 218 hostString = prepareUserHost(hostString, context); 219 final String host = hostString; 220 final String dirLocation = execute(new Callable<String>() { 221 public String call() throws Exception { 222 return setupRemote(host, context, action); 223 } 224 225 }); 226 227 String runningPid = execute(new Callable<String>() { 228 public String call() throws Exception { 229 return checkIfRunning(host, context, action); 230 } 231 }); 232 String pid = ""; 233 234 LOG.trace("runningPid={0}", runningPid); 235 236 if (runningPid == null) { 237 final Element commandElement = conf.getChild("command", nameSpace); 238 final boolean ignoreOutput = conf.getChild("capture-output", nameSpace) == null; 239 240 boolean preserve = false; 241 if (commandElement != null) { 242 String[] args = null; 243 // Will either have <args>, <arg>, or neither (but not both) 244 List<Element> argsList = conf.getChildren("args", nameSpace); 245 // Arguments in an <args> are "flattened" (spaces are delimiters) 246 if (argsList != null && argsList.size() > 0) { 247 StringBuilder argsString = new StringBuilder(""); 248 for (Element argsElement : argsList) { 249 argsString = argsString.append(argsElement.getValue()).append(" "); 250 } 251 args = new String[]{argsString.toString()}; 252 } 253 else { 254 // Arguments in an <arg> are preserved, even with spaces 255 argsList = conf.getChildren("arg", nameSpace); 256 if (argsList != null && argsList.size() > 0) { 257 preserve = true; 258 args = new String[argsList.size()]; 259 for (int i = 0; i < argsList.size(); i++) { 260 Element argsElement = argsList.get(i); 261 args[i] = argsElement.getValue(); 262 // Even though we're keeping the args as an array, if they contain a space we still have to either quote 263 // them or escape their space (because the scripts will split them up otherwise) 264 if (args[i].contains(" ") && 265 !(args[i].startsWith("\"") && args[i].endsWith("\"") || 266 args[i].startsWith("'") && args[i].endsWith("'"))) { 267 args[i] = StringUtils.escapeString(args[i], '\\', ' '); 268 } 269 } 270 } 271 } 272 final String[] argsF = args; 273 final String recoveryId = context.getRecoveryId(); 274 final boolean preserveF = preserve; 275 pid = execute(new Callable<String>() { 276 277 @Override 278 public String call() throws Exception { 279 return doExecute(host, dirLocation, commandElement.getValue(), argsF, ignoreOutput, action, recoveryId, 280 preserveF); 281 } 282 283 }); 284 } 285 context.setStartData(pid, host, host); 286 } 287 else { 288 pid = runningPid; 289 context.setStartData(pid, host, host); 290 check(context, action); 291 } 292 log.info("start() ends"); 293 } 294 295 private String checkIfRunning(String host, final Context context, final WorkflowAction action) { 296 String pid = null; 297 String outFile = getRemoteFileName(context, action, "pid", false, false); 298 String getOutputCmd = SSH_COMMAND_BASE + host + " cat " + outFile; 299 try { 300 Process process = Runtime.getRuntime().exec(getOutputCmd.split("\\s")); 301 StringBuffer buffer = new StringBuffer(); 302 drainBuffers(process, buffer, null, maxLen); 303 pid = getFirstLine(buffer); 304 305 if (Long.valueOf(pid) > 0) { 306 return pid; 307 } 308 else { 309 return null; 310 } 311 } 312 catch (Exception e) { 313 return null; 314 } 315 } 316 317 /** 318 * Get remote host working location. 319 * 320 * @param context action execution context 321 * @param action Action 322 * @param fileExtension Extension to be added to file name 323 * @param dirOnly Get the Directory only 324 * @param useExtId Flag to use external ID in the path 325 * @return remote host file name/Directory. 326 */ 327 public String getRemoteFileName(Context context, WorkflowAction action, String fileExtension, boolean dirOnly, 328 boolean useExtId) { 329 String path = getActionDirPath(context.getWorkflow().getId(), action, ACTION_TYPE, false) + "/"; 330 if (dirOnly) { 331 return path; 332 } 333 if (useExtId) { 334 path = path + action.getExternalId() + "."; 335 } 336 path = path + context.getRecoveryId() + "." + fileExtension; 337 return path; 338 } 339 340 /** 341 * Utility method to execute command. 342 * 343 * @param command Command to execute as String. 344 * @return exit status of the execution. 345 * @throws IOException if process exits with status nonzero. 346 * @throws InterruptedException if process does not run properly. 347 */ 348 public int executeCommand(String command) throws IOException, InterruptedException { 349 Runtime runtime = Runtime.getRuntime(); 350 Process p = runtime.exec(command.split("\\s")); 351 352 StringBuffer errorBuffer = new StringBuffer(); 353 int exitValue = drainBuffers(p, null, errorBuffer, maxLen); 354 355 String error = null; 356 if (exitValue != 0) { 357 error = getTruncatedString(errorBuffer); 358 throw new IOException(XLog.format("Not able to perform operation [{0}]", command) + " | " + "ErrorStream: " 359 + error); 360 } 361 return exitValue; 362 } 363 364 /** 365 * Do ssh action execution setup on remote host. 366 * 367 * @param host host name. 368 * @param context action execution context. 369 * @param action action object. 370 * @return remote host working directory. 371 * @throws IOException thrown if failed to setup. 372 * @throws InterruptedException thrown if any interruption happens. 373 */ 374 protected String setupRemote(String host, Context context, WorkflowAction action) throws IOException, InterruptedException { 375 LOG.info("Attempting to copy ssh base scripts to remote host [{0}]", host); 376 String localDirLocation = Services.get().getRuntimeDir() + "/ssh"; 377 if (localDirLocation.endsWith("/")) { 378 localDirLocation = localDirLocation.substring(0, localDirLocation.length() - 1); 379 } 380 File file = new File(localDirLocation + "/ssh-base.sh"); 381 if (!file.exists()) { 382 throw new IOException("Required Local file " + file.getAbsolutePath() + " not present."); 383 } 384 file = new File(localDirLocation + "/ssh-wrapper.sh"); 385 if (!file.exists()) { 386 throw new IOException("Required Local file " + file.getAbsolutePath() + " not present."); 387 } 388 String remoteDirLocation = getRemoteFileName(context, action, null, true, true); 389 String command = XLog.format("{0}{1} mkdir -p {2} ", SSH_COMMAND_BASE, host, remoteDirLocation).toString(); 390 executeCommand(command); 391 command = XLog.format("{0}{1}/ssh-base.sh {2}/ssh-wrapper.sh {3}:{4}", SCP_COMMAND_BASE, localDirLocation, 392 localDirLocation, host, remoteDirLocation); 393 executeCommand(command); 394 command = XLog.format("{0}{1} chmod +x {2}ssh-base.sh {3}ssh-wrapper.sh ", SSH_COMMAND_BASE, host, 395 remoteDirLocation, remoteDirLocation); 396 executeCommand(command); 397 return remoteDirLocation; 398 } 399 400 /** 401 * Execute the ssh command. 402 * 403 * @param host hostname. 404 * @param dirLocation location of the base and wrapper scripts. 405 * @param cmnd command to be executed. 406 * @param args command arguments. 407 * @param ignoreOutput ignore output option. 408 * @param action action object. 409 * @param recoveryId action id + run number to enable recovery in rerun 410 * @param preserveArgs tell the ssh scripts to preserve or flatten the arguments 411 * @return process id of the running command. 412 * @throws IOException thrown if failed to run the command. 413 * @throws InterruptedException thrown if any interruption happens. 414 */ 415 protected String doExecute(String host, String dirLocation, String cmnd, String[] args, boolean ignoreOutput, 416 WorkflowAction action, String recoveryId, boolean preserveArgs) 417 throws IOException, InterruptedException { 418 XLog log = XLog.getLog(getClass()); 419 Runtime runtime = Runtime.getRuntime(); 420 String callbackPost = ignoreOutput ? "_" : ConfigurationService.get(HTTP_COMMAND_OPTIONS).replace(" ", "%%%"); 421 String preserveArgsS = preserveArgs ? "PRESERVE_ARGS" : "FLATTEN_ARGS"; 422 // TODO check 423 String callBackUrl = Services.get().get(CallbackService.class) 424 .createCallBackUrl(action.getId(), EXT_STATUS_VAR); 425 String command = XLog.format("{0}{1} {2}ssh-base.sh {3} {4} \"{5}\" \"{6}\" {7} {8} ", SSH_COMMAND_BASE, host, dirLocation, 426 preserveArgsS, ConfigurationService.get(HTTP_COMMAND), callBackUrl, callbackPost, recoveryId, cmnd) 427 .toString(); 428 String[] commandArray = command.split("\\s"); 429 String[] finalCommand; 430 if (args == null) { 431 finalCommand = commandArray; 432 } 433 else { 434 finalCommand = new String[commandArray.length + args.length]; 435 System.arraycopy(commandArray, 0, finalCommand, 0, commandArray.length); 436 System.arraycopy(args, 0, finalCommand, commandArray.length, args.length); 437 } 438 439 LOG.trace("Executing SSH command [finalCommand={0}]", Arrays.toString(finalCommand)); 440 final Process p = runtime.exec(finalCommand); 441 final String pid; 442 443 final StringBuffer inputBuffer = new StringBuffer(); 444 final StringBuffer errorBuffer = new StringBuffer(); 445 final int exitValue = drainBuffers(p, inputBuffer, errorBuffer, maxLen); 446 447 pid = getFirstLine(inputBuffer); 448 449 String error = null; 450 if (exitValue != 0) { 451 error = getTruncatedString(errorBuffer); 452 throw new IOException(XLog.format("Not able to execute ssh-base.sh on {0}", host) + " | " + "ErrorStream: " 453 + error); 454 } 455 456 LOG.trace("After execution pid={0}", pid); 457 458 return pid; 459 } 460 461 /** 462 * End action execution. 463 * 464 * @param context action execution context. 465 * @param action action object. 466 * @throws ActionExecutorException thrown if action end execution fails. 467 */ 468 public void end(final Context context, final WorkflowAction action) throws ActionExecutorException { 469 if (action.getExternalStatus().equals("OK")) { 470 context.setEndData(WorkflowAction.Status.OK, WorkflowAction.Status.OK.toString()); 471 } 472 else { 473 context.setEndData(WorkflowAction.Status.ERROR, WorkflowAction.Status.ERROR.toString()); 474 } 475 boolean deleteTmpDir = ConfigurationService.getBoolean(DELETE_TMP_DIR); 476 if (deleteTmpDir) { 477 String tmpDir = getRemoteFileName(context, action, null, true, false); 478 String removeTmpDirCmd = SSH_COMMAND_BASE + action.getTrackerUri() + " rm -rf " + tmpDir; 479 int retVal = getReturnValue(removeTmpDirCmd); 480 if (retVal != 0) { 481 XLog.getLog(getClass()).warn("Cannot delete temp dir {0}", tmpDir); 482 } 483 } 484 } 485 486 /** 487 * Get the return value of a process. 488 * 489 * @param command command to be executed. 490 * @return zero if execution is successful and any non zero value for failure. 491 * @throws ActionExecutorException 492 */ 493 private int getReturnValue(String command) throws ActionExecutorException { 494 LOG.trace("Getting return value for command={0}", command); 495 496 int returnValue; 497 Process ps = null; 498 try { 499 ps = Runtime.getRuntime().exec(command.split("\\s")); 500 returnValue = drainBuffers(ps, null, null, 0); 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 * Drains the inputStream and errorStream of the Process being executed. The contents of the streams are stored if a 698 * buffer is provided for the stream. 699 * 700 * @param p The Process instance. 701 * @param inputBuffer The buffer into which STDOUT is to be read. Can be null if only draining is required. 702 * @param errorBuffer The buffer into which STDERR is to be read. Can be null if only draining is required. 703 * @param maxLength The maximum data length to be stored in these buffers. This is an indicative value, and the 704 * store content may exceed this length. 705 * @return the exit value of the process. 706 * @throws IOException 707 */ 708 private int drainBuffers(final Process p, final StringBuffer inputBuffer, final StringBuffer errorBuffer, final int maxLength) 709 throws IOException { 710 LOG.trace("drainBuffers() start"); 711 712 int exitValue = -1; 713 int inBytesRead = 0; 714 int errBytesRead = 0; 715 716 boolean processEnded = false; 717 718 try (final BufferedReader ir = new BufferedReader(new InputStreamReader(p.getInputStream(), Charsets.UTF_8)); 719 final BufferedReader er = new BufferedReader(new InputStreamReader(p.getErrorStream(), Charsets.UTF_8))) { 720 // Here we do some kind of busy waiting, checking whether the process has finished by calling Process#exitValue(). 721 // If not yet finished, an IllegalThreadStateException is thrown and ignored, the progress on stdout and stderr read, 722 // and retried until the process has ended. 723 // Note that Process#waitFor() may block sometimes, that's why we do a polling mechanism using Process#exitValue() 724 // instead. Until we extend unit and integration test coverage for SSH action, and we can introduce a more sophisticated 725 // error handling based on the extended coverage, this solution should stay in place. 726 while (!processEnded) { 727 try { 728 // Doesn't block but throws IllegalThreadStateException if the process hasn't finished yet 729 exitValue = p.exitValue(); 730 processEnded = true; 731 } 732 catch (final IllegalThreadStateException itse) { 733 // Continue to drain 734 } 735 736 // Drain input and error streams 737 inBytesRead += drainBuffer(ir, inputBuffer, maxLength, inBytesRead, processEnded); 738 errBytesRead += drainBuffer(er, errorBuffer, maxLength, errBytesRead, processEnded); 739 740 // Necessary evil: sleep and retry 741 if (!processEnded) { 742 try { 743 Thread.sleep(500); 744 } 745 catch (final InterruptedException ie) { 746 // Sleep a little, then check again 747 } 748 } 749 } 750 } 751 752 LOG.trace("drainBuffers() end [exitValue={0}]", exitValue); 753 754 return exitValue; 755 } 756 757 /** 758 * Reads the contents of a stream and stores them into the provided buffer. 759 * 760 * @param br The stream to be read. 761 * @param storageBuf The buffer into which the contents of the stream are to be stored. 762 * @param maxLength The maximum number of bytes to be stored in the buffer. An indicative value and may be 763 * exceeded. 764 * @param bytesRead The number of bytes read from this stream to date. 765 * @param readAll If true, the stream is drained while their is data available in it. Otherwise, only a single chunk 766 * of data is read, irrespective of how much is available. 767 * @return 768 * @throws IOException 769 */ 770 private int drainBuffer(BufferedReader br, StringBuffer storageBuf, int maxLength, int bytesRead, boolean readAll) 771 throws IOException { 772 int bReadSession = 0; 773 if (br.ready()) { 774 char[] buf = new char[1024]; 775 do { 776 int bReadCurrent = br.read(buf, 0, 1024); 777 if (storageBuf != null && bytesRead < maxLength) { 778 storageBuf.append(buf, 0, bReadCurrent); 779 } 780 bReadSession += bReadCurrent; 781 } while (br.ready() && readAll); 782 } 783 return bReadSession; 784 } 785 786 /** 787 * Returns the first line from a StringBuffer, recognized by the new line character \n. 788 * 789 * @param buffer The StringBuffer from which the first line is required. 790 * @return The first line of the buffer. 791 */ 792 private String getFirstLine(StringBuffer buffer) { 793 int newLineIndex = 0; 794 newLineIndex = buffer.indexOf("\n"); 795 if (newLineIndex == -1) { 796 return buffer.toString(); 797 } 798 else { 799 return buffer.substring(0, newLineIndex); 800 } 801 } 802}