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.oozie.WorkflowJobBean;
023import org.apache.oozie.ErrorCode;
024import org.apache.oozie.service.JPAService;
025import org.apache.oozie.service.WorkflowStoreService;
026import org.apache.oozie.service.WorkflowAppService;
027import org.apache.oozie.service.Services;
028import org.apache.oozie.service.DagXLogInfoService;
029import org.apache.oozie.util.InstrumentUtils;
030import org.apache.oozie.util.LogUtils;
031import org.apache.oozie.util.XLog;
032import org.apache.oozie.util.ParamChecker;
033import org.apache.oozie.util.XConfiguration;
034import org.apache.oozie.util.XmlUtils;
035import org.apache.oozie.command.CommandException;
036import org.apache.oozie.executor.jpa.WorkflowJobInsertJPAExecutor;
037import org.apache.oozie.store.StoreException;
038import org.apache.oozie.workflow.WorkflowApp;
039import org.apache.oozie.workflow.WorkflowException;
040import org.apache.oozie.workflow.WorkflowInstance;
041import org.apache.oozie.workflow.WorkflowLib;
042import org.apache.oozie.util.PropertiesUtils;
043import org.apache.oozie.client.OozieClient;
044import org.apache.oozie.client.WorkflowJob;
045import org.apache.oozie.client.XOozieClient;
046import org.jdom.Element;
047import org.jdom.Namespace;
048
049import java.util.Date;
050import java.util.List;
051import java.util.Map;
052import java.util.Set;
053import java.util.HashSet;
054
055public abstract class SubmitHttpXCommand extends WorkflowXCommand<String> {
056
057    protected static final Set<String> MANDATORY_OOZIE_CONFS = new HashSet<String>();
058    protected static final Set<String> OPTIONAL_OOZIE_CONFS = new HashSet<String>();
059
060    static {
061        MANDATORY_OOZIE_CONFS.add(XOozieClient.JT);
062        MANDATORY_OOZIE_CONFS.add(XOozieClient.NN);
063        MANDATORY_OOZIE_CONFS.add(OozieClient.LIBPATH);
064
065        OPTIONAL_OOZIE_CONFS.add(XOozieClient.FILES);
066        OPTIONAL_OOZIE_CONFS.add(XOozieClient.ARCHIVES);
067    }
068
069    private Configuration conf;
070
071    public SubmitHttpXCommand(String name, String type, Configuration conf) {
072        super(name, type, 1);
073        this.conf = ParamChecker.notNull(conf, "conf");
074    }
075
076    private static final Set<String> DISALLOWED_USER_PROPERTIES = new HashSet<String>();
077
078    static {
079        String[] badUserProps = { PropertiesUtils.DAYS, PropertiesUtils.HOURS, PropertiesUtils.MINUTES,
080                PropertiesUtils.KB, PropertiesUtils.MB, PropertiesUtils.GB, PropertiesUtils.TB, PropertiesUtils.PB,
081                PropertiesUtils.RECORDS, PropertiesUtils.MAP_IN, PropertiesUtils.MAP_OUT, PropertiesUtils.REDUCE_IN,
082                PropertiesUtils.REDUCE_OUT, PropertiesUtils.GROUPS };
083        PropertiesUtils.createPropertySet(badUserProps, DISALLOWED_USER_PROPERTIES);
084    }
085
086    abstract protected Element generateSection(Configuration conf, Namespace ns);
087
088    abstract protected Namespace getSectionNamespace();
089
090    abstract protected String getWorkflowName();
091
092    protected void checkMandatoryConf(Configuration conf) {
093        for (String key : MANDATORY_OOZIE_CONFS) {
094            String value = conf.get(key);
095            if (value == null) {
096                throw new RuntimeException(key + " is not specified");
097            }
098        }
099    }
100
101    protected Namespace getWorkflowNamespace() {
102        return Namespace.getNamespace("uri:oozie:workflow:0.2");
103    }
104    /**
105     * Generate workflow xml from conf object
106     *
107     * @param conf the configuration object
108     * @return workflow xml def string representation
109     */
110    protected String getWorkflowXml(Configuration conf) {
111        checkMandatoryConf(conf);
112
113        Namespace ns = getWorkflowNamespace();
114        Element root = new Element("workflow-app", ns);
115        String name = getWorkflowName();
116        root.setAttribute("name", "oozie-" + name);
117
118        Element start = new Element("start", ns);
119        String nodeName = name + "1";
120        start.setAttribute("to", nodeName);
121        root.addContent(start);
122
123        Element action = new Element("action", ns);
124        action.setAttribute("name", nodeName);
125
126        Element ele = generateSection(conf, getSectionNamespace());
127        action.addContent(ele);
128
129        Element ok = new Element("ok", ns);
130        ok.setAttribute("to", "end");
131        action.addContent(ok);
132
133        Element error = new Element("error", ns);
134        error.setAttribute("to", "fail");
135        action.addContent(error);
136
137        root.addContent(action);
138
139        Element kill = new Element("kill", ns);
140        kill.setAttribute("name", "fail");
141        Element message = new Element("message", ns);
142        message.addContent(name + " failed, error message[${wf:errorMessage(wf:lastErrorNode())}]");
143        kill.addContent(message);
144        root.addContent(kill);
145
146        Element end = new Element("end", ns);
147        end.setAttribute("name", "end");
148        root.addContent(end);
149
150        return XmlUtils.prettyPrint(root).toString();
151    };
152
153    protected Element generateConfigurationSection(List<String> Dargs, Namespace ns) {
154        Element configuration = new Element("configuration", ns);
155        for (String arg : Dargs) {
156            String name = null, value = null;
157            int pos = arg.indexOf("=");
158            if (pos == -1) { // "-D<name>" or "-D" only
159                name = arg.substring(2, arg.length());
160                value = "";
161            }
162            else { // "-D<name>=<value>"
163                name = arg.substring(2, pos);
164                value = arg.substring(pos + 1, arg.length());
165            }
166
167            Element property = new Element("property", ns);
168            Element nameElement = new Element("name", ns);
169            nameElement.addContent(name);
170            property.addContent(nameElement);
171            Element valueElement = new Element("value", ns);
172            valueElement.addContent(value);
173            property.addContent(valueElement);
174            configuration.addContent(property);
175        }
176
177        return configuration;
178    }
179
180    /* (non-Javadoc)
181     * @see org.apache.oozie.command.XCommand#execute()
182     */
183    @Override
184    protected String execute() throws CommandException {
185        InstrumentUtils.incrJobCounter(getName(), 1, getInstrumentation());
186        WorkflowAppService wps = Services.get().get(WorkflowAppService.class);
187        try {
188            XLog.Info.get().setParameter(DagXLogInfoService.TOKEN, conf.get(OozieClient.LOG_TOKEN));
189            String wfXml = getWorkflowXml(conf);
190            LOG.debug("workflow xml created on the server side is :\n");
191            LOG.debug(wfXml);
192            WorkflowApp app = wps.parseDef(wfXml, conf);
193            XConfiguration protoActionConf = wps.createProtoActionConf(conf, false);
194            WorkflowLib workflowLib = Services.get().get(WorkflowStoreService.class).getWorkflowLibWithNoDB();
195
196            PropertiesUtils.checkDisallowedProperties(conf, DISALLOWED_USER_PROPERTIES);
197            PropertiesUtils.checkDefaultDisallowedProperties(conf);
198
199            // Resolving all variables in the job properties.
200            // This ensures the Hadoop Configuration semantics is preserved.
201            XConfiguration resolvedVarsConf = new XConfiguration();
202            for (Map.Entry<String, String> entry : conf) {
203                resolvedVarsConf.set(entry.getKey(), conf.get(entry.getKey()));
204            }
205            conf = resolvedVarsConf;
206
207            WorkflowInstance wfInstance;
208            try {
209                wfInstance = workflowLib.createInstance(app, conf);
210            }
211            catch (WorkflowException e) {
212                throw new StoreException(e);
213            }
214
215            Configuration conf = wfInstance.getConf();
216
217            WorkflowJobBean workflow = new WorkflowJobBean();
218            workflow.setId(wfInstance.getId());
219            workflow.setAppName(app.getName());
220            workflow.setAppPath(conf.get(OozieClient.APP_PATH));
221            workflow.setConf(XmlUtils.prettyPrint(conf).toString());
222            workflow.setProtoActionConf(protoActionConf.toXmlString());
223            workflow.setCreatedTime(new Date());
224            workflow.setLastModifiedTime(new Date());
225            workflow.setLogToken(conf.get(OozieClient.LOG_TOKEN, ""));
226            workflow.setStatus(WorkflowJob.Status.PREP);
227            workflow.setRun(0);
228            workflow.setUser(conf.get(OozieClient.USER_NAME));
229            workflow.setGroup(conf.get(OozieClient.GROUP_NAME));
230            workflow.setWorkflowInstance(wfInstance);
231            workflow.setExternalId(conf.get(OozieClient.EXTERNAL_ID));
232
233            LogUtils.setLogInfo(workflow);
234            JPAService jpaService = Services.get().get(JPAService.class);
235            if (jpaService != null) {
236                jpaService.execute(new WorkflowJobInsertJPAExecutor(workflow));
237            }
238            else {
239                LOG.error(ErrorCode.E0610);
240                return null;
241            }
242
243            return workflow.getId();
244        }
245        catch (WorkflowException ex) {
246            throw new CommandException(ex);
247        }
248        catch (Exception ex) {
249            throw new CommandException(ErrorCode.E0803, ex.getMessage(), ex);
250        }
251    }
252
253    static private void addSection(Element X, Namespace ns, String filesStr, String tagName) {
254        if (filesStr != null) {
255            String[] files = filesStr.split(",");
256            for (String f : files) {
257                Element tagElement = new Element(tagName, ns);
258                if (f.contains("#")) {
259                    tagElement.addContent(f);
260                }
261                else {
262                    String filename = f.substring(f.lastIndexOf("/") + 1, f.length());
263                    if (filename == null || filename.isEmpty()) {
264                        tagElement.addContent(f);
265                    }
266                    else {
267                        tagElement.addContent(f + "#" + filename);
268                    }
269                }
270                X.addContent(tagElement);
271            }
272        }
273    }
274
275    /**
276     * Add file section in X.
277     *
278     * @param X XML element to be appended
279     * @param conf Configuration object
280     * @param ns XML element namespace
281     */
282    static void addFileSection(Element X, Configuration conf, Namespace ns) {
283        String filesStr = conf.get(XOozieClient.FILES);
284        addSection(X, ns, filesStr, "file");
285    }
286
287    /**
288     * Add archive section in X.
289     *
290     * @param X XML element to be appended
291     * @param conf Configuration object
292     * @param ns XML element namespace
293     */
294    static void addArchiveSection(Element X, Configuration conf, Namespace ns) {
295        String archivesStr = conf.get(XOozieClient.ARCHIVES);
296        addSection(X, ns, archivesStr, "archive");
297    }
298}