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.io.StringReader;
023import java.net.URI;
024import java.net.URISyntaxException;
025import java.util.Date;
026import java.util.HashMap;
027import java.util.Map;
028import java.util.Properties;
029import java.util.Set;
030
031import org.apache.hadoop.conf.Configuration;
032import org.apache.hadoop.fs.FileSystem;
033import org.apache.hadoop.fs.Path;
034import org.apache.oozie.DagELFunctions;
035import org.apache.oozie.ErrorCode;
036import org.apache.oozie.WorkflowActionBean;
037import org.apache.oozie.WorkflowJobBean;
038import org.apache.oozie.action.ActionExecutor;
039import org.apache.oozie.client.Job;
040import org.apache.oozie.client.WorkflowAction;
041import org.apache.oozie.client.WorkflowJob;
042import org.apache.oozie.command.CommandException;
043import org.apache.oozie.service.CallbackService;
044import org.apache.oozie.service.ELService;
045import org.apache.oozie.service.HadoopAccessorException;
046import org.apache.oozie.service.HadoopAccessorService;
047import org.apache.oozie.service.JPAService;
048import org.apache.oozie.service.LiteWorkflowStoreService;
049import org.apache.oozie.service.Services;
050import org.apache.oozie.util.ELEvaluator;
051import org.apache.oozie.util.InstrumentUtils;
052import org.apache.oozie.util.Instrumentation;
053import org.apache.oozie.util.XConfiguration;
054import org.apache.oozie.workflow.WorkflowException;
055import org.apache.oozie.workflow.WorkflowInstance;
056import org.apache.oozie.workflow.lite.LiteWorkflowInstance;
057
058/**
059 * Base class for Action execution commands. Provides common functionality to handle different types of errors while
060 * attempting to start or end an action.
061 */
062public abstract class ActionXCommand<T> extends WorkflowXCommand<T> {
063    private static final String INSTRUMENTATION_GROUP = "action.executors";
064
065    protected static final String RECOVERY_ID_SEPARATOR = "@";
066
067    public ActionXCommand(String name, String type, int priority) {
068        super(name, type, priority);
069    }
070
071    /**
072     * Takes care of Transient failures. Sets the action status to retry and increments the retry count if not enough
073     * attempts have been made. Otherwise returns false.
074     *
075     * @param context the execution context.
076     * @param executor the executor instance being used.
077     * @param status the status to be set for the action.
078     * @return true if the action is scheduled for another retry. false if the number of retries has exceeded the
079     *         maximum number of configured retries.
080     * @throws CommandException thrown if unable to handle transient
081     */
082    protected boolean handleTransient(ActionExecutor.Context context, ActionExecutor executor,
083            WorkflowAction.Status status) throws CommandException {
084        LOG.debug("Attempting to retry");
085        ActionExecutorContext aContext = (ActionExecutorContext) context;
086        WorkflowActionBean action = (WorkflowActionBean) aContext.getAction();
087        incrActionErrorCounter(action.getType(), "transient", 1);
088
089        int actionRetryCount = action.getRetries();
090        if (actionRetryCount >= executor.getMaxRetries()) {
091            LOG.warn("Exceeded max retry count [{0}]. Suspending Job", executor.getMaxRetries());
092            return false;
093        }
094        else {
095            action.setStatus(status);
096            action.setPending();
097            action.incRetries();
098            long retryDelayMillis = getRetryDelay(actionRetryCount, executor.getRetryInterval(), executor.getRetryPolicy());
099            action.setPendingAge(new Date(System.currentTimeMillis() + retryDelayMillis));
100            LOG.info("Next Retry, Attempt Number [{0}] in [{1}] milliseconds", actionRetryCount + 1, retryDelayMillis);
101            this.resetUsed();
102            queueCommandForTransientFailure(retryDelayMillis);
103            return true;
104        }
105    }
106
107    protected void queueCommandForTransientFailure(long retryDelayMillis){
108        queue(this, retryDelayMillis);
109    }
110    /**
111     * Takes care of non transient failures. The job is suspended, and the state of the action is changed to *MANUAL and
112     * set pending flag of action to false
113     *
114     * @param context the execution context.
115     * @param executor the executor instance being used.
116     * @param status the status to be set for the action.
117     * @throws CommandException thrown if unable to suspend job
118     */
119    protected void handleNonTransient(ActionExecutor.Context context, ActionExecutor executor,
120            WorkflowAction.Status status) throws CommandException {
121        ActionExecutorContext aContext = (ActionExecutorContext) context;
122        WorkflowActionBean action = (WorkflowActionBean) aContext.getAction();
123        incrActionErrorCounter(action.getType(), "nontransient", 1);
124        WorkflowJobBean workflow = (WorkflowJobBean) context.getWorkflow();
125        String id = workflow.getId();
126        action.setStatus(status);
127        action.resetPendingOnly();
128        LOG.warn("Suspending Workflow Job id=" + id);
129        try {
130            SuspendXCommand.suspendJob(Services.get().get(JPAService.class), workflow, id, action.getId(), null);
131        }
132        catch (Exception e) {
133            throw new CommandException(ErrorCode.E0727, id, e.getMessage());
134        }
135        finally {
136            updateParentIfNecessary(workflow, 3);
137        }
138    }
139
140    /**
141     * Takes care of errors. </p> For errors while attempting to start the action, the job state is updated and an
142     * {@link ActionEndCommand} is queued. </p> For errors while attempting to end the action, the job state is updated.
143     * </p>
144     *
145     * @param context the execution context.
146     * @param executor the executor instance being used.
147     * @param message
148     * @param isStart whether the error was generated while starting or ending an action.
149     * @param status the status to be set for the action.
150     * @throws CommandException thrown if unable to handle action error
151     */
152    protected void handleError(ActionExecutor.Context context, ActionExecutor executor, String message,
153            boolean isStart, WorkflowAction.Status status) throws CommandException {
154        LOG.warn("Setting Action Status to [{0}]", status);
155        ActionExecutorContext aContext = (ActionExecutorContext) context;
156        WorkflowActionBean action = (WorkflowActionBean) aContext.getAction();
157
158        if (!handleUserRetry(action)) {
159            incrActionErrorCounter(action.getType(), "error", 1);
160            action.setPending();
161            if (isStart) {
162                action.setExecutionData(message, null);
163                queue(new ActionEndXCommand(action.getId(), action.getType()));
164            }
165            else {
166                action.setEndData(status, WorkflowAction.Status.ERROR.toString());
167            }
168        }
169    }
170
171    /**
172     * Fail the job due to failed action
173     *
174     * @param context the execution context.
175     * @throws CommandException thrown if unable to fail job
176     */
177    public void failJob(ActionExecutor.Context context) throws CommandException {
178        ActionExecutorContext aContext = (ActionExecutorContext) context;
179        WorkflowActionBean action = (WorkflowActionBean) aContext.getAction();
180        failJob(context, action);
181    }
182
183    /**
184     * Fail the job due to failed action
185     *
186     * @param context the execution context.
187     * @param action the action that caused the workflow to fail
188     * @throws CommandException thrown if unable to fail job
189     */
190    public void failJob(ActionExecutor.Context context, WorkflowActionBean action) throws CommandException {
191        WorkflowJobBean workflow = (WorkflowJobBean) context.getWorkflow();
192        if (!handleUserRetry(action)) {
193            incrActionErrorCounter(action.getType(), "failed", 1);
194            LOG.warn("Failing Job due to failed action [{0}]", action.getName());
195            try {
196                workflow.getWorkflowInstance().fail(action.getName());
197                WorkflowInstance wfInstance = workflow.getWorkflowInstance();
198                ((LiteWorkflowInstance) wfInstance).setStatus(WorkflowInstance.Status.FAILED);
199                workflow.setWorkflowInstance(wfInstance);
200                workflow.setStatus(WorkflowJob.Status.FAILED);
201                action.setStatus(WorkflowAction.Status.FAILED);
202                action.resetPending();
203                queue(new WorkflowNotificationXCommand(workflow, action));
204                queue(new KillXCommand(workflow.getId()));
205                InstrumentUtils.incrJobCounter(INSTR_FAILED_JOBS_COUNTER_NAME, 1, getInstrumentation());
206            }
207            catch (WorkflowException ex) {
208                throw new CommandException(ex);
209            }
210        }
211    }
212
213    /**
214     * Execute retry for action if this action is eligible for user-retry
215     *
216     * @param context the execution context.
217     * @return true if user-retry has to be handled for this action
218     * @throws CommandException thrown if unable to fail job
219     */
220    public boolean handleUserRetry(WorkflowActionBean action) throws CommandException {
221        String errorCode = action.getErrorCode();
222        Set<String> allowedRetryCode = LiteWorkflowStoreService.getUserRetryErrorCode();
223
224        if ((allowedRetryCode.contains(LiteWorkflowStoreService.USER_ERROR_CODE_ALL) || allowedRetryCode.contains(errorCode))
225                && action.getUserRetryCount() < action.getUserRetryMax()) {
226            LOG.info("Preparing retry this action [{0}], errorCode [{1}], userRetryCount [{2}], "
227                    + "userRetryMax [{3}], userRetryInterval [{4}]", action.getId(), errorCode, action
228                    .getUserRetryCount(), action.getUserRetryMax(), action.getUserRetryInterval());
229            int interval = action.getUserRetryInterval() * 60 * 1000;
230            action.setStatus(WorkflowAction.Status.USER_RETRY);
231            action.incrmentUserRetryCount();
232            action.setPending();
233            queue(new ActionStartXCommand(action.getId(), action.getType()), interval);
234            return true;
235        }
236        return false;
237    }
238
239        /*
240         * In case of action error increment the error count for instrumentation
241         */
242    private void incrActionErrorCounter(String type, String error, int count) {
243        getInstrumentation().incr(INSTRUMENTATION_GROUP, type + "#ex." + error, count);
244    }
245
246        /**
247         * Increment the action counter in the instrumentation log. indicating how
248         * many times the action was executed since the start Oozie server
249         */
250    protected void incrActionCounter(String type, int count) {
251        getInstrumentation().incr(INSTRUMENTATION_GROUP, type + "#" + getName(), count);
252    }
253
254        /**
255         * Adding a cron for the instrumentation time for the given Instrumentation
256         * group
257         */
258    protected void addActionCron(String type, Instrumentation.Cron cron) {
259        getInstrumentation().addCron(INSTRUMENTATION_GROUP, type + "#" + getName(), cron);
260    }
261
262    /*
263     * Returns the next retry time in milliseconds, based on retry policy algorithm.
264     */
265    private long getRetryDelay(int retryCount, long retryInterval, ActionExecutor.RETRYPOLICY retryPolicy) {
266        switch (retryPolicy) {
267            case EXPONENTIAL:
268                long retryTime = ((long) Math.pow(2, retryCount) * retryInterval * 1000L);
269                return retryTime;
270            case PERIODIC:
271                return retryInterval * 1000L;
272            default:
273                throw new UnsupportedOperationException("Retry policy not supported");
274        }
275    }
276
277    /**
278     * Workflow action executor context
279     *
280     */
281    public static class ActionExecutorContext implements ActionExecutor.Context {
282        protected final WorkflowJobBean workflow;
283        private Configuration protoConf;
284        protected final WorkflowActionBean action;
285        private final boolean isRetry;
286        private final boolean isUserRetry;
287        private boolean started;
288        private boolean ended;
289        private boolean executed;
290        private boolean shouldEndWF;
291        private Job.Status jobStatus;
292
293        /**
294                 * Constructing the ActionExecutorContext, setting the private members
295                 * and constructing the proto configuration
296                 */
297        public ActionExecutorContext(WorkflowJobBean workflow, WorkflowActionBean action, boolean isRetry, boolean isUserRetry) {
298            this.workflow = workflow;
299            this.action = action;
300            this.isRetry = isRetry;
301            this.isUserRetry = isUserRetry;
302            try {
303                protoConf = new XConfiguration(new StringReader(workflow.getProtoActionConf()));
304            }
305            catch (IOException ex) {
306                throw new RuntimeException("It should not happen", ex);
307            }
308        }
309
310        /*
311         * (non-Javadoc)
312         * @see org.apache.oozie.action.ActionExecutor.Context#getCallbackUrl(java.lang.String)
313         */
314        public String getCallbackUrl(String externalStatusVar) {
315            return Services.get().get(CallbackService.class).createCallBackUrl(action.getId(), externalStatusVar);
316        }
317
318        /*
319         * (non-Javadoc)
320         * @see org.apache.oozie.action.ActionExecutor.Context#getProtoActionConf()
321         */
322        public Configuration getProtoActionConf() {
323            return protoConf;
324        }
325
326        /*
327         * (non-Javadoc)
328         * @see org.apache.oozie.action.ActionExecutor.Context#getWorkflow()
329         */
330        public WorkflowJob getWorkflow() {
331            return workflow;
332        }
333
334        /**
335         * Returns the workflow action of the given action context
336         *
337         * @return the workflow action of the given action context
338         */
339        public WorkflowAction getAction() {
340            return action;
341        }
342
343        /*
344         * (non-Javadoc)
345         * @see org.apache.oozie.action.ActionExecutor.Context#getELEvaluator()
346         */
347        public ELEvaluator getELEvaluator() {
348            ELEvaluator evaluator = Services.get().get(ELService.class).createEvaluator("workflow");
349            DagELFunctions.configureEvaluator(evaluator, workflow, action);
350            return evaluator;
351        }
352
353        /*
354         * (non-Javadoc)
355         * @see org.apache.oozie.action.ActionExecutor.Context#setVar(java.lang.String, java.lang.String)
356         */
357        public void setVar(String name, String value) {
358            setVarToWorkflow(name, value);
359        }
360
361        /**
362         * This is not thread safe, don't use if workflowjob is shared among multiple actions command
363         * @param name
364         * @param value
365         */
366        public void setVarToWorkflow(String name, String value) {
367            name = action.getName() + WorkflowInstance.NODE_VAR_SEPARATOR + name;
368            WorkflowInstance wfInstance = workflow.getWorkflowInstance();
369            wfInstance.setVar(name, value);
370            workflow.setWorkflowInstance(wfInstance);
371        }
372
373        /*
374         * (non-Javadoc)
375         * @see org.apache.oozie.action.ActionExecutor.Context#getVar(java.lang.String)
376         */
377        public String getVar(String name) {
378            name = action.getName() + WorkflowInstance.NODE_VAR_SEPARATOR + name;
379            return workflow.getWorkflowInstance().getVar(name);
380        }
381
382        /*
383         * (non-Javadoc)
384         * @see org.apache.oozie.action.ActionExecutor.Context#setStartData(java.lang.String, java.lang.String, java.lang.String)
385         */
386        public void setStartData(String externalId, String trackerUri, String consoleUrl) {
387            action.setStartData(externalId, trackerUri, consoleUrl);
388            started = true;
389        }
390
391        /**
392         * Setting the start time of the action
393         */
394        public void setStartTime() {
395            Date now = new Date();
396            action.setStartTime(now);
397        }
398
399        /*
400         * (non-Javadoc)
401         * @see org.apache.oozie.action.ActionExecutor.Context#setExecutionData(java.lang.String, java.util.Properties)
402         */
403        public void setExecutionData(String externalStatus, Properties actionData) {
404            action.setExecutionData(externalStatus, actionData);
405            executed = true;
406        }
407
408        /*
409         * (non-Javadoc)
410         * @see org.apache.oozie.action.ActionExecutor.Context#setExecutionStats(java.lang.String)
411         */
412        public void setExecutionStats(String jsonStats) {
413            action.setExecutionStats(jsonStats);
414            executed = true;
415        }
416
417        /*
418         * (non-Javadoc)
419         * @see org.apache.oozie.action.ActionExecutor.Context#setExternalChildIDs(java.lang.String)
420         */
421        public void setExternalChildIDs(String externalChildIDs) {
422            action.setExternalChildIDs(externalChildIDs);
423            executed = true;
424        }
425
426        /*
427         * (non-Javadoc)
428         * @see org.apache.oozie.action.ActionExecutor.Context#setEndData(org.apache.oozie.client.WorkflowAction.Status, java.lang.String)
429         */
430        public void setEndData(WorkflowAction.Status status, String signalValue) {
431            action.setEndData(status, signalValue);
432            ended = true;
433        }
434
435        /*
436         * (non-Javadoc)
437         * @see org.apache.oozie.action.ActionExecutor.Context#isRetry()
438         */
439        public boolean isRetry() {
440            return isRetry;
441        }
442
443        /**
444         * Return if the executor invocation is a user retry or not.
445         *
446         * @return if the executor invocation is a user retry or not.
447         */
448        public boolean isUserRetry() {
449            return isUserRetry;
450        }
451
452        /**
453         * Returns whether setStartData has been called or not.
454         *
455         * @return true if start completion info has been set.
456         */
457        public boolean isStarted() {
458            return started;
459        }
460
461        /**
462         * Returns whether setExecutionData has been called or not.
463         *
464         * @return true if execution completion info has been set, otherwise false.
465         */
466        public boolean isExecuted() {
467            return executed;
468        }
469
470        /**
471         * Returns whether setEndData has been called or not.
472         *
473         * @return true if end completion info has been set.
474         */
475        public boolean isEnded() {
476            return ended;
477        }
478
479        public void setExternalStatus(String externalStatus) {
480            action.setExternalStatus(externalStatus);
481        }
482
483        @Override
484        public String getRecoveryId() {
485            return action.getId() + RECOVERY_ID_SEPARATOR + workflow.getRun();
486        }
487
488        /* (non-Javadoc)
489         * @see org.apache.oozie.action.ActionExecutor.Context#getActionDir()
490         */
491        public Path getActionDir() throws HadoopAccessorException, IOException, URISyntaxException {
492            String name = getWorkflow().getId() + "/" + action.getName() + "--" + action.getType();
493            FileSystem fs = getAppFileSystem();
494            String actionDirPath = Services.get().getSystemId() + "/" + name;
495            Path fqActionDir = new Path(fs.getHomeDirectory(), actionDirPath);
496            return fqActionDir;
497        }
498
499        /* (non-Javadoc)
500         * @see org.apache.oozie.action.ActionExecutor.Context#getAppFileSystem()
501         */
502        public FileSystem getAppFileSystem() throws HadoopAccessorException, IOException, URISyntaxException {
503            WorkflowJob workflow = getWorkflow();
504            URI uri = new URI(getWorkflow().getAppPath());
505            HadoopAccessorService has = Services.get().get(HadoopAccessorService.class);
506            Configuration fsConf = has.createJobConf(uri.getAuthority());
507            return has.createFileSystem(workflow.getUser(), uri, fsConf);
508
509        }
510
511        /* (non-Javadoc)
512         * @see org.apache.oozie.action.ActionExecutor.Context#setErrorInfo(java.lang.String, java.lang.String)
513         */
514        @Override
515        public void setErrorInfo(String str, String exMsg) {
516            action.setErrorInfo(str, exMsg);
517        }
518
519        public boolean isShouldEndWF() {
520            return shouldEndWF;
521        }
522
523        public void setShouldEndWF(boolean shouldEndWF) {
524            this.shouldEndWF = shouldEndWF;
525        }
526
527        public Job.Status getJobStatus() {
528            return jobStatus;
529        }
530
531        public void setJobStatus(Job.Status jobStatus) {
532            this.jobStatus = jobStatus;
533        }
534    }
535
536    public static class ForkedActionExecutorContext extends ActionExecutorContext {
537        private Map<String, String> contextVariableMap = new HashMap<String, String>();
538
539        public ForkedActionExecutorContext(WorkflowJobBean workflow, WorkflowActionBean action, boolean isRetry,
540                boolean isUserRetry) {
541            super(workflow, action, isRetry, isUserRetry);
542        }
543
544        public void setVar(String name, String value) {
545            if (value != null) {
546                contextVariableMap.remove(name);
547            }
548            else {
549                contextVariableMap.put(name, value);
550            }
551        }
552
553        public String getVar(String name) {
554            return contextVariableMap.get(name);
555        }
556
557        public Map<String, String> getContextMap() {
558            return contextVariableMap;
559        }
560    }
561
562}