feat(creator-community): 게시물 언어 감지 파이프라인을 추가한다

This commit is contained in:
2026-09-10 15:12:16 +09:00
parent 0284070667
commit 1b19ade69e
2 changed files with 876 additions and 2 deletions
@@ -7,8 +7,11 @@ import kr.co.vividnext.sodalive.content.category.CategoryRepository
import kr.co.vividnext.sodalive.content.comment.AudioContentCommentRepository
import kr.co.vividnext.sodalive.content.series.ContentSeriesRepository
import kr.co.vividnext.sodalive.explorer.profile.CreatorCheersRepository
import kr.co.vividnext.sodalive.explorer.profile.creatorCommunity.CreatorCommunityRepository
import kr.co.vividnext.sodalive.i18n.translation.LanguageTranslationEvent
import kr.co.vividnext.sodalive.i18n.translation.LanguageTranslationTargetType
import kr.co.vividnext.sodalive.i18n.translation.PapagoTranslationService.Companion.getTranslatableLanguageCodes
import kr.co.vividnext.sodalive.v2.creator.channel.community.translation.application.CreatorCommunityTranslationService
import org.slf4j.LoggerFactory
import org.springframework.beans.factory.annotation.Value
import org.springframework.context.ApplicationEventPublisher
@@ -22,8 +25,12 @@ import org.springframework.transaction.annotation.Propagation
import org.springframework.transaction.annotation.Transactional
import org.springframework.transaction.event.TransactionPhase
import org.springframework.transaction.event.TransactionalEventListener
import org.springframework.transaction.support.TransactionSynchronization
import org.springframework.transaction.support.TransactionSynchronizationManager
import org.springframework.util.LinkedMultiValueMap
import org.springframework.web.client.RestTemplate
import javax.persistence.EntityManager
import javax.persistence.LockModeType
/**
* 텍스트 기반 데이터(콘텐츠, 댓글 등)에 대해 파파고 언어 감지를 요청하기 위한 이벤트.
@@ -37,13 +44,16 @@ enum class LanguageDetectTargetType {
SERIES,
ORIGINAL_WORK,
CREATOR_CONTENT_CATEGORY
CREATOR_CONTENT_CATEGORY,
CREATOR_COMMUNITY
}
class LanguageDetectEvent(
val id: Long,
val query: String,
val targetType: LanguageDetectTargetType = LanguageDetectTargetType.CONTENT
val targetType: LanguageDetectTargetType = LanguageDetectTargetType.CONTENT,
val sourceRevision: Long? = null,
val targetLanguage: String? = null
)
data class PapagoLanguageDetectResponse(
@@ -60,7 +70,10 @@ class LanguageDetectListener(
private val seriesRepository: ContentSeriesRepository,
private val originalWorkRepository: OriginalWorkRepository,
private val categoryRepository: CategoryRepository,
private val creatorCommunityRepository: CreatorCommunityRepository,
private val languageDetectionCacheService: LanguageDetectionCacheService,
private val creatorCommunityTranslationService: CreatorCommunityTranslationService,
private val entityManager: EntityManager,
private val applicationEventPublisher: ApplicationEventPublisher,
@@ -95,6 +108,7 @@ class LanguageDetectListener(
LanguageDetectTargetType.SERIES -> handleSeriesLanguageDetect(event)
LanguageDetectTargetType.ORIGINAL_WORK -> handleOriginalWorkLanguageDetect(event)
LanguageDetectTargetType.CREATOR_CONTENT_CATEGORY -> handleCreatorContentCategoryLanguageDetect(event)
LanguageDetectTargetType.CREATOR_COMMUNITY -> handleCreatorCommunityLanguageDetect(event)
}
}
@@ -366,6 +380,53 @@ class LanguageDetectListener(
)
}
private fun handleCreatorCommunityLanguageDetect(event: LanguageDetectEvent) {
val sourceRevision = event.sourceRevision ?: return
val creatorCommunity = creatorCommunityRepository.findByIdAndIsActiveTrue(event.id) ?: return
if (
creatorCommunity.contentRevision != sourceRevision ||
creatorCommunity.content != event.query
) {
return
}
if (!creatorCommunity.languageCode.isNullOrBlank()) {
requestTranslationsAfterCommit(event)
return
}
val languageCode = detectLanguageCode(event, event.id) ?: return
if (languageCode !in getTranslatableLanguageCodes(null)) return
val currentCreatorCommunity = creatorCommunityRepository.findByIdAndIsActiveTrueForUpdate(event.id) ?: return
entityManager.refresh(currentCreatorCommunity, LockModeType.PESSIMISTIC_WRITE)
if (
currentCreatorCommunity.contentRevision != sourceRevision ||
currentCreatorCommunity.content != event.query
) {
return
}
if (!currentCreatorCommunity.languageCode.isNullOrBlank()) {
requestTranslationsAfterCommit(event)
return
}
currentCreatorCommunity.languageCode = languageCode
creatorCommunityRepository.save(currentCreatorCommunity)
requestTranslationsAfterCommit(event)
}
private fun requestTranslationsAfterCommit(event: LanguageDetectEvent) {
TransactionSynchronizationManager.registerSynchronization(
object : TransactionSynchronization {
override fun afterCommit() {
creatorCommunityTranslationService.requestTranslations(event.id, event.targetLanguage)
}
}
)
}
private fun detectLanguageCode(event: LanguageDetectEvent, targetIdForLog: Long): String? {
return languageDetectionCacheService.detectWithCache(event.query) {
requestPapagoLanguageCode(event.query, targetIdForLog)
@@ -0,0 +1,813 @@
package kr.co.vividnext.sodalive.content
import com.fasterxml.jackson.databind.ObjectMapper
import com.fasterxml.jackson.module.kotlin.KotlinModule
import kr.co.vividnext.sodalive.aws.cloudfront.AudioContentCloudFront
import kr.co.vividnext.sodalive.aws.s3.S3Uploader
import kr.co.vividnext.sodalive.can.payment.CanPaymentService
import kr.co.vividnext.sodalive.can.use.UseCanRepository
import kr.co.vividnext.sodalive.configs.QueryDslConfig
import kr.co.vividnext.sodalive.content.theme.AudioContentThemeQueryRepository
import kr.co.vividnext.sodalive.explorer.profile.creatorCommunity.CreatorCommunity
import kr.co.vividnext.sodalive.explorer.profile.creatorCommunity.CreatorCommunityRepository
import kr.co.vividnext.sodalive.explorer.profile.creatorCommunity.CreatorCommunityService
import kr.co.vividnext.sodalive.explorer.profile.creatorCommunity.comment.CreatorCommunityCommentRepository
import kr.co.vividnext.sodalive.explorer.profile.creatorCommunity.like.CreatorCommunityLikeRepository
import kr.co.vividnext.sodalive.i18n.LangContext
import kr.co.vividnext.sodalive.i18n.SodaMessageSource
import kr.co.vividnext.sodalive.i18n.translation.LanguageTranslationTargetType
import kr.co.vividnext.sodalive.i18n.translation.ResourceTranslationJobScheduler
import kr.co.vividnext.sodalive.i18n.translation.SourceTextNormalizer
import kr.co.vividnext.sodalive.i18n.translation.TranslationJobRepository
import kr.co.vividnext.sodalive.i18n.translation.TranslationJobScheduler
import kr.co.vividnext.sodalive.i18n.translation.TranslationMemory
import kr.co.vividnext.sodalive.i18n.translation.TranslationMemoryRepository
import kr.co.vividnext.sodalive.i18n.translation.TranslationReadModelMaterializer
import kr.co.vividnext.sodalive.i18n.translation.TranslationSourceExtractor
import kr.co.vividnext.sodalive.member.Member
import kr.co.vividnext.sodalive.member.MemberRepository
import kr.co.vividnext.sodalive.member.MemberRole
import kr.co.vividnext.sodalive.member.block.BlockMemberRepository
import kr.co.vividnext.sodalive.v2.creator.channel.community.translation.adapter.out.persistence.CreatorCommunityTranslation
import kr.co.vividnext.sodalive.v2.creator.channel.community.translation.adapter.out.persistence.CreatorCommunityTranslationRepository
import kr.co.vividnext.sodalive.v2.creator.channel.community.translation.application.CreatorCommunityTranslationService
import kr.co.vividnext.sodalive.v2.home.following.application.HomeFollowingNewsPublishService
import org.junit.jupiter.api.Assertions.assertEquals
import org.junit.jupiter.api.Assertions.assertNull
import org.junit.jupiter.api.Assertions.assertTrue
import org.junit.jupiter.api.BeforeEach
import org.junit.jupiter.api.DisplayName
import org.junit.jupiter.api.Test
import org.mockito.Mockito
import org.springframework.aop.support.AopUtils
import org.springframework.beans.factory.annotation.Autowired
import org.springframework.boot.test.autoconfigure.jdbc.AutoConfigureTestDatabase
import org.springframework.boot.test.autoconfigure.orm.jpa.DataJpaTest
import org.springframework.boot.test.context.TestConfiguration
import org.springframework.context.ApplicationEventPublisher
import org.springframework.context.annotation.Bean
import org.springframework.context.annotation.Import
import org.springframework.context.annotation.Primary
import org.springframework.scheduling.annotation.EnableAsync
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor
import org.springframework.test.util.AopTestUtils
import org.springframework.transaction.PlatformTransactionManager
import org.springframework.transaction.annotation.Propagation
import org.springframework.transaction.annotation.Transactional
import org.springframework.transaction.support.TransactionTemplate
import java.util.UUID
import java.util.concurrent.Callable
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.ConcurrentLinkedQueue
import java.util.concurrent.CountDownLatch
import java.util.concurrent.Executors
import java.util.concurrent.TimeUnit
import java.util.concurrent.atomic.AtomicInteger
import java.util.concurrent.atomic.AtomicReference
import javax.persistence.EntityManager
@DataJpaTest(
properties = [
"spring.cache.type=none",
"spring.datasource.url=" +
"jdbc:h2:mem:creator-community-language-detect;MODE=MySQL;NON_KEYWORDS=VALUE;DB_CLOSE_ON_EXIT=FALSE",
"cloud.naver.papago-client-id=test-client-id",
"cloud.naver.papago-client-secret=test-client-secret"
]
)
@AutoConfigureTestDatabase(replace = AutoConfigureTestDatabase.Replace.NONE)
@Import(
QueryDslConfig::class,
AudioContentThemeQueryRepository::class,
LanguageDetectListener::class,
CreatorCommunityTranslationService::class,
TranslationSourceExtractor::class,
TranslationReadModelMaterializer::class,
ResourceTranslationJobScheduler::class,
TranslationJobScheduler::class,
CreatorCommunityLanguageDetectTest.AsyncConfiguration::class
)
@Transactional(propagation = Propagation.NOT_SUPPORTED)
class CreatorCommunityLanguageDetectTest @Autowired constructor(
private val creatorCommunityRepository: CreatorCommunityRepository,
private val memberRepository: MemberRepository,
private val translationJobRepository: TranslationJobRepository,
private val applicationEventPublisher: ApplicationEventPublisher,
private val languageDetectListener: LanguageDetectListener,
private val creatorCommunityTranslationService: CreatorCommunityTranslationService,
private val creatorCommunityTranslationRepository: CreatorCommunityTranslationRepository,
private val creatorCommunityService: CreatorCommunityService,
private val languageDetectionCacheService: ControlledLanguageDetectionCacheService,
private val translationMemoryRepository: TranslationMemoryRepository,
private val translationReadModelMaterializer: TranslationReadModelMaterializer,
private val detailPipelineTaskExecutor: DetailPipelineTaskExecutor,
transactionManager: PlatformTransactionManager
) {
private val transactionTemplate = TransactionTemplate(transactionManager)
@BeforeEach
fun setUp() {
languageDetectionCacheService.clear()
}
@Test
@DisplayName("커뮤니티 감지 이벤트는 기존 기본값을 유지하면서 현재 본문 메타데이터를 전달한다")
fun shouldKeepLegacyDefaultsAndDeclareCreatorCommunityMetadata() {
val fields = LanguageDetectEvent::class.java.declaredFields.map { it.name }
assertTrue(LanguageDetectTargetType.values().any { it.name == "CREATOR_COMMUNITY" })
assertTrue(fields.contains("sourceRevision"))
assertTrue(fields.contains("targetLanguage"))
assertEquals(LanguageDetectTargetType.CONTENT, LanguageDetectEvent(1L, "legacy").targetType)
}
@Test
@DisplayName("감지 후 번역 요청 서비스는 독립 프록시 트랜잭션을 사용한다")
fun shouldExposeProxiedRequiresNewTranslationRequestEntryPoint() {
val annotation = AopUtils.getTargetClass(creatorCommunityTranslationService)
.getMethod("requestTranslations", Long::class.javaPrimitiveType, String::class.java)
.getAnnotation(Transactional::class.java)
assertTrue(AopUtils.isAopProxy(creatorCommunityTranslationService))
assertEquals(Propagation.REQUIRES_NEW, annotation.propagation)
}
@Test
@DisplayName("상세 파이프라인 executor는 이전 작업 완료로 다음 작업 latch를 해제하지 않는다")
fun shouldNotCompleteNextTaskWhenPreviousTaskFinishesAfterLatchIsReassigned() {
val executor = DetailPipelineTaskExecutor()
val firstStarted = CountDownLatch(1)
val firstRelease = CountDownLatch(1)
val secondStarted = CountDownLatch(1)
val secondRelease = CountDownLatch(1)
executor.initialize()
try {
executor.prepareNextTask()
executor.execute {
firstStarted.countDown()
check(firstRelease.await(5, TimeUnit.SECONDS))
}
assertTrue(firstStarted.await(5, TimeUnit.SECONDS))
executor.prepareNextTask()
executor.execute {
secondStarted.countDown()
check(secondRelease.await(5, TimeUnit.SECONDS))
}
firstRelease.countDown()
assertTrue(secondStarted.await(5, TimeUnit.SECONDS))
assertEquals(false, executor.awaitTaskCompletion(0, TimeUnit.MILLISECONDS))
} finally {
firstRelease.countDown()
secondRelease.countDown()
executor.shutdown()
assertTrue(executor.threadPoolExecutor.awaitTermination(5, TimeUnit.SECONDS))
}
}
@Test
@DisplayName("커밋 후 감지는 요청 언어의 번역 작업을 실제로 저장한다")
fun shouldPersistRequestedTargetJobAfterCommittedDetection() {
val postId = persistPost(content = "current source", contentRevision = 3)
givenDetectedLanguage("current source", "ko")
publishAfterCommit(
LanguageDetectEvent(
id = postId,
query = "current source",
targetType = LanguageDetectTargetType.CREATOR_COMMUNITY,
sourceRevision = 3,
targetLanguage = "ja"
)
)
assertJobTargets(postId, setOf("ja"))
assertPost(postId, content = "current source", languageCode = "ko")
}
@Test
@DisplayName("상세 조회는 실제 프록시 감지와 작업 저장 뒤 재료화한 요청 언어 번역문을 반환한다")
fun shouldRunDetailTranslationPipelineThroughProxyListenerAndScheduler() {
val sourceContent = "detail source content"
val postId = persistPost(content = sourceContent, contentRevision = 0, isCommentAvailable = false)
val blockingDetection = languageDetectionCacheService.block(sourceContent, "en")
detailPipelineTaskExecutor.prepareNextTask()
try {
assertTrue(AopUtils.isAopProxy(creatorCommunityService))
assertTrue(AopUtils.isAopProxy(creatorCommunityTranslationService))
val initialDetail = creatorCommunityService.getCommunityPostDetail(
postId = postId,
memberId = 999L,
timezone = "Asia/Seoul",
isAdult = true
)
assertEquals(sourceContent, initialDetail.content)
assertTrue(blockingDetection.started.await(5, TimeUnit.SECONDS))
assertEquals(emptySet<String>(), jobTargets(postId))
} finally {
blockingDetection.release.countDown()
}
assertTrue(detailPipelineTaskExecutor.awaitTaskCompletion(5, TimeUnit.SECONDS))
assertPost(postId, content = sourceContent, languageCode = "en")
assertEquals(setOf("ko"), jobTargets(postId))
assertEquals(1, jobCount(postId))
val repeatedDetail = creatorCommunityService.getCommunityPostDetail(
postId = postId,
memberId = 999L,
timezone = "Asia/Seoul",
isAdult = true
)
assertEquals(sourceContent, repeatedDetail.content)
val executor = Executors.newFixedThreadPool(2)
try {
val repeatedDetails = listOf(
executor.submit(
Callable {
creatorCommunityService.getCommunityPostDetail(postId, 999L, "Asia/Seoul", true).content
}
),
executor.submit(
Callable {
creatorCommunityService.getCommunityPostDetail(postId, 999L, "Asia/Seoul", true).content
}
)
).map { it.get(5, TimeUnit.SECONDS) }
assertEquals(listOf(sourceContent, sourceContent), repeatedDetails)
} finally {
executor.shutdownNow()
}
assertEquals(1, jobCount(postId))
transactionTemplate.executeWithoutResult {
translationMemoryRepository.save(
TranslationMemory(
sourceHash = SourceTextNormalizer.hash(sourceContent),
sourceText = sourceContent,
sourceLanguage = "en",
targetLanguage = "ko",
translatedText = "translated detail content",
provider = TranslationReadModelMaterializer.DEFAULT_PROVIDER,
providerVersion = "nmt-v1"
)
)
}
assertTrue(
translationReadModelMaterializer.materialize(
LanguageTranslationTargetType.CREATOR_COMMUNITY,
postId,
"ko"
)
)
val translatedDetail = creatorCommunityService.getCommunityPostDetail(
postId = postId,
memberId = 999L,
timezone = "Asia/Seoul",
isAdult = true
)
assertEquals("translated detail content", translatedDetail.content)
assertEquals(1, jobCount(postId))
}
@Test
@DisplayName("전체 범위 감지는 원문 언어를 제외한 실제 번역 작업을 저장한다")
fun shouldPersistAllNonSourceTargetJobsAfterCommittedDetection() {
val postId = persistPost(content = "English source", contentRevision = 1)
givenDetectedLanguage("English source", "en")
publishAfterCommit(
LanguageDetectEvent(
id = postId,
query = "English source",
targetType = LanguageDetectTargetType.CREATOR_COMMUNITY,
sourceRevision = 1
)
)
assertJobTargets(postId, setOf("ja", "ko"))
assertPost(postId, content = "English source", languageCode = "en")
}
@Test
@DisplayName("감지 중 변경된 본문은 refresh 뒤에도 이전 언어와 작업을 적용하지 않는다")
fun shouldRejectStalePostAfterBlockedDetection() {
val postId = persistPost(content = "before", contentRevision = 0)
val blockingDetection = languageDetectionCacheService.block("before", "ko")
val executor = Executors.newSingleThreadExecutor()
try {
val completion = executor.submit {
detectInTransaction(
LanguageDetectEvent(
id = postId,
query = "before",
targetType = LanguageDetectTargetType.CREATOR_COMMUNITY,
sourceRevision = 0,
targetLanguage = "ja"
)
)
}
assertTrue(blockingDetection.started.await(5, TimeUnit.SECONDS))
updatePost(postId, content = "after", contentRevision = 1)
blockingDetection.release.countDown()
completion.get(5, TimeUnit.SECONDS)
assertPost(postId, content = "after", languageCode = null)
assertEquals(emptySet<String>(), jobTargets(postId))
} finally {
blockingDetection.release.countDown()
executor.shutdownNow()
}
}
@Test
@DisplayName("동시 감지 중 첫 감지가 언어를 저장해도 두 번째 요청 언어를 예약한다")
fun shouldPreserveSecondTargetAfterConcurrentDetectionSetsLanguage() {
val postId = persistPost(content = "concurrent source", contentRevision = 0)
val concurrentDetection = languageDetectionCacheService.blockConcurrent("concurrent source", "ko")
val executor = Executors.newFixedThreadPool(2)
try {
val first = executor.submit {
detectInTransaction(
LanguageDetectEvent(
id = postId,
query = "concurrent source",
targetType = LanguageDetectTargetType.CREATOR_COMMUNITY,
sourceRevision = 0,
targetLanguage = "ja"
)
)
}
assertTrue(concurrentDetection.firstStarted.await(5, TimeUnit.SECONDS))
val second = executor.submit {
detectInTransaction(
LanguageDetectEvent(
id = postId,
query = "concurrent source",
targetType = LanguageDetectTargetType.CREATOR_COMMUNITY,
sourceRevision = 0,
targetLanguage = "en"
)
)
}
assertTrue(concurrentDetection.secondStarted.await(5, TimeUnit.SECONDS))
concurrentDetection.firstRelease.countDown()
first.get(5, TimeUnit.SECONDS)
concurrentDetection.secondRelease.countDown()
second.get(5, TimeUnit.SECONDS)
assertPost(postId, content = "concurrent source", languageCode = "ko")
assertJobTargets(postId, setOf("en", "ja"))
} finally {
concurrentDetection.firstRelease.countDown()
concurrentDetection.secondRelease.countDown()
executor.shutdownNow()
}
}
@Test
@DisplayName("언어가 이미 확정된 동시 목표 이벤트는 둘 다 번역 작업으로 이어진다")
fun shouldContinueBothRequestedTargetsWhenLanguageIsAlreadyKnown() {
val postId = persistPost(content = "known source", contentRevision = 2, languageCode = "ko")
publishAfterCommit(
LanguageDetectEvent(
id = postId,
query = "known source",
targetType = LanguageDetectTargetType.CREATOR_COMMUNITY,
sourceRevision = 2,
targetLanguage = "ja"
)
)
publishAfterCommit(
LanguageDetectEvent(
id = postId,
query = "known source",
targetType = LanguageDetectTargetType.CREATOR_COMMUNITY,
sourceRevision = 2,
targetLanguage = "en"
)
)
assertJobTargets(postId, setOf("en", "ja"))
assertTrue(languageDetectionCacheService.queries.isEmpty())
}
@Test
@DisplayName("빈 본문과 미지원 감지값과 감지 실패는 원문과 작업 상태를 유지한다")
fun shouldKeepOriginalWhenDetectionCannotProduceSupportedLanguage() {
val blankPostId = persistPost(content = " ", contentRevision = 0)
val unsupportedPostId = persistPost(content = "unsupported", contentRevision = 0)
val failedPostId = persistPost(content = "provider failure", contentRevision = 0)
givenDetectedLanguage("unsupported", "fr")
publishAfterCommit(communityEvent(blankPostId, " ", 0))
publishAfterCommit(communityEvent(unsupportedPostId, "unsupported", 0))
publishAfterCommit(communityEvent(failedPostId, "provider failure", 0))
assertTrue(await { languageDetectionCacheService.queries.contains("unsupported") })
assertTrue(await { languageDetectionCacheService.queries.contains("provider failure") })
assertPost(blankPostId, content = " ", languageCode = null)
assertPost(unsupportedPostId, content = "unsupported", languageCode = null)
assertPost(failedPostId, content = "provider failure", languageCode = null)
assertEquals(emptySet<String>(), jobTargets(blankPostId))
assertEquals(emptySet<String>(), jobTargets(unsupportedPostId))
assertEquals(emptySet<String>(), jobTargets(failedPostId))
}
@Test
@DisplayName("실제 프록시 생성 트랜잭션이 롤백되면 게시글과 번역 후속 작업을 남기지 않는다")
fun shouldLeaveNoPostOrTranslationArtifactsWhenProxiedCreateRollsBack() {
val creator = persistCreator()
val content = "rolled back community ${UUID.randomUUID()}"
val postId = transactionTemplate.execute { status ->
creatorCommunityService.createCommunityPost(
audioFile = null,
postImage = null,
requestString = createRequest(content),
member = creator
)
val post = creatorCommunityRepository.findAll().single { it.content == content }
status.setRollbackOnly()
post.id!!
}!!
assertTrue(AopUtils.isAopProxy(creatorCommunityService))
assertTrue(transactionTemplate.execute { creatorCommunityRepository.findById(postId).isEmpty }!!)
assertEquals(emptySet<String>(), jobTargets(postId))
assertTrue(
transactionTemplate.execute {
creatorCommunityTranslationRepository.findAll().none { it.creatorCommunityId == postId }
}!!
)
assertTrue(languageDetectionCacheService.queries.isEmpty())
}
@Test
@DisplayName("실제 영문 본문 수정은 이전 번역을 숨기고 커밋 후 영문 감지와 새 작업으로 이어진다")
fun shouldHideStaleTranslationAndScheduleNewTargetsAfterCommittedEnglishModification() {
val creator = persistCreator()
val koreanContent = "기존 한국어 본문"
val englishContent = "current English community source"
val postId = transactionTemplate.execute {
val post = creatorCommunityRepository.save(
CreatorCommunity(
content = koreanContent,
price = 0,
isCommentAvailable = true,
isAdult = false,
languageCode = "ko",
contentRevision = 0
).apply {
member = creator
}
)
creatorCommunityTranslationRepository.save(
CreatorCommunityTranslation(
creatorCommunityId = post.id!!,
locale = "en",
content = "old English translation",
sourceRevision = 0,
sourceHash = SourceTextNormalizer.hash(koreanContent),
sourceLanguage = "ko"
)
)
post.id!!
}!!
val blockingDetection = languageDetectionCacheService.block(englishContent, "en")
val executor = Executors.newSingleThreadExecutor()
try {
val modification = executor.submit {
transactionTemplate.executeWithoutResult {
creatorCommunityService.modifyCommunityPost(
postImage = null,
requestString = modifyRequest(postId, englishContent),
member = creator
)
}
}
assertTrue(blockingDetection.started.await(5, TimeUnit.SECONDS))
transactionTemplate.executeWithoutResult {
val post = creatorCommunityRepository.findById(postId).orElseThrow()
assertEquals(englishContent, post.content)
assertNull(post.languageCode)
assertEquals(1, post.contentRevision)
assertEquals(
englishContent,
creatorCommunityTranslationService.findDisplayContents(listOf(postId), "en")[postId]?.content
)
assertEquals(
"old English translation",
creatorCommunityTranslationRepository
.findByCreatorCommunityIdAndLocale(postId, "en")
?.content
)
assertEquals(emptySet<String>(), jobTargets(postId))
}
blockingDetection.release.countDown()
modification.get(5, TimeUnit.SECONDS)
assertTrue(
await {
transactionTemplate.execute {
creatorCommunityRepository.findById(postId).orElseThrow().languageCode == "en"
} == true
}
)
assertPost(postId, content = englishContent, languageCode = "en")
assertJobTargets(postId, setOf("ja", "ko"))
} finally {
blockingDetection.release.countDown()
executor.shutdownNow()
}
}
private fun communityEvent(postId: Long, content: String, sourceRevision: Long): LanguageDetectEvent {
return LanguageDetectEvent(
id = postId,
query = content,
targetType = LanguageDetectTargetType.CREATOR_COMMUNITY,
sourceRevision = sourceRevision
)
}
private fun givenDetectedLanguage(content: String, detectedLanguage: String) {
languageDetectionCacheService.detections[content] = detectedLanguage
}
private fun createRequest(content: String): String {
return """{"content":"$content","price":0,"isCommentAvailable":true,"isAdult":false}"""
}
private fun modifyRequest(postId: Long, content: String): String {
return """{"creatorCommunityId":$postId,"content":"$content"}"""
}
private fun publishAfterCommit(event: LanguageDetectEvent) {
transactionTemplate.executeWithoutResult {
applicationEventPublisher.publishEvent(event)
}
}
private fun detectInTransaction(event: LanguageDetectEvent) {
val listener: LanguageDetectListener = AopTestUtils.getTargetObject(languageDetectListener)
transactionTemplate.executeWithoutResult {
listener.detectLanguage(event)
}
}
private fun persistPost(
content: String,
contentRevision: Long,
languageCode: String? = null,
isCommentAvailable: Boolean = true
): Long {
return transactionTemplate.execute {
val member = memberRepository.save(
Member(
password = "password",
nickname = "language-detect-${UUID.randomUUID()}",
role = MemberRole.CREATOR
)
)
creatorCommunityRepository.save(
CreatorCommunity(
content = content,
price = 0,
isCommentAvailable = isCommentAvailable,
isAdult = false,
languageCode = languageCode,
contentRevision = contentRevision
).apply {
this.member = member
}
).id!!
}!!
}
private fun persistCreator(): Member {
return transactionTemplate.execute {
memberRepository.save(
Member(
password = "password",
nickname = "community-service-${UUID.randomUUID()}",
role = MemberRole.CREATOR
)
)
}!!
}
private fun updatePost(postId: Long, content: String, contentRevision: Long) {
transactionTemplate.executeWithoutResult {
val post = creatorCommunityRepository.findById(postId).orElseThrow()
post.content = content
post.contentRevision = contentRevision
post.languageCode = null
}
}
private fun assertPost(postId: Long, content: String, languageCode: String?) {
val post = transactionTemplate.execute {
creatorCommunityRepository.findById(postId).orElseThrow()
}!!
assertEquals(content, post.content)
if (languageCode == null) {
assertNull(post.languageCode)
} else {
assertEquals(languageCode, post.languageCode)
}
}
private fun assertJobTargets(postId: Long, expectedTargets: Set<String>) {
assertTrue(await { jobTargets(postId) == expectedTargets })
}
private fun jobTargets(postId: Long): Set<String> {
return transactionTemplate.execute {
translationJobRepository.findAll()
.filter {
it.resourceType == LanguageTranslationTargetType.CREATOR_COMMUNITY && it.resourceId == postId
}
.map { it.targetLanguage }
.toSet()
}.orEmpty()
}
private fun jobCount(postId: Long): Int {
return transactionTemplate.execute {
translationJobRepository.findAll().count {
it.resourceType == LanguageTranslationTargetType.CREATOR_COMMUNITY && it.resourceId == postId
}
}!!
}
private fun await(condition: () -> Boolean): Boolean {
repeat(50) {
if (condition()) return true
Thread.sleep(100)
}
return condition()
}
@TestConfiguration
@EnableAsync
class AsyncConfiguration {
@Bean(name = ["taskExecutor"])
fun taskExecutor(): DetailPipelineTaskExecutor {
return DetailPipelineTaskExecutor()
}
@Bean
@Primary
fun languageDetectionCacheService(
languageDetectionResultRepository: LanguageDetectionResultRepository
): ControlledLanguageDetectionCacheService {
return ControlledLanguageDetectionCacheService(languageDetectionResultRepository)
}
@Bean
fun creatorCommunityService(
repository: CreatorCommunityRepository,
blockMemberRepository: BlockMemberRepository,
likeRepository: CreatorCommunityLikeRepository,
commentRepository: CreatorCommunityCommentRepository,
useCanRepository: UseCanRepository,
applicationEventPublisher: ApplicationEventPublisher,
creatorCommunityTranslationService: CreatorCommunityTranslationService,
entityManager: EntityManager
): CreatorCommunityService {
return CreatorCommunityService(
canPaymentService = Mockito.mock(CanPaymentService::class.java),
repository = repository,
blockMemberRepository = blockMemberRepository,
likeRepository = likeRepository,
commentRepository = commentRepository,
useCanRepository = useCanRepository,
s3Uploader = Mockito.mock(S3Uploader::class.java),
objectMapper = ObjectMapper().registerModule(KotlinModule.Builder().build()),
audioContentCloudFront = Mockito.mock(AudioContentCloudFront::class.java),
applicationEventPublisher = applicationEventPublisher,
messageSource = SodaMessageSource(),
langContext = LangContext(),
homeFollowingNewsPublishService = Mockito.mock(HomeFollowingNewsPublishService::class.java),
creatorCommunityTranslationService = creatorCommunityTranslationService,
entityManager = entityManager,
imageBucket = "test-image-bucket",
contentBucket = "test-content-bucket",
imageHost = "https://test.cloudfront.net"
)
}
}
}
class ControlledLanguageDetectionCacheService(
languageDetectionResultRepository: LanguageDetectionResultRepository
) : LanguageDetectionCacheService(languageDetectionResultRepository) {
val detections = ConcurrentHashMap<String, String>()
val queries = ConcurrentLinkedQueue<String>()
private val blockingDetections = ConcurrentHashMap<String, BlockingDetection>()
private val concurrentDetections = ConcurrentHashMap<String, ConcurrentDetection>()
fun clear() {
detections.clear()
queries.clear()
blockingDetections.clear()
concurrentDetections.clear()
}
fun block(query: String, detectedLanguage: String): BlockingDetection {
return BlockingDetection(detectedLanguage).also { blockingDetections[query] = it }
}
fun blockConcurrent(query: String, detectedLanguage: String): ConcurrentDetection {
return ConcurrentDetection(detectedLanguage).also { concurrentDetections[query] = it }
}
override fun detectWithCache(query: String, provider: String, detector: () -> String?): String? {
queries.add(query)
val concurrentDetection = concurrentDetections[query]
if (concurrentDetection != null) {
return concurrentDetection.detect()
}
val blockingDetection = blockingDetections[query]
if (blockingDetection != null) {
blockingDetection.started.countDown()
check(blockingDetection.release.await(5, TimeUnit.SECONDS))
return blockingDetection.detectedLanguage
}
return detections[query]
}
}
class BlockingDetection(val detectedLanguage: String) {
val started = CountDownLatch(1)
val release = CountDownLatch(1)
}
class ConcurrentDetection(private val detectedLanguage: String) {
private val order = AtomicInteger()
val firstStarted = CountDownLatch(1)
val secondStarted = CountDownLatch(1)
val firstRelease = CountDownLatch(1)
val secondRelease = CountDownLatch(1)
fun detect(): String {
when (order.incrementAndGet()) {
1 -> {
firstStarted.countDown()
check(firstRelease.await(5, TimeUnit.SECONDS))
}
2 -> {
secondStarted.countDown()
check(secondRelease.await(5, TimeUnit.SECONDS))
}
}
return detectedLanguage
}
}
class DetailPipelineTaskExecutor : ThreadPoolTaskExecutor() {
private val completion = AtomicReference(CountDownLatch(0))
init {
corePoolSize = 1
maxPoolSize = 1
setThreadNamePrefix("detail-pipeline-")
setTaskDecorator { task ->
val taskCompletion = completion.get()
Runnable {
try {
task.run()
} finally {
taskCompletion.countDown()
}
}
}
}
fun prepareNextTask() {
completion.set(CountDownLatch(1))
}
fun awaitTaskCompletion(timeout: Long, unit: TimeUnit): Boolean {
return completion.get().await(timeout, unit)
}
}