diff --git a/atproto/netclient/cid.go b/atproto/heap/cid.go similarity index 95% rename from atproto/netclient/cid.go rename to atproto/heap/cid.go index a849ac9..b029e91 100644 --- a/atproto/netclient/cid.go +++ b/atproto/heap/cid.go @@ -1,4 +1,4 @@ -package netclient +package heap import ( "github.com/ipfs/go-cid" diff --git a/atproto/netclient/examples_test.go b/atproto/heap/examples_test.go similarity index 60% rename from atproto/netclient/examples_test.go rename to atproto/heap/examples_test.go index a438df9..f433ad7 100644 --- a/atproto/netclient/examples_test.go +++ b/atproto/heap/examples_test.go @@ -1,8 +1,9 @@ -package netclient +package heap import ( "bytes" "context" + "encoding/json" "fmt" "github.com/bluesky-social/indigo/atproto/repo" @@ -66,5 +67,42 @@ func ExampleNetClient_GetAccountStatus() { } fmt.Printf("active=%t status=%s\n", active, status) - // Output: active=true status= + // active=true status= +} + +func ExampleNetClient_GetRecordUnverified() { + + ctx := context.Background() + nc := NewNetClient() + did := syntax.DID("did:plc:ewvi7nxzyoun6zhxrhs64oiz") + collection := syntax.NSID("app.bsky.actor.profile") + rkey := syntax.RecordKey("self") + + raw, _, err := nc.GetRecordUnverified(ctx, did, collection, rkey) + if err != nil { + panic("failed to fetch record: " + err.Error()) + } + var record map[string]any + _ = json.Unmarshal(*raw, &record) + + fmt.Println(record["displayName"]) + // AT Protocol Developers +} + +func ExampleNetClient_GetRecord() { + + ctx := context.Background() + nc := NewNetClient() + did := syntax.DID("did:plc:ewvi7nxzyoun6zhxrhs64oiz") + collection := syntax.NSID("app.bsky.actor.profile") + rkey := syntax.RecordKey("self") + + var record map[string]any + _, err := nc.GetRecord(ctx, did, collection, rkey, &record) + if err != nil { + panic("failed to fetch record: " + err.Error()) + } + + fmt.Println(record["displayName"]) + // Output: AT Protocol Developers } diff --git a/atproto/netclient/netclient.go b/atproto/heap/netclient.go similarity index 99% rename from atproto/netclient/netclient.go rename to atproto/heap/netclient.go index cb5f3be..7da1c03 100644 --- a/atproto/netclient/netclient.go +++ b/atproto/heap/netclient.go @@ -1,4 +1,4 @@ -package netclient +package heap import ( "bytes" diff --git a/atproto/heap/record.go b/atproto/heap/record.go new file mode 100644 index 0000000..6709774 --- /dev/null +++ b/atproto/heap/record.go @@ -0,0 +1,156 @@ +package heap + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "log/slog" + "net/http" + + "github.com/bluesky-social/indigo/atproto/data" + "github.com/bluesky-social/indigo/atproto/repo" + "github.com/bluesky-social/indigo/atproto/syntax" +) + +type repoRecordResp struct { + URI string `json:"uri"` + CID syntax.CID `json:"cid"` + Value json.RawMessage `json:"value"` +} + +// Fetches record JSON using com.atproto.repo.getRecord, and returns record as [json.RawMessage] and the CID (as string). +func (nc *NetClient) GetRecordUnverified(ctx context.Context, did syntax.DID, collection syntax.NSID, rkey syntax.RecordKey) (*json.RawMessage, syntax.CID, error) { + ident, err := nc.Dir.LookupDID(ctx, did) + if err != nil { + return nil, "", err + } + host := ident.PDSEndpoint() + if host == "" { + return nil, "", fmt.Errorf("account has no PDS host registered: %s", did.String()) + } + // TODO: validate host + // TODO: DID escaping (?) + u := fmt.Sprintf("%s/xrpc/com.atproto.repo.getRecord?repo=%s&collection=%s&rkey=%s", host, did, collection, rkey) + + slog.Debug("fetching record JSON", "did", did, "url", u) + req, err := http.NewRequestWithContext(ctx, "GET", u, nil) + if err != nil { + return nil, "", err + } + if nc.UserAgent != "" { + req.Header.Set("User-Agent", nc.UserAgent) + } + req.Header.Set("Accept", "application/json") + + resp, err := nc.Client.Do(req) + if err != nil { + return nil, "", fmt.Errorf("fetching record JSON (%s): %w", did, err) + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return nil, "", fmt.Errorf("HTTP error fetching record JSON (%s): %d", did, resp.StatusCode) + } + + var rrr repoRecordResp + if err := json.NewDecoder(resp.Body).Decode(&rrr); err != nil { + return nil, "", fmt.Errorf("failed decoding account status response: %w", err) + } + + return &rrr.Value, rrr.CID, nil +} + +// Fetches a record "proof" using com.atproto.sync.getRecord. Verifies signature and merkel chain. Copies record content in out 'out' parameter. +// +// If out is nil, record data is not returned. If it is [bytes.Buffer], the record CBOR is copied in. Otherwise, the record is transformed to JSON and Unmarshalled in to provided output, which could be a pointer to a struct, [json.RawMessage], `map[string]any`, etc. +// +// TODO: this might not be fully validating MST tree and record CID hashes or encoding yet +func (nc *NetClient) GetRecord(ctx context.Context, did syntax.DID, collection syntax.NSID, rkey syntax.RecordKey, out any) (syntax.CID, error) { + // TODO: "GetRecordProof" variant, which just returns CAR as io.ReadCloser? + ident, err := nc.Dir.LookupDID(ctx, did) + if err != nil { + return "", err + } + pub, err := ident.PublicKey() + if err != nil { + return "", err + } + host := ident.PDSEndpoint() + if host == "" { + return "", fmt.Errorf("account has no PDS host registered: %s", did.String()) + } + // TODO: validate host + // TODO: DID escaping (?) + u := fmt.Sprintf("%s/xrpc/com.atproto.sync.getRecord?did=%s&collection=%s&rkey=%s", host, did, collection, rkey) + + slog.Debug("fetching record proof", "did", did, "url", u) + req, err := http.NewRequestWithContext(ctx, "GET", u, nil) + if err != nil { + return "", err + } + if nc.UserAgent != "" { + req.Header.Set("User-Agent", nc.UserAgent) + } + req.Header.Set("Accept", "application/vnd.ipld.car") + + resp, err := nc.Client.Do(req) + if err != nil { + return "", fmt.Errorf("fetching record proof (%s): %w", did, err) + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return "", fmt.Errorf("HTTP error fetching record proof (%s): %d", did, resp.StatusCode) + } + + // TODO: re-confirm if loading tree re-checks all CIDs; or if we need to re-compute the tree data CID + commit, rp, err := repo.LoadRepoFromCAR(ctx, resp.Body) + if err != nil { + return "", fmt.Errorf("failed to parse record proof CAR (%s): %w", did, err) + } + + // NOTE: LoadRepoFromCAR calls commit.VerifyStructure() internally + + if err := commit.VerifySignature(pub); err != nil { + return "", fmt.Errorf("failed to verify record proof signature (%s): %w", did, err) + } + + rbytes, rcid, err := rp.GetRecordBytes(ctx, collection, rkey) + if err != nil { + return "", fmt.Errorf("failed to read record from proof CAR (%s): %w", did, err) + } + cidStr := syntax.CID(rcid.String()) + + // TODO: `GetRecordBytes` does not currently verify record CID, but unpacking CAR file should have done that? but need to confirm CAR implementation does this + + // check that record CBOR is valid, even if we don't return it + rdata, err := data.UnmarshalCBOR(rbytes) + if err != nil { + return "", fmt.Errorf("failed to parse record CBOR (%s): %w", did, err) + } + + switch out := out.(type) { + case nil: + // if output isn't captured, bail out early + return cidStr, nil + case *bytes.Buffer: + // simply copy data over + out.Reset() + _, err := out.Write(rbytes) + if err != nil { + return "", err + } + return cidStr, nil + default: + // attempt to unmarshal from json + jsonBytes, err := json.Marshal(rdata) + if err != nil { + return "", err + } + if err := json.Unmarshal(jsonBytes, out); err != nil { + return "", fmt.Errorf("failed unmarhsaling record (%s): %w", did, err) + } + return cidStr, nil + } +}