From b91be663a09ffbea44b0dd0763e27b02132515db Mon Sep 17 00:00:00 2001 From: Lewis Date: Mon, 24 Aug 2026 09:56:02 +0300 Subject: [PATCH] lexicons,api,knotfeed: swap refUpdates methods for git.ref records Lewis: May this revision serve well! --- .../notificationdeleteNotification.go | 2 +- .../notificationlistNotifications.go | 2 +- api/org_tangled/notificationupdateSeen.go | 2 +- api/tangled/cbor_gen.go | 130 +++++++++ api/tangled/gitcountRefUpdates.go | 38 --- api/tangled/gitcountRefUpdatesBy.go | 38 --- api/tangled/gitlistRefUpdates.go | 61 ---- api/tangled/gitlistRefUpdatesBy.go | 53 ---- api/tangled/gitref.go | 23 ++ api/tangled/knotsubscribeRepos.go | 34 --- cmd/cborgen/cborgen.go | 1 + knotfeed/frame.go | 260 ++++++++++++++++++ knotfeed/frame_test.go | 236 ++++++++++++++++ knotfeed/record.go | 170 ++++++++++++ knotfeed/record_test.go | 124 +++++++++ knotfeed/subscribe.go | 227 +++++++++++++++ knotfeed/subscribe_test.go | 116 ++++++++ lexicons/git/countRefUpdates.json | 39 --- lexicons/git/countRefUpdatesBy.json | 39 --- lexicons/git/listRefUpdates.json | 71 ----- lexicons/git/listRefUpdatesBy.json | 59 ---- lexicons/git/ref.json | 28 ++ lexicons/knot/subscribeRepos.json | 57 ---- .../temp/notification/deleteNotification.json | 3 +- .../temp/notification/listNotifications.json | 3 +- lexicons/temp/notification/updateSeen.json | 3 +- web/src/lib/api/count.ts | 2 - web/src/lib/api/lexicons/index.ts | 6 +- .../temp/notification/deleteNotification.ts | 4 +- .../temp/notification/listNotifications.ts | 4 +- .../tangled/temp/notification/updateSeen.ts | 4 +- .../types/sh/tangled/git/countRefUpdates.ts | 42 --- .../types/sh/tangled/git/countRefUpdatesBy.ts | 42 --- .../types/sh/tangled/git/listRefUpdates.ts | 84 ------ .../types/sh/tangled/git/listRefUpdatesBy.ts | 69 ----- .../api/lexicons/types/sh/tangled/git/ref.ts | 32 +++ .../types/sh/tangled/knot/subscribeRepos.ts | 90 ------ 37 files changed, 1360 insertions(+), 838 deletions(-) delete mode 100644 api/tangled/gitcountRefUpdates.go delete mode 100644 api/tangled/gitcountRefUpdatesBy.go delete mode 100644 api/tangled/gitlistRefUpdates.go delete mode 100644 api/tangled/gitlistRefUpdatesBy.go create mode 100644 api/tangled/gitref.go delete mode 100644 api/tangled/knotsubscribeRepos.go create mode 100644 knotfeed/frame.go create mode 100644 knotfeed/frame_test.go create mode 100644 knotfeed/record.go create mode 100644 knotfeed/record_test.go create mode 100644 knotfeed/subscribe.go create mode 100644 knotfeed/subscribe_test.go delete mode 100644 lexicons/git/countRefUpdates.json delete mode 100644 lexicons/git/countRefUpdatesBy.json delete mode 100644 lexicons/git/listRefUpdates.json delete mode 100644 lexicons/git/listRefUpdatesBy.json create mode 100644 lexicons/git/ref.json delete mode 100644 lexicons/knot/subscribeRepos.json delete mode 100644 web/src/lib/api/lexicons/types/sh/tangled/git/countRefUpdates.ts delete mode 100644 web/src/lib/api/lexicons/types/sh/tangled/git/countRefUpdatesBy.ts delete mode 100644 web/src/lib/api/lexicons/types/sh/tangled/git/listRefUpdates.ts delete mode 100644 web/src/lib/api/lexicons/types/sh/tangled/git/listRefUpdatesBy.ts create mode 100644 web/src/lib/api/lexicons/types/sh/tangled/git/ref.ts delete mode 100644 web/src/lib/api/lexicons/types/sh/tangled/knot/subscribeRepos.ts diff --git a/api/org_tangled/notificationdeleteNotification.go b/api/org_tangled/notificationdeleteNotification.go index 17b09b2dd..c9f778fdb 100644 --- a/api/org_tangled/notificationdeleteNotification.go +++ b/api/org_tangled/notificationdeleteNotification.go @@ -16,7 +16,7 @@ const ( // TempNotificationDeleteNotification_Input is the input argument to a org.tangled.temp.notification.deleteNotification call. type TempNotificationDeleteNotification_Input struct { - // uri: at-uri of the notification to delete. + // uri: notification key as returned by listNotifications; pass it back unchanged. Uri string `json:"uri" cborgen:"uri"` } diff --git a/api/org_tangled/notificationlistNotifications.go b/api/org_tangled/notificationlistNotifications.go index c07499f8b..f15fd5682 100644 --- a/api/org_tangled/notificationlistNotifications.go +++ b/api/org_tangled/notificationlistNotifications.go @@ -30,7 +30,7 @@ type TempNotificationListNotifications_Notification struct { RepoDid *string `json:"repoDid,omitempty" cborgen:"repoDid,omitempty"` // type: Notification type: repo_starred, issue_created, issue_commented, issue_closed, issue_reopen, issue_assigned, issue_unassigned, pull_created, pull_commented, pull_merged, pull_closed, pull_reopen, pull_assigned, pull_unassigned, followed, user_mentioned. Type string `json:"type" cborgen:"type"` - // uri: at-uri of this notification; the stable key for read/unread state. + // uri: notification key; pass it back to updateSeen and deleteNotification. Uri string `json:"uri" cborgen:"uri"` } diff --git a/api/org_tangled/notificationupdateSeen.go b/api/org_tangled/notificationupdateSeen.go index a2de7102d..e3d42a67f 100644 --- a/api/org_tangled/notificationupdateSeen.go +++ b/api/org_tangled/notificationupdateSeen.go @@ -17,7 +17,7 @@ const ( // TempNotificationUpdateSeen_Input is the input argument to a org.tangled.temp.notification.updateSeen call. type TempNotificationUpdateSeen_Input struct { Read bool `json:"read" cborgen:"read"` - // uri: at-uri of the notification to update. + // uri: notification key as returned by listNotifications; pass it back unchanged. Uri string `json:"uri" cborgen:"uri"` } diff --git a/api/tangled/cbor_gen.go b/api/tangled/cbor_gen.go index 7b39f0b3a..27dd6f204 100644 --- a/api/tangled/cbor_gen.go +++ b/api/tangled/cbor_gen.go @@ -4564,6 +4564,136 @@ func (t *GitOid) UnmarshalCBOR(r io.Reader) (err error) { return nil } +func (t *GitRef) MarshalCBOR(w io.Writer) error { + if t == nil { + _, err := w.Write(cbg.CborNull) + return err + } + + cw := cbg.NewCborWriter(w) + + if _, err := cw.Write([]byte{162}); err != nil { + return err + } + + // t.Sha (string) (string) + if len("sha") > 1000000 { + return xerrors.Errorf("Value in field \"sha\" was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("sha"))); err != nil { + return err + } + if _, err := cw.WriteString(string("sha")); err != nil { + return err + } + + if len(t.Sha) > 1000000 { + return xerrors.Errorf("Value in field t.Sha was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.Sha))); err != nil { + return err + } + if _, err := cw.WriteString(string(t.Sha)); err != nil { + return err + } + + // t.LexiconTypeID (string) (string) + if len("$type") > 1000000 { + return xerrors.Errorf("Value in field \"$type\" was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("$type"))); err != nil { + return err + } + if _, err := cw.WriteString(string("$type")); err != nil { + return err + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("sh.tangled.git.ref"))); err != nil { + return err + } + if _, err := cw.WriteString(string("sh.tangled.git.ref")); err != nil { + return err + } + return nil +} + +func (t *GitRef) UnmarshalCBOR(r io.Reader) (err error) { + *t = GitRef{} + + cr := cbg.NewCborReader(r) + + maj, extra, err := cr.ReadHeader() + if err != nil { + return err + } + defer func() { + if err == io.EOF { + err = io.ErrUnexpectedEOF + } + }() + + if maj != cbg.MajMap { + return fmt.Errorf("cbor input should be of type map") + } + + if extra > cbg.MaxLength { + return fmt.Errorf("GitRef: map struct too large (%d)", extra) + } + + n := extra + + nameBuf := make([]byte, 5) + for i := uint64(0); i < n; i++ { + nameLen, ok, err := cbg.ReadFullStringIntoBuf(cr, nameBuf, 1000000) + if err != nil { + return err + } + + if !ok { + // Field doesn't exist on this type, so ignore it + if err := cbg.ScanForLinks(cr, func(cid.Cid) {}); err != nil { + return err + } + continue + } + + switch string(nameBuf[:nameLen]) { + // t.Sha (string) (string) + case "sha": + + { + sval, err := cbg.ReadStringWithMax(cr, 1000000) + if err != nil { + return err + } + + t.Sha = string(sval) + } + // t.LexiconTypeID (string) (string) + case "$type": + + { + sval, err := cbg.ReadStringWithMax(cr, 1000000) + if err != nil { + return err + } + + t.LexiconTypeID = string(sval) + } + + default: + // Field doesn't exist on this type, so ignore it + if err := cbg.ScanForLinks(r, func(cid.Cid) {}); err != nil { + return err + } + } + } + + return nil +} func (t *GitRefUpdate) MarshalCBOR(w io.Writer) error { if t == nil { _, err := w.Write(cbg.CborNull) diff --git a/api/tangled/gitcountRefUpdates.go b/api/tangled/gitcountRefUpdates.go deleted file mode 100644 index 902bf4449..000000000 --- a/api/tangled/gitcountRefUpdates.go +++ /dev/null @@ -1,38 +0,0 @@ -// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. - -package tangled - -// schema: sh.tangled.git.countRefUpdates - -import ( - "context" - - "github.com/bluesky-social/indigo/lex/util" -) - -const ( - GitCountRefUpdatesNSID = "sh.tangled.git.countRefUpdates" -) - -// GitCountRefUpdates_Output is the output of a sh.tangled.git.countRefUpdates call. -type GitCountRefUpdates_Output struct { - // count: Total number of matching records. - Count int64 `json:"count" cborgen:"count"` - // distinctAuthors: Number of distinct authors among the matching records. - DistinctAuthors int64 `json:"distinctAuthors" cborgen:"distinctAuthors"` -} - -// GitCountRefUpdates calls the XRPC method "sh.tangled.git.countRefUpdates". -// -// subject: Repo DID whose ref-update records to list. -func GitCountRefUpdates(ctx context.Context, c util.LexClient, subject string) (*GitCountRefUpdates_Output, error) { - var out GitCountRefUpdates_Output - - params := map[string]interface{}{} - params["subject"] = subject - if err := c.LexDo(ctx, util.Query, "", "sh.tangled.git.countRefUpdates", params, nil, &out); err != nil { - return nil, err - } - - return &out, nil -} diff --git a/api/tangled/gitcountRefUpdatesBy.go b/api/tangled/gitcountRefUpdatesBy.go deleted file mode 100644 index 1dc358dc8..000000000 --- a/api/tangled/gitcountRefUpdatesBy.go +++ /dev/null @@ -1,38 +0,0 @@ -// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. - -package tangled - -// schema: sh.tangled.git.countRefUpdatesBy - -import ( - "context" - - "github.com/bluesky-social/indigo/lex/util" -) - -const ( - GitCountRefUpdatesByNSID = "sh.tangled.git.countRefUpdatesBy" -) - -// GitCountRefUpdatesBy_Output is the output of a sh.tangled.git.countRefUpdatesBy call. -type GitCountRefUpdatesBy_Output struct { - // count: Total number of matching records. - Count int64 `json:"count" cborgen:"count"` - // distinctAuthors: Number of distinct authors among the matching records. - DistinctAuthors int64 `json:"distinctAuthors" cborgen:"distinctAuthors"` -} - -// GitCountRefUpdatesBy calls the XRPC method "sh.tangled.git.countRefUpdatesBy". -// -// subject: Actor DID whose ref-update authorings to list. -func GitCountRefUpdatesBy(ctx context.Context, c util.LexClient, subject string) (*GitCountRefUpdatesBy_Output, error) { - var out GitCountRefUpdatesBy_Output - - params := map[string]interface{}{} - params["subject"] = subject - if err := c.LexDo(ctx, util.Query, "", "sh.tangled.git.countRefUpdatesBy", params, nil, &out); err != nil { - return nil, err - } - - return &out, nil -} diff --git a/api/tangled/gitlistRefUpdates.go b/api/tangled/gitlistRefUpdates.go deleted file mode 100644 index baa9f51bf..000000000 --- a/api/tangled/gitlistRefUpdates.go +++ /dev/null @@ -1,61 +0,0 @@ -// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. - -package tangled - -// schema: sh.tangled.git.listRefUpdates - -import ( - "context" - - "github.com/bluesky-social/indigo/lex/util" -) - -const ( - GitListRefUpdatesNSID = "sh.tangled.git.listRefUpdates" -) - -// GitListRefUpdates_ListItem is a "listItem" in the sh.tangled.git.listRefUpdates schema. -type GitListRefUpdates_ListItem struct { - Cid *string `json:"cid,omitempty" cborgen:"cid,omitempty"` - Uri string `json:"uri" cborgen:"uri"` - // value: Embedded sh.tangled.git.refUpdate record - Value *util.LexiconTypeDecoder `json:"value" cborgen:"value"` -} - -// GitListRefUpdates_Output is the output of a sh.tangled.git.listRefUpdates call. -type GitListRefUpdates_Output struct { - Cursor *string `json:"cursor,omitempty" cborgen:"cursor,omitempty"` - Items []*GitListRefUpdates_ListItem `json:"items" cborgen:"items"` - // total: Total items in the full list; omitted for filtered or merged views - Total *int64 `json:"total,omitempty" cborgen:"total,omitempty"` -} - -// GitListRefUpdates calls the XRPC method "sh.tangled.git.listRefUpdates". -// -// cursor: Pagination cursor -// offset: Absolute offset for random-access pagination. Mutually exclusive with cursor; offsets drift under concurrent writes, so follow up with the returned cursor. -// order: Sort direction by createdAt. -// subject: Repo DID whose ref-update records to list. -func GitListRefUpdates(ctx context.Context, c util.LexClient, cursor string, limit int64, offset int64, order string, subject string) (*GitListRefUpdates_Output, error) { - var out GitListRefUpdates_Output - - params := map[string]interface{}{} - if cursor != "" { - params["cursor"] = cursor - } - if limit != 0 { - params["limit"] = limit - } - if offset != 0 { - params["offset"] = offset - } - if order != "" { - params["order"] = order - } - params["subject"] = subject - if err := c.LexDo(ctx, util.Query, "", "sh.tangled.git.listRefUpdates", params, nil, &out); err != nil { - return nil, err - } - - return &out, nil -} diff --git a/api/tangled/gitlistRefUpdatesBy.go b/api/tangled/gitlistRefUpdatesBy.go deleted file mode 100644 index fef424ab1..000000000 --- a/api/tangled/gitlistRefUpdatesBy.go +++ /dev/null @@ -1,53 +0,0 @@ -// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. - -package tangled - -// schema: sh.tangled.git.listRefUpdatesBy - -import ( - "context" - - "github.com/bluesky-social/indigo/lex/util" -) - -const ( - GitListRefUpdatesByNSID = "sh.tangled.git.listRefUpdatesBy" -) - -// GitListRefUpdatesBy_Output is the output of a sh.tangled.git.listRefUpdatesBy call. -type GitListRefUpdatesBy_Output struct { - Cursor *string `json:"cursor,omitempty" cborgen:"cursor,omitempty"` - Items []*GitListRefUpdates_ListItem `json:"items" cborgen:"items"` - // total: Total items in the full list; omitted for filtered or merged views - Total *int64 `json:"total,omitempty" cborgen:"total,omitempty"` -} - -// GitListRefUpdatesBy calls the XRPC method "sh.tangled.git.listRefUpdatesBy". -// -// cursor: Pagination cursor -// offset: Absolute offset for random-access pagination. Mutually exclusive with cursor; offsets drift under concurrent writes, so follow up with the returned cursor. -// order: Sort direction by createdAt. -// subject: Actor DID whose ref-update authorings to list. -func GitListRefUpdatesBy(ctx context.Context, c util.LexClient, cursor string, limit int64, offset int64, order string, subject string) (*GitListRefUpdatesBy_Output, error) { - var out GitListRefUpdatesBy_Output - - params := map[string]interface{}{} - if cursor != "" { - params["cursor"] = cursor - } - if limit != 0 { - params["limit"] = limit - } - if offset != 0 { - params["offset"] = offset - } - if order != "" { - params["order"] = order - } - params["subject"] = subject - if err := c.LexDo(ctx, util.Query, "", "sh.tangled.git.listRefUpdatesBy", params, nil, &out); err != nil { - return nil, err - } - - return &out, nil -} diff --git a/api/tangled/gitref.go b/api/tangled/gitref.go new file mode 100644 index 000000000..803c88849 --- /dev/null +++ b/api/tangled/gitref.go @@ -0,0 +1,23 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package tangled + +// schema: sh.tangled.git.ref + +import ( + "github.com/bluesky-social/indigo/lex/util" +) + +const ( + GitRefNSID = "sh.tangled.git.ref" +) + +func init() { + util.RegisterType("sh.tangled.git.ref", &GitRef{}) +} // +// RECORDTYPE: GitRef +type GitRef struct { + LexiconTypeID string `json:"$type,const=sh.tangled.git.ref" cborgen:"$type,const=sh.tangled.git.ref"` + // sha: Object the ref points at + Sha string `json:"sha" cborgen:"sha"` +} diff --git a/api/tangled/knotsubscribeRepos.go b/api/tangled/knotsubscribeRepos.go deleted file mode 100644 index 0f4d22a93..000000000 --- a/api/tangled/knotsubscribeRepos.go +++ /dev/null @@ -1,34 +0,0 @@ -// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. - -package tangled - -// schema: sh.tangled.knot.subscribeRepos - -const ( - KnotSubscribeReposNSID = "sh.tangled.knot.subscribeRepos" -) - -// KnotSubscribeRepos_GitSync1 is a "gitSync1" in the sh.tangled.knot.subscribeRepos schema. -type KnotSubscribeRepos_GitSync1 struct { - // did: Repository DID identifier - Did string `json:"did" cborgen:"did"` - // seq: The stream sequence number of this message. - Seq int64 `json:"seq" cborgen:"seq"` -} - -// KnotSubscribeRepos_GitSync2 is a "gitSync2" in the sh.tangled.knot.subscribeRepos schema. -type KnotSubscribeRepos_GitSync2 struct { - // repo: Repository AT-URI identifier - Repo string `json:"repo" cborgen:"repo"` - // seq: The stream sequence number of this message. - Seq int64 `json:"seq" cborgen:"seq"` -} - -// KnotSubscribeRepos_Identity is a "identity" in the sh.tangled.knot.subscribeRepos schema. -type KnotSubscribeRepos_Identity struct { - // did: Repository DID identifier - Did string `json:"did" cborgen:"did"` - // seq: The stream sequence number of this message. - Seq int64 `json:"seq" cborgen:"seq"` - Time string `json:"time" cborgen:"time"` -} diff --git a/cmd/cborgen/cborgen.go b/cmd/cborgen/cborgen.go index 7c501a48a..fdaaec81a 100644 --- a/cmd/cborgen/cborgen.go +++ b/cmd/cborgen/cborgen.go @@ -31,6 +31,7 @@ func main() { tangled.FeedStar_Repo{}, tangled.FeedStar_String{}, tangled.GitOid{}, + tangled.GitRef{}, tangled.GitRefUpdate{}, tangled.GitRefUpdate_CommitCountBreakdown{}, tangled.GitRefUpdate_IndividualEmailCommitCount{}, diff --git a/knotfeed/frame.go b/knotfeed/frame.go new file mode 100644 index 000000000..86643be32 --- /dev/null +++ b/knotfeed/frame.go @@ -0,0 +1,260 @@ +package knotfeed + +import ( + "bytes" + "context" + "errors" + "fmt" + "io" + "log/slog" + "strings" + + comatproto "github.com/bluesky-social/indigo/api/atproto" + atprotorepo "github.com/bluesky-social/indigo/atproto/repo" + "github.com/bluesky-social/indigo/atproto/syntax" + cbg "github.com/whyrusleeping/cbor-gen" +) + +const ( + TypeCommit = "#commit" + TypeSync = "#sync" + TypeIdentity = "#identity" + TypeAccount = "#account" + TypeInfo = "#info" + TypeError = "#error" +) + +type Message struct { + Type string + Commit *Commit + InfoName string + Error string + Detail string +} + +type RecordOp struct { + Action string + Collection string + Rkey string + Bytes []byte +} + +func (op RecordOp) Deleted() bool { + return op.Action == "delete" +} + +type Commit struct { + Repo string + Seq int64 + Rev string + Records []RecordOp +} + +func Decode(data []byte, log *slog.Logger) (Message, error) { + r := bytes.NewReader(data) + head, err := readHeader(r) + if err != nil { + return Message{}, fmt.Errorf("frame header: %w", err) + } + if head.isError { + return decodeError(r, head) + } + if head.msgType != "" { + return decodeTyped(head.msgType, r, log) + } + return Message{}, fmt.Errorf("frame header names neither a type or an error") +} + +func decodeError(r io.Reader, head frameHeader) (Message, error) { + msg := Message{Type: TypeError, Error: head.errName, Detail: head.errDetail} + maj, count, err := cbg.CborReadHeader(r) + if errors.Is(err, io.EOF) { + return msg, nil + } + if err != nil { + return Message{}, fmt.Errorf("error frame body: %w", err) + } + if maj != cbg.MajMap { + return Message{}, fmt.Errorf("error frame body is major type %d, not a map", maj) + } + for range count { + key, err := cbg.ReadString(r) + if err != nil { + return Message{}, fmt.Errorf("error frame body: %w", err) + } + switch key { + case "error": + if msg.Error, err = cbg.ReadString(r); err != nil { + return Message{}, fmt.Errorf("error frame body: %w", err) + } + case "message": + if msg.Detail, err = cbg.ReadString(r); err != nil { + return Message{}, fmt.Errorf("error frame body: %w", err) + } + default: + if err := skipValue(r, 0); err != nil { + return Message{}, err + } + } + } + return msg, nil +} + +type frameHeader struct { + msgType string + isError bool + errName string + errDetail string +} + +func readHeader(r io.Reader) (frameHeader, error) { + maj, count, err := cbg.CborReadHeader(r) + if err != nil { + return frameHeader{}, err + } + if maj != cbg.MajMap { + return frameHeader{}, fmt.Errorf("expected a cbor map header, got major type %d", maj) + } + head := frameHeader{} + for range count { + key, err := cbg.ReadString(r) + if err != nil { + return frameHeader{}, err + } + switch key { + case "t": + if head.msgType, err = cbg.ReadString(r); err != nil { + return frameHeader{}, err + } + case "op": + maj, val, err := cbg.CborReadHeader(r) + if err != nil { + return frameHeader{}, err + } + head.isError = maj == cbg.MajNegativeInt && val == 0 + case "error": + if head.errName, err = cbg.ReadString(r); err != nil { + return frameHeader{}, err + } + case "message": + if head.errDetail, err = cbg.ReadString(r); err != nil { + return frameHeader{}, err + } + default: + if err := skipValue(r, 0); err != nil { + return frameHeader{}, err + } + } + } + return head, nil +} + +func decodeTyped(msgType string, r io.Reader, log *slog.Logger) (Message, error) { + switch msgType { + case TypeCommit: + var evt comatproto.SyncSubscribeRepos_Commit + if err := evt.UnmarshalCBOR(r); err != nil { + return Message{}, fmt.Errorf("commit frame: %w", err) + } + commit, err := resolveCommit(&evt, log) + if err != nil { + return Message{}, err + } + return Message{Type: TypeCommit, Commit: commit}, nil + case TypeSync, TypeIdentity, TypeAccount: + return Message{Type: msgType}, nil + case TypeInfo: + var evt comatproto.SyncSubscribeRepos_Info + if err := evt.UnmarshalCBOR(r); err != nil { + return Message{}, fmt.Errorf("info frame: %w", err) + } + return Message{Type: TypeInfo, InfoName: evt.Name}, nil + default: + if err := skipValue(r, 0); err != nil { + return Message{}, err + } + return Message{Type: msgType}, nil + } +} + +func resolveCommit(evt *comatproto.SyncSubscribeRepos_Commit, log *slog.Logger) (*Commit, error) { + commit := &Commit{ + Repo: evt.Repo, + Seq: evt.Seq, + Rev: evt.Rev, + } + if len(evt.Ops) == 0 { + return commit, nil + } + ctx := context.Background() + _, repo, err := atprotorepo.LoadRepoFromCAR(ctx, bytes.NewReader(evt.Blocks)) + if err != nil { + log.Error("commit frame blocks didn't decode as a car, dropping the commit", "repo", evt.Repo, "seq", evt.Seq, "err", err) + return commit, nil + } + records := make([]RecordOp, 0, len(evt.Ops)) + for _, op := range evt.Ops { + if op == nil { + continue + } + collection, rkey, found := strings.Cut(op.Path, "/") + if !found { + log.Warn("commit op path has no rkey, dropping the op", "repo", evt.Repo, "path", op.Path) + continue + } + record := RecordOp{ + Action: op.Action, + Collection: collection, + Rkey: rkey, + } + if !record.Deleted() && repo != nil { + payload, _, err := repo.GetRecordBytes(ctx, syntax.NSID(collection), syntax.RecordKey(rkey)) + if err != nil { + log.Error("record bytes missing from the frame car, dropping the op", "repo", evt.Repo, "path", op.Path, "err", err) + continue + } + record.Bytes = payload + } + records = append(records, record) + } + commit.Records = records + return commit, nil +} + +const maxCborDepth = 64 + +func skipValue(r io.Reader, depth int) error { + if depth > maxCborDepth { + return fmt.Errorf("cbor value nests past %d levels", maxCborDepth) + } + maj, val, err := cbg.CborReadHeader(r) + if err != nil { + return err + } + switch maj { + case cbg.MajByteString, cbg.MajTextString: + _, err = io.CopyN(io.Discard, r, int64(val)) + return err + case cbg.MajArray: + for range val { + if err := skipValue(r, depth+1); err != nil { + return err + } + } + return nil + case cbg.MajMap: + for range val { + if err := skipValue(r, depth+1); err != nil { + return err + } + if err := skipValue(r, depth+1); err != nil { + return err + } + } + return nil + case cbg.MajTag: + return skipValue(r, depth+1) + default: + return nil + } +} diff --git a/knotfeed/frame_test.go b/knotfeed/frame_test.go new file mode 100644 index 000000000..3702ce2e9 --- /dev/null +++ b/knotfeed/frame_test.go @@ -0,0 +1,236 @@ +package 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 + "/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") + } +} diff --git a/knotfeed/record.go b/knotfeed/record.go new file mode 100644 index 000000000..69a6c3fd1 --- /dev/null +++ b/knotfeed/record.go @@ -0,0 +1,170 @@ +package knotfeed + +import ( + "bytes" + "fmt" + "strings" + "unicode/utf8" + + cbg "github.com/whyrusleeping/cbor-gen" +) + +const maxPushOptions = 64 +const GitRefCollection = "sh.tangled.git.ref" + +const mstKeyBudget = 256 + +const MaxRefRkeyBytes = mstKeyBudget - len(GitRefCollection) - 1 + +const hexDigits = "0123456789abcdef" + +func plainByte(b byte) bool { + switch { + case b >= 'A' && b <= 'Z', b >= 'a' && b <= 'z', b >= '0' && b <= '9': + return true + case b == '-', b == '_': + return true + default: + return false + } +} + +func allPlain(s string) bool { + for i := range len(s) { + if !plainByte(s[i]) { + return false + } + } + return true +} + +func EscapeRefname(refname string) (string, bool) { + out := make([]byte, 0, len(refname)) + for i := range len(refname) { + b := refname[i] + if plainByte(b) { + out = append(out, b) + continue + } + out = append(out, '~', hexDigits[b>>4], hexDigits[b&0x0f]) + } + rkey := string(out) + if len(rkey) > MaxRefRkeyBytes { + return "", false + } + return rkey, true +} + +func UnescapeRkey(rkey string) (string, bool) { + if !strings.Contains(rkey, "~") { + return rkey, allPlain(rkey) + } + var out []byte + rest := rkey + for { + idx := strings.IndexByte(rest, '~') + if idx < 0 { + if !allPlain(rest) { + return "", false + } + if !utf8.Valid(out) { + return "", false + } + return string(append(out, rest...)), true + } + head, chunk := rest[:idx], rest[idx+1:] + if !allPlain(head) || len(chunk) < 2 { + return "", false + } + hi, ok := nibble(chunk[0]) + if !ok { + return "", false + } + lo, ok := nibble(chunk[1]) + if !ok { + return "", false + } + out = append(out, head...) + out = append(out, hi<<4|lo) + rest = chunk[2:] + } +} + +func nibble(b byte) (byte, bool) { + switch { + case b >= '0' && b <= '9': + return b - '0', true + case b >= 'a' && b <= 'f': + return b - 'a' + 10, true + default: + return 0, false + } +} + +type RefRecord struct { + Sha string + Editor string + PushOptions []string +} + +func DecodeRefRecord(data []byte) (RefRecord, error) { + var rec RefRecord + r := bytes.NewReader(data) + maj, count, err := cbg.CborReadHeader(r) + if err != nil { + return rec, err + } + if maj != cbg.MajMap { + return rec, fmt.Errorf("expected a cbor map, got major type %d", maj) + } + for range count { + key, err := cbg.ReadString(r) + if err != nil { + return rec, err + } + switch key { + case "sha": + if rec.Sha, err = cbg.ReadString(r); err != nil { + return rec, err + } + case "x-tngl-editor": + if rec.Editor, err = cbg.ReadString(r); err != nil { + return rec, err + } + case "x-tngl-push-options": + if rec.PushOptions, err = readStringArray(r); err != nil { + return rec, err + } + default: + if err := skipValue(r, 0); err != nil { + return rec, err + } + } + } + if rec.Sha == "" { + return rec, fmt.Errorf("ref record is missing its sha") + } + return rec, nil +} + +func readStringArray(r *bytes.Reader) ([]string, error) { + maj, count, err := cbg.CborReadHeader(r) + if err != nil { + return nil, err + } + if maj != cbg.MajArray { + return nil, fmt.Errorf("expected a cbor array, got major type %d", maj) + } + if count > maxPushOptions { + return nil, fmt.Errorf("push options array claims %d entries", count) + } + out := make([]string, 0, count) + for range count { + s, err := cbg.ReadString(r) + if err != nil { + return nil, err + } + out = append(out, s) + } + return out, nil +} diff --git a/knotfeed/record_test.go b/knotfeed/record_test.go new file mode 100644 index 000000000..2ece8b69b --- /dev/null +++ b/knotfeed/record_test.go @@ -0,0 +1,124 @@ +package knotfeed + +import ( + "bytes" + "encoding/hex" + "slices" + "strings" + "testing" + + cbg "github.com/whyrusleeping/cbor-gen" +) + +const frozenPushedHex = "a4637368617828616261626162616261626162616261626162616261626162616261626162616261626162616261626524747970657273682e74616e676c65642e6769742e7265666d782d746e676c2d656469746f726e6469643a706c633a6c696d70657473782d746e676c2d707573682d6f7074696f6e738167736b69702d6369" + +const frozenBareHex = "a2637368617840636463646364636463646364636463646364636463646364636463646364636463646364636463646364636463646364636463646364636463646364636463646364636463646524747970657273682e74616e676c65642e6769742e726566" + +func TestDecodeRefRecordReadsTheFrozenRecordBytes(t *testing.T) { + for _, tt := range []struct { + name string + hex string + wantSha string + wantEditor string + wantOptions []string + }{ + {"pushed", frozenPushedHex, strings.Repeat("ab", 20), "did:plc:limpet", []string{"skip-ci"}}, + {"catch-up", frozenBareHex, strings.Repeat("cd", 32), "", nil}, + } { + t.Run(tt.name, func(t *testing.T) { + data, err := hex.DecodeString(tt.hex) + if err != nil { + t.Fatalf("frozen hex: %v", err) + } + rec, err := DecodeRefRecord(data) + if err != nil { + t.Fatalf("DecodeRefRecord: %v", err) + } + if rec.Sha != tt.wantSha { + t.Fatalf("Sha = %q, want %q", rec.Sha, tt.wantSha) + } + if rec.Editor != tt.wantEditor { + t.Fatalf("Editor = %q, want %q", rec.Editor, tt.wantEditor) + } + if !slices.Equal(rec.PushOptions, tt.wantOptions) { + t.Fatalf("PushOptions = %v, want %v", rec.PushOptions, tt.wantOptions) + } + }) + } +} + +func TestDecodeRefRecordRefusesMalformedRecords(t *testing.T) { + var out bytes.Buffer + if err := cbg.CborWriteHeader(&out, cbg.MajMap, 1); err != nil { + t.Fatalf("map header: %v", err) + } + writeText(&out, "x-tngl-push-options") + if err := cbg.CborWriteHeader(&out, cbg.MajArray, 1<<32); err != nil { + t.Fatalf("array header: %v", err) + } + for _, data := range [][]byte{ + {0xa1, 0x63, 'a', 'b', 0x00}, + out.Bytes(), + } { + if _, err := DecodeRefRecord(data); err == nil { + t.Fatalf("DecodeRefRecord accepted % x", data) + } + } +} + +func TestEscapeRefnameRulesTheRkey(t *testing.T) { + for _, tt := range []struct { + refname string + want string + wantOK bool + }{ + {"refs/heads/main", "refs~2fheads~2fmain", true}, + {"refs/heads/~x", "refs~2fheads~2f~7ex", true}, + {strings.Repeat("/", 79), strings.Repeat("~2f", 79), true}, + {strings.Repeat("/", 80), "", false}, + } { + rkey, ok := EscapeRefname(tt.refname) + if ok != tt.wantOK || (ok && rkey != tt.want) { + t.Fatalf("EscapeRefname(%q) = %q, %v", tt.refname, rkey, ok) + } + } +} + +func TestUnescapeRkeyRoundTripsPlainAndEscapedRefnames(t *testing.T) { + for _, refname := range []string{ + "refs/heads/main", + "main", + "refs/tags/v1.0.0+build", + "refs/heads/feature/foo_bar-squid.conch", + "refs/heads/~tilde", + "refs/heads/ branch with spaces ", + "refs/heads/日本語ブランチ", + "refs/heads/emoji-🍀", + } { + rkey, ok := EscapeRefname(refname) + if !ok { + t.Fatalf("EscapeRefname refused %q", refname) + } + decoded, ok := UnescapeRkey(rkey) + if !ok { + t.Fatalf("UnescapeRkey refused %q", rkey) + } + if decoded != refname { + t.Fatalf("round trip of %q answered %q", refname, decoded) + } + } +} + +func TestUnescapeRkeyRefusesMalformedKeys(t *testing.T) { + for _, rkey := range []string{ + "refs~2fheads~", + "refs~2fheads~2", + "refs~2fheads~zz", + "refs~2fhea ds", + "refs~ff", + } { + if decoded, ok := UnescapeRkey(rkey); ok { + t.Fatalf("UnescapeRkey(%q) = %q, want refusal", rkey, decoded) + } + } +} diff --git a/knotfeed/subscribe.go b/knotfeed/subscribe.go new file mode 100644 index 000000000..7ce514d90 --- /dev/null +++ b/knotfeed/subscribe.go @@ -0,0 +1,227 @@ +package knotfeed + +import ( + "cmp" + "context" + "errors" + "fmt" + "log/slog" + "net/url" + "time" + + "github.com/gorilla/websocket" +) + +const ( + livenessTimeout = 90 * time.Second + pongWriteTimeout = 5 * time.Second + initialBackoff = time.Second + maxBackoff = 60 * time.Second + maxFrameBytes = 8 << 20 + subscribeReposNS = "com.atproto.sync.subscribeRepos" + + defaultMaxHandlerAttempts = 3 + defaultPoisonGrace = 10 * time.Minute +) + +type Consumer struct { + Host string + NoTLS bool + Dialer *websocket.Dialer + Logger *slog.Logger + + LoadCursor func(context.Context) (int64, error) + StoreCursor func(context.Context, int64) error + Handle func(context.Context, Message) error + OutdatedReplay func(context.Context) int64 + OnConnectError func(error) + + ReplayFromStart bool + + MaxHandlerAttempts int + PoisonGrace time.Duration + + now func() time.Time + + poisonSeq int64 + poisonAttempts int + poisonSince time.Time +} + +func (c *Consumer) Run(ctx context.Context) error { + c.fill() + logger := c.Logger + backoff := initialBackoff + for { + if ctx.Err() != nil { + return nil + } + advanced, resync, err := c.session(ctx) + switch { + case err != nil: + if c.OnConnectError != nil && isConnectError(err) { + c.OnConnectError(err) + } + logger.Error("firehose session failed", "host", c.Host, "err", err) + case resync: + advanced = true + } + if advanced { + backoff = initialBackoff + } else { + backoff = min(backoff*2, maxBackoff) + } + if resync { + continue + } + timer := time.NewTimer(backoff) + select { + case <-ctx.Done(): + timer.Stop() + return nil + case <-timer.C: + } + } +} + +type connectError struct{ err error } + +func (e connectError) Error() string { return e.err.Error() } + +func (e connectError) Unwrap() error { return e.err } + +func isConnectError(err error) bool { + var ce connectError + return errors.As(err, &ce) +} + +func (c *Consumer) session(ctx context.Context) (advanced bool, resync bool, err error) { + cursor, err := c.LoadCursor(ctx) + if err != nil { + return false, false, fmt.Errorf("loading cursor: %w", err) + } + conn, _, err := c.Dialer.DialContext(ctx, c.url(cursor), nil) + if err != nil { + return false, false, connectError{err} + } + defer conn.Close() + watcher := context.AfterFunc(ctx, func() { conn.Close() }) + defer watcher() + conn.SetReadLimit(maxFrameBytes) + c.Logger.Info("subscribed to the firehose", "host", c.Host, "cursor", cursor) + + conn.SetPingHandler(func(payload string) error { + conn.SetReadDeadline(time.Now().Add(livenessTimeout)) + return conn.WriteControl(websocket.PongMessage, []byte(payload), time.Now().Add(pongWriteTimeout)) + }) + conn.SetReadDeadline(time.Now().Add(livenessTimeout)) + for { + _, data, err := conn.ReadMessage() + if err != nil { + return advanced, false, err + } + conn.SetReadDeadline(time.Now().Add(livenessTimeout)) + message, err := Decode(data, c.Logger) + if err != nil { + c.Logger.Error("undecodable firehose frame", "host", c.Host, "err", err) + continue + } + resync, err := c.dispatch(ctx, message, &cursor) + if resync { + return advanced, true, nil + } + if err != nil { + return false, false, err + } + if message.Type == TypeCommit { + advanced = true + } + } +} + +func (c *Consumer) dispatch(ctx context.Context, message Message, cursor *int64) (bool, error) { + if message.Type == TypeCommit { + return c.commit(ctx, message, cursor) + } + if c.unreachable(message) { + c.Logger.Warn("firehose can't serve our cursor, resuming live", "host", c.Host, "error", c.resumeReason(message)) + replay := int64(0) + if c.OutdatedReplay != nil { + replay = c.OutdatedReplay(ctx) + } + *cursor = replay + if err := c.StoreCursor(ctx, replay); err != nil { + return false, fmt.Errorf("storing cursor: %w", err) + } + return true, nil + } + if c.Handle != nil { + if err := c.Handle(ctx, message); err != nil { + c.Logger.Error("firehose handler failed", "host", c.Host, "frame", message.Type, "err", err) + } + } + return false, nil +} + +func (c *Consumer) commit(ctx context.Context, message Message, cursor *int64) (bool, error) { + if c.Handle != nil { + if err := c.Handle(ctx, message); err != nil { + if !c.poison(message.Commit.Seq) { + return false, fmt.Errorf("frame seq %d: %w", message.Commit.Seq, err) + } + c.Logger.Error("handler keeps failing at one frame, skipping it", + "host", c.Host, "seq", message.Commit.Seq, "attempts", c.poisonAttempts, "err", err) + } + } + *cursor = message.Commit.Seq + if err := c.StoreCursor(ctx, *cursor); err != nil { + return false, fmt.Errorf("storing cursor: %w", err) + } + return false, nil +} + +func (c *Consumer) poison(seq int64) bool { + if seq != c.poisonSeq { + c.poisonSeq = seq + c.poisonAttempts = 0 + c.poisonSince = c.now() + } + c.poisonAttempts++ + return c.poisonAttempts >= c.MaxHandlerAttempts && c.now().Sub(c.poisonSince) >= c.PoisonGrace +} + +func (c *Consumer) unreachable(message Message) bool { + if message.Type == TypeError { + return message.Error == "FutureCursor" + } + return message.Type == TypeInfo && message.InfoName == "OutdatedCursor" +} + +func (c *Consumer) resumeReason(message Message) string { + if message.Type == TypeError { + return message.Error + } + return message.InfoName +} + +func (c *Consumer) url(cursor int64) string { + scheme := "wss" + if c.NoTLS { + scheme = "ws" + } + endpoint := url.URL{Scheme: scheme, Host: c.Host, Path: "/xrpc/" + subscribeReposNS} + if cursor > 0 || (cursor == 0 && c.ReplayFromStart) { + endpoint.RawQuery = "cursor=" + fmt.Sprint(cursor) + } + return endpoint.String() +} + +func (c *Consumer) fill() { + c.MaxHandlerAttempts = cmp.Or(c.MaxHandlerAttempts, defaultMaxHandlerAttempts) + c.PoisonGrace = cmp.Or(c.PoisonGrace, defaultPoisonGrace) + c.Dialer = cmp.Or(c.Dialer, websocket.DefaultDialer) + c.Logger = cmp.Or(c.Logger, slog.Default()) + if c.now == nil { + c.now = time.Now + } +} diff --git a/knotfeed/subscribe_test.go b/knotfeed/subscribe_test.go new file mode 100644 index 000000000..bbdfe67e8 --- /dev/null +++ b/knotfeed/subscribe_test.go @@ -0,0 +1,116 @@ +package knotfeed + +import ( + "bytes" + "context" + "errors" + "io" + "log/slog" + "net/http" + "net/http/httptest" + "slices" + "sync/atomic" + "testing" + "time" + + comatproto "github.com/bluesky-social/indigo/api/atproto" + lexutil "github.com/bluesky-social/indigo/lex/util" + "github.com/gorilla/websocket" +) + +func poisonTestLogger() *slog.Logger { + return slog.New(slog.NewTextHandler(io.Discard, nil)) +} + +func commitFrameForSeq(t *testing.T, seq int64) []byte { + t.Helper() + var payload bytes.Buffer + evt := comatproto.SyncSubscribeRepos_Commit{ + Repo: "did:plc:scallop", + Seq: seq, + Rev: "3lb2xkw2qrs2j", + Commit: lexutil.LexLink(testCid(t)), + } + if err := evt.MarshalCBOR(&payload); err != nil { + t.Fatalf("MarshalCBOR: %v", err) + } + return append(headerFrame(t, "t", "#commit"), payload.Bytes()...) +} + +func frameOnceServer(t *testing.T, frame []byte) *httptest.Server { + t.Helper() + upgrader := websocket.Upgrader{} + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + conn, err := upgrader.Upgrade(w, r, nil) + if err != nil { + return + } + if err := conn.WriteMessage(websocket.BinaryMessage, frame); err != nil { + return + } + _ = conn.WriteControl(websocket.CloseMessage, + websocket.FormatCloseMessage(websocket.CloseNormalClosure, ""), time.Now().Add(time.Second)) + conn.Close() + })) + t.Cleanup(srv.Close) + return srv +} + +func poisonConsumer(srv *httptest.Server, cursor *atomic.Int64, stored *atomic.Int64, now *atomic.Int64, grace time.Duration) *Consumer { + return &Consumer{ + Host: srv.Listener.Addr().String(), + NoTLS: true, + Logger: poisonTestLogger(), + + LoadCursor: func(context.Context) (int64, error) { + return cursor.Load(), nil + }, + StoreCursor: func(_ context.Context, seq int64) error { + cursor.Store(seq) + stored.Store(seq) + return nil + }, + Handle: func(context.Context, Message) error { + return errors.New("boom") + }, + + MaxHandlerAttempts: 3, + PoisonGrace: grace, + now: func() time.Time { + return time.Unix(0, now.Load()) + }, + } +} + +func poisonSequence(t *testing.T, c *Consumer, now, stored *atomic.Int64) []int64 { + t.Helper() + var seen []int64 + for attempt := range 3 { + now.Store(int64(attempt+1) * 6 * int64(time.Minute)) + _, _, _ = c.session(context.Background()) + seen = append(seen, stored.Load()) + } + return seen +} + +func TestPoisonFrameSkipsAfterThreeFailuresAcrossTheGrace(t *testing.T) { + srv := frameOnceServer(t, commitFrameForSeq(t, 100)) + + var cursor, stored, now atomic.Int64 + c := poisonConsumer(srv, &cursor, &stored, &now, defaultPoisonGrace) + + if seen := poisonSequence(t, c, &now, &stored); !slices.Equal(seen, []int64{0, 0, 100}) { + t.Fatalf("cursor stores per attempt = %v, want [0 0 100] once the poison lands", seen) + } +} + +func TestPoisonFrameHoldsWhileFailuresStayInsideTheGrace(t *testing.T) { + srv := frameOnceServer(t, commitFrameForSeq(t, 100)) + + var cursor, stored, now atomic.Int64 + c := poisonConsumer(srv, &cursor, &stored, &now, time.Hour) + + if seen := poisonSequence(t, c, &now, &stored); !slices.Equal(seen, []int64{0, 0, 0}) { + t.Fatalf("cursor stores per attempt = %v, want all zero while inside the grace", seen) + } +} diff --git a/lexicons/git/countRefUpdates.json b/lexicons/git/countRefUpdates.json deleted file mode 100644 index 3e457bcce..000000000 --- a/lexicons/git/countRefUpdates.json +++ /dev/null @@ -1,39 +0,0 @@ -{ - "lexicon": 1, - "id": "sh.tangled.git.countRefUpdates", - "defs": { - "main": { - "type": "query", - "parameters": { - "type": "params", - "required": ["subject"], - "properties": { - "subject": { - "type": "string", - "format": "did", - "description": "Repo DID whose ref-update records to list." - } - } - }, - "output": { - "encoding": "application/json", - "schema": { - "type": "object", - "required": ["count", "distinctAuthors"], - "properties": { - "count": { - "type": "integer", - "minimum": 0, - "description": "Total number of matching records." - }, - "distinctAuthors": { - "type": "integer", - "minimum": 0, - "description": "Number of distinct authors among the matching records." - } - } - } - } - } - } -} diff --git a/lexicons/git/countRefUpdatesBy.json b/lexicons/git/countRefUpdatesBy.json deleted file mode 100644 index 06f8208a7..000000000 --- a/lexicons/git/countRefUpdatesBy.json +++ /dev/null @@ -1,39 +0,0 @@ -{ - "lexicon": 1, - "id": "sh.tangled.git.countRefUpdatesBy", - "defs": { - "main": { - "type": "query", - "parameters": { - "type": "params", - "required": ["subject"], - "properties": { - "subject": { - "type": "string", - "format": "did", - "description": "Actor DID whose ref-update authorings to list." - } - } - }, - "output": { - "encoding": "application/json", - "schema": { - "type": "object", - "required": ["count", "distinctAuthors"], - "properties": { - "count": { - "type": "integer", - "minimum": 0, - "description": "Total number of matching records." - }, - "distinctAuthors": { - "type": "integer", - "minimum": 0, - "description": "Number of distinct authors among the matching records." - } - } - } - } - } - } -} diff --git a/lexicons/git/listRefUpdates.json b/lexicons/git/listRefUpdates.json deleted file mode 100644 index c5e1337f0..000000000 --- a/lexicons/git/listRefUpdates.json +++ /dev/null @@ -1,71 +0,0 @@ -{ - "lexicon": 1, - "id": "sh.tangled.git.listRefUpdates", - "defs": { - "main": { - "type": "query", - "parameters": { - "type": "params", - "required": ["subject"], - "properties": { - "subject": { - "type": "string", - "format": "did", - "description": "Repo DID whose ref-update records to list." - }, - "cursor": { - "type": "string", - "description": "Pagination cursor" - }, - "offset": { - "type": "integer", - "minimum": 0, - "description": "Absolute offset for random-access pagination. Mutually exclusive with cursor; offsets drift under concurrent writes, so follow up with the returned cursor." - }, - "limit": { - "type": "integer", - "minimum": 1, - "maximum": 1000, - "default": 50 - }, - "order": { - "type": "string", - "knownValues": ["asc", "desc"], - "default": "desc", - "description": "Sort direction by createdAt." - } - } - }, - "output": { - "encoding": "application/json", - "schema": { - "type": "object", - "required": ["items"], - "properties": { - "items": { - "type": "array", - "items": { "type": "ref", "ref": "#listItem" } - }, - "cursor": { "type": "string" }, - "total": { - "type": "integer", - "description": "Total items in the full list; omitted for filtered or merged views" - } - } - } - } - }, - "listItem": { - "type": "object", - "required": ["uri", "value"], - "properties": { - "uri": { "type": "string", "format": "at-uri" }, - "cid": { "type": "string", "format": "cid" }, - "value": { - "type": "unknown", - "description": "Embedded sh.tangled.git.refUpdate record" - } - } - } - } -} diff --git a/lexicons/git/listRefUpdatesBy.json b/lexicons/git/listRefUpdatesBy.json deleted file mode 100644 index 8f6dedb2e..000000000 --- a/lexicons/git/listRefUpdatesBy.json +++ /dev/null @@ -1,59 +0,0 @@ -{ - "lexicon": 1, - "id": "sh.tangled.git.listRefUpdatesBy", - "defs": { - "main": { - "type": "query", - "parameters": { - "type": "params", - "required": ["subject"], - "properties": { - "subject": { - "type": "string", - "format": "did", - "description": "Actor DID whose ref-update authorings to list." - }, - "cursor": { - "type": "string", - "description": "Pagination cursor" - }, - "offset": { - "type": "integer", - "minimum": 0, - "description": "Absolute offset for random-access pagination. Mutually exclusive with cursor; offsets drift under concurrent writes, so follow up with the returned cursor." - }, - "limit": { - "type": "integer", - "minimum": 1, - "maximum": 1000, - "default": 50 - }, - "order": { - "type": "string", - "knownValues": ["asc", "desc"], - "default": "desc", - "description": "Sort direction by createdAt." - } - } - }, - "output": { - "encoding": "application/json", - "schema": { - "type": "object", - "required": ["items"], - "properties": { - "items": { - "type": "array", - "items": { "type": "ref", "ref": "sh.tangled.git.listRefUpdates#listItem" } - }, - "cursor": { "type": "string" }, - "total": { - "type": "integer", - "description": "Total items in the full list; omitted for filtered or merged views" - } - } - } - } - } - } -} diff --git a/lexicons/git/ref.json b/lexicons/git/ref.json new file mode 100644 index 000000000..62c69a10d --- /dev/null +++ b/lexicons/git/ref.json @@ -0,0 +1,28 @@ +{ + "lexicon": 1, + "id": "sh.tangled.git.ref", + "needsCbor": true, + "needsType": true, + "defs": { + "main": { + "type": "record", + "description": "Present state of one git ref on this repository, w/ the escaped refname as rkey (every byte outside [A-Za-z0-9_-] as ~xx in hex, ~ itself as ~7e; refs/heads/main is refs~2fheads~2fmain)", + "key": "any", + "record": { + "type": "object", + "required": [ + "sha" + ], + "properties": { + "sha": { + "type": "string", + "minLength": 40, + "maxLength": 64, + "pattern": "^[0-9a-f]+$", + "description": "Object the ref points at" + } + } + } + } + } +} diff --git a/lexicons/knot/subscribeRepos.json b/lexicons/knot/subscribeRepos.json deleted file mode 100644 index b06dca1b1..000000000 --- a/lexicons/knot/subscribeRepos.json +++ /dev/null @@ -1,57 +0,0 @@ -{ - "lexicon": 1, - "id": "sh.tangled.knot.subscribeRepos", - "defs": { - "main": { - "type": "subscription", - "description": "Repository event stream, aka Firehose endpoint. Outputs repo commits with diff data, and identity update events, for all repositories on the current server. See the atproto specifications for details around stream sequencing, repo versioning, CAR diff format, and more. Public and does not require auth; implemented by PDS and Relay.", - "parameters": { - "type": "params", - "properties": { - "cursor": { - "type": "integer", - "description": "The last known event seq number to backfill from." - } - } - }, - "message": { - "schema": { - "type": "union", - "refs": ["#identity", "sh.tangled.git.refUpdate"] - } - }, - "errors": [ - { "name": "FutureCursor" }, - { - "name": "ConsumerTooSlow", - "description": "If the consumer of the stream can not keep up with events, and a backlog gets too large, the server will drop the connection." - } - ] - }, - "identity": { - "type": "object", - "required": ["seq", "did", "time"], - "properties": { - "seq": { "type": "integer", "description": "The stream sequence number of this message." }, - "did": { "type": "string", "format": "did", "description": "Repository DID identifier" }, - "time": { "type": "string", "format": "datetime" } - } - }, - "gitSync1": { - "type": "object", - "required": ["seq", "did"], - "properties": { - "seq": { "type": "integer", "description": "The stream sequence number of this message." }, - "did": { "type": "string", "format": "did", "description": "Repository DID identifier" } - } - }, - "gitSync2": { - "type": "object", - "required": ["seq", "repo"], - "properties": { - "seq": { "type": "integer", "description": "The stream sequence number of this message." }, - "repo": { "type": "string", "format": "at-uri", "description": "Repository AT-URI identifier" } - } - } - } -} diff --git a/lexicons/temp/notification/deleteNotification.json b/lexicons/temp/notification/deleteNotification.json index d41fce867..89376f02b 100644 --- a/lexicons/temp/notification/deleteNotification.json +++ b/lexicons/temp/notification/deleteNotification.json @@ -13,8 +13,7 @@ "properties": { "uri": { "type": "string", - "format": "at-uri", - "description": "at-uri of the notification to delete." + "description": "notification key as returned by listNotifications; pass it back unchanged." } } } diff --git a/lexicons/temp/notification/listNotifications.json b/lexicons/temp/notification/listNotifications.json index 06d319c98..d6fcbee61 100644 --- a/lexicons/temp/notification/listNotifications.json +++ b/lexicons/temp/notification/listNotifications.json @@ -59,8 +59,7 @@ "properties": { "uri": { "type": "string", - "format": "at-uri", - "description": "at-uri of this notification; the stable key for read/unread state." + "description": "notification key; pass it back to updateSeen and deleteNotification." }, "type": { "type": "string", diff --git a/lexicons/temp/notification/updateSeen.json b/lexicons/temp/notification/updateSeen.json index a72e51ce3..adc93b415 100644 --- a/lexicons/temp/notification/updateSeen.json +++ b/lexicons/temp/notification/updateSeen.json @@ -13,8 +13,7 @@ "properties": { "uri": { "type": "string", - "format": "at-uri", - "description": "at-uri of the notification to update." + "description": "notification key as returned by listNotifications; pass it back unchanged." }, "read": { "type": "boolean" diff --git a/web/src/lib/api/count.ts b/web/src/lib/api/count.ts index a862b924b..fb85c6b8b 100644 --- a/web/src/lib/api/count.ts +++ b/web/src/lib/api/count.ts @@ -13,8 +13,6 @@ export type CountName = | "sh.tangled.graph.countFollowsBy" | "sh.tangled.graph.countVouches" | "sh.tangled.graph.countVouchesBy" - | "sh.tangled.git.countRefUpdates" - | "sh.tangled.git.countRefUpdatesBy" | "sh.tangled.knot.countKnots" | "sh.tangled.knot.countMembers" | "sh.tangled.knot.countMembersBy" diff --git a/web/src/lib/api/lexicons/index.ts b/web/src/lib/api/lexicons/index.ts index 53bb25964..737ca2580 100644 --- a/web/src/lib/api/lexicons/index.ts +++ b/web/src/lib/api/lexicons/index.ts @@ -79,16 +79,13 @@ export * as ShTangledFeedListStars from "./types/sh/tangled/feed/listStars.js"; export * as ShTangledFeedListStarsBy from "./types/sh/tangled/feed/listStarsBy.js"; export * as ShTangledFeedReaction from "./types/sh/tangled/feed/reaction.js"; export * as ShTangledFeedStar from "./types/sh/tangled/feed/star.js"; -export * as ShTangledGitCountRefUpdates from "./types/sh/tangled/git/countRefUpdates.js"; -export * as ShTangledGitCountRefUpdatesBy from "./types/sh/tangled/git/countRefUpdatesBy.js"; export * as ShTangledGitDefs from "./types/sh/tangled/git/defs.js"; export * as ShTangledGitKeepCommit from "./types/sh/tangled/git/keepCommit.js"; -export * as ShTangledGitListRefUpdates from "./types/sh/tangled/git/listRefUpdates.js"; -export * as ShTangledGitListRefUpdatesBy from "./types/sh/tangled/git/listRefUpdatesBy.js"; export * as ShTangledGitListRefs from "./types/sh/tangled/git/listRefs.js"; export * as ShTangledGitMergeCheck from "./types/sh/tangled/git/mergeCheck.js"; export * as ShTangledGitMergeCommit from "./types/sh/tangled/git/mergeCommit.js"; export * as ShTangledGitOid from "./types/sh/tangled/git/oid.js"; +export * as ShTangledGitRef from "./types/sh/tangled/git/ref.js"; export * as ShTangledGitRefUpdate from "./types/sh/tangled/git/refUpdate.js"; export * as ShTangledGitTempAnalyzeMerge from "./types/sh/tangled/git/temp/analyzeMerge.js"; export * as ShTangledGitTempDefs from "./types/sh/tangled/git/temp/defs.js"; @@ -138,7 +135,6 @@ export * as ShTangledKnotMember from "./types/sh/tangled/knot/member.js"; export * as ShTangledKnotMemberAcceptance from "./types/sh/tangled/knot/memberAcceptance.js"; export * as ShTangledKnotMemberInvite from "./types/sh/tangled/knot/memberInvite.js"; export * as ShTangledKnotRemoveMember from "./types/sh/tangled/knot/removeMember.js"; -export * as ShTangledKnotSubscribeRepos from "./types/sh/tangled/knot/subscribeRepos.js"; export * as ShTangledKnotVersion from "./types/sh/tangled/knot/version.js"; export * as ShTangledLabelCountDefinitions from "./types/sh/tangled/label/countDefinitions.js"; export * as ShTangledLabelCountOps from "./types/sh/tangled/label/countOps.js"; diff --git a/web/src/lib/api/lexicons/types/org/tangled/temp/notification/deleteNotification.ts b/web/src/lib/api/lexicons/types/org/tangled/temp/notification/deleteNotification.ts index 28a11048a..e58d612c9 100644 --- a/web/src/lib/api/lexicons/types/org/tangled/temp/notification/deleteNotification.ts +++ b/web/src/lib/api/lexicons/types/org/tangled/temp/notification/deleteNotification.ts @@ -10,9 +10,9 @@ const _mainSchema = /*#__PURE__*/ v.procedure( type: "lex", schema: /*#__PURE__*/ v.object({ /** - * at-uri of the notification to delete. + * notification key as returned by listNotifications; pass it back unchanged. */ - uri: /*#__PURE__*/ v.resourceUriString(), + uri: /*#__PURE__*/ v.string(), }), }, output: null, diff --git a/web/src/lib/api/lexicons/types/org/tangled/temp/notification/listNotifications.ts b/web/src/lib/api/lexicons/types/org/tangled/temp/notification/listNotifications.ts index 5e745c510..c64dd1a38 100644 --- a/web/src/lib/api/lexicons/types/org/tangled/temp/notification/listNotifications.ts +++ b/web/src/lib/api/lexicons/types/org/tangled/temp/notification/listNotifications.ts @@ -72,9 +72,9 @@ const _notificationSchema = /*#__PURE__*/ v.object({ */ type: /*#__PURE__*/ v.string(), /** - * at-uri of this notification; the stable key for read/unread state. + * notification key; pass it back to updateSeen and deleteNotification. */ - uri: /*#__PURE__*/ v.resourceUriString(), + uri: /*#__PURE__*/ v.string(), }); type main$schematype = typeof _mainSchema; diff --git a/web/src/lib/api/lexicons/types/org/tangled/temp/notification/updateSeen.ts b/web/src/lib/api/lexicons/types/org/tangled/temp/notification/updateSeen.ts index 1d63a1171..1d39600a1 100644 --- a/web/src/lib/api/lexicons/types/org/tangled/temp/notification/updateSeen.ts +++ b/web/src/lib/api/lexicons/types/org/tangled/temp/notification/updateSeen.ts @@ -11,9 +11,9 @@ const _mainSchema = /*#__PURE__*/ v.procedure( schema: /*#__PURE__*/ v.object({ read: /*#__PURE__*/ v.boolean(), /** - * at-uri of the notification to update. + * notification key as returned by listNotifications; pass it back unchanged. */ - uri: /*#__PURE__*/ v.resourceUriString(), + uri: /*#__PURE__*/ v.string(), }), }, output: null, diff --git a/web/src/lib/api/lexicons/types/sh/tangled/git/countRefUpdates.ts b/web/src/lib/api/lexicons/types/sh/tangled/git/countRefUpdates.ts deleted file mode 100644 index 8306704d5..000000000 --- a/web/src/lib/api/lexicons/types/sh/tangled/git/countRefUpdates.ts +++ /dev/null @@ -1,42 +0,0 @@ -import type {} from "@atcute/lexicons"; -import * as v from "@atcute/lexicons/validations"; -import type {} from "@atcute/lexicons/ambient"; - -const _mainSchema = /*#__PURE__*/ v.query("sh.tangled.git.countRefUpdates", { - params: /*#__PURE__*/ v.object({ - /** - * Repo DID whose ref-update records to list. - */ - subject: /*#__PURE__*/ v.didString(), - }), - output: { - type: "lex", - schema: /*#__PURE__*/ v.object({ - /** - * Total number of matching records. - * @minimum 0 - */ - count: /*#__PURE__*/ v.integer(), - /** - * Number of distinct authors among the matching records. - * @minimum 0 - */ - distinctAuthors: /*#__PURE__*/ v.integer(), - }), - }, -}); - -type main$schematype = typeof _mainSchema; - -export interface mainSchema extends main$schematype {} - -export const mainSchema = _mainSchema as mainSchema; - -export interface $params extends v.InferInput {} -export interface $output extends v.InferXRPCBodyInput {} - -declare module "@atcute/lexicons/ambient" { - interface XRPCQueries { - "sh.tangled.git.countRefUpdates": mainSchema; - } -} diff --git a/web/src/lib/api/lexicons/types/sh/tangled/git/countRefUpdatesBy.ts b/web/src/lib/api/lexicons/types/sh/tangled/git/countRefUpdatesBy.ts deleted file mode 100644 index 77c3c5ddd..000000000 --- a/web/src/lib/api/lexicons/types/sh/tangled/git/countRefUpdatesBy.ts +++ /dev/null @@ -1,42 +0,0 @@ -import type {} from "@atcute/lexicons"; -import * as v from "@atcute/lexicons/validations"; -import type {} from "@atcute/lexicons/ambient"; - -const _mainSchema = /*#__PURE__*/ v.query("sh.tangled.git.countRefUpdatesBy", { - params: /*#__PURE__*/ v.object({ - /** - * Actor DID whose ref-update authorings to list. - */ - subject: /*#__PURE__*/ v.didString(), - }), - output: { - type: "lex", - schema: /*#__PURE__*/ v.object({ - /** - * Total number of matching records. - * @minimum 0 - */ - count: /*#__PURE__*/ v.integer(), - /** - * Number of distinct authors among the matching records. - * @minimum 0 - */ - distinctAuthors: /*#__PURE__*/ v.integer(), - }), - }, -}); - -type main$schematype = typeof _mainSchema; - -export interface mainSchema extends main$schematype {} - -export const mainSchema = _mainSchema as mainSchema; - -export interface $params extends v.InferInput {} -export interface $output extends v.InferXRPCBodyInput {} - -declare module "@atcute/lexicons/ambient" { - interface XRPCQueries { - "sh.tangled.git.countRefUpdatesBy": mainSchema; - } -} diff --git a/web/src/lib/api/lexicons/types/sh/tangled/git/listRefUpdates.ts b/web/src/lib/api/lexicons/types/sh/tangled/git/listRefUpdates.ts deleted file mode 100644 index dc6590a9c..000000000 --- a/web/src/lib/api/lexicons/types/sh/tangled/git/listRefUpdates.ts +++ /dev/null @@ -1,84 +0,0 @@ -import type {} from "@atcute/lexicons"; -import * as v from "@atcute/lexicons/validations"; -import type {} from "@atcute/lexicons/ambient"; - -const _listItemSchema = /*#__PURE__*/ v.object({ - $type: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.literal("sh.tangled.git.listRefUpdates#listItem"), - ), - cid: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.cidString()), - uri: /*#__PURE__*/ v.resourceUriString(), - /** - * Embedded sh.tangled.git.refUpdate record - */ - value: /*#__PURE__*/ v.unknown(), -}); -const _mainSchema = /*#__PURE__*/ v.query("sh.tangled.git.listRefUpdates", { - params: /*#__PURE__*/ v.object({ - /** - * Pagination cursor - */ - cursor: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.string()), - /** - * @minimum 1 - * @maximum 1000 - * @default 50 - */ - limit: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.constrain(/*#__PURE__*/ v.integer(), [ - /*#__PURE__*/ v.integerRange(1, 1000), - ]), - 50, - ), - /** - * Absolute offset for random-access pagination. Mutually exclusive with cursor; offsets drift under concurrent writes, so follow up with the returned cursor. - * @minimum 0 - */ - offset: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.integer()), - /** - * Sort direction by createdAt. - * @default "desc" - */ - order: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.string<"asc" | "desc" | (string & {})>(), - "desc", - ), - /** - * Repo DID whose ref-update records to list. - */ - subject: /*#__PURE__*/ v.didString(), - }), - output: { - type: "lex", - schema: /*#__PURE__*/ v.object({ - cursor: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.string()), - get items() { - return /*#__PURE__*/ v.array(listItemSchema); - }, - /** - * Total items in the full list; omitted for filtered or merged views - */ - total: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.integer()), - }), - }, -}); - -type listItem$schematype = typeof _listItemSchema; -type main$schematype = typeof _mainSchema; - -export interface listItemSchema extends listItem$schematype {} -export interface mainSchema extends main$schematype {} - -export const listItemSchema = _listItemSchema as listItemSchema; -export const mainSchema = _mainSchema as mainSchema; - -export interface ListItem extends v.InferInput {} - -export interface $params extends v.InferInput {} -export interface $output extends v.InferXRPCBodyInput {} - -declare module "@atcute/lexicons/ambient" { - interface XRPCQueries { - "sh.tangled.git.listRefUpdates": mainSchema; - } -} diff --git a/web/src/lib/api/lexicons/types/sh/tangled/git/listRefUpdatesBy.ts b/web/src/lib/api/lexicons/types/sh/tangled/git/listRefUpdatesBy.ts deleted file mode 100644 index c3d71ee9d..000000000 --- a/web/src/lib/api/lexicons/types/sh/tangled/git/listRefUpdatesBy.ts +++ /dev/null @@ -1,69 +0,0 @@ -import type {} from "@atcute/lexicons"; -import * as v from "@atcute/lexicons/validations"; -import type {} from "@atcute/lexicons/ambient"; -import * as ShTangledGitListRefUpdates from "./listRefUpdates.js"; - -const _mainSchema = /*#__PURE__*/ v.query("sh.tangled.git.listRefUpdatesBy", { - params: /*#__PURE__*/ v.object({ - /** - * Pagination cursor - */ - cursor: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.string()), - /** - * @minimum 1 - * @maximum 1000 - * @default 50 - */ - limit: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.constrain(/*#__PURE__*/ v.integer(), [ - /*#__PURE__*/ v.integerRange(1, 1000), - ]), - 50, - ), - /** - * Absolute offset for random-access pagination. Mutually exclusive with cursor; offsets drift under concurrent writes, so follow up with the returned cursor. - * @minimum 0 - */ - offset: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.integer()), - /** - * Sort direction by createdAt. - * @default "desc" - */ - order: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.string<"asc" | "desc" | (string & {})>(), - "desc", - ), - /** - * Actor DID whose ref-update authorings to list. - */ - subject: /*#__PURE__*/ v.didString(), - }), - output: { - type: "lex", - schema: /*#__PURE__*/ v.object({ - cursor: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.string()), - get items() { - return /*#__PURE__*/ v.array(ShTangledGitListRefUpdates.listItemSchema); - }, - /** - * Total items in the full list; omitted for filtered or merged views - */ - total: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.integer()), - }), - }, -}); - -type main$schematype = typeof _mainSchema; - -export interface mainSchema extends main$schematype {} - -export const mainSchema = _mainSchema as mainSchema; - -export interface $params extends v.InferInput {} -export interface $output extends v.InferXRPCBodyInput {} - -declare module "@atcute/lexicons/ambient" { - interface XRPCQueries { - "sh.tangled.git.listRefUpdatesBy": mainSchema; - } -} diff --git a/web/src/lib/api/lexicons/types/sh/tangled/git/ref.ts b/web/src/lib/api/lexicons/types/sh/tangled/git/ref.ts new file mode 100644 index 000000000..5b7bc4497 --- /dev/null +++ b/web/src/lib/api/lexicons/types/sh/tangled/git/ref.ts @@ -0,0 +1,32 @@ +import type {} from "@atcute/lexicons"; +import * as v from "@atcute/lexicons/validations"; +import type {} from "@atcute/lexicons/ambient"; + +const _mainSchema = /*#__PURE__*/ v.record( + /*#__PURE__*/ v.string(), + /*#__PURE__*/ v.object({ + $type: /*#__PURE__*/ v.literal("sh.tangled.git.ref"), + /** + * Object the ref points at + * @minLength 40 + * @maxLength 64 + */ + sha: /*#__PURE__*/ v.constrain(/*#__PURE__*/ v.string(), [ + /*#__PURE__*/ v.stringLength(40, 64), + ]), + }), +); + +type main$schematype = typeof _mainSchema; + +export interface mainSchema extends main$schematype {} + +export const mainSchema = _mainSchema as mainSchema; + +export interface Main extends v.InferInput {} + +declare module "@atcute/lexicons/ambient" { + interface Records { + "sh.tangled.git.ref": mainSchema; + } +} diff --git a/web/src/lib/api/lexicons/types/sh/tangled/knot/subscribeRepos.ts b/web/src/lib/api/lexicons/types/sh/tangled/knot/subscribeRepos.ts deleted file mode 100644 index 4ad550395..000000000 --- a/web/src/lib/api/lexicons/types/sh/tangled/knot/subscribeRepos.ts +++ /dev/null @@ -1,90 +0,0 @@ -import type {} from "@atcute/lexicons"; -import * as v from "@atcute/lexicons/validations"; -import type {} from "@atcute/lexicons/ambient"; -import * as ShTangledGitRefUpdate from "../git/refUpdate.js"; - -const _gitSync1Schema = /*#__PURE__*/ v.object({ - $type: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.literal("sh.tangled.knot.subscribeRepos#gitSync1"), - ), - /** - * Repository DID identifier - */ - did: /*#__PURE__*/ v.didString(), - /** - * The stream sequence number of this message. - */ - seq: /*#__PURE__*/ v.integer(), -}); -const _gitSync2Schema = /*#__PURE__*/ v.object({ - $type: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.literal("sh.tangled.knot.subscribeRepos#gitSync2"), - ), - /** - * Repository AT-URI identifier - */ - repo: /*#__PURE__*/ v.resourceUriString(), - /** - * The stream sequence number of this message. - */ - seq: /*#__PURE__*/ v.integer(), -}); -const _identitySchema = /*#__PURE__*/ v.object({ - $type: /*#__PURE__*/ v.optional( - /*#__PURE__*/ v.literal("sh.tangled.knot.subscribeRepos#identity"), - ), - /** - * Repository DID identifier - */ - did: /*#__PURE__*/ v.didString(), - /** - * The stream sequence number of this message. - */ - seq: /*#__PURE__*/ v.integer(), - time: /*#__PURE__*/ v.datetimeString(), -}); -const _mainSchema = /*#__PURE__*/ v.subscription( - "sh.tangled.knot.subscribeRepos", - { - params: /*#__PURE__*/ v.object({ - /** - * The last known event seq number to backfill from. - */ - cursor: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.integer()), - }), - get message() { - return /*#__PURE__*/ v.variant([ - ShTangledGitRefUpdate.mainSchema, - identitySchema, - ]); - }, - }, -); - -type gitSync1$schematype = typeof _gitSync1Schema; -type gitSync2$schematype = typeof _gitSync2Schema; -type identity$schematype = typeof _identitySchema; -type main$schematype = typeof _mainSchema; - -export interface gitSync1Schema extends gitSync1$schematype {} -export interface gitSync2Schema extends gitSync2$schematype {} -export interface identitySchema extends identity$schematype {} -export interface mainSchema extends main$schematype {} - -export const gitSync1Schema = _gitSync1Schema as gitSync1Schema; -export const gitSync2Schema = _gitSync2Schema as gitSync2Schema; -export const identitySchema = _identitySchema as identitySchema; -export const mainSchema = _mainSchema as mainSchema; - -export interface GitSync1 extends v.InferInput {} -export interface GitSync2 extends v.InferInput {} -export interface Identity extends v.InferInput {} - -export interface $params extends v.InferInput {} -export type $message = v.InferInput; - -declare module "@atcute/lexicons/ambient" { - interface XRPCSubscriptions { - "sh.tangled.knot.subscribeRepos": mainSchema; - } -} -- 2.51.2