|
1 | 1 | package com.springqprobackend.springqpro.service; |
2 | 2 |
|
3 | 3 | import com.springqprobackend.springqpro.domain.entity.TaskEntity; |
| 4 | +import com.springqprobackend.springqpro.domain.event.TaskCreatedEvent; |
4 | 5 | import com.springqprobackend.springqpro.enums.TaskStatus; |
5 | 6 | import com.springqprobackend.springqpro.handlers.TaskHandler; |
6 | 7 | import com.springqprobackend.springqpro.mapper.TaskMapper; |
|
13 | 14 | import jakarta.persistence.PersistenceContext; |
14 | 15 | import org.slf4j.Logger; |
15 | 16 | import org.slf4j.LoggerFactory; |
| 17 | +import org.springframework.beans.factory.annotation.Autowired; |
16 | 18 | import org.springframework.beans.factory.annotation.Qualifier; |
17 | 19 | import org.springframework.context.annotation.Lazy; |
| 20 | +import org.springframework.context.ApplicationEventPublisher; |
18 | 21 | import org.springframework.stereotype.Service; |
19 | 22 | import org.springframework.transaction.annotation.Transactional; |
20 | 23 | import io.micrometer.core.instrument.Counter; |
@@ -127,6 +130,9 @@ public class ProcessingService { |
127 | 130 | private final RedisDistributedLock redisLock; // 2025-11-23-DEBUG: REDIS INTEGRATION PHASE! |
128 | 131 | private final TaskRedisRepository cache; // 2025-11-23-DEBUG: Refactoring for TaskRedisRepository.java |
129 | 132 |
|
| 133 | + @Autowired |
| 134 | + private ApplicationEventPublisher publisher; |
| 135 | + |
130 | 136 | // 2025-11-26-NOTE:+DEBUG: METRICS PHASE FIELD ADDITIONS: |
131 | 137 | private final Counter tasksSubmittedCounter; |
132 | 138 | private final Counter tasksClaimedCounter; |
@@ -264,25 +270,28 @@ public void claimAndProcess(String taskId) { |
264 | 270 | @Transactional |
265 | 271 | public boolean manuallyRequeue(String taskId) { |
266 | 272 | Optional<TaskEntity> opt = taskRepository.findById(taskId); |
267 | | - if(opt.isEmpty()) return false; |
268 | | - |
| 273 | + if(opt.isEmpty()) { |
| 274 | + logger.warn("[ManualRequeue] Task {} not found.", taskId); |
| 275 | + return false; |
| 276 | + } |
269 | 277 | TaskEntity task = opt.get(); |
270 | 278 | // Only FAILED tasks may be manually requeued |
271 | 279 | if (task.getStatus() != TaskStatus.FAILED) { |
272 | 280 | logger.warn("[ManualRequeue] Task {} is not FAILED, ignoring.", taskId); |
273 | 281 | return false; |
274 | 282 | } |
275 | | - // Reset status to QUEUED for the sake of auto-retry: |
| 283 | + // Reset status to QUEUED (and attempts to 0) for the sake of auto-retry: |
276 | 284 | int updated = taskRepository.transitionStatus(taskId, TaskStatus.FAILED, TaskStatus.QUEUED, 0); |
277 | 285 | if(updated == 0) { |
278 | 286 | logger.warn("[ManualRequeue] DB Transition FAILED->QUEUED did not update any rows."); |
279 | 287 | return false; |
280 | 288 | } |
281 | | - // Update the cache: |
282 | | - cache.put(task); |
283 | | - // Push into in-memory queue immediately: |
284 | | - queueService.enqueueById(taskId); |
285 | | - logger.info("[ManualRequeue] Task {} successfully re-enqueued manually.", taskId); |
| 289 | + // Re-read the fresh entity and update the cache: |
| 290 | + TaskEntity refreshed = taskRepository.findById(taskId).orElseThrow(); |
| 291 | + cache.put(refreshed); |
| 292 | + // Publish event - the listener will call enqueueById after the AFTER COMMIT (this is deliberate, can't directly call enqueueById, it's a bad idea): |
| 293 | + publisher.publishEvent(new TaskCreatedEvent(this, taskId)); |
| 294 | + logger.info("[ManualRequeue] Task {} successfully re-enqueued manually (AFTER_COMMIT will enqueue worker).", taskId); |
286 | 295 | return true; |
287 | 296 | } |
288 | 297 |
|
|
0 commit comments