How Quarkus transactional works with reactive sql client

Viewed 1267

Currently i'm trying to build some simple prototype application with Quarkus (1.3.1). I'm trying to use a reactive approach with reactive database connection to PostgreSQL using https://quarkus.io/guides/reactive-sql-clients. Here I'm using kotlin and mutiny to simplify things a bit.

I have the following controller:

package me.invitation.infrastructure.rest

import io.smallrye.mutiny.Uni
import me.invitation.application.service.InvitationService
import me.invitation.infrastructure.rest.api.InvitationDto
import me.invitation.infrastructure.rest.system.MicroserviceContentType
import org.eclipse.microprofile.context.ThreadContext
import org.slf4j.LoggerFactory
import javax.inject.Inject
import javax.ws.rs.*
import javax.ws.rs.core.Response


@Path("/invitation")
class Resource @Inject constructor(
        private val invitationService: InvitationService,
        private val threadContext: ThreadContext) {


    @POST
    @Consumes(MicroserviceContentType.PUBLIC_RSC_JSON)
    fun postInvitation(invitation: InvitationDto): Uni<Response> {
        return Uni.createFrom()
                .completionStage(
                        threadContext.withContextCapture(
                                invitationService.sendInvite(invitation.email)))
                .map { Response.status(Response.Status.CREATED).build() }
    }

    companion object {
        private val log = LoggerFactory.getLogger(Resource::class.java)
    }
}

The InvitationService implementation is as follows:

package me.invitation.application.service

import io.smallrye.mutiny.Uni
import me.invitation.application.registry.InvitationRegistry
import me.invitation.application.service.mail.MailService
import org.apache.commons.codec.binary.Base32
import org.eclipse.microprofile.config.inject.ConfigProperty
import org.eclipse.microprofile.context.ManagedExecutor
import org.eclipse.microprofile.context.ThreadContext
import org.slf4j.LoggerFactory
import java.security.NoSuchAlgorithmException
import java.security.SecureRandom
import java.time.Duration
import java.time.LocalDateTime
import java.util.*
import java.util.concurrent.CompletableFuture
import java.util.function.Supplier
import javax.enterprise.context.ApplicationScoped
import javax.inject.Inject
import javax.transaction.Transactional

@ApplicationScoped
class ExternalRegistryInvitationService @Inject constructor(
        private val invitationRepository: InvitationRegistry,
        private val mailService: MailService,
        private val threadContext: ThreadContext,
        private val managedExecutor: ManagedExecutor,
        @ConfigProperty(name ="invitation.expiration-ttl-seconds") private val invitationTTLSeconds: Long
) : InvitationService {

    @Transactional
    override fun sendInvite(email: String): CompletableFuture<Unit> {
        return Uni.createFrom()
                  .completionStage(threadContext.withContextCapture(invitationRepository.save(id, email)))
                  .flatMap {wasSaved ->
                                if(wasSaved) {
                                   Uni.createFrom().completionStage(threadContext.withContextCapture(mailService.sendInvitationEmail(email)))
                                } else {
                                    throw InvitationAlreadyPresentException()
                                }
                            }.subscribeAsCompletionStage()
    }
}

Here I'm trying to save a record in database and after that send an email

InvitationRegistry implementation is as follows:

package me.invitation.infrastructure.db

import io.smallrye.mutiny.Uni
import io.vertx.mutiny.pgclient.PgPool
import io.vertx.mutiny.sqlclient.Tuple
import me.invitation.application.registry.InvitationRegistry
import java.time.LocalDateTime
import java.util.concurrent.CompletableFuture
import javax.enterprise.context.ApplicationScoped
import javax.inject.Inject
import javax.transaction.Transactional

@ApplicationScoped
class PostgresInvitationRegistry @Inject constructor(private val client: PgPool) : InvitationRegistry {

    private final val insertQuery = "INSERT INTO invitation(id, email, aud_ts_created, aud_ts_last_modified) VALUES ($1, $2, $3, $4) ON CONFLICT DO NOTHING RETURNING id"

    @Transactional(value = Transactional.TxType.MANDATORY)
    override fun save(id: String, email: String): CompletableFuture<Boolean> {
        return client
                .preparedQuery(insertQuery, Tuple.of(id, email, LocalDateTime.now(), LocalDateTime.now()))
                .flatMap { result -> Uni.createFrom().item(result.rowCount() > 0) }
                .subscribeAsCompletionStage()
    }
}

MailService implementation is irrelevant because everything it does is throwing RuntimeException.

Here I'm expecting that top-level transaction to rollback. Because MailService.sendInvitationEmail is throwing an exception. But that's not the case. Even in case of exception the record is persisted. I have read this guide https://quarkus.io/guides/transaction#reactive-extensions and https://quarkus.io/guides/context-propagation#usage-example-for-completionstage and I thought that declarative transactions works well with reactive sql in Quarkus.

I think that there is an autocommit setting for reactive connection which is set to true. But didn't find anything related to autocommit, neither in docs, nor in the sources of quarkus and its extensions. Suprisingly, there is even no autocommit setting even for quarkus main jdbc pool = agroal!

Myabe someone knows how to make declarative transactions work with reactive sql client in Quarkus

0 Answers
Related