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.servlet; 020 021import java.io.IOException; 022import java.util.ArrayList; 023import java.util.HashSet; 024import java.util.List; 025import java.util.Set; 026 027import javax.servlet.http.HttpServletRequest; 028import javax.servlet.http.HttpServletResponse; 029 030import org.apache.hadoop.conf.Configuration; 031import org.apache.oozie.BaseEngineException; 032import org.apache.oozie.BulkResponseInfo; 033import org.apache.oozie.BundleJobBean; 034import org.apache.oozie.BundleJobInfo; 035import org.apache.oozie.CoordinatorEngine; 036import org.apache.oozie.BundleEngine; 037import org.apache.oozie.CoordinatorEngineException; 038import org.apache.oozie.BundleEngineException; 039import org.apache.oozie.CoordinatorJobBean; 040import org.apache.oozie.CoordinatorJobInfo; 041import org.apache.oozie.DagEngine; 042import org.apache.oozie.DagEngineException; 043import org.apache.oozie.ErrorCode; 044import org.apache.oozie.WorkflowJobBean; 045import org.apache.oozie.WorkflowsInfo; 046import org.apache.oozie.cli.OozieCLI; 047import org.apache.oozie.client.OozieClient; 048import org.apache.oozie.client.rest.BulkResponseImpl; 049import org.apache.oozie.client.rest.JsonTags; 050import org.apache.oozie.client.rest.RestConstants; 051import org.apache.oozie.service.CoordinatorEngineService; 052import org.apache.oozie.service.DagEngineService; 053import org.apache.oozie.service.BundleEngineService; 054import org.apache.oozie.service.Services; 055import org.apache.oozie.util.ConfigUtils; 056import org.apache.oozie.util.IOUtils; 057import org.apache.oozie.util.XLog; 058import org.apache.oozie.util.XmlUtils; 059import org.json.simple.JSONArray; 060import org.json.simple.JSONObject; 061 062public class V1JobsServlet extends BaseJobsServlet { 063 064 private static final String INSTRUMENTATION_NAME = "v1jobs"; 065 private static final Set<String> httpJobType = new HashSet<String>(){{ 066 this.add(OozieCLI.HIVE_CMD); 067 this.add(OozieCLI.SQOOP_CMD); 068 this.add(OozieCLI.PIG_CMD); 069 this.add(OozieCLI.MR_CMD); 070 }}; 071 072 public V1JobsServlet() { 073 super(INSTRUMENTATION_NAME); 074 } 075 076 /** 077 * v1 service implementation to submit a job, either workflow or coordinator 078 */ 079 @Override 080 protected JSONObject submitJob(HttpServletRequest request, Configuration conf) throws XServletException, 081 IOException { 082 JSONObject json = null; 083 084 String jobType = request.getParameter(RestConstants.JOBTYPE_PARAM); 085 086 if (!getUser(request).equals(UNDEF)) { 087 ConfigUtils.checkAndSetDisallowedProperties(conf, 088 getUser(request), 089 new XServletException(HttpServletResponse.SC_BAD_REQUEST, 090 ErrorCode.E0303, 091 "configuration", 092 OozieClient.USER_NAME), 093 false); 094 } 095 096 if (jobType == null) { 097 String wfPath = conf.get(OozieClient.APP_PATH); 098 String coordPath = conf.get(OozieClient.COORDINATOR_APP_PATH); 099 String bundlePath = conf.get(OozieClient.BUNDLE_APP_PATH); 100 101 ServletUtilities.ValidateAppPath(wfPath, coordPath, bundlePath); 102 103 if (wfPath != null) { 104 json = submitWorkflowJob(request, conf); 105 } 106 else if (coordPath != null) { 107 json = submitCoordinatorJob(request, conf); 108 } 109 else { 110 json = submitBundleJob(request, conf); 111 } 112 } 113 else { // This is a http submission job 114 if (httpJobType.contains(jobType)) { 115 json = submitHttpJob(request, conf, jobType); 116 } 117 else { 118 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ErrorCode.E0303, 119 RestConstants.JOBTYPE_PARAM, jobType); 120 } 121 } 122 return json; 123 } 124 125 /** 126 * v1 service implementation to get a JSONObject representation of a job from its external ID 127 */ 128 @Override 129 protected JSONObject getJobIdForExternalId(HttpServletRequest request, String externalId) throws XServletException, 130 IOException { 131 JSONObject json = null; 132 /* 133 * Configuration conf = new XConfiguration(); String wfPath = 134 * conf.get(OozieClient.APP_PATH); String coordPath = 135 * conf.get(OozieClient.COORDINATOR_APP_PATH); 136 * 137 * ServletUtilities.ValidateAppPath(wfPath, coordPath); 138 */ 139 String jobtype = request.getParameter(RestConstants.JOBTYPE_PARAM); 140 jobtype = (jobtype != null) ? jobtype : "wf"; 141 if (jobtype.contains("wf")) { 142 json = getWorkflowJobIdForExternalId(request, externalId); 143 } 144 else { 145 json = getCoordinatorJobIdForExternalId(request, externalId); 146 } 147 return json; 148 } 149 150 /** 151 * v1 service implementation to get a list of workflows, coordinators, or bundles, with filtering or interested 152 * windows embedded in the request object 153 */ 154 @Override 155 protected JSONObject getJobs(HttpServletRequest request) throws XServletException, IOException { 156 JSONObject json = null; 157 String isBulk = request.getParameter(RestConstants.JOBS_BULK_PARAM); 158 if(isBulk != null) { 159 json = getBulkJobs(request); 160 } else { 161 String jobtype = request.getParameter(RestConstants.JOBTYPE_PARAM); 162 jobtype = (jobtype != null) ? jobtype : "wf"; 163 164 if (jobtype.contains("wf")) { 165 json = getWorkflowJobs(request); 166 } 167 else if (jobtype.contains("coord")) { 168 json = getCoordinatorJobs(request); 169 } 170 else if (jobtype.contains("bundle")) { 171 json = getBundleJobs(request); 172 } 173 } 174 return json; 175 } 176 177 /** 178 * v1 service implementation to submit a workflow job 179 */ 180 @SuppressWarnings("unchecked") 181 private JSONObject submitWorkflowJob(HttpServletRequest request, Configuration conf) throws XServletException { 182 183 JSONObject json = new JSONObject(); 184 185 try { 186 String action = request.getParameter(RestConstants.ACTION_PARAM); 187 if (action != null && !action.equals(RestConstants.JOB_ACTION_START) 188 && !action.equals(RestConstants.JOB_ACTION_DRYRUN)) { 189 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ErrorCode.E0303, 190 RestConstants.ACTION_PARAM, action); 191 } 192 boolean startJob = (action != null); 193 String user = conf.get(OozieClient.USER_NAME); 194 DagEngine dagEngine = Services.get().get(DagEngineService.class).getDagEngine(user); 195 String id; 196 boolean dryrun = false; 197 if (action != null) { 198 dryrun = (action.equals(RestConstants.JOB_ACTION_DRYRUN)); 199 } 200 if (dryrun) { 201 id = dagEngine.dryRunSubmit(conf); 202 } 203 else { 204 id = dagEngine.submitJob(conf, startJob); 205 } 206 json.put(JsonTags.JOB_ID, id); 207 } 208 catch (BaseEngineException ex) { 209 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 210 } 211 212 return json; 213 } 214 215 /** 216 * v1 service implementation to submit a coordinator job 217 */ 218 @SuppressWarnings("unchecked") 219 private JSONObject submitCoordinatorJob(HttpServletRequest request, Configuration conf) throws XServletException { 220 221 JSONObject json = new JSONObject(); 222 XLog.getLog(getClass()).warn("submitCoordinatorJob " + XmlUtils.prettyPrint(conf).toString()); 223 try { 224 String action = request.getParameter(RestConstants.ACTION_PARAM); 225 if (action != null && !action.equals(RestConstants.JOB_ACTION_START) 226 && !action.equals(RestConstants.JOB_ACTION_DRYRUN)) { 227 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ErrorCode.E0303, 228 RestConstants.ACTION_PARAM, action); 229 } 230 boolean startJob = (action != null); 231 String user = conf.get(OozieClient.USER_NAME); 232 CoordinatorEngine coordEngine = Services.get().get(CoordinatorEngineService.class).getCoordinatorEngine( 233 user); 234 String id = null; 235 boolean dryrun = false; 236 if (action != null) { 237 dryrun = (action.equals(RestConstants.JOB_ACTION_DRYRUN)); 238 } 239 if (dryrun) { 240 id = coordEngine.dryRunSubmit(conf); 241 } 242 else { 243 id = coordEngine.submitJob(conf, startJob); 244 } 245 json.put(JsonTags.JOB_ID, id); 246 } 247 catch (CoordinatorEngineException ex) { 248 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 249 } 250 251 return json; 252 } 253 254 /** 255 * v1 service implementation to submit a bundle job 256 */ 257 @SuppressWarnings("unchecked") 258 private JSONObject submitBundleJob(HttpServletRequest request, Configuration conf) throws XServletException { 259 JSONObject json = new JSONObject(); 260 XLog.getLog(getClass()).warn("submitBundleJob " + XmlUtils.prettyPrint(conf).toString()); 261 try { 262 String action = request.getParameter(RestConstants.ACTION_PARAM); 263 if (action != null && !action.equals(RestConstants.JOB_ACTION_START) 264 && !action.equals(RestConstants.JOB_ACTION_DRYRUN)) { 265 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ErrorCode.E0303, 266 RestConstants.ACTION_PARAM, action); 267 } 268 boolean startJob = (action != null); 269 String user = conf.get(OozieClient.USER_NAME); 270 BundleEngine bundleEngine = Services.get().get(BundleEngineService.class).getBundleEngine(user); 271 String id = null; 272 boolean dryrun = false; 273 if (action != null) { 274 dryrun = (action.equals(RestConstants.JOB_ACTION_DRYRUN)); 275 } 276 if (dryrun) { 277 id = bundleEngine.dryRunSubmit(conf); 278 } 279 else { 280 id = bundleEngine.submitJob(conf, startJob); 281 } 282 json.put(JsonTags.JOB_ID, id); 283 } 284 catch (BundleEngineException ex) { 285 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 286 } 287 288 return json; 289 } 290 291 /** 292 * v1 service implementation to get a JSONObject representation of a job from its external ID 293 */ 294 @SuppressWarnings("unchecked") 295 private JSONObject getWorkflowJobIdForExternalId(HttpServletRequest request, String externalId) 296 throws XServletException { 297 JSONObject json = new JSONObject(); 298 try { 299 DagEngine dagEngine = Services.get().get(DagEngineService.class).getDagEngine(getUser(request)); 300 String jobId = dagEngine.getJobIdForExternalId(externalId); 301 json.put(JsonTags.JOB_ID, jobId); 302 } 303 catch (DagEngineException ex) { 304 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 305 } 306 return json; 307 } 308 309 /** 310 * v1 service implementation to get a JSONObject representation of a job from its external ID 311 */ 312 private JSONObject getCoordinatorJobIdForExternalId(HttpServletRequest request, String externalId) 313 throws XServletException { 314 JSONObject json = new JSONObject(); 315 return json; 316 } 317 318 /** 319 * v1 service implementation to get a list of workflows, with filtering or interested windows embedded in the 320 * request object 321 */ 322 private JSONObject getWorkflowJobs(HttpServletRequest request) throws XServletException { 323 JSONObject json = new JSONObject(); 324 try { 325 String filter = request.getParameter(RestConstants.JOBS_FILTER_PARAM); 326 String startStr = request.getParameter(RestConstants.OFFSET_PARAM); 327 String lenStr = request.getParameter(RestConstants.LEN_PARAM); 328 String timeZoneId = request.getParameter(RestConstants.TIME_ZONE_PARAM) == null 329 ? "GMT" : request.getParameter(RestConstants.TIME_ZONE_PARAM); 330 int start = (startStr != null) ? Integer.parseInt(startStr) : 1; 331 start = (start < 1) ? 1 : start; 332 int len = (lenStr != null) ? Integer.parseInt(lenStr) : 50; 333 len = (len < 1) ? 50 : len; 334 DagEngine dagEngine = Services.get().get(DagEngineService.class).getDagEngine(getUser(request)); 335 WorkflowsInfo jobs = dagEngine.getJobs(filter, start, len); 336 List<WorkflowJobBean> jsonWorkflows = jobs.getWorkflows(); 337 json.put(JsonTags.WORKFLOWS_JOBS, WorkflowJobBean.toJSONArray(jsonWorkflows, timeZoneId)); 338 json.put(JsonTags.WORKFLOWS_TOTAL, jobs.getTotal()); 339 json.put(JsonTags.WORKFLOWS_OFFSET, jobs.getStart()); 340 json.put(JsonTags.WORKFLOWS_LEN, jobs.getLen()); 341 342 } 343 catch (DagEngineException ex) { 344 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 345 } 346 347 return json; 348 } 349 350 /** 351 * v1 service implementation to get a list of workflows, with filtering or interested windows embedded in the 352 * request object 353 */ 354 @SuppressWarnings("unchecked") 355 private JSONObject getCoordinatorJobs(HttpServletRequest request) throws XServletException { 356 JSONObject json = new JSONObject(); 357 try { 358 String filter = request.getParameter(RestConstants.JOBS_FILTER_PARAM); 359 String startStr = request.getParameter(RestConstants.OFFSET_PARAM); 360 String lenStr = request.getParameter(RestConstants.LEN_PARAM); 361 String timeZoneId = request.getParameter(RestConstants.TIME_ZONE_PARAM) == null 362 ? "GMT" : request.getParameter(RestConstants.TIME_ZONE_PARAM); 363 int start = (startStr != null) ? Integer.parseInt(startStr) : 1; 364 start = (start < 1) ? 1 : start; 365 int len = (lenStr != null) ? Integer.parseInt(lenStr) : 50; 366 len = (len < 1) ? 50 : len; 367 CoordinatorEngine coordEngine = Services.get().get(CoordinatorEngineService.class).getCoordinatorEngine( 368 getUser(request)); 369 CoordinatorJobInfo jobs = coordEngine.getCoordJobs(filter, start, len); 370 List<CoordinatorJobBean> jsonJobs = jobs.getCoordJobs(); 371 json.put(JsonTags.COORDINATOR_JOBS, CoordinatorJobBean.toJSONArray(jsonJobs, timeZoneId)); 372 json.put(JsonTags.COORD_JOB_TOTAL, jobs.getTotal()); 373 json.put(JsonTags.COORD_JOB_OFFSET, jobs.getStart()); 374 json.put(JsonTags.COORD_JOB_LEN, jobs.getLen()); 375 376 } 377 catch (CoordinatorEngineException ex) { 378 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 379 } 380 return json; 381 } 382 383 @SuppressWarnings("unchecked") 384 private JSONObject getBundleJobs(HttpServletRequest request) throws XServletException { 385 JSONObject json = new JSONObject(); 386 try { 387 String filter = request.getParameter(RestConstants.JOBS_FILTER_PARAM); 388 String startStr = request.getParameter(RestConstants.OFFSET_PARAM); 389 String lenStr = request.getParameter(RestConstants.LEN_PARAM); 390 String timeZoneId = request.getParameter(RestConstants.TIME_ZONE_PARAM) == null 391 ? "GMT" : request.getParameter(RestConstants.TIME_ZONE_PARAM); 392 int start = (startStr != null) ? Integer.parseInt(startStr) : 1; 393 start = (start < 1) ? 1 : start; 394 int len = (lenStr != null) ? Integer.parseInt(lenStr) : 50; 395 len = (len < 1) ? 50 : len; 396 397 BundleEngine bundleEngine = Services.get().get(BundleEngineService.class).getBundleEngine(getUser(request)); 398 BundleJobInfo jobs = bundleEngine.getBundleJobs(filter, start, len); 399 List<BundleJobBean> jsonJobs = jobs.getBundleJobs(); 400 401 json.put(JsonTags.BUNDLE_JOBS, BundleJobBean.toJSONArray(jsonJobs, timeZoneId)); 402 json.put(JsonTags.BUNDLE_JOB_TOTAL, jobs.getTotal()); 403 json.put(JsonTags.BUNDLE_JOB_OFFSET, jobs.getStart()); 404 json.put(JsonTags.BUNDLE_JOB_LEN, jobs.getLen()); 405 406 } 407 catch (BundleEngineException ex) { 408 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 409 } 410 return json; 411 } 412 413 @SuppressWarnings("unchecked") 414 private JSONObject getBulkJobs(HttpServletRequest request) throws XServletException, IOException { 415 JSONObject json = new JSONObject(); 416 try { 417 String bulkFilter = request.getParameter(RestConstants.JOBS_BULK_PARAM); //REST API 418 String startStr = request.getParameter(RestConstants.OFFSET_PARAM); 419 String lenStr = request.getParameter(RestConstants.LEN_PARAM); 420 String timeZoneId = request.getParameter(RestConstants.TIME_ZONE_PARAM) == null 421 ? "GMT" : request.getParameter(RestConstants.TIME_ZONE_PARAM); 422 int start = (startStr != null) ? Integer.parseInt(startStr) : 1; 423 start = (start < 1) ? 1 : start; 424 int len = (lenStr != null) ? Integer.parseInt(lenStr) : 50; 425 len = (len < 1) ? 50 : len; 426 427 BundleEngine bundleEngine = Services.get().get(BundleEngineService.class).getBundleEngine(getUser(request)); 428 BulkResponseInfo bulkResponse = bundleEngine.getBulkJobs(bulkFilter, start, len); 429 List<BulkResponseImpl> responsesToJson = bulkResponse.getResponses(); 430 431 json.put(JsonTags.BULK_RESPONSES, BulkResponseImpl.toJSONArray(responsesToJson, timeZoneId)); 432 json.put(JsonTags.BULK_RESPONSE_TOTAL, bulkResponse.getTotal()); 433 json.put(JsonTags.BULK_RESPONSE_OFFSET, bulkResponse.getStart()); 434 json.put(JsonTags.BULK_RESPONSE_LEN, bulkResponse.getLen()); 435 436 } 437 catch (BaseEngineException ex) { 438 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 439 } 440 return json; 441 } 442 443 /** 444 * service implementation to submit a http job 445 */ 446 private JSONObject submitHttpJob(HttpServletRequest request, Configuration conf, String jobType) 447 throws XServletException { 448 JSONObject json = new JSONObject(); 449 450 try { 451 String user = conf.get(OozieClient.USER_NAME); 452 DagEngine dagEngine = Services.get().get(DagEngineService.class).getDagEngine(user); 453 String id = dagEngine.submitHttpJob(conf, jobType); 454 json.put(JsonTags.JOB_ID, id); 455 } 456 catch (DagEngineException ex) { 457 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 458 } 459 460 return json; 461 } 462 463 /** 464 * service implementation to bulk kill jobs 465 * @param request 466 * @param response 467 * @return 468 * @throws XServletException 469 * @throws IOException 470 */ 471 @Override 472 protected JSONObject killJobs(HttpServletRequest request, HttpServletResponse response) throws XServletException, 473 IOException { 474 return bulkModifyJobs(request, response); 475 } 476 477 /** 478 * service implementation to bulk suspend jobs 479 * @param request 480 * @param response 481 * @return 482 * @throws XServletException 483 * @throws IOException 484 */ 485 @Override 486 protected JSONObject suspendJobs(HttpServletRequest request, HttpServletResponse response) throws XServletException, 487 IOException { 488 return bulkModifyJobs(request, response); 489 } 490 491 /** 492 * service implementation to bulk resume jobs 493 * @param request 494 * @param response 495 * @return 496 * @throws XServletException 497 * @throws IOException 498 */ 499 @Override 500 protected JSONObject resumeJobs(HttpServletRequest request, HttpServletResponse response) throws XServletException, 501 IOException { 502 return bulkModifyJobs(request, response); 503 } 504 505 private JSONObject bulkModifyJobs(HttpServletRequest request, HttpServletResponse response) throws XServletException, 506 IOException { 507 String action = request.getParameter(RestConstants.ACTION_PARAM); 508 String jobType = request.getParameter(RestConstants.JOBTYPE_PARAM); 509 String filter = request.getParameter(RestConstants.JOBS_FILTER_PARAM); 510 String startStr = request.getParameter(RestConstants.OFFSET_PARAM); 511 String lenStr = request.getParameter(RestConstants.LEN_PARAM); 512 String timeZoneId = request.getParameter(RestConstants.TIME_ZONE_PARAM) == null 513 ? "GMT" : request.getParameter(RestConstants.TIME_ZONE_PARAM); 514 515 int start = (startStr != null) ? Integer.parseInt(startStr) : 1; 516 start = (start < 1) ? 1 : start; 517 int len = (lenStr != null) ? Integer.parseInt(lenStr) : 50; 518 len = (len < 1) ? 50 : len; 519 520 JSONObject json = new JSONObject(); 521 List<String> ids = new ArrayList<String>(); 522 523 if (jobType.equals("wf")) { 524 WorkflowsInfo jobs = null; 525 DagEngine dagEngine = Services.get().get(DagEngineService.class).getDagEngine(getUser(request)); 526 if (action.equals(RestConstants.JOB_ACTION_KILL)) { 527 try { 528 jobs = dagEngine.killJobs(filter, start, len); 529 } 530 catch (DagEngineException ex) { 531 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 532 } 533 } else if (action.equals(RestConstants.JOB_ACTION_SUSPEND)) { 534 try { 535 jobs = dagEngine.suspendJobs(filter, start, len); 536 } 537 catch (DagEngineException ex) { 538 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 539 } 540 } else if (action.equals(RestConstants.JOB_ACTION_RESUME)) { 541 try { 542 jobs = dagEngine.resumeJobs(filter, start, len); 543 } 544 catch (DagEngineException ex) { 545 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 546 } 547 } 548 549 json.put(JsonTags.WORKFLOWS_JOBS, WorkflowJobBean.toJSONArray(jobs.getWorkflows(), timeZoneId)); 550 json.put(JsonTags.WORKFLOWS_TOTAL, jobs.getTotal()); 551 json.put(JsonTags.WORKFLOWS_OFFSET, jobs.getStart()); 552 json.put(JsonTags.WORKFLOWS_LEN, jobs.getLen()); 553 554 } 555 else if (jobType.equals("bundle")) { 556 BundleJobInfo jobs = null; 557 BundleEngine bundleEngine = Services.get().get(BundleEngineService.class). 558 getBundleEngine(getUser(request)); 559 if (action.equals(RestConstants.JOB_ACTION_KILL)) { 560 try { 561 jobs = bundleEngine.killJobs(filter, start, len); 562 } 563 catch (BundleEngineException ex) { 564 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 565 } 566 } 567 else if (action.equals(RestConstants.JOB_ACTION_SUSPEND)) { 568 try { 569 jobs = bundleEngine.suspendJobs(filter, start, len); 570 } 571 catch (BundleEngineException ex) { 572 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 573 } 574 } 575 else if (action.equals(RestConstants.JOB_ACTION_RESUME)) { 576 try { 577 jobs = bundleEngine.resumeJobs(filter, start, len); 578 } 579 catch (BundleEngineException ex) { 580 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 581 } 582 } 583 584 json.put(JsonTags.BUNDLE_JOBS, BundleJobBean.toJSONArray(jobs.getBundleJobs(), timeZoneId)); 585 json.put(JsonTags.BUNDLE_JOB_TOTAL, jobs.getTotal()); 586 json.put(JsonTags.BUNDLE_JOB_OFFSET, jobs.getStart()); 587 json.put(JsonTags.BUNDLE_JOB_LEN, jobs.getLen()); 588 } 589 else { 590 CoordinatorJobInfo jobs = null; 591 CoordinatorEngine coordEngine = Services.get().get(CoordinatorEngineService.class). 592 getCoordinatorEngine(getUser(request)); 593 if (action.equals(RestConstants.JOB_ACTION_KILL)) { 594 try { 595 jobs = coordEngine.killJobs(filter, start, len); 596 } 597 catch (CoordinatorEngineException ex) { 598 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 599 } 600 } 601 else if (action.equals(RestConstants.JOB_ACTION_SUSPEND)) { 602 try { 603 jobs = coordEngine.suspendJobs(filter, start, len); 604 } 605 catch (CoordinatorEngineException ex) { 606 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 607 } 608 } 609 else if (action.equals(RestConstants.JOB_ACTION_RESUME)) { 610 try { 611 jobs = coordEngine.resumeJobs(filter, start, len); 612 } 613 catch (CoordinatorEngineException ex) { 614 throw new XServletException(HttpServletResponse.SC_BAD_REQUEST, ex); 615 } 616 } 617 618 json.put(JsonTags.COORDINATOR_JOBS, CoordinatorJobBean.toJSONArray(jobs.getCoordJobs(), timeZoneId)); 619 json.put(JsonTags.COORD_JOB_TOTAL, jobs.getTotal()); 620 json.put(JsonTags.COORD_JOB_OFFSET, jobs.getStart()); 621 json.put(JsonTags.COORD_JOB_LEN, jobs.getLen()); 622 } 623 624 json.put(JsonTags.JOB_IDS, toJSONArray(ids)); 625 return json; 626 } 627 628 private static JSONArray toJSONArray(List<String> ids) { 629 JSONArray array = new JSONArray(); 630 for (String id : ids) { 631 array.add(id); 632 } 633 return array; 634 } 635}