diff --git a/atproto/lexicon/cmd/lextool/main.go b/atproto/lexicon/cmd/lextool/main.go index 85f37936..1718f50f 100644 --- a/atproto/lexicon/cmd/lextool/main.go +++ b/atproto/lexicon/cmd/lextool/main.go @@ -38,6 +38,11 @@ func main() { Usage: "subscribe to a firehose, validate every known record against catalog", Action: runValidateFirehose, }, + &cli.Command{ + Name: "resolve", + Usage: "resolves an NSID to a lexicon schema", + Action: runResolve, + }, } h := slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelDebug}) slog.SetDefault(slog.New(h)) @@ -88,3 +93,23 @@ func runLoadDirectory(cctx *cli.Context) error { fmt.Println("success!") return nil } + +func runResolve(cctx *cli.Context) error { + ref := cctx.Args().First() + if ref == "" { + return fmt.Errorf("need to provide NSID as an argument") + } + + c := lexicon.NewResolvingCatalog() + schema, err := c.Resolve(ref) + if err != nil { + return err + } + + out, err := json.MarshalIndent(schema, "", " ") + if err != nil { + return err + } + fmt.Println(string(out)) + return nil +} diff --git a/atproto/lexicon/repogetRecord.go b/atproto/lexicon/repogetRecord.go new file mode 100644 index 00000000..3e4e7f83 --- /dev/null +++ b/atproto/lexicon/repogetRecord.go @@ -0,0 +1,42 @@ +// Copied from indigo:api/atproto/repolistRecords.go + +package lexicon + +// schema: com.atproto.repo.getRecord + +import ( + "context" + "encoding/json" + + "github.com/bluesky-social/indigo/xrpc" +) + +// RepoGetRecord_Output is the output of a com.atproto.repo.getRecord call. +type RepoGetRecord_Output struct { + Cid *string `json:"cid,omitempty" cborgen:"cid,omitempty"` + Uri string `json:"uri" cborgen:"uri"` + // NOTE: changed from lex decoder to json.RawMessage + Value *json.RawMessage `json:"value" cborgen:"value"` +} + +// RepoGetRecord calls the XRPC method "com.atproto.repo.getRecord". +// +// cid: The CID of the version of the record. If not specified, then return the most recent version. +// collection: The NSID of the record collection. +// repo: The handle or DID of the repo. +// rkey: The Record Key. +func RepoGetRecord(ctx context.Context, c *xrpc.Client, cid string, collection string, repo string, rkey string) (*RepoGetRecord_Output, error) { + var out RepoGetRecord_Output + + params := map[string]interface{}{ + "cid": cid, + "collection": collection, + "repo": repo, + "rkey": rkey, + } + if err := c.Do(ctx, xrpc.Query, "", "com.atproto.repo.getRecord", params, nil, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/atproto/lexicon/resolving_catalog.go b/atproto/lexicon/resolving_catalog.go new file mode 100644 index 00000000..b3194532 --- /dev/null +++ b/atproto/lexicon/resolving_catalog.go @@ -0,0 +1,149 @@ +package lexicon + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "log/slog" + "net" + "strings" + "time" + + "github.com/bluesky-social/indigo/atproto/data" + "github.com/bluesky-social/indigo/atproto/identity" + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/bluesky-social/indigo/xrpc" +) + +// Catalog which supplements an in-memory BaseCatalog with live resolution from the network +type ResolvingCatalog struct { + Base BaseCatalog + Resolver net.Resolver + Directory identity.Directory +} + +// TODO: maybe this should take a base catalog as an arg? +func NewResolvingCatalog() ResolvingCatalog { + return ResolvingCatalog{ + Base: NewBaseCatalog(), + Resolver: net.Resolver{ + Dial: func(ctx context.Context, network, address string) (net.Conn, error) { + d := net.Dialer{Timeout: time.Second * 5} + return d.DialContext(ctx, network, address) + }, + }, + Directory: identity.DefaultDirectory(), + } +} + +func (rc *ResolvingCatalog) Resolve(ref string) (*Schema, error) { + if ref == "" { + return nil, fmt.Errorf("tried to resolve empty string name") + } + + existing, err := rc.Base.Resolve(ref) + if nil == err { + return existing, nil + } + + // TODO: was the idea to split the API and have "ensure" as a method? + ctx := context.Background() + + // TODO: split on '#' + nsid, err := syntax.ParseNSID(ref) + if err != nil { + return nil, err + } + + did, err := rc.ResolveNSID(ctx, nsid) + if err != nil { + return nil, err + } + slog.Info("resolved NSID", "nsid", nsid, "did", did) + + ident, err := rc.Directory.LookupDID(ctx, did) + if err != nil { + return nil, err + } + + aturi := syntax.ATURI(fmt.Sprintf("at://%s/com.atproto.lexicon.record/%s", did, nsid)) + record, err := fetchRecord(ctx, *ident, aturi) + if err != nil { + return nil, err + } + + recordJSON, err := json.Marshal(record) + if err != nil { + return nil, err + } + + var sf SchemaFile + if err = json.Unmarshal(recordJSON, &sf); err != nil { + return nil, err + } + if err = rc.Base.AddSchemaFile(sf); err != nil { + return nil, err + } + + return rc.Base.Resolve(ref) +} + +var ( + ErrResolutionFailed = fmt.Errorf("NSID resolution mechanism failed") + ErrNotFound = fmt.Errorf("NSID not associated with a DID") +) + +func parseTXTResp(res []string) (syntax.DID, error) { + for _, s := range res { + if strings.HasPrefix(s, "did=") { + parts := strings.SplitN(s, "=", 2) + did, err := syntax.ParseDID(parts[1]) + if err != nil { + return "", fmt.Errorf("%w: invalid DID in handle DNS record: %w", ErrResolutionFailed, err) + } + return did, nil + } + } + return "", ErrNotFound +} + +// resolves an NSID to a DID, using lexicon DNS TXT record +func (rc *ResolvingCatalog) ResolveNSID(ctx context.Context, nsid syntax.NSID) (syntax.DID, error) { + + domain := nsid.Authority() + res, err := rc.Resolver.LookupTXT(ctx, "_lexicon."+domain) + // check for NXDOMAIN + var dnsErr *net.DNSError + if errors.As(err, &dnsErr) { + if dnsErr.IsNotFound { + return "", ErrNotFound + } + } + if err != nil { + return "", fmt.Errorf("%w: %w", ErrResolutionFailed, err) + } + return parseTXTResp(res) +} + +func fetchRecord(ctx context.Context, ident identity.Identity, aturi syntax.ATURI) (any, error) { + + slog.Debug("fetching record", "did", ident.DID.String(), "collection", aturi.Collection().String(), "rkey", aturi.RecordKey().String()) + xrpcc := xrpc.Client{ + Host: ident.PDSEndpoint(), + } + resp, err := RepoGetRecord(ctx, &xrpcc, "", aturi.Collection().String(), ident.DID.String(), aturi.RecordKey().String()) + if err != nil { + return nil, err + } + + if nil == resp.Value { + return nil, fmt.Errorf("empty record in response") + } + record, err := data.UnmarshalJSON(*resp.Value) + if err != nil { + return nil, fmt.Errorf("fetched record was invalid data: %w", err) + } + + return record, nil +} -- 2.51.2 From 99b97be52b1082ec054e08c0e4c5f90fc95f13d6 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Sun, 22 Dec 2024 20:13:14 -0800 Subject: [PATCH 02/13] switch to agnostic package --- atproto/lexicon/repogetRecord.go | 42 ---------------------------- atproto/lexicon/resolving_catalog.go | 3 +- 2 files changed, 2 insertions(+), 43 deletions(-) delete mode 100644 atproto/lexicon/repogetRecord.go diff --git a/atproto/lexicon/repogetRecord.go b/atproto/lexicon/repogetRecord.go deleted file mode 100644 index 3e4e7f83..00000000 --- a/atproto/lexicon/repogetRecord.go +++ /dev/null @@ -1,42 +0,0 @@ -// Copied from indigo:api/atproto/repolistRecords.go - -package lexicon - -// schema: com.atproto.repo.getRecord - -import ( - "context" - "encoding/json" - - "github.com/bluesky-social/indigo/xrpc" -) - -// RepoGetRecord_Output is the output of a com.atproto.repo.getRecord call. -type RepoGetRecord_Output struct { - Cid *string `json:"cid,omitempty" cborgen:"cid,omitempty"` - Uri string `json:"uri" cborgen:"uri"` - // NOTE: changed from lex decoder to json.RawMessage - Value *json.RawMessage `json:"value" cborgen:"value"` -} - -// RepoGetRecord calls the XRPC method "com.atproto.repo.getRecord". -// -// cid: The CID of the version of the record. If not specified, then return the most recent version. -// collection: The NSID of the record collection. -// repo: The handle or DID of the repo. -// rkey: The Record Key. -func RepoGetRecord(ctx context.Context, c *xrpc.Client, cid string, collection string, repo string, rkey string) (*RepoGetRecord_Output, error) { - var out RepoGetRecord_Output - - params := map[string]interface{}{ - "cid": cid, - "collection": collection, - "repo": repo, - "rkey": rkey, - } - if err := c.Do(ctx, xrpc.Query, "", "com.atproto.repo.getRecord", params, nil, &out); err != nil { - return nil, err - } - - return &out, nil -} diff --git a/atproto/lexicon/resolving_catalog.go b/atproto/lexicon/resolving_catalog.go index b3194532..81710e03 100644 --- a/atproto/lexicon/resolving_catalog.go +++ b/atproto/lexicon/resolving_catalog.go @@ -14,6 +14,7 @@ import ( "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" + "github.com/bluesky-social/indigo/api/agnostic" ) // Catalog which supplements an in-memory BaseCatalog with live resolution from the network @@ -132,7 +133,7 @@ func fetchRecord(ctx context.Context, ident identity.Identity, aturi syntax.ATUR xrpcc := xrpc.Client{ Host: ident.PDSEndpoint(), } - resp, err := RepoGetRecord(ctx, &xrpcc, "", aturi.Collection().String(), ident.DID.String(), aturi.RecordKey().String()) + resp, err := agnostic.RepoGetRecord(ctx, &xrpcc, "", aturi.Collection().String(), ident.DID.String(), aturi.RecordKey().String()) if err != nil { return nil, err } -- 2.51.2 From ad007530d4415bac6530e5975cfd4bc36f165b45 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Sun, 22 Dec 2024 20:29:12 -0800 Subject: [PATCH 03/13] add NSID resolution to BaseDirectory --- atproto/identity/live_test.go | 14 ++++++++++++++ atproto/identity/nsid.go | 36 +++++++++++++++++++++++++++++++++++ 2 files changed, 50 insertions(+) create mode 100644 atproto/identity/nsid.go diff --git a/atproto/identity/live_test.go b/atproto/identity/live_test.go index 22d5b8c0..c763bf09 100644 --- a/atproto/identity/live_test.go +++ b/atproto/identity/live_test.go @@ -138,3 +138,17 @@ func TestFallbackDNS(t *testing.T) { assert.Error(err) assert.ErrorIs(err, ErrHandleResolutionFailed) } + +func TestResolveNSID(t *testing.T) { + t.Skip("TODO: skipping live network test") + assert := assert.New(t) + ctx := context.Background() + + dir := BaseDirectory{} + // NOTE: this was a very short temporary NSID when rkey restriction was short + nsid := syntax.NSID("net.bnewbold.m") + did, err := dir.ResolveNSID(ctx, nsid) + + assert.NoError(err) + assert.Equal(did, syntax.DID("did:plc:nhxcyu4ewwhl5pqil4dotqjo")) +} diff --git a/atproto/identity/nsid.go b/atproto/identity/nsid.go new file mode 100644 index 00000000..567d0cf6 --- /dev/null +++ b/atproto/identity/nsid.go @@ -0,0 +1,36 @@ +package identity + +import ( + "context" + "errors" + "fmt" + "net" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + +var ( + ErrNSIDResolutionFailed = fmt.Errorf("NSID resolution mechanism failed") + ErrNSIDNotFound = fmt.Errorf("NSID not associated with a DID") +) + +// Resolves an NSID to a DID, as used for Lexicon resolution (using "_lexicon" DNS TXT record) +func (d *BaseDirectory) ResolveNSID(ctx context.Context, nsid syntax.NSID) (syntax.DID, error) { + + domain := nsid.Authority() + res, err := d.Resolver.LookupTXT(ctx, "_lexicon."+domain) + + // check for NXDOMAIN + var dnsErr *net.DNSError + if errors.As(err, &dnsErr) { + if dnsErr.IsNotFound { + return "", ErrNSIDNotFound + } + } + + if err != nil { + return "", fmt.Errorf("%w: %w", ErrNSIDResolutionFailed, err) + } + + return parseTXTResp(res) +} -- 2.51.2 From 292a4c5faed55c659d5b2143c3e9c9d3be7235e2 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Sun, 22 Dec 2024 21:19:15 -0800 Subject: [PATCH 04/13] refactor lex resolution --- atproto/lexicon/resolve.go | 62 +++++++++++++++++++ atproto/lexicon/resolving_catalog.go | 90 ++-------------------------- 2 files changed, 66 insertions(+), 86 deletions(-) create mode 100644 atproto/lexicon/resolve.go diff --git a/atproto/lexicon/resolve.go b/atproto/lexicon/resolve.go new file mode 100644 index 00000000..50acc3b8 --- /dev/null +++ b/atproto/lexicon/resolve.go @@ -0,0 +1,62 @@ +package lexicon + +import ( + "context" + "fmt" + "log/slog" + + "github.com/bluesky-social/indigo/atproto/data" + "github.com/bluesky-social/indigo/atproto/identity" + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/bluesky-social/indigo/xrpc" + "github.com/bluesky-social/indigo/api/agnostic" +) + +// Low-level routine for resolving an NSID to full Lexicon data record (as stored in a repository). +// +// The current implementation uses a naive 'getRepo' fetch to the relevant PDS instance, without validating MST proof chain. +// +// Calling code should usually use ResolvingCatalog, which handles basing caching. +func ResolveLexiconData(ctx context.Context, dir identity.Directory, nsid syntax.NSID) (map[string]any, error) { + + baseDir := identity.BaseDirectory{} + did, err := baseDir.ResolveNSID(ctx, nsid) + if err != nil { + return nil, err + } + slog.Info("resolved NSID", "nsid", nsid, "did", did) + + ident, err := dir.LookupDID(ctx, did) + if err != nil { + return nil, err + } + + aturi := syntax.ATURI(fmt.Sprintf("at://%s/com.atproto.lexicon.schema/%s", did, nsid)) + record, err := fetchRecord(ctx, *ident, aturi) + if err != nil { + return nil, err + } + return record, nil +} + +func fetchRecord(ctx context.Context, ident identity.Identity, aturi syntax.ATURI) (map[string]any, error) { + + slog.Debug("fetching record", "did", ident.DID.String(), "collection", aturi.Collection().String(), "rkey", aturi.RecordKey().String()) + xrpcc := xrpc.Client{ + Host: ident.PDSEndpoint(), + } + resp, err := agnostic.RepoGetRecord(ctx, &xrpcc, "", aturi.Collection().String(), ident.DID.String(), aturi.RecordKey().String()) + if err != nil { + return nil, err + } + + if nil == resp.Value { + return nil, fmt.Errorf("empty record in response") + } + record, err := data.UnmarshalJSON(*resp.Value) + if err != nil { + return nil, fmt.Errorf("fetched record was invalid data: %w", err) + } + + return record, nil +} diff --git a/atproto/lexicon/resolving_catalog.go b/atproto/lexicon/resolving_catalog.go index 81710e03..364834b5 100644 --- a/atproto/lexicon/resolving_catalog.go +++ b/atproto/lexicon/resolving_catalog.go @@ -3,18 +3,12 @@ package lexicon import ( "context" "encoding/json" - "errors" "fmt" - "log/slog" "net" - "strings" "time" - "github.com/bluesky-social/indigo/atproto/data" "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" - "github.com/bluesky-social/indigo/xrpc" - "github.com/bluesky-social/indigo/api/agnostic" ) // Catalog which supplements an in-memory BaseCatalog with live resolution from the network @@ -39,37 +33,20 @@ func NewResolvingCatalog() ResolvingCatalog { } func (rc *ResolvingCatalog) Resolve(ref string) (*Schema, error) { + // TODO: not passed through! + ctx := context.Background() + if ref == "" { return nil, fmt.Errorf("tried to resolve empty string name") } - existing, err := rc.Base.Resolve(ref) - if nil == err { - return existing, nil - } - - // TODO: was the idea to split the API and have "ensure" as a method? - ctx := context.Background() - // TODO: split on '#' nsid, err := syntax.ParseNSID(ref) if err != nil { return nil, err } - did, err := rc.ResolveNSID(ctx, nsid) - if err != nil { - return nil, err - } - slog.Info("resolved NSID", "nsid", nsid, "did", did) - - ident, err := rc.Directory.LookupDID(ctx, did) - if err != nil { - return nil, err - } - - aturi := syntax.ATURI(fmt.Sprintf("at://%s/com.atproto.lexicon.record/%s", did, nsid)) - record, err := fetchRecord(ctx, *ident, aturi) + record, err := ResolveLexiconData(ctx, rc.Directory, nsid) if err != nil { return nil, err } @@ -89,62 +66,3 @@ func (rc *ResolvingCatalog) Resolve(ref string) (*Schema, error) { return rc.Base.Resolve(ref) } - -var ( - ErrResolutionFailed = fmt.Errorf("NSID resolution mechanism failed") - ErrNotFound = fmt.Errorf("NSID not associated with a DID") -) - -func parseTXTResp(res []string) (syntax.DID, error) { - for _, s := range res { - if strings.HasPrefix(s, "did=") { - parts := strings.SplitN(s, "=", 2) - did, err := syntax.ParseDID(parts[1]) - if err != nil { - return "", fmt.Errorf("%w: invalid DID in handle DNS record: %w", ErrResolutionFailed, err) - } - return did, nil - } - } - return "", ErrNotFound -} - -// resolves an NSID to a DID, using lexicon DNS TXT record -func (rc *ResolvingCatalog) ResolveNSID(ctx context.Context, nsid syntax.NSID) (syntax.DID, error) { - - domain := nsid.Authority() - res, err := rc.Resolver.LookupTXT(ctx, "_lexicon."+domain) - // check for NXDOMAIN - var dnsErr *net.DNSError - if errors.As(err, &dnsErr) { - if dnsErr.IsNotFound { - return "", ErrNotFound - } - } - if err != nil { - return "", fmt.Errorf("%w: %w", ErrResolutionFailed, err) - } - return parseTXTResp(res) -} - -func fetchRecord(ctx context.Context, ident identity.Identity, aturi syntax.ATURI) (any, error) { - - slog.Debug("fetching record", "did", ident.DID.String(), "collection", aturi.Collection().String(), "rkey", aturi.RecordKey().String()) - xrpcc := xrpc.Client{ - Host: ident.PDSEndpoint(), - } - resp, err := agnostic.RepoGetRecord(ctx, &xrpcc, "", aturi.Collection().String(), ident.DID.String(), aturi.RecordKey().String()) - if err != nil { - return nil, err - } - - if nil == resp.Value { - return nil, fmt.Errorf("empty record in response") - } - record, err := data.UnmarshalJSON(*resp.Value) - if err != nil { - return nil, fmt.Errorf("fetched record was invalid data: %w", err) - } - - return record, nil -} -- 2.51.2 From b7232ecd5fa5f0a98946104f9d8954417a33b733 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Sun, 22 Dec 2024 21:27:03 -0800 Subject: [PATCH 05/13] make fmt --- atproto/lexicon/resolve.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/atproto/lexicon/resolve.go b/atproto/lexicon/resolve.go index 50acc3b8..d2ea8e48 100644 --- a/atproto/lexicon/resolve.go +++ b/atproto/lexicon/resolve.go @@ -5,11 +5,11 @@ import ( "fmt" "log/slog" + "github.com/bluesky-social/indigo/api/agnostic" "github.com/bluesky-social/indigo/atproto/data" "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" - "github.com/bluesky-social/indigo/api/agnostic" ) // Low-level routine for resolving an NSID to full Lexicon data record (as stored in a repository). -- 2.51.2 From 8aa9f113902ff0b58b9c700f68496acbfa93961e Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Mon, 23 Dec 2024 16:46:02 -0800 Subject: [PATCH 06/13] lextool: remove unimplemented firehose validator Will add this to 'goat' instead --- atproto/lexicon/cmd/lextool/main.go | 5 ----- atproto/lexicon/cmd/lextool/net.go | 15 --------------- 2 files changed, 20 deletions(-) diff --git a/atproto/lexicon/cmd/lextool/main.go b/atproto/lexicon/cmd/lextool/main.go index 1718f50f..9ea8d592 100644 --- a/atproto/lexicon/cmd/lextool/main.go +++ b/atproto/lexicon/cmd/lextool/main.go @@ -33,11 +33,6 @@ func main() { Usage: "fetch from network, validate against catalog", Action: runValidateRecord, }, - &cli.Command{ - Name: "validate-firehose", - Usage: "subscribe to a firehose, validate every known record against catalog", - Action: runValidateFirehose, - }, &cli.Command{ Name: "resolve", Usage: "resolves an NSID to a lexicon schema", diff --git a/atproto/lexicon/cmd/lextool/net.go b/atproto/lexicon/cmd/lextool/net.go index 1600aaee..99a85365 100644 --- a/atproto/lexicon/cmd/lextool/net.go +++ b/atproto/lexicon/cmd/lextool/net.go @@ -78,18 +78,3 @@ func runValidateRecord(cctx *cli.Context) error { fmt.Println("success!") return nil } - -func runValidateFirehose(cctx *cli.Context) error { - p := cctx.Args().First() - if p == "" { - return fmt.Errorf("need to provide directory path as an argument") - } - - cat := lexicon.NewBaseCatalog() - err := cat.LoadDirectory(p) - if err != nil { - return err - } - - return fmt.Errorf("UNIMPLEMENTED") -} -- 2.51.2 From 92a0ce3316a5dba8016911b752e6fb481cf1f4ab Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Mon, 23 Dec 2024 16:46:36 -0800 Subject: [PATCH 07/13] less verbose logging --- atproto/lexicon/resolve.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/atproto/lexicon/resolve.go b/atproto/lexicon/resolve.go index d2ea8e48..980d4978 100644 --- a/atproto/lexicon/resolve.go +++ b/atproto/lexicon/resolve.go @@ -24,7 +24,7 @@ func ResolveLexiconData(ctx context.Context, dir identity.Directory, nsid syntax if err != nil { return nil, err } - slog.Info("resolved NSID", "nsid", nsid, "did", did) + slog.Debug("resolved NSID", "nsid", nsid, "did", did) ident, err := dir.LookupDID(ctx, did) if err != nil { -- 2.51.2 From f38e4801fe9aab7684ba0de173d90ee3952337d1 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Mon, 23 Dec 2024 17:53:23 -0800 Subject: [PATCH 08/13] remove 'revision' from lexicon package --- atproto/lexicon/catalog.go | 5 ++--- atproto/lexicon/language.go | 1 - atproto/lexicon/lexicon.go | 5 ++--- 3 files changed, 4 insertions(+), 7 deletions(-) diff --git a/atproto/lexicon/catalog.go b/atproto/lexicon/catalog.go index 7797fad9..524ef81a 100644 --- a/atproto/lexicon/catalog.go +++ b/atproto/lexicon/catalog.go @@ -72,9 +72,8 @@ func (c *BaseCatalog) AddSchemaFile(sf SchemaFile) error { return err } s := Schema{ - ID: name, - Revision: sf.Revision, - Def: def.Inner, + ID: name, + Def: def.Inner, } c.schemas[name] = s } diff --git a/atproto/lexicon/language.go b/atproto/lexicon/language.go index 1eb5bdf0..b2243076 100644 --- a/atproto/lexicon/language.go +++ b/atproto/lexicon/language.go @@ -16,7 +16,6 @@ import ( type SchemaFile struct { Lexicon int `json:"lexicon,const=1"` ID string `json:"id"` - Revision *int `json:"revision,omitempty"` Description *string `json:"description,omitempty"` Defs map[string]SchemaDef `json:"defs"` } diff --git a/atproto/lexicon/lexicon.go b/atproto/lexicon/lexicon.go index 840d2ae7..df0909da 100644 --- a/atproto/lexicon/lexicon.go +++ b/atproto/lexicon/lexicon.go @@ -22,9 +22,8 @@ var LenientMode ValidateFlags = AllowLegacyBlob | AllowLenientDatetime // Represents a Lexicon schema definition type Schema struct { - ID string - Revision *int - Def any + ID string + Def any } // Checks Lexicon schema (fetched from the catalog) for the given record, with optional flags tweaking default validation rules. -- 2.51.2 From 98c6f96e10cf9ac86696955ee215ece2f8f85966 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Mon, 23 Dec 2024 17:54:15 -0800 Subject: [PATCH 09/13] resolving catalog: remove unused net.Resolver, validate some fields --- atproto/lexicon/resolve.go | 2 +- atproto/lexicon/resolving_catalog.go | 34 ++++++++++++++++------------ 2 files changed, 21 insertions(+), 15 deletions(-) diff --git a/atproto/lexicon/resolve.go b/atproto/lexicon/resolve.go index 980d4978..1ea380dd 100644 --- a/atproto/lexicon/resolve.go +++ b/atproto/lexicon/resolve.go @@ -16,7 +16,7 @@ import ( // // The current implementation uses a naive 'getRepo' fetch to the relevant PDS instance, without validating MST proof chain. // -// Calling code should usually use ResolvingCatalog, which handles basing caching. +// Calling code should usually use ResolvingCatalog, which handles basic caching and validation of the Lexicon language itself. func ResolveLexiconData(ctx context.Context, dir identity.Directory, nsid syntax.NSID) (map[string]any, error) { baseDir := identity.BaseDirectory{} diff --git a/atproto/lexicon/resolving_catalog.go b/atproto/lexicon/resolving_catalog.go index 364834b5..9ae6d51d 100644 --- a/atproto/lexicon/resolving_catalog.go +++ b/atproto/lexicon/resolving_catalog.go @@ -4,8 +4,7 @@ import ( "context" "encoding/json" "fmt" - "net" - "time" + "strings" "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" @@ -14,34 +13,33 @@ import ( // Catalog which supplements an in-memory BaseCatalog with live resolution from the network type ResolvingCatalog struct { Base BaseCatalog - Resolver net.Resolver Directory identity.Directory } -// TODO: maybe this should take a base catalog as an arg? func NewResolvingCatalog() ResolvingCatalog { return ResolvingCatalog{ - Base: NewBaseCatalog(), - Resolver: net.Resolver{ - Dial: func(ctx context.Context, network, address string) (net.Conn, error) { - d := net.Dialer{Timeout: time.Second * 5} - return d.DialContext(ctx, network, address) - }, - }, + Base: NewBaseCatalog(), Directory: identity.DefaultDirectory(), } } func (rc *ResolvingCatalog) Resolve(ref string) (*Schema, error) { - // TODO: not passed through! + // NOTE: not passed through! ctx := context.Background() if ref == "" { return nil, fmt.Errorf("tried to resolve empty string name") } - // TODO: split on '#' - nsid, err := syntax.ParseNSID(ref) + // first try existing catalog + schema, err := rc.Base.Resolve(ref) + if nil == err { // no error: found a hit + return schema, nil + } + + // split any ref from the end '#' + parts := strings.SplitN(ref, "#", 2) + nsid, err := syntax.ParseNSID(parts[0]) if err != nil { return nil, err } @@ -60,9 +58,17 @@ func (rc *ResolvingCatalog) Resolve(ref string) (*Schema, error) { if err = json.Unmarshal(recordJSON, &sf); err != nil { return nil, err } + + if sf.Lexicon != 1 { + return nil, fmt.Errorf("unsupported lexicon language version: %d", sf.Lexicon) + } + if sf.ID != nsid.String() { + return nil, fmt.Errorf("lexicon ID does not match NSID: %s != %s", sf.ID, nsid) + } if err = rc.Base.AddSchemaFile(sf); err != nil { return nil, err } + // re-resolving from the raw ref ensures that fragments are handled return rc.Base.Resolve(ref) } -- 2.51.2 From 6ac44d91383485d4430f4d6898e6bdf31a9105c1 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Mon, 23 Dec 2024 18:04:39 -0800 Subject: [PATCH 10/13] helper for loading schemas from an embed.FS --- atproto/lexicon/catalog.go | 72 ++++++++++++++++++++++---------------- 1 file changed, 42 insertions(+), 30 deletions(-) diff --git a/atproto/lexicon/catalog.go b/atproto/lexicon/catalog.go index 524ef81a..b057070f 100644 --- a/atproto/lexicon/catalog.go +++ b/atproto/lexicon/catalog.go @@ -1,6 +1,7 @@ package lexicon import ( + "embed" "encoding/json" "fmt" "io" @@ -46,6 +47,9 @@ func (c *BaseCatalog) Resolve(ref string) (*Schema, error) { // Inserts a schema loaded from a JSON file in to the catalog. func (c *BaseCatalog) AddSchemaFile(sf SchemaFile) error { + if sf.Lexicon != 1 { + return fmt.Errorf("unsupported lexicon language version: %d", sf.Lexicon) + } base := sf.ID for frag, def := range sf.Defs { if len(frag) == 0 || strings.Contains(frag, "#") || strings.Contains(frag, ".") { @@ -80,37 +84,45 @@ func (c *BaseCatalog) AddSchemaFile(sf SchemaFile) error { return nil } +// internal helper for loading file paths (either real filesystem or embed.FS) +func (c *BaseCatalog) addDirEntry(p string, d fs.DirEntry, err error) error { + if err != nil { + return err + } + if d.IsDir() { + return nil + } + if !strings.HasSuffix(p, ".json") { + return nil + } + slog.Debug("loading Lexicon schema file", "path", p) + f, err := os.Open(p) + if err != nil { + return err + } + defer func() { _ = f.Close() }() + + b, err := io.ReadAll(f) + if err != nil { + return err + } + + var sf SchemaFile + if err = json.Unmarshal(b, &sf); err != nil { + return err + } + if err = c.AddSchemaFile(sf); err != nil { + return err + } + return nil +} + // Recursively loads all '.json' files from a directory in to the catalog. func (c *BaseCatalog) LoadDirectory(dirPath string) error { - return filepath.WalkDir(dirPath, func(p string, d fs.DirEntry, err error) error { - if err != nil { - return err - } - if d.IsDir() { - return nil - } - if !strings.HasSuffix(p, ".json") { - return nil - } - slog.Debug("loading Lexicon schema file", "path", p) - f, err := os.Open(p) - if err != nil { - return err - } - defer func() { _ = f.Close() }() - - b, err := io.ReadAll(f) - if err != nil { - return err - } + return filepath.WalkDir(dirPath, c.addDirEntry) +} - var sf SchemaFile - if err = json.Unmarshal(b, &sf); err != nil { - return err - } - if err = c.AddSchemaFile(sf); err != nil { - return err - } - return nil - }) +// Recursively loads all '.json' files from an embed.FS +func (c *BaseCatalog) LoadEmbedFS(efs embed.FS) error { + return fs.WalkDir(efs, ".", c.addDirEntry) } -- 2.51.2 From 4498ba011b5bf19abee38a683df90f58481c0288 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Thu, 16 Jan 2025 18:30:27 -0800 Subject: [PATCH 11/13] refactor catalog recursive loading for embed fs --- atproto/lexicon/catalog.go | 74 +++++++++++++++++++++------------ atproto/lexicon/catalog_test.go | 41 ++++++++++++++++++ 2 files changed, 88 insertions(+), 27 deletions(-) create mode 100644 atproto/lexicon/catalog_test.go diff --git a/atproto/lexicon/catalog.go b/atproto/lexicon/catalog.go index b057070f..5101b613 100644 --- a/atproto/lexicon/catalog.go +++ b/atproto/lexicon/catalog.go @@ -84,34 +84,13 @@ func (c *BaseCatalog) AddSchemaFile(sf SchemaFile) error { return nil } -// internal helper for loading file paths (either real filesystem or embed.FS) -func (c *BaseCatalog) addDirEntry(p string, d fs.DirEntry, err error) error { - if err != nil { - return err - } - if d.IsDir() { - return nil - } - if !strings.HasSuffix(p, ".json") { - return nil - } - slog.Debug("loading Lexicon schema file", "path", p) - f, err := os.Open(p) - if err != nil { - return err - } - defer func() { _ = f.Close() }() - - b, err := io.ReadAll(f) - if err != nil { - return err - } - +// internal helper for loading JSON files from bytes +func (c *BaseCatalog) addSchemaFromBytes(b []byte) error { var sf SchemaFile - if err = json.Unmarshal(b, &sf); err != nil { + if err := json.Unmarshal(b, &sf); err != nil { return err } - if err = c.AddSchemaFile(sf); err != nil { + if err := c.AddSchemaFile(sf); err != nil { return err } return nil @@ -119,10 +98,51 @@ func (c *BaseCatalog) addDirEntry(p string, d fs.DirEntry, err error) error { // Recursively loads all '.json' files from a directory in to the catalog. func (c *BaseCatalog) LoadDirectory(dirPath string) error { - return filepath.WalkDir(dirPath, c.addDirEntry) + walkFunc := func(p string, d fs.DirEntry, err error) error { + if err != nil { + return err + } + if d.IsDir() { + return nil + } + if !strings.HasSuffix(p, ".json") { + return nil + } + slog.Debug("loading Lexicon schema file", "path", p) + f, err := os.Open(p) + if err != nil { + return err + } + defer func() { _ = f.Close() }() + + b, err := io.ReadAll(f) + if err != nil { + return err + } + return c.addSchemaFromBytes(b) + } + return filepath.WalkDir(dirPath, walkFunc) } // Recursively loads all '.json' files from an embed.FS func (c *BaseCatalog) LoadEmbedFS(efs embed.FS) error { - return fs.WalkDir(efs, ".", c.addDirEntry) + walkFunc := func(p string, d fs.DirEntry, err error) error { + if err != nil { + return err + } + if d.IsDir() { + return nil + } + if !strings.HasSuffix(p, ".json") { + return nil + } + + slog.Debug("loading embedded Lexicon schema file", "path", p) + b, err := efs.ReadFile(p) + if err != nil { + return err + } + return c.addSchemaFromBytes(b) + } + return fs.WalkDir(efs, ".", walkFunc) } diff --git a/atproto/lexicon/catalog_test.go b/atproto/lexicon/catalog_test.go new file mode 100644 index 00000000..b2420306 --- /dev/null +++ b/atproto/lexicon/catalog_test.go @@ -0,0 +1,41 @@ +package lexicon + +import ( + "embed" + "testing" + + "github.com/stretchr/testify/assert" +) + +//go:embed testdata/catalog +var embedDir embed.FS + +func TestEmbedCatalog(t *testing.T) { + assert := assert.New(t) + + cat := NewBaseCatalog() + + err := cat.LoadEmbedFS(embedDir) + assert.NoError(err) + + _, err = cat.Resolve("example.lexicon.query") + assert.NoError(err) + + _, err = cat.Resolve("example.lexicon.notThere") + assert.Error(err) +} + +func TestDirCatalog(t *testing.T) { + assert := assert.New(t) + + cat := NewBaseCatalog() + + err := cat.LoadDirectory("testdata/catalog") + assert.NoError(err) + + _, err = cat.Resolve("example.lexicon.query") + assert.NoError(err) + + _, err = cat.Resolve("example.lexicon.notThere") + assert.Error(err) +} -- 2.51.2 From ad991e41b4581e43cde245d5aca72df62639cb69 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Fri, 17 Jan 2025 12:16:47 -0800 Subject: [PATCH 12/13] refactor lexicon record fetching a bit --- atproto/lexicon/resolve.go | 43 +++++++++++++++++++++++++++++++------- 1 file changed, 35 insertions(+), 8 deletions(-) diff --git a/atproto/lexicon/resolve.go b/atproto/lexicon/resolve.go index 1ea380dd..c987093f 100644 --- a/atproto/lexicon/resolve.go +++ b/atproto/lexicon/resolve.go @@ -2,6 +2,7 @@ package lexicon import ( "context" + "encoding/json" "fmt" "log/slog" @@ -19,6 +20,36 @@ import ( // Calling code should usually use ResolvingCatalog, which handles basic caching and validation of the Lexicon language itself. func ResolveLexiconData(ctx context.Context, dir identity.Directory, nsid syntax.NSID) (map[string]any, error) { + record, err := resolveLexiconJSON(ctx, dir, nsid) + if err != nil { + return nil, err + } + + d, err := data.UnmarshalJSON(*record) + if err != nil { + return nil, fmt.Errorf("fetched Lexicon schema record was invalid: %w", err) + } + return d, nil +} + +// Low-level routine for resolving an NSID to `SchemaFile`. +// +// Same as `ResolveLexiconData`, but returns a parsed `SchemaFile` struct. +func ResolveLexiconSchemaFile(ctx context.Context, dir identity.Directory, nsid syntax.NSID) (*SchemaFile, error) { + record, err := resolveLexiconJSON(ctx, dir, nsid) + if err != nil { + return nil, err + } + + var sf SchemaFile + if err := json.Unmarshal(*record, &sf); err != nil { + return nil, fmt.Errorf("fetched Lexicon schema record was invalid: %w", err) + } + return &sf, nil +} + +// internal helper for fetching lexicon record as JSON bytes +func resolveLexiconJSON(ctx context.Context, dir identity.Directory, nsid syntax.NSID) (*json.RawMessage, error) { baseDir := identity.BaseDirectory{} did, err := baseDir.ResolveNSID(ctx, nsid) if err != nil { @@ -32,14 +63,14 @@ func ResolveLexiconData(ctx context.Context, dir identity.Directory, nsid syntax } aturi := syntax.ATURI(fmt.Sprintf("at://%s/com.atproto.lexicon.schema/%s", did, nsid)) - record, err := fetchRecord(ctx, *ident, aturi) + msg, err := fetchRecordJSON(ctx, *ident, aturi) if err != nil { return nil, err } - return record, nil + return msg, err } -func fetchRecord(ctx context.Context, ident identity.Identity, aturi syntax.ATURI) (map[string]any, error) { +func fetchRecordJSON(ctx context.Context, ident identity.Identity, aturi syntax.ATURI) (*json.RawMessage, error) { slog.Debug("fetching record", "did", ident.DID.String(), "collection", aturi.Collection().String(), "rkey", aturi.RecordKey().String()) xrpcc := xrpc.Client{ @@ -53,10 +84,6 @@ func fetchRecord(ctx context.Context, ident identity.Identity, aturi syntax.ATUR if nil == resp.Value { return nil, fmt.Errorf("empty record in response") } - record, err := data.UnmarshalJSON(*resp.Value) - if err != nil { - return nil, fmt.Errorf("fetched record was invalid data: %w", err) - } - return record, nil + return resp.Value, nil } -- 2.51.2 From 465bcb9fea31b5ae2b3a065a346d602ec9fd3335 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Fri, 17 Jan 2025 12:20:06 -0800 Subject: [PATCH 13/13] add comment about 'ref' argument --- atproto/lexicon/catalog.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/atproto/lexicon/catalog.go b/atproto/lexicon/catalog.go index 5101b613..69c45800 100644 --- a/atproto/lexicon/catalog.go +++ b/atproto/lexicon/catalog.go @@ -30,6 +30,9 @@ func NewBaseCatalog() BaseCatalog { } } +// Returns a scheman definition (`Schema` struct) for a Lexicon reference. +// +// A Lexicon ref string is an NSID with an optional #-separated fragment. If the fragment isn't specified, '#main' is used by default. func (c *BaseCatalog) Resolve(ref string) (*Schema, error) { if ref == "" { return nil, fmt.Errorf("tried to resolve empty string name")