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 org.apache.hadoop.conf.Configuration;
022import org.apache.hadoop.fs.Path;
023import org.apache.hadoop.fs.FileSystem;
024import org.apache.oozie.AppType;
025import org.apache.oozie.SLAEventBean;
026import org.apache.oozie.command.wf.ActionXCommand.ActionExecutorContext;
027import org.apache.oozie.WorkflowJobBean;
028import org.apache.oozie.ErrorCode;
029import org.apache.oozie.action.oozie.SubWorkflowActionExecutor;
030import org.apache.oozie.service.HadoopAccessorException;
031import org.apache.oozie.service.JPAService;
032import org.apache.oozie.service.UUIDService;
033import org.apache.oozie.service.WorkflowStoreService;
034import org.apache.oozie.service.WorkflowAppService;
035import org.apache.oozie.service.HadoopAccessorService;
036import org.apache.oozie.service.Services;
037import org.apache.oozie.service.DagXLogInfoService;
038import org.apache.oozie.util.ELUtils;
039import org.apache.oozie.util.LogUtils;
040import org.apache.oozie.sla.SLAOperations;
041import org.apache.oozie.util.XLog;
042import org.apache.oozie.util.ParamChecker;
043import org.apache.oozie.util.XConfiguration;
044import org.apache.oozie.util.XmlUtils;
045import org.apache.oozie.command.CommandException;
046import org.apache.oozie.executor.jpa.BatchQueryExecutor;
047import org.apache.oozie.executor.jpa.JPAExecutorException;
048import org.apache.oozie.service.ELService;
049import org.apache.oozie.store.StoreException;
050import org.apache.oozie.workflow.WorkflowApp;
051import org.apache.oozie.workflow.WorkflowException;
052import org.apache.oozie.workflow.WorkflowInstance;
053import org.apache.oozie.workflow.WorkflowLib;
054import org.apache.oozie.util.ELEvaluator;
055import org.apache.oozie.util.InstrumentUtils;
056import org.apache.oozie.util.PropertiesUtils;
057import org.apache.oozie.util.db.SLADbOperations;
058import org.apache.oozie.service.SchemaService.SchemaName;
059import org.apache.oozie.client.OozieClient;
060import org.apache.oozie.client.WorkflowJob;
061import org.apache.oozie.client.SLAEvent.SlaAppType;
062import org.apache.oozie.client.rest.JsonBean;
063import org.jdom.Element;
064import org.jdom.filter.ElementFilter;
065
066import java.util.ArrayList;
067import java.util.Date;
068import java.util.Iterator;
069import java.util.List;
070import java.util.Map;
071import java.util.Set;
072import java.util.HashSet;
073import java.io.IOException;
074import java.net.URI;
075
076@SuppressWarnings("deprecation")
077public class SubmitXCommand extends WorkflowXCommand<String> {
078    public static final String CONFIG_DEFAULT = "config-default.xml";
079
080    private Configuration conf;
081    private List<JsonBean> insertList = new ArrayList<JsonBean>();
082    private String parentId;
083
084    /**
085     * Constructor to create the workflow Submit Command.
086     *
087     * @param conf : Configuration for workflow job
088     */
089    public SubmitXCommand(Configuration conf) {
090        super("submit", "submit", 1);
091        this.conf = ParamChecker.notNull(conf, "conf");
092    }
093
094    /**
095     * Constructor for submitting wf through coordinator
096     *
097     * @param conf : Configuration for workflow job
098     * @param parentId: the coord action id
099     */
100    public SubmitXCommand(Configuration conf, String parentId) {
101        this(conf);
102        this.parentId = parentId;
103    }
104
105    /**
106     * Constructor to create the workflow Submit Command.
107     *
108     * @param dryrun : if dryrun
109     * @param conf : Configuration for workflow job
110     * @param authToken : To be used for authentication
111     */
112    public SubmitXCommand(boolean dryrun, Configuration conf) {
113        this(conf);
114        this.dryrun = dryrun;
115    }
116
117    private static final Set<String> DISALLOWED_USER_PROPERTIES = new HashSet<String>();
118
119    static {
120        String[] badUserProps = {PropertiesUtils.DAYS, PropertiesUtils.HOURS, PropertiesUtils.MINUTES,
121                PropertiesUtils.KB, PropertiesUtils.MB, PropertiesUtils.GB, PropertiesUtils.TB, PropertiesUtils.PB,
122                PropertiesUtils.RECORDS, PropertiesUtils.MAP_IN, PropertiesUtils.MAP_OUT, PropertiesUtils.REDUCE_IN,
123                PropertiesUtils.REDUCE_OUT, PropertiesUtils.GROUPS};
124        PropertiesUtils.createPropertySet(badUserProps, DISALLOWED_USER_PROPERTIES);
125    }
126
127    @Override
128    protected String execute() throws CommandException {
129        InstrumentUtils.incrJobCounter(getName(), 1, getInstrumentation());
130        WorkflowAppService wps = Services.get().get(WorkflowAppService.class);
131        try {
132            XLog.Info.get().setParameter(DagXLogInfoService.TOKEN, conf.get(OozieClient.LOG_TOKEN));
133            String user = conf.get(OozieClient.USER_NAME);
134            URI uri = new URI(conf.get(OozieClient.APP_PATH));
135            HadoopAccessorService has = Services.get().get(HadoopAccessorService.class);
136            Configuration fsConf = has.createJobConf(uri.getAuthority());
137            FileSystem fs = has.createFileSystem(user, uri, fsConf);
138
139            Path configDefault = null;
140            Configuration defaultConf = null;
141            // app path could be a directory
142            Path path = new Path(uri.getPath());
143            if (!fs.isFile(path)) {
144                configDefault = new Path(path, CONFIG_DEFAULT);
145            } else {
146                configDefault = new Path(path.getParent(), CONFIG_DEFAULT);
147            }
148
149            if (fs.exists(configDefault)) {
150                try {
151                    defaultConf = new XConfiguration(fs.open(configDefault));
152                    PropertiesUtils.checkDisallowedProperties(defaultConf, DISALLOWED_USER_PROPERTIES);
153                    PropertiesUtils.checkDefaultDisallowedProperties(defaultConf);
154                    XConfiguration.injectDefaults(defaultConf, conf);
155                }
156                catch (IOException ex) {
157                    throw new IOException("default configuration file, " + ex.getMessage(), ex);
158                }
159            }
160
161            WorkflowApp app = wps.parseDef(conf, defaultConf);
162            XConfiguration protoActionConf = wps.createProtoActionConf(conf, true);
163            WorkflowLib workflowLib = Services.get().get(WorkflowStoreService.class).getWorkflowLibWithNoDB();
164
165            PropertiesUtils.checkDisallowedProperties(conf, DISALLOWED_USER_PROPERTIES);
166
167            // Resolving all variables in the job properties.
168            // This ensures the Hadoop Configuration semantics is preserved.
169            XConfiguration resolvedVarsConf = new XConfiguration();
170            for (Map.Entry<String, String> entry : conf) {
171                resolvedVarsConf.set(entry.getKey(), conf.get(entry.getKey()));
172            }
173            conf = resolvedVarsConf;
174
175            WorkflowInstance wfInstance;
176            try {
177                wfInstance = workflowLib.createInstance(app, conf);
178            }
179            catch (WorkflowException e) {
180                throw new StoreException(e);
181            }
182
183            Configuration conf = wfInstance.getConf();
184            // System.out.println("WF INSTANCE CONF:");
185            // System.out.println(XmlUtils.prettyPrint(conf).toString());
186
187            WorkflowJobBean workflow = new WorkflowJobBean();
188            workflow.setId(wfInstance.getId());
189            workflow.setAppName(ELUtils.resolveAppName(app.getName(), conf));
190            workflow.setAppPath(conf.get(OozieClient.APP_PATH));
191            workflow.setConf(XmlUtils.prettyPrint(conf).toString());
192            workflow.setProtoActionConf(protoActionConf.toXmlString());
193            workflow.setCreatedTime(new Date());
194            workflow.setLastModifiedTime(new Date());
195            workflow.setLogToken(conf.get(OozieClient.LOG_TOKEN, ""));
196            workflow.setStatus(WorkflowJob.Status.PREP);
197            workflow.setRun(0);
198            workflow.setUser(conf.get(OozieClient.USER_NAME));
199            workflow.setGroup(conf.get(OozieClient.GROUP_NAME));
200            workflow.setWorkflowInstance(wfInstance);
201            workflow.setExternalId(conf.get(OozieClient.EXTERNAL_ID));
202            // Set parent id if it doesn't already have one (for subworkflows)
203            if (workflow.getParentId() == null) {
204                workflow.setParentId(conf.get(SubWorkflowActionExecutor.PARENT_ID));
205            }
206            // Set to coord action Id if workflow submitted through coordinator
207            if (workflow.getParentId() == null) {
208                workflow.setParentId(parentId);
209            }
210
211            LogUtils.setLogInfo(workflow);
212            LOG.debug("Workflow record created, Status [{0}]", workflow.getStatus());
213            Element wfElem = XmlUtils.parseXml(app.getDefinition());
214            ELEvaluator evalSla = createELEvaluatorForGroup(conf, "wf-sla-submit");
215            String jobSlaXml = verifySlaElements(wfElem, evalSla);
216            if (!dryrun) {
217                writeSLARegistration(wfElem, jobSlaXml, workflow.getId(), workflow.getParentId(), workflow.getUser(),
218                        workflow.getGroup(), workflow.getAppName(), LOG, evalSla);
219                workflow.setSlaXml(jobSlaXml);
220                // System.out.println("SlaXml :"+ slaXml);
221
222                //store.insertWorkflow(workflow);
223                insertList.add(workflow);
224                JPAService jpaService = Services.get().get(JPAService.class);
225                if (jpaService != null) {
226                    try {
227                        BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(insertList, null, null);
228                    }
229                    catch (JPAExecutorException je) {
230                        throw new CommandException(je);
231                    }
232                }
233                else {
234                    LOG.error(ErrorCode.E0610);
235                    return null;
236                }
237
238                return workflow.getId();
239            }
240            else {
241                // Checking variable substitution for dryrun
242                ActionExecutorContext context = new ActionXCommand.ActionExecutorContext(workflow, null, false, false);
243                Element workflowXml = XmlUtils.parseXml(app.getDefinition());
244                removeSlaElements(workflowXml);
245                String workflowXmlString = XmlUtils.removeComments(XmlUtils.prettyPrint(workflowXml).toString());
246                workflowXmlString = context.getELEvaluator().evaluate(workflowXmlString, String.class);
247                workflowXml = XmlUtils.parseXml(workflowXmlString);
248
249                Iterator<Element> it = workflowXml.getDescendants(new ElementFilter("job-xml"));
250
251                // Checking all variable substitutions in job-xml files
252                while (it.hasNext()) {
253                    Element e = it.next();
254                    String jobXml = e.getTextTrim();
255                    Path xmlPath = new Path(workflow.getAppPath(), jobXml);
256                    Configuration jobXmlConf = new XConfiguration(fs.open(xmlPath));
257
258
259                    String jobXmlConfString = XmlUtils.prettyPrint(jobXmlConf).toString();
260                    jobXmlConfString = XmlUtils.removeComments(jobXmlConfString);
261                    context.getELEvaluator().evaluate(jobXmlConfString, String.class);
262                }
263
264                return "OK";
265            }
266        }
267        catch (WorkflowException ex) {
268            throw new CommandException(ex);
269        }
270        catch (HadoopAccessorException ex) {
271            throw new CommandException(ex);
272        }
273        catch (Exception ex) {
274            throw new CommandException(ErrorCode.E0803, ex.getMessage(), ex);
275        }
276    }
277
278    private void removeSlaElements(Element eWfJob) {
279        Element sla = XmlUtils.getSLAElement(eWfJob);
280        if (sla != null) {
281            eWfJob.removeChildren(sla.getName(), sla.getNamespace());
282        }
283
284        for (Element action : (List<Element>) eWfJob.getChildren("action", eWfJob.getNamespace())) {
285            sla = XmlUtils.getSLAElement(action);
286            if (sla != null) {
287                action.removeChildren(sla.getName(), sla.getNamespace());
288            }
289        }
290    }
291    private String verifySlaElements(Element eWfJob, ELEvaluator evalSla) throws CommandException {
292        String jobSlaXml = "";
293        // Validate WF job
294        Element eSla = XmlUtils.getSLAElement(eWfJob);
295        if (eSla != null) {
296            jobSlaXml = resolveSla(eSla, evalSla);
297        }
298
299        // Validate all actions
300        for (Element action : (List<Element>) eWfJob.getChildren("action", eWfJob.getNamespace())) {
301            eSla = XmlUtils.getSLAElement(action);
302            if (eSla != null) {
303                resolveSla(eSla, evalSla);
304            }
305        }
306        return jobSlaXml;
307    }
308
309    private void writeSLARegistration(Element eWfJob, String slaXml, String jobId, String parentId, String user,
310            String group, String appName, XLog log, ELEvaluator evalSla) throws CommandException {
311        try {
312            if (slaXml != null && slaXml.length() > 0) {
313                Element eSla = XmlUtils.parseXml(slaXml);
314                SLAEventBean slaEvent = SLADbOperations.createSlaRegistrationEvent(eSla, jobId,
315                        SlaAppType.WORKFLOW_JOB, user, group, log);
316                if(slaEvent != null) {
317                    insertList.add(slaEvent);
318                }
319                // insert into new table
320                SLAOperations.createSlaRegistrationEvent(eSla, jobId, parentId, AppType.WORKFLOW_JOB, user, appName,
321                        log, false);
322            }
323            // Add sla for wf actions
324            for (Element action : (List<Element>) eWfJob.getChildren("action", eWfJob.getNamespace())) {
325                Element actionSla = XmlUtils.getSLAElement(action);
326                if (actionSla != null) {
327                    String actionSlaXml = SubmitXCommand.resolveSla(actionSla, evalSla);
328                    actionSla = XmlUtils.parseXml(actionSlaXml);
329                    String actionId = Services.get().get(UUIDService.class)
330                            .generateChildId(jobId, action.getAttributeValue("name") + "");
331                    SLAOperations.createSlaRegistrationEvent(actionSla, actionId, jobId, AppType.WORKFLOW_ACTION,
332                            user, appName, log, false);
333                }
334            }
335        }
336        catch (Exception e) {
337            e.printStackTrace();
338            throw new CommandException(ErrorCode.E1007, "workflow " + jobId, e.getMessage(), e);
339        }
340    }
341
342    /**
343     * Resolve variables in sla xml element.
344     *
345     * @param eSla sla xml element
346     * @param evalSla sla evaluator
347     * @return sla xml string after evaluation
348     * @throws CommandException
349     */
350    public static String resolveSla(Element eSla, ELEvaluator evalSla) throws CommandException {
351        // EL evaluation
352        String slaXml = XmlUtils.prettyPrint(eSla).toString();
353        try {
354            slaXml = XmlUtils.removeComments(slaXml);
355            slaXml = evalSla.evaluate(slaXml, String.class);
356            XmlUtils.validateData(slaXml, SchemaName.SLA_ORIGINAL);
357            return slaXml;
358        }
359        catch (Exception e) {
360            throw new CommandException(ErrorCode.E1004, "Validation error :" + e.getMessage(), e);
361        }
362    }
363
364    /**
365     * Create an EL evaluator for a given group.
366     *
367     * @param conf configuration variable
368     * @param group group variable
369     * @return the evaluator created for the group
370     */
371    public static ELEvaluator createELEvaluatorForGroup(Configuration conf, String group) {
372        ELEvaluator eval = Services.get().get(ELService.class).createEvaluator(group);
373        for (Map.Entry<String, String> entry : conf) {
374            eval.setVariable(entry.getKey(), entry.getValue());
375        }
376        return eval;
377    }
378
379    @Override
380    public String getEntityKey() {
381        return null;
382    }
383
384    @Override
385    protected boolean isLockRequired() {
386        return false;
387    }
388
389    @Override
390    protected void loadState() {
391
392    }
393
394    @Override
395    protected void verifyPrecondition() throws CommandException {
396
397    }
398
399}