SPB Git forge
15commits 1branches 0releases
29.7 MBsize
maindefault branch
10 days agolast push
TypeScript 36.3% Python 31.8% Go 18% JavaScript 9.8% Shell 1.9% SQL 1.4% CSS 0.5%
5.2 KB · 172 lines go
Raw Blame History
1package client23import (4	"bytes"5	"compress/gzip"6	"context"7	"encoding/json"8	"errors"9	"net/http"10	"net/http/httptest"11	"strconv"12	"strings"13	"testing"14	"time"1516	"internetpressure.io/probe-agent/internal/protocol"17	"internetpressure.io/probe-agent/internal/signer"18)1920const key = "000102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f"2122// fakeServer verifies signatures exactly like the Python ingest API should.23func fakeServer(t *testing.T, skew time.Duration) (*httptest.Server, *int) {24	verifier, _ := signer.New("ca-qc-01", key)25	calls := 026	h := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {27		calls++28		body, _ := readAll(r)29		ts, _ := strconv.ParseInt(r.Header.Get(signer.HeaderTimestamp), 10, 64)30		if r.Header.Get(signer.HeaderProbe) != "ca-qc-01" ||31			!verifier.Verify(r.Method, r.URL.Path, ts, body, r.Header.Get(signer.HeaderSignature)) {32			http.Error(w, `{"detail":"bad signature"}`, 401)33			return34		}35		if !strings.HasPrefix(r.UserAgent(), "InternetPressureProbe/0.1.0 (+https://www.internetpressure.io/probes)") {36			http.Error(w, "ua", 400)37			return38		}39		serverNow := time.Now().Add(-skew)40		w.Header().Set("Date", serverNow.UTC().Format(http.TimeFormat))41		if abs(time.Duration(ts)*time.Second-time.Duration(serverNow.Unix())*time.Second) > 300*time.Second {42			http.Error(w, `{"detail":"timestamp skew too large"}`, 401)43			return44		}45		w.Header().Set("Content-Type", "application/json")46		switch r.URL.Path {47		case "/ingest/v1/config":48			json.NewEncoder(w).Encode(map[string]any{"server_time": protocol.FormatTime(serverNow),49				"config_version": "v1", "probe": map[string]any{"probe_id": "ca-qc-01", "enabled": true},50				"schedule": map[string]any{"tiers": map[string]int{"1": 20}}, "targets": []any{}})51		case "/ingest/v1/batch":52			if r.Header.Get("Content-Encoding") != "gzip" {53				http.Error(w, "not gzip", 400)54				return55			}56			zr, err := gzip.NewReader(bytes.NewReader(body))57			if err != nil {58				http.Error(w, "bad gzip", 400)59				return60			}61			var b protocol.Batch62			if err := json.NewDecoder(zr).Decode(&b); err != nil {63				http.Error(w, "bad json", 400)64				return65			}66			json.NewEncoder(w).Encode(protocol.BatchResponse{Accepted: len(b.Measurements),67				ConfigVersion: "v2", ServerTime: protocol.FormatTime(serverNow)})68		default:69			http.NotFound(w, r)70		}71	})72	return httptest.NewServer(h), &calls73}7475func abs(d time.Duration) time.Duration {76	if d < 0 {77		return -d78	}79	return d80}8182func readAll(r *http.Request) ([]byte, error) {83	var buf bytes.Buffer84	_, err := buf.ReadFrom(r.Body)85	return buf.Bytes(), err86}8788func newClient(t *testing.T, base string) *Client {89	s, _ := signer.New("ca-qc-01", key)90	c, err := New(base+"/ingest/v1", s, "0.1.0")91	if err != nil {92		t.Fatal(err)93	}94	return c95}9697func TestConfigAndBatchSigned(t *testing.T) {98	srv, _ := fakeServer(t, 0)99	defer srv.Close()100	c := newClient(t, srv.URL)101102	cfg, err := c.GetConfig(context.Background())103	if err != nil {104		t.Fatal(err)105	}106	if cfg.ConfigVersion != "v1" || !cfg.Probe.Enabled {107		t.Fatalf("config: %+v", cfg)108	}109	if c.Clock.Samples() != 1 || abs(c.Clock.Offset()) > time.Second {110		t.Fatalf("clock not fed: samples=%d offset=%v", c.Clock.Samples(), c.Clock.Offset())111	}112113	var buf bytes.Buffer114	zw := gzip.NewWriter(&buf)115	json.NewEncoder(zw).Encode(protocol.Batch{ProbeID: "ca-qc-01", AgentVersion: "0.1.0",116		Measurements: []protocol.Measurement{{Kind: "http"}, {Kind: "dns"}}})117	zw.Close()118	resp, err := c.PostBatch(context.Background(), buf.Bytes(), 0)119	if err != nil {120		t.Fatal(err)121	}122	if resp.Accepted != 2 || resp.ConfigVersion != "v2" {123		t.Fatalf("batch resp: %+v", resp)124	}125}126127func TestSkewResyncViaDateHeader(t *testing.T) {128	// Server clock is 10 minutes behind us → first request is rejected with "skew", Date header resyncs us.129	srv, calls := fakeServer(t, 10*time.Minute)130	defer srv.Close()131	c := newClient(t, srv.URL)132	_, err := c.GetConfig(context.Background())133	var he *HTTPError134	if !errors.As(err, &he) || !he.IsSkew() {135		t.Fatalf("expected skew 401, got %v", err)136	}137	if off := c.Clock.Offset(); off < 9*time.Minute || off > 11*time.Minute {138		t.Fatalf("offset after resync = %v", off)139	}140	if _, err := c.GetConfig(context.Background()); err != nil {141		t.Fatalf("retry after resync should pass: %v", err)142	}143	if *calls != 2 {144		t.Fatalf("calls = %d", *calls)145	}146}147148func TestRetryable(t *testing.T) {149	if !Retryable(errors.New("dial tcp: connection refused")) {150		t.Error("network error should be retryable")151	}152	if !Retryable(&HTTPError{Status: 503}) || Retryable(&HTTPError{Status: 401}) || Retryable(&HTTPError{Status: 400}) {153		t.Error("status classification wrong")154	}155	if Retryable(nil) {156		t.Error("nil is not retryable")157	}158}159160func TestClockEWMA(t *testing.T) {161	c := NewClock()162	now := time.Now()163	c.Observe(now, 100*time.Millisecond, now.Add(50*time.Millisecond).Add(-2*time.Second)) // local ahead by 2 s164	if got := c.Offset(); got < 1900*time.Millisecond || got > 2100*time.Millisecond {165		t.Fatalf("first sample should be taken as-is: %v", got)166	}167	c.Observe(now, 100*time.Millisecond, now.Add(50*time.Millisecond))                 // offset 0 sample168	if got := c.Offset(); got < 1300*time.Millisecond || got > 1500*time.Millisecond { // 2 s × 0.7169		t.Fatalf("ewma: %v", got)170	}171}172