diff --git a/build.gradle b/build.gradle index afc28f4d..ff45ec4c 100644 --- a/build.gradle +++ b/build.gradle @@ -45,6 +45,9 @@ dependencies { implementation 'org.springframework.boot:spring-boot-starter-security-oauth2-client' implementation 'org.springdoc:springdoc-openapi-starter-webmvc-ui:3.0.3' implementation 'org.springframework.boot:spring-boot-starter-flyway' + // Context→AI 요청 경로의 전달 수단(BD-48). 재시도 토픽 체인·DLT는 spring-kafka의 + // @RetryableTopic을 쓴다. 버전은 Boot BOM이 관리한다. + implementation 'org.springframework.boot:spring-boot-starter-kafka' implementation 'io.micrometer:micrometer-registry-prometheus' implementation 'org.flywaydb:flyway-database-postgresql' compileOnly 'org.projectlombok:lombok' @@ -58,6 +61,7 @@ dependencies { testImplementation 'org.springframework.boot:spring-boot-testcontainers' testImplementation 'org.testcontainers:testcontainers-junit-jupiter' testImplementation 'org.testcontainers:testcontainers-postgresql' + testImplementation 'org.testcontainers:testcontainers-kafka' testCompileOnly 'org.projectlombok:lombok' testRuntimeOnly 'org.junit.platform:junit-platform-launcher' testAnnotationProcessor 'org.projectlombok:lombok' diff --git a/compose.yaml b/compose.yaml index db913d13..aa8e65cc 100644 --- a/compose.yaml +++ b/compose.yaml @@ -33,5 +33,33 @@ services: timeout: 3s retries: 10 + kafka: + # Context→AI 요청 큐(BD-48). 테스트 컨테이너(IntegrationContainerSupport)도 같은 태그를 쓴다. + # KRaft 단일 노드. 호스트(애플리케이션)는 HOST 리스너(19092)로, 컨테이너 내부는 9092로 붙는다 — + # Kafka는 접속 후 advertised listener로 재접속하므로 호스트용 주소를 따로 광고해야 한다. + image: apache/kafka:4.1.0 + labels: + # Boot의 compose 자동 연결을 끄고 application-local.yml의 명시 주소를 쓴다. 자동 연결은 + # 9092의 매핑 포트를 기대하는데 이 구성은 호스트 리스너가 19092라 서로 어긋난다. + org.springframework.boot.ignore: true + environment: + KAFKA_NODE_ID: 1 + KAFKA_PROCESS_ROLES: broker,controller + KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093 + KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093,HOST://:19092 + KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,HOST://localhost:${KAFKA_PORT:-19092} + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT,HOST:PLAINTEXT + KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT + KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 + ports: + - "${KAFKA_PORT:-19092}:19092" + healthcheck: + test: ["CMD-SHELL", "/opt/kafka/bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092 > /dev/null 2>&1"] + interval: 5s + timeout: 10s + retries: 10 + start_period: 15s + volumes: postgres-data: diff --git a/docs/backend/decisions/BD-17-async-without-message-queue.md b/docs/backend/decisions/BD-17-async-without-message-queue.md index 18b0fc3e..f3dd7b59 100644 --- a/docs/backend/decisions/BD-17-async-without-message-queue.md +++ b/docs/backend/decisions/BD-17-async-without-message-queue.md @@ -1,6 +1,6 @@ # BD-17. 비동기 AI 처리를 메시지 큐 없이 DB State + Scheduler로 -- **상태**: Accepted +- **상태**: Superseded by [BD-48](BD-48-context-ai-kafka-queue.md) — 요청 전달 수단(인메모리 큐→Kafka)에 한한 대체. DB State가 진실의 원본이라는 결정과 재스캔 안전망은 유지된다 - **날짜**: 2026-07-22 (AI 설계 확립 시점. 단일 커밋으로 특정 불가) - **작성 시점**: 2026-07-27 — 결정 이후에 정리 - **관련**: S15P11A705-76 diff --git a/docs/backend/decisions/BD-48-context-ai-kafka-queue.md b/docs/backend/decisions/BD-48-context-ai-kafka-queue.md new file mode 100644 index 00000000..cc8b6e7c --- /dev/null +++ b/docs/backend/decisions/BD-48-context-ai-kafka-queue.md @@ -0,0 +1,70 @@ +# BD-48. Context→AI 요청 경로의 인메모리 큐를 Kafka로 교체한다 + +- **상태**: Accepted +- **날짜**: 2026-08-04 +- **관련**: S15P11A705-290 · [back#129](https://github.com/Team-PinLog/back/issues/129) · [BD-17](BD-17-async-without-message-queue.md) + +## 맥락 + +Record 저장 커밋 후 FastAPI 호출은 `@Async` 리스너의 인메모리 큐(`aiCallExecutor`)를 탔다. +이 큐는 포화되면 요청을 버리고(재스캔이 다음 주기에 줍는다), 프로세스가 죽으면 담긴 힌트가 +사라지며, 브로커 수준의 재시도·백오프·DLQ 같은 실패 격리 장치가 없다. + +back#129에서 안 A(요청 방향 전환)·안 B(완료 이벤트 역방향)·안 C(back 내부에만 도입)를 비교했고, +AI 파트가 안 C 착수에 이견 없음을 밝히며 다섯 가지 질문을 남겼다. 핵심 지적은 정확하다 — +**유실 방지는 이미 DB State가 해결하고 있으므로, 이 경로에서 MQ가 실제로 개선하는 것은 +"재스캔 주기(최대 5분)만큼의 복구 지연"이다.** 재스캔 주기 단축·큐 용량 확대 같은 더 싼 +대안도 함께 제시됐다. + +동시에 이 도입의 명분 절반은 처음부터 기술 경험·포트폴리오였고, back#129 본문과 AI 파트 +코멘트 모두 "학습 가치는 그 자체로 팀이 도입을 결정할 수 있는 사유"로 인정했다. + +## 선택지 + +| 안 | 장점 | 단점 | +|---|---|---| +| (a) 현상 유지 + 설정 조정(재스캔 주기 단축 등) | 비용 0에 가까움 | 복구 지연 개선뿐, 실패 격리·학습 목표는 얻지 못함 | +| (b) RabbitMQ | 작업 큐 패턴과 도구 성격이 가장 일치(ack·DLX 기본 제공) | 컨테이너 추가. 이력서·확장성 관점 매력은 Kafka보다 낮다고 판단 | +| (c) Redis Streams | 인프라 추가 0 | 백오프·DLQ 직접 구현 — "MQ를 배웠다"기보다 "Redis로 흉내"에 가까움 | +| (d) Kafka | 학습·포트폴리오 가치 최대. 안 B(완료 이벤트)를 넘어 이벤트 백본으로 키울 여지 | 이 워크로드에는 과투자. 개별 재전달이 없어 재시도 토픽 체인을 직접 구성. 운영 부담 최대 | + +## 결정 + +**(d) Kafka를 채택한다. 능동적 선택이며, back#129가 열어 둔 "학습 목적으로 과투자를 감수한다"는 +선택지를 집는 것이다.** 신뢰성 필수론으로 포장하지 않는다 — 아키텍처가 얻는 실익은 복구 지연 +단축(5분 → 백오프 수 초)과 실패 격리(DLT)이고, 나머지는 학습 가치다. + +구현 경계는 다음을 지킨다. + +- **DB State(`ai.context_ai_state`)가 진실의 원본이라는 BD-17의 핵심은 유지한다.** Kafka는 전달 + 수단만 교체하며, 소비자의 멱등 가드도 브로커가 아니라 상태 행을 근거로 판정한다. +- **FastAPI와의 계약은 바꾸지 않는다.** 소비자는 재스캔과 같은 조립기·클라이언트로 같은 HTTP를 + 보낸다. AI 파트 코드 변경 없음(안 C의 정의). +- **재스캔은 최종 안전망으로 유지한다.** 브로커와 독립인 HTTP 직접 호출 경로다. 발행 실패·브로커 + 장애·DLT 격리 어느 경우에도 PENDING이 남아 재스캔·Finalizer가 최종 처분한다. +- 실패 분류: 5xx·연결 실패는 재시도 토픽 체인(`@RetryableTopic`, 지수 백오프), 4xx·역직렬화 + 실패는 DLT 직행. DLT 격리는 ERROR 로그로 남긴다. +- 발행은 커밋 후 리스너에서 하며 실패를 삼킨다(기존 의미론). `max.block.ms=1000`으로 브로커 + 장애가 요청 스레드를 붙잡는 상한을 건다. + +back#129의 다섯 질문에 대한 답은 이 문서와 티켓 S15P11A705-290 본문이 겸한다 — ① 포화 실측 +없음(예방+학습), ② 5분 지연은 시연 체감 문제, ③ 주기 단축 대비 추가 이득은 실패 격리와 학습, +④ 단일 인스턴스에서 브로커가 주는 것은 프로세스 독립 내구성, ⑤ DLQ 관측은 ERROR 로그가 +최소선이고 지표는 `/metrics` 승인 뒤 후속. + +## 결과 + +- 이 결정으로 감수하는 것: + - **재시도 주체가 둘이 된다.** 체인(수 초 단위 순간 장애)과 재스캔(그 밖의 모든 유실, 예산 + `retry_count`)의 역할을 위처럼 갈랐고, FAILED 종결 권한은 여전히 재스캔·Finalizer에만 있다. + - **INFRA에 브로커 운영 부담이 생긴다.** dev/prod 배포는 INFRA 소유이며 이 작업의 머지 + 선행 조건이다. + - **DLT 관측이 로그뿐이다.** 지표·알림은 `/metrics` prod 승격 승인(infra `docs/ai-serving.md` + 검증 7) 이후 후속 작업. + - 공용 계약 `10_MVP_기능범위` §2(MQ 제외)와 어긋난다 — back#129 합의로 대체하며 docs 레포 + 갱신이 후속으로 필요하다. +- 재검토 트리거(이 조건이 오면 다시 논의): + - 안 B(FastAPI 완료 이벤트 역방향)를 진행하게 되면 — 이벤트 스키마·발행 실패 정책을 AI 파트와 + 합의하고, 이 브로커를 그 백본으로 쓸지 결정한다. + - 브로커 운영 부담이 실익을 넘어서면 — 안 (a)로 회귀하는 것도 열려 있다. 재스캔이 그대로 + 있으므로 회귀 비용은 발행·소비 코드 제거뿐이다. diff --git a/docs/backend/worklog/2026-08-04-S15P11A705-290-context-ai-kafka.md b/docs/backend/worklog/2026-08-04-S15P11A705-290-context-ai-kafka.md new file mode 100644 index 00000000..ee0a7a8c --- /dev/null +++ b/docs/backend/worklog/2026-08-04-S15P11A705-290-context-ai-kafka.md @@ -0,0 +1,26 @@ +# Context→AI 요청 경로의 인메모리 큐를 Kafka로 교체 + +- **날짜**: 2026-08-04 +- **추적**: S15P11A705-290 +- **관련**: [BD-48](../decisions/BD-48-context-ai-kafka-queue.md) · [back#129](https://github.com/Team-PinLog/back/issues/129) · [BD-17](../decisions/BD-17-async-without-message-queue.md) + +back#129의 안 C를 구현했다. 커밋 후 리스너가 `@Async` 인메모리 큐 대신 Kafka 토픽에 발행하고, +back 내부 컨슈머가 소비해 기존과 동일한 FastAPI 호출을 보낸다. FastAPI 계약과 재스캔 안전망은 +그대로다. 브로커를 Kafka로 정한 근거(학습 가치로 과투자를 감수)와 그때 검토한 트레이드오프, +AI 파트의 다섯 질문에 대한 답은 BD-48에 있다. + +내린 판단들: + +- 발행 실패는 기존처럼 삼킨다. 대신 `max.block.ms=1000`으로 브로커 장애가 요청 스레드를 붙잡는 + 상한을 걸었다 — 커밋 후 리스너가 요청 스레드에서 돌기 때문이다. +- 실패 분류를 예외 타입으로 나눴다. 5xx·연결 실패는 `AiProcessRetryableException`(재시도 토픽 + 체인), 4xx·역직렬화 실패는 `AiProcessFatalException`(DLT 직행). `AiProcessClient`의 삼키는 + `process()`는 재스캔용으로 남기고, 던지는 `processOrThrow()`를 소비자용으로 추가했다. +- 멱등 가드는 브로커가 아니라 상태 행 기준이다 — 두 단계 모두 PENDING을 지난 Context는 생략한다. +- `aiCallExecutor`·`@EnableAsync`는 이 리스너 전용이었으므로 함께 제거했다. +- 5xx·연결 끊김 통합 테스트는 재시도 체인이 생기며 호출이 3회로 늘었다. 각 테스트가 자기 + 재시도를 끝까지 소비하도록 고쳤다 — 남기면 뒤 테스트의 "호출 없음" 검증 구간에 흘러든다. + +머지 선행 조건: INFRA의 Kafka 브로커 배포(dev). AI 파트 소유 명세 +`docs/ai/spec/ai-integration.md` 4장(인메모리 큐 서술)과 공용 계약 `10_MVP_기능범위` §2(MQ 제외)가 +이 변경과 어긋난다 — 각각 갱신 요청과 docs 레포 후속 PR이 필요하다. diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/AiIntegrationConfig.java b/src/main/java/com/pinlog/pinlogback/domain/ai/AiIntegrationConfig.java index 65fb62a2..547be04c 100644 --- a/src/main/java/com/pinlog/pinlogback/domain/ai/AiIntegrationConfig.java +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/AiIntegrationConfig.java @@ -1,20 +1,17 @@ package com.pinlog.pinlogback.domain.ai; -import java.util.concurrent.Executor; - -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; +import org.apache.kafka.clients.admin.NewTopic; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.http.client.SimpleClientHttpRequestFactory; +import org.springframework.kafka.config.TopicBuilder; import org.springframework.scheduling.TaskScheduler; -import org.springframework.scheduling.annotation.EnableAsync; import org.springframework.scheduling.annotation.EnableScheduling; -import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler; import org.springframework.web.client.RestClient; +import com.pinlog.pinlogback.domain.ai.queue.AiQueueProperties; import com.pinlog.pinlogback.domain.ai.service.AiRescanProperties; /** @@ -22,17 +19,26 @@ * 패키지에 두는 기준은 {@code docs/development/package-structure.md}의 보안 설정과 같다 — 설정과 * 그 설정이 조립하는 구현이 떨어져 있으면 한쪽만 고치게 된다. * - *
{@link EnableScheduling}이 여기 있는 것은 {@link EnableAsync}와 같은 사정이다. 둘 다 애플리케이션 - * 전역 스위치인데, 켜야 하는 이유가 이 연동에만 있다. {@code global/config}로 올리면 스위치와 - * 그 스위치의 유일한 소비자가 떨어져 앉는다. + *
{@link EnableScheduling}이 여기 있는 것도 같은 사정이다. 애플리케이션 전역 스위치인데, + * 켜야 하는 이유가 이 연동(재스캔)에만 있다. {@code global/config}로 올리면 스위치와 그 스위치의 + * 유일한 소비자가 떨어져 앉는다. 한때 여기 있던 {@code @EnableAsync}와 {@code aiCallExecutor} + * (커밋 후 FastAPI 호출용 인메모리 큐)는 Kafka 발행으로 대체되어 사라졌다(BD-48). */ @Configuration -@EnableAsync @EnableScheduling -@EnableConfigurationProperties({AiProperties.class, AiRescanProperties.class}) +@EnableConfigurationProperties({AiProperties.class, AiRescanProperties.class, AiQueueProperties.class}) public class AiIntegrationConfig { - private static final Logger log = LoggerFactory.getLogger(AiIntegrationConfig.class); + /** + * Context→AI 요청의 본 토픽(BD-48). 재시도 토픽과 DLT는 {@code @RetryableTopic}이 여기서 + * 파생해 만들므로 본 토픽만 선언한다. 파티션 1인 이유: 처리량이 Record 저장 빈도(초당 수 건)를 + * 넘지 않고, 컨슈머도 단일 인스턴스라 병렬화로 얻을 것이 없다. 복제 1은 단일 브로커 전제다 — + * 브로커 구성이 커지면 INFRA와 함께 올린다. + */ + @Bean + public NewTopic contextAiProcessTopic(AiQueueProperties queueProperties) { + return TopicBuilder.name(queueProperties.topic()).partitions(1).replicas(1).build(); + } /** {@code process}용 전용 인스턴스(AI 파트 소유 명세 {@code docs/ai/spec/ai-integration.md} 2·3장). */ @Bean @@ -102,26 +108,4 @@ public TaskScheduler taskScheduler() { return scheduler; } - /** - * AI 호출 전용 풀(명세 4.2). 공용 executor를 쓰지 않는 이유는 FastAPI 장애가 다른 비동기 작업까지 - * 굶기지 않게 하기 위해서다. - * - *
큐가 차면 버린다. {@code CallerRunsPolicy}를 쓰면 요청 스레드가 외부 호출을 대신 - * 수행해 FastAPI 장애가 그대로 Core 처리량 저하가 된다. 버려진 요청은 {@code PENDING}으로 남아 - * 재스캔 대상이 되므로 유실이 아니다 — 그래서 버리는 쪽이 안전하다. - */ - @Bean - public Executor aiCallExecutor() { - ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); - executor.setThreadNamePrefix("ai-call-"); - executor.setCorePoolSize(2); - executor.setMaxPoolSize(4); - executor.setQueueCapacity(100); - // 버린 사실은 남긴다. 조용히 버리면 "왜 PENDING만 쌓이는가"를 알 길이 없다. 예외를 던지는 - // AbortPolicy는 쓰지 않는다 — 커밋 이후 리스너를 타고 올라가 정상 응답을 오류로 만든다. - executor.setRejectedExecutionHandler((task, pool) -> - log.debug("AI 호출 큐 포화로 요청을 버렸다. PENDING이 남아 재스캔이 다시 집는다.")); - executor.initialize(); - return executor; - } } diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/client/AiProcessClient.java b/src/main/java/com/pinlog/pinlogback/domain/ai/client/AiProcessClient.java index 7997754c..8a92858b 100644 --- a/src/main/java/com/pinlog/pinlogback/domain/ai/client/AiProcessClient.java +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/client/AiProcessClient.java @@ -13,6 +13,8 @@ import org.springframework.web.client.RestClientResponseException; import com.pinlog.pinlogback.domain.ai.AiProperties; +import com.pinlog.pinlogback.domain.ai.exception.AiProcessFatalException; +import com.pinlog.pinlogback.domain.ai.exception.AiProcessRetryableException; import com.pinlog.pinlogback.global.web.TraceIdFilter; /** @@ -23,11 +25,12 @@ * 통보용 웹훅·콜백이 없으므로 이 클래스는 응답 본문을 읽지 않는다. 정합성의 근거는 이 호출이 * 아니라 DB에 영속된 {@code ai.context_ai_state}다 — 이 호출은 "지금 처리하면 조금 빨라지는 힌트"다. * - *
모든 실패를 삼킨다. 예외를 밖으로 던지면 호출자(커밋 이후 리스너)를 타고 올라가 - * 사용자 응답을 오류로 만들 수 있는데, Core 데이터는 이미 커밋되어 정상이므로 그것은 거짓말이다. - * 상태를 {@code FAILED}로 바꾸지도, {@code retry_count}를 올리지도 않는다 — 호출 실패는 FastAPI - * 내부 작업의 실패가 아니고, 두 주체가 같은 사유로 FAILED를 각각 기록하면 원인 추적이 불가능해진다. - * {@code PENDING}이 남아 있으므로 재스캔이 같은 Context를 다시 집는다. + *
실패의 취급이 호출자마다 다르므로 진입점이 둘이다. {@link #process}는 모든 실패를 + * 삼킨다 — 재스캔 경로에서 예외는 아무것도 복구하지 못하고, {@code PENDING}이 남아 다음 + * 회차가 같은 Context를 다시 집는다. {@link #processOrThrow}는 실패를 분류해 던진다 — 큐 소비 + * 경로에서는 던지는 것이 곧 재시도 체인·DLT 격리의 신호다(BD-48). 어느 쪽도 상태를 + * {@code FAILED}로 바꾸거나 {@code retry_count}를 올리지 않는다 — 호출 실패는 FastAPI 내부 + * 작업의 실패가 아니고, 두 주체가 같은 사유로 FAILED를 각각 기록하면 원인 추적이 불가능해진다. */ @Component public class AiProcessClient { @@ -75,6 +78,20 @@ private static String requireSecret(String secret, Environment environment) { /** 호출 결과를 반환하지 않는다. 호출자가 분기할 수 있으면 그 분기가 곧 Core 트랜잭션 결과에 스며든다. */ public void process(ContextProcessRequest request) { + try { + processOrThrow(request); + } catch (AiProcessRetryableException | AiProcessFatalException e) { + // 로그는 processOrThrow가 이미 남겼다. 재스캔 경로에서 실패는 여기서 끝난다 — + // PENDING이 남아 있으므로 다음 회차가 같은 Context를 다시 집는다. + } + } + + /** + * 실패를 분류해 던진다. 5xx·연결 계열은 {@link AiProcessRetryableException}(다시 보내면 성공할 + * 수 있다), 4xx는 {@link AiProcessFatalException}(요청 자체의 문제라 몇 번을 보내도 같다). + * 401·403도 fatal이다 — 시크릿·헤더 설정이 고쳐지기 전에는 재시도가 전부 같은 답을 받는다. + */ + public void processOrThrow(ContextProcessRequest request) { String requestId = currentRequestId(); try { restClient.post() @@ -86,12 +103,20 @@ public void process(ContextProcessRequest request) { .toBodilessEntity(); log.debug("AI process 접수됨: contextId={}, requestId={}", request.contextId(), requestId); } catch (RestClientResponseException e) { - // 4xx는 요청 payload 형식 문제일 가능성이 있어 사람이 봐야 한다. 5xx는 상대 장애이므로 - // 재스캔이 흡수한다. 응답 본문은 남기지 않는다 — 내부 API라도 로그로 새어 나갈 이유가 없다. + // 4xx는 요청 payload 형식 문제일 가능성이 있어 사람이 봐야 한다. 응답 본문은 남기지 + // 않는다 — 내부 API라도 로그로 새어 나갈 이유가 없다. logByStatus(e, request, requestId); + if (e.getStatusCode().is4xxClientError()) { + throw new AiProcessFatalException( + "AI process 호출이 " + e.getStatusCode().value() + "로 거절됐다: contextId=" + request.contextId(), e); + } + throw new AiProcessRetryableException( + "AI process 호출이 " + e.getStatusCode().value() + "를 받았다: contextId=" + request.contextId(), e); } catch (RuntimeException e) { log.warn("AI process 호출 실패(연결·타임아웃): contextId={}, requestId={}, cause={}", request.contextId(), requestId, e.toString()); + throw new AiProcessRetryableException( + "AI process 호출이 연결 단계에서 실패했다: contextId=" + request.contextId(), e); } } diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/event/ContextAiRequestedListener.java b/src/main/java/com/pinlog/pinlogback/domain/ai/event/ContextAiRequestedListener.java index d7096768..59735f87 100644 --- a/src/main/java/com/pinlog/pinlogback/domain/ai/event/ContextAiRequestedListener.java +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/event/ContextAiRequestedListener.java @@ -1,54 +1,36 @@ package com.pinlog.pinlogback.domain.ai.event; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import org.springframework.scheduling.annotation.Async; import org.springframework.stereotype.Component; import org.springframework.transaction.event.TransactionPhase; import org.springframework.transaction.event.TransactionalEventListener; -import com.pinlog.pinlogback.domain.ai.client.AiProcessClient; -import com.pinlog.pinlogback.domain.ai.service.ContextProcessRequestAssembler; +import com.pinlog.pinlogback.domain.ai.queue.ContextAiProcessPublisher; /** - * Core 커밋 이후에 FastAPI를 호출한다(AI 파트 소유 명세 {@code docs/ai/spec/ai-integration.md} 4장). + * Core 커밋 이후에 처리 요청을 큐에 발행한다(BD-48). * - *
{@code AFTER_COMMIT}이 이 클래스의 존재 이유다. 트랜잭션 안에서 호출하면 FastAPI가 - * 아직 커밋되지 않은 {@code ai.context_ai_state}를 조회해 처리 대상을 찾지 못한다. 워커는 별 - * 프로세스이므로 우리 트랜잭션의 미커밋 스냅샷을 볼 수 없다 — 202를 받고도 아무 일도 일어나지 - * 않는다. 외부 호출 지연만큼 DB 커넥션과 행 잠금이 유지되는 문제도 함께 사라진다. + *
{@code AFTER_COMMIT}이 이 클래스의 존재 이유다. 커밋 전에 발행하면 컨슈머가 아직 + * 커밋되지 않은 {@code ai.context_ai_state}를 조회해 처리 대상을 찾지 못한다 — 컨슈머와 FastAPI + * 워커는 별 프로세스이므로 우리 트랜잭션의 미커밋 스냅샷을 볼 수 없다. 롤백된 트랜잭션에서는 + * 이 리스너가 아예 실행되지 않으므로, 저장되지 않은 Context의 메시지가 큐에 실리는 경로가 없다. * - *
롤백된 트랜잭션에서는 이 리스너가 아예 실행되지 않으므로, 저장되지 않은 Context로 호출이 - * 나가는 경로가 없다. 반대 방향도 마찬가지다 — 여기서 무슨 일이 나든 Core는 이미 커밋되어 있어 - * 되돌아가지 않는다. - * - *
{@code @Async}는 응답 지연을 사용자 요청 시간에서 떼기 위한 것이다. 동기로 두면 FastAPI - * 응답 시간이 그대로 Record 저장 API의 응답 시간이 된다. + *
인메모리 큐 시절의 {@code @Async}는 여기 없다. FastAPI 응답을 기다리던 HTTP 호출과 달리 + * {@code send()}는 프로듀서 버퍼에 넣고 곧바로 돌아오므로 떼어 낼 지연이 없고, 브로커 장애로 + * 블록되는 최악의 경우도 {@code max.block.ms}(1s)가 끊는다. 발행 실패는 + * {@link ContextAiProcessPublisher}가 삼킨다 — {@code PENDING}이 이미 커밋되어 있어 재스캔이 + * 같은 Context를 다시 집는다. */ @Component public class ContextAiRequestedListener { - private static final Logger log = LoggerFactory.getLogger(ContextAiRequestedListener.class); - - private final ContextProcessRequestAssembler assembler; - private final AiProcessClient client; + private final ContextAiProcessPublisher publisher; - public ContextAiRequestedListener(ContextProcessRequestAssembler assembler, AiProcessClient client) { - this.assembler = assembler; - this.client = client; + public ContextAiRequestedListener(ContextAiProcessPublisher publisher) { + this.publisher = publisher; } - @Async("aiCallExecutor") @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) public void on(ContextAiRequested event) { - try { - assembler.assemble(event).ifPresentOrElse( - client::process, - () -> log.debug("이미 삭제된 Context라 AI process 호출을 생략한다: contextId={}", event.contextId())); - } catch (RuntimeException e) { - // 조립 단계(DB 조회)의 실패까지 삼킨다. 여기서 던지면 @Async의 uncaught handler로 갈 뿐 - // 아무 것도 복구되지 않고, PENDING은 이미 커밋되어 있어 재스캔이 같은 Context를 다시 집는다. - log.warn("AI process 요청 조립 실패: contextId={}, cause={}", event.contextId(), e.toString()); - } + publisher.publish(event); } } diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/exception/AiProcessFatalException.java b/src/main/java/com/pinlog/pinlogback/domain/ai/exception/AiProcessFatalException.java new file mode 100644 index 00000000..ff2a5287 --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/exception/AiProcessFatalException.java @@ -0,0 +1,13 @@ +package com.pinlog.pinlogback.domain.ai.exception; + +/** + * FastAPI {@code process} 호출의 영구 실패 — 4xx(요청 형식·시크릿 문제)와 역직렬화 실패. 같은 + * 메시지를 몇 번을 다시 보내도 결과가 같으므로 재시도 체인을 타지 않고 DLT로 직행한다(BD-48). + * 재시도로 시간을 끌면 사람이 봐야 할 문제의 발견만 늦어진다. + */ +public class AiProcessFatalException extends RuntimeException { + + public AiProcessFatalException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/exception/AiProcessRetryableException.java b/src/main/java/com/pinlog/pinlogback/domain/ai/exception/AiProcessRetryableException.java new file mode 100644 index 00000000..ced9e75a --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/exception/AiProcessRetryableException.java @@ -0,0 +1,12 @@ +package com.pinlog.pinlogback.domain.ai.exception; + +/** + * FastAPI {@code process} 호출의 일시 장애 — 5xx, 연결 거부, 타임아웃. 나중에 다시 보내면 성공할 + * 수 있으므로 재시도 토픽 체인의 대상이다(BD-48). 소진되면 DLT로 격리된다. + */ +public class AiProcessRetryableException extends RuntimeException { + + public AiProcessRetryableException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/queue/AiQueueProperties.java b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/AiQueueProperties.java new file mode 100644 index 00000000..7409ba7f --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/AiQueueProperties.java @@ -0,0 +1,27 @@ +package com.pinlog.pinlogback.domain.ai.queue; + +import org.springframework.boot.context.properties.ConfigurationProperties; + +/** + * Context→AI 요청을 나르는 큐 설정(BD-48). + * + *
재시도 값이 상수가 아니라 설정인 이유는 재스캔({@code pinlog.ai.rescan})과 같다 — 체인 소진을 + * 테스트에서 만들려면 시도 횟수와 백오프를 줄일 수 있어야 한다. {@code @RetryableTopic}은 이 record가 + * 아니라 같은 키의 프로퍼티 플레이스홀더를 직접 읽는다(애노테이션 속성이라 Bean 주입이 안 된다). + * 이 record는 토픽 Bean 조립과 "이 키들이 바인딩 가능한 형태로 존재한다"는 기동 시 검증을 맡는다. + * + * @param topic 본 토픽 이름. 재시도 토픽({@code -retry-N})과 DLT({@code -dlt})는 여기서 파생된다 + * @param group 컨슈머 그룹. back 인스턴스들이 하나의 그룹으로 작업을 나눠 갖는다 + * @param retryAttempts 총 시도 횟수(본 토픽 1회 포함). 소진되면 DLT로 격리된다 + * @param retryInitialDelayMs 첫 재시도까지의 지연(밀리초) + * @param retryMultiplier 재시도마다 지연에 곱하는 배수 + */ +@ConfigurationProperties("pinlog.ai.queue") +public record AiQueueProperties( + String topic, + String group, + int retryAttempts, + long retryInitialDelayMs, + double retryMultiplier +) { +} diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessConsumer.java b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessConsumer.java new file mode 100644 index 00000000..7d4f0a5e --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessConsumer.java @@ -0,0 +1,113 @@ +package com.pinlog.pinlogback.domain.ai.queue; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.kafka.annotation.BackOff; +import org.springframework.kafka.annotation.DltHandler; +import org.springframework.kafka.annotation.KafkaListener; +import org.springframework.kafka.annotation.RetryableTopic; +import org.springframework.kafka.support.KafkaHeaders; +import org.springframework.messaging.handler.annotation.Header; +import org.springframework.stereotype.Component; + +import com.pinlog.pinlogback.domain.ai.client.AiProcessClient; +import com.pinlog.pinlogback.domain.ai.exception.AiProcessFatalException; +import com.pinlog.pinlogback.domain.ai.service.AiRescanCandidateService; +import com.pinlog.pinlogback.domain.ai.service.ContextProcessRequestAssembler; + +/** + * 큐에서 처리 요청을 꺼내 FastAPI를 호출한다(BD-48). 인메모리 큐 시절 {@code @Async} 리스너가 + * 하던 일의 소비 측 절반이며, FastAPI와의 계약은 그대로다 — 조립기도 클라이언트도 재스캔이 + * 쓰는 것과 같은 것을 쓴다. + * + *
본문을 메시지에서 읽지 않고 조립기로 다시 읽는 이유는 이벤트 시절과 같다 + * ({@code docs/ai/spec/ai-integration.md} 4.2·4.4). 조립기가 비면 그 Context는 소비 시점에 이미 + * 삭제·교체된 것이므로 호출을 생략한다 — 그것이 삭제 확인 그 자체다. + * + *
여기서는 실패를 삼키지 않고 던진다. 인메모리 큐에서는 던져도 아무것도 복구되지 + * 않았지만, 여기서 던지는 것은 재시도 체인의 신호다 — {@code @RetryableTopic}이 백오프가 다른 + * 재시도 토픽으로 옮겨 다시 전달하고, 소진되면 DLT로 격리한다. 단 {@code AiProcessFatalException} + * (4xx·역직렬화 실패)은 체인을 태우지 않고 DLT로 직행한다 — 같은 메시지는 몇 번을 보내도 같다. + */ +@Component +public class ContextAiProcessConsumer { + + private static final Logger log = LoggerFactory.getLogger(ContextAiProcessConsumer.class); + + private static final String PENDING = "PENDING"; + + private final ContextProcessRequestAssembler assembler; + private final AiProcessClient client; + private final AiRescanCandidateService candidates; + + public ContextAiProcessConsumer(ContextProcessRequestAssembler assembler, AiProcessClient client, + AiRescanCandidateService candidates) { + this.assembler = assembler; + this.client = client; + this.candidates = candidates; + } + + /** + * 재시도 값은 {@code pinlog.ai.queue.*}가 정본이다({@link AiQueueProperties}). 애노테이션이 + * record 대신 플레이스홀더를 읽는 것은 애노테이션 속성에 Bean을 주입할 수 없다는 제약 때문이다. + * + *
{@code traversingCauses}를 켠 이유: 컨테이너가 리스너 예외를 + * {@code ListenerExecutionFailedException}으로 감싸는 경우가 있어, 원인 사슬을 타고 내려가야 + * fatal 분류가 우리 예외를 찾는다. + */ + @RetryableTopic( + attempts = "${pinlog.ai.queue.retry-attempts}", + backOff = @BackOff( + delayString = "${pinlog.ai.queue.retry-initial-delay-ms}", + multiplierString = "${pinlog.ai.queue.retry-multiplier}"), + exclude = AiProcessFatalException.class, + traversingCauses = "true") + @KafkaListener(topics = "${pinlog.ai.queue.topic}", groupId = "${pinlog.ai.queue.group}") + public void consume(String payload) { + ContextAiProcessMessage message = parse(payload); + if (!stillWaiting(message.contextId())) { + return; + } + assembler.assemble(message.contextId()).ifPresentOrElse( + client::processOrThrow, + () -> log.debug("이미 삭제된 Context라 소비를 생략한다: contextId={}", message.contextId())); + } + + /** + * 멱등 가드. at-least-once 전달에서 중복은 전제이고, 걸러 내는 근거는 브로커가 아니라 진실의 + * 원본인 {@code ai.context_ai_state}다 — 어느 한쪽 상태라도 {@code PENDING}이면 힌트가 아직 + * 유효하고, 둘 다 지났으면(PROCESSING·DONE·FAILED·CANCELLED) 보낼 이유가 사라진 것이다. + * 상태 행이 없으면 보내지 않는다 — 행 없이 도착한 메시지는 이미 정리된 Context다. + */ + private boolean stillWaiting(long contextId) { + boolean waiting = candidates.findCurrentState(contextId) + .map(state -> PENDING.equals(state.embeddingStatus()) || PENDING.equals(state.keywordStatus())) + .orElse(false); + if (!waiting) { + log.debug("이미 처리 단계를 지났거나 상태 행이 없는 Context라 소비를 생략한다: contextId={}", contextId); + } + return waiting; + } + + /** + * DLQ는 쌓이기만 하면 보이지 않는다(back#129 논의의 관측 우려). ERROR 레벨 + 원문 페이로드가 + * 관측의 최소선이고, 지표 노출은 {@code /metrics} 승인(infra {@code docs/ai-serving.md} 검증 7) + * 뒤의 후속이다. 여기 격리된 Context도 상태는 {@code PENDING}으로 남아 있으므로 재스캔·Finalizer가 + * 최종 처분(재시도 소진 시 FAILED 종결)을 맡는다 — DLT는 증거 보존이지 별도 복구 경로가 아니다. + */ + @DltHandler + public void deadLetter(String payload, + @Header(name = KafkaHeaders.EXCEPTION_MESSAGE, required = false) String reason) { + log.error("AI process 메시지가 재시도 소진·영구 실패로 DLT에 격리됐다: payload={}, reason={}", + payload, reason); + } + + /** 형식이 틀린 메시지는 몇 번을 다시 읽어도 같다 — 재시도 없이 DLT로 보낸다. */ + private ContextAiProcessMessage parse(String payload) { + try { + return ContextAiProcessMessage.fromJson(payload); + } catch (RuntimeException e) { + throw new AiProcessFatalException("큐 메시지를 역직렬화하지 못했다", e); + } + } +} diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessMessage.java b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessMessage.java new file mode 100644 index 00000000..502e897c --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessMessage.java @@ -0,0 +1,28 @@ +package com.pinlog.pinlogback.domain.ai.queue; + +import tools.jackson.databind.json.JsonMapper; + +/** + * 큐에 실리는 페이로드. {@code ContextAiRequested} 이벤트와 같은 식별자 셋이다 — 본문을 싣지 않는 + * 이유도 같다({@code docs/ai/spec/ai-integration.md} 4.2: 같은 {@code context_id}로 다른 본문을 + * 보내는 것은 계약 위반이므로, 소비 시점에 Core에서 다시 읽는 것만이 그 계약을 구조로 보장한다). + * + *
직렬화를 Kafka serializer 설정이 아니라 이 record가 갖는 이유: 페이로드 스키마의 정본을 + * 한 파일로 모으기 위해서다. 필드가 바뀌면 여기와 소비자 가드만 보면 된다. + * + * @param contextId 처리 대상 Context + * @param memberId 소유 회원 + * @param recordId 소속 Record + */ +public record ContextAiProcessMessage(long contextId, long memberId, long recordId) { + + private static final JsonMapper JSON = JsonMapper.builder().build(); + + public String toJson() { + return JSON.writeValueAsString(this); + } + + public static ContextAiProcessMessage fromJson(String json) { + return JSON.readValue(json, ContextAiProcessMessage.class); + } +} diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessPublisher.java b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessPublisher.java new file mode 100644 index 00000000..0f5e5507 --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessPublisher.java @@ -0,0 +1,52 @@ +package com.pinlog.pinlogback.domain.ai.queue; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.stereotype.Component; + +import com.pinlog.pinlogback.domain.ai.event.ContextAiRequested; + +/** + * 커밋된 Context의 처리 요청을 큐에 발행한다(BD-48). + * + *
모든 실패를 삼킨다. 인메모리 큐 시절 {@code AiProcessClient}가 그랬던 이유와 같다 — + * 예외를 올리면 커밋 이후 리스너를 타고 사용자 응답을 오류로 만드는데, Core는 이미 커밋되어 + * 정상이다. 발행이 실패해도 {@code PENDING}이 남아 재스캔이 같은 Context를 다시 집는다. + * + *
동기 예외까지 잡는 이유: 브로커가 죽어 있으면 {@code send()}가 메타데이터를 기다리다 + * {@code max.block.ms}(1s) 초과로 호출 스레드에서 던진다. 이 클래스는 요청 스레드에서 + * 불리므로 그 예외가 곧 사용자 오류 응답이 된다. + * + *
키가 {@code contextId}인 것은 같은 Context의 메시지를 같은 파티션에 몰기 위한 것이다.
+ * 지금은 파티션이 1이라 효과가 없지만, 키 없는 발행은 파티션을 늘리는 순간 순서가 흩어진다.
+ */
+@Component
+public class ContextAiProcessPublisher {
+
+ private static final Logger log = LoggerFactory.getLogger(ContextAiProcessPublisher.class);
+
+ private final KafkaTemplate 커밋 후 리스너가 발행하고({@code ContextAiProcessPublisher}), back 내부 컨슈머가 소비해
+ * 기존과 동일한 FastAPI 호출을 수행한다({@code ContextAiProcessConsumer}) — FastAPI와의 계약은
+ * 바뀌지 않는다. 실패는 재시도 토픽 체인을 거쳐 DLT로 격리되고, 그와 별개로
+ * {@code ai.context_ai_state}가 진실의 원본이라는 구조(BD-17의 핵심)와 재스캔 안전망은 유지된다.
+ */
+package com.pinlog.pinlogback.domain.ai.queue;
diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/scheduler/AiRescanScheduler.java b/src/main/java/com/pinlog/pinlogback/domain/ai/scheduler/AiRescanScheduler.java
index 12b64b65..95c5cf3c 100644
--- a/src/main/java/com/pinlog/pinlogback/domain/ai/scheduler/AiRescanScheduler.java
+++ b/src/main/java/com/pinlog/pinlogback/domain/ai/scheduler/AiRescanScheduler.java
@@ -19,12 +19,13 @@
* (AI 파트 소유 명세 {@code docs/ai/spec/ai-rescan-scheduler.md} 3.1).
*
* 이 클래스가 없으면 한 번 실패한 Context는 영구히 {@code PENDING}으로 남는다. 실패 경로
- * 네 곳이 모두 "재스캔이 복구한다"를 안전망으로 전제한다 — 큐 포화로 버려진 요청
- * ({@link com.pinlog.pinlogback.domain.ai.AiIntegrationConfig}), 삼켜진 호출 실패
- * ({@link AiProcessClient}), 커밋 이후 리스너의 조립 실패
- * ({@link com.pinlog.pinlogback.domain.ai.event.ContextAiRequestedListener}), 그리고 FastAPI가 202
- * 이후 내부에서 실패한 경우. 상태만 보면 "처리 대기 중"이라 정상과 구별되지 않는 것이 이 문제의
- * 성질이다.
+ * 네 곳이 모두 "재스캔이 복구한다"를 안전망으로 전제한다 — 발행 실패·브로커 장애로 큐에 실리지
+ * 못한 요청({@link com.pinlog.pinlogback.domain.ai.queue.ContextAiProcessPublisher}), 재시도 체인
+ * 소진으로 DLT에 격리된 메시지
+ * ({@link com.pinlog.pinlogback.domain.ai.queue.ContextAiProcessConsumer}), 재스캔 자신의 호출
+ * 실패({@link AiProcessClient}가 삼킨다), 그리고 FastAPI가 202 이후 내부에서 실패한 경우. 상태만
+ * 보면 "처리 대기 중"이라 정상과 구별되지 않는 것이 이 문제의 성질이다. 이 안전망은 브로커와
+ * 독립이다 — FastAPI를 큐 없이 직접 호출하므로 브로커가 통째로 죽어도 복구가 돈다(BD-48).
*
* 순서가 계약이다. {@code Finalize → 후보 잠금·retry 증가 → 커밋 → Context 재조회 → 삭제
* 확인 → FastAPI 호출}. Finalize를 먼저 두는 이유는 나중에 두면 같은 회차에서 방금
diff --git a/src/main/resources/application-local.yml b/src/main/resources/application-local.yml
index 6ced6e27..ccbd9be3 100644
--- a/src/main/resources/application-local.yml
+++ b/src/main/resources/application-local.yml
@@ -1,4 +1,4 @@
-# 로컬 개발 override. compose.yaml이 띄운 Postgres/Redis에 접속한다.
+# 로컬 개발 override. compose.yaml이 띄운 Postgres/Redis/Kafka에 접속한다.
# 로컬 전용 값만 두고, 실제 비밀번호·토큰은 넣지 않는다(configuration.md).
spring:
datasource:
@@ -9,6 +9,8 @@ spring:
redis:
host: localhost
port: 16379
+ kafka:
+ bootstrap-servers: localhost:19092
# JWT_PRIVATE_KEY를 주지 않으면 기동할 때 임시 키쌍을 만든다(BD-31). 재시작하면 기존 토큰이
# 무효가 되니 고정하고 싶으면 .env에 PEM을 넣는다 — 절차는 authentication.md 8장.
diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml
index 9fae3e65..301ba359 100644
--- a/src/main/resources/application.yml
+++ b/src/main/resources/application.yml
@@ -82,6 +82,20 @@ spring:
# 없으므로 짧게 끊고 실패로 처리한다.
timeout: 2s
connect-timeout: 1s
+ kafka:
+ # Context→AI 요청 경로의 브로커(BD-48). 브로커 배포는 INFRA 소유이며, 주소는 infra가 주입하는
+ # 환경변수를 따른다(BD-46과 같은 원칙). 기본값은 로컬 compose.yaml의 HOST 리스너다.
+ bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:localhost:19092}
+ producer:
+ properties:
+ # 브로커 장애·버퍼 포화 시 send()가 호출 스레드를 붙잡는 상한. 기본 60초는 AFTER_COMMIT
+ # 리스너가 요청 스레드에서 돌므로 사용자 응답을 그만큼 지연시킨다. 초과하면 발행 실패로
+ # 삼켜지고 PENDING이 남아 재스캔이 복구한다 — Redis timeout을 짧게 끊는 것과 같은 사정이다.
+ "[max.block.ms]": 1000
+ consumer:
+ # 컨슈머 그룹이 처음 붙거나 offset이 만료됐을 때 놓친 메시지를 버리지 않고 따라잡는다.
+ # 중복 소비는 소비자의 멱등 가드(상태 기준)가 흡수한다.
+ auto-offset-reset: earliest
jpa:
hibernate:
ddl-auto: validate
@@ -139,6 +153,17 @@ pinlog:
# 스레드 점유만 늘린다(AI 파트 소유 명세 docs/ai/spec/ai-integration.md 3장).
connect-timeout: 1s
read-timeout: 5s
+ # Context→AI 요청을 나르는 큐(BD-48). 커밋 후 발행 → back 내부 컨슈머가 소비해 FastAPI를
+ # 호출한다. 재시도 값이 설정인 이유는 rescan과 같다 — 테스트가 체인 소진을 기다릴 수 있어야 한다.
+ queue:
+ topic: context-ai.process
+ group: pinlog-back
+ # 재시도 토픽 체인의 총 시도 횟수(본 토픽 1회 + 재시도 토픽 3회). 소진되면 DLT로 격리된다.
+ # 재스캔(max-retry 3)과 별개 예산이다 — 체인은 순간 장애를, 재스캔은 그 밖의 모든 유실을 맡는다.
+ retry-attempts: 4
+ # 지수 백오프: 1s → 2s → 4s. 재스캔 주기(5m)보다 훨씬 짧은 구간을 맡는다.
+ retry-initial-delay-ms: 1000
+ retry-multiplier: 2.0
# 유실·정지된 AI 처리를 복구하는 재스캔과 FAILED Finalizer. 정본은 AI 파트가 소유한
# docs/ai/spec/ai-rescan-scheduler.md 2장이며 여기서 임의로 바꾸지 않는다. 상수로 박지 않는
# 이유는 튜닝 대상인 것도 있지만, 무엇보다 "만료됐다"를 테스트에서 만들 수 없기 때문이다.
diff --git a/src/test/java/com/pinlog/pinlogback/domain/ai/ContextAiEnqueueTests.java b/src/test/java/com/pinlog/pinlogback/domain/ai/ContextAiEnqueueTests.java
index 27c4c5b2..abd6d7ca 100644
--- a/src/test/java/com/pinlog/pinlogback/domain/ai/ContextAiEnqueueTests.java
+++ b/src/test/java/com/pinlog/pinlogback/domain/ai/ContextAiEnqueueTests.java
@@ -197,6 +197,10 @@ void replacingTheOnlyContextEnqueuesTheNewOneAndLeavesTheOldStateToTheDeletionPa
/**
* FastAPI가 5xx를 돌려줘도 Core와 PENDING은 남아야 한다. 롤백시키면 사용자가 저장한 기록이
* 상대 서버 장애 때문에 사라진다.
+ *
+ * 재시도까지 전부 받아 두는 이유: 5xx는 이제 재시도 토픽 체인을 탄다(BD-48,
+ * 테스트 프로파일 attempts=3). 여기서 소비하지 않으면 남은 재시도가 뒤 테스트의 "호출이
+ * 없어야 한다" 검증 구간에 도착해 엉뚱한 테스트를 깨뜨린다.
*/
@Test
void serverErrorFromFastApiDoesNotRollBackTheContextOrItsPendingRow() throws Exception {
@@ -206,6 +210,8 @@ void serverErrorFromFastApiDoesNotRollBackTheContextOrItsPendingRow() throws Exc
RecordCreateResponse created = recordService.create(memberId, createRequest("ai-enqueue-5", "5xx여도 남아야 한다"));
assertThat(STUB.awaitCall()).isNotNull();
+ assertThat(STUB.awaitCall()).as("재시도 1").isNotNull();
+ assertThat(STUB.awaitCall()).as("재시도 2 — 체인의 마지막").isNotNull();
long contextId = onlyContextId(created.recordId());
assertThat(contextRepository.findById(contextId)).isPresent();
assertThat(stateOf(contextId))
@@ -213,7 +219,7 @@ void serverErrorFromFastApiDoesNotRollBackTheContextOrItsPendingRow() throws Exc
.containsEntry("embedding_status", "PENDING");
}
- /** 연결이 끊기는 실패(다운·타임아웃 계열)도 같다. 클라이언트가 예외를 삼키지 않으면 여기서 드러난다. */
+ /** 연결이 끊기는 실패(다운·타임아웃 계열)도 같다. 재시도를 전부 받아 두는 이유는 5xx 테스트와 같다. */
@Test
void droppedConnectionDoesNotRollBackTheContextOrItsPendingRow() throws Exception {
STUB.reset(FastApiProcessStub.Mode.HANG_UP);
@@ -222,6 +228,8 @@ void droppedConnectionDoesNotRollBackTheContextOrItsPendingRow() throws Exceptio
RecordCreateResponse created = recordService.create(memberId, createRequest("ai-enqueue-6", "끊겨도 남아야 한다"));
assertThat(STUB.awaitCall()).isNotNull();
+ assertThat(STUB.awaitCall()).as("재시도 1").isNotNull();
+ assertThat(STUB.awaitCall()).as("재시도 2 — 체인의 마지막").isNotNull();
long contextId = onlyContextId(created.recordId());
assertThat(contextRepository.findById(contextId)).isPresent();
assertThat(stateOf(contextId)).containsEntry("embedding_status", "PENDING");
diff --git a/src/test/java/com/pinlog/pinlogback/domain/ai/FastApiProcessStub.java b/src/test/java/com/pinlog/pinlogback/domain/ai/FastApiProcessStub.java
index 869afd6b..319822ca 100644
--- a/src/test/java/com/pinlog/pinlogback/domain/ai/FastApiProcessStub.java
+++ b/src/test/java/com/pinlog/pinlogback/domain/ai/FastApiProcessStub.java
@@ -38,9 +38,9 @@
* 코드(별 커넥션 DB 조회)를 돌려야 하는데, 요청·응답 기록만 해 주는 mock 서버로는 그 지점을 잡을
* 수 없다. 새 테스트 의존성도 필요 없다.
*/
-final class FastApiProcessStub {
+public final class FastApiProcessStub {
- static final String PATH = "/internal/v1/context/process";
+ public static final String PATH = "/internal/v1/context/process";
private static final JsonMapper JSON = JsonMapper.builder().build();
@@ -54,7 +54,7 @@ final class FastApiProcessStub {
* @param embeddingStatus 보였다면 그 값. 아니면 {@code null}
* @param keywordStatus 보였다면 그 값. 아니면 {@code null}
*/
- record Received(
+ public record Received(
long contextId,
String internalSecret,
String text,
@@ -65,11 +65,13 @@ record Received(
}
/** 대역의 응답 방식. 실패 경로가 롤백을 유발하지 않는 것을 확인하려면 실패를 만들 수 있어야 한다. */
- enum Mode {
+ public enum Mode {
/** 정상 접수. */
ACCEPTED,
/** 5xx — 상대 장애. */
SERVER_ERROR,
+ /** 4xx — 요청 자체의 문제. 몇 번을 다시 보내도 같으므로 재시도 대상이 아니다. */
+ BAD_REQUEST,
/** 응답 없이 연결을 끊는다 — 연결·타임아웃 계열 실패에 해당한다. */
HANG_UP
}
@@ -79,7 +81,7 @@ enum Mode {
private final BlockingQueue {@code ContextAiEnqueueTests}가 "커밋과 호출의 순서"를 고정한다면, 이 클래스는 전달 수단이
+ * 브로커를 경유한다는 사실을 고정한다 — 인메모리 큐 시절에는 프로세스가 죽으면 힌트가 사라졌지만,
+ * 이제 발행까지만 성공하면 소비는 브로커가 보증한다. 실패 경로(재시도 체인·DLT)와 멱등 가드는
+ * 뒤 태스크에서 이 클래스에 추가된다.
+ */
+@SpringBootTest
+class ContextAiKafkaPipelineTests extends IntegrationContainerSupport {
+
+ private static final FastApiProcessStub STUB = new FastApiProcessStub(POSTGRES);
+
+ @DynamicPropertySource
+ static void aiServerPointsAtTheStub(DynamicPropertyRegistry registry) {
+ registry.add("pinlog.ai.base-url", STUB::baseUrl);
+ }
+
+ @Autowired
+ private RecordService recordService;
+
+ @Autowired
+ private ContextRepository contextRepository;
+
+ @Autowired
+ private MemberRepository memberRepository;
+
+ @Autowired
+ private JdbcTemplate jdbcTemplate;
+
+ @Autowired
+ private KafkaTemplate 이 메서드는 컨텍스트 생성 시 한 번 돌므로 suffix는 컨텍스트 단위로 고정된다. 같은 설정을
+ * 공유해 컨텍스트를 재사용하는 클래스들은 같은 토픽을 이어 쓴다 — 그것이 캐시의 의미다.
+ */
+ @DynamicPropertySource
+ static void isolatedQueuePerContext(DynamicPropertyRegistry registry) {
+ String suffix = UUID.randomUUID().toString().substring(0, 8);
+ registry.add("pinlog.ai.queue.topic", () -> "context-ai.process-test-" + suffix);
+ registry.add("pinlog.ai.queue.group", () -> "pinlog-back-test-" + suffix);
+ }
+
/** {@code org.testcontainers.containers.PostgreSQLContainer}는 2.x에서 deprecated다. */
@ServiceConnection
protected static final PostgreSQLContainer POSTGRES =
@@ -67,6 +98,11 @@ public abstract class IntegrationContainerSupport {
protected static final GenericContainer> REDIS =
new GenericContainer<>(DockerImageName.parse("redis:7.4.5-alpine")).withExposedPorts(6379);
+ /** {@code compose.yaml}의 kafka와 같은 태그. Context→AI 큐(BD-48)가 이 브로커를 쓴다. */
+ @ServiceConnection
+ protected static final KafkaContainer KAFKA =
+ new KafkaContainer(DockerImageName.parse("apache/kafka:4.1.0"));
+
static {
// JVM 전체에서 한 번만 띄운다(Testcontainers 싱글턴 컨테이너 패턴).
//
@@ -80,6 +116,7 @@ public abstract class IntegrationContainerSupport {
// 여기서 수동으로 시작하면 컨테이너 하나가 실행 내내 살아 있고, 정리는 JVM 종료 시 Ryuk가 한다.
POSTGRES.start();
REDIS.start();
+ KAFKA.start();
}
}