refactor: move interfaces closer to usage
1 file changed, 32 insertions(+), 7 deletions(-)
changed files
M internal/importer/main.go → internal/importer/main.go
@@ -7,6 +7,7 @@ "maps" "os/exec" "slices" "strings" + "sync" "time" "alin.ovh/x/log"@@ -22,6 +23,10 @@ "alin.ovh/searchix/internal/index/meta" "alin.ovh/searchix/internal/manpages" "alin.ovh/searchix/internal/programs" ) + +type Processor interface { + Process(context.Context, chan<- index.Indexable, chan<- error) +} type Options struct { LowMemory bool@@ -378,25 +383,45 @@ source.Repo.Revision, ) switch source.Importer { case config.Options: - processor, err = NewOptionProcessor( + processor = NewOptionProcessor( files.Options, source, logger.Named("processor"), ) case config.Packages: - processor, err = NewPackageProcessor( + processor = NewPackageProcessor( files.Packages, source, logger.Named("processor"), pdb, ) } - if err != nil { - return fault.Wrap(err, fmsg.Withf("failed to create processor")) - } - hadWarnings, err := imp.process(ctx, processor) - if err != nil { + var ( + hadWarnings bool + criticalError error + ) + wg := sync.WaitGroup{} + objects := make(chan index.Indexable, 1) + errs := make(chan error) + wg.Go(func() { + processor.Process(ctx, objects, errs) + close(objects) + }) + wg.Go(func() { + err := imp.options.WriteIndex.Import(ctx, objects, errs) + if err != nil { + criticalError = fault.Wrap(err, fmsg.With("error writing batch")) + } + close(errs) + }) + for err := range errs { + hadWarnings = true + imp.options.Logger.Warn("error processing object", "error", err) + } + wg.Wait() + imp.options.Logger.Debug("ingest completed") + if criticalError != nil { return fault.Wrap(err, fmsg.Withf("failed to process source")) }