feature: init kotlin compose multiplatforme
Tests / build (pull_request) Successful in 7m44s
Tests / test (pull_request) Failing after 10m37s
Tests / lint (pull_request) Successful in 14m30s

This commit is contained in:
2026-08-06 21:28:57 +02:00
parent 505cfe38f0
commit 774e80d9d5
196 changed files with 1297 additions and 688 deletions
+105
View File
@@ -0,0 +1,105 @@
import org.jlleitschuh.gradle.ktlint.KtlintExtension
val ktorVersion: Provider<String> = providers.gradleProperty("ktor_version")
val kotlinVersion: Provider<String> = providers.gradleProperty("kotlin_version")
val kotlinSerializationVersion: Provider<String> = providers.gradleProperty("kotlin_serialization_version")
val logbackVersion: Provider<String> = providers.gradleProperty("logback_version")
val koinVersion: Provider<String> = providers.gradleProperty("koin_version")
val kotlinLoggingVersion: Provider<String> = providers.gradleProperty("kotlin_logging_version")
val kotestVersion: Provider<String> = providers.gradleProperty("kotest_version")
plugins {
application
kotlin("jvm")
id("io.ktor.plugin") version "3.5.1"
id("org.jetbrains.kotlin.plugin.serialization")
id("org.jlleitschuh.gradle.ktlint") version "14.2.0"
}
group = "io.github.flecomte"
application {
mainClass.set("eventDemo.ApplicationKt")
val isDevelopment: Boolean = project.ext.has("development")
applicationDefaultJvmArgs = listOf("-Dio.ktor.development=$isDevelopment")
}
configure<KtlintExtension> {
version.set("1.8.0")
}
ktlint {
reporters {
reporter(org.jlleitschuh.gradle.ktlint.reporter.ReporterType.CHECKSTYLE)
}
}
repositories {
mavenCentral()
}
java {
toolchain {
languageVersion = JavaLanguageVersion.of(21)
}
}
kotlin {
compilerOptions {
freeCompilerArgs.add("-opt-in=kotlin.uuid.ExperimentalUuidApi")
}
}
tasks.withType<Test>().configureEach {
useJUnitPlatform()
jvmArgs("-Djdk.attach.allowAttachSelf=true", "-XX:+EnableDynamicAgentLoading")
// Dynamic self-attach (used by MockK/ByteBuddy) times out in Docker containers because the
// SIGQUIT-triggered AttachListener handshake never completes there. Loading the byte-buddy
// agent jar statically via -javaagent avoids the attach handshake entirely: MockK detects the
// pre-installed Instrumentation instance and skips dynamic attach.
doFirst {
val agentJar =
classpath.files.firstOrNull { it.name.startsWith("byte-buddy-agent") }
?: error("byte-buddy-agent jar not found on test classpath")
jvmArgs("-javaagent:$agentJar")
}
}
dependencies {
implementation(project(":shared"))
implementation("io.ktor:ktor-server-core-jvm")
implementation("io.ktor:ktor-server-auth-jvm")
implementation("io.ktor:ktor-server-auth-jwt-jvm")
implementation("io.ktor:ktor-server-auto-head-response-jvm")
implementation("io.ktor:ktor-server-resources")
implementation("io.ktor:ktor-server-content-negotiation-jvm")
implementation("io.ktor:ktor-serialization-kotlinx-json-jvm")
implementation("io.ktor:ktor-server-websockets-jvm")
implementation("io.ktor:ktor-server-cors-jvm")
implementation("io.ktor:ktor-server-host-common-jvm")
implementation("io.ktor:ktor-server-status-pages-jvm")
implementation("io.ktor:ktor-server-netty-jvm")
implementation("io.ktor:ktor-server-data-conversion")
implementation("io.ktor:ktor-client-content-negotiation")
implementation("io.ktor:ktor-client-auth")
implementation("ch.qos.logback:logback-classic:${logbackVersion.get()}")
implementation("io.insert-koin:koin-ktor:${koinVersion.get()}")
implementation("io.insert-koin:koin-logger-slf4j:${koinVersion.get()}")
implementation("io.github.oshai:kotlin-logging-jvm:${kotlinLoggingVersion.get()}")
implementation("org.jetbrains.kotlinx:kotlinx-serialization-json-jvm:${kotlinSerializationVersion.get()}")
implementation("org.jetbrains.kotlinx:kotlinx-datetime:0.6.2")
implementation("org.postgresql:postgresql:42.7.13")
implementation("com.zaxxer:HikariCP:6.3.0")
implementation("com.rabbitmq:amqp-client:5.25.0")
implementation("com.password4j:password4j:1.8.4")
// Force version of sub library (for security)
implementation("commons-codec:commons-codec:1.13")
testImplementation("io.kotest:kotest-extensions-koin:${kotestVersion.get()}")
testImplementation("org.jetbrains.kotlin:kotlin-test-junit:${kotlinVersion.get()}")
testImplementation("io.ktor:ktor-server-test-host-jvm:${ktorVersion.get()}")
testImplementation("io.kotest:kotest-runner-junit5:${kotestVersion.get()}")
testImplementation("io.mockk:mockk:1.14.11")
testImplementation("com.tngtech.archunit:archunit-junit5:1.3.0")
}
@@ -0,0 +1,6 @@
package eventDemo
import io.ktor.server.netty.EngineMain
fun main(args: Array<String>): Unit =
EngineMain.main(args)
@@ -0,0 +1,44 @@
package eventDemo.configuration
import io.ktor.server.config.ApplicationConfig
data class Configuration(
val jwtSecret: String,
val postgresql: Postgresql,
val rabbitmq: RabbitMQ,
) {
data class Postgresql(
val url: String,
val username: String,
val password: String,
)
data class RabbitMQ(
val url: String,
val port: Int,
val username: String,
val password: String,
)
}
val ApplicationConfig.configuration
get() =
Configuration(
jwtSecret = getProperty("jwt.secret"),
postgresql =
Configuration.Postgresql(
url = getProperty("postgresql.url"),
username = getProperty("postgresql.username"),
password = getProperty("postgresql.password"),
),
rabbitmq =
Configuration.RabbitMQ(
url = getProperty("rabbitmq.url"),
port = getProperty("rabbitmq.port").toInt(),
username = getProperty("rabbitmq.username"),
password = getProperty("rabbitmq.password"),
),
)
private fun ApplicationConfig.getProperty(path: String): String =
propertyOrNull(path)?.getString() ?: error("You must set the $path")
@@ -0,0 +1,14 @@
package eventDemo.configuration
import eventDemo.contexts.auth.infrastructure.configure.configureAuthDi
import eventDemo.contexts.game.infrastructure.configuration.injections.application.configureGameDIApplication
import eventDemo.contexts.game.infrastructure.configuration.injections.infrastructure.configureGameDIInfrastructure
import org.koin.dsl.module
fun appKoinModule(config: Configuration) =
module {
configureDIDataSource(config)
configureAuthDi()
configureGameDIInfrastructure()
configureGameDIApplication()
}
@@ -0,0 +1,48 @@
package eventDemo.configuration
import com.rabbitmq.client.ConnectionFactory
import com.zaxxer.hikari.HikariConfig
import com.zaxxer.hikari.HikariDataSource
import org.koin.core.module.Module
import org.koin.core.scope.Scope
import org.koin.core.scope.ScopeCallback
import org.koin.dsl.bind
import javax.sql.DataSource
fun Module.configureDIDataSource(config: Configuration) {
// PostgreSQL (for EventStore)
single {
hikariDataSource(config)
.apply {
registerCallback(
object : ScopeCallback {
override fun onScopeClose(scope: Scope) {
close()
}
},
)
}
} bind DataSource::class
// RabbitMQ (for EventBus)
factory {
ConnectionFactory().apply {
host = config.rabbitmq.url
port = config.rabbitmq.port
username = config.rabbitmq.username
password = config.rabbitmq.password
}
}
}
private fun hikariDataSource(config: Configuration): HikariDataSource =
HikariConfig()
.apply {
jdbcUrl = config.postgresql.url
username = config.postgresql.username
password = config.postgresql.password
maximumPoolSize = 10
minimumIdle = 10
}.let {
HikariDataSource(it)
}
@@ -0,0 +1,16 @@
package eventDemo.configuration
import io.ktor.server.application.Application
import io.ktor.server.application.install
import org.koin.ktor.plugin.Koin
import org.koin.logger.slf4jLogger
fun Application.configureKoin() {
install(Koin) {
slf4jLogger()
modules(
appKoinModule(environment.config.configuration),
)
}
}
@@ -0,0 +1,11 @@
package eventDemo.configuration
import eventDemo.contexts.auth.infrastructure.configure.configureAuth
import eventDemo.contexts.game.infrastructure.configuration.ktor.configureUno
import io.ktor.server.application.Application
fun Application.configure() {
configureKoin()
configureAuth()
configureUno()
}
@@ -0,0 +1,24 @@
package eventDemo.contexts.auth.application.eventStores
import eventDemo.contexts.auth.application.ports.UserEventStore
import eventDemo.contexts.auth.domain.User
import eventDemo.shared.ids.UserId
class UserEventStoreRepository(
val eventStore: UserEventStore,
) : UserRepository {
override fun get(id: UserId): User? {
val events =
eventStore
.getStream(id)
.readAll()
if (events.isEmpty()) {
return null
}
return events.let { User.loadFromHistory(it) }
}
override fun save(user: User) {
eventStore.append(user.recordedEvents)
}
}
@@ -0,0 +1,10 @@
package eventDemo.contexts.auth.application.eventStores
import eventDemo.contexts.auth.domain.User
import eventDemo.shared.ids.UserId
interface UserRepository {
fun get(id: UserId): User?
fun save(user: User)
}
@@ -0,0 +1,7 @@
package eventDemo.contexts.auth.application.ports
import eventDemo.contexts.auth.domain.events.UserEvent
import eventDemo.libs.eventSource.eventStore.EventStore
import eventDemo.shared.ids.UserId
interface UserEventStore : EventStore<UserEvent, UserId>
@@ -0,0 +1,14 @@
package eventDemo.contexts.auth.application.ports
import eventDemo.contexts.auth.infrastructure.persistence.projection.UserProjection
interface UserProjectionRepository {
fun getByUsername(username: String): UserProjection?
fun save(user: UserProjection)
fun getUserIfPasswordIsValid(
username: String,
rawPassword: String,
): UserProjection?
}
@@ -0,0 +1,39 @@
package eventDemo.contexts.auth.domain
import eventDemo.contexts.auth.domain.events.NewUserCreatedEvent
import eventDemo.contexts.auth.domain.events.UserEvent
import eventDemo.shared.ids.UserId
import kotlinx.serialization.Serializable
@Serializable
data class User(
val id: UserId,
val username: String,
val password: String,
val version: Int,
val recordedEvents: Set<UserEvent>,
) {
companion object {
fun createNewUser(
username: String,
password: String,
): User =
apply(NewUserCreatedEvent(username, password, version = 1))
fun apply(event: NewUserCreatedEvent): User =
User(
id = event.aggregateId,
username = event.username,
password = event.password,
version = event.version,
recordedEvents = setOf(event),
)
fun loadFromHistory(events: Set<UserEvent>): User? =
events.fold(null as User?) { acc, event ->
when (event) {
is NewUserCreatedEvent -> apply(event)
}
}
}
}
@@ -0,0 +1,18 @@
package eventDemo.contexts.auth.domain.events
import eventDemo.shared.ids.EventId
import eventDemo.shared.ids.UserId
import kotlinx.datetime.Clock
import kotlinx.datetime.Instant
import kotlinx.serialization.Serializable
@Serializable
class NewUserCreatedEvent(
val username: String,
val password: String,
override val version: Int,
override val createdAt: Instant = Clock.System.now(),
override val aggregateId: UserId = UserId(),
) : UserEvent {
override val eventId: EventId = EventId()
}
@@ -0,0 +1,8 @@
package eventDemo.contexts.auth.domain.events
import eventDemo.libs.eventSource.Event
import eventDemo.shared.ids.UserId
import kotlinx.serialization.Serializable
@Serializable
sealed interface UserEvent : Event<UserId>
@@ -0,0 +1,13 @@
package eventDemo.contexts.auth.infrastructure
import com.password4j.Hash
import com.password4j.Password
internal fun hashPassword(password: String): Hash =
Password.hash(password).addRandomSalt().withArgon2()
internal fun checkPassword(
password: String,
hash: Hash,
): Boolean =
Password.check(password, hash)
@@ -0,0 +1,8 @@
package eventDemo.contexts.auth.infrastructure.configure
import io.ktor.server.application.Application
fun Application.configureAuth() {
configureKtorAuth()
configureAuthRoutes()
}
@@ -0,0 +1,17 @@
package eventDemo.contexts.auth.infrastructure.configure
import eventDemo.contexts.auth.application.eventStores.UserEventStoreRepository
import eventDemo.contexts.auth.application.eventStores.UserRepository
import eventDemo.contexts.auth.application.ports.UserEventStore
import eventDemo.contexts.auth.application.ports.UserProjectionRepository
import eventDemo.contexts.auth.infrastructure.persistence.eventStore.UserEventStoreInPostgresql
import eventDemo.contexts.auth.infrastructure.persistence.projection.UserProjectionRepositoryInPostgresql
import org.koin.core.module.Module
import org.koin.core.module.dsl.singleOf
import org.koin.dsl.bind
fun Module.configureAuthDi() {
singleOf(::UserEventStoreRepository) bind UserRepository::class
singleOf(::UserEventStoreInPostgresql) bind UserEventStore::class
singleOf(::UserProjectionRepositoryInPostgresql) bind UserProjectionRepository::class
}
@@ -0,0 +1,19 @@
package eventDemo.contexts.auth.infrastructure.configure
import eventDemo.configuration.configuration
import eventDemo.contexts.auth.application.eventStores.UserEventStoreRepository
import eventDemo.contexts.auth.application.ports.UserProjectionRepository
import eventDemo.contexts.auth.infrastructure.rest.createUserRoute
import eventDemo.contexts.auth.infrastructure.rest.loginRoute
import io.ktor.server.application.Application
import io.ktor.server.routing.routing
import org.koin.ktor.ext.get
fun Application.configureAuthRoutes() {
val userRepository = get<UserEventStoreRepository>()
val userProjectionRepository = get<UserProjectionRepository>()
routing {
createUserRoute(userRepository)
loginRoute(environment.config.configuration.jwtSecret, userProjectionRepository)
}
}
@@ -0,0 +1,65 @@
package eventDemo.contexts.auth.infrastructure.configure
import com.auth0.jwt.JWT
import com.auth0.jwt.algorithms.Algorithm
import eventDemo.configuration.configuration
import eventDemo.contexts.auth.domain.User
import eventDemo.contexts.auth.infrastructure.persistence.projection.UserProjection
import eventDemo.shared.ids.UserId
import io.ktor.http.HttpStatusCode
import io.ktor.server.application.Application
import io.ktor.server.auth.authentication
import io.ktor.server.auth.jwt.JWTPrincipal
import io.ktor.server.auth.jwt.jwt
import io.ktor.server.response.respond
import java.util.Date
fun Application.configureKtorAuth() {
val jwtSecret = environment.config.configuration.jwtSecret
authentication {
jwt {
realm = "Play card game"
verifier(
JWT
.require(Algorithm.HMAC256(jwtSecret))
.withIssuer(JWT_ISSUER)
.build(),
)
validate { credential ->
if (credential.payload
.getClaim("username")
.asString()
.isNotEmpty()
) {
JWTPrincipal(credential.payload)
} else {
null
}
}
challenge { _, _ ->
call.respond(HttpStatusCode.Unauthorized, "Token is not valid or has expired")
}
}
}
}
private const val JWT_ISSUER = "PlayCardGame"
fun UserProjection.makeJwt(jwtSecret: String): String =
makeJwt(jwtSecret, id, username)
fun User.makeJwt(jwtSecret: String): String =
makeJwt(jwtSecret, id, username)
fun makeJwt(
jwtSecret: String,
id: UserId,
username: String,
): String =
JWT
.create()
.withIssuer(JWT_ISSUER)
.withClaim("username", username)
.withClaim("userid", id.toString())
.withExpiresAt(Date(System.currentTimeMillis() + 60000))
.sign(Algorithm.HMAC256(jwtSecret))
@@ -0,0 +1,14 @@
package eventDemo.contexts.auth.infrastructure.persistence.eventStore
import eventDemo.contexts.auth.application.ports.UserEventStore
import eventDemo.contexts.auth.domain.events.UserEvent
import eventDemo.libs.eventSource.eventStore.EventStore
import eventDemo.libs.eventSource.eventStore.EventStoreInMemory
import eventDemo.shared.ids.UserId
/**
* A stream to publish and read the user events.
*/
class UserEventStoreInMemory :
UserEventStore,
EventStore<UserEvent, UserId> by EventStoreInMemory()
@@ -0,0 +1,22 @@
package eventDemo.contexts.auth.infrastructure.persistence.eventStore
import eventDemo.contexts.auth.application.ports.UserEventStore
import eventDemo.contexts.auth.domain.events.UserEvent
import eventDemo.libs.eventSource.eventStore.EventStore
import eventDemo.libs.eventSource.eventStore.EventStoreInPostgresql
import eventDemo.shared.ids.UserId
import kotlinx.serialization.json.Json
import javax.sql.DataSource
/**
* A stream to publish and read the user events.
*/
class UserEventStoreInPostgresql(
dataSource: DataSource,
) : UserEventStore,
EventStore<UserEvent, UserId> by EventStoreInPostgresql(
dataSource,
{ Json.encodeToString(it) },
{ Json.decodeFromString(it) },
"auth.user_event_stream",
)
@@ -0,0 +1,9 @@
package eventDemo.contexts.auth.infrastructure.persistence.projection
import eventDemo.shared.ids.UserId
data class UserProjection(
val id: UserId,
val username: String,
val password: String,
)
@@ -0,0 +1,62 @@
package eventDemo.contexts.auth.infrastructure.persistence.projection
import eventDemo.contexts.auth.application.ports.UserProjectionRepository
import eventDemo.contexts.auth.infrastructure.checkPassword
import eventDemo.contexts.auth.infrastructure.hashPassword
import eventDemo.shared.ids.UserId
import javax.sql.DataSource
import kotlin.uuid.Uuid
import kotlin.uuid.toJavaUuid
class UserProjectionRepositoryInPostgresql(
val dataSource: DataSource,
) : UserProjectionRepository {
override fun getByUsername(username: String): UserProjection? =
dataSource.connection
.prepareStatement(
"""
select id, username
from auth."user"
where id = ?;
""".trimIndent(),
).use {
it.setObject(1, username)
it.executeQuery()
}.use { resultSet ->
if (resultSet.next()) {
UserProjection(
id = UserId(Uuid.parse(resultSet.getString("id"))),
username = resultSet.getString("username"),
password = resultSet.getString("password"),
)
} else {
null
}
}
override fun save(user: UserProjection) {
dataSource.connection.use { connection ->
connection
.prepareStatement(
"""
insert into auth.user (id, username)
values (?, ?)
""".trimIndent(),
).use {
it.setObject(1, user.id.id.toJavaUuid())
it.setString(2, user.username)
it.executeUpdate()
}
}
}
override fun getUserIfPasswordIsValid(
username: String,
rawPassword: String,
): UserProjection? {
val user = getByUsername(username) ?: return null
val isValid = checkPassword(rawPassword, hashPassword(user.password))
if (!isValid) return null
return user
}
}
@@ -0,0 +1,24 @@
package eventDemo.contexts.auth.infrastructure.rest
import eventDemo.contexts.auth.application.ports.UserProjectionRepository
import eventDemo.contexts.auth.infrastructure.configure.makeJwt
import io.ktor.http.HttpStatusCode
import io.ktor.server.response.respond
import io.ktor.server.routing.Route
import io.ktor.server.routing.post
fun Route.loginRoute(
jwtSecret: String,
userProjectionRepository: UserProjectionRepository,
) {
post("login/{username}") {
val username = call.parameters["username"]!!
val rawPassword = call.parameters["password"]!!
val userProjection =
userProjectionRepository.getUserIfPasswordIsValid(username, rawPassword)
?: return@post call.respond(HttpStatusCode.BadRequest)
call.respond(hashMapOf("token" to userProjection.makeJwt(jwtSecret)))
}
}
@@ -0,0 +1,41 @@
package eventDemo.contexts.auth.infrastructure.rest
import eventDemo.contexts.auth.application.eventStores.UserRepository
import eventDemo.contexts.auth.domain.User
import eventDemo.contexts.auth.infrastructure.hashPassword
import io.ktor.resources.Resource
import io.ktor.server.auth.authenticate
import io.ktor.server.response.respond
import io.ktor.server.routing.Route
import io.ktor.server.routing.post
import kotlinx.serialization.Serializable
@Serializable
@Resource("/users")
class Users {
@Serializable
@Resource("/create")
class Create(
val username: String,
val password: String,
)
}
/**
* API routes to show all games.
*/
fun Route.createUserRoute(userRepository: UserRepository) {
authenticate {
// Create a new User, and return there ID
post<Users.Create> {
val passwordHash = hashPassword(it.password)
val user = User.createNewUser(it.username, passwordHash.result)
userRepository.save(user)
call.respond(
object {
val id = user.id.toString()
},
)
}
}
}
@@ -0,0 +1,37 @@
package eventDemo.contexts.game.application.channels
import eventDemo.contexts.game.application.notification.CommandSubscriber
import eventDemo.contexts.game.application.notification.EventToNotificationSubscriber
import eventDemo.shared.game.command.GameCommand
import eventDemo.shared.game.notification.Notification
import eventDemo.shared.ids.GameId
import eventDemo.shared.ids.UserId
import kotlinx.coroutines.DelicateCoroutinesApi
import kotlinx.coroutines.channels.ReceiveChannel
import kotlinx.coroutines.channels.SendChannel
class GameChannelsSubscriber(
private val eventToNotificationSubscriber: EventToNotificationSubscriber,
private val commandSubscriber: CommandSubscriber,
) {
@DelicateCoroutinesApi
fun subscribePlayerToGameChannels(
gameId: GameId,
userId: UserId,
incomingCommandChannel: ReceiveChannel<GameCommand>,
sendNotificationChannel: SendChannel<Notification>,
) {
val sub =
eventToNotificationSubscriber.subscribeToEventsAndSendNotification(
gameId = gameId,
currentUserId = userId,
outgoingFrameChannel = sendNotificationChannel,
)
commandSubscriber
.subscribe(
currentUserId = userId,
incomingFrameChannel = incomingCommandChannel,
).invokeOnCompletion { sub.close() }
}
}
@@ -0,0 +1,5 @@
package eventDemo.contexts.game.application.command.handlers
class CommandException(
override val message: String,
) : Exception(message)
@@ -0,0 +1,31 @@
package eventDemo.contexts.game.application.command.handlers
import eventDemo.shared.game.command.GameCommand
import eventDemo.shared.game.command.JoinTheGameCommand
import eventDemo.shared.game.command.PlayCardCommand
import eventDemo.shared.game.command.ReadyToPlayCommand
import eventDemo.shared.game.command.TakeCartFromDrawPileCommand
import eventDemo.shared.ids.GameId
import java.util.Collections
class GameCommandHandlerDispatcher(
private val playCardHandler: PlayCardHandler,
private val readyToPlayHandler: ReadyToPlayHandler,
private val joinTheGameHandler: JoinTheGameHandler,
private val takeCartFromDrawPileHandler: TakeCartFromDrawPileHandler,
) {
companion object {
val lock: MutableMap<GameId, String> = Collections.synchronizedMap(mutableMapOf())
}
fun dispatch(command: GameCommand) {
synchronized(lock.getOrPut(command.payload.aggregateId) { command.payload.aggregateId.toString() }) {
when (command) {
is JoinTheGameCommand -> joinTheGameHandler.handle(command)
is ReadyToPlayCommand -> readyToPlayHandler.handle(command)
is PlayCardCommand -> playCardHandler.handle(command)
is TakeCartFromDrawPileCommand -> takeCartFromDrawPileHandler.handle(command)
}
}
}
}
@@ -0,0 +1,58 @@
package eventDemo.contexts.game.application.command.handlers
import eventDemo.contexts.game.application.eventStores.GameRepository
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.contexts.game.domain.game.gameState.Game
import eventDemo.libs.eventSource.eventStore.VersionConflictException
import eventDemo.shared.command.Command
import eventDemo.shared.game.command.GameCommand
import io.github.oshai.kotlinlogging.KotlinLogging
import kotlin.reflect.KClass
sealed interface CommandHandler<C : Command> {
fun handle(command: C)
}
abstract class GameEventManager(
private val gameRepository: GameRepository,
private val gameEventBus: GameEventBus,
) {
private val logger = KotlinLogging.logger {}
fun GameCommand.getGame(): Game =
gameRepository.get(payload.aggregateId) ?: error("Game not found")
fun GameEvent.getGame(): Game =
gameRepository.get(aggregateId) ?: error("Game not found")
@Throws(VersionConflictException::class)
protected fun Game.saveEvents(): Game {
gameRepository.save(this)
return this
}
protected fun Game.publishEvents(): Game {
gameEventBus.publish(recordedEvents)
return this
}
protected inline fun <reified G : Game> Game.isStatusOrFail(message: String): G =
this as? G ?: throw CommandException(message)
protected fun <T> retry(
mapAttempts: Int = 5,
block: () -> T,
): T =
try {
block()
} catch (e: VersionConflictException) {
if (mapAttempts > 0) {
logger.warn { "retry after version conflict (attempts left: $mapAttempts)" }
retry(mapAttempts - 1, block)
} else {
logger.error { "Version conflict retry failed" }
throw e
}
}
}
@@ -0,0 +1,31 @@
package eventDemo.contexts.game.application.command.handlers
import eventDemo.contexts.auth.application.eventStores.UserRepository
import eventDemo.contexts.game.application.eventStores.GameRepository
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.domain.game.gameState.GameCreated
import eventDemo.shared.game.command.JoinTheGameCommand
/**
* A command to perform an action to play a new card
*/
class JoinTheGameHandler(
gameRepository: GameRepository,
gameEventBus: GameEventBus,
private val userRepository: UserRepository,
) : GameEventManager(gameRepository, gameEventBus),
CommandHandler<JoinTheGameCommand> {
override fun handle(command: JoinTheGameCommand) {
val user = userRepository.get(command.userId) ?: error("User with id ${command.userId} doesn't exist")
retry {
command
.getGame()
.isStatusOrFail<GameCreated>("The game is started")
.userJoinTheGame(
userId = command.userId,
name = user.username,
).saveEvents()
.publishEvents()
}
}
}
@@ -0,0 +1,27 @@
package eventDemo.contexts.game.application.command.handlers
import eventDemo.contexts.game.application.eventStores.GameRepository
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.domain.game.gameState.GameStarted
import eventDemo.shared.game.command.PlayCardCommand
/**
* A command to perform an action to play a new card
*/
class PlayCardHandler(
gameRepository: GameRepository,
gameEventBus: GameEventBus,
) : GameEventManager(gameRepository, gameEventBus),
CommandHandler<PlayCardCommand> {
override fun handle(command: PlayCardCommand) {
command
.getGame()
.isStatusOrFail<GameStarted>("The game is not started")
.playTheCard(
card = command.payload.card,
playerId = command.payload.playerId,
chosenColor = command.payload.chosenColor,
).saveEvents()
.publishEvents()
}
}
@@ -0,0 +1,24 @@
package eventDemo.contexts.game.application.command.handlers
import eventDemo.contexts.game.application.eventStores.GameRepository
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.domain.game.gameState.GameCreated
import eventDemo.shared.game.command.ReadyToPlayCommand
/**
* A command to set as ready to play
*/
class ReadyToPlayHandler(
gameRepository: GameRepository,
gameEventBus: GameEventBus,
) : GameEventManager(gameRepository, gameEventBus),
CommandHandler<ReadyToPlayCommand> {
override fun handle(command: ReadyToPlayCommand) {
command
.getGame()
.isStatusOrFail<GameCreated>("The game is started")
.setReadyPlayer(command.payload.playerId)
.saveEvents()
.publishEvents()
}
}
@@ -0,0 +1,26 @@
package eventDemo.contexts.game.application.command.handlers
import eventDemo.contexts.game.application.eventStores.GameRepository
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.domain.game.gameState.GameStarted
import eventDemo.shared.game.command.TakeCartFromDrawPileCommand
/**
* A command to draw card on draw pile.
*
* Is can be triggered when you cannot play any card in your hand.
*/
class TakeCartFromDrawPileHandler(
gameRepository: GameRepository,
gameEventBus: GameEventBus,
) : GameEventManager(gameRepository, gameEventBus),
CommandHandler<TakeCartFromDrawPileCommand> {
override fun handle(command: TakeCartFromDrawPileCommand) {
command
.getGame()
.isStatusOrFail<GameStarted>("The game is not started")
.playerTakeCartFromDrawPile(command.payload.playerId, 1)
.saveEvents()
.publishEvents()
}
}
@@ -0,0 +1,26 @@
package eventDemo.contexts.game.application.eventStores
import eventDemo.contexts.game.application.ports.GameEventStore
import eventDemo.contexts.game.domain.game.gameState.Game
import eventDemo.libs.eventSource.eventStore.VersionConflictException
import eventDemo.shared.ids.GameId
class GameEventStoreRepository(
val eventStore: GameEventStore,
) : GameRepository {
override fun get(id: GameId): Game? {
val events =
eventStore
.getStream(id)
.readAll()
if (events.isEmpty()) {
return null
}
return events.let { Game.loadFromHistory(it) }
}
@Throws(VersionConflictException::class)
override fun save(game: Game) {
eventStore.append(game.recordedEvents)
}
}
@@ -0,0 +1,20 @@
package eventDemo.contexts.game.application.eventStores
import eventDemo.contexts.game.domain.game.gameState.Game
import eventDemo.contexts.game.domain.game.gameState.GameCreated
import eventDemo.contexts.game.domain.game.gameState.GameInit
import eventDemo.libs.eventSource.eventStore.VersionConflictException
import eventDemo.shared.ids.GameId
interface GameRepository {
fun get(id: GameId): Game?
@Throws(VersionConflictException::class)
fun save(game: Game)
fun getOrCreate(gameId: GameId): Game =
get(gameId) ?: create(gameId)
fun create(gameId: GameId = GameId()): GameCreated =
GameInit.createNewGame(gameId).also { save(it) }
}
@@ -0,0 +1,29 @@
package eventDemo.contexts.game.application.logging
import io.github.oshai.kotlinlogging.withLoggingContext
inline fun <T> withLoggingContext(
vararg pair: Pair<LoggingContextKeys, *>,
body: () -> T,
): T =
withLoggingContext(
*pair
.map {
it.first.name to it.second.toString()
}.toTypedArray(),
restorePrevious = true,
body = body,
)
// inline fun withLoggingContext(
// vararg pair: Pair<LoggingContextKeys, *>,
// body: () -> Unit,
// ) =
// withLoggingContext(
// *pair
// .map {
// it.first.name to it.second.toString()
// }.toTypedArray(),
// restorePrevious = true,
// body = body,
// )
@@ -0,0 +1,9 @@
package eventDemo.contexts.game.application.logging
enum class LoggingContextKeys {
CurrentUserId,
Notification,
Game,
Event,
Command,
}
@@ -0,0 +1,127 @@
package eventDemo.contexts.game.application.notification
import eventDemo.contexts.game.domain.events.CardIsPlayedEvent
import eventDemo.contexts.game.domain.events.DrawFilledWithDiscardEvent
import eventDemo.contexts.game.domain.events.GameCreatedEvent
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.contexts.game.domain.events.GameStartedEvent
import eventDemo.contexts.game.domain.events.NewPlayerEvent
import eventDemo.contexts.game.domain.events.PlayerActionEvent
import eventDemo.contexts.game.domain.events.PlayerHaveDrawCardEvent
import eventDemo.contexts.game.domain.events.PlayerReadyEvent
import eventDemo.contexts.game.domain.events.PlayerWinEvent
import eventDemo.contexts.game.domain.game.gameState.Game
import eventDemo.contexts.game.domain.game.gameState.GameStarted
import eventDemo.shared.game.notification.ItsTheTurnOfNotification
import eventDemo.shared.game.notification.Notification
import eventDemo.shared.game.notification.PilesShuffledNotification
import eventDemo.shared.game.notification.PlayerAsJoinTheGameNotification
import eventDemo.shared.game.notification.PlayerAsPlayACardNotification
import eventDemo.shared.game.notification.PlayerHavePassNotification
import eventDemo.shared.game.notification.PlayerWasReadyNotification
import eventDemo.shared.game.notification.PlayerWinNotification
import eventDemo.shared.game.notification.TheGameWasStartedNotification
import eventDemo.shared.game.notification.WelcomeToTheGameNotification
import eventDemo.shared.game.notification.YourNewCardNotification
import eventDemo.shared.ids.UserId
import io.github.oshai.kotlinlogging.KotlinLogging
import io.github.oshai.kotlinlogging.withLoggingContext
private val logger = KotlinLogging.logger {}
fun GameEvent.toNotification(
game: Game,
currentUserId: UserId,
): Iterable<Notification> =
Iterable {
iterator {
context(iterator: SequenceScope<Notification>)
suspend fun Notification.send() {
withLoggingContext("notification" to (this).toString()) {
logger.info { "Notification sent" }
iterator.yield(this)
}
}
fun PlayerActionEvent.isFromCurrentUser(): Boolean =
game.players.get(currentUserId).id == playerId
when (this@toNotification) {
is GameCreatedEvent -> {
// Nothing to send
}
is DrawFilledWithDiscardEvent -> {
PilesShuffledNotification().send()
}
is NewPlayerEvent -> {
if (this@toNotification.player.userId != currentUserId) {
PlayerAsJoinTheGameNotification(
player = this@toNotification.player,
).send()
} else {
WelcomeToTheGameNotification(
players = game.players,
).send()
}
}
is CardIsPlayedEvent -> {
PlayerAsPlayACardNotification(
playerId = this@toNotification.playerId,
card = this@toNotification.card,
).send()
if (game is GameStarted) {
ItsTheTurnOfNotification(
player = game.nextPlayer,
).send()
}
}
is GameStartedEvent -> {
TheGameWasStartedNotification(
hand =
game.players
.get(currentUserId)
.hand.cards,
).send()
if (game is GameStarted) {
ItsTheTurnOfNotification(player = game.nextPlayer)
.send()
}
}
is PlayerHaveDrawCardEvent -> {
if (this@toNotification.isFromCurrentUser()) {
YourNewCardNotification(
cards = this@toNotification.takenCards,
).send()
} else {
PlayerHavePassNotification(
playerId = this@toNotification.playerId,
).send()
}
if (game is GameStarted) {
ItsTheTurnOfNotification(player = game.nextPlayer)
.send()
}
}
is PlayerReadyEvent -> {
PlayerWasReadyNotification(
playerId = this@toNotification.playerId,
).send()
}
is PlayerWinEvent -> {
PlayerWinNotification(
playerId = this@toNotification.playerId,
).send()
}
}
}
}
@@ -0,0 +1,72 @@
package eventDemo.contexts.game.application.notification
import eventDemo.contexts.game.application.command.handlers.GameCommandHandlerDispatcher
import eventDemo.contexts.game.application.eventStores.GameRepository
import eventDemo.contexts.game.application.logging.LoggingContextKeys.Command
import eventDemo.contexts.game.application.logging.LoggingContextKeys.CurrentUserId
import eventDemo.contexts.game.application.logging.LoggingContextKeys.Event
import eventDemo.contexts.game.application.logging.LoggingContextKeys.Game
import eventDemo.contexts.game.application.logging.LoggingContextKeys.Notification
import eventDemo.contexts.game.application.logging.withLoggingContext
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.libs.bus.Bus
import eventDemo.libs.command.CommandUnicityChecker
import eventDemo.shared.game.command.GameCommand
import eventDemo.shared.game.notification.Notification
import eventDemo.shared.ids.GameId
import eventDemo.shared.ids.UserId
import kotlinx.coroutines.DelicateCoroutinesApi
import kotlinx.coroutines.GlobalScope
import kotlinx.coroutines.Job
import kotlinx.coroutines.channels.ReceiveChannel
import kotlinx.coroutines.channels.SendChannel
import kotlinx.coroutines.channels.trySendBlocking
import kotlinx.coroutines.launch
class EventToNotificationSubscriber(
private val gameEventBus: GameEventBus,
private val gameRepository: GameRepository,
) {
fun subscribeToEventsAndSendNotification(
gameId: GameId,
currentUserId: UserId,
outgoingFrameChannel: SendChannel<Notification>,
): Bus.Subscription =
withLoggingContext(CurrentUserId to currentUserId) {
gameEventBus.subscribe { event ->
val game = gameRepository.get(gameId) ?: error("Game not found")
withLoggingContext(Event to event, Game to game) {
event
.toNotification(
game = game,
currentUserId = currentUserId,
).forEach { notification ->
withLoggingContext(Notification to notification) {
outgoingFrameChannel.trySendBlocking(notification)
}
}
}
}
}
}
class CommandSubscriber(
private val gameCommandHandlerDispatcher: GameCommandHandlerDispatcher,
) {
private val controller = CommandUnicityChecker<GameCommand>()
@DelicateCoroutinesApi
fun subscribe(
currentUserId: UserId,
incomingFrameChannel: ReceiveChannel<GameCommand>,
): Job =
GlobalScope.launch {
for (command in incomingFrameChannel) {
withLoggingContext(CurrentUserId to currentUserId, Command to command) {
controller.runOnlyOnce(command) {
gameCommandHandlerDispatcher.dispatch(command)
}
}
}
}
}
@@ -0,0 +1,6 @@
package eventDemo.contexts.game.application.ports
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.libs.bus.Bus
interface GameEventBus : Bus<GameEvent>
@@ -0,0 +1,7 @@
package eventDemo.contexts.game.application.ports
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.libs.eventSource.eventStore.EventStore
import eventDemo.shared.ids.GameId
interface GameEventStore : EventStore<GameEvent, GameId>
@@ -0,0 +1,14 @@
package eventDemo.domain.event.projection
import eventDemo.shared.game.projection.GameList
interface GameListRepository {
fun getList(
limit: Int = 100,
offset: Int = 0,
): List<GameList>
fun save(gameList: GameList)
fun subscribeToBus()
}
@@ -0,0 +1,6 @@
package eventDemo.contexts.game.application.ports
import eventDemo.libs.bus.Bus
import eventDemo.shared.game.projection.GameProjection
interface GameProjectionBus : Bus<GameProjection>
@@ -0,0 +1,55 @@
package eventDemo.contexts.game.application.projections
import eventDemo.contexts.game.domain.events.CardIsPlayedEvent
import eventDemo.contexts.game.domain.events.DrawFilledWithDiscardEvent
import eventDemo.contexts.game.domain.events.GameCreatedEvent
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.contexts.game.domain.events.GameStartedEvent
import eventDemo.contexts.game.domain.events.NewPlayerEvent
import eventDemo.contexts.game.domain.events.PlayerHaveDrawCardEvent
import eventDemo.contexts.game.domain.events.PlayerReadyEvent
import eventDemo.contexts.game.domain.events.PlayerWinEvent
import eventDemo.shared.game.projection.GameList
fun GameList.applyEvent(event: GameEvent): GameList =
when (event) {
is GameCreatedEvent -> {
this
}
is NewPlayerEvent -> {
copy(
players = players + event.player,
status = GameList.Status.OPENING,
)
}
is GameStartedEvent -> {
copy(
status = GameList.Status.IS_STARTED,
)
}
is PlayerWinEvent -> {
copy(
winners = winners,
status = GameList.Status.FINISH,
)
}
is CardIsPlayedEvent -> {
this
}
is PlayerHaveDrawCardEvent -> {
this
}
is PlayerReadyEvent -> {
this
}
is DrawFilledWithDiscardEvent -> {
this
}
}
@@ -0,0 +1,65 @@
package eventDemo.contexts.game.application.reaction
import eventDemo.contexts.game.application.command.handlers.GameEventManager
import eventDemo.contexts.game.application.eventStores.GameRepository
import eventDemo.contexts.game.application.logging.LoggingContextKeys
import eventDemo.contexts.game.application.logging.withLoggingContext
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.domain.game.gameState.Game
import eventDemo.contexts.game.domain.game.gameState.GameCreated
import eventDemo.contexts.game.domain.game.gameState.GameStarted
import io.github.oshai.kotlinlogging.KotlinLogging
import java.util.concurrent.ConcurrentSkipListSet
class ReactionListener(
gameRepository: GameRepository,
private val gameEventBus: GameEventBus,
) : GameEventManager(gameRepository, gameEventBus) {
private companion object Config {
val registeredListeners = ConcurrentSkipListSet<GameEventBus>()
}
private val logger = KotlinLogging.logger { }
fun subscribeToBus() {
if (registeredListeners.add(gameEventBus)) {
gameEventBus.subscribe { event ->
val game = event.getGame()
withLoggingContext(LoggingContextKeys.Game to game) {
sendStartGameEvent(game)
sendWinnerEvent(game)
}
}
} else {
"${this::class.simpleName} is already init for this bus".let {
logger.error { it }
error(it)
}
}
}
private fun sendStartGameEvent(game: Game) {
if (game is GameCreated && game.allPlayerIsReady) {
game
.startGame()
.saveEvents()
.publishEvents()
}
}
private fun sendWinnerEvent(game: Game) {
if (game is GameStarted && game.lastPlayerId != null) {
val lastPlayerWin =
game
.players
.get(game.lastPlayerId)
.hand.size == 0
if (lastPlayerWin) {
game
.playerWin(game.lastPlayerId)
.saveEvents()
.publishEvents()
}
}
}
}
@@ -0,0 +1,31 @@
package eventDemo.contexts.game.domain.events
import eventDemo.shared.game.Card
import eventDemo.shared.game.Card.Color
import eventDemo.shared.game.Player
import eventDemo.shared.ids.EventId
import eventDemo.shared.ids.GameId
import eventDemo.shared.serializers.EventIdSerializer
import kotlinx.datetime.Clock
import kotlinx.datetime.Instant
import kotlinx.serialization.Serializable
import kotlin.uuid.Uuid
/**
* An [GameEvent] to represent a played card.
*/
@Serializable
data class CardIsPlayedEvent(
override val aggregateId: GameId,
val card: Card,
override val playerId: Player.PlayerId,
val chosenColor: Color? = null,
override val version: Int,
) : GameEvent,
PlayerActionEvent {
@Serializable(with = EventIdSerializer::class)
override val eventId: EventId = EventId(Uuid.random())
override val createdAt: Instant = Clock.System.now()
val theColorCard get() = if (card is Card.CardWithColor) card.color else chosenColor
}
@@ -0,0 +1,26 @@
package eventDemo.contexts.game.domain.events
import eventDemo.shared.game.DiscardPile
import eventDemo.shared.game.DrawPile
import eventDemo.shared.ids.EventId
import eventDemo.shared.ids.GameId
import eventDemo.shared.serializers.EventIdSerializer
import kotlinx.datetime.Clock
import kotlinx.datetime.Instant
import kotlinx.serialization.Serializable
import kotlin.uuid.Uuid
/**
* When the Pile are shuffled after the draw pille was empty
*/
@Serializable
class DrawFilledWithDiscardEvent(
override val aggregateId: GameId,
val newDrawPile: DrawPile,
val newDiscardPile: DiscardPile,
override val version: Int,
) : GameEvent {
@Serializable(with = EventIdSerializer::class)
override val eventId: EventId = EventId(Uuid.random())
override val createdAt: Instant = Clock.System.now()
}
@@ -0,0 +1,22 @@
package eventDemo.contexts.game.domain.events
import eventDemo.shared.ids.EventId
import eventDemo.shared.ids.GameId
import eventDemo.shared.serializers.EventIdSerializer
import kotlinx.datetime.Clock
import kotlinx.datetime.Instant
import kotlinx.serialization.Serializable
import kotlin.uuid.Uuid
/**
* This [GameEvent] is sent when all players are ready.
*/
@Serializable
data class GameCreatedEvent(
override val aggregateId: GameId,
override val version: Int,
) : GameEvent {
@Serializable(with = EventIdSerializer::class)
override val eventId: EventId = EventId(Uuid.random())
override val createdAt: Instant = Clock.System.now()
}
@@ -0,0 +1,21 @@
package eventDemo.contexts.game.domain.events
import eventDemo.libs.eventSource.Event
import eventDemo.shared.ids.EventId
import eventDemo.shared.ids.GameId
import eventDemo.shared.serializers.EventIdSerializer
import eventDemo.shared.serializers.GameIdSerializer
import kotlinx.serialization.Serializable
/**
* An [Event] of a Game.
*/
@Serializable
sealed interface GameEvent : Event<GameId> {
@Serializable(with = EventIdSerializer::class)
override val eventId: EventId
@Serializable(with = GameIdSerializer::class)
override val aggregateId: GameId
override val version: Int
}
@@ -0,0 +1,32 @@
package eventDemo.contexts.game.domain.events
import eventDemo.shared.game.DiscardPile
import eventDemo.shared.game.DrawPile
import eventDemo.shared.game.Player
import eventDemo.shared.game.PlayerHand
import eventDemo.shared.ids.EventId
import eventDemo.shared.ids.GameId
import eventDemo.shared.serializers.EventIdSerializer
import eventDemo.shared.serializers.PlayerIdSerializer
import kotlinx.datetime.Clock
import kotlinx.datetime.Instant
import kotlinx.serialization.Serializable
import kotlin.uuid.Uuid
/**
* This [GameEvent] is sent when all players are ready.
*/
@Serializable
data class GameStartedEvent(
override val aggregateId: GameId,
@Serializable(with = PlayerIdSerializer::class)
val firstPlayer: Player.PlayerId,
val playersHans: Map<Player.PlayerId, PlayerHand>,
val drawPile: DrawPile,
val discardPile: DiscardPile,
override val version: Int,
) : GameEvent {
@Serializable(with = EventIdSerializer::class)
override val eventId: EventId = EventId(Uuid.random())
override val createdAt: Instant = Clock.System.now()
}
@@ -0,0 +1,27 @@
package eventDemo.contexts.game.domain.events
import eventDemo.shared.game.Player
import eventDemo.shared.ids.EventId
import eventDemo.shared.ids.GameId
import eventDemo.shared.serializers.EventIdSerializer
import kotlinx.datetime.Clock
import kotlinx.datetime.Instant
import kotlinx.serialization.Serializable
import kotlin.uuid.Uuid
/**
* An [GameEvent] to represent a new player joining the game.
*/
@Serializable
data class NewPlayerEvent(
override val aggregateId: GameId,
val player: Player,
override val version: Int,
) : GameEvent,
PlayerActionEvent {
override val playerId: Player.PlayerId get() = player.id
@Serializable(with = EventIdSerializer::class)
override val eventId: EventId = EventId(Uuid.random())
override val createdAt: Instant = Clock.System.now()
}
@@ -0,0 +1,9 @@
package eventDemo.contexts.game.domain.events
import eventDemo.shared.game.Player
import kotlinx.serialization.Serializable
@Serializable
sealed interface PlayerActionEvent : GameEvent {
val playerId: Player.PlayerId
}
@@ -0,0 +1,29 @@
package eventDemo.contexts.game.domain.events
import eventDemo.shared.game.Card
import eventDemo.shared.game.Player
import eventDemo.shared.ids.EventId
import eventDemo.shared.ids.GameId
import eventDemo.shared.serializers.EventIdSerializer
import eventDemo.shared.serializers.PlayerIdSerializer
import kotlinx.datetime.Clock
import kotlinx.datetime.Instant
import kotlinx.serialization.Serializable
import kotlin.uuid.Uuid
/**
* This [GameEvent] is sent when a player can play.
*/
@Serializable
data class PlayerHaveDrawCardEvent(
override val aggregateId: GameId,
@Serializable(with = PlayerIdSerializer::class)
override val playerId: Player.PlayerId,
val takenCards: Set<Card>,
override val version: Int,
) : GameEvent,
PlayerActionEvent {
@Serializable(with = EventIdSerializer::class)
override val eventId: EventId = EventId(Uuid.random())
override val createdAt: Instant = Clock.System.now()
}
@@ -0,0 +1,27 @@
package eventDemo.contexts.game.domain.events
import eventDemo.shared.game.Player
import eventDemo.shared.ids.EventId
import eventDemo.shared.ids.GameId
import eventDemo.shared.serializers.EventIdSerializer
import eventDemo.shared.serializers.PlayerIdSerializer
import kotlinx.datetime.Clock
import kotlinx.datetime.Instant
import kotlinx.serialization.Serializable
import kotlin.uuid.Uuid
/**
* This [GameEvent] is sent when a player is ready.
*/
@Serializable
data class PlayerReadyEvent(
override val aggregateId: GameId,
@Serializable(with = PlayerIdSerializer::class)
override val playerId: Player.PlayerId,
override val version: Int,
) : GameEvent,
PlayerActionEvent {
@Serializable(with = EventIdSerializer::class)
override val eventId: EventId = EventId(Uuid.random())
override val createdAt: Instant = Clock.System.now()
}
@@ -0,0 +1,27 @@
package eventDemo.contexts.game.domain.events
import eventDemo.shared.game.Player
import eventDemo.shared.ids.EventId
import eventDemo.shared.ids.GameId
import eventDemo.shared.serializers.EventIdSerializer
import eventDemo.shared.serializers.PlayerIdSerializer
import kotlinx.datetime.Clock
import kotlinx.datetime.Instant
import kotlinx.serialization.Serializable
import kotlin.uuid.Uuid
/**
* This [GameEvent] is sent when a player is ready.
*/
@Serializable
data class PlayerWinEvent(
override val aggregateId: GameId,
@Serializable(with = PlayerIdSerializer::class)
override val playerId: Player.PlayerId,
override val version: Int,
) : GameEvent,
PlayerActionEvent {
@Serializable(with = EventIdSerializer::class)
override val eventId: EventId = EventId(Uuid.random())
override val createdAt: Instant = Clock.System.now()
}
@@ -0,0 +1,68 @@
package eventDemo.contexts.game.domain.game.errors
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.contexts.game.domain.game.gameState.Deck
import eventDemo.contexts.game.domain.game.gameState.Game
import eventDemo.shared.game.Card
import eventDemo.shared.game.Player
import eventDemo.shared.game.PlayerList
abstract class GameException(
message: String,
) : Exception(message)
abstract class IllegalActionException(
message: String,
) : GameException(message)
class ItsNotTheTurnException(
val playerId: Player.PlayerId,
) : IllegalActionException("It is not the turn of the player") {
constructor(playerId: Player.PlayerId, event: GameEvent) : this(playerId)
}
class TheCardIsAColorCardException(
val playerId: Player.PlayerId,
) : IllegalActionException("The card is a color card")
class TheCardHasNoColorException(
val playerId: Player.PlayerId,
) : IllegalActionException("The card has no color, you must chose a color")
class ThePlayerHasRemainingCardsException(
val player: Player,
) : IllegalActionException("The player has remaining cards")
class ThePlayerHasAlreadyWinException(
val playerId: Player.PlayerId,
) : IllegalActionException("The player has already win")
class ThePlayerIsNotInTheGameException(
val playerId: Player.PlayerId,
) : IllegalActionException("The player is not in the game")
class ThePlayerMustPlayACardException(
val playerId: Player.PlayerId,
val playableCards: Set<Card>,
) : IllegalActionException("The player must be play a card")
class NeedMorePlayersToStartGameException(
val players: PlayerList,
) : IllegalActionException("You cannot start a game with less than 2 players!")
class AllPlayerNotReadyException(
val players: PlayerList,
) : IllegalActionException("All players not ready!")
class DeckMissingCardsException(
val players: PlayerList,
deck: Deck,
) : IllegalActionException("The deck missing cards")
class InconsistentGameException(
val game: Game,
) : GameException("Inconsistent game state")
class InconsistentEventVersionException(
val game: Set<GameEvent>,
) : GameException("Inconsistent event version")
@@ -0,0 +1,81 @@
package eventDemo.contexts.game.domain.game.gameState
import eventDemo.contexts.game.domain.events.CardIsPlayedEvent
import eventDemo.contexts.game.domain.events.DrawFilledWithDiscardEvent
import eventDemo.contexts.game.domain.events.GameCreatedEvent
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.contexts.game.domain.events.GameStartedEvent
import eventDemo.contexts.game.domain.events.NewPlayerEvent
import eventDemo.contexts.game.domain.events.PlayerHaveDrawCardEvent
import eventDemo.contexts.game.domain.events.PlayerReadyEvent
import eventDemo.contexts.game.domain.events.PlayerWinEvent
import eventDemo.contexts.game.domain.game.errors.GameException
import eventDemo.contexts.game.domain.game.errors.InconsistentEventVersionException
import eventDemo.shared.game.PlayerList
import eventDemo.shared.ids.GameId
sealed interface Game {
val aggregateId: GameId
val players: PlayerList
/**
* On each modification, an event is put in their
*/
val recordedEvents: Set<GameEvent>
val version: Int
enum class Direction {
CLOCKWISE,
COUNTER_CLOCKWISE,
;
fun revert(): Direction =
if (this === CLOCKWISE) {
COUNTER_CLOCKWISE
} else {
CLOCKWISE
}
}
companion object {
fun loadFromHistory(events: Set<GameEvent>): Game =
events
.fold(GameInit(events.first().aggregateId)) { game: Game, event ->
game.run {
when (event) {
is GameCreatedEvent if this is GameInit -> applyEvent(event)
is GameCreatedEvent -> error("Game is already created")
is NewPlayerEvent if this is GameCreated -> applyEvent(event)
is NewPlayerEvent -> error("Game is already stared")
is PlayerReadyEvent if this is GameCreated -> applyEvent(event)
is PlayerReadyEvent -> error("Game is already stared")
is GameStartedEvent if this is GameCreated -> applyEvent(event)
is GameStartedEvent -> error("Game is already started")
is CardIsPlayedEvent if this is GameStarted -> applyEvent(event)
is CardIsPlayedEvent -> error("Game is end")
is PlayerHaveDrawCardEvent if this is GameStarted -> applyEvent(event)
is PlayerHaveDrawCardEvent -> error("Game is end")
is PlayerWinEvent if this is GameStarted -> applyEvent(event)
is PlayerWinEvent -> error("Game is end")
is DrawFilledWithDiscardEvent if this is GameStarted -> applyEvent(event)
is DrawFilledWithDiscardEvent -> error("Game is end")
}
}
}.let {
when (it) {
is GameInit -> it
is GameCreated -> it.copy(recordedEvents = emptySet())
is GameEnded -> it.copy(recordedEvents = emptySet())
is GameStarted -> it.copy(recordedEvents = emptySet())
}
}
}
}
internal fun <T : GameEvent> T.checkState(
block: (T) -> Boolean,
exception: (T) -> GameException,
): T {
if (!block(this)) throw exception(this)
return this
}
@@ -0,0 +1,179 @@
package eventDemo.contexts.game.domain.game.gameState
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.contexts.game.domain.events.GameStartedEvent
import eventDemo.contexts.game.domain.events.NewPlayerEvent
import eventDemo.contexts.game.domain.events.PlayerReadyEvent
import eventDemo.contexts.game.domain.game.errors.AllPlayerNotReadyException
import eventDemo.contexts.game.domain.game.errors.DeckMissingCardsException
import eventDemo.contexts.game.domain.game.errors.NeedMorePlayersToStartGameException
import eventDemo.contexts.game.domain.game.errors.ThePlayerIsNotInTheGameException
import eventDemo.shared.game.Card
import eventDemo.shared.game.DiscardPile
import eventDemo.shared.game.DrawPile
import eventDemo.shared.game.Player
import eventDemo.shared.game.PlayerHand
import eventDemo.shared.game.PlayerList
import eventDemo.shared.ids.GameId
import eventDemo.shared.ids.UserId
data class GameCreated(
override val aggregateId: GameId,
override val players: PlayerList = PlayerList(),
val playersStatus: Map<Player.PlayerId, PlayerStatus> = emptyMap(),
override val recordedEvents: Set<GameEvent>,
override val version: Int,
) : Game {
val allPlayerIsReady: Boolean
get() {
return playersStatus.isNotEmpty() && playersStatus.values.all { it == PlayerStatus.Ready }
}
fun startGame(deck: Deck = newDeck().shuffleDeck()): GameStarted {
val (drawPile, discardPile, playersHands) =
initPiles(deck)
.let { (drawPile, discardPile) ->
createHandsFromDrawPile(drawPile)
.let { (drawPile, playersHands) ->
Triple(drawPile, discardPile, playersHands)
}
}
return GameStartedEvent(
aggregateId = aggregateId,
firstPlayer = players.randomPlayer().id,
version = version + 1,
drawPile = drawPile,
discardPile = discardPile,
playersHans = playersHands,
).checkState(
{ players.size > 1 },
{ NeedMorePlayersToStartGameException(players) },
).checkState(
{ allPlayerIsReady },
{ AllPlayerNotReadyException(players) },
).checkState(
{ deck.size == 108 },
{ DeckMissingCardsException(players, deck) },
).also { if (it.drawPile.size + it.discardPile.size + playersHands.values.sumOf { it.size } != 108) error("missing cards!") }
.run(::applyEvent)
}
private fun initPiles(deck: Set<Card>): Pair<DrawPile, DiscardPile> =
DrawPile(deck)
.generateValidDrawPile()
.take(1)
.let { (draw, cards) ->
draw to DiscardPile(cards)
}
private fun createHandsFromDrawPile(drawPile: DrawPile) =
players
.map { it.id }
.fold(Pair(drawPile, emptyMap<Player.PlayerId, PlayerHand>())) { (drawAcc, handsAcc), playerId ->
drawAcc
.take(7)
.let { (draw, hand) ->
Pair(
draw,
handsAcc + (playerId to PlayerHand(hand)),
)
}
}
fun userJoinTheGame(
userId: UserId,
name: String,
): GameCreated {
if (players.map { it.userId }.contains(userId)) {
throw IllegalStateException("User $userId already in party")
}
val player = Player(name, userId)
return applyEvent(
NewPlayerEvent(
aggregateId = aggregateId,
player = player,
version = version + 1,
),
)
}
fun setReadyPlayer(playerId: Player.PlayerId): GameCreated {
if (!players.map { it.id }.contains(playerId)) {
throw ThePlayerIsNotInTheGameException(playerId)
}
return PlayerReadyEvent(aggregateId, playerId, version + 1)
.run(::applyEvent)
}
internal fun applyEvent(event: NewPlayerEvent): GameCreated =
copy(
players = players + (event.player),
playersStatus = playersStatus + (event.player.id to PlayerStatus.Waiting),
recordedEvents = recordedEvents + event,
version = version + 1,
)
internal fun applyEvent(event: PlayerReadyEvent): GameCreated =
copy(
playersStatus = playersStatus + (event.playerId to PlayerStatus.Ready),
recordedEvents = recordedEvents + event,
version = version + 1,
)
internal fun applyEvent(event: GameStartedEvent): GameStarted =
GameStarted(
aggregateId = event.aggregateId,
players =
players
.map {
it.copy(hand = event.playersHans[it.id] ?: error("Player ${it.id} not found"))
}.let { PlayerList(it.toSet()) },
lastPlayerId = null,
nextPlayerId = event.firstPlayer,
drawPile = event.drawPile,
discardPile = event.discardPile,
version = event.version + 1,
recordedEvents = recordedEvents + event,
currentColor = event.discardPile.topCardColor ?: error("The discard pile was not initialized!"),
)
enum class PlayerStatus {
Ready,
Waiting,
}
}
typealias Deck = Set<Card>
fun newDeck(): Deck =
listOf(Card.Color.Red, Card.Color.Blue, Card.Color.Yellow, Card.Color.Green)
.flatMap { color ->
((0..9) + (1..9)).map { Card.NumericCard(it, color) } +
(1..2).map { Card.Plus2Card(color) } +
(1..2).map { Card.ReverseCard(color) } +
(1..2).map { Card.PassCard(color) }
}.let {
it + (1..4).map { Card.Plus4Card() }
}.let {
it + (1..4).map { Card.ChangeColorCard() }
}.toSet()
fun Set<Card>.shuffleDeck(): Set<Card> {
if (isDisabled) return this
return shuffled().toSet()
}
private fun PlayerList.randomPlayer(): Player {
if (isDisabled) return first()
return random()
}
private var isDisabled = false
fun disableRandomForTest() {
isDisabled = true
}
@@ -0,0 +1,20 @@
package eventDemo.contexts.game.domain.game.gameState
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.shared.game.Player
import eventDemo.shared.game.PlayerList
import eventDemo.shared.ids.GameId
data class GameEnded(
override val aggregateId: GameId,
override val players: PlayerList,
val playerWins: Set<Player.PlayerId> = emptySet(),
override val version: Int,
override val recordedEvents: Set<GameEvent>,
) : Game {
init {
if (!players.map { it.id }.containsAll(playerWins)) {
throw IllegalArgumentException("Player ${players.map { it.id }} were not in players")
}
}
}
@@ -0,0 +1,28 @@
package eventDemo.contexts.game.domain.game.gameState
import eventDemo.contexts.game.domain.events.GameCreatedEvent
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.shared.game.PlayerList
import eventDemo.shared.ids.GameId
data class GameInit(
override val aggregateId: GameId,
) : Game {
override val players: PlayerList = PlayerList()
override var recordedEvents: Set<GameEvent> = emptySet()
// 0 = no events; not persisted; not really exist.
override var version: Int = 0
companion object {
fun createNewGame(gameId: GameId = GameId()): GameCreated {
val event = GameCreatedEvent(gameId, 1)
return GameInit(event.aggregateId).run {
event.run(::applyEvent)
}
}
}
internal fun applyEvent(event: GameCreatedEvent): GameCreated =
GameCreated(aggregateId, recordedEvents = setOf(event), version = event.version)
}
@@ -0,0 +1,264 @@
package eventDemo.contexts.game.domain.game.gameState
import eventDemo.contexts.game.domain.events.CardIsPlayedEvent
import eventDemo.contexts.game.domain.events.DrawFilledWithDiscardEvent
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.contexts.game.domain.events.PlayerActionEvent
import eventDemo.contexts.game.domain.events.PlayerHaveDrawCardEvent
import eventDemo.contexts.game.domain.events.PlayerWinEvent
import eventDemo.contexts.game.domain.game.errors.InconsistentGameException
import eventDemo.contexts.game.domain.game.errors.ItsNotTheTurnException
import eventDemo.contexts.game.domain.game.errors.TheCardHasNoColorException
import eventDemo.contexts.game.domain.game.errors.TheCardIsAColorCardException
import eventDemo.contexts.game.domain.game.errors.ThePlayerHasAlreadyWinException
import eventDemo.contexts.game.domain.game.errors.ThePlayerHasRemainingCardsException
import eventDemo.contexts.game.domain.game.errors.ThePlayerIsNotInTheGameException
import eventDemo.contexts.game.domain.game.errors.ThePlayerMustPlayACardException
import eventDemo.contexts.game.domain.game.gameState.Game.Direction
import eventDemo.shared.game.Card
import eventDemo.shared.game.Card.Color
import eventDemo.shared.game.DiscardPile
import eventDemo.shared.game.DrawPile
import eventDemo.shared.game.Player
import eventDemo.shared.game.PlayerList
import eventDemo.shared.ids.GameId
fun PlayerList.nextPlayerTurn(
lastPlayerId: Player.PlayerId,
direction: Direction,
): Player.PlayerId {
val lastPlayer = get(lastPlayerId)
val playersLastTurn = filter { it.hand.cards.isNotEmpty() || it == lastPlayer }
return playersLastTurn
.indexOf(lastPlayer)
.let { lastPlayerIndex ->
if (direction == Direction.CLOCKWISE) {
if (lastPlayerIndex == playersLastTurn.size - 1) {
0
} else {
lastPlayerIndex + 1
}
} else {
if (lastPlayerIndex == 0) {
playersLastTurn.size - 1
} else {
lastPlayerIndex - 1
}
}
}.let { nextPlayerIndex -> elementAt(nextPlayerIndex).id }
}
data class GameStarted(
override val aggregateId: GameId,
override val players: PlayerList,
val drawPile: DrawPile,
val discardPile: DiscardPile,
val lastPlayerId: Player.PlayerId?,
val nextPlayerId: Player.PlayerId,
val currentColor: Color,
val playedTurnHistory: List<History> = emptyList(),
val direction: Direction = Direction.CLOCKWISE,
val playerWins: Set<Player.PlayerId> = emptySet(),
override val version: Int,
override val recordedEvents: Set<GameEvent>,
) : Game {
val playersInGame by lazy { players.filter { it.hand.cards.isNotEmpty() } }
val lastPlayedCard: Card? by lazy { discardPile.topCard }
data class History(
val playerId: Player.PlayerId,
val event: GameEvent,
val direction: Direction,
)
val lastPlayer: Player? by lazy { lastPlayerId?.let { players.get(it) } }
val nextPlayer: Player by lazy { players.get(nextPlayerId) }
fun canBePlayThisCard(card: Card): Boolean {
val cardOnBoard = discardPile.topCard ?: return false
return when (cardOnBoard) {
is Card.NumericCard -> {
when (card) {
is Card.CardWith4Color -> true
is Card.NumericCard -> card.number == cardOnBoard.number || card.color == cardOnBoard.color
is Card.CardWithColor -> card.color == cardOnBoard.color
}
}
is Card.ReverseCard -> {
when (card) {
is Card.ReverseCard -> true
is Card.CardWith4Color -> true
is Card.CardWithColor -> card.color == cardOnBoard.color
}
}
is Card.PassCard -> {
when (card) {
is Card.CardWith4Color -> true
is Card.CardWithColor -> card.color == cardOnBoard.color
}
}
is Card.ChangeColorCard -> {
when (card) {
is Card.CardWith4Color -> true
is Card.CardWithColor -> card.color == currentColor
}
}
is Card.Plus2Card -> {
when (card) {
is Card.Plus2Card -> true
::isPlayedLastTurn -> false
is Card.CardWith4Color -> true
is Card.CardWithColor -> card.color == currentColor
}
}
is Card.Plus4Card -> {
when (card) {
is Card.Plus4Card -> true
::isPlayedLastTurn -> false
is Card.CardWith4Color -> true
is Card.CardWithColor -> card.color == currentColor
}
}
}
}
fun playableCards(playerId: Player.PlayerId): Set<Card> =
players
.get(playerId)
.hand
.cards
.filter(::canBePlayThisCard)
.toSet()
fun playTheCard(
playerId: Player.PlayerId,
card: Card,
chosenColor: Color? = null,
): GameStarted =
CardIsPlayedEvent(aggregateId, card, playerId, chosenColor, version + 1)
.checkPlayerTurn()
.checkState({
(card is Card.CardWithColor && chosenColor == null) || card is Card.CardWith4Color
}, { TheCardIsAColorCardException(playerId) })
.checkState({
(card is Card.CardWith4Color && chosenColor != null) || card is Card.CardWithColor
}, { TheCardHasNoColorException(playerId) })
.run(::applyEvent)
internal fun applyEvent(event: CardIsPlayedEvent): GameStarted =
run {
val nextDirectionAfterPlay =
when (event.card) {
is Card.ReverseCard -> direction.revert()
else -> direction
}
val color =
when (event.card) {
is Card.CardWithColor -> event.card.color
is Card.CardWith4Color -> event.chosenColor!!
}
copy(
players = players.withDropCardOnPlayerHand(event.playerId, event.card),
discardPile = discardPile.withNewCard(card = event.card),
currentColor = color,
lastPlayerId = event.playerId,
nextPlayerId = players.nextPlayerTurn(event.playerId, nextDirectionAfterPlay),
playedTurnHistory = playedTurnHistory - History(event.playerId, event, direction),
direction = nextDirectionAfterPlay,
version = event.version,
recordedEvents = recordedEvents + event,
)
}
fun playerTakeCartFromDrawPile(
playerId: Player.PlayerId,
number: Int,
): GameStarted {
val takenCards = drawPile.take(number).second
return PlayerHaveDrawCardEvent(aggregateId, playerId, takenCards, version + 1)
.checkPlayerTurn()
.checkState({
playableCards(playerId).isEmpty()
}, {
ThePlayerMustPlayACardException(playerId, playableCards(playerId))
})
.run(::applyEvent)
.run {
val missingCardsCount = number - takenCards.size
if (missingCardsCount > 0) {
fillDrawWithDiscard()
.playerTakeCartFromDrawPile(playerId, missingCardsCount)
} else {
this
}
}
}
internal fun applyEvent(event: PlayerHaveDrawCardEvent): GameStarted =
copy(
players = players.withNewCardOnPlayerHand(event.playerId, event.takenCards),
drawPile = drawPile.take(event.takenCards.size).first,
lastPlayerId = event.playerId,
nextPlayerId = players.nextPlayerTurn(event.playerId, direction),
version = event.version,
recordedEvents = recordedEvents + event,
)
/**
* Filling the draw pile with the discard pile while excluding the top card
*/
private fun fillDrawWithDiscard(): GameStarted =
run {
val topCard = discardPile.topCard ?: throw InconsistentGameException(this)
DrawPile(discardPile.cards - topCard).shuffled() to DiscardPile(setOf(topCard))
}.let { (newDrawPile, newDiscardPile) ->
DrawFilledWithDiscardEvent(aggregateId, newDrawPile, newDiscardPile, version + 1)
.run(::applyEvent)
}
internal fun applyEvent(event: DrawFilledWithDiscardEvent): GameStarted =
copy(
drawPile = event.newDrawPile,
discardPile = event.newDiscardPile,
version = event.version,
recordedEvents = recordedEvents + event,
)
fun playerWin(playerId: Player.PlayerId): GameStarted =
PlayerWinEvent(aggregateId, playerId, version + 1)
.checkState({
players
.get(playerId)
.hand.cards
.isEmpty()
}, { ThePlayerHasRemainingCardsException(players.get(playerId)) })
.checkState({ playerWins.contains(playerId) }, { ThePlayerHasAlreadyWinException(playerId) })
.checkState({ !players.map { it.id }.contains(playerId) }, { ThePlayerIsNotInTheGameException(playerId) })
.run(::applyEvent)
internal fun applyEvent(event: PlayerWinEvent): GameStarted =
copy(
playerWins = playerWins + event.playerId,
version = event.version,
recordedEvents = recordedEvents + event,
)
private fun <T : PlayerActionEvent> T.checkPlayerTurn(): T =
checkState(
{ nextPlayer.id == playerId },
{ ItsNotTheTurnException(playerId, this) },
)
private fun isPlayedLastTurn(card: Card): Boolean =
(playedTurnHistory.last().event as? CardIsPlayedEvent)?.card == card
}
@@ -0,0 +1,17 @@
package eventDemo.contexts.game.infrastructure.configuration.injections.application
import eventDemo.contexts.game.application.eventStores.GameEventStoreRepository
import eventDemo.contexts.game.application.eventStores.GameRepository
import eventDemo.contexts.game.application.notification.EventToNotificationSubscriber
import eventDemo.contexts.game.application.reaction.ReactionListener
import org.koin.core.module.Module
import org.koin.core.module.dsl.singleOf
import org.koin.dsl.bind
fun Module.configureGameDIApplication() {
configureDICommandHandlers()
singleOf(::ReactionListener)
singleOf(::EventToNotificationSubscriber)
singleOf(::GameEventStoreRepository) bind GameRepository::class
}
@@ -0,0 +1,18 @@
package eventDemo.contexts.game.infrastructure.configuration.injections.application
import eventDemo.contexts.game.application.command.handlers.JoinTheGameHandler
import eventDemo.contexts.game.application.command.handlers.PlayCardHandler
import eventDemo.contexts.game.application.command.handlers.ReadyToPlayHandler
import eventDemo.contexts.game.application.command.handlers.TakeCartFromDrawPileHandler
import org.koin.core.module.Module
import org.koin.core.module.dsl.singleOf
/**
* Configure all actions
*/
fun Module.configureDICommandHandlers() {
singleOf(::PlayCardHandler)
singleOf(::ReadyToPlayHandler)
singleOf(::JoinTheGameHandler)
singleOf(::TakeCartFromDrawPileHandler)
}
@@ -0,0 +1,26 @@
package eventDemo.contexts.game.infrastructure.configuration.injections.infrastructure
import eventDemo.contexts.game.application.channels.GameChannelsSubscriber
import eventDemo.contexts.game.application.command.handlers.GameCommandHandlerDispatcher
import eventDemo.contexts.game.application.notification.CommandSubscriber
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.application.ports.GameEventStore
import eventDemo.contexts.game.application.ports.GameProjectionBus
import eventDemo.contexts.game.infrastructure.persistence.eventBus.GameEventBusInRabbinMQ
import eventDemo.contexts.game.infrastructure.persistence.eventStore.GameEventStoreInPostgresql
import eventDemo.contexts.game.infrastructure.persistence.projections.GameListRepositoryInMemory
import eventDemo.contexts.game.infrastructure.persistence.projections.bus.GameProjectionBusInRabbitMQ
import eventDemo.domain.event.projection.GameListRepository
import org.koin.core.module.Module
import org.koin.core.module.dsl.singleOf
import org.koin.dsl.bind
fun Module.configureGameDIInfrastructure() {
singleOf(::GameEventStoreInPostgresql) bind GameEventStore::class
singleOf(::GameEventBusInRabbinMQ) bind GameEventBus::class
singleOf(::GameProjectionBusInRabbitMQ) bind GameProjectionBus::class
singleOf(::CommandSubscriber)
singleOf(::GameChannelsSubscriber)
singleOf(::GameCommandHandlerDispatcher)
singleOf(::GameListRepositoryInMemory) bind GameListRepository::class
}
@@ -0,0 +1,40 @@
package eventDemo.contexts.game.infrastructure.configuration.ktor
import eventDemo.shared.http.HttpErrorBadRequest
import io.ktor.http.HttpHeaders
import io.ktor.http.HttpMethod
import io.ktor.http.HttpStatusCode
import io.ktor.server.application.Application
import io.ktor.server.application.install
import io.ktor.server.plugins.autohead.AutoHeadResponse
import io.ktor.server.plugins.cors.routing.CORS
import io.ktor.server.plugins.statuspages.StatusPages
import io.ktor.server.resources.Resources
import io.ktor.server.response.respondText
fun Application.configureHttpRouting() {
install(CORS) {
allowMethod(HttpMethod.Options)
allowMethod(HttpMethod.Put)
allowMethod(HttpMethod.Post)
allowMethod(HttpMethod.Delete)
allowMethod(HttpMethod.Patch)
allowHeader(HttpHeaders.Authorization)
allowHeader("MyCustomHeader")
anyHost() // @TODO: Don't do this in production if possible. Try to limit it.
}
install(AutoHeadResponse)
install(Resources)
install(StatusPages) {
exception<BadRequestException> { call, cause ->
call.respondText(text = "400: $cause", status = HttpStatusCode.BadRequest)
}
exception<Throwable> { call, cause ->
call.respondText(text = "500: $cause", status = HttpStatusCode.InternalServerError)
}
}
}
class BadRequestException(
val httpError: HttpErrorBadRequest,
) : Exception()
@@ -0,0 +1,38 @@
package eventDemo.contexts.game.infrastructure.configuration.ktor
import eventDemo.shared.game.Player
import eventDemo.shared.ids.CommandId
import eventDemo.shared.ids.EventId
import eventDemo.shared.ids.GameId
import eventDemo.shared.serializers.CommandIdSerializer
import eventDemo.shared.serializers.EventIdSerializer
import eventDemo.shared.serializers.GameIdSerializer
import eventDemo.shared.serializers.PlayerIdSerializer
import eventDemo.shared.serializers.UUIDSerializer
import io.ktor.serialization.kotlinx.json.json
import io.ktor.server.application.Application
import io.ktor.server.application.install
import io.ktor.server.plugins.contentnegotiation.ContentNegotiation
import kotlinx.serialization.json.Json
import kotlinx.serialization.modules.SerializersModule
import kotlin.uuid.Uuid
fun Application.configureSerialization() {
install(ContentNegotiation) {
json(
defaultJsonSerializer(),
)
}
}
fun defaultJsonSerializer(): Json =
Json {
serializersModule =
SerializersModule {
contextual(Uuid::class) { UUIDSerializer }
contextual(GameId::class) { GameIdSerializer }
contextual(EventId::class) { EventIdSerializer }
contextual(CommandId::class) { CommandIdSerializer }
contextual(Player.PlayerId::class) { PlayerIdSerializer }
}
}
@@ -0,0 +1,21 @@
package eventDemo.contexts.game.infrastructure.configuration.ktor
import eventDemo.contexts.game.infrastructure.configuration.listener.configureProjectionListener
import eventDemo.contexts.game.infrastructure.configuration.listener.configureReactionListener
import io.ktor.server.application.Application
import org.koin.ktor.plugin.koin
fun Application.configureUno() {
configureSerialization()
configureWebSockets()
declareWebSocketsRoute()
configureHttpRouting()
declareHttpGameRoute()
koin().run {
configureProjectionListener()
configureReactionListener()
}
}
@@ -0,0 +1,17 @@
package eventDemo.contexts.game.infrastructure.configuration.ktor
import io.ktor.server.application.Application
import io.ktor.server.application.install
import io.ktor.server.websocket.WebSockets
import io.ktor.server.websocket.pingPeriod
import io.ktor.server.websocket.timeout
import kotlin.time.Duration.Companion.seconds
fun Application.configureWebSockets() {
install(WebSockets) {
pingPeriod = 15.seconds
timeout = 15.seconds
maxFrameSize = Long.MAX_VALUE
masking = false
}
}
@@ -0,0 +1,14 @@
package eventDemo.contexts.game.infrastructure.configuration.ktor
import eventDemo.contexts.game.infrastructure.rest.gamesListRoute
import eventDemo.contexts.game.infrastructure.rest.getFullNotificationsRoute
import io.ktor.server.application.Application
import io.ktor.server.routing.routing
import org.koin.ktor.ext.get as getDi
fun Application.declareHttpGameRoute() {
routing {
gamesListRoute(getDi())
getFullNotificationsRoute(getDi(), getDi())
}
}
@@ -0,0 +1,16 @@
package eventDemo.contexts.game.infrastructure.configuration.ktor
import eventDemo.contexts.game.infrastructure.websocket.gameWebSocket
import io.ktor.server.application.Application
import io.ktor.server.routing.routing
import kotlinx.coroutines.DelicateCoroutinesApi
import org.koin.ktor.ext.get as getDi
@OptIn(DelicateCoroutinesApi::class)
fun Application.declareWebSocketsRoute() {
routing {
gameWebSocket(
getDi(),
)
}
}
@@ -0,0 +1,9 @@
package eventDemo.contexts.game.infrastructure.configuration.listener
import eventDemo.domain.event.projection.GameListRepository
import org.koin.core.Koin
fun Koin.configureProjectionListener() {
get<GameListRepository>()
.subscribeToBus()
}
@@ -0,0 +1,9 @@
package eventDemo.contexts.game.infrastructure.configuration.listener
import eventDemo.contexts.game.application.reaction.ReactionListener
import org.koin.core.Koin
fun Koin.configureReactionListener() {
get<ReactionListener>()
.subscribeToBus()
}
@@ -0,0 +1,17 @@
package eventDemo.contexts.game.infrastructure.persistence.eventBus
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.libs.bus.Bus
import eventDemo.libs.bus.BusInMemory
import java.util.UUID
class GameEventBusInMemory :
GameEventBus,
Bus<GameEvent> by BusInMemory(GameEventBusInMemory::class),
Comparable<GameEventBusInMemory> {
private val instanceId: UUID = UUID.randomUUID()
override fun compareTo(other: GameEventBusInMemory): Int =
compareValues(instanceId, other.instanceId)
}
@@ -0,0 +1,25 @@
package eventDemo.contexts.game.infrastructure.persistence.eventBus
import com.rabbitmq.client.ConnectionFactory
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.libs.bus.Bus
import eventDemo.libs.bus.BusInRabbitMQ
import kotlinx.serialization.json.Json
import java.util.UUID
class GameEventBusInRabbinMQ(
private val connectionFactory: ConnectionFactory,
) : GameEventBus,
Bus<GameEvent> by BusInRabbitMQ(
connectionFactory,
"GameEvent",
{ Json.encodeToString(it) },
{ Json.decodeFromString<GameEvent>(it) },
),
Comparable<GameEventBusInRabbinMQ> {
private val instanceId: UUID = UUID.randomUUID()
override fun compareTo(other: GameEventBusInRabbinMQ): Int =
compareValues(instanceId, other.instanceId)
}
@@ -0,0 +1,14 @@
package eventDemo.contexts.game.infrastructure.persistence.eventStore
import eventDemo.contexts.game.application.ports.GameEventStore
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.libs.eventSource.eventStore.EventStore
import eventDemo.libs.eventSource.eventStore.EventStoreInMemory
import eventDemo.shared.ids.GameId
/**
* A stream to publish and read the played card event.
*/
class GameEventStoreInMemory :
GameEventStore,
EventStore<GameEvent, GameId> by EventStoreInMemory()
@@ -0,0 +1,22 @@
package eventDemo.contexts.game.infrastructure.persistence.eventStore
import eventDemo.contexts.game.application.ports.GameEventStore
import eventDemo.contexts.game.domain.events.GameEvent
import eventDemo.libs.eventSource.eventStore.EventStore
import eventDemo.libs.eventSource.eventStore.EventStoreInPostgresql
import eventDemo.shared.ids.GameId
import kotlinx.serialization.json.Json
import javax.sql.DataSource
/**
* A stream to publish and read the played card event.
*/
class GameEventStoreInPostgresql(
dataSource: DataSource,
) : GameEventStore,
EventStore<GameEvent, GameId> by EventStoreInPostgresql(
dataSource,
{ Json.encodeToString(it) },
{ Json.decodeFromString(it) },
"game.game_event_stream",
)
@@ -0,0 +1,49 @@
package eventDemo.contexts.game.infrastructure.persistence.projections
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.application.ports.GameEventStore
import eventDemo.contexts.game.application.ports.GameProjectionBus
import eventDemo.contexts.game.application.projections.applyEvent
import eventDemo.domain.event.projection.GameListRepository
import eventDemo.shared.game.projection.GameList
import eventDemo.shared.ids.GameId
import io.github.oshai.kotlinlogging.withLoggingContext
/**
* Manages [projections][GameList], their building and publication in the [bus][GameProjectionBus].
*/
class GameListRepositoryInMemory(
val gameEventStore: GameEventStore,
val projectionBus: GameProjectionBus,
val eventBus: GameEventBus,
) : GameListRepository {
val projections: MutableMap<GameId, GameList> = mutableMapOf()
override fun getList(
limit: Int,
offset: Int,
): List<GameList> =
projections
.values
.drop(offset)
.take(limit)
override fun save(gameList: GameList) {
projections[gameList.aggregateId] = gameList
}
override fun subscribeToBus() {
// On new event was received, build projection and publish it to the projection bus
eventBus.subscribe { event ->
withLoggingContext("event" to event.toString()) {
gameEventStore
.getStream(event.aggregateId)
.readAll()
.fold(GameList(event.aggregateId)) { acc, event ->
acc.applyEvent(event)
}.also { save(it) }
.also { projectionBus.publish(it) }
}
}
}
}
@@ -0,0 +1,17 @@
package eventDemo.contexts.game.infrastructure.persistence.projections.bus
import eventDemo.contexts.game.application.ports.GameProjectionBus
import eventDemo.libs.bus.Bus
import eventDemo.libs.bus.BusInMemory
import eventDemo.shared.game.projection.GameProjection
import java.util.UUID
class GameProjectionBusInMemory :
GameProjectionBus,
Bus<GameProjection> by BusInMemory(GameProjectionBusInMemory::class),
Comparable<GameProjectionBusInMemory> {
private val instanceId: UUID = UUID.randomUUID()
override fun compareTo(other: GameProjectionBusInMemory): Int =
compareValues(instanceId, other.instanceId)
}
@@ -0,0 +1,25 @@
package eventDemo.contexts.game.infrastructure.persistence.projections.bus
import com.rabbitmq.client.ConnectionFactory
import eventDemo.contexts.game.application.ports.GameProjectionBus
import eventDemo.libs.bus.Bus
import eventDemo.libs.bus.BusInRabbitMQ
import eventDemo.shared.game.projection.GameProjection
import kotlinx.serialization.json.Json
import java.util.UUID
class GameProjectionBusInRabbitMQ(
private val connectionFactory: ConnectionFactory,
) : GameProjectionBus,
Bus<GameProjection> by BusInRabbitMQ(
connectionFactory,
"GameProjection",
{ Json.encodeToString(it) },
{ Json.decodeFromString<GameProjection>(it) },
),
Comparable<GameProjectionBusInRabbitMQ> {
private val instanceId: UUID = UUID.randomUUID()
override fun compareTo(other: GameProjectionBusInRabbitMQ): Int =
compareValues(instanceId, other.instanceId)
}
@@ -0,0 +1,26 @@
package eventDemo.contexts.game.infrastructure.rest
import eventDemo.domain.event.projection.GameListRepository
import io.ktor.resources.Resource
import io.ktor.server.auth.authenticate
import io.ktor.server.resources.get
import io.ktor.server.response.respond
import io.ktor.server.routing.Route
import kotlinx.serialization.Serializable
@Serializable
@Resource("/games")
class Games
/**
* API routes to show all games.
*/
fun Route.gamesListRoute(gameListRepository: GameListRepository) {
authenticate {
// Read the last played card on the game.
get<Games> {
val gameList = gameListRepository.getList()
call.respond(gameList)
}
}
}
@@ -0,0 +1,45 @@
package eventDemo.contexts.game.infrastructure.rest
import eventDemo.contexts.game.application.eventStores.GameRepository
import eventDemo.contexts.game.application.notification.toNotification
import eventDemo.contexts.game.application.ports.GameEventStore
import eventDemo.shared.ids.GameId
import eventDemo.shared.serializers.GameIdSerializer
import eventDemo.sharedKernel.currentUserId
import io.ktor.http.HttpStatusCode
import io.ktor.resources.Resource
import io.ktor.server.auth.authenticate
import io.ktor.server.resources.get
import io.ktor.server.response.respond
import io.ktor.server.routing.Route
import kotlinx.serialization.Serializable
@Serializable
@Resource("/games/{id}")
class Game(
@Serializable(with = GameIdSerializer::class)
val id: GameId,
)
/**
* API routes to read the game state.
*/
fun Route.getFullNotificationsRoute(
gameRepository: GameRepository,
gameEventStore: GameEventStore,
) {
authenticate {
get<Game> { body ->
val game =
gameRepository.get(body.id)
?: return@get call.respond(HttpStatusCode.NotFound)
val notifications =
gameEventStore
.getStream(body.id)
.readAll()
.flatMap { it.toNotification(game, call.currentUserId) }
call.respond(notifications)
}
}
}
@@ -0,0 +1,26 @@
package eventDemo.contexts.game.infrastructure.websocket
import eventDemo.contexts.game.application.channels.GameChannelsSubscriber
import eventDemo.libs.helpers.fromFrameChannel
import eventDemo.libs.helpers.toObjectChannel
import eventDemo.shared.ids.GameId
import eventDemo.sharedKernel.currentUserId
import io.ktor.server.auth.authenticate
import io.ktor.server.routing.Route
import io.ktor.server.websocket.webSocket
import kotlinx.coroutines.DelicateCoroutinesApi
import kotlin.uuid.Uuid
@DelicateCoroutinesApi
fun Route.gameWebSocket(channelSubscriber: GameChannelsSubscriber) {
authenticate {
webSocket("/games/{id}") {
channelSubscriber.subscribePlayerToGameChannels(
gameId = GameId(Uuid.parse(call.parameters["id"]!!)),
userId = call.currentUserId,
incomingCommandChannel = toObjectChannel(incoming),
sendNotificationChannel = fromFrameChannel(outgoing),
)
}
}
}
@@ -0,0 +1,27 @@
package eventDemo.libs.bus
interface Bus<T> {
/**
* Publish a new [message][item] to the bus.
*/
fun publish(item: T)
fun publish(items: Collection<T>) {
items.forEach { publish(it) }
}
/**
* Subscribe a [lambda][block] to the bus.
*
* When a message is sent to the bus, the [block] is executed.
*/
fun subscribe(block: (T) -> Unit): Subscription
/**
* The returns of the [subscribe] method.
* It can be called to [cancel][close] the subscription.
*/
interface Subscription : AutoCloseable {
override fun close()
}
}
@@ -0,0 +1,31 @@
package eventDemo.libs.bus
import io.github.oshai.kotlinlogging.KotlinLogging
import io.github.oshai.kotlinlogging.withLoggingContext
import kotlin.reflect.KClass
class BusInMemory<E>(
val name: KClass<*> = BusInMemory::class,
) : Bus<E> {
private val logger = KotlinLogging.logger(name.qualifiedName.toString())
private val subscribers: MutableList<(E) -> Unit> = mutableListOf()
override fun publish(item: E) {
withLoggingContext("busItem" to item.toString()) {
logger.info { "Item sent to the bus" }
subscribers
.forEach {
it(item)
}
}
}
override fun subscribe(block: (E) -> Unit): Bus.Subscription {
subscribers.add(block)
return object : Bus.Subscription {
override fun close() {
subscribers.remove(block)
}
}
}
}
@@ -0,0 +1,99 @@
package eventDemo.libs.bus
import com.rabbitmq.client.AMQP
import com.rabbitmq.client.BuiltinExchangeType
import com.rabbitmq.client.Connection
import com.rabbitmq.client.ConnectionFactory
import com.rabbitmq.client.DefaultConsumer
import com.rabbitmq.client.Envelope
import io.github.oshai.kotlinlogging.KotlinLogging
import io.github.oshai.kotlinlogging.withLoggingContext
import io.ktor.utils.io.core.toByteArray
import kotlinx.coroutines.runBlocking
class BusInRabbitMQ<E>(
private val connectionFactory: ConnectionFactory,
private val exchangeName: String,
private val objectToString: (E) -> String,
private val stringToObject: (String) -> E,
) : Bus<E> {
private val logger = KotlinLogging.logger { }
private val connection: Connection = connectionFactory.newConnection()
get() {
return if (field.isOpen) {
field
} else {
connectionFactory.newConnection()
}
}
private val routingKey = ""
init {
connection
.createChannel()
.use {
it.exchangeDeclare(
exchangeName,
BuiltinExchangeType.FANOUT,
true,
false,
emptyMap(),
)
}
}
override fun publish(item: E) {
withLoggingContext("item" to item.toString()) {
connection
.createChannel()
.basicPublish(
exchangeName,
routingKey,
AMQP.BasicProperties(),
objectToString(item).toByteArray(),
)
logger.info { "Item sent to the bus" }
}
}
override fun subscribe(block: (E) -> Unit): Bus.Subscription {
connection
.createChannel()
.also { channel ->
val queue =
channel
.queueDeclare()
.queue
.also { channel.queueBind(it, exchangeName, routingKey) }
channel
.basicConsume(
queue,
object : DefaultConsumer(channel) {
override fun handleDelivery(
consumerTag: String,
envelope: Envelope,
properties: AMQP.BasicProperties,
body: ByteArray,
) {
runBlocking {
val obj = stringToObject(body.toString(Charsets.UTF_8))
withLoggingContext("item" to obj.toString()) {
logger.info { "Received delivery of $exchangeName" }
}
block(obj)
}
channel.basicAck(envelope.deliveryTag, false)
}
},
)
}.let {
return object : Bus.Subscription {
override fun close() {
it.close()
}
}
}
}
}
@@ -0,0 +1,53 @@
package eventDemo.libs.command
import eventDemo.shared.command.Command
import eventDemo.shared.ids.CommandId
import kotlinx.datetime.Clock
import kotlinx.datetime.Instant
import java.util.concurrent.ConcurrentHashMap
import kotlin.time.Duration
import kotlin.time.Duration.Companion.minutes
/**
* Controls the execution of a command to prevent it from being executed more than once.
*/
class CommandUnicityChecker<C : Command>(
private val maxCacheTime: Duration = 10.minutes,
) {
private val executedCommand: ConcurrentHashMap<CommandId, Pair<Boolean, Instant>> = ConcurrentHashMap()
fun runOnlyOnce(
command: C,
action: (C) -> Unit,
) {
if (!isAlreadyExecuted(command)) {
action(command)
setAsExecuted(command)
removeOldCache()
} else {
throw UnicityException("Command already executed", command)
}
}
private fun setAsExecuted(command: C) {
executedCommand.computeIfAbsent(command.id) { Pair(true, Clock.System.now()) }
}
private fun removeOldCache() {
executedCommand
.filterValues { (_, date) ->
(date + maxCacheTime) < Clock.System.now()
}.keys
.forEach {
executedCommand.remove(it)
}
}
private fun isAlreadyExecuted(command: C): Boolean =
executedCommand[command.id]?.first ?: false
class UnicityException(
override val message: String,
val command: Command,
) : Exception(message)
}
@@ -0,0 +1,16 @@
package eventDemo.libs.eventSource
import eventDemo.shared.ids.AggregateId
import eventDemo.shared.ids.EventId
import kotlinx.datetime.Instant
/**
* The basic interface for an Event
* @see eventDemo.libs.eventSource.eventStore.EventStream
*/
interface Event<ID : AggregateId> {
val eventId: EventId
val aggregateId: ID
val createdAt: Instant
val version: Int
}
@@ -0,0 +1,19 @@
package eventDemo.libs.eventSource.eventStore
import eventDemo.libs.eventSource.Event
import eventDemo.shared.ids.AggregateId
import io.github.oshai.kotlinlogging.withLoggingContext
interface EventStore<E : Event<ID>, ID : AggregateId> {
fun getStream(aggregateId: ID): EventStream<E, ID>
@Throws(VersionConflictException::class)
fun append(event: E) =
withLoggingContext("event" to event.toString()) {
getStream(event.aggregateId).append(event)
}
@Throws(VersionConflictException::class)
fun append(events: Set<E>) =
events.forEach { append(it) }
}
@@ -0,0 +1,13 @@
package eventDemo.libs.eventSource.eventStore
import eventDemo.libs.eventSource.Event
import eventDemo.shared.ids.AggregateId
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.ConcurrentMap
class EventStoreInMemory<E : Event<ID>, ID : AggregateId> : EventStore<E, ID> {
private val streams: ConcurrentMap<ID, EventStream<E, ID>> = ConcurrentHashMap()
override fun getStream(aggregateId: ID): EventStream<E, ID> =
streams.computeIfAbsent(aggregateId) { EventStreamInMemory(aggregateId) }
}
@@ -0,0 +1,15 @@
package eventDemo.libs.eventSource.eventStore
import eventDemo.libs.eventSource.Event
import eventDemo.shared.ids.AggregateId
import javax.sql.DataSource
class EventStoreInPostgresql<E : Event<ID>, ID : AggregateId>(
private val dataSource: DataSource,
private val objectToString: (E) -> String,
private val stringToObject: (String) -> E,
private val tableName: String,
) : EventStore<E, ID> {
override fun getStream(aggregateId: ID): EventStream<E, ID> =
EventStreamInPostgresql(aggregateId, dataSource, objectToString, stringToObject, tableName)
}
@@ -0,0 +1,42 @@
package eventDemo.libs.eventSource.eventStore
import eventDemo.libs.eventSource.Event
import eventDemo.shared.ids.AggregateId
import io.github.oshai.kotlinlogging.withLoggingContext
/**
* Interface representing an event stream for publishing and reading domain events
*/
interface EventStream<E : Event<ID>, ID : AggregateId> {
val aggregateId: ID
/** Publishes a single event to the event stream */
@Throws(VersionConflictException::class)
fun append(event: E)
/** Publishes multiple events to the event stream */
fun append(vararg events: E) {
events.forEach {
withLoggingContext("event" to it.toString()) {
append(it)
}
}
}
/** Reads all events */
fun readAll(): Set<E>
fun readGreaterOfVersion(version: Int): Set<E> =
readVersionBetween(version + 1..Int.MAX_VALUE)
fun readVersionBetween(version: IntRange): Set<E>
fun getByVersion(version: Int): E? =
readVersionBetween(version..version).firstOrNull()
fun exist(): Boolean
}
class VersionConflictException(
event: Event<*>,
) : RuntimeException("Version conflict: ${event.version}")
@@ -0,0 +1,42 @@
package eventDemo.libs.eventSource.eventStore
import eventDemo.libs.eventSource.Event
import eventDemo.shared.ids.AggregateId
import io.github.oshai.kotlinlogging.KotlinLogging
import java.util.Queue
import java.util.concurrent.ConcurrentLinkedQueue
/**
* An In-Memory implementation of an event stream.
*
* All methods are implemented.
*/
class EventStreamInMemory<E : Event<ID>, ID : AggregateId>(
override val aggregateId: ID,
) : EventStream<E, ID> {
private val logger = KotlinLogging.logger {}
private val events: Queue<E> = ConcurrentLinkedQueue()
override fun append(event: E) {
if (event.aggregateId != aggregateId) {
throw EventStreamPublishException(
"You cannot publish this event in this stream because it has a different aggregateId!",
)
}
if (events.none { it.eventId == event.eventId }) {
events.add(event)
logger.info { "Event published" }
}
}
override fun readAll(): Set<E> =
events.toSet()
override fun readVersionBetween(version: IntRange): Set<E> =
events
.filter { version.contains(it.version) }
.toSet()
override fun exist(): Boolean =
events.isNotEmpty()
}
@@ -0,0 +1,142 @@
package eventDemo.libs.eventSource.eventStore
import eventDemo.libs.eventSource.Event
import eventDemo.shared.ids.AggregateId
import io.github.oshai.kotlinlogging.KotlinLogging
import io.github.oshai.kotlinlogging.withLoggingContext
import org.postgresql.util.PGobject
import org.postgresql.util.PSQLException
import javax.sql.DataSource
import kotlin.uuid.toJavaUuid
/**
* An In-Memory implementation of an event stream.
*
* All methods are implemented.
*/
class EventStreamInPostgresql<E : Event<ID>, ID : AggregateId>(
override val aggregateId: ID,
private val dataSource: DataSource,
private val objectToString: (E) -> String,
private val stringToObject: (String) -> E,
private val tableName: String,
) : EventStream<E, ID> {
private val logger = KotlinLogging.logger {}
override fun append(event: E) {
withLoggingContext("event" to event.toString()) {
if (event.aggregateId != aggregateId) {
throw EventStreamPublishException(
"You cannot publish this event in this stream because it has a different aggregateId!",
)
}
try {
dataSource.connection.use { connection ->
connection
.prepareStatement(
"""
insert into $tableName (id, aggregate_id, version, data)
values (?, ?, ?, ?)
on conflict (id) do nothing
""".trimIndent(),
).use {
it.setObject(1, event.eventId.id.toJavaUuid())
it.setObject(2, event.aggregateId.id.toJavaUuid())
it.setInt(3, event.version)
it.setObject(4, PGJsonb(objectToString(event)))
it.executeUpdate()
}
}
} catch (e: PSQLException) {
if (e.serverErrorMessage?.constraint == "game_event_stream_aggregate_id_version_key") {
logger.warn { "duplicate version" }
throw VersionConflictException(event)
} else {
throw e
}
}
logger.info { "Event appended" }
}
}
override fun readAll(): Set<E> =
dataSource.connection.use { connection ->
connection
.prepareStatement(
"""
select data
from $tableName
where aggregate_id = ?
order by version asc
""".trimIndent(),
).use {
it.setObject(1, aggregateId.id.toJavaUuid())
it.executeQuery().use { resultSet ->
buildSet {
while (resultSet.next()) {
resultSet
.getString("data")
.let(stringToObject)
.let { add(it) }
}
}
}
}
}
override fun exist(): Boolean =
dataSource.connection.use { connection ->
connection
.prepareStatement(
"""
select 1
from $tableName
where aggregate_id = ?
limit 1
order by version asc
""".trimIndent(),
).use {
it.setObject(1, aggregateId.id.toJavaUuid())
it.executeQuery().use { resultSet ->
resultSet.next()
}
}
}
override fun readVersionBetween(version: IntRange): Set<E> =
dataSource.connection.use { connection ->
connection
.prepareStatement(
"""
select data
from $tableName
where version between ? and ?
and aggregate_id = ?
order by version asc
""".trimIndent(),
).use { stmt ->
stmt.setInt(1, version.first)
stmt.setInt(2, version.last)
stmt.setObject(3, aggregateId.id.toJavaUuid())
stmt.executeQuery().use { resultSet ->
buildSet {
while (resultSet.next()) {
resultSet
.getString("data")
.let(stringToObject)
.let { add(it) }
}
}
}
}
}
}
class PGJsonb(
value: String,
) : PGobject() {
init {
this.value = value
this.type = "jsonb"
}
}
@@ -0,0 +1,5 @@
package eventDemo.libs.eventSource.eventStore
class EventStreamPublishException(
override val message: String,
) : Exception(message)
@@ -0,0 +1,49 @@
package eventDemo.libs.helpers
import io.github.oshai.kotlinlogging.KotlinLogging
import io.ktor.websocket.Frame
import io.ktor.websocket.readText
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.channels.ReceiveChannel
import kotlinx.coroutines.channels.SendChannel
import kotlinx.coroutines.channels.consumeEach
import kotlinx.coroutines.channels.produce
import kotlinx.coroutines.launch
import kotlinx.serialization.json.Json
/**
* Convert a [ReceiveChannel] of [Frame] to another [ReceiveChannel] of [object][T]
*/
@OptIn(ExperimentalCoroutinesApi::class)
inline fun <reified T> CoroutineScope.toObjectChannel(
frames: ReceiveChannel<Frame>,
bufferSize: Int = 0,
): ReceiveChannel<T> {
val logger = KotlinLogging.logger { }
return produce(capacity = bufferSize) {
frames.consumeEach { frame ->
if (frame is Frame.Text) {
val frameText = frame.readText()
logger.debug { "Conversion of the Frame: $frameText to ${T::class.simpleName}" }
send(Json.decodeFromString(frameText))
} else {
logger.warn { "The frame is not a text frame" }
}
}
}
}
/**
* Convert a [SendChannel] of [Frame] to another [SendChannel] of [object][T]
*/
inline fun <reified T> CoroutineScope.fromFrameChannel(frames: SendChannel<Frame>): SendChannel<T> {
val channel = Channel<T>()
launch {
channel.consumeEach { obj ->
frames.send(Frame.Text(Json.encodeToString(obj)))
}
}
return channel
}
@@ -0,0 +1,11 @@
package eventDemo.libs.helpers
fun List<Int>.toRanges(): List<IntRange> =
fold(listOf()) { acc, i ->
val last = acc.lastOrNull()
if (last != null && last.max() + 1 == i) {
(acc - setOf(last)) + setOf(IntRange(last.min(), i))
} else {
acc + setOf(IntRange(i, i))
}
}
@@ -0,0 +1,13 @@
package eventDemo.sharedKernel
import eventDemo.shared.ids.UserId
import io.ktor.server.application.ApplicationCall
import io.ktor.server.auth.jwt.JWTPrincipal
import io.ktor.server.auth.principal
import kotlin.uuid.Uuid
internal val ApplicationCall.currentUserId: UserId
get() =
principal<JWTPrincipal>()!!.run {
UserId(Uuid.parse(payload.getClaim("userid").asString()))
}
@@ -0,0 +1,38 @@
ktor {
deployment {
port = 8080
}
application {
modules = [ eventDemo.configuration.ConfigureKtorKt.configure ]
}
}
jwt {
secret = "secret"
secret = ${?JWT_SECRET}
}
postgresql {
url = "jdbc:postgresql://localhost:5432/event-demo"
url = ${?POSTGRESQL_URL}
username = "event-demo"
username = ${?POSTGRESQL_USERNAME}
password = "changeit"
password = ${?POSTGRESQL_PASSWORD}
}
rabbitmq {
url = "localhost"
url = ${?RABBITMQ_URL}
port = "5672"
port = ${?RABBITMQ_PORT}
username = "event-demo"
username = ${?RABBITMQ_USERNAME}
password = "changeit"
password = ${?RABBITMQ_PASSWORD}
}
+13
View File
@@ -0,0 +1,13 @@
<configuration>
<appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{YYYY-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n > MDC=%mdc%n</pattern>
</encoder>
</appender>
<root level="trace">
<appender-ref ref="STDOUT"/>
</root>
<logger name="org.eclipse.jetty" level="INFO"/>
<logger name="io.netty" level="INFO"/>
<Logger name="com.zaxxer.hikari" level="WARN" additivity="true"/>
</configuration>
+11
View File
@@ -0,0 +1,11 @@
package eventDemo
import io.kotest.core.Tag
object Tag {
object Postgresql : Tag()
object RabbitMQ : Tag()
object Concurrence : Tag()
}

Some files were not shown because too many files have changed in this diff Show More