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
13 changes: 13 additions & 0 deletions .github/workflows/pull_request_event.yml
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,16 @@ jobs:
if: ${{ github.event_name == 'pull_request'
&& (github.event.action == 'opened' || github.event.action == 'synchronize' ||
github.event.action == 'reopened' || github.event.action == 'ready_for_review') }}
services:
redis:
image: redis:6.2.7-alpine
ports:
- 6379:6379
options: >-
--health-cmd "redis-cli ping"
--health-interval 10s
--health-timeout 5s
--health-retries 5
steps:
- uses: actions/checkout@v4
with:
Expand All @@ -28,4 +38,7 @@ jobs:
distribution: 'temurin'
java-version: 17
- name: Test
env:
TEST_REDIS_HOST: localhost
TEST_REDIS_PORT: 6379
run: ./gradlew clean test --info
28 changes: 28 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
# Repository Guidelines

## Project Structure & Module Organization
- Core libraries live in module directories such as `core`, `core-reactor`, and `core-kotlin-coroutine`; each follows the Gradle layout `src/main` and `src/test`.
- Spring adapters sit under `core-spring*` modules, while runnable samples are in `req-shield-*example` projects.
- Shared utilities and constants are centralized in `support`. Generated build outputs stay under each module's `build/` folder.

## Build, Test, and Development Commands
- `./gradlew build` compiles all modules, runs unit tests, and assembles artifacts; pass `--parallel` for faster local feedback.
- `./gradlew test` executes Kotlin/JVM unit tests across every enabled module.
- `./gradlew ktlintCheck` enforces the project's formatting contract before you open a PR.
- Use `./gradlew :req-shield-spring-boot3-example:bootRun` (or another sample module) to manually exercise integration paths.

## Coding Style & Naming Conventions
- Kotlin sources use 4-space indentation, `UpperCamelCase` for types, and `lowerCamelCase` for functions and properties.
- Keep package names lowercase and aligned with module boundaries (e.g., `com.linecorp.reqshield.core`).
- Always add the Apache 2.0 copyright header shown in `CONTRIBUTING.md` to new files.
- Prefer early-return patterns and meaningful exception messages; align with the `ErrorCode` enums already defined.

## Testing Guidelines
- Write tests with JUnit 5 (`org.junit.jupiter`) and place them under `src/test/kotlin` mirroring the `src/main` package.
- Use descriptive method names such as `shouldCollapseConcurrentRequests()` and cover both success and failure paths.
- When adding integration behaviour, extend the corresponding example module and run `./gradlew test` before submission.

## Commit & Pull Request Guidelines
- Follow the repository's history of concise, imperative commits (e.g., `Add cache invalidation helper`).
- Reference related issues in the body, summarise motivation, modifications, and results, and include screenshots/logs for behaviour changes.
- Verify CLS (Contributor License Agreement) status, ensure CI passes locally, and request review from a maintainer familiar with your module.
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:
- `isLocalLock`: Use local vs distributed locking (default: true)
- `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: 10)
- `maxAttemptGetCache`: Max retry attempts when waiting for cache (default: 60)
- `reqShieldWorkMode`: CREATE_AND_UPDATE_CACHE | ONLY_CREATE_CACHE | ONLY_UPDATE_CACHE

### Work Modes
Expand Down
38 changes: 38 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,44 @@ A lib that regulates the cache-based requests an application receives in terms o
`implementation("com.linecorp.cse.reqshield:core-spring-webflux:{version}")`<br>
`implementation("com.linecorp.cse.reqshield:core-spring-webflux-kotlin-coroutine:{version}")`<br>

## Testing & Integration Tips

### Integration tests with Redis (Testcontainers)

- Redis-backed integration tests using Testcontainers always run as part of module test tasks.
- Requirements:
- A working local Docker daemon with network access to pull `redis:6.2.7-alpine` on first run.
- Sufficient permissions to start containers from tests.
- If you need to temporarily bypass Redis ITs locally (e.g., no Docker), run specific unit-test-only tasks or exclude the example modules when invoking Gradle.

### WebFlux null handling

- `@ReqShieldCacheable(nullHandling = ...)` controls how `null` values are emitted in WebFlux:
- `EMIT_EMPTY` (default): map `null` to `Mono.empty()`.
- `ERROR`: throw an `IllegalStateException` if a `null` value is produced.

### Global lock guidance

- When `isLocalLock = false`, you must provide real global lock/unlock implementations.
- Recommended approach with Redis:
- Lock: `SETNX lock:{key} 1` + `PEXPIRE lock:{key} {ttlMillis}`
- Unlock: `DEL lock:{key}`
- The provided defaults return `true` and are only suitable for local/dev usage.

### Reactor Scheduler tuning

- Reactor-based modules accept a `Scheduler` (e.g., `boundedElastic`) through configuration.
- Spring WebFlux adapter exposes a `reqShieldScheduler` bean you can override for tuning thread usage.

### Kotlin Coroutine Parallelism Configuration

| Property | Default | Description |
|----------|---------|-------------|
| `reqshield.blocking.parallelism` | `availableProcessors * 2` (clamped 4-256) | Controls parallelism for blocking calls in the coroutine aspect |

**Note**: This feature uses `Dispatchers.IO.limitedParallelism()` which is marked as `@ExperimentalCoroutinesApi`.
The API may change in future Kotlin Coroutines versions.

## Contributing

Pull requests are welcome. For major changes, please open an issue first to discuss what you would like to
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,57 +27,153 @@ import kotlinx.coroutines.launch
import org.slf4j.LoggerFactory
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.Semaphore
import java.util.concurrent.atomic.AtomicBoolean
import kotlin.coroutines.CoroutineContext

private val log = LoggerFactory.getLogger(KeyLocalLock::class.java)

class KeyLocalLock(private val lockTimeoutMillis: Long) : KeyLock, CoroutineScope {
private data class LockInfo(val semaphore: Semaphore, val createdAt: Long)
/**
* Internal lock state holder.
* Using class instead of data class to allow mutable expiresAt for atomic updates.
*/
private class LockInfo(
val semaphore: Semaphore,
/**
* Expiration timestamp in milliseconds.
* @Volatile ensures visibility across threads when updated inside compute() and read by monitor.
*/
@Volatile var expiresAt: Long,
/**
* Tracks whether the lock is currently held.
* Uses AtomicBoolean with CAS operations to prevent over-release
* when multiple threads race to release the same lock (e.g., tryLock expiration
* check vs unLock, or monitor cleanup vs unLock).
*/
val isHeld: AtomicBoolean = AtomicBoolean(false),
)

private val lockMap = ConcurrentHashMap<String, LockInfo>()
companion object {
private val lockMap = ConcurrentHashMap<String, LockInfo>()

@Volatile
private var monitorJob: Job? = null

private fun ensureMonitorStarted() {
if (monitorJob?.isActive == true) return
synchronized(this) {
if (monitorJob?.isActive == true) return
monitorJob =
CoroutineScope(Dispatchers.IO).launch {
while (isActive) {
runCatching {
val now = System.currentTimeMillis()
// Remove expired locks using compute() for atomic check-and-remove.
// This prevents TOCTOU race condition where removeIf's lambda returns true
// but the actual removal happens after a new lock is acquired.
// compute() guarantees atomic execution per key, so cleanup and tryLock
// are mutually exclusive for the same key.
lockMap.keys.forEach { key ->
lockMap.compute(key) { _, lockInfo ->
if (lockInfo == null) return@compute null

if (now > lockInfo.expiresAt) {
// Expired lock: force release regardless of isHeld state.
// This handles the case where unlock() was missed due to exception.
// CAS ensures safe release (no-op if already released).
if (lockInfo.isHeld.compareAndSet(true, false)) {
lockInfo.semaphore.release()
}
null // Atomic removal
} else {
lockInfo // Keep the entry
}
}
}
delay(LOCK_MONITOR_INTERVAL_MILLIS)
}.onFailure { e ->
log.error("Error in lock lifecycle monitoring: {}", e.message, e)
}
}
}
}
}

// For testing and resource cleanup
internal fun stopMonitoring() {
synchronized(this) {
monitorJob?.cancel()
monitorJob = null
}
}
}

private val job = Job()
override val coroutineContext: CoroutineContext
get() = Dispatchers.IO + job

init {
launch {
while (isActive) {
runCatching {
val now = System.currentTimeMillis()
lockMap.entries.removeIf { now - it.value.createdAt > lockTimeoutMillis } // 특정 시간이 지나면 lock 여부와 상관없이 map에서 삭제한다.
delay(LOCK_MONITOR_INTERVAL_MILLIS)
}.onFailure { e ->
log.error("Error in lock lifecycle monitoring : {}", e.message)
}
}
}
ensureMonitorStarted()
}

override suspend fun tryLock(
key: String,
lockType: LockType,
): Boolean {
val completeKey = "${key}_${lockType.name}"
val lockInfo = lockMap.computeIfAbsent(completeKey) { LockInfo(Semaphore(1), nowToEpochTime()) }
val now = nowToEpochTime()
val result = AtomicBoolean(false)

// Use compute() for atomic lock acquisition.
// This ensures mutual exclusion with cleanup - they cannot race on the same key.
lockMap.compute(completeKey) { _, existing ->
if (existing != null) {
// Force-release expired locks to allow reacquisition.
// Use CAS to prevent race condition with concurrent unLock().
// Without CAS, if unLock() executes between isHeld.get() and release(),
// both threads would call release(), causing over-release (permits > 1).
if (now > existing.expiresAt && existing.isHeld.compareAndSet(true, false)) {
existing.semaphore.release()
}

return lockInfo.semaphore.tryAcquire()
// Existing entry: try to acquire semaphore
if (existing.semaphore.tryAcquire()) {
existing.isHeld.set(true)
existing.expiresAt = now + lockTimeoutMillis
result.set(true)
}
existing
} else {
// New entry: create and acquire
val newLock = LockInfo(Semaphore(1), now + lockTimeoutMillis)
newLock.semaphore.tryAcquire() // Always succeeds for new semaphore
newLock.isHeld.set(true)
result.set(true)
newLock
}
}
return result.get()
}

override suspend fun unLock(
key: String,
lockType: LockType,
): Boolean {
val completeKey = "${key}_${lockType.name}"
val lockInfo = lockMap[completeKey]
lockInfo?.let {
it.semaphore.release()
lockMap.remove(completeKey)
val lockInfo = lockMap[completeKey] ?: return false

// Use CAS to prevent over-release: only release if we actually hold the lock
return if (lockInfo.isHeld.compareAndSet(true, false)) {
lockInfo.semaphore.release()
true
} else {
log.debug("Attempted to unlock key '{}' that is not held", completeKey)
false
}
return true
}

fun cancel() {
job.cancel()
// Monitor cleanup is handled via stopMonitoring() in tests
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -43,11 +43,16 @@ class KeyGlobalLockTest :

@BeforeEach
fun init() {
val redisUrl = "redis://localhost:6379" // testContainer url
val host = AbstractRedisTest.redisHost
val port = AbstractRedisTest.redisPort
val redisUrl = "redis://$host:$port"
val redisClient = RedisClient.create(redisUrl)
val connection = redisClient.connect()
redisCommands = connection.async()

// Clean up all keys from previous tests for proper test isolation
connection.sync().flushdb()

globalLockFunc = { key, timeToLiveMillis ->
redisCommands.setnx(key, key).toCompletableFuture().await()
}
Expand Down
Loading
Loading