Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237package knotfeed
import ( "bytes" "log/slog" "testing"
comatproto "github.com/bluesky-social/indigo/api/atproto" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/ipfs/go-cid" "github.com/multiformats/go-multihash" cbg "github.com/whyrusleeping/cbor-gen")
func writeText(w *bytes.Buffer, s string) { if err := cbg.CborWriteHeader(w, cbg.MajTextString, uint64(len(s))); err != nil { w.Reset() } w.WriteString(s)}
func headerFrame(t *testing.T, pairs ...any) []byte { t.Helper() if len(pairs)%2 != 0 { t.Fatal("header pairs must alternate key and value") } var out bytes.Buffer if err := cbg.CborWriteHeader(&out, cbg.MajMap, uint64(len(pairs)/2)); err != nil { t.Fatalf("map header: %v", err) } for i := 0; i < len(pairs); i += 2 { key, ok := pairs[i].(string) if !ok { t.Fatal("header keys must be strings") } writeText(&out, key) switch value := pairs[i+1].(type) { case string: writeText(&out, value) case int: if value < 0 { if err := cbg.CborWriteHeader(&out, cbg.MajNegativeInt, uint64(-1-value)); err != nil { t.Fatalf("value for %q: %v", key, err) } } else { if err := cbg.CborWriteHeader(&out, cbg.MajUnsignedInt, uint64(value)); err != nil { t.Fatalf("value for %q: %v", key, err) } } case []string: if err := cbg.CborWriteHeader(&out, cbg.MajArray, uint64(len(value))); err != nil { t.Fatalf("value for %q: %v", key, err) } for _, s := range value { writeText(&out, s) } case []byte: out.Write(value) default: t.Fatalf("unsupported header value for %q", key) } } return out.Bytes()}
func testCid(t *testing.T) cid.Cid { t.Helper() sum, err := multihash.Sum([]byte("knotfeed"), multihash.SHA2_256, -1) if err != nil { t.Fatalf("multihash: %v", err) } return cid.NewCidV1(cid.DagCBOR, sum)}
func TestDecodeReadsACommitFrameWithoutBlocks(t *testing.T) { var payload bytes.Buffer evt := comatproto.SyncSubscribeRepos_Commit{ Repo: "did:plc:scallop", Seq: 42, Rev: "3lb2xkw2qrs2j", Commit: lexutil.LexLink(testCid(t)), } if err := evt.MarshalCBOR(&payload); err != nil { t.Fatalf("MarshalCBOR: %v", err) } frame := append(headerFrame(t, "t", "#commit"), payload.Bytes()...)
message, err := Decode(frame, slog.Default()) if err != nil { t.Fatalf("Decode: %v", err) } if message.Type != TypeCommit { t.Fatalf("Type = %q, want %q", message.Type, TypeCommit) } if message.Commit == nil { t.Fatal("Decode answered no commit") } if message.Commit.Repo != "did:plc:scallop" || message.Commit.Seq != 42 || message.Commit.Rev != "3lb2xkw2qrs2j" { t.Fatalf("Commit = %+v", message.Commit) } if message.Commit.Records != nil { t.Fatalf("a frame with no ops answered records: %+v", message.Commit.Records) }}
func TestDecodeDropsUnresolvableOpsAndKeepsFrame(t *testing.T) { var payload bytes.Buffer evt := comatproto.SyncSubscribeRepos_Commit{ Repo: "did:plc:scallop", Seq: 43, Commit: lexutil.LexLink(testCid(t)), Ops: []*comatproto.SyncSubscribeRepos_RepoOp{ {Action: "create", Path: GitRefCollection.String() + "/refs~2fheads~2fmain"}, }, } if err := evt.MarshalCBOR(&payload); err != nil { t.Fatalf("MarshalCBOR: %v", err) } frame := append(headerFrame(t, "t", "#commit"), payload.Bytes()...)
message, err := Decode(frame, slog.Default()) if err != nil { t.Fatalf("Decode: %v", err) } if message.Commit == nil { t.Fatal("Decode answered no commit") } if len(message.Commit.Records) != 0 { t.Fatalf("a frame without a car answered records: %+v", message.Commit.Records) } if message.Commit.Seq != 43 { t.Fatalf("Seq = %d, want 43", message.Commit.Seq) }}
func errorFrameBody(t *testing.T, name, detail string) []byte { t.Helper() var body bytes.Buffer if err := cbg.CborWriteHeader(&body, cbg.MajMap, 2); err != nil { t.Fatalf("map header: %v", err) } writeText(&body, "error") writeText(&body, name) writeText(&body, "message") writeText(&body, detail) return body.Bytes()}
func errorFrame(t *testing.T, name, detail string) []byte { t.Helper() return append(headerFrame(t, "op", -1), errorFrameBody(t, name, detail)...)}
func TestDecodeReadsErrorFrames(t *testing.T) { for _, tt := range []struct { name string frame []byte wantError string wantDetail string }{ {"two-item frame", errorFrame(t, "FutureCursor", "the cursor is ahead of the knot"), "FutureCursor", "the cursor is ahead of the knot"}, {"single-map frame", headerFrame(t, "op", -1, "t", "#error", "error", "FutureCursor", "message", "the cursor is ahead of the knot"), "FutureCursor", "the cursor is ahead of the knot"}, {"body wins over the header", append( headerFrame(t, "op", -1, "error", "StaleHeader", "message", "from the header"), errorFrameBody(t, "FutureCursor", "from the body")..., ), "FutureCursor", "from the body"}, } { t.Run(tt.name, func(t *testing.T) { message, err := Decode(tt.frame, slog.Default()) if err != nil { t.Fatalf("Decode: %v", err) } if message.Type != TypeError || message.Error != tt.wantError || message.Detail != tt.wantDetail { t.Fatalf("Message = %+v", message) } }) }}
func TestDecodeRefusesDeeplyNestedUnknownFields(t *testing.T) { var nested bytes.Buffer for range maxCborDepth + 10 { if err := cbg.CborWriteHeader(&nested, cbg.MajArray, 1); err != nil { t.Fatalf("array header: %v", err) } } if err := cbg.CborWriteHeader(&nested, cbg.MajUnsignedInt, 0); err != nil { t.Fatalf("leaf: %v", err) } if _, err := Decode(headerFrame(t, "t", "#account", "extra", nested.Bytes()), slog.Default()); err == nil { t.Fatal("Decode accepted a frame nesting unknown fields past the depth budget") }}
func TestDecodeReadsAnInfoFrame(t *testing.T) { var payload bytes.Buffer evt := comatproto.SyncSubscribeRepos_Info{ Name: "OutdatedCursor", } if err := evt.MarshalCBOR(&payload); err != nil { t.Fatalf("MarshalCBOR: %v", err) } frame := append(headerFrame(t, "t", "#info"), payload.Bytes()...)
message, err := Decode(frame, slog.Default()) if err != nil { t.Fatalf("Decode: %v", err) } if message.Type != TypeInfo || message.InfoName != "OutdatedCursor" { t.Fatalf("Message = %+v", message) }}
func TestDecodeSkipsUnknownHeaderFields(t *testing.T) { var payload bytes.Buffer evt := comatproto.SyncSubscribeRepos_Account{Did: "did:plc:scallop", Active: true} if err := evt.MarshalCBOR(&payload); err != nil { t.Fatalf("MarshalCBOR: %v", err) } frame := append(headerFrame(t, "t", "#account", "extra", []string{"a", "b"}), payload.Bytes()...)
message, err := Decode(frame, slog.Default()) if err != nil { t.Fatalf("Decode: %v", err) } if message.Type != TypeAccount { t.Fatalf("Type = %q, want %q", message.Type, TypeAccount) }}
func TestDecodeRefusesFrameNamingNeitherTypeOrError(t *testing.T) { frame := headerFrame(t, "ops", []string{"nothing"}) if _, err := Decode(frame, slog.Default()); err == nil { t.Fatal("Decode accepted a frame naming neither a type or an error") }}