diff --git a/build.gradle b/build.gradle index 0d06b3d..87675fa 100644 --- a/build.gradle +++ b/build.gradle @@ -1,6 +1,6 @@ plugins { id 'java' - id 'org.springframework.boot' version '3.3.18' + id 'org.springframework.boot' version '3.5.14' id 'io.spring.dependency-management' version '1.1.7' } @@ -23,6 +23,9 @@ dependencies { implementation 'org.springframework.boot:spring-boot-starter-data-jpa' implementation 'org.springframework.boot:spring-boot-starter-data-redis' implementation 'org.springframework.boot:spring-boot-starter-actuator' + implementation 'org.springframework.kafka:spring-kafka' + implementation 'org.springframework.boot:spring-boot-starter-web' + implementation 'org.redisson:redisson:3.27.2' implementation 'io.github.cdimascio:dotenv-java:3.0.0' compileOnly 'org.projectlombok:lombok' annotationProcessor 'org.projectlombok:lombok' diff --git a/gradle/wrapper/gradle-wrapper.properties b/gradle/wrapper/gradle-wrapper.properties index a441313..aaaabb3 100644 --- a/gradle/wrapper/gradle-wrapper.properties +++ b/gradle/wrapper/gradle-wrapper.properties @@ -1,6 +1,6 @@ distributionBase=GRADLE_USER_HOME distributionPath=wrapper/dists -distributionUrl=https\://services.gradle.org/distributions/gradle-8.8-bin.zip +distributionUrl=https\://services.gradle.org/distributions/gradle-8.14.4-bin.zip networkTimeout=10000 validateDistributionUrl=true zipStoreBase=GRADLE_USER_HOME diff --git a/src/main/java/com/rocketcrew/pocatbatch/client/MainAppBuyoutClient.java b/src/main/java/com/rocketcrew/pocatbatch/client/MainAppBuyoutClient.java new file mode 100644 index 0000000..1b91733 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/client/MainAppBuyoutClient.java @@ -0,0 +1,67 @@ +package com.rocketcrew.pocatbatch.client; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.http.HttpEntity; +import org.springframework.http.HttpHeaders; +import org.springframework.http.HttpMethod; +import org.springframework.stereotype.Component; +import org.springframework.web.client.HttpClientErrorException; +import org.springframework.web.client.RestTemplate; + +@Component +@RequiredArgsConstructor +@Slf4j +public class MainAppBuyoutClient { + + private final RestTemplate restTemplate; + + @Value("${pocat.main-app.base-url}") + private String baseUrl; + + @Value("${pocat.main-app.internal-token}") + private String internalToken; + + /** + * 경매 구매 확정 실패 복구 요청 + * POST {baseUrl}/internal/auctions/{id}/recover-buyout + * 최대 3회 재시도 (지수 백오프) + * 4xx 오류는 SkipException 발생 + */ + public void recoverBuyout(Long auctionId, Long jobExecutionId) { + String url = String.format("%s/internal/auctions/%d/recover-buyout", baseUrl, auctionId); + + HttpHeaders headers = new HttpHeaders(); + headers.set("X-Internal-Token", internalToken); + headers.set("Idempotency-Key", String.format("buyout-%d-%d", auctionId, jobExecutionId)); + + HttpEntity request = new HttpEntity<>(headers); + + int maxRetries = 3; + long delayMs = 1000; + + for (int attempt = 1; attempt <= maxRetries; attempt++) { + try { + restTemplate.exchange(url, HttpMethod.POST, request, String.class); + log.info("경매 구매 확정 복구 성공: auctionId={}", auctionId); + return; + } catch (HttpClientErrorException e) { + log.warn("경매 구매 확정 복구 4xx 오류: auctionId={}, status={}", auctionId, e.getStatusCode()); + throw new RuntimeException("4xx 오류로 스킵", e); + } catch (Exception e) { + if (attempt == maxRetries) { + log.error("경매 구매 확정 복구 실패: auctionId={}", auctionId, e); + throw e; + } + try { + Thread.sleep(delayMs); + delayMs *= 2; + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + throw new RuntimeException(ie); + } + } + } + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/client/MainAppRefundClient.java b/src/main/java/com/rocketcrew/pocatbatch/client/MainAppRefundClient.java new file mode 100644 index 0000000..cffc071 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/client/MainAppRefundClient.java @@ -0,0 +1,67 @@ +package com.rocketcrew.pocatbatch.client; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.http.HttpEntity; +import org.springframework.http.HttpHeaders; +import org.springframework.http.HttpMethod; +import org.springframework.stereotype.Component; +import org.springframework.web.client.HttpClientErrorException; +import org.springframework.web.client.RestTemplate; + +@Component +@RequiredArgsConstructor +@Slf4j +public class MainAppRefundClient { + + private final RestTemplate restTemplate; + + @Value("${pocat.main-app.base-url}") + private String baseUrl; + + @Value("${pocat.main-app.internal-token}") + private String internalToken; + + /** + * 환불 재시도 요청 + * POST {baseUrl}/internal/refunds/{id}/retry + * 최대 3회 재시도 (지수 백오프) + * 4xx 오류는 SkipException 발생 + */ + public void retryRefund(Long refundId, Long jobExecutionId) { + String url = String.format("%s/internal/refunds/%d/retry", baseUrl, refundId); + + HttpHeaders headers = new HttpHeaders(); + headers.set("X-Internal-Token", internalToken); + headers.set("Idempotency-Key", String.format("refund-%d-%d", refundId, jobExecutionId)); + + HttpEntity request = new HttpEntity<>(headers); + + int maxRetries = 3; + long delayMs = 1000; + + for (int attempt = 1; attempt <= maxRetries; attempt++) { + try { + restTemplate.exchange(url, HttpMethod.POST, request, String.class); + log.info("환불 재시도 성공: refundId={}", refundId); + return; + } catch (HttpClientErrorException e) { + log.warn("환불 재시도 4xx 오류: refundId={}, status={}", refundId, e.getStatusCode()); + throw new RuntimeException("4xx 오류로 스킵", e); + } catch (Exception e) { + if (attempt == maxRetries) { + log.error("환불 재시도 실패: refundId={}", refundId, e); + throw e; + } + try { + Thread.sleep(delayMs); + delayMs *= 2; + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + throw new RuntimeException(ie); + } + } + } + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/config/KafkaProducerConfig.java b/src/main/java/com/rocketcrew/pocatbatch/config/KafkaProducerConfig.java new file mode 100644 index 0000000..d188dce --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/config/KafkaProducerConfig.java @@ -0,0 +1,49 @@ +package com.rocketcrew.pocatbatch.config; + +import lombok.RequiredArgsConstructor; +import org.springframework.boot.autoconfigure.kafka.KafkaProperties; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Profile; +import org.springframework.kafka.core.DefaultKafkaProducerFactory; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.ProducerFactory; + +import java.util.HashMap; +import java.util.Map; +import java.util.Set; + +@Configuration +@Profile("!test") +@RequiredArgsConstructor +public class KafkaProducerConfig { + + private final KafkaProperties kafkaProperties; + + public static final Set FINANCIAL_TOPICS = Set.of("payment", "refund", "settlement"); + + /** + * 기본 KafkaTemplate (acks=1) + * 일반 이벤트 발행용 + */ + @Bean + public KafkaTemplate kafkaTemplate() { + Map props = new HashMap<>(kafkaProperties.buildProducerProperties(null)); + props.put("acks", "1"); + ProducerFactory factory = new DefaultKafkaProducerFactory<>(props); + return new KafkaTemplate<>(factory); + } + + /** + * 금융 관련 KafkaTemplate (acks=all, enable.idempotence=true) + * 환불, 결제, 정산 이벤트 발행용 (중복 방지 + 높은 안정성) + */ + @Bean + public KafkaTemplate paymentKafkaTemplate() { + Map props = new HashMap<>(kafkaProperties.buildProducerProperties(null)); + props.put("acks", "all"); + props.put("enable.idempotence", true); + ProducerFactory factory = new DefaultKafkaProducerFactory<>(props); + return new KafkaTemplate<>(factory); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/config/RedissonConfig.java b/src/main/java/com/rocketcrew/pocatbatch/config/RedissonConfig.java new file mode 100644 index 0000000..78e2c27 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/config/RedissonConfig.java @@ -0,0 +1,39 @@ +package com.rocketcrew.pocatbatch.config; + +import org.redisson.Redisson; +import org.redisson.api.RedissonClient; +import org.redisson.config.Config; +import org.redisson.config.SingleServerConfig; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Profile; + +@Profile("!test") +@Configuration +public class RedissonConfig { + + @Value("${spring.data.redis.host:localhost}") + private String redisHost; + + @Value("${spring.data.redis.port:6379}") + private int redisPort; + + @Value("${spring.data.redis.password:}") + private String redisPassword; + + /** + * Redisson Client Bean + * 경매 락(auction activation, expiration) 및 일반적인 Redis 분산 락 사용 + */ + @Bean + public RedissonClient redissonClient() { + Config config = new Config(); + SingleServerConfig serverConfig = config.useSingleServer() + .setAddress(String.format("redis://%s:%d", redisHost, redisPort)); + if (redisPassword != null && !redisPassword.isEmpty()) { + serverConfig.setPassword(redisPassword); + } + return Redisson.create(config); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/config/RestClientConfig.java b/src/main/java/com/rocketcrew/pocatbatch/config/RestClientConfig.java new file mode 100644 index 0000000..66ba8c7 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/config/RestClientConfig.java @@ -0,0 +1,26 @@ +package com.rocketcrew.pocatbatch.config; + +import org.springframework.boot.web.client.RestTemplateBuilder; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Profile; +import org.springframework.web.client.RestTemplate; + +import java.time.Duration; + +@Configuration +@Profile("!test") +public class RestClientConfig { + + /** + * RestTemplate Bean + * Main App 내부 API 호출용 (buyout recovery, refund retry) + */ + @Bean + public RestTemplate restTemplate(RestTemplateBuilder builder) { + return builder + .connectTimeout(Duration.ofSeconds(5)) + .readTimeout(Duration.ofSeconds(10)) + .build(); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/ai/entity/AiChatSession.java b/src/main/java/com/rocketcrew/pocatbatch/domain/ai/entity/AiChatSession.java new file mode 100644 index 0000000..0ebb7ae --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/ai/entity/AiChatSession.java @@ -0,0 +1,67 @@ +package com.rocketcrew.pocatbatch.domain.ai.entity; + +import com.rocketcrew.pocatbatch.domain.freepost.entity.BaseEntity; +import jakarta.persistence.*; +import lombok.*; +import org.hibernate.annotations.SQLDelete; +import org.hibernate.annotations.SQLRestriction; + +import java.time.LocalDateTime; + +/** + * AI 채팅 세션 엔티티. + * 사용자별 멀티턴 대화 세션 관리. + */ +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +@Builder +@AllArgsConstructor +@Entity +@Table(name = "ai_chat_sessions", + indexes = { + @Index(name = "idx_ai_chat_session_user_id", columnList = "user_id"), + @Index(name = "idx_ai_chat_session_uuid", columnList = "session_uuid"), + @Index(name = "idx_ai_chat_session_last_active", columnList = "last_active_at") + } +) +@SQLDelete(sql = "UPDATE ai_chat_sessions SET deleted_at = NOW() WHERE id = ?") +@SQLRestriction("deleted_at IS NULL") +public class AiChatSession extends BaseEntity { + + @Column(name = "user_id", nullable = false) + private Long userId; + + @Column(name = "session_uuid", nullable = false, length = 36, unique = true) + private String sessionUuid; + + @Column(name = "total_tokens", nullable = false) + private Integer totalTokens; + + @Column(name = "is_expired", nullable = false) + private Boolean isExpired; + + @Column(name = "last_active_at", nullable = false) + private LocalDateTime lastActiveAt; + + /** + * 토큰 추가. + */ + public void addTokens(int tokens) { + this.totalTokens += tokens; + } + + /** + * 마지막 활동 시간 업데이트. + */ + public void updateLastActiveAt(LocalDateTime now) { + this.lastActiveAt = now; + } + + /** + * 만료된 세션 재활성화. + */ + public void reactivate(LocalDateTime now) { + this.isExpired = false; + this.lastActiveAt = now; + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/ai/repository/AiChatSessionRepository.java b/src/main/java/com/rocketcrew/pocatbatch/domain/ai/repository/AiChatSessionRepository.java new file mode 100644 index 0000000..3a5a01d --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/ai/repository/AiChatSessionRepository.java @@ -0,0 +1,24 @@ +package com.rocketcrew.pocatbatch.domain.ai.repository; + +import com.rocketcrew.pocatbatch.domain.ai.entity.AiChatSession; +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Modifying; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; +import org.springframework.stereotype.Repository; + +import java.time.LocalDateTime; + +@Repository +public interface AiChatSessionRepository extends JpaRepository { + + /** + * 지정된 시간 이전의 비활성 세션을 만료 처리. + * + * @param threshold 기준 시간 + * @return 만료된 세션 수 + */ + @Modifying + @Query("UPDATE AiChatSession s SET s.isExpired = true WHERE s.lastActiveAt < :threshold AND s.isExpired = false") + int expireSessionsBeforeTime(@Param("threshold") LocalDateTime threshold); +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/auction/entity/Auction.java b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/entity/Auction.java new file mode 100644 index 0000000..85f8c25 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/entity/Auction.java @@ -0,0 +1,81 @@ +package com.rocketcrew.pocatbatch.domain.auction.entity; + +import com.rocketcrew.pocatbatch.domain.auction.enums.AuctionStatus; +import com.rocketcrew.pocatbatch.domain.freepost.entity.BaseEntity; +import jakarta.persistence.*; +import lombok.*; +import org.hibernate.annotations.SQLDelete; +import org.hibernate.annotations.SQLRestriction; + +import java.time.LocalDateTime; + +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +@AllArgsConstructor +@Builder +@Entity +@Table(name = "auctions") +@SQLDelete(sql = "UPDATE auctions SET deleted_at = NOW() WHERE id = ?") +@SQLRestriction("deleted_at IS NULL") +public class Auction extends BaseEntity { + + @Column(name = "card_id", nullable = false) + private Long cardId; + + @Column(name = "seller_id", nullable = false) + private Long sellerId; + + @Column(name = "highest_bidder_id") + private Long highestBidderId; + + @Column(name = "title", nullable = false) + private String title; + + @Column(name = "description", columnDefinition = "TEXT") + private String description; + + @Column(name = "starting_price", nullable = false) + private Long startingPrice; + + @Column(name = "buyout_price") + private Long buyoutPrice; + + @Column(name = "highest_price") + private Long highestPrice; + + @Enumerated(EnumType.STRING) + @Column(name = "status", nullable = false, length = 30) + private AuctionStatus status; + + @Column(name = "started_at") + private LocalDateTime startedAt; + + @Column(name = "ended_at") + private LocalDateTime endedAt; + + @Column(name = "reason", columnDefinition = "TEXT") + private String reason; + + @Column(name = "inspected_at") + private LocalDateTime inspectedAt; + + @Column(name = "inspected_by") + private Long inspectedBy; + + /** + * 승인된 경매를 실제 진행 상태로 전환하고 시작/종료 시각을 확정한다. + */ + public void activate(LocalDateTime startedAt, LocalDateTime endedAt) { + this.status = AuctionStatus.ACTIVE; + this.startedAt = startedAt; + this.endedAt = endedAt; + this.reason = null; + } + + /** + * 진행 중인 경매를 정상 종료 상태로 전환한다. + */ + public void end() { + this.status = AuctionStatus.ENDED; + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/auction/enums/AuctionStatus.java b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/enums/AuctionStatus.java new file mode 100644 index 0000000..e379f2b --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/enums/AuctionStatus.java @@ -0,0 +1,13 @@ +package com.rocketcrew.pocatbatch.domain.auction.enums; + +public enum AuctionStatus { + PENDING, + INSPECTING, + APPROVED, + REJECTED, + ACTIVE, + ENDED, + NO_BIDDER, + CANCELLED, + PAYMENT_PENDING +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/config/AuctionRankingProperties.java b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/config/AuctionRankingProperties.java new file mode 100644 index 0000000..3f84f5d --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/config/AuctionRankingProperties.java @@ -0,0 +1,17 @@ +package com.rocketcrew.pocatbatch.domain.auction.ranking.config; + +import lombok.Data; +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.stereotype.Component; + +@Component +@ConfigurationProperties(prefix = "pocat.batch.ranking.auction") +@Data +public class AuctionRankingProperties { + + private int cacheSize = 100; + private int ttlSeconds = 70; + private double likeWeight = 1.0; + private double bidWeight = 2.0; + private int maxResponseSize = 100; +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/dto/AuctionCountProjection.java b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/dto/AuctionCountProjection.java new file mode 100644 index 0000000..6333826 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/dto/AuctionCountProjection.java @@ -0,0 +1,6 @@ +package com.rocketcrew.pocatbatch.domain.auction.ranking.dto; + +public interface AuctionCountProjection { + Long getAuctionId(); + Long getCnt(); +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/repository/AuctionBidRepository.java b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/repository/AuctionBidRepository.java new file mode 100644 index 0000000..593d0b9 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/repository/AuctionBidRepository.java @@ -0,0 +1,15 @@ +package com.rocketcrew.pocatbatch.domain.auction.ranking.repository; + +import com.rocketcrew.pocatbatch.domain.auction.ranking.dto.AuctionCountProjection; +import com.rocketcrew.pocatbatch.domain.bid.entity.AuctionBid; +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; + +import java.util.List; + +public interface AuctionBidRepository extends JpaRepository { + + @Query("SELECT b.auctionId AS auctionId, COUNT(b) AS cnt FROM AuctionBid b WHERE b.auctionId IN :auctionIds GROUP BY b.auctionId") + List countByAuctionIdIn(@Param("auctionIds") List auctionIds); +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/repository/LikeRepository.java b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/repository/LikeRepository.java new file mode 100644 index 0000000..8b74a99 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/repository/LikeRepository.java @@ -0,0 +1,15 @@ +package com.rocketcrew.pocatbatch.domain.auction.ranking.repository; + +import com.rocketcrew.pocatbatch.domain.auction.ranking.dto.AuctionCountProjection; +import com.rocketcrew.pocatbatch.domain.like.entity.Like; +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; + +import java.util.List; + +public interface LikeRepository extends JpaRepository { + + @Query("SELECT l.auctionId AS auctionId, COUNT(l) AS cnt FROM Like l WHERE l.auctionId IN :auctionIds GROUP BY l.auctionId") + List countByAuctionIdIn(@Param("auctionIds") List auctionIds); +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/service/AuctionRankingService.java b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/service/AuctionRankingService.java new file mode 100644 index 0000000..2f4c946 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/ranking/service/AuctionRankingService.java @@ -0,0 +1,84 @@ +package com.rocketcrew.pocatbatch.domain.auction.ranking.service; + +import com.rocketcrew.pocatbatch.domain.auction.entity.Auction; +import com.rocketcrew.pocatbatch.domain.auction.enums.AuctionStatus; +import com.rocketcrew.pocatbatch.domain.auction.ranking.config.AuctionRankingProperties; +import com.rocketcrew.pocatbatch.domain.auction.ranking.dto.AuctionCountProjection; +import com.rocketcrew.pocatbatch.domain.auction.ranking.repository.AuctionBidRepository; +import com.rocketcrew.pocatbatch.domain.auction.ranking.repository.LikeRepository; +import com.rocketcrew.pocatbatch.domain.auction.repository.AuctionRepository; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.data.redis.core.StringRedisTemplate; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; + +@Slf4j +@Service +@RequiredArgsConstructor +@Transactional(readOnly = true) +public class AuctionRankingService { + + private final StringRedisTemplate redisTemplate; + private final AuctionRepository auctionRepository; + private final LikeRepository likeRepository; + private final AuctionBidRepository auctionBidRepository; + private final AuctionRankingProperties properties; + + static final String RANKING_KEY = "ranking:auction:popular"; + + public void refreshRanking() { + try { + List activeAuctions = auctionRepository.findAllByStatus(AuctionStatus.ACTIVE); + if (activeAuctions.isEmpty()) { + redisTemplate.delete(RANKING_KEY); + log.debug("활성 경매 없음 — 랭킹 캐시 삭제"); + return; + } + + List auctionIds = activeAuctions.stream().map(Auction::getId).toList(); + + Map likeCounts = toLongMap(likeRepository.countByAuctionIdIn(auctionIds)); + Map bidCounts = toLongMap(auctionBidRepository.countByAuctionIdIn(auctionIds)); + + String stagingKey = RANKING_KEY + ":staging:" + UUID.randomUUID().toString().replace("-", "").substring(0, 8); + try { + for (Auction auction : activeAuctions) { + long likeCount = likeCounts.getOrDefault(auction.getId(), 0L); + long bidCount = bidCounts.getOrDefault(auction.getId(), 0L); + double score = likeCount * properties.getLikeWeight() + bidCount * properties.getBidWeight(); + redisTemplate.opsForZSet().add(stagingKey, auction.getId().toString(), score); + } + + trimToCacheSize(stagingKey); + redisTemplate.rename(stagingKey, RANKING_KEY); + redisTemplate.expire(RANKING_KEY, properties.getTtlSeconds(), TimeUnit.SECONDS); + log.debug("경매 랭킹 갱신 완료: {} 개 경매", activeAuctions.size()); + } finally { + redisTemplate.delete(stagingKey); + } + + } catch (Exception e) { + log.warn("경매 랭킹 갱신 실패", e); + throw new RuntimeException("경매 랭킹 갱신 실패", e); + } + } + + private void trimToCacheSize(String key) { + Long size = redisTemplate.opsForZSet().zCard(key); + if (size != null && size > properties.getCacheSize()) { + redisTemplate.opsForZSet().removeRange(key, 0, size - properties.getCacheSize() - 1); + } + } + + private Map toLongMap(List projections) { + return projections.stream() + .collect(Collectors.toMap(AuctionCountProjection::getAuctionId, AuctionCountProjection::getCnt)); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/auction/repository/AuctionRepository.java b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/repository/AuctionRepository.java new file mode 100644 index 0000000..9d84d7f --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/repository/AuctionRepository.java @@ -0,0 +1,23 @@ +package com.rocketcrew.pocatbatch.domain.auction.repository; + +import com.rocketcrew.pocatbatch.domain.auction.entity.Auction; +import com.rocketcrew.pocatbatch.domain.auction.enums.AuctionStatus; +import org.springframework.data.jpa.repository.JpaRepository; + +import java.time.LocalDateTime; +import java.util.List; + +public interface AuctionRepository extends JpaRepository { + + List findAllByStatus(AuctionStatus status); + + List findAllByStatusAndEndedAtLessThanEqualOrderByEndedAtAsc( + AuctionStatus status, + LocalDateTime endedAt + ); + + List findAllByStatusAndUpdatedAtLessThanEqualOrderByUpdatedAtAsc( + AuctionStatus status, + LocalDateTime updatedAt + ); +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/auction/service/AuctionBatchService.java b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/service/AuctionBatchService.java new file mode 100644 index 0000000..2cdc974 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/auction/service/AuctionBatchService.java @@ -0,0 +1,50 @@ +package com.rocketcrew.pocatbatch.domain.auction.service; + +import com.rocketcrew.pocatbatch.domain.auction.entity.Auction; +import com.rocketcrew.pocatbatch.domain.auction.repository.AuctionRepository; +import com.rocketcrew.pocatbatch.domain.outbox.service.OutboxEventWriter; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +import java.time.LocalDateTime; + +@Slf4j +@Service +@RequiredArgsConstructor +public class AuctionBatchService { + + private final AuctionRepository auctionRepository; + private final OutboxEventWriter outboxEventWriter; + + @Transactional + public void activateAuction(Auction auction) { + LocalDateTime now = LocalDateTime.now(); + LocalDateTime startedAt = now; + LocalDateTime endedAt = now.plusHours(7); + + auction.activate(startedAt, endedAt); + auctionRepository.save(auction); + + String payload = String.format( + "{\"auctionId\":%d,\"status\":\"ACTIVE\",\"startedAt\":\"%s\",\"endedAt\":\"%s\"}", + auction.getId(), startedAt, endedAt + ); + outboxEventWriter.write("auction", String.valueOf(auction.getId()), "auction.activated", payload); + log.info("경매 활성화 처리: auctionId={}", auction.getId()); + } + + @Transactional + public void endAuction(Auction auction) { + auction.end(); + auctionRepository.save(auction); + + String payload = String.format( + "{\"auctionId\":%d,\"status\":\"ENDED\",\"endedAt\":\"%s\"}", + auction.getId(), auction.getEndedAt() + ); + outboxEventWriter.write("auction", String.valueOf(auction.getId()), "auction.ended", payload); + log.info("경매 종료 처리: auctionId={}", auction.getId()); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/bid/entity/AuctionBid.java b/src/main/java/com/rocketcrew/pocatbatch/domain/bid/entity/AuctionBid.java new file mode 100644 index 0000000..cc3ed4b --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/bid/entity/AuctionBid.java @@ -0,0 +1,51 @@ +package com.rocketcrew.pocatbatch.domain.bid.entity; + +import com.rocketcrew.pocatbatch.domain.bid.enums.BidStatus; +import com.rocketcrew.pocatbatch.domain.freepost.entity.BaseEntity; +import jakarta.persistence.*; +import lombok.*; +import org.hibernate.annotations.SQLDelete; +import org.hibernate.annotations.SQLRestriction; + +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +@Builder +@AllArgsConstructor +@Entity +@Table(name = "auction_bids") +@SQLDelete(sql = "UPDATE auction_bids SET deleted_at = NOW() WHERE id = ?") +@SQLRestriction("deleted_at IS NULL") +public class AuctionBid extends BaseEntity { + + @Column(name = "user_id", nullable = false) + private Long userId; + + @Column(name = "auction_id", nullable = false) + private Long auctionId; + + @Column(name = "bid_price", nullable = false) + private Long bidPrice; + + @Enumerated(EnumType.STRING) + @Column(name = "status", nullable = false, length = 20) + private BidStatus status; + + public void markOutbid() { + this.status = BidStatus.OUTBID; + } + + public void cancel() { + this.status = BidStatus.CANCELLED; + } + + public void markWon() { + this.status = BidStatus.WON; + } + + public void markLost() { + if (this.status == BidStatus.LOST) { + return; + } + this.status = BidStatus.LOST; + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/bid/enums/BidStatus.java b/src/main/java/com/rocketcrew/pocatbatch/domain/bid/enums/BidStatus.java new file mode 100644 index 0000000..716e071 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/bid/enums/BidStatus.java @@ -0,0 +1,9 @@ +package com.rocketcrew.pocatbatch.domain.bid.enums; + +public enum BidStatus { + LEADING, // 현재 최고 입찰 + OUTBID, // 다른 입찰자에게 밀림 + WON, // 최종 낙찰 + LOST, // 최종 패찰 + CANCELLED +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/Card.java b/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/Card.java new file mode 100644 index 0000000..1d7b59c --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/Card.java @@ -0,0 +1,76 @@ +package com.rocketcrew.pocatbatch.domain.card.entity; + +import com.rocketcrew.pocatbatch.domain.card.entity.enums.CardCategory; +import com.rocketcrew.pocatbatch.domain.card.entity.enums.CardGrade; +import com.rocketcrew.pocatbatch.domain.card.entity.enums.CardSource; +import com.rocketcrew.pocatbatch.domain.card.entity.enums.CardStatus; +import com.rocketcrew.pocatbatch.domain.freepost.entity.BaseEntity; +import jakarta.persistence.*; +import lombok.*; +import org.hibernate.annotations.SQLDelete; +import org.hibernate.annotations.SQLRestriction; + +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +@Builder +@AllArgsConstructor +@Entity +@Table(name = "cards", + indexes = { + @Index(name = "idx_cards_status", columnList = "status"), + @Index(name = "idx_cards_user_id_status", columnList = "user_id, status"), + @Index(name = "idx_cards_grade", columnList = "grade"), + @Index(name = "idx_cards_category", columnList = "category"), + @Index(name = "idx_cards_created_at", columnList = "created_at") + } +) +@SQLDelete(sql = "UPDATE cards SET deleted_at = NOW() WHERE id = ?") +@SQLRestriction("deleted_at IS NULL") +public class Card extends BaseEntity { + + @Column(name = "user_id", nullable = false) + private Long userId; + + @Column(name = "tcgdex_id", length = 100) + private String tcgdexId; + + @Column(name = "name", nullable = false) + private String name; + + @Column(name = "series_id") + private Long seriesId; + + @Column(name = "pokemon_set_id") + private Long pokemonSetId; + + @Column(name = "pokemon_id") + private Long pokemonId; + + @Column(name = "card_number", length = 20, nullable = false) + private String cardNumber; + + @Column(name = "rarity", length = 50, nullable = false) + private String rarity; + + @Enumerated(EnumType.STRING) + @Column(name = "category", length = 20, nullable = false) + private CardCategory category; + + @Enumerated(EnumType.STRING) + @Column(name = "grade", nullable = false, length = 20) + private CardGrade grade; + + @Column(name = "image_url", length = 500) + private String imageUrl; + + @Enumerated(EnumType.STRING) + @Column(name = "source", nullable = false, length = 20) + private CardSource source; + + @Enumerated(EnumType.STRING) + @Column(name = "status", nullable = false, length = 20) + private CardStatus status; + + @Column(name = "reject_reason", columnDefinition = "TEXT") + private String rejectReason; +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/enums/CardCategory.java b/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/enums/CardCategory.java new file mode 100644 index 0000000..c0693eb --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/enums/CardCategory.java @@ -0,0 +1,8 @@ +package com.rocketcrew.pocatbatch.domain.card.entity.enums; + +public enum CardCategory { + POKEMON, + TRAINERS, + ENERGY, + UNKNOWN +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/enums/CardGrade.java b/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/enums/CardGrade.java new file mode 100644 index 0000000..a917574 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/enums/CardGrade.java @@ -0,0 +1,15 @@ +package com.rocketcrew.pocatbatch.domain.card.entity.enums; + +public enum CardGrade { + PSA_10, + PSA_9, + PSA_8, + PSA_7, + PSA_6, + PSA_5, + PSA_4, + PSA_3, + PSA_2, + PSA_1, + UNGRADED +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/enums/CardSource.java b/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/enums/CardSource.java new file mode 100644 index 0000000..ad9bf2e --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/enums/CardSource.java @@ -0,0 +1,6 @@ +package com.rocketcrew.pocatbatch.domain.card.entity.enums; + +public enum CardSource { + TCGDEX, + MANUAL +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/enums/CardStatus.java b/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/enums/CardStatus.java new file mode 100644 index 0000000..eaacaba --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/card/entity/enums/CardStatus.java @@ -0,0 +1,7 @@ +package com.rocketcrew.pocatbatch.domain.card.entity.enums; + +public enum CardStatus { + PENDING, + ACTIVE, + REJECTED +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/card/repository/CardRepository.java b/src/main/java/com/rocketcrew/pocatbatch/domain/card/repository/CardRepository.java new file mode 100644 index 0000000..a8bf4e2 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/card/repository/CardRepository.java @@ -0,0 +1,16 @@ +package com.rocketcrew.pocatbatch.domain.card.repository; + +import com.rocketcrew.pocatbatch.domain.card.entity.Card; +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; + +import java.util.Optional; + +public interface CardRepository extends JpaRepository { + + boolean existsByTcgdexId(String tcgdexId); + + @Query(value = "SELECT * FROM cards WHERE tcgdex_id = :tcgdexId LIMIT 1", nativeQuery = true) + Optional findByTcgdexIdIncludingDeleted(@Param("tcgdexId") String tcgdexId); +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/like/entity/Like.java b/src/main/java/com/rocketcrew/pocatbatch/domain/like/entity/Like.java new file mode 100644 index 0000000..c43e949 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/like/entity/Like.java @@ -0,0 +1,25 @@ +package com.rocketcrew.pocatbatch.domain.like.entity; + +import com.rocketcrew.pocatbatch.domain.freepost.entity.BaseEntity; +import jakarta.persistence.*; +import lombok.*; +import org.hibernate.annotations.SQLDelete; +import org.hibernate.annotations.SQLRestriction; + +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +@Builder +@AllArgsConstructor +@Entity +@Table(name = "likes", + indexes = @Index(name = "idx_likes_user_auction", columnList = "user_id, auction_id")) +@SQLDelete(sql = "UPDATE likes SET deleted_at = NOW() WHERE id = ?") +@SQLRestriction("deleted_at IS NULL") +public class Like extends BaseEntity { + + @Column(name = "user_id", nullable = false) + private Long userId; + + @Column(name = "auction_id") + private Long auctionId; +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/entity/OutboxEvent.java b/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/entity/OutboxEvent.java new file mode 100644 index 0000000..d86c1f2 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/entity/OutboxEvent.java @@ -0,0 +1,86 @@ +package com.rocketcrew.pocatbatch.domain.outbox.entity; + +import com.rocketcrew.pocatbatch.domain.outbox.enums.OutboxStatus; +import jakarta.persistence.*; +import lombok.AccessLevel; +import lombok.Getter; +import lombok.NoArgsConstructor; +import org.springframework.data.annotation.CreatedDate; +import org.springframework.data.jpa.domain.support.AuditingEntityListener; + +import java.time.LocalDateTime; + +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +@Entity +@Table( + name = "outbox_events", + indexes = @Index(name = "idx_outbox_status_created", columnList = "status, created_at") +) +@EntityListeners(AuditingEntityListener.class) +public class OutboxEvent { + + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @Column(nullable = false, length = 50) + private String topic; + + @Column(name = "partition_key", nullable = false, length = 100) + private String partitionKey; + + @Column(name = "event_type", nullable = false, length = 100) + private String eventType; + + @Column(nullable = false, columnDefinition = "TEXT") + private String payload; + + @Enumerated(EnumType.STRING) + @Column(nullable = false, length = 20) + private OutboxStatus status; + + @Column(name = "retry_count", nullable = false) + private int retryCount = 0; + + @CreatedDate + @Column(name = "created_at", updatable = false) + private LocalDateTime createdAt; + + @Column(name = "processed_at") + private LocalDateTime processedAt; + + public static OutboxEvent pending(String topic, String partitionKey, String eventType, String payload) { + OutboxEvent event = new OutboxEvent(); + event.topic = topic; + event.partitionKey = partitionKey; + event.eventType = eventType; + event.payload = payload; + event.status = OutboxStatus.PENDING; + return event; + } + + public void markSent() { + this.status = OutboxStatus.SENT; + this.processedAt = LocalDateTime.now(); + } + + public void markPendingForRetry() { + this.retryCount++; + if (this.retryCount >= 5) { + this.status = OutboxStatus.FAILED; + this.processedAt = LocalDateTime.now(); + } else { + this.status = OutboxStatus.PENDING; + } + } + + public void changeStatusToProcessing() { + this.status = OutboxStatus.PROCESSING; + } + + public void resetToPending() { + this.status = OutboxStatus.PENDING; + this.processedAt = null; + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/enums/OutboxStatus.java b/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/enums/OutboxStatus.java new file mode 100644 index 0000000..8c6e492 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/enums/OutboxStatus.java @@ -0,0 +1,8 @@ +package com.rocketcrew.pocatbatch.domain.outbox.enums; + +public enum OutboxStatus { + PENDING, + PROCESSING, + SENT, + FAILED +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/repository/OutboxRepository.java b/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/repository/OutboxRepository.java new file mode 100644 index 0000000..82a1070 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/repository/OutboxRepository.java @@ -0,0 +1,45 @@ +package com.rocketcrew.pocatbatch.domain.outbox.repository; + +import com.rocketcrew.pocatbatch.domain.outbox.entity.OutboxEvent; +import com.rocketcrew.pocatbatch.domain.outbox.enums.OutboxStatus; +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Modifying; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; + +import java.time.LocalDateTime; +import java.util.List; + +public interface OutboxRepository extends JpaRepository { + + List findTop100ByStatusOrderByCreatedAtAsc(OutboxStatus status); + + List findTop100ByStatusAndCreatedAtBeforeOrderByCreatedAtAsc(OutboxStatus status, LocalDateTime before); + + /** + * 조건부 UPDATE (선점용) + */ + @Modifying(clearAutomatically = true) + @Query("UPDATE OutboxEvent o SET o.status = :to, o.processedAt = CURRENT_TIMESTAMP " + + "WHERE o.id = :id AND o.status = :from") + int markProcessingIfPending( + @Param("id") Long id, + @Param("from") OutboxStatus from, + @Param("to") OutboxStatus to + ); + + /** + * PROCESSING 상태에서 일정 시간 이상 멈춘 이벤트 조회 (Reaper용, 최대 100건) + */ + List findTop100ByStatusAndProcessedAtBefore(OutboxStatus status, LocalDateTime before); + + /** + * 오래된 SENT 이벤트 삭제 + */ + @Modifying(clearAutomatically = true) + @Query("DELETE FROM OutboxEvent o WHERE o.status = :status AND o.createdAt < :cutoff") + int deleteOldEvents( + @Param("status") OutboxStatus status, + @Param("cutoff") LocalDateTime cutoff + ); +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/service/OutboxEventWriter.java b/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/service/OutboxEventWriter.java new file mode 100644 index 0000000..8e28c54 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/service/OutboxEventWriter.java @@ -0,0 +1,46 @@ +package com.rocketcrew.pocatbatch.domain.outbox.service; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.rocketcrew.pocatbatch.domain.outbox.entity.OutboxEvent; +import com.rocketcrew.pocatbatch.domain.outbox.repository.OutboxRepository; +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Component; + +@Component +@RequiredArgsConstructor +public class OutboxEventWriter { + + private final OutboxRepository outboxRepository; + private final ObjectMapper objectMapper; + + /** + * 비즈니스 이벤트를 JSON 페이로드로 변환하여 아웃박스 테이블에 PENDING 상태로 저장합니다. + * + * @param topic 카프카 토픽명 + * @param partitionKey 카프카 파티션 키 + * @param eventType 이벤트 타입 + * @param payload JSON 페이로드 + */ + public void write(String topic, String partitionKey, String eventType, String payload) { + try { + OutboxEvent outboxEvent = OutboxEvent.pending(topic, partitionKey, eventType, payload); + outboxRepository.save(outboxEvent); + } catch (Exception e) { + throw new RuntimeException("아웃박스 이벤트 저장 중 치명적 에러 발생: " + eventType, e); + } + } + + /** + * 객체를 JSON으로 변환하여 아웃박스에 저장합니다. + */ + public void write(String topic, String partitionKey, String eventType, Object eventPayload) { + String payload; + try { + payload = objectMapper.writeValueAsString(eventPayload); + } catch (JsonProcessingException e) { + throw new RuntimeException("이벤트 직렬화 실패: " + eventType, e); + } + write(topic, partitionKey, eventType, payload); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/service/OutboxProcessor.java b/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/service/OutboxProcessor.java new file mode 100644 index 0000000..db85d80 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/outbox/service/OutboxProcessor.java @@ -0,0 +1,61 @@ +package com.rocketcrew.pocatbatch.domain.outbox.service; + +import com.rocketcrew.pocatbatch.config.KafkaProducerConfig; +import com.rocketcrew.pocatbatch.domain.outbox.entity.OutboxEvent; +import com.rocketcrew.pocatbatch.domain.outbox.enums.OutboxStatus; +import com.rocketcrew.pocatbatch.domain.outbox.repository.OutboxRepository; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; + +import java.util.concurrent.TimeUnit; + +@Component +@RequiredArgsConstructor +@Slf4j +public class OutboxProcessor { + + private final OutboxRepository outboxRepository; + + /** + * REQUIRES_NEW를 통해 건당 독립적인 트랜잭션을 보장합니다. + * 하나의 이벤트가 실패하거나 롤백되어도 다른 이벤트에 영향을 주지 않습니다. + */ + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void processEvent(OutboxEvent event, KafkaTemplate template) { + + // 1. 조건부 UPDATE로 선점 (Atomic) + int updated = outboxRepository.markProcessingIfPending( + event.getId(), OutboxStatus.PENDING, OutboxStatus.PROCESSING + ); + + if (updated == 0) return; // 이미 다른 서버가 선점 + + // relay()에 @Transactional이 없어 findTop100... 이후 엔티티가 detached 상태임. + // REQUIRES_NEW 트랜잭션 안에서 재조회해야 dirty checking이 정상 동작함. + OutboxEvent managed = outboxRepository.findById(event.getId()).orElse(null); + if (managed == null) return; + + try { + // 3. Kafka 발행 (동기 대기) + template.send(managed.getTopic(), managed.getPartitionKey(), managed.getPayload()) + .get(5, TimeUnit.SECONDS); + + // 4. 발행 성공 ➡️ SENT + managed.markSent(); + log.info("릴레이 성공: id={}, topic={}", managed.getId(), managed.getTopic()); + + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + log.warn("릴레이 인터럽트 발생: id={}", managed.getId()); + managed.markPendingForRetry(); + } catch (Exception e) { + managed.markPendingForRetry(); + log.error("릴레이 실패: id={}, retryCount={}", managed.getId(), managed.getRetryCount(), e); + } + // managed 엔티티이므로 트랜잭션 커밋 시 dirty checking으로 자동 UPDATE됨. + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/pokemon/entity/Pokemon.java b/src/main/java/com/rocketcrew/pocatbatch/domain/pokemon/entity/Pokemon.java new file mode 100644 index 0000000..a35d3ae --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/pokemon/entity/Pokemon.java @@ -0,0 +1,29 @@ +package com.rocketcrew.pocatbatch.domain.pokemon.entity; + +import com.rocketcrew.pocatbatch.domain.freepost.entity.BaseEntity; +import jakarta.persistence.*; +import lombok.*; +import org.hibernate.annotations.SQLDelete; +import org.hibernate.annotations.SQLRestriction; + +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +@Builder +@AllArgsConstructor +@Entity +@SQLDelete(sql = "UPDATE pokemon SET deleted_at = NOW() WHERE id = ?") +@SQLRestriction("deleted_at IS NULL") +@Table(name = "pokemon", + uniqueConstraints = @UniqueConstraint(columnNames = "name")) +public class Pokemon extends BaseEntity { + + @Column(name = "name", nullable = false, length = 100) + private String name; // "Charizard" + + @Column(name = "name_ko", length = 100) + private String nameKo; // "리자몽" + + public void updateNameKo(String nameKo) { + this.nameKo = nameKo; + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/refund/entity/Refund.java b/src/main/java/com/rocketcrew/pocatbatch/domain/refund/entity/Refund.java new file mode 100644 index 0000000..c400896 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/refund/entity/Refund.java @@ -0,0 +1,53 @@ +package com.rocketcrew.pocatbatch.domain.refund.entity; + +import com.rocketcrew.pocatbatch.domain.freepost.entity.BaseEntity; +import jakarta.persistence.*; +import lombok.*; +import org.hibernate.annotations.SQLDelete; +import org.hibernate.annotations.SQLRestriction; + +import java.time.LocalDateTime; + +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +@Builder +@AllArgsConstructor +@Entity +@Table(name = "refunds") +@SQLDelete(sql = "UPDATE refunds SET deleted_at = NOW() WHERE id = ?") +@SQLRestriction("deleted_at IS NULL") +public class Refund extends BaseEntity { + + @Column(name = "order_id", nullable = false) + private Long orderId; + + @Column(name = "payment_id", nullable = false) + private Long paymentId; + + @Column(name = "amount", nullable = false) + private Long amount; + + @Column(name = "reason", nullable = false, length = 100) + private String reason; + + @Column(name = "reject_reason", length = 100) + private String rejectReason; + + @Enumerated(EnumType.STRING) + @Column(name = "status", nullable = false, length = 20) + private RefundStatus status; + + @Builder.Default + @Column(name = "retry_count", nullable = false) + private int retryCount = 0; + + @Column(name = "next_retry_at") + private LocalDateTime nextRetryAt; + + @Column(name = "failure_reason", length = 255) + private String failureReason; + + public boolean isRetryDue(LocalDateTime now) { + return nextRetryAt == null || !now.isBefore(nextRetryAt); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/refund/entity/RefundStatus.java b/src/main/java/com/rocketcrew/pocatbatch/domain/refund/entity/RefundStatus.java new file mode 100644 index 0000000..4d17bb5 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/refund/entity/RefundStatus.java @@ -0,0 +1,10 @@ +package com.rocketcrew.pocatbatch.domain.refund.entity; + +public enum RefundStatus { + REQUESTED, + PROCESSING, + COMPLETED, + REJECTED, + FAILED_RETRYABLE, + FAILED_FINAL +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/refund/repository/RefundRepository.java b/src/main/java/com/rocketcrew/pocatbatch/domain/refund/repository/RefundRepository.java new file mode 100644 index 0000000..922ffbe --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/refund/repository/RefundRepository.java @@ -0,0 +1,27 @@ +package com.rocketcrew.pocatbatch.domain.refund.repository; + +import com.rocketcrew.pocatbatch.domain.refund.entity.Refund; +import com.rocketcrew.pocatbatch.domain.refund.entity.RefundStatus; +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; + +import java.time.LocalDateTime; +import java.util.List; + +public interface RefundRepository extends JpaRepository { + + /** + * 자동 재시도 대상 조회: + * 1) FAILED_RETRYABLE 중 nextRetryAt이 지난 것 + * 2) PROCESSING 중 updatedAt이 stuckBefore보다 오래된 것 (approveRefund DB 실패로 방치된 건) + */ + @Query("SELECT r FROM Refund r WHERE " + + "(r.status = :retryable AND (r.nextRetryAt IS NULL OR r.nextRetryAt <= :now)) OR " + + "(r.status = :processing AND r.updatedAt < :stuckBefore)") + List findRetryableTargets( + @Param("retryable") RefundStatus retryable, + @Param("now") LocalDateTime now, + @Param("processing") RefundStatus processing, + @Param("stuckBefore") LocalDateTime stuckBefore); +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/series/entity/Series.java b/src/main/java/com/rocketcrew/pocatbatch/domain/series/entity/Series.java new file mode 100644 index 0000000..677a370 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/series/entity/Series.java @@ -0,0 +1,29 @@ +package com.rocketcrew.pocatbatch.domain.series.entity; + +import com.rocketcrew.pocatbatch.domain.freepost.entity.BaseEntity; +import jakarta.persistence.*; +import lombok.*; +import org.hibernate.annotations.SQLDelete; +import org.hibernate.annotations.SQLRestriction; + +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +@Builder +@AllArgsConstructor +@Entity +@SQLDelete(sql = "UPDATE series SET deleted_at = NOW() WHERE id = ?") +@SQLRestriction("deleted_at IS NULL") +@Table(name = "series", + uniqueConstraints = @UniqueConstraint(columnNames = "name")) +public class Series extends BaseEntity { + + @Column(name = "name", nullable = false, length = 100) + private String name; // "Sword & Shield" (TCGdex 영문명) + + @Column(name = "name_ko", length = 500) + private String nameKo; // "검과방패 소드실드 소드앤실드" (공백 구분 한글 별칭) + + public void updateNameKo(String nameKo) { + this.nameKo = nameKo; + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/domain/set/entity/PokemonSet.java b/src/main/java/com/rocketcrew/pocatbatch/domain/set/entity/PokemonSet.java new file mode 100644 index 0000000..e942b18 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/domain/set/entity/PokemonSet.java @@ -0,0 +1,37 @@ +package com.rocketcrew.pocatbatch.domain.set.entity; + +import com.rocketcrew.pocatbatch.domain.freepost.entity.BaseEntity; +import com.rocketcrew.pocatbatch.domain.series.entity.Series; +import jakarta.persistence.*; +import lombok.*; +import org.hibernate.annotations.SQLDelete; +import org.hibernate.annotations.SQLRestriction; + +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +@Builder +@AllArgsConstructor +@Entity +@SQLDelete(sql = "UPDATE pokemon_sets SET deleted_at = NOW() WHERE id = ?") +@SQLRestriction("deleted_at IS NULL") +@Table(name = "pokemon_sets", + uniqueConstraints = @UniqueConstraint(columnNames = "set_id")) +public class PokemonSet extends BaseEntity { + + @ManyToOne(fetch = FetchType.LAZY) + @JoinColumn(name = "series_id") + private Series series; + + @Column(name = "set_id", nullable = false, length = 50) + private String setId; // "swsh5" (TCGdex ID) + + @Column(name = "name", nullable = false, length = 100) + private String name; // "Rebel Clash" + + @Column(name = "name_ko", length = 1000) + private String nameKo; // "반역크래시" (공백 구분 한글 별칭) + + public void updateNameKo(String nameKo) { + this.nameKo = nameKo; + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/aisession/AiSessionCleanupJobConfig.java b/src/main/java/com/rocketcrew/pocatbatch/job/aisession/AiSessionCleanupJobConfig.java new file mode 100644 index 0000000..e36459b --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/aisession/AiSessionCleanupJobConfig.java @@ -0,0 +1,37 @@ +package com.rocketcrew.pocatbatch.job.aisession; + +import lombok.RequiredArgsConstructor; +import org.springframework.batch.core.Job; +import org.springframework.batch.core.Step; +import org.springframework.batch.core.job.builder.JobBuilder; +import org.springframework.batch.core.repository.JobRepository; +import org.springframework.batch.core.step.builder.StepBuilder; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.transaction.PlatformTransactionManager; + +@Configuration +@RequiredArgsConstructor +public class AiSessionCleanupJobConfig { + + public static final String JOB_NAME = "aiSessionCleanupJob"; + public static final String STEP_NAME = "aiSessionCleanupStep"; + + private final JobRepository jobRepository; + private final PlatformTransactionManager transactionManager; + private final AiSessionCleanupTasklet aiSessionCleanupTasklet; + + @Bean(name = JOB_NAME) + public Job aiSessionCleanupJob() { + return new JobBuilder(JOB_NAME, jobRepository) + .start(aiSessionCleanupStep()) + .build(); + } + + @Bean(name = STEP_NAME) + public Step aiSessionCleanupStep() { + return new StepBuilder(STEP_NAME, jobRepository) + .tasklet(aiSessionCleanupTasklet, transactionManager) + .build(); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/aisession/AiSessionCleanupTasklet.java b/src/main/java/com/rocketcrew/pocatbatch/job/aisession/AiSessionCleanupTasklet.java new file mode 100644 index 0000000..9e041c6 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/aisession/AiSessionCleanupTasklet.java @@ -0,0 +1,34 @@ +package com.rocketcrew.pocatbatch.job.aisession; + +import com.rocketcrew.pocatbatch.domain.ai.repository.AiChatSessionRepository; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.stereotype.Component; + +import java.time.LocalDateTime; + +@Slf4j +@Component +@RequiredArgsConstructor +public class AiSessionCleanupTasklet implements Tasklet { + + private final AiChatSessionRepository aiChatSessionRepository; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + try { + // 30분 이상 비활성인 세션 만료 처리 + LocalDateTime threshold = LocalDateTime.now().minusMinutes(30); + int expiredCount = aiChatSessionRepository.expireSessionsBeforeTime(threshold); + log.info("AI 세션 만료 처리 완료: {}개 세션 만료됨", expiredCount); + return RepeatStatus.FINISHED; + } catch (Exception e) { + log.error("AI 세션 만료 처리 실패", e); + throw new RuntimeException("AI 세션 만료 처리 중 오류 발생", e); + } + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/auctionactivation/AuctionActivationJobConfig.java b/src/main/java/com/rocketcrew/pocatbatch/job/auctionactivation/AuctionActivationJobConfig.java new file mode 100644 index 0000000..aa18342 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/auctionactivation/AuctionActivationJobConfig.java @@ -0,0 +1,37 @@ +package com.rocketcrew.pocatbatch.job.auctionactivation; + +import lombok.RequiredArgsConstructor; +import org.springframework.batch.core.Job; +import org.springframework.batch.core.Step; +import org.springframework.batch.core.job.builder.JobBuilder; +import org.springframework.batch.core.repository.JobRepository; +import org.springframework.batch.core.step.builder.StepBuilder; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.transaction.PlatformTransactionManager; + +@Configuration +@RequiredArgsConstructor +public class AuctionActivationJobConfig { + + public static final String JOB_NAME = "auctionActivationJob"; + public static final String STEP_NAME = "auctionActivationStep"; + + private final JobRepository jobRepository; + private final PlatformTransactionManager transactionManager; + private final AuctionActivationTasklet auctionActivationTasklet; + + @Bean(name = JOB_NAME) + public Job auctionActivationJob() { + return new JobBuilder(JOB_NAME, jobRepository) + .start(auctionActivationStep()) + .build(); + } + + @Bean(name = STEP_NAME) + public Step auctionActivationStep() { + return new StepBuilder(STEP_NAME, jobRepository) + .tasklet(auctionActivationTasklet, transactionManager) + .build(); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/auctionactivation/AuctionActivationTasklet.java b/src/main/java/com/rocketcrew/pocatbatch/job/auctionactivation/AuctionActivationTasklet.java new file mode 100644 index 0000000..0e9189b --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/auctionactivation/AuctionActivationTasklet.java @@ -0,0 +1,69 @@ +package com.rocketcrew.pocatbatch.job.auctionactivation; + +import com.rocketcrew.pocatbatch.domain.auction.entity.Auction; +import com.rocketcrew.pocatbatch.domain.auction.enums.AuctionStatus; +import com.rocketcrew.pocatbatch.domain.auction.repository.AuctionRepository; +import com.rocketcrew.pocatbatch.domain.auction.service.AuctionBatchService; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.redisson.api.RLock; +import org.redisson.api.RedissonClient; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.stereotype.Component; + +import java.util.List; +import java.util.concurrent.TimeUnit; + +@Slf4j +@Component +@RequiredArgsConstructor +public class AuctionActivationTasklet implements Tasklet { + + private final AuctionRepository auctionRepository; + private final AuctionBatchService auctionBatchService; + private final RedissonClient redissonClient; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + try { + List approvedAuctions = auctionRepository.findAllByStatus(AuctionStatus.APPROVED); + int activatedCount = 0; + + for (Auction auction : approvedAuctions) { + RLock lock = redissonClient.getLock("auction:lock:" + auction.getId()); + boolean locked; + try { + locked = lock.tryLock(0, 30, TimeUnit.SECONDS); + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + log.warn("경매 활성화 락 획득 중 인터럽트: auctionId={}", auction.getId()); + throw new RuntimeException("경매 활성화 인터럽트 발생", ie); + } + if (!locked) { + log.debug("경매 활성화 락 실패 (이미 처리 중): auctionId={}", auction.getId()); + continue; + } + + try { + auctionBatchService.activateAuction(auction); + activatedCount++; + } catch (Exception e) { + log.error("경매 활성화 실패: auctionId={}", auction.getId(), e); + } finally { + if (lock.isHeldByCurrentThread()) { + lock.unlock(); + } + } + } + + log.info("경매 활성화 완료: {} 개 경매 활성화됨", activatedCount); + return RepeatStatus.FINISHED; + } catch (Exception e) { + log.error("경매 활성화 작업 실패", e); + throw new RuntimeException("경매 활성화 중 오류 발생", e); + } + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/auctionexpiration/AuctionExpirationJobConfig.java b/src/main/java/com/rocketcrew/pocatbatch/job/auctionexpiration/AuctionExpirationJobConfig.java new file mode 100644 index 0000000..3d6770a --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/auctionexpiration/AuctionExpirationJobConfig.java @@ -0,0 +1,37 @@ +package com.rocketcrew.pocatbatch.job.auctionexpiration; + +import lombok.RequiredArgsConstructor; +import org.springframework.batch.core.Job; +import org.springframework.batch.core.Step; +import org.springframework.batch.core.job.builder.JobBuilder; +import org.springframework.batch.core.repository.JobRepository; +import org.springframework.batch.core.step.builder.StepBuilder; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.transaction.PlatformTransactionManager; + +@Configuration +@RequiredArgsConstructor +public class AuctionExpirationJobConfig { + + public static final String JOB_NAME = "auctionExpirationJob"; + public static final String STEP_NAME = "auctionExpirationStep"; + + private final JobRepository jobRepository; + private final PlatformTransactionManager transactionManager; + private final AuctionExpirationTasklet auctionExpirationTasklet; + + @Bean(name = JOB_NAME) + public Job auctionExpirationJob() { + return new JobBuilder(JOB_NAME, jobRepository) + .start(auctionExpirationStep()) + .build(); + } + + @Bean(name = STEP_NAME) + public Step auctionExpirationStep() { + return new StepBuilder(STEP_NAME, jobRepository) + .tasklet(auctionExpirationTasklet, transactionManager) + .build(); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/auctionexpiration/AuctionExpirationTasklet.java b/src/main/java/com/rocketcrew/pocatbatch/job/auctionexpiration/AuctionExpirationTasklet.java new file mode 100644 index 0000000..a9ea907 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/auctionexpiration/AuctionExpirationTasklet.java @@ -0,0 +1,74 @@ +package com.rocketcrew.pocatbatch.job.auctionexpiration; + +import com.rocketcrew.pocatbatch.domain.auction.entity.Auction; +import com.rocketcrew.pocatbatch.domain.auction.enums.AuctionStatus; +import com.rocketcrew.pocatbatch.domain.auction.repository.AuctionRepository; +import com.rocketcrew.pocatbatch.domain.auction.service.AuctionBatchService; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.redisson.api.RLock; +import org.redisson.api.RedissonClient; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.stereotype.Component; + +import java.time.LocalDateTime; +import java.util.List; +import java.util.concurrent.TimeUnit; + +@Slf4j +@Component +@RequiredArgsConstructor +public class AuctionExpirationTasklet implements Tasklet { + + private final AuctionRepository auctionRepository; + private final AuctionBatchService auctionBatchService; + private final RedissonClient redissonClient; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + try { + LocalDateTime now = LocalDateTime.now(); + List expiredAuctions = auctionRepository.findAllByStatusAndEndedAtLessThanEqualOrderByEndedAtAsc( + AuctionStatus.ACTIVE, now + ); + + int expiredCount = 0; + + for (Auction auction : expiredAuctions) { + RLock lock = redissonClient.getLock("auction:lock:" + auction.getId()); + boolean locked; + try { + locked = lock.tryLock(0, 30, TimeUnit.SECONDS); + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + log.warn("경매 종료 락 획득 중 인터럽트: auctionId={}", auction.getId()); + throw new RuntimeException("경매 종료 인터럽트 발생", ie); + } + if (!locked) { + log.debug("경매 종료 락 실패 (이미 처리 중): auctionId={}", auction.getId()); + continue; + } + + try { + auctionBatchService.endAuction(auction); + expiredCount++; + } catch (Exception e) { + log.error("경매 종료 실패: auctionId={}", auction.getId(), e); + } finally { + if (lock.isHeldByCurrentThread()) { + lock.unlock(); + } + } + } + + log.info("경매 종료 완료: {} 개 경매 종료됨", expiredCount); + return RepeatStatus.FINISHED; + } catch (Exception e) { + log.error("경매 종료 작업 실패", e); + throw new RuntimeException("경매 종료 중 오류 발생", e); + } + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/auctionranking/AuctionRankingJobConfig.java b/src/main/java/com/rocketcrew/pocatbatch/job/auctionranking/AuctionRankingJobConfig.java new file mode 100644 index 0000000..f0b5cc9 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/auctionranking/AuctionRankingJobConfig.java @@ -0,0 +1,37 @@ +package com.rocketcrew.pocatbatch.job.auctionranking; + +import lombok.RequiredArgsConstructor; +import org.springframework.batch.core.Job; +import org.springframework.batch.core.Step; +import org.springframework.batch.core.job.builder.JobBuilder; +import org.springframework.batch.core.repository.JobRepository; +import org.springframework.batch.core.step.builder.StepBuilder; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.transaction.PlatformTransactionManager; + +@Configuration +@RequiredArgsConstructor +public class AuctionRankingJobConfig { + + public static final String JOB_NAME = "auctionRankingJob"; + public static final String STEP_NAME = "auctionRankingStep"; + + private final JobRepository jobRepository; + private final PlatformTransactionManager transactionManager; + private final AuctionRankingTasklet auctionRankingTasklet; + + @Bean(name = JOB_NAME) + public Job auctionRankingJob() { + return new JobBuilder(JOB_NAME, jobRepository) + .start(auctionRankingStep()) + .build(); + } + + @Bean(name = STEP_NAME) + public Step auctionRankingStep() { + return new StepBuilder(STEP_NAME, jobRepository) + .tasklet(auctionRankingTasklet, transactionManager) + .build(); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/auctionranking/AuctionRankingTasklet.java b/src/main/java/com/rocketcrew/pocatbatch/job/auctionranking/AuctionRankingTasklet.java new file mode 100644 index 0000000..2aa2a08 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/auctionranking/AuctionRankingTasklet.java @@ -0,0 +1,30 @@ +package com.rocketcrew.pocatbatch.job.auctionranking; + +import com.rocketcrew.pocatbatch.domain.auction.ranking.service.AuctionRankingService; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.stereotype.Component; + +@Slf4j +@Component +@RequiredArgsConstructor +public class AuctionRankingTasklet implements Tasklet { + + private final AuctionRankingService auctionRankingService; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + try { + auctionRankingService.refreshRanking(); + log.info("경매 랭킹 갱신 완료"); + return RepeatStatus.FINISHED; + } catch (Exception e) { + log.error("경매 랭킹 갱신 실패", e); + throw new RuntimeException("경매 랭킹 갱신 중 오류 발생", e); + } + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/buyoutrecovery/BuyoutRecoveryJobConfig.java b/src/main/java/com/rocketcrew/pocatbatch/job/buyoutrecovery/BuyoutRecoveryJobConfig.java new file mode 100644 index 0000000..48bdb3a --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/buyoutrecovery/BuyoutRecoveryJobConfig.java @@ -0,0 +1,37 @@ +package com.rocketcrew.pocatbatch.job.buyoutrecovery; + +import lombok.RequiredArgsConstructor; +import org.springframework.batch.core.Job; +import org.springframework.batch.core.Step; +import org.springframework.batch.core.job.builder.JobBuilder; +import org.springframework.batch.core.repository.JobRepository; +import org.springframework.batch.core.step.builder.StepBuilder; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.transaction.PlatformTransactionManager; + +@Configuration +@RequiredArgsConstructor +public class BuyoutRecoveryJobConfig { + + public static final String JOB_NAME = "buyoutRecoveryJob"; + public static final String STEP_NAME = "buyoutRecoveryStep"; + + private final JobRepository jobRepository; + private final PlatformTransactionManager transactionManager; + private final BuyoutRecoveryTasklet buyoutRecoveryTasklet; + + @Bean(name = JOB_NAME) + public Job buyoutRecoveryJob() { + return new JobBuilder(JOB_NAME, jobRepository) + .start(buyoutRecoveryStep()) + .build(); + } + + @Bean(name = STEP_NAME) + public Step buyoutRecoveryStep() { + return new StepBuilder(STEP_NAME, jobRepository) + .tasklet(buyoutRecoveryTasklet, transactionManager) + .build(); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/buyoutrecovery/BuyoutRecoveryTasklet.java b/src/main/java/com/rocketcrew/pocatbatch/job/buyoutrecovery/BuyoutRecoveryTasklet.java new file mode 100644 index 0000000..4716847 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/buyoutrecovery/BuyoutRecoveryTasklet.java @@ -0,0 +1,54 @@ +package com.rocketcrew.pocatbatch.job.buyoutrecovery; + +import com.rocketcrew.pocatbatch.client.MainAppBuyoutClient; +import com.rocketcrew.pocatbatch.domain.auction.entity.Auction; +import com.rocketcrew.pocatbatch.domain.auction.enums.AuctionStatus; +import com.rocketcrew.pocatbatch.domain.auction.repository.AuctionRepository; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.stereotype.Component; + +import java.time.LocalDateTime; +import java.util.List; + +@Slf4j +@Component +@RequiredArgsConstructor +public class BuyoutRecoveryTasklet implements Tasklet { + + private final AuctionRepository auctionRepository; + private final MainAppBuyoutClient mainAppBuyoutClient; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + try { + // PAYMENT_PENDING + updatedAt <= now-2분 경매 조회 + LocalDateTime twoMinutesAgo = LocalDateTime.now().minusMinutes(2); + List stuckAuctions = auctionRepository.findAllByStatusAndUpdatedAtLessThanEqualOrderByUpdatedAtAsc( + AuctionStatus.PAYMENT_PENDING, twoMinutesAgo + ); + + int recoveredCount = 0; + long jobExecutionId = chunkContext.getStepContext().getStepExecution().getJobExecution().getId(); + + for (Auction auction : stuckAuctions) { + try { + mainAppBuyoutClient.recoverBuyout(auction.getId(), jobExecutionId); + recoveredCount++; + } catch (Exception e) { + log.warn("구매 확정 복구 실패: auctionId={}", auction.getId(), e); + } + } + + log.info("구매 확정 복구 완료: {} 개 경매 복구됨", recoveredCount); + return RepeatStatus.FINISHED; + } catch (Exception e) { + log.error("구매 확정 복구 작업 실패", e); + throw new RuntimeException("구매 확정 복구 중 오류 발생", e); + } + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/cardsync/CardSyncJobConfig.java b/src/main/java/com/rocketcrew/pocatbatch/job/cardsync/CardSyncJobConfig.java new file mode 100644 index 0000000..0ff773c --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/cardsync/CardSyncJobConfig.java @@ -0,0 +1,37 @@ +package com.rocketcrew.pocatbatch.job.cardsync; + +import lombok.RequiredArgsConstructor; +import org.springframework.batch.core.Job; +import org.springframework.batch.core.Step; +import org.springframework.batch.core.job.builder.JobBuilder; +import org.springframework.batch.core.repository.JobRepository; +import org.springframework.batch.core.step.builder.StepBuilder; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.transaction.PlatformTransactionManager; + +@Configuration +@RequiredArgsConstructor +public class CardSyncJobConfig { + + public static final String JOB_NAME = "cardSyncJob"; + public static final String STEP_NAME = "cardSyncStep"; + + private final JobRepository jobRepository; + private final PlatformTransactionManager transactionManager; + private final CardSyncTasklet cardSyncTasklet; + + @Bean(name = JOB_NAME) + public Job cardSyncJob() { + return new JobBuilder(JOB_NAME, jobRepository) + .start(cardSyncStep()) + .build(); + } + + @Bean(name = STEP_NAME) + public Step cardSyncStep() { + return new StepBuilder(STEP_NAME, jobRepository) + .tasklet(cardSyncTasklet, transactionManager) + .build(); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/cardsync/CardSyncTasklet.java b/src/main/java/com/rocketcrew/pocatbatch/job/cardsync/CardSyncTasklet.java new file mode 100644 index 0000000..a6672e0 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/cardsync/CardSyncTasklet.java @@ -0,0 +1,146 @@ +package com.rocketcrew.pocatbatch.job.cardsync; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.rocketcrew.pocatbatch.domain.card.entity.Card; +import com.rocketcrew.pocatbatch.domain.card.entity.enums.CardCategory; +import com.rocketcrew.pocatbatch.domain.card.entity.enums.CardGrade; +import com.rocketcrew.pocatbatch.domain.card.entity.enums.CardSource; +import com.rocketcrew.pocatbatch.domain.card.entity.enums.CardStatus; +import com.rocketcrew.pocatbatch.domain.card.repository.CardRepository; +import com.rocketcrew.pocatbatch.domain.outbox.service.OutboxEventWriter; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.dao.DataIntegrityViolationException; +import org.springframework.stereotype.Component; +import org.springframework.web.client.RestTemplate; + +import java.util.Optional; + +@Slf4j +@Component +@RequiredArgsConstructor +public class CardSyncTasklet implements Tasklet { + + private final CardRepository cardRepository; + private final OutboxEventWriter outboxEventWriter; + private final ObjectMapper objectMapper; + private final RestTemplate restTemplate; + + @Value("${pocat.batch.card-sync.admin-user-id:1}") + private Long adminUserId; + + private static final String TCGDEX_SETS_URL = "https://api.tcgdex.net/v2/en/sets"; + private static final String TCGDEX_SET_URL = "https://api.tcgdex.net/v2/en/sets/"; + private static final String TCGDEX_CARD_URL = "https://api.tcgdex.net/v2/en/cards/"; + + private static final CardGrade[] GRADES = CardGrade.values(); + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + log.info("[CardSync] 배치 동기화 시작 (adminUserId={})", adminUserId); + + int totalSynced = 0; + + try { + String setsJson = restTemplate.getForObject(TCGDEX_SETS_URL, String.class); + if (setsJson == null) { + log.warn("[CardSync] TCGdex 응답 null — 동기화 스킵"); + return RepeatStatus.FINISHED; + } + JsonNode setsArray = objectMapper.readTree(setsJson); + + for (JsonNode setNode : setsArray) { + String setId = setNode.path("id").asText(); + try { + totalSynced += syncSet(setId, totalSynced); + } catch (Exception e) { + log.warn("[CardSync] 세트 동기화 실패 ({}): {}", setId, e.getMessage()); + } + } + + } catch (Exception e) { + log.error("[CardSync] 배치 동기화 실패 — 스킵 처리", e); + } + + log.info("[CardSync] 배치 동기화 완료 — 신규 카드 총 {}개", totalSynced); + return RepeatStatus.FINISHED; + } + + private int syncSet(String setId, int offset) throws Exception { + String setJson = restTemplate.getForObject(TCGDEX_SET_URL + setId, String.class); + if (setJson == null) return 0; + JsonNode setRoot = objectMapper.readTree(setJson); + + String setName = setRoot.path("name").asText(""); + JsonNode cardNodes = setRoot.path("cards"); + + int synced = 0; + for (JsonNode cardNode : cardNodes) { + String tcgdexId = cardNode.path("id").asText(); + + Optional existing = cardRepository.findByTcgdexIdIncludingDeleted(tcgdexId); + if (existing.isPresent()) { + // soft-deleted 포함 이미 존재 → skip (soft-deleted여도 tcgdex_id 중복이므로 skip) + continue; + } + + try { + String cardJson = restTemplate.getForObject(TCGDEX_CARD_URL + tcgdexId, String.class); + JsonNode cardRoot = objectMapper.readTree(cardJson); + + String name = cardRoot.path("name").asText(); + String localId = cardRoot.path("localId").asText(); + String imageBase = cardRoot.path("image").asText(""); + String imageUrl = imageBase.isEmpty() ? null : imageBase + "/high.webp"; + String rarity = cardRoot.path("rarity").asText(""); + CardCategory category = parseCategory(cardRoot.path("category").asText("")); + CardGrade grade = GRADES[(offset + synced) % GRADES.length]; + + Card card = Card.builder() + .userId(adminUserId) + .tcgdexId(tcgdexId) + .name(name) + .seriesId(null) + .pokemonSetId(null) + .pokemonId(null) + .cardNumber(localId) + .rarity(rarity.isEmpty() ? "UNKNOWN" : rarity) + .category(category) + .grade(grade) + .imageUrl(imageUrl) + .source(CardSource.TCGDEX) + .status(CardStatus.ACTIVE) + .build(); + + try { + Card saved = cardRepository.save(card); + outboxEventWriter.write("card", String.valueOf(saved.getId()), "card.synced", saved.getId()); + synced++; + } catch (DataIntegrityViolationException e) { + log.warn("[CardSync] 중복 카드 스킵 (race): {}", tcgdexId); + } + + } catch (Exception e) { + log.warn("[CardSync] 카드 처리 실패 ({}): {}", tcgdexId, e.getMessage()); + } + } + + return synced; + } + + private CardCategory parseCategory(String category) { + if (category == null || category.isBlank()) return CardCategory.UNKNOWN; + return switch (category.toUpperCase()) { + case "POKEMON" -> CardCategory.POKEMON; + case "TRAINER", "TRAINERS" -> CardCategory.TRAINERS; + case "ENERGY" -> CardCategory.ENERGY; + default -> CardCategory.UNKNOWN; + }; + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxCleanupTasklet.java b/src/main/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxCleanupTasklet.java new file mode 100644 index 0000000..4ad5fe9 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxCleanupTasklet.java @@ -0,0 +1,38 @@ +package com.rocketcrew.pocatbatch.job.outboxrelay; + +import com.rocketcrew.pocatbatch.domain.outbox.enums.OutboxStatus; +import com.rocketcrew.pocatbatch.domain.outbox.repository.OutboxRepository; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Transactional; + +import java.time.LocalDateTime; + +@Slf4j +@Component +@RequiredArgsConstructor +public class OutboxCleanupTasklet implements Tasklet { + + private final OutboxRepository outboxRepository; + + @Override + @Transactional + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + try { + // SENT 상태이고 7일 이상 된 이벤트 삭제 + LocalDateTime cutoff = LocalDateTime.now().minusDays(7); + int deletedCount = outboxRepository.deleteOldEvents(OutboxStatus.SENT, cutoff); + + log.info("아웃박스 정리 완료: {} 개 오래된 이벤트 삭제", deletedCount); + return RepeatStatus.FINISHED; + } catch (Exception e) { + log.error("아웃박스 정리 실패", e); + throw new RuntimeException("아웃박스 정리 중 오류 발생", e); + } + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxReaperTasklet.java b/src/main/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxReaperTasklet.java new file mode 100644 index 0000000..87e7f5a --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxReaperTasklet.java @@ -0,0 +1,51 @@ +package com.rocketcrew.pocatbatch.job.outboxrelay; + +import com.rocketcrew.pocatbatch.domain.outbox.entity.OutboxEvent; +import com.rocketcrew.pocatbatch.domain.outbox.enums.OutboxStatus; +import com.rocketcrew.pocatbatch.domain.outbox.repository.OutboxRepository; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.stereotype.Component; + +import java.time.LocalDateTime; +import java.util.List; + +@Slf4j +@Component +@RequiredArgsConstructor +public class OutboxReaperTasklet implements Tasklet { + + private final OutboxRepository outboxRepository; + private static final int STUCK_MINUTES = 5; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + try { + LocalDateTime stuckBefore = LocalDateTime.now().minusMinutes(STUCK_MINUTES); + + List stuckEvents = outboxRepository.findTop100ByStatusAndProcessedAtBefore( + OutboxStatus.PROCESSING, stuckBefore + ); + + for (OutboxEvent event : stuckEvents) { + try { + event.resetToPending(); + outboxRepository.save(event); + log.warn("Outbox stuck 이벤트 PENDING 복구: id={}", event.getId()); + } catch (Exception e) { + log.error("Outbox stuck 이벤트 복구 실패: id={}, error={}", event.getId(), e.getMessage()); + } + } + + log.info("OutboxReaper 완료: {} 개 복구", stuckEvents.size()); + return RepeatStatus.FINISHED; + } catch (Exception e) { + log.error("아웃박스 Reaper 실패", e); + throw new RuntimeException("아웃박스 Reaper 중 오류 발생", e); + } + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxRelayJobConfig.java b/src/main/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxRelayJobConfig.java new file mode 100644 index 0000000..445e1b0 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxRelayJobConfig.java @@ -0,0 +1,61 @@ +package com.rocketcrew.pocatbatch.job.outboxrelay; + +import lombok.RequiredArgsConstructor; +import org.springframework.batch.core.Job; +import org.springframework.batch.core.Step; +import org.springframework.batch.core.job.builder.JobBuilder; +import org.springframework.batch.core.repository.JobRepository; +import org.springframework.batch.core.step.builder.StepBuilder; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.transaction.PlatformTransactionManager; + +@Configuration +@RequiredArgsConstructor +public class OutboxRelayJobConfig { + + public static final String JOB_NAME = "outboxRelayJob"; + public static final String RELAY_STEP_NAME = "outboxRelayStep"; + public static final String REAPER_STEP_NAME = "outboxReaperStep"; + public static final String CLEANUP_STEP_NAME = "outboxCleanupStep"; + + private final JobRepository jobRepository; + private final PlatformTransactionManager transactionManager; + private final OutboxRelayTasklet outboxRelayTasklet; + private final OutboxReaperTasklet outboxReaperTasklet; + private final OutboxCleanupTasklet outboxCleanupTasklet; + + @Bean(name = JOB_NAME) + public Job outboxRelayJob() { + return new JobBuilder(JOB_NAME, jobRepository) + .start(outboxRelayStep()) + .on("FAILED").to(outboxReaperStep()) + .from(outboxRelayStep()) + .on("*").to(outboxReaperStep()) + .from(outboxReaperStep()) + .on("*").to(outboxCleanupStep()) + .end() + .build(); + } + + @Bean(name = RELAY_STEP_NAME) + public Step outboxRelayStep() { + return new StepBuilder(RELAY_STEP_NAME, jobRepository) + .tasklet(outboxRelayTasklet, transactionManager) + .build(); + } + + @Bean(name = REAPER_STEP_NAME) + public Step outboxReaperStep() { + return new StepBuilder(REAPER_STEP_NAME, jobRepository) + .tasklet(outboxReaperTasklet, transactionManager) + .build(); + } + + @Bean(name = CLEANUP_STEP_NAME) + public Step outboxCleanupStep() { + return new StepBuilder(CLEANUP_STEP_NAME, jobRepository) + .tasklet(outboxCleanupTasklet, transactionManager) + .build(); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxRelayTasklet.java b/src/main/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxRelayTasklet.java new file mode 100644 index 0000000..a3bd498 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxRelayTasklet.java @@ -0,0 +1,71 @@ +package com.rocketcrew.pocatbatch.job.outboxrelay; + +import com.rocketcrew.pocatbatch.config.KafkaProducerConfig; +import com.rocketcrew.pocatbatch.domain.outbox.entity.OutboxEvent; +import com.rocketcrew.pocatbatch.domain.outbox.enums.OutboxStatus; +import com.rocketcrew.pocatbatch.domain.outbox.repository.OutboxRepository; +import com.rocketcrew.pocatbatch.domain.outbox.service.OutboxProcessor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.stereotype.Component; + +import java.time.LocalDateTime; +import java.util.List; + +@Slf4j +@Component +public class OutboxRelayTasklet implements Tasklet { + + private final OutboxRepository outboxRepository; + private final OutboxProcessor outboxProcessor; + private final KafkaTemplate kafkaTemplate; + private final KafkaTemplate paymentKafkaTemplate; + + public OutboxRelayTasklet( + OutboxRepository outboxRepository, + OutboxProcessor outboxProcessor, + @Qualifier("kafkaTemplate") KafkaTemplate kafkaTemplate, + @Qualifier("paymentKafkaTemplate") KafkaTemplate paymentKafkaTemplate) { + this.outboxRepository = outboxRepository; + this.outboxProcessor = outboxProcessor; + this.kafkaTemplate = kafkaTemplate; + this.paymentKafkaTemplate = paymentKafkaTemplate; + } + + @Value("${outbox.relay.age-threshold-seconds:10}") + private int ageThresholdSeconds; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + try { + LocalDateTime ageThreshold = LocalDateTime.now().minusSeconds(ageThresholdSeconds); + List pendingEvents = outboxRepository.findTop100ByStatusAndCreatedAtBeforeOrderByCreatedAtAsc( + OutboxStatus.PENDING, ageThreshold + ); + + for (OutboxEvent event : pendingEvents) { + try { + // 금융 관련 토픽은 paymentKafkaTemplate 사용 + KafkaTemplate template = KafkaProducerConfig.FINANCIAL_TOPICS.contains(event.getTopic()) + ? paymentKafkaTemplate + : kafkaTemplate; + outboxProcessor.processEvent(event, template); + } catch (Exception e) { + log.error("아웃박스 이벤트 처리 실패: id={}", event.getId(), e); + } + } + + log.info("아웃박스 릴레이 완료: {} 개 이벤트 처리", pendingEvents.size()); + return RepeatStatus.FINISHED; + } catch (Exception e) { + log.error("아웃박스 릴레이 실패", e); + throw new RuntimeException("아웃박스 릴레이 중 오류 발생", e); + } + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/refundretry/RefundRetryJobConfig.java b/src/main/java/com/rocketcrew/pocatbatch/job/refundretry/RefundRetryJobConfig.java new file mode 100644 index 0000000..12963e8 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/refundretry/RefundRetryJobConfig.java @@ -0,0 +1,37 @@ +package com.rocketcrew.pocatbatch.job.refundretry; + +import lombok.RequiredArgsConstructor; +import org.springframework.batch.core.Job; +import org.springframework.batch.core.Step; +import org.springframework.batch.core.job.builder.JobBuilder; +import org.springframework.batch.core.repository.JobRepository; +import org.springframework.batch.core.step.builder.StepBuilder; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.transaction.PlatformTransactionManager; + +@Configuration +@RequiredArgsConstructor +public class RefundRetryJobConfig { + + public static final String JOB_NAME = "refundRetryJob"; + public static final String STEP_NAME = "refundRetryStep"; + + private final JobRepository jobRepository; + private final PlatformTransactionManager transactionManager; + private final RefundRetryTasklet refundRetryTasklet; + + @Bean(name = JOB_NAME) + public Job refundRetryJob() { + return new JobBuilder(JOB_NAME, jobRepository) + .start(refundRetryStep()) + .build(); + } + + @Bean(name = STEP_NAME) + public Step refundRetryStep() { + return new StepBuilder(STEP_NAME, jobRepository) + .tasklet(refundRetryTasklet, transactionManager) + .build(); + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/job/refundretry/RefundRetryTasklet.java b/src/main/java/com/rocketcrew/pocatbatch/job/refundretry/RefundRetryTasklet.java new file mode 100644 index 0000000..38be922 --- /dev/null +++ b/src/main/java/com/rocketcrew/pocatbatch/job/refundretry/RefundRetryTasklet.java @@ -0,0 +1,56 @@ +package com.rocketcrew.pocatbatch.job.refundretry; + +import com.rocketcrew.pocatbatch.client.MainAppRefundClient; +import com.rocketcrew.pocatbatch.domain.refund.entity.Refund; +import com.rocketcrew.pocatbatch.domain.refund.entity.RefundStatus; +import com.rocketcrew.pocatbatch.domain.refund.repository.RefundRepository; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.batch.core.StepContribution; +import org.springframework.batch.core.scope.context.ChunkContext; +import org.springframework.batch.core.step.tasklet.Tasklet; +import org.springframework.batch.repeat.RepeatStatus; +import org.springframework.stereotype.Component; + +import java.time.LocalDateTime; +import java.util.List; + +@Slf4j +@Component +@RequiredArgsConstructor +public class RefundRetryTasklet implements Tasklet { + + private final RefundRepository refundRepository; + private final MainAppRefundClient mainAppRefundClient; + + @Override + public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) { + try { + LocalDateTime now = LocalDateTime.now(); + LocalDateTime stuckBefore = now.minusMinutes(5); + + List retryableRefunds = refundRepository.findRetryableTargets( + RefundStatus.FAILED_RETRYABLE, now, + RefundStatus.PROCESSING, stuckBefore + ); + + int retriedCount = 0; + long jobExecutionId = chunkContext.getStepContext().getStepExecution().getJobExecution().getId(); + + for (Refund refund : retryableRefunds) { + try { + mainAppRefundClient.retryRefund(refund.getId(), jobExecutionId); + retriedCount++; + } catch (Exception e) { + log.warn("환불 재시도 실패: refundId={}", refund.getId(), e); + } + } + + log.info("환불 재시도 완료: {} 개 환불 처리됨", retriedCount); + return RepeatStatus.FINISHED; + } catch (Exception e) { + log.error("환불 재시도 작업 실패", e); + throw new RuntimeException("환불 재시도 중 오류 발생", e); + } + } +} diff --git a/src/main/java/com/rocketcrew/pocatbatch/scheduler/BatchScheduler.java b/src/main/java/com/rocketcrew/pocatbatch/scheduler/BatchScheduler.java index 40ede18..9f8f829 100644 --- a/src/main/java/com/rocketcrew/pocatbatch/scheduler/BatchScheduler.java +++ b/src/main/java/com/rocketcrew/pocatbatch/scheduler/BatchScheduler.java @@ -1,6 +1,14 @@ package com.rocketcrew.pocatbatch.scheduler; +import com.rocketcrew.pocatbatch.job.aisession.AiSessionCleanupJobConfig; +import com.rocketcrew.pocatbatch.job.auctionactivation.AuctionActivationJobConfig; +import com.rocketcrew.pocatbatch.job.auctionexpiration.AuctionExpirationJobConfig; +import com.rocketcrew.pocatbatch.job.auctionranking.AuctionRankingJobConfig; +import com.rocketcrew.pocatbatch.job.buyoutrecovery.BuyoutRecoveryJobConfig; +import com.rocketcrew.pocatbatch.job.cardsync.CardSyncJobConfig; +import com.rocketcrew.pocatbatch.job.outboxrelay.OutboxRelayJobConfig; import com.rocketcrew.pocatbatch.job.ranking.FreePostRankingJobConfig; +import com.rocketcrew.pocatbatch.job.refundretry.RefundRetryJobConfig; import com.rocketcrew.pocatbatch.job.viewcount.ViewCountFlushJobConfig; import lombok.extern.slf4j.Slf4j; import org.springframework.batch.core.BatchStatus; @@ -10,25 +18,51 @@ import org.springframework.batch.core.JobParametersBuilder; import org.springframework.batch.core.launch.JobLauncher; import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; @Slf4j @Component +@ConditionalOnProperty(name = "pocat.batch.scheduler.enabled", havingValue = "true", matchIfMissing = true) public class BatchScheduler { private final JobLauncher jobLauncher; private final Job freePostRankingJob; private final Job viewCountFlushJob; + private final Job aiSessionCleanupJob; + private final Job auctionRankingJob; + private final Job outboxRelayJob; + private final Job cardSyncJob; + private final Job auctionActivationJob; + private final Job auctionExpirationJob; + private final Job buyoutRecoveryJob; + private final Job refundRetryJob; // @RequiredArgsConstructor는 @Qualifier 미지원 → 수동 생성자 필수 public BatchScheduler( JobLauncher jobLauncher, @Qualifier(FreePostRankingJobConfig.JOB_NAME) Job freePostRankingJob, - @Qualifier(ViewCountFlushJobConfig.JOB_NAME) Job viewCountFlushJob) { - this.jobLauncher = jobLauncher; - this.freePostRankingJob = freePostRankingJob; - this.viewCountFlushJob = viewCountFlushJob; + @Qualifier(ViewCountFlushJobConfig.JOB_NAME) Job viewCountFlushJob, + @Qualifier(AiSessionCleanupJobConfig.JOB_NAME) Job aiSessionCleanupJob, + @Qualifier(AuctionRankingJobConfig.JOB_NAME) Job auctionRankingJob, + @Qualifier(OutboxRelayJobConfig.JOB_NAME) Job outboxRelayJob, + @Qualifier(CardSyncJobConfig.JOB_NAME) Job cardSyncJob, + @Qualifier(AuctionActivationJobConfig.JOB_NAME) Job auctionActivationJob, + @Qualifier(AuctionExpirationJobConfig.JOB_NAME) Job auctionExpirationJob, + @Qualifier(BuyoutRecoveryJobConfig.JOB_NAME) Job buyoutRecoveryJob, + @Qualifier(RefundRetryJobConfig.JOB_NAME) Job refundRetryJob) { + this.jobLauncher = jobLauncher; + this.freePostRankingJob = freePostRankingJob; + this.viewCountFlushJob = viewCountFlushJob; + this.aiSessionCleanupJob = aiSessionCleanupJob; + this.auctionRankingJob = auctionRankingJob; + this.outboxRelayJob = outboxRelayJob; + this.cardSyncJob = cardSyncJob; + this.auctionActivationJob = auctionActivationJob; + this.auctionExpirationJob = auctionExpirationJob; + this.buyoutRecoveryJob = buyoutRecoveryJob; + this.refundRetryJob = refundRetryJob; } @Scheduled(fixedDelay = 60_000) @@ -41,6 +75,46 @@ public void runViewCountFlush() { launch(viewCountFlushJob, "viewCountFlushJob"); } + @Scheduled(fixedRate = 300_000) + public void runAiSessionCleanup() { + launch(aiSessionCleanupJob, "aiSessionCleanupJob"); + } + + @Scheduled(fixedDelay = 60_000) + public void runAuctionRanking() { + launch(auctionRankingJob, "auctionRankingJob"); + } + + @Scheduled(fixedDelay = 5_000) + public void runOutboxRelay() { + launch(outboxRelayJob, "outboxRelayJob"); + } + + @Scheduled(cron = "0 0 0 * * SUN", zone = "Asia/Seoul") + public void runCardSync() { + launch(cardSyncJob, "cardSyncJob"); + } + + @Scheduled(cron = "0 0 19 * * *", zone = "Asia/Seoul") + public void runAuctionActivation() { + launch(auctionActivationJob, "auctionActivationJob"); + } + + @Scheduled(cron = "0 5-30 19 * * *", zone = "Asia/Seoul") + public void runAuctionExpiration() { + launch(auctionExpirationJob, "auctionExpirationJob"); + } + + @Scheduled(fixedDelay = 60_000) + public void runBuyoutRecovery() { + launch(buyoutRecoveryJob, "buyoutRecoveryJob"); + } + + @Scheduled(fixedDelay = 60_000) + public void runRefundRetry() { + launch(refundRetryJob, "refundRetryJob"); + } + private void launch(Job job, String label) { try { JobParameters params = new JobParametersBuilder() diff --git a/src/main/resources/application.yaml b/src/main/resources/application.yaml index 0f1a9f2..98b44c2 100644 --- a/src/main/resources/application.yaml +++ b/src/main/resources/application.yaml @@ -23,6 +23,9 @@ spring: port: ${REDIS_PORT:6379} password: ${REDIS_PASSWORD:} + kafka: + bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:localhost:9092} + batch: jdbc: initialize-schema: always @@ -39,6 +42,9 @@ management: include: health,info,metrics pocat: + main-app: + base-url: ${POCAT_API_BASE_URL:http://localhost:8080} + internal-token: ${POCAT_INTERNAL_TOKEN} batch: ranking: free: @@ -46,3 +52,9 @@ pocat: ttl-seconds: 70 popular-days: 7 comment-weight: 3 + card-sync: + admin-user-id: ${CARD_SYNC_ADMIN_USER_ID:1} + +outbox: + relay: + age-threshold-seconds: ${OUTBOX_AGE_THRESHOLD_SECONDS:10} diff --git a/src/test/java/com/rocketcrew/pocatbatch/config/BatchTestConfig.java b/src/test/java/com/rocketcrew/pocatbatch/config/BatchTestConfig.java new file mode 100644 index 0000000..306bbdd --- /dev/null +++ b/src/test/java/com/rocketcrew/pocatbatch/config/BatchTestConfig.java @@ -0,0 +1,57 @@ +package com.rocketcrew.pocatbatch.config; + +import org.mockito.Mockito; +import org.redisson.api.RLock; +import org.redisson.api.RedissonClient; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Primary; +import org.springframework.context.annotation.Profile; +import org.springframework.data.redis.core.StringRedisTemplate; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.web.client.RestTemplate; + +@Configuration +@Profile("test") +public class BatchTestConfig { + + @Bean + @Primary + public RedissonClient redissonClient() { + RedissonClient mock = Mockito.mock(RedissonClient.class); + RLock mockLock = Mockito.mock(RLock.class); + try { + Mockito.when(mock.getLock(Mockito.anyString())).thenReturn(mockLock); + Mockito.when(mockLock.tryLock(Mockito.anyLong(), Mockito.anyLong(), Mockito.any())).thenReturn(true); + Mockito.when(mockLock.isHeldByCurrentThread()).thenReturn(true); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new RuntimeException(e); + } + return mock; + } + + @Bean("kafkaTemplate") + @Primary + public KafkaTemplate kafkaTemplate() { + return Mockito.mock(KafkaTemplate.class); + } + + @Bean("paymentKafkaTemplate") + @Primary + public KafkaTemplate paymentKafkaTemplate() { + return Mockito.mock(KafkaTemplate.class); + } + + @Bean + @Primary + public RestTemplate restTemplate() { + return Mockito.mock(RestTemplate.class); + } + + @Bean + @Primary + public StringRedisTemplate stringRedisTemplate() { + return Mockito.mock(StringRedisTemplate.class); + } +} diff --git a/src/test/java/com/rocketcrew/pocatbatch/job/aisession/AiSessionCleanupJobConfigTest.java b/src/test/java/com/rocketcrew/pocatbatch/job/aisession/AiSessionCleanupJobConfigTest.java new file mode 100644 index 0000000..2ee4285 --- /dev/null +++ b/src/test/java/com/rocketcrew/pocatbatch/job/aisession/AiSessionCleanupJobConfigTest.java @@ -0,0 +1,38 @@ +package com.rocketcrew.pocatbatch.job.aisession; + +import org.junit.jupiter.api.Test; +import org.springframework.batch.core.*; +import org.springframework.batch.test.JobLauncherTestUtils; +import org.springframework.batch.test.context.SpringBatchTest; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.ActiveProfiles; + +import static org.assertj.core.api.Assertions.assertThat; + +@SpringBatchTest +@SpringBootTest +@ActiveProfiles("test") +class AiSessionCleanupJobConfigTest { + + @Autowired + private JobLauncherTestUtils jobLauncherTestUtils; + + @Autowired + @Qualifier(AiSessionCleanupJobConfig.JOB_NAME) + private Job aiSessionCleanupJob; + + @Test + void aiSessionCleanupJob_실행_성공() throws Exception { + jobLauncherTestUtils.setJob(aiSessionCleanupJob); + + JobExecution jobExecution = jobLauncherTestUtils.launchJob( + new JobParametersBuilder() + .addLong("ts", System.currentTimeMillis()) + .toJobParameters()); + + assertThat(jobExecution.getStatus()).isEqualTo(BatchStatus.COMPLETED); + assertThat(jobExecution.getExitStatus().getExitCode()).isEqualTo("COMPLETED"); + } +} diff --git a/src/test/java/com/rocketcrew/pocatbatch/job/auctionactivation/AuctionActivationJobConfigTest.java b/src/test/java/com/rocketcrew/pocatbatch/job/auctionactivation/AuctionActivationJobConfigTest.java new file mode 100644 index 0000000..5cba2b7 --- /dev/null +++ b/src/test/java/com/rocketcrew/pocatbatch/job/auctionactivation/AuctionActivationJobConfigTest.java @@ -0,0 +1,38 @@ +package com.rocketcrew.pocatbatch.job.auctionactivation; + +import org.junit.jupiter.api.Test; +import org.springframework.batch.core.*; +import org.springframework.batch.test.JobLauncherTestUtils; +import org.springframework.batch.test.context.SpringBatchTest; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.ActiveProfiles; + +import static org.assertj.core.api.Assertions.assertThat; + +@SpringBatchTest +@SpringBootTest +@ActiveProfiles("test") +class AuctionActivationJobConfigTest { + + @Autowired + private JobLauncherTestUtils jobLauncherTestUtils; + + @Autowired + @Qualifier(AuctionActivationJobConfig.JOB_NAME) + private Job auctionActivationJob; + + @Test + void auctionActivationJob_실행_성공() throws Exception { + jobLauncherTestUtils.setJob(auctionActivationJob); + + JobExecution jobExecution = jobLauncherTestUtils.launchJob( + new JobParametersBuilder() + .addLong("ts", System.currentTimeMillis()) + .toJobParameters()); + + assertThat(jobExecution.getStatus()).isEqualTo(BatchStatus.COMPLETED); + assertThat(jobExecution.getExitStatus().getExitCode()).isEqualTo("COMPLETED"); + } +} diff --git a/src/test/java/com/rocketcrew/pocatbatch/job/auctionexpiration/AuctionExpirationJobConfigTest.java b/src/test/java/com/rocketcrew/pocatbatch/job/auctionexpiration/AuctionExpirationJobConfigTest.java new file mode 100644 index 0000000..2566b79 --- /dev/null +++ b/src/test/java/com/rocketcrew/pocatbatch/job/auctionexpiration/AuctionExpirationJobConfigTest.java @@ -0,0 +1,38 @@ +package com.rocketcrew.pocatbatch.job.auctionexpiration; + +import org.junit.jupiter.api.Test; +import org.springframework.batch.core.*; +import org.springframework.batch.test.JobLauncherTestUtils; +import org.springframework.batch.test.context.SpringBatchTest; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.ActiveProfiles; + +import static org.assertj.core.api.Assertions.assertThat; + +@SpringBatchTest +@SpringBootTest +@ActiveProfiles("test") +class AuctionExpirationJobConfigTest { + + @Autowired + private JobLauncherTestUtils jobLauncherTestUtils; + + @Autowired + @Qualifier(AuctionExpirationJobConfig.JOB_NAME) + private Job auctionExpirationJob; + + @Test + void auctionExpirationJob_실행_성공() throws Exception { + jobLauncherTestUtils.setJob(auctionExpirationJob); + + JobExecution jobExecution = jobLauncherTestUtils.launchJob( + new JobParametersBuilder() + .addLong("ts", System.currentTimeMillis()) + .toJobParameters()); + + assertThat(jobExecution.getStatus()).isEqualTo(BatchStatus.COMPLETED); + assertThat(jobExecution.getExitStatus().getExitCode()).isEqualTo("COMPLETED"); + } +} diff --git a/src/test/java/com/rocketcrew/pocatbatch/job/auctionranking/AuctionRankingJobConfigTest.java b/src/test/java/com/rocketcrew/pocatbatch/job/auctionranking/AuctionRankingJobConfigTest.java new file mode 100644 index 0000000..74a55c1 --- /dev/null +++ b/src/test/java/com/rocketcrew/pocatbatch/job/auctionranking/AuctionRankingJobConfigTest.java @@ -0,0 +1,38 @@ +package com.rocketcrew.pocatbatch.job.auctionranking; + +import org.junit.jupiter.api.Test; +import org.springframework.batch.core.*; +import org.springframework.batch.test.JobLauncherTestUtils; +import org.springframework.batch.test.context.SpringBatchTest; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.ActiveProfiles; + +import static org.assertj.core.api.Assertions.assertThat; + +@SpringBatchTest +@SpringBootTest +@ActiveProfiles("test") +class AuctionRankingJobConfigTest { + + @Autowired + private JobLauncherTestUtils jobLauncherTestUtils; + + @Autowired + @Qualifier(AuctionRankingJobConfig.JOB_NAME) + private Job auctionRankingJob; + + @Test + void auctionRankingJob_실행_성공() throws Exception { + jobLauncherTestUtils.setJob(auctionRankingJob); + + JobExecution jobExecution = jobLauncherTestUtils.launchJob( + new JobParametersBuilder() + .addLong("ts", System.currentTimeMillis()) + .toJobParameters()); + + assertThat(jobExecution.getStatus()).isEqualTo(BatchStatus.COMPLETED); + assertThat(jobExecution.getExitStatus().getExitCode()).isEqualTo("COMPLETED"); + } +} diff --git a/src/test/java/com/rocketcrew/pocatbatch/job/buyoutrecovery/BuyoutRecoveryJobConfigTest.java b/src/test/java/com/rocketcrew/pocatbatch/job/buyoutrecovery/BuyoutRecoveryJobConfigTest.java new file mode 100644 index 0000000..65f5e39 --- /dev/null +++ b/src/test/java/com/rocketcrew/pocatbatch/job/buyoutrecovery/BuyoutRecoveryJobConfigTest.java @@ -0,0 +1,38 @@ +package com.rocketcrew.pocatbatch.job.buyoutrecovery; + +import org.junit.jupiter.api.Test; +import org.springframework.batch.core.*; +import org.springframework.batch.test.JobLauncherTestUtils; +import org.springframework.batch.test.context.SpringBatchTest; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.ActiveProfiles; + +import static org.assertj.core.api.Assertions.assertThat; + +@SpringBatchTest +@SpringBootTest +@ActiveProfiles("test") +class BuyoutRecoveryJobConfigTest { + + @Autowired + private JobLauncherTestUtils jobLauncherTestUtils; + + @Autowired + @Qualifier(BuyoutRecoveryJobConfig.JOB_NAME) + private Job buyoutRecoveryJob; + + @Test + void buyoutRecoveryJob_실행_성공() throws Exception { + jobLauncherTestUtils.setJob(buyoutRecoveryJob); + + JobExecution jobExecution = jobLauncherTestUtils.launchJob( + new JobParametersBuilder() + .addLong("ts", System.currentTimeMillis()) + .toJobParameters()); + + assertThat(jobExecution.getStatus()).isEqualTo(BatchStatus.COMPLETED); + assertThat(jobExecution.getExitStatus().getExitCode()).isEqualTo("COMPLETED"); + } +} diff --git a/src/test/java/com/rocketcrew/pocatbatch/job/cardsync/CardSyncJobConfigTest.java b/src/test/java/com/rocketcrew/pocatbatch/job/cardsync/CardSyncJobConfigTest.java new file mode 100644 index 0000000..2ba6150 --- /dev/null +++ b/src/test/java/com/rocketcrew/pocatbatch/job/cardsync/CardSyncJobConfigTest.java @@ -0,0 +1,38 @@ +package com.rocketcrew.pocatbatch.job.cardsync; + +import org.junit.jupiter.api.Test; +import org.springframework.batch.core.*; +import org.springframework.batch.test.JobLauncherTestUtils; +import org.springframework.batch.test.context.SpringBatchTest; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.ActiveProfiles; + +import static org.assertj.core.api.Assertions.assertThat; + +@SpringBatchTest +@SpringBootTest +@ActiveProfiles("test") +class CardSyncJobConfigTest { + + @Autowired + private JobLauncherTestUtils jobLauncherTestUtils; + + @Autowired + @Qualifier(CardSyncJobConfig.JOB_NAME) + private Job cardSyncJob; + + @Test + void cardSyncJob_실행_성공() throws Exception { + jobLauncherTestUtils.setJob(cardSyncJob); + + JobExecution jobExecution = jobLauncherTestUtils.launchJob( + new JobParametersBuilder() + .addLong("ts", System.currentTimeMillis()) + .toJobParameters()); + + assertThat(jobExecution.getStatus()).isEqualTo(BatchStatus.COMPLETED); + assertThat(jobExecution.getExitStatus().getExitCode()).isEqualTo("COMPLETED"); + } +} diff --git a/src/test/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxRelayJobConfigTest.java b/src/test/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxRelayJobConfigTest.java new file mode 100644 index 0000000..f8e4c4d --- /dev/null +++ b/src/test/java/com/rocketcrew/pocatbatch/job/outboxrelay/OutboxRelayJobConfigTest.java @@ -0,0 +1,38 @@ +package com.rocketcrew.pocatbatch.job.outboxrelay; + +import org.junit.jupiter.api.Test; +import org.springframework.batch.core.*; +import org.springframework.batch.test.JobLauncherTestUtils; +import org.springframework.batch.test.context.SpringBatchTest; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.ActiveProfiles; + +import static org.assertj.core.api.Assertions.assertThat; + +@SpringBatchTest +@SpringBootTest +@ActiveProfiles("test") +class OutboxRelayJobConfigTest { + + @Autowired + private JobLauncherTestUtils jobLauncherTestUtils; + + @Autowired + @Qualifier(OutboxRelayJobConfig.JOB_NAME) + private Job outboxRelayJob; + + @Test + void outboxRelayJob_실행_성공() throws Exception { + jobLauncherTestUtils.setJob(outboxRelayJob); + + JobExecution jobExecution = jobLauncherTestUtils.launchJob( + new JobParametersBuilder() + .addLong("ts", System.currentTimeMillis()) + .toJobParameters()); + + assertThat(jobExecution.getStatus()).isEqualTo(BatchStatus.COMPLETED); + assertThat(jobExecution.getExitStatus().getExitCode()).isEqualTo("COMPLETED"); + } +} diff --git a/src/test/java/com/rocketcrew/pocatbatch/job/refundretry/RefundRetryJobConfigTest.java b/src/test/java/com/rocketcrew/pocatbatch/job/refundretry/RefundRetryJobConfigTest.java new file mode 100644 index 0000000..b235829 --- /dev/null +++ b/src/test/java/com/rocketcrew/pocatbatch/job/refundretry/RefundRetryJobConfigTest.java @@ -0,0 +1,38 @@ +package com.rocketcrew.pocatbatch.job.refundretry; + +import org.junit.jupiter.api.Test; +import org.springframework.batch.core.*; +import org.springframework.batch.test.JobLauncherTestUtils; +import org.springframework.batch.test.context.SpringBatchTest; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.ActiveProfiles; + +import static org.assertj.core.api.Assertions.assertThat; + +@SpringBatchTest +@SpringBootTest +@ActiveProfiles("test") +class RefundRetryJobConfigTest { + + @Autowired + private JobLauncherTestUtils jobLauncherTestUtils; + + @Autowired + @Qualifier(RefundRetryJobConfig.JOB_NAME) + private Job refundRetryJob; + + @Test + void refundRetryJob_실행_성공() throws Exception { + jobLauncherTestUtils.setJob(refundRetryJob); + + JobExecution jobExecution = jobLauncherTestUtils.launchJob( + new JobParametersBuilder() + .addLong("ts", System.currentTimeMillis()) + .toJobParameters()); + + assertThat(jobExecution.getStatus()).isEqualTo(BatchStatus.COMPLETED); + assertThat(jobExecution.getExitStatus().getExitCode()).isEqualTo("COMPLETED"); + } +} diff --git a/src/test/resources/application.yaml b/src/test/resources/application.yaml index 1383be4..7a670c5 100644 --- a/src/test/resources/application.yaml +++ b/src/test/resources/application.yaml @@ -15,8 +15,37 @@ spring: host: localhost port: 6379 + kafka: + bootstrap-servers: localhost:9092 + + autoconfigure: + exclude: + - org.springframework.boot.autoconfigure.data.redis.RedisAutoConfiguration + - org.springframework.boot.autoconfigure.data.redis.RedisRepositoriesAutoConfiguration + - org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration + batch: jdbc: initialize-schema: always job: enabled: false + +pocat: + main-app: + base-url: http://localhost:8080 + internal-token: test-token + batch: + scheduler: + enabled: false + card-sync: + admin-user-id: 1 + ranking: + free: + cache-size: 100 + ttl-seconds: 70 + popular-days: 7 + comment-weight: 3 + +outbox: + relay: + age-threshold-seconds: 10