Initial commit

This commit is contained in:
Fokinator
2026-03-09 20:17:47 +03:00
commit b1cab9de46
44 changed files with 3176 additions and 0 deletions
+2
View File
@@ -0,0 +1,2 @@
/mvnw text eol=lf
*.cmd text eol=crlf
+34
View File
@@ -0,0 +1,34 @@
HELP.md
target/
!.mvn/wrapper/maven-wrapper.jar
!**/src/main/**/target/
!**/src/test/**/target/
### STS ###
.apt_generated
.classpath
.factorypath
.project
.settings
.springBeans
.sts4-cache
### IntelliJ IDEA ###
.idea
.gigaide
*.iws
*.iml
*.ipr
### NetBeans ###
/nbproject/private/
/nbbuild/
/dist/
/nbdist/
/.nb-gradle/
build/
!**/src/main/**/build/
!**/src/test/**/build/
### VS Code ###
.vscode/
+19
View File
@@ -0,0 +1,19 @@
# Licensed to the Apache Software Foundation (ASF) under one
# or more contributor license agreements. See the NOTICE file
# distributed with this work for additional information
# regarding copyright ownership. The ASF licenses this file
# to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance
# with the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing,
# software distributed under the License is distributed on an
# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.
wrapperVersion=3.3.2
distributionType=only-script
distributionUrl=https://repo.maven.apache.org/maven2/org/apache/maven/apache-maven/3.9.9/apache-maven-3.9.9-bin.zip
+144
View File
@@ -0,0 +1,144 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>3.4.5</version>
<relativePath/> <!-- lookup parent from repository -->
</parent>
<groupId>ru.fokinatorr</groupId>
<artifactId>chess.server</artifactId>
<version>0.0.1-SNAPSHOT</version>
<name>Ch4ss Server</name>
<description>Chess game server that syncs clients&apos; positions</description>
<url/>
<licenses>
<license/>
</licenses>
<developers>
<developer/>
</developers>
<scm>
<connection/>
<developerConnection/>
<tag/>
<url/>
</scm>
<properties>
<java.version>23</java.version>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-r2dbc</artifactId>
</dependency>
<dependency>
<groupId>org.postgresql</groupId>
<artifactId>r2dbc-postgresql</artifactId>
<!--scope>runtime</scope-->
</dependency>
<dependency>
<groupId>org.postgresql</groupId>
<artifactId>postgresql</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>ru.fokinatorr</groupId>
<artifactId>chess.api</artifactId>
<version>0.0.2-SNAPSHOT</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-rsocket</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webflux</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-devtools</artifactId>
<scope>runtime</scope>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.springframework.security</groupId>
<artifactId>spring-security-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-security</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.security</groupId>
<artifactId>spring-security-rsocket</artifactId>
</dependency>
<dependency>
<groupId>org.jetbrains</groupId>
<artifactId>annotations</artifactId>
<version>26.0.1</version>
<scope>compile</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<annotationProcessorPaths>
<path>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
</path>
</annotationProcessorPaths>
</configuration>
</plugin>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
<configuration>
<excludes>
<exclude>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
</exclude>
</excludes>
</configuration>
</plugin>
<plugin>
<groupId>org.liquibase</groupId>
<artifactId>liquibase-maven-plugin</artifactId>
<version>4.29.2</version>
<configuration>
<changeLogFile>src/main/resources/liquibase/master.xml</changeLogFile>
<url>jdbc:postgresql://localhost:5432/chessserverdb</url>
<username>postgres</username>
<password>3h6F$t3_6e&amp;`0L//</password>
<driver>org.postgresql.Driver</driver>
</configuration>
</plugin>
</plugins>
</build>
</project>
@@ -0,0 +1,174 @@
package ru.fokinatorr.chess.server;
import lombok.NonNull;
import lombok.experimental.UtilityClass;
import org.jetbrains.annotations.NotNull;
import java.util.Objects;
import java.util.function.BiPredicate;
import java.util.function.BinaryOperator;
import java.util.function.Predicate;
import java.util.function.UnaryOperator;
/// Provides `boolean` operators as `Function`s.
@UtilityClass
public class BooleanOperators {
/// Logical AND operator.
public static final BoolBinaryOperator AND = (b1, b2) -> b1 && b2;
/// Logical OR operator.
public static final BoolBinaryOperator OR = (b1, b2) -> b1 || b2;
/// Logical XOR operator.
public static final BoolBinaryOperator XOR = (b1, b2) -> b1 ^ b2;
/// Reverse of logical AND.
public static final BoolBinaryOperator NAND = AND.negate();
/// Reverse of logical OR.
public static final BoolBinaryOperator NOR = OR.negate();
/// Reverse of logical XOR.
public static final BoolBinaryOperator XNOR = (b1, b2) -> b1 == b2;
/// Logical NOT operator.
public static final BoolUnaryOperator NOT = b -> !b;
/// Identity function.
public static final BoolUnaryOperator IDENTITY = b -> b;
/// A combination of `BiPredicate<Boolean, Boolean>` and `BinaryOperator<Boolean>`.
/// <p>This is a functional interface whose functional method is {@link #test(Boolean, Boolean)}. {@link #apply(Boolean, Boolean)} does the same.
@FunctionalInterface
public interface BoolBinaryOperator extends BiPredicate<Boolean, Boolean>, BinaryOperator<Boolean> {
/// {@inheritDoc}
@Override
default Boolean apply(Boolean b1, Boolean b2) {
return test(b1, b2);
}
/// {@inheritDoc}
@Override
boolean test(Boolean b1, Boolean b2);
/// {@inheritDoc}
@NonNull
default BooleanOperators.BoolBinaryOperator negate() {
return (b1, b2) -> !apply(b1, b2);
}
/// Returns a composed operator that represents a short-circuiting logical AND of this operator and another.
default BoolBinaryOperator and(BoolBinaryOperator operator) {
Objects.requireNonNull(operator);
return (b1, b2) -> apply(b1, b2) && operator.apply(b1, b2);
}
/// Returns a composed operator that represents a short-circuiting logical OR of this operator and another.
default BoolBinaryOperator or(BoolBinaryOperator operator) {
Objects.requireNonNull(operator);
return (b1, b2) -> apply(b1, b2) || operator.apply(b1, b2);
}
/// Returns a composed operator that represents a logical XOR of this operator and another.
default BoolBinaryOperator xor(BoolBinaryOperator operator) {
Objects.requireNonNull(operator);
return (b1, b2) -> apply(b1, b2) ^ operator.apply(b1, b2);
}
/// Returns a composed operator that represents a short-circuiting logical NAND of this operator and another.
default BoolBinaryOperator nand(BoolBinaryOperator operator) {
Objects.requireNonNull(operator);
return (b1, b2) -> !(apply(b1, b2) && operator.apply(b1, b2));
}
/// Returns a composed operator that represents a short-circuiting logical NOR of this operator and another.
default BoolBinaryOperator nor(BoolBinaryOperator operator) {
Objects.requireNonNull(operator);
return (b1, b2) -> !(apply(b1, b2) || operator.apply(b1, b2));
}
/// Returns a composed operator that represents a logical XNOR of this operator and another.
default BoolBinaryOperator xnor(BoolBinaryOperator operator) {
Objects.requireNonNull(operator);
return (b1, b2) -> apply(b1, b2) == operator.apply(b1, b2);
}
/// Returns a `BiPredicateOperator` delegating to the given predicate.
static BoolBinaryOperator valueOf(BiPredicate<Boolean, Boolean> predicate) {
Objects.requireNonNull(predicate);
return predicate instanceof BoolBinaryOperator boolBinaryOperator ? boolBinaryOperator : predicate::test;
}
/// Returns a `BiPredicateOperator` delegating to the given binary operator.
static BoolBinaryOperator valueOf(BinaryOperator<Boolean> binaryOperator) {
Objects.requireNonNull(binaryOperator);
return binaryOperator instanceof BoolBinaryOperator boolBinaryOperator ? boolBinaryOperator : binaryOperator::apply;
}
}
/// A combination of `Predicate<Boolean>` and `UnaryOperator<Boolean>`.
/// <p>This is a functional interface whose functional method is {@link #test(Boolean)}. {@link #apply(Boolean)} does the same.
@FunctionalInterface
public interface BoolUnaryOperator extends Predicate<Boolean>, UnaryOperator<Boolean> {
/// {@inheritDoc}
@Override
default Boolean apply(Boolean b) {
return test(b);
}
/// {@inheritDoc}
@Override
boolean test(Boolean aBoolean);
/// {@inheritDoc}
@Override
@NotNull
default BooleanOperators.BoolUnaryOperator negate() {
return b -> !apply(b);
}
/// Returns a composed operator that represents a short-circuiting logical AND of this operator and another.
default BoolUnaryOperator and(BoolUnaryOperator operator) {
Objects.requireNonNull(operator);
return b -> apply(b) && operator.apply(b);
}
/// Returns a composed operator that represents a short-circuiting logical OR of this operator and another.
default BoolUnaryOperator or(BoolUnaryOperator operator) {
Objects.requireNonNull(operator);
return b -> apply(b) || operator.apply(b);
}
/// Returns a composed operator that represents a logical XOR of this operator and another.
default BoolUnaryOperator xor(BoolUnaryOperator operator) {
Objects.requireNonNull(operator);
return b1 -> apply(b1) ^ operator.apply(b1);
}
/// Returns a composed operator that represents a short-circuiting logical NAND of this operator and another.
default BoolUnaryOperator nand(BoolUnaryOperator operator) {
Objects.requireNonNull(operator);
return (b1) -> !(apply(b1) && operator.apply(b1));
}
/// Returns a composed operator that represents a short-circuiting logical NOR of this operator and another.
default BoolUnaryOperator nor(BoolUnaryOperator operator) {
Objects.requireNonNull(operator);
return (b1) -> !(apply(b1) || operator.apply(b1));
}
/// Returns a composed operator that represents a logical XNOR of this operator and another.
default BoolUnaryOperator xnor(BoolUnaryOperator operator) {
Objects.requireNonNull(operator);
return (b1) -> apply(b1) == operator.apply(b1);
}
/// Returns a `PredicateOperator` delegating to the given predicate.
static BoolUnaryOperator valueOf(Predicate<Boolean> predicate) {
Objects.requireNonNull(predicate);
return predicate instanceof BoolUnaryOperator boolUnaryOperator ? boolUnaryOperator : predicate::test;
}
/// Returns a `PredicateOperator` delegating to the given unary operator.
static BoolUnaryOperator valueOf(UnaryOperator<Boolean> unaryOperator) {
Objects.requireNonNull(unaryOperator);
return unaryOperator instanceof BoolUnaryOperator boolUnaryOperator ? boolUnaryOperator : unaryOperator::apply;
}
}
}
@@ -0,0 +1,15 @@
package ru.fokinatorr.chess.server;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
/// The entrypoint class.
@SpringBootApplication
public class Ch4ssServerApplication {
/// The main method, executed from JNI
public static void main(String[] args) {
SpringApplication.run(Ch4ssServerApplication.class, args);
}
}
@@ -0,0 +1,29 @@
package ru.fokinatorr.chess.server;
import lombok.AccessLevel;
import lombok.RequiredArgsConstructor;
import lombok.experimental.FieldDefaults;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.stereotype.Component;
import ru.fokinatorr.chess.server.data.UserService;
/// Default `ApplicationRunner`.
@Component
@RequiredArgsConstructor
@FieldDefaults(level = AccessLevel.PRIVATE, makeFinal = true)
@Slf4j
public class DefaultApplicationRunner implements ApplicationRunner {
UserService userService;
@Override
public void run(ApplicationArguments args) throws Exception {
userService.findAll()
.collectList()
.subscribe(users -> log.info("Users: {}", users));
}
}
@@ -0,0 +1,30 @@
package ru.fokinatorr.chess.server.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.convert.CustomConversions;
import org.springframework.data.r2dbc.convert.R2dbcCustomConversions;
import org.springframework.data.r2dbc.dialect.PostgresDialect;
import org.springframework.data.r2dbc.dialect.R2dbcDialect;
import ru.fokinatorr.chess.server.data.converter.JsonNodeToJsonConverter;
import ru.fokinatorr.chess.server.data.converter.JsonToJsonNodeConverter;
/// DB-related configuration.
@Configuration
public class DataConfig {
@Bean
public CustomConversions customConversions(R2dbcDialect sqlDialect,
JsonNodeToJsonConverter jsonNodeToJsonConverter,
JsonToJsonNodeConverter jsonToJsonNodeConverter) {
return R2dbcCustomConversions.of(sqlDialect,
jsonNodeToJsonConverter, jsonToJsonNodeConverter
);
}
@Bean
public R2dbcDialect sqlDialect() {
return PostgresDialect.INSTANCE;
}
}
@@ -0,0 +1,16 @@
package ru.fokinatorr.chess.server.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.Random;
@Configuration
public class MiscBeansConfig {
@Bean
public Random randomGenerator() {
return new Random();
}
}
@@ -0,0 +1,46 @@
package ru.fokinatorr.chess.server.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.http.codec.cbor.Jackson2CborDecoder;
import org.springframework.http.codec.cbor.Jackson2CborEncoder;
import org.springframework.messaging.rsocket.RSocketStrategies;
import org.springframework.messaging.rsocket.annotation.support.RSocketMessageHandler;
import org.springframework.security.authentication.ReactiveAuthenticationManager;
import org.springframework.security.config.Customizer;
import org.springframework.security.config.annotation.rsocket.EnableRSocketSecurity;
import org.springframework.security.config.annotation.rsocket.RSocketSecurity;
import org.springframework.security.rsocket.core.PayloadSocketAcceptorInterceptor;
import org.springframework.web.util.pattern.PathPatternRouteMatcher;
/// Configures RSocket.
@Configuration
@EnableRSocketSecurity
public class RSocketConfig {
@Bean
public RSocketStrategies rsocketStrategies() {
return RSocketStrategies.builder()
.encoder(new Jackson2CborEncoder())
.decoder(new Jackson2CborDecoder())
.routeMatcher(new PathPatternRouteMatcher())
.build();
}
@Bean
public RSocketMessageHandler rsocketMessageHandler(RSocketStrategies rsocketStrategies) {
RSocketMessageHandler handler = new RSocketMessageHandler();
handler.setRSocketStrategies(rsocketStrategies);
return handler;
}
@Bean
public PayloadSocketAcceptorInterceptor rsocketInterceptor(RSocketSecurity rsocket) {
return rsocket
.authorizePayload(authorizePayload -> authorizePayload
.route("game.*").authenticated()
.anyExchange().permitAll())
.simpleAuthentication(Customizer.withDefaults())
.build();
}
}
@@ -0,0 +1,87 @@
package ru.fokinatorr.chess.server.config;
import lombok.extern.slf4j.Slf4j;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.http.HttpStatus;
import org.springframework.http.HttpStatusCode;
import org.springframework.http.MediaType;
import org.springframework.http.server.reactive.ServerHttpResponse;
import org.springframework.security.config.annotation.web.reactive.EnableWebFluxSecurity;
import org.springframework.security.config.web.server.ServerHttpSecurity;
import org.springframework.security.core.userdetails.ReactiveUserDetailsService;
import org.springframework.security.core.userdetails.UserDetails;
import org.springframework.security.core.userdetails.UsernameNotFoundException;
import org.springframework.security.crypto.bcrypt.BCryptPasswordEncoder;
import org.springframework.security.crypto.password.PasswordEncoder;
import org.springframework.security.web.server.SecurityWebFilterChain;
import reactor.core.publisher.Mono;
import ru.fokinatorr.chess.server.data.UserService;
import java.nio.charset.StandardCharsets;
/// Spring Security configuration.
@Configuration
@EnableWebFluxSecurity
@Slf4j
public class SecurityConfig {
@Bean
public PasswordEncoder passwordEncoder() {
return new BCryptPasswordEncoder();
}
@Bean
public SecurityWebFilterChain securityFilterChain(ServerHttpSecurity http) {
return http
.httpBasic(httpBasic -> httpBasic
.authenticationFailureHandler((webFilterExchange, exception) ->
statusCodeContentTypeContent(
webFilterExchange.getExchange().getResponse(),
HttpStatus.BAD_REQUEST,
MediaType.TEXT_PLAIN,
exception.toString()
))
)
.exceptionHandling(exceptionHandling -> exceptionHandling
.authenticationEntryPoint((exchange, ex) -> statusCodeContentTypeContent(
exchange.getResponse(),
HttpStatus.UNAUTHORIZED,
MediaType.TEXT_PLAIN,
ex.toString()
))
.accessDeniedHandler((exchange, ex) -> statusCodeContentTypeContent(
exchange.getResponse(),
HttpStatus.FORBIDDEN,
MediaType.TEXT_PLAIN,
ex.toString()
)))
.authorizeExchange(authorizeExchange -> authorizeExchange
.pathMatchers("/game/**").authenticated()
.anyExchange().permitAll())
// disable everything browser-only related
.formLogin(ServerHttpSecurity.FormLoginSpec::disable)
.logout(ServerHttpSecurity.LogoutSpec::disable)
.csrf(ServerHttpSecurity.CsrfSpec::disable)
.build();
}
@Bean
public ReactiveUserDetailsService userDetailsService(UserService userService) {
return username -> userService.findByUsername(username)
.switchIfEmpty(Mono.error(() -> new UsernameNotFoundException(username)))
.cast(UserDetails.class);
}
private static Mono<Void> statusCodeContentTypeContent(@NotNull ServerHttpResponse response, @NotNull HttpStatusCode statusCode, @Nullable MediaType contentType, @Nullable String body) {
response.setStatusCode(statusCode);
if (contentType == null || body == null) {
return response.setComplete();
} else {
response.getHeaders().setContentType(contentType);
return response.writeWith(Mono.just(response.bufferFactory().wrap(body.getBytes(StandardCharsets.UTF_8))));
}
}
}
@@ -0,0 +1,278 @@
package ru.fokinatorr.chess.server.controller.http;
import com.fasterxml.jackson.databind.JsonNode;
import lombok.AccessLevel;
import lombok.RequiredArgsConstructor;
import lombok.experimental.FieldDefaults;
import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
import org.springframework.security.core.annotation.AuthenticationPrincipal;
import org.springframework.web.bind.annotation.*;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import ru.fokinatorr.chess.impl.PieceColor;
import ru.fokinatorr.chess.server.entity.User;
import ru.fokinatorr.chess.server.exception.GameNotFoundException;
import ru.fokinatorr.chess.server.exception.NotParticipantException;
import ru.fokinatorr.chess.server.game.*;
import ru.fokinatorr.chess.server.misc.DTOServiceExtension;
import ru.fokinatorr.chess.server.misc.RecordAttributesService;
import java.util.UUID;
/// REST controller for handling game-related HTTP operations.
///
/// This controller provides a comprehensive set of endpoints for managing
/// chess games through HTTP requests. It handles game creation, joining,
/// cancellation, and various game actions like making moves and managing
/// draw offers. The controller is organized with nested inner classes
/// to group related endpoints by functionality.
///
/// The main functionality is divided into several sections:
/// - Game lifecycle management (`/game/create`, `/game/join`, `/game/cancel`)
/// - Game listing (`/game/list/pending`, `/game/list/ongoing`, `/game/list/finished`)
/// - Move execution (`/game/move/perform`)
/// - Draw offer management (`/game/draw/offer`, `/game/draw/refuse`, `/game/draw/accept`)
/// - Game surrender (`/game/giveup`)
///
/// The controller uses reactive programming with [Mono] and [Flux]
/// for non-blocking operations and integrates with Spring Security for
/// authentication via the `@`[AuthenticationPrincipal] annotation.
///
/// @see GameService
/// @see CreateGameOptions
/// @see PendingGameDTO
/// @see OngoingGameDTO
/// @see FinishedGameDTO
@RestController
@FieldDefaults(makeFinal = true, level = AccessLevel.PRIVATE)
@RequiredArgsConstructor
@RequestMapping("/game")
public class GameControllerRest {
/// The main game service that handles game logic and state management
GameService gameService;
/// Service for converting record classes to JSON attributes
RecordAttributesService recordAttributesService;
/// Service extension for converting game entities to DTOs
DTOServiceExtension dtoService;
/// Retrieves the available options for creating a new game.
///
/// This endpoint returns the attributes and constraints for the [CreateGameOptions]
/// record class, which clients can use to understand what parameters are available
/// when creating a new game (such as time controls, game variants, etc.).
///
/// @return a {@link Mono} containing JsonNode with the attributes of `CreateGameOptions`
@GetMapping("/create")
public Mono<JsonNode> createGameOptions() {
return Mono.fromCallable(() -> recordAttributesService.toAttributes(CreateGameOptions.class));
}
/// Creates a new chess game with the specified options.
///
/// This endpoint creates a new game with the given configuration and associates
/// it with the authenticated user (principal). The game is initially in a pending
/// state until another player joins.
///
/// @param principal the authenticated user creating the game
/// @param createGameOptions the configuration options for the new game
/// @return a Mono containing ResponseEntity with the created game's UUID and CREATED status
/// @see CreateGameOptions for available game configuration options
@PostMapping(value = "/create", consumes = MediaType.APPLICATION_JSON_VALUE)
public Mono<ResponseEntity<UUID>> createGame(@AuthenticationPrincipal(errorOnInvalidType = true) User principal, @RequestBody CreateGameOptions createGameOptions) {
return gameService.createGame(principal.getId(), createGameOptions)
.map(gameId -> ResponseEntity.status(HttpStatus.CREATED).body(gameId));
}
/// Joins an existing pending game.
///
/// This endpoint allows a user to join a game that was created by another player
/// and is currently in the pending state. Once joined, the game transitions to
/// ongoing state and becomes playable.
///
/// @param id the UUID of the game to join
/// @param principal the authenticated user joining the game
/// @return a {@link Mono} that signals when the join operation is successful
@PostMapping("/join")
@ResponseStatus(HttpStatus.NO_CONTENT)
public Mono<Void> joinGame(@RequestParam("id") UUID id, @AuthenticationPrincipal(errorOnInvalidType = true) User principal) {
return gameService.joinGame(id, principal.getId());
}
/// Cancels a game, either pending or ongoing.
///
/// This endpoint allows a user to cancel a game they created. It first attempts
/// to cancel a pending game, and if that fails with [GameNotFoundException],
/// it tries to cancel an ongoing game. This provides a unified interface for
/// game cancellation regardless of the game's current state.
///
/// @return a {@link Mono} that signals when the cancellation is successful
/// @throws GameNotFoundException if the game doesn't exist
@DeleteMapping("/cancel")
@ResponseStatus(HttpStatus.NO_CONTENT)
public Mono<Void> cancelGame(@RequestParam("id") UUID id, @AuthenticationPrincipal(errorOnInvalidType = true) User principal) {
return gameService.cancelPendingGame(id, principal.getId())
.onErrorResume(GameNotFoundException.class, _ -> gameService.cancelOngoingGame(id, principal.getId()));
}
/// Subcontroller for listing games in different states.
///
/// Provides endpoints to retrieve lists of games
/// based on their current state: pending, ongoing, or finished. Each endpoint
/// returns a stream of DTOs containing relevant game information for the
/// respective game state.
///
/// The full paths are `/game/list/pending`, `/game/list/ongoing`,
/// and `/game/list/finished`.
///
/// @see PendingGameDTO
/// @see OngoingGameDTO
/// @see FinishedGameDTO
@RestController
@RequestMapping("/game/list")
public class ListGamesController {
/// Retrieves a list of all pending games that are waiting for an opponent.
///
/// Pending games are those that have been created but not yet joined by
/// a second player. Clients can use this endpoint to find available games
/// to join.
///
/// @return a {@link Flux} stream of {@link PendingGameDTO} objects containing information about pending games
@GetMapping("/pending")
public Flux<PendingGameDTO> listPendingGames() {
return gameService.getPendingGames().map(dtoService::toDTO);
}
/// Retrieves a list of all currently ongoing games.
///
/// Ongoing games are those that have two players and are actively being played.
/// This endpoint can be used to monitor active games or for spectator functionality.
///
/// @return a {@link Flux} stream of {@link OngoingGameDTO} objects containing information about ongoing games
@GetMapping("/ongoing")
public Flux<OngoingGameDTO> listOngoingGames() {
return gameService.getOngoingGames().map(dtoService::toDTO);
}
/// Retrieves a list of all finished games.
///
/// Finished games are those that have reached a terminal state (checkmate,
/// stalemate, draw, or resignation). This endpoint provides access to game
/// history and results.
///
/// @return a {@link Flux} stream of {@link FinishedGameDTO} objects containing information about finished games
@GetMapping("/finished")
public Flux<FinishedGameDTO> listFinishedGames() {
return gameService.getFinishedGames().map(dtoService::toDTO);
}
}
/// Subcontroller for handling move-related operations.
///
/// Provides endpoints to execute moves in ongoing games.
/// It handles the validation and processing of chess moves in UCI (Universal
/// Chess Interface) format.
///
/// The full paths are `/game/move/perform`. Undoing and redoing to be added.
///
/// @see <a href="https://en.wikipedia.org/wiki/Universal_Chess_Interface">Universal Chess Interface (UCI)</a> for move format specification
@RestController
@RequestMapping("/game/move")
public class MovesController {
/// Executes a chess move in the specified game.
///
/// This endpoint processes a move in UCI format (e.g., "e2e4", "g1f3") for
/// the authenticated player in the specified game. The move is validated
/// against the current game state and executed if legal.
///
/// @return a {@link Mono} that signals when the move is successfully processed
/// @throws IllegalArgumentException if the move is invalid or illegal in the current position
@PutMapping("/perform")
@ResponseStatus(HttpStatus.NO_CONTENT)
public Mono<Void> performMove(@RequestParam("gameId") UUID gameId,
@RequestParam("move") String uciMove,
@AuthenticationPrincipal(errorOnInvalidType = true) User principal) {
return gameService.performMove(gameId, principal.getId(), Mono.just(uciMove));
}
}
/// Subcontroller for managing draw offers in games.
///
/// Provides endpoints to offer, accept, and refuse
/// draw offers in ongoing games. It implements the complete workflow for
/// draw negotiation between players.
///
/// The full paths are `/game/draw/offer`, `/game/draw/accept`,
/// and `/game/draw/refuse`.
///
/// The draw negotiation flow is:
/// 1. One player offers a draw using the offer endpoint
/// 2. The opponent can either accept (ending the game in a draw) or refuse
/// 3. If refused, the game continues normally
@RestController
@RequestMapping("/game/draw")
public class DrawController {
/// Offers a draw to the opponent in the specified game.
///
/// This endpoint allows a player to offer a draw to their opponent in an
/// ongoing game. Once offered, the opponent can choose to accept or refuse
/// the draw offer through the corresponding endpoints.
///
/// @return a {@link Mono} that signals when the draw offer is successfully registered
@PutMapping("/offer")
@ResponseStatus(HttpStatus.NO_CONTENT)
public Mono<Void> offerDraw(@RequestParam("gameId") UUID gameId, @AuthenticationPrincipal(errorOnInvalidType = true) User principal) {
return gameService.offerDraw(gameId, principal.getId());
}
/// Accepts a draw offer in the specified game.
///
/// This endpoint allows a player to accept a draw offer that was made by
/// their opponent, resulting in the game ending in a draw. The game state
/// is updated accordingly.
///
/// @return a {@link Mono} that signals when the draw is accepted and the game is finalized
@PutMapping("/accept")
@ResponseStatus(HttpStatus.NO_CONTENT)
public Mono<Void> acceptDraw(@RequestParam("gameId") UUID gameId, @AuthenticationPrincipal(errorOnInvalidType = true) User principal) {
return gameService.acceptDraw(gameId, principal.getId());
}
/// Refuses a draw offer in the specified game.
///
/// This endpoint allows a player to refuse a draw offer made by their
/// opponent. After refusal, the game continues normally and the opponent
/// may offer a draw again later if desired.
///
/// @return a {@link Mono} that signals when the draw offer is successfully refused
@PutMapping("/refuse")
@ResponseStatus(HttpStatus.NO_CONTENT)
public Mono<Void> refuseDraw(@RequestParam("gameId") UUID gameId, @AuthenticationPrincipal(errorOnInvalidType = true) User principal) {
return gameService.refuseDraw(gameId, principal.getId());
}
}
/// Allows a player to resign from a game.
///
/// This endpoint enables a player to surrender the game, resulting in a
/// loss for the resigning player and a win for their opponent. The game
/// is immediately finalized with the appropriate result.
///
/// @return a {@link Mono} that signals when the resignation is processed and the game is finalized
@PutMapping("/giveup")
@ResponseStatus(HttpStatus.NO_CONTENT)
public Mono<Void> giveUp(@RequestParam("gameId") UUID gameId, @AuthenticationPrincipal(errorOnInvalidType = true) User principal) {
return gameService.giveUp(gameId, principal.getId());
}
}
@@ -0,0 +1,68 @@
package ru.fokinatorr.chess.server.controller.http;
import com.fasterxml.jackson.databind.JsonNode;
import lombok.AccessLevel;
import lombok.RequiredArgsConstructor;
import lombok.experimental.FieldDefaults;
import lombok.extern.slf4j.Slf4j;
import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.server.ResponseStatusException;
import reactor.core.publisher.Mono;
import ru.fokinatorr.chess.server.BooleanOperators;
import ru.fokinatorr.chess.server.data.UserService;
import ru.fokinatorr.chess.server.entity.User;
import ru.fokinatorr.chess.server.misc.RecordAttributesService;
import ru.fokinatorr.chess.server.security.LoginRequirements;
import ru.fokinatorr.chess.server.security.RegisterForm;
import java.util.Set;
import java.util.function.Predicate;
/// Controller for login-related URLs.
@Slf4j
@RestController
@FieldDefaults(makeFinal = true, level = AccessLevel.PRIVATE)
@RequiredArgsConstructor
public class LoginController {
RecordAttributesService recordAttributesService;
UserService userService;
@GetMapping(value = "/login", produces = MediaType.APPLICATION_JSON_VALUE)
public Mono<LoginRequirements> login() {
return Mono.just(LoginRequirements.builder()
.requiredFields(Set.of("username", "password"))
.build());
}
@GetMapping(value = "/register", produces = MediaType.APPLICATION_JSON_VALUE)
public Mono<JsonNode> registerForm() {
return Mono.fromCallable(() -> recordAttributesService.toAttributes(RegisterForm.class));
}
@PostMapping(value = "/register", consumes = MediaType.APPLICATION_JSON_VALUE)
public Mono<ResponseEntity<Long>> register(@RequestBody RegisterForm registerForm) {
return userService.existsByUsername(registerForm.username())
.filter(BooleanOperators.NOT)
.switchIfEmpty(badRequest("User with username " + registerForm.username() + " already exists."))
.filter(_ -> registerForm.password().equals(registerForm.confirmPassword()))
.switchIfEmpty(badRequest("Password and confirmation password don't match."))
.flatMap(_ -> userService.create(registerForm.username(), registerForm.password(), false))
.flatMap(userService::save)
.map(User::getId)
.map(id -> ResponseEntity.status(HttpStatus.CREATED).body(id));
}
private static <T> Mono<T> badRequest(String reason) {
return Mono.defer(() -> {
log.info("Bad request: {}", reason);
return Mono.error(() -> new ResponseStatusException(HttpStatus.BAD_REQUEST, reason));
});
}
}
@@ -0,0 +1,8 @@
package ru.fokinatorr.chess.server.controller.http;
import org.springframework.web.bind.annotation.RestController;
/// Controller related to notifications.
@RestController
public class NotificationController { // ToDo
}
@@ -0,0 +1,69 @@
package ru.fokinatorr.chess.server.controller.rsocket;
import lombok.AccessLevel;
import lombok.RequiredArgsConstructor;
import lombok.experimental.FieldDefaults;
import org.springframework.messaging.handler.annotation.MessageMapping;
import org.springframework.messaging.handler.annotation.Payload;
import org.springframework.stereotype.Controller;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import ru.fokinatorr.chess.server.game.event.GameJoinedEvent;
import ru.fokinatorr.chess.server.game.event.GameEvent;
import ru.fokinatorr.chess.server.game.GameService;
import java.util.UUID;
/// Game-related controller for handling RSocket connections.
///
/// This controller provides endpoints for clients to subscribe to game events
/// and wait for game join notifications via RSocket messaging. It acts as the
/// server-side router for real-time game event streaming.
///
/// The controller handles two main operations:
/// - `subscribeToGame` (route: `"game.join"`): Allows clients to subscribe to a stream of game events
/// for a specific game, enabling real-time updates about moves, game state
/// changes, and other game-related events.
/// - `waitForJoin` (route: `"game.waitForJoin"`): Allows clients to wait for a notification when a game
/// is joined by another player, useful for matchmaking scenarios.
///
/// @see GameService GameService
/// @see GameEvent GameEvent
/// @see GameJoinedEvent GameJoinedEvent
@Controller
@MessageMapping("game")
@FieldDefaults(makeFinal = true, level = AccessLevel.PRIVATE)
@RequiredArgsConstructor
public class GameControllerRSocket {
/// The game service that manages game state and events.
GameService gameService;
/// Subscribes a client to receive a stream of game events for a specific game.
///
/// This method establishes a reactive stream ([Flux]) that continuously sends
/// game events to the client as they occur in the specified game. The client
/// will receive events like moves, game finish, draw offers, etc.
///
/// @param gameId the UUID of the game to subscribe to
/// @return a `Flux` stream of {@link GameEvent} objects that emits events as they occur
@MessageMapping("join")
public Flux<GameEvent> subscribeToGame(@Payload Mono<UUID> gameId) {
return gameId.flatMapMany(gameService::subscribeToGame);
}
/// Waits for a notification when a game is joined by another player.
///
/// This method returns a [Mono] that completes when another player joins
/// the specified game. It's typically used in matchmaking scenarios where
/// a player creates a game and waits for an opponent to join.
///
/// @param gameId the UUID of the game to wait for
/// @return a `Mono` that emits a {@link GameJoinedEvent} when another player joins
/// @see GameJoinedEvent for the event structure containing join information
@MessageMapping("waitForJoin")
public Mono<GameJoinedEvent> waitForJoin(@Payload Mono<UUID> gameId) {
return gameId.flatMap(gameService::waitForJoin);
}
}
@@ -0,0 +1,12 @@
package ru.fokinatorr.chess.server.data;
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
import org.springframework.stereotype.Repository;
import ru.fokinatorr.chess.server.game.FinishedGame;
import java.util.UUID;
/// Spring Data [FinishedGame] reactive repository.
@Repository
interface FinishedGameRepository extends ReactiveCrudRepository<FinishedGame, UUID> {
}
@@ -0,0 +1,124 @@
package ru.fokinatorr.chess.server.data;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.AccessLevel;
import lombok.RequiredArgsConstructor;
import lombok.experimental.FieldDefaults;
import org.jetbrains.annotations.Nullable;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.util.function.Tuples;
import ru.fokinatorr.chess.impl.ChessBoard;
import ru.fokinatorr.chess.impl.ChessUtil;
import ru.fokinatorr.chess.impl.PieceColor;
import ru.fokinatorr.chess.impl.misc.ChessBoardAndMoveStack;
import ru.fokinatorr.chess.impl.movestack.Move;
import ru.fokinatorr.chess.server.exception.UserNotFoundException;
import ru.fokinatorr.chess.server.game.FinishedGame;
import ru.fokinatorr.chess.server.game.event.GameFinishedEvent;
import java.util.UUID;
/// This service restores relations returned from the actual `FinishedGameRepository`.
@Service
@RequiredArgsConstructor
@FieldDefaults(makeFinal = true, level = AccessLevel.PRIVATE)
public class FinishedGameService {
FinishedGameRepository internalRepository;
ObjectMapper objectMapper;
UserService userService;
private Mono<FinishedGame> chessBoardToJson(FinishedGame source) {
return Mono.fromCallable(() -> {
ChessBoardAndMoveStack chessBoardAndMoveStack = source.getChessBoardAndMoveStack();
ObjectNode root = new ObjectNode(objectMapper.getNodeFactory());
root.set("chessBoard", objectMapper.valueToTree(source.getInitialChessBoard()));
ObjectNode moveStackNode = new ObjectNode(objectMapper.getNodeFactory());
ArrayNode doneMovesNode = new ArrayNode(objectMapper.getNodeFactory());
for (Move move : chessBoardAndMoveStack.getMoveStack().getDoneMoves()) {
doneMovesNode.add(move.toUCI());
}
root.<ObjectNode>set("moveStack", moveStackNode.set("done", doneMovesNode))
.put("currentlyMoving", chessBoardAndMoveStack.getChessBoard().getCurrentlyMoving().name());
source.setRawChessBoardAndMoveStack(root);
return source;
});
}
private Mono<FinishedGame> jsonToChessBoard(FinishedGame source) {
return Mono.fromCallable(() -> {
ObjectNode root = (ObjectNode) source.getRawChessBoardAndMoveStack();
ChessBoard initialBoard = objectMapper.treeToValue(root.get("chessBoard"), ChessBoard.class);
ChessBoardAndMoveStack chessBoardAndMoveStack = new ChessBoardAndMoveStack(initialBoard);
for (var iterator = root.get("moveStack").get("done").elements(); iterator.hasNext();) {
chessBoardAndMoveStack.perform(ChessUtil.parseUCIMove(iterator.next().asText(), chessBoardAndMoveStack));
}
chessBoardAndMoveStack.getChessBoard().setCurrentlyMoving(PieceColor.valueOf(root.get("currentlyMoving").asText()));
source.setInitialChessBoard(initialBoard);
source.setChessBoardAndMoveStack(chessBoardAndMoveStack);
return source;
});
}
private Mono<FinishedGame> restoreRelations(FinishedGame source) {
return jsonToChessBoard(source)
.zipWhen(game -> userService.findById(game.getWhiteUserId()))
.map(tuple -> {
tuple.getT1().setWhite(tuple.getT2());
return tuple.getT1();
})
.zipWhen(game -> userService.findById(game.getBlackUserId()))
.map(tuple -> {
tuple.getT1().setBlack(tuple.getT2());
return tuple.getT1();
})
.doOnNext(game -> game.setNew(false));
}
/// Returns a `Publisher` that creates and emits an `FinishedGame` with the given data.
public Mono<FinishedGame> create(UUID id, long whiteUserId, long blackUserId,
ChessBoard initialBoard, ChessBoardAndMoveStack endBoard,
@Nullable PieceColor winner,
GameFinishedEvent.GameFinishCause finishCause) {
return userService.findById(whiteUserId)
.switchIfEmpty(Mono.error(() -> new UserNotFoundException(whiteUserId)))
.zipWith(userService.findById(blackUserId)
.switchIfEmpty(Mono.error(() -> new UserNotFoundException(blackUserId))))
.map(tuple -> {
FinishedGame finishedGame = new FinishedGame(id, tuple.getT1(), tuple.getT2(), endBoard.clone(), initialBoard.clone(), finishCause);
finishedGame.setWinner(winner);
finishedGame.setWhiteUserId(whiteUserId);
finishedGame.setBlackUserId(blackUserId);
return finishedGame;
});
}
/// Saves the entity emitted by the given `Publisher` and returns a `Publisher` emitting the saved entity.
public Mono<FinishedGame> save(Mono<FinishedGame> entity) {
return entity
.flatMap(this::chessBoardToJson)
.flatMap(internalRepository::save)
.flatMap(this::restoreRelations);
}
/// Convenience method to allow `flatMap()` calls instead of `as()` calls
public Mono<FinishedGame> save(FinishedGame entity) {
return save(Mono.just(entity));
}
/// Returns a `Publisher` emitting an [FinishedGame] with the given ID, or nothing if there's no entity
/// with such ID.
public Mono<FinishedGame> findById(UUID id) {
return internalRepository.findById(id).flatMap(this::restoreRelations);
}
/// Returns a `Publisher` emitting all available entities.
public Flux<FinishedGame> findAll() {
return internalRepository.findAll().flatMap(this::restoreRelations);
}
}
@@ -0,0 +1,19 @@
package ru.fokinatorr.chess.server.data;
import com.fasterxml.jackson.databind.JsonNode;
import org.springframework.data.r2dbc.repository.Query;
import org.springframework.data.repository.query.Param;
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
import org.springframework.stereotype.Repository;
import reactor.core.publisher.Mono;
import ru.fokinatorr.chess.server.game.OngoingGame;
import java.util.UUID;
/// Spring Data [OngoingGame] reactive repository.
@Repository
interface OngoingGameRepository extends ReactiveCrudRepository<OngoingGame, UUID> {
@Query("UPDATE ongoing_game SET column raw_chess_board_and_move_stack = $1 WHERE id = $2")
Mono<Integer> updatePosition(JsonNode rawChessBoardAndMoveStack, UUID id);
}
@@ -0,0 +1,132 @@
package ru.fokinatorr.chess.server.data;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import io.r2dbc.postgresql.codec.Json;
import lombok.AccessLevel;
import lombok.RequiredArgsConstructor;
import lombok.experimental.FieldDefaults;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.util.function.Tuples;
import ru.fokinatorr.chess.impl.ChessBoard;
import ru.fokinatorr.chess.impl.PieceColor;
import ru.fokinatorr.chess.impl.ChessUtil;
import ru.fokinatorr.chess.impl.misc.ChessBoardAndMoveStack;
import ru.fokinatorr.chess.impl.movestack.Move;
import ru.fokinatorr.chess.server.exception.UserNotFoundException;
import ru.fokinatorr.chess.server.game.OngoingGame;
import java.util.UUID;
/// This service restores relations returned from the actual `OngoingGameRepository`.
@Service
@RequiredArgsConstructor
@FieldDefaults(makeFinal = true, level = AccessLevel.PRIVATE)
public class OngoingGameService {
OngoingGameRepository internalRepository;
ObjectMapper objectMapper;
UserService userService;
private Mono<OngoingGame> chessBoardToJson(OngoingGame source) {
return Mono.fromCallable(() -> {
ChessBoardAndMoveStack chessBoardAndMoveStack = source.getChessBoardAndMoveStack();
ObjectNode root = objectMapper.createObjectNode();
root.set("chessBoard", objectMapper.valueToTree(source.getInitialChessBoard()));
ObjectNode moveStackNode = objectMapper.createObjectNode();
ArrayNode doneMovesNode = objectMapper.createArrayNode();
for (Move move : chessBoardAndMoveStack.getMoveStack().getDoneMoves()) {
doneMovesNode.add(move.toUCI());
}
root.<ObjectNode>set("moveStack", moveStackNode.set("done", doneMovesNode))
.put("currentlyMoving", chessBoardAndMoveStack.getChessBoard().getCurrentlyMoving().name());
source.setRawChessBoardAndMoveStack(root);
return source;
});
}
private Mono<OngoingGame> jsonToChessBoard(OngoingGame source) {
return Mono.fromCallable(() -> {
ObjectNode root = (ObjectNode) source.getRawChessBoardAndMoveStack();
ChessBoard initialBoard = objectMapper.treeToValue(root.get("chessBoard"), ChessBoard.class);
ChessBoardAndMoveStack chessBoardAndMoveStack = new ChessBoardAndMoveStack(initialBoard);
for (var iterator = root.get("moveStack").get("done").elements(); iterator.hasNext();) {
chessBoardAndMoveStack.perform(ChessUtil.parseUCIMove(iterator.next().asText(), chessBoardAndMoveStack));
}
chessBoardAndMoveStack.getChessBoard().setCurrentlyMoving(PieceColor.valueOf(root.get("currentlyMoving").asText()));
source.setInitialChessBoard(initialBoard);
source.setChessBoardAndMoveStack(chessBoardAndMoveStack);
return source;
});
}
private Mono<OngoingGame> restoreRelations(OngoingGame source) {
return jsonToChessBoard(source)
.zipWhen(game -> userService.findById(game.getWhiteUserId()))
.map(tuple -> {
tuple.getT1().setWhite(tuple.getT2());
return tuple.getT1();
})
.zipWhen(game -> userService.findById(game.getBlackUserId()))
.map(tuple -> {
tuple.getT1().setBlack(tuple.getT2());
return tuple.getT1();
})
.doOnNext(game -> game.setNew(false));
}
/// Returns a `Publisher` that creates and emits an `OngoingGame` with the given data.
public Mono<OngoingGame> create(UUID id, long whiteUserId, long blackUserId, ChessBoardAndMoveStack initialBoard) {
return userService.findById(whiteUserId)
.switchIfEmpty(Mono.error(() -> new UserNotFoundException(whiteUserId)))
.zipWith(userService.findById(blackUserId)
.switchIfEmpty(Mono.error(() -> new UserNotFoundException(blackUserId))))
.zipWith(Mono.just(initialBoard), (tuple, chessBoardAndMoveStack) -> Tuples.of(tuple.getT1(), tuple.getT2(), chessBoardAndMoveStack))
.map(tuple -> {
OngoingGame ongoingGame = new OngoingGame(id, tuple.getT1(), tuple.getT2(), tuple.getT3().clone());
ongoingGame.setInitialChessBoard(tuple.getT3().getChessBoard().clone());
ongoingGame.setWhiteUserId(whiteUserId);
ongoingGame.setBlackUserId(blackUserId);
return ongoingGame;
});
}
/// Saves the entity emitted by the given `Publisher` and returns a `Publisher` emitting the saved entity.
public Mono<OngoingGame> save(Mono<OngoingGame> entity) {
return entity
.flatMap(this::chessBoardToJson)
.flatMap(internalRepository::save)
.flatMap(this::restoreRelations);
}
/// Convenience method to allow `flatMap()` calls instead of `as()` calls
public Mono<OngoingGame> save(OngoingGame entity) {
return save(Mono.just(entity));
}
/// Returns a `Publisher` emitting an [OngoingGame] with the given ID, or nothing if there's no entity
/// with such ID.
public Mono<OngoingGame> findById(UUID id) {
return internalRepository.findById(id).flatMap(this::restoreRelations);
}
/// Returns a `Publisher` emitting `true` if an entity with the given ID exists, `false` otherwise.
public Mono<Boolean> existsById(UUID id) {
return internalRepository.existsById(id);
}
/// Returns a `Publisher` signaling when an entity with the given ID was deleted, or it was detected that there's no
/// entity with such ID.
public Mono<Void> deleteById(UUID id) {
return internalRepository.deleteById(id);
}
/// Returns a `Publisher` emitting all available entities.
public Flux<OngoingGame> findAll() {
return internalRepository.findAll().flatMap(this::restoreRelations);
}
}
@@ -0,0 +1,10 @@
package ru.fokinatorr.chess.server.data;
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
import ru.fokinatorr.chess.server.game.PendingGame;
import java.util.UUID;
/// Spring Data [PendingGame] reactive repository.
interface PendingGameRepository extends ReactiveCrudRepository<PendingGame, UUID> {
}
@@ -0,0 +1,110 @@
package ru.fokinatorr.chess.server.data;
import lombok.AccessLevel;
import lombok.RequiredArgsConstructor;
import lombok.experimental.FieldDefaults;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import ru.fokinatorr.chess.impl.PieceColorOrRandom;
import ru.fokinatorr.chess.server.exception.UserNotFoundException;
import ru.fokinatorr.chess.server.game.PendingGame;
import java.util.Optional;
import java.util.UUID;
/// This service restores relations returned from the actual `PendingGameRepository`.
@Slf4j
@Service
@RequiredArgsConstructor
@FieldDefaults(makeFinal = true, level = AccessLevel.PRIVATE)
public class PendingGameService {
PendingGameRepository internalRepository;
UserService userService;
private Mono<PendingGame> restoreRelations(PendingGame source) {
return Mono.just(source)
.zipWhen(game -> userService.findById(game.getWithUserId()))
.map(tuple -> {
tuple.getT1().setWithUser(tuple.getT2());
return tuple.getT1();
})
.doOnNext(game -> game.setNew(false));
}
/// Returns a `Publisher` that creates and emits an `PendingGame` with the given data.
public Mono<PendingGame> create(UUID id, long withUserId, PieceColorOrRandom ownerPlaysAs) {
return Mono.defer(() -> {
log.info("Creating pending game");
return userService.findById(withUserId)
.switchIfEmpty(Mono.error(() -> new UserNotFoundException(withUserId)))
.map(withUser -> {
PendingGame pendingGame = new PendingGame(id, withUser, ownerPlaysAs);
pendingGame.setWithUserId(withUserId);
log.info("Created pending game: {}", pendingGame);
return pendingGame;
});
});
}
/// Saves the given entity and returns a `Publisher` emitting the saved entity.
public Mono<PendingGame> save(PendingGame entity) {
return Mono.defer(() -> {
log.info("Saving pending game: {}", entity);
return Mono.just(entity)
.flatMap(internalRepository::save)
.flatMap(this::restoreRelations)
.doOnNext(saved -> log.info("Saved pending game: {}", saved));
});
}
/// Returns a `Publisher` emitting an [PendingGame] with the given ID, or nothing if there's no entity
/// with such ID.
public Mono<PendingGame> findById(UUID id) {
return Mono.defer(() -> {
log.info("Finding pending game by ID: {}", id);
return internalRepository.findById(id)
.flatMap(this::restoreRelations)
.doOnNext(game -> log.info("Retrieved pending game by ID: {}", game))
.switchIfEmpty(Mono.fromRunnable(() -> log.info("No pending game with ID {}", id)));
});
}
/// Returns a `Publisher` emitting `true` if an entity with the given ID exists, `false` otherwise.
public Mono<Boolean> existsById(UUID id) {
return Mono.defer(() -> {
log.info("Determining if pending game with ID {} exists...", id);
return internalRepository.existsById(id)
.doOnNext(b -> {
if (b) {
log.info("Pending game with ID {} exists", id);
} else {
log.info("Pending game with ID {} does not exist", id);
}
});
});
}
/// Returns a `Publisher` signaling when an entity with the given ID was deleted, or it was detected that there's no
/// entity with such ID.
public Mono<Void> deleteById(UUID id) {
return Mono.defer(() -> {
log.info("Deleting pending game by ID: {}", id);
return internalRepository.deleteById(id)
.doOnSuccess(_ -> log.info("Deleted pending game by ID: {}", id));
});
}
/// Returns a `Publisher` emitting all available entities.
public Flux<PendingGame> findAll() {
return Flux.defer(() -> {
log.info("Finding all pending games...");
return internalRepository.findAll()
.flatMap(this::restoreRelations)
.doOnNext(game -> log.info("Found pending game: {}", game))
.switchIfEmpty(Mono.fromRunnable(() -> log.info("No pending games available")));
});
}
}
@@ -0,0 +1,15 @@
package ru.fokinatorr.chess.server.data;
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
import org.springframework.stereotype.Repository;
import reactor.core.publisher.Flux;
import ru.fokinatorr.chess.server.entity.UserNotification;
/// Spring Data [UserNotification] reactive repository.
@Repository
interface UserNotificationRepository extends ReactiveCrudRepository<UserNotification, Long> {
Flux<UserNotification> findAllBySenderId(long senderId);
Flux<UserNotification> findAllByReceiverId(long receiverId);
}
@@ -0,0 +1,91 @@
package ru.fokinatorr.chess.server.data;
import lombok.AccessLevel;
import lombok.RequiredArgsConstructor;
import lombok.experimental.FieldDefaults;
import lombok.extern.slf4j.Slf4j;
import org.reactivestreams.Publisher;
import org.springframework.dao.OptimisticLockingFailureException;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import ru.fokinatorr.chess.server.entity.UserNotification;
/// This service restores relations returned from the actual `UserNotificationRepository`.
@Slf4j
@Service
@RequiredArgsConstructor
@FieldDefaults(makeFinal = true, level = AccessLevel.PRIVATE)
public class UserNotificationService {
UserNotificationRepository internalRepository;
UserService userService;
/// Saves a given entity. Use the returned instance for further operations as the save operation might have changed the
/// entity instance completely.
///
/// @param entity must not be {@literal null}.
/// @return {@link Mono} emitting the saved entity.
/// @throws IllegalArgumentException in case the given {@literal entity} is {@literal null}.
/// @throws OptimisticLockingFailureException when the entity uses optimistic locking and has a version attribute with
/// a different value from that found in the persistence store. Also thrown if the entity is assumed to be
/// present but does not exist in the database.
public Mono<UserNotification> save(UserNotification entity) {
return Mono.defer(() -> {
log.info("Saving notification: {}", entity);
return internalRepository.save(entity)
.doOnNext(savedEntity -> log.info("Saved notification: {}", savedEntity));
});
}
/// Returns a [Mono] emitting a [UserNotification] with the given ID, or nothing if there's no notification with such ID.
public Mono<UserNotification> findById(long id) {
return Mono.defer(() -> {
log.info("Finding notification by ID: {}", id);
return internalRepository.findById(id)
.transformDeferred(this::restoreRelations)
.doOnNext(notification -> log.info("Retrieved notification by ID: {}", notification))
.switchIfEmpty(Mono.fromRunnable(() -> log.info("No notification with ID {}", id)));
});
}
/// Returns a [Flux] emitting all [UserNotification]s with the given sender user ID,
/// or nothing if there's no user with such ID.
public Flux<UserNotification> findAllBySenderId(long senderId) {
return Flux.defer(() -> {
log.info("Finding notifications by sender user ID: {}", senderId);
return internalRepository.findAllBySenderId(senderId)
.flatMap(this::restoreRelations)
.doOnNext(notification -> log.info("Found notification with sender user ID {}: {}", senderId, notification))
.switchIfEmpty(Mono.fromRunnable(() -> log.info("No notifications with sender user ID {}", senderId)));
});
}
/// Returns a [Flux] emitting all [UserNotification]s with the given receiver user ID,
/// or nothing if there's no user with such ID.
public Flux<UserNotification> findAllByReceiverId(long receiverId) {
return Flux.defer(() -> {
log.info("Finding notifications by receiver user ID: {}", receiverId);
return internalRepository.findAllByReceiverId(receiverId)
.flatMap(this::restoreRelations)
.doOnNext(notification -> log.info("Found notification with receiver user ID {}: {}", receiverId, notification))
.switchIfEmpty(Mono.fromRunnable(() -> log.info("No notifications with receiver user ID {}", receiverId)));
});
}
private Mono<UserNotification> restoreRelations(Mono<UserNotification> source) {
return Mono.from(source).flatMap(notification -> Mono.zip(
Mono.just(notification),
userService.findById(notification.getSenderId()),
userService.findById(notification.getReceiverId())
)).map(tuple -> {
tuple.getT1().setSender(tuple.getT2());
tuple.getT1().setReceiver(tuple.getT3());
return tuple.getT1();
});
}
private Mono<UserNotification> restoreRelations(UserNotification notification) {
return restoreRelations(Mono.just(notification));
}
}
@@ -0,0 +1,18 @@
package ru.fokinatorr.chess.server.data;
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
import org.springframework.stereotype.Repository;
import reactor.core.publisher.Mono;
import ru.fokinatorr.chess.server.entity.User;
/// Spring Data [User] reactive repository.
@Repository
interface UserRepository extends ReactiveCrudRepository<User, Long> {
/// Returns a [Mono] emitting a [User] with the given username, or nothing if there's no user with such name.
Mono<User> findByUsername(String username);
Mono<Void> deleteByUsername(String username);
Mono<Boolean> existsByUsername(String username);
}
@@ -0,0 +1,139 @@
package ru.fokinatorr.chess.server.data;
import lombok.AccessLevel;
import lombok.RequiredArgsConstructor;
import lombok.experimental.FieldDefaults;
import lombok.extern.slf4j.Slf4j;
import org.springframework.dao.OptimisticLockingFailureException;
import org.springframework.security.crypto.password.PasswordEncoder;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import ru.fokinatorr.chess.server.entity.User;
import ru.fokinatorr.chess.server.security.GrantedAuthorityImpl;
import java.util.stream.Collectors;
/// This service restores relations returned from the actual `UserRepository`.
@Service
@RequiredArgsConstructor
@FieldDefaults(makeFinal = true, level = AccessLevel.PRIVATE)
@Slf4j
public class UserService {
UserRepository internalRepository;
PasswordEncoder encoder;
/// Saves a given entity. Use the returned instance for further operations as the save operation might have changed the
/// entity instance completely.
///
/// @param entity must not be {@literal null}.
/// @return {@link Mono} emitting the saved entity.
/// @throws IllegalArgumentException in case the given {@literal entity} is {@literal null}.
/// @throws OptimisticLockingFailureException when the entity uses optimistic locking and has a version attribute with
/// a different value from that found in the persistence store. Also thrown if the entity is assumed to be
/// present but does not exist in the database.
public Mono<User> save(User entity) {
return Mono.defer(() -> {
log.info("Saving user: {}", entity);
return internalRepository.save(entity)
.flatMap(this::restoreRelations)
.doOnNext(savedEntity -> log.info("Saved user: {}", savedEntity));
});
}
/// Returns a [Mono] emitting a [User] with the given ID, or nothing if there's no user with such ID.
public Mono<User> findById(long id) {
return Mono.defer(() -> {
log.info("Finding user by ID: {}", id);
return internalRepository.findById(id)
.flatMap(this::restoreRelations)
.doOnNext(user -> log.info("Retrieved user by ID: {}", user))
.switchIfEmpty(Mono.fromRunnable(() -> log.info("No user with ID {}", id)));
});
}
/// Returns a [Mono] emitting `true` if the [User] with the given ID exists, `false` otherwise.
public Mono<Boolean> existsById(long id) {
return Mono.defer(() -> {
log.info("Determining if user with ID {} exists...", id);
return internalRepository.existsById(id)
.doOnNext(b -> {
if (b) {
log.info("User with ID {} exists", id);
} else {
log.info("User with ID {} does not exist", id);
}
});
});
}
/// Returns a [Mono] emitting a [User] with the given username, or nothing if there's no user with such name.
public Mono<User> findByUsername(String username) {
return Mono.defer(() -> {
log.info("Finding user by username: {}", username);
return internalRepository.findByUsername(username)
.flatMap(this::restoreRelations)
.doOnNext(user -> log.info("Retrieved user by username: {}", user))
.switchIfEmpty(Mono.fromRunnable(() -> log.info("No user with username {}", username)));
});
}
/// Returns a [Flux] emitting all available [User]s.
public Flux<User> findAll() {
return Flux.defer(() -> {
log.info("Finding all users...");
return internalRepository.findAll()
.flatMap(this::restoreRelations)
.doOnNext(user -> log.info("Found user: {}", user))
.switchIfEmpty(Mono.fromRunnable(() -> log.info("No users available")));
});
}
/// Returns a [Mono] that emits `true` if there is a user with the given username, `false` otherwise.
public Mono<Boolean> existsByUsername(String username) {
return Mono.defer(() -> {
log.info("Determining if user with username {} exists...", username);
return internalRepository.existsByUsername(username)
.doOnNext(b -> {
if (b) {
log.info("User with username {} exists", username);
} else {
log.info("User with username {} does not exist", username);
}
});
});
}
/// Returns a [Mono] signaling when the user with the given ID was deleted in the database.
public Mono<Void> deleteByUsername(String username) {
return Mono.defer(() -> {
log.info("Deleting user by username: {}", username);
return internalRepository.deleteByUsername(username)
.doOnSuccess(_ -> log.info("Deleted user by username: {}", username));
});
}
/// Returns a `Publisher` that creates and emits a `User` with the given data.
/// @implNote ***The password should be passed in an UNENCODED view, it's encoded automatically.***
public Mono<User> create(String username, String password, boolean isAdmin) {
return Mono.fromCallable(() -> new User(username, encoder.encode(password)))
.doOnNext(user -> {
user.addAuthority(GrantedAuthorityImpl.ROLE_USER);
if (isAdmin) {
user.addAuthority(GrantedAuthorityImpl.ROLE_ADMIN);
}
log.info("Created user: {}", user);
});
}
private Mono<User> restoreRelations(User source) {
return Mono.just(source)
.map(user -> {
user.setAuthorities(user.getAuthorityIds().stream()
.map(GrantedAuthorityImpl::valueOf)
.collect(Collectors.toUnmodifiableSet()));
return user;
});
}
}
@@ -0,0 +1,30 @@
package ru.fokinatorr.chess.server.data.converter;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import io.r2dbc.postgresql.codec.Json;
import lombok.AccessLevel;
import lombok.RequiredArgsConstructor;
import lombok.SneakyThrows;
import lombok.experimental.FieldDefaults;
import org.jetbrains.annotations.NotNull;
import org.springframework.core.convert.converter.Converter;
import org.springframework.data.convert.WritingConverter;
import org.springframework.stereotype.Component;
/// Converts Jackson's [JsonNode] to PostgreSQL's jsonb
@WritingConverter
@Component
@RequiredArgsConstructor
@FieldDefaults(makeFinal = true, level = AccessLevel.PRIVATE)
public class JsonNodeToJsonConverter implements Converter<JsonNode, Json> {
ObjectMapper objectMapper;
@Override
@SneakyThrows(JsonProcessingException.class)
public Json convert(@NotNull JsonNode source) {
return Json.of(objectMapper.writeValueAsBytes(source));
}
}
@@ -0,0 +1,31 @@
package ru.fokinatorr.chess.server.data.converter;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import io.r2dbc.postgresql.codec.Json;
import lombok.AccessLevel;
import lombok.RequiredArgsConstructor;
import lombok.SneakyThrows;
import lombok.experimental.FieldDefaults;
import org.jetbrains.annotations.NotNull;
import org.springframework.core.convert.converter.Converter;
import org.springframework.data.convert.ReadingConverter;
import org.springframework.stereotype.Component;
import java.io.IOException;
/// Converts PostgreSQL's jsonb to Jackson's [JsonNode]
@ReadingConverter
@Component
@RequiredArgsConstructor
@FieldDefaults(makeFinal = true, level = AccessLevel.PRIVATE)
public class JsonToJsonNodeConverter implements Converter<Json, JsonNode> {
ObjectMapper objectMapper;
@Override
@SneakyThrows(IOException.class)
public JsonNode convert(@NotNull Json source) {
return objectMapper.readTree(source.asArray());
}
}
@@ -0,0 +1,37 @@
package ru.fokinatorr.chess.server.entity;
import lombok.*;
import lombok.experimental.FieldDefaults;
import org.springframework.data.annotation.Id;
import org.springframework.data.annotation.Transient;
import org.springframework.data.relational.core.mapping.Table;
import org.springframework.security.core.userdetails.UserDetails;
import ru.fokinatorr.chess.server.security.GrantedAuthorityImpl;
import java.util.HashSet;
import java.util.Set;
/// Chess Server [UserDetails] implementation.
@Data
@Table(name = "users")
@RequiredArgsConstructor
@NoArgsConstructor(access = AccessLevel.PROTECTED, force = true)
@FieldDefaults(level = AccessLevel.PRIVATE)
public final class User implements UserDetails {
@Id
Long id;
@NonNull
String username;
@NonNull
String password;
@Transient
transient Set<GrantedAuthorityImpl> authorities = new HashSet<>();
Set<Integer> authorityIds = new HashSet<>();
public void addAuthority(GrantedAuthorityImpl authority) {
authorities.add(authority);
authorityIds.add(authority.getId());
}
}
@@ -0,0 +1,40 @@
package ru.fokinatorr.chess.server.entity;
import lombok.*;
import lombok.experimental.FieldDefaults;
import org.springframework.data.annotation.Id;
import org.springframework.data.annotation.Transient;
import org.springframework.data.relational.core.mapping.Table;
import java.io.Serializable;
import java.util.HashMap;
import java.util.Map;
/// A notification that was sent from a user to another user.
@Data
@Table(name = "user_notification")
@RequiredArgsConstructor
@NoArgsConstructor(access = AccessLevel.PROTECTED, force = true)
@FieldDefaults(level = AccessLevel.PRIVATE)
public class UserNotification implements Serializable {
@Id
Long id;
@NonNull
@Transient
transient User sender;
@NonNull
@Transient
transient User receiver;
@NonNull
Type type;
@NonNull
Map<String, Object> additionalData = new HashMap<>();
long senderId;
long receiverId;
/// The type of notification.
public enum Type {
GAME_INVITATION
}
}
@@ -0,0 +1,49 @@
package ru.fokinatorr.chess.server.game;
import com.fasterxml.jackson.databind.JsonNode;
import lombok.*;
import lombok.experimental.FieldDefaults;
import org.springframework.data.annotation.Id;
import org.springframework.data.annotation.Transient;
import org.springframework.data.domain.Persistable;
import org.springframework.data.relational.core.mapping.Table;
import ru.fokinatorr.chess.impl.ChessBoard;
import ru.fokinatorr.chess.impl.PieceColor;
import ru.fokinatorr.chess.impl.misc.ChessBoardAndMoveStack;
import ru.fokinatorr.chess.server.entity.User;
import ru.fokinatorr.chess.server.game.event.GameFinishedEvent;
import java.io.Serializable;
import java.util.UUID;
/// A finished chess game.
@Data
@Table(name = "finished_game")
@RequiredArgsConstructor
@NoArgsConstructor(access = AccessLevel.PROTECTED, force = true)
@FieldDefaults(level = AccessLevel.PRIVATE)
public class FinishedGame implements Serializable, Persistable<UUID> {
@Id
@NonNull
UUID id;
@Transient
@NonNull
transient User white;
long whiteUserId;
@Transient
@NonNull
transient User black;
long blackUserId;
@Transient
@NonNull
transient ChessBoardAndMoveStack chessBoardAndMoveStack;
@Transient
@NonNull
transient ChessBoard initialChessBoard;
JsonNode rawChessBoardAndMoveStack;
PieceColor winner;
@NonNull
GameFinishedEvent.GameFinishCause finishCause;
@Transient
transient boolean isNew = true;
}
@@ -0,0 +1,413 @@
package ru.fokinatorr.chess.server.game;
import lombok.AccessLevel;
import lombok.RequiredArgsConstructor;
import lombok.experimental.FieldDefaults;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.core.publisher.Sinks;
import reactor.core.scheduler.Schedulers;
import reactor.util.Loggers;
import reactor.util.function.Tuple2;
import ru.fokinatorr.chess.impl.ChessBoard;
import ru.fokinatorr.chess.impl.ChessUtil;
import ru.fokinatorr.chess.impl.PieceColor;
import ru.fokinatorr.chess.impl.misc.ChessBoardAndMoveStack;
import ru.fokinatorr.chess.server.BooleanOperators;
import ru.fokinatorr.chess.server.data.FinishedGameService;
import ru.fokinatorr.chess.server.data.OngoingGameService;
import ru.fokinatorr.chess.server.data.PendingGameService;
import ru.fokinatorr.chess.server.data.UserService;
import ru.fokinatorr.chess.server.exception.*;
import ru.fokinatorr.chess.server.game.event.*;
import ru.fokinatorr.chess.server.misc.DTOServiceExtension;
import ru.fokinatorr.chess.server.misc.DefaultEmitFailureHandlingSinksMany;
import ru.fokinatorr.chess.server.service.DTOService;
import java.util.Optional;
import java.util.Random;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
/// This service controls the chess games.
@Service
@FieldDefaults(level = AccessLevel.PRIVATE, makeFinal = true)
@RequiredArgsConstructor
public class GameService {
private static final Sinks.EmitFailureHandler CONTINUE_IF_NO_SUB = (_, emitResult) -> emitResult == Sinks.EmitResult.FAIL_ZERO_SUBSCRIBER;
ConcurrentMap<UUID, DefaultEmitFailureHandlingSinksMany<GameEvent>> ongoingGameEventBroadcasters = new ConcurrentHashMap<>();
ConcurrentMap<UUID, Sinks.One<GameJoinedEvent>> gameJoinEventBroadcasters = new ConcurrentHashMap<>();
UserService userService;
PendingGameService pendingGameService;
OngoingGameService ongoingGameService;
FinishedGameService finishedGameService;
DTOServiceExtension dtoService;
Random random;
// /////////////////////////////////////////////
// /////////////////////////////////////////////
// ////////////// HTTP methods /////////////////
// /////////////////////////////////////////////
// /////////////////////////////////////////////
/// Creates a new pending game as the given principal and return its ID.
public Mono<UUID> createGame(long principalUserId, CreateGameOptions options) {
return assertUserExists(principalUserId, Mono.fromSupplier(UUID::randomUUID))
.subscribeOn(Schedulers.boundedElastic())
.flatMap(id -> pendingGameService.create(id, principalUserId, options.playAs()))
.flatMap(pendingGameService::save)
.map(PendingGame::getId)
.doOnNext(gameId -> gameJoinEventBroadcasters.put(gameId, Sinks.one()));
}
/// Tries to cancel a pending game with the given ID as the given principal.
public Mono<Void> cancelPendingGame(UUID gameId, long principalUserId) {
return assertUserExists(principalUserId, pendingGameService.findById(gameId))
.switchIfEmpty(Mono.error(() -> new GameNotFoundException(gameId)))
.filter(game -> game.getWithUserId() == principalUserId)
.switchIfEmpty(Mono.error(() -> new NotParticipantException(principalUserId, gameId)))
.flatMap(_ -> {
getOrCreateJoinBroadcaster(gameId).emitError(new GameCancelledException(gameId), CONTINUE_IF_NO_SUB);
return pendingGameService.deleteById(gameId);
});
}
/// Tries to cancel an ongoing game with the given ID as the given principal. Only possible if no move was done and `principalUserId == <any participant id>`
public Mono<Void> cancelOngoingGame(UUID gameId, long principalUserId) {
return assertUserExists(principalUserId, ongoingGameService.findById(gameId))
.switchIfEmpty(Mono.error(() -> new GameNotFoundException(gameId)))
.filter(game -> game.getWhiteUserId() == principalUserId || game.getBlackUserId() == principalUserId)
.switchIfEmpty(Mono.error(() -> new NotParticipantException(principalUserId, gameId)))
.filter(game -> game.getChessBoardAndMoveStack().peekDone() == null)
.switchIfEmpty(Mono.error(() -> new MoveAlreadyDoneException(gameId)))
.flatMap(game -> {
getOrCreateEventBroadcaster(gameId).emitNextDefault(
new GameFinishedEvent(
null,
GameFinishedEvent.GameFinishCause.CANCELLED,
dtoService.toDTO(game.getChessBoardAndMoveStack(), game.getInitialChessBoard())
)
);
getOrCreateEventBroadcaster(gameId).emitCompleteDefault();
return ongoingGameService.deleteById(gameId);
});
}
/// Creates a new ongoing game that is played between the before stored `User` and the given one.
public Mono<Void> joinGame(UUID id, long joiningUserId) {
return pendingGameService.findById(id)
.switchIfEmpty(Mono.error(() -> new GameNotFoundException(id)))
.zipWith(assertUserExists(joiningUserId, userService.findById(joiningUserId)))
.flatMap(tuple -> Mono.defer(() -> {
long whiteUserId, blackUserId;
PieceColor color = tuple.getT1().getOwnerPlaysAs().getActualPieceColor(random);
switch (color) {
case WHITE -> {
whiteUserId = tuple.getT1().getWithUserId();
blackUserId = tuple.getT2().getId();
}
case BLACK -> {
whiteUserId = tuple.getT2().getId();
blackUserId = tuple.getT1().getWithUserId();
}
default ->
throw new IllegalArgumentException("Unexpected value: " + color); // should never happen
}
ChessBoard startPos = ChessBoard.startPosition();
GameJoinedEvent gameJoinedEvent = new GameJoinedEvent(
startPos,
color,
dtoService.toDTO(tuple.getT2())
);
ongoingGameEventBroadcasters.put(id, new DefaultEmitFailureHandlingSinksMany<>(Sinks.many().replay().all(), CONTINUE_IF_NO_SUB));
getOrCreateJoinBroadcaster(id).emitValue(gameJoinedEvent, CONTINUE_IF_NO_SUB);
getOrCreateEventBroadcaster(id).emitNextDefault(gameJoinedEvent);
return ongoingGameService.create(id, whiteUserId, blackUserId, new ChessBoardAndMoveStack(startPos));
}))
.flatMap(ongoingGameService::save)
.thenEmpty(pendingGameService.deleteById(id));
}
/// Try to perform the given move in the given game as the given user.
public Mono<Void> performMove(UUID gameId, long principalUserId, Mono<String> uciMove) {
return checkGameAccess(gameId, principalUserId, true)
.zipWith(uciMove)
.flatMap(tuple -> doPerformMove(tuple.getT1(), tuple.getT2()))
.zipWhen(tuple -> ongoingGameService.save(tuple.getT2()), (tuple, game) -> tuple.mapT2(_ -> game))
.doOnNext(tuple -> getOrCreateEventBroadcaster(gameId).emitNextDefault(tuple.getT1()))
.flatMap(tuple -> checkTerminalState(tuple.getT2()))
.flatMap(gameFinishedEvent -> finishGame(gameFinishedEvent, gameId))
.then();
}
/// Try to offer draw in the given game as the given user.
public Mono<Void> offerDraw(UUID gameId, long principalUserId) {
return checkGameAccess(gameId, principalUserId, true)
.filter(game -> game.getDrawOfferedBy() == null)
.switchIfEmpty(Mono.error(() -> new DrawAlreadyOfferedException(gameId)))
.map(game -> {
// at this point we know that principalUserId is currently moving user
game.setDrawOfferedBy(game.getChessBoardAndMoveStack().getChessBoard().getCurrentlyMoving());
return game;
})
.flatMap(ongoingGameService::save)
.doOnNext(game -> getOrCreateEventBroadcaster(gameId).emitNextDefault(new DrawOfferedEvent(game.getDrawOfferedBy())))
.then();
}
/// Try to accept draw in the given game as the given user.
public Mono<Void> acceptDraw(UUID gameId, long principalUserId) {
return checkGameAccess(gameId, principalUserId, false)
.filter(game -> game.getDrawOfferedBy() != null)
.switchIfEmpty(Mono.error(() -> new DrawNotOfferedException(gameId)))
.filter(game -> canAcceptOrRefuseDraw(principalUserId, game))
.switchIfEmpty(Mono.error(() -> new CannotAcceptOrRefuseDrawException(principalUserId, gameId)))
.flatMap(game -> {
GameFinishedEvent gameFinishedEvent = new GameFinishedEvent(
null,
GameFinishedEvent.GameFinishCause.DRAW_ACCEPTED,
dtoService.toDTO(game.getChessBoardAndMoveStack(), game.getInitialChessBoard())
);
var broadcaster = getOrCreateEventBroadcaster(gameId);
broadcaster.emitNextDefault(gameFinishedEvent);
broadcaster.emitCompleteDefault();
return finishGame(gameFinishedEvent, gameId);
})
.then();
}
/// Try to refuse draw in the given game as the given user.
public Mono<Void> refuseDraw(UUID gameId, long principalUserId) {
return checkGameAccess(gameId, principalUserId, false)
.filter(game -> game.getDrawOfferedBy() != null)
.switchIfEmpty(Mono.error(() -> new DrawNotOfferedException(gameId)))
.filter(game -> canAcceptOrRefuseDraw(principalUserId, game))
.switchIfEmpty(Mono.error(() -> new CannotAcceptOrRefuseDrawException(principalUserId, gameId)))
.flatMap(game -> {
DrawRefusedEvent drawRefusedEvent = new DrawRefusedEvent(game.getDrawOfferedBy().reverse());
game.setDrawOfferedBy(null);
getOrCreateEventBroadcaster(gameId).emitNextDefault(drawRefusedEvent);
return ongoingGameService.save(game);
})
.then();
}
/// Try to give up in the given game as the given user.
public Mono<Void> giveUp(UUID gameId, long principalUserId) {
return checkGameAccess(gameId, principalUserId, true)
.flatMap(game -> {
GameFinishedEvent gameFinishedEvent = new GameFinishedEvent(
game.getChessBoardAndMoveStack().getChessBoard().getCurrentlyMoving().reverse(),
GameFinishedEvent.GameFinishCause.GIVE_UP,
dtoService.toDTO(game.getChessBoardAndMoveStack(), game.getInitialChessBoard())
);
var broadcaster = getOrCreateEventBroadcaster(gameId);
broadcaster.emitNextDefault(gameFinishedEvent);
broadcaster.emitCompleteDefault();
return finishGame(gameFinishedEvent, gameId);
})
.then();
}
/// Returns all the pending games.
public Flux<PendingGame> getPendingGames() {
return pendingGameService.findAll();
}
/// Returns all the ongoing games.
public Flux<OngoingGame> getOngoingGames() {
return ongoingGameService.findAll();
}
/// Returns all the finished games.
public Flux<FinishedGame> getFinishedGames() {
return finishedGameService.findAll();
}
// /////////////////////////////////////////////
// /////////////////////////////////////////////
// ///////////// RSocket methods ///////////////
// /////////////////////////////////////////////
// /////////////////////////////////////////////
/// Returns a `Publisher` that emits events connected to a game with the given ID.
public Flux<GameEvent> subscribeToGame(UUID gameId) {
return ongoingGameService.existsById(gameId)
.filter(Boolean::booleanValue)
.switchIfEmpty(Mono.error(() -> new GameNotFoundException(gameId)))
.thenMany(getOrCreateEventBroadcaster(gameId).asFlux())
.log(Loggers.getLogger(GameService.class));
}
/// Returns a `Publisher` emitting the color that the owner user of the game with the given ID will play as, when someone joins the game.
public Mono<GameJoinedEvent> waitForJoin(UUID gameId) {
return ongoingGameService.existsById(gameId)
.filter(Boolean::booleanValue)
.switchIfEmpty(Mono.error(() -> new GameNotFoundException(gameId)))
.then(getOrCreateJoinBroadcaster(gameId).asMono())
.log(Loggers.getLogger(GameService.class));
}
// /////////////////////////////////////////////
// /////////////////////////////////////////////
// ////////////// Private Impl /////////////////
// /////////////////////////////////////////////
// /////////////////////////////////////////////
private Mono<Tuple2<MovePerformedEvent, OngoingGame>> doPerformMove(OngoingGame game, String uciMove) {
ChessBoardAndMoveStack chessBoardAndMoveStack = game.getChessBoardAndMoveStack();
ChessBoard initialBoard = game.getInitialChessBoard();
return Mono.fromCallable(() -> ChessUtil.parseUCIMove(uciMove, chessBoardAndMoveStack))
.filter(move -> ChessUtil.isValidMove(chessBoardAndMoveStack, move))
.switchIfEmpty(Mono.error(() -> new InvalidMoveException(uciMove, dtoService.toDTO(chessBoardAndMoveStack, initialBoard))))
.map(move -> {
chessBoardAndMoveStack.perform(move);
chessBoardAndMoveStack.reverseCurrentlyMoving();
return new MovePerformedEvent(uciMove, dtoService.toDTO(chessBoardAndMoveStack, initialBoard));
})
.zipWith(Mono.just(game));
}
/// Checks if the game is in terminal state (checkmate or stalemate), and if so, creates a [GameFinishedEvent],
/// emits it into the event broadcaster and completes it. If the game isn't in terminal state, returns an empty [Mono].
private Mono<GameFinishedEvent> checkTerminalState(OngoingGame game) {
// return Mono.fromCallable(() -> { // keep the original code in case the new doesn't work
// PieceColor winner = null;
// GameFinishedEvent.GameFinishCause gameFinishCause = null;
// if (game.getChessBoardAndMoveStack().isCheckmate(game.getChessBoardAndMoveStack().getChessBoard().getCurrentlyMoving())) {
// gameFinishCause = GameFinishedEvent.GameFinishCause.CHECKMATE;
// winner = game.getChessBoardAndMoveStack().getChessBoard().getCurrentlyMoving().reverse();
// } else if (game.getChessBoardAndMoveStack().isStalemate(game.getChessBoardAndMoveStack().getChessBoard().getCurrentlyMoving())) {
// gameFinishCause = GameFinishedEvent.GameFinishCause.STALEMATE;
// }
// return Tuples.of(Optional.ofNullable(winner), Optional.ofNullable(gameFinishCause));
// })
// .filter(tuple -> tuple.getT2().isPresent())
// .map(tuple -> new GameFinishedEvent(
// tuple.getT1().orElse(null),
// tuple.getT2().get(), // should not be null
// dtoService.toDTO(game.getChessBoardAndMoveStack(), game.getInitialChessBoard())
// ))
// .doOnNext(gameFinishedEvent -> {
// var broadcaster = getOrCreateEventBroadcaster(game.getId());
// broadcaster.emitNextDefault(gameFinishedEvent);
// broadcaster.emitCompleteDefault();
// });
return Mono.justOrEmpty(detectTerminalState(game, dtoService))
.doOnNext(gameFinishedEvent -> {
var broadcaster = getOrCreateEventBroadcaster(game.getId());
broadcaster.emitNextDefault(gameFinishedEvent);
broadcaster.emitCompleteDefault();
});
}
/// Internal implementation for [#checkTerminalState(OngoingGame)] that doesn't emit the event.
static Optional<GameFinishedEvent> detectTerminalState(OngoingGame game, DTOService dtoService) {
ChessBoardAndMoveStack stack = game.getChessBoardAndMoveStack();
PieceColor current = stack.getChessBoard().getCurrentlyMoving();
GameFinishedEvent.GameFinishCause cause;
PieceColor winner = null;
if (stack.isCheckmate(current)) {
cause = GameFinishedEvent.GameFinishCause.CHECKMATE;
winner = current.reverse(); // winner is the one who delivered mate
} else if (stack.isStalemate(current)) {
cause = GameFinishedEvent.GameFinishCause.STALEMATE;
} else {
return Optional.empty();
}
return Optional.of(new GameFinishedEvent(
winner,
cause,
dtoService.toDTO(stack, game.getInitialChessBoard())
));
}
private Mono<FinishedGame> finishGame(GameFinishedEvent gameFinishedEvent, UUID gameId) {
// Assume that the game-exists check was done
return Mono.just(gameFinishedEvent)
.zipWith(ongoingGameService.findById(gameId))
.flatMap(tuple -> {
GameFinishedEvent e = tuple.getT1();
OngoingGame game = tuple.getT2();
return finishedGameService.create(
game.getId(),
game.getWhiteUserId(),
game.getBlackUserId(),
game.getInitialChessBoard(),
game.getChessBoardAndMoveStack(),
e.winner(),
e.cause()
);
})
.flatMap(finishedGameService::save)
.flatMap(finishedGame -> ongoingGameService.deleteById(gameId).thenReturn(finishedGame))
.doOnSuccess(_ -> {
ongoingGameEventBroadcasters.remove(gameId);
gameJoinEventBroadcasters.remove(gameId);
});
}
private <T> Mono<T> assertUserExists(long userId, Mono<T> next) {
return userService.existsById(userId)
.filter(BooleanOperators.IDENTITY)
.switchIfEmpty(Mono.error(() -> new UserNotFoundException(userId)))
.then(next);
}
/// Check whether the given user has access to the given ongoing game. Also checks if it's the user's turn if specified by parameter.
///
/// @param gameId the ongoing game ID
/// @param userId the user ID
/// @param checkMoveOrder if `true`, verify that it's users turn
/// @return The `Publisher` will emit the ongoing game if all checks passed; otherwise an error
/// @throws UserNotFoundException if the user doesn't exist
/// @throws GameNotFoundException if the ongoing game doesn't exist
/// @throws NotParticipantException if the user is not participating in the game
/// @throws MoveOrderException if `checkMoveOrder == true` and it's not the user's turn in the game
private Mono<OngoingGame> checkGameAccess(UUID gameId, long userId, boolean checkMoveOrder) {
return Mono.defer(() -> {
// check if the user exists
Mono<OngoingGame> result = assertUserExists(userId, ongoingGameService.findById(gameId))
// check if the game exists
.switchIfEmpty(Mono.error(() -> new GameNotFoundException(gameId)))
// check that the user is a participant of the game
.filter(game -> userId == game.getWhiteUserId() || userId == game.getBlackUserId())
.switchIfEmpty(Mono.error(() -> new NotParticipantException(userId, gameId)));
if (checkMoveOrder) {
result = result
// check that it's their move now
.filter(game ->
userId == (game.getChessBoardAndMoveStack().getChessBoard().getCurrentlyMoving() == PieceColor.WHITE
? game.getWhiteUserId()
: game.getBlackUserId())
).switchIfEmpty(Mono.error(() -> new MoveOrderException(userId, gameId)));
}
return result;
});
}
private boolean canAcceptOrRefuseDraw(long userId, OngoingGame game) {
return userId == (game.getDrawOfferedBy() == PieceColor.WHITE
? game.getBlackUserId()
: game.getWhiteUserId());
}
private DefaultEmitFailureHandlingSinksMany<GameEvent> getOrCreateEventBroadcaster(UUID gameId) {
return ongoingGameEventBroadcasters.computeIfAbsent(gameId, _ -> new DefaultEmitFailureHandlingSinksMany<>(Sinks.many().replay().all(), CONTINUE_IF_NO_SUB));
}
private Sinks.One<GameJoinedEvent> getOrCreateJoinBroadcaster(UUID gameId) {
return gameJoinEventBroadcasters.computeIfAbsent(gameId, _ -> Sinks.one());
}
}
@@ -0,0 +1,45 @@
package ru.fokinatorr.chess.server.game;
import com.fasterxml.jackson.databind.JsonNode;
import lombok.*;
import lombok.experimental.FieldDefaults;
import org.springframework.data.annotation.Id;
import org.springframework.data.annotation.Transient;
import org.springframework.data.domain.Persistable;
import org.springframework.data.relational.core.mapping.Table;
import ru.fokinatorr.chess.impl.ChessBoard;
import ru.fokinatorr.chess.impl.PieceColor;
import ru.fokinatorr.chess.impl.misc.ChessBoardAndMoveStack;
import ru.fokinatorr.chess.server.entity.User;
import java.io.Serializable;
import java.util.UUID;
/// An ongoing chess game.
@Data
@Table(name = "ongoing_game")
@RequiredArgsConstructor
@NoArgsConstructor(access = AccessLevel.PROTECTED, force = true)
@FieldDefaults(level = AccessLevel.PRIVATE)
public class OngoingGame implements Serializable, Persistable<UUID> {
@Id
@NonNull
UUID id;
@Transient
@NonNull
transient User white;
long whiteUserId;
@Transient
@NonNull
transient User black;
long blackUserId;
@Transient
@NonNull
transient ChessBoardAndMoveStack chessBoardAndMoveStack;
@Transient
transient ChessBoard initialChessBoard;
JsonNode rawChessBoardAndMoveStack;
PieceColor drawOfferedBy;
@Transient
transient boolean isNew = true;
}
@@ -0,0 +1,35 @@
package ru.fokinatorr.chess.server.game;
import lombok.*;
import lombok.experimental.FieldDefaults;
import org.springframework.data.annotation.Id;
import org.springframework.data.annotation.Transient;
import org.springframework.data.domain.Persistable;
import org.springframework.data.relational.core.mapping.Table;
import ru.fokinatorr.chess.impl.PieceColorOrRandom;
import ru.fokinatorr.chess.server.entity.User;
import java.io.Serializable;
import java.util.UUID;
/// A pending chess game.
@Data
@Table(name = "pending_game")
@RequiredArgsConstructor
@NoArgsConstructor(access = AccessLevel.PROTECTED, force = true)
@FieldDefaults(level = AccessLevel.PRIVATE)
public class PendingGame implements Serializable, Persistable<UUID> {
@Id
@NonNull
UUID id;
@Transient
@NonNull
transient User withUser;
long withUserId;
/// Which color the user `withUser` will play.
@NonNull
PieceColorOrRandom ownerPlaysAs;
@Transient
transient boolean isNew = true;
}
@@ -0,0 +1,120 @@
package ru.fokinatorr.chess.server.misc;
import lombok.RequiredArgsConstructor;
import lombok.experimental.FieldDefaults;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Mono;
import ru.fokinatorr.chess.server.UserDTO;
import ru.fokinatorr.chess.server.data.FinishedGameService;
import ru.fokinatorr.chess.server.data.OngoingGameService;
import ru.fokinatorr.chess.server.data.PendingGameService;
import ru.fokinatorr.chess.server.data.UserService;
import ru.fokinatorr.chess.server.entity.User;
import ru.fokinatorr.chess.server.exception.GameNotFoundException;
import ru.fokinatorr.chess.server.exception.UserNotFoundException;
import ru.fokinatorr.chess.server.game.*;
import ru.fokinatorr.chess.server.service.DTOService;
import java.util.NoSuchElementException;
import java.util.function.Supplier;
import static lombok.AccessLevel.PRIVATE;
/// Server extension of `DTOService`.
@Service
@RequiredArgsConstructor
@FieldDefaults(makeFinal = true, level = PRIVATE)
public class DTOServiceExtension extends DTOService {
PendingGameService pendingGameService;
OngoingGameService ongoingGameService;
FinishedGameService finishedGameService;
UserService userService;
/// Creates and returns a DTO for the given `PendingGame` entity.
public PendingGameDTO toDTO(PendingGame source) {
return PendingGameDTO.builder()
.id(source.getId())
.withUser(toDTO(source.getWithUser()))
.ownerPlaysAs(source.getOwnerPlaysAs())
.build();
}
/// Either retrieves a `PendingGame` from the database, if there is one with such ID, or, if creating is allowed, creates a new
/// one, saves it, and returns, or else throws `NoSuchElementException`.
public Mono<PendingGame> fromDTO(PendingGameDTO dto, boolean allowCreating) {
return pendingGameService.findById(dto.id())
.switchIfEmpty(createOrError(pendingGameService.create(
dto.id(),
dto.withUser().id(),
dto.ownerPlaysAs()
)
.flatMap(pendingGameService::save), allowCreating, () -> noEntityMatching("PendingGame", dto)));
}
/// Creates and returns a DTO for the given `OngoingGame` entity.
public OngoingGameDTO toDTO(OngoingGame source) {
return OngoingGameDTO.builder()
.id(source.getId())
.white(toDTO(source.getWhite()))
.black(toDTO(source.getBlack()))
.chessBoardAndMoveStack(toDTO(source.getChessBoardAndMoveStack(), source.getInitialChessBoard()))
.build();
}
/// Either retrieves an `OngoingGame` from the database, if there is one with such ID, or, if creating is allowed, creates a new
/// one, saves it, and returns, or else throws `NoSuchElementException`.
public Mono<OngoingGame> fromDTO(OngoingGameDTO dto, boolean allowCreating) {
return ongoingGameService.findById(dto.id())
.switchIfEmpty(createOrError(ongoingGameService.create(
dto.id(),
dto.white().id(),
dto.black().id(),
fromDTO(dto.chessBoardAndMoveStack())
)
.flatMap(ongoingGameService::save), allowCreating, () -> noEntityMatching("OngoingGame", dto)));
}
/// Creates and returns a DTO for the given `FinishedGame` entity.
public FinishedGameDTO toDTO(FinishedGame source) {
return FinishedGameDTO.builder()
.id(source.getId())
.white(toDTO(source.getWhite()))
.black(toDTO(source.getBlack()))
.chessBoardAndMoveStack(toDTO(source.getChessBoardAndMoveStack(), source.getInitialChessBoard()))
.winner(source.getWinner())
.finishCause(source.getFinishCause())
.build();
}
/// Either retrieves a `FinishedGame` from the database, if there is one with such ID, or else throws `GameNotFoundException`.
public Mono<FinishedGame> fromDTO(FinishedGameDTO dto) {
return finishedGameService.findById(dto.id())
.switchIfEmpty(Mono.error(() -> new GameNotFoundException(dto.id())));
}
/// Creates and returns a DTO for the given `User` entity.
public UserDTO toDTO(User source) {
return UserDTO.builder()
.id(source.getId())
.username(source.getUsername())
.build();
}
/// Either retrieves a `User` from the database, if there is one with such ID, or throws a `UserNotFoundException`.
public Mono<User> fromDTO(UserDTO dto) {
return userService.findById(dto.id())
.switchIfEmpty(Mono.error(() -> new UserNotFoundException(dto.id())));
}
private static <S> Mono<S> createOrError(Mono<S> creator, boolean allowCreating, Supplier<? extends Throwable> exception) {
return Mono.defer(() -> allowCreating ? creator : Mono.error(exception));
}
private static RuntimeException noEntityMatching(String className, Object dto) {
return new NoSuchElementException("No " + className + " matching " + dto + " found");
}
}
@@ -0,0 +1,38 @@
package ru.fokinatorr.chess.server.misc;
import lombok.AccessLevel;
import lombok.NonNull;
import lombok.RequiredArgsConstructor;
import lombok.experimental.Delegate;
import lombok.experimental.FieldDefaults;
import reactor.core.publisher.Sinks;
import reactor.core.publisher.Sinks.EmitFailureHandler;
/// A [Sinks.Many] that delegates to another `Sinks.Many` and has a default [EmitFailureHandler].
@RequiredArgsConstructor
@FieldDefaults(makeFinal = true, level = AccessLevel.PRIVATE)
public class DefaultEmitFailureHandlingSinksMany<T> implements Sinks.Many<T> {
@NonNull
@Delegate
Sinks.Many<T> delegate;
@NonNull
EmitFailureHandler emitFailureHandler;
// Methods with default EmitFailureHandler
/// Calls [#emitNext(T, EmitFailureHandler)] with the default provided [EmitFailureHandler].
public void emitNextDefault(T t) {
emitNext(t, emitFailureHandler);
}
/// Calls [#emitError(Throwable, EmitFailureHandler)] with the default provided [EmitFailureHandler].
public void emitErrorDefault(Throwable error) {
emitError(error, emitFailureHandler);
}
/// Calls [#emitComplete(EmitFailureHandler)] with the default provided [EmitFailureHandler].
public void emitCompleteDefault() {
emitComplete(emitFailureHandler);
}
}
@@ -0,0 +1,43 @@
package ru.fokinatorr.chess.server.misc;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.AccessLevel;
import lombok.RequiredArgsConstructor;
import lombok.experimental.FieldDefaults;
import org.springframework.stereotype.Service;
import java.lang.reflect.RecordComponent;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
/// Converts a record class to a List of its attributes.
@Service
@FieldDefaults(makeFinal = true, level = AccessLevel.PRIVATE)
@RequiredArgsConstructor
public class RecordAttributesService {
private static final ConcurrentMap<Class<? extends Record>, RecordComponent[]> RECORD_COMPONENTS_CACHE = new ConcurrentHashMap<>();
ObjectMapper objectMapper;
/// Converts the given record class into a JSON list of attributes.
public <R extends Record> JsonNode toAttributes(Class<R> recordClass) {
ArrayNode attrs = new ArrayNode(objectMapper.getNodeFactory());
for (RecordComponent recordComponent : RECORD_COMPONENTS_CACHE.computeIfAbsent(recordClass, Class::getRecordComponents)) {
attrs.add(new ObjectNode(objectMapper.getNodeFactory())
.put("name", recordComponent.getName())
.put("type", recordComponent.getType().getName()));
}
return new ObjectNode(objectMapper.getNodeFactory())
.put("type", recordClass.getName())
.set("attributes", attrs);
}
// /// Clears the record component cache.
// public void clearCache() {
// RECORD_COMPONENTS_CACHE.clear();
// }
}
@@ -0,0 +1,36 @@
package ru.fokinatorr.chess.server.security;
import lombok.AccessLevel;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import lombok.experimental.FieldDefaults;
import org.springframework.security.core.GrantedAuthority;
@Getter
@RequiredArgsConstructor
@FieldDefaults(makeFinal = true, level = AccessLevel.PRIVATE)
public enum GrantedAuthorityImpl implements GrantedAuthority {
/// User with standard privileges
ROLE_USER(0),
/// Administrator
ROLE_ADMIN(1);
/// The ID of this `GrantedAuthorityImpl`.
int id;
/// {@inheritDoc}
@Override
public String getAuthority() {
return name();
}
/// Returns a `GrantedAuthorityImpl` instance with the given ID.
public static GrantedAuthorityImpl valueOf(int id) {
for (GrantedAuthorityImpl grantedAuthority : values()) {
if (grantedAuthority.id == id) {
return grantedAuthority;
}
}
throw new IllegalArgumentException("Unexpected value: " + id);
}
}
@@ -0,0 +1,61 @@
package ru.fokinatorr.chess.server.security;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.AccessLevel;
import lombok.NonNull;
import lombok.RequiredArgsConstructor;
import lombok.Setter;
import lombok.experimental.FieldDefaults;
import org.springframework.core.io.buffer.DataBufferUtils;
import org.springframework.security.authentication.UsernamePasswordAuthenticationToken;
import org.springframework.security.core.Authentication;
import org.springframework.security.web.server.authentication.ServerAuthenticationConverter;
import org.springframework.stereotype.Component;
import org.springframework.web.server.ServerWebExchange;
import reactor.core.Exceptions;
import reactor.core.publisher.Mono;
import java.nio.charset.StandardCharsets;
import java.util.Map;
import java.util.function.Function;
/// Creates [Authentication] object from a JSON form.
@Component
@FieldDefaults(level = AccessLevel.PRIVATE)
@RequiredArgsConstructor
@Setter
public class ServerJsonAuthenticationConverter implements ServerAuthenticationConverter, Function<ServerWebExchange, Mono<Authentication>> {
@NonNull final ObjectMapper objectMapper;
@NonNull String usernameParameter = "username";
@NonNull String passwordParameter = "password";
@Override
public Mono<Authentication> convert(ServerWebExchange exchange) {
return Mono.justOrEmpty(exchange.getRequest().getHeaders().getContentType())
.filter(contentType -> contentType.getSubtype().equals("json"))
.flatMap(_ -> DataBufferUtils.join(exchange.getRequest().getBody()))
.map(buf -> {
byte[] data = new byte[buf.readableByteCount()];
buf.read(data);
return new String(data, StandardCharsets.UTF_8);
})
.map(s -> {
Map<String, String> json;
try {
json = (Map<String, String>) objectMapper.readValue(s, Map.class);
} catch (JsonProcessingException e) {
throw Exceptions.propagate(e);
}
String username = json.get(usernameParameter);
String password = json.get(passwordParameter);
return UsernamePasswordAuthenticationToken.unauthenticated(username, password);
});
}
@Override
public Mono<Authentication> apply(ServerWebExchange serverWebExchange) {
return convert(serverWebExchange);
}
}
+23
View File
@@ -0,0 +1,23 @@
spring:
application:
name: Ch4ss Server
rsocket:
server:
transport: websocket
mapping-path: /ws
r2dbc:
url: r2dbc:postgres://localhost:5432/chessserverdb
username: postgres
password: 3h6F$t3_6e&`0L//
properties:
globally-quoted-identifiers: true
server:
error:
whitelabel:
enabled: false
logging:
level:
web: trace
io.r2dbc.postgresql:
QUERY: DEBUG
PARAM: DEBUG
+90
View File
@@ -0,0 +1,90 @@
<databaseChangeLog xmlns="http://www.liquibase.org/xml/ns/dbchangelog"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://www.liquibase.org/xml/ns/dbchangelog https://www.liquibase.org/xml/ns/dbchangelog/dbchangelog-latest.xsd">
<changeSet id="20250426_184445" author="Fokinatorr_">
<!-- Added entity User -->
<createTable tableName="users">
<column name="id" type="bigint" autoIncrement="true">
<constraints primaryKey="true" nullable="false"/>
</column>
<column name="username" type="varchar(50)">
<constraints nullable="false" unique="true"/>
</column>
<column name="password" type="varchar(72)">
<constraints nullable="false"/>
</column>
<column name="authority_ids" type="integer[]">
<constraints nullable="false"/>
</column>
</createTable>
<!-- Added entity UserNotification -->
<createTable tableName="user_notification">
<column name="id" type="bigint" autoIncrement="true">
<constraints primaryKey="true" nullable="false"/>
</column>
<column name="type" type="varchar(30)">
<constraints nullable="false"/>
</column>
<column name="sender_id" type="bigint">
<constraints nullable="false" foreignKeyName="fk_sender_user_id" references="users(id)"/>
</column>
<column name="receiver_id" type="bigint">
<constraints nullable="false" foreignKeyName="fk_receiver_user_id" references="users(id)"/>
</column>
<column name="additional_data" type="jsonb">
<constraints nullable="false"/>
</column>
</createTable>
<!-- Added entity OngoingGame -->
<createTable tableName="ongoing_game">
<column name="id" type="uuid">
<constraints primaryKey="true" nullable="false"/>
</column>
<column name="white_user_id" type="bigint">
<constraints nullable="false" foreignKeyName="fk_white_user_id" references="users(id)"/>
</column>
<column name="black_user_id" type="bigint">
<constraints nullable="false" foreignKeyName="fk_black_user_id" references="users(id)"/>
</column>
<column name="draw_offered_by" type="varchar(5)">
<constraints nullable="true"/>
</column>
<column name="raw_chess_board_and_move_stack" type="jsonb">
<constraints nullable="false"/>
</column>
</createTable>
<!-- Added entity PendingGame -->
<createTable tableName="pending_game">
<column name="id" type="uuid">
<constraints primaryKey="true" nullable="false"/>
</column>
<column name="with_user_id" type="bigint">
<constraints nullable="false" foreignKeyName="fk_white_user_id" references="users(id)"/>
</column>
<column name="owner_plays_as" type="varchar(6)">
<constraints nullable="false"/>
</column>
</createTable>
<!-- Added entity FinishedGame -->
<createTable tableName="finished_game">
<column name="id" type="uuid">
<constraints primaryKey="true" nullable="false"/>
</column>
<column name="white_user_id" type="bigint">
<constraints nullable="false" foreignKeyName="fk_white_user_id" references="users(id)"/>
</column>
<column name="black_user_id" type="bigint">
<constraints nullable="false" foreignKeyName="fk_black_user_id" references="users(id)"/>
</column>
<column name="raw_chess_board_and_move_stack" type="jsonb">
<constraints nullable="false"/>
</column>
<column name="winner" type="varchar(5)">
<constraints nullable="true"/>
</column>
<column name="finish_cause" type="varchar(15)">
<constraints nullable="false"/>
</column>
</createTable>
</changeSet>
</databaseChangeLog>
@@ -0,0 +1,13 @@
package ru.fokinatorr.chess.server;
import org.junit.jupiter.api.Test;
import org.springframework.boot.test.context.SpringBootTest;
@SpringBootTest
class Ch4ssServerApplicationTests {
@Test
void contextLoads() {
}
}
@@ -0,0 +1,322 @@
package ru.fokinatorr.chess.server;
import lombok.extern.slf4j.Slf4j;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.boot.test.web.server.LocalServerPort;
import org.springframework.context.ApplicationContext;
import org.springframework.core.io.ClassPathResource;
import org.springframework.messaging.rsocket.RSocketRequester;
import org.springframework.messaging.rsocket.RSocketStrategies;
import org.springframework.security.authentication.AbstractAuthenticationToken;
import org.springframework.security.core.Authentication;
import org.springframework.security.core.context.SecurityContextHolder;
import org.springframework.security.crypto.password.PasswordEncoder;
import org.springframework.security.test.context.support.WithMockUser;
import org.springframework.security.test.web.reactive.server.SecurityMockServerConfigurers;
import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.web.reactive.server.WebTestClient;
import org.springframework.test.web.reactive.server.WebTestClientConfigurer;
import org.springframework.util.StreamUtils;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
import ru.fokinatorr.chess.impl.ChessBoard;
import ru.fokinatorr.chess.impl.PieceColor;
import ru.fokinatorr.chess.impl.PieceColorOrRandom;
import ru.fokinatorr.chess.impl.Position;
import ru.fokinatorr.chess.impl.misc.ChessBoardAndMoveStack;
import ru.fokinatorr.chess.impl.movestack.SimpleMove;
import ru.fokinatorr.chess.server.data.UserService;
import ru.fokinatorr.chess.server.entity.User;
import ru.fokinatorr.chess.server.game.*;
import ru.fokinatorr.chess.server.game.event.GameEvent;
import ru.fokinatorr.chess.server.game.event.GameFinishedEvent;
import ru.fokinatorr.chess.server.game.event.GameJoinedEvent;
import ru.fokinatorr.chess.server.game.event.MovePerformedEvent;
import ru.fokinatorr.chess.server.misc.DTOServiceExtension;
import ru.fokinatorr.chess.server.security.GrantedAuthorityImpl;
import java.io.IOException;
import java.net.URI;
import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.function.Consumer;
import java.util.function.Predicate;
import static org.junit.jupiter.api.Assertions.*;
import static org.springframework.security.test.web.reactive.server.SecurityMockServerConfigurers.springSecurity;
@Slf4j
@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
@DirtiesContext
class GameControllerTests {
@Autowired
ApplicationContext applicationContext;
@Autowired
UserService userService;
@Autowired
GameService gameService;
RSocketRequester.Builder testRequesterBuilder;
@Autowired
PasswordEncoder passwordEncoder;
@Autowired
RSocketStrategies rsocketStrategies;
@Autowired
DTOServiceExtension dtoService;
WebTestClient rest;
@LocalServerPort
int port;
@Value("${spring.rsocket.server.mapping-path}")
String basePath;
@BeforeEach
void init() {
rest = WebTestClient.bindToApplicationContext(applicationContext)
.apply(springSecurity())
.configureClient()
.build();
testRequesterBuilder = RSocketRequester.builder().rsocketStrategies(rsocketStrategies);
}
RSocketRequester createRequester() {
return testRequesterBuilder.websocket(URI.create("ws://localhost:" + port + basePath));
}
@Test
void shouldReturnUnauthorized() {
rest.get()
.uri("/game/list/pending")
.exchange()
.expectStatus().isUnauthorized();
}
@Test
@WithMockUser
void shouldReturnOkAndEmptyArray() {
rest.get()
.uri("/game/list/pending")
.exchange()
.expectStatus().isOk()
.expectBody().json("[]");
}
@Test
@WithMockUser
void testPendingGameCreation() throws IOException {
String createGameOptionsJson = StreamUtils.copyToString(new ClassPathResource("/create_game_options.json").getInputStream(), StandardCharsets.UTF_8);
rest.get()
.uri("/game/create")
.exchange()
.expectStatus().isOk()
.expectBody().json(createGameOptionsJson);
UUID newGameId = rest
.mutateWith(injectPrincipal("test", "pass123"))
.post()
.uri("/game/create")
.body(Mono.just(new CreateGameOptions(PieceColorOrRandom.WHITE)), CreateGameOptions.class)
.exchange()
.expectStatus().isCreated()
.expectBody(UUID.class).returnResult().getResponseBody();
log.info("New pending game ID: {}", newGameId);
rest.get()
.uri("/game/list/pending")
.exchange()
.expectStatus().isOk()
.returnResult(PendingGameDTO.class).getResponseBody()
.as(StepVerifier::create)
.consumeNextWith(pendingGameDTO -> {
assertEquals(newGameId, pendingGameDTO.id());
assertEquals("test", pendingGameDTO.withUser().username());
assertEquals(PieceColorOrRandom.WHITE, pendingGameDTO.ownerPlaysAs());
})
.verifyComplete();
rest.mutateWith(injectPrincipal("test2", "pass321")).delete() // a user that is not the game owner
.uri("/game/cancel?id={newGameId}", newGameId)
.exchange()
.expectStatus().isForbidden();
rest.mutateWith(injectPrincipal("test", "pass123")).delete()
.uri("/game/cancel?id={newGameId}", newGameId)
.exchange()
.expectStatus().isNoContent()
.expectBody().isEmpty();
rest.get()
.uri("/game/list/pending")
.exchange()
.expectStatus().isOk()
.expectBody().json("[]");
}
@Test
@WithMockUser
void testOngoingGame() throws InterruptedException {
UUID newGameId = rest
.mutateWith(injectPrincipal("test", "pass123"))
.post()
.uri("/game/create")
.body(Mono.just(new CreateGameOptions(PieceColorOrRandom.BLACK)), CreateGameOptions.class)
.exchange()
.expectStatus().isCreated()
.expectBody(UUID.class).returnResult().getResponseBody();
assertNotNull(newGameId);
log.info("New pending game ID (ongoing game test): {}", newGameId);
CountDownLatch latch = new CountDownLatch(1);
Thread.ofVirtual()
.uncaughtExceptionHandler(Thread.currentThread().getThreadGroup())
.name("Join notifier test")
.start(() -> {
try {
RSocketRequester requester = createRequester();
requester.route("game.waitForJoin")
.data(Mono.just(newGameId), UUID.class)
.retrieveMono(GameJoinedEvent.class)
.doOnTerminate(requester::dispose)
.as(StepVerifier::create)
.consumeNextWith(e -> {
assertEquals(ChessBoard.startPosition(), e.startPosition());
assertEquals(PieceColor.BLACK, e.playingAs());
assertEquals("test2", e.playingWith().username());
})
.verifyComplete();
} finally {
latch.countDown();
}
});
rest.mutateWith(injectPrincipal("test2", "pass321"))
.post()
.uri("/game/join?id={newGameId}", newGameId)
.exchange()
.expectStatus().isNoContent()
.expectBody().isEmpty();
latch.await();
rest.get()
.uri("/game/list/pending")
.exchange()
.expectStatus().isOk()
.expectBody().json("[]");
Thread.sleep(3000);
rest.get()
.uri("/game/list/ongoing")
.exchange()
.expectStatus().isOk()
.expectBodyList(OngoingGameDTO.class).consumeWith(entityExchangeResult -> {
List<OngoingGameDTO> ongoingGames = entityExchangeResult.getResponseBody();
assertNotNull(ongoingGames);
assertEquals(1, ongoingGames.size());
OngoingGameDTO ongoingGame = ongoingGames.getFirst();
assertEquals(newGameId, ongoingGame.id());
assertEquals("test2", ongoingGame.white().username());
assertEquals("test", ongoingGame.black().username());
assertEquals(dtoService.toDTO(ChessBoardAndMoveStack.startPosition(), ChessBoard.startPosition()), ongoingGame.chessBoardAndMoveStack());
});
RSocketRequester requester = createRequester();
Flux<GameEvent> gameEventPub = requester.route("game.join")
.data(Mono.just(newGameId), UUID.class)
.retrieveFlux(GameEvent.class)
.doOnTerminate(requester::dispose);
// Play Fool's mate for shorter results
playMove(newGameId, "g2g4", "test2", "pass321");
playMove(newGameId, "e7e5", "test", "pass123");
playMove(newGameId, "f2f3", "test2", "pass321");
playMove(newGameId, "d8h4", "test", "pass123");
StepVerifier.create(gameEventPub)
.expectNextMatches(movePerformed("g2g4"))
.expectNextMatches(movePerformed("e7e5"))
.expectNextMatches(movePerformed("f2f3"))
.expectNextMatches(movePerformed("d8h4"))
.expectNextMatches(gameFinished(PieceColor.BLACK, GameFinishedEvent.GameFinishCause.CHECKMATE))
.verifyComplete();
rest.get()
.uri("/game/list/ongoing")
.exchange()
.expectStatus().isOk()
.expectBody().json("[]");
rest.get()
.uri("/game/list/finished")
.exchange()
.expectStatus().isOk()
.expectBodyList(FinishedGameDTO.class).consumeWith(entityExchangeResult -> {
List<FinishedGameDTO> finishedGames = entityExchangeResult.getResponseBody();
assertNotNull(finishedGames);
assertEquals(1, finishedGames.size());
FinishedGameDTO finishedGame = finishedGames.getFirst();
assertEquals(newGameId, finishedGame.id());
assertEquals("test2", finishedGame.white().username());
assertEquals("test", finishedGame.black().username());
assertEquals(dtoService.toDTO(make(ChessBoardAndMoveStack.startPosition(), cbams -> {
cbams.perform(new SimpleMove(Position.valueOf("g2"), Position.valueOf("g4")));
cbams.perform(new SimpleMove(Position.valueOf("e7"), Position.valueOf("e5")));
cbams.perform(new SimpleMove(Position.valueOf("f2"), Position.valueOf("f3")));
cbams.perform(new SimpleMove(Position.valueOf("d8"), Position.valueOf("h4")));
}), ChessBoard.startPosition()), finishedGame.chessBoardAndMoveStack());
assertSame(PieceColor.BLACK, finishedGame.winner());
assertSame(GameFinishedEvent.GameFinishCause.CHECKMATE, finishedGame.finishCause());
});
}
private void playMove(UUID gameId, String move, String username, String password) {
rest.mutateWith(injectPrincipal(username, password))
.put()
.uri("/game/move/perform?gameId={gameId}&move={move}", gameId, move)
.exchange()
.expectStatus().isNoContent()
.expectBody().isEmpty();
}
private static Predicate<? super GameEvent> movePerformed(String theMove) {
return gameEvent -> gameEvent instanceof MovePerformedEvent(
String move, ChessBoardAndMoveStackDTO _
) && move.equals(theMove);
}
private static Predicate<? super GameEvent> gameFinished(PieceColor theWinner, GameFinishedEvent.GameFinishCause theFinishCause) {
return gameEvent -> gameEvent instanceof GameFinishedEvent(
PieceColor winner, GameFinishedEvent.GameFinishCause cause,
ChessBoardAndMoveStackDTO _
) && winner == theWinner && cause == theFinishCause;
}
private static <T> T make(T initial, Consumer<? super T> customizer) {
customizer.accept(initial);
return initial;
}
private WebTestClientConfigurer injectPrincipal(String username, String password) {
User principal = userService.existsByUsername(username)
.flatMap(b -> {
if (b) {
return userService.findByUsername(username);
} else {
return userService.create(username, password, false).flatMap(userService::save);
}
})
.block();
assertNotNull(principal);
assertEquals(username, principal.getUsername());
assertTrue(passwordEncoder.matches(password, principal.getPassword()));
assertEquals(Set.of(GrantedAuthorityImpl.ROLE_USER), principal.getAuthorities());
Authentication authentication = new AbstractAuthenticationToken(principal.getAuthorities()) {
@Override
public Object getCredentials() {
return principal.getPassword();
}
@Override
public Object getPrincipal() {
return principal;
}
};
authentication.setAuthenticated(true);
SecurityContextHolder.getContext().setAuthentication(authentication);
return SecurityMockServerConfigurers.mockAuthentication(authentication);
}
}
@@ -0,0 +1,52 @@
package ru.fokinatorr.chess.server;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import ru.fokinatorr.chess.impl.ChessBoard;
import ru.fokinatorr.chess.impl.PieceColor;
import ru.fokinatorr.chess.impl.Position;
import ru.fokinatorr.chess.impl.piece.*;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
@Slf4j
@SpringBootTest
class ObjectMapperTest {
@Autowired
private ObjectMapper objectMapper;
@Test
void writesPosition() {
for (Position.Column column : Position.Column.values()) {
for (Position.Row row : Position.Row.values()) {
commonWrites(Position.valueOf(column, row), Position.class);
}
}
}
@Test
void writesPiece() {
commonWrites(new Pawn(PieceColor.WHITE), Piece.class);
commonWrites(new Knight(PieceColor.BLACK), Piece.class);
commonWrites(new King(PieceColor.BLACK), Piece.class);
commonWrites(new Queen(PieceColor.WHITE), Piece.class);
}
@Test
void writesChessBoard() {
commonWrites(ChessBoard.startPosition(), ChessBoard.class);
}
<T> void commonWrites(T initial, Class<T> tClass) {
String tAsJson = assertDoesNotThrow(() -> objectMapper.writeValueAsString(initial));
log.info("JSON {}: {}", tClass.getSimpleName(), tAsJson);
T decoded = assertDoesNotThrow(() -> objectMapper.readValue(tAsJson, tClass));
assertEquals(initial, decoded);
}
}
@@ -0,0 +1,9 @@
{
"type": "ru.fokinatorr.chess.server.game.CreateGameOptions",
"attributes": [
{
"name": "playAs",
"type": "ru.fokinatorr.chess.impl.PieceColorOrRandom"
}
]
}