Skip to content

Commit 989fbe4

Browse files
committed
feat: Implement Distributed & Reactive CacheFlow Strategy
This commit implements the full CacheFlow strategy including: - Redis and Edge Cache integration - Russian Doll caching with tag-based invalidation - Cache warming capabilities - Touch propagation for parent-child relationships - Comprehensive testing for all new components
1 parent 1657e20 commit 989fbe4

23 files changed

Lines changed: 1244 additions & 314 deletions

‎build.gradle.kts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,9 @@ dependencies {
5858
implementation("org.jetbrains.kotlin:kotlin-reflect")
5959
implementation("org.jetbrains.kotlin:kotlin-stdlib-jdk8")
6060
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-core")
61+
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-core")
6162
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-reactor")
63+
implementation("com.fasterxml.jackson.module:jackson-module-kotlin")
6264

6365
implementation("software.amazon.awssdk:cloudfront:2.21.29")
6466

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,23 @@
1+
package io.cacheflow.spring.annotation
2+
3+
import java.lang.annotation.Inherited
4+
5+
/**
6+
* Annotation to trigger an update (touch) on a parent entity when a method is executed.
7+
*
8+
* This is useful for "Russian Doll" caching where updating a child entity should invalidate
9+
* or update the parent entity's cache key (e.g. by updating its updatedAt timestamp).
10+
*
11+
* @property parent SpEL expression to evaluate the parent ID (e.g., "#entity.parentId" or "#args[0]").
12+
* @property entityType The type of the parent entity (e.g., "user", "organization").
13+
* @property condition SpEL expression to verify if the update should proceed.
14+
*/
15+
@Target(AnnotationTarget.FUNCTION)
16+
@Retention(AnnotationRetention.RUNTIME)
17+
@Inherited
18+
@MustBeDocumented
19+
annotation class CacheFlowUpdate(
20+
val parent: String,
21+
val entityType: String,
22+
val condition: String = "",
23+
)
Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,21 @@
1+
package io.cacheflow.spring.aspect
2+
3+
/**
4+
* Interface to define how to "touch" a parent entity to update its timestamp.
5+
*
6+
* Implementations should update the 'updatedAt' (or equivalent) timestamp of the
7+
* specified entity, triggering a cache invalidation or refresh for any Russian Doll
8+
* caches that depend on that parent.
9+
*/
10+
interface ParentToucher {
11+
/**
12+
* Touches the specified parent entity.
13+
*
14+
* @param entityType The type string from @CacheFlowUpdate
15+
* @param parentId The ID of the parent entity
16+
*/
17+
fun touch(
18+
entityType: String,
19+
parentId: String,
20+
)
21+
}
Lines changed: 83 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,83 @@
1+
package io.cacheflow.spring.aspect
2+
3+
import io.cacheflow.spring.annotation.CacheFlowUpdate
4+
import org.aspectj.lang.JoinPoint
5+
import org.aspectj.lang.annotation.AfterReturning
6+
import org.aspectj.lang.annotation.Aspect
7+
import org.aspectj.lang.reflect.MethodSignature
8+
import org.slf4j.LoggerFactory
9+
import org.springframework.context.expression.MethodBasedEvaluationContext
10+
import org.springframework.core.DefaultParameterNameDiscoverer
11+
import org.springframework.expression.ExpressionParser
12+
import org.springframework.expression.spel.standard.SpelExpressionParser
13+
import org.springframework.expression.spel.support.StandardEvaluationContext
14+
import org.springframework.stereotype.Component
15+
16+
/**
17+
* Aspect to handle [CacheFlowUpdate] annotations.
18+
*
19+
* This aspect intercepts methods annotated with @CacheFlowUpdate and executes the
20+
* [ParentToucher.touch] method for the resolved parent entity.
21+
*/
22+
@Aspect
23+
@Component
24+
class TouchPropagationAspect(
25+
private val parentToucher: ParentToucher?,
26+
) {
27+
private val logger = LoggerFactory.getLogger(TouchPropagationAspect::class.java)
28+
private val parser: ExpressionParser = SpelExpressionParser()
29+
private val parameterNameDiscoverer = DefaultParameterNameDiscoverer()
30+
31+
@AfterReturning("@annotation(io.cacheflow.spring.annotation.CacheFlowUpdate)")
32+
fun handleUpdate(joinPoint: JoinPoint) {
33+
if (parentToucher == null) {
34+
logger.debug("No ParentToucher bean found. Skipping @CacheFlowUpdate processing.")
35+
return
36+
}
37+
38+
val signature = joinPoint.signature as MethodSignature
39+
var method = signature.method
40+
var annotation = method.getAnnotation(CacheFlowUpdate::class.java)
41+
42+
// If annotation is not on the interface method, check the implementation class
43+
if (annotation == null && joinPoint.target != null) {
44+
try {
45+
val targetMethod =
46+
joinPoint.target.javaClass.getMethod(method.name, *method.parameterTypes)
47+
annotation = targetMethod.getAnnotation(CacheFlowUpdate::class.java)
48+
method = targetMethod // Use the target method for context evaluation
49+
} catch (e: NoSuchMethodException) {
50+
// Ignore, keep original method
51+
}
52+
}
53+
54+
if (annotation == null) return
55+
56+
try {
57+
val context =
58+
MethodBasedEvaluationContext(
59+
joinPoint.target,
60+
method,
61+
joinPoint.args,
62+
parameterNameDiscoverer,
63+
)
64+
65+
// Check condition if present
66+
if (annotation.condition.isNotBlank()) {
67+
val conditionMet =
68+
parser.parseExpression(annotation.condition).getValue(context, Boolean::class.java)
69+
if (conditionMet != true) return
70+
}
71+
72+
// Resolve parent ID
73+
val parentId =
74+
parser.parseExpression(annotation.parent).getValue(context, String::class.java)
75+
76+
if (!parentId.isNullOrBlank()) {
77+
parentToucher.touch(annotation.entityType, parentId)
78+
}
79+
} catch (e: Exception) {
80+
logger.error("Error processing @CacheFlowUpdate", e)
81+
}
82+
}
83+
}

‎src/main/kotlin/io/cacheflow/spring/autoconfigure/CacheFlowAspectConfiguration.kt‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -75,4 +75,16 @@ class CacheFlowAspectConfiguration {
7575
dependencyResolver: DependencyResolver,
7676
tagManager: FragmentTagManager,
7777
): FragmentCacheAspect = FragmentCacheAspect(fragmentCacheService, dependencyResolver, tagManager)
78+
79+
/**
80+
* Creates the touch propagation aspect bean.
81+
*
82+
* @param parentToucher The parent toucher (optional)
83+
* @return The touch propagation aspect
84+
*/
85+
@Bean
86+
@ConditionalOnMissingBean
87+
fun touchPropagationAspect(
88+
@org.springframework.beans.factory.annotation.Autowired(required = false) parentToucher: io.cacheflow.spring.aspect.ParentToucher?,
89+
): io.cacheflow.spring.aspect.TouchPropagationAspect = io.cacheflow.spring.aspect.TouchPropagationAspect(parentToucher)
7890
}
Lines changed: 6 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,10 @@
11
package io.cacheflow.spring.autoconfigure
22

33
import io.cacheflow.spring.config.CacheFlowProperties
4+
import io.cacheflow.spring.autoconfigure.CacheFlowWarmingConfiguration
5+
import org.springframework.boot.autoconfigure.AutoConfiguration
46
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty
57
import org.springframework.boot.context.properties.EnableConfigurationProperties
6-
import org.springframework.context.annotation.Configuration
78
import org.springframework.context.annotation.Import
89

910
/**
@@ -13,19 +14,15 @@ import org.springframework.context.annotation.Import
1314
* configuration properties.
1415
*/
1516

16-
@Configuration
17-
@ConditionalOnProperty(
18-
prefix = "cacheflow",
19-
name = ["enabled"],
20-
havingValue = "true",
21-
matchIfMissing = true,
22-
)
17+
@AutoConfiguration
18+
@ConditionalOnProperty(prefix = "cacheflow", name = ["enabled"], havingValue = "true", matchIfMissing = true)
2319
@EnableConfigurationProperties(CacheFlowProperties::class)
2420
@Import(
2521
CacheFlowCoreConfiguration::class,
2622
CacheFlowFragmentConfiguration::class,
23+
CacheFlowRedisConfiguration::class,
2724
CacheFlowAspectConfiguration::class,
2825
CacheFlowManagementConfiguration::class,
29-
CacheFlowRedisConfiguration::class,
26+
CacheFlowWarmingConfiguration::class,
3027
)
3128
class CacheFlowAutoConfiguration

‎src/main/kotlin/io/cacheflow/spring/autoconfigure/CacheFlowCoreConfiguration.kt‎

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -42,7 +42,8 @@ class CacheFlowCoreConfiguration {
4242
@Autowired(required = false) @Qualifier("cacheFlowRedisTemplate") redisTemplate: RedisTemplate<String, Any>?,
4343
@Autowired(required = false) edgeCacheService: EdgeCacheIntegrationService?,
4444
@Autowired(required = false) meterRegistry: MeterRegistry?,
45-
): CacheFlowService = CacheFlowServiceImpl(properties, redisTemplate, edgeCacheService, meterRegistry)
45+
@Autowired(required = false) redisCacheInvalidator: io.cacheflow.spring.messaging.RedisCacheInvalidator?,
46+
): CacheFlowService = CacheFlowServiceImpl(properties, redisTemplate, edgeCacheService, meterRegistry, redisCacheInvalidator)
4647

4748
/**
4849
* Creates the dependency resolver bean.
@@ -51,7 +52,10 @@ class CacheFlowCoreConfiguration {
5152
*/
5253
@Bean
5354
@ConditionalOnMissingBean
54-
fun dependencyResolver(): DependencyResolver = CacheDependencyTracker()
55+
fun dependencyResolver(
56+
properties: CacheFlowProperties,
57+
@Autowired(required = false) redisTemplate: org.springframework.data.redis.core.StringRedisTemplate?,
58+
): DependencyResolver = CacheDependencyTracker(properties, redisTemplate)
5559

5660
/**
5761
* Creates the timestamp extractor bean.

‎src/main/kotlin/io/cacheflow/spring/autoconfigure/CacheFlowRedisConfiguration.kt‎

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,4 +28,46 @@ class CacheFlowRedisConfiguration {
2828
template.afterPropertiesSet()
2929
return template
3030
}
31+
32+
@Bean
33+
@ConditionalOnMissingBean
34+
fun redisCacheInvalidator(
35+
properties: io.cacheflow.spring.config.CacheFlowProperties,
36+
redisTemplate: org.springframework.data.redis.core.StringRedisTemplate,
37+
@org.springframework.context.annotation.Lazy cacheFlowService: io.cacheflow.spring.service.CacheFlowService,
38+
objectMapper: ObjectMapper,
39+
): io.cacheflow.spring.messaging.RedisCacheInvalidator {
40+
return io.cacheflow.spring.messaging.RedisCacheInvalidator(
41+
properties,
42+
redisTemplate,
43+
cacheFlowService,
44+
objectMapper
45+
)
46+
}
47+
48+
@Bean
49+
@ConditionalOnMissingBean
50+
fun cacheInvalidationListenerAdapter(
51+
redisCacheInvalidator: io.cacheflow.spring.messaging.RedisCacheInvalidator
52+
): org.springframework.data.redis.listener.adapter.MessageListenerAdapter {
53+
return org.springframework.data.redis.listener.adapter.MessageListenerAdapter(
54+
redisCacheInvalidator,
55+
"handleMessage"
56+
)
57+
}
58+
59+
@Bean
60+
@ConditionalOnMissingBean
61+
fun redisMessageListenerContainer(
62+
connectionFactory: RedisConnectionFactory,
63+
cacheInvalidationListenerAdapter: org.springframework.data.redis.listener.adapter.MessageListenerAdapter
64+
): org.springframework.data.redis.listener.RedisMessageListenerContainer {
65+
val container = org.springframework.data.redis.listener.RedisMessageListenerContainer()
66+
container.setConnectionFactory(connectionFactory)
67+
container.addMessageListener(
68+
cacheInvalidationListenerAdapter,
69+
org.springframework.data.redis.listener.ChannelTopic("cacheflow:invalidation")
70+
)
71+
return container
72+
}
3173
}
Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,23 @@
1+
package io.cacheflow.spring.autoconfigure
2+
3+
import io.cacheflow.spring.config.CacheFlowProperties
4+
import io.cacheflow.spring.warming.CacheWarmer
5+
import io.cacheflow.spring.warming.CacheWarmupProvider
6+
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean
7+
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty
8+
import org.springframework.context.annotation.Bean
9+
import org.springframework.context.annotation.Configuration
10+
11+
@Configuration
12+
@ConditionalOnProperty(prefix = "cacheflow.warming", name = ["enabled"], havingValue = "true", matchIfMissing = true)
13+
class CacheFlowWarmingConfiguration {
14+
15+
@Bean
16+
@ConditionalOnMissingBean
17+
fun cacheWarmer(
18+
properties: CacheFlowProperties,
19+
warmupProviders: List<CacheWarmupProvider>,
20+
): CacheWarmer {
21+
return CacheWarmer(properties, warmupProviders)
22+
}
23+
}

‎src/main/kotlin/io/cacheflow/spring/config/CacheFlowProperties.kt‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@ data class CacheFlowProperties(
2929
val awsCloudFront: AwsCloudFrontProperties = AwsCloudFrontProperties(),
3030
val fastly: FastlyProperties = FastlyProperties(),
3131
val metrics: MetricsProperties = MetricsProperties(),
32+
val warming: WarmingProperties = WarmingProperties(),
3233
val baseUrl: String = "https://yourdomain.com",
3334
) {
3435
/**
@@ -163,4 +164,13 @@ data class CacheFlowProperties(
163164
val enabled: Boolean = true,
164165
val exportInterval: Long = 60,
165166
)
167+
168+
/**
169+
* Cache warming configuration.
170+
*
171+
* @property enabled Whether cache warming is enabled
172+
*/
173+
data class WarmingProperties(
174+
val enabled: Boolean = true,
175+
)
166176
}

0 commit comments

Comments
 (0)