Add cluster backplane abstraction and interfaces (#6449)

This commit is contained in:
Anthony Stirling
2026-05-27 11:54:59 +01:00
committed by GitHub
parent 7743720f0a
commit 4564ed5bec
40 changed files with 2167 additions and 247 deletions
@@ -56,4 +56,17 @@ class ArchitectureTest {
.resideInAPackage("stirling.software.saas..");
rule.check(commonClasses);
}
@Test
void clusterInterfacesHaveNoImplementationDependencies() {
ArchRule rule =
noClasses()
.that()
.resideInAPackage("stirling.software.common.cluster..")
.should()
.dependOnClassesThat()
.resideInAnyPackage(
"stirling.software.proprietary..", "stirling.software.saas..");
rule.check(commonClasses);
}
}
@@ -0,0 +1,51 @@
package stirling.software.common.cluster;
import static org.junit.jupiter.api.Assertions.assertEquals;
import java.time.Instant;
import java.util.List;
import java.util.Map;
import org.junit.jupiter.api.Test;
class BackplaneContractCompilationTest {
@Test
void jobStoreEntryRecordRoundTrips() {
Instant now = Instant.now();
JobStoreEntry entry =
new JobStoreEntry(
"job-1",
JobStoreEntry.JobState.PENDING,
"node-a",
now,
null,
null,
List.of("file-1"),
Map.of("k", "v"));
assertEquals("job-1", entry.jobId());
assertEquals(JobStoreEntry.JobState.PENDING, entry.state());
assertEquals("node-a", entry.owningNodeId());
assertEquals(now, entry.createdAt());
assertEquals(List.of("file-1"), entry.fileIds());
assertEquals("v", entry.resultMeta().get("k"));
}
@Test
void clusterNodeRecordRoundTrips() {
Instant heartbeat = Instant.now();
ClusterNode node = new ClusterNode("node-a", "10.0.0.1:8080", heartbeat, "BOTH");
assertEquals("node-a", node.nodeId());
assertEquals("10.0.0.1:8080", node.internalAddress());
assertEquals(heartbeat, node.lastHeartbeat());
assertEquals("BOTH", node.role());
}
@Test
void rateLimitDecisionRecordRoundTrips() {
RateLimitStore.RateLimitDecision d = new RateLimitStore.RateLimitDecision(true, 7, 0L);
assertEquals(true, d.allowed());
assertEquals(7, d.remainingTokens());
assertEquals(0L, d.nanosToWaitForRefill());
}
}
@@ -0,0 +1,65 @@
package stirling.software.common.cluster;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertThrows;
import java.lang.reflect.Method;
import org.junit.jupiter.api.Test;
import stirling.software.common.model.ApplicationProperties;
import stirling.software.common.model.ApplicationProperties.Cluster;
class ClusterConfigValidationTest {
@Test
void validationPassesWhenDisabled() {
ApplicationProperties props = new ApplicationProperties();
ClusterConfig config = new ClusterConfig(props);
assertDoesNotThrow(() -> invokeValidate(config));
}
@Test
void validationFailsWhenValkeyEnabledWithoutUrl() {
ApplicationProperties props = new ApplicationProperties();
Cluster cluster = props.getCluster();
cluster.setEnabled(true);
cluster.setBackplane("valkey");
ClusterConfig config = new ClusterConfig(props);
assertThrows(IllegalStateException.class, () -> invokeValidate(config));
}
@Test
void validationPassesWhenValkeyEnabledWithUrl() {
ApplicationProperties props = new ApplicationProperties();
Cluster cluster = props.getCluster();
cluster.setEnabled(true);
cluster.setBackplane("valkey");
cluster.getValkey().setUrl("redis://localhost:6379");
ClusterConfig config = new ClusterConfig(props);
assertDoesNotThrow(() -> invokeValidate(config));
}
@Test
void validationPassesWhenInProcessEnabled() {
ApplicationProperties props = new ApplicationProperties();
Cluster cluster = props.getCluster();
cluster.setEnabled(true);
cluster.setBackplane("inprocess");
ClusterConfig config = new ClusterConfig(props);
assertDoesNotThrow(() -> invokeValidate(config));
}
private void invokeValidate(ClusterConfig config) throws Exception {
Method m = ClusterConfig.class.getDeclaredMethod("validate");
m.setAccessible(true);
try {
m.invoke(config);
} catch (java.lang.reflect.InvocationTargetException ex) {
if (ex.getCause() instanceof RuntimeException re) {
throw re;
}
throw ex;
}
}
}
@@ -0,0 +1,62 @@
package stirling.software.common.cluster;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import org.junit.jupiter.api.Test;
import stirling.software.common.model.ApplicationProperties;
import stirling.software.common.model.ApplicationProperties.Cluster;
class ClusterPropertiesTest {
@Test
void defaultsAreDisabledAndInprocess() {
Cluster props = new ApplicationProperties().getCluster();
assertFalse(props.isEnabled());
assertEquals("inprocess", props.getBackplane());
assertEquals("local", props.getArtifactStore());
assertEquals(Cluster.NodeRole.BOTH, props.resolvedRole());
assertEquals("", props.getValkey().getUrl());
assertFalse(props.getValkey().getTls().isSkipCertVerification());
assertEquals("both", props.getNode().getRole());
assertEquals("http", props.getNode().getScheme());
assertEquals(5000L, props.getNode().getHeartbeatIntervalMs());
}
@Test
void resolvedRoleParsesCaseInsensitively() {
Cluster props = new ApplicationProperties().getCluster();
props.getNode().setRole("WEB");
assertEquals(Cluster.NodeRole.WEB, props.resolvedRole());
props.getNode().setRole("web");
assertEquals(Cluster.NodeRole.WEB, props.resolvedRole());
props.getNode().setRole("Worker");
assertEquals(Cluster.NodeRole.WORKER, props.resolvedRole());
props.getNode().setRole("garbage");
assertEquals(Cluster.NodeRole.BOTH, props.resolvedRole());
props.getNode().setRole(null);
assertEquals(Cluster.NodeRole.BOTH, props.resolvedRole());
}
@Test
void resolvedNodeIdIsStableAcrossCalls() {
Cluster props = new ApplicationProperties().getCluster();
String first = props.resolvedNodeId();
String second = props.resolvedNodeId();
assertNotNull(first);
assertEquals(first, second);
}
@Test
void resolvedNodeIdHonoursExplicitId() {
Cluster props = new ApplicationProperties().getCluster();
props.getNode().setId("abc");
assertEquals("abc", props.resolvedNodeId());
}
}
@@ -0,0 +1,90 @@
package stirling.software.common.cluster;
import static org.assertj.core.api.Assertions.assertThat;
import org.junit.jupiter.api.Test;
import org.springframework.boot.autoconfigure.context.PropertyPlaceholderAutoConfiguration;
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import stirling.software.common.cluster.inprocess.InProcessClusterConfiguration;
import stirling.software.common.model.ApplicationProperties;
/**
* Verifies the {@link InProcessClusterConfiguration} conditional wiring: in-process beans wire when
* cluster mode is off or {@code backplane=inprocess}, and are skipped when {@code
* backplane=valkey}.
*/
class InProcessConfigurationConditionalTest {
private final ApplicationContextRunner runner =
new ApplicationContextRunner()
.withConfiguration(
org.springframework.boot.autoconfigure.AutoConfigurations.of(
PropertyPlaceholderAutoConfiguration.class))
.withUserConfiguration(
TestAppPropertiesConfig.class,
ClusterConfig.class,
InProcessClusterConfiguration.class);
@Test
void inProcessBeansWireWhenClusterDisabled() {
runner.run(
context ->
assertThat(context)
.hasNotFailed()
.hasSingleBean(ClusterBackplane.class)
.hasSingleBean(JobStore.class)
.hasSingleBean(RateLimitStore.class)
.hasSingleBean(DistributedLock.class)
.hasSingleBean(KeyValueCache.class)
.hasSingleBean(InstanceRegistry.class));
}
@Test
void inProcessBeansWireWhenEnabledWithInProcessBackplane() {
runner.withPropertyValues("cluster.enabled=true", "cluster.backplane=inprocess")
.run(
context ->
assertThat(context)
.hasNotFailed()
.hasSingleBean(ClusterBackplane.class)
.hasSingleBean(JobStore.class)
.hasSingleBean(RateLimitStore.class)
.hasSingleBean(DistributedLock.class)
.hasSingleBean(KeyValueCache.class)
.hasSingleBean(InstanceRegistry.class));
}
@Test
void inProcessBeansSkippedWhenEnabledWithDistributedBackplane() {
runner.withPropertyValues(
"cluster.enabled=true",
"cluster.backplane=valkey",
"cluster.valkey.url=redis://localhost:6379")
.run(
context ->
assertThat(context)
.hasNotFailed()
.doesNotHaveBean(ClusterBackplane.class)
.doesNotHaveBean(JobStore.class)
.doesNotHaveBean(RateLimitStore.class)
.doesNotHaveBean(DistributedLock.class)
.doesNotHaveBean(KeyValueCache.class)
.doesNotHaveBean(InstanceRegistry.class));
}
/**
* Hand-rolled {@link ApplicationProperties} bean: the production class loads YAML at startup
* via a {@code @PostConstruct} hook that isn't appropriate for the slice runner, so we wire a
* defaults-only instance here.
*/
@Configuration
static class TestAppPropertiesConfig {
@Bean
ApplicationProperties applicationProperties() {
return new ApplicationProperties();
}
}
}
@@ -0,0 +1,167 @@
package stirling.software.common.cluster.inprocess;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.time.Duration;
import java.util.Optional;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import org.junit.jupiter.api.Test;
import stirling.software.common.cluster.DistributedLock;
class InProcessDistributedLockTest {
@Test
void acquireReleaseAcquire() {
DistributedLock lock = new InProcessDistributedLock();
DistributedLock.LockHandle h1 = lock.tryAcquire("k", Duration.ofSeconds(30)).orElseThrow();
h1.release();
assertTrue(lock.tryAcquire("k", Duration.ofSeconds(30)).isPresent());
}
@Test
void reentryFromSameThreadFails() {
DistributedLock lock = new InProcessDistributedLock();
DistributedLock.LockHandle h1 = lock.tryAcquire("k", Duration.ofSeconds(30)).orElseThrow();
Optional<DistributedLock.LockHandle> reentry = lock.tryAcquire("k", Duration.ofSeconds(30));
assertFalse(reentry.isPresent(), "in-process lock must be non-reentrant");
h1.release();
// After release, anyone can acquire again.
assertTrue(lock.tryAcquire("k", Duration.ofSeconds(30)).isPresent());
}
@Test
void secondAcquireFromAnotherThreadFails() throws InterruptedException {
DistributedLock lock = new InProcessDistributedLock();
DistributedLock.LockHandle h1 = lock.tryAcquire("k", Duration.ofSeconds(30)).orElseThrow();
CountDownLatch done = new CountDownLatch(1);
AtomicBoolean acquired = new AtomicBoolean(true);
Thread t =
new Thread(
() -> {
Optional<DistributedLock.LockHandle> attempt =
lock.tryAcquire("k", Duration.ofSeconds(30));
acquired.set(attempt.isPresent());
attempt.ifPresent(DistributedLock.LockHandle::release);
done.countDown();
});
t.start();
assertTrue(done.await(2, TimeUnit.SECONDS));
assertFalse(acquired.get());
h1.release();
}
@Test
void leaseExpiryAllowsTakeoverEvenWithoutRelease() throws InterruptedException {
// Acquire with a short lease, never call release, then try to acquire again after the
// lease has elapsed. Matches Redis SET-NX-EX semantics - the second caller gets the lock
// because the first lease auto-expired. 250ms lease + 350ms wait gives CI generous slack.
DistributedLock lock = new InProcessDistributedLock();
DistributedLock.LockHandle h1 = lock.tryAcquire("k", Duration.ofMillis(250)).orElseThrow();
Thread.sleep(350);
Optional<DistributedLock.LockHandle> takeover =
lock.tryAcquire("k", Duration.ofSeconds(30));
assertTrue(
takeover.isPresent(),
"expired lease must release the lock so a new caller can take over");
// Calling release() on the original handle after takeover must be a no-op (token check).
h1.release();
// The takeover holder is still the legitimate owner.
assertFalse(lock.tryAcquire("k", Duration.ofSeconds(30)).isPresent());
takeover.get().release();
}
@Test
void renewExtendsLease() throws InterruptedException {
// Acquire with a short lease, renew it before it expires, then verify the lock is still
// held past the original expiry point. 200ms initial + renew to 2s + wait 350ms.
DistributedLock lock = new InProcessDistributedLock();
DistributedLock.LockHandle h1 = lock.tryAcquire("k", Duration.ofMillis(200)).orElseThrow();
assertTrue(h1.renew(Duration.ofSeconds(2)), "renew on a held lease must succeed");
Thread.sleep(350);
assertFalse(
lock.tryAcquire("k", Duration.ofSeconds(30)).isPresent(),
"renew should have pushed expiry well past the original 200ms");
h1.release();
}
@Test
void renewAfterReleaseFails() {
DistributedLock lock = new InProcessDistributedLock();
DistributedLock.LockHandle h1 = lock.tryAcquire("k", Duration.ofSeconds(30)).orElseThrow();
h1.release();
assertFalse(h1.renew(Duration.ofSeconds(30)), "renew on a released handle must fail");
}
/**
* Concurrency stress: many threads contending on the same key with each holder respecting the
* lease (hold &lt;&lt; lease). The lock behaves as a strict mutex in this regime so asserting
* mutual exclusion is meaningful. A separate test ({@link
* #leaseExpiryAllowsTakeoverEvenWithoutRelease}) covers the takeover-across-expiry branch,
* which legitimately allows two holders momentarily and is split-brain behaviour inherent to
* any lease-based lock.
*/
@Test
void concurrentContentionPreservesMutualExclusion() throws InterruptedException {
DistributedLock lock = new InProcessDistributedLock();
int threads = 16;
int attemptsPerThread = 200;
// Lease far exceeds any plausible hold time, so the takeover branch never triggers in
// this test and the lock acts as a strict mutex.
Duration lease = Duration.ofSeconds(5);
java.util.concurrent.atomic.AtomicInteger concurrentHolders =
new java.util.concurrent.atomic.AtomicInteger();
java.util.concurrent.atomic.AtomicInteger maxConcurrent =
new java.util.concurrent.atomic.AtomicInteger();
java.util.concurrent.atomic.AtomicInteger acquires =
new java.util.concurrent.atomic.AtomicInteger();
java.util.concurrent.atomic.AtomicReference<Throwable> firstFailure =
new java.util.concurrent.atomic.AtomicReference<>();
CountDownLatch start = new CountDownLatch(1);
CountDownLatch done = new CountDownLatch(threads);
for (int i = 0; i < threads; i++) {
new Thread(
() -> {
try {
start.await();
for (int j = 0; j < attemptsPerThread; j++) {
Optional<DistributedLock.LockHandle> h =
lock.tryAcquire("hot", lease);
if (h.isPresent()) {
int now = concurrentHolders.incrementAndGet();
maxConcurrent.accumulateAndGet(now, Math::max);
acquires.incrementAndGet();
// Trivial critical section; well within lease.
concurrentHolders.decrementAndGet();
h.get().release();
}
}
} catch (Throwable t) {
firstFailure.compareAndSet(null, t);
} finally {
done.countDown();
}
},
"lock-stress-" + i)
.start();
}
start.countDown();
assertTrue(done.await(30, TimeUnit.SECONDS), "stress workers must finish in time");
org.junit.jupiter.api.Assertions.assertNull(firstFailure.get(), "no worker may throw");
org.junit.jupiter.api.Assertions.assertEquals(
1,
maxConcurrent.get(),
"mutual exclusion violated: more than one holder observed simultaneously");
assertTrue(
acquires.get() > 0,
"at least some acquires must succeed under contention (saw "
+ acquires.get()
+ ")");
}
}
@@ -0,0 +1,28 @@
package stirling.software.common.cluster.inprocess;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.time.Duration;
import java.time.Instant;
import org.junit.jupiter.api.Test;
import stirling.software.common.cluster.ClusterNode;
class InProcessInstanceRegistryTest {
@Test
void registerThenLookupAndActiveNodes() {
InProcessInstanceRegistry registry = new InProcessInstanceRegistry();
ClusterNode node = new ClusterNode("node-1", "127.0.0.1:8080", Instant.now(), "BOTH");
registry.register(node, Duration.ofSeconds(30));
assertTrue(registry.lookup("node-1").isPresent());
assertEquals("node-1", registry.lookup("node-1").get().nodeId());
assertEquals(1, registry.activeNodes().size());
registry.deregister("node-1");
assertTrue(registry.lookup("node-1").isEmpty());
}
}
@@ -0,0 +1,91 @@
package stirling.software.common.cluster.inprocess;
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 java.time.Duration;
import java.time.Instant;
import java.util.List;
import java.util.Map;
import org.junit.jupiter.api.Test;
import stirling.software.common.cluster.JobStoreEntry;
class InProcessJobStoreTest {
private final InProcessJobStore store = new InProcessJobStore();
@Test
void putGetDeleteExistsRoundTrip() {
JobStoreEntry entry = entry("job-1");
store.put(entry, Duration.ofMinutes(30));
assertTrue(store.exists("job-1"));
assertEquals(entry, store.get("job-1").orElseThrow());
store.delete("job-1");
assertFalse(store.exists("job-1"));
}
@Test
void ttlExpiry() throws InterruptedException {
store.put(entry("job-2"), Duration.ofMillis(50));
Thread.sleep(100);
assertFalse(store.get("job-2").isPresent());
}
@Test
void purgeExpiredRemovesOnlyStaleEntries() throws InterruptedException {
store.put(entry("job-fresh"), Duration.ofMinutes(30));
store.put(entry("job-stale"), Duration.ofMillis(20));
Thread.sleep(80);
int removed = store.purgeExpired();
assertEquals(1, removed);
assertTrue(store.exists("job-fresh"));
}
@Test
void findJobIdByFileIdReturnsTheRightJob() {
store.put(
new JobStoreEntry(
"job-a",
JobStoreEntry.JobState.COMPLETE,
"node-1",
Instant.now(),
Instant.now(),
null,
List.of("file-1", "file-2"),
Map.of()),
Duration.ofMinutes(30));
store.put(
new JobStoreEntry(
"job-b",
JobStoreEntry.JobState.COMPLETE,
"node-1",
Instant.now(),
Instant.now(),
null,
List.of("file-3"),
Map.of()),
Duration.ofMinutes(30));
assertEquals("job-a", store.findJobIdByFileId("file-1").orElseThrow());
assertEquals("job-b", store.findJobIdByFileId("file-3").orElseThrow());
assertFalse(store.findJobIdByFileId("missing").isPresent());
}
private JobStoreEntry entry(String id) {
return new JobStoreEntry(
id,
JobStoreEntry.JobState.PENDING,
"node-1",
Instant.now(),
null,
null,
List.of(),
Map.of());
}
}
@@ -0,0 +1,41 @@
package stirling.software.common.cluster.inprocess;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import java.time.Duration;
import org.junit.jupiter.api.Test;
import stirling.software.common.cluster.KeyValueCache;
class InProcessKeyValueCacheTest {
@Test
void putGetEvict() {
KeyValueCache cache = new InProcessKeyValueCache();
cache.put("apikey", "a", "userA", Duration.ofMinutes(1));
assertEquals("userA", cache.get("apikey", "a").orElseThrow());
cache.evict("apikey", "a");
assertFalse(cache.get("apikey", "a").isPresent());
}
@Test
void ttlExpiry() throws InterruptedException {
KeyValueCache cache = new InProcessKeyValueCache();
cache.put("ns", "k", "v", Duration.ofMillis(40));
Thread.sleep(80);
assertFalse(cache.get("ns", "k").isPresent());
}
@Test
void evictNamespace() {
KeyValueCache cache = new InProcessKeyValueCache();
cache.put("ns", "a", "1", Duration.ofMinutes(1));
cache.put("ns", "b", "2", Duration.ofMinutes(1));
cache.evictNamespace("ns");
assertFalse(cache.get("ns", "a").isPresent());
assertFalse(cache.get("ns", "b").isPresent());
}
}
@@ -0,0 +1,56 @@
package stirling.software.common.cluster.inprocess;
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 java.time.Duration;
import org.junit.jupiter.api.Test;
import stirling.software.common.cluster.RateLimitStore;
import stirling.software.common.cluster.RateLimitStore.RateLimitDecision;
class InProcessRateLimitStoreTest {
@Test
void firstNConsumesAllowed() {
RateLimitStore store = new InProcessRateLimitStore();
for (int i = 0; i < 5; i++) {
assertTrue(store.tryConsume("k", 5, Duration.ofSeconds(60)).allowed(), "i=" + i);
}
assertFalse(store.tryConsume("k", 5, Duration.ofSeconds(60)).allowed());
}
@Test
void remainingTokensDecrements() {
RateLimitStore store = new InProcessRateLimitStore();
RateLimitDecision d1 = store.tryConsume("k", 5, Duration.ofSeconds(60));
RateLimitDecision d2 = store.tryConsume("k", 5, Duration.ofSeconds(60));
assertTrue(d1.allowed());
assertTrue(d2.allowed());
assertEquals(4, d1.remainingTokens());
assertEquals(3, d2.remainingTokens());
}
@Test
void refillRestoresTokens() throws InterruptedException {
RateLimitStore store = new InProcessRateLimitStore();
// Capacity 2 with smooth refill over 100 ms -> ~1 token per 50 ms.
for (int i = 0; i < 2; i++) {
assertTrue(store.tryConsume("k", 2, Duration.ofMillis(100)).allowed());
}
assertFalse(store.tryConsume("k", 2, Duration.ofMillis(100)).allowed());
Thread.sleep(150);
assertTrue(store.tryConsume("k", 2, Duration.ofMillis(100)).allowed());
}
@Test
void deniedConsumeReportsWaitNanos() {
RateLimitStore store = new InProcessRateLimitStore();
assertTrue(store.tryConsume("wait", 1, Duration.ofSeconds(10)).allowed());
RateLimitDecision denied = store.tryConsume("wait", 1, Duration.ofSeconds(10));
assertFalse(denied.allowed());
assertTrue(denied.nanosToWaitForRefill() > 0L);
}
}
@@ -0,0 +1,42 @@
package stirling.software.common.cluster.inprocess;
import static org.junit.jupiter.api.Assertions.assertArrayEquals;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.io.ByteArrayInputStream;
import java.io.IOException;
import java.nio.file.Path;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import stirling.software.common.cluster.FileStore;
class LocalDiskFileStoreTest {
@Test
void storeRetrieveSizeDeleteExistsRoundTrip(@TempDir Path dir) throws IOException {
LocalDiskFileStore store = new LocalDiskFileStore(dir.toString());
byte[] payload = "hello-bytes".getBytes();
FileStore.Stored stored = store.store(new ByteArrayInputStream(payload), "x.txt");
assertEquals(payload.length, stored.size());
assertTrue(store.exists(stored.fileId()));
assertEquals(payload.length, store.size(stored.fileId()));
assertArrayEquals(payload, store.retrieveBytes(stored.fileId()));
assertTrue(store.delete(stored.fileId()));
assertFalse(store.exists(stored.fileId()));
}
@Test
void traversalIdsAreRejected(@TempDir Path dir) {
LocalDiskFileStore store = new LocalDiskFileStore(dir.toString());
assertThrows(IllegalArgumentException.class, () -> store.resolve("../foo"));
assertThrows(IllegalArgumentException.class, () -> store.resolve("a/b"));
assertThrows(IllegalArgumentException.class, () -> store.resolve("a\\b"));
}
}
@@ -0,0 +1,27 @@
package stirling.software.common.service;
import static org.junit.jupiter.api.Assertions.assertArrayEquals;
import static org.mockito.Mockito.mock;
import java.io.IOException;
import java.nio.file.Path;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import stirling.software.common.cluster.inprocess.LocalDiskFileStore;
class FileStorageDelegationTest {
@Test
void storeBytesThenRetrieveBytesRoundTripsThroughFileStore(@TempDir Path tempDir)
throws IOException {
FileStorage fs =
new FileStorage(
mock(FileOrUploadService.class),
new LocalDiskFileStore(tempDir.toString()));
byte[] payload = "round-trip".getBytes();
String id = fs.storeBytes(payload, "x.bin");
assertArrayEquals(payload, fs.retrieveBytes(id));
}
}
@@ -3,6 +3,7 @@ package stirling.software.common.service;
import static org.junit.jupiter.api.Assertions.*;
import static org.mockito.Mockito.*;
import java.io.ByteArrayInputStream;
import java.io.IOException;
import java.io.InputStream;
import java.nio.charset.StandardCharsets;
@@ -13,29 +14,30 @@ import java.util.stream.Stream;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.MockitoAnnotations;
import org.springframework.core.io.ByteArrayResource;
import org.springframework.core.io.Resource;
import org.springframework.http.MediaType;
import org.springframework.test.util.ReflectionTestUtils;
import org.springframework.web.multipart.MultipartFile;
import stirling.software.common.cluster.inprocess.LocalDiskFileStore;
class FileStorageTest {
@TempDir Path tempDir;
@Mock private FileOrUploadService fileOrUploadService;
@InjectMocks private FileStorage fileStorage;
private FileStorage fileStorage;
private MultipartFile mockFile;
@BeforeEach
void setUp() {
void setUp() throws IOException {
MockitoAnnotations.openMocks(this);
ReflectionTestUtils.setField(fileStorage, "tempDirPath", tempDir.toString());
fileStorage =
new FileStorage(fileOrUploadService, new LocalDiskFileStore(tempDir.toString()));
// Create a mock MultipartFile
mockFile = mock(MultipartFile.class);
@@ -47,17 +49,7 @@ class FileStorageTest {
void testStoreFile() throws IOException {
// Arrange
byte[] fileContent = "Test PDF content".getBytes();
when(mockFile.getBytes()).thenReturn(fileContent);
// Set up mock to handle transferTo by writing the file
doAnswer(
invocation -> {
java.io.File file = invocation.getArgument(0);
Files.write(file.toPath(), fileContent);
return null;
})
.when(mockFile)
.transferTo(any(java.io.File.class));
when(mockFile.getInputStream()).thenReturn(new ByteArrayInputStream(fileContent));
// Act
String fileId = fileStorage.storeFile(mockFile);
@@ -65,7 +57,7 @@ class FileStorageTest {
// Assert
assertNotNull(fileId);
assertTrue(Files.exists(tempDir.resolve(fileId)));
verify(mockFile).transferTo(any(java.io.File.class));
assertArrayEquals(fileContent, Files.readAllBytes(tempDir.resolve(fileId)));
}
@Test
@@ -247,11 +239,11 @@ class FileStorageTest {
filesBefore = s.count();
}
// Act + Assert: IOException must propagate out — not be swallowed.
// Act + Assert: IOException must propagate out - not be swallowed.
assertThrows(
IOException.class, () -> fileStorage.storeFromResource(flakyResource, "n.pdf"));
// Assert: no partial file lingers under the storage directory — the finally
// Assert: no partial file lingers under the storage directory - the finally
// branch's deleteIfExists must have cleaned it up.
long filesAfter;
try (Stream<Path> s = Files.list(tempDir)) {
@@ -0,0 +1,120 @@
package stirling.software.common.service;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.verify;
import java.time.LocalDateTime;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.mockito.Mock;
import org.mockito.MockitoAnnotations;
import org.springframework.test.util.ReflectionTestUtils;
import stirling.software.common.cluster.ClusterBackplane;
import stirling.software.common.cluster.JobStoreEntry;
import stirling.software.common.cluster.JobStoreEntry.JobState;
import stirling.software.common.cluster.inprocess.InProcessClusterBackplane;
import stirling.software.common.cluster.inprocess.InProcessJobStore;
import stirling.software.common.model.ApplicationProperties;
import stirling.software.common.model.job.JobResult;
class TaskManagerJobStoreDelegationTest {
@Mock private FileStorage fileStorage;
private InProcessJobStore jobStore;
private ClusterBackplane backplane;
private TaskManager taskManager;
@BeforeEach
void setUp() {
MockitoAnnotations.openMocks(this);
jobStore = spy(new InProcessJobStore());
backplane = new InProcessClusterBackplane(new ApplicationProperties());
taskManager = new TaskManager(fileStorage, jobStore, backplane);
ReflectionTestUtils.setField(taskManager, "jobResultExpiryMinutes", 30);
}
@Test
void createTaskWritesPendingEntry() {
taskManager.createTask("job-1");
JobStoreEntry entry = jobStore.get("job-1").orElseThrow();
assertEquals(JobState.PENDING, entry.state());
assertEquals(backplane.localNodeId(), entry.owningNodeId());
}
@Test
void setCompleteFlipsToComplete() {
taskManager.createTask("job-2");
taskManager.setResult("job-2", "ok");
taskManager.setComplete("job-2");
JobStoreEntry entry = jobStore.get("job-2").orElseThrow();
assertEquals(JobState.COMPLETE, entry.state());
}
@Test
void setErrorFlipsToFailed() {
taskManager.createTask("job-3");
taskManager.setError("job-3", "boom");
JobStoreEntry entry = jobStore.get("job-3").orElseThrow();
assertEquals(JobState.FAILED, entry.state());
assertEquals("boom", entry.error());
}
@Test
void cleanupOldJobsIsNoopWhenBackplaneIsNotInProcess() {
ClusterBackplane mockedValkeyBackplane =
new ClusterBackplane() {
@Override
public boolean isHealthy() {
return true;
}
@Override
public String backplaneType() {
return "valkey";
}
@Override
public String localNodeId() {
return "node-1";
}
@Override
public boolean shouldRunLocalCleanup() {
return false;
}
};
TaskManager tm = new TaskManager(fileStorage, jobStore, mockedValkeyBackplane);
ReflectionTestUtils.setField(tm, "jobResultExpiryMinutes", 30);
tm.createTask("job-4");
tm.setComplete("job-4");
ageJobPastExpiry(tm, "job-4");
tm.cleanupOldJobs();
// cleanup must short-circuit before touching jobStore in cluster mode; the backplane
// TTL owns expiry there. If the gate fired correctly, delete is never called.
verify(jobStore, never()).delete(any());
}
@Test
void cleanupOldJobsDeletesFromJobStoreWhenBackplaneIsInProcess() {
taskManager.createTask("job-5");
taskManager.setComplete("job-5");
ageJobPastExpiry(taskManager, "job-5");
taskManager.cleanupOldJobs();
verify(jobStore).delete("job-5");
}
@SuppressWarnings("unchecked")
private static void ageJobPastExpiry(TaskManager tm, String jobId) {
var jobResults =
(java.util.Map<String, JobResult>) ReflectionTestUtils.getField(tm, "jobResults");
JobResult result = jobResults.get(jobId);
ReflectionTestUtils.setField(result, "completedAt", LocalDateTime.now().minusHours(2));
ReflectionTestUtils.setField(result, "complete", true);
}
}
@@ -1,20 +1,28 @@
package stirling.software.common.service;
import static org.junit.jupiter.api.Assertions.*;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.*;
import java.time.LocalDateTime;
import java.util.Map;
import java.util.Optional;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.MockitoAnnotations;
import org.springframework.http.MediaType;
import org.springframework.test.util.ReflectionTestUtils;
import stirling.software.common.cluster.ClusterBackplane;
import stirling.software.common.cluster.JobStore;
import stirling.software.common.cluster.JobStoreEntry;
import stirling.software.common.cluster.JobStoreEntry.JobState;
import stirling.software.common.model.job.JobResult;
import stirling.software.common.model.job.JobStats;
import stirling.software.common.model.job.ResultFile;
@@ -22,6 +30,8 @@ import stirling.software.common.model.job.ResultFile;
class TaskManagerTest {
@Mock private FileStorage fileStorage;
@Mock private JobStore jobStore;
@Mock private ClusterBackplane clusterBackplane;
@InjectMocks private TaskManager taskManager;
@@ -30,6 +40,10 @@ class TaskManagerTest {
@BeforeEach
void setUp() {
closeable = MockitoAnnotations.openMocks(this);
// Treat the backplane as in-process so cleanupOldJobs is not short-circuited.
lenient().when(clusterBackplane.backplaneType()).thenReturn("inprocess");
lenient().when(clusterBackplane.localNodeId()).thenReturn("test-node");
lenient().when(clusterBackplane.shouldRunLocalCleanup()).thenReturn(true);
ReflectionTestUtils.setField(taskManager, "jobResultExpiryMinutes", 30);
}
@@ -270,6 +284,33 @@ class TaskManagerTest {
verify(fileStorage).deleteFile("file-id");
}
@Test
void testCleanupOldJobs_NoOpWhenBackplaneOwnsExpiry() {
// When the backplane reports it should NOT run local cleanup (e.g. a distributed
// backplane with its own TTL), the cleanup loop must leave local state untouched.
when(clusterBackplane.shouldRunLocalCleanup()).thenReturn(false);
// Seed an old completed job that would normally be removed.
String oldJobId = "old-job-distributed";
taskManager.createTask(oldJobId);
JobResult oldJob = taskManager.getJobResult(oldJobId);
ReflectionTestUtils.setField(oldJob, "completedAt", LocalDateTime.now().minusHours(1));
ReflectionTestUtils.setField(oldJob, "complete", true);
Map<String, JobResult> jobResultsMap =
(Map<String, JobResult>) ReflectionTestUtils.getField(taskManager, "jobResults");
assertNotNull(jobResultsMap);
assertTrue(jobResultsMap.containsKey(oldJobId));
// Act
taskManager.cleanupOldJobs();
// Assert: nothing was removed locally, and no jobStore.delete was issued.
assertTrue(jobResultsMap.containsKey(oldJobId));
verify(jobStore, never()).delete(anyString());
verify(fileStorage, never()).deleteFile(anyString());
}
@Test
void testShutdown() {
// This mainly tests that the shutdown method doesn't throw exceptions
@@ -310,4 +351,33 @@ class TaskManagerTest {
// Assert
assertFalse(result);
}
@Test
void testWriteThroughOnUpdate() {
// Mutating calls must write through to the injected JobStore.
String jobId = "write-through-job";
taskManager.createTask(jobId);
taskManager.setResult(jobId, "done");
ArgumentCaptor<JobStoreEntry> captor = ArgumentCaptor.forClass(JobStoreEntry.class);
verify(jobStore, atLeast(2)).put(captor.capture(), any());
JobStoreEntry last = captor.getValue();
assertEquals(jobId, last.jobId());
assertEquals(JobState.COMPLETE, last.state());
assertEquals("test-node", last.owningNodeId());
}
@Test
void testFindJobKeyByFileId_FallsBackToJobStore() {
// When the file id is not in the local map, TaskManager delegates to JobStore.
String fileId = "remote-file-id";
String expectedJobKey = "remote-job-key";
when(jobStore.findJobIdByFileId(fileId)).thenReturn(Optional.of(expectedJobKey));
String actual = taskManager.findJobKeyByFileId(fileId);
assertEquals(expectedJobKey, actual);
verify(jobStore).findJobIdByFileId(fileId);
}
}