From aa77df5a1667f57523dd178c3fd3722c1a3921ee Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Thu, 29 Jan 2026 12:33:26 +0000 Subject: [PATCH] knotmirror: introduce knotmirror KnotMirror is an external service that is intended to be used by appview. It will ingest all known git repos and provide xrpc methods to inspect them, so appview won't need to fetch individual knots on every page render. Using postgres exclusively instead of sqlite to support A LOT of concurrent writes. Signed-off-by: Seongmin Lee --- flake.nix | 4 +++- go.mod | 19 ++++++++++++------- go.sum | 44 ++++++++++++++++++++++++++++---------------- knotmirror/adminpage.go | 182 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/crawler.go | 25 +++++++++++++++++++++++++ knotmirror/git.go | 305 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/knotmirror.go | 117 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/metrics.go | 29 +++++++++++++++++++++++++++++ knotmirror/readme.md | 8 ++++++++ knotmirror/resyncer.go | 273 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/tapclient.go | 152 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ nix/gomod2nix.toml | 43 +++++++++++++++++++++++++++++-------------- cmd/knotmirror/main.go | 58 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/config/config.go | 34 ++++++++++++++++++++++++++++++++++ knotmirror/db/db.go | 100 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/db/hosts.go | 102 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/db/repos.go | 275 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/knotstream/knotstream.go | 88 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/knotstream/metrics.go | 28 ++++++++++++++++++++++++++++ knotmirror/knotstream/scheduler.go | 102 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/knotstream/slurper.go | 334 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/knotstream/subscription.go | 22 ++++++++++++++++++++++ knotmirror/models/models.go | 110 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/templates/base.html | 55 +++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/templates/hosts.html | 44 ++++++++++++++++++++++++++++++++++++++++++++ knotmirror/templates/repos.html | 86 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ nix/pkgs/knot-mirror.nix | 18 ++++++++++++++++++ 27 file(s) changed, 2619 insertion(s)(+), 38 deletion(s)(-) diff --git a/flake.nix b/flake.nix --- a/flake.nix +++ b/flake.nix @@ -107,10 +107,11 @@ knot = self.callPackage ./nix/pkgs/knot.nix {}; dolly = self.callPackage ./nix/pkgs/dolly.nix {}; tap = self.callPackage ./nix/pkgs/tap.nix {}; + knotmirror = self.callPackage ./nix/pkgs/knot-mirror.nix {}; }); in { overlays.default = final: prev: { - inherit (mkPackageSet final) lexgen goat sqlite-lib spindle knot-unwrapped knot appview docs dolly tap; + inherit (mkPackageSet final) lexgen goat sqlite-lib spindle knot-unwrapped knot appview docs dolly tap knotmirror; }; packages = forAllSystems (system: let @@ -206,6 +207,7 @@ pkgs.coreutils # for those of us who are on systems that use busybox (alpine) packages'.lexgen packages'.treefmt-wrapper + packages'.tap ]; shellHook = '' mkdir -p appview/pages/static diff --git a/go.mod b/go.mod --- a/go.mod +++ b/go.mod @@ -35,16 +35,18 @@ github.com/hiddeco/sshsig v0.2.0 github.com/hpcloud/tail v1.0.0 github.com/ipfs/go-cid v0.5.0 + github.com/jackc/pgx/v5 v5.8.0 github.com/mattn/go-sqlite3 v1.14.24 github.com/microcosm-cc/bluemonday v1.0.27 github.com/openbao/openbao/api/v2 v2.3.0 github.com/posthog/posthog-go v1.5.5 + github.com/prometheus/client_golang v1.23.2 github.com/redis/go-redis/v9 v9.7.3 github.com/resend/resend-go/v2 v2.15.0 github.com/sethvargo/go-envconfig v1.1.0 github.com/srwiley/oksvg v0.0.0-20221011165216-be6e8873101c github.com/srwiley/rasterx v0.0.0-20220730225603-2ab79fcdd4ef - github.com/stretchr/testify v1.10.0 + github.com/stretchr/testify v1.11.1 github.com/urfave/cli/v3 v3.4.1 github.com/whyrusleeping/cbor-gen v0.3.1 github.com/yuin/goldmark v1.7.13 @@ -52,9 +54,9 @@ github.com/yuin/goldmark-highlighting/v2 v2.0.0-20230729083705-37449abec8cc gitlab.com/staticnoise/goldmark-callout v0.0.0-20240609120641-6366b799e4ab go.abhg.dev/goldmark/mermaid v0.6.0 - golang.org/x/crypto v0.40.0 + golang.org/x/crypto v0.41.0 golang.org/x/image v0.31.0 - golang.org/x/net v0.42.0 + golang.org/x/net v0.43.0 golang.org/x/xerrors v0.0.0-20240903120638-7835f813f4da gopkg.in/yaml.v3 v3.0.1 ) @@ -161,6 +163,9 @@ github.com/ipfs/go-log v1.0.5 // indirect github.com/ipfs/go-log/v2 v2.6.0 // indirect github.com/ipfs/go-metrics-interface v0.3.0 // indirect + github.com/jackc/pgpassfile v1.0.0 // indirect + github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect + github.com/jackc/puddle/v2 v2.2.2 // indirect github.com/json-iterator/go v1.1.12 // indirect github.com/kevinburke/ssh_config v1.2.0 // indirect github.com/klauspost/compress v1.18.0 // indirect @@ -193,9 +198,8 @@ github.com/pkg/errors v0.9.1 // indirect github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect github.com/polydawn/refmt v0.89.1-0.20221221234430-40501e09de1f // indirect - github.com/prometheus/client_golang v1.22.0 // indirect github.com/prometheus/client_model v0.6.2 // indirect - github.com/prometheus/common v0.64.0 // indirect + github.com/prometheus/common v0.66.1 // indirect github.com/prometheus/procfs v0.16.1 // indirect github.com/rivo/uniseg v0.4.7 // indirect github.com/ryanuber/go-glob v1.0.0 // indirect @@ -222,15 +226,16 @@ go.uber.org/atomic v1.11.0 // indirect go.uber.org/multierr v1.11.0 // indirect go.uber.org/zap v1.27.0 // indirect + go.yaml.in/yaml/v2 v2.4.2 // indirect golang.org/x/exp v0.0.0-20250620022241-b7579e27df2b // indirect golang.org/x/sync v0.17.0 // indirect - golang.org/x/sys v0.34.0 // indirect + golang.org/x/sys v0.35.0 // indirect golang.org/x/text v0.29.0 // indirect golang.org/x/time v0.12.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20250603155806-513f23925822 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20250603155806-513f23925822 // indirect google.golang.org/grpc v1.73.0 // indirect - google.golang.org/protobuf v1.36.6 // indirect + google.golang.org/protobuf v1.36.8 // indirect gopkg.in/fsnotify.v1 v1.4.7 // indirect gopkg.in/tomb.v1 v1.0.0-20141024135613-dd632973f1e7 // indirect gopkg.in/warnings.v0 v0.1.2 // indirect diff --git a/go.sum b/go.sum --- a/go.sum +++ b/go.sum @@ -350,6 +350,14 @@ github.com/ipfs/go-log/v2 v2.6.0/go.mod h1:p+Efr3qaY5YXpx9TX7MoLCSEZX5boSWj9wh86P5HJa8= github.com/ipfs/go-metrics-interface v0.3.0 h1:YwG7/Cy4R94mYDUuwsBfeziJCVm9pBMJ6q/JR9V40TU= github.com/ipfs/go-metrics-interface v0.3.0/go.mod h1:OxxQjZDGocXVdyTPocns6cOLwHieqej/jos7H4POwoY= +github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= +github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= +github.com/jackc/pgx/v5 v5.8.0 h1:TYPDoleBBme0xGSAX3/+NujXXtpZn9HBONkQC7IEZSo= +github.com/jackc/pgx/v5 v5.8.0/go.mod h1:QVeDInX2m9VyzvNeiCJVjCkNFqzsNb43204HshNSZKw= +github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= +github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= github.com/jtolds/gls v4.20.0+incompatible h1:xdiiI2gbIgH/gLH7ADydsJ1uDOEzR8yvV7C0MuV77Wo= @@ -371,6 +379,8 @@ github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= +github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= github.com/lucasb-eyer/go-colorful v1.2.0 h1:1nnpGOrhyZZuNyfu1QjKiUICQ74+3FNCN69Aj6K7nkY= github.com/lucasb-eyer/go-colorful v1.2.0/go.mod h1:R4dSotOR9KMtayYi1e77YzuveK+i7ruzyGqttikkLy0= github.com/mattn/go-colorable v0.1.14 h1:9A9LHSqF/7dyVVX6g0U9cwm9pG3kP9gSzcuIPHPsaIE= @@ -472,12 +482,12 @@ github.com/polydawn/refmt v0.89.1-0.20221221234430-40501e09de1f/go.mod h1:/zvteZs/GwLtCgZ4BL6CBsk9IKIlexP43ObX9AxTqTw= github.com/posthog/posthog-go v1.5.5 h1:2o3j7IrHbTIfxRtj4MPaXKeimuTYg49onNzNBZbwksM= github.com/posthog/posthog-go v1.5.5/go.mod h1:3RqUmSnPuwmeVj/GYrS75wNGqcAKdpODiwc83xZWgdE= -github.com/prometheus/client_golang v1.22.0 h1:rb93p9lokFEsctTys46VnV1kLCDpVZ0a/Y92Vm0Zc6Q= -github.com/prometheus/client_golang v1.22.0/go.mod h1:R7ljNsLXhuQXYZYtw6GAE9AZg8Y7vEW5scdCXrWRXC0= +github.com/prometheus/client_golang v1.23.2 h1:Je96obch5RDVy3FDMndoUsjAhG5Edi49h0RJWRi/o0o= +github.com/prometheus/client_golang v1.23.2/go.mod h1:Tb1a6LWHB3/SPIzCoaDXI4I8UHKeFTEQ1YCr+0Gyqmg= github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk= github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE= -github.com/prometheus/common v0.64.0 h1:pdZeA+g617P7oGv1CzdTzyeShxAGrTBsolKNOLQPGO4= -github.com/prometheus/common v0.64.0/go.mod h1:0gZns+BLRQ3V6NdaerOhMbwwRbNh9hkGINtQAsP5GS8= +github.com/prometheus/common v0.66.1 h1:h5E0h5/Y8niHc5DlaLlWLArTQI7tMrsfQjHV+d9ZoGs= +github.com/prometheus/common v0.66.1/go.mod h1:gcaUsgf3KfRSwHY4dIMXLPV0K/Wg1oZ8+SbZk/HH/dA= github.com/prometheus/procfs v0.16.1 h1:hZ15bTNuirocR6u0JZ6BAHHmwS1p8B4P6MRqxtzMyRg= github.com/prometheus/procfs v0.16.1/go.mod h1:teAbpZRB1iIAJYREa1LsoWUXykVXA1KlTmWl8x/U+Is= github.com/redis/go-redis/v9 v9.0.0-rc.4/go.mod h1:Vo3EsyWnicKnSKCA7HhgnvnyA74wOA69Cd2Meli5mmA= @@ -523,8 +533,8 @@ github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= -github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= -github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/tidwall/gjson v1.14.2/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk= github.com/tidwall/gjson v1.18.0 h1:FIDeeyB800efLX89e5a8Y0BNH+LOngJyGrIWxG2FKQY= github.com/tidwall/gjson v1.18.0/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk= @@ -608,6 +618,8 @@ go.uber.org/zap v1.16.0/go.mod h1:MA8QOfq0BHJwdXa996Y4dYkAqRKB8/1K1QMMZVaNZjQ= go.uber.org/zap v1.27.0 h1:aJMhYGrd5QSmlpLMr2MftRKl7t8J8PTZPA732ud/XR8= go.uber.org/zap v1.27.0/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E= +go.yaml.in/yaml/v2 v2.4.2 h1:DzmwEr2rDGHl7lsFgAHxmNz/1NlQ7xLIrlN2h5d1eGI= +go.yaml.in/yaml/v2 v2.4.2/go.mod h1:081UH+NErpNdqlCXm3TtEran0rJZGxAYx9hb/ELlsPU= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20190510104115-cbcb75029529/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= @@ -615,8 +627,8 @@ golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc= golang.org/x/crypto v0.1.0/go.mod h1:RecgLatLF4+eUMCP1PoPZQb+cVrJcOPbHkTkbkB9sbw= golang.org/x/crypto v0.19.0/go.mod h1:Iy9bg/ha4yyC70EfRS8jz+B6ybOBKMaSxLj6P6oBDfU= -golang.org/x/crypto v0.40.0 h1:r4x+VvoG5Fm+eJcxMaY8CQM7Lb0l1lsmjGBQ6s8BfKM= -golang.org/x/crypto v0.40.0/go.mod h1:Qr1vMER5WyS2dfPHAlsOj01wgLbsyWtFn/aY+5+ZdxY= +golang.org/x/crypto v0.41.0 h1:WKYxWedPGCTVVl5+WHSSrOBT0O8lx32+zxmHxijgXp4= +golang.org/x/crypto v0.41.0/go.mod h1:pO5AFd7FA68rFak7rOAGVuygIISepHftHnr8dr6+sUc= golang.org/x/exp v0.0.0-20250620022241-b7579e27df2b h1:M2rDM6z3Fhozi9O7NWsxAkg/yqS/lQJ6PmkyIV3YP+o= golang.org/x/exp v0.0.0-20250620022241-b7579e27df2b/go.mod h1:3//PLf8L/X+8b4vuAfHzxeRUl04Adcb341+IGKfnqS8= golang.org/x/image v0.31.0 h1:mLChjE2MV6g1S7oqbXC0/UcKijjm5fnJLUYKIYrLESA= @@ -651,8 +663,8 @@ golang.org/x/net v0.5.0/go.mod h1:DivGGAXEgPSlEBzxGzZI+ZLohi+xUj054jfeKui00ws= golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs= golang.org/x/net v0.10.0/go.mod h1:0qNGK6F8kojg2nk9dLZ2mShWaEBan6FAoqfSigmmuDg= -golang.org/x/net v0.42.0 h1:jzkYrhi3YQWD6MLBJcsklgQsoAcw89EcZbJw8Z614hs= -golang.org/x/net v0.42.0/go.mod h1:FF1RA5d3u7nAYA4z2TkclSCKh68eSXtiFwcWQpPXdt8= +golang.org/x/net v0.43.0 h1:lat02VYK2j4aLzMzecihNvTlJNQUq316m2Mr9rnM6YE= +golang.org/x/net v0.43.0/go.mod h1:vhO1fvI4dGsIjh73sWfUVjj3N7CA9WkKJNQm2svM6Jg= golang.org/x/sync v0.0.0-20180314180146-1d60e4601c6f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20190911185100-cd5d95a43a6e/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= @@ -692,8 +704,8 @@ golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.17.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= -golang.org/x/sys v0.34.0 h1:H5Y5sJ2L2JRdyv7ROF1he/lPdvFsd0mJHFw2ThKHxLA= -golang.org/x/sys v0.34.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= +golang.org/x/sys v0.35.0 h1:vz1N37gP5bs89s7He8XuIYXpyY0+QlsKmzipCbUtyxI= +golang.org/x/sys v0.35.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= golang.org/x/term v0.1.0/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= @@ -703,8 +715,8 @@ golang.org/x/term v0.5.0/go.mod h1:jMB1sMXY+tzblOD4FWmEbocvup2/aLOaQEp7JmGp78k= golang.org/x/term v0.8.0/go.mod h1:xPskH00ivmX89bAKVGSKKtLOWNx2+17Eiy94tnKShWo= golang.org/x/term v0.17.0/go.mod h1:lLRBjIVuehSbZlaOtGMbcMncT+aqLLLmKrsjNrUguwk= -golang.org/x/term v0.33.0 h1:NuFncQrRcaRvVmgRkvM3j/F00gWIAlcmlB8ACEKmGIg= -golang.org/x/term v0.33.0/go.mod h1:s18+ql9tYWp1IfpV9DmCtQDDSRBUjKaw9M1eAv5UeF0= +golang.org/x/term v0.34.0 h1:O/2T7POpk0ZZ7MAzMeWFSg6S5IpWd/RXDlM9hgM3DR4= +golang.org/x/term v0.34.0/go.mod h1:5jC53AEywhIVebHgPVeg0mj8OD3VO9OzclacVrqpaAw= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= @@ -757,8 +769,8 @@ google.golang.org/protobuf v1.26.0-rc.1/go.mod h1:jlhhOSvTdKEhbULTjvd4ARK9grFBp09yW+WbY/TyQbw= google.golang.org/protobuf v1.26.0/go.mod h1:9q0QmTI4eRPtz6boOQmLYwt+qCgq0jsYwAQnmE0givc= google.golang.org/protobuf v1.28.0/go.mod h1:HV8QOd/L58Z+nl8r43ehVNZIU/HEI6OcFqwMG9pJV4I= -google.golang.org/protobuf v1.36.6 h1:z1NpPI8ku2WgiWnf+t9wTPsn6eP1L7ksHUlkfLvd9xY= -google.golang.org/protobuf v1.36.6/go.mod h1:jduwjTPXsFjZGTmRluh+L6NjiWu7pchiJ2/5YcXBHnY= +google.golang.org/protobuf v1.36.8 h1:xHScyCOEuuwZEc6UtSOvPbAT4zRh0xcNRYekJwfqyMc= +google.golang.org/protobuf v1.36.8/go.mod h1:fuxRtAxBytpl4zzqUh6/eyUujkJdNiuEkXntxiD/uRU= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= diff --git a/knotmirror/adminpage.go b/knotmirror/adminpage.go new file mode 100644 --- /dev/null +++ b/knotmirror/adminpage.go @@ -0,0 +1,182 @@ +package knotmirror + +import ( + "database/sql" + "embed" + "fmt" + "html" + "html/template" + "log/slog" + "net/http" + "strconv" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/go-chi/chi/v5" + "tangled.org/core/appview/pagination" + "tangled.org/core/knotmirror/db" + "tangled.org/core/knotmirror/models" +) + +//go:embed templates/*.html +var templateFS embed.FS + +const repoPageSize = 20 + +type AdminServer struct { + db *sql.DB + resyncer *Resyncer + logger *slog.Logger +} + +func NewAdminServer(l *slog.Logger, database *sql.DB, resyncer *Resyncer) *AdminServer { + return &AdminServer{ + db: database, + resyncer: resyncer, + logger: l, + } +} + +func (s *AdminServer) Router() http.Handler { + r := chi.NewRouter() + r.Get("/repos", s.handleRepos()) + r.Get("/hosts", s.handleHosts()) + + r.Post("/api/triggerRepoResync", s.handleRepoResyncTrigger()) + r.Post("/api/cancelRepoResync", s.handleRepoResyncCancel()) + return r +} + +func funcmap() template.FuncMap { + return template.FuncMap{ + "add": func(a, b int) int { return a + b }, + "sub": func(a, b int) int { return a - b }, + "readt": func(ts int64) string { + if ts <= 0 { + return "n/a" + } + return time.Unix(ts, 0).Format("2006-01-02 15:04") + }, + "const": func() map[string]any { + return map[string]any{ + "AllRepoStates": models.AllRepoStates, + "AllHostStatuses": models.AllHostStatuses, + } + }, + } +} + +func (s *AdminServer) handleRepos() http.HandlerFunc { + tpl := template.Must(template.New("").Funcs(funcmap()).ParseFS(templateFS, "templates/base.html", "templates/repos.html")) + return func(w http.ResponseWriter, r *http.Request) { + pageNum, _ := strconv.Atoi(r.URL.Query().Get("page")) + if pageNum < 1 { + pageNum = 1 + } + page := pagination.Page{ + Offset: (pageNum - 1) * repoPageSize, + Limit: repoPageSize, + } + + var ( + did = r.URL.Query().Get("did") + knot = r.URL.Query().Get("knot") + state = r.URL.Query().Get("state") + ) + + repos, err := db.ListRepos(r.Context(), s.db, page, did, knot, state) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + } + counts, err := db.GetRepoCountsByState(r.Context(), s.db) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + } + err = tpl.ExecuteTemplate(w, "base", map[string]any{ + "Repos": repos, + "RepoCounts": counts, + "Page": pageNum, + "FilterByDid": did, + "FilterByKnot": knot, + "FilterByState": models.RepoState(state), + }) + if err != nil { + slog.Error("failed to render", "err", err) + } + } +} + +func (s *AdminServer) handleHosts() http.HandlerFunc { + tpl := template.Must(template.New("").Funcs(funcmap()).ParseFS(templateFS, "templates/base.html", "templates/hosts.html")) + return func(w http.ResponseWriter, r *http.Request) { + var status = models.HostStatus(r.URL.Query().Get("status")) + if status == "" { + status = models.HostStatusActive + } + + hosts, err := db.ListHosts(r.Context(), s.db, status) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + } + err = tpl.ExecuteTemplate(w, "base", map[string]any{ + "Hosts": hosts, + "FilterByStatus": models.HostStatus(status), + }) + if err != nil { + slog.Error("failed to render", "err", err) + } + } +} + +func (s *AdminServer) handleRepoResyncTrigger() http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + var repoQuery = r.FormValue("repo") + + repo, err := syntax.ParseATURI(repoQuery) + if err != nil || repo.RecordKey() == "" { + writeNotif(w, http.StatusBadRequest, fmt.Sprintf("repo parameter invalid: %s", repoQuery)) + return + } + + if err := s.resyncer.TriggerResyncJob(r.Context(), repo); err != nil { + s.logger.Error("failed to trigger resync job", "err", err) + writeNotif(w, http.StatusInternalServerError, fmt.Sprintf("repo parameter invalid: %s", repoQuery)) + return + } + writeNotif(w, http.StatusOK, "success") + } +} + +func (s *AdminServer) handleRepoResyncCancel() http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + var repoQuery = r.FormValue("repo") + + repo, err := syntax.ParseATURI(repoQuery) + if err != nil || repo.RecordKey() == "" { + writeNotif(w, http.StatusBadRequest, fmt.Sprintf("repo parameter invalid: %s", repoQuery)) + return + } + + s.resyncer.CancelResyncJob(repo) + writeNotif(w, http.StatusOK, "success") + } +} + +func writeNotif(w http.ResponseWriter, status int, msg string) { + w.Header().Set("Content-Type", "text/html") + w.WriteHeader(status) + + class := "info" + switch { + case status >= 500: + class = "error" + case status >= 400: + class = "warn" + } + + fmt.Fprintf(w, + `
%s
`, + class, + html.EscapeString(msg), + ) +} diff --git a/knotmirror/crawler.go b/knotmirror/crawler.go new file mode 100644 --- /dev/null +++ b/knotmirror/crawler.go @@ -0,0 +1,25 @@ +package knotmirror + +import ( + "context" + "database/sql" + "log/slog" + + "tangled.org/core/log" +) + +type Crawler struct { + logger *slog.Logger + db *sql.DB +} + +func NewCrawler(l *slog.Logger, db *sql.DB) *Crawler { + return &Crawler{ + logger: log.SubLogger(l, "crawler"), + db: db, + } +} + +func (c *Crawler) Start(ctx context.Context) { + // TODO: repository crawler +} diff --git a/knotmirror/git.go b/knotmirror/git.go new file mode 100644 --- /dev/null +++ b/knotmirror/git.go @@ -0,0 +1,305 @@ +package knotmirror + +import ( + "context" + "errors" + "fmt" + "net/url" + "os" + "os/exec" + "path/filepath" + "regexp" + "strings" + + "github.com/go-git/go-git/v5" + gitconfig "github.com/go-git/go-git/v5/config" + "github.com/go-git/go-git/v5/plumbing/transport" + "tangled.org/core/knotmirror/models" +) + +type GitMirrorManager interface { + Exist(repo *models.Repo) (bool, error) + // RemoteSetUrl updates git repository 'origin' remote + RemoteSetUrl(ctx context.Context, repo *models.Repo) error + // Clone clones the repository as a mirror + Clone(ctx context.Context, repo *models.Repo) error + // Fetch fetches the repository + Fetch(ctx context.Context, repo *models.Repo) error + // Sync mirrors the repository. It will clone the repository if repository doesn't exist. + Sync(ctx context.Context, repo *models.Repo) error +} + +type CliGitMirrorManager struct { + repoBasePath string + knotUseSSL bool +} + +func NewCliGitMirrorManager(repoBasePath string, knotUseSSL bool) *CliGitMirrorManager { + return &CliGitMirrorManager{ + repoBasePath, + knotUseSSL, + } +} + +var _ GitMirrorManager = new(CliGitMirrorManager) + +func (c *CliGitMirrorManager) makeRepoPath(repo *models.Repo) string { + return filepath.Join(c.repoBasePath, repo.Did.String(), repo.Rkey.String()) +} + +func (c *CliGitMirrorManager) Exist(repo *models.Repo) (bool, error) { + return isDir(c.makeRepoPath(repo)) +} + +func (c *CliGitMirrorManager) RemoteSetUrl(ctx context.Context, repo *models.Repo) error { + path := c.makeRepoPath(repo) + url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), c.knotUseSSL) + if err != nil { + return fmt.Errorf("constructing repo remote url: %w", err) + } + cmd := exec.CommandContext(ctx, "git", "-C", path, "remote", "set-url", "origin", url) + if out, err := cmd.CombinedOutput(); err != nil { + if ctx.Err() != nil { + return ctx.Err() + } + msg := string(out) + return fmt.Errorf("running 'git remote set-url origin %s': %w\n%s", url, err, msg) + } + return nil +} + +func (c *CliGitMirrorManager) Clone(ctx context.Context, repo *models.Repo) error { + path := c.makeRepoPath(repo) + url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), c.knotUseSSL) + if err != nil { + return fmt.Errorf("constructing repo remote url: %w", err) + } + return c.clone(ctx, path, url) +} + +func (c *CliGitMirrorManager) clone(ctx context.Context, path, url string) error { + cmd := exec.CommandContext(ctx, "git", "clone", "--mirror", url, path) + if out, err := cmd.CombinedOutput(); err != nil { + if ctx.Err() != nil { + return ctx.Err() + } + msg := string(out) + if classification := classifyCliError(msg); classification != nil { + return classification + } + return fmt.Errorf("running 'git clone --mirror %s': %w\n%s", url, err, msg) + } + return nil +} + +func (c *CliGitMirrorManager) Fetch(ctx context.Context, repo *models.Repo) error { + path := c.makeRepoPath(repo) + return c.fetch(ctx, path) +} + +func (c *CliGitMirrorManager) fetch(ctx context.Context, path string) error { + // TODO: use `repo.Knot` instead of depending on origin + cmd := exec.CommandContext(ctx, "git", "-C", path, "fetch", "--prune", "origin") + if out, err := cmd.CombinedOutput(); err != nil { + if ctx.Err() != nil { + return ctx.Err() + } + return fmt.Errorf("running 'git fetch': %w\n%s", err, string(out)) + } + return nil +} + +func (c *CliGitMirrorManager) Sync(ctx context.Context, repo *models.Repo) error { + path := c.makeRepoPath(repo) + url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), c.knotUseSSL) + if err != nil { + return fmt.Errorf("constructing repo remote url: %w", err) + } + + exist, err := isDir(path) + if err != nil { + return fmt.Errorf("checking repo path: %w", err) + } + if !exist { + if err := c.clone(ctx, path, url); err != nil { + return fmt.Errorf("cloning repo: %w", err) + } + } else { + if err := c.fetch(ctx, path); err != nil { + return fmt.Errorf("fetching repo: %w", err) + } + } + return nil +} + +var ( + ErrDNSFailure = errors.New("git: knot: dns failure (could not resolve host)") + ErrCertExpired = errors.New("git: knot: certificate has expired") + ErrCertMismatch = errors.New("git: knot: certificate hostname mismatch") + ErrTLSHandshake = errors.New("git: knot: tls handshake failure") + ErrHTTPStatus = errors.New("git: knot: request url returned error") + ErrUnreachable = errors.New("git: knot: could not connect to server") + ErrRepoNotFound = errors.New("git: repo: repository not found") +) + +var ( + reDNSFailure = regexp.MustCompile(`Could not resolve host:`) + reCertExpired = regexp.MustCompile(`SSL certificate OpenSSL verify result: certificate has expired`) + reCertMismatch = regexp.MustCompile(`SSL: no alternative certificate subject name matches target hostname`) + reTLSHandshake = regexp.MustCompile(`TLS connect error: (.*)`) + reHTTPStatus = regexp.MustCompile(`The requested URL returned error: (\d\d\d)`) + reUnreachable = regexp.MustCompile(`Could not connect to server`) + reRepoNotFound = regexp.MustCompile(`repository '.*?' not found`) +) + +// classifyCliError classifies git cli error message. It will return nil for unknown error messages +func classifyCliError(stderr string) error { + msg := strings.TrimSpace(stderr) + if m := reTLSHandshake.FindStringSubmatch(msg); len(m) > 1 { + return fmt.Errorf("%w: %s", ErrTLSHandshake, m[1]) + } + if m := reHTTPStatus.FindStringSubmatch(msg); len(m) > 1 { + return fmt.Errorf("%w: %s", ErrHTTPStatus, m[1]) + } + switch { + case reDNSFailure.MatchString(msg): + return ErrDNSFailure + case reCertExpired.MatchString(msg): + return ErrCertExpired + case reCertMismatch.MatchString(msg): + return ErrCertMismatch + case reUnreachable.MatchString(msg): + return ErrUnreachable + case reRepoNotFound.MatchString(msg): + return ErrRepoNotFound + } + return nil +} + +type GoGitMirrorManager struct { + repoBasePath string + knotUseSSL bool +} + +func NewGoGitMirrorClient(repoBasePath string, knotUseSSL bool) *GoGitMirrorManager { + return &GoGitMirrorManager{ + repoBasePath, + knotUseSSL, + } +} + +var _ GitMirrorManager = new(GoGitMirrorManager) + +func (c *GoGitMirrorManager) makeRepoPath(repo *models.Repo) string { + return filepath.Join(c.repoBasePath, repo.Did.String(), repo.Rkey.String()) +} + +func (c *GoGitMirrorManager) Exist(repo *models.Repo) (bool, error) { + return isDir(c.makeRepoPath(repo)) +} + +func (c *GoGitMirrorManager) RemoteSetUrl(ctx context.Context, repo *models.Repo) error { + panic("unimplemented") +} + +func (c *GoGitMirrorManager) Clone(ctx context.Context, repo *models.Repo) error { + path := c.makeRepoPath(repo) + url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), c.knotUseSSL) + if err != nil { + return fmt.Errorf("constructing repo remote url: %w", err) + } + return c.clone(ctx, path, url) +} + +func (c *GoGitMirrorManager) clone(ctx context.Context, path, url string) error { + _, err := git.PlainCloneContext(ctx, path, true, &git.CloneOptions{ + URL: url, + Mirror: true, + }) + if err != nil && !errors.Is(err, transport.ErrEmptyRemoteRepository) { + return fmt.Errorf("cloning repo: %w", err) + } + return nil +} + +func (c *GoGitMirrorManager) Fetch(ctx context.Context, repo *models.Repo) error { + path := c.makeRepoPath(repo) + url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), c.knotUseSSL) + if err != nil { + return fmt.Errorf("constructing repo remote url: %w", err) + } + + return c.fetch(ctx, path, url) +} + +func (c *GoGitMirrorManager) fetch(ctx context.Context, path, url string) error { + gr, err := git.PlainOpen(path) + if err != nil { + return fmt.Errorf("opening local repo: %w", err) + } + if err := gr.FetchContext(ctx, &git.FetchOptions{ + RemoteURL: url, + RefSpecs: []gitconfig.RefSpec{gitconfig.RefSpec("+refs/*:refs/*")}, + Force: true, + Prune: true, + }); err != nil { + return fmt.Errorf("fetching reppo: %w", err) + } + return nil +} + +func (c *GoGitMirrorManager) Sync(ctx context.Context, repo *models.Repo) error { + path := c.makeRepoPath(repo) + url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), c.knotUseSSL) + if err != nil { + return fmt.Errorf("constructing repo remote url: %w", err) + } + + exist, err := isDir(path) + if err != nil { + return fmt.Errorf("checking repo path: %w", err) + } + if !exist { + if err := c.clone(ctx, path, url); err != nil { + return fmt.Errorf("cloning repo: %w", err) + } + } else { + if err := c.fetch(ctx, path, url); err != nil { + return fmt.Errorf("fetching repo: %w", err) + } + } + return nil +} + +func makeRepoRemoteUrl(knot, didSlashRepo string, knotUseSSL bool) (string, error) { + if !strings.Contains(knot, "://") { + if knotUseSSL { + knot = "https://" + knot + } else { + knot = "http://" + knot + } + } + + u, err := url.Parse(knot) + if err != nil { + return "", err + } + + if u.Scheme != "http" && u.Scheme != "https" { + return "", fmt.Errorf("unsupported scheme: %s", u.Scheme) + } + + u = u.JoinPath(didSlashRepo) + return u.String(), nil +} + +func isDir(path string) (bool, error) { + info, err := os.Stat(path) + if err == nil && info.IsDir() { + return true, nil + } + if os.IsNotExist(err) { + return false, nil + } + return false, err +} diff --git a/knotmirror/knotmirror.go b/knotmirror/knotmirror.go new file mode 100644 --- /dev/null +++ b/knotmirror/knotmirror.go @@ -0,0 +1,117 @@ +package knotmirror + +import ( + "context" + "fmt" + "net/http" + _ "net/http/pprof" + "time" + + "github.com/prometheus/client_golang/prometheus/promhttp" + "tangled.org/core/knotmirror/config" + "tangled.org/core/knotmirror/db" + "tangled.org/core/knotmirror/knotstream" + "tangled.org/core/knotmirror/models" + "tangled.org/core/log" +) + +func Run(ctx context.Context, cfg *config.Config) error { + // make sure every services are cleaned up on fast return + ctx, cancel := context.WithCancel(ctx) + defer cancel() + + logger := log.FromContext(ctx) + + db, err := db.Make(ctx, cfg.DbUrl, 32) + if err != nil { + return fmt.Errorf("initializing db: %w", err) + } + + // NOTE: using plain git-cli for clone/fetch as go-git is too memory-intensive. + gitm := NewCliGitMirrorManager(cfg.GitRepoBasePath, cfg.KnotUseSSL) + + res, err := db.ExecContext(ctx, + `update repos set state = $1 where state = $2`, + models.RepoStateDesynchronized, + models.RepoStateResyncing, + ) + if err != nil { + return fmt.Errorf("clearing resyning states: %w", err) + } + rows, err := res.RowsAffected() + if err != nil { + return fmt.Errorf("getting affected rows: %w", err) + } + logger.Info(fmt.Sprintf("clearing resyning states: %d records updated", rows)) + + knotstream := knotstream.NewKnotStream(logger, db, cfg) + crawler := NewCrawler(logger, db) + resyncer := NewResyncer(logger, db, gitm, cfg) + adminpage := NewAdminServer(logger, db, resyncer) + + // maintain repository list with tap + // NOTE: this can be removed once we introduce did-for-repo because then we can just listen to KnotStream for #identity events. + tap := NewTapClient(logger, cfg, db, gitm, knotstream) + + // start metrics endpoint + go func() { + metricsAddr := cfg.MetricsListen + logger.Info("starting metrics server", "addr", metricsAddr) + http.Handle("/metrics", promhttp.Handler()) + if err := http.ListenAndServe(metricsAddr, nil); err != nil { + logger.Error("metrics server failed", "error", err) + } + }() + + // start admin page endpoint + go func() { + logger.Info("starting admin server", "addr", cfg.AdminListen) + if err := http.ListenAndServe(cfg.AdminListen, adminpage.Router()); err != nil { + logger.Error("admin server failed", "error", err) + } + }() + + tap.Start(ctx) + + resyncer.Start(ctx) + + // periodically crawl the entire network to mirror the repositories + crawler.Start(ctx) + + // listen to knotstream (currently we don't have relay for knots, so subscribe every known knots) + knotstream.Start(ctx) + + svcErr := make(chan error, 1) + if err := knotstream.ResubscribeAllHosts(ctx); err != nil { + svcErr <- fmt.Errorf("resubscribing known hosts: %w", err) + } + + logger.Info("startup complete") + select { + case <-ctx.Done(): + logger.Info("received shutdown signal", "reason", ctx.Err()) + case err := <-svcErr: + if err != nil { + logger.Error("service error", "error", err) + } + cancel() + } + + logger.Info("shutting down knotmirror") + shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 10*time.Second) + defer shutdownCancel() + + var errs []error + if err := knotstream.Shutdown(shutdownCtx); err != nil { + errs = append(errs, err) + } + if err := db.Close(); err != nil { + errs = append(errs, err) + } + for _, err := range errs { + logger.Error("error during shutdown", "err", err) + } + + logger.Info("shutdown complete") + return nil +} diff --git a/knotmirror/metrics.go b/knotmirror/metrics.go new file mode 100644 --- /dev/null +++ b/knotmirror/metrics.go @@ -0,0 +1,29 @@ +package knotmirror + +import ( + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promauto" +) + +// Resync metrics +var ( + // TODO: + // - working / waiting resycner counts + resyncsStarted = promauto.NewCounter(prometheus.CounterOpts{ + Name: "knotmirror_resyncs_started_total", + Help: "Total number of repo resyncs started", + }) + resyncsCompleted = promauto.NewCounter(prometheus.CounterOpts{ + Name: "knotmirror_resyncs_completed_total", + Help: "Total number of repo resyncs completed", + }) + resyncsFailed = promauto.NewCounter(prometheus.CounterOpts{ + Name: "knotmirror_resyncs_failed_total", + Help: "Total number of repo resyncs failed", + }) + resyncDuration = promauto.NewHistogram(prometheus.HistogramOpts{ + Name: "knotmirror_resync_duration_seconds", + Help: "Duration of repo resync operations", + Buckets: prometheus.ExponentialBuckets(0.1, 2, 12), + }) +) diff --git a/knotmirror/readme.md b/knotmirror/readme.md new file mode 100644 --- /dev/null +++ b/knotmirror/readme.md @@ -0,0 +1,8 @@ +# KnotMirror + +KnotMirror is a git mirror service for all known repos. Heavily inspired by [indigo/relay] and [indigo/tap]. + +KnotMirror syncs repo list using tap and subscribe to all known knots as KnotStream. + +[indigo/relay]: https://github.com/bluesky-social/indigo/tree/main/cmd/relay +[indigo/tap]: https://github.com/bluesky-social/indigo/tree/main/cmd/tap diff --git a/knotmirror/resyncer.go b/knotmirror/resyncer.go new file mode 100644 --- /dev/null +++ b/knotmirror/resyncer.go @@ -0,0 +1,273 @@ +package knotmirror + +import ( + "context" + "database/sql" + "errors" + "fmt" + "log/slog" + "math/rand" + "strings" + "sync" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/knotmirror/config" + "tangled.org/core/knotmirror/db" + "tangled.org/core/knotmirror/models" + "tangled.org/core/log" +) + +type Resyncer struct { + logger *slog.Logger + db *sql.DB + gitm GitMirrorManager + + claimJobMu sync.Mutex + + runningJobs map[syntax.ATURI]context.CancelFunc + runningJobsMu sync.Mutex + + repoFetchTimeout time.Duration + manualResyncTimeout time.Duration + parallelism int +} + +func NewResyncer(l *slog.Logger, db *sql.DB, gitm GitMirrorManager, cfg *config.Config) *Resyncer { + return &Resyncer{ + logger: log.SubLogger(l, "resyncer"), + db: db, + gitm: gitm, + + runningJobs: make(map[syntax.ATURI]context.CancelFunc), + + repoFetchTimeout: cfg.GitRepoFetchTimeout, + manualResyncTimeout: 30 * time.Minute, + parallelism: cfg.ResyncParallelism, + } +} + +func (r *Resyncer) Start(ctx context.Context) { + for i := 0; i < r.parallelism; i++ { + go r.runResyncWorker(ctx, i) + } +} + +func (r *Resyncer) runResyncWorker(ctx context.Context, workerID int) { + l := r.logger.With("worker", workerID) + for { + select { + case <-ctx.Done(): + l.Info("resync worker shutting down", "error", ctx.Err()) + return + default: + } + repoAt, found, err := r.claimResyncJob(ctx) + if err != nil { + l.Error("failed to claim resync job", "error", err) + time.Sleep(time.Second) + continue + } + if !found { + time.Sleep(time.Second) + continue + } + l.Info("processing resync", "aturi", repoAt) + if err := r.resyncRepo(ctx, repoAt); err != nil { + l.Error("resync failed", "aturi", repoAt, "error", err) + } + } +} + +func (r *Resyncer) registerRunning(repo syntax.ATURI, cancel context.CancelFunc) { + r.runningJobsMu.Lock() + defer r.runningJobsMu.Unlock() + + if _, exists := r.runningJobs[repo]; exists { + return + } + r.runningJobs[repo] = cancel +} + +func (r *Resyncer) unregisterRunning(repo syntax.ATURI) { + r.runningJobsMu.Lock() + defer r.runningJobsMu.Unlock() + + delete(r.runningJobs, repo) +} + +func (r *Resyncer) CancelResyncJob(repo syntax.ATURI) { + r.runningJobsMu.Lock() + defer r.runningJobsMu.Unlock() + + cancel, ok := r.runningJobs[repo] + if !ok { + return + } + delete(r.runningJobs, repo) + cancel() +} + +// TriggerResyncJob manually triggers the resync job +func (r *Resyncer) TriggerResyncJob(ctx context.Context, repoAt syntax.ATURI) error { + repo, err := db.GetRepoByAtUri(ctx, r.db, repoAt) + if err != nil { + return fmt.Errorf("failed to get repo: %w", err) + } + if repo == nil { + return fmt.Errorf("repo not found: %s", repoAt) + } + + if repo.State == models.RepoStateResyncing { + return fmt.Errorf("repo already resyncing") + } + + repo.State = models.RepoStatePending + repo.RetryAfter = -1 // resyncer will prioritize this + + if err := db.UpsertRepo(ctx, r.db, repo); err != nil { + return fmt.Errorf("updating repo state to pending %w", err) + } + return nil +} + +func (r *Resyncer) claimResyncJob(ctx context.Context) (syntax.ATURI, bool, error) { + // use mutex to prevent duplicated jobs + r.claimJobMu.Lock() + defer r.claimJobMu.Unlock() + + var repoAt syntax.ATURI + now := time.Now().Unix() + if err := r.db.QueryRowContext(ctx, + `update repos + set state = $1 + where at_uri = ( + select at_uri from repos + where state in ($2, $3, $4) + and (retry_after = -1 or retry_after = 0 or retry_after < $5) + order by + (retry_after = -1) desc, + (retry_after = 0) desc, + retry_after + limit 1 + ) + returning at_uri + `, + models.RepoStateResyncing, + models.RepoStatePending, models.RepoStateDesynchronized, models.RepoStateError, + now, + ).Scan(&repoAt); err != nil { + if errors.Is(err, sql.ErrNoRows) { + return "", false, nil + } + return "", false, err + } + + return repoAt, true, nil +} + +func (r *Resyncer) resyncRepo(ctx context.Context, repoAt syntax.ATURI) error { + // ctx, span := tracer.Start(ctx, "resyncRepo") + // span.SetAttributes(attribute.String("aturi", repoAt)) + // defer span.End() + + resyncsStarted.Inc() + startTime := time.Now() + + jobCtx, cancel := context.WithCancel(ctx) + r.registerRunning(repoAt, cancel) + defer r.unregisterRunning(repoAt) + + success, err := r.doResync(jobCtx, repoAt) + if !success { + resyncsFailed.Inc() + resyncDuration.Observe(time.Since(startTime).Seconds()) + return r.handleResyncFailure(ctx, repoAt, err) + } + + resyncsCompleted.Inc() + resyncDuration.Observe(time.Since(startTime).Seconds()) + return nil +} + +func (r *Resyncer) doResync(ctx context.Context, repoAt syntax.ATURI) (bool, error) { + // ctx, span := tracer.Start(ctx, "doResync") + // span.SetAttributes(attribute.String("aturi", repoAt)) + // defer span.End() + + repo, err := db.GetRepoByAtUri(ctx, r.db, repoAt) + if err != nil { + return false, fmt.Errorf("failed to get repo: %w", err) + } + if repo == nil { // untracked repo, skip + return false, nil + } + + // TODO: check if Knot is on backoff list. If so, return (false, nil) + // TODO: detect rate limit error (http.StatusTooManyRequests) to put Knot in backoff list + + timeout := r.repoFetchTimeout + if repo.RetryAfter == -1 { + timeout = r.manualResyncTimeout + } + fetchCtx, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + + if err := r.gitm.Sync(fetchCtx, repo); err != nil { + return false, err + } + + // repo.GitRev = + // repo.RepoSha = + repo.State = models.RepoStateActive + repo.ErrorMsg = "" + repo.RetryCount = 0 + repo.RetryAfter = 0 + if err := db.UpsertRepo(ctx, r.db, repo); err != nil { + return false, fmt.Errorf("updating repo state to active %w", err) + } + return true, nil +} + +func (r *Resyncer) handleResyncFailure(ctx context.Context, repoAt syntax.ATURI, err error) error { + r.logger.Debug("handleResyncFailure", "at_uri", repoAt, "err", err) + var state models.RepoState + var errMsg string + if err == nil { + state = models.RepoStateDesynchronized + errMsg = "" + } else { + state = models.RepoStateError + errMsg = err.Error() + } + + repo, err := db.GetRepoByAtUri(ctx, r.db, repoAt) + if err != nil { + return fmt.Errorf("failed to get repo: %w", err) + } + if repo == nil { + return fmt.Errorf("failed to get repo. repo '%s' doesn't exist in db", repoAt) + } + + // start a 1 min & go up to 1 hr between retries + var retryCount = repo.RetryCount + 1 + var retryAfter = time.Now().Add(backoff(retryCount, 60) * 60).Unix() + + // remove null bytes + errMsg = strings.ReplaceAll(errMsg, "\x00", "") + + repo.State = state + repo.ErrorMsg = errMsg + repo.RetryCount = retryCount + repo.RetryAfter = retryAfter + if err := db.UpsertRepo(ctx, r.db, repo); err != nil { + return fmt.Errorf("failed to update repo state: %w", err) + } + return nil +} + +func backoff(retries int, max int) time.Duration { + dur := min(1< 0 { + pageClause = " limit $1 offset $2 " + args = append(args, page.Limit, page.Offset) + } + + whereClause := "" + if did != "" { + conditions = append(conditions, fmt.Sprintf("did = $%d", len(args)+1)) + args = append(args, did) + } + if knot != "" { + conditions = append(conditions, fmt.Sprintf("knot_domain = $%d", len(args)+1)) + args = append(args, knot) + } + if state != "" { + conditions = append(conditions, fmt.Sprintf("state = $%d", len(args)+1)) + args = append(args, state) + } + if len(conditions) > 0 { + whereClause = "WHERE " + conditions[0] + for _, condition := range conditions[1:] { + whereClause += " AND " + condition + } + } + + query := ` + select + did, + rkey, + cid, + name, + knot_domain, + git_rev, + repo_sha, + state, + error_msg, + retry_count, + retry_after + from repos + ` + whereClause + pageClause + rows, err := e.QueryContext(ctx, query, args...) + if err != nil { + return nil, err + } + defer rows.Close() + + var repos []models.Repo + for rows.Next() { + var repo models.Repo + if err := rows.Scan( + &repo.Did, + &repo.Rkey, + &repo.Cid, + &repo.Name, + &repo.KnotDomain, + &repo.GitRev, + &repo.RepoSha, + &repo.State, + &repo.ErrorMsg, + &repo.RetryCount, + &repo.RetryAfter, + ); err != nil { + return nil, fmt.Errorf("scanning row: %w", err) + } + repos = append(repos, repo) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("scanning rows: %w ", err) + } + + return repos, nil +} + +func GetRepoCountsByState(ctx context.Context, e *sql.DB) (map[models.RepoState]int64, error) { + const q = ` + SELECT state, COUNT(*) + FROM repos + GROUP BY state + ` + + rows, err := e.QueryContext(ctx, q) + if err != nil { + return nil, err + } + defer rows.Close() + + counts := make(map[models.RepoState]int64) + + for rows.Next() { + var state string + var count int64 + + if err := rows.Scan(&state, &count); err != nil { + return nil, err + } + + counts[models.RepoState(state)] = count + } + + if err := rows.Err(); err != nil { + return nil, err + } + + for _, s := range models.AllRepoStates { + if _, ok := counts[s]; !ok { + counts[s] = 0 + } + } + + return counts, nil +} diff --git a/knotmirror/knotstream/knotstream.go b/knotmirror/knotstream/knotstream.go new file mode 100644 --- /dev/null +++ b/knotmirror/knotstream/knotstream.go @@ -0,0 +1,88 @@ +package knotstream + +import ( + "context" + "database/sql" + "fmt" + "log/slog" + "time" + + "tangled.org/core/knotmirror/config" + "tangled.org/core/knotmirror/db" + "tangled.org/core/knotmirror/models" + "tangled.org/core/log" +) + +type KnotStream struct { + logger *slog.Logger + db *sql.DB + slurper *KnotSlurper +} + +func NewKnotStream(l *slog.Logger, db *sql.DB, cfg *config.Config) *KnotStream { + l = log.SubLogger(l, "knotstream") + return &KnotStream{ + logger: l, + db: db, + slurper: NewKnotSlurper(l, db, cfg.Slurper), + } +} + +func (s *KnotStream) Start(ctx context.Context) { + go s.slurper.Run(ctx) +} + +func (s *KnotStream) Shutdown(ctx context.Context) error { + return s.slurper.Shutdown(ctx) +} + +func (s *KnotStream) CheckIfSubscribed(hostname string) bool { + return s.slurper.CheckIfSubscribed(hostname) +} + +func (s *KnotStream) SubscribeHost(ctx context.Context, hostname string, noSSL bool) error { + l := s.logger.With("hostname", hostname, "nossl", noSSL) + l.Debug("subscribe") + host, err := db.GetHost(ctx, s.db, hostname) + if err != nil { + return fmt.Errorf("loading host from db: %w", err) + } + + if host == nil { + host = &models.Host{ + Hostname: hostname, + NoSSL: noSSL, + Status: models.HostStatusActive, + LastSeq: 0, + } + + if err := db.UpsertHost(ctx, s.db, host); err != nil { + return fmt.Errorf("adding host to db: %w", err) + } + + l.Info("adding new host subscription") + } + + if host.Status == models.HostStatusBanned { + return fmt.Errorf("cannot subscribe to banned knot") + } + return s.slurper.Subscribe(ctx, *host) +} + +func (s *KnotStream) ResubscribeAllHosts(ctx context.Context) error { + hosts, err := db.ListHosts(ctx, s.db, models.HostStatusActive) + if err != nil { + return fmt.Errorf("listing hosts: %w", err) + } + + for _, host := range hosts { + l := s.logger.With("hostname", host.Hostname) + l.Info("re-subscribing to active host") + if err := s.slurper.Subscribe(ctx, host); err != nil { + l.Warn("failed to re-subscribe to host", "err", err) + } + // sleep for a very short period, so we don't open tons of sockets at the same time + time.Sleep(1 * time.Millisecond) + } + return nil +} diff --git a/knotmirror/knotstream/metrics.go b/knotmirror/knotstream/metrics.go new file mode 100644 --- /dev/null +++ b/knotmirror/knotstream/metrics.go @@ -0,0 +1,28 @@ +package knotstream + +import ( + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promauto" +) + +// KnotStream metrics +var ( + knotstreamEventsReceived = promauto.NewCounter(prometheus.CounterOpts{ + Name: "knotmirror_knotstream_events_received_total", + Help: "Total number of events received from knotstream", + }) + knotstreamEventsProcessed = promauto.NewCounter(prometheus.CounterOpts{ + Name: "knotmirror_knotstream_events_processed_total", + Help: "Total number of events successfully processed", + }) + knotstreamEventsSkipped = promauto.NewCounter(prometheus.CounterOpts{ + Name: "knotmirror_knotstream_events_skipped_total", + Help: "Total number of events skipped (not tracked)", + }) +) + +// slurper metrics +var connectedInbound = promauto.NewGauge(prometheus.GaugeOpts{ + Name: "knotmirror_connected_inbound", + Help: "Number of inbound knotstream we are consuming", +}) diff --git a/knotmirror/knotstream/scheduler.go b/knotmirror/knotstream/scheduler.go new file mode 100644 --- /dev/null +++ b/knotmirror/knotstream/scheduler.go @@ -0,0 +1,102 @@ +package knotstream + +import ( + "context" + "log/slog" + "sync" + "sync/atomic" + "time" + + "tangled.org/core/log" +) + +type ParallelScheduler struct { + concurrency int + + do func(ctx context.Context, task *Task) error + + feeder chan *Task + lk sync.Mutex + scheduled map[string][]*Task + lastSeq atomic.Int64 + + logger *slog.Logger +} + +type Task struct { + key string + message []byte +} + +func NewParallelScheduler(maxC int, ident string, do func(context.Context, *Task) error) *ParallelScheduler { + return &ParallelScheduler{ + concurrency: maxC, + do: do, + feeder: make(chan *Task), + scheduled: make(map[string][]*Task), + logger: log.New("parallel-scheduler"), + } +} + +func (s *ParallelScheduler) Start(ctx context.Context) { + for range s.concurrency { + go s.ForEach(ctx, s.do) + } +} + +func (s *ParallelScheduler) AddTask(ctx context.Context, task *Task) { + s.lk.Lock() + if st, ok := s.scheduled[task.key]; ok { + // schedule task + s.scheduled[task.key] = append(st, task) + s.lk.Unlock() + return + } + s.scheduled[task.key] = []*Task{} + s.lk.Unlock() + + select { + case <-ctx.Done(): + return + case s.feeder <- task: + return + } +} + +func (s *ParallelScheduler) ForEach(ctx context.Context, fn func(context.Context, *Task) error) { + for task := range s.feeder { + for task != nil { + select { + case <-ctx.Done(): + return + default: + } + if err := fn(ctx, task); err != nil { + s.logger.Error("event handler failed", "err", err) + } + + s.lk.Lock() + func() { + rem, ok := s.scheduled[task.key] + if !ok { + s.logger.Error("should always have an 'active' entry if a worker is processing a job") + } + if len(rem) == 0 { + delete(s.scheduled, task.key) + task = nil + } else { + task = rem[0] + s.scheduled[task.key] = rem[1:] + } + + // TODO: update seq from received message + s.lastSeq.Store(time.Now().UnixNano()) + }() + s.lk.Unlock() + } + } +} + +func (s *ParallelScheduler) LastSeq() int64 { + return s.lastSeq.Load() +} diff --git a/knotmirror/knotstream/slurper.go b/knotmirror/knotstream/slurper.go new file mode 100644 --- /dev/null +++ b/knotmirror/knotstream/slurper.go @@ -0,0 +1,334 @@ +package knotstream + +import ( + "context" + "database/sql" + "encoding/json" + "fmt" + "log/slog" + "math/rand" + "net/http" + "sync" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/bluesky-social/indigo/util/ssrf" + "github.com/carlmjohnson/versioninfo" + "github.com/gorilla/websocket" + "tangled.org/core/api/tangled" + "tangled.org/core/knotmirror/config" + "tangled.org/core/knotmirror/db" + "tangled.org/core/knotmirror/models" + "tangled.org/core/log" +) + +type KnotSlurper struct { + logger *slog.Logger + db *sql.DB + cfg config.SlurperConfig + + subsLk sync.Mutex + subs map[string]*subscription +} + +func NewKnotSlurper(l *slog.Logger, db *sql.DB, cfg config.SlurperConfig) *KnotSlurper { + return &KnotSlurper{ + logger: log.SubLogger(l, "slurper"), + db: db, + cfg: cfg, + subs: make(map[string]*subscription), + } +} + +func (s *KnotSlurper) Run(ctx context.Context) { + for { + select { + case <-ctx.Done(): + return + case <-time.After(s.cfg.PersistCursorPeriod): + if err := s.persistCursors(ctx); err != nil { + s.logger.Error("failed to flush cursors", "err", err) + } + } + } +} + +func (s *KnotSlurper) CheckIfSubscribed(hostname string) bool { + s.subsLk.Lock() + defer s.subsLk.Unlock() + + _, ok := s.subs[hostname] + return ok +} + +func (s *KnotSlurper) Shutdown(ctx context.Context) error { + s.logger.Info("starting shutdown host cursor flush") + err := s.persistCursors(ctx) + if err != nil { + s.logger.Error("shutdown error", "err", err) + } + s.logger.Info("slurper shutdown complete") + return err +} + +func (s *KnotSlurper) persistCursors(ctx context.Context) error { + // // gather cursor list from subscriptions and store them to DB + // start := time.Now() + + s.subsLk.Lock() + cursors := make([]models.HostCursor, len(s.subs)) + i := 0 + for _, sub := range s.subs { + cursors[i] = sub.HostCursor() + i++ + } + s.subsLk.Unlock() + + err := db.StoreCursors(ctx, s.db, cursors) + // s.logger.Info("finished persisting cursors", "count", len(cursors), "duration", time.Since(start).String(), "err", err) + return err +} + +func (s *KnotSlurper) Subscribe(ctx context.Context, host models.Host) error { + s.subsLk.Lock() + defer s.subsLk.Unlock() + + _, ok := s.subs[host.Hostname] + if ok { + return fmt.Errorf("already subscribed: %s", host.Hostname) + } + + // TODO: include `cancel` function to kill subscription by hostname + sub := &subscription{ + hostname: host.Hostname, + scheduler: NewParallelScheduler( + s.cfg.ConcurrencyPerHost, + host.Hostname, + s.ProcessEvent, + ), + } + s.subs[host.Hostname] = sub + + sub.scheduler.Start(ctx) + go s.subscribeWithRedialer(ctx, host, sub) + return nil +} + +func (s *KnotSlurper) subscribeWithRedialer(ctx context.Context, host models.Host, sub *subscription) { + l := s.logger.With("host", host.Hostname) + + dialer := websocket.Dialer{ + HandshakeTimeout: time.Second * 5, + } + + // if this isn't a localhost / private connection, then we should enable SSRF protections + if !host.NoSSL { + netDialer := ssrf.PublicOnlyDialer() + dialer.NetDialContext = netDialer.DialContext + } + + cursor := host.LastSeq + + connectedInbound.Inc() + defer connectedInbound.Dec() + + var backoff int + for { + select { + case <-ctx.Done(): + return + default: + } + u := host.LegacyEventsURL(cursor) + l.Debug("made url with cursor", "cursor", cursor, "url", u) + + // NOTE: manual backoff retry implementation to explicitly handle fails + hdr := make(http.Header) + hdr.Add("User-Agent", userAgent()) + conn, resp, err := dialer.DialContext(ctx, u, hdr) + if err != nil { + l.Warn("dialing failed", "err", err, "backoff", backoff) + time.Sleep(sleepForBackoff(backoff)) + backoff++ + if backoff > 30 { + l.Warn("host does not appear to be online, disabling for now") + host.Status = models.HostStatusOffline + if err := db.UpsertHost(ctx, s.db, &host); err != nil { + l.Error("failed to update host status", "err", err) + } + return + } + continue + } + + l.Debug("knot event subscription response", "code", resp.StatusCode, "url", u) + + if err := s.handleConnection(ctx, conn, sub); err != nil { + // TODO: measure the last N connection error times and if they're coming too fast reconnect slower or don't reconnect and wait for requestCrawl + l.Warn("host connection failed", "err", err, "backoff", backoff) + } + + updatedCursor := sub.LastSeq() + didProgress := updatedCursor > cursor + l.Debug("cursor compare", "cursor", cursor, "updatedCursor", updatedCursor, "didProgress", didProgress) + if cursor == 0 || didProgress { + cursor = updatedCursor + backoff = 0 + + batch := []models.HostCursor{sub.HostCursor()} + if err := db.StoreCursors(ctx, s.db, batch); err != nil { + l.Error("failed to store cursors", "err", err) + } + } + } +} + +// handleConnection handles websocket connection. +// Schedules task from received event and return when connection is closed +func (s *KnotSlurper) handleConnection(ctx context.Context, conn *websocket.Conn, sub *subscription) error { + // ping on every 30s + ctx, cancel := context.WithCancel(ctx) + defer cancel() // close the background ping job on connection close + go func() { + t := time.NewTicker(30 * time.Second) + defer t.Stop() + failcount := 0 + + for { + select { + case <-t.C: + if err := conn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second*10)); err != nil { + s.logger.Warn("failed to ping", "err", err) + failcount++ + if failcount >= 4 { + s.logger.Error("too many ping fails", "count", failcount) + _ = conn.Close() + return + } + } else { + failcount = 0 // ok ping + } + case <-ctx.Done(): + _ = conn.Close() + return + } + } + }() + + conn.SetPingHandler(func(message string) error { + err := conn.WriteControl(websocket.PongMessage, []byte(message), time.Now().Add(time.Minute)) + if err == websocket.ErrCloseSent { + return nil + } + return err + }) + conn.SetPongHandler(func(_ string) error { + if err := conn.SetReadDeadline(time.Now().Add(time.Minute)); err != nil { + s.logger.Error("failed to set read deadline", "err", err) + } + return nil + }) + + for { + select { + case <-ctx.Done(): + return ctx.Err() + default: + } + msgType, msg, err := conn.ReadMessage() + if err != nil { + return err + } + + if msgType != websocket.TextMessage { + continue + } + + sub.scheduler.AddTask(ctx, &Task{ + key: sub.hostname, // TODO: replace to repository AT-URI for better concurrency + message: msg, + }) + } +} + +type LegacyGitEvent struct { + Rkey string + Nsid string + Event tangled.GitRefUpdate +} + +func (s *KnotSlurper) ProcessEvent(ctx context.Context, task *Task) error { + var legacyMessage LegacyGitEvent + if err := json.Unmarshal(task.message, &legacyMessage); err != nil { + return fmt.Errorf("unmarshaling message: %w", err) + } + + if err := s.ProcessLegacyGitRefUpdate(ctx, &legacyMessage); err != nil { + return fmt.Errorf("processing gitRefUpdate: %w", err) + } + return nil +} + +func (s *KnotSlurper) ProcessLegacyGitRefUpdate(ctx context.Context, evt *LegacyGitEvent) error { + knotstreamEventsReceived.Inc() + + curr, err := db.GetRepoByName(ctx, s.db, syntax.DID(evt.Event.RepoDid), evt.Event.RepoName) + if err != nil { + return fmt.Errorf("failed to get repo '%s': %w", evt.Event.RepoDid+"/"+evt.Event.RepoName, err) + } + if curr == nil { + // if repo doesn't exist in DB, just ignore the event. That repo is unknown. + // + // Normally did+name is already enough to perform git-fetch as that's + // what needed to fetch the repository. + // But we want to store that in did/rkey in knot-mirror. + // Therefore, we should ignore when the repository is unknown. + // Hopefully crawler will sync it later. + s.logger.Warn("skipping event from unknown repo", "did/repo", evt.Event.RepoDid+"/"+evt.Event.RepoName) + knotstreamEventsSkipped.Inc() + return nil + } + l := s.logger.With("repoAt", curr.AtUri()) + + // TODO: should plan resync to resyncBuffer on RepoStateResyncing + if curr.State != models.RepoStateActive { + l.Debug("skipping non-active repo") + knotstreamEventsSkipped.Inc() + return nil + } + + if curr.GitRev != "" && evt.Rkey <= curr.GitRev.String() { + l.Debug("skipping replayed event", "event.Rkey", evt.Rkey, "currentRev", curr.GitRev) + knotstreamEventsSkipped.Inc() + return nil + } + + // if curr.State == models.RepoStateResyncing { + // firehoseEventsSkipped.Inc() + // return fp.events.addToResyncBuffer(ctx, commit) + // } + + // can't skip anything, update repo state + if err := db.UpdateRepoState(ctx, s.db, curr.Did, curr.Rkey, models.RepoStateDesynchronized); err != nil { + return err + } + + l.Info("event processed", "eventRev", evt.Rkey) + + knotstreamEventsProcessed.Inc() + return nil +} + +func userAgent() string { + return fmt.Sprintf("knotmirror/%s", versioninfo.Short()) +} + +func sleepForBackoff(b int) time.Duration { + if b == 0 { + return 0 + } + if b < 10 { + return time.Millisecond * time.Duration((50*b)+rand.Intn(500)) + } + return time.Second * 30 +} diff --git a/knotmirror/knotstream/subscription.go b/knotmirror/knotstream/subscription.go new file mode 100644 --- /dev/null +++ b/knotmirror/knotstream/subscription.go @@ -0,0 +1,22 @@ +package knotstream + +import "tangled.org/core/knotmirror/models" + +// subscription represents websocket connection with that host +type subscription struct { + hostname string + + // embedded parallel job scheduler + scheduler *ParallelScheduler +} + +func (s *subscription) LastSeq() int64 { + return s.scheduler.LastSeq() +} + +func (s *subscription) HostCursor() models.HostCursor { + return models.HostCursor{ + Hostname: s.hostname, + LastSeq: s.LastSeq(), + } +} diff --git a/knotmirror/models/models.go b/knotmirror/models/models.go new file mode 100644 --- /dev/null +++ b/knotmirror/models/models.go @@ -0,0 +1,110 @@ +package models + +import ( + "fmt" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" +) + +type Repo struct { + Did syntax.DID + Rkey syntax.RecordKey + Cid *syntax.CID + // content of tangled.Repo + Name string + KnotDomain string + + GitRev syntax.TID // last processed git.refUpdate revision + RepoSha string // sha256 sum of git refs (to avoid no-op git fetch) + State RepoState + ErrorMsg string + RetryCount int + RetryAfter int64 // Unix timestamp (seconds) +} + +func (r *Repo) AtUri() syntax.ATURI { + return syntax.ATURI(fmt.Sprintf("at://%s/%s/%s", r.Did, tangled.RepoNSID, r.Rkey)) +} + +func (r *Repo) DidSlashRepo() string { + return fmt.Sprintf("%s/%s", r.Did, r.Name) +} + +type RepoState string + +const ( + RepoStatePending RepoState = "pending" + RepoStateDesynchronized RepoState = "desynchronized" + RepoStateResyncing RepoState = "resyncing" + RepoStateActive RepoState = "active" + RepoStateSuspended RepoState = "suspended" + RepoStateError RepoState = "error" +) + +var AllRepoStates = []RepoState{ + RepoStatePending, + RepoStateDesynchronized, + RepoStateResyncing, + RepoStateActive, + RepoStateSuspended, + RepoStateError, +} + +func (s RepoState) IsResyncing() bool { + return s == RepoStateResyncing +} + +type HostCursor struct { + Hostname string + LastSeq int64 +} + +type Host struct { + Hostname string + NoSSL bool + Status HostStatus + LastSeq int64 +} + +type HostStatus string + +const ( + HostStatusActive HostStatus = "active" + HostStatusIdle HostStatus = "idle" + HostStatusOffline HostStatus = "offline" + HostStatusThrottled HostStatus = "throttled" + HostStatusBanned HostStatus = "banned" +) + +var AllHostStatuses = []HostStatus{ + HostStatusActive, + HostStatusIdle, + HostStatusOffline, + HostStatusThrottled, + HostStatusBanned, +} + +// func (h *Host) SubscribeGitRefsURL(cursor int64) string { +// scheme := "wss" +// if h.NoSSL { +// scheme = "ws" +// } +// u := fmt.Sprintf("%s://%s/xrpc/%s", scheme, h.Hostname, tangled.SubscribeGitRefsNSID) +// if cursor > 0 { +// u = fmt.Sprintf("%s?cursor=%d", u, h.LastSeq) +// } +// return u +// } + +func (h *Host) LegacyEventsURL(cursor int64) string { + scheme := "wss" + if h.NoSSL { + scheme = "ws" + } + u := fmt.Sprintf("%s://%s/events", scheme, h.Hostname) + if cursor > 0 { + u = fmt.Sprintf("%s?cursor=%d", u, cursor) + } + return u +} diff --git a/knotmirror/templates/base.html b/knotmirror/templates/base.html new file mode 100644 --- /dev/null +++ b/knotmirror/templates/base.html @@ -0,0 +1,55 @@ +{{define "base"}} + + + + KnotMirror Admin + + + + + +
+ {{template "content" .}} +
+
+ + + +{{end}} diff --git a/knotmirror/templates/hosts.html b/knotmirror/templates/hosts.html new file mode 100644 --- /dev/null +++ b/knotmirror/templates/hosts.html @@ -0,0 +1,44 @@ +{{template "base" .}} +{{define "content"}} +

Knot Hosts

+ +
+
+ + +
+
+ + + + + + + + + + + + {{range .Hosts}} + + + + + + + {{else}} + + {{end}} + +
HostnameSSLStatusLast Seq
{{.Hostname}}{{if .NoSSL}}False{{else}}True{{end}}{{.Status}}{{.LastSeq}}
No hosts registered.
+{{end}} diff --git a/knotmirror/templates/repos.html b/knotmirror/templates/repos.html new file mode 100644 --- /dev/null +++ b/knotmirror/templates/repos.html @@ -0,0 +1,86 @@ +{{template "base" .}} +{{define "content"}} +

Repositories

+ +
+
+ + + + + Clear +
+
+ +
+
+ {{range const.AllRepoStates}} + + {{.}}: {{index $.RepoCounts .}} + + {{end}} +
+ + + + + + + + + + + + + + + {{range .Repos}} + + + + + + + + + + + {{else}} + + {{end}} + +
DIDNameKnotStateRetryRetry AfterError MessageAction
{{.Did}}{{.Name}}{{.KnotDomain}}{{.State}}{{.RetryCount}}{{readt .RetryAfter}}{{.ErrorMsg}} +
+ + +
+
No repositories found.
+
+ + +{{end}} diff --git a/nix/pkgs/knot-mirror.nix b/nix/pkgs/knot-mirror.nix new file mode 100644 --- /dev/null +++ b/nix/pkgs/knot-mirror.nix @@ -0,0 +1,18 @@ +{ + buildGoApplication, + modules, + src, +}: +buildGoApplication { + pname = "knotmirror"; + version = "0.1.0"; + inherit src modules; + + doCheck = false; + + subPackages = ["cmd/knotmirror"]; + + meta = { + mainProgram = "knotmirror"; + }; +} -- tangled.sh