Skip to content

Commit 0afec01

Browse files
committed
jobs: fix corner cases, add NPE checks
Signed-off-by: Rohit Yadav <[email protected]>
1 parent d5538fb commit 0afec01

3 files changed

Lines changed: 5 additions & 2 deletions

File tree

framework/jobs/src/org/apache/cloudstack/framework/jobs/dao/AsyncJobDaoImpl.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,7 @@ public AsyncJobDaoImpl() {
6464
pendingAsyncJobsSearch.done();
6565

6666
expiringUnfinishedAsyncJobSearch = createSearchBuilder();
67+
expiringUnfinishedAsyncJobSearch.and("jobDispatcher", expiringUnfinishedAsyncJobSearch.entity().getDispatcher(), SearchCriteria.Op.NEQ);
6768
expiringUnfinishedAsyncJobSearch.and("created", expiringUnfinishedAsyncJobSearch.entity().getCreated(), SearchCriteria.Op.LTEQ);
6869
expiringUnfinishedAsyncJobSearch.and("completeMsId", expiringUnfinishedAsyncJobSearch.entity().getCompleteMsid(), SearchCriteria.Op.NULL);
6970
expiringUnfinishedAsyncJobSearch.and("jobStatus", expiringUnfinishedAsyncJobSearch.entity().getStatus(), SearchCriteria.Op.EQ);
@@ -159,6 +160,7 @@ public List<AsyncJobVO> getExpiredJobs(Date cutTime, int limit) {
159160
@Override
160161
public List<AsyncJobVO> getExpiredUnfinishedJobs(Date cutTime, int limit) {
161162
SearchCriteria<AsyncJobVO> sc = expiringUnfinishedAsyncJobSearch.create();
163+
sc.setParameters("jobDispatcher", AsyncJobVO.JOB_DISPATCHER_PSEUDO);
162164
sc.setParameters("created", cutTime);
163165
sc.setParameters("jobStatus", JobInfo.Status.IN_PROGRESS);
164166
Filter filter = new Filter(AsyncJobVO.class, "created", true, 0L, (long)limit);

framework/jobs/src/org/apache/cloudstack/framework/jobs/impl/AsyncJobManagerImpl.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -673,7 +673,7 @@ public boolean waitAndCheck(AsyncJob job, String[] wakeupTopicsOnMessageBus, lon
673673
while (timeoutInMiliseconds < 0 || System.currentTimeMillis() - startTick < timeoutInMiliseconds) {
674674
msgDetector.waitAny(checkIntervalInMilliSeconds);
675675
job = _jobDao.findById(job.getId());
676-
if (job.getStatus().done()) {
676+
if (job != null && job.getStatus().done()) {
677677
return true;
678678
}
679679

framework/jobs/src/org/apache/cloudstack/framework/jobs/impl/SyncQueueManagerImpl.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -142,7 +142,7 @@ public void doInTransactionWithoutResult(TransactionStatus status) {
142142
for(SyncQueueItemVO item : l) {
143143
SyncQueueVO queueVO = _syncQueueDao.findById(item.getQueueId());
144144
SyncQueueItemVO itemVO = _syncQueueItemDao.findById(item.getId());
145-
if(queueReadyToProcess(queueVO) && itemVO.getLastProcessNumber() == null) {
145+
if(queueReadyToProcess(queueVO) && itemVO != null && itemVO.getLastProcessNumber() == null) {
146146
Long processNumber = queueVO.getLastProcessNumber();
147147
if (processNumber == null)
148148
processNumber = new Long(1);
@@ -220,6 +220,7 @@ public void doInTransactionWithoutResult(TransactionStatus status) {
220220
itemVO.setLastProcessTime(null);
221221
_syncQueueItemDao.update(queueItemId, itemVO);
222222

223+
queueVO.setQueueSize(queueVO.getQueueSize() - 1);
223224
queueVO.setLastUpdated(DateUtil.currentGMTTime());
224225
_syncQueueDao.update(queueVO.getId(), queueVO);
225226
}

0 commit comments

Comments
 (0)