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 }