11 Commits
Author SHA1 Message Date
flecomte da13b0d2b8 chore: update docker compose in CI/CD
Tests / build (push) Successful in 6m45s
Tests / test (push) Failing after 11m30s
Tests / lint (push) Successful in 14m32s
2026-08-05 00:06:23 +02:00
flecomte d632dc0f6b chore: refactor docker compose env's
Tests / build (push) Successful in 6m41s
Tests / test (push) Failing after 10m1s
Tests / lint (push) Successful in 13m29s
2026-08-04 23:38:01 +02:00
flecomte b60e4aa457 chore: run CI test in docker
Tests / build (push) Successful in 7m13s
Tests / test (push) Failing after 9m19s
Tests / lint (push) Successful in 13m33s
2026-08-04 20:28:49 +02:00
flecomte e262f35b27 chore: eol=lf 2026-08-01 22:13:15 +02:00
flecomte 77e4cab8ad chore: split docker-compose-tools-local 2026-07-31 22:09:46 +02:00
flecomte 7291419c9f chore: fix postgresql.secret in ci
Tests / build (push) Successful in 6m17s
Tests / test (push) Failing after 11m40s
Tests / lint (push) Successful in 14m20s
2026-07-31 21:51:27 +02:00
flecomte b2b8fcf92f refactor: fix cast warning 2026-07-31 21:36:21 +02:00
flecomte 4e4b307275 chore: fix cache for copyEnv gradle task
Tests / build (push) Successful in 9m51s
Tests / test (push) Failing after 11m0s
Tests / lint (push) Successful in 18m7s
2026-07-31 00:58:27 +02:00
flecomte e87a36caa5 chore: update CI actions versions 2026-07-31 00:24:17 +02:00
flecomte e2d7942c7e refactoring: Masive refactor to build the V2
Tests / build (push) Successful in 7m14s
Tests / lint (push) Successful in 8m35s
Tests / test (push) Failing after 11m0s
2026-07-30 23:01:50 +02:00
flecomte b313b39cf4 chore: config db in pgadmin 2026-07-30 22:42:54 +02:00
280 changed files with 5349 additions and 4931 deletions
+26
View File
@@ -0,0 +1,26 @@
# Version control
.git
.github
# Gradle build outputs / caches (must always be rebuilt fresh inside the image)
.gradle
build/
.kotlin
!gradle/wrapper/gradle-wrapper.jar
# IDE
.idea
.vscode
.run
*.iml
*.iws
*.ipr
# Docker-only local files (secrets/env must never be baked into the image)
docker/.env
docker/*.env.docker
docker/*.secret
# Misc
*.hprof
.gradle-docker-cache/
+7
View File
@@ -0,0 +1,7 @@
* text=auto
* eol=lf
*.sh text eol=lf
*.png binary
*.jar binary
gradlew.bat eol=crlf
gradlew text eol=lf
+41 -29
View File
@@ -18,10 +18,10 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@v4
uses: actions/checkout@v6
- name: Set up JDK 21
uses: actions/setup-java@v4
uses: actions/setup-java@v5
with:
distribution: 'temurin'
java-version: '21'
@@ -31,7 +31,7 @@ jobs:
run: echo "key=gradle-${{ runner.os }}-${{ hashFiles('**/*.gradle*', '**/gradle-wrapper.properties') }}" >> $GITHUB_OUTPUT
- name: Cache Gradle dependencies
uses: actions/cache@v3
uses: actions/cache@v6
with:
path: |
~/.gradle/caches
@@ -48,16 +48,16 @@ jobs:
runs-on: ubuntu-latest
steps:
- name: Checkout code
uses: actions/checkout@v4
uses: actions/checkout@v6
- name: Set up JDK 21
uses: actions/setup-java@v4
uses: actions/setup-java@v5
with:
distribution: 'temurin'
java-version: '21'
- name: Restore Gradle cache
uses: actions/cache@v3
uses: actions/cache@v6
with:
path: |
~/.gradle/caches
@@ -73,53 +73,65 @@ jobs:
run: ./gradlew ktlintCheck
- name: Publish ktlint report
uses: yutailang0119/action-ktlint@v4
uses: yutailang0119/action-ktlint@v5
if: always()
with:
report-path: build/reports/ktlint/**/*.xml
continue-on-error: false
test:
needs: build
runs-on: ubuntu-latest
env:
GRADLE_CACHE_DIR: ${{ github.workspace }}/.gradle-docker-cache
steps:
- name: Checkout code
uses: actions/checkout@v4
uses: actions/checkout@v6
- name: Set up JDK 21
uses: actions/setup-java@v4
with:
distribution: 'temurin'
java-version: '21'
- name: Install a pinned Docker Compose version
run: |
mkdir -p ~/.docker/cli-plugins
curl -fSL https://github.com/docker/compose/releases/download/v5.1.4/docker-compose-linux-x86_64 \
-o ~/.docker/cli-plugins/docker-compose
chmod +x ~/.docker/cli-plugins/docker-compose
docker compose version
- name: Restore Gradle cache
uses: actions/cache@v3
- name: Prepare docker secrets
run: |
[ -f docker/postgresql.secret ] || echo -n "changeit" > docker/postgresql.secret
- name: Generate cache key
id: cache-key-generator
run: echo "key=gradle-docker-${{ runner.os }}-${{ hashFiles('**/*.gradle*', '**/gradle-wrapper.properties') }}" >> $GITHUB_OUTPUT
- name: Restore Gradle cache (Docker)
uses: actions/cache@v6
with:
path: |
~/.gradle/caches
~/.gradle/wrapper
key: ${{ needs.build.outputs.cache-key }}
path: ${{ env.GRADLE_CACHE_DIR }}
key: ${{ steps.cache-key-generator.outputs.key }}
restore-keys: |
gradle-${{ runner.os }}-
gradle-docker-${{ runner.os }}-
- name: Grant execute permission to Gradle wrapper
run: chmod +x gradlew
- name: Prepare cache directory permissions
run: |
mkdir -p "$GRADLE_CACHE_DIR"
chmod -R 777 "$GRADLE_CACHE_DIR"
- name: Start CI Docker Compose services
run: ./gradlew composeUp -Pci
- name: Run tests in Docker
run: docker compose -f docker/docker-compose-test.yaml run tests
- name: Run tests
run: ./gradlew test -x composeUp --no-daemon
- name: Shut down Docker services
if: always()
run: docker compose -f docker/docker-compose-test.yaml down -v
- name: Upload test reports
if: always()
uses: actions/upload-artifact@v4
uses: actions/upload-artifact@v7
with:
name: test-results
path: build/reports/tests/test
- name: Publish Test Report
uses: dorny/test-reporter@v1
uses: dorny/test-reporter@v3
if: always()
with:
name: JUnit Tests
+1
View File
@@ -37,3 +37,4 @@ out/
/docker/.env
/docker/*.secret
*.hprof
/.gradle-docker-cache/
+32
View File
@@ -0,0 +1,32 @@
<DataSourcesHistory>
<DataSourceFromHistory isRemovedFromProject="false">
<data-source source="LOCAL" name="event-demo@localhost" uuid="af2eabb1-64f7-49de-a94f-be1560baa96a">
<database-info product="PostgreSQL" version="18.4 (Debian 18.4-1.pgdg13+1)" jdbc-version="4.2" driver-name="PostgreSQL JDBC Driver" driver-version="42.7.3" dbms="POSTGRES" exact-version="18.4" exact-driver-version="42.7">
<identifier-quote-string>&quot;</identifier-quote-string>
</database-info>
<case-sensitivity plain-identifiers="lower" quoted-identifiers="exact" />
<driver-ref>postgresql</driver-ref>
<synchronize>true</synchronize>
<jdbc-driver>org.postgresql.Driver</jdbc-driver>
<jdbc-url>jdbc:postgresql://localhost:5432/event-demo</jdbc-url>
<secret-storage>master_key</secret-storage>
<user-name>event-demo</user-name>
<schema-mapping>
<introspection-scope>
<node negative="1">
<node kind="database" qname="@">
<node kind="schema" qname="@" />
</node>
<node kind="database" qname="event-demo">
<node kind="schema">
<name qname="auth" />
<name qname="game" />
</node>
</node>
</node>
</introspection-scope>
</schema-mapping>
<working-dir>$ProjectFileDir$</working-dir>
</data-source>
</DataSourceFromHistory>
</DataSourcesHistory>
+27
View File
@@ -0,0 +1,27 @@
<component name="ProjectRunConfigurationManager">
<configuration default="false" name="docker composeUp" type="GradleRunConfiguration" factoryName="Gradle">
<ExternalSystemSettings>
<option name="executionName" />
<option name="externalProjectPath" value="$PROJECT_DIR$" />
<option name="externalSystemIdString" value="GRADLE" />
<option name="scriptParameters" value="" />
<option name="taskDescriptions">
<list />
</option>
<option name="taskNames">
<list>
<option value="composeUp" />
</list>
</option>
<option name="vmOptions" />
</ExternalSystemSettings>
<ExternalSystemDebugServerProcess>true</ExternalSystemDebugServerProcess>
<ExternalSystemReattachDebugProcess>true</ExternalSystemReattachDebugProcess>
<ExternalSystemDebugDisabled>false</ExternalSystemDebugDisabled>
<DebugAllEnabled>false</DebugAllEnabled>
<RunAsTest>false</RunAsTest>
<GradleProfilingDisabled>false</GradleProfilingDisabled>
<GradleCoverageDisabled>false</GradleCoverageDisabled>
<method v="2" />
</configuration>
</component>
+32 -79
View File
@@ -1,22 +1,19 @@
@file:Suppress("PropertyName")
import org.jlleitschuh.gradle.ktlint.KtlintExtension
val ktor_version: String by project
val kotlin_version: String by project
val kotlin_serialization_version: String by project
val logback_version: String by project
val koin_version: String by project
val kotlin_logging_version: String by project
val kotest_version: String by project
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") version "2.1.21"
id("io.ktor.plugin") version "3.5.1"
id("org.jetbrains.kotlin.plugin.serialization") version "2.4.10"
id("org.jlleitschuh.gradle.ktlint") version "12.2.0"
id("com.avast.gradle.docker-compose") version "0.17.12"
id("org.jlleitschuh.gradle.ktlint") version "14.2.0"
}
group = "io.github.flecomte"
@@ -29,7 +26,7 @@ application {
}
configure<KtlintExtension> {
version.set("1.5.0")
version.set("1.8.0")
}
ktlint {
reporters {
@@ -49,63 +46,17 @@ java {
tasks.withType<Test>().configureEach {
useJUnitPlatform()
}
dockerCompose {
val composeFile =
if (project.hasProperty("ci")) {
// Use docker-compose-ci.yaml for the CI
"docker/docker-compose-ci.yaml"
} else {
// Use docker-compose-test.yaml for local tests
"docker/docker-compose-test.yaml"
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")
}
useComposeFiles.set(listOf(composeFile))
setProjectName("event-demo-test")
}
tasks.test {
dependsOn("composeUp")
dockerCompose.useComposeFiles.set(listOf("docker/docker-compose-test.yaml"))
dockerCompose.setProjectName("event-demo-test")
}
tasks.named("run") {
dependsOn("composeUp")
dockerCompose.useComposeFiles.set(listOf("docker/docker-compose-test.yaml"))
dockerCompose.setProjectName("event-demo-dev")
}
tasks.register<Copy>("copyEnv") {
group = "docker"
description = "copy the default dotenv file"
from("docker")
into("docker")
rename {
it.removeSuffix(".template")
}
include(".env.template")
eachFile {
if (File("docker/$name").exists()) {
exclude()
}
}
doLast {
val files =
listOf(
File("docker/pgadmin.secret"),
File("docker/postgresql.secret"),
)
files.forEach {
if (!it.exists()) {
it.writeText("changeit")
}
}
}
}
tasks.composeUp {
dependsOn("copyEnv")
}
dependencies {
@@ -124,23 +75,25 @@ dependencies {
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:$logback_version")
implementation("io.insert-koin:koin-ktor:$koin_version")
implementation("io.insert-koin:koin-logger-slf4j:$koin_version")
implementation("io.github.oshai:kotlin-logging-jvm:$kotlin_logging_version")
implementation("org.jetbrains.kotlinx:kotlinx-serialization-json-jvm:$kotlin_serialization_version")
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("redis.clients:jedis:5.2.0")
implementation("org.postgresql:postgresql:42.7.5")
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:$kotest_version")
testImplementation("org.jetbrains.kotlin:kotlin-test-junit:$kotlin_version")
testImplementation("io.ktor:ktor-server-test-host-jvm:$ktor_version")
testImplementation("io.kotest:kotest-runner-junit5:$kotest_version")
testImplementation("io.mockk:mockk:1.13.17")
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")
}
+67
View File
@@ -0,0 +1,67 @@
# Exemple de structure
Les couches, du plus interne au plus externe
```
Domain (le cœur, ne dépend de RIEN d'externe)
Application (orchestre le Domain, ne connaît pas l'infra concrète)
Infrastructure (WebSocket, DB, event store — dépend de tout le reste)
```
```
src/
└── contexts/
├── auth/
└── ...
└── game/
├── domain/ ← Le cœur métier, zéro dépendance externe
│ ├── game/
│ │ ├── Game.ts ← Aggregate Root
│ │ ├── Player.ts ← Entity interne
│ │ ├── Card.ts ← Entity
│ │ ├── Color.ts ← Value Object
│ │ ├── Deck.ts ← VO ou petite structure
│ │ └── errors/
│ │ ├── InvalidMoveError.ts
│ │ └── ColorChoiceRequiredError.ts
│ └── events/ ← Events de DOMAINE (internes)
│ ├── CardPlayed.ts
│ ├── CardDrawn.ts
│ ├── TurnPassed.ts
│ └── DomainEvent.ts ← interface/type de base
├── application/ ← Orchestration, cas d'usage
│ ├── commands/ ← Les Commandes (intentions)
│ │ ├── PlayCardCommand.ts
│ │ └── DrawCardCommand.ts
│ ├── handlers/ ← Un handler par commande
│ │ ├── PlayCardHandler.ts ← charge l'aggregate, appelle game.playCard(), save
│ │ └── DrawCardHandler.ts
│ ├── projections/ ← LA LOGIQUE de construction des projections
│ │ ├── GameSummaryProjector.kt ← écoute les events, met à jour la vue
│ │ └── PlayerStatsProjector.kt
│ └── ports/ ← INTERFACES seulement (le "hexagone")
│ ├── GameRepository.ts ← interface, pas d'implémentation
│ ├── EventPublisher.ts ← interface, pas d'implémentation
│ └── ProjectionStore.kt ← interface, où lire/écrire la projection
├── infrastructure/ ← Tout ce qui est technique/externe
│ ├── persistence/
│ │ ├── EventStoreGameRepository.ts ← implémente GameRepository
│ │ ├── EventStore.ts
│ │ ├── projections/
│ │ │ ├── GameSummaryProjectionStore.kt ← implémentation concrète (DB, table dédiée)
│ │ │ └── models/
│ │ │ └── GameSummaryView.kt ← structure de la vue elle-même
│ ├── websocket/
│ │ ├── WebSocketServer.ts
│ │ ├── connectionManager.ts ← Map<gameId, Map<playerId, WebSocket>>
│ │ └── commandRouter.ts ← reçoit le message brut, dispatch vers le bon handler
│ └── eventPublisher/
│ └── WebSocketEventPublisher.ts ← implémente EventPublisher, fait le broadcast
└── presentation/ ← Traduction vers/depuis le client (le fameux DTO layer)
├── clientEvents/
│ ├── ClientEvent.ts ← types des events envoyés au front
│ └── toClientEvent.ts ← fonction de traduction domain event → client event
└── clientCommands/
└── parseIncomingCommand.ts ← valide/parse le message brut du client → Command
```
+13 -1
View File
@@ -1,12 +1,24 @@
Installation
============
To run the stack:
To run the stack in production:
```shell
docker compose -f docker\docker-compose-prod.yaml -p event-demo up -d
```
To run only the app dependencies in development mode and run the app localy (not in docker):
```shell
docker compose -f docker\docker-compose-dev.yaml up -d
```
To run the tests in docker (it's designed for the CI):
```shell
docker compose -f docker\docker-compose-test.yaml up -d
```
Api url:
- [Backend API](http://api.traefik.me/)
- [Frontend web site](http://app.traefik.me/) (WIP)
+3
View File
@@ -0,0 +1,3 @@
REDIS_URL=redis://redis:6379
POSTGRESQL_URL=jdbc:postgresql://postgresql/event-demo
RABBITMQ_URL=rabbitmq
-1
View File
@@ -1 +0,0 @@
PGADMIN_DEFAULT_EMAIL=
+4 -4
View File
@@ -1,5 +1,5 @@
# Stage 1: Cache Gradle dependencies
FROM gradle:latest AS cache
FROM gradle:9.6.1-jdk21-alpine AS cache
RUN mkdir -p /home/gradle/cache_home
ENV GRADLE_USER_HOME=/home/gradle/cache_home
COPY build.gradle.* gradle.properties /home/gradle/app/
@@ -7,7 +7,7 @@ WORKDIR /home/gradle/app
RUN gradle build -i -x check
# Stage 2: Build Application
FROM gradle:latest AS build
FROM gradle:9.6.1-jdk21-alpine AS build
COPY --from=cache /home/gradle/cache_home /home/gradle/.gradle
COPY --chown=gradle:gradle . /home/gradle/src
WORKDIR /home/gradle/src
@@ -16,8 +16,8 @@ WORKDIR /home/gradle/src
RUN gradle buildFatJar --no-daemon
# Stage 3: Create the Runtime Image
FROM amazoncorretto:21 AS runtime
FROM eclipse-temurin:21-jre-alpine AS runtime
EXPOSE 8080
RUN mkdir /app
COPY --from=build /home/gradle/src/build/libs/*.jar /app/event-demo-all.jar
COPY --from=build /home/gradle/src/build/libs/*-all.jar /app/event-demo-all.jar
ENTRYPOINT ["java","-jar","/app/event-demo-all.jar"]
+10
View File
@@ -0,0 +1,10 @@
# Image officielle Gradle avec JDK 21 déjà installé
FROM gradle:9.6.1-jdk21
WORKDIR /app
# Copie du wrapper et des fichiers de config en premier pour profiter du cache Docker
COPY build.gradle.kts settings.gradle.kts ./
# Lance les tests Kotlin
CMD ["gradle", "test", "--no-daemon"]
-6
View File
@@ -1,6 +0,0 @@
name: event-demo-test
include:
- path:
- parts/docker-compose-databases.yaml
- parts/docker-compose-databases-expose.yaml
- parts/docker-compose-traefik.yaml
+17
View File
@@ -0,0 +1,17 @@
name: event-demo-dev
include:
- path:
- parts/docker-compose-databases.yaml
- parts/docker-compose-databases-expose.yaml
- parts/docker-compose-tools.yaml
- parts/docker-compose-tools-local.yaml
- parts/docker-compose-traefik.yaml
services:
postgresql:
environment:
POSTGRES_PASSWORD: "changeit"
pgadmin:
environment:
PGADMIN_DEFAULT_PASSWORD: "changeit"
+13
View File
@@ -5,3 +5,16 @@ include:
- parts/docker-compose-app.yaml
- parts/docker-compose-tools.yaml
- parts/docker-compose-traefik.yaml
services:
postgresql:
environment:
POSTGRES_PASSWORD_FILE: /run/secrets/postgresql_password
volumes:
- ./postgresql.secret:/run/secrets/postgresql_password:ro
pgadmin:
environment:
PGADMIN_DEFAULT_PASSWORD_FILE: /run/secrets/pgadmin_password
volumes:
- ./pgadmin.secret:/run/secrets/pgadmin_password:ro
+6 -2
View File
@@ -2,6 +2,10 @@ name: event-demo-test
include:
- path:
- parts/docker-compose-databases.yaml
- parts/docker-compose-databases-expose.yaml
- parts/docker-compose-tools.yaml
- parts/docker-compose-test.yaml
- parts/docker-compose-traefik.yaml
services:
postgresql:
environment:
POSTGRES_PASSWORD: "changeit"
+2
View File
@@ -12,6 +12,8 @@ services:
condition: service_healthy
redis:
condition: service_healthy
env_file:
- ../.env.docker
labels:
- "traefik.http.routers.api.rule=Host(`api.traefik.me`)"
- "traefik.http.services.api.loadbalancer.server.port=8080"
@@ -22,10 +22,7 @@ services:
image: postgres:18.4
command: postgres -c 'max_connections=500'
environment:
POSTGRES_PASSWORD_FILE: /run/secrets/postgresql_password
POSTGRES_USER: event-demo
secrets:
- postgresql_password
healthcheck:
test: ["CMD-SHELL", "sh -c 'pg_isready -U event-demo'"]
interval: 1s
@@ -47,10 +44,6 @@ services:
volumes:
- rabbitmq_data:/var/lib/rabbitmq/
secrets:
postgresql_password:
file: ../postgresql.secret
volumes:
redis_data:
redisinsight_data:
+22
View File
@@ -0,0 +1,22 @@
services:
tests:
build:
context: ../..
dockerfile: docker/DockerfileTest
volumes:
- ${GRADLE_CACHE_DIR:-gradle-cache}:/home/gradle/.gradle
- ../..:/app
depends_on:
flyway:
condition: service_completed_successfully
postgresql:
condition: service_healthy
rabbitmq:
condition: service_healthy
redis:
condition: service_healthy
env_file:
- ../.env.docker
volumes:
gradle-cache:
@@ -0,0 +1,18 @@
services:
pgadmin:
environment:
PGADMIN_CONFIG_SERVER_MODE: 'False'
PGADMIN_CONFIG_MASTER_PASSWORD_REQUIRED: 'False'
configs:
- source: pgpass
target: /pgpass
mode: 0600
uid: "5050"
gid: "5050"
- source: servers_json
target: /pgadmin4/servers.json
configs:
pgpass:
content: |
*:*:*:event-demo:changeit
+21 -7
View File
@@ -2,12 +2,12 @@ services:
pgadmin:
image: dpage/pgadmin4
environment:
PGADMIN_DEFAULT_EMAIL: $PGADMIN_DEFAULT_EMAIL
PGADMIN_DEFAULT_PASSWORD_FILE: /run/secrets/pgadmin_password
secrets:
- pgadmin_password
PGADMIN_DEFAULT_EMAIL: ${PGADMIN_DEFAULT_EMAIL:-admin@event-demo.dev}
volumes:
- pgadmin_data:/var/lib/pgadmin
configs:
- source: servers_json
target: /pgadmin4/servers.json
labels:
- "traefik.http.routers.pgadmin.rule=Host(`pgadmin.postgresql.traefik.me`)"
- "traefik.http.services.pgadmin.loadbalancer.server.port=80"
@@ -24,9 +24,23 @@ services:
- "traefik.http.routers.rabbitmq-management.service=rabbitmq-management"
- "traefik.http.services.rabbitmq-management.loadbalancer.server.port=15672"
secrets:
pgadmin_password:
file: ../pgadmin.secret
configs:
servers_json:
content: |
{
"Servers": {
"1": {
"Name": "Event demo",
"Group": "Servers",
"Host": "postgresql",
"Port": 5432,
"MaintenanceDB": "event-demo",
"Username": "event-demo",
"PassFile": "/pgpass",
"SSLMode": "prefer"
}
}
}
volumes:
pgadmin_data:
+1 -1
View File
@@ -1,6 +1,6 @@
services:
traefik:
image: traefik:3.3.4
image: traefik:3.7.9
command:
- "--api.insecure=true"
- "--api.dashboard=true"
@@ -1,4 +1,5 @@
create table event_stream (
create schema game;
create table game.game_event_stream (
id uuid not null primary key,
aggregate_id uuid not null,
version int not null,
@@ -0,0 +1,8 @@
create schema auth;
create table auth.user_event_stream (
id uuid not null primary key,
aggregate_id uuid not null,
version int not null,
data jsonb not null,
unique(aggregate_id, version)
);
@@ -0,0 +1,6 @@
create table auth.user (
id uuid not null primary key,
username text not null,
unique(id),
unique(username)
);
@@ -1,14 +0,0 @@
package eventDemo.adapter.infrastructure.event
import eventDemo.domain.entity.GameId
import eventDemo.domain.event.GameEventStore
import eventDemo.domain.event.event.GameEvent
import eventDemo.libs.event.EventStore
import eventDemo.libs.event.EventStoreInMemory
/**
* A stream to publish and read the played card event.
*/
class GameEventStoreInMemory :
GameEventStore,
EventStore<GameEvent, GameId> by EventStoreInMemory()
@@ -1,21 +0,0 @@
package eventDemo.adapter.infrastructure.event
import eventDemo.domain.entity.GameId
import eventDemo.domain.event.GameEventStore
import eventDemo.domain.event.event.GameEvent
import eventDemo.libs.event.EventStore
import eventDemo.libs.event.EventStoreInPostgresql
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) },
)
@@ -1,44 +0,0 @@
package eventDemo.adapter.infrastructure.event.projection
import eventDemo.domain.entity.GameId
import eventDemo.domain.event.GameEventBus
import eventDemo.domain.event.projection.GameList
import eventDemo.domain.event.projection.GameListRepository
import eventDemo.domain.event.projection.GameProjectionBus
import eventDemo.domain.event.projection.GameState
import eventDemo.domain.event.projection.apply
import eventDemo.libs.event.projection.ProjectionRepositoryInMemory
import io.github.oshai.kotlinlogging.withLoggingContext
/**
* Manages [projections][GameList], their building and publication in the [bus][GameProjectionBus].
*/
class GameListRepositoryInMemory : GameListRepository {
private val projectionsRepository =
ProjectionRepositoryInMemory(
applyToProjection = GameList::apply,
initialStateBuilder = { aggregateId: GameId -> GameList(aggregateId) },
)
fun subscribeToBus(
projectionBus: GameProjectionBus,
eventBus: GameEventBus,
) {
// On new event was received, build projection and publish it to the projection bus
eventBus.subscribe { event ->
withLoggingContext("event" to event.toString()) {
projectionsRepository
.applyAndSave(event)
.also { projectionBus.publish(it) }
}
}
}
/**
* Get the last version of the [GameState] from the all eventStream.
*
* It fetches it from the local cache if possible, otherwise it builds it.
*/
override fun getList(): List<GameList> =
projectionsRepository.getList()
}
@@ -1,51 +0,0 @@
package eventDemo.adapter.infrastructure.event.projection
import eventDemo.domain.entity.GameId
import eventDemo.domain.event.GameEventBus
import eventDemo.domain.event.projection.GameList
import eventDemo.domain.event.projection.GameListRepository
import eventDemo.domain.event.projection.GameProjectionBus
import eventDemo.domain.event.projection.GameState
import eventDemo.domain.event.projection.apply
import eventDemo.libs.event.projection.ProjectionRepositoryInRedis
import io.github.oshai.kotlinlogging.withLoggingContext
import kotlinx.serialization.json.Json
import redis.clients.jedis.UnifiedJedis
/**
* Manages [projections][GameList], their building and publication in the [bus][GameProjectionBus].
*/
class GameListRepositoryInRedis(
jedis: UnifiedJedis,
) : GameListRepository {
private val projectionsRepository =
ProjectionRepositoryInRedis(
initialStateBuilder = { aggregateId: GameId -> GameList(aggregateId) },
projectionClass = GameList::class,
projectionToJson = { Json.encodeToString(GameList.serializer(), it) },
jsonToProjection = { Json.decodeFromString(GameList.serializer(), it) },
applyToProjection = GameList::apply,
jedis = jedis,
)
fun subscribeToBus(
projectionBus: GameProjectionBus,
eventBus: GameEventBus,
) {
eventBus.subscribe { event ->
withLoggingContext("event" to event.toString()) {
projectionsRepository
.applyAndSave(event)
.also { projectionBus.publish(it) }
}
}
}
/**
* Get the last version of the [GameState] from the all eventStream.
*
* It fetches it from the local cache if possible, otherwise it builds it.
*/
override fun getList(): List<GameList> =
projectionsRepository.getList()
}
@@ -1,41 +0,0 @@
package eventDemo.adapter.infrastructure.event.projection
import eventDemo.domain.entity.GameId
import eventDemo.domain.event.GameEventBus
import eventDemo.domain.event.projection.GameProjectionBus
import eventDemo.domain.event.projection.GameState
import eventDemo.domain.event.projection.GameStateRepository
import eventDemo.domain.event.projection.apply
import eventDemo.libs.event.projection.ProjectionRepositoryInMemory
import io.github.oshai.kotlinlogging.withLoggingContext
/**
* Manages [projections][GameState], their building and publication in the [bus][GameProjectionBus].
*/
class GameStateRepositoryInMemory : GameStateRepository {
private val projectionsRepository =
ProjectionRepositoryInMemory(
applyToProjection = GameState::apply,
initialStateBuilder = { aggregateId: GameId -> GameState(aggregateId) },
)
fun subscribeToBus(
projectionBus: GameProjectionBus,
eventBus: GameEventBus,
) {
// On new event was received, build projection and publish it to the projection bus
eventBus.subscribe { event ->
withLoggingContext("event" to event.toString()) {
projectionsRepository
.applyAndSave(event)
.also { projectionBus.publish(it) }
}
}
}
/**
* Get the [GameState].
*/
override fun get(gameId: GameId): GameState =
projectionsRepository.get(gameId)
}
@@ -1,49 +0,0 @@
package eventDemo.adapter.infrastructure.event.projection
import eventDemo.domain.entity.GameId
import eventDemo.domain.event.GameEventBus
import eventDemo.domain.event.projection.GameProjectionBus
import eventDemo.domain.event.projection.GameState
import eventDemo.domain.event.projection.GameStateRepository
import eventDemo.domain.event.projection.apply
import eventDemo.libs.event.projection.ProjectionRepositoryInRedis
import io.github.oshai.kotlinlogging.withLoggingContext
import kotlinx.serialization.json.Json
import redis.clients.jedis.UnifiedJedis
/**
* Manages [projections][GameState], their building and publication in the [bus][GameProjectionBus].
*/
class GameStateRepositoryInRedis(
jedis: UnifiedJedis,
) : GameStateRepository {
private val projectionsRepository =
ProjectionRepositoryInRedis(
initialStateBuilder = { aggregateId: GameId -> GameState(aggregateId) },
projectionClass = GameState::class,
projectionToJson = { Json.encodeToString(GameState.serializer(), it) },
jsonToProjection = { Json.decodeFromString(GameState.serializer(), it) },
applyToProjection = GameState::apply,
jedis = jedis,
)
fun subscribeToBus(
projectionBus: GameProjectionBus,
eventBus: GameEventBus,
) {
// On new event was received, build projection and publish it to the projection bus
eventBus.subscribe { event ->
withLoggingContext("event" to event.toString()) {
projectionsRepository
.applyAndSave(event)
.also { projectionBus.publish(it) }
}
}
}
/**
* Get the [GameState].
*/
override fun get(gameId: GameId): GameState =
projectionsRepository.get(gameId)
}
@@ -1,67 +0,0 @@
package eventDemo.adapter.presenter.query
import eventDemo.domain.command.GameCommandHandler
import eventDemo.domain.command.command.GameCommand
import eventDemo.domain.entity.GameId
import eventDemo.domain.event.projection.projectionListener.PlayerNotificationListener
import eventDemo.domain.notification.Notification
import eventDemo.libs.fromFrameChannel
import eventDemo.libs.toObjectChannel
import io.github.oshai.kotlinlogging.withLoggingContext
import io.ktor.server.auth.authenticate
import io.ktor.server.routing.Route
import io.ktor.server.websocket.DefaultWebSocketServerSession
import io.ktor.server.websocket.webSocket
import kotlinx.coroutines.DelicateCoroutinesApi
import kotlinx.coroutines.GlobalScope
import kotlinx.coroutines.channels.ReceiveChannel
import kotlinx.coroutines.channels.SendChannel
import kotlinx.coroutines.channels.trySendBlocking
import kotlinx.coroutines.launch
import java.util.UUID
@DelicateCoroutinesApi
fun Route.gameWebSocket(
playerNotificationListener: PlayerNotificationListener,
commandHandler: GameCommandHandler,
) {
authenticate {
webSocket("/games/new") {
runWebSocket(GameId(), commandHandler, playerNotificationListener)
}
webSocket("/games/{id}") {
val gameId = GameId(UUID.fromString(call.parameters["id"]!!))
runWebSocket(gameId, commandHandler, playerNotificationListener)
}
}
}
@DelicateCoroutinesApi
private fun DefaultWebSocketServerSession.runWebSocket(
gameId: GameId,
commandHandler: GameCommandHandler,
playerNotificationListener: PlayerNotificationListener,
) {
val currentPlayer = call.getPlayerCredentials()
val incomingFrameChannel: ReceiveChannel<GameCommand> = toObjectChannel(incoming)
val outgoingFrameChannel: SendChannel<Notification> = fromFrameChannel(outgoing)
withLoggingContext("currentPlayer" to currentPlayer.toString()) {
val notificationListener =
playerNotificationListener.startListening(
currentPlayer,
gameId,
) { outgoingFrameChannel.trySendBlocking(it) }
// TODO change GlobalScope
GlobalScope.launch {
commandHandler.handleIncomingPlayerCommands(
currentPlayer,
gameId,
incomingFrameChannel,
outgoingFrameChannel,
)
notificationListener.close()
}
}
}
@@ -1,14 +0,0 @@
package eventDemo.adapter.presenter.query
import eventDemo.domain.entity.Player
import io.ktor.server.application.ApplicationCall
import io.ktor.server.auth.jwt.JWTPrincipal
import io.ktor.server.auth.principal
internal fun ApplicationCall.getPlayerCredentials() =
principal<JWTPrincipal>()!!.run {
Player(
id = payload.getClaim("playerid").asString(),
name = payload.getClaim("username").asString(),
)
}
@@ -1,53 +0,0 @@
package eventDemo.adapter.presenter.query
import eventDemo.domain.entity.GameId
import eventDemo.domain.event.projection.GameStateRepository
import eventDemo.configuration.serializer.GameIdSerializer
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,
) {
@Serializable
@Resource("card/last")
class Card(
val game: Game,
)
@Serializable
@Resource("state")
class State(
val game: Game,
)
}
/**
* API routes to read the game state.
*/
fun Route.readTheGameState(gameStateRepository: GameStateRepository) {
authenticate {
// Read the last played card on the game.
get<Game.Card> { body ->
gameStateRepository
.get(body.game.id)
.cardOnCurrentStack
?.let { call.respond(it) }
?: call.response.status(HttpStatusCode.BadRequest)
}
// Read the last played card on the game.
get<Game.State> { body ->
val state = gameStateRepository.get(body.game.id)
call.respond(state)
}
}
}
@@ -0,0 +1,46 @@
package eventDemo.configuration
import io.ktor.server.config.ApplicationConfig
data class Configuration(
val redisUrl: String,
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(
redisUrl = getProperty("redis.url"),
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")
@@ -1,29 +0,0 @@
package eventDemo.configuration
import eventDemo.configuration.domain.configureGameListener
import eventDemo.configuration.ktor.configureHttpRouting
import eventDemo.configuration.ktor.configureKoin
import eventDemo.configuration.ktor.configureSecurity
import eventDemo.configuration.ktor.configureSerialization
import eventDemo.configuration.ktor.configureWebSockets
import eventDemo.configuration.route.declareHttpGameRoute
import eventDemo.configuration.route.declareWebSocketsGameRoute
import io.ktor.server.application.Application
import org.koin.ktor.ext.get
import org.koin.ktor.ext.getKoin
fun Application.configure() {
configureKoin()
configureSecurity()
configureSerialization()
configureWebSockets()
declareWebSocketsGameRoute(get(), get())
configureHttpRouting()
declareHttpGameRoute()
getKoin().configureGameListener()
}
@@ -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,55 @@
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 redis.clients.jedis.JedisPooled
import redis.clients.jedis.UnifiedJedis
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
// Redis (for Projections)
single {
JedisPooled(config.redisUrl)
} bind UnifiedJedis::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()
}
@@ -1,21 +0,0 @@
package eventDemo.configuration.domain
import eventDemo.adapter.infrastructure.event.projection.GameListRepositoryInRedis
import eventDemo.adapter.infrastructure.event.projection.GameStateRepositoryInRedis
import eventDemo.domain.command.GameCommandHandler
import eventDemo.domain.event.projection.projectionListener.ReactionListener
import org.koin.core.Koin
fun Koin.configureGameListener() {
get<GameCommandHandler>()
.subscribeToBus(get())
get<GameStateRepositoryInRedis>()
.subscribeToBus(get(), get())
get<GameListRepositoryInRedis>()
.subscribeToBus(get(), get())
get<ReactionListener>()
.subscribeToBus(get())
}
@@ -1,30 +0,0 @@
package eventDemo.configuration.injection
import org.koin.dsl.module
fun appKoinModule(config: Configuration) =
module {
configureDIBusiness()
configureDIInfrastructure(config)
configureDILibs()
configureDICommandActions()
}
data class Configuration(
val redisUrl: 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,
)
}
@@ -1,18 +0,0 @@
package eventDemo.configuration.injection
import eventDemo.domain.command.action.ICantPlay
import eventDemo.domain.command.action.IWantToJoinTheGame
import eventDemo.domain.command.action.IWantToPlayCard
import eventDemo.domain.command.action.IamReadyToPlay
import org.koin.core.module.Module
import org.koin.core.module.dsl.singleOf
/**
* Configure all actions
*/
fun Module.configureDICommandActions() {
singleOf(::IWantToPlayCard)
singleOf(::IamReadyToPlay)
singleOf(::IWantToJoinTheGame)
singleOf(::ICantPlay)
}
@@ -1,19 +0,0 @@
package eventDemo.configuration.injection
import eventDemo.domain.command.GameCommandActionRunner
import eventDemo.domain.command.GameCommandHandler
import eventDemo.domain.event.GameEventHandler
import eventDemo.domain.event.projection.projectionListener.PlayerNotificationListener
import eventDemo.domain.event.projection.projectionListener.ReactionListener
import org.koin.core.module.Module
import org.koin.core.module.dsl.singleOf
fun Module.configureDIBusiness() {
single {
GameCommandHandler(get(), get(), get(), get())
}
singleOf(::GameEventHandler)
singleOf(::GameCommandActionRunner)
singleOf(::PlayerNotificationListener)
singleOf(::ReactionListener)
}
@@ -1,73 +0,0 @@
package eventDemo.configuration.injection
import com.rabbitmq.client.ConnectionFactory
import com.zaxxer.hikari.HikariConfig
import com.zaxxer.hikari.HikariDataSource
import eventDemo.adapter.infrastructure.event.GameEventBusInRabbinMQ
import eventDemo.adapter.infrastructure.event.GameEventStoreInPostgresql
import eventDemo.adapter.infrastructure.event.projection.GameListRepositoryInRedis
import eventDemo.adapter.infrastructure.event.projection.GameProjectionBusInRabbitMQ
import eventDemo.adapter.infrastructure.event.projection.GameStateRepositoryInRedis
import eventDemo.domain.event.GameEventBus
import eventDemo.domain.event.GameEventStore
import eventDemo.domain.event.projection.GameListRepository
import eventDemo.domain.event.projection.GameProjectionBus
import eventDemo.domain.event.projection.GameStateRepository
import org.koin.core.module.Module
import org.koin.core.module.dsl.singleOf
import org.koin.core.scope.Scope
import org.koin.core.scope.ScopeCallback
import org.koin.dsl.bind
import redis.clients.jedis.JedisPooled
import redis.clients.jedis.UnifiedJedis
import javax.sql.DataSource
fun Module.configureDIInfrastructure(config: Configuration) {
// Postgresql config
single {
JedisPooled(config.redisUrl)
} bind UnifiedJedis::class
single {
HikariConfig()
.apply {
jdbcUrl = config.postgresql.url
username = config.postgresql.username
password = config.postgresql.password
maximumPoolSize = 10
minimumIdle = 10
}.let {
HikariDataSource(it)
}.also { datasource ->
registerCallback(
object : ScopeCallback {
override fun onScopeClose(scope: Scope) {
datasource.close()
}
},
)
}
} bind DataSource::class
// RabbitMQ config
factory {
ConnectionFactory().apply {
host = config.rabbitmq.url
port = config.rabbitmq.port
username = config.rabbitmq.username
password = config.rabbitmq.password
}
}
singleOf(::GameEventBusInRabbinMQ) bind GameEventBus::class
singleOf(::GameEventStoreInPostgresql) bind GameEventStore::class
singleOf(::GameProjectionBusInRabbitMQ) bind GameProjectionBus::class
single {
GameStateRepositoryInRedis(get())
} bind GameStateRepository::class
single {
GameListRepositoryInRedis(get())
} bind GameListRepository::class
}
@@ -1,11 +0,0 @@
package eventDemo.configuration.injection
import eventDemo.libs.event.VersionBuilder
import eventDemo.libs.event.VersionBuilderLocal
import org.koin.core.module.Module
import org.koin.core.module.dsl.singleOf
import org.koin.dsl.bind
fun Module.configureDILibs() {
singleOf(::VersionBuilderLocal) bind VersionBuilder::class
}
@@ -1,42 +0,0 @@
package eventDemo.configuration.ktor
import eventDemo.configuration.injection.Configuration
import eventDemo.configuration.injection.appKoinModule
import io.ktor.server.application.Application
import io.ktor.server.application.install
import io.ktor.server.config.ApplicationConfig
import org.koin.ktor.plugin.Koin
import org.koin.logger.slf4jLogger
fun Application.configureKoin() {
install(Koin) {
slf4jLogger()
modules(
appKoinModule(
environment.config.configuration(),
),
)
}
}
fun ApplicationConfig.configuration() =
Configuration(
redisUrl = getProperty("redis.url"),
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")
@@ -1,14 +0,0 @@
package eventDemo.configuration.route
import eventDemo.adapter.presenter.query.readGamesList
import eventDemo.adapter.presenter.query.readTheGameState
import io.ktor.server.application.Application
import io.ktor.server.routing.routing
import org.koin.ktor.ext.get
fun Application.declareHttpGameRoute() {
routing {
readTheGameState(this@declareHttpGameRoute.get())
readGamesList(this@declareHttpGameRoute.get())
}
}
@@ -1,18 +0,0 @@
package eventDemo.configuration.route
import eventDemo.adapter.presenter.query.gameWebSocket
import eventDemo.domain.command.GameCommandHandler
import eventDemo.domain.event.projection.projectionListener.PlayerNotificationListener
import io.ktor.server.application.Application
import io.ktor.server.routing.routing
import kotlinx.coroutines.DelicateCoroutinesApi
@OptIn(DelicateCoroutinesApi::class)
fun Application.declareWebSocketsGameRoute(
playerNotificationListener: PlayerNotificationListener,
commandHandler: GameCommandHandler,
) {
routing {
gameWebSocket(playerNotificationListener, commandHandler)
}
}
@@ -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.sharedKernel.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.sharedKernel.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.sharedKernel.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.sharedKernel.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.libs.eventSource.EventId
import eventDemo.sharedKernel.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.sharedKernel.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)
}
}
@@ -1,24 +1,21 @@
package eventDemo.configuration.ktor
package eventDemo.contexts.auth.infrastructure.configure
import com.auth0.jwt.JWT
import com.auth0.jwt.algorithms.Algorithm
import eventDemo.domain.entity.Player
import eventDemo.configuration.configuration
import eventDemo.contexts.auth.domain.User
import eventDemo.contexts.auth.infrastructure.persistence.projection.UserProjection
import eventDemo.sharedKernel.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 io.ktor.server.routing.post
import io.ktor.server.routing.routing
import kotlinx.serialization.json.Json
import java.util.Date
private const val JWT_ISSUER = "PlayCardGame"
fun Application.configureSecurity() {
val jwtSecret = environment.config.propertyOrNull("jwt.secret")?.getString() ?: error("You must set a jwt secret")
fun Application.configureKtorAuth() {
val jwtSecret = environment.config.configuration.jwtSecret
authentication {
jwt {
realm = "Play card game"
@@ -29,7 +26,11 @@ fun Application.configureSecurity() {
.build(),
)
validate { credential ->
if (credential.payload.getClaim("username").asString() != "") {
if (credential.payload
.getClaim("username")
.asString()
.isNotEmpty()
) {
JWTPrincipal(credential.payload)
} else {
null
@@ -40,22 +41,25 @@ fun Application.configureSecurity() {
}
}
}
routing {
post("login/{username}") {
val username = call.parameters["username"]!!
val player = Player(name = username)
call.respond(hashMapOf("token" to player.makeJwt(jwtSecret)))
}
}
}
fun Player.makeJwt(jwtSecret: String): String =
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", name)
.withPayload(Json.encodeToString(this))
.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.sharedKernel.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.sharedKernel.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.sharedKernel.UserId
data class UserProjection(
val id: UserId,
val username: String,
val password: String,
)
@@ -0,0 +1,61 @@
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.sharedKernel.UserId
import java.util.UUID
import javax.sql.DataSource
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.fromString(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)
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.command.models.GameCommand
import eventDemo.contexts.game.application.notification.CommandSubscriber
import eventDemo.contexts.game.application.notification.EventToNotificationSubscriber
import eventDemo.contexts.game.application.notification.models.Notification
import eventDemo.contexts.game.domain.game.GameId
import eventDemo.sharedKernel.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() }
}
}
@@ -1,4 +1,4 @@
package eventDemo.domain.command
package eventDemo.contexts.game.application.command.handlers
class CommandException(
override val message: String,
@@ -0,0 +1,31 @@
package eventDemo.contexts.game.application.command.handlers
import eventDemo.contexts.game.application.command.models.GameCommand
import eventDemo.contexts.game.application.command.models.JoinTheGameCommand
import eventDemo.contexts.game.application.command.models.PlayCardCommand
import eventDemo.contexts.game.application.command.models.ReadyToPlayCommand
import eventDemo.contexts.game.application.command.models.TakeCartFromDrawPileCommand
import eventDemo.contexts.game.domain.game.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.command.models.GameCommand
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.command.Command
import eventDemo.libs.eventSource.eventStore.VersionConflictException
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.command.models.JoinTheGameCommand
import eventDemo.contexts.game.application.eventStores.GameRepository
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.domain.game.gameState.GameCreated
/**
* 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.command.models.PlayCardCommand
import eventDemo.contexts.game.application.eventStores.GameRepository
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.domain.game.gameState.GameStarted
/**
* 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.command.models.ReadyToPlayCommand
import eventDemo.contexts.game.application.eventStores.GameRepository
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.domain.game.gameState.GameCreated
/**
* 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.command.models.TakeCartFromDrawPileCommand
import eventDemo.contexts.game.application.eventStores.GameRepository
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.domain.game.gameState.GameStarted
/**
* 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,19 @@
package eventDemo.contexts.game.application.command.models
import eventDemo.contexts.game.domain.game.GameId
import eventDemo.contexts.game.infrastructure.persistence.serializers.GameIdSerializer
import eventDemo.libs.command.Command
import eventDemo.sharedKernel.UserId
import kotlinx.serialization.Serializable
@Serializable
sealed interface GameCommand : Command {
val userId: UserId
val payload: Payload
@Serializable
sealed interface Payload {
@Serializable(with = GameIdSerializer::class)
val aggregateId: GameId
}
}
@@ -1,22 +1,24 @@
package eventDemo.domain.command.command
package eventDemo.contexts.game.application.command.models
import eventDemo.domain.entity.GameId
import eventDemo.domain.entity.Player
import eventDemo.contexts.game.domain.game.GameId
import eventDemo.contexts.game.infrastructure.persistence.serializers.GameIdSerializer
import eventDemo.libs.command.CommandId
import eventDemo.sharedKernel.UserId
import kotlinx.serialization.Serializable
/**
* A command to perform an action to play a new card
*/
@Serializable
data class IWantToJoinTheGameCommand(
data class JoinTheGameCommand(
override val userId: UserId,
override val payload: Payload,
) : GameCommand {
override val id: CommandId = CommandId()
@Serializable
data class Payload(
@Serializable(with = GameIdSerializer::class)
override val aggregateId: GameId,
override val player: Player,
) : GameCommand.Payload
}
@@ -0,0 +1,31 @@
package eventDemo.contexts.game.application.command.models
import eventDemo.contexts.game.domain.game.Card
import eventDemo.contexts.game.domain.game.GameId
import eventDemo.contexts.game.domain.game.Player
import eventDemo.contexts.game.infrastructure.persistence.serializers.GameIdSerializer
import eventDemo.contexts.game.infrastructure.persistence.serializers.PlayerIdSerializer
import eventDemo.libs.command.CommandId
import eventDemo.sharedKernel.UserId
import kotlinx.serialization.Serializable
/**
* A command to perform an action to play a new card
*/
@Serializable
data class PlayCardCommand(
override val userId: UserId,
override val payload: Payload,
) : GameCommand {
override val id: CommandId = CommandId()
@Serializable
data class Payload(
@Serializable(with = GameIdSerializer::class)
override val aggregateId: GameId,
@Serializable(with = PlayerIdSerializer::class)
val playerId: Player.PlayerId,
val card: Card,
val chosenColor: Card.Color?,
) : GameCommand.Payload
}
@@ -0,0 +1,28 @@
package eventDemo.contexts.game.application.command.models
import eventDemo.contexts.game.domain.game.GameId
import eventDemo.contexts.game.domain.game.Player
import eventDemo.contexts.game.infrastructure.persistence.serializers.GameIdSerializer
import eventDemo.contexts.game.infrastructure.persistence.serializers.PlayerIdSerializer
import eventDemo.libs.command.CommandId
import eventDemo.sharedKernel.UserId
import kotlinx.serialization.Serializable
/**
* A command to set as ready to play
*/
@Serializable
data class ReadyToPlayCommand(
override val userId: UserId,
override val payload: Payload,
) : GameCommand {
override val id: CommandId = CommandId()
@Serializable
data class Payload(
@Serializable(with = GameIdSerializer::class)
override val aggregateId: GameId,
@Serializable(with = PlayerIdSerializer::class)
val playerId: Player.PlayerId,
) : GameCommand.Payload
}
@@ -0,0 +1,28 @@
package eventDemo.contexts.game.application.command.models
import eventDemo.contexts.game.domain.game.GameId
import eventDemo.contexts.game.domain.game.Player
import eventDemo.contexts.game.infrastructure.persistence.serializers.GameIdSerializer
import eventDemo.contexts.game.infrastructure.persistence.serializers.PlayerIdSerializer
import eventDemo.libs.command.CommandId
import eventDemo.sharedKernel.UserId
import kotlinx.serialization.Serializable
/**
* A command to perform an action to play a new card
*/
@Serializable
data class TakeCartFromDrawPileCommand(
override val userId: UserId,
override val payload: Payload,
) : GameCommand {
override val id: CommandId = CommandId()
@Serializable
data class Payload(
@Serializable(with = GameIdSerializer::class)
override val aggregateId: GameId,
@Serializable(with = PlayerIdSerializer::class)
val playerId: Player.PlayerId,
) : GameCommand.Payload
}
@@ -0,0 +1,26 @@
package eventDemo.contexts.game.application.eventStores
import eventDemo.contexts.game.application.ports.GameEventStore
import eventDemo.contexts.game.domain.game.GameId
import eventDemo.contexts.game.domain.game.gameState.Game
import eventDemo.libs.eventSource.eventStore.VersionConflictException
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.GameId
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
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.application.notification.models.ItsTheTurnOfNotification
import eventDemo.contexts.game.application.notification.models.Notification
import eventDemo.contexts.game.application.notification.models.PilesShuffledNotification
import eventDemo.contexts.game.application.notification.models.PlayerAsJoinTheGameNotification
import eventDemo.contexts.game.application.notification.models.PlayerAsPlayACardNotification
import eventDemo.contexts.game.application.notification.models.PlayerHavePassNotification
import eventDemo.contexts.game.application.notification.models.PlayerWasReadyNotification
import eventDemo.contexts.game.application.notification.models.PlayerWinNotification
import eventDemo.contexts.game.application.notification.models.TheGameWasStartedNotification
import eventDemo.contexts.game.application.notification.models.WelcomeToTheGameNotification
import eventDemo.contexts.game.application.notification.models.YourNewCardNotification
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.sharedKernel.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.command.models.GameCommand
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.notification.models.Notification
import eventDemo.contexts.game.application.ports.GameEventBus
import eventDemo.contexts.game.domain.game.GameId
import eventDemo.libs.bus.Bus
import eventDemo.libs.command.CommandUnicityChecker
import eventDemo.sharedKernel.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)
}
}
}
}
}
@@ -1,7 +1,7 @@
package eventDemo.domain.notification
package eventDemo.contexts.game.application.notification.models
import eventDemo.domain.entity.Player
import eventDemo.configuration.serializer.UUIDSerializer
import eventDemo.contexts.game.domain.game.Player
import eventDemo.libs.serializer.UUIDSerializer
import kotlinx.serialization.Serializable
import java.util.UUID
@@ -1,6 +1,6 @@
package eventDemo.domain.notification
package eventDemo.contexts.game.application.notification.models
import eventDemo.configuration.serializer.UUIDSerializer
import eventDemo.libs.serializer.UUIDSerializer
import kotlinx.serialization.Serializable
import java.util.UUID
@@ -0,0 +1,11 @@
package eventDemo.contexts.game.application.notification.models
import eventDemo.libs.serializer.UUIDSerializer
import kotlinx.serialization.Serializable
import java.util.UUID
@Serializable
data class PilesShuffledNotification(
@Serializable(with = UUIDSerializer::class)
override val id: UUID = UUID.randomUUID(),
) : Notification
@@ -1,7 +1,7 @@
package eventDemo.domain.notification
package eventDemo.contexts.game.application.notification.models
import eventDemo.domain.entity.Player
import eventDemo.configuration.serializer.UUIDSerializer
import eventDemo.contexts.game.domain.game.Player
import eventDemo.libs.serializer.UUIDSerializer
import kotlinx.serialization.Serializable
import java.util.UUID
@@ -1,8 +1,8 @@
package eventDemo.domain.notification
package eventDemo.contexts.game.application.notification.models
import eventDemo.domain.entity.Card
import eventDemo.domain.entity.Player
import eventDemo.configuration.serializer.UUIDSerializer
import eventDemo.contexts.game.domain.game.Card
import eventDemo.contexts.game.domain.game.Player
import eventDemo.libs.serializer.UUIDSerializer
import kotlinx.serialization.Serializable
import java.util.UUID
@@ -10,6 +10,6 @@ import java.util.UUID
data class PlayerAsPlayACardNotification(
@Serializable(with = UUIDSerializer::class)
override val id: UUID = UUID.randomUUID(),
val player: Player,
val playerId: Player.PlayerId,
val card: Card,
) : Notification
@@ -1,7 +1,7 @@
package eventDemo.domain.notification
package eventDemo.contexts.game.application.notification.models
import eventDemo.domain.entity.Player
import eventDemo.configuration.serializer.UUIDSerializer
import eventDemo.contexts.game.domain.game.Player
import eventDemo.libs.serializer.UUIDSerializer
import kotlinx.serialization.Serializable
import java.util.UUID
@@ -9,5 +9,5 @@ import java.util.UUID
data class PlayerHavePassNotification(
@Serializable(with = UUIDSerializer::class)
override val id: UUID = UUID.randomUUID(),
val player: Player,
val playerId: Player.PlayerId,
) : Notification
@@ -1,7 +1,7 @@
package eventDemo.domain.notification
package eventDemo.contexts.game.application.notification.models
import eventDemo.domain.entity.Player
import eventDemo.configuration.serializer.UUIDSerializer
import eventDemo.contexts.game.domain.game.Player
import eventDemo.libs.serializer.UUIDSerializer
import kotlinx.serialization.Serializable
import java.util.UUID
@@ -9,5 +9,5 @@ import java.util.UUID
data class PlayerWasReadyNotification(
@Serializable(with = UUIDSerializer::class)
override val id: UUID = UUID.randomUUID(),
val player: Player,
val playerId: Player.PlayerId,
) : Notification
@@ -1,7 +1,7 @@
package eventDemo.domain.notification
package eventDemo.contexts.game.application.notification.models
import eventDemo.domain.entity.Player
import eventDemo.configuration.serializer.UUIDSerializer
import eventDemo.contexts.game.domain.game.Player
import eventDemo.libs.serializer.UUIDSerializer
import kotlinx.serialization.Serializable
import java.util.UUID
@@ -9,5 +9,5 @@ import java.util.UUID
data class PlayerWinNotification(
@Serializable(with = UUIDSerializer::class)
override val id: UUID = UUID.randomUUID(),
val player: Player,
val playerId: Player.PlayerId,
) : Notification
@@ -1,7 +1,7 @@
package eventDemo.domain.notification
package eventDemo.contexts.game.application.notification.models
import eventDemo.domain.entity.Card
import eventDemo.configuration.serializer.UUIDSerializer
import eventDemo.contexts.game.domain.game.Card
import eventDemo.libs.serializer.UUIDSerializer
import kotlinx.serialization.Serializable
import java.util.UUID
@@ -9,5 +9,5 @@ import java.util.UUID
data class TheGameWasStartedNotification(
@Serializable(with = UUIDSerializer::class)
override val id: UUID = UUID.randomUUID(),
val hand: List<Card>,
val hand: Set<Card>,
) : Notification
@@ -1,7 +1,7 @@
package eventDemo.domain.notification
package eventDemo.contexts.game.application.notification.models
import eventDemo.domain.entity.Player
import eventDemo.configuration.serializer.UUIDSerializer
import eventDemo.contexts.game.domain.game.Player
import eventDemo.libs.serializer.UUIDSerializer
import kotlinx.serialization.Serializable
import java.util.UUID
@@ -1,7 +1,7 @@
package eventDemo.domain.notification
package eventDemo.contexts.game.application.notification.models
import eventDemo.domain.entity.Card
import eventDemo.configuration.serializer.UUIDSerializer
import eventDemo.contexts.game.domain.game.Card
import eventDemo.libs.serializer.UUIDSerializer
import kotlinx.serialization.Serializable
import java.util.UUID
@@ -9,5 +9,5 @@ import java.util.UUID
data class YourNewCardNotification(
@Serializable(with = UUIDSerializer::class)
override val id: UUID = UUID.randomUUID(),
val card: Card,
val cards: Set<Card>,
) : Notification
@@ -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.contexts.game.domain.game.GameId
import eventDemo.libs.eventSource.eventStore.EventStore
interface GameEventStore : EventStore<GameEvent, GameId>

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