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
28 changes: 14 additions & 14 deletions cmd/issuetracker/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"time"

goredis "github.com/redis/go-redis/v9"
"issuetracker/internal/bus"
"issuetracker/internal/locks"
"issuetracker/internal/processor"
"issuetracker/internal/processor/fetcher"
Expand All @@ -25,7 +26,6 @@ import (
parserStage "issuetracker/internal/processor/parser/stage"
parserWorker "issuetracker/internal/processor/parser/worker"
"issuetracker/internal/processor/validate"
"issuetracker/internal/publisher"
"issuetracker/internal/scheduler"
"issuetracker/internal/storage"
pgstore "issuetracker/internal/storage/postgres"
Expand Down Expand Up @@ -108,10 +108,10 @@ func main() {
// 이슈 #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))
resolver := bus.NewCompositeResolver(core.PriorityNormal)
resolver.Add(&bus.ExplicitPriorityResolver{})
resolver.Add(bus.NewSourcePriorityResolver(core.PriorityNormal))
resolver.Add(bus.NewRuleBasedPriorityResolver(core.PriorityNormal))

highConsumer := queue.NewConsumer(crawlerKafkaCfg, queue.TopicCrawlHigh)
defer highConsumer.Close()
Expand All @@ -133,7 +133,7 @@ func main() {
}
defer pool.Close()

jobPublisher := publisher.New(crawlerProducer, resolver, log)
jobPublisher := bus.New(crawlerProducer, resolver, log)

// rule.Parser: parser_rules 테이블 기반 단일 파서 엔진.
// 사이트별 NaverParser/CNNParser/... 를 대체 — 모든 사이트가 본 단일 인스턴스를 공유.
Expand Down Expand Up @@ -278,7 +278,7 @@ func main() {
// Redis 초기화 실패 시에도 크롤링이 중단되지 않도록 graceful degrade 합니다.
var procLock locks.ProcessingLock
var ingestionLock locks.IngestionLock
var retryScheduler publisher.RetryScheduler
var retryScheduler bus.RetryScheduler
var retrySchedulerStop func()
var redisClientShared *redis.Client // failure counter wiring 에서 재사용
redisCfg, err := config.LoadRedis()
Expand All @@ -301,14 +301,14 @@ func main() {
// Delayed retry queue: retry 를 Redis ZSET 에 보관하고 별도
// goroutine 이 ScheduledAt 도달 시 Kafka 에 발행 — worker 슬롯 점유 회피.
// Redis 부재 시 worker 가 lazy 로 KafkaImmediateRetryScheduler 를 사용 (기존 동작).
retryCfg := publisher.DefaultRedisRetrySchedulerConfig()
retryCfg := bus.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 := publisher.NewRedisDelayedRetryScheduler(
redisRetry := bus.NewRedisDelayedRetryScheduler(
redisClient, jobPublisher,
retryCfg,
log,
Expand Down Expand Up @@ -530,7 +530,7 @@ func main() {
fetcherResolver,
rawIDTracker,
rawSvc,
jobPublisher, // 이슈 #388 — publisher.UpgradePublisher (단일 facade)
jobPublisher, // 이슈 #388 — bus.UpgradePublisher (단일 facade)
redisRaw, // nil 허용 — in-flight lock 비활성 (단일 인스턴스)
log,
)
Expand Down Expand Up @@ -569,7 +569,7 @@ func main() {

pw := parserWorker.NewParserWorker(
parserConsumer,
jobPublisher, // 이슈 #392 — 구 producer + JobPublisher 두 인자 통합 (publisher.Publisher 가 Forward + PublishChained 모두 제공)
jobPublisher, // 이슈 #392 — 구 producer + JobPublisher 두 인자 통합 (bus.Publisher 가 Forward + PublishChained 모두 제공)
rawSvc,
contentSvc,
ruleParser,
Expand All @@ -593,7 +593,7 @@ func main() {
}

// ── Page-parse 블랙리스트 ───────────────────────────────────────
// 카테고리에서 추출된 article URL 중 blacklist 매칭은 publisher.Publish 직전 drop.
// 카테고리에서 추출된 article URL 중 blacklist 매칭은 bus.Publish 직전 drop.
// Matcher 가 nil (BLACKLIST_ENABLED=false) 이면 setter noop — 모든 링크 그대로 발행.
if blacklistMatcher != nil {
pw.SetBlacklist(blacklistMatcher)
Expand Down Expand Up @@ -724,7 +724,7 @@ func main() {
log.WithError(err).Fatal("failed to load scheduler config")
}

// 이슈 #387 — scheduler.JobEmitter 제거. scheduler 가 publisher.Publisher 직접 의존.
// 이슈 #387 — scheduler.JobEmitter 제거. scheduler 가 bus.Publisher 직접 의존.
// guard / normalizer 는 jobPublisher 에 이미 wiring 되어 있어 시드 / chained 둘 다 동일
// 정책 적용 (메타 #385 — 단일 facade).
// Scheduler entries 는 두 source-of-truth 중 하나에서 결정 (이슈 #328):
Expand Down Expand Up @@ -795,7 +795,7 @@ func main() {

// 이슈 #393 — validate worker 가 publisher facade 의존. validate 전용 producer 를
// thin publisher 로 wrap (resolver/guard 불필요 — validate 는 Forward 만 사용).
validatePublisher := publisher.New(validateProducer, nil, log)
validatePublisher := bus.New(validateProducer, nil, log)

// Validate StageGate (이슈 #356) — ProcessingLock + per-stage Semaphore 합성.
validateCap := config.CapPerStage(workerCountsCfg.Validate, stageGateCfg.ValidateMaxConcurrentPerStage)
Expand Down
4 changes: 2 additions & 2 deletions cmd/processor/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,9 @@ import (
"os/signal"
"syscall"

"issuetracker/internal/bus"
"issuetracker/internal/locks"
"issuetracker/internal/processor/validate"
"issuetracker/internal/publisher"
pgstore "issuetracker/internal/storage/postgres"
"issuetracker/internal/storage/service"
"issuetracker/pkg/config"
Expand Down Expand Up @@ -97,7 +97,7 @@ func main() {

// 이슈 #393 — validate worker 가 publisher facade 의존으로 변경. processor 단독 실행은
// resolver / guard 가 필요 없으므로 nil resolver 로 thin publisher 만 생성.
pub := publisher.New(producer, nil, log)
pub := bus.New(producer, nil, log)

// ── 5. Validate Worker 시작 ───────────────────────────────────────────────
// validator 결과 (passed/rejected) 는 contentSvc.UpdateValidationStatus 로 contents 에 기록.
Expand Down
4 changes: 2 additions & 2 deletions examples/kafka_pipeline/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,9 @@ import (
"sync"
"time"

"issuetracker/internal/bus"
"issuetracker/internal/processor/fetcher/core"
"issuetracker/internal/processor/fetcher/worker"
"issuetracker/internal/publisher"
"issuetracker/internal/storage"
"issuetracker/internal/storage/service"
"issuetracker/pkg/logger"
Expand Down Expand Up @@ -368,7 +368,7 @@ func main() {
handler := &testCrawlerHandler{log: log}
contentSvc := newMockContentService()

pub := publisher.New(producer, nil, log)
pub := bus.New(producer, nil, log)
pool := worker.NewKafkaConsumerPool(consumer, pub, handler, contentSvc, workerCount)

start := time.Now()
Expand Down
4 changes: 2 additions & 2 deletions internal/publisher/chain.go → internal/bus/chain.go
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package publisher
package bus

import (
"context"
Expand Down Expand Up @@ -92,7 +92,7 @@ func (p *Publisher) PublishChained(
// buildJobMessages 는 url 목록을 CrawlJob 으로 변환하여 priority resolve 후 Kafka Message
// 슬라이스로 반환합니다 (CodeRabbit PR #394 피드백 — PublishChained 함수 분리).
//
// MaxRetries 는 publisher.DefaultMaxRetries 상수 사용.
// MaxRetries 는 bus.DefaultMaxRetries 상수 사용.
func (p *Publisher) buildJobMessages(
crawlerName string,
urls []string,
Expand Down
2 changes: 1 addition & 1 deletion internal/publisher/guard.go → internal/bus/guard.go
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package publisher
package bus

import (
"context"
Expand Down
10 changes: 5 additions & 5 deletions internal/publisher/resolver.go → internal/bus/resolver.go
Original file line number Diff line number Diff line change
@@ -1,10 +1,10 @@
// Package publisher 의 priority resolver 모듈 (이슈 #391, 메타 #385 Sub 6).
// Package bus 의 priority resolver 모듈 (이슈 #391, 메타 #385 Sub 6).
//
// 구 internal/processor/fetcher/worker/resolver.go 위치에서 이동 — Kafka I/O 단일 책임
// 원칙에 따라 PublishX 메소드의 priority routing 도 publisher 가 단일 출처.
//
// publisher.go 의 PriorityResolver 인터페이스를 본 파일의 chain/impl 들이 만족.
package publisher
// worker.go 의 PriorityResolver 인터페이스를 본 파일의 chain/impl 들이 만족.
package bus

import (
"issuetracker/internal/processor/fetcher/core"
Expand Down Expand Up @@ -196,8 +196,8 @@ func (r *DefaultPriorityResolver) CanResolve(_ *core.CrawlJob) bool {
//
// 사용 예:
//
// composite := publisher.NewCompositeResolver(core.PriorityNormal)
// composite.Add(&publisher.ExplicitPriorityResolver{}) // 1순위: 명시 priority 보존
// composite := bus.NewCompositeResolver(core.PriorityNormal)
// composite.Add(&bus.ExplicitPriorityResolver{}) // 1순위: 명시 priority 보존
// composite.Add(sourceResolver) // 2순위: 등록된 소스 매핑
// composite.Add(ruleResolver) // 3순위: 조건 규칙
// // 4순위 (자동): DefaultPriorityResolver → Normal
Expand Down
2 changes: 1 addition & 1 deletion internal/publisher/retry.go → internal/bus/retry.go
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package publisher
package bus

import (
"context"
Expand Down
6 changes: 3 additions & 3 deletions internal/publisher/seed.go → internal/bus/seed.go
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package publisher
package bus

import (
"context"
Expand All @@ -11,15 +11,15 @@ import (
// ErrPublishSkipped 는 PipelineGuard 가 "이미 in-pipeline" 으로 판단해 publish 를 건너뛴 경우
// PublishSeed 가 반환하는 sentinel error 입니다 (이슈 #387 — 구 scheduler.ErrEmitSkipped).
//
// 호출자는 errors.Is(err, publisher.ErrPublishSkipped) 로 분기하여 "failed to publish" /
// 호출자는 errors.Is(err, bus.ErrPublishSkipped) 로 분기하여 "failed to publish" /
// "scheduled" 로그를 모두 생략 — 실제로 발행되지 않은 job 이 발행된 것처럼 보이는 misleading
// 로그 회피.
var ErrPublishSkipped = errors.New("publish skipped — url already in pipeline")

// SeedPublisher 는 시드 CrawlJob 발행 책임을 정의하는 인터페이스입니다
// (이슈 #396 — 메타 #385 의 Kafka I/O 단일 책임 원칙에 따라 publisher 패키지에서 계약을 정의).
//
// publisher.Publisher 가 본 인터페이스를 만족하며, 외부 모듈은 본 인터페이스를 통해 시드
// bus.Publisher 가 본 인터페이스를 만족하며, 외부 모듈은 본 인터페이스를 통해 시드
// 발행 기능을 주입받아 사용합니다. 인터페이스 정의를 publisher 측에 두는 이유:
// - Kafka I/O 책임 = publisher 단일 진실 원천 — 시그니처 / 계약 변경 시 publisher 측만 갱신
// - 다른 sub (UpgradePublisher / RetryPublisher 등) 도 동일 원칙 적용 — 정합 일관
Expand Down
6 changes: 3 additions & 3 deletions internal/publisher/upgrade.go → internal/bus/upgrade.go
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package publisher
package bus

import (
"context"
Expand All @@ -11,7 +11,7 @@ import (
// 인터페이스입니다 (이슈 #388 — 메타 #385 의 Kafka I/O 단일 책임 원칙에 따라 publisher
// 패키지에서 계약을 정의 / 이슈 #396 의 원칙 적용).
//
// publisher.Publisher 가 본 인터페이스를 만족하며, 외부 모듈은 본 인터페이스를 통해 upgrade
// bus.Publisher 가 본 인터페이스를 만족하며, 외부 모듈은 본 인터페이스를 통해 upgrade
// 발행 기능을 주입받아 사용합니다.
//
// UpgradePublisher dispatches a batch of pre-built upgrade republish messages to Kafka.
Expand Down Expand Up @@ -43,7 +43,7 @@ func (p *Publisher) PublishUpgrade(ctx context.Context, host string, msgs []queu
if err := p.producer.PublishBatch(ctx, msgs); err != nil {
return fmt.Errorf("upgrade publish batch (host=%s, count=%d): %w", host, len(msgs), err)
}
// gemini PR #398 — defensive nil check (publisher.New 가 log 검증 안 하므로 caller 보호).
// gemini PR #398 — defensive nil check (bus.New 가 log 검증 안 하므로 caller 보호).
if p.log != nil {
p.log.WithFields(map[string]interface{}{
"host": host,
Expand Down
8 changes: 4 additions & 4 deletions internal/publisher/publisher.go → internal/bus/worker.go
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
// Package publisher 는 Kafka crawl 토픽 발행 책임을 단일 hub 로 통합합니다 (이슈 #385).
// Package bus 는 Kafka crawl 토픽 발행 책임을 단일 hub 로 통합합니다 (이슈 #385).
//
// 역할 (메타 #385 — Publisher 통합 모듈화):
// - PublishChained : 크롤된 페이지에서 발견된 URL 을 다음 CrawlJob 으로 연결 (chain.go)
Expand All @@ -7,14 +7,14 @@
// - PublishUpgrade : auto-upgrade (goquery → chromedp) republish (Sub 3 — pending)
//
// 외부 facade 단일화 — caller 는 *Publisher 의 메소드만 사용하면 됨. 내부 file 분리:
// - publisher.go : facade struct + 생성자 + 공통 Kafka helpers (buildMessage / CrawlTopic / newJobID)
// - worker.go : facade struct + 생성자 + 공통 Kafka helpers (buildMessage / CrawlTopic / newJobID)
// - chain.go : PublishChained 메소드 + 정규화 / guard / ingestion lock helper
// - guard.go : IngestionLock / PipelineGuard / atomic wrapper + Set* setters
//
// 의존 관계 (이슈 #385 책임 분리 원칙):
// - 본 패키지 = Kafka I/O + 라우팅 (priority resolver) + guard/lock 책임
// - caller 의 stage 핵심 로직 (parsing rule / validation / fetch decision) 은 본 패키지 의존성 없음
package publisher
package bus

import (
"context"
Expand Down Expand Up @@ -69,7 +69,7 @@ type PriorityResolver interface {
//
// 사용 흐름:
//
// pub := publisher.New(producer, resolver, log)
// pub := bus.New(producer, resolver, log)
// pub.SetNormalizer(...) // 선택
// pub.SetPipelineGuard(...) // 선택 (또는 SetIngestionLock fallback)
// pub.SetGate(...) // 선택
Expand Down
6 changes: 3 additions & 3 deletions internal/processor/fetcher/domain/general/chain_handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,9 +7,9 @@ import (
"fmt"
"net/url"

"issuetracker/internal/bus"
"issuetracker/internal/processor/fetcher/core"
"issuetracker/internal/processor/fetcher/rule"
"issuetracker/internal/publisher"
"issuetracker/internal/storage"
"issuetracker/internal/storage/service"
"issuetracker/pkg/logger"
Expand Down Expand Up @@ -49,7 +49,7 @@ type ChainHandler struct {
ChromedpChains []Handler // worker_id 별 browser only chain
Resolver rule.Resolver // optional. nil 이면 항상 DefaultChain 사용
RawSvc service.RawContentService
Pub *publisher.Publisher
Pub *bus.Publisher
Log *logger.Logger
}

Expand All @@ -70,7 +70,7 @@ func NewChainHandler(
chromedpChains []Handler,
resolver rule.Resolver,
rawSvc service.RawContentService,
pub *publisher.Publisher,
pub *bus.Publisher,
log *logger.Logger,
) *ChainHandler {
return &ChainHandler{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import (
"sort"
"time"

"issuetracker/internal/bus"
"issuetracker/internal/processor/fetcher/core"
"issuetracker/internal/processor/fetcher/domain/general"
"issuetracker/internal/processor/fetcher/domain/general/fetcher"
Expand All @@ -19,7 +20,6 @@ import (
"issuetracker/internal/processor/fetcher/implementation/goquery"
"issuetracker/internal/processor/fetcher/rate_limiter"
"issuetracker/internal/processor/fetcher/rule"
"issuetracker/internal/publisher"
"issuetracker/internal/storage"
"issuetracker/internal/storage/service"
"issuetracker/pkg/logger"
Expand Down Expand Up @@ -55,7 +55,7 @@ func RegisterAll(
fetcherRuleRepo storage.FetcherRuleRepository,
baseConfig core.Config,
rawSvc service.RawContentService,
pub *publisher.Publisher,
pub *bus.Publisher,
resolver rule.Resolver,
chromedpRemoteURLs []string,
log *logger.Logger,
Expand Down Expand Up @@ -138,7 +138,7 @@ func buildHandler(
rec *storage.FetcherRuleRecord,
baseConfig core.Config,
rawSvc service.RawContentService,
pub *publisher.Publisher,
pub *bus.Publisher,
resolver rule.Resolver,
chromedpRemoteURLs []string,
dnsResolver core.IPResolver,
Expand Down
2 changes: 1 addition & 1 deletion internal/processor/fetcher/domain/general/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ type Fetcher interface {
}

// JobPublisher 는 카테고리/목록 페이지에서 발견된 URL 을 다음 CrawlJob 으로 연결하는 인터페이스입니다.
// publisher.Publisher 가 본 인터페이스를 구현하며 (PublishChained 메소드 — 이슈 #386),
// bus.Publisher 가 본 인터페이스를 구현하며 (PublishChained 메소드 — 이슈 #386),
// ChainHandler 에 주입됩니다.
//
// JobPublisher dispatches CrawlJobs discovered from list/category pages.
Expand Down
4 changes: 2 additions & 2 deletions internal/processor/fetcher/domain/search/handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ import (

// JobPublisher 는 SearchHandler 가 fanout chained article job 을 발행할 때 사용하는 인터페이스입니다.
//
// 실제 구현은 internal/publisher.Publisher — host 별 batch 발행을 위해 host group 마다 호출.
// 실제 구현은 internal/bus.Publisher — host 별 batch 발행을 위해 host group 마다 호출.
// 이슈 #386 — Publisher.PublishChained 메소드명 일치 (구 Publish 에서 rename).
type JobPublisher interface {
PublishChained(
Expand All @@ -33,7 +33,7 @@ type JobPublisher interface {
// 1. job.Target.Metadata 에서 engine / per_query_max_results / date_range_days 추출
// 2. SearchKeywordRepository.ListEnabled 로 enabled keyword 전체 조회
// 3. 각 keyword 별 CSEClient.Search 호출 → URL 누적 (keyword-level 실패는 skip + warn)
// 4. URL 들을 host 단위로 group → 각 host 별 publisher.Publish(crawlerName=host, TargetTypeArticle)
// 4. URL 들을 host 단위로 group → 각 host 별 bus.PublishChained(crawlerName=host, TargetTypeArticle)
// fanout — 다운스트림 fetcher 가 host-specific handler 로 라우팅
// 5. 성공 keyword 의 last_searched_at 갱신
//
Expand Down
Loading
Loading