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