diff --git a/Makefile b/Makefile index 5330957b5..91d9be459 100644 --- a/Makefile +++ b/Makefile @@ -359,7 +359,7 @@ js-lexicons: && sed -i.bak 's/AppBskyGraphBlock\.Main/AppBskyGraphBlock\.Record/' $$(find ./js/streamplace/src/lexicons/types/place/stream -type f) \ && sed -i.bak 's/PlaceStreamChatProfile\.Main/PlaceStreamChatProfile\.Record/' $$(find ./js/streamplace/src/lexicons/types/place/stream -type f) \ && for x in $$(find ./js/streamplace/src/lexicons -type f -name '*.ts'); do \ - echo 'import { AppBskyRichtextFacet, AppBskyGraphBlock, ComAtprotoRepoStrongRef, AppBskyActorDefs, ComAtprotoSyncListRepos, AppBskyActorGetProfile, AppBskyFeedGetFeedSkeleton, ComAtprotoIdentityResolveHandle, ComAtprotoModerationCreateReport, ComAtprotoRepoCreateRecord, ComAtprotoRepoDeleteRecord, ComAtprotoRepoDescribeRepo, ComAtprotoRepoGetRecord, ComAtprotoRepoListRecords, ComAtprotoRepoPutRecord, ComAtprotoRepoUploadBlob, ComAtprotoServerDescribeServer, ComAtprotoSyncGetRecord, ComAtprotoSyncListReposComAtprotoRepoCreateRecord, ComAtprotoRepoDeleteRecord, ComAtprotoRepoGetRecord, ComAtprotoRepoListRecords } from "@atproto/api"' >> $$x; \ + echo 'import { AppBskyRichtextFacet, AppBskyGraphBlock, ComAtprotoRepoStrongRef, AppBskyActorDefs, ComAtprotoSyncListRepos, AppBskyActorGetProfile, AppBskyFeedGetFeedSkeleton, ComAtprotoIdentityResolveHandle, ComAtprotoModerationCreateReport, ComAtprotoRepoCreateRecord, ComAtprotoRepoDeleteRecord, ComAtprotoRepoDescribeRepo, ComAtprotoRepoGetRecord, ComAtprotoRepoListRecords, ComAtprotoRepoPutRecord, ComAtprotoRepoUploadBlob, ComAtprotoServerDescribeServer, ComAtprotoSyncGetRecord, ComAtprotoSyncListReposComAtprotoRepoCreateRecord, ComAtprotoRepoDeleteRecord, ComAtprotoRepoGetRecord, ComAtprotoRepoListRecords, ComAtprotoIdentityRefreshIdentity } from "@atproto/api"' >> $$x; \ done \ && npx prettier --write $$(find ./js/streamplace/src/lexicons -type f -name '*.ts') \ && find . | grep bak$$ | xargs rm diff --git a/js/docs/src/content/docs/lex-reference/openapi.json b/js/docs/src/content/docs/lex-reference/openapi.json index 84f908a7a..7aefb163b 100644 --- a/js/docs/src/content/docs/lex-reference/openapi.json +++ b/js/docs/src/content/docs/lex-reference/openapi.json @@ -849,6 +849,72 @@ } } }, + "/xrpc/com.atproto.identity.refreshIdentity": { + "post": { + "summary": "Request that the server re-resolve an identity (DID and handle). The server may ignore this request, or require authentication, depending on the role, implementation, and policy of the server.", + "operationId": "com.atproto.identity.refreshIdentity", + "tags": ["com.atproto.identity"], + "responses": { + "200": { + "description": "Success", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/com.atproto.identity.defs_identityInfo" + } + } + } + }, + "400": { + "description": "Bad Request", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["error", "message"], + "properties": { + "error": { + "type": "string", + "oneOf": [ + { + "const": "HandleNotFound" + }, + { + "const": "DidNotFound" + }, + { + "const": "DidDeactivated" + } + ] + }, + "message": { + "type": "string" + } + } + } + } + } + } + }, + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "identifier": { + "type": "string", + "format": "at-identifier" + } + }, + "required": ["identifier"] + } + } + } + } + } + }, "/xrpc/com.atproto.identity.resolveHandle": { "get": { "summary": "Resolves an atproto handle (hostname) to a DID. Does not necessarily bi-directionally verify against the the DID document.", @@ -1721,6 +1787,24 @@ }, "required": ["name"] }, + "com.atproto.identity.defs_identityInfo": { + "type": "object", + "properties": { + "did": { + "type": "string", + "format": "did" + }, + "handle": { + "type": "string", + "description": "The validated handle of the account; or 'handle.invalid' if the handle did not bi-directionally match the DID document.", + "format": "handle" + }, + "didDoc": { + "description": "The complete DID document for the identity." + } + }, + "required": ["did", "handle", "didDoc"] + }, "app.bsky.feed.defs_skeletonFeedPost": { "type": "object", "properties": { diff --git a/lexicons/com/atproto/identity/refreshIdentity.json b/lexicons/com/atproto/identity/refreshIdentity.json new file mode 100644 index 000000000..4c56084b9 --- /dev/null +++ b/lexicons/com/atproto/identity/refreshIdentity.json @@ -0,0 +1,44 @@ +{ + "lexicon": 1, + "id": "com.atproto.identity.refreshIdentity", + "defs": { + "main": { + "type": "procedure", + "description": "Request that the server re-resolve an identity (DID and handle). The server may ignore this request, or require authentication, depending on the role, implementation, and policy of the server.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["identifier"], + "properties": { + "identifier": { + "type": "string", + "format": "at-identifier" + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "ref", + "ref": "com.atproto.identity.defs#identityInfo" + } + }, + "errors": [ + { + "name": "HandleNotFound", + "description": "The resolution process confirmed that the handle does not resolve to any DID." + }, + { + "name": "DidNotFound", + "description": "The DID resolution process confirmed that there is no current DID." + }, + { + "name": "DidDeactivated", + "description": "The DID previously existed, but has been deactivated." + } + ] + } + } +} diff --git a/pkg/atproto/atproto.go b/pkg/atproto/atproto.go index da73ec060..320cbdb3b 100644 --- a/pkg/atproto/atproto.go +++ b/pkg/atproto/atproto.go @@ -42,7 +42,7 @@ type mstNode struct { } func (atsync *ATProtoSynchronizer) SyncBlueskyRepo(ctx context.Context, handle string, mod model.Model) (*model.Repo, error) { - ident, err := atsync.resolveIdent(ctx, handle) + ident, err := atsync.resolveIdent(ctx, handle, true) if err != nil { return nil, fmt.Errorf("failed to resolve Bluesky handle %s: %w", handle, err) } @@ -165,16 +165,41 @@ func (atsync *ATProtoSynchronizer) SyncBlueskyRepo(ctx context.Context, handle s return &newRepo, nil } -func (atsync *ATProtoSynchronizer) resolveIdent(ctx context.Context, arg string) (*identity.Identity, error) { +func (atsync *ATProtoSynchronizer) RefreshIdentity(ctx context.Context, did string) (*identity.Identity, error) { + id, err := atsync.resolveIdent(ctx, did, false) + if err != nil { + return nil, fmt.Errorf("failed to resolve ident: %w", err) + } + newRepo := model.Repo{ + DID: id.DID.String(), + PDS: id.PDSEndpoint(), + Handle: id.Handle.String(), + } + err = atsync.Model.UpdateRepo(&newRepo) + if err != nil { + return nil, fmt.Errorf("failed to update repo: %w", err) + } + return id, nil +} + +func (atsync *ATProtoSynchronizer) resolveIdent(ctx context.Context, arg string, cached bool) (*identity.Identity, error) { if atsync.PLCDirectory == nil { atsync.PLCDirectory = CustomDirectory(atsync.CLI.PLCURL) } + if atsync.CachedPLCDirectory == nil { + cachedDir := identity.NewCacheDirectory(atsync.PLCDirectory, 250_000, time.Hour*24, time.Minute*2, time.Minute*5) + atsync.CachedPLCDirectory = &cachedDir + } + dir := atsync.PLCDirectory + if cached { + dir = atsync.CachedPLCDirectory + } id, err := syntax.ParseAtIdentifier(arg) if err != nil { return nil, err } - resolvedID, err := atsync.PLCDirectory.Lookup(ctx, *id) + resolvedID, err := dir.Lookup(ctx, *id) if err != nil { return nil, err } @@ -182,7 +207,6 @@ func (atsync *ATProtoSynchronizer) resolveIdent(ctx context.Context, arg string) return resolvedID, nil } - func CustomDirectory(plcURL string) identity.Directory { base := identity.BaseDirectory{ PLCURL: plcURL, @@ -197,8 +221,7 @@ func CustomDirectory(plcURL string) identity.Directory { // primary Bluesky PDS instance only supports HTTP resolution method SkipDNSDomainSuffixes: []string{".bsky.social"}, } - cached := identity.NewCacheDirectory(&base, 250_000, time.Hour*24, time.Minute*2, time.Minute*5) - return &cached + return &base } func DIDDoc(host string) map[string]any { diff --git a/pkg/atproto/firehose.go b/pkg/atproto/firehose.go index c4e6774b6..e3459be2c 100644 --- a/pkg/atproto/firehose.go +++ b/pkg/atproto/firehose.go @@ -34,14 +34,15 @@ import ( ) type ATProtoSynchronizer struct { - CLI *config.CLI - Model model.Model - StatefulDB *statedb.StatefulDB - LastSeen time.Time - LastEvent time.Time - Noter notificationpkg.FirebaseNotifier - Bus *bus.Bus - PLCDirectory identity.Directory + CLI *config.CLI + Model model.Model + StatefulDB *statedb.StatefulDB + LastSeen time.Time + LastEvent time.Time + Noter notificationpkg.FirebaseNotifier + Bus *bus.Bus + PLCDirectory identity.Directory + CachedPLCDirectory identity.Directory } func (atsync *ATProtoSynchronizer) StartFirehose(ctx context.Context) error { @@ -99,6 +100,10 @@ func (atsync *ATProtoSynchronizer) StartFirehoseRetry(ctx context.Context) error go atsync.handleCommitEventOps(ctx, evt) return nil }, + RepoIdentity: func(evt *comatproto.SyncSubscribeRepos_Identity) error { + go atsync.handleIdentityEventOps(ctx, evt) + return nil + }, Error: func(evt *events.ErrorFrame) error { log.Error(ctx, "firehose error", "err", evt.Error, "message", evt.Message) cancel() @@ -306,3 +311,25 @@ func (atsync *ATProtoSynchronizer) handleCommitEventOps(ctx context.Context, evt } } } + +func (atsync *ATProtoSynchronizer) handleIdentityEventOps(ctx context.Context, evt *comatproto.SyncSubscribeRepos_Identity) { + handle := "" + if evt.Handle != nil { + handle = *evt.Handle + } + ctx = log.WithLogValues(ctx, "event", "identity", "did", evt.Did, "handle", handle, "func", "handleIdentityEventOps") + r, err := atsync.Model.GetRepo(evt.Did) + if err != nil { + log.Error(ctx, "failed to get repo", "err", err) + return + } + if r == nil { + log.Debug(ctx, "no repo found for identity", "did", evt.Did) + return + } + _, err = atsync.RefreshIdentity(ctx, evt.Did) + if err != nil { + log.Error(ctx, "failed to refresh ident", "err", err) + return + } +} diff --git a/pkg/atproto/handle_test.go b/pkg/atproto/handle_test.go index 1e21643a3..4ac391a88 100644 --- a/pkg/atproto/handle_test.go +++ b/pkg/atproto/handle_test.go @@ -14,6 +14,7 @@ import ( "stream.place/streamplace/pkg/bus" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/devenv" + "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/model" "stream.place/streamplace/pkg/statedb" "stream.place/streamplace/pkg/streamplace" @@ -86,6 +87,7 @@ func TestHandleChange(t *testing.T) { require.NoError(t, err) require.Equal(t, user.Handle, message.Author.Handle) + log.Log(ctx, "updating handle", "handle", "new-handle.test") err = comatproto.IdentityUpdateHandle(context.Background(), user.XRPC, &comatproto.IdentityUpdateHandle_Input{ Handle: "new-handle.test", }) diff --git a/pkg/atproto/labeler_firehose.go b/pkg/atproto/labeler_firehose.go index c99fc3d31..8a19ce7c1 100644 --- a/pkg/atproto/labeler_firehose.go +++ b/pkg/atproto/labeler_firehose.go @@ -54,7 +54,7 @@ func (atsync *ATProtoSynchronizer) StartLabelerFirehose(ctx context.Context, did func (atsync *ATProtoSynchronizer) StartLabelerFirehoseRetry(ctx context.Context, did string) error { ctx = log.WithLogValues(ctx, "func", "StartLabelerFirehose") - ident, err := atsync.resolveIdent(ctx, did) + ident, err := atsync.resolveIdent(ctx, did, true) if err != nil { return fmt.Errorf("failed to resolve DID %s: %w", did, err) } diff --git a/pkg/config/config.go b/pkg/config/config.go index e63566633..b211a1165 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -250,7 +250,7 @@ var GormLogger = slogGorm.New( slogGorm.WithHandler(tint.NewHandler(os.Stderr, &tint.Options{ TimeFormat: time.RFC3339, })), - slogGorm.WithTraceAll(), + // slogGorm.WithTraceAll(), ) func (cli *CLI) Parse(fs *flag.FlagSet, args []string) error { diff --git a/pkg/devenv/devenv.go b/pkg/devenv/devenv.go index ed6abaca0..3d7e9632d 100644 --- a/pkg/devenv/devenv.go +++ b/pkg/devenv/devenv.go @@ -182,6 +182,5 @@ func (d *DevEnv) TestDirectory() identity.Directory { TryAuthoritativeDNS: true, SkipDNSDomainSuffixes: []string{".bsky.social"}, } - cached := identity.NewCacheDirectory(&base, 250_000, time.Hour*24, time.Minute*2, time.Minute*5) - return &cached + return &base } diff --git a/pkg/spxrpc/com_atproto_identity.go b/pkg/spxrpc/com_atproto_identity.go index 45c0b2ce6..77bd92166 100644 --- a/pkg/spxrpc/com_atproto_identity.go +++ b/pkg/spxrpc/com_atproto_identity.go @@ -14,3 +14,15 @@ func (s *Server) handleComAtprotoIdentityResolveHandle(ctx context.Context, hand } return &comatprototypes.IdentityResolveHandle_Output{Did: did}, nil } + +func (s *Server) handleComAtprotoIdentityRefreshIdentity(ctx context.Context, body *comatprototypes.IdentityRefreshIdentity_Input) (*comatprototypes.IdentityDefs_IdentityInfo, error) { + ident, err := s.ATSync.RefreshIdentity(ctx, body.Identifier) + if err != nil { + return nil, err + } + return &comatprototypes.IdentityDefs_IdentityInfo{ + Did: ident.DID.String(), + Handle: ident.Handle.String(), + DidDoc: ident.DIDDocument(), + }, nil +} diff --git a/pkg/spxrpc/stubs.go b/pkg/spxrpc/stubs.go index c18676d6e..d83a92dc7 100644 --- a/pkg/spxrpc/stubs.go +++ b/pkg/spxrpc/stubs.go @@ -62,6 +62,7 @@ func (s *Server) RegisterHandlersChatBsky(e *echo.Echo) error { } func (s *Server) RegisterHandlersComAtproto(e *echo.Echo) error { + e.POST("/xrpc/com.atproto.identity.refreshIdentity", s.HandleComAtprotoIdentityRefreshIdentity) e.GET("/xrpc/com.atproto.identity.resolveHandle", s.HandleComAtprotoIdentityResolveHandle) e.POST("/xrpc/com.atproto.moderation.createReport", s.HandleComAtprotoModerationCreateReport) e.GET("/xrpc/com.atproto.repo.describeRepo", s.HandleComAtprotoRepoDescribeRepo) @@ -74,6 +75,24 @@ func (s *Server) RegisterHandlersComAtproto(e *echo.Echo) error { return nil } +func (s *Server) HandleComAtprotoIdentityRefreshIdentity(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandleComAtprotoIdentityRefreshIdentity") + defer span.End() + + var body comatprototypes.IdentityRefreshIdentity_Input + if err := c.Bind(&body); err != nil { + return err + } + var out *comatprototypes.IdentityDefs_IdentityInfo + var handleErr error + // func (s *Server) handleComAtprotoIdentityRefreshIdentity(ctx context.Context,body *comatprototypes.IdentityRefreshIdentity_Input) (*comatprototypes.IdentityDefs_IdentityInfo, error) + out, handleErr = s.handleComAtprotoIdentityRefreshIdentity(ctx, &body) + if handleErr != nil { + return handleErr + } + return c.JSON(200, out) +} + func (s *Server) HandleComAtprotoIdentityResolveHandle(c echo.Context) error { ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandleComAtprotoIdentityResolveHandle") defer span.End()