diff --git a/CHANGELOG.md b/CHANGELOG.md
index 3c6c2169..c8f0d7d5 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -11,6 +11,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- Job list: filter jobs matching any of the selected exact tags. [PR #548](https://github.com/riverqueue/riverui/pull/548).
+### Changed
+
+- Job state sidebar: keep large counts responsive while preserving useful magnitude with bounded live counts, adaptively cached exact snapshots, and PostgreSQL estimates. [PR #XXX](https://github.com/riverqueue/riverui/pull/XXX).
+
### Fixed
- Job args: preserve large numeric JSON values exactly when displaying and copying args, while keeping object keys sorted. [Fixes #593](https://github.com/riverqueue/riverui/issues/593). [PR #594](https://github.com/riverqueue/riverui/pull/594).
diff --git a/handler_api_endpoint.go b/handler_api_endpoint.go
index dce312f1..f74dbad9 100644
--- a/handler_api_endpoint.go
+++ b/handler_api_endpoint.go
@@ -12,6 +12,7 @@ import (
"github.com/jackc/pgerrcode"
"github.com/jackc/pgx/v5/pgconn"
+ "github.com/jackc/pgx/v5/pgtype"
"github.com/riverqueue/apiframe/apiendpoint"
"github.com/riverqueue/apiframe/apierror"
@@ -878,23 +879,36 @@ type stateAndCountGetEndpoint[TTx any] struct {
apibundle.APIBundle[TTx]
apiendpoint.Endpoint[jobCancelRequest, stateAndCountGetResponse]
- queryCacheSkipThreshold int // constant normally, but settable for testing
- queryCacher *querycacher.QueryCacher[map[rivertype.JobState]int]
+ boundedQueryCacher *querycacher.QueryCacher[stateCountSnapshot]
+ countMax int
+ estimateCounts func(ctx context.Context, states []rivertype.JobState) (map[rivertype.JobState]stateCountEstimate, error)
+ exactQueryCacher *querycacher.QueryCacher[stateCountSnapshot]
}
-func newStateAndCountGetEndpoint[TTx any](bundle apibundle.APIBundle[TTx]) *stateAndCountGetEndpoint[TTx] {
- runQuery := func(ctx context.Context) (map[rivertype.JobState]int, error) {
- return dbutil.WithTxV(ctx, bundle.DB, func(ctx context.Context, execTx riverdriver.ExecutorTx) (map[rivertype.JobState]int, error) {
- tx := bundle.Driver.UnwrapTx(execTx)
+const (
+ stateAndCountDefaultMax = 10_000
+
+ // Two missed maximum-interval refreshes make an estimate preferable to an
+ // increasingly misleading exact snapshot.
+ stateCountExactMaxAge = 1 * time.Hour
+ stateCountExactRefreshMin = 1 * time.Minute
+ stateCountExactRefreshMax = 30 * time.Minute
+ stateCountExactRefreshCostMul = 100
+)
- return bundle.Driver.UnwrapExecutor(tx).JobCountByAllStates(ctx, &riverdriver.JobCountByAllStatesParams{Schema: bundle.Client.Schema()})
- })
- }
- return &stateAndCountGetEndpoint[TTx]{
- APIBundle: bundle,
- queryCacheSkipThreshold: 1_000_000,
- queryCacher: querycacher.NewQueryCacher(bundle.Archetype, runQuery),
+func newStateAndCountGetEndpoint[TTx any](bundle apibundle.APIBundle[TTx]) *stateAndCountGetEndpoint[TTx] {
+ endpoint := &stateAndCountGetEndpoint[TTx]{
+ APIBundle: bundle,
+ countMax: stateAndCountDefaultMax,
}
+ endpoint.boundedQueryCacher = querycacher.NewQueryCacher(bundle.Archetype, endpoint.queryBoundedCounts)
+ endpoint.exactQueryCacher = querycacher.NewQueryCacherWithOpts(
+ bundle.Archetype,
+ endpoint.queryExactCounts,
+ &querycacher.QueryCacherOpts{NextTickPeriod: stateCountExactRefreshPeriod},
+ )
+ endpoint.estimateCounts = endpoint.queryEstimatedCounts
+ return endpoint
}
func (*stateAndCountGetEndpoint[TTx]) Meta() *apiendpoint.EndpointMeta {
@@ -905,63 +919,377 @@ func (*stateAndCountGetEndpoint[TTx]) Meta() *apiendpoint.EndpointMeta {
}
func (a *stateAndCountGetEndpoint[TTx]) SubServices() []startstop.Service {
- return []startstop.Service{a.queryCacher}
+ return []startstop.Service{a.boundedQueryCacher, a.exactQueryCacher}
}
type stateAndCountGetRequest struct{}
-type stateAndCountGetResponse struct {
- Available int `json:"available"`
- Cancelled int `json:"cancelled"`
- Completed int `json:"completed"`
- Discarded int `json:"discarded"`
- Pending int `json:"pending"`
- Retryable int `json:"retryable"`
- Running int `json:"running"`
- Scheduled int `json:"scheduled"`
+type stateCountAccuracy string
+
+const (
+ stateCountAccuracyEstimated stateCountAccuracy = "estimated" // uses Postgres planner estimate (Postgres only)
+ stateCountAccuracyExact stateCountAccuracy = "exact" // exact
+ stateCountAccuracyExactCached stateCountAccuracy = "exact_cached" // exact (cached)
+ stateCountAccuracyLowerBound stateCountAccuracy = "lower_bound" // constrained to stateAndCountDefaultMax
+)
+
+type stateCountResponse struct {
+ Accuracy stateCountAccuracy `json:"accuracy"`
+ Count int `json:"count"`
+ ObservedAt *time.Time `json:"observed_at,omitempty"`
}
+type stateAndCountGetResponse struct {
+ Available stateCountResponse `json:"available"`
+ Cancelled stateCountResponse `json:"cancelled"`
+ Completed stateCountResponse `json:"completed"`
+ Discarded stateCountResponse `json:"discarded"`
+ Pending stateCountResponse `json:"pending"`
+ Retryable stateCountResponse `json:"retryable"`
+ Running stateCountResponse `json:"running"`
+ Scheduled stateCountResponse `json:"scheduled"`
+}
+
+// Execute resolves every state's count from the cheapest sufficiently useful
+// source. A bounded index scan gives fresh exact values for small states. Large
+// states prefer a recent exact snapshot refreshed adaptively in the background,
+// then a PostgreSQL planner estimate, and finally the bound proven by the index
+// scan. Full exact scans are never part of request latency.
func (a *stateAndCountGetEndpoint[TTx]) Execute(ctx context.Context, _ *stateAndCountGetRequest) (*stateAndCountGetResponse, error) {
- // Counts the total number of jobs in a state and count result.
- totalJobs := func(stateAndCountRes map[rivertype.JobState]int) int {
- var totalJobs int
- for _, count := range stateAndCountRes {
- totalJobs += count
+ countsAreExact := func(snapshot stateCountSnapshot) bool {
+ for _, count := range snapshot.Counts {
+ if count > a.countMax {
+ return false
+ }
+ }
+ return true
+ }
+
+ // Prefer fresh counts while every state is below the cap. Once any state is
+ // capped, serve the periodically refreshed result to collapse queries from
+ // multiple UI clients. Both paths use the same bounded query.
+ boundedSnapshot, ok := a.boundedQueryCacher.CachedRes()
+ if !ok || countsAreExact(boundedSnapshot) {
+ var err error
+ boundedSnapshot, err = a.queryBoundedCounts(ctx)
+ if err != nil {
+ return nil, fmt.Errorf("error getting states and counts: %w", err)
+ }
+ }
+
+ cappedStates := make([]rivertype.JobState, 0, len(allJobStates))
+ for _, state := range allJobStates {
+ if boundedSnapshot.Counts[state] > a.countMax {
+ cappedStates = append(cappedStates, state)
}
- return totalJobs
}
- // Counting jobs can be an expensive operation given a large table, so in
- // the presence of such, prefer to use a result that's cached periodically
- // instead of querying inline with the API request. In case we don't have a
- // cached result yet or there's a relatively small number of job rows, run
- // the query directly (in the case of the latter so we present the freshest
- // possible information).
- stateAndCountRes, ok := a.queryCacher.CachedRes()
- if !ok || totalJobs(stateAndCountRes) < a.queryCacheSkipThreshold {
+ var (
+ exactSnapshot, hasExactSnapshot = a.exactQueryCacher.CachedRes()
+ exactSnapshotIsFresh = hasExactSnapshot && time.Since(exactSnapshot.ObservedAt) <= stateCountExactMaxAge
+ )
+
+ statesNeedingEstimate := make([]rivertype.JobState, 0, len(cappedStates))
+ for _, state := range cappedStates {
+ if !exactSnapshotIsFresh || exactSnapshot.Counts[state] <= a.countMax {
+ statesNeedingEstimate = append(statesNeedingEstimate, state)
+ }
+ }
+
+ estimates := make(map[rivertype.JobState]stateCountEstimate)
+ if len(statesNeedingEstimate) > 0 {
var err error
- stateAndCountRes, err = dbutil.WithTxV(ctx, a.DB, func(ctx context.Context, execTx riverdriver.ExecutorTx) (map[rivertype.JobState]int, error) {
- tx := a.Driver.UnwrapTx(execTx)
+ estimates, err = a.estimateCounts(ctx, statesNeedingEstimate)
+ if err != nil {
+ // Estimates are an optional telemetry enhancement. The bounded count
+ // is still trustworthy, so degrade to a lower bound instead of failing
+ // the entire sidebar when planner statistics can't be read.
+ a.Logger.WarnContext(ctx, "Unable to estimate large job counts", "err", err)
+ estimates = make(map[rivertype.JobState]stateCountEstimate)
+ }
+ }
+
+ resolvedCounts := make(map[rivertype.JobState]stateCountResponse, len(allJobStates))
+ for _, state := range allJobStates {
+ boundedCount := boundedSnapshot.Counts[state]
+
+ if boundedCount <= a.countMax {
+ // The bounded scan reached the end of this state's index range, so the
+ // value is exact and fresh even if another, larger state was capped.
+ resolvedCounts[state] = stateCountResponse{
+ Accuracy: stateCountAccuracyExact,
+ Count: boundedCount,
+ ObservedAt: &boundedSnapshot.ObservedAt,
+ }
+ continue
+ }
+
+ if exactSnapshotIsFresh && exactSnapshot.Counts[state] > a.countMax {
+ // A recent full scan preserves the useful magnitude for common large
+ // states. Its timestamp makes the deliberate staleness visible.
+ resolvedCounts[state] = stateCountResponse{
+ Accuracy: stateCountAccuracyExactCached,
+ Count: exactSnapshot.Counts[state],
+ ObservedAt: &exactSnapshot.ObservedAt,
+ }
+ continue
+ }
+
+ if estimate, ok := estimates[state]; ok && estimate.Count > a.countMax {
+ // Planner statistics are cheap and retain an order of magnitude during
+ // cold start or when the last exact snapshot has become too old.
+ resolvedCounts[state] = stateCountResponse{
+ Accuracy: stateCountAccuracyEstimated,
+ Count: estimate.Count,
+ ObservedAt: estimate.ObservedAt,
+ }
+ continue
+ }
- return a.Driver.UnwrapExecutor(tx).JobCountByAllStates(ctx, &riverdriver.JobCountByAllStatesParams{Schema: a.Client.Schema()})
+ // The bounded scan proves only that there are more than countMax rows.
+ // Never present a stale planner estimate below that known lower bound.
+ resolvedCounts[state] = stateCountResponse{
+ Accuracy: stateCountAccuracyLowerBound,
+ Count: a.countMax,
+ ObservedAt: &boundedSnapshot.ObservedAt,
+ }
+ }
+
+ resp := &stateAndCountGetResponse{
+ Available: resolvedCounts[rivertype.JobStateAvailable],
+ Cancelled: resolvedCounts[rivertype.JobStateCancelled],
+ Completed: resolvedCounts[rivertype.JobStateCompleted],
+ Discarded: resolvedCounts[rivertype.JobStateDiscarded],
+ Pending: resolvedCounts[rivertype.JobStatePending],
+ Retryable: resolvedCounts[rivertype.JobStateRetryable],
+ Running: resolvedCounts[rivertype.JobStateRunning],
+ Scheduled: resolvedCounts[rivertype.JobStateScheduled],
+ }
+
+ return resp, nil
+}
+
+type stateCountSnapshot struct {
+ Counts map[rivertype.JobState]int
+ ObservedAt time.Time
+}
+
+type stateCountEstimate struct {
+ Count int
+ ObservedAt *time.Time
+}
+
+var allJobStates = []rivertype.JobState{ //nolint:gochecknoglobals
+ rivertype.JobStateAvailable,
+ rivertype.JobStateCancelled,
+ rivertype.JobStateCompleted,
+ rivertype.JobStateDiscarded,
+ rivertype.JobStatePending,
+ rivertype.JobStateRetryable,
+ rivertype.JobStateRunning,
+ rivertype.JobStateScheduled,
+}
+
+func jobStateSQLLiteral(state rivertype.JobState) (string, error) {
+ // These are deliberately explicit rather than quoting an arbitrary string.
+ // queryEstimatedCounts embeds the result in SQL so PostgreSQL always plans
+ // against a state constant, even if its prepared statement cache later
+ // chooses a generic plan.
+ switch state {
+ case rivertype.JobStateAvailable:
+ return "'available'", nil
+ case rivertype.JobStateCancelled:
+ return "'cancelled'", nil
+ case rivertype.JobStateCompleted:
+ return "'completed'", nil
+ case rivertype.JobStateDiscarded:
+ return "'discarded'", nil
+ case rivertype.JobStatePending:
+ return "'pending'", nil
+ case rivertype.JobStateRetryable:
+ return "'retryable'", nil
+ case rivertype.JobStateRunning:
+ return "'running'", nil
+ case rivertype.JobStateScheduled:
+ return "'scheduled'", nil
+ default:
+ return "", fmt.Errorf("invalid job state for count estimate: %q", state)
+ }
+}
+
+func (a *stateAndCountGetEndpoint[TTx]) queryBoundedCounts(ctx context.Context) (stateCountSnapshot, error) {
+ return dbutil.WithTxV(ctx, a.DB, func(ctx context.Context, execTx riverdriver.ExecutorTx) (stateCountSnapshot, error) {
+ counts, err := jobCountByAllStatesCapped(ctx, execTx, a.Driver.ArgPlaceholder(), a.Client.Schema(), a.countMax)
+ if err != nil {
+ return stateCountSnapshot{}, err
+ }
+ return stateCountSnapshot{Counts: counts, ObservedAt: time.Now()}, nil
+ })
+}
+
+func (a *stateAndCountGetEndpoint[TTx]) queryExactCounts(ctx context.Context) (stateCountSnapshot, error) {
+ return dbutil.WithTxV(ctx, a.DB, func(ctx context.Context, execTx riverdriver.ExecutorTx) (stateCountSnapshot, error) {
+ counts, err := execTx.JobCountByAllStates(ctx, &riverdriver.JobCountByAllStatesParams{
+ Schema: a.Client.Schema(),
})
if err != nil {
- return nil, fmt.Errorf("error getting states and counts: %w", err)
+ return stateCountSnapshot{}, fmt.Errorf("error counting all jobs by state exactly: %w", err)
}
+ return stateCountSnapshot{Counts: counts, ObservedAt: time.Now()}, nil
+ })
+}
+
+func stateCountExactRefreshPeriod(queryDuration time.Duration, queryErr error) time.Duration {
+ if queryErr != nil {
+ // A failed full-table count is likely load-related. Back off to the
+ // maximum interval instead of repeatedly adding pressure to the database.
+ return stateCountExactRefreshMax
+ }
+
+ // Target about one percent of wall time for full exact counts. Fast counts
+ // still wait at least a minute, while the maximum keeps exact telemetry from
+ // disappearing entirely on very large installations.
+ refreshPeriod := queryDuration * stateCountExactRefreshCostMul
+ return min(max(refreshPeriod, stateCountExactRefreshMin), stateCountExactRefreshMax)
+}
+
+// jobCountByAllStatesCapped counts at most countMax+1 rows for every job state.
+// A result at or below countMax is exact; countMax+1 is a sentinel proving that
+// more rows exist without making the request scan the state's entire index.
+func jobCountByAllStatesCapped(ctx context.Context, exec riverdriver.Executor, argPlaceholder, schema string, countMax int) (map[rivertype.JobState]int, error) {
+ if countMax < 1 {
+ return nil, errors.New("count max must be positive")
}
- return &stateAndCountGetResponse{
- Available: stateAndCountRes[rivertype.JobStateAvailable],
- Cancelled: stateAndCountRes[rivertype.JobStateCancelled],
- Completed: stateAndCountRes[rivertype.JobStateCompleted],
- Discarded: stateAndCountRes[rivertype.JobStateDiscarded],
- Pending: stateAndCountRes[rivertype.JobStatePending],
- Retryable: stateAndCountRes[rivertype.JobStateRetryable],
- Running: stateAndCountRes[rivertype.JobStateRunning],
- Scheduled: stateAndCountRes[rivertype.JobStateScheduled],
+ jobsTable := dbutil.SafeIdentifier("river_job")
+ if schema != "" {
+ jobsTable = dbutil.SafeIdentifier(schema) + "." + jobsTable
+ }
+
+ // Each subquery returns at most countMax+1 rows. The extra entry lets the
+ // caller distinguish an exact count of countMax from a capped count.
+ // Ordering by the remaining columns in river_job_prioritized_fetching_index
+ // encourages an index-only scan that can stop as soon as the limit is met.
+ query := fmt.Sprintf(`
+SELECT
+ (SELECT count(*) FROM (SELECT 1 FROM %[1]s WHERE state = 'available' ORDER BY queue, priority, scheduled_at, id LIMIT %[2]s) AS limited_available),
+ (SELECT count(*) FROM (SELECT 1 FROM %[1]s WHERE state = 'cancelled' ORDER BY queue, priority, scheduled_at, id LIMIT %[2]s) AS limited_cancelled),
+ (SELECT count(*) FROM (SELECT 1 FROM %[1]s WHERE state = 'completed' ORDER BY queue, priority, scheduled_at, id LIMIT %[2]s) AS limited_completed),
+ (SELECT count(*) FROM (SELECT 1 FROM %[1]s WHERE state = 'discarded' ORDER BY queue, priority, scheduled_at, id LIMIT %[2]s) AS limited_discarded),
+ (SELECT count(*) FROM (SELECT 1 FROM %[1]s WHERE state = 'pending' ORDER BY queue, priority, scheduled_at, id LIMIT %[2]s) AS limited_pending),
+ (SELECT count(*) FROM (SELECT 1 FROM %[1]s WHERE state = 'retryable' ORDER BY queue, priority, scheduled_at, id LIMIT %[2]s) AS limited_retryable),
+ (SELECT count(*) FROM (SELECT 1 FROM %[1]s WHERE state = 'running' ORDER BY queue, priority, scheduled_at, id LIMIT %[2]s) AS limited_running),
+ (SELECT count(*) FROM (SELECT 1 FROM %[1]s WHERE state = 'scheduled' ORDER BY queue, priority, scheduled_at, id LIMIT %[2]s) AS limited_scheduled)`,
+ jobsTable,
+ argPlaceholder+"1",
+ )
+
+ var (
+ available int64
+ cancelled int64
+ completed int64
+ discarded int64
+ pending int64
+ retryable int64
+ running int64
+ scheduled int64
+ )
+ if err := exec.QueryRow(ctx, query, countMax+1).Scan(
+ &available,
+ &cancelled,
+ &completed,
+ &discarded,
+ &pending,
+ &retryable,
+ &running,
+ &scheduled,
+ ); err != nil {
+ return nil, fmt.Errorf("error counting jobs by state: %w", err)
+ }
+
+ return map[rivertype.JobState]int{
+ rivertype.JobStateAvailable: int(available),
+ rivertype.JobStateCancelled: int(cancelled),
+ rivertype.JobStateCompleted: int(completed),
+ rivertype.JobStateDiscarded: int(discarded),
+ rivertype.JobStatePending: int(pending),
+ rivertype.JobStateRetryable: int(retryable),
+ rivertype.JobStateRunning: int(running),
+ rivertype.JobStateScheduled: int(scheduled),
}, nil
}
+// queryEstimatedCounts asks PostgreSQL to plan, but not execute, one query per
+// state and returns each plan's estimated row count. Estimates are used only
+// when a bounded count is known to exceed countMax and no recent exact snapshot
+// is available; non-PostgreSQL databases fall back to that known lower bound.
+func (a *stateAndCountGetEndpoint[TTx]) queryEstimatedCounts(ctx context.Context, states []rivertype.JobState) (map[rivertype.JobState]stateCountEstimate, error) {
+ if a.Driver.DatabaseName() != riverdriver.DatabaseNamePostgres {
+ return nil, errors.New("job count estimates are only available for PostgreSQL")
+ }
+
+ return dbutil.WithTxV(ctx, a.DB, func(ctx context.Context, execTx riverdriver.ExecutorTx) (map[rivertype.JobState]stateCountEstimate, error) {
+ jobsTable := dbutil.SafeIdentifier("river_job")
+ if schema := a.Client.Schema(); schema != "" {
+ jobsTable = dbutil.SafeIdentifier(schema) + "." + jobsTable
+ }
+
+ // EXPLAIN's Plan Rows comes from PostgreSQL's existing ANALYZE statistics,
+ // so it gives us order-of-magnitude telemetry without reading every
+ // matching row. last_analyze makes that estimate's freshness visible.
+ var analyzedAt pgtype.Timestamptz
+ _ = execTx.QueryRow(ctx, `
+SELECT GREATEST(last_analyze, last_autoanalyze)
+FROM pg_stat_all_tables
+WHERE schemaname = COALESCE(NULLIF(`+a.Driver.ArgPlaceholder()+`1, ''), current_schema())
+ AND relname = 'river_job'`, a.Client.Schema()).Scan(&analyzedAt)
+
+ var observedAt *time.Time
+ if analyzedAt.Valid {
+ observedAtCopy := analyzedAt.Time
+ observedAt = &observedAtCopy
+ }
+
+ type explainPlan struct {
+ Plan struct {
+ Rows int `json:"Plan Rows"` //nolint:tagliatelle // PostgreSQL owns this JSON key.
+ } `json:"Plan"` //nolint:tagliatelle // PostgreSQL owns this JSON key.
+ }
+
+ estimates := make(map[rivertype.JobState]stateCountEstimate, len(states))
+ for _, state := range states {
+ stateLiteral, err := jobStateSQLLiteral(state)
+ if err != nil {
+ return nil, err
+ }
+
+ // A literal makes the statement text state-specific. A parameter here
+ // could eventually receive PostgreSQL's generic prepared plan, losing
+ // the per-state selectivity that makes this estimate useful.
+ query := fmt.Sprintf("EXPLAIN (FORMAT JSON) SELECT 1 FROM %s WHERE state = %s", jobsTable, stateLiteral)
+ var rawPlan []byte
+ if err := execTx.QueryRow(ctx, query).Scan(&rawPlan); err != nil {
+ return nil, fmt.Errorf("error explaining job count for state %q: %w", state, err)
+ }
+
+ var plans []explainPlan
+ if err := json.Unmarshal(rawPlan, &plans); err != nil {
+ return nil, fmt.Errorf("error decoding job count estimate for state %q: %w", state, err)
+ }
+ if len(plans) != 1 {
+ return nil, fmt.Errorf("expected one job count estimate plan for state %q, got %d", state, len(plans))
+ }
+
+ estimates[state] = stateCountEstimate{
+ Count: plans[0].Plan.Rows,
+ ObservedAt: observedAt,
+ }
+ }
+
+ return estimates, nil
+ })
+}
+
func NewNotFoundJob(jobID int64) *apierror.NotFound {
return apierror.NewNotFoundf("Job not found: %d.", jobID)
}
diff --git a/handler_api_endpoint_test.go b/handler_api_endpoint_test.go
index a5b918a9..d47c8f70 100644
--- a/handler_api_endpoint_test.go
+++ b/handler_api_endpoint_test.go
@@ -3,6 +3,7 @@ package riverui
import (
"context"
"encoding/json"
+ "errors"
"log/slog"
"net/http"
"net/http/httptest"
@@ -1054,6 +1055,25 @@ func TestStateAndCountGetEndpoint(t *testing.T) {
t.Parallel()
ctx := context.Background()
+ stateCountsFromResponse := func(resp *stateAndCountGetResponse) map[rivertype.JobState]*stateCountResponse {
+ return map[rivertype.JobState]*stateCountResponse{
+ rivertype.JobStateAvailable: &resp.Available,
+ rivertype.JobStateCancelled: &resp.Cancelled,
+ rivertype.JobStateCompleted: &resp.Completed,
+ rivertype.JobStateDiscarded: &resp.Discarded,
+ rivertype.JobStatePending: &resp.Pending,
+ rivertype.JobStateRetryable: &resp.Retryable,
+ rivertype.JobStateRunning: &resp.Running,
+ rivertype.JobStateScheduled: &resp.Scheduled,
+ }
+ }
+ requireExactCounts := func(t *testing.T, resp *stateAndCountGetResponse) {
+ t.Helper()
+ for state, stateCount := range stateCountsFromResponse(resp) {
+ require.Equal(t, stateCountAccuracyExact, stateCount.Accuracy, state)
+ require.NotNil(t, stateCount.ObservedAt, state)
+ }
+ }
t.Run("Success", func(t *testing.T) {
t.Parallel()
@@ -1092,55 +1112,231 @@ func TestStateAndCountGetEndpoint(t *testing.T) {
resp, err := apitest.InvokeHandler(ctx, endpoint.Execute, testMountOpts(t), &stateAndCountGetRequest{})
require.NoError(t, err)
- require.Equal(t, &stateAndCountGetResponse{
- Available: 1,
- Cancelled: 2,
- Completed: 3,
- Discarded: 4,
- Pending: 5,
- Retryable: 6,
- Running: 7,
- Scheduled: 8,
- }, resp)
+ requireExactCounts(t, resp)
+ require.Equal(t, 1, resp.Available.Count)
+ require.Equal(t, 2, resp.Cancelled.Count)
+ require.Equal(t, 3, resp.Completed.Count)
+ require.Equal(t, 4, resp.Discarded.Count)
+ require.Equal(t, 5, resp.Pending.Count)
+ require.Equal(t, 6, resp.Retryable.Count)
+ require.Equal(t, 7, resp.Running.Count)
+ require.Equal(t, 8, resp.Scheduled.Count)
})
- t.Run("WithCachedQueryAboveSkipThreshold", func(t *testing.T) {
+ t.Run("AtCountMaxIsExact", func(t *testing.T) {
t.Parallel()
- endpoint, bundle := setupEndpoint(ctx, t, newStateAndCountGetEndpoint)
+ const countMax = 3
+ endpoint, bundle := setupEndpoint(ctx, t, func(bundle apibundle.APIBundle[pgx.Tx]) *stateAndCountGetEndpoint[pgx.Tx] {
+ endpoint := newStateAndCountGetEndpoint(bundle)
+ endpoint.countMax = countMax
+ return endpoint
+ })
- const queryCacheSkipThreshold = 3
- for range queryCacheSkipThreshold + 1 {
+ for range countMax {
_ = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateAvailable)})
}
- _, err := endpoint.queryCacher.RunQuery(ctx)
+ resp, err := apitest.InvokeHandler(ctx, endpoint.Execute, testMountOpts(t), &stateAndCountGetRequest{})
require.NoError(t, err)
+ requireExactCounts(t, resp)
+ require.Equal(t, countMax, resp.Available.Count)
+ })
+
+ t.Run("WithExactCachedSnapshot", func(t *testing.T) {
+ t.Parallel()
+
+ const countMax = 3
+ endpoint, bundle := setupEndpoint(ctx, t, func(bundle apibundle.APIBundle[pgx.Tx]) *stateAndCountGetEndpoint[pgx.Tx] {
+ endpoint := newStateAndCountGetEndpoint(bundle)
+ endpoint.countMax = countMax
+ return endpoint
+ })
+
+ for range countMax + 1 {
+ _ = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateAvailable)})
+ }
+
+ _, err := endpoint.boundedQueryCacher.RunQuery(ctx)
+ require.NoError(t, err)
+ _, err = endpoint.exactQueryCacher.RunQuery(ctx)
+ require.NoError(t, err)
+
+ // Once a state is capped, both caches are reused instead of making an
+ // exact count part of the request's latency.
+ _ = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateCancelled), FinalizedAt: ptrutil.Ptr(time.Now())})
resp, err := apitest.InvokeHandler(ctx, endpoint.Execute, testMountOpts(t), &stateAndCountGetRequest{})
require.NoError(t, err)
- require.Equal(t, &stateAndCountGetResponse{
- Available: queryCacheSkipThreshold + 1,
- }, resp)
+ require.Equal(t, countMax+1, resp.Available.Count)
+ require.Equal(t, stateCountAccuracyExactCached, resp.Available.Accuracy)
+ require.NotNil(t, resp.Available.ObservedAt)
+ require.Equal(t, 0, resp.Cancelled.Count)
+ require.Equal(t, stateCountAccuracyExact, resp.Cancelled.Accuracy)
})
- t.Run("WithCachedQueryBelowSkipThreshold", func(t *testing.T) {
+ t.Run("WithExactCachedCount", func(t *testing.T) {
t.Parallel()
- endpoint, bundle := setupEndpoint(ctx, t, newStateAndCountGetEndpoint)
+ const countMax = 3
+ endpoint, bundle := setupEndpoint(ctx, t, func(bundle apibundle.APIBundle[pgx.Tx]) *stateAndCountGetEndpoint[pgx.Tx] {
+ endpoint := newStateAndCountGetEndpoint(bundle)
+ endpoint.countMax = countMax
+ return endpoint
+ })
- const queryCacheSkipThreshold = 3
- for range queryCacheSkipThreshold - 1 {
+ for range countMax - 1 {
_ = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateAvailable)})
}
- _, err := endpoint.queryCacher.RunQuery(ctx)
+ _, err := endpoint.boundedQueryCacher.RunQuery(ctx)
require.NoError(t, err)
+ // An exact cache result is refreshed inline for the latest counts.
+ _ = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateCancelled), FinalizedAt: ptrutil.Ptr(time.Now())})
+
resp, err := apitest.InvokeHandler(ctx, endpoint.Execute, testMountOpts(t), &stateAndCountGetRequest{})
require.NoError(t, err)
- require.Equal(t, &stateAndCountGetResponse{
- Available: queryCacheSkipThreshold - 1,
- }, resp)
+ requireExactCounts(t, resp)
+ require.Equal(t, countMax-1, resp.Available.Count)
+ require.Equal(t, 1, resp.Cancelled.Count)
+ })
+
+ t.Run("WithPlannerEstimate", func(t *testing.T) {
+ t.Parallel()
+
+ const countMax = 3
+ endpoint, bundle := setupEndpoint(ctx, t, func(bundle apibundle.APIBundle[pgx.Tx]) *stateAndCountGetEndpoint[pgx.Tx] {
+ endpoint := newStateAndCountGetEndpoint(bundle)
+ endpoint.countMax = countMax
+ return endpoint
+ })
+
+ for range countMax + 1 {
+ _ = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateCompleted), FinalizedAt: ptrutil.Ptr(time.Now())})
+ }
+ _, err := endpoint.boundedQueryCacher.RunQuery(ctx)
+ require.NoError(t, err)
+
+ observedAt := time.Now().Add(-5 * time.Minute)
+ endpoint.estimateCounts = func(_ context.Context, states []rivertype.JobState) (map[rivertype.JobState]stateCountEstimate, error) {
+ require.Equal(t, []rivertype.JobState{rivertype.JobStateCompleted}, states)
+ return map[rivertype.JobState]stateCountEstimate{
+ rivertype.JobStateCompleted: {Count: 1_000_000, ObservedAt: &observedAt},
+ }, nil
+ }
+
+ resp, err := apitest.InvokeHandler(ctx, endpoint.Execute, testMountOpts(t), &stateAndCountGetRequest{})
+ require.NoError(t, err)
+ require.Equal(t, stateCountResponse{
+ Accuracy: stateCountAccuracyEstimated,
+ Count: 1_000_000,
+ ObservedAt: &observedAt,
+ }, resp.Completed)
+ })
+
+ t.Run("ReadsPlannerEstimateFromPostgres", func(t *testing.T) {
+ t.Parallel()
+
+ endpoint, bundle := setupEndpoint(ctx, t, newStateAndCountGetEndpoint)
+ for range 100 {
+ _ = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateCompleted), FinalizedAt: ptrutil.Ptr(time.Now())})
+ }
+ for range 10 {
+ _ = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateRunning)})
+ }
+ require.NoError(t, bundle.exec.Exec(ctx, "ANALYZE river_job"))
+
+ estimates, err := endpoint.queryEstimatedCounts(ctx, []rivertype.JobState{
+ rivertype.JobStateCompleted,
+ rivertype.JobStateRunning,
+ })
+ require.NoError(t, err)
+ require.Positive(t, estimates[rivertype.JobStateCompleted].Count)
+ require.NotNil(t, estimates[rivertype.JobStateCompleted].ObservedAt)
+ require.Greater(t, estimates[rivertype.JobStateCompleted].Count, estimates[rivertype.JobStateRunning].Count)
})
+
+ t.Run("WithLowerBoundForStaleEstimate", func(t *testing.T) {
+ t.Parallel()
+
+ const countMax = 3
+ endpoint, bundle := setupEndpoint(ctx, t, func(bundle apibundle.APIBundle[pgx.Tx]) *stateAndCountGetEndpoint[pgx.Tx] {
+ endpoint := newStateAndCountGetEndpoint(bundle)
+ endpoint.countMax = countMax
+ return endpoint
+ })
+
+ for range countMax + 1 {
+ _ = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateAvailable)})
+ }
+ _, err := endpoint.boundedQueryCacher.RunQuery(ctx)
+ require.NoError(t, err)
+ endpoint.estimateCounts = func(_ context.Context, _ []rivertype.JobState) (map[rivertype.JobState]stateCountEstimate, error) {
+ return map[rivertype.JobState]stateCountEstimate{
+ rivertype.JobStateAvailable: {Count: countMax - 1},
+ }, nil
+ }
+
+ resp, err := apitest.InvokeHandler(ctx, endpoint.Execute, testMountOpts(t), &stateAndCountGetRequest{})
+ require.NoError(t, err)
+ require.Equal(t, countMax, resp.Available.Count)
+ require.Equal(t, stateCountAccuracyLowerBound, resp.Available.Accuracy)
+ require.NotNil(t, resp.Available.ObservedAt)
+ })
+}
+
+func TestAllJobStates(t *testing.T) {
+ t.Parallel()
+
+ // Keep the endpoint's exhaustive response and SQL allowlist synchronized
+ // with River when a job state is added or reordered upstream.
+ require.Equal(t, rivertype.JobStates(), allJobStates)
+}
+
+func TestStateCountExactRefreshPeriod(t *testing.T) {
+ t.Parallel()
+
+ require.Equal(t, stateCountExactRefreshMin, stateCountExactRefreshPeriod(100*time.Millisecond, nil))
+ require.Equal(t, 200*time.Second, stateCountExactRefreshPeriod(2*time.Second, nil))
+ require.Equal(t, stateCountExactRefreshMax, stateCountExactRefreshPeriod(time.Hour, nil))
+ require.Equal(t, stateCountExactRefreshMax, stateCountExactRefreshPeriod(time.Second, errors.New("database busy")))
+}
+
+func TestJobStateSQLLiteral(t *testing.T) {
+ t.Parallel()
+
+ expected := map[rivertype.JobState]string{
+ rivertype.JobStateAvailable: "'available'",
+ rivertype.JobStateCancelled: "'cancelled'",
+ rivertype.JobStateCompleted: "'completed'",
+ rivertype.JobStateDiscarded: "'discarded'",
+ rivertype.JobStatePending: "'pending'",
+ rivertype.JobStateRetryable: "'retryable'",
+ rivertype.JobStateRunning: "'running'",
+ rivertype.JobStateScheduled: "'scheduled'",
+ }
+ for state, expectedLiteral := range expected {
+ literal, err := jobStateSQLLiteral(state)
+ require.NoError(t, err)
+ require.Equal(t, expectedLiteral, literal)
+ }
+
+ _, err := jobStateSQLLiteral(rivertype.JobState("completed'; DROP TABLE river_job; --"))
+ require.EqualError(t, err, `invalid job state for count estimate: "completed'; DROP TABLE river_job; --"`)
+}
+
+func TestStateAndCountGetEndpointCustomSchema(t *testing.T) {
+ t.Parallel()
+
+ ctx := context.Background()
+ endpoint, bundle := setupEndpointWithCustomSchema(ctx, t, newStateAndCountGetEndpoint)
+ jobParams := testfactory.Job_Build(t, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateRunning)})
+ jobParams.Schema = bundle.client.Schema()
+ _, err := bundle.exec.JobInsertFull(ctx, jobParams)
+ require.NoError(t, err)
+
+ resp, err := apitest.InvokeHandler(ctx, endpoint.Execute, testMountOpts(t), &stateAndCountGetRequest{})
+ require.NoError(t, err)
+ require.Equal(t, 1, resp.Running.Count)
}
diff --git a/internal/querycacher/query_cacher.go b/internal/querycacher/query_cacher.go
index e7ff0b28..357d471e 100644
--- a/internal/querycacher/query_cacher.go
+++ b/internal/querycacher/query_cacher.go
@@ -23,20 +23,37 @@ type QueryCacher[TRes any] struct {
cachedRes TRes
cachedResSet bool
mu sync.RWMutex
+ nextTickPeriod func(queryDuration time.Duration, queryErr error) time.Duration
runQuery func(ctx context.Context) (TRes, error)
runQueryTestChan chan struct{} // closed when query is run; for testing
tickPeriod time.Duration // constant normally, but settable for testing
}
+type QueryCacherOpts struct {
+ // NextTickPeriod makes the interval adaptive to the cost and result of the
+ // preceding query. The period starts after the query finishes, so an
+ // expensive query can never cause this service to run continuously.
+ NextTickPeriod func(queryDuration time.Duration, queryErr error) time.Duration
+}
+
func NewQueryCacher[TRes any](archetype *baseservice.Archetype, runQuery func(ctx context.Context) (TRes, error)) *QueryCacher[TRes] {
+ return NewQueryCacherWithOpts(archetype, runQuery, nil)
+}
+
+func NewQueryCacherWithOpts[TRes any](archetype *baseservice.Archetype, runQuery func(ctx context.Context) (TRes, error), opts *QueryCacherOpts) *QueryCacher[TRes] {
// +/- 1s random variance to ticker interval. Makes sure that given multiple
// query caches running simultaneously, they all start and are scheduled a
// little differently to make a thundering herd problem less likely.
randomTickVariance := time.Duration(rand.Float64()*float64(2*time.Second)) - 1*time.Second
+ var nextTickPeriod func(queryDuration time.Duration, queryErr error) time.Duration
+ if opts != nil {
+ nextTickPeriod = opts.NextTickPeriod
+ }
queryCacher := baseservice.Init(archetype, &QueryCacher[TRes]{
- runQuery: runQuery,
- tickPeriod: 10*time.Second + randomTickVariance,
+ nextTickPeriod: nextTickPeriod,
+ runQuery: runQuery,
+ tickPeriod: 10*time.Second + randomTickVariance,
})
// TODO(brandur): Push this up into baseservice.
@@ -76,7 +93,7 @@ func (s *QueryCacher[TRes]) RunQuery(ctx context.Context) (TRes, error) {
return emptyRes, err
}
- s.Logger.DebugContext(ctx, s.Name+": Ran query and cached result", "duration", time.Since(start), "tick_period", s.tickPeriod)
+ s.Logger.DebugContext(ctx, s.Name+": Ran query and cached result", "duration", time.Since(start))
s.mu.Lock()
s.cachedRes = res
@@ -104,20 +121,36 @@ func (s *QueryCacher[TRes]) Start(ctx context.Context) error {
started()
defer stopped()
- // In case a query runs long and exceeds tickPeriod, time.Ticker will
- // drop ticks to compensate.
- ticker := time.NewTicker(s.tickPeriod)
- defer ticker.Stop()
+ // A timer is reset only after each query finishes. Unlike a ticker, this
+ // prevents a slow query from leaving a pending tick that starts another
+ // expensive query immediately.
+ timer := time.NewTimer(s.tickPeriod)
+ defer timer.Stop()
for {
select {
case <-ctx.Done():
return
- case <-ticker.C:
- if _, err := s.RunQuery(ctx); err != nil {
+ case <-timer.C:
+ start := time.Now()
+ _, err := s.RunQuery(ctx)
+ queryDuration := time.Since(start)
+ if err != nil {
s.Logger.ErrorContext(ctx, s.Name+": Error running query", "err", err)
}
+
+ nextTickPeriod := s.tickPeriod
+ if s.nextTickPeriod != nil {
+ nextTickPeriod = s.nextTickPeriod(queryDuration, err)
+ }
+ if nextTickPeriod <= 0 {
+ // A non-positive period would make the service spin. Falling
+ // back to the base interval is safer than treating bad options
+ // as permission to continuously query the database.
+ nextTickPeriod = s.tickPeriod
+ }
+ timer.Reset(nextTickPeriod)
}
}
}()
diff --git a/internal/querycacher/query_cacher_test.go b/internal/querycacher/query_cacher_test.go
index df813462..67bd0b51 100644
--- a/internal/querycacher/query_cacher_test.go
+++ b/internal/querycacher/query_cacher_test.go
@@ -115,6 +115,33 @@ func TestQueryCacher(t *testing.T) {
}, res)
})
+ t.Run("UsesAdaptivePeriodAfterQueryFinishes", func(t *testing.T) {
+ t.Parallel()
+
+ var queryFinishedAt time.Time
+ nextPeriodCalled := make(chan struct{})
+ queryCacher := NewQueryCacherWithOpts(
+ riversharedtest.BaseServiceArchetype(t),
+ func(_ context.Context) (int, error) {
+ queryFinishedAt = time.Now()
+ return 1, nil
+ },
+ &QueryCacherOpts{
+ NextTickPeriod: func(_ time.Duration, queryErr error) time.Duration {
+ require.NoError(t, queryErr)
+ require.False(t, queryFinishedAt.IsZero())
+ close(nextPeriodCalled)
+ return time.Hour
+ },
+ },
+ )
+ queryCacher.tickPeriod = time.Millisecond
+
+ require.NoError(t, queryCacher.Start(ctx))
+ t.Cleanup(queryCacher.Stop)
+ riversharedtest.WaitOrTimeout(t, nextPeriodCalled)
+ })
+
t.Run("StartStopStress", func(t *testing.T) {
t.Parallel()
diff --git a/src/components/JobList.tsx b/src/components/JobList.tsx
index 8a7ab4f9..ecdd9186 100644
--- a/src/components/JobList.tsx
+++ b/src/components/JobList.tsx
@@ -482,9 +482,9 @@ const JobList = (props: JobListProps) => {
const stateFormatted = state.charAt(0).toUpperCase() + state.slice(1);
const jobsInState = useMemo(() => {
if (!statesAndCounts) {
- return 0;
+ return BigInt(0);
}
- return statesAndCounts[state] || 0;
+ return statesAndCounts[state].count;
}, [state, statesAndCounts]);
const filterItems = useMemo(
diff --git a/src/components/JobStateFilters.test.tsx b/src/components/JobStateFilters.test.tsx
index fb76ba35..59552bbc 100644
--- a/src/components/JobStateFilters.test.tsx
+++ b/src/components/JobStateFilters.test.tsx
@@ -1,3 +1,4 @@
+import { StatesAndCounts } from "@services/states";
import { JobState } from "@services/types";
import {
createMemoryHistory,
@@ -14,30 +15,29 @@ import { describe, expect, test } from "vitest";
import { defaultValues, jobSearchSchema } from "../routes/jobs/index.schema";
import { JobStateFilters } from "./JobStateFilters";
-const rootRoute = createRootRoute({
- component: () => ,
-});
-
-const jobsRoute = createRoute({
- component: () => ,
- getParentRoute: () => rootRoute,
- path: "/jobs",
- search: {
- middlewares: [stripSearchParams(defaultValues)],
- },
- validateSearch: jobSearchSchema,
-});
-
-const routeTree = rootRoute.addChildren([jobsRoute]);
-
-const renderWithLocation = async (location: string) => {
+const renderWithLocation = async (
+ location: string,
+ statesAndCounts?: StatesAndCounts,
+) => {
+ const rootRoute = createRootRoute({
+ component: () => ,
+ });
+ const jobsRoute = createRoute({
+ component: () => ,
+ getParentRoute: () => rootRoute,
+ path: "/jobs",
+ search: {
+ middlewares: [stripSearchParams(defaultValues)],
+ },
+ validateSearch: jobSearchSchema,
+ });
const history = createMemoryHistory({
initialEntries: [location],
});
const router = createRouter({
history,
- routeTree,
+ routeTree: rootRoute.addChildren([jobsRoute]),
});
await router.load();
@@ -45,6 +45,20 @@ const renderWithLocation = async (location: string) => {
return render();
};
+const statesAndCounts = (
+ overrides: Partial,
+): StatesAndCounts => ({
+ available: { accuracy: "exact", count: 0n },
+ cancelled: { accuracy: "exact", count: 0n },
+ completed: { accuracy: "exact", count: 0n },
+ discarded: { accuracy: "exact", count: 0n },
+ pending: { accuracy: "exact", count: 0n },
+ retryable: { accuracy: "exact", count: 0n },
+ running: { accuracy: "exact", count: 0n },
+ scheduled: { accuracy: "exact", count: 0n },
+ ...overrides,
+});
+
describe("JobStateFilters", () => {
test("only the selected state link is active", async () => {
await renderWithLocation(`/jobs?state=${JobState.Discarded}`);
@@ -64,4 +78,43 @@ describe("JobStateFilters", () => {
const runningLink = await screen.findByRole("link", { name: "Running" });
expect(runningLink).toHaveAttribute("data-status", "active");
});
+
+ test("shows exact, cached, estimated, and lower-bound telemetry", async () => {
+ const observedAt = new Date("2026-08-10T12:00:00Z");
+ await renderWithLocation(
+ "/jobs",
+ statesAndCounts({
+ available: {
+ accuracy: "lower_bound",
+ count: 10_000n,
+ observedAt,
+ },
+ completed: {
+ accuracy: "exact_cached",
+ count: 12_345_678n,
+ observedAt,
+ },
+ discarded: {
+ accuracy: "estimated",
+ count: 987_654n,
+ observedAt,
+ },
+ running: { accuracy: "exact", count: 2n, observedAt },
+ }),
+ );
+
+ expect(await screen.findByText("10K+")).toHaveAttribute(
+ "title",
+ expect.stringContaining("At least 10,000 jobs"),
+ );
+ expect(screen.getByText("12.3M")).toHaveAttribute(
+ "title",
+ expect.stringContaining("12,345,678 jobs (exact snapshot"),
+ );
+ expect(screen.getByText("≈987.7K")).toHaveAttribute(
+ "title",
+ expect.stringContaining("Approximately 987,654 jobs"),
+ );
+ expect(screen.getByText("2")).toBeInTheDocument();
+ });
});
diff --git a/src/components/JobStateFilters.tsx b/src/components/JobStateFilters.tsx
index 757119ed..7018f6a9 100644
--- a/src/components/JobStateFilters.tsx
+++ b/src/components/JobStateFilters.tsx
@@ -6,10 +6,54 @@ import React, { useMemo } from "react";
import { Badge } from "./Badge";
+const compactCountFormatter = new Intl.NumberFormat("en-US", {
+ maximumFractionDigits: 1,
+ notation: "compact",
+});
+
type JobStateFiltersProps = {
statesAndCounts?: StatesAndCounts;
};
+const formatFilterItemCount = (
+ item: ReturnType[number],
+): string => {
+ switch (item.accuracy) {
+ case "estimated":
+ // The approximation marker prevents a planner estimate from looking
+ // indistinguishable from an exact snapshot.
+ return `≈${compactCountFormatter.format(item.count)}`;
+ case "exact":
+ // Small exact values are easiest to scan without abbreviation.
+ return item.count.toString();
+ case "exact_cached":
+ // Compact notation retains the useful order of magnitude in a narrow
+ // sidebar; the tooltip below keeps the full exact snapshot available.
+ return compactCountFormatter.format(item.count);
+ case "lower_bound":
+ // A plus is the strongest claim supported by the bounded index scan.
+ return `${compactCountFormatter.format(item.count)}+`;
+ }
+};
+
+const filterItemCountTitle = (
+ item: ReturnType[number],
+): string => {
+ const fullCount = item.count.toLocaleString("en-US");
+ const observedAt = item.observedAt?.toLocaleString();
+
+ switch (item.accuracy) {
+ case "estimated":
+ return `Approximately ${fullCount} jobs (PostgreSQL statistics${observedAt ? ` from ${observedAt}` : ""})`;
+ case "exact":
+ return `${fullCount} jobs (exact)`;
+ case "exact_cached":
+ return `${fullCount} jobs (exact snapshot${observedAt ? ` from ${observedAt}` : ""})`;
+ case "lower_bound":
+ return `At least ${fullCount} jobs; an exact snapshot or useful PostgreSQL estimate is not available yet`;
+ }
+};
+
export const JobStateFilters: (
props: JobStateFiltersProps,
) => React.JSX.Element = ({ statesAndCounts }) => {
@@ -56,8 +100,9 @@ export const JobStateFilters: (
- {item.count.toString()}
+ {formatFilterItemCount(item)}
) : null}
diff --git a/src/services/states.ts b/src/services/states.ts
index 96aafc92..307d6980 100644
--- a/src/services/states.ts
+++ b/src/services/states.ts
@@ -2,14 +2,41 @@ import type { QueryFunction } from "@tanstack/react-query";
import { API } from "@utils/api";
-import type { JobState, SnakeToCamelCase } from "./types";
+import { JobState } from "./types";
+
+export type StateCount = {
+ accuracy: StateCountAccuracy;
+ count: bigint;
+ observedAt?: Date;
+};
+
+export type StateCountAccuracy =
+ "estimated" | "exact_cached" | "exact" | "lower_bound";
export type StatesAndCounts = {
- [Key in JobState as SnakeToCamelCase]: bigint;
+ [Key in JobState]: StateCount;
};
type CountsByStateKey = ["countsByState"];
+type StatesAndCountsFromAPI = {
+ [Key in JobState]: {
+ accuracy: StateCountAccuracy;
+ count: number;
+ observed_at?: string;
+ };
+};
+
+const stateCountFromAPI = (
+ stateCount: StatesAndCountsFromAPI[JobState],
+): StateCount => ({
+ accuracy: stateCount.accuracy,
+ count: BigInt(stateCount.count),
+ observedAt: stateCount.observed_at
+ ? new Date(stateCount.observed_at)
+ : undefined,
+});
+
export const countsByStateKey = (): CountsByStateKey => {
return ["countsByState"];
};
@@ -18,7 +45,16 @@ export const countsByState: QueryFunction<
StatesAndCounts,
CountsByStateKey
> = async ({ signal }) => {
- return API.get({ path: "/states" }, { signal }).then(
- (response) => response,
+ return API.get({ path: "/states" }, { signal }).then(
+ (response) => ({
+ available: stateCountFromAPI(response.available),
+ cancelled: stateCountFromAPI(response.cancelled),
+ completed: stateCountFromAPI(response.completed),
+ discarded: stateCountFromAPI(response.discarded),
+ pending: stateCountFromAPI(response.pending),
+ retryable: stateCountFromAPI(response.retryable),
+ running: stateCountFromAPI(response.running),
+ scheduled: stateCountFromAPI(response.scheduled),
+ }),
);
};
diff --git a/src/utils/jobStateFilterItems.ts b/src/utils/jobStateFilterItems.ts
index a934dcb0..e1ca5266 100644
--- a/src/utils/jobStateFilterItems.ts
+++ b/src/utils/jobStateFilterItems.ts
@@ -1,60 +1,64 @@
-import { StatesAndCounts } from "@services/states";
+import { StateCountAccuracy, StatesAndCounts } from "@services/states";
import { JobState } from "@services/types";
export type JobStateFilterItem = {
+ accuracy: StateCountAccuracy;
count: bigint;
name: string;
+ observedAt?: Date;
state: JobState;
};
export const jobStateFilterItems: (
statesAndCounts: StatesAndCounts | undefined,
) => JobStateFilterItem[] = (statesAndCounts) => {
- const getCount = (state: JobState): bigint => {
- if (statesAndCounts) {
- return BigInt(statesAndCounts[state]);
- }
- return BigInt(0);
+ const getStateCount = (state: JobState) => {
+ return (
+ statesAndCounts?.[state] ?? {
+ accuracy: "exact" as const,
+ count: BigInt(0),
+ }
+ );
};
return [
{
- count: getCount(JobState.Pending),
+ ...getStateCount(JobState.Pending),
name: "Pending",
state: JobState.Pending,
},
{
- count: getCount(JobState.Scheduled),
+ ...getStateCount(JobState.Scheduled),
name: "Scheduled",
state: JobState.Scheduled,
},
{
- count: getCount(JobState.Available),
+ ...getStateCount(JobState.Available),
name: "Available",
state: JobState.Available,
},
{
- count: getCount(JobState.Running),
+ ...getStateCount(JobState.Running),
name: "Running",
state: JobState.Running,
},
{
- count: getCount(JobState.Retryable),
+ ...getStateCount(JobState.Retryable),
name: "Retryable",
state: JobState.Retryable,
},
{
- count: getCount(JobState.Cancelled),
+ ...getStateCount(JobState.Cancelled),
name: "Cancelled",
state: JobState.Cancelled,
},
{
- count: getCount(JobState.Discarded),
+ ...getStateCount(JobState.Discarded),
name: "Discarded",
state: JobState.Discarded,
},
{
- count: getCount(JobState.Completed),
+ ...getStateCount(JobState.Completed),
name: "Completed",
state: JobState.Completed,
},