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,
+ `
`,
+ 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
+
+
+
+
+
+
+
+
+ | Hostname |
+ SSL |
+ Status |
+ Last Seq |
+
+
+
+ {{range .Hosts}}
+
+ | {{.Hostname}} |
+ {{if .NoSSL}}False{{else}}True{{end}} |
+ {{.Status}} |
+ {{.LastSeq}} |
+
+ {{else}}
+ | No hosts registered. |
+ {{end}}
+
+
+{{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
+
+
+
+
+
+ {{range const.AllRepoStates}}
+
+ {{.}}: {{index $.RepoCounts .}}
+
+ {{end}}
+
+
+
+
+ | DID |
+ Name |
+ Knot |
+ State |
+ Retry |
+ Retry After |
+ Error Message |
+ Action |
+
+
+
+ {{range .Repos}}
+
+ {{.Did}} |
+ {{.Name}} |
+ {{.KnotDomain}} |
+ {{.State}} |
+ {{.RetryCount}} |
+ {{readt .RetryAfter}} |
+ {{.ErrorMsg}} |
+
+
+ |
+
+ {{else}}
+ | No repositories found. |
+ {{end}}
+
+
+
+
+
+{{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";
+ };
+}