From 34d7157cc275616826b32ab28e78a177883e6820 Mon Sep 17 00:00:00 2001 From: Lewis Date: Thu, 6 Aug 2026 14:43:27 +0300 Subject: [PATCH] repoverify,xrpc,appview: reject repo whose knot won't ID an owner Lewis: May this revision serve well! --- appview/ingester.go | 2 +- appview/ingester_repo.go | 42 ++++++++-- appview/ingester_repo_test.go | 123 +++++++++++++++------------- appview/knots/knots.go | 5 ++ appview/spindles/spindles.go | 5 ++ cmd/zoekt-tngl-indexserver/index.go | 8 +- knotmirror/xrpc/proxy.go | 12 ++- repoverify/verify.go | 67 ++++++++++++--- repoverify/verify_test.go | 66 ++++++++++++++- spindle/server.go | 6 +- xrpc/xrpcclient/xrpc.go | 12 +++ xrpc/xrpcclient/xrpc_test.go | 43 ++++++++++ 12 files changed, 309 insertions(+), 82 deletions(-) create mode 100644 xrpc/xrpcclient/xrpc_test.go diff --git a/appview/ingester.go b/appview/ingester.go index cf2d6037..48e7cc1d 100644 --- a/appview/ingester.go +++ b/appview/ingester.go @@ -131,7 +131,7 @@ func (i *Ingester) Ingest() processFunc { } if err != nil { - l.Warn("failed to ingest record, skipping", "err", err) + l.Error("failed to ingest record, dropping it without retry", "err", err) } return nil diff --git a/appview/ingester_repo.go b/appview/ingester_repo.go index 3b5924ad..ce090769 100644 --- a/appview/ingester_repo.go +++ b/appview/ingester_repo.go @@ -9,7 +9,9 @@ import ( "log/slog" "slices" "strings" + "time" + "github.com/avast/retry-go/v4" "github.com/bluesky-social/indigo/atproto/syntax" jmodels "github.com/bluesky-social/jetstream/pkg/models" "tangled.org/core/api/tangled" @@ -17,6 +19,7 @@ import ( "tangled.org/core/appview/models" "tangled.org/core/orm" "tangled.org/core/repoident" + "tangled.org/core/repoverify" ) func (i *Ingester) ingestRepo(ctx context.Context, e *jmodels.Event, l *slog.Logger) error { @@ -418,6 +421,14 @@ func derefString(s *string) string { return *s } +const ownershipVerifyAttempts = 3 + +var rejectionReasons = map[repoverify.Answer]string{ + repoverify.AnswerNoRoute: "the knot doesn't serve describeRepo, so upgrade it to 1.14 or later", + repoverify.AnswerAbsent: "the knot doesn't host this repoDid", + repoverify.AnswerUnset: "the knot didn't identify an owner for this repoDid", +} + func (i *Ingester) verifyOwnership(ctx context.Context, l *slog.Logger, repoDid, eventDid, recordKnot string) (bool, error) { if i.Verifier == nil { return false, fmt.Errorf("ingester has no repo ownership verifier configured") @@ -427,18 +438,37 @@ func (i *Ingester) verifyOwnership(ctx context.Context, l *slog.Logger, repoDid, l.Warn("rejecting repo event: invalid repoDid on record", "repoDid", repoDid, "err", err) return false, nil } - result, err := i.Verifier(ctx, rd) + + verifyCtx, cancel := context.WithTimeout(ctx, 15*time.Second) + defer cancel() + + var result repoverify.Result + err = retry.Do( + func() (attemptErr error) { + result, attemptErr = i.Verifier(verifyCtx, rd) + return + }, + retry.Context(verifyCtx), + retry.Attempts(ownershipVerifyAttempts), + retry.Delay(500*time.Millisecond), + retry.DelayType(retry.FixedDelay), + retry.RetryIf(repoverify.Retriable), + retry.LastErrorOnly(true), + ) if err != nil { return false, fmt.Errorf("verify repo ownership: %w", err) } - if result.OwnerDid == "" { - l.Warn("knot lacks RepoDescribeRepo, skipping owner check; upgrade knot to 1.14+", - "repoDid", repoDid, "knot", result.KnotURL.String()) - } else if result.OwnerDid.String() != eventDid { + ownership, ok := result.Ownership() + if !ok { + l.Warn("rejecting repo event", "reason", rejectionReasons[result.Answer()], + "repoDid", repoDid, "knot", result.KnotURL.String(), "answer", result.Answer()) + return false, nil + } + if ownership.OwnerDid.String() != eventDid { l.Warn("rejecting repo event: owner mismatch", "repoDid", repoDid, "claimedOwner", eventDid, - "knotOwner", result.OwnerDid.String(), + "knotOwner", ownership.OwnerDid.String(), "knot", result.KnotURL.String(), ) return false, nil diff --git a/appview/ingester_repo_test.go b/appview/ingester_repo_test.go index e597c3e8..73d97936 100644 --- a/appview/ingester_repo_test.go +++ b/appview/ingester_repo_test.go @@ -34,20 +34,34 @@ func acceptOwner(t *testing.T, e *jmodels.Event) repoverify.Verifier { t.Helper() knot := mustKnotURL(t, "https://knot.example") return func(_ context.Context, repoDid repoident.RepoDid) (repoverify.Result, error) { - return repoverify.Result{ - RepoDid: repoDid, + return repoverify.Owned(repoDid, knot, repoverify.Ownership{ OwnerDid: repoident.OwnerDid(e.Did), - KnotURL: knot, - }, nil + }), nil } } +func ownedBy(t *testing.T, repoDid, ownerDid string) repoverify.Result { + t.Helper() + return repoverify.Owned( + repoident.RepoDid(repoDid), + mustKnotURL(t, "https://knot.example"), + repoverify.Ownership{OwnerDid: repoident.OwnerDid(ownerDid)}, + ) +} + func stubVerifier(result repoverify.Result, err error) repoverify.Verifier { return func(_ context.Context, _ repoident.RepoDid) (repoverify.Result, error) { return result, err } } +func countingVerifier(result repoverify.Result, err error, attempts *int) repoverify.Verifier { + return func(_ context.Context, _ repoident.RepoDid) (repoverify.Result, error) { + *attempts++ + return result, err + } +} + type spyNotifier struct { notify.BaseNotifier creates int @@ -559,11 +573,7 @@ func TestIngestRepo_CreateSquatRejected(t *testing.T) { RepoDid: ptr("did:plc:akshays-repo"), }) - withVerifier(ing, stubVerifier(repoverify.Result{ - RepoDid: "did:plc:akshays-repo", - OwnerDid: "did:plc:akshay", - KnotURL: mustKnotURL(t, "https://knot.example"), - }, nil)) + withVerifier(ing, stubVerifier(ownedBy(t, "did:plc:akshays-repo", "did:plc:akshay"), nil)) if err := ing.ingestRepo(context.Background(), e, ing.Logger); err != nil { t.Fatalf("ingestRepo: %v", err) @@ -589,11 +599,7 @@ func TestIngestRepo_CreateHijackExistingRepoRejected(t *testing.T) { RepoDid: ptr("did:plc:akshays-repo"), }) - withVerifier(ing, stubVerifier(repoverify.Result{ - RepoDid: "did:plc:akshays-repo", - OwnerDid: "did:plc:akshay", - KnotURL: mustKnotURL(t, "https://knot.example"), - }, nil)) + withVerifier(ing, stubVerifier(ownedBy(t, "did:plc:akshays-repo", "did:plc:akshay"), nil)) if err := ing.ingestRepo(context.Background(), e, ing.Logger); err != nil { t.Fatalf("ingestRepo: %v", err) @@ -618,11 +624,7 @@ func TestIngestRepo_CreateRenameIgnoresRkeyDrift(t *testing.T) { RepoDid: ptr("did:plc:akshays-repo"), }) - withVerifier(ing, stubVerifier(repoverify.Result{ - RepoDid: "did:plc:akshays-repo", - OwnerDid: "did:plc:akshay", - KnotURL: mustKnotURL(t, "https://knot.example"), - }, nil)) + withVerifier(ing, stubVerifier(ownedBy(t, "did:plc:akshays-repo", "did:plc:akshay"), nil)) if err := ing.ingestRepo(context.Background(), e, ing.Logger); err != nil { t.Fatalf("ingestRepo: %v", err) @@ -637,22 +639,49 @@ func TestIngestRepo_CreateRenameIgnoresRkeyDrift(t *testing.T) { } } -func TestIngestRepo_CreateVerifierTransientErrorPropagates(t *testing.T) { - ing, spy := newTestIngester(t) - - e := makeEvent(t, jmodels.CommitOperationCreate, "did:plc:akshay", "myrepo", tangled.Repo{ - Knot: "knot.example", - RepoDid: ptr("did:plc:akshays-repo"), - }) - - withVerifier(ing, stubVerifier(repoverify.Result{}, errors.New("knot unreachable"))) - - err := ing.ingestRepo(context.Background(), e, ing.Logger) - if err == nil { - t.Fatalf("expected error on transient verifier failure, got nil") - } - if spy.creates != 0 { - t.Errorf("NewRepo called %d times despite verifier error", spy.creates) +func TestIngestRepo_CreateTakesOnlyAnIdentifiedOwnerAndRetriesOnlyTheUndecided(t *testing.T) { + repoDid := repoident.RepoDid("did:plc:akshays-repo") + knot := mustKnotURL(t, "https://knot.example") + refused := func(a repoverify.Answer) repoverify.Result { return repoverify.Refused(repoDid, knot, a) } + cases := map[string]struct { + result repoverify.Result + verifyErr error + wantErr bool + wantAttempts int + wantCreates int + }{ + "we'll retry an unreachable knot before the failure propagates": {verifyErr: errors.New("knot unreachable"), wantErr: true, wantAttempts: ownershipVerifyAttempts}, + "we won't retry a knot that already answered": {verifyErr: repoverify.ErrKnotAnswer, wantErr: true, wantAttempts: 1}, + "a knot without describeRepo can't confirm the claimed owner": {result: refused(repoverify.AnswerNoRoute), wantAttempts: 1}, + "we won't create when the knot doesn't host the repoDid": {result: refused(repoverify.AnswerAbsent), wantAttempts: 1}, + "we'll fail closed on a verifier that doesn't pong": {wantAttempts: 1}, + "we'll take an identified owner on the first round-trip": {result: ownedBy(t, repoDid.String(), "did:plc:akshay"), wantAttempts: 1, wantCreates: 1}, + } + for name, tc := range cases { + t.Run(name, func(t *testing.T) { + ing, spy := newTestIngester(t) + e := makeEvent(t, jmodels.CommitOperationCreate, "did:plc:akshay", "myrepo", tangled.Repo{ + Knot: "knot.example", + RepoDid: ptr(repoDid.String()), + }) + + attempts := 0 + withVerifier(ing, countingVerifier(tc.result, tc.verifyErr, &attempts)) + + if err := ing.ingestRepo(context.Background(), e, ing.Logger); (err != nil) != tc.wantErr { + t.Fatalf("ingestRepo err = %v, wantErr %v", err, tc.wantErr) + } + if attempts != tc.wantAttempts { + t.Errorf("verifier called %d times, want %d", attempts, tc.wantAttempts) + } + if spy.creates != tc.wantCreates { + t.Errorf("NewRepo called %d times, want %d", spy.creates, tc.wantCreates) + } + _, err := db.GetRepo(ing.Db, orm.FilterEq("did", "did:plc:akshay"), orm.FilterEq("rkey", "myrepo")) + if got, want := err == nil, tc.wantCreates == 1; got != want { + t.Fatalf("repo row exists = %v, want %v (err=%v)", got, want, err) + } + }) } } @@ -666,11 +695,7 @@ func TestIngestRepo_UpdateRejectsOwnerMismatch(t *testing.T) { RepoDid: ptr("did:plc:akshays-repo"), }) - withVerifier(ing, stubVerifier(repoverify.Result{ - RepoDid: "did:plc:akshays-repo", - OwnerDid: "did:plc:akshay", - KnotURL: mustKnotURL(t, "https://knot.example"), - }, nil)) + withVerifier(ing, stubVerifier(ownedBy(t, "did:plc:akshays-repo", "did:plc:akshay"), nil)) if err := ing.ingestRepo(context.Background(), e, ing.Logger); err != nil { t.Fatalf("ingestRepo: %v", err) @@ -732,11 +757,7 @@ func TestIngestRepo_CreateRejectsKnotMismatch(t *testing.T) { RepoDid: ptr("did:plc:akshays-repo"), }) - withVerifier(ing, stubVerifier(repoverify.Result{ - RepoDid: "did:plc:akshays-repo", - OwnerDid: "did:plc:akshay", - KnotURL: mustKnotURL(t, "https://knot.example"), - }, nil)) + withVerifier(ing, stubVerifier(ownedBy(t, "did:plc:akshays-repo", "did:plc:akshay"), nil)) if err := ing.ingestRepo(context.Background(), e, ing.Logger); err != nil { t.Fatalf("ingestRepo: %v", err) @@ -762,11 +783,7 @@ func TestIngestRepo_UpdateRejectsKnotMismatch(t *testing.T) { RepoDid: ptr("did:plc:akshays-repo"), }) - withVerifier(ing, stubVerifier(repoverify.Result{ - RepoDid: "did:plc:akshays-repo", - OwnerDid: "did:plc:akshay", - KnotURL: mustKnotURL(t, "https://knot.example"), - }, nil)) + withVerifier(ing, stubVerifier(ownedBy(t, "did:plc:akshays-repo", "did:plc:akshay"), nil)) if err := ing.ingestRepo(context.Background(), e, ing.Logger); err != nil { t.Fatalf("ingestRepo: %v", err) @@ -790,11 +807,7 @@ func TestIngestRepo_UpdateRejectsRepoDidMutation(t *testing.T) { RepoDid: ptr("did:plc:other-repo"), }) - withVerifier(ing, stubVerifier(repoverify.Result{ - RepoDid: "did:plc:other-repo", - OwnerDid: "did:plc:akshay", - KnotURL: mustKnotURL(t, "https://knot.example"), - }, nil)) + withVerifier(ing, stubVerifier(ownedBy(t, "did:plc:other-repo", "did:plc:akshay"), nil)) if err := ing.ingestRepo(context.Background(), e, ing.Logger); err != nil { t.Fatalf("ingestRepo: %v", err) diff --git a/appview/knots/knots.go b/appview/knots/knots.go index d47506b6..b7357389 100644 --- a/appview/knots/knots.go +++ b/appview/knots/knots.go @@ -403,6 +403,11 @@ func (k *Knots) retry(w http.ResponseWriter, r *http.Request) { return } + if errors.Is(err, xrpcclient.ErrXrpcNotFound) { + k.Pages.Notice(w, noticeId, "Failed to verify knot, its owner query returned 404. Check that this domain reaches your knot.") + return + } + if e, ok := err.(*serververify.OwnerMismatch); ok { k.Pages.Notice(w, noticeId, e.Error()) return diff --git a/appview/spindles/spindles.go b/appview/spindles/spindles.go index 6840fb39..4524759e 100644 --- a/appview/spindles/spindles.go +++ b/appview/spindles/spindles.go @@ -386,6 +386,11 @@ func (s *Spindles) retry(w http.ResponseWriter, r *http.Request) { return } + if errors.Is(err, xrpcclient.ErrXrpcNotFound) { + s.Pages.Notice(w, noticeId, "Failed to verify spindle, its owner query returned 404. Check that this instance reaches your spindle.") + return + } + if e, ok := err.(*serververify.OwnerMismatch); ok { s.Pages.Notice(w, noticeId, e.Error()) return diff --git a/cmd/zoekt-tngl-indexserver/index.go b/cmd/zoekt-tngl-indexserver/index.go index 6b0e5a70..cd16082f 100644 --- a/cmd/zoekt-tngl-indexserver/index.go +++ b/cmd/zoekt-tngl-indexserver/index.go @@ -63,11 +63,15 @@ func loadRepo(ctx context.Context, cfg *Config, dir identity.Directory, repoDID if err != nil { return nil, err } + ownership, ok := described.Ownership() + if !ok { + return nil, fmt.Errorf("knot %s answered %s for repoDid %s", knot, described.Answer(), repoDID) + } return &Repo{ Did: repoDID, - Owner: described.OwnerDid, - Slug: described.Rkey, + Owner: ownership.OwnerDid, + Slug: ownership.Rkey, Knot: knot, }, nil } diff --git a/knotmirror/xrpc/proxy.go b/knotmirror/xrpc/proxy.go index 2794979b..7ba7e2f1 100644 --- a/knotmirror/xrpc/proxy.go +++ b/knotmirror/xrpc/proxy.go @@ -95,11 +95,17 @@ func (x *Xrpc) resolveKnot(ctx context.Context, repoDid syntax.DID) (*knotInfo, return &knotInfo{baseURL: knotURL, repoIdentifier: repoDid.String()}, nil } + ownership, ok := described.Ownership() + if !ok { + x.logger.Warn("serving without a metadata upsert, since describeRepo didn't identify an owner", "knot", knotURL, "repo", repoDid, "answer", described.Answer()) + return &knotInfo{baseURL: knotURL, repoIdentifier: repoDid.String()}, nil + } + go func() { pending := &models.Repo{ - Did: syntax.DID(described.OwnerDid), - Rkey: described.Rkey, - Name: string(described.Rkey), + Did: syntax.DID(ownership.OwnerDid), + Rkey: ownership.Rkey, + Name: string(ownership.Rkey), KnotDomain: knotURL, RepoDid: repoDid, State: models.RepoStatePending, diff --git a/repoverify/verify.go b/repoverify/verify.go index dc492b69..24e158a8 100644 --- a/repoverify/verify.go +++ b/repoverify/verify.go @@ -9,6 +9,7 @@ import ( "github.com/bluesky-social/indigo/atproto/syntax" indigoxrpc "github.com/bluesky-social/indigo/xrpc" + "github.com/samber/lo" "tangled.org/core/api/tangled" "tangled.org/core/hostutil" "tangled.org/core/idresolver" @@ -16,21 +17,67 @@ import ( "tangled.org/core/xrpc/xrpcclient" ) -type Result struct { - RepoDid repoident.RepoDid +type Answer string + +const ( + AnswerUnset Answer = "unset" + AnswerOwner Answer = "owner" + AnswerNoRoute Answer = "noRoute" + AnswerAbsent Answer = "absent" +) + +type Ownership struct { OwnerDid repoident.OwnerDid - KnotURL repoident.KnotURL - // Rkey of the sh.tangled.repo record tracked by the knot; empty when the - // knot does not support describeRepo. - Rkey syntax.RecordKey + Rkey syntax.RecordKey +} + +type Result struct { + RepoDid repoident.RepoDid + KnotURL repoident.KnotURL + + answer Answer + ownership Ownership +} + +func Owned(repoDid repoident.RepoDid, knot repoident.KnotURL, ownership Ownership) Result { + return Result{RepoDid: repoDid, KnotURL: knot, answer: AnswerOwner, ownership: ownership} +} + +func Refused(repoDid repoident.RepoDid, knot repoident.KnotURL, answer Answer) Result { + return Result{RepoDid: repoDid, KnotURL: knot, answer: answer} +} + +func (r Result) Answer() Answer { + return lo.Ternary(r.answer == "", AnswerUnset, r.answer) +} + +func (r Result) Ownership() (Ownership, bool) { + return r.ownership, r.answer == AnswerOwner } var ErrKnotAnswer = errors.New("knot returned an invalid describeRepo answer") +const repoNotFound = "RepoNotFound" + +var terminalErrors = []error{ + ErrKnotAnswer, + repoident.ErrNoKnotService, + repoident.ErrNilIdentity, + xrpcclient.ErrXrpcUnauthorized, + xrpcclient.ErrXrpcForbidden, +} + +func Retriable(err error) bool { + return err != nil && !lo.ContainsBy(terminalErrors, func(t error) bool { return errors.Is(err, t) }) +} + func Describe(ctx context.Context, httpClient *http.Client, knot repoident.KnotURL, repoDid repoident.RepoDid) (Result, error) { client := &indigoxrpc.Client{Host: knot.String(), Client: httpClient} out, err := tangled.RepoDescribeRepo(ctx, client, repoDid.String()) if xrpcErr := xrpcclient.HandleXrpcErr(err); xrpcErr != nil { + if errors.Is(xrpcErr, xrpcclient.ErrXrpcUnsupported) || errors.Is(xrpcErr, xrpcclient.ErrXrpcNotFound) { + return Refused(repoDid, knot, lo.Ternary(xrpcclient.ErrorName(err) == repoNotFound, AnswerAbsent, AnswerNoRoute)), nil + } return Result{}, fmt.Errorf("describeRepo on %s: %w (%v)", knot, xrpcErr, err) } if out.RepoDid != repoDid.String() { @@ -44,7 +91,7 @@ func Describe(ctx context.Context, httpClient *http.Client, knot repoident.KnotU if err != nil { return Result{}, fmt.Errorf("%w: knot %s returned rkey %q: %w", ErrKnotAnswer, knot, out.Rkey, err) } - return Result{RepoDid: repoDid, OwnerDid: ownerDid, KnotURL: knot, Rkey: rkey}, nil + return Owned(repoDid, knot, Ownership{OwnerDid: ownerDid, Rkey: rkey}), nil } type Verifier func(ctx context.Context, repoDid repoident.RepoDid) (Result, error) @@ -69,10 +116,6 @@ func New(resolver *idresolver.Resolver, dev bool) Verifier { return Result{}, fmt.Errorf("repoDid %s: %w", repoDid, err) } - result, err := Describe(ctx, httpClient, knot, repoDid) - if errors.Is(err, xrpcclient.ErrXrpcUnsupported) { - return Result{RepoDid: repoDid, KnotURL: knot}, nil - } - return result, err + return Describe(ctx, httpClient, knot, repoDid) } } diff --git a/repoverify/verify_test.go b/repoverify/verify_test.go index 14fadc06..00109eaf 100644 --- a/repoverify/verify_test.go +++ b/repoverify/verify_test.go @@ -4,15 +4,18 @@ import ( "context" "encoding/json" "errors" + "fmt" "net/http" "net/http/httptest" "strings" "testing" "github.com/bluesky-social/indigo/atproto/identity" + "github.com/samber/lo" "tangled.org/core/api/tangled" "tangled.org/core/idresolver" "tangled.org/core/repoident" + "tangled.org/core/xrpc/xrpcclient" ) const ( @@ -55,14 +58,53 @@ func TestNew_DevModeAcceptsHttpKnotEndpointAndStripsThePath(t *testing.T) { if err != nil { t.Fatalf("dev mode should accept an http knot endpoint: %v", err) } - if result.OwnerDid.String() != testOwnerDid { - t.Errorf("OwnerDid = %q, want %q", result.OwnerDid, testOwnerDid) + ownership, ok := result.Ownership() + if !ok { + t.Fatalf("Answer = %s, want %s", result.Answer(), AnswerOwner) + } + if ownership.OwnerDid.String() != testOwnerDid { + t.Errorf("OwnerDid = %q, want %q", ownership.OwnerDid, testOwnerDid) } if result.KnotURL.String() != srv.URL { t.Errorf("KnotURL = %q, want %q", result.KnotURL, srv.URL) } } +func TestNew_AnswersFromA404(t *testing.T) { + cases := map[string]struct { + body string + want Answer + }{ + "only RepoNotFound will refute the repoDid": {`{"error":"RepoNotFound","message":"no such repo"}`, AnswerAbsent}, + "a knot without the route will answer with an empty body": {"", AnswerNoRoute}, + "a 404 body without an error name can't claim the route": {`{"message":"not found"}`, AnswerNoRoute}, + "a proxy json-ifying its own 404 will never answer for the knot": {`{"error":"NotFound","message":"no upstream"}`, AnswerNoRoute}, + } + for name, tc := range cases { + t.Run(name, func(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + if tc.body != "" { + w.Header().Set("Content-Type", "application/json") + } + w.WriteHeader(http.StatusNotFound) + _, _ = w.Write([]byte(tc.body)) + })) + t.Cleanup(srv.Close) + + result, err := verifyKnot(t, srv.URL, true) + if err != nil { + t.Fatalf("a 404 from describeRepo must be an answer: %v", err) + } + if result.Answer() != tc.want { + t.Errorf("Answer = %s, want %s", result.Answer(), tc.want) + } + if _, ok := result.Ownership(); ok { + t.Error("a 404 won't identify an owner") + } + }) + } +} + func TestNew_Rejections(t *testing.T) { target := describeRepoServer(t, testRepoDid) otherRepo := describeRepoServer(t, "did:plc:anemone") @@ -97,3 +139,23 @@ func TestNew_Rejections(t *testing.T) { }) } } + +func TestRetriable(t *testing.T) { + retriable := []error{ + errors.New("connection refused"), + fmt.Errorf("describeRepo: %w", xrpcclient.ErrXrpcFailed), + } + terminal := append([]error{nil}, lo.Map(terminalErrors, func(e error, _ int) error { + return fmt.Errorf("describeRepo: %w", e) + })...) + for _, err := range retriable { + if !Retriable(err) { + t.Errorf("Retriable(%v) = false, want true, since the knot may answer the next call", err) + } + } + for _, err := range terminal { + if Retriable(err) { + t.Errorf("Retriable(%v) = true, want false, since the answer won't change on a retry", err) + } + } +} diff --git a/spindle/server.go b/spindle/server.go index 51dcc119..f3136d87 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -581,7 +581,11 @@ func (s *Spindle) resolveSourceRepoInfo(ctx context.Context, repoDid syntax.DID) if err != nil { return nil, fmt.Errorf("verify sourceRepo %s: %w", repoDid, err) } - return s.buildTriggerRepoFrom(ctx, res.KnotURL.Host(), res.OwnerDid.String(), res.Rkey.String(), repoDid.String()), nil + ownership, ok := res.Ownership() + if !ok { + return nil, fmt.Errorf("verify sourceRepo %s: knot %s answered %s", repoDid, res.KnotURL, res.Answer()) + } + return s.buildTriggerRepoFrom(ctx, res.KnotURL.Host(), ownership.OwnerDid.String(), ownership.Rkey.String(), repoDid.String()), nil } // runPipeline compiles and enqueues the pipeline for the given revision. diff --git a/xrpc/xrpcclient/xrpc.go b/xrpc/xrpcclient/xrpc.go index 3d492e07..7910d7a3 100644 --- a/xrpc/xrpcclient/xrpc.go +++ b/xrpc/xrpcclient/xrpc.go @@ -9,6 +9,7 @@ import ( var ( ErrXrpcUnsupported = errors.New("xrpc not supported on this knot") + ErrXrpcNotFound = errors.New("not found on this knot") ErrXrpcUnauthorized = errors.New("unauthorized xrpc request") ErrXrpcForbidden = errors.New("forbidden xrpc request") ErrXrpcFailed = errors.New("xrpc request failed") @@ -28,6 +29,9 @@ func HandleXrpcErr(err error) error { switch xrpcerr.StatusCode { case http.StatusNotFound: + if ErrorName(err) != "" { + return ErrXrpcNotFound + } return ErrXrpcUnsupported case http.StatusUnauthorized: return ErrXrpcUnauthorized @@ -37,3 +41,11 @@ func HandleXrpcErr(err error) error { return ErrXrpcFailed } } + +func ErrorName(err error) string { + var named *indigoxrpc.XRPCError + if !errors.As(err, &named) { + return "" + } + return named.ErrStr +} diff --git a/xrpc/xrpcclient/xrpc_test.go b/xrpc/xrpcclient/xrpc_test.go new file mode 100644 index 00000000..e499dc91 --- /dev/null +++ b/xrpc/xrpcclient/xrpc_test.go @@ -0,0 +1,43 @@ +package xrpcclient + +import ( + "errors" + "fmt" + "testing" + + indigoxrpc "github.com/bluesky-social/indigo/xrpc" +) + +func TestA404IsARouteAnswerOnlyWithAnErrorName(t *testing.T) { + cases := map[string]struct { + wrapped error + wantName string + want error + }{ + "an error name means the knot served the route": { + &indigoxrpc.XRPCError{ErrStr: "RepoNotFound"}, "RepoNotFound", ErrXrpcNotFound, + }, + "a 404 whose body won't decode means the route is missing": { + fmt.Errorf("failed to decode xrpc error message: unexpected end of JSON input"), "", ErrXrpcUnsupported, + }, + "a 404 body without an error name can't claim the route": { + &indigoxrpc.XRPCError{Message: "not found"}, "", ErrXrpcUnsupported, + }, + } + for name, tc := range cases { + t.Run(name, func(t *testing.T) { + err := &indigoxrpc.Error{StatusCode: 404, Wrapped: tc.wrapped} + if got := ErrorName(err); got != tc.wantName { + t.Errorf("ErrorName = %q, want %q", got, tc.wantName) + } + if got := HandleXrpcErr(err); !errors.Is(got, tc.want) { + t.Errorf("HandleXrpcErr = %v, want %v", got, tc.want) + } + }) + } + for _, err := range []error{errors.New("connection refused"), nil} { + if got := ErrorName(err); got != "" { + t.Errorf("ErrorName(%v) = %q, want empty, since only a lexicon error has an error name", err, got) + } + } +} -- 2.51.2