From 06b64b695ad956ebeb671cf3b1aed1a9b256834a Mon Sep 17 00:00:00 2001 From: MINYONG PARK Date: Tue, 4 Aug 2026 09:06:18 +0900 Subject: [PATCH 1/8] =?UTF-8?q?feat(S15P11A705-290):=20Kafka=20=EC=9D=98?= =?UTF-8?q?=EC=A1=B4=EC=84=B1=20=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 --- build.gradle | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/build.gradle b/build.gradle index afc28f4d..2035627c 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,9 @@ dependencies { testImplementation 'org.springframework.boot:spring-boot-testcontainers' testImplementation 'org.testcontainers:testcontainers-junit-jupiter' testImplementation 'org.testcontainers:testcontainers-postgresql' + testImplementation 'org.testcontainers:testcontainers-kafka' + // DLT 검증에서 특정 토픽의 레코드 하나를 기다리는 헬퍼(KafkaTestUtils)용. + testImplementation 'org.springframework.kafka:spring-kafka-test' testCompileOnly 'org.projectlombok:lombok' testRuntimeOnly 'org.junit.platform:junit-platform-launcher' testAnnotationProcessor 'org.projectlombok:lombok' From f014e9ccf3c1a69e60c14e6e6598b15ccb7a6f2b Mon Sep 17 00:00:00 2001 From: MINYONG PARK Date: Tue, 4 Aug 2026 09:09:00 +0900 Subject: [PATCH 2/8] =?UTF-8?q?feat(S15P11A705-290):=20Kafka=20=EB=B8=8C?= =?UTF-8?q?=EB=A1=9C=EC=BB=A4=C2=B7=ED=81=90=20=EC=84=A4=EC=A0=95=20?= =?UTF-8?q?=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 --- compose.yaml | 28 +++++++++++++++++++ .../domain/ai/AiIntegrationConfig.java | 16 ++++++++++- .../domain/ai/queue/AiQueueProperties.java | 27 ++++++++++++++++++ .../domain/ai/queue/package-info.java | 9 ++++++ src/main/resources/application-local.yml | 4 ++- src/main/resources/application.yml | 25 +++++++++++++++++ 6 files changed, 107 insertions(+), 2 deletions(-) create mode 100644 src/main/java/com/pinlog/pinlogback/domain/ai/queue/AiQueueProperties.java create mode 100644 src/main/java/com/pinlog/pinlogback/domain/ai/queue/package-info.java 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/src/main/java/com/pinlog/pinlogback/domain/ai/AiIntegrationConfig.java b/src/main/java/com/pinlog/pinlogback/domain/ai/AiIntegrationConfig.java index 65fb62a2..37a812d1 100644 --- a/src/main/java/com/pinlog/pinlogback/domain/ai/AiIntegrationConfig.java +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/AiIntegrationConfig.java @@ -2,12 +2,14 @@ import java.util.concurrent.Executor; +import org.apache.kafka.clients.admin.NewTopic; import org.slf4j.Logger; import org.slf4j.LoggerFactory; 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; @@ -15,6 +17,7 @@ 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; /** @@ -29,11 +32,22 @@ @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 public RestClient aiProcessRestClient(AiProperties properties) { 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/package-info.java b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/package-info.java new file mode 100644 index 00000000..4df72ef8 --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/package-info.java @@ -0,0 +1,9 @@ +/** + * Context→AI 요청을 나르는 Kafka 큐(BD-48). + * + *

커밋 후 리스너가 발행하고({@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/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장이며 여기서 임의로 바꾸지 않는다. 상수로 박지 않는 # 이유는 튜닝 대상인 것도 있지만, 무엇보다 "만료됐다"를 테스트에서 만들 수 없기 때문이다. From 6914fdd292f12ead589e25fad7c774c6cc0a0b9b Mon Sep 17 00:00:00 2001 From: MINYONG PARK Date: Tue, 4 Aug 2026 09:10:56 +0900 Subject: [PATCH 3/8] =?UTF-8?q?test(S15P11A705-290):=20=ED=86=B5=ED=95=A9?= =?UTF-8?q?=20=ED=85=8C=EC=8A=A4=ED=8A=B8=EC=97=90=20Kafka=20=EC=BB=A8?= =?UTF-8?q?=ED=85=8C=EC=9D=B4=EB=84=88=20=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 --- .../integration/IntegrationContainerSupport.java | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/src/test/java/com/pinlog/pinlogback/integration/IntegrationContainerSupport.java b/src/test/java/com/pinlog/pinlogback/integration/IntegrationContainerSupport.java index db010c7a..476e75e5 100644 --- a/src/test/java/com/pinlog/pinlogback/integration/IntegrationContainerSupport.java +++ b/src/test/java/com/pinlog/pinlogback/integration/IntegrationContainerSupport.java @@ -3,6 +3,7 @@ import org.springframework.boot.testcontainers.service.connection.ServiceConnection; import org.springframework.test.context.TestPropertySource; import org.testcontainers.containers.GenericContainer; +import org.testcontainers.kafka.KafkaContainer; import org.testcontainers.postgresql.PostgreSQLContainer; import org.testcontainers.utility.DockerImageName; @@ -47,10 +48,14 @@ // // 끄지 않고 늘리는 이유: @Scheduled 등록 자체가 검증 대상이다(fixedDelay 인지, 전용 스케줄러를 // 쓰는지). 조건부로 끄면 그 계약을 볼 수 없다. +// 큐 재시도도 줄인다. 기본값(4회, 1s부터 지수 백오프)이면 재시도 체인 소진(DLT 격리)을 검증하는 +// 테스트가 회차마다 합계 7초를 기다린다. 줄여도 검증 대상(체인을 타고 DLT에 도달한다)은 같다. @TestPropertySource(properties = { "pinlog.ai.base-url=http://127.0.0.1:1", "pinlog.ai.internal-secret=test-internal-secret", - "pinlog.ai.rescan.interval=PT1H" + "pinlog.ai.rescan.interval=PT1H", + "pinlog.ai.queue.retry-attempts=3", + "pinlog.ai.queue.retry-initial-delay-ms=100" }) public abstract class IntegrationContainerSupport { @@ -67,6 +72,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 +90,7 @@ public abstract class IntegrationContainerSupport { // 여기서 수동으로 시작하면 컨테이너 하나가 실행 내내 살아 있고, 정리는 JVM 종료 시 Ryuk가 한다. POSTGRES.start(); REDIS.start(); + KAFKA.start(); } } From 5f0084a1dc59d015879e775c3202e608fafa0acb Mon Sep 17 00:00:00 2001 From: MINYONG PARK Date: Tue, 4 Aug 2026 09:16:03 +0900 Subject: [PATCH 4/8] =?UTF-8?q?feat(S15P11A705-290):=20=EC=9D=B8=EB=A9=94?= =?UTF-8?q?=EB=AA=A8=EB=A6=AC=20=ED=81=90=EB=A5=BC=20Kafka=20=EB=B0=9C?= =?UTF-8?q?=ED=96=89=C2=B7=EC=86=8C=EB=B9=84=20=EA=B2=BD=EB=A1=9C=EB=A1=9C?= =?UTF-8?q?=20=EA=B5=90=EC=B2=B4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 발행 실패는 삼킨다 — PENDING이 커밋되어 있어 재스캔이 복구한다(기존 의미론 유지). aiCallExecutor·@EnableAsync는 이 리스너 전용이었으므로 함께 제거한다. Co-Authored-By: Claude Fable 5 --- .../domain/ai/AiIntegrationConfig.java | 38 +------ .../ai/event/ContextAiRequestedListener.java | 48 +++------ .../ai/queue/ContextAiProcessConsumer.java | 40 +++++++ .../ai/queue/ContextAiProcessMessage.java | 28 +++++ .../ai/queue/ContextAiProcessPublisher.java | 52 +++++++++ .../domain/ai/FastApiProcessStub.java | 20 ++-- .../ai/queue/ContextAiKafkaPipelineTests.java | 101 ++++++++++++++++++ 7 files changed, 250 insertions(+), 77 deletions(-) create mode 100644 src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessConsumer.java create mode 100644 src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessMessage.java create mode 100644 src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessPublisher.java create mode 100644 src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java 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 37a812d1..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,19 +1,13 @@ package com.pinlog.pinlogback.domain.ai; -import java.util.concurrent.Executor; - import org.apache.kafka.clients.admin.NewTopic; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; 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; @@ -25,18 +19,16 @@ * 패키지에 두는 기준은 {@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, AiQueueProperties.class}) public class AiIntegrationConfig { - private static final Logger log = LoggerFactory.getLogger(AiIntegrationConfig.class); - /** * Context→AI 요청의 본 토픽(BD-48). 재시도 토픽과 DLT는 {@code @RetryableTopic}이 여기서 * 파생해 만들므로 본 토픽만 선언한다. 파티션 1인 이유: 처리량이 Record 저장 빈도(초당 수 건)를 @@ -116,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/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/queue/ContextAiProcessConsumer.java b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessConsumer.java new file mode 100644 index 00000000..f717be76 --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessConsumer.java @@ -0,0 +1,40 @@ +package com.pinlog.pinlogback.domain.ai.queue; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.kafka.annotation.KafkaListener; +import org.springframework.stereotype.Component; + +import com.pinlog.pinlogback.domain.ai.client.AiProcessClient; +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는 소비 시점에 이미 + * 삭제·교체된 것이므로 호출을 생략한다 — 그것이 삭제 확인 그 자체다. + */ +@Component +public class ContextAiProcessConsumer { + + private static final Logger log = LoggerFactory.getLogger(ContextAiProcessConsumer.class); + + private final ContextProcessRequestAssembler assembler; + private final AiProcessClient client; + + public ContextAiProcessConsumer(ContextProcessRequestAssembler assembler, AiProcessClient client) { + this.assembler = assembler; + this.client = client; + } + + @KafkaListener(topics = "${pinlog.ai.queue.topic}", groupId = "${pinlog.ai.queue.group}") + public void consume(String payload) { + ContextAiProcessMessage message = ContextAiProcessMessage.fromJson(payload); + assembler.assemble(message.contextId()).ifPresentOrElse( + client::process, + () -> log.debug("이미 삭제된 Context라 소비를 생략한다: contextId={}", message.contextId())); + } +} 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 kafkaTemplate; + private final AiQueueProperties properties; + + public ContextAiProcessPublisher(KafkaTemplate kafkaTemplate, AiQueueProperties properties) { + this.kafkaTemplate = kafkaTemplate; + this.properties = properties; + } + + public void publish(ContextAiRequested event) { + ContextAiProcessMessage message = + new ContextAiProcessMessage(event.contextId(), event.memberId(), event.recordId()); + try { + kafkaTemplate.send(properties.topic(), Long.toString(event.contextId()), message.toJson()) + .whenComplete((result, failure) -> { + if (failure != null) { + log.warn("AI process 메시지 발행 실패(브로커 응답 대기 중): contextId={}, cause={}", + event.contextId(), failure.toString()); + } + }); + } catch (RuntimeException e) { + log.warn("AI process 메시지 발행 실패: contextId={}, cause={}", event.contextId(), e.toString()); + } + } +} 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..5d6d7ea8 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,7 +65,7 @@ record Received( } /** 대역의 응답 방식. 실패 경로가 롤백을 유발하지 않는 것을 확인하려면 실패를 만들 수 있어야 한다. */ - enum Mode { + public enum Mode { /** 정상 접수. */ ACCEPTED, /** 5xx — 상대 장애. */ @@ -79,7 +79,7 @@ enum Mode { private final BlockingQueue received = new LinkedBlockingQueue<>(); private final AtomicReference mode = new AtomicReference<>(Mode.ACCEPTED); - FastApiProcessStub(PostgreSQLContainer postgres) { + public FastApiProcessStub(PostgreSQLContainer postgres) { this.postgres = postgres; try { this.server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); @@ -91,22 +91,22 @@ enum Mode { server.start(); } - String baseUrl() { + public String baseUrl() { return "http://127.0.0.1:" + server.getAddress().getPort(); } - void reset(Mode next) { + public void reset(Mode next) { mode.set(next); received.clear(); } /** 호출이 오지 않으면 {@code null}. 비동기 호출이라 폴링이 아니라 대기로 받는다. */ - Received awaitCall() throws InterruptedException { + public Received awaitCall() throws InterruptedException { return received.poll(15, TimeUnit.SECONDS); } /** 호출이 오지 않았음을 확인할 때 쓴다. 짧게 기다린 뒤 비어 있으면 참으로 본다. */ - boolean noCallWithin(long millis) throws InterruptedException { + public boolean noCallWithin(long millis) throws InterruptedException { return received.poll(millis, TimeUnit.MILLISECONDS) == null; } @@ -114,7 +114,7 @@ boolean noCallWithin(long millis) throws InterruptedException { * 대역을 쓰는 클래스가 끝날 때 닫는다. {@code stop}은 실행자를 건드리지 않으므로 직접 내린다 — * 그러지 않으면 non-daemon 스레드 둘이 JVM 끝까지 남는다. */ - void stop() { + public void stop() { server.stop(0); if (server.getExecutor() instanceof ExecutorService executor) { executor.shutdownNow(); diff --git a/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java b/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java new file mode 100644 index 00000000..1a7e1948 --- /dev/null +++ b/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java @@ -0,0 +1,101 @@ +package com.pinlog.pinlogback.domain.ai.queue; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.math.BigDecimal; +import java.util.List; + +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.DynamicPropertyRegistry; +import org.springframework.test.context.DynamicPropertySource; + +import com.pinlog.pinlogback.domain.ai.FastApiProcessStub; +import com.pinlog.pinlogback.domain.member.entity.Member; +import com.pinlog.pinlogback.domain.member.repository.MemberRepository; +import com.pinlog.pinlogback.domain.record.dto.PlacePayload; +import com.pinlog.pinlogback.domain.record.dto.RecordCreateRequest; +import com.pinlog.pinlogback.domain.record.entity.Context; +import com.pinlog.pinlogback.domain.record.repository.ContextRepository; +import com.pinlog.pinlogback.domain.record.service.RecordService; +import com.pinlog.pinlogback.integration.IntegrationContainerSupport; + +/** + * Record 저장 → Kafka 발행 → back 컨슈머 → FastAPI 호출이 완주하는지 확인한다(BD-48). + * + *

{@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; + + @BeforeEach + void resetStub() { + STUB.reset(FastApiProcessStub.Mode.ACCEPTED); + } + + @AfterAll + static void stopStub() { + STUB.stop(); + } + + @Test + void messageOnTheTopicReachesFastApiThroughTheConsumer() throws Exception { + long memberId = newMemberId(); + + recordService.create(memberId, createRequest("kafka-pipe-1", "카프카로 간다")); + + FastApiProcessStub.Received call = STUB.awaitCall(); + assertThat(call).as("발행→소비→HTTP 경로가 완주해야 한다").isNotNull(); + assertThat(call.stateVisible()) + .as("PENDING 커밋 후에만 발행되므로 소비 시점에는 반드시 보인다") + .isTrue(); + assertThat(call.text()).isEqualTo("카프카로 간다"); + } + + private long newMemberId() { + return memberRepository.save(Member.create()).getId(); + } + + private long onlyContextId(Long recordId) { + List contexts = contextRepository.findByRecordIdOrderByOriginCreatedAtAscIdAsc(recordId); + assertThat(contexts).hasSize(1); + return contexts.get(0).getId(); + } + + private RecordCreateRequest createRequest(String kakaoPlaceId, String body) { + return new RecordCreateRequest( + new PlacePayload( + kakaoPlaceId, + "앤트러사이트 성수", + "서울 성동구 성수이로 7길 30", + "서울 성동구 성수이로7길 30", + null, + null, + new BigDecimal("37.5445000"), + new BigDecimal("127.0557000") + ), + body); + } +} From e9812784d0af41ecf19634e95208b20762676796 Mon Sep 17 00:00:00 2001 From: MINYONG PARK Date: Tue, 4 Aug 2026 09:24:58 +0900 Subject: [PATCH 5/8] =?UTF-8?q?feat(S15P11A705-290):=20=EC=86=8C=EB=B9=84?= =?UTF-8?q?=20=EC=8B=A4=ED=8C=A8=EB=A5=BC=20=EC=9E=AC=EC=8B=9C=EB=8F=84=20?= =?UTF-8?q?=ED=86=A0=ED=94=BD=20=EC=B2=B4=EC=9D=B8=EA=B3=BC=20DLT=EB=A1=9C?= =?UTF-8?q?=20=EA=B2=A9=EB=A6=AC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 5xx·연결 실패는 백오프를 늘려 가며 재시도하고, 4xx·역직렬화 실패는 몇 번을 보내도 같으므로 재시도 없이 DLT로 직행한다. DLT 격리는 ERROR 로그로 남긴다 — 상태는 PENDING으로 남아 있어 최종 처분은 여전히 재스캔·Finalizer의 몫이다. Co-Authored-By: Claude Fable 5 --- build.gradle | 2 - .../domain/ai/client/AiProcessClient.java | 39 +++++++++-- .../ai/exception/AiProcessFatalException.java | 13 ++++ .../AiProcessRetryableException.java | 12 ++++ .../ai/queue/ContextAiProcessConsumer.java | 52 ++++++++++++++- .../domain/ai/ContextAiEnqueueTests.java | 10 ++- .../domain/ai/FastApiProcessStub.java | 13 +++- .../ai/queue/ContextAiKafkaPipelineTests.java | 65 +++++++++++++++++++ 8 files changed, 193 insertions(+), 13 deletions(-) create mode 100644 src/main/java/com/pinlog/pinlogback/domain/ai/exception/AiProcessFatalException.java create mode 100644 src/main/java/com/pinlog/pinlogback/domain/ai/exception/AiProcessRetryableException.java diff --git a/build.gradle b/build.gradle index 2035627c..ff45ec4c 100644 --- a/build.gradle +++ b/build.gradle @@ -62,8 +62,6 @@ dependencies { testImplementation 'org.testcontainers:testcontainers-junit-jupiter' testImplementation 'org.testcontainers:testcontainers-postgresql' testImplementation 'org.testcontainers:testcontainers-kafka' - // DLT 검증에서 특정 토픽의 레코드 하나를 기다리는 헬퍼(KafkaTestUtils)용. - testImplementation 'org.springframework.kafka:spring-kafka-test' testCompileOnly 'org.projectlombok:lombok' testRuntimeOnly 'org.junit.platform:junit-platform-launcher' testAnnotationProcessor 'org.projectlombok:lombok' 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/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/ContextAiProcessConsumer.java b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessConsumer.java index f717be76..1939747a 100644 --- a/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessConsumer.java +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessConsumer.java @@ -2,10 +2,16 @@ 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.ContextProcessRequestAssembler; /** @@ -16,6 +22,11 @@ *

본문을 메시지에서 읽지 않고 조립기로 다시 읽는 이유는 이벤트 시절과 같다 * ({@code docs/ai/spec/ai-integration.md} 4.2·4.4). 조립기가 비면 그 Context는 소비 시점에 이미 * 삭제·교체된 것이므로 호출을 생략한다 — 그것이 삭제 확인 그 자체다. + * + *

여기서는 실패를 삼키지 않고 던진다. 인메모리 큐에서는 던져도 아무것도 복구되지 + * 않았지만, 여기서 던지는 것은 재시도 체인의 신호다 — {@code @RetryableTopic}이 백오프가 다른 + * 재시도 토픽으로 옮겨 다시 전달하고, 소진되면 DLT로 격리한다. 단 {@code AiProcessFatalException} + * (4xx·역직렬화 실패)은 체인을 태우지 않고 DLT로 직행한다 — 같은 메시지는 몇 번을 보내도 같다. */ @Component public class ContextAiProcessConsumer { @@ -30,11 +41,48 @@ public ContextAiProcessConsumer(ContextProcessRequestAssembler assembler, AiProc this.client = client; } + /** + * 재시도 값은 {@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 = ContextAiProcessMessage.fromJson(payload); + ContextAiProcessMessage message = parse(payload); assembler.assemble(message.contextId()).ifPresentOrElse( - client::process, + client::processOrThrow, () -> log.debug("이미 삭제된 Context라 소비를 생략한다: contextId={}", message.contextId())); } + + /** + * 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/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 5d6d7ea8..319822ca 100644 --- a/src/test/java/com/pinlog/pinlogback/domain/ai/FastApiProcessStub.java +++ b/src/test/java/com/pinlog/pinlogback/domain/ai/FastApiProcessStub.java @@ -70,6 +70,8 @@ public enum Mode { ACCEPTED, /** 5xx — 상대 장애. */ SERVER_ERROR, + /** 4xx — 요청 자체의 문제. 몇 번을 다시 보내도 같으므로 재시도 대상이 아니다. */ + BAD_REQUEST, /** 응답 없이 연결을 끊는다 — 연결·타임아웃 계열 실패에 해당한다. */ HANG_UP } @@ -131,10 +133,19 @@ private void handle(HttpExchange exchange) throws IOException { exchange.close(); return; } - exchange.sendResponseHeaders(current == Mode.ACCEPTED ? 202 : 503, -1); + exchange.sendResponseHeaders(statusOf(current), -1); exchange.close(); } + private int statusOf(Mode current) { + return switch (current) { + case ACCEPTED -> 202; + case BAD_REQUEST -> 400; + case SERVER_ERROR -> 503; + case HANG_UP -> throw new IllegalStateException("HANG_UP은 응답을 보내지 않는다"); + }; + } + private Received observe(HttpExchange exchange, JsonNode body, long contextId) { String secret = exchange.getRequestHeaders().getFirst("X-Internal-Secret"); String text = body.get("text").asString(); diff --git a/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java b/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java index 1a7e1948..d3861f52 100644 --- a/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java +++ b/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java @@ -3,8 +3,15 @@ import static org.assertj.core.api.Assertions.assertThat; import java.math.BigDecimal; +import java.time.Duration; import java.util.List; +import java.util.Map; +import java.util.UUID; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.common.serialization.StringDeserializer; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -74,6 +81,64 @@ void messageOnTheTopicReachesFastApiThroughTheConsumer() throws Exception { assertThat(call.text()).isEqualTo("카프카로 간다"); } + /** + * 5xx는 재시도 체인을 소진한 뒤 DLT로 격리된다. 시도 횟수는 + * {@code IntegrationContainerSupport}가 3으로 줄여 두었다(본 토픽 1 + 재시도 토픽 2). + */ + @Test + void serverErrorsExhaustTheRetryChainAndLandInTheDlt() throws Exception { + STUB.reset(FastApiProcessStub.Mode.SERVER_ERROR); + long memberId = newMemberId(); + + long recordId = recordService.create(memberId, createRequest("kafka-dlt-1", "5xx는 재시도 후 DLT")).recordId(); + + assertThat(STUB.awaitCall()).as("1차 시도").isNotNull(); + assertThat(STUB.awaitCall()).as("재시도 1").isNotNull(); + assertThat(STUB.awaitCall()).as("재시도 2 — attempts=3의 마지막").isNotNull(); + assertThat(dltReceivesMessageFor(onlyContextId(recordId), Duration.ofSeconds(15))) + .as("소진된 메시지는 DLT로 격리된다") + .isTrue(); + } + + /** 4xx는 요청 자체의 문제라 몇 번을 다시 보내도 같다 — 재시도 없이 DLT로 직행한다. */ + @Test + void badRequestGoesStraightToTheDltWithoutRetry() throws Exception { + STUB.reset(FastApiProcessStub.Mode.BAD_REQUEST); + long memberId = newMemberId(); + + long recordId = recordService.create(memberId, createRequest("kafka-dlt-2", "4xx는 재시도 없이 DLT")).recordId(); + + assertThat(STUB.awaitCall()).as("한 번은 호출된다").isNotNull(); + assertThat(dltReceivesMessageFor(onlyContextId(recordId), Duration.ofSeconds(15))).isTrue(); + assertThat(STUB.noCallWithin(1000)).as("4xx는 재시도하지 않는다").isTrue(); + } + + /** + * 자기만의 그룹으로 DLT를 처음부터 읽되, 이 테스트의 contextId만 자기 것으로 친다. + * {@code earliest}로 읽으면 앞 테스트가 격리시킨 메시지도 함께 오므로, 대조 없이 "무언가 왔다"로 + * 판정하면 앞 테스트의 잔여물로도 통과해 버린다. + */ + private boolean dltReceivesMessageFor(long contextId, Duration timeout) { + Map props = Map.of( + ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, KAFKA.getBootstrapServers(), + ConsumerConfig.GROUP_ID_CONFIG, "dlt-probe-" + UUID.randomUUID(), + ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest", + ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class, + ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + long deadline = System.nanoTime() + timeout.toNanos(); + try (KafkaConsumer probe = new KafkaConsumer<>(props)) { + probe.subscribe(List.of("context-ai.process-dlt")); + while (System.nanoTime() < deadline) { + for (ConsumerRecord record : probe.poll(Duration.ofMillis(500))) { + if (ContextAiProcessMessage.fromJson(record.value()).contextId() == contextId) { + return true; + } + } + } + } + return false; + } + private long newMemberId() { return memberRepository.save(Member.create()).getId(); } From a1ce08712846c598381f24c21d6ba433db5283f1 Mon Sep 17 00:00:00 2001 From: MINYONG PARK Date: Tue, 4 Aug 2026 09:30:18 +0900 Subject: [PATCH 6/8] =?UTF-8?q?feat(S15P11A705-290):=20=EC=A4=91=EB=B3=B5?= =?UTF-8?q?=20=EC=A0=84=EB=8B=AC=EC=9D=84=20=EC=83=81=ED=83=9C=20=EA=B8=B0?= =?UTF-8?q?=EC=A4=80=EC=9C=BC=EB=A1=9C=20=EA=B1=B8=EB=9F=AC=EB=82=B4?= =?UTF-8?q?=EB=8A=94=20=EB=A9=B1=EB=93=B1=20=EA=B0=80=EB=93=9C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit at-least-once 전달에서 중복은 전제다. 걸러 내는 근거는 브로커가 아니라 진실의 원본인 ai.context_ai_state다 — 두 단계 모두 PENDING을 지났으면 생략한다. Co-Authored-By: Claude Fable 5 --- .../ai/queue/ContextAiProcessConsumer.java | 27 ++++++++- .../ai/queue/ContextAiKafkaPipelineTests.java | 55 +++++++++++++++++++ 2 files changed, 81 insertions(+), 1 deletion(-) 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 index 1939747a..7d4f0a5e 100644 --- a/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessConsumer.java +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiProcessConsumer.java @@ -12,6 +12,7 @@ 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; /** @@ -33,12 +34,17 @@ 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) { + public ContextAiProcessConsumer(ContextProcessRequestAssembler assembler, AiProcessClient client, + AiRescanCandidateService candidates) { this.assembler = assembler; this.client = client; + this.candidates = candidates; } /** @@ -59,11 +65,30 @@ public ContextAiProcessConsumer(ContextProcessRequestAssembler assembler, AiProc @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) diff --git a/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java b/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java index d3861f52..9545d200 100644 --- a/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java +++ b/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java @@ -17,6 +17,8 @@ import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.kafka.core.KafkaTemplate; import org.springframework.test.context.DynamicPropertyRegistry; import org.springframework.test.context.DynamicPropertySource; @@ -25,6 +27,7 @@ import com.pinlog.pinlogback.domain.member.repository.MemberRepository; import com.pinlog.pinlogback.domain.record.dto.PlacePayload; import com.pinlog.pinlogback.domain.record.dto.RecordCreateRequest; +import com.pinlog.pinlogback.domain.record.dto.RecordCreateResponse; import com.pinlog.pinlogback.domain.record.entity.Context; import com.pinlog.pinlogback.domain.record.repository.ContextRepository; import com.pinlog.pinlogback.domain.record.service.RecordService; @@ -57,6 +60,12 @@ static void aiServerPointsAtTheStub(DynamicPropertyRegistry registry) { @Autowired private MemberRepository memberRepository; + @Autowired + private JdbcTemplate jdbcTemplate; + + @Autowired + private KafkaTemplate kafkaTemplate; + @BeforeEach void resetStub() { STUB.reset(FastApiProcessStub.Mode.ACCEPTED); @@ -113,6 +122,52 @@ void badRequestGoesStraightToTheDltWithoutRetry() throws Exception { assertThat(STUB.noCallWithin(1000)).as("4xx는 재시도하지 않는다").isTrue(); } + /** + * 같은 메시지가 두 번 와도 처리 단계를 이미 지난 Context에는 힌트를 또 보내지 않는다. + * at-least-once 전달에서 중복은 결함이 아니라 전제다 — 걸러 내는 근거는 브로커가 아니라 + * 진실의 원본인 {@code ai.context_ai_state}다. + */ + @Test + void duplicateDeliveryForAnAlreadyHandledContextIsSkipped() throws Exception { + long memberId = newMemberId(); + RecordCreateResponse created = recordService.create(memberId, createRequest("kafka-dup-1", "중복은 걸러진다")); + assertThat(STUB.awaitCall()).isNotNull(); + long contextId = onlyContextId(created.recordId()); + jdbcTemplate.update( + "UPDATE ai.context_ai_state SET embedding_status = 'COMPLETED', keyword_status = 'COMPLETED' " + + "WHERE context_id = ?", + contextId); + STUB.reset(FastApiProcessStub.Mode.ACCEPTED); + + kafkaTemplate.send("context-ai.process", Long.toString(contextId), + new ContextAiProcessMessage(contextId, memberId, created.recordId()).toJson()).get(); + + assertThat(STUB.noCallWithin(2000)) + .as("이미 처리 단계를 지난 Context에 힌트를 또 보내지 않는다") + .isTrue(); + } + + /** 삭제된 Context의 메시지는 실패가 아니라 정상 생략이다 — 재시도도 DLT 격리도 일어나지 않는다. */ + @Test + void messageForADeletedContextCompletesWithoutCallOrDlt() throws Exception { + long memberId = newMemberId(); + RecordCreateResponse created = recordService.create(memberId, createRequest("kafka-del-1", "교체 전")); + assertThat(STUB.awaitCall()).isNotNull(); + long recordId = created.recordId(); + long oldContextId = onlyContextId(recordId); + recordService.replaceContext(memberId, recordId, oldContextId, "교체 후"); + assertThat(STUB.awaitCall()).isNotNull(); + STUB.reset(FastApiProcessStub.Mode.ACCEPTED); + + kafkaTemplate.send("context-ai.process", Long.toString(oldContextId), + new ContextAiProcessMessage(oldContextId, memberId, recordId).toJson()).get(); + + assertThat(STUB.noCallWithin(2000)).as("삭제된 Context는 호출을 생략한다").isTrue(); + assertThat(dltReceivesMessageFor(oldContextId, Duration.ofSeconds(3))) + .as("정상 생략이지 실패가 아니므로 DLT로 가지 않는다") + .isFalse(); + } + /** * 자기만의 그룹으로 DLT를 처음부터 읽되, 이 테스트의 contextId만 자기 것으로 친다. * {@code earliest}로 읽으면 앞 테스트가 격리시킨 메시지도 함께 오므로, 대조 없이 "무언가 왔다"로 From b54c0bb88805b7a20c042e91522162377b54955f Mon Sep 17 00:00:00 2001 From: MINYONG PARK Date: Tue, 4 Aug 2026 09:32:01 +0900 Subject: [PATCH 7/8] =?UTF-8?q?docs(S15P11A705-290):=20BD-48=20=EA=B8=B0?= =?UTF-8?q?=EB=A1=9D,=20BD-17=20=EB=8C=80=EC=B2=B4=20=EC=B2=98=EB=A6=AC,?= =?UTF-8?q?=20=EC=9E=91=EC=97=85=20=EB=A1=9C=EA=B7=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 --- .../BD-17-async-without-message-queue.md | 2 +- .../decisions/BD-48-context-ai-kafka-queue.md | 70 +++++++++++++++++++ ...6-08-04-S15P11A705-290-context-ai-kafka.md | 26 +++++++ .../ai/scheduler/AiRescanScheduler.java | 13 ++-- 4 files changed, 104 insertions(+), 7 deletions(-) create mode 100644 docs/backend/decisions/BD-48-context-ai-kafka-queue.md create mode 100644 docs/backend/worklog/2026-08-04-S15P11A705-290-context-ai-kafka.md 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/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를 먼저 두는 이유는 나중에 두면 같은 회차에서 방금 From 39d1b0f8c6a655796cd200f5b45fa779d7400bae Mon Sep 17 00:00:00 2001 From: MINYONG PARK Date: Tue, 4 Aug 2026 09:47:56 +0900 Subject: [PATCH 8/8] =?UTF-8?q?test(S15P11A705-290):=20=ED=81=90=20?= =?UTF-8?q?=ED=86=A0=ED=94=BD=C2=B7=EA=B7=B8=EB=A3=B9=EC=9D=84=20=EC=BB=A8?= =?UTF-8?q?=ED=85=8D=EC=8A=A4=ED=8A=B8=EB=B3=84=EB=A1=9C=20=EA=B2=A9?= =?UTF-8?q?=EB=A6=AC=ED=95=B4=20=EC=8A=A4=EC=9C=84=ED=8A=B8=20=EA=B0=84?= =?UTF-8?q?=EC=84=AD=20=EC=A0=9C=EA=B1=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 캐시된 컨텍스트들이 토픽·그룹을 공유하면 리밸런스 재전달이 다른 클래스의 '호출 없음' 검증 구간에 흘러든다. Hikari 풀도 5로 줄인다 — 컨텍스트가 하나 늘며 Postgres max_connections(100)를 넘겼다. Co-Authored-By: Claude Fable 5 --- .../ai/queue/ContextAiKafkaPipelineTests.java | 10 +++++-- .../IntegrationContainerSupport.java | 28 ++++++++++++++++++- 2 files changed, 34 insertions(+), 4 deletions(-) diff --git a/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java b/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java index 9545d200..81cb8bdc 100644 --- a/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java +++ b/src/test/java/com/pinlog/pinlogback/domain/ai/queue/ContextAiKafkaPipelineTests.java @@ -66,6 +66,10 @@ static void aiServerPointsAtTheStub(DynamicPropertyRegistry registry) { @Autowired private KafkaTemplate kafkaTemplate; + /** 토픽은 컨텍스트마다 격리된다({@code IntegrationContainerSupport}) — 하드코딩하면 남의 토픽을 본다. */ + @Autowired + private AiQueueProperties queueProperties; + @BeforeEach void resetStub() { STUB.reset(FastApiProcessStub.Mode.ACCEPTED); @@ -139,7 +143,7 @@ void duplicateDeliveryForAnAlreadyHandledContextIsSkipped() throws Exception { contextId); STUB.reset(FastApiProcessStub.Mode.ACCEPTED); - kafkaTemplate.send("context-ai.process", Long.toString(contextId), + kafkaTemplate.send(queueProperties.topic(), Long.toString(contextId), new ContextAiProcessMessage(contextId, memberId, created.recordId()).toJson()).get(); assertThat(STUB.noCallWithin(2000)) @@ -159,7 +163,7 @@ void messageForADeletedContextCompletesWithoutCallOrDlt() throws Exception { assertThat(STUB.awaitCall()).isNotNull(); STUB.reset(FastApiProcessStub.Mode.ACCEPTED); - kafkaTemplate.send("context-ai.process", Long.toString(oldContextId), + kafkaTemplate.send(queueProperties.topic(), Long.toString(oldContextId), new ContextAiProcessMessage(oldContextId, memberId, recordId).toJson()).get(); assertThat(STUB.noCallWithin(2000)).as("삭제된 Context는 호출을 생략한다").isTrue(); @@ -182,7 +186,7 @@ private boolean dltReceivesMessageFor(long contextId, Duration timeout) { ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); long deadline = System.nanoTime() + timeout.toNanos(); try (KafkaConsumer probe = new KafkaConsumer<>(props)) { - probe.subscribe(List.of("context-ai.process-dlt")); + probe.subscribe(List.of(queueProperties.topic() + "-dlt")); while (System.nanoTime() < deadline) { for (ConsumerRecord record : probe.poll(Duration.ofMillis(500))) { if (ContextAiProcessMessage.fromJson(record.value()).contextId() == contextId) { diff --git a/src/test/java/com/pinlog/pinlogback/integration/IntegrationContainerSupport.java b/src/test/java/com/pinlog/pinlogback/integration/IntegrationContainerSupport.java index 476e75e5..165a9935 100644 --- a/src/test/java/com/pinlog/pinlogback/integration/IntegrationContainerSupport.java +++ b/src/test/java/com/pinlog/pinlogback/integration/IntegrationContainerSupport.java @@ -1,6 +1,10 @@ package com.pinlog.pinlogback.integration; +import java.util.UUID; + import org.springframework.boot.testcontainers.service.connection.ServiceConnection; +import org.springframework.test.context.DynamicPropertyRegistry; +import org.springframework.test.context.DynamicPropertySource; import org.springframework.test.context.TestPropertySource; import org.testcontainers.containers.GenericContainer; import org.testcontainers.kafka.KafkaContainer; @@ -50,15 +54,37 @@ // 쓰는지). 조건부로 끄면 그 계약을 볼 수 없다. // 큐 재시도도 줄인다. 기본값(4회, 1s부터 지수 백오프)이면 재시도 체인 소진(DLT 격리)을 검증하는 // 테스트가 회차마다 합계 7초를 기다린다. 줄여도 검증 대상(체인을 타고 DLT에 도달한다)은 같다. +// +// Hikari 풀도 줄인다. 스위트가 캐시하는 Spring 컨텍스트 하나마다 풀(기본 10)이 통째로 살아 +// 남는데, Kafka 큐 도입으로 전용 컨텍스트가 하나 더 생기며 Postgres 기본 max_connections(100)를 +// 넘겼다 — 실제로 too many clients로 죽었다. 테스트는 순차 실행이라 5로도 남는다. @TestPropertySource(properties = { "pinlog.ai.base-url=http://127.0.0.1:1", "pinlog.ai.internal-secret=test-internal-secret", "pinlog.ai.rescan.interval=PT1H", "pinlog.ai.queue.retry-attempts=3", - "pinlog.ai.queue.retry-initial-delay-ms=100" + "pinlog.ai.queue.retry-initial-delay-ms=100", + "spring.datasource.hikari.maximum-pool-size=5" }) public abstract class IntegrationContainerSupport { + /** + * 큐 토픽·그룹을 Spring 컨텍스트마다 격리한다. 캐시된 컨텍스트들은 스위트가 끝날 때까지 + * 살아서 컨슈머를 계속 돌리는데, 토픽·그룹을 공유하면 새 컨텍스트가 뜰 때마다 리밸런스가 나고 + * 커밋 안 된 offset이 재전달되어 다른 클래스의 "호출이 없어야 한다" 검증 구간에 남의 + * 호출이 흘러든다(실측: rollingBackTheTransactionNeverReachesFastApi가 그렇게 깨졌다). + * 토픽이 갈리면 컨슈머가 서로의 메시지를 볼 수 없어 재생·리밸런스 간섭이 구조적으로 사라진다. + * + *

이 메서드는 컨텍스트 생성 시 한 번 돌므로 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 =