From 8ba39017a0524fde65b72250d6e9099927522b7c Mon Sep 17 00:00:00 2001 From: Hailey Date: Tue, 20 May 2025 15:18:06 -0700 Subject: [PATCH 1/3] codegen --- api/bsky/actorstatus.go | 2 +- api/bsky/feeddefs.go | 4 ++++ api/bsky/feedgetFeedSkeleton.go | 2 ++ api/bsky/unspeccedgetConfig.go | 9 ++++++++- api/ozone/moderationdefs.go | 12 ++++++------ api/ozone/verificationdefs.go | 12 ++++++------ 6 files changed, 27 insertions(+), 14 deletions(-) diff --git a/api/bsky/actorstatus.go b/api/bsky/actorstatus.go index 2dfbc0f1..54d08d55 100644 --- a/api/bsky/actorstatus.go +++ b/api/bsky/actorstatus.go @@ -21,7 +21,7 @@ func init() { type ActorStatus struct { LexiconTypeID string `json:"$type,const=app.bsky.actor.status" cborgen:"$type,const=app.bsky.actor.status"` CreatedAt string `json:"createdAt" cborgen:"createdAt"` - // durationMinutes: The duration of the status in minutes. Applications can choose to limit the duration. + // durationMinutes: The duration of the status in minutes. Applications can choose to impose minimum and maximum limits. DurationMinutes *int64 `json:"durationMinutes,omitempty" cborgen:"durationMinutes,omitempty"` // embed: An optional embed associated with the status. Embed *ActorStatus_Embed `json:"embed,omitempty" cborgen:"embed,omitempty"` diff --git a/api/bsky/feeddefs.go b/api/bsky/feeddefs.go index 5c7e64a4..bd42e260 100644 --- a/api/bsky/feeddefs.go +++ b/api/bsky/feeddefs.go @@ -35,6 +35,8 @@ type FeedDefs_FeedViewPost struct { Post *FeedDefs_PostView `json:"post" cborgen:"post"` Reason *FeedDefs_FeedViewPost_Reason `json:"reason,omitempty" cborgen:"reason,omitempty"` Reply *FeedDefs_ReplyRef `json:"reply,omitempty" cborgen:"reply,omitempty"` + // reqId: Unique identifier per request that may be passed back alongside interactions. + ReqId *string `json:"reqId,omitempty" cborgen:"reqId,omitempty"` } type FeedDefs_FeedViewPost_Reason struct { @@ -104,6 +106,8 @@ type FeedDefs_Interaction struct { // feedContext: Context on a feed item that was originally supplied by the feed generator on getFeedSkeleton. FeedContext *string `json:"feedContext,omitempty" cborgen:"feedContext,omitempty"` Item *string `json:"item,omitempty" cborgen:"item,omitempty"` + // reqId: Unique identifier per request that may be passed back alongside interactions. + ReqId *string `json:"reqId,omitempty" cborgen:"reqId,omitempty"` } // FeedDefs_NotFoundPost is a "notFoundPost" in the app.bsky.feed.defs schema. diff --git a/api/bsky/feedgetFeedSkeleton.go b/api/bsky/feedgetFeedSkeleton.go index f4f94b9a..7f2a5ba2 100644 --- a/api/bsky/feedgetFeedSkeleton.go +++ b/api/bsky/feedgetFeedSkeleton.go @@ -14,6 +14,8 @@ import ( type FeedGetFeedSkeleton_Output struct { Cursor *string `json:"cursor,omitempty" cborgen:"cursor,omitempty"` Feed []*FeedDefs_SkeletonFeedPost `json:"feed" cborgen:"feed"` + // reqId: Unique identifier per request that may be passed back alongside interactions. + ReqId *string `json:"reqId,omitempty" cborgen:"reqId,omitempty"` } // FeedGetFeedSkeleton calls the XRPC method "app.bsky.feed.getFeedSkeleton". diff --git a/api/bsky/unspeccedgetConfig.go b/api/bsky/unspeccedgetConfig.go index 7bc72834..d88e0c0e 100644 --- a/api/bsky/unspeccedgetConfig.go +++ b/api/bsky/unspeccedgetConfig.go @@ -10,9 +10,16 @@ import ( "github.com/bluesky-social/indigo/xrpc" ) +// UnspeccedGetConfig_LiveNowConfig is a "liveNowConfig" in the app.bsky.unspecced.getConfig schema. +type UnspeccedGetConfig_LiveNowConfig struct { + Did string `json:"did" cborgen:"did"` + Domains []string `json:"domains" cborgen:"domains"` +} + // UnspeccedGetConfig_Output is the output of a app.bsky.unspecced.getConfig call. type UnspeccedGetConfig_Output struct { - CheckEmailConfirmed *bool `json:"checkEmailConfirmed,omitempty" cborgen:"checkEmailConfirmed,omitempty"` + CheckEmailConfirmed *bool `json:"checkEmailConfirmed,omitempty" cborgen:"checkEmailConfirmed,omitempty"` + LiveNow []*UnspeccedGetConfig_LiveNowConfig `json:"liveNow,omitempty" cborgen:"liveNow,omitempty"` } // UnspeccedGetConfig calls the XRPC method "app.bsky.unspecced.getConfig". diff --git a/api/ozone/moderationdefs.go b/api/ozone/moderationdefs.go index 19f48191..1f0263d0 100644 --- a/api/ozone/moderationdefs.go +++ b/api/ozone/moderationdefs.go @@ -1043,12 +1043,12 @@ func (t *ModerationDefs_SubjectStatusView_Subject) UnmarshalJSON(b []byte) error // // Detailed view of a subject. For record subjects, the author's repo and profile will be returned. type ModerationDefs_SubjectView struct { - //Profile *ModerationDefs_SubjectView_Profile `json:"profile,omitempty" cborgen:"profile,omitempty"` - Record *ModerationDefs_RecordViewDetail `json:"record,omitempty" cborgen:"record,omitempty"` - Repo *ModerationDefs_RepoViewDetail `json:"repo,omitempty" cborgen:"repo,omitempty"` - Status *ModerationDefs_SubjectStatusView `json:"status,omitempty" cborgen:"status,omitempty"` - Subject string `json:"subject" cborgen:"subject"` - Type *string `json:"type" cborgen:"type"` + Profile *ModerationDefs_SubjectView_Profile `json:"profile,omitempty" cborgen:"profile,omitempty"` + Record *ModerationDefs_RecordViewDetail `json:"record,omitempty" cborgen:"record,omitempty"` + Repo *ModerationDefs_RepoViewDetail `json:"repo,omitempty" cborgen:"repo,omitempty"` + Status *ModerationDefs_SubjectStatusView `json:"status,omitempty" cborgen:"status,omitempty"` + Subject string `json:"subject" cborgen:"subject"` + Type *string `json:"type" cborgen:"type"` } // ModerationDefs_VideoDetails is a "videoDetails" in the tools.ozone.moderation.defs schema. diff --git a/api/ozone/verificationdefs.go b/api/ozone/verificationdefs.go index c615fb12..85d3943d 100644 --- a/api/ozone/verificationdefs.go +++ b/api/ozone/verificationdefs.go @@ -22,9 +22,9 @@ type VerificationDefs_VerificationView struct { // handle: Handle of the subject the verification applies to at the moment of verifying, which might not be the same at the time of viewing. The verification is only valid if the current handle matches the one at the time of verifying. Handle string `json:"handle" cborgen:"handle"` // issuer: The user who issued this verification. - Issuer string `json:"issuer" cborgen:"issuer"` - //IssuerProfile *VerificationDefs_VerificationView_IssuerProfile `json:"issuerProfile,omitempty" cborgen:"issuerProfile,omitempty"` - //IssuerRepo *VerificationDefs_VerificationView_IssuerRepo `json:"issuerRepo,omitempty" cborgen:"issuerRepo,omitempty"` + Issuer string `json:"issuer" cborgen:"issuer"` + IssuerProfile *VerificationDefs_VerificationView_IssuerProfile `json:"issuerProfile,omitempty" cborgen:"issuerProfile,omitempty"` + IssuerRepo *VerificationDefs_VerificationView_IssuerRepo `json:"issuerRepo,omitempty" cborgen:"issuerRepo,omitempty"` // revokeReason: Describes the reason for revocation, also indicating that the verification is no longer valid. RevokeReason *string `json:"revokeReason,omitempty" cborgen:"revokeReason,omitempty"` // revokedAt: Timestamp when the verification was revoked. @@ -32,9 +32,9 @@ type VerificationDefs_VerificationView struct { // revokedBy: The user who revoked this verification. RevokedBy *string `json:"revokedBy,omitempty" cborgen:"revokedBy,omitempty"` // subject: The subject of the verification. - Subject string `json:"subject" cborgen:"subject"` - //SubjectProfile *VerificationDefs_VerificationView_SubjectProfile `json:"subjectProfile,omitempty" cborgen:"subjectProfile,omitempty"` - SubjectRepo *VerificationDefs_VerificationView_SubjectRepo `json:"subjectRepo,omitempty" cborgen:"subjectRepo,omitempty"` + Subject string `json:"subject" cborgen:"subject"` + SubjectProfile *VerificationDefs_VerificationView_SubjectProfile `json:"subjectProfile,omitempty" cborgen:"subjectProfile,omitempty"` + SubjectRepo *VerificationDefs_VerificationView_SubjectRepo `json:"subjectRepo,omitempty" cborgen:"subjectRepo,omitempty"` // uri: The AT-URI of the verification record. Uri string `json:"uri" cborgen:"uri"` } -- 2.51.2 From e203ea0e3c89ef04627a32e557c5d3f123038693 Mon Sep 17 00:00:00 2001 From: whyrusleeping Date: Tue, 20 May 2025 17:21:20 -0700 Subject: [PATCH 2/3] remove all the unused event fields --- api/atproto/cbor_gen.go | 697 +----------------- api/atproto/syncsubscribeRepos.go | 29 - api/ozone/moderationdefs.go | 1 - api/ozone/verificationdefs.go | 2 - bgs/bgs.go | 78 -- bgs/fedmgr.go | 45 -- cmd/goat/firehose.go | 18 - cmd/gosky/debug.go | 16 - cmd/gosky/main.go | 41 -- cmd/gosky/streamdiff.go | 10 - cmd/relay/relay/ingest.go | 9 - cmd/relay/relay/slurper.go | 27 - cmd/relay/stream/consumer.go | 79 +- cmd/relay/stream/events.go | 61 +- .../stream/persist/diskpersist/diskpersist.go | 34 - cmd/relay/testing/consumer.go | 8 - cmd/sonar/sonar.go | 36 - cmd/supercollider/main.go | 9 - events/consumer.go | 76 +- events/dbpersist/dbpersist.go | 88 --- events/diskpersist/diskpersist.go | 34 - events/events.go | 61 +- events/persist.go | 6 - events/yolopersist/yolopersist.go | 6 - gen/main.go | 3 - pds/server.go | 20 - search/firehose.go | 16 - testing/integ_test.go | 2 - testing/utils.go | 7 - 29 files changed, 67 insertions(+), 1452 deletions(-) diff --git a/api/atproto/cbor_gen.go b/api/atproto/cbor_gen.go index dec76f18..98fd5c87 100644 --- a/api/atproto/cbor_gen.go +++ b/api/atproto/cbor_gen.go @@ -1206,222 +1206,7 @@ func (t *SyncSubscribeRepos_Sync) UnmarshalCBOR(r io.Reader) (err error) { return nil } -func (t *SyncSubscribeRepos_Handle) MarshalCBOR(w io.Writer) error { - if t == nil { - _, err := w.Write(cbg.CborNull) - return err - } - - cw := cbg.NewCborWriter(w) - - if _, err := cw.Write([]byte{164}); err != nil { - return err - } - - // t.Did (string) (string) - if len("did") > 1000000 { - return xerrors.Errorf("Value in field \"did\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("did"))); err != nil { - return err - } - if _, err := cw.WriteString(string("did")); err != nil { - return err - } - - if len(t.Did) > 1000000 { - return xerrors.Errorf("Value in field t.Did was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.Did))); err != nil { - return err - } - if _, err := cw.WriteString(string(t.Did)); err != nil { - return err - } - - // t.Seq (int64) (int64) - if len("seq") > 1000000 { - return xerrors.Errorf("Value in field \"seq\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("seq"))); err != nil { - return err - } - if _, err := cw.WriteString(string("seq")); err != nil { - return err - } - - if t.Seq >= 0 { - if err := cw.WriteMajorTypeHeader(cbg.MajUnsignedInt, uint64(t.Seq)); err != nil { - return err - } - } else { - if err := cw.WriteMajorTypeHeader(cbg.MajNegativeInt, uint64(-t.Seq-1)); err != nil { - return err - } - } - - // t.Time (string) (string) - if len("time") > 1000000 { - return xerrors.Errorf("Value in field \"time\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("time"))); err != nil { - return err - } - if _, err := cw.WriteString(string("time")); err != nil { - return err - } - - if len(t.Time) > 1000000 { - return xerrors.Errorf("Value in field t.Time was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.Time))); err != nil { - return err - } - if _, err := cw.WriteString(string(t.Time)); err != nil { - return err - } - - // t.Handle (string) (string) - if len("handle") > 1000000 { - return xerrors.Errorf("Value in field \"handle\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("handle"))); err != nil { - return err - } - if _, err := cw.WriteString(string("handle")); err != nil { - return err - } - - if len(t.Handle) > 1000000 { - return xerrors.Errorf("Value in field t.Handle was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.Handle))); err != nil { - return err - } - if _, err := cw.WriteString(string(t.Handle)); err != nil { - return err - } - return nil -} - -func (t *SyncSubscribeRepos_Handle) UnmarshalCBOR(r io.Reader) (err error) { - *t = SyncSubscribeRepos_Handle{} - - 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("SyncSubscribeRepos_Handle: map struct too large (%d)", extra) - } - - n := extra - - nameBuf := make([]byte, 6) - 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.Did (string) (string) - case "did": - - { - sval, err := cbg.ReadStringWithMax(cr, 1000000) - if err != nil { - return err - } - - t.Did = string(sval) - } - // t.Seq (int64) (int64) - case "seq": - { - maj, extra, err := cr.ReadHeader() - if err != nil { - return err - } - var extraI int64 - switch maj { - case cbg.MajUnsignedInt: - extraI = int64(extra) - if extraI < 0 { - return fmt.Errorf("int64 positive overflow") - } - case cbg.MajNegativeInt: - extraI = int64(extra) - if extraI < 0 { - return fmt.Errorf("int64 negative overflow") - } - extraI = -1 - extraI - default: - return fmt.Errorf("wrong type for int64 field: %d", maj) - } - - t.Seq = int64(extraI) - } - // t.Time (string) (string) - case "time": - - { - sval, err := cbg.ReadStringWithMax(cr, 1000000) - if err != nil { - return err - } - - t.Time = string(sval) - } - // t.Handle (string) (string) - case "handle": - - { - sval, err := cbg.ReadStringWithMax(cr, 1000000) - if err != nil { - return err - } - - t.Handle = 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 *SyncSubscribeRepos_Identity) MarshalCBOR(w io.Writer) error { if t == nil { _, err := w.Write(cbg.CborNull) @@ -2094,305 +1879,74 @@ func (t *SyncSubscribeRepos_Info) UnmarshalCBOR(r io.Reader) (err error) { return nil } -func (t *SyncSubscribeRepos_Migrate) MarshalCBOR(w io.Writer) error { + +func (t *SyncSubscribeRepos_RepoOp) MarshalCBOR(w io.Writer) error { if t == nil { _, err := w.Write(cbg.CborNull) return err } cw := cbg.NewCborWriter(w) + fieldCount := 4 + + if t.Prev == nil { + fieldCount-- + } - if _, err := cw.Write([]byte{164}); err != nil { + if _, err := cw.Write(cbg.CborEncodeMajorType(cbg.MajMap, uint64(fieldCount))); err != nil { return err } - // t.Did (string) (string) - if len("did") > 1000000 { - return xerrors.Errorf("Value in field \"did\" was too long") + // t.Cid (util.LexLink) (struct) + if len("cid") > 1000000 { + return xerrors.Errorf("Value in field \"cid\" was too long") } - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("did"))); err != nil { + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("cid"))); err != nil { return err } - if _, err := cw.WriteString(string("did")); err != nil { + if _, err := cw.WriteString(string("cid")); err != nil { return err } - if len(t.Did) > 1000000 { - return xerrors.Errorf("Value in field t.Did was too long") + if err := t.Cid.MarshalCBOR(cw); err != nil { + return err } - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.Did))); err != nil { + // t.Path (string) (string) + if len("path") > 1000000 { + return xerrors.Errorf("Value in field \"path\" was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("path"))); err != nil { return err } - if _, err := cw.WriteString(string(t.Did)); err != nil { + if _, err := cw.WriteString(string("path")); err != nil { return err } - // t.Seq (int64) (int64) - if len("seq") > 1000000 { - return xerrors.Errorf("Value in field \"seq\" was too long") + if len(t.Path) > 1000000 { + return xerrors.Errorf("Value in field t.Path was too long") } - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("seq"))); err != nil { + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.Path))); err != nil { return err } - if _, err := cw.WriteString(string("seq")); err != nil { + if _, err := cw.WriteString(string(t.Path)); err != nil { return err } - if t.Seq >= 0 { - if err := cw.WriteMajorTypeHeader(cbg.MajUnsignedInt, uint64(t.Seq)); err != nil { + // t.Prev (util.LexLink) (struct) + if t.Prev != nil { + + if len("prev") > 1000000 { + return xerrors.Errorf("Value in field \"prev\" was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("prev"))); err != nil { return err } - } else { - if err := cw.WriteMajorTypeHeader(cbg.MajNegativeInt, uint64(-t.Seq-1)); err != nil { - return err - } - } - - // t.Time (string) (string) - if len("time") > 1000000 { - return xerrors.Errorf("Value in field \"time\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("time"))); err != nil { - return err - } - if _, err := cw.WriteString(string("time")); err != nil { - return err - } - - if len(t.Time) > 1000000 { - return xerrors.Errorf("Value in field t.Time was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.Time))); err != nil { - return err - } - if _, err := cw.WriteString(string(t.Time)); err != nil { - return err - } - - // t.MigrateTo (string) (string) - if len("migrateTo") > 1000000 { - return xerrors.Errorf("Value in field \"migrateTo\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("migrateTo"))); err != nil { - return err - } - if _, err := cw.WriteString(string("migrateTo")); err != nil { - return err - } - - if t.MigrateTo == nil { - if _, err := cw.Write(cbg.CborNull); err != nil { - return err - } - } else { - if len(*t.MigrateTo) > 1000000 { - return xerrors.Errorf("Value in field t.MigrateTo was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(*t.MigrateTo))); err != nil { - return err - } - if _, err := cw.WriteString(string(*t.MigrateTo)); err != nil { - return err - } - } - return nil -} - -func (t *SyncSubscribeRepos_Migrate) UnmarshalCBOR(r io.Reader) (err error) { - *t = SyncSubscribeRepos_Migrate{} - - 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("SyncSubscribeRepos_Migrate: map struct too large (%d)", extra) - } - - n := extra - - nameBuf := make([]byte, 9) - 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.Did (string) (string) - case "did": - - { - sval, err := cbg.ReadStringWithMax(cr, 1000000) - if err != nil { - return err - } - - t.Did = string(sval) - } - // t.Seq (int64) (int64) - case "seq": - { - maj, extra, err := cr.ReadHeader() - if err != nil { - return err - } - var extraI int64 - switch maj { - case cbg.MajUnsignedInt: - extraI = int64(extra) - if extraI < 0 { - return fmt.Errorf("int64 positive overflow") - } - case cbg.MajNegativeInt: - extraI = int64(extra) - if extraI < 0 { - return fmt.Errorf("int64 negative overflow") - } - extraI = -1 - extraI - default: - return fmt.Errorf("wrong type for int64 field: %d", maj) - } - - t.Seq = int64(extraI) - } - // t.Time (string) (string) - case "time": - - { - sval, err := cbg.ReadStringWithMax(cr, 1000000) - if err != nil { - return err - } - - t.Time = string(sval) - } - // t.MigrateTo (string) (string) - case "migrateTo": - - { - b, err := cr.ReadByte() - if err != nil { - return err - } - if b != cbg.CborNull[0] { - if err := cr.UnreadByte(); err != nil { - return err - } - - sval, err := cbg.ReadStringWithMax(cr, 1000000) - if err != nil { - return err - } - - t.MigrateTo = (*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 *SyncSubscribeRepos_RepoOp) MarshalCBOR(w io.Writer) error { - if t == nil { - _, err := w.Write(cbg.CborNull) - return err - } - - cw := cbg.NewCborWriter(w) - fieldCount := 4 - - if t.Prev == nil { - fieldCount-- - } - - if _, err := cw.Write(cbg.CborEncodeMajorType(cbg.MajMap, uint64(fieldCount))); err != nil { - return err - } - - // t.Cid (util.LexLink) (struct) - if len("cid") > 1000000 { - return xerrors.Errorf("Value in field \"cid\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("cid"))); err != nil { - return err - } - if _, err := cw.WriteString(string("cid")); err != nil { - return err - } - - if err := t.Cid.MarshalCBOR(cw); err != nil { - return err - } - - // t.Path (string) (string) - if len("path") > 1000000 { - return xerrors.Errorf("Value in field \"path\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("path"))); err != nil { - return err - } - if _, err := cw.WriteString(string("path")); err != nil { - return err - } - - if len(t.Path) > 1000000 { - return xerrors.Errorf("Value in field t.Path was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.Path))); err != nil { - return err - } - if _, err := cw.WriteString(string(t.Path)); err != nil { - return err - } - - // t.Prev (util.LexLink) (struct) - if t.Prev != nil { - - if len("prev") > 1000000 { - return xerrors.Errorf("Value in field \"prev\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("prev"))); err != nil { - return err - } - if _, err := cw.WriteString(string("prev")); err != nil { + if _, err := cw.WriteString(string("prev")); err != nil { return err } @@ -2540,188 +2094,7 @@ func (t *SyncSubscribeRepos_RepoOp) UnmarshalCBOR(r io.Reader) (err error) { return nil } -func (t *SyncSubscribeRepos_Tombstone) MarshalCBOR(w io.Writer) error { - if t == nil { - _, err := w.Write(cbg.CborNull) - return err - } - - cw := cbg.NewCborWriter(w) - - if _, err := cw.Write([]byte{163}); err != nil { - return err - } - - // t.Did (string) (string) - if len("did") > 1000000 { - return xerrors.Errorf("Value in field \"did\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("did"))); err != nil { - return err - } - if _, err := cw.WriteString(string("did")); err != nil { - return err - } - - if len(t.Did) > 1000000 { - return xerrors.Errorf("Value in field t.Did was too long") - } - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.Did))); err != nil { - return err - } - if _, err := cw.WriteString(string(t.Did)); err != nil { - return err - } - - // t.Seq (int64) (int64) - if len("seq") > 1000000 { - return xerrors.Errorf("Value in field \"seq\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("seq"))); err != nil { - return err - } - if _, err := cw.WriteString(string("seq")); err != nil { - return err - } - - if t.Seq >= 0 { - if err := cw.WriteMajorTypeHeader(cbg.MajUnsignedInt, uint64(t.Seq)); err != nil { - return err - } - } else { - if err := cw.WriteMajorTypeHeader(cbg.MajNegativeInt, uint64(-t.Seq-1)); err != nil { - return err - } - } - - // t.Time (string) (string) - if len("time") > 1000000 { - return xerrors.Errorf("Value in field \"time\" was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("time"))); err != nil { - return err - } - if _, err := cw.WriteString(string("time")); err != nil { - return err - } - - if len(t.Time) > 1000000 { - return xerrors.Errorf("Value in field t.Time was too long") - } - - if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len(t.Time))); err != nil { - return err - } - if _, err := cw.WriteString(string(t.Time)); err != nil { - return err - } - return nil -} - -func (t *SyncSubscribeRepos_Tombstone) UnmarshalCBOR(r io.Reader) (err error) { - *t = SyncSubscribeRepos_Tombstone{} - - 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("SyncSubscribeRepos_Tombstone: map struct too large (%d)", extra) - } - - n := extra - - nameBuf := make([]byte, 4) - 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.Did (string) (string) - case "did": - - { - sval, err := cbg.ReadStringWithMax(cr, 1000000) - if err != nil { - return err - } - - t.Did = string(sval) - } - // t.Seq (int64) (int64) - case "seq": - { - maj, extra, err := cr.ReadHeader() - if err != nil { - return err - } - var extraI int64 - switch maj { - case cbg.MajUnsignedInt: - extraI = int64(extra) - if extraI < 0 { - return fmt.Errorf("int64 positive overflow") - } - case cbg.MajNegativeInt: - extraI = int64(extra) - if extraI < 0 { - return fmt.Errorf("int64 negative overflow") - } - extraI = -1 - extraI - default: - return fmt.Errorf("wrong type for int64 field: %d", maj) - } - - t.Seq = int64(extraI) - } - // t.Time (string) (string) - case "time": - - { - sval, err := cbg.ReadStringWithMax(cr, 1000000) - if err != nil { - return err - } - - t.Time = 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 *LabelDefs_SelfLabels) MarshalCBOR(w io.Writer) error { if t == nil { _, err := w.Write(cbg.CborNull) diff --git a/api/atproto/syncsubscribeRepos.go b/api/atproto/syncsubscribeRepos.go index ea460415..5bb08125 100644 --- a/api/atproto/syncsubscribeRepos.go +++ b/api/atproto/syncsubscribeRepos.go @@ -49,16 +49,6 @@ type SyncSubscribeRepos_Commit struct { TooBig bool `json:"tooBig" cborgen:"tooBig"` } -// SyncSubscribeRepos_Handle is a "handle" in the com.atproto.sync.subscribeRepos schema. -// -// DEPRECATED -- Use #identity event instead -type SyncSubscribeRepos_Handle struct { - Did string `json:"did" cborgen:"did"` - Handle string `json:"handle" cborgen:"handle"` - Seq int64 `json:"seq" cborgen:"seq"` - Time string `json:"time" cborgen:"time"` -} - // SyncSubscribeRepos_Identity is a "identity" in the com.atproto.sync.subscribeRepos schema. // // Represents a change to an account's identity. Could be an updated handle, signing key, or pds hosting endpoint. Serves as a prod to all downstream services to refresh their identity cache. @@ -76,16 +66,6 @@ type SyncSubscribeRepos_Info struct { Name string `json:"name" cborgen:"name"` } -// SyncSubscribeRepos_Migrate is a "migrate" in the com.atproto.sync.subscribeRepos schema. -// -// DEPRECATED -- Use #account event instead -type SyncSubscribeRepos_Migrate struct { - Did string `json:"did" cborgen:"did"` - MigrateTo *string `json:"migrateTo" cborgen:"migrateTo"` - Seq int64 `json:"seq" cborgen:"seq"` - Time string `json:"time" cborgen:"time"` -} - // SyncSubscribeRepos_RepoOp is a "repoOp" in the com.atproto.sync.subscribeRepos schema. // // A repo operation, ie a mutation of a single record. @@ -113,12 +93,3 @@ type SyncSubscribeRepos_Sync struct { // time: Timestamp of when this message was originally broadcast. Time string `json:"time" cborgen:"time"` } - -// SyncSubscribeRepos_Tombstone is a "tombstone" in the com.atproto.sync.subscribeRepos schema. -// -// DEPRECATED -- Use #account event instead -type SyncSubscribeRepos_Tombstone struct { - Did string `json:"did" cborgen:"did"` - Seq int64 `json:"seq" cborgen:"seq"` - Time string `json:"time" cborgen:"time"` -} diff --git a/api/ozone/moderationdefs.go b/api/ozone/moderationdefs.go index 1f0263d0..eb0784a0 100644 --- a/api/ozone/moderationdefs.go +++ b/api/ozone/moderationdefs.go @@ -1043,7 +1043,6 @@ func (t *ModerationDefs_SubjectStatusView_Subject) UnmarshalJSON(b []byte) error // // Detailed view of a subject. For record subjects, the author's repo and profile will be returned. type ModerationDefs_SubjectView struct { - Profile *ModerationDefs_SubjectView_Profile `json:"profile,omitempty" cborgen:"profile,omitempty"` Record *ModerationDefs_RecordViewDetail `json:"record,omitempty" cborgen:"record,omitempty"` Repo *ModerationDefs_RepoViewDetail `json:"repo,omitempty" cborgen:"repo,omitempty"` Status *ModerationDefs_SubjectStatusView `json:"status,omitempty" cborgen:"status,omitempty"` diff --git a/api/ozone/verificationdefs.go b/api/ozone/verificationdefs.go index 85d3943d..b3bbec84 100644 --- a/api/ozone/verificationdefs.go +++ b/api/ozone/verificationdefs.go @@ -23,7 +23,6 @@ type VerificationDefs_VerificationView struct { Handle string `json:"handle" cborgen:"handle"` // issuer: The user who issued this verification. Issuer string `json:"issuer" cborgen:"issuer"` - IssuerProfile *VerificationDefs_VerificationView_IssuerProfile `json:"issuerProfile,omitempty" cborgen:"issuerProfile,omitempty"` IssuerRepo *VerificationDefs_VerificationView_IssuerRepo `json:"issuerRepo,omitempty" cborgen:"issuerRepo,omitempty"` // revokeReason: Describes the reason for revocation, also indicating that the verification is no longer valid. RevokeReason *string `json:"revokeReason,omitempty" cborgen:"revokeReason,omitempty"` @@ -33,7 +32,6 @@ type VerificationDefs_VerificationView struct { RevokedBy *string `json:"revokedBy,omitempty" cborgen:"revokedBy,omitempty"` // subject: The subject of the verification. Subject string `json:"subject" cborgen:"subject"` - SubjectProfile *VerificationDefs_VerificationView_SubjectProfile `json:"subjectProfile,omitempty" cborgen:"subjectProfile,omitempty"` SubjectRepo *VerificationDefs_VerificationView_SubjectRepo `json:"subjectRepo,omitempty" cborgen:"subjectRepo,omitempty"` // uri: The AT-URI of the verification record. Uri string `json:"uri" cborgen:"uri"` diff --git a/bgs/bgs.go b/bgs/bgs.go index 435c042c..5ab61c5f 100644 --- a/bgs/bgs.go +++ b/bgs/bgs.go @@ -902,37 +902,6 @@ func (bgs *BGS) handleFedEvent(ctx context.Context, host *models.PDS, env *event } repoCommitsResultCounter.WithLabelValues(host.Host, "ok").Inc() - return nil - case env.RepoHandle != nil: - bgs.log.Info("bgs got repo handle event", "did", env.RepoHandle.Did, "handle", env.RepoHandle.Handle) - // Flush any cached DID documents for this user - bgs.didr.FlushCacheFor(env.RepoHandle.Did) - - // TODO: ignoring the data in the message and just going out to the DID doc - act, err := bgs.createExternalUser(ctx, env.RepoHandle.Did) - if err != nil { - return err - } - - if act.Handle.String != env.RepoHandle.Handle { - bgs.log.Warn("handle update did not update handle to asserted value", "did", env.RepoHandle.Did, "expected", env.RepoHandle.Handle, "actual", act.Handle) - } - - // TODO: Update the ReposHandle event type to include "verified" or something - - // Broadcast the handle update to all consumers - err = bgs.events.AddEvent(ctx, &events.XRPCStreamEvent{ - RepoHandle: &comatproto.SyncSubscribeRepos_Handle{ - Did: env.RepoHandle.Did, - Handle: env.RepoHandle.Handle, - Time: env.RepoHandle.Time, - }, - }) - if err != nil { - bgs.log.Error("failed to broadcast RepoHandle event", "error", err, "did", env.RepoHandle.Did, "handle", env.RepoHandle.Handle) - return fmt.Errorf("failed to broadcast RepoHandle event: %w", err) - } - return nil case env.RepoIdentity != nil: bgs.log.Info("bgs got identity event", "did", env.RepoIdentity.Did) @@ -1034,59 +1003,12 @@ func (bgs *BGS) handleFedEvent(ctx context.Context, host *models.PDS, env *event return fmt.Errorf("failed to broadcast Account event: %w", err) } - return nil - case env.RepoMigrate != nil: - if _, err := bgs.createExternalUser(ctx, env.RepoMigrate.Did); err != nil { - return err - } - - return nil - case env.RepoTombstone != nil: - if err := bgs.handleRepoTombstone(ctx, host, env.RepoTombstone); err != nil { - return err - } - return nil default: return fmt.Errorf("invalid fed event") } } -func (bgs *BGS) handleRepoTombstone(ctx context.Context, pds *models.PDS, evt *atproto.SyncSubscribeRepos_Tombstone) error { - u, err := bgs.lookupUserByDid(ctx, evt.Did) - if err != nil { - return err - } - - if u.PDS != pds.ID { - return fmt.Errorf("unauthoritative tombstone event from %s for %s", pds.Host, evt.Did) - } - - if err := bgs.db.Model(&User{}).Where("id = ?", u.ID).UpdateColumns(map[string]any{ - "tombstoned": true, - "handle": nil, - }).Error; err != nil { - return err - } - u.SetTombstoned(true) - - if err := bgs.db.Model(&models.ActorInfo{}).Where("uid = ?", u.ID).UpdateColumns(map[string]any{ - "handle": nil, - }).Error; err != nil { - return err - } - - // delete data from carstore - if err := bgs.repoman.TakeDownRepo(ctx, u.ID); err != nil { - // don't let a failure here prevent us from propagating this event - bgs.log.Error("failed to delete user data from carstore", "err", err) - } - - return bgs.events.AddEvent(ctx, &events.XRPCStreamEvent{ - RepoTombstone: evt, - }) -} - // TODO: rename? This also updates users, and 'external' is an old phrasing func (s *BGS) createExternalUser(ctx context.Context, did string) (*models.ActorInfo, error) { ctx, span := tracer.Start(ctx, "createExternalUser") diff --git a/bgs/fedmgr.go b/bgs/fedmgr.go index 710c15cb..617a5f11 100644 --- a/bgs/fedmgr.go +++ b/bgs/fedmgr.go @@ -569,51 +569,6 @@ func (s *Slurper) handleConnection(ctx context.Context, host *models.PDS, con *w return nil }, - RepoHandle: func(evt *comatproto.SyncSubscribeRepos_Handle) error { - log.Info("got remote handle update event", "pdsHost", host.Host, "did", evt.Did, "handle", evt.Handle) - if err := s.cb(context.TODO(), host, &events.XRPCStreamEvent{ - RepoHandle: evt, - }); err != nil { - log.Error("failed handling event", "host", host.Host, "seq", evt.Seq, "err", err) - } - *lastCursor = evt.Seq - - if err := s.updateCursor(sub, *lastCursor); err != nil { - return fmt.Errorf("updating cursor: %w", err) - } - - return nil - }, - RepoMigrate: func(evt *comatproto.SyncSubscribeRepos_Migrate) error { - log.Info("got remote repo migrate event", "pdsHost", host.Host, "did", evt.Did, "migrateTo", evt.MigrateTo) - if err := s.cb(context.TODO(), host, &events.XRPCStreamEvent{ - RepoMigrate: evt, - }); err != nil { - log.Error("failed handling event", "host", host.Host, "seq", evt.Seq, "err", err) - } - *lastCursor = evt.Seq - - if err := s.updateCursor(sub, *lastCursor); err != nil { - return fmt.Errorf("updating cursor: %w", err) - } - - return nil - }, - RepoTombstone: func(evt *comatproto.SyncSubscribeRepos_Tombstone) error { - log.Info("got remote repo tombstone event", "pdsHost", host.Host, "did", evt.Did) - if err := s.cb(context.TODO(), host, &events.XRPCStreamEvent{ - RepoTombstone: evt, - }); err != nil { - log.Error("failed handling event", "host", host.Host, "seq", evt.Seq, "err", err) - } - *lastCursor = evt.Seq - - if err := s.updateCursor(sub, *lastCursor); err != nil { - return fmt.Errorf("updating cursor: %w", err) - } - - return nil - }, RepoInfo: func(info *comatproto.SyncSubscribeRepos_Info) error { log.Info("info event", "name", info.Name, "message", info.Message, "pdsHost", host.Host) return nil diff --git a/cmd/goat/firehose.go b/cmd/goat/firehose.go index 7201ab9c..d006ab97 100644 --- a/cmd/goat/firehose.go +++ b/cmd/goat/firehose.go @@ -189,24 +189,6 @@ func runFirehose(cctx *cli.Context) error { } return nil }, - RepoHandle: func(evt *comatproto.SyncSubscribeRepos_Handle) error { - if gfc.VerifyBasic { - slog.Info("deprecated event type", "eventType", "handle", "did", evt.Did, "seq", evt.Seq) - } - return nil - }, - RepoMigrate: func(evt *comatproto.SyncSubscribeRepos_Migrate) error { - if gfc.VerifyBasic { - slog.Info("deprecated event type", "eventType", "migrate", "did", evt.Did, "seq", evt.Seq) - } - return nil - }, - RepoTombstone: func(evt *comatproto.SyncSubscribeRepos_Tombstone) error { - if gfc.VerifyBasic { - slog.Info("deprecated event type", "eventType", "handle", "did", evt.Did, "seq", evt.Seq) - } - return nil - }, } scheduler := parallel.NewScheduler( diff --git a/cmd/gosky/debug.go b/cmd/gosky/debug.go index 159cf988..8d8a55b1 100644 --- a/cmd/gosky/debug.go +++ b/cmd/gosky/debug.go @@ -262,22 +262,6 @@ var debugStreamCmd = &cli.Command{ return nil }, - RepoHandle: func(evt *comatproto.SyncSubscribeRepos_Handle) error { - fmt.Printf("\rChecking seq: %d ", evt.Seq) - if lastSeq > 0 && evt.Seq != lastSeq+1 { - fmt.Println("Gap in sequence numbers: ", lastSeq, evt.Seq) - } - lastSeq = evt.Seq - return nil - }, - RepoTombstone: func(evt *comatproto.SyncSubscribeRepos_Tombstone) error { - fmt.Printf("\rChecking seq: %d ", evt.Seq) - if lastSeq > 0 && evt.Seq != lastSeq+1 { - fmt.Println("Gap in sequence numbers: ", lastSeq, evt.Seq) - } - lastSeq = evt.Seq - return nil - }, RepoInfo: func(evt *comatproto.SyncSubscribeRepos_Info) error { return nil }, diff --git a/cmd/gosky/main.go b/cmd/gosky/main.go index 3b47c29e..d37679cc 100644 --- a/cmd/gosky/main.go +++ b/cmd/gosky/main.go @@ -284,33 +284,6 @@ var readRepoStreamCmd = &cli.Command{ return nil }, - RepoMigrate: func(migrate *comatproto.SyncSubscribeRepos_Migrate) error { - if jsonfmt { - b, err := json.Marshal(migrate) - if err != nil { - return err - } - fmt.Println(string(b)) - } else { - fmt.Printf("(%d) RepoMigrate: %s moving to: %s\n", migrate.Seq, migrate.Did, *migrate.MigrateTo) - } - - return nil - }, - RepoHandle: func(handle *comatproto.SyncSubscribeRepos_Handle) error { - if jsonfmt { - b, err := json.Marshal(handle) - if err != nil { - return err - } - fmt.Println(string(b)) - } else { - fmt.Printf("(%d) RepoHandle: %s (changed to: %s)\n", handle.Seq, handle.Did, handle.Handle) - } - - return nil - - }, RepoInfo: func(info *comatproto.SyncSubscribeRepos_Info) error { if jsonfmt { b, err := json.Marshal(info) @@ -324,20 +297,6 @@ var readRepoStreamCmd = &cli.Command{ return nil }, - RepoTombstone: func(tomb *comatproto.SyncSubscribeRepos_Tombstone) error { - if jsonfmt { - b, err := json.Marshal(tomb) - if err != nil { - return err - } - fmt.Println(string(b)) - } else { - fmt.Printf("(%d) Tombstone: %s\n", tomb.Seq, tomb.Did) - } - - return nil - - }, // TODO: all the other event types Error: func(errf *events.ErrorFrame) error { return fmt.Errorf("error frame: %s: %s", errf.Error, errf.Message) diff --git a/cmd/gosky/streamdiff.go b/cmd/gosky/streamdiff.go index 89b8a8b1..508e21d1 100644 --- a/cmd/gosky/streamdiff.go +++ b/cmd/gosky/streamdiff.go @@ -127,14 +127,8 @@ func evtOp(evt *events.XRPCStreamEvent) string { return "ERROR" case evt.RepoCommit != nil: return "#commit" - case evt.RepoHandle != nil: - return "#handle" case evt.RepoInfo != nil: return "#info" - case evt.RepoMigrate != nil: - return "#migrate" - case evt.RepoTombstone != nil: - return "#tombstone" default: return "unknown" } @@ -157,10 +151,6 @@ func findEvt(evt *events.XRPCStreamEvent, list []*events.XRPCStreamEvent) int { if sameCommit(evt.RepoCommit, oe.RepoCommit) { return i } - case evt.RepoHandle != nil: - panic("not handling handle updates yet") - case evt.RepoMigrate != nil: - panic("not handling repo migrates yet") default: panic("unhandled event type: " + evtop) } diff --git a/cmd/relay/relay/ingest.go b/cmd/relay/relay/ingest.go index 03583d23..e507cba7 100644 --- a/cmd/relay/relay/ingest.go +++ b/cmd/relay/relay/ingest.go @@ -43,15 +43,6 @@ func (r *Relay) processRepoEvent(ctx context.Context, evt *stream.XRPCStreamEven case evt.RepoAccount != nil: //repoAccountReceivedCounter.WithLabelValues(hostname).Add(1) return r.processAccountEvent(ctx, evt.RepoAccount, hostname, hostID) - case evt.RepoHandle != nil: // DEPRECATED - eventsWarningsCounter.WithLabelValues(hostname, "handle").Add(1) - return nil - case evt.RepoMigrate != nil: // DEPRECATED - eventsWarningsCounter.WithLabelValues(hostname, "migrate").Add(1) - return nil - case evt.RepoTombstone != nil: // DEPRECATED - eventsWarningsCounter.WithLabelValues(hostname, "tombstone").Add(1) - return nil default: return fmt.Errorf("unhandled repo stream event type") } diff --git a/cmd/relay/relay/slurper.go b/cmd/relay/relay/slurper.go index bf4ab332..bc87f9ea 100644 --- a/cmd/relay/relay/slurper.go +++ b/cmd/relay/relay/slurper.go @@ -454,33 +454,6 @@ func (s *Slurper) handleConnection(ctx context.Context, conn *websocket.Conn, su s.logger.Debug("info event", "name", info.Name, "message", info.Message, "host", sub.Hostname) return nil }, - RepoHandle: func(evt *comatproto.SyncSubscribeRepos_Handle) error { // DEPRECATED - logger := s.logger.With("host", sub.Hostname, "did", evt.Did, "seq", evt.Seq, "eventType", "handle") - logger.Debug("got remote handle update event", "handle", evt.Handle) - if err := s.processCallback(context.Background(), &stream.XRPCStreamEvent{RepoHandle: evt}, sub.Hostname, sub.HostID); err != nil { - logger.Error("failed handling event", "err", err) - } - sub.UpdateSeq() - return nil - }, - RepoMigrate: func(evt *comatproto.SyncSubscribeRepos_Migrate) error { // DEPRECATED - logger := s.logger.With("host", sub.Hostname, "did", evt.Did, "seq", evt.Seq, "eventType", "migrate") - logger.Debug("got remote repo migrate event", "migrateTo", evt.MigrateTo) - if err := s.processCallback(context.Background(), &stream.XRPCStreamEvent{RepoMigrate: evt}, sub.Hostname, sub.HostID); err != nil { - logger.Error("failed handling event", "err", err) - } - sub.UpdateSeq() - return nil - }, - RepoTombstone: func(evt *comatproto.SyncSubscribeRepos_Tombstone) error { // DEPRECATED - logger := s.logger.With("host", sub.Hostname, "did", evt.Did, "seq", evt.Seq, "eventType", "tombstone") - logger.Debug("got remote repo tombstone event") - if err := s.processCallback(context.Background(), &stream.XRPCStreamEvent{RepoTombstone: evt}, sub.Hostname, sub.HostID); err != nil { - logger.Error("failed handling event", "err", err) - } - sub.UpdateSeq() - return nil - }, } limiters := []*slidingwindow.Limiter{ diff --git a/cmd/relay/stream/consumer.go b/cmd/relay/stream/consumer.go index 872b5b46..b6a571aa 100644 --- a/cmd/relay/stream/consumer.go +++ b/cmd/relay/stream/consumer.go @@ -18,17 +18,14 @@ import ( const MaxMessageBytes = 5_000_000 type RepoStreamCallbacks struct { - RepoCommit func(evt *comatproto.SyncSubscribeRepos_Commit) error - RepoSync func(evt *comatproto.SyncSubscribeRepos_Sync) error - RepoHandle func(evt *comatproto.SyncSubscribeRepos_Handle) error - RepoIdentity func(evt *comatproto.SyncSubscribeRepos_Identity) error - RepoAccount func(evt *comatproto.SyncSubscribeRepos_Account) error - RepoInfo func(evt *comatproto.SyncSubscribeRepos_Info) error - RepoMigrate func(evt *comatproto.SyncSubscribeRepos_Migrate) error - RepoTombstone func(evt *comatproto.SyncSubscribeRepos_Tombstone) error - LabelLabels func(evt *comatproto.LabelSubscribeLabels_Labels) error - LabelInfo func(evt *comatproto.LabelSubscribeLabels_Info) error - Error func(evt *ErrorFrame) error + RepoCommit func(evt *comatproto.SyncSubscribeRepos_Commit) error + RepoSync func(evt *comatproto.SyncSubscribeRepos_Sync) error + RepoIdentity func(evt *comatproto.SyncSubscribeRepos_Identity) error + RepoAccount func(evt *comatproto.SyncSubscribeRepos_Account) error + RepoInfo func(evt *comatproto.SyncSubscribeRepos_Info) error + LabelLabels func(evt *comatproto.LabelSubscribeLabels_Labels) error + LabelInfo func(evt *comatproto.LabelSubscribeLabels_Info) error + Error func(evt *ErrorFrame) error } func (rsc *RepoStreamCallbacks) EventHandler(ctx context.Context, xev *XRPCStreamEvent) error { @@ -37,18 +34,12 @@ func (rsc *RepoStreamCallbacks) EventHandler(ctx context.Context, xev *XRPCStrea return rsc.RepoCommit(xev.RepoCommit) case xev.RepoSync != nil && rsc.RepoSync != nil: return rsc.RepoSync(xev.RepoSync) - case xev.RepoHandle != nil && rsc.RepoHandle != nil: - return rsc.RepoHandle(xev.RepoHandle) case xev.RepoInfo != nil && rsc.RepoInfo != nil: return rsc.RepoInfo(xev.RepoInfo) - case xev.RepoMigrate != nil && rsc.RepoMigrate != nil: - return rsc.RepoMigrate(xev.RepoMigrate) case xev.RepoIdentity != nil && rsc.RepoIdentity != nil: return rsc.RepoIdentity(xev.RepoIdentity) case xev.RepoAccount != nil && rsc.RepoAccount != nil: return rsc.RepoAccount(xev.RepoAccount) - case xev.RepoTombstone != nil && rsc.RepoTombstone != nil: - return rsc.RepoTombstone(xev.RepoTombstone) case xev.LabelLabels != nil && rsc.LabelLabels != nil: return rsc.LabelLabels(xev.LabelLabels) case xev.LabelInfo != nil && rsc.LabelInfo != nil: @@ -248,24 +239,6 @@ func HandleRepoStream(ctx context.Context, con *websocket.Conn, sched Scheduler, }); err != nil { return err } - case "#handle": - // TODO: DEPRECATED message; warning/counter; drop message - var evt comatproto.SyncSubscribeRepos_Handle - if err := evt.UnmarshalCBOR(r); err != nil { - return err - } - - if evt.Seq <= lastSeq { - logger.Error("got events out of order from stream", "seq", evt.Seq, "prev", lastSeq) - continue - } - lastSeq = evt.Seq - - if err := sched.AddWork(ctx, evt.Did, &XRPCStreamEvent{ - RepoHandle: &evt, - }); err != nil { - return err - } case "#identity": var evt comatproto.SyncSubscribeRepos_Identity if err := evt.UnmarshalCBOR(r); err != nil { @@ -312,42 +285,6 @@ func HandleRepoStream(ctx context.Context, con *websocket.Conn, sched Scheduler, }); err != nil { return err } - case "#migrate": - // TODO: DEPRECATED message; warning/counter; drop message - var evt comatproto.SyncSubscribeRepos_Migrate - if err := evt.UnmarshalCBOR(r); err != nil { - return err - } - - if evt.Seq <= lastSeq { - logger.Error("got events out of order from stream", "seq", evt.Seq, "prev", lastSeq) - continue - } - lastSeq = evt.Seq - - if err := sched.AddWork(ctx, evt.Did, &XRPCStreamEvent{ - RepoMigrate: &evt, - }); err != nil { - return err - } - case "#tombstone": - // TODO: DEPRECATED message; warning/counter; drop message - var evt comatproto.SyncSubscribeRepos_Tombstone - if err := evt.UnmarshalCBOR(r); err != nil { - return err - } - - if evt.Seq <= lastSeq { - logger.Error("got events out of order from stream", "seq", evt.Seq, "prev", lastSeq) - continue - } - lastSeq = evt.Seq - - if err := sched.AddWork(ctx, evt.Did, &XRPCStreamEvent{ - RepoTombstone: &evt, - }); err != nil { - return err - } case "#labels": var evt comatproto.LabelSubscribeLabels_Labels if err := evt.UnmarshalCBOR(r); err != nil { diff --git a/cmd/relay/stream/events.go b/cmd/relay/stream/events.go index 7b38e34f..e327aad3 100644 --- a/cmd/relay/stream/events.go +++ b/cmd/relay/stream/events.go @@ -23,17 +23,14 @@ type EventHeader struct { } type XRPCStreamEvent struct { - Error *ErrorFrame - RepoCommit *comatproto.SyncSubscribeRepos_Commit - RepoSync *comatproto.SyncSubscribeRepos_Sync - RepoHandle *comatproto.SyncSubscribeRepos_Handle // DEPRECATED - RepoIdentity *comatproto.SyncSubscribeRepos_Identity - RepoInfo *comatproto.SyncSubscribeRepos_Info - RepoMigrate *comatproto.SyncSubscribeRepos_Migrate // DEPRECATED - RepoTombstone *comatproto.SyncSubscribeRepos_Tombstone // DEPRECATED - RepoAccount *comatproto.SyncSubscribeRepos_Account - LabelLabels *comatproto.LabelSubscribeLabels_Labels - LabelInfo *comatproto.LabelSubscribeLabels_Info + Error *ErrorFrame + RepoCommit *comatproto.SyncSubscribeRepos_Commit + RepoSync *comatproto.SyncSubscribeRepos_Sync + RepoIdentity *comatproto.SyncSubscribeRepos_Identity + RepoInfo *comatproto.SyncSubscribeRepos_Info + RepoAccount *comatproto.SyncSubscribeRepos_Account + LabelLabels *comatproto.LabelSubscribeLabels_Labels + LabelInfo *comatproto.LabelSubscribeLabels_Info // some private fields for internal routing perf PrivUid uint64 `json:"-" cborgen:"-"` @@ -56,9 +53,6 @@ func (evt *XRPCStreamEvent) Serialize(wc io.Writer) error { case evt.RepoSync != nil: header.MsgType = "#sync" obj = evt.RepoSync - case evt.RepoHandle != nil: - header.MsgType = "#handle" - obj = evt.RepoHandle case evt.RepoIdentity != nil: header.MsgType = "#identity" obj = evt.RepoIdentity @@ -68,12 +62,6 @@ func (evt *XRPCStreamEvent) Serialize(wc io.Writer) error { case evt.RepoInfo != nil: header.MsgType = "#info" obj = evt.RepoInfo - case evt.RepoMigrate != nil: - header.MsgType = "#migrate" - obj = evt.RepoMigrate - case evt.RepoTombstone != nil: - header.MsgType = "#tombstone" - obj = evt.RepoTombstone default: return fmt.Errorf("unrecognized event kind") } @@ -105,13 +93,6 @@ func (xevt *XRPCStreamEvent) Deserialize(r io.Reader) error { return fmt.Errorf("reading repoSync event: %w", err) } xevt.RepoSync = &evt - case "#handle": - // TODO: DEPRECATED message; warning/counter; drop message - var evt comatproto.SyncSubscribeRepos_Handle - if err := evt.UnmarshalCBOR(r); err != nil { - return err - } - xevt.RepoHandle = &evt case "#identity": var evt comatproto.SyncSubscribeRepos_Identity if err := evt.UnmarshalCBOR(r); err != nil { @@ -131,20 +112,6 @@ func (xevt *XRPCStreamEvent) Deserialize(r io.Reader) error { return err } xevt.RepoInfo = &evt - case "#migrate": - // TODO: DEPRECATED message; warning/counter; drop message - var evt comatproto.SyncSubscribeRepos_Migrate - if err := evt.UnmarshalCBOR(r); err != nil { - return err - } - xevt.RepoMigrate = &evt - case "#tombstone": - // TODO: DEPRECATED message; warning/counter; drop message - var evt comatproto.SyncSubscribeRepos_Tombstone - if err := evt.UnmarshalCBOR(r); err != nil { - return err - } - xevt.RepoTombstone = &evt case "#labels": var evt comatproto.LabelSubscribeLabels_Labels if err := evt.UnmarshalCBOR(r); err != nil { @@ -193,12 +160,6 @@ func (evt *XRPCStreamEvent) Sequence() int64 { return evt.RepoCommit.Seq case evt.RepoSync != nil: return evt.RepoSync.Seq - case evt.RepoHandle != nil: - return evt.RepoHandle.Seq - case evt.RepoMigrate != nil: - return evt.RepoMigrate.Seq - case evt.RepoTombstone != nil: - return evt.RepoTombstone.Seq case evt.RepoIdentity != nil: return evt.RepoIdentity.Seq case evt.RepoAccount != nil: @@ -220,12 +181,6 @@ func (evt *XRPCStreamEvent) GetSequence() (int64, bool) { return evt.RepoCommit.Seq, true case evt.RepoSync != nil: return evt.RepoSync.Seq, true - case evt.RepoHandle != nil: - return evt.RepoHandle.Seq, true - case evt.RepoMigrate != nil: - return evt.RepoMigrate.Seq, true - case evt.RepoTombstone != nil: - return evt.RepoTombstone.Seq, true case evt.RepoIdentity != nil: return evt.RepoIdentity.Seq, true case evt.RepoAccount != nil: diff --git a/cmd/relay/stream/persist/diskpersist/diskpersist.go b/cmd/relay/stream/persist/diskpersist/diskpersist.go index df2f52ce..516c984d 100644 --- a/cmd/relay/stream/persist/diskpersist/diskpersist.go +++ b/cmd/relay/stream/persist/diskpersist/diskpersist.go @@ -484,14 +484,10 @@ func (dp *DiskPersistence) doPersist(ctx context.Context, pjob persistJob) error pjob.Evt.RepoCommit.Seq = seq case pjob.Evt.RepoSync != nil: pjob.Evt.RepoSync.Seq = seq - case pjob.Evt.RepoHandle != nil: - pjob.Evt.RepoHandle.Seq = seq case pjob.Evt.RepoIdentity != nil: pjob.Evt.RepoIdentity.Seq = seq case pjob.Evt.RepoAccount != nil: pjob.Evt.RepoAccount.Seq = seq - case pjob.Evt.RepoTombstone != nil: - pjob.Evt.RepoTombstone.Seq = seq default: // only those three get peristed right now // we should not actually ever get here... @@ -547,12 +543,6 @@ func (dp *DiskPersistence) Persist(ctx context.Context, xevt *stream.XRPCStreamE if err := xevt.RepoSync.MarshalCBOR(cw); err != nil { return fmt.Errorf("failed to marshal: %w", err) } - case xevt.RepoHandle != nil: - evtKind = evtKindHandle - did = xevt.RepoHandle.Did - if err := xevt.RepoHandle.MarshalCBOR(cw); err != nil { - return fmt.Errorf("failed to marshal: %w", err) - } case xevt.RepoIdentity != nil: evtKind = evtKindIdentity did = xevt.RepoIdentity.Did @@ -565,12 +555,6 @@ func (dp *DiskPersistence) Persist(ctx context.Context, xevt *stream.XRPCStreamE if err := xevt.RepoAccount.MarshalCBOR(cw); err != nil { return fmt.Errorf("failed to marshal: %w", err) } - case xevt.RepoTombstone != nil: - evtKind = evtKindTombstone - did = xevt.RepoTombstone.Did - if err := xevt.RepoTombstone.MarshalCBOR(cw); err != nil { - return fmt.Errorf("failed to marshal: %w", err) - } default: return nil // only those two get peristed right now @@ -810,15 +794,6 @@ func (dp *DiskPersistence) readEventsFrom(ctx context.Context, since int64, fn s if err := cb(&stream.XRPCStreamEvent{RepoSync: &evt}); err != nil { return nil, err } - case evtKindHandle: - var evt atproto.SyncSubscribeRepos_Handle - if err := evt.UnmarshalCBOR(io.LimitReader(bufr, h.Len64())); err != nil { - return nil, err - } - evt.Seq = h.Seq - if err := cb(&stream.XRPCStreamEvent{RepoHandle: &evt}); err != nil { - return nil, err - } case evtKindIdentity: var evt atproto.SyncSubscribeRepos_Identity if err := evt.UnmarshalCBOR(io.LimitReader(bufr, h.Len64())); err != nil { @@ -837,15 +812,6 @@ func (dp *DiskPersistence) readEventsFrom(ctx context.Context, since int64, fn s if err := cb(&stream.XRPCStreamEvent{RepoAccount: &evt}); err != nil { return nil, err } - case evtKindTombstone: - var evt atproto.SyncSubscribeRepos_Tombstone - if err := evt.UnmarshalCBOR(io.LimitReader(bufr, h.Len64())); err != nil { - return nil, err - } - evt.Seq = h.Seq - if err := cb(&stream.XRPCStreamEvent{RepoTombstone: &evt}); err != nil { - return nil, err - } default: dp.log.Warn("unrecognized event kind coming from log file", "seq", h.Seq, "kind", h.Kind) return nil, fmt.Errorf("halting on unrecognized event kind") diff --git a/cmd/relay/testing/consumer.go b/cmd/relay/testing/consumer.go index 2c517e62..25a3c08f 100644 --- a/cmd/relay/testing/consumer.go +++ b/cmd/relay/testing/consumer.go @@ -62,14 +62,6 @@ func (c *Consumer) eventCallbacks() *stream.RepoStreamCallbacks { c.LastSeq = evt.Seq return nil }, - // NOTE: this is included to test that the events are *not* passed through; can be removed in the near future - RepoHandle: func(evt *comatproto.SyncSubscribeRepos_Handle) error { - c.eventsLk.Lock() - defer c.eventsLk.Unlock() - c.Events = append(c.Events, &stream.XRPCStreamEvent{RepoHandle: evt}) - c.LastSeq = evt.Seq - return nil - }, } return rsc } diff --git a/cmd/sonar/sonar.go b/cmd/sonar/sonar.go index 8b08e1d2..17056ea6 100644 --- a/cmd/sonar/sonar.go +++ b/cmd/sonar/sonar.go @@ -109,23 +109,6 @@ func (s *Sonar) HandleStreamEvent(ctx context.Context, xe *events.XRPCStreamEven case xe.RepoCommit != nil: eventsProcessedCounter.WithLabelValues("repo_commit", s.SocketURL).Inc() return s.HandleRepoCommit(ctx, xe.RepoCommit) - case xe.RepoHandle != nil: - eventsProcessedCounter.WithLabelValues("repo_handle", s.SocketURL).Inc() - now := time.Now() - s.ProgMux.Lock() - s.Progress.LastSeq = xe.RepoHandle.Seq - s.Progress.LastSeqProcessedAt = now - s.ProgMux.Unlock() - // Parse time from the event time string - t, err := time.Parse(time.RFC3339, xe.RepoHandle.Time) - if err != nil { - s.Logger.Error("error parsing time", "err", err) - return nil - } - lastEvtCreatedAtGauge.WithLabelValues(s.SocketURL).Set(float64(t.UnixNano())) - lastEvtProcessedAtGauge.WithLabelValues(s.SocketURL).Set(float64(now.UnixNano())) - lastEvtCreatedEvtProcessedGapGauge.WithLabelValues(s.SocketURL).Set(float64(now.Sub(t).Seconds())) - lastSeqGauge.WithLabelValues(s.SocketURL).Set(float64(xe.RepoHandle.Seq)) case xe.RepoIdentity != nil: eventsProcessedCounter.WithLabelValues("identity", s.SocketURL).Inc() now := time.Now() @@ -142,25 +125,6 @@ func (s *Sonar) HandleStreamEvent(ctx context.Context, xe *events.XRPCStreamEven s.ProgMux.Unlock() case xe.RepoInfo != nil: eventsProcessedCounter.WithLabelValues("repo_info", s.SocketURL).Inc() - case xe.RepoMigrate != nil: - eventsProcessedCounter.WithLabelValues("repo_migrate", s.SocketURL).Inc() - now := time.Now() - s.ProgMux.Lock() - s.Progress.LastSeq = xe.RepoMigrate.Seq - s.Progress.LastSeqProcessedAt = time.Now() - s.ProgMux.Unlock() - // Parse time from the event time string - t, err := time.Parse(time.RFC3339, xe.RepoMigrate.Time) - if err != nil { - s.Logger.Error("error parsing time", "err", err) - return nil - } - lastEvtCreatedAtGauge.WithLabelValues(s.SocketURL).Set(float64(t.UnixNano())) - lastEvtProcessedAtGauge.WithLabelValues(s.SocketURL).Set(float64(now.UnixNano())) - lastEvtCreatedEvtProcessedGapGauge.WithLabelValues(s.SocketURL).Set(float64(now.Sub(t).Seconds())) - lastSeqGauge.WithLabelValues(s.SocketURL).Set(float64(xe.RepoHandle.Seq)) - case xe.RepoTombstone != nil: - eventsProcessedCounter.WithLabelValues("repo_tombstone", s.SocketURL).Inc() case xe.LabelInfo != nil: eventsProcessedCounter.WithLabelValues("label_info", s.SocketURL).Inc() case xe.LabelLabels != nil: diff --git a/cmd/supercollider/main.go b/cmd/supercollider/main.go index 57a4d61c..928342a0 100644 --- a/cmd/supercollider/main.go +++ b/cmd/supercollider/main.go @@ -338,18 +338,9 @@ func Reload(cctx *cli.Context) error { case evt.RepoCommit != nil: header.MsgType = "#commit" obj = evt.RepoCommit - case evt.RepoHandle != nil: - header.MsgType = "#handle" - obj = evt.RepoHandle case evt.RepoInfo != nil: header.MsgType = "#info" obj = evt.RepoInfo - case evt.RepoMigrate != nil: - header.MsgType = "#migrate" - obj = evt.RepoMigrate - case evt.RepoTombstone != nil: - header.MsgType = "#tombstone" - obj = evt.RepoTombstone default: logger.Error("unrecognized event kind") continue diff --git a/events/consumer.go b/events/consumer.go index 6c832afc..b9180cd4 100644 --- a/events/consumer.go +++ b/events/consumer.go @@ -16,17 +16,14 @@ import ( ) type RepoStreamCallbacks struct { - RepoCommit func(evt *comatproto.SyncSubscribeRepos_Commit) error - RepoSync func(evt *comatproto.SyncSubscribeRepos_Sync) error - RepoHandle func(evt *comatproto.SyncSubscribeRepos_Handle) error - RepoIdentity func(evt *comatproto.SyncSubscribeRepos_Identity) error - RepoAccount func(evt *comatproto.SyncSubscribeRepos_Account) error - RepoInfo func(evt *comatproto.SyncSubscribeRepos_Info) error - RepoMigrate func(evt *comatproto.SyncSubscribeRepos_Migrate) error - RepoTombstone func(evt *comatproto.SyncSubscribeRepos_Tombstone) error - LabelLabels func(evt *comatproto.LabelSubscribeLabels_Labels) error - LabelInfo func(evt *comatproto.LabelSubscribeLabels_Info) error - Error func(evt *ErrorFrame) error + RepoCommit func(evt *comatproto.SyncSubscribeRepos_Commit) error + RepoSync func(evt *comatproto.SyncSubscribeRepos_Sync) error + RepoIdentity func(evt *comatproto.SyncSubscribeRepos_Identity) error + RepoAccount func(evt *comatproto.SyncSubscribeRepos_Account) error + RepoInfo func(evt *comatproto.SyncSubscribeRepos_Info) error + LabelLabels func(evt *comatproto.LabelSubscribeLabels_Labels) error + LabelInfo func(evt *comatproto.LabelSubscribeLabels_Info) error + Error func(evt *ErrorFrame) error } func (rsc *RepoStreamCallbacks) EventHandler(ctx context.Context, xev *XRPCStreamEvent) error { @@ -35,18 +32,12 @@ func (rsc *RepoStreamCallbacks) EventHandler(ctx context.Context, xev *XRPCStrea return rsc.RepoCommit(xev.RepoCommit) case xev.RepoSync != nil && rsc.RepoSync != nil: return rsc.RepoSync(xev.RepoSync) - case xev.RepoHandle != nil && rsc.RepoHandle != nil: - return rsc.RepoHandle(xev.RepoHandle) case xev.RepoInfo != nil && rsc.RepoInfo != nil: return rsc.RepoInfo(xev.RepoInfo) - case xev.RepoMigrate != nil && rsc.RepoMigrate != nil: - return rsc.RepoMigrate(xev.RepoMigrate) case xev.RepoIdentity != nil && rsc.RepoIdentity != nil: return rsc.RepoIdentity(xev.RepoIdentity) case xev.RepoAccount != nil && rsc.RepoAccount != nil: return rsc.RepoAccount(xev.RepoAccount) - case xev.RepoTombstone != nil && rsc.RepoTombstone != nil: - return rsc.RepoTombstone(xev.RepoTombstone) case xev.LabelLabels != nil && rsc.LabelLabels != nil: return rsc.LabelLabels(xev.LabelLabels) case xev.LabelInfo != nil && rsc.LabelInfo != nil: @@ -241,23 +232,6 @@ func HandleRepoStream(ctx context.Context, con *websocket.Conn, sched Scheduler, }); err != nil { return err } - case "#handle": - // TODO: DEPRECATED message; warning/counter; drop message - var evt comatproto.SyncSubscribeRepos_Handle - if err := evt.UnmarshalCBOR(r); err != nil { - return err - } - - if evt.Seq < lastSeq { - log.Error("Got events out of order from stream", "seq", evt.Seq, "prev", lastSeq) - } - lastSeq = evt.Seq - - if err := sched.AddWork(ctx, evt.Did, &XRPCStreamEvent{ - RepoHandle: &evt, - }); err != nil { - return err - } case "#identity": var evt comatproto.SyncSubscribeRepos_Identity if err := evt.UnmarshalCBOR(r); err != nil { @@ -302,40 +276,6 @@ func HandleRepoStream(ctx context.Context, con *websocket.Conn, sched Scheduler, }); err != nil { return err } - case "#migrate": - // TODO: DEPRECATED message; warning/counter; drop message - var evt comatproto.SyncSubscribeRepos_Migrate - if err := evt.UnmarshalCBOR(r); err != nil { - return err - } - - if evt.Seq < lastSeq { - log.Error("Got events out of order from stream", "seq", evt.Seq, "prev", lastSeq) - } - lastSeq = evt.Seq - - if err := sched.AddWork(ctx, evt.Did, &XRPCStreamEvent{ - RepoMigrate: &evt, - }); err != nil { - return err - } - case "#tombstone": - // TODO: DEPRECATED message; warning/counter; drop message - var evt comatproto.SyncSubscribeRepos_Tombstone - if err := evt.UnmarshalCBOR(r); err != nil { - return err - } - - if evt.Seq < lastSeq { - log.Error("Got events out of order from stream", "seq", evt.Seq, "prev", lastSeq) - } - lastSeq = evt.Seq - - if err := sched.AddWork(ctx, evt.Did, &XRPCStreamEvent{ - RepoTombstone: &evt, - }); err != nil { - return err - } case "#labels": var evt comatproto.LabelSubscribeLabels_Labels if err := evt.UnmarshalCBOR(r); err != nil { diff --git a/events/dbpersist/dbpersist.go b/events/dbpersist/dbpersist.go index 3b24bd5b..efd0ba2b 100644 --- a/events/dbpersist/dbpersist.go +++ b/events/dbpersist/dbpersist.go @@ -171,14 +171,10 @@ func (p *DbPersistence) flushBatchLocked(ctx context.Context) error { switch { case e.RepoCommit != nil: e.RepoCommit.Seq = int64(item.Seq) - case e.RepoHandle != nil: - e.RepoHandle.Seq = int64(item.Seq) case e.RepoIdentity != nil: e.RepoIdentity.Seq = int64(item.Seq) case e.RepoAccount != nil: e.RepoAccount.Seq = int64(item.Seq) - case e.RepoTombstone != nil: - e.RepoTombstone.Seq = int64(item.Seq) default: return fmt.Errorf("unknown event type") } @@ -218,11 +214,6 @@ func (p *DbPersistence) Persist(ctx context.Context, e *events.XRPCStreamEvent) if err != nil { return err } - case e.RepoHandle != nil: - rer, err = p.RecordFromHandleChange(ctx, e.RepoHandle) - if err != nil { - return err - } case e.RepoIdentity != nil: rer, err = p.RecordFromRepoIdentity(ctx, e.RepoIdentity) if err != nil { @@ -233,11 +224,6 @@ func (p *DbPersistence) Persist(ctx context.Context, e *events.XRPCStreamEvent) if err != nil { return err } - case e.RepoTombstone != nil: - rer, err = p.RecordFromTombstone(ctx, e.RepoTombstone) - if err != nil { - return err - } default: return nil } @@ -249,25 +235,6 @@ func (p *DbPersistence) Persist(ctx context.Context, e *events.XRPCStreamEvent) return nil } -func (p *DbPersistence) RecordFromHandleChange(ctx context.Context, evt *comatproto.SyncSubscribeRepos_Handle) (*RepoEventRecord, error) { - t, err := time.Parse(util.ISO8601, evt.Time) - if err != nil { - return nil, err - } - - uid, err := p.uidForDid(ctx, evt.Did) - if err != nil { - return nil, err - } - - return &RepoEventRecord{ - Repo: uid, - Type: "repo_handle", - Time: t, - NewHandle: &evt.Handle, - }, nil -} - func (p *DbPersistence) RecordFromRepoIdentity(ctx context.Context, evt *comatproto.SyncSubscribeRepos_Identity) (*RepoEventRecord, error) { t, err := time.Parse(util.ISO8601, evt.Time) if err != nil { @@ -306,24 +273,6 @@ func (p *DbPersistence) RecordFromRepoAccount(ctx context.Context, evt *comatpro }, nil } -func (p *DbPersistence) RecordFromTombstone(ctx context.Context, evt *comatproto.SyncSubscribeRepos_Tombstone) (*RepoEventRecord, error) { - t, err := time.Parse(util.ISO8601, evt.Time) - if err != nil { - return nil, err - } - - uid, err := p.uidForDid(ctx, evt.Did) - if err != nil { - return nil, err - } - - return &RepoEventRecord{ - Repo: uid, - Type: "repo_tombstone", - Time: t, - }, nil -} - func (p *DbPersistence) RecordFromRepoCommit(ctx context.Context, evt *comatproto.SyncSubscribeRepos_Commit) (*RepoEventRecord, error) { // TODO: hack hack hack if len(evt.Ops) > 8192 { @@ -449,14 +398,10 @@ func (p *DbPersistence) hydrateBatch(ctx context.Context, batch []*RepoEventReco switch { case record.Commit != nil: streamEvent, err = p.hydrateCommit(ctx, record) - case record.NewHandle != nil: - streamEvent, err = p.hydrateHandleChange(ctx, record) case record.Type == "repo_identity": streamEvent, err = p.hydrateIdentityEvent(ctx, record) case record.Type == "repo_account": streamEvent, err = p.hydrateAccountEvent(ctx, record) - case record.Type == "repo_tombstone": - streamEvent, err = p.hydrateTombstone(ctx, record) default: err = fmt.Errorf("unknown event type: %s", record.Type) } @@ -519,25 +464,6 @@ func (p *DbPersistence) didForUid(ctx context.Context, uid models.Uid) (string, return u.Did, nil } -func (p *DbPersistence) hydrateHandleChange(ctx context.Context, rer *RepoEventRecord) (*events.XRPCStreamEvent, error) { - if rer.NewHandle == nil { - return nil, fmt.Errorf("NewHandle is nil") - } - - did, err := p.didForUid(ctx, rer.Repo) - if err != nil { - return nil, err - } - - return &events.XRPCStreamEvent{ - RepoHandle: &comatproto.SyncSubscribeRepos_Handle{ - Did: did, - Handle: *rer.NewHandle, - Time: rer.Time.Format(util.ISO8601), - }, - }, nil -} - func (p *DbPersistence) hydrateIdentityEvent(ctx context.Context, rer *RepoEventRecord) (*events.XRPCStreamEvent, error) { did, err := p.didForUid(ctx, rer.Repo) if err != nil { @@ -568,20 +494,6 @@ func (p *DbPersistence) hydrateAccountEvent(ctx context.Context, rer *RepoEventR }, nil } -func (p *DbPersistence) hydrateTombstone(ctx context.Context, rer *RepoEventRecord) (*events.XRPCStreamEvent, error) { - did, err := p.didForUid(ctx, rer.Repo) - if err != nil { - return nil, err - } - - return &events.XRPCStreamEvent{ - RepoTombstone: &comatproto.SyncSubscribeRepos_Tombstone{ - Did: did, - Time: rer.Time.Format(util.ISO8601), - }, - }, nil -} - func (p *DbPersistence) hydrateCommit(ctx context.Context, rer *RepoEventRecord) (*events.XRPCStreamEvent, error) { if rer.Commit == nil { return nil, fmt.Errorf("commit is nil") diff --git a/events/diskpersist/diskpersist.go b/events/diskpersist/diskpersist.go index 069185fd..b662067f 100644 --- a/events/diskpersist/diskpersist.go +++ b/events/diskpersist/diskpersist.go @@ -455,14 +455,10 @@ func (dp *DiskPersistence) doPersist(ctx context.Context, j persistJob) error { switch { case e.RepoCommit != nil: e.RepoCommit.Seq = seq - case e.RepoHandle != nil: - e.RepoHandle.Seq = seq case e.RepoIdentity != nil: e.RepoIdentity.Seq = seq case e.RepoAccount != nil: e.RepoAccount.Seq = seq - case e.RepoTombstone != nil: - e.RepoTombstone.Seq = seq default: // only those three get peristed right now // we should not actually ever get here... @@ -509,12 +505,6 @@ func (dp *DiskPersistence) Persist(ctx context.Context, e *events.XRPCStreamEven if err := e.RepoCommit.MarshalCBOR(cw); err != nil { return fmt.Errorf("failed to marshal: %w", err) } - case e.RepoHandle != nil: - evtKind = evtKindHandle - did = e.RepoHandle.Did - if err := e.RepoHandle.MarshalCBOR(cw); err != nil { - return fmt.Errorf("failed to marshal: %w", err) - } case e.RepoIdentity != nil: evtKind = evtKindIdentity did = e.RepoIdentity.Did @@ -527,12 +517,6 @@ func (dp *DiskPersistence) Persist(ctx context.Context, e *events.XRPCStreamEven if err := e.RepoAccount.MarshalCBOR(cw); err != nil { return fmt.Errorf("failed to marshal: %w", err) } - case e.RepoTombstone != nil: - evtKind = evtKindTombstone - did = e.RepoTombstone.Did - if err := e.RepoTombstone.MarshalCBOR(cw); err != nil { - return fmt.Errorf("failed to marshal: %w", err) - } default: return nil // only those two get peristed right now @@ -745,15 +729,6 @@ func (dp *DiskPersistence) readEventsFrom(ctx context.Context, since int64, fn s if err := cb(&events.XRPCStreamEvent{RepoCommit: &evt}); err != nil { return nil, err } - case evtKindHandle: - var evt atproto.SyncSubscribeRepos_Handle - if err := evt.UnmarshalCBOR(io.LimitReader(bufr, h.Len64())); err != nil { - return nil, err - } - evt.Seq = h.Seq - if err := cb(&events.XRPCStreamEvent{RepoHandle: &evt}); err != nil { - return nil, err - } case evtKindIdentity: var evt atproto.SyncSubscribeRepos_Identity if err := evt.UnmarshalCBOR(io.LimitReader(bufr, h.Len64())); err != nil { @@ -772,15 +747,6 @@ func (dp *DiskPersistence) readEventsFrom(ctx context.Context, since int64, fn s if err := cb(&events.XRPCStreamEvent{RepoAccount: &evt}); err != nil { return nil, err } - case evtKindTombstone: - var evt atproto.SyncSubscribeRepos_Tombstone - if err := evt.UnmarshalCBOR(io.LimitReader(bufr, h.Len64())); err != nil { - return nil, err - } - evt.Seq = h.Seq - if err := cb(&events.XRPCStreamEvent{RepoTombstone: &evt}); err != nil { - return nil, err - } default: log.Warn("unrecognized event kind coming from log file", "seq", h.Seq, "kind", h.Kind) return nil, fmt.Errorf("halting on unrecognized event kind") diff --git a/events/events.go b/events/events.go index 5eca9f74..b056e2f1 100644 --- a/events/events.go +++ b/events/events.go @@ -187,17 +187,14 @@ func init() { } type XRPCStreamEvent struct { - Error *ErrorFrame - RepoCommit *comatproto.SyncSubscribeRepos_Commit - RepoSync *comatproto.SyncSubscribeRepos_Sync - RepoHandle *comatproto.SyncSubscribeRepos_Handle // DEPRECATED - RepoIdentity *comatproto.SyncSubscribeRepos_Identity - RepoInfo *comatproto.SyncSubscribeRepos_Info - RepoMigrate *comatproto.SyncSubscribeRepos_Migrate // DEPRECATED - RepoTombstone *comatproto.SyncSubscribeRepos_Tombstone // DEPRECATED - RepoAccount *comatproto.SyncSubscribeRepos_Account - LabelLabels *comatproto.LabelSubscribeLabels_Labels - LabelInfo *comatproto.LabelSubscribeLabels_Info + Error *ErrorFrame + RepoCommit *comatproto.SyncSubscribeRepos_Commit + RepoSync *comatproto.SyncSubscribeRepos_Sync + RepoIdentity *comatproto.SyncSubscribeRepos_Identity + RepoInfo *comatproto.SyncSubscribeRepos_Info + RepoAccount *comatproto.SyncSubscribeRepos_Account + LabelLabels *comatproto.LabelSubscribeLabels_Labels + LabelInfo *comatproto.LabelSubscribeLabels_Info // some private fields for internal routing perf PrivUid models.Uid `json:"-" cborgen:"-"` @@ -220,9 +217,6 @@ func (evt *XRPCStreamEvent) Serialize(wc io.Writer) error { case evt.RepoSync != nil: header.MsgType = "#sync" obj = evt.RepoSync - case evt.RepoHandle != nil: - header.MsgType = "#handle" - obj = evt.RepoHandle case evt.RepoIdentity != nil: header.MsgType = "#identity" obj = evt.RepoIdentity @@ -232,12 +226,6 @@ func (evt *XRPCStreamEvent) Serialize(wc io.Writer) error { case evt.RepoInfo != nil: header.MsgType = "#info" obj = evt.RepoInfo - case evt.RepoMigrate != nil: - header.MsgType = "#migrate" - obj = evt.RepoMigrate - case evt.RepoTombstone != nil: - header.MsgType = "#tombstone" - obj = evt.RepoTombstone default: return fmt.Errorf("unrecognized event kind") } @@ -269,13 +257,6 @@ func (xevt *XRPCStreamEvent) Deserialize(r io.Reader) error { return fmt.Errorf("reading repoSync event: %w", err) } xevt.RepoSync = &evt - case "#handle": - // TODO: DEPRECATED message; warning/counter; drop message - var evt comatproto.SyncSubscribeRepos_Handle - if err := evt.UnmarshalCBOR(r); err != nil { - return err - } - xevt.RepoHandle = &evt case "#identity": var evt comatproto.SyncSubscribeRepos_Identity if err := evt.UnmarshalCBOR(r); err != nil { @@ -295,20 +276,6 @@ func (xevt *XRPCStreamEvent) Deserialize(r io.Reader) error { return err } xevt.RepoInfo = &evt - case "#migrate": - // TODO: DEPRECATED message; warning/counter; drop message - var evt comatproto.SyncSubscribeRepos_Migrate - if err := evt.UnmarshalCBOR(r); err != nil { - return err - } - xevt.RepoMigrate = &evt - case "#tombstone": - // TODO: DEPRECATED message; warning/counter; drop message - var evt comatproto.SyncSubscribeRepos_Tombstone - if err := evt.UnmarshalCBOR(r); err != nil { - return err - } - xevt.RepoTombstone = &evt case "#labels": var evt comatproto.LabelSubscribeLabels_Labels if err := evt.UnmarshalCBOR(r); err != nil { @@ -476,12 +443,6 @@ func (evt *XRPCStreamEvent) Sequence() int64 { return evt.RepoCommit.Seq case evt.RepoSync != nil: return evt.RepoSync.Seq - case evt.RepoHandle != nil: - return evt.RepoHandle.Seq - case evt.RepoMigrate != nil: - return evt.RepoMigrate.Seq - case evt.RepoTombstone != nil: - return evt.RepoTombstone.Seq case evt.RepoIdentity != nil: return evt.RepoIdentity.Seq case evt.RepoAccount != nil: @@ -503,12 +464,6 @@ func (evt *XRPCStreamEvent) GetSequence() (int64, bool) { return evt.RepoCommit.Seq, true case evt.RepoSync != nil: return evt.RepoSync.Seq, true - case evt.RepoHandle != nil: - return evt.RepoHandle.Seq, true - case evt.RepoMigrate != nil: - return evt.RepoMigrate.Seq, true - case evt.RepoTombstone != nil: - return evt.RepoTombstone.Seq, true case evt.RepoIdentity != nil: return evt.RepoIdentity.Seq, true case evt.RepoAccount != nil: diff --git a/events/persist.go b/events/persist.go index 0d41db12..274d9787 100644 --- a/events/persist.go +++ b/events/persist.go @@ -42,16 +42,10 @@ func (mp *MemPersister) Persist(ctx context.Context, e *XRPCStreamEvent) error { switch { case e.RepoCommit != nil: e.RepoCommit.Seq = mp.seq - case e.RepoHandle != nil: - e.RepoHandle.Seq = mp.seq case e.RepoIdentity != nil: e.RepoIdentity.Seq = mp.seq case e.RepoAccount != nil: e.RepoAccount.Seq = mp.seq - case e.RepoMigrate != nil: - e.RepoMigrate.Seq = mp.seq - case e.RepoTombstone != nil: - e.RepoTombstone.Seq = mp.seq case e.LabelLabels != nil: e.LabelLabels.Seq = mp.seq default: diff --git a/events/yolopersist/yolopersist.go b/events/yolopersist/yolopersist.go index d6b71363..cce2faaa 100644 --- a/events/yolopersist/yolopersist.go +++ b/events/yolopersist/yolopersist.go @@ -28,16 +28,10 @@ func (yp *YoloPersister) Persist(ctx context.Context, e *events.XRPCStreamEvent) switch { case e.RepoCommit != nil: e.RepoCommit.Seq = yp.seq - case e.RepoHandle != nil: - e.RepoHandle.Seq = yp.seq case e.RepoIdentity != nil: e.RepoIdentity.Seq = yp.seq case e.RepoAccount != nil: e.RepoAccount.Seq = yp.seq - case e.RepoMigrate != nil: - e.RepoMigrate.Seq = yp.seq - case e.RepoTombstone != nil: - e.RepoTombstone.Seq = yp.seq case e.LabelLabels != nil: e.LabelLabels.Seq = yp.seq default: diff --git a/gen/main.go b/gen/main.go index 6204c49e..3fa2f2d5 100644 --- a/gen/main.go +++ b/gen/main.go @@ -98,13 +98,10 @@ func main() { atproto.RepoStrongRef{}, atproto.SyncSubscribeRepos_Commit{}, atproto.SyncSubscribeRepos_Sync{}, - atproto.SyncSubscribeRepos_Handle{}, atproto.SyncSubscribeRepos_Identity{}, atproto.SyncSubscribeRepos_Account{}, atproto.SyncSubscribeRepos_Info{}, - atproto.SyncSubscribeRepos_Migrate{}, atproto.SyncSubscribeRepos_RepoOp{}, - atproto.SyncSubscribeRepos_Tombstone{}, atproto.LabelDefs_SelfLabels{}, atproto.LabelDefs_SelfLabel{}, atproto.LabelDefs_Label{}, diff --git a/pds/server.go b/pds/server.go index 54430f72..bf61cabc 100644 --- a/pds/server.go +++ b/pds/server.go @@ -598,9 +598,6 @@ func (s *Server) EventsHandler(c echo.Context) error { case evt.RepoCommit != nil: header.MsgType = "#commit" obj = evt.RepoCommit - case evt.RepoHandle != nil: - header.MsgType = "#handle" - obj = evt.RepoHandle case evt.RepoIdentity != nil: header.MsgType = "#identity" obj = evt.RepoIdentity @@ -610,12 +607,6 @@ func (s *Server) EventsHandler(c echo.Context) error { case evt.RepoInfo != nil: header.MsgType = "#info" obj = evt.RepoInfo - case evt.RepoMigrate != nil: - header.MsgType = "#migrate" - obj = evt.RepoMigrate - case evt.RepoTombstone != nil: - header.MsgType = "#tombstone" - obj = evt.RepoTombstone default: return fmt.Errorf("unrecognized event kind") } @@ -660,17 +651,6 @@ func (s *Server) UpdateUserHandle(ctx context.Context, u *User, handle string) e return fmt.Errorf("failed to update handle: %w", err) } - if err := s.events.AddEvent(ctx, &events.XRPCStreamEvent{ - RepoHandle: &comatproto.SyncSubscribeRepos_Handle{ - Did: u.Did, - Handle: handle, - Time: time.Now().Format(util.ISO8601), - }, - }); err != nil { - return fmt.Errorf("failed to push event: %s", err) - } - - // Also push an Identity event if err := s.events.AddEvent(ctx, &events.XRPCStreamEvent{ RepoIdentity: &comatproto.SyncSubscribeRepos_Identity{ Did: u.Did, diff --git a/search/firehose.go b/search/firehose.go index 6db633cc..3d73c89e 100644 --- a/search/firehose.go +++ b/search/firehose.go @@ -115,22 +115,6 @@ func (idx *Indexer) RunIndexer(ctx context.Context) error { return nil }, - RepoHandle: func(evt *comatproto.SyncSubscribeRepos_Handle) error { - ctx := context.Background() - ctx, span := tracer.Start(ctx, "RepoHandle") - defer span.End() - - did, err := syntax.ParseDID(evt.Did) - if err != nil { - idx.logger.Error("bad DID in RepoHandle event", "did", evt.Did, "handle", evt.Handle, "seq", evt.Seq, "err", err) - return nil - } - if err := idx.updateUserHandle(ctx, did, evt.Handle); err != nil { - // TODO: handle this case (instead of return nil) - idx.logger.Error("failed to update user handle", "did", evt.Did, "handle", evt.Handle, "seq", evt.Seq, "err", err) - } - return nil - }, } return events.HandleRepoStream( diff --git a/testing/integ_test.go b/testing/integ_test.go index 45298396..1e29d8b9 100644 --- a/testing/integ_test.go +++ b/testing/integ_test.go @@ -299,8 +299,6 @@ func TestHandleChange(t *testing.T) { initevt := evts.Next() t.Log(initevt.RepoCommit) - hcevt := evts.Next() - t.Log(hcevt.RepoHandle) idevt := evts.Next() t.Log(idevt.RepoIdentity) } diff --git a/testing/utils.go b/testing/utils.go index e6362702..573a4a25 100644 --- a/testing/utils.go +++ b/testing/utils.go @@ -666,13 +666,6 @@ func (b *TestRelay) Events(t *testing.T, since int64) *EventStream { es.Lk.Unlock() return nil }, - RepoHandle: func(evt *atproto.SyncSubscribeRepos_Handle) error { - fmt.Println("received handle event: ", evt.Seq, evt.Did) - es.Lk.Lock() - es.Events = append(es.Events, &events.XRPCStreamEvent{RepoHandle: evt}) - es.Lk.Unlock() - return nil - }, RepoIdentity: func(evt *atproto.SyncSubscribeRepos_Identity) error { fmt.Println("received identity event: ", evt.Seq, evt.Did) es.Lk.Lock() -- 2.51.2 From 82b8a2993c934a589dd754383a729d232d63f589 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Tue, 3 Jun 2025 12:30:07 -0700 Subject: [PATCH 3/3] make fmt --- api/ozone/moderationdefs.go | 10 +++++----- api/ozone/verificationdefs.go | 8 ++++---- 2 files changed, 9 insertions(+), 9 deletions(-) diff --git a/api/ozone/moderationdefs.go b/api/ozone/moderationdefs.go index eb0784a0..c647bb60 100644 --- a/api/ozone/moderationdefs.go +++ b/api/ozone/moderationdefs.go @@ -1043,11 +1043,11 @@ func (t *ModerationDefs_SubjectStatusView_Subject) UnmarshalJSON(b []byte) error // // Detailed view of a subject. For record subjects, the author's repo and profile will be returned. type ModerationDefs_SubjectView struct { - Record *ModerationDefs_RecordViewDetail `json:"record,omitempty" cborgen:"record,omitempty"` - Repo *ModerationDefs_RepoViewDetail `json:"repo,omitempty" cborgen:"repo,omitempty"` - Status *ModerationDefs_SubjectStatusView `json:"status,omitempty" cborgen:"status,omitempty"` - Subject string `json:"subject" cborgen:"subject"` - Type *string `json:"type" cborgen:"type"` + Record *ModerationDefs_RecordViewDetail `json:"record,omitempty" cborgen:"record,omitempty"` + Repo *ModerationDefs_RepoViewDetail `json:"repo,omitempty" cborgen:"repo,omitempty"` + Status *ModerationDefs_SubjectStatusView `json:"status,omitempty" cborgen:"status,omitempty"` + Subject string `json:"subject" cborgen:"subject"` + Type *string `json:"type" cborgen:"type"` } // ModerationDefs_VideoDetails is a "videoDetails" in the tools.ozone.moderation.defs schema. diff --git a/api/ozone/verificationdefs.go b/api/ozone/verificationdefs.go index b3bbec84..3ece2afc 100644 --- a/api/ozone/verificationdefs.go +++ b/api/ozone/verificationdefs.go @@ -22,8 +22,8 @@ type VerificationDefs_VerificationView struct { // handle: Handle of the subject the verification applies to at the moment of verifying, which might not be the same at the time of viewing. The verification is only valid if the current handle matches the one at the time of verifying. Handle string `json:"handle" cborgen:"handle"` // issuer: The user who issued this verification. - Issuer string `json:"issuer" cborgen:"issuer"` - IssuerRepo *VerificationDefs_VerificationView_IssuerRepo `json:"issuerRepo,omitempty" cborgen:"issuerRepo,omitempty"` + Issuer string `json:"issuer" cborgen:"issuer"` + IssuerRepo *VerificationDefs_VerificationView_IssuerRepo `json:"issuerRepo,omitempty" cborgen:"issuerRepo,omitempty"` // revokeReason: Describes the reason for revocation, also indicating that the verification is no longer valid. RevokeReason *string `json:"revokeReason,omitempty" cborgen:"revokeReason,omitempty"` // revokedAt: Timestamp when the verification was revoked. @@ -31,8 +31,8 @@ type VerificationDefs_VerificationView struct { // revokedBy: The user who revoked this verification. RevokedBy *string `json:"revokedBy,omitempty" cborgen:"revokedBy,omitempty"` // subject: The subject of the verification. - Subject string `json:"subject" cborgen:"subject"` - SubjectRepo *VerificationDefs_VerificationView_SubjectRepo `json:"subjectRepo,omitempty" cborgen:"subjectRepo,omitempty"` + Subject string `json:"subject" cborgen:"subject"` + SubjectRepo *VerificationDefs_VerificationView_SubjectRepo `json:"subjectRepo,omitempty" cborgen:"subjectRepo,omitempty"` // uri: The AT-URI of the verification record. Uri string `json:"uri" cborgen:"uri"` } -- 2.51.2