diff --git a/appview/db/jetstream.go b/appview/db/jetstream.go index cdfbcda5..91fe8cb8 100644 --- a/appview/db/jetstream.go +++ b/appview/db/jetstream.go @@ -1,10 +1,6 @@ package db -type DbWrapper struct { - Execer -} - -func (db DbWrapper) SaveLastTimeUs(lastTimeUs int64) error { +func (db *DB) SaveLastTimeUs(lastTimeUs int64) error { _, err := db.Exec(` insert into _jetstream (id, last_time_us) values (1, ?) @@ -13,7 +9,7 @@ func (db DbWrapper) SaveLastTimeUs(lastTimeUs int64) error { return err } -func (db DbWrapper) GetLastTimeUs() (int64, error) { +func (db *DB) GetLastTimeUs() (int64, error) { var lastTimeUs int64 row := db.QueryRow(`select last_time_us from _jetstream where id = 1;`) err := row.Scan(&lastTimeUs) diff --git a/appview/ingester.go b/appview/ingester.go index 5be37ab3..ecf666d7 100644 --- a/appview/ingester.go +++ b/appview/ingester.go @@ -35,7 +35,7 @@ import ( ) type Ingester struct { - Db db.DbWrapper + Db *db.DB Enforcer *rbac.Enforcer IdResolver *idresolver.Resolver Cache *cache.Cache @@ -494,12 +494,7 @@ func (i *Ingester) ingestProfile(ctx context.Context, e *jmodels.Event) error { PreferredHandle: preferredHandle, } - ddb, ok := i.Db.Execer.(*db.DB) - if !ok { - return fmt.Errorf("failed to index profile record, invalid db cast") - } - - tx, err := ddb.Begin() + tx, err := i.Db.Begin() if err != nil { return fmt.Errorf("failed to start transaction") } @@ -566,12 +561,7 @@ func (i *Ingester) ingestSpindleMember(ctx context.Context, e *jmodels.Event) er return err } - ddb, ok := i.Db.Execer.(*db.DB) - if !ok { - return fmt.Errorf("invalid db cast") - } - - err = db.AddSpindleMember(ddb, models.SpindleMember{ + err = db.AddSpindleMember(i.Db, models.SpindleMember{ Did: syntax.DID(did), Rkey: e.Commit.RKey, Instance: record.Instance, @@ -590,14 +580,9 @@ func (i *Ingester) ingestSpindleMember(ctx context.Context, e *jmodels.Event) er case jmodels.CommitOperationDelete: rkey := e.Commit.RKey - ddb, ok := i.Db.Execer.(*db.DB) - if !ok { - return fmt.Errorf("failed to index profile record, invalid db cast") - } - // get record from db first members, err := db.GetSpindleMembers( - ddb, + i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey), ) @@ -606,7 +591,7 @@ func (i *Ingester) ingestSpindleMember(ctx context.Context, e *jmodels.Event) er } member := members[0] - tx, err := ddb.Begin() + tx, err := i.Db.Begin() if err != nil { return fmt.Errorf("failed to start txn: %w", err) } @@ -659,12 +644,7 @@ func (i *Ingester) ingestSpindle(ctx context.Context, e *jmodels.Event) error { instance := e.Commit.RKey - ddb, ok := i.Db.Execer.(*db.DB) - if !ok { - return fmt.Errorf("failed to index profile record, invalid db cast") - } - - err := db.AddSpindle(ddb, models.Spindle{ + err := db.AddSpindle(i.Db, models.Spindle{ Owner: syntax.DID(did), Instance: instance, }) @@ -683,7 +663,7 @@ func (i *Ingester) ingestSpindle(ctx context.Context, e *jmodels.Event) error { return err } - _, err = serververify.MarkSpindleVerified(ddb, i.Enforcer, instance, did) + _, err = serververify.MarkSpindleVerified(i.Db, i.Enforcer, instance, did) if err != nil { return fmt.Errorf("failed to mark verified: %w", err) } @@ -693,15 +673,10 @@ func (i *Ingester) ingestSpindle(ctx context.Context, e *jmodels.Event) error { case jmodels.CommitOperationDelete: instance := e.Commit.RKey - ddb, ok := i.Db.Execer.(*db.DB) - if !ok { - return fmt.Errorf("failed to index profile record, invalid db cast") - } - // get record from db first spindles, err := db.GetSpindles( ctx, - ddb, + i.Db, orm.FilterEq("owner", did), orm.FilterEq("instance", instance), ) @@ -710,7 +685,7 @@ func (i *Ingester) ingestSpindle(ctx context.Context, e *jmodels.Event) error { } spindle := spindles[0] - tx, err := ddb.Begin() + tx, err := i.Db.Begin() if err != nil { return err } @@ -768,11 +743,6 @@ func (i *Ingester) ingestString(e *jmodels.Event) error { l := i.Logger.With("handler", "ingestString", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) l.Info("ingesting record") - ddb, ok := i.Db.Execer.(*db.DB) - if !ok { - return fmt.Errorf("failed to index string record, invalid db cast") - } - switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) @@ -790,7 +760,7 @@ func (i *Ingester) ingestString(e *jmodels.Event) error { return err } - if err = db.AddString(ddb, string); err != nil { + if err = db.AddString(i.Db, string); err != nil { l.Error("failed to add string", "err", err) return err } @@ -799,7 +769,7 @@ func (i *Ingester) ingestString(e *jmodels.Event) error { case jmodels.CommitOperationDelete: if err := db.DeleteString( - ddb, + i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey), ); err != nil { @@ -884,12 +854,7 @@ func (i *Ingester) ingestKnot(e *jmodels.Event) error { domain := e.Commit.RKey - ddb, ok := i.Db.Execer.(*db.DB) - if !ok { - return fmt.Errorf("failed to index profile record, invalid db cast") - } - - err := db.AddKnot(ddb, domain, did) + err := db.AddKnot(i.Db, domain, did) if err != nil { l.Error("failed to add knot to db", "err", err, "domain", domain) return err @@ -907,7 +872,7 @@ func (i *Ingester) ingestKnot(e *jmodels.Event) error { return err } - err = serververify.MarkKnotVerified(ddb, i.Enforcer, domain, did) + err = serververify.MarkKnotVerified(i.Db, i.Enforcer, domain, did) if err != nil { return fmt.Errorf("failed to mark verified: %w", err) } @@ -917,14 +882,9 @@ func (i *Ingester) ingestKnot(e *jmodels.Event) error { case jmodels.CommitOperationDelete: domain := e.Commit.RKey - ddb, ok := i.Db.Execer.(*db.DB) - if !ok { - return fmt.Errorf("failed to index knot record, invalid db cast") - } - // get record from db first registrations, err := db.GetRegistrations( - ddb, + i.Db, orm.FilterEq("domain", domain), orm.FilterEq("did", did), ) @@ -936,7 +896,7 @@ func (i *Ingester) ingestKnot(e *jmodels.Event) error { } registration := registrations[0] - tx, err := ddb.Begin() + tx, err := i.Db.Begin() if err != nil { return err } @@ -983,11 +943,6 @@ func (i *Ingester) ingestIssue(ctx context.Context, e *jmodels.Event) error { l := i.Logger.With("handler", "ingestIssue", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) l.Info("ingesting record") - ddb, ok := i.Db.Execer.(*db.DB) - if !ok { - return fmt.Errorf("failed to index issue record, invalid db cast") - } - switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) @@ -1008,16 +963,7 @@ func (i *Ingester) ingestIssue(ctx context.Context, e *jmodels.Event) error { return fmt.Errorf("failed to validate issue: %w", err) } - if record.Repo != nil { - repo, repoErr := db.GetRepoByAtUri(i.Db, *record.Repo) - if repoErr == nil && repo.RepoDid != "" { - if enqErr := db.EnqueuePdsRewrite(i.Db, did, repo.RepoDid, tangled.RepoIssueNSID, rkey, *record.Repo); enqErr != nil { - l.Warn("failed to enqueue PDS rewrite for issue", "err", enqErr, "did", did, "repoDid", repo.RepoDid) - } - } - } - - tx, err := ddb.BeginTx(ctx, nil) + tx, err := i.Db.BeginTx(ctx, nil) if err != nil { l.Error("failed to begin transaction", "err", err) return err @@ -1039,7 +985,7 @@ func (i *Ingester) ingestIssue(ctx context.Context, e *jmodels.Event) error { return nil case jmodels.CommitOperationDelete: - tx, err := ddb.BeginTx(ctx, nil) + tx, err := i.Db.BeginTx(ctx, nil) if err != nil { l.Error("failed to begin transaction", "err", err) return err @@ -1074,11 +1020,6 @@ func (i *Ingester) ingestPull(ctx context.Context, e *jmodels.Event) error { l := i.Logger.With("handler", "ingestPull", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) l.Info("ingesting record") - ddb, ok := i.Db.Execer.(*db.DB) - if !ok { - return fmt.Errorf("failed to index pull record, invalid db cast") - } - switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) @@ -1161,7 +1102,7 @@ func (i *Ingester) ingestPull(ctx context.Context, e *jmodels.Event) error { return fmt.Errorf("failed to validate pull: %w", err) } - tx, err := ddb.BeginTx(ctx, nil) + tx, err := i.Db.BeginTx(ctx, nil) if err != nil { l.Error("failed to begin transaction", "err", err) return err @@ -1183,7 +1124,7 @@ func (i *Ingester) ingestPull(ctx context.Context, e *jmodels.Event) error { return nil case jmodels.CommitOperationDelete: - tx, err := ddb.BeginTx(ctx, nil) + tx, err := i.Db.BeginTx(ctx, nil) if err != nil { l.Error("failed to begin transaction", "err", err) return err @@ -1218,11 +1159,6 @@ func (i *Ingester) ingestIssueComment(e *jmodels.Event) error { l := i.Logger.With("handler", "ingestIssueComment", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) l.Info("ingesting record") - ddb, ok := i.Db.Execer.(*db.DB) - if !ok { - return fmt.Errorf("failed to index issue comment record, invalid db cast") - } - switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) @@ -1241,7 +1177,7 @@ func (i *Ingester) ingestIssueComment(e *jmodels.Event) error { return fmt.Errorf("failed to validate comment: %w", err) } - tx, err := ddb.Begin() + tx, err := i.Db.Begin() if err != nil { return fmt.Errorf("failed to start transaction: %w", err) } @@ -1256,7 +1192,7 @@ func (i *Ingester) ingestIssueComment(e *jmodels.Event) error { case jmodels.CommitOperationDelete: if err := db.DeleteIssueComments( - ddb, + i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey), ); err != nil { @@ -1278,11 +1214,6 @@ func (i *Ingester) ingestLabelDefinition(e *jmodels.Event) error { l := i.Logger.With("handler", "ingestLabelDefinition", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) l.Info("ingesting record") - ddb, ok := i.Db.Execer.(*db.DB) - if !ok { - return fmt.Errorf("failed to index label definition, invalid db cast") - } - switch e.Commit.Operation { case jmodels.CommitOperationCreate, jmodels.CommitOperationUpdate: raw := json.RawMessage(e.Commit.Record) @@ -1301,7 +1232,7 @@ func (i *Ingester) ingestLabelDefinition(e *jmodels.Event) error { return fmt.Errorf("failed to validate labeldef: %w", err) } - _, err = db.AddLabelDefinition(ddb, def) + _, err = db.AddLabelDefinition(i.Db, def) if err != nil { return fmt.Errorf("failed to create labeldef: %w", err) } @@ -1310,7 +1241,7 @@ func (i *Ingester) ingestLabelDefinition(e *jmodels.Event) error { case jmodels.CommitOperationDelete: if err := db.DeleteLabelDefinition( - ddb, + i.Db, orm.FilterEq("did", did), orm.FilterEq("rkey", rkey), ); err != nil { @@ -1332,11 +1263,6 @@ func (i *Ingester) ingestLabelOp(e *jmodels.Event) error { l := i.Logger.With("handler", "ingestLabelOp", "nsid", e.Commit.Collection, "did", did, "rkey", rkey) l.Info("ingesting record") - ddb, ok := i.Db.Execer.(*db.DB) - if !ok { - return fmt.Errorf("failed to index label op, invalid db cast") - } - switch e.Commit.Operation { case jmodels.CommitOperationCreate: raw := json.RawMessage(e.Commit.Record) @@ -1352,7 +1278,7 @@ func (i *Ingester) ingestLabelOp(e *jmodels.Event) error { var repo *models.Repo switch collection { case tangled.RepoIssueNSID: - i, err := db.GetIssues(ddb, orm.FilterEq("at_uri", subject)) + i, err := db.GetIssues(i.Db, orm.FilterEq("at_uri", subject)) if err != nil || len(i) != 1 { return fmt.Errorf("failed to find subject: %w || subject count %d", err, len(i)) } @@ -1361,7 +1287,7 @@ func (i *Ingester) ingestLabelOp(e *jmodels.Event) error { return fmt.Errorf("unsupported label subject: %s", collection) } - actx, err := db.NewLabelApplicationCtx(ddb, orm.FilterIn("at_uri", repo.Labels)) + actx, err := db.NewLabelApplicationCtx(i.Db, orm.FilterIn("at_uri", repo.Labels)) if err != nil { return fmt.Errorf("failed to build label application ctx: %w", err) } @@ -1378,7 +1304,7 @@ func (i *Ingester) ingestLabelOp(e *jmodels.Event) error { } } - tx, err := ddb.Begin() + tx, err := i.Db.Begin() if err != nil { return err } diff --git a/appview/state/state.go b/appview/state/state.go index fc18104f..1138027e 100644 --- a/appview/state/state.go +++ b/appview/state/state.go @@ -118,7 +118,6 @@ func Make(ctx context.Context, config *config.Config) (*State, error) { mentionsResolver := mentions.New(config, res, d, log.SubLogger(logger, "mentionsResolver")) - wrapper := db.DbWrapper{Execer: d} jc, err := jetstream.NewJetstreamClient( config.Jetstream.Endpoint, "appview", @@ -142,7 +141,7 @@ func Make(ctx context.Context, config *config.Config) (*State, error) { }, nil, tlog.SubLogger(logger, "jetstream"), - wrapper, + d, false, // in-memory filter is inapplicable to appview so