From 3a592f15dd9dae829486eaf4839939b5b2525c29 Mon Sep 17 00:00:00 2001 From: realvic22 Date: Sat, 25 Jul 2026 19:28:24 +0100 Subject: [PATCH 1/6] feat: add on-chain ingestion pipeline --- .../api/GuildWorkmanApplication.java | 2 + .../api/chain/api/ChainEventController.java | 21 ++++++ .../api/chain/api/ChainEventResponse.java | 8 +++ .../chain/api/IngestChainEventRequest.java | 12 ++++ .../api/chain/api/ReplayRequest.java | 7 ++ .../api/chain/model/ChainEventStatus.java | 5 ++ .../api/chain/model/OnChainEvent.java | 46 +++++++++++++ .../api/chain/model/OutboxEvent.java | 31 +++++++++ .../api/chain/model/OutboxStatus.java | 5 ++ .../repository/OnChainEventRepository.java | 21 ++++++ .../repository/OutboxEventRepository.java | 16 +++++ .../api/chain/service/ChainEventService.java | 68 +++++++++++++++++++ .../api/config/SecurityConfig.java | 1 + .../src/main/resources/application.properties | 3 + .../api/chain/ChainEventServiceTest.java | 28 ++++++++ 15 files changed, 274 insertions(+) create mode 100644 backend-api/src/main/java/com/guildworkman/api/chain/api/ChainEventController.java create mode 100644 backend-api/src/main/java/com/guildworkman/api/chain/api/ChainEventResponse.java create mode 100644 backend-api/src/main/java/com/guildworkman/api/chain/api/IngestChainEventRequest.java create mode 100644 backend-api/src/main/java/com/guildworkman/api/chain/api/ReplayRequest.java create mode 100644 backend-api/src/main/java/com/guildworkman/api/chain/model/ChainEventStatus.java create mode 100644 backend-api/src/main/java/com/guildworkman/api/chain/model/OnChainEvent.java create mode 100644 backend-api/src/main/java/com/guildworkman/api/chain/model/OutboxEvent.java create mode 100644 backend-api/src/main/java/com/guildworkman/api/chain/model/OutboxStatus.java create mode 100644 backend-api/src/main/java/com/guildworkman/api/chain/repository/OnChainEventRepository.java create mode 100644 backend-api/src/main/java/com/guildworkman/api/chain/repository/OutboxEventRepository.java create mode 100644 backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java create mode 100644 backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java diff --git a/backend-api/src/main/java/com/guildworkman/api/GuildWorkmanApplication.java b/backend-api/src/main/java/com/guildworkman/api/GuildWorkmanApplication.java index 4256000..000b625 100644 --- a/backend-api/src/main/java/com/guildworkman/api/GuildWorkmanApplication.java +++ b/backend-api/src/main/java/com/guildworkman/api/GuildWorkmanApplication.java @@ -2,10 +2,12 @@ import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.scheduling.annotation.EnableScheduling; import java.time.LocalDateTime; @SpringBootApplication +@EnableScheduling public class GuildWorkmanApplication { public static void main(String[] args) { diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/api/ChainEventController.java b/backend-api/src/main/java/com/guildworkman/api/chain/api/ChainEventController.java new file mode 100644 index 0000000..99771d9 --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/api/ChainEventController.java @@ -0,0 +1,21 @@ +package com.guildworkman.api.chain.api; + +import com.guildworkman.api.chain.service.ChainEventService; +import jakarta.validation.Valid; +import lombok.RequiredArgsConstructor; +import org.springframework.http.HttpStatus; +import org.springframework.web.bind.annotation.*; + +import java.util.Map; + +@RestController @RequestMapping("/api/v1/chain/events") @RequiredArgsConstructor +public class ChainEventController { + private final ChainEventService service; + + @PostMapping + @ResponseStatus(HttpStatus.ACCEPTED) + public ChainEventResponse ingest(@Valid @RequestBody IngestChainEventRequest request) { return service.ingest(request); } + + @PostMapping("/replay") + public Map replay(@Valid @RequestBody ReplayRequest request) { return Map.of("replayed", service.replay(request)); } +} diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/api/ChainEventResponse.java b/backend-api/src/main/java/com/guildworkman/api/chain/api/ChainEventResponse.java new file mode 100644 index 0000000..d149ff1 --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/api/ChainEventResponse.java @@ -0,0 +1,8 @@ +package com.guildworkman.api.chain.api; + +import com.guildworkman.api.chain.model.OnChainEvent; +import java.time.Instant; + +public record ChainEventResponse(Long id, String eventKey, String contractId, long ledger, int eventIndex, String topics, String payload, String status, int attempts, Instant processedAt) { + public static ChainEventResponse from(OnChainEvent e) { return new ChainEventResponse(e.getId(), e.getEventKey(), e.getContractId(), e.getLedger(), e.getEventIndex(), e.getTopics(), e.getPayload(), e.getStatus().name(), e.getAttempts(), e.getProcessedAt()); } +} diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/api/IngestChainEventRequest.java b/backend-api/src/main/java/com/guildworkman/api/chain/api/IngestChainEventRequest.java new file mode 100644 index 0000000..df80a5a --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/api/IngestChainEventRequest.java @@ -0,0 +1,12 @@ +package com.guildworkman.api.chain.api; + +import jakarta.validation.constraints.*; +import java.util.List; + +public record IngestChainEventRequest( + @NotBlank @Size(max = 128) String eventKey, + @NotBlank @Size(max = 128) String contractId, + @PositiveOrZero long ledger, + @PositiveOrZero int eventIndex, + @NotEmpty List<@NotBlank @Size(max = 128) String> topics, + @NotBlank String payload) {} diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/api/ReplayRequest.java b/backend-api/src/main/java/com/guildworkman/api/chain/api/ReplayRequest.java new file mode 100644 index 0000000..0ebfab5 --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/api/ReplayRequest.java @@ -0,0 +1,7 @@ +package com.guildworkman.api.chain.api; + +import jakarta.validation.constraints.PositiveOrZero; + +public record ReplayRequest(@PositiveOrZero long fromLedger, @PositiveOrZero long toLedger) { + public ReplayRequest { if (toLedger < fromLedger) throw new IllegalArgumentException("toLedger must be greater than or equal to fromLedger"); } +} diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/model/ChainEventStatus.java b/backend-api/src/main/java/com/guildworkman/api/chain/model/ChainEventStatus.java new file mode 100644 index 0000000..9f93982 --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/model/ChainEventStatus.java @@ -0,0 +1,5 @@ +package com.guildworkman.api.chain.model; + +public enum ChainEventStatus { + PENDING, PROCESSING, PROCESSED, DEAD_LETTER +} diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/model/OnChainEvent.java b/backend-api/src/main/java/com/guildworkman/api/chain/model/OnChainEvent.java new file mode 100644 index 0000000..f9eaaa1 --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/model/OnChainEvent.java @@ -0,0 +1,46 @@ +package com.guildworkman.api.chain.model; + +import jakarta.persistence.*; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.Setter; + +import java.time.Instant; + +@Entity +@Table(name = "on_chain_events", uniqueConstraints = @UniqueConstraint(name = "uk_chain_event_key", columnNames = "event_key"), indexes = { + @Index(name = "idx_chain_event_stream_order", columnList = "contract_id,ledger,event_index"), + @Index(name = "idx_chain_event_status", columnList = "status,next_attempt_at") +}) +@Getter +@Setter +@NoArgsConstructor +public class OnChainEvent { + @Id @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + @Column(name = "event_key", nullable = false, updatable = false, length = 128) + private String eventKey; + @Column(name = "contract_id", nullable = false, length = 128) + private String contractId; + @Column(nullable = false) + private long ledger; + @Column(name = "event_index", nullable = false) + private int eventIndex; + @Column(nullable = false, length = 256) + private String topics; + @Lob @Column(nullable = false) + private String payload; + @Enumerated(EnumType.STRING) @Column(nullable = false, length = 20) + private ChainEventStatus status = ChainEventStatus.PENDING; + @Column(nullable = false) + private int attempts; + @Column(name = "next_attempt_at", nullable = false) + private Instant nextAttemptAt = Instant.now(); + @Column(length = 1000) + private String lastError; + @Column(nullable = false, updatable = false) + private Instant createdAt = Instant.now(); + private Instant processedAt; + @Version + private long version; +} diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/model/OutboxEvent.java b/backend-api/src/main/java/com/guildworkman/api/chain/model/OutboxEvent.java new file mode 100644 index 0000000..9a0cca7 --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/model/OutboxEvent.java @@ -0,0 +1,31 @@ +package com.guildworkman.api.chain.model; + +import jakarta.persistence.*; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.Setter; + +import java.time.Instant; + +@Entity +@Table(name = "chain_event_outbox", uniqueConstraints = @UniqueConstraint(name = "uk_outbox_event", columnNames = "event_id"), indexes = @Index(name = "idx_outbox_status", columnList = "status,next_attempt_at")) +@Getter @Setter @NoArgsConstructor +public class OutboxEvent { + @Id @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + @Column(name = "event_id", nullable = false, updatable = false) + private Long eventId; + @Enumerated(EnumType.STRING) @Column(nullable = false, length = 20) + private OutboxStatus status = OutboxStatus.PENDING; + @Column(nullable = false) + private int attempts; + @Column(name = "next_attempt_at", nullable = false) + private Instant nextAttemptAt = Instant.now(); + @Column(length = 1000) + private String lastError; + @Column(nullable = false, updatable = false) + private Instant createdAt = Instant.now(); + private Instant completedAt; + @Version + private long version; +} diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/model/OutboxStatus.java b/backend-api/src/main/java/com/guildworkman/api/chain/model/OutboxStatus.java new file mode 100644 index 0000000..301458f --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/model/OutboxStatus.java @@ -0,0 +1,5 @@ +package com.guildworkman.api.chain.model; + +public enum OutboxStatus { + PENDING, PROCESSING, COMPLETED, DEAD_LETTER +} diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/repository/OnChainEventRepository.java b/backend-api/src/main/java/com/guildworkman/api/chain/repository/OnChainEventRepository.java new file mode 100644 index 0000000..aa99505 --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/repository/OnChainEventRepository.java @@ -0,0 +1,21 @@ +package com.guildworkman.api.chain.repository; + +import com.guildworkman.api.chain.model.ChainEventStatus; +import com.guildworkman.api.chain.model.OnChainEvent; +import org.springframework.data.domain.Pageable; +import org.springframework.data.jpa.repository.*; +import org.springframework.data.repository.query.Param; + +import jakarta.persistence.LockModeType; +import java.time.Instant; +import java.util.*; + +public interface OnChainEventRepository extends JpaRepository { + Optional findByEventKey(String eventKey); + boolean existsByEventKey(String eventKey); + @Lock(LockModeType.PESSIMISTIC_WRITE) + @Query("select e from OnChainEvent e where e.status in :statuses and e.nextAttemptAt <= :now order by e.contractId, e.ledger, e.eventIndex, e.id") + List claimNext(@Param("statuses") Set statuses, @Param("now") Instant now, Pageable pageable); + List findByLedgerBetweenOrderByContractIdAscLedgerAscEventIndexAsc(long fromLedger, long toLedger); + long countByStatus(ChainEventStatus status); +} diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/repository/OutboxEventRepository.java b/backend-api/src/main/java/com/guildworkman/api/chain/repository/OutboxEventRepository.java new file mode 100644 index 0000000..d32955f --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/repository/OutboxEventRepository.java @@ -0,0 +1,16 @@ +package com.guildworkman.api.chain.repository; + +import com.guildworkman.api.chain.model.*; +import org.springframework.data.domain.Pageable; +import org.springframework.data.jpa.repository.*; +import org.springframework.data.repository.query.Param; +import jakarta.persistence.LockModeType; +import java.time.Instant; +import java.util.*; + +public interface OutboxEventRepository extends JpaRepository { + @Lock(LockModeType.PESSIMISTIC_WRITE) + @Query("select o from OutboxEvent o where o.status in :statuses and o.nextAttemptAt <= :now order by o.id") + List claimNext(@Param("statuses") Set statuses, @Param("now") Instant now, Pageable pageable); + Optional findByEventId(Long eventId); +} diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java b/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java new file mode 100644 index 0000000..8537758 --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java @@ -0,0 +1,68 @@ +package com.guildworkman.api.chain.service; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.guildworkman.api.chain.api.*; +import com.guildworkman.api.chain.model.*; +import com.guildworkman.api.chain.repository.*; +import jakarta.transaction.Transactional; +import lombok.RequiredArgsConstructor; +import org.springframework.data.domain.PageRequest; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Service; + +import java.time.*; +import java.util.*; + +@Service @RequiredArgsConstructor +public class ChainEventService { + private final OnChainEventRepository events; + private final OutboxEventRepository outbox; + private final ObjectMapper objectMapper; + private static final int MAX_ATTEMPTS = 5; + + @Transactional + public ChainEventResponse ingest(IngestChainEventRequest request) { + return events.findByEventKey(request.eventKey()).map(ChainEventResponse::from).orElseGet(() -> { + OnChainEvent event = new OnChainEvent(); + event.setEventKey(request.eventKey()); event.setContractId(request.contractId()); + event.setLedger(request.ledger()); event.setEventIndex(request.eventIndex()); + try { event.setTopics(objectMapper.writeValueAsString(request.topics())); } + catch (Exception ex) { throw new IllegalArgumentException("topics must be serializable", ex); } + event.setPayload(request.payload()); event.setStatus(ChainEventStatus.PENDING); + event.setNextAttemptAt(Instant.now()); + OnChainEvent saved = events.save(event); + OutboxEvent message = new OutboxEvent(); message.setEventId(saved.getId()); outbox.save(message); + return ChainEventResponse.from(saved); + }); + } + + @Transactional + public int replay(ReplayRequest request) { + int count = 0; + for (OnChainEvent event : events.findByLedgerBetweenOrderByContractIdAscLedgerAscEventIndexAsc(request.fromLedger(), request.toLedger())) { + event.setStatus(ChainEventStatus.PENDING); event.setAttempts(0); event.setLastError(null); event.setProcessedAt(null); event.setNextAttemptAt(Instant.now()); + outbox.findByEventId(event.getId()).ifPresent(message -> { message.setStatus(OutboxStatus.PENDING); message.setAttempts(0); message.setLastError(null); message.setCompletedAt(null); message.setNextAttemptAt(Instant.now()); }); + count++; + } + return count; + } + + @Scheduled(fixedDelayString = "${chain.events.poll-delay-ms:1000}") + @Transactional + public void processOne() { + events.claimNext(EnumSet.of(ChainEventStatus.PENDING, ChainEventStatus.PROCESSING), Instant.now(), PageRequest.of(0, 1)).stream().findFirst().ifPresent(this::process); + } + + private void process(OnChainEvent event) { + try { + event.setStatus(ChainEventStatus.PROCESSING); event.setAttempts(event.getAttempts() + 1); + // The event row is the durable projection. Keeping the transition in + // the same transaction as the outbox acknowledgement makes retries safe. + event.setStatus(ChainEventStatus.PROCESSED); event.setProcessedAt(Instant.now()); event.setLastError(null); + outbox.findByEventId(event.getId()).ifPresent(message -> { message.setStatus(OutboxStatus.COMPLETED); message.setCompletedAt(Instant.now()); message.setLastError(null); }); + } catch (RuntimeException ex) { + event.setLastError(ex.getMessage()); + if (event.getAttempts() >= MAX_ATTEMPTS) event.setStatus(ChainEventStatus.DEAD_LETTER); else { event.setStatus(ChainEventStatus.PENDING); event.setNextAttemptAt(Instant.now().plusSeconds(1L << Math.min(event.getAttempts(), 6))); } + } + } +} diff --git a/backend-api/src/main/java/com/guildworkman/api/config/SecurityConfig.java b/backend-api/src/main/java/com/guildworkman/api/config/SecurityConfig.java index 20c8fff..5f3860a 100644 --- a/backend-api/src/main/java/com/guildworkman/api/config/SecurityConfig.java +++ b/backend-api/src/main/java/com/guildworkman/api/config/SecurityConfig.java @@ -58,6 +58,7 @@ public class SecurityConfig { "/api/v1/auth/logout", "/api/v1/client/**", "/api/v1/skilledWorker/**", + "/api/v1/chain/events/**", "/v3/api-docs/**", "/swagger-ui/**", "/swagger-ui.html" diff --git a/backend-api/src/main/resources/application.properties b/backend-api/src/main/resources/application.properties index c0a182c..918a5c2 100644 --- a/backend-api/src/main/resources/application.properties +++ b/backend-api/src/main/resources/application.properties @@ -33,6 +33,9 @@ spring.mvc.problemdetails.enabled=true # staging/prod can point at their own canonical docs host. guildworkman.problem-details.type-base=${PROBLEM_TYPE_BASE:https://guildworkman.dev/problems/} +# On-chain ingestion worker. Set to 0 to disable polling in a batch/replay job. +chain.events.poll-delay-ms=${CHAIN_EVENTS_POLL_DELAY_MS:1000} + # JWT authentication. # - secret: HMAC-256 signing key. MUST be overridden in every real # environment via JWT_SECRET (>= 32 bytes). The baked-in default exists diff --git a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java new file mode 100644 index 0000000..4cec348 --- /dev/null +++ b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java @@ -0,0 +1,28 @@ +package com.guildworkman.api.chain; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.guildworkman.api.chain.api.IngestChainEventRequest; +import com.guildworkman.api.chain.repository.*; +import com.guildworkman.api.chain.service.ChainEventService; +import org.junit.jupiter.api.Test; +import org.mockito.Mockito; + +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; + +class ChainEventServiceTest { + @Test + void ingestIsIdempotentForTheSameEventKey() { + var events = Mockito.mock(OnChainEventRepository.class); + var outbox = Mockito.mock(OutboxEventRepository.class); + var service = new ChainEventService(events, outbox, new ObjectMapper()); + var request = new IngestChainEventRequest("evt-1", "CABC", 10, 0, List.of("Transfer"), "{\"amount\":1}"); + Mockito.when(events.findByEventKey("evt-1")).thenReturn(java.util.Optional.empty()); + var event = new com.guildworkman.api.chain.model.OnChainEvent(); event.setId(1L); event.setEventKey("evt-1"); event.setContractId("CABC"); event.setTopics("[\"Transfer\"]"); event.setPayload(request.payload()); + Mockito.when(events.save(Mockito.any())).thenReturn(event); + assertThat(service.ingest(request).eventKey()).isEqualTo("evt-1"); + Mockito.verify(events).save(Mockito.any()); + Mockito.verify(outbox).save(Mockito.any()); + } +} From 54234e585585cb863c5bcac6bf0e342fce43afa8 Mon Sep 17 00:00:00 2001 From: realvic22 Date: Mon, 27 Jul 2026 10:56:33 +0100 Subject: [PATCH 2/6] fix: address review feedback - persist replay changes, improve ingest error handling, add comprehensive tests - Fix replay() to explicitly call saveAll() to persist event/outbox state changes - Improve ingest() error handling: catch DataIntegrityViolationException and perform fallback lookup to guarantee idempotency even during race conditions - Add ChainEventHandler functional interface for extensible and testable event processing pipeline - Add comprehensive integration tests (ChainEventServiceIntegrationTest): - Transactional outbox semantics (event + outbox created atomically) - Happy-path processing (PENDING -> PROCESSED, outbox COMPLETED) - Failure path: dead-letter transition after MAX_ATTEMPTS (5) retries - Exponential backoff on transient failures - Replay persistence and reprocessing verification - Concurrency test: pessimistic locking prevents duplicate processing - Already-processed events are not claimed again - Expand unit tests (ChainEventServiceTest) covering: - Idempotency when event already exists - DataIntegrityViolationException with fallback lookup - Unexpected exception propagation - Replay explicit saveAll verification --- .../api/chain/service/ChainEventHandler.java | 8 + .../api/chain/service/ChainEventService.java | 29 +- .../ChainEventServiceIntegrationTest.java | 266 ++++++++++++++++++ .../api/chain/ChainEventServiceTest.java | 110 +++++++- 4 files changed, 402 insertions(+), 11 deletions(-) create mode 100644 backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventHandler.java create mode 100644 backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventHandler.java b/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventHandler.java new file mode 100644 index 0000000..c5b5db8 --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventHandler.java @@ -0,0 +1,8 @@ +package com.guildworkman.api.chain.service; + +import com.guildworkman.api.chain.model.OnChainEvent; + +@FunctionalInterface +public interface ChainEventHandler { + void handle(OnChainEvent event); +} diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java b/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java index 8537758..0040b3e 100644 --- a/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java +++ b/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java @@ -6,6 +6,7 @@ import com.guildworkman.api.chain.repository.*; import jakarta.transaction.Transactional; import lombok.RequiredArgsConstructor; +import org.springframework.dao.DataIntegrityViolationException; import org.springframework.data.domain.PageRequest; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; @@ -18,11 +19,18 @@ public class ChainEventService { private final OnChainEventRepository events; private final OutboxEventRepository outbox; private final ObjectMapper objectMapper; + private final List handlers; private static final int MAX_ATTEMPTS = 5; @Transactional public ChainEventResponse ingest(IngestChainEventRequest request) { - return events.findByEventKey(request.eventKey()).map(ChainEventResponse::from).orElseGet(() -> { + return events.findByEventKey(request.eventKey()) + .map(ChainEventResponse::from) + .orElseGet(() -> createEvent(request)); + } + + private ChainEventResponse createEvent(IngestChainEventRequest request) { + try { OnChainEvent event = new OnChainEvent(); event.setEventKey(request.eventKey()); event.setContractId(request.contractId()); event.setLedger(request.ledger()); event.setEventIndex(request.eventIndex()); @@ -33,17 +41,27 @@ public ChainEventResponse ingest(IngestChainEventRequest request) { OnChainEvent saved = events.save(event); OutboxEvent message = new OutboxEvent(); message.setEventId(saved.getId()); outbox.save(message); return ChainEventResponse.from(saved); - }); + } catch (DataIntegrityViolationException ex) { + return events.findByEventKey(request.eventKey()) + .map(ChainEventResponse::from) + .orElseThrow(() -> new IllegalStateException("Event not found after idempotent-guard violation for key=" + request.eventKey(), ex)); + } catch (RuntimeException ex) { + return events.findByEventKey(request.eventKey()) + .map(ChainEventResponse::from) + .orElseThrow(() -> ex); + } } @Transactional public int replay(ReplayRequest request) { int count = 0; - for (OnChainEvent event : events.findByLedgerBetweenOrderByContractIdAscLedgerAscEventIndexAsc(request.fromLedger(), request.toLedger())) { + List batch = events.findByLedgerBetweenOrderByContractIdAscLedgerAscEventIndexAsc(request.fromLedger(), request.toLedger()); + for (OnChainEvent event : batch) { event.setStatus(ChainEventStatus.PENDING); event.setAttempts(0); event.setLastError(null); event.setProcessedAt(null); event.setNextAttemptAt(Instant.now()); outbox.findByEventId(event.getId()).ifPresent(message -> { message.setStatus(OutboxStatus.PENDING); message.setAttempts(0); message.setLastError(null); message.setCompletedAt(null); message.setNextAttemptAt(Instant.now()); }); count++; } + events.saveAll(batch); return count; } @@ -53,11 +71,10 @@ public void processOne() { events.claimNext(EnumSet.of(ChainEventStatus.PENDING, ChainEventStatus.PROCESSING), Instant.now(), PageRequest.of(0, 1)).stream().findFirst().ifPresent(this::process); } - private void process(OnChainEvent event) { + void process(OnChainEvent event) { try { event.setStatus(ChainEventStatus.PROCESSING); event.setAttempts(event.getAttempts() + 1); - // The event row is the durable projection. Keeping the transition in - // the same transaction as the outbox acknowledgement makes retries safe. + for (ChainEventHandler handler : handlers) handler.handle(event); event.setStatus(ChainEventStatus.PROCESSED); event.setProcessedAt(Instant.now()); event.setLastError(null); outbox.findByEventId(event.getId()).ifPresent(message -> { message.setStatus(OutboxStatus.COMPLETED); message.setCompletedAt(Instant.now()); message.setLastError(null); }); } catch (RuntimeException ex) { diff --git a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java new file mode 100644 index 0000000..025925e --- /dev/null +++ b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java @@ -0,0 +1,266 @@ +package com.guildworkman.api.chain; + +import com.guildworkman.api.chain.api.ChainEventResponse; +import com.guildworkman.api.chain.api.IngestChainEventRequest; +import com.guildworkman.api.chain.api.ReplayRequest; +import com.guildworkman.api.chain.model.ChainEventStatus; +import com.guildworkman.api.chain.model.OnChainEvent; +import com.guildworkman.api.chain.model.OutboxEvent; +import com.guildworkman.api.chain.model.OutboxStatus; +import com.guildworkman.api.chain.repository.OnChainEventRepository; +import com.guildworkman.api.chain.repository.OutboxEventRepository; +import com.guildworkman.api.chain.service.ChainEventHandler; +import com.guildworkman.api.chain.service.ChainEventService; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.mock.mockito.MockBean; + +import java.util.List; +import java.util.concurrent.Callable; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doThrow; + +@SpringBootTest(properties = { "chain.events.poll-delay-ms=60000" }) +class ChainEventServiceIntegrationTest { + + @Autowired + private OnChainEventRepository events; + + @Autowired + private OutboxEventRepository outbox; + + @Autowired + private ChainEventService service; + + @MockBean + private ChainEventHandler chainEventHandler; + + @BeforeEach + void cleanSlate() { + outbox.deleteAll(); + events.deleteAll(); + } + + private OnChainEvent saveEvent(String eventKey, ChainEventStatus status, int attempts) { + OnChainEvent e = new OnChainEvent(); + e.setEventKey(eventKey); + e.setContractId("CTEST"); + e.setLedger(100); + e.setEventIndex(0); + e.setTopics("[\"Test\"]"); + e.setPayload("{\"x\":1}"); + e.setStatus(status); + e.setAttempts(attempts); + e.setNextAttemptAt(java.time.Instant.now()); + return events.save(e); + } + + private OutboxEvent saveOutbox(Long eventId, OutboxStatus status) { + OutboxEvent o = new OutboxEvent(); + o.setEventId(eventId); + o.setStatus(status); + return outbox.save(o); + } + + // ------------------------------------------------------------------------- + // 1. Transactional outbox: ingest creates both rows atomically + // ------------------------------------------------------------------------- + + @Test + void ingestCreatesBothEventAndOutboxRow() { + IngestChainEventRequest req = new IngestChainEventRequest("k1", "C001", 1, 0, List.of("T"), "{}"); + ChainEventResponse resp = service.ingest(req); + + OnChainEvent saved = events.findById(resp.id()).orElseThrow(); + assertThat(saved.getStatus()).isEqualTo(ChainEventStatus.PENDING); + + OutboxEvent msg = outbox.findByEventId(resp.id()).orElseThrow(); + assertThat(msg.getStatus()).isEqualTo(OutboxStatus.PENDING); + assertThat(msg.getEventId()).isEqualTo(saved.getId()); + } + + // ------------------------------------------------------------------------- + // 2. Ingest idempotency + // ------------------------------------------------------------------------- + + @Test + void ingestIsIdempotentForDuplicateEventKey() { + IngestChainEventRequest req = new IngestChainEventRequest("dup", "C002", 2, 0, List.of("T"), "{}"); + ChainEventResponse first = service.ingest(req); + ChainEventResponse second = service.ingest(req); + + assertThat(second.id()).isEqualTo(first.id()); + assertThat(events.count()).isEqualTo(1); + assertThat(outbox.count()).isEqualTo(1); + } + + // ------------------------------------------------------------------------- + // 3. Happy-path processing + // ------------------------------------------------------------------------- + + @Test + void processOneTransitionsEventToProcessedAndOutboxToCompleted() { + OnChainEvent event = saveEvent("happy", ChainEventStatus.PENDING, 0); + saveOutbox(event.getId(), OutboxStatus.PENDING); + + service.processOne(); + + OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); + assertThat(reloaded.getStatus()).isEqualTo(ChainEventStatus.PROCESSED); + assertThat(reloaded.getAttempts()).isEqualTo(1); + assertThat(reloaded.getProcessedAt()).isNotNull(); + + OutboxEvent msg = outbox.findByEventId(event.getId()).orElseThrow(); + assertThat(msg.getStatus()).isEqualTo(OutboxStatus.COMPLETED); + assertThat(msg.getCompletedAt()).isNotNull(); + } + + // ------------------------------------------------------------------------- + // 4. Failure path: handler exception -> DEAD_LETTER after MAX_ATTEMPTS + // ------------------------------------------------------------------------- + + @Test + void failureExhaustingRetriesMovesToDeadLetter() { + OnChainEvent event = saveEvent("dead", ChainEventStatus.PENDING, 0); + saveOutbox(event.getId(), OutboxStatus.PENDING); + + doThrow(new RuntimeException("simulated processing failure")) + .when(chainEventHandler).handle(any()); + + for (int i = 0; i < 5; i++) { + try { + service.processOne(); + } catch (RuntimeException ignored) { + } + } + + OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); + assertThat(reloaded.getStatus()).isEqualTo(ChainEventStatus.DEAD_LETTER); + assertThat(reloaded.getAttempts()).isEqualTo(5); + assertThat(reloaded.getLastError()).isEqualTo("simulated processing failure"); + } + + // ------------------------------------------------------------------------- + // 5. Backoff after failure + // ------------------------------------------------------------------------- + + @Test + void failureAppliesExponentialBackoff() { + OnChainEvent event = saveEvent("backoff", ChainEventStatus.PENDING, 0); + saveOutbox(event.getId(), OutboxStatus.PENDING); + + doThrow(new RuntimeException("transient error")) + .when(chainEventHandler).handle(any()); + + try { + service.processOne(); + } catch (RuntimeException ignored) { + } + + OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); + assertThat(reloaded.getStatus()).isEqualTo(ChainEventStatus.PENDING); + assertThat(reloaded.getAttempts()).isEqualTo(1); + assertThat(reloaded.getNextAttemptAt()).isAfter(java.time.Instant.now()); + } + + // ------------------------------------------------------------------------- + // 6. Replay persists changes + // ------------------------------------------------------------------------- + + @Test + void replayResetsEventAndOutboxStateAndPersists() { + OnChainEvent event = saveEvent("replay-me", ChainEventStatus.PROCESSED, 3); + event.setLastError("some error"); + events.save(event); + saveOutbox(event.getId(), OutboxStatus.COMPLETED); + + ReplayRequest req = new ReplayRequest(100, 100); + int count = service.replay(req); + + assertThat(count).isEqualTo(1); + + OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); + assertThat(reloaded.getStatus()).isEqualTo(ChainEventStatus.PENDING); + assertThat(reloaded.getAttempts()).isZero(); + assertThat(reloaded.getLastError()).isNull(); + assertThat(reloaded.getProcessedAt()).isNull(); + + OutboxEvent msg = outbox.findByEventId(event.getId()).orElseThrow(); + assertThat(msg.getStatus()).isEqualTo(OutboxStatus.PENDING); + assertThat(msg.getAttempts()).isZero(); + assertThat(msg.getLastError()).isNull(); + assertThat(msg.getCompletedAt()).isNull(); + } + + // ------------------------------------------------------------------------- + // 7. Replayed events can be reprocessed + // ------------------------------------------------------------------------- + + @Test + void replayedEventsCanBeReprocessed() { + OnChainEvent event = saveEvent("reprocess-me", ChainEventStatus.PROCESSED, 5); + saveOutbox(event.getId(), OutboxStatus.COMPLETED); + + service.replay(new ReplayRequest(100, 100)); + service.processOne(); + + OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); + assertThat(reloaded.getStatus()).isEqualTo(ChainEventStatus.PROCESSED); + assertThat(reloaded.getAttempts()).isEqualTo(1); + assertThat(reloaded.getProcessedAt()).isNotNull(); + } + + // ------------------------------------------------------------------------- + // 8. Concurrency: pessimistic locking prevents duplicate processing + // ------------------------------------------------------------------------- + + @Test + void pessimisticLockingPreventsDuplicateProcessing() throws Exception { + OnChainEvent event = saveEvent("concurrent", ChainEventStatus.PENDING, 0); + saveOutbox(event.getId(), OutboxStatus.PENDING); + + int threads = 8; + ExecutorService pool = Executors.newFixedThreadPool(threads); + + List> tasks = java.util.stream.IntStream.range(0, threads) + .mapToObj(i -> (Callable) () -> { service.processOne(); return null; }) + .toList(); + + for (Future f : pool.invokeAll(tasks)) { + f.get(); + } + pool.shutdown(); + + OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); + assertThat(reloaded.getStatus()).isEqualTo(ChainEventStatus.PROCESSED); + assertThat(reloaded.getAttempts()).isEqualTo(1); + } + + // ------------------------------------------------------------------------- + // 9. Processed events are not claimed again + // ------------------------------------------------------------------------- + + @Test + void alreadyProcessedEventsAreNotClaimedAgain() { + OnChainEvent event = saveEvent("claimed", ChainEventStatus.PENDING, 0); + saveOutbox(event.getId(), OutboxStatus.PENDING); + + service.processOne(); + + OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); + assertThat(reloaded.getStatus()).isEqualTo(ChainEventStatus.PROCESSED); + assertThat(reloaded.getAttempts()).isEqualTo(1); + + service.processOne(); + + reloaded = events.findById(event.getId()).orElseThrow(); + assertThat(reloaded.getAttempts()).isEqualTo(1); + } +} diff --git a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java index 4cec348..77ca413 100644 --- a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java +++ b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java @@ -2,27 +2,127 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.guildworkman.api.chain.api.IngestChainEventRequest; +import com.guildworkman.api.chain.model.ChainEventStatus; +import com.guildworkman.api.chain.model.OnChainEvent; import com.guildworkman.api.chain.repository.*; import com.guildworkman.api.chain.service.ChainEventService; import org.junit.jupiter.api.Test; import org.mockito.Mockito; +import org.springframework.dao.DataIntegrityViolationException; import java.util.List; +import java.util.Optional; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.*; class ChainEventServiceTest { + @Test void ingestIsIdempotentForTheSameEventKey() { var events = Mockito.mock(OnChainEventRepository.class); var outbox = Mockito.mock(OutboxEventRepository.class); var service = new ChainEventService(events, outbox, new ObjectMapper()); var request = new IngestChainEventRequest("evt-1", "CABC", 10, 0, List.of("Transfer"), "{\"amount\":1}"); - Mockito.when(events.findByEventKey("evt-1")).thenReturn(java.util.Optional.empty()); - var event = new com.guildworkman.api.chain.model.OnChainEvent(); event.setId(1L); event.setEventKey("evt-1"); event.setContractId("CABC"); event.setTopics("[\"Transfer\"]"); event.setPayload(request.payload()); - Mockito.when(events.save(Mockito.any())).thenReturn(event); + when(events.findByEventKey("evt-1")).thenReturn(Optional.empty()); + var event = new OnChainEvent(); + event.setId(1L); + event.setEventKey("evt-1"); + event.setContractId("CABC"); + event.setStatus(ChainEventStatus.PENDING); + event.setTopics("[\"Transfer\"]"); + event.setPayload(request.payload()); + when(events.save(any())).thenReturn(event); assertThat(service.ingest(request).eventKey()).isEqualTo("evt-1"); - Mockito.verify(events).save(Mockito.any()); - Mockito.verify(outbox).save(Mockito.any()); + verify(events).save(any()); + verify(outbox).save(any()); + } + + @Test + void ingestReturnsExistingEventOnDuplicateEventKey() { + var events = Mockito.mock(OnChainEventRepository.class); + var outbox = Mockito.mock(OutboxEventRepository.class); + var service = new ChainEventService(events, outbox, new ObjectMapper()); + var request = new IngestChainEventRequest("evt-dup", "CABC", 10, 0, List.of("Transfer"), "{}"); + var existing = new OnChainEvent(); + existing.setId(1L); + existing.setEventKey("evt-dup"); + existing.setContractId("CABC"); + existing.setStatus(ChainEventStatus.PENDING); + existing.setTopics("[\"Transfer\"]"); + existing.setPayload("{}"); + when(events.findByEventKey("evt-dup")).thenReturn(Optional.of(existing)); + var resp = service.ingest(request); + assertThat(resp.id()).isEqualTo(1L); + verify(events, never()).save(any()); + verify(outbox, never()).save(any()); + } + + @Test + void ingestHandlesDataIntegrityViolationWithFallbackLookup() { + var events = Mockito.mock(OnChainEventRepository.class); + var outbox = Mockito.mock(OutboxEventRepository.class); + var service = new ChainEventService(events, outbox, new ObjectMapper()); + var request = new IngestChainEventRequest("race", "CABC", 10, 0, List.of("T"), "{}"); + + when(events.findByEventKey("race")).thenReturn(Optional.empty()); + when(events.save(any())).thenThrow(new DataIntegrityViolationException("dup key")); + + var afterSave = new OnChainEvent(); + afterSave.setId(42L); + afterSave.setEventKey("race"); + afterSave.setContractId("CABC"); + afterSave.setStatus(ChainEventStatus.PENDING); + afterSave.setTopics("[\"T\"]"); + afterSave.setPayload("{}"); + + // After the save fails, createEvent calls findByIdempotentKey again + when(events.findByEventKey("race")).thenReturn(Optional.empty(), Optional.of(afterSave)); + + var resp = service.ingest(request); + assertThat(resp.id()).isEqualTo(42L); + assertThat(resp.eventKey()).isEqualTo("race"); + } + + @Test + void ingestPropagatesUnexpectedExceptionWhenNoEventFound() { + var events = Mockito.mock(OnChainEventRepository.class); + var outbox = Mockito.mock(OutboxEventRepository.class); + var service = new ChainEventService(events, outbox, new ObjectMapper()); + var request = new IngestChainEventRequest("boom", "CABC", 10, 0, List.of("T"), "{}"); + + when(events.findByEventKey("boom")).thenReturn(Optional.empty()); + when(events.save(any())).thenThrow(new RuntimeException("db connection lost")); + + assertThatThrownBy(() -> service.ingest(request)) + .isInstanceOf(RuntimeException.class) + .hasMessage("db connection lost"); + } + + @Test + void replayPersistsChangesExplicitly() { + var events = Mockito.mock(OnChainEventRepository.class); + var outbox = Mockito.mock(OutboxEventRepository.class); + var service = new ChainEventService(events, outbox, new ObjectMapper()); + + var event = new OnChainEvent(); + event.setId(1L); + event.setEventKey("r"); + event.setStatus(ChainEventStatus.PROCESSED); + event.setAttempts(3); + event.setLastError("err"); + event.setProcessedAt(java.time.Instant.now()); + event.setLedger(10); + + when(events.findByLedgerBetweenOrderByContractIdAscLedgerAscEventIndexAsc(5, 15)) + .thenReturn(List.of(event)); + + var replay = new com.guildworkman.api.chain.api.ReplayRequest(5, 15); + int count = service.replay(replay); + + assertThat(count).isEqualTo(1); + verify(events).saveAll(any()); } } From 4e8daa876e84baa9cf9c4c7670c506f6e894fe38 Mon Sep 17 00:00:00 2001 From: realvic22 Date: Mon, 27 Jul 2026 11:22:15 +0100 Subject: [PATCH 3/6] fix: update unit test constructor calls to include empty handler list --- .../guildworkman/api/chain/ChainEventServiceTest.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java index 77ca413..f2b5f3f 100644 --- a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java +++ b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java @@ -24,7 +24,7 @@ class ChainEventServiceTest { void ingestIsIdempotentForTheSameEventKey() { var events = Mockito.mock(OnChainEventRepository.class); var outbox = Mockito.mock(OutboxEventRepository.class); - var service = new ChainEventService(events, outbox, new ObjectMapper()); + var service = new ChainEventService(events, outbox, new ObjectMapper(), List.of()); var request = new IngestChainEventRequest("evt-1", "CABC", 10, 0, List.of("Transfer"), "{\"amount\":1}"); when(events.findByEventKey("evt-1")).thenReturn(Optional.empty()); var event = new OnChainEvent(); @@ -44,7 +44,7 @@ void ingestIsIdempotentForTheSameEventKey() { void ingestReturnsExistingEventOnDuplicateEventKey() { var events = Mockito.mock(OnChainEventRepository.class); var outbox = Mockito.mock(OutboxEventRepository.class); - var service = new ChainEventService(events, outbox, new ObjectMapper()); + var service = new ChainEventService(events, outbox, new ObjectMapper(), List.of()); var request = new IngestChainEventRequest("evt-dup", "CABC", 10, 0, List.of("Transfer"), "{}"); var existing = new OnChainEvent(); existing.setId(1L); @@ -64,7 +64,7 @@ void ingestReturnsExistingEventOnDuplicateEventKey() { void ingestHandlesDataIntegrityViolationWithFallbackLookup() { var events = Mockito.mock(OnChainEventRepository.class); var outbox = Mockito.mock(OutboxEventRepository.class); - var service = new ChainEventService(events, outbox, new ObjectMapper()); + var service = new ChainEventService(events, outbox, new ObjectMapper(), List.of()); var request = new IngestChainEventRequest("race", "CABC", 10, 0, List.of("T"), "{}"); when(events.findByEventKey("race")).thenReturn(Optional.empty()); @@ -90,7 +90,7 @@ void ingestHandlesDataIntegrityViolationWithFallbackLookup() { void ingestPropagatesUnexpectedExceptionWhenNoEventFound() { var events = Mockito.mock(OnChainEventRepository.class); var outbox = Mockito.mock(OutboxEventRepository.class); - var service = new ChainEventService(events, outbox, new ObjectMapper()); + var service = new ChainEventService(events, outbox, new ObjectMapper(), List.of()); var request = new IngestChainEventRequest("boom", "CABC", 10, 0, List.of("T"), "{}"); when(events.findByEventKey("boom")).thenReturn(Optional.empty()); @@ -105,7 +105,7 @@ void ingestPropagatesUnexpectedExceptionWhenNoEventFound() { void replayPersistsChangesExplicitly() { var events = Mockito.mock(OnChainEventRepository.class); var outbox = Mockito.mock(OutboxEventRepository.class); - var service = new ChainEventService(events, outbox, new ObjectMapper()); + var service = new ChainEventService(events, outbox, new ObjectMapper(), List.of()); var event = new OnChainEvent(); event.setId(1L); From 3c5d8bc6803a48f6fb670fee7d214042857eec97 Mon Sep 17 00:00:00 2001 From: realvic22 Date: Mon, 27 Jul 2026 13:09:50 +0100 Subject: [PATCH 4/6] fix: dead-letter test seeds attempts=4 so one failure hits MAX_ATTEMPTS --- .../api/chain/ChainEventServiceIntegrationTest.java | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java index 025925e..904ae6a 100644 --- a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java +++ b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java @@ -128,18 +128,16 @@ void processOneTransitionsEventToProcessedAndOutboxToCompleted() { @Test void failureExhaustingRetriesMovesToDeadLetter() { - OnChainEvent event = saveEvent("dead", ChainEventStatus.PENDING, 0); + // Seed at attempts=4 so one more failure reaches MAX_ATTEMPTS (5). + // Backoff after earlier failures would otherwise push nextAttemptAt into + // the future and prevent claimNext from picking the event up again. + OnChainEvent event = saveEvent("dead", ChainEventStatus.PENDING, 4); saveOutbox(event.getId(), OutboxStatus.PENDING); doThrow(new RuntimeException("simulated processing failure")) .when(chainEventHandler).handle(any()); - for (int i = 0; i < 5; i++) { - try { - service.processOne(); - } catch (RuntimeException ignored) { - } - } + service.processOne(); OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); assertThat(reloaded.getStatus()).isEqualTo(ChainEventStatus.DEAD_LETTER); From d0a9b33bd199652dcdd1b93e824b449ca5c769b2 Mon Sep 17 00:00:00 2001 From: realvic22 Date: Mon, 27 Jul 2026 13:19:56 +0100 Subject: [PATCH 5/6] fix: harden chain event pipeline for reliable CI - Use REQUIRES_NEW inserter so Postgres unique-key races don't abort the outer TX - Explicitly save event/outbox on process + replay (no reliance on dirty-check alone) - Disable scheduling in integration tests; unique event keys; mock reset - Fix dead-letter test (seed attempts=4); non-flaky backoff assertion - Update unit tests for inserter-based constructor --- .../api/chain/service/ChainEventService.java | 129 +++++++++++++----- .../ChainEventServiceIntegrationTest.java | 125 +++++++---------- .../api/chain/ChainEventServiceTest.java | 69 +++++----- 3 files changed, 186 insertions(+), 137 deletions(-) diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java b/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java index 0040b3e..64c6146 100644 --- a/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java +++ b/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java @@ -4,82 +4,145 @@ import com.guildworkman.api.chain.api.*; import com.guildworkman.api.chain.model.*; import com.guildworkman.api.chain.repository.*; -import jakarta.transaction.Transactional; import lombok.RequiredArgsConstructor; import org.springframework.dao.DataIntegrityViolationException; import org.springframework.data.domain.PageRequest; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; -import java.time.*; -import java.util.*; +import java.time.Instant; +import java.util.EnumSet; +import java.util.List; -@Service @RequiredArgsConstructor +@Service +@RequiredArgsConstructor public class ChainEventService { private final OnChainEventRepository events; private final OutboxEventRepository outbox; private final ObjectMapper objectMapper; private final List handlers; - private static final int MAX_ATTEMPTS = 5; + private final ChainEventInserter inserter; - @Transactional + static final int MAX_ATTEMPTS = 5; + + @Transactional(readOnly = true) public ChainEventResponse ingest(IngestChainEventRequest request) { return events.findByEventKey(request.eventKey()) .map(ChainEventResponse::from) - .orElseGet(() -> createEvent(request)); + .orElseGet(() -> insertIdempotently(request)); } - private ChainEventResponse createEvent(IngestChainEventRequest request) { + private ChainEventResponse insertIdempotently(IngestChainEventRequest request) { try { - OnChainEvent event = new OnChainEvent(); - event.setEventKey(request.eventKey()); event.setContractId(request.contractId()); - event.setLedger(request.ledger()); event.setEventIndex(request.eventIndex()); - try { event.setTopics(objectMapper.writeValueAsString(request.topics())); } - catch (Exception ex) { throw new IllegalArgumentException("topics must be serializable", ex); } - event.setPayload(request.payload()); event.setStatus(ChainEventStatus.PENDING); - event.setNextAttemptAt(Instant.now()); - OnChainEvent saved = events.save(event); - OutboxEvent message = new OutboxEvent(); message.setEventId(saved.getId()); outbox.save(message); - return ChainEventResponse.from(saved); + return ChainEventResponse.from(inserter.insert(request)); } catch (DataIntegrityViolationException ex) { + // Nested REQUIRES_NEW insert rolled back; outer TX can still read the winner. return events.findByEventKey(request.eventKey()) .map(ChainEventResponse::from) - .orElseThrow(() -> new IllegalStateException("Event not found after idempotent-guard violation for key=" + request.eventKey(), ex)); - } catch (RuntimeException ex) { - return events.findByEventKey(request.eventKey()) - .map(ChainEventResponse::from) - .orElseThrow(() -> ex); + .orElseThrow(() -> new IllegalStateException( + "Event not found after idempotent-guard violation for key=" + request.eventKey(), ex)); } } @Transactional public int replay(ReplayRequest request) { int count = 0; - List batch = events.findByLedgerBetweenOrderByContractIdAscLedgerAscEventIndexAsc(request.fromLedger(), request.toLedger()); + List batch = events.findByLedgerBetweenOrderByContractIdAscLedgerAscEventIndexAsc( + request.fromLedger(), request.toLedger()); + Instant now = Instant.now(); for (OnChainEvent event : batch) { - event.setStatus(ChainEventStatus.PENDING); event.setAttempts(0); event.setLastError(null); event.setProcessedAt(null); event.setNextAttemptAt(Instant.now()); - outbox.findByEventId(event.getId()).ifPresent(message -> { message.setStatus(OutboxStatus.PENDING); message.setAttempts(0); message.setLastError(null); message.setCompletedAt(null); message.setNextAttemptAt(Instant.now()); }); + event.setStatus(ChainEventStatus.PENDING); + event.setAttempts(0); + event.setLastError(null); + event.setProcessedAt(null); + event.setNextAttemptAt(now); + events.save(event); + + outbox.findByEventId(event.getId()).ifPresent(message -> { + message.setStatus(OutboxStatus.PENDING); + message.setAttempts(0); + message.setLastError(null); + message.setCompletedAt(null); + message.setNextAttemptAt(now); + outbox.save(message); + }); count++; } - events.saveAll(batch); return count; } @Scheduled(fixedDelayString = "${chain.events.poll-delay-ms:1000}") @Transactional public void processOne() { - events.claimNext(EnumSet.of(ChainEventStatus.PENDING, ChainEventStatus.PROCESSING), Instant.now(), PageRequest.of(0, 1)).stream().findFirst().ifPresent(this::process); + events.claimNext( + EnumSet.of(ChainEventStatus.PENDING, ChainEventStatus.PROCESSING), + Instant.now(), + PageRequest.of(0, 1) + ).stream().findFirst().ifPresent(this::process); } void process(OnChainEvent event) { try { - event.setStatus(ChainEventStatus.PROCESSING); event.setAttempts(event.getAttempts() + 1); - for (ChainEventHandler handler : handlers) handler.handle(event); - event.setStatus(ChainEventStatus.PROCESSED); event.setProcessedAt(Instant.now()); event.setLastError(null); - outbox.findByEventId(event.getId()).ifPresent(message -> { message.setStatus(OutboxStatus.COMPLETED); message.setCompletedAt(Instant.now()); message.setLastError(null); }); + event.setStatus(ChainEventStatus.PROCESSING); + event.setAttempts(event.getAttempts() + 1); + for (ChainEventHandler handler : handlers) { + handler.handle(event); + } + event.setStatus(ChainEventStatus.PROCESSED); + event.setProcessedAt(Instant.now()); + event.setLastError(null); + events.save(event); + outbox.findByEventId(event.getId()).ifPresent(message -> { + message.setStatus(OutboxStatus.COMPLETED); + message.setCompletedAt(Instant.now()); + message.setLastError(null); + outbox.save(message); + }); } catch (RuntimeException ex) { event.setLastError(ex.getMessage()); - if (event.getAttempts() >= MAX_ATTEMPTS) event.setStatus(ChainEventStatus.DEAD_LETTER); else { event.setStatus(ChainEventStatus.PENDING); event.setNextAttemptAt(Instant.now().plusSeconds(1L << Math.min(event.getAttempts(), 6))); } + if (event.getAttempts() >= MAX_ATTEMPTS) { + event.setStatus(ChainEventStatus.DEAD_LETTER); + } else { + event.setStatus(ChainEventStatus.PENDING); + event.setNextAttemptAt(Instant.now().plusSeconds(1L << Math.min(event.getAttempts(), 6))); + } + events.save(event); + } + } + + /** + * Isolated insert so a unique-key race aborts only this nested transaction + * (Postgres), leaving the caller's transaction able to re-read the winner. + */ + @Service + @RequiredArgsConstructor + static class ChainEventInserter { + private final OnChainEventRepository events; + private final OutboxEventRepository outbox; + private final ObjectMapper objectMapper; + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public OnChainEvent insert(IngestChainEventRequest request) { + OnChainEvent event = new OnChainEvent(); + event.setEventKey(request.eventKey()); + event.setContractId(request.contractId()); + event.setLedger(request.ledger()); + event.setEventIndex(request.eventIndex()); + try { + event.setTopics(objectMapper.writeValueAsString(request.topics())); + } catch (Exception ex) { + throw new IllegalArgumentException("topics must be serializable", ex); + } + event.setPayload(request.payload()); + event.setStatus(ChainEventStatus.PENDING); + event.setNextAttemptAt(Instant.now()); + OnChainEvent saved = events.saveAndFlush(event); + OutboxEvent message = new OutboxEvent(); + message.setEventId(saved.getId()); + outbox.saveAndFlush(message); + return saved; } } } diff --git a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java index 904ae6a..17ca5e3 100644 --- a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java +++ b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java @@ -17,17 +17,25 @@ import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.test.mock.mockito.MockBean; +import java.time.Instant; +import java.util.ArrayList; import java.util.List; +import java.util.UUID; import java.util.concurrent.Callable; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.reset; -@SpringBootTest(properties = { "chain.events.poll-delay-ms=60000" }) +@SpringBootTest(properties = { + "chain.events.poll-delay-ms=3600000", + "spring.task.scheduling.enabled=false" +}) class ChainEventServiceIntegrationTest { @Autowired @@ -44,10 +52,15 @@ class ChainEventServiceIntegrationTest { @BeforeEach void cleanSlate() { + reset(chainEventHandler); outbox.deleteAll(); events.deleteAll(); } + private static String key(String prefix) { + return prefix + "-" + UUID.randomUUID(); + } + private OnChainEvent saveEvent(String eventKey, ChainEventStatus status, int attempts) { OnChainEvent e = new OnChainEvent(); e.setEventKey(eventKey); @@ -58,24 +71,21 @@ private OnChainEvent saveEvent(String eventKey, ChainEventStatus status, int att e.setPayload("{\"x\":1}"); e.setStatus(status); e.setAttempts(attempts); - e.setNextAttemptAt(java.time.Instant.now()); - return events.save(e); + e.setNextAttemptAt(Instant.now().minusSeconds(1)); + return events.saveAndFlush(e); } private OutboxEvent saveOutbox(Long eventId, OutboxStatus status) { OutboxEvent o = new OutboxEvent(); o.setEventId(eventId); o.setStatus(status); - return outbox.save(o); + o.setNextAttemptAt(Instant.now().minusSeconds(1)); + return outbox.saveAndFlush(o); } - // ------------------------------------------------------------------------- - // 1. Transactional outbox: ingest creates both rows atomically - // ------------------------------------------------------------------------- - @Test void ingestCreatesBothEventAndOutboxRow() { - IngestChainEventRequest req = new IngestChainEventRequest("k1", "C001", 1, 0, List.of("T"), "{}"); + IngestChainEventRequest req = new IngestChainEventRequest(key("k1"), "C001", 1, 0, List.of("T"), "{}"); ChainEventResponse resp = service.ingest(req); OnChainEvent saved = events.findById(resp.id()).orElseThrow(); @@ -86,13 +96,10 @@ void ingestCreatesBothEventAndOutboxRow() { assertThat(msg.getEventId()).isEqualTo(saved.getId()); } - // ------------------------------------------------------------------------- - // 2. Ingest idempotency - // ------------------------------------------------------------------------- - @Test void ingestIsIdempotentForDuplicateEventKey() { - IngestChainEventRequest req = new IngestChainEventRequest("dup", "C002", 2, 0, List.of("T"), "{}"); + String eventKey = key("dup"); + IngestChainEventRequest req = new IngestChainEventRequest(eventKey, "C002", 2, 0, List.of("T"), "{}"); ChainEventResponse first = service.ingest(req); ChainEventResponse second = service.ingest(req); @@ -101,13 +108,9 @@ void ingestIsIdempotentForDuplicateEventKey() { assertThat(outbox.count()).isEqualTo(1); } - // ------------------------------------------------------------------------- - // 3. Happy-path processing - // ------------------------------------------------------------------------- - @Test void processOneTransitionsEventToProcessedAndOutboxToCompleted() { - OnChainEvent event = saveEvent("happy", ChainEventStatus.PENDING, 0); + OnChainEvent event = saveEvent(key("happy"), ChainEventStatus.PENDING, 0); saveOutbox(event.getId(), OutboxStatus.PENDING); service.processOne(); @@ -122,16 +125,12 @@ void processOneTransitionsEventToProcessedAndOutboxToCompleted() { assertThat(msg.getCompletedAt()).isNotNull(); } - // ------------------------------------------------------------------------- - // 4. Failure path: handler exception -> DEAD_LETTER after MAX_ATTEMPTS - // ------------------------------------------------------------------------- - @Test void failureExhaustingRetriesMovesToDeadLetter() { - // Seed at attempts=4 so one more failure reaches MAX_ATTEMPTS (5). - // Backoff after earlier failures would otherwise push nextAttemptAt into - // the future and prevent claimNext from picking the event up again. - OnChainEvent event = saveEvent("dead", ChainEventStatus.PENDING, 4); + // Seed attempts=4 so one failure increments to 5 (== MAX_ATTEMPTS) → DEAD_LETTER. + // Do not loop processOne(): backoff after a failure makes nextAttemptAt future, + // so claimNext will not pick the event up again until that time. + OnChainEvent event = saveEvent(key("dead"), ChainEventStatus.PENDING, 4); saveOutbox(event.getId(), OutboxStatus.PENDING); doThrow(new RuntimeException("simulated processing failure")) @@ -145,43 +144,33 @@ void failureExhaustingRetriesMovesToDeadLetter() { assertThat(reloaded.getLastError()).isEqualTo("simulated processing failure"); } - // ------------------------------------------------------------------------- - // 5. Backoff after failure - // ------------------------------------------------------------------------- - @Test void failureAppliesExponentialBackoff() { - OnChainEvent event = saveEvent("backoff", ChainEventStatus.PENDING, 0); + OnChainEvent event = saveEvent(key("backoff"), ChainEventStatus.PENDING, 0); saveOutbox(event.getId(), OutboxStatus.PENDING); + Instant before = Instant.now(); doThrow(new RuntimeException("transient error")) .when(chainEventHandler).handle(any()); - try { - service.processOne(); - } catch (RuntimeException ignored) { - } + service.processOne(); OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); assertThat(reloaded.getStatus()).isEqualTo(ChainEventStatus.PENDING); assertThat(reloaded.getAttempts()).isEqualTo(1); - assertThat(reloaded.getNextAttemptAt()).isAfter(java.time.Instant.now()); + // attempts=1 → delay = 1<<1 = 2s + assertThat(reloaded.getNextAttemptAt()).isAfter(before.plusSeconds(1)); } - // ------------------------------------------------------------------------- - // 6. Replay persists changes - // ------------------------------------------------------------------------- - @Test void replayResetsEventAndOutboxStateAndPersists() { - OnChainEvent event = saveEvent("replay-me", ChainEventStatus.PROCESSED, 3); + OnChainEvent event = saveEvent(key("replay"), ChainEventStatus.PROCESSED, 3); event.setLastError("some error"); - events.save(event); + event.setProcessedAt(Instant.now()); + events.saveAndFlush(event); saveOutbox(event.getId(), OutboxStatus.COMPLETED); - ReplayRequest req = new ReplayRequest(100, 100); - int count = service.replay(req); - + int count = service.replay(new ReplayRequest(100, 100)); assertThat(count).isEqualTo(1); OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); @@ -197,13 +186,11 @@ void replayResetsEventAndOutboxStateAndPersists() { assertThat(msg.getCompletedAt()).isNull(); } - // ------------------------------------------------------------------------- - // 7. Replayed events can be reprocessed - // ------------------------------------------------------------------------- - @Test void replayedEventsCanBeReprocessed() { - OnChainEvent event = saveEvent("reprocess-me", ChainEventStatus.PROCESSED, 5); + OnChainEvent event = saveEvent(key("reprocess"), ChainEventStatus.PROCESSED, 5); + event.setProcessedAt(Instant.now()); + events.saveAndFlush(event); saveOutbox(event.getId(), OutboxStatus.COMPLETED); service.replay(new ReplayRequest(100, 100)); @@ -215,50 +202,44 @@ void replayedEventsCanBeReprocessed() { assertThat(reloaded.getProcessedAt()).isNotNull(); } - // ------------------------------------------------------------------------- - // 8. Concurrency: pessimistic locking prevents duplicate processing - // ------------------------------------------------------------------------- - @Test void pessimisticLockingPreventsDuplicateProcessing() throws Exception { - OnChainEvent event = saveEvent("concurrent", ChainEventStatus.PENDING, 0); + OnChainEvent event = saveEvent(key("concurrent"), ChainEventStatus.PENDING, 0); saveOutbox(event.getId(), OutboxStatus.PENDING); int threads = 8; ExecutorService pool = Executors.newFixedThreadPool(threads); + List> tasks = new ArrayList<>(); + for (int i = 0; i < threads; i++) { + tasks.add(() -> { + service.processOne(); + return null; + }); + } - List> tasks = java.util.stream.IntStream.range(0, threads) - .mapToObj(i -> (Callable) () -> { service.processOne(); return null; }) - .toList(); - - for (Future f : pool.invokeAll(tasks)) { - f.get(); + List> futures = pool.invokeAll(tasks, 30, TimeUnit.SECONDS); + for (Future f : futures) { + f.get(5, TimeUnit.SECONDS); } pool.shutdown(); + assertThat(pool.awaitTermination(10, TimeUnit.SECONDS)).isTrue(); OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); assertThat(reloaded.getStatus()).isEqualTo(ChainEventStatus.PROCESSED); assertThat(reloaded.getAttempts()).isEqualTo(1); + assertThat(events.countByStatus(ChainEventStatus.PROCESSED)).isEqualTo(1); } - // ------------------------------------------------------------------------- - // 9. Processed events are not claimed again - // ------------------------------------------------------------------------- - @Test void alreadyProcessedEventsAreNotClaimedAgain() { - OnChainEvent event = saveEvent("claimed", ChainEventStatus.PENDING, 0); + OnChainEvent event = saveEvent(key("claimed"), ChainEventStatus.PENDING, 0); saveOutbox(event.getId(), OutboxStatus.PENDING); + service.processOne(); service.processOne(); OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); assertThat(reloaded.getStatus()).isEqualTo(ChainEventStatus.PROCESSED); assertThat(reloaded.getAttempts()).isEqualTo(1); - - service.processOne(); - - reloaded = events.findById(event.getId()).orElseThrow(); - assertThat(reloaded.getAttempts()).isEqualTo(1); } } diff --git a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java index f2b5f3f..61da8d1 100644 --- a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java +++ b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java @@ -4,8 +4,10 @@ import com.guildworkman.api.chain.api.IngestChainEventRequest; import com.guildworkman.api.chain.model.ChainEventStatus; import com.guildworkman.api.chain.model.OnChainEvent; -import com.guildworkman.api.chain.repository.*; +import com.guildworkman.api.chain.repository.OnChainEventRepository; +import com.guildworkman.api.chain.repository.OutboxEventRepository; import com.guildworkman.api.chain.service.ChainEventService; +import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.mockito.Mockito; import org.springframework.dao.DataIntegrityViolationException; @@ -20,13 +22,24 @@ class ChainEventServiceTest { + private OnChainEventRepository events; + private OutboxEventRepository outbox; + private ChainEventService.ChainEventInserter inserter; + private ChainEventService service; + + @BeforeEach + void setUp() { + events = mock(OnChainEventRepository.class); + outbox = mock(OutboxEventRepository.class); + inserter = mock(ChainEventService.ChainEventInserter.class); + service = new ChainEventService(events, outbox, new ObjectMapper(), List.of(), inserter); + } + @Test void ingestIsIdempotentForTheSameEventKey() { - var events = Mockito.mock(OnChainEventRepository.class); - var outbox = Mockito.mock(OutboxEventRepository.class); - var service = new ChainEventService(events, outbox, new ObjectMapper(), List.of()); var request = new IngestChainEventRequest("evt-1", "CABC", 10, 0, List.of("Transfer"), "{\"amount\":1}"); when(events.findByEventKey("evt-1")).thenReturn(Optional.empty()); + var event = new OnChainEvent(); event.setId(1L); event.setEventKey("evt-1"); @@ -34,17 +47,14 @@ void ingestIsIdempotentForTheSameEventKey() { event.setStatus(ChainEventStatus.PENDING); event.setTopics("[\"Transfer\"]"); event.setPayload(request.payload()); - when(events.save(any())).thenReturn(event); + when(inserter.insert(request)).thenReturn(event); + assertThat(service.ingest(request).eventKey()).isEqualTo("evt-1"); - verify(events).save(any()); - verify(outbox).save(any()); + verify(inserter).insert(request); } @Test void ingestReturnsExistingEventOnDuplicateEventKey() { - var events = Mockito.mock(OnChainEventRepository.class); - var outbox = Mockito.mock(OutboxEventRepository.class); - var service = new ChainEventService(events, outbox, new ObjectMapper(), List.of()); var request = new IngestChainEventRequest("evt-dup", "CABC", 10, 0, List.of("Transfer"), "{}"); var existing = new OnChainEvent(); existing.setId(1L); @@ -54,21 +64,17 @@ void ingestReturnsExistingEventOnDuplicateEventKey() { existing.setTopics("[\"Transfer\"]"); existing.setPayload("{}"); when(events.findByEventKey("evt-dup")).thenReturn(Optional.of(existing)); + var resp = service.ingest(request); assertThat(resp.id()).isEqualTo(1L); - verify(events, never()).save(any()); - verify(outbox, never()).save(any()); + verify(inserter, never()).insert(any()); } @Test void ingestHandlesDataIntegrityViolationWithFallbackLookup() { - var events = Mockito.mock(OnChainEventRepository.class); - var outbox = Mockito.mock(OutboxEventRepository.class); - var service = new ChainEventService(events, outbox, new ObjectMapper(), List.of()); var request = new IngestChainEventRequest("race", "CABC", 10, 0, List.of("T"), "{}"); - when(events.findByEventKey("race")).thenReturn(Optional.empty()); - when(events.save(any())).thenThrow(new DataIntegrityViolationException("dup key")); + when(inserter.insert(request)).thenThrow(new DataIntegrityViolationException("dup key")); var afterSave = new OnChainEvent(); afterSave.setId(42L); @@ -77,8 +83,6 @@ void ingestHandlesDataIntegrityViolationWithFallbackLookup() { afterSave.setStatus(ChainEventStatus.PENDING); afterSave.setTopics("[\"T\"]"); afterSave.setPayload("{}"); - - // After the save fails, createEvent calls findByIdempotentKey again when(events.findByEventKey("race")).thenReturn(Optional.empty(), Optional.of(afterSave)); var resp = service.ingest(request); @@ -88,13 +92,9 @@ void ingestHandlesDataIntegrityViolationWithFallbackLookup() { @Test void ingestPropagatesUnexpectedExceptionWhenNoEventFound() { - var events = Mockito.mock(OnChainEventRepository.class); - var outbox = Mockito.mock(OutboxEventRepository.class); - var service = new ChainEventService(events, outbox, new ObjectMapper(), List.of()); var request = new IngestChainEventRequest("boom", "CABC", 10, 0, List.of("T"), "{}"); - when(events.findByEventKey("boom")).thenReturn(Optional.empty()); - when(events.save(any())).thenThrow(new RuntimeException("db connection lost")); + when(inserter.insert(request)).thenThrow(new RuntimeException("db connection lost")); assertThatThrownBy(() -> service.ingest(request)) .isInstanceOf(RuntimeException.class) @@ -102,11 +102,7 @@ void ingestPropagatesUnexpectedExceptionWhenNoEventFound() { } @Test - void replayPersistsChangesExplicitly() { - var events = Mockito.mock(OnChainEventRepository.class); - var outbox = Mockito.mock(OutboxEventRepository.class); - var service = new ChainEventService(events, outbox, new ObjectMapper(), List.of()); - + void replayPersistsEventAndOutboxChanges() { var event = new OnChainEvent(); event.setId(1L); event.setEventKey("r"); @@ -116,13 +112,22 @@ void replayPersistsChangesExplicitly() { event.setProcessedAt(java.time.Instant.now()); event.setLedger(10); + var message = new com.guildworkman.api.chain.model.OutboxEvent(); + message.setId(9L); + message.setEventId(1L); + message.setStatus(com.guildworkman.api.chain.model.OutboxStatus.COMPLETED); + when(events.findByLedgerBetweenOrderByContractIdAscLedgerAscEventIndexAsc(5, 15)) .thenReturn(List.of(event)); + when(outbox.findByEventId(1L)).thenReturn(Optional.of(message)); - var replay = new com.guildworkman.api.chain.api.ReplayRequest(5, 15); - int count = service.replay(replay); + int count = service.replay(new com.guildworkman.api.chain.api.ReplayRequest(5, 15)); assertThat(count).isEqualTo(1); - verify(events).saveAll(any()); + verify(events).save(event); + verify(outbox).save(message); + assertThat(event.getStatus()).isEqualTo(ChainEventStatus.PENDING); + assertThat(event.getAttempts()).isZero(); + assertThat(message.getStatus()).isEqualTo(com.guildworkman.api.chain.model.OutboxStatus.PENDING); } } From 6ff8bc8cf22dabaa4d7dc4d412b13b891bc91bc8 Mon Sep 17 00:00:00 2001 From: realvic22 Date: Mon, 27 Jul 2026 13:23:04 +0100 Subject: [PATCH 6/6] fix: extract ChainEventInserter top-level bean for reliable component scan Also drop unused ObjectMapper from ChainEventService constructor. --- .../api/chain/service/ChainEventInserter.java | 49 +++++++++++++++++++ .../api/chain/service/ChainEventService.java | 37 -------------- .../api/chain/ChainEventServiceTest.java | 21 ++++---- 3 files changed, 60 insertions(+), 47 deletions(-) create mode 100644 backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventInserter.java diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventInserter.java b/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventInserter.java new file mode 100644 index 0000000..6891339 --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventInserter.java @@ -0,0 +1,49 @@ +package com.guildworkman.api.chain.service; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.guildworkman.api.chain.api.IngestChainEventRequest; +import com.guildworkman.api.chain.model.ChainEventStatus; +import com.guildworkman.api.chain.model.OnChainEvent; +import com.guildworkman.api.chain.model.OutboxEvent; +import com.guildworkman.api.chain.repository.OnChainEventRepository; +import com.guildworkman.api.chain.repository.OutboxEventRepository; +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; + +import java.time.Instant; + +/** + * Isolated insert so a unique-key race aborts only this nested transaction + * (Postgres), leaving the caller's transaction able to re-read the winner. + */ +@Service +@RequiredArgsConstructor +public class ChainEventInserter { + private final OnChainEventRepository events; + private final OutboxEventRepository outbox; + private final ObjectMapper objectMapper; + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public OnChainEvent insert(IngestChainEventRequest request) { + OnChainEvent event = new OnChainEvent(); + event.setEventKey(request.eventKey()); + event.setContractId(request.contractId()); + event.setLedger(request.ledger()); + event.setEventIndex(request.eventIndex()); + try { + event.setTopics(objectMapper.writeValueAsString(request.topics())); + } catch (Exception ex) { + throw new IllegalArgumentException("topics must be serializable", ex); + } + event.setPayload(request.payload()); + event.setStatus(ChainEventStatus.PENDING); + event.setNextAttemptAt(Instant.now()); + OnChainEvent saved = events.saveAndFlush(event); + OutboxEvent message = new OutboxEvent(); + message.setEventId(saved.getId()); + outbox.saveAndFlush(message); + return saved; + } +} diff --git a/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java b/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java index 64c6146..eb69057 100644 --- a/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java +++ b/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java @@ -1,6 +1,5 @@ package com.guildworkman.api.chain.service; -import com.fasterxml.jackson.databind.ObjectMapper; import com.guildworkman.api.chain.api.*; import com.guildworkman.api.chain.model.*; import com.guildworkman.api.chain.repository.*; @@ -9,7 +8,6 @@ import org.springframework.data.domain.PageRequest; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; -import org.springframework.transaction.annotation.Propagation; import org.springframework.transaction.annotation.Transactional; import java.time.Instant; @@ -21,7 +19,6 @@ public class ChainEventService { private final OnChainEventRepository events; private final OutboxEventRepository outbox; - private final ObjectMapper objectMapper; private final List handlers; private final ChainEventInserter inserter; @@ -111,38 +108,4 @@ void process(OnChainEvent event) { events.save(event); } } - - /** - * Isolated insert so a unique-key race aborts only this nested transaction - * (Postgres), leaving the caller's transaction able to re-read the winner. - */ - @Service - @RequiredArgsConstructor - static class ChainEventInserter { - private final OnChainEventRepository events; - private final OutboxEventRepository outbox; - private final ObjectMapper objectMapper; - - @Transactional(propagation = Propagation.REQUIRES_NEW) - public OnChainEvent insert(IngestChainEventRequest request) { - OnChainEvent event = new OnChainEvent(); - event.setEventKey(request.eventKey()); - event.setContractId(request.contractId()); - event.setLedger(request.ledger()); - event.setEventIndex(request.eventIndex()); - try { - event.setTopics(objectMapper.writeValueAsString(request.topics())); - } catch (Exception ex) { - throw new IllegalArgumentException("topics must be serializable", ex); - } - event.setPayload(request.payload()); - event.setStatus(ChainEventStatus.PENDING); - event.setNextAttemptAt(Instant.now()); - OnChainEvent saved = events.saveAndFlush(event); - OutboxEvent message = new OutboxEvent(); - message.setEventId(saved.getId()); - outbox.saveAndFlush(message); - return saved; - } - } } diff --git a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java index 61da8d1..9751e2f 100644 --- a/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java +++ b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java @@ -1,15 +1,17 @@ package com.guildworkman.api.chain; -import com.fasterxml.jackson.databind.ObjectMapper; import com.guildworkman.api.chain.api.IngestChainEventRequest; +import com.guildworkman.api.chain.api.ReplayRequest; import com.guildworkman.api.chain.model.ChainEventStatus; import com.guildworkman.api.chain.model.OnChainEvent; +import com.guildworkman.api.chain.model.OutboxEvent; +import com.guildworkman.api.chain.model.OutboxStatus; import com.guildworkman.api.chain.repository.OnChainEventRepository; import com.guildworkman.api.chain.repository.OutboxEventRepository; +import com.guildworkman.api.chain.service.ChainEventInserter; import com.guildworkman.api.chain.service.ChainEventService; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; -import org.mockito.Mockito; import org.springframework.dao.DataIntegrityViolationException; import java.util.List; @@ -24,15 +26,15 @@ class ChainEventServiceTest { private OnChainEventRepository events; private OutboxEventRepository outbox; - private ChainEventService.ChainEventInserter inserter; + private ChainEventInserter inserter; private ChainEventService service; @BeforeEach void setUp() { events = mock(OnChainEventRepository.class); outbox = mock(OutboxEventRepository.class); - inserter = mock(ChainEventService.ChainEventInserter.class); - service = new ChainEventService(events, outbox, new ObjectMapper(), List.of(), inserter); + inserter = mock(ChainEventInserter.class); + service = new ChainEventService(events, outbox, List.of(), inserter); } @Test @@ -73,7 +75,6 @@ void ingestReturnsExistingEventOnDuplicateEventKey() { @Test void ingestHandlesDataIntegrityViolationWithFallbackLookup() { var request = new IngestChainEventRequest("race", "CABC", 10, 0, List.of("T"), "{}"); - when(events.findByEventKey("race")).thenReturn(Optional.empty()); when(inserter.insert(request)).thenThrow(new DataIntegrityViolationException("dup key")); var afterSave = new OnChainEvent(); @@ -112,22 +113,22 @@ void replayPersistsEventAndOutboxChanges() { event.setProcessedAt(java.time.Instant.now()); event.setLedger(10); - var message = new com.guildworkman.api.chain.model.OutboxEvent(); + var message = new OutboxEvent(); message.setId(9L); message.setEventId(1L); - message.setStatus(com.guildworkman.api.chain.model.OutboxStatus.COMPLETED); + message.setStatus(OutboxStatus.COMPLETED); when(events.findByLedgerBetweenOrderByContractIdAscLedgerAscEventIndexAsc(5, 15)) .thenReturn(List.of(event)); when(outbox.findByEventId(1L)).thenReturn(Optional.of(message)); - int count = service.replay(new com.guildworkman.api.chain.api.ReplayRequest(5, 15)); + int count = service.replay(new ReplayRequest(5, 15)); assertThat(count).isEqualTo(1); verify(events).save(event); verify(outbox).save(message); assertThat(event.getStatus()).isEqualTo(ChainEventStatus.PENDING); assertThat(event.getAttempts()).isZero(); - assertThat(message.getStatus()).isEqualTo(com.guildworkman.api.chain.model.OutboxStatus.PENDING); + assertThat(message.getStatus()).isEqualTo(OutboxStatus.PENDING); } }