spb/internetpressure
Public
TypeScript 36.3%
Python 31.8%
Go 18%
JavaScript 9.8%
Shell 1.9%
SQL 1.4%
CSS 0.5%
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