diff --git a/cmd/server/main.go b/cmd/server/main.go index 6aab376..bc20743 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -53,15 +53,16 @@ func run() error { disc := discovery.New(dockerCli, cfg.StacksRoot, cfg.OptInLabel) exec := updater.NewComposeExecutor() - queue := updater.NewQueue(exec) - queue.Start(context.Background()) - defer queue.Stop() reg := prometheus.NewRegistry() m := metrics.New(reg) m.BuildInfo.WithLabelValues(version, commit).Set(1) - handlers := api.NewHandlers(disc, &submitterAdapter{queue: queue, timeout: cfg.UpdateTimeout}, dockerCli, version, commit, buildTime) + queue := updater.NewQueue(exec, m) + queue.Start(context.Background()) + defer queue.Stop() + + handlers := api.NewHandlers(disc, &submitterAdapter{queue: queue, timeout: cfg.UpdateTimeout}, dockerCli, version, commit, buildTime, m) mux := http.NewServeMux() mux.HandleFunc("POST /update", handlers.Update) @@ -70,7 +71,7 @@ func run() error { mux.Handle("GET /metrics", metrics.Handler(reg)) authed := api.Auth(cfg.APIKey) - handler := api.RequestID(api.RequestLogger(logger)(routeAuth(mux, authed))) + handler := api.RequestID(api.RequestLogger(logger, m)(routeAuth(mux, authed))) srv := &http.Server{ Addr: ":" + cfg.Port, diff --git a/internal/api/handlers.go b/internal/api/handlers.go index b3392b8..cf251a3 100644 --- a/internal/api/handlers.go +++ b/internal/api/handlers.go @@ -8,6 +8,7 @@ import ( "github.com/docker/docker/api/types" "github.com/shcizo/package-updater/internal/discovery" "github.com/shcizo/package-updater/internal/logging" + "github.com/shcizo/package-updater/internal/metrics" "github.com/shcizo/package-updater/internal/updater" ) @@ -34,11 +35,13 @@ type Handlers struct { version string commit string buildTime string + metrics *metrics.Metrics } // NewHandlers constructs a Handlers value with all dependencies injected. -func NewHandlers(f Finder, s Submitter, p Pinger, version, commit, buildTime string) *Handlers { - return &Handlers{finder: f, submitter: s, pinger: p, version: version, commit: commit, buildTime: buildTime} +// m may be nil, in which case metrics recording is silently skipped. +func NewHandlers(f Finder, s Submitter, p Pinger, version, commit, buildTime string, m *metrics.Metrics) *Handlers { + return &Handlers{finder: f, submitter: s, pinger: p, version: version, commit: commit, buildTime: buildTime, metrics: m} } // Update implements POST /update. @@ -107,9 +110,15 @@ func (h *Handlers) Update(w http.ResponseWriter, r *http.Request) { // Healthz implements GET /healthz. func (h *Handlers) Healthz(w http.ResponseWriter, r *http.Request) { if _, err := h.pinger.Ping(r.Context()); err != nil { + if h.metrics != nil { + h.metrics.DockerPingUp.Set(0) + } writeJSON(w, http.StatusServiceUnavailable, HealthResponse{Status: "unhealthy", Docker: "unreachable"}) return } + if h.metrics != nil { + h.metrics.DockerPingUp.Set(1) + } writeJSON(w, http.StatusOK, HealthResponse{Status: "ok", Docker: "ok"}) } diff --git a/internal/api/handlers_test.go b/internal/api/handlers_test.go index 3de30e7..b62f830 100644 --- a/internal/api/handlers_test.go +++ b/internal/api/handlers_test.go @@ -56,7 +56,7 @@ func decode[T any](t *testing.T, body io.Reader) T { } func TestUpdate_ValidationError(t *testing.T) { - h := api.NewHandlers(&fakeFinder{}, &fakeSubmitter{}, &fakePinger{}, "v0.0.0", "abc", "now") + h := api.NewHandlers(&fakeFinder{}, &fakeSubmitter{}, &fakePinger{}, "v0.0.0", "abc", "now", nil) req := httptest.NewRequest(http.MethodPost, "/update", strings.NewReader(`{"tag":"v1.2.3"}`)) w := httptest.NewRecorder() @@ -65,7 +65,7 @@ func TestUpdate_ValidationError(t *testing.T) { } func TestUpdate_BadJSON(t *testing.T) { - h := api.NewHandlers(&fakeFinder{}, &fakeSubmitter{}, &fakePinger{}, "v0.0.0", "abc", "now") + h := api.NewHandlers(&fakeFinder{}, &fakeSubmitter{}, &fakePinger{}, "v0.0.0", "abc", "now", nil) req := httptest.NewRequest(http.MethodPost, "/update", strings.NewReader(`not json`)) w := httptest.NewRecorder() h.Update(w, req) @@ -74,7 +74,7 @@ func TestUpdate_BadJSON(t *testing.T) { func TestUpdate_DiscoveryFailureReturns500(t *testing.T) { finder := &fakeFinder{err: errors.New("daemon unreachable")} - h := api.NewHandlers(finder, &fakeSubmitter{}, &fakePinger{}, "v0.0.0", "abc", "now") + h := api.NewHandlers(finder, &fakeSubmitter{}, &fakePinger{}, "v0.0.0", "abc", "now", nil) body, _ := json.Marshal(api.UpdateRequest{Image: "r/x"}) req := httptest.NewRequest(http.MethodPost, "/update", bytes.NewReader(body)) w := httptest.NewRecorder() @@ -83,7 +83,7 @@ func TestUpdate_DiscoveryFailureReturns500(t *testing.T) { } func TestUpdate_ZeroMatchesReturns200(t *testing.T) { - h := api.NewHandlers(&fakeFinder{}, &fakeSubmitter{}, &fakePinger{}, "v0.0.0", "abc", "now") + h := api.NewHandlers(&fakeFinder{}, &fakeSubmitter{}, &fakePinger{}, "v0.0.0", "abc", "now", nil) body, _ := json.Marshal(api.UpdateRequest{Image: "r/x"}) req := httptest.NewRequest(http.MethodPost, "/update", bytes.NewReader(body)) w := httptest.NewRecorder() @@ -97,7 +97,7 @@ func TestUpdate_AllSucceeded200(t *testing.T) { finder := &fakeFinder{jobs: []discovery.Job{ {Project: "p", Service: "s", WorkingDir: "/x", ConfigFiles: []string{"/x/c.yml"}}, }} - h := api.NewHandlers(finder, &fakeSubmitter{}, &fakePinger{}, "v0.0.0", "abc", "now") + h := api.NewHandlers(finder, &fakeSubmitter{}, &fakePinger{}, "v0.0.0", "abc", "now", nil) body, _ := json.Marshal(api.UpdateRequest{Image: "r/x", Tag: "v1"}) req := httptest.NewRequest(http.MethodPost, "/update", bytes.NewReader(body)) w := httptest.NewRecorder() @@ -119,7 +119,7 @@ func TestUpdate_MixedReturns207(t *testing.T) { {Job: jobs[0], Status: updater.StatusUpdated}, {Job: jobs[1], Status: updater.StatusFailed, Error: "boom"}, }} - h := api.NewHandlers(finder, submitter, &fakePinger{}, "v0.0.0", "abc", "now") + h := api.NewHandlers(finder, submitter, &fakePinger{}, "v0.0.0", "abc", "now", nil) body, _ := json.Marshal(api.UpdateRequest{Image: "r/x"}) req := httptest.NewRequest(http.MethodPost, "/update", bytes.NewReader(body)) w := httptest.NewRecorder() @@ -133,7 +133,7 @@ func TestUpdate_AllFailedReturns500(t *testing.T) { submitter := &fakeSubmitter{results: []updater.Result{ {Job: jobs[0], Status: updater.StatusFailed, Error: "boom"}, }} - h := api.NewHandlers(finder, submitter, &fakePinger{}, "v0.0.0", "abc", "now") + h := api.NewHandlers(finder, submitter, &fakePinger{}, "v0.0.0", "abc", "now", nil) body, _ := json.Marshal(api.UpdateRequest{Image: "r/x"}) req := httptest.NewRequest(http.MethodPost, "/update", bytes.NewReader(body)) w := httptest.NewRecorder() @@ -142,7 +142,7 @@ func TestUpdate_AllFailedReturns500(t *testing.T) { } func TestHealthz_OKWhenDockerUp(t *testing.T) { - h := api.NewHandlers(&fakeFinder{}, &fakeSubmitter{}, &fakePinger{}, "v0.0.0", "abc", "now") + h := api.NewHandlers(&fakeFinder{}, &fakeSubmitter{}, &fakePinger{}, "v0.0.0", "abc", "now", nil) req := httptest.NewRequest(http.MethodGet, "/healthz", nil) w := httptest.NewRecorder() h.Healthz(w, req) @@ -150,7 +150,7 @@ func TestHealthz_OKWhenDockerUp(t *testing.T) { } func TestHealthz_503WhenDockerDown(t *testing.T) { - h := api.NewHandlers(&fakeFinder{}, &fakeSubmitter{}, &fakePinger{err: errors.New("ping fail")}, "v0.0.0", "abc", "now") + h := api.NewHandlers(&fakeFinder{}, &fakeSubmitter{}, &fakePinger{err: errors.New("ping fail")}, "v0.0.0", "abc", "now", nil) req := httptest.NewRequest(http.MethodGet, "/healthz", nil) w := httptest.NewRecorder() h.Healthz(w, req) @@ -158,7 +158,7 @@ func TestHealthz_503WhenDockerDown(t *testing.T) { } func TestVersion(t *testing.T) { - h := api.NewHandlers(&fakeFinder{}, &fakeSubmitter{}, &fakePinger{}, "v1.2.3", "abcdef", "2026-05-22T00:00:00Z") + h := api.NewHandlers(&fakeFinder{}, &fakeSubmitter{}, &fakePinger{}, "v1.2.3", "abcdef", "2026-05-22T00:00:00Z", nil) req := httptest.NewRequest(http.MethodGet, "/version", nil) w := httptest.NewRecorder() h.Version(w, req) diff --git a/internal/api/middleware.go b/internal/api/middleware.go index 0423bd4..8564890 100644 --- a/internal/api/middleware.go +++ b/internal/api/middleware.go @@ -7,10 +7,12 @@ import ( "fmt" "log/slog" "net/http" + "strconv" "strings" "time" "github.com/shcizo/package-updater/internal/logging" + "github.com/shcizo/package-updater/internal/metrics" ) // Auth returns middleware that requires a matching bearer token. @@ -66,8 +68,9 @@ func RequestID(next http.Handler) http.Handler { }) } -// RequestLogger emits a structured access-log line per request. -func RequestLogger(base *slog.Logger) func(http.Handler) http.Handler { +// RequestLogger emits a structured access-log line per request. m may be nil, +// in which case metrics recording is silently skipped. +func RequestLogger(base *slog.Logger, m *metrics.Metrics) func(http.Handler) http.Handler { return func(next http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { start := time.Now() @@ -81,6 +84,9 @@ func RequestLogger(base *slog.Logger) func(http.Handler) http.Handler { "duration_ms", time.Since(start).Milliseconds(), "client_ip", clientIP(r), ) + if m != nil { + m.HTTPRequests.WithLabelValues(r.URL.Path, strconv.Itoa(sw.status)).Inc() + } }) } } diff --git a/internal/api/middleware_test.go b/internal/api/middleware_test.go index d042ba3..e973aad 100644 --- a/internal/api/middleware_test.go +++ b/internal/api/middleware_test.go @@ -80,7 +80,7 @@ func TestRequestID_UsesIncoming(t *testing.T) { func TestRequestLogger_LogsAndDelegates(t *testing.T) { logger := newTestLogger() called := false - h := api.RequestLogger(logger)(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + h := api.RequestLogger(logger, nil)(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { called = true w.WriteHeader(http.StatusTeapot) })) diff --git a/internal/updater/queue.go b/internal/updater/queue.go index 9f02da6..c7b3a38 100644 --- a/internal/updater/queue.go +++ b/internal/updater/queue.go @@ -5,6 +5,7 @@ import ( "time" "github.com/shcizo/package-updater/internal/discovery" + "github.com/shcizo/package-updater/internal/metrics" ) // Status describes the outcome of executing a single Job. @@ -30,10 +31,11 @@ type Result struct { // It also satisfies the "single global FIFO worker" design choice in // spec section 5.7. type Queue struct { - exec Executor - ch chan submission - stop chan struct{} - done chan struct{} + exec Executor + metrics *metrics.Metrics + ch chan submission + stop chan struct{} + done chan struct{} } type submission struct { @@ -42,13 +44,15 @@ type submission struct { results chan []Result } -// NewQueue constructs a queue bound to the given executor. -func NewQueue(exec Executor) *Queue { +// NewQueue constructs a queue bound to the given executor. m may be nil, +// in which case metrics recording is silently skipped. +func NewQueue(exec Executor, m *metrics.Metrics) *Queue { return &Queue{ - exec: exec, - ch: make(chan submission, 16), - stop: make(chan struct{}), - done: make(chan struct{}), + exec: exec, + metrics: m, + ch: make(chan submission, 16), + stop: make(chan struct{}), + done: make(chan struct{}), } } @@ -70,6 +74,9 @@ func (q *Queue) Stop() { func (q *Queue) Submit(ctx context.Context, jobs []discovery.Job) []Result { resCh := make(chan []Result, 1) q.ch <- submission{ctx: ctx, jobs: jobs, results: resCh} + if q.metrics != nil { + q.metrics.QueueDepth.Set(float64(len(q.ch))) + } return <-resCh } @@ -97,6 +104,10 @@ func (q *Queue) runOne(ctx context.Context, job discovery.Job) Result { r.Status = StatusRefused r.Error = job.RefusedReason r.DurationMs = time.Since(start).Milliseconds() + if q.metrics != nil { + q.metrics.UpdateJobs.WithLabelValues(job.Project, job.Service, string(StatusRefused)).Inc() + q.metrics.UpdateDuration.WithLabelValues(job.Project, job.Service).Observe(time.Since(start).Seconds()) + } return r } @@ -112,6 +123,15 @@ func (q *Queue) runOne(ctx context.Context, job discovery.Job) Result { r.Status = StatusFailed r.Error = err.Error() } + + if q.metrics != nil { + elapsed := time.Since(start).Seconds() + q.metrics.UpdateJobs.WithLabelValues(job.Project, job.Service, string(r.Status)).Inc() + q.metrics.UpdateDuration.WithLabelValues(job.Project, job.Service).Observe(elapsed) + if r.Status == StatusUpdated { + q.metrics.LastUpdateTime.WithLabelValues(job.Project, job.Service).SetToCurrentTime() + } + } return r } diff --git a/internal/updater/queue_test.go b/internal/updater/queue_test.go index d346a46..c8f99d8 100644 --- a/internal/updater/queue_test.go +++ b/internal/updater/queue_test.go @@ -7,7 +7,10 @@ import ( "testing" "time" + dto "github.com/prometheus/client_model/go" + "github.com/prometheus/client_golang/prometheus" "github.com/shcizo/package-updater/internal/discovery" + "github.com/shcizo/package-updater/internal/metrics" "github.com/shcizo/package-updater/internal/updater" "github.com/stretchr/testify/require" ) @@ -44,7 +47,7 @@ func (f *fakeExec) callsCopy() []discovery.Job { func TestQueue_ProcessesFIFO(t *testing.T) { exec := &fakeExec{delay: 20 * time.Millisecond} - q := updater.NewQueue(exec) + q := updater.NewQueue(exec, nil) q.Start(context.Background()) defer q.Stop() @@ -69,7 +72,7 @@ func TestQueue_ProcessesFIFO(t *testing.T) { func TestQueue_ReturnsPerJobResults(t *testing.T) { exec := &fakeExec{errForSvc: map[string]error{"failing": errors.New("boom")}} - q := updater.NewQueue(exec) + q := updater.NewQueue(exec, nil) q.Start(context.Background()) defer q.Stop() @@ -90,7 +93,7 @@ func TestQueue_ReturnsPerJobResults(t *testing.T) { func TestQueue_Timeout(t *testing.T) { exec := &fakeExec{delay: 200 * time.Millisecond} - q := updater.NewQueue(exec) + q := updater.NewQueue(exec, nil) q.Start(context.Background()) defer q.Stop() @@ -103,3 +106,24 @@ func TestQueue_Timeout(t *testing.T) { require.Len(t, results, 1) require.Equal(t, updater.StatusTimeout, results[0].Status) } + +func TestQueue_RecordsMetrics(t *testing.T) { + reg := prometheus.NewRegistry() + m := metrics.New(reg) + exec := &fakeExec{} + q := updater.NewQueue(exec, m) + q.Start(context.Background()) + defer q.Stop() + + q.Submit(context.Background(), []discovery.Job{{Project: "p", Service: "s"}}) + + // Verify UpdateJobs counter was incremented for the successful job. + metric := &dto.Metric{} + require.NoError(t, m.UpdateJobs.WithLabelValues("p", "s", "updated").Write(metric)) + require.Equal(t, 1.0, metric.GetCounter().GetValue()) + + // Verify LastUpdateTime was set (non-zero). + tsMetric := &dto.Metric{} + require.NoError(t, m.LastUpdateTime.WithLabelValues("p", "s").Write(tsMetric)) + require.Greater(t, tsMetric.GetGauge().GetValue(), 0.0) +} diff --git a/internal/updater/worker_test.go b/internal/updater/worker_test.go index 7dcc8d4..9522b8f 100644 --- a/internal/updater/worker_test.go +++ b/internal/updater/worker_test.go @@ -21,7 +21,7 @@ func (c *counterExec) Execute(_ context.Context, _ discovery.Job) error { func TestWorker_RunsExactlyOneAtATime(t *testing.T) { exec := &counterExec{} - q := updater.NewQueue(exec) + q := updater.NewQueue(exec, nil) q.Start(context.Background()) defer q.Stop()