fix(security): 扫描任务改为事务提交后发布 (Closes #612) (#733)

* fix(security): publish scan task after transaction commit

SecurityScanService.triggerScan is @Transactional but published the Redis
Stream scan task inline, before the transaction committed. The stream
consumer could receive the task before the skill_version / security_audit
rows were visible, fail with "SkillVersion not found" / "SecurityAudit not
found", exhaust its immediate retries while the publishing transaction was
still open, and leave the committed version stuck in SCANNING.

Defer the publish to an afterCommit transaction synchronization so the
consumer only ever sees the task once the rows are committed and visible; on
rollback the task is never published. Falls back to an inline publish when
called outside a transaction.

Closes #612

Signed-off-by: FenjuFu <fufenjupku@gmail.com>

* test(security): cover scan task after-commit publishing

Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com>

* refactor(security): hide scan publish transaction callback

Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com>

---------

Signed-off-by: FenjuFu <fufenjupku@gmail.com>
Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com>
Co-authored-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com>
This commit is contained in:
FenjuFu 2026-08-21 14:18:40 +08:00 committed by GitHub
parent bbdc0f7a0c
commit 51457bfa2c
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 103 additions and 3 deletions

View file

@ -75,7 +75,7 @@ public class SecurityScanService {
}
// Always create a new audit record supports multiple rounds per version
auditRepository.save(new SecurityAudit(versionId, ScannerType.SKILL_SCANNER));
scanTaskProducer.publishScanTask(new ScanTask(
final ScanTask scanTask = new ScanTask(
UUID.randomUUID().toString(),
versionId,
packagePath,
@ -83,7 +83,10 @@ public class SecurityScanService {
publisherId,
System.currentTimeMillis(),
Map.of("scannerType", ScannerType.SKILL_SCANNER.getValue())
));
);
// The stream consumer must not observe this task before skill_version /
// security_audit rows are committed and visible.
TransactionCommitCallbacks.afterCommitOrNow(() -> scanTaskProducer.publishScanTask(scanTask));
// Only transition to SCANNING if the version is not already published (auto-publish flow)
if (version.getStatus() != SkillVersionStatus.PUBLISHED) {
version.setStatus(SkillVersionStatus.SCANNING);

View file

@ -0,0 +1,31 @@
package com.iflytek.skillhub.domain.security;
import org.springframework.transaction.support.TransactionSynchronization;
import org.springframework.transaction.support.TransactionSynchronizationManager;
/**
* Keeps Spring transaction callback plumbing out of domain workflows.
*/
final class TransactionCommitCallbacks {
private TransactionCommitCallbacks() {
}
/**
* Runs the callback after the current transaction commits, or immediately when no transaction
* synchronization is active.
*/
static void afterCommitOrNow(Runnable callback) {
if (!TransactionSynchronizationManager.isSynchronizationActive()) {
callback.run();
return;
}
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
@Override
public void afterCommit() {
callback.run();
}
});
}
}

View file

@ -5,21 +5,26 @@ import com.iflytek.skillhub.domain.skill.SkillVersion;
import com.iflytek.skillhub.domain.skill.SkillVersionRepository;
import com.iflytek.skillhub.domain.skill.SkillVersionStatus;
import com.iflytek.skillhub.domain.skill.validation.PackageEntry;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
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.transaction.support.TransactionSynchronization;
import org.springframework.transaction.support.TransactionSynchronizationManager;
import java.lang.reflect.Field;
import java.nio.file.Path;
import java.util.List;
import java.util.Optional;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.BDDMockito.given;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
@ExtendWith(MockitoExtension.class)
@ -48,6 +53,13 @@ class SecurityScanServiceTest {
);
}
@AfterEach
void clearTransactionSynchronization() {
if (TransactionSynchronizationManager.isSynchronizationActive()) {
TransactionSynchronizationManager.clearSynchronization();
}
}
@Test
void securityAudit_startsWithSuspiciousUnsafeDefaults() {
SecurityAudit audit = new SecurityAudit(42L, ScannerType.SKILL_SCANNER);
@ -124,6 +136,54 @@ class SecurityScanServiceTest {
assertThat(task.bundleKey()).isEqualTo("packages/8/42/bundle.zip");
}
@Test
void triggerScan_defersTaskPublishingUntilTransactionCommit() throws Exception {
SkillVersion version = new SkillVersion(8L, "1.0.0", "publisher-1");
setId(version, 42L);
PackageEntry entry = new PackageEntry(
"README.md",
"# demo".getBytes(),
6L,
"text/markdown"
);
given(skillVersionRepository.findById(42L)).willReturn(Optional.of(version));
TransactionSynchronizationManager.initSynchronization();
service.triggerScan(42L, List.of(entry), "publisher-1");
verify(auditRepository).save(any(SecurityAudit.class));
verify(skillVersionRepository).save(version);
verify(scanTaskProducer, never()).publishScanTask(any(ScanTask.class));
assertThat(version.getStatus()).isEqualTo(SkillVersionStatus.SCANNING);
commitRegisteredSynchronizations();
ArgumentCaptor<ScanTask> taskCaptor = ArgumentCaptor.forClass(ScanTask.class);
verify(scanTaskProducer).publishScanTask(taskCaptor.capture());
assertThat(taskCaptor.getValue().versionId()).isEqualTo(42L);
}
@Test
void triggerScan_doesNotPublishTaskWhenTransactionNeverCommits() throws Exception {
SkillVersion version = new SkillVersion(8L, "1.0.0", "publisher-1");
setId(version, 42L);
PackageEntry entry = new PackageEntry(
"README.md",
"# demo".getBytes(),
6L,
"text/markdown"
);
given(skillVersionRepository.findById(42L)).willReturn(Optional.of(version));
TransactionSynchronizationManager.initSynchronization();
service.triggerScan(42L, List.of(entry), "publisher-1");
verify(auditRepository).save(any(SecurityAudit.class));
verify(scanTaskProducer, never()).publishScanTask(any(ScanTask.class));
}
@Test
void triggerScan_rejectsDirectoryTraversalEntries() throws Exception {
SkillVersion version = new SkillVersion(8L, "1.0.0", "publisher-1");
@ -264,4 +324,10 @@ class SecurityScanServiceTest {
field.setAccessible(true);
field.set(target, id);
}
private void commitRegisteredSynchronizations() {
for (TransactionSynchronization synchronization : TransactionSynchronizationManager.getSynchronizations()) {
synchronization.afterCommit();
}
}
}