diff --git a/events/dbpersist/dbpersist.go b/events/dbpersist/dbpersist.go index 3b24bd5b..a34e787a 100644 --- a/events/dbpersist/dbpersist.go +++ b/events/dbpersist/dbpersist.go @@ -171,6 +171,8 @@ func (p *DbPersistence) flushBatchLocked(ctx context.Context) error { switch { case e.RepoCommit != nil: e.RepoCommit.Seq = int64(item.Seq) + case e.RepoSync != nil: + e.RepoSync.Seq = int64(item.Seq) case e.RepoHandle != nil: e.RepoHandle.Seq = int64(item.Seq) case e.RepoIdentity != nil: @@ -218,6 +220,11 @@ func (p *DbPersistence) Persist(ctx context.Context, e *events.XRPCStreamEvent) if err != nil { return err } + case e.RepoSync != nil: + rer, err = p.RecordFromRepoSync(ctx, e.RepoSync) + if err != nil { + return err + } case e.RepoHandle != nil: rer, err = p.RecordFromHandleChange(ctx, e.RepoHandle) if err != nil { @@ -371,6 +378,28 @@ func (p *DbPersistence) RecordFromRepoCommit(ctx context.Context, evt *comatprot return &rer, nil } +func (p *DbPersistence) RecordFromRepoSync(ctx context.Context, evt *comatproto.SyncSubscribeRepos_Sync) (*RepoEventRecord, error) { + + uid, err := p.uidForDid(ctx, evt.Did) + if err != nil { + return nil, err + } + + t, err := time.Parse(util.ISO8601, evt.Time) + if err != nil { + return nil, err + } + + rer := RepoEventRecord{ + Repo: uid, + Type: "repo_sync", + Time: t, + Rev: evt.Rev, + } + + return &rer, nil +} + func (p *DbPersistence) Playback(ctx context.Context, since int64, cb func(*events.XRPCStreamEvent) error) error { pageSize := 1000 @@ -449,6 +478,8 @@ func (p *DbPersistence) hydrateBatch(ctx context.Context, batch []*RepoEventReco switch { case record.Commit != nil: streamEvent, err = p.hydrateCommit(ctx, record) + case record.Type == "repo_sync": + streamEvent, err = p.hydrateSyncEvent(ctx, record) case record.NewHandle != nil: streamEvent, err = p.hydrateHandleChange(ctx, record) case record.Type == "repo_identity": @@ -639,6 +670,29 @@ func (p *DbPersistence) hydrateCommit(ctx context.Context, rer *RepoEventRecord) return &events.XRPCStreamEvent{RepoCommit: out}, nil } +func (p *DbPersistence) hydrateSyncEvent(ctx context.Context, rer *RepoEventRecord) (*events.XRPCStreamEvent, error) { + + did, err := p.didForUid(ctx, rer.Repo) + if err != nil { + return nil, err + } + + evt := &comatproto.SyncSubscribeRepos_Sync{ + Seq: int64(rer.Seq), + Did: did, + Time: rer.Time.Format(util.ISO8601), + Rev: rer.Rev, + } + + cs, err := p.readCarSlice(ctx, rer) + if err != nil { + return nil, fmt.Errorf("read car slice: %w", err) + } + evt.Blocks = cs + + return &events.XRPCStreamEvent{RepoSync: evt}, nil +} + func (p *DbPersistence) readCarSlice(ctx context.Context, rer *RepoEventRecord) ([]byte, error) { buf := new(bytes.Buffer) diff --git a/events/diskpersist/diskpersist.go b/events/diskpersist/diskpersist.go index 069185fd..403035da 100644 --- a/events/diskpersist/diskpersist.go +++ b/events/diskpersist/diskpersist.go @@ -282,6 +282,7 @@ const ( evtKindTombstone = 3 evtKindIdentity = 4 evtKindAccount = 5 + evtKindSync = 6 ) var emptyHeader = make([]byte, headerSize) @@ -455,6 +456,8 @@ func (dp *DiskPersistence) doPersist(ctx context.Context, j persistJob) error { switch { case e.RepoCommit != nil: e.RepoCommit.Seq = seq + case e.RepoSync != nil: + e.RepoSync.Seq = seq case e.RepoHandle != nil: e.RepoHandle.Seq = seq case e.RepoIdentity != nil: @@ -509,6 +512,12 @@ func (dp *DiskPersistence) Persist(ctx context.Context, e *events.XRPCStreamEven if err := e.RepoCommit.MarshalCBOR(cw); err != nil { return fmt.Errorf("failed to marshal: %w", err) } + case e.RepoSync != nil: + evtKind = evtKindSync + did = e.RepoSync.Did + if err := e.RepoSync.MarshalCBOR(cw); err != nil { + return fmt.Errorf("failed to marshal: %w", err) + } case e.RepoHandle != nil: evtKind = evtKindHandle did = e.RepoHandle.Did @@ -745,6 +754,15 @@ func (dp *DiskPersistence) readEventsFrom(ctx context.Context, since int64, fn s if err := cb(&events.XRPCStreamEvent{RepoCommit: &evt}); err != nil { return nil, err } + case evtKindSync: + var evt atproto.SyncSubscribeRepos_Sync + if err := evt.UnmarshalCBOR(io.LimitReader(bufr, h.Len64())); err != nil { + return nil, err + } + evt.Seq = h.Seq + if err := cb(&events.XRPCStreamEvent{RepoSync: &evt}); err != nil { + return nil, err + } case evtKindHandle: var evt atproto.SyncSubscribeRepos_Handle if err := evt.UnmarshalCBOR(io.LimitReader(bufr, h.Len64())); err != nil { diff --git a/events/yolopersist/yolopersist.go b/events/yolopersist/yolopersist.go index d6b71363..1fbd461d 100644 --- a/events/yolopersist/yolopersist.go +++ b/events/yolopersist/yolopersist.go @@ -28,6 +28,8 @@ func (yp *YoloPersister) Persist(ctx context.Context, e *events.XRPCStreamEvent) switch { case e.RepoCommit != nil: e.RepoCommit.Seq = yp.seq + case e.RepoSync != nil: + e.RepoSync.Seq = yp.seq case e.RepoHandle != nil: e.RepoHandle.Seq = yp.seq case e.RepoIdentity != nil: -- 2.51.2 From f16258882a1c6a82c25465547ae824afbeec9486 Mon Sep 17 00:00:00 2001 From: jcalabro Date: Mon, 14 Apr 2025 10:02:52 -0400 Subject: [PATCH 02/17] install /debug/pprof/ handlers on the metrics server --- cmd/collectiondir/serve.go | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/cmd/collectiondir/serve.go b/cmd/collectiondir/serve.go index 1bb89f9b..f3f93396 100644 --- a/cmd/collectiondir/serve.go +++ b/cmd/collectiondir/serve.go @@ -9,6 +9,7 @@ import ( "log/slog" "net" "net/http" + _ "net/http/pprof" "net/url" "os" "os/signal" @@ -388,9 +389,14 @@ func (cs *collectionServer) handleCommit(commit *comatproto.SyncSubscribeRepos_C func (cs *collectionServer) StartMetricsServer(ctx context.Context, addr string) error { defer cs.wg.Done() defer cs.log.Info("metrics server exit") + + e := echo.New() + e.GET("/metrics", echo.WrapHandler(promhttp.Handler())) + e.Any("/debug/pprof/*", echo.WrapHandler(http.DefaultServeMux)) + cs.metricsServer = &http.Server{ Addr: addr, - Handler: promhttp.Handler(), + Handler: e, } return cs.metricsServer.ListenAndServe() } -- 2.51.2 From 82656e88f5ed5ff706c00cfabb98701c3e8e6d0a Mon Sep 17 00:00:00 2001 From: jcalabro Date: Mon, 14 Apr 2025 10:05:47 -0400 Subject: [PATCH 03/17] fix startup race condition --- cmd/collectiondir/serve.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/cmd/collectiondir/serve.go b/cmd/collectiondir/serve.go index f3f93396..00e1e249 100644 --- a/cmd/collectiondir/serve.go +++ b/cmd/collectiondir/serve.go @@ -194,12 +194,12 @@ func (cs *collectionServer) run(cctx *cli.Context) error { if err != nil { return fmt.Errorf("lru init, %w", err) } - cs.wg.Add(1) - go cs.ingestReceiver() cs.log = log cs.ctx = cctx.Context cs.AdminToken = cctx.String("admin-token") cs.ExepctedAuthHeader = "Bearer " + cs.AdminToken + cs.wg.Add(1) + go cs.ingestReceiver() pebblePath := cctx.String("pebble") cs.pcd = &PebbleCollectionDirectory{ log: cs.log, -- 2.51.2 From 8c00259a165833ad0e8e8870691b5f55e3725ea8 Mon Sep 17 00:00:00 2001 From: jcalabro Date: Mon, 14 Apr 2025 10:17:14 -0400 Subject: [PATCH 04/17] fix infinite hang on service shutdown --- cmd/collectiondir/serve.go | 62 ++++++++++++++++---------------------- 1 file changed, 26 insertions(+), 36 deletions(-) diff --git a/cmd/collectiondir/serve.go b/cmd/collectiondir/serve.go index 00e1e249..03733ca2 100644 --- a/cmd/collectiondir/serve.go +++ b/cmd/collectiondir/serve.go @@ -216,17 +216,14 @@ func (cs *collectionServer) run(cctx *cli.Context) error { } } cs.statsCacheFresh.L = &cs.statsCacheLock - errchan := make(chan error, 3) + apiAddr := cctx.String("api-listen") cs.wg.Add(1) - go func() { - errchan <- cs.StartApiServer(cctx.Context, apiAddr) - }() + go func() { cs.StartApiServer(cctx.Context, apiAddr) }() + metricsAddr := cctx.String("metrics-listen") cs.wg.Add(1) - go func() { - errchan <- cs.StartMetricsServer(cctx.Context, metricsAddr) - }() + go func() { cs.StartMetricsServer(cctx.Context, metricsAddr) }() upstream := cctx.String("upstream") if upstream != "" { @@ -248,19 +245,9 @@ func (cs *collectionServer) run(cctx *cli.Context) error { go cs.handleFirehose(fhevents) } - select { - case <-signals: - log.Info("received shutdown signal") - go errchanlog(cs.log, "server error", errchan) - return cs.Shutdown() - case err := <-errchan: - if err != nil { - log.Error("server error", "err", err) - go errchanlog(cs.log, "server error", errchan) - return cs.Shutdown() - } - } - return nil + <-signals + log.Info("received shutdown signal") + return cs.Shutdown() } func (cs *collectionServer) openDau() error { @@ -283,28 +270,31 @@ func (cs *collectionServer) openDau() error { return nil } -func errchanlog(log *slog.Logger, msg string, errchan <-chan error) { - for err := range errchan { - log.Error(msg, "err", err) - } -} - func (cs *collectionServer) Shutdown() error { close(cs.shutdown) - go func() { + + func() { + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + cs.log.Info("metrics shutdown start") - sherr := cs.metricsServer.Shutdown(context.Background()) + sherr := cs.metricsServer.Shutdown(ctx) cs.log.Info("metrics shutdown", "err", sherr) }() - cs.log.Info("api shutdown start...") - err := cs.apiServer.Shutdown(context.Background()) - //err := cs.esrv.Shutdown(context.Background()) - cs.log.Info("api shutdown, thread wait...", "err", err) - cs.wg.Wait() + + func() { + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + + cs.log.Info("api shutdown start...") + err := cs.apiServer.Shutdown(ctx) + cs.log.Info("api shutdown, thread wait...", "err", err) + }() + cs.log.Info("threads done, db close...") - ee := cs.pcd.Close() - if ee != nil { - cs.log.Error("failed to shutdown pebble", "err", ee) + err := cs.pcd.Close() + if err != nil { + cs.log.Error("failed to shutdown pebble", "err", err) } cs.log.Info("db done. done.") return err -- 2.51.2 From 5191ddd755928d5763cee92d7510c5999583d006 Mon Sep 17 00:00:00 2001 From: jcalabro Date: Mon, 14 Apr 2025 10:23:48 -0400 Subject: [PATCH 05/17] fix server create/shutdown data race --- cmd/collectiondir/serve.go | 21 ++++++++++++++------- 1 file changed, 14 insertions(+), 7 deletions(-) diff --git a/cmd/collectiondir/serve.go b/cmd/collectiondir/serve.go index 03733ca2..3b34ec35 100644 --- a/cmd/collectiondir/serve.go +++ b/cmd/collectiondir/serve.go @@ -217,9 +217,12 @@ func (cs *collectionServer) run(cctx *cli.Context) error { } cs.statsCacheFresh.L = &cs.statsCacheLock - apiAddr := cctx.String("api-listen") + apiServerEcho, err := cs.createApiServer(cctx.Context, cctx.String("api-listen")) + if err != nil { + return err + } cs.wg.Add(1) - go func() { cs.StartApiServer(cctx.Context, apiAddr) }() + go func() { cs.StartApiServer(cctx.Context, apiServerEcho) }() metricsAddr := cctx.String("metrics-listen") cs.wg.Add(1) @@ -391,13 +394,11 @@ func (cs *collectionServer) StartMetricsServer(ctx context.Context, addr string) return cs.metricsServer.ListenAndServe() } -func (cs *collectionServer) StartApiServer(ctx context.Context, addr string) error { - defer cs.wg.Done() - defer cs.log.Info("api server exit") +func (cs *collectionServer) createApiServer(ctx context.Context, addr string) (*echo.Echo, error) { var lc net.ListenConfig li, err := lc.Listen(ctx, "tcp", addr) if err != nil { - return err + return nil, err } e := echo.New() e.HideBanner = true @@ -427,7 +428,13 @@ func (cs *collectionServer) StartApiServer(ctx context.Context, addr string) err Handler: e, } cs.apiServer = srv - return srv.Serve(li) + return e, nil +} + +func (cs *collectionServer) StartApiServer(ctx context.Context, e *echo.Echo) error { + defer cs.wg.Done() + defer cs.log.Info("api server exit") + return cs.apiServer.Serve(e.Listener) } const statsCacheDuration = time.Second * 300 -- 2.51.2 From c90e279b78482bea54451cbf6d0520faa432fa8d Mon Sep 17 00:00:00 2001 From: jcalabro Date: Mon, 14 Apr 2025 10:26:57 -0400 Subject: [PATCH 06/17] fix final race condition --- cmd/collectiondir/serve.go | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/cmd/collectiondir/serve.go b/cmd/collectiondir/serve.go index 3b34ec35..ebcdc99c 100644 --- a/cmd/collectiondir/serve.go +++ b/cmd/collectiondir/serve.go @@ -224,9 +224,9 @@ func (cs *collectionServer) run(cctx *cli.Context) error { cs.wg.Add(1) go func() { cs.StartApiServer(cctx.Context, apiServerEcho) }() - metricsAddr := cctx.String("metrics-listen") + cs.createMetricsServer(cctx.String("metrics-listen")) cs.wg.Add(1) - go func() { cs.StartMetricsServer(cctx.Context, metricsAddr) }() + go func() { cs.StartMetricsServer(cctx.Context) }() upstream := cctx.String("upstream") if upstream != "" { @@ -379,10 +379,7 @@ func (cs *collectionServer) handleCommit(commit *comatproto.SyncSubscribeRepos_C } } -func (cs *collectionServer) StartMetricsServer(ctx context.Context, addr string) error { - defer cs.wg.Done() - defer cs.log.Info("metrics server exit") - +func (cs *collectionServer) createMetricsServer(addr string) { e := echo.New() e.GET("/metrics", echo.WrapHandler(promhttp.Handler())) e.Any("/debug/pprof/*", echo.WrapHandler(http.DefaultServeMux)) @@ -391,6 +388,12 @@ func (cs *collectionServer) StartMetricsServer(ctx context.Context, addr string) Addr: addr, Handler: e, } +} + +func (cs *collectionServer) StartMetricsServer(ctx context.Context) error { + defer cs.wg.Done() + defer cs.log.Info("metrics server exit") + return cs.metricsServer.ListenAndServe() } -- 2.51.2 From 444a6aa911744316ac491b404a3a793311bfddfe Mon Sep 17 00:00:00 2001 From: jcalabro Date: Tue, 15 Apr 2025 09:09:15 -0400 Subject: [PATCH 07/17] handle server errors as fatal --- cmd/collectiondir/serve.go | 17 +++++++++++++---- 1 file changed, 13 insertions(+), 4 deletions(-) diff --git a/cmd/collectiondir/serve.go b/cmd/collectiondir/serve.go index ebcdc99c..965d10e0 100644 --- a/cmd/collectiondir/serve.go +++ b/cmd/collectiondir/serve.go @@ -5,6 +5,7 @@ import ( "context" "encoding/csv" "encoding/json" + "errors" "fmt" "log/slog" "net" @@ -390,11 +391,15 @@ func (cs *collectionServer) createMetricsServer(addr string) { } } -func (cs *collectionServer) StartMetricsServer(ctx context.Context) error { +func (cs *collectionServer) StartMetricsServer(ctx context.Context) { defer cs.wg.Done() defer cs.log.Info("metrics server exit") - return cs.metricsServer.ListenAndServe() + err := cs.metricsServer.ListenAndServe() + if err != nil && !errors.Is(err, http.ErrServerClosed) { + slog.Error("error in metrics server", "err", err) + os.Exit(1) + } } func (cs *collectionServer) createApiServer(ctx context.Context, addr string) (*echo.Echo, error) { @@ -434,10 +439,14 @@ func (cs *collectionServer) createApiServer(ctx context.Context, addr string) (* return e, nil } -func (cs *collectionServer) StartApiServer(ctx context.Context, e *echo.Echo) error { +func (cs *collectionServer) StartApiServer(ctx context.Context, e *echo.Echo) { defer cs.wg.Done() defer cs.log.Info("api server exit") - return cs.apiServer.Serve(e.Listener) + err := cs.apiServer.Serve(e.Listener) + if err != nil && !errors.Is(err, http.ErrServerClosed) { + slog.Error("error in api server", "err", err) + os.Exit(1) + } } const statsCacheDuration = time.Second * 300 -- 2.51.2 From 5fd978367e271c14f315e2df95fde7adc76e11bb Mon Sep 17 00:00:00 2001 From: jcalabro Date: Thu, 17 Apr 2025 16:18:40 -0400 Subject: [PATCH 08/17] still wait for server shutdown --- cmd/collectiondir/serve.go | 1 + 1 file changed, 1 insertion(+) diff --git a/cmd/collectiondir/serve.go b/cmd/collectiondir/serve.go index 91819aa0..0abf074f 100644 --- a/cmd/collectiondir/serve.go +++ b/cmd/collectiondir/serve.go @@ -301,6 +301,7 @@ func (cs *collectionServer) Shutdown() error { cs.log.Error("failed to shutdown pebble", "err", err) } cs.log.Info("db done. done.") + cs.wg.Wait() return err } -- 2.51.2 From a55e998453f0829e5783c92299eaca740cc11b47 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Mon, 28 Apr 2025 20:39:00 -0700 Subject: [PATCH 09/17] MVP JWK support for atproto/crypto package --- atproto/crypto/jwk.go | 124 +++++++++++++++++++++++++++++++++++++ atproto/crypto/jwk_test.go | 86 +++++++++++++++++++++++++ atproto/crypto/keys.go | 3 + 3 files changed, 213 insertions(+) create mode 100644 atproto/crypto/jwk.go create mode 100644 atproto/crypto/jwk_test.go diff --git a/atproto/crypto/jwk.go b/atproto/crypto/jwk.go new file mode 100644 index 00000000..24995a2f --- /dev/null +++ b/atproto/crypto/jwk.go @@ -0,0 +1,124 @@ +package crypto + +import ( + "crypto/ecdsa" + "crypto/elliptic" + "encoding/base64" + "encoding/json" + "fmt" + "math/big" + + secp256k1 "gitlab.com/yawning/secp256k1-voi" + secp256k1secec "gitlab.com/yawning/secp256k1-voi/secec" +) + +// Representation of a JSON Web Key (JWK), as relevant to the keys supported by this package. +// +// Expected to be marshalled/unmarshalled as JSON. +type JWK struct { + Algorithm string `json:"alg"` + Curve string `json:"crv"` + X string `json:"x"` // base64url, no padding + Y string `json:"y"` // base64url, no padding + Use string `json:"use,omitempty"` + KeyID *string `json:"kid,omitempty"` +} + +// Loads a [PublicKey] from JWK (serialized as JSON bytes) +func ParsePublicJWKBytes(jwkBytes []byte) (PublicKey, error) { + var jwk JWK + if err := json.Unmarshal(jwkBytes, &jwk); err != nil { + return nil, fmt.Errorf("parsing JWK JSON: %w", err) + } + return ParsePublicJWK(jwk) +} + +// Loads a [PublicKey] from JWK struct. +func ParsePublicJWK(jwk JWK) (PublicKey, error) { + + if jwk.Algorithm != "EC" { + return nil, fmt.Errorf("unsupported JWK cryptography: %s", jwk.Algorithm) + } + + // base64url with no encoding + xbuf, err := base64.RawURLEncoding.DecodeString(jwk.X) + if err != nil { + return nil, fmt.Errorf("invalid JWK base64 encoding: %w", err) + } + ybuf, err := base64.RawURLEncoding.DecodeString(jwk.Y) + if err != nil { + return nil, fmt.Errorf("invalid JWK base64 encoding: %w", err) + } + + switch jwk.Curve { + case "P-256": + curve := elliptic.P256() + + var x, y big.Int + x.SetBytes(xbuf) + y.SetBytes(ybuf) + + if !curve.Params().IsOnCurve(&x, &y) { + return nil, fmt.Errorf("invalid P-256 public key (not on curve)") + } + pubECDSA := &ecdsa.PublicKey{ + Curve: curve, + X: &x, + Y: &y, + } + pub := PublicKeyP256{pubP256: *pubECDSA} + err := pub.checkCurve() + if err != nil { + return nil, err + } + return &pub, nil + case "secp256k1": // K-256 + if len(xbuf) != 32 || len(ybuf) != 32 { + return nil, fmt.Errorf("invalid K-256 coordinates") + } + xarr := ([32]byte)(xbuf[:32]) + yarr := ([32]byte)(ybuf[:32]) + p, err := secp256k1.NewPointFromCoords(&xarr, &yarr) + if err != nil { + return nil, fmt.Errorf("invalid K-256 coordinates: %w", err) + } + pubK, err := secp256k1secec.NewPublicKeyFromPoint(p) + if err != nil { + return nil, fmt.Errorf("invalid K-256/secp256k1 public key: %w", err) + } + pub := PublicKeyK256{pubK256: pubK} + err = pub.ensureBytes() + if err != nil { + return nil, err + } + return &pub, nil + default: + return nil, fmt.Errorf("unsupported JWK cryptography: %s", jwk.Curve) + } +} + +func (k *PublicKeyP256) JWK() (*JWK, error) { + jwk := JWK{ + Algorithm: "EC", + Curve: "P-256", + X: base64.RawURLEncoding.EncodeToString(k.pubP256.X.Bytes()), + Y: base64.RawURLEncoding.EncodeToString(k.pubP256.Y.Bytes()), + } + return &jwk, nil +} + +func (k *PublicKeyK256) JWK() (*JWK, error) { + raw := k.UncompressedBytes() + if len(raw) != 65 { + return nil, fmt.Errorf("unexpected K-256 bytes size") + } + xbytes := raw[1:33] + ybytes := raw[33:65] + jwk := JWK{ + Algorithm: "EC", + Curve: "secp256k1", + X: base64.RawURLEncoding.EncodeToString(xbytes), + Y: base64.RawURLEncoding.EncodeToString(ybytes), + } + return &jwk, nil +} diff --git a/atproto/crypto/jwk_test.go b/atproto/crypto/jwk_test.go new file mode 100644 index 00000000..444b7670 --- /dev/null +++ b/atproto/crypto/jwk_test.go @@ -0,0 +1,86 @@ +package crypto + +import ( + "testing" + + "github.com/stretchr/testify/assert" +) + +func TestParseJWK(t *testing.T) { + assert := assert.New(t) + + jwkTestFixtures := []string{ + // https://openid.net/specs/draft-jones-json-web-key-03.html + `{ + "alg":"EC", + "crv":"P-256", + "x":"MKBCTNIcKUSDii11ySs3526iDZ8AiTo7Tu6KPAqv7D4", + "y":"4Etl6SRW2YiLUrN5vfvVHuhp7x8PxltmWWlbbM4IFyM", + "use":"enc", + "kid":"1" + }`, + // https://w3c-ccg.github.io/lds-ecdsa-secp256k1-2019/ + `{ + "alg": "EC", + "crv": "secp256k1", + "kid": "JUvpllMEYUZ2joO59UNui_XYDqxVqiFLLAJ8klWuPBw", + "x": "dWCvM4fTdeM0KmloF57zxtBPXTOythHPMm1HCLrdd3A", + "y": "36uMVGM7hnw-N6GnjFcihWE3SkrhMLzzLCdPMXPEXlA" + }`, + } + + for _, jwkBytes := range jwkTestFixtures { + _, err := ParsePublicJWKBytes([]byte(jwkBytes)) + assert.NoError(err) + } +} + +func TestP256GenJWK(t *testing.T) { + assert := assert.New(t) + + priv, err := GeneratePrivateKeyP256() + if err != nil { + t.Fatal(err) + } + pub, err := priv.PublicKey() + if err != nil { + t.Fatal(err) + } + + pk, ok := pub.(*PublicKeyP256) + if !ok { + t.Fatal() + } + jwk, err := pk.JWK() + if err != nil { + t.Fatal(err) + } + + _, err = ParsePublicJWK(*jwk) + assert.NoError(err) +} + +func TestK256GenJWK(t *testing.T) { + assert := assert.New(t) + + priv, err := GeneratePrivateKeyK256() + if err != nil { + t.Fatal(err) + } + pub, err := priv.PublicKey() + if err != nil { + t.Fatal(err) + } + + pk, ok := pub.(*PublicKeyK256) + if !ok { + t.Fatal() + } + jwk, err := pk.JWK() + if err != nil { + t.Fatal(err) + } + + _, err = ParsePublicJWK(*jwk) + assert.NoError(err) +} diff --git a/atproto/crypto/keys.go b/atproto/crypto/keys.go index c540a1b9..d8192fe6 100644 --- a/atproto/crypto/keys.go +++ b/atproto/crypto/keys.go @@ -63,6 +63,9 @@ type PublicKey interface { // For systems with no compressed/uncompressed distinction, returns the same // value as Bytes(). UncompressedBytes() []byte + + // Serialization as JWK struct (which can be marshalled to JSON) + JWK() (*JWK, error) } var ErrInvalidSignature = errors.New("crytographic signature invalid") -- 2.51.2 From a1b106cc1806633ffcfb3d6f3c8dc6d50a67b7a4 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Fri, 9 May 2025 01:55:47 -0700 Subject: [PATCH 10/17] crypto: fix JWK serialization (kty) --- atproto/crypto/jwk.go | 32 ++++++++++++++++---------------- 1 file changed, 16 insertions(+), 16 deletions(-) diff --git a/atproto/crypto/jwk.go b/atproto/crypto/jwk.go index 24995a2f..44956610 100644 --- a/atproto/crypto/jwk.go +++ b/atproto/crypto/jwk.go @@ -16,12 +16,12 @@ import ( // // Expected to be marshalled/unmarshalled as JSON. type JWK struct { - Algorithm string `json:"alg"` - Curve string `json:"crv"` - X string `json:"x"` // base64url, no padding - Y string `json:"y"` // base64url, no padding - Use string `json:"use,omitempty"` - KeyID *string `json:"kid,omitempty"` + KeyType string `json:"kty"` + Curve string `json:"crv"` + X string `json:"x"` // base64url, no padding + Y string `json:"y"` // base64url, no padding + Use string `json:"use,omitempty"` + KeyID *string `json:"kid,omitempty"` } // Loads a [PublicKey] from JWK (serialized as JSON bytes) @@ -36,8 +36,8 @@ func ParsePublicJWKBytes(jwkBytes []byte) (PublicKey, error) { // Loads a [PublicKey] from JWK struct. func ParsePublicJWK(jwk JWK) (PublicKey, error) { - if jwk.Algorithm != "EC" { - return nil, fmt.Errorf("unsupported JWK cryptography: %s", jwk.Algorithm) + if jwk.KeyType != "EC" { + return nil, fmt.Errorf("unsupported JWK key type: %s", jwk.KeyType) } // base64url with no encoding @@ -99,10 +99,10 @@ func ParsePublicJWK(jwk JWK) (PublicKey, error) { func (k *PublicKeyP256) JWK() (*JWK, error) { jwk := JWK{ - Algorithm: "EC", - Curve: "P-256", - X: base64.RawURLEncoding.EncodeToString(k.pubP256.X.Bytes()), - Y: base64.RawURLEncoding.EncodeToString(k.pubP256.Y.Bytes()), + KeyType: "EC", + Curve: "P-256", + X: base64.RawURLEncoding.EncodeToString(k.pubP256.X.Bytes()), + Y: base64.RawURLEncoding.EncodeToString(k.pubP256.Y.Bytes()), } return &jwk, nil } @@ -115,10 +115,10 @@ func (k *PublicKeyK256) JWK() (*JWK, error) { xbytes := raw[1:33] ybytes := raw[33:65] jwk := JWK{ - Algorithm: "EC", - Curve: "secp256k1", - X: base64.RawURLEncoding.EncodeToString(xbytes), - Y: base64.RawURLEncoding.EncodeToString(ybytes), + KeyType: "EC", + Curve: "secp256k1", + X: base64.RawURLEncoding.EncodeToString(xbytes), + Y: base64.RawURLEncoding.EncodeToString(ybytes), } return &jwk, nil } -- 2.51.2 From 920f5025b95ebc9df521cea053d27bbb0fe373f8 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Fri, 9 May 2025 21:28:14 -0700 Subject: [PATCH 11/17] fix crypto JWK tests --- atproto/crypto/jwk_test.go | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) diff --git a/atproto/crypto/jwk_test.go b/atproto/crypto/jwk_test.go index 444b7670..23a219de 100644 --- a/atproto/crypto/jwk_test.go +++ b/atproto/crypto/jwk_test.go @@ -10,18 +10,30 @@ func TestParseJWK(t *testing.T) { assert := assert.New(t) jwkTestFixtures := []string{ - // https://openid.net/specs/draft-jones-json-web-key-03.html + // https://datatracker.ietf.org/doc/html/rfc7517#appendix-A.1 + `{ + "kty":"EC", + "crv":"P-256", + "x":"MKBCTNIcKUSDii11ySs3526iDZ8AiTo7Tu6KPAqv7D4", + "y":"4Etl6SRW2YiLUrN5vfvVHuhp7x8PxltmWWlbbM4IFyM", + "d":"870MB6gfuTJ4HtUnUvYMyJpr5eUZNP4Bk43bVdj3eAE", + "use":"enc", + "kid":"1" + }`, + // https://openid.net/specs/draft-jones-json-web-key-03.html; with kty in addition to alg `{ "alg":"EC", + "kty":"EC", "crv":"P-256", "x":"MKBCTNIcKUSDii11ySs3526iDZ8AiTo7Tu6KPAqv7D4", "y":"4Etl6SRW2YiLUrN5vfvVHuhp7x8PxltmWWlbbM4IFyM", "use":"enc", "kid":"1" }`, - // https://w3c-ccg.github.io/lds-ecdsa-secp256k1-2019/ + // https://w3c-ccg.github.io/lds-ecdsa-secp256k1-2019/; with kty in addition to alg `{ "alg": "EC", + "kty": "EC", "crv": "secp256k1", "kid": "JUvpllMEYUZ2joO59UNui_XYDqxVqiFLLAJ8klWuPBw", "x": "dWCvM4fTdeM0KmloF57zxtBPXTOythHPMm1HCLrdd3A", -- 2.51.2 From a9ac41723e3fef22895d1a16bf6d76878a1612cc Mon Sep 17 00:00:00 2001 From: Larry Person Date: Mon, 26 May 2025 10:22:34 -0400 Subject: [PATCH 12/17] Remove extraneous word from error message --- atproto/data/parse.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/atproto/data/parse.go b/atproto/data/parse.go index 2f1f4305..e8e7412d 100644 --- a/atproto/data/parse.go +++ b/atproto/data/parse.go @@ -11,7 +11,7 @@ import ( func parseFloat(f float64) (int64, error) { if f != float64(int64(f)) { - return 0, fmt.Errorf("number was is not a safe integer: %f", f) + return 0, fmt.Errorf("number is not a safe integer: %f", f) } return int64(f), nil } -- 2.51.2 From 4854778cb97a3376bbe70069ffcedf718a4b6178 Mon Sep 17 00:00:00 2001 From: jcalabro Date: Mon, 26 May 2025 15:54:22 -0400 Subject: [PATCH 13/17] Lexgen Latest --- api/bsky/feeddefs.go | 6 ++++++ api/bsky/feedgetFeedSkeleton.go | 2 ++ api/bsky/feedlike.go | 1 + api/bsky/feedrepost.go | 1 + api/bsky/notificationlistNotifications.go | 2 +- 5 files changed, 11 insertions(+), 1 deletion(-) diff --git a/api/bsky/feeddefs.go b/api/bsky/feeddefs.go index 5c7e64a4..14eedd7e 100644 --- a/api/bsky/feeddefs.go +++ b/api/bsky/feeddefs.go @@ -35,6 +35,8 @@ type FeedDefs_FeedViewPost struct { Post *FeedDefs_PostView `json:"post" cborgen:"post"` Reason *FeedDefs_FeedViewPost_Reason `json:"reason,omitempty" cborgen:"reason,omitempty"` Reply *FeedDefs_ReplyRef `json:"reply,omitempty" cborgen:"reply,omitempty"` + // reqId: Unique identifier per request that may be passed back alongside interactions. + ReqId *string `json:"reqId,omitempty" cborgen:"reqId,omitempty"` } type FeedDefs_FeedViewPost_Reason struct { @@ -104,6 +106,8 @@ type FeedDefs_Interaction struct { // feedContext: Context on a feed item that was originally supplied by the feed generator on getFeedSkeleton. FeedContext *string `json:"feedContext,omitempty" cborgen:"feedContext,omitempty"` Item *string `json:"item,omitempty" cborgen:"item,omitempty"` + // reqId: Unique identifier per request that may be passed back alongside interactions. + ReqId *string `json:"reqId,omitempty" cborgen:"reqId,omitempty"` } // FeedDefs_NotFoundPost is a "notFoundPost" in the app.bsky.feed.defs schema. @@ -207,7 +211,9 @@ type FeedDefs_ReasonPin struct { type FeedDefs_ReasonRepost struct { LexiconTypeID string `json:"$type,const=app.bsky.feed.defs#reasonRepost" cborgen:"$type,const=app.bsky.feed.defs#reasonRepost"` By *ActorDefs_ProfileViewBasic `json:"by" cborgen:"by"` + Cid *string `json:"cid,omitempty" cborgen:"cid,omitempty"` IndexedAt string `json:"indexedAt" cborgen:"indexedAt"` + Uri *string `json:"uri,omitempty" cborgen:"uri,omitempty"` } // FeedDefs_ReplyRef is a "replyRef" in the app.bsky.feed.defs schema. diff --git a/api/bsky/feedgetFeedSkeleton.go b/api/bsky/feedgetFeedSkeleton.go index f4f94b9a..7f2a5ba2 100644 --- a/api/bsky/feedgetFeedSkeleton.go +++ b/api/bsky/feedgetFeedSkeleton.go @@ -14,6 +14,8 @@ import ( type FeedGetFeedSkeleton_Output struct { Cursor *string `json:"cursor,omitempty" cborgen:"cursor,omitempty"` Feed []*FeedDefs_SkeletonFeedPost `json:"feed" cborgen:"feed"` + // reqId: Unique identifier per request that may be passed back alongside interactions. + ReqId *string `json:"reqId,omitempty" cborgen:"reqId,omitempty"` } // FeedGetFeedSkeleton calls the XRPC method "app.bsky.feed.getFeedSkeleton". diff --git a/api/bsky/feedlike.go b/api/bsky/feedlike.go index 08a3246b..a481bf9f 100644 --- a/api/bsky/feedlike.go +++ b/api/bsky/feedlike.go @@ -17,4 +17,5 @@ type FeedLike struct { LexiconTypeID string `json:"$type,const=app.bsky.feed.like" cborgen:"$type,const=app.bsky.feed.like"` CreatedAt string `json:"createdAt" cborgen:"createdAt"` Subject *comatprototypes.RepoStrongRef `json:"subject" cborgen:"subject"` + Via *comatprototypes.RepoStrongRef `json:"via,omitempty" cborgen:"via,omitempty"` } diff --git a/api/bsky/feedrepost.go b/api/bsky/feedrepost.go index 3216bd33..a107dada 100644 --- a/api/bsky/feedrepost.go +++ b/api/bsky/feedrepost.go @@ -17,4 +17,5 @@ type FeedRepost struct { LexiconTypeID string `json:"$type,const=app.bsky.feed.repost" cborgen:"$type,const=app.bsky.feed.repost"` CreatedAt string `json:"createdAt" cborgen:"createdAt"` Subject *comatprototypes.RepoStrongRef `json:"subject" cborgen:"subject"` + Via *comatprototypes.RepoStrongRef `json:"via,omitempty" cborgen:"via,omitempty"` } diff --git a/api/bsky/notificationlistNotifications.go b/api/bsky/notificationlistNotifications.go index d396b340..2ed53e1f 100644 --- a/api/bsky/notificationlistNotifications.go +++ b/api/bsky/notificationlistNotifications.go @@ -19,7 +19,7 @@ type NotificationListNotifications_Notification struct { IndexedAt string `json:"indexedAt" cborgen:"indexedAt"` IsRead bool `json:"isRead" cborgen:"isRead"` Labels []*comatprototypes.LabelDefs_Label `json:"labels,omitempty" cborgen:"labels,omitempty"` - // reason: Expected values are 'like', 'repost', 'follow', 'mention', 'reply', 'quote', 'starterpack-joined', 'verified', and 'unverified'. + // reason: The reason why this notification was delivered - e.g. your post was liked, or you received a new follower. Reason string `json:"reason" cborgen:"reason"` ReasonSubject *string `json:"reasonSubject,omitempty" cborgen:"reasonSubject,omitempty"` Record *util.LexiconTypeDecoder `json:"record" cborgen:"record"` -- 2.51.2 From dee24e74ff358b5f05b7c95d19ced190a25e1c33 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Sat, 3 May 2025 19:14:29 -0700 Subject: [PATCH 14/17] new lexutil.LexClient interface which xrpc.Client implements --- lex/gen.go | 1 - lex/type_schema.go | 8 ++++---- lex/util/client.go | 17 +++++++++++++++++ xrpc/xrpc.go | 23 ++++++++++++++--------- 4 files changed, 35 insertions(+), 14 deletions(-) create mode 100644 lex/util/client.go diff --git a/lex/gen.go b/lex/gen.go index 63ed56e8..7d6d3adc 100644 --- a/lex/gen.go +++ b/lex/gen.go @@ -135,7 +135,6 @@ func GenCodeForSchema(pkg Package, reqcode bool, s *Schema, packages []Package, pf("\t\"fmt\"\n") pf("\t\"encoding/json\"\n") pf("\tcbg \"github.com/whyrusleeping/cbor-gen\"\n") - pf("\t\"github.com/bluesky-social/indigo/xrpc\"\n") pf("\t\"github.com/bluesky-social/indigo/lex/util\"\n") for _, xpkg := range packages { if xpkg.Prefix != pkg.Prefix { diff --git a/lex/type_schema.go b/lex/type_schema.go index 62c5e670..78b8f769 100644 --- a/lex/type_schema.go +++ b/lex/type_schema.go @@ -54,7 +54,7 @@ func (s *TypeSchema) WriteRPC(w io.Writer, typename, inputname string) error { pf := printerf(w) fname := typename - params := "ctx context.Context, c *xrpc.Client" + params := "ctx context.Context, c util.LexClient" inpvar := "nil" inpenc := "" @@ -161,14 +161,14 @@ func (s *TypeSchema) WriteRPC(w io.Writer, typename, inputname string) error { var reqtype string switch s.Type { case "procedure": - reqtype = "xrpc.Procedure" + reqtype = "util.Procedure" case "query": - reqtype = "xrpc.Query" + reqtype = "util.Query" default: return fmt.Errorf("can only generate RPC for Query or Procedure (got %s)", s.Type) } - pf("\tif err := c.Do(ctx, %s, %q, \"%s\", %s, %s, %s); err != nil {\n", reqtype, inpenc, s.id, queryparams, inpvar, outvar) + pf("\tif err := c.LexDo(ctx, %s, %q, \"%s\", %s, %s, %s); err != nil {\n", reqtype, inpenc, s.id, queryparams, inpvar, outvar) pf("\t\treturn %s\n", errRet) pf("\t}\n\n") pf("\treturn %s\n", outRet) diff --git a/lex/util/client.go b/lex/util/client.go new file mode 100644 index 00000000..f40a78f7 --- /dev/null +++ b/lex/util/client.go @@ -0,0 +1,17 @@ +package util + +import ( + "context" +) + +type XRPCRequestType int + +const ( + Query = XRPCRequestType(iota) + Procedure +) + +// API client interface used in lexgen. +type LexClient interface { + LexDo(ctx context.Context, kind XRPCRequestType, inpenc string, method string, params map[string]any, bodyobj any, out any) error +} diff --git a/xrpc/xrpc.go b/xrpc/xrpc.go index 740835d2..2d939fe5 100644 --- a/xrpc/xrpc.go +++ b/xrpc/xrpc.go @@ -13,6 +13,7 @@ import ( "strings" "time" + lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/util" "github.com/carlmjohnson/versioninfo" ) @@ -34,7 +35,12 @@ func (c *Client) getClient() *http.Client { return c.Client } -type XRPCRequestType int +type XRPCRequestType = lexutil.XRPCRequestType + +var ( + Query = lexutil.Query + Procedure = lexutil.Procedure +) type AuthInfo struct { AccessJwt string `json:"accessJwt"` @@ -110,11 +116,6 @@ type RatelimitInfo struct { Reset time.Time } -const ( - Query = XRPCRequestType(iota) - Procedure -) - // makeParams converts a map of string keys and any values into a URL-encoded string. // If a value is a slice of strings, it will be joined with commas. // Generally the values will be strings, numbers, booleans, or slices of strings @@ -133,7 +134,7 @@ func makeParams(p map[string]any) string { return params.Encode() } -func (c *Client) Do(ctx context.Context, kind XRPCRequestType, inpenc string, method string, params map[string]interface{}, bodyobj interface{}, out interface{}) error { +func (c *Client) Do(ctx context.Context, kind lexutil.XRPCRequestType, inpenc string, method string, params map[string]interface{}, bodyobj interface{}, out interface{}) error { var body io.Reader if bodyobj != nil { if rr, ok := bodyobj.(io.Reader); ok { @@ -150,9 +151,9 @@ func (c *Client) Do(ctx context.Context, kind XRPCRequestType, inpenc string, me var m string switch kind { - case Query: + case lexutil.Query: m = "GET" - case Procedure: + case lexutil.Procedure: m = "POST" default: return fmt.Errorf("unsupported request kind: %d", kind) @@ -227,3 +228,7 @@ func (c *Client) Do(ctx context.Context, kind XRPCRequestType, inpenc string, me return nil } + +func (c *Client) LexDo(ctx context.Context, kind lexutil.XRPCRequestType, inpenc string, method string, params map[string]any, bodyobj any, out any) error { + return c.Do(ctx, kind, inpenc, method, params, bodyobj, out) +} -- 2.51.2 From 9c41058b4406e474efcee543c4da326ba771cf3d Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Sat, 3 May 2025 19:14:53 -0700 Subject: [PATCH 15/17] manually update record-agnostic API helpers as an example --- api/agnostic/actorgetPreferences.go | 6 +++--- api/agnostic/actorputPreferences.go | 6 +++--- api/agnostic/identitygetRecommendedDidCredentials.go | 6 +++--- api/agnostic/identitysignPlcOperation.go | 6 +++--- api/agnostic/identitysubmitPlcOperation.go | 6 +++--- api/agnostic/repoapplyWrites.go | 5 ++--- api/agnostic/repocreateRecord.go | 6 +++--- api/agnostic/repogetRecord.go | 6 +++--- api/agnostic/repolistRecords.go | 6 +++--- api/agnostic/repoputRecord.go | 6 +++--- 10 files changed, 29 insertions(+), 30 deletions(-) diff --git a/api/agnostic/actorgetPreferences.go b/api/agnostic/actorgetPreferences.go index 9ee242c0..150e3e79 100644 --- a/api/agnostic/actorgetPreferences.go +++ b/api/agnostic/actorgetPreferences.go @@ -7,7 +7,7 @@ package agnostic import ( "context" - "github.com/bluesky-social/indigo/xrpc" + "github.com/bluesky-social/indigo/lex/util" ) // ActorGetPreferences_Output is the output of a app.bsky.actor.getPreferences call. @@ -16,11 +16,11 @@ type ActorGetPreferences_Output struct { } // ActorGetPreferences calls the XRPC method "app.bsky.actor.getPreferences". -func ActorGetPreferences(ctx context.Context, c *xrpc.Client) (*ActorGetPreferences_Output, error) { +func ActorGetPreferences(ctx context.Context, c util.LexClient) (*ActorGetPreferences_Output, error) { var out ActorGetPreferences_Output params := map[string]interface{}{} - if err := c.Do(ctx, xrpc.Query, "", "app.bsky.actor.getPreferences", params, nil, &out); err != nil { + if err := c.LexDo(ctx, util.Query, "", "app.bsky.actor.getPreferences", params, nil, &out); err != nil { return nil, err } diff --git a/api/agnostic/actorputPreferences.go b/api/agnostic/actorputPreferences.go index 3648e5bd..02276da2 100644 --- a/api/agnostic/actorputPreferences.go +++ b/api/agnostic/actorputPreferences.go @@ -7,7 +7,7 @@ package agnostic import ( "context" - "github.com/bluesky-social/indigo/xrpc" + "github.com/bluesky-social/indigo/lex/util" ) // ActorPutPreferences_Input is the input argument to a app.bsky.actor.putPreferences call. @@ -16,8 +16,8 @@ type ActorPutPreferences_Input struct { } // ActorPutPreferences calls the XRPC method "app.bsky.actor.putPreferences". -func ActorPutPreferences(ctx context.Context, c *xrpc.Client, input *ActorPutPreferences_Input) error { - if err := c.Do(ctx, xrpc.Procedure, "application/json", "app.bsky.actor.putPreferences", nil, input, nil); err != nil { +func ActorPutPreferences(ctx context.Context, c util.LexClient, input *ActorPutPreferences_Input) error { + if err := c.LexDo(ctx, util.Procedure, "application/json", "app.bsky.actor.putPreferences", nil, input, nil); err != nil { return err } diff --git a/api/agnostic/identitygetRecommendedDidCredentials.go b/api/agnostic/identitygetRecommendedDidCredentials.go index ea0a70a9..85fabcf6 100644 --- a/api/agnostic/identitygetRecommendedDidCredentials.go +++ b/api/agnostic/identitygetRecommendedDidCredentials.go @@ -8,14 +8,14 @@ import ( "context" "encoding/json" - "github.com/bluesky-social/indigo/xrpc" + "github.com/bluesky-social/indigo/lex/util" ) // IdentityGetRecommendedDidCredentials calls the XRPC method "com.atproto.identity.getRecommendedDidCredentials". -func IdentityGetRecommendedDidCredentials(ctx context.Context, c *xrpc.Client) (*json.RawMessage, error) { +func IdentityGetRecommendedDidCredentials(ctx context.Context, c util.LexClient) (*json.RawMessage, error) { var out json.RawMessage - if err := c.Do(ctx, xrpc.Query, "", "com.atproto.identity.getRecommendedDidCredentials", nil, nil, &out); err != nil { + if err := c.LexDo(ctx, util.Query, "", "com.atproto.identity.getRecommendedDidCredentials", nil, nil, &out); err != nil { return nil, err } diff --git a/api/agnostic/identitysignPlcOperation.go b/api/agnostic/identitysignPlcOperation.go index 59d83a23..44c6f732 100644 --- a/api/agnostic/identitysignPlcOperation.go +++ b/api/agnostic/identitysignPlcOperation.go @@ -8,7 +8,7 @@ import ( "context" "encoding/json" - "github.com/bluesky-social/indigo/xrpc" + "github.com/bluesky-social/indigo/lex/util" ) // IdentitySignPlcOperation_Input is the input argument to a com.atproto.identity.signPlcOperation call. @@ -28,9 +28,9 @@ type IdentitySignPlcOperation_Output struct { } // IdentitySignPlcOperation calls the XRPC method "com.atproto.identity.signPlcOperation". -func IdentitySignPlcOperation(ctx context.Context, c *xrpc.Client, input *IdentitySignPlcOperation_Input) (*IdentitySignPlcOperation_Output, error) { +func IdentitySignPlcOperation(ctx context.Context, c util.LexClient, input *IdentitySignPlcOperation_Input) (*IdentitySignPlcOperation_Output, error) { var out IdentitySignPlcOperation_Output - if err := c.Do(ctx, xrpc.Procedure, "application/json", "com.atproto.identity.signPlcOperation", nil, input, &out); err != nil { + if err := c.LexDo(ctx, util.Procedure, "application/json", "com.atproto.identity.signPlcOperation", nil, input, &out); err != nil { return nil, err } diff --git a/api/agnostic/identitysubmitPlcOperation.go b/api/agnostic/identitysubmitPlcOperation.go index 91f9cf0c..ef7e8a0a 100644 --- a/api/agnostic/identitysubmitPlcOperation.go +++ b/api/agnostic/identitysubmitPlcOperation.go @@ -8,7 +8,7 @@ import ( "context" "encoding/json" - "github.com/bluesky-social/indigo/xrpc" + "github.com/bluesky-social/indigo/lex/util" ) // IdentitySubmitPlcOperation_Input is the input argument to a com.atproto.identity.submitPlcOperation call. @@ -17,8 +17,8 @@ type IdentitySubmitPlcOperation_Input struct { } // IdentitySubmitPlcOperation calls the XRPC method "com.atproto.identity.submitPlcOperation". -func IdentitySubmitPlcOperation(ctx context.Context, c *xrpc.Client, input *IdentitySubmitPlcOperation_Input) error { - if err := c.Do(ctx, xrpc.Procedure, "application/json", "com.atproto.identity.submitPlcOperation", nil, input, nil); err != nil { +func IdentitySubmitPlcOperation(ctx context.Context, c util.LexClient, input *IdentitySubmitPlcOperation_Input) error { + if err := c.LexDo(ctx, util.Procedure, "application/json", "com.atproto.identity.submitPlcOperation", nil, input, nil); err != nil { return err } diff --git a/api/agnostic/repoapplyWrites.go b/api/agnostic/repoapplyWrites.go index de0799ac..0dbe127b 100644 --- a/api/agnostic/repoapplyWrites.go +++ b/api/agnostic/repoapplyWrites.go @@ -10,7 +10,6 @@ import ( "fmt" "github.com/bluesky-social/indigo/lex/util" - "github.com/bluesky-social/indigo/xrpc" ) // RepoApplyWrites_Create is a "create" in the com.atproto.repo.applyWrites schema. @@ -179,9 +178,9 @@ type RepoApplyWrites_UpdateResult struct { } // RepoApplyWrites calls the XRPC method "com.atproto.repo.applyWrites". -func RepoApplyWrites(ctx context.Context, c *xrpc.Client, input *RepoApplyWrites_Input) (*RepoApplyWrites_Output, error) { +func RepoApplyWrites(ctx context.Context, c util.LexClient, input *RepoApplyWrites_Input) (*RepoApplyWrites_Output, error) { var out RepoApplyWrites_Output - if err := c.Do(ctx, xrpc.Procedure, "application/json", "com.atproto.repo.applyWrites", nil, input, &out); err != nil { + if err := c.LexDo(ctx, util.Procedure, "application/json", "com.atproto.repo.applyWrites", nil, input, &out); err != nil { return nil, err } diff --git a/api/agnostic/repocreateRecord.go b/api/agnostic/repocreateRecord.go index 9dcddadb..44680b76 100644 --- a/api/agnostic/repocreateRecord.go +++ b/api/agnostic/repocreateRecord.go @@ -7,7 +7,7 @@ package agnostic import ( "context" - "github.com/bluesky-social/indigo/xrpc" + "github.com/bluesky-social/indigo/lex/util" ) // RepoDefs_CommitMeta is a "commitMeta" in the com.atproto.repo.defs schema. @@ -41,9 +41,9 @@ type RepoCreateRecord_Output struct { } // RepoCreateRecord calls the XRPC method "com.atproto.repo.createRecord". -func RepoCreateRecord(ctx context.Context, c *xrpc.Client, input *RepoCreateRecord_Input) (*RepoCreateRecord_Output, error) { +func RepoCreateRecord(ctx context.Context, c util.LexClient, input *RepoCreateRecord_Input) (*RepoCreateRecord_Output, error) { var out RepoCreateRecord_Output - if err := c.Do(ctx, xrpc.Procedure, "application/json", "com.atproto.repo.createRecord", nil, input, &out); err != nil { + if err := c.LexDo(ctx, util.Procedure, "application/json", "com.atproto.repo.createRecord", nil, input, &out); err != nil { return nil, err } diff --git a/api/agnostic/repogetRecord.go b/api/agnostic/repogetRecord.go index 6025c959..0417c79e 100644 --- a/api/agnostic/repogetRecord.go +++ b/api/agnostic/repogetRecord.go @@ -8,7 +8,7 @@ import ( "context" "encoding/json" - "github.com/bluesky-social/indigo/xrpc" + "github.com/bluesky-social/indigo/lex/util" ) // RepoGetRecord_Output is the output of a com.atproto.repo.getRecord call. @@ -25,7 +25,7 @@ type RepoGetRecord_Output struct { // 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) { +func RepoGetRecord(ctx context.Context, c util.LexClient, cid string, collection string, repo string, rkey string) (*RepoGetRecord_Output, error) { var out RepoGetRecord_Output params := map[string]interface{}{ @@ -34,7 +34,7 @@ func RepoGetRecord(ctx context.Context, c *xrpc.Client, cid string, collection s "repo": repo, "rkey": rkey, } - if err := c.Do(ctx, xrpc.Query, "", "com.atproto.repo.getRecord", params, nil, &out); err != nil { + if err := c.LexDo(ctx, util.Query, "", "com.atproto.repo.getRecord", params, nil, &out); err != nil { return nil, err } diff --git a/api/agnostic/repolistRecords.go b/api/agnostic/repolistRecords.go index 2e339250..74b6ba92 100644 --- a/api/agnostic/repolistRecords.go +++ b/api/agnostic/repolistRecords.go @@ -8,7 +8,7 @@ import ( "context" "encoding/json" - "github.com/bluesky-social/indigo/xrpc" + "github.com/bluesky-social/indigo/lex/util" ) // RepoListRecords_Output is the output of a com.atproto.repo.listRecords call. @@ -31,7 +31,7 @@ type RepoListRecords_Record struct { // limit: The number of records to return. // repo: The handle or DID of the repo. // reverse: Flag to reverse the order of the returned records. -func RepoListRecords(ctx context.Context, c *xrpc.Client, collection string, cursor string, limit int64, repo string, reverse bool) (*RepoListRecords_Output, error) { +func RepoListRecords(ctx context.Context, c util.LexClient, collection string, cursor string, limit int64, repo string, reverse bool) (*RepoListRecords_Output, error) { var out RepoListRecords_Output params := map[string]interface{}{ @@ -41,7 +41,7 @@ func RepoListRecords(ctx context.Context, c *xrpc.Client, collection string, cur "repo": repo, "reverse": reverse, } - if err := c.Do(ctx, xrpc.Query, "", "com.atproto.repo.listRecords", params, nil, &out); err != nil { + if err := c.LexDo(ctx, util.Query, "", "com.atproto.repo.listRecords", params, nil, &out); err != nil { return nil, err } diff --git a/api/agnostic/repoputRecord.go b/api/agnostic/repoputRecord.go index 2448793a..0d2d80c4 100644 --- a/api/agnostic/repoputRecord.go +++ b/api/agnostic/repoputRecord.go @@ -7,7 +7,7 @@ package agnostic import ( "context" - "github.com/bluesky-social/indigo/xrpc" + "github.com/bluesky-social/indigo/lex/util" ) // RepoPutRecord_Input is the input argument to a com.atproto.repo.putRecord call. @@ -37,9 +37,9 @@ type RepoPutRecord_Output struct { } // RepoPutRecord calls the XRPC method "com.atproto.repo.putRecord". -func RepoPutRecord(ctx context.Context, c *xrpc.Client, input *RepoPutRecord_Input) (*RepoPutRecord_Output, error) { +func RepoPutRecord(ctx context.Context, c util.LexClient, input *RepoPutRecord_Input) (*RepoPutRecord_Output, error) { var out RepoPutRecord_Output - if err := c.Do(ctx, xrpc.Procedure, "application/json", "com.atproto.repo.putRecord", nil, input, &out); err != nil { + if err := c.LexDo(ctx, util.Procedure, "application/json", "com.atproto.repo.putRecord", nil, input, &out); err != nil { return nil, err } -- 2.51.2 From 054cbe16538fcf7f1cd9ad3fd09455471601351e Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Sat, 31 May 2025 19:29:46 -0700 Subject: [PATCH 16/17] switch to using net/http consts for request types --- lex/util/client.go | 9 ++++----- xrpc/xrpc.go | 17 +++++++---------- 2 files changed, 11 insertions(+), 15 deletions(-) diff --git a/lex/util/client.go b/lex/util/client.go index f40a78f7..3d702aa2 100644 --- a/lex/util/client.go +++ b/lex/util/client.go @@ -2,16 +2,15 @@ package util import ( "context" + "net/http" ) -type XRPCRequestType int - const ( - Query = XRPCRequestType(iota) - Procedure + Query = http.MethodGet + Procedure = http.MethodPost ) // API client interface used in lexgen. type LexClient interface { - LexDo(ctx context.Context, kind XRPCRequestType, inpenc string, method string, params map[string]any, bodyobj any, out any) error + LexDo(ctx context.Context, kind string, inpenc string, method string, params map[string]any, bodyobj any, out any) error } diff --git a/xrpc/xrpc.go b/xrpc/xrpc.go index 2d939fe5..16ea812c 100644 --- a/xrpc/xrpc.go +++ b/xrpc/xrpc.go @@ -13,7 +13,6 @@ import ( "strings" "time" - lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/util" "github.com/carlmjohnson/versioninfo" ) @@ -35,11 +34,9 @@ func (c *Client) getClient() *http.Client { return c.Client } -type XRPCRequestType = lexutil.XRPCRequestType - var ( - Query = lexutil.Query - Procedure = lexutil.Procedure + Query = http.MethodGet + Procedure = http.MethodPost ) type AuthInfo struct { @@ -134,7 +131,7 @@ func makeParams(p map[string]any) string { return params.Encode() } -func (c *Client) Do(ctx context.Context, kind lexutil.XRPCRequestType, inpenc string, method string, params map[string]interface{}, bodyobj interface{}, out interface{}) error { +func (c *Client) Do(ctx context.Context, kind string, inpenc string, method string, params map[string]interface{}, bodyobj interface{}, out interface{}) error { var body io.Reader if bodyobj != nil { if rr, ok := bodyobj.(io.Reader); ok { @@ -151,12 +148,12 @@ func (c *Client) Do(ctx context.Context, kind lexutil.XRPCRequestType, inpenc st var m string switch kind { - case lexutil.Query: + case Query: m = "GET" - case lexutil.Procedure: + case Procedure: m = "POST" default: - return fmt.Errorf("unsupported request kind: %d", kind) + return fmt.Errorf("unsupported request kind: %s", kind) } var paramStr string @@ -229,6 +226,6 @@ func (c *Client) Do(ctx context.Context, kind lexutil.XRPCRequestType, inpenc st return nil } -func (c *Client) LexDo(ctx context.Context, kind lexutil.XRPCRequestType, inpenc string, method string, params map[string]any, bodyobj any, out any) error { +func (c *Client) LexDo(ctx context.Context, kind string, inpenc string, method string, params map[string]any, bodyobj any, out any) error { return c.Do(ctx, kind, inpenc, method, params, bodyobj, out) } -- 2.51.2 From 66566a8f73caa7af0a5d8cbef45e2cf889c5f7bf Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Mon, 2 Jun 2025 23:31:29 -0700 Subject: [PATCH 17/17] better method param names for LexDo() --- lex/util/client.go | 4 +++- xrpc/xrpc.go | 4 ++-- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/lex/util/client.go b/lex/util/client.go index 3d702aa2..3ee3e84f 100644 --- a/lex/util/client.go +++ b/lex/util/client.go @@ -11,6 +11,8 @@ const ( ) // API client interface used in lexgen. +// +// 'method' is the HTTP method type. 'inputEncoding' is the Content-Type for bodyData in Procedure calls. 'params' are query parameters. 'bodyData' should be either 'nil', an [io.Reader], or a type which can be marshalled to JSON. 'out' is optional; if not nil it should be a pointer to a type which can be un-Marshaled as JSON, for the response body. type LexClient interface { - LexDo(ctx context.Context, kind string, inpenc string, method string, params map[string]any, bodyobj any, out any) error + LexDo(ctx context.Context, method string, inputEncoding string, endpoint string, params map[string]any, bodyData any, out any) error } diff --git a/xrpc/xrpc.go b/xrpc/xrpc.go index 16ea812c..bf67897e 100644 --- a/xrpc/xrpc.go +++ b/xrpc/xrpc.go @@ -226,6 +226,6 @@ func (c *Client) Do(ctx context.Context, kind string, inpenc string, method stri return nil } -func (c *Client) LexDo(ctx context.Context, kind string, inpenc string, method string, params map[string]any, bodyobj any, out any) error { - return c.Do(ctx, kind, inpenc, method, params, bodyobj, out) +func (c *Client) LexDo(ctx context.Context, method string, inputEncoding string, endpoint string, params map[string]any, bodyData any, out any) error { + return c.Do(ctx, method, inputEncoding, endpoint, params, bodyData, out) }