diff --git a/docs/tasks/NOOK-162/README.md b/docs/tasks/NOOK-162/README.md new file mode 100644 index 00000000..e20ab7ac --- /dev/null +++ b/docs/tasks/NOOK-162/README.md @@ -0,0 +1,61 @@ +# NOOK-162 게시물·장소 파싱 구조화 단계 로그 표준화 + +## 목적 + +Prometheus metric과 별개로 게시물 저장부터 콘텐츠 수집, 장소 파싱과 후처리까지 원본 게시물 ID로 +검색 가능한 구조화 애플리케이션 로그를 제공한다. + +이번 리팩터링 기간에만 사용하는 임시 추적 로그이며 모든 새 메시지는 `[PostParcingTracker]`로 시작한다. + +## 범위 + +- 공통 lifecycle event와 로그 필드를 정의한다. +- 콘텐츠·장소 비동기 실행에 `source_post.id`와 `processing.flow` MDC를 설정한다. +- Instagram, OpenAI, Kakao, Naver, 미디어, Google 사진과 태그 처리 결과를 요약해 기록한다. +- Instagram 캐시 적중 여부와 실제 provider 호출, fallback 및 응답 시간을 기록한다. +- Google 장소 매칭, 사진 목록, 사진 URI 조회와 스토리지 저장을 분리하고 누락 사유를 기록한다. +- 사용자 저장 게시물 ID와 원본 게시물 ID의 매핑을 기록한다. +- 새 구조화 로그에는 API key, 인증 헤더, 본문 원문, 전체 외부 응답과 전체 미디어 URL을 추가하지 않는다. +- 기존 로그의 메시지와 레벨은 호환성을 위해 변경하지 않는다. + +## 제외 범위 + +- Prometheus metric 변경 +- 장소 검색 및 후보 선택 알고리즘 변경 +- 비동기 Job 및 DB 스키마 변경 + +## 주요 검색 필드 + +- `event.action` +- `event.outcome` +- `processing.flow` +- `processing.stage` +- `processing.attempt` +- `source_post.id` +- `saved_post.id` +- `provider.name` +- `failure.type` +- `failure.reason` + +## 사진 누락 판별 + +- `place_not_matched`: Google 검색 결과에서 원본 장소와 일치하는 장소가 없음 +- `no_photos`: 일치 장소는 있지만 제공된 사진이 없음 +- `photo_uri_missing`: 사진 메타데이터 응답에 다운로드 URI가 없음 +- `google.photo.media.failed`: 사진 URI 조회 요청 실패 +- `google.photo.store.failed`: 스토리지 저장 실패 +- `google.photo.pipeline.completed`: 선택·저장·실패한 사진 수의 최종 집계 + +## 로그 레벨 + +- 새 추적 로그는 기존 운영 로그와 Grafana 알림에 영향을 주지 않도록 모두 `DEBUG`로 기록한다. +- `event.outcome`과 `failure.type`으로 성공, 빈 결과, 재시도와 실패를 구분한다. +- 기존 로그의 레벨과 메시지는 그대로 유지한다. + +고카디널리티 값은 Loki label로 추가하지 않고 JSON 검색 필드로만 기록한다. + +## 검증 + +- `./gradlew detekt` +- `./gradlew test` +- `./gradlew check` diff --git a/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/PlaceParsingLog.kt b/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/PlaceParsingLog.kt new file mode 100644 index 00000000..86701947 --- /dev/null +++ b/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/PlaceParsingLog.kt @@ -0,0 +1,24 @@ +package org.every.nook.api.application.place + +import org.every.nook.api.application.processing.ProcessingLogEvent + +internal fun ClaimedPlaceParsingJob.event( + action: String, + stage: String, + outcome: String, + durationMs: Long? = null, + fields: Map = emptyMap(), +) = ProcessingLogEvent(action, PLACE_FLOW, stage, outcome, postId, attempt, durationMs, fields) + +internal fun failureFields(exception: Throwable, reason: String): Map = mapOf( + "failure.type" to exception::class.simpleName, + "failure.reason" to reason, +) + +internal fun placeFailureReason(exception: Throwable): String = exception.message.orEmpty() + .ifBlank { DEFAULT_FAILURE_REASON } + .take(MAX_FAILURE_REASON_LENGTH) + +private const val PLACE_FLOW = "place" +private const val DEFAULT_FAILURE_REASON = "Place parsing failed" +private const val MAX_FAILURE_REASON_LENGTH = 500 diff --git a/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/ProcessPlaceParsingJobUseCase.kt b/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/ProcessPlaceParsingJobUseCase.kt index 0ae6f63a..7946b1b7 100644 --- a/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/ProcessPlaceParsingJobUseCase.kt +++ b/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/ProcessPlaceParsingJobUseCase.kt @@ -3,7 +3,11 @@ package org.every.nook.api.application.place import mu.KotlinLogging import org.every.nook.api.application.processing.NoOpProcessingMetrics import org.every.nook.api.application.processing.ProcessingMetrics +import org.every.nook.api.application.processing.error +import org.every.nook.api.application.processing.info import org.every.nook.api.application.processing.measure +import org.every.nook.api.application.processing.warn +import org.slf4j.LoggerFactory import java.time.Clock import java.time.Duration import java.time.Instant @@ -22,46 +26,59 @@ class ProcessPlaceParsingJobUseCase( operator fun invoke(postId: Long): Result { val job = jobPort.claim(postId, processingTimeout) ?: return Result.Skipped val startedAt = clock.instant() + eventLogger.info(job.event("place.job.claimed", JOB_STAGE, SUCCESS_OUTCOME)) logger.info { "Place parsing started: postId=${job.postId}, attempt=${job.attempt}" } - return runCatching { - val expectedPlaceCount = expectedPlaceCount(job.body) - val textClues = (job.textClues ?: extractClues(job)).filter { clue -> - clue.isGroundedIn(job).also { grounded -> - if (!grounded) { - logger.warn { - "Ungrounded text place clue skipped: postId=${job.postId}, attempt=${job.attempt}, " + - "placeName=${clue.name}, region=${clue.region}, queries=${clue.queries}" - } + return runCatching { process(job, startedAt) }.getOrElse { exception -> + handleFailure(job, exception, startedAt) + } + } + + private fun process(job: ClaimedPlaceParsingJob, startedAt: Instant): Result { + val expectedPlaceCount = expectedPlaceCount(job.body) + val textClues = (job.textClues ?: extractClues(job)).filter { clue -> + clue.isGroundedIn(job).also { grounded -> + if (!grounded) { + logger.warn { + "Ungrounded text place clue skipped: postId=${job.postId}, attempt=${job.attempt}, " + + "placeName=${clue.name}, region=${clue.region}, queries=${clue.queries}" } } } - val textResolution = resolveClues(job, textClues) - val imageResolution = resolveImageClues(job, textResolution.places.size, expectedPlaceCount) - val places = (textResolution.places + imageResolution?.places.orEmpty()) - .distinctBy { it.provider to it.externalPlaceId } - if (places.isEmpty()) { - val failure = imageResolution?.failure ?: textResolution.failure - terminalFailure( - failure?.message ?: if (imageResolution == null) { - NO_PLACE_RESOLVED_REASON - } else { - NO_PLACE_RESOLVED_AFTER_IMAGE_REASON - }, - ) - } - measure(job, COMPLETE_STAGE) { - jobPort.complete(job.postId, places) - } - val duration = Duration.between(startedAt, clock.instant()).toMillis() - logger.info { - "Place parsing completed: postId=${job.postId}, attempt=${job.attempt}, " + - "placeCount=${places.size}, durationMs=$duration" - } - Result.Completed - }.getOrElse { exception -> - handleFailure(job, exception, startedAt) } + val textResolution = resolveClues(job, textClues) + logOcrDecision(eventLogger, job, textResolution.places.size, expectedPlaceCount) + val imageResolution = resolveImageClues(job, textResolution.places.size, expectedPlaceCount) + val places = (textResolution.places + imageResolution?.places.orEmpty()) + .distinctBy { it.provider to it.externalPlaceId } + if (places.isEmpty()) { + val failure = imageResolution?.failure ?: textResolution.failure + terminalFailure( + failure?.message ?: if (imageResolution == null) { + NO_PLACE_RESOLVED_REASON + } else { + NO_PLACE_RESOLVED_AFTER_IMAGE_REASON + }, + ) + } + measure(job, COMPLETE_STAGE) { + jobPort.complete(job.postId, places) + } + val duration = Duration.between(startedAt, clock.instant()).toMillis() + logger.info { + "Place parsing completed: postId=${job.postId}, attempt=${job.attempt}, " + + "placeCount=${places.size}, durationMs=$duration" + } + eventLogger.info( + job.event( + "place.job.completed", + JOB_STAGE, + SUCCESS_OUTCOME, + duration, + mapOf("place.resolved_count" to places.size), + ), + ) + return Result.Completed } private fun resolveImageClues( @@ -187,6 +204,20 @@ class ProcessPlaceParsingJobUseCase( if (!clue.isSupportedBy(resolved)) { failResolution("Selected place is not grounded in image evidence: ${clue.name}") } + eventLogger.info( + job.event( + "place.candidate.selected", + SELECT_STAGE, + SUCCESS_OUTCOME, + fields = mapOf( + "provider.name" to resolved.provider, + "place.external_id" to resolved.externalPlaceId, + "place.selection_method" to if (matches.size == 1) "strict_match" else "openai", + "place.candidate_count" to candidates.size, + "place.strict_match_count" to matches.size, + ), + ), + ) logger.info { "Place resolved: provider=${resolved.provider}, externalPlaceId=${resolved.externalPlaceId}, " + "name=${resolved.name}, address=${resolved.address}" @@ -234,10 +265,14 @@ class ProcessPlaceParsingJobUseCase( private fun failResolution(message: String): Nothing = throw PlaceResolutionException(message) private fun handleFailure(job: ClaimedPlaceParsingJob, exception: Throwable, startedAt: Instant): Result { - val reason = failureReason(exception) + val reason = placeFailureReason(exception) val duration = Duration.between(startedAt, clock.instant()).toMillis() if (exception is TerminalPlaceParsingException) { jobPort.fail(job.postId, reason) + eventLogger.warn( + job.event("place.job.failed", JOB_STAGE, FAILURE_OUTCOME, duration, failureFields(exception, reason)), + exception, + ) logger.warn { "Place parsing failed without retry: postId=${job.postId}, attempt=${job.attempt}, " + "durationMs=$duration, reason=$reason" @@ -249,6 +284,16 @@ class ProcessPlaceParsingJobUseCase( if (backoff != null) { val nextAttemptAt = clock.instant().plus(backoff) jobPort.retry(job.postId, nextAttemptAt, reason) + eventLogger.warn( + job.event( + "place.job.retry_scheduled", + JOB_STAGE, + FAILURE_OUTCOME, + duration, + failureFields(exception, reason) + ("retry.next_attempt_at" to nextAttemptAt), + ), + exception, + ) logger.warn(exception) { "Place parsing retry scheduled: postId=${job.postId}, attempt=${job.attempt}, " + "nextAttemptAt=$nextAttemptAt, durationMs=$duration, reason=$reason" @@ -257,6 +302,10 @@ class ProcessPlaceParsingJobUseCase( } jobPort.fail(job.postId, reason) + eventLogger.error( + job.event("place.job.failed", JOB_STAGE, FAILURE_OUTCOME, duration, failureFields(exception, reason)), + exception, + ) logger.error(exception) { "Place parsing failed permanently: postId=${job.postId}, attempt=${job.attempt}, " + "durationMs=$duration, reason=$reason" @@ -264,10 +313,6 @@ class ProcessPlaceParsingJobUseCase( return Result.Failed } - private fun failureReason(exception: Throwable): String = exception.message.orEmpty() - .ifBlank { DEFAULT_FAILURE_REASON } - .take(MAX_FAILURE_REASON_LENGTH) - private fun terminalFailure(message: String): Nothing = throw TerminalPlaceParsingException(message) sealed interface Result { @@ -282,13 +327,12 @@ class ProcessPlaceParsingJobUseCase( private companion object { val logger = KotlinLogging.logger {} + val eventLogger = LoggerFactory.getLogger(ProcessPlaceParsingJobUseCase::class.java) const val MAX_PLACE_COUNT = 20 const val MAX_QUERY_COUNT = 4 const val MAX_IMAGE_COUNT = 20 const val CANDIDATE_LOG_LIMIT = 5 - const val MAX_FAILURE_REASON_LENGTH = 500 - const val DEFAULT_FAILURE_REASON = "Place parsing failed" const val NO_PLACE_RESOLVED_REASON = "No place could be resolved from text" const val NO_PLACE_RESOLVED_AFTER_IMAGE_REASON = "No place could be resolved after image analysis" const val PLACE_FLOW = "place" @@ -298,6 +342,10 @@ class ProcessPlaceParsingJobUseCase( const val SEARCH_STAGE = "search" const val SELECT_STAGE = "select" const val COMPLETE_STAGE = "complete" + const val JOB_STAGE = "job" + const val OCR_STAGE = "ocr" + const val SUCCESS_OUTCOME = "success" + const val FAILURE_OUTCOME = "failure" } private class PlaceResolutionException(message: String) : IllegalStateException(message) @@ -307,6 +355,28 @@ class ProcessPlaceParsingJobUseCase( private data class ClueResolution(val places: List, val failure: PlaceResolutionException?) } +private fun logOcrDecision( + logger: org.slf4j.Logger, + job: ClaimedPlaceParsingJob, + textPlaceCount: Int, + expectedPlaceCount: Int?, +) { + logger.info( + job.event( + "place.ocr.decision", + "ocr", + "success", + fields = mapOf( + "ocr.required" to requiresImageAnalysis(textPlaceCount, expectedPlaceCount), + "ocr.reason" to ocrReason(textPlaceCount, expectedPlaceCount, job.imageUrls.isEmpty()), + "place.text_resolved_count" to textPlaceCount, + "place.expected_count" to expectedPlaceCount, + "content.image_count" to job.imageUrls.size, + ), + ), + ) +} + private fun PlaceClue.isGroundedIn(job: ClaimedPlaceParsingJob): Boolean { val sources = buildList { job.body?.let(::add) @@ -323,6 +393,13 @@ private fun PlaceClue.isGroundedIn(job: ClaimedPlaceParsingJob): Boolean { private fun requiresImageAnalysis(textPlaceCount: Int, expectedPlaceCount: Int?): Boolean = textPlaceCount == 0 || expectedPlaceCount?.let { textPlaceCount < it } == true +private fun ocrReason(textPlaceCount: Int, expectedPlaceCount: Int?, imagesEmpty: Boolean): String = when { + imagesEmpty -> "no_images" + textPlaceCount == 0 -> "no_text_place_resolved" + expectedPlaceCount != null && textPlaceCount < expectedPlaceCount -> "expected_place_count_shortfall" + else -> "text_places_sufficient" +} + private fun expectedPlaceCount(body: String?): Int? = body?.let { content -> EXPECTED_PLACE_COUNT_PATTERN.findAll(content) .mapNotNull { match -> match.groupValues[1].toIntOrNull() } diff --git a/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/StorePlaceTagsUseCase.kt b/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/StorePlaceTagsUseCase.kt index a49927bf..043b65b6 100644 --- a/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/StorePlaceTagsUseCase.kt +++ b/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/StorePlaceTagsUseCase.kt @@ -1,12 +1,27 @@ package org.every.nook.api.application.place +import org.every.nook.api.application.processing.ProcessingLogEvent +import org.every.nook.api.application.processing.info +import org.slf4j.LoggerFactory + class StorePlaceTagsUseCase( private val sourcePort: PlaceTagSourcePort, private val extractor: PlaceTagExtractor, private val updatePort: PlaceTagUpdatePort, ) { operator fun invoke(event: PlaceTagsRequestedEvent) { - val source = sourcePort.find(event.postId) ?: return + val source = sourcePort.find(event.postId) ?: run { + logger.info( + logEvent( + event, + "place.tags.skipped", + "load_source", + "skipped", + mapOf("skip.reason" to "source_not_found"), + ), + ) + return + } val tags = extractor.extract( PlaceTagExtractor.Request( place = event.place, @@ -16,9 +31,35 @@ class StorePlaceTagsUseCase( ), ).distinctBy(InferredPlaceTag::tag).take(MAX_TAG_COUNT) updatePort.replace(event.postId, event.placeId, tags) + logger.info( + logEvent( + event, + "place.tags.completed", + "complete", + "success", + mapOf("place.tag_count" to tags.size), + ), + ) } + private fun logEvent( + event: PlaceTagsRequestedEvent, + action: String, + stage: String, + outcome: String, + fields: Map, + ) = ProcessingLogEvent( + action = action, + flow = TAG_FLOW, + stage = stage, + outcome = outcome, + sourcePostId = event.postId, + fields = fields + mapOf("place.id" to event.placeId), + ) + private companion object { + val logger = LoggerFactory.getLogger(StorePlaceTagsUseCase::class.java) + const val TAG_FLOW = "place-tags" const val MAX_TAG_COUNT = 4 } } diff --git a/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/StorePlaceThumbnailUseCase.kt b/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/StorePlaceThumbnailUseCase.kt index b02825d0..220ce296 100644 --- a/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/StorePlaceThumbnailUseCase.kt +++ b/nook-api-application/src/main/kotlin/org/every/nook/api/application/place/StorePlaceThumbnailUseCase.kt @@ -1,9 +1,13 @@ package org.every.nook.api.application.place import org.every.nook.api.application.processing.NoOpProcessingMetrics +import org.every.nook.api.application.processing.ProcessingLogEvent import org.every.nook.api.application.processing.ProcessingMetrics +import org.every.nook.api.application.processing.error +import org.every.nook.api.application.processing.info import org.every.nook.api.application.processing.measure import org.every.nook.api.domain.place.PlaceThumbnailParsingStatus +import org.slf4j.LoggerFactory import java.time.Clock class StorePlaceThumbnailUseCase( @@ -13,6 +17,8 @@ class StorePlaceThumbnailUseCase( private val clock: Clock = Clock.systemUTC(), ) { operator fun invoke(postId: Long, place: PlaceCandidate) { + val startedAt = clock.millis() + logger.info(event(postId, place, "place.thumbnail.started", FETCH_STAGE, "started")) runCatching { updatePort.update(place.provider, place.externalPlaceId, PlaceThumbnailParsingStatus.PROCESSING) val supplement = metrics.measure(THUMBNAIL_FLOW, FETCH_STAGE, postId, null, clock) { @@ -26,17 +32,57 @@ class StorePlaceThumbnailUseCase( supplement, ) } + logger.info( + event( + postId, + place, + "place.thumbnail.completed", + COMPLETE_STAGE, + "success", + startedAt, + mapOf( + "place.photo_count" to supplement?.photoUrls?.size, + "place.opening_hours_found" to (supplement?.openingHours != null), + ), + ), + ) }.getOrElse { exception -> runCatching { updatePort.update(place.provider, place.externalPlaceId, PlaceThumbnailParsingStatus.FAILED) }.onFailure { statusException -> exception.addSuppressed(statusException) } + logger.error( + event(postId, place, "place.thumbnail.failed", FETCH_STAGE, "failure", startedAt), + exception, + ) throw exception } } + private fun event( + postId: Long, + place: PlaceCandidate, + action: String, + stage: String, + outcome: String, + startedAt: Long? = null, + fields: Map = emptyMap(), + ) = ProcessingLogEvent( + action = action, + flow = THUMBNAIL_FLOW, + stage = stage, + outcome = outcome, + sourcePostId = postId, + durationMs = startedAt?.let { clock.millis() - it }, + fields = fields + mapOf( + "provider.name" to place.provider, + "place.external_id" to place.externalPlaceId, + ), + ) + private companion object { + val logger = LoggerFactory.getLogger(StorePlaceThumbnailUseCase::class.java) const val THUMBNAIL_FLOW = "place-thumbnail" const val FETCH_STAGE = "fetch" const val COMPLETE_STAGE = "complete" diff --git a/nook-api-application/src/main/kotlin/org/every/nook/api/application/post/ProcessPostContentParsingJobUseCase.kt b/nook-api-application/src/main/kotlin/org/every/nook/api/application/post/ProcessPostContentParsingJobUseCase.kt index 9a22649a..63b5873b 100644 --- a/nook-api-application/src/main/kotlin/org/every/nook/api/application/post/ProcessPostContentParsingJobUseCase.kt +++ b/nook-api-application/src/main/kotlin/org/every/nook/api/application/post/ProcessPostContentParsingJobUseCase.kt @@ -6,9 +6,14 @@ import org.every.nook.api.application.content.ExtractedPostContent import org.every.nook.api.application.content.PostContentNotFoundException import org.every.nook.api.application.content.UnsupportedPostUrlException import org.every.nook.api.application.processing.NoOpProcessingMetrics +import org.every.nook.api.application.processing.ProcessingLogEvent import org.every.nook.api.application.processing.ProcessingMetrics +import org.every.nook.api.application.processing.error +import org.every.nook.api.application.processing.info import org.every.nook.api.application.processing.measure +import org.every.nook.api.application.processing.warn import org.every.nook.api.domain.post.Post +import org.slf4j.LoggerFactory import java.time.Clock import java.time.Duration import java.time.Instant @@ -25,6 +30,7 @@ class ProcessPostContentParsingJobUseCase( operator fun invoke(postId: Long): Result { val job = jobPort.claim(postId, processingTimeout) ?: return Result.Skipped val startedAt = clock.instant() + eventLogger.info(job.event("content.job.claimed", JOB_STAGE, SUCCESS_OUTCOME)) logger.info { "Post content parsing started: postId=${job.postId}, attempt=${job.attempt}" } return runCatching { @@ -55,6 +61,15 @@ class ProcessPostContentParsingJobUseCase( "Post content parsing completed: postId=${job.postId}, attempt=${job.attempt}, " + "mediaCount=${completedPost.media.size}, durationMs=$duration" } + eventLogger.info( + job.event( + action = "content.job.completed", + stage = JOB_STAGE, + outcome = SUCCESS_OUTCOME, + durationMs = duration, + fields = mapOf("content.media_count" to completedPost.media.size), + ), + ) Result.Completed }.getOrElse { exception -> handleFailure(job, exception, startedAt) @@ -77,6 +92,10 @@ class ProcessPostContentParsingJobUseCase( val duration = Duration.between(startedAt, clock.instant()).toMillis() if (exception is PostContentNotFoundException || exception is UnsupportedPostUrlException) { jobPort.fail(job.postId, reason) + eventLogger.warn( + job.event("content.job.failed", JOB_STAGE, FAILURE_OUTCOME, duration, failureFields(exception, reason)), + exception, + ) logger.warn { "Post content parsing failed without retry: postId=${job.postId}, attempt=${job.attempt}, " + "durationMs=$duration, reason=$reason" @@ -87,6 +106,16 @@ class ProcessPostContentParsingJobUseCase( if (backoff != null) { val nextAttemptAt = clock.instant().plus(backoff) jobPort.retry(job.postId, nextAttemptAt, reason) + eventLogger.warn( + job.event( + "content.job.retry_scheduled", + JOB_STAGE, + FAILURE_OUTCOME, + duration, + failureFields(exception, reason) + ("retry.next_attempt_at" to nextAttemptAt), + ), + exception, + ) logger.warn(exception) { "Post content parsing retry scheduled: postId=${job.postId}, attempt=${job.attempt}, " + "nextAttemptAt=$nextAttemptAt, durationMs=$duration, reason=$reason" @@ -95,6 +124,10 @@ class ProcessPostContentParsingJobUseCase( } jobPort.fail(job.postId, reason) + eventLogger.error( + job.event("content.job.failed", JOB_STAGE, FAILURE_OUTCOME, duration, failureFields(exception, reason)), + exception, + ) logger.error(exception) { "Post content parsing failed permanently: postId=${job.postId}, attempt=${job.attempt}, " + "durationMs=$duration, reason=$reason" @@ -114,6 +147,19 @@ class ProcessPostContentParsingJobUseCase( .distinct() .toList() + private fun ClaimedPostContentParsingJob.event( + action: String, + stage: String, + outcome: String, + durationMs: Long? = null, + fields: Map = emptyMap(), + ) = ProcessingLogEvent(action, CONTENT_FLOW, stage, outcome, postId, attempt, durationMs, fields) + + private fun failureFields(exception: Throwable, reason: String): Map = mapOf( + "failure.type" to exception::class.simpleName, + "failure.reason" to reason, + ) + sealed interface Result { data object Completed : Result @@ -126,11 +172,15 @@ class ProcessPostContentParsingJobUseCase( private companion object { val logger = KotlinLogging.logger {} + val eventLogger = LoggerFactory.getLogger(ProcessPostContentParsingJobUseCase::class.java) const val MAX_FAILURE_REASON_LENGTH = 500 const val DEFAULT_FAILURE_REASON = "Post content parsing failed" const val CONTENT_FLOW = "post-content" const val EXTRACT_STAGE = "extract" const val INFERENCE_STAGE = "inference" const val COMPLETE_STAGE = "complete" + const val JOB_STAGE = "job" + const val SUCCESS_OUTCOME = "success" + const val FAILURE_OUTCOME = "failure" } } diff --git a/nook-api-application/src/main/kotlin/org/every/nook/api/application/processing/ProcessingLog.kt b/nook-api-application/src/main/kotlin/org/every/nook/api/application/processing/ProcessingLog.kt new file mode 100644 index 00000000..b00689e3 --- /dev/null +++ b/nook-api-application/src/main/kotlin/org/every/nook/api/application/processing/ProcessingLog.kt @@ -0,0 +1,68 @@ +package org.every.nook.api.application.processing + +import org.slf4j.Logger +import org.slf4j.MDC + +object ProcessingLogFields { + const val EVENT_ACTION = "event.action" + const val EVENT_OUTCOME = "event.outcome" + const val EVENT_DURATION_MS = "event.duration_ms" + const val FLOW = "processing.flow" + const val STAGE = "processing.stage" + const val ATTEMPT = "processing.attempt" + const val SOURCE_POST_ID = "source_post.id" + const val SAVED_POST_ID = "saved_post.id" + const val PROVIDER = "provider.name" + const val FAILURE_TYPE = "failure.type" + const val FAILURE_REASON = "failure.reason" +} + +data class ProcessingLogEvent( + val action: String, + val flow: String, + val stage: String, + val outcome: String? = null, + val sourcePostId: Long? = null, + val attempt: Int? = null, + val durationMs: Long? = null, + val fields: Map = emptyMap(), +) + +fun Logger.info(event: ProcessingLogEvent) = trackerBuilder(event).log(trackerMessage(event.action)) + +fun Logger.debug(event: ProcessingLogEvent) = trackerBuilder(event).log(trackerMessage(event.action)) + +fun Logger.warn(event: ProcessingLogEvent, cause: Throwable? = null) = + trackerBuilder(event).setCause(cause).log(trackerMessage(event.action)) + +fun Logger.error(event: ProcessingLogEvent, cause: Throwable? = null) = + trackerBuilder(event).setCause(cause).log(trackerMessage(event.action)) + +private fun Logger.trackerBuilder(event: ProcessingLogEvent) = eventBuilder(event, atDebug()) + +private fun trackerMessage(action: String) = "$TRACKER_PREFIX $action" + +fun withProcessingLogContext(sourcePostId: Long, flow: String, action: () -> T): T { + val previous = MDC.getCopyOfContextMap() + return try { + MDC.put(ProcessingLogFields.SOURCE_POST_ID, sourcePostId.toString()) + MDC.put(ProcessingLogFields.FLOW, flow) + action() + } finally { + if (previous == null) MDC.clear() else MDC.setContextMap(previous) + } +} + +private fun eventBuilder(event: ProcessingLogEvent, builder: org.slf4j.spi.LoggingEventBuilder) = builder + .addKeyValue(ProcessingLogFields.EVENT_ACTION, event.action) + .addKeyValue(ProcessingLogFields.FLOW, event.flow) + .addKeyValue(ProcessingLogFields.STAGE, event.stage) + .apply { + event.outcome?.let { addKeyValue(ProcessingLogFields.EVENT_OUTCOME, it) } + event.sourcePostId?.let { addKeyValue(ProcessingLogFields.SOURCE_POST_ID, it) } + event.attempt?.let { addKeyValue(ProcessingLogFields.ATTEMPT, it) } + event.durationMs?.let { addKeyValue(ProcessingLogFields.EVENT_DURATION_MS, it) } + event.fields.filterValues { it != null }.forEach { (key, value) -> addKeyValue(key, value) } + } + +private const val TRACKER_PREFIX = "[PostParcingTracker]" diff --git a/nook-api-application/src/main/kotlin/org/every/nook/api/application/processing/ProcessingMetrics.kt b/nook-api-application/src/main/kotlin/org/every/nook/api/application/processing/ProcessingMetrics.kt index b2b2a97c..cea30dfa 100644 --- a/nook-api-application/src/main/kotlin/org/every/nook/api/application/processing/ProcessingMetrics.kt +++ b/nook-api-application/src/main/kotlin/org/every/nook/api/application/processing/ProcessingMetrics.kt @@ -1,5 +1,6 @@ package org.every.nook.api.application.processing +import org.slf4j.LoggerFactory import java.time.Clock import java.time.Duration @@ -35,6 +36,7 @@ inline fun ProcessingMetrics.measure( ): T { val startedAt = clock.instant() val result = runCatching(action) + val duration = Duration.between(startedAt, clock.instant()) record( ProcessingMetrics.Measurement( flow = flow, @@ -42,8 +44,40 @@ inline fun ProcessingMetrics.measure( postId = postId, attempt = attempt, outcome = if (result.isSuccess) ProcessingMetrics.Outcome.SUCCESS else ProcessingMetrics.Outcome.FAILURE, - duration = Duration.between(startedAt, clock.instant()), + duration = duration, ), ) + val event = ProcessingLogEvent( + action = "processing.stage.completed", + flow = flow, + stage = stage, + outcome = if (result.isSuccess) "success" else "failure", + sourcePostId = postId, + attempt = attempt, + durationMs = duration.toMillis(), + fields = result.exceptionOrNull()?.let { exception -> + mapOf( + ProcessingLogFields.FAILURE_TYPE to exception::class.simpleName, + ProcessingLogFields.FAILURE_REASON to exception.message?.take(MAX_FAILURE_REASON_LENGTH), + ) + }.orEmpty(), + ) + if (result.isSuccess && isHighVolumeStage(flow, stage)) { + processingStageLogger.debug(event) + } else if (result.isSuccess) { + processingStageLogger.info(event) + } else { + processingStageLogger.warn(event, result.exceptionOrNull()) + } return result.getOrThrow() } + +@PublishedApi +internal val processingStageLogger = LoggerFactory.getLogger("processing.stage") + +@PublishedApi +internal const val MAX_FAILURE_REASON_LENGTH = 500 + +@PublishedApi +internal fun isHighVolumeStage(flow: String, stage: String): Boolean = + flow == "post-media" || (flow == "place" && stage == "search") diff --git a/nook-api-application/src/test/kotlin/org/every/nook/api/application/processing/ProcessingLogTest.kt b/nook-api-application/src/test/kotlin/org/every/nook/api/application/processing/ProcessingLogTest.kt new file mode 100644 index 00000000..b974c8b5 --- /dev/null +++ b/nook-api-application/src/test/kotlin/org/every/nook/api/application/processing/ProcessingLogTest.kt @@ -0,0 +1,24 @@ +package org.every.nook.api.application.processing + +import org.slf4j.LoggerFactory +import kotlin.test.Test +import kotlin.test.assertEquals + +class ProcessingLogTest { + @Test + fun `processing context returns action result when logging backend has no MDC adapter`() { + assertEquals(7, withProcessingLogContext(42, "place") { 7 }) + } + + @Test + fun `structured event can omit nullable fields`() { + LoggerFactory.getLogger(ProcessingLogTest::class.java).info( + ProcessingLogEvent( + action = "place.job.started", + flow = "place", + stage = "job", + fields = mapOf("nullable" to null), + ), + ) + } +} diff --git a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/ApifyInstagramPostContentExtractor.kt b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/ApifyInstagramPostContentExtractor.kt index 281b6edf..060cec6a 100644 --- a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/ApifyInstagramPostContentExtractor.kt +++ b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/ApifyInstagramPostContentExtractor.kt @@ -25,14 +25,17 @@ class ApifyInstagramPostContentExtractor( override fun supports(url: String): Boolean = InstagramContentUrl.supports(url) override fun extract(url: String): ExtractedPostContent { + val startedAt = System.nanoTime() val instagramUrl = InstagramContentUrl.parse(url) responseCache.find(PROVIDER, SOURCE_TYPE, instagramUrl.shortcode)?.let { cached -> + logger.logCacheHit(PROVIDER, startedAt) return mapResponse(instagramUrl, parseResponse(cached)) } if (properties.apiToken.isBlank()) { providerFailure() } val responseBody = try { + logger.logProviderRequestStarted(PROVIDER) restClient.post() .uri { builder -> builder.path(RUN_SYNC_PATH) @@ -52,16 +55,19 @@ class ApifyInstagramPostContentExtractor( .retrieve() .body(String::class.java) } catch (exception: RestClientResponseException) { + logger.logProviderRequestFailed(PROVIDER, startedAt, exception, exception.statusCode.value()) if (exception.statusCode.value() in TIMEOUT_STATUSES) { providerTimeout(exception) } providerFailure(exception) } catch (exception: ResourceAccessException) { + logger.logProviderRequestFailed(PROVIDER, startedAt, exception) handleResourceAccessException(exception) } val extracted = mapResponse(instagramUrl, parseResponse(responseBody)) responseBody?.let { body -> responseCache.save(PROVIDER, SOURCE_TYPE, instagramUrl.shortcode, body) } + logger.logProviderRequestCompleted(PROVIDER, startedAt, extracted.post.media.size) return extracted } diff --git a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/BrightDataInstagramPostContentExtractor.kt b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/BrightDataInstagramPostContentExtractor.kt index 80df2d57..2ffda906 100644 --- a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/BrightDataInstagramPostContentExtractor.kt +++ b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/BrightDataInstagramPostContentExtractor.kt @@ -25,14 +25,17 @@ class BrightDataInstagramPostContentExtractor( override fun supports(url: String): Boolean = InstagramContentUrl.supports(url) override fun extract(url: String): ExtractedPostContent { + val startedAt = System.nanoTime() val instagramUrl = InstagramContentUrl.parse(url) responseCache.find(PROVIDER, SOURCE_TYPE, instagramUrl.shortcode)?.let { cached -> + logger.logCacheHit(PROVIDER, startedAt) return mapResponse(instagramUrl, parseResponse(cached)) } if (properties.apiToken.isBlank()) { providerFailure() } val responseBody = try { + logger.logProviderRequestStarted(PROVIDER) restClient.post() .uri { builder -> builder.path(SCRAPE_PATH) @@ -46,8 +49,10 @@ class BrightDataInstagramPostContentExtractor( .retrieve() .body(String::class.java) } catch (exception: RestClientResponseException) { + logger.logProviderRequestFailed(PROVIDER, startedAt, exception, exception.statusCode.value()) handleResponseException(exception) } catch (exception: ResourceAccessException) { + logger.logProviderRequestFailed(PROVIDER, startedAt, exception) handleResourceAccessException(exception) } @@ -55,6 +60,7 @@ class BrightDataInstagramPostContentExtractor( responseBody?.let { body -> responseCache.save(PROVIDER, SOURCE_TYPE, instagramUrl.shortcode, body) } + logger.logProviderRequestCompleted(PROVIDER, startedAt, extracted.post.media.size) return extracted } diff --git a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/InstagramPostContentExtractor.kt b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/InstagramPostContentExtractor.kt index f24fe213..4488b9bb 100644 --- a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/InstagramPostContentExtractor.kt +++ b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/InstagramPostContentExtractor.kt @@ -5,6 +5,9 @@ import org.every.nook.api.application.content.ExtractedPostContent import org.every.nook.api.application.content.PostContentExtractor import org.every.nook.api.application.content.PostContentProviderException import org.every.nook.api.application.content.PostContentProviderTimeoutException +import org.every.nook.api.application.processing.ProcessingLogEvent +import org.every.nook.api.application.processing.info +import org.every.nook.api.application.processing.warn import org.slf4j.LoggerFactory class InstagramPostContentExtractor( @@ -24,6 +27,15 @@ class InstagramPostContentExtractor( InstagramScrapingProviderMode.DEFAULT, ) } + logger.info( + ProcessingLogEvent( + action = "instagram.provider.mode.selected", + flow = "post-content", + stage = "extract", + outcome = "success", + fields = mapOf("provider.mode" to mode.name), + ), + ) return when (mode) { InstagramScrapingProviderMode.BRIGHT_DATA_ONLY -> brightDataExtractor.extract(url) InstagramScrapingProviderMode.APIFY_ONLY -> apifyExtractor.extract(url) @@ -35,9 +47,11 @@ class InstagramPostContentExtractor( private fun extractBrightDataWithFallback(url: String): ExtractedPostContent = try { brightDataExtractor.extract(url) } catch (exception: PostContentProviderTimeoutException) { + logFallback("bright_data", "apify", "timeout", exception) logger.warn("Bright Data timed out; falling back to Apify", exception) apifyExtractor.extract(url) } catch (exception: PostContentProviderException) { + logFallback("bright_data", "apify", "failure", exception) logger.warn("Bright Data failed; falling back to Apify", exception) apifyExtractor.extract(url) } @@ -45,13 +59,32 @@ class InstagramPostContentExtractor( private fun extractApifyWithFallback(url: String): ExtractedPostContent = try { apifyExtractor.extract(url) } catch (exception: PostContentProviderTimeoutException) { + logFallback("apify", "bright_data", "timeout", exception) logger.warn("Apify timed out; falling back to Bright Data", exception) brightDataExtractor.extract(url) } catch (exception: PostContentProviderException) { + logFallback("apify", "bright_data", "failure", exception) logger.warn("Apify failed; falling back to Bright Data", exception) brightDataExtractor.extract(url) } + private fun logFallback(from: String, to: String, reason: String, exception: Throwable) { + logger.warn( + ProcessingLogEvent( + action = "instagram.provider.fallback", + flow = "post-content", + stage = "extract", + outcome = "fallback", + fields = mapOf( + "provider.name" to from, + "provider.fallback_to" to to, + "failure.reason" to reason, + ), + ), + exception, + ) + } + private companion object { val logger = LoggerFactory.getLogger(InstagramPostContentExtractor::class.java) } diff --git a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/InstagramProviderLog.kt b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/InstagramProviderLog.kt new file mode 100644 index 00000000..79643572 --- /dev/null +++ b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/instagram/InstagramProviderLog.kt @@ -0,0 +1,68 @@ +package org.every.nook.api.infrastructure.instagram + +import org.every.nook.api.application.processing.ProcessingLogEvent +import org.every.nook.api.application.processing.info +import org.every.nook.api.application.processing.warn +import org.slf4j.Logger + +internal fun instagramProviderEvent( + provider: String, + action: String, + outcome: String, + startedAt: Long? = null, + fields: Map = emptyMap(), +) = ProcessingLogEvent( + action = action, + flow = "post-content", + stage = "extract", + outcome = outcome, + durationMs = startedAt?.let { (System.nanoTime() - it) / NANOS_PER_MILLISECOND }, + fields = fields + mapOf("provider.name" to provider.lowercase()), +) + +private const val NANOS_PER_MILLISECOND = 1_000_000 + +internal fun Logger.logCacheHit(provider: String, startedAt: Long) = info( + instagramProviderEvent( + provider, + "instagram.provider.cache.hit", + "success", + startedAt, + mapOf("cache.hit" to true), + ), +) + +internal fun Logger.logProviderRequestStarted(provider: String) = info( + instagramProviderEvent( + provider, + "instagram.provider.request.started", + "started", + fields = mapOf("cache.hit" to false), + ), +) + +internal fun Logger.logProviderRequestFailed( + provider: String, + startedAt: Long, + exception: Throwable, + statusCode: Int? = null, +) = warn( + instagramProviderEvent( + provider, + "instagram.provider.request.failed", + "failure", + startedAt, + mapOf("http.status_code" to statusCode), + ), + exception, +) + +internal fun Logger.logProviderRequestCompleted(provider: String, startedAt: Long, mediaCount: Int) = info( + instagramProviderEvent( + provider, + "instagram.provider.request.completed", + "success", + startedAt, + mapOf("content.media_count" to mediaCount), + ), +) diff --git a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/openai/OpenAiContentInferenceAdapter.kt b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/openai/OpenAiContentInferenceAdapter.kt index 81e21266..a4c1e757 100644 --- a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/openai/OpenAiContentInferenceAdapter.kt +++ b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/openai/OpenAiContentInferenceAdapter.kt @@ -9,7 +9,10 @@ import org.every.nook.api.application.place.PlaceClueExtractor import org.every.nook.api.application.place.PlaceTagEvidenceSource import org.every.nook.api.application.place.PlaceTagExtractor import org.every.nook.api.application.post.PostContentInference +import org.every.nook.api.application.processing.ProcessingLogEvent +import org.every.nook.api.application.processing.info import org.every.nook.api.domain.place.PlaceTag +import org.slf4j.LoggerFactory import org.springframework.web.client.RestClient import tools.jackson.databind.JsonNode import tools.jackson.databind.ObjectMapper @@ -108,6 +111,7 @@ class OpenAiContentInferenceAdapter( schema: Map, maxOutputTokens: Int, ): JsonNode { + val startedAt = System.nanoTime() require(properties.apiKey.isNotBlank()) { "OpenAI API key is not configured" } val request = mapOf( "model" to properties.model, @@ -143,6 +147,16 @@ class OpenAiContentInferenceAdapter( ?: error("OpenAI returned no structured output") return objectMapper.readTree(text).also { result -> logger.info { "OpenAI structured output received: name=$name, output=$result" } + eventLogger.info( + ProcessingLogEvent( + action = "openai.response.completed", + flow = "content-inference", + stage = name, + outcome = "success", + durationMs = (System.nanoTime() - startedAt) / NANOS_PER_MILLISECOND, + fields = mapOf("provider.name" to "openai", "openai.model" to properties.model), + ), + ) } } @@ -204,6 +218,8 @@ class OpenAiContentInferenceAdapter( private companion object { val logger = KotlinLogging.logger {} + val eventLogger = LoggerFactory.getLogger(OpenAiContentInferenceAdapter::class.java) + const val NANOS_PER_MILLISECOND = 1_000_000 const val MAX_TITLE_LENGTH = 25 const val MAX_PLACE_COUNT = 20 diff --git a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/openai/OpenAiImageTextExtractor.kt b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/openai/OpenAiImageTextExtractor.kt index 26a033a1..5a1cbdaa 100644 --- a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/openai/OpenAiImageTextExtractor.kt +++ b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/openai/OpenAiImageTextExtractor.kt @@ -2,6 +2,9 @@ package org.every.nook.api.infrastructure.openai import org.every.nook.api.application.place.ImageTextExtractor import org.every.nook.api.application.place.ImageTranscript +import org.every.nook.api.application.processing.ProcessingLogEvent +import org.every.nook.api.application.processing.info +import org.slf4j.LoggerFactory import org.springframework.web.client.RestClient import tools.jackson.databind.JsonNode import tools.jackson.databind.ObjectMapper @@ -12,6 +15,7 @@ class OpenAiImageTextExtractor( private val properties: OpenAiProperties, ) : ImageTextExtractor { override fun extract(request: ImageTextExtractor.Request): List { + val startedAt = System.nanoTime() require(request.images.isNotEmpty()) { "At least one image is required" } require(request.images.size <= MAX_BATCH_SIZE) { "Too many images in transcript request" } require(properties.apiKey.isNotBlank()) { "OpenAI API key is not configured" } @@ -37,6 +41,22 @@ class OpenAiImageTextExtractor( texts = image.path("texts").toList().map(JsonNode::asText).map(String::trim) .filter(String::isNotEmpty), ) + }.also { transcripts -> + logger.info( + ProcessingLogEvent( + action = "openai.response.completed", + flow = "place", + stage = "image-transcript", + outcome = "success", + durationMs = (System.nanoTime() - startedAt) / NANOS_PER_MILLISECOND, + fields = mapOf( + "provider.name" to "openai", + "openai.model" to properties.model, + "content.image_count" to request.images.size, + "ocr.transcript_count" to transcripts.size, + ), + ), + ) } } @@ -73,6 +93,8 @@ class OpenAiImageTextExtractor( ) private companion object { + val logger = LoggerFactory.getLogger(OpenAiImageTextExtractor::class.java) + const val NANOS_PER_MILLISECOND = 1_000_000 const val MAX_BATCH_SIZE = 5 const val MAX_IMAGE_COUNT = 20 const val MAX_OUTPUT_TOKENS = 4000 diff --git a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/persistence/save/PostPersistenceAdapter.kt b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/persistence/save/PostPersistenceAdapter.kt index 510b56e7..0f1473b2 100644 --- a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/persistence/save/PostPersistenceAdapter.kt +++ b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/persistence/save/PostPersistenceAdapter.kt @@ -13,6 +13,8 @@ import org.every.nook.api.application.post.port.FindPostPlaceParsingPort import org.every.nook.api.application.post.port.PostPlaceParsingSnapshot import org.every.nook.api.application.post.port.ReusePostPort import org.every.nook.api.application.post.port.UpdatePostMemoPort +import org.every.nook.api.application.processing.ProcessingLogEvent +import org.every.nook.api.application.processing.info import org.every.nook.api.domain.place.GeoPoint import org.every.nook.api.domain.place.Place import org.every.nook.api.domain.place.PlaceParsingStatus @@ -33,6 +35,7 @@ import org.every.nook.api.infrastructure.persistence.post.PostContentParsingJobJ import org.every.nook.api.infrastructure.persistence.post.PostEntity import org.every.nook.api.infrastructure.persistence.post.PostJpaRepository import org.every.nook.api.infrastructure.persistence.post.PostPlaceJpaRepository +import org.slf4j.LoggerFactory import org.springframework.context.ApplicationEventPublisher import org.springframework.stereotype.Component import org.springframework.transaction.annotation.Transactional @@ -86,7 +89,8 @@ class PostPersistenceAdapter( nextAttemptAt = clock.instant(), ), ) - val userPost = findOrCreateUserPost(userId, sourcePostId, memo).entity + val userPostCreation = findOrCreateUserPost(userId, sourcePostId, memo) + val userPost = userPostCreation.entity addToGroups(requireNotNull(userPost.id), groupIds) if (existingContentJob == null) { eventPublisher.publishEvent(PostContentParsingJobRequestedEvent(sourcePostId, clock.instant())) @@ -94,11 +98,13 @@ class PostPersistenceAdapter( restartFailedJob(sourcePostId, contentJob) } - return CreatedPost( + val createdPost = CreatedPost( postId = requireNotNull(userPost.id), contentParsingStatus = contentJob.status, placeParsingStatus = placeParsingJobJpaRepository.findByPostId(sourcePostId)?.status, ) + logSavedPostMapping(sourcePostId, createdPost.postId, userPostCreation.created, "post.save.completed") + return createdPost } @Transactional @@ -119,11 +125,13 @@ class PostPersistenceAdapter( addToGroups(requireNotNull(userPost.id), groupIds) restartFailedJob(sourcePostId, contentJob) val placeParsingJob = placeParsingJobJpaRepository.findByPostId(sourcePostId) - return CreatedPost( + val createdPost = CreatedPost( postId = requireNotNull(userPost.id), contentParsingStatus = contentJob.status, placeParsingStatus = placeParsingJob?.status, ) + logSavedPostMapping(sourcePostId, createdPost.postId, userPostCreation.created, "post.reuse.completed") + return createdPost } @Transactional @@ -275,3 +283,21 @@ class PostPersistenceAdapter( id = id, ) } + +private fun logSavedPostMapping(sourcePostId: Long, savedPostId: Long, created: Boolean, action: String) { + postPersistenceEventLogger.info( + ProcessingLogEvent( + action = action, + flow = "post-save", + stage = "persist", + outcome = "success", + sourcePostId = sourcePostId, + fields = mapOf( + "saved_post.id" to savedPostId, + "saved_post.created" to created, + ), + ), + ) +} + +private val postPersistenceEventLogger = LoggerFactory.getLogger(PostPersistenceAdapter::class.java) diff --git a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/CompositePlaceSearchProvider.kt b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/CompositePlaceSearchProvider.kt index 7008b569..62522a9a 100644 --- a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/CompositePlaceSearchProvider.kt +++ b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/CompositePlaceSearchProvider.kt @@ -5,6 +5,10 @@ import org.every.nook.api.application.place.PlaceCandidate import org.every.nook.api.application.place.PlaceSearchProvider import org.every.nook.api.application.place.PlaceSearchProviderException import org.every.nook.api.application.place.PlaceSearchProviderTimeoutException +import org.every.nook.api.application.processing.ProcessingLogEvent +import org.every.nook.api.application.processing.info +import org.every.nook.api.application.processing.warn +import org.slf4j.LoggerFactory import java.util.concurrent.ExecutorCompletionService import java.util.concurrent.ExecutorService @@ -26,6 +30,7 @@ class CompositePlaceSearchProvider( pending.asSequence() .filterNot { future -> future.isDone } .forEach { future -> future.cancel(true) } + eventLogger.info(result.event("place.search.selected", "success")) return result.candidates } } @@ -40,6 +45,7 @@ class CompositePlaceSearchProvider( val emptySuccess = results.filterIsInstance().firstOrNull() if (emptySuccess != null) { + eventLogger.info(emptySuccess.event("place.search.completed", "empty")) return emptySuccess.candidates } @@ -53,13 +59,26 @@ class CompositePlaceSearchProvider( private fun NamedPlaceSearchProvider.searchSafely(request: PlaceSearchProvider.Request): ProviderResult = try { ProviderResult.Success(name, provider.search(request)) } catch (exception: PlaceSearchProviderTimeoutException) { + eventLogger.warn(ProviderResult.Timeout(name).event("place.provider.search.failed", "timeout"), exception) logger.warn(exception) { "Place search provider timed out: provider=$name, query=${request.query}" } ProviderResult.Timeout(name) } catch (exception: PlaceSearchProviderException) { + eventLogger.warn(ProviderResult.Failure(name).event("place.provider.search.failed", "failure"), exception) logger.warn(exception) { "Place search provider failed: provider=$name, query=${request.query}" } ProviderResult.Failure(name) } + private fun ProviderResult.event(action: String, outcome: String) = ProcessingLogEvent( + action = action, + flow = "place", + stage = "search", + outcome = outcome, + fields = mapOf( + "provider.name" to provider, + "provider.result_count" to (this as? ProviderResult.Success)?.candidates?.size, + ), + ) + data class NamedPlaceSearchProvider(val name: String, val provider: PlaceSearchProvider) private sealed interface ProviderResult { @@ -72,5 +91,6 @@ class CompositePlaceSearchProvider( private companion object { val logger = KotlinLogging.logger {} + val eventLogger = LoggerFactory.getLogger(CompositePlaceSearchProvider::class.java) } } diff --git a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/GooglePlacePhotoLog.kt b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/GooglePlacePhotoLog.kt new file mode 100644 index 00000000..16f304f7 --- /dev/null +++ b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/GooglePlacePhotoLog.kt @@ -0,0 +1,97 @@ +package org.every.nook.api.infrastructure.place + +import org.every.nook.api.application.place.PlaceCandidate +import org.every.nook.api.application.processing.ProcessingLogEvent +import org.every.nook.api.application.processing.info +import org.every.nook.api.application.processing.warn +import org.slf4j.Logger + +internal fun PlaceCandidate.event( + action: String, + stage: String, + outcome: String, + fields: Map = emptyMap(), +) = ProcessingLogEvent( + action = action, + flow = "place-thumbnail", + stage = stage, + outcome = outcome, + fields = fields + mapOf( + "provider.name" to "google", + "place.source_provider" to provider, + "place.external_id" to externalPlaceId, + ), +) + +internal fun PlaceCandidate.photoEvent( + action: String, + stage: String, + sequence: Int, + exception: Throwable, + startedAt: Long? = null, +) = event( + action, + stage, + "failure", + mapOf( + "event.duration_ms" to startedAt?.let(::elapsedMillis), + "media.sequence" to sequence, + "failure.type" to exception::class.simpleName, + "failure.reason" to exception.message?.take(MAX_FAILURE_REASON_LENGTH), + ), +) + +internal fun failureFields(exception: Throwable): Map = mapOf( + "failure.type" to exception::class.simpleName, + "failure.reason" to exception.message?.take(MAX_FAILURE_REASON_LENGTH), +) + +internal fun elapsedMillis(startedAt: Long): Long = (System.nanoTime() - startedAt) / NANOS_PER_MILLISECOND + +internal fun Logger.logGoogleSkipped(place: PlaceCandidate) = warn( + place.event( + "google.place.skipped", + "configuration", + "skipped", + mapOf("skip.reason" to "invalid_configuration"), + ), +) + +internal fun Logger.logGooglePhotoList(place: PlaceCandidate, availableCount: Int, selectedCount: Int) = info( + place.event( + "google.photo.list.completed", + "google-photo-list", + if (selectedCount == 0) "empty" else "success", + mapOf( + "google.photo_available_count" to availableCount, + "google.photo_selected_count" to selectedCount, + "empty.reason" to if (selectedCount == 0) "no_photos" else null, + ), + ), +) + +internal fun Logger.logGooglePhotoPipeline(place: PlaceCandidate, selectedCount: Int, storedCount: Int) = info( + place.event( + "google.photo.pipeline.completed", + "google-photo-store", + if (storedCount == 0 && selectedCount > 0) "failure" else "success", + mapOf( + "google.photo_selected_count" to selectedCount, + "google.photo_stored_count" to storedCount, + "google.photo_failed_count" to selectedCount - storedCount, + ), + ), +) + +internal fun Logger.logGoogleFetchFailure(place: PlaceCandidate, exception: Throwable, startedAt: Long) = warn( + place.event( + "google.place.fetch.failed", + "google-place-match", + "failure", + failureFields(exception) + ("event.duration_ms" to elapsedMillis(startedAt)), + ), + exception, +) + +private const val MAX_FAILURE_REASON_LENGTH = 500 +private const val NANOS_PER_MILLISECOND = 1_000_000 diff --git a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/GooglePlacePhotoProvider.kt b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/GooglePlacePhotoProvider.kt index c5ccf064..4b7a2671 100644 --- a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/GooglePlacePhotoProvider.kt +++ b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/GooglePlacePhotoProvider.kt @@ -9,7 +9,11 @@ import org.every.nook.api.application.place.PlaceOpeningPoint import org.every.nook.api.application.place.PlaceSupplement import org.every.nook.api.application.place.PlaceThumbnailProvider import org.every.nook.api.application.post.port.PostMediaStoragePort +import org.every.nook.api.application.processing.debug +import org.every.nook.api.application.processing.info +import org.every.nook.api.application.processing.warn import org.every.nook.api.domain.post.PostMedia +import org.slf4j.LoggerFactory import org.springframework.http.MediaType import org.springframework.web.client.RestClient import java.math.BigDecimal @@ -27,12 +31,14 @@ class GooglePlacePhotoProvider( val shouldSkip = !properties.enabled || properties.apiKey.isBlank() return if (shouldSkip) { + eventLogger.logGoogleSkipped(place) logger.warn { "Google place photo skipped: reason=invalid_configuration, enabled=${properties.enabled}, " + "apiKeyConfigured=${properties.apiKey.isNotBlank()}" } null } else { + val startedAt = System.nanoTime() runCatching { logger.info { "Google place photo search started: provider=${place.provider}, " + @@ -46,11 +52,14 @@ class GooglePlacePhotoProvider( } null } else { - val photoUrls = googlePlace.photos.orEmpty() + val availablePhotos = googlePlace.photos.orEmpty() .mapNotNull(GooglePhoto::name) .distinct() .take(PlaceSupplement.MAX_PHOTO_COUNT) + eventLogger.logGooglePhotoList(place, googlePlace.photos.orEmpty().size, availablePhotos.size) + val photoUrls = availablePhotos .mapIndexedNotNull { sequence, photoName -> storePhoto(photoName, sequence, place) } + eventLogger.logGooglePhotoPipeline(place, availablePhotos.size, photoUrls.size) PlaceSupplement( openingHours = googlePlace.toOpeningHours(), photoUrls = photoUrls, @@ -58,6 +67,7 @@ class GooglePlacePhotoProvider( ) } }.onFailure { exception -> + eventLogger.logGoogleFetchFailure(place, exception, startedAt) logger.warn(exception) { "Google place photo fetch failed: provider=${place.provider}, " + "externalPlaceId=${place.externalPlaceId}, name=${place.name}" @@ -68,6 +78,7 @@ class GooglePlacePhotoProvider( private fun searchPlace(place: PlaceCandidate): GooglePlace? { place.googlePlaceId?.let { return getPlace(it) } + val startedAt = System.nanoTime() val response = restClient.post() .uri("/v1/places:searchText") .contentType(MediaType.APPLICATION_JSON) @@ -93,12 +104,21 @@ class GooglePlacePhotoProvider( val matched = scored.maxByOrNull { it.second } ?.takeIf { it.second >= MIN_MATCH_SCORE } ?.first - logger.debug { - "[PostParcingTracker] stage=GOOGLE_PLACE_MATCH status=COMPLETED " + - "provider=${place.provider} externalPlaceId=${place.externalPlaceId} " + - "candidateScores=${scored.map { "${it.first.placeId()}:${it.second}" }} " + - "selectedId=${matched?.placeId()}" - } + eventLogger.info( + place.event( + "google.place.match.completed", + SEARCH_STAGE, + if (matched == null) "empty" else "success", + mapOf( + "event.duration_ms" to elapsedMillis(startedAt), + "google.place_candidate_count" to response?.places.orEmpty().size, + "google.place_candidate_scores" to scored.map { "${it.first.placeId()}:${it.second}" }, + "google.place_selected_id" to matched?.placeId(), + "google.place_matched" to (matched != null), + "empty.reason" to if (matched == null) "place_not_matched" else null, + ), + ), + ) logger.info { "Google place photo search completed: provider=${place.provider}, " + "externalPlaceId=${place.externalPlaceId}, googlePlaceCount=${response?.places.orEmpty().size}, " + @@ -156,16 +176,37 @@ class GooglePlacePhotoProvider( private fun String.normalize(): String = lowercase().filter(Char::isLetterOrDigit) - private fun storePhoto(photoName: String, sequence: Int, place: PlaceCandidate): String? = runCatching { - fetchPhotoUri(photoName)?.let { photoUri -> + private fun storePhoto(photoName: String, sequence: Int, place: PlaceCandidate): String? { + val photoUri = runCatching { fetchPhotoUri(photoName, sequence, place) } + .onFailure { exception -> + eventLogger.warn( + place.photoEvent("google.photo.media.failed", PHOTO_MEDIA_STAGE, sequence, exception), + exception, + ) + }.getOrNull() ?: return null + val startedAt = System.nanoTime() + return runCatching { mediaStorage.store(PostMedia(PostMedia.MediaType.IMAGE, photoUri, sequence)).url - } - }.onFailure { exception -> - logger.warn(exception) { - "Google place photo storage failed: provider=${place.provider}, " + - "externalPlaceId=${place.externalPlaceId}, sequence=$sequence" - } - }.getOrNull() + }.onSuccess { + eventLogger.debug( + place.event( + "google.photo.store.completed", + PHOTO_STORE_STAGE, + "success", + mapOf("event.duration_ms" to elapsedMillis(startedAt), "media.sequence" to sequence), + ), + ) + }.onFailure { exception -> + eventLogger.warn( + place.photoEvent("google.photo.store.failed", PHOTO_STORE_STAGE, sequence, exception, startedAt), + exception, + ) + logger.warn(exception) { + "Google place photo storage failed: provider=${place.provider}, " + + "externalPlaceId=${place.externalPlaceId}, sequence=$sequence" + } + }.getOrNull() + } private fun GooglePlace.toOpeningHours(): PlaceOpeningHours? { val hours = regularOpeningHours ?: return null @@ -185,7 +226,8 @@ class GooglePlacePhotoProvider( return runCatching { PlaceOpeningPoint(validDay, hour ?: 0, minute ?: 0) }.getOrNull() } - private fun fetchPhotoUri(photoName: String): String? { + private fun fetchPhotoUri(photoName: String, sequence: Int, place: PlaceCandidate): String? { + val startedAt = System.nanoTime() val photoUri = restClient.get() .uri { builder -> builder @@ -198,6 +240,19 @@ class GooglePlacePhotoProvider( .retrieve() .body(PhotoMediaResponse::class.java) ?.photoUri + eventLogger.debug( + place.event( + "google.photo.media.completed", + PHOTO_MEDIA_STAGE, + if (photoUri == null) "empty" else "success", + mapOf( + "event.duration_ms" to elapsedMillis(startedAt), + "media.sequence" to sequence, + "google.photo_uri_found" to (photoUri != null), + "empty.reason" to if (photoUri == null) "photo_uri_missing" else null, + ), + ), + ) logger.info { "Google place photo media completed: photoUriFound=${photoUri != null}" } return photoUri } @@ -266,6 +321,7 @@ class GooglePlacePhotoProvider( private companion object { val logger = KotlinLogging.logger {} + val eventLogger = LoggerFactory.getLogger(GooglePlacePhotoProvider::class.java) const val API_KEY_HEADER = "X-Goog-Api-Key" const val FIELD_MASK_HEADER = "X-Goog-FieldMask" const val DETAIL_FIELD_MASK = @@ -285,6 +341,10 @@ class GooglePlacePhotoProvider( const val FAR_MATCH_DISTANCE_METERS = 2_000.0 const val CLOSE_MATCH_DISTANCE_METERS = 100.0 const val EARTH_RADIUS_METERS = 6_371_000.0 + const val SEARCH_STAGE = "google-place-match" + const val PHOTO_LIST_STAGE = "google-photo-list" + const val PHOTO_MEDIA_STAGE = "google-photo-media" + const val PHOTO_STORE_STAGE = "google-photo-store" fun distanceMeters(lat1: BigDecimal, lon1: BigDecimal, lat2: BigDecimal, lon2: BigDecimal): Double { val latitudeDelta = Math.toRadians(lat2.toDouble() - lat1.toDouble()) diff --git a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/KakaoPlaceSearchProvider.kt b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/KakaoPlaceSearchProvider.kt index 72d2f307..fa9fa0b9 100644 --- a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/KakaoPlaceSearchProvider.kt +++ b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/KakaoPlaceSearchProvider.kt @@ -7,6 +7,9 @@ import org.every.nook.api.application.place.PlaceCandidatePage import org.every.nook.api.application.place.PlaceSearchProvider import org.every.nook.api.application.place.PlaceSearchProviderException import org.every.nook.api.application.place.PlaceSearchProviderTimeoutException +import org.every.nook.api.application.processing.ProcessingLogEvent +import org.every.nook.api.application.processing.info +import org.slf4j.LoggerFactory import org.springframework.web.client.ResourceAccessException import org.springframework.web.client.RestClient import org.springframework.web.client.RestClientResponseException @@ -24,6 +27,7 @@ class KakaoPlaceSearchProvider( searchPage(request.copy(page = FIRST_PAGE, size = DEFAULT_RESULT_SIZE)).items override fun searchPage(request: PlaceSearchProvider.Request): PlaceCandidatePage { + val startedAt = System.nanoTime() ensureConfigured() val responseBody = try { restClient.get() @@ -57,6 +61,21 @@ class KakaoPlaceSearchProvider( size = request.size, hasNext = !response.meta.isEnd, ).also { page -> + eventLogger.info( + ProcessingLogEvent( + action = "place.provider.search.completed", + flow = "place", + stage = "search", + outcome = "success", + durationMs = (System.nanoTime() - startedAt) / NANOS_PER_MILLISECOND, + fields = mapOf( + "provider.name" to "kakao", + "provider.result_count" to page.items.size, + "provider.page" to request.page, + "provider.has_next" to page.hasNext, + ), + ), + ) logger.info { "Kakao place search completed: query=${request.query}, page=${request.page}, " + "candidateCount=${page.items.size}" @@ -83,6 +102,8 @@ class KakaoPlaceSearchProvider( private companion object { val logger = KotlinLogging.logger {} + val eventLogger = LoggerFactory.getLogger(KakaoPlaceSearchProvider::class.java) + const val NANOS_PER_MILLISECOND = 1_000_000 const val SEARCH_PATH = "/v2/local/search/keyword.json" const val QUERY = "query" diff --git a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/NaverPlaceSearchProvider.kt b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/NaverPlaceSearchProvider.kt index a1b7e6fb..38be9ba0 100644 --- a/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/NaverPlaceSearchProvider.kt +++ b/nook-api-infrastructure/src/main/kotlin/org/every/nook/api/infrastructure/place/NaverPlaceSearchProvider.kt @@ -5,6 +5,9 @@ import org.every.nook.api.application.place.PlaceCandidate import org.every.nook.api.application.place.PlaceSearchProvider import org.every.nook.api.application.place.PlaceSearchProviderException import org.every.nook.api.application.place.PlaceSearchProviderTimeoutException +import org.every.nook.api.application.processing.ProcessingLogEvent +import org.every.nook.api.application.processing.info +import org.slf4j.LoggerFactory import org.springframework.http.MediaType import org.springframework.web.client.ResourceAccessException import org.springframework.web.client.RestClient @@ -19,6 +22,7 @@ class NaverPlaceSearchProvider( private val mapper: NaverPlaceMapper, ) : PlaceSearchProvider { override fun search(request: PlaceSearchProvider.Request): List { + val startedAt = System.nanoTime() ensureConfigured() val responseBody = try { restClient.get() @@ -46,6 +50,19 @@ class NaverPlaceSearchProvider( return runCatching { mapper.map(request.query, objectMapper.readValue(responseBody, NaverPlaceResponse::class.java)) .also { candidates -> + eventLogger.info( + ProcessingLogEvent( + action = "place.provider.search.completed", + flow = "place", + stage = "search", + outcome = "success", + durationMs = (System.nanoTime() - startedAt) / NANOS_PER_MILLISECOND, + fields = mapOf( + "provider.name" to "naver", + "provider.result_count" to candidates.size, + ), + ), + ) logger.info { "Naver local search completed: query=${request.query}, candidateCount=${candidates.size}" } @@ -71,6 +88,8 @@ class NaverPlaceSearchProvider( private companion object { val logger = KotlinLogging.logger {} + val eventLogger = LoggerFactory.getLogger(NaverPlaceSearchProvider::class.java) + const val NANOS_PER_MILLISECOND = 1_000_000 const val SEARCH_PATH = "/search/v1/local" const val QUERY = "query" diff --git a/nook-api-presentation/src/main/kotlin/org/every/nook/api/place/PlaceParsingEventListener.kt b/nook-api-presentation/src/main/kotlin/org/every/nook/api/place/PlaceParsingEventListener.kt index 634e2c47..e1bbc926 100644 --- a/nook-api-presentation/src/main/kotlin/org/every/nook/api/place/PlaceParsingEventListener.kt +++ b/nook-api-presentation/src/main/kotlin/org/every/nook/api/place/PlaceParsingEventListener.kt @@ -10,6 +10,7 @@ import org.every.nook.api.application.place.StorePlaceTagsUseCase import org.every.nook.api.application.place.StorePlaceThumbnailUseCase import org.every.nook.api.application.processing.NoOpProcessingMetrics import org.every.nook.api.application.processing.ProcessingMetrics +import org.every.nook.api.application.processing.withProcessingLogContext import org.springframework.beans.factory.annotation.Qualifier import org.springframework.boot.context.event.ApplicationReadyEvent import org.springframework.context.ApplicationEventPublisher @@ -69,34 +70,38 @@ class PlaceParsingEventListener( @Async("placeParsingTaskExecutor") @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT, fallbackExecution = true) fun process(event: PlaceParsingJobRequestedEvent) { - logger.info { - "Place parsing event received: postId=${event.postId}, availableAt=${event.availableAt}" - } - if (scheduleWhenUnavailable(event)) { - return - } - recordQueueDelay(event.availableAt, event.postId) - when (val result = processPlaceParsingJob(event.postId)) { - is ProcessPlaceParsingJobUseCase.Result.Retry -> schedule( - PlaceParsingJobRequestedEvent(event.postId, result.nextAttemptAt), - ) + withProcessingLogContext(event.postId, PLACE_FLOW) { + logger.info { + "Place parsing event received: postId=${event.postId}, availableAt=${event.availableAt}" + } + if (scheduleWhenUnavailable(event)) { + return@withProcessingLogContext + } + recordQueueDelay(event.availableAt, event.postId) + when (val result = processPlaceParsingJob(event.postId)) { + is ProcessPlaceParsingJobUseCase.Result.Retry -> schedule( + PlaceParsingJobRequestedEvent(event.postId, result.nextAttemptAt), + ) - ProcessPlaceParsingJobUseCase.Result.Completed, - ProcessPlaceParsingJobUseCase.Result.Failed, - ProcessPlaceParsingJobUseCase.Result.Skipped, - -> Unit + ProcessPlaceParsingJobUseCase.Result.Completed, + ProcessPlaceParsingJobUseCase.Result.Failed, + ProcessPlaceParsingJobUseCase.Result.Skipped, + -> Unit + } } } @Async("placeSupplementTaskExecutor") @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT, fallbackExecution = true) fun storeThumbnail(event: PlaceThumbnailRequestedEvent) { - runCatching { - storePlaceThumbnail(event.postId, event.place) - }.onFailure { exception -> - logger.warn(exception) { - "Place thumbnail storage failed: postId=${event.postId}, provider=${event.place.provider}, " + - "externalPlaceId=${event.place.externalPlaceId}" + withProcessingLogContext(event.postId, THUMBNAIL_FLOW) { + runCatching { + storePlaceThumbnail(event.postId, event.place) + }.onFailure { exception -> + logger.warn(exception) { + "Place thumbnail storage failed: postId=${event.postId}, provider=${event.place.provider}, " + + "externalPlaceId=${event.place.externalPlaceId}" + } } } } @@ -104,12 +109,14 @@ class PlaceParsingEventListener( @Async("placeParsingTaskExecutor") @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT, fallbackExecution = true) fun storeTags(event: PlaceTagsRequestedEvent) { - runCatching { - storePlaceTags(event) - }.onFailure { exception -> - logger.warn(exception) { - "Place tag storage failed: postId=${event.postId}, placeId=${event.placeId}, " + - "provider=${event.place.provider}, externalPlaceId=${event.place.externalPlaceId}" + withProcessingLogContext(event.postId, TAG_FLOW) { + runCatching { + storePlaceTags(event) + }.onFailure { exception -> + logger.warn(exception) { + "Place tag storage failed: postId=${event.postId}, placeId=${event.placeId}, " + + "provider=${event.place.provider}, externalPlaceId=${event.place.externalPlaceId}" + } } } } @@ -148,6 +155,8 @@ class PlaceParsingEventListener( private companion object { val logger = KotlinLogging.logger {} const val PLACE_FLOW = "place" + const val THUMBNAIL_FLOW = "place-thumbnail" + const val TAG_FLOW = "place-tags" const val QUEUE_STAGE = "queue" } } diff --git a/nook-api-presentation/src/main/kotlin/org/every/nook/api/post/PostContentParsingEventListener.kt b/nook-api-presentation/src/main/kotlin/org/every/nook/api/post/PostContentParsingEventListener.kt index 5cbdbb9b..946ef9f9 100644 --- a/nook-api-presentation/src/main/kotlin/org/every/nook/api/post/PostContentParsingEventListener.kt +++ b/nook-api-presentation/src/main/kotlin/org/every/nook/api/post/PostContentParsingEventListener.kt @@ -8,6 +8,7 @@ import org.every.nook.api.application.post.ProcessPostContentParsingJobUseCase import org.every.nook.api.application.post.StorePostMediaUseCase import org.every.nook.api.application.processing.NoOpProcessingMetrics import org.every.nook.api.application.processing.ProcessingMetrics +import org.every.nook.api.application.processing.withProcessingLogContext import org.springframework.beans.factory.annotation.Qualifier import org.springframework.beans.factory.annotation.Value import org.springframework.boot.context.event.ApplicationReadyEvent @@ -69,45 +70,49 @@ class PostContentParsingEventListener( @Async("postContentParsingTaskExecutor") @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT, fallbackExecution = true) fun process(event: PostContentParsingJobRequestedEvent) { - logger.info { - "Post content parsing event received: postId=${event.postId}, availableAt=${event.availableAt}" - } - if (scheduleWhenUnavailable(event)) { - return - } - recordQueueDelay(event.availableAt, event.postId) - when (val result = processPostContentParsingJob(event.postId)) { - is ProcessPostContentParsingJobUseCase.Result.Retry -> schedule( - PostContentParsingJobRequestedEvent(event.postId, result.nextAttemptAt), - ) + withProcessingLogContext(event.postId, CONTENT_FLOW) { + logger.info { + "Post content parsing event received: postId=${event.postId}, availableAt=${event.availableAt}" + } + if (scheduleWhenUnavailable(event)) { + return@withProcessingLogContext + } + recordQueueDelay(event.availableAt, event.postId) + when (val result = processPostContentParsingJob(event.postId)) { + is ProcessPostContentParsingJobUseCase.Result.Retry -> schedule( + PostContentParsingJobRequestedEvent(event.postId, result.nextAttemptAt), + ) - ProcessPostContentParsingJobUseCase.Result.Completed, - ProcessPostContentParsingJobUseCase.Result.Failed, - ProcessPostContentParsingJobUseCase.Result.Skipped, - -> Unit + ProcessPostContentParsingJobUseCase.Result.Completed, + ProcessPostContentParsingJobUseCase.Result.Failed, + ProcessPostContentParsingJobUseCase.Result.Skipped, + -> Unit + } } } @Async("postContentParsingTaskExecutor") @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT, fallbackExecution = true) fun storeMedia(event: PostMediaStorageRequestedEvent) { - runCatching { - storePostMedia( - event.postId, - StorePostMediaUseCase.Command(event.mediaType, event.sourceUrl, event.sequence), - ) - }.onFailure { exception -> - if (event.attempt < MEDIA_MAX_ATTEMPTS) { - val retryAt = clock.instant().plus(mediaRetryBackoff) - scheduleMedia(event.copy(attempt = event.attempt + 1, availableAt = retryAt)) - logger.warn(exception) { - "Post media storage retry scheduled: postId=${event.postId}, sequence=${event.sequence}, " + - "attempt=${event.attempt}, nextAttemptAt=$retryAt" - } - } else { - logger.error(exception) { - "Post media storage failed permanently: postId=${event.postId}, " + - "sequence=${event.sequence}, attempt=${event.attempt}" + withProcessingLogContext(event.postId, MEDIA_FLOW) { + runCatching { + storePostMedia( + event.postId, + StorePostMediaUseCase.Command(event.mediaType, event.sourceUrl, event.sequence), + ) + }.onFailure { exception -> + if (event.attempt < MEDIA_MAX_ATTEMPTS) { + val retryAt = clock.instant().plus(mediaRetryBackoff) + scheduleMedia(event.copy(attempt = event.attempt + 1, availableAt = retryAt)) + logger.warn(exception) { + "Post media storage retry scheduled: postId=${event.postId}, sequence=${event.sequence}, " + + "attempt=${event.attempt}, nextAttemptAt=$retryAt" + } + } else { + logger.error(exception) { + "Post media storage failed permanently: postId=${event.postId}, " + + "sequence=${event.sequence}, attempt=${event.attempt}" + } } } } @@ -154,6 +159,7 @@ class PostContentParsingEventListener( private companion object { val logger = KotlinLogging.logger {} const val CONTENT_FLOW = "post-content" + const val MEDIA_FLOW = "post-media" const val QUEUE_STAGE = "queue" const val MEDIA_MAX_ATTEMPTS = 4 const val DEFAULT_MEDIA_RETRY_BACKOFF_SECONDS = 3L