Migrate from Firehorse to JetstreamUnverified
aea953a parent: 442e266 modified
.env.example +1 -1 | @@ -1,6 +1,6 @@ | ||
| 1 | 1 | GLEAN_ADDR=:8080 |
| 2 | 2 | GLEAN_DB=glean.db |
| 3 | -GLEAN_RELAY=wss://bsky.network | |
| 3 | +GLEAN_JETSTREAM=wss://jetstream2.fr.hose.cam | |
| 4 | 4 | # Leave empty for localhost OAuth (development) |
| 5 | 5 | # GLEAN_OAUTH_CLIENT_ID=https://glean.at/oauth/client-metadata |
| 6 | 6 | # GLEAN_OAUTH_REDIRECT_URL=https://glean.at/auth/callback |
| @@ -1,6 +1,6 @@ | |||
| 1 | GLEAN_ADDR=:8080 | 1 | GLEAN_ADDR=:8080 |
| 2 | GLEAN_DB=glean.db | 2 | GLEAN_DB=glean.db |
| 3 | -GLEAN_RELAY=wss://bsky.network | 3 | +GLEAN_JETSTREAM=wss://jetstream2.fr.hose.cam |
| 4 | # Leave empty for localhost OAuth (development) | 4 | # Leave empty for localhost OAuth (development) |
| 5 | # GLEAN_OAUTH_CLIENT_ID=https://glean.at/oauth/client-metadata | 5 | # GLEAN_OAUTH_CLIENT_ID=https://glean.at/oauth/client-metadata |
| 6 | # GLEAN_OAUTH_REDIRECT_URL=https://glean.at/auth/callback | 6 | # GLEAN_OAUTH_REDIRECT_URL=https://glean.at/auth/callback |
modified
Makefile +11 -3 | @@ -1,3 +1,7 @@ | ||
| 1 | +.PHONY: tools-install | |
| 2 | +tools-install: | |
| 3 | + go install github.com/golangci/golangci-lint/v2/cmd/golangci-lint@latest | |
| 4 | + | |
| 1 | 5 | .PHONY: lint |
| 2 | 6 | lint: |
| 3 | 7 | go vet ./... |
| @@ -5,9 +9,13 @@ lint: | ||
| 5 | 9 | test -z "$(shell gofmt -l ./...)" |
| 6 | 10 | golangci-lint run ./... --fix |
| 7 | 11 | |
| 8 | -.PHONY: lint-install | |
| 9 | -lint-install: | |
| 10 | - go install github.com/golangci/golangci-lint/v2/cmd/golangci-lint@latest | |
| 12 | +.PHONY: lex-lint | |
| 13 | +lex-lint: | |
| 14 | + goat lex lint | |
| 15 | + | |
| 16 | +.PHONY: lex-parse | |
| 17 | +lex-parse: | |
| 18 | + goat lex parse $(shell find lexicons -name '*.json' 2>/dev/null) | |
| 11 | 19 | |
| 12 | 20 | .PHONY: dev |
| 13 | 21 | dev: css |
| @@ -1,3 +1,7 @@ | |||
| 1 | +.PHONY: tools-install | ||
| 2 | +tools-install: | ||
| 3 | + go install github.com/golangci/golangci-lint/v2/cmd/golangci-lint@latest | ||
| 4 | + | ||
| 1 | .PHONY: lint | 5 | .PHONY: lint |
| 2 | lint: | 6 | lint: |
| 3 | go vet ./... | 7 | go vet ./... |
| @@ -5,9 +9,13 @@ lint: | |||
| 5 | test -z "$(shell gofmt -l ./...)" | 9 | test -z "$(shell gofmt -l ./...)" |
| 6 | golangci-lint run ./... --fix | 10 | golangci-lint run ./... --fix |
| 7 | 11 | ||
| 8 | -.PHONY: lint-install | 12 | +.PHONY: lex-lint |
| 9 | -lint-install: | 13 | +lex-lint: |
| 10 | - go install github.com/golangci/golangci-lint/v2/cmd/golangci-lint@latest | 14 | + goat lex lint |
| 15 | + | ||
| 16 | +.PHONY: lex-parse | ||
| 17 | +lex-parse: | ||
| 18 | + goat lex parse $(shell find lexicons -name '*.json' 2>/dev/null) | ||
| 11 | 19 | ||
| 12 | .PHONY: dev | 20 | .PHONY: dev |
| 13 | dev: css | 21 | dev: css |
modified
go.mod +25 -10 | @@ -4,32 +4,47 @@ go 1.26.2 | ||
| 4 | 4 | |
| 5 | 5 | require ( |
| 6 | 6 | github.com/bluesky-social/indigo v0.0.0-20260417172304-7da09df6081d |
| 7 | + github.com/bluesky-social/jetstream v0.0.0-20260415170838-8a65de4eda28 | |
| 7 | 8 | github.com/go-chi/chi/v5 v5.2.5 |
| 8 | 9 | github.com/go-chi/cors v1.2.2 |
| 9 | - github.com/gorilla/websocket v1.5.3 | |
| 10 | 10 | github.com/mattn/go-sqlite3 v1.14.22 |
| 11 | - github.com/prometheus/client_golang v1.17.0 | |
| 11 | + github.com/prometheus/client_golang v1.19.1 | |
| 12 | + go.uber.org/atomic v1.11.0 | |
| 12 | 13 | golang.org/x/net v0.53.0 |
| 13 | 14 | gotest.tools/v3 v3.5.2 |
| 14 | 15 | ) |
| 15 | 16 | |
| 16 | 17 | require ( |
| 17 | 18 | github.com/beorn7/perks v1.0.1 // indirect |
| 18 | - github.com/cespare/xxhash/v2 v2.2.0 // indirect | |
| 19 | + github.com/cespare/xxhash/v2 v2.3.0 // indirect | |
| 19 | 20 | github.com/earthboundkid/versioninfo/v2 v2.24.1 // indirect |
| 21 | + github.com/goccy/go-json v0.10.2 // indirect | |
| 20 | 22 | github.com/golang-jwt/jwt/v5 v5.2.2 // indirect |
| 21 | - github.com/google/go-cmp v0.5.9 // indirect | |
| 23 | + github.com/google/go-cmp v0.6.0 // indirect | |
| 22 | 24 | github.com/google/go-querystring v1.1.0 // indirect |
| 25 | + github.com/gorilla/websocket v1.5.3 // indirect | |
| 23 | 26 | github.com/hashicorp/golang-lru/v2 v2.0.7 // indirect |
| 24 | - github.com/matttproud/golang_protobuf_extensions/v2 v2.0.0 // indirect | |
| 27 | + github.com/ipfs/go-cid v0.4.1 // indirect | |
| 28 | + github.com/klauspost/compress v1.17.9 // indirect | |
| 29 | + github.com/klauspost/cpuid/v2 v2.2.7 // indirect | |
| 30 | + github.com/minio/sha256-simd v1.0.1 // indirect | |
| 25 | 31 | github.com/mr-tron/base58 v1.2.0 // indirect |
| 26 | - github.com/prometheus/client_model v0.5.0 // indirect | |
| 27 | - github.com/prometheus/common v0.45.0 // indirect | |
| 28 | - github.com/prometheus/procfs v0.12.0 // indirect | |
| 32 | + github.com/multiformats/go-base32 v0.1.0 // indirect | |
| 33 | + github.com/multiformats/go-base36 v0.2.0 // indirect | |
| 34 | + github.com/multiformats/go-multibase v0.2.0 // indirect | |
| 35 | + github.com/multiformats/go-multihash v0.2.3 // indirect | |
| 36 | + github.com/multiformats/go-varint v0.0.7 // indirect | |
| 37 | + github.com/prometheus/client_model v0.6.1 // indirect | |
| 38 | + github.com/prometheus/common v0.54.0 // indirect | |
| 39 | + github.com/prometheus/procfs v0.15.1 // indirect | |
| 40 | + github.com/spaolacci/murmur3 v1.1.0 // indirect | |
| 41 | + github.com/whyrusleeping/cbor-gen v0.2.1-0.20241030202151-b7a6831be65e // indirect | |
| 29 | 42 | gitlab.com/yawning/secp256k1-voi v0.0.0-20230925100816-f2616030848b // indirect |
| 30 | 43 | gitlab.com/yawning/tuplehash v0.0.0-20230713102510-df83abbf9a02 // indirect |
| 31 | 44 | golang.org/x/crypto v0.50.0 // indirect |
| 32 | 45 | golang.org/x/sys v0.43.0 // indirect |
| 33 | - golang.org/x/time v0.3.0 // indirect | |
| 34 | - google.golang.org/protobuf v1.33.0 // indirect | |
| 46 | + golang.org/x/time v0.5.0 // indirect | |
| 47 | + golang.org/x/xerrors v0.0.0-20231012003039-104605ab7028 // indirect | |
| 48 | + google.golang.org/protobuf v1.34.2 // indirect | |
| 49 | + lukechampine.com/blake3 v1.2.1 // indirect | |
| 35 | 50 | ) |
| @@ -4,32 +4,47 @@ go 1.26.2 | |||
| 4 | 4 | ||
| 5 | require ( | 5 | require ( |
| 6 | github.com/bluesky-social/indigo v0.0.0-20260417172304-7da09df6081d | 6 | github.com/bluesky-social/indigo v0.0.0-20260417172304-7da09df6081d |
| 7 | + github.com/bluesky-social/jetstream v0.0.0-20260415170838-8a65de4eda28 | ||
| 7 | github.com/go-chi/chi/v5 v5.2.5 | 8 | github.com/go-chi/chi/v5 v5.2.5 |
| 8 | github.com/go-chi/cors v1.2.2 | 9 | github.com/go-chi/cors v1.2.2 |
| 9 | - github.com/gorilla/websocket v1.5.3 | ||
| 10 | github.com/mattn/go-sqlite3 v1.14.22 | 10 | github.com/mattn/go-sqlite3 v1.14.22 |
| 11 | - github.com/prometheus/client_golang v1.17.0 | 11 | + github.com/prometheus/client_golang v1.19.1 |
| 12 | + go.uber.org/atomic v1.11.0 | ||
| 12 | golang.org/x/net v0.53.0 | 13 | golang.org/x/net v0.53.0 |
| 13 | gotest.tools/v3 v3.5.2 | 14 | gotest.tools/v3 v3.5.2 |
| 14 | ) | 15 | ) |
| 15 | 16 | ||
| 16 | require ( | 17 | require ( |
| 17 | github.com/beorn7/perks v1.0.1 // indirect | 18 | github.com/beorn7/perks v1.0.1 // indirect |
| 18 | - github.com/cespare/xxhash/v2 v2.2.0 // indirect | 19 | + github.com/cespare/xxhash/v2 v2.3.0 // indirect |
| 19 | github.com/earthboundkid/versioninfo/v2 v2.24.1 // indirect | 20 | github.com/earthboundkid/versioninfo/v2 v2.24.1 // indirect |
| 21 | + github.com/goccy/go-json v0.10.2 // indirect | ||
| 20 | github.com/golang-jwt/jwt/v5 v5.2.2 // indirect | 22 | github.com/golang-jwt/jwt/v5 v5.2.2 // indirect |
| 21 | - github.com/google/go-cmp v0.5.9 // indirect | 23 | + github.com/google/go-cmp v0.6.0 // indirect |
| 22 | github.com/google/go-querystring v1.1.0 // indirect | 24 | github.com/google/go-querystring v1.1.0 // indirect |
| 25 | + github.com/gorilla/websocket v1.5.3 // indirect | ||
| 23 | github.com/hashicorp/golang-lru/v2 v2.0.7 // indirect | 26 | github.com/hashicorp/golang-lru/v2 v2.0.7 // indirect |
| 24 | - github.com/matttproud/golang_protobuf_extensions/v2 v2.0.0 // indirect | 27 | + github.com/ipfs/go-cid v0.4.1 // indirect |
| 28 | + github.com/klauspost/compress v1.17.9 // indirect | ||
| 29 | + github.com/klauspost/cpuid/v2 v2.2.7 // indirect | ||
| 30 | + github.com/minio/sha256-simd v1.0.1 // indirect | ||
| 25 | github.com/mr-tron/base58 v1.2.0 // indirect | 31 | github.com/mr-tron/base58 v1.2.0 // indirect |
| 26 | - github.com/prometheus/client_model v0.5.0 // indirect | 32 | + github.com/multiformats/go-base32 v0.1.0 // indirect |
| 27 | - github.com/prometheus/common v0.45.0 // indirect | 33 | + github.com/multiformats/go-base36 v0.2.0 // indirect |
| 28 | - github.com/prometheus/procfs v0.12.0 // indirect | 34 | + github.com/multiformats/go-multibase v0.2.0 // indirect |
| 35 | + github.com/multiformats/go-multihash v0.2.3 // indirect | ||
| 36 | + github.com/multiformats/go-varint v0.0.7 // indirect | ||
| 37 | + github.com/prometheus/client_model v0.6.1 // indirect | ||
| 38 | + github.com/prometheus/common v0.54.0 // indirect | ||
| 39 | + github.com/prometheus/procfs v0.15.1 // indirect | ||
| 40 | + github.com/spaolacci/murmur3 v1.1.0 // indirect | ||
| 41 | + github.com/whyrusleeping/cbor-gen v0.2.1-0.20241030202151-b7a6831be65e // indirect | ||
| 29 | gitlab.com/yawning/secp256k1-voi v0.0.0-20230925100816-f2616030848b // indirect | 42 | gitlab.com/yawning/secp256k1-voi v0.0.0-20230925100816-f2616030848b // indirect |
| 30 | gitlab.com/yawning/tuplehash v0.0.0-20230713102510-df83abbf9a02 // indirect | 43 | gitlab.com/yawning/tuplehash v0.0.0-20230713102510-df83abbf9a02 // indirect |
| 31 | golang.org/x/crypto v0.50.0 // indirect | 44 | golang.org/x/crypto v0.50.0 // indirect |
| 32 | golang.org/x/sys v0.43.0 // indirect | 45 | golang.org/x/sys v0.43.0 // indirect |
| 33 | - golang.org/x/time v0.3.0 // indirect | 46 | + golang.org/x/time v0.5.0 // indirect |
| 34 | - google.golang.org/protobuf v1.33.0 // indirect | 47 | + golang.org/x/xerrors v0.0.0-20231012003039-104605ab7028 // indirect |
| 48 | + google.golang.org/protobuf v1.34.2 // indirect | ||
| 49 | + lukechampine.com/blake3 v1.2.1 // indirect | ||
| 35 | ) | 50 | ) |
modified
go.sum +25 -18 | @@ -2,8 +2,10 @@ github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= | ||
| 2 | 2 | github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= |
| 3 | 3 | github.com/bluesky-social/indigo v0.0.0-20260417172304-7da09df6081d h1:ThKFUrkm2/IZwbvmIKLJYr0wPHibtCkIVmuZCWmdIHM= |
| 4 | 4 | github.com/bluesky-social/indigo v0.0.0-20260417172304-7da09df6081d/go.mod h1:JqQkz8lrOI6YZivP38GHmtVOTtzsNToITKj1gMpU5Jo= |
| 5 | -github.com/cespare/xxhash/v2 v2.2.0 h1:DC2CZ1Ep5Y4k3ZQ899DldepgrayRUGE6BBZ/cd9Cj44= | |
| 6 | -github.com/cespare/xxhash/v2 v2.2.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= | |
| 5 | +github.com/bluesky-social/jetstream v0.0.0-20260415170838-8a65de4eda28 h1:cQ7kasyLcEuh/Zd7g0h8kmaWz14SEjJvbeEPU9rCyx0= | |
| 6 | +github.com/bluesky-social/jetstream v0.0.0-20260415170838-8a65de4eda28/go.mod h1:1TEGvYje9ONndPpaGqW6WqyLiEvYm2iKd2GrpzLvGhg= | |
| 7 | +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= | |
| 8 | +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= | |
| 7 | 9 | github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= |
| 8 | 10 | github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= |
| 9 | 11 | github.com/earthboundkid/versioninfo/v2 v2.24.1 h1:SJTMHaoUx3GzjjnUO1QzP3ZXK6Ee/nbWyCm58eY3oUg= |
| @@ -12,11 +14,13 @@ github.com/go-chi/chi/v5 v5.2.5 h1:Eg4myHZBjyvJmAFjFvWgrqDTXFyOzjj7YIm3L3mu6Ug= | ||
| 12 | 14 | github.com/go-chi/chi/v5 v5.2.5/go.mod h1:X7Gx4mteadT3eDOMTsXzmI4/rwUpOwBHLpAfupzFJP0= |
| 13 | 15 | github.com/go-chi/cors v1.2.2 h1:Jmey33TE+b+rB7fT8MUy1u0I4L+NARQlK6LhzKPSyQE= |
| 14 | 16 | github.com/go-chi/cors v1.2.2/go.mod h1:sSbTewc+6wYHBBCW7ytsFSn836hqM7JxpglAy2Vzc58= |
| 17 | +github.com/goccy/go-json v0.10.2 h1:CrxCmQqYDkv1z7lO7Wbh2HN93uovUHgrECaO5ZrCXAU= | |
| 18 | +github.com/goccy/go-json v0.10.2/go.mod h1:6MelG93GURQebXPDq3khkgXZkazVtN9CRI+MGFi0w8I= | |
| 15 | 19 | github.com/golang-jwt/jwt/v5 v5.2.2 h1:Rl4B7itRWVtYIHFrSNd7vhTiz9UpLdi6gZhZ3wEeDy8= |
| 16 | 20 | github.com/golang-jwt/jwt/v5 v5.2.2/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk= |
| 17 | 21 | github.com/google/go-cmp v0.5.2/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= |
| 18 | -github.com/google/go-cmp v0.5.9 h1:O2Tfq5qg4qc4AmwVlvv0oLiVAGB7enBSJ2x2DqQFi38= | |
| 19 | -github.com/google/go-cmp v0.5.9/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= | |
| 22 | +github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= | |
| 23 | +github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= | |
| 20 | 24 | github.com/google/go-querystring v1.1.0 h1:AnCroh3fv4ZBgVIf1Iwtovgjaw/GiKJo8M8yD/fhyJ8= |
| 21 | 25 | github.com/google/go-querystring v1.1.0/go.mod h1:Kcdr2DB4koayq7X8pmAG4sNG59So17icRSOU623lUBU= |
| 22 | 26 | github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= |
| @@ -25,12 +29,12 @@ github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs | ||
| 25 | 29 | github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= |
| 26 | 30 | github.com/ipfs/go-cid v0.4.1 h1:A/T3qGvxi4kpKWWcPC/PgbvDA2bjVLO7n4UeVwnbs/s= |
| 27 | 31 | github.com/ipfs/go-cid v0.4.1/go.mod h1:uQHwDeX4c6CtyrFwdqyhpNcxVewur1M7l7fNU7LKwZk= |
| 32 | +github.com/klauspost/compress v1.17.9 h1:6KIumPrER1LHsvBVuDa0r5xaG0Es51mhhB9BQB2qeMA= | |
| 33 | +github.com/klauspost/compress v1.17.9/go.mod h1:Di0epgTjJY877eYKx5yC51cX2A2Vl2ibi7bDH9ttBbw= | |
| 28 | 34 | github.com/klauspost/cpuid/v2 v2.2.7 h1:ZWSB3igEs+d0qvnxR/ZBzXVmxkgt8DdzP6m9pfuVLDM= |
| 29 | 35 | github.com/klauspost/cpuid/v2 v2.2.7/go.mod h1:Lcz8mBdAVJIBVzewtcLocK12l3Y+JytZYpaMropDUws= |
| 30 | 36 | github.com/mattn/go-sqlite3 v1.14.22 h1:2gZY6PC6kBnID23Tichd1K+Z0oS6nE/XwU+Vz/5o4kU= |
| 31 | 37 | github.com/mattn/go-sqlite3 v1.14.22/go.mod h1:Uh1q+B4BYcTPb+yiD3kU8Ct7aC0hY9fxUwlHK0RXw+Y= |
| 32 | -github.com/matttproud/golang_protobuf_extensions/v2 v2.0.0 h1:jWpvCLoY8Z/e3VKvlsiIGKtc+UG6U5vzxaoagmhXfyg= | |
| 33 | -github.com/matttproud/golang_protobuf_extensions/v2 v2.0.0/go.mod h1:QUyp042oQthUoa9bqDv0ER0wrtXnBruoNd7aNjkbP+k= | |
| 34 | 38 | github.com/minio/sha256-simd v1.0.1 h1:6kaan5IFmwTNynnKKpDHe6FWHohJOHhCPchzK49dzMM= |
| 35 | 39 | github.com/minio/sha256-simd v1.0.1/go.mod h1:Pz6AKMiUdngCLpeTL/RJY1M9rUuPMYujV5xJjtbRSN8= |
| 36 | 40 | github.com/mr-tron/base58 v1.2.0 h1:T/HDJBh4ZCPbU39/+c3rRvE0uKBQlU27+QI8LJ4t64o= |
| @@ -47,14 +51,14 @@ github.com/multiformats/go-varint v0.0.7 h1:sWSGR+f/eu5ABZA2ZpYKBILXTTs9JWpdEM/n | ||
| 47 | 51 | github.com/multiformats/go-varint v0.0.7/go.mod h1:r8PUYw/fD/SjBCiKOoDlGF6QawOELpZAu9eioSos/OU= |
| 48 | 52 | github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= |
| 49 | 53 | github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= |
| 50 | -github.com/prometheus/client_golang v1.17.0 h1:rl2sfwZMtSthVU752MqfjQozy7blglC+1SOtjMAMh+Q= | |
| 51 | -github.com/prometheus/client_golang v1.17.0/go.mod h1:VeL+gMmOAxkS2IqfCq0ZmHSL+LjWfWDUmp1mBz9JgUY= | |
| 52 | -github.com/prometheus/client_model v0.5.0 h1:VQw1hfvPvk3Uv6Qf29VrPF32JB6rtbgI6cYPYQjL0Qw= | |
| 53 | -github.com/prometheus/client_model v0.5.0/go.mod h1:dTiFglRmd66nLR9Pv9f0mZi7B7fk5Pm3gvsjB5tr+kI= | |
| 54 | -github.com/prometheus/common v0.45.0 h1:2BGz0eBc2hdMDLnO/8n0jeB3oPrt2D08CekT0lneoxM= | |
| 55 | -github.com/prometheus/common v0.45.0/go.mod h1:YJmSTw9BoKxJplESWWxlbyttQR4uaEcGyv9MZjVOJsY= | |
| 56 | -github.com/prometheus/procfs v0.12.0 h1:jluTpSng7V9hY0O2R9DzzJHYb2xULk9VTR1V1R/k6Bo= | |
| 57 | -github.com/prometheus/procfs v0.12.0/go.mod h1:pcuDEFsWDnvcgNzo4EEweacyhjeA9Zk3cnaOZAZEfOo= | |
| 54 | +github.com/prometheus/client_golang v1.19.1 h1:wZWJDwK+NameRJuPGDhlnFgx8e8HN3XHQeLaYJFJBOE= | |
| 55 | +github.com/prometheus/client_golang v1.19.1/go.mod h1:mP78NwGzrVks5S2H6ab8+ZZGJLZUq1hoULYBAYBw1Ho= | |
| 56 | +github.com/prometheus/client_model v0.6.1 h1:ZKSh/rekM+n3CeS952MLRAdFwIKqeY8b62p8ais2e9E= | |
| 57 | +github.com/prometheus/client_model v0.6.1/go.mod h1:OrxVMOVHjw3lKMa8+x6HeMGkHMQyHDk9E3jmP2AmGiY= | |
| 58 | +github.com/prometheus/common v0.54.0 h1:ZlZy0BgJhTwVZUn7dLOkwCZHUkrAqd3WYtcFCWnM1D8= | |
| 59 | +github.com/prometheus/common v0.54.0/go.mod h1:/TQgMJP5CuVYveyT7n/0Ix8yLNNXy9yRSkhnLTHPDIQ= | |
| 60 | +github.com/prometheus/procfs v0.15.1 h1:YagwOFzUgYfKKHX6Dr+sHT7km/hxC76UB0learggepc= | |
| 61 | +github.com/prometheus/procfs v0.15.1/go.mod h1:fB45yRUv8NstnjriLhBQLuOUt+WW4BsoGhij/e3PBqk= | |
| 58 | 62 | github.com/spaolacci/murmur3 v1.1.0 h1:7c1g84S4BPRrfL5Xrdp6fOJ206sU9y293DDHaoy0bLI= |
| 59 | 63 | github.com/spaolacci/murmur3 v1.1.0/go.mod h1:JwIasOWyU6f++ZhiEuf87xNszmSA2myDM2Kzu9HwQUA= |
| 60 | 64 | github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= |
| @@ -65,19 +69,22 @@ gitlab.com/yawning/secp256k1-voi v0.0.0-20230925100816-f2616030848b h1:CzigHMRyS | ||
| 65 | 69 | gitlab.com/yawning/secp256k1-voi v0.0.0-20230925100816-f2616030848b/go.mod h1:/y/V339mxv2sZmYYR64O07VuCpdNZqCTwO8ZcouTMI8= |
| 66 | 70 | gitlab.com/yawning/tuplehash v0.0.0-20230713102510-df83abbf9a02 h1:qwDnMxjkyLmAFgcfgTnfJrmYKWhHnci3GjDqcZp1M3Q= |
| 67 | 71 | gitlab.com/yawning/tuplehash v0.0.0-20230713102510-df83abbf9a02/go.mod h1:JTnUj0mpYiAsuZLmKjTx/ex3AtMowcCgnE7YNyCEP0I= |
| 72 | +go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= | |
| 73 | +go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= | |
| 68 | 74 | golang.org/x/crypto v0.50.0 h1:zO47/JPrL6vsNkINmLoo/PH1gcxpls50DNogFvB5ZGI= |
| 69 | 75 | golang.org/x/crypto v0.50.0/go.mod h1:3muZ7vA7PBCE6xgPX7nkzzjiUq87kRItoJQM1Yo8S+Q= |
| 70 | 76 | golang.org/x/net v0.53.0 h1:d+qAbo5L0orcWAr0a9JweQpjXF19LMXJE8Ey7hwOdUA= |
| 71 | 77 | golang.org/x/net v0.53.0/go.mod h1:JvMuJH7rrdiCfbeHoo3fCQU24Lf5JJwT9W3sJFulfgs= |
| 78 | +golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= | |
| 72 | 79 | golang.org/x/sys v0.43.0 h1:Rlag2XtaFTxp19wS8MXlJwTvoh8ArU6ezoyFsMyCTNI= |
| 73 | 80 | golang.org/x/sys v0.43.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= |
| 74 | -golang.org/x/time v0.3.0 h1:rg5rLMjNzMS1RkNLzCG38eapWhnYLFYXDXj2gOlr8j4= | |
| 75 | -golang.org/x/time v0.3.0/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= | |
| 81 | +golang.org/x/time v0.5.0 h1:o7cqy6amK/52YcAKIPlM3a+Fpj35zvRj2TP+e1xFSfk= | |
| 82 | +golang.org/x/time v0.5.0/go.mod h1:3BpzKBy/shNhVucY/MWOyx10tF3SFh9QdLuxbVysPQM= | |
| 76 | 83 | golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= |
| 77 | 84 | golang.org/x/xerrors v0.0.0-20231012003039-104605ab7028 h1:+cNy6SZtPcJQH3LJVLOSmiC7MMxXNOb3PU/VUEz+EhU= |
| 78 | 85 | golang.org/x/xerrors v0.0.0-20231012003039-104605ab7028/go.mod h1:NDW/Ps6MPRej6fsCIbMTohpP40sJ/P/vI1MoTEGwX90= |
| 79 | -google.golang.org/protobuf v1.33.0 h1:uNO2rsAINq/JlFpSdYEKIZ0uKD/R9cpdv0T+yoGwGmI= | |
| 80 | -google.golang.org/protobuf v1.33.0/go.mod h1:c6P6GXX6sHbq/GpV6MGZEdwhWPcYBgnhAHhKbcUYpos= | |
| 86 | +google.golang.org/protobuf v1.34.2 h1:6xV6lTsCfpGD21XK49h7MhtcApnLqkfYgPcdHftf6hg= | |
| 87 | +google.golang.org/protobuf v1.34.2/go.mod h1:qYOHts0dSfpeUzUFpOMr/WGzszTmLH+DiWniOlNbLDw= | |
| 81 | 88 | gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= |
| 82 | 89 | gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= |
| 83 | 90 | gotest.tools/v3 v3.5.2 h1:7koQfIKdy+I8UTetycgUqXWSDwpgv193Ka+qRsmBY8Q= |
| @@ -2,8 +2,10 @@ github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= | |||
| 2 | github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= | 2 | github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= |
| 3 | github.com/bluesky-social/indigo v0.0.0-20260417172304-7da09df6081d h1:ThKFUrkm2/IZwbvmIKLJYr0wPHibtCkIVmuZCWmdIHM= | 3 | github.com/bluesky-social/indigo v0.0.0-20260417172304-7da09df6081d h1:ThKFUrkm2/IZwbvmIKLJYr0wPHibtCkIVmuZCWmdIHM= |
| 4 | github.com/bluesky-social/indigo v0.0.0-20260417172304-7da09df6081d/go.mod h1:JqQkz8lrOI6YZivP38GHmtVOTtzsNToITKj1gMpU5Jo= | 4 | github.com/bluesky-social/indigo v0.0.0-20260417172304-7da09df6081d/go.mod h1:JqQkz8lrOI6YZivP38GHmtVOTtzsNToITKj1gMpU5Jo= |
| 5 | -github.com/cespare/xxhash/v2 v2.2.0 h1:DC2CZ1Ep5Y4k3ZQ899DldepgrayRUGE6BBZ/cd9Cj44= | 5 | +github.com/bluesky-social/jetstream v0.0.0-20260415170838-8a65de4eda28 h1:cQ7kasyLcEuh/Zd7g0h8kmaWz14SEjJvbeEPU9rCyx0= |
| 6 | -github.com/cespare/xxhash/v2 v2.2.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= | 6 | +github.com/bluesky-social/jetstream v0.0.0-20260415170838-8a65de4eda28/go.mod h1:1TEGvYje9ONndPpaGqW6WqyLiEvYm2iKd2GrpzLvGhg= |
| 7 | +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= | ||
| 8 | +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= | ||
| 7 | github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= | 9 | github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= |
| 8 | github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= | 10 | github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= |
| 9 | github.com/earthboundkid/versioninfo/v2 v2.24.1 h1:SJTMHaoUx3GzjjnUO1QzP3ZXK6Ee/nbWyCm58eY3oUg= | 11 | github.com/earthboundkid/versioninfo/v2 v2.24.1 h1:SJTMHaoUx3GzjjnUO1QzP3ZXK6Ee/nbWyCm58eY3oUg= |
| @@ -12,11 +14,13 @@ github.com/go-chi/chi/v5 v5.2.5 h1:Eg4myHZBjyvJmAFjFvWgrqDTXFyOzjj7YIm3L3mu6Ug= | |||
| 12 | github.com/go-chi/chi/v5 v5.2.5/go.mod h1:X7Gx4mteadT3eDOMTsXzmI4/rwUpOwBHLpAfupzFJP0= | 14 | github.com/go-chi/chi/v5 v5.2.5/go.mod h1:X7Gx4mteadT3eDOMTsXzmI4/rwUpOwBHLpAfupzFJP0= |
| 13 | github.com/go-chi/cors v1.2.2 h1:Jmey33TE+b+rB7fT8MUy1u0I4L+NARQlK6LhzKPSyQE= | 15 | github.com/go-chi/cors v1.2.2 h1:Jmey33TE+b+rB7fT8MUy1u0I4L+NARQlK6LhzKPSyQE= |
| 14 | github.com/go-chi/cors v1.2.2/go.mod h1:sSbTewc+6wYHBBCW7ytsFSn836hqM7JxpglAy2Vzc58= | 16 | github.com/go-chi/cors v1.2.2/go.mod h1:sSbTewc+6wYHBBCW7ytsFSn836hqM7JxpglAy2Vzc58= |
| 17 | +github.com/goccy/go-json v0.10.2 h1:CrxCmQqYDkv1z7lO7Wbh2HN93uovUHgrECaO5ZrCXAU= | ||
| 18 | +github.com/goccy/go-json v0.10.2/go.mod h1:6MelG93GURQebXPDq3khkgXZkazVtN9CRI+MGFi0w8I= | ||
| 15 | github.com/golang-jwt/jwt/v5 v5.2.2 h1:Rl4B7itRWVtYIHFrSNd7vhTiz9UpLdi6gZhZ3wEeDy8= | 19 | github.com/golang-jwt/jwt/v5 v5.2.2 h1:Rl4B7itRWVtYIHFrSNd7vhTiz9UpLdi6gZhZ3wEeDy8= |
| 16 | github.com/golang-jwt/jwt/v5 v5.2.2/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk= | 20 | github.com/golang-jwt/jwt/v5 v5.2.2/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk= |
| 17 | github.com/google/go-cmp v0.5.2/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= | 21 | github.com/google/go-cmp v0.5.2/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= |
| 18 | -github.com/google/go-cmp v0.5.9 h1:O2Tfq5qg4qc4AmwVlvv0oLiVAGB7enBSJ2x2DqQFi38= | 22 | +github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= |
| 19 | -github.com/google/go-cmp v0.5.9/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= | 23 | +github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= |
| 20 | github.com/google/go-querystring v1.1.0 h1:AnCroh3fv4ZBgVIf1Iwtovgjaw/GiKJo8M8yD/fhyJ8= | 24 | github.com/google/go-querystring v1.1.0 h1:AnCroh3fv4ZBgVIf1Iwtovgjaw/GiKJo8M8yD/fhyJ8= |
| 21 | github.com/google/go-querystring v1.1.0/go.mod h1:Kcdr2DB4koayq7X8pmAG4sNG59So17icRSOU623lUBU= | 25 | github.com/google/go-querystring v1.1.0/go.mod h1:Kcdr2DB4koayq7X8pmAG4sNG59So17icRSOU623lUBU= |
| 22 | github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= | 26 | github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= |
| @@ -25,12 +29,12 @@ github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs | |||
| 25 | github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= | 29 | github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= |
| 26 | github.com/ipfs/go-cid v0.4.1 h1:A/T3qGvxi4kpKWWcPC/PgbvDA2bjVLO7n4UeVwnbs/s= | 30 | github.com/ipfs/go-cid v0.4.1 h1:A/T3qGvxi4kpKWWcPC/PgbvDA2bjVLO7n4UeVwnbs/s= |
| 27 | github.com/ipfs/go-cid v0.4.1/go.mod h1:uQHwDeX4c6CtyrFwdqyhpNcxVewur1M7l7fNU7LKwZk= | 31 | github.com/ipfs/go-cid v0.4.1/go.mod h1:uQHwDeX4c6CtyrFwdqyhpNcxVewur1M7l7fNU7LKwZk= |
| 32 | +github.com/klauspost/compress v1.17.9 h1:6KIumPrER1LHsvBVuDa0r5xaG0Es51mhhB9BQB2qeMA= | ||
| 33 | +github.com/klauspost/compress v1.17.9/go.mod h1:Di0epgTjJY877eYKx5yC51cX2A2Vl2ibi7bDH9ttBbw= | ||
| 28 | github.com/klauspost/cpuid/v2 v2.2.7 h1:ZWSB3igEs+d0qvnxR/ZBzXVmxkgt8DdzP6m9pfuVLDM= | 34 | github.com/klauspost/cpuid/v2 v2.2.7 h1:ZWSB3igEs+d0qvnxR/ZBzXVmxkgt8DdzP6m9pfuVLDM= |
| 29 | github.com/klauspost/cpuid/v2 v2.2.7/go.mod h1:Lcz8mBdAVJIBVzewtcLocK12l3Y+JytZYpaMropDUws= | 35 | github.com/klauspost/cpuid/v2 v2.2.7/go.mod h1:Lcz8mBdAVJIBVzewtcLocK12l3Y+JytZYpaMropDUws= |
| 30 | github.com/mattn/go-sqlite3 v1.14.22 h1:2gZY6PC6kBnID23Tichd1K+Z0oS6nE/XwU+Vz/5o4kU= | 36 | github.com/mattn/go-sqlite3 v1.14.22 h1:2gZY6PC6kBnID23Tichd1K+Z0oS6nE/XwU+Vz/5o4kU= |
| 31 | github.com/mattn/go-sqlite3 v1.14.22/go.mod h1:Uh1q+B4BYcTPb+yiD3kU8Ct7aC0hY9fxUwlHK0RXw+Y= | 37 | github.com/mattn/go-sqlite3 v1.14.22/go.mod h1:Uh1q+B4BYcTPb+yiD3kU8Ct7aC0hY9fxUwlHK0RXw+Y= |
| 32 | -github.com/matttproud/golang_protobuf_extensions/v2 v2.0.0 h1:jWpvCLoY8Z/e3VKvlsiIGKtc+UG6U5vzxaoagmhXfyg= | ||
| 33 | -github.com/matttproud/golang_protobuf_extensions/v2 v2.0.0/go.mod h1:QUyp042oQthUoa9bqDv0ER0wrtXnBruoNd7aNjkbP+k= | ||
| 34 | github.com/minio/sha256-simd v1.0.1 h1:6kaan5IFmwTNynnKKpDHe6FWHohJOHhCPchzK49dzMM= | 38 | github.com/minio/sha256-simd v1.0.1 h1:6kaan5IFmwTNynnKKpDHe6FWHohJOHhCPchzK49dzMM= |
| 35 | github.com/minio/sha256-simd v1.0.1/go.mod h1:Pz6AKMiUdngCLpeTL/RJY1M9rUuPMYujV5xJjtbRSN8= | 39 | github.com/minio/sha256-simd v1.0.1/go.mod h1:Pz6AKMiUdngCLpeTL/RJY1M9rUuPMYujV5xJjtbRSN8= |
| 36 | github.com/mr-tron/base58 v1.2.0 h1:T/HDJBh4ZCPbU39/+c3rRvE0uKBQlU27+QI8LJ4t64o= | 40 | github.com/mr-tron/base58 v1.2.0 h1:T/HDJBh4ZCPbU39/+c3rRvE0uKBQlU27+QI8LJ4t64o= |
| @@ -47,14 +51,14 @@ github.com/multiformats/go-varint v0.0.7 h1:sWSGR+f/eu5ABZA2ZpYKBILXTTs9JWpdEM/n | |||
| 47 | github.com/multiformats/go-varint v0.0.7/go.mod h1:r8PUYw/fD/SjBCiKOoDlGF6QawOELpZAu9eioSos/OU= | 51 | github.com/multiformats/go-varint v0.0.7/go.mod h1:r8PUYw/fD/SjBCiKOoDlGF6QawOELpZAu9eioSos/OU= |
| 48 | github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= | 52 | github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= |
| 49 | github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= | 53 | github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= |
| 50 | -github.com/prometheus/client_golang v1.17.0 h1:rl2sfwZMtSthVU752MqfjQozy7blglC+1SOtjMAMh+Q= | 54 | +github.com/prometheus/client_golang v1.19.1 h1:wZWJDwK+NameRJuPGDhlnFgx8e8HN3XHQeLaYJFJBOE= |
| 51 | -github.com/prometheus/client_golang v1.17.0/go.mod h1:VeL+gMmOAxkS2IqfCq0ZmHSL+LjWfWDUmp1mBz9JgUY= | 55 | +github.com/prometheus/client_golang v1.19.1/go.mod h1:mP78NwGzrVks5S2H6ab8+ZZGJLZUq1hoULYBAYBw1Ho= |
| 52 | -github.com/prometheus/client_model v0.5.0 h1:VQw1hfvPvk3Uv6Qf29VrPF32JB6rtbgI6cYPYQjL0Qw= | 56 | +github.com/prometheus/client_model v0.6.1 h1:ZKSh/rekM+n3CeS952MLRAdFwIKqeY8b62p8ais2e9E= |
| 53 | -github.com/prometheus/client_model v0.5.0/go.mod h1:dTiFglRmd66nLR9Pv9f0mZi7B7fk5Pm3gvsjB5tr+kI= | 57 | +github.com/prometheus/client_model v0.6.1/go.mod h1:OrxVMOVHjw3lKMa8+x6HeMGkHMQyHDk9E3jmP2AmGiY= |
| 54 | -github.com/prometheus/common v0.45.0 h1:2BGz0eBc2hdMDLnO/8n0jeB3oPrt2D08CekT0lneoxM= | 58 | +github.com/prometheus/common v0.54.0 h1:ZlZy0BgJhTwVZUn7dLOkwCZHUkrAqd3WYtcFCWnM1D8= |
| 55 | -github.com/prometheus/common v0.45.0/go.mod h1:YJmSTw9BoKxJplESWWxlbyttQR4uaEcGyv9MZjVOJsY= | 59 | +github.com/prometheus/common v0.54.0/go.mod h1:/TQgMJP5CuVYveyT7n/0Ix8yLNNXy9yRSkhnLTHPDIQ= |
| 56 | -github.com/prometheus/procfs v0.12.0 h1:jluTpSng7V9hY0O2R9DzzJHYb2xULk9VTR1V1R/k6Bo= | 60 | +github.com/prometheus/procfs v0.15.1 h1:YagwOFzUgYfKKHX6Dr+sHT7km/hxC76UB0learggepc= |
| 57 | -github.com/prometheus/procfs v0.12.0/go.mod h1:pcuDEFsWDnvcgNzo4EEweacyhjeA9Zk3cnaOZAZEfOo= | 61 | +github.com/prometheus/procfs v0.15.1/go.mod h1:fB45yRUv8NstnjriLhBQLuOUt+WW4BsoGhij/e3PBqk= |
| 58 | github.com/spaolacci/murmur3 v1.1.0 h1:7c1g84S4BPRrfL5Xrdp6fOJ206sU9y293DDHaoy0bLI= | 62 | github.com/spaolacci/murmur3 v1.1.0 h1:7c1g84S4BPRrfL5Xrdp6fOJ206sU9y293DDHaoy0bLI= |
| 59 | github.com/spaolacci/murmur3 v1.1.0/go.mod h1:JwIasOWyU6f++ZhiEuf87xNszmSA2myDM2Kzu9HwQUA= | 63 | github.com/spaolacci/murmur3 v1.1.0/go.mod h1:JwIasOWyU6f++ZhiEuf87xNszmSA2myDM2Kzu9HwQUA= |
| 60 | github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= | 64 | github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= |
| @@ -65,19 +69,22 @@ gitlab.com/yawning/secp256k1-voi v0.0.0-20230925100816-f2616030848b h1:CzigHMRyS | |||
| 65 | gitlab.com/yawning/secp256k1-voi v0.0.0-20230925100816-f2616030848b/go.mod h1:/y/V339mxv2sZmYYR64O07VuCpdNZqCTwO8ZcouTMI8= | 69 | gitlab.com/yawning/secp256k1-voi v0.0.0-20230925100816-f2616030848b/go.mod h1:/y/V339mxv2sZmYYR64O07VuCpdNZqCTwO8ZcouTMI8= |
| 66 | gitlab.com/yawning/tuplehash v0.0.0-20230713102510-df83abbf9a02 h1:qwDnMxjkyLmAFgcfgTnfJrmYKWhHnci3GjDqcZp1M3Q= | 70 | gitlab.com/yawning/tuplehash v0.0.0-20230713102510-df83abbf9a02 h1:qwDnMxjkyLmAFgcfgTnfJrmYKWhHnci3GjDqcZp1M3Q= |
| 67 | gitlab.com/yawning/tuplehash v0.0.0-20230713102510-df83abbf9a02/go.mod h1:JTnUj0mpYiAsuZLmKjTx/ex3AtMowcCgnE7YNyCEP0I= | 71 | gitlab.com/yawning/tuplehash v0.0.0-20230713102510-df83abbf9a02/go.mod h1:JTnUj0mpYiAsuZLmKjTx/ex3AtMowcCgnE7YNyCEP0I= |
| 72 | +go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= | ||
| 73 | +go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= | ||
| 68 | golang.org/x/crypto v0.50.0 h1:zO47/JPrL6vsNkINmLoo/PH1gcxpls50DNogFvB5ZGI= | 74 | golang.org/x/crypto v0.50.0 h1:zO47/JPrL6vsNkINmLoo/PH1gcxpls50DNogFvB5ZGI= |
| 69 | golang.org/x/crypto v0.50.0/go.mod h1:3muZ7vA7PBCE6xgPX7nkzzjiUq87kRItoJQM1Yo8S+Q= | 75 | golang.org/x/crypto v0.50.0/go.mod h1:3muZ7vA7PBCE6xgPX7nkzzjiUq87kRItoJQM1Yo8S+Q= |
| 70 | golang.org/x/net v0.53.0 h1:d+qAbo5L0orcWAr0a9JweQpjXF19LMXJE8Ey7hwOdUA= | 76 | golang.org/x/net v0.53.0 h1:d+qAbo5L0orcWAr0a9JweQpjXF19LMXJE8Ey7hwOdUA= |
| 71 | golang.org/x/net v0.53.0/go.mod h1:JvMuJH7rrdiCfbeHoo3fCQU24Lf5JJwT9W3sJFulfgs= | 77 | golang.org/x/net v0.53.0/go.mod h1:JvMuJH7rrdiCfbeHoo3fCQU24Lf5JJwT9W3sJFulfgs= |
| 78 | +golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= | ||
| 72 | golang.org/x/sys v0.43.0 h1:Rlag2XtaFTxp19wS8MXlJwTvoh8ArU6ezoyFsMyCTNI= | 79 | golang.org/x/sys v0.43.0 h1:Rlag2XtaFTxp19wS8MXlJwTvoh8ArU6ezoyFsMyCTNI= |
| 73 | golang.org/x/sys v0.43.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= | 80 | golang.org/x/sys v0.43.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= |
| 74 | -golang.org/x/time v0.3.0 h1:rg5rLMjNzMS1RkNLzCG38eapWhnYLFYXDXj2gOlr8j4= | 81 | +golang.org/x/time v0.5.0 h1:o7cqy6amK/52YcAKIPlM3a+Fpj35zvRj2TP+e1xFSfk= |
| 75 | -golang.org/x/time v0.3.0/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= | 82 | +golang.org/x/time v0.5.0/go.mod h1:3BpzKBy/shNhVucY/MWOyx10tF3SFh9QdLuxbVysPQM= |
| 76 | golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= | 83 | golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= |
| 77 | golang.org/x/xerrors v0.0.0-20231012003039-104605ab7028 h1:+cNy6SZtPcJQH3LJVLOSmiC7MMxXNOb3PU/VUEz+EhU= | 84 | golang.org/x/xerrors v0.0.0-20231012003039-104605ab7028 h1:+cNy6SZtPcJQH3LJVLOSmiC7MMxXNOb3PU/VUEz+EhU= |
| 78 | golang.org/x/xerrors v0.0.0-20231012003039-104605ab7028/go.mod h1:NDW/Ps6MPRej6fsCIbMTohpP40sJ/P/vI1MoTEGwX90= | 85 | golang.org/x/xerrors v0.0.0-20231012003039-104605ab7028/go.mod h1:NDW/Ps6MPRej6fsCIbMTohpP40sJ/P/vI1MoTEGwX90= |
| 79 | -google.golang.org/protobuf v1.33.0 h1:uNO2rsAINq/JlFpSdYEKIZ0uKD/R9cpdv0T+yoGwGmI= | 86 | +google.golang.org/protobuf v1.34.2 h1:6xV6lTsCfpGD21XK49h7MhtcApnLqkfYgPcdHftf6hg= |
| 80 | -google.golang.org/protobuf v1.33.0/go.mod h1:c6P6GXX6sHbq/GpV6MGZEdwhWPcYBgnhAHhKbcUYpos= | 87 | +google.golang.org/protobuf v1.34.2/go.mod h1:qYOHts0dSfpeUzUFpOMr/WGzszTmLH+DiWniOlNbLDw= |
| 81 | gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= | 88 | gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= |
| 82 | gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= | 89 | gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= |
| 83 | gotest.tools/v3 v3.5.2 h1:7koQfIKdy+I8UTetycgUqXWSDwpgv193Ka+qRsmBY8Q= | 90 | gotest.tools/v3 v3.5.2 h1:7koQfIKdy+I8UTetycgUqXWSDwpgv193Ka+qRsmBY8Q= |
deleted
internal/atproto/firehose.go +0 -209 | deleted file mode 100644 | ||
| @@ -1,209 +0,0 @@ | ||
| 1 | -package atproto | |
| 2 | - | |
| 3 | -import ( | |
| 4 | - "context" | |
| 5 | - "encoding/json" | |
| 6 | - "fmt" | |
| 7 | - "log/slog" | |
| 8 | - "net/http" | |
| 9 | - "net/url" | |
| 10 | - "time" | |
| 11 | - | |
| 12 | - "github.com/gorilla/websocket" | |
| 13 | - | |
| 14 | - "pkg.rbrt.fr/glean/internal/metrics" | |
| 15 | -) | |
| 16 | - | |
| 17 | -type FirehoseEvent struct { | |
| 18 | - Type string | |
| 19 | - DID string | |
| 20 | - Collection string | |
| 21 | - RKey string | |
| 22 | - URI string | |
| 23 | - CID string | |
| 24 | - Value json.RawMessage | |
| 25 | -} | |
| 26 | - | |
| 27 | -type FirehoseHandler func(ctx context.Context, event *FirehoseEvent) error | |
| 28 | - | |
| 29 | -type FirehoseConsumer struct { | |
| 30 | - relayURL string | |
| 31 | - handler FirehoseHandler | |
| 32 | - cursor int64 | |
| 33 | - logger *slog.Logger | |
| 34 | - collections map[string]bool | |
| 35 | -} | |
| 36 | - | |
| 37 | -func NewFirehoseConsumer(relayURL string, handler FirehoseHandler, logger *slog.Logger) *FirehoseConsumer { | |
| 38 | - return &FirehoseConsumer{ | |
| 39 | - relayURL: relayURL, | |
| 40 | - handler: handler, | |
| 41 | - logger: logger, | |
| 42 | - collections: map[string]bool{ | |
| 43 | - "at.glean.subscription": true, | |
| 44 | - "at.glean.annotation": true, | |
| 45 | - "at.glean.like": true, | |
| 46 | - "app.bsky.graph.follow": true, | |
| 47 | - "sh.tangled.graph.follow": true, | |
| 48 | - }, | |
| 49 | - } | |
| 50 | -} | |
| 51 | - | |
| 52 | -func (fc *FirehoseConsumer) Start(ctx context.Context) error { | |
| 53 | - for { | |
| 54 | - err := fc.connect(ctx) | |
| 55 | - if ctx.Err() != nil { | |
| 56 | - return ctx.Err() | |
| 57 | - } | |
| 58 | - if err != nil { | |
| 59 | - fc.logger.Error("firehose connection error", "error", err) | |
| 60 | - metrics.FirehoseReconnects.Inc() | |
| 61 | - } | |
| 62 | - | |
| 63 | - select { | |
| 64 | - case <-ctx.Done(): | |
| 65 | - return ctx.Err() | |
| 66 | - case <-time.After(5 * time.Second): | |
| 67 | - } | |
| 68 | - } | |
| 69 | -} | |
| 70 | - | |
| 71 | -func (fc *FirehoseConsumer) connect(ctx context.Context) error { | |
| 72 | - u, err := url.Parse(fc.relayURL) | |
| 73 | - if err != nil { | |
| 74 | - return fmt.Errorf("parsing relay URL: %w", err) | |
| 75 | - } | |
| 76 | - | |
| 77 | - scheme := "wss" | |
| 78 | - if u.Scheme == "http" || u.Scheme == "ws" { | |
| 79 | - scheme = "ws" | |
| 80 | - } | |
| 81 | - | |
| 82 | - wsURL := fmt.Sprintf("%s://%s/xrpc/com.atproto.sync.subscribeRepos", scheme, u.Host) | |
| 83 | - if fc.cursor > 0 { | |
| 84 | - wsURL += fmt.Sprintf("?cursor=%d", fc.cursor) | |
| 85 | - } | |
| 86 | - | |
| 87 | - fc.logger.Info("connecting to firehose", "url", wsURL) | |
| 88 | - | |
| 89 | - conn, _, err := websocket.DefaultDialer.DialContext(ctx, wsURL, http.Header{}) | |
| 90 | - if err != nil { | |
| 91 | - return fmt.Errorf("dialing firehose: %w", err) | |
| 92 | - } | |
| 93 | - defer conn.Close() | |
| 94 | - | |
| 95 | - fc.logger.Info("firehose connected") | |
| 96 | - | |
| 97 | - for { | |
| 98 | - select { | |
| 99 | - case <-ctx.Done(): | |
| 100 | - return ctx.Err() | |
| 101 | - default: | |
| 102 | - } | |
| 103 | - | |
| 104 | - _, msg, err := conn.ReadMessage() | |
| 105 | - if err != nil { | |
| 106 | - return fmt.Errorf("reading firehose: %w", err) | |
| 107 | - } | |
| 108 | - | |
| 109 | - fc.handleMessage(ctx, msg) | |
| 110 | - } | |
| 111 | -} | |
| 112 | - | |
| 113 | -func (fc *FirehoseConsumer) handleMessage(ctx context.Context, msg []byte) { | |
| 114 | - var frame struct { | |
| 115 | - Type string `json:"type"` | |
| 116 | - Commit json.RawMessage `json:"commit"` | |
| 117 | - Seq int64 `json:"seq"` | |
| 118 | - } | |
| 119 | - | |
| 120 | - if err := json.Unmarshal(msg, &frame); err != nil { | |
| 121 | - return | |
| 122 | - } | |
| 123 | - | |
| 124 | - if frame.Seq > 0 { | |
| 125 | - fc.cursor = frame.Seq | |
| 126 | - } | |
| 127 | - | |
| 128 | - if frame.Type == "#commit" { | |
| 129 | - fc.parseCommit(ctx, frame.Commit) | |
| 130 | - } | |
| 131 | -} | |
| 132 | - | |
| 133 | -func (fc *FirehoseConsumer) parseCommit(ctx context.Context, raw json.RawMessage) { | |
| 134 | - var commit struct { | |
| 135 | - Did string `json:"did"` | |
| 136 | - Ops []struct { | |
| 137 | - Action string `json:"action"` | |
| 138 | - Path string `json:"path"` | |
| 139 | - CID json.RawMessage `json:"cid"` | |
| 140 | - Record json.RawMessage `json:"record"` | |
| 141 | - } `json:"ops"` | |
| 142 | - } | |
| 143 | - | |
| 144 | - if err := json.Unmarshal(raw, &commit); err != nil { | |
| 145 | - return | |
| 146 | - } | |
| 147 | - | |
| 148 | - for _, op := range commit.Ops { | |
| 149 | - parts := splitPath(op.Path) | |
| 150 | - if len(parts) != 2 { | |
| 151 | - continue | |
| 152 | - } | |
| 153 | - | |
| 154 | - collection := parts[0] | |
| 155 | - rkey := parts[1] | |
| 156 | - | |
| 157 | - if !fc.collections[collection] { | |
| 158 | - continue | |
| 159 | - } | |
| 160 | - | |
| 161 | - var action string | |
| 162 | - switch op.Action { | |
| 163 | - case "create": | |
| 164 | - action = "create" | |
| 165 | - case "update": | |
| 166 | - action = "update" | |
| 167 | - case "delete": | |
| 168 | - action = "delete" | |
| 169 | - default: | |
| 170 | - continue | |
| 171 | - } | |
| 172 | - | |
| 173 | - evt := &FirehoseEvent{ | |
| 174 | - Type: action, | |
| 175 | - DID: commit.Did, | |
| 176 | - Collection: collection, | |
| 177 | - RKey: rkey, | |
| 178 | - URI: fmt.Sprintf("at://%s/%s/%s", commit.Did, collection, rkey), | |
| 179 | - Value: op.Record, | |
| 180 | - } | |
| 181 | - | |
| 182 | - if op.CID != nil { | |
| 183 | - evt.CID = string(op.CID) | |
| 184 | - } | |
| 185 | - | |
| 186 | - if err := fc.handler(ctx, evt); err != nil { | |
| 187 | - fc.logger.Error("firehose handler error", "error", err) | |
| 188 | - metrics.FirehoseErrors.Inc() | |
| 189 | - } | |
| 190 | - | |
| 191 | - metrics.FirehoseEvents.WithLabelValues(collection, action).Inc() | |
| 192 | - } | |
| 193 | -} | |
| 194 | - | |
| 195 | -func splitPath(p string) []string { | |
| 196 | - if p == "" { | |
| 197 | - return nil | |
| 198 | - } | |
| 199 | - var parts []string | |
| 200 | - start := 0 | |
| 201 | - for i := 0; i < len(p); i++ { | |
| 202 | - if p[i] == '/' { | |
| 203 | - parts = append(parts, p[start:i]) | |
| 204 | - start = i + 1 | |
| 205 | - } | |
| 206 | - } | |
| 207 | - parts = append(parts, p[start:]) | |
| 208 | - return parts | |
| 209 | -} | |
| deleted file mode 100644 | |||
| @@ -1,209 +0,0 @@ | |||
| 1 | -package atproto | ||
| 2 | - | ||
| 3 | -import ( | ||
| 4 | - "context" | ||
| 5 | - "encoding/json" | ||
| 6 | - "fmt" | ||
| 7 | - "log/slog" | ||
| 8 | - "net/http" | ||
| 9 | - "net/url" | ||
| 10 | - "time" | ||
| 11 | - | ||
| 12 | - "github.com/gorilla/websocket" | ||
| 13 | - | ||
| 14 | - "pkg.rbrt.fr/glean/internal/metrics" | ||
| 15 | -) | ||
| 16 | - | ||
| 17 | -type FirehoseEvent struct { | ||
| 18 | - Type string | ||
| 19 | - DID string | ||
| 20 | - Collection string | ||
| 21 | - RKey string | ||
| 22 | - URI string | ||
| 23 | - CID string | ||
| 24 | - Value json.RawMessage | ||
| 25 | -} | ||
| 26 | - | ||
| 27 | -type FirehoseHandler func(ctx context.Context, event *FirehoseEvent) error | ||
| 28 | - | ||
| 29 | -type FirehoseConsumer struct { | ||
| 30 | - relayURL string | ||
| 31 | - handler FirehoseHandler | ||
| 32 | - cursor int64 | ||
| 33 | - logger *slog.Logger | ||
| 34 | - collections map[string]bool | ||
| 35 | -} | ||
| 36 | - | ||
| 37 | -func NewFirehoseConsumer(relayURL string, handler FirehoseHandler, logger *slog.Logger) *FirehoseConsumer { | ||
| 38 | - return &FirehoseConsumer{ | ||
| 39 | - relayURL: relayURL, | ||
| 40 | - handler: handler, | ||
| 41 | - logger: logger, | ||
| 42 | - collections: map[string]bool{ | ||
| 43 | - "at.glean.subscription": true, | ||
| 44 | - "at.glean.annotation": true, | ||
| 45 | - "at.glean.like": true, | ||
| 46 | - "app.bsky.graph.follow": true, | ||
| 47 | - "sh.tangled.graph.follow": true, | ||
| 48 | - }, | ||
| 49 | - } | ||
| 50 | -} | ||
| 51 | - | ||
| 52 | -func (fc *FirehoseConsumer) Start(ctx context.Context) error { | ||
| 53 | - for { | ||
| 54 | - err := fc.connect(ctx) | ||
| 55 | - if ctx.Err() != nil { | ||
| 56 | - return ctx.Err() | ||
| 57 | - } | ||
| 58 | - if err != nil { | ||
| 59 | - fc.logger.Error("firehose connection error", "error", err) | ||
| 60 | - metrics.FirehoseReconnects.Inc() | ||
| 61 | - } | ||
| 62 | - | ||
| 63 | - select { | ||
| 64 | - case <-ctx.Done(): | ||
| 65 | - return ctx.Err() | ||
| 66 | - case <-time.After(5 * time.Second): | ||
| 67 | - } | ||
| 68 | - } | ||
| 69 | -} | ||
| 70 | - | ||
| 71 | -func (fc *FirehoseConsumer) connect(ctx context.Context) error { | ||
| 72 | - u, err := url.Parse(fc.relayURL) | ||
| 73 | - if err != nil { | ||
| 74 | - return fmt.Errorf("parsing relay URL: %w", err) | ||
| 75 | - } | ||
| 76 | - | ||
| 77 | - scheme := "wss" | ||
| 78 | - if u.Scheme == "http" || u.Scheme == "ws" { | ||
| 79 | - scheme = "ws" | ||
| 80 | - } | ||
| 81 | - | ||
| 82 | - wsURL := fmt.Sprintf("%s://%s/xrpc/com.atproto.sync.subscribeRepos", scheme, u.Host) | ||
| 83 | - if fc.cursor > 0 { | ||
| 84 | - wsURL += fmt.Sprintf("?cursor=%d", fc.cursor) | ||
| 85 | - } | ||
| 86 | - | ||
| 87 | - fc.logger.Info("connecting to firehose", "url", wsURL) | ||
| 88 | - | ||
| 89 | - conn, _, err := websocket.DefaultDialer.DialContext(ctx, wsURL, http.Header{}) | ||
| 90 | - if err != nil { | ||
| 91 | - return fmt.Errorf("dialing firehose: %w", err) | ||
| 92 | - } | ||
| 93 | - defer conn.Close() | ||
| 94 | - | ||
| 95 | - fc.logger.Info("firehose connected") | ||
| 96 | - | ||
| 97 | - for { | ||
| 98 | - select { | ||
| 99 | - case <-ctx.Done(): | ||
| 100 | - return ctx.Err() | ||
| 101 | - default: | ||
| 102 | - } | ||
| 103 | - | ||
| 104 | - _, msg, err := conn.ReadMessage() | ||
| 105 | - if err != nil { | ||
| 106 | - return fmt.Errorf("reading firehose: %w", err) | ||
| 107 | - } | ||
| 108 | - | ||
| 109 | - fc.handleMessage(ctx, msg) | ||
| 110 | - } | ||
| 111 | -} | ||
| 112 | - | ||
| 113 | -func (fc *FirehoseConsumer) handleMessage(ctx context.Context, msg []byte) { | ||
| 114 | - var frame struct { | ||
| 115 | - Type string `json:"type"` | ||
| 116 | - Commit json.RawMessage `json:"commit"` | ||
| 117 | - Seq int64 `json:"seq"` | ||
| 118 | - } | ||
| 119 | - | ||
| 120 | - if err := json.Unmarshal(msg, &frame); err != nil { | ||
| 121 | - return | ||
| 122 | - } | ||
| 123 | - | ||
| 124 | - if frame.Seq > 0 { | ||
| 125 | - fc.cursor = frame.Seq | ||
| 126 | - } | ||
| 127 | - | ||
| 128 | - if frame.Type == "#commit" { | ||
| 129 | - fc.parseCommit(ctx, frame.Commit) | ||
| 130 | - } | ||
| 131 | -} | ||
| 132 | - | ||
| 133 | -func (fc *FirehoseConsumer) parseCommit(ctx context.Context, raw json.RawMessage) { | ||
| 134 | - var commit struct { | ||
| 135 | - Did string `json:"did"` | ||
| 136 | - Ops []struct { | ||
| 137 | - Action string `json:"action"` | ||
| 138 | - Path string `json:"path"` | ||
| 139 | - CID json.RawMessage `json:"cid"` | ||
| 140 | - Record json.RawMessage `json:"record"` | ||
| 141 | - } `json:"ops"` | ||
| 142 | - } | ||
| 143 | - | ||
| 144 | - if err := json.Unmarshal(raw, &commit); err != nil { | ||
| 145 | - return | ||
| 146 | - } | ||
| 147 | - | ||
| 148 | - for _, op := range commit.Ops { | ||
| 149 | - parts := splitPath(op.Path) | ||
| 150 | - if len(parts) != 2 { | ||
| 151 | - continue | ||
| 152 | - } | ||
| 153 | - | ||
| 154 | - collection := parts[0] | ||
| 155 | - rkey := parts[1] | ||
| 156 | - | ||
| 157 | - if !fc.collections[collection] { | ||
| 158 | - continue | ||
| 159 | - } | ||
| 160 | - | ||
| 161 | - var action string | ||
| 162 | - switch op.Action { | ||
| 163 | - case "create": | ||
| 164 | - action = "create" | ||
| 165 | - case "update": | ||
| 166 | - action = "update" | ||
| 167 | - case "delete": | ||
| 168 | - action = "delete" | ||
| 169 | - default: | ||
| 170 | - continue | ||
| 171 | - } | ||
| 172 | - | ||
| 173 | - evt := &FirehoseEvent{ | ||
| 174 | - Type: action, | ||
| 175 | - DID: commit.Did, | ||
| 176 | - Collection: collection, | ||
| 177 | - RKey: rkey, | ||
| 178 | - URI: fmt.Sprintf("at://%s/%s/%s", commit.Did, collection, rkey), | ||
| 179 | - Value: op.Record, | ||
| 180 | - } | ||
| 181 | - | ||
| 182 | - if op.CID != nil { | ||
| 183 | - evt.CID = string(op.CID) | ||
| 184 | - } | ||
| 185 | - | ||
| 186 | - if err := fc.handler(ctx, evt); err != nil { | ||
| 187 | - fc.logger.Error("firehose handler error", "error", err) | ||
| 188 | - metrics.FirehoseErrors.Inc() | ||
| 189 | - } | ||
| 190 | - | ||
| 191 | - metrics.FirehoseEvents.WithLabelValues(collection, action).Inc() | ||
| 192 | - } | ||
| 193 | -} | ||
| 194 | - | ||
| 195 | -func splitPath(p string) []string { | ||
| 196 | - if p == "" { | ||
| 197 | - return nil | ||
| 198 | - } | ||
| 199 | - var parts []string | ||
| 200 | - start := 0 | ||
| 201 | - for i := 0; i < len(p); i++ { | ||
| 202 | - if p[i] == '/' { | ||
| 203 | - parts = append(parts, p[start:i]) | ||
| 204 | - start = i + 1 | ||
| 205 | - } | ||
| 206 | - } | ||
| 207 | - parts = append(parts, p[start:]) | ||
| 208 | - return parts | ||
| 209 | -} | ||
added
internal/atproto/jetstream.go +142 -0 | new file mode 100644 | ||
| @@ -0,0 +1,142 @@ | ||
| 1 | +package atproto | |
| 2 | + | |
| 3 | +import ( | |
| 4 | + "context" | |
| 5 | + "encoding/json" | |
| 6 | + "fmt" | |
| 7 | + "log/slog" | |
| 8 | + "strings" | |
| 9 | + "time" | |
| 10 | + | |
| 11 | + jsc "github.com/bluesky-social/jetstream/pkg/client" | |
| 12 | + "github.com/bluesky-social/jetstream/pkg/models" | |
| 13 | + "go.uber.org/atomic" | |
| 14 | + | |
| 15 | + "pkg.rbrt.fr/glean/internal/metrics" | |
| 16 | +) | |
| 17 | + | |
| 18 | +type Event struct { | |
| 19 | + Type string | |
| 20 | + DID string | |
| 21 | + Collection string | |
| 22 | + RKey string | |
| 23 | + URI string | |
| 24 | + CID string | |
| 25 | + Value json.RawMessage | |
| 26 | +} | |
| 27 | + | |
| 28 | +type EventHandler func(ctx context.Context, event *Event) error | |
| 29 | + | |
| 30 | +type jetstreamScheduler struct { | |
| 31 | + handler EventHandler | |
| 32 | + logger *slog.Logger | |
| 33 | + cursor atomic.Int64 | |
| 34 | +} | |
| 35 | + | |
| 36 | +func (s *jetstreamScheduler) AddWork(ctx context.Context, _ string, evt *models.Event) error { | |
| 37 | + if evt.TimeUS > 0 { | |
| 38 | + s.cursor.Store(evt.TimeUS) | |
| 39 | + } | |
| 40 | + | |
| 41 | + if evt.Kind != models.EventKindCommit || evt.Commit == nil { | |
| 42 | + return nil | |
| 43 | + } | |
| 44 | + | |
| 45 | + c := evt.Commit | |
| 46 | + if c.Operation != models.CommitOperationCreate && | |
| 47 | + c.Operation != models.CommitOperationUpdate && | |
| 48 | + c.Operation != models.CommitOperationDelete { | |
| 49 | + return nil | |
| 50 | + } | |
| 51 | + | |
| 52 | + e := &Event{ | |
| 53 | + Type: c.Operation, | |
| 54 | + DID: evt.Did, | |
| 55 | + Collection: c.Collection, | |
| 56 | + RKey: c.RKey, | |
| 57 | + URI: fmt.Sprintf("at://%s/%s/%s", evt.Did, c.Collection, c.RKey), | |
| 58 | + CID: c.CID, | |
| 59 | + Value: json.RawMessage(c.Record), | |
| 60 | + } | |
| 61 | + | |
| 62 | + if err := s.handler(ctx, e); err != nil { | |
| 63 | + s.logger.Error("jetstream handler error", "error", err) | |
| 64 | + metrics.JetstreamErrors.Inc() | |
| 65 | + } | |
| 66 | + | |
| 67 | + metrics.JetstreamEvents.WithLabelValues(c.Collection, c.Operation).Inc() | |
| 68 | + return nil | |
| 69 | +} | |
| 70 | + | |
| 71 | +func (s *jetstreamScheduler) Shutdown() {} | |
| 72 | + | |
| 73 | +type JetstreamConsumer struct { | |
| 74 | + client *jsc.Client | |
| 75 | + logger *slog.Logger | |
| 76 | + sched *jetstreamScheduler | |
| 77 | +} | |
| 78 | + | |
| 79 | +func NewJetstreamConsumer(jetstreamURL string, handler EventHandler, logger *slog.Logger) *JetstreamConsumer { | |
| 80 | + sched := &jetstreamScheduler{ | |
| 81 | + handler: handler, | |
| 82 | + logger: logger, | |
| 83 | + } | |
| 84 | + | |
| 85 | + wsURL := jetstreamURL | |
| 86 | + if !strings.HasSuffix(wsURL, "/subscribe") { | |
| 87 | + wsURL += "/subscribe" | |
| 88 | + } | |
| 89 | + | |
| 90 | + config := &jsc.ClientConfig{ | |
| 91 | + Compress: true, | |
| 92 | + WebsocketURL: wsURL, | |
| 93 | + ExtraHeaders: map[string]string{ | |
| 94 | + "User-Agent": "glean/1.0", | |
| 95 | + }, | |
| 96 | + WantedCollections: []string{ | |
| 97 | + "at.glean.subscription", | |
| 98 | + "at.glean.annotation", | |
| 99 | + "at.glean.like", | |
| 100 | + "app.bsky.graph.follow", | |
| 101 | + "sh.tangled.graph.follow", | |
| 102 | + }, | |
| 103 | + } | |
| 104 | + | |
| 105 | + c, err := jsc.NewClient(config, logger, sched) | |
| 106 | + if err != nil { | |
| 107 | + logger.Error("failed to create jetstream client", "error", err) | |
| 108 | + return nil | |
| 109 | + } | |
| 110 | + | |
| 111 | + return &JetstreamConsumer{ | |
| 112 | + client: c, | |
| 113 | + logger: logger, | |
| 114 | + sched: sched, | |
| 115 | + } | |
| 116 | +} | |
| 117 | + | |
| 118 | +func (jc *JetstreamConsumer) Start(ctx context.Context) error { | |
| 119 | + for { | |
| 120 | + cursor := jc.sched.cursor.Load() | |
| 121 | + var cursorPtr *int64 | |
| 122 | + if cursor > 0 { | |
| 123 | + adjusted := cursor - int64(5*time.Second/time.Microsecond) | |
| 124 | + cursorPtr = &adjusted | |
| 125 | + } | |
| 126 | + | |
| 127 | + err := jc.client.ConnectAndRead(ctx, cursorPtr) | |
| 128 | + if ctx.Err() != nil { | |
| 129 | + return ctx.Err() | |
| 130 | + } | |
| 131 | + if err != nil { | |
| 132 | + jc.logger.Error("jetstream connection error", "error", err) | |
| 133 | + metrics.JetstreamReconnects.Inc() | |
| 134 | + } | |
| 135 | + | |
| 136 | + select { | |
| 137 | + case <-ctx.Done(): | |
| 138 | + return ctx.Err() | |
| 139 | + case <-time.After(5 * time.Second): | |
| 140 | + } | |
| 141 | + } | |
| 142 | +} | |
| new file mode 100644 | |||
| @@ -0,0 +1,142 @@ | |||
| 1 | +package atproto | ||
| 2 | + | ||
| 3 | +import ( | ||
| 4 | + "context" | ||
| 5 | + "encoding/json" | ||
| 6 | + "fmt" | ||
| 7 | + "log/slog" | ||
| 8 | + "strings" | ||
| 9 | + "time" | ||
| 10 | + | ||
| 11 | + jsc "github.com/bluesky-social/jetstream/pkg/client" | ||
| 12 | + "github.com/bluesky-social/jetstream/pkg/models" | ||
| 13 | + "go.uber.org/atomic" | ||
| 14 | + | ||
| 15 | + "pkg.rbrt.fr/glean/internal/metrics" | ||
| 16 | +) | ||
| 17 | + | ||
| 18 | +type Event struct { | ||
| 19 | + Type string | ||
| 20 | + DID string | ||
| 21 | + Collection string | ||
| 22 | + RKey string | ||
| 23 | + URI string | ||
| 24 | + CID string | ||
| 25 | + Value json.RawMessage | ||
| 26 | +} | ||
| 27 | + | ||
| 28 | +type EventHandler func(ctx context.Context, event *Event) error | ||
| 29 | + | ||
| 30 | +type jetstreamScheduler struct { | ||
| 31 | + handler EventHandler | ||
| 32 | + logger *slog.Logger | ||
| 33 | + cursor atomic.Int64 | ||
| 34 | +} | ||
| 35 | + | ||
| 36 | +func (s *jetstreamScheduler) AddWork(ctx context.Context, _ string, evt *models.Event) error { | ||
| 37 | + if evt.TimeUS > 0 { | ||
| 38 | + s.cursor.Store(evt.TimeUS) | ||
| 39 | + } | ||
| 40 | + | ||
| 41 | + if evt.Kind != models.EventKindCommit || evt.Commit == nil { | ||
| 42 | + return nil | ||
| 43 | + } | ||
| 44 | + | ||
| 45 | + c := evt.Commit | ||
| 46 | + if c.Operation != models.CommitOperationCreate && | ||
| 47 | + c.Operation != models.CommitOperationUpdate && | ||
| 48 | + c.Operation != models.CommitOperationDelete { | ||
| 49 | + return nil | ||
| 50 | + } | ||
| 51 | + | ||
| 52 | + e := &Event{ | ||
| 53 | + Type: c.Operation, | ||
| 54 | + DID: evt.Did, | ||
| 55 | + Collection: c.Collection, | ||
| 56 | + RKey: c.RKey, | ||
| 57 | + URI: fmt.Sprintf("at://%s/%s/%s", evt.Did, c.Collection, c.RKey), | ||
| 58 | + CID: c.CID, | ||
| 59 | + Value: json.RawMessage(c.Record), | ||
| 60 | + } | ||
| 61 | + | ||
| 62 | + if err := s.handler(ctx, e); err != nil { | ||
| 63 | + s.logger.Error("jetstream handler error", "error", err) | ||
| 64 | + metrics.JetstreamErrors.Inc() | ||
| 65 | + } | ||
| 66 | + | ||
| 67 | + metrics.JetstreamEvents.WithLabelValues(c.Collection, c.Operation).Inc() | ||
| 68 | + return nil | ||
| 69 | +} | ||
| 70 | + | ||
| 71 | +func (s *jetstreamScheduler) Shutdown() {} | ||
| 72 | + | ||
| 73 | +type JetstreamConsumer struct { | ||
| 74 | + client *jsc.Client | ||
| 75 | + logger *slog.Logger | ||
| 76 | + sched *jetstreamScheduler | ||
| 77 | +} | ||
| 78 | + | ||
| 79 | +func NewJetstreamConsumer(jetstreamURL string, handler EventHandler, logger *slog.Logger) *JetstreamConsumer { | ||
| 80 | + sched := &jetstreamScheduler{ | ||
| 81 | + handler: handler, | ||
| 82 | + logger: logger, | ||
| 83 | + } | ||
| 84 | + | ||
| 85 | + wsURL := jetstreamURL | ||
| 86 | + if !strings.HasSuffix(wsURL, "/subscribe") { | ||
| 87 | + wsURL += "/subscribe" | ||
| 88 | + } | ||
| 89 | + | ||
| 90 | + config := &jsc.ClientConfig{ | ||
| 91 | + Compress: true, | ||
| 92 | + WebsocketURL: wsURL, | ||
| 93 | + ExtraHeaders: map[string]string{ | ||
| 94 | + "User-Agent": "glean/1.0", | ||
| 95 | + }, | ||
| 96 | + WantedCollections: []string{ | ||
| 97 | + "at.glean.subscription", | ||
| 98 | + "at.glean.annotation", | ||
| 99 | + "at.glean.like", | ||
| 100 | + "app.bsky.graph.follow", | ||
| 101 | + "sh.tangled.graph.follow", | ||
| 102 | + }, | ||
| 103 | + } | ||
| 104 | + | ||
| 105 | + c, err := jsc.NewClient(config, logger, sched) | ||
| 106 | + if err != nil { | ||
| 107 | + logger.Error("failed to create jetstream client", "error", err) | ||
| 108 | + return nil | ||
| 109 | + } | ||
| 110 | + | ||
| 111 | + return &JetstreamConsumer{ | ||
| 112 | + client: c, | ||
| 113 | + logger: logger, | ||
| 114 | + sched: sched, | ||
| 115 | + } | ||
| 116 | +} | ||
| 117 | + | ||
| 118 | +func (jc *JetstreamConsumer) Start(ctx context.Context) error { | ||
| 119 | + for { | ||
| 120 | + cursor := jc.sched.cursor.Load() | ||
| 121 | + var cursorPtr *int64 | ||
| 122 | + if cursor > 0 { | ||
| 123 | + adjusted := cursor - int64(5*time.Second/time.Microsecond) | ||
| 124 | + cursorPtr = &adjusted | ||
| 125 | + } | ||
| 126 | + | ||
| 127 | + err := jc.client.ConnectAndRead(ctx, cursorPtr) | ||
| 128 | + if ctx.Err() != nil { | ||
| 129 | + return ctx.Err() | ||
| 130 | + } | ||
| 131 | + if err != nil { | ||
| 132 | + jc.logger.Error("jetstream connection error", "error", err) | ||
| 133 | + metrics.JetstreamReconnects.Inc() | ||
| 134 | + } | ||
| 135 | + | ||
| 136 | + select { | ||
| 137 | + case <-ctx.Done(): | ||
| 138 | + return ctx.Err() | ||
| 139 | + case <-time.After(5 * time.Second): | ||
| 140 | + } | ||
| 141 | + } | ||
| 142 | +} | ||
modified
internal/atproto/lexicon.go +2 -0 | @@ -1,3 +1,5 @@ | ||
| 1 | +// Lexicon record types are maintained by hand (no lexgen). | |
| 2 | +// See lexicon_test.go for the test ensuring these stay in sync with lexicons/. | |
| 1 | 3 | package atproto |
| 2 | 4 | |
| 3 | 5 | import ( |
| @@ -1,3 +1,5 @@ | |||
| 1 | +// Lexicon record types are maintained by hand (no lexgen). | ||
| 2 | +// See lexicon_test.go for the test ensuring these stay in sync with lexicons/. | ||
| 1 | package atproto | 3 | package atproto |
| 2 | 4 | ||
| 3 | import ( | 5 | import ( |
added
internal/atproto/lexicon_test.go +74 -0 | new file mode 100644 | ||
| @@ -0,0 +1,74 @@ | ||
| 1 | +package atproto | |
| 2 | + | |
| 3 | +import ( | |
| 4 | + "encoding/json" | |
| 5 | + "os" | |
| 6 | + "path/filepath" | |
| 7 | + "reflect" | |
| 8 | + "strings" | |
| 9 | + "testing" | |
| 10 | + | |
| 11 | + "gotest.tools/v3/assert" | |
| 12 | +) | |
| 13 | + | |
| 14 | +func lexiconPath(filename string) string { | |
| 15 | + return filepath.Join("..", "..", "lexicons", "at", "glean", filename) | |
| 16 | +} | |
| 17 | + | |
| 18 | +func readLexiconProperties(t *testing.T, filename string) map[string]any { | |
| 19 | + t.Helper() | |
| 20 | + data, err := os.ReadFile(lexiconPath(filename)) | |
| 21 | + assert.NilError(t, err) | |
| 22 | + | |
| 23 | + var schema struct { | |
| 24 | + Defs struct { | |
| 25 | + Main struct { | |
| 26 | + Record struct { | |
| 27 | + Properties map[string]any `json:"properties"` | |
| 28 | + } `json:"record"` | |
| 29 | + } `json:"main"` | |
| 30 | + } `json:"defs"` | |
| 31 | + } | |
| 32 | + assert.NilError(t, json.Unmarshal(data, &schema)) | |
| 33 | + assert.Assert(t, len(schema.Defs.Main.Record.Properties) > 0, "no properties found in %s", filename) | |
| 34 | + return schema.Defs.Main.Record.Properties | |
| 35 | +} | |
| 36 | + | |
| 37 | +func assertStructMatchesLexicon[T any](t *testing.T, filename string) { | |
| 38 | + t.Helper() | |
| 39 | + properties := readLexiconProperties(t, filename) | |
| 40 | + | |
| 41 | + var zero T | |
| 42 | + typ := reflect.TypeOf(zero) | |
| 43 | + | |
| 44 | + jsonTags := make(map[string]bool) | |
| 45 | + for field := range typ.Fields() { | |
| 46 | + tag := field.Tag.Get("json") | |
| 47 | + name := strings.Split(tag, ",")[0] | |
| 48 | + if name == "" || name == "-" { | |
| 49 | + continue | |
| 50 | + } | |
| 51 | + jsonTags[name] = true | |
| 52 | + } | |
| 53 | + | |
| 54 | + for prop := range properties { | |
| 55 | + assert.Assert(t, jsonTags[prop], "lexicon property %q missing from %s (lexicon file: %s)", prop, typ.Name(), filename) | |
| 56 | + } | |
| 57 | + | |
| 58 | + for name := range jsonTags { | |
| 59 | + _, exists := properties[name] | |
| 60 | + assert.Assert(t, exists, "Go field %q in %s missing from lexicon %s", name, typ.Name(), filename) | |
| 61 | + } | |
| 62 | +} | |
| 63 | + | |
| 64 | +func TestSubscriptionRecordMatchesLexicon(t *testing.T) { | |
| 65 | + assertStructMatchesLexicon[SubscriptionRecord](t, "subscription.json") | |
| 66 | +} | |
| 67 | + | |
| 68 | +func TestAnnotationRecordMatchesLexicon(t *testing.T) { | |
| 69 | + assertStructMatchesLexicon[AnnotationRecord](t, "annotation.json") | |
| 70 | +} | |
| 71 | + | |
| 72 | +func TestLikeRecordMatchesLexicon(t *testing.T) { | |
| 73 | + assertStructMatchesLexicon[LikeRecord](t, "like.json") | |
| 74 | +} | |
| new file mode 100644 | |||
| @@ -0,0 +1,74 @@ | |||
| 1 | +package atproto | ||
| 2 | + | ||
| 3 | +import ( | ||
| 4 | + "encoding/json" | ||
| 5 | + "os" | ||
| 6 | + "path/filepath" | ||
| 7 | + "reflect" | ||
| 8 | + "strings" | ||
| 9 | + "testing" | ||
| 10 | + | ||
| 11 | + "gotest.tools/v3/assert" | ||
| 12 | +) | ||
| 13 | + | ||
| 14 | +func lexiconPath(filename string) string { | ||
| 15 | + return filepath.Join("..", "..", "lexicons", "at", "glean", filename) | ||
| 16 | +} | ||
| 17 | + | ||
| 18 | +func readLexiconProperties(t *testing.T, filename string) map[string]any { | ||
| 19 | + t.Helper() | ||
| 20 | + data, err := os.ReadFile(lexiconPath(filename)) | ||
| 21 | + assert.NilError(t, err) | ||
| 22 | + | ||
| 23 | + var schema struct { | ||
| 24 | + Defs struct { | ||
| 25 | + Main struct { | ||
| 26 | + Record struct { | ||
| 27 | + Properties map[string]any `json:"properties"` | ||
| 28 | + } `json:"record"` | ||
| 29 | + } `json:"main"` | ||
| 30 | + } `json:"defs"` | ||
| 31 | + } | ||
| 32 | + assert.NilError(t, json.Unmarshal(data, &schema)) | ||
| 33 | + assert.Assert(t, len(schema.Defs.Main.Record.Properties) > 0, "no properties found in %s", filename) | ||
| 34 | + return schema.Defs.Main.Record.Properties | ||
| 35 | +} | ||
| 36 | + | ||
| 37 | +func assertStructMatchesLexicon[T any](t *testing.T, filename string) { | ||
| 38 | + t.Helper() | ||
| 39 | + properties := readLexiconProperties(t, filename) | ||
| 40 | + | ||
| 41 | + var zero T | ||
| 42 | + typ := reflect.TypeOf(zero) | ||
| 43 | + | ||
| 44 | + jsonTags := make(map[string]bool) | ||
| 45 | + for field := range typ.Fields() { | ||
| 46 | + tag := field.Tag.Get("json") | ||
| 47 | + name := strings.Split(tag, ",")[0] | ||
| 48 | + if name == "" || name == "-" { | ||
| 49 | + continue | ||
| 50 | + } | ||
| 51 | + jsonTags[name] = true | ||
| 52 | + } | ||
| 53 | + | ||
| 54 | + for prop := range properties { | ||
| 55 | + assert.Assert(t, jsonTags[prop], "lexicon property %q missing from %s (lexicon file: %s)", prop, typ.Name(), filename) | ||
| 56 | + } | ||
| 57 | + | ||
| 58 | + for name := range jsonTags { | ||
| 59 | + _, exists := properties[name] | ||
| 60 | + assert.Assert(t, exists, "Go field %q in %s missing from lexicon %s", name, typ.Name(), filename) | ||
| 61 | + } | ||
| 62 | +} | ||
| 63 | + | ||
| 64 | +func TestSubscriptionRecordMatchesLexicon(t *testing.T) { | ||
| 65 | + assertStructMatchesLexicon[SubscriptionRecord](t, "subscription.json") | ||
| 66 | +} | ||
| 67 | + | ||
| 68 | +func TestAnnotationRecordMatchesLexicon(t *testing.T) { | ||
| 69 | + assertStructMatchesLexicon[AnnotationRecord](t, "annotation.json") | ||
| 70 | +} | ||
| 71 | + | ||
| 72 | +func TestLikeRecordMatchesLexicon(t *testing.T) { | ||
| 73 | + assertStructMatchesLexicon[LikeRecord](t, "like.json") | ||
| 74 | +} | ||
renamed
internal/atproto/stream_handler.go +8 -8 | similarity index 84% | ||
| rename from internal/atproto/firehose_handler.go | ||
| rename to internal/atproto/stream_handler.go | ||
| @@ -10,16 +10,16 @@ import ( | ||
| 10 | 10 | "pkg.rbrt.fr/glean/internal/db" |
| 11 | 11 | ) |
| 12 | 12 | |
| 13 | -type FirehoseDBHandler struct { | |
| 13 | +type StreamDBHandler struct { | |
| 14 | 14 | db *db.DB |
| 15 | 15 | logger *slog.Logger |
| 16 | 16 | } |
| 17 | 17 | |
| 18 | -func NewFirehoseDBHandler(database *db.DB, logger *slog.Logger) *FirehoseDBHandler { | |
| 19 | - return &FirehoseDBHandler{db: database, logger: logger} | |
| 18 | +func NewStreamDBHandler(database *db.DB, logger *slog.Logger) *StreamDBHandler { | |
| 19 | + return &StreamDBHandler{db: database, logger: logger} | |
| 20 | 20 | } |
| 21 | 21 | |
| 22 | -func (h *FirehoseDBHandler) Handle(ctx context.Context, event *FirehoseEvent) error { | |
| 22 | +func (h *StreamDBHandler) Handle(ctx context.Context, event *Event) error { | |
| 23 | 23 | switch event.Collection { |
| 24 | 24 | case "at.glean.subscription": |
| 25 | 25 | return h.handleSubscription(ctx, event) |
| @@ -33,7 +33,7 @@ func (h *FirehoseDBHandler) Handle(ctx context.Context, event *FirehoseEvent) er | ||
| 33 | 33 | return nil |
| 34 | 34 | } |
| 35 | 35 | |
| 36 | -func (h *FirehoseDBHandler) handleSubscription(ctx context.Context, event *FirehoseEvent) error { | |
| 36 | +func (h *StreamDBHandler) handleSubscription(ctx context.Context, event *Event) error { | |
| 37 | 37 | switch event.Type { |
| 38 | 38 | case "create", "update": |
| 39 | 39 | var rec SubscriptionRecord |
| @@ -72,7 +72,7 @@ func (h *FirehoseDBHandler) handleSubscription(ctx context.Context, event *Fireh | ||
| 72 | 72 | return nil |
| 73 | 73 | } |
| 74 | 74 | |
| 75 | -func (h *FirehoseDBHandler) handleLike(ctx context.Context, event *FirehoseEvent) error { | |
| 75 | +func (h *StreamDBHandler) handleLike(ctx context.Context, event *Event) error { | |
| 76 | 76 | switch event.Type { |
| 77 | 77 | case "create": |
| 78 | 78 | var rec LikeRecord |
| @@ -104,7 +104,7 @@ func (h *FirehoseDBHandler) handleLike(ctx context.Context, event *FirehoseEvent | ||
| 104 | 104 | return nil |
| 105 | 105 | } |
| 106 | 106 | |
| 107 | -func (h *FirehoseDBHandler) handleAnnotation(ctx context.Context, event *FirehoseEvent) error { | |
| 107 | +func (h *StreamDBHandler) handleAnnotation(ctx context.Context, event *Event) error { | |
| 108 | 108 | switch event.Type { |
| 109 | 109 | case "create": |
| 110 | 110 | var rec AnnotationRecord |
| @@ -138,7 +138,7 @@ func (h *FirehoseDBHandler) handleAnnotation(ctx context.Context, event *Firehos | ||
| 138 | 138 | return nil |
| 139 | 139 | } |
| 140 | 140 | |
| 141 | -func (h *FirehoseDBHandler) handleFollow(ctx context.Context, event *FirehoseEvent) error { | |
| 141 | +func (h *StreamDBHandler) handleFollow(ctx context.Context, event *Event) error { | |
| 142 | 142 | switch event.Type { |
| 143 | 143 | case "create": |
| 144 | 144 | var rec FollowRecord |
| similarity index 84% | |||
| rename from internal/atproto/firehose_handler.go | |||
| rename to internal/atproto/stream_handler.go | |||
| @@ -10,16 +10,16 @@ import ( | |||
| 10 | "pkg.rbrt.fr/glean/internal/db" | 10 | "pkg.rbrt.fr/glean/internal/db" |
| 11 | ) | 11 | ) |
| 12 | 12 | ||
| 13 | -type FirehoseDBHandler struct { | 13 | +type StreamDBHandler struct { |
| 14 | db *db.DB | 14 | db *db.DB |
| 15 | logger *slog.Logger | 15 | logger *slog.Logger |
| 16 | } | 16 | } |
| 17 | 17 | ||
| 18 | -func NewFirehoseDBHandler(database *db.DB, logger *slog.Logger) *FirehoseDBHandler { | 18 | +func NewStreamDBHandler(database *db.DB, logger *slog.Logger) *StreamDBHandler { |
| 19 | - return &FirehoseDBHandler{db: database, logger: logger} | 19 | + return &StreamDBHandler{db: database, logger: logger} |
| 20 | } | 20 | } |
| 21 | 21 | ||
| 22 | -func (h *FirehoseDBHandler) Handle(ctx context.Context, event *FirehoseEvent) error { | 22 | +func (h *StreamDBHandler) Handle(ctx context.Context, event *Event) error { |
| 23 | switch event.Collection { | 23 | switch event.Collection { |
| 24 | case "at.glean.subscription": | 24 | case "at.glean.subscription": |
| 25 | return h.handleSubscription(ctx, event) | 25 | return h.handleSubscription(ctx, event) |
| @@ -33,7 +33,7 @@ func (h *FirehoseDBHandler) Handle(ctx context.Context, event *FirehoseEvent) er | |||
| 33 | return nil | 33 | return nil |
| 34 | } | 34 | } |
| 35 | 35 | ||
| 36 | -func (h *FirehoseDBHandler) handleSubscription(ctx context.Context, event *FirehoseEvent) error { | 36 | +func (h *StreamDBHandler) handleSubscription(ctx context.Context, event *Event) error { |
| 37 | switch event.Type { | 37 | switch event.Type { |
| 38 | case "create", "update": | 38 | case "create", "update": |
| 39 | var rec SubscriptionRecord | 39 | var rec SubscriptionRecord |
| @@ -72,7 +72,7 @@ func (h *FirehoseDBHandler) handleSubscription(ctx context.Context, event *Fireh | |||
| 72 | return nil | 72 | return nil |
| 73 | } | 73 | } |
| 74 | 74 | ||
| 75 | -func (h *FirehoseDBHandler) handleLike(ctx context.Context, event *FirehoseEvent) error { | 75 | +func (h *StreamDBHandler) handleLike(ctx context.Context, event *Event) error { |
| 76 | switch event.Type { | 76 | switch event.Type { |
| 77 | case "create": | 77 | case "create": |
| 78 | var rec LikeRecord | 78 | var rec LikeRecord |
| @@ -104,7 +104,7 @@ func (h *FirehoseDBHandler) handleLike(ctx context.Context, event *FirehoseEvent | |||
| 104 | return nil | 104 | return nil |
| 105 | } | 105 | } |
| 106 | 106 | ||
| 107 | -func (h *FirehoseDBHandler) handleAnnotation(ctx context.Context, event *FirehoseEvent) error { | 107 | +func (h *StreamDBHandler) handleAnnotation(ctx context.Context, event *Event) error { |
| 108 | switch event.Type { | 108 | switch event.Type { |
| 109 | case "create": | 109 | case "create": |
| 110 | var rec AnnotationRecord | 110 | var rec AnnotationRecord |
| @@ -138,7 +138,7 @@ func (h *FirehoseDBHandler) handleAnnotation(ctx context.Context, event *Firehos | |||
| 138 | return nil | 138 | return nil |
| 139 | } | 139 | } |
| 140 | 140 | ||
| 141 | -func (h *FirehoseDBHandler) handleFollow(ctx context.Context, event *FirehoseEvent) error { | 141 | +func (h *StreamDBHandler) handleFollow(ctx context.Context, event *Event) error { |
| 142 | switch event.Type { | 142 | switch event.Type { |
| 143 | case "create": | 143 | case "create": |
| 144 | var rec FollowRecord | 144 | var rec FollowRecord |
modified
internal/atproto/sync.go +14 -3 | @@ -1,3 +1,14 @@ | ||
| 1 | +// Sync implements per-user reconciliation using com.atproto.repo.listRecords. | |
| 2 | +// This is a lighter alternative to full ATProto backfilling (which uses | |
| 3 | +// com.atproto.sync.getRepo with revision tracking and event buffering). | |
| 4 | +// The full approach is unnecessary here because: | |
| 5 | +// - we only sync known users (not the entire network) | |
| 6 | +// - the Jetstream consumer handles real-time events concurrently | |
| 7 | +// - all reconcile operations are idempotent | |
| 8 | +// | |
| 9 | +// Known trade-off: syncFollows atomically replaces all follows for a user. | |
| 10 | +// A Jetstream follow event arriving mid-sync could be lost, but self-heals | |
| 11 | +// on the next sync cycle or Jetstream event. | |
| 1 | 12 | package atproto |
| 2 | 13 | |
| 3 | 14 | import ( |
| @@ -72,6 +83,9 @@ func (s *Sync) reconcileSubscription(ctx context.Context, userDID, uri, cid stri | ||
| 72 | 83 | return nil |
| 73 | 84 | } |
| 74 | 85 | |
| 86 | + f := &db.Feed{FeedURL: rec.FeedURL, Title: db.NullStr(rec.Title)} | |
| 87 | + _ = s.db.UpsertFeed(ctx, f) | |
| 88 | + | |
| 75 | 89 | existing, err := s.db.GetSubscription(ctx, userDID, rec.FeedURL) |
| 76 | 90 | if err == nil && existing != nil { |
| 77 | 91 | if !existing.URI.Valid || existing.URI.String == "" { |
| @@ -80,9 +94,6 @@ func (s *Sync) reconcileSubscription(ctx context.Context, userDID, uri, cid stri | ||
| 80 | 94 | return nil |
| 81 | 95 | } |
| 82 | 96 | |
| 83 | - f := &db.Feed{FeedURL: rec.FeedURL, Title: db.NullStr(rec.Title)} | |
| 84 | - _ = s.db.UpsertFeed(ctx, f) | |
| 85 | - | |
| 86 | 97 | return s.db.CreateSubscription(ctx, userDID, rec.FeedURL, rec.Title, rec.Category, uri, cid) |
| 87 | 98 | } |
| 88 | 99 | |
| @@ -1,3 +1,14 @@ | |||
| 1 | +// Sync implements per-user reconciliation using com.atproto.repo.listRecords. | ||
| 2 | +// This is a lighter alternative to full ATProto backfilling (which uses | ||
| 3 | +// com.atproto.sync.getRepo with revision tracking and event buffering). | ||
| 4 | +// The full approach is unnecessary here because: | ||
| 5 | +// - we only sync known users (not the entire network) | ||
| 6 | +// - the Jetstream consumer handles real-time events concurrently | ||
| 7 | +// - all reconcile operations are idempotent | ||
| 8 | +// | ||
| 9 | +// Known trade-off: syncFollows atomically replaces all follows for a user. | ||
| 10 | +// A Jetstream follow event arriving mid-sync could be lost, but self-heals | ||
| 11 | +// on the next sync cycle or Jetstream event. | ||
| 1 | package atproto | 12 | package atproto |
| 2 | 13 | ||
| 3 | import ( | 14 | import ( |
| @@ -72,6 +83,9 @@ func (s *Sync) reconcileSubscription(ctx context.Context, userDID, uri, cid stri | |||
| 72 | return nil | 83 | return nil |
| 73 | } | 84 | } |
| 74 | 85 | ||
| 86 | + f := &db.Feed{FeedURL: rec.FeedURL, Title: db.NullStr(rec.Title)} | ||
| 87 | + _ = s.db.UpsertFeed(ctx, f) | ||
| 88 | + | ||
| 75 | existing, err := s.db.GetSubscription(ctx, userDID, rec.FeedURL) | 89 | existing, err := s.db.GetSubscription(ctx, userDID, rec.FeedURL) |
| 76 | if err == nil && existing != nil { | 90 | if err == nil && existing != nil { |
| 77 | if !existing.URI.Valid || existing.URI.String == "" { | 91 | if !existing.URI.Valid || existing.URI.String == "" { |
| @@ -80,9 +94,6 @@ func (s *Sync) reconcileSubscription(ctx context.Context, userDID, uri, cid stri | |||
| 80 | return nil | 94 | return nil |
| 81 | } | 95 | } |
| 82 | 96 | ||
| 83 | - f := &db.Feed{FeedURL: rec.FeedURL, Title: db.NullStr(rec.Title)} | ||
| 84 | - _ = s.db.UpsertFeed(ctx, f) | ||
| 85 | - | ||
| 86 | return s.db.CreateSubscription(ctx, userDID, rec.FeedURL, rec.Title, rec.Category, uri, cid) | 97 | return s.db.CreateSubscription(ctx, userDID, rec.FeedURL, rec.Title, rec.Category, uri, cid) |
| 87 | } | 98 | } |
| 88 | 99 | ||
modified
internal/cluster/jaccard.go +1 -0 | @@ -128,6 +128,7 @@ func (e *Engine) ComputeUserSimilarity(ctx context.Context) error { | ||
| 128 | 128 | 0.5, |
| 129 | 129 | 0 |
| 130 | 130 | FROM follows f |
| 131 | + GROUP BY MIN(f.user_did, f.target_did), MAX(f.user_did, f.target_did) | |
| 131 | 132 | ON CONFLICT(user_a, user_b) DO UPDATE SET |
| 132 | 133 | jaccard = jaccard + 0.5 |
| 133 | 134 | `) |
| @@ -128,6 +128,7 @@ func (e *Engine) ComputeUserSimilarity(ctx context.Context) error { | |||
| 128 | 0.5, | 128 | 0.5, |
| 129 | 0 | 129 | 0 |
| 130 | FROM follows f | 130 | FROM follows f |
| 131 | + GROUP BY MIN(f.user_did, f.target_did), MAX(f.user_did, f.target_did) | ||
| 131 | ON CONFLICT(user_a, user_b) DO UPDATE SET | 132 | ON CONFLICT(user_a, user_b) DO UPDATE SET |
| 132 | jaccard = jaccard + 0.5 | 133 | jaccard = jaccard + 0.5 |
| 133 | `) | 134 | `) |
modified
internal/metrics/metrics.go +9 -9 | @@ -22,19 +22,19 @@ var ( | ||
| 22 | 22 | Help: "Total number of articles upserted", |
| 23 | 23 | }) |
| 24 | 24 | |
| 25 | - FirehoseEvents = promauto.NewCounterVec(prometheus.CounterOpts{ | |
| 26 | - Name: "glean_firehose_events_total", | |
| 27 | - Help: "Total number of firehose events processed", | |
| 25 | + JetstreamEvents = promauto.NewCounterVec(prometheus.CounterOpts{ | |
| 26 | + Name: "glean_jetstream_events_total", | |
| 27 | + Help: "Total number of jetstream events processed", | |
| 28 | 28 | }, []string{"collection", "action"}) |
| 29 | 29 | |
| 30 | - FirehoseErrors = promauto.NewCounter(prometheus.CounterOpts{ | |
| 31 | - Name: "glean_firehose_errors_total", | |
| 32 | - Help: "Total number of firehose handler errors", | |
| 30 | + JetstreamErrors = promauto.NewCounter(prometheus.CounterOpts{ | |
| 31 | + Name: "glean_jetstream_errors_total", | |
| 32 | + Help: "Total number of jetstream handler errors", | |
| 33 | 33 | }) |
| 34 | 34 | |
| 35 | - FirehoseReconnects = promauto.NewCounter(prometheus.CounterOpts{ | |
| 36 | - Name: "glean_firehose_reconnects_total", | |
| 37 | - Help: "Number of firehose reconnections", | |
| 35 | + JetstreamReconnects = promauto.NewCounter(prometheus.CounterOpts{ | |
| 36 | + Name: "glean_jetstream_reconnects_total", | |
| 37 | + Help: "Number of jetstream reconnections", | |
| 38 | 38 | }) |
| 39 | 39 | |
| 40 | 40 | HTTPRequests = promauto.NewCounterVec(prometheus.CounterOpts{ |
| @@ -22,19 +22,19 @@ var ( | |||
| 22 | Help: "Total number of articles upserted", | 22 | Help: "Total number of articles upserted", |
| 23 | }) | 23 | }) |
| 24 | 24 | ||
| 25 | - FirehoseEvents = promauto.NewCounterVec(prometheus.CounterOpts{ | 25 | + JetstreamEvents = promauto.NewCounterVec(prometheus.CounterOpts{ |
| 26 | - Name: "glean_firehose_events_total", | 26 | + Name: "glean_jetstream_events_total", |
| 27 | - Help: "Total number of firehose events processed", | 27 | + Help: "Total number of jetstream events processed", |
| 28 | }, []string{"collection", "action"}) | 28 | }, []string{"collection", "action"}) |
| 29 | 29 | ||
| 30 | - FirehoseErrors = promauto.NewCounter(prometheus.CounterOpts{ | 30 | + JetstreamErrors = promauto.NewCounter(prometheus.CounterOpts{ |
| 31 | - Name: "glean_firehose_errors_total", | 31 | + Name: "glean_jetstream_errors_total", |
| 32 | - Help: "Total number of firehose handler errors", | 32 | + Help: "Total number of jetstream handler errors", |
| 33 | }) | 33 | }) |
| 34 | 34 | ||
| 35 | - FirehoseReconnects = promauto.NewCounter(prometheus.CounterOpts{ | 35 | + JetstreamReconnects = promauto.NewCounter(prometheus.CounterOpts{ |
| 36 | - Name: "glean_firehose_reconnects_total", | 36 | + Name: "glean_jetstream_reconnects_total", |
| 37 | - Help: "Number of firehose reconnections", | 37 | + Help: "Number of jetstream reconnections", |
| 38 | }) | 38 | }) |
| 39 | 39 | ||
| 40 | HTTPRequests = promauto.NewCounterVec(prometheus.CounterOpts{ | 40 | HTTPRequests = promauto.NewCounterVec(prometheus.CounterOpts{ |
modified
internal/server/server.go +8 -12 | @@ -8,6 +8,7 @@ import ( | ||
| 8 | 8 | "log/slog" |
| 9 | 9 | "net/http" |
| 10 | 10 | "net/url" |
| 11 | + "slices" | |
| 11 | 12 | "strconv" |
| 12 | 13 | "strings" |
| 13 | 14 | "time" |
| @@ -26,8 +27,8 @@ import ( | ||
| 26 | 27 | "pkg.rbrt.fr/glean/internal/metrics" |
| 27 | 28 | "pkg.rbrt.fr/glean/internal/sanitize" |
| 28 | 29 | "pkg.rbrt.fr/glean/internal/scraper" |
| 29 | - "pkg.rbrt.fr/glean/static" | |
| 30 | 30 | "pkg.rbrt.fr/glean/internal/tmpl" |
| 31 | + "pkg.rbrt.fr/glean/static" | |
| 31 | 32 | ) |
| 32 | 33 | |
| 33 | 34 | var oauthScopes = []string{"atproto", "transition:generic"} |
| @@ -256,14 +257,14 @@ func (s *Server) loadTemplates() { | ||
| 256 | 257 | return id |
| 257 | 258 | } |
| 258 | 259 | } |
| 259 | - if strings.HasPrefix(u.Path, "/embed/") { | |
| 260 | - id := strings.TrimPrefix(u.Path, "/embed/") | |
| 260 | + if after, ok := strings.CutPrefix(u.Path, "/embed/"); ok { | |
| 261 | + id := after | |
| 261 | 262 | if id != "" { |
| 262 | 263 | return id |
| 263 | 264 | } |
| 264 | 265 | } |
| 265 | - if strings.HasPrefix(u.Path, "/shorts/") { | |
| 266 | - id := strings.TrimPrefix(u.Path, "/shorts/") | |
| 266 | + if after, ok := strings.CutPrefix(u.Path, "/shorts/"); ok { | |
| 267 | + id := after | |
| 267 | 268 | if id != "" { |
| 268 | 269 | return id |
| 269 | 270 | } |
| @@ -277,18 +278,13 @@ func (s *Server) loadTemplates() { | ||
| 277 | 278 | return false |
| 278 | 279 | } |
| 279 | 280 | host := strings.ToLower(u.Hostname()) |
| 280 | - for _, h := range []string{ | |
| 281 | + return slices.Contains([]string{ | |
| 281 | 282 | "www.youtube.com", "youtube.com", "m.youtube.com", "youtu.be", |
| 282 | 283 | "vimeo.com", "player.vimeo.com", |
| 283 | 284 | "open.spotify.com", "embed.spotify.com", |
| 284 | 285 | "w.soundcloud.com", |
| 285 | 286 | "bandcamp.com", |
| 286 | - } { | |
| 287 | - if host == h { | |
| 288 | - return true | |
| 289 | - } | |
| 290 | - } | |
| 291 | - return false | |
| 287 | + }, host) | |
| 292 | 288 | }, |
| 293 | 289 | "sanitizeHTML": func(input string) template.HTML { |
| 294 | 290 | return template.HTML(sanitize.HTML(input)) |
| @@ -8,6 +8,7 @@ import ( | |||
| 8 | "log/slog" | 8 | "log/slog" |
| 9 | "net/http" | 9 | "net/http" |
| 10 | "net/url" | 10 | "net/url" |
| 11 | + "slices" | ||
| 11 | "strconv" | 12 | "strconv" |
| 12 | "strings" | 13 | "strings" |
| 13 | "time" | 14 | "time" |
| @@ -26,8 +27,8 @@ import ( | |||
| 26 | "pkg.rbrt.fr/glean/internal/metrics" | 27 | "pkg.rbrt.fr/glean/internal/metrics" |
| 27 | "pkg.rbrt.fr/glean/internal/sanitize" | 28 | "pkg.rbrt.fr/glean/internal/sanitize" |
| 28 | "pkg.rbrt.fr/glean/internal/scraper" | 29 | "pkg.rbrt.fr/glean/internal/scraper" |
| 29 | - "pkg.rbrt.fr/glean/static" | ||
| 30 | "pkg.rbrt.fr/glean/internal/tmpl" | 30 | "pkg.rbrt.fr/glean/internal/tmpl" |
| 31 | + "pkg.rbrt.fr/glean/static" | ||
| 31 | ) | 32 | ) |
| 32 | 33 | ||
| 33 | var oauthScopes = []string{"atproto", "transition:generic"} | 34 | var oauthScopes = []string{"atproto", "transition:generic"} |
| @@ -256,14 +257,14 @@ func (s *Server) loadTemplates() { | |||
| 256 | return id | 257 | return id |
| 257 | } | 258 | } |
| 258 | } | 259 | } |
| 259 | - if strings.HasPrefix(u.Path, "/embed/") { | 260 | + if after, ok := strings.CutPrefix(u.Path, "/embed/"); ok { |
| 260 | - id := strings.TrimPrefix(u.Path, "/embed/") | 261 | + id := after |
| 261 | if id != "" { | 262 | if id != "" { |
| 262 | return id | 263 | return id |
| 263 | } | 264 | } |
| 264 | } | 265 | } |
| 265 | - if strings.HasPrefix(u.Path, "/shorts/") { | 266 | + if after, ok := strings.CutPrefix(u.Path, "/shorts/"); ok { |
| 266 | - id := strings.TrimPrefix(u.Path, "/shorts/") | 267 | + id := after |
| 267 | if id != "" { | 268 | if id != "" { |
| 268 | return id | 269 | return id |
| 269 | } | 270 | } |
| @@ -277,18 +278,13 @@ func (s *Server) loadTemplates() { | |||
| 277 | return false | 278 | return false |
| 278 | } | 279 | } |
| 279 | host := strings.ToLower(u.Hostname()) | 280 | host := strings.ToLower(u.Hostname()) |
| 280 | - for _, h := range []string{ | 281 | + return slices.Contains([]string{ |
| 281 | "www.youtube.com", "youtube.com", "m.youtube.com", "youtu.be", | 282 | "www.youtube.com", "youtube.com", "m.youtube.com", "youtu.be", |
| 282 | "vimeo.com", "player.vimeo.com", | 283 | "vimeo.com", "player.vimeo.com", |
| 283 | "open.spotify.com", "embed.spotify.com", | 284 | "open.spotify.com", "embed.spotify.com", |
| 284 | "w.soundcloud.com", | 285 | "w.soundcloud.com", |
| 285 | "bandcamp.com", | 286 | "bandcamp.com", |
| 286 | - } { | 287 | + }, host) |
| 287 | - if host == h { | ||
| 288 | - return true | ||
| 289 | - } | ||
| 290 | - } | ||
| 291 | - return false | ||
| 292 | }, | 288 | }, |
| 293 | "sanitizeHTML": func(input string) template.HTML { | 289 | "sanitizeHTML": func(input string) template.HTML { |
| 294 | return template.HTML(sanitize.HTML(input)) | 290 | return template.HTML(sanitize.HTML(input)) |
modified
main.go +5 -5 | @@ -21,7 +21,7 @@ import ( | ||
| 21 | 21 | func main() { |
| 22 | 22 | addr := flag.String("addr", envOr("GLEAN_ADDR", ":8080"), "listen address") |
| 23 | 23 | dbPath := flag.String("db", envOr("GLEAN_DB", "glean.db"), "database path") |
| 24 | - relayURL := flag.String("relay", envOr("GLEAN_RELAY", "wss://bsky.network"), "AT Relay URL") | |
| 24 | + jetstreamURL := flag.String("jetstream", envOr("GLEAN_JETSTREAM", "wss://jetstream2.fr.hose.cam"), "Jetstream URL") | |
| 25 | 25 | flag.Parse() |
| 26 | 26 | |
| 27 | 27 | logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo})) |
| @@ -44,8 +44,8 @@ func main() { | ||
| 44 | 44 | engine := cluster.NewEngine(database.DB, logger) |
| 45 | 45 | cron := cluster.NewCron(engine, 6*time.Hour, logger) |
| 46 | 46 | |
| 47 | - firehoseHandler := atproto.NewFirehoseDBHandler(database, logger) | |
| 48 | - firehose := atproto.NewFirehoseConsumer(*relayURL, firehoseHandler.Handle, logger) | |
| 47 | + handler := atproto.NewStreamDBHandler(database, logger) | |
| 48 | + jetstream := atproto.NewJetstreamConsumer(*jetstreamURL, handler.Handle, logger) | |
| 49 | 49 | |
| 50 | 50 | ctx, cancel := context.WithCancel(context.Background()) |
| 51 | 51 | defer cancel() |
| @@ -64,8 +64,8 @@ func main() { | ||
| 64 | 64 | srv.PeriodicSync(ctx, 1*time.Hour) |
| 65 | 65 | }() |
| 66 | 66 | go func() { |
| 67 | - if err := firehose.Start(ctx); err != nil && ctx.Err() == nil { | |
| 68 | - logger.Error("firehose error", "error", err) | |
| 67 | + if err := jetstream.Start(ctx); err != nil && ctx.Err() == nil { | |
| 68 | + logger.Error("jetstream error", "error", err) | |
| 69 | 69 | } |
| 70 | 70 | }() |
| 71 | 71 | |
| @@ -21,7 +21,7 @@ import ( | |||
| 21 | func main() { | 21 | func main() { |
| 22 | addr := flag.String("addr", envOr("GLEAN_ADDR", ":8080"), "listen address") | 22 | addr := flag.String("addr", envOr("GLEAN_ADDR", ":8080"), "listen address") |
| 23 | dbPath := flag.String("db", envOr("GLEAN_DB", "glean.db"), "database path") | 23 | dbPath := flag.String("db", envOr("GLEAN_DB", "glean.db"), "database path") |
| 24 | - relayURL := flag.String("relay", envOr("GLEAN_RELAY", "wss://bsky.network"), "AT Relay URL") | 24 | + jetstreamURL := flag.String("jetstream", envOr("GLEAN_JETSTREAM", "wss://jetstream2.fr.hose.cam"), "Jetstream URL") |
| 25 | flag.Parse() | 25 | flag.Parse() |
| 26 | 26 | ||
| 27 | logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo})) | 27 | logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo})) |
| @@ -44,8 +44,8 @@ func main() { | |||
| 44 | engine := cluster.NewEngine(database.DB, logger) | 44 | engine := cluster.NewEngine(database.DB, logger) |
| 45 | cron := cluster.NewCron(engine, 6*time.Hour, logger) | 45 | cron := cluster.NewCron(engine, 6*time.Hour, logger) |
| 46 | 46 | ||
| 47 | - firehoseHandler := atproto.NewFirehoseDBHandler(database, logger) | 47 | + handler := atproto.NewStreamDBHandler(database, logger) |
| 48 | - firehose := atproto.NewFirehoseConsumer(*relayURL, firehoseHandler.Handle, logger) | 48 | + jetstream := atproto.NewJetstreamConsumer(*jetstreamURL, handler.Handle, logger) |
| 49 | 49 | ||
| 50 | ctx, cancel := context.WithCancel(context.Background()) | 50 | ctx, cancel := context.WithCancel(context.Background()) |
| 51 | defer cancel() | 51 | defer cancel() |
| @@ -64,8 +64,8 @@ func main() { | |||
| 64 | srv.PeriodicSync(ctx, 1*time.Hour) | 64 | srv.PeriodicSync(ctx, 1*time.Hour) |
| 65 | }() | 65 | }() |
| 66 | go func() { | 66 | go func() { |
| 67 | - if err := firehose.Start(ctx); err != nil && ctx.Err() == nil { | 67 | + if err := jetstream.Start(ctx); err != nil && ctx.Err() == nil { |
| 68 | - logger.Error("firehose error", "error", err) | 68 | + logger.Error("jetstream error", "error", err) |
| 69 | } | 69 | } |
| 70 | }() | 70 | }() |
| 71 | 71 | ||
modified
readme.md +7 -7 | @@ -36,13 +36,13 @@ Then open `http://localhost:8080`. | ||
| 36 | 36 | |
| 37 | 37 | ## Configuration |
| 38 | 38 | |
| 39 | -| Variable | Default | What it does | | |
| 40 | -| -------------------------- | -------------------- | --------------------------------------------------------- | | |
| 41 | -| `GLEAN_ADDR` | `:8080` | Listen address | | |
| 42 | -| `GLEAN_DB` | `glean.db` | SQLite database path | | |
| 43 | -| `GLEAN_RELAY` | `wss://bsky.network` | AT Relay WebSocket URL | | |
| 44 | -| `GLEAN_OAUTH_CLIENT_ID` | _(empty)_ | OAuth client metadata URL (leave empty for localhost dev) | | |
| 45 | -| `GLEAN_OAUTH_REDIRECT_URL` | _(empty)_ | OAuth redirect URL (leave empty for localhost dev) | | |
| 39 | +| Variable | Default | What it does | | |
| 40 | +| -------------------------- | ------------------------------ | --------------------------------------------------------- | | |
| 41 | +| `GLEAN_ADDR` | `:8080` | Listen address | | |
| 42 | +| `GLEAN_DB` | `glean.db` | SQLite database path | | |
| 43 | +| `GLEAN_JETSTREAM` | `wss://jetstream2.fr.hose.cam` | Jetstream WebSocket URL | | |
| 44 | +| `GLEAN_OAUTH_CLIENT_ID` | _(empty)_ | OAuth client metadata URL (leave empty for localhost dev) | | |
| 45 | +| `GLEAN_OAUTH_REDIRECT_URL` | _(empty)_ | OAuth redirect URL (leave empty for localhost dev) | | |
| 46 | 46 | |
| 47 | 47 | For production: |
| 48 | 48 | |
| @@ -36,13 +36,13 @@ Then open `http://localhost:8080`. | |||
| 36 | 36 | ||
| 37 | ## Configuration | 37 | ## Configuration |
| 38 | 38 | ||
| 39 | -| Variable | Default | What it does | | 39 | +| Variable | Default | What it does | |
| 40 | -| -------------------------- | -------------------- | --------------------------------------------------------- | | 40 | +| -------------------------- | ------------------------------ | --------------------------------------------------------- | |
| 41 | -| `GLEAN_ADDR` | `:8080` | Listen address | | 41 | +| `GLEAN_ADDR` | `:8080` | Listen address | |
| 42 | -| `GLEAN_DB` | `glean.db` | SQLite database path | | 42 | +| `GLEAN_DB` | `glean.db` | SQLite database path | |
| 43 | -| `GLEAN_RELAY` | `wss://bsky.network` | AT Relay WebSocket URL | | 43 | +| `GLEAN_JETSTREAM` | `wss://jetstream2.fr.hose.cam` | Jetstream WebSocket URL | |
| 44 | -| `GLEAN_OAUTH_CLIENT_ID` | _(empty)_ | OAuth client metadata URL (leave empty for localhost dev) | | 44 | +| `GLEAN_OAUTH_CLIENT_ID` | _(empty)_ | OAuth client metadata URL (leave empty for localhost dev) | |
| 45 | -| `GLEAN_OAUTH_REDIRECT_URL` | _(empty)_ | OAuth redirect URL (leave empty for localhost dev) | | 45 | +| `GLEAN_OAUTH_REDIRECT_URL` | _(empty)_ | OAuth redirect URL (leave empty for localhost dev) | |
| 46 | 46 | ||
| 47 | For production: | 47 | For production: |
| 48 | 48 | ||