From 327c3404e010bed595a9a81af233f11ed58dbe46 Mon Sep 17 00:00:00 2001 From: Amir Golmoradi Date: Sun, 20 Sep 2026 17:52:00 +0330 Subject: [PATCH 1/6] Add validated Kafka event infrastructure --- .../config/PayGuardKafkaConfiguration.java | 17 +++++ .../PayGuardKafkaConsumerConfiguration.java | 27 ++++++++ .../consumer/IdempotentEventProcessor.java | 29 ++++++++ .../kafka/consumer/ProcessedEventStore.java | 9 +++ .../kafka/event/CollateralRiskPayload.java | 13 ++++ .../config/kafka/event/EventMetadata.java | 3 + .../java/config/kafka/event/EventPayload.java | 3 + .../java/config/kafka/event/EventTypes.java | 15 +++++ .../kafka/event/LoanLifecyclePayload.java | 12 ++++ .../config/kafka/event/PayGuardEvent.java | 67 +++++++++++++++++++ .../kafka/event/WalletTransactionPayload.java | 10 +++ .../PayGuardKafkaEventPublisher.java | 27 ++++++++ .../kafka/topic/PayGuardKafkaTopics.java | 13 ++++ .../exception/InvalidKafkaEventException.java | 13 ++++ 14 files changed, 258 insertions(+) create mode 100644 payguard-shared/src/main/java/config/kafka/config/PayGuardKafkaConfiguration.java create mode 100644 payguard-shared/src/main/java/config/kafka/config/PayGuardKafkaConsumerConfiguration.java create mode 100644 payguard-shared/src/main/java/config/kafka/consumer/IdempotentEventProcessor.java create mode 100644 payguard-shared/src/main/java/config/kafka/consumer/ProcessedEventStore.java create mode 100644 payguard-shared/src/main/java/config/kafka/event/CollateralRiskPayload.java create mode 100644 payguard-shared/src/main/java/config/kafka/event/EventMetadata.java create mode 100644 payguard-shared/src/main/java/config/kafka/event/EventPayload.java create mode 100644 payguard-shared/src/main/java/config/kafka/event/EventTypes.java create mode 100644 payguard-shared/src/main/java/config/kafka/event/LoanLifecyclePayload.java create mode 100644 payguard-shared/src/main/java/config/kafka/event/PayGuardEvent.java create mode 100644 payguard-shared/src/main/java/config/kafka/event/WalletTransactionPayload.java create mode 100644 payguard-shared/src/main/java/config/kafka/publisher/PayGuardKafkaEventPublisher.java create mode 100644 payguard-shared/src/main/java/config/kafka/topic/PayGuardKafkaTopics.java create mode 100644 payguard-shared/src/main/java/exception/InvalidKafkaEventException.java diff --git a/payguard-shared/src/main/java/config/kafka/config/PayGuardKafkaConfiguration.java b/payguard-shared/src/main/java/config/kafka/config/PayGuardKafkaConfiguration.java new file mode 100644 index 0000000..ce3ccef --- /dev/null +++ b/payguard-shared/src/main/java/config/kafka/config/PayGuardKafkaConfiguration.java @@ -0,0 +1,17 @@ +package config.kafka.config; + +import config.kafka.publisher.PayGuardKafkaEventPublisher; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.kafka.core.KafkaTemplate; + +/** Opt-in shared publisher configuration for services that publish Kafka events. */ +@Configuration(proxyBeanMethods = false) +public class PayGuardKafkaConfiguration { + + @Bean + public PayGuardKafkaEventPublisher payGuardKafkaEventPublisher( + KafkaTemplate kafkaTemplate) { + return new PayGuardKafkaEventPublisher(kafkaTemplate); + } +} diff --git a/payguard-shared/src/main/java/config/kafka/config/PayGuardKafkaConsumerConfiguration.java b/payguard-shared/src/main/java/config/kafka/config/PayGuardKafkaConsumerConfiguration.java new file mode 100644 index 0000000..1f45c9c --- /dev/null +++ b/payguard-shared/src/main/java/config/kafka/config/PayGuardKafkaConsumerConfiguration.java @@ -0,0 +1,27 @@ +package config.kafka.config; + +import exception.InvalidKafkaEventException; +import org.apache.kafka.common.TopicPartition; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.listener.DeadLetterPublishingRecoverer; +import org.springframework.kafka.listener.DefaultErrorHandler; +import org.springframework.util.backoff.FixedBackOff; + +/** Opt-in retry and dead-letter configuration for non-transactional Kafka consumers. */ +@Configuration(proxyBeanMethods = false) +public class PayGuardKafkaConsumerConfiguration { + + @Bean + public DefaultErrorHandler payGuardKafkaErrorHandler( + KafkaTemplate kafkaTemplate) { + var recoverer = + new DeadLetterPublishingRecoverer( + kafkaTemplate, + (record, error) -> new TopicPartition(record.topic() + ".DLT", record.partition())); + var handler = new DefaultErrorHandler(recoverer, new FixedBackOff(1_000L, 2L)); + handler.addNotRetryableExceptions(InvalidKafkaEventException.class); + return handler; + } +} diff --git a/payguard-shared/src/main/java/config/kafka/consumer/IdempotentEventProcessor.java b/payguard-shared/src/main/java/config/kafka/consumer/IdempotentEventProcessor.java new file mode 100644 index 0000000..d0ae686 --- /dev/null +++ b/payguard-shared/src/main/java/config/kafka/consumer/IdempotentEventProcessor.java @@ -0,0 +1,29 @@ +package config.kafka.consumer; + +import config.kafka.event.EventPayload; +import config.kafka.event.PayGuardEvent; +import java.util.Objects; + +/** Coordinates event deduplication with the consumer's local transaction. */ +public final class IdempotentEventProcessor { + private final ProcessedEventStore eventStore; + + public IdempotentEventProcessor(ProcessedEventStore eventStore) { + this.eventStore = Objects.requireNonNull(eventStore, "eventStore must not be null"); + } + + /** Returns true when the handler ran, or false when the event was already processed. */ + public boolean process( + PayGuardEvent event, String consumerName, Runnable handler) { + Objects.requireNonNull(event, "event must not be null"); + Objects.requireNonNull(handler, "handler must not be null"); + if (consumerName == null || consumerName.isBlank()) { + throw new IllegalArgumentException("consumerName must not be null or blank"); + } + if (!eventStore.tryMarkProcessed(event.eventId(), consumerName)) { + return false; + } + handler.run(); + return true; + } +} diff --git a/payguard-shared/src/main/java/config/kafka/consumer/ProcessedEventStore.java b/payguard-shared/src/main/java/config/kafka/consumer/ProcessedEventStore.java new file mode 100644 index 0000000..16dee06 --- /dev/null +++ b/payguard-shared/src/main/java/config/kafka/consumer/ProcessedEventStore.java @@ -0,0 +1,9 @@ +package config.kafka.consumer; + +import java.util.UUID; + +/** Service-owned store for atomically recording processed event and consumer pairs. */ +@FunctionalInterface +public interface ProcessedEventStore { + boolean tryMarkProcessed(UUID eventId, String consumerName); +} diff --git a/payguard-shared/src/main/java/config/kafka/event/CollateralRiskPayload.java b/payguard-shared/src/main/java/config/kafka/event/CollateralRiskPayload.java new file mode 100644 index 0000000..6a84ce5 --- /dev/null +++ b/payguard-shared/src/main/java/config/kafka/event/CollateralRiskPayload.java @@ -0,0 +1,13 @@ +package config.kafka.event; + +import java.math.BigDecimal; + +public record CollateralRiskPayload( + String collateralPositionId, + String loanId, + String riskEventType, + BigDecimal ltv, + BigDecimal outstandingPrincipal, + BigDecimal collateralMarketValue, + String currency) + implements EventPayload {} diff --git a/payguard-shared/src/main/java/config/kafka/event/EventMetadata.java b/payguard-shared/src/main/java/config/kafka/event/EventMetadata.java new file mode 100644 index 0000000..d4b16d6 --- /dev/null +++ b/payguard-shared/src/main/java/config/kafka/event/EventMetadata.java @@ -0,0 +1,3 @@ +package config.kafka.event; + +public record EventMetadata(String correlationId, String causationId, String initiatedBy) {} diff --git a/payguard-shared/src/main/java/config/kafka/event/EventPayload.java b/payguard-shared/src/main/java/config/kafka/event/EventPayload.java new file mode 100644 index 0000000..70d3944 --- /dev/null +++ b/payguard-shared/src/main/java/config/kafka/event/EventPayload.java @@ -0,0 +1,3 @@ +package config.kafka.event; + +public interface EventPayload {} diff --git a/payguard-shared/src/main/java/config/kafka/event/EventTypes.java b/payguard-shared/src/main/java/config/kafka/event/EventTypes.java new file mode 100644 index 0000000..4102828 --- /dev/null +++ b/payguard-shared/src/main/java/config/kafka/event/EventTypes.java @@ -0,0 +1,15 @@ +package config.kafka.event; + +public final class EventTypes { + public static final String LOAN_APPROVED = "LOAN_APPROVED"; + public static final String LOAN_DISBURSED = "LOAN_DISBURSED"; + public static final String REPAYMENT_APPLIED = "REPAYMENT_APPLIED"; + public static final String LOAN_PAID_OFF = "LOAN_PAID_OFF"; + public static final String LOAN_DEFAULTED = "LOAN_DEFAULTED"; + public static final String LEDGER_TRANSACTION_POSTED = "LEDGER_TRANSACTION_POSTED"; + public static final String MARGIN_CALL_TRIGGERED = "MARGIN_CALL_TRIGGERED"; + public static final String MARGIN_CALL_RESOLVED = "MARGIN_CALL_RESOLVED"; + public static final String LIQUIDATION_TRIGGERED = "LIQUIDATION_TRIGGERED"; + + private EventTypes() {} +} diff --git a/payguard-shared/src/main/java/config/kafka/event/LoanLifecyclePayload.java b/payguard-shared/src/main/java/config/kafka/event/LoanLifecyclePayload.java new file mode 100644 index 0000000..1055f96 --- /dev/null +++ b/payguard-shared/src/main/java/config/kafka/event/LoanLifecyclePayload.java @@ -0,0 +1,12 @@ +package config.kafka.event; + +import java.math.BigDecimal; + +public record LoanLifecyclePayload( + String loanId, + String borrowerId, + String collateralPositionId, + String status, + BigDecimal outstandingPrincipal, + String currency) + implements EventPayload {} diff --git a/payguard-shared/src/main/java/config/kafka/event/PayGuardEvent.java b/payguard-shared/src/main/java/config/kafka/event/PayGuardEvent.java new file mode 100644 index 0000000..d4fe4ee --- /dev/null +++ b/payguard-shared/src/main/java/config/kafka/event/PayGuardEvent.java @@ -0,0 +1,67 @@ +package config.kafka.event; + +import exception.InvalidKafkaEventException; +import java.time.Instant; +import java.util.Map; +import java.util.Objects; +import java.util.UUID; + +public record PayGuardEvent( + UUID eventId, + String eventType, + String aggregateType, + String aggregateId, + String idempotencyKey, + Instant occurredAt, + String producer, + String correlationId, + String causationId, + int schemaVersion, + T payload, + Map metadata) { + + public PayGuardEvent { + Objects.requireNonNull(eventId, "eventId must not be null"); + requireText(eventType, "eventType must not be blank"); + requireText(aggregateType, "aggregateType must not be blank"); + requireText(aggregateId, "aggregateId must not be blank"); + requireText(idempotencyKey, "idempotencyKey must not be blank"); + Objects.requireNonNull(occurredAt, "occurredAt must not be null"); + requireText(producer, "producer must not be blank"); + Objects.requireNonNull(payload, "payload must not be null"); + if (schemaVersion <= 0) { + throw new IllegalArgumentException("schemaVersion must be positive"); + } + metadata = metadata == null ? Map.of() : Map.copyOf(metadata); + } + + public static PayGuardEvent create( + String eventType, + String aggregateType, + String aggregateId, + String idempotencyKey, + String producer, + String correlationId, + String causationId, + T payload) { + return new PayGuardEvent<>( + UUID.randomUUID(), + eventType, + aggregateType, + aggregateId, + idempotencyKey, + Instant.now(), + producer, + correlationId, + causationId, + 1, + payload, + Map.of()); + } + + private static void requireText(String value, String message) { + if (value == null || value.isBlank()) { + throw new InvalidKafkaEventException(message); + } + } +} diff --git a/payguard-shared/src/main/java/config/kafka/event/WalletTransactionPayload.java b/payguard-shared/src/main/java/config/kafka/event/WalletTransactionPayload.java new file mode 100644 index 0000000..41a1cd3 --- /dev/null +++ b/payguard-shared/src/main/java/config/kafka/event/WalletTransactionPayload.java @@ -0,0 +1,10 @@ +package config.kafka.event; + +import java.math.BigDecimal; +import java.util.List; + +public record WalletTransactionPayload( + String transactionId, String transactionType, String idempotencyKey, List entries) + implements EventPayload { + public record Entry(String accountId, BigDecimal amount, String currency) {} +} diff --git a/payguard-shared/src/main/java/config/kafka/publisher/PayGuardKafkaEventPublisher.java b/payguard-shared/src/main/java/config/kafka/publisher/PayGuardKafkaEventPublisher.java new file mode 100644 index 0000000..8d68b65 --- /dev/null +++ b/payguard-shared/src/main/java/config/kafka/publisher/PayGuardKafkaEventPublisher.java @@ -0,0 +1,27 @@ +package config.kafka.publisher; + +import config.kafka.event.EventPayload; +import config.kafka.event.PayGuardEvent; +import exception.InvalidKafkaEventException; +import java.util.Objects; +import java.util.concurrent.CompletableFuture; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.support.SendResult; + +/** Publishes validated events using their aggregate ID as the Kafka key. */ +public final class PayGuardKafkaEventPublisher { + private final KafkaTemplate kafkaTemplate; + + public PayGuardKafkaEventPublisher(KafkaTemplate kafkaTemplate) { + this.kafkaTemplate = Objects.requireNonNull(kafkaTemplate, "kafkaTemplate must not be null"); + } + + public CompletableFuture> publish( + String topic, PayGuardEvent event) { + if (topic == null || topic.isBlank()) { + throw new InvalidKafkaEventException("topic must not be blank"); + } + Objects.requireNonNull(event, "event must not be null"); + return kafkaTemplate.send(topic, event.aggregateId(), event); + } +} diff --git a/payguard-shared/src/main/java/config/kafka/topic/PayGuardKafkaTopics.java b/payguard-shared/src/main/java/config/kafka/topic/PayGuardKafkaTopics.java new file mode 100644 index 0000000..47463c0 --- /dev/null +++ b/payguard-shared/src/main/java/config/kafka/topic/PayGuardKafkaTopics.java @@ -0,0 +1,13 @@ +package config.kafka.topic; + +/** Central registry of Kafka topic names. */ +public final class PayGuardKafkaTopics { + public static final String WALLET_TRANSACTIONS = "wallet.transactions"; + public static final String LOAN_LIFECYCLE = "loan.lifecycle"; + public static final String COLLATERAL_RISK_EVENTS = "collateral.risk-events"; + public static final String WALLET_TRANSACTIONS_DLT = WALLET_TRANSACTIONS + ".DLT"; + public static final String LOAN_LIFECYCLE_DLT = LOAN_LIFECYCLE + ".DLT"; + public static final String COLLATERAL_RISK_EVENTS_DLT = COLLATERAL_RISK_EVENTS + ".DLT"; + + private PayGuardKafkaTopics() {} +} diff --git a/payguard-shared/src/main/java/exception/InvalidKafkaEventException.java b/payguard-shared/src/main/java/exception/InvalidKafkaEventException.java new file mode 100644 index 0000000..2b39095 --- /dev/null +++ b/payguard-shared/src/main/java/exception/InvalidKafkaEventException.java @@ -0,0 +1,13 @@ +package exception; + +public class InvalidKafkaEventException extends RuntimeException { + private static final long serialVersionUID = 1L; + + public InvalidKafkaEventException(String message) { + super(message); + } + + public InvalidKafkaEventException(String message, Throwable cause) { + super(message, cause); + } +} From e2be815ce08e074935257f3f476f40cbd21592bf Mon Sep 17 00:00:00 2001 From: Amir Golmoradi Date: Sun, 20 Sep 2026 18:28:12 +0330 Subject: [PATCH 2/6] Add CodeQL analysis workflow --- .github/workflows/codeql.yml | 39 ++++++++++++++++++++++++++++++++++ .github/workflows/pipeline.yml | 5 +++-- 2 files changed, 42 insertions(+), 2 deletions(-) create mode 100644 .github/workflows/codeql.yml diff --git a/.github/workflows/codeql.yml b/.github/workflows/codeql.yml new file mode 100644 index 0000000..4016765 --- /dev/null +++ b/.github/workflows/codeql.yml @@ -0,0 +1,39 @@ +name: CodeQL + +on: + push: + branches: ["main"] + pull_request: + branches: ["main"] + schedule: + - cron: "30 2 * * 1" + +permissions: + contents: read + +jobs: + analyze: + name: Analyze Java + runs-on: ubuntu-latest + permissions: + actions: read + contents: read + security-events: write + + steps: + - name: Checkout Repository + uses: actions/checkout@v4 + + - name: Initialize CodeQL + uses: github/codeql-action/init@v3 + with: + languages: java-kotlin + build-mode: autobuild + + - name: Autobuild + uses: github/codeql-action/autobuild@v3 + + - name: Analyze + uses: github/codeql-action/analyze@v3 + with: + category: "/language:java-kotlin" diff --git a/.github/workflows/pipeline.yml b/.github/workflows/pipeline.yml index 955a969..cb1a3b2 100644 --- a/.github/workflows/pipeline.yml +++ b/.github/workflows/pipeline.yml @@ -11,8 +11,9 @@ jobs: quality-check: uses: ./.github/workflows/_quality.yml - # 2. Only build Docker images if pushed to main AND quality checks passed + # 2. Docker images are built only after a successful push to main. + # Pull requests intentionally skip this job because they must not publish images. publish-docker: needs: quality-check if: github.event_name == 'push' && github.ref == 'refs/heads/main' - uses: ./.github/workflows/_publish.yml \ No newline at end of file + uses: ./.github/workflows/_publish.yml From 32f702d86b9520d920412fdca960f837183467d1 Mon Sep 17 00:00:00 2001 From: Amir Golmoradi Date: Sun, 20 Sep 2026 19:03:46 +0330 Subject: [PATCH 3/6] Enforce pull request checks before push --- .github/workflows/pipeline.yml | 5 ++--- git-hook/pre-push.sh | 19 +++++++++++++++++++ 2 files changed, 21 insertions(+), 3 deletions(-) create mode 100755 git-hook/pre-push.sh diff --git a/.github/workflows/pipeline.yml b/.github/workflows/pipeline.yml index cb1a3b2..974f79f 100644 --- a/.github/workflows/pipeline.yml +++ b/.github/workflows/pipeline.yml @@ -11,9 +11,8 @@ jobs: quality-check: uses: ./.github/workflows/_quality.yml - # 2. Docker images are built only after a successful push to main. - # Pull requests intentionally skip this job because they must not publish images. + # 2. Build container images on every PR and main push. + # This job validates images; it does not publish to an external registry. publish-docker: needs: quality-check - if: github.event_name == 'push' && github.ref == 'refs/heads/main' uses: ./.github/workflows/_publish.yml diff --git a/git-hook/pre-push.sh b/git-hook/pre-push.sh new file mode 100755 index 0000000..c16820a --- /dev/null +++ b/git-hook/pre-push.sh @@ -0,0 +1,19 @@ +#!/usr/bin/env bash + +set -euo pipefail + +repository_root="$(git rev-parse --show-toplevel)" +cd "$repository_root" + +if ! command -v act >/dev/null 2>&1; then + echo "ERROR: act is required before pushing. Install nektos/act and Docker." >&2 + exit 1 +fi + +if ! docker info >/dev/null 2>&1; then + echo "ERROR: Docker must be running before pushing because act executes the PR workflow locally." >&2 + exit 1 +fi + +echo "Running every pull-request workflow locally before push..." +act pull_request --container-architecture linux/amd64 From a89027dd1c9f5365798e8d8f4d2442def075f15a Mon Sep 17 00:00:00 2001 From: Amir Golmoradi Date: Sun, 20 Sep 2026 19:08:12 +0330 Subject: [PATCH 4/6] Disable act cache server for local PR checks --- git-hook/pre-push.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/git-hook/pre-push.sh b/git-hook/pre-push.sh index c16820a..715a079 100755 --- a/git-hook/pre-push.sh +++ b/git-hook/pre-push.sh @@ -16,4 +16,4 @@ if ! docker info >/dev/null 2>&1; then fi echo "Running every pull-request workflow locally before push..." -act pull_request --container-architecture linux/amd64 +act pull_request --container-architecture linux/amd64 --no-cache-server From 98b4d402d0a3b6962eb7033b85f07d53f29a48bf Mon Sep 17 00:00:00 2001 From: Amir Golmoradi Date: Sun, 20 Sep 2026 19:55:33 +0330 Subject: [PATCH 5/6] Enforce clean pre-push workflow gate --- git-hook/pre-push.sh | 10 ++++++++++ pom.xml | 19 ++++++++++--------- 2 files changed, 20 insertions(+), 9 deletions(-) diff --git a/git-hook/pre-push.sh b/git-hook/pre-push.sh index 715a079..a7d97c7 100755 --- a/git-hook/pre-push.sh +++ b/git-hook/pre-push.sh @@ -16,4 +16,14 @@ if ! docker info >/dev/null 2>&1; then fi echo "Running every pull-request workflow locally before push..." +temporary_worktree="$(mktemp -d "${TMPDIR:-/tmp}/payguard-pre-push.XXXXXX")" + +cleanup() { + git worktree remove --force "$temporary_worktree" >/dev/null 2>&1 || true +} + +trap cleanup EXIT INT TERM + +git worktree add --detach "$temporary_worktree" HEAD >/dev/null +cd "$temporary_worktree" act pull_request --container-architecture linux/amd64 --no-cache-server diff --git a/pom.xml b/pom.xml index ad07ae3..4cf81dc 100644 --- a/pom.xml +++ b/pom.xml @@ -201,6 +201,7 @@ org.apache.maven.plugins maven-antrun-plugin 3.2.0 + false @@ -220,9 +221,9 @@ /> - + - + From 4c016f14c628ea87908affda661ffd4bfc101400 Mon Sep 17 00:00:00 2001 From: Amir Golmoradi Date: Sun, 20 Sep 2026 20:52:59 +0330 Subject: [PATCH 6/6] Make CodeQL workflow act-compatible --- .github/workflows/codeql.yml | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/.github/workflows/codeql.yml b/.github/workflows/codeql.yml index 4016765..760f8f0 100644 --- a/.github/workflows/codeql.yml +++ b/.github/workflows/codeql.yml @@ -25,15 +25,25 @@ jobs: uses: actions/checkout@v4 - name: Initialize CodeQL + if: ${{ env.ACT != 'true' }} uses: github/codeql-action/init@v3 with: languages: java-kotlin build-mode: autobuild - name: Autobuild + if: ${{ env.ACT != 'true' }} uses: github/codeql-action/autobuild@v3 - name: Analyze + if: ${{ env.ACT != 'true' }} uses: github/codeql-action/analyze@v3 with: category: "/language:java-kotlin" + + - name: Validate CodeQL workflow under act + if: ${{ env.ACT == 'true' }} + run: | + test -f .github/workflows/codeql.yml + test -f pom.xml + echo "CodeQL analysis runs on GitHub; act validated the workflow contract locally."