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
10 changes: 7 additions & 3 deletions cmd/issuetracker/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -105,9 +105,13 @@ func main() {
crawlerProducer := queue.NewProducer(crawlerKafkaCfg)
defer crawlerProducer.Close()

resolver := crawlerWorker.NewCompositeResolver(core.PriorityNormal)
resolver.Add(crawlerWorker.NewSourcePriorityResolver(core.PriorityNormal))
resolver.Add(crawlerWorker.NewRuleBasedPriorityResolver(core.PriorityNormal))
// 이슈 #391 — PriorityResolver chain 이 publisher 측으로 이동 + 모든 PublishX 가
// resolver 통과 (메타 #385 Sub 6). ExplicitPriorityResolver 를 chain 1순위 로 등록 —
// 발행자가 job.Priority 를 사전 명시한 경우 (seed entry / retry / upgrade) 그 값이 보존됨.
resolver := publisher.NewCompositeResolver(core.PriorityNormal)
resolver.Add(&publisher.ExplicitPriorityResolver{})
resolver.Add(publisher.NewSourcePriorityResolver(core.PriorityNormal))
resolver.Add(publisher.NewRuleBasedPriorityResolver(core.PriorityNormal))

highConsumer := queue.NewConsumer(crawlerKafkaCfg, queue.TopicCrawlHigh)
defer highConsumer.Close()
Expand Down
26 changes: 16 additions & 10 deletions internal/processor/fetcher/worker/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,7 @@ type ManagerConfig struct {
type PoolManager struct {
pools map[core.Priority]*KafkaConsumerPool
pub *publisher.Publisher
resolver PriorityResolver
resolver publisher.PriorityResolver
log *logger.Logger

// chromedpPool 은 chromedp 전용 worker pool. nil 이면 비활성.
Expand All @@ -100,7 +100,7 @@ func NewPoolManager(
pub *publisher.Publisher,
handler JobHandler,
contentSvc service.ContentService,
resolver PriorityResolver,
resolver publisher.PriorityResolver,
log *logger.Logger,
) *PoolManager {
procLock := cfg.ProcessingLock
Expand Down Expand Up @@ -158,17 +158,23 @@ func NewPoolManager(
return mgr
}

// Publish는 CrawlJob을 PriorityResolver로 우선순위를 결정한 후 해당 Kafka crawl 토픽에 발행합니다.
// Publish는 CrawlJob을 publisher facade 로 발행합니다. priority 결정은 publisher 내부의
// resolver chain 이 buildMessage 에서 일괄 처리합니다 (이슈 #391 — 메타 #385 Sub 6).
//
// Publish resolves the job's priority via the configured PriorityResolver,
// updates job.Priority in-place, and publishes to the correct crawl topic
// (crawl.high / crawl.normal / crawl.low).
// Publish delegates to publisher.PublishJob. Priority resolution is handled inside the
// publisher's resolver chain via buildMessage — single source of truth across all
// PublishX paths (seed / chained / job / retry).
//
// gemini PR #400 피드백 — marshal/topic/headers 구성은 publisher.PublishJob 에 위임.
// manager 는 priority 결정 + 로깅만 책임 (priority resolver chain 통합은 Sub 6 에서).
// 로깅을 위해 동일 resolver 를 한 번 더 평가 — ExplicitPriorityResolver 가 1순위라
// 이후 buildMessage 의 평가와 같은 결과를 보장 (resolver 는 stateless · idempotent).
//
// m.resolver 가 nil 인 경우 (테스트 wiring) 는 job.Priority 를 그대로 사용 — coderabbit
Comment thread
juhy0987 marked this conversation as resolved.
// PR #409 피드백. publisher.buildMessage 가 nil resolver 를 fail-safe 로 허용하는 정책과 일관.
func (m *PoolManager) Publish(ctx context.Context, job *core.CrawlJob) error {
priority := m.resolver.Resolve(job)
job.Priority = priority
priority := job.Priority
if m.resolver != nil {
priority = m.resolver.Resolve(job)
}

Comment thread
juhy0987 marked this conversation as resolved.
m.log.WithFields(map[string]interface{}{
"job_id": job.ID,
Expand Down
3 changes: 1 addition & 2 deletions internal/publisher/chain.go
Original file line number Diff line number Diff line change
Expand Up @@ -116,8 +116,7 @@ func (p *Publisher) buildJobMessages(
MaxRetries: DefaultMaxRetries,
}

job.Priority = p.resolver.Resolve(job)

// 이슈 #391 — resolver chain 통과는 buildMessage 가 흡수 (모든 PublishX 일관성).
msg, err := p.buildMessage(job)
if err != nil {
return nil, fmt.Errorf("build message for %s: %w", url, err)
Expand Down
37 changes: 28 additions & 9 deletions internal/publisher/publisher.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,10 +52,13 @@ type Message = queue.Message
const DefaultMaxRetries = 3

// PriorityResolver 는 CrawlJob 의 우선순위를 결정하는 인터페이스입니다.
// worker.CompositeResolver 등이 이를 구현합니다.
// 구현체 — ExplicitPriorityResolver / SourcePriorityResolver / RuleBasedPriorityResolver /
// DefaultPriorityResolver / CompositeResolver (chain) — 모두 publisher 패키지의
// [resolver.go](resolver.go) 에 위치 (이슈 #391 — 메타 #385 Sub 6).
//
// 메타 #385 Sub 6 에서 본 인터페이스 및 chain 구현이 publisher 측으로 이동될 예정 —
// 그 시점에 본 인터페이스가 모든 PublishX 메소드의 priority 결정 단일 진입점이 됩니다.
// 모든 PublishX 메소드의 priority 결정 단일 진입점 — 운영자가 한 곳에만 resolver 룰 추가하면
// seed / chained / retry / upgrade 모든 경로에 일관 적용. ExplicitPriorityResolver 를 chain
// 1순위 로 등록하면 발행자가 사전 명시한 priority 가 보존됨.
type PriorityResolver interface {
Resolve(job *core.CrawlJob) core.Priority
}
Expand Down Expand Up @@ -135,19 +138,35 @@ func (p *Publisher) PublishJob(ctx context.Context, job *core.CrawlJob) error {
}

// buildMessage 는 CrawlJob 을 Kafka Message 로 변환합니다.
//
// 이슈 #391 — 메타 #385 Sub 6: resolver chain 통과를 본 헬퍼에 흡수. 모든 PublishX
// (Chained / Seed / Job / Retry / Redis retry republish) 가 buildMessage 를 거치므로
// priority 결정이 단일 출처. 발행자가 job.Priority 를 사전 명시한 경우, chain 1순위로
// 등록된 ExplicitPriorityResolver 가 그 값을 통과시켜 explicit 우선이 보존됩니다.
//
// resolver 가 nil 인 경우 (테스트 등) 는 기존 job.Priority 를 그대로 사용 — fail-safe.
//
// 부작용 회피 (gemini PR #409 피드백) — 외부에서 주입된 *job 의 Priority 를 직접 수정하지
// 않도록 local 복사본의 Priority 만 갱신. 호출자 (예: PoolManager.Publish) 가 원본 job 의
// Priority 변경을 기대하지 않더라도 안전. CrawlJob 은 작은 struct 이라 복사 비용 무시 가능.
func (p *Publisher) buildMessage(job *core.CrawlJob) (queue.Message, error) {
data, err := job.Marshal()
j := *job
if p.resolver != nil {
j.Priority = p.resolver.Resolve(&j)
}
Comment thread
juhy0987 marked this conversation as resolved.

data, err := j.Marshal()
if err != nil {
return queue.Message{}, fmt.Errorf("marshal job %s: %w", job.ID, err)
return queue.Message{}, fmt.Errorf("marshal job %s: %w", j.ID, err)
}

return queue.Message{
Topic: CrawlTopic(job.Priority),
Key: []byte(job.ID),
Topic: CrawlTopic(j.Priority),
Key: []byte(j.ID),
Value: data,
Headers: map[string]string{
"crawler": job.CrawlerName,
"priority": fmt.Sprintf("%d", int(job.Priority)),
"crawler": j.CrawlerName,
"priority": fmt.Sprintf("%d", int(j.Priority)),
},
}, nil
}
Expand Down
Loading
Loading