From 716e4f9a617575c26651b8a6b2ce690a5ce6d03b Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Thu, 8 Oct 2026 15:26:47 -0700 Subject: [PATCH] Prevent recursive repository reads through stale PDS endpoints --- .../installing/downloading-streamplace.md | 12 ++ pkg/spxrpc/com_atproto_repo_test.go | 129 ++++++++++++++++++ pkg/spxrpc/spxrpc.go | 11 ++ 3 files changed, 152 insertions(+) create mode 100644 pkg/spxrpc/com_atproto_repo_test.go diff --git a/js/docs/src/content/docs/guides/installing/downloading-streamplace.md b/js/docs/src/content/docs/guides/installing/downloading-streamplace.md index 52286177e..82d85087a 100644 --- a/js/docs/src/content/docs/guides/installing/downloading-streamplace.md +++ b/js/docs/src/content/docs/guides/installing/downloading-streamplace.md @@ -54,6 +54,18 @@ SP_TLS_CERT=/tls/tls.crt SP_TLS_KEY=/tls/tls.key ``` +### Reusing a former PDS hostname + +A Streamplace node does not host the accounts from a PDS that previously ran at +the same hostname. If their DID documents still point there, repository reads +return HTTP 508 (Loop Detected) rather than repeatedly proxying back to the node. +Restore the PDS or update those accounts' PDS endpoints to make their records +available again. + +Repository reads forwarded by Streamplace carry `X-Streamplace-Repo-Proxy`. +Reverse proxies must preserve this header: a node can serve its own repository +to a forwarded request, but will not forward that request again. + ### Docker Running Streamplace from a Docker image works great except for Docker diff --git a/pkg/spxrpc/com_atproto_repo_test.go b/pkg/spxrpc/com_atproto_repo_test.go new file mode 100644 index 000000000..ebadc5669 --- /dev/null +++ b/pkg/spxrpc/com_atproto_repo_test.go @@ -0,0 +1,129 @@ +package spxrpc + +import ( + "context" + "encoding/json" + "net" + "net/http" + "net/http/httptest" + "net/url" + "sync/atomic" + "testing" + "time" + + "github.com/labstack/echo/v4" + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/aqhttp" + "stream.place/streamplace/pkg/config" +) + +// A former user's DID still points at a PDS hostname now serving Streamplace. +// Exercise DID resolution and the real HTTP proxy, not a mocked upstream call. +func TestRepoProxyLoop(t *testing.T) { + for _, method := range []string{"getRecord", "listRecords", "describeRepo"} { + t.Run(method, func(t *testing.T) { + s := &Server{cli: &config.CLI{ServerHost: "node.example", BroadcasterHost: "node.example"}} + e := echo.New() + e.Use(s.ContextPreservingMiddleware()) + require.NoError(t, s.RegisterHandlersComatproto(e)) + + var did, service string + var requests atomic.Int32 + upstream := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path == "/"+did { + _ = json.NewEncoder(w).Encode(map[string]any{ + "id": did, + "service": []map[string]string{{ + "id": "#atproto_pds", "type": "AtprotoPersonalDataServer", "serviceEndpoint": service, + }}, + }) + return + } + // Bound the broken implementation so a regression never exhausts + // sockets or leaves recursive requests running after the test. + if requests.Add(1) > 4 { + w.WriteHeader(http.StatusLoopDetected) + return + } + e.ServeHTTP(w, r) + })) + defer upstream.Close() + service = upstream.URL + did = "did:plc:ho26ynsdw2ey7l56xx4owjzd" + repoTestClients(t, upstream) + params := url.Values{"repo": {did}, "collection": {"app.bsky.actor.profile"}, "rkey": {"self"}} + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + req, err := http.NewRequestWithContext(ctx, http.MethodGet, service+"/xrpc/com.atproto.repo."+method+"?"+params.Encode(), nil) + require.NoError(t, err) + resp, err := upstream.Client().Do(req) + require.NoError(t, err) + defer resp.Body.Close() + require.Equal(t, http.StatusLoopDetected, resp.StatusCode) + require.EqualValues(t, 2, requests.Load(), "the forwarded request must not forward again") + }) + } +} + +// A proxied read of a node's actual repo must still be served locally. Rejecting +// every marked request at ingress would break reads between Streamplace nodes. +func TestRepoProxyLocalRecord(t *testing.T) { + var e *echo.Echo + var did, service string + upstream := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path == "/.well-known/did.json" { + _ = json.NewEncoder(w).Encode(map[string]any{ + "id": did, + "service": []map[string]string{{ + "id": "#atproto_pds", "type": "AtprotoPersonalDataServer", "serviceEndpoint": service, + }}, + }) + return + } + e.ServeHTTP(w, r) + })) + defer upstream.Close() + service = upstream.URL + cli, router := newSyncTestNode(t, "node.example", "node.example") + e = router + did = cli.ServerDID() + commitTestRecord(t, cli, "place.stream.live.viewerCount", "streamer") + + repoTestClients(t, upstream) + params := url.Values{"repo": {did}, "collection": {"place.stream.live.viewerCount"}, "rkey": {"streamer"}} + req, err := http.NewRequest(http.MethodGet, service+"/xrpc/com.atproto.repo.getRecord?"+params.Encode(), nil) + require.NoError(t, err) + req.Header.Set("X-Streamplace-Repo-Proxy", "true") + resp, err := upstream.Client().Do(req) + require.NoError(t, err) + defer resp.Body.Close() + require.Equal(t, http.StatusOK, resp.StatusCode) + var record struct { + URI string `json:"uri"` + Value struct { + Count int `json:"count"` + } `json:"value"` + } + require.NoError(t, json.NewDecoder(resp.Body).Decode(&record)) + require.Equal(t, "at://"+did+"/place.stream.live.viewerCount/streamer", record.URI) + require.Equal(t, 7, record.Value.Count) +} + +func repoTestClients(t *testing.T, upstream *httptest.Server) { + t.Helper() + originalClient, originalDefault := aqhttp.Client, http.DefaultClient + aqhttp.Client = *upstream.Client() + // oatproxy's service resolver currently uses http.DefaultClient rather + // than its client argument. Route its real HTTPS lookups to this fixture. + transport := upstream.Client().Transport.(*http.Transport).Clone() + transport.TLSClientConfig = transport.TLSClientConfig.Clone() + transport.TLSClientConfig.ServerName = "127.0.0.1" + transport.DialContext = func(ctx context.Context, network, _ string) (net.Conn, error) { + return (&net.Dialer{}).DialContext(ctx, network, upstream.Listener.Addr().String()) + } + http.DefaultClient = &http.Client{Transport: transport, Timeout: 5 * time.Second} + t.Cleanup(func() { + aqhttp.Client, http.DefaultClient = originalClient, originalDefault + transport.CloseIdleConnections() + }) +} diff --git a/pkg/spxrpc/spxrpc.go b/pkg/spxrpc/spxrpc.go index 6ad4d5f41..456302e0b 100644 --- a/pkg/spxrpc/spxrpc.go +++ b/pkg/spxrpc/spxrpc.go @@ -189,7 +189,14 @@ func (s *Server) isServerPDS(ctx context.Context) bool { return ec.Request().Host == s.cli.ServerHost } +const repoProxyHeader = "X-Streamplace-Repo-Proxy" + func makeUnauthenticatedRequest(ctx context.Context, service, method string, params map[string]interface{}, out interface{}) error { + // Only forward once: a stale PDS endpoint may route back to this node + // under an alias, or to another Streamplace node that would proxy again. + if ec, ok := ctx.Value(echoContextKey).(echo.Context); ok && ec.Request().Header.Get(repoProxyHeader) != "" { + return echo.NewHTTPError(http.StatusLoopDetected, "repository proxy loop detected") + } ctx, cancel := context.WithTimeout(ctx, 10*time.Second) defer cancel() u, err := url.Parse(fmt.Sprintf("%s/xrpc/%s", service, method)) @@ -210,6 +217,7 @@ func makeUnauthenticatedRequest(ctx context.Context, service, method string, par if err != nil { return fmt.Errorf("failed to create request: %w", err) } + req.Header.Set(repoProxyHeader, "true") resp, err := aqhttp.Client.Do(req) if err != nil { @@ -217,6 +225,9 @@ func makeUnauthenticatedRequest(ctx context.Context, service, method string, par } defer resp.Body.Close() + if resp.StatusCode == http.StatusLoopDetected { + return echo.NewHTTPError(http.StatusLoopDetected, "repository proxy loop detected") + } if resp.StatusCode != http.StatusOK { return fmt.Errorf("upstream request failed with status %d", resp.StatusCode) } -- 2.51.2