Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,7 @@ Contains shared:
### ReqShieldConfiguration Parameters
- `isLocalLock`: Use local vs distributed locking (default: true)
- `globalLockFunction` / `globalUnLockFunction`: `(lockKey, token, ttlMillis) -> Boolean` / `(lockKey, token) -> Boolean`, required when `isLocalLock = false`
- `executor` (core) / `scheduler` (reactor) / `scope` (coroutine): where background cache writes run; defaults are shared, the Spring adapters expose them as `reqShieldExecutor` / `reqShieldScheduler` / `reqShieldCoroutineScope` beans
- `executor` (core) / `scheduler` (reactor) / `scope` (coroutine): where background cache writes run; defaults are shared. The Spring adapters register none of them as beans (an `Executor` bean would make Spring Boot drop its `applicationTaskExecutor`, and any library bean would clash with an application bean of the same name): each aspect uses an application bean named `reqShieldExecutor` / `reqShieldScheduler` / `reqShieldCoroutineScope` if one exists (a bean of that name with another type fails the context refresh), otherwise its own default (an aspect-owned pool / the shared `boundedElastic` / an aspect-owned scope), and only shuts down what it owns
- `lockTimeoutMillis`: Lock acquisition timeout (default: 3000ms)
- `decisionForUpdate`: Percentage of TTL after which to trigger async cache refresh (default: 80)
- `maxAttemptGetCache`: Max retry attempts when waiting for cache (default: 60, 50ms apart). Three consecutive cache-read failures while waiting fall back to the supplier immediately; supplier failures propagate as `ClientException(SUPPLIER_ERROR)`
Expand Down
26 changes: 20 additions & 6 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -155,13 +155,27 @@ refresh dependencies with `./gradlew build --refresh-dependencies`.
### Thread pools and schedulers

- `core` accepts any `java.util.concurrent.Executor` for its background cache writes; only `execute` is called and the
library never shuts a caller-supplied pool down. The default is a shared daemon pool; the Spring adapter exposes it as
the `reqShieldExecutor` bean (an `ExecutorService` the context shuts down), which you can override.
- `core-reactor` accepts a `Scheduler` (default `boundedElastic`). The Spring WebFlux adapter exposes it as the
`reqShieldScheduler` bean.
library never shuts a caller-supplied pool down. The default is a shared daemon pool. The Spring adapter runs the
writes on a daemon pool of its own that is shut down with the context, and registers no `Executor` bean, so Spring
Boot keeps its `applicationTaskExecutor`. To use your own pool, define an `Executor` bean named `reqShieldExecutor`;
its lifecycle stays yours, and like any `Executor` bean it makes Spring Boot skip `applicationTaskExecutor`.
- `core-reactor` accepts a `Scheduler` (default `boundedElastic`). The Spring WebFlux adapter uses the shared
`boundedElastic` too and never disposes it; to use your own, define a `Scheduler` bean named `reqShieldScheduler`.
- `core-kotlin-coroutine` accepts a `CoroutineScope` for background cache writes (default: a shared supervisor scope on
`Dispatchers.IO`). The coroutine Spring adapter exposes it as the `reqShieldCoroutineScope` bean and cancels it on
context shutdown.
`Dispatchers.IO`). The coroutine Spring adapter runs them on a supervisor scope of its own on `Dispatchers.IO` and
cancels it on context shutdown; to use your own, define a `CoroutineScope` bean named `reqShieldCoroutineScope`.
- The Spring adapters register none of these beans themselves, so an application bean of the same name replaces the
default instead of clashing with it, and its lifecycle stays with the application. A bean of one of these names
that is not of the expected type fails the application startup instead of being ignored.

### WebFlux adapter

- `@ReqShieldCacheable` / `@ReqShieldCacheEvict` from `core-spring-webflux` require methods that return `Mono`. The
aspect hands back a lazy `Mono`, and other return types are not rejected up front:
- On a `void` (`Unit`) method Spring discards that `Mono`, so the method body never runs and nothing is evicted,
without any error. Return `Mono<Void>` from eviction methods instead.
- Any other return type fails with a `ClassCastException` at the call site.
- Use `core-spring` for blocking methods and `core-spring-webflux-kotlin-coroutine` for `suspend` functions.

### Kotlin coroutine adapter

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -151,6 +151,7 @@ class ReqShield<T>(
},
).doFinally {
// A cache hit, read failure, or cancellation during the recheck must release our token.
// Once creation started, createReqShieldData releases it (error, cancel or after the write).
if (!cacheCreationStarted.get()) {
releaseLock(key, lockType, token)
}
Expand All @@ -169,24 +170,28 @@ class ReqShield<T>(
timeToLiveMillis: Long,
lockType: LockType,
token: String?,
): Mono<ReqShieldData<T>> =
executeCallable({ callable.call() }, key, lockType, token)
): Mono<ReqShieldData<T>> {
// Set once the lock is handed to the asynchronous cache write, which releases it after writing
val lockHandedToCacheWrite = AtomicBoolean(false)

return executeCallable({ callable.call() }, key, lockType, token)
.map { data -> buildReqShieldData(data, timeToLiveMillis) }
.doOnNext { reqShieldData ->
lockHandedToCacheWrite.set(true)
// Async fire-and-forget cache storage (matches coroutine implementation)
setReqShieldData(
reqShieldConfig.setCacheFunction,
key,
reqShieldData,
lockType,
token,
).subscribeOn(reqShieldConfig.scheduler)
.subscribe(
{ /* success - no action needed */ },
{ e -> log.error("Failed to set cache for key '{}': {}", key, e.message, e) },
)
).subscribe(
{ /* success - no action needed */ },
{ e -> log.error("Failed to set cache for key '{}': {}", key, e.message, e) },
)
}.switchIfEmpty(
Mono.defer {
lockHandedToCacheWrite.set(true)
val reqShieldData = buildReqShieldData(null, timeToLiveMillis)
// Async fire-and-forget cache storage (matches coroutine implementation)
setReqShieldData(
Expand All @@ -195,14 +200,20 @@ class ReqShield<T>(
reqShieldData,
lockType,
token,
).subscribeOn(reqShieldConfig.scheduler)
.subscribe(
{ /* success - no action needed */ },
{ e -> log.error("Failed to set cache for key '{}': {}", key, e.message, e) },
)
).subscribe(
{ /* success - no action needed */ },
{ e -> log.error("Failed to set cache for key '{}': {}", key, e.message, e) },
)
Mono.just(reqShieldData)
},
)
).doOnCancel {
// A caller cancelled while the supplier runs (client disconnect, timeout()) never reaches
// the cache write, so the lock is released here instead of lingering until it expires.
if (token != null && !lockHandedToCacheWrite.get()) {
releaseLock(key, lockType, token)
}
}
}

/**
* Waits for the request that owns the lock to fill the cache.
Expand Down Expand Up @@ -303,12 +314,15 @@ class ReqShield<T>(
Mono
.defer { setFunction(key, value, value.timeToLiveMillis) }
.onErrorMap { e -> ClientException(ErrorCode.SET_CACHE_ERROR, cause = e) }
.subscribeOn(reqShieldConfig.scheduler)
.doFinally {
// Only the holder of a token took a lock, so only it may release one.
// Placed after subscribeOn: a scheduler that rejects the write (disposed or saturated) never
// subscribes the operators above it, so a release there would never run.
if (token != null) {
releaseLock(key, lockType, token)
}
}.subscribeOn(reqShieldConfig.scheduler)
}

private fun releaseLock(
key: String,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,19 +29,25 @@ import io.mockk.every
import io.mockk.mockk
import io.mockk.slot
import io.mockk.verify
import org.awaitility.Awaitility.await
import org.junit.jupiter.api.Assertions.assertEquals
import org.junit.jupiter.api.Assertions.assertTrue
import org.junit.jupiter.api.BeforeEach
import org.junit.jupiter.api.Test
import org.junit.jupiter.api.assertThrows
import reactor.core.Disposable
import reactor.core.publisher.Mono
import reactor.core.publisher.MonoSink
import reactor.core.scheduler.Scheduler
import reactor.core.scheduler.Schedulers
import reactor.test.StepVerifier
import java.lang.reflect.InvocationTargetException
import java.lang.reflect.Method
import java.time.Duration
import java.util.concurrent.Callable
import java.util.concurrent.RejectedExecutionException
import java.util.concurrent.TimeoutException
import java.util.concurrent.atomic.AtomicBoolean
import java.util.concurrent.atomic.AtomicInteger
import java.util.concurrent.atomic.AtomicReference
import kotlin.test.assertNotNull
Expand Down Expand Up @@ -252,6 +258,62 @@ class ReqShieldTest : BaseReqShieldTest {
verify(exactly = 0) { callable.call() }
}

@Test
fun shouldReleaseLockWhenCallerCancelsWhileSupplierRuns() {
every { cacheGetter(key) } returns Mono.empty()
every { keyLock.tryLock(key, LockType.CREATE) } returns Mono.just(token)
every { keyLock.unLock(key, LockType.CREATE, token) } returns Mono.just(true)
every { callable.call() } returns Mono.never()

StepVerifier.create(reqShield.getAndSetReqShieldData(key, callable, timeToLiveMillis))
.then { verify(timeout = 1000, exactly = 1) { callable.call() } }
.thenCancel()
.verify()

verify(timeout = 1000, exactly = 1) { keyLock.unLock(key, LockType.CREATE, token) }
verify(exactly = 0) { cacheSetter(any(), any(), any()) }
}

@Test
fun shouldFreeLocalLockForNextRequestWhenCallerTimesOutWhileSupplierRuns() {
val cancelKey = "cancel-${System.nanoTime()}"
val localLock = KeyLocalLock(60_000)
val shield = ReqShield(ReqShieldConfiguration<Product>({ _, _, _ -> Mono.just(true) }, { Mono.empty() }, keyLock = localLock))

StepVerifier.create(
shield
.getAndSetReqShieldData(cancelKey, Callable { Mono.never<Product?>() }, timeToLiveMillis)
.timeout(Duration.ofMillis(100)),
).expectError(TimeoutException::class.java)
.verify()

// Without the release, the lock would stay held for its whole 60s timeout
await().atMost(Duration.ofSeconds(1)).until { localLock.tryLock(cancelKey, LockType.CREATE).block() != null }
}

@Test
fun shouldReleaseLockOnlyAfterCacheWriteWhenCallerCancelsAfterValueIsEmitted() {
lateinit var pendingWrite: MonoSink<Boolean>
every { cacheGetter(key) } returns Mono.empty()
every { cacheSetter(key, any(), any()) } returns Mono.create { pendingWrite = it }
every { keyLock.tryLock(key, LockType.CREATE) } returns Mono.just(token)
every { keyLock.unLock(key, LockType.CREATE, token) } returns Mono.just(true)
val shield =
ReqShield(
ReqShieldConfiguration(cacheSetter, cacheGetter, keyLock = keyLock, scheduler = Schedulers.immediate()),
)

StepVerifier.create(shield.getAndSetReqShieldData(key, callable, timeToLiveMillis))
.expectNextMatches { it.value == value }
.thenCancel()
.verify()

// The write already owns the lock, so the cancellation must not release it early
verify(exactly = 0) { keyLock.unLock(key, LockType.CREATE, token) }
pendingWrite.success(true)
verify(exactly = 1) { keyLock.unLock(key, LockType.CREATE, token) }
}

@Test
fun shouldKeepLockUntilAsyncCacheWriteCompletesAfterRecheckMiss() {
lateinit var pendingWrite: MonoSink<Boolean>
Expand All @@ -275,6 +337,66 @@ class ReqShieldTest : BaseReqShieldTest {
verify(exactly = 1) { callable.call() }
}

@Test
fun shouldReturnComputedDataAndReleaseLockWhenSchedulerRejectsTheCacheWrite() {
val scheduler = RejectingScheduler()
every { cacheGetter(key) } returns Mono.empty()
every { keyLock.tryLock(key, LockType.CREATE) } returns Mono.just(token)
every { keyLock.unLock(key, LockType.CREATE, token) } returns Mono.just(true)
// Built outside answers {}, whose scope has its own `value`
val supplied = Mono.just<Product?>(value)
// The supplier runs right before the write is scheduled, so only the write is rejected
every { callable.call() } answers {
scheduler.rejecting.set(true)
supplied
}
val shield = ReqShield(ReqShieldConfiguration(cacheSetter, cacheGetter, keyLock = keyLock, scheduler = scheduler))

StepVerifier.create(shield.getAndSetReqShieldData(key, callable, timeToLiveMillis))
.expectNextMatches { it.value == value }
.verifyComplete()

verify(exactly = 1) { keyLock.unLock(key, LockType.CREATE, token) }
verify(exactly = 0) { cacheSetter(any(), any(), any()) }
}

@Test
fun shouldReturnCachedDataAndReleaseLockWhenSchedulerRejectsTheRefreshWrite() {
val scheduler = RejectingScheduler()
val cached = updateTargetReqShieldData(oldValue)
every { cacheGetter(key) } returns Mono.just(cached)
every { keyLock.tryLock(key, LockType.UPDATE) } returns Mono.just(token)
every { keyLock.unLock(key, LockType.UPDATE, token) } returns Mono.just(true)
val supplied = Mono.just<Product?>(value)
every { callable.call() } answers {
scheduler.rejecting.set(true)
supplied
}
val shield = ReqShield(ReqShieldConfiguration(cacheSetter, cacheGetter, keyLock = keyLock, scheduler = scheduler))

StepVerifier.create(shield.getAndSetReqShieldData(key, callable, timeToLiveMillis))
.expectNext(cached)
.verifyComplete()

verify(exactly = 1) { keyLock.unLock(key, LockType.UPDATE, token) }
verify(exactly = 0) { cacheSetter(any(), any(), any()) }
}

/** Runs tasks inline until [rejecting] is set, then rejects them the way a disposed or saturated scheduler does. */
private class RejectingScheduler : Scheduler by Schedulers.immediate() {
val rejecting = AtomicBoolean(false)

override fun createWorker(): Scheduler.Worker {
val delegate = Schedulers.immediate().createWorker()
return object : Scheduler.Worker by delegate {
override fun schedule(task: Runnable): Disposable {
if (rejecting.get()) throw RejectedExecutionException("rejected by test")
return delegate.schedule(task)
}
}
}
}

/** Cached entry that has passed the decisionForUpdate threshold (90% of its TTL). */
private fun updateTargetReqShieldData(cachedValue: Product?): ReqShieldData<Product> =
ReqShieldData(
Expand Down
Loading
Loading