新需求更新 同步更新

This commit is contained in:
supernijia
2026-08-06 01:11:54 +08:00
parent 9048bbb7f8
commit 28e7fce11c
112 changed files with 7739 additions and 637 deletions
@@ -0,0 +1,228 @@
package com.nanri.aiimage.modules.task.service;
import com.baomidou.mybatisplus.core.MybatisConfiguration;
import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.baomidou.mybatisplus.core.metadata.TableInfoHelper;
import com.nanri.aiimage.modules.task.mapper.TaskFileJobMapper;
import com.nanri.aiimage.modules.task.model.dto.TaskFileJobDispatchEvent;
import com.nanri.aiimage.modules.task.model.entity.TaskFileJobEntity;
import org.apache.ibatis.builder.MapperBuilderAssistant;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.context.ApplicationEventPublisher;
import java.time.LocalDateTime;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
class TaskFileJobServiceTest {
@Mock private TaskFileJobMapper taskFileJobMapper;
@Mock private ApplicationEventPublisher applicationEventPublisher;
@BeforeAll
static void initializeTableInfo() {
TableInfoHelper.initTableInfo(
new MapperBuilderAssistant(new MybatisConfiguration(), ""),
TaskFileJobEntity.class);
}
@Test
void stuckJobIsRequeuedAndRetryCountIsIncremented() {
TaskFileJobEntity running = runningJob(101L, 3, LocalDateTime.now().minusHours(1));
TaskFileJobEntity pending = runningJob(101L, 4, LocalDateTime.now());
pending.setStatus("PENDING");
when(taskFileJobMapper.selectList(any())).thenReturn(List.of(running));
when(taskFileJobMapper.update(isNull(), any(LambdaUpdateWrapper.class))).thenReturn(1);
when(taskFileJobMapper.selectById(101L)).thenReturn(pending);
TaskFileJobService service = new TaskFileJobService(taskFileJobMapper, applicationEventPublisher);
TaskFileJobService.StuckJobResetResult result = service.resetStuckRunningJobsDetailed(30, 20);
assertEquals(1, result.resetCount());
assertTrue(result.exhaustedJobs().isEmpty());
verify(applicationEventPublisher).publishEvent(any(TaskFileJobDispatchEvent.class));
assertUpdateContains(4, running.getUpdatedAt());
}
@Test
void stuckJobAtLastRetryWaitsForDurableFailureFinalization() {
TaskFileJobEntity running = runningJob(102L, 4, LocalDateTime.now().minusHours(1));
TaskFileJobEntity failed = runningJob(102L, TaskFileJobService.MAX_RETRY_COUNT, LocalDateTime.now());
failed.setStatus("FAILED");
/*
failed.setErrorMessage("文件生成任务运行超时,已达到最大重试次数");
*/
failed.setErrorMessage("result file job timeout");
when(taskFileJobMapper.selectList(any())).thenReturn(List.of(running));
when(taskFileJobMapper.update(isNull(), any(LambdaUpdateWrapper.class))).thenReturn(1);
when(taskFileJobMapper.selectById(102L)).thenReturn(failed);
TaskFileJobService service = new TaskFileJobService(taskFileJobMapper, applicationEventPublisher);
TaskFileJobService.StuckJobResetResult result = service.resetStuckRunningJobsDetailed(30, 20);
assertEquals(0, result.resetCount());
assertEquals(List.of(failed), result.exhaustedJobs());
verify(applicationEventPublisher, never()).publishEvent(any());
assertUpdateContains(TaskFileJobService.MAX_RETRY_COUNT, running.getUpdatedAt());
}
@Test
void pendingFailureFinalizationIsReturnedAgainWithoutRequeueing() {
TaskFileJobEntity pending = runningJob(106L, TaskFileJobService.MAX_RETRY_COUNT, LocalDateTime.now());
pending.setStatus("FAILED");
when(taskFileJobMapper.selectList(any())).thenReturn(List.of(pending));
TaskFileJobService service = new TaskFileJobService(taskFileJobMapper, applicationEventPublisher);
TaskFileJobService.StuckJobResetResult result = service.resetStuckRunningJobsDetailed(30, 20);
assertEquals(0, result.resetCount());
assertEquals(List.of(pending), result.exhaustedJobs());
verify(taskFileJobMapper, never()).update(any(), any());
verify(applicationEventPublisher, never()).publishEvent(any());
}
@Test
void queuedClaimActivationUsesUpdatedAtAsFencingToken() {
TaskFileJobEntity claim = runningJob(107L, 1, LocalDateTime.now().minusMinutes(1));
when(taskFileJobMapper.update(isNull(), any(LambdaUpdateWrapper.class))).thenReturn(1);
TaskFileJobService service = new TaskFileJobService(taskFileJobMapper, applicationEventPublisher);
assertTrue(service.activateRunningClaim(claim));
ArgumentCaptor<LambdaUpdateWrapper<TaskFileJobEntity>> update = updateCaptor();
verify(taskFileJobMapper).update(isNull(), update.capture());
assertTrue(update.getValue().getSqlSegment().contains("updated_at"));
assertTrue(update.getValue().getParamNameValuePairs().containsValue(claim.getUpdatedAt()));
}
@Test
void heartbeatRacePreventsStuckJobReset() {
TaskFileJobEntity running = runningJob(103L, 4, LocalDateTime.now().minusHours(1));
when(taskFileJobMapper.selectList(any())).thenReturn(List.of(running));
when(taskFileJobMapper.update(isNull(), any(LambdaUpdateWrapper.class))).thenReturn(0);
TaskFileJobService service = new TaskFileJobService(taskFileJobMapper, applicationEventPublisher);
TaskFileJobService.StuckJobResetResult result = service.resetStuckRunningJobsDetailed(30, 20);
assertEquals(0, result.resetCount());
assertTrue(result.exhaustedJobs().isEmpty());
verify(taskFileJobMapper, never()).selectById(any());
verify(applicationEventPublisher, never()).publishEvent(any());
assertUpdateContains(TaskFileJobService.MAX_RETRY_COUNT, running.getUpdatedAt());
}
@Test
void exhaustedJobCannotBeClaimedOrRequeued() {
when(taskFileJobMapper.update(isNull(), any(LambdaUpdateWrapper.class))).thenReturn(0);
TaskFileJobService service = new TaskFileJobService(taskFileJobMapper, applicationEventPublisher);
assertFalse(service.markRunning(104L));
assertFalse(service.requeue(104L, "retry"));
ArgumentCaptor<LambdaUpdateWrapper<TaskFileJobEntity>> updates = updateCaptor();
verify(taskFileJobMapper, org.mockito.Mockito.times(2)).update(isNull(), updates.capture());
for (LambdaUpdateWrapper<TaskFileJobEntity> update : updates.getAllValues()) {
assertTrue(update.getSqlSegment().contains("retry_count"));
assertTrue(update.getParamNameValuePairs().containsValue(TaskFileJobService.MAX_RETRY_COUNT));
}
verify(taskFileJobMapper, never()).selectById(any());
verify(applicationEventPublisher, never()).publishEvent(any());
}
@Test
void markFailedUsesCurrentRetryCountWithCompareAndSet() {
TaskFileJobEntity stale = runningJob(105L, 0, LocalDateTime.now().minusMinutes(5));
TaskFileJobEntity current = runningJob(105L, 4, LocalDateTime.now());
when(taskFileJobMapper.selectById(105L)).thenReturn(current);
when(taskFileJobMapper.update(isNull(), any(LambdaUpdateWrapper.class))).thenReturn(1);
TaskFileJobService service = new TaskFileJobService(taskFileJobMapper, applicationEventPublisher);
service.markFailed(stale, "failed");
ArgumentCaptor<LambdaUpdateWrapper<TaskFileJobEntity>> update = updateCaptor();
verify(taskFileJobMapper).update(isNull(), update.capture());
assertTrue(update.getValue().getSqlSegment().contains("retry_count"));
assertTrue(update.getValue().getParamNameValuePairs().containsValue(4));
assertTrue(update.getValue().getParamNameValuePairs().containsValue(TaskFileJobService.MAX_RETRY_COUNT));
}
@Test
void successWriteIsFencedToRunningJob() {
TaskFileJobEntity job = runningJob(108L, 1, LocalDateTime.now());
TaskFileJobService service = new TaskFileJobService(taskFileJobMapper, applicationEventPublisher);
service.markSuccess(job, "result.xlsx");
ArgumentCaptor<LambdaUpdateWrapper<TaskFileJobEntity>> update = updateCaptor();
verify(taskFileJobMapper).update(isNull(), update.capture());
assertTrue(update.getValue().getSqlSegment().contains("status"));
assertTrue(update.getValue().getParamNameValuePairs().containsValue("RUNNING"));
}
@Test
void similarAsinHeartbeatOnlyTouchesStaleRunningAssembleJobs() {
TaskFileJobService service = new TaskFileJobService(taskFileJobMapper, applicationEventPublisher);
service.touchRunningAssembleJobsIfStale(20553L, "SIMILAR_ASIN", 60000L);
ArgumentCaptor<LambdaUpdateWrapper<TaskFileJobEntity>> update = updateCaptor();
verify(taskFileJobMapper).update(isNull(), update.capture());
assertTrue(update.getValue().getSqlSegment().contains("task_id"));
assertTrue(update.getValue().getParamNameValuePairs().containsValue("SIMILAR_ASIN"));
assertTrue(update.getValue().getParamNameValuePairs().containsValue("ASSEMBLE_RESULT"));
assertTrue(update.getValue().getParamNameValuePairs().containsValue("RUNNING"));
}
@Test
void retryExhaustedOnlyMeansFailedTerminalJob() {
TaskFileJobEntity success = runningJob(109L, TaskFileJobService.MAX_RETRY_COUNT, LocalDateTime.now());
success.setStatus("SUCCESS");
TaskFileJobEntity failed = runningJob(110L, TaskFileJobService.MAX_RETRY_COUNT, LocalDateTime.now());
failed.setStatus("FAILED");
when(taskFileJobMapper.selectById(109L)).thenReturn(success);
when(taskFileJobMapper.selectById(110L)).thenReturn(failed);
TaskFileJobService service = new TaskFileJobService(taskFileJobMapper, applicationEventPublisher);
assertFalse(service.isRetryExhausted(109L));
assertTrue(service.isRetryExhausted(110L));
}
private void assertUpdateContains(int nextRetry, LocalDateTime updatedAt) {
ArgumentCaptor<LambdaUpdateWrapper<TaskFileJobEntity>> update = updateCaptor();
verify(taskFileJobMapper).update(isNull(), update.capture());
assertTrue(update.getValue().getSqlSegment().contains("updated_at"));
assertTrue(update.getValue().getParamNameValuePairs().containsValue(updatedAt));
assertTrue(update.getValue().getParamNameValuePairs().containsValue(nextRetry));
}
@SuppressWarnings({"rawtypes", "unchecked"})
private static ArgumentCaptor<LambdaUpdateWrapper<TaskFileJobEntity>> updateCaptor() {
return ArgumentCaptor.forClass((Class) LambdaUpdateWrapper.class);
}
private static TaskFileJobEntity runningJob(Long id, int retryCount, LocalDateTime updatedAt) {
TaskFileJobEntity job = new TaskFileJobEntity();
job.setId(id);
job.setTaskId(20553L);
job.setResultId(23110L);
job.setModuleType("SIMILAR_ASIN");
job.setStatus("RUNNING");
job.setRetryCount(retryCount);
job.setUpdatedAt(updatedAt);
return job;
}
}
@@ -8,6 +8,7 @@ import com.nanri.aiimage.config.SimilarAsinProperties;
import com.nanri.aiimage.modules.appearancepatent.service.AppearancePatentTaskCacheService;
import com.nanri.aiimage.modules.brand.mapper.BrandCrawlTaskMapper;
import com.nanri.aiimage.modules.brand.service.BrandTaskProgressCacheService;
import com.nanri.aiimage.modules.collectdata.service.CollectDataService;
import com.nanri.aiimage.modules.deletebrand.service.DeleteBrandTaskCacheService;
import com.nanri.aiimage.modules.patroldelete.service.PatrolDeleteTaskCacheService;
import com.nanri.aiimage.modules.pricetrack.service.PriceTrackTaskCacheService;
@@ -67,8 +68,10 @@ class TaskHeartbeatServiceTest {
@Mock private AppearancePatentTaskCacheService appearancePatentTaskCacheService;
@Mock private SimilarAsinTaskCacheService similarAsinTaskCacheService;
@Mock private SimilarAsinProperties similarAsinProperties;
@Mock private TaskFileJobService taskFileJobService;
@Mock private DeleteBrandTaskCacheService deleteBrandTaskCacheService;
@Mock private BrandTaskProgressCacheService brandTaskProgressCacheService;
@Mock private CollectDataService collectDataService;
@InjectMocks private TaskHeartbeatService service;
@@ -124,6 +127,28 @@ class TaskHeartbeatServiceTest {
verify(shopDataCrawlTaskCacheService).saveTaskCache(task);
}
@Test
@SuppressWarnings("unchecked")
void collectDataHeartbeatForwardsProcessedKeywordProgress() {
long taskId = 21016L;
FileTaskEntity task = new FileTaskEntity();
task.setId(taskId);
task.setModuleType("COLLECT_DATA");
task.setStatus("RUNNING");
TaskHeartbeatRequest request = new TaskHeartbeatRequest();
request.setCurrent(4);
request.setTotal(19);
when(fileTaskMapper.selectOne(any(LambdaQueryWrapper.class))).thenReturn(task);
when(brandCrawlTaskMapper.selectOne(any(LambdaQueryWrapper.class))).thenReturn(null);
when(fileTaskMapper.update(isNull(), any(LambdaUpdateWrapper.class))).thenReturn(1);
TaskHeartbeatVo result = service.heartbeat(taskId, request);
assertTrue(result.isAlive());
verify(collectDataService).updateProgress(taskId, request);
}
@Test
@SuppressWarnings("unchecked")
void similarAsinHeartbeatUsesRedisWithoutRefreshingRecentDatabaseCheckpoint() {
@@ -142,6 +167,7 @@ class TaskHeartbeatServiceTest {
assertTrue(result.isAlive());
assertEquals("SIMILAR_ASIN", result.getModuleType());
verify(similarAsinTaskCacheService).touchTaskHeartbeat(taskId);
verify(taskFileJobService).touchRunningAssembleJobsIfStale(taskId, "SIMILAR_ASIN", 120000L);
verify(fileTaskMapper, never()).update(isNull(), any(LambdaUpdateWrapper.class));
}
}
@@ -17,6 +17,10 @@ import com.nanri.aiimage.modules.task.mapper.FileResultMapper;
import com.nanri.aiimage.modules.task.model.entity.FileResultEntity;
import com.nanri.aiimage.modules.task.model.entity.TaskFileJobEntity;
import com.nanri.aiimage.modules.withdraw.service.WithdrawTaskService;
import java.time.LocalDateTime;
import java.util.List;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.InOrder;
@@ -24,9 +28,13 @@ import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.inOrder;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
@@ -72,7 +80,7 @@ class TaskResultFileJobWorkerTest {
result.setResultFileUrl("result/withdraw/20140.xlsx");
TaskDistributedLockService.LockHandle lock = mock(TaskDistributedLockService.LockHandle.class);
when(taskFileJobService.markRunning(jobId)).thenReturn(true);
allowClaim(job);
when(taskDistributedLockService.acquire("WITHDRAW", taskId, TaskDistributedLockService.DEFAULT_WAIT_MILLIS))
.thenReturn(lock);
when(fileResultMapper.selectById(resultId)).thenReturn(result);
@@ -102,7 +110,7 @@ class TaskResultFileJobWorkerTest {
result.setResultFileUrl("result/publish/20141.xlsx");
TaskDistributedLockService.LockHandle lock = mock(TaskDistributedLockService.LockHandle.class);
when(taskFileJobService.markRunning(jobId)).thenReturn(true);
allowClaim(job);
when(instanceMetadata.getInstanceId()).thenReturn("instance-a");
when(taskDistributedLockService.acquire(
PublishTaskService.MODULE_TYPE,
@@ -150,7 +158,7 @@ class TaskResultFileJobWorkerTest {
TaskDistributedLockService.LockHandle lock = mock(TaskDistributedLockService.LockHandle.class);
when(instanceMetadata.getInstanceId()).thenReturn("instance-a");
when(taskFileJobService.markRunning(jobId)).thenReturn(true);
allowClaim(job);
when(taskDistributedLockService.acquire("SHOP_DATA_CRAWL", taskId,
TaskDistributedLockService.DEFAULT_WAIT_MILLIS)).thenReturn(lock);
when(fileResultMapper.selectById(resultId)).thenReturn(result);
@@ -173,7 +181,7 @@ class TaskResultFileJobWorkerTest {
job.setScopeKey("task:20144:owner:instance-a");
TaskDistributedLockService.LockHandle lock = mock(TaskDistributedLockService.LockHandle.class);
when(instanceMetadata.getInstanceId()).thenReturn("instance-a");
when(taskFileJobService.markRunning(job.getId())).thenReturn(true);
allowClaim(job);
when(taskDistributedLockService.acquire("SHOP_DATA_CRAWL", job.getTaskId(),
TaskDistributedLockService.DEFAULT_WAIT_MILLIS)).thenReturn(lock);
doThrow(new IllegalStateException("upload failed"))
@@ -184,5 +192,86 @@ class TaskResultFileJobWorkerTest {
verify(taskFileJobService).markFailed(job, "upload failed");
verify(shopDataCrawlTaskService).handleResultFileJobFailure(job, "upload failed");
verify(taskFileJobService).markFailureFinalized(job.getId(), "upload failed");
}
@Test
void stuckSimilarAsinJobAtRetryLimitFailsOwningTask() {
TaskFileJobEntity job = new TaskFileJobEntity();
job.setId(13858L);
job.setTaskId(20553L);
job.setResultId(23110L);
job.setModuleType("SIMILAR_ASIN");
job.setRetryCount(TaskFileJobService.MAX_RETRY_COUNT);
job.setErrorMessage("文件生成任务运行超时,已达到最大重试次数");
TaskFileJobService.StuckJobResetResult resetResult =
new TaskFileJobService.StuckJobResetResult(0, List.of(job));
when(taskFileJobService.resetStuckRunningJobsDetailed(0, 0)).thenReturn(resetResult);
worker.resetStuckJobs();
verify(similarAsinTaskService).handleResultFileJobFailure(job, job.getErrorMessage());
verify(taskFileJobService).markFailureFinalized(job.getId(), job.getErrorMessage());
}
@Test
void stuckJobFailureCallbackDoesNotBlockRemainingJobs() {
TaskFileJobEntity first = exhaustedSimilarAsinJob(13858L, 20553L);
TaskFileJobEntity second = exhaustedSimilarAsinJob(13859L, 20554L);
TaskFileJobService.StuckJobResetResult resetResult =
new TaskFileJobService.StuckJobResetResult(0, List.of(first, second));
when(taskFileJobService.resetStuckRunningJobsDetailed(0, 0)).thenReturn(resetResult);
doThrow(new IllegalStateException("owner mismatch"))
.doNothing()
.when(similarAsinTaskService)
.handleResultFileJobFailure(any(), anyString());
worker.resetStuckJobs();
verify(similarAsinTaskService).handleResultFileJobFailure(first, first.getErrorMessage());
verify(similarAsinTaskService).handleResultFileJobFailure(second, second.getErrorMessage());
verify(taskFileJobService, never()).markFailureFinalized(first.getId(), first.getErrorMessage());
verify(taskFileJobService).markFailureFinalized(second.getId(), second.getErrorMessage());
}
@Test
void pendingFailureCallbackIsRetriedOnNextScan() {
TaskFileJobEntity job = exhaustedSimilarAsinJob(13860L, 20555L);
TaskFileJobService.StuckJobResetResult resetResult =
new TaskFileJobService.StuckJobResetResult(0, List.of(job));
when(taskFileJobService.resetStuckRunningJobsDetailed(0, 0))
.thenReturn(resetResult, resetResult);
doThrow(new IllegalStateException("temporary database failure"))
.doNothing()
.when(similarAsinTaskService)
.handleResultFileJobFailure(job, job.getErrorMessage());
worker.resetStuckJobs();
worker.resetStuckJobs();
verify(similarAsinTaskService, times(2))
.handleResultFileJobFailure(job, job.getErrorMessage());
verify(taskFileJobService).markFailureFinalized(job.getId(), job.getErrorMessage());
}
private void allowClaim(TaskFileJobEntity job) {
TaskFileJobEntity claim = new TaskFileJobEntity();
claim.setId(job.getId());
claim.setTaskId(job.getTaskId());
claim.setModuleType(job.getModuleType());
claim.setStatus("RUNNING");
claim.setUpdatedAt(LocalDateTime.now());
when(taskFileJobService.claimRunning(job.getId())).thenReturn(claim);
when(taskFileJobService.activateRunningClaim(claim)).thenReturn(true);
}
private static TaskFileJobEntity exhaustedSimilarAsinJob(long jobId, long taskId) {
TaskFileJobEntity job = new TaskFileJobEntity();
job.setId(jobId);
job.setTaskId(taskId);
job.setModuleType("SIMILAR_ASIN");
job.setRetryCount(TaskFileJobService.MAX_RETRY_COUNT);
job.setErrorMessage("文件生成任务运行超时,已达到最大重试次数");
return job;
}
}