From f1da7f6f5da5bb47d297b5fa14d582085ba12440 Mon Sep 17 00:00:00 2001 From: dawn Date: Wed, 2 Sep 2026 00:16:58 +0900 Subject: [PATCH] knotmirror: index commit activity by month Signed-off-by: dawn --- api/tangled/temp2getCommitStats.go | 46 +++ knotmirror/db/commit_activity.go | 310 ++++++++++++++++ knotmirror/db/commit_activity_test.go | 223 ++++++++++++ knotmirror/db/migrations_list.go | 45 +++ knotmirror/knotmirror.go | 6 +- knotmirror/repoindexer/activity.go | 344 ++++++++++++++++++ knotmirror/repoindexer/activity_test.go | 242 ++++++++++++ knotmirror/resyncer.go | 36 +- knotmirror/xrpc/git_get_commit_stats.go | 184 ++++++++++ knotmirror/xrpc/git_get_commit_stats_test.go | 129 +++++++ knotmirror/xrpc/metrics.go | 5 + knotmirror/xrpc/xrpc.go | 54 +-- lexicons/git/temp2/getCommitStats.json | 73 ++++ localinfra/knotmirror.Dockerfile | 2 +- nix/modules/knotmirror.nix | 1 + web/src/lib/api/lexicons/index.ts | 2 +- .../sh/tangled/git/temp2/getCommitStats.ts | 73 ++++ 17 files changed, 1738 insertions(+), 37 deletions(-) create mode 100644 api/tangled/temp2getCommitStats.go create mode 100644 knotmirror/db/commit_activity.go create mode 100644 knotmirror/db/commit_activity_test.go create mode 100644 knotmirror/repoindexer/activity.go create mode 100644 knotmirror/repoindexer/activity_test.go create mode 100644 knotmirror/xrpc/git_get_commit_stats.go create mode 100644 knotmirror/xrpc/git_get_commit_stats_test.go create mode 100644 lexicons/git/temp2/getCommitStats.json create mode 100644 web/src/lib/api/lexicons/types/sh/tangled/git/temp2/getCommitStats.ts diff --git a/api/tangled/temp2getCommitStats.go b/api/tangled/temp2getCommitStats.go new file mode 100644 index 000000000..057f4ff93 --- /dev/null +++ b/api/tangled/temp2getCommitStats.go @@ -0,0 +1,46 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package tangled + +// schema: sh.tangled.git.temp2.getCommitStats + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +const ( + GitTemp2GetCommitStatsNSID = "sh.tangled.git.temp2.getCommitStats" +) + +// GitTemp2GetCommitStats_Month is a "month" in the sh.tangled.git.temp2.getCommitStats schema. +type GitTemp2GetCommitStats_Month struct { + // commits: Commits attributed to the actor during this month. + Commits int64 `json:"commits" cborgen:"commits"` + // start: Start of this UTC calendar month. + Start string `json:"start" cborgen:"start"` +} + +// GitTemp2GetCommitStats_Output is the output of a sh.tangled.git.temp2.getCommitStats call. +type GitTemp2GetCommitStats_Output struct { + Months []*GitTemp2GetCommitStats_Month `json:"months" cborgen:"months"` +} + +// GitTemp2GetCommitStats calls the XRPC method "sh.tangled.git.temp2.getCommitStats". +// +// actor: Actor whose commits to count. +func GitTemp2GetCommitStats(ctx context.Context, c util.LexClient, actor string, months int64) (*GitTemp2GetCommitStats_Output, error) { + var out GitTemp2GetCommitStats_Output + + params := map[string]interface{}{} + params["actor"] = actor + if months != 0 { + params["months"] = months + } + if err := c.LexDo(ctx, util.Query, "", "sh.tangled.git.temp2.getCommitStats", params, nil, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/knotmirror/db/commit_activity.go b/knotmirror/db/commit_activity.go new file mode 100644 index 000000000..a4108ccf3 --- /dev/null +++ b/knotmirror/db/commit_activity.go @@ -0,0 +1,310 @@ +package db + +import ( + "context" + "database/sql" + "encoding/json" + "fmt" + "time" +) + +type CommitActivityJob struct { + RepoDid string + RetryCount int +} + +type CommitActivity struct { + OID string `json:"oid"` + Committer string `json:"committer"` + Month string `json:"month"` +} + +type CommitRef struct { + Name string `json:"name"` + OID string `json:"oid"` +} + +type CommitCount struct { + Month time.Time + Commits int64 +} + +func EnqueueCommitActivity(ctx context.Context, e DBTX, repoDid string) error { + _, err := e.ExecContext(ctx, ` + insert into git_commit_jobs (repo_did) + values ($1) + on conflict (repo_did) do update set + dirty = git_commit_jobs.dirty or git_commit_jobs.state = 'processing', + state = case + when git_commit_jobs.state = 'processing' then 'processing' + else 'pending' + end, + retry_after = case + when git_commit_jobs.state = 'processing' then git_commit_jobs.retry_after + else now() + end, + requested_at = now()`, + repoDid, + ) + if err != nil { + return fmt.Errorf("enqueueing commit activity: %w", err) + } + return nil +} + +func ResetAbandonedCommitActivityJobs(ctx context.Context, e DBTX) error { + _, err := e.ExecContext(ctx, ` + update git_commit_jobs + set state = 'pending', dirty = false, claimed_at = null, retry_after = now() + where state = 'processing'`) + if err != nil { + return fmt.Errorf("resetting abandoned commit activity jobs: %w", err) + } + return nil +} + +func ClaimCommitActivityJob(ctx context.Context, e DBTX) (*CommitActivityJob, error) { + row := e.QueryRowContext(ctx, ` + with next as ( + select j.repo_did + from git_commit_jobs j + join repos r on r.repo_did = j.repo_did + where j.state = 'pending' + and j.retry_after <= now() + and r.state = 'active' + order by j.requested_at, j.repo_did + for update of j skip locked + limit 1 + ) + update git_commit_jobs j + set state = 'processing', dirty = false, claimed_at = now() + from next + where j.repo_did = next.repo_did + returning j.repo_did, j.retry_count`) + var job CommitActivityJob + if err := row.Scan(&job.RepoDid, &job.RetryCount); err != nil { + if err == sql.ErrNoRows { + return nil, nil + } + return nil, fmt.Errorf("claiming commit activity: %w", err) + } + return &job, nil +} + +func GetCommitActivityRefs(ctx context.Context, e DBTX, repoDid string) ([]CommitRef, error) { + rows, err := e.QueryContext(ctx, ` + select ref_name, oid + from git_commit_refs + where repo_did = $1 + order by ref_name`, + repoDid, + ) + if err != nil { + return nil, fmt.Errorf("listing commit activity refs: %w", err) + } + defer rows.Close() + + var refs []CommitRef + for rows.Next() { + var ref CommitRef + if err := rows.Scan(&ref.Name, &ref.OID); err != nil { + return nil, fmt.Errorf("scanning commit activity ref: %w", err) + } + refs = append(refs, ref) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("listing commit activity ref rows: %w", err) + } + return refs, nil +} + +func SaveCommitActivityBatch(ctx context.Context, e DBTX, commits []CommitActivity) (int64, error) { + if len(commits) == 0 { + return 0, nil + } + payload, err := json.Marshal(commits) + if err != nil { + return 0, fmt.Errorf("encoding commit activity: %w", err) + } + var inserted int64 + err = e.QueryRowContext(ctx, ` + with input as ( + select oid, committer, month + from jsonb_to_recordset($1::jsonb) + as row(oid text, committer text, month date) + ), inserted as ( + insert into git_seen_commits (oid, committer, committed_month) + select decode(oid, 'hex'), committer, month + from input + on conflict (oid) do nothing + returning committer, committed_month + ), grouped as ( + select committer, committed_month, count(*) as commits + from inserted + group by committer, committed_month + ), updated as ( + insert into git_commit_counts (committer, month, commits) + select committer, committed_month, commits + from grouped + on conflict (committer, month) do update set + commits = git_commit_counts.commits + excluded.commits, + updated_at = now() + returning commits + ) + select count(*) from inserted`, + payload, + ).Scan(&inserted) + if err != nil { + return 0, fmt.Errorf("saving commit activity batch: %w", err) + } + return inserted, nil +} + +func CompleteCommitActivityJob( + ctx context.Context, + e *sql.DB, + job CommitActivityJob, + refs []CommitRef, +) error { + payload, err := json.Marshal(refs) + if err != nil { + return fmt.Errorf("encoding commit activity refs: %w", err) + } + tx, err := e.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("beginning commit activity completion: %w", err) + } + defer tx.Rollback() + + var dirty bool + if err := tx.QueryRowContext(ctx, ` + select dirty + from git_commit_jobs + where repo_did = $1 and state = 'processing' + for update`, + job.RepoDid, + ).Scan(&dirty); err != nil { + return fmt.Errorf("locking commit activity job: %w", err) + } + + if _, err := tx.ExecContext(ctx, `delete from git_commit_refs where repo_did = $1`, job.RepoDid); err != nil { + return fmt.Errorf("clearing commit activity refs: %w", err) + } + if _, err := tx.ExecContext(ctx, ` + insert into git_commit_refs (repo_did, ref_name, oid) + select $1, name, oid + from jsonb_to_recordset($2::jsonb) as row(name text, oid text)`, + job.RepoDid, payload, + ); err != nil { + return fmt.Errorf("saving commit activity refs: %w", err) + } + if dirty { + if _, err := tx.ExecContext(ctx, ` + update git_commit_jobs + set state = 'pending', dirty = false, retry_count = 0, + retry_after = now(), claimed_at = null, error_msg = '' + where repo_did = $1 and state = 'processing'`, + job.RepoDid, + ); err != nil { + return fmt.Errorf("rescheduling dirty commit activity job: %w", err) + } + } else { + res, err := tx.ExecContext(ctx, ` + delete from git_commit_jobs + where repo_did = $1 and state = 'processing'`, + job.RepoDid, + ) + if err != nil { + return fmt.Errorf("finishing commit activity job: %w", err) + } + finished, err := res.RowsAffected() + if err != nil || finished != 1 { + return fmt.Errorf("finishing commit activity job: deleted %d rows: %w", finished, err) + } + } + if err := tx.Commit(); err != nil { + return fmt.Errorf("committing commit activity completion: %w", err) + } + return nil +} + +func RetryCommitActivityJob(ctx context.Context, e DBTX, job CommitActivityJob, retryAt time.Time, cause error) error { + message := "" + if cause != nil { + message = cause.Error() + } + if len(message) > 1000 { + message = message[:1000] + } + _, err := e.ExecContext(ctx, ` + update git_commit_jobs + set state = 'pending', dirty = false, claimed_at = null, + retry_count = retry_count + 1, retry_after = $2, error_msg = $3 + where repo_did = $1 and state = 'processing'`, + job.RepoDid, retryAt.UTC(), message, + ) + if err != nil { + return fmt.Errorf("retrying commit activity: %w", err) + } + return nil +} + +func PruneCommitActivity(ctx context.Context, e *sql.DB, before time.Time) error { + tx, err := e.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("beginning commit activity pruning: %w", err) + } + defer tx.Rollback() + if _, err := tx.ExecContext(ctx, `delete from git_seen_commits where committed_month < $1::date`, before.UTC().Format(time.DateOnly)); err != nil { + return fmt.Errorf("pruning seen commits: %w", err) + } + if _, err := tx.ExecContext(ctx, `delete from git_commit_counts where month < $1::date`, before.UTC().Format(time.DateOnly)); err != nil { + return fmt.Errorf("pruning commit counts: %w", err) + } + if err := tx.Commit(); err != nil { + return fmt.Errorf("committing commit activity pruning: %w", err) + } + return nil +} + +func ListCommitCounts( + ctx context.Context, + e DBTX, + committers []string, + since time.Time, + before time.Time, + limit int, +) ([]CommitCount, error) { + if limit < 1 || limit > 24 { + return nil, fmt.Errorf("listing commit counts with invalid limit: %d", limit) + } + if len(committers) < 1 || len(committers) > 251 { + return nil, fmt.Errorf("listing commit counts with invalid committer count: %d", len(committers)) + } + rows, err := e.QueryContext(ctx, ` + select month, sum(commits)::bigint + from git_commit_counts + where committer = any($1::text[]) and month >= $2::date and month < $3::date + group by month + order by month desc + limit $4`, + committers, since.UTC().Format(time.DateOnly), before.UTC().Format(time.DateOnly), limit, + ) + if err != nil { + return nil, fmt.Errorf("listing commit counts: %w", err) + } + defer rows.Close() + + counts := make([]CommitCount, 0, limit) + for rows.Next() { + var count CommitCount + if err := rows.Scan(&count.Month, &count.Commits); err != nil { + return nil, fmt.Errorf("scanning commit count: %w", err) + } + counts = append(counts, count) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("listing commit count rows: %w", err) + } + return counts, nil +} diff --git a/knotmirror/db/commit_activity_test.go b/knotmirror/db/commit_activity_test.go new file mode 100644 index 000000000..77b19c9f5 --- /dev/null +++ b/knotmirror/db/commit_activity_test.go @@ -0,0 +1,223 @@ +package db + +import ( + "context" + "database/sql" + "fmt" + "net/url" + "os" + "strings" + "testing" + "time" +) + +func commitActivityTestDB(tb testing.TB) *sql.DB { + tb.Helper() + databaseURL := os.Getenv("KNOTMIRROR_TEST_DATABASE_URL") + if databaseURL == "" { + tb.Skip("KNOTMIRROR_TEST_DATABASE_URL is not set") + } + ctx := context.Background() + admin, err := sql.Open("pgx", databaseURL) + if err != nil { + tb.Fatal(err) + } + schema := fmt.Sprintf("commit_activity_test_%d", time.Now().UnixNano()) + if _, err := admin.ExecContext(ctx, "create schema "+schema); err != nil { + admin.Close() + tb.Fatal(err) + } + parsed, err := url.Parse(databaseURL) + if err != nil { + tb.Fatal(err) + } + query := parsed.Query() + query.Set("search_path", schema) + query.Set("timezone", "Australia/Darwin") + parsed.RawQuery = query.Encode() + database, err := sql.Open("pgx", parsed.String()) + if err != nil { + tb.Fatal(err) + } + database.SetMaxOpenConns(4) + if err := database.PingContext(ctx); err != nil { + tb.Fatal(err) + } + tb.Cleanup(func() { + database.Close() + if _, err := admin.ExecContext(context.Background(), "drop schema "+schema+" cascade"); err != nil { + tb.Errorf("dropping test schema: %v", err) + } + admin.Close() + }) + if _, err := database.ExecContext(ctx, ` + create table repos ( + repo_did text primary key, + state text not null + )`); err != nil { + tb.Fatal(err) + } + return database +} + +func applyCommitActivityMigration(tb testing.TB, database *sql.DB) { + tb.Helper() + tx, err := database.BeginTx(context.Background(), nil) + if err != nil { + tb.Fatal(err) + } + defer tx.Rollback() + if err := gitCommitActivity(context.Background(), tx); err != nil { + tb.Fatal(err) + } + if err := tx.Commit(); err != nil { + tb.Fatal(err) + } +} + +func TestCommitActivityPostgresLifecycle(t *testing.T) { + ctx := context.Background() + database := commitActivityTestDB(t) + const repoDID = "did:plc:repo" + if _, err := database.ExecContext(ctx, `insert into repos (repo_did, state) values ($1, 'active')`, repoDID); err != nil { + t.Fatal(err) + } + applyCommitActivityMigration(t, database) + + job, err := ClaimCommitActivityJob(ctx, database) + if err != nil || job == nil || job.RepoDid != repoDID { + t.Fatalf("initial migration job = %#v, err = %v", job, err) + } + if err := EnqueueCommitActivity(ctx, database, repoDID); err != nil { + t.Fatal(err) + } + + current := time.Date(2026, 9, 1, 0, 0, 0, 0, time.UTC) + commits := []CommitActivity{ + {OID: strings.Repeat("a", 40), Committer: "did:plc:actor", Month: current.Format(time.DateOnly)}, + {OID: strings.Repeat("b", 40), Committer: "person@example.com", Month: current.Format(time.DateOnly)}, + {OID: strings.Repeat("c", 40), Committer: "person@example.com", Month: current.Format(time.DateOnly)}, + } + inserted, err := SaveCommitActivityBatch(ctx, database, commits) + if err != nil || inserted != 3 { + t.Fatalf("first batch inserted = %d, err = %v", inserted, err) + } + inserted, err = SaveCommitActivityBatch(ctx, database, commits) + if err != nil || inserted != 0 { + t.Fatalf("replayed batch inserted = %d, err = %v", inserted, err) + } + counts, err := ListCommitCounts( + ctx, + database, + []string{"did:plc:actor", "person@example.com"}, + current, + current.AddDate(0, 1, 0), + 24, + ) + if err != nil || len(counts) != 1 || counts[0].Commits != 3 { + t.Fatalf("aggregated counts = %#v, err = %v", counts, err) + } + + refs := []CommitRef{{Name: "refs/heads/main", OID: strings.Repeat("d", 40)}} + if err := CompleteCommitActivityJob(ctx, database, *job, refs); err != nil { + t.Fatal(err) + } + job, err = ClaimCommitActivityJob(ctx, database) + if err != nil || job == nil { + t.Fatalf("dirty replay job = %#v, err = %v", job, err) + } + storedRefs, err := GetCommitActivityRefs(ctx, database, repoDID) + if err != nil || len(storedRefs) != 1 || storedRefs[0] != refs[0] { + t.Fatalf("stored refs = %#v, err = %v", storedRefs, err) + } + if err := CompleteCommitActivityJob(ctx, database, *job, refs); err != nil { + t.Fatal(err) + } + job, err = ClaimCommitActivityJob(ctx, database) + if err != nil || job != nil { + t.Fatalf("finished job = %#v, err = %v", job, err) + } +} + +func TestCommitActivityPostgresPruningAndBounds(t *testing.T) { + ctx := context.Background() + database := commitActivityTestDB(t) + applyCommitActivityMigration(t, database) + current := time.Date(2026, 9, 1, 0, 0, 0, 0, time.UTC) + old := current.AddDate(0, -24, 0) + if _, err := SaveCommitActivityBatch(ctx, database, []CommitActivity{ + {OID: strings.Repeat("e", 40), Committer: "did:plc:actor", Month: old.Format(time.DateOnly)}, + {OID: strings.Repeat("f", 40), Committer: "did:plc:actor", Month: current.Format(time.DateOnly)}, + }); err != nil { + t.Fatal(err) + } + if err := PruneCommitActivity(ctx, database, current.AddDate(0, -23, 0)); err != nil { + t.Fatal(err) + } + var seen, counts int + if err := database.QueryRowContext(ctx, `select count(*) from git_seen_commits`).Scan(&seen); err != nil { + t.Fatal(err) + } + if err := database.QueryRowContext(ctx, `select count(*) from git_commit_counts`).Scan(&counts); err != nil { + t.Fatal(err) + } + if seen != 1 || counts != 1 { + t.Fatalf("after prune: seen = %d, counts = %d", seen, counts) + } + if _, err := ListCommitCounts(ctx, database, nil, old, current.AddDate(0, 1, 0), 24); err == nil { + t.Fatal("empty committer list was accepted") + } + if _, err := ListCommitCounts(ctx, database, []string{"did:plc:actor"}, old, current.AddDate(0, 1, 0), 25); err == nil { + t.Fatal("limit above API maximum was accepted") + } +} + +func BenchmarkListCommitCounts250kIdentities(b *testing.B) { + ctx := context.Background() + database := commitActivityTestDB(b) + applyCommitActivityMigration(b, database) + current := time.Date(2026, 9, 1, 0, 0, 0, 0, time.UTC) + if _, err := database.ExecContext(ctx, ` + insert into git_commit_counts (committer, month, commits) + select 'did:plc:actor' || actor, + $1::date - ((actor % 24) || ' months')::interval, + 1 + from generate_series(1, 250000) actor`, current); err != nil { + b.Fatal(err) + } + committers := make([]string, 0, 251) + committers = append(committers, "did:plc:target") + for index := range 250 { + committers = append(committers, fmt.Sprintf("person%03d@example.com", index)) + } + if _, err := database.ExecContext(ctx, ` + insert into git_commit_counts (committer, month, commits) + select committer, + $1::date - (month || ' months')::interval, + 1 + from unnest($2::text[]) committer + cross join generate_series(0, 23) month`, current, committers); err != nil { + b.Fatal(err) + } + var relationBytes int64 + if err := database.QueryRowContext(ctx, `select pg_total_relation_size('git_commit_counts')`).Scan(&relationBytes); err != nil { + b.Fatal(err) + } + b.ReportMetric(250000+251*24, "dataset_rows") + b.ReportMetric(float64(relationBytes), "dataset_bytes") + b.ReportAllocs() + b.ResetTimer() + for range b.N { + counts, err := ListCommitCounts( + ctx, + database, + committers, + current.AddDate(0, -23, 0), + current.AddDate(0, 1, 0), + 24, + ) + if err != nil || len(counts) != 24 { + b.Fatalf("counts = %d months, err = %v", len(counts), err) + } + } +} diff --git a/knotmirror/db/migrations_list.go b/knotmirror/db/migrations_list.go index 5016f9c12..d78b23c0c 100644 --- a/knotmirror/db/migrations_list.go +++ b/knotmirror/db/migrations_list.go @@ -17,6 +17,10 @@ var Migrations = []Migration{ Name: "add_last_feed_to_hosts", Fn: addLastFeedToHosts, }, + { + Name: "git_commit_activity", + Fn: gitCommitActivity, + }, } func addLastFeedToHosts(ctx context.Context, tx *sql.Tx) error { @@ -64,3 +68,44 @@ func execAll(ctx context.Context, tx *sql.Tx, stmts ...string) error { } return execAll(ctx, tx, stmts[1:]...) } + +func gitCommitActivity(ctx context.Context, tx *sql.Tx) error { + return execAll(ctx, tx, + `create table git_commit_jobs ( + repo_did text primary key references repos (repo_did) on delete cascade, + state text not null default 'pending' check (state in ('pending', 'processing')), + dirty boolean not null default false, + retry_count integer not null default 0, + retry_after timestamptz not null default now(), + requested_at timestamptz not null default now(), + claimed_at timestamptz, + error_msg text not null default '' + )`, + `create index git_commit_jobs_pending + on git_commit_jobs (retry_after, requested_at) + where state = 'pending'`, + `create table git_seen_commits ( + oid bytea primary key, + committer text not null, + committed_month date not null + )`, + `create index git_seen_commits_month + on git_seen_commits (committed_month)`, + `create table git_commit_counts ( + committer text not null, + month date not null, + commits bigint not null check (commits >= 0), + updated_at timestamptz not null default now(), + primary key (committer, month) + )`, + `create table git_commit_refs ( + repo_did text not null references repos (repo_did) on delete cascade, + ref_name text not null, + oid text not null, + primary key (repo_did, ref_name) + )`, + `insert into git_commit_jobs (repo_did) + select repo_did from repos + on conflict (repo_did) do nothing`, + ) +} diff --git a/knotmirror/knotmirror.go b/knotmirror/knotmirror.go index 188138310..2157f64ca 100644 --- a/knotmirror/knotmirror.go +++ b/knotmirror/knotmirror.go @@ -73,10 +73,14 @@ func Run(ctx context.Context, cfg *config.Config) error { indexer := repoindexer.NewIndexer(logger, cfg, rdb) indexScheduler := repoindexer.NewBackgroundIndexScheduler(logger, cfg, db, indexer) indexScheduler.Start(ctx) + activityScheduler := repoindexer.NewCommitActivityScheduler(logger, db, cfg.GitRepoBasePath) + if err := activityScheduler.Start(ctx); err != nil { + return fmt.Errorf("starting commit activity indexer: %w", err) + } knotstream := knotstream.NewKnotStream(logger, db, cfg) crawler := NewCrawler(logger, db) - resyncer := NewResyncer(logger, db, gitm, indexScheduler, cfg) + resyncer := NewResyncer(logger, db, gitm, indexScheduler, activityScheduler, cfg) xrpc := xrpc.New(logger, cfg, db, rdb, indexer, resolver, knotstream) xrpc.SetServiceSigner(serviceSigner) adminpage := NewAdminServer(logger, db, resyncer, xrpc, resolver) diff --git a/knotmirror/repoindexer/activity.go b/knotmirror/repoindexer/activity.go new file mode 100644 index 000000000..7ab466f62 --- /dev/null +++ b/knotmirror/repoindexer/activity.go @@ -0,0 +1,344 @@ +package repoindexer + +import ( + "bufio" + "bytes" + "context" + "database/sql" + "errors" + "fmt" + "io" + "log/slog" + "os/exec" + "path/filepath" + "regexp" + "strings" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/knotmirror/db" + "tangled.org/core/log" +) + +const ( + commitActivityPollInterval = time.Minute + commitActivityBatchPause = 100 * time.Millisecond + commitActivityBatchSize = 500 + commitActivityMonths = 24 +) + +var ( + commitObjectIDPattern = regexp.MustCompile(`^[0-9a-f]{40}([0-9a-f]{24})?$`) + errCommitWalk = errors.New("commit walk failed") +) + +type CommitActivityScheduler struct { + logger *slog.Logger + db *sql.DB + repoBase string + wake chan struct{} +} + +func NewCommitActivityScheduler(l *slog.Logger, database *sql.DB, repoBase string) *CommitActivityScheduler { + return &CommitActivityScheduler{ + logger: log.SubLogger(l, "commit-activity"), + db: database, + repoBase: repoBase, + wake: make(chan struct{}, 1), + } +} + +func (s *CommitActivityScheduler) Start(ctx context.Context) error { + if err := db.ResetAbandonedCommitActivityJobs(ctx, s.db); err != nil { + return err + } + go s.run(ctx) + return nil +} + +func (s *CommitActivityScheduler) Enqueue(ctx context.Context, repoDid syntax.DID) error { + if err := db.EnqueueCommitActivity(ctx, s.db, repoDid.String()); err != nil { + return err + } + select { + case s.wake <- struct{}{}: + default: + } + return nil +} + +func (s *CommitActivityScheduler) run(ctx context.Context) { + nextPrune := time.Now().Add(5 * time.Minute) + for { + s.drain(ctx) + if ctx.Err() == nil && !time.Now().Before(nextPrune) { + before := monthStart(time.Now().UTC()).AddDate(0, -commitActivityMonths, 0) + if err := db.PruneCommitActivity(ctx, s.db, before); err != nil { + s.logger.Warn("pruning commit activity failed", "err", err) + nextPrune = time.Now().Add(commitActivityPollInterval) + } else { + nextPrune = time.Now().Add(24 * time.Hour) + } + } + timer := time.NewTimer(commitActivityPollInterval) + select { + case <-ctx.Done(): + timer.Stop() + return + case <-s.wake: + timer.Stop() + case <-timer.C: + } + } +} + +func (s *CommitActivityScheduler) drain(ctx context.Context) { + for ctx.Err() == nil { + job, err := db.ClaimCommitActivityJob(ctx, s.db) + if err != nil { + s.logger.Error("claiming commit activity failed", "err", err) + return + } + if job == nil { + return + } + + start := time.Now() + refs, inserted, err := s.index(ctx, *job) + if err == nil { + err = db.CompleteCommitActivityJob(ctx, s.db, *job, refs) + } + if err != nil { + retryAt := time.Now().Add(commitActivityBackoff(job.RetryCount)) + if retryErr := db.RetryCommitActivityJob(ctx, s.db, *job, retryAt, err); retryErr != nil { + s.logger.Error("retrying commit activity failed", "repo", job.RepoDid, "err", retryErr) + return + } + s.logger.Warn("commit activity indexing deferred", "repo", job.RepoDid, "retry_at", retryAt, "err", err) + } else { + s.logger.Debug("indexed commit activity", "repo", job.RepoDid, "commits", inserted, "duration", time.Since(start)) + } + } +} + +func (s *CommitActivityScheduler) index(ctx context.Context, job db.CommitActivityJob) ([]db.CommitRef, int64, error) { + repoDid, err := syntax.ParseDID(job.RepoDid) + if err != nil { + return nil, 0, fmt.Errorf("invalid repo DID: %w", err) + } + repoPath := filepath.Join(s.repoBase, repoDid.String()) + refs, err := currentBranchRefs(ctx, repoPath) + if err != nil { + return nil, 0, err + } + previous, err := db.GetCommitActivityRefs(ctx, s.db, repoDid.String()) + if err != nil { + return nil, 0, err + } + currentMonth := monthStart(time.Now().UTC()) + cutoff := currentMonth.AddDate(0, -(commitActivityMonths - 1), 0) + before := currentMonth.AddDate(0, 1, 0) + + inserted, err := scanCommitActivity(ctx, repoPath, refs, previous, cutoff, before, func(batch []db.CommitActivity) (int64, error) { + count, err := db.SaveCommitActivityBatch(ctx, s.db, batch) + if err != nil { + return 0, err + } + timer := time.NewTimer(commitActivityBatchPause) + select { + case <-ctx.Done(): + timer.Stop() + return 0, ctx.Err() + case <-timer.C: + } + return count, nil + }) + return refs, inserted, err +} + +func commitActivityBackoff(retryCount int) time.Duration { + retryCount = min(max(retryCount, 0), 10) + return min(time.Second*time.Duration(1< 2_048 || parseErr != nil { + continue + } + at = at.UTC() + if at.Before(cutoff) || !at.Before(before) { + continue + } + batch = append(batch, db.CommitActivity{ + OID: oid, + Committer: identity, + Month: monthStart(at).Format(time.DateOnly), + }) + if len(batch) == cap(batch) { + count, saveErr := save(batch) + if saveErr != nil { + return inserted, saveErr + } + inserted += count + batch = batch[:0] + } + } + waitErr := cmd.Wait() + waited = true + if waitErr != nil { + return inserted, fmt.Errorf("%w: %v: %s", errCommitWalk, waitErr, strings.TrimSpace(stderr.String())) + } + if len(batch) > 0 { + count, err := save(batch) + if err != nil { + return inserted, err + } + inserted += count + } + return inserted, nil +} + +func readNullField(r *bufio.Reader) (string, error) { + field, err := r.ReadString(0) + return strings.TrimSuffix(field, "\x00"), err +} + +func normalizeCommitter(value string) string { + value = strings.TrimSpace(value) + if did, err := syntax.ParseDID(value); err == nil { + return did.String() + } + return strings.ToLower(value) +} + +func monthStart(value time.Time) time.Time { + value = value.UTC() + return time.Date(value.Year(), value.Month(), 1, 0, 0, 0, 0, time.UTC) +} + +func lowestPriorityGitCommand(ctx context.Context, repoPath string, args ...string) *exec.Cmd { + gitArgs := append([]string{"-C", repoPath}, args...) + if _, err := exec.LookPath("ionice"); err == nil { + return exec.CommandContext(ctx, "ionice", append([]string{"-c", "3", "nice", "-n", "19", "git"}, gitArgs...)...) + } + return exec.CommandContext(ctx, "nice", append([]string{"-n", "19", "git"}, gitArgs...)...) +} + +func runLowestPriorityGit(ctx context.Context, repoPath string, args ...string) ([]byte, error) { + cmd := lowestPriorityGitCommand(ctx, repoPath, args...) + out, err := cmd.CombinedOutput() + if err != nil { + return nil, fmt.Errorf("%s: %w", strings.TrimSpace(string(out)), err) + } + return out, nil +} diff --git a/knotmirror/repoindexer/activity_test.go b/knotmirror/repoindexer/activity_test.go new file mode 100644 index 000000000..f0615d002 --- /dev/null +++ b/knotmirror/repoindexer/activity_test.go @@ -0,0 +1,242 @@ +package repoindexer + +import ( + "bytes" + "context" + "fmt" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" + "time" + + "tangled.org/core/knotmirror/db" +) + +func activityGit(t testing.TB, dir string, env map[string]string, args ...string) { + t.Helper() + cmd := exec.Command("git", append([]string{"-C", dir}, args...)...) + cmd.Env = os.Environ() + for key, value := range env { + cmd.Env = append(cmd.Env, key+"="+value) + } + if out, err := cmd.CombinedOutput(); err != nil { + t.Fatalf("git %v: %v: %s", args, err, out) + } +} + +func activityCommit(t testing.TB, dir, filename, committer string, at time.Time) { + t.Helper() + if err := os.WriteFile(filepath.Join(dir, filename), []byte(filename), 0o600); err != nil { + t.Fatal(err) + } + activityGit(t, dir, nil, "add", filename) + date := at.Format(time.RFC3339) + activityGit(t, dir, map[string]string{ + "GIT_AUTHOR_NAME": "author", + "GIT_AUTHOR_EMAIL": "not-the-committer@example.com", + "GIT_AUTHOR_DATE": date, + "GIT_COMMITTER_NAME": "committer", + "GIT_COMMITTER_EMAIL": committer, + "GIT_COMMITTER_DATE": date, + }, "commit", "-q", "-m", filename) +} + +func activityFixture(t testing.TB) (string, time.Time) { + t.Helper() + dir := t.TempDir() + activityGit(t, dir, nil, "init", "-q", "-b", "main") + current := monthStart(time.Now().UTC()) + activityCommit(t, dir, "base", "did:plc:alice", current.AddDate(0, -2, 1)) + activityGit(t, dir, nil, "branch", "feature") + activityCommit(t, dir, "main", "did:plc:alice", current.AddDate(0, 0, 1)) + activityGit(t, dir, nil, "checkout", "-q", "feature") + activityCommit(t, dir, "feature", "did:plc:bob", current.AddDate(0, -1, 1)) + activityGit(t, dir, nil, "checkout", "-q", "main") + return dir, current +} + +func TestScanCommitActivityUsesCommittersAndUTCMonths(t *testing.T) { + dir, current := activityFixture(t) + refs, err := currentBranchRefs(context.Background(), dir) + if err != nil { + t.Fatal(err) + } + var got []db.CommitActivity + inserted, err := scanCommitActivity( + context.Background(), + dir, + refs, + nil, + current.AddDate(0, -23, 0), + current.AddDate(0, 1, 0), + func(batch []db.CommitActivity) (int64, error) { + got = append(got, batch...) + return int64(len(batch)), nil + }, + ) + if err != nil { + t.Fatal(err) + } + if inserted != 3 || len(got) != 3 { + t.Fatalf("inserted %d commits: %#v", inserted, got) + } + counts := make(map[string]map[string]int) + for _, commit := range got { + if counts[commit.Committer] == nil { + counts[commit.Committer] = make(map[string]int) + } + counts[commit.Committer][commit.Month]++ + } + if counts["did:plc:alice"][current.Format(time.DateOnly)] != 1 { + t.Fatalf("alice current month counts: %#v", counts) + } + if counts["did:plc:alice"][current.AddDate(0, -2, 0).Format(time.DateOnly)] != 1 { + t.Fatalf("alice old month counts: %#v", counts) + } + if counts["did:plc:bob"][current.AddDate(0, -1, 0).Format(time.DateOnly)] != 1 { + t.Fatalf("bob counts: %#v", counts) + } + if _, ok := counts["not-the-committer@example.com"]; ok { + t.Fatalf("author was counted instead of committer: %#v", counts) + } +} + +func TestScanCommitActivityExcludesPreviouslyMirroredReachability(t *testing.T) { + dir, current := activityFixture(t) + previous, err := currentBranchRefs(context.Background(), dir) + if err != nil { + t.Fatal(err) + } + activityCommit(t, dir, "new", "did:plc:carol", current.AddDate(0, 0, 2)) + refs, err := currentBranchRefs(context.Background(), dir) + if err != nil { + t.Fatal(err) + } + var got []db.CommitActivity + inserted, err := scanCommitActivity( + context.Background(), + dir, + refs, + previous, + current.AddDate(0, -23, 0), + current.AddDate(0, 1, 0), + func(batch []db.CommitActivity) (int64, error) { + got = append(got, batch...) + return int64(len(batch)), nil + }, + ) + if err != nil { + t.Fatal(err) + } + if inserted != 1 || len(got) != 1 || got[0].Committer != "did:plc:carol" { + t.Fatalf("delta = %#v, inserted = %d", got, inserted) + } +} + +func TestScanCommitActivityFallsBackWhenPreviousObjectsAreGone(t *testing.T) { + dir, current := activityFixture(t) + refs, err := currentBranchRefs(context.Background(), dir) + if err != nil { + t.Fatal(err) + } + var got int64 + inserted, err := scanCommitActivity( + context.Background(), dir, refs, + []db.CommitRef{{Name: "refs/heads/old", OID: strings.Repeat("f", 40)}}, + current.AddDate(0, -23, 0), current.AddDate(0, 1, 0), + func(batch []db.CommitActivity) (int64, error) { + got += int64(len(batch)) + return int64(len(batch)), nil + }, + ) + if err != nil { + t.Fatal(err) + } + if inserted != 3 || got != 3 { + t.Fatalf("fallback inserted %d, saved %d", inserted, got) + } +} + +func TestScanCommitActivityDeduplicatesBranchesAndRetries(t *testing.T) { + dir, current := activityFixture(t) + activityGit(t, dir, nil, "branch", "same-as-main", "main") + refs, err := currentBranchRefs(context.Background(), dir) + if err != nil { + t.Fatal(err) + } + seen := make(map[string]struct{}) + save := func(batch []db.CommitActivity) (int64, error) { + var inserted int64 + for _, commit := range batch { + if _, ok := seen[commit.OID]; ok { + continue + } + seen[commit.OID] = struct{}{} + inserted++ + } + return inserted, nil + } + for run, want := range []int64{3, 0} { + got, err := scanCommitActivity( + context.Background(), dir, refs, nil, + current.AddDate(0, -23, 0), current.AddDate(0, 1, 0), save, + ) + if err != nil { + t.Fatal(err) + } + if got != want { + t.Fatalf("run %d inserted %d, want %d", run, got, want) + } + } +} + +func fastImportActivityRepo(b *testing.B, commits int) string { + b.Helper() + dir := b.TempDir() + activityGit(b, dir, nil, "init", "-q", "--bare") + var stream bytes.Buffer + stream.WriteString("blob\nmark :1\ndata 1\nx\n") + base := time.Now().UTC().Add(-30 * 24 * time.Hour).Unix() + for i := range commits { + fmt.Fprintf(&stream, "commit refs/heads/main\nmark :%d\n", i+2) + fmt.Fprintf(&stream, "author author %d +0000\n", base+int64(i)) + fmt.Fprintf(&stream, "committer committer %d +0000\n", i%100, base+int64(i)) + stream.WriteString("data 1\nx\n") + if i > 0 { + fmt.Fprintf(&stream, "from :%d\n", i+1) + } + stream.WriteString("M 100644 :1 file\n\n") + } + cmd := exec.Command("git", "-C", dir, "fast-import", "--quiet") + cmd.Stdin = &stream + if out, err := cmd.CombinedOutput(); err != nil { + b.Fatalf("fast-import: %v: %s", err, out) + } + return dir +} + +func BenchmarkScanCommitActivity10k(b *testing.B) { + dir := fastImportActivityRepo(b, 10_000) + refs, err := currentBranchRefs(context.Background(), dir) + if err != nil { + b.Fatal(err) + } + current := monthStart(time.Now().UTC()) + b.ReportAllocs() + b.ResetTimer() + for range b.N { + count, err := scanCommitActivity( + context.Background(), dir, refs, nil, + current.AddDate(0, -23, 0), current.AddDate(0, 1, 0), + func(batch []db.CommitActivity) (int64, error) { return int64(len(batch)), nil }, + ) + if err != nil { + b.Fatal(err) + } + if count != 10_000 { + b.Fatalf("got %d commits", count) + } + } +} diff --git a/knotmirror/resyncer.go b/knotmirror/resyncer.go index 449a6faf6..e2045ab6d 100644 --- a/knotmirror/resyncer.go +++ b/knotmirror/resyncer.go @@ -21,15 +21,17 @@ import ( "tangled.org/core/knotmirror/db" "tangled.org/core/knotmirror/knotstream" "tangled.org/core/knotmirror/models" + "tangled.org/core/knotmirror/repoindexer" "tangled.org/core/log" ) type Resyncer struct { - logger *slog.Logger - db *sql.DB - gitm GitMirrorManager - cfg *config.Config - indexer *knotstream.ParallelScheduler + logger *slog.Logger + db *sql.DB + gitm GitMirrorManager + cfg *config.Config + indexer *knotstream.ParallelScheduler + activityIndexer *repoindexer.CommitActivityScheduler claimJobMu sync.Mutex @@ -46,13 +48,21 @@ type Resyncer struct { httpClient *http.Client } -func NewResyncer(l *slog.Logger, db *sql.DB, gitm GitMirrorManager, indexer *knotstream.ParallelScheduler, cfg *config.Config) *Resyncer { +func NewResyncer( + l *slog.Logger, + db *sql.DB, + gitm GitMirrorManager, + indexer *knotstream.ParallelScheduler, + activityIndexer *repoindexer.CommitActivityScheduler, + cfg *config.Config, +) *Resyncer { return &Resyncer{ - logger: log.SubLogger(l, "resyncer"), - db: db, - gitm: gitm, - cfg: cfg, - indexer: indexer, + logger: log.SubLogger(l, "resyncer"), + db: db, + gitm: gitm, + cfg: cfg, + indexer: indexer, + activityIndexer: activityIndexer, runningJobs: make(map[syntax.DID]context.CancelFunc), @@ -280,6 +290,10 @@ func (r *Resyncer) doResync(ctx context.Context, repoDid syntax.DID) (bool, erro // queue repo_stats_update job r.indexer.AddTask(context.TODO(), &knotstream.Task{Key: repo.RepoDid.String()}) + if err := r.activityIndexer.Enqueue(ctx, repo.RepoDid); err != nil { + return false, fmt.Errorf("queueing commit activity: %w", err) + } + // repo.GitRev = // repo.RepoSha = repo.State = models.RepoStateActive diff --git a/knotmirror/xrpc/git_get_commit_stats.go b/knotmirror/xrpc/git_get_commit_stats.go new file mode 100644 index 000000000..c31dcc570 --- /dev/null +++ b/knotmirror/xrpc/git_get_commit_stats.go @@ -0,0 +1,184 @@ +package xrpc + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/url" + "strconv" + "strings" + "time" + + "github.com/bluesky-social/indigo/atproto/atclient" + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" + "tangled.org/core/knotmirror/db" +) + +const ( + defaultCommitStatsMonths = 7 + maxCommitStatsMonths = 24 + committerCacheCapacity = 4_096 + committerCacheFreshness = 30 * time.Second +) + +func (x *Xrpc) GetCommitStats(w http.ResponseWriter, r *http.Request) { + actor, err := syntax.ParseDID(r.URL.Query().Get("actor")) + if err != nil { + _ = writeJson(w, http.StatusBadRequest, atclient.ErrorBody{Name: "InvalidRequest", Message: "actor must be a DID"}) + return + } + + months := defaultCommitStatsMonths + if raw := r.URL.Query().Get("months"); raw != "" { + months, err = strconv.Atoi(raw) + if err != nil || months < 1 || months > maxCommitStatsMonths { + _ = writeJson(w, http.StatusBadRequest, atclient.ErrorBody{Name: "InvalidRequest", Message: "months must be between 1 and 24"}) + return + } + } + + now := time.Now().UTC() + currentMonth := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, time.UTC) + committers := x.resolveCommitters(r.Context(), actor) + counts, err := db.ListCommitCounts( + r.Context(), + x.db, + committers, + currentMonth.AddDate(0, -(months-1), 0), + currentMonth.AddDate(0, 1, 0), + months, + ) + if err != nil { + x.logger.Error("listing commit stats failed", "actor", actor, "err", err) + _ = writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalServerError", Message: "internal server error"}) + return + } + + result := make([]*tangled.GitTemp2GetCommitStats_Month, len(counts)) + for i, count := range counts { + result[i] = &tangled.GitTemp2GetCommitStats_Month{ + Start: count.Month.UTC().Format(time.RFC3339), + Commits: count.Commits, + } + } + _ = writeJson(w, http.StatusOK, tangled.GitTemp2GetCommitStats_Output{Months: result}) +} + +// resolveCommitters maps an actor DID to its verified committer emails with a +// short-lived LRU cache; concurrent lookups of the same actor share one call. +// Failures are not cached: the caller falls back to the actor DID itself. +func (x *Xrpc) resolveCommitters(ctx context.Context, actor syntax.DID) []string { + key := actor.String() + if committers, ok := x.cachedCommitters(key); ok { + return committers + } + result, _, _ := x.committerGroup.Do(key, func() (any, error) { + if committers, ok := x.cachedCommitters(key); ok { + return committers, nil + } + committers, err := x.fetchResolvedCommitters(ctx, actor) + if err != nil { + committerResolverFailures.Inc() + x.logger.Warn("resolving verified committer emails failed", "actor", actor, "err", err) + return []string{key}, nil + } + if x.committers != nil { + x.committers.Add(key, committers) + } + return committers, nil + }) + return result.([]string) +} + +func (x *Xrpc) cachedCommitters(key string) ([]string, bool) { + if x.committers == nil { + return nil, false + } + return x.committers.Get(key) +} + +func (x *Xrpc) fetchResolvedCommitters(ctx context.Context, actor syntax.DID) ([]string, error) { + if x.cfg.DeliberiURL == "" || x.serviceSigner == nil { + return []string{actor.String()}, nil + } + endpoint, err := url.JoinPath(x.cfg.DeliberiURL, "xrpc", tangled.IdentityResolveCommittersNSID) + if err != nil { + return nil, fmt.Errorf("building Deliberi endpoint: %w", err) + } + payload, err := json.Marshal(map[string]string{"actor": actor.String()}) + if err != nil { + return nil, err + } + requestCtx, cancel := context.WithTimeout(ctx, time.Second) + defer cancel() + req, err := http.NewRequestWithContext(requestCtx, http.MethodPost, endpoint, bytes.NewReader(payload)) + if err != nil { + return nil, fmt.Errorf("building committer resolver request: %w", err) + } + audience, err := syntax.ParseDID(x.cfg.DeliberiDID) + if err != nil { + return nil, fmt.Errorf("parsing committer resolver DID: %w", err) + } + token, err := x.serviceSigner.Sign(audience, syntax.NSID(tangled.IdentityResolveCommittersNSID)) + if err != nil { + return nil, fmt.Errorf("signing committer resolver request: %w", err) + } + req.Header.Set("Authorization", "Bearer "+token) + req.Header.Set("Content-Type", "application/json") + client := x.internalClient + if client == nil { + client = x.httpClient + } + if client == nil { + client = http.DefaultClient + } + resp, err := client.Do(req) + if err != nil { + return nil, fmt.Errorf("calling committer resolver: %w", err) + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("committer resolver returned %s", resp.Status) + } + body, err := io.ReadAll(io.LimitReader(resp.Body, 64*1024+1)) + if err != nil { + return nil, fmt.Errorf("reading committer resolver response: %w", err) + } + if len(body) > 64*1024 { + return nil, fmt.Errorf("committer resolver response is too large") + } + var output tangled.IdentityResolveCommitters_Output + if err := json.Unmarshal(body, &output); err != nil { + return nil, fmt.Errorf("decoding committer resolver response: %w", err) + } + + committers := make([]string, 0, min(len(output.Committers)+1, 251)) + committers = append(committers, actor.String()) + seen := map[string]struct{}{actor.String(): {}} + for _, value := range output.Committers { + value = strings.ToLower(strings.TrimSpace(value)) + if value == "" || len(value) > 320 || strings.ContainsRune(value, 0) { + continue + } + if did, err := syntax.ParseDID(value); err == nil { + if did.String() != actor.String() { + continue + } + } else if !strings.Contains(value, "@") { + continue + } + if _, ok := seen[value]; ok { + continue + } + seen[value] = struct{}{} + committers = append(committers, value) + if len(committers) == 251 { + break + } + } + return committers, nil +} diff --git a/knotmirror/xrpc/git_get_commit_stats_test.go b/knotmirror/xrpc/git_get_commit_stats_test.go new file mode 100644 index 000000000..5ec1a60f1 --- /dev/null +++ b/knotmirror/xrpc/git_get_commit_stats_test.go @@ -0,0 +1,129 @@ +package xrpc + +import ( + "context" + "encoding/json" + "io" + "log/slog" + "net/http" + "net/http/httptest" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/bluesky-social/indigo/atproto/atcrypto" + "github.com/bluesky-social/indigo/atproto/identity" + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/hashicorp/golang-lru/v2/expirable" + "tangled.org/core/api/tangled" + "tangled.org/core/knotmirror/config" + "tangled.org/core/xrpc/serviceauth" +) + +const committerResolverAudience = "did:web:deliberi.example" + +func newCommitterResolverServer(t *testing.T, next http.Handler) (*httptest.Server, *serviceauth.Signer) { + t.Helper() + privateKey, err := atcrypto.GeneratePrivateKeyP256() + if err != nil { + t.Fatal(err) + } + signer, err := serviceauth.NewSigner(syntax.DID("did:web:mirror.example"), privateKey.Multibase()) + if err != nil { + t.Fatal(err) + } + document := signer.DIDDocument() + directory := identity.NewMockDirectory() + directory.Insert(identity.ParseIdentity(&document)) + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + verified := serviceauth.NewServiceAuth(logger, directory, committerResolverAudience).VerifyServiceAuth(next) + return httptest.NewServer(verified), signer +} + +func TestFetchResolvedCommitters(t *testing.T) { + const actor = "did:plc:profileactor" + server, signer := newCommitterResolverServer(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if got, want := r.URL.Path, "/xrpc/"+tangled.IdentityResolveCommittersNSID; got != want { + t.Errorf("path = %q, want %q", got, want) + } + var input tangled.IdentityResolveCommitters_Input + if err := json.NewDecoder(r.Body).Decode(&input); err != nil { + t.Fatal(err) + } + if input.Actor != actor { + t.Errorf("actor = %q", input.Actor) + } + _ = json.NewEncoder(w).Encode(tangled.IdentityResolveCommitters_Output{Committers: []string{ + actor, + "Person@Example.com", + "person@example.com", + "not-an-email", + "did:plc:someoneelse", + }}) + })) + defer server.Close() + x := &Xrpc{ + cfg: &config.Config{ + DeliberiURL: server.URL, + DeliberiDID: committerResolverAudience, + }, + logger: slog.New(slog.NewTextHandler(io.Discard, nil)), + internalClient: server.Client(), + serviceSigner: signer, + } + did, _ := syntax.ParseDID(actor) + committers, err := x.fetchResolvedCommitters(context.Background(), did) + if err != nil { + t.Fatal(err) + } + if got, want := strings.Join(committers, ","), actor+",person@example.com"; got != want { + t.Fatalf("committers = %q, want %q", got, want) + } +} + +func TestFetchResolvedCommittersFallsBackWhenDisabled(t *testing.T) { + const actor = "did:plc:profileactor" + x := &Xrpc{cfg: &config.Config{}} + did, _ := syntax.ParseDID(actor) + committers, err := x.fetchResolvedCommitters(context.Background(), did) + if err != nil { + t.Fatal(err) + } + if got := strings.Join(committers, ","); got != actor { + t.Fatalf("committers = %q, want %q", got, actor) + } +} + +func TestResolveCommittersCachesTheInternalResponse(t *testing.T) { + const actor = "did:plc:profileactor" + var requests atomic.Int32 + server, signer := newCommitterResolverServer(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + requests.Add(1) + _ = json.NewEncoder(w).Encode(tangled.IdentityResolveCommitters_Output{Committers: []string{ + actor, + "person@example.com", + }}) + })) + defer server.Close() + x := &Xrpc{ + cfg: &config.Config{ + DeliberiURL: server.URL, + DeliberiDID: committerResolverAudience, + }, + logger: slog.New(slog.NewTextHandler(io.Discard, nil)), + internalClient: server.Client(), + committers: expirable.NewLRU[string, []string](4, nil, time.Minute), + serviceSigner: signer, + } + did, _ := syntax.ParseDID(actor) + for range 2 { + committers := x.resolveCommitters(context.Background(), did) + if got, want := strings.Join(committers, ","), actor+",person@example.com"; got != want { + t.Fatalf("committers = %q, want %q", got, want) + } + } + if requests.Load() != 1 { + t.Fatalf("resolver requests = %d, want 1", requests.Load()) + } +} diff --git a/knotmirror/xrpc/metrics.go b/knotmirror/xrpc/metrics.go index 1618f340d..aa9ccbf0d 100644 --- a/knotmirror/xrpc/metrics.go +++ b/knotmirror/xrpc/metrics.go @@ -21,6 +21,11 @@ var ( Help: "HTTP request duration in seconds", Buckets: prometheus.DefBuckets, }, []string{"method", "path", "status", "repo"}) + + committerResolverFailures = promauto.NewCounter(prometheus.CounterOpts{ + Name: "knotmirror_committer_resolver_failures_total", + Help: "Total failed verified committer identity resolutions", + }) ) type statusRecorder struct { diff --git a/knotmirror/xrpc/xrpc.go b/knotmirror/xrpc/xrpc.go index 988e7d784..6c42f023d 100644 --- a/knotmirror/xrpc/xrpc.go +++ b/knotmirror/xrpc/xrpc.go @@ -13,7 +13,9 @@ import ( "github.com/bluesky-social/indigo/atproto/atclient" "github.com/bluesky-social/indigo/util/ssrf" "github.com/go-chi/chi/v5" + "github.com/hashicorp/golang-lru/v2/expirable" "github.com/redis/go-redis/v9" + "golang.org/x/sync/singleflight" "tangled.org/core/api/tangled" "tangled.org/core/idresolver" "tangled.org/core/knotmirror/config" @@ -24,19 +26,22 @@ import ( ) type Xrpc struct { - cfg *config.Config - db *sql.DB - rdb *redis.Client - indexer *repoindexer.Indexer - resolver *idresolver.Resolver - ks *knotstream.KnotStream - logger *slog.Logger - httpClient *http.Client - v2Client *http.Client + cfg *config.Config + db *sql.DB + rdb *redis.Client + indexer *repoindexer.Indexer + resolver *idresolver.Resolver + ks *knotstream.KnotStream + logger *slog.Logger + httpClient *http.Client + internalClient *http.Client + v2Client *http.Client // write usually takes longer. for example; merge - v2WriteClient *http.Client - inflight *inflightTracker - serviceSigner *serviceauth.Signer + v2WriteClient *http.Client + committers *expirable.LRU[string, []string] + committerGroup singleflight.Group + inflight *inflightTracker + serviceSigner *serviceauth.Signer } func New(logger *slog.Logger, cfg *config.Config, db *sql.DB, rdb *redis.Client, indexer *repoindexer.Indexer, resolver *idresolver.Resolver, ks *knotstream.KnotStream) *Xrpc { @@ -47,17 +52,19 @@ func New(logger *slog.Logger, cfg *config.Config, db *sql.DB, rdb *redis.Client, httpClient.Transport = ssrf.PublicOnlyTransport() } return &Xrpc{ - cfg: cfg, - db: db, - rdb: rdb, - indexer: indexer, - resolver: resolver, - ks: ks, - logger: log.SubLogger(logger, "xrpc"), - httpClient: httpClient, - v2Client: &http.Client{Timeout: 10 * time.Second}, - v2WriteClient: &http.Client{Timeout: 5 * time.Minute}, - inflight: newInflightTracker(), + cfg: cfg, + db: db, + rdb: rdb, + indexer: indexer, + resolver: resolver, + ks: ks, + logger: log.SubLogger(logger, "xrpc"), + httpClient: httpClient, + internalClient: &http.Client{Timeout: time.Second}, + v2Client: &http.Client{Timeout: 10 * time.Second}, + v2WriteClient: &http.Client{Timeout: 5 * time.Minute}, + committers: expirable.NewLRU[string, []string](committerCacheCapacity, nil, committerCacheFreshness), + inflight: newInflightTracker(), } } @@ -68,6 +75,7 @@ func (x *Xrpc) SetServiceSigner(signer *serviceauth.Signer) { func (x *Xrpc) Router() http.Handler { r := chi.NewRouter() r.Use(metricsMiddleware) + r.With(x.inflight.middleware).Get("/"+tangled.GitTemp2GetCommitStatsNSID, x.GetCommitStats) r.Group(func(r chi.Router) { r.Use(x.inflight.middleware) diff --git a/lexicons/git/temp2/getCommitStats.json b/lexicons/git/temp2/getCommitStats.json new file mode 100644 index 000000000..132e072b3 --- /dev/null +++ b/lexicons/git/temp2/getCommitStats.json @@ -0,0 +1,73 @@ +{ + "lexicon": 1, + "id": "sh.tangled.git.temp2.getCommitStats", + "defs": { + "main": { + "type": "query", + "description": "Get an actor's monthly commit counts.", + "parameters": { + "type": "params", + "required": [ + "actor" + ], + "properties": { + "actor": { + "type": "string", + "format": "did", + "description": "Actor whose commits to count." + }, + "months": { + "type": "integer", + "minimum": 1, + "maximum": 24, + "default": 7 + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [ + "months" + ], + "properties": { + "months": { + "type": "array", + "maxLength": 24, + "items": { + "type": "ref", + "ref": "#month" + } + } + } + } + }, + "errors": [ + { + "name": "InvalidRequest", + "description": "Invalid request parameters" + } + ] + }, + "month": { + "type": "object", + "required": [ + "start", + "commits" + ], + "properties": { + "start": { + "type": "string", + "format": "datetime", + "description": "Start of this UTC calendar month." + }, + "commits": { + "type": "integer", + "minimum": 0, + "description": "Commits attributed to the actor during this month." + } + } + } + } +} diff --git a/localinfra/knotmirror.Dockerfile b/localinfra/knotmirror.Dockerfile index b7efebef2..28600b841 100644 --- a/localinfra/knotmirror.Dockerfile +++ b/localinfra/knotmirror.Dockerfile @@ -12,7 +12,7 @@ RUN CGO_ENABLED=0 go build -o /knotmirror ./cmd/knotmirror FROM alpine:3.22 -RUN apk add --no-cache git tini ca-certificates +RUN apk add --no-cache git tini ca-certificates util-linux-misc COPY --from=build /knotmirror /usr/local/bin/knotmirror diff --git a/nix/modules/knotmirror.nix b/nix/modules/knotmirror.nix index 4e732d06a..461fbe90b 100644 --- a/nix/modules/knotmirror.nix +++ b/nix/modules/knotmirror.nix @@ -188,6 +188,7 @@ in wantedBy = ["multi-user.target"]; path = [ pkgs.git + pkgs.util-linux ]; serviceConfig = { LogsDirectory = "knotmirror"; diff --git a/web/src/lib/api/lexicons/index.ts b/web/src/lib/api/lexicons/index.ts index bcbd6e21d..500edc964 100644 --- a/web/src/lib/api/lexicons/index.ts +++ b/web/src/lib/api/lexicons/index.ts @@ -105,7 +105,7 @@ export * as ShTangledGitTempListCommits from "./types/sh/tangled/git/temp/listCo export * as ShTangledGitTempListLanguages from "./types/sh/tangled/git/temp/listLanguages.js"; export * as ShTangledGitTempListTags from "./types/sh/tangled/git/temp/listTags.js"; export * as ShTangledGitTemp2GetBlame from "./types/sh/tangled/git/temp2/getBlame.js"; -export * as ShTangledGitTemp2GetDiff from "./types/sh/tangled/git/temp2/getDiff.js"; +export * as ShTangledGitTemp2GetCommitStats from "./types/sh/tangled/git/temp2/getCommitStats.js";export * as ShTangledGitTemp2GetDiff from "./types/sh/tangled/git/temp2/getDiff.js"; export * as ShTangledGitTemp2GetInterdiff from "./types/sh/tangled/git/temp2/getInterdiff.js"; export * as ShTangledGitTemp2ListCommits from "./types/sh/tangled/git/temp2/listCommits.js"; export * as ShTangledGitTemp2MergeCheck from "./types/sh/tangled/git/temp2/mergeCheck.js"; diff --git a/web/src/lib/api/lexicons/types/sh/tangled/git/temp2/getCommitStats.ts b/web/src/lib/api/lexicons/types/sh/tangled/git/temp2/getCommitStats.ts new file mode 100644 index 000000000..2ab9172dc --- /dev/null +++ b/web/src/lib/api/lexicons/types/sh/tangled/git/temp2/getCommitStats.ts @@ -0,0 +1,73 @@ +import type {} from "@atcute/lexicons"; +import * as v from "@atcute/lexicons/validations"; +import type {} from "@atcute/lexicons/ambient"; + +const _mainSchema = /*#__PURE__*/ v.query( + "sh.tangled.git.temp2.getCommitStats", + { + params: /*#__PURE__*/ v.object({ + /** + * Actor whose commits to count. + */ + actor: /*#__PURE__*/ v.didString(), + /** + * @minimum 1 + * @maximum 24 + * @default 7 + */ + months: /*#__PURE__*/ v.optional( + /*#__PURE__*/ v.constrain(/*#__PURE__*/ v.integer(), [ + /*#__PURE__*/ v.integerRange(1, 24), + ]), + 7, + ), + }), + output: { + type: "lex", + schema: /*#__PURE__*/ v.object({ + /** + * @maxLength 24 + */ + get months() { + return /*#__PURE__*/ v.constrain(/*#__PURE__*/ v.array(monthSchema), [ + /*#__PURE__*/ v.arrayLength(0, 24), + ]); + }, + }), + }, + }, +); +const _monthSchema = /*#__PURE__*/ v.object({ + $type: /*#__PURE__*/ v.optional( + /*#__PURE__*/ v.literal("sh.tangled.git.temp2.getCommitStats#month"), + ), + /** + * Commits attributed to the actor during this month. + * @minimum 0 + */ + commits: /*#__PURE__*/ v.integer(), + /** + * Start of this UTC calendar month. + */ + start: /*#__PURE__*/ v.datetimeString(), +}); + +type main$schematype = typeof _mainSchema; +type month$schematype = typeof _monthSchema; + +export interface mainSchema extends main$schematype {} +export interface monthSchema extends month$schematype {} + +export const mainSchema = _mainSchema as mainSchema; +export const monthSchema = _monthSchema as monthSchema; + +export interface Month extends v.InferInput {} + +export interface $params extends v.InferInput {} +export interface $output extends v.InferXRPCBodyInput {} + +declare module "@atcute/lexicons/ambient" { + interface XRPCQueries { + "sh.tangled.git.temp2.getCommitStats": mainSchema; + } +} -- 2.51.2