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
 			}