diff --git a/go.mod b/go.mod index c171ff26..56dbc0df 100644 --- a/go.mod +++ b/go.mod @@ -10,6 +10,8 @@ replace github.com/gocql/gocql => github.com/scylladb/gocql v1.14.4 replace github.com/AxisCommunications/go-dpop => github.com/streamplace/go-dpop v0.0.0-20250510031900-c897158a8ad4 +replace github.com/bluesky-social/indigo => ../indigo + require ( firebase.google.com/go/v4 v4.14.1 git.stream.place/streamplace/c2pa-go v0.7.0 diff --git a/go.sum b/go.sum index a620fb41..6c973b22 100644 --- a/go.sum +++ b/go.sum @@ -119,8 +119,6 @@ github.com/bkielbasa/cyclop v1.2.3 h1:faIVMIGDIANuGPWH031CZJTi2ymOQBULs9H21HSMa5 github.com/bkielbasa/cyclop v1.2.3/go.mod h1:kHTwA9Q0uZqOADdupvcFJQtp/ksSnytRMe8ztxG8Fuo= github.com/blizzy78/varnamelen v0.8.0 h1:oqSblyuQvFsW1hbBHh1zfwrKe3kcSj0rnXkKzsQ089M= github.com/blizzy78/varnamelen v0.8.0/go.mod h1:V9TzQZ4fLJ1DSrjVDfl89H7aMnTvKkApdHeyESmyR7k= -github.com/bluesky-social/indigo v0.0.0-20250617211950-336ebe49427b h1:qc08nIOSRS3Ue5rqRjI3SgxBPS30QvGkBrMoV2+z2KA= -github.com/bluesky-social/indigo v0.0.0-20250617211950-336ebe49427b/go.mod h1:8FlFpF5cIq3DQG0kEHqyTkPV/5MDQoaWLcVwza5ZPJU= github.com/bmizerany/assert v0.0.0-20160611221934-b7ed37b82869 h1:DDGfHa7BWjL4YnC6+E63dPcxHo2sUxDIu8g3QgEJdRY= github.com/bmizerany/assert v0.0.0-20160611221934-b7ed37b82869/go.mod h1:Ekp36dRnpXw/yCqJaO+ZrUyxD+3VXMFFr56k5XYrpB4= github.com/bombsimon/wsl/v4 v4.7.0 h1:1Ilm9JBPRczjyUs6hvOPKvd7VL1Q++PL8M0SXBDf+jQ= @@ -492,8 +490,6 @@ github.com/ipfs/go-block-format v0.2.1 h1:96kW71XGNNa+mZw/MTzJrCpMhBWCrd9kBLoKm9 github.com/ipfs/go-block-format v0.2.1/go.mod h1:frtvXHMQhM6zn7HvEQu+Qz5wSTj+04oEH/I+NjDgEjk= github.com/ipfs/go-blockservice v0.5.2 h1:in9Bc+QcXwd1apOVM7Un9t8tixPKdaHQFdLSUM1Xgk8= github.com/ipfs/go-blockservice v0.5.2/go.mod h1:VpMblFEqG67A/H2sHKAemeH9vlURVavlysbdUI632yk= -github.com/ipfs/go-bs-sqlite3 v0.0.0-20221122195556-bfcee1be620d h1:9V+GGXCuOfDiFpdAHz58q9mKLg447xp0cQKvqQrAwYE= -github.com/ipfs/go-bs-sqlite3 v0.0.0-20221122195556-bfcee1be620d/go.mod h1:pMbnFyNAGjryYCLCe59YDLRv/ujdN+zGJBT1umlvYRM= github.com/ipfs/go-cid v0.5.0 h1:goEKKhaGm0ul11IHA7I6p1GmKz8kEYniqFopaB5Otwg= github.com/ipfs/go-cid v0.5.0/go.mod h1:0L7vmeNXpQpUS9vt+yEARkJ8rOg43DF3iPgn4GIN0mk= github.com/ipfs/go-datastore v0.8.2 h1:Jy3wjqQR6sg/LhyY0NIePZC3Vux19nLtg7dx0TVqr6U= diff --git a/pkg/atproto/labeler_firehose.go b/pkg/atproto/labeler_firehose.go index 9619e07c..dae283c0 100644 --- a/pkg/atproto/labeler_firehose.go +++ b/pkg/atproto/labeler_firehose.go @@ -1,6 +1,7 @@ package atproto import ( + "bytes" "context" "fmt" "net/http" @@ -15,6 +16,7 @@ import ( "golang.org/x/sync/errgroup" "stream.place/streamplace/pkg/aqhttp" "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/model" ) func (atsync *ATProtoSynchronizer) StartLabelerFirehose(ctx context.Context, did string) error { @@ -83,8 +85,18 @@ func (atsync *ATProtoSynchronizer) StartLabelerFirehoseRetry(ctx context.Context } else { return fmt.Errorf("invalid labeler URI scheme: %s", labeler.URL) } + dbLabeler, err := atsync.Model.GetLabeler(did) + if err != nil { + return fmt.Errorf("failed to get labeler %s: %w", did, err) + } + if dbLabeler == nil { + dbLabeler, err = atsync.Model.CreateLabeler(did) + if err != nil { + return fmt.Errorf("failed to create labeler %s: %w", did, err) + } + } query := u.Query() - query.Set("cursor", "0") + query.Set("cursor", fmt.Sprintf("%d", dbLabeler.Cursor)) u.RawQuery = query.Encode() con, _, err := dialer.Dial(u.String(), http.Header{ @@ -97,6 +109,10 @@ func (atsync *ATProtoSynchronizer) StartLabelerFirehoseRetry(ctx context.Context rsc := &events.RepoStreamCallbacks{ LabelLabels: func(evt *comatproto.LabelSubscribeLabels_Labels) error { log.Log(ctx, "labeler labels", "labels", evt.Labels, "seq", evt.Seq) + err = atsync.Model.UpdateLabelerCursor(did, evt.Seq) + if err != nil { + log.Error(ctx, "failed to update labeler cursor", "err", err) + } for _, labelLex := range evt.Labels { l := label.FromLexicon(labelLex) err = l.VerifySignature(pub) @@ -109,7 +125,28 @@ func (atsync *ATProtoSynchronizer) StartLabelerFirehoseRetry(ctx context.Context log.Error(ctx, "failed to verify label syntax", "err", err) continue } - log.Log(ctx, "labeler label", "cid", l.CID, "createdAt", l.CreatedAt, "expiresAt", l.ExpiresAt, "negated", l.Negated, "sourceDID", l.SourceDID, "uri", l.URI, "val", l.Val, "version", l.Version) + bs := bytes.Buffer{} + err = labelLex.MarshalCBOR(&bs) + if err != nil { + log.Error(ctx, "failed to marshal label", "err", err) + continue + } + err = atsync.Model.CreateLabel(&model.Label{ + Cid: l.CID, + Cts: l.CreatedAt, + Exp: l.ExpiresAt, + Neg: l.Negated, + Sig: l.Sig, + Src: l.SourceDID, + Uri: l.URI, + Val: l.Val, + Ver: &l.Version, + Record: bs.Bytes(), + }) + if err != nil { + log.Error(ctx, "failed to create label", "err", err) + continue + } } return nil }, diff --git a/pkg/model/label.go b/pkg/model/label.go new file mode 100644 index 00000000..33f2c3e3 --- /dev/null +++ b/pkg/model/label.go @@ -0,0 +1,37 @@ +package model + +import "gorm.io/gorm/clause" + +type Label struct { + // cid: Optionally, CID specifying the specific version of 'uri' resource this label applies to. + Cid *string `json:"cid,omitempty" cborgen:"cid,omitempty" gorm:"column:cid"` + // cts: Timestamp when this label was created. + Cts string `json:"cts" cborgen:"cts" gorm:"column:cts"` + // exp: Timestamp at which this label expires (no longer applies). + Exp *string `json:"exp,omitempty" cborgen:"exp,omitempty" gorm:"column:exp"` + // neg: If true, this is a negation label, overwriting a previous label. + Neg *bool `json:"neg,omitempty" cborgen:"neg,omitempty" gorm:"column:neg"` + // sig: Signature of dag-cbor encoded label. + Sig []byte `json:"sig,omitempty" cborgen:"sig,omitempty" gorm:"column:sig"` + // src: DID of the actor who created this label. + Src string `json:"src" cborgen:"src" gorm:"primaryKey;column:src"` + // uri: AT URI of the record, repository (account), or other resource that this label applies to. + Uri string `json:"uri" cborgen:"uri" gorm:"primaryKey;column:uri;index"` + // val: The short string name of the value or type of this label. + Val string `json:"val" cborgen:"val" gorm:"primaryKey;column:val"` + // ver: The AT Protocol version of the label object. + Ver *int64 `json:"ver,omitempty" cborgen:"ver,omitempty" gorm:"column:ver"` + + Record []byte `json:"record,omitempty" cborgen:"record,omitempty" gorm:"column:record"` +} + +func (m *DBModel) CreateLabel(label *Label) error { + return m.DB.Clauses(clause.OnConflict{ + Columns: []clause.Column{ + {Name: "src"}, + {Name: "uri"}, + {Name: "val"}, + }, + UpdateAll: true, + }).Create(label).Error +} diff --git a/pkg/model/labeler.go b/pkg/model/labeler.go new file mode 100644 index 00000000..02697f80 --- /dev/null +++ b/pkg/model/labeler.go @@ -0,0 +1,40 @@ +package model + +import ( + "errors" + + "gorm.io/gorm" +) + +type Labeler struct { + DID string `gorm:"primaryKey;column:did"` + Cursor int64 `gorm:"column:cursor"` +} + +func (m *DBModel) GetLabeler(did string) (*Labeler, error) { + var labeler Labeler + err := m.DB.Where("did = ?", did).First(&labeler).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, nil + } + if err != nil { + return nil, err + } + return &labeler, nil +} + +func (m *DBModel) CreateLabeler(did string) (*Labeler, error) { + labeler := &Labeler{ + DID: did, + Cursor: 0, + } + + if err := m.DB.Create(labeler).Error; err != nil { + return nil, err + } + return labeler, nil +} + +func (m *DBModel) UpdateLabelerCursor(did string, cursor int64) error { + return m.DB.Model(&Labeler{}).Where("did = ?", did).Update("cursor", cursor).Error +} diff --git a/pkg/model/model.go b/pkg/model/model.go index f3d71cd7..38d1ca34 100644 --- a/pkg/model/model.go +++ b/pkg/model/model.go @@ -107,6 +107,12 @@ type Model interface { GetCommitEventsSince(repoDID string, t time.Time) ([]*XrpcStreamEvent, error) GetCommitEventsSinceSeq(repoDID string, seq int64) ([]*XrpcStreamEvent, error) GetMostRecentCommitEvent(repoDID string) (*XrpcStreamEvent, error) + + CreateLabeler(did string) (*Labeler, error) + GetLabeler(did string) (*Labeler, error) + UpdateLabelerCursor(did string, cursor int64) error + + CreateLabel(label *Label) error } func MakeDB(dbURL string) (Model, error) { @@ -169,6 +175,8 @@ func MakeDB(dbURL string) (Model, error) { oatproxy.OAuthSession{}, ServerSettings{}, XrpcStreamEvent{}, + Labeler{}, + Label{}, } { err = db.AutoMigrate(model) if err != nil {