all repos — searchix @ f0cee06cdb2390b873a9fc37b77e5e912d0412ae

Search engine for NixOS, nix-darwin, home-manager and NUR users

refactor: make importer responsible for spawning goroutines

Alan Pearce
commit

f0cee06cdb2390b873a9fc37b77e5e912d0412ae

parent

7971216b5bd8d6df7f4a27379b16b783c5ce0132

M internal/importer/importer.gointernal/importer/importer.go
@@ -9,29 +9,31 @@ "alin.ovh/searchix/internal/nix"
) type Processor interface { - Process(context.Context) (<-chan nix.Importable, <-chan error) + Process(context.Context, chan<- nix.Importable, chan<- error) } func (imp *Importer) process( ctx context.Context, processor Processor, -) (bool, error) { +) (hadObjectErrors bool, criticalError error) { wg := sync.WaitGroup{} - wg.Add(1) - objects, pErrs := processor.Process(ctx) + objects := make(chan nix.Importable, 1) + pErrs := make(chan error) + wg.Go(func() { + processor.Process(ctx, objects, pErrs) + }) - wg.Add(1) - iErrs := imp.options.WriteIndex.Import(ctx, objects) + iErrs := make(chan error) + wg.Go(func() { + imp.options.WriteIndex.Import(ctx, objects, iErrs) + }) - var hadObjectErrors bool - var criticalError error go func() { for { select { case err, running := <-iErrs: if !running { - wg.Done() iErrs = nil imp.options.Logger.Debug("ingest completed")
@@ -47,7 +49,6 @@ hadObjectErrors = true
imp.options.Logger.Warn("error ingesting object", "error", err) case err, running := <-pErrs: if !running { - wg.Done() pErrs = nil continue
M internal/importer/options.gointernal/importer/options.go
@@ -98,95 +98,92 @@
return &i, nil } -func (i *OptionIngester) Process(ctx context.Context) (<-chan nix.Importable, <-chan error) { - results := make(chan nix.Importable, 1) - errs := make(chan error) +func (i *OptionIngester) Process( + ctx context.Context, + results chan<- nix.Importable, + errs chan<- error, +) { + defer i.infile.Close() + defer close(results) + defer close(errs) - go func() { - defer i.infile.Close() - defer close(results) - defer close(errs) - - 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")) +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")) - continue - } - if mv.ValueType != jstream.Object { - 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) - var decls []*nix.Link - for _, decl := range x["declarations"].([]any) { - switch decl := reflect.ValueOf(decl); decl.Kind() { - case reflect.String: - s := decl.String() - 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, - )) - - continue - } - decls = append(decls, link) - case reflect.Map: - v := decl.Interface().(map[string]any) - link := nix.Link{ - Name: v["name"].(string), - URL: v["url"].(string), - } - decls = append(decls, &link) - default: - errs <- fault.Newf("unexpected declaration type %s", decl.Kind().String()) + var decls []*nix.Link + for _, decl := range x["declarations"].([]any) { + switch decl := reflect.ValueOf(decl); decl.Kind() { + case reflect.String: + s := decl.String() + 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, + )) continue } - } - if len(decls) > 0 { - x["declarations"] = decls - } - - i.optJSON = nixOptionJSON{} - err := i.ms.Decode(x) // stores in optJSON - if err != nil { - errs <- fault.Wrap(err, fmsg.Withf("failed to decode option %#v", x)) + decls = append(decls, link) + case reflect.Map: + v := decl.Interface().(map[string]any) + link := nix.Link{ + Name: v["name"].(string), + URL: v["url"].(string), + } + decls = append(decls, &link) + default: + errs <- fault.Newf("unexpected declaration type %s", decl.Kind().String()) continue } + } + if len(decls) > 0 { + x["declarations"] = decls + } - decs := make([]nix.Link, len(i.optJSON.Declarations)) - for i, d := range i.optJSON.Declarations { - decs[i] = nix.Link(d) - } + i.optJSON = nixOptionJSON{} + err := i.ms.Decode(x) // stores in optJSON + if err != nil { + errs <- fault.Wrap(err, fmsg.Withf("failed to decode option %#v", x)) + + continue + } - // log.Debug("sending option", "name", kv.Key) - results <- nix.Option{ - Name: kv.Key, - Source: i.source.Key, - Declarations: decs, - Default: i.convertDocsValue(i.optJSON.Default), - Description: nix.Markdown(i.optJSON.Description), - Example: i.convertDocsValue(i.optJSON.Example), - RelatedPackages: nix.Markdown(i.optJSON.RelatedPackages), - Loc: i.optJSON.Loc, - Parents: strings.Join(i.optJSON.Loc[:len(i.optJSON.Loc)-1], "."), - Type: i.optJSON.Type, - ImportedAt: time.Now(), - } + decs := make([]nix.Link, len(i.optJSON.Declarations)) + for i, d := range i.optJSON.Declarations { + decs[i] = nix.Link(d) } - }() - return results, errs + // log.Debug("sending option", "name", kv.Key) + results <- nix.Option{ + Name: kv.Key, + Source: i.source.Key, + Declarations: decs, + Default: i.convertDocsValue(i.optJSON.Default), + Description: nix.Markdown(i.optJSON.Description), + Example: i.convertDocsValue(i.optJSON.Example), + RelatedPackages: nix.Markdown(i.optJSON.RelatedPackages), + Loc: i.optJSON.Loc, + Parents: strings.Join(i.optJSON.Loc[:len(i.optJSON.Loc)-1], "."), + Type: i.optJSON.Type, + ImportedAt: time.Now(), + } + } }
M internal/importer/package.gointernal/importer/package.go
@@ -120,10 +120,11 @@
return l } -func (i *PackageIngester) Process(ctx context.Context) (<-chan nix.Importable, <-chan error) { - results := make(chan nix.Importable, 1) - errs := make(chan error) - +func (i *PackageIngester) Process( + ctx context.Context, + results chan<- nix.Importable, + errs chan<- error, +) { if i.programs != nil { err := i.programs.Open(ctx) if err != nil {
@@ -132,217 +133,213 @@ i.programs = nil
} } - go func() { - defer i.infile.Close() - defer close(results) - defer close(errs) + defer i.infile.Close() + defer close(results) + defer close(errs) + + if i.programs != nil { + defer i.programs.Close() + } - if i.programs != nil { - defer i.programs.Close() +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")) - 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 { + 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) - - 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() { - case reflect.String: - licenses[i] = makeAdhocLicense(v.String()) - case reflect.Map: - licenses[i] = *convertToLicense(v.Interface().(map[string]any)) - default: - 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: - errs <- fault.Newf( - "don't know how to handle license of type %s: %v", - v.Kind().String(), - meta["license"], - ) - } - delete(meta, "license") - } + meta := x["meta"].(map[string]any) - 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() { + 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() { case reflect.String: - plats[i] = v.String() + licenses[i] = makeAdhocLicense(v.String()) case reflect.Map: - plats[i] = 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...) + licenses[i] = *convertToLicense(v.Interface().(map[string]any)) default: errs <- fault.Newf( - "don't know how to convert platform type %s: %v", + "don't know how to handle sublicense of type %s: %v", v.Kind().String(), - v.Interface(), + v, ) } - i++ } - meta["platforms"] = plats + 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"], + ) } + delete(meta, "license") + } - if meta["homepage"] == nil && meta["homePage"] != nil { - meta["homepage"] = meta["homePage"] - } - if meta["homepage"] != nil { - switch v := reflect.ValueOf(meta["homepage"]); v.Kind() { + 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() { case reflect.String: - meta["homepage"] = []string{v.String()} + plats[i] = v.String() + case reflect.Map: + plats[i] = makeAdhocPlatform(v.Interface()) case reflect.Slice: - // already fine + 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 interpret homepage type %s'", + "don't know how to convert platform type %s: %v", v.Kind().String(), + v.Interface(), ) } + i++ } + meta["platforms"] = plats + } - 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) - } - 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, - ) + 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) + } + 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: - errs <- fault.Newf( - "don't know how to interpret maintainers type %s'", - maint.Kind().String(), - ) } - meta["maintainers"] = maints + default: + errs <- fault.Newf( + "don't know how to interpret maintainers type %s'", + maint.Kind().String(), + ) } + meta["maintainers"] = maints + } + + i.pkg = packageJSON{} + if err := i.ms.Decode(x); err != nil { // stores in i.pkg + errs <- fault.Wrap(err, fmsg.Withf("failed to decode package %#v", x)) - i.pkg = packageJSON{} - if err := i.ms.Decode(x); err != nil { // stores in i.pkg - errs <- fault.Wrap(err, fmsg.Withf("failed to decode package %#v", x)) + continue + } - continue + 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", i.pkg.Name)) } + } - 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", i.pkg.Name)) - } + maintainers := make([]nix.Maintainer, len(i.pkg.Meta.Maintainers)) + for i, m := range i.pkg.Meta.Maintainers { + maintainers[i] = nix.Maintainer{ + Name: m.Name, + Github: m.Github, } + } - maintainers := make([]nix.Maintainer, len(i.pkg.Meta.Maintainers)) - for i, m := range i.pkg.Meta.Maintainers { - maintainers[i] = nix.Maintainer{ - Name: m.Name, - Github: m.Github, - } - } + pkgSet, _, found := strings.Cut(kv.Key, ".") + if !found { + pkgSet = "" + } - pkgSet, _, found := strings.Cut(kv.Key, ".") - if !found { - pkgSet = "" + var definition string + if i.pkg.Meta.Position != "" { + defURL, err := url.Parse(i.pkg.Meta.Position) + if err != nil { + errs <- fault.Wrap(err, fmsg.Withf("failed to parse source URL %s", definition)) } - - var definition string - if i.pkg.Meta.Position != "" { - defURL, err := url.Parse(i.pkg.Meta.Position) + if defURL.IsAbs() { + definition = i.pkg.Meta.Position + } else { + subpath, line, _ := strings.Cut(i.pkg.Meta.Position, ":") + definition, err = i.source.Repo.GetFileURL(subpath, line) if err != nil { - errs <- fault.Wrap(err, fmsg.Withf("failed to parse source URL %s", definition)) - } - if defURL.IsAbs() { - definition = i.pkg.Meta.Position - } else { - subpath, line, _ := strings.Cut(i.pkg.Meta.Position, ":") - definition, err = i.source.Repo.GetFileURL(subpath, line) - if err != nil { - errs <- fault.Wrap(err, fmsg.Withf("failed to make repo URL for package %s", i.pkg.Name)) - } + errs <- fault.Wrap(err, fmsg.Withf("failed to make repo URL for package %s", i.pkg.Name)) } } + } - results <- nix.Package{ - Name: i.pkg.Name, - Attribute: strings.TrimPrefix(kv.Key, "nur.repos."), - Source: i.source.Key, - PackageSet: pkgSet, - Version: i.pkg.Version, - Broken: i.pkg.Meta.Broken, - Description: i.pkg.Meta.Description, - LongDescription: nix.Markdown(i.pkg.Meta.LongDescription), - Homepages: i.pkg.Meta.Homepages, - MainProgram: i.pkg.Meta.MainProgram, - Platforms: i.pkg.Meta.Platforms, - Licenses: licenses, - Maintainers: maintainers, - Definition: definition, - Programs: programs, - ImportedAt: time.Now(), - } + results <- nix.Package{ + Name: i.pkg.Name, + Attribute: strings.TrimPrefix(kv.Key, "nur.repos."), + Source: i.source.Key, + PackageSet: pkgSet, + Version: i.pkg.Version, + Broken: i.pkg.Meta.Broken, + Description: i.pkg.Meta.Description, + LongDescription: nix.Markdown(i.pkg.Meta.LongDescription), + Homepages: i.pkg.Meta.Homepages, + MainProgram: i.pkg.Meta.MainProgram, + Platforms: i.pkg.Meta.Platforms, + Licenses: licenses, + Maintainers: maintainers, + Definition: definition, + Programs: programs, + ImportedAt: time.Now(), } - }() - - return results, errs + } }
M internal/index/indexer.gointernal/index/indexer.go
@@ -318,10 +318,11 @@
func (i *WriteIndex) Import( ctx context.Context, objects <-chan nix.Importable, -) <-chan error { + errs chan<- error, +) { indexMapping := i.index.Mapping() - return i.WithBatchObjects(ctx, objects, func(batch *bleve.Batch, obj nix.Importable) error { + i.WithBatchObjects(ctx, objects, errs, func(batch *bleve.Batch, obj nix.Importable) error { doc := document.NewDocument(nix.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()))
@@ -369,49 +370,45 @@
func (i *WriteIndex) WithBatchObjects( ctx context.Context, objects <-chan nix.Importable, + errs chan<- error, processor func(batch *bleve.Batch, obj nix.Importable) error, -) <-chan error { +) { var err error - errs := make(chan error) - go func() { - defer close(errs) - k := 0 - batch := i.index.NewBatch() + defer close(errs) + k := 0 + batch := i.index.NewBatch() - outer: - for obj := range objects { - select { - case <-ctx.Done(): - i.log.Warn("batch process aborted") +outer: + for obj := range objects { + select { + case <-ctx.Done(): + i.log.Warn("batch process aborted") - break outer - default: - } + break outer + default: + } - if err := processor(batch, obj); err != nil { - errs <- fault.Wrap(err, fmsg.With("could not process object")) + if err := processor(batch, obj); err != nil { + errs <- fault.Wrap(err, fmsg.With("could not process object")) - continue - } + continue + } - if k++; k%i.batchSize == 0 { - err = i.Flush(batch) - if err != nil { - errs <- err + if k++; k%i.batchSize == 0 { + err = i.Flush(batch) + if err != nil { + errs <- err - return - } + return } } - - err := i.Flush(batch) - if err != nil { - errs <- err - } - }() + } - return errs + err = i.Flush(batch) + if err != nil { + errs <- err + } } func (i *WriteIndex) Flush(batch *bleve.Batch) error {