diff --git a/docs/partials/metrics.md b/docs/partials/metrics.md index 71412189..5ba04aa9 100644 --- a/docs/partials/metrics.md +++ b/docs/partials/metrics.md @@ -271,6 +271,9 @@ github_status_pull_requests_up{} github_status_webhooks_up{} : Current health status of Webhooks on GitHub status page +github_workflow_job_completed_total{owner, repo, workflow_name, name, conclusion} +: Total number of completed workflow jobs + github_workflow_job_created_timestamp{owner, repo, name, title, branch, sha, identifier, run_id, run_attempt, labels, runner_id, runner_name, runner_group_id, runner_group_name, workflow_name, conclusion} : Timestamp when the workflow job have been created @@ -280,6 +283,9 @@ github_workflow_job_duration_ms{owner, repo, name, title, branch, sha, identifie github_workflow_job_duration_run_created_minutes{owner, repo, name, title, branch, sha, identifier, run_id, run_attempt, labels, runner_id, runner_name, runner_group_id, runner_group_name, workflow_name, conclusion} : Duration since the workflow run creation time in minutes +github_workflow_job_duration_seconds_total{owner, repo, workflow_name, name, conclusion} +: Total duration of completed workflow jobs in seconds + github_workflow_job_started_timestamp{owner, repo, name, title, branch, sha, identifier, run_id, run_attempt, labels, runner_id, runner_name, runner_group_id, runner_group_name, workflow_name, conclusion} : Timestamp when the workflow job have been started diff --git a/pkg/exporter/workflow_job.go b/pkg/exporter/workflow_job.go index 6aca17a1..f9a709ed 100644 --- a/pkg/exporter/workflow_job.go +++ b/pkg/exporter/workflow_job.go @@ -19,11 +19,13 @@ type WorkflowJobCollector struct { duration *prometheus.HistogramVec config config.Target - Status *prometheus.Desc - Duration *prometheus.Desc - Creation *prometheus.Desc - Created *prometheus.Desc - Started *prometheus.Desc + Status *prometheus.Desc + Duration *prometheus.Desc + Creation *prometheus.Desc + Created *prometheus.Desc + Started *prometheus.Desc + CompletedTotal *prometheus.Desc + DurationTotal *prometheus.Desc } // NewWorkflowJobCollector returns a new WorkflowCollector. @@ -33,6 +35,15 @@ func NewWorkflowJobCollector(logger *slog.Logger, client *github.Client, db stor } labels := cfg.WorkflowJobs.Labels + + completionLabels := []string{ + "owner", + "repo", + "workflow_name", + "name", + "conclusion", + } + return &WorkflowJobCollector{ client: client, logger: logger.With("collector", "workflow_job"), @@ -71,6 +82,18 @@ func NewWorkflowJobCollector(logger *slog.Logger, client *github.Client, db stor labels, nil, ), + CompletedTotal: prometheus.NewDesc( + "github_workflow_job_completed_total", + "Total number of completed workflow jobs", + completionLabels, + nil, + ), + DurationTotal: prometheus.NewDesc( + "github_workflow_job_duration_seconds_total", + "Total duration of completed workflow jobs in seconds", + completionLabels, + nil, + ), } } @@ -82,6 +105,8 @@ func (c *WorkflowJobCollector) Metrics() []*prometheus.Desc { c.Creation, c.Created, c.Started, + c.CompletedTotal, + c.DurationTotal, } } @@ -92,6 +117,8 @@ func (c *WorkflowJobCollector) Describe(ch chan<- *prometheus.Desc) { ch <- c.Creation ch <- c.Created ch <- c.Started + ch <- c.CompletedTotal + ch <- c.DurationTotal } // Collect is called by the Prometheus registry when collecting metrics. @@ -180,6 +207,45 @@ func (c *WorkflowJobCollector) Collect(ch chan<- prometheus.Metric) { labels..., ) } + + completions, err := c.db.GetWorkflowJobCompletions() + + if err != nil { + c.logger.Error("Failed to fetch workflow job completions", + "err", err, + ) + + c.failures.WithLabelValues("workflow_job").Inc() + return + } + + c.logger.Debug("Fetched workflow job completions", + "count", len(completions), + ) + + for _, completion := range completions { + labels := []string{ + completion.Owner, + completion.Repo, + completion.WorkflowName, + completion.Name, + completion.Conclusion, + } + + ch <- prometheus.MustNewConstMetric( + c.CompletedTotal, + prometheus.CounterValue, + float64(completion.Count), + labels..., + ) + + ch <- prometheus.MustNewConstMetric( + c.DurationTotal, + prometheus.CounterValue, + completion.DurationSecondsTotal, + labels..., + ) + } } func jobStatusToGauge(conclusion string) float64 { diff --git a/pkg/exporter/workflow_job_test.go b/pkg/exporter/workflow_job_test.go index e66dc638..e951c1d9 100644 --- a/pkg/exporter/workflow_job_test.go +++ b/pkg/exporter/workflow_job_test.go @@ -1,7 +1,6 @@ package exporter import ( - "fmt" "log/slog" "os" "reflect" @@ -10,24 +9,13 @@ import ( "github.com/google/go-github/v92/github" "github.com/prometheus/client_golang/prometheus" + dto "github.com/prometheus/client_model/go" "github.com/promhippie/github_exporter/pkg/config" "github.com/promhippie/github_exporter/pkg/store" ) type StaticStore struct{} -func (s StaticStore) GetWorkflowJobRuns(owner, repo, workflow string) ([]*store.WorkflowRun, error) { - _, _ = fmt.Fprintf( - os.Stdout, - "GetWorkflowJobRuns for %s/%s %s \n", - owner, - repo, - workflow, - ) - - return nil, nil -} - func (s StaticStore) StoreWorkflowRunEvent(*github.WorkflowRunEvent) error { return nil } @@ -52,6 +40,10 @@ func (s StaticStore) PruneWorkflowJobs(time.Duration) error { return nil } +func (s StaticStore) GetWorkflowJobCompletions() ([]*store.WorkflowJobCompletionAggregate, error) { + return nil, nil +} + func (s StaticStore) Open() (bool, error) { return true, nil } @@ -99,23 +91,38 @@ func TestWorkflowJobCollector(t *testing.T) { duration: mockDuration, config: mockConfig, Status: prometheus.NewDesc( - "workflow_job_status", + "github_workflow_job_status", "Status of the workflow job", nil, nil, ), Duration: prometheus.NewDesc( - "workflow_job_duration_seconds", + "github_workflow_job_duration_ms", "Duration of the workflow job", nil, nil, ), Creation: prometheus.NewDesc( - "workflow_job_creation_timestamp_seconds", - "Creation time of the workflow job", + "github_workflow_job_duration_run_created_minutes", + "Duration since the workflow run creation time in minutes", nil, nil, ), Created: prometheus.NewDesc( - "workflow_job_created_timestamp_seconds", - "Created time of the workflow job", + "github_workflow_job_created_timestamp", + "Timestamp when the workflow job have been created", + nil, nil, + ), + Started: prometheus.NewDesc( + "github_workflow_job_started_timestamp", + "Timestamp when the workflow job have been started", + nil, nil, + ), + CompletedTotal: prometheus.NewDesc( + "github_workflow_job_completed_total", + "Total number of completed workflow jobs", + nil, nil, + ), + DurationTotal: prometheus.NewDesc( + "github_workflow_job_duration_seconds_total", + "Total duration of completed workflow jobs in seconds", nil, nil, ), } @@ -139,3 +146,104 @@ func TestWorkflowJobCollector(t *testing.T) { t.Errorf("Expected config to be %v, got %v", mockConfig, collector.config) } } + +type completionStore struct { + StaticStore + completions []*store.WorkflowJobCompletionAggregate +} + +func (s completionStore) GetWorkflowJobCompletions() ([]*store.WorkflowJobCompletionAggregate, error) { + return s.completions, nil +} + +func TestWorkflowJobCollectorCounters(t *testing.T) { + mockLogger := slog.New( + slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{ + Level: slog.LevelDebug, + }), + ) + + mockFailures := prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "test_failures_total", + Help: "Total number of test failures", + }, []string{"type"}) + + mockDuration := prometheus.NewHistogramVec(prometheus.HistogramOpts{ + Name: "test_duration_seconds", + Help: "Duration of test", + }, []string{"type"}) + + completions := []*store.WorkflowJobCompletionAggregate{ + { + Owner: "promhippie", + Repo: "github_exporter", + WorkflowName: "CI", + Name: "test", + Conclusion: "success", + Count: 2, + DurationSecondsTotal: 42.5, + }, + { + Owner: "promhippie", + Repo: "github_exporter", + WorkflowName: "CI", + Name: "test", + Conclusion: "failure", + Count: 1, + DurationSecondsTotal: 10.0, + }, + } + + store := completionStore{completions: completions} + collector := NewWorkflowJobCollector( + mockLogger, + nil, + store, + mockFailures, + mockDuration, + config.Target{}, + ) + + registry := prometheus.NewRegistry() + registry.MustRegister(collector) + + metrics, err := registry.Gather() + if err != nil { + t.Fatalf("failed to gather metrics: %v", err) + } + + expected := map[string]float64{ + "github_workflow_job_completed_total": 3, + "github_workflow_job_duration_seconds_total": 52.5, + } + + for name, expectedValue := range expected { + value := metricFamilyValue(t, metrics, name) + if value != expectedValue { + t.Errorf("expected %s to be %v, got %v", name, expectedValue, value) + } + } +} + +func metricFamilyValue(t *testing.T, metrics []*dto.MetricFamily, name string) float64 { + t.Helper() + + for _, mf := range metrics { + if mf.GetName() != name { + continue + } + + var total float64 + + for _, m := range mf.GetMetric() { + if m.Counter != nil { + total += m.Counter.GetValue() + } + } + + return total + } + + t.Errorf("metric family %s not found", name) + return 0 +} diff --git a/pkg/store/chai.go b/pkg/store/chai.go index 6a16025f..1026afe2 100644 --- a/pkg/store/chai.go +++ b/pkg/store/chai.go @@ -73,6 +73,28 @@ var ( PRIMARY KEY(owner, repo, identifier) );`, }, + { + Version: 4, + Description: "Creating table workflow_job_completions", + Script: `CREATE TABLE workflow_job_completions ( + owner TEXT NOT NULL, + repo TEXT NOT NULL, + identifier BIGINT NOT NULL, + run_attempt INTEGER NOT NULL, + workflow_name TEXT, + name TEXT, + conclusion TEXT, + duration_seconds DOUBLE PRECISION, + recorded_at INTEGER, + PRIMARY KEY(owner, repo, identifier, run_attempt) + );`, + }, + { + Version: 5, + Description: "Creating index for workflow_job_completions aggregate", + Script: `CREATE INDEX idx_workflow_job_completions_aggregate + ON workflow_job_completions(owner, repo, workflow_name, name, conclusion);`, + }, } ) @@ -166,6 +188,11 @@ func (s *chaiStore) PruneWorkflowJobs(timeframe time.Duration) error { return pruneWorkflowJobs(s.handle, timeframe) } +// GetWorkflowJobCompletions implements the Store interface. +func (s *chaiStore) GetWorkflowJobCompletions() ([]*WorkflowJobCompletionAggregate, error) { + return getWorkflowJobCompletions(s.handle) +} + func (s *chaiStore) dsn() string { if len(s.meta) > 0 { return fmt.Sprintf( diff --git a/pkg/store/generic_workflow_job.go b/pkg/store/generic_workflow_job.go index 5ba002aa..d0db98c9 100644 --- a/pkg/store/generic_workflow_job.go +++ b/pkg/store/generic_workflow_job.go @@ -37,7 +37,11 @@ func storeWorkflowJobEvent(handle *sqlx.DB, event *github.WorkflowJobEvent) erro WorkflowName: job.GetWorkflowName(), } - return createOrUpdateWorkflowJob(handle, record) + if err := createOrUpdateWorkflowJob(handle, record); err != nil { + return err + } + + return recordWorkflowJobCompletion(handle, record) } // createOrUpdateWorkflowJob creates or updates the record. @@ -81,6 +85,60 @@ func createOrUpdateWorkflowJob(handle *sqlx.DB, record *WorkflowJob) error { return nil } +// recordWorkflowJobCompletion records a terminal workflow job event in an +// append-only table. The first terminal payload for a given +// (owner, repo, identifier, run_attempt) tuple wins; later payloads with the +// same key are silently ignored to keep the Prometheus counters monotonic. +func recordWorkflowJobCompletion(handle *sqlx.DB, record *WorkflowJob) error { + if record.Status != "completed" || record.Conclusion == "" { + return nil + } + + duration := 0.0 + startedAt := record.StartedAt + completedAt := record.CompletedAt + + if startedAt > 0 && completedAt > 0 { + duration = float64(completedAt - startedAt) + + if duration < 0 { + duration = 0 + } + } + + completion := &WorkflowJobCompletion{ + Owner: record.Owner, + Repo: record.Repo, + Identifier: record.Identifier, + RunAttempt: record.RunAttempt, + WorkflowName: record.WorkflowName, + Name: record.Name, + Conclusion: record.Conclusion, + DurationSeconds: duration, + RecordedAt: time.Now().Unix(), + } + + if _, err := handle.NamedExec( + createWorkflowJobCompletionQuery(handle.DriverName()), + completion, + ); err != nil { + return fmt.Errorf("failed to record completion: %w", err) + } + + return nil +} + +// createWorkflowJobCompletionQuery returns the dialect-specific idempotent +// insert for a workflow job completion. +func createWorkflowJobCompletionQuery(driver string) string { + switch driver { + case "mysql", "mariadb": + return createWorkflowJobCompletionQueryMySQL + default: + return createWorkflowJobCompletionQueryDefault + } +} + // getWorkflowJobs retrieves the workflow jobs from the database. func getWorkflowJobs(handle *sqlx.DB, window time.Duration) ([]*WorkflowJob, error) { records := make([]*WorkflowJob, 0) @@ -120,6 +178,20 @@ func getWorkflowJobs(handle *sqlx.DB, window time.Duration) ([]*WorkflowJob, err return records, nil } +// getWorkflowJobCompletions retrieves aggregated workflow job completions. +func getWorkflowJobCompletions(handle *sqlx.DB) ([]*WorkflowJobCompletionAggregate, error) { + records := make([]*WorkflowJobCompletionAggregate, 0) + + if err := handle.Select( + &records, + selectWorkflowJobCompletionsQuery, + ); err != nil { + return records, err + } + + return records, nil +} + // pruneWorkflowJobs prunes older workflow job records. func pruneWorkflowJobs(handle *sqlx.DB, timeframe time.Duration) error { if _, err := handle.NamedExec( @@ -241,3 +313,69 @@ DELETE FROM workflow_jobs WHERE created_at < :timeframe;` + +var createWorkflowJobCompletionQueryDefault = ` +INSERT INTO workflow_job_completions ( + owner, + repo, + identifier, + run_attempt, + workflow_name, + name, + conclusion, + duration_seconds, + recorded_at +) VALUES ( + :owner, + :repo, + :identifier, + :run_attempt, + :workflow_name, + :name, + :conclusion, + :duration_seconds, + :recorded_at +) +ON CONFLICT DO NOTHING;` + +var createWorkflowJobCompletionQueryMySQL = ` +INSERT INTO workflow_job_completions ( + owner, + repo, + identifier, + run_attempt, + workflow_name, + name, + conclusion, + duration_seconds, + recorded_at +) VALUES ( + :owner, + :repo, + :identifier, + :run_attempt, + :workflow_name, + :name, + :conclusion, + :duration_seconds, + :recorded_at +) +ON DUPLICATE KEY UPDATE owner=owner;` + +var selectWorkflowJobCompletionsQuery = ` +SELECT + owner, + repo, + workflow_name, + name, + conclusion, + COUNT(*) AS count, + COALESCE(SUM(duration_seconds), 0.0) AS duration_seconds_total +FROM + workflow_job_completions +GROUP BY + owner, + repo, + workflow_name, + name, + conclusion;` diff --git a/pkg/store/mysql.go b/pkg/store/mysql.go index 2e923777..e190a046 100644 --- a/pkg/store/mysql.go +++ b/pkg/store/mysql.go @@ -73,6 +73,28 @@ var ( PRIMARY KEY(owner, repo, identifier) );`, }, + { + Version: 4, + Description: "Creating table workflow_job_completions", + Script: `CREATE TABLE workflow_job_completions ( + owner VARCHAR(64) NOT NULL, + repo VARCHAR(100) NOT NULL, + identifier BIGINT NOT NULL, + run_attempt INTEGER NOT NULL, + workflow_name VARCHAR(128), + name VARCHAR(128), + conclusion VARCHAR(64), + duration_seconds DOUBLE, + recorded_at BIGINT, + PRIMARY KEY(owner, repo, identifier, run_attempt) + ) ENGINE=InnoDB CHARACTER SET=utf8;`, + }, + { + Version: 5, + Description: "Creating index for workflow_job_completions aggregate", + Script: `CREATE INDEX idx_workflow_job_completions_aggregate + ON workflow_job_completions(owner, repo, workflow_name, name, conclusion);`, + }, } ) @@ -177,6 +199,11 @@ func (s *mysqlStore) PruneWorkflowJobs(timeframe time.Duration) error { return pruneWorkflowJobs(s.handle, timeframe) } +// GetWorkflowJobCompletions implements the Store interface. +func (s *mysqlStore) GetWorkflowJobCompletions() ([]*WorkflowJobCompletionAggregate, error) { + return getWorkflowJobCompletions(s.handle) +} + func (s *mysqlStore) dsn() string { if s.password != "" { return fmt.Sprintf( diff --git a/pkg/store/postgres.go b/pkg/store/postgres.go index de48eb46..efb25ec3 100644 --- a/pkg/store/postgres.go +++ b/pkg/store/postgres.go @@ -83,6 +83,28 @@ var ( Description: "Fix run_id be BIGINT", Script: `ALTER TABLE workflow_jobs ALTER COLUMN run_id TYPE BIGINT USING run_id::BIGINT;`, }, + { + Version: 6, + Description: "Creating table workflow_job_completions", + Script: `CREATE TABLE workflow_job_completions ( + owner TEXT NOT NULL, + repo TEXT NOT NULL, + identifier BIGINT NOT NULL, + run_attempt INTEGER NOT NULL, + workflow_name TEXT, + name TEXT, + conclusion TEXT, + duration_seconds DOUBLE PRECISION, + recorded_at BIGINT, + PRIMARY KEY(owner, repo, identifier, run_attempt) + );`, + }, + { + Version: 7, + Description: "Creating index for workflow_job_completions aggregate", + Script: `CREATE INDEX idx_workflow_job_completions_aggregate + ON workflow_job_completions(owner, repo, workflow_name, name, conclusion);`, + }, } ) @@ -187,6 +209,11 @@ func (s *postgresStore) PruneWorkflowJobs(timeframe time.Duration) error { return pruneWorkflowJobs(s.handle, timeframe) } +// GetWorkflowJobCompletions implements the Store interface. +func (s *postgresStore) GetWorkflowJobCompletions() ([]*WorkflowJobCompletionAggregate, error) { + return getWorkflowJobCompletions(s.handle) +} + func (s *postgresStore) dsn() string { dsn := fmt.Sprintf( "host=%s port=%s dbname=%s user=%s", diff --git a/pkg/store/sqlite.go b/pkg/store/sqlite.go index f968d832..58c87953 100644 --- a/pkg/store/sqlite.go +++ b/pkg/store/sqlite.go @@ -74,6 +74,28 @@ var ( PRIMARY KEY(owner, repo, identifier) );`, }, + { + Version: 4, + Description: "Creating table workflow_job_completions", + Script: `CREATE TABLE workflow_job_completions ( + owner TEXT NOT NULL, + repo TEXT NOT NULL, + identifier BIGINT NOT NULL, + run_attempt INTEGER NOT NULL, + workflow_name TEXT, + name TEXT, + conclusion TEXT, + duration_seconds REAL, + recorded_at BIGINT, + PRIMARY KEY(owner, repo, identifier, run_attempt) + );`, + }, + { + Version: 5, + Description: "Creating index for workflow_job_completions aggregate", + Script: `CREATE INDEX idx_workflow_job_completions_aggregate + ON workflow_job_completions(owner, repo, workflow_name, name, conclusion);`, + }, } ) @@ -174,6 +196,11 @@ func (s *sqliteStore) PruneWorkflowJobs(timeframe time.Duration) error { return pruneWorkflowJobs(s.handle, timeframe) } +// GetWorkflowJobCompletions implements the Store interface. +func (s *sqliteStore) GetWorkflowJobCompletions() ([]*WorkflowJobCompletionAggregate, error) { + return getWorkflowJobCompletions(s.handle) +} + func (s *sqliteStore) dsn() string { if len(s.meta) > 0 { return fmt.Sprintf( diff --git a/pkg/store/store.go b/pkg/store/store.go index 648ec725..25b1350d 100644 --- a/pkg/store/store.go +++ b/pkg/store/store.go @@ -30,6 +30,7 @@ type Store interface { StoreWorkflowJobEvent(*github.WorkflowJobEvent) error GetWorkflowJobs(time.Duration) ([]*WorkflowJob, error) PruneWorkflowJobs(time.Duration) error + GetWorkflowJobCompletions() ([]*WorkflowJobCompletionAggregate, error) Open() (bool, error) Close() error diff --git a/pkg/store/types.go b/pkg/store/types.go index fb381d55..418674e1 100644 --- a/pkg/store/types.go +++ b/pkg/store/types.go @@ -141,3 +141,30 @@ func (r *WorkflowJob) ByLabel(label string) string { return "" } + +// WorkflowJobCompletion defines the append-only record of a terminal workflow +// job event. Each (owner, repo, identifier, run_attempt) tuple is recorded at +// most once; the first terminal payload wins. +type WorkflowJobCompletion struct { + Owner string `db:"owner"` + Repo string `db:"repo"` + Identifier int64 `db:"identifier"` + RunAttempt int `db:"run_attempt"` + WorkflowName string `db:"workflow_name"` + Name string `db:"name"` + Conclusion string `db:"conclusion"` + DurationSeconds float64 `db:"duration_seconds"` + RecordedAt int64 `db:"recorded_at"` +} + +// WorkflowJobCompletionAggregate groups completion records for counter +// emission. +type WorkflowJobCompletionAggregate struct { + Owner string `db:"owner"` + Repo string `db:"repo"` + WorkflowName string `db:"workflow_name"` + Name string `db:"name"` + Conclusion string `db:"conclusion"` + Count int64 `db:"count"` + DurationSecondsTotal float64 `db:"duration_seconds_total"` +}