diff --git a/.gitignore b/.gitignore
index 94afd7b..44b19e6 100644
--- a/.gitignore
+++ b/.gitignore
@@ -1,6 +1,7 @@
# Created by .ignore support plugin (hsz.mobi)
### Maven template
target/
+temp_context/
pom.xml.tag
pom.xml.releaseBackup
pom.xml.versionsBackup
diff --git a/README.md b/README.md
index 7432fa1..4016db6 100644
--- a/README.md
+++ b/README.md
@@ -1 +1,79 @@
# CC Reporter
+
+`cc-reporter` асинхронно строит CSV-отчёты по платежам и выводам. Сервис читает доменные события из Kafka,
+поддерживает в PostgreSQL актуальное состояние транзакций и формирует отчёт по согласованному снимку данных.
+
+Жизненный цикл задания:
+
+```text
+pending -> processing -> created -> expired
+ | | |
+ | | +-> failed / timed_out
+ | +----> pending (retry)
+ +----------------> canceled
+```
+
+Готовый файл хранится во внешнем файловом хранилище и выдаётся по временной подписанной ссылке.
+
+## CSV
+
+Формат одинаков для всех отчётов:
+
+| Параметр | Значение |
+|---|---|
+| Кодировка | UTF-8 |
+| Разделитель | `,` |
+| Конец строки | CRLF |
+| Экранирование | RFC 4180 |
+| `null` | Пустое поле |
+| Дата | `yyyy-MM-dd` |
+| Время | `HH:mm:ss` |
+| Timezone | `CreateReportRequest.timezone`, по умолчанию UTC |
+| Денежные значения | Decimal по экспоненте соответствующей валюты |
+| `exchange_rate_internal` | Decimal без экспоненциальной записи |
+
+`finalized_date` и `finalized_time` соответствуют текущему терминальному статусу. Если более новое событие
+корректирует терминальный статус, время финализации также обновляется.
+
+### Payments
+
+```csv
+created_date,created_time,finalized_date,finalized_time,invoice_id,payment_id,status,amount,currency,trx_id,provider_id,terminal_id,shop_id,exchange_rate_internal,provider_amount,provider_currency,original_amount,original_currency,converted_amount
+2026-08-20,10:15:00,2026-08-20,10:15:04,invoice-1,payment-1,captured,1000.00,RUB,trx-1,12,34,shop-1,1.0000000000,1000.00,RUB,1000.00,RUB,1000.00
+```
+
+| CSV-поле | Источник |
+|---|---|
+| `created_date`, `created_time` | `payment_txn_current.created_at` |
+| `finalized_date`, `finalized_time` | `payment_txn_current.finalized_at` |
+| `invoice_id`, `payment_id`, `status` | `payment_txn_current` |
+| `amount`, `currency` | `payment_txn_current` |
+| `trx_id` | `payment_txn_current.trx_id` |
+| `provider_id`, `terminal_id`, `shop_id` | `payment_txn_current` |
+| `exchange_rate_internal` | `payment_txn_current.exchange_rate_internal` |
+| `provider_amount`, `provider_currency` | `payment_txn_current` |
+| `original_amount`, `original_currency`, `converted_amount` | `payment_txn_current` |
+
+### Withdrawals
+
+```csv
+created_date,created_time,finalized_date,finalized_time,withdrawal_id,status,amount,currency,trx_id,provider_id,terminal_id,wallet_id,exchange_rate_internal,provider_amount,provider_currency,original_amount,original_currency,converted_amount
+2026-08-20,11:20:00,2026-08-20,11:20:05,withdrawal-1,succeeded,5000.00,RUB,trx-2,12,34,wallet-1,83.3333333333,5000.00,RUB,60.00,USD,5000.00
+```
+
+| CSV-поле | Источник |
+|---|---|
+| `created_date`, `created_time` | `withdrawal_txn_current.created_at` |
+| `finalized_date`, `finalized_time` | `withdrawal_txn_current.finalized_at` |
+| `withdrawal_id`, `status` | `withdrawal_txn_current` |
+| `amount`, `currency` | `withdrawal_txn_current`; обновляются при `body_changed` |
+| `trx_id` | последняя `withdrawal_session` |
+| `provider_id`, `terminal_id`, `wallet_id` | `withdrawal_txn_current` |
+| `exchange_rate_internal` | `withdrawal_txn_current.exchange_rate_internal` |
+| `provider_amount`, `provider_currency` | `withdrawal_txn_current` |
+| `original_amount`, `original_currency`, `converted_amount` | `withdrawal_txn_current` |
+
+## Документация
+
+- [Модель данных, ingestion и жизненный цикл](docs/docs.md)
+- [Требования и архитектурные решения](docs/PLAN.md)
diff --git a/docs/CSV_REPORT_FORMAT.md b/docs/CSV_REPORT_FORMAT.md
deleted file mode 100644
index 19ba0fb..0000000
--- a/docs/CSV_REPORT_FORMAT.md
+++ /dev/null
@@ -1,188 +0,0 @@
-# CSV Report Format (CC Reporter, final candidate)
-
-Документ описывает, как будет выглядеть итоговый CSV-файл для выгрузки.
-В терминах API:
-1. `report_type` задает бизнес-сущность отчета (`payments` или `withdrawals`);
-2. `file_type` задает формат файла (сейчас это только `csv`).
-
-Для каждой пары `report_type + file_type` используется один фиксированный формат без дополнительных профилей.
-
-Здесь зафиксированы:
-
-1. состав колонок в файле;
-2. понятный смысл этих колонок;
-3. технический источник значений для технарей.
-
-## Общие правила
-
-1. Encoding: UTF-8.
-2. Escaping: RFC4180.
-3. Decimal separator: `.`.
-4. Timestamps форматируются в timezone отчета (`CreateReportRequest.timezone`, по умолчанию `UTC`).
-5. `created_*` и `finalized_*` всегда разделены на отдельные колонки `date` и `time`.
-6. `finalized_*` пустые для non-terminal статусов.
-7. `null` сериализуется как пустая CSV-ячейка.
-8. Порядок колонок фиксирован и меняется только новой версией контракта.
-9. Понятные человеку названия (`shop_name`, `wallet_name`, `provider_name`, `terminal_name`) могут использоваться для поиска внутри системы, но в текущем CSV-контракте не считаются окончательно зафиксированными.
-10. В current-state схеме эти display-name поля уже существуют, но ingestion пока не имеет подтвержденного event-native mapping для всех
- таких значений; до отдельного закрытия этого гэпа они не должны восприниматься как надежно заполненные данные.
-
-## Payments CSV (`report_type = payments`, `file_type = csv`)
-
-> TODO: временное проектное решение для первой реализации.
-> До отдельного подтверждения источника `payments` FX-полей
-> (`original_amount`, `original_currency`, `converted_amount`, `exchange_rate_internal`,
-> `provider_amount`, `provider_currency`) эти колонки сохраняются в контракте и
-> временно могут заполняться mock-значениями. Это сделано специально, чтобы не
-> схлопнуть их до постоянного `null` и не потерять из вида незавершенный участок.
-
-### Порядок колонок в файле
-
-1. `created_date`
-2. `created_time`
-3. `finalized_date`
-4. `finalized_time`
-5. `invoice_id`
-6. `payment_id`
-7. `status`
-8. `amount`
-9. `currency`
-10. `trx_id`
-11. `provider_id`
-12. `terminal_id`
-13. `shop_id`
-14. `exchange_rate_internal`
-15. `provider_amount`
-16. `provider_currency`
-17. `original_amount`
-18. `original_currency`
-19. `converted_amount`
-
-### Что означают ключевые колонки
-
-1. `amount`:
- Основная сумма платежа в валюте самой операции.
-2. `currency`:
- Валюта основной суммы платежа.
-3. `trx_id`:
- Идентификатор транзакции на стороне внешнего провайдера или платежного канала.
-4. `exchange_rate_internal`:
- Внутренний курс конвертации, который использовала наша система при пересчете суммы между валютами. Это значение показывает, по какому курсу система рассчитала итоговую сумму при валютной операции.
-5. `provider_amount`:
- Сумма, которая фактически была передана провайдеру для обработки платежа.
-6. `provider_currency`:
- Валюта суммы, которая была передана провайдеру.
-7. `original_amount`:
- Исходная сумма до конвертации, если операция была валютной.
-8. `original_currency`:
- Валюта исходной суммы до конвертации.
-9. `converted_amount`:
- Сумма после конвертации в валюте `currency`.
-
-### Технический источник значения (для технарей)
-
-| CSV | CCR source |
-|---|---|
-| `created_date` / `created_time` | `payment_txn_current.created_at` |
-| `finalized_date` / `finalized_time` | `payment_txn_current.finalized_at` |
-| `invoice_id` | `payment_txn_current.invoice_id` |
-| `payment_id` | `payment_txn_current.payment_id` |
-| `status` | `payment_txn_current.status` |
-| `amount` | `payment_txn_current.amount` |
-| `currency` | `payment_txn_current.currency` |
-| `trx_id` | `payment_txn_current.trx_id` |
-| `provider_id` | `payment_txn_current.provider_id` |
-| `terminal_id` | `payment_txn_current.terminal_id` |
-| `shop_id` | `payment_txn_current.shop_id` |
-| `exchange_rate_internal` | `payment_txn_current.exchange_rate_internal` |
-| `provider_amount` | `payment_txn_current.provider_amount` |
-| `provider_currency` | `payment_txn_current.provider_currency` |
-| `original_amount` | `payment_txn_current.original_amount` |
-| `original_currency` | `payment_txn_current.original_currency` |
-| `converted_amount` | `payment_txn_current.converted_amount` |
-
-Для первой реализации:
-1. `trx_id` по `payments` планируется получать из `TransactionInfo.id` в `SessionTransactionBound`.
-2. FX-блок по `payments` пока считается незавершенным участком.
-3. Пока источник FX-данных не подтвержден, генератор CSV может временно писать туда mock-значения.
-4. Такое заполнение должно быть явно помечено в коде как `TODO`, чтобы затем заменить его на нормальный маппинг, а не закрепить как итоговое поведение.
-
-## Withdrawals CSV (`report_type = withdrawals`, `file_type = csv`)
-
-### Порядок колонок в файле
-
-1. `created_date`
-2. `created_time`
-3. `finalized_date`
-4. `finalized_time`
-5. `withdrawal_id`
-6. `status`
-7. `amount`
-8. `currency`
-9. `trx_id`
-10. `provider_id`
-11. `terminal_id`
-12. `wallet_id`
-13. `exchange_rate_internal`
-14. `provider_amount`
-15. `provider_currency`
-16. `original_amount`
-17. `original_currency`
-18. `converted_amount`
-
-### Что означают ключевые колонки
-
-1. `amount`:
- Основная сумма выплаты в валюте самой выплаты.
-2. `currency`:
- Валюта основной суммы выплаты.
-3. `trx_id`:
- Идентификатор операции на стороне внешнего провайдера или канала выплаты.
-4. `exchange_rate_internal`:
- Внутренний курс конвертации, который использовала наша система, если сумма выплаты пересчитывалась между валютами.
-5. `provider_amount`:
- Сумма, которая фактически была передана провайдеру для выполнения выплаты.
-6. `provider_currency`:
- Валюта суммы, которая была передана провайдеру.
-7. `original_amount`:
- Исходная сумма до конвертации, если выплата была валютной.
-8. `original_currency`:
- Валюта исходной суммы до конвертации.
-9. `converted_amount`:
- Сумма после конвертации в валюте `currency`.
-
-### Технический источник значения (для backend/QA)
-
-| CSV | CCR source |
-|---|---|
-| `created_date` / `created_time` | `withdrawal_txn_current.created_at` |
-| `finalized_date` / `finalized_time` | `withdrawal_txn_current.finalized_at` |
-| `withdrawal_id` | `withdrawal_txn_current.withdrawal_id` |
-| `status` | `withdrawal_txn_current.status` |
-| `amount` | `withdrawal_txn_current.amount` |
-| `currency` | `withdrawal_txn_current.currency` |
-| `trx_id` | `withdrawal_txn_current.trx_id` |
-| `provider_id` | `withdrawal_txn_current.provider_id` |
-| `terminal_id` | `withdrawal_txn_current.terminal_id` |
-| `wallet_id` | `withdrawal_txn_current.wallet_id` |
-| `exchange_rate_internal` | `withdrawal_txn_current.exchange_rate_internal` |
-| `provider_amount` | `withdrawal_txn_current.provider_amount` |
-| `provider_currency` | `withdrawal_txn_current.provider_currency` |
-| `original_amount` | `withdrawal_txn_current.original_amount` |
-| `original_currency` | `withdrawal_txn_current.original_currency` |
-| `converted_amount` | `withdrawal_txn_current.converted_amount` |
-
-## Как показываются денежные значения
-
-1. `amount` и `converted_amount` форматируются по exponent валюты из `currency`.
-2. `original_amount` форматируется по exponent валюты из `original_currency`.
-3. `provider_amount` форматируется по exponent валюты из `provider_currency`; если `provider_currency` пустая, используется exponent из `currency`.
-4. `exchange_rate_internal` показывается как обычное десятичное число, без инженерной записи через степень. Например: `1.25`, а не `1.25E0`.
-
-## Как заполняются поля в нестандартных случаях
-
-1. `trx_id` заполняется только если исходная система действительно передает идентификатор операции; иначе поле может остаться пустым.
-2. Для `payments` FX/conversion поля на первом проходе реализации могут временно заполняться mock-значениями с пометкой `TODO`; это временная проектная заглушка, а не финальный контракт.
-3. `provider_currency` должно быть заполнено, если `provider_amount` указана в валюте, отличной от `currency`; если исходная система не передает валюту провайдера отдельно, поле может остаться пустым.
-4. `finalized_date` / `finalized_time` заполняются только для конечных статусов операции.
-5. `finalized_at` фиксируется как момент первого перехода в конечный статус и после этого не должен изменяться последующими уточняющими обновлениями.
diff --git a/docs/PLAN.md b/docs/PLAN.md
index deaf1a6..6e3dfe0 100644
--- a/docs/PLAN.md
+++ b/docs/PLAN.md
@@ -25,7 +25,10 @@
### Что входит в первую версию
-1. Жизненный цикл отчета: `Create -> Pending -> Processing -> Created | Failed | TimedOut | Canceled | Expired`.
+1. Жизненный цикл отчета:
+ `pending -> processing -> created -> expired`,
+ `processing -> pending | failed | timed_out | canceled`,
+ `pending -> canceled`.
2. Загрузка данных из Kafka в таблицы актуального состояния сущностей по платежам и выплатам.
3. API на Thrift для клиентской (фронт) части.
4. Фоновый обработчик и планировщик для построения отчетов.
@@ -67,8 +70,8 @@
3. Для `insert` и `update` в lookup-таблицах сохраняется актуальное имя сущности и `dominant_version_id`.
4. Для `remove` сохраняется tombstone-состояние (`deleted = true`), чтобы более старый commit не мог вернуть устаревшее
имя.
-5. Эти lookup-таблицы используются для denormalized search по имени и/или ID при построении отчетов, но сами
- display-name поля не являются частью бизнес-ключа транзакционных current-state таблиц.
+5. Эти таблицы используются при поиске по имени и идентификатору во время построения отчёта.
+ Имена сущностей не входят в бизнес-ключи транзакционных таблиц актуального состояния.
### 3.4 Актуальное состояние `payments`
@@ -78,7 +81,11 @@
4. Контракт `upsert`:
- `INSERT ... ON CONFLICT (invoice_id, payment_id) DO UPDATE`
- `... WHERE payment_txn_current.domain_event_id < EXCLUDED.domain_event_id`
-5. `finalized_at` заполняется только при первом терминальном статусе и не должен затираться последующим нетерминальным обогащением данных.
+5. При событии смены статуса `finalized_at` синхронизируется с новым статусом:
+ - для терминального статуса записывается время этого status-event;
+ - для нетерминального статуса поле очищается;
+ - события, которые статус не меняют, `finalized_at` не затрагивают.
+ Это позволяет корректно обработать в том числе смену одного финального статуса на другой после корректировки.
### 3.5 Актуальное состояние `withdrawals`
@@ -129,7 +136,7 @@ Trx ID - идентификатор транзакции со стороны п
В формате выгружаемого файла нужно
-Пофиксить Проблему, при которой при выгрузке транзакций значение в поле Amount имеет некорректное значение (пример 2В 000,00В в‚ё на валюте Тенге)
+Пофиксить Проблему, при которой при выгрузке транзакций значение в поле Amount имеет некорректное значение (пример 2 000,00 ₸ на валюте Тенге)
Разделить дату и время на отдельные столбцы
Добавить столбцы выгрузки trx_id (идентификатор транзакции со стороны провайдера) и Currency
Для обработки транзакции через валютные каналы, добавить столбцы Курсов с нашей стороны и фактическая сумма в валюте передаваемая провайдеру
@@ -165,7 +172,7 @@ Trx ID - идентификатор транзакции со стороны п
1. Источник данных: Kafka, по тому же принципу, что и в текущих доменных сервисах.
2. Процесс чтения подтверждает `offset` только после успешного `commit` транзакции БД по всему пакету.
3. Повторное чтение `topic` должно быть безопасным и не приводить к дубликатам или откату актуального состояния.
-4. `CreateReport` должен поддерживать идемпотентность по `(created_by, idempotency_key)`. (`created_by` это идентификатор субьъекта из jwt токена который приходит в `Wachter`, `idempotency_key` uid с фронт энда)
+4. `CreateReport` должен поддерживать идемпотентность по `(created_by, idempotency_key)`. (`created_by` это идентификатор субъекта из jwt токена который приходит в `Wachter`, `idempotency_key` uid с фронтенда)
5. Нужен управляемый механизм восстановления для повторных попыток и зависших заданий.
6. Нужны индексы под реальные фильтры интерфейса и длинные диапазоны.
7. При ошибке генерации не должно оставаться поврежденных локальных артефактов.
@@ -177,8 +184,8 @@ Trx ID - идентификатор транзакции со стороны п
1. `Control Center Frontend`
2. `Wachter` (аутентификация, авторизация и маршрутизация по JWT)
3. `CC Reporter API` (Thrift)
-4. `CC Reporter Schedulator`: обработчик и планировщик
-5. `CC Reporter Kakfa Listener`: процессы чтения Kafka
+4. `CC Reporter Scheduler`: обработчик и планировщик
+5. `CC Reporter Kafka Listener`: процессы чтения Kafka
6. `PostgreSQL` (актуальное состояние и жизненный цикл отчетов)
7. `Minio`: S3-совместимое хранилище
@@ -187,10 +194,12 @@ Trx ID - идентификатор транзакции со стороны п
1. Пользователь задает фильтры и нажимает `Download report`.
2. Клиентская часть вызывает `CreateReport`.
3. API проверяет соответствие пары `report_type + file_type` и ветки `query`, сохраняет `report_job(status = pending)` и возвращает `report_id`.
-4. Обработчик забирает задание со статусом `pending`, переводит его в `processing` и увеличивает `attempt`.
+4. Планировщик атомарно забирает до `report.worker-concurrency` заданий со статусом `pending`, переводит их в
+ `processing`, увеличивает `attempt` и запускает ограниченным пулом worker-ов.
5. Обработчик открывает отдельную транзакцию `READ ONLY REPEATABLE READ` для чтения данных отчета.
6. Сразу после открытия транзакции он фиксирует `data_snapshot_fixed_at = transaction_timestamp()`.
-7. `started_at` отражает момент старта обработки задания worker-ом, а `data_snapshot_fixed_at` отражает момент фиксации MVCC-снимка данных для всего отчета.
+7. `started_at` отражает старт текущей попытки worker-а, а `data_snapshot_fixed_at` — момент фиксации MVCC-снимка данных
+ для всего отчета.
8. Внутри этой транзакции обработчик потоково читает актуальное состояние через серверный курсор, порциями.
9. CSV записывается во временный артефакт:
- локальный временный файл или
@@ -198,16 +207,25 @@ Trx ID - идентификатор транзакции со стороны п
10. После полной записи файл хешируется, загружается в итоговый ключ объекта, затем создается `report_file` (одна запись на один `report_job`).
11. Только после успешной публикации файла задание завершается в статусе `created`.
12. Если возникает ошибка, обработчик удаляет временный файл или объект и не создает `report_file`.
-13. При временной ошибке задание возвращается в `pending` с новым `next_attempt_at`.
-14. Отдельный процесс контроля таймаутов переводит зависшие задания в `timed_out`.
-15. Клиентская часть получает статусы через `GetReports` и `GetReport`.
-16. Для скачивания клиентская часть вызывает `GeneratePresignedUrl` с `file_id` и получает ссылку, для которой TTL принудительно ограничивается на стороне сервиса.
+13. При временной ошибке задание возвращается в `pending`; `next_attempt_at` рассчитывается от фактического времени
+ ошибки.
+14. Каждая попытка имеет hard timeout. По deadline worker task сначала отменяется через interrupt, затем отчёт условно
+ переводится из `processing` в `timed_out`.
+15. Stale cleanup переводит оставшиеся `processing` в `timed_out` после остановки процесса или потери instance.
+16. Все конкурирующие переходы содержат предикат исходного статуса. Поздний worker не может переписать `canceled`,
+ `timed_out` или другой терминальный статус.
+17. SQL генерации ограничивается локальным PostgreSQL `statement_timeout`. JDBC connect, socket read и cancel signal
+ имеют отдельные таймауты; TCP keepalive включён.
+18. Клиентская часть получает статусы через `GetReports` и `GetReport`.
+19. Для скачивания клиентская часть вызывает `GeneratePresignedUrl` с `file_id` и получает ссылку, для которой TTL
+ принудительно ограничивается на стороне сервиса.
### 5.3 Согласованность данных при генерации
`CCR` гарантирует два уровня фиксации:
-1. Логическое окно данных задается фильтрами пользователя (`requested_time_from`, `requested_time_to`).
+1. Логическое окно данных задается `time_range` внутри сохраненного `query_json`; отдельные дублирующие колонки не
+ нужны.
2. Физическая согласованность чтения обеспечивается транзакцией `READ ONLY REPEATABLE READ` на все время построения файла.
Это означает:
@@ -219,7 +237,7 @@ Trx ID - идентификатор транзакции со стороны п
1. SQL DDL
2. Thrift API
-3. Формат CSV
+3. Формат CSV в `README.md`
## 7. Сценарий в админке
@@ -272,9 +290,12 @@ Trx ID - идентификатор транзакции со стороны п
1. Использовать потоковую запись без накопления всего файла в памяти.
2. Читать данные через серверный курсор, порциями.
-3. Выделить отдельный пул обработчиков для тяжелых заданий.
+3. Использовать ограниченный пул обработчиков; число worker-ов задаётся через `report.worker-concurrency`.
4. Настроить политику повторных попыток через `next_attempt_at`.
-5. Отдельно отслеживать и переводить в таймаут зависшие задания со статусом `processing`.
+5. Ограничивать каждую попытку hard timeout и сохранять stale cleanup как восстановление после потери процесса.
+6. Поддерживать `spring.datasource.hikari.maximum-pool-size >= worker-concurrency + 2`; для смешанной нагрузки
+ использовать резерв `+4`. Соединения сверх worker pool нужны для переходов статусов и остальных транзакций
+ instance.
### 8.3 Позднее дозаполнение (`trx_id`, FX, данные провайдера)
@@ -299,8 +320,9 @@ Trx ID - идентификатор транзакции со стороны п
1. Ограничить число одновременно строящихся длинных отчетов.
2. Использовать только потоковое чтение, без помещения всего набора результатов в память.
-3. Ввести отдельный таймаут на построение отчета.
-4. Проводить нагрузочное тестирование именно для длительных сценариев с `REPEATABLE READ`.
+3. Ограничивать построение общим hard timeout и PostgreSQL `statement_timeout`.
+4. Ограничивать JDBC connect, socket read и cancel signal; включать TCP keepalive.
+5. Проводить нагрузочное тестирование именно для длительных сценариев с `REPEATABLE READ`.
### 8.5 Ограничение модели актуального состояния
@@ -350,3 +372,5 @@ Trx ID - идентификатор транзакции со стороны п
4. Во время генерации параллельные обновления не смешиваются в одном отчете.
5. Поврежденные временные артефакты не публикуются наружу.
6. CSV соответствует обязательным полям по требованиям.
+7. Несколько worker-ов не забирают один отчёт повторно и не превышают настроенный предел параллелизма.
+8. Поздний worker не перезаписывает `canceled`, `timed_out`, `failed` или `expired`.
diff --git a/docs/docs.md b/docs/docs.md
index 03795a2..4be01ee 100644
--- a/docs/docs.md
+++ b/docs/docs.md
@@ -1,93 +1,88 @@
-## Как живёт отчёт
+# Устройство CC Reporter
-- Сразу после `createReport` отчёт попадает в `pending`. Это просто запись о том, что отчёт заказан и ждёт воркера.
+## Жизненный цикл отчёта
-- Когда воркер забирает его в работу, статус меняется на `processing`. В этот момент отчёт реально строится.
+```text
+ ┌── временная ошибка ──> pending(next_attempt_at)
+ │
+pending ── claim ──> processing ──────┼── успех ──> created ── TTL ──> expired
+ │ │
+ └── cancel ──> canceled ├── закончились попытки ──> failed
+ ├── hard timeout ──> timed_out
+ └── cancel ──> canceled
+```
-- Если всё прошло нормально, отчёт переходит в `created`. Это значит, что файл уже собран, загружен в storage и его
- можно скачать.
+Переход `pending -> processing` выполняется атомарно. DAO выбирает только задания, для которых наступил
+`next_attempt_at`, блокирует строку через `FOR UPDATE SKIP LOCKED`, увеличивает `attempt`, записывает `started_at`
+и очищает ошибку предыдущей попытки. Поэтому несколько экземпляров сервиса могут разбирать одну очередь без
+двойной обработки одного задания.
-- Если во время построения что-то сломалось, отчёт либо вернётся в ожидание новой попытки, либо закончится в `failed`,
- если попытки закончились.
+Число одновременно выполняемых отчётов задаёт `report.worker-concurrency`. Одна попытка ограничена
+`report.processing-timeout-ms`. При превышении лимита worker получает interrupt, а запись условно переводится
+из `processing` в `timed_out`. Позднее завершение worker не может перезаписать уже установленный терминальный статус.
-- Если воркер завис или пропал и отчёт слишком долго висит в работе, его переводят в `timed_out`.
+При временной ошибке отчёт возвращается в `pending` с новым `next_attempt_at`. После исчерпания попыток он
+переходит в `failed`. Завершение отчёта и добавление `report_file` выполняются в одной транзакции.
-- Если пользователь успел отменить отчёт, пока он ещё не дошёл до готового файла, он становится `canceled`.
+Scheduler переводит готовые отчёты в `expired` после `expires_at`. `GetReport` и `GetReports` перед чтением также
+выполняют идемпотентное истечение просроченных `created`-отчётов, поэтому корректность API не зависит от точности
+срабатывания фонового scheduler. `GeneratePresignedUrl` дополнительно разрешает скачивание только пока отчёт
+не просрочен.
-- У готового отчёта есть срок жизни. Когда он заканчивается, статус меняется с `created` на `expired`. Это нормальный
- финал для успешного отчёта: он был доступен какое-то время, потом протух.
+## Согласованность current-state
-## Вычитывание полей `amount`, `provider`, `original` из потока событий
+`payment_txn_current` и `withdrawal_txn_current` содержат не историю, а последнее известное состояние сущности.
+Обновление принимается только если `MachineEvent.eventId` больше уже сохранённого. Повторные и более старые события
+не откатывают состояние назад.
-### Платежи
+Событие смены статуса является авторитетным для полей, которые зависят от статуса:
-Класс: `PaymentEventProjector`
+- `status` заменяется значением более нового status-event;
+- для терминального статуса `finalized_at` становится временем этого события;
+- если новый статус нетерминальный, `finalized_at` очищается;
+- `error_summary` заменяется значением нового status-event и очищается, если новая ошибка отсутствует.
-- `InvoicePaymentStarted`
- - при первом старте платежа заполняются основные денежные поля;
- - `amount` и `currency` берутся из `payment.cost`;
- - `providerId` и `terminalId` берутся из `started.route`, если маршрут уже есть;
- - `originalAmount` и `originalCurrency` на этом шаге тоже ставятся из `payment.cost`.
+События, которые статус не меняют, эти поля сохраняют. Это важно для корректировок, меняющих один финальный статус
+на другой.
-- `InvoicePaymentRouteChanged`
- - обновляет только маршрут;
- - `providerId` и `terminalId` перечитываются из нового `route`.
+## Payments
-- `InvoicePaymentCashChanged`
- - обновляет сумму платежа;
- - `amount` и `currency` берутся из `newCash`.
+`PaymentEventProjector` обрабатывает изменения платежа по порядку внутри `MachineEvent`.
-- `InvoicePaymentCashFlowChanged`
- - пересчитывает денежные значения по проводкам;
- - `amount` берётся через `DomainCashFlowExtractor.extractPaymentAmount(...)`;
- - дополнительно здесь же обновляется `fee`;
- - `originalAmount` и `originalCurrency` это событие не меняет.
+- `InvoicePaymentStarted` задаёт исходные идентификаторы, сумму, валюту, маршрут и статус `pending`.
+- `InvoicePaymentRouteChanged` обновляет provider/terminal.
+- `InvoicePaymentCashChanged` обновляет `amount` и `currency`.
+- `InvoicePaymentCashFlowChanged` пересчитывает `amount` и `fee` через `CashFlowAmountExtractor`.
+- `InvoicePaymentStatusChanged` обновляет статус, `finalized_at`, `error_summary` и данные `capturedCost`.
+- `SessionTransactionBound` сохраняет `trx_id`, RRN, approval code и данные конвертации.
+- `SessionProxyStateChanged` используется как fallback для `trx_id`.
-- `InvoicePaymentStatusChanged`
- - не меняет пользовательскую сумму платежа;
- - обновляет провайдерскую сторону расчёта:
- - `providerAmount` и `providerCurrency` берутся из `capturedCost`, если он есть в статусе.
+Если в одном `MachineEvent` есть несколько изменений одного платежа, они объединяются в порядке поступления.
+При нескольких status-event последнее изменение статуса определяет `status`, `finalized_at` и `error_summary`.
-Итого по платежу:
+## Withdrawals
-- обычная сумма платежа сначала приходит из `InvoicePaymentStarted`, потом может поменяться через
- `InvoicePaymentCashChanged` или `InvoicePaymentCashFlowChanged`;
-- провайдер определяется сначала в `InvoicePaymentStarted`, потом может быть переопределён через
- `InvoicePaymentRouteChanged`;
-- `originalAmount` и `originalCurrency` выставляются только на старте платежа и дальше этим проектором не обновляются.
+`WithdrawalEventProjector` обрабатывает:
-### Выводы
+- `created`: исходные данные вывода, body, маршрут и quote;
+- `body_changed`: заменяет текущие `amount` и `currency` значениями из `new_body`;
+- `route`: обновляет provider/terminal;
+- `status_changed`: обновляет статус, `finalized_at` и `error_summary`;
+- `transfer.payload.created.transfer.cashflow`: обновляет `fee`.
-Класс: `WithdrawalEventProjector`
+`original_amount`, `original_currency`, `provider_amount`, `provider_currency` и внутренний курс первоначально
+вычисляются из quote события создания. Текущая сумма вывода при этом может измениться отдельным `body_changed`.
-- `Created`
- - на создании вывода заполняются и текущая сумма, и исходная сумма, и маршрут;
- - `amount` и `currency` берутся из `withdrawal.body`;
- - `providerId` и `terminalId` берутся из `withdrawal.route`, если маршрут уже известен;
- - `originalAmount` и `originalCurrency` берутся из `quote.cashFrom`, если есть котировка;
- - дополнительно провайдерская сумма заполняется из `quote.cashTo`;
- - курс `exchangeRateInternal` считается как отношение `cashTo / cashFrom`.
+`WithdrawalSessionEventProjector` хранит связь с выводом и транзакционные данные сессии (`trx_id`, `trx_search`).
+В CSV используется последняя подходящая сессия.
-- `Route`
- - обновляет только маршрут;
- - `providerId` и `terminalId` перечитываются из нового `route`.
+## Построение CSV
-- `StatusChanged`
- - сумму, маршрут и исходную сумму не меняет.
+Один отчёт выполняется в транзакции `READ ONLY REPEATABLE READ` и использует один согласованный снимок PostgreSQL.
+Данные читаются курсором, поэтому весь набор строк не загружается в память.
-- `Transfer -> payload.created.transfer.cashflow`
- - обновляет только `fee`;
- - `amount`, `providerId`, `originalAmount` и связанные валюты не трогает.
+Денежные значения хранятся в minor units и при записи CSV переводятся в decimal по exponent валюты.
+`exchange_rate_internal` записывается как обычное десятичное число без экспоненциальной формы.
-Итого по выводу:
-
-- `amount` приходит из события создания и дальше этим проектором не меняется;
-- провайдер приходит из события создания и может обновиться отдельным событием смены маршрута;
-- `originalAmount` и `originalCurrency` приходят из котировки в событии создания и дальше не меняются.
-
-### Сессия вывода
-
-Класс: `WithdrawalSessionEventProjector`
-
-Этот проектор не работает с `amount`, `provider` и `original`. Он сохраняет связь с выводом и данные по транзакции (
-`trxId`, `trxSearch`).
+Локальный временный CSV удаляется при ошибке генерации. После успешной загрузки жизненный цикл файла контролируется
+через запись отчёта и TTL внешнего хранилища.
diff --git a/mockgit.txt b/mockgit.txt
new file mode 100644
index 0000000..d019ed9
--- /dev/null
+++ b/mockgit.txt
@@ -0,0 +1,25 @@
+```
+set -e
+
+git fetch origin
+
+# Пустой базовый commit ветки mock
+MOCK_COMMIT=$(git rev-parse origin/mock)
+
+# Актуальное дерево ветки init
+INIT_TREE=$(git rev-parse 'origin/init^{tree}')
+
+# Создаём новый review-commit:
+# содержимое = актуальный init
+# родитель = пустой mock
+REVIEW_COMMIT=$(
+ printf '%s\n' "Full init code snapshot for review" |
+ git commit-tree "$INIT_TREE" -p "$MOCK_COMMIT"
+)
+
+# Обновляем техническую ветку
+git branch -f init-full-review "$REVIEW_COMMIT"
+
+# Обновляем PR
+git push --force-with-lease origin init-full-review
+```
\ No newline at end of file
diff --git a/pom.xml b/pom.xml
index 6093bec..4cb3830 100644
--- a/pom.xml
+++ b/pom.xml
@@ -6,7 +6,7 @@
dev.vality
service-parent-pom
- 4.0.0-BETA-1
+ 4.0.1
cc-reporter
@@ -36,7 +36,7 @@
org.springframework.boot
- spring-boot-starter-web
+ spring-boot-starter-webmvc
org.springframework.boot
@@ -79,6 +79,13 @@
org.projectlombok
lombok
+ provided
+
+
+
+
+ io.opentelemetry
+ opentelemetry-api
@@ -113,10 +120,14 @@
dev.vality
damsel
+
+ dev.vality.geck
+ serializer
+
dev.vality
fistful-proto
- 1.188-f7ce08e
+ 1.194-360c737
jakarta.annotation
@@ -157,6 +168,7 @@
true
Dockerfile
+ opentelemetry-javaagent.jar
diff --git a/renovate.json b/renovate.json
deleted file mode 100644
index a20bfd6..0000000
--- a/renovate.json
+++ /dev/null
@@ -1,4 +0,0 @@
-{
- "$schema": "https://docs.renovatebot.com/renovate-schema.json",
- "extends": ["local>valitydev/.github:renovate-config"]
-}
diff --git a/src/main/java/dev/vality/ccreporter/config/FileStorageConfig.java b/src/main/java/dev/vality/ccreporter/config/FileStorageConfig.java
index cfd5a78..58fc513 100644
--- a/src/main/java/dev/vality/ccreporter/config/FileStorageConfig.java
+++ b/src/main/java/dev/vality/ccreporter/config/FileStorageConfig.java
@@ -3,21 +3,21 @@
import dev.vality.ccreporter.config.properties.FileStorageProperties;
import dev.vality.file.storage.FileStorageSrv;
import dev.vality.woody.thrift.impl.http.THSpawnClientBuilder;
-import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
-import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.net.URI;
import java.net.http.HttpClient;
-import java.lang.reflect.Proxy;
+import java.time.Duration;
@Configuration
public class FileStorageConfig {
@Bean
- public HttpClient httpClient() {
- return HttpClient.newHttpClient();
+ public HttpClient httpClient(FileStorageProperties fileStorageProperties) {
+ return HttpClient.newBuilder()
+ .connectTimeout(Duration.ofMillis(fileStorageProperties.getNetworkTimeout()))
+ .build();
}
@Bean
diff --git a/src/main/java/dev/vality/ccreporter/config/ReportWorkerConfig.java b/src/main/java/dev/vality/ccreporter/config/ReportWorkerConfig.java
new file mode 100644
index 0000000..7da729a
--- /dev/null
+++ b/src/main/java/dev/vality/ccreporter/config/ReportWorkerConfig.java
@@ -0,0 +1,36 @@
+package dev.vality.ccreporter.config;
+
+import com.zaxxer.hikari.HikariDataSource;
+import dev.vality.ccreporter.config.properties.ReportProperties;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+
+import javax.sql.DataSource;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+
+@Configuration
+public class ReportWorkerConfig {
+
+ private static final int MINIMUM_CONNECTION_RESERVE = 2;
+
+ @Bean(destroyMethod = "shutdownNow")
+ public ExecutorService reportWorkerExecutor(ReportProperties reportProperties, DataSource dataSource) {
+ validatePoolCapacity(reportProperties, dataSource);
+ var threadFactory = Thread.ofPlatform()
+ .name("ccr-report-worker-", 0)
+ .factory();
+ return Executors.newFixedThreadPool(reportProperties.getWorkerConcurrency(), threadFactory);
+ }
+
+ private void validatePoolCapacity(ReportProperties reportProperties, DataSource dataSource) {
+ var minimumPoolSize = reportProperties.getWorkerConcurrency() + MINIMUM_CONNECTION_RESERVE;
+ if (dataSource instanceof HikariDataSource hikariDataSource
+ && hikariDataSource.getMaximumPoolSize() < minimumPoolSize) {
+ throw new IllegalStateException(
+ "spring.datasource.hikari.maximum-pool-size must be at least " + minimumPoolSize +
+ " for report.worker-concurrency=" + reportProperties.getWorkerConcurrency()
+ );
+ }
+ }
+}
diff --git a/src/main/java/dev/vality/ccreporter/config/properties/CcrApiProperties.java b/src/main/java/dev/vality/ccreporter/config/properties/CcrApiProperties.java
index a6e3cc8..8b9d0e0 100644
--- a/src/main/java/dev/vality/ccreporter/config/properties/CcrApiProperties.java
+++ b/src/main/java/dev/vality/ccreporter/config/properties/CcrApiProperties.java
@@ -6,7 +6,7 @@
@Data
@Configuration
-@ConfigurationProperties(prefix = "ccr.api")
+@ConfigurationProperties(prefix = "api")
public class CcrApiProperties {
private String path;
diff --git a/src/main/java/dev/vality/ccreporter/config/properties/CcrKafkaProperties.java b/src/main/java/dev/vality/ccreporter/config/properties/CcrKafkaProperties.java
index c6f9c42..cad4889 100644
--- a/src/main/java/dev/vality/ccreporter/config/properties/CcrKafkaProperties.java
+++ b/src/main/java/dev/vality/ccreporter/config/properties/CcrKafkaProperties.java
@@ -4,7 +4,7 @@
import org.springframework.boot.context.properties.ConfigurationProperties;
@Data
-@ConfigurationProperties(prefix = "ccr.kafka")
+@ConfigurationProperties(prefix = "kafka")
public class CcrKafkaProperties {
private Consumer consumer;
diff --git a/src/main/java/dev/vality/ccreporter/config/properties/FileStorageProperties.java b/src/main/java/dev/vality/ccreporter/config/properties/FileStorageProperties.java
index 3afd146..6343427 100644
--- a/src/main/java/dev/vality/ccreporter/config/properties/FileStorageProperties.java
+++ b/src/main/java/dev/vality/ccreporter/config/properties/FileStorageProperties.java
@@ -6,7 +6,7 @@
@Data
@Configuration
-@ConfigurationProperties(prefix = "ccr.storage.file-storage")
+@ConfigurationProperties(prefix = "storage.file-storage")
public class FileStorageProperties {
private String url;
diff --git a/src/main/java/dev/vality/ccreporter/config/properties/ReportProperties.java b/src/main/java/dev/vality/ccreporter/config/properties/ReportProperties.java
index 9b4fb87..b582d80 100644
--- a/src/main/java/dev/vality/ccreporter/config/properties/ReportProperties.java
+++ b/src/main/java/dev/vality/ccreporter/config/properties/ReportProperties.java
@@ -1,15 +1,22 @@
package dev.vality.ccreporter.config.properties;
+import jakarta.validation.constraints.Positive;
import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Configuration;
+import org.springframework.validation.annotation.Validated;
@Data
+@Validated
@Configuration
-@ConfigurationProperties(prefix = "ccr.report")
+@ConfigurationProperties(prefix = "report")
public class ReportProperties {
private int maxAttempts;
+ @Positive
+ private int workerConcurrency;
+ @Positive
+ private long processingTimeoutMs;
private int presignedUrlTtlSec;
private long expirationSec;
diff --git a/src/main/java/dev/vality/ccreporter/config/properties/ReportSchedulerProperties.java b/src/main/java/dev/vality/ccreporter/config/properties/ReportSchedulerProperties.java
deleted file mode 100644
index 3e525cf..0000000
--- a/src/main/java/dev/vality/ccreporter/config/properties/ReportSchedulerProperties.java
+++ /dev/null
@@ -1,14 +0,0 @@
-package dev.vality.ccreporter.config.properties;
-
-import lombok.Data;
-import org.springframework.boot.context.properties.ConfigurationProperties;
-import org.springframework.context.annotation.Configuration;
-
-@Data
-@Configuration
-@ConfigurationProperties(prefix = "ccr.scheduler")
-public class ReportSchedulerProperties {
-
- private long staleProcessingTimeoutMs;
-
-}
diff --git a/src/main/java/dev/vality/ccreporter/dao/PaymentTxnCurrentDao.java b/src/main/java/dev/vality/ccreporter/dao/PaymentTxnCurrentDao.java
index 2844337..b066604 100644
--- a/src/main/java/dev/vality/ccreporter/dao/PaymentTxnCurrentDao.java
+++ b/src/main/java/dev/vality/ccreporter/dao/PaymentTxnCurrentDao.java
@@ -4,6 +4,7 @@
import lombok.RequiredArgsConstructor;
import org.jooq.DSLContext;
import org.jooq.Field;
+import org.jooq.impl.DSL;
import org.springframework.stereotype.Repository;
import java.util.Map;
@@ -19,8 +20,7 @@ public class PaymentTxnCurrentDao {
private static final Set> IMMUTABLE_FIELDS = Set.of(
PAYMENT_TXN_CURRENT.ID,
PAYMENT_TXN_CURRENT.INVOICE_ID,
- PAYMENT_TXN_CURRENT.PAYMENT_ID,
- PAYMENT_TXN_CURRENT.UPDATED_AT
+ PAYMENT_TXN_CURRENT.PAYMENT_ID
);
private static final Set> OVERWRITE_FIELDS = Set.of(
@@ -43,12 +43,22 @@ var record = dslContext.newRecord(PAYMENT_TXN_CURRENT, update);
PAYMENT_TXN_CURRENT,
IMMUTABLE_FIELDS,
OVERWRITE_FIELDS,
- Map.of(PAYMENT_TXN_CURRENT.UPDATED_AT, UTC_NOW)
- ))
- .where(isIncomingEventNewer(
- PAYMENT_TXN_CURRENT.DOMAIN_EVENT_CREATED_AT,
- PAYMENT_TXN_CURRENT.DOMAIN_EVENT_ID
+ Map.of(
+ PAYMENT_TXN_CURRENT.UPDATED_AT,
+ UTC_NOW,
+ PAYMENT_TXN_CURRENT.FINALIZED_AT,
+ DSL.when(
+ DSL.excluded(PAYMENT_TXN_CURRENT.STATUS).isNotNull(),
+ DSL.excluded(PAYMENT_TXN_CURRENT.FINALIZED_AT)
+ ).otherwise(PAYMENT_TXN_CURRENT.FINALIZED_AT),
+ PAYMENT_TXN_CURRENT.ERROR_SUMMARY,
+ DSL.when(
+ DSL.excluded(PAYMENT_TXN_CURRENT.STATUS).isNotNull(),
+ DSL.excluded(PAYMENT_TXN_CURRENT.ERROR_SUMMARY)
+ ).otherwise(PAYMENT_TXN_CURRENT.ERROR_SUMMARY)
+ )
))
+ .where(isIncomingEventNewer(PAYMENT_TXN_CURRENT.DOMAIN_EVENT_ID))
.execute();
}
}
diff --git a/src/main/java/dev/vality/ccreporter/dao/ReportCommandDao.java b/src/main/java/dev/vality/ccreporter/dao/ReportCommandDao.java
index c67baa4..eca87e0 100644
--- a/src/main/java/dev/vality/ccreporter/dao/ReportCommandDao.java
+++ b/src/main/java/dev/vality/ccreporter/dao/ReportCommandDao.java
@@ -3,17 +3,14 @@
import dev.vality.ccreporter.FileType;
import dev.vality.ccreporter.ReportQuery;
import dev.vality.ccreporter.ReportType;
-import dev.vality.ccreporter.dao.mapper.ReportCommandMapper;
-import dev.vality.ccreporter.report.ReportQueryService;
+import dev.vality.ccreporter.dao.mapper.ReportRecordMapper;
import dev.vality.ccreporter.serde.json.ThriftJsonCodec;
import lombok.RequiredArgsConstructor;
import org.jooq.DSLContext;
-import org.jooq.exception.IntegrityConstraintViolationException;
-import org.springframework.dao.DuplicateKeyException;
+import org.jooq.JSONB;
import org.springframework.stereotype.Repository;
import org.springframework.util.StringUtils;
-import java.util.Objects;
import java.util.Optional;
import static dev.vality.ccreporter.domain.Tables.REPORT_JOB;
@@ -23,21 +20,9 @@
public class ReportCommandDao {
private final DSLContext dslContext;
- private final ReportQueryService reportQueryService;
private final ThriftJsonCodec thriftJsonCodec;
- public Optional findByIdempotencyKey(String createdBy, String idempotencyKey) {
- if (!StringUtils.hasText(idempotencyKey)) {
- return Optional.empty();
- }
- return dslContext.select(REPORT_JOB.ID)
- .from(REPORT_JOB)
- .where(REPORT_JOB.CREATED_BY.eq(createdBy))
- .and(REPORT_JOB.IDEMPOTENCY_KEY.eq(idempotencyKey))
- .fetchOptional(REPORT_JOB.ID);
- }
-
- public long createReport(
+ public CreateResult createReport(
String createdBy,
ReportType reportType,
FileType fileType,
@@ -45,32 +30,43 @@ public long createReport(
String timezone,
String idempotencyKey
) {
- var querySpec = reportQueryService.resolveQuerySpec(query);
- var queryJson = thriftJsonCodec.serialize(query);
- var queryHash = reportQueryService.hash(queryJson);
- var reportJob = ReportCommandMapper.mapReportJob(
- createdBy,
- reportType,
- fileType,
- queryJson,
- queryHash,
- querySpec.timeRange().from(),
- querySpec.timeRange().to(),
- timezone,
- idempotencyKey
- );
+ var normalizedIdempotencyKey = StringUtils.hasText(idempotencyKey) ? idempotencyKey : null;
+ var insertedId = dslContext.insertInto(REPORT_JOB)
+ .columns(
+ REPORT_JOB.REPORT_TYPE,
+ REPORT_JOB.FILE_TYPE,
+ REPORT_JOB.QUERY_JSON,
+ REPORT_JOB.TIMEZONE,
+ REPORT_JOB.CREATED_BY,
+ REPORT_JOB.IDEMPOTENCY_KEY
+ )
+ .values(
+ ReportRecordMapper.mapEnum(
+ reportType,
+ dev.vality.ccreporter.domain.enums.ReportType.class
+ ),
+ ReportRecordMapper.mapEnum(
+ fileType,
+ dev.vality.ccreporter.domain.enums.FileType.class
+ ),
+ JSONB.jsonb(thriftJsonCodec.serialize(query)),
+ timezone,
+ createdBy,
+ normalizedIdempotencyKey
+ )
+ .onConflictDoNothing()
+ .returningResult(REPORT_JOB.ID)
+ .fetchOptional(REPORT_JOB.ID);
- try {
- return Objects.requireNonNull(
- dslContext.insertInto(REPORT_JOB)
- .set(ReportCommandMapper.newInsertableRecord(dslContext, reportJob))
- .returningResult(REPORT_JOB.ID)
- .fetchOne(REPORT_JOB.ID),
- "Report creation must return an id"
- );
- } catch (IntegrityConstraintViolationException ex) {
- throw new DuplicateKeyException("Report idempotency key already exists", ex);
+ if (insertedId.isPresent()) {
+ return new CreateResult(insertedId.get(), true);
}
+ if (normalizedIdempotencyKey == null) {
+ throw new IllegalStateException("Report insert returned no id");
+ }
+ return findByIdempotencyKey(createdBy, normalizedIdempotencyKey)
+ .map(reportId -> new CreateResult(reportId, false))
+ .orElseThrow(() -> new IllegalStateException("Conflicting report was not found"));
}
public boolean reportExists(String createdBy, long reportId) {
@@ -81,4 +77,15 @@ public boolean reportExists(String createdBy, long reportId) {
.and(REPORT_JOB.CREATED_BY.eq(createdBy))
);
}
+
+ private Optional findByIdempotencyKey(String createdBy, String idempotencyKey) {
+ return dslContext.select(REPORT_JOB.ID)
+ .from(REPORT_JOB)
+ .where(REPORT_JOB.CREATED_BY.eq(createdBy))
+ .and(REPORT_JOB.IDEMPOTENCY_KEY.eq(idempotencyKey))
+ .fetchOptional(REPORT_JOB.ID);
+ }
+
+ public record CreateResult(long reportId, boolean created) {
+ }
}
diff --git a/src/main/java/dev/vality/ccreporter/dao/ReportCsvDao.java b/src/main/java/dev/vality/ccreporter/dao/ReportCsvDao.java
index 434aa9b..c29b53b 100644
--- a/src/main/java/dev/vality/ccreporter/dao/ReportCsvDao.java
+++ b/src/main/java/dev/vality/ccreporter/dao/ReportCsvDao.java
@@ -1,14 +1,11 @@
package dev.vality.ccreporter.dao;
import dev.vality.ccreporter.PaymentsQuery;
-import dev.vality.ccreporter.PaymentsSearchFilter;
import dev.vality.ccreporter.WithdrawalsQuery;
-import dev.vality.ccreporter.WithdrawalsSearchFilter;
import lombok.RequiredArgsConstructor;
import org.jooq.*;
import org.jooq.Record;
import org.jooq.impl.DSL;
-import org.jooq.impl.SQLDataType;
import org.springframework.stereotype.Repository;
import java.time.Instant;
@@ -17,6 +14,7 @@
import java.util.List;
import static dev.vality.ccreporter.domain.Tables.*;
+import static dev.vality.ccreporter.util.SearchValueNormalizer.normalize;
import static dev.vality.ccreporter.util.TimestampUtils.parse;
import static dev.vality.ccreporter.util.TimestampUtils.toLocalDateTime;
@@ -39,11 +37,18 @@ public Instant currentSnapshot() {
.toInstant();
}
+ public void setLocalStatementTimeout(long timeoutMs) {
+ dslContext.fetchSingle(
+ "SELECT set_config('statement_timeout', ?, true)",
+ timeoutMs + "ms"
+ );
+ }
+
public Cursor extends Record> fetchPayments(PaymentsQuery query) {
var conditions = buildPaymentsConditions(query);
return dslContext.select(
- timestampField(PAYMENT_TXN_CURRENT.CREATED_AT, CREATED_AT),
- timestampField(PAYMENT_TXN_CURRENT.FINALIZED_AT, FINALIZED_AT),
+ PAYMENT_TXN_CURRENT.CREATED_AT.as(CREATED_AT),
+ PAYMENT_TXN_CURRENT.FINALIZED_AT.as(FINALIZED_AT),
PAYMENT_TXN_CURRENT.INVOICE_ID.as("invoice_id"),
PAYMENT_TXN_CURRENT.PAYMENT_ID.as("payment_id"),
PAYMENT_TXN_CURRENT.STATUS.as("status"),
@@ -92,8 +97,8 @@ public Cursor extends Record> fetchWithdrawals(WithdrawalsQuery query) {
var latestSessionSessionId = latestSession.field("session_id", String.class);
var conditions = buildWithdrawalConditions(query, latestSessionTrxId, latestSessionTrxSearch);
return dslContext.select(
- timestampField(WITHDRAWAL_TXN_CURRENT.CREATED_AT, CREATED_AT),
- timestampField(WITHDRAWAL_TXN_CURRENT.FINALIZED_AT, FINALIZED_AT),
+ WITHDRAWAL_TXN_CURRENT.CREATED_AT.as(CREATED_AT),
+ WITHDRAWAL_TXN_CURRENT.FINALIZED_AT.as(FINALIZED_AT),
WITHDRAWAL_TXN_CURRENT.WITHDRAWAL_ID.as("withdrawal_id"),
WITHDRAWAL_TXN_CURRENT.STATUS.as("status"),
WITHDRAWAL_TXN_CURRENT.AMOUNT.as("amount"),
@@ -106,7 +111,8 @@ public Cursor extends Record> fetchWithdrawals(WithdrawalsQuery query) {
WITHDRAWAL_TXN_CURRENT.PROVIDER_AMOUNT.as("provider_amount"),
WITHDRAWAL_TXN_CURRENT.PROVIDER_CURRENCY.as(PROVIDER_CURRENCY),
WITHDRAWAL_TXN_CURRENT.ORIGINAL_AMOUNT.as("original_amount"),
- WITHDRAWAL_TXN_CURRENT.ORIGINAL_CURRENCY.as(ORIGINAL_CURRENCY)
+ WITHDRAWAL_TXN_CURRENT.ORIGINAL_CURRENCY.as(ORIGINAL_CURRENCY),
+ WITHDRAWAL_TXN_CURRENT.CONVERTED_AMOUNT.as("converted_amount")
)
.from(WITHDRAWAL_TXN_CURRENT)
.leftJoin(latestSession).on(DSL.trueCondition())
@@ -134,24 +140,26 @@ private List buildPaymentsConditions(PaymentsQuery query) {
appendInCondition(conditions, PAYMENT_TXN_CURRENT.TRX_ID, query.getTrxIds());
appendInCondition(conditions, PAYMENT_TXN_CURRENT.CURRENCY, query.getCurrencies());
appendInCondition(conditions, PAYMENT_TXN_CURRENT.STATUS, query.getStatuses());
- appendSearchCondition(conditions, lowercaseSearchField(SHOP_LOOKUP.SHOP_SEARCH), query.getFilter(), "shop");
+ var filter = query.getFilter();
+ appendSearchCondition(
+ conditions,
+ SHOP_LOOKUP.SHOP_SEARCH,
+ filter == null ? null : filter.getShopTerm()
+ );
appendSearchCondition(
conditions,
- lowercaseSearchField(PROVIDER_LOOKUP.PROVIDER_SEARCH),
- query.getFilter(),
- "provider"
+ PROVIDER_LOOKUP.PROVIDER_SEARCH,
+ filter == null ? null : filter.getProviderTerm()
);
appendSearchCondition(
conditions,
- lowercaseSearchField(TERMINAL_LOOKUP.TERMINAL_SEARCH),
- query.getFilter(),
- "terminal"
+ TERMINAL_LOOKUP.TERMINAL_SEARCH,
+ filter == null ? null : filter.getTerminalTerm()
);
appendSearchCondition(
conditions,
- lowercaseSearchField(PAYMENT_TXN_CURRENT.TRX_SEARCH),
- query.getFilter(),
- "trx"
+ PAYMENT_TXN_CURRENT.TRX_SEARCH,
+ filter == null ? null : filter.getTrxTerm()
);
return conditions;
}
@@ -175,25 +183,27 @@ private List buildWithdrawalConditions(
appendInCondition(conditions, latestSessionTrxId, query.getTrxIds());
appendInCondition(conditions, WITHDRAWAL_TXN_CURRENT.CURRENCY, query.getCurrencies());
appendInCondition(conditions, WITHDRAWAL_TXN_CURRENT.STATUS, query.getStatuses());
+ var filter = query.getFilter();
appendSearchCondition(
conditions,
- lowercaseSearchField(WALLET_LOOKUP.WALLET_SEARCH),
- query.getFilter(),
- "wallet"
+ WALLET_LOOKUP.WALLET_SEARCH,
+ filter == null ? null : filter.getWalletTerm()
);
appendSearchCondition(
conditions,
- lowercaseSearchField(PROVIDER_LOOKUP.PROVIDER_SEARCH),
- query.getFilter(),
- "provider"
+ PROVIDER_LOOKUP.PROVIDER_SEARCH,
+ filter == null ? null : filter.getProviderTerm()
);
appendSearchCondition(
conditions,
- lowercaseSearchField(TERMINAL_LOOKUP.TERMINAL_SEARCH),
- query.getFilter(),
- "terminal"
+ TERMINAL_LOOKUP.TERMINAL_SEARCH,
+ filter == null ? null : filter.getTerminalTerm()
+ );
+ appendSearchCondition(
+ conditions,
+ latestSessionTrxSearch,
+ filter == null ? null : filter.getTrxTerm()
);
- appendSearchCondition(conditions, lowercaseSearchField(latestSessionTrxSearch), query.getFilter(), "trx");
return conditions;
}
@@ -204,54 +214,19 @@ private void appendInCondition(List conditions, Field field,
conditions.add(field.in(values));
}
- private void appendSearchCondition(
- List conditions,
- Field field,
- PaymentsSearchFilter filter,
- String filterType
- ) {
- if (filter == null) {
- return;
- }
- appendSearchCondition(conditions, field, switch (filterType) {
- case "shop" -> filter.getShopTerm();
- case "provider" -> filter.getProviderTerm();
- case "terminal" -> filter.getTerminalTerm();
- case "trx" -> filter.getTrxTerm();
- default -> null;
- });
- }
-
- private void appendSearchCondition(
- List conditions,
- Field field,
- WithdrawalsSearchFilter filter,
- String filterType
- ) {
- if (filter == null) {
- return;
- }
- appendSearchCondition(conditions, field, switch (filterType) {
- case "wallet" -> filter.getWalletTerm();
- case "provider" -> filter.getProviderTerm();
- case "terminal" -> filter.getTerminalTerm();
- case "trx" -> filter.getTrxTerm();
- default -> null;
- });
- }
-
private void appendSearchCondition(List conditions, Field field, String value) {
if (value == null || value.isBlank()) {
return;
}
- conditions.add(field.like("%" + value.toLowerCase() + "%"));
- }
-
- private Field lowercaseSearchField(Field field) {
- return DSL.lower(DSL.coalesce(field, ""));
+ var normalizedValue = normalize(value);
+ var pattern = "%" + escapeLikeLiteral(normalizedValue) + "%";
+ conditions.add(DSL.condition("{0} LIKE {1} ESCAPE '!'", field, DSL.val(pattern)));
}
- private Field timestampField(Field field, String alias) {
- return field.cast(SQLDataType.TIMESTAMP).as(alias);
+ private String escapeLikeLiteral(String value) {
+ return value
+ .replace("!", "!!")
+ .replace("%", "!%")
+ .replace("_", "!_");
}
}
diff --git a/src/main/java/dev/vality/ccreporter/dao/ReportLifecycleDao.java b/src/main/java/dev/vality/ccreporter/dao/ReportLifecycleDao.java
index 87996a8..f08f2f9 100644
--- a/src/main/java/dev/vality/ccreporter/dao/ReportLifecycleDao.java
+++ b/src/main/java/dev/vality/ccreporter/dao/ReportLifecycleDao.java
@@ -1,197 +1,188 @@
package dev.vality.ccreporter.dao;
-import dev.vality.ccreporter.dao.mapper.ReportLifecycleMapper;
-import dev.vality.ccreporter.dao.mapper.ReportRecordMapper;
import dev.vality.ccreporter.domain.enums.ReportStatus;
import dev.vality.ccreporter.domain.tables.pojos.ReportFile;
-import dev.vality.ccreporter.domain.tables.pojos.ReportJob;
+import dev.vality.ccreporter.model.ReportTask;
import lombok.RequiredArgsConstructor;
-import org.jooq.Condition;
import org.jooq.DSLContext;
import org.jooq.Field;
import org.jooq.impl.DSL;
import org.springframework.stereotype.Repository;
+import org.springframework.transaction.annotation.Transactional;
import java.time.Instant;
import java.time.LocalDateTime;
import java.util.Optional;
-import static dev.vality.ccreporter.dao.support.ReportDaoSupport.*;
import static dev.vality.ccreporter.domain.Tables.REPORT_FILE;
import static dev.vality.ccreporter.domain.Tables.REPORT_JOB;
+import static dev.vality.ccreporter.util.TimestampUtils.toLocalDateTime;
@Repository
@RequiredArgsConstructor
public class ReportLifecycleDao {
private static final Field CANDIDATE_ID = DSL.field(DSL.name("candidate", "id"), Long.class);
- private static final ReportStatus PENDING = ReportStatus.pending;
- private static final ReportStatus PROCESSING = ReportStatus.processing;
- private static final ReportStatus CREATED = ReportStatus.created;
- private static final ReportStatus CANCELED = ReportStatus.canceled;
- private static final ReportStatus TIMED_OUT = ReportStatus.timed_out;
- private static final ReportStatus EXPIRED = ReportStatus.expired;
private static final String WORKER_TIMEOUT_CODE = "worker_timeout";
- private static final String WORKER_TIMEOUT_MESSAGE = "Report processing exceeded stale timeout";
+ private static final String WORKER_TIMEOUT_MESSAGE = "Report processing exceeded maximum duration";
private final DSLContext dslContext;
- public Optional claimNextPendingReport(Instant now) {
+ public Optional claimNextPendingReport(Instant now) {
var candidate = DSL.table(DSL.name("candidate"));
+ var claimTime = toLocalDateTime(now);
return dslContext.with("candidate").as(
dslContext.select(REPORT_JOB.ID)
.from(REPORT_JOB)
- .where(REPORT_JOB.STATUS.eq(PENDING))
- .and(isReadyForClaim(now))
+ .where(REPORT_JOB.STATUS.eq(ReportStatus.pending))
+ .and(REPORT_JOB.NEXT_ATTEMPT_AT.isNull()
+ .or(REPORT_JOB.NEXT_ATTEMPT_AT.le(claimTime)))
.orderBy(REPORT_JOB.CREATED_AT.asc(), REPORT_JOB.ID.asc())
.limit(1)
.forUpdate()
.skipLocked()
)
.update(REPORT_JOB)
- .set(REPORT_JOB.STATUS, PROCESSING)
+ .set(REPORT_JOB.STATUS, ReportStatus.processing)
.set(REPORT_JOB.ATTEMPT, REPORT_JOB.ATTEMPT.plus(1))
- .set(REPORT_JOB.STARTED_AT, firstTimestampOrExisting(now, REPORT_JOB.STARTED_AT))
+ .set(REPORT_JOB.STARTED_AT, claimTime)
.set(REPORT_JOB.NEXT_ATTEMPT_AT, (LocalDateTime) null)
- .set(REPORT_JOB.UPDATED_AT, timestampValue(now, REPORT_JOB.UPDATED_AT))
+ .set(REPORT_JOB.ERROR_CODE, (String) null)
+ .set(REPORT_JOB.ERROR_MESSAGE, (String) null)
.from(candidate)
.where(REPORT_JOB.ID.eq(CANDIDATE_ID))
- .returning(REPORT_JOB.fields())
- .fetchOptional(ReportRecordMapper::mapReportJob);
+ .returning(
+ REPORT_JOB.ID,
+ REPORT_JOB.REPORT_TYPE,
+ REPORT_JOB.QUERY_JSON,
+ REPORT_JOB.TIMEZONE,
+ REPORT_JOB.ATTEMPT
+ )
+ .fetchOptional(record -> new ReportTask(
+ record.get(REPORT_JOB.ID),
+ record.get(REPORT_JOB.REPORT_TYPE),
+ record.get(REPORT_JOB.QUERY_JSON).data(),
+ record.get(REPORT_JOB.TIMEZONE),
+ record.get(REPORT_JOB.ATTEMPT)
+ ));
}
public boolean cancelReport(String createdBy, long reportId, Instant now) {
- var updated = dslContext.update(REPORT_JOB)
- .set(REPORT_JOB.STATUS, CANCELED)
- .set(REPORT_JOB.FINISHED_AT, firstTimestampOrExisting(now, REPORT_JOB.FINISHED_AT))
- .set(REPORT_JOB.UPDATED_AT, timestampValue(now, REPORT_JOB.UPDATED_AT))
+ return dslContext.update(REPORT_JOB)
+ .set(REPORT_JOB.STATUS, ReportStatus.canceled)
+ .set(REPORT_JOB.FINISHED_AT, toLocalDateTime(now))
+ .set(REPORT_JOB.NEXT_ATTEMPT_AT, (LocalDateTime) null)
.where(REPORT_JOB.ID.eq(reportId))
.and(REPORT_JOB.CREATED_BY.eq(createdBy))
- .and(REPORT_JOB.STATUS.in(PENDING, PROCESSING))
- .execute();
- return updated > 0;
+ .and(REPORT_JOB.STATUS.in(ReportStatus.pending, ReportStatus.processing))
+ .execute() == 1;
}
public boolean rescheduleForRetry(long reportId, Instant nextAttemptAt, String errorCode, String errorMessage) {
- var updated = dslContext.update(REPORT_JOB)
- .set(REPORT_JOB.STATUS, PENDING)
- .set(REPORT_JOB.NEXT_ATTEMPT_AT, timestampValue(nextAttemptAt, REPORT_JOB.NEXT_ATTEMPT_AT))
+ return dslContext.update(REPORT_JOB)
+ .set(REPORT_JOB.STATUS, ReportStatus.pending)
+ .set(REPORT_JOB.STARTED_AT, (LocalDateTime) null)
+ .set(REPORT_JOB.NEXT_ATTEMPT_AT, toLocalDateTime(nextAttemptAt))
.set(REPORT_JOB.ERROR_CODE, errorCode)
.set(REPORT_JOB.ERROR_MESSAGE, errorMessage)
- .set(REPORT_JOB.UPDATED_AT, timestampValue(nextAttemptAt, REPORT_JOB.UPDATED_AT))
.where(REPORT_JOB.ID.eq(reportId))
- .and(REPORT_JOB.STATUS.eq(PROCESSING))
- .execute();
- return updated > 0;
+ .and(REPORT_JOB.STATUS.eq(ReportStatus.processing))
+ .execute() == 1;
}
- public boolean markCreated(
+ public boolean markFailed(long reportId, Instant finishedAt, String code, String message) {
+ return dslContext.update(REPORT_JOB)
+ .set(REPORT_JOB.STATUS, ReportStatus.failed)
+ .set(REPORT_JOB.FINISHED_AT, toLocalDateTime(finishedAt))
+ .set(REPORT_JOB.ERROR_CODE, code)
+ .set(REPORT_JOB.ERROR_MESSAGE, message)
+ .set(REPORT_JOB.NEXT_ATTEMPT_AT, (LocalDateTime) null)
+ .where(REPORT_JOB.ID.eq(reportId))
+ .and(REPORT_JOB.STATUS.eq(ReportStatus.processing))
+ .execute() == 1;
+ }
+
+ @Transactional
+ public boolean completeReport(
long reportId,
+ ReportFile reportFile,
Instant dataSnapshotFixedAt,
Instant finishedAt,
Instant expiresAt,
long rowsCount
) {
- var updated = dslContext.update(REPORT_JOB)
- .set(REPORT_JOB.STATUS, CREATED)
- .set(REPORT_JOB.DATA_SNAPSHOT_FIXED_AT,
- firstTimestampOrExisting(dataSnapshotFixedAt, REPORT_JOB.DATA_SNAPSHOT_FIXED_AT))
- .set(REPORT_JOB.FINISHED_AT, firstTimestampOrExisting(finishedAt, REPORT_JOB.FINISHED_AT))
- .set(REPORT_JOB.EXPIRES_AT, timestampValue(expiresAt, REPORT_JOB.EXPIRES_AT))
+ var completed = dslContext.update(REPORT_JOB)
+ .set(REPORT_JOB.STATUS, ReportStatus.created)
+ .set(REPORT_JOB.DATA_SNAPSHOT_FIXED_AT, toLocalDateTime(dataSnapshotFixedAt))
+ .set(REPORT_JOB.FINISHED_AT, toLocalDateTime(finishedAt))
+ .set(REPORT_JOB.EXPIRES_AT, toLocalDateTime(expiresAt))
.set(REPORT_JOB.ROWS_COUNT, rowsCount)
.set(REPORT_JOB.ERROR_CODE, (String) null)
.set(REPORT_JOB.ERROR_MESSAGE, (String) null)
.set(REPORT_JOB.NEXT_ATTEMPT_AT, (LocalDateTime) null)
- .set(REPORT_JOB.UPDATED_AT, timestampValue(finishedAt, REPORT_JOB.UPDATED_AT))
.where(REPORT_JOB.ID.eq(reportId))
- .and(REPORT_JOB.STATUS.eq(PROCESSING))
- .execute();
- return updated > 0;
- }
-
- public boolean publishFileRecord(long reportId, ReportFile reportFile, Instant createdAt) {
- var updated = dslContext.insertInto(REPORT_FILE)
- .set(ReportLifecycleMapper.newInsertableFileRecord(dslContext, reportId, reportFile, createdAt))
+ .and(REPORT_JOB.STATUS.eq(ReportStatus.processing))
.execute();
- return updated > 0;
- }
-
- public boolean markFailed(long reportId, Instant dataSnapshotFixedAt, Instant finishedAt, String code,
- String message) {
- return markTerminal(reportId, ReportStatus.failed, dataSnapshotFixedAt, finishedAt, code, message);
- }
-
- public boolean markTimedOut(long reportId, Instant dataSnapshotFixedAt, Instant finishedAt, String code,
- String message) {
- return markTerminal(reportId, ReportStatus.timed_out, dataSnapshotFixedAt, finishedAt, code, message);
- }
-
- public boolean expireReport(long reportId, Instant expiredAt) {
- var updated = dslContext.update(REPORT_JOB)
- .set(REPORT_JOB.STATUS, EXPIRED)
- .set(REPORT_JOB.FINISHED_AT, firstTimestampOrExisting(expiredAt, REPORT_JOB.FINISHED_AT))
- .set(REPORT_JOB.UPDATED_AT, timestampValue(expiredAt, REPORT_JOB.UPDATED_AT))
- .where(REPORT_JOB.ID.eq(reportId))
- .and(REPORT_JOB.STATUS.eq(CREATED))
+ if (completed == 0) {
+ return false;
+ }
+
+ dslContext.insertInto(REPORT_FILE)
+ .columns(
+ REPORT_FILE.REPORT_ID,
+ REPORT_FILE.FILE_ID,
+ REPORT_FILE.FILE_TYPE,
+ REPORT_FILE.FILENAME,
+ REPORT_FILE.CONTENT_TYPE,
+ REPORT_FILE.SIZE_BYTES,
+ REPORT_FILE.MD5,
+ REPORT_FILE.SHA256,
+ REPORT_FILE.CREATED_AT
+ )
+ .values(
+ reportId,
+ reportFile.getFileId(),
+ reportFile.getFileType(),
+ reportFile.getFilename(),
+ reportFile.getContentType(),
+ reportFile.getSizeBytes(),
+ reportFile.getMd5(),
+ reportFile.getSha256(),
+ toLocalDateTime(finishedAt)
+ )
.execute();
- return updated > 0;
+ return true;
}
public int timeoutStaleProcessingReports(Instant staleBefore, Instant finishedAt) {
return dslContext.update(REPORT_JOB)
- .set(REPORT_JOB.STATUS, TIMED_OUT)
- .set(REPORT_JOB.FINISHED_AT, firstTimestampOrExisting(finishedAt, REPORT_JOB.FINISHED_AT))
- .set(REPORT_JOB.ERROR_CODE, firstValueOrExisting(WORKER_TIMEOUT_CODE, REPORT_JOB.ERROR_CODE))
- .set(REPORT_JOB.ERROR_MESSAGE, firstValueOrExisting(WORKER_TIMEOUT_MESSAGE, REPORT_JOB.ERROR_MESSAGE))
+ .set(REPORT_JOB.STATUS, ReportStatus.timed_out)
+ .set(REPORT_JOB.FINISHED_AT, toLocalDateTime(finishedAt))
+ .set(REPORT_JOB.ERROR_CODE, WORKER_TIMEOUT_CODE)
+ .set(REPORT_JOB.ERROR_MESSAGE, WORKER_TIMEOUT_MESSAGE)
.set(REPORT_JOB.NEXT_ATTEMPT_AT, (LocalDateTime) null)
- .set(REPORT_JOB.UPDATED_AT, timestampValue(finishedAt, REPORT_JOB.UPDATED_AT))
- .where(REPORT_JOB.STATUS.eq(PROCESSING))
- .and(REPORT_JOB.UPDATED_AT.le(timestampValue(staleBefore, REPORT_JOB.UPDATED_AT)))
+ .where(REPORT_JOB.STATUS.eq(ReportStatus.processing))
+ .and(REPORT_JOB.STARTED_AT.le(toLocalDateTime(staleBefore)))
.execute();
}
- public int expireReports(Instant now) {
+ public boolean markTimedOut(long reportId, Instant finishedAt) {
return dslContext.update(REPORT_JOB)
- .set(REPORT_JOB.STATUS, EXPIRED)
- .set(
- REPORT_JOB.FINISHED_AT,
- DSL.coalesce(
- REPORT_JOB.FINISHED_AT,
- REPORT_JOB.EXPIRES_AT,
- timestampValue(now, REPORT_JOB.FINISHED_AT)
- )
- )
- .set(REPORT_JOB.UPDATED_AT, timestampValue(now, REPORT_JOB.UPDATED_AT))
- .where(REPORT_JOB.STATUS.eq(CREATED))
- .and(REPORT_JOB.EXPIRES_AT.isNotNull())
- .and(REPORT_JOB.EXPIRES_AT.le(timestampValue(now, REPORT_JOB.EXPIRES_AT)))
- .execute();
- }
-
- private boolean markTerminal(
- long reportId,
- ReportStatus terminalStatus,
- Instant dataSnapshotFixedAt,
- Instant finishedAt,
- String code,
- String message
- ) {
- var updated = dslContext.update(REPORT_JOB)
- .set(REPORT_JOB.STATUS, terminalStatus)
- .set(REPORT_JOB.DATA_SNAPSHOT_FIXED_AT,
- firstTimestampOrExisting(dataSnapshotFixedAt, REPORT_JOB.DATA_SNAPSHOT_FIXED_AT))
- .set(REPORT_JOB.FINISHED_AT, firstTimestampOrExisting(finishedAt, REPORT_JOB.FINISHED_AT))
- .set(REPORT_JOB.ERROR_CODE, code)
- .set(REPORT_JOB.ERROR_MESSAGE, message)
+ .set(REPORT_JOB.STATUS, ReportStatus.timed_out)
+ .set(REPORT_JOB.FINISHED_AT, toLocalDateTime(finishedAt))
+ .set(REPORT_JOB.ERROR_CODE, WORKER_TIMEOUT_CODE)
+ .set(REPORT_JOB.ERROR_MESSAGE, WORKER_TIMEOUT_MESSAGE)
.set(REPORT_JOB.NEXT_ATTEMPT_AT, (LocalDateTime) null)
- .set(REPORT_JOB.UPDATED_AT, timestampValue(finishedAt, REPORT_JOB.UPDATED_AT))
.where(REPORT_JOB.ID.eq(reportId))
- .and(REPORT_JOB.STATUS.eq(PROCESSING))
- .execute();
- return updated > 0;
+ .and(REPORT_JOB.STATUS.eq(ReportStatus.processing))
+ .execute() == 1;
}
- private static Condition isReadyForClaim(Instant now) {
- return isReadyAt(REPORT_JOB.NEXT_ATTEMPT_AT, now);
+ public int expireReports(Instant now) {
+ return dslContext.update(REPORT_JOB)
+ .set(REPORT_JOB.STATUS, ReportStatus.expired)
+ .where(REPORT_JOB.STATUS.eq(ReportStatus.created))
+ .and(REPORT_JOB.EXPIRES_AT.le(toLocalDateTime(now)))
+ .execute();
}
}
diff --git a/src/main/java/dev/vality/ccreporter/dao/ReportQueryDao.java b/src/main/java/dev/vality/ccreporter/dao/ReportQueryDao.java
index d974c65..2d4cf81 100644
--- a/src/main/java/dev/vality/ccreporter/dao/ReportQueryDao.java
+++ b/src/main/java/dev/vality/ccreporter/dao/ReportQueryDao.java
@@ -2,7 +2,8 @@
import dev.vality.ccreporter.GetReportsFilter;
import dev.vality.ccreporter.dao.mapper.ReportRecordMapper;
-import dev.vality.ccreporter.domain.tables.pojos.ReportFile;
+import dev.vality.ccreporter.domain.enums.ReportStatus;
+import dev.vality.ccreporter.model.DownloadableFile;
import dev.vality.ccreporter.model.ReportProjection;
import dev.vality.ccreporter.serde.json.ContinuationTokenJsonSerializer.PageCursor;
import lombok.RequiredArgsConstructor;
@@ -12,14 +13,14 @@
import org.jooq.SelectJoinStep;
import org.springframework.stereotype.Repository;
+import java.time.Instant;
import java.util.ArrayList;
import java.util.List;
import java.util.Optional;
import static dev.vality.ccreporter.domain.Tables.REPORT_FILE;
import static dev.vality.ccreporter.domain.Tables.REPORT_JOB;
-import static dev.vality.ccreporter.util.TimestampUtils.parse;
-import static dev.vality.ccreporter.util.TimestampUtils.toLocalDateTime;
+import static dev.vality.ccreporter.util.TimestampUtils.*;
@Repository
@RequiredArgsConstructor
@@ -42,11 +43,19 @@ public List getReports(String createdBy, GetReportsFilter filt
.fetch(ReportRecordMapper::mapReportProjection);
}
- public Optional getFile(String createdBy, String fileId) {
- return baseFileSelect()
+ public Optional getDownloadableFile(String createdBy, String fileId, Instant now) {
+ return dslContext.select(REPORT_FILE.fields())
+ .select(REPORT_JOB.EXPIRES_AT)
+ .from(REPORT_FILE)
+ .join(REPORT_JOB).on(REPORT_JOB.ID.eq(REPORT_FILE.REPORT_ID))
.where(REPORT_FILE.FILE_ID.eq(fileId))
.and(REPORT_JOB.CREATED_BY.eq(createdBy))
- .fetchOptional(ReportRecordMapper::mapReportFile);
+ .and(REPORT_JOB.STATUS.eq(ReportStatus.created))
+ .and(REPORT_JOB.EXPIRES_AT.gt(toLocalDateTime(now)))
+ .fetchOptional(record -> new DownloadableFile(
+ ReportRecordMapper.mapReportFile(record),
+ toInstant(record.get(REPORT_JOB.EXPIRES_AT))
+ ));
}
private List buildReportConditions(String createdBy, GetReportsFilter filter, PageCursor cursor) {
@@ -114,9 +123,4 @@ private SelectJoinStep baseReportSelect() {
.leftJoin(REPORT_FILE).on(REPORT_FILE.REPORT_ID.eq(REPORT_JOB.ID));
}
- private SelectJoinStep baseFileSelect() {
- return dslContext.select(REPORT_FILE.fields())
- .from(REPORT_FILE)
- .join(REPORT_JOB).on(REPORT_JOB.ID.eq(REPORT_FILE.REPORT_ID));
- }
}
diff --git a/src/main/java/dev/vality/ccreporter/dao/WithdrawalSessionDao.java b/src/main/java/dev/vality/ccreporter/dao/WithdrawalSessionDao.java
index d0dbb96..773ad49 100644
--- a/src/main/java/dev/vality/ccreporter/dao/WithdrawalSessionDao.java
+++ b/src/main/java/dev/vality/ccreporter/dao/WithdrawalSessionDao.java
@@ -42,10 +42,7 @@ var record = dslContext.newRecord(WITHDRAWAL_SESSION, update);
OVERWRITE_FIELDS,
Map.of(WITHDRAWAL_SESSION.UPDATED_AT, UTC_NOW)
))
- .where(isIncomingEventNewer(
- WITHDRAWAL_SESSION.DOMAIN_EVENT_CREATED_AT,
- WITHDRAWAL_SESSION.DOMAIN_EVENT_ID
- ))
+ .where(isIncomingEventNewer(WITHDRAWAL_SESSION.DOMAIN_EVENT_ID))
.execute();
}
}
diff --git a/src/main/java/dev/vality/ccreporter/dao/WithdrawalTxnCurrentDao.java b/src/main/java/dev/vality/ccreporter/dao/WithdrawalTxnCurrentDao.java
index 8a3ba68..5408a65 100644
--- a/src/main/java/dev/vality/ccreporter/dao/WithdrawalTxnCurrentDao.java
+++ b/src/main/java/dev/vality/ccreporter/dao/WithdrawalTxnCurrentDao.java
@@ -4,6 +4,7 @@
import lombok.RequiredArgsConstructor;
import org.jooq.DSLContext;
import org.jooq.Field;
+import org.jooq.impl.DSL;
import org.springframework.stereotype.Repository;
import java.util.Map;
@@ -17,8 +18,7 @@
public class WithdrawalTxnCurrentDao {
private static final Set> IMMUTABLE_FIELDS = Set.of(
- WITHDRAWAL_TXN_CURRENT.WITHDRAWAL_ID,
- WITHDRAWAL_TXN_CURRENT.UPDATED_AT
+ WITHDRAWAL_TXN_CURRENT.WITHDRAWAL_ID
);
private static final Set> OVERWRITE_FIELDS = Set.of(
@@ -40,12 +40,22 @@ var record = dslContext.newRecord(WITHDRAWAL_TXN_CURRENT, update);
WITHDRAWAL_TXN_CURRENT,
IMMUTABLE_FIELDS,
OVERWRITE_FIELDS,
- Map.of(WITHDRAWAL_TXN_CURRENT.UPDATED_AT, UTC_NOW)
- ))
- .where(isIncomingEventNewer(
- WITHDRAWAL_TXN_CURRENT.DOMAIN_EVENT_CREATED_AT,
- WITHDRAWAL_TXN_CURRENT.DOMAIN_EVENT_ID
+ Map.of(
+ WITHDRAWAL_TXN_CURRENT.UPDATED_AT,
+ UTC_NOW,
+ WITHDRAWAL_TXN_CURRENT.FINALIZED_AT,
+ DSL.when(
+ DSL.excluded(WITHDRAWAL_TXN_CURRENT.STATUS).isNotNull(),
+ DSL.excluded(WITHDRAWAL_TXN_CURRENT.FINALIZED_AT)
+ ).otherwise(WITHDRAWAL_TXN_CURRENT.FINALIZED_AT),
+ WITHDRAWAL_TXN_CURRENT.ERROR_SUMMARY,
+ DSL.when(
+ DSL.excluded(WITHDRAWAL_TXN_CURRENT.STATUS).isNotNull(),
+ DSL.excluded(WITHDRAWAL_TXN_CURRENT.ERROR_SUMMARY)
+ ).otherwise(WITHDRAWAL_TXN_CURRENT.ERROR_SUMMARY)
+ )
))
+ .where(isIncomingEventNewer(WITHDRAWAL_TXN_CURRENT.DOMAIN_EVENT_ID))
.execute();
}
-}
\ No newline at end of file
+}
diff --git a/src/main/java/dev/vality/ccreporter/dao/mapper/ReportCommandMapper.java b/src/main/java/dev/vality/ccreporter/dao/mapper/ReportCommandMapper.java
deleted file mode 100644
index 2f65025..0000000
--- a/src/main/java/dev/vality/ccreporter/dao/mapper/ReportCommandMapper.java
+++ /dev/null
@@ -1,53 +0,0 @@
-package dev.vality.ccreporter.dao.mapper;
-
-import dev.vality.ccreporter.FileType;
-import dev.vality.ccreporter.ReportType;
-import dev.vality.ccreporter.domain.tables.pojos.ReportJob;
-import dev.vality.ccreporter.domain.tables.records.ReportJobRecord;
-import lombok.experimental.UtilityClass;
-import org.jooq.DSLContext;
-import org.jooq.JSONB;
-import org.springframework.util.StringUtils;
-
-import java.time.Instant;
-
-import static dev.vality.ccreporter.domain.Tables.REPORT_JOB;
-import static dev.vality.ccreporter.util.TimestampUtils.toLocalDateTime;
-
-@UtilityClass
-public class ReportCommandMapper {
-
- public static ReportJob mapReportJob(
- String createdBy,
- ReportType reportType,
- FileType fileType,
- String queryJson,
- String queryHash,
- Instant requestedTimeFrom,
- Instant requestedTimeTo,
- String timezone,
- String idempotencyKey
- ) {
- return new ReportJob()
- .setReportType(
- ReportRecordMapper.mapEnum(reportType, dev.vality.ccreporter.domain.enums.ReportType.class))
- .setFileType(ReportRecordMapper.mapEnum(fileType, dev.vality.ccreporter.domain.enums.FileType.class))
- .setQueryJson(JSONB.jsonb(queryJson))
- .setQueryHash(queryHash)
- .setRequestedTimeFrom(toLocalDateTime(requestedTimeFrom))
- .setRequestedTimeTo(toLocalDateTime(requestedTimeTo))
- .setTimezone(timezone)
- .setCreatedBy(createdBy)
- .setIdempotencyKey(StringUtils.hasText(idempotencyKey) ? idempotencyKey : null);
- }
-
- public static ReportJobRecord newInsertableRecord(DSLContext dslContext, ReportJob reportJob) {
- var record = dslContext.newRecord(REPORT_JOB, reportJob);
- record.changed(REPORT_JOB.ID, false);
- record.changed(REPORT_JOB.STATUS, false);
- record.changed(REPORT_JOB.ATTEMPT, false);
- record.changed(REPORT_JOB.CREATED_AT, false);
- record.changed(REPORT_JOB.UPDATED_AT, false);
- return record;
- }
-}
diff --git a/src/main/java/dev/vality/ccreporter/dao/mapper/ReportLifecycleMapper.java b/src/main/java/dev/vality/ccreporter/dao/mapper/ReportLifecycleMapper.java
deleted file mode 100644
index 89f570e..0000000
--- a/src/main/java/dev/vality/ccreporter/dao/mapper/ReportLifecycleMapper.java
+++ /dev/null
@@ -1,33 +0,0 @@
-package dev.vality.ccreporter.dao.mapper;
-
-import dev.vality.ccreporter.domain.tables.pojos.ReportFile;
-import dev.vality.ccreporter.domain.tables.records.ReportFileRecord;
-import lombok.experimental.UtilityClass;
-import org.jooq.DSLContext;
-
-import java.time.Instant;
-import java.time.ZoneOffset;
-
-import static dev.vality.ccreporter.domain.Tables.REPORT_FILE;
-
-@UtilityClass
-public class ReportLifecycleMapper {
-
- private static final dev.vality.ccreporter.domain.enums.FileType CSV_FILE_TYPE =
- dev.vality.ccreporter.domain.enums.FileType.csv;
-
- public static ReportFileRecord newInsertableFileRecord(
- DSLContext dslContext,
- long reportId,
- ReportFile reportFile,
- Instant createdAt
- ) {
- var insertableReportFile = new ReportFile(reportFile)
- .setReportId(reportId)
- .setFileType(CSV_FILE_TYPE)
- .setCreatedAt(createdAt.atZone(ZoneOffset.UTC).toLocalDateTime());
- var record = dslContext.newRecord(REPORT_FILE, insertableReportFile);
- record.changed(REPORT_FILE.ID, false);
- return record;
- }
-}
diff --git a/src/main/java/dev/vality/ccreporter/dao/support/DaoUpsertUtils.java b/src/main/java/dev/vality/ccreporter/dao/support/DaoUpsertUtils.java
index 5242fdb..1a7ec8a 100644
--- a/src/main/java/dev/vality/ccreporter/dao/support/DaoUpsertUtils.java
+++ b/src/main/java/dev/vality/ccreporter/dao/support/DaoUpsertUtils.java
@@ -59,14 +59,7 @@ public static Map, Object> buildLookupUpsertMap(
);
}
- public static org.jooq.Condition isIncomingEventNewer(
- Field createdAtField,
- Field eventIdField
- ) {
- return DSL.excluded(createdAtField).gt(createdAtField)
- .or(
- DSL.excluded(createdAtField).eq(createdAtField)
- .and(DSL.excluded(eventIdField).gt(eventIdField))
- );
+ public static org.jooq.Condition isIncomingEventNewer(Field eventIdField) {
+ return DSL.excluded(eventIdField).gt(eventIdField);
}
}
diff --git a/src/main/java/dev/vality/ccreporter/dao/support/ReportDaoSupport.java b/src/main/java/dev/vality/ccreporter/dao/support/ReportDaoSupport.java
deleted file mode 100644
index ab8230c..0000000
--- a/src/main/java/dev/vality/ccreporter/dao/support/ReportDaoSupport.java
+++ /dev/null
@@ -1,37 +0,0 @@
-package dev.vality.ccreporter.dao.support;
-
-import lombok.experimental.UtilityClass;
-import org.jooq.Condition;
-import org.jooq.Field;
-import org.jooq.impl.DSL;
-
-import java.sql.Timestamp;
-import java.time.Instant;
-import java.time.LocalDateTime;
-
-@UtilityClass
-public class ReportDaoSupport {
-
- public static Field firstValueOrExisting(T value, Field field) {
- return DSL.coalesce(field, DSL.val(value, field));
- }
-
- public static Field firstTimestampOrExisting(Instant value, Field field) {
- return firstValueOrExisting(toLocalDateTime(value), field);
- }
-
- public static Condition isReadyAt(Field field, Instant now) {
- return field.isNull().or(field.le(timestampValue(now, field)));
- }
-
- public static Field timestampValue(Instant value, Field field) {
- if (value == null) {
- return DSL.castNull(field.getDataType());
- }
- return DSL.val(Timestamp.from(value)).cast(field.getDataType());
- }
-
- private static LocalDateTime toLocalDateTime(Instant value) {
- return value == null ? null : Timestamp.from(value).toLocalDateTime();
- }
-}
diff --git a/src/main/java/dev/vality/ccreporter/ingestion/payment/PaymentEventProjector.java b/src/main/java/dev/vality/ccreporter/ingestion/payment/PaymentEventProjector.java
index 1e74122..bf226e4 100644
--- a/src/main/java/dev/vality/ccreporter/ingestion/payment/PaymentEventProjector.java
+++ b/src/main/java/dev/vality/ccreporter/ingestion/payment/PaymentEventProjector.java
@@ -15,14 +15,15 @@
import org.springframework.stereotype.Component;
import java.time.Instant;
-import java.util.ArrayList;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Optional;
+import java.util.function.Consumer;
import static dev.vality.ccreporter.ingestion.shared.status.StatusDetailExtractor.*;
import static dev.vality.ccreporter.util.SearchValueNormalizer.normalize;
import static dev.vality.ccreporter.util.TimestampUtils.toLocalDateTime;
-import static dev.vality.ccreporter.util.TimestampUtils.toOptionalLocalDateTime;
+import static dev.vality.ccreporter.util.TimestampUtils.toNullableLocalDateTime;
@Component
@RequiredArgsConstructor
@@ -31,16 +32,20 @@ public class PaymentEventProjector {
private final ProxyStateExtractor proxyStateExtractor;
public List project(MachineEvent event, EventPayload payload) {
- var updates = new ArrayList();
- if (!payload.isSetInvoiceChanges()) {
- return updates;
+ if (payload == null || !payload.isSetInvoiceChanges()) {
+ return List.of();
}
+ var updatesByPaymentId = new LinkedHashMap();
for (InvoiceChange change : payload.getInvoiceChanges()) {
if (change.isSetInvoicePaymentChange()) {
- projectPaymentChange(event, change).ifPresent(updates::add);
+ projectPaymentChange(event, change).ifPresent(update -> updatesByPaymentId.merge(
+ update.getPaymentId(),
+ update,
+ this::mergeUpdates
+ ));
}
}
- return updates;
+ return List.copyOf(updatesByPaymentId.values());
}
private Optional projectPaymentChange(MachineEvent event, InvoiceChange change) {
@@ -134,7 +139,7 @@ private Optional paymentStatusChangedUpdate(
var capturedCost = extractCapturedCost(status);
return Optional.of(baseUpdate(event, paymentChange)
.setFinalizedAt(
- toOptionalLocalDateTime(terminalFinalizedAt(status, Instant.parse(event.getCreatedAt()))))
+ toNullableLocalDateTime(terminalFinalizedAt(status, Instant.parse(event.getCreatedAt()))))
.setStatus(status.getSetField().getFieldName())
.setErrorSummary(extractErrorSummary(status))
.setProviderAmount(capturedCost != null ? capturedCost.getAmount() : null)
@@ -187,6 +192,41 @@ private PaymentTxnCurrent baseUpdate(MachineEvent event, InvoicePaymentChange pa
.setDomainEventCreatedAt(toLocalDateTime(event.getCreatedAt()));
}
+ private PaymentTxnCurrent mergeUpdates(PaymentTxnCurrent accumulated, PaymentTxnCurrent update) {
+ copyIfPresent(update.getPartyId(), accumulated::setPartyId);
+ copyIfPresent(update.getShopId(), accumulated::setShopId);
+ copyIfPresent(update.getCreatedAt(), accumulated::setCreatedAt);
+ if (update.getStatus() != null) {
+ accumulated.setStatus(update.getStatus());
+ accumulated.setFinalizedAt(update.getFinalizedAt());
+ accumulated.setErrorSummary(update.getErrorSummary());
+ }
+ copyIfPresent(update.getProviderId(), accumulated::setProviderId);
+ copyIfPresent(update.getTerminalId(), accumulated::setTerminalId);
+ copyIfPresent(update.getAmount(), accumulated::setAmount);
+ copyIfPresent(update.getFee(), accumulated::setFee);
+ copyIfPresent(update.getCurrency(), accumulated::setCurrency);
+ copyIfPresent(update.getTrxId(), accumulated::setTrxId);
+ copyIfPresent(update.getExternalId(), accumulated::setExternalId);
+ copyIfPresent(update.getRrn(), accumulated::setRrn);
+ copyIfPresent(update.getApprovalCode(), accumulated::setApprovalCode);
+ copyIfPresent(update.getPaymentToolType(), accumulated::setPaymentToolType);
+ copyIfPresent(update.getOriginalAmount(), accumulated::setOriginalAmount);
+ copyIfPresent(update.getOriginalCurrency(), accumulated::setOriginalCurrency);
+ copyIfPresent(update.getConvertedAmount(), accumulated::setConvertedAmount);
+ copyIfPresent(update.getExchangeRateInternal(), accumulated::setExchangeRateInternal);
+ copyIfPresent(update.getProviderAmount(), accumulated::setProviderAmount);
+ copyIfPresent(update.getProviderCurrency(), accumulated::setProviderCurrency);
+ copyIfPresent(update.getTrxSearch(), accumulated::setTrxSearch);
+ return accumulated;
+ }
+
+ private void copyIfPresent(T value, Consumer setter) {
+ if (value != null) {
+ setter.accept(value);
+ }
+ }
+
private SessionChangePayload sessionPayload(InvoicePaymentChange paymentChange) {
if (!paymentChange.getPayload().isSetInvoicePaymentSessionChange()) {
return null;
diff --git a/src/main/java/dev/vality/ccreporter/ingestion/withdrawal/WithdrawalEventProjector.java b/src/main/java/dev/vality/ccreporter/ingestion/withdrawal/WithdrawalEventProjector.java
index cb681be..89aada5 100644
--- a/src/main/java/dev/vality/ccreporter/ingestion/withdrawal/WithdrawalEventProjector.java
+++ b/src/main/java/dev/vality/ccreporter/ingestion/withdrawal/WithdrawalEventProjector.java
@@ -19,7 +19,7 @@
import static dev.vality.ccreporter.ingestion.shared.status.StatusDetailExtractor.PENDING_STATUS;
import static dev.vality.ccreporter.ingestion.shared.status.StatusDetailExtractor.extractErrorSummary;
import static dev.vality.ccreporter.util.TimestampUtils.toLocalDateTime;
-import static dev.vality.ccreporter.util.TimestampUtils.toOptionalLocalDateTime;
+import static dev.vality.ccreporter.util.TimestampUtils.toNullableLocalDateTime;
@Component
public class WithdrawalEventProjector {
@@ -35,6 +35,7 @@ public List project(MachineEvent event, TimestampedChange
private Optional projectChange(MachineEvent event, Change change) {
return createdUpdate(event, change)
+ .or(() -> bodyChangedUpdate(event, change))
.or(() -> routeChangedUpdate(event, change))
.or(() -> statusChangedUpdate(event, change))
.or(() -> transferCashFlowUpdate(event, change));
@@ -62,11 +63,22 @@ private Optional createdUpdate(MachineEvent event, Change
.setExternalId(withdrawal.getExternalId())
.setOriginalAmount(quote != null ? quote.getCashFrom().getAmount() : null)
.setOriginalCurrency(quote != null ? quote.getCashFrom().getCurrency().getSymbolicCode() : null)
+ .setConvertedAmount(quote != null ? body.getAmount() : null)
.setExchangeRateInternal(toRate(quote))
.setProviderAmount(quote != null ? quote.getCashTo().getAmount() : null)
.setProviderCurrency(quote != null ? quote.getCashTo().getCurrency().getSymbolicCode() : null));
}
+ private Optional bodyChangedUpdate(MachineEvent event, Change change) {
+ if (!change.isSetBodyChanged()) {
+ return Optional.empty();
+ }
+ var body = change.getBodyChanged().getNewBody();
+ return Optional.of(baseUpdate(event)
+ .setAmount(body.getAmount())
+ .setCurrency(body.getCurrency().getSymbolicCode()));
+ }
+
private Optional routeChangedUpdate(MachineEvent event, Change change) {
if (!change.isSetRoute()) {
return Optional.empty();
@@ -85,7 +97,7 @@ private Optional statusChangedUpdate(MachineEvent event, C
return Optional.of(baseUpdate(event)
.setStatus(status.getSetField().getFieldName())
.setFinalizedAt(
- toOptionalLocalDateTime(terminalFinalizedAt(status, Instant.parse(event.getCreatedAt()))))
+ toNullableLocalDateTime(terminalFinalizedAt(status, Instant.parse(event.getCreatedAt()))))
.setErrorSummary(extractErrorSummary(status)));
}
diff --git a/src/main/java/dev/vality/ccreporter/kafka/listener/DominantEventListener.java b/src/main/java/dev/vality/ccreporter/kafka/listener/DominantEventListener.java
index 2133e30..b47c3a5 100644
--- a/src/main/java/dev/vality/ccreporter/kafka/listener/DominantEventListener.java
+++ b/src/main/java/dev/vality/ccreporter/kafka/listener/DominantEventListener.java
@@ -14,13 +14,13 @@
@Component
@RequiredArgsConstructor
-@ConditionalOnProperty(prefix = "ccr.kafka.topics.dominant", name = "enabled", havingValue = "true")
+@ConditionalOnProperty(prefix = "kafka.topics.dominant", name = "enabled", havingValue = "true")
public class DominantEventListener implements BatchLoggingKafkaListener {
private final DominantLookupIngestionService dominantLookupIngestionService;
@KafkaListener(
- topics = "${ccr.kafka.topics.dominant.id}",
+ topics = "${kafka.topics.dominant.id}",
containerFactory = "dominantKafkaListenerContainerFactory"
)
public void listen(List> batch, Acknowledgment acknowledgment) {
diff --git a/src/main/java/dev/vality/ccreporter/kafka/listener/PaymentEventListener.java b/src/main/java/dev/vality/ccreporter/kafka/listener/PaymentEventListener.java
index f52eb2e..8c03108 100644
--- a/src/main/java/dev/vality/ccreporter/kafka/listener/PaymentEventListener.java
+++ b/src/main/java/dev/vality/ccreporter/kafka/listener/PaymentEventListener.java
@@ -14,13 +14,13 @@
@Component
@RequiredArgsConstructor
-@ConditionalOnProperty(prefix = "ccr.kafka.topics.payments", name = "enabled", havingValue = "true")
+@ConditionalOnProperty(prefix = "kafka.topics.payments", name = "enabled", havingValue = "true")
public class PaymentEventListener implements BatchLoggingKafkaListener {
private final PaymentIngestionService paymentIngestionService;
@KafkaListener(
- topics = "${ccr.kafka.topics.payments.id}",
+ topics = "${kafka.topics.payments.id}",
containerFactory = "paymentsKafkaListenerContainerFactory"
)
public void listen(List> batch, Acknowledgment acknowledgment) {
diff --git a/src/main/java/dev/vality/ccreporter/kafka/listener/WithdrawalEventListener.java b/src/main/java/dev/vality/ccreporter/kafka/listener/WithdrawalEventListener.java
index adeb33e..e225a96 100644
--- a/src/main/java/dev/vality/ccreporter/kafka/listener/WithdrawalEventListener.java
+++ b/src/main/java/dev/vality/ccreporter/kafka/listener/WithdrawalEventListener.java
@@ -14,13 +14,13 @@
@Component
@RequiredArgsConstructor
-@ConditionalOnProperty(prefix = "ccr.kafka.topics.withdrawals", name = "enabled", havingValue = "true")
+@ConditionalOnProperty(prefix = "kafka.topics.withdrawals", name = "enabled", havingValue = "true")
public class WithdrawalEventListener implements BatchLoggingKafkaListener {
private final WithdrawalIngestionService withdrawalIngestionService;
@KafkaListener(
- topics = "${ccr.kafka.topics.withdrawals.id}",
+ topics = "${kafka.topics.withdrawals.id}",
containerFactory = "withdrawalsKafkaListenerContainerFactory"
)
public void listen(List> batch, Acknowledgment acknowledgment) {
diff --git a/src/main/java/dev/vality/ccreporter/kafka/listener/WithdrawalSessionEventListener.java b/src/main/java/dev/vality/ccreporter/kafka/listener/WithdrawalSessionEventListener.java
index 433b03c..37e12bb 100644
--- a/src/main/java/dev/vality/ccreporter/kafka/listener/WithdrawalSessionEventListener.java
+++ b/src/main/java/dev/vality/ccreporter/kafka/listener/WithdrawalSessionEventListener.java
@@ -14,13 +14,13 @@
@Component
@RequiredArgsConstructor
-@ConditionalOnProperty(prefix = "ccr.kafka.topics.withdrawal-sessions", name = "enabled", havingValue = "true")
+@ConditionalOnProperty(prefix = "kafka.topics.withdrawal-sessions", name = "enabled", havingValue = "true")
public class WithdrawalSessionEventListener implements BatchLoggingKafkaListener {
private final WithdrawalSessionIngestionService withdrawalSessionIngestionService;
@KafkaListener(
- topics = "${ccr.kafka.topics.withdrawal-sessions.id}",
+ topics = "${kafka.topics.withdrawal-sessions.id}",
containerFactory = "withdrawalSessionsKafkaListenerContainerFactory"
)
public void listen(List> batch, Acknowledgment acknowledgment) {
diff --git a/src/main/java/dev/vality/ccreporter/model/DownloadableFile.java b/src/main/java/dev/vality/ccreporter/model/DownloadableFile.java
new file mode 100644
index 0000000..c9980eb
--- /dev/null
+++ b/src/main/java/dev/vality/ccreporter/model/DownloadableFile.java
@@ -0,0 +1,8 @@
+package dev.vality.ccreporter.model;
+
+import dev.vality.ccreporter.domain.tables.pojos.ReportFile;
+
+import java.time.Instant;
+
+public record DownloadableFile(ReportFile file, Instant reportExpiresAt) {
+}
diff --git a/src/main/java/dev/vality/ccreporter/model/ReportTask.java b/src/main/java/dev/vality/ccreporter/model/ReportTask.java
new file mode 100644
index 0000000..f532de1
--- /dev/null
+++ b/src/main/java/dev/vality/ccreporter/model/ReportTask.java
@@ -0,0 +1,12 @@
+package dev.vality.ccreporter.model;
+
+import dev.vality.ccreporter.domain.enums.ReportType;
+
+public record ReportTask(
+ long id,
+ ReportType reportType,
+ String queryJson,
+ String timezone,
+ int attempt
+) {
+}
diff --git a/src/main/java/dev/vality/ccreporter/report/ReportCsvService.java b/src/main/java/dev/vality/ccreporter/report/ReportCsvService.java
index ac70903..916211d 100644
--- a/src/main/java/dev/vality/ccreporter/report/ReportCsvService.java
+++ b/src/main/java/dev/vality/ccreporter/report/ReportCsvService.java
@@ -2,15 +2,13 @@
import dev.vality.ccreporter.PaymentsQuery;
import dev.vality.ccreporter.ReportQuery;
-import dev.vality.ccreporter.ReportType;
import dev.vality.ccreporter.WithdrawalsQuery;
+import dev.vality.ccreporter.config.properties.ReportProperties;
import dev.vality.ccreporter.dao.ReportCsvDao;
-import dev.vality.ccreporter.dao.mapper.ReportRecordMapper;
-import dev.vality.ccreporter.domain.tables.pojos.ReportJob;
import dev.vality.ccreporter.model.GeneratedCsvReport;
+import dev.vality.ccreporter.model.ReportTask;
import dev.vality.ccreporter.serde.json.ThriftJsonCodec;
import lombok.RequiredArgsConstructor;
-import lombok.SneakyThrows;
import org.jooq.Cursor;
import org.jooq.Record;
import org.springframework.stereotype.Service;
@@ -27,77 +25,106 @@
import java.nio.file.Path;
import java.security.DigestOutputStream;
import java.security.MessageDigest;
-import java.sql.Timestamp;
+import java.security.NoSuchAlgorithmException;
+import java.time.LocalDateTime;
import java.time.ZoneId;
+import java.time.ZoneOffset;
import java.time.format.DateTimeFormatter;
import java.util.Currency;
import java.util.HexFormat;
import java.util.List;
import java.util.Locale;
+import java.util.concurrent.CancellationException;
@Service
@RequiredArgsConstructor
public class ReportCsvService {
+ private static final String CSV_LINE_ENDING = "\r\n";
private static final DateTimeFormatter CSV_DATE_FORMATTER = DateTimeFormatter.ISO_LOCAL_DATE;
private static final DateTimeFormatter CSV_TIME_FORMATTER = DateTimeFormatter.ofPattern("HH:mm:ss");
+ private static final String CREATED_DATE_COLUMN = "created_date";
+ private static final String CREATED_TIME_COLUMN = "created_time";
+ private static final String FINALIZED_DATE_COLUMN = "finalized_date";
+ private static final String FINALIZED_TIME_COLUMN = "finalized_time";
+ private static final String INVOICE_ID_COLUMN = "invoice_id";
+ private static final String PAYMENT_ID_COLUMN = "payment_id";
+ private static final String WITHDRAWAL_ID_COLUMN = "withdrawal_id";
+ private static final String STATUS_COLUMN = "status";
+ private static final String AMOUNT_COLUMN = "amount";
+ private static final String CURRENCY_COLUMN = "currency";
+ private static final String TRX_ID_COLUMN = "trx_id";
+ private static final String PROVIDER_ID_COLUMN = "provider_id";
+ private static final String TERMINAL_ID_COLUMN = "terminal_id";
+ private static final String SHOP_ID_COLUMN = "shop_id";
+ private static final String WALLET_ID_COLUMN = "wallet_id";
+ private static final String EXCHANGE_RATE_INTERNAL_COLUMN = "exchange_rate_internal";
+ private static final String PROVIDER_AMOUNT_COLUMN = "provider_amount";
+ private static final String PROVIDER_CURRENCY_COLUMN = "provider_currency";
+ private static final String ORIGINAL_AMOUNT_COLUMN = "original_amount";
+ private static final String ORIGINAL_CURRENCY_COLUMN = "original_currency";
+ private static final String CONVERTED_AMOUNT_COLUMN = "converted_amount";
+
private static final List PAYMENT_COLUMNS = List.of(
- "created_date",
- "created_time",
- "finalized_date",
- "finalized_time",
- "invoice_id",
- "payment_id",
- "status",
- "amount",
- "currency",
- "trx_id",
- "provider_id",
- "terminal_id",
- "shop_id",
- "exchange_rate_internal",
- "provider_amount",
- "provider_currency",
- "original_amount",
- "original_currency",
- "converted_amount"
+ CREATED_DATE_COLUMN,
+ CREATED_TIME_COLUMN,
+ FINALIZED_DATE_COLUMN,
+ FINALIZED_TIME_COLUMN,
+ INVOICE_ID_COLUMN,
+ PAYMENT_ID_COLUMN,
+ STATUS_COLUMN,
+ AMOUNT_COLUMN,
+ CURRENCY_COLUMN,
+ TRX_ID_COLUMN,
+ PROVIDER_ID_COLUMN,
+ TERMINAL_ID_COLUMN,
+ SHOP_ID_COLUMN,
+ EXCHANGE_RATE_INTERNAL_COLUMN,
+ PROVIDER_AMOUNT_COLUMN,
+ PROVIDER_CURRENCY_COLUMN,
+ ORIGINAL_AMOUNT_COLUMN,
+ ORIGINAL_CURRENCY_COLUMN,
+ CONVERTED_AMOUNT_COLUMN
);
private static final List WITHDRAWAL_COLUMNS = List.of(
- "created_date",
- "created_time",
- "finalized_date",
- "finalized_time",
- "withdrawal_id",
- "status",
- "amount",
- "currency",
- "trx_id",
- "provider_id",
- "terminal_id",
- "wallet_id",
- "exchange_rate_internal",
- "provider_amount",
- "provider_currency",
- "original_amount",
- "original_currency"
+ CREATED_DATE_COLUMN,
+ CREATED_TIME_COLUMN,
+ FINALIZED_DATE_COLUMN,
+ FINALIZED_TIME_COLUMN,
+ WITHDRAWAL_ID_COLUMN,
+ STATUS_COLUMN,
+ AMOUNT_COLUMN,
+ CURRENCY_COLUMN,
+ TRX_ID_COLUMN,
+ PROVIDER_ID_COLUMN,
+ TERMINAL_ID_COLUMN,
+ WALLET_ID_COLUMN,
+ EXCHANGE_RATE_INTERNAL_COLUMN,
+ PROVIDER_AMOUNT_COLUMN,
+ PROVIDER_CURRENCY_COLUMN,
+ ORIGINAL_AMOUNT_COLUMN,
+ ORIGINAL_CURRENCY_COLUMN,
+ CONVERTED_AMOUNT_COLUMN
);
private final ReportCsvDao reportCsvDao;
private final ThriftJsonCodec thriftJsonCodec;
+ private final ReportProperties reportProperties;
@Transactional(readOnly = true, isolation = Isolation.REPEATABLE_READ)
- public GeneratedCsvReport generate(ReportJob reportJob) {
+ public GeneratedCsvReport generate(ReportTask reportTask) {
+ reportCsvDao.setLocalStatementTimeout(reportProperties.getProcessingTimeoutMs());
var snapshotFixedAt = reportCsvDao.currentSnapshot();
- var reportQuery = thriftJsonCodec.deserialize(reportJob.getQueryJson().data(), ReportQuery.class);
- var zoneId = ZoneId.of(reportJob.getTimezone());
- var reportType = ReportRecordMapper.mapEnum(reportJob.getReportType(), ReportType.class);
- var fileName = reportType.name() + "-report-" + reportJob.getId() + ".csv";
- var stagedFile = createTempFile(reportJob.getId());
+ var reportQuery = thriftJsonCodec.deserialize(reportTask.queryJson(), ReportQuery.class);
+ var zoneId = ZoneId.of(reportTask.timezone());
+ var reportType = reportTask.reportType();
+ var fileName = reportType.name() + "-report-" + reportTask.id() + ".csv";
+ var stagedFile = createTempFile(reportTask.id());
try {
var md5 = createDigest("MD5");
var sha256 = createDigest("SHA-256");
- var rowsCount = 0L;
+ long rowsCount;
try (
var fileOutputStream = Files.newOutputStream(stagedFile);
var bufferedOutputStream = new BufferedOutputStream(fileOutputStream);
@@ -111,7 +138,6 @@ public GeneratedCsvReport generate(ReportJob reportJob) {
case payments -> writePaymentsCsv(writer, reportQuery.getPayments(), zoneId);
case withdrawals -> writeWithdrawalsCsv(writer, reportQuery.getWithdrawals(), zoneId);
};
- writer.flush();
}
return new GeneratedCsvReport(
fileName,
@@ -134,7 +160,7 @@ public GeneratedCsvReport generate(ReportJob reportJob) {
private long writePaymentsCsv(BufferedWriter writer, PaymentsQuery query, ZoneId zoneId) throws IOException {
writer.write(String.join(",", PAYMENT_COLUMNS));
- writer.newLine();
+ writer.write(CSV_LINE_ENDING);
try (var rows = reportCsvDao.fetchPayments(query)) {
return writeRows(writer, rows, PAYMENT_COLUMNS, zoneId);
}
@@ -146,27 +172,33 @@ private long writeWithdrawalsCsv(
ZoneId zoneId
) throws IOException {
writer.write(String.join(",", WITHDRAWAL_COLUMNS));
- writer.newLine();
+ writer.write(CSV_LINE_ENDING);
try (var rows = reportCsvDao.fetchWithdrawals(query)) {
return writeRows(writer, rows, WITHDRAWAL_COLUMNS, zoneId);
}
}
- @SneakyThrows
private long writeRows(
BufferedWriter writer,
Cursor extends Record> rows,
List columns,
ZoneId zoneId
- ) {
+ ) throws IOException {
var rowCount = 0L;
for (var row : rows) {
+ throwIfInterrupted();
writeRow(writer, row, columns, zoneId);
rowCount++;
}
return rowCount;
}
+ private void throwIfInterrupted() {
+ if (Thread.currentThread().isInterrupted()) {
+ throw new CancellationException("Report CSV generation was interrupted");
+ }
+ }
+
private void writeRow(
BufferedWriter writer,
Record row,
@@ -180,25 +212,25 @@ private void writeRow(
var column = columns.get(i);
writer.write(escapeCsv(renderValue(row, column, zoneId)));
}
- writer.newLine();
+ writer.write(CSV_LINE_ENDING);
}
private String renderValue(Record row, String column, ZoneId zoneId) {
return switch (column) {
- case "created_date" -> renderTimestampDate(row.get("created_at", Timestamp.class), zoneId);
- case "created_time" -> renderTimestampTime(row.get("created_at", Timestamp.class), zoneId);
- case "finalized_date" -> renderTimestampDate(row.get("finalized_at", Timestamp.class), zoneId);
- case "finalized_time" -> renderTimestampTime(row.get("finalized_at", Timestamp.class), zoneId);
- case "amount" -> renderMinorUnits(row.get("amount"), row.get("currency", String.class));
- case "provider_amount" -> renderMinorUnits(
+ case CREATED_DATE_COLUMN -> renderTimestampDate(row.get("created_at", LocalDateTime.class), zoneId);
+ case CREATED_TIME_COLUMN -> renderTimestampTime(row.get("created_at", LocalDateTime.class), zoneId);
+ case FINALIZED_DATE_COLUMN -> renderTimestampDate(row.get("finalized_at", LocalDateTime.class), zoneId);
+ case FINALIZED_TIME_COLUMN -> renderTimestampTime(row.get("finalized_at", LocalDateTime.class), zoneId);
+ case AMOUNT_COLUMN -> renderMinorUnits(row.get("amount"), row.get("currency", String.class));
+ case PROVIDER_AMOUNT_COLUMN -> renderMinorUnits(
row.get("provider_amount"),
firstNonBlank(row.get("provider_currency", String.class), row.get("currency", String.class))
);
- case "original_amount" -> renderMinorUnits(
+ case ORIGINAL_AMOUNT_COLUMN -> renderMinorUnits(
row.get("original_amount"),
row.get("original_currency", String.class)
);
- case "converted_amount" -> renderMinorUnits(
+ case CONVERTED_AMOUNT_COLUMN -> renderMinorUnits(
row.get("converted_amount"),
row.get("currency", String.class)
);
@@ -216,19 +248,19 @@ private String renderScalarValue(Object value) {
return value.toString();
}
- private String renderTimestampDate(Timestamp timestamp, ZoneId zoneId) {
+ private String renderTimestampDate(LocalDateTime timestamp, ZoneId zoneId) {
if (timestamp == null) {
return "";
}
- var localDateTime = timestamp.toInstant().atZone(zoneId).toLocalDateTime();
+ var localDateTime = timestamp.atZone(ZoneOffset.UTC).withZoneSameInstant(zoneId).toLocalDateTime();
return CSV_DATE_FORMATTER.format(localDateTime.toLocalDate());
}
- private String renderTimestampTime(Timestamp timestamp, ZoneId zoneId) {
+ private String renderTimestampTime(LocalDateTime timestamp, ZoneId zoneId) {
if (timestamp == null) {
return "";
}
- var localDateTime = timestamp.toInstant().atZone(zoneId).toLocalDateTime();
+ var localDateTime = timestamp.atZone(ZoneOffset.UTC).withZoneSameInstant(zoneId).toLocalDateTime();
return CSV_TIME_FORMATTER.format(localDateTime.toLocalTime());
}
@@ -267,7 +299,10 @@ private String firstNonBlank(String first, String second) {
}
private String escapeCsv(String value) {
- if (!value.contains(",") && !value.contains("\"") && !value.contains("\n")) {
+ if (!value.contains(",")
+ && !value.contains("\"")
+ && !value.contains("\n")
+ && !value.contains("\r")) {
return value;
}
return "\"" + value.replace("\"", "\"\"") + "\"";
@@ -284,7 +319,7 @@ private Path createTempFile(long reportId) {
private MessageDigest createDigest(String algorithm) {
try {
return MessageDigest.getInstance(algorithm);
- } catch (Exception ex) {
+ } catch (NoSuchAlgorithmException ex) {
throw new IllegalStateException("Failed to initialize " + algorithm + " digest", ex);
}
}
diff --git a/src/main/java/dev/vality/ccreporter/report/ReportLifecycleScheduler.java b/src/main/java/dev/vality/ccreporter/report/ReportLifecycleScheduler.java
index 87c542f..1880485 100644
--- a/src/main/java/dev/vality/ccreporter/report/ReportLifecycleScheduler.java
+++ b/src/main/java/dev/vality/ccreporter/report/ReportLifecycleScheduler.java
@@ -7,15 +7,13 @@
@Component
@RequiredArgsConstructor
-@ConditionalOnProperty(prefix = "ccr.scheduler", name = "enabled", havingValue = "true")
+@ConditionalOnProperty(prefix = "scheduler", name = "enabled", havingValue = "true")
public class ReportLifecycleScheduler {
private final ReportLifecycleService reportLifecycleService;
- @Scheduled(fixedDelayString = "${ccr.scheduler.poll-interval-ms:10000}")
+ @Scheduled(fixedDelayString = "${scheduler.poll-interval-ms:10000}")
public void runLifecycleTick() {
- reportLifecycleService.timeoutStaleProcessingReports();
- reportLifecycleService.expireReadyReports();
- reportLifecycleService.processNextPendingReport();
+ reportLifecycleService.runLifecycleTick();
}
}
diff --git a/src/main/java/dev/vality/ccreporter/report/ReportLifecycleService.java b/src/main/java/dev/vality/ccreporter/report/ReportLifecycleService.java
index a8c19c0..f565008 100644
--- a/src/main/java/dev/vality/ccreporter/report/ReportLifecycleService.java
+++ b/src/main/java/dev/vality/ccreporter/report/ReportLifecycleService.java
@@ -1,20 +1,23 @@
package dev.vality.ccreporter.report;
import dev.vality.ccreporter.config.properties.ReportProperties;
-import dev.vality.ccreporter.config.properties.ReportSchedulerProperties;
import dev.vality.ccreporter.dao.ReportLifecycleDao;
import dev.vality.ccreporter.domain.tables.pojos.ReportFile;
-import dev.vality.ccreporter.domain.tables.pojos.ReportJob;
import dev.vality.ccreporter.model.GeneratedCsvReport;
+import dev.vality.ccreporter.model.ReportTask;
import dev.vality.ccreporter.storage.FileStorageService;
import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.io.IOException;
import java.nio.file.Files;
import java.time.Duration;
import java.time.Instant;
+import java.util.ArrayList;
+import java.util.concurrent.*;
+@Slf4j
@Service
@RequiredArgsConstructor
public class ReportLifecycleService {
@@ -24,101 +27,247 @@ public class ReportLifecycleService {
private final ReportLifecycleDao reportLifecycleDao;
private final ReportCsvService reportCsvService;
private final FileStorageService fileStorageService;
- private final ReportLifecycleTransactionService reportLifecycleTransactionService;
private final ReportProperties reportProperties;
- private final ReportSchedulerProperties reportSchedulerProperties;
+ private final ExecutorService reportWorkerExecutor;
- public void timeoutStaleProcessingReports() {
- timeoutStaleProcessingReports(Instant.now());
+ public void runLifecycleTick() {
+ var now = Instant.now();
+ timeoutStaleProcessingReports(now);
+ expireReadyReports(now);
+ while (!Thread.currentThread().isInterrupted()
+ && processPendingBatch(Instant.now()) == reportProperties.getWorkerConcurrency()) {
+ // Drain ready reports in bounded parallel batches.
+ }
}
public int timeoutStaleProcessingReports(Instant now) {
- var staleBefore = now.minusMillis(reportSchedulerProperties.getStaleProcessingTimeoutMs());
- return reportLifecycleDao.timeoutStaleProcessingReports(staleBefore, now);
+ var staleBefore = now.minusMillis(reportProperties.getProcessingTimeoutMs());
+ var timedOutReports = reportLifecycleDao.timeoutStaleProcessingReports(staleBefore, now);
+ if (timedOutReports > 0) {
+ log.warn("Timed out {} stale processing report(s)", timedOutReports);
+ }
+ return timedOutReports;
}
- public void expireReadyReports() {
- expireReadyReports(Instant.now());
+ public int expireReadyReports(Instant now) {
+ var expiredReports = reportLifecycleDao.expireReports(now);
+ if (expiredReports > 0) {
+ log.info("Expired {} report(s)", expiredReports);
+ }
+ return expiredReports;
}
- public int expireReadyReports(Instant now) {
- return reportLifecycleDao.expireReports(now);
+ public boolean processNextPendingReport(Instant now) {
+ var reportTask = reportLifecycleDao.claimNextPendingReport(now);
+ if (reportTask.isEmpty()) {
+ return false;
+ }
+ var runningReport = startReport(reportTask.get());
+ return runningReport != null && awaitReport(runningReport);
+ }
+
+ private int processPendingBatch(Instant now) {
+ var runningReports = new ArrayList(reportProperties.getWorkerConcurrency());
+ for (int worker = 0; worker < reportProperties.getWorkerConcurrency(); worker++) {
+ var reportTask = reportLifecycleDao.claimNextPendingReport(now);
+ if (reportTask.isEmpty()) {
+ break;
+ }
+ var runningReport = startReport(reportTask.get());
+ if (runningReport == null) {
+ break;
+ }
+ runningReports.add(runningReport);
+ }
+ for (int reportIndex = 0; reportIndex < runningReports.size(); reportIndex++) {
+ if (!awaitReport(runningReports.get(reportIndex))) {
+ cancelRemainingReports(runningReports, reportIndex + 1);
+ break;
+ }
+ }
+ return runningReports.size();
}
- public void processNextPendingReport() {
- processNextPendingReport(Instant.now());
+ private RunningReport startReport(ReportTask reportTask) {
+ try {
+ var processing = reportWorkerExecutor.submit(() -> processReportTask(reportTask));
+ var deadlineNanos = System.nanoTime() +
+ TimeUnit.MILLISECONDS.toNanos(reportProperties.getProcessingTimeoutMs());
+ log.info("Started report {} processing attempt {}", reportTask.id(), reportTask.attempt());
+ return new RunningReport(reportTask, processing, deadlineNanos);
+ } catch (RejectedExecutionException ex) {
+ handleProcessingFailure(reportTask, Instant.now(), ex);
+ return null;
+ }
}
- public boolean processNextPendingReport(Instant now) {
- var reportJob = reportLifecycleDao.claimNextPendingReport(now);
- if (reportJob.isEmpty()) {
+ private boolean awaitReport(RunningReport runningReport) {
+ var reportTask = runningReport.reportTask();
+ var processing = runningReport.processing();
+ try {
+ var remainingNanos = runningReport.deadlineNanos() - System.nanoTime();
+ if (remainingNanos <= 0) {
+ throw new TimeoutException("Report processing deadline elapsed");
+ }
+ processing.get(remainingNanos, TimeUnit.NANOSECONDS);
+ return true;
+ } catch (TimeoutException ex) {
+ timeoutReport(reportTask, processing, "maximum processing time exceeded");
+ return true;
+ } catch (CancellationException ex) {
+ timeoutReport(reportTask, processing, "worker execution was canceled");
return false;
+ } catch (InterruptedException ex) {
+ timeoutReport(reportTask, processing, "scheduler thread was interrupted");
+ Thread.currentThread().interrupt();
+ return false;
+ } catch (ExecutionException ex) {
+ handleProcessingFailure(reportTask, Instant.now(), ex.getCause());
+ return true;
}
- processReportJob(reportJob.get(), now);
- return true;
}
- private void processReportJob(ReportJob reportJob, Instant processingTime) {
- var generatedCsvReport = (GeneratedCsvReport) null;
+ private void cancelRemainingReports(ArrayList runningReports, int firstReportIndex) {
+ for (int reportIndex = firstReportIndex; reportIndex < runningReports.size(); reportIndex++) {
+ var runningReport = runningReports.get(reportIndex);
+ timeoutReport(
+ runningReport.reportTask(),
+ runningReport.processing(),
+ "scheduler stopped while processing batch"
+ );
+ }
+ }
+
+ private void processReportTask(ReportTask reportTask) {
+ GeneratedCsvReport generatedReport = null;
try {
- generatedCsvReport = reportCsvService.generate(reportJob);
- var expiresAt = processingTime.plusSeconds(reportProperties.getExpirationSec());
+ generatedReport = reportCsvService.generate(reportTask);
+ throwIfInterrupted();
+ var expiresAt = Instant.now().plusSeconds(reportProperties.getExpirationSec());
var fileId = fileStorageService.storeFile(
- generatedCsvReport.fileName(),
- generatedCsvReport.contentType(),
- generatedCsvReport.contentPath(),
+ generatedReport.fileName(),
+ generatedReport.contentType(),
+ generatedReport.contentPath(),
expiresAt
);
+ throwIfInterrupted();
+ var reportFile = buildReportFile(fileId, generatedReport);
var finishedAt = Instant.now();
- var reportFile = buildReportFile(fileId, generatedCsvReport);
- reportLifecycleTransactionService.publishCompletedReport(
- reportJob.getId(),
+ var completed = reportLifecycleDao.completeReport(
+ reportTask.id(),
reportFile,
+ generatedReport.dataSnapshotFixedAt(),
finishedAt,
expiresAt,
- generatedCsvReport
+ generatedReport.rowsCount()
);
- } catch (Exception ex) {
- handleProcessingFailure(reportJob, processingTime, ex);
+ if (!completed) {
+ log.info(
+ "Report {} changed state while it was being generated; uploaded file will expire",
+ reportTask.id()
+ );
+ } else {
+ log.info(
+ "Completed report {} with {} row(s)",
+ reportTask.id(),
+ generatedReport.rowsCount()
+ );
+ }
} finally {
- deleteStagedFile(generatedCsvReport);
+ deleteStagedFile(generatedReport);
}
}
- private void handleProcessingFailure(ReportJob reportJob, Instant now, Exception ex) {
+ private void timeoutReport(ReportTask reportTask, Future> processing, String reason) {
+ processing.cancel(true);
+ var finishedAt = Instant.now();
+ try {
+ var timedOut = reportLifecycleDao.markTimedOut(reportTask.id(), finishedAt);
+ if (timedOut) {
+ log.warn("Report {} timed out: {}", reportTask.id(), reason);
+ } else {
+ log.info("Report {} changed state before timeout transition", reportTask.id());
+ }
+ } catch (RuntimeException ex) {
+ log.error(
+ "Failed to mark report {} as timed out after worker cancellation",
+ reportTask.id(),
+ ex
+ );
+ }
+ }
+
+ private void handleProcessingFailure(ReportTask reportTask, Instant now, Throwable ex) {
var errorCode = "report_processing_error";
var errorMessage = ex.getMessage() == null ? ex.getClass().getSimpleName() : ex.getMessage();
- if (reportJob.getAttempt() >= reportProperties.getMaxAttempts()) {
- reportLifecycleDao.markFailed(reportJob.getId(), now, now, errorCode, errorMessage);
- return;
+ if (reportTask.attempt() >= reportProperties.getMaxAttempts()) {
+ var failed = reportLifecycleDao.markFailed(reportTask.id(), now, errorCode, errorMessage);
+ if (failed) {
+ log.error(
+ "Report {} failed after {} attempt(s): {}",
+ reportTask.id(),
+ reportTask.attempt(),
+ errorMessage,
+ ex
+ );
+ } else {
+ log.info("Report {} changed state before failed transition", reportTask.id());
+ }
+ } else {
+ var nextAttemptAt = now.plus(RETRY_BACKOFF);
+ var rescheduled = reportLifecycleDao.rescheduleForRetry(
+ reportTask.id(),
+ nextAttemptAt,
+ errorCode,
+ errorMessage
+ );
+ if (rescheduled) {
+ log.warn(
+ "Report {} attempt {} failed; next attempt at {}: {}",
+ reportTask.id(),
+ reportTask.attempt(),
+ nextAttemptAt,
+ errorMessage,
+ ex
+ );
+ } else {
+ log.info("Report {} changed state before retry transition", reportTask.id());
+ }
}
- reportLifecycleDao.rescheduleForRetry(reportJob.getId(), now.plus(RETRY_BACKOFF), errorCode, errorMessage);
}
- private ReportFile buildReportFile(
- String fileId,
- GeneratedCsvReport generatedCsvReport
- ) {
+ private void throwIfInterrupted() {
+ if (Thread.currentThread().isInterrupted()) {
+ throw new CancellationException("Report processing was interrupted");
+ }
+ }
+
+ private ReportFile buildReportFile(String fileId, GeneratedCsvReport generatedReport) {
return new ReportFile()
.setFileId(fileId)
.setFileType(dev.vality.ccreporter.domain.enums.FileType.csv)
- .setBucket("file-storage")
- .setObjectKey(fileId)
- .setFilename(generatedCsvReport.fileName())
- .setContentType(generatedCsvReport.contentType())
- .setSizeBytes(generatedCsvReport.sizeBytes())
- .setMd5(generatedCsvReport.md5())
- .setSha256(generatedCsvReport.sha256());
+ .setFilename(generatedReport.fileName())
+ .setContentType(generatedReport.contentType())
+ .setSizeBytes(generatedReport.sizeBytes())
+ .setMd5(generatedReport.md5())
+ .setSha256(generatedReport.sha256());
}
- private void deleteStagedFile(GeneratedCsvReport generatedCsvReport) {
- if (generatedCsvReport == null) {
+ private void deleteStagedFile(GeneratedCsvReport generatedReport) {
+ if (generatedReport == null) {
return;
}
try {
- Files.deleteIfExists(generatedCsvReport.contentPath());
- } catch (IOException ignored) {
- // Best-effort cleanup for staged files after upload/publication.
+ Files.deleteIfExists(generatedReport.contentPath());
+ } catch (IOException ex) {
+ log.warn("Failed to delete staged report file {}", generatedReport.contentPath(), ex);
}
}
+
+ private record RunningReport(
+ ReportTask reportTask,
+ Future> processing,
+ long deadlineNanos
+ ) {
+ }
}
diff --git a/src/main/java/dev/vality/ccreporter/report/ReportLifecycleTransactionService.java b/src/main/java/dev/vality/ccreporter/report/ReportLifecycleTransactionService.java
deleted file mode 100644
index 3c91140..0000000
--- a/src/main/java/dev/vality/ccreporter/report/ReportLifecycleTransactionService.java
+++ /dev/null
@@ -1,38 +0,0 @@
-package dev.vality.ccreporter.report;
-
-import dev.vality.ccreporter.dao.ReportLifecycleDao;
-import dev.vality.ccreporter.domain.tables.pojos.ReportFile;
-import dev.vality.ccreporter.model.GeneratedCsvReport;
-import lombok.RequiredArgsConstructor;
-import org.springframework.stereotype.Service;
-import org.springframework.transaction.annotation.Transactional;
-
-import java.time.Instant;
-
-@Service
-@RequiredArgsConstructor
-public class ReportLifecycleTransactionService {
-
- private final ReportLifecycleDao reportLifecycleDao;
-
- @Transactional
- public void publishCompletedReport(
- long reportId,
- ReportFile reportFile,
- Instant finishedAt,
- Instant expiresAt,
- GeneratedCsvReport generatedCsvReport
- ) {
- var published = reportLifecycleDao.publishFileRecord(reportId, reportFile, finishedAt);
- var markedCreated = reportLifecycleDao.markCreated(
- reportId,
- generatedCsvReport.dataSnapshotFixedAt(),
- finishedAt,
- expiresAt,
- generatedCsvReport.rowsCount()
- );
- if (!published || !markedCreated) {
- throw new IllegalStateException("Failed to publish created report " + reportId);
- }
- }
-}
diff --git a/src/main/java/dev/vality/ccreporter/report/ReportManagementService.java b/src/main/java/dev/vality/ccreporter/report/ReportManagementService.java
index 2e1f4f7..bd01300 100644
--- a/src/main/java/dev/vality/ccreporter/report/ReportManagementService.java
+++ b/src/main/java/dev/vality/ccreporter/report/ReportManagementService.java
@@ -12,8 +12,6 @@
import dev.vality.ccreporter.storage.FileStorageService;
import dev.vality.ccreporter.util.TimestampUtils;
import lombok.RequiredArgsConstructor;
-import lombok.SneakyThrows;
-import org.springframework.dao.DuplicateKeyException;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.util.StringUtils;
@@ -28,7 +26,6 @@ public class ReportManagementService {
private final ReportCommandDao reportCommandDao;
private final ReportQueryDao reportQueryDao;
private final ReportLifecycleDao reportLifecycleDao;
- private final ReportManagementTransactionService reportManagementTransactionService;
private final ReportAuditService reportAuditService;
private final ReportRequestValidator reportRequestValidator;
private final ReportThriftMapper reportThriftMapper;
@@ -38,48 +35,58 @@ public class ReportManagementService {
private final ReportProperties reportProperties;
private final FileStorageService fileStorageService;
- @SneakyThrows
- public long createReport(CreateReportRequest request) {
+ @Transactional
+ public long createReport(CreateReportRequest request) throws InvalidRequest {
reportRequestValidator.validateCreate(request);
var auditMetadata = requestAuditMetadataResolver.resolve();
var timezone = StringUtils.hasText(request.getTimezone()) ? request.getTimezone() : "UTC";
var createdBy = auditMetadata.email();
- try {
- return reportManagementTransactionService.createReport(createdBy, auditMetadata, request, timezone);
- } catch (DuplicateKeyException ex) {
- return reportCommandDao.findByIdempotencyKey(createdBy, request.getIdempotencyKey())
- .orElseThrow(() -> ex);
+ var result = reportCommandDao.createReport(
+ createdBy,
+ request.getReportType(),
+ request.getFileType(),
+ request.getQuery(),
+ timezone,
+ request.getIdempotencyKey()
+ );
+ if (result.created()) {
+ reportAuditService.writeReportCreated(result.reportId(), createdBy, auditMetadata, request, timezone);
}
+ return result.reportId();
}
- @SneakyThrows
- public Report getReport(GetReportRequest request) {
+ @Transactional
+ public Report getReport(GetReportRequest request) throws InvalidRequest, ReportNotFound {
if (request == null) {
throw invalidRequest("request is required");
}
var createdBy = requestAuditMetadataResolver.resolve().email();
+ reportLifecycleDao.expireReports(Instant.now());
return reportQueryDao.getReport(createdBy, request.getReportId())
.map(reportThriftMapper::mapReport)
.orElseThrow(ReportNotFound::new);
}
- @SneakyThrows
- public GetReportsResponse getReports(GetReportsRequest request) {
+ @Transactional
+ public GetReportsResponse getReports(GetReportsRequest request) throws InvalidRequest, BadContinuationToken {
var createdBy = requestAuditMetadataResolver.resolve().email();
var safeRequest = request == null ? new GetReportsRequest() : request;
reportRequestValidator.validateGetReports(safeRequest);
+ reportLifecycleDao.expireReports(Instant.now());
var meta = safeRequest.getMeta();
var limit = resolveLimit(meta);
var cursor = meta != null && meta.isSetContinuationToken()
? continuationTokenJsonSerializer.deserialize(meta.getContinuationToken())
: null;
- var storedReports = reportQueryDao.getReports(createdBy, safeRequest.getFilter(), cursor, limit);
+ var storedReports = reportQueryDao.getReports(createdBy, safeRequest.getFilter(), cursor, limit + 1);
+ var hasNextPage = storedReports.size() > limit;
+ var page = hasNextPage ? storedReports.subList(0, limit) : storedReports;
var response = new GetReportsResponse();
- response.setReports(storedReports.stream().map(reportThriftMapper::mapReport).toList());
- if (storedReports.size() == limit) {
- var lastReport = storedReports.getLast();
+ response.setReports(page.stream().map(reportThriftMapper::mapReport).toList());
+ if (hasNextPage) {
+ var lastReport = page.getLast();
response.setContinuationToken(
continuationTokenJsonSerializer.serialize(
TimestampUtils.toInstant(lastReport.job().getCreatedAt()),
@@ -91,8 +98,7 @@ public GetReportsResponse getReports(GetReportsRequest request) {
}
@Transactional
- @SneakyThrows
- public void cancelReport(CancelReportRequest request) {
+ public void cancelReport(CancelReportRequest request) throws InvalidRequest, ReportNotFound {
if (request == null) {
throw invalidRequest("request is required");
}
@@ -105,37 +111,44 @@ public void cancelReport(CancelReportRequest request) {
reportAuditService.writeReportCanceled(request.getReportId(), createdBy, auditMetadata, updated);
}
- @SneakyThrows
- public String generatePresignedUrl(GeneratePresignedUrlRequest request) {
+ public String generatePresignedUrl(GeneratePresignedUrlRequest request) throws InvalidRequest, FileNotFound {
if (request == null) {
throw invalidRequest("request is required");
}
var auditMetadata = requestAuditMetadataResolver.resolve();
var createdBy = auditMetadata.email();
- var fileData = reportQueryDao.getFile(createdBy, request.getFileId());
+ var now = Instant.now();
+ var fileData = reportQueryDao.getDownloadableFile(createdBy, request.getFileId(), now);
if (fileData.isEmpty()) {
throw new FileNotFound();
}
- var effectiveExpiresAt = resolveEffectivePresignedUrlExpiresAt(request);
+ var downloadableFile = fileData.get();
+ var effectiveExpiresAt = resolveEffectivePresignedUrlExpiresAt(
+ request,
+ downloadableFile.reportExpiresAt(),
+ now
+ );
var url = fileStorageService.generateDownloadUrl(
- fileData.get().getFileId(),
+ downloadableFile.file().getFileId(),
effectiveExpiresAt
);
reportAuditService.writePresignedUrlGenerated(
- fileData.get().getReportId(),
+ downloadableFile.file().getReportId(),
createdBy,
auditMetadata,
request,
effectiveExpiresAt,
- fileData.get().getFileId()
+ downloadableFile.file().getFileId()
);
return url;
}
- @SneakyThrows
- private Instant resolveEffectivePresignedUrlExpiresAt(GeneratePresignedUrlRequest request) {
- var now = Instant.now();
+ private Instant resolveEffectivePresignedUrlExpiresAt(
+ GeneratePresignedUrlRequest request,
+ Instant reportExpiresAt,
+ Instant now
+ ) throws InvalidRequest {
var ttlCap = now.plusSeconds(reportProperties.getPresignedUrlTtlSec());
var requestedExpiresAt = request.isSetRequestedExpiresAt()
? TimestampUtils.parse(request.getRequestedExpiresAt())
@@ -143,7 +156,8 @@ private Instant resolveEffectivePresignedUrlExpiresAt(GeneratePresignedUrlReques
if (!requestedExpiresAt.isAfter(now)) {
throw invalidRequest("requested_expires_at must be in the future");
}
- return requestedExpiresAt.isAfter(ttlCap) ? ttlCap : requestedExpiresAt;
+ var requestAndConfigCap = requestedExpiresAt.isAfter(ttlCap) ? ttlCap : requestedExpiresAt;
+ return requestAndConfigCap.isAfter(reportExpiresAt) ? reportExpiresAt : requestAndConfigCap;
}
private int resolveLimit(GetReportsMeta meta) {
diff --git a/src/main/java/dev/vality/ccreporter/report/ReportManagementTransactionService.java b/src/main/java/dev/vality/ccreporter/report/ReportManagementTransactionService.java
deleted file mode 100644
index 5db675e..0000000
--- a/src/main/java/dev/vality/ccreporter/report/ReportManagementTransactionService.java
+++ /dev/null
@@ -1,35 +0,0 @@
-package dev.vality.ccreporter.report;
-
-import dev.vality.ccreporter.CreateReportRequest;
-import dev.vality.ccreporter.dao.ReportCommandDao;
-import dev.vality.ccreporter.model.RequestAuditMetadata;
-import lombok.RequiredArgsConstructor;
-import org.springframework.stereotype.Service;
-import org.springframework.transaction.annotation.Transactional;
-
-@Service
-@RequiredArgsConstructor
-public class ReportManagementTransactionService {
-
- private final ReportCommandDao reportCommandDao;
- private final ReportAuditService reportAuditService;
-
- @Transactional
- public long createReport(
- String createdBy,
- RequestAuditMetadata auditMetadata,
- CreateReportRequest request,
- String timezone
- ) {
- var reportId = reportCommandDao.createReport(
- createdBy,
- request.getReportType(),
- request.getFileType(),
- request.getQuery(),
- timezone,
- request.getIdempotencyKey()
- );
- reportAuditService.writeReportCreated(reportId, createdBy, auditMetadata, request, timezone);
- return reportId;
- }
-}
diff --git a/src/main/java/dev/vality/ccreporter/report/ReportQueryService.java b/src/main/java/dev/vality/ccreporter/report/ReportQueryService.java
index 44eccee..b58ace5 100644
--- a/src/main/java/dev/vality/ccreporter/report/ReportQueryService.java
+++ b/src/main/java/dev/vality/ccreporter/report/ReportQueryService.java
@@ -2,16 +2,15 @@
import dev.vality.ccreporter.ReportQuery;
import dev.vality.ccreporter.ReportType;
+import dev.vality.ccreporter.TimeRange;
import dev.vality.ccreporter.util.TimestampUtils;
-import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Service;
+import org.springframework.util.StringUtils;
-import java.nio.charset.StandardCharsets;
-import java.security.MessageDigest;
+import java.time.DateTimeException;
import java.time.Instant;
@Service
-@RequiredArgsConstructor
public class ReportQueryService {
public QuerySpec resolveQuerySpec(ReportQuery query) {
@@ -19,38 +18,27 @@ public QuerySpec resolveQuerySpec(ReportQuery query) {
throw new IllegalArgumentException("query is required");
}
if (query.isSetPayments()) {
- var paymentsQuery = query.getPayments();
- var timeRange = paymentsQuery.getTimeRange();
- return new QuerySpec(
- ReportType.payments,
- new QueryTimeRange(
- TimestampUtils.parse(timeRange.getFromTime()),
- TimestampUtils.parse(timeRange.getToTime())
- )
- );
+ return new QuerySpec(ReportType.payments, parseTimeRange(query.getPayments().getTimeRange()));
+ }
+ if (query.isSetWithdrawals()) {
+ return new QuerySpec(ReportType.withdrawals, parseTimeRange(query.getWithdrawals().getTimeRange()));
}
- var withdrawalsQuery = query.getWithdrawals();
- var timeRange = withdrawalsQuery.getTimeRange();
- return new QuerySpec(
- ReportType.withdrawals,
- new QueryTimeRange(
- TimestampUtils.parse(timeRange.getFromTime()),
- TimestampUtils.parse(timeRange.getToTime())
- )
- );
+ throw new IllegalArgumentException("query must select one branch");
}
- public String hash(String value) {
+ private QueryTimeRange parseTimeRange(TimeRange timeRange) {
+ if (timeRange == null
+ || !StringUtils.hasText(timeRange.getFromTime())
+ || !StringUtils.hasText(timeRange.getToTime())) {
+ throw new IllegalArgumentException("time range is required");
+ }
try {
- var digest = MessageDigest.getInstance("SHA-256");
- var hashBytes = digest.digest(value.getBytes(StandardCharsets.UTF_8));
- var stringBuilder = new StringBuilder(hashBytes.length * 2);
- for (byte hashByte : hashBytes) {
- stringBuilder.append(String.format("%02x", hashByte));
- }
- return stringBuilder.toString();
- } catch (Exception ex) {
- throw new IllegalStateException("Failed to hash report query", ex);
+ return new QueryTimeRange(
+ TimestampUtils.parse(timeRange.getFromTime()),
+ TimestampUtils.parse(timeRange.getToTime())
+ );
+ } catch (DateTimeException ex) {
+ throw new IllegalArgumentException("time range must contain ISO-8601 timestamps", ex);
}
}
diff --git a/src/main/java/dev/vality/ccreporter/report/ReportRequestValidator.java b/src/main/java/dev/vality/ccreporter/report/ReportRequestValidator.java
index 0a1f0e5..dc8ad39 100644
--- a/src/main/java/dev/vality/ccreporter/report/ReportRequestValidator.java
+++ b/src/main/java/dev/vality/ccreporter/report/ReportRequestValidator.java
@@ -8,10 +8,11 @@
import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
+import java.time.DateTimeException;
+import java.time.Instant;
import java.time.ZoneId;
import java.util.ArrayList;
import java.util.List;
-import java.util.Objects;
@Component
@RequiredArgsConstructor
@@ -24,14 +25,14 @@ public void validateCreate(CreateReportRequest request) throws InvalidRequest {
if (request == null) {
errors.add("request is required");
} else {
- validateQuery(request, errors);
- if (StringUtils.hasText(request.getTimezone())) {
- try {
- Objects.requireNonNull(ZoneId.of(request.getTimezone()));
- } catch (Exception ex) {
- errors.add("timezone must be a valid IANA timezone");
- }
+ if (!request.isSetReportType()) {
+ errors.add("report_type is required");
+ }
+ if (!request.isSetFileType()) {
+ errors.add("file_type is required");
}
+ validateQuery(request, errors);
+ validateTimezone(request.getTimezone(), errors);
}
if (!errors.isEmpty()) {
throw new InvalidRequest(errors);
@@ -45,10 +46,20 @@ public void validateGetReports(GetReportsRequest request) throws InvalidRequest
errors.add("meta.limit must be positive");
}
var filter = request.getFilter();
- if (filter != null && filter.isSetCreatedFrom() && filter.isSetCreatedTo()) {
- var createdFrom = TimestampUtils.parse(filter.getCreatedFrom());
- var createdTo = TimestampUtils.parse(filter.getCreatedTo());
- if (createdFrom.isAfter(createdTo)) {
+ if (filter != null) {
+ var createdFrom = parseFilterTimestamp(
+ filter.isSetCreatedFrom(),
+ filter.getCreatedFrom(),
+ "filter.created_from",
+ errors
+ );
+ var createdTo = parseFilterTimestamp(
+ filter.isSetCreatedTo(),
+ filter.getCreatedTo(),
+ "filter.created_to",
+ errors
+ );
+ if (createdFrom != null && createdTo != null && createdFrom.isAfter(createdTo)) {
errors.add("filter.created_from must be before or equal to filter.created_to");
}
}
@@ -57,12 +68,29 @@ public void validateGetReports(GetReportsRequest request) throws InvalidRequest
}
}
+ private Instant parseFilterTimestamp(
+ boolean isSet,
+ String value,
+ String fieldName,
+ List errors
+ ) {
+ if (!isSet) {
+ return null;
+ }
+ try {
+ return TimestampUtils.parse(value);
+ } catch (DateTimeException ex) {
+ errors.add(fieldName + " must use ISO-8601 format");
+ return null;
+ }
+ }
+
private void validateQuery(CreateReportRequest request, List errors) {
ReportQueryService.QuerySpec querySpec;
try {
querySpec = reportQueryService.resolveQuerySpec(request.getQuery());
} catch (IllegalArgumentException ex) {
- errors.add("query must select exactly one branch");
+ errors.add(ex.getMessage());
return;
}
if (request.isSetReportType() && request.getReportType() != querySpec.reportType()) {
@@ -72,4 +100,15 @@ private void validateQuery(CreateReportRequest request, List errors) {
errors.add("time_range.from_time must be before time_range.to_time");
}
}
+
+ private void validateTimezone(String timezone, List errors) {
+ if (!StringUtils.hasText(timezone)) {
+ return;
+ }
+ try {
+ ZoneId.of(timezone);
+ } catch (DateTimeException ex) {
+ errors.add("timezone must be a valid IANA timezone");
+ }
+ }
}
diff --git a/src/main/java/dev/vality/ccreporter/handler/ReportingHandler.java b/src/main/java/dev/vality/ccreporter/resource/ReportingHandler.java
similarity index 92%
rename from src/main/java/dev/vality/ccreporter/handler/ReportingHandler.java
rename to src/main/java/dev/vality/ccreporter/resource/ReportingHandler.java
index 09f6803..838b5d4 100644
--- a/src/main/java/dev/vality/ccreporter/handler/ReportingHandler.java
+++ b/src/main/java/dev/vality/ccreporter/resource/ReportingHandler.java
@@ -1,8 +1,8 @@
-package dev.vality.ccreporter.handler;
+package dev.vality.ccreporter.resource;
import dev.vality.ccreporter.*;
-import dev.vality.ccreporter.handler.support.ReportingHandlerLogSupport;
-import dev.vality.ccreporter.handler.support.ThriftLoggingHandler;
+import dev.vality.ccreporter.resource.util.ReportingHandlerLogSupport;
+import dev.vality.ccreporter.resource.util.ThriftLoggingHandler;
import dev.vality.ccreporter.report.ReportManagementService;
import lombok.RequiredArgsConstructor;
import lombok.SneakyThrows;
diff --git a/src/main/java/dev/vality/ccreporter/handler/support/ReportingHandlerLogSupport.java b/src/main/java/dev/vality/ccreporter/resource/util/ReportingHandlerLogSupport.java
similarity index 98%
rename from src/main/java/dev/vality/ccreporter/handler/support/ReportingHandlerLogSupport.java
rename to src/main/java/dev/vality/ccreporter/resource/util/ReportingHandlerLogSupport.java
index 8241dad..040784b 100644
--- a/src/main/java/dev/vality/ccreporter/handler/support/ReportingHandlerLogSupport.java
+++ b/src/main/java/dev/vality/ccreporter/resource/util/ReportingHandlerLogSupport.java
@@ -1,4 +1,4 @@
-package dev.vality.ccreporter.handler.support;
+package dev.vality.ccreporter.resource.util;
import dev.vality.ccreporter.*;
import lombok.experimental.UtilityClass;
diff --git a/src/main/java/dev/vality/ccreporter/handler/support/ThriftLoggingHandler.java b/src/main/java/dev/vality/ccreporter/resource/util/ThriftLoggingHandler.java
similarity index 97%
rename from src/main/java/dev/vality/ccreporter/handler/support/ThriftLoggingHandler.java
rename to src/main/java/dev/vality/ccreporter/resource/util/ThriftLoggingHandler.java
index e333b65..1f34b89 100644
--- a/src/main/java/dev/vality/ccreporter/handler/support/ThriftLoggingHandler.java
+++ b/src/main/java/dev/vality/ccreporter/resource/util/ThriftLoggingHandler.java
@@ -1,4 +1,4 @@
-package dev.vality.ccreporter.handler.support;
+package dev.vality.ccreporter.resource.util;
import org.slf4j.LoggerFactory;
diff --git a/src/main/java/dev/vality/ccreporter/security/RequestAuditMetadataResolver.java b/src/main/java/dev/vality/ccreporter/security/RequestAuditMetadataResolver.java
index 0b9a9cb..5140d7a 100644
--- a/src/main/java/dev/vality/ccreporter/security/RequestAuditMetadataResolver.java
+++ b/src/main/java/dev/vality/ccreporter/security/RequestAuditMetadataResolver.java
@@ -2,19 +2,17 @@
import dev.vality.ccreporter.model.RequestAuditMetadata;
import dev.vality.woody.api.trace.Metadata;
-import dev.vality.woody.api.trace.TraceData;
import dev.vality.woody.api.trace.context.TraceContext;
import dev.vality.woody.api.trace.context.metadata.user.UserIdentityEmailExtensionKit;
import dev.vality.woody.api.trace.context.metadata.user.UserIdentityIdExtensionKit;
import dev.vality.woody.api.trace.context.metadata.user.UserIdentityRealmExtensionKit;
import dev.vality.woody.api.trace.context.metadata.user.UserIdentityUsernameExtensionKit;
-import io.opentelemetry.api.GlobalOpenTelemetry;
-import io.opentelemetry.context.propagation.TextMapSetter;
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.api.trace.SpanContext;
import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
-import java.util.HashMap;
-import java.util.Map;
+import java.util.stream.Collectors;
@Component
public class RequestAuditMetadataResolver {
@@ -23,36 +21,45 @@ public class RequestAuditMetadataResolver {
private static final String WOODY_USERNAME = UserIdentityUsernameExtensionKit.KEY;
private static final String WOODY_EMAIL = UserIdentityEmailExtensionKit.KEY;
private static final String WOODY_REALM = UserIdentityRealmExtensionKit.KEY;
- private static final String TRACE_PARENT = "traceparent";
- private static final String TRACE_STATE = "tracestate";
public RequestAuditMetadata resolve() {
var traceData = TraceContext.getCurrentTraceData();
var activeSpan = traceData.getActiveSpan();
var metadata = activeSpan.getCustomMetadata();
- var traceHeaders = extractTraceHeaders(traceData);
+ var spanContext = Span.current().getSpanContext();
return new RequestAuditMetadata(
metadataValue(metadata, WOODY_USER_ID),
metadataValue(metadata, WOODY_USERNAME),
metadataValue(metadata, WOODY_EMAIL),
metadataValue(metadata, WOODY_REALM),
- activeSpan.getSpan().getTraceId(),
- traceHeaders.get(TRACE_PARENT),
- traceHeaders.get(TRACE_STATE)
+ resolveTraceId(spanContext, activeSpan.getSpan().getTraceId()),
+ resolveTraceparent(spanContext),
+ resolveTracestate(spanContext)
);
}
- private Map extractTraceHeaders(TraceData traceData) {
- var headers = new HashMap();
- var otelSpan = traceData.getOtelSpan();
- if (otelSpan != null && otelSpan.getSpanContext().isValid()) {
- GlobalOpenTelemetry.getPropagators()
- .getTextMapPropagator()
- .inject(traceData.getOtelContext(), headers, MAP_SETTER);
+ private String resolveTraceId(SpanContext spanContext, String woodyTraceId) {
+ return spanContext.isValid() ? spanContext.getTraceId() : woodyTraceId;
+ }
+
+ private String resolveTraceparent(SpanContext spanContext) {
+ if (!spanContext.isValid()) {
+ return null;
}
- putIfHasText(headers, TRACE_PARENT, traceData.getInboundTraceParent());
- putIfHasText(headers, TRACE_STATE, traceData.getInboundTraceState());
- return headers;
+ return "00-%s-%s-%s".formatted(
+ spanContext.getTraceId(),
+ spanContext.getSpanId(),
+ spanContext.getTraceFlags().asHex()
+ );
+ }
+
+ private String resolveTracestate(SpanContext spanContext) {
+ if (!spanContext.isValid() || spanContext.getTraceState().isEmpty()) {
+ return null;
+ }
+ return spanContext.getTraceState().asMap().entrySet().stream()
+ .map(entry -> entry.getKey() + "=" + entry.getValue())
+ .collect(Collectors.joining(","));
}
private String metadataValue(Metadata metadata, String key) {
@@ -73,16 +80,4 @@ private String normalize(Object value) {
return null;
}
- private void putIfHasText(Map headers, String key, String value) {
- if (StringUtils.hasText(value) && !headers.containsKey(key)) {
- headers.put(key, value.trim());
- }
- }
-
- private static final TextMapSetter