From 4cc99295a953243df7bee4f45cd1528a6e698033 Mon Sep 17 00:00:00 2001 From: devxb Date: Fri, 2 Feb 2024 15:14:45 +0900 Subject: [PATCH 01/19] =?UTF-8?q?feat:=20transaction=20api=20=EB=A5=BC=20?= =?UTF-8?q?=EC=A0=95=EC=9D=98=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/main/kotlin/org/rooftop/netx/.gitkeep | 0 .../org/rooftop/netx/api/TransactionManager.kt | 17 +++++++++++++++++ .../org/rooftop/netx/core/EventPublisher.kt | 6 ++++++ 3 files changed, 23 insertions(+) delete mode 100644 src/main/kotlin/org/rooftop/netx/.gitkeep create mode 100644 src/main/kotlin/org/rooftop/netx/api/TransactionManager.kt create mode 100644 src/main/kotlin/org/rooftop/netx/core/EventPublisher.kt diff --git a/src/main/kotlin/org/rooftop/netx/.gitkeep b/src/main/kotlin/org/rooftop/netx/.gitkeep deleted file mode 100644 index e69de29..0000000 diff --git a/src/main/kotlin/org/rooftop/netx/api/TransactionManager.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionManager.kt new file mode 100644 index 0000000..74fc6d7 --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/api/TransactionManager.kt @@ -0,0 +1,17 @@ +package org.rooftop.netx.api + +import reactor.core.publisher.Mono + +interface TransactionManager { + + fun start(replay: String): Mono + + fun exists(transactionId: String): Mono + + fun join(transactionId: String, replay: String): Mono + + fun commit(transactionId: String): Mono + + fun rollback(transactionId: String, cause: String): Mono + +} diff --git a/src/main/kotlin/org/rooftop/netx/core/EventPublisher.kt b/src/main/kotlin/org/rooftop/netx/core/EventPublisher.kt new file mode 100644 index 0000000..af1dd7a --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/core/EventPublisher.kt @@ -0,0 +1,6 @@ +package org.rooftop.netx.core + +fun interface EventPublisher { + + fun publish(event: TransactionJoinedEvent) +} From 5fac3ae4f234d9f38b1f288de90eb32184e84616 Mon Sep 17 00:00:00 2001 From: devxb Date: Fri, 2 Feb 2024 15:16:36 +0900 Subject: [PATCH 02/19] =?UTF-8?q?feat:=20=EC=8B=A4=EC=A0=9C=20=EB=8F=99?= =?UTF-8?q?=EC=9E=91=EC=9D=B4=20=EA=B5=AC=ED=98=84=EB=90=98=EC=96=B4?= =?UTF-8?q?=EC=9E=88=EB=8A=94=20engine=20=EB=A0=88=EC=9D=B4=EC=96=B4?= =?UTF-8?q?=EB=A5=BC=20=EC=A0=95=EC=9D=98=ED=95=98=EA=B3=A0=20=EB=8F=99?= =?UTF-8?q?=EC=9E=91=EC=9D=84=20=EA=B5=AC=ED=98=84=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- idl | 2 +- .../netx/api/TransactionIdGenerator.kt | 6 ++ .../netx/engine/AbstractTransactionManager.kt | 89 +++++++++++++++++++ .../netx/{core => engine}/EventPublisher.kt | 2 +- .../netx/engine/TransactionJoinedEvent.kt | 5 ++ .../netx/engine/TsidTransactionIdGenerator.kt | 13 +++ 6 files changed, 115 insertions(+), 2 deletions(-) create mode 100644 src/main/kotlin/org/rooftop/netx/api/TransactionIdGenerator.kt create mode 100644 src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt rename src/main/kotlin/org/rooftop/netx/{core => engine}/EventPublisher.kt (71%) create mode 100644 src/main/kotlin/org/rooftop/netx/engine/TransactionJoinedEvent.kt create mode 100644 src/main/kotlin/org/rooftop/netx/engine/TsidTransactionIdGenerator.kt diff --git a/idl b/idl index 573ba4a..8f2c1b2 160000 --- a/idl +++ b/idl @@ -1 +1 @@ -Subproject commit 573ba4ae4a9e9f7d1d4372f98b6dbc3b287c7ba2 +Subproject commit 8f2c1b286b481f30aa577f9205bd1638b3e685da diff --git a/src/main/kotlin/org/rooftop/netx/api/TransactionIdGenerator.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionIdGenerator.kt new file mode 100644 index 0000000..4e3d15f --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/api/TransactionIdGenerator.kt @@ -0,0 +1,6 @@ +package org.rooftop.netx.api + +fun interface TransactionIdGenerator { + + fun generate(): String +} diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt new file mode 100644 index 0000000..885c134 --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt @@ -0,0 +1,89 @@ +package org.rooftop.netx.engine + +import org.rooftop.netx.api.TransactionIdGenerator +import org.rooftop.netx.api.TransactionManager +import org.rooftop.netx.idl.Transaction +import org.rooftop.netx.idl.TransactionState +import org.rooftop.netx.idl.transaction +import reactor.core.publisher.Mono + +abstract class AbstractTransactionManager( + private val appServerId: String, + private val eventPublisher: EventPublisher, + private val transactionIdGenerator: TransactionIdGenerator, +) : TransactionManager { + + override fun start(replay: String): Mono { + return startTransaction(replay) + .publishJoinedEvent() + .contextWrite { it.put(CONTEXT_TX_KEY, transactionIdGenerator.generate()) } + } + + private fun startTransaction(replay: String): Mono { + return Mono.deferContextual { Mono.just(it[CONTEXT_TX_KEY]) } + .flatMap { transactionId -> + publishTransaction(transactionId, transaction { + id = transactionId + serverId = appServerId + this.replay = replay + }) + } + } + + override fun join(transactionId: String, replay: String): Mono { + return exists(transactionId) + .joinTransaction(replay) + .publishJoinedEvent() + .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) } + } + + private fun Mono.joinTransaction(replay: String): Mono { + return flatMap { transactionId -> + publishTransaction(transactionId, transaction { + id = transactionId + serverId = appServerId + this.replay = replay + state = TransactionState.TRANSACTION_STATE_JOIN + }) + } + } + + private fun Mono.publishJoinedEvent(): Mono { + return this.doOnSuccess { + eventPublisher.publish(TransactionJoinedEvent(it)) + } + } + + override fun rollback(transactionId: String, cause: String): Mono { + return exists(transactionId) + .publishTransaction(transaction { + id = transactionId + serverId = appServerId + state = TransactionState.TRANSACTION_STATE_ROLLBACK + this.cause = cause + }) + .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) } + } + + override fun commit(transactionId: String): Mono { + return exists(transactionId) + .publishTransaction(transaction { + id = transactionId + serverId = appServerId + state = TransactionState.TRANSACTION_STATE_COMMIT + }) + .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) } + } + + private fun Mono.publishTransaction(transaction: Transaction): Mono { + return this.flatMap { + publishTransaction(it, transaction) + } + } + + abstract fun publishTransaction(transactionId: String, transaction: Transaction): Mono + + private companion object { + private const val CONTEXT_TX_KEY = "transactionId" + } +} diff --git a/src/main/kotlin/org/rooftop/netx/core/EventPublisher.kt b/src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt similarity index 71% rename from src/main/kotlin/org/rooftop/netx/core/EventPublisher.kt rename to src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt index af1dd7a..506c1ad 100644 --- a/src/main/kotlin/org/rooftop/netx/core/EventPublisher.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt @@ -1,4 +1,4 @@ -package org.rooftop.netx.core +package org.rooftop.netx.engine fun interface EventPublisher { diff --git a/src/main/kotlin/org/rooftop/netx/engine/TransactionJoinedEvent.kt b/src/main/kotlin/org/rooftop/netx/engine/TransactionJoinedEvent.kt new file mode 100644 index 0000000..9f2d385 --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/engine/TransactionJoinedEvent.kt @@ -0,0 +1,5 @@ +package org.rooftop.netx.engine + +data class TransactionJoinedEvent( + val transactionId: String, +) diff --git a/src/main/kotlin/org/rooftop/netx/engine/TsidTransactionIdGenerator.kt b/src/main/kotlin/org/rooftop/netx/engine/TsidTransactionIdGenerator.kt new file mode 100644 index 0000000..a519127 --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/engine/TsidTransactionIdGenerator.kt @@ -0,0 +1,13 @@ +package org.rooftop.netx.engine + +import com.github.f4b6a3.tsid.TsidFactory +import org.rooftop.netx.api.TransactionIdGenerator + +class TsidTransactionIdGenerator : TransactionIdGenerator { + + override fun generate(): String = tsidFactory.create().toLong().toString() + + private companion object { + private val tsidFactory = TsidFactory.newInstance256(110) + } +} From e1de4faea5ea653c8ff7a27b469dd1ba2095f8a9 Mon Sep 17 00:00:00 2001 From: devxb Date: Fri, 2 Feb 2024 15:59:02 +0900 Subject: [PATCH 03/19] =?UTF-8?q?feat:=20Transaction=EC=9D=98=20=EC=83=81?= =?UTF-8?q?=ED=83=9C=EB=B3=84=EB=A1=9C=20event=20=EB=A5=BC=20=EB=B0=9C?= =?UTF-8?q?=ED=96=89=ED=95=98=EA=B2=8C=20=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../kotlin/org/rooftop/netx/api/TransactionCommitEvent.kt | 5 +++++ .../kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt | 6 ++++++ .../org/rooftop/netx/api/TransactionRollbackEvent.kt | 7 +++++++ .../kotlin/org/rooftop/netx/api/TransactionStartEvent.kt | 6 ++++++ .../org/rooftop/netx/engine/TransactionJoinedEvent.kt | 5 ----- 5 files changed, 24 insertions(+), 5 deletions(-) create mode 100644 src/main/kotlin/org/rooftop/netx/api/TransactionCommitEvent.kt create mode 100644 src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt create mode 100644 src/main/kotlin/org/rooftop/netx/api/TransactionRollbackEvent.kt create mode 100644 src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt delete mode 100644 src/main/kotlin/org/rooftop/netx/engine/TransactionJoinedEvent.kt diff --git a/src/main/kotlin/org/rooftop/netx/api/TransactionCommitEvent.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionCommitEvent.kt new file mode 100644 index 0000000..b6692f3 --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/api/TransactionCommitEvent.kt @@ -0,0 +1,5 @@ +package org.rooftop.netx.api + +data class TransactionCommitEvent( + private val transactionId: String, +) diff --git a/src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt new file mode 100644 index 0000000..f32568b --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt @@ -0,0 +1,6 @@ +package org.rooftop.netx.api + +data class TransactionJoinEvent( + val transactionId: String, + val replay: String, +) diff --git a/src/main/kotlin/org/rooftop/netx/api/TransactionRollbackEvent.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionRollbackEvent.kt new file mode 100644 index 0000000..8b306da --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/api/TransactionRollbackEvent.kt @@ -0,0 +1,7 @@ +package org.rooftop.netx.api + +data class TransactionRollbackEvent( + val transactionId: String, + val replay: String, + val cause: String?, +) diff --git a/src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt new file mode 100644 index 0000000..87e43d2 --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt @@ -0,0 +1,6 @@ +package org.rooftop.netx.api + +data class TransactionStartEvent( + private val transactionId: String, + private val replay: String, +) diff --git a/src/main/kotlin/org/rooftop/netx/engine/TransactionJoinedEvent.kt b/src/main/kotlin/org/rooftop/netx/engine/TransactionJoinedEvent.kt deleted file mode 100644 index 9f2d385..0000000 --- a/src/main/kotlin/org/rooftop/netx/engine/TransactionJoinedEvent.kt +++ /dev/null @@ -1,5 +0,0 @@ -package org.rooftop.netx.engine - -data class TransactionJoinedEvent( - val transactionId: String, -) From fed01371b38282e5c0bb5925affc87eab1ad5171 Mon Sep 17 00:00:00 2001 From: devxb Date: Fri, 2 Feb 2024 15:59:22 +0900 Subject: [PATCH 04/19] =?UTF-8?q?refactor:=20Transaction=20start=20?= =?UTF-8?q?=ED=98=B9=EC=9D=80=20join=EC=9D=BC=EB=95=8C,=20transaction?= =?UTF-8?q?=EC=9D=84=20=EA=B5=AC=EB=8F=85=ED=95=98=EB=8F=84=EB=A1=9D?= =?UTF-8?q?=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt | 2 +- .../org/rooftop/netx/engine/SubscribeTransactionEvent.kt | 5 +++++ 2 files changed, 6 insertions(+), 1 deletion(-) create mode 100644 src/main/kotlin/org/rooftop/netx/engine/SubscribeTransactionEvent.kt diff --git a/src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt b/src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt index 506c1ad..17d7ff0 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt @@ -2,5 +2,5 @@ package org.rooftop.netx.engine fun interface EventPublisher { - fun publish(event: TransactionJoinedEvent) + fun publish(event: SubscribeTransactionEvent) } diff --git a/src/main/kotlin/org/rooftop/netx/engine/SubscribeTransactionEvent.kt b/src/main/kotlin/org/rooftop/netx/engine/SubscribeTransactionEvent.kt new file mode 100644 index 0000000..0fe1e79 --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/engine/SubscribeTransactionEvent.kt @@ -0,0 +1,5 @@ +package org.rooftop.netx.engine + +data class SubscribeTransactionEvent( + private val transactionId: String +) From 7d5226a9a1ea44e83d9aac5e43ee0c74dc68dbd6 Mon Sep 17 00:00:00 2001 From: devxb Date: Fri, 2 Feb 2024 16:02:01 +0900 Subject: [PATCH 05/19] =?UTF-8?q?feat:=20redis-stream=20=EA=B8=B0=EB=B0=98?= =?UTF-8?q?=EC=9D=98=20transation-manager=20=EB=A5=BC=20=EA=B5=AC=ED=98=84?= =?UTF-8?q?=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../netx/engine/AbstractTransactionManager.kt | 10 ++- .../netx/engine/SubscribeTransactionEvent.kt | 2 +- .../redis/RedisStreamTransactionDispatcher.kt | 84 +++++++++++++++++++ .../redis/RedisStreamTransactionManager.kt | 50 +++++++++++ .../netx/redis/SpringEventPublisher.kt | 12 +++ 5 files changed, 153 insertions(+), 5 deletions(-) create mode 100644 src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt create mode 100644 src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt create mode 100644 src/main/kotlin/org/rooftop/netx/redis/SpringEventPublisher.kt diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt index 885c134..ced6a89 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt @@ -1,6 +1,7 @@ package org.rooftop.netx.engine import org.rooftop.netx.api.TransactionIdGenerator +import org.rooftop.netx.api.TransactionJoinEvent import org.rooftop.netx.api.TransactionManager import org.rooftop.netx.idl.Transaction import org.rooftop.netx.idl.TransactionState @@ -15,7 +16,7 @@ abstract class AbstractTransactionManager( override fun start(replay: String): Mono { return startTransaction(replay) - .publishJoinedEvent() + .subscribeTransaction() .contextWrite { it.put(CONTEXT_TX_KEY, transactionIdGenerator.generate()) } } @@ -26,6 +27,7 @@ abstract class AbstractTransactionManager( id = transactionId serverId = appServerId this.replay = replay + this.state = TransactionState.TRANSACTION_STATE_START }) } } @@ -33,7 +35,7 @@ abstract class AbstractTransactionManager( override fun join(transactionId: String, replay: String): Mono { return exists(transactionId) .joinTransaction(replay) - .publishJoinedEvent() + .subscribeTransaction() .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) } } @@ -48,9 +50,9 @@ abstract class AbstractTransactionManager( } } - private fun Mono.publishJoinedEvent(): Mono { + private fun Mono.subscribeTransaction(): Mono { return this.doOnSuccess { - eventPublisher.publish(TransactionJoinedEvent(it)) + eventPublisher.publish(SubscribeTransactionEvent(it)) } } diff --git a/src/main/kotlin/org/rooftop/netx/engine/SubscribeTransactionEvent.kt b/src/main/kotlin/org/rooftop/netx/engine/SubscribeTransactionEvent.kt index 0fe1e79..79e8447 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/SubscribeTransactionEvent.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/SubscribeTransactionEvent.kt @@ -1,5 +1,5 @@ package org.rooftop.netx.engine data class SubscribeTransactionEvent( - private val transactionId: String + val transactionId: String ) diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt new file mode 100644 index 0000000..36b6c4c --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt @@ -0,0 +1,84 @@ +package org.rooftop.netx.redis + +import org.rooftop.netx.api.TransactionCommitEvent +import org.rooftop.netx.api.TransactionJoinEvent +import org.rooftop.netx.api.TransactionRollbackEvent +import org.rooftop.netx.api.TransactionStartEvent +import org.rooftop.netx.engine.SubscribeTransactionEvent +import org.rooftop.netx.idl.Transaction +import org.rooftop.netx.idl.TransactionState +import org.springframework.context.ApplicationEventPublisher +import org.springframework.context.event.EventListener +import org.springframework.data.redis.connection.ReactiveRedisConnectionFactory +import org.springframework.data.redis.connection.stream.Consumer +import org.springframework.data.redis.connection.stream.StreamOffset +import org.springframework.data.redis.stream.StreamReceiver +import reactor.core.publisher.Flux +import reactor.core.publisher.Mono +import reactor.core.scheduler.Schedulers + +class RedisStreamTransactionDispatcher( + private val eventPublisher: ApplicationEventPublisher, + private val streamGroup: String, + private val name: String, + connectionFactory: ReactiveRedisConnectionFactory, +) { + + private val options = StreamReceiver.StreamReceiverOptions.builder() + .pollTimeout(java.time.Duration.ofMillis(100)) + .build() + + private val receiver = StreamReceiver.create(connectionFactory, options) + + @EventListener(SubscribeTransactionEvent::class) + fun subscribeStream(subscribeTransactionEvent: SubscribeTransactionEvent): Flux { + return receiver.receiveAutoAck( + Consumer.from(streamGroup, name), + StreamOffset.fromStart(subscribeTransactionEvent.transactionId) + ).publishOn(Schedulers.parallel()) + .map { Transaction.parseFrom(it.value["data"]?.toByteArray()) } + .dispatch() + } + + private fun Flux.dispatch(): Flux { + return this.flatMap { + when (it.state) { + TransactionState.TRANSACTION_STATE_JOIN -> publishJoin(it) + TransactionState.TRANSACTION_STATE_COMMIT -> publishCommit(it) + TransactionState.TRANSACTION_STATE_ROLLBACK -> publishRollback(it) + TransactionState.TRANSACTION_STATE_START -> publishStart(it) + else -> error("Cannot find matched transaction state \"${it.state}\"") + } + } + } + + private fun publishJoin(it: Transaction): Mono { + return Mono.just(it) + .doOnNext { eventPublisher.publishEvent(TransactionJoinEvent(it.id, it.replay)) } + } + + private fun publishCommit(it: Transaction): Mono { + return Mono.just(it) + .doOnNext { eventPublisher.publishEvent(TransactionCommitEvent(it.id)) } + } + + private fun publishRollback(it: Transaction): Mono { + return Mono.just(it) + .doOnNext { + eventPublisher.publishEvent( + TransactionRollbackEvent( + it.id, + it.replay, + it.cause + ) + ) + } + } + + private fun publishStart(it: Transaction): Mono { + return Mono.just(it) + .doOnNext { + eventPublisher.publishEvent(TransactionStartEvent(it.id, it.replay)) + } + } +} diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt new file mode 100644 index 0000000..84812f9 --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt @@ -0,0 +1,50 @@ +package org.rooftop.netx.redis + +import org.rooftop.netx.api.TransactionIdGenerator +import org.rooftop.netx.engine.AbstractTransactionManager +import org.rooftop.netx.engine.EventPublisher +import org.rooftop.netx.idl.Transaction +import org.springframework.data.domain.Range +import org.springframework.data.redis.connection.stream.Record +import org.springframework.data.redis.core.ReactiveRedisTemplate +import reactor.core.publisher.Mono + +class RedisStreamTransactionManager( + appServerId: String, + transactionIdGenerator: TransactionIdGenerator, + eventPublisher: EventPublisher, + private val transactionServer: ReactiveRedisTemplate, +) : AbstractTransactionManager(appServerId, eventPublisher, transactionIdGenerator) { + + override fun exists(transactionId: String): Mono { + return transactionServer.opsForStream() + .range(transactionId, Range.open("-", "+")) + .map { Transaction.parseFrom(it.value[DATA].toString().toByteArray()) } + .next() + .switchIfEmpty( + Mono.error { + IllegalStateException("Cannot find exists transaction id \"$transactionId\"") + } + ) + .transformTransactionId() + } + + private fun Mono<*>.transformTransactionId(): Mono { + return this.flatMap { + Mono.deferContextual { Mono.just(it["transactionId"]) } + } + } + + override fun publishTransaction(transactionId: String, transaction: Transaction): Mono { + return transactionServer.opsForStream() + .add( + Record.of(mapOf(DATA to transaction.toByteArray())) + .withStreamKey(transactionId) + ) + .transformTransactionId() + } + + private companion object { + private const val DATA = "data" + } +} diff --git a/src/main/kotlin/org/rooftop/netx/redis/SpringEventPublisher.kt b/src/main/kotlin/org/rooftop/netx/redis/SpringEventPublisher.kt new file mode 100644 index 0000000..e09abde --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/redis/SpringEventPublisher.kt @@ -0,0 +1,12 @@ +package org.rooftop.netx.redis + +import org.rooftop.netx.engine.EventPublisher +import org.rooftop.netx.engine.SubscribeTransactionEvent +import org.springframework.context.ApplicationEventPublisher + +class SpringEventPublisher(private val eventPublisher: ApplicationEventPublisher) : EventPublisher { + + override fun publish(event: SubscribeTransactionEvent) { + eventPublisher.publishEvent(event) + } +} From ba67fe8ae0c4cdf7f0541779e7019cedacf7fcfc Mon Sep 17 00:00:00 2001 From: devxb Date: Fri, 2 Feb 2024 19:26:51 +0900 Subject: [PATCH 06/19] =?UTF-8?q?feat:=20AutoConfig=20=EB=A5=BC=20?= =?UTF-8?q?=EC=A0=95=EC=9D=98=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- build.gradle | 10 +++- gradle/spring.gradle | 7 +-- idl | 2 +- .../netx/api/TransactionIdGenerator.kt | 6 -- .../netx/engine/AbstractTransactionManager.kt | 14 ++--- ...Generator.kt => TransactionIdGenerator.kt} | 3 +- .../netx/redis/ByteArrayRedisSerializer.kt | 10 ++++ .../redis/RedisStreamTransactionDispatcher.kt | 4 +- .../redis/RedisStreamTransactionManager.kt | 10 ++-- .../netx/redis/RedisTransactionAutoConfig.kt | 59 +++++++++++++++++++ ...ot.autoconfigure.AutoConfiguration.imports | 1 + 11 files changed, 95 insertions(+), 31 deletions(-) delete mode 100644 src/main/kotlin/org/rooftop/netx/api/TransactionIdGenerator.kt rename src/main/kotlin/org/rooftop/netx/engine/{TsidTransactionIdGenerator.kt => TransactionIdGenerator.kt} (70%) create mode 100644 src/main/kotlin/org/rooftop/netx/redis/ByteArrayRedisSerializer.kt create mode 100644 src/main/kotlin/org/rooftop/netx/redis/RedisTransactionAutoConfig.kt diff --git a/build.gradle b/build.gradle index a1553c4..9829923 100644 --- a/build.gradle +++ b/build.gradle @@ -1,8 +1,10 @@ +import org.springframework.boot.gradle.plugin.SpringBootPlugin + plugins { id "application" id "org.jetbrains.kotlin.jvm" version "${jetbrainKotlinVersion}" id "org.jetbrains.kotlin.plugin.spring" version "${jetbrainKotlinVersion}" - id "org.springframework.boot" version "${springbootVersion}" + id "org.springframework.boot" version "${springbootVersion}" apply false id "io.spring.dependency-management" version "${springDependencyManagementVersion}" id "org.sonarqube" version "${sonarcloudVersion}" id "com.google.protobuf" version "${protobufPluginVersion}" @@ -15,6 +17,12 @@ repositories { mavenCentral() } +dependencyManagement { + imports { + mavenBom SpringBootPlugin.BOM_COORDINATES + } +} + apply from: "gradle/mq.gradle" apply from: "gradle/test.gradle" apply from: "gradle/core.gradle" diff --git a/gradle/spring.gradle b/gradle/spring.gradle index 10a602b..b5a482d 100644 --- a/gradle/spring.gradle +++ b/gradle/spring.gradle @@ -1,9 +1,6 @@ -jar { - enabled = false -} - dependencies { - implementation "org.springframework.boot:spring-boot-starter" + implementation 'org.springframework:spring-context-support' + implementation 'org.springframework.boot:spring-boot-autoconfigure' implementation "org.springframework.boot:spring-boot-starter-data-redis-reactive" testImplementation "org.springframework.boot:spring-boot-starter-test" diff --git a/idl b/idl index 8f2c1b2..e1f481d 160000 --- a/idl +++ b/idl @@ -1 +1 @@ -Subproject commit 8f2c1b286b481f30aa577f9205bd1638b3e685da +Subproject commit e1f481d6f34e879b82487603c11cb7a57cc6a6ab diff --git a/src/main/kotlin/org/rooftop/netx/api/TransactionIdGenerator.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionIdGenerator.kt deleted file mode 100644 index 4e3d15f..0000000 --- a/src/main/kotlin/org/rooftop/netx/api/TransactionIdGenerator.kt +++ /dev/null @@ -1,6 +0,0 @@ -package org.rooftop.netx.api - -fun interface TransactionIdGenerator { - - fun generate(): String -} diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt index ced6a89..5c9e8d7 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt @@ -1,7 +1,5 @@ package org.rooftop.netx.engine -import org.rooftop.netx.api.TransactionIdGenerator -import org.rooftop.netx.api.TransactionJoinEvent import org.rooftop.netx.api.TransactionManager import org.rooftop.netx.idl.Transaction import org.rooftop.netx.idl.TransactionState @@ -9,9 +7,9 @@ import org.rooftop.netx.idl.transaction import reactor.core.publisher.Mono abstract class AbstractTransactionManager( - private val appServerId: String, + private val nodeName: String, private val eventPublisher: EventPublisher, - private val transactionIdGenerator: TransactionIdGenerator, + private val transactionIdGenerator: TransactionIdGenerator = TransactionIdGenerator(), ) : TransactionManager { override fun start(replay: String): Mono { @@ -25,7 +23,7 @@ abstract class AbstractTransactionManager( .flatMap { transactionId -> publishTransaction(transactionId, transaction { id = transactionId - serverId = appServerId + serverId = nodeName this.replay = replay this.state = TransactionState.TRANSACTION_STATE_START }) @@ -43,7 +41,7 @@ abstract class AbstractTransactionManager( return flatMap { transactionId -> publishTransaction(transactionId, transaction { id = transactionId - serverId = appServerId + serverId = nodeName this.replay = replay state = TransactionState.TRANSACTION_STATE_JOIN }) @@ -60,7 +58,7 @@ abstract class AbstractTransactionManager( return exists(transactionId) .publishTransaction(transaction { id = transactionId - serverId = appServerId + serverId = nodeName state = TransactionState.TRANSACTION_STATE_ROLLBACK this.cause = cause }) @@ -71,7 +69,7 @@ abstract class AbstractTransactionManager( return exists(transactionId) .publishTransaction(transaction { id = transactionId - serverId = appServerId + serverId = nodeName state = TransactionState.TRANSACTION_STATE_COMMIT }) .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) } diff --git a/src/main/kotlin/org/rooftop/netx/engine/TsidTransactionIdGenerator.kt b/src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt similarity index 70% rename from src/main/kotlin/org/rooftop/netx/engine/TsidTransactionIdGenerator.kt rename to src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt index a519127..e02f016 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/TsidTransactionIdGenerator.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt @@ -1,9 +1,8 @@ package org.rooftop.netx.engine import com.github.f4b6a3.tsid.TsidFactory -import org.rooftop.netx.api.TransactionIdGenerator -class TsidTransactionIdGenerator : TransactionIdGenerator { +class TransactionIdGenerator : TransactionIdGenerator { override fun generate(): String = tsidFactory.create().toLong().toString() diff --git a/src/main/kotlin/org/rooftop/netx/redis/ByteArrayRedisSerializer.kt b/src/main/kotlin/org/rooftop/netx/redis/ByteArrayRedisSerializer.kt new file mode 100644 index 0000000..884998e --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/redis/ByteArrayRedisSerializer.kt @@ -0,0 +1,10 @@ +package org.rooftop.pay.infra.transaction + +import org.springframework.data.redis.serializer.RedisSerializer + +class ByteArrayRedisSerializer : RedisSerializer { + + override fun serialize(t: ByteArray?): ByteArray? = t + + override fun deserialize(bytes: ByteArray?): ByteArray? = bytes +} diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt index 36b6c4c..6b17b4f 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt @@ -20,7 +20,7 @@ import reactor.core.scheduler.Schedulers class RedisStreamTransactionDispatcher( private val eventPublisher: ApplicationEventPublisher, private val streamGroup: String, - private val name: String, + private val nodeName: String, connectionFactory: ReactiveRedisConnectionFactory, ) { @@ -33,7 +33,7 @@ class RedisStreamTransactionDispatcher( @EventListener(SubscribeTransactionEvent::class) fun subscribeStream(subscribeTransactionEvent: SubscribeTransactionEvent): Flux { return receiver.receiveAutoAck( - Consumer.from(streamGroup, name), + Consumer.from(streamGroup, nodeName), StreamOffset.fromStart(subscribeTransactionEvent.transactionId) ).publishOn(Schedulers.parallel()) .map { Transaction.parseFrom(it.value["data"]?.toByteArray()) } diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt index 84812f9..4d8cc56 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt @@ -1,20 +1,18 @@ package org.rooftop.netx.redis -import org.rooftop.netx.api.TransactionIdGenerator import org.rooftop.netx.engine.AbstractTransactionManager -import org.rooftop.netx.engine.EventPublisher import org.rooftop.netx.idl.Transaction +import org.springframework.context.ApplicationEventPublisher import org.springframework.data.domain.Range import org.springframework.data.redis.connection.stream.Record import org.springframework.data.redis.core.ReactiveRedisTemplate import reactor.core.publisher.Mono class RedisStreamTransactionManager( - appServerId: String, - transactionIdGenerator: TransactionIdGenerator, - eventPublisher: EventPublisher, + nodeName: String, + applicationEventPublisher: ApplicationEventPublisher, private val transactionServer: ReactiveRedisTemplate, -) : AbstractTransactionManager(appServerId, eventPublisher, transactionIdGenerator) { +) : AbstractTransactionManager(nodeName, SpringEventPublisher(applicationEventPublisher)) { override fun exists(transactionId: String): Mono { return transactionServer.opsForStream() diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionAutoConfig.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionAutoConfig.kt new file mode 100644 index 0000000..6197b3d --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionAutoConfig.kt @@ -0,0 +1,59 @@ +package org.rooftop.netx.redis + +import org.rooftop.netx.api.TransactionManager +import org.rooftop.pay.infra.transaction.ByteArrayRedisSerializer +import org.springframework.beans.factory.annotation.Value +import org.springframework.boot.autoconfigure.AutoConfiguration +import org.springframework.context.ApplicationEventPublisher +import org.springframework.context.annotation.Bean +import org.springframework.data.redis.connection.ReactiveRedisConnectionFactory +import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory +import org.springframework.data.redis.core.ReactiveRedisTemplate +import org.springframework.data.redis.serializer.RedisSerializationContext +import org.springframework.data.redis.serializer.StringRedisSerializer + +@AutoConfiguration +class RedisTransactionAutoConfiguration( + @Value("\${netx.host}") private val host: String, + @Value("\${netx.port}") private val port: String, + @Value("\${netx.group}") private val group: String, + @Value("\${netx.node-name}") private val nodeName: String, + private val applicationEventPublisher: ApplicationEventPublisher, +) { + + @Bean + fun redisStreamTransactionManager(): TransactionManager = + RedisStreamTransactionManager(nodeName, applicationEventPublisher, transactionServer()) + + @Bean + fun redisStreamTransactionDispatcher(): RedisStreamTransactionDispatcher = + RedisStreamTransactionDispatcher( + applicationEventPublisher, + group, + nodeName, + transactionServerConnectionFactory() + ) + + @Bean + fun transactionServer(): ReactiveRedisTemplate { + val builder = RedisSerializationContext.newSerializationContext( + StringRedisSerializer() + ) + + val context = builder.value(byteArrayRedisSerializer()).build() + + return ReactiveRedisTemplate(transactionServerConnectionFactory(), context) + } + + @Bean + fun byteArrayRedisSerializer(): ByteArrayRedisSerializer { + return ByteArrayRedisSerializer() + } + + @Bean + fun transactionServerConnectionFactory(): ReactiveRedisConnectionFactory { + val port: String = System.getProperty("netx.port") ?: port + + return LettuceConnectionFactory(host, port.toInt()) + } +} diff --git a/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports b/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports index e69de29..b180427 100644 --- a/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports +++ b/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports @@ -0,0 +1 @@ +org.rooftop.netx.redis.RedisTransactionAutoConfiguration From 9b3dbf9c173da23555ceae77111ea186a8030e56 Mon Sep 17 00:00:00 2001 From: devxb Date: Fri, 2 Feb 2024 19:27:22 +0900 Subject: [PATCH 07/19] =?UTF-8?q?fix:=20TransactionIdGenerator=EC=97=90?= =?UTF-8?q?=EC=84=9C=20=EC=A1=B4=EC=9E=AC=ED=95=98=EC=A7=80=20=EC=95=8A?= =?UTF-8?q?=EB=8A=94=20=ED=81=B4=EB=9E=98=EC=8A=A4=20=EA=B5=AC=ED=98=84?= =?UTF-8?q?=ED=95=98=EA=B3=A0=20=EC=9E=88=EB=8A=94=20=EB=B2=84=EA=B7=B8=20?= =?UTF-8?q?=EC=88=98=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt b/src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt index e02f016..6a61a5b 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt @@ -2,9 +2,9 @@ package org.rooftop.netx.engine import com.github.f4b6a3.tsid.TsidFactory -class TransactionIdGenerator : TransactionIdGenerator { +class TransactionIdGenerator { - override fun generate(): String = tsidFactory.create().toLong().toString() + fun generate(): String = tsidFactory.create().toLong().toString() private companion object { private val tsidFactory = TsidFactory.newInstance256(110) From b69b3a8cb96bc625c65528de51ee93570e35ef7b Mon Sep 17 00:00:00 2001 From: devxb Date: Fri, 2 Feb 2024 19:34:31 +0900 Subject: [PATCH 08/19] =?UTF-8?q?refactor:=20transaction=20id=20=EB=A5=BC?= =?UTF-8?q?=20node-id=20=EA=B8=B0=EB=B0=98=EC=9C=BC=EB=A1=9C=20=EC=83=9D?= =?UTF-8?q?=EC=84=B1=ED=95=98=EB=8F=84=EB=A1=9D=20=EC=88=98=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../rooftop/netx/engine/AbstractTransactionManager.kt | 3 ++- .../org/rooftop/netx/engine/TransactionIdGenerator.kt | 10 +++++----- .../netx/redis/RedisStreamTransactionManager.kt | 3 ++- .../rooftop/netx/redis/RedisTransactionAutoConfig.kt | 8 +++++++- 4 files changed, 16 insertions(+), 8 deletions(-) diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt index 5c9e8d7..f5033ae 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt @@ -7,9 +7,10 @@ import org.rooftop.netx.idl.transaction import reactor.core.publisher.Mono abstract class AbstractTransactionManager( + nodeId: Int, private val nodeName: String, private val eventPublisher: EventPublisher, - private val transactionIdGenerator: TransactionIdGenerator = TransactionIdGenerator(), + private val transactionIdGenerator: TransactionIdGenerator = TransactionIdGenerator(nodeId), ) : TransactionManager { override fun start(replay: String): Mono { diff --git a/src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt b/src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt index 6a61a5b..968a13a 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt @@ -2,11 +2,11 @@ package org.rooftop.netx.engine import com.github.f4b6a3.tsid.TsidFactory -class TransactionIdGenerator { +class TransactionIdGenerator( + nodeId: Int, + private val tsidFactory: TsidFactory = TsidFactory.newInstance256(nodeId), +) { fun generate(): String = tsidFactory.create().toLong().toString() - - private companion object { - private val tsidFactory = TsidFactory.newInstance256(110) - } } + diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt index 4d8cc56..a72e9e6 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt @@ -9,10 +9,11 @@ import org.springframework.data.redis.core.ReactiveRedisTemplate import reactor.core.publisher.Mono class RedisStreamTransactionManager( + nodeId: Int, nodeName: String, applicationEventPublisher: ApplicationEventPublisher, private val transactionServer: ReactiveRedisTemplate, -) : AbstractTransactionManager(nodeName, SpringEventPublisher(applicationEventPublisher)) { +) : AbstractTransactionManager(nodeId, nodeName, SpringEventPublisher(applicationEventPublisher)) { override fun exists(transactionId: String): Mono { return transactionServer.opsForStream() diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionAutoConfig.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionAutoConfig.kt index 6197b3d..9be5a2a 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionAutoConfig.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionAutoConfig.kt @@ -17,13 +17,19 @@ class RedisTransactionAutoConfiguration( @Value("\${netx.host}") private val host: String, @Value("\${netx.port}") private val port: String, @Value("\${netx.group}") private val group: String, + @Value("\${netx.node-id}") private val nodeId: Int, @Value("\${netx.node-name}") private val nodeName: String, private val applicationEventPublisher: ApplicationEventPublisher, ) { @Bean fun redisStreamTransactionManager(): TransactionManager = - RedisStreamTransactionManager(nodeName, applicationEventPublisher, transactionServer()) + RedisStreamTransactionManager( + nodeId, + nodeName, + applicationEventPublisher, + transactionServer() + ) @Bean fun redisStreamTransactionDispatcher(): RedisStreamTransactionDispatcher = From a6a5c05957dd6c5aad58357b5c2ff9e25dd1ab80 Mon Sep 17 00:00:00 2001 From: devxb Date: Fri, 2 Feb 2024 19:50:04 +0900 Subject: [PATCH 09/19] =?UTF-8?q?refactor:=20=EC=9E=90=EB=8F=99=EA=B5=AC?= =?UTF-8?q?=EC=84=B1=EC=9D=B4=20=EC=96=B4=EB=85=B8=ED=85=8C=EC=9D=B4?= =?UTF-8?q?=EC=85=98=20=EA=B8=B0=EB=B0=98=EC=9C=BC=EB=A1=9C=20=EB=8F=99?= =?UTF-8?q?=EC=9E=91=20=EA=B0=80=EB=8A=A5=ED=95=98=EB=8F=84=EB=A1=9D=20?= =?UTF-8?q?=EC=88=98=EC=A0=95=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../rooftop/netx/redis/AutoConfigureRedisTransaction.kt | 8 ++++++++ ...sactionAutoConfig.kt => RedisTransactionConfigurer.kt} | 2 +- ...framework.boot.autoconfigure.AutoConfiguration.imports | 1 - src/main/resources/META-INF/spring/spring.factories | 2 ++ 4 files changed, 11 insertions(+), 2 deletions(-) create mode 100644 src/main/kotlin/org/rooftop/netx/redis/AutoConfigureRedisTransaction.kt rename src/main/kotlin/org/rooftop/netx/redis/{RedisTransactionAutoConfig.kt => RedisTransactionConfigurer.kt} (98%) delete mode 100644 src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports create mode 100644 src/main/resources/META-INF/spring/spring.factories diff --git a/src/main/kotlin/org/rooftop/netx/redis/AutoConfigureRedisTransaction.kt b/src/main/kotlin/org/rooftop/netx/redis/AutoConfigureRedisTransaction.kt new file mode 100644 index 0000000..44ca2fa --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/redis/AutoConfigureRedisTransaction.kt @@ -0,0 +1,8 @@ +package org.rooftop.netx.redis + +import org.springframework.boot.autoconfigure.ImportAutoConfiguration + +@ImportAutoConfiguration +@Target(AnnotationTarget.CLASS) +@Retention(AnnotationRetention.RUNTIME) +annotation class AutoConfigureRedisTransaction diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionAutoConfig.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt similarity index 98% rename from src/main/kotlin/org/rooftop/netx/redis/RedisTransactionAutoConfig.kt rename to src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt index 9be5a2a..5a655c0 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionAutoConfig.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt @@ -13,7 +13,7 @@ import org.springframework.data.redis.serializer.RedisSerializationContext import org.springframework.data.redis.serializer.StringRedisSerializer @AutoConfiguration -class RedisTransactionAutoConfiguration( +class RedisTransactionConfigurer( @Value("\${netx.host}") private val host: String, @Value("\${netx.port}") private val port: String, @Value("\${netx.group}") private val group: String, diff --git a/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports b/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports deleted file mode 100644 index b180427..0000000 --- a/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports +++ /dev/null @@ -1 +0,0 @@ -org.rooftop.netx.redis.RedisTransactionAutoConfiguration diff --git a/src/main/resources/META-INF/spring/spring.factories b/src/main/resources/META-INF/spring/spring.factories new file mode 100644 index 0000000..18fe703 --- /dev/null +++ b/src/main/resources/META-INF/spring/spring.factories @@ -0,0 +1,2 @@ +org.rooftop.netx.redis.AutoConfigureRedisTransaction=\ +org.rooftop.netx.redis.RedisTransactionConfigurer From fb3ba4b4a6d097f274d42de1b94139c2d1cb2f00 Mon Sep 17 00:00:00 2001 From: devxb Date: Sat, 3 Feb 2024 00:03:48 +0900 Subject: [PATCH 10/19] =?UTF-8?q?refactor:=20TransactionEvent=EB=93=A4?= =?UTF-8?q?=EC=97=90=20=EC=96=B4=EB=96=A4=20=EC=84=9C=EB=B2=84=EC=97=90?= =?UTF-8?q?=EC=84=9C=20=EB=B0=9C=ED=96=89=EB=90=98=EC=97=88=EB=8A=94?= =?UTF-8?q?=EC=A7=80=20=ED=99=95=EC=9D=B8=ED=95=A0=20=EC=88=98=20=EC=9E=88?= =?UTF-8?q?=EB=8F=84=EB=A1=9D=20nodeName=20=ED=95=84=EB=93=9C=EB=A5=BC=20?= =?UTF-8?q?=EC=B6=94=EA=B0=80=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../rooftop/netx/api/TransactionCommitEvent.kt | 3 ++- .../rooftop/netx/api/TransactionJoinEvent.kt | 1 + .../netx/api/TransactionRollbackEvent.kt | 1 + .../rooftop/netx/api/TransactionStartEvent.kt | 5 +++-- .../redis/RedisStreamTransactionDispatcher.kt | 17 +++++++++++++---- 5 files changed, 20 insertions(+), 7 deletions(-) diff --git a/src/main/kotlin/org/rooftop/netx/api/TransactionCommitEvent.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionCommitEvent.kt index b6692f3..ef5e477 100644 --- a/src/main/kotlin/org/rooftop/netx/api/TransactionCommitEvent.kt +++ b/src/main/kotlin/org/rooftop/netx/api/TransactionCommitEvent.kt @@ -1,5 +1,6 @@ package org.rooftop.netx.api data class TransactionCommitEvent( - private val transactionId: String, + val transactionId: String, + val nodeName: String, ) diff --git a/src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt index f32568b..8899322 100644 --- a/src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt +++ b/src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt @@ -3,4 +3,5 @@ package org.rooftop.netx.api data class TransactionJoinEvent( val transactionId: String, val replay: String, + val nodeName: String, ) diff --git a/src/main/kotlin/org/rooftop/netx/api/TransactionRollbackEvent.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionRollbackEvent.kt index 8b306da..ffa9d2a 100644 --- a/src/main/kotlin/org/rooftop/netx/api/TransactionRollbackEvent.kt +++ b/src/main/kotlin/org/rooftop/netx/api/TransactionRollbackEvent.kt @@ -3,5 +3,6 @@ package org.rooftop.netx.api data class TransactionRollbackEvent( val transactionId: String, val replay: String, + val nodeName: String, val cause: String?, ) diff --git a/src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt index 87e43d2..cfdc7bf 100644 --- a/src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt +++ b/src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt @@ -1,6 +1,7 @@ package org.rooftop.netx.api data class TransactionStartEvent( - private val transactionId: String, - private val replay: String, + val transactionId: String, + val replay: String, + val nodeName: String, ) diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt index 6b17b4f..1f4522b 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt @@ -54,12 +54,20 @@ class RedisStreamTransactionDispatcher( private fun publishJoin(it: Transaction): Mono { return Mono.just(it) - .doOnNext { eventPublisher.publishEvent(TransactionJoinEvent(it.id, it.replay)) } + .doOnNext { + eventPublisher.publishEvent( + TransactionJoinEvent( + it.id, + it.replay, + it.serverId + ) + ) + } } private fun publishCommit(it: Transaction): Mono { return Mono.just(it) - .doOnNext { eventPublisher.publishEvent(TransactionCommitEvent(it.id)) } + .doOnNext { eventPublisher.publishEvent(TransactionCommitEvent(it.id, it.serverId)) } } private fun publishRollback(it: Transaction): Mono { @@ -69,7 +77,8 @@ class RedisStreamTransactionDispatcher( TransactionRollbackEvent( it.id, it.replay, - it.cause + it.serverId, + it.cause, ) ) } @@ -78,7 +87,7 @@ class RedisStreamTransactionDispatcher( private fun publishStart(it: Transaction): Mono { return Mono.just(it) .doOnNext { - eventPublisher.publishEvent(TransactionStartEvent(it.id, it.replay)) + eventPublisher.publishEvent(TransactionStartEvent(it.id, it.replay, it.serverId)) } } } From fc4e07600068c71019e0cd36d9f3c0c95addb301 Mon Sep 17 00:00:00 2001 From: devxb Date: Sat, 3 Feb 2024 00:09:52 +0900 Subject: [PATCH 11/19] =?UTF-8?q?refactor:=20TransactionDispatcher?= =?UTF-8?q?=EB=A5=BC=20=EC=B6=94=EC=83=81=ED=99=94=EC=8B=9C=ED=82=A4?= =?UTF-8?q?=EA=B3=A0,=20=EB=A1=9C=EC=A7=81=EC=9D=84=20=EB=94=B0=EB=A5=B4?= =?UTF-8?q?=EB=8F=84=EB=A1=9D=20=EA=B0=95=EC=A0=9C=ED=99=94=20=ED=95=9C?= =?UTF-8?q?=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../engine/AbstractTransactionDispatcher.kt | 74 +++++++++++++++++++ .../org/rooftop/netx/engine/EventPublisher.kt | 2 +- .../redis/RedisStreamTransactionDispatcher.kt | 73 ++---------------- .../netx/redis/SpringEventPublisher.kt | 2 +- 4 files changed, 84 insertions(+), 67 deletions(-) create mode 100644 src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt new file mode 100644 index 0000000..bcf4e20 --- /dev/null +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt @@ -0,0 +1,74 @@ +package org.rooftop.netx.engine + +import org.rooftop.netx.api.TransactionCommitEvent +import org.rooftop.netx.api.TransactionJoinEvent +import org.rooftop.netx.api.TransactionRollbackEvent +import org.rooftop.netx.api.TransactionStartEvent +import org.rooftop.netx.idl.Transaction +import org.rooftop.netx.idl.TransactionState +import org.springframework.context.event.EventListener +import reactor.core.publisher.Flux +import reactor.core.publisher.Mono + +abstract class AbstractTransactionDispatcher( + private val eventPublisher: EventPublisher, +) { + + @EventListener(SubscribeTransactionEvent::class) + fun subscribeStream(event: SubscribeTransactionEvent): Flux { + return receive(event).dispatch() + } + + protected abstract fun receive(event: SubscribeTransactionEvent): Flux + + private fun Flux.dispatch(): Flux { + return this.flatMap { + when (it.state) { + TransactionState.TRANSACTION_STATE_JOIN -> publishJoin(it) + TransactionState.TRANSACTION_STATE_COMMIT -> publishCommit(it) + TransactionState.TRANSACTION_STATE_ROLLBACK -> publishRollback(it) + TransactionState.TRANSACTION_STATE_START -> publishStart(it) + else -> error("Cannot find matched transaction state \"${it.state}\"") + } + } + } + + private fun publishJoin(it: Transaction): Mono { + return Mono.just(it) + .doOnNext { + eventPublisher.publish( + TransactionJoinEvent( + it.id, + it.replay, + it.serverId + ) + ) + } + } + + private fun publishCommit(it: Transaction): Mono { + return Mono.just(it) + .doOnNext { eventPublisher.publish(TransactionCommitEvent(it.id, it.serverId)) } + } + + private fun publishRollback(it: Transaction): Mono { + return Mono.just(it) + .doOnNext { + eventPublisher.publish( + TransactionRollbackEvent( + it.id, + it.replay, + it.serverId, + it.cause, + ) + ) + } + } + + private fun publishStart(it: Transaction): Mono { + return Mono.just(it) + .doOnNext { + eventPublisher.publish(TransactionStartEvent(it.id, it.replay, it.serverId)) + } + } +} diff --git a/src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt b/src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt index 17d7ff0..c57d744 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt @@ -2,5 +2,5 @@ package org.rooftop.netx.engine fun interface EventPublisher { - fun publish(event: SubscribeTransactionEvent) + fun publish(event: Any) } diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt index 1f4522b..59e9d89 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt @@ -1,93 +1,36 @@ package org.rooftop.netx.redis -import org.rooftop.netx.api.TransactionCommitEvent -import org.rooftop.netx.api.TransactionJoinEvent -import org.rooftop.netx.api.TransactionRollbackEvent -import org.rooftop.netx.api.TransactionStartEvent +import org.rooftop.netx.engine.AbstractTransactionDispatcher import org.rooftop.netx.engine.SubscribeTransactionEvent import org.rooftop.netx.idl.Transaction -import org.rooftop.netx.idl.TransactionState import org.springframework.context.ApplicationEventPublisher -import org.springframework.context.event.EventListener import org.springframework.data.redis.connection.ReactiveRedisConnectionFactory import org.springframework.data.redis.connection.stream.Consumer import org.springframework.data.redis.connection.stream.StreamOffset import org.springframework.data.redis.stream.StreamReceiver import reactor.core.publisher.Flux -import reactor.core.publisher.Mono import reactor.core.scheduler.Schedulers +import kotlin.time.Duration.Companion.milliseconds +import kotlin.time.toJavaDuration class RedisStreamTransactionDispatcher( - private val eventPublisher: ApplicationEventPublisher, + eventPublisher: ApplicationEventPublisher, private val streamGroup: String, private val nodeName: String, connectionFactory: ReactiveRedisConnectionFactory, -) { +) : AbstractTransactionDispatcher(SpringEventPublisher(eventPublisher)) { private val options = StreamReceiver.StreamReceiverOptions.builder() - .pollTimeout(java.time.Duration.ofMillis(100)) + .pollTimeout(100.milliseconds.toJavaDuration()) .build() private val receiver = StreamReceiver.create(connectionFactory, options) - @EventListener(SubscribeTransactionEvent::class) - fun subscribeStream(subscribeTransactionEvent: SubscribeTransactionEvent): Flux { + override fun receive(event: SubscribeTransactionEvent): Flux { return receiver.receiveAutoAck( Consumer.from(streamGroup, nodeName), - StreamOffset.fromStart(subscribeTransactionEvent.transactionId) + StreamOffset.fromStart(event.transactionId) ).publishOn(Schedulers.parallel()) .map { Transaction.parseFrom(it.value["data"]?.toByteArray()) } - .dispatch() - } - - private fun Flux.dispatch(): Flux { - return this.flatMap { - when (it.state) { - TransactionState.TRANSACTION_STATE_JOIN -> publishJoin(it) - TransactionState.TRANSACTION_STATE_COMMIT -> publishCommit(it) - TransactionState.TRANSACTION_STATE_ROLLBACK -> publishRollback(it) - TransactionState.TRANSACTION_STATE_START -> publishStart(it) - else -> error("Cannot find matched transaction state \"${it.state}\"") - } - } - } - - private fun publishJoin(it: Transaction): Mono { - return Mono.just(it) - .doOnNext { - eventPublisher.publishEvent( - TransactionJoinEvent( - it.id, - it.replay, - it.serverId - ) - ) - } - } - - private fun publishCommit(it: Transaction): Mono { - return Mono.just(it) - .doOnNext { eventPublisher.publishEvent(TransactionCommitEvent(it.id, it.serverId)) } - } - - private fun publishRollback(it: Transaction): Mono { - return Mono.just(it) - .doOnNext { - eventPublisher.publishEvent( - TransactionRollbackEvent( - it.id, - it.replay, - it.serverId, - it.cause, - ) - ) - } - } - - private fun publishStart(it: Transaction): Mono { - return Mono.just(it) - .doOnNext { - eventPublisher.publishEvent(TransactionStartEvent(it.id, it.replay, it.serverId)) - } } } diff --git a/src/main/kotlin/org/rooftop/netx/redis/SpringEventPublisher.kt b/src/main/kotlin/org/rooftop/netx/redis/SpringEventPublisher.kt index e09abde..28ad8ad 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/SpringEventPublisher.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/SpringEventPublisher.kt @@ -6,7 +6,7 @@ import org.springframework.context.ApplicationEventPublisher class SpringEventPublisher(private val eventPublisher: ApplicationEventPublisher) : EventPublisher { - override fun publish(event: SubscribeTransactionEvent) { + override fun publish(event: Any) { eventPublisher.publishEvent(event) } } From 8e176f9cdcb96f1fb542f3c8ae860ddd1e66bf90 Mon Sep 17 00:00:00 2001 From: devxb Date: Sat, 3 Feb 2024 01:33:36 +0900 Subject: [PATCH 12/19] =?UTF-8?q?fix:=20RedisStream=20=EA=B5=AC=EB=8F=85?= =?UTF-8?q?=EC=9D=84=20poll=20=EB=B0=A9=EC=8B=9D=EC=9D=B4=20=EC=95=84?= =?UTF-8?q?=EB=8B=8C=20subscirbe=20=EB=B0=A9=EC=8B=9D=EC=9C=BC=EB=A1=9C=20?= =?UTF-8?q?=EC=88=98=EC=A0=95=ED=95=98=EA=B3=A0=20=EC=B2=98=EC=9D=8C?= =?UTF-8?q?=EB=B6=80=ED=84=B0=20=EB=AA=A8=EB=93=A0=20=EB=8D=B0=EC=9D=B4?= =?UTF-8?q?=ED=84=B0=EB=A5=BC=20=EC=9D=BD=EC=96=B4=EC=98=A4=EB=8F=84?= =?UTF-8?q?=EB=A1=9D=20=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../redis/AutoConfigureRedisTransaction.kt | 2 +- .../redis/RedisStreamTransactionDispatcher.kt | 28 +++++++++++++------ .../redis/RedisStreamTransactionManager.kt | 6 ++-- .../netx/redis/RedisTransactionConfigurer.kt | 16 ++++++----- 4 files changed, 33 insertions(+), 19 deletions(-) diff --git a/src/main/kotlin/org/rooftop/netx/redis/AutoConfigureRedisTransaction.kt b/src/main/kotlin/org/rooftop/netx/redis/AutoConfigureRedisTransaction.kt index 44ca2fa..12cefa0 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/AutoConfigureRedisTransaction.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/AutoConfigureRedisTransaction.kt @@ -2,7 +2,7 @@ package org.rooftop.netx.redis import org.springframework.boot.autoconfigure.ImportAutoConfiguration -@ImportAutoConfiguration +@ImportAutoConfiguration(RedisTransactionConfigurer::class) @Target(AnnotationTarget.CLASS) @Retention(AnnotationRetention.RUNTIME) annotation class AutoConfigureRedisTransaction diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt index 59e9d89..3390a3c 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt @@ -6,31 +6,43 @@ import org.rooftop.netx.idl.Transaction import org.springframework.context.ApplicationEventPublisher import org.springframework.data.redis.connection.ReactiveRedisConnectionFactory import org.springframework.data.redis.connection.stream.Consumer +import org.springframework.data.redis.connection.stream.ReadOffset import org.springframework.data.redis.connection.stream.StreamOffset +import org.springframework.data.redis.core.ReactiveRedisTemplate import org.springframework.data.redis.stream.StreamReceiver import reactor.core.publisher.Flux import reactor.core.scheduler.Schedulers -import kotlin.time.Duration.Companion.milliseconds +import kotlin.time.Duration.Companion.hours import kotlin.time.toJavaDuration class RedisStreamTransactionDispatcher( eventPublisher: ApplicationEventPublisher, + connectionFactory: ReactiveRedisConnectionFactory, private val streamGroup: String, private val nodeName: String, - connectionFactory: ReactiveRedisConnectionFactory, + private val reactiveRedisTemplate: ReactiveRedisTemplate, ) : AbstractTransactionDispatcher(SpringEventPublisher(eventPublisher)) { private val options = StreamReceiver.StreamReceiverOptions.builder() - .pollTimeout(100.milliseconds.toJavaDuration()) + .pollTimeout(1.hours.toJavaDuration()) .build() private val receiver = StreamReceiver.create(connectionFactory, options) override fun receive(event: SubscribeTransactionEvent): Flux { - return receiver.receiveAutoAck( - Consumer.from(streamGroup, nodeName), - StreamOffset.fromStart(event.transactionId) - ).publishOn(Schedulers.parallel()) - .map { Transaction.parseFrom(it.value["data"]?.toByteArray()) } + return createGroupIfNotExists(event) + .flatMap { + receiver.receiveAutoAck( + Consumer.from(streamGroup, nodeName), + StreamOffset.create(event.transactionId, ReadOffset.from(">")) + ).publishOn(Schedulers.parallel()) + .map { Transaction.parseFrom(it.value["data"]?.toByteArray()) } + } + } + + private fun createGroupIfNotExists(event: SubscribeTransactionEvent): Flux { + return reactiveRedisTemplate.opsForStream() + .createGroup(event.transactionId, ReadOffset.from("0"), streamGroup) + .flatMapMany { Flux.just(it) } } } diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt index a72e9e6..6dac7c7 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt @@ -12,11 +12,11 @@ class RedisStreamTransactionManager( nodeId: Int, nodeName: String, applicationEventPublisher: ApplicationEventPublisher, - private val transactionServer: ReactiveRedisTemplate, + private val reactiveRedisTemplate: ReactiveRedisTemplate, ) : AbstractTransactionManager(nodeId, nodeName, SpringEventPublisher(applicationEventPublisher)) { override fun exists(transactionId: String): Mono { - return transactionServer.opsForStream() + return reactiveRedisTemplate.opsForStream() .range(transactionId, Range.open("-", "+")) .map { Transaction.parseFrom(it.value[DATA].toString().toByteArray()) } .next() @@ -35,7 +35,7 @@ class RedisStreamTransactionManager( } override fun publishTransaction(transactionId: String, transaction: Transaction): Mono { - return transactionServer.opsForStream() + return reactiveRedisTemplate.opsForStream() .add( Record.of(mapOf(DATA to transaction.toByteArray())) .withStreamKey(transactionId) diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt index 5a655c0..eb16a91 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt @@ -3,16 +3,16 @@ package org.rooftop.netx.redis import org.rooftop.netx.api.TransactionManager import org.rooftop.pay.infra.transaction.ByteArrayRedisSerializer import org.springframework.beans.factory.annotation.Value -import org.springframework.boot.autoconfigure.AutoConfiguration import org.springframework.context.ApplicationEventPublisher import org.springframework.context.annotation.Bean +import org.springframework.context.annotation.Configuration import org.springframework.data.redis.connection.ReactiveRedisConnectionFactory import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory import org.springframework.data.redis.core.ReactiveRedisTemplate import org.springframework.data.redis.serializer.RedisSerializationContext import org.springframework.data.redis.serializer.StringRedisSerializer -@AutoConfiguration +@Configuration class RedisTransactionConfigurer( @Value("\${netx.host}") private val host: String, @Value("\${netx.port}") private val port: String, @@ -22,33 +22,35 @@ class RedisTransactionConfigurer( private val applicationEventPublisher: ApplicationEventPublisher, ) { + @Bean fun redisStreamTransactionManager(): TransactionManager = RedisStreamTransactionManager( nodeId, nodeName, applicationEventPublisher, - transactionServer() + reactiveRedisTemplate() ) @Bean fun redisStreamTransactionDispatcher(): RedisStreamTransactionDispatcher = RedisStreamTransactionDispatcher( applicationEventPublisher, + reactiveRedisConnectionFactory(), group, nodeName, - transactionServerConnectionFactory() + reactiveRedisTemplate() ) @Bean - fun transactionServer(): ReactiveRedisTemplate { + fun reactiveRedisTemplate(): ReactiveRedisTemplate { val builder = RedisSerializationContext.newSerializationContext( StringRedisSerializer() ) val context = builder.value(byteArrayRedisSerializer()).build() - return ReactiveRedisTemplate(transactionServerConnectionFactory(), context) + return ReactiveRedisTemplate(reactiveRedisConnectionFactory(), context) } @Bean @@ -57,7 +59,7 @@ class RedisTransactionConfigurer( } @Bean - fun transactionServerConnectionFactory(): ReactiveRedisConnectionFactory { + fun reactiveRedisConnectionFactory(): ReactiveRedisConnectionFactory { val port: String = System.getProperty("netx.port") ?: port return LettuceConnectionFactory(host, port.toInt()) From 84a70fc235ab6356de06366cf9df50abb38a8017 Mon Sep 17 00:00:00 2001 From: devxb Date: Sun, 4 Feb 2024 10:10:25 +0900 Subject: [PATCH 13/19] =?UTF-8?q?refactor:=20replay=20=ED=95=84=EB=93=9C?= =?UTF-8?q?=EA=B0=80=20rollback=20=EC=9D=B4=EB=B2=A4=ED=8A=B8=EC=97=90?= =?UTF-8?q?=EB=A7=8C=20=EC=A1=B4=EC=9E=AC=ED=95=98=EB=8F=84=EB=A1=9D=20?= =?UTF-8?q?=EC=88=98=EC=A0=95=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt | 1 - src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt | 1 - .../org/rooftop/netx/engine/AbstractTransactionDispatcher.kt | 3 +-- 3 files changed, 1 insertion(+), 4 deletions(-) diff --git a/src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt index 8899322..9e46e09 100644 --- a/src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt +++ b/src/main/kotlin/org/rooftop/netx/api/TransactionJoinEvent.kt @@ -2,6 +2,5 @@ package org.rooftop.netx.api data class TransactionJoinEvent( val transactionId: String, - val replay: String, val nodeName: String, ) diff --git a/src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt index cfdc7bf..738514a 100644 --- a/src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt +++ b/src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt @@ -2,6 +2,5 @@ package org.rooftop.netx.api data class TransactionStartEvent( val transactionId: String, - val replay: String, val nodeName: String, ) diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt index bcf4e20..b15af24 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt @@ -39,7 +39,6 @@ abstract class AbstractTransactionDispatcher( eventPublisher.publish( TransactionJoinEvent( it.id, - it.replay, it.serverId ) ) @@ -68,7 +67,7 @@ abstract class AbstractTransactionDispatcher( private fun publishStart(it: Transaction): Mono { return Mono.just(it) .doOnNext { - eventPublisher.publish(TransactionStartEvent(it.id, it.replay, it.serverId)) + eventPublisher.publish(TransactionStartEvent(it.id, it.serverId)) } } } From c4e9dcbffc73a480e643fe3c361ee80e95763828 Mon Sep 17 00:00:00 2001 From: devxb Date: Sun, 4 Feb 2024 10:50:59 +0900 Subject: [PATCH 14/19] =?UTF-8?q?refactor:=20exists=20=EB=A9=94=EC=86=8C?= =?UTF-8?q?=EB=93=9C=20=EA=B5=AC=ED=98=84=EC=9D=84=20engine=EB=A0=88?= =?UTF-8?q?=EC=9D=B4=EC=96=B4=EB=A1=9C=20=EC=98=AC=EB=A6=AC=EA=B3=A0,=20co?= =?UTF-8?q?ntext=EB=A5=BC=20=EC=84=B8=ED=8C=85=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../netx/engine/AbstractTransactionManager.kt | 45 ++++++++++++++----- .../redis/RedisStreamTransactionManager.kt | 14 +----- 2 files changed, 34 insertions(+), 25 deletions(-) diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt index f5033ae..00b66ba 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt @@ -13,7 +13,7 @@ abstract class AbstractTransactionManager( private val transactionIdGenerator: TransactionIdGenerator = TransactionIdGenerator(nodeId), ) : TransactionManager { - override fun start(replay: String): Mono { + final override fun start(replay: String): Mono { return startTransaction(replay) .subscribeTransaction() .contextWrite { it.put(CONTEXT_TX_KEY, transactionIdGenerator.generate()) } @@ -31,7 +31,7 @@ abstract class AbstractTransactionManager( } } - override fun join(transactionId: String, replay: String): Mono { + final override fun join(transactionId: String, replay: String): Mono { return exists(transactionId) .joinTransaction(replay) .subscribeTransaction() @@ -40,13 +40,13 @@ abstract class AbstractTransactionManager( private fun Mono.joinTransaction(replay: String): Mono { return flatMap { transactionId -> - publishTransaction(transactionId, transaction { - id = transactionId - serverId = nodeName - this.replay = replay - state = TransactionState.TRANSACTION_STATE_JOIN - }) - } + publishTransaction(transactionId, transaction { + id = transactionId + serverId = nodeName + this.replay = replay + state = TransactionState.TRANSACTION_STATE_JOIN + }) + } } private fun Mono.subscribeTransaction(): Mono { @@ -55,7 +55,7 @@ abstract class AbstractTransactionManager( } } - override fun rollback(transactionId: String, cause: String): Mono { + final override fun rollback(transactionId: String, cause: String): Mono { return exists(transactionId) .publishTransaction(transaction { id = transactionId @@ -66,7 +66,7 @@ abstract class AbstractTransactionManager( .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) } } - override fun commit(transactionId: String): Mono { + final override fun commit(transactionId: String): Mono { return exists(transactionId) .publishTransaction(transaction { id = transactionId @@ -76,13 +76,34 @@ abstract class AbstractTransactionManager( .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) } } + final override fun exists(transactionId: String): Mono { + return findAnyTransaction(transactionId) + .switchIfEmpty( + Mono.error { + IllegalStateException("Cannot find exists transaction id \"$transactionId\"") + } + ).transformTransactionId() + .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) } + } + + protected abstract fun findAnyTransaction(transactionId: String): Mono + + protected fun Mono<*>.transformTransactionId(): Mono { + return this.flatMap { + Mono.deferContextual { Mono.just(it["transactionId"]) } + } + } + private fun Mono.publishTransaction(transaction: Transaction): Mono { return this.flatMap { publishTransaction(it, transaction) } } - abstract fun publishTransaction(transactionId: String, transaction: Transaction): Mono + protected abstract fun publishTransaction( + transactionId: String, + transaction: Transaction, + ): Mono private companion object { private const val CONTEXT_TX_KEY = "transactionId" diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt index 6dac7c7..e2604d9 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt @@ -15,23 +15,11 @@ class RedisStreamTransactionManager( private val reactiveRedisTemplate: ReactiveRedisTemplate, ) : AbstractTransactionManager(nodeId, nodeName, SpringEventPublisher(applicationEventPublisher)) { - override fun exists(transactionId: String): Mono { + override fun findAnyTransaction(transactionId: String): Mono { return reactiveRedisTemplate.opsForStream() .range(transactionId, Range.open("-", "+")) .map { Transaction.parseFrom(it.value[DATA].toString().toByteArray()) } .next() - .switchIfEmpty( - Mono.error { - IllegalStateException("Cannot find exists transaction id \"$transactionId\"") - } - ) - .transformTransactionId() - } - - private fun Mono<*>.transformTransactionId(): Mono { - return this.flatMap { - Mono.deferContextual { Mono.just(it["transactionId"]) } - } } override fun publishTransaction(transactionId: String, transaction: Transaction): Mono { From 6ffdf9f080e9d29f4f761aa1fe4187c438cdabfb Mon Sep 17 00:00:00 2001 From: devxb Date: Sun, 4 Feb 2024 10:51:58 +0900 Subject: [PATCH 15/19] =?UTF-8?q?refactor:=20transformTransactionId=20?= =?UTF-8?q?=EC=9D=B4=EB=A6=84=EC=9D=84=20mapTransactionId=EB=A1=9C=20?= =?UTF-8?q?=EB=B3=80=EA=B2=BD=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../org/rooftop/netx/engine/AbstractTransactionManager.kt | 4 ++-- .../org/rooftop/netx/redis/RedisStreamTransactionManager.kt | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt index 00b66ba..160a0d7 100644 --- a/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt @@ -82,13 +82,13 @@ abstract class AbstractTransactionManager( Mono.error { IllegalStateException("Cannot find exists transaction id \"$transactionId\"") } - ).transformTransactionId() + ).mapTransactionId() .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) } } protected abstract fun findAnyTransaction(transactionId: String): Mono - protected fun Mono<*>.transformTransactionId(): Mono { + protected fun Mono<*>.mapTransactionId(): Mono { return this.flatMap { Mono.deferContextual { Mono.just(it["transactionId"]) } } diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt index e2604d9..8d0c31f 100644 --- a/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt +++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt @@ -28,7 +28,7 @@ class RedisStreamTransactionManager( Record.of(mapOf(DATA to transaction.toByteArray())) .withStreamKey(transactionId) ) - .transformTransactionId() + .mapTransactionId() } private companion object { From 583c92d028cac4f13bb88aa83b16ec4211522403 Mon Sep 17 00:00:00 2001 From: devxb Date: Sun, 4 Feb 2024 10:54:06 +0900 Subject: [PATCH 16/19] =?UTF-8?q?test:=20RedisStreamTransactionManager?= =?UTF-8?q?=EC=9D=98=20Test=EB=A5=BC=20=EC=9E=91=EC=84=B1=ED=95=9C?= =?UTF-8?q?=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../org/rooftop/netx/redis/EventCapture.kt | 24 +++ .../org/rooftop/netx/redis/RedisContainer.kt | 20 +++ .../RedisStreamTransactionManagerTest.kt | 151 ++++++++++++++++++ src/test/resources/application.properties | 5 + 4 files changed, 200 insertions(+) create mode 100644 src/test/kotlin/org/rooftop/netx/redis/EventCapture.kt create mode 100644 src/test/kotlin/org/rooftop/netx/redis/RedisContainer.kt create mode 100644 src/test/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManagerTest.kt create mode 100644 src/test/resources/application.properties diff --git a/src/test/kotlin/org/rooftop/netx/redis/EventCapture.kt b/src/test/kotlin/org/rooftop/netx/redis/EventCapture.kt new file mode 100644 index 0000000..54d23d9 --- /dev/null +++ b/src/test/kotlin/org/rooftop/netx/redis/EventCapture.kt @@ -0,0 +1,24 @@ +package org.rooftop.netx.redis + +import org.springframework.boot.test.context.TestComponent +import org.springframework.context.event.EventListener +import kotlin.reflect.KClass + +@TestComponent +class EventCapture { + + private val eventCapture: MutableMap, Long> = mutableMapOf() + + fun clear() { + eventCapture.clear() + } + + fun capturedCount(type: KClass<*>): Long { + return eventCapture[type] ?: 0 + } + + @EventListener(Any::class) + fun captureEvent(type: Any) { + eventCapture[type::class] = eventCapture.getOrDefault(type::class, 0) + 1 + } +} diff --git a/src/test/kotlin/org/rooftop/netx/redis/RedisContainer.kt b/src/test/kotlin/org/rooftop/netx/redis/RedisContainer.kt new file mode 100644 index 0000000..29e5625 --- /dev/null +++ b/src/test/kotlin/org/rooftop/netx/redis/RedisContainer.kt @@ -0,0 +1,20 @@ +package org.rooftop.netx.redis + +import org.springframework.boot.test.context.TestConfiguration +import org.testcontainers.containers.GenericContainer +import org.testcontainers.utility.DockerImageName + +@TestConfiguration +class RedisContainer { + init { + val redis: GenericContainer<*> = GenericContainer(DockerImageName.parse("redis:7.2.3")) + .withExposedPorts(6379) + + redis.start() + + System.setProperty( + "netx.port", + redis.getMappedPort(6379).toString() + ) + } +} diff --git a/src/test/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManagerTest.kt b/src/test/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManagerTest.kt new file mode 100644 index 0000000..5272ee1 --- /dev/null +++ b/src/test/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManagerTest.kt @@ -0,0 +1,151 @@ +package org.rooftop.netx.redis + +import io.kotest.assertions.nondeterministic.eventually +import io.kotest.core.annotation.DisplayName +import io.kotest.core.spec.style.DescribeSpec +import io.kotest.matchers.shouldBe +import org.rooftop.netx.api.* +import org.springframework.test.context.ContextConfiguration +import org.springframework.test.context.TestPropertySource +import reactor.test.StepVerifier +import kotlin.time.Duration.Companion.minutes + +@AutoConfigureRedisTransaction +@ContextConfiguration( + classes = [ + EventCapture::class, + RedisContainer::class, + ] +) +@DisplayName("RedisStreamTransactionManager 클래스의") +@TestPropertySource("classpath:application.properties") +internal class RedisStreamTransactionManagerTest( + private val eventCapture: EventCapture, + private val transactionManager: TransactionManager, +) : DescribeSpec({ + + beforeEach { + eventCapture.clear() + } + + describe("start 메소드는") { + context("replay 를 입력받으면,") { + it("트랜잭션을 시작하고 transaction-id를 반환한다.") { + transactionManager.start(REPLAY).subscribe() + + eventually(5.minutes) { + eventCapture.capturedCount(TransactionStartEvent::class) shouldBe 1 + } + } + } + + context("서로 다른 id의 트랜잭션이 여러번 시작되어도") { + it("모두 읽을 수 있다.") { + transactionManager.start(REPLAY).block() + transactionManager.start(REPLAY).block() + + eventually(5.minutes) { + eventCapture.capturedCount(TransactionStartEvent::class) shouldBe 2 + } + } + } + } + + describe("join 메소드는") { + context("존재하는 transactionId를 입력받으면,") { + val transactionId = transactionManager.start(REPLAY).block()!! + + it("트랜잭션에 참여한다.") { + transactionManager.join(transactionId, REPLAY).subscribe() + + eventually(5.minutes) { + eventCapture.capturedCount(TransactionJoinEvent::class) shouldBe 1 + } + } + } + + context("존재하지 않는 transactionId를 입력받으면,") { + it("IllegalStateException 을 던진다.") { + val result = transactionManager.join(NOT_EXIST_TX_ID, REPLAY) + + StepVerifier.create(result) + .verifyErrorMessage("Cannot find exists transaction id \"$NOT_EXIST_TX_ID\"") + } + } + } + + describe("exists 메소드는") { + context("존재하는 transactionId를 입력받으면,") { + val transactionId = transactionManager.start(REPLAY).block()!! + + it("트랜잭션 id를 반환한다.") { + val result = transactionManager.exists(transactionId) + + StepVerifier.create(result) + .expectNext(transactionId) + .verifyComplete() + } + } + + context("존재하지 않는 transactionId를 입력받으면,") { + it("IllegalStateException 을 던진다.") { + val result = transactionManager.exists(NOT_EXIST_TX_ID) + + StepVerifier.create(result) + .verifyErrorMessage("Cannot find exists transaction id \"$NOT_EXIST_TX_ID\"") + } + } + } + + describe("commit 메소드는") { + context("존재하는 transactionId를 입력받으면,") { + val transactionId = transactionManager.start(REPLAY).block()!! + + it("commit 메시지를 publish 한다") { + transactionManager.commit(transactionId).block() + + eventually(5.minutes) { + eventCapture.capturedCount(TransactionCommitEvent::class) + } + } + } + + context("존재하지 않는 transactionId를 입력받으면,") { + it("IllegalStateException 을 던진다.") { + val result = transactionManager.commit(NOT_EXIST_TX_ID) + + StepVerifier.create(result) + .verifyErrorMessage("Cannot find exists transaction id \"$NOT_EXIST_TX_ID\"") + } + } + } + + describe("rollback 메소드는") { + context("존재하는 transactionId를 입력받으면,") { + val transactionId = transactionManager.start(REPLAY).block()!! + + it("rollback 메시지를 publish 한다") { + transactionManager.rollback(transactionId, "rollback occured for test").block() + + eventually(5.minutes) { + eventCapture.capturedCount(TransactionRollbackEvent::class) + } + } + } + + context("존재하지 않는 transactionId를 입력받으면,") { + it("IllegalStateException 을 던진다.") { + val result = transactionManager.commit(NOT_EXIST_TX_ID) + + StepVerifier.create(result) + .verifyErrorMessage("Cannot find exists transaction id \"$NOT_EXIST_TX_ID\"") + } + } + } +}) { + + private companion object { + private const val REPLAY = "REPLAY" + private const val NOT_EXIST_TX_ID = "NOT_EXISTS_TX_ID" + } +} diff --git a/src/test/resources/application.properties b/src/test/resources/application.properties new file mode 100644 index 0000000..8b3384d --- /dev/null +++ b/src/test/resources/application.properties @@ -0,0 +1,5 @@ +netx.host=localhost +netx.port=6379 +netx.group=netx-group +netx.node-id=1 +netx.node-name=netx-node From 968d7668551a331f7d1d370ee02498a9de6fe740 Mon Sep 17 00:00:00 2001 From: devxb Date: Sun, 4 Feb 2024 13:32:41 +0900 Subject: [PATCH 17/19] =?UTF-8?q?build:=20jitpack=20=EB=B0=B0=ED=8F=AC=20?= =?UTF-8?q?=ED=94=8C=EB=9F=AC=EA=B7=B8=EC=9D=B8=EC=9D=84=20=EC=84=A4?= =?UTF-8?q?=EC=A0=95=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- build.gradle | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/build.gradle b/build.gradle index 9829923..b37c964 100644 --- a/build.gradle +++ b/build.gradle @@ -8,6 +8,7 @@ plugins { id "io.spring.dependency-management" version "${springDependencyManagementVersion}" id "org.sonarqube" version "${sonarcloudVersion}" id "com.google.protobuf" version "${protobufPluginVersion}" + id "maven-publish" } group = "${group}" @@ -23,6 +24,14 @@ dependencyManagement { } } +publishing { + publications { + maven(MavenPublication) { + from components.java + } + } +} + apply from: "gradle/mq.gradle" apply from: "gradle/test.gradle" apply from: "gradle/core.gradle" From 84edd4f95af2adc02800071b72bcc8b9f1f20345 Mon Sep 17 00:00:00 2001 From: devxb Date: Sun, 4 Feb 2024 17:48:00 +0900 Subject: [PATCH 18/19] =?UTF-8?q?docs:=20=EC=82=AC=EC=9A=A9=EB=B2=95?= =?UTF-8?q?=EA=B3=BC=20=EB=8B=A4=EC=9A=B4=EB=A1=9C=EB=93=9C=20=EB=B0=A9?= =?UTF-8?q?=EB=B2=95=EC=9D=84=20=EC=9E=91=EC=84=B1=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- README.md | 121 ++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 121 insertions(+) diff --git a/README.md b/README.md index 7fcece4..290d592 100644 --- a/README.md +++ b/README.md @@ -1,3 +1,124 @@ # Netx > Distributed transaction library based on Choreography + +![version 0.1.0](https://img.shields.io/badge/version-0.1.0-black?labelColor=black&style=flat-square) + +Choreography 방식으로 구현된 분산 트랜잭션 라이브러리 입니다. +`Netx` 는 다음 기능을 제공합니다. +1. [Reactor](https://projectreactor.io/) 기반의 완전한 비동기 트랜잭션 관리 +2. Redis-stream 기반의 트랜잭션 관리 +3. 여러 노드가 중복 트랜잭션 이벤트를 수신하는 문제 방지 + +## How to use + +Netx는 스프링 환경에서 사용할 수 있으며, 아래와 같이 `@AutoConfigureRedisTransaction` 어노테이션을 붙이는것으로 손쉽게 사용할 수 있습니다. + +```kotlin +@SpringBootApplication +@AutoConfigureRedisTransaction +@EnableAutoConfiguration(exclude = [RedisReactiveAutoConfiguration::class]) +class Application { + + companion object { + @JvmStatic + fun main(vararg args: String) { + SpringApplication.run(Application::class.java, *args) + } + } +} +``` + +`@AutoconfigureRedisTransaction` 어노테이션으로 자동 구성할 경우 netx는 아래 프로퍼티를 사용해 메시지 큐와 커넥션을 맺습니다. + +#### Properties +| key | example | description | +|------------------|----------|------------------------------------------------------------------------------------------------------------------------------------| +| **netx.host** | localhost | 트랜잭션 관리에 사용할 메시지 큐 의 host url 입니다. (ex. redis host) | +| **netx.port** | 6379 | 트랜잭션 관리에 사용할 메시지 큐의 port 입니다. | +| **netx.group** | pay-group | 분산 노드의 그룹입니다. 트랜잭션 이벤트는 같은 그룹내 하나의 노드로만 전송됩니다. | +| **netx.node-id** | 1 | id 생성에 사용될 식별자입니다. 모든 서버는 반드시 다른 id를 할당받아야 하며, 1~256 만큼의 id를 설정할 수 있습니다. _`중복된 id 생성을 방지하기위해 twitter snowflake 알고리즘으로 id를 생성합니다.`_ | +| **netx.node-name** | pay-1 | _`$netx.group`_ 에 참여할 서버의 이름입니다. 같은 그룹내에 중복된 이름이 존재하면 안됩니다. | + + +### Usage example + +#### Scenario1. Start pay transaction + +```kotlin +fun pay(param: Any): Mono { + return transactionManager.start("paid=1000") // Start distributed transaction and publish transaction start event + .flatMap { transactionId -> + service.pay(param) + .doOnError { throwable -> + transactionManager.rollback(transactionId, throwable.message) // Publish rollback event to all transaction joined node + } + }.doOnSuccess { transactionId -> + transactionManager.commit(transactionId) // Publish commit event to all transaction joined node + } +} +``` + +#### Scenario2. Join order transaction + +```kotlin +fun order(param: Any): Mono { + return transactionManager.join(param.transactionId, "orderId=1:state=PENDING") // join exists distributed transaction and publish transaction join event + .flatMap { transactionId -> + service.order(param) + .doOnError { throwable -> + transactionManager.rollback(transactionId, throwable.message) + } + }.doOnSuccess { transactionId -> + transactionManager.commit(transactionId) + } +} +``` + +#### Scenario3. Check exists transaction + +```kotlin +fun exists(param: Any): Mono { + return transactionManager.exists(param.transactionId) // Find any transaction has ever been started +} +``` + +#### Scenario4. Handle transaction event + +다른 분산서버가 (혹은 자기자신이) transactionManager를 통해서 트랜잭션을 시작하거나 트랜잭션 상태를 변경했을때, 호출한 메소드에 맞는 트랜잭션 이벤트를 발행합니다. + +```kotlin + +@EventListener(TransactionStartEvent::class) +fun handleTransactionStartEvent(event: TransactionStartEvent) { + // ... +} + +@EventListener(TransactionJoinEvent::class) +fun handleTransactionJoinEvent(event: TransactionJoinEvent) { + // ... +} + +@EventListener(TransactionCommitEvent::class) +fun handleTransactionCommitEvent(event: TransactionCommitEvent) { + // ... +} + +@EventListener(TransactionRollbackEvent::class) +fun handleTransactionRollbackEvent(event: TransactionRollbackEvent) { + // ... +} +``` + +## Download + +```groovy + +repositories { + maven { url "https://jitpack.io" } +} + +dependencies { + implementation "com.github.rooftop-msa:netx:${version}" +} +``` From a6baa7df021dc0af8d26ed1e8f35fbf10cbf73da Mon Sep 17 00:00:00 2001 From: xb205 <62425964+devxb@users.noreply.github.com> Date: Sun, 4 Feb 2024 17:49:46 +0900 Subject: [PATCH 19/19] =?UTF-8?q?docs:=20=EB=B6=84=EC=82=B0=20=ED=8A=B8?= =?UTF-8?q?=EB=9E=9C=EC=9E=AD=EC=85=98=20=EC=9E=91=EB=8F=99=EA=B3=BC?= =?UTF-8?q?=EC=A0=95=20gif=EB=A5=BC=20=EC=B6=94=EA=B0=80=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- README.md | 2 ++ 1 file changed, 2 insertions(+) diff --git a/README.md b/README.md index 290d592..f3b78e9 100644 --- a/README.md +++ b/README.md @@ -4,6 +4,8 @@ ![version 0.1.0](https://img.shields.io/badge/version-0.1.0-black?labelColor=black&style=flat-square) + + Choreography 방식으로 구현된 분산 트랜잭션 라이브러리 입니다. `Netx` 는 다음 기능을 제공합니다. 1. [Reactor](https://projectreactor.io/) 기반의 완전한 비동기 트랜잭션 관리