배경
라이브 분석 + #380 (resolver chain 진화) 진행 도중 발견:
현재 Kafka 발행 책임이 5+ 모듈에 분산 — `scheduler.JobEmitter` (시드 entry), `publisher.Publisher` (chained URL), `fetcher/rule/upgrader` (auto-upgrade), `fetcher/domain/general/chain_handler` (chained job), `fetcher/worker/retry_scheduler` (retry republish), `fetcher/worker/{manager,pool}` (worker pool I/O). 이중 publisher 만 PriorityResolver 통과하고, 다른 모듈은 `entry.Priority` / hardcoded priority 직접 사용.
→ 동일 URL 이 시드와 chained 로 들어올 때 priority 결정 로직 다름. 새 resolver chain (#380) 의 의도가 시드 / retry / upgrade 경로에서 무효화.
또한 `fetcher/worker/{manager,pool}` 의 Kafka Consumer + Publisher 도 같은 자리에 있어 책임 혼재 — fetcher 의 worker pool 책임 (job dispatch, stage gate) 과 Kafka I/O 책임이 분리되지 않음.
목표
- Publisher 패키지를 Kafka 발행 단일 책임 hub 로 통합 — 모든 발행 경로가 동일 resolver chain + guard + producer 통과
- fetcher 쪽 publish/consume 모듈 이동 — Kafka I/O 는 publisher 책임, fetcher 는 fetch 핵심 로직만
- fetcher 로직이 강한 부분은 fetcher 책임 유지 — Kafka I/O 만 분리, worker pool 의 dispatch / stage gate / pool manager 등은 그대로
- 단일 facade — caller 는 `publisher.Publisher` 한 객체에 접근, 내부는 파일별 모듈화
구조 (옵션 A + C 결합)
패키지 구조
```
internal/publisher/
├── publisher.go # 외부 facade (PublishSeed / PublishChained / PublishRetry / PublishUpgrade)
├── seed.go # 구 scheduler.JobEmitter (시드 entry 발행)
├── chain.go # 구 publisher.Publisher (chained URL 발행)
├── upgrade.go # 구 fetcher/rule/upgrader (auto-upgrade republish)
├── retry.go # 구 fetcher/worker/retry_scheduler (Redis delayed retry)
├── resolver.go # PriorityResolver chain (현 worker/resolver.go 이동)
├── guard.go # PipelineGuard / IngestionLock wiring
├── consumer.go # Kafka consumer wrapper (구 fetcher/worker 의 consumer 부분)
└── producer.go # Kafka producer wrapper
```
외부 facade
```go
type Publisher struct { ... }
func (p *Publisher) PublishSeed(ctx, entry ScheduleEntry) error
func (p *Publisher) PublishChained(ctx, crawlerName string, urls []string, targetType core.TargetType, timeout time.Duration) error
func (p *Publisher) PublishRetry(ctx, job *core.CrawlJob, lastErr error) error
func (p *Publisher) PublishUpgrade(ctx, host string, raws []*core.RawContent) error
// Consumer 측
func (p *Publisher) Subscribe(topic string) Consumer
```
내부적으로 모든 메소드가 동일 resolver chain + guard + producer 공유.
fetcher 책임 유지
- `fetcher/worker/{pool, manager}` 의 dispatch 로직 / stage gate / circuit breaker — fetcher 책임
- Kafka I/O 만 publisher 로 위임 — `pool.Subscribe(consumer)` 가 아니라 `pool 이 publisher.Consumer 를 주입받음`
Sub-issues (6개 PR 분할)
| Sub |
작업 |
의존 |
| Sub 1 |
기존 `internal/publisher/publisher.go` 파일 분리 (chain.go) + 외부 facade Publisher 정의 |
없음 |
| Sub 2 |
`scheduler.JobEmitter` → `publisher/seed.go` 이동 + scheduler 가 publisher.PublishSeed 호출 |
Sub 1 |
| Sub 3 |
`fetcher/rule/upgrader` → `publisher/upgrade.go` 이동 + 호출처 갱신 |
Sub 1 |
| Sub 4 |
`fetcher/worker/retry_scheduler` → `publisher/retry.go` 이동 |
Sub 1 |
| Sub 5 |
Kafka Consumer 책임을 publisher 로 통합 — `fetcher/worker/{pool,manager}` 의 Kafka I/O 부분 분리, fetcher 로직은 그대로 |
Sub 1 |
| Sub 6 |
PriorityResolver chain 을 `worker/resolver.go` → `publisher/resolver.go` 이동 + 모든 PublishX 가 통과 |
Sub 1~5 |
각 Sub 는 독립 PR, reviewable diff.
완료 후 효과
영향 / 위험
- High — 대규모 refactoring (10+ 파일). Sub 분할로 review 부담 분산.
- 라이브 영향 없음 — 각 Sub 는 wiring 시점에 caller 교체, 동작 동일.
- 롤백 — 각 PR 단위로 revert 가능
완료 조건 (전체)
후속 이슈 (이미 등록)
배경
라이브 분석 + #380 (resolver chain 진화) 진행 도중 발견:
현재 Kafka 발행 책임이 5+ 모듈에 분산 — `scheduler.JobEmitter` (시드 entry), `publisher.Publisher` (chained URL), `fetcher/rule/upgrader` (auto-upgrade), `fetcher/domain/general/chain_handler` (chained job), `fetcher/worker/retry_scheduler` (retry republish), `fetcher/worker/{manager,pool}` (worker pool I/O). 이중 publisher 만 PriorityResolver 통과하고, 다른 모듈은 `entry.Priority` / hardcoded priority 직접 사용.
→ 동일 URL 이 시드와 chained 로 들어올 때 priority 결정 로직 다름. 새 resolver chain (#380) 의 의도가 시드 / retry / upgrade 경로에서 무효화.
또한 `fetcher/worker/{manager,pool}` 의 Kafka Consumer + Publisher 도 같은 자리에 있어 책임 혼재 — fetcher 의 worker pool 책임 (job dispatch, stage gate) 과 Kafka I/O 책임이 분리되지 않음.
목표
구조 (옵션 A + C 결합)
패키지 구조
```
internal/publisher/
├── publisher.go # 외부 facade (PublishSeed / PublishChained / PublishRetry / PublishUpgrade)
├── seed.go # 구 scheduler.JobEmitter (시드 entry 발행)
├── chain.go # 구 publisher.Publisher (chained URL 발행)
├── upgrade.go # 구 fetcher/rule/upgrader (auto-upgrade republish)
├── retry.go # 구 fetcher/worker/retry_scheduler (Redis delayed retry)
├── resolver.go # PriorityResolver chain (현 worker/resolver.go 이동)
├── guard.go # PipelineGuard / IngestionLock wiring
├── consumer.go # Kafka consumer wrapper (구 fetcher/worker 의 consumer 부분)
└── producer.go # Kafka producer wrapper
```
외부 facade
```go
type Publisher struct { ... }
func (p *Publisher) PublishSeed(ctx, entry ScheduleEntry) error
func (p *Publisher) PublishChained(ctx, crawlerName string, urls []string, targetType core.TargetType, timeout time.Duration) error
func (p *Publisher) PublishRetry(ctx, job *core.CrawlJob, lastErr error) error
func (p *Publisher) PublishUpgrade(ctx, host string, raws []*core.RawContent) error
// Consumer 측
func (p *Publisher) Subscribe(topic string) Consumer
```
내부적으로 모든 메소드가 동일 resolver chain + guard + producer 공유.
fetcher 책임 유지
Sub-issues (6개 PR 분할)
각 Sub 는 독립 PR, reviewable diff.
완료 후 효과
영향 / 위험
완료 조건 (전체)
후속 이슈 (이미 등록)