diff --git a/appview/db/repos.go b/appview/db/repos.go index 88ec4c13..1cbd78ba 100644 --- a/appview/db/repos.go +++ b/appview/db/repos.go @@ -409,8 +409,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, rkey string) error { + _, err := e.Exec(`delete from repos where did = ? and rkey = ?`, did, rkey) return err } diff --git a/appview/ingester.go b/appview/ingester.go index 7be25938..2d3a3fd2 100644 --- a/appview/ingester.go +++ b/appview/ingester.go @@ -83,6 +83,8 @@ func (i *Ingester) Ingest() processFunc { err = i.ingestLabelDefinition(e) case tangled.LabelOpNSID: err = i.ingestLabelOp(e) + case tangled.RepoNSID: + err = i.ingestRepo(ctx, e) } l = i.Logger.With("nsid", e.Commit.Collection) } @@ -95,6 +97,52 @@ func (i *Ingester) Ingest() processFunc { } } +func (i *Ingester) ingestRepo(ctx context.Context, e *jmodels.Event) error { + l := i.Logger.With("handler", "ingestStar") + l = l.With("nsid", e.Commit.Collection) + + switch e.Commit.Operation { + case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: + record := tangled.Repo{} + if err := json.Unmarshal(e.Commit.Record, &record); err != nil { + l.Error("invalid record", "err", err) + return err + } + + repo := &models.Repo{ + Did: e.Did, + Rkey: e.Commit.RKey, + Name: record.Name, + Knot: record.Knot, + Description: "", + } + ddb, _ := i.Db.Execer.(*db.DB) + + tx, err := ddb.Begin() + if err != nil { + return fmt.Errorf("failed to start transaction") + } + defer tx.Rollback() + // tx, err := i.Db.SaveLastTimeUs + + if err := db.AddRepo(tx, repo); err != nil { + return fmt.Errorf("adding repo: %w", err) + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("commiting repo add transaction: %w", err) + } + case jmodels.CommitOperationDelete: + if err := db.RemoveRepo(i.Db, e.Did, e.Commit.RKey); err != nil { + return fmt.Errorf("deleting repo: %w", err) + } + } + + l.Info("processed repo", "operation", e.Commit.Operation, "rkey", e.Commit.RKey) + + return nil +} + func (i *Ingester) ingestStar(e *jmodels.Event) error { var err error did := e.Did @@ -361,7 +409,7 @@ func (i *Ingester) ingestProfile(e *jmodels.Event) error { err = db.ValidateProfile(tx, &profile) if err != nil { - return fmt.Errorf("invalid profile record") + return fmt.Errorf("invalid profile record: %w", err) } err = db.UpsertProfile(tx, &profile) diff --git a/appview/repo/repo.go b/appview/repo/repo.go index 3568b3f2..f899151d 100644 --- a/appview/repo/repo.go +++ b/appview/repo/repo.go @@ -906,7 +906,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, f.Did, f.Rkey) if err != nil { rp.pages.Notice(w, noticeId, "Failed to update appview") return diff --git a/appview/state/state.go b/appview/state/state.go index 150a2010..00cbaca0 100644 --- a/appview/state/state.go +++ b/appview/state/state.go @@ -117,6 +117,7 @@ func Make(ctx context.Context, config *config.Config) (*State, error) { tangled.RepoIssueCommentNSID, tangled.LabelDefinitionNSID, tangled.LabelOpNSID, + tangled.RepoNSID, }, nil, tlog.SubLogger(logger, "jetstream"),