cmd/ejobs: add -since flag to filter by date This change adds a (optional) -since flag to ejobs start command to filter modules by their version_updated_at timestamp before enqueuing tasks. This helps to skip analyzing older modules and speed up jobs. It also updates internal/pkgsitedb and internal/worker to support this. Change-Id: I13f1b3ae2373ceceb919581225a7de77e186deeb Reviewed-on: https://go-review.googlesource.com/c/pkgsite-metrics/+/774280 Reviewed-by: Jonathan Amsterdam <jba@google.com> LUCI-TryBot-Result: golang-scoped@luci-project-accounts.iam.gserviceaccount.com <golang-scoped@luci-project-accounts.iam.gserviceaccount.com> Reviewed-by: Ethan Lee <ethanalee@google.com> Auto-Submit: Hyang-Ah Hana Kim <hyangah@gmail.com>
diff --git a/cmd/ejobs/main.go b/cmd/ejobs/main.go index d017c29..332096a 100644 --- a/cmd/ejobs/main.go +++ b/cmd/ejobs/main.go
@@ -54,6 +54,7 @@ maxImporters int // for start noDeps bool // for start moduleFile string // for start + since string // for start waitInterval time.Duration // for wait outfile string // for results stream bool // for results @@ -75,7 +76,7 @@ {"cancel", "JOBID...", "cancel the jobs", doCancel, nil}, - {"start", "[-min MIN_IMPORTERS] [-file MODULE_FILE] [-nodeps] BINARY ARGS...", + {"start", "[-min MIN_IMPORTERS] [-file MODULE_FILE] [-nodeps] [-since DATE] BINARY ARGS...", "start a job", doStart, func(fs *flag.FlagSet) { @@ -86,6 +87,8 @@ fs.StringVar(&moduleFile, "file", "", "file with modules to use: each line is MODULE_PATH VERSION NUM_IMPORTERS") fs.BoolVar(&noDeps, "nodeps", false, "do not download dependencies for modules") + fs.StringVar(&since, "since", "", + "only analyze modules with version_updated_at >= this date (YYYY-MM-DD or RFC3339). Uses the timestamp when the module version was ingested by the system.") }, }, {"wait", "JOBID", @@ -383,6 +386,9 @@ if maxImporters >= 0 { u += fmt.Sprintf("&max=%d", maxImporters) } + if since != "" { + u += fmt.Sprintf("&since=%s", url.QueryEscape(since)) + } if gcsPath != "" { gurl := "gs://" + gcsPath u += fmt.Sprintf("&file=%s", url.QueryEscape(gurl))
diff --git a/internal/pkgsitedb/db.go b/internal/pkgsitedb/db.go index 4d90f79..27f270d 100644 --- a/internal/pkgsitedb/db.go +++ b/internal/pkgsitedb/db.go
@@ -13,6 +13,7 @@ "database/sql" "fmt" "regexp" + "time" _ "github.com/lib/pq" @@ -52,15 +53,27 @@ // ModuleSpecs retrieves all modules that contain packages that are // imported by minImportedByCount or more packages. // It looks for the information in the search_documents table of the given pkgsite DB. -func ModuleSpecs(ctx context.Context, db *sql.DB, minImports, maxImports int32) (specs []scan.ModuleSpec, err error) { +func ModuleSpecs(ctx context.Context, db *sql.DB, minImports, maxImports int32, since time.Time) (specs []scan.ModuleSpec, err error) { defer derrors.Wrap(&err, "moduleSpecsFromDB") query := ` SELECT module_path, version, max(imported_by_count) - FROM search_documents + FROM search_documents` + + var args []any + args = append(args, minImports, maxImports) + + if !since.IsZero() { + query += ` + WHERE version_updated_at >= $3` + args = append(args, since) + } + + query += ` GROUP BY module_path, version HAVING max(imported_by_count) >= $1 AND max(imported_by_count) <= $2 ORDER BY max(imported_by_count) desc` - rows, err := db.QueryContext(ctx, query, minImports, maxImports) + + rows, err := db.QueryContext(ctx, query, args...) if err != nil { return nil, err }
diff --git a/internal/pkgsitedb/db_plan9.go b/internal/pkgsitedb/db_plan9.go index b0a1dc7..752c476 100644 --- a/internal/pkgsitedb/db_plan9.go +++ b/internal/pkgsitedb/db_plan9.go
@@ -10,6 +10,7 @@ "context" "database/sql" "errors" + "time" "golang.org/x/pkgsite-metrics/internal/config" "golang.org/x/pkgsite-metrics/internal/scan" @@ -21,6 +22,6 @@ return nil, errDoesNotCompile } -func ModuleSpecs(ctx context.Context, db *sql.DB, minImportedByCount int) (specs []scan.ModuleSpec, err error) { +func ModuleSpecs(ctx context.Context, db *sql.DB, minImports, maxImports int32, since time.Time) (specs []scan.ModuleSpec, err error) { return nil, errDoesNotCompile }
diff --git a/internal/pkgsitedb/db_test.go b/internal/pkgsitedb/db_test.go index e4005f9..119326d 100644 --- a/internal/pkgsitedb/db_test.go +++ b/internal/pkgsitedb/db_test.go
@@ -16,6 +16,7 @@ "net/url" "strings" "testing" + "time" _ "github.com/lib/pq" ) @@ -51,7 +52,7 @@ if err := db.PingContext(ctx); err != nil { t.Fatal(err) } - got, err := ModuleSpecs(ctx, db, 1000, math.MaxInt32) + got, err := ModuleSpecs(ctx, db, 1000, math.MaxInt32, time.Time{}) if err != nil { t.Fatal(err) }
diff --git a/internal/worker/analysis.go b/internal/worker/analysis.go index e73bee7..d213548 100644 --- a/internal/worker/analysis.go +++ b/internal/worker/analysis.go
@@ -465,7 +465,20 @@ if err != nil { return err } - mods, err := readModules(ctx, s.cfg, params.File, params.Min, params.Max) + + var since time.Time + if s := r.FormValue("since"); s != "" { + var err error + since, err = time.Parse(time.RFC3339, s) + if err != nil { + since, err = time.Parse("2006-01-02", s) + if err != nil { + return fmt.Errorf("%w: analysis: invalid since time: %v", derrors.InvalidArgument, err) + } + } + } + + mods, err := readModules(ctx, s.cfg, params.File, params.Min, params.Max, since) if err != nil { return err }
diff --git a/internal/worker/enqueue.go b/internal/worker/enqueue.go index 599a880..6d9a615 100644 --- a/internal/worker/enqueue.go +++ b/internal/worker/enqueue.go
@@ -8,6 +8,7 @@ "context" "math" "sync" + "time" "golang.org/x/pkgsite-metrics/internal/config" "golang.org/x/pkgsite-metrics/internal/derrors" @@ -22,22 +23,22 @@ defaultMaxImportedByCount = math.MaxInt32 ) -func readModules(ctx context.Context, cfg *config.Config, file string, minImports, maxImports int32) ([]scan.ModuleSpec, error) { +func readModules(ctx context.Context, cfg *config.Config, file string, minImports, maxImports int32, since time.Time) ([]scan.ModuleSpec, error) { if file != "" { log.Infof(ctx, "reading modules from file %s", file) return scan.ParseCorpusFile(file, minImports, maxImports) } log.Infof(ctx, "reading modules from DB %s", cfg.PkgsiteDBName) - return readFromDB(ctx, cfg, minImports, maxImports) + return readFromDB(ctx, cfg, minImports, maxImports, since) } -func readFromDB(ctx context.Context, cfg *config.Config, minImports, maxImports int32) ([]scan.ModuleSpec, error) { +func readFromDB(ctx context.Context, cfg *config.Config, minImports, maxImports int32, since time.Time) ([]scan.ModuleSpec, error) { db, err := pkgsitedb.Open(ctx, cfg) if err != nil { return nil, err } defer db.Close() - return pkgsitedb.ModuleSpecs(ctx, db, minImports, maxImports) + return pkgsitedb.ModuleSpecs(ctx, db, minImports, maxImports, since) } func enqueueTasks(ctx context.Context, tasks []queue.Task, q queue.Queue, opts *queue.Options) (err error) {
diff --git a/internal/worker/govulncheck_enqueue.go b/internal/worker/govulncheck_enqueue.go index 2410882..bfb81ad 100644 --- a/internal/worker/govulncheck_enqueue.go +++ b/internal/worker/govulncheck_enqueue.go
@@ -12,6 +12,7 @@ "net/http" "sort" "strings" + "time" "golang.org/x/pkgsite-metrics/internal/config" "golang.org/x/pkgsite-metrics/internal/derrors" @@ -81,7 +82,7 @@ ) for _, mode := range modes { if modspecs == nil { - modspecs, err = readModules(ctx, cfg, params.File, params.Min, math.MaxInt32) + modspecs, err = readModules(ctx, cfg, params.File, params.Min, math.MaxInt32, time.Time{}) if err != nil { return nil, err }