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
6 changes: 6 additions & 0 deletions docker-compose.lifecycle.yml
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@ services:
- 0.0.0.0:7070
- -health
- 0.0.0.0:8080
- -health-auth-token
- "${STREAMHIVE_LIFECYCLE_HEALTH_TOKEN:-streamhive-lifecycle-health-demo-token}"
- -replicate
- -store-dir
- /data
Expand Down Expand Up @@ -43,6 +45,8 @@ services:
- 0.0.0.0:7070
- -health
- 0.0.0.0:8080
- -health-auth-token
- "${STREAMHIVE_LIFECYCLE_HEALTH_TOKEN:-streamhive-lifecycle-health-demo-token}"
- -replicate
- -store-dir
- /data
Expand Down Expand Up @@ -79,6 +83,8 @@ services:
- 0.0.0.0:7070
- -health
- 0.0.0.0:8080
- -health-auth-token
- "${STREAMHIVE_LIFECYCLE_HEALTH_TOKEN:-streamhive-lifecycle-health-demo-token}"
- -replicate
- -store-dir
- /data
Expand Down
6 changes: 6 additions & 0 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@ services:
- 0.0.0.0:7070
- -health
- 0.0.0.0:8080
- -health-auth-token
- "${STREAMHIVE_HEALTH_TOKEN:-streamhive-health-demo-token}"
- -replicate
- -store-dir
- /data
Expand Down Expand Up @@ -37,6 +39,8 @@ services:
- 0.0.0.0:7070
- -health
- 0.0.0.0:8080
- -health-auth-token
- "${STREAMHIVE_HEALTH_TOKEN:-streamhive-health-demo-token}"
- -replicate
- -store-dir
- /data
Expand Down Expand Up @@ -67,6 +71,8 @@ services:
- 0.0.0.0:7070
- -health
- 0.0.0.0:8080
- -health-auth-token
- "${STREAMHIVE_HEALTH_TOKEN:-streamhive-health-demo-token}"
- -replicate
- -store-dir
- /data
Expand Down
24 changes: 19 additions & 5 deletions docs/DEPLOYMENT.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,15 +6,23 @@ Build and run (example flags):

```bash
docker build -t streamhive:local .
docker run --rm -p 7070:7070 -p 8080:8080 streamhive:local \
docker run --rm -e STREAMHIVE_HEALTH_TOKEN=change-me -p 7070:7070 -p 8080:8080 streamhive:local \
-listen 0.0.0.0:7070 \
-health 0.0.0.0:8080
-health 0.0.0.0:8080 \
-health-auth-token change-me
```

- **7070** — P2P TCP listener (example).
- **8080** — HTTP `/livez`, `/readyz`, `/version` (aggregate runtime identity), `/peers` (JSON peer metadata), `/metrics` (JSON counters), `/metrics/prometheus` (Prometheus text), `/inventory/status` (aggregate live key fingerprint), `/storage/status` (aggregate live blob integrity), and `/lifecycle/status` (aggregate opt-in lifecycle state).

Use TLS flags (`-tls-cert`, `-tls-key`, `-tls-ca`, `-tls-server-name`) when exposing services beyond a lab network. The verified certificate and application-auth ordering is documented in [TLS_AUTH.md](TLS_AUTH.md) and exercised by `make test-tls-auth`. Reserve `-tls-insecure-skip-verify` for local development. For CLI mTLS, add `-tls-client-ca -tls-require-client-cert` on the listener and `-tls-client-cert -tls-client-key` on outbound peers. Use `-tls-expiry-warning` to set the aggregate short-lived-credential warning window (`720h` by default, `0` disables the warning). For custom trust policy, configure `p2p.TCPTransport.TLSServerConfig` and `TLSClientConfig` in library code.
The health listener is plain HTTP and loopback-first. `/livez` and `/readyz` remain unauthenticated for local and orchestrator probes; diagnostic routes require `Authorization: Bearer <token>` whenever `-health-auth-token` is configured. A non-loopback `-health` address requires `-health-auth-token` and fails closed without one. For example, an explicitly network-scoped lab listener is:

```bash
go run . -listen 0.0.0.0:7070 -health 0.0.0.0:8080 \
-health-auth-token "$STREAMHIVE_HEALTH_TOKEN"
```

Use a TLS-terminating reverse proxy or a private, access-controlled network before exposing health beyond a trusted lab; the P2P TLS flags protect the P2P listener, not this HTTP listener. The verified certificate and application-auth ordering is documented in [TLS_AUTH.md](TLS_AUTH.md) and exercised by `make test-tls-auth`. Reserve `-tls-insecure-skip-verify` for local development. For CLI mTLS, add `-tls-client-ca -tls-require-client-cert` on the listener and `-tls-client-cert -tls-client-key` on outbound peers. Use `-tls-expiry-warning` to set the aggregate short-lived-credential warning window (`720h` by default, `0` disables the warning). For custom trust policy, configure `p2p.TCPTransport.TLSServerConfig` and `TLSClientConfig` in library code.

Certificate files are loaded at process startup. Use the [TLS rotation runbook](TLS_ROTATION.md)
for replacement material, trust overlap, restart ordering, rollback, and reconnect checks.
Expand Down Expand Up @@ -187,7 +195,13 @@ spec:
containers:
- name: streamhive
image: streamhive:local
args: ["-listen", "0.0.0.0:7070", "-health", "0.0.0.0:8080"]
args: ["-listen", "0.0.0.0:7070", "-health", "0.0.0.0:8080", "-health-auth-token", "$(STREAMHIVE_HEALTH_TOKEN)"]
env:
- name: STREAMHIVE_HEALTH_TOKEN
valueFrom:
secretKeyRef:
name: streamhive-health
key: token
ports:
- containerPort: 7070
name: p2p
Expand All @@ -207,7 +221,7 @@ spec:
periodSeconds: 10
```

Add a `Service` for the health port and (separately) headless or load-balanced service for P2P depending on your topology. Tune resource requests/limits and pod anti-affinity for HA; this manifest is illustrative only.
The health token is intentionally supplied from a Secret. `/livez` and `/readyz` do not require the token, so the illustrative probes remain usable; diagnostic clients must send the bearer header. Add a `Service` for the health port and (separately) headless or load-balanced service for P2P depending on your topology. Tune resource requests/limits and pod anti-affinity for HA; this manifest is illustrative only.

## SLOs

Expand Down
6 changes: 5 additions & 1 deletion docs/PROTOCOL.md
Original file line number Diff line number Diff line change
Expand Up @@ -426,7 +426,11 @@ exemplars, peer addresses, blob keys, or certificate metadata.
The optional HTTP health server also uses fixed bounds rather than unbounded request handling:
5-second header reads, 10-second request reads, 10-second response writes, 60-second idle
connections, and 1 MiB maximum headers. Process cancellation gracefully shuts down the server;
the P2P wire protocol and endpoint paths are unchanged.
the P2P wire protocol and endpoint paths are unchanged. Health binds are loopback-first; a
non-loopback address requires the explicit `-health-auth-token` bearer boundary. `/livez` and
`/readyz` remain probe-safe, while diagnostic routes reject requests without the configured token.
The P2P TLS flags do not apply to this HTTP listener, so deployments beyond a trusted lab must use
a TLS-terminating proxy or an equivalent private network boundary.

## Shutdown and Drain

Expand Down
8 changes: 7 additions & 1 deletion docs/TLS_AUTH.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,13 +35,19 @@ The listener loads its certificate and private key:
```bash
go run . \
-listen 0.0.0.0:7070 \
-health 0.0.0.0:8080 \
-health 127.0.0.1:8080 \
-replicate -store-dir ./streamhive-data \
-tls-cert ./server-cert.pem -tls-key ./server-key.pem \
-peer-auth-token "$STREAMHIVE_PEER_TOKEN" \
-peer-id server -peer-allow-ids client
```

The P2P TLS flags above do not encrypt or authenticate the separate HTTP health listener. Keep
health on loopback, or provide `-health-auth-token` for a deliberately network-scoped lab bind and
place a TLS-terminating reverse proxy or private access-controlled network in front of it. Diagnostic
health requests then send `Authorization: Bearer <token>`; `/livez` and `/readyz` remain available to
orchestrator probes without the token.

The outbound peer supplies the trusted CA and the DNS name present in the certificate:

```bash
Expand Down
63 changes: 63 additions & 0 deletions health_server_acceptance_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,69 @@ func TestRun_healthServerReadOnlyHandlerRejectsBeforeWork(t *testing.T) {
assert.False(t, called)
}

func TestRun_healthServerRequiresAuthForNonLoopbackAddress(t *testing.T) {
var out, stderr safeBuffer
err := run(context.Background(), []string{
"-listen", "127.0.0.1:0",
"-health", "0.0.0.0:0",
}, &out, &stderr)
require.Error(t, err)
assert.Contains(t, err.Error(), "requires -health-auth-token")
assert.Empty(t, out.String())
}

func TestRun_healthServerBearerAuth(t *testing.T) {
called := false
next := healthAuthHandler("test-token", http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
called = true
w.WriteHeader(http.StatusNoContent)
}))

for _, test := range []struct {
name string
path string
authority string
statusCode int
wantCalled bool
}{
{name: "missing", path: "/metrics", statusCode: http.StatusUnauthorized},
{name: "wrong", path: "/metrics", authority: "Bearer wrong-token", statusCode: http.StatusUnauthorized},
{name: "valid", path: "/metrics", authority: "Bearer test-token", statusCode: http.StatusNoContent, wantCalled: true},
{name: "liveness probe", path: "/livez", statusCode: http.StatusNoContent, wantCalled: true},
} {
t.Run(test.name, func(t *testing.T) {
called = false
req := httptest.NewRequest(http.MethodGet, "http://127.0.0.1"+test.path, nil)
if test.authority != "" {
req.Header.Set("Authorization", test.authority)
}
resp := httptest.NewRecorder()
next.ServeHTTP(resp, req)
assert.Equal(t, test.statusCode, resp.Code)
assert.Equal(t, test.wantCalled, called)
})
}
}

func TestRun_healthAddrIsLoopback(t *testing.T) {
for _, test := range []struct {
addr string
loopback bool
}{
{addr: "127.0.0.1:8080", loopback: true},
{addr: "[::1]:8080", loopback: true},
{addr: "localhost:8080", loopback: true},
{addr: ":8080"},
{addr: "0.0.0.0:8080"},
{addr: "[::]:8080"},
{addr: "health.internal:8080"},
} {
t.Run(test.addr, func(t *testing.T) {
assert.Equal(t, test.loopback, healthAddrIsLoopback(test.addr))
})
}
}

func TestRun_healthServerReadOnlyEndpoints(t *testing.T) {
node := startHealthAcceptanceNode(t)
t.Cleanup(func() { node.stop(t) })
Expand Down
44 changes: 41 additions & 3 deletions main.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package main
import (
"bytes"
"context"
"crypto/subtle"
"crypto/tls"
"crypto/x509"
"encoding/hex"
Expand Down Expand Up @@ -80,6 +81,7 @@ func run(ctx context.Context, args []string, stdout, stderr io.Writer) error {
peerReconnectMax := fs.Duration("peer-reconnect-max", 30*time.Second, "maximum reconnect backoff for -peer-reconnect")
syncInterval := fs.Duration("sync-interval", 0, "periodically advertise local blob keys to connected peers (0 = startup only)")
health := fs.String("health", "", "optional HTTP listen addr for /livez /readyz /peers /metrics (e.g. :8080)")
healthAuthToken := fs.String("health-auth-token", "", "bearer token required when -health binds a non-loopback address")
shutdownGrace := fs.Duration("shutdown-grace", defaultShutdownGrace, "bounded graceful shutdown deadline for health and P2P drain")
maxPeers := fs.Int("max-peers", 0, "max simultaneous peers (0 = unlimited)")
peerAuthToken := fs.String("peer-auth-token", "", "optional shared token required before peer registration")
Expand Down Expand Up @@ -136,6 +138,9 @@ func run(ctx context.Context, args []string, stdout, stderr io.Writer) error {
if err := fs.Parse(args); err != nil {
return err
}
if *health != "" && !healthAddrIsLoopback(*health) && strings.TrimSpace(*healthAuthToken) == "" {
return fmt.Errorf("health: non-loopback address %q requires -health-auth-token", *health)
}

log := slog.New(slog.NewTextHandler(stderr, &slog.HandlerOptions{Level: slog.LevelInfo}))

Expand Down Expand Up @@ -700,7 +705,7 @@ func run(ctx context.Context, args []string, stdout, stderr io.Writer) error {

if *health != "" {
var err error
hsrv, err = startHealth(*health, tr, replMetrics, blobStore, keyLister, tlsHealth, lifecycleState, reconnectMetrics, log)
hsrv, err = startHealth(*health, *healthAuthToken, tr, replMetrics, blobStore, keyLister, tlsHealth, lifecycleState, reconnectMetrics, log)
if err != nil {
return fmt.Errorf("health: %w", err)
}
Expand Down Expand Up @@ -2422,7 +2427,39 @@ func readOnlyHealthHandler(handler http.HandlerFunc) http.HandlerFunc {
}
}

func startHealth(addr string, tr *p2p.TCPTransport, replMetrics *replicationMetrics, blobStore storage.BlobStore, keyLister storage.BlobKeyLister, tlsHealth *tlsCredentialHealth, lifecycleState *lifecycleRuntime, reconnectMetrics *peerReconnectMetrics, log *slog.Logger) (*http.Server, error) {
func healthAddrIsLoopback(addr string) bool {
host, _, err := net.SplitHostPort(addr)
if err != nil {
return false
}
host = strings.Trim(host, "[]")
if strings.EqualFold(host, "localhost") {
return true
}
ip := net.ParseIP(host)
return ip != nil && ip.IsLoopback()
}

func healthAuthHandler(token string, next http.Handler) http.Handler {
if token == "" {
return next
}
expected := "Bearer " + token
return http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) {
if req.URL.Path == "/livez" || req.URL.Path == "/readyz" {
next.ServeHTTP(w, req)
return
}
if subtle.ConstantTimeCompare([]byte(req.Header.Get("Authorization")), []byte(expected)) != 1 {
w.Header().Set("WWW-Authenticate", `Bearer realm="streamhive-health"`)
http.Error(w, "unauthorized", http.StatusUnauthorized)
return
}
next.ServeHTTP(w, req)
})
}

func startHealth(addr, authToken string, tr *p2p.TCPTransport, replMetrics *replicationMetrics, blobStore storage.BlobStore, keyLister storage.BlobKeyLister, tlsHealth *tlsCredentialHealth, lifecycleState *lifecycleRuntime, reconnectMetrics *peerReconnectMetrics, log *slog.Logger) (*http.Server, error) {
mux := http.NewServeMux()
mux.HandleFunc("/livez", readOnlyHealthHandler(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
Expand Down Expand Up @@ -2512,14 +2549,15 @@ func startHealth(addr string, tr *p2p.TCPTransport, replMetrics *replicationMetr
}
writePrometheusMetrics(w, snapshot)
}))
protected := healthAuthHandler(authToken, mux)

ln, err := net.Listen("tcp", addr)
if err != nil {
return nil, err
}

srv := &http.Server{
Handler: mux,
Handler: protected,
ReadHeaderTimeout: defaultHealthReadHeaderTimeout,
ReadTimeout: defaultHealthReadTimeout,
WriteTimeout: defaultHealthWriteTimeout,
Expand Down
18 changes: 12 additions & 6 deletions scripts/demo-auth.sh
Original file line number Diff line number Diff line change
Expand Up @@ -6,12 +6,14 @@ DATA_DIR="${STREAMHIVE_DATA_DIR:-$ROOT_DIR/.streamhive-auth}"
COMPOSE="docker compose"
EXPECTED_KEY="cd13ac0817f0f8ba2f29fba23617ef0191a6193ed0311298163834199398ee05"
TOKEN="${STREAMHIVE_PEER_TOKEN:-streamhive-compose-demo-token}"
HEALTH_TOKEN="${STREAMHIVE_HEALTH_TOKEN:-streamhive-health-demo-token}"
WRONG_TOKEN="${STREAMHIVE_WRONG_PEER_TOKEN:-streamhive-invalid-token}"
if [ "$TOKEN" = "$WRONG_TOKEN" ]; then
WRONG_TOKEN="${TOKEN}!"
fi
export STREAMHIVE_DATA_DIR="$DATA_DIR"
export STREAMHIVE_PEER_TOKEN="$TOKEN"
export STREAMHIVE_HEALTH_TOKEN="$HEALTH_TOKEN"
export STREAMHIVE_NODE1_ID=node1
export STREAMHIVE_NODE1_ALLOW_IDS=node2,node3,seed
export STREAMHIVE_NODE2_ID=node2
Expand All @@ -20,6 +22,10 @@ export STREAMHIVE_NODE3_ID=node3
export STREAMHIVE_NODE3_ALLOW_IDS=node1,node2
export STREAMHIVE_SEED_ID=seed

curl_health() {
curl -fsS -H "Authorization: Bearer $HEALTH_TOKEN" "$@"
}

cleanup() {
$COMPOSE -f "$ROOT_DIR/docker-compose.yml" down --remove-orphans >/dev/null 2>&1 || true
}
Expand All @@ -29,7 +35,7 @@ wait_ready() {
name="$1"
url="$2"
i=0
until curl -fsS "$url/readyz" >/dev/null 2>&1; do
until curl_health "$url/readyz" >/dev/null 2>&1; do
i=$((i + 1))
if [ "$i" -gt 80 ]; then
echo "$name did not become ready" >&2
Expand All @@ -45,11 +51,11 @@ wait_metric() {
url="$2"
metric="$3"
i=0
until curl -fsS "$url/metrics" | grep "\"$metric\": [1-9]" >/dev/null; do
until curl_health "$url/metrics" | grep "\"$metric\": [1-9]" >/dev/null; do
i=$((i + 1))
if [ "$i" -gt 80 ]; then
echo "$name did not report a positive $metric counter" >&2
curl -fsS "$url/metrics" >&2 || true
curl_health "$url/metrics" >&2 || true
$COMPOSE -f "$ROOT_DIR/docker-compose.yml" logs "$name" >&2 || true
exit 1
fi
Expand Down Expand Up @@ -80,19 +86,19 @@ wait_key_present() {
metric_value() {
url="$1"
metric="$2"
curl -fsS "$url/metrics" | awk -F': ' -v metric="\"$metric\"" 'index($1, metric) > 0 {gsub(/,/, "", $2); print $2; exit}'
curl_health "$url/metrics" | awk -F': ' -v metric="\"$metric\"" 'index($1, metric) > 0 {gsub(/,/, "", $2); print $2; exit}'
}

wait_identity() {
name="$1"
url="$2"
identity="$3"
i=0
until curl -fsS "$url/peers" | grep "\"auth_identity\": \"$identity\"" >/dev/null; do
until curl_health "$url/peers" | grep "\"auth_identity\": \"$identity\"" >/dev/null; do
i=$((i + 1))
if [ "$i" -gt 80 ]; then
echo "$name did not expose authenticated peer identity $identity" >&2
curl -fsS "$url/peers" >&2 || true
curl_health "$url/peers" >&2 || true
$COMPOSE -f "$ROOT_DIR/docker-compose.yml" logs "$name" >&2 || true
exit 1
fi
Expand Down
Loading