Skip to content

Commit 35de6d3

Browse files
authored
Merge pull request #381 from BINBEAN777/volume-8
Volume 8
2 parents 7d3cfa4 + 8746cac commit 35de6d3

21 files changed

Lines changed: 975 additions & 4 deletions
Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,64 @@
1+
package com.loopers.application.queue;
2+
3+
import com.loopers.domain.queue.EntryTokenStore;
4+
import com.loopers.domain.queue.WaitingQueueRepository;
5+
import com.loopers.support.error.CoreException;
6+
import com.loopers.support.error.ErrorType;
7+
import lombok.RequiredArgsConstructor;
8+
import org.springframework.beans.factory.annotation.Value;
9+
import org.springframework.stereotype.Component;
10+
11+
import java.util.Optional;
12+
13+
@RequiredArgsConstructor
14+
@Component
15+
public class QueueFacade {
16+
17+
private final WaitingQueueRepository waitingQueueRepository;
18+
private final EntryTokenStore entryTokenStore;
19+
20+
@Value("${loopers.queue.gate.enabled:false}") // 관문 토글 (기본 꺼짐)
21+
private boolean gateEnabled;
22+
23+
/** 대기열에 진입시키고, 현재 순번/예상 대기시간을 돌려준다. */
24+
public QueueInfo enter(Long userId) {
25+
waitingQueueRepository.enter(userId);
26+
long rank = waitingQueueRepository.findRank(userId).orElse(0L);
27+
return QueueInfo.of(rank);
28+
}
29+
30+
/** 현재 순번/예상 대기시간을 조회한다. 큐에 없지만 토큰이 있으면 '내 차례', 둘 다 없으면 404. */
31+
public QueueInfo position(Long userId) {
32+
Optional<Long> rank = waitingQueueRepository.findRank(userId);
33+
if (rank.isPresent()) {
34+
return QueueInfo.of(rank.get()); // 아직 대기 중
35+
}
36+
// 큐에 없음 → 토큰이 있으면 내 차례(입장 허가), 없으면 이탈/미진입(404)
37+
return entryTokenStore.find(userId)
38+
.map(QueueInfo::admitted)
39+
.orElseThrow(() -> new CoreException(ErrorType.NOT_FOUND, "대기열에 없습니다."));
40+
}
41+
42+
/** 주문 API 진입 관문: 토큰을 검증한다. 관문이 꺼져 있으면 그냥 통과. */
43+
public void validateEntry(Long userId, String token) {
44+
if (!gateEnabled) {
45+
return; // 관문 OFF → 통과 (평소 / Redis 장애 시 bypass)
46+
}
47+
if (token == null || token.isBlank()) {
48+
throw new CoreException(ErrorType.FORBIDDEN, "입장 토큰이 필요합니다.");
49+
}
50+
String stored = entryTokenStore.find(userId)
51+
.orElseThrow(() -> new CoreException(ErrorType.FORBIDDEN, "유효한 입장 토큰이 없습니다."));
52+
if (!stored.equals(token)) {
53+
throw new CoreException(ErrorType.FORBIDDEN, "입장 토큰이 일치하지 않습니다.");
54+
}
55+
}
56+
57+
/** 주문 성공 후 토큰을 소진(삭제)한다. 관문이 꺼져 있으면 아무 것도 하지 않는다. */
58+
public void consumeEntry(Long userId) {
59+
if (!gateEnabled) {
60+
return;
61+
}
62+
entryTokenStore.remove(userId);
63+
}
64+
}
Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
package com.loopers.application.queue;
2+
3+
import java.util.Optional;
4+
5+
public record QueueInfo(long position, long estimatedWaitSeconds, Optional<String> token) {
6+
7+
// Phase 0 처리량 산정 근거: 커넥션풀 50 / 200ms → 250 TPS → 안전마진 70% → 175 TPS
8+
private static final long THROUGHPUT_PER_SEC = 175;
9+
10+
/** 대기 중: 순번으로 예상시간 계산, 토큰 없음. */
11+
public static QueueInfo of(long rank) {
12+
return new QueueInfo(rank, rank / THROUGHPUT_PER_SEC, Optional.empty());
13+
}
14+
15+
/** 입장 허가(내 차례): 순번 0, 토큰 발급됨. */
16+
public static QueueInfo admitted(String token) {
17+
return new QueueInfo(0L, 0L, Optional.of(token));
18+
}
19+
}
Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,43 @@
1+
package com.loopers.application.queue;
2+
3+
import com.loopers.domain.queue.EntryTokenStore;
4+
import com.loopers.domain.queue.WaitingQueueRepository;
5+
import lombok.RequiredArgsConstructor;
6+
import lombok.extern.slf4j.Slf4j;
7+
import org.springframework.beans.factory.annotation.Value;
8+
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
9+
import org.springframework.scheduling.annotation.Scheduled;
10+
import org.springframework.stereotype.Component;
11+
12+
import java.util.List;
13+
14+
@Slf4j
15+
@RequiredArgsConstructor
16+
@Component
17+
@ConditionalOnProperty(prefix = "loopers.queue.scheduler", name = "enabled", havingValue = "true", matchIfMissing = true)
18+
public class TokenIssueScheduler {
19+
20+
private final WaitingQueueRepository waitingQueueRepository;
21+
private final EntryTokenStore entryTokenStore;
22+
23+
@Value("${loopers.queue.scheduler.batch-size:18}")
24+
private int batchSize;
25+
26+
/**
27+
* 100ms 마다 대기열 앞 N명을 꺼내 입장 토큰을 발급한다.
28+
* (Thundering Herd 완화: 1초에 몰지 않고 100ms 단위로 ~18명씩 분산)
29+
*/
30+
@Scheduled(fixedDelayString = "${loopers.queue.scheduler.interval-ms:100}")
31+
public void issueTokens() {
32+
try {
33+
List<Long> admitted = waitingQueueRepository.popMin(batchSize);
34+
for (Long userId : admitted) {
35+
entryTokenStore.issue(userId);
36+
}
37+
} catch (Exception e) {
38+
// ★ @Scheduled 메서드에서 예외가 밖으로 나가면 이후 실행이 조용히 취소된다.
39+
// 반드시 여기서 삼켜서 스케줄러가 계속 돌게 한다.
40+
log.error("토큰 발급 스케줄러 실행 중 오류", e);
41+
}
42+
}
43+
}
Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,15 @@
1+
package com.loopers.domain.queue;
2+
3+
import java.util.Optional;
4+
5+
public interface EntryTokenStore {
6+
7+
/** userId 에게 입장 토큰을 발급하고, 발급된 토큰 값을 반환한다. TTL 후 자동 만료. */
8+
String issue(Long userId);
9+
10+
/** userId 의 현재 유효한 토큰을 조회한다. 없거나 만료됐으면 empty. */
11+
Optional<String> find(Long userId);
12+
13+
/** 사용 완료된 토큰을 삭제한다. */
14+
void remove(Long userId);
15+
}
Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
package com.loopers.domain.queue;
2+
3+
import java.util.List;
4+
import java.util.Optional;
5+
6+
public interface WaitingQueueRepository {
7+
8+
/** 대기열에 진입한다. 이미 대기 중이면 순번을 유지한다(중복 진입 무시). */
9+
void enter(Long userId);
10+
11+
/** 현재 순번(0-based). 대기열에 없으면 empty. */
12+
Optional<Long> findRank(Long userId);
13+
14+
/** 현재 대기 중인 전체 인원. */
15+
long size();
16+
17+
/** 대기열 맨 앞 N명을 꺼낸다(제거 + 반환). 스케줄러가 입장 처리에 사용. */
18+
List<Long> popMin(int count);
19+
}
Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,51 @@
1+
package com.loopers.infrastructure.queue;
2+
3+
import com.loopers.config.redis.RedisConfig;
4+
import com.loopers.domain.queue.EntryTokenStore;
5+
import org.springframework.beans.factory.annotation.Qualifier;
6+
import org.springframework.beans.factory.annotation.Value;
7+
import org.springframework.data.redis.core.RedisTemplate;
8+
import org.springframework.stereotype.Repository;
9+
10+
import java.time.Duration;
11+
import java.util.Optional;
12+
import java.util.UUID;
13+
14+
@Repository
15+
public class RedisEntryTokenStore implements EntryTokenStore {
16+
17+
private static final String TOKEN_KEY_PREFIX = "queue:token:";
18+
19+
private final RedisTemplate<String, String> redisTemplate;
20+
private final Duration ttl;
21+
22+
public RedisEntryTokenStore(
23+
@Qualifier(RedisConfig.REDIS_TEMPLATE_MASTER) RedisTemplate<String, String> redisTemplate,
24+
@Value("${loopers.queue.token.ttl-seconds:300}") long ttlSeconds // 기본 5분, 오버라이드 가능
25+
) {
26+
this.redisTemplate = redisTemplate;
27+
this.ttl = Duration.ofSeconds(ttlSeconds);
28+
}
29+
30+
@Override
31+
public String issue(Long userId) {
32+
String token = UUID.randomUUID().toString();
33+
// SET queue:token:{userId} {token} EX {ttl} → TTL 지나면 Redis 가 알아서 삭제
34+
redisTemplate.opsForValue().set(key(userId), token, ttl);
35+
return token;
36+
}
37+
38+
@Override
39+
public Optional<String> find(Long userId) {
40+
return Optional.ofNullable(redisTemplate.opsForValue().get(key(userId)));
41+
}
42+
43+
@Override
44+
public void remove(Long userId) {
45+
redisTemplate.delete(key(userId));
46+
}
47+
48+
private String key(Long userId) {
49+
return TOKEN_KEY_PREFIX + userId;
50+
}
51+
}
Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
1+
package com.loopers.infrastructure.queue;
2+
3+
import com.loopers.config.redis.RedisConfig;
4+
import com.loopers.domain.queue.WaitingQueueRepository;
5+
import org.springframework.beans.factory.annotation.Qualifier;
6+
import org.springframework.data.redis.core.RedisTemplate;
7+
import org.springframework.data.redis.core.ZSetOperations;
8+
import org.springframework.stereotype.Repository;
9+
10+
import java.util.List;
11+
import java.util.Objects;
12+
import java.util.Optional;
13+
import java.util.Set;
14+
15+
@Repository
16+
public class RedisWaitingQueueRepository implements WaitingQueueRepository {
17+
18+
private static final String WAITING_QUEUE_KEY = "queue:waiting";
19+
20+
// 순서(순번)는 단일 기준점에서만 정해야 하므로, 쓰기/읽기 모두 마스터 템플릿을 쓴다.
21+
// (Replica 로 읽으면 복제 지연 때문에 순번이 뒤처지거나 토큰을 놓칠 수 있다.)
22+
private final RedisTemplate<String, String> redisTemplate;
23+
24+
public RedisWaitingQueueRepository(
25+
@Qualifier(RedisConfig.REDIS_TEMPLATE_MASTER) RedisTemplate<String, String> redisTemplate
26+
) {
27+
this.redisTemplate = redisTemplate;
28+
}
29+
30+
@Override
31+
public void enter(Long userId) {
32+
// ZADD NX: 이미 대기 중이면 score(진입시각)를 덮어쓰지 않는다 → 최초 순번 유지(중복 진입 무시)
33+
redisTemplate.opsForZSet()
34+
.addIfAbsent(WAITING_QUEUE_KEY, String.valueOf(userId), System.currentTimeMillis());
35+
}
36+
37+
@Override
38+
public Optional<Long> findRank(Long userId) {
39+
// ZRANK: score 오름차순에서 몇 번째인지(0-based). 없으면 null → empty
40+
Long rank = redisTemplate.opsForZSet().rank(WAITING_QUEUE_KEY, String.valueOf(userId));
41+
return Optional.ofNullable(rank);
42+
}
43+
44+
@Override
45+
public long size() {
46+
// ZCARD: 대기열 전체 인원. 키가 없으면 null 이므로 0 으로 보정
47+
Long count = redisTemplate.opsForZSet().zCard(WAITING_QUEUE_KEY);
48+
return count == null ? 0L : count;
49+
}
50+
51+
@Override
52+
public List<Long> popMin(int count) {
53+
// ZPOPMIN: score 가 가장 낮은(=먼저 온) N명을 꺼내면서 동시에 제거한다 (atomic)
54+
Set<ZSetOperations.TypedTuple<String>> popped =
55+
redisTemplate.opsForZSet().popMin(WAITING_QUEUE_KEY, count);
56+
if (popped == null) {
57+
return List.of();
58+
}
59+
return popped.stream()
60+
.map(ZSetOperations.TypedTuple::getValue)
61+
.filter(Objects::nonNull)
62+
.map(Long::valueOf)
63+
.toList();
64+
}
65+
}

apps/commerce-api/src/main/java/com/loopers/interfaces/api/order/OrderV1ApiSpec.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,5 +8,5 @@
88
public interface OrderV1ApiSpec {
99

1010
@Operation(summary = "주문 생성", description = "여러 상품을 한 번에 주문한다. 재고 부족 시 전체 실패(All-or-Nothing).")
11-
ApiResponse<OrderV1Dto.OrderResponse> createOrder(Long userId, OrderV1Dto.OrderRequest request);
11+
ApiResponse<OrderV1Dto.OrderResponse> createOrder(Long userId, String entryToken, OrderV1Dto.OrderRequest request);
1212
}

apps/commerce-api/src/main/java/com/loopers/interfaces/api/order/OrderV1Controller.java

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22

33
import com.loopers.application.order.OrderFacade;
44
import com.loopers.application.order.OrderInfo;
5+
import com.loopers.application.queue.QueueFacade;
56
import com.loopers.interfaces.api.ApiResponse;
67
import com.loopers.support.error.CoreException;
78
import com.loopers.support.error.ErrorType;
@@ -14,22 +15,30 @@
1415
public class OrderV1Controller implements OrderV1ApiSpec {
1516

1617
private final OrderFacade orderFacade;
18+
private final QueueFacade queueFacade;
1719

1820
@PostMapping
1921
@Override
2022
public ApiResponse<OrderV1Dto.OrderResponse> createOrder(
2123
@RequestHeader(value = "X-Loopers-UserId", required = false) Long userId,
24+
@RequestHeader(value = "X-Entry-Token", required = false) String entryToken,
2225
@RequestBody OrderV1Dto.OrderRequest request
2326
) {
2427
// ① 인증 식별자 확인 (헤더 필수)
2528
if (userId == null) {
2629
throw new CoreException(ErrorType.BAD_REQUEST, "X-Loopers-UserId 헤더가 필요합니다.");
2730
}
2831

29-
// ② 헤더 userId + 바디 request 를 합쳐 응용 입력(OrderCriteria)으로 변환 → Facade 호출
32+
// ② 대기열 관문: 입장 토큰 검증 (관문이 꺼져 있으면 통과)
33+
queueFacade.validateEntry(userId, entryToken);
34+
35+
// ③ 헤더 userId + 바디 request 를 합쳐 응용 입력(OrderCriteria)으로 변환 → Facade 호출
3036
OrderInfo info = orderFacade.placeOrder(request.toCriteria(userId));
3137

32-
// ③ 응용 출력(OrderInfo) → API 응답(OrderResponse) 으로 변환
38+
// ④ 주문 성공 후 입장 토큰 소진(삭제)
39+
queueFacade.consumeEntry(userId);
40+
41+
// ⑤ 응용 출력(OrderInfo) → API 응답(OrderResponse) 으로 변환
3342
return ApiResponse.success(OrderV1Dto.OrderResponse.from(info));
3443
}
3544
}
Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,15 @@
1+
package com.loopers.interfaces.api.queue;
2+
3+
import com.loopers.interfaces.api.ApiResponse;
4+
import io.swagger.v3.oas.annotations.Operation;
5+
import io.swagger.v3.oas.annotations.tags.Tag;
6+
7+
@Tag(name = "Queue V1 API", description = "주문 대기열 관련 API")
8+
public interface QueueV1ApiSpec {
9+
10+
@Operation(summary = "대기열 진입", description = "대기열에 진입하고 현재 순번과 예상 대기시간을 반환한다.")
11+
ApiResponse<QueueV1Dto.EnterResponse> enter(Long userId);
12+
13+
@Operation(summary = "순번 조회", description = "현재 순번과 예상 대기시간을 반환한다. 대기열에 없으면 404.")
14+
ApiResponse<QueueV1Dto.PositionResponse> position(Long userId);
15+
}

0 commit comments

Comments
 (0)