From 3ef425a5a452ea4a02d44d6eefb1d0ae144deb11 Mon Sep 17 00:00:00 2001 From: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> Date: Wed, 26 Aug 2026 17:22:35 +0800 Subject: [PATCH] fix(scan): retry reclaimed task lock contention Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> --- .../skillhub/stream/ScanTaskConsumer.java | 12 ++++++- .../skillhub/stream/ScanTaskConsumerTest.java | 32 +++++++++++++++++-- 2 files changed, 41 insertions(+), 3 deletions(-) diff --git a/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/ScanTaskConsumer.java b/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/ScanTaskConsumer.java index f32a2876..3b1b2c66 100644 --- a/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/ScanTaskConsumer.java +++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/ScanTaskConsumer.java @@ -144,7 +144,11 @@ public class ScanTaskConsumer extends AbstractStreamConsumer 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();