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/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/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 new file mode 100644 index 0000000..eb69057 --- /dev/null +++ b/backend-api/src/main/java/com/guildworkman/api/chain/service/ChainEventService.java @@ -0,0 +1,111 @@ +package com.guildworkman.api.chain.service; + +import com.guildworkman.api.chain.api.*; +import com.guildworkman.api.chain.model.*; +import com.guildworkman.api.chain.repository.*; +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.Transactional; + +import java.time.Instant; +import java.util.EnumSet; +import java.util.List; + +@Service +@RequiredArgsConstructor +public class ChainEventService { + private final OnChainEventRepository events; + private final OutboxEventRepository outbox; + private final List handlers; + private final ChainEventInserter inserter; + + static final int MAX_ATTEMPTS = 5; + + @Transactional(readOnly = true) + public ChainEventResponse ingest(IngestChainEventRequest request) { + return events.findByEventKey(request.eventKey()) + .map(ChainEventResponse::from) + .orElseGet(() -> insertIdempotently(request)); + } + + private ChainEventResponse insertIdempotently(IngestChainEventRequest request) { + try { + 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)); + } + } + + @Transactional + public int replay(ReplayRequest request) { + int count = 0; + 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(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++; + } + 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); + } + + 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); + 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))); + } + events.save(event); + } + } +} 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/ChainEventServiceIntegrationTest.java b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java new file mode 100644 index 0000000..17ca5e3 --- /dev/null +++ b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceIntegrationTest.java @@ -0,0 +1,245 @@ +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.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=3600000", + "spring.task.scheduling.enabled=false" +}) +class ChainEventServiceIntegrationTest { + + @Autowired + private OnChainEventRepository events; + + @Autowired + private OutboxEventRepository outbox; + + @Autowired + private ChainEventService service; + + @MockBean + private ChainEventHandler chainEventHandler; + + @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); + e.setContractId("CTEST"); + e.setLedger(100); + e.setEventIndex(0); + e.setTopics("[\"Test\"]"); + e.setPayload("{\"x\":1}"); + e.setStatus(status); + e.setAttempts(attempts); + 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); + o.setNextAttemptAt(Instant.now().minusSeconds(1)); + return outbox.saveAndFlush(o); + } + + @Test + void ingestCreatesBothEventAndOutboxRow() { + IngestChainEventRequest req = new IngestChainEventRequest(key("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()); + } + + @Test + void ingestIsIdempotentForDuplicateEventKey() { + 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); + + assertThat(second.id()).isEqualTo(first.id()); + assertThat(events.count()).isEqualTo(1); + assertThat(outbox.count()).isEqualTo(1); + } + + @Test + void processOneTransitionsEventToProcessedAndOutboxToCompleted() { + OnChainEvent event = saveEvent(key("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(); + } + + @Test + void failureExhaustingRetriesMovesToDeadLetter() { + // 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")) + .when(chainEventHandler).handle(any()); + + service.processOne(); + + 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"); + } + + @Test + void failureAppliesExponentialBackoff() { + 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()); + + service.processOne(); + + OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); + assertThat(reloaded.getStatus()).isEqualTo(ChainEventStatus.PENDING); + assertThat(reloaded.getAttempts()).isEqualTo(1); + // attempts=1 → delay = 1<<1 = 2s + assertThat(reloaded.getNextAttemptAt()).isAfter(before.plusSeconds(1)); + } + + @Test + void replayResetsEventAndOutboxStateAndPersists() { + OnChainEvent event = saveEvent(key("replay"), ChainEventStatus.PROCESSED, 3); + event.setLastError("some error"); + event.setProcessedAt(Instant.now()); + events.saveAndFlush(event); + saveOutbox(event.getId(), OutboxStatus.COMPLETED); + + int count = service.replay(new ReplayRequest(100, 100)); + 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(); + } + + @Test + void replayedEventsCanBeReprocessed() { + 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)); + service.processOne(); + + OnChainEvent reloaded = events.findById(event.getId()).orElseThrow(); + assertThat(reloaded.getStatus()).isEqualTo(ChainEventStatus.PROCESSED); + assertThat(reloaded.getAttempts()).isEqualTo(1); + assertThat(reloaded.getProcessedAt()).isNotNull(); + } + + @Test + void pessimisticLockingPreventsDuplicateProcessing() throws Exception { + 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> 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); + } + + @Test + void alreadyProcessedEventsAreNotClaimedAgain() { + 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); + } +} 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..9751e2f --- /dev/null +++ b/backend-api/src/test/java/com/guildworkman/api/chain/ChainEventServiceTest.java @@ -0,0 +1,134 @@ +package com.guildworkman.api.chain; + +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.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 { + + private OnChainEventRepository events; + private OutboxEventRepository outbox; + private ChainEventInserter inserter; + private ChainEventService service; + + @BeforeEach + void setUp() { + events = mock(OnChainEventRepository.class); + outbox = mock(OutboxEventRepository.class); + inserter = mock(ChainEventInserter.class); + service = new ChainEventService(events, outbox, List.of(), inserter); + } + + @Test + void ingestIsIdempotentForTheSameEventKey() { + 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"); + event.setContractId("CABC"); + event.setStatus(ChainEventStatus.PENDING); + event.setTopics("[\"Transfer\"]"); + event.setPayload(request.payload()); + when(inserter.insert(request)).thenReturn(event); + + assertThat(service.ingest(request).eventKey()).isEqualTo("evt-1"); + verify(inserter).insert(request); + } + + @Test + void ingestReturnsExistingEventOnDuplicateEventKey() { + 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(inserter, never()).insert(any()); + } + + @Test + void ingestHandlesDataIntegrityViolationWithFallbackLookup() { + var request = new IngestChainEventRequest("race", "CABC", 10, 0, List.of("T"), "{}"); + when(inserter.insert(request)).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("{}"); + 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 request = new IngestChainEventRequest("boom", "CABC", 10, 0, List.of("T"), "{}"); + when(events.findByEventKey("boom")).thenReturn(Optional.empty()); + when(inserter.insert(request)).thenThrow(new RuntimeException("db connection lost")); + + assertThatThrownBy(() -> service.ingest(request)) + .isInstanceOf(RuntimeException.class) + .hasMessage("db connection lost"); + } + + @Test + void replayPersistsEventAndOutboxChanges() { + 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); + + var message = new OutboxEvent(); + message.setId(9L); + message.setEventId(1L); + 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 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(OutboxStatus.PENDING); + } +}