From ee00a1a1fd10ca43a14149c5e79a38d4e2deed18 Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Fri, 16 Jan 2026 19:41:29 +0900 Subject: [PATCH] wip: appview: migrate to tap ingester with partial backfill support still WIP DO NOT MERGE Signed-off-by: Seongmin Lee --- appview/db/artifact.go | 4 + appview/db/issues.go | 4 + appview/db/label.go | 4 + appview/db/profile.go | 5 + appview/db/pulls.go | 4 + appview/db/repos.go | 8 +- appview/db/strings.go | 4 + appview/ingester/ingester.go | 808 +++++++++++++++++++++++++++++++++++ appview/models/artifact.go | 26 +- appview/models/follow.go | 21 + appview/models/profile.go | 84 ++++ appview/models/pull.go | 4 + appview/models/reaction.go | 26 ++ appview/models/repo.go | 4 + appview/models/star.go | 20 + appview/repo/repo.go | 2 +- 16 files changed, 1024 insertions(+), 4 deletions(-) create mode 100644 appview/ingester/ingester.go diff --git a/appview/db/artifact.go b/appview/db/artifact.go index 68042b78..a8232e64 100644 --- a/appview/db/artifact.go +++ b/appview/db/artifact.go @@ -11,6 +11,10 @@ import ( "tangled.org/core/orm" ) +func UpsertArtifact(e Execer, artifact models.Artifact) error { + panic("unimplemented") +} + func AddArtifact(e Execer, artifact models.Artifact) error { _, err := e.Exec( `insert or ignore into artifacts ( diff --git a/appview/db/issues.go b/appview/db/issues.go index 7ab3f6aa..7c8ebd8a 100644 --- a/appview/db/issues.go +++ b/appview/db/issues.go @@ -295,6 +295,10 @@ func GetIssues(e Execer, filters ...orm.Filter) ([]models.Issue, error) { return GetIssuesPaginated(e, pagination.Page{}, filters...) } +func UpsertIssueComment(tx *sql.Tx, c models.IssueComment) (int64, error) { + panic("unimplemented") +} + func AddIssueComment(tx *sql.Tx, c models.IssueComment) (int64, error) { result, err := tx.Exec( `insert into issue_comments ( diff --git a/appview/db/label.go b/appview/db/label.go index 18d5780c..bb2c89cb 100644 --- a/appview/db/label.go +++ b/appview/db/label.go @@ -13,6 +13,10 @@ import ( "tangled.org/core/orm" ) +func UpsertLabelDefinition(e Execer, l *models.LabelDefinition) (int64, error) { + panic("unimplemented") +} + // no updating type for now func AddLabelDefinition(e Execer, l *models.LabelDefinition) (int64, error) { result, err := e.Exec( diff --git a/appview/db/profile.go b/appview/db/profile.go index a2a8ead2..9f763e3c 100644 --- a/appview/db/profile.go +++ b/appview/db/profile.go @@ -229,6 +229,11 @@ func UpsertProfile(tx *sql.Tx, profile *models.Profile) error { return nil } +func DeleteProfile(e Execer, did syntax.DID) error { + _, err := e.Exec(`delete from profiles where did = ?`, did) + return err +} + func GetProfiles(e Execer, filters ...orm.Filter) (map[string]*models.Profile, error) { var conditions []string var args []any diff --git a/appview/db/pulls.go b/appview/db/pulls.go index a4a78a27..e54a7f90 100644 --- a/appview/db/pulls.go +++ b/appview/db/pulls.go @@ -585,6 +585,10 @@ func GetPullsByOwnerDid(e Execer, did, timeframe string) ([]models.Pull, error) return pulls, nil } +func UpsertPullComment(tx *sql.Tx, comment *models.PullComment) error { + panic("unimplemented") +} + func NewPullComment(tx *sql.Tx, comment *models.PullComment) (int64, error) { query := `insert into pull_comments (owner_did, repo_at, submission_id, comment_at, pull_id, body) values (?, ?, ?, ?, ?, ?)` res, err := tx.Exec( diff --git a/appview/db/repos.go b/appview/db/repos.go index 88ec4c13..f00acc52 100644 --- a/appview/db/repos.go +++ b/appview/db/repos.go @@ -386,6 +386,10 @@ func PutRepo(tx *sql.Tx, repo models.Repo) error { return err } +func UpsertRepo(tx *sql.Tx, repo *models.Repo) error { + panic("unimplemented") +} + func AddRepo(tx *sql.Tx, repo *models.Repo) error { _, err := tx.Exec( `insert into repos @@ -409,8 +413,8 @@ func AddRepo(tx *sql.Tx, repo *models.Repo) error { return nil } -func RemoveRepo(e Execer, did, name string) error { - _, err := e.Exec(`delete from repos where did = ? and name = ?`, did, name) +func RemoveRepo(e Execer, did syntax.DID, rkey syntax.RecordKey) error { + _, err := e.Exec(`delete from repos where did = ? and rkey = ?`, did, rkey) return err } diff --git a/appview/db/strings.go b/appview/db/strings.go index dc20070d..b1df2fed 100644 --- a/appview/db/strings.go +++ b/appview/db/strings.go @@ -11,6 +11,10 @@ import ( "tangled.org/core/orm" ) +func UpsertString(e Execer, s models.String) error { + panic("unimplemented") +} + func AddString(e Execer, s models.String) error { _, err := e.Exec( `insert into strings ( diff --git a/appview/ingester/ingester.go b/appview/ingester/ingester.go new file mode 100644 index 00000000..bf286c8e --- /dev/null +++ b/appview/ingester/ingester.go @@ -0,0 +1,808 @@ +package ingester + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" + "tangled.org/core/appview/config" + "tangled.org/core/appview/db" + "tangled.org/core/appview/models" + "tangled.org/core/appview/notify" + "tangled.org/core/log" + "tangled.org/core/orm" + "tangled.org/core/rbac2" + "tangled.org/core/tapc" +) + +type Ingester struct { + cfg *config.Config + db *db.DB + e *rbac2.Enforcer + notifier notify.Notifier + l *slog.Logger +} + +// TODO: finish with rbac/v1 +// TODO: just don't notify state events for now (we need full object) + +func (i *Ingester) ProcessEvent(ctx context.Context, evt tapc.Event) error { + var err error + switch evt.Type { + case tapc.EvtRecord: + revt := evt.Record + ctx = log.IntoContext(ctx, i.l.With("record", revt.AtUri())) + // NOTE: sort by alphabetical order + switch revt.Collection.String() { + case tangled.ActorProfileNSID: + err = i.ingestActorProfile(ctx, revt) + case tangled.FeedReactionNSID: + err = i.ingestFeedReaction(ctx, revt) + case tangled.FeedStarNSID: + err = i.ingestFeedStar(ctx, revt) + case tangled.GraphFollowNSID: + err = i.ingestGraphFollow(ctx, revt) + case tangled.KnotMemberNSID: + err = i.ingestKnotMember(ctx, revt) + case tangled.KnotNSID: + err = i.ingestKnot(ctx, revt) + case tangled.LabelDefinitionNSID: + err = i.ingestLabelDefinition(ctx, revt) + case tangled.LabelOpNSID: + err = i.ingestLabelOp(ctx, revt) + case tangled.PublicKeyNSID: + err = i.ingestPublicKey(ctx, revt) + case tangled.RepoArtifactNSID: + err = i.ingestRepoArtifact(ctx, revt) + case tangled.RepoIssueCommentNSID: + err = i.ingestRepoIssueComment(ctx, revt) + case tangled.RepoIssueNSID: + err = i.ingestRepoIssue(ctx, revt) + case tangled.RepoIssueStateNSID: + err = i.ingestRepoIssueState(ctx, revt) + case tangled.RepoNSID: + err = i.ingestRepo(ctx, revt) + case tangled.RepoPullCommentNSID: + err = i.ingestRepoPullComment(ctx, revt) + case tangled.RepoPullNSID: + err = i.ingestRepoPull(ctx, revt) + case tangled.RepoPullStatusNSID: + err = i.ingestRepoPullStatus(ctx, revt) + case tangled.SpindleMemberNSID: + err = i.ingestSpindleMember(ctx, revt) + case tangled.SpindleNSID: + err = i.ingestSpindle(ctx, revt) + case tangled.StringNSID: + err = i.ingestString(ctx, revt) + } + case tapc.EvtIdentity: + // no-op + } + + if err != nil { + i.l.Error("failed to process message. will retry later", "event.ID", evt.ID, "err", err) + return err + } + return nil +} + +func (i *Ingester) ingestActorProfile(ctx context.Context, evt *tapc.RecordEventData) error { + // ignore invalid rkey + if evt.Rkey.String() != "self" { + return nil + } + + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.ActorProfile + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + profile, err := models.ProfileFromRecord(evt.Did, record) + if err != nil { + i.l.Warn("ignoring invalid profile record", "err", err) + return nil + } + + tx, err := i.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("starting transaction: %w", err) + } + defer tx.Rollback() + + if err := db.UpsertProfile(tx, &profile); err != nil { + return fmt.Errorf("upserting profile: %w", err) + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("commiting profile upsert: %w", err) + } + case tapc.RecordDeleteAction: + if err := db.DeleteProfile(i.db, evt.Did); err != nil { + return fmt.Errorf("deleting profile from db: %w", err) + } + } + return nil +} + +func (i *Ingester) ingestFeedReaction(_ context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.FeedReaction + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + reaction, err := models.ReactionFromRecord(evt.Did, evt.Rkey, record) + if err != nil { + i.l.Warn("ignoring invalid reaction record", "err", err) + return nil + } + + if err := db.UpsertReaction(i.db, reaction); err != nil { + return fmt.Errorf("upserting reaction record") + } + case tapc.RecordDeleteAction: + if err := db.DeleteReactionByRkey( + i.db, + evt.Did.String(), + evt.Rkey.String(), + ); err != nil { + return fmt.Errorf("deleting reaction from db: %w", err) + } + } + return nil +} + +func (i *Ingester) ingestFeedStar(_ context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.FeedStar + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + star, err := models.StarFromRecord(evt.Did, evt.Rkey, record) + if err != nil { + i.l.Warn("ignoring invalid star record", "err", err) + return nil + } + + if err := db.UpsertStar(i.db, star); err != nil { + return fmt.Errorf("upserting star record") + } + case tapc.RecordDeleteAction: + if err := db.DeleteStarByRkey( + i.db, + evt.Did.String(), + evt.Rkey.String(), + ); err != nil { + return fmt.Errorf("deleting record from db: %w", err) + } + } + return nil +} + +func (i *Ingester) ingestGraphFollow(_ context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.GraphFollow + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + follow, err := models.FollowFromRecord(evt.Did, evt.Rkey, record) + if err != nil { + i.l.Warn("ignoring invalid follow record", "err", err) + return nil + } + + if err := db.UpsertFollow(i.db, follow); err != nil { + return fmt.Errorf("upserting follow record") + } + case tapc.RecordDeleteAction: + if err := db.DeleteFollowByRkey( + i.db, + evt.Did.String(), + evt.Rkey.String(), + ); err != nil { + return fmt.Errorf("deleting record from db: %w", err) + } + } + return nil +} + +// TODO: let's just remove the knot.member record +func (i *Ingester) ingestKnotMember(_ context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.KnotMember + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + panic("unimplemented") + case tapc.RecordDeleteAction: + panic("unimplemented") + } + return nil +} + +func (i *Ingester) ingestKnot(ctx context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.Knot + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + domain := evt.Rkey.String() + + if err := db.AddKnot(i.db, domain, evt.Did.String()); err != nil { + return fmt.Errorf("upserting knot: %w", err) + } + + // TODO: hmmm should we run verification here? + // There can be unverified knot in user profile. + panic("unimplemented") + case tapc.RecordDeleteAction: + domain := evt.Rkey.String() + + // get record from db first + registration, err := func(domain string, did syntax.DID) (models.Registration, error) { + registrations, err := db.GetRegistrations( + i.db, + orm.FilterEq("domain", domain), + orm.FilterEq("did", evt.Did), + ) + if err != nil { + return models.Registration{}, err + } + if len(registrations) != 1 { + return models.Registration{}, fmt.Errorf("got incorret number of registrations: %d, expected 1", len(registrations)) + } + return registrations[0], nil + }(domain, evt.Did) + if err != nil { + return fmt.Errorf("getting registration: %w", err) + } + + tx, err := i.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("starting transaction: %w", err) + } + defer tx.Rollback() + // TODO: rollback enforcer + + if err := db.DeleteKnot( + tx, + orm.FilterEq("did", evt.Did), + orm.FilterEq("domain", domain), + ); err != nil { + return fmt.Errorf("deleting knot: %w", err) + } + + if registration.Registered != nil { + // TODO: clear from enforcer + panic("unimplemented") + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("commiting transaction: %w", err) + } + } + return nil +} + +func (i *Ingester) ingestLabelDefinition(_ context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.LabelDefinition + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + def, err := models.LabelDefinitionFromRecord(evt.Did.String(), evt.Rkey.String(), record) + if err != nil { + return fmt.Errorf("failed to parse labeldef from record: %w", err) + } + + if err := def.Validate(); err != nil { + i.l.Warn("ignoring invalid label def record", "err", err) + return nil + } + + if _, err := db.UpsertLabelDefinition(i.db, def); err != nil { + return fmt.Errorf("upserting label definition") + } + case tapc.RecordDeleteAction: + if err := db.DeleteLabelDefinition( + i.db, + orm.FilterEq("did", evt.Did), + orm.FilterEq("rkey", evt.Rkey), + ); err != nil { + return fmt.Errorf("deleting record from db: %w", err) + } + } + return nil +} + +// TODO: label.op record is not designed to be mutable. should be reimplemented +func (i *Ingester) ingestLabelOp(ctx context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction: + var record tangled.LabelOp + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + // TODO: + // 1. validate permissions + // 2. labelOp.Subject -> (subject).Repo + // 3. get all label definition for that repo, constructing actx + ops := models.LabelOpsFromRecord(evt.Did.String(), evt.Rkey.String(), record) + for _, o := range ops { + // 4. find label def based on o.OperandKey (AT-URI to the label definition) + // 5. validate labelOp from def + panic("unimplemented") + } + + tx, err := i.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("starting transaction: %w", err) + } + defer tx.Rollback() + + for _, o := range ops { + _, err := db.AddLabelOp(tx, &o) + if err != nil { + return fmt.Errorf("adding label op: %w", err) + } + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("commiting transaction: %w", err) + } + case tapc.RecordUpdateAction: + // no-op. we are ignoring update action for label.op records + case tapc.RecordDeleteAction: + // no-op. we are ignoring delete action for label.op records + } + return nil +} + +func (i *Ingester) ingestPublicKey(_ context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.PublicKey + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + pubKey, err := models.PublicKeyFromRecord(evt.Did, evt.Rkey, record) + if err != nil { + i.l.Warn("ignoring invalid publicKey record", "err", err) + return nil + } + if err := pubKey.Validate(); err != nil { + i.l.Warn("ignoring invalid publicKey record", "err", err) + return nil + } + + if err := db.UpsertPublicKey(i.db, pubKey); err != nil { + return fmt.Errorf("upserting publicKey record") + } + case tapc.RecordDeleteAction: + if err := db.DeletePublicKeyByRkey( + i.db, + evt.Did.String(), + evt.Rkey.String(), + ); err != nil { + return fmt.Errorf("deleting record from db: %w", err) + } + } + return nil +} + +// so this is one of the reasons why we need repo-DID +// and possibly its own PDS. +// for things need permission like repo artifact, even its not collaborative, +// we need to pass the Knot to check the authority. +// We cannot really distribute the arbitrary RBAC rules in meaningful way. +// That's pretty hard and fragile. +// Instead, when any data is belongs to the repo (no matter who created it), it should be owned +// by the repo. I mean the actual data should be. +// If someone removed the original artifact from their PDS, or if their PDS goes down, we lost +// a way to backfill relateed data to construct the repo view. + +func (i *Ingester) ingestRepoArtifact(_ context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.RepoArtifact + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + artifact, err := models.ArtifactFromRecord(evt.Did, evt.Rkey, record) + if err != nil { + i.l.Warn("ignoring invalid artifact record", "err", err) + return nil + } + + if err := db.UpsertArtifact(i.db, artifact); err != nil { + return fmt.Errorf("upserting artifact: %w", err) + } + case tapc.RecordDeleteAction: + if err := db.DeleteArtifact( + i.db, + orm.FilterEq("did", evt.Did), + orm.FilterEq("rkey", evt.Rkey), + ); err != nil { + return fmt.Errorf("deleting record from db: %w", err) + } + } + return nil +} + +func (i *Ingester) ingestRepoIssueComment(ctx context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.RepoIssueComment + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + comment, err := models.IssueCommentFromRecord(evt.Did.String(), evt.Rkey.String(), record) + if err != nil { + i.l.Warn("ignoring invalid issue.comment record", "err", err) + return nil + } + + tx, err := i.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("starting transaction: %w", err) + } + defer tx.Rollback() + + if _, err = db.UpsertIssueComment(tx, *comment); err != nil { + return fmt.Errorf("upserting issue comment: %w", err) + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("commiting transaction: %w", err) + } + case tapc.RecordDeleteAction: + if err := db.DeleteIssueComments( + i.db, + orm.FilterEq("did", evt.Did), + orm.FilterEq("rkey", evt.Rkey), + ); err != nil { + return fmt.Errorf("deleting issue comment: %w", err) + } + } + return nil +} + +func (i *Ingester) ingestRepoIssue(ctx context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.RepoIssue + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + issue := models.IssueFromRecord(evt.Did.String(), evt.Rkey.String(), record) + if err := issue.Validate(); err != nil { + i.l.Warn("ignoring invalid issue record", "err", err) + return nil + } + + tx, err := i.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("starting transaction: %w", err) + } + defer tx.Rollback() + + if err := db.PutIssue(tx, &issue); err != nil { + return fmt.Errorf("upserting issue: %w", err) + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("commiting issue upsert: %w", err) + } + + if evt.Action == tapc.RecordCreateAction { + i.notifier.NewIssue(ctx, &issue, issue.Mentions) + } + case tapc.RecordDeleteAction: + tx, err := i.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("starting transaction: %w", err) + } + defer tx.Rollback() + + if err := db.DeleteIssues( + tx, + evt.Did.String(), + evt.Rkey.String(), + ); err != nil { + return fmt.Errorf("deleting issue: %w", err) + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("commiting issue delete: %w", err) + } + } + return nil +} + +func (i *Ingester) ingestRepoIssueState(_ context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction: + var record tangled.RepoIssueState + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + // TODO: check permission + + switch record.State { + case tangled.RepoIssueStateOpen: + if err := db.ReopenIssues( + i.db, + orm.FilterEq("at_uri", record.Issue), + ); err != nil { + return fmt.Errorf("opening issue: %w", err) + } + case tangled.RepoIssueStateClosed: + if err := db.CloseIssues( + i.db, + orm.FilterEq("at_uri", record.Issue), + ); err != nil { + return fmt.Errorf("closing issue: %w", err) + } + default: + return nil + } + + // TODO: notify + case tapc.RecordUpdateAction: + // no-op + case tapc.RecordDeleteAction: + // no-op + } + return nil +} + +func (i *Ingester) ingestRepo(ctx context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.Repo + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + repo, err := models.RepoFromRecord(evt.Did, evt.Rkey, record) + if err != nil { + i.l.Warn("ignoring invalid repo record", "err", err) + return nil + } + + tx, err := i.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("starting transaction: %w", err) + } + defer tx.Rollback() + + if err := db.UpsertRepo(tx, &repo); err != nil { + return fmt.Errorf("upserting repo: %w", err) + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("commiting repo upsert: %w", err) + } + case tapc.RecordDeleteAction: + if err := db.RemoveRepo(i.db, evt.Did, evt.Rkey); err != nil { + return fmt.Errorf("deleting repo: %w", err) + } + } + return nil +} + +func (i *Ingester) ingestRepoPullComment(ctx context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.RepoPullComment + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + comment, err := models.PullCommentFromRecord(evt.Did, evt.Rkey, record) + if err != nil { + i.l.Warn("ignoring invalid issue.comment record", "err", err) + return nil + } + + tx, err := i.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("starting transaction: %w", err) + } + defer tx.Rollback() + + if err := db.UpsertPullComment(tx, &comment); err != nil { + return fmt.Errorf("upserting pull comment: %w", err) + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("commiting transaction: %w", err) + } + + if evt.Action == tapc.RecordCreateAction { + i.notifier.NewPullComment(ctx, &comment, comment.Mentions) + } + case tapc.RecordDeleteAction: + if err := db.DeletePullComments( + i.db, + orm.FilterEq("did", evt.Did), + orm.FilterEq("rkey", evt.Rkey), + ); err != nil { + return fmt.Errorf("deleting pull comment: %w", err) + } + } + return nil +} + +func (i *Ingester) ingestRepoPull(_ context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.RepoPull + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + panic("unimplemented") + case tapc.RecordDeleteAction: + panic("unimplemented") + } + return nil +} + +func (i *Ingester) ingestRepoPullStatus(_ context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction: + var record tangled.RepoPullStatus + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + panic("unimplemented") + case tapc.RecordUpdateAction: + // no-op + case tapc.RecordDeleteAction: + // no-op + } + return nil +} + +func (i *Ingester) ingestSpindleMember(_ context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.SpindleMember + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + // TODO: let's just remove the spindle.member record + panic("unimplemented") + case tapc.RecordDeleteAction: + panic("unimplemented") + } + return nil +} + +func (i *Ingester) ingestSpindle(ctx context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction: + var record tangled.Spindle + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + instance := evt.Rkey.String() + + if err := db.AddSpindle(i.db, models.Spindle{ + Owner: evt.Did, + Instance: instance, + }); err != nil { + return fmt.Errorf("adding spindle: %w", err) + } + + panic("unimplemented") + case tapc.RecordDeleteAction: + instance := evt.Rkey.String() + + // get record from db first + spindle, err := func(instance string, did syntax.DID) (models.Spindle, error) { + spindles, err := db.GetSpindles( + i.db, + orm.FilterEq("owner", did), + orm.FilterEq("instance", instance), + ) + if err != nil { + return models.Spindle{}, err + } + if len(spindles) != 1 { + return models.Spindle{}, fmt.Errorf("got incorret number of spindles: %d, expected 1", len(spindles)) + } + return spindles[0], nil + }(instance, evt.Did) + if err != nil { + return fmt.Errorf("getting spindle: %w", err) + } + + tx, err := i.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("starting transaction: %w", err) + } + defer tx.Rollback() + // TODO: rollback enforcer + + // remove spindle members first + if err := db.RemoveSpindleMember( + tx, + orm.FilterEq("owner", evt.Did), + orm.FilterEq("instance", instance), + ); err != nil { + return fmt.Errorf("deleting spindle members: %w", err) + } + if err := db.DeleteSpindle( + tx, + orm.FilterEq("owner", evt.Did), + orm.FilterEq("instance", instance), + ); err != nil { + return fmt.Errorf("deleting spindle: %w", err) + } + + if spindle.Verified != nil { + // TODO: clear from enforcer + panic("unimplemented") + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("commiting transaction: %w", err) + } + } + return nil +} + +func (i *Ingester) ingestString(_ context.Context, evt *tapc.RecordEventData) error { + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + var record tangled.String + if err := json.Unmarshal(evt.Record, &record); err != nil { + return fmt.Errorf("parsing record json: %w", err) + } + + string := models.StringFromRecord(evt.Did.String(), evt.Rkey.String(), record) + if err := string.Validate(); err != nil { + i.l.Warn("invalid record", "err", err) + return nil + } + + if err := db.UpsertString(i.db, string); err != nil { + return fmt.Errorf("upserting string: %w", err) + } + + return nil + case tapc.RecordDeleteAction: + if err := db.DeleteString( + i.db, + orm.FilterEq("did", evt.Did), + orm.FilterEq("rkey", evt.Rkey), + ); err != nil { + return fmt.Errorf("deleting string: %w", err) + } + + return nil + } + return nil +} diff --git a/appview/models/artifact.go b/appview/models/artifact.go index e3c14121..5b75d316 100644 --- a/appview/models/artifact.go +++ b/appview/models/artifact.go @@ -25,6 +25,30 @@ type Artifact struct { MimeType string } -func (a *Artifact) ArtifactAt() syntax.ATURI { +func (a *Artifact) AtUri() syntax.ATURI { return syntax.ATURI(fmt.Sprintf("at://%s/%s/%s", a.Did, tangled.RepoArtifactNSID, a.Rkey)) } + +func ArtifactFromRecord(did syntax.DID, rkey syntax.RecordKey, record tangled.RepoArtifact) (Artifact, error) { + // validate atproto record + repoAt, err := syntax.ParseATURI(record.Repo) + if err != nil { + return Artifact{}, fmt.Errorf("invalid record %T: %w", record, fmt.Errorf("repo should be valid at-uri: %w", err)) + } + created, err := time.Parse(time.RFC3339, record.CreatedAt) + if err != nil { + return Artifact{}, fmt.Errorf("invalid record %T: %w", record, fmt.Errorf("invalid time format '%s'", record.CreatedAt)) + } + + return Artifact{ + Did: did.String(), + Rkey: rkey.String(), + RepoAt: repoAt, + Tag: plumbing.Hash(record.Tag), + CreatedAt: created, + BlobCid: cid.Cid(record.Artifact.Ref), + Name: record.Name, + Size: uint64(record.Artifact.Size), + MimeType: record.Artifact.MimeType, + }, nil +} diff --git a/appview/models/follow.go b/appview/models/follow.go index d371226a..2d8375f4 100644 --- a/appview/models/follow.go +++ b/appview/models/follow.go @@ -1,8 +1,10 @@ package models import ( + "fmt" "time" + "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/api/tangled" ) @@ -20,6 +22,25 @@ func (f *Follow) AsRecord() tangled.GraphFollow { } } +func FollowFromRecord(did syntax.DID, rkey syntax.RecordKey, record tangled.GraphFollow) (Follow, error) { + subjectDid, err := syntax.ParseDID(record.Subject) + if err != nil { + return Follow{}, fmt.Errorf("subject should be valid did: %w", err) + } + + created, err := time.Parse(time.RFC3339, record.CreatedAt) + if err != nil { + return Follow{}, fmt.Errorf("invalid time format '%s'", record.CreatedAt) + } + + return Follow{ + UserDid: did.String(), + Rkey: rkey.String(), + SubjectDid: subjectDid.String(), + FollowedAt: created, + }, nil +} + type FollowStats struct { Followers int64 Following int64 diff --git a/appview/models/profile.go b/appview/models/profile.go index 83dc2607..7c931c10 100644 --- a/appview/models/profile.go +++ b/appview/models/profile.go @@ -50,6 +50,90 @@ func (p Profile) IsPinnedReposEmpty() bool { return true } +// ProfileFromRecord will validate the atproto record and convert it to [Profile]. +// It can return error for invalid records. +func ProfileFromRecord(did syntax.DID, record tangled.ActorProfile) (Profile, error) { + // validate atproto record + if err := func(record tangled.ActorProfile) error { + // validate description + if record.Description != nil && len(*record.Description) > 256 { + return fmt.Errorf("bio is too long") + } + + // validate links + if len(record.Links) > 5 { + return fmt.Errorf("links cannot be more than 5") + } + + // validate location + if record.Location != nil && len(*record.Location) > 256 { + return fmt.Errorf("location is too long") + } + + // validate pinnedRepositories + if len(record.PinnedRepositories) >= 5 { + return fmt.Errorf("pinnedRepositories cannot be more than 6") + } + for i, v := range record.PinnedRepositories { + if _, err := syntax.ParseATURI(v); err != nil { + return fmt.Errorf("invalid at-uri at pinnedRepositories[%d]: %w", i, err) + } + } + + // validate pronouns + if record.Pronouns != nil && len(*record.Pronouns) > 40 { + return fmt.Errorf("pronouns are too long") + } + + // validate stats + if len(record.Stats) > 2 { + return fmt.Errorf("stats cannot be more than 2") + } + for i, v := range record.Stats { + if VanityStatKind(v).String() == "" { + return fmt.Errorf("unknown stat kind '%s' at stats[%d]", v, i) + } + } + return nil + }(record); err != nil { + return Profile{}, fmt.Errorf("invalid record %T: %w", record, err) + } + + p := Profile{Did: did.String()} + + if record.Description != nil { + p.Description = *record.Description + } + + p.IncludeBluesky = record.Bluesky + + if record.Location != nil { + p.Location = *record.Location + } + + copy(p.Links[:], record.Links) + + for i, s := range record.Stats { + if i >= 2 { + break + } + p.Stats[i].Kind = VanityStatKind(s) + } + + for i, r := range record.PinnedRepositories { + if i >= 6 { + break + } + p.PinnedRepos[i] = syntax.ATURI(r) + } + + if record.Pronouns != nil { + p.Pronouns = *record.Pronouns + } + + return p, nil +} + type VanityStatKind string const ( diff --git a/appview/models/pull.go b/appview/models/pull.go index 4a2e0ee7..6bb5b64b 100644 --- a/appview/models/pull.go +++ b/appview/models/pull.go @@ -171,6 +171,10 @@ func (p *PullComment) AtUri() syntax.ATURI { return syntax.ATURI(p.CommentAt) } +func PullCommentFromRecord(did syntax.DID, rkey syntax.RecordKey, record tangled.RepoPullComment) (PullComment, error) { + panic("unimplemented") +} + func (p *Pull) TotalComments() int { total := 0 for _, s := range p.Submissions { diff --git a/appview/models/reaction.go b/appview/models/reaction.go index 748800f4..8ab0807f 100644 --- a/appview/models/reaction.go +++ b/appview/models/reaction.go @@ -1,6 +1,7 @@ package models import ( + "fmt" "time" "github.com/bluesky-social/indigo/atproto/syntax" @@ -65,6 +66,31 @@ func (r *Reaction) AsRecord() tangled.FeedReaction { } } +func ReactionFromRecord(did syntax.DID, rkey syntax.RecordKey, record tangled.FeedReaction) (Reaction, error) { + subjectAt, err := syntax.ParseATURI(record.Subject) + if err != nil { + return Reaction{}, fmt.Errorf("subject should be valid at-uri: %w", err) + } + + kind, ok := ParseReactionKind(record.Reaction) + if !ok { + return Reaction{}, fmt.Errorf("invalid reaction kind '%s'", record.Reaction) + } + + created, err := time.Parse(time.RFC3339, record.CreatedAt) + if err != nil { + return Reaction{}, fmt.Errorf("invalid time format '%s'", record.CreatedAt) + } + + return Reaction{ + ReactedByDid: did.String(), + Rkey: rkey.String(), + ThreadAt: subjectAt, + Kind: kind, + Created: created, + }, nil +} + type ReactionDisplayData struct { Count int Users []string diff --git a/appview/models/repo.go b/appview/models/repo.go index 4a7d96b0..da84253b 100644 --- a/appview/models/repo.go +++ b/appview/models/repo.go @@ -30,6 +30,10 @@ type Repo struct { Source string } +func RepoFromRecord(did syntax.DID, rkey syntax.RecordKey, record tangled.Repo) (Repo, error) { + panic("unimplemented") +} + func (r *Repo) AsRecord() tangled.Repo { var source, spindle, description, website *string diff --git a/appview/models/star.go b/appview/models/star.go index b0c140d1..014b74c3 100644 --- a/appview/models/star.go +++ b/appview/models/star.go @@ -1,6 +1,7 @@ package models import ( + "fmt" "time" "github.com/bluesky-social/indigo/atproto/syntax" @@ -21,6 +22,25 @@ func (s *Star) AsRecord() tangled.FeedStar { } } +func StarFromRecord(did syntax.DID, rkey syntax.RecordKey, record tangled.FeedStar) (Star, error) { + subjectAt, err := syntax.ParseATURI(record.Subject) + if err != nil { + return Star{}, fmt.Errorf("subject should be valid at-uri: %w", err) + } + + created, err := time.Parse(time.RFC3339, record.CreatedAt) + if err != nil { + return Star{}, fmt.Errorf("invalid time format '%s'", record.CreatedAt) + } + + return Star{ + Did: did.String(), + Rkey: rkey.String(), + RepoAt: subjectAt, + Created: created, + }, nil +} + // RepoStar is used for reverse mapping to repos type RepoStar struct { Star diff --git a/appview/repo/repo.go b/appview/repo/repo.go index 6fe5dc5d..9a49156f 100644 --- a/appview/repo/repo.go +++ b/appview/repo/repo.go @@ -912,7 +912,7 @@ func (rp *Repo) DeleteRepo(w http.ResponseWriter, r *http.Request) { } // remove repo from db - err = db.RemoveRepo(tx, f.Did, f.Name) + err = db.RemoveRepo(tx, syntax.DID(f.Did), syntax.RecordKey(f.Rkey)) if err != nil { rp.pages.Notice(w, noticeId, "Failed to update appview") return -- 2.51.2