mirror of
https://github.com/immich-app/yucca.git
synced 2026-09-30 13:33:00 +08:00
feat(monk): add monk, a measured ceph scrub-backlog exporter (#587)
* feat(monk): add monk, a measured ceph scrub-backlog exporter * fix(monk): harden collection failure paths * docs(monk): defer the deploy role reference to its follow-up PR * fix(monk): emit the age distribution as a real histogram and document operational contracts * feat(monk): read overdue targets from the cluster instead of flags * ci(monk): register the monk image build filter * chore(monk): cut comments the types already carry
This commit is contained in:
@@ -59,6 +59,7 @@ jobs:
|
||||
web: *node
|
||||
michael: ['packages/michael/**', '.dockerignore']
|
||||
columbo: ['packages/columbo/**', '.dockerignore']
|
||||
monk: ['packages/monk/**', '.dockerignore']
|
||||
|
||||
build:
|
||||
name: Build ${{ matrix.app }}
|
||||
|
||||
@@ -56,6 +56,7 @@ jobs:
|
||||
- { name: web, dockerfile: packages/web/Dockerfile }
|
||||
- { name: michael, dockerfile: packages/michael/Dockerfile }
|
||||
- { name: columbo, dockerfile: packages/columbo/Dockerfile }
|
||||
- { name: monk, dockerfile: packages/monk/Dockerfile }
|
||||
steps:
|
||||
- name: Checkout
|
||||
uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6.0.2
|
||||
|
||||
+1
-1
@@ -3,7 +3,7 @@
|
||||
set -e
|
||||
|
||||
status=0
|
||||
for pkg in michael yuctl columbo; do
|
||||
for pkg in michael yuctl columbo monk; do
|
||||
(cd "packages/$pkg" && golangci-lint run ./... "$@") || status=1
|
||||
done
|
||||
exit $status
|
||||
|
||||
Executable
+5
@@ -0,0 +1,5 @@
|
||||
#!/usr/bin/env bash
|
||||
#MISE description="Build monk"
|
||||
set -e
|
||||
|
||||
cd packages/monk && go build -o ../../dist/monk .
|
||||
Executable
+6
@@ -0,0 +1,6 @@
|
||||
#!/usr/bin/env bash
|
||||
#MISE description="Run unit tests for monk"
|
||||
#MISE dir="{{config_root}}/packages/monk"
|
||||
set -e
|
||||
|
||||
go test ./... "$@"
|
||||
@@ -107,6 +107,7 @@ Zod-validated `env.ts`, JWT auth guards via `@AuthRoute()`, OTel from `@common/s
|
||||
| `yucca-admin-api` | NestJS | Admin API (user/session/repository management). Same DB + JWT validation. |
|
||||
| `michael` | Go | **Production** restic REST backend — S3 proxy with JWT verification, WORM enforcement, backend pooling. |
|
||||
| `columbo` | Go | Ticket investigation agent: LLM loop (OpenRouter) over per-user-scoped o11y queries, answers only into staff threads. See `docs/columbo.md`. |
|
||||
| `monk` | Go | Ceph scrub-backlog exporter: polls `pg ls`, serves measured per-pool scrub-age metrics on :9284. Deployed onto mon hosts by ansible, not K8s. |
|
||||
| `restic-api` | NestJS | Earlier TS implementation of the restic backend, kept as **reference**; not deployed. |
|
||||
| `yucca-metrics-worker` | NestJS | 5-min cron: RadosGW usage → meter tables → per-connection rollup (`connectionMetrics`, billing floor), OTel gauges. |
|
||||
| `redis` (valkey) | | Shared platform cache (ephemeral; keys `yucca:<service>:<purpose>:*`). Primary-region only. |
|
||||
|
||||
@@ -0,0 +1,24 @@
|
||||
ARG ALPINE_VERSION=3.23
|
||||
# The runtime base is the cluster's own ceph image, not alpine: monk execs the
|
||||
# ceph CLI, and the mon hosts already hold every layer of this base for the
|
||||
# cluster daemons, so the pull cost is only the binary layer. Bump together
|
||||
# with the cluster ceph version (ansible ceph_upgrade_target_image).
|
||||
ARG CEPH_IMAGE=quay.io/ceph/ceph:v20.2.2@sha256:6b4b5ae33acd3d736eb26d2a19238bce71a22f9cfb99cca887ba6312d0957644
|
||||
|
||||
FROM golang:1.27-alpine${ALPINE_VERSION} AS builder
|
||||
WORKDIR /app
|
||||
|
||||
COPY packages/monk/go.mod packages/monk/go.sum ./
|
||||
RUN go mod download
|
||||
|
||||
COPY packages/monk/ ./
|
||||
RUN CGO_ENABLED=0 go build -o /monk .
|
||||
|
||||
FROM ${CEPH_IMAGE}
|
||||
|
||||
USER ceph
|
||||
|
||||
COPY --from=builder /monk /usr/local/bin/monk
|
||||
|
||||
EXPOSE 9284
|
||||
ENTRYPOINT ["/usr/local/bin/monk"]
|
||||
@@ -0,0 +1,58 @@
|
||||
# monk
|
||||
|
||||
Prometheus exporter for measured Ceph scrub backlog. Ceph tracks per-PG
|
||||
`last_scrub_stamp` / `last_deep_scrub_stamp` but exports neither as metrics
|
||||
(the mgr prometheus module ships only PG state counts and the
|
||||
`PG_NOT_SCRUBBED` health booleans; upstream PR ceph/ceph#68925, which
|
||||
implemented exactly this, was stale-bot-closed unmerged in Aug 2026). monk
|
||||
polls `ceph pg ls -f json` and serves pool-level aggregates so scrub-cycle
|
||||
dashboards report ground truth instead of estimates derived from scrub
|
||||
read-byte counters.
|
||||
|
||||
Runs on the cluster's mon hosts, not Kubernetes; the ansible role that deploys
|
||||
it lands in a follow-up PR. The image is the cluster ceph image plus the monk
|
||||
binary, and the container needs `/etc/ceph` with a read-only keyring (mon r,
|
||||
mgr r) mounted. The container runs as the image's `ceph` user (uid 167), so
|
||||
the keyring file must be readable by that uid; a root-owned 0600 keyring fails
|
||||
as `no keyring found`.
|
||||
|
||||
## Run
|
||||
|
||||
```
|
||||
monk # cluster host with ceph CLI + keyring
|
||||
monk -ceph-cmd "cephadm shell -- ceph"
|
||||
```
|
||||
|
||||
Flags: `-listen :9284`, `-refresh 2m`, `-timeout 90s`, `-ceph-cmd ceph`
|
||||
(space-split command prefix). Overdue targets follow the cluster: each
|
||||
refresh reads `osd_scrub_max_interval` / `osd_deep_scrub_interval` from
|
||||
`ceph config get osd` plus per-pool overrides from `osd pool ls detail`, so
|
||||
the thresholds cannot drift from what the scrub scheduler targets.
|
||||
`-shallow-interval` / `-deep-interval` pin a depth explicitly instead (pins
|
||||
also suppress that depth's pool overrides); a failed interval read keeps the
|
||||
last-known targets and logs a warning.
|
||||
|
||||
## Metrics
|
||||
|
||||
| metric | labels | meaning |
|
||||
|---|---|---|
|
||||
| `ceph_pg_last_scrub_stamp` | pool_id | oldest per-PG shallow stamp in the pool, epoch seconds (name follows ceph PR #68925) |
|
||||
| `ceph_pg_last_deep_scrub_stamp` | pool_id | oldest per-PG deep stamp in the pool |
|
||||
| `ceph_scrub_pool_pgs` / `ceph_scrub_pool_bytes` | pool_id | PG count / logical (data) bytes per pool; do not mix with raw-capacity metrics like `ceph_osd_stat_bytes_used` |
|
||||
| `ceph_scrub_overdue_pgs` / `ceph_scrub_overdue_bytes` | pool_id, depth | PGs / bytes whose stamp is older than the target interval; PGs with unparsable stamps count here |
|
||||
| `ceph_scrub_age_seconds` | pool_id, depth | histogram of scrub age weighted by bytes (buckets 1d..49d); `_sum/_count` gives mean data age, `histogram_quantile` the age of the Nth-percentile byte |
|
||||
| `ceph_scrub_schedule_pgs` | state | PGs by scrub_schedule state (scheduled, queued, scrubbing, blocked, reserving, none, other) |
|
||||
| `ceph_scrub_target_interval_seconds` | pool_id, depth | the interval each pool's overdue numbers were judged against (cluster-read, or the pin) |
|
||||
| `ceph_scrub_collect_success` / `_duration_seconds` / `_timestamp_seconds`, `ceph_scrub_parse_errors` | | collection health; alert on success == 0, or on a timestamp older than about three refresh intervals |
|
||||
|
||||
Pool names come from joining `ceph_pool_metadata` (mgr module) on `pool_id`.
|
||||
Several monk instances may be scraped for availability; dashboards dedup with
|
||||
`max by (pool_id)`.
|
||||
|
||||
```
|
||||
# deep-scrub cycle coverage
|
||||
1 - sum(ceph_scrub_overdue_bytes{depth="deep"}) / sum(ceph_scrub_pool_bytes)
|
||||
|
||||
# work outstanding, bytes
|
||||
sum(ceph_scrub_overdue_bytes{depth="deep"})
|
||||
```
|
||||
@@ -0,0 +1,24 @@
|
||||
module monk
|
||||
|
||||
go 1.27.0
|
||||
|
||||
require (
|
||||
github.com/prometheus/client_golang v1.23.2
|
||||
github.com/rs/zerolog v1.35.1
|
||||
)
|
||||
|
||||
require (
|
||||
github.com/beorn7/perks v1.0.1 // indirect
|
||||
github.com/cespare/xxhash/v2 v2.3.0 // indirect
|
||||
github.com/kr/text v0.2.0 // indirect
|
||||
github.com/kylelemons/godebug v1.1.0 // indirect
|
||||
github.com/mattn/go-colorable v0.1.14 // indirect
|
||||
github.com/mattn/go-isatty v0.0.20 // indirect
|
||||
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
|
||||
github.com/prometheus/client_model v0.6.2 // indirect
|
||||
github.com/prometheus/common v0.66.1 // indirect
|
||||
github.com/prometheus/procfs v0.16.1 // indirect
|
||||
go.yaml.in/yaml/v2 v2.4.2 // indirect
|
||||
golang.org/x/sys v0.35.0 // indirect
|
||||
google.golang.org/protobuf v1.36.8 // indirect
|
||||
)
|
||||
@@ -0,0 +1,53 @@
|
||||
github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
|
||||
github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw=
|
||||
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
|
||||
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
|
||||
github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E=
|
||||
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
|
||||
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||
github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8=
|
||||
github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU=
|
||||
github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo=
|
||||
github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ=
|
||||
github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE=
|
||||
github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk=
|
||||
github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
|
||||
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
|
||||
github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc=
|
||||
github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw=
|
||||
github.com/mattn/go-colorable v0.1.14 h1:9A9LHSqF/7dyVVX6g0U9cwm9pG3kP9gSzcuIPHPsaIE=
|
||||
github.com/mattn/go-colorable v0.1.14/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8=
|
||||
github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY=
|
||||
github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y=
|
||||
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA=
|
||||
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ=
|
||||
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
|
||||
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
|
||||
github.com/prometheus/client_golang v1.23.2 h1:Je96obch5RDVy3FDMndoUsjAhG5Edi49h0RJWRi/o0o=
|
||||
github.com/prometheus/client_golang v1.23.2/go.mod h1:Tb1a6LWHB3/SPIzCoaDXI4I8UHKeFTEQ1YCr+0Gyqmg=
|
||||
github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk=
|
||||
github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE=
|
||||
github.com/prometheus/common v0.66.1 h1:h5E0h5/Y8niHc5DlaLlWLArTQI7tMrsfQjHV+d9ZoGs=
|
||||
github.com/prometheus/common v0.66.1/go.mod h1:gcaUsgf3KfRSwHY4dIMXLPV0K/Wg1oZ8+SbZk/HH/dA=
|
||||
github.com/prometheus/procfs v0.16.1 h1:hZ15bTNuirocR6u0JZ6BAHHmwS1p8B4P6MRqxtzMyRg=
|
||||
github.com/prometheus/procfs v0.16.1/go.mod h1:teAbpZRB1iIAJYREa1LsoWUXykVXA1KlTmWl8x/U+Is=
|
||||
github.com/rogpeppe/go-internal v1.10.0 h1:TMyTOH3F/DB16zRVcYyreMH6GnZZrwQVAoYjRBZyWFQ=
|
||||
github.com/rogpeppe/go-internal v1.10.0/go.mod h1:UQnix2H7Ngw/k4C5ijL5+65zddjncjaFoBhdsK/akog=
|
||||
github.com/rs/zerolog v1.35.1 h1:m7xQeoiLIiV0BCEY4Hs+j2NG4Gp2o2KPKmhnnLiazKI=
|
||||
github.com/rs/zerolog v1.35.1/go.mod h1:EjML9kdfa/RMA7h/6z6pYmq1ykOuA8/mjWaEvGI+jcw=
|
||||
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
|
||||
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
|
||||
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
|
||||
go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
|
||||
go.yaml.in/yaml/v2 v2.4.2 h1:DzmwEr2rDGHl7lsFgAHxmNz/1NlQ7xLIrlN2h5d1eGI=
|
||||
go.yaml.in/yaml/v2 v2.4.2/go.mod h1:081UH+NErpNdqlCXm3TtEran0rJZGxAYx9hb/ELlsPU=
|
||||
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.35.0 h1:vz1N37gP5bs89s7He8XuIYXpyY0+QlsKmzipCbUtyxI=
|
||||
golang.org/x/sys v0.35.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k=
|
||||
google.golang.org/protobuf v1.36.8 h1:xHScyCOEuuwZEc6UtSOvPbAT4zRh0xcNRYekJwfqyMc=
|
||||
google.golang.org/protobuf v1.36.8/go.mod h1:fuxRtAxBytpl4zzqUh6/eyUujkJdNiuEkXntxiD/uRU=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk=
|
||||
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q=
|
||||
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
|
||||
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||
@@ -0,0 +1,266 @@
|
||||
package collector
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os/exec"
|
||||
"slices"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Ceph stamps carry +0000 rather than Z, and microsecond precision.
|
||||
const stampLayout = "2006-01-02T15:04:05.999999-0700"
|
||||
|
||||
var AgeBuckets = []time.Duration{
|
||||
24 * time.Hour,
|
||||
3 * 24 * time.Hour,
|
||||
7 * 24 * time.Hour,
|
||||
14 * 24 * time.Hour,
|
||||
21 * 24 * time.Hour,
|
||||
28 * 24 * time.Hour,
|
||||
35 * 24 * time.Hour,
|
||||
49 * 24 * time.Hour,
|
||||
}
|
||||
|
||||
type Depth string
|
||||
|
||||
const (
|
||||
Shallow Depth = "shallow"
|
||||
Deep Depth = "deep"
|
||||
)
|
||||
|
||||
var Depths = []Depth{Shallow, Deep}
|
||||
|
||||
type Intervals struct {
|
||||
Global map[Depth]time.Duration
|
||||
PerPool map[string]map[Depth]time.Duration
|
||||
}
|
||||
|
||||
func (iv Intervals) For(pool string, d Depth) time.Duration {
|
||||
if overrides, ok := iv.PerPool[pool]; ok {
|
||||
if v, ok := overrides[d]; ok && v > 0 {
|
||||
return v
|
||||
}
|
||||
}
|
||||
return iv.Global[d]
|
||||
}
|
||||
|
||||
var intervalOptions = map[Depth]string{
|
||||
Shallow: "osd_scrub_max_interval",
|
||||
Deep: "osd_deep_scrub_interval",
|
||||
}
|
||||
|
||||
type pgStat struct {
|
||||
PGID string `json:"pgid"`
|
||||
State string `json:"state"`
|
||||
LastScrubStamp string `json:"last_scrub_stamp"`
|
||||
LastDeepScrubStamp string `json:"last_deep_scrub_stamp"`
|
||||
ScrubSchedule string `json:"scrub_schedule"`
|
||||
StatSum struct {
|
||||
NumBytes int64 `json:"num_bytes"`
|
||||
} `json:"stat_sum"`
|
||||
}
|
||||
|
||||
type pgLs struct {
|
||||
PGStats []pgStat `json:"pg_stats"`
|
||||
}
|
||||
|
||||
type poolDetail struct {
|
||||
PoolID int64 `json:"pool_id"`
|
||||
Options struct {
|
||||
ScrubMaxInterval float64 `json:"scrub_max_interval"`
|
||||
DeepScrubInterval float64 `json:"deep_scrub_interval"`
|
||||
} `json:"options"`
|
||||
}
|
||||
|
||||
type PoolStats struct {
|
||||
PGs int
|
||||
Bytes int64
|
||||
Interval map[Depth]time.Duration
|
||||
OldestStamp map[Depth]time.Time
|
||||
OverduePGs map[Depth]int
|
||||
OverdueBytes map[Depth]int64
|
||||
// The age histogram observes each stored byte at its PG's scrub age; PGs
|
||||
// with unparsable stamps are excluded, and AgeBucketBytes is indexed like
|
||||
// AgeBuckets.
|
||||
ParsedBytes map[Depth]int64
|
||||
AgeSum map[Depth]float64
|
||||
AgeBucketBytes map[Depth][]int64
|
||||
}
|
||||
|
||||
type Snapshot struct {
|
||||
Taken time.Time
|
||||
Pools map[string]*PoolStats
|
||||
ScheduleStates map[string]int
|
||||
ParseErrors int
|
||||
}
|
||||
|
||||
func cephOutput(ctx context.Context, cephCmd []string, args ...string) ([]byte, error) {
|
||||
full := slices.Concat(cephCmd[1:], args)
|
||||
cmd := exec.CommandContext(ctx, cephCmd[0], full...)
|
||||
// A wrapper cephCmd (cephadm shell) leaves a grandchild holding stdout past
|
||||
// the context kill; WaitDelay lets Output return anyway.
|
||||
cmd.WaitDelay = time.Second
|
||||
out, err := cmd.Output()
|
||||
if err != nil {
|
||||
if ee, ok := err.(*exec.ExitError); ok {
|
||||
return nil, fmt.Errorf("%s %s: %w: %s", cephCmd[0], args[0], err, strings.TrimSpace(string(ee.Stderr)))
|
||||
}
|
||||
return nil, fmt.Errorf("%s %s: %w", cephCmd[0], args[0], err)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func Fetch(ctx context.Context, cephCmd []string) ([]byte, error) {
|
||||
return cephOutput(ctx, cephCmd, "pg", "ls", "-f", "json")
|
||||
}
|
||||
|
||||
// FetchIntervals reads the overdue policy from the cluster itself: the osd
|
||||
// section's effective intervals plus per-pool option overrides, so monk's
|
||||
// thresholds cannot drift from what the scrub scheduler actually targets.
|
||||
func FetchIntervals(ctx context.Context, cephCmd []string) (Intervals, error) {
|
||||
iv := Intervals{Global: map[Depth]time.Duration{}, PerPool: map[string]map[Depth]time.Duration{}}
|
||||
for depth, option := range intervalOptions {
|
||||
out, err := cephOutput(ctx, cephCmd, "config", "get", "osd", option)
|
||||
if err != nil {
|
||||
return Intervals{}, err
|
||||
}
|
||||
d, err := parseIntervalSeconds(string(out))
|
||||
if err != nil {
|
||||
return Intervals{}, fmt.Errorf("%s: %w", option, err)
|
||||
}
|
||||
iv.Global[depth] = d
|
||||
}
|
||||
out, err := cephOutput(ctx, cephCmd, "osd", "pool", "ls", "detail", "-f", "json")
|
||||
if err != nil {
|
||||
return Intervals{}, err
|
||||
}
|
||||
perPool, err := parsePoolIntervals(out)
|
||||
if err != nil {
|
||||
return Intervals{}, err
|
||||
}
|
||||
iv.PerPool = perPool
|
||||
return iv, nil
|
||||
}
|
||||
|
||||
func parseIntervalSeconds(s string) (time.Duration, error) {
|
||||
secs, err := strconv.ParseFloat(strings.TrimSpace(s), 64)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("parse interval %q: %w", s, err)
|
||||
}
|
||||
if secs <= 0 {
|
||||
return 0, fmt.Errorf("interval %q is not positive", s)
|
||||
}
|
||||
return time.Duration(secs * float64(time.Second)), nil
|
||||
}
|
||||
|
||||
func parsePoolIntervals(raw []byte) (map[string]map[Depth]time.Duration, error) {
|
||||
var pools []poolDetail
|
||||
if err := json.Unmarshal(raw, &pools); err != nil {
|
||||
return nil, fmt.Errorf("parse pool ls detail: %w", err)
|
||||
}
|
||||
perPool := map[string]map[Depth]time.Duration{}
|
||||
for _, p := range pools {
|
||||
overrides := map[Depth]time.Duration{}
|
||||
if p.Options.ScrubMaxInterval > 0 {
|
||||
overrides[Shallow] = time.Duration(p.Options.ScrubMaxInterval * float64(time.Second))
|
||||
}
|
||||
if p.Options.DeepScrubInterval > 0 {
|
||||
overrides[Deep] = time.Duration(p.Options.DeepScrubInterval * float64(time.Second))
|
||||
}
|
||||
if len(overrides) > 0 {
|
||||
perPool[strconv.FormatInt(p.PoolID, 10)] = overrides
|
||||
}
|
||||
}
|
||||
return perPool, nil
|
||||
}
|
||||
|
||||
func Compute(raw []byte, now time.Time, intervals Intervals) (*Snapshot, error) {
|
||||
var data pgLs
|
||||
if err := json.Unmarshal(raw, &data); err != nil {
|
||||
return nil, fmt.Errorf("parse pg ls: %w", err)
|
||||
}
|
||||
if len(data.PGStats) == 0 {
|
||||
return nil, fmt.Errorf("pg ls returned no pg_stats")
|
||||
}
|
||||
|
||||
snap := &Snapshot{
|
||||
Taken: now,
|
||||
Pools: map[string]*PoolStats{},
|
||||
ScheduleStates: map[string]int{},
|
||||
}
|
||||
for _, pg := range data.PGStats {
|
||||
pool, _, ok := strings.Cut(pg.PGID, ".")
|
||||
if !ok {
|
||||
snap.ParseErrors++
|
||||
continue
|
||||
}
|
||||
ps := snap.Pools[pool]
|
||||
if ps == nil {
|
||||
ps = &PoolStats{
|
||||
Interval: map[Depth]time.Duration{
|
||||
Shallow: intervals.For(pool, Shallow),
|
||||
Deep: intervals.For(pool, Deep),
|
||||
},
|
||||
OldestStamp: map[Depth]time.Time{},
|
||||
OverduePGs: map[Depth]int{},
|
||||
OverdueBytes: map[Depth]int64{},
|
||||
ParsedBytes: map[Depth]int64{},
|
||||
AgeSum: map[Depth]float64{},
|
||||
AgeBucketBytes: map[Depth][]int64{Shallow: make([]int64, len(AgeBuckets)), Deep: make([]int64, len(AgeBuckets))},
|
||||
}
|
||||
snap.Pools[pool] = ps
|
||||
}
|
||||
ps.PGs++
|
||||
ps.Bytes += pg.StatSum.NumBytes
|
||||
snap.ScheduleStates[scheduleState(pg.ScrubSchedule)]++
|
||||
|
||||
for depth, stampStr := range map[Depth]string{Shallow: pg.LastScrubStamp, Deep: pg.LastDeepScrubStamp} {
|
||||
stamp, err := time.Parse(stampLayout, stampStr)
|
||||
if err != nil {
|
||||
snap.ParseErrors++
|
||||
ps.OverduePGs[depth]++
|
||||
ps.OverdueBytes[depth] += pg.StatSum.NumBytes
|
||||
continue
|
||||
}
|
||||
if old, ok := ps.OldestStamp[depth]; !ok || stamp.Before(old) {
|
||||
ps.OldestStamp[depth] = stamp
|
||||
}
|
||||
age := now.Sub(stamp)
|
||||
if age > ps.Interval[depth] {
|
||||
ps.OverduePGs[depth]++
|
||||
ps.OverdueBytes[depth] += pg.StatSum.NumBytes
|
||||
}
|
||||
ps.ParsedBytes[depth] += pg.StatSum.NumBytes
|
||||
ps.AgeSum[depth] += age.Seconds() * float64(pg.StatSum.NumBytes)
|
||||
for i, le := range AgeBuckets {
|
||||
if age <= le {
|
||||
ps.AgeBucketBytes[depth][i] += pg.StatSum.NumBytes
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return snap, nil
|
||||
}
|
||||
|
||||
func scheduleState(s string) string {
|
||||
switch {
|
||||
case s == "" || s == "--" || strings.Contains(s, "no scrub"):
|
||||
return "none"
|
||||
case strings.Contains(s, "scheduled @"):
|
||||
return "scheduled"
|
||||
case strings.HasPrefix(s, "queued"):
|
||||
return "queued"
|
||||
case strings.Contains(s, "scrubbing"):
|
||||
return "scrubbing"
|
||||
case strings.HasPrefix(s, "Blocked"):
|
||||
return "blocked"
|
||||
case strings.HasPrefix(s, "Reserving"):
|
||||
return "reserving"
|
||||
default:
|
||||
return "other"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,285 @@
|
||||
package collector
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/prometheus/client_golang/prometheus/testutil"
|
||||
)
|
||||
|
||||
var testIntervals = Intervals{Global: map[Depth]time.Duration{
|
||||
Shallow: 7 * 24 * time.Hour,
|
||||
Deep: 28 * 24 * time.Hour,
|
||||
}}
|
||||
|
||||
func testNow(t *testing.T) time.Time {
|
||||
t.Helper()
|
||||
now, err := time.Parse(time.RFC3339, "2026-09-01T00:00:00Z")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return now
|
||||
}
|
||||
|
||||
func loadFixture(t *testing.T) []byte {
|
||||
t.Helper()
|
||||
raw, err := os.ReadFile("testdata/pg_ls.json")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return raw
|
||||
}
|
||||
|
||||
func TestComputeAgainstBruteForce(t *testing.T) {
|
||||
raw := loadFixture(t)
|
||||
now := testNow(t)
|
||||
snap, err := Compute(raw, now, testIntervals)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if snap.ParseErrors != 0 {
|
||||
t.Fatalf("parse errors on fixture: %d", snap.ParseErrors)
|
||||
}
|
||||
|
||||
var data pgLs
|
||||
if err := json.Unmarshal(raw, &data); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
totalPGs, totalBytes := 0, int64(0)
|
||||
overdue := map[Depth]int{}
|
||||
for _, pg := range data.PGStats {
|
||||
totalPGs++
|
||||
totalBytes += pg.StatSum.NumBytes
|
||||
for depth, s := range map[Depth]string{Shallow: pg.LastScrubStamp, Deep: pg.LastDeepScrubStamp} {
|
||||
stamp, err := time.Parse(stampLayout, s)
|
||||
if err != nil {
|
||||
t.Fatalf("fixture stamp %q: %v", s, err)
|
||||
}
|
||||
if now.Sub(stamp) > testIntervals.Global[depth] {
|
||||
overdue[depth]++
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
gotPGs, gotBytes := 0, int64(0)
|
||||
gotOverdue := map[Depth]int{}
|
||||
for _, ps := range snap.Pools {
|
||||
gotPGs += ps.PGs
|
||||
gotBytes += ps.Bytes
|
||||
for _, d := range Depths {
|
||||
gotOverdue[d] += ps.OverduePGs[d]
|
||||
}
|
||||
}
|
||||
if gotPGs != totalPGs || gotBytes != totalBytes {
|
||||
t.Errorf("totals: got %d PGs / %d bytes, want %d / %d", gotPGs, gotBytes, totalPGs, totalBytes)
|
||||
}
|
||||
for _, d := range Depths {
|
||||
if gotOverdue[d] != overdue[d] {
|
||||
t.Errorf("overdue[%s]: got %d, want %d", d, gotOverdue[d], overdue[d])
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgeBucketsCumulative(t *testing.T) {
|
||||
snap, err := Compute(loadFixture(t), testNow(t), testIntervals)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for pool, ps := range snap.Pools {
|
||||
for _, d := range Depths {
|
||||
buckets := ps.AgeBucketBytes[d]
|
||||
for i := 1; i < len(buckets); i++ {
|
||||
if buckets[i] < buckets[i-1] {
|
||||
t.Errorf("pool %s %s: bucket %d (%d) < bucket %d (%d)", pool, d, i, buckets[i], i-1, buckets[i-1])
|
||||
}
|
||||
}
|
||||
if last := buckets[len(buckets)-1]; last > ps.Bytes {
|
||||
t.Errorf("pool %s %s: largest bucket %d exceeds pool bytes %d", pool, d, last, ps.Bytes)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestScheduleState(t *testing.T) {
|
||||
cases := map[string]string{
|
||||
"periodic scrub scheduled @ 2026-08-29T02:13:18.699625+0000": "scheduled",
|
||||
"periodic deep scrub scheduled @ 2026-09-12T01:00:00.000000+0000": "scheduled",
|
||||
"queued for deep scrub": "queued",
|
||||
"deep scrubbing for 123s": "scrubbing",
|
||||
"Blocked! locked objects (for 5s)": "blocked",
|
||||
"Reserving. Waiting 3s for OSD.12 (2/20)": "reserving",
|
||||
"no scrub is scheduled": "none",
|
||||
"--": "none",
|
||||
"": "none",
|
||||
"user requested, deferred until 2026-09-01": "other",
|
||||
}
|
||||
for in, want := range cases {
|
||||
if got := scheduleState(in); got != want {
|
||||
t.Errorf("scheduleState(%q) = %q, want %q", in, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestComputeRejectsEmpty(t *testing.T) {
|
||||
if _, err := Compute([]byte(`{"pg_stats": []}`), testNow(t), testIntervals); err == nil {
|
||||
t.Error("empty pg_stats should error, not report zero work outstanding")
|
||||
}
|
||||
}
|
||||
|
||||
func TestComputeCountsUnparsableStampsOverdue(t *testing.T) {
|
||||
now := testNow(t)
|
||||
goodStamp, err := time.Parse(stampLayout, "2026-08-31T00:00:00.000000+0000")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
raw := []byte(`{"pg_stats": [
|
||||
{"pgid": "7.a", "last_scrub_stamp": "2026-08-31T00:00:00.000000+0000", "last_deep_scrub_stamp": "2026-08-31T00:00:00.000000+0000", "stat_sum": {"num_bytes": 100}},
|
||||
{"pgid": "7.b", "last_scrub_stamp": "", "last_deep_scrub_stamp": "", "stat_sum": {"num_bytes": 40}}
|
||||
]}`)
|
||||
snap, err := Compute(raw, now, testIntervals)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if snap.ParseErrors != 2 {
|
||||
t.Errorf("ParseErrors: got %d, want 2", snap.ParseErrors)
|
||||
}
|
||||
ps := snap.Pools["7"]
|
||||
if ps == nil {
|
||||
t.Fatal("pool 7 missing")
|
||||
}
|
||||
if ps.PGs != 2 || ps.Bytes != 140 {
|
||||
t.Errorf("pool totals: got %d PGs / %d bytes, want 2 / 140", ps.PGs, ps.Bytes)
|
||||
}
|
||||
for _, d := range Depths {
|
||||
if ps.OverduePGs[d] != 1 {
|
||||
t.Errorf("OverduePGs[%s]: got %d, want 1", d, ps.OverduePGs[d])
|
||||
}
|
||||
if ps.OverdueBytes[d] != 40 {
|
||||
t.Errorf("OverdueBytes[%s]: got %d, want 40", d, ps.OverdueBytes[d])
|
||||
}
|
||||
if !ps.OldestStamp[d].Equal(goodStamp) {
|
||||
t.Errorf("OldestStamp[%s]: got %v, want %v", d, ps.OldestStamp[d], goodStamp)
|
||||
}
|
||||
for i, b := range ps.AgeBucketBytes[d] {
|
||||
if b != 100 {
|
||||
t.Errorf("AgeBucketBytes[%s][%d]: got %d, want 100", d, i, b)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestExporterMetricSurface(t *testing.T) {
|
||||
now := testNow(t)
|
||||
deepStamp, err := time.Parse(stampLayout, "2026-07-01T00:00:00.000000+0000")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
raw := []byte(`{"pg_stats": [
|
||||
{"pgid": "7.a", "last_scrub_stamp": "2026-08-31T00:00:00.000000+0000", "last_deep_scrub_stamp": "2026-07-01T00:00:00.000000+0000", "scrub_schedule": "periodic scrub scheduled @ 2026-09-02T00:00:00.000000+0000", "stat_sum": {"num_bytes": 100}}
|
||||
]}`)
|
||||
snap, err := Compute(raw, now, testIntervals)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
exporter := &Exporter{}
|
||||
exporter.Store(snap, time.Second)
|
||||
|
||||
if problems, err := testutil.CollectAndLint(exporter); err != nil || len(problems) > 0 {
|
||||
t.Fatalf("lint: %v %v", problems, err)
|
||||
}
|
||||
|
||||
expected := fmt.Sprintf(`
|
||||
# HELP ceph_pg_last_deep_scrub_stamp Oldest per-PG last_deep_scrub_stamp in the pool (seconds since epoch)
|
||||
# TYPE ceph_pg_last_deep_scrub_stamp gauge
|
||||
ceph_pg_last_deep_scrub_stamp{pool_id="7"} %g
|
||||
# HELP ceph_scrub_collect_success Whether the last pg ls collection succeeded
|
||||
# TYPE ceph_scrub_collect_success gauge
|
||||
ceph_scrub_collect_success 1
|
||||
# HELP ceph_scrub_overdue_bytes Bytes in PGs whose last scrub at this depth is older than the target interval
|
||||
# TYPE ceph_scrub_overdue_bytes gauge
|
||||
ceph_scrub_overdue_bytes{depth="deep",pool_id="7"} 100
|
||||
ceph_scrub_overdue_bytes{depth="shallow",pool_id="7"} 0
|
||||
# HELP ceph_scrub_target_interval_seconds Scrub target interval the pool's overdue numbers were judged against
|
||||
# TYPE ceph_scrub_target_interval_seconds gauge
|
||||
ceph_scrub_target_interval_seconds{depth="deep",pool_id="7"} 2.4192e+06
|
||||
ceph_scrub_target_interval_seconds{depth="shallow",pool_id="7"} 604800
|
||||
`, float64(deepStamp.UnixMicro())/1e6)
|
||||
err = testutil.CollectAndCompare(exporter, strings.NewReader(expected),
|
||||
"ceph_pg_last_deep_scrub_stamp", "ceph_scrub_collect_success",
|
||||
"ceph_scrub_overdue_bytes", "ceph_scrub_target_interval_seconds")
|
||||
if err != nil {
|
||||
t.Error(err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseIntervalSeconds(t *testing.T) {
|
||||
d, err := parseIntervalSeconds("2419200.000000\n")
|
||||
if err != nil || d != 28*24*time.Hour {
|
||||
t.Errorf("deep: got %v %v, want 672h", d, err)
|
||||
}
|
||||
d, err = parseIntervalSeconds("604800.000000\n")
|
||||
if err != nil || d != 7*24*time.Hour {
|
||||
t.Errorf("shallow: got %v %v, want 168h", d, err)
|
||||
}
|
||||
for _, bad := range []string{"", "abc", "0.000000", "-1"} {
|
||||
if _, err := parseIntervalSeconds(bad); err == nil {
|
||||
t.Errorf("parseIntervalSeconds(%q) should error", bad)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestParsePoolIntervals(t *testing.T) {
|
||||
raw := []byte(`[
|
||||
{"pool_id": 2, "pool_name": "data", "options": {}},
|
||||
{"pool_id": 7, "pool_name": "ctl", "options": {"deep_scrub_interval": 1209600.0}},
|
||||
{"pool_id": 8, "pool_name": "meta", "options": {"scrub_max_interval": 86400.0, "deep_scrub_interval": 0}}
|
||||
]`)
|
||||
perPool, err := parsePoolIntervals(raw)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, ok := perPool["2"]; ok {
|
||||
t.Error("pool 2 has no overrides but appeared")
|
||||
}
|
||||
if got := perPool["7"][Deep]; got != 14*24*time.Hour {
|
||||
t.Errorf("pool 7 deep override: got %v, want 336h", got)
|
||||
}
|
||||
if _, ok := perPool["7"][Shallow]; ok {
|
||||
t.Error("pool 7 has no shallow override but one appeared")
|
||||
}
|
||||
if got := perPool["8"][Shallow]; got != 24*time.Hour {
|
||||
t.Errorf("pool 8 shallow override: got %v, want 24h", got)
|
||||
}
|
||||
if _, ok := perPool["8"][Deep]; ok {
|
||||
t.Error("pool 8 deep override is 0 (unset) but appeared")
|
||||
}
|
||||
}
|
||||
|
||||
func TestComputeAppliesPoolOverride(t *testing.T) {
|
||||
now := testNow(t)
|
||||
iv := Intervals{
|
||||
Global: map[Depth]time.Duration{Shallow: 7 * 24 * time.Hour, Deep: 28 * 24 * time.Hour},
|
||||
PerPool: map[string]map[Depth]time.Duration{"7": {Deep: 24 * time.Hour}},
|
||||
}
|
||||
raw := []byte(`{"pg_stats": [
|
||||
{"pgid": "7.a", "last_scrub_stamp": "2026-08-29T00:00:00.000000+0000", "last_deep_scrub_stamp": "2026-08-29T00:00:00.000000+0000", "stat_sum": {"num_bytes": 100}},
|
||||
{"pgid": "2.a", "last_scrub_stamp": "2026-08-29T00:00:00.000000+0000", "last_deep_scrub_stamp": "2026-08-29T00:00:00.000000+0000", "stat_sum": {"num_bytes": 100}}
|
||||
]}`)
|
||||
snap, err := Compute(raw, now, iv)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := snap.Pools["7"].OverduePGs[Deep]; got != 1 {
|
||||
t.Errorf("pool 7 deep overdue under 24h override: got %d, want 1 (stamp is 3d old)", got)
|
||||
}
|
||||
if got := snap.Pools["2"].OverduePGs[Deep]; got != 0 {
|
||||
t.Errorf("pool 2 deep overdue under 28d global: got %d, want 0", got)
|
||||
}
|
||||
if got := snap.Pools["7"].Interval[Deep]; got != 24*time.Hour {
|
||||
t.Errorf("pool 7 recorded interval: got %v, want 24h", got)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,87 @@
|
||||
package collector
|
||||
|
||||
import (
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
)
|
||||
|
||||
var (
|
||||
// Names track ceph PR #68925 so a future upstream implementation supersedes these.
|
||||
descLastScrub = prometheus.NewDesc("ceph_pg_last_scrub_stamp", "Oldest per-PG last_scrub_stamp in the pool (seconds since epoch)", []string{"pool_id"}, nil)
|
||||
descLastDeepScrub = prometheus.NewDesc("ceph_pg_last_deep_scrub_stamp", "Oldest per-PG last_deep_scrub_stamp in the pool (seconds since epoch)", []string{"pool_id"}, nil)
|
||||
|
||||
descPoolPGs = prometheus.NewDesc("ceph_scrub_pool_pgs", "PGs in the pool", []string{"pool_id"}, nil)
|
||||
descPoolBytes = prometheus.NewDesc("ceph_scrub_pool_bytes", "Stored bytes in the pool (sum of PG stat_sum.num_bytes)", []string{"pool_id"}, nil)
|
||||
descOverduePGs = prometheus.NewDesc("ceph_scrub_overdue_pgs", "PGs whose last scrub at this depth is older than the target interval", []string{"pool_id", "depth"}, nil)
|
||||
descOverdueBytes = prometheus.NewDesc("ceph_scrub_overdue_bytes", "Bytes in PGs whose last scrub at this depth is older than the target interval", []string{"pool_id", "depth"}, nil)
|
||||
descAgeHist = prometheus.NewDesc("ceph_scrub_age_seconds", "Scrub age distribution weighted by bytes: each stored byte observes its PG's age", []string{"pool_id", "depth"}, nil)
|
||||
descSchedule = prometheus.NewDesc("ceph_scrub_schedule_pgs", "PGs by scrub_schedule state", []string{"state"}, nil)
|
||||
descInterval = prometheus.NewDesc("ceph_scrub_target_interval_seconds", "Scrub target interval the pool's overdue numbers were judged against", []string{"pool_id", "depth"}, nil)
|
||||
|
||||
descSuccess = prometheus.NewDesc("ceph_scrub_collect_success", "Whether the last pg ls collection succeeded", nil, nil)
|
||||
descDuration = prometheus.NewDesc("ceph_scrub_collect_duration_seconds", "Duration of the last successful collection", nil, nil)
|
||||
descTimestamp = prometheus.NewDesc("ceph_scrub_collect_timestamp_seconds", "Time of the last successful collection", nil, nil)
|
||||
descParseErrors = prometheus.NewDesc("ceph_scrub_parse_errors", "Records skipped during the last successful collection", nil, nil)
|
||||
)
|
||||
|
||||
type Exporter struct {
|
||||
snapshot atomic.Pointer[Snapshot]
|
||||
lastDuration atomic.Int64
|
||||
failed atomic.Bool
|
||||
}
|
||||
|
||||
func (e *Exporter) Store(s *Snapshot, took time.Duration) {
|
||||
e.snapshot.Store(s)
|
||||
e.lastDuration.Store(int64(took))
|
||||
e.failed.Store(false)
|
||||
}
|
||||
|
||||
func (e *Exporter) MarkFailed() {
|
||||
e.failed.Store(true)
|
||||
}
|
||||
|
||||
func (e *Exporter) Describe(ch chan<- *prometheus.Desc) {
|
||||
prometheus.DescribeByCollect(e, ch)
|
||||
}
|
||||
|
||||
func (e *Exporter) Collect(ch chan<- prometheus.Metric) {
|
||||
success := 0.0
|
||||
if !e.failed.Load() {
|
||||
success = 1.0
|
||||
}
|
||||
ch <- prometheus.MustNewConstMetric(descSuccess, prometheus.GaugeValue, success)
|
||||
|
||||
snap := e.snapshot.Load()
|
||||
if snap == nil {
|
||||
return
|
||||
}
|
||||
ch <- prometheus.MustNewConstMetric(descDuration, prometheus.GaugeValue, time.Duration(e.lastDuration.Load()).Seconds())
|
||||
ch <- prometheus.MustNewConstMetric(descTimestamp, prometheus.GaugeValue, float64(snap.Taken.Unix()))
|
||||
ch <- prometheus.MustNewConstMetric(descParseErrors, prometheus.GaugeValue, float64(snap.ParseErrors))
|
||||
|
||||
for state, n := range snap.ScheduleStates {
|
||||
ch <- prometheus.MustNewConstMetric(descSchedule, prometheus.GaugeValue, float64(n), state)
|
||||
}
|
||||
for pool, ps := range snap.Pools {
|
||||
ch <- prometheus.MustNewConstMetric(descPoolPGs, prometheus.GaugeValue, float64(ps.PGs), pool)
|
||||
ch <- prometheus.MustNewConstMetric(descPoolBytes, prometheus.GaugeValue, float64(ps.Bytes), pool)
|
||||
if s, ok := ps.OldestStamp[Shallow]; ok {
|
||||
ch <- prometheus.MustNewConstMetric(descLastScrub, prometheus.GaugeValue, float64(s.UnixMicro())/1e6, pool)
|
||||
}
|
||||
if s, ok := ps.OldestStamp[Deep]; ok {
|
||||
ch <- prometheus.MustNewConstMetric(descLastDeepScrub, prometheus.GaugeValue, float64(s.UnixMicro())/1e6, pool)
|
||||
}
|
||||
for _, depth := range Depths {
|
||||
ch <- prometheus.MustNewConstMetric(descInterval, prometheus.GaugeValue, ps.Interval[depth].Seconds(), pool, string(depth))
|
||||
ch <- prometheus.MustNewConstMetric(descOverduePGs, prometheus.GaugeValue, float64(ps.OverduePGs[depth]), pool, string(depth))
|
||||
ch <- prometheus.MustNewConstMetric(descOverdueBytes, prometheus.GaugeValue, float64(ps.OverdueBytes[depth]), pool, string(depth))
|
||||
buckets := make(map[float64]uint64, len(AgeBuckets))
|
||||
for i, le := range AgeBuckets {
|
||||
buckets[le.Seconds()] = uint64(ps.AgeBucketBytes[depth][i])
|
||||
}
|
||||
ch <- prometheus.MustNewConstHistogram(descAgeHist, uint64(ps.ParsedBytes[depth]), ps.AgeSum[depth], buckets, pool, string(depth))
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,290 @@
|
||||
{
|
||||
"pg_stats": [
|
||||
{
|
||||
"pgid": "1.0",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T17:21:29.956752+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-21T10:11:56.049721+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 440,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T02:13:18.699625+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 1883701808,
|
||||
"num_objects": 451
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "2.0",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-28T10:46:31.252225+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-24T17:43:41.738507+0000",
|
||||
"last_scrub_duration": 200,
|
||||
"objects_scrubbed": 530196,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T16:47:21.152387+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 2212516228864,
|
||||
"num_objects": 530196
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "2.1",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T09:00:34.940514+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-17T05:30:11.743003+0000",
|
||||
"last_scrub_duration": 207,
|
||||
"objects_scrubbed": 530388,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-28T19:01:02.402388+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 2213451628745,
|
||||
"num_objects": 530387
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "2.2",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T22:18:10.315448+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-23T19:27:00.859775+0000",
|
||||
"last_scrub_duration": 199,
|
||||
"objects_scrubbed": 530267,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-28T22:54:53.156730+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 2212961245624,
|
||||
"num_objects": 530267
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "3.0",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-28T01:01:45.257141+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-17T21:53:16.465546+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 688,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T10:23:13.278334+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 0,
|
||||
"num_objects": 688
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "3.1",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-28T07:43:32.041802+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-17T22:02:35.040955+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 680,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T18:31:58.192529+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 0,
|
||||
"num_objects": 680
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "3.2",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T13:57:50.569609+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-16T00:07:09.531128+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 716,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-28T19:24:11.325451+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 0,
|
||||
"num_objects": 716
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "4.0",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T23:12:13.363017+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-20T07:05:29.182890+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 35,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T03:19:21.777444+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 2415,
|
||||
"num_objects": 35
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "4.1",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T17:45:45.186565+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-25T01:37:27.651221+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 38,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T01:26:30.517739+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 2622,
|
||||
"num_objects": 38
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "4.2",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-28T15:19:12.392950+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-14T11:41:09.135850+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 27,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T16:49:49.131901+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 1863,
|
||||
"num_objects": 27
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "5.0",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-28T16:32:32.292060+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-19T21:34:19.067197+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 3,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T17:14:26.177300+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 1188,
|
||||
"num_objects": 3
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "5.1",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T18:19:04.430717+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-24T12:12:19.113097+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 0,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-28T23:56:19.098763+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 0,
|
||||
"num_objects": 0
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "5.2",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T18:16:38.017516+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-16T04:09:15.999326+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 1,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T00:16:22.427448+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 46,
|
||||
"num_objects": 1
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "6.0",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T21:11:51.367181+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-18T23:25:57.472621+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 18,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T03:42:56.158465+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 40268876,
|
||||
"num_objects": 18
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "6.1",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T16:33:01.498452+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-14T19:00:51.641852+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 16,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-28T22:46:41.362714+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 1248,
|
||||
"num_objects": 16
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "6.2",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-28T06:47:32.082745+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-13T21:34:53.985785+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 16,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T08:29:18.978331+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 1554,
|
||||
"num_objects": 16
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "7.0",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T10:57:47.564210+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-14T14:09:28.660259+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 1,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-28T22:34:15.820830+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 0,
|
||||
"num_objects": 1
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "7.1",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T22:55:17.096136+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-23T05:24:28.263044+0000",
|
||||
"last_scrub_duration": 0,
|
||||
"objects_scrubbed": 0,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T05:40:41.091974+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 0,
|
||||
"num_objects": 0
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "7.2",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T12:08:56.770232+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-16T13:07:27.578722+0000",
|
||||
"last_scrub_duration": 0,
|
||||
"objects_scrubbed": 0,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-28T18:47:21.697049+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 0,
|
||||
"num_objects": 0
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "8.0",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T18:20:40.182316+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-16T07:44:47.720307+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 14,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T01:58:44.836369+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 3840,
|
||||
"num_objects": 14
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "8.1",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-27T19:09:19.772513+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-16T05:42:55.050154+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 14,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T00:38:47.801912+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 4943,
|
||||
"num_objects": 14
|
||||
}
|
||||
},
|
||||
{
|
||||
"pgid": "8.2",
|
||||
"state": "active+clean",
|
||||
"last_scrub_stamp": "2026-08-28T03:41:06.684834+0000",
|
||||
"last_deep_scrub_stamp": "2026-08-17T10:05:04.126426+0000",
|
||||
"last_scrub_duration": 1,
|
||||
"objects_scrubbed": 22,
|
||||
"scrub_schedule": "periodic scrub scheduled @ 2026-08-29T09:43:27.122526+0000",
|
||||
"stat_sum": {
|
||||
"num_bytes": 7273,
|
||||
"num_objects": 22
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -0,0 +1,5 @@
|
||||
// Package version holds the monorepo release version, stamped into this file
|
||||
// by release-please (extra-files in release-please-config.json).
|
||||
package version
|
||||
|
||||
const Version = "0.37.1" // x-release-please-version
|
||||
@@ -0,0 +1,106 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"flag"
|
||||
"net/http"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"github.com/prometheus/client_golang/prometheus/promhttp"
|
||||
"github.com/rs/zerolog"
|
||||
|
||||
"monk/internal/collector"
|
||||
"monk/internal/version"
|
||||
)
|
||||
|
||||
func main() {
|
||||
listen := flag.String("listen", ":9284", "metrics listen address")
|
||||
refresh := flag.Duration("refresh", 2*time.Minute, "refresh interval")
|
||||
timeout := flag.Duration("timeout", 90*time.Second, "ceph command timeout")
|
||||
shallowPin := flag.String("shallow-interval", "", "pin the shallow target interval (e.g. 168h); empty follows the cluster")
|
||||
deepPin := flag.String("deep-interval", "", "pin the deep target interval (e.g. 672h); empty follows the cluster")
|
||||
cephCmd := flag.String("ceph-cmd", "ceph", "command prefix to reach the ceph CLI, split on spaces")
|
||||
flag.Parse()
|
||||
|
||||
log := zerolog.New(os.Stderr).With().Timestamp().Logger()
|
||||
cmd := strings.Fields(*cephCmd)
|
||||
if len(cmd) == 0 {
|
||||
log.Fatal().Msg("empty -ceph-cmd")
|
||||
}
|
||||
pins := map[collector.Depth]time.Duration{}
|
||||
for depth, pin := range map[collector.Depth]string{collector.Shallow: *shallowPin, collector.Deep: *deepPin} {
|
||||
if pin == "" {
|
||||
continue
|
||||
}
|
||||
d, err := time.ParseDuration(pin)
|
||||
if err != nil || d <= 0 {
|
||||
log.Fatal().Str("value", pin).Msg("invalid interval pin")
|
||||
}
|
||||
pins[depth] = d
|
||||
}
|
||||
|
||||
// Until the first successful cluster read, unpinned depths fall back to
|
||||
// ceph's own defaults so a cold start with an unreachable mon still serves
|
||||
// sane thresholds; the target-interval metric shows what was used.
|
||||
intervals := collector.Intervals{Global: map[collector.Depth]time.Duration{
|
||||
collector.Shallow: 7 * 24 * time.Hour,
|
||||
collector.Deep: 28 * 24 * time.Hour,
|
||||
}}
|
||||
applyPins := func(iv collector.Intervals) collector.Intervals {
|
||||
for depth, d := range pins {
|
||||
iv.Global[depth] = d
|
||||
for _, overrides := range iv.PerPool {
|
||||
delete(overrides, depth)
|
||||
}
|
||||
}
|
||||
return iv
|
||||
}
|
||||
intervals = applyPins(intervals)
|
||||
|
||||
exporter := &collector.Exporter{}
|
||||
exporter.MarkFailed()
|
||||
registry := prometheus.NewRegistry()
|
||||
registry.MustRegister(exporter)
|
||||
|
||||
collect := func() {
|
||||
start := time.Now()
|
||||
ctx, cancel := context.WithTimeout(context.Background(), *timeout)
|
||||
defer cancel()
|
||||
if len(pins) < len(collector.Depths) {
|
||||
if iv, err := collector.FetchIntervals(ctx, cmd); err == nil {
|
||||
intervals = applyPins(iv)
|
||||
} else {
|
||||
log.Warn().Err(err).Msg("interval read failed, keeping previous targets")
|
||||
}
|
||||
}
|
||||
raw, err := collector.Fetch(ctx, cmd)
|
||||
if err == nil {
|
||||
var snap *collector.Snapshot
|
||||
if snap, err = collector.Compute(raw, start, intervals); err == nil {
|
||||
exporter.Store(snap, time.Since(start))
|
||||
log.Info().Dur("took", time.Since(start)).Int("pools", len(snap.Pools)).Int("parse_errors", snap.ParseErrors).Msg("collected")
|
||||
return
|
||||
}
|
||||
}
|
||||
exporter.MarkFailed()
|
||||
log.Error().Err(err).Msg("collection failed")
|
||||
}
|
||||
|
||||
go func() {
|
||||
collect()
|
||||
for range time.Tick(*refresh) {
|
||||
collect()
|
||||
}
|
||||
}()
|
||||
|
||||
mux := http.NewServeMux()
|
||||
mux.Handle("/metrics", promhttp.HandlerFor(registry, promhttp.HandlerOpts{}))
|
||||
server := &http.Server{Addr: *listen, Handler: mux, ReadHeaderTimeout: 5 * time.Second}
|
||||
log.Info().Str("version", version.Version).Str("listen", *listen).Msg("serving")
|
||||
if err := server.ListenAndServe(); err != nil {
|
||||
log.Fatal().Err(err).Msg("listen failed")
|
||||
}
|
||||
}
|
||||
@@ -75,6 +75,10 @@
|
||||
"type": "generic",
|
||||
"path": "packages/columbo/internal/version/version.go"
|
||||
},
|
||||
{
|
||||
"type": "generic",
|
||||
"path": "packages/monk/internal/version/version.go"
|
||||
},
|
||||
{
|
||||
"type": "generic",
|
||||
"path": "charts/apps/michael/Chart.yaml"
|
||||
|
||||
Reference in New Issue
Block a user