refactor: use iterator instead of channel for import
4 files changed, 359 insertions(+), 344 deletions(-)
M internal/importer/main.go → internal/importer/main.go
@@ -3,11 +3,11 @@ import ( "context" "errors" + "iter" "maps" "os/exec" "slices" "strings" - "sync" "time" "alin.ovh/x/log"@@ -25,7 +25,8 @@ "alin.ovh/searchix/internal/programs" ) type Processor interface { - Process(context.Context, chan<- index.Indexable, chan<- error) + Process(context.Context) iter.Seq[index.Indexable] + Err() error } type Options struct {@@ -397,33 +398,19 @@ pdb, ) } - 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 { + objects := processor.Process(ctx) + + err = imp.options.WriteIndex.Import(ctx, objects) + if err != nil { + return fault.Wrap(err, fmsg.With("error writing batch")) + } + + hadWarnings := false + if err := processor.Err(); err != nil { hadWarnings = true - imp.options.Logger.Warn("error processing object", "error", err) + imp.options.Logger.Warn("error processing objects", "error", err) } - wg.Wait() imp.options.Logger.Debug("ingest completed") - if criticalError != nil { - return fault.Wrap(err, fmsg.Withf("failed to process source")) - } sourceMeta.StoredAt = time.Now()
M internal/importer/options.go → internal/importer/options.go
@@ -3,6 +3,7 @@ import ( "context" "io" + "iter" "reflect" "strings" "time"@@ -55,6 +56,7 @@ dec *jstream.Decoder log *log.Logger infile io.ReadCloser source config.Source + errs []error } func NewOptionProcessor(@@ -72,112 +74,120 @@ } func (i *OptionIngester) Process( ctx context.Context, - results chan<- index.Indexable, - errs chan<- error, -) { - defer i.infile.Close() +) iter.Seq[index.Indexable] { + return func(yield func(index.Indexable) bool) { + defer i.infile.Close() -outer: - for mv := range i.dec.Stream() { - select { - case <-ctx.Done(): - break outer - default: - } - if err := i.dec.Err(); err != nil { - errs <- fault.Wrap(err, fmsg.With("could not decode JSON")) + for mv := range i.dec.Stream() { + if err := i.dec.Err(); err != nil { + i.errs = append(i.errs, fault.Wrap(err, fmsg.With("could not decode JSON"))) - continue - } - if mv.ValueType != jstream.Object { - errs <- fault.Newf("unexpected object type %s", ValueTypeToString(mv.ValueType)) + continue + } + if mv.ValueType != jstream.Object { + i.errs = append( + i.errs, + fault.Newf("unexpected object type %s", ValueTypeToString(mv.ValueType)), + ) - continue - } - kv := mv.Value.(jstream.KV) - x := kv.Value.(map[string]any) + continue + } + kv := mv.Value.(jstream.KV) + x := kv.Value.(map[string]any) - var decls []nix.Link - for _, decl := range x["declarations"].([]any) { - switch decl := decl.(type) { - case string: - s := decl - link, err := MakeChannelLink(i.source.Repo, s) - if err != nil { - errs <- fault.Wrap(err, fmsg.Withf("could not make a channel link for channel %s, revision %s and subpath %s", - i.source.Channel, i.source.Repo.Revision, s, - )) + var decls []nix.Link + for _, decl := range x["declarations"].([]any) { + switch decl := decl.(type) { + case string: + s := decl + link, err := MakeChannelLink(i.source.Repo, s) + if err != nil { + i.errs = append(i.errs, fault.Wrap(err, fmsg.Withf("could not make a channel link for channel %s, revision %s and subpath %s", + i.source.Channel, i.source.Repo.Revision, s, + ))) + + continue + } + decls = append(decls, *link) + case map[string]any: + v := decl + link := nix.Link{ + Name: v["name"].(string), + URL: v["url"].(string), + } + decls = append(decls, link) + default: + i.errs = append(i.errs, fault.Newf("unexpected declaration type %s", reflect.TypeOf(decl).String())) continue } - decls = append(decls, *link) - case map[string]any: - v := decl - link := nix.Link{ - Name: v["name"].(string), - URL: v["url"].(string), - } - decls = append(decls, link) - default: - errs <- fault.Newf("unexpected declaration type %s", reflect.TypeOf(decl).String()) - - continue + } + var description string + if v, ok := x["description"].(string); ok { + description = v } - } - var description string - if v, ok := x["description"].(string); ok { - description = v - } - var optType string - if v, ok := x["type"].(string); ok { - optType = v - } + var optType string + if v, ok := x["type"].(string); ok { + optType = v + } - var relatedPackages string - if v, ok := x["relatedPackages"].(string); ok { - relatedPackages = v - } + var relatedPackages string + if v, ok := x["relatedPackages"].(string); ok { + relatedPackages = v + } - var loc []string - if v, ok := x["loc"].([]any); ok { - loc = make([]string, len(v)) - for i, item := range v { - if s, ok := item.(string); ok { - loc[i] = s + var loc []string + if v, ok := x["loc"].([]any); ok { + loc = make([]string, len(v)) + for i, item := range v { + if s, ok := item.(string); ok { + loc[i] = s + } } } - } - var defaultDocs *nix.Docs - if v, ok := x["default"].(map[string]any); ok { - defaultDocs = i.convertDocsValue(v) - } + var defaultDocs *nix.Docs + if v, ok := x["default"].(map[string]any); ok { + defaultDocs = i.convertDocsValue(v) + } - var exampleDocs *nix.Docs - if v, ok := x["example"].(map[string]any); ok { - exampleDocs = i.convertDocsValue(v) - } + var exampleDocs *nix.Docs + if v, ok := x["example"].(map[string]any); ok { + exampleDocs = i.convertDocsValue(v) + } - // Calculate parents - var parents string - if len(loc) > 1 { - parents = strings.Join(loc[:len(loc)-1], ".") + var parents string + if len(loc) > 1 { + parents = strings.Join(loc[:len(loc)-1], ".") + } + + // log.Debug("sending option", "name", kv.Key) + if !yield(nix.Option{ + Name: kv.Key, + Source: i.source.Key, + Declarations: decls, + Default: defaultDocs, + Description: nixdocs.Markdown(description), + Example: exampleDocs, + RelatedPackages: nixdocs.Markdown(relatedPackages), + Loc: loc, + Parents: parents, + Type: optType, + ImportedAt: time.Now(), + }) { + return + } } + } +} - // log.Debug("sending option", "name", kv.Key) - results <- nix.Option{ - Name: kv.Key, - Source: i.source.Key, - Declarations: decls, - Default: defaultDocs, - Description: nixdocs.Markdown(description), - Example: exampleDocs, - RelatedPackages: nixdocs.Markdown(relatedPackages), - Loc: loc, - Parents: parents, - Type: optType, - ImportedAt: time.Now(), - } +func (i *OptionIngester) Err() error { + if len(i.errs) == 0 { + return nil + } + if len(i.errs) == 1 { + return i.errs[0] } + return fault.Newf("encountered %d errors during processing", len(i.errs)) }
M internal/importer/package.go → internal/importer/package.go
@@ -4,6 +4,7 @@ import ( "context" "encoding/json" "io" + "iter" "net/url" "reflect" "strings"@@ -28,6 +29,7 @@ log *log.Logger infile io.ReadCloser source config.Source programs *programs.DB + errs []error } func makeAdhocLicense(name string) nix.License {@@ -83,267 +85,286 @@ } func (i *PackageIngester) Process( ctx context.Context, - results chan<- index.Indexable, - errs chan<- error, -) { - if i.programs != nil { - err := i.programs.Open(ctx) - if err != nil { - errs <- fault.Wrap(err, fmsg.With("could not open programs database")) - i.programs = nil +) iter.Seq[index.Indexable] { + return func(yield func(index.Indexable) bool) { + defer i.infile.Close() + + if i.programs != nil { + err := i.programs.Open(ctx) + if err != nil { + i.errs = append(i.errs, fault.Wrap(err, fmsg.With("could not open programs database"))) + i.programs = nil + } } - } - defer i.infile.Close() + if i.programs != nil { + defer i.programs.Close() + } - if i.programs != nil { - defer i.programs.Close() - } + for mv := range i.dec.Stream() { + if ctx.Err() != nil { + break + } + if err := i.dec.Err(); err != nil { + i.errs = append(i.errs, fault.Wrap(err, fmsg.With("could not decode JSON"))) -outer: - for mv := range i.dec.Stream() { - var err error - var programs []string - select { - case <-ctx.Done(): - break outer - default: - } - if err := i.dec.Err(); err != nil { - errs <- fault.Wrap(err, fmsg.With("could not decode JSON")) + continue + } + if mv.ValueType != jstream.Object { + i.errs = append( + i.errs, + fault.Newf("unexpected object type %s", ValueTypeToString(mv.ValueType)), + ) - continue - } - if mv.ValueType != jstream.Object { - errs <- fault.Newf("unexpected object type %s", ValueTypeToString(mv.ValueType)) + continue + } + kv := mv.Value.(jstream.KV) + x := kv.Value.(map[string]any) - continue - } - kv := mv.Value.(jstream.KV) - x := kv.Value.(map[string]any) + meta := x["meta"].(map[string]any) - meta := x["meta"].(map[string]any) + var licenses []nix.License + if meta["license"] != nil { + switch v := reflect.ValueOf(meta["license"]); v.Kind() { + case reflect.Map: + licenses = append(licenses, *convertToLicense(v.Interface().(map[string]any))) + case reflect.Array, reflect.Slice: + licenses = make([]nix.License, v.Len()) + for idx, v := range v.Interface().([]any) { + switch v := reflect.ValueOf(v); v.Kind() { + case reflect.String: + licenses[idx] = makeAdhocLicense(v.String()) + case reflect.Map: + licenses[idx] = *convertToLicense(v.Interface().(map[string]any)) + default: + i.errs = append(i.errs, fault.Newf( + "don't know how to handle sublicense of type %s: %v", + v.Kind().String(), + v, + )) + } + } + case reflect.String: + licenses = append(licenses, makeAdhocLicense(v.String())) + default: + i.errs = append(i.errs, fault.Newf( + "don't know how to handle license of type %s: %v", + v.Kind().String(), + meta["license"], + )) + } + delete(meta, "license") + } - var licenses []nix.License - if meta["license"] != nil { - switch v := reflect.ValueOf(meta["license"]); v.Kind() { - case reflect.Map: - licenses = append(licenses, *convertToLicense(v.Interface().(map[string]any))) - case reflect.Array, reflect.Slice: - licenses = make([]nix.License, v.Len()) - for i, v := range v.Interface().([]any) { - switch v := reflect.ValueOf(v); v.Kind() { + if meta["platforms"] != nil { + plats := make([]any, len(meta["platforms"].([]any))) + idx := 0 + for _, plat := range meta["platforms"].([]any) { + switch v := reflect.ValueOf(plat); v.Kind() { case reflect.String: - licenses[i] = makeAdhocLicense(v.String()) + plats[idx] = v.String() case reflect.Map: - licenses[i] = *convertToLicense(v.Interface().(map[string]any)) + plats[idx] = makeAdhocPlatform(v.Interface()) + case reflect.Slice: + ps := make([]any, v.Len()) + for j, item := range v.Slice(0, v.Len()).Interface().([]any) { + ps[j] = item.(string) + } + plats = append(plats, ps...) default: - errs <- fault.Newf( - "don't know how to handle sublicense of type %s: %v", + i.errs = append(i.errs, fault.Newf( + "don't know how to convert platform type %s: %v", v.Kind().String(), - v, - ) + v.Interface(), + )) } + idx++ } - case reflect.String: - licenses = append(licenses, makeAdhocLicense(v.String())) - default: - errs <- fault.Newf( - "don't know how to handle license of type %s: %v", - v.Kind().String(), - meta["license"], - ) + meta["platforms"] = plats } - delete(meta, "license") - } - if meta["platforms"] != nil { - plats := make([]any, len(meta["platforms"].([]any))) - i := 0 - for _, plat := range meta["platforms"].([]any) { - switch v := reflect.ValueOf(plat); v.Kind() { + if meta["homepage"] == nil && meta["homePage"] != nil { + meta["homepage"] = meta["homePage"] + } + if meta["homepage"] != nil { + switch v := reflect.ValueOf(meta["homepage"]); v.Kind() { case reflect.String: - plats[i] = v.String() - case reflect.Map: - plats[i] = makeAdhocPlatform(v.Interface()) + meta["homepage"] = []string{v.String()} case reflect.Slice: - ps := make([]any, v.Len()) - for j, item := range v.Slice(0, v.Len()).Interface().([]any) { - ps[j] = item.(string) - } - plats = append(plats, ps...) + // already fine default: - errs <- fault.Newf( - "don't know how to convert platform type %s: %v", + i.errs = append(i.errs, fault.Newf( + "don't know how to interpret homepage type %s'", v.Kind().String(), - v.Interface(), - ) + )) } - i++ } - meta["platforms"] = plats - } - if meta["homepage"] == nil && meta["homePage"] != nil { - meta["homepage"] = meta["homePage"] - } - if meta["homepage"] != nil { - switch v := reflect.ValueOf(meta["homepage"]); v.Kind() { - case reflect.String: - meta["homepage"] = []string{v.String()} - case reflect.Slice: - // already fine - default: - errs <- fault.Newf( - "don't know how to interpret homepage type %s'", - v.Kind().String(), - ) - } - } - - var maints []nix.Maintainer - if meta["maintainers"] != nil { - switch maint := reflect.ValueOf(meta["maintainers"]); maint.Kind() { - case reflect.String: - maints = []nix.Maintainer{{Name: maint.String(), Github: maint.String()}} - case reflect.Slice, reflect.Array: - maints = make([]nix.Maintainer, maint.Len()) - for i, val := range maint.Slice(0, maint.Len()).Interface().([]any) { - switch v := reflect.ValueOf(val); v.Kind() { - case reflect.String: - maints[i] = nix.Maintainer{Name: v.String(), Github: v.String()} - case reflect.Map: - m := v.Interface().(map[string]any) - maints[i] = nix.Maintainer{} - if m["name"] != nil && m["name"].(string) != "" { - maints[i].Name = m["name"].(string) + var maints []nix.Maintainer + if meta["maintainers"] != nil { + switch maint := reflect.ValueOf(meta["maintainers"]); maint.Kind() { + case reflect.String: + maints = []nix.Maintainer{{Name: maint.String(), Github: maint.String()}} + case reflect.Slice, reflect.Array: + maints = make([]nix.Maintainer, maint.Len()) + for idx, val := range maint.Slice(0, maint.Len()).Interface().([]any) { + switch v := reflect.ValueOf(val); v.Kind() { + case reflect.String: + maints[idx] = nix.Maintainer{Name: v.String(), Github: v.String()} + case reflect.Map: + m := v.Interface().(map[string]any) + maints[idx] = nix.Maintainer{} + if m["name"] != nil && m["name"].(string) != "" { + maints[idx].Name = m["name"].(string) + } + if m["github"] != nil && m["github"].(string) != "" { + maints[idx].Github = m["github"].(string) + } + default: + i.errs = append(i.errs, fault.Newf( + "don't know how to handle maintainer entry of type %s: %v", + v.Kind().String(), + v, + )) } - if m["github"] != nil && m["github"].(string) != "" { - maints[i].Github = m["github"].(string) - } - default: - errs <- fault.Newf( - "don't know how to handle maintainer entry of type %s: %v", - v.Kind().String(), - v, - ) } + default: + i.errs = append(i.errs, fault.Newf( + "don't know how to interpret maintainers type %s'", + maint.Kind().String(), + )) } - default: - errs <- fault.Newf( - "don't know how to interpret maintainers type %s'", - maint.Kind().String(), - ) + meta["maintainers"] = maints } - meta["maintainers"] = maints - } - // Extract package name - var pkgName string - if pname, ok := x["pname"].(string); ok { - pkgName = pname - } + // Extract package name + var pkgName string + if pname, ok := x["pname"].(string); ok { + pkgName = pname + } - // Extract version - var version string - if v, ok := x["version"].(string); ok { - version = v - } + // Extract version + var version string + if v, ok := x["version"].(string); ok { + version = v + } - // Extract meta fields - var broken bool - if v, ok := meta["broken"].(bool); ok { - broken = v - } + // Extract meta fields + var broken bool + if v, ok := meta["broken"].(bool); ok { + broken = v + } - var description string - if v, ok := meta["description"].(string); ok { - description = v - } + var description string + if v, ok := meta["description"].(string); ok { + description = v + } - var longDescription string - if v, ok := meta["longDescription"].(string); ok { - longDescription = v - } + var longDescription string + if v, ok := meta["longDescription"].(string); ok { + longDescription = v + } - var homepages []string - if v, ok := meta["homepage"].([]string); ok { - homepages = v - } else if v, ok := meta["homepage"].([]any); ok { - homepages = make([]string, len(v)) - for i, h := range v { - if s, ok := h.(string); ok { - homepages[i] = s + var homepages []string + if v, ok := meta["homepage"].([]string); ok { + homepages = v + } else if v, ok := meta["homepage"].([]any); ok { + homepages = make([]string, len(v)) + for idx, h := range v { + if s, ok := h.(string); ok { + homepages[idx] = s + } } } - } - var mainProgram string - if v, ok := meta["mainProgram"].(string); ok { - mainProgram = v - } + var mainProgram string + if v, ok := meta["mainProgram"].(string); ok { + mainProgram = v + } - var platforms []string - if v, ok := meta["platforms"].([]any); ok { - platforms = make([]string, len(v)) - for i, p := range v { - if s, ok := p.(string); ok { - platforms[i] = s + var platforms []string + if v, ok := meta["platforms"].([]any); ok { + platforms = make([]string, len(v)) + for idx, p := range v { + if s, ok := p.(string); ok { + platforms[idx] = s + } } } - } - var position string - if v, ok := meta["position"].(string); ok { - position = v - } - - if i.source.Programs.Enable { - programs, err = i.programs.GetPackagePrograms(ctx, kv.Key) - if err != nil { - errs <- fault.Wrap(err, fmsg.Withf("failed to get programs for package %s", pkgName)) + var position string + if v, ok := meta["position"].(string); ok { + position = v } - } - pkgSet, _, found := strings.Cut(kv.Key, ".") - if !found { - pkgSet = "" - } + var programs []string + var err error + if i.source.Programs.Enable { + programs, err = i.programs.GetPackagePrograms(ctx, kv.Key) + if err != nil { + i.errs = append( + i.errs, + fault.Wrap(err, fmsg.Withf("failed to get programs for package %s", pkgName)), + ) + } + } - var definition string - if position != "" { - defURL, err := url.Parse(position) - if err != nil { - errs <- fault.Wrap(err, fmsg.Withf("failed to parse source URL %s", definition)) + pkgSet, _, found := strings.Cut(kv.Key, ".") + if !found { + pkgSet = "" } - if defURL.IsAbs() { - definition = position - } else { - subpath, line, _ := strings.Cut(position, ":") - definition, err = i.source.Repo.GetFileURL(subpath, line) + + var definition string + if position != "" { + defURL, err := url.Parse(position) if err != nil { - errs <- fault.Wrap(err, fmsg.Withf("failed to make repo URL for package %s", pkgName)) + i.errs = append( + i.errs, + fault.Wrap(err, fmsg.Withf("failed to parse source URL %s", definition)), + ) + } + if defURL.IsAbs() { + definition = position + } else { + subpath, line, _ := strings.Cut(position, ":") + definition, err = i.source.Repo.GetFileURL(subpath, line) + if err != nil { + i.errs = append(i.errs, fault.Wrap(err, fmsg.Withf("failed to make repo URL for package %s", pkgName))) + } } } + + if !yield(nix.Package{ + Name: pkgName, + Attribute: strings.TrimPrefix(kv.Key, "nur.repos."), + Source: i.source.Key, + PackageSet: pkgSet, + Version: version, + Broken: broken, + Description: description, + LongDescription: nixdocs.Markdown(longDescription), + Homepages: homepages, + MainProgram: mainProgram, + Platforms: platforms, + Licenses: licenses, + Maintainers: maints, + Definition: definition, + Programs: programs, + ImportedAt: time.Now(), + }) { + return + } } + } +} - results <- nix.Package{ - Name: pkgName, - Attribute: strings.TrimPrefix(kv.Key, "nur.repos."), - Source: i.source.Key, - PackageSet: pkgSet, - Version: version, - Broken: broken, - Description: description, - LongDescription: nixdocs.Markdown(longDescription), - Homepages: homepages, - MainProgram: mainProgram, - Platforms: platforms, - Licenses: licenses, - Maintainers: maints, - Definition: definition, - Programs: programs, - ImportedAt: time.Now(), - } +func (i *PackageIngester) Err() error { + if len(i.errs) == 0 { + return nil + } + if len(i.errs) == 1 { + return i.errs[0] } + + return fault.Newf("encountered %d errors during processing", len(i.errs)) }
M internal/index/indexer.go → internal/index/indexer.go
@@ -4,6 +4,7 @@ import ( "bytes" "context" "encoding/gob" + "iter" "math" "alin.ovh/searchix/internal/config"@@ -327,21 +328,24 @@ } func (i *WriteIndex) Import( ctx context.Context, - objects <-chan Indexable, - errs chan<- error, + objects iter.Seq[Indexable], ) error { indexMapping := i.index.Mapping() - return i.WithBatchObjects(ctx, objects, errs, func(batch *bleve.Batch, obj Indexable) error { + return i.WithBatchObjects(ctx, objects, func(batch *bleve.Batch, obj Indexable) { doc := document.NewDocument(GetKey(obj)) if err := indexMapping.MapDocument(doc, obj); err != nil { - return fault.Wrap(err, fmsg.Withf("could not map document for object: %s", obj.GetName())) + i.log.Warn("could not map document for object", "name", obj.GetName()) + + return } var data bytes.Buffer enc := gob.NewEncoder(&data) if err := enc.Encode(&obj); err != nil { - return fault.Wrap(err, fmsg.With("could not store object in search index")) + i.log.Error("could not store object in search index", "name", obj.GetName()) + + return } field := document.NewTextFieldWithIndexingOptions("_data", nil, data.Bytes(), index.StoreField) doc.AddField(field)@@ -354,10 +358,8 @@ doc.AddField(idField) // log.Debug("adding object to index", "name", opt.Name) if err := batch.IndexAdvanced(doc); err != nil { - return fault.Wrap(err, fmsg.Withf("could not index object %s", obj.GetName())) + i.log.Error("could not index object", "name", obj.GetName()) } - - return nil }) }@@ -379,9 +381,8 @@ } func (i *WriteIndex) WithBatchObjects( ctx context.Context, - objects <-chan Indexable, - errs chan<- error, - processor func(batch *bleve.Batch, obj Indexable) error, + objects iter.Seq[Indexable], + processor func(batch *bleve.Batch, obj Indexable), ) error { k := 0 batch := i.index.NewBatch()@@ -395,11 +396,7 @@ return ctx.Err() default: } - if err := processor(batch, obj); err != nil { - errs <- fault.Wrap(err, fmsg.With("could not process object")) - - continue - } + processor(batch, obj) if k++; k%i.batchSize == 0 { err := i.Flush(batch)