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.util.ArrayList;
022import java.util.Date;
023import java.util.List;
024import javax.servlet.jsp.el.ELException;
025
026import org.apache.hadoop.conf.Configuration;
027import org.apache.oozie.ErrorCode;
028import org.apache.oozie.FaultInjection;
029import org.apache.oozie.SLAEventBean;
030import org.apache.oozie.WorkflowActionBean;
031import org.apache.oozie.WorkflowJobBean;
032import org.apache.oozie.XException;
033import org.apache.oozie.action.ActionExecutor;
034import org.apache.oozie.action.ActionExecutorException;
035import org.apache.oozie.action.control.ControlNodeActionExecutor;
036import org.apache.oozie.client.OozieClient;
037import org.apache.oozie.client.WorkflowAction;
038import org.apache.oozie.client.WorkflowJob;
039import org.apache.oozie.client.SLAEvent.SlaAppType;
040import org.apache.oozie.client.SLAEvent.Status;
041import org.apache.oozie.client.rest.JsonBean;
042import org.apache.oozie.command.CommandException;
043import org.apache.oozie.command.PreconditionException;
044import org.apache.oozie.executor.jpa.BatchQueryExecutor.UpdateEntry;
045import org.apache.oozie.executor.jpa.BatchQueryExecutor;
046import org.apache.oozie.executor.jpa.JPAExecutorException;
047import org.apache.oozie.executor.jpa.WorkflowActionQueryExecutor;
048import org.apache.oozie.executor.jpa.WorkflowActionQueryExecutor.WorkflowActionQuery;
049import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor;
050import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor.WorkflowJobQuery;
051import org.apache.oozie.service.ActionService;
052import org.apache.oozie.service.EventHandlerService;
053import org.apache.oozie.service.JPAService;
054import org.apache.oozie.service.Services;
055import org.apache.oozie.service.UUIDService;
056import org.apache.oozie.util.ELEvaluationException;
057import org.apache.oozie.util.Instrumentation;
058import org.apache.oozie.util.LogUtils;
059import org.apache.oozie.util.XLog;
060import org.apache.oozie.util.XmlUtils;
061import org.apache.oozie.util.db.SLADbXOperations;
062
063@SuppressWarnings("deprecation")
064public class ActionStartXCommand extends ActionXCommand<org.apache.oozie.command.wf.ActionXCommand.ActionExecutorContext> {
065    public static final String EL_ERROR = "EL_ERROR";
066    public static final String EL_EVAL_ERROR = "EL_EVAL_ERROR";
067    public static final String COULD_NOT_START = "COULD_NOT_START";
068    public static final String START_DATA_MISSING = "START_DATA_MISSING";
069    public static final String EXEC_DATA_MISSING = "EXEC_DATA_MISSING";
070
071    private String jobId = null;
072    protected String actionId = null;
073    protected WorkflowJobBean wfJob = null;
074    protected WorkflowActionBean wfAction = null;
075    private JPAService jpaService = null;
076    private ActionExecutor executor = null;
077    private List<UpdateEntry> updateList = new ArrayList<UpdateEntry>();
078    private List<JsonBean> insertList = new ArrayList<JsonBean>();
079    protected ActionExecutorContext context = null;
080
081    public ActionStartXCommand(String actionId, String type) {
082        super("action.start", type, 0);
083        this.actionId = actionId;
084        this.jobId = Services.get().get(UUIDService.class).getId(actionId);
085    }
086
087    public ActionStartXCommand(WorkflowJobBean job, String actionId, String type) {
088        super("action.start", type, 0);
089        this.actionId = actionId;
090        this.wfJob = job;
091        this.jobId = wfJob.getId();
092    }
093
094    @Override
095    protected void setLogInfo() {
096        LogUtils.setLogInfo(actionId);
097    }
098
099    @Override
100    protected boolean isLockRequired() {
101        return true;
102    }
103
104    @Override
105    public String getEntityKey() {
106        return this.jobId;
107    }
108
109    @Override
110    protected void loadState() throws CommandException {
111        try {
112            jpaService = Services.get().get(JPAService.class);
113            if (jpaService != null) {
114                if (wfJob == null) {
115                    this.wfJob = WorkflowJobQueryExecutor.getInstance().get(WorkflowJobQuery.GET_WORKFLOW, jobId);
116                }
117                this.wfAction = WorkflowActionQueryExecutor.getInstance().get(WorkflowActionQuery.GET_ACTION, actionId);
118                LogUtils.setLogInfo( wfJob);
119                LogUtils.setLogInfo(wfAction);
120            }
121            else {
122                throw new CommandException(ErrorCode.E0610);
123            }
124        }
125        catch (XException ex) {
126            throw new CommandException(ex);
127        }
128    }
129
130    @Override
131    protected void verifyPrecondition() throws CommandException, PreconditionException {
132        if (wfJob == null) {
133            throw new PreconditionException(ErrorCode.E0604, jobId);
134        }
135        if (wfAction == null) {
136            throw new PreconditionException(ErrorCode.E0605, actionId);
137        }
138        if (wfAction.isPending()
139                && (wfAction.getStatus() == WorkflowActionBean.Status.PREP
140                        || wfAction.getStatus() == WorkflowActionBean.Status.START_RETRY
141                        || wfAction.getStatus() == WorkflowActionBean.Status.START_MANUAL
142                        || wfAction.getStatus() == WorkflowActionBean.Status.USER_RETRY
143                        )) {
144            if (wfJob.getStatus() != WorkflowJob.Status.RUNNING) {
145                throw new PreconditionException(ErrorCode.E0810, WorkflowJob.Status.RUNNING.toString());
146            }
147        }
148        else {
149            throw new PreconditionException(ErrorCode.E0816, wfAction.isPending(), wfAction.getStatusStr());
150        }
151
152        executor = Services.get().get(ActionService.class).getExecutor(wfAction.getType());
153        if (executor == null) {
154            throw new CommandException(ErrorCode.E0802, wfAction.getType());
155        }
156    }
157
158    @Override
159    protected ActionExecutorContext execute() throws CommandException {
160        LOG.debug("STARTED ActionStartXCommand for wf actionId=" + actionId);
161        Configuration conf = wfJob.getWorkflowInstance().getConf();
162
163        int maxRetries = 0;
164        long retryInterval = 0;
165        boolean execSynchronous = false;
166
167        if (!(executor instanceof ControlNodeActionExecutor)) {
168            maxRetries = conf.getInt(OozieClient.ACTION_MAX_RETRIES, executor.getMaxRetries());
169            retryInterval = conf.getLong(OozieClient.ACTION_RETRY_INTERVAL, executor.getRetryInterval());
170        }
171
172        executor.setMaxRetries(maxRetries);
173        executor.setRetryInterval(retryInterval);
174
175        try {
176            boolean isRetry = false;
177            if (wfAction.getStatus() == WorkflowActionBean.Status.START_RETRY
178                    || wfAction.getStatus() == WorkflowActionBean.Status.START_MANUAL) {
179                isRetry = true;
180                prepareForRetry(wfAction);
181            }
182            boolean isUserRetry = false;
183            if (wfAction.getStatus() == WorkflowActionBean.Status.USER_RETRY) {
184                isUserRetry = true;
185                prepareForRetry(wfAction);
186            }
187            context = getContext(isRetry, isUserRetry);
188            boolean caught = false;
189            try {
190                if (!(executor instanceof ControlNodeActionExecutor)) {
191                    String tmpActionConf = XmlUtils.removeComments(wfAction.getConf());
192                    String actionConf = context.getELEvaluator().evaluate(tmpActionConf, String.class);
193                    wfAction.setConf(actionConf);
194                    LOG.debug("Start, name [{0}] type [{1}] configuration{E}{E}{2}{E}", wfAction.getName(), wfAction
195                            .getType(), actionConf);
196                }
197            }
198            catch (ELEvaluationException ex) {
199                caught = true;
200                throw new ActionExecutorException(ActionExecutorException.ErrorType.TRANSIENT, EL_EVAL_ERROR, ex
201                        .getMessage(), ex);
202            }
203            catch (ELException ex) {
204                caught = true;
205                context.setErrorInfo(EL_ERROR, ex.getMessage());
206                LOG.warn("ELException in ActionStartXCommand ", ex.getMessage(), ex);
207                handleError(context, wfJob, wfAction);
208            }
209            catch (org.jdom.JDOMException je) {
210                caught = true;
211                context.setErrorInfo("ParsingError", je.getMessage());
212                LOG.warn("JDOMException in ActionStartXCommand ", je.getMessage(), je);
213                handleError(context, wfJob, wfAction);
214            }
215            catch (Exception ex) {
216                caught = true;
217                context.setErrorInfo(EL_ERROR, ex.getMessage());
218                LOG.warn("Exception in ActionStartXCommand ", ex.getMessage(), ex);
219                handleError(context, wfJob, wfAction);
220            }
221            if(!caught) {
222                wfAction.setErrorInfo(null, null);
223                incrActionCounter(wfAction.getType(), 1);
224
225                LOG.info("Start action [{0}] with user-retry state : userRetryCount [{1}], userRetryMax [{2}], userRetryInterval [{3}]",
226                                wfAction.getId(), wfAction.getUserRetryCount(), wfAction.getUserRetryMax(), wfAction
227                                        .getUserRetryInterval());
228
229                Instrumentation.Cron cron = new Instrumentation.Cron();
230                cron.start();
231                context.setStartTime();
232                executor.start(context, wfAction);
233                cron.stop();
234                FaultInjection.activate("org.apache.oozie.command.SkipCommitFaultInjection");
235                addActionCron(wfAction.getType(), cron);
236
237                wfAction.setRetries(0);
238                if (wfAction.isExecutionComplete()) {
239                    if (!context.isExecuted()) {
240                        LOG.warn(XLog.OPS, "Action Completed, ActionExecutor [{0}] must call setExecutionData()", executor
241                                .getType());
242                        wfAction.setErrorInfo(EXEC_DATA_MISSING,
243                                "Execution Complete, but Execution Data Missing from Action");
244                        failJob(context);
245                    } else {
246                        wfAction.setPending();
247                        if (!(executor instanceof ControlNodeActionExecutor)) {
248                            queue(new ActionEndXCommand(wfAction.getId(), wfAction.getType()));
249                        }
250                        else {
251                            execSynchronous = true;
252                        }
253                    }
254                }
255                else {
256                    if (!context.isStarted()) {
257                        LOG.warn(XLog.OPS, "Action Started, ActionExecutor [{0}] must call setStartData()", executor
258                                .getType());
259                        wfAction.setErrorInfo(START_DATA_MISSING, "Execution Started, but Start Data Missing from Action");
260                        failJob(context);
261                    } else {
262                        queue(new WorkflowNotificationXCommand(wfJob, wfAction));
263                    }
264                }
265
266                LOG.info(XLog.STD, "[***" + wfAction.getId() + "***]" + "Action status=" + wfAction.getStatusStr());
267
268                updateList.add(new UpdateEntry<WorkflowActionQuery>(WorkflowActionQuery.UPDATE_ACTION_START, wfAction));
269                updateJobLastModified();
270                // Add SLA status event (STARTED) for WF_ACTION
271                SLAEventBean slaEvent = SLADbXOperations.createStatusEvent(wfAction.getSlaXml(), wfAction.getId(), Status.STARTED,
272                        SlaAppType.WORKFLOW_ACTION);
273                if(slaEvent != null) {
274                    insertList.add(slaEvent);
275                }
276                LOG.info(XLog.STD, "[***" + wfAction.getId() + "***]" + "Action updated in DB!");
277            }
278        }
279        catch (ActionExecutorException ex) {
280            LOG.warn("Error starting action [{0}]. ErrorType [{1}], ErrorCode [{2}], Message [{3}]",
281                    wfAction.getName(), ex.getErrorType(), ex.getErrorCode(), ex.getMessage(), ex);
282            wfAction.setErrorInfo(ex.getErrorCode(), ex.getMessage());
283            switch (ex.getErrorType()) {
284                case TRANSIENT:
285                    if (!handleTransient(context, executor, WorkflowAction.Status.START_RETRY)) {
286                        handleNonTransient(context, executor, WorkflowAction.Status.START_MANUAL);
287                        wfAction.setPendingAge(new Date());
288                        wfAction.setRetries(0);
289                        wfAction.setStartTime(null);
290                    }
291                    break;
292                case NON_TRANSIENT:
293                    handleNonTransient(context, executor, WorkflowAction.Status.START_MANUAL);
294                    break;
295                case ERROR:
296                    handleError(context, executor, WorkflowAction.Status.ERROR.toString(), true,
297                            WorkflowAction.Status.DONE);
298                    break;
299                case FAILED:
300                    try {
301                        failJob(context);
302                        endWF();
303                        SLAEventBean slaEvent1 = SLADbXOperations.createStatusEvent(wfAction.getSlaXml(), wfAction.getId(), Status.FAILED,
304                                SlaAppType.WORKFLOW_ACTION);
305                        if(slaEvent1 != null) {
306                            insertList.add(slaEvent1);
307                        }
308                    }
309                    catch (XException x) {
310                        LOG.warn("ActionStartXCommand - case:FAILED ", x.getMessage());
311                    }
312                    break;
313            }
314            updateList.add(new UpdateEntry<WorkflowActionQuery>(WorkflowActionQuery.UPDATE_ACTION_START, wfAction));
315            updateJobLastModified();
316        }
317        finally {
318            try {
319                BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(insertList, updateList, null);
320                if (!(executor instanceof ControlNodeActionExecutor) && EventHandlerService.isEnabled()) {
321                    generateEvent(wfAction, wfJob.getUser());
322                }
323                if (execSynchronous) {
324                    // Changing to synchronous call from asynchronous queuing to prevent
325                    // undue delay from ::start:: to action due to queuing
326                    callActionEnd();
327                }
328            }
329            catch (JPAExecutorException e) {
330                throw new CommandException(e);
331            }
332        }
333
334        LOG.debug("ENDED ActionStartXCommand for wf actionId=" + actionId + ", jobId=" + jobId);
335
336        return null;
337    }
338
339    protected void callActionEnd() throws CommandException {
340        new ActionEndXCommand(wfAction.getId(), wfAction.getType()).call(getEntityKey());
341    }
342
343    /**
344     * Get action executor context
345     * @param isRetry
346     * @param isUserRetry
347     * @return
348     */
349    protected ActionExecutorContext getContext(boolean isRetry, boolean isUserRetry) {
350        return new ActionXCommand.ActionExecutorContext(wfJob, wfAction, isRetry, isUserRetry);
351    }
352
353    protected void updateJobLastModified(){
354        wfJob.setLastModifiedTime(new Date());
355        updateList.add(new UpdateEntry<WorkflowJobQuery>(WorkflowJobQuery.UPDATE_WORKFLOW_STATUS_INSTANCE_MODIFIED, wfJob));
356    }
357
358    protected void endWF() throws CommandException{
359        updateParentIfNecessary(wfJob, 3);
360        new WfEndXCommand(wfJob).call(); // To delete the WF temp dir
361        SLAEventBean slaEvent2 = SLADbXOperations.createStatusEvent(wfJob.getSlaXml(), wfJob.getId(), Status.FAILED,
362                SlaAppType.WORKFLOW_JOB);
363        if(slaEvent2 != null) {
364            insertList.add(slaEvent2);
365        }
366    }
367
368    protected void handleError(ActionExecutorContext context, WorkflowJobBean workflow, WorkflowActionBean action)
369            throws CommandException {
370        failJob(context);
371        updateList.add(new UpdateEntry<WorkflowActionQuery>(WorkflowActionQuery.UPDATE_ACTION_START, wfAction));
372        updateJobLastModified();
373        SLAEventBean slaEvent1 = SLADbXOperations.createStatusEvent(action.getSlaXml(), action.getId(),
374                Status.FAILED, SlaAppType.WORKFLOW_ACTION);
375        if(slaEvent1 != null) {
376            insertList.add(slaEvent1);
377        }
378        endWF();
379        return;
380    }
381
382    /* (non-Javadoc)
383     * @see org.apache.oozie.command.XCommand#getKey()
384     */
385    @Override
386    public String getKey(){
387        return getName() + "_" + actionId;
388    }
389
390    private void prepareForRetry(WorkflowActionBean wfAction) {
391        if (wfAction.getType().equals("map-reduce")) {
392            // need to delete child job id of original run
393            wfAction.setExternalChildIDs("");
394        }
395    }
396
397    @Override
398    protected void queueCommandForTransientFailure(long retryDelayMillis){
399        queue(new ActionStartXCommand(wfAction.getId(), wfAction.getType()), retryDelayMillis);
400    }
401
402}