mirror of
https://github.com/iflytek/skillhub.git
synced 2026-08-27 11:14:59 +00:00
fix(scan): retry reclaimed task lock contention
Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com>
This commit is contained in:
parent
15566650f0
commit
3ef425a5a4
2 changed files with 41 additions and 3 deletions
|
|
@ -144,7 +144,11 @@ public class ScanTaskConsumer extends AbstractStreamConsumer<ScanTaskConsumer.Sc
|
|||
log.info("Skipping concurrently processed security scan task: taskId={}, versionId={}",
|
||||
payload.taskId(), payload.versionId());
|
||||
payload.skipCleanup();
|
||||
return;
|
||||
// A normal return is treated as success by AbstractStreamConsumer and ACKs
|
||||
// the Redis entry. Requeue through the common failure path instead, so a
|
||||
// reclaimed duplicate cannot erase the only durable delivery while the active
|
||||
// scanner still owns the task lock.
|
||||
throw new ConcurrentScanInProgressException(payload.taskId());
|
||||
}
|
||||
if (securityScanService.isTaskAlreadyProcessed(payload.taskId())) {
|
||||
return;
|
||||
|
|
@ -165,6 +169,12 @@ public class ScanTaskConsumer extends AbstractStreamConsumer<ScanTaskConsumer.Sc
|
|||
securityScanService.processScanResult(payload.versionId(), payload.scannerType(), response);
|
||||
}
|
||||
|
||||
private static final class ConcurrentScanInProgressException extends RuntimeException {
|
||||
private ConcurrentScanInProgressException(String taskId) {
|
||||
super("Security scan is already in progress: taskId=" + taskId);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void markCompleted(ScanTaskPayload payload) {
|
||||
cleanupTempPath(payload.cleanupPath());
|
||||
|
|
|
|||
|
|
@ -293,8 +293,9 @@ class ScanTaskConsumerTest {
|
|||
"task-inflight", 42L, tempDir.toString(), null, ScannerType.SKILL_SCANNER);
|
||||
|
||||
try {
|
||||
consumer.invokeProcessBusiness(payload);
|
||||
consumer.invokeMarkCompleted(payload);
|
||||
assertThatThrownBy(() -> consumer.invokeProcessBusiness(payload))
|
||||
.isInstanceOf(RuntimeException.class)
|
||||
.hasMessage("Security scan is already in progress: taskId=task-inflight");
|
||||
|
||||
assertThat(securityScanner.lastRequest).isNull();
|
||||
assertThat(skillFile).exists();
|
||||
|
|
@ -305,6 +306,33 @@ class ScanTaskConsumerTest {
|
|||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void handleMessage_whenTaskLockIsHeld_republishesInsteadOfDroppingDelivery() {
|
||||
StubSecurityScanner securityScanner = new StubSecurityScanner();
|
||||
InMemoryScanTaskProducer producer = new InMemoryScanTaskProducer();
|
||||
RLock processingLock = mock(RLock.class);
|
||||
when(processingLock.tryLock()).thenReturn(false);
|
||||
TestableScanTaskConsumer consumer = new TestableScanTaskConsumer(
|
||||
securityScanner,
|
||||
new StubSecurityScanService(),
|
||||
new InMemorySkillVersionRepository(),
|
||||
producer,
|
||||
new InMemoryObjectStorageService(),
|
||||
redissonClient(processingLock)
|
||||
);
|
||||
|
||||
consumer.handleMessage(new StreamMessageId(11, 0), Map.of(
|
||||
"taskId", "task-reclaimed",
|
||||
"versionId", "42",
|
||||
"skillPath", "/tmp/skillhub-scans/42",
|
||||
"scannerType", ScannerType.SKILL_SCANNER.getValue()
|
||||
));
|
||||
|
||||
assertThat(producer.publishedTask.taskId()).isEqualTo("task-reclaimed");
|
||||
assertThat(producer.publishedTask.metadata()).containsEntry("retryCount", "1");
|
||||
verify(consumer.stream).ack("skillhub-scanners", new StreamMessageId(11, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void processBusiness_whenScannerFails_releasesProcessingLock() {
|
||||
StubSecurityScanner securityScanner = new StubSecurityScanner();
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue