diff --git a/pom.xml b/pom.xml
index 208ce5a8..cdcc61b1 100644
--- a/pom.xml
+++ b/pom.xml
@@ -6,7 +6,7 @@
dev.vality
service-parent-pom
- 3.1.9
+ 3.1.11
daway
@@ -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());
+ }
+}