diff --git a/framework/jobs/src/main/java/org/apache/cloudstack/framework/jobs/impl/AsyncJobManagerImpl.java b/framework/jobs/src/main/java/org/apache/cloudstack/framework/jobs/impl/AsyncJobManagerImpl.java index b9c9b22d9eaa..8c3d965f9766 100644 --- a/framework/jobs/src/main/java/org/apache/cloudstack/framework/jobs/impl/AsyncJobManagerImpl.java +++ b/framework/jobs/src/main/java/org/apache/cloudstack/framework/jobs/impl/AsyncJobManagerImpl.java @@ -523,7 +523,7 @@ public String obfuscatePassword(String result, boolean hidePassword) { return StringUtils.obfuscatePasswordInJsonLikeString(result); } - private void scheduleExecution(final AsyncJobVO job) { + protected void scheduleExecution(final AsyncJobVO job) { scheduleExecution(job, false); } @@ -701,58 +701,58 @@ private int getAndResetPendingSignals(AsyncJob job) { return signals; } - private void executeQueueItem(SyncQueueItemVO item, boolean fromPreviousSession) { + protected void executeQueueItem(SyncQueueItemVO item, boolean fromPreviousSession) { AsyncJobVO job = _jobDao.findById(item.getContentId()); - if (job != null) { + if (job == null) { if (logger.isDebugEnabled()) { - logger.debug("Schedule queued job-" + job.getId()); - } - - job.setSyncSource(item); - - // - // TODO: a temporary solution to work-around DB deadlock situation - // - // to live with DB deadlocks, we will give a chance for job to be rescheduled - // in case of exceptions (most-likely DB deadlock exceptions) - try { - job.setExecutingMsid(getMsid()); - _jobDao.update(job.getId(), job); - } catch (Exception e) { - logger.warn("Unexpected exception while dispatching job-" + item.getContentId(), e); - - try { - _queueMgr.returnItem(item.getId()); - } catch (Throwable thr) { - logger.error("Unexpected exception while returning job-" + item.getContentId() + " to queue", thr); - } + logger.debug("Unable to find related job for queue item: " + item.toString()); } + _queueMgr.purgeItem(item.getId()); + return; + } - try { - scheduleExecution(job); - } catch (RejectedExecutionException e) { - logger.warn("Execution for job-" + job.getId() + " is rejected, return it to the queue for next turn"); + if (logger.isDebugEnabled()) { + logger.debug("Schedule queued job-" + job.getId()); + } + job.setSyncSource(item); - try { - _queueMgr.returnItem(item.getId()); - } catch (Exception e2) { - logger.error("Unexpected exception while returning job-" + item.getContentId() + " to queue", e2); - } + // + // TODO: a temporary solution to work-around DB deadlock situation + // + // to live with DB deadlocks, we will give a chance for job to be rescheduled + // in case of exceptions (most-likely DB deadlock exceptions) + try { + job.setExecutingMsid(getMsid()); + _jobDao.update(job.getId(), job); + } catch (Exception e) { + logger.warn("Unexpected exception while dispatching job-" + item.getContentId(), e); + returnItemToQueue(item); + return; + } - try { - job.setExecutingMsid(null); - _jobDao.update(job.getId(), job); - } catch (Exception e3) { - logger.warn("Unexpected exception while update job-" + item.getContentId() + " msid for bookkeeping"); - } - } + try { + scheduleExecution(job); + } catch (RejectedExecutionException e) { + logger.warn("Execution for job-" + job.getId() + " is rejected, return it to the queue for next turn"); + returnItemToQueue(item); + clearExecutingMsid(job, item); + } + } - } else { - if (logger.isDebugEnabled()) { - logger.debug("Unable to find related job for queue item: " + item.toString()); - } + private void returnItemToQueue(SyncQueueItemVO item) { + try { + _queueMgr.returnItem(item.getId()); + } catch (Throwable thr) { + logger.error("Unexpected exception while returning job-" + item.getContentId() + " to queue", thr); + } + } - _queueMgr.purgeItem(item.getId()); + private void clearExecutingMsid(AsyncJobVO job, SyncQueueItemVO item) { + try { + job.setExecutingMsid(null); + _jobDao.update(job.getId(), job); + } catch (Exception e) { + logger.warn("Unexpected exception while update job-" + item.getContentId() + " msid for bookkeeping"); } } diff --git a/framework/jobs/src/test/java/org/apache/cloudstack/framework/jobs/impl/AsyncJobManagerImplExecuteQueueItemTest.java b/framework/jobs/src/test/java/org/apache/cloudstack/framework/jobs/impl/AsyncJobManagerImplExecuteQueueItemTest.java new file mode 100644 index 000000000000..a197d7b4f0a7 --- /dev/null +++ b/framework/jobs/src/test/java/org/apache/cloudstack/framework/jobs/impl/AsyncJobManagerImplExecuteQueueItemTest.java @@ -0,0 +1,70 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +package org.apache.cloudstack.framework.jobs.impl; + +import org.apache.cloudstack.framework.jobs.dao.AsyncJobDao; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.Spy; +import org.mockito.junit.MockitoJUnitRunner; + +import com.cloud.utils.exception.CloudRuntimeException; + +@RunWith(MockitoJUnitRunner.Silent.class) +public class AsyncJobManagerImplExecuteQueueItemTest { + + @Mock + AsyncJobDao _jobDao; + @Mock + SyncQueueManager _queueMgr; + + @Spy + @InjectMocks + AsyncJobManagerImpl asyncJobManager = new AsyncJobManagerImpl(); + + @Test + public void executeQueueItemDoesNotScheduleWhenTheJobUpdateFailsAndItemIsReturned() { + long contentId = 10L; + long itemId = 20L; + long jobId = 1L; + + SyncQueueItemVO item = Mockito.mock(SyncQueueItemVO.class); + Mockito.when(item.getContentId()).thenReturn(contentId); + Mockito.when(item.getId()).thenReturn(itemId); + + AsyncJobVO job = Mockito.mock(AsyncJobVO.class); + Mockito.when(job.getId()).thenReturn(jobId); + Mockito.when(_jobDao.findById(contentId)).thenReturn(job); + + // Simulate the DB deadlock the catch block was written to survive. + Mockito.doThrow(new CloudRuntimeException("simulated DB deadlock")) + .when(_jobDao).update(Mockito.anyLong(), Mockito.any(AsyncJobVO.class)); + + // Stub the executor path so we can assert whether it is reached (and avoid the real submit). + Mockito.doNothing().when(asyncJobManager).scheduleExecution(Mockito.any(AsyncJobVO.class)); + + asyncJobManager.executeQueueItem(item, false); + + // The queue item was returned for a later retry; the job must NOT also be scheduled now, or it + // would run twice (once here and once when the heartbeat re-dequeues the returned item). + Mockito.verify(_queueMgr).returnItem(itemId); + Mockito.verify(asyncJobManager, Mockito.never()).scheduleExecution(Mockito.any(AsyncJobVO.class)); + } +}