Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
61 changes: 61 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -11,19 +11,80 @@ jobs:
defaults:
run:
working-directory: ingestion-service
# This service talks to Postgres and Redis, so its tests need both to be
# real. Without them the db and cache suites skip themselves and the job
# still goes green while testing almost nothing.
services:
postgres:
image: postgres:16
env:
POSTGRES_USER: pipelineops
POSTGRES_PASSWORD: pipelineops
POSTGRES_DB: pipelineops
ports:
- 5432:5432
options: >-
--health-cmd pg_isready
--health-interval 5s
--health-timeout 5s
--health-retries 5
redis:
image: redis:7
ports:
- 6379:6379
options: >-
--health-cmd "redis-cli ping"
--health-interval 5s
--health-timeout 5s
--health-retries 5
env:
INGESTION_TEST_DATABASE_URL: postgres://pipelineops:pipelineops@localhost:5432/pipelineops
INGESTION_TEST_REDIS_URL: redis://localhost:6379/0
steps:
- uses: actions/checkout@v4
- uses: actions/setup-go@v5
with:
go-version: "1.24"
cache-dependency-path: ingestion-service/go.sum
- uses: actions/setup-python@v5
with:
python-version: "3.12"
cache: pip
cache-dependency-path: core-api/requirements.txt

# The ingestion service reads and writes tables Django owns the
# migrations for, and it deliberately never runs migrations itself. So
# the schema its tests need has to come from the same place production's
# does — running Django's migrations — rather than from a copy of the
# DDL kept in this service, which would drift the first time a model
# changed.
- name: apply the schema Django owns
working-directory: core-api
env:
DATABASE_URL: postgres://pipelineops:pipelineops@localhost:5432/pipelineops
DJANGO_SECRET_KEY: ci-test-key
run: |
pip install -r requirements.txt
python manage.py migrate --noinput

- name: golangci-lint
uses: golangci/golangci-lint-action@v7
with:
version: v2.12.2
working-directory: ingestion-service
- run: go build ./...
- run: go test ./... -v -coverprofile=coverage.out -covermode=atomic

# The suites skip themselves when Postgres or Redis is unreachable,
# which is the right behaviour locally and a silent failure here. If the
# containers above ever stop being wired up correctly, this fails the
# job instead of quietly reporting a much smaller test run.
- name: no test skipped for a missing dependency
run: |
if go test ./... -v 2>&1 | grep -E "skipping: no reachable (postgres|redis)"; then
echo "::error::integration tests skipped — the postgres/redis service containers are not reachable"
exit 1
fi
- name: Upload coverage
uses: codecov/codecov-action@v5
with:
Expand Down
9 changes: 9 additions & 0 deletions codecov.yml
Original file line number Diff line number Diff line change
Expand Up @@ -22,3 +22,12 @@ flags:
frontend:
paths:
- frontend/src/

ignore:
# main() is now nothing but wiring: read the environment, dial Postgres and
# Redis, serve, wait for a signal, shut down. Everything it used to decide
# moved to cmd/server/bootstrap.go, which is tested directly — the retry
# loops, the request logger, the env fallbacks and the route table. What is
# left here could only be covered by starting the process and killing it,
# which would measure the test harness rather than the code.
- "ingestion-service/cmd/server/main.go"
90 changes: 90 additions & 0 deletions ingestion-service/cmd/server/bootstrap.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,90 @@
package main

import (
"log/slog"
"os"
"time"

"github.com/gin-gonic/gin"
"github.com/prometheus/client_golang/prometheus/promhttp"

"pipelineops/ingestion-service/internal/cache"
"pipelineops/ingestion-service/internal/db"
"pipelineops/ingestion-service/internal/handlers"
)

// Everything main() decides lives here rather than inside main() itself, so it
// can be tested. What's left in main.go is the part that genuinely can't be:
// binding a port, waiting on a signal, and exiting.

func getenv(key, fallback string) string {
if v := os.Getenv(key); v != "" {
return v
}
return fallback
}

// newRouter builds the full route table. It's separated from main so the
// wiring is covered by a test — a route registered at the wrong path or bound
// to the wrong handler is a real bug, and one that a compiler can't catch and
// a running service only reveals in production.
func newRouter(deps handlers.Deps, logger *slog.Logger) *gin.Engine {
router := gin.New()
router.Use(gin.Recovery(), slogRequestLogger(logger))

router.GET("/healthz", handlers.Healthz(deps))
router.GET("/metrics", gin.WrapH(promhttp.Handler()))

v1 := router.Group("/v1")
{
v1.POST("/heartbeat", deps.PostHeartbeat)
v1.GET("/jobs/:job/heartbeat/latest", deps.GetLatestHeartbeat)
}
return router
}

// connectWithRetry calls connect up to 15 times, sleeping 2 seconds between
// attempts, until it succeeds or the attempts are exhausted. connect and
// sleep are injected so tests can substitute a fake connector and a
// non-blocking sleep; main passes the real db.Connect and time.Sleep.
func connectWithRetry(logger *slog.Logger, connect func() (*db.Pool, error), sleep func(time.Duration)) (*db.Pool, error) {
var lastErr error
for i := 0; i < 15; i++ {
pool, err := connect()
if err == nil {
return pool, nil
}
lastErr = err
logger.Warn("postgres not ready, retrying", "attempt", i+1, "error", err)
sleep(2 * time.Second)
}
return nil, lastErr
}

// connectRedisWithRetry mirrors connectWithRetry for the Redis client.
func connectRedisWithRetry(logger *slog.Logger, connect func() (*cache.Client, error), sleep func(time.Duration)) (*cache.Client, error) {
var lastErr error
for i := 0; i < 15; i++ {
client, err := connect()
if err == nil {
return client, nil
}
lastErr = err
logger.Warn("redis not ready, retrying", "attempt", i+1, "error", err)
sleep(2 * time.Second)
}
return nil, lastErr
}

func slogRequestLogger(logger *slog.Logger) gin.HandlerFunc {
return func(c *gin.Context) {
start := time.Now()
c.Next()
logger.Info("request",
"method", c.Request.Method,
"path", c.Request.URL.Path,
"status", c.Writer.Status(),
"duration_ms", time.Since(start).Milliseconds(),
)
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package main

import (
"bytes"
"context"
"encoding/json"
"errors"
"log/slog"
Expand All @@ -13,9 +14,11 @@ import (
"time"

"github.com/gin-gonic/gin"
"github.com/google/uuid"

"pipelineops/ingestion-service/internal/cache"
"pipelineops/ingestion-service/internal/db"
"pipelineops/ingestion-service/internal/handlers"
)

func init() {
Expand Down Expand Up @@ -270,3 +273,158 @@ func TestConnectRedisWithRetry_ExhaustsRetriesAndReturnsLastError(t *testing.T)
t.Fatalf("sleep called %d times, want 15", sleeps)
}
}

// --- route wiring -----------------------------------------------------------
//
// newRouter is the one part of startup where a mistake is invisible until
// production: a route registered at the wrong path, or bound to the wrong
// handler, compiles perfectly and fails only when a real client calls it.
// These tests assert the route table and then prove each route reaches the
// handler it claims to, by checking behaviour only that handler produces.

type stubDB struct {
job *db.Job
findErr error
pingErr error
}

func (s stubDB) FindJobByNameOrID(_ context.Context, _ string) (*db.Job, error) {
return s.job, s.findErr
}

func (s stubDB) InsertHeartbeat(_ context.Context, _ db.InsertHeartbeatParams) (int64, time.Time, error) {
return 1, time.Unix(0, 0).UTC(), nil
}

func (s stubDB) Ping(_ context.Context) error { return s.pingErr }

type stubCache struct {
hb *cache.LastHeartbeat
pingErr error
}

func (s stubCache) SetLastHeartbeat(_ context.Context, _ cache.LastHeartbeat) error { return nil }

func (s stubCache) GetLastHeartbeat(_ context.Context, _ string) (*cache.LastHeartbeat, error) {
return s.hb, nil
}

func (s stubCache) Ping(_ context.Context) error { return s.pingErr }

func testDeps(d stubDB, c stubCache) handlers.Deps {
return handlers.Deps{DB: d, Cache: c, Logger: discardLogger()}
}

func TestNewRouter_RegistersEveryRoute(t *testing.T) {
router := newRouter(testDeps(stubDB{}, stubCache{}), discardLogger())

registered := map[string]bool{}
for _, r := range router.Routes() {
registered[r.Method+" "+r.Path] = true
}

for _, want := range []string{
"GET /healthz",
"GET /metrics",
"POST /v1/heartbeat",
"GET /v1/jobs/:job/heartbeat/latest",
} {
if !registered[want] {
t.Errorf("route %q is not registered; got %v", want, registered)
}
}
if len(router.Routes()) != 4 {
t.Errorf("router has %d routes, want exactly 4 — an unintended route is as much a bug as a missing one", len(router.Routes()))
}
}

func TestNewRouter_HealthzReachesTheHealthHandler(t *testing.T) {
t.Run("healthy when both dependencies answer", func(t *testing.T) {
router := newRouter(testDeps(stubDB{}, stubCache{}), discardLogger())
rec := httptest.NewRecorder()
router.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/healthz", nil))

if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusOK)
}
if !strings.Contains(rec.Body.String(), `"status":"ok"`) {
t.Fatalf("body = %s", rec.Body.String())
}
})

// A 503 here can only come from Healthz itself, which is what makes this
// a test of the binding rather than of the path.
t.Run("degraded when postgres is down", func(t *testing.T) {
router := newRouter(testDeps(stubDB{pingErr: errors.New("down")}, stubCache{}), discardLogger())
rec := httptest.NewRecorder()
router.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/healthz", nil))

if rec.Code != http.StatusServiceUnavailable {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusServiceUnavailable)
}
if !strings.Contains(rec.Body.String(), `"db":"down"`) {
t.Fatalf("body = %s", rec.Body.String())
}
})
}

func TestNewRouter_MetricsServesPrometheus(t *testing.T) {
router := newRouter(testDeps(stubDB{}, stubCache{}), discardLogger())
rec := httptest.NewRecorder()
router.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/metrics", nil))

if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusOK)
}
if !strings.Contains(rec.Body.String(), "# HELP") {
t.Fatalf("/metrics did not return a Prometheus exposition body: %s", rec.Body.String()[:200])
}
}

func TestNewRouter_HeartbeatRouteReachesPostHeartbeat(t *testing.T) {
router := newRouter(testDeps(stubDB{}, stubCache{}), discardLogger())

// An empty body fails PostHeartbeat's binding on the required "job"
// field. Only that handler can produce this 400.
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodPost, "/v1/heartbeat", strings.NewReader(`{}`))
req.Header.Set("Content-Type", "application/json")
router.ServeHTTP(rec, req)

if rec.Code != http.StatusBadRequest {
t.Fatalf("status = %d, want %d (body: %s)", rec.Code, http.StatusBadRequest, rec.Body.String())
}
}

func TestNewRouter_LatestHeartbeatRouteBindsTheJobParam(t *testing.T) {
jobID := uuid.New()
router := newRouter(
testDeps(stubDB{job: &db.Job{ID: jobID, Name: "nightly-etl"}}, stubCache{}),
discardLogger(),
)

rec := httptest.NewRecorder()
router.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/v1/jobs/nightly-etl/heartbeat/latest", nil))

if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusOK)
}
// The handler echoes the resolved job id, so seeing it here proves the
// :job path segment was actually bound and passed through.
if !strings.Contains(rec.Body.String(), jobID.String()) {
t.Fatalf("body = %s, want it to contain the resolved job id %s", rec.Body.String(), jobID)
}
if got := rec.Header().Get("X-Cache"); got != "MISS" {
t.Fatalf("X-Cache = %q, want MISS", got)
}
}

func TestNewRouter_UnknownRouteIs404(t *testing.T) {
router := newRouter(testDeps(stubDB{}, stubCache{}), discardLogger())
rec := httptest.NewRecorder()
router.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/v1/nope", nil))

if rec.Code != http.StatusNotFound {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusNotFound)
}
}
Loading
Loading