diff --git a/docs/backend/WORKLOG.md b/docs/backend/WORKLOG.md index 517c4206..0b10a568 100644 --- a/docs/backend/WORKLOG.md +++ b/docs/backend/WORKLOG.md @@ -86,3 +86,4 @@ | 2026-07-29 | Kakao·Naver 소셜 로그인을 붙였다(#33). **새로 들인 것은 응답 형태뿐이다** — 토큰 발급·회전·쿠키·회원 확정은 공급자와 무관한 경로라 BI-18이 이미 만들어 뒀고, 이 티켓이 더한 것은 공급자마다 다른 사용자 정보를 하나로 옮기는 일이다. 셋이 다 다르다: Google `sub`(최상위), Kakao `id`(최상위, **숫자**)+`kakao_account.email`, Naver `response.id`+`response.email`. **드러난 것 넷.** ① **Naver만 Spring의 식별자 검사를 우회한다** — `user-name-attribute: response`가 감싼 Map을 지목하므로 Spring은 그 키의 존재만 확인하고 안의 `id`는 보지 않는다. Google·Kakao는 최상위 스칼라라 앞단에서 걸러 주는데 Naver만 그 보증이 없어, `required(nested(attributes,"response","id"))`가 유일한 방어선이다(없으면 `provider_user_id` NOT NULL 위반이 되어 원인이 DB까지 내려간다) ② **Kakao는 client secret을 본문으로 받는다** — 기본값 `basic`이면 토큰 교환이 401이라 `client_secret_post`를 명시했다. 게다가 콘솔에서 활성화해야 검사되는 **선택 항목**이라 콘솔 상태와 설정이 어긋나면 콜백 마지막 단계에서 실패한다 ③ **이메일 없는 가입이 기본 경로일 수 있었다** — Kakao는 동의항목을 콘솔에 설정하지 않으면 인가 요청 자체가 `KOE205`로 거절되고, 이메일 수집은 비즈 앱 전환을 요구한다. 네 층(DB·엔티티·팩토리·정규화)이 모두 nullable이라 통과하며 `callbackSucceedsWithoutEmail`로 고정했다 — 선택 동의라 **동의한 사용자도 철회할 수 있어** 비즈 앱 전환 후에도 유효한 경로다 ④ 등록정보를 추가하자 **무관한 테스트 11개가 컨텍스트 실패**했는데 원인은 이 티켓이 아니라 `.env.example`이 자격증명을 빈 값으로 정의하던 기존 함정이었다(BT-05). **실제 Kakao·Naver 계정으로 수동 검증했다**: Kakao `provider_user_id=5013244578`(숫자→문자열), Naver는 43자 식별자로 저장돼 ①의 방어선이 실제로 값을 했다(Map `toString()`이 아니다), 회원 2건 분리, `auth:refresh-index` 회원별 생성·TTL, 쿠키 4종. `08 §3.1`은 이미 세 공급자를 적고 있어 **구현이 명세를 따라잡은 것**이라 공용 문서 변경은 없다. `clean check` **326개 통과** (S15P11A705-64) | [BI-24](implements/BI-24-2026-07-29-kakao-naver-login.md) · [BT-05](troubleshooting/BT-05-dotenv-empty-value-overrides-default.md) | | 2026-07-29 | 운영 런타임 Secret 5개를 `pinlog-secrets-prod` Environment 경계에서만 읽어 SHA 고정 Infra action으로 넘기는 수동 workflow를 추가했다. bridge token을 포함한 참조 6개 집합·최소 권한·`github.sha` checkout/revision은 정적 계약 테스트로 고정했고 일반 `backend-ci`에는 Secret 접근을 추가하지 않았다 (S15P11A705-154) | [BI-26](implements/BI-26-2026-07-29-runtime-secret-workflow.md) | | 2026-07-30 | back#98 리뷰 반영. `dev`가 `required_conversation_resolution`이라 **미해결 스레드가 그대로 병합 게이트**여서, 판정이 `COMMENTED`·내용이 `nit`인 줄 단위 지적 2건과 리뷰가 지목한 테스트 구멍 3건을 닫았다. 고친 둘은 **javadoc이 약속한 범위와 실제 방어 범위가 어긋난 자리**라는 점에서 성격이 같다 — ① `AiSearchClient` 생성자 javadoc이 "시크릿 검사를 여기 두지 않는 이유"로 든 논거는 "같은 키를 읽는 `AiProcessClient`가 이미 검사한다"인데, `embedding-profile`은 **이 클라이언트만 읽는 새 키**라 그 논거가 적용되지 않는다. 빈 문자열이면 FastAPI가 Profile 대조에서 422를 주므로 결과는 모든 검색이 503이고, `application.yml` 기본값은 변수를 **설정하지 않은** 경우만 막는다(`PINLOG_AI_EMBEDDING_PROFILE=`처럼 빈 값으로 정의하면 빈 문자열이 이긴다 — BT-05로 이미 겪은 형태). `requireSecret`과 같은 기준(운영 기동 실패/그 외 경고)으로 기동 시점에 끊었다. 기각된 (c)(기동 시 FastAPI 조회)가 **아니다** — 상대에게 묻지 않고 우리 값 유무만 본다. ② `distinctByRecord`는 "상대 결함이 우리 500이 되지 않게 한다"고 적어 두고 `match` **자체가 null**인 경우가 빠져 있었다. `AiSearchClient`가 최상위 `results == null`을 이미 방어하는데 그것도 계약상 올 수 없는 형태라 **층이 어긋난** 것이라, 원소 쪽 층을 맞췄다. 테스트 구멍 셋(`SEARCH_QUERY_MAX` 501자 미검증 · `PRIVATE_ONLY` 한 번도 미투입 · `insertPreset`의 `active`가 죽은 파라미터)은 **코드는 이미 맞는데 지키는 단언이 없던** 자리라 평소의 RED가 안 나온다 — 세 가드를 일부러 부순 뒤(화이트리스트를 `IN ('PUBLIC')`으로 좁히고 `is_active` 조건을 지우고 `@Size`를 떼고) **새 테스트만 실패하고 기존 24개는 전부 통과하는 것**을 관측해 RED를 대신했다. 리뷰가 "그렇게 고쳐도 전부 초록"이라고 한 것이 그대로 재현됐다. `match == null`은 보통의 RED였다(대역 `{"results":[null]}` → **500** 관측 → 가드 → 200). `docs/ai/spec/ai-integration.md` §2의 낡은 줄 2개(패키지 위치 · "Bean 하나")는 **위임 범위 밖이라 고치지 않고 남겼다**(CLAUDE.md 9). `clean check` **370개 통과**(checkstyle·jacoco 포함) (S15P11A705-135) | [BD-39](decisions/BD-39-embedding-profile-in-application-config.md) · [BI-25](implements/BI-25-2026-07-29-personal-search-backend-integration.md) | +| 2026-07-30 | 유실·정지된 AI 처리를 복구하는 재스캔 Scheduler와 FAILED Finalizer를 붙였다(S15P11A705-159). **이 저장소에 스케줄링이 처음 들어온다** — `@Scheduled`가 0건이었다. AI 연동의 실패 경로 네 곳이 모두 *"재스캔이 복구한다"*를 안전망으로 전제하고 있었는데 그 재스캔이 없어, 한 번 실패한 Context가 영구히 `PENDING`으로 남았다(상태만 보면 정상과 구별되지 않는다). Bean을 셋으로 가른 것은 **트랜잭션 프록시 때문**이다 — 한 클래스에 두면 자기 메서드 호출이 프록시를 지나지 않아 트랜잭션 없이 돌고, 그러면 `FOR UPDATE SKIP LOCKED`의 잠금이 조회 직후 풀려 중복 방어가 조용히 사라진다. 명세가 근거로 든 것을 **실측으로 뒤집은 지점이 하나 있다**: "Finalize를 먼저 두는 이유는 방금 `retry_count`를 3으로 올린 행이 곧바로 종결되기 때문"이라는데, `runOnce`의 두 줄을 맞바꿔도 테스트가 통과했다 — 증가가 `updated_at`을 함께 갱신해 그 행이 **만료 상태에서 벗어나** Finalizer 후보 조건에 걸리지 않는다. 창을 실제로 확보하는 것은 순서가 아니라 만료 조건 + `updated_at` 갱신이고, 순서는 심층 방어로 남겨 `InOrder` 단위 테스트로 고정했다(그 사실을 BI-28에 적었다). `SKIP LOCKED`는 주장으로 두지 않고 **다른 커넥션이 행을 붙잡은 채 회차를 돌려** 실제로 건너뛰는지 봤다 — 없으면 테스트가 매달리므로 별 스레드 + 15초 타임아웃으로 실패로 드러나게 했다. 함정 둘: `@Scheduled(fixedDelayString)`은 Boot의 완화된 바인딩을 쓰지 않아 `5m`이면 기동이 실패한다(`PT5M`로 두고 `ConfigurationContractTests`가 고정), Spring은 스케줄러를 **작업별로 고르지 않아** 전용 스케줄러라도 Bean 이름 `taskScheduler`를 점유해야 해석이 확정된다(BD-40). RED 6건 확인. `clean check` 386개 통과 | [BD-40](decisions/BD-40-scheduling-with-dedicated-scheduler-and-no-distributed-lock.md) · [BI-28](implements/BI-28-2026-07-30-ai-rescan-scheduler.md) · [패키지 구조](../development/package-structure.md) | diff --git a/docs/backend/decisions/BD-40-scheduling-with-dedicated-scheduler-and-no-distributed-lock.md b/docs/backend/decisions/BD-40-scheduling-with-dedicated-scheduler-and-no-distributed-lock.md new file mode 100644 index 00000000..5d5b22ff --- /dev/null +++ b/docs/backend/decisions/BD-40-scheduling-with-dedicated-scheduler-and-no-distributed-lock.md @@ -0,0 +1,74 @@ +# BD-40. 스케줄링을 `@EnableScheduling` + 자체 `taskScheduler` Bean으로 들이고, 다중 인스턴스 조정은 분산 락 없이 `SKIP LOCKED`에 맡긴다 + +- **상태**: Accepted +- **날짜**: 2026-07-30 +- **관련**: S15P11A705-159, + [BI-28](../implements/BI-28-2026-07-30-ai-rescan-scheduler.md), + [BD-17](BD-17-async-without-message-queue.md)(메시지 큐 없이 비동기), + [BD-36](BD-36-pending-insert-in-transaction-process-call-after-commit.md), + AI 파트 소유 명세 `docs/ai/spec/ai-rescan-scheduler.md` §3·§4.2, + 공용 계약 `Team-PinLog/docs` `static/05_AI_설계.md` §10.3·§10.4 + +## 맥락 + +**이 저장소에 스케줄링이 없었다.** `@Scheduled`·`@EnableScheduling`이 0건이다. 재스캔 회차가 첫 스케줄 +작업이므로, 이 티켓은 "재스캔을 만든다"에 더해 **"이 레포가 주기 작업을 어떻게 돌리는가"**를 함께 +정하게 된다. 뒤에 붙는 배치가 이 자리를 그대로 물려받으므로 되돌리기 비용이 코드 몇 줄이 아니다. + +정해야 할 것이 둘이다. + +**하나. 무엇이 회차를 돌리는가.** Boot는 `@EnableScheduling`이 켜지면 `TaskSchedulingAutoConfiguration`이 +**단일 스레드** `taskScheduler`를 만들어 준다. AI 파트 명세 §3은 그것을 쓰지 말고 "전용 +`ThreadPoolTaskScheduler`"를 두라고 정한다. + +**둘. 인스턴스가 여러 대일 때 같은 행을 두 번 집는 것을 무엇이 막는가.** 명세 §3은 "리더 선출을 두지 +않는다"고 정하고, 그 근거로 후보 선택이 `FOR UPDATE SKIP LOCKED` 기반이라는 점을 든다. 우리 쪽에서 +확인할 것은 그 전제가 실제로 성립하느냐다 — 잠금은 트랜잭션이 살아 있는 동안만 유효하므로, 후보 +조회를 트랜잭션 밖에서 부르면 이 결정은 근거를 잃는다. + +## 선택지 + +| 안 | 장점 | 단점 | +|---|---|---| +| (a) `@EnableScheduling` + 자체 `taskScheduler` Bean, 조정은 `SKIP LOCKED` | 의존성이 늘지 않는다. 조정이 이미 필요한 DB 트랜잭션 안에서 공짜로 끝난다. 인스턴스가 서로 다른 행을 집으므로 처리량이 대수로 늘어난다 | `taskScheduler`라는 Bean 이름을 우리가 점유한다. 조정의 정확성이 SQL 한 구절(`SKIP LOCKED`)과 트랜잭션 경계에 달려 있다 | +| (b) `@EnableScheduling` + Boot 기본 스케줄러 | 코드가 없다 | 스레드 하나를 앞으로 생기는 모든 배치가 공유한다. 한 회차가 배치 크기(100)만큼의 HTTP 호출을 순차로 내보내므로 **다른 배치를 굶긴다.** 명세 §3이 명시적으로 거부한다 | +| (c) ShedLock 등 분산 락 라이브러리 | "한 번에 한 인스턴스"가 라이브러리 계약으로 보장된다 | 의존성과 락 테이블이 늘고, **처리량이 인스턴스 수와 무관해진다**(항상 한 대만 일한다). 무엇보다 `SKIP LOCKED`가 이미 해결한 문제를 한 번 더 해결한다 | +| (d) Kubernetes CronJob으로 외부화 | 애플리케이션에 스케줄링을 안 들인다 | Core 도메인 코드(조립기·클라이언트)를 재사용할 수 없어 재구현이 필요하다. 인프라 소관이 늘고 배포 단위가 둘이 된다 | +| (e) 작업별 전용 스케줄러 | 이름이 사실과 일치한다 | **Spring이 지원하지 않는다.** 스케줄러 선택은 `ScheduledTaskRegistrar` 단위이지 작업 단위가 아니다 | + +## 결정 + +**(a).** 성격이 두 층으로 갈린다. + +- **리더 선출을 두지 않는 것은 제약으로 주어졌다.** AI 파트 명세 §3이 정한다. 제약 안에서 우리가 한 + 선택은 그 전제를 코드로 지키는 방식이다 — 후보 조회를 `@Transactional` 서비스 안에서만 부르고, + 외부 호출은 그 트랜잭션이 커밋된 **뒤에** 한다(명세 §3.1의 커밋 경계). 잠금은 트랜잭션과 수명이 + 같으므로 이 경계가 곧 조정 장치다. 그래서 스케줄러 진입점과 후보 잠금이 **서로 다른 Bean**이다 — + 같은 Bean이면 자기 메서드 호출이 프록시를 지나지 않아 트랜잭션 없이 돌고, 조정이 조용히 사라진다. +- **자체 `taskScheduler` Bean은 능동적 선택이다.** (b)·(c)·(d)를 두고 골랐다. + +Bean 이름을 `taskScheduler`로 둔 것도 선택이다. (e)가 불가능하므로 어차피 애플리케이션 전체가 하나를 +공유하는데, 이름을 다르게 두면 **타입이 유일할 때만** 해석되고 두 번째 `TaskScheduler` Bean이 생기는 +순간 조용히 로컬 단일 스레드 실행자로 떨어진다. 이름을 점유해 해석을 확정적으로 만들고, 대신 그 상태를 +테스트로 고정했다(`TaskScheduler` Bean이 하나뿐임을 단언한다). + +## 결과 + +- 이 결정으로 감수하는 것: + - **Boot의 스케줄링 자동설정이 물러난다.** `spring.task.scheduling.*` 프로퍼티가 더 이상 효력이 + 없다. 풀 크기·스레드 이름을 바꾸려면 `AiIntegrationConfig`를 고쳐야 한다. + - **스레드 이름 `ai-rescan-`이 언젠가 거짓말이 된다.** 두 번째 배치가 붙으면 같은 풀을 공유한다. + 그때 이름을 중립적으로 바꾸거나 그 배치에 자기 registrar를 주는 판단이 필요하다. + - **조정의 정확성이 SQL 한 구절에 달려 있다.** `SKIP LOCKED`를 잃으면 후보 조회가 잠금을 기다리며 + 멈추고, `FOR UPDATE` 자체를 잃으면 두 인스턴스가 같은 Context를 중복 처리한다. 둘 다 조용한 + 실패다 — 그래서 다른 커넥션이 행을 붙잡은 채 회차를 돌려 실제로 건너뛰는지 보는 테스트를 두었다. + - **처리량 상한이 인스턴스 수 × 배치 크기다.** 락 기반((c))보다 낫지만, 한 인스턴스가 병렬로 + 처리하지는 않는다(회차 안에서 HTTP 호출은 순차다). + - **`@EnableScheduling`은 전역이라 모든 `@SpringBootTest`가 스케줄러를 함께 띄운다.** 테스트에서 + 주기를 늘려 배경 실행을 없앴다. 끄지 않은 이유는 등록 자체가 검증 대상이기 때문이다. +- 재검토 트리거(이 조건이 오면 다시 논의): + - **두 번째 스케줄 배치가 생길 때.** 풀 공유·스레드 이름·registrar 분리를 그 시점에 함께 본다. + - **회차 소요가 주기(5분)에 근접할 때.** `fixedDelay`라 겹치지는 않지만 만료 행이 밀린다는 뜻이며, + 회차 안 병렬화나 배치 크기 조정이 먼저다. + - **인스턴스를 늘려도 잔량이 줄지 않을 때.** `SKIP LOCKED`가 기대대로 분산되지 않는다는 신호다. + - **배치를 애플리케이션 밖으로 빼야 할 때**(예: 배포 단위 분리). (d)를 다시 본다. diff --git a/docs/backend/decisions/README.md b/docs/backend/decisions/README.md index 298b173d..dae43c74 100644 --- a/docs/backend/decisions/README.md +++ b/docs/backend/decisions/README.md @@ -101,3 +101,4 @@ | [BD-37](BD-37-ai-derived-invalidation-inside-deletion-transaction.md) | AI 파생 데이터 무효화를 삭제 트랜잭션 안에서 백엔드가 직접 쓴다 | Accepted | S15P11A705-124 | | [BD-38](BD-38-published-predicate-per-layer.md) | 발행 여부 판정은 층마다 자체 보유하고(아홉 곳), 대가를 각 자리를 지키는 테스트로 갚는다 | Accepted | S15P11A705-148 | | [BD-39](BD-39-embedding-profile-in-application-config.md) | Embedding Profile을 `application.yml` 리터럴로 두고 환경변수는 덮어쓰기로만 — 런타임 대조가 성립하려면 Spring도 값을 가져야 한다 | Accepted | S15P11A705-135 | +| [BD-40](BD-40-scheduling-with-dedicated-scheduler-and-no-distributed-lock.md) | 스케줄링을 `@EnableScheduling` + 자체 `taskScheduler` Bean으로 들이고, 다중 인스턴스 조정은 분산 락 없이 `SKIP LOCKED`에 맡긴다 | Accepted | S15P11A705-159 | diff --git a/docs/backend/implements/BI-28-2026-07-30-ai-rescan-scheduler.md b/docs/backend/implements/BI-28-2026-07-30-ai-rescan-scheduler.md new file mode 100644 index 00000000..2f34d8a5 --- /dev/null +++ b/docs/backend/implements/BI-28-2026-07-30-ai-rescan-scheduler.md @@ -0,0 +1,222 @@ +# BI-28. 재스캔 Scheduler와 FAILED Finalizer + +- **상태**: ✅ 완료 +- **날짜**: 2026-07-30 +- **관련**: S15P11A705-159, + [BD-40](../decisions/BD-40-scheduling-with-dedicated-scheduler-and-no-distributed-lock.md), + [BI-22](BI-22-2026-07-29-context-ai-enqueue.md)(접수 방향), + [BI-23](BI-23-2026-07-29-ai-derived-invalidation-on-delete.md)(무효화 방향), + AI 파트 소유 명세 `docs/ai/spec/ai-rescan-scheduler.md`, + 공용 계약 `Team-PinLog/docs` `static/05_AI_설계.md` §10.3·§10.4 + +정책의 정본은 AI 파트가 소유한 `docs/ai/spec/ai-rescan-scheduler.md`다. 이 문서는 **Spring에서 어떻게 +구현했고 무엇을 검증했는가**만 다룬다. + +## 왜 필요했나 + +AI 연동의 실패 경로 **네 곳이 모두 "재스캔이 복구한다"를 안전망으로 전제**하는데 그 재스캔이 없었다. + +| 실패 지점 | 코드가 하는 말 | +|---|---| +| `AiIntegrationConfig` 큐 포화 | *"버려진 요청은 PENDING으로 남아 재스캔 대상이 되므로 유실이 아니다"* | +| `AiProcessClient` (401 포함 전부 삼킴) | *"PENDING이 남아 있으므로 재스캔이 같은 Context를 다시 집는다"* | +| `ContextAiRequestedListener` 조립 실패 | *"PENDING은 이미 커밋되어 있어 재스캔이 같은 Context를 다시 집는다"* | +| FastAPI가 `202` 이후 내부에서 실패 | 통보 경로가 없다 | + +증상은 **한 번 실패한 Context가 영구히 `PENDING`으로 남는 것**이다. 상태만 보면 "처리 대기 중"이라 +정상과 구별되지 않는다. + +## 산출 + +### 한 회차 + +```text +AiRescanScheduler#runOnce @Scheduled(fixedDelayString = "${pinlog.ai.rescan.interval}") + 1. AiFailedFinalizer#finalizeExpired @Transactional ← retry_count >= 3 만료 → FAILED + 2. AiRescanCandidateService#claimStale @Transactional ← retry_count < 3 만료 잠금 + retry 증가 + ────────── 커밋 ────────── + 3. ContextProcessRequestAssembler#assemble(contextId) ← Core 재조회 = 삭제 확인 + 4. AiProcessClient#process ← 트랜잭션 밖 +``` + +**Bean이 셋으로 갈린 것은 트랜잭션 프록시 때문이다.** 한 클래스에 두면 자기 메서드 호출이 프록시를 +지나지 않아 트랜잭션 없이 돌고, 그러면 `FOR UPDATE SKIP LOCKED`의 잠금이 조회 직후 풀려 **중복 방어가 +조용히 사라진다.** 스케줄러 자신은 트랜잭션을 열지 않는다 — 열면 뒤의 HTTP 호출 시간만큼 행 잠금과 DB +커넥션이 붙잡힌다(명세 3.1의 커밋 경계). + +### `AiRescanStateRepository` — 만료 술어를 하나로 공유한다 + +재스캔(4.1)과 Finalizer(6.1)가 **같은 만료 술어**를 쓴다. 명세 6.1이 "만료 기준은 재스캔과 동일하다"고 +정하므로 두 SQL에 따로 적으면 한쪽만 고쳐 기준이 갈라진다. `String.formatted`로 한 상수를 두 쿼리에 +끼운다. + +```sql +WHERE retry_count < :maxRetry -- Finalizer는 >= + AND ( + (embedding_status = 'PENDING' AND updated_at < now() - make_interval(secs => :pendingSecs)) + OR (keyword_status = 'PENDING' AND updated_at < now() - make_interval(secs => :pendingSecs)) + OR (embedding_status = 'PROCESSING' AND updated_at < now() - make_interval(secs => :procSecs)) + OR (keyword_status = 'PROCESSING' AND updated_at < now() - make_interval(secs => :procSecs)) + ) +ORDER BY updated_at +LIMIT :batchSize +FOR UPDATE SKIP LOCKED +``` + +- **시각 비교를 DB `now()`로 한다.** 자바에서 컷오프를 계산해 넘기면 인스턴스마다 시계가 달라 만료 + 시점이 갈라진다. 만료값은 `make_interval(secs => ...)`로 넘기고 파라미터를 `double`로 보낸다 — 그 + 함수의 인자 타입이 `double precision`이라 값도 그 타입이면 함수 해석에 추론이 끼어들지 않는다. +- **`FAILED`·`CANCELLED`를 조건에 적지 않는다.** 만료 술어가 `PENDING`·`PROCESSING` 화이트리스트라 + 자동으로 빠진다. 블랙리스트(`<> 'COMPLETED'`)로 쓰면 나중에 추가되는 status가 조용히 후보가 된다. +- `retry_count` 증가는 status를 손대지 않는다. **만료된 `PROCESSING`을 `PENDING`으로 되돌리지 + 않는다** — Spring은 `PROCESSING`을 쓰지도 해제하지도 않고, FastAPI의 선점 UPDATE가 `PROCESSING`을 + 허용 조건에 포함하므로 재요청만으로 재개된다(명세 5장). 되돌리면 두 주체가 같은 컬럼을 경쟁적으로 + 쓰게 되어 소유권 경계가 무너진다. +- 종결 UPDATE의 `CASE`는 **화이트리스트**다. 후보를 잡은 뒤 UPDATE 직전에 삭제·교체가 끼어들 수 있고, + 블랙리스트면 그 창에서 `CANCELLED`가 `FAILED`로 뒤집혀 "사용자가 지웠다"와 "AI가 실패했다"를 구별할 + 수 없게 된다(명세 6.3). + +### `ContextProcessRequestAssembler` — 진입점을 하나 더 두었다 + +재스캔은 상태 행에서 후보를 집으므로 `context_id`뿐이다. `assemble(long contextId)`를 더해 `member_id`· +`record_id`를 Context의 비정규화 컬럼에서 읽는다. **기존 이벤트 경로(`assemble(ContextAiRequested)`)의 +동작은 바꾸지 않았다** — 그쪽은 발행 시점의 값을 그대로 쓰고, 이미 검증된 경로다. 공통 본문만 private +메서드로 뺐다. + +`Context`에 `@SQLRestriction("deleted_at IS NULL")`이 걸려 있어 **이 재조회가 곧 삭제 확인이다**(명세 +5.1). 삭제된 것과 수정으로 교체된 구버전이 같은 경로로 함께 빠지므로 둘을 구분하지 않는다. + +### `warnUnlessCancelled` — 명세 5.1의 정합성 경고 + +Context가 사라졌으면 그 상태는 `CANCELLED`여야 한다. 아니라면 삭제·수정 트랜잭션이 무효화를 빠뜨렸다는 +뜻이므로 경고를 남긴다. **후보 조회 결과로는 판정할 수 없다** — 후보 조회는 `CANCELLED`를 애초에 잡지 +않으므로(화이트리스트), 삭제가 그 뒤에 일어났을 가능성을 보려면 그 시점의 값을 다시 읽어야 한다. +그래서 단건 조회 메서드가 하나 더 있다. + +### 스케줄링 — 이 레포에 처음 들어온다 + +`@EnableScheduling`을 `AiIntegrationConfig`에 두었다. `@EnableAsync`와 같은 사정이다 — 둘 다 전역 +스위치인데 켜야 하는 이유가 이 연동에만 있고, `global/config`로 올리면 스위치와 유일한 소비자가 떨어져 +앉는다. 전용 `ThreadPoolTaskScheduler`(`ai-rescan-`, 풀 2)를 두어 Boot의 기본 단일 스레드 스케줄러를 +쓰지 않는다. Bean 이름을 `taskScheduler`로 점유한 이유와 감수하는 것은 BD-40에 있다. + +`domain/ai/scheduler` 패키지를 새로 만들었고 `package-structure.md`의 `ai` 행을 먼저 갱신했다. +`event`와 가른 기준은 **무엇이 호출을 촉발하는가**이고, 그에 따라 트랜잭션 경계도 다르다. + +`AiRescanProperties`는 그 패키지가 아니라 `service`에 있다. 처음에 `scheduler`에 두었는데 **소비자 +둘이 모두 `service`에 있어 패키지 의존이 순환했다**(`service` → `scheduler` → `service`). 설정은 +소비자와 같은 패키지에 둔다는 규약(`package-structure.md`)을 따르면 순환도 함께 사라진다 — 스케줄러가 +읽는 것은 주기 하나이고, 그것은 애노테이션 문자열이라 타입 참조가 아니다. + +### 설정 + +```yaml +pinlog: + ai: + rescan: + interval: PT5M # ← ISO-8601. 아래 함정 참조 + pending-expiry: 5m + processing-expiry: 10m + max-retry: 3 + batch-size: 100 +``` + +`interval`만 ISO-8601이다. 이 값은 Boot의 완화된 바인딩이 아니라 **`@Scheduled(fixedDelayString)`이 +직접 파싱**하고, 그쪽은 숫자(밀리초)나 ISO-8601만 받는다. 다른 키들처럼 `5m`으로 적으면 +`NumberFormatException`으로 **기동이 실패한다.** `AiRescanProperties`에는 같은 값이 `Duration`으로도 +들어 있다 — 애노테이션의 문자열이 무슨 값인지 타입으로 드러나지 않기 때문이다. + +`max-retry`의 정본은 **DB의 `CHECK (retry_count BETWEEN 0 AND 3)`**(`V100__ai_tables.sql`)이다. 이 +값만 올리면 증가 UPDATE가 제약 위반으로 실패한다. Backoff는 두지 않는다(명세 2장). + +## 검증 + +`AiRescanSchedulerTests` — PostgreSQL Testcontainers(`pgvector/pgvector:0.8.5-pg16`), 테스트 13개. +`AiRescanSchedulerOrderTest` — 대역·리플렉션, 2개. `ConfigurationContractTests`에 파라미터 계약 1개. + +### 만료를 만드는 방법과 임계값을 덮는 이유 + +5분을 기다릴 수 없으므로 `updated_at`을 과거로 밀어 넣는다. 임계값을 기본값(5분·10분)이 아닌 +**2분·4분으로 덮고** 세 행의 나이를 그 사이에 배치해 한 회차에서 갈라지는 것을 본다 — 임계값이 코드 +상수라면 이 배치가 성립하지 않으므로, 이 테스트가 곧 "설정으로 주입된다"의 확인이다. + +기본값 자체는 어느 통합 테스트도 지키지 못한다(각자 덮으므로). 그래서 +`ConfigurationContractTests`가 `application.yml`의 다섯 값을 파일로 고정한다. + +### 회차마다 모든 상태 행의 나이를 0으로 돌린다 + +컨테이너는 JVM이 공유하고 이 클래스는 실제로 커밋하므로, 앞선 테스트가 남긴 만료 행이 다음 회차에 +후보로 섞인다. `@BeforeEach`·`@AfterEach`에서 `UPDATE ai.context_ai_state SET updated_at = now()`를 +돌려 앞뒤로 격리했다 — 뒤도 정리하는 이유는 다른 테스트 클래스의 배경 회차에 만료 행을 물려주지 않기 +위해서다. + +Core에 없는 `context_id`는 **음수**로 만든다. `ai.context_ai_state`에는 `core.context`로 향하는 FK가 +없어(V100) 그런 행을 만들 수 있고, 음수면 IDENTITY가 만드는 실제 id와 절대 겹치지 않는다. 상태 전이만 +보는 테스트는 Record·Place 조립이 필요 없다. + +### `SKIP LOCKED`를 실제로 관측한다 + +리더 선출을 두지 않는 근거가 이 동작이므로 주장으로 남기지 않았다. 다른 커넥션(`DriverManager`)이 한 +행을 `FOR UPDATE`로 붙잡은 채 회차를 돌리고, 그 행의 `retry_count`가 그대로이며 다른 만료 행은 처리된 +것을 확인한다. + +**회차를 별 스레드에서 돌린다.** `SKIP LOCKED`가 빠지면 후보 조회가 잠금을 기다리며 멈추는데, 같은 +스레드에서 부르면 테스트가 실패하는 대신 영원히 매달린다. 15초 타임아웃을 걸어 **실패로** 드러나게 +했다. + +### "Finalize를 먼저" — 명세의 근거가 실측에서 관측되지 않았다 + +명세 3.1은 Finalize를 먼저 두는 이유를 *"나중에 두면 같은 회차에서 방금 `retry_count`를 3으로 올린 행을 +곧바로 FAILED로 종결한다"*로 든다. 그것을 결과로 고정하려고 `retry_count = 2` 만료 행으로 테스트를 +썼는데, **`runOnce`의 두 줄을 맞바꿔도 통과했다.** + +원인은 재시도 증가가 `updated_at`을 함께 갱신하는 것이다. 증가 직후 그 행은 **만료 상태에서 벗어나** +Finalizer 후보 조건(6.1의 만료 조건)에 걸리지 않는다. 즉 마지막 재시도의 창을 실제로 확보하는 것은 +순서가 아니라 **Finalizer의 만료 조건 + 증가 시 `updated_at` 갱신**이다. 명세 6.1도 그 둘이 "함께 그 +창을 확보한다"고 쓰고 있으니 명세와 어긋나는 관측은 아니지만, **순서만으로 그 창이 생긴다는 읽기는 +사실이 아니다.** + +그래서 검증을 셋으로 나눴다. + +| 고정하는 것 | 어디서 | +|---|---| +| 3회차 요청이 실제로 나가고 그 회차에서 종결되지 않는다 | `theLastRetryActuallyGoesOutAndIsNotFinalizedInTheSameRound` | +| Finalizer에 만료 조건이 붙어 있다(창을 만드는 실제 장치) | `anExhaustedRowThatIsNotExpiredYetIsLeftAlone` | +| 호출 순서 자체(심층 방어) | `AiRescanSchedulerOrderTest` — Mockito `InOrder` | + +순서를 심층 방어로 남기는 이유: 누군가 증가 UPDATE에서 `updated_at` 갱신을 빼면 순서가 유일한 보호가 +된다. + +### RED 확인 + +구현을 하나씩 되돌려 테스트가 실제로 잡는지 확인했다. + +| 되돌린 것 | 실패한 테스트 | +|---|---| +| `runOnce`의 Finalize를 재스캔 뒤로 | `finalizeRunsBeforeTheCandidateClaimInEveryRound` (**통합 테스트는 통과 — 위 절 참조**) | +| Finalizer 후보 조회에서 만료 조건 제거 | `anExhaustedRowThatIsNotExpiredYetIsLeftAlone` | +| `FOR UPDATE SKIP LOCKED` → `FOR UPDATE` | `skipLockedLeavesALockedRowToWhoeverHoldsIt` (15초 타임아웃) | +| 종결 `CASE`를 화이트리스트 → 블랙리스트(`<> 'COMPLETED'`) | `theFinalizerKeepsCompletedStagesAndNeverOverwritesCancelled` (`CANCELLED`가 `FAILED`로 뒤집힌다) | +| `PROCESSING` 만료값을 `PENDING`과 같게 | `expiryThresholdsComeFromConfigurationAndDifferByStage` | +| `fixedDelayString` → `fixedRateString` | `theRoundIsScheduledWithFixedDelayAndReadsTheConfiguredInterval` | + +### 실행 결과 + +```text +./gradlew clean check --no-daemon → BUILD SUCCESSFUL +386 tests, 0 failures, 0 skipped +jacoco LINE 96.45% BRANCH 82.59% (게이트 80%) +``` + +## 남은 것 + +- **관측·알림 파이프라인**(명세 8장). 회차별 후보·종결 건수는 `INFO` 로그 한 줄로만 남는다. 지표 + (`retry_count` 분포, 상태별 잔량, `retry_count >= 3`이면서 `PENDING`/`PROCESSING`인 잔량)는 없다. + 마지막 항목이 Finalizer의 건강 상태이며 정상이면 한 주기 안에 0으로 수렴해야 한다. +- **`FAILED` 전환 로그에 "마지막 실패 사유"가 없다.** 명세 6.3이 요구하지만 **그 값이 Spring 쪽에 + 존재하지 않는다** — 202 이후의 실패는 FastAPI 내부에서 일어나고 통보 경로가 없다. 가진 단서(종결 직전 + 단계별 status·소진한 재시도 횟수)를 남기고 사유를 어디서 찾아야 하는지 문장으로 가리켰다. 사유를 + 실을 수 있으려면 FastAPI가 실패를 어딘가 기록해야 하고, 그것은 `S15P11A705-121`(`ai#44`) 소관이다. +- **회차 안 HTTP 호출은 순차다.** 배치 100건이 전부 재요청 대상이면 한 회차가 길어진다. `fixedDelay`라 + 겹치지는 않지만 만료 행이 밀린다. 그 상황이 실제로 보이면 BD-40의 재검토 트리거를 따른다. +- **스레드 이름 `ai-rescan-`은 두 번째 배치가 붙는 순간 거짓말이 된다.** BD-40의 감수 목록에 있다. diff --git a/docs/backend/implements/README.md b/docs/backend/implements/README.md index 40e2071a..b264edd6 100644 --- a/docs/backend/implements/README.md +++ b/docs/backend/implements/README.md @@ -41,3 +41,4 @@ | BI-24 | Kakao·Naver 소셜 로그인 추가 — 공급자별 사용자 정보 정규화 (S15P11A705-64) | ✅ 완료 | [BI-24](BI-24-2026-07-29-kakao-naver-login.md) | | BI-25 | 개인 자연어 검색 백엔드 연동 — FastAPI 호출·Core 재검증·Record 단위 조립 (S15P11A705-135) | ✅ 완료 | [BI-25](BI-25-2026-07-29-personal-search-backend-integration.md) | | BI-26 | 운영 Secret 5개를 Infra SealedSecret PR로 전달하는 수동 workflow (S15P11A705-154) | ✅ 완료 | [BI-26](BI-26-2026-07-29-runtime-secret-workflow.md) | +| BI-28 | 재스캔 Scheduler와 FAILED Finalizer — 유실·정지된 AI 처리 복구 (S15P11A705-159) | ✅ 완료 | [BI-28](BI-28-2026-07-30-ai-rescan-scheduler.md) | diff --git a/docs/development/package-structure.md b/docs/development/package-structure.md index c0d5136c..a3b3b590 100644 --- a/docs/development/package-structure.md +++ b/docs/development/package-structure.md @@ -47,7 +47,7 @@ com.pinlog.pinlogback | `feed` | 피드 조회·서빙 API | 관측 로그 `core.feed_event` 테이블은 **AI 소유(V102)** — 재정의 금지, 조회만 | | `auth` | 인증·인가 | **별도 인증 PR에서 생성.** 그 전에는 만들지 않음 | | `search` | 개인 자연어 검색 조회 API | Record를 돌려주지만 `record`에 두지 않습니다 — 진입 경로(`/v1/search/records`)와 조립 규칙(FastAPI 응답의 Core 재검증)이 Record CRUD와 다릅니다. FastAPI 호출 자체는 `ai`가 맡고 이 도메인은 그 결과를 검증·조립만 합니다 | -| `ai` | FastAPI AI Server 연동과 `ai` 스키마 접근 | 애그리거트가 아니라 **파트 경계**입니다. 하위 계층은 `repository`(백엔드가 `ai`에 쓰는 SQL과 응답 조립용 읽기) · `client`(내부 API 호출) · `event`(커밋 이후 훅) · `service`(요청 조립) · `exception`(호출 실패를 도메인 오류로 옮김)이며, `controller`·`entity`는 없습니다 — 외부 진입점이 아니고 남의 스키마를 엔티티로 고정하지 않습니다 | +| `ai` | FastAPI AI Server 연동과 `ai` 스키마 접근 | 애그리거트가 아니라 **파트 경계**입니다. 하위 계층은 `repository`(백엔드가 `ai`에 쓰는 SQL과 응답 조립용 읽기) · `client`(내부 API 호출) · `event`(커밋 이후 훅) · `scheduler`(시간이 촉발하는 훅) · `service`(요청 조립·트랜잭션 경계) · `exception`(호출 실패를 도메인 오류로 옮김)이며, `controller`·`entity`는 없습니다 — 외부 진입점이 아니고 남의 스키마를 엔티티로 고정하지 않습니다. `event`와 `scheduler`를 가른 기준은 **무엇이 호출을 촉발하는가**이고, 그에 따라 트랜잭션 경계도 다릅니다 — `event`는 남의 트랜잭션이 커밋된 뒤에 얹히고, `scheduler`는 자기 트랜잭션을 열고 닫습니다(S15P11A705-159) | > 검토 필요: `member` 행의 "공개 프로필·소개를 포함합니다"는 `docs/static/06_데이터모델_및_무결성.md` 2.1(익명 서비스이므로 저장하는 개인정보가 없다)과 실제 구현체(`domain/member/entity/Member` — 개인정보·프로필 컬럼 없음)에 모두 반합니다. CLAUDE.md 9번 규칙에 따라 임의로 고치지 않고 충돌로 기록만 남깁니다. > 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 9ecbd08b..65fb62a2 100644 --- a/src/main/java/com/pinlog/pinlogback/domain/ai/AiIntegrationConfig.java +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/AiIntegrationConfig.java @@ -8,18 +8,28 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.http.client.SimpleClientHttpRequestFactory; +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.service.AiRescanProperties; + /** * FastAPI 연동에 필요한 인프라 조립. 설정 클래스를 {@code global/config}가 아니라 소비자와 같은 * 패키지에 두는 기준은 {@code docs/development/package-structure.md}의 보안 설정과 같다 — 설정과 * 그 설정이 조립하는 구현이 떨어져 있으면 한쪽만 고치게 된다. + * + *

{@link EnableScheduling}이 여기 있는 것은 {@link EnableAsync}와 같은 사정이다. 둘 다 애플리케이션 + * 전역 스위치인데, 켜야 하는 이유가 이 연동에만 있다. {@code global/config}로 올리면 스위치와 + * 그 스위치의 유일한 소비자가 떨어져 앉는다. */ @Configuration @EnableAsync -@EnableConfigurationProperties(AiProperties.class) +@EnableScheduling +@EnableConfigurationProperties({AiProperties.class, AiRescanProperties.class}) public class AiIntegrationConfig { private static final Logger log = LoggerFactory.getLogger(AiIntegrationConfig.class); @@ -64,6 +74,34 @@ private RestClient restClient(String baseUrl, AiProperties.Timeouts timeouts) { .build(); } + /** + * 재스캔 회차를 돌리는 스케줄러(AI 파트 소유 명세 {@code docs/ai/spec/ai-rescan-scheduler.md} 3장). + * Boot의 기본 단일 스레드 스케줄러를 쓰지 않는다. + * + *

Bean 이름이 {@code taskScheduler}인 것은 의도다. Spring은 스케줄러를 작업별로 고르지 + * 않고 {@code ScheduledTaskRegistrar} 단위로 하나 고르므로, {@code @Scheduled} 메서드가 + * "자기 스케줄러"를 지목할 방법이 없다. 이름을 다르게 두면 타입이 유일할 때만 해석되고, 두 번째 + * {@link TaskScheduler} Bean이 생기는 순간 조용히 로컬 단일 스레드 실행자로 떨어진다. + * + *

따라서 이것은 애플리케이션의 스케줄러이며, 스레드 이름이 {@code ai-rescan-}인 것은 + * 지금 이 작업이 유일한 입주자라는 사실을 스레드 덤프에서 읽히게 하려는 것이다. 두 번째 배치가 + * 생기면 이름을 중립적으로 바꾸거나 그 배치에 자기 registrar를 주는 판단이 함께 필요하다(BD-40). + * + *

스레드가 2개인 이유: 회차는 {@code fixedDelay}라 겹치지 않으므로 하나로 충분하지만, 두 번째 + * 스케줄 작업이 붙었을 때 재스캔 한 회차가 그 작업을 굶기는 것을 막는 여유다. 한 회차는 + * 배치 크기만큼의 HTTP 호출을 순차로 내보내 길어질 수 있다. + */ + @Bean + public TaskScheduler taskScheduler() { + ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler(); + scheduler.setThreadNamePrefix("ai-rescan-"); + scheduler.setPoolSize(2); + // 종료를 기다리지 않는다(기본값). 기다리게 하면 진행 중인 FastAPI 호출의 read-timeout만큼 + // 종료가 늦어져 graceful shutdown 예산(20s)을 잠식한다. 중단된 회차는 다음 회차가 다시 집는다. + scheduler.initialize(); + return scheduler; + } + /** * AI 호출 전용 풀(명세 4.2). 공용 executor를 쓰지 않는 이유는 FastAPI 장애가 다른 비동기 작업까지 * 굶기지 않게 하기 위해서다. diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/repository/AiRescanStateRepository.java b/src/main/java/com/pinlog/pinlogback/domain/ai/repository/AiRescanStateRepository.java new file mode 100644 index 00000000..f5d76513 --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/repository/AiRescanStateRepository.java @@ -0,0 +1,189 @@ +package com.pinlog.pinlogback.domain.ai.repository; + +import java.time.Duration; +import java.util.List; +import java.util.Map; +import java.util.Optional; + +import org.springframework.jdbc.core.RowMapper; +import org.springframework.jdbc.core.namedparam.NamedParameterJdbcTemplate; +import org.springframework.stereotype.Repository; + +/** + * 유실·정지된 AI 처리를 찾아내는 조회와 그 뒤처리 UPDATE + * (AI 파트 소유 명세 {@code docs/ai/spec/ai-rescan-scheduler.md} 4·5·6장). + * + *

후보 조회는 반드시 호출자의 트랜잭션 안에서 부른다. {@code FOR UPDATE SKIP LOCKED}가 + * 잡는 잠금은 트랜잭션이 끝날 때 풀리므로, 트랜잭션 없이 부르면(auto-commit) 조회가 끝나는 즉시 + * 잠금이 사라져 같은 행을 두 인스턴스가 동시에 집는다. 그러면 이 클래스가 존재하는 이유가 + * 없어진다. + * + *

같은 테이블을 쓰는 {@link ContextAiStateRepository}(생성)·{@link AiDerivedDataRepository} + * (삭제 무효화)와 파일을 나눈 기준은 호출 시점과 트랜잭션 경계다(BD-36과 같은 기준). 저 둘은 + * 사용자 요청 트랜잭션에 얹혀 돌고, 이쪽은 스케줄러가 자기 트랜잭션을 열어 부른다. + */ +@Repository +public class AiRescanStateRepository { + + /** + * 만료 판정. 재스캔(4.1)과 Finalizer(6.1)가 이 술어를 공유한다 — 명세 6.1이 "만료 기준은 + * 재스캔과 동일하다"고 정하므로, 두 SQL에 따로 적으면 한쪽만 고쳐 두 경로의 기준이 갈라진다. + * + *

기준 컬럼은 {@code updated_at}이다(명세 4.1). 시각 비교를 DB {@code now()}로 하는 이유는 + * 애플리케이션과 DB의 시계가 어긋나도 판정이 흔들리지 않게 하기 위해서다 — 인스턴스가 여러 대면 + * 각자의 시계로 자른 만료 시점이 서로 달라진다. + * + *

초를 {@code double}로 넘기는 이유: {@code make_interval(secs ...)}의 파라미터 타입이 + * {@code double precision}이라 값도 그 타입으로 보내면 함수 해석에 추론이 끼어들지 않는다. + */ + private static final String EXPIRED_STAGE_PREDICATE = """ + ( + (embedding_status = 'PENDING' AND updated_at < now() - make_interval(secs => :pendingSecs)) + OR (keyword_status = 'PENDING' AND updated_at < now() - make_interval(secs => :pendingSecs)) + OR (embedding_status = 'PROCESSING' AND updated_at < now() - make_interval(secs => :procSecs)) + OR (keyword_status = 'PROCESSING' AND updated_at < now() - make_interval(secs => :procSecs)) + )"""; + + /** + * 재스캔 후보(명세 4.2). + * + *

{@code FOR UPDATE SKIP LOCKED}가 하는 일은 두 가지다. 잠금은 다중 인스턴스·실행 겹침에서 같은 + * Context를 두 번 집는 것을 막고, {@code SKIP LOCKED}는 잠긴 행을 기다리지 않고 건너뛴다 — + * 기다리면 배치 전체가 느린 한 행에 묶인다. 이것이 중복 방어의 유일한 장치는 아니며, FastAPI의 + * {@code PROCESSING} 조건부 UPDATE가 최종 방어선이다(명세 4.2). + * + *

{@code ORDER BY updated_at}으로 가장 오래 멈춘 것부터 처리한다. + * + *

FAILED·CANCELLED를 조건에 적지 않는다. 만료 술어가 {@code PENDING}·{@code PROCESSING} + * 화이트리스트이므로 자동으로 빠진다. 블랙리스트({@code <> 'COMPLETED'} 등)로 쓰면 나중에 추가되는 + * status가 조용히 후보가 된다. + */ + private static final String LOCK_STALE_SQL = """ + SELECT context_id, embedding_status, keyword_status, retry_count + FROM ai.context_ai_state + WHERE retry_count < :maxRetry + AND %s + ORDER BY updated_at + LIMIT :batchSize + FOR UPDATE SKIP LOCKED + """.formatted(EXPIRED_STAGE_PREDICATE); + + /** + * Finalizer 후보(명세 6.2). 재스캔과 {@code retry_count} 비교 방향만 다르다. + * + *

만료 조건이 여기에도 붙는 것이 핵심이다. 없으면 {@code retry_count}를 3으로 올린 그 + * 회차에서 곧바로 종결되어, 마지막 재시도 요청이 처리될 시간을 갖지 못한다(명세 6.1). + */ + private static final String LOCK_EXHAUSTED_SQL = """ + SELECT context_id, embedding_status, keyword_status, retry_count + FROM ai.context_ai_state + WHERE retry_count >= :maxRetry + AND %s + ORDER BY updated_at + LIMIT :batchSize + FOR UPDATE SKIP LOCKED + """.formatted(EXPIRED_STAGE_PREDICATE); + + /** + * 재시도 소진(명세 5장). status는 손대지 않는다 — 만료된 {@code PROCESSING}을 {@code PENDING}으로 + * 되돌리지 않는다. Spring은 {@code PROCESSING}을 쓰지도, 해제하지도 않으며, FastAPI의 선점 UPDATE가 + * {@code PROCESSING}을 허용 조건에 포함하므로 재요청만으로 재개된다. Spring이 상태를 손대면 두 + * 주체가 같은 컬럼을 경쟁적으로 쓰게 되어 소유권 경계가 무너진다. + */ + private static final String INCREMENT_RETRY_SQL = """ + UPDATE ai.context_ai_state + SET retry_count = retry_count + 1, + updated_at = now() + WHERE context_id IN (:contextIds) + """; + + /** + * 미완료 단계 종결(명세 6.2). + * + *

{@code IN ('PENDING','PROCESSING')} 화이트리스트를 블랙리스트로 바꾸지 않는다. 후보를 + * 잡은 뒤 UPDATE 직전에 Context가 삭제·교체될 수 있고, 그때 블랙리스트({@code <> 'COMPLETED'})면 + * {@code CANCELLED}가 {@code FAILED}로 뒤집힌다. 삭제 표시가 실패로 바뀌면 "사용자가 지웠다"와 + * "AI가 실패했다"를 구별할 수 없게 된다(명세 6.3). + * + *

이미 {@code COMPLETED}인 단계도 그대로 둔다. 부분 성공을 지우면 Embedding을 다시 만들어야 + * 하므로 부분 재사용의 이점이 사라진다. + */ + private static final String FINALIZE_SQL = """ + UPDATE ai.context_ai_state + SET embedding_status = + CASE WHEN embedding_status IN ('PENDING','PROCESSING') THEN 'FAILED' ELSE embedding_status END, + keyword_status = + CASE WHEN keyword_status IN ('PENDING','PROCESSING') THEN 'FAILED' ELSE keyword_status END, + updated_at = now() + WHERE context_id IN (:contextIds) + """; + + private static final String SELECT_ONE_SQL = """ + SELECT context_id, embedding_status, keyword_status, retry_count + FROM ai.context_ai_state + WHERE context_id = :contextId + """; + + private static final RowMapper ROW_MAPPER = (rows, index) -> new ContextAiStateRow( + rows.getLong("context_id"), + rows.getString("embedding_status"), + rows.getString("keyword_status"), + rows.getInt("retry_count")); + + private final NamedParameterJdbcTemplate jdbc; + + public AiRescanStateRepository(NamedParameterJdbcTemplate jdbc) { + this.jdbc = jdbc; + } + + /** 재스캔 후보를 잠근다. 호출자의 트랜잭션 안에서 부른다. */ + public List lockStale(Duration pendingExpiry, Duration processingExpiry, + int maxRetry, int batchSize) { + return jdbc.query(LOCK_STALE_SQL, + candidateParameters(pendingExpiry, processingExpiry, maxRetry, batchSize), ROW_MAPPER); + } + + /** 재시도를 소진한 만료 후보를 잠근다. 호출자의 트랜잭션 안에서 부른다. */ + public List lockRetryExhausted(Duration pendingExpiry, Duration processingExpiry, + int maxRetry, int batchSize) { + return jdbc.query(LOCK_EXHAUSTED_SQL, + candidateParameters(pendingExpiry, processingExpiry, maxRetry, batchSize), ROW_MAPPER); + } + + /** {@code retry_count}를 1 올린다. 상한은 DB {@code CHECK} 제약이 지킨다. */ + public void incrementRetryCount(List contextIds) { + if (contextIds.isEmpty()) { + return; + } + jdbc.update(INCREMENT_RETRY_SQL, Map.of("contextIds", contextIds)); + } + + /** 미완료 단계만 {@code FAILED}로 바꾼다. */ + public void failIncompleteStages(List contextIds) { + if (contextIds.isEmpty()) { + return; + } + jdbc.update(FINALIZE_SQL, Map.of("contextIds", contextIds)); + } + + /** + * 단건 조회. 삭제된 Context의 상태가 {@code CANCELLED}인지 확인하는 정합성 검사에만 쓴다(명세 5.1) + * — 후보 조회는 {@code CANCELLED}를 애초에 잡지 않으므로 그 시점의 값을 다시 읽어야 한다. + */ + public Optional findOne(long contextId) { + return jdbc.query(SELECT_ONE_SQL, Map.of("contextId", contextId), ROW_MAPPER).stream().findFirst(); + } + + private static Map candidateParameters(Duration pendingExpiry, Duration processingExpiry, + int maxRetry, int batchSize) { + return Map.of( + "pendingSecs", seconds(pendingExpiry), + "procSecs", seconds(processingExpiry), + "maxRetry", maxRetry, + "batchSize", batchSize); + } + + private static double seconds(Duration expiry) { + return expiry.toMillis() / 1000.0; + } +} diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/repository/ContextAiStateRow.java b/src/main/java/com/pinlog/pinlogback/domain/ai/repository/ContextAiStateRow.java new file mode 100644 index 00000000..1873399b --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/repository/ContextAiStateRow.java @@ -0,0 +1,21 @@ +package com.pinlog.pinlogback.domain.ai.repository; + +/** + * {@code ai.context_ai_state} 한 행의 재스캔 판정에 필요한 부분. + * + *

엔티티로 매핑하지 않는다 — {@code ai} 스키마의 소유는 AI 파트이고, 남의 스키마를 JPA 엔티티로 + * 고정하면 저쪽의 컬럼 추가가 우리 기동 실패({@code ddl-auto=validate})가 된다(package-structure.md). + * + * @param contextId 대상 Context + * @param embeddingStatus 임베딩 단계 status + * @param keywordStatus Keyword 단계 status. 두 단계는 독립 전이한다 — 한쪽만 + * {@code PENDING}인 행이 정상적으로 존재한다(명세 7장) + * @param retryCount 소진한 재시도 횟수 + */ +public record ContextAiStateRow( + long contextId, + String embeddingStatus, + String keywordStatus, + int retryCount +) { +} diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/repository/package-info.java b/src/main/java/com/pinlog/pinlogback/domain/ai/repository/package-info.java index 224781b4..c5c18f97 100644 --- a/src/main/java/com/pinlog/pinlogback/domain/ai/repository/package-info.java +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/repository/package-info.java @@ -1,5 +1,5 @@ /** - * 백엔드가 {@code ai} 스키마에 수행하는 쓰기. + * 백엔드가 {@code ai} 스키마에 수행하는 쓰기와, 그 쓰기를 판단하기 위한 조회. * *

{@link org.jspecify.annotations.NullMarked}로 선언한다(BD-29과 같은 이유·같은 단위). */ 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 new file mode 100644 index 00000000..12b64b65 --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/scheduler/AiRescanScheduler.java @@ -0,0 +1,140 @@ +package com.pinlog.pinlogback.domain.ai.scheduler; + +import java.util.List; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +import com.pinlog.pinlogback.domain.ai.client.AiProcessClient; +import com.pinlog.pinlogback.domain.ai.client.ContextProcessRequest; +import com.pinlog.pinlogback.domain.ai.repository.ContextAiStateRow; +import com.pinlog.pinlogback.domain.ai.service.AiFailedFinalizer; +import com.pinlog.pinlogback.domain.ai.service.AiRescanCandidateService; +import com.pinlog.pinlogback.domain.ai.service.ContextProcessRequestAssembler; + +/** + * 유실·정지된 AI 처리를 복구하는 한 회차 + * (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 + * 이후 내부에서 실패한 경우. 상태만 보면 "처리 대기 중"이라 정상과 구별되지 않는 것이 이 문제의 + * 성질이다. + * + *

순서가 계약이다. {@code Finalize → 후보 잠금·retry 증가 → 커밋 → Context 재조회 → 삭제 + * 확인 → FastAPI 호출}. Finalize를 먼저 두는 이유는 나중에 두면 같은 회차에서 방금 + * {@code retry_count}를 3으로 올린 행을 곧바로 {@code FAILED}로 종결해, 마지막 재시도가 실행되기도 + * 전에 사망 선고를 내리기 때문이다(명세 3.1·6.1). + * + *

커밋 경계와 외부 호출이 갈리는 지점은 {@link AiRescanCandidateService#claimStale()}의 반환이다. + * 이 클래스 자체는 트랜잭션을 열지 않는다 — 열면 뒤의 HTTP 호출 시간만큼 행 잠금이 붙잡힌다. + */ +@Component +public class AiRescanScheduler { + + private static final Logger log = LoggerFactory.getLogger(AiRescanScheduler.class); + + private static final String CANCELLED = "CANCELLED"; + + private final AiFailedFinalizer finalizer; + private final AiRescanCandidateService candidates; + private final ContextProcessRequestAssembler assembler; + private final AiProcessClient client; + + public AiRescanScheduler(AiFailedFinalizer finalizer, AiRescanCandidateService candidates, + ContextProcessRequestAssembler assembler, AiProcessClient client) { + this.finalizer = finalizer; + this.candidates = candidates; + this.assembler = assembler; + this.client = client; + } + + /** + * 한 회차. + * + *

{@code fixedRate}가 아니라 {@code fixedDelay}다(명세 3장). 한 회차가 배치 크기만큼의 + * HTTP 호출을 순차로 내보내므로 실행 시간이 주기를 넘길 수 있고, {@code fixedRate}면 그때 다음 + * 회차가 겹쳐 돈다. 겹쳐도 {@code SKIP LOCKED} 덕에 같은 행을 두 번 집지는 않지만, 밀린 회차가 + * 계속 쌓여 FastAPI에 부하를 더한다. + * + *

주기 값을 애노테이션에 문자열로 두는 것은 {@code @Scheduled}의 제약이다. 타입 있는 값은 + * {@link com.pinlog.pinlogback.domain.ai.service.AiRescanProperties#interval()}에 같은 키로 있다. + * + *

예외를 잡지 않는다. Spring의 기본 오류 처리기가 로그를 남기고 {@code fixedDelay} 일정은 + * 유지되므로, 여기서 삼키면 실패가 회차 집계 로그에 "0건 처리"로 위장될 뿐이다. + */ + @Scheduled(fixedDelayString = "${pinlog.ai.rescan.interval}") + public void runOnce() { + int finalized = finalizer.finalizeExpired().size(); + List claimed = candidates.claimStale(); + int requested = 0; + for (ContextAiStateRow candidate : claimed) { + if (requestProcessing(candidate)) { + requested++; + } + } + if (finalized > 0 || !claimed.isEmpty()) { + log.info("AI 재스캔 회차 종료: FAILED 종결 {}건, 후보 {}건, 재요청 {}건", + finalized, claimed.size(), requested); + } + } + + /** + * 명세 5장의 {@code 재조회 → 삭제 확인 → 호출}. + * + *

이전 요청 본문을 보관했다가 재전송하지 않는다. Context는 불변이라 재조회한 본문은 첫 + * 시도와 같은데도 다시 읽는 이유는 본문을 얻기 위해서가 아니라 그 Context가 아직 살아 있는지 + * 확인하기 위해서다(명세 5.1). 조립기가 소프트 삭제된 Context를 걸러 내므로, 삭제된 것과 수정으로 + * 교체된 구버전이 같은 경로로 함께 빠진다 — 둘을 구분할 필요가 없다. + * + *

조립 실패를 삼키는 이유는 {@code ContextAiRequestedListener}와 같다. 한 후보의 DB 조회 실패로 + * 나머지 후보까지 버리면 회차 전체가 한 행에 묶인다. {@code retry_count}는 이미 커밋됐으므로 + * 이 실패로 예산이 되돌아가지도 않는다. + * + * @return FastAPI에 요청을 보냈는가 + */ + private boolean requestProcessing(ContextAiStateRow candidate) { + long contextId = candidate.contextId(); + try { + ContextProcessRequest request = assembler.assemble(contextId).orElse(null); + if (request == null) { + warnUnlessCancelled(contextId); + return false; + } + client.process(request); + return true; + } catch (RuntimeException e) { + log.warn("재스캔 요청 조립 실패: contextId={}, cause={}", contextId, e.toString()); + return false; + } + } + + /** + * Context가 사라졌는데 상태가 {@code CANCELLED}가 아니면 정합성 경고를 남긴다(명세 5.1). 삭제·수정 + * 트랜잭션이 {@code CANCELLED} 기록을 빠뜨렸다는 뜻이고, 그대로 두면 늦게 도착한 FastAPI 결과가 + * 지워진 Context에 저장된다 — {@code CANCELLED}가 막아야 할 바로 그 누출이다. + * + *

두 단계 중 하나라도 {@code CANCELLED}가 아니면 경고한다. 무효화는 두 컬럼을 함께 + * 쓰므로({@code AiDerivedDataRepository}) 한쪽만 남았다는 것도 같은 종류의 누락이다. + */ + private void warnUnlessCancelled(long contextId) { + ContextAiStateRow current = candidates.findCurrentState(contextId).orElse(null); + if (current == null) { + log.warn("재스캔 후보의 상태 행이 사라졌다: contextId={}. 상태 행 없이는 이 Context를 " + + "다시 집을 수도 없다", contextId); + return; + } + if (CANCELLED.equals(current.embeddingStatus()) && CANCELLED.equals(current.keywordStatus())) { + log.debug("이미 삭제된 Context라 재스캔 호출을 생략한다: contextId={}", contextId); + return; + } + log.warn("삭제된 Context의 상태가 CANCELLED가 아니다(삭제·수정 트랜잭션이 무효화를 빠뜨렸다): " + + "contextId={}, embedding={}, keyword={}", + contextId, current.embeddingStatus(), current.keywordStatus()); + } +} diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/scheduler/package-info.java b/src/main/java/com/pinlog/pinlogback/domain/ai/scheduler/package-info.java new file mode 100644 index 00000000..0e40eb5e --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/scheduler/package-info.java @@ -0,0 +1,15 @@ +/** + * 시간이 기동시키는 AI 연동 작업 — 유실·정지된 처리의 복구. + * + *

{@code event}가 "커밋이 기동시키는 작업"이라면 여기는 "시간이 기동시키는 작업"이다. 두 축을 + * 가른 기준은 무엇이 호출을 촉발하는가이며, 그에 따라 트랜잭션 경계도 다르다 — 저쪽은 남의 트랜잭션이 + * 끝난 뒤에 얹히고, 이쪽은 자기 트랜잭션을 열고 닫는다. + * + *

{@link org.jspecify.annotations.NullMarked}로 선언한다(BD-29과 같은 이유·같은 단위). + * {@code @NullMarked}는 하위 패키지로 상속되지 않으므로 상위 {@code domain.ai}에 선언돼 있어도 + * 여기에 다시 적어야 한다. + */ +@NullMarked +package com.pinlog.pinlogback.domain.ai.scheduler; + +import org.jspecify.annotations.NullMarked; diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/service/AiFailedFinalizer.java b/src/main/java/com/pinlog/pinlogback/domain/ai/service/AiFailedFinalizer.java new file mode 100644 index 00000000..757e67fa --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/service/AiFailedFinalizer.java @@ -0,0 +1,75 @@ +package com.pinlog.pinlogback.domain.ai.service; + +import java.util.List; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +import com.pinlog.pinlogback.domain.ai.repository.AiRescanStateRepository; +import com.pinlog.pinlogback.domain.ai.repository.ContextAiStateRow; + +/** + * 재시도를 소진한 만료 작업의 미완료 단계를 {@code FAILED}로 종결한다 + * (AI 파트 소유 명세 {@code docs/ai/spec/ai-rescan-scheduler.md} 6장). + * + *

선택 사항이 아니다. {@code process} 호출의 응답은 {@code 202 Accepted}이고 그것은 접수만 + * 뜻한다. 완료 통보용 웹훅이 없으므로 202 이후 FastAPI 내부에서 일어난 실패를 Spring은 알 수 없고, + * {@code retry_count >= 3}인 행은 재스캔 후보 조건({@code retry_count < 3})에서 제외되기만 할 뿐 + * PROCESSING으로 영원히 남는다. 그러면 상태 지표상 "처리 중"으로 오인되어 장애 관측 자체가 + * 불가능해진다. + * + *

FastAPI를 호출하지 않는다. 순수한 상태 정리 단계다. + * + *

{@link com.pinlog.pinlogback.domain.ai.scheduler.AiRescanScheduler}와 별 Bean인 이유는 + * {@code @Transactional}이 프록시로 걸리기 때문이다. 스케줄러가 자기 메서드를 직접 부르면 프록시를 + * 지나지 않아 트랜잭션이 열리지 않고, 그러면 {@code FOR UPDATE}가 잡은 잠금이 조회 직후 풀려 + * 다중 인스턴스에서 같은 행을 두 번 종결한다. + */ +@Service +public class AiFailedFinalizer { + + private static final Logger log = LoggerFactory.getLogger(AiFailedFinalizer.class); + + private final AiRescanStateRepository repository; + private final AiRescanProperties properties; + + public AiFailedFinalizer(AiRescanStateRepository repository, AiRescanProperties properties) { + this.repository = repository; + this.properties = properties; + } + + /** + * 후보를 잠그고 미완료 단계를 종결한다. 스케줄러 회차의 첫 단계로 부른다(명세 3.1) — + * 나중에 두면 같은 회차에서 방금 {@code retry_count}를 3으로 올린 행을 곧바로 종결해, 마지막 + * 재시도가 실행되기도 전에 사망 선고를 내린다. + * + * @return 종결한 행들의 종결 직전 상태. 로그와 회차 집계에 쓴다 + */ + @Transactional + public List finalizeExpired() { + List exhausted = repository.lockRetryExhausted( + properties.pendingExpiry(), properties.processingExpiry(), + properties.maxRetry(), properties.batchSize()); + if (exhausted.isEmpty()) { + return List.of(); + } + repository.failIncompleteStages(exhausted.stream().map(ContextAiStateRow::contextId).toList()); + exhausted.forEach(AiFailedFinalizer::logFinalized); + return exhausted; + } + + /** + * 명세 6.3은 "{@code contextId}와 마지막 실패 사유"를 남기라고 하는데, 사유는 Spring 쪽에 + * 존재하지 않는다. 202 이후의 실패는 FastAPI 내부에서 일어나고 통보 경로가 없다. 그래서 우리가 + * 가진 단서 — 종결 직전의 단계별 status와 소진한 재시도 횟수 — 를 남기고, 사유를 어디서 찾아야 + * 하는지 문장으로 가리킨다. 이것이 Preset·모델 설정 문제를 판별할 출발점이다. + */ + private static void logFinalized(ContextAiStateRow row) { + log.warn("AI 처리를 FAILED로 종결한다: contextId={}, embedding={}, keyword={}, retryCount={}. " + + "재시도를 모두 소진하고 만료됐다 — 실패 사유는 FastAPI 로그에서 확인해야 한다" + + "(202 이후의 실패는 Spring에 통보되지 않는다).", + row.contextId(), row.embeddingStatus(), row.keywordStatus(), row.retryCount()); + } +} diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/service/AiRescanCandidateService.java b/src/main/java/com/pinlog/pinlogback/domain/ai/service/AiRescanCandidateService.java new file mode 100644 index 00000000..600ce9f3 --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/service/AiRescanCandidateService.java @@ -0,0 +1,61 @@ +package com.pinlog.pinlogback.domain.ai.service; + +import java.util.List; +import java.util.Optional; + +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +import com.pinlog.pinlogback.domain.ai.repository.AiRescanStateRepository; +import com.pinlog.pinlogback.domain.ai.repository.ContextAiStateRow; + +/** + * 재스캔 후보를 잠그고 재시도 예산을 소진시킨다 + * (AI 파트 소유 명세 {@code docs/ai/spec/ai-rescan-scheduler.md} 4장·5장). + * + *

이 클래스의 트랜잭션 경계가 명세 3.1의 커밋 경계다. 후보 잠금과 {@code retry_count} 증가는 + * 한 트랜잭션에서 끝내고, FastAPI 호출은 이 메서드가 반환한 뒤 트랜잭션 밖에서 한다. 외부 호출을 + * 안에 두면 호출 지연만큼 행 잠금과 DB 커넥션이 붙잡혀 있게 된다. + * + *

같은 이유로 호출자와 별 Bean이다. {@code @Transactional}은 프록시로 걸리므로 스케줄러가 자기 + * 메서드를 부르면 트랜잭션이 열리지 않고, 그러면 {@code FOR UPDATE SKIP LOCKED}의 잠금이 조회 직후 + * 풀려 중복 방어가 사라진다. + */ +@Service +public class AiRescanCandidateService { + + private final AiRescanStateRepository repository; + private final AiRescanProperties properties; + + public AiRescanCandidateService(AiRescanStateRepository repository, AiRescanProperties properties) { + this.repository = repository; + this.properties = properties; + } + + /** + * 만료된 미완료 작업을 잠그고 {@code retry_count}를 올린다. + * + *

반환 시점에 증가는 커밋된다. 그래야 뒤이은 외부 호출이 실패하더라도 예산이 실제로 + * 줄어들고, 같은 행이 다음 회차에 다시 잡혀 무한 재시도가 되지 않는다. + * + * @return 잠근 행들의 증가 전 상태. 만료된 단계가 무엇이었는지가 이 값에만 남는다 + */ + @Transactional + public List claimStale() { + List stale = repository.lockStale( + properties.pendingExpiry(), properties.processingExpiry(), + properties.maxRetry(), properties.batchSize()); + repository.incrementRetryCount(stale.stream().map(ContextAiStateRow::contextId).toList()); + return stale; + } + + /** + * 지금의 상태 행. 삭제된 Context를 만났을 때 그 상태가 {@code CANCELLED}인지 확인하는 데만 쓴다 + * (명세 5.1). 후보 조회는 {@code CANCELLED}를 잡지 않으므로 {@link #claimStale()}이 준 값으로는 + * 판정할 수 없다 — 삭제는 그 뒤에 일어났을 수 있다. + */ + @Transactional(readOnly = true) + public Optional findCurrentState(long contextId) { + return repository.findOne(contextId); + } +} diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/service/AiRescanProperties.java b/src/main/java/com/pinlog/pinlogback/domain/ai/service/AiRescanProperties.java new file mode 100644 index 00000000..8eb26897 --- /dev/null +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/service/AiRescanProperties.java @@ -0,0 +1,42 @@ +package com.pinlog.pinlogback.domain.ai.service; + +import java.time.Duration; + +import org.springframework.boot.context.properties.ConfigurationProperties; + +/** + * 재스캔·Finalizer 파라미터(AI 파트 소유 명세 {@code docs/ai/spec/ai-rescan-scheduler.md} 2장). + * + *

모두 설정값이다. 상수로 박으면 만료 기준을 바꿀 때 재배포가 필요하고, 무엇보다 + * "만료됐다"를 테스트에서 만들 수 없다 — 5분을 기다리는 테스트는 쓸 수 없다. 값의 정본은 명세이며 + * 여기서 임의로 바꾸지 않는다. + * + *

{@code pinlog.ai} 아래에 있지만 {@link com.pinlog.pinlogback.domain.ai.AiProperties}에 합치지 + * 않는다. 저쪽은 FastAPI 연결 계약(주소·시크릿·타임아웃)이고 이쪽은 복구 정책이라 + * 바뀌는 이유가 다르다. 접두어가 겹쳐도 충돌하지 않는다 — {@code @ConfigurationProperties}는 모르는 + * 키를 무시하므로 {@code AiProperties}가 {@code rescan}을 보고 실패하지 않는다. + * + *

{@code scheduler}가 아니라 여기 있는 이유는 값을 읽는 곳이 여기이기 때문이다 + * (package-structure.md — 설정은 소비자와 같은 패키지에 둔다). {@link AiFailedFinalizer}와 + * {@link AiRescanCandidateService} 둘이 소비자이고, 스케줄러는 주기 하나를 애노테이션 문자열로만 + * 읽는다. 반대로 두면 {@code service}가 {@code scheduler}를 참조해 패키지 의존이 순환한다. + * + * @param interval 실행 주기. 이 값을 읽는 것은 Java 코드가 아니라 + * {@code AiRescanScheduler}의 {@code @Scheduled(fixedDelayString)}이다. 여기 둔 이유는 그 + * 애노테이션이 문자열 placeholder만 받아 어떤 값인지 타입으로 드러나지 않기 때문이다 + * @param pendingExpiry {@code PENDING}이 이 시간을 넘기면 유실로 본다 + * @param processingExpiry {@code PROCESSING}이 이 시간을 넘기면 프로세스 종료로 유실된 것으로 본다. + * {@code pendingExpiry}보다 긴 이유는 실제로 처리 중일 가능성을 고려하기 때문이다(명세 2장) + * @param maxRetry 재시도 상한. 정본은 DB의 {@code CHECK (retry_count BETWEEN 0 AND 3)}이다 + * ({@code V100__ai_tables.sql}) — 이 값을 올리면 증가 UPDATE가 제약 위반으로 실패한다 + * @param batchSize 한 회차에 집는 후보 수 상한. 재스캔과 Finalizer가 각각 이 값을 쓴다 + */ +@ConfigurationProperties("pinlog.ai.rescan") +public record AiRescanProperties( + Duration interval, + Duration pendingExpiry, + Duration processingExpiry, + int maxRetry, + int batchSize +) { +} diff --git a/src/main/java/com/pinlog/pinlogback/domain/ai/service/ContextProcessRequestAssembler.java b/src/main/java/com/pinlog/pinlogback/domain/ai/service/ContextProcessRequestAssembler.java index 1c338547..1c925e6c 100644 --- a/src/main/java/com/pinlog/pinlogback/domain/ai/service/ContextProcessRequestAssembler.java +++ b/src/main/java/com/pinlog/pinlogback/domain/ai/service/ContextProcessRequestAssembler.java @@ -47,23 +47,39 @@ public ContextProcessRequestAssembler(ContextRepository contextRepository, Recor */ @Transactional(readOnly = true) public Optional assemble(ContextAiRequested event) { - Optional context = contextRepository.findById(event.contextId()); - if (context.isEmpty()) { - return Optional.empty(); - } + return contextRepository.findById(event.contextId()) + .flatMap(context -> assembleFrom(context, event.memberId(), event.recordId())); + } + + /** + * {@code context_id}만 들고 조립한다. 재스캔이 쓰는 진입점이다 + * (AI 파트 소유 명세 {@code docs/ai/spec/ai-rescan-scheduler.md} 5.1). + * + *

이벤트 경로와 갈라 둔 이유는 가진 정보가 다르기 때문이다. 재스캔은 상태 행에서 후보를 + * 집으므로 {@code context_id}뿐이고, {@code member_id}·{@code record_id}는 Context의 비정규화 + * 컬럼에서 읽는다. 반대로 이벤트 경로는 발행 시점의 값을 그대로 쓴다 — 그쪽을 이 메서드로 바꾸면 + * 같은 값을 두 번 읽게 되고, 무엇보다 이미 검증된 경로의 동작을 이 티켓이 건드리게 된다. + * + *

Context가 없으면(소프트 삭제·수정 교체) {@link Optional#empty()}다. 이것이 재스캔에서 + * 삭제 확인 그 자체다 — 별도 질의를 두지 않는다. + */ + @Transactional(readOnly = true) + public Optional assemble(long contextId) { + return contextRepository.findById(contextId) + .flatMap(context -> assembleFrom(context, context.getMemberId(), context.getRecordId())); + } + + private Optional assembleFrom(Context context, Long memberId, Long recordId) { // Place는 Record를 거쳐야 나온다. Record까지 지워졌으면 placeMeta 없이 보내지 않고 생략한다 — // Record가 없다는 것은 이 Context도 이미 함께 지워졌다는 뜻이다(연쇄 삭제). - Optional record = recordRepository.findById(event.recordId()); - if (record.isEmpty()) { - return Optional.empty(); - } - return Optional.of(new ContextProcessRequest( - event.contextId(), - event.memberId(), - event.recordId(), - context.get().getBody(), - placeMetaOf(record.get()) - )); + return recordRepository.findById(recordId) + .map(record -> new ContextProcessRequest( + context.getId(), + memberId, + recordId, + context.getBody(), + placeMetaOf(record) + )); } private ContextProcessRequest.PlaceMeta placeMetaOf(Record record) { diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index cabf0bfb..8d17934a 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -131,6 +131,26 @@ pinlog: # 스레드 점유만 늘린다(AI 파트 소유 명세 docs/ai/spec/ai-integration.md 3장). connect-timeout: 1s read-timeout: 5s + # 유실·정지된 AI 처리를 복구하는 재스캔과 FAILED Finalizer. 정본은 AI 파트가 소유한 + # docs/ai/spec/ai-rescan-scheduler.md 2장이며 여기서 임의로 바꾸지 않는다. 상수로 박지 않는 + # 이유는 튜닝 대상인 것도 있지만, 무엇보다 "만료됐다"를 테스트에서 만들 수 없기 때문이다. + rescan: + # ISO-8601이다. 다른 Duration 값들처럼 5m으로 적을 수 없다 — 이 값만은 Boot의 완화된 바인딩이 + # 아니라 @Scheduled(fixedDelayString)이 직접 파싱하고, 그쪽은 숫자(밀리초)나 ISO-8601만 받는다. + # 5m으로 두면 NumberFormatException으로 기동이 실패한다. + interval: PT5M + # 이 시간을 넘긴 PENDING은 요청이 유실된 것으로 본다(큐 포화·호출 실패·프로세스 종료). + pending-expiry: 5m + # PENDING보다 길다. 실제로 처리 중일 가능성을 고려한 값이며, 10분은 Embedding + LLM 판정의 + # 정상 소요를 크게 웃돈다. 이 시간을 넘긴 PROCESSING은 워커 프로세스 종료로 유실된 것이다. + processing-expiry: 10m + # 정본은 DB의 CHECK (retry_count BETWEEN 0 AND 3)이다(V100__ai_tables.sql). 이 값만 올리면 + # 증가 UPDATE가 제약 위반으로 실패한다. Backoff는 두지 않는다 — 주기 5분이 그 자체로 최소 + # 간격이고, 3회뿐이라 지수 backoff의 이득이 없다. + max-retry: 3 + # 재스캔과 Finalizer가 각각 한 회차에 집는 상한. 후보 건수가 계속 이 값에 붙어 있으면 + # FastAPI가 처리량을 못 따라가고 있다는 신호다(명세 8장). + batch-size: 100 # Feed 추천 정책값. 정본은 AI 파트가 소유한 docs/ai/spec/feed-scoring.md이며 여기서 임의로 # 바꾸지 않는다. 상수로 박지 않고 설정으로 두는 이유는 튜닝 대상이기 때문이다 — 재배포 없이 # 조정할 수 있어야 하고, 가중치를 바꿔도 순위가 안 바뀌는 회귀를 테스트가 잡을 수 있어야 한다. diff --git a/src/test/java/com/pinlog/pinlogback/ConfigurationContractTests.java b/src/test/java/com/pinlog/pinlogback/ConfigurationContractTests.java index 6e37b457..f5e9501c 100644 --- a/src/test/java/com/pinlog/pinlogback/ConfigurationContractTests.java +++ b/src/test/java/com/pinlog/pinlogback/ConfigurationContractTests.java @@ -85,6 +85,32 @@ void theEmbeddingProfileIsCommittedAsALiteralWithAnOptionalEnvironmentOverride() .isEqualTo("${PINLOG_AI_EMBEDDING_PROFILE:openai-text-embedding-3-small-1536-cosine-v1}"); } + /** + * 재스캔 파라미터는 AI 파트 소유 명세 {@code docs/ai/spec/ai-rescan-scheduler.md} 2장이 정본이다. + * 값을 파일 자체로 고정하는 이유는 어긋나도 아무 테스트가 깨지지 않기 때문이다 — 만료를 + * 검증하는 통합 테스트는 자기 임계값을 덮어 쓰므로 기본값이 무엇이든 통과한다. + * + *

{@code interval}만 ISO-8601인 것은 실수가 아니다. 이 값은 Boot의 완화된 바인딩이 아니라 + * {@code @Scheduled(fixedDelayString)}이 직접 파싱하고, 그쪽은 숫자(밀리초)나 ISO-8601만 받는다. + * {@code 5m}으로 적으면 {@code NumberFormatException}으로 기동이 실패한다. + */ + @Test + void theRescanParametersMatchTheOwningSpec() throws IOException { + Map defaults = load("application.yml"); + + assertThat(defaults.get("pinlog.ai.rescan.interval")) + .as("@Scheduled가 직접 파싱하므로 5m이 아니라 ISO-8601이어야 한다") + .isEqualTo("PT5M"); + assertThat(defaults.get("pinlog.ai.rescan.pending-expiry")).isEqualTo("5m"); + assertThat(defaults.get("pinlog.ai.rescan.processing-expiry")) + .as("실제로 처리 중일 가능성을 고려해 PENDING보다 길다") + .isEqualTo("10m"); + assertThat(String.valueOf(defaults.get("pinlog.ai.rescan.max-retry"))) + .as("정본은 DB의 CHECK (retry_count BETWEEN 0 AND 3)이다 — 올리면 증가 UPDATE가 실패한다") + .isEqualTo("3"); + assertThat(String.valueOf(defaults.get("pinlog.ai.rescan.batch-size"))).isEqualTo("100"); + } + @Test void prodProfileStillHidesApiDocumentation() throws IOException { Map prod = load("application-prod.yml"); diff --git a/src/test/java/com/pinlog/pinlogback/domain/ai/AiRescanSchedulerTests.java b/src/test/java/com/pinlog/pinlogback/domain/ai/AiRescanSchedulerTests.java new file mode 100644 index 00000000..c4ee777e --- /dev/null +++ b/src/test/java/com/pinlog/pinlogback/domain/ai/AiRescanSchedulerTests.java @@ -0,0 +1,522 @@ +package com.pinlog.pinlogback.domain.ai; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.math.BigDecimal; +import java.sql.Connection; +import java.sql.DriverManager; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.time.Duration; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; + +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.AfterEach; +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.context.ApplicationContext; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.scheduling.TaskScheduler; +import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler; +import org.springframework.test.context.DynamicPropertyRegistry; +import org.springframework.test.context.DynamicPropertySource; + +import com.pinlog.pinlogback.domain.ai.scheduler.AiRescanScheduler; +import com.pinlog.pinlogback.domain.ai.service.AiRescanProperties; +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.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; +import com.pinlog.pinlogback.integration.IntegrationContainerSupport; + +/** + * 유실·정지된 AI 처리를 복구하는 회차를 검증한다(S15P11A705-159, AI 파트 소유 명세 + * {@code docs/ai/spec/ai-rescan-scheduler.md} 3~6장). + * + *

주기를 기다리지 않고 {@link AiRescanScheduler#runOnce()}를 직접 부른다. 스케줄 실행을 + * 기다리면 테스트가 시계에 의존해 느려지고 불안정해진다. {@code @Scheduled} 등록 자체(fixedDelay인지, + * 전용 스케줄러를 쓰는지)는 별도 테스트가 본다 — 이 클래스가 보는 것은 회차가 무엇을 하는가다. + * + *

만료를 만드는 방법은 {@code updated_at}을 과거로 밀어 넣는 것이다. 5분을 기다릴 수는 없다. + * 만료 임계값을 기본값(5분·10분)이 아닌 값으로 덮어 두는 이유는 그 값이 설정에서 오는지를 + * 함께 확인하기 위해서다 — 코드에 상수로 박혀 있으면 아래 임계 테스트가 깨진다. + * + *

{@code @Transactional} 롤백 테스트가 아니다. {@code FOR UPDATE SKIP LOCKED}와 커밋 경계가 + * 검증 대상이라 실제로 커밋해야 한다({@code ContextAiEnqueueTests}와 같은 사정). + */ +@SpringBootTest +class AiRescanSchedulerTests extends IntegrationContainerSupport { + + /** Spring Context보다 먼저 떠야 {@code @DynamicPropertySource}가 포트를 알 수 있다. */ + private static final FastApiProcessStub STUB = new FastApiProcessStub(POSTGRES); + + private static final Duration PENDING_EXPIRY = Duration.ofMinutes(2); + private static final Duration PROCESSING_EXPIRY = Duration.ofMinutes(4); + /** 임계값 판정과 무관하게 "확실히 만료됐다"를 만들 때 쓰는 나이. */ + private static final Duration LONG_AGO = Duration.ofMinutes(30); + + /** + * Core에 존재하지 않는 {@code context_id}. {@code ai.context_ai_state}에는 {@code core.context}로 + * 향하는 FK가 없어(V100) 이런 행을 만들 수 있고, 음수를 쓰면 IDENTITY가 만드는 실제 id와 절대 + * 겹치지 않는다. 상태 전이만 보는 테스트는 Record·Place 조립이 필요 없다. + */ + private static final AtomicLong SYNTHETIC_IDS = new AtomicLong(-1_000L); + + private static final String SEED_SQL = """ + INSERT INTO ai.context_ai_state (context_id, embedding_status, keyword_status, retry_count, updated_at) + VALUES (?, ?, ?, ?, now() - make_interval(secs => ?)) + ON CONFLICT (context_id) DO UPDATE + SET embedding_status = EXCLUDED.embedding_status, + keyword_status = EXCLUDED.keyword_status, + retry_count = EXCLUDED.retry_count, + updated_at = EXCLUDED.updated_at + """; + + @DynamicPropertySource + static void aiServerPointsAtTheStub(DynamicPropertyRegistry registry) { + registry.add("pinlog.ai.base-url", STUB::baseUrl); + registry.add("pinlog.ai.rescan.pending-expiry", PENDING_EXPIRY::toString); + registry.add("pinlog.ai.rescan.processing-expiry", PROCESSING_EXPIRY::toString); + } + + @Autowired + private AiRescanScheduler scheduler; + + @Autowired + private AiRescanProperties properties; + + @Autowired + private RecordService recordService; + + @Autowired + private ContextRepository contextRepository; + + @Autowired + private MemberRepository memberRepository; + + @Autowired + private JdbcTemplate jdbcTemplate; + + @Autowired + private ApplicationContext applicationContext; + + /** + * 모든 상태 행의 나이를 0으로 돌린 뒤 시작한다. 컨테이너는 JVM이 공유하고 이 클래스의 각 + * 테스트는 실제로 커밋하므로, 앞선 테스트가 남긴 만료 행이 다음 회차에 후보로 섞인다. 그러면 + * "호출이 오지 않아야 한다" 같은 단언이 남의 행 때문에 깨진다. 뒤에도 같은 정리를 해서 다른 + * 테스트 클래스의 배경 회차에 만료 행을 물려주지 않는다. + */ + @BeforeEach + void freshenEveryStateRowAndResetTheStub() { + unstaleEveryStateRow(); + STUB.reset(FastApiProcessStub.Mode.ACCEPTED); + } + + @AfterEach + void freshenEveryStateRowAgain() { + unstaleEveryStateRow(); + } + + @AfterAll + static void stopStub() { + STUB.stop(); + } + + /** + * 재스캔의 본래 목적. 요청이 유실돼 {@code PENDING}으로 남은 Context를 다시 FastAPI에 보낸다. + * + *

본문을 함께 단언하는 이유: 명세 5.1은 이전 요청 본문을 보관해 재전송하지 말고 Core를 다시 + * 조회해 만들라고 정한다. 대역이 받은 {@code text}가 Core 본문과 같다는 것이 그 확인이다. + */ + @Test + void staleRowsAreSentToFastApiAgainAndTheRetryBudgetIsSpent() throws Exception { + long contextId = newContextIdFor("rescan-retry", "재스캔이 다시 집어야 한다"); + seedState(contextId, "PENDING", "PENDING", 0, LONG_AGO); + + scheduler.runOnce(); + + FastApiProcessStub.Received call = STUB.awaitCall(); + assertThat(call).as("만료된 PENDING은 다시 요청돼야 한다").isNotNull(); + assertThat(call.contextId()).isEqualTo(contextId); + assertThat(call.text()) + .as("요청 본문은 보관해 둔 것이 아니라 Core에서 다시 읽은 것이다") + .isEqualTo("재스캔이 다시 집어야 한다"); + assertThat(stateOf(contextId)).containsEntry("retry_count", 1); + } + + /** + * {@code retry_count} 증가는 외부 호출과 다른 트랜잭션이다(명세 3.1의 커밋 경계). 호출이 + * 실패했다고 예산이 되돌아가면 같은 행이 영원히 재시도되고, 그러면 상한 3회가 무의미해진다. + */ + @Test + void theRetryBudgetIsSpentEvenWhenTheCallFails() throws Exception { + long contextId = newContextIdFor("rescan-failing-call", "호출은 실패해도 예산은 줄어든다"); + seedState(contextId, "PENDING", "PENDING", 0, LONG_AGO); + STUB.reset(FastApiProcessStub.Mode.SERVER_ERROR); + + scheduler.runOnce(); + + assertThat(STUB.awaitCall()).as("호출은 나갔고 5xx를 받았다").isNotNull(); + assertThat(stateOf(contextId)) + .as("증가는 호출 전에 커밋됐으므로 실패가 되돌리지 못한다") + .containsEntry("retry_count", 1); + } + + /** + * 두 단계는 독립 전이한다(명세 7장). 한쪽만 {@code PENDING}인 행이 정상적으로 존재하고, 그것도 + * 재스캔 대상이다. 두 단계를 한 덩어리로 취급하면 Keyword만 밀린 Context가 영원히 방치된다. + * + *

재요청이 나가도 {@code COMPLETED}인 단계는 그대로 남아야 한다 — 어느 단계부터 재개할지는 + * FastAPI가 판단하고(명세 7장) Spring은 지시하지 않는다. + */ + @Test + void oneStageCompletedAndTheOtherStaleIsStillACandidate() throws Exception { + long contextId = newContextIdFor("rescan-half-done", "임베딩만 끝난 상태"); + seedState(contextId, "COMPLETED", "PENDING", 0, LONG_AGO); + + scheduler.runOnce(); + + FastApiProcessStub.Received call = STUB.awaitCall(); + assertThat(call).as("Keyword만 밀린 행도 재스캔 대상이다").isNotNull(); + assertThat(call.contextId()).isEqualTo(contextId); + assertThat(stateOf(contextId)) + .as("끝난 단계를 되돌리지 않는다 — 되돌리면 Embedding을 다시 만들어야 한다") + .containsEntry("embedding_status", "COMPLETED") + .containsEntry("keyword_status", "PENDING") + .containsEntry("retry_count", 1); + } + + /** + * 종결 상태는 후보가 아니다(명세 4.1). {@code FAILED}가 자동으로 되살아나는 경로는 없고, + * {@code CANCELLED}는 사용자가 지운 Context이므로 되살리면 삭제한 기록으로 임베딩이 만들어진다. + */ + @Test + void completedFailedAndCancelledRowsAreNeverCandidates() throws Exception { + long completed = seedSynthetic("COMPLETED", "COMPLETED", 0); + long failed = seedSynthetic("FAILED", "FAILED", 0); + long cancelled = seedSynthetic("CANCELLED", "CANCELLED", 0); + + scheduler.runOnce(); + + assertThat(STUB.noCallWithin(300)).as("종결된 행으로는 호출이 나가지 않는다").isTrue(); + assertThat(retryCountOf(completed)).isZero(); + assertThat(retryCountOf(failed)).isZero(); + assertThat(retryCountOf(cancelled)).isZero(); + assertThat(stateOf(cancelled)) + .as("CANCELLED를 손대지 않는다") + .containsEntry("embedding_status", "CANCELLED"); + } + + /** + * 만료 임계값은 설정에서 오고 단계 상태별로 다르다(명세 2장). {@code PROCESSING}이 더 긴 + * 이유는 실제로 처리 중일 가능성을 고려하기 때문이다. + * + *

세 행의 나이를 두 임계값 사이에 배치해 한 회차에서 갈라지는 것을 본다. 임계값이 코드 + * 상수라면 이 배치가 성립하지 않으므로, 이 테스트가 곧 "설정으로 주입된다"의 확인이다. + */ + @Test + void expiryThresholdsComeFromConfigurationAndDifferByStage() { + long stalePending = seedSynthetic("COMPLETED", "PENDING", 0, Duration.ofMinutes(3)); + long freshProcessing = seedSynthetic("COMPLETED", "PROCESSING", 0, Duration.ofMinutes(3)); + long staleProcessing = seedSynthetic("COMPLETED", "PROCESSING", 0, Duration.ofMinutes(5)); + + scheduler.runOnce(); + + assertThat(properties.pendingExpiry()).isEqualTo(PENDING_EXPIRY); + assertThat(properties.processingExpiry()).isEqualTo(PROCESSING_EXPIRY); + assertThat(retryCountOf(stalePending)).as("3분 > PENDING 임계 2분").isEqualTo(1); + assertThat(retryCountOf(freshProcessing)).as("3분 < PROCESSING 임계 4분").isZero(); + assertThat(retryCountOf(staleProcessing)).as("5분 > PROCESSING 임계 4분").isEqualTo(1); + } + + /** + * 마지막 재시도는 그 회차에서 종결되지 않는다(명세 3.1·6.1). {@code retry_count = 2}인 만료 + * 행은 이 회차에서 3이 되고, 3회차 요청이 실제로 나가야 하며 상태는 아직 {@code FAILED}가 아니어야 + * 한다. + * + *

이 단언은 "Finalize가 먼저 돈다"의 프록시가 아니다. 실측에서 확인했다 — {@code runOnce}의 + * 두 줄을 맞바꿔도 이 테스트는 통과한다. 재시도 증가가 {@code updated_at}을 함께 갱신하므로 그 행이 + * 만료 상태에서 벗어나 Finalizer 후보 조건에 걸리지 않기 때문이다. 즉 이 창을 실제로 확보하는 + * 것은 순서가 아니라 Finalizer의 만료 조건 + 증가 시 {@code updated_at} 갱신이고, 그 둘은 + * 아래 {@code anExhaustedRowThatIsNotExpiredYetIsLeftAlone}이 고정한다. 순서 자체는 심층 방어이며 + * {@code AiRescanSchedulerOrderTest}가 호출 순서로 고정한다. + */ + @Test + void theLastRetryActuallyGoesOutAndIsNotFinalizedInTheSameRound() throws Exception { + long contextId = newContextIdFor("rescan-last-try", "마지막 재시도"); + seedState(contextId, "PENDING", "PENDING", 2, LONG_AGO); + + scheduler.runOnce(); + + FastApiProcessStub.Received call = STUB.awaitCall(); + assertThat(call).as("3회차 요청은 실제로 나가야 한다").isNotNull(); + assertThat(call.contextId()).isEqualTo(contextId); + assertThat(stateOf(contextId)) + .as("같은 회차에서 종결되면 마지막 재시도가 실행되기도 전에 사망 선고가 된다") + .containsEntry("retry_count", 3) + .containsEntry("embedding_status", "PENDING") + .containsEntry("keyword_status", "PENDING"); + } + + /** + * Finalizer 없이는 {@code retry_count >= 3}인 행이 PROCESSING으로 영원히 남는다 — 재스캔 후보 + * 조건에서 빠지기만 할 뿐이다. 상태 지표상 "처리 중"으로 오인되어 장애 관측이 불가능해진다(명세 6장). + * + *

Finalizer는 상태 정리만 하고 FastAPI를 부르지 않는다. + */ + @Test + void retryExhaustedAndExpiredRowsAreFinalizedToFailedWithoutCallingFastApi() throws Exception { + long contextId = seedSynthetic("PROCESSING", "PENDING", 3); + + scheduler.runOnce(); + + assertThat(STUB.noCallWithin(300)).as("Finalizer는 순수한 상태 정리 단계다").isTrue(); + assertThat(stateOf(contextId)) + .containsEntry("embedding_status", "FAILED") + .containsEntry("keyword_status", "FAILED") + .containsEntry("retry_count", 3); + } + + /** + * Finalizer에도 만료 조건이 붙는다(명세 6.1). 이것이 마지막 재시도에게 창을 주는 장치다 — + * 없으면 {@code retry_count}를 3으로 올린 직후의 행이 다음 회차에서 곧바로 종결되어, 방금 나간 + * 3회차 요청이 처리될 시간을 갖지 못한다. + * + *

{@code retry_count = 3}이지만 아직 만료되지 않은 행이 그 상태다. 재스캔 후보도 아니고 + * ({@code retry_count < 3}이 아니다) 종결 대상도 아니어서, 이 회차는 이 행을 건드리지 않아야 + * 한다. 위 {@code theLastRetryActuallyGoesOutAndIsNotFinalizedInTheSameRound}가 확인하지 못하는 + * 바로 그 부분을 여기서 본다. + */ + @Test + void anExhaustedRowThatIsNotExpiredYetIsLeftAlone() { + long recentlyTouched = seedSynthetic("PENDING", "PENDING", 3, Duration.ofMinutes(1)); + + scheduler.runOnce(); + + assertThat(stateOf(recentlyTouched)) + .as("만료 조건이 없으면 방금 나간 3회차 요청이 처리되기도 전에 사망 선고가 된다") + .containsEntry("embedding_status", "PENDING") + .containsEntry("keyword_status", "PENDING") + .containsEntry("retry_count", 3); + } + + /** + * 종결은 미완료 단계만 건드린다(명세 6.3·6.4). + * + *

+ */ + @Test + void theFinalizerKeepsCompletedStagesAndNeverOverwritesCancelled() { + long halfDone = seedSynthetic("COMPLETED", "PROCESSING", 3); + long halfCancelled = seedSynthetic("CANCELLED", "PROCESSING", 3); + long fullyCancelled = seedSynthetic("CANCELLED", "CANCELLED", 3); + + scheduler.runOnce(); + + assertThat(stateOf(halfDone)) + .containsEntry("embedding_status", "COMPLETED") + .containsEntry("keyword_status", "FAILED"); + assertThat(stateOf(halfCancelled)) + .as("삭제 표시가 실패로 뒤집히면 '사용자가 지웠다'와 'AI가 실패했다'를 구별할 수 없다") + .containsEntry("embedding_status", "CANCELLED") + .containsEntry("keyword_status", "FAILED"); + assertThat(stateOf(fullyCancelled)) + .as("두 단계 모두 종결 상태라 후보조차 되지 않는다") + .containsEntry("embedding_status", "CANCELLED") + .containsEntry("keyword_status", "CANCELLED"); + } + + /** + * 삭제된 Context와 수정으로 교체된 구버전은 호출 대상이 아니다(명세 5.1·5.2). + * + *

구 Context의 상태 행을 일부러 {@code PENDING}으로 되돌려 후보가 되게 만든다. 정상 경로라면 + * 삭제 트랜잭션이 {@code CANCELLED}로 바꿔 두므로 후보조차 되지 않지만, 그러면 호출 직전의 삭제 + * 확인이 실제로 동작하는지를 볼 수 없다 — 이 재현은 삭제 트랜잭션이 무효화를 빠뜨린 상태 + * (명세 5.1이 정합성 경고를 남기라고 하는 그 상태)이기도 하다. + * + *

후보로 잡혀 {@code retry_count}는 올라가고 호출만 생략되는 것이 기대 동작이다. 예산을 + * 소진시키므로 지워진 Context가 3회차 뒤에는 후보에서도 사라진다. + */ + @Test + void deletedAndSupersededContextsAreClaimedButNeverSentToFastApi() throws Exception { + long memberId = newMemberId(); + RecordCreateResponse created = recordService.create(memberId, createRequest("rescan-deleted", "교체 전 이유")); + STUB.awaitCall(); + long recordId = created.recordId(); + long oldContextId = onlyContextId(recordId); + recordService.replaceContext(memberId, recordId, oldContextId, "교체 후 이유"); + STUB.awaitCall(); + assertThat(contextRepository.findById(oldContextId)) + .as("구 Context는 소프트 삭제돼 조회되지 않아야 한다 — 이 테스트의 전제") + .isEmpty(); + seedState(oldContextId, "PENDING", "PENDING", 0, LONG_AGO); + STUB.reset(FastApiProcessStub.Mode.ACCEPTED); + + scheduler.runOnce(); + + assertThat(STUB.noCallWithin(500)) + .as("지워진 Context의 본문을 보내면 삭제된 기록으로 임베딩이 만들어진다") + .isTrue(); + assertThat(retryCountOf(oldContextId)) + .as("호출만 생략한다 — 후보 선택과 예산 소진은 그대로 일어난다") + .isEqualTo(1); + } + + /** + * {@code SKIP LOCKED}가 실제로 건너뛰는지 관측한다(명세 4.2). 리더 선출을 두지 않는 근거가 이 + * 동작이므로, 이것이 깨지면 다중 인스턴스에서 같은 Context가 두 번 처리된다. + * + *

다른 커넥션이 한 행을 {@code FOR UPDATE}로 붙잡은 채 회차를 돌린다. 회차를 별 스레드에서 + * 돌리는 이유: {@code SKIP LOCKED}가 빠지면 후보 조회가 잠금을 기다리며 멈추는데, 같은 스레드에서 + * 부르면 테스트가 실패하는 대신 영원히 매달린다. 타임아웃을 걸어 실패로 드러나게 한다. + */ + @Test + void skipLockedLeavesALockedRowToWhoeverHoldsIt() throws Exception { + long locked = seedSynthetic("PENDING", "PENDING", 0); + long free = seedSynthetic("PENDING", "PENDING", 0); + + try (Connection holder = DriverManager.getConnection( + POSTGRES.getJdbcUrl(), POSTGRES.getUsername(), POSTGRES.getPassword())) { + holder.setAutoCommit(false); + lockRow(holder, locked); + + CompletableFuture round = CompletableFuture.runAsync(scheduler::runOnce); + round.get(15, TimeUnit.SECONDS); + + holder.rollback(); + } + + assertThat(retryCountOf(locked)) + .as("잠긴 행은 기다리지 않고 건너뛴다 — 기다리면 배치 전체가 한 행에 묶인다") + .isZero(); + assertThat(retryCountOf(free)) + .as("건너뛴 것은 잠긴 행뿐이고 회차는 계속 진행한다") + .isEqualTo(1); + } + + /** + * 스케줄링은 전용 {@link ThreadPoolTaskScheduler}를 쓴다(명세 3장). Boot의 기본 스케줄러는 + * 단일 스레드이고 앞으로 붙는 모든 배치가 그것을 공유한다. + * + *

{@link TaskScheduler} Bean이 하나뿐인 것까지 단언한다. Spring은 스케줄러를 작업별로 고르지 + * 않으므로, 두 번째 Bean이 생기면 이름이 {@code taskScheduler}가 아닌 쪽은 조용히 무시되고 최악의 + * 경우 로컬 단일 스레드 실행자로 떨어진다. 그 변경이 생기면 여기서 먼저 걸린다. + */ + @Test + void schedulingRunsOnADedicatedThreadPoolTaskScheduler() { + Map schedulers = applicationContext.getBeansOfType(TaskScheduler.class); + + assertThat(schedulers) + .as("Spring은 스케줄러를 작업별로 고르지 않는다 — 유일해야 해석이 흔들리지 않는다") + .hasSize(1); + TaskScheduler dedicated = schedulers.values().iterator().next(); + assertThat(dedicated).isInstanceOf(ThreadPoolTaskScheduler.class); + assertThat(((ThreadPoolTaskScheduler)dedicated).getThreadNamePrefix()).isEqualTo("ai-rescan-"); + assertThat(((ThreadPoolTaskScheduler)dedicated).getPoolSize()).isEqualTo(2); + } + + /** + * {@code @Scheduled(fixedDelayString)}가 읽는 키와 {@link AiRescanProperties}가 읽는 키가 같은 + * {@code pinlog.ai.rescan.interval}임을 확인한다. 두 곳이 갈라지면 "설정을 바꿨는데 주기가 그대로"가 + * 된다. 값은 {@code IntegrationContainerSupport}가 테스트용으로 덮은 것이다. + * + *

이 컨텍스트가 떴다는 사실 자체가 애노테이션 쪽 파싱의 확인이기도 하다 — 형식이 맞지 않으면 + * {@code @Scheduled} 등록 단계에서 기동이 실패한다. + */ + @Test + void theIntervalPropertyBindsIntoTheRecordThatDocumentsIt() { + assertThat(properties.interval()).isEqualTo(Duration.ofHours(1)); + assertThat(properties.maxRetry()) + .as("정본은 DB의 CHECK (retry_count BETWEEN 0 AND 3)이다") + .isEqualTo(3); + assertThat(properties.batchSize()).isEqualTo(100); + } + + private static void lockRow(Connection holder, long contextId) throws SQLException { + try (PreparedStatement lock = holder.prepareStatement( + "SELECT context_id FROM ai.context_ai_state WHERE context_id = ? FOR UPDATE")) { + lock.setLong(1, contextId); + try (ResultSet rows = lock.executeQuery()) { + assertThat(rows.next()).as("잠글 행이 있어야 한다 — 이 테스트의 전제").isTrue(); + } + } + } + + private void unstaleEveryStateRow() { + jdbcTemplate.update("UPDATE ai.context_ai_state SET updated_at = now()"); + } + + private long seedSynthetic(String embeddingStatus, String keywordStatus, int retryCount) { + return seedSynthetic(embeddingStatus, keywordStatus, retryCount, LONG_AGO); + } + + private long seedSynthetic(String embeddingStatus, String keywordStatus, int retryCount, Duration age) { + long contextId = SYNTHETIC_IDS.decrementAndGet(); + seedState(contextId, embeddingStatus, keywordStatus, retryCount, age); + return contextId; + } + + private void seedState(long contextId, String embeddingStatus, String keywordStatus, + int retryCount, Duration age) { + jdbcTemplate.update(SEED_SQL, contextId, embeddingStatus, keywordStatus, retryCount, + (double)age.toSeconds()); + } + + /** Record를 만들고 접수 호출까지 소진한 뒤 그 Context id를 준다. */ + private long newContextIdFor(String kakaoPlaceId, String body) throws InterruptedException { + RecordCreateResponse created = recordService.create(newMemberId(), createRequest(kakaoPlaceId, body)); + assertThat(STUB.awaitCall()).as("생성 시 접수 호출이 먼저 일어난다 — 이 테스트의 전제").isNotNull(); + STUB.reset(FastApiProcessStub.Mode.ACCEPTED); + return onlyContextId(created.recordId()); + } + + 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 int retryCountOf(long contextId) { + return (int)stateOf(contextId).get("retry_count"); + } + + private Map stateOf(long contextId) { + return jdbcTemplate.queryForMap( + "SELECT embedding_status, keyword_status, retry_count FROM ai.context_ai_state WHERE context_id = ?", + contextId); + } + + 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); + } +} diff --git a/src/test/java/com/pinlog/pinlogback/domain/ai/scheduler/AiRescanSchedulerOrderTest.java b/src/test/java/com/pinlog/pinlogback/domain/ai/scheduler/AiRescanSchedulerOrderTest.java new file mode 100644 index 00000000..ec63d0c7 --- /dev/null +++ b/src/test/java/com/pinlog/pinlogback/domain/ai/scheduler/AiRescanSchedulerOrderTest.java @@ -0,0 +1,83 @@ +package com.pinlog.pinlogback.domain.ai.scheduler; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +import java.util.List; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.InOrder; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.scheduling.annotation.Scheduled; + +import com.pinlog.pinlogback.domain.ai.client.AiProcessClient; +import com.pinlog.pinlogback.domain.ai.service.AiFailedFinalizer; +import com.pinlog.pinlogback.domain.ai.service.AiRescanCandidateService; +import com.pinlog.pinlogback.domain.ai.service.ContextProcessRequestAssembler; + +/** + * 회차의 단계 순서스케줄 방식을 대역·리플렉션으로 고정한다(AI 파트 소유 명세 + * {@code docs/ai/spec/ai-rescan-scheduler.md} 3장). + * + *

{@code AiRescanSchedulerTests}가 같은 두 계약을 관측 가능한 결과로도 고정한다 — 그쪽이 + * 더 강한 검증이다. 여기를 따로 두는 이유는 두 가지다. 순서 쪽은 결과가 아니라 호출 순서 자체를 + * 남겨 두어야 "왜 Finalize가 먼저인가"를 잃지 않고, 스케줄 방식 쪽은 {@code fixedRate}로 바뀌었을 때 + * DB 결과로는 아무 차이가 나지 않아 결과 검증으로 잡을 수 없다. + */ +@ExtendWith(MockitoExtension.class) +class AiRescanSchedulerOrderTest { + + @Mock + private AiFailedFinalizer finalizer; + + @Mock + private AiRescanCandidateService candidates; + + @Mock + private ContextProcessRequestAssembler assembler; + + @Mock + private AiProcessClient client; + + /** + * Finalize가 후보 선택보다 먼저 돈다. 뒤에 두면 같은 회차에서 방금 {@code retry_count}를 3으로 올린 + * 행을 곧바로 {@code FAILED}로 종결해, 마지막 재시도가 실행되기도 전에 사망 선고를 내린다(명세 3.1). + */ + @Test + void finalizeRunsBeforeTheCandidateClaimInEveryRound() { + when(finalizer.finalizeExpired()).thenReturn(List.of()); + when(candidates.claimStale()).thenReturn(List.of()); + AiRescanScheduler scheduler = new AiRescanScheduler(finalizer, candidates, assembler, client); + + scheduler.runOnce(); + + InOrder order = inOrder(finalizer, candidates); + order.verify(finalizer).finalizeExpired(); + order.verify(candidates).claimStale(); + verifyNoInteractions(client); + } + + /** + * {@code fixedRate}가 아니라 {@code fixedDelay}다(명세 3장). 한 회차가 배치 크기만큼의 HTTP 호출을 + * 순차로 내보내므로 실행이 주기를 넘길 수 있고, {@code fixedRate}면 그때 회차가 겹쳐 돈다. + * + *

주기를 리터럴이 아니라 placeholder로 두는 것도 함께 고정한다. 값을 애노테이션에 박으면 + * {@code pinlog.ai.rescan.interval} 설정이 있어도 아무 효력이 없다. + */ + @Test + void theRoundIsScheduledWithFixedDelayAndReadsTheConfiguredInterval() throws Exception { + Scheduled scheduled = AiRescanScheduler.class.getMethod("runOnce").getAnnotation(Scheduled.class); + + assertThat(scheduled).as("@Scheduled가 없으면 이 회차는 아무도 돌리지 않는다").isNotNull(); + assertThat(scheduled.fixedDelayString()).isEqualTo("${pinlog.ai.rescan.interval}"); + assertThat(scheduled.fixedRateString()) + .as("fixedRate면 이전 회차가 길어졌을 때 다음 회차가 겹쳐 돈다") + .isEmpty(); + assertThat(scheduled.fixedRate()).isEqualTo(-1); + assertThat(scheduled.cron()).isEmpty(); + } +} diff --git a/src/test/java/com/pinlog/pinlogback/integration/IntegrationContainerSupport.java b/src/test/java/com/pinlog/pinlogback/integration/IntegrationContainerSupport.java index 44ba6700..db010c7a 100644 --- a/src/test/java/com/pinlog/pinlogback/integration/IntegrationContainerSupport.java +++ b/src/test/java/com/pinlog/pinlogback/integration/IntegrationContainerSupport.java @@ -39,9 +39,18 @@ // 하는데, 두 쪽 다 @DynamicPropertySource 면 상위 클래스 쪽이 나중에 등록돼 하위를 덮어버린다 // (실측: 그 클래스 테스트 5개가 전부 "호출이 오지 않음"으로 깨졌다). DynamicValuesPropertySource // 는 우선순위가 가장 높으므로, 기본값을 @TestPropertySource 로 한 단계 낮춰 두면 하위가 이긴다. +// +// 재스캔 주기도 늘린다. @EnableScheduling 은 전역이라 모든 @SpringBootTest 가 스케줄러를 함께 +// 띄우는데, 5분 주기로 두면 컨텍스트를 오래 공유하는 스위트에서 회차가 배경에서 돌아 다른 테스트가 +// 만든 상태 행을 건드린다. 재스캔 자체를 검증하는 테스트는 주기를 기다리지 않고 runOnce() 를 직접 +// 부르므로(그래야 결정적이다) 이 값이 크면 배경 실행만 사라지고 검증은 그대로다. +// +// 끄지 않고 늘리는 이유: @Scheduled 등록 자체가 검증 대상이다(fixedDelay 인지, 전용 스케줄러를 +// 쓰는지). 조건부로 끄면 그 계약을 볼 수 없다. @TestPropertySource(properties = { "pinlog.ai.base-url=http://127.0.0.1:1", - "pinlog.ai.internal-secret=test-internal-secret" + "pinlog.ai.internal-secret=test-internal-secret", + "pinlog.ai.rescan.interval=PT1H" }) public abstract class IntegrationContainerSupport {