Skip to content
Merged
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
6 changes: 3 additions & 3 deletions cmd/issuetracker/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -274,7 +274,7 @@ func main() {
// Redis 초기화 실패 시에도 크롤링이 중단되지 않도록 graceful degrade 합니다.
var procLock locks.ProcessingLock
var ingestionLock locks.IngestionLock
var retryScheduler crawlerWorker.RetryScheduler
var retryScheduler publisher.RetryScheduler
var retrySchedulerStop func()
var redisClientShared *redis.Client // failure counter wiring 에서 재사용
redisCfg, err := config.LoadRedis()
Expand All @@ -297,14 +297,14 @@ func main() {
// Delayed retry queue: retry 를 Redis ZSET 에 보관하고 별도
// goroutine 이 ScheduledAt 도달 시 Kafka 에 발행 — worker 슬롯 점유 회피.
// Redis 부재 시 worker 가 lazy 로 KafkaImmediateRetryScheduler 를 사용 (기존 동작).
retryCfg := crawlerWorker.DefaultRedisRetrySchedulerConfig()
retryCfg := publisher.DefaultRedisRetrySchedulerConfig()
// idle heartbeat 압축 (이슈 #370) — pkg/config 로 env 로드 일관성 유지.
retrySchedCfg, retrySchedErr := config.LoadRetryScheduler()
if retrySchedErr != nil {
log.WithError(retrySchedErr).Fatal("RETRY_HEARTBEAT_EVERY_N_IDLE_TICKS 로드 실패")
}
retryCfg.HeartbeatEveryNIdleTicks = retrySchedCfg.HeartbeatEveryNIdleTicks
redisRetry := crawlerWorker.NewRedisDelayedRetryScheduler(
redisRetry := publisher.NewRedisDelayedRetryScheduler(
redisClient, crawlerProducer,
retryCfg,
log,
Expand Down
5 changes: 3 additions & 2 deletions internal/processor/fetcher/worker/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (

"issuetracker/internal/locks"
"issuetracker/internal/processor/fetcher/core"
"issuetracker/internal/publisher"
"issuetracker/internal/storage/service"
"issuetracker/pkg/config"
"issuetracker/pkg/logger"
Expand Down Expand Up @@ -49,7 +50,7 @@ type ManagerConfig struct {
Normal PoolConfig
Low PoolConfig
ProcessingLock locks.ProcessingLock
RetryScheduler RetryScheduler
RetryScheduler publisher.RetryScheduler

// MaxConcurrentPerStage: fetcher stage 의 Semaphore capacity 설정값 (이슈 #356).
// 0 이하 → 각 pool 의 WorkerCount/2 (floor) 자동.
Expand Down Expand Up @@ -168,7 +169,7 @@ func (m *PoolManager) Publish(ctx context.Context, job *core.CrawlJob) error {
return fmt.Errorf("marshal job %s: %w", job.ID, err)
}

topic := topicForPriority(priority)
topic := publisher.CrawlTopic(priority)

msg := queue.Message{
Topic: topic,
Expand Down
27 changes: 8 additions & 19 deletions internal/processor/fetcher/worker/pool.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (

"issuetracker/internal/locks"
"issuetracker/internal/processor/fetcher/core"
"issuetracker/internal/publisher"
"issuetracker/internal/storage/service"
"issuetracker/pkg/links"
"issuetracker/pkg/logger"
Expand Down Expand Up @@ -84,7 +85,7 @@ type KafkaConsumerPool struct {
// 미설정(nil) 이면 lazy 로 KafkaImmediateRetryScheduler 가 사용되어 기존 동작 유지 —
// 즉시 priority 토픽에 publish + worker 가 ScheduledAt 까지 sleep.
// atomic.Pointer 로 race-safe 설정 — Start 이후 SetRetryScheduler 호출에도 안전.
retryScheduler atomic.Pointer[retrySchedulerHolder]
retryScheduler atomic.Pointer[publisher.RetrySchedulerHolder]
// pollDone은 pollMessages goroutine이 완전히 종료됐음을 알리는 신호입니다.
// close(p.jobs) 전에 반드시 이 채널이 닫혔음을 확인해야 합니다.
pollDone chan struct{}
Expand Down Expand Up @@ -177,12 +178,12 @@ func NewKafkaConsumerPoolWithOptions(
// SetRetryScheduler 는 requeueWithRetry 시 사용할 RetryScheduler 를 설정합니다.
// nil 전달 시 fallback 인 KafkaImmediateRetryScheduler (lazy 생성) 가 사용됩니다 —
// 기존 동작 보존. atomic 교체로 Start 이후에도 race-safe.
func (p *KafkaConsumerPool) SetRetryScheduler(rs RetryScheduler) {
func (p *KafkaConsumerPool) SetRetryScheduler(rs publisher.RetryScheduler) {
if rs == nil {
p.retryScheduler.Store(nil)
return
}
p.retryScheduler.Store(&retrySchedulerHolder{s: rs})
p.retryScheduler.Store(&publisher.RetrySchedulerHolder{Scheduler: rs})
}

// SetGate 는 processJob 진입 시 URL 검사에 사용할 urlguard.Gate 를 설정합니다.
Expand Down Expand Up @@ -762,7 +763,7 @@ func (p *KafkaConsumerPool) requeueWithRetry(ctx context.Context, job *core.Craw
job.ScheduledAt = time.Now().Add(backoffDelay)

scheduler := p.resolveRetryScheduler()
topic := topicForPriority(job.Priority)
topic := publisher.CrawlTopic(job.Priority)

log.WithFields(map[string]interface{}{
"job_id": job.ID,
Expand Down Expand Up @@ -796,11 +797,11 @@ func (p *KafkaConsumerPool) requeueWithRetry(ctx context.Context, job *core.Craw
// resolveRetryScheduler 는 SetRetryScheduler 로 주입된 구현체를 반환하고, 미설정 시
// 기존 동작 (즉시 Kafka publish + worker sleep) 을 보존하는 KafkaImmediateRetryScheduler
// 를 lazy 생성합니다. 매 호출마다 새 인스턴스이지만 stateless 이므로 비용 무시 가능.
func (p *KafkaConsumerPool) resolveRetryScheduler() RetryScheduler {
func (p *KafkaConsumerPool) resolveRetryScheduler() publisher.RetryScheduler {
if h := p.retryScheduler.Load(); h != nil {
return h.s
return h.Scheduler
}
return NewKafkaImmediateRetryScheduler(p.producer)
return publisher.NewKafkaImmediateRetryScheduler(p.producer)
}

// logShutdownAware 는 graceful shutdown 으로 발생한 컨텍스트성 에러를 DEBUG 로,
Expand All @@ -824,15 +825,3 @@ func logShutdownAware(ctx context.Context, log *logger.Logger, err error, msg st
}
log.WithError(err).Error(msg)
}

// topicForPriority는 우선순위에 맞는 Kafka 토픽 이름을 반환합니다.
func topicForPriority(p core.Priority) string {
switch p {
case core.PriorityHigh:
return queue.TopicCrawlHigh
case core.PriorityLow:
return queue.TopicCrawlLow
default:
return queue.TopicCrawlNormal
}
}
10 changes: 6 additions & 4 deletions internal/publisher/publisher.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
// - PublishUpgrade : auto-upgrade (goquery → chromedp) republish (Sub 3 — pending)
//
// 외부 facade 단일화 — caller 는 *Publisher 의 메소드만 사용하면 됨. 내부 file 분리:
// - publisher.go : facade struct + 생성자 + 공통 Kafka helpers (buildMessage / crawlTopic / newJobID)
// - publisher.go : facade struct + 생성자 + 공통 Kafka helpers (buildMessage / CrawlTopic / newJobID)
// - chain.go : PublishChained 메소드 + 정규화 / guard / ingestion lock helper
// - guard.go : IngestionLock / PipelineGuard / atomic wrapper + Set* setters
//
Expand Down Expand Up @@ -84,7 +84,7 @@ func (p *Publisher) buildMessage(job *core.CrawlJob) (queue.Message, error) {
}

return queue.Message{
Topic: crawlTopic(job.Priority),
Topic: CrawlTopic(job.Priority),
Key: []byte(job.ID),
Value: data,
Headers: map[string]string{
Expand All @@ -94,8 +94,10 @@ func (p *Publisher) buildMessage(job *core.CrawlJob) (queue.Message, error) {
}, nil
}

// crawlTopic 은 Priority 에 대응하는 Kafka crawl 토픽 이름을 반환합니다.
func crawlTopic(p core.Priority) string {
// CrawlTopic 은 Priority 에 대응하는 Kafka crawl 토픽 이름을 반환합니다 (이슈 #389
// 피드백 — Kafka I/O 단일 책임 원칙에 따라 publisher 가 priority → topic 매핑의
// 유일한 출처. worker.topicForPriority / scheduler.crawlTopic 중복은 본 함수로 통합).
func CrawlTopic(p core.Priority) string {
switch p {
case core.PriorityHigh:
return queue.TopicCrawlHigh
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package worker
package publisher

import (
"context"
Expand All @@ -14,21 +14,30 @@ import (
pkgredis "issuetracker/pkg/redis"
)

// retrySchedulerHolder 는 atomic.Pointer 가 인터페이스를 직접 저장하지 못하므로
// retryDrainTimeout 은 graceful shutdown 으로 ctx 가 canceled 된 뒤에도 Redis enqueue
// 같은 cleanup 작업을 잠깐 더 허용하기 위한 별도 timeout 입니다 (이슈 #389 — 구
// fetcher/worker 의 drainTimeout 동등 값).
const retryDrainTimeout = 5 * time.Second

// RetrySchedulerHolder 는 atomic.Pointer 가 인터페이스를 직접 저장하지 못하므로
// RetryScheduler 인터페이스 값을 감싸 atomic 교체를 지원하는 wrapper 입니다.
type retrySchedulerHolder struct {
s RetryScheduler
//
// 호출자 (fetcher/worker pool) 가 atomic.Pointer[RetrySchedulerHolder] 필드로 보유 후
// SetRetryScheduler / 조회 시 사용.
type RetrySchedulerHolder struct {
Scheduler RetryScheduler
}
Comment thread
juhy0987 marked this conversation as resolved.

// RetryScheduler 는 처리 실패한 CrawlJob 의 재시도 발행 시점을 관리하는 인터페이스입니다
// (이슈 #389 — 메타 #385 의 Kafka I/O 단일 책임 원칙에 따라 publisher 패키지에서 정의).
//
// 두 가지 구현 전략을 추상화합니다:
// - KafkaImmediateRetryScheduler: 즉시 Kafka 에 재발행하고 worker 가 ScheduledAt 까지
// 대기 — 기존 동작 보존 (worker 슬롯 점유 문제 그대로)
// - RedisDelayedRetryScheduler: Redis ZSET 에 보관하고 별도 goroutine 이 ScheduledAt
// 도달 시 Kafka 에 발행 — worker 슬롯 점유 회피
//
// 호출자 (worker pool) 는 ScheduledAt 과 RetryCount 를 미리 셋팅한 job 을 전달합니다.
// 호출자 (fetcher/worker pool) 는 ScheduledAt 과 RetryCount 를 미리 셋팅한 job 을 전달합니다.
// 구현체는 lastErr 의 메시지를 last-error 헤더로 보존해야 합니다.
type RetryScheduler interface {
Enqueue(ctx context.Context, job *core.CrawlJob, lastErr error) error
Expand Down Expand Up @@ -61,7 +70,7 @@ func (s *KafkaImmediateRetryScheduler) Enqueue(ctx context.Context, job *core.Cr
}

msg := queue.Message{
Topic: topicForPriority(job.Priority),
Topic: CrawlTopic(job.Priority),
Key: []byte(job.ID),
Value: data,
Headers: retryHeaders(job, lastErr),
Expand Down Expand Up @@ -329,7 +338,7 @@ func (s *RedisDelayedRetryScheduler) republish(ctx context.Context, item pkgredi
}

msg := queue.Message{
Topic: topicForPriority(job.Priority),
Topic: CrawlTopic(job.Priority),
Key: []byte(item.JobID),
Value: entry.JobBytes,
Headers: retryHeaders(&job, lastErr),
Expand All @@ -353,7 +362,7 @@ func (s *RedisDelayedRetryScheduler) republish(ctx context.Context, item pkgredi
enqueueCtx := ctx
var cancelDrain context.CancelFunc
if ctx.Err() != nil {
enqueueCtx, cancelDrain = context.WithTimeout(context.WithoutCancel(ctx), drainTimeout)
enqueueCtx, cancelDrain = context.WithTimeout(context.WithoutCancel(ctx), retryDrainTimeout)
defer cancelDrain()
}
if reErr := s.client.EnqueueRetry(enqueueCtx, item.JobID, item.Payload, retryAt); reErr != nil {
Expand Down
16 changes: 2 additions & 14 deletions internal/scheduler/throttle.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,23 +5,11 @@ import (
"time"

"issuetracker/internal/processor/fetcher/core"
"issuetracker/internal/publisher"
"issuetracker/pkg/logger"
"issuetracker/pkg/queue"
)

// crawlTopic 은 Priority 에 대응하는 Kafka crawl 토픽 이름을 반환합니다.
// 이슈 #387 — 구 emitter.go 에서 throttle.go 의 backlog 조회용 helper 만 잔존.
func crawlTopic(p core.Priority) string {
switch p {
case core.PriorityHigh:
return queue.TopicCrawlHigh
case core.PriorityLow:
return queue.TopicCrawlLow
default:
return queue.TopicCrawlNormal
}
}

// Throttler 는 publish 직전에 호출되어 throttle 여부를 결정합니다.
// true 반환 시 Scheduler 는 emit 호출 없이 silent drop 합니다.
//
Expand Down Expand Up @@ -82,7 +70,7 @@ func (t *BacklogThrottler) ShouldThrottle(ctx context.Context, job *core.CrawlJo
return false
}

topic := crawlTopic(job.Priority)
topic := publisher.CrawlTopic(job.Priority)

checkCtx := ctx
if t.timeout > 0 {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,85 @@
package worker_test

import (
"context"
"errors"
"sync"
"testing"
"time"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/mock"

"issuetracker/internal/processor/fetcher/worker"
"issuetracker/internal/publisher"
"issuetracker/pkg/logger"
pkgredis "issuetracker/pkg/redis"
)

// poolRetryFakeQueue 는 publisher 측 retryQueueClient 인터페이스를 구조적으로 만족하는
// in-memory 더블입니다. pool.SetRetryScheduler 통합 시나리오 전용 (다른 publisher 단위 테스트는
// test/internal/publisher/retry_test.go 의 자체 fake 사용).
type poolRetryFakeQueue struct {
mu sync.Mutex
items []poolRetryItem
}

type poolRetryItem struct {
jobID string
payload []byte
scheduledAt time.Time
}

func newPoolRetryFakeQueue() *poolRetryFakeQueue { return &poolRetryFakeQueue{} }

func (f *poolRetryFakeQueue) EnqueueRetry(_ context.Context, jobID string, payload []byte, scheduledAt time.Time) error {
f.mu.Lock()
defer f.mu.Unlock()
f.items = append(f.items, poolRetryItem{jobID: jobID, payload: payload, scheduledAt: scheduledAt})
return nil
}

func (f *poolRetryFakeQueue) PeekDueRetries(_ context.Context, _ time.Time, _ int) ([]pkgredis.DueRetry, error) {
return nil, nil
}

func (f *poolRetryFakeQueue) AckRetry(_ context.Context, _ string) error { return nil }

func (f *poolRetryFakeQueue) snapshotLen() int {
f.mu.Lock()
defer f.mu.Unlock()
return len(f.items)
}

// TestKafkaConsumerPool_SetRetryScheduler_BypassesInlinePublish 는 SetRetryScheduler 로
// 주입한 publisher.RetryScheduler 구현체가 호출되어 inline producer.Publish 가 우회됨을 검증.
func TestKafkaConsumerPool_SetRetryScheduler_BypassesInlinePublish(t *testing.T) {
consumer := new(mockConsumer)
producer := new(mockProducer)
handler := new(mockJobHandler)
contentSvc := new(mockContentService)

pool := worker.NewKafkaConsumerPool(consumer, producer, handler, contentSvc, 1)

q := newPoolRetryFakeQueue()
customScheduler := publisher.NewRedisDelayedRetryScheduler(q, producer, publisher.RedisRetrySchedulerConfig{
PollInterval: time.Hour,
BatchSize: 1,
RepublishFailureBackoff: time.Hour,
}, logger.New(logger.DefaultConfig()))
pool.SetRetryScheduler(customScheduler)

job := newTestJob()
job.RetryCount = 1
msg := marshaledJobMsg(t, job)

handler.On("Handle", mock.Anything, job).Return(nil, errors.New("transient"))
consumer.On("CommitMessages", mock.Anything, mock.Anything).Return(nil)

runPool(t, consumer, pool, msg)

// inline producer.Publish 는 호출되지 않아야 함 — 모든 retry 가 Redis 경로로
producer.AssertNotCalled(t, "Publish", mock.Anything, mock.Anything)
assert.Equal(t, 1, q.snapshotLen(), "주입된 RedisDelayedRetryScheduler 가 enqueue 받았음")
consumer.AssertCalled(t, "CommitMessages", mock.Anything, mock.Anything)
}
Loading
Loading