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
019
020package org.apache.oozie.command.coord;
021
022import java.io.ByteArrayOutputStream;
023import java.io.IOException;
024import java.io.StringReader;
025import java.util.Date;
026
027import org.apache.commons.lang.StringUtils;
028import org.apache.hadoop.conf.Configuration;
029import org.apache.oozie.CoordinatorJobBean;
030import org.apache.oozie.ErrorCode;
031import org.apache.oozie.client.CoordinatorJob;
032import org.apache.oozie.client.OozieClient;
033import org.apache.oozie.command.CommandException;
034import org.apache.oozie.executor.jpa.CoordJobQueryExecutor;
035import org.apache.oozie.executor.jpa.JPAExecutorException;
036import org.apache.oozie.executor.jpa.CoordJobQueryExecutor.CoordJobQuery;
037import org.apache.oozie.service.JPAService;
038import org.apache.oozie.service.Services;
039import org.apache.oozie.util.ConfigUtils;
040import org.apache.oozie.util.LogUtils;
041import org.apache.oozie.util.XConfiguration;
042import org.apache.oozie.util.XmlUtils;
043import org.eclipse.jgit.diff.DiffFormatter;
044import org.eclipse.jgit.diff.EditList;
045import org.eclipse.jgit.diff.HistogramDiff;
046import org.eclipse.jgit.diff.RawText;
047import org.eclipse.jgit.diff.RawTextComparator;
048import org.jdom.Element;
049
050/**
051 * This class provides the functionalities to update coordinator job XML and properties. It uses CoordSubmitXCommand
052 * functionality to validate XML and resolve all the variables or properties using job configurations.
053 */
054public class CoordUpdateXCommand extends CoordSubmitXCommand {
055
056    private final String jobId;
057    private boolean showDiff = true;
058    private boolean isConfChange = false;
059
060    //This properties are set in coord jobs by bundle. An update command should not overide it.
061    final static String[] bundleConfList = new String[] { OozieClient.BUNDLE_ID };
062
063    StringBuffer diff = new StringBuffer();
064    CoordinatorJobBean oldCoordJob = null;
065
066    public CoordUpdateXCommand(boolean dryrun, Configuration conf, String jobId) {
067        super(dryrun, conf);
068        this.jobId = jobId;
069        isConfChange = conf.size() == 0 ? false : true;
070    }
071
072    public CoordUpdateXCommand(boolean dryrun, Configuration conf, String jobId, boolean showDiff) {
073        super(dryrun, conf);
074        this.jobId = jobId;
075        this.showDiff = showDiff;
076        isConfChange = conf.size() == 0 ? false : true;
077    }
078
079    @Override
080    protected String storeToDB(String xmlElement, Element eJob, CoordinatorJobBean coordJob) throws CommandException {
081        check(oldCoordJob, coordJob);
082
083        ConfigUtils.checkAndSetDisallowedProperties(conf,
084                this.oldCoordJob.getUser(),
085                new CommandException(ErrorCode.E1003,
086                        String.format("%s=%s", OozieClient.USER_NAME, conf.get(OozieClient.USER_NAME))),
087                true);
088
089        computeDiff(eJob);
090        oldCoordJob.setAppPath(conf.get(OozieClient.COORDINATOR_APP_PATH));
091        if (isConfChange) {
092            oldCoordJob.setConf(XmlUtils.prettyPrint(conf).toString());
093        }
094        oldCoordJob.setMatThrottling(coordJob.getMatThrottling());
095        oldCoordJob.setOrigJobXml(xmlElement);
096        oldCoordJob.setConcurrency(coordJob.getConcurrency());
097        oldCoordJob.setExecution(coordJob.getExecution());
098        oldCoordJob.setTimeout(coordJob.getTimeout());
099        oldCoordJob.setJobXml(XmlUtils.prettyPrint(eJob).toString());
100
101
102        if (!dryrun) {
103            oldCoordJob.setLastModifiedTime(new Date());
104            // Should log the changes, this should be useful for debugging.
105            LOG.info("Coord update changes : " + diff.toString());
106            try {
107                CoordJobQueryExecutor.getInstance().executeUpdate(CoordJobQuery.UPDATE_COORD_JOB, oldCoordJob);
108            }
109            catch (JPAExecutorException jpaee) {
110                throw new CommandException(jpaee);
111            }
112        }
113        return jobId;
114    }
115
116    @Override
117    protected void loadState() throws CommandException {
118        super.loadState();
119        jpaService = Services.get().get(JPAService.class);
120        if (jpaService == null) {
121            throw new CommandException(ErrorCode.E0610);
122        }
123        coordJob = new CoordinatorJobBean();
124        try {
125            oldCoordJob = CoordJobQueryExecutor.getInstance().get(CoordJobQuery.GET_COORD_JOB, jobId);
126        }
127        catch (JPAExecutorException e) {
128            throw new CommandException(e);
129        }
130
131        LogUtils.setLogInfo(oldCoordJob);
132        if (!isConfChange || StringUtils.isEmpty(conf.get(OozieClient.COORDINATOR_APP_PATH))) {
133            try {
134                XConfiguration jobConf = new XConfiguration(new StringReader(oldCoordJob.getConf()));
135
136                if (!isConfChange) {
137                    conf = jobConf;
138                }
139                else {
140                    for (String bundleConfKey : bundleConfList) {
141                        if (jobConf.get(bundleConfKey) != null) {
142                            conf.set(bundleConfKey, jobConf.get(bundleConfKey));
143                        }
144                    }
145                    if (StringUtils.isEmpty(conf.get(OozieClient.COORDINATOR_APP_PATH))) {
146                        conf.set(OozieClient.COORDINATOR_APP_PATH, jobConf.get(OozieClient.COORDINATOR_APP_PATH));
147                    }
148                }
149            }
150            catch (Exception e) {
151                throw new CommandException(ErrorCode.E1023, e.getMessage(), e);
152            }
153        }
154        coordJob.setConf(XmlUtils.prettyPrint(conf).toString());
155        setJob(coordJob);
156    }
157
158    @Override
159    protected void verifyPrecondition() throws CommandException {
160        if (coordJob.getStatus() == CoordinatorJob.Status.SUCCEEDED
161                || coordJob.getStatus() == CoordinatorJob.Status.DONEWITHERROR) {
162            LOG.info("Can't update coord job. Job has finished processing");
163            throw new CommandException(ErrorCode.E1023, "Can't update coord job. Job has finished processing");
164        }
165    }
166
167    /**
168     * Gets the difference of job definition and properties.
169     *
170     * @param eJob the e job
171     * @return the diff
172     */
173
174    private void computeDiff(Element eJob) {
175        try {
176            diff.append("**********Job definition changes**********").append(System.getProperty("line.separator"));
177            diff.append(getDiffinGitFormat(oldCoordJob.getJobXml(), XmlUtils.prettyPrint(eJob).toString()));
178            diff.append("******************************************").append(System.getProperty("line.separator"));
179            diff.append("**********Job conf changes****************").append(System.getProperty("line.separator"));
180            if (isConfChange) {
181                diff.append(getDiffinGitFormat(oldCoordJob.getConf(), XmlUtils.prettyPrint(conf).toString()));
182            }
183            else {
184                diff.append("No conf update requested").append(System.getProperty("line.separator"));
185            }
186            diff.append("******************************************").append(System.getProperty("line.separator"));
187        }
188        catch (IOException e) {
189            diff.append("Error computing diff. Error " + e.getMessage());
190            LOG.warn("Error computing diff.", e);
191        }
192    }
193
194    /**
195     * Get the differences in git format.
196     *
197     * @param string1 the string1
198     * @param string2 the string2
199     * @return the diff
200     * @throws IOException Signals that an I/O exception has occurred.
201     */
202    private String getDiffinGitFormat(String string1, String string2) throws IOException {
203        ByteArrayOutputStream out = new ByteArrayOutputStream();
204        RawText rt1 = new RawText(string1.getBytes());
205        RawText rt2 = new RawText(string2.getBytes());
206        EditList diffList = new EditList();
207        diffList.addAll(new HistogramDiff().diff(RawTextComparator.DEFAULT, rt1, rt2));
208        new DiffFormatter(out).format(diffList, rt1, rt2);
209        return out.toString();
210    }
211
212    @Override
213    protected String submit() throws CommandException {
214        LOG.info("STARTED Coordinator update");
215        submitJob();
216        LOG.info("ENDED Coordinator update");
217        if (showDiff) {
218            return diff.toString();
219        }
220        else {
221            return "";
222        }
223    }
224
225    /**
226     * Check. Frequency can't be changed. EndTime can't be changed. StartTime can't be changed. AppName can't be changed
227     * Timeunit can't be changed. Timezone can't be changed
228     *
229     * @param oldCoord the old coord
230     * @param newCoord the new coord
231     * @throws CommandException the command exception
232     */
233    public void check(CoordinatorJobBean oldCoord, CoordinatorJobBean newCoord) throws CommandException {
234        if (!oldCoord.getFrequency().equals(newCoord.getFrequency())) {
235            throw new CommandException(ErrorCode.E1023, "Frequency can't be changed. Old frequency = "
236                    + oldCoord.getFrequency() + " new frequency = " + newCoord.getFrequency());
237        }
238
239        if (!oldCoord.getEndTime().equals(newCoord.getEndTime())) {
240            throw new CommandException(ErrorCode.E1023, "End time can't be changed. Old end time = "
241                    + oldCoord.getEndTime() + " new end time = " + newCoord.getEndTime());
242        }
243
244        if (!oldCoord.getStartTime().equals(newCoord.getStartTime())) {
245            throw new CommandException(ErrorCode.E1023, "Start time can't be changed. Old start time = "
246                    + oldCoord.getStartTime() + " new start time = " + newCoord.getStartTime());
247        }
248
249        if (!oldCoord.getAppName().equals(newCoord.getAppName())) {
250            throw new CommandException(ErrorCode.E1023, "Coord name can't be changed. Old name = "
251                    + oldCoord.getAppName() + " new name = " + newCoord.getAppName());
252        }
253
254        if (!oldCoord.getTimeUnitStr().equals(newCoord.getTimeUnitStr())) {
255            throw new CommandException(ErrorCode.E1023, "Timeunit can't be changed. Old Timeunit = "
256                    + oldCoord.getTimeUnitStr() + " new Timeunit = " + newCoord.getTimeUnitStr());
257        }
258
259        if (!oldCoord.getTimeZone().equals(newCoord.getTimeZone())) {
260            throw new CommandException(ErrorCode.E1023, "TimeZone can't be changed. Old timeZone = "
261                    + oldCoord.getTimeZone() + " new timeZone = " + newCoord.getTimeZone());
262        }
263
264    }
265
266    @Override
267    protected void queueMaterializeTransitionXCommand(String jobId) {
268    }
269
270    @Override
271    public void notifyParent() throws CommandException {
272    }
273
274    @Override
275    protected boolean isLockRequired() {
276        return true;
277    }
278
279    @Override
280    public String getEntityKey() {
281        return jobId;
282    }
283
284    @Override
285    public void transitToNext() {
286    }
287
288    @Override
289    public String getKey() {
290        return getName() + "_" + jobId;
291    }
292
293    @Override
294    public String getDryRun(CoordinatorJobBean job) throws Exception{
295        return super.getDryRun(oldCoordJob);
296    }
297}