From ff26589147c86cbed1b83c95044b421f96c125d1 Mon Sep 17 00:00:00 2001 From: Egor Cherniak Date: Tue, 28 Jul 2026 12:50:31 +0300 Subject: [PATCH 1/2] add withdrawal cash changed support --- pom.xml | 4 +- .../iface/WithdrawalCashChangeDao.java | 13 ++++ .../impl/WithdrawalCashChangeDaoImpl.java | 63 +++++++++++++++++++ .../WithdrawalBodyChangedHandler.java | 62 ++++++++++++++++++ .../V54__add_withdrawal_cash_change.sql | 18 ++++++ src/test/java/dev/vality/daway/TestData.java | 18 ++++++ .../WithdrawalBodyChangedHandlerTest.java | 52 +++++++++++++++ 7 files changed, 228 insertions(+), 2 deletions(-) create mode 100644 src/main/java/dev/vality/daway/dao/withdrawal/iface/WithdrawalCashChangeDao.java create mode 100644 src/main/java/dev/vality/daway/dao/withdrawal/impl/WithdrawalCashChangeDaoImpl.java create mode 100644 src/main/java/dev/vality/daway/handler/event/stock/impl/withdrawal/WithdrawalBodyChangedHandler.java create mode 100644 src/main/resources/db/migration/V54__add_withdrawal_cash_change.sql create mode 100644 src/test/java/dev/vality/daway/handler/event/stock/impl/withdrawal/WithdrawalBodyChangedHandlerTest.java diff --git a/pom.xml b/pom.xml index 208ce5a8..cfa183c1 100644 --- a/pom.xml +++ b/pom.xml @@ -177,12 +177,12 @@ dev.vality damsel - 1.696-b07c077 + 1.698-99c91c9 dev.vality fistful-proto - 1.188-f7ce08e + 1.191-94b75ce dev.vality diff --git a/src/main/java/dev/vality/daway/dao/withdrawal/iface/WithdrawalCashChangeDao.java b/src/main/java/dev/vality/daway/dao/withdrawal/iface/WithdrawalCashChangeDao.java new file mode 100644 index 00000000..f5c87354 --- /dev/null +++ b/src/main/java/dev/vality/daway/dao/withdrawal/iface/WithdrawalCashChangeDao.java @@ -0,0 +1,13 @@ +package dev.vality.daway.dao.withdrawal.iface; + +import dev.vality.dao.GenericDao; +import dev.vality.daway.domain.tables.pojos.WithdrawalCashChange; +import dev.vality.daway.exception.DaoException; + +public interface WithdrawalCashChangeDao extends GenericDao { + + void save(WithdrawalCashChange withdrawalCashChange) throws DaoException; + + void switchCurrent(String withdrawalId) throws DaoException; + +} diff --git a/src/main/java/dev/vality/daway/dao/withdrawal/impl/WithdrawalCashChangeDaoImpl.java b/src/main/java/dev/vality/daway/dao/withdrawal/impl/WithdrawalCashChangeDaoImpl.java new file mode 100644 index 00000000..298ef9b9 --- /dev/null +++ b/src/main/java/dev/vality/daway/dao/withdrawal/impl/WithdrawalCashChangeDaoImpl.java @@ -0,0 +1,63 @@ +package dev.vality.daway.dao.withdrawal.impl; + +import dev.vality.dao.impl.AbstractGenericDao; +import dev.vality.daway.dao.withdrawal.iface.WithdrawalCashChangeDao; +import dev.vality.daway.domain.tables.pojos.WithdrawalCashChange; +import dev.vality.daway.domain.tables.records.WithdrawalCashChangeRecord; +import dev.vality.daway.exception.DaoException; +import org.jooq.Query; +import org.jooq.impl.DSL; +import org.springframework.stereotype.Component; + +import javax.sql.DataSource; + +import static dev.vality.daway.domain.Tables.WITHDRAWAL_CASH_CHANGE; + +@Component +public class WithdrawalCashChangeDaoImpl extends AbstractGenericDao implements WithdrawalCashChangeDao { + + public WithdrawalCashChangeDaoImpl(DataSource dataSource) { + super(dataSource); + } + + @Override + public void save(WithdrawalCashChange withdrawalCashChange) throws DaoException { + WithdrawalCashChangeRecord record = getDslContext().newRecord(WITHDRAWAL_CASH_CHANGE, withdrawalCashChange); + execute(prepareInsertQuery(record)); + } + + private Query prepareInsertQuery(WithdrawalCashChangeRecord record) { + return getDslContext().insertInto(WITHDRAWAL_CASH_CHANGE) + .set(record) + .onConflict( + WITHDRAWAL_CASH_CHANGE.WITHDRAWAL_ID, + WITHDRAWAL_CASH_CHANGE.SEQUENCE_ID + ) + .doNothing(); + } + + @Override + public void switchCurrent(String withdrawalId) throws DaoException { + setOldWithdrawalCashChangeNotCurrent(withdrawalId); + setLatestWithdrawalCashChangeCurrent(withdrawalId); + } + + private void setOldWithdrawalCashChangeNotCurrent(String withdrawalId) { + execute(getDslContext().update(WITHDRAWAL_CASH_CHANGE) + .set(WITHDRAWAL_CASH_CHANGE.CURRENT, false) + .where(WITHDRAWAL_CASH_CHANGE.WITHDRAWAL_ID.eq(withdrawalId) + .and(WITHDRAWAL_CASH_CHANGE.CURRENT)) + ); + } + + private void setLatestWithdrawalCashChangeCurrent(String withdrawalId) { + execute(getDslContext().update(WITHDRAWAL_CASH_CHANGE) + .set(WITHDRAWAL_CASH_CHANGE.CURRENT, true) + .where(WITHDRAWAL_CASH_CHANGE.ID.eq( + DSL.select(DSL.max(WITHDRAWAL_CASH_CHANGE.ID)) + .from(WITHDRAWAL_CASH_CHANGE) + .where(WITHDRAWAL_CASH_CHANGE.WITHDRAWAL_ID.eq(withdrawalId)) + )) + ); + } +} diff --git a/src/main/java/dev/vality/daway/handler/event/stock/impl/withdrawal/WithdrawalBodyChangedHandler.java b/src/main/java/dev/vality/daway/handler/event/stock/impl/withdrawal/WithdrawalBodyChangedHandler.java new file mode 100644 index 00000000..0c28dd52 --- /dev/null +++ b/src/main/java/dev/vality/daway/handler/event/stock/impl/withdrawal/WithdrawalBodyChangedHandler.java @@ -0,0 +1,62 @@ +package dev.vality.daway.handler.event.stock.impl.withdrawal; + +import dev.vality.daway.dao.withdrawal.iface.WithdrawalCashChangeDao; +import dev.vality.daway.domain.tables.pojos.WithdrawalCashChange; +import dev.vality.fistful.base.Cash; +import dev.vality.fistful.withdrawal.BodyChange; +import dev.vality.fistful.withdrawal.Change; +import dev.vality.fistful.withdrawal.TimestampedChange; +import dev.vality.geck.common.util.TypeUtil; +import dev.vality.geck.filter.Filter; +import dev.vality.geck.filter.PathConditionFilter; +import dev.vality.geck.filter.condition.IsNullCondition; +import dev.vality.geck.filter.rule.PathConditionRule; +import dev.vality.machinegun.eventsink.MachineEvent; +import lombok.Getter; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; + +@Slf4j +@Component +@RequiredArgsConstructor +public class WithdrawalBodyChangedHandler implements WithdrawalHandler { + + private final WithdrawalCashChangeDao withdrawalCashChangeDao; + + @Getter + private final Filter filter = new PathConditionFilter( + new PathConditionRule("change.body_changed", new IsNullCondition().not())); + + @Override + @Transactional(propagation = Propagation.REQUIRED) + public void handle(TimestampedChange timestampedChange, MachineEvent event) { + Change change = timestampedChange.getChange(); + long sequenceId = event.getEventId(); + String withdrawalId = event.getSourceId(); + log.info("Start withdrawal body changed handling, sequenceId={}, withdrawalId={}", sequenceId, withdrawalId); + + BodyChange bodyChange = change.getBodyChanged(); + WithdrawalCashChange cashChange = new WithdrawalCashChange(); + cashChange.setWtime(null); + cashChange.setId(null); + cashChange.setSequenceId(sequenceId); + cashChange.setWithdrawalId(withdrawalId); + cashChange.setCurrent(true); + cashChange.setEventCreatedAt(TypeUtil.stringToLocalDateTime(event.getCreatedAt())); + + Cash newCash = bodyChange.getNewBody(); + cashChange.setNewAmount(newCash.getAmount()); + cashChange.setNewCurrencyCode(newCash.getCurrency().getSymbolicCode()); + Cash oldCash = bodyChange.getOldBody(); + cashChange.setOldAmount(oldCash.getAmount()); + cashChange.setOldCurrencyCode(oldCash.getCurrency().getSymbolicCode()); + + withdrawalCashChangeDao.save(cashChange); + withdrawalCashChangeDao.switchCurrent(withdrawalId); + log.info("Withdrawal body change has been saved, sequenceId={}, withdrawalId={}", sequenceId, withdrawalId); + } + +} diff --git a/src/main/resources/db/migration/V54__add_withdrawal_cash_change.sql b/src/main/resources/db/migration/V54__add_withdrawal_cash_change.sql new file mode 100644 index 00000000..f967e5cd --- /dev/null +++ b/src/main/resources/db/migration/V54__add_withdrawal_cash_change.sql @@ -0,0 +1,18 @@ +CREATE TABLE dw.withdrawal_cash_change +( + id bigserial NOT NULL, + event_created_at timestamp without time zone NOT NULL, + withdrawal_id character varying NOT NULL, + + new_amount bigint NOT NULL, + new_currency_code character varying NOT NULL, + old_amount bigint NOT NULL, + old_currency_code character varying NOT NULL, + + current BOOLEAN NOT NULL DEFAULT false, + wtime timestamp without time zone NOT NULL DEFAULT (now() AT TIME ZONE 'utc'::text), + sequence_id bigint, + + CONSTRAINT withdrawal_cash_change_pkey PRIMARY KEY (id), + CONSTRAINT withdrawal_cash_change_uniq UNIQUE (withdrawal_id, sequence_id) +); diff --git a/src/test/java/dev/vality/daway/TestData.java b/src/test/java/dev/vality/daway/TestData.java index 105c6135..69d5aa25 100644 --- a/src/test/java/dev/vality/daway/TestData.java +++ b/src/test/java/dev/vality/daway/TestData.java @@ -625,6 +625,24 @@ public static TimestampedChange createWithdrawalCreatedChange(String id) { return timestampedChange; } + public static TimestampedChange createWithdrawalBodyChangedChange() { + BodyChange bodyChange = new BodyChange() + .setOldBody(new dev.vality.fistful.base.Cash() + .setAmount(100L) + .setCurrency(new dev.vality.fistful.base.CurrencyRef() + .setSymbolicCode("RUB"))) + .setNewBody(new dev.vality.fistful.base.Cash() + .setAmount(200L) + .setCurrency(new dev.vality.fistful.base.CurrencyRef() + .setSymbolicCode("USD"))); + Change change = new Change(); + change.setBodyChanged(bodyChange); + TimestampedChange timestampedChange = new TimestampedChange(); + timestampedChange.setOccuredAt(OCCURED_AT); + timestampedChange.setChange(change); + return timestampedChange; + } + public static MachineEvent createInvoice(InvoicePaymentChangePayload invoicePaymentChangePayload) { PaymentEventPayloadSerializer paymentEventPayloadSerializer = new PaymentEventPayloadSerializer(); MachineEvent message = new MachineEvent(); diff --git a/src/test/java/dev/vality/daway/handler/event/stock/impl/withdrawal/WithdrawalBodyChangedHandlerTest.java b/src/test/java/dev/vality/daway/handler/event/stock/impl/withdrawal/WithdrawalBodyChangedHandlerTest.java new file mode 100644 index 00000000..9e2fa941 --- /dev/null +++ b/src/test/java/dev/vality/daway/handler/event/stock/impl/withdrawal/WithdrawalBodyChangedHandlerTest.java @@ -0,0 +1,52 @@ +package dev.vality.daway.handler.event.stock.impl.withdrawal; + +import dev.vality.daway.TestData; +import dev.vality.daway.config.PostgresqlJooqSpringBootITest; +import dev.vality.daway.dao.withdrawal.impl.WithdrawalCashChangeDaoImpl; +import dev.vality.daway.domain.tables.records.WithdrawalCashChangeRecord; +import dev.vality.fistful.withdrawal.TimestampedChange; +import dev.vality.machinegun.eventsink.MachineEvent; +import org.jooq.DSLContext; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.test.context.ContextConfiguration; + +import static dev.vality.daway.domain.tables.WithdrawalCashChange.WITHDRAWAL_CASH_CHANGE; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +@PostgresqlJooqSpringBootITest +@ContextConfiguration(classes = {WithdrawalCashChangeDaoImpl.class, WithdrawalBodyChangedHandler.class}) +class WithdrawalBodyChangedHandlerTest { + + @Autowired + private WithdrawalBodyChangedHandler handler; + + @Autowired + private DSLContext dslContext; + + @BeforeEach + void setUp() { + dslContext.deleteFrom(WITHDRAWAL_CASH_CHANGE).execute(); + } + + @Test + void handle() { + TimestampedChange timestampedChange = TestData.createWithdrawalBodyChangedChange(); + MachineEvent event = TestData.createMachineEvent(timestampedChange); + + handler.handle(timestampedChange, event); + + WithdrawalCashChangeRecord record = dslContext.fetchAny(WITHDRAWAL_CASH_CHANGE); + assertNotNull(record); + assertEquals(event.getSourceId(), record.getWithdrawalId()); + assertEquals(event.getEventId(), record.getSequenceId()); + assertEquals(100L, record.getOldAmount()); + assertEquals("RUB", record.getOldCurrencyCode()); + assertEquals(200L, record.getNewAmount()); + assertEquals("USD", record.getNewCurrencyCode()); + assertTrue(record.getCurrent()); + } +} From 155f3468e0b9d50efbd88230f33e4928ca6168af Mon Sep 17 00:00:00 2001 From: Egor Cherniak Date: Tue, 28 Jul 2026 14:44:23 +0300 Subject: [PATCH 2/2] bump --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index cfa183c1..cdcc61b1 100644 --- a/pom.xml +++ b/pom.xml @@ -6,7 +6,7 @@ dev.vality service-parent-pom - 3.1.9 + 3.1.11 daway