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}