(READ ONLY) Margin is an open annotation layer for the internet. Powered by the AT Protocol.
margin.at
extension web atproto comments
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407package xrpc
import ( "bytes" "context" "crypto/ecdsa" "crypto/rand" "crypto/sha256" "encoding/base64" "encoding/json" "fmt" "io" "net/http" "time"
"github.com/go-jose/go-jose/v4")
type Client struct { PDS string AccessToken string DPoPKey *ecdsa.PrivateKey DPoPNonce string}
func NewClient(pds, accessToken string, dpopKey *ecdsa.PrivateKey) *Client { return &Client{ PDS: pds, AccessToken: accessToken, DPoPKey: dpopKey, }}
func (c *Client) createDPoPProof(method, uri string) (string, error) { now := time.Now() jti := make([]byte, 16) if _, err := io.ReadFull(rand.Reader, jti); err != nil {
for i := range jti { jti[i] = byte(now.UnixNano() >> (i * 8)) } }
publicJWK := jose.JSONWebKey{ Key: &c.DPoPKey.PublicKey, Algorithm: string(jose.ES256), }
ath := "" if c.AccessToken != "" { hash := sha256.Sum256([]byte(c.AccessToken)) ath = base64.RawURLEncoding.EncodeToString(hash[:]) }
claims := map[string]interface{}{ "jti": base64.RawURLEncoding.EncodeToString(jti), "htm": method, "htu": uri, "iat": now.Unix(), "exp": now.Add(5 * time.Minute).Unix(), } if c.DPoPNonce != "" { claims["nonce"] = c.DPoPNonce } if ath != "" { claims["ath"] = ath }
signer, err := jose.NewSigner(jose.SigningKey{Algorithm: jose.ES256, Key: c.DPoPKey}, &jose.SignerOptions{ ExtraHeaders: map[jose.HeaderKey]interface{}{ "typ": "dpop+jwt", "jwk": publicJWK, }, }) if err != nil { return "", err }
claimsBytes, _ := json.Marshal(claims) sig, err := signer.Sign(claimsBytes) if err != nil { return "", err }
return sig.CompactSerialize()}
func (c *Client) Call(ctx context.Context, method, nsid string, input, output interface{}) error { url := fmt.Sprintf("%s/xrpc/%s", c.PDS, nsid)
maxRetries := 2 for i := 0; i < maxRetries; i++ { var reqBody io.Reader if input != nil {
data, err := json.Marshal(input) if err != nil { return err } reqBody = bytes.NewReader(data) }
req, err := http.NewRequestWithContext(ctx, method, url, reqBody) if err != nil { return err }
if input != nil { req.Header.Set("Content-Type", "application/json") }
dpopProof, err := c.createDPoPProof(method, url) if err != nil { return fmt.Errorf("failed to create DPoP proof: %w", err) }
req.Header.Set("Authorization", "DPoP "+c.AccessToken) req.Header.Set("DPoP", dpopProof)
resp, err := http.DefaultClient.Do(req) if err != nil { return err } defer resp.Body.Close()
if nonce := resp.Header.Get("DPoP-Nonce"); nonce != "" { c.DPoPNonce = nonce }
if resp.StatusCode < 400 { if output != nil { return json.NewDecoder(resp.Body).Decode(output) } return nil }
bodyBytes, _ := io.ReadAll(resp.Body) bodyStr := string(bodyBytes)
if resp.StatusCode == 401 && (bytes.Contains(bodyBytes, []byte("use_dpop_nonce")) || bytes.Contains(bodyBytes, []byte("UseDpopNonce"))) { continue }
return fmt.Errorf("XRPC error %d: %s", resp.StatusCode, bodyStr) }
return fmt.Errorf("XRPC failed after retries")}
type CreateRecordInput struct { Repo string `json:"repo"` Collection string `json:"collection"` RKey string `json:"rkey,omitempty"` Record interface{} `json:"record"`}
type CreateRecordOutput struct { URI string `json:"uri"` CID string `json:"cid"`}
func (c *Client) CreateRecord(ctx context.Context, repo, collection string, record interface{}) (*CreateRecordOutput, error) { input := CreateRecordInput{ Repo: repo, Collection: collection, Record: record, }
var output CreateRecordOutput err := c.Call(ctx, "POST", "com.atproto.repo.createRecord", input, &output) if err != nil { return nil, err }
return &output, nil}
type DeleteRecordInput struct { Repo string `json:"repo"` Collection string `json:"collection"` RKey string `json:"rkey"`}
func (c *Client) DeleteRecord(ctx context.Context, repo, collection, rkey string) error { input := DeleteRecordInput{ Repo: repo, Collection: collection, RKey: rkey, }
return c.Call(ctx, "POST", "com.atproto.repo.deleteRecord", input, nil)}
func (c *Client) DeleteRecordByURI(ctx context.Context, uri string) error { parsed, err := ParseATURI(uri) if err != nil { return err }
if parsed.Collection == "" || parsed.RKey == "" { return fmt.Errorf("invalid AT-URI: must include collection and rkey") }
return c.DeleteRecord(ctx, parsed.DID, parsed.Collection, parsed.RKey)}
type PutRecordInput struct { Repo string `json:"repo"` Collection string `json:"collection"` RKey string `json:"rkey"` Record interface{} `json:"record"`}
type PutRecordOutput struct { URI string `json:"uri"` CID string `json:"cid"`}
func (c *Client) PutRecord(ctx context.Context, repo, collection, rkey string, record interface{}) (*PutRecordOutput, error) { input := PutRecordInput{ Repo: repo, Collection: collection, RKey: rkey, Record: record, }
var output PutRecordOutput err := c.Call(ctx, "POST", "com.atproto.repo.putRecord", input, &output) if err != nil { return nil, err }
return &output, nil}
type GetRecordOutput struct { URI string `json:"uri"` CID string `json:"cid"` Value json.RawMessage `json:"value"`}
func (c *Client) GetRecord(ctx context.Context, repo, collection, rkey string) (*GetRecordOutput, error) { url := fmt.Sprintf("%s/xrpc/com.atproto.repo.getRecord?repo=%s&collection=%s&rkey=%s", c.PDS, repo, collection, rkey)
req, err := http.NewRequestWithContext(ctx, "GET", url, nil) if err != nil { return nil, err }
dpopProof, err := c.createDPoPProof("GET", url) if err != nil { return nil, err }
req.Header.Set("Authorization", "DPoP "+c.AccessToken) req.Header.Set("DPoP", dpopProof)
resp, err := http.DefaultClient.Do(req) if err != nil { return nil, err } defer resp.Body.Close()
if resp.StatusCode >= 400 { bodyBytes, _ := io.ReadAll(resp.Body) return nil, fmt.Errorf("XRPC error %d: %s", resp.StatusCode, string(bodyBytes)) }
var output GetRecordOutput if err := json.NewDecoder(resp.Body).Decode(&output); err != nil { return nil, err }
return &output, nil}
type ListRecordsRecord struct { URI string `json:"uri"` CID string `json:"cid"` Value json.RawMessage `json:"value"`}
type ListRecordsOutput struct { Cursor string `json:"cursor"` Records []ListRecordsRecord `json:"records"`}
func (c *Client) ListRecords(ctx context.Context, repo, collection string, limit int) (*ListRecordsOutput, error) { url := fmt.Sprintf("%s/xrpc/com.atproto.repo.listRecords?repo=%s&collection=%s&limit=%d", c.PDS, repo, collection, limit)
req, err := http.NewRequestWithContext(ctx, "GET", url, nil) if err != nil { return nil, err }
dpopProof, err := c.createDPoPProof("GET", url) if err != nil { return nil, err }
req.Header.Set("Authorization", "DPoP "+c.AccessToken) req.Header.Set("DPoP", dpopProof)
resp, err := http.DefaultClient.Do(req) if err != nil { return nil, err } defer resp.Body.Close()
if resp.StatusCode >= 400 { bodyBytes, _ := io.ReadAll(resp.Body) return nil, fmt.Errorf("XRPC error %d: %s", resp.StatusCode, string(bodyBytes)) }
var output ListRecordsOutput if err := json.NewDecoder(resp.Body).Decode(&output); err != nil { return nil, err }
return &output, nil}
type ResolveHandleOutput struct { Did string `json:"did"`}
func (c *Client) ResolveHandle(ctx context.Context, handle string) (string, error) { url := fmt.Sprintf("%s/xrpc/com.atproto.identity.resolveHandle?handle=%s", c.PDS, handle)
req, err := http.NewRequestWithContext(ctx, "GET", url, nil) if err != nil { return "", err }
resp, err := http.DefaultClient.Do(req) if err != nil { return "", err } defer resp.Body.Close()
if resp.StatusCode >= 400 { return "", fmt.Errorf("XRPC error %d", resp.StatusCode) }
var output ResolveHandleOutput if err := json.NewDecoder(resp.Body).Decode(&output); err != nil { return "", err }
return output.Did, nil}
type UploadBlobOutput struct { Blob BlobRef `json:"blob"`}
func (c *Client) UploadBlob(ctx context.Context, data []byte, contentType string) (*BlobRef, error) { url := fmt.Sprintf("%s/xrpc/com.atproto.repo.uploadBlob", c.PDS)
maxRetries := 2 for i := 0; i < maxRetries; i++ { req, err := http.NewRequestWithContext(ctx, "POST", url, bytes.NewReader(data)) if err != nil { return nil, err }
req.Header.Set("Content-Type", contentType)
dpopProof, err := c.createDPoPProof("POST", url) if err != nil { return nil, fmt.Errorf("failed to create DPoP proof: %w", err) }
req.Header.Set("Authorization", "DPoP "+c.AccessToken) req.Header.Set("DPoP", dpopProof)
resp, err := http.DefaultClient.Do(req) if err != nil { return nil, err } defer resp.Body.Close()
if nonce := resp.Header.Get("DPoP-Nonce"); nonce != "" { c.DPoPNonce = nonce }
if resp.StatusCode < 400 { var output UploadBlobOutput if err := json.NewDecoder(resp.Body).Decode(&output); err != nil { return nil, err } return &output.Blob, nil }
bodyBytes, _ := io.ReadAll(resp.Body) if resp.StatusCode == 401 && (bytes.Contains(bodyBytes, []byte("use_dpop_nonce")) || bytes.Contains(bodyBytes, []byte("UseDpopNonce"))) { continue }
return nil, fmt.Errorf("XRPC error %d: %s", resp.StatusCode, string(bodyBytes)) }
return nil, fmt.Errorf("upload blob failed after retries")}