xprobe

package module
v1.0.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: May 21, 2026 License: MIT Imports: 4 Imported by: 0

README

xprobe

CI

Composable, transport-agnostic health probes for Go services. Pull-mode HTTP (Kubernetes liveness/readiness/startup), cached HTTP, and grpc.health.v1 — all backed by the same probe primitives.

  • Stdlib-only core. gRPC transport lives in a separate module so non-gRPC users don't pull google.golang.org/grpc.
  • Status taxonomy distinguishes Up / Down / Timeout / Unknown — slow probes don't masquerade as failures.
  • Composite probes with clear All (AND) / Any (OR) semantics, concurrent execution.
  • Cached state with pub/sub — required for gRPC Watch, useful for expensive checks under burst traffic.
  • Pluggable reporters (slog included) invoked only on status transitions.

Minimum Go version: 1.25.0


Install

Core (probes + HTTP transport):

go get github.com/gopherex/xprobe@latest

gRPC transport (optional, separate module):

go get github.com/gopherex/xprobe/pkg/transport/grpc@latest

Quickstart — Kubernetes liveness/readiness

package main

import (
	"context"
	"errors"
	"log"
	"net/http"

	"github.com/gopherex/xprobe"
)

func main() {
	live := xprobe.NewBool()
	live.Set(true)

	ready := xprobe.All(
		xprobe.FromError(pingDB),
		xprobe.FromError(pingCache),
	)

	mux := xprobe.Mux(
		xprobe.Liveness(live),
		xprobe.Readiness(ready),
	)

	log.Fatal(http.ListenAndServe(":8080", mux))
}

func pingDB(ctx context.Context) error    { return nil }
func pingCache(ctx context.Context) error { return errors.New("offline") }

Endpoints exposed:

Path Source
/healthz/liveness NewBool().Set(true)
/healthz/readiness DB AND cache up
/healthz/startup (if you add one)

Response codes: 200 Healthy / 503 Unhealthy <name> / 504 Unhealthy <name> (timeout).


Status taxonomy

xprobe.StatusUnknown  // never checked (zero value)
xprobe.StatusUp       // healthy
xprobe.StatusDown     // probe returned a failure
xprobe.StatusTimeout  // probe exceeded its deadline

Status.OK() is true only for StatusUp.


Defining probes

From an error-returning func
p := xprobe.FromError(func(ctx context.Context) error {
	return db.PingContext(ctx)
})
Toggle flag
b := xprobe.NewBool()
b.Set(true)   // mark up
b.Set(false)  // mark down

Useful as a liveness flag flipped during graceful shutdown.

Custom
p := xprobe.Func(func(ctx context.Context) xprobe.Status {
	if queueDepth() > 1000 {
		return xprobe.StatusDown
	}
	return xprobe.StatusUp
})
Composition
xprobe.All(p1, p2, p3)  // AND — worst status wins
xprobe.Any(p1, p2)      // OR  — best status wins

Children run concurrently. Customize concurrency:

import "github.com/gopherex/xprobe/pkg/probe"

c := probe.New(probe.ModeAll, probe.WithAsyncer(probe.PoolFactory(4)))
c.Add(p1).Add(p2)

HTTP — pull mode (default)

Handler runs the probe synchronously per request. Suitable for cheap checks (in-memory flag, simple ping).

import httpprobe "github.com/gopherex/xprobe/pkg/transport/http"

h := httpprobe.Handler(p,
	httpprobe.WithName("db"),
	httpprobe.WithTimeout(2*time.Second),
	httpprobe.AsJSON(),  // optional JSON body
)
http.Handle("/healthz/db", h)

JSON output:

{"name":"db","status":"up"}

HTTP — cached mode

For expensive probes, or to share state with gRPC. A Runner polls the probe in the background and writes into a State; the handler reads from the State without ever invoking the probe.

import (
	httpprobe "github.com/gopherex/xprobe/pkg/transport/http"
	"github.com/gopherex/xprobe/pkg/runner"
	"github.com/gopherex/xprobe/pkg/state"
)

s := state.New()
r := runner.New(myProbe, s,
	runner.WithName("db"),
	runner.WithInterval(5*time.Second),
	runner.WithTimeout(2*time.Second),
	runner.RunImmediately(),
)
r.Start(ctx)  // non-blocking

http.Handle("/healthz/db", httpprobe.CachedHandler(s, httpprobe.WithName("db")))

The probe runs at most once per interval, regardless of request volume.


Reporters

Side-effect hooks invoked when cached status changes (never on every tick).

import (
	"log/slog"
	"github.com/gopherex/xprobe/pkg/reporter"
	"github.com/gopherex/xprobe/pkg/runner"
)

r := runner.New(myProbe, s,
	runner.WithReporter(reporter.Slog(slog.Default())),
)

Built-in: reporter.Slog, reporter.Func, reporter.Multi, reporter.Nop.

Implement reporter.Reporter to plug Prometheus, OpenTelemetry, alerting, etc.

type Event struct {
	Name string
	Prev probe.Status
	Cur  probe.Status
}

type Reporter interface {
	OnStatus(ctx context.Context, ev Event)
}

Reporters are called synchronously from the runner tick goroutine — a slow reporter blocks the next tick. Dispatch to a worker pool if needed.


gRPC — grpc.health.v1

The gRPC health protocol requires the server to know status synchronously and stream changes. The gRPC transport reads from a state.Registry; back it with runner.Runner per service.

package main

import (
	"context"
	"net"

	"google.golang.org/grpc"
	hv1 "google.golang.org/grpc/health/grpc_health_v1"

	"github.com/gopherex/xprobe/pkg/probe"
	"github.com/gopherex/xprobe/pkg/runner"
	"github.com/gopherex/xprobe/pkg/state"
	grpcprobe "github.com/gopherex/xprobe/pkg/transport/grpc"
)

func main() {
	ctx := context.Background()

	reg := state.NewRegistry()

	// Overall service status: empty name "" per grpc.health.v1 convention.
	runner.New(myCompositeProbe, reg.Get(""),
		runner.WithInterval(5*time.Second),
		runner.RunImmediately(),
	).Start(ctx)

	// Per-component status.
	runner.New(probe.FromError(pingDB), reg.Get("db.Service"),
		runner.WithInterval(2*time.Second),
		runner.RunImmediately(),
	).Start(ctx)

	srv := grpc.NewServer()
	hv1.RegisterHealthServer(srv, grpcprobe.New(reg))

	lis, _ := net.Listen("tcp", ":9000")
	srv.Serve(lis)
}

Client usage (any standard gRPC health client works):

c := hv1.NewHealthClient(conn)
resp, _ := c.Check(ctx, &hv1.HealthCheckRequest{Service: "db.Service"})
// resp.Status: SERVING / NOT_SERVING / UNKNOWN

stream, _ := c.Watch(ctx, &hv1.HealthCheckRequest{Service: ""})
for {
	msg, err := stream.Recv()
	if err != nil { break }
	// react to msg.Status
}

Status mapping:

xprobe grpc.health.v1
StatusUp SERVING
StatusDown NOT_SERVING
StatusTimeout NOT_SERVING
StatusUnknown (post-Set) UNKNOWN
never Set SERVICE_UNKNOWN

Check on a service that has never been registered returns codes.NotFound. Watch on an unknown service emits SERVICE_UNKNOWN then auto-registers, so a subsequent Set is streamed on the same stream — matches google.golang.org/grpc/health reference behavior.

⚠ Each Watch RPC holds a goroutine + buffered channel until the client cancels. Rate-limit public-facing endpoints.


Layout

xprobe/
├── xprobe.go                       facade — aliases + quick-start
├── go.mod                          stdlib only
└── pkg/
    ├── probe/                      Status, Probe, Func, Bool, Composite, Asyncer
    ├── state/                      State (cache + pub/sub), Registry
    ├── runner/                     Background poller
    ├── reporter/                   Reporter interface + slog adapter
    └── transport/
        ├── http/                   pull (Handler) + cached (CachedHandler)
        └── grpc/                   grpc.health.v1 server — SEPARATE go.mod

Why a separate go.mod for gRPC: keeps the core dependency-free for users who only need HTTP probes.


Architecture

                 ┌── reporter.Slog / metrics / alerts
                 │
Probe ──poll──▶ Runner ──Set──▶ State ──Subscribe──▶ gRPC Watch stream
                                  ▲
                                  └── HTTP CachedHandler (instant read)

Probe ────────────────────────────▶ HTTP Handler (pull, on-request)

Pick pull-mode HTTP when your check is cheap. Switch to cached HTTP / gRPC when checks are expensive, traffic is bursty, or you need streaming updates.


License

See LICENSE.

Documentation

Overview

Package xprobe provides composable, transport-agnostic health probes.

Layout:

pkg/probe     — core types (Status, Probe, Composite)
pkg/state     — cached status with pub/sub (used by gRPC Watch and cached HTTP)
pkg/runner    — periodic poller pushing Probe results into State
pkg/reporter  — side-effect hooks invoked on status transitions
pkg/transport/http  — HTTP handlers (pull or cached)
pkg/transport/grpc  — grpc.health.v1 server (separate go.mod)

The root re-exports the most common identifiers for ergonomic use.

Index

Constants

View Source
const (
	StatusUnknown = probe.StatusUnknown
	StatusUp      = probe.StatusUp
	StatusDown    = probe.StatusDown
	StatusTimeout = probe.StatusTimeout
)

Variables

This section is empty.

Functions

func Mux

func Mux(probes ...*HTTPProbe) *http.ServeMux

Types

type Bool

type Bool = probe.Bool

func NewBool

func NewBool() *Bool

type Composite

type Composite = probe.Composite

func All

func All(probes ...Probe) *Composite

func Any

func Any(probes ...Probe) *Composite

type Func

type Func = probe.Func

type HTTPOption

type HTTPOption = httpprobe.Option

type HTTPProbe

type HTTPProbe = httpprobe.HTTPProbe

func Liveness

func Liveness(p Probe, opts ...HTTPOption) *HTTPProbe

func Readiness

func Readiness(p Probe, opts ...HTTPOption) *HTTPProbe

func Startup

func Startup(p Probe, opts ...HTTPOption) *HTTPProbe

type Probe

type Probe = probe.Probe

func FromError

func FromError(f func(ctx context.Context) error) Probe

type Status

type Status = probe.Status

Directories

Path Synopsis
pkg
probe
Package probe defines the core probe abstractions: a status enum, the Probe interface, function adapters and composition helpers.
Package probe defines the core probe abstractions: a status enum, the Probe interface, function adapters and composition helpers.
reporter
Package reporter defines a side-effect interface invoked on probe status transitions.
Package reporter defines a side-effect interface invoked on probe status transitions.
runner
Package runner periodically executes a Probe and pushes the result into a State, invoking Reporters on transitions.
Package runner periodically executes a Probe and pushes the result into a State, invoking Reporters on transitions.
state
Package state holds the cached health status of a probe along with pub/sub semantics so transports (gRPC Watch, log/metric reporters, cached HTTP) can react to changes without re-running the underlying probe.
Package state holds the cached health status of a probe along with pub/sub semantics so transports (gRPC Watch, log/metric reporters, cached HTTP) can react to changes without re-running the underlying probe.
transport/http
Package httpprobe exposes probes over HTTP with sensible defaults for Kubernetes liveness/readiness/startup checks, but generic enough for any service.
Package httpprobe exposes probes over HTTP with sensible defaults for Kubernetes liveness/readiness/startup checks, but generic enough for any service.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL