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}