diff --git a/README.md b/README.md
index 7fcece4..f3b78e9 100644
--- a/README.md
+++ b/README.md
@@ -1,3 +1,126 @@
# Netx
> Distributed transaction library based on Choreography
+
+
+
+
+
+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}"
+}
+```
diff --git a/build.gradle b/build.gradle
index a1553c4..b37c964 100644
--- a/build.gradle
+++ b/build.gradle
@@ -1,11 +1,14 @@
+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}"
+ id "maven-publish"
}
group = "${group}"
@@ -15,6 +18,20 @@ repositories {
mavenCentral()
}
+dependencyManagement {
+ imports {
+ mavenBom SpringBootPlugin.BOM_COORDINATES
+ }
+}
+
+publishing {
+ publications {
+ maven(MavenPublication) {
+ from components.java
+ }
+ }
+}
+
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 573ba4a..e1f481d 160000
--- a/idl
+++ b/idl
@@ -1 +1 @@
-Subproject commit 573ba4ae4a9e9f7d1d4372f98b6dbc3b287c7ba2
+Subproject commit e1f481d6f34e879b82487603c11cb7a57cc6a6ab
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/TransactionCommitEvent.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionCommitEvent.kt
new file mode 100644
index 0000000..ef5e477
--- /dev/null
+++ b/src/main/kotlin/org/rooftop/netx/api/TransactionCommitEvent.kt
@@ -0,0 +1,6 @@
+package org.rooftop.netx.api
+
+data class TransactionCommitEvent(
+ 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
new file mode 100644
index 0000000..9e46e09
--- /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 nodeName: String,
+)
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/api/TransactionRollbackEvent.kt b/src/main/kotlin/org/rooftop/netx/api/TransactionRollbackEvent.kt
new file mode 100644
index 0000000..ffa9d2a
--- /dev/null
+++ b/src/main/kotlin/org/rooftop/netx/api/TransactionRollbackEvent.kt
@@ -0,0 +1,8 @@
+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
new file mode 100644
index 0000000..738514a
--- /dev/null
+++ b/src/main/kotlin/org/rooftop/netx/api/TransactionStartEvent.kt
@@ -0,0 +1,6 @@
+package org.rooftop.netx.api
+
+data class TransactionStartEvent(
+ val transactionId: 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
new file mode 100644
index 0000000..b15af24
--- /dev/null
+++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionDispatcher.kt
@@ -0,0 +1,73 @@
+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.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.serverId))
+ }
+ }
+}
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..160a0d7
--- /dev/null
+++ b/src/main/kotlin/org/rooftop/netx/engine/AbstractTransactionManager.kt
@@ -0,0 +1,111 @@
+package org.rooftop.netx.engine
+
+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(
+ nodeId: Int,
+ private val nodeName: String,
+ private val eventPublisher: EventPublisher,
+ private val transactionIdGenerator: TransactionIdGenerator = TransactionIdGenerator(nodeId),
+) : TransactionManager {
+
+ final override fun start(replay: String): Mono {
+ return startTransaction(replay)
+ .subscribeTransaction()
+ .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 = nodeName
+ this.replay = replay
+ this.state = TransactionState.TRANSACTION_STATE_START
+ })
+ }
+ }
+
+ final override fun join(transactionId: String, replay: String): Mono {
+ return exists(transactionId)
+ .joinTransaction(replay)
+ .subscribeTransaction()
+ .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) }
+ }
+
+ 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
+ })
+ }
+ }
+
+ private fun Mono.subscribeTransaction(): Mono {
+ return this.doOnSuccess {
+ eventPublisher.publish(SubscribeTransactionEvent(it))
+ }
+ }
+
+ final override fun rollback(transactionId: String, cause: String): Mono {
+ return exists(transactionId)
+ .publishTransaction(transaction {
+ id = transactionId
+ serverId = nodeName
+ state = TransactionState.TRANSACTION_STATE_ROLLBACK
+ this.cause = cause
+ })
+ .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) }
+ }
+
+ final override fun commit(transactionId: String): Mono {
+ return exists(transactionId)
+ .publishTransaction(transaction {
+ id = transactionId
+ serverId = nodeName
+ state = TransactionState.TRANSACTION_STATE_COMMIT
+ })
+ .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\"")
+ }
+ ).mapTransactionId()
+ .contextWrite { it.put(CONTEXT_TX_KEY, transactionId) }
+ }
+
+ protected abstract fun findAnyTransaction(transactionId: String): Mono
+
+ protected fun Mono<*>.mapTransactionId(): Mono {
+ return this.flatMap {
+ Mono.deferContextual { Mono.just(it["transactionId"]) }
+ }
+ }
+
+ private fun Mono.publishTransaction(transaction: Transaction): Mono {
+ return this.flatMap {
+ publishTransaction(it, transaction)
+ }
+ }
+
+ 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/engine/EventPublisher.kt b/src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt
new file mode 100644
index 0000000..c57d744
--- /dev/null
+++ b/src/main/kotlin/org/rooftop/netx/engine/EventPublisher.kt
@@ -0,0 +1,6 @@
+package org.rooftop.netx.engine
+
+fun interface EventPublisher {
+
+ fun publish(event: Any)
+}
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..79e8447
--- /dev/null
+++ b/src/main/kotlin/org/rooftop/netx/engine/SubscribeTransactionEvent.kt
@@ -0,0 +1,5 @@
+package org.rooftop.netx.engine
+
+data class SubscribeTransactionEvent(
+ val transactionId: String
+)
diff --git a/src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt b/src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt
new file mode 100644
index 0000000..968a13a
--- /dev/null
+++ b/src/main/kotlin/org/rooftop/netx/engine/TransactionIdGenerator.kt
@@ -0,0 +1,12 @@
+package org.rooftop.netx.engine
+
+import com.github.f4b6a3.tsid.TsidFactory
+
+class TransactionIdGenerator(
+ nodeId: Int,
+ private val tsidFactory: TsidFactory = TsidFactory.newInstance256(nodeId),
+) {
+
+ fun generate(): String = tsidFactory.create().toLong().toString()
+}
+
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..12cefa0
--- /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(RedisTransactionConfigurer::class)
+@Target(AnnotationTarget.CLASS)
+@Retention(AnnotationRetention.RUNTIME)
+annotation class AutoConfigureRedisTransaction
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
new file mode 100644
index 0000000..3390a3c
--- /dev/null
+++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionDispatcher.kt
@@ -0,0 +1,48 @@
+package org.rooftop.netx.redis
+
+import org.rooftop.netx.engine.AbstractTransactionDispatcher
+import org.rooftop.netx.engine.SubscribeTransactionEvent
+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.hours
+import kotlin.time.toJavaDuration
+
+class RedisStreamTransactionDispatcher(
+ eventPublisher: ApplicationEventPublisher,
+ connectionFactory: ReactiveRedisConnectionFactory,
+ private val streamGroup: String,
+ private val nodeName: String,
+ private val reactiveRedisTemplate: ReactiveRedisTemplate,
+) : AbstractTransactionDispatcher(SpringEventPublisher(eventPublisher)) {
+
+ private val options = StreamReceiver.StreamReceiverOptions.builder()
+ .pollTimeout(1.hours.toJavaDuration())
+ .build()
+
+ private val receiver = StreamReceiver.create(connectionFactory, options)
+
+ override fun receive(event: SubscribeTransactionEvent): Flux {
+ 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
new file mode 100644
index 0000000..8d0c31f
--- /dev/null
+++ b/src/main/kotlin/org/rooftop/netx/redis/RedisStreamTransactionManager.kt
@@ -0,0 +1,37 @@
+package org.rooftop.netx.redis
+
+import org.rooftop.netx.engine.AbstractTransactionManager
+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(
+ nodeId: Int,
+ nodeName: String,
+ applicationEventPublisher: ApplicationEventPublisher,
+ private val reactiveRedisTemplate: ReactiveRedisTemplate,
+) : AbstractTransactionManager(nodeId, nodeName, SpringEventPublisher(applicationEventPublisher)) {
+
+ override fun findAnyTransaction(transactionId: String): Mono {
+ return reactiveRedisTemplate.opsForStream()
+ .range(transactionId, Range.open("-", "+"))
+ .map { Transaction.parseFrom(it.value[DATA].toString().toByteArray()) }
+ .next()
+ }
+
+ override fun publishTransaction(transactionId: String, transaction: Transaction): Mono {
+ return reactiveRedisTemplate.opsForStream()
+ .add(
+ Record.of(mapOf(DATA to transaction.toByteArray()))
+ .withStreamKey(transactionId)
+ )
+ .mapTransactionId()
+ }
+
+ private companion object {
+ private const val DATA = "data"
+ }
+}
diff --git a/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt
new file mode 100644
index 0000000..eb16a91
--- /dev/null
+++ b/src/main/kotlin/org/rooftop/netx/redis/RedisTransactionConfigurer.kt
@@ -0,0 +1,67 @@
+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.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
+
+@Configuration
+class RedisTransactionConfigurer(
+ @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(
+ nodeId,
+ nodeName,
+ applicationEventPublisher,
+ reactiveRedisTemplate()
+ )
+
+ @Bean
+ fun redisStreamTransactionDispatcher(): RedisStreamTransactionDispatcher =
+ RedisStreamTransactionDispatcher(
+ applicationEventPublisher,
+ reactiveRedisConnectionFactory(),
+ group,
+ nodeName,
+ reactiveRedisTemplate()
+ )
+
+ @Bean
+ fun reactiveRedisTemplate(): ReactiveRedisTemplate {
+ val builder = RedisSerializationContext.newSerializationContext(
+ StringRedisSerializer()
+ )
+
+ val context = builder.value(byteArrayRedisSerializer()).build()
+
+ return ReactiveRedisTemplate(reactiveRedisConnectionFactory(), context)
+ }
+
+ @Bean
+ fun byteArrayRedisSerializer(): ByteArrayRedisSerializer {
+ return ByteArrayRedisSerializer()
+ }
+
+ @Bean
+ fun reactiveRedisConnectionFactory(): ReactiveRedisConnectionFactory {
+ val port: String = System.getProperty("netx.port") ?: port
+
+ return LettuceConnectionFactory(host, port.toInt())
+ }
+}
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..28ad8ad
--- /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: Any) {
+ eventPublisher.publishEvent(event)
+ }
+}
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 e69de29..0000000
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
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