reenable queue new_persistence
parent
4fb26ec775
commit
e5e98e2890
|
|
@ -26,7 +26,6 @@ func NewModelToPersistencePipeline(ctx context.Context, cfg Config) (Pipeline, e
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return Pipeline{}, err
|
return Pipeline{}, err
|
||||||
}
|
}
|
||||||
writer = NewNoopQueue()
|
|
||||||
return Pipeline{
|
return Pipeline{
|
||||||
writer: writer,
|
writer: writer,
|
||||||
reader: reader,
|
reader: reader,
|
||||||
|
|
|
||||||
|
|
@ -23,15 +23,17 @@ type SlackScrape struct {
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewSlackScrapePipeline(ctx context.Context, cfg Config) (Pipeline, error) {
|
func NewSlackScrapePipeline(ctx context.Context, cfg Config) (Pipeline, error) {
|
||||||
reader, err := NewQueue(ctx, "slack_channels_to_scrape", cfg.driver)
|
writer, err := NewQueue(ctx, "new_persistence", cfg.driver)
|
||||||
|
if err != nil {
|
||||||
|
return Pipeline{}, err
|
||||||
|
}
|
||||||
|
cfg.slackScrapePipeline.reader, err = NewQueue(ctx, "slack_channels_to_scrape", cfg.driver)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return Pipeline{}, err
|
return Pipeline{}, err
|
||||||
}
|
}
|
||||||
cfg.slackScrapePipeline.reader = reader
|
|
||||||
writer := NewNoopQueue()
|
|
||||||
return Pipeline{
|
return Pipeline{
|
||||||
writer: writer,
|
writer: writer,
|
||||||
reader: reader,
|
reader: cfg.slackScrapePipeline.reader,
|
||||||
process: newSlackScrapeProcess(cfg),
|
process: newSlackScrapeProcess(cfg),
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue