diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt index 3a351e6..b4f5efc 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt @@ -61,14 +61,14 @@ abstract class AbstractTransactionDispatcher { ) ) - TransactionState.TRANSACTION_STATE_ROLLBACK -> findOwnTransaction(transaction) + TransactionState.TRANSACTION_STATE_ROLLBACK -> findOwnUndo(transaction) .map { TransactionRollbackEvent( transaction.id, transaction.serverId, transaction.group, transaction.cause, - it.undo, + it, ) } @@ -76,7 +76,7 @@ abstract class AbstractTransactionDispatcher { } } - protected abstract fun findOwnTransaction(transaction: Transaction): Mono + protected abstract fun findOwnUndo(transaction: Transaction): Mono protected abstract fun ack( transaction: Transaction, diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt index a816fe0..7263cd0 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt @@ -10,8 +10,6 @@ import org.rooftop.netx.idl.Transaction import org.rooftop.netx.idl.TransactionState import org.rooftop.netx.meta.TransactionHandler import org.springframework.context.ApplicationContext -import org.springframework.data.redis.connection.stream.ReadOffset -import org.springframework.data.redis.connection.stream.StreamOffset import org.springframework.data.redis.core.ReactiveRedisTemplate import reactor.core.publisher.Mono import kotlin.reflect.KClass @@ -64,19 +62,15 @@ class RedisStreamTransactionDispatcher( } } - override fun findOwnTransaction(transaction: Transaction): Mono { - return reactiveRedisTemplate.opsForStream() - .read(StreamOffset.create(STREAM_KEY, ReadOffset.from("0"))) - .map { Transaction.parseFrom(it.value["data"]?.toByteArray()) } - .filter { it.group == nodeGroup } - .filter { hasUndo(it) } - .next() + override fun findOwnUndo(transaction: Transaction): Mono { + return reactiveRedisTemplate.opsForHash()[transaction.id, nodeGroup] + .switchIfEmpty( + Mono.error { + error("Cannot find undo state in transaction hashes key \"${transaction.id}\"") + } + ) } - private fun hasUndo(transaction: Transaction): Boolean = - transaction.state == TransactionState.TRANSACTION_STATE_JOIN - || transaction.state == TransactionState.TRANSACTION_STATE_START - override fun ack(transaction: Transaction, messageId: String): Mono> { return reactiveRedisTemplate.opsForStream() .acknowledge(STREAM_KEY, nodeGroup, messageId) diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt index 96b91f8..429071b 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt @@ -1,23 +1,17 @@ package org.rooftop.netx.redis -import org.redisson.api.RedissonReactiveClient import org.rooftop.netx.engine.AbstractTransactionManager -import org.rooftop.netx.engine.AbstractTransactionRetrySupporter import org.rooftop.netx.idl.Transaction import org.rooftop.netx.idl.TransactionState import org.springframework.data.redis.connection.stream.Record import org.springframework.data.redis.core.ReactiveRedisTemplate import reactor.core.publisher.Mono -import reactor.core.scheduler.Schedulers -import java.util.concurrent.TimeUnit class RedisStreamTransactionManager( nodeId: Int, nodeName: String, - nodeGroup: String, - transactionRetrySupporter: AbstractTransactionRetrySupporter, + private val nodeGroup: String, private val reactiveRedisTemplate: ReactiveRedisTemplate, - private val redissonReactiveClient: RedissonReactiveClient, ) : AbstractTransactionManager( nodeId = nodeId, nodeName = nodeName, @@ -26,13 +20,13 @@ class RedisStreamTransactionManager( override fun findAnyTransaction(transactionId: String): Mono { return reactiveRedisTemplate - .opsForValue()[transactionId] + .opsForHash()[transactionId, STATE_KEY] .switchIfEmpty( Mono.error { error("Cannot find exists transaction id \"$transactionId\"") } ) - .map { TransactionState.valueOf(String(it)) } + .map { TransactionState.valueOf(it) } } override fun publishTransaction(transactionId: String, transaction: Transaction): Mono { @@ -42,24 +36,32 @@ class RedisStreamTransactionManager( .withStreamKey(STREAM_KEY) ) .flatMap { - redissonReactiveClient.getLock("$transactionId-key") - .tryLock(10, TimeUnit.MINUTES) - } - .flatMap { - reactiveRedisTemplate.opsForValue() - .set(transactionId, transaction.state.name.toByteArray()) - } - .doFinally { - redissonReactiveClient.getLock("$transactionId-key") - .forceUnlock() - .subscribeOn(Schedulers.parallel()) - .subscribe() + if (hasUndo(transaction)) { + return@flatMap reactiveRedisTemplate.opsForHash() + .putAll( + transactionId, mapOf( + STATE_KEY to transaction.state.name, + nodeGroup to transaction.undo + ) + ) + } + reactiveRedisTemplate.opsForHash() + .putAll( + transactionId, mapOf( + STATE_KEY to transaction.state.name, + ) + ) } .map { transactionId } } + private fun hasUndo(transaction: Transaction): Boolean = + transaction.state == TransactionState.TRANSACTION_STATE_JOIN + || transaction.state == TransactionState.TRANSACTION_STATE_START + private companion object { private const val DATA = "data" private const val STREAM_KEY = "NETX_STREAM" + private const val STATE_KEY = "TX_STATE" } } diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt index fb2939c..7d4fe88 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt @@ -34,9 +34,7 @@ class RedisTransactionConfigurer( nodeId = nodeId, nodeName = nodeName, nodeGroup = nodeGroup, - transactionRetrySupporter = redisTransactionRetrySupporter(), reactiveRedisTemplate = reactiveRedisTemplate(), - redissonReactiveClient = redissonReactiveClient(), ) @Bean diff --git a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt index e7506b6..00a060b 100644 --- a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt +++ b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisStreamTransactionDispatcher.kt @@ -61,19 +61,15 @@ class NoAckRedisStreamTransactionDispatcher( } } - override fun findOwnTransaction(transaction: Transaction): Mono { - return reactiveRedisTemplate.opsForStream() - .read(StreamOffset.create(transaction.id, ReadOffset.from("0"))) - .map { Transaction.parseFrom(it.value["data"]?.toByteArray()) } - .filter { it.group == nodeGroup } - .filter { hasUndo(it) } - .next() + override fun findOwnUndo(transaction: Transaction): Mono { + return reactiveRedisTemplate.opsForHash()[transaction.id, nodeGroup] + .switchIfEmpty( + Mono.error { + error("Cannot find undo state in transaction hashes key \"${transaction.id}\"") + } + ) } - private fun hasUndo(transaction: Transaction): Boolean = - transaction.state == TransactionState.TRANSACTION_STATE_JOIN - || transaction.state == TransactionState.TRANSACTION_STATE_START - override fun ack(transaction: Transaction, messageId: String): Mono> = Mono.just(transaction to messageId) diff --git a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisTransactionConfigurer.kt b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisTransactionConfigurer.kt index 429eb40..88cd6c9 100644 --- a/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisTransactionConfigurer.kt +++ b/src/test/kotlin/org/rooftop/netx/redis/NoAckRedisTransactionConfigurer.kt @@ -34,9 +34,7 @@ class NoAckRedisTransactionConfigurer( nodeId = nodeId, nodeName = nodeName, nodeGroup = nodeGroup, - transactionRetrySupporter = redisTransactionRetrySupporter(), reactiveRedisTemplate = reactiveRedisTemplate(), - redissonReactiveClient = redissonReactiveClient(), ) @Bean