Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -3,9 +3,11 @@ package com.wafflestudio.spring2025
import org.springframework.boot.autoconfigure.SpringBootApplication
import org.springframework.boot.context.properties.ConfigurationPropertiesScan
import org.springframework.boot.runApplication
import org.springframework.scheduling.annotation.EnableScheduling

@SpringBootApplication
@ConfigurationPropertiesScan
@EnableScheduling
class SeminarApplication

fun main(args: Array<String>) {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
package com.wafflestudio.spring2025.common.email.outbox.event

data class EmailOutboxCreatedEvent(
val outboxId: Long,
)
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
package com.wafflestudio.spring2025.common.email.outbox.model

import org.springframework.data.annotation.CreatedDate
import org.springframework.data.annotation.Id
import org.springframework.data.relational.core.mapping.Column
import org.springframework.data.relational.core.mapping.Table
import java.time.Instant

@Table("email_outbox")
class EmailOutbox(
@Id var id: Long? = null,
var messageKey: String,
var eventType: EmailOutboxEventType,
var recipientEmail: String,
@Column("payload_json")
var payloadJson: String,
var status: EmailOutboxStatus = EmailOutboxStatus.PENDING,
var retryCount: Int = 0,
var maxRetryCount: Int = 5,
var nextRetryAt: Instant = Instant.now(),
var lastError: String? = null,
var sentAt: Instant? = null,
@CreatedDate
var createdAt: Instant? = null,
var updatedAt: Instant? = null,
)
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
package com.wafflestudio.spring2025.common.email.outbox.model

enum class EmailOutboxEventType {
REGISTRATION_STATUS,
REGISTRATION_DELETE,
WAITLIST_PROMOTION,
REGISTRATION_DEMOTION,
EVENT_CANCELLATION,
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
package com.wafflestudio.spring2025.common.email.outbox.model

enum class EmailOutboxStatus {
PENDING,
PROCESSING,
SENT,
FAILED,
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,105 @@
package com.wafflestudio.spring2025.common.email.outbox.repository

import com.wafflestudio.spring2025.common.email.outbox.model.EmailOutboxStatus
import org.springframework.jdbc.core.namedparam.MapSqlParameterSource
import org.springframework.jdbc.core.namedparam.NamedParameterJdbcTemplate
import org.springframework.stereotype.Repository
import java.sql.Timestamp
import java.time.Instant

@Repository
class EmailOutboxCommandRepository(
private val jdbcTemplate: NamedParameterJdbcTemplate,
) {
fun findProcessableIds(
limit: Int,
now: Instant = Instant.now(),
): List<Long> {
val params =
MapSqlParameterSource()
.addValue("limit", limit)
.addValue("now", Timestamp.from(now))
return jdbcTemplate.query(
"""
SELECT id
FROM email_outbox
WHERE status IN ('PENDING', 'FAILED')
AND retry_count < max_retry_count
AND next_retry_at <= :now
ORDER BY id ASC
LIMIT :limit
""".trimIndent(),
params,
) { rs, _ -> rs.getLong("id") }
}

fun claimForProcessing(
id: Long,
now: Instant = Instant.now(),
): Boolean {
val params =
MapSqlParameterSource()
.addValue("id", id)
.addValue("processing", EmailOutboxStatus.PROCESSING.name)
.addValue("now", Timestamp.from(now))
val updated =
jdbcTemplate.update(
"""
UPDATE email_outbox
SET status = :processing,
updated_at = CURRENT_TIMESTAMP(6)
WHERE id = :id
AND status IN ('PENDING', 'FAILED')
AND retry_count < max_retry_count
AND next_retry_at <= :now
""".trimIndent(),
params,
)
return updated == 1
}

fun markSent(id: Long) {
val params =
MapSqlParameterSource()
.addValue("id", id)
.addValue("sent", EmailOutboxStatus.SENT.name)
jdbcTemplate.update(
"""
UPDATE email_outbox
SET status = :sent,
sent_at = CURRENT_TIMESTAMP(6),
last_error = NULL,
updated_at = CURRENT_TIMESTAMP(6)
WHERE id = :id
""".trimIndent(),
params,
)
}

fun markFailed(
id: Long,
retryCount: Int,
nextRetryAt: Instant,
lastError: String,
) {
val params =
MapSqlParameterSource()
.addValue("id", id)
.addValue("failed", EmailOutboxStatus.FAILED.name)
.addValue("retryCount", retryCount)
.addValue("nextRetryAt", Timestamp.from(nextRetryAt))
.addValue("lastError", lastError)
jdbcTemplate.update(
"""
UPDATE email_outbox
SET status = :failed,
retry_count = :retryCount,
next_retry_at = :nextRetryAt,
last_error = :lastError,
updated_at = CURRENT_TIMESTAMP(6)
WHERE id = :id
""".trimIndent(),
params,
)
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
package com.wafflestudio.spring2025.common.email.outbox.repository

import com.wafflestudio.spring2025.common.email.outbox.model.EmailOutbox
import org.springframework.data.repository.ListCrudRepository

interface EmailOutboxRepository : ListCrudRepository<EmailOutbox, Long>
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
package com.wafflestudio.spring2025.common.email.outbox.service

import com.fasterxml.jackson.databind.ObjectMapper
import com.wafflestudio.spring2025.common.email.outbox.event.EmailOutboxCreatedEvent
import com.wafflestudio.spring2025.common.email.outbox.model.EmailOutbox
import com.wafflestudio.spring2025.common.email.outbox.model.EmailOutboxEventType
import com.wafflestudio.spring2025.common.email.outbox.repository.EmailOutboxRepository
import com.wafflestudio.spring2025.common.email.service.EmailService
import org.slf4j.LoggerFactory
import org.springframework.beans.factory.annotation.Value
import org.springframework.context.ApplicationEventPublisher
import org.springframework.dao.DataIntegrityViolationException
import org.springframework.dao.DuplicateKeyException
import org.springframework.stereotype.Component
import java.nio.charset.StandardCharsets
import java.security.MessageDigest
import java.sql.SQLIntegrityConstraintViolationException
import java.time.Instant

@Component
class EmailOutboxProducer(
private val emailOutboxRepository: EmailOutboxRepository,
private val objectMapper: ObjectMapper,
private val eventPublisher: ApplicationEventPublisher,
@Value("\${email.outbox.max-retry:5}")
private val maxRetryCount: Int,
) {
private val logger = LoggerFactory.getLogger(javaClass)

fun enqueueRegistrationStatus(data: EmailService.RegistrationStatusEmailData) {
val registrationPublicId = data.registrationPublicId ?: return
saveOutbox(
eventType = EmailOutboxEventType.REGISTRATION_STATUS,
messageKey = "$registrationPublicId:${data.status.name}",
recipientEmail = data.toEmail,
payload = data,
)
}

fun enqueueRegistrationDelete(data: EmailService.RegistrationDeleteEmailData) {
saveOutbox(
eventType = EmailOutboxEventType.REGISTRATION_DELETE,
messageKey = "${data.registrationPublicId}:DELETE",
recipientEmail = data.toEmail,
payload = data,
)
}

fun enqueueWaitlistPromotion(data: EmailService.WaitlistPromotionEmailData) {
saveOutbox(
eventType = EmailOutboxEventType.WAITLIST_PROMOTION,
messageKey = "${data.registrationPublicId}:PROMOTION",
recipientEmail = data.toEmail,
payload = data,
)
}

fun enqueueRegistrationDemotion(data: EmailService.DemotionEmailData) {
saveOutbox(
eventType = EmailOutboxEventType.REGISTRATION_DEMOTION,
messageKey = "${data.registrationPublicId}:DEMOTION",
recipientEmail = data.toEmail,
payload = data,
)
}

fun enqueueEventCancellation(data: EmailService.EventCancellationEmailData) {
saveOutbox(
eventType = EmailOutboxEventType.EVENT_CANCELLATION,
messageKey = "${data.eventPublicId}:${data.toEmail}:EVENT_CANCELLATION",
recipientEmail = data.toEmail,
payload = data,
)
}

private fun saveOutbox(
eventType: EmailOutboxEventType,
messageKey: String,
recipientEmail: String,
payload: Any,
) {
val payloadJson = objectMapper.writeValueAsString(payload)
val normalizedMessageKey = "${eventType.name}:${sha256(messageKey)}"
val now = Instant.now()
val outbox =
EmailOutbox(
messageKey = normalizedMessageKey,
eventType = eventType,
recipientEmail = recipientEmail,
payloadJson = payloadJson,
maxRetryCount = maxRetryCount,
nextRetryAt = now,
createdAt = now,
updatedAt = now,
)

try {
val saved = emailOutboxRepository.save(outbox)
saved.id?.let { id ->
eventPublisher.publishEvent(EmailOutboxCreatedEvent(id))
}
} catch (ex: Exception) {
if (isDuplicateOutboxException(ex)) {
logger.info("중복 outbox 메시지 스킵: {}", normalizedMessageKey)
} else {
throw ex
}
}
}

private fun isDuplicateOutboxException(ex: Throwable): Boolean {
var current: Throwable? = ex
while (current != null) {
if (current is DuplicateKeyException ||
current is DataIntegrityViolationException ||
current is SQLIntegrityConstraintViolationException
) {
return true
}
current = current.cause
}
return false
}

private fun sha256(source: String): String {
val digest = MessageDigest.getInstance("SHA-256")
val bytes = digest.digest(source.toByteArray(StandardCharsets.UTF_8))
return bytes.joinToString("") { "%02x".format(it) }
}
}
Loading
Loading