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.command.wf;
020
021import java.io.IOException;
022import java.net.URI;
023import java.net.URISyntaxException;
024import java.util.ArrayList;
025import java.util.Collection;
026import java.util.Date;
027import java.util.HashMap;
028import java.util.HashSet;
029import java.util.List;
030import java.util.Map;
031import java.util.Set;
032
033import org.apache.hadoop.conf.Configuration;
034import org.apache.hadoop.fs.FileSystem;
035import org.apache.hadoop.fs.Path;
036import org.apache.oozie.AppType;
037import org.apache.oozie.ErrorCode;
038import org.apache.oozie.WorkflowActionBean;
039import org.apache.oozie.WorkflowJobBean;
040import org.apache.oozie.action.oozie.SubWorkflowActionExecutor;
041import org.apache.oozie.client.OozieClient;
042import org.apache.oozie.client.WorkflowAction;
043import org.apache.oozie.client.WorkflowJob;
044import org.apache.oozie.client.rest.JsonBean;
045import org.apache.oozie.command.CommandException;
046import org.apache.oozie.command.PreconditionException;
047import org.apache.oozie.executor.jpa.JPAExecutorException;
048import org.apache.oozie.executor.jpa.WorkflowActionQueryExecutor;
049import org.apache.oozie.executor.jpa.WorkflowActionQueryExecutor.WorkflowActionQuery;
050import org.apache.oozie.executor.jpa.BatchQueryExecutor;
051import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor;
052import org.apache.oozie.executor.jpa.BatchQueryExecutor.UpdateEntry;
053import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor.WorkflowJobQuery;
054import org.apache.oozie.service.ConfigurationService;
055import org.apache.oozie.service.DagXLogInfoService;
056import org.apache.oozie.service.HadoopAccessorException;
057import org.apache.oozie.service.HadoopAccessorService;
058import org.apache.oozie.service.Services;
059import org.apache.oozie.service.UUIDService;
060import org.apache.oozie.service.WorkflowAppService;
061import org.apache.oozie.service.WorkflowStoreService;
062import org.apache.oozie.sla.SLAOperations;
063import org.apache.oozie.sla.service.SLAService;
064import org.apache.oozie.util.ConfigUtils;
065import org.apache.oozie.util.ELEvaluator;
066import org.apache.oozie.util.ELUtils;
067import org.apache.oozie.util.InstrumentUtils;
068import org.apache.oozie.util.LogUtils;
069import org.apache.oozie.util.ParamChecker;
070import org.apache.oozie.util.PropertiesUtils;
071import org.apache.oozie.util.XConfiguration;
072import org.apache.oozie.util.XLog;
073import org.apache.oozie.util.XmlUtils;
074import org.apache.oozie.workflow.WorkflowApp;
075import org.apache.oozie.workflow.WorkflowException;
076import org.apache.oozie.workflow.WorkflowInstance;
077import org.apache.oozie.workflow.WorkflowLib;
078import org.apache.oozie.workflow.lite.NodeHandler;
079import org.jdom.Element;
080import org.jdom.JDOMException;
081
082/**
083 * This is a RerunXCommand which is used for rerunn.
084 *
085 */
086public class ReRunXCommand extends WorkflowXCommand<Void> {
087    private final String jobId;
088    private Configuration conf;
089    private final Set<String> nodesToSkip = new HashSet<String>();
090    public static final String TO_SKIP = "TO_SKIP";
091    private WorkflowJobBean wfBean;
092    private List<WorkflowActionBean> actions;
093    private List<UpdateEntry> updateList = new ArrayList<UpdateEntry>();
094    private List<JsonBean> deleteList = new ArrayList<JsonBean>();
095
096    private static final Set<String> DISALLOWED_USER_PROPERTIES = new HashSet<String>();
097    public static final String DISABLE_CHILD_RERUN = "oozie.wf.rerun.disablechild";
098
099    static {
100        String[] badUserProps = { PropertiesUtils.DAYS, PropertiesUtils.HOURS, PropertiesUtils.MINUTES,
101                PropertiesUtils.KB, PropertiesUtils.MB, PropertiesUtils.GB, PropertiesUtils.TB, PropertiesUtils.PB,
102                PropertiesUtils.RECORDS, PropertiesUtils.MAP_IN, PropertiesUtils.MAP_OUT, PropertiesUtils.REDUCE_IN,
103                PropertiesUtils.REDUCE_OUT, PropertiesUtils.GROUPS };
104        PropertiesUtils.createPropertySet(badUserProps, DISALLOWED_USER_PROPERTIES);
105    }
106
107    public ReRunXCommand(String jobId, Configuration conf) {
108        super("rerun", "rerun", 1);
109        this.jobId = ParamChecker.notEmpty(jobId, "jobId");
110        this.conf = ParamChecker.notNull(conf, "conf");
111    }
112
113    @Override
114    protected void setLogInfo() {
115        LogUtils.setLogInfo(jobId);
116    }
117
118    /* (non-Javadoc)
119     * @see org.apache.oozie.command.XCommand#execute()
120     */
121    @Override
122    protected Void execute() throws CommandException {
123        setupReRun();
124        startWorkflow(jobId);
125        return null;
126    }
127
128    private void startWorkflow(String jobId) throws CommandException {
129        new StartXCommand(jobId).call();
130    }
131
132    private void setupReRun() throws CommandException {
133        InstrumentUtils.incrJobCounter(getName(), 1, getInstrumentation());
134        LogUtils.setLogInfo(wfBean);
135        WorkflowInstance oldWfInstance = this.wfBean.getWorkflowInstance();
136        WorkflowInstance newWfInstance;
137        String appPath = null;
138
139        WorkflowAppService wps = Services.get().get(WorkflowAppService.class);
140        try {
141            XLog.Info.get().setParameter(DagXLogInfoService.TOKEN, conf.get(OozieClient.LOG_TOKEN));
142            WorkflowApp app = wps.parseDef(conf, null);
143            XConfiguration protoActionConf = wps.createProtoActionConf(conf, true);
144            WorkflowLib workflowLib = Services.get().get(WorkflowStoreService.class).getWorkflowLibWithNoDB();
145
146            appPath = conf.get(OozieClient.APP_PATH);
147            URI uri = new URI(appPath);
148            HadoopAccessorService has = Services.get().get(HadoopAccessorService.class);
149            Configuration fsConf = has.createJobConf(uri.getAuthority());
150            FileSystem fs = has.createFileSystem(wfBean.getUser(), uri, fsConf);
151
152            Path configDefault = null;
153            // app path could be a directory
154            Path path = new Path(uri.getPath());
155            if (!fs.isFile(path)) {
156                configDefault = new Path(path, SubmitXCommand.CONFIG_DEFAULT);
157            }
158            else {
159                configDefault = new Path(path.getParent(), SubmitXCommand.CONFIG_DEFAULT);
160            }
161
162            if (fs.exists(configDefault)) {
163                Configuration defaultConf = new XConfiguration(fs.open(configDefault));
164                PropertiesUtils.checkDisallowedProperties(defaultConf, DISALLOWED_USER_PROPERTIES);
165                PropertiesUtils.checkDefaultDisallowedProperties(defaultConf);
166                XConfiguration.injectDefaults(defaultConf, conf);
167            }
168
169            PropertiesUtils.checkDisallowedProperties(conf, DISALLOWED_USER_PROPERTIES);
170
171            // Resolving all variables in the job properties. This ensures the Hadoop Configuration semantics are
172            // preserved. The Configuration.get function within XConfiguration.resolve() works recursively to get the
173            // final value corresponding to a key in the map Resetting the conf to contain all the resolved values is
174            // necessary to ensure propagation of Oozie properties to Hadoop calls downstream
175            conf = ((XConfiguration) conf).resolve();
176
177            try {
178                newWfInstance = workflowLib.createInstance(app, conf, jobId);
179            }
180            catch (WorkflowException e) {
181                throw new CommandException(e);
182            }
183            String appName = ELUtils.resolveAppName(app.getName(), conf);
184            if (SLAService.isEnabled()) {
185                Element wfElem = XmlUtils.parseXml(app.getDefinition());
186                ELEvaluator evalSla = SubmitXCommand.createELEvaluatorForGroup(conf, "wf-sla-submit");
187                Element eSla = XmlUtils.getSLAElement(wfElem);
188                String jobSlaXml = null;
189                if (eSla != null) {
190                    jobSlaXml = SubmitXCommand.resolveSla(eSla, evalSla);
191                }
192                writeSLARegistration(wfElem, jobSlaXml, newWfInstance.getId(),
193                        conf.get(SubWorkflowActionExecutor.PARENT_ID), conf.get(OozieClient.USER_NAME), appName,
194                        evalSla);
195            }
196            wfBean.setAppName(appName);
197            wfBean.setProtoActionConf(protoActionConf.toXmlString());
198        }
199        catch (WorkflowException ex) {
200            throw new CommandException(ex);
201        }
202        catch (IOException ex) {
203            throw new CommandException(ErrorCode.E0803, ex.getMessage(), ex);
204        }
205        catch (HadoopAccessorException ex) {
206            throw new CommandException(ex);
207        }
208        catch (URISyntaxException ex) {
209            throw new CommandException(ErrorCode.E0711, appPath, ex.getMessage(), ex);
210        }
211        catch (Exception ex) {
212            throw new CommandException(ErrorCode.E1007, ex.getMessage(), ex);
213        }
214
215        for (int i = 0; i < actions.size(); i++) {
216            // Skipping to delete the sub workflow when rerun failed node option has been provided. As same
217            // action will be used to rerun the job.
218            if (!nodesToSkip.contains(actions.get(i).getName()) &&
219                    !(conf.getBoolean(OozieClient.RERUN_FAIL_NODES, false) &&
220                    SubWorkflowActionExecutor.ACTION_TYPE.equals(actions.get(i).getType()))) {
221                deleteList.add(actions.get(i));
222                LOG.info("Deleting Action[{0}] for re-run", actions.get(i).getId());
223            }
224            else {
225                copyActionData(newWfInstance, oldWfInstance);
226            }
227        }
228
229        wfBean.setAppPath(conf.get(OozieClient.APP_PATH));
230        wfBean.setConf(XmlUtils.prettyPrint(conf).toString());
231        wfBean.setLogToken(conf.get(OozieClient.LOG_TOKEN, ""));
232        wfBean.setUser(conf.get(OozieClient.USER_NAME));
233        String group = ConfigUtils.getWithDeprecatedCheck(conf, OozieClient.JOB_ACL, OozieClient.GROUP_NAME, null);
234        wfBean.setGroup(group);
235        wfBean.setExternalId(conf.get(OozieClient.EXTERNAL_ID));
236        wfBean.setEndTime(null);
237        wfBean.setRun(wfBean.getRun() + 1);
238        wfBean.setStatus(WorkflowJob.Status.PREP);
239        wfBean.setWorkflowInstance(newWfInstance);
240
241        try {
242            wfBean.setLastModifiedTime(new Date());
243            updateList.add(new UpdateEntry<WorkflowJobQuery>(WorkflowJobQuery.UPDATE_WORKFLOW_RERUN, wfBean));
244            // call JPAExecutor to do the bulk writes
245            BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(null, updateList, deleteList);
246        }
247        catch (JPAExecutorException je) {
248            throw new CommandException(je);
249        }
250        finally {
251            updateParentIfNecessary(wfBean);
252        }
253
254    }
255
256    @SuppressWarnings("unchecked")
257        private void writeSLARegistration(Element wfElem, String jobSlaXml, String id, String parentId, String user,
258            String appName, ELEvaluator evalSla) throws JDOMException, CommandException {
259        if (jobSlaXml != null && jobSlaXml.length() > 0) {
260            Element eSla = XmlUtils.parseXml(jobSlaXml);
261            // insert into new table
262            SLAOperations.createSlaRegistrationEvent(eSla, jobId, parentId, AppType.WORKFLOW_JOB, user, appName, LOG,
263                    true);
264        }
265        // Add sla for wf actions
266        for (Element action : (List<Element>) wfElem.getChildren("action", wfElem.getNamespace())) {
267            Element actionSla = XmlUtils.getSLAElement(action);
268            if (actionSla != null) {
269                String actionSlaXml = SubmitXCommand.resolveSla(actionSla, evalSla);
270                actionSla = XmlUtils.parseXml(actionSlaXml);
271                if (!nodesToSkip.contains(action.getAttributeValue("name"))) {
272                    String actionId = Services.get().get(UUIDService.class)
273                            .generateChildId(jobId, action.getAttributeValue("name") + "");
274                    SLAOperations.createSlaRegistrationEvent(actionSla, actionId, jobId, AppType.WORKFLOW_ACTION, user,
275                            appName, LOG, true);
276                }
277            }
278        }
279
280    }
281
282    /**
283     * Loading the Wfjob and workflow actions. Parses the config and adds the nodes that are to be skipped to the
284     * skipped node list
285     *
286     * @throws CommandException
287     */
288    @Override
289    protected void eagerLoadState() throws CommandException {
290        try {
291            this.wfBean = WorkflowJobQueryExecutor.getInstance().get(WorkflowJobQuery.GET_WORKFLOW_STATUS, this.jobId);
292            this.actions = WorkflowActionQueryExecutor.getInstance().getList(
293                    WorkflowActionQuery.GET_ACTIONS_FOR_WORKFLOW_RERUN, this.jobId);
294
295            if (conf != null) {
296                if (conf.getBoolean(OozieClient.RERUN_FAIL_NODES, false) == false) { //Rerun with skipNodes
297                    Collection<String> skipNodes = conf.getStringCollection(OozieClient.RERUN_SKIP_NODES);
298                    for (String str : skipNodes) {
299                        // trimming is required
300                        nodesToSkip.add(str.trim());
301                    }
302                    LOG.debug("Skipnode size :" + nodesToSkip.size());
303                }
304                else {
305                    for (WorkflowActionBean action : actions) { // Rerun from failed nodes
306                        if (action.getStatus() == WorkflowAction.Status.OK) {
307                            nodesToSkip.add(action.getName());
308                        }
309                    }
310                    LOG.debug("Skipnode size are to rerun from FAIL nodes :" + nodesToSkip.size());
311                }
312                StringBuilder tmp = new StringBuilder();
313                for (String node : nodesToSkip) {
314                    tmp.append(node).append(",");
315                }
316                LOG.debug("SkipNode List :" + tmp);
317            }
318        }
319        catch (Exception ex) {
320            throw new CommandException(ErrorCode.E0603, ex.getMessage(), ex);
321        }
322    }
323
324    /**
325     * Checks the pre-conditions that are required for workflow to recover - Last run of Workflow should be completed -
326     * The nodes that are to be skipped are to be completed successfully in the base run.
327     *
328     * @throws org.apache.oozie.command.CommandException,PreconditionException On failure of pre-conditions
329     */
330    @Override
331    protected void eagerVerifyPrecondition() throws CommandException, PreconditionException {
332        // Throwing error if parent exist and same workflow trying to rerun, when running child workflow disabled
333        // through conf.
334        if (wfBean.getParentId() != null && !conf.getBoolean(SubWorkflowActionExecutor.SUBWORKFLOW_RERUN, false)
335                && ConfigurationService.getBoolean(DISABLE_CHILD_RERUN)) {
336            throw new PreconditionException(ErrorCode.E0755, " Rerun is not allowed through child workflow, please" +
337                    " re-run through the parent " + wfBean.getParentId());
338        }
339
340        if (!(wfBean.getStatus().equals(WorkflowJob.Status.FAILED)
341                || wfBean.getStatus().equals(WorkflowJob.Status.KILLED) || wfBean.getStatus().equals(
342                        WorkflowJob.Status.SUCCEEDED))) {
343            throw new CommandException(ErrorCode.E0805, wfBean.getStatus());
344        }
345        Set<String> unmachedNodes = new HashSet<String>(nodesToSkip);
346        for (WorkflowActionBean action : actions) {
347            if (nodesToSkip.contains(action.getName())) {
348                if (!action.getStatus().equals(WorkflowAction.Status.OK)
349                        && !action.getStatus().equals(WorkflowAction.Status.ERROR)) {
350                    throw new CommandException(ErrorCode.E0806, action.getName());
351                }
352                unmachedNodes.remove(action.getName());
353            }
354        }
355        if (unmachedNodes.size() > 0) {
356            StringBuilder sb = new StringBuilder();
357            String separator = "";
358            for (String s : unmachedNodes) {
359                sb.append(separator).append(s);
360                separator = ",";
361            }
362            throw new CommandException(ErrorCode.E0807, sb);
363        }
364    }
365
366    /**
367     * Copys the variables for skipped nodes from the old wfInstance to new one.
368     *
369     * @param newWfInstance : Source WF instance object
370     * @param oldWfInstance : Update WF instance
371     */
372    private void copyActionData(WorkflowInstance newWfInstance, WorkflowInstance oldWfInstance) {
373        Map<String, String> oldVars = new HashMap<String, String>();
374        Map<String, String> newVars = new HashMap<String, String>();
375        oldVars = oldWfInstance.getAllVars();
376        for (String var : oldVars.keySet()) {
377            String actionName = var.split(WorkflowInstance.NODE_VAR_SEPARATOR)[0];
378            if (nodesToSkip.contains(actionName)) {
379                newVars.put(var, oldVars.get(var));
380            }
381        }
382        for (String node : nodesToSkip) {
383            // Setting the TO_SKIP variable to true. This will be used by
384            // SignalCommand and LiteNodeHandler to skip the action.
385            newVars.put(node + WorkflowInstance.NODE_VAR_SEPARATOR + TO_SKIP, "true");
386            String visitedFlag = NodeHandler.getLoopFlag(node);
387            // Removing the visited flag so that the action won't be considered
388            // a loop.
389            if (newVars.containsKey(visitedFlag)) {
390                newVars.remove(visitedFlag);
391            }
392        }
393        newWfInstance.setAllVars(newVars);
394    }
395
396    /* (non-Javadoc)
397     * @see org.apache.oozie.command.XCommand#getEntityKey()
398     */
399    @Override
400    public String getEntityKey() {
401        return this.jobId;
402    }
403
404    /* (non-Javadoc)
405     * @see org.apache.oozie.command.XCommand#isLockRequired()
406     */
407    @Override
408    protected boolean isLockRequired() {
409        return true;
410    }
411
412    /* (non-Javadoc)
413     * @see org.apache.oozie.command.XCommand#loadState()
414     */
415    @Override
416    protected void loadState() throws CommandException {
417        try {
418            this.wfBean = WorkflowJobQueryExecutor.getInstance().get(WorkflowJobQuery.GET_WORKFLOW_RERUN, this.jobId);
419            this.actions = WorkflowActionQueryExecutor.getInstance().getList(
420                    WorkflowActionQuery.GET_ACTIONS_FOR_WORKFLOW_RERUN, this.jobId);
421        }
422        catch (JPAExecutorException jpe) {
423            throw new CommandException(jpe);
424        }
425    }
426
427    /* (non-Javadoc)
428     * @see org.apache.oozie.command.XCommand#verifyPrecondition()
429     */
430    @Override
431    protected void verifyPrecondition() throws CommandException, PreconditionException {
432        eagerVerifyPrecondition();
433    }
434}