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.sla;
020
021import java.sql.Timestamp;
022import java.util.ArrayList;
023import java.util.Collections;
024import java.util.Date;
025import java.util.HashSet;
026import java.util.Iterator;
027import java.util.List;
028import java.util.Map;
029import java.util.Set;
030import java.util.concurrent.ConcurrentHashMap;
031
032import org.apache.hadoop.conf.Configuration;
033import org.apache.oozie.AppType;
034import org.apache.oozie.CoordinatorActionBean;
035import org.apache.oozie.CoordinatorJobBean;
036import org.apache.oozie.ErrorCode;
037import org.apache.oozie.WorkflowActionBean;
038import org.apache.oozie.WorkflowJobBean;
039import org.apache.oozie.XException;
040import org.apache.oozie.client.CoordinatorAction;
041import org.apache.oozie.client.WorkflowAction;
042import org.apache.oozie.client.WorkflowJob;
043import org.apache.oozie.client.event.JobEvent;
044import org.apache.oozie.client.event.SLAEvent.EventStatus;
045import org.apache.oozie.client.event.SLAEvent.SLAStatus;
046import org.apache.oozie.client.rest.JsonBean;
047import org.apache.oozie.executor.jpa.BatchQueryExecutor;
048import org.apache.oozie.executor.jpa.CoordActionGetForSLAJPAExecutor;
049import org.apache.oozie.executor.jpa.CoordActionQueryExecutor;
050import org.apache.oozie.executor.jpa.CoordActionQueryExecutor.CoordActionQuery;
051import org.apache.oozie.executor.jpa.CoordJobQueryExecutor;
052import org.apache.oozie.executor.jpa.CoordJobQueryExecutor.CoordJobQuery;
053import org.apache.oozie.executor.jpa.JPAExecutorException;
054import org.apache.oozie.executor.jpa.SLARegistrationQueryExecutor;
055import org.apache.oozie.executor.jpa.WorkflowActionQueryExecutor;
056import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor;
057import org.apache.oozie.executor.jpa.SLARegistrationQueryExecutor.SLARegQuery;
058import org.apache.oozie.executor.jpa.SLASummaryQueryExecutor;
059import org.apache.oozie.executor.jpa.WorkflowActionGetForSLAJPAExecutor;
060import org.apache.oozie.executor.jpa.WorkflowJobGetForSLAJPAExecutor;
061import org.apache.oozie.executor.jpa.WorkflowActionQueryExecutor.WorkflowActionQuery;
062import org.apache.oozie.executor.jpa.WorkflowJobQueryExecutor.WorkflowJobQuery;
063import org.apache.oozie.executor.jpa.sla.SLASummaryGetRecordsOnRestartJPAExecutor;
064import org.apache.oozie.executor.jpa.SLASummaryQueryExecutor.SLASummaryQuery;
065import org.apache.oozie.executor.jpa.BatchQueryExecutor.UpdateEntry;
066import org.apache.oozie.lock.LockToken;
067import org.apache.oozie.service.ConfigurationService;
068import org.apache.oozie.service.EventHandlerService;
069import org.apache.oozie.service.InstrumentationService;
070import org.apache.oozie.service.JPAService;
071import org.apache.oozie.service.JobsConcurrencyService;
072import org.apache.oozie.service.MemoryLocksService;
073import org.apache.oozie.service.SchedulerService;
074import org.apache.oozie.service.ServiceException;
075import org.apache.oozie.service.Services;
076import org.apache.oozie.sla.service.SLAService;
077import org.apache.oozie.util.DateUtils;
078import org.apache.oozie.util.Instrumentation;
079import org.apache.oozie.util.LogUtils;
080import org.apache.oozie.util.XLog;
081
082import com.google.common.annotations.VisibleForTesting;
083
084
085/**
086 * Implementation class for SLACalculator that calculates SLA related to
087 * start/end/duration of jobs using a memory-based map
088 */
089public class SLACalculatorMemory implements SLACalculator {
090
091    private static XLog LOG = XLog.getLog(SLACalculatorMemory.class);
092    // TODO optimization priority based insertion/processing/bumping up-down
093    private Map<String, SLACalcStatus> slaMap;
094    protected Set<String> historySet;
095    private static int capacity;
096    private static JPAService jpaService;
097    protected EventHandlerService eventHandler;
098    private static int modifiedAfter;
099    private static long jobEventLatency;
100    private Instrumentation instrumentation;
101    public static final String INSTRUMENTATION_GROUP = "sla-calculator";
102    public static final String SLA_MAP = "sla-map";
103    private int maxRetryCount;
104
105    @Override
106    public void init(Configuration conf) throws ServiceException {
107        capacity = ConfigurationService.getInt(conf, SLAService.CONF_CAPACITY);
108        jobEventLatency = ConfigurationService.getInt(conf, SLAService.CONF_JOB_EVENT_LATENCY);
109        maxRetryCount = ConfigurationService.getInt(conf, SLAService.CONF_MAXIMUM_RETRY_COUNT);
110        slaMap = new ConcurrentHashMap<String, SLACalcStatus>();
111        historySet = Collections.synchronizedSet(new HashSet<String>());
112        jpaService = Services.get().get(JPAService.class);
113        eventHandler = Services.get().get(EventHandlerService.class);
114        instrumentation = Services.get().get(InstrumentationService.class).get();
115        // load events modified after
116        modifiedAfter = conf.getInt(SLAService.CONF_EVENTS_MODIFIED_AFTER, 7);
117        loadOnRestart();
118        Runnable purgeThread = new HistoryPurgeWorker();
119        // schedule runnable by default 1 day
120        Services.get()
121                .get(SchedulerService.class)
122                .schedule(purgeThread, 86400, Services.get().getConf().getInt(SLAService.CONF_SLA_HISTORY_PURGE_INTERVAL, 86400),
123                        SchedulerService.Unit.SEC);
124    }
125
126    public class HistoryPurgeWorker implements Runnable {
127
128        public HistoryPurgeWorker() {
129        }
130
131        @Override
132        public void run() {
133            if (Thread.currentThread().isInterrupted()) {
134                return;
135            }
136            Iterator<String> jobItr = historySet.iterator();
137            while (jobItr.hasNext()) {
138                String jobId = jobItr.next();
139
140                if (jobId.endsWith("-W")) {
141                    WorkflowJobBean wfJob = null;
142                    try {
143                        wfJob = WorkflowJobQueryExecutor.getInstance().get(WorkflowJobQuery.GET_WORKFLOW_STATUS, jobId);
144                    }
145                    catch (JPAExecutorException e) {
146                        if (e.getErrorCode().equals(ErrorCode.E0604)) {
147                            jobItr.remove();
148                        }
149                        else {
150                            LOG.info("Failed to fetch the workflow job: " + jobId, e);
151                        }
152                    }
153                    if (wfJob != null && wfJob.inTerminalState()) {
154                        try {
155                            updateSLASummary(wfJob.getId(), wfJob.getStartTime(), wfJob.getEndTime());
156                            jobItr.remove();
157                        }
158                        catch (JPAExecutorException e) {
159                            LOG.info("Failed to update SLASummaryBean when purging history set entry for " + jobId, e);
160                        }
161
162                    }
163                }
164                else if (jobId.contains("-W@")) {
165                    WorkflowActionBean wfAction = null;
166                    try {
167                        wfAction = WorkflowActionQueryExecutor.getInstance().get(
168                                WorkflowActionQuery.GET_ACTION_COMPLETED, jobId);
169                    }
170                    catch (JPAExecutorException e) {
171                        if (e.getErrorCode().equals(ErrorCode.E0605)) {
172                            jobItr.remove();
173                        }
174                        else {
175                            LOG.info("Failed to fetch the workflow action: " + jobId, e);
176                        }
177                    }
178                    if (wfAction != null && (wfAction.isComplete() || wfAction.isTerminalWithFailure())) {
179                        try {
180                            updateSLASummary(wfAction.getId(), wfAction.getStartTime(), wfAction.getEndTime());
181                            jobItr.remove();
182                        }
183                        catch (JPAExecutorException e) {
184                            LOG.info("Failed to update SLASummaryBean when purging history set entry for " + jobId, e);
185                        }
186                    }
187                }
188                else if (jobId.contains("-C@")) {
189                    CoordinatorActionBean cAction = null;
190                    try {
191                        cAction = CoordActionQueryExecutor.getInstance().get(CoordActionQuery.GET_COORD_ACTION, jobId);
192                    }
193                    catch (JPAExecutorException e) {
194                        if (e.getErrorCode().equals(ErrorCode.E0605)) {
195                            jobItr.remove();
196                        }
197                        else {
198                            LOG.info("Failed to fetch the coord action: " + jobId, e);
199                        }
200                    }
201                    if (cAction != null && cAction.isTerminalStatus()) {
202                        try {
203                            updateSLASummaryForCoordAction(cAction);
204                            jobItr.remove();
205                        }
206                        catch (JPAExecutorException e) {
207                            XLog.getLog(SLACalculatorMemory.class).info(
208                                    "Failed to update SLASummaryBean when purging history set entry for " + jobId, e);
209                        }
210
211                    }
212                }
213                else if (jobId.endsWith("-C")) {
214                    CoordinatorJobBean cJob = null;
215                    try {
216                        cJob = CoordJobQueryExecutor.getInstance().get(CoordJobQuery.GET_COORD_JOB_STATUS_PARENTID,
217                                jobId);
218                    }
219                    catch (JPAExecutorException e) {
220                        if (e.getErrorCode().equals(ErrorCode.E0604)) {
221                            jobItr.remove();
222                        }
223                        else {
224                            LOG.info("Failed to fetch the coord job: " + jobId, e);
225                        }
226                    }
227                    if (cJob != null && cJob.isTerminalStatus()) {
228                        try {
229                            updateSLASummary(cJob.getId(), cJob.getStartTime(), cJob.getEndTime());
230                            jobItr.remove();
231                        }
232                        catch (JPAExecutorException e) {
233                            LOG.info("Failed to update SLASummaryBean when purging history set entry for " + jobId, e);
234                        }
235
236                    }
237                }
238            }
239        }
240
241        private void updateSLASummary(String id, Date startTime, Date endTime) throws JPAExecutorException {
242            SLASummaryBean sla = SLASummaryQueryExecutor.getInstance().get(SLASummaryQuery.GET_SLA_SUMMARY, id);
243            if (sla != null) {
244                sla.setActualStart(startTime);
245                sla.setActualEnd(endTime);
246                if (startTime != null && endTime != null) {
247                    sla.setActualDuration(endTime.getTime() - startTime.getTime());
248                }
249                sla.setLastModifiedTime(new Date());
250                sla.setEventProcessed(8);
251                SLASummaryQueryExecutor.getInstance().executeUpdate(
252                        SLASummaryQuery.UPDATE_SLA_SUMMARY_FOR_ACTUAL_TIMES, sla);
253            }
254        }
255
256        private void updateSLASummaryForCoordAction(CoordinatorActionBean bean) throws JPAExecutorException {
257            String wrkflowId = bean.getExternalId();
258            if (wrkflowId != null) {
259                WorkflowJobBean wrkflow = WorkflowJobQueryExecutor.getInstance().get(
260                        WorkflowJobQuery.GET_WORKFLOW_START_END_TIME, wrkflowId);
261                if (wrkflow != null) {
262                    updateSLASummary(bean.getId(), wrkflow.getStartTime(), wrkflow.getEndTime());
263                }
264            }
265        }
266    }
267
268    private void loadOnRestart() {
269        boolean isJobModified = false;
270        try {
271            long slaPendingCount = 0;
272            long statusPendingCount = 0;
273            List<SLASummaryBean> summaryBeans = jpaService.execute(new SLASummaryGetRecordsOnRestartJPAExecutor(
274                    modifiedAfter));
275            for (SLASummaryBean summaryBean : summaryBeans) {
276                String jobId = summaryBean.getId();
277                LockToken lock = null;
278                try {
279                    switch (summaryBean.getAppType()) {
280                        case COORDINATOR_ACTION:
281                            isJobModified = processSummaryBeanForCoordAction(summaryBean, jobId);
282                            break;
283                        case WORKFLOW_ACTION:
284                            isJobModified = processSummaryBeanForWorkflowAction(summaryBean, jobId);
285                            break;
286                        case WORKFLOW_JOB:
287                            isJobModified = processSummaryBeanForWorkflowJob(summaryBean, jobId);
288                            break;
289                        default:
290                            break;
291                    }
292                }
293                catch (final JPAExecutorException | IllegalArgumentException e) {
294                    LOG.warn("Failed to process SLA summary bean " + jobId, e);
295                }
296                if (isJobModified) {
297                    try {
298                        boolean update = true;
299                        if (Services.get().get(JobsConcurrencyService.class).isHighlyAvailableMode()) {
300                            lock = Services
301                                    .get()
302                                    .get(MemoryLocksService.class)
303                                    .getWriteLock(
304                                            SLACalcStatus.SLA_ENTITYKEY_PREFIX + jobId,
305                                            Services.get().getConf()
306                                                    .getLong(SLAService.CONF_SLA_CALC_LOCK_TIMEOUT, 5 * 1000));
307                            if (lock == null) {
308                                update = false;
309                            }
310                        }
311                        if (update) {
312                            summaryBean.setLastModifiedTime(new Date());
313                            SLASummaryQueryExecutor.getInstance().executeUpdate(
314                                    SLASummaryQuery.UPDATE_SLA_SUMMARY_FOR_STATUS_ACTUAL_TIMES, summaryBean);
315                        }
316                    }
317                    catch (Exception e) {
318                        LOG.warn("Failed to load records for " + jobId, e);
319                    }
320                    finally {
321                        if (lock != null) {
322                            lock.release();
323                            lock = null;
324                        }
325                    }
326                }
327                try {
328                    if (summaryBean.getEventProcessed() == 7) {
329                        historySet.add(jobId);
330                        statusPendingCount++;
331                    }
332                    else if (summaryBean.getEventProcessed() <= 7) {
333                        SLARegistrationBean slaRegBean = SLARegistrationQueryExecutor.getInstance().get(
334                                SLARegQuery.GET_SLA_REG_ON_RESTART, jobId);
335                        SLACalcStatus slaCalcStatus = new SLACalcStatus(summaryBean, slaRegBean);
336                        putAndIncrement(jobId, slaCalcStatus);
337                        slaPendingCount++;
338                    }
339                }
340                catch (Exception e) {
341                    LOG.warn("Failed to fetch/update records for " + jobId, e);
342                }
343
344            }
345            LOG.info("Loaded SLASummary pendingSLA=" + slaPendingCount + ", pendingStatusUpdate=" + statusPendingCount);
346
347        }
348        catch (Exception e) {
349            LOG.warn("Failed to retrieve SLASummary records on restart", e);
350        }
351    }
352
353    private boolean processSummaryBeanForCoordAction(SLASummaryBean summaryBean, String jobId)
354            throws JPAExecutorException {
355        boolean isJobModified = false;
356        CoordinatorActionBean coordAction = null;
357        coordAction = jpaService.execute(new CoordActionGetForSLAJPAExecutor(jobId));
358        if (!coordAction.getStatusStr().equals(summaryBean.getJobStatus())) {
359            LOG.trace("Coordinator action status is " + coordAction.getStatusStr() + " and summary bean status is "
360                    + summaryBean.getJobStatus());
361            isJobModified = true;
362            summaryBean.setJobStatus(coordAction.getStatusStr());
363            if (coordAction.isTerminalStatus()) {
364                WorkflowJobBean wfJob = jpaService.execute(new WorkflowJobGetForSLAJPAExecutor(coordAction
365                        .getExternalId()));
366                setEndForSLASummaryBean(summaryBean, wfJob.getStartTime(), coordAction.getLastModifiedTime(),
367                        coordAction.getStatusStr());
368            }
369            else if (coordAction.getStatus() != CoordinatorAction.Status.WAITING) {
370                WorkflowJobBean wfJob = jpaService.execute(new WorkflowJobGetForSLAJPAExecutor(coordAction
371                        .getExternalId()));
372                setStartForSLASummaryBean(summaryBean, summaryBean.getEventProcessed(), wfJob.getStartTime());
373            }
374        }
375        return isJobModified;
376    }
377
378    private boolean processSummaryBeanForWorkflowAction(SLASummaryBean summaryBean, String jobId)
379            throws JPAExecutorException {
380        boolean isJobModified = false;
381        WorkflowActionBean wfAction = null;
382        wfAction = jpaService.execute(new WorkflowActionGetForSLAJPAExecutor(jobId));
383        if (!wfAction.getStatusStr().equals(summaryBean.getJobStatus())) {
384            LOG.trace("Workflow action status is " + wfAction.getStatusStr() + "and summary bean status is "
385                    + summaryBean.getJobStatus());
386            isJobModified = true;
387            summaryBean.setJobStatus(wfAction.getStatusStr());
388            if (wfAction.inTerminalState()) {
389                setEndForSLASummaryBean(summaryBean, wfAction.getStartTime(), wfAction.getEndTime(), wfAction.getStatusStr());
390            }
391            else if (wfAction.getStatus() != WorkflowAction.Status.PREP) {
392                setStartForSLASummaryBean(summaryBean, summaryBean.getEventProcessed(), wfAction.getStartTime());
393            }
394        }
395        return isJobModified;
396    }
397
398    private boolean processSummaryBeanForWorkflowJob(SLASummaryBean summaryBean, String jobId)
399            throws JPAExecutorException {
400        boolean isJobModified = false;
401        WorkflowJobBean wfJob = null;
402        wfJob = jpaService.execute(new WorkflowJobGetForSLAJPAExecutor(jobId));
403        if (!wfJob.getStatusStr().equals(summaryBean.getJobStatus())) {
404            LOG.trace("Workflow job status is " + wfJob.getStatusStr() + "and summary bean status is "
405                    + summaryBean.getJobStatus());
406            isJobModified = true;
407            summaryBean.setJobStatus(wfJob.getStatusStr());
408            if (wfJob.inTerminalState()) {
409                setEndForSLASummaryBean(summaryBean, wfJob.getStartTime(), wfJob.getEndTime(), wfJob.getStatusStr());
410            }
411            else if (wfJob.getStatus() != WorkflowJob.Status.PREP) {
412                setStartForSLASummaryBean(summaryBean, summaryBean.getEventProcessed(), wfJob.getStartTime());
413            }
414        }
415        return isJobModified;
416    }
417
418    private void setEndForSLASummaryBean(SLASummaryBean summaryBean, Date startTime, Date endTime, String status) {
419        byte eventProc = summaryBean.getEventProcessed();
420        summaryBean.setEventProcessed(8);
421        summaryBean.setActualStart(startTime);
422        summaryBean.setActualEnd(endTime);
423        long actualDuration = endTime.getTime() - startTime.getTime();
424        summaryBean.setActualDuration(actualDuration);
425        if (eventProc < 4) {
426            if (status.equals(WorkflowJob.Status.SUCCEEDED.name()) || status.equals(WorkflowAction.Status.OK.name())
427                    || status.equals(CoordinatorAction.Status.SUCCEEDED.name())) {
428                if (endTime.getTime() <= summaryBean.getExpectedEnd().getTime()) {
429                    summaryBean.setSLAStatus(SLAStatus.MET);
430                }
431                else {
432                    summaryBean.setSLAStatus(SLAStatus.MISS);
433                }
434            }
435            else {
436                summaryBean.setSLAStatus(SLAStatus.MISS);
437            }
438        }
439
440    }
441
442    private void setStartForSLASummaryBean(SLASummaryBean summaryBean, byte eventProc, Date startTime) {
443        if (((eventProc & 1) == 0)) {
444            eventProc += 1;
445            summaryBean.setEventProcessed(eventProc);
446        }
447        if (summaryBean.getSLAStatus().equals(SLAStatus.NOT_STARTED)) {
448            summaryBean.setSLAStatus(SLAStatus.IN_PROCESS);
449        }
450        summaryBean.setActualStart(startTime);
451    }
452
453    @Override
454    public int size() {
455        return slaMap.size();
456    }
457
458    @Override
459    public SLACalcStatus get(String jobId) throws JPAExecutorException {
460        SLACalcStatus memObj;
461        memObj = slaMap.get(jobId);
462        if (memObj == null && historySet.contains(jobId)) {
463            memObj = new SLACalcStatus(SLASummaryQueryExecutor.getInstance().get(SLASummaryQuery.GET_SLA_SUMMARY, jobId),
464                    SLARegistrationQueryExecutor.getInstance().get(SLARegQuery.GET_SLA_REG_ON_RESTART, jobId));
465        }
466        return memObj;
467    }
468
469    @Override
470    public Iterator<String> iterator() {
471        return slaMap.keySet().iterator();
472    }
473
474    @Override
475    public boolean isEmpty() {
476        return slaMap.isEmpty();
477    }
478
479    @Override
480    public void clear() {
481        final int originalSize = slaMap.size();
482        LOG.trace("Clearing SLA map and SLA history set. [slaMap.size={};slaHistorySet.size={1}]",
483                slaMap.size(),
484                historySet.size());
485        slaMap.clear();
486        historySet.clear();
487        instrumentation.decr(INSTRUMENTATION_GROUP, SLA_MAP, originalSize);
488    }
489
490    /**
491     * Invoked via periodic run, update the SLA for registered jobs.
492     * <p>
493     * Track the number of times the {@link SLACalcStatus} entry has not been processed successfully, and when a preconfigured
494     * {code oozie.sla.service.SLAService.maximum.retry.count} is reached, remove any {@link SLACalculatorMemory#slaMap} entries
495     * that are causing {@code JPAExecutorException}s of certain {@link ErrorCode}s.
496     * @param jobId the workflow or coordinator job or action ID the SLA is tracked against
497     */
498    void updateJobSla(String jobId) throws Exception {
499        SLACalcStatus slaCalc = slaMap.get(jobId);
500        synchronized (slaCalc) {
501            boolean change = false;
502            Object eventProcObj = null;
503            try {
504                // get eventProcessed on DB for validation in HA
505                eventProcObj = ((SLASummaryQueryExecutor) SLASummaryQueryExecutor.getInstance()).getSingleValue(
506                        SLASummaryQuery.GET_SLA_SUMMARY_EVENTPROCESSED, jobId);
507                resetRetryCount(jobId);
508            } catch (final JPAExecutorException e) {
509                if (e.getErrorCode().equals(ErrorCode.E0603)
510                        || e.getErrorCode().equals(ErrorCode.E0604)
511                        || e.getErrorCode().equals(ErrorCode.E0605)) {
512                    LOG.debug("job [{0}] is not in DB, removing from Memory", jobId);
513                    incrementRetryCountAndRemove(jobId);
514                    return;
515                }
516                throw e;
517            }
518            byte eventProc = ((Byte) eventProcObj).byteValue();
519            if (eventProc >= 7) {
520                if (eventProc == 7) {
521                    historySet.add(jobId);
522                }
523                removeAndDecrement(jobId);
524                LOG.trace("Removed Job [{0}] from map as SLA processed", jobId);
525            }
526            else {
527                slaCalc.setEventProcessed(eventProc);
528                SLARegistrationBean reg = slaCalc.getSLARegistrationBean();
529                // calculation w.r.t current time and status
530                if ((eventProc & 1) == 0) { // first bit (start-processed) unset
531                    if (reg.getExpectedStart() != null) {
532                        if (reg.getExpectedStart().getTime() + jobEventLatency < System.currentTimeMillis()) {
533                            confirmWithDB(slaCalc);
534                            eventProc = slaCalc.getEventProcessed();
535                            if (eventProc != 8 && (eventProc & 1) == 0) {
536                                // Some DB exception
537                                slaCalc.setEventStatus(EventStatus.START_MISS);
538                                eventHandler.queueEvent(new SLACalcStatus(slaCalc));
539                                eventProc++;
540                            }
541                            change = true;
542                        }
543                    }
544                    else {
545                        eventProc++; // disable further processing for optional start sla condition
546                        change = true;
547                    }
548                }
549                // check if second bit (duration-processed) is unset
550                if (eventProc != 8 && ((eventProc >> 1) & 1) == 0) {
551                    if (reg.getExpectedDuration() == -1) {
552                        eventProc += 2;
553                        change = true;
554                    }
555                    else if (slaCalc.getActualStart() != null) {
556                        if ((reg.getExpectedDuration() + jobEventLatency) < (System.currentTimeMillis() - slaCalc
557                                .getActualStart().getTime())) {
558                            slaCalc.setEventProcessed(eventProc);
559                            confirmWithDB(slaCalc);
560                            eventProc = slaCalc.getEventProcessed();
561                            if (eventProc != 8 && ((eventProc >> 1) & 1) == 0) {
562                                // Some DB exception
563                                slaCalc.setEventStatus(EventStatus.DURATION_MISS);
564                                eventHandler.queueEvent(new SLACalcStatus(slaCalc));
565                                eventProc += 2;
566                            }
567                            change = true;
568                        }
569                    }
570                }
571                if (eventProc < 4) {
572                    if (reg.getExpectedEnd().getTime() + jobEventLatency < System.currentTimeMillis()) {
573                        slaCalc.setEventProcessed(eventProc);
574                        confirmWithDB(slaCalc);
575                        eventProc = slaCalc.getEventProcessed();
576                        change = true;
577                    }
578                }
579                if (change) {
580                    try {
581                        boolean locked = true;
582                        slaCalc.acquireLock();
583                        locked = slaCalc.isLocked();
584                        if (locked) {
585                            // no more processing, no transfer to history set
586                            if (slaCalc.getEventProcessed() >= 8) {
587                                eventProc = 8;
588                                // Should not be > 8. But to handle any corner cases
589                                slaCalc.setEventProcessed(8);
590                                removeAndDecrement(jobId);
591                            }
592                            else {
593                                slaCalc.setEventProcessed(eventProc);
594                            }
595                            SLASummaryBean slaSummaryBean = new SLASummaryBean();
596                            slaSummaryBean.setId(slaCalc.getId());
597                            slaSummaryBean.setEventProcessed(eventProc);
598                            slaSummaryBean.setSLAStatus(slaCalc.getSLAStatus());
599                            slaSummaryBean.setEventStatus(slaCalc.getEventStatus());
600                            slaSummaryBean.setActualEnd(slaCalc.getActualEnd());
601                            slaSummaryBean.setActualStart(slaCalc.getActualStart());
602                            slaSummaryBean.setActualDuration(slaCalc.getActualDuration());
603                            slaSummaryBean.setJobStatus(slaCalc.getJobStatus());
604                            slaSummaryBean.setLastModifiedTime(new Date());
605                            SLASummaryQueryExecutor.getInstance().executeUpdate(
606                                    SLASummaryQuery.UPDATE_SLA_SUMMARY_FOR_STATUS_ACTUAL_TIMES, slaSummaryBean);
607                            if (eventProc == 7) {
608                                historySet.add(jobId);
609                                removeAndDecrement(jobId);
610                                LOG.trace("Removed Job [{0}] from map after End-processed", jobId);
611                            }
612                        }
613                    }
614                    catch (InterruptedException e) {
615                        incrementRetryCountAndRemove(jobId);
616                        throw new XException(ErrorCode.E0606, slaCalc.getId(), slaCalc.getLockTimeOut());
617                    }
618                    finally {
619                        slaCalc.releaseLock();
620                    }
621                }
622            }
623        }
624    }
625
626    /**
627     * Periodically run by the SLAService worker threads to update SLA status by
628     * iterating through all the jobs in the map
629     */
630    @Override
631    public void updateAllSlaStatus() {
632        LOG.info("Running periodic SLA check");
633        Iterator<String> iterator = slaMap.keySet().iterator();
634        while (iterator.hasNext()) {
635            String jobId = iterator.next();
636            try {
637                LOG.trace("Processing SLA for jobid={0}", jobId);
638                updateJobSla(jobId);
639            }
640            catch (Exception e) {
641                setLogPrefix(jobId);
642                LOG.error("Exception in SLA processing for job [{0}]", jobId, e);
643                LogUtils.clearLogPrefix();
644            }
645        }
646    }
647
648    /**
649     * Register a new job into the map for SLA tracking
650     */
651    @Override
652    public boolean addRegistration(String jobId, SLARegistrationBean reg) throws JPAExecutorException {
653        try {
654            if (slaMap.size() < capacity) {
655                SLACalcStatus slaCalc = new SLACalcStatus(reg);
656                slaCalc.setSLAStatus(SLAStatus.NOT_STARTED);
657                slaCalc.setJobStatus(getJobStatus(reg.getAppType()));
658                List<JsonBean> insertList = new ArrayList<JsonBean>();
659                final SLASummaryBean summaryBean = new SLASummaryBean(slaCalc);
660                final Timestamp currentTime = DateUtils.convertDateToTimestamp(new Date());
661                reg.setCreatedTimestamp(currentTime);
662                summaryBean.setCreatedTimestamp(currentTime);
663                insertList.add(reg);
664                insertList.add(summaryBean);
665                BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(insertList, null, null);
666                putAndIncrement(jobId, slaCalc);
667                LOG.trace("SLA Registration Event - Job:" + jobId);
668                return true;
669            }
670            else {
671                setLogPrefix(reg.getId());
672                LOG.error(
673                        "SLACalculator memory capacity reached. Cannot add or update new SLA Registration entry for job [{0}]",
674                        reg.getId());
675                LogUtils.clearLogPrefix();
676            }
677        }
678        catch (JPAExecutorException jpa) {
679            throw jpa;
680        }
681        return false;
682    }
683
684    private String getJobStatus(AppType appType) {
685        String status = null;
686        switch (appType) {
687            case COORDINATOR_ACTION:
688                status = CoordinatorAction.Status.WAITING.name();
689                break;
690            case WORKFLOW_ACTION:
691                status = WorkflowAction.Status.PREP.name();
692                break;
693            case WORKFLOW_JOB:
694                status = WorkflowJob.Status.PREP.name();
695                break;
696            default:
697                break;
698        }
699        return status;
700    }
701
702    /**
703     * Update job into the map for SLA tracking
704     */
705    @Override
706    public boolean updateRegistration(String jobId, SLARegistrationBean reg) throws JPAExecutorException {
707        try {
708            if (slaMap.size() < capacity) {
709                SLACalcStatus slaCalc = new SLACalcStatus(reg);
710                slaCalc.setSLAStatus(SLAStatus.NOT_STARTED);
711                slaCalc.setJobStatus(getJobStatus(reg.getAppType()));
712                @SuppressWarnings("rawtypes")
713                List<UpdateEntry> updateList = new ArrayList<UpdateEntry>();
714                updateList.add(new UpdateEntry<SLARegQuery>(SLARegQuery.UPDATE_SLA_REG_ALL, reg));
715                updateList.add(new UpdateEntry<SLASummaryQuery>(SLASummaryQuery.UPDATE_SLA_SUMMARY_ALL,
716                        new SLASummaryBean(slaCalc)));
717                BatchQueryExecutor.getInstance().executeBatchInsertUpdateDelete(null, updateList, null);
718                putAndIncrement(jobId, slaCalc);
719                LOG.trace("SLA Registration Event - Job:" + jobId);
720                return true;
721            }
722            else {
723                setLogPrefix(reg.getId());
724                LOG.error(
725                        "SLACalculator memory capacity reached. Cannot add or update new SLA Registration entry for job [{0}]",
726                        reg.getId());
727                LogUtils.clearLogPrefix();
728            }
729        }
730        catch (JPAExecutorException jpa) {
731            throw jpa;
732        }
733        return false;
734    }
735
736    /**
737     * Remove job from being tracked in map
738     */
739    @Override
740    public void removeRegistration(String jobId) {
741        if (!removeAndDecrement(jobId)) {
742            historySet.remove(jobId);
743        }
744    }
745
746    /**
747     * Triggered after receiving Job status change event, update SLA status
748     * accordingly
749     */
750    @Override
751    public boolean addJobStatus(String jobId, String jobStatus, JobEvent.EventStatus jobEventStatus, Date startTime,
752            Date endTime) throws JPAExecutorException, ServiceException {
753        SLACalcStatus slaCalc = slaMap.get(jobId);
754        SLASummaryBean slaInfo = null;
755        boolean hasSla = false;
756        if (slaCalc == null) {
757            if (historySet.contains(jobId)) {
758                slaInfo = SLASummaryQueryExecutor.getInstance().get(SLASummaryQuery.GET_SLA_SUMMARY, jobId);
759                if (slaInfo == null) {
760                    throw new JPAExecutorException(ErrorCode.E0604, jobId);
761                }
762                slaInfo.setJobStatus(jobStatus);
763                slaInfo.setActualStart(startTime);
764                slaInfo.setActualEnd(endTime);
765                if (endTime != null) {
766                    slaInfo.setActualDuration(endTime.getTime() - startTime.getTime());
767                }
768                slaInfo.setEventProcessed(8);
769                historySet.remove(jobId);
770                slaInfo.setLastModifiedTime(new Date());
771                SLASummaryQueryExecutor.getInstance().executeUpdate(
772                        SLASummaryQuery.UPDATE_SLA_SUMMARY_FOR_STATUS_ACTUAL_TIMES, slaInfo);
773                hasSla = true;
774            }
775            else if (Services.get().get(JobsConcurrencyService.class).isHighlyAvailableMode()) {
776                // jobid might not exist in slaMap in HA Setting
777                SLARegistrationBean slaRegBean = SLARegistrationQueryExecutor.getInstance().get(
778                        SLARegQuery.GET_SLA_REG_ALL, jobId);
779                if (slaRegBean != null) { // filter out jobs picked by SLA job event listener
780                                          // but not actually configured for SLA
781                    SLASummaryBean slaSummaryBean = SLASummaryQueryExecutor.getInstance().get(
782                            SLASummaryQuery.GET_SLA_SUMMARY, jobId);
783                    if (slaSummaryBean.getEventProcessed() < 7) {
784                        slaCalc = new SLACalcStatus(slaSummaryBean, slaRegBean);
785                        putAndIncrement(jobId, slaCalc);
786                    }
787                }
788            }
789        }
790        if (slaCalc != null) {
791            synchronized (slaCalc) {
792                try {
793                    // only get ZK lock when multiple servers running
794                    boolean locked = true;
795                    slaCalc.acquireLock();
796                    locked = slaCalc.isLocked();
797                    if (locked) {
798                        // get eventProcessed on DB for validation in HA
799                        Object eventProcObj = ((SLASummaryQueryExecutor) SLASummaryQueryExecutor.getInstance())
800                                .getSingleValue(SLASummaryQuery.GET_SLA_SUMMARY_EVENTPROCESSED, jobId);
801                        byte eventProc = ((Byte) eventProcObj).byteValue();
802                        slaCalc.setEventProcessed(eventProc);
803                        slaCalc.setJobStatus(jobStatus);
804                        switch (jobEventStatus) {
805                            case STARTED:
806                                slaInfo = processJobStartSLA(slaCalc, startTime);
807                                break;
808                            case SUCCESS:
809                                slaInfo = processJobEndSuccessSLA(slaCalc, startTime, endTime);
810                                break;
811                            case FAILURE:
812                                slaInfo = processJobEndFailureSLA(slaCalc, startTime, endTime);
813                                break;
814                            default:
815                                LOG.debug("Unknown Job Status for SLA purpose[{0}]", jobEventStatus);
816                                slaInfo = getSLASummaryBean(slaCalc);
817                        }
818                        if (slaCalc.getEventProcessed() == 7) {
819                            slaInfo.setEventProcessed(8);
820                            slaMap.remove(jobId);
821                        }
822                        slaInfo.setLastModifiedTime(new Date());
823                        SLASummaryQueryExecutor.getInstance().executeUpdate(
824                                SLASummaryQuery.UPDATE_SLA_SUMMARY_FOR_STATUS_ACTUAL_TIMES, slaInfo);
825                        hasSla = true;
826                    }
827                }
828                catch (InterruptedException e) {
829                    throw new ServiceException(ErrorCode.E0606, slaCalc.getEntityKey(), slaCalc.getLockTimeOut());
830                }
831                finally {
832                    slaCalc.releaseLock();
833                }
834            }
835            LOG.trace("SLA Status Event - Job:" + jobId + " Status:" + slaCalc.getSLAStatus());
836        }
837
838        return hasSla;
839    }
840
841    /**
842     * Process SLA for jobs that started running. Also update actual-start time
843     *
844     * @param slaCalc
845     * @param actualStart
846     * @return SLASummaryBean
847     */
848    private SLASummaryBean processJobStartSLA(SLACalcStatus slaCalc, Date actualStart) {
849        slaCalc.setActualStart(actualStart);
850        if (slaCalc.getSLAStatus().equals(SLAStatus.NOT_STARTED)) {
851            slaCalc.setSLAStatus(SLAStatus.IN_PROCESS);
852        }
853        SLARegistrationBean reg = slaCalc.getSLARegistrationBean();
854        Date expecStart = reg.getExpectedStart();
855        byte eventProc = slaCalc.getEventProcessed();
856        // set event proc here
857        if (((eventProc & 1) == 0)) {
858            if (expecStart != null) {
859                if (actualStart.getTime() > expecStart.getTime()) {
860                    slaCalc.setEventStatus(EventStatus.START_MISS);
861                }
862                else {
863                    slaCalc.setEventStatus(EventStatus.START_MET);
864                }
865                eventHandler.queueEvent(new SLACalcStatus(slaCalc));
866            }
867            eventProc += 1;
868            slaCalc.setEventProcessed(eventProc);
869        }
870        return getSLASummaryBean(slaCalc);
871    }
872
873    /**
874     * Process SLA for jobs that ended successfully. Also update actual-start
875     * and end time
876     *
877     * @param slaCalc
878     * @param actualStart
879     * @param actualEnd
880     * @return SLASummaryBean
881     * @throws JPAExecutorException
882     */
883    private SLASummaryBean processJobEndSuccessSLA(SLACalcStatus slaCalc, Date actualStart, Date actualEnd) throws JPAExecutorException {
884        SLARegistrationBean reg = slaCalc.getSLARegistrationBean();
885        slaCalc.setActualStart(actualStart);
886        slaCalc.setActualEnd(actualEnd);
887        long expectedDuration = reg.getExpectedDuration();
888        long actualDuration = actualEnd.getTime() - actualStart.getTime();
889        slaCalc.setActualDuration(actualDuration);
890        //check event proc
891        byte eventProc = slaCalc.getEventProcessed();
892        if (((eventProc >> 1) & 1) == 0) {
893            processDurationSLA(expectedDuration, actualDuration, slaCalc);
894            eventProc += 2;
895            slaCalc.setEventProcessed(eventProc);
896        }
897
898        if (eventProc < 4) {
899            Date expectedEnd = reg.getExpectedEnd();
900            if (actualEnd.getTime() > expectedEnd.getTime()) {
901                slaCalc.setEventStatus(EventStatus.END_MISS);
902                slaCalc.setSLAStatus(SLAStatus.MISS);
903            }
904            else {
905                slaCalc.setEventStatus(EventStatus.END_MET);
906                slaCalc.setSLAStatus(SLAStatus.MET);
907            }
908            eventProc += 4;
909            slaCalc.setEventProcessed(eventProc);
910            eventHandler.queueEvent(new SLACalcStatus(slaCalc));
911        }
912        return getSLASummaryBean(slaCalc);
913    }
914
915    /**
916     * Process SLA for jobs that ended in failure. Also update actual-start and
917     * end time
918     *
919     * @param slaCalc
920     * @param actualStart
921     * @param actualEnd
922     * @return SLASummaryBean
923     * @throws JPAExecutorException
924     */
925    private SLASummaryBean processJobEndFailureSLA(SLACalcStatus slaCalc, Date actualStart, Date actualEnd) throws JPAExecutorException {
926        slaCalc.setActualStart(actualStart);
927        slaCalc.setActualEnd(actualEnd);
928        if (actualStart == null) { // job failed before starting
929            if (slaCalc.getEventProcessed() < 4) {
930                slaCalc.setEventStatus(EventStatus.END_MISS);
931                slaCalc.setSLAStatus(SLAStatus.MISS);
932                eventHandler.queueEvent(new SLACalcStatus(slaCalc));
933                slaCalc.setEventProcessed(7);
934                return getSLASummaryBean(slaCalc);
935            }
936        }
937        SLARegistrationBean reg = slaCalc.getSLARegistrationBean();
938        long expectedDuration = reg.getExpectedDuration();
939        long actualDuration = actualEnd.getTime() - actualStart.getTime();
940        slaCalc.setActualDuration(actualDuration);
941
942        byte eventProc = slaCalc.getEventProcessed();
943        if (((eventProc >> 1) & 1) == 0) {
944            if (expectedDuration != -1) {
945                slaCalc.setEventStatus(EventStatus.DURATION_MISS);
946                eventHandler.queueEvent(new SLACalcStatus(slaCalc));
947            }
948            eventProc += 2;
949            slaCalc.setEventProcessed(eventProc);
950        }
951        if (eventProc < 4) {
952            slaCalc.setEventStatus(EventStatus.END_MISS);
953            slaCalc.setSLAStatus(SLAStatus.MISS);
954            eventProc += 4;
955            slaCalc.setEventProcessed(eventProc);
956            eventHandler.queueEvent(new SLACalcStatus(slaCalc));
957        }
958        return getSLASummaryBean(slaCalc);
959    }
960
961    private SLASummaryBean getSLASummaryBean (SLACalcStatus slaCalc) {
962        SLASummaryBean slaSummaryBean = new SLASummaryBean();
963        slaSummaryBean.setActualStart(slaCalc.getActualStart());
964        slaSummaryBean.setActualEnd(slaCalc.getActualEnd());
965        slaSummaryBean.setActualDuration(slaCalc.getActualDuration());
966        slaSummaryBean.setSLAStatus(slaCalc.getSLAStatus());
967        slaSummaryBean.setEventStatus(slaCalc.getEventStatus());
968        slaSummaryBean.setEventProcessed(slaCalc.getEventProcessed());
969        slaSummaryBean.setId(slaCalc.getId());
970        slaSummaryBean.setJobStatus(slaCalc.getJobStatus());
971        return slaSummaryBean;
972    }
973
974    private void processDurationSLA(long expected, long actual, SLACalcStatus slaCalc) {
975        if (expected != -1 && actual > expected) {
976            slaCalc.setEventStatus(EventStatus.DURATION_MISS);
977            eventHandler.queueEvent(new SLACalcStatus(slaCalc));
978        }
979        else if (expected != -1 && actual <= expected) {
980            slaCalc.setEventStatus(EventStatus.DURATION_MET);
981            eventHandler.queueEvent(new SLACalcStatus(slaCalc));
982        }
983    }
984
985    /*
986     * Confirm alerts against source of truth - DB. Also required in case of High Availability
987     */
988    private void confirmWithDB(SLACalcStatus slaCalc) {
989        boolean ended = false, isEndMiss = false;
990        try {
991            switch (slaCalc.getAppType()) {
992                case WORKFLOW_JOB:
993                    WorkflowJobBean wf = jpaService.execute(new WorkflowJobGetForSLAJPAExecutor(slaCalc.getId()));
994                    if (wf.getEndTime() != null) {
995                        ended = true;
996                        if (wf.getStatus() == WorkflowJob.Status.KILLED || wf.getStatus() == WorkflowJob.Status.FAILED
997                                || wf.getEndTime().getTime() > slaCalc.getExpectedEnd().getTime()) {
998                            isEndMiss = true;
999                        }
1000                    }
1001                    slaCalc.setActualStart(wf.getStartTime());
1002                    slaCalc.setActualEnd(wf.getEndTime());
1003                    slaCalc.setJobStatus(wf.getStatusStr());
1004                    break;
1005                case WORKFLOW_ACTION:
1006                    WorkflowActionBean wa = jpaService.execute(new WorkflowActionGetForSLAJPAExecutor(slaCalc.getId()));
1007                    if (wa.getEndTime() != null) {
1008                        ended = true;
1009                        if (wa.isTerminalWithFailure()
1010                                || wa.getEndTime().getTime() > slaCalc.getExpectedEnd().getTime()) {
1011                            isEndMiss = true;
1012                        }
1013                    }
1014                    slaCalc.setActualStart(wa.getStartTime());
1015                    slaCalc.setActualEnd(wa.getEndTime());
1016                    slaCalc.setJobStatus(wa.getStatusStr());
1017                    break;
1018                case COORDINATOR_ACTION:
1019                    CoordinatorActionBean ca = jpaService.execute(new CoordActionGetForSLAJPAExecutor(slaCalc.getId()));
1020                    if (ca.isTerminalWithFailure()) {
1021                        isEndMiss = ended = true;
1022                        slaCalc.setActualEnd(ca.getLastModifiedTime());
1023                    }
1024                    if (ca.getExternalId() != null) {
1025                        wf = jpaService.execute(new WorkflowJobGetForSLAJPAExecutor(ca.getExternalId()));
1026                        if (wf.getEndTime() != null) {
1027                            ended = true;
1028                            if (wf.getEndTime().getTime() > slaCalc.getExpectedEnd().getTime()) {
1029                                isEndMiss = true;
1030                            }
1031                        }
1032                        slaCalc.setActualEnd(wf.getEndTime());
1033                        slaCalc.setActualStart(wf.getStartTime());
1034                    }
1035                    slaCalc.setJobStatus(ca.getStatusStr());
1036                    break;
1037                default:
1038                    LOG.debug("Unsupported App-type for SLA - " + slaCalc.getAppType());
1039            }
1040
1041            byte eventProc = slaCalc.getEventProcessed();
1042            if (ended) {
1043                if (isEndMiss) {
1044                    slaCalc.setSLAStatus(SLAStatus.MISS);
1045                }
1046                else {
1047                    slaCalc.setSLAStatus(SLAStatus.MET);
1048                }
1049                if (slaCalc.getActualStart() != null) {
1050                    if ((eventProc & 1) == 0) {
1051                        if (slaCalc.getExpectedStart().getTime() < slaCalc.getActualStart().getTime()) {
1052                            slaCalc.setEventStatus(EventStatus.START_MISS);
1053                        }
1054                        else {
1055                            slaCalc.setEventStatus(EventStatus.START_MET);
1056                        }
1057                        eventHandler.queueEvent(new SLACalcStatus(slaCalc));
1058                    }
1059                    slaCalc.setActualDuration(slaCalc.getActualEnd().getTime() - slaCalc.getActualStart().getTime());
1060                    if (((eventProc >> 1) & 1) == 0) {
1061                        processDurationSLA(slaCalc.getExpectedDuration(), slaCalc.getActualDuration(), slaCalc);
1062                    }
1063                }
1064                if (eventProc < 4) {
1065                    if (isEndMiss) {
1066                        slaCalc.setEventStatus(EventStatus.END_MISS);
1067                    }
1068                    else {
1069                        slaCalc.setEventStatus(EventStatus.END_MET);
1070                    }
1071                    eventHandler.queueEvent(new SLACalcStatus(slaCalc));
1072                }
1073                slaCalc.setEventProcessed(8);
1074            }
1075            else {
1076                if (slaCalc.getActualStart() != null) {
1077                    slaCalc.setSLAStatus(SLAStatus.IN_PROCESS);
1078                }
1079                if ((eventProc & 1) == 0) {
1080                    if (slaCalc.getActualStart() != null) {
1081                        if (slaCalc.getExpectedStart().getTime() < slaCalc.getActualStart().getTime()) {
1082                            slaCalc.setEventStatus(EventStatus.START_MISS);
1083                        }
1084                        else {
1085                            slaCalc.setEventStatus(EventStatus.START_MET);
1086                        }
1087                        eventHandler.queueEvent(new SLACalcStatus(slaCalc));
1088                        eventProc++;
1089                    }
1090                    else if (slaCalc.getExpectedStart().getTime() < System.currentTimeMillis()) {
1091                        slaCalc.setEventStatus(EventStatus.START_MISS);
1092                        eventHandler.queueEvent(new SLACalcStatus(slaCalc));
1093                        eventProc++;
1094                    }
1095                }
1096                if (((eventProc >> 1) & 1) == 0 && slaCalc.getActualStart() != null
1097                        && slaCalc.getExpectedDuration() != -1) {
1098                    if (System.currentTimeMillis() - slaCalc.getActualStart().getTime() > slaCalc.getExpectedDuration()) {
1099                        slaCalc.setEventStatus(EventStatus.DURATION_MISS);
1100                        eventHandler.queueEvent(new SLACalcStatus(slaCalc));
1101                        eventProc += 2;
1102                    }
1103                }
1104                if (eventProc < 4 && slaCalc.getExpectedEnd().getTime() < System.currentTimeMillis()) {
1105                    slaCalc.setEventStatus(EventStatus.END_MISS);
1106                    slaCalc.setSLAStatus(SLAStatus.MISS);
1107                    eventHandler.queueEvent(new SLACalcStatus(slaCalc));
1108                    eventProc += 4;
1109                }
1110                slaCalc.setEventProcessed(eventProc);
1111            }
1112        }
1113        catch (Exception e) {
1114            LOG.warn("Error while confirming SLA against DB for jobid= " + slaCalc.getId() + ". Exception is "
1115                    + e.getClass().getName() + ": " + e.getMessage());
1116            if (slaCalc.getEventProcessed() < 4 && slaCalc.getExpectedEnd().getTime() < System.currentTimeMillis()) {
1117                slaCalc.setEventStatus(EventStatus.END_MISS);
1118                slaCalc.setSLAStatus(SLAStatus.MISS);
1119                eventHandler.queueEvent(new SLACalcStatus(slaCalc));
1120                slaCalc.setEventProcessed(slaCalc.getEventProcessed() + 4);
1121            }
1122        }
1123    }
1124
1125    @VisibleForTesting
1126    public boolean isJobIdInSLAMap(String jobId) {
1127        return this.slaMap.containsKey(jobId);
1128    }
1129
1130    @VisibleForTesting
1131    public boolean isJobIdInHistorySet(String jobId) {
1132        return this.historySet.contains(jobId);
1133    }
1134
1135    private void setLogPrefix(String jobId) {
1136        LOG = LogUtils.setLogInfo(LOG, jobId, null, null);
1137    }
1138
1139    private boolean putAndIncrement(final String jobId, final SLACalcStatus newStatus) {
1140        if (slaMap.put(jobId, newStatus) == null) {
1141            LOG.trace("Added a new item to SLA map. [jobId={0}]", jobId);
1142            instrumentation.incr(INSTRUMENTATION_GROUP, SLA_MAP, 1);
1143            return true;
1144        }
1145
1146        LOG.trace("Updated an existing item in SLA map. [jobId={0}]", jobId);
1147        return false;
1148    }
1149
1150    private boolean removeAndDecrement(final String jobId) {
1151        if (slaMap.remove(jobId) != null) {
1152            LOG.trace("Removed an existing item from SLA map. [jobId={0}]", jobId);
1153            instrumentation.decr(INSTRUMENTATION_GROUP, SLA_MAP, 1);
1154            return true;
1155        }
1156
1157        LOG.trace("Tried to remove a non-existing item from SLA map. [jobId={0}]", jobId);
1158        return false;
1159    }
1160
1161    private void resetRetryCount(final String jobId) {
1162        if (slaMap.containsKey(jobId)) {
1163            LOG.debug("Resetting retry count on [{0}]", jobId);
1164            final SLACalcStatus existingStatus = slaMap.get(jobId);
1165            existingStatus.resetRetryCount();
1166            putAndIncrement(jobId, existingStatus);
1167        }
1168    }
1169
1170    private void incrementRetryCountAndRemove(final String jobId) {
1171        LOG.debug("Checking SLA calculator status [{0}] for retry count", jobId);
1172        if (slaMap.containsKey(jobId)) {
1173            final SLACalcStatus existingStatus = slaMap.get(jobId);
1174            if (existingStatus.getRetryCount() < maxRetryCount) {
1175                existingStatus.incrementRetryCount();
1176                LOG.debug("Retrying with SLA calculator status [{0}] retry count [{1}]", jobId, existingStatus.getRetryCount());
1177                putAndIncrement(jobId, existingStatus);
1178            }
1179            else {
1180                LOG.debug("Removing [{0}] from SLA map as maximum retry count reached", jobId);
1181                removeAndDecrement(jobId);
1182            }
1183        }
1184    }
1185}