diff --git a/bgs/bgs.go b/bgs/bgs.go index 50bec20a..d1acdaa4 100644 --- a/bgs/bgs.go +++ b/bgs/bgs.go @@ -522,8 +522,11 @@ func (bgs *BGS) handleFedEvent(ctx context.Context, host *models.PDS, env *event ctx, span := otel.Tracer("bgs").Start(ctx, "handleFedEvent") defer span.End() + eventsReceivedCounter.WithLabelValues(host.Host).Add(1) + switch { case env.RepoCommit != nil: + repoCommitsReceivedCounter.WithLabelValues(host.Host).Add(1) evt := env.RepoCommit log.Infow("bgs got repo append event", "seq", evt.Seq, "host", host.Host, "repo", evt.Repo) u, err := bgs.lookupUserByDid(ctx, evt.Repo) @@ -549,6 +552,7 @@ func (bgs *BGS) handleFedEvent(ctx context.Context, host *models.PDS, env *event // skip the fast path for rebases or if the user is already in the slow path if evt.Rebase || bgs.Index.Crawler.RepoInSlowPath(ctx, host, u.ID) { + rebasesCounter.WithLabelValues(host.Host).Add(1) ai, err := bgs.Index.LookupUser(ctx, u.ID) if err != nil { return err diff --git a/bgs/metrics.go b/bgs/metrics.go new file mode 100644 index 00000000..9f9853de --- /dev/null +++ b/bgs/metrics.go @@ -0,0 +1,21 @@ +package bgs + +import ( + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promauto" +) + +var eventsReceivedCounter = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "events_received_counter", + Help: "The total number of events received", +}, []string{"pds"}) + +var repoCommitsReceivedCounter = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "repo_commits_received_counter", + Help: "The total number of events received", +}, []string{"pds"}) + +var rebasesCounter = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "event_rebases", + Help: "The total number of rebase events received", +}, []string{"pds"}) -- 2.51.2 From 71f33361c5b594e0afddd9a0fad7b7d4749dc208 Mon Sep 17 00:00:00 2001 From: Jaz Volpert Date: Wed, 12 Jul 2023 22:34:21 +0000 Subject: [PATCH 2/7] Move bigsky util package to util/cliutil since lots of stuff uses it --- cmd/bigsky/main.go | 2 +- cmd/fakermaker/main.go | 2 +- cmd/gosky/admin.go | 2 +- cmd/gosky/debug.go | 2 +- cmd/gosky/main.go | 2 +- cmd/labelmaker/main.go | 2 +- cmd/laputa/main.go | 2 +- cmd/palomar/main.go | 2 +- cmd/stress/main.go | 2 +- pds/handlers_test.go | 2 +- {cmd/gosky/util => util/cliutil}/key.go | 0 {cmd/gosky/util => util/cliutil}/key_test.go | 0 {cmd/gosky/util => util/cliutil}/util.go | 0 13 files changed, 10 insertions(+), 10 deletions(-) rename {cmd/gosky/util => util/cliutil}/key.go (100%) rename {cmd/gosky/util => util/cliutil}/key_test.go (100%) rename {cmd/gosky/util => util/cliutil}/util.go (100%) diff --git a/cmd/bigsky/main.go b/cmd/bigsky/main.go index d6222076..c4088b6c 100644 --- a/cmd/bigsky/main.go +++ b/cmd/bigsky/main.go @@ -11,12 +11,12 @@ import ( "github.com/bluesky-social/indigo/bgs" "github.com/bluesky-social/indigo/blobs" "github.com/bluesky-social/indigo/carstore" - cliutil "github.com/bluesky-social/indigo/cmd/gosky/util" "github.com/bluesky-social/indigo/events" "github.com/bluesky-social/indigo/indexer" "github.com/bluesky-social/indigo/notifs" "github.com/bluesky-social/indigo/plc" "github.com/bluesky-social/indigo/repomgr" + "github.com/bluesky-social/indigo/util/cliutil" "github.com/bluesky-social/indigo/version" "github.com/bluesky-social/indigo/xrpc" diff --git a/cmd/fakermaker/main.go b/cmd/fakermaker/main.go index 07ae8129..eb2ec43e 100644 --- a/cmd/fakermaker/main.go +++ b/cmd/fakermaker/main.go @@ -10,8 +10,8 @@ import ( "os" "runtime" - cliutil "github.com/bluesky-social/indigo/cmd/gosky/util" "github.com/bluesky-social/indigo/fakedata" + "github.com/bluesky-social/indigo/util/cliutil" "github.com/bluesky-social/indigo/version" _ "github.com/joho/godotenv/autoload" diff --git a/cmd/gosky/admin.go b/cmd/gosky/admin.go index d5f72778..c871011f 100644 --- a/cmd/gosky/admin.go +++ b/cmd/gosky/admin.go @@ -12,7 +12,7 @@ import ( "github.com/bluesky-social/indigo/api" "github.com/bluesky-social/indigo/api/atproto" - cliutil "github.com/bluesky-social/indigo/cmd/gosky/util" + "github.com/bluesky-social/indigo/util/cliutil" cli "github.com/urfave/cli/v2" ) diff --git a/cmd/gosky/debug.go b/cmd/gosky/debug.go index e1a8e7bb..ba1f7de8 100644 --- a/cmd/gosky/debug.go +++ b/cmd/gosky/debug.go @@ -16,13 +16,13 @@ import ( "github.com/bluesky-social/indigo/api/atproto" comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/api/bsky" - cliutil "github.com/bluesky-social/indigo/cmd/gosky/util" "github.com/bluesky-social/indigo/did" "github.com/bluesky-social/indigo/events" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/repomgr" "github.com/bluesky-social/indigo/util" + "github.com/bluesky-social/indigo/util/cliutil" "github.com/bluesky-social/indigo/xrpc" "github.com/gorilla/websocket" diff --git a/cmd/gosky/main.go b/cmd/gosky/main.go index 192ce36f..56c2dd69 100644 --- a/cmd/gosky/main.go +++ b/cmd/gosky/main.go @@ -19,11 +19,11 @@ import ( comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/api/bsky" appbsky "github.com/bluesky-social/indigo/api/bsky" - cliutil "github.com/bluesky-social/indigo/cmd/gosky/util" "github.com/bluesky-social/indigo/events" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/util" + "github.com/bluesky-social/indigo/util/cliutil" "github.com/bluesky-social/indigo/version" "github.com/bluesky-social/indigo/xrpc" lru "github.com/hashicorp/golang-lru" diff --git a/cmd/labelmaker/main.go b/cmd/labelmaker/main.go index 6e417c6a..c8bc9e54 100644 --- a/cmd/labelmaker/main.go +++ b/cmd/labelmaker/main.go @@ -6,8 +6,8 @@ import ( "path/filepath" "github.com/bluesky-social/indigo/carstore" - cliutil "github.com/bluesky-social/indigo/cmd/gosky/util" "github.com/bluesky-social/indigo/labeler" + "github.com/bluesky-social/indigo/util/cliutil" "github.com/bluesky-social/indigo/version" "github.com/urfave/cli/v2" diff --git a/cmd/laputa/main.go b/cmd/laputa/main.go index 447708db..afdc00d4 100644 --- a/cmd/laputa/main.go +++ b/cmd/laputa/main.go @@ -6,9 +6,9 @@ import ( "github.com/bluesky-social/indigo/api" "github.com/bluesky-social/indigo/carstore" - cliutil "github.com/bluesky-social/indigo/cmd/gosky/util" "github.com/bluesky-social/indigo/pds" "github.com/bluesky-social/indigo/plc" + "github.com/bluesky-social/indigo/util/cliutil" "github.com/bluesky-social/indigo/version" _ "github.com/joho/godotenv/autoload" diff --git a/cmd/palomar/main.go b/cmd/palomar/main.go index 55021d93..e5873057 100644 --- a/cmd/palomar/main.go +++ b/cmd/palomar/main.go @@ -10,8 +10,8 @@ import ( _ "github.com/joho/godotenv/autoload" - cliutil "github.com/bluesky-social/indigo/cmd/gosky/util" "github.com/bluesky-social/indigo/search" + "github.com/bluesky-social/indigo/util/cliutil" "github.com/bluesky-social/indigo/version" logging "github.com/ipfs/go-log" diff --git a/cmd/stress/main.go b/cmd/stress/main.go index 5982f963..1d147cad 100644 --- a/cmd/stress/main.go +++ b/cmd/stress/main.go @@ -12,10 +12,10 @@ import ( comatproto "github.com/bluesky-social/indigo/api/atproto" appbsky "github.com/bluesky-social/indigo/api/bsky" "github.com/bluesky-social/indigo/carstore" - cliutil "github.com/bluesky-social/indigo/cmd/gosky/util" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/testing" + "github.com/bluesky-social/indigo/util/cliutil" "github.com/bluesky-social/indigo/version" "github.com/bluesky-social/indigo/xrpc" diff --git a/pds/handlers_test.go b/pds/handlers_test.go index 9c5d5396..82cc70ad 100644 --- a/pds/handlers_test.go +++ b/pds/handlers_test.go @@ -11,8 +11,8 @@ import ( "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/carstore" - cliutil "github.com/bluesky-social/indigo/cmd/gosky/util" "github.com/bluesky-social/indigo/plc" + "github.com/bluesky-social/indigo/util/cliutil" "github.com/whyrusleeping/go-did" "gorm.io/gorm" ) diff --git a/cmd/gosky/util/key.go b/util/cliutil/key.go similarity index 100% rename from cmd/gosky/util/key.go rename to util/cliutil/key.go diff --git a/cmd/gosky/util/key_test.go b/util/cliutil/key_test.go similarity index 100% rename from cmd/gosky/util/key_test.go rename to util/cliutil/key_test.go diff --git a/cmd/gosky/util/util.go b/util/cliutil/util.go similarity index 100% rename from cmd/gosky/util/util.go rename to util/cliutil/util.go -- 2.51.2 From 8586af4d8a7d5dbd41dc41c00448fc2b5e1b8b0b Mon Sep 17 00:00:00 2001 From: Jaz Volpert Date: Wed, 12 Jul 2023 22:42:50 +0000 Subject: [PATCH 3/7] Move version and testscripts --- cmd/beemo/main.go | 2 +- cmd/bigsky/main.go | 2 +- cmd/fakermaker/main.go | 2 +- cmd/gosky/main.go | 2 +- cmd/labelmaker/main.go | 2 +- cmd/laputa/main.go | 2 +- cmd/palomar/main.go | 2 +- cmd/sonar/main.go | 2 +- cmd/stress/main.go | 2 +- fakedata/accounts.go | 2 +- labeler/hiveai.go | 2 +- labeler/micro_nsfw_img.go | 2 +- labeler/sqrl.go | 2 +- labeler/xrpc_endpoints.go | 2 +- search/server.go | 2 +- {testscripts => testing/testscripts}/cleanup.sh | 0 {testscripts => testing/testscripts}/pdstest.sh | 0 {version => util/version}/version.go | 0 xrpc/xrpc.go | 2 +- 19 files changed, 16 insertions(+), 16 deletions(-) rename {testscripts => testing/testscripts}/cleanup.sh (100%) rename {testscripts => testing/testscripts}/pdstest.sh (100%) rename {version => util/version}/version.go (100%) diff --git a/cmd/beemo/main.go b/cmd/beemo/main.go index e8025549..35025a2f 100644 --- a/cmd/beemo/main.go +++ b/cmd/beemo/main.go @@ -15,7 +15,7 @@ import ( comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/util" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" "github.com/bluesky-social/indigo/xrpc" _ "github.com/joho/godotenv/autoload" diff --git a/cmd/bigsky/main.go b/cmd/bigsky/main.go index c4088b6c..de5a818d 100644 --- a/cmd/bigsky/main.go +++ b/cmd/bigsky/main.go @@ -17,7 +17,7 @@ import ( "github.com/bluesky-social/indigo/plc" "github.com/bluesky-social/indigo/repomgr" "github.com/bluesky-social/indigo/util/cliutil" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" "github.com/bluesky-social/indigo/xrpc" _ "net/http/pprof" diff --git a/cmd/fakermaker/main.go b/cmd/fakermaker/main.go index eb2ec43e..41f985e2 100644 --- a/cmd/fakermaker/main.go +++ b/cmd/fakermaker/main.go @@ -12,7 +12,7 @@ import ( "github.com/bluesky-social/indigo/fakedata" "github.com/bluesky-social/indigo/util/cliutil" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" _ "github.com/joho/godotenv/autoload" diff --git a/cmd/gosky/main.go b/cmd/gosky/main.go index 56c2dd69..8f9a023a 100644 --- a/cmd/gosky/main.go +++ b/cmd/gosky/main.go @@ -24,7 +24,7 @@ import ( "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/util" "github.com/bluesky-social/indigo/util/cliutil" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" "github.com/bluesky-social/indigo/xrpc" lru "github.com/hashicorp/golang-lru" diff --git a/cmd/labelmaker/main.go b/cmd/labelmaker/main.go index c8bc9e54..22d56c1f 100644 --- a/cmd/labelmaker/main.go +++ b/cmd/labelmaker/main.go @@ -8,7 +8,7 @@ import ( "github.com/bluesky-social/indigo/carstore" "github.com/bluesky-social/indigo/labeler" "github.com/bluesky-social/indigo/util/cliutil" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" "github.com/urfave/cli/v2" _ "github.com/joho/godotenv/autoload" diff --git a/cmd/laputa/main.go b/cmd/laputa/main.go index afdc00d4..80adcadd 100644 --- a/cmd/laputa/main.go +++ b/cmd/laputa/main.go @@ -9,7 +9,7 @@ import ( "github.com/bluesky-social/indigo/pds" "github.com/bluesky-social/indigo/plc" "github.com/bluesky-social/indigo/util/cliutil" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" _ "github.com/joho/godotenv/autoload" diff --git a/cmd/palomar/main.go b/cmd/palomar/main.go index e5873057..19de9183 100644 --- a/cmd/palomar/main.go +++ b/cmd/palomar/main.go @@ -13,7 +13,7 @@ import ( "github.com/bluesky-social/indigo/search" "github.com/bluesky-social/indigo/util/cliutil" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" logging "github.com/ipfs/go-log" es "github.com/opensearch-project/opensearch-go/v2" cli "github.com/urfave/cli/v2" diff --git a/cmd/sonar/main.go b/cmd/sonar/main.go index 8c653056..6d87af01 100644 --- a/cmd/sonar/main.go +++ b/cmd/sonar/main.go @@ -14,7 +14,7 @@ import ( "github.com/bluesky-social/indigo/events" "github.com/bluesky-social/indigo/sonar" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" "github.com/gorilla/websocket" "github.com/prometheus/client_golang/prometheus/promhttp" "go.uber.org/zap" diff --git a/cmd/stress/main.go b/cmd/stress/main.go index 1d147cad..61b1a8fd 100644 --- a/cmd/stress/main.go +++ b/cmd/stress/main.go @@ -16,7 +16,7 @@ import ( "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/testing" "github.com/bluesky-social/indigo/util/cliutil" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" "github.com/bluesky-social/indigo/xrpc" "github.com/ipfs/go-cid" diff --git a/fakedata/accounts.go b/fakedata/accounts.go index aeb5a43d..c8b9892e 100644 --- a/fakedata/accounts.go +++ b/fakedata/accounts.go @@ -8,7 +8,7 @@ import ( comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/util" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" "github.com/bluesky-social/indigo/xrpc" ) diff --git a/labeler/hiveai.go b/labeler/hiveai.go index e6780d04..0e8da740 100644 --- a/labeler/hiveai.go +++ b/labeler/hiveai.go @@ -11,7 +11,7 @@ import ( lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/util" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" ) type HiveAILabeler struct { diff --git a/labeler/micro_nsfw_img.go b/labeler/micro_nsfw_img.go index 3a3912b2..f64a3854 100644 --- a/labeler/micro_nsfw_img.go +++ b/labeler/micro_nsfw_img.go @@ -11,7 +11,7 @@ import ( lexutil "github.com/bluesky-social/indigo/lex/util" util "github.com/bluesky-social/indigo/util" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" ) type MicroNSFWImgLabeler struct { diff --git a/labeler/sqrl.go b/labeler/sqrl.go index 139f0ffe..ea50b953 100644 --- a/labeler/sqrl.go +++ b/labeler/sqrl.go @@ -10,7 +10,7 @@ import ( appbsky "github.com/bluesky-social/indigo/api/bsky" "github.com/bluesky-social/indigo/util" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" ) type SQRLLabeler struct { diff --git a/labeler/xrpc_endpoints.go b/labeler/xrpc_endpoints.go index bf663b24..6c05f87c 100644 --- a/labeler/xrpc_endpoints.go +++ b/labeler/xrpc_endpoints.go @@ -6,7 +6,7 @@ import ( atproto "github.com/bluesky-social/indigo/api/atproto" label "github.com/bluesky-social/indigo/api/label" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" "github.com/labstack/echo/v4" "go.opentelemetry.io/otel" diff --git a/search/server.go b/search/server.go index ae9420e4..3efc1f80 100644 --- a/search/server.go +++ b/search/server.go @@ -17,7 +17,7 @@ import ( lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/repomgr" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" "github.com/bluesky-social/indigo/xrpc" "github.com/gorilla/websocket" diff --git a/testscripts/cleanup.sh b/testing/testscripts/cleanup.sh similarity index 100% rename from testscripts/cleanup.sh rename to testing/testscripts/cleanup.sh diff --git a/testscripts/pdstest.sh b/testing/testscripts/pdstest.sh similarity index 100% rename from testscripts/pdstest.sh rename to testing/testscripts/pdstest.sh diff --git a/version/version.go b/util/version/version.go similarity index 100% rename from version/version.go rename to util/version/version.go diff --git a/xrpc/xrpc.go b/xrpc/xrpc.go index 7bed73e2..50e95aac 100644 --- a/xrpc/xrpc.go +++ b/xrpc/xrpc.go @@ -12,7 +12,7 @@ import ( "strings" "github.com/bluesky-social/indigo/util" - "github.com/bluesky-social/indigo/version" + "github.com/bluesky-social/indigo/util/version" ) type Client struct { -- 2.51.2 From 02876216182e17d209f7372d2e8dcb110787625f Mon Sep 17 00:00:00 2001 From: whyrusleeping Date: Thu, 13 Jul 2023 14:04:49 -0700 Subject: [PATCH 4/7] better logging when handle resolution fails --- api/extra.go | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/api/extra.go b/api/extra.go index 265717d5..b1c3caf9 100644 --- a/api/extra.go +++ b/api/extra.go @@ -109,7 +109,11 @@ func (dr *ProdHandleResolver) ResolveHandleToDid(ctx context.Context, handle str return parsed.String(), nil } - log.Infof("failed to resolve handle (%s) through HTTP well-known route: %s", handle, wkerr) + if wkerr != nil { + log.Infof("failed to resolve handle (%s) through HTTP well-known route: %s", handle, wkerr) + } else if resp.StatusCode != 200 { + log.Infof("failed to resolve handle (%s) through HTTP well-known route: status=%d", handle, resp.StatusCode) + } res, err := net.LookupTXT("_atproto." + handle) if err != nil { -- 2.51.2 From 01375c30fb09e37b2eab5987b5bec07523358a75 Mon Sep 17 00:00:00 2001 From: whyrusleeping Date: Thu, 13 Jul 2023 15:49:51 -0700 Subject: [PATCH 5/7] add in ability to edit handle resolve request --- api/extra.go | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/api/extra.go b/api/extra.go index b1c3caf9..1c108ad4 100644 --- a/api/extra.go +++ b/api/extra.go @@ -73,6 +73,7 @@ type HandleResolver interface { } type ProdHandleResolver struct { + ReqMod func(*http.Request, string) error } func (dr *ProdHandleResolver) ResolveHandleToDid(ctx context.Context, handle string) (string, error) { @@ -89,6 +90,12 @@ func (dr *ProdHandleResolver) ResolveHandleToDid(ctx context.Context, handle str return "", err } + if dr.ReqMod != nil { + if err := dr.ReqMod(req, handle); err != nil { + return "", err + } + } + req = req.WithContext(ctx) resp, wkerr := c.Do(req) -- 2.51.2 From d6682e8f93ca2387ac1018236564e805a60f5482 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Wed, 12 Jul 2023 15:38:54 -0700 Subject: [PATCH 6/7] mv util/dbcid.go to models/dbcid.go --- carstore/bs.go | 21 +++++++++++---------- events/dbpersist.go | 10 +++++----- {util => models}/dbcid.go | 2 +- 3 files changed, 17 insertions(+), 16 deletions(-) rename {util => models}/dbcid.go (98%) diff --git a/carstore/bs.go b/carstore/bs.go index 1919d049..488f0c98 100644 --- a/carstore/bs.go +++ b/carstore/bs.go @@ -13,6 +13,7 @@ import ( "sync" "time" + "github.com/bluesky-social/indigo/models" util "github.com/bluesky-social/indigo/util" blockformat "github.com/ipfs/go-block-format" @@ -67,7 +68,7 @@ type CarShard struct { ID uint `gorm:"primarykey"` CreatedAt time.Time - Root util.DbCID + Root models.DbCID DataStart int64 Seq int `gorm:"index"` Path string @@ -76,8 +77,8 @@ type CarShard struct { } type blockRef struct { - ID uint `gorm:"primarykey"` - Cid util.DbCID `gorm:"index"` + ID uint `gorm:"primarykey"` + Cid models.DbCID `gorm:"index"` Shard uint Offset int64 //User uint `gorm:"index"` @@ -103,7 +104,7 @@ func (uv *userView) Has(ctx context.Context, k cid.Cid) (bool, error) { Model(blockRef{}). Select("path, block_refs.offset"). Joins("left join car_shards on block_refs.shard = car_shards.id"). - Where("usr = ? AND cid = ?", uv.user, util.DbCID{k}). + Where("usr = ? AND cid = ?", uv.user, models.DbCID{k}). Count(&count).Error; err != nil { return false, err } @@ -133,7 +134,7 @@ func (uv *userView) Get(ctx context.Context, k cid.Cid) (blockformat.Block, erro Model(blockRef{}). Select("path, block_refs.offset"). Joins("left join car_shards on block_refs.shard = car_shards.id"). - Where("usr = ? AND cid = ?", uv.user, util.DbCID{k}). + Where("usr = ? AND cid = ?", uv.user, models.DbCID{k}). Find(&info).Error; err != nil { return nil, err } @@ -351,7 +352,7 @@ func (cs *CarStore) ReadUserCar(ctx context.Context, user util.Uid, earlyCid, la if earlyCid.Defined() { var untilShard CarShard - if err := cs.meta.First(&untilShard, "root = ? AND usr = ?", util.DbCID{earlyCid}, user).Error; err != nil { + if err := cs.meta.First(&untilShard, "root = ? AND usr = ?", models.DbCID{earlyCid}, user).Error; err != nil { return fmt.Errorf("finding early shard: %w", err) } earlySeq = untilShard.Seq @@ -359,7 +360,7 @@ func (cs *CarStore) ReadUserCar(ctx context.Context, user util.Uid, earlyCid, la if lateCid.Defined() { var fromShard CarShard - if err := cs.meta.First(&fromShard, "root = ? AND usr = ?", util.DbCID{lateCid}, user).Error; err != nil { + if err := cs.meta.First(&fromShard, "root = ? AND usr = ?", models.DbCID{lateCid}, user).Error; err != nil { return fmt.Errorf("finding late shard: %w", err) } lateSeq = fromShard.Seq @@ -616,7 +617,7 @@ func (ds *DeltaSession) closeWithRoot(ctx context.Context, root cid.Cid, rebase // adding things to the db by map is the only way to get gorm to not // add the 'returning' clause, which costs a lot of time brefs = append(brefs, map[string]interface{}{ - "cid": util.DbCID{k}, + "cid": models.DbCID{k}, "offset": offset, }) @@ -630,7 +631,7 @@ func (ds *DeltaSession) closeWithRoot(ctx context.Context, root cid.Cid, rebase // TODO: all this database work needs to be in a single transaction shard := CarShard{ - Root: util.DbCID{root}, + Root: models.DbCID{root}, DataStart: hnw, Seq: ds.seq, Path: path, @@ -835,7 +836,7 @@ func (cs *CarStore) checkFork(ctx context.Context, user util.Uid, prev cid.Cid) } var maybeShard CarShard - if err := cs.meta.WithContext(ctx).Model(CarShard{}).Find(&maybeShard, "usr = ? AND root = ?", user, &util.DbCID{prev}).Error; err != nil { + if err := cs.meta.WithContext(ctx).Model(CarShard{}).Find(&maybeShard, "usr = ? AND root = ?", user, &models.DbCID{prev}).Error; err != nil { return false, err } diff --git a/events/dbpersist.go b/events/dbpersist.go index 55bee282..5f0cb4ca 100644 --- a/events/dbpersist.go +++ b/events/dbpersist.go @@ -67,8 +67,8 @@ type DbPersistence struct { type RepoEventRecord struct { Seq uint `gorm:"primarykey"` - Commit *util.DbCID - Prev *util.DbCID + Commit *models.DbCID + Prev *models.DbCID NewHandle *string // NewHandle is only set if this is a handle change event Time time.Time @@ -249,9 +249,9 @@ func (p *DbPersistence) RecordFromRepoCommit(ctx context.Context, evt *comatprot return nil, err } - var prev *util.DbCID + var prev *models.DbCID if evt.Prev != nil && evt.Prev.Defined() { - prev = &util.DbCID{cid.Cid(*evt.Prev)} + prev = &models.DbCID{cid.Cid(*evt.Prev)} } var blobs []byte @@ -269,7 +269,7 @@ func (p *DbPersistence) RecordFromRepoCommit(ctx context.Context, evt *comatprot } rer := RepoEventRecord{ - Commit: &util.DbCID{cid.Cid(evt.Commit)}, + Commit: &models.DbCID{cid.Cid(evt.Commit)}, Prev: prev, Repo: uid, Type: "repo_append", // TODO: refactor to "#commit"? can "rebase" come through this path? diff --git a/util/dbcid.go b/models/dbcid.go similarity index 98% rename from util/dbcid.go rename to models/dbcid.go index cb490f4a..366a0e82 100644 --- a/util/dbcid.go +++ b/models/dbcid.go @@ -1,4 +1,4 @@ -package util +package models import ( "database/sql/driver" -- 2.51.2 From ebc8cbc8404800a10d0877774c83e770749a38f5 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Wed, 12 Jul 2023 15:48:44 -0700 Subject: [PATCH 7/7] mv util/uid.go models/uid.go --- bgs/bgs.go | 6 +++--- carstore/bs.go | 39 ++++++++++++++++----------------- events/dbpersist.go | 14 ++++++------ events/events.go | 12 +++++------ events/persist.go | 10 ++++----- indexer/crawler.go | 15 ++++++------- indexer/indexer.go | 8 +++---- labeler/service.go | 3 +-- models/models.go | 19 ++++++++-------- {util => models}/uid.go | 2 +- notifs/notifs.go | 43 ++++++++++++++++++------------------ notifs/null.go | 17 +++++++-------- pds/feedgen.go | 10 ++++----- pds/server.go | 4 ++-- repomgr/dbheadstore.go | 9 ++++---- repomgr/memheadstore.go | 13 +++++------ repomgr/repomgr.go | 48 +++++++++++++++++++++-------------------- 17 files changed, 135 insertions(+), 137 deletions(-) rename {util => models}/uid.go (50%) diff --git a/bgs/bgs.go b/bgs/bgs.go index 50bec20a..3537a36b 100644 --- a/bgs/bgs.go +++ b/bgs/bgs.go @@ -27,8 +27,8 @@ import ( "github.com/bluesky-social/indigo/models" "github.com/bluesky-social/indigo/repomgr" "github.com/bluesky-social/indigo/util" - bsutil "github.com/bluesky-social/indigo/util" "github.com/bluesky-social/indigo/xrpc" + "github.com/gorilla/websocket" "github.com/ipfs/go-cid" logging "github.com/ipfs/go-log" @@ -337,7 +337,7 @@ func (bgs *BGS) checkAdminAuth(next echo.HandlerFunc) echo.HandlerFunc { } type User struct { - ID bsutil.Uid `gorm:"primarykey"` + ID models.Uid `gorm:"primarykey"` CreatedAt time.Time UpdatedAt time.Time DeletedAt gorm.DeletedAt `gorm:"index"` @@ -619,7 +619,7 @@ func (bgs *BGS) handleFedEvent(ctx context.Context, host *models.PDS, env *event } } -func (s *BGS) syncUserBlobs(ctx context.Context, pds *models.PDS, user bsutil.Uid, blobs []string) error { +func (s *BGS) syncUserBlobs(ctx context.Context, pds *models.PDS, user models.Uid, blobs []string) error { if s.blobs == nil { log.Debugf("blob syncing disabled") return nil diff --git a/carstore/bs.go b/carstore/bs.go index 488f0c98..822c0a2e 100644 --- a/carstore/bs.go +++ b/carstore/bs.go @@ -14,7 +14,6 @@ import ( "time" "github.com/bluesky-social/indigo/models" - util "github.com/bluesky-social/indigo/util" blockformat "github.com/ipfs/go-block-format" "github.com/ipfs/go-cid" @@ -36,7 +35,7 @@ type CarStore struct { rootDir string lscLk sync.Mutex - lastShardCache map[util.Uid]*CarShard + lastShardCache map[models.Uid]*CarShard } func NewCarStore(meta *gorm.DB, root string) (*CarStore, error) { @@ -55,7 +54,7 @@ func NewCarStore(meta *gorm.DB, root string) (*CarStore, error) { return &CarStore{ meta: meta, rootDir: root, - lastShardCache: make(map[util.Uid]*CarShard), + lastShardCache: make(map[models.Uid]*CarShard), }, nil } @@ -72,7 +71,7 @@ type CarShard struct { DataStart int64 Seq int `gorm:"index"` Path string - Usr util.Uid `gorm:"index"` + Usr models.Uid `gorm:"index"` Rebase bool } @@ -86,7 +85,7 @@ type blockRef struct { type userView struct { cs *CarStore - user util.Uid + user models.Uid cache map[cid.Cid]blockformat.Block prefetch bool @@ -240,13 +239,13 @@ type DeltaSession struct { fresh blockstore.Blockstore blks map[cid.Cid]blockformat.Block base blockstore.Blockstore - user util.Uid + user models.Uid seq int readonly bool cs *CarStore } -func (cs *CarStore) checkLastShardCache(user util.Uid) *CarShard { +func (cs *CarStore) checkLastShardCache(user models.Uid) *CarShard { cs.lscLk.Lock() defer cs.lscLk.Unlock() @@ -258,14 +257,14 @@ func (cs *CarStore) checkLastShardCache(user util.Uid) *CarShard { return nil } -func (cs *CarStore) putLastShardCache(user util.Uid, ls *CarShard) { +func (cs *CarStore) putLastShardCache(user models.Uid, ls *CarShard) { cs.lscLk.Lock() defer cs.lscLk.Unlock() cs.lastShardCache[user] = ls } -func (cs *CarStore) getLastShard(ctx context.Context, user util.Uid) (*CarShard, error) { +func (cs *CarStore) getLastShard(ctx context.Context, user models.Uid) (*CarShard, error) { maybeLs := cs.checkLastShardCache(user) if maybeLs != nil { return maybeLs, nil @@ -289,7 +288,7 @@ var ErrRepoBaseMismatch = fmt.Errorf("attempted a delta session on top of the wr var ErrRepoFork = fmt.Errorf("repo fork detected") -func (cs *CarStore) NewDeltaSession(ctx context.Context, user util.Uid, prev *cid.Cid) (*DeltaSession, error) { +func (cs *CarStore) NewDeltaSession(ctx context.Context, user models.Uid, prev *cid.Cid) (*DeltaSession, error) { ctx, span := otel.Tracer("carstore").Start(ctx, "NewSession") defer span.End() @@ -330,7 +329,7 @@ func (cs *CarStore) NewDeltaSession(ctx context.Context, user util.Uid, prev *ci }, nil } -func (cs *CarStore) ReadOnlySession(user util.Uid) (*DeltaSession, error) { +func (cs *CarStore) ReadOnlySession(user models.Uid) (*DeltaSession, error) { return &DeltaSession{ base: &userView{ user: user, @@ -344,7 +343,7 @@ func (cs *CarStore) ReadOnlySession(user util.Uid) (*DeltaSession, error) { }, nil } -func (cs *CarStore) ReadUserCar(ctx context.Context, user util.Uid, earlyCid, lateCid cid.Cid, incremental bool, w io.Writer) error { +func (cs *CarStore) ReadUserCar(ctx context.Context, user models.Uid, earlyCid, lateCid cid.Cid, incremental bool, w io.Writer) error { ctx, span := otel.Tracer("carstore").Start(ctx, "ReadUserCar") defer span.End() @@ -525,10 +524,10 @@ func (ds *DeltaSession) GetSize(ctx context.Context, c cid.Cid) (int, error) { return ds.base.GetSize(ctx, c) } -func fnameForShard(user util.Uid, seq int) string { +func fnameForShard(user models.Uid, seq int) string { return fmt.Sprintf("sh-%d-%d", user, seq) } -func (cs *CarStore) openNewShardFile(ctx context.Context, user util.Uid, seq int) (*os.File, string, error) { +func (cs *CarStore) openNewShardFile(ctx context.Context, user models.Uid, seq int) (*os.File, string, error) { // TODO: some overwrite protections fname := filepath.Join(cs.rootDir, fnameForShard(user, seq)) fi, err := os.Create(fname) @@ -539,7 +538,7 @@ func (cs *CarStore) openNewShardFile(ctx context.Context, user util.Uid, seq int return fi, fname, nil } -func (cs *CarStore) writeNewShardFile(ctx context.Context, user util.Uid, seq int, data []byte) (string, error) { +func (cs *CarStore) writeNewShardFile(ctx context.Context, user models.Uid, seq int, data []byte) (string, error) { _, span := otel.Tracer("carstore").Start(ctx, "writeNewShardFile") defer span.End() @@ -758,7 +757,7 @@ func LdWrite(w io.Writer, d ...[]byte) (int64, error) { return int64(nw), nil } -func (cs *CarStore) ImportSlice(ctx context.Context, uid util.Uid, prev *cid.Cid, carslice []byte) (cid.Cid, *DeltaSession, error) { +func (cs *CarStore) ImportSlice(ctx context.Context, uid models.Uid, prev *cid.Cid, carslice []byte) (cid.Cid, *DeltaSession, error) { ctx, span := otel.Tracer("carstore").Start(ctx, "ImportSlice") defer span.End() @@ -793,7 +792,7 @@ func (cs *CarStore) ImportSlice(ctx context.Context, uid util.Uid, prev *cid.Cid return carr.Header.Roots[0], ds, nil } -func (cs *CarStore) GetUserRepoHead(ctx context.Context, user util.Uid) (cid.Cid, error) { +func (cs *CarStore) GetUserRepoHead(ctx context.Context, user models.Uid) (cid.Cid, error) { lastShard, err := cs.getLastShard(ctx, user) if err != nil { return cid.Undef, err @@ -811,7 +810,7 @@ type UserStat struct { Created time.Time } -func (cs *CarStore) Stat(ctx context.Context, usr util.Uid) ([]UserStat, error) { +func (cs *CarStore) Stat(ctx context.Context, usr models.Uid) ([]UserStat, error) { var shards []CarShard if err := cs.meta.Order("seq asc").Find(&shards, "usr = ?", usr).Error; err != nil { return nil, err @@ -829,7 +828,7 @@ func (cs *CarStore) Stat(ctx context.Context, usr util.Uid) ([]UserStat, error) return out, nil } -func (cs *CarStore) checkFork(ctx context.Context, user util.Uid, prev cid.Cid) (bool, error) { +func (cs *CarStore) checkFork(ctx context.Context, user models.Uid, prev cid.Cid) (bool, error) { lastShard, err := cs.getLastShard(ctx, user) if err != nil { return false, err @@ -852,7 +851,7 @@ func (cs *CarStore) checkFork(ctx context.Context, user util.Uid, prev cid.Cid) return true, nil } -func (cs *CarStore) TakeDownRepo(ctx context.Context, user util.Uid) error { +func (cs *CarStore) TakeDownRepo(ctx context.Context, user models.Uid) error { var shards []CarShard if err := cs.meta.Find(&shards, "usr = ?", user).Error; err != nil { return err diff --git a/events/dbpersist.go b/events/dbpersist.go index 5f0cb4ca..76ae8a1b 100644 --- a/events/dbpersist.go +++ b/events/dbpersist.go @@ -73,7 +73,7 @@ type RepoEventRecord struct { Time time.Time Blobs []byte - Repo util.Uid + Repo models.Uid Type string Rebase bool @@ -388,9 +388,9 @@ func (p *DbPersistence) hydrateBatch(ctx context.Context, batch []*RepoEventReco return nil } -func (p *DbPersistence) uidForDid(ctx context.Context, did string) (util.Uid, error) { +func (p *DbPersistence) uidForDid(ctx context.Context, did string) (models.Uid, error) { if uid, ok := p.didCache.Get(did); ok { - return uid.(util.Uid), nil + return uid.(models.Uid), nil } var u models.ActorInfo @@ -403,7 +403,7 @@ func (p *DbPersistence) uidForDid(ctx context.Context, did string) (util.Uid, er return u.Uid, nil } -func (p *DbPersistence) didForUid(ctx context.Context, uid util.Uid) (string, error) { +func (p *DbPersistence) didForUid(ctx context.Context, uid models.Uid) (string, error) { if did, ok := p.uidCache.Get(uid); ok { return did.(string), nil } @@ -513,11 +513,11 @@ func (p *DbPersistence) readCarSlice(ctx context.Context, rer *RepoEventRecord) return buf.Bytes(), nil } -func (p *DbPersistence) TakeDownRepo(ctx context.Context, usr util.Uid) error { +func (p *DbPersistence) TakeDownRepo(ctx context.Context, usr models.Uid) error { return p.deleteAllEventsForUser(ctx, usr) } -func (p *DbPersistence) deleteAllEventsForUser(ctx context.Context, usr util.Uid) error { +func (p *DbPersistence) deleteAllEventsForUser(ctx context.Context, usr models.Uid) error { if err := p.db.Where("repo = ?", usr).Delete(&RepoEventRecord{}).Error; err != nil { return err } @@ -525,7 +525,7 @@ func (p *DbPersistence) deleteAllEventsForUser(ctx context.Context, usr util.Uid return nil } -func (p *DbPersistence) RebaseRepoEvents(ctx context.Context, usr util.Uid) error { +func (p *DbPersistence) RebaseRepoEvents(ctx context.Context, usr models.Uid) error { // a little weird that this is the same action as a takedown return p.deleteAllEventsForUser(ctx, usr) } diff --git a/events/events.go b/events/events.go index 4196f5d1..39511d6a 100644 --- a/events/events.go +++ b/events/events.go @@ -8,8 +8,8 @@ import ( comatproto "github.com/bluesky-social/indigo/api/atproto" label "github.com/bluesky-social/indigo/api/label" + "github.com/bluesky-social/indigo/models" - "github.com/bluesky-social/indigo/util" logging "github.com/ipfs/go-log" "go.opentelemetry.io/otel" ) @@ -107,9 +107,9 @@ type XRPCStreamEvent struct { LabelInfo *label.SubscribeLabels_Info // some private fields for internal routing perf - PrivUid util.Uid `json:"-" cborgen:"-"` - PrivPdsId uint `json:"-" cborgen:"-"` - PrivRelevantPds []uint `json:"-" cborgen:"-"` + PrivUid models.Uid `json:"-" cborgen:"-"` + PrivPdsId uint `json:"-" cborgen:"-"` + PrivRelevantPds []uint `json:"-" cborgen:"-"` } type ErrorFrame struct { @@ -195,10 +195,10 @@ func (em *EventManager) addSubscriber(sub *Subscriber) { em.subs = append(em.subs, sub) } -func (em *EventManager) TakeDownRepo(ctx context.Context, user util.Uid) error { +func (em *EventManager) TakeDownRepo(ctx context.Context, user models.Uid) error { return em.persister.TakeDownRepo(ctx, user) } -func (em *EventManager) HandleRebase(ctx context.Context, user util.Uid) error { +func (em *EventManager) HandleRebase(ctx context.Context, user models.Uid) error { return em.persister.RebaseRepoEvents(ctx, user) } diff --git a/events/persist.go b/events/persist.go index bed82018..8f77082e 100644 --- a/events/persist.go +++ b/events/persist.go @@ -5,15 +5,15 @@ import ( "fmt" "sync" - "github.com/bluesky-social/indigo/util" + "github.com/bluesky-social/indigo/models" ) // Note that this interface looks generic, but some persisters might only work with RepoAppend or LabelLabels type EventPersistence interface { Persist(ctx context.Context, e *XRPCStreamEvent) error Playback(ctx context.Context, since int64, cb func(*XRPCStreamEvent) error) error - TakeDownRepo(ctx context.Context, usr util.Uid) error - RebaseRepoEvents(ctx context.Context, usr util.Uid) error + TakeDownRepo(ctx context.Context, usr models.Uid) error + RebaseRepoEvents(ctx context.Context, usr models.Uid) error SetEventBroadcaster(func(*XRPCStreamEvent)) } @@ -77,11 +77,11 @@ func (mp *MemPersister) Playback(ctx context.Context, since int64, cb func(*XRPC return nil } -func (mp *MemPersister) TakeDownRepo(ctx context.Context, uid util.Uid) error { +func (mp *MemPersister) TakeDownRepo(ctx context.Context, uid models.Uid) error { return fmt.Errorf("repo takedowns not currently supported by memory persister, test usage only") } -func (mp *MemPersister) RebaseRepoEvents(ctx context.Context, usr util.Uid) error { +func (mp *MemPersister) RebaseRepoEvents(ctx context.Context, usr models.Uid) error { return fmt.Errorf("repo rebases not currently supported by memory persister, test usage only") } diff --git a/indexer/crawler.go b/indexer/crawler.go index bb673ed3..7ac30886 100644 --- a/indexer/crawler.go +++ b/indexer/crawler.go @@ -7,7 +7,6 @@ import ( comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/models" - "github.com/bluesky-social/indigo/util" "go.opentelemetry.io/otel" ) @@ -19,11 +18,11 @@ type CrawlDispatcher struct { catchup chan *catchupJob - complete chan util.Uid + complete chan models.Uid maplk sync.Mutex - todo map[util.Uid]*crawlWork - inProgress map[util.Uid]*crawlWork + todo map[models.Uid]*crawlWork + inProgress map[models.Uid]*crawlWork doRepoCrawl func(context.Context, *crawlWork) error @@ -38,12 +37,12 @@ func NewCrawlDispatcher(repoFn func(context.Context, *crawlWork) error, concurre return &CrawlDispatcher{ ingest: make(chan *models.ActorInfo), repoSync: make(chan *crawlWork), - complete: make(chan util.Uid), + complete: make(chan models.Uid), catchup: make(chan *catchupJob), doRepoCrawl: repoFn, concurrency: concurrency, - todo: make(map[util.Uid]*crawlWork), - inProgress: make(map[util.Uid]*crawlWork), + todo: make(map[models.Uid]*crawlWork), + inProgress: make(map[models.Uid]*crawlWork), }, nil } @@ -227,7 +226,7 @@ func (c *CrawlDispatcher) AddToCatchupQueue(ctx context.Context, host *models.PD } } -func (c *CrawlDispatcher) RepoInSlowPath(ctx context.Context, host *models.PDS, uid util.Uid) bool { +func (c *CrawlDispatcher) RepoInSlowPath(ctx context.Context, host *models.PDS, uid models.Uid) bool { c.maplk.Lock() defer c.maplk.Unlock() if _, ok := c.todo[uid]; ok { diff --git a/indexer/indexer.go b/indexer/indexer.go index 8aac49b9..5400222b 100644 --- a/indexer/indexer.go +++ b/indexer/indexer.go @@ -568,7 +568,7 @@ func (ix *Indexer) GetPostOrMissing(ctx context.Context, uri string) (*models.Fe return &post, nil } -func (ix *Indexer) handleRecordCreateFeedPost(ctx context.Context, user util.Uid, rkey string, rcid cid.Cid, rec *bsky.FeedPost) error { +func (ix *Indexer) handleRecordCreateFeedPost(ctx context.Context, user models.Uid, rkey string, rcid cid.Cid, rec *bsky.FeedPost) error { var replyid uint if rec.Reply != nil { replyto, err := ix.GetPostOrMissing(ctx, rec.Reply.Parent.Uri) @@ -699,7 +699,7 @@ func (ix *Indexer) addUserToCrawler(ctx context.Context, ai *models.ActorInfo) e return ix.Crawler.Crawl(ctx, ai) } -func (ix *Indexer) DidForUser(ctx context.Context, uid util.Uid) (string, error) { +func (ix *Indexer) DidForUser(ctx context.Context, uid models.Uid) (string, error) { var ai models.ActorInfo if err := ix.db.First(&ai, "uid = ?", uid).Error; err != nil { return "", err @@ -708,7 +708,7 @@ func (ix *Indexer) DidForUser(ctx context.Context, uid util.Uid) (string, error) return ai.Did, nil } -func (ix *Indexer) LookupUser(ctx context.Context, id util.Uid) (*models.ActorInfo, error) { +func (ix *Indexer) LookupUser(ctx context.Context, id models.Uid) (*models.ActorInfo, error) { var ai models.ActorInfo if err := ix.db.First(&ai, "uid = ?", id).Error; err != nil { return nil, err @@ -765,7 +765,7 @@ func (ix *Indexer) addNewPostNotification(ctx context.Context, post *bsky.FeedPo return nil } -func (ix *Indexer) addNewVoteNotification(ctx context.Context, postauthor util.Uid, vr *models.VoteRecord) error { +func (ix *Indexer) addNewVoteNotification(ctx context.Context, postauthor models.Uid, vr *models.VoteRecord) error { return ix.notifman.AddUpVote(ctx, vr.Voter, vr.Post, vr.ID, postauthor) } diff --git a/labeler/service.go b/labeler/service.go index 13dea74d..4cd0d23d 100644 --- a/labeler/service.go +++ b/labeler/service.go @@ -24,7 +24,6 @@ import ( "github.com/bluesky-social/indigo/pds" "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/repomgr" - util "github.com/bluesky-social/indigo/util" cbg "github.com/whyrusleeping/cbor-gen" logging "github.com/ipfs/go-log" @@ -58,7 +57,7 @@ type RepoConfig struct { Did string Password string SigningKey *did.PrivKey - UserId util.Uid + UserId models.Uid } // In addition to configuring the service, will connect to upstream BGS and start processing events. Won't handle HTTP or WebSocket endpoints until RunAPI() is called. diff --git a/models/models.go b/models/models.go index 92578874..3b7b1cbe 100644 --- a/models/models.go +++ b/models/models.go @@ -6,14 +6,13 @@ import ( "gorm.io/gorm" bsky "github.com/bluesky-social/indigo/api/bsky" - "github.com/bluesky-social/indigo/util" "github.com/bluesky-social/indigo/xrpc" ) type FeedPost struct { gorm.Model - Author util.Uid `gorm:"index:idx_feedpost_rkey,unique"` - Rkey string `gorm:"index:idx_feedpost_rkey,unique"` + Author Uid `gorm:"index:idx_feedpost_rkey,unique"` + Rkey string `gorm:"index:idx_feedpost_rkey,unique"` Cid string UpCount int64 ReplyCount int64 @@ -28,16 +27,16 @@ type RepostRecord struct { CreatedAt time.Time RecCreated string Post uint - Reposter util.Uid - Author util.Uid + Reposter Uid + Author Uid RecCid string Rkey string } type ActorInfo struct { gorm.Model - Uid util.Uid `gorm:"uniqueindex"` - Handle string `gorm:"uniqueindex"` + Uid Uid `gorm:"uniqueindex"` + Handle string `gorm:"uniqueindex"` DisplayName string Did string `gorm:"uniqueindex"` Following int64 @@ -85,7 +84,7 @@ const ( type VoteRecord struct { gorm.Model Dir VoteDir - Voter util.Uid + Voter Uid Post uint Created string Rkey string @@ -94,8 +93,8 @@ type VoteRecord struct { type FollowRecord struct { gorm.Model - Follower util.Uid - Target util.Uid + Follower Uid + Target Uid Rkey string Cid string } diff --git a/util/uid.go b/models/uid.go similarity index 50% rename from util/uid.go rename to models/uid.go index bc803943..903978ab 100644 --- a/util/uid.go +++ b/models/uid.go @@ -1,3 +1,3 @@ -package util +package models type Uid uint diff --git a/notifs/notifs.go b/notifs/notifs.go index b613d809..bdafb4b5 100644 --- a/notifs/notifs.go +++ b/notifs/notifs.go @@ -8,7 +8,6 @@ import ( appbskytypes "github.com/bluesky-social/indigo/api/bsky" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/models" - bsutil "github.com/bluesky-social/indigo/util" "github.com/ipfs/go-cid" cbg "github.com/whyrusleeping/cbor-gen" "gorm.io/gorm" @@ -16,14 +15,14 @@ import ( ) type NotificationManager interface { - GetNotifications(ctx context.Context, user bsutil.Uid) ([]*appbskytypes.NotificationListNotifications_Notification, error) - GetCount(ctx context.Context, user bsutil.Uid) (int64, error) - UpdateSeen(ctx context.Context, usr bsutil.Uid, seen time.Time) error - AddReplyTo(ctx context.Context, user bsutil.Uid, replyid uint, replyto *models.FeedPost) error - AddMention(ctx context.Context, user bsutil.Uid, postid uint, mentioned bsutil.Uid) error - AddUpVote(ctx context.Context, voter bsutil.Uid, postid uint, voteid uint, postauthor bsutil.Uid) error - AddFollow(ctx context.Context, follower, followed bsutil.Uid, recid uint) error - AddRepost(ctx context.Context, op bsutil.Uid, repost uint, reposter bsutil.Uid) error + GetNotifications(ctx context.Context, user models.Uid) ([]*appbskytypes.NotificationListNotifications_Notification, error) + GetCount(ctx context.Context, user models.Uid) (int64, error) + UpdateSeen(ctx context.Context, usr models.Uid, seen time.Time) error + AddReplyTo(ctx context.Context, user models.Uid, replyid uint, replyto *models.FeedPost) error + AddMention(ctx context.Context, user models.Uid, postid uint, mentioned models.Uid) error + AddUpVote(ctx context.Context, voter models.Uid, postid uint, voteid uint, postauthor models.Uid) error + AddFollow(ctx context.Context, follower, followed models.Uid, recid uint) error + AddRepost(ctx context.Context, op models.Uid, repost uint, reposter models.Uid) error } var _ NotificationManager = (*DBNotifMan)(nil) @@ -33,7 +32,7 @@ type DBNotifMan struct { getRecord GetRecord } -type GetRecord func(ctx context.Context, user bsutil.Uid, collection string, rkey string, maybeCid cid.Cid) (cid.Cid, cbg.CBORMarshaler, error) +type GetRecord func(ctx context.Context, user models.Uid, collection string, rkey string, maybeCid cid.Cid) (cid.Cid, cbg.CBORMarshaler, error) func NewNotificationManager(db *gorm.DB, getrec GetRecord) *DBNotifMan { @@ -56,16 +55,16 @@ const ( type NotifRecord struct { gorm.Model - For bsutil.Uid + For models.Uid Kind int64 Record uint - Who bsutil.Uid + Who models.Uid ReplyTo uint } type NotifSeen struct { ID uint `gorm:"primarykey"` - Usr bsutil.Uid `gorm:"uniqueIndex"` + Usr models.Uid `gorm:"uniqueIndex"` LastSeen time.Time } @@ -80,7 +79,7 @@ type HydratedNotification struct { ReasonSubject *string } -func (nm *DBNotifMan) GetNotifications(ctx context.Context, user bsutil.Uid) ([]*appbskytypes.NotificationListNotifications_Notification, error) { +func (nm *DBNotifMan) GetNotifications(ctx context.Context, user models.Uid) ([]*appbskytypes.NotificationListNotifications_Notification, error) { var lastSeen time.Time if err := nm.db.Model(NotifSeen{}).Where("usr = ?", user).Select("last_seen").Scan(&lastSeen).Error; err != nil { return nil, err @@ -137,7 +136,7 @@ func (nm *DBNotifMan) hydrateNotification(ctx context.Context, nrec *NotifRecord return nil, fmt.Errorf("attempted to hydrate unknown notif kind: %d", nrec.Kind) } } -func (nm *DBNotifMan) getActor(ctx context.Context, act bsutil.Uid) (*models.ActorInfo, error) { +func (nm *DBNotifMan) getActor(ctx context.Context, act models.Uid) (*models.ActorInfo, error) { var ai models.ActorInfo if err := nm.db.First(&ai, "uid = ?", act).Error; err != nil { return nil, err @@ -294,7 +293,7 @@ func (nm *DBNotifMan) hydrateNotificationFollow(ctx context.Context, nrec *Notif } -func (nm *DBNotifMan) GetCount(ctx context.Context, user bsutil.Uid) (int64, error) { +func (nm *DBNotifMan) GetCount(ctx context.Context, user models.Uid) (int64, error) { // TODO: sql count is inefficient var lseen time.Time if err := nm.db.Model(NotifSeen{}).Where("usr = ?", user).Select("last_seen").Scan(&lseen).Error; err != nil { @@ -310,7 +309,7 @@ func (nm *DBNotifMan) GetCount(ctx context.Context, user bsutil.Uid) (int64, err return c, nil } -func (nm *DBNotifMan) UpdateSeen(ctx context.Context, usr bsutil.Uid, seen time.Time) error { +func (nm *DBNotifMan) UpdateSeen(ctx context.Context, usr models.Uid, seen time.Time) error { if err := nm.db.Clauses(clause.OnConflict{ Columns: []clause.Column{{Name: "usr"}}, DoUpdates: clause.AssignmentColumns([]string{"last_seen"}), @@ -324,7 +323,7 @@ func (nm *DBNotifMan) UpdateSeen(ctx context.Context, usr bsutil.Uid, seen time. return nil } -func (nm *DBNotifMan) AddReplyTo(ctx context.Context, user bsutil.Uid, replyid uint, replyto *models.FeedPost) error { +func (nm *DBNotifMan) AddReplyTo(ctx context.Context, user models.Uid, replyid uint, replyto *models.FeedPost) error { return nm.db.Create(&NotifRecord{ Kind: NotifKindReply, For: replyto.Author, @@ -334,7 +333,7 @@ func (nm *DBNotifMan) AddReplyTo(ctx context.Context, user bsutil.Uid, replyid u }).Error } -func (nm *DBNotifMan) AddMention(ctx context.Context, user bsutil.Uid, postid uint, mentioned bsutil.Uid) error { +func (nm *DBNotifMan) AddMention(ctx context.Context, user models.Uid, postid uint, mentioned models.Uid) error { return nm.db.Create(&NotifRecord{ For: mentioned, Kind: NotifKindMention, @@ -343,7 +342,7 @@ func (nm *DBNotifMan) AddMention(ctx context.Context, user bsutil.Uid, postid ui }).Error } -func (nm *DBNotifMan) AddUpVote(ctx context.Context, voter bsutil.Uid, postid uint, voteid uint, postauthor bsutil.Uid) error { +func (nm *DBNotifMan) AddUpVote(ctx context.Context, voter models.Uid, postid uint, voteid uint, postauthor models.Uid) error { return nm.db.Create(&NotifRecord{ For: postauthor, Kind: NotifKindUpVote, @@ -353,7 +352,7 @@ func (nm *DBNotifMan) AddUpVote(ctx context.Context, voter bsutil.Uid, postid ui }).Error } -func (nm *DBNotifMan) AddFollow(ctx context.Context, follower, followed bsutil.Uid, recid uint) error { +func (nm *DBNotifMan) AddFollow(ctx context.Context, follower, followed models.Uid, recid uint) error { return nm.db.Create(&NotifRecord{ Kind: NotifKindFollow, For: followed, @@ -362,7 +361,7 @@ func (nm *DBNotifMan) AddFollow(ctx context.Context, follower, followed bsutil.U }).Error } -func (nm *DBNotifMan) AddRepost(ctx context.Context, op bsutil.Uid, repost uint, reposter bsutil.Uid) error { +func (nm *DBNotifMan) AddRepost(ctx context.Context, op models.Uid, repost uint, reposter models.Uid) error { return nm.db.Create(&NotifRecord{ Kind: NotifKindRepost, For: op, diff --git a/notifs/null.go b/notifs/null.go index f29f8c5b..2be93e86 100644 --- a/notifs/null.go +++ b/notifs/null.go @@ -7,7 +7,6 @@ import ( appbskytypes "github.com/bluesky-social/indigo/api/bsky" "github.com/bluesky-social/indigo/models" - "github.com/bluesky-social/indigo/util" ) type NullNotifs struct { @@ -15,34 +14,34 @@ type NullNotifs struct { var _ NotificationManager = (*NullNotifs)(nil) -func (nn *NullNotifs) GetNotifications(ctx context.Context, user util.Uid) ([]*appbskytypes.NotificationListNotifications_Notification, error) { +func (nn *NullNotifs) GetNotifications(ctx context.Context, user models.Uid) ([]*appbskytypes.NotificationListNotifications_Notification, error) { return nil, fmt.Errorf("no notifications engine loaded") } -func (nn *NullNotifs) GetCount(ctx context.Context, user util.Uid) (int64, error) { +func (nn *NullNotifs) GetCount(ctx context.Context, user models.Uid) (int64, error) { return 0, fmt.Errorf("no notifications engine loaded") } -func (nn *NullNotifs) UpdateSeen(ctx context.Context, usr util.Uid, seen time.Time) error { +func (nn *NullNotifs) UpdateSeen(ctx context.Context, usr models.Uid, seen time.Time) error { return nil } -func (nn *NullNotifs) AddReplyTo(ctx context.Context, user util.Uid, replyid uint, replyto *models.FeedPost) error { +func (nn *NullNotifs) AddReplyTo(ctx context.Context, user models.Uid, replyid uint, replyto *models.FeedPost) error { return nil } -func (nn *NullNotifs) AddMention(ctx context.Context, user util.Uid, postid uint, mentioned util.Uid) error { +func (nn *NullNotifs) AddMention(ctx context.Context, user models.Uid, postid uint, mentioned models.Uid) error { return nil } -func (nn *NullNotifs) AddUpVote(ctx context.Context, voter util.Uid, postid uint, voteid uint, postauthor util.Uid) error { +func (nn *NullNotifs) AddUpVote(ctx context.Context, voter models.Uid, postid uint, voteid uint, postauthor models.Uid) error { return nil } -func (nn *NullNotifs) AddFollow(ctx context.Context, follower, followed util.Uid, recid uint) error { +func (nn *NullNotifs) AddFollow(ctx context.Context, follower, followed models.Uid, recid uint) error { return nil } -func (nn *NullNotifs) AddRepost(ctx context.Context, op util.Uid, repost uint, reposter util.Uid) error { +func (nn *NullNotifs) AddRepost(ctx context.Context, op models.Uid, repost uint, reposter models.Uid) error { return nil } diff --git a/pds/feedgen.go b/pds/feedgen.go index 49d7984b..af5a8791 100644 --- a/pds/feedgen.go +++ b/pds/feedgen.go @@ -11,7 +11,7 @@ import ( "github.com/bluesky-social/indigo/indexer" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/models" - bsutil "github.com/bluesky-social/indigo/util" + "github.com/ipfs/go-cid" "go.opentelemetry.io/otel" "gorm.io/gorm" @@ -32,7 +32,7 @@ func NewFeedGenerator(db *gorm.DB, ix *indexer.Indexer, readRecord ReadRecordFun }, nil } -type ReadRecordFunc func(context.Context, bsutil.Uid, cid.Cid) (lexutil.CBOR, error) +type ReadRecordFunc func(context.Context, models.Uid, cid.Cid) (lexutil.CBOR, error) /* type HydratedFeedItem struct { @@ -94,7 +94,7 @@ func (fg *FeedGenerator) hydrateFeed(ctx context.Context, items []*models.FeedPo return out, nil } -func (fg *FeedGenerator) didForUser(ctx context.Context, user bsutil.Uid) (string, error) { +func (fg *FeedGenerator) didForUser(ctx context.Context, user models.Uid) (string, error) { // TODO: cache the shit out of this var ai models.ActorInfo if err := fg.db.First(&ai, "uid = ?", user).Error; err != nil { @@ -104,7 +104,7 @@ func (fg *FeedGenerator) didForUser(ctx context.Context, user bsutil.Uid) (strin return ai.Did, nil } -func (fg *FeedGenerator) getActorRefInfo(ctx context.Context, user bsutil.Uid) (*bsky.ActorDefs_ProfileViewBasic, error) { +func (fg *FeedGenerator) getActorRefInfo(ctx context.Context, user models.Uid) (*bsky.ActorDefs_ProfileViewBasic, error) { // TODO: cache the shit out of this too var ai models.ActorInfo if err := fg.db.First(&ai, "uid = ?", user).Error; err != nil { @@ -153,7 +153,7 @@ func (fg *FeedGenerator) hydrateItem(ctx context.Context, item *models.FeedPost) return out, nil } -func (fg *FeedGenerator) getPostViewerState(ctx context.Context, item uint, viewer bsutil.Uid, viewerDid string) (*bsky.FeedDefs_ViewerState, error) { +func (fg *FeedGenerator) getPostViewerState(ctx context.Context, item uint, viewer models.Uid, viewerDid string) (*bsky.FeedDefs_ViewerState, error) { var out bsky.FeedDefs_ViewerState var vote models.VoteRecord diff --git a/pds/server.go b/pds/server.go index 6911ae36..4814d34d 100644 --- a/pds/server.go +++ b/pds/server.go @@ -263,7 +263,7 @@ func (s *Server) repoEventToFedEvent(ctx context.Context, evt *repomgr.RepoEvent return out, nil } -func (s *Server) readRecordFunc(ctx context.Context, user bsutil.Uid, c cid.Cid) (lexutil.CBOR, error) { +func (s *Server) readRecordFunc(ctx context.Context, user models.Uid, c cid.Cid) (lexutil.CBOR, error) { bs, err := s.cs.ReadOnlySession(user) if err != nil { return nil, err @@ -396,7 +396,7 @@ func (s *Server) HandleResolveDid(c echo.Context) error { } type User struct { - ID bsutil.Uid `gorm:"primarykey"` + ID models.Uid `gorm:"primarykey"` CreatedAt time.Time UpdatedAt time.Time DeletedAt gorm.DeletedAt `gorm:"index"` diff --git a/repomgr/dbheadstore.go b/repomgr/dbheadstore.go index c2ff7cb7..d1f9fa30 100644 --- a/repomgr/dbheadstore.go +++ b/repomgr/dbheadstore.go @@ -3,7 +3,8 @@ package repomgr import ( "context" - "github.com/bluesky-social/indigo/util" + "github.com/bluesky-social/indigo/models" + "github.com/ipfs/go-cid" "go.opentelemetry.io/otel" "gorm.io/gorm" @@ -20,7 +21,7 @@ type DbHeadStore struct { db *gorm.DB } -func (hs *DbHeadStore) InitUser(ctx context.Context, user util.Uid, root cid.Cid) error { +func (hs *DbHeadStore) InitUser(ctx context.Context, user models.Uid, root cid.Cid) error { if err := hs.db.Create(&RepoHead{ Usr: user, Root: root.String(), @@ -31,7 +32,7 @@ func (hs *DbHeadStore) InitUser(ctx context.Context, user util.Uid, root cid.Cid return nil } -func (hs *DbHeadStore) UpdateUserRepoHead(ctx context.Context, user util.Uid, root cid.Cid) error { +func (hs *DbHeadStore) UpdateUserRepoHead(ctx context.Context, user models.Uid, root cid.Cid) error { if err := hs.db.WithContext(ctx).Clauses(clause.OnConflict{ Columns: []clause.Column{{Name: "usr"}}, DoUpdates: clause.AssignmentColumns([]string{"root"}), @@ -45,7 +46,7 @@ func (hs *DbHeadStore) UpdateUserRepoHead(ctx context.Context, user util.Uid, ro return nil } -func (hs *DbHeadStore) GetUserRepoHead(ctx context.Context, user util.Uid) (cid.Cid, error) { +func (hs *DbHeadStore) GetUserRepoHead(ctx context.Context, user models.Uid) (cid.Cid, error) { ctx, span := otel.Tracer("repoman").Start(ctx, "GetUserRepoHead") defer span.End() diff --git a/repomgr/memheadstore.go b/repomgr/memheadstore.go index 1b7be585..2a08462b 100644 --- a/repomgr/memheadstore.go +++ b/repomgr/memheadstore.go @@ -4,21 +4,22 @@ import ( "context" "fmt" - "github.com/bluesky-social/indigo/util" + "github.com/bluesky-social/indigo/models" + "github.com/ipfs/go-cid" ) type MemHeadStore struct { - heads map[util.Uid]cid.Cid + heads map[models.Uid]cid.Cid } func NewMemHeadStore() *MemHeadStore { return &MemHeadStore{ - heads: make(map[util.Uid]cid.Cid), + heads: make(map[models.Uid]cid.Cid), } } -func (hs *MemHeadStore) GetUserRepoHead(ctx context.Context, user util.Uid) (cid.Cid, error) { +func (hs *MemHeadStore) GetUserRepoHead(ctx context.Context, user models.Uid) (cid.Cid, error) { h, ok := hs.heads[user] if !ok { return cid.Undef, fmt.Errorf("user head not found") @@ -27,7 +28,7 @@ func (hs *MemHeadStore) GetUserRepoHead(ctx context.Context, user util.Uid) (cid return h, nil } -func (hs *MemHeadStore) UpdateUserRepoHead(ctx context.Context, user util.Uid, root cid.Cid) error { +func (hs *MemHeadStore) UpdateUserRepoHead(ctx context.Context, user models.Uid, root cid.Cid) error { _, ok := hs.heads[user] if !ok { return fmt.Errorf("cannot update user head if it doesnt exist already") @@ -37,7 +38,7 @@ func (hs *MemHeadStore) UpdateUserRepoHead(ctx context.Context, user util.Uid, r return nil } -func (hs *MemHeadStore) InitUser(ctx context.Context, user util.Uid, root cid.Cid) error { +func (hs *MemHeadStore) InitUser(ctx context.Context, user models.Uid, root cid.Cid) error { _, ok := hs.heads[user] if ok { return fmt.Errorf("cannot init user head if it exists already") diff --git a/repomgr/repomgr.go b/repomgr/repomgr.go index fa4b92b9..8f491c44 100644 --- a/repomgr/repomgr.go +++ b/repomgr/repomgr.go @@ -13,9 +13,11 @@ import ( bsky "github.com/bluesky-social/indigo/api/bsky" "github.com/bluesky-social/indigo/carstore" lexutil "github.com/bluesky-social/indigo/lex/util" + "github.com/bluesky-social/indigo/models" "github.com/bluesky-social/indigo/mst" "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/util" + "github.com/ipfs/go-cid" "github.com/ipfs/go-datastore" blockstore "github.com/ipfs/go-ipfs-blockstore" @@ -35,15 +37,15 @@ func NewRepoManager(hs HeadStore, cs *carstore.CarStore, kmgr KeyManager) *RepoM return &RepoManager{ hs: hs, cs: cs, - userLocks: make(map[util.Uid]*userLock), + userLocks: make(map[models.Uid]*userLock), kmgr: kmgr, } } type HeadStore interface { - GetUserRepoHead(ctx context.Context, user util.Uid) (cid.Cid, error) - UpdateUserRepoHead(ctx context.Context, user util.Uid, root cid.Cid) error - InitUser(ctx context.Context, user util.Uid, root cid.Cid) error + GetUserRepoHead(ctx context.Context, user models.Uid) (cid.Cid, error) + UpdateUserRepoHead(ctx context.Context, user models.Uid, root cid.Cid) error + InitUser(ctx context.Context, user models.Uid, root cid.Cid) error } type KeyManager interface { @@ -61,7 +63,7 @@ type RepoManager struct { kmgr KeyManager lklk sync.Mutex - userLocks map[util.Uid]*userLock + userLocks map[models.Uid]*userLock events func(context.Context, *RepoEvent) } @@ -74,7 +76,7 @@ type ActorInfo struct { } type RepoEvent struct { - User util.Uid + User models.Uid OldRoot *cid.Cid NewRoot cid.Cid RepoSlice []byte @@ -102,7 +104,7 @@ const ( type RepoHead struct { gorm.Model - Usr util.Uid `gorm:"uniqueIndex"` + Usr models.Uid `gorm:"uniqueIndex"` Root string } @@ -111,7 +113,7 @@ type userLock struct { count int } -func (rm *RepoManager) lockUser(ctx context.Context, user util.Uid) func() { +func (rm *RepoManager) lockUser(ctx context.Context, user models.Uid) func() { ctx, span := otel.Tracer("repoman").Start(ctx, "userLock") defer span.End() @@ -146,7 +148,7 @@ func (rm *RepoManager) CarStore() *carstore.CarStore { return rm.cs } -func (rm *RepoManager) CreateRecord(ctx context.Context, user util.Uid, collection string, rec cbg.CBORMarshaler) (string, cid.Cid, error) { +func (rm *RepoManager) CreateRecord(ctx context.Context, user models.Uid, collection string, rec cbg.CBORMarshaler) (string, cid.Cid, error) { ctx, span := otel.Tracer("repoman").Start(ctx, "CreateRecord") defer span.End() @@ -212,7 +214,7 @@ func (rm *RepoManager) CreateRecord(ctx context.Context, user util.Uid, collecti return collection + "/" + tid, cc, nil } -func (rm *RepoManager) UpdateRecord(ctx context.Context, user util.Uid, collection, rkey string, rec cbg.CBORMarshaler) (cid.Cid, error) { +func (rm *RepoManager) UpdateRecord(ctx context.Context, user models.Uid, collection, rkey string, rec cbg.CBORMarshaler) (cid.Cid, error) { ctx, span := otel.Tracer("repoman").Start(ctx, "UpdateRecord") defer span.End() @@ -279,7 +281,7 @@ func (rm *RepoManager) UpdateRecord(ctx context.Context, user util.Uid, collecti return cc, nil } -func (rm *RepoManager) DeleteRecord(ctx context.Context, user util.Uid, collection, rkey string) error { +func (rm *RepoManager) DeleteRecord(ctx context.Context, user models.Uid, collection, rkey string) error { ctx, span := otel.Tracer("repoman").Start(ctx, "DeleteRecord") defer span.End() @@ -343,7 +345,7 @@ func (rm *RepoManager) DeleteRecord(ctx context.Context, user util.Uid, collecti } -func (rm *RepoManager) InitNewActor(ctx context.Context, user util.Uid, handle, did, displayname string, declcid, actortype string) error { +func (rm *RepoManager) InitNewActor(ctx context.Context, user models.Uid, handle, did, displayname string, declcid, actortype string) error { unlock := rm.lockUser(ctx, user) defer unlock() @@ -402,18 +404,18 @@ func (rm *RepoManager) InitNewActor(ctx context.Context, user util.Uid, handle, return nil } -func (rm *RepoManager) GetRepoRoot(ctx context.Context, user util.Uid) (cid.Cid, error) { +func (rm *RepoManager) GetRepoRoot(ctx context.Context, user models.Uid) (cid.Cid, error) { unlock := rm.lockUser(ctx, user) defer unlock() return rm.hs.GetUserRepoHead(ctx, user) } -func (rm *RepoManager) ReadRepo(ctx context.Context, user util.Uid, earlyCid, lateCid cid.Cid, w io.Writer) error { +func (rm *RepoManager) ReadRepo(ctx context.Context, user models.Uid, earlyCid, lateCid cid.Cid, w io.Writer) error { return rm.cs.ReadUserCar(ctx, user, earlyCid, lateCid, true, w) } -func (rm *RepoManager) GetRecord(ctx context.Context, user util.Uid, collection string, rkey string, maybeCid cid.Cid) (cid.Cid, cbg.CBORMarshaler, error) { +func (rm *RepoManager) GetRecord(ctx context.Context, user models.Uid, collection string, rkey string, maybeCid cid.Cid) (cid.Cid, cbg.CBORMarshaler, error) { bs, err := rm.cs.ReadOnlySession(user) if err != nil { return cid.Undef, nil, err @@ -441,7 +443,7 @@ func (rm *RepoManager) GetRecord(ctx context.Context, user util.Uid, collection return ocid, val, nil } -func (rm *RepoManager) GetProfile(ctx context.Context, uid util.Uid) (*bsky.ActorProfile, error) { +func (rm *RepoManager) GetProfile(ctx context.Context, uid models.Uid) (*bsky.ActorProfile, error) { bs, err := rm.cs.ReadOnlySession(uid) if err != nil { return nil, err @@ -472,7 +474,7 @@ func (rm *RepoManager) GetProfile(ctx context.Context, uid util.Uid) (*bsky.Acto var ErrUncleanRebase = fmt.Errorf("unclean rebase") -func (rm *RepoManager) HandleRebase(ctx context.Context, pdsid uint, uid util.Uid, did string, prev *cid.Cid, commit cid.Cid, carslice []byte) error { +func (rm *RepoManager) HandleRebase(ctx context.Context, pdsid uint, uid models.Uid, did string, prev *cid.Cid, commit cid.Cid, carslice []byte) error { ctx, span := otel.Tracer("repoman").Start(ctx, "HandleRebase") defer span.End() @@ -553,7 +555,7 @@ func (rm *RepoManager) HandleRebase(ctx context.Context, pdsid uint, uid util.Ui return nil } -func (rm *RepoManager) DoRebase(ctx context.Context, uid util.Uid) error { +func (rm *RepoManager) DoRebase(ctx context.Context, uid models.Uid) error { ctx, span := otel.Tracer("repoman").Start(ctx, "DoRebase") defer span.End() @@ -640,7 +642,7 @@ func (rm *RepoManager) CheckRepoSig(ctx context.Context, r *repo.Repo, expdid st return nil } -func (rm *RepoManager) HandleExternalUserEvent(ctx context.Context, pdsid uint, uid util.Uid, did string, prev *cid.Cid, carslice []byte, ops []*atproto.SyncSubscribeRepos_RepoOp) error { +func (rm *RepoManager) HandleExternalUserEvent(ctx context.Context, pdsid uint, uid models.Uid, did string, prev *cid.Cid, carslice []byte, ops []*atproto.SyncSubscribeRepos_RepoOp) error { ctx, span := otel.Tracer("repoman").Start(ctx, "HandleExternalUserEvent") defer span.End() @@ -737,7 +739,7 @@ func rkeyForCollection(collection string) string { return repo.NextTID() } -func (rm *RepoManager) BatchWrite(ctx context.Context, user util.Uid, writes []*atproto.RepoApplyWrites_Input_Writes_Elem) error { +func (rm *RepoManager) BatchWrite(ctx context.Context, user models.Uid, writes []*atproto.RepoApplyWrites_Input_Writes_Elem) error { ctx, span := otel.Tracer("repoman").Start(ctx, "BatchWrite") defer span.End() @@ -849,7 +851,7 @@ func (rm *RepoManager) BatchWrite(ctx context.Context, user util.Uid, writes []* return nil } -func (rm *RepoManager) ImportNewRepo(ctx context.Context, user util.Uid, repoDid string, r io.Reader, oldest cid.Cid) error { +func (rm *RepoManager) ImportNewRepo(ctx context.Context, user models.Uid, repoDid string, r io.Reader, oldest cid.Cid) error { ctx, span := otel.Tracer("repoman").Start(ctx, "ImportNewRepo") defer span.End() @@ -992,7 +994,7 @@ func processOp(ctx context.Context, bs blockstore.Blockstore, op *mst.DiffOp) (* } } -func (rm *RepoManager) processNewRepo(ctx context.Context, user util.Uid, r io.Reader, until cid.Cid, cb func(ctx context.Context, old, nu cid.Cid, finish func(context.Context) ([]byte, error), bs blockstore.Blockstore) error) error { +func (rm *RepoManager) processNewRepo(ctx context.Context, user models.Uid, r io.Reader, until cid.Cid, cb func(ctx context.Context, old, nu cid.Cid, finish func(context.Context) ([]byte, error), bs blockstore.Blockstore) error) error { ctx, span := otel.Tracer("repoman").Start(ctx, "ImportNewRepo") defer span.End() @@ -1165,7 +1167,7 @@ func walkTree(ctx context.Context, skip map[cid.Cid]bool, root cid.Cid, bs block return out, nil } -func (rm *RepoManager) TakeDownRepo(ctx context.Context, uid util.Uid) error { +func (rm *RepoManager) TakeDownRepo(ctx context.Context, uid models.Uid) error { unlock := rm.lockUser(ctx, uid) defer unlock()