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.workflow.lite;
020
021import org.apache.commons.codec.binary.Base64;
022import org.apache.hadoop.io.Writable;
023import org.apache.oozie.action.hadoop.FsActionExecutor;
024import org.apache.oozie.action.oozie.SubWorkflowActionExecutor;
025import org.apache.oozie.service.ConfigurationService;
026import org.apache.oozie.util.ELUtils;
027import org.apache.oozie.util.IOUtils;
028import org.apache.oozie.util.XConfiguration;
029import org.apache.oozie.util.XmlUtils;
030import org.apache.oozie.util.ParamChecker;
031import org.apache.oozie.util.ParameterVerifier;
032import org.apache.oozie.util.ParameterVerifierException;
033import org.apache.oozie.util.WritableUtils;
034import org.apache.oozie.ErrorCode;
035import org.apache.oozie.workflow.WorkflowException;
036import org.apache.oozie.action.ActionExecutor;
037import org.apache.oozie.service.Services;
038import org.apache.oozie.service.ActionService;
039import org.apache.commons.lang.StringUtils;
040import org.apache.hadoop.conf.Configuration;
041import org.jdom.Element;
042import org.jdom.JDOMException;
043import org.jdom.Namespace;
044import org.xml.sax.SAXException;
045
046import javax.xml.transform.stream.StreamSource;
047import javax.xml.validation.Schema;
048import javax.xml.validation.Validator;
049
050import java.io.IOException;
051import java.io.Reader;
052import java.io.StringReader;
053import java.io.StringWriter;
054import java.io.ByteArrayOutputStream;
055import java.io.ByteArrayInputStream;
056import java.io.DataInputStream;
057import java.io.DataInput;
058import java.io.DataOutput;
059import java.io.DataOutputStream;
060import java.util.ArrayList;
061import java.util.Arrays;
062import java.util.Deque;
063import java.util.HashMap;
064import java.util.HashSet;
065import java.util.LinkedList;
066import java.util.List;
067import java.util.Map;
068import java.util.zip.*;
069
070/**
071 * Class to parse and validate workflow xml
072 */
073public class LiteWorkflowAppParser {
074
075    private static final String DECISION_E = "decision";
076    private static final String ACTION_E = "action";
077    private static final String END_E = "end";
078    private static final String START_E = "start";
079    private static final String JOIN_E = "join";
080    private static final String FORK_E = "fork";
081    private static final Object KILL_E = "kill";
082
083    private static final String SLA_INFO = "info";
084    private static final String CREDENTIALS = "credentials";
085    private static final String GLOBAL = "global";
086    private static final String PARAMETERS = "parameters";
087
088    private static final String NAME_A = "name";
089    private static final String CRED_A = "cred";
090    private static final String USER_RETRY_MAX_A = "retry-max";
091    private static final String USER_RETRY_INTERVAL_A = "retry-interval";
092    private static final String TO_A = "to";
093
094    private static final String FORK_PATH_E = "path";
095    private static final String FORK_START_A = "start";
096
097    private static final String ACTION_OK_E = "ok";
098    private static final String ACTION_ERROR_E = "error";
099
100    private static final String DECISION_SWITCH_E = "switch";
101    private static final String DECISION_CASE_E = "case";
102    private static final String DECISION_DEFAULT_E = "default";
103
104    private static final String SUBWORKFLOW_E = "sub-workflow";
105
106    private static final String KILL_MESSAGE_E = "message";
107    public static final String VALIDATE_FORK_JOIN = "oozie.validate.ForkJoin";
108    public static final String WF_VALIDATE_FORK_JOIN = "oozie.wf.validate.ForkJoin";
109
110    public static final String DEFAULT_NAME_NODE = "oozie.actions.default.name-node";
111    public static final String DEFAULT_JOB_TRACKER = "oozie.actions.default.job-tracker";
112    public static final String OOZIE_GLOBAL = "oozie.wf.globalconf";
113
114    private static final String JOB_TRACKER = "job-tracker";
115    private static final String NAME_NODE = "name-node";
116    private static final String JOB_XML = "job-xml";
117    private static final String CONFIGURATION = "configuration";
118
119    private Schema schema;
120    private Class<? extends ControlNodeHandler> controlNodeHandler;
121    private Class<? extends DecisionNodeHandler> decisionHandlerClass;
122    private Class<? extends ActionNodeHandler> actionHandlerClass;
123
124    private static enum VisitStatus {
125        VISITING, VISITED
126    }
127
128    /**
129     * We use this to store a node name and its top (eldest) decision parent node name for the forkjoin validation
130     */
131    class NodeAndTopDecisionParent {
132        String node;
133        String topDecisionParent;
134
135        public NodeAndTopDecisionParent(String node, String topDecisionParent) {
136            this.node = node;
137            this.topDecisionParent = topDecisionParent;
138        }
139    }
140
141    private List<String> forkList = new ArrayList<String>();
142    private List<String> joinList = new ArrayList<String>();
143    private StartNodeDef startNode;
144    private List<NodeAndTopDecisionParent> visitedOkNodes = new ArrayList<NodeAndTopDecisionParent>();
145    private List<String> visitedJoinNodes = new ArrayList<String>();
146
147    private String defaultNameNode;
148    private String defaultJobTracker;
149
150    public LiteWorkflowAppParser(Schema schema,
151                                 Class<? extends ControlNodeHandler> controlNodeHandler,
152                                 Class<? extends DecisionNodeHandler> decisionHandlerClass,
153                                 Class<? extends ActionNodeHandler> actionHandlerClass) throws WorkflowException {
154        this.schema = schema;
155        this.controlNodeHandler = controlNodeHandler;
156        this.decisionHandlerClass = decisionHandlerClass;
157        this.actionHandlerClass = actionHandlerClass;
158
159        defaultNameNode = ConfigurationService.get(DEFAULT_NAME_NODE);
160        if (defaultNameNode != null) {
161            defaultNameNode = defaultNameNode.trim();
162            if (defaultNameNode.isEmpty()) {
163                defaultNameNode = null;
164            }
165        }
166        defaultJobTracker = ConfigurationService.get(DEFAULT_JOB_TRACKER);
167        if (defaultJobTracker != null) {
168            defaultJobTracker = defaultJobTracker.trim();
169            if (defaultJobTracker.isEmpty()) {
170                defaultJobTracker = null;
171            }
172        }
173    }
174
175    public LiteWorkflowApp validateAndParse(Reader reader, Configuration jobConf) throws WorkflowException {
176        return validateAndParse(reader, jobConf, null);
177    }
178
179    /**
180     * Parse and validate xml to {@link LiteWorkflowApp}
181     *
182     * @param reader
183     * @return LiteWorkflowApp
184     * @throws WorkflowException
185     */
186    public LiteWorkflowApp validateAndParse(Reader reader, Configuration jobConf, Configuration configDefault)
187            throws WorkflowException {
188        try {
189            StringWriter writer = new StringWriter();
190            IOUtils.copyCharStream(reader, writer);
191            String strDef = writer.toString();
192
193            if (schema != null) {
194                Validator validator = schema.newValidator();
195                validator.validate(new StreamSource(new StringReader(strDef)));
196            }
197
198            Element wfDefElement = XmlUtils.parseXml(strDef);
199            ParameterVerifier.verifyParameters(jobConf, wfDefElement);
200            LiteWorkflowApp app = parse(strDef, wfDefElement, configDefault, jobConf);
201            Map<String, VisitStatus> traversed = new HashMap<String, VisitStatus>();
202            traversed.put(app.getNode(StartNodeDef.START).getName(), VisitStatus.VISITING);
203            validate(app, app.getNode(StartNodeDef.START), traversed);
204            //Validate whether fork/join are in pair or not
205            if (jobConf.getBoolean(WF_VALIDATE_FORK_JOIN, true)
206                    && ConfigurationService.getBoolean(VALIDATE_FORK_JOIN)) {
207                validateForkJoin(app);
208            }
209            return app;
210        }
211        catch (ParameterVerifierException ex) {
212            throw new WorkflowException(ex);
213        }
214        catch (JDOMException ex) {
215            throw new WorkflowException(ErrorCode.E0700, ex.getMessage(), ex);
216        }
217        catch (SAXException ex) {
218            throw new WorkflowException(ErrorCode.E0701, ex.getMessage(), ex);
219        }
220        catch (IOException ex) {
221            throw new WorkflowException(ErrorCode.E0702, ex.getMessage(), ex);
222        }
223    }
224
225    /**
226     * Validate whether fork/join are in pair or not
227     * @param app LiteWorkflowApp
228     * @throws WorkflowException
229     */
230    private void validateForkJoin(LiteWorkflowApp app) throws WorkflowException {
231        // Make sure the number of forks and joins in wf are equal
232        if (forkList.size() != joinList.size()) {
233            throw new WorkflowException(ErrorCode.E0730);
234        }
235
236        // No need to bother going through all of this if there are no fork/join nodes
237        if (!forkList.isEmpty()) {
238            visitedOkNodes.clear();
239            visitedJoinNodes.clear();
240            validateForkJoin(startNode, app, new LinkedList<String>(), new LinkedList<String>(), new LinkedList<String>(), true,
241                    null);
242        }
243    }
244
245    /*
246     * Recursively walk through the DAG and make sure that all fork paths are valid.
247     * This should be called from validateForkJoin(LiteWorkflowApp app).  It assumes that visitedOkNodes and visitedJoinNodes are
248     * both empty ArrayLists on the first call.
249     *
250     * @param node the current node; use the startNode on the first call
251     * @param app the WorkflowApp
252     * @param forkNodes a stack of the current fork nodes
253     * @param joinNodes a stack of the current join nodes
254     * @param path a stack of the current path
255     * @param okTo false if node (or an ancestor of node) was gotten to via an "error to" transition or via a join node that has
256     * already been visited at least once before
257     * @param topDecisionParent The top (eldest) decision node along the path to this node, or null if there isn't one
258     * @throws WorkflowException
259     */
260    private void validateForkJoin(NodeDef node, LiteWorkflowApp app, Deque<String> forkNodes, Deque<String> joinNodes,
261            Deque<String> path, boolean okTo, String topDecisionParent) throws WorkflowException {
262        if (path.contains(node.getName())) {
263            // cycle
264            throw new WorkflowException(ErrorCode.E0741, node.getName(), Arrays.toString(path.toArray()));
265        }
266        path.push(node.getName());
267
268        // Make sure that we're not revisiting a node (that's not a Kill, Join, or End type) that's been visited before from an
269        // "ok to" transition; if its from an "error to" transition, then its okay to visit it multiple times.  Also, because we
270        // traverse through join nodes multiple times, we have to make sure not to throw an exception here when we're really just
271        // re-walking the same execution path (this is why we need the visitedJoinNodes list used later)
272        if (okTo && !(node instanceof KillNodeDef) && !(node instanceof JoinNodeDef) && !(node instanceof EndNodeDef)) {
273            NodeAndTopDecisionParent natdp = findInVisitedOkNodes(node.getName());
274            if (natdp != null) {
275                // However, if we've visited the node and it's under a decision node, we may be seeing it again and it's only
276                // illegal if that decision node is not the same as what we're seeing now (because during execution we only go
277                // down one path of the decision node, so while we're seeing the node multiple times here, during runtime it will
278                // only be executed once).  Also, this decision node should be the top (eldest) decision node.  As null indicates
279                // that there isn't a decision node, when this happens they must both be null to be valid.  Here is a good example
280                // to visualize a node ("actionX") that has three "ok to" paths to it, but should still be a valid workflow (it may
281                // be easier to see if you draw it):
282                    // decisionA --> {actionX, decisionB}
283                    // decisionB --> {actionX, actionY}
284                    // actionY   --> {actionX}
285                // And, if we visit this node twice under the same decision node in an invalid way, the path cycle checking code
286                // will catch it, so we don't have to worry about that here.
287                if ((natdp.topDecisionParent == null && topDecisionParent == null)
288                     || (natdp.topDecisionParent == null && topDecisionParent != null)
289                     || (natdp.topDecisionParent != null && topDecisionParent == null)
290                     || !natdp.topDecisionParent.equals(topDecisionParent)) {
291                    // If we get here, then we've seen this node before from an "ok to" transition but they don't have the same
292                    // decision node top parent, which means that this node will be executed twice, which is illegal
293                    throw new WorkflowException(ErrorCode.E0743, node.getName());
294                }
295            }
296            else {
297                // If we haven't transitioned to this node before, add it and its top decision parent node
298                visitedOkNodes.add(new NodeAndTopDecisionParent(node.getName(), topDecisionParent));
299            }
300        }
301
302        if (node instanceof StartNodeDef) {
303            String transition = node.getTransitions().get(0);   // start always has only 1 transition
304            NodeDef tranNode = app.getNode(transition);
305            validateForkJoin(tranNode, app, forkNodes, joinNodes, path, okTo, topDecisionParent);
306        }
307        else if (node instanceof ActionNodeDef) {
308            String transition = node.getTransitions().get(0);   // "ok to" transition
309            NodeDef tranNode = app.getNode(transition);
310            validateForkJoin(tranNode, app, forkNodes, joinNodes, path, okTo, topDecisionParent);  // propogate okTo
311            transition = node.getTransitions().get(1);          // "error to" transition
312            tranNode = app.getNode(transition);
313            validateForkJoin(tranNode, app, forkNodes, joinNodes, path, false, topDecisionParent); // use false
314        }
315        else if (node instanceof DecisionNodeDef) {
316            for(String transition : (new HashSet<String>(node.getTransitions()))) {
317                NodeDef tranNode = app.getNode(transition);
318                // if there currently isn't a topDecisionParent (i.e. null), then use this node instead of propagating null
319                String parentDecisionNode = topDecisionParent;
320                if (parentDecisionNode == null) {
321                    parentDecisionNode = node.getName();
322                }
323                validateForkJoin(tranNode, app, forkNodes, joinNodes, path, okTo, parentDecisionNode);
324            }
325        }
326        else if (node instanceof ForkNodeDef) {
327            forkNodes.push(node.getName());
328            List<String> transitionsList = node.getTransitions();
329            HashSet<String> transitionsSet = new HashSet<String>(transitionsList);
330            // Check that a fork doesn't go to the same node more than once
331            if (!transitionsList.isEmpty() && transitionsList.size() != transitionsSet.size()) {
332                // Now we have to figure out which node is the problem and what type of node they are (join and kill are ok)
333                for (int i = 0; i < transitionsList.size(); i++) {
334                    String a = transitionsList.get(i);
335                    NodeDef aNode = app.getNode(a);
336                    if (!(aNode instanceof JoinNodeDef) && !(aNode instanceof KillNodeDef)) {
337                        for (int k = i+1; k < transitionsList.size(); k++) {
338                            String b = transitionsList.get(k);
339                            if (a.equals(b)) {
340                                throw new WorkflowException(ErrorCode.E0744, node.getName(), a);
341                            }
342                        }
343                    }
344                }
345            }
346            for(String transition : transitionsSet) {
347                NodeDef tranNode = app.getNode(transition);
348                validateForkJoin(tranNode, app, forkNodes, joinNodes, path, okTo, topDecisionParent);
349            }
350            forkNodes.pop();
351            if (!joinNodes.isEmpty()) {
352                joinNodes.pop();
353            }
354        }
355        else if (node instanceof JoinNodeDef) {
356            if (forkNodes.isEmpty()) {
357                // no fork for join to match with
358                throw new WorkflowException(ErrorCode.E0742, node.getName());
359            }
360            if (forkNodes.size() > joinNodes.size() && (joinNodes.isEmpty() || !joinNodes.peek().equals(node.getName()))) {
361                joinNodes.push(node.getName());
362            }
363            if (!joinNodes.peek().equals(node.getName())) {
364                // join doesn't match fork
365                throw new WorkflowException(ErrorCode.E0732, forkNodes.peek(), node.getName(), joinNodes.peek());
366            }
367            joinNodes.pop();
368            String currentForkNode = forkNodes.pop();
369            String transition = node.getTransitions().get(0);   // join always has only 1 transition
370            NodeDef tranNode = app.getNode(transition);
371            // If we're already under a situation where okTo is false, use false (propogate it)
372            // Or if we've already visited this join node, use false (because we've already traversed this path before and we don't
373            // want to throw an exception from the check against visitedOkNodes)
374            if (!okTo || visitedJoinNodes.contains(node.getName())) {
375                validateForkJoin(tranNode, app, forkNodes, joinNodes, path, false, topDecisionParent);
376            // Else, use true because this is either the first time we've gone through this join node or okTo was already false
377            } else {
378                visitedJoinNodes.add(node.getName());
379                validateForkJoin(tranNode, app, forkNodes, joinNodes, path, true, topDecisionParent);
380            }
381            forkNodes.push(currentForkNode);
382            joinNodes.push(node.getName());
383        }
384        else if (node instanceof KillNodeDef) {
385            // do nothing
386        }
387        else if (node instanceof EndNodeDef) {
388            if (!forkNodes.isEmpty()) {
389                path.pop();     // = node
390                String parent = path.peek();
391                // can't go to an end node in a fork
392                throw new WorkflowException(ErrorCode.E0737, parent, node.getName());
393            }
394        }
395        else {
396            // invalid node type (shouldn't happen)
397            throw new WorkflowException(ErrorCode.E0740, node.getName());
398        }
399        path.pop();
400    }
401
402    /**
403     * Return a {@link NodeAndTopDecisionParent} whose {@link NodeAndTopDecisionParent#node} is equal to the passed in name, or null
404     * if it isn't in the {@link LiteWorkflowAppParser#visitedOkNodes} list.
405     *
406     * @param name The name to search for
407     * @return a NodeAndTopDecisionParent or null
408     */
409    private NodeAndTopDecisionParent findInVisitedOkNodes(String name) {
410        NodeAndTopDecisionParent natdp = null;
411        for (NodeAndTopDecisionParent v : visitedOkNodes) {
412            if (v.node.equals(name)) {
413                natdp = v;
414                break;
415            }
416        }
417        return natdp;
418    }
419
420    /**
421     * Parse xml to {@link LiteWorkflowApp}
422     *
423     * @param strDef
424     * @param root
425     * @param configDefault
426     * @param jobConf
427     * @return LiteWorkflowApp
428     * @throws WorkflowException
429     */
430    @SuppressWarnings({"unchecked"})
431    private LiteWorkflowApp parse(String strDef, Element root, Configuration configDefault, Configuration jobConf)
432            throws WorkflowException {
433        Namespace ns = root.getNamespace();
434        LiteWorkflowApp def = null;
435        GlobalSectionData gData = jobConf.get(OOZIE_GLOBAL) == null ?
436                null : getGlobalFromString(jobConf.get(OOZIE_GLOBAL));
437        boolean serializedGlobalConf = false;
438        for (Element eNode : (List<Element>) root.getChildren()) {
439            if (eNode.getName().equals(START_E)) {
440                def = new LiteWorkflowApp(root.getAttributeValue(NAME_A), strDef,
441                                          new StartNodeDef(controlNodeHandler, eNode.getAttributeValue(TO_A)));
442            } else if (eNode.getName().equals(END_E)) {
443                def.addNode(new EndNodeDef(eNode.getAttributeValue(NAME_A), controlNodeHandler));
444            } else if (eNode.getName().equals(KILL_E)) {
445                def.addNode(new KillNodeDef(eNode.getAttributeValue(NAME_A),
446                                            eNode.getChildText(KILL_MESSAGE_E, ns), controlNodeHandler));
447            } else if (eNode.getName().equals(FORK_E)) {
448                List<String> paths = new ArrayList<String>();
449                for (Element tran : (List<Element>) eNode.getChildren(FORK_PATH_E, ns)) {
450                    paths.add(tran.getAttributeValue(FORK_START_A));
451                }
452                def.addNode(new ForkNodeDef(eNode.getAttributeValue(NAME_A), controlNodeHandler, paths));
453            } else if (eNode.getName().equals(JOIN_E)) {
454                def.addNode(new JoinNodeDef(eNode.getAttributeValue(NAME_A), controlNodeHandler, eNode.getAttributeValue(TO_A)));
455            } else if (eNode.getName().equals(DECISION_E)) {
456                Element eSwitch = eNode.getChild(DECISION_SWITCH_E, ns);
457                List<String> transitions = new ArrayList<String>();
458                for (Element e : (List<Element>) eSwitch.getChildren(DECISION_CASE_E, ns)) {
459                    transitions.add(e.getAttributeValue(TO_A));
460                }
461                transitions.add(eSwitch.getChild(DECISION_DEFAULT_E, ns).getAttributeValue(TO_A));
462
463                String switchStatement = XmlUtils.prettyPrint(eSwitch).toString();
464                def.addNode(new DecisionNodeDef(eNode.getAttributeValue(NAME_A), switchStatement, decisionHandlerClass,
465                                                transitions));
466            } else if (ACTION_E.equals(eNode.getName())) {
467                String[] transitions = new String[2];
468                Element eActionConf = null;
469                for (Element elem : (List<Element>) eNode.getChildren()) {
470                    if (ACTION_OK_E.equals(elem.getName())) {
471                        transitions[0] = elem.getAttributeValue(TO_A);
472                    } else if (ACTION_ERROR_E.equals(elem.getName())) {
473                        transitions[1] = elem.getAttributeValue(TO_A);
474                    } else if (SLA_INFO.equals(elem.getName()) || CREDENTIALS.equals(elem.getName())) {
475                        continue;
476                    } else {
477                        if (!serializedGlobalConf && elem.getName().equals(SubWorkflowActionExecutor.ACTION_TYPE) &&
478                                elem.getChild(("propagate-configuration"), ns) != null && gData != null) {
479                            serializedGlobalConf = true;
480                            jobConf.set(OOZIE_GLOBAL, getGlobalString(gData));
481                        }
482                        eActionConf = elem;
483                        if (SUBWORKFLOW_E.equals(elem.getName())) {
484                            handleDefaultsAndGlobal(gData, null, elem);
485                        }
486                        else {
487                            handleDefaultsAndGlobal(gData, configDefault, elem);
488                        }
489                    }
490                }
491
492                String credStr = eNode.getAttributeValue(CRED_A);
493                String userRetryMaxStr = eNode.getAttributeValue(USER_RETRY_MAX_A);
494                String userRetryIntervalStr = eNode.getAttributeValue(USER_RETRY_INTERVAL_A);
495                try {
496                    if (!StringUtils.isEmpty(userRetryMaxStr)) {
497                        userRetryMaxStr = ELUtils.resolveAppName(userRetryMaxStr, jobConf);
498                    }
499                    if (!StringUtils.isEmpty(userRetryIntervalStr)) {
500                        userRetryIntervalStr = ELUtils.resolveAppName(userRetryIntervalStr, jobConf);
501                    }
502                }
503                catch (Exception e) {
504                    throw new WorkflowException(ErrorCode.E0703, e.getMessage());
505                }
506
507                String actionConf = XmlUtils.prettyPrint(eActionConf).toString();
508                def.addNode(new ActionNodeDef(eNode.getAttributeValue(NAME_A), actionConf, actionHandlerClass,
509                                              transitions[0], transitions[1], credStr,
510                                              userRetryMaxStr, userRetryIntervalStr));
511            } else if (SLA_INFO.equals(eNode.getName()) || CREDENTIALS.equals(eNode.getName())) {
512                // No operation is required
513            } else if (eNode.getName().equals(GLOBAL)) {
514                if(jobConf.get(OOZIE_GLOBAL) != null) {
515                    gData = getGlobalFromString(jobConf.get(OOZIE_GLOBAL));
516                    handleDefaultsAndGlobal(gData, null, eNode);
517                }
518                gData = parseGlobalSection(ns, eNode);
519            } else if (eNode.getName().equals(PARAMETERS)) {
520                // No operation is required
521            } else {
522                throw new WorkflowException(ErrorCode.E0703, eNode.getName());
523            }
524        }
525        return def;
526    }
527
528    /**
529     * Read the GlobalSectionData from Base64 string.
530     * @param globalStr
531     * @return GlobalSectionData
532     * @throws WorkflowException
533     */
534    private GlobalSectionData getGlobalFromString(String globalStr) throws WorkflowException {
535        GlobalSectionData globalSectionData = new GlobalSectionData();
536        try {
537            byte[] data = Base64.decodeBase64(globalStr);
538            Inflater inflater = new Inflater();
539            DataInputStream ois = new DataInputStream(new InflaterInputStream(new ByteArrayInputStream(data), inflater));
540            globalSectionData.readFields(ois);
541            ois.close();
542        } catch (Exception ex) {
543            throw new WorkflowException(ErrorCode.E0700, "Error while processing global section conf");
544        }
545        return globalSectionData;
546    }
547
548
549    /**
550     * Write the GlobalSectionData to a Base64 string.
551     * @param globalSectionData
552     * @return String
553     * @throws WorkflowException
554     */
555    private String getGlobalString(GlobalSectionData globalSectionData) throws WorkflowException {
556        ByteArrayOutputStream baos = new ByteArrayOutputStream();
557        DataOutputStream oos = null;
558        try {
559            Deflater def = new Deflater();
560            oos = new DataOutputStream(new DeflaterOutputStream(baos, def));
561            globalSectionData.write(oos);
562            oos.close();
563        } catch (IOException e) {
564            throw new WorkflowException(ErrorCode.E0700, "Error while processing global section conf");
565        }
566        return Base64.encodeBase64String(baos.toByteArray());
567    }
568
569    /**
570     * Validate workflow xml
571     *
572     * @param app
573     * @param node
574     * @param traversed
575     * @throws WorkflowException
576     */
577    private void validate(LiteWorkflowApp app, NodeDef node, Map<String, VisitStatus> traversed) throws WorkflowException {
578        if (node instanceof StartNodeDef) {
579            startNode = (StartNodeDef) node;
580        }
581        else {
582            try {
583                ParamChecker.validateActionName(node.getName());
584            }
585            catch (IllegalArgumentException ex) {
586                throw new WorkflowException(ErrorCode.E0724, ex.getMessage());
587            }
588        }
589        if (node instanceof ActionNodeDef) {
590            try {
591                Element action = XmlUtils.parseXml(node.getConf());
592                boolean supportedAction = Services.get().get(ActionService.class).getExecutor(action.getName()) != null;
593                if (!supportedAction) {
594                    throw new WorkflowException(ErrorCode.E0723, node.getName(), action.getName());
595                }
596            }
597            catch (JDOMException ex) {
598                throw new RuntimeException("It should never happen, " + ex.getMessage(), ex);
599            }
600        }
601
602        if(node instanceof ForkNodeDef){
603            forkList.add(node.getName());
604        }
605
606        if(node instanceof JoinNodeDef){
607            joinList.add(node.getName());
608        }
609
610        if (node instanceof EndNodeDef) {
611            traversed.put(node.getName(), VisitStatus.VISITED);
612            return;
613        }
614        if (node instanceof KillNodeDef) {
615            traversed.put(node.getName(), VisitStatus.VISITED);
616            return;
617        }
618        for (String transition : node.getTransitions()) {
619
620            if (app.getNode(transition) == null) {
621                throw new WorkflowException(ErrorCode.E0708, node.getName(), transition);
622            }
623
624            //check if it is a cycle
625            if (traversed.get(app.getNode(transition).getName()) == VisitStatus.VISITING) {
626                throw new WorkflowException(ErrorCode.E0707, app.getNode(transition).getName());
627            }
628            //ignore validated one
629            if (traversed.get(app.getNode(transition).getName()) == VisitStatus.VISITED) {
630                continue;
631            }
632
633            traversed.put(app.getNode(transition).getName(), VisitStatus.VISITING);
634            validate(app, app.getNode(transition), traversed);
635        }
636        traversed.put(node.getName(), VisitStatus.VISITED);
637    }
638
639    private void addChildElement(Element parent, Namespace ns, String childName, String childValue) {
640        Element child = new Element(childName, ns);
641        child.setText(childValue);
642        parent.addContent(child);
643    }
644
645    private class GlobalSectionData implements Writable {
646        String jobTracker;
647        String nameNode;
648        List<String> jobXmls;
649        Configuration conf;
650
651        public GlobalSectionData() {
652        }
653
654        public GlobalSectionData(String jobTracker, String nameNode, List<String> jobXmls, Configuration conf) {
655            this.jobTracker = jobTracker;
656            this.nameNode = nameNode;
657            this.jobXmls = jobXmls;
658            this.conf = conf;
659        }
660
661        @Override
662        public void write(DataOutput dataOutput) throws IOException {
663            WritableUtils.writeStr(dataOutput, jobTracker);
664            WritableUtils.writeStr(dataOutput, nameNode);
665
666            if(jobXmls != null && !jobXmls.isEmpty()) {
667                dataOutput.writeInt(jobXmls.size());
668                for (String content : jobXmls) {
669                    WritableUtils.writeStr(dataOutput, content);
670                }
671            } else {
672                dataOutput.writeInt(0);
673            }
674            if(conf != null) {
675                WritableUtils.writeStr(dataOutput, XmlUtils.prettyPrint(conf).toString());
676            } else {
677                WritableUtils.writeStr(dataOutput, null);
678            }
679        }
680
681        @Override
682        public void readFields(DataInput dataInput) throws IOException {
683            jobTracker = WritableUtils.readStr(dataInput);
684            nameNode = WritableUtils.readStr(dataInput);
685            int length = dataInput.readInt();
686            if (length > 0) {
687                jobXmls = new ArrayList<String>();
688                for (int i = 0; i < length; i++) {
689                    jobXmls.add(WritableUtils.readStr(dataInput));
690                }
691            }
692            String confString = WritableUtils.readStr(dataInput);
693            if(confString != null) {
694                conf = new XConfiguration(new StringReader(confString));
695            }
696        }
697    }
698
699    private GlobalSectionData parseGlobalSection(Namespace ns, Element global) throws WorkflowException {
700        GlobalSectionData gData = null;
701        if (global != null) {
702            String globalJobTracker = null;
703            Element globalJobTrackerElement = global.getChild(JOB_TRACKER, ns);
704            if (globalJobTrackerElement != null) {
705                globalJobTracker = globalJobTrackerElement.getValue();
706            }
707
708            String globalNameNode = null;
709            Element globalNameNodeElement = global.getChild(NAME_NODE, ns);
710            if (globalNameNodeElement != null) {
711                globalNameNode = globalNameNodeElement.getValue();
712            }
713
714            List<String> globalJobXmls = null;
715            @SuppressWarnings("unchecked")
716            List<Element> globalJobXmlElements = global.getChildren(JOB_XML, ns);
717            if (!globalJobXmlElements.isEmpty()) {
718                globalJobXmls = new ArrayList<String>(globalJobXmlElements.size());
719                for(Element jobXmlElement: globalJobXmlElements) {
720                    globalJobXmls.add(jobXmlElement.getText());
721                }
722            }
723
724            Configuration globalConf = null;
725            Element globalConfigurationElement = global.getChild(CONFIGURATION, ns);
726            if (globalConfigurationElement != null) {
727                try {
728                    globalConf = new XConfiguration(new StringReader(XmlUtils.prettyPrint(globalConfigurationElement).toString()));
729                } catch (IOException ioe) {
730                    throw new WorkflowException(ErrorCode.E0700, "Error while processing global section conf");
731                }
732            }
733            gData = new GlobalSectionData(globalJobTracker, globalNameNode, globalJobXmls, globalConf);
734        }
735        return gData;
736    }
737
738    private void handleDefaultsAndGlobal(GlobalSectionData gData, Configuration configDefault, Element actionElement)
739            throws WorkflowException {
740
741        ActionExecutor ae = Services.get().get(ActionService.class).getExecutor(actionElement.getName());
742        if (ae == null && !GLOBAL.equals(actionElement.getName())) {
743            throw new WorkflowException(ErrorCode.E0723, actionElement.getName(), ActionService.class.getName());
744        }
745
746        Namespace actionNs = actionElement.getNamespace();
747
748        // If this is the global section or ActionExecutor.requiresNameNodeJobTracker() returns true, we parse the action's
749        // <name-node> and <job-tracker> fields.  If those aren't defined, we take them from the <global> section.  If those
750        // aren't defined, we take them from the oozie-site defaults.  If those aren't defined, we throw a WorkflowException.
751        // However, for the SubWorkflow and FS Actions, as well as the <global> section, we don't throw the WorkflowException.
752        // Also, we only parse the NN (not the JT) for the FS Action.
753        if (SubWorkflowActionExecutor.ACTION_TYPE.equals(actionElement.getName()) ||
754                FsActionExecutor.ACTION_TYPE.equals(actionElement.getName()) ||
755                GLOBAL.equals(actionElement.getName()) || ae.requiresNameNodeJobTracker()) {
756            if (actionElement.getChild(NAME_NODE, actionNs) == null) {
757                if (gData != null && gData.nameNode != null) {
758                    addChildElement(actionElement, actionNs, NAME_NODE, gData.nameNode);
759                } else if (defaultNameNode != null) {
760                    addChildElement(actionElement, actionNs, NAME_NODE, defaultNameNode);
761                } else if (!(SubWorkflowActionExecutor.ACTION_TYPE.equals(actionElement.getName()) ||
762                        FsActionExecutor.ACTION_TYPE.equals(actionElement.getName()) ||
763                        GLOBAL.equals(actionElement.getName()))) {
764                    throw new WorkflowException(ErrorCode.E0701, "No " + NAME_NODE + " defined");
765                }
766            }
767            if (actionElement.getChild(JOB_TRACKER, actionNs) == null &&
768                    !FsActionExecutor.ACTION_TYPE.equals(actionElement.getName())) {
769                if (gData != null && gData.jobTracker != null) {
770                    addChildElement(actionElement, actionNs, JOB_TRACKER, gData.jobTracker);
771                } else if (defaultJobTracker != null) {
772                    addChildElement(actionElement, actionNs, JOB_TRACKER, defaultJobTracker);
773                } else if (!(SubWorkflowActionExecutor.ACTION_TYPE.equals(actionElement.getName()) ||
774                        GLOBAL.equals(actionElement.getName()))) {
775                    throw new WorkflowException(ErrorCode.E0701, "No " + JOB_TRACKER + " defined");
776                }
777            }
778        }
779
780        // If this is the global section or ActionExecutor.supportsConfigurationJobXML() returns true, we parse the action's
781        // <configuration> and <job-xml> fields.  We also merge this with those from the <global> section, if given.  If none are
782        // defined, empty values are placed.  Exceptions are thrown if there's an error parsing, but not if they're not given.
783        if ( GLOBAL.equals(actionElement.getName()) || ae.supportsConfigurationJobXML()) {
784            @SuppressWarnings("unchecked")
785            List<Element> actionJobXmls = actionElement.getChildren(JOB_XML, actionNs);
786            if (gData != null && gData.jobXmls != null) {
787                for(String gJobXml : gData.jobXmls) {
788                    boolean alreadyExists = false;
789                    for (Element actionXml : actionJobXmls) {
790                        if (gJobXml.equals(actionXml.getText())) {
791                            alreadyExists = true;
792                            break;
793                        }
794                    }
795                    if (!alreadyExists) {
796                        Element ejobXml = new Element(JOB_XML, actionNs);
797                        ejobXml.setText(gJobXml);
798                        actionElement.addContent(ejobXml);
799                    }
800                }
801            }
802
803            try {
804                XConfiguration actionConf = new XConfiguration();
805                if (configDefault != null)
806                    XConfiguration.copy(configDefault, actionConf);
807                if (gData != null && gData.conf != null) {
808                    XConfiguration.copy(gData.conf, actionConf);
809                }
810                Element actionConfiguration = actionElement.getChild(CONFIGURATION, actionNs);
811                if (actionConfiguration != null) {
812                    //copy and override
813                    XConfiguration.copy(new XConfiguration(new StringReader(XmlUtils.prettyPrint(actionConfiguration).toString())),
814                            actionConf);
815                }
816                int position = actionElement.indexOf(actionConfiguration);
817                actionElement.removeContent(actionConfiguration); //replace with enhanced one
818                Element eConfXml = XmlUtils.parseXml(actionConf.toXmlString(false));
819                eConfXml.detach();
820                eConfXml.setNamespace(actionNs);
821                if (position > 0) {
822                    actionElement.addContent(position, eConfXml);
823                }
824                else {
825                    actionElement.addContent(eConfXml);
826                }
827            }
828            catch (IOException e) {
829                throw new WorkflowException(ErrorCode.E0700, "Error while processing action conf");
830            }
831            catch (JDOMException e) {
832                throw new WorkflowException(ErrorCode.E0700, "Error while processing action conf");
833            }
834        }
835    }
836}