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.hadoop; 020 021import java.io.IOException; 022import java.net.URISyntaxException; 023import java.util.ArrayList; 024import java.util.HashMap; 025import java.util.List; 026import java.util.Map; 027 028import org.apache.hadoop.fs.FSDataOutputStream; 029import org.apache.hadoop.fs.FileStatus; 030import org.apache.hadoop.fs.FileSystem; 031import org.apache.hadoop.fs.FileUtil; 032import org.apache.hadoop.fs.Path; 033import org.apache.hadoop.fs.permission.FsPermission; 034import org.apache.hadoop.mapred.JobConf; 035import org.apache.hadoop.security.AccessControlException; 036import org.apache.oozie.action.ActionExecutor; 037import org.apache.oozie.action.ActionExecutorException; 038import org.apache.oozie.client.WorkflowAction; 039import org.apache.oozie.command.wf.WorkflowXCommand; 040import org.apache.oozie.dependency.FSURIHandler; 041import org.apache.oozie.dependency.URIHandler; 042import org.apache.oozie.service.ConfigurationService; 043import org.apache.oozie.service.HadoopAccessorException; 044import org.apache.oozie.service.HadoopAccessorService; 045import org.apache.oozie.service.Services; 046import org.apache.oozie.util.XConfiguration; 047import org.apache.oozie.util.XmlUtils; 048import org.jdom.Element; 049 050/** 051 * File system action executor. <p/> This executes the file system mkdir, move and delete commands 052 */ 053public class FsActionExecutor extends ActionExecutor { 054 055 public static final String ACTION_TYPE = "fs"; 056 057 private final int maxGlobCount; 058 059 public FsActionExecutor() { 060 super(ACTION_TYPE); 061 maxGlobCount = ConfigurationService.getInt(LauncherMapper.CONF_OOZIE_ACTION_FS_GLOB_MAX); 062 } 063 064 /** 065 * Initialize Action. 066 */ 067 @Override 068 public void initActionType() { 069 super.initActionType(); 070 registerError(AccessControlException.class.getName(), ActionExecutorException.ErrorType.ERROR, "FS014"); 071 } 072 073 Path getPath(Element element, String attribute) { 074 String str = element.getAttributeValue(attribute).trim(); 075 return new Path(str); 076 } 077 078 void validatePath(Path path, boolean withScheme) throws ActionExecutorException { 079 try { 080 String scheme = path.toUri().getScheme(); 081 if (withScheme) { 082 if (scheme == null) { 083 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS001", 084 "Missing scheme in path [{0}]", path); 085 } 086 else { 087 Services.get().get(HadoopAccessorService.class).checkSupportedFilesystem(path.toUri()); 088 } 089 } 090 else { 091 if (scheme != null) { 092 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS002", 093 "Scheme [{0}] not allowed in path [{1}]", scheme, path); 094 } 095 } 096 } 097 catch (HadoopAccessorException hex) { 098 throw convertException(hex); 099 } 100 } 101 102 Path resolveToFullPath(Path nameNode, Path path, boolean withScheme) throws ActionExecutorException { 103 Path fullPath; 104 105 // If no nameNode is given, validate the path as-is and return it as-is 106 if (nameNode == null) { 107 validatePath(path, withScheme); 108 fullPath = path; 109 } else { 110 // If the path doesn't have a scheme or authority, use the nameNode which should have already been verified earlier 111 String pathScheme = path.toUri().getScheme(); 112 String pathAuthority = path.toUri().getAuthority(); 113 if (pathScheme == null || pathAuthority == null) { 114 if (path.isAbsolute()) { 115 String nameNodeSchemeAuthority = nameNode.toUri().getScheme() + "://" + nameNode.toUri().getAuthority(); 116 fullPath = new Path(nameNodeSchemeAuthority + path.toString()); 117 } else { 118 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS011", 119 "Path [{0}] cannot be relative", path); 120 } 121 } else { 122 // If the path has a scheme and authority, but its not the nameNode then validate the path as-is and return it as-is 123 // If it is the nameNode, then it should have already been verified earlier so return it as-is 124 if (!nameNode.toUri().getScheme().equals(pathScheme) || !nameNode.toUri().getAuthority().equals(pathAuthority)) { 125 validatePath(path, withScheme); 126 } 127 fullPath = path; 128 } 129 } 130 return fullPath; 131 } 132 133 void validateSameNN(Path source, Path dest) throws ActionExecutorException { 134 Path destPath = new Path(source, dest); 135 String t = destPath.toUri().getScheme() + destPath.toUri().getAuthority(); 136 String s = source.toUri().getScheme() + source.toUri().getAuthority(); 137 138 //checking whether NN prefix of source and target is same. can modify this to adjust for a set of multiple whitelisted NN 139 if(!t.equals(s)) { 140 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS007", 141 "move, target NN URI different from that of source", dest); 142 } 143 } 144 145 @SuppressWarnings("unchecked") 146 void doOperations(Context context, Element element) throws ActionExecutorException { 147 try { 148 FileSystem fs = context.getAppFileSystem(); 149 boolean recovery = fs.exists(getRecoveryPath(context)); 150 if (!recovery) { 151 fs.mkdirs(getRecoveryPath(context)); 152 } 153 154 Path nameNodePath = null; 155 Element nameNodeElement = element.getChild("name-node", element.getNamespace()); 156 if (nameNodeElement != null) { 157 String nameNode = nameNodeElement.getTextTrim(); 158 if (nameNode != null) { 159 nameNodePath = new Path(nameNode); 160 // Verify the name node now 161 validatePath(nameNodePath, true); 162 } 163 } 164 165 XConfiguration fsConf = new XConfiguration(); 166 Path appPath = new Path(context.getWorkflow().getAppPath()); 167 // app path could be a file 168 if (fs.isFile(appPath)) { 169 appPath = appPath.getParent(); 170 } 171 JavaActionExecutor.parseJobXmlAndConfiguration(context, element, appPath, fsConf); 172 173 for (Element commandElement : (List<Element>) element.getChildren()) { 174 String command = commandElement.getName(); 175 if (command.equals("mkdir")) { 176 Path path = getPath(commandElement, "path"); 177 mkdir(context, fsConf, nameNodePath, path); 178 } 179 else { 180 if (command.equals("delete")) { 181 Path path = getPath(commandElement, "path"); 182 delete(context, fsConf, nameNodePath, path); 183 } 184 else { 185 if (command.equals("move")) { 186 Path source = getPath(commandElement, "source"); 187 Path target = getPath(commandElement, "target"); 188 move(context, fsConf, nameNodePath, source, target, recovery); 189 } 190 else { 191 if (command.equals("chmod")) { 192 Path path = getPath(commandElement, "path"); 193 boolean recursive = commandElement.getChild("recursive", commandElement.getNamespace()) != null; 194 String str = commandElement.getAttributeValue("dir-files"); 195 boolean dirFiles = (str == null) || Boolean.parseBoolean(str); 196 String permissionsMask = commandElement.getAttributeValue("permissions").trim(); 197 chmod(context, fsConf, nameNodePath, path, permissionsMask, dirFiles, recursive); 198 } 199 else { 200 if (command.equals("touchz")) { 201 Path path = getPath(commandElement, "path"); 202 touchz(context, fsConf, nameNodePath, path); 203 } 204 else { 205 if (command.equals("chgrp")) { 206 Path path = getPath(commandElement, "path"); 207 boolean recursive = commandElement.getChild("recursive", 208 commandElement.getNamespace()) != null; 209 String group = commandElement.getAttributeValue("group"); 210 String str = commandElement.getAttributeValue("dir-files"); 211 boolean dirFiles = (str == null) || Boolean.parseBoolean(str); 212 chgrp(context, fsConf, nameNodePath, path, context.getWorkflow().getUser(), 213 group, dirFiles, recursive); 214 } 215 } 216 } 217 } 218 } 219 } 220 } 221 } 222 catch (Exception ex) { 223 throw convertException(ex); 224 } 225 } 226 227 void chgrp(Context context, XConfiguration fsConf, Path nameNodePath, Path path, String user, String group, 228 boolean dirFiles, boolean recursive) throws ActionExecutorException { 229 230 HashMap<String, String> argsMap = new HashMap<String, String>(); 231 argsMap.put("user", user); 232 argsMap.put("group", group); 233 try { 234 FileSystem fs = getFileSystemFor(path, context, fsConf); 235 path = resolveToFullPath(nameNodePath, path, true); 236 Path[] pathArr = FileUtil.stat2Paths(fs.globStatus(path)); 237 if (pathArr == null || pathArr.length == 0) { 238 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS009", "chgrp" 239 + ", path(s) that matches [{0}] does not exist", path); 240 } 241 checkGlobMax(pathArr); 242 for (Path p : pathArr) { 243 recursiveFsOperation("chgrp", fs, nameNodePath, p, argsMap, dirFiles, recursive, true); 244 } 245 } 246 catch (Exception ex) { 247 throw convertException(ex); 248 } 249 } 250 251 private void recursiveFsOperation(String op, FileSystem fs, Path nameNodePath, Path path, 252 Map<String, String> argsMap, boolean dirFiles, boolean recursive, boolean isRoot) 253 throws ActionExecutorException { 254 255 try { 256 FileStatus pathStatus = fs.getFileStatus(path); 257 List<Path> paths = new ArrayList<Path>(); 258 259 if (dirFiles && pathStatus.isDir()) { 260 if (isRoot) { 261 paths.add(path); 262 } 263 FileStatus[] filesStatus = fs.listStatus(path); 264 for (int i = 0; i < filesStatus.length; i++) { 265 Path p = filesStatus[i].getPath(); 266 paths.add(p); 267 if (recursive && filesStatus[i].isDir()) { 268 recursiveFsOperation(op, fs, null, p, argsMap, dirFiles, recursive, false); 269 } 270 } 271 } 272 else { 273 paths.add(path); 274 } 275 for (Path p : paths) { 276 doFsOperation(op, fs, p, argsMap); 277 } 278 } 279 catch (Exception ex) { 280 throw convertException(ex); 281 } 282 } 283 284 private void doFsOperation(String op, FileSystem fs, Path p, Map<String, String> argsMap) 285 throws ActionExecutorException, IOException { 286 if (op.equals("chmod")) { 287 String permissions = argsMap.get("permissions"); 288 FsPermission newFsPermission = createShortPermission(permissions, p); 289 fs.setPermission(p, newFsPermission); 290 } 291 else if (op.equals("chgrp")) { 292 String user = argsMap.get("user"); 293 String group = argsMap.get("group"); 294 fs.setOwner(p, user, group); 295 } 296 } 297 298 /** 299 * @param path 300 * @param context 301 * @param fsConf 302 * @return FileSystem 303 * @throws HadoopAccessorException 304 */ 305 private FileSystem getFileSystemFor(Path path, Context context, XConfiguration fsConf) throws HadoopAccessorException { 306 String user = context.getWorkflow().getUser(); 307 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 308 JobConf conf = has.createJobConf(path.toUri().getAuthority()); 309 XConfiguration.copy(context.getProtoActionConf(), conf); 310 if (fsConf != null) { 311 XConfiguration.copy(fsConf, conf); 312 } 313 return has.createFileSystem(user, path.toUri(), conf); 314 } 315 316 /** 317 * @param path 318 * @param user 319 * @param group 320 * @return FileSystem 321 * @throws HadoopAccessorException 322 */ 323 private FileSystem getFileSystemFor(Path path, String user) throws HadoopAccessorException { 324 HadoopAccessorService has = Services.get().get(HadoopAccessorService.class); 325 JobConf jobConf = has.createJobConf(path.toUri().getAuthority()); 326 return has.createFileSystem(user, path.toUri(), jobConf); 327 } 328 329 void mkdir(Context context, Path path) throws ActionExecutorException { 330 mkdir(context, null, null, path); 331 } 332 333 void mkdir(Context context, XConfiguration fsConf, Path nameNodePath, Path path) throws ActionExecutorException { 334 try { 335 path = resolveToFullPath(nameNodePath, path, true); 336 FileSystem fs = getFileSystemFor(path, context, fsConf); 337 338 if (!fs.exists(path)) { 339 if (!fs.mkdirs(path)) { 340 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS004", 341 "mkdir, path [{0}] could not create directory", path); 342 } 343 } 344 } 345 catch (Exception ex) { 346 throw convertException(ex); 347 } 348 } 349 350 /** 351 * Delete path 352 * 353 * @param context 354 * @param path 355 * @throws ActionExecutorException 356 */ 357 public void delete(Context context, Path path) throws ActionExecutorException { 358 delete(context, null, null, path); 359 } 360 361 /** 362 * Delete path 363 * 364 * @param context 365 * @param fsConf 366 * @param nameNodePath 367 * @param path 368 * @throws ActionExecutorException 369 */ 370 public void delete(Context context, XConfiguration fsConf, Path nameNodePath, Path path) throws ActionExecutorException { 371 try { 372 path = resolveToFullPath(nameNodePath, path, true); 373 FileSystem fs = getFileSystemFor(path, context, fsConf); 374 Path[] pathArr = FileUtil.stat2Paths(fs.globStatus(path)); 375 if (pathArr != null && pathArr.length > 0) { 376 checkGlobMax(pathArr); 377 for (Path p : pathArr) { 378 if (fs.exists(p)) { 379 if (!fs.delete(p, true)) { 380 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS005", 381 "delete, path [{0}] could not delete path", p); 382 } 383 } 384 } 385 } 386 } 387 catch (Exception ex) { 388 throw convertException(ex); 389 } 390 } 391 392 /** 393 * Delete path 394 * 395 * @param user 396 * @param group 397 * @param path 398 * @throws ActionExecutorException 399 */ 400 public void delete(String user, String group, Path path) throws ActionExecutorException { 401 try { 402 validatePath(path, true); 403 FileSystem fs = getFileSystemFor(path, user); 404 405 if (fs.exists(path)) { 406 if (!fs.delete(path, true)) { 407 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS005", 408 "delete, path [{0}] could not delete path", path); 409 } 410 } 411 } 412 catch (Exception ex) { 413 throw convertException(ex); 414 } 415 } 416 417 /** 418 * Move source to target 419 * 420 * @param context 421 * @param source 422 * @param target 423 * @param recovery 424 * @throws ActionExecutorException 425 */ 426 public void move(Context context, Path source, Path target, boolean recovery) throws ActionExecutorException { 427 move(context, null, null, source, target, recovery); 428 } 429 430 /** 431 * Move source to target 432 * 433 * @param context 434 * @param fsConf 435 * @param nameNodePath 436 * @param source 437 * @param target 438 * @param recovery 439 * @throws ActionExecutorException 440 */ 441 public void move(Context context, XConfiguration fsConf, Path nameNodePath, Path source, Path target, boolean recovery) 442 throws ActionExecutorException { 443 try { 444 source = resolveToFullPath(nameNodePath, source, true); 445 validateSameNN(source, target); 446 FileSystem fs = getFileSystemFor(source, context, fsConf); 447 Path[] pathArr = FileUtil.stat2Paths(fs.globStatus(source)); 448 if (( pathArr == null || pathArr.length == 0 ) ){ 449 if (!recovery) { 450 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS006", 451 "move, source path [{0}] does not exist", source); 452 } else { 453 return; 454 } 455 } 456 if (pathArr.length > 1 && (!fs.exists(target) || fs.isFile(target))) { 457 if(!recovery) { 458 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS012", 459 "move, could not rename multiple sources to the same target name"); 460 } else { 461 return; 462 } 463 } 464 checkGlobMax(pathArr); 465 for (Path p : pathArr) { 466 if (!fs.rename(p, target) && !recovery) { 467 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS008", 468 "move, could not move [{0}] to [{1}]", p, target); 469 } 470 } 471 } 472 catch (Exception ex) { 473 throw convertException(ex); 474 } 475 } 476 477 void chmod(Context context, Path path, String permissions, boolean dirFiles, boolean recursive) throws ActionExecutorException { 478 chmod(context, null, null, path, permissions, dirFiles, recursive); 479 } 480 481 void chmod(Context context, XConfiguration fsConf, Path nameNodePath, Path path, String permissions, 482 boolean dirFiles, boolean recursive) throws ActionExecutorException { 483 484 HashMap<String, String> argsMap = new HashMap<String, String>(); 485 argsMap.put("permissions", permissions); 486 try { 487 FileSystem fs = getFileSystemFor(path, context, fsConf); 488 path = resolveToFullPath(nameNodePath, path, true); 489 Path[] pathArr = FileUtil.stat2Paths(fs.globStatus(path)); 490 if (pathArr == null || pathArr.length == 0) { 491 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS009", "chmod" 492 + ", path(s) that matches [{0}] does not exist", path); 493 } 494 checkGlobMax(pathArr); 495 for (Path p : pathArr) { 496 recursiveFsOperation("chmod", fs, nameNodePath, p, argsMap, dirFiles, recursive, true); 497 } 498 499 } 500 catch (Exception ex) { 501 throw convertException(ex); 502 } 503 } 504 505 void touchz(Context context, Path path) throws ActionExecutorException { 506 touchz(context, null, null, path); 507 } 508 509 void touchz(Context context, XConfiguration fsConf, Path nameNodePath, Path path) throws ActionExecutorException { 510 try { 511 path = resolveToFullPath(nameNodePath, path, true); 512 FileSystem fs = getFileSystemFor(path, context, fsConf); 513 514 FileStatus st; 515 if (fs.exists(path)) { 516 st = fs.getFileStatus(path); 517 if (st.isDir()) { 518 throw new Exception(path.toString() + " is a directory"); 519 } else if (st.getLen() != 0) 520 throw new Exception(path.toString() + " must be a zero-length file"); 521 } 522 FSDataOutputStream out = fs.create(path); 523 out.close(); 524 } 525 catch (Exception ex) { 526 throw convertException(ex); 527 } 528 } 529 530 FsPermission createShortPermission(String permissions, Path path) throws ActionExecutorException { 531 if (permissions.length() == 3) { 532 char user = permissions.charAt(0); 533 char group = permissions.charAt(1); 534 char other = permissions.charAt(2); 535 int useri = user - '0'; 536 int groupi = group - '0'; 537 int otheri = other - '0'; 538 int mask = useri * 100 + groupi * 10 + otheri; 539 short omask = Short.parseShort(Integer.toString(mask), 8); 540 return new FsPermission(omask); 541 } 542 else { 543 if (permissions.length() == 10) { 544 return FsPermission.valueOf(permissions); 545 } 546 else { 547 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS010", 548 "chmod, path [{0}] invalid permissions mask [{1}]", path, permissions); 549 } 550 } 551 } 552 553 @Override 554 public void check(Context context, WorkflowAction action) throws ActionExecutorException { 555 } 556 557 @Override 558 public void kill(Context context, WorkflowAction action) throws ActionExecutorException { 559 } 560 561 @Override 562 public void start(Context context, WorkflowAction action) throws ActionExecutorException { 563 try { 564 context.setStartData("-", "-", "-"); 565 Element actionXml = XmlUtils.parseXml(action.getConf()); 566 doOperations(context, actionXml); 567 context.setExecutionData("OK", null); 568 } 569 catch (Exception ex) { 570 throw convertException(ex); 571 } 572 } 573 574 @Override 575 public void end(Context context, WorkflowAction action) throws ActionExecutorException { 576 String externalStatus = action.getExternalStatus(); 577 WorkflowAction.Status status = externalStatus.equals("OK") ? WorkflowAction.Status.OK : 578 WorkflowAction.Status.ERROR; 579 context.setEndData(status, getActionSignal(status)); 580 if (!context.getProtoActionConf().getBoolean(WorkflowXCommand.KEEP_WF_ACTION_DIR, false)) { 581 try { 582 FileSystem fs = context.getAppFileSystem(); 583 fs.delete(context.getActionDir(), true); 584 } 585 catch (Exception ex) { 586 throw convertException(ex); 587 } 588 } 589 } 590 591 @Override 592 public boolean isCompleted(String externalStatus) { 593 return true; 594 } 595 596 /** 597 * @param context 598 * @return 599 * @throws HadoopAccessorException 600 * @throws IOException 601 * @throws URISyntaxException 602 */ 603 public Path getRecoveryPath(Context context) throws HadoopAccessorException, IOException, URISyntaxException { 604 return new Path(context.getActionDir(), "fs-" + context.getRecoveryId()); 605 } 606 607 private void checkGlobMax(Path[] pathArr) throws ActionExecutorException { 608 if(pathArr.length > maxGlobCount) { 609 throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "FS013", 610 "too many globbed files/dirs to do FS operation"); 611 } 612 } 613 614 public boolean supportsConfigurationJobXML() { 615 return true; 616 } 617}