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}