001/**
002 * Licensed to the Apache Software Foundation (ASF) under one
003 * or more contributor license agreements.  See the NOTICE file
004 * distributed with this work for additional information
005 * regarding copyright ownership.  The ASF licenses this file
006 * to you under the Apache License, Version 2.0 (the
007 * "License"); you may not use this file except in compliance
008 * with the License.  You may obtain a copy of the License at
009 *
010 *      http://www.apache.org/licenses/LICENSE-2.0
011 *
012 * Unless required by applicable law or agreed to in writing, software
013 * distributed under the License is distributed on an "AS IS" BASIS,
014 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
015 * See the License for the specific language governing permissions and
016 * limitations under the License.
017 */
018
019package org.apache.oozie.action.oozie;
020
021import java.io.IOException;
022import java.io.StringReader;
023import java.util.HashSet;
024import java.util.Set;
025
026import org.apache.hadoop.conf.Configuration;
027import org.apache.oozie.DagEngine;
028import org.apache.oozie.LocalOozieClient;
029import org.apache.oozie.WorkflowJobBean;
030import org.apache.oozie.action.ActionExecutor;
031import org.apache.oozie.action.ActionExecutorException;
032import org.apache.oozie.action.hadoop.OozieJobInfo;
033import org.apache.oozie.client.OozieClient;
034import org.apache.oozie.client.OozieClientException;
035import org.apache.oozie.client.WorkflowAction;
036import org.apache.oozie.client.WorkflowJob;
037import org.apache.oozie.command.CommandException;
038import org.apache.oozie.command.wf.ActionStartXCommand;
039import org.apache.oozie.service.ConfigurationService;
040import org.apache.oozie.service.DagEngineService;
041import org.apache.oozie.service.Services;
042import org.apache.oozie.util.ConfigUtils;
043import org.apache.oozie.util.JobUtils;
044import org.apache.oozie.util.PropertiesUtils;
045import org.apache.oozie.util.XConfiguration;
046import org.apache.oozie.util.XLog;
047import org.apache.oozie.util.XmlUtils;
048import org.jdom.Element;
049import org.jdom.Namespace;
050
051public class SubWorkflowActionExecutor extends ActionExecutor {
052    public static final String ACTION_TYPE = "sub-workflow";
053    public static final String LOCAL = "local";
054    public static final String PARENT_ID = "oozie.wf.parent.id";
055    public static final String SUPER_PARENT_ID = "oozie.wf.superparent.id";
056    public static final String SUBWORKFLOW_MAX_DEPTH = "oozie.action.subworkflow.max.depth";
057    public static final String SUBWORKFLOW_DEPTH = "oozie.action.subworkflow.depth";
058    public static final String SUBWORKFLOW_RERUN = "oozie.action.subworkflow.rerun";
059
060    private static final Set<String> DISALLOWED_DEFAULT_PROPERTIES = new HashSet<String>();
061    public XLog LOG = XLog.getLog(getClass());
062
063
064    static {
065        String[] badUserProps = {PropertiesUtils.DAYS, PropertiesUtils.HOURS, PropertiesUtils.MINUTES,
066                PropertiesUtils.KB, PropertiesUtils.MB, PropertiesUtils.GB, PropertiesUtils.TB, PropertiesUtils.PB,
067                PropertiesUtils.RECORDS, PropertiesUtils.MAP_IN, PropertiesUtils.MAP_OUT, PropertiesUtils.REDUCE_IN,
068                PropertiesUtils.REDUCE_OUT, PropertiesUtils.GROUPS};
069
070        String[] badDefaultProps = {PropertiesUtils.HADOOP_USER};
071        PropertiesUtils.createPropertySet(badUserProps, DISALLOWED_DEFAULT_PROPERTIES);
072        PropertiesUtils.createPropertySet(badDefaultProps, DISALLOWED_DEFAULT_PROPERTIES);
073    }
074
075    protected SubWorkflowActionExecutor() {
076        super(ACTION_TYPE);
077    }
078
079    public void initActionType() {
080        super.initActionType();
081    }
082
083    protected OozieClient getWorkflowClient(Context context, String oozieUri) {
084        OozieClient oozieClient;
085        if (oozieUri.equals(LOCAL)) {
086            WorkflowJobBean workflow = (WorkflowJobBean) context.getWorkflow();
087            String user = workflow.getUser();
088            String group = workflow.getGroup();
089            DagEngine dagEngine = Services.get().get(DagEngineService.class).getDagEngine(user);
090            oozieClient = new LocalOozieClient(dagEngine);
091        }
092        else {
093            // TODO we need to add authToken to the WC for the remote case
094            oozieClient = new OozieClient(oozieUri);
095        }
096        return oozieClient;
097    }
098
099    protected void injectInline(Element eConf, Configuration subWorkflowConf) throws IOException,
100            ActionExecutorException {
101        if (eConf != null) {
102            String strConf = XmlUtils.prettyPrint(eConf).toString();
103            Configuration conf = new XConfiguration(new StringReader(strConf));
104            try {
105                PropertiesUtils.checkDisallowedProperties(conf, DISALLOWED_DEFAULT_PROPERTIES);
106            }
107            catch (CommandException ex) {
108                throw convertException(ex);
109            }
110            XConfiguration.copy(conf, subWorkflowConf);
111        }
112    }
113
114    @SuppressWarnings("unchecked")
115    protected void injectCallback(Context context, Configuration conf) {
116        String callback = context.getCallbackUrl("$status");
117        if (conf.get(OozieClient.WORKFLOW_NOTIFICATION_URL) != null) {
118            XLog.getLog(getClass())
119                    .warn("Sub-Workflow configuration has a custom job end notification URI, overriding");
120        }
121        conf.set(OozieClient.WORKFLOW_NOTIFICATION_URL, callback);
122    }
123
124    protected void injectRecovery(String externalId, Configuration conf) {
125        conf.set(OozieClient.EXTERNAL_ID, externalId);
126    }
127
128    protected void injectParent(String parentId, Configuration conf) {
129        conf.set(PARENT_ID, parentId);
130    }
131
132    protected void injectSuperParent(WorkflowJob parentWorkflow, Configuration parentConf, Configuration conf) {
133        String superParentId = parentConf.get(SUPER_PARENT_ID);
134        if (superParentId == null) {
135            // This is a sub-workflow at depth 1
136            superParentId = parentWorkflow.getParentId();
137
138            // If the parent workflow is not submitted through a coordinator then the parentId will be the super parent id.
139            if (superParentId == null) {
140                superParentId = parentWorkflow.getId();
141            }
142            conf.set(SUPER_PARENT_ID, superParentId);
143        } else {
144            // Sub-workflow at depth 2 or more.
145            conf.set(SUPER_PARENT_ID, superParentId);
146        }
147    }
148
149    protected void verifyAndInjectSubworkflowDepth(Configuration parentConf, Configuration conf) throws ActionExecutorException {
150        int depth = parentConf.getInt(SUBWORKFLOW_DEPTH, 0);
151        int maxDepth = ConfigurationService.getInt(SUBWORKFLOW_MAX_DEPTH);
152        if (depth >= maxDepth) {
153            throw new ActionExecutorException(ActionExecutorException.ErrorType.ERROR, "SUBWF001",
154                    "Depth [{0}] cannot exceed maximum subworkflow depth [{1}]", (depth + 1), maxDepth);
155        }
156        conf.setInt(SUBWORKFLOW_DEPTH, depth + 1);
157    }
158
159    protected String checkIfRunning(OozieClient oozieClient, String extId) throws OozieClientException {
160        String jobId = oozieClient.getJobId(extId);
161        if (jobId.equals("")) {
162            return null;
163        }
164        return jobId;
165    }
166
167    public void start(Context context, WorkflowAction action) throws ActionExecutorException {
168        try {
169            Element eConf = XmlUtils.parseXml(action.getConf());
170            Namespace ns = eConf.getNamespace();
171            Element e = eConf.getChild("oozie", ns);
172            String oozieUri = (e == null) ? LOCAL : e.getTextTrim();
173            OozieClient oozieClient = getWorkflowClient(context, oozieUri);
174            String subWorkflowId = null;
175            String extId = context.getRecoveryId();
176            String runningJobId = null;
177            if (extId != null) {
178                runningJobId = checkIfRunning(oozieClient, extId);
179            }
180            if (runningJobId == null) {
181                String appPath = eConf.getChild("app-path", ns).getTextTrim();
182
183                XConfiguration subWorkflowConf = new XConfiguration();
184
185                Configuration parentConf = new XConfiguration(new StringReader(context.getWorkflow().getConf()));
186
187                if (eConf.getChild(("propagate-configuration"), ns) != null) {
188                    XConfiguration.copy(parentConf, subWorkflowConf);
189                }
190
191                // Propagate coordinator and bundle info to subworkflow
192                if (OozieJobInfo.isJobInfoEnabled()) {
193                  if (parentConf.get(OozieJobInfo.COORD_ID) != null) {
194                    subWorkflowConf.set(OozieJobInfo.COORD_ID, parentConf.get(OozieJobInfo.COORD_ID));
195                    subWorkflowConf.set(OozieJobInfo.COORD_NAME, parentConf.get(OozieJobInfo.COORD_NAME));
196                    subWorkflowConf.set(OozieJobInfo.COORD_NOMINAL_TIME, parentConf.get(OozieJobInfo.COORD_NOMINAL_TIME));
197                  }
198                  if (parentConf.get(OozieJobInfo.BUNDLE_ID) != null) {
199                    subWorkflowConf.set(OozieJobInfo.BUNDLE_ID, parentConf.get(OozieJobInfo.BUNDLE_ID));
200                    subWorkflowConf.set(OozieJobInfo.BUNDLE_NAME, parentConf.get(OozieJobInfo.BUNDLE_NAME));
201                  }
202                }
203
204                // the proto has the necessary credentials
205                Configuration protoActionConf = context.getProtoActionConf();
206                XConfiguration.copy(protoActionConf, subWorkflowConf);
207                subWorkflowConf.set(OozieClient.APP_PATH, appPath);
208                String group = ConfigUtils.getWithDeprecatedCheck(parentConf, OozieClient.JOB_ACL, OozieClient.GROUP_NAME, null);
209                if(group != null) {
210                    subWorkflowConf.set(OozieClient.GROUP_NAME, group);
211                }
212
213                injectInline(eConf.getChild("configuration", ns), subWorkflowConf);
214                injectCallback(context, subWorkflowConf);
215                injectRecovery(extId, subWorkflowConf);
216                injectParent(context.getWorkflow().getId(), subWorkflowConf);
217                injectSuperParent(context.getWorkflow(), parentConf, subWorkflowConf);
218                verifyAndInjectSubworkflowDepth(parentConf, subWorkflowConf);
219
220                //TODO: this has to be refactored later to be done in a single place for REST calls and this
221                JobUtils.normalizeAppPath(context.getWorkflow().getUser(), context.getWorkflow().getGroup(),
222                                          subWorkflowConf);
223
224                subWorkflowConf.set(OOZIE_ACTION_YARN_TAG, getActionYarnTag(parentConf, context.getWorkflow(), action));
225
226                // if the rerun failed node option is provided during the time of rerun command, old subworkflow will
227                // rerun again.
228                if(action.getExternalId() != null && parentConf.getBoolean(OozieClient.RERUN_FAIL_NODES, false)) {
229                    subWorkflowConf.setBoolean(SUBWORKFLOW_RERUN, true);
230                    oozieClient.reRun(action.getExternalId(), subWorkflowConf.toProperties());
231                    subWorkflowId = action.getExternalId();
232                } else {
233                    subWorkflowId = oozieClient.run(subWorkflowConf.toProperties());
234                }
235            }
236            else {
237                subWorkflowId = runningJobId;
238            }
239            WorkflowJob workflow = oozieClient.getJobInfo(subWorkflowId);
240            String consoleUrl = workflow.getConsoleUrl();
241            context.setStartData(subWorkflowId, oozieUri, consoleUrl);
242            if (runningJobId != null) {
243                check(context, action);
244            }
245        }
246        catch (Exception ex) {
247            LOG.error(ex);
248            throw convertException(ex);
249        }
250    }
251
252    public void end(Context context, WorkflowAction action) throws ActionExecutorException {
253        try {
254            String externalStatus = action.getExternalStatus();
255            WorkflowAction.Status status = externalStatus.equals("SUCCEEDED") ? WorkflowAction.Status.OK
256                                           : WorkflowAction.Status.ERROR;
257            context.setEndData(status, getActionSignal(status));
258        }
259        catch (Exception ex) {
260            throw convertException(ex);
261        }
262    }
263
264    public void check(Context context, WorkflowAction action) throws ActionExecutorException {
265        try {
266            String subWorkflowId = action.getExternalId();
267            String oozieUri = action.getTrackerUri();
268            OozieClient oozieClient = getWorkflowClient(context, oozieUri);
269            WorkflowJob subWorkflow = oozieClient.getJobInfo(subWorkflowId);
270            WorkflowJob.Status status = subWorkflow.getStatus();
271            switch (status) {
272                case FAILED:
273                case KILLED:
274                case SUCCEEDED:
275                    context.setExecutionData(status.toString(), null);
276                    break;
277                default:
278                    context.setExternalStatus(status.toString());
279                    break;
280            }
281        }
282        catch (Exception ex) {
283            throw convertException(ex);
284        }
285    }
286
287    public void kill(Context context, WorkflowAction action) throws ActionExecutorException {
288        try {
289            String subWorkflowId = action.getExternalId();
290            String oozieUri = action.getTrackerUri();
291            if (subWorkflowId != null && oozieUri != null) {
292                OozieClient oozieClient = getWorkflowClient(context, oozieUri);
293                oozieClient.kill(subWorkflowId);
294            }
295            context.setEndData(WorkflowAction.Status.KILLED, getActionSignal(WorkflowAction.Status.KILLED));
296        }
297        catch (Exception ex) {
298            throw convertException(ex);
299        }
300    }
301
302    private static Set<String> FINAL_STATUS = new HashSet<String>();
303
304    static {
305        FINAL_STATUS.add("SUCCEEDED");
306        FINAL_STATUS.add("KILLED");
307        FINAL_STATUS.add("FAILED");
308    }
309
310    public boolean isCompleted(String externalStatus) {
311        return FINAL_STATUS.contains(externalStatus);
312    }
313
314    public boolean supportsConfigurationJobXML() {
315        return true;
316    }
317}