diff --git a/.tangled/workflows/test.yml b/.tangled/workflows/test.yml index 69a1b62c..a31a616d 100644 --- a/.tangled/workflows/test.yml +++ b/.tangled/workflows/test.yml @@ -8,12 +8,28 @@ dependencies: nixpkgs: - go - gcc + - postgresql + - shadow + - util-linux steps: - name: patch static dir command: | mkdir -p appview/pages/static; touch appview/pages/static/x + - name: start postgres + command: | + set -euo pipefail + mkdir -p /pg/socket + useradd -r -m -d /pg/home pguser + chown -R pguser /pg + uid=$(id -u pguser); gid=$(id -g pguser) + setpriv --reuid "$uid" --regid "$gid" --init-groups -- initdb -D /pg/home/data --auth=trust --username=postgres + setpriv --reuid "$uid" --regid "$gid" --init-groups -- pg_ctl -D /pg/home/data -l /pg/home/pg.log -o "-p 5432 -k /pg/socket -h 127.0.0.1" start + for i in $(seq 1 20); do pg_isready -h 127.0.0.1 -p 5432 -U postgres && break; sleep 0.5; done + pg_isready -h 127.0.0.1 -p 5432 -U postgres + createdb -h 127.0.0.1 -p 5432 -U postgres mirror + - name: run linter environment: CGO_ENABLED: 1 @@ -23,5 +39,6 @@ steps: - name: run all tests environment: CGO_ENABLED: 1 + TEST_POSTGRES_URL: postgresql://postgres@127.0.0.1:5432/mirror?sslmode=disable command: | go test -v ./... diff --git a/go.mod b/go.mod index 1653f3b6..72214ff6 100644 --- a/go.mod +++ b/go.mod @@ -63,7 +63,8 @@ require ( ) require ( - dario.cat/mergo v1.0.1 // indirect + dario.cat/mergo v1.0.2 // indirect + github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c // indirect github.com/BurntSushi/toml v0.3.1 // indirect github.com/Microsoft/go-winio v0.6.2 // indirect github.com/ProtonMail/go-crypto v1.3.0 // indirect @@ -120,13 +121,16 @@ require ( github.com/containerd/errdefs v1.0.0 // indirect github.com/containerd/errdefs/pkg v0.3.0 // indirect github.com/containerd/log v0.1.0 // indirect + github.com/containerd/platforms v0.2.1 // indirect + github.com/cpuguy83/dockercfg v0.3.2 // indirect github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect github.com/distribution/reference v0.6.0 // indirect github.com/dlclark/regexp2 v1.11.5 // indirect - github.com/docker/go-connections v0.5.0 // indirect + github.com/docker/go-connections v0.6.0 // indirect github.com/docker/go-units v0.5.0 // indirect github.com/earthboundkid/versioninfo/v2 v2.24.1 // indirect + github.com/ebitengine/purego v0.10.0 // indirect github.com/emirpasic/gods v1.18.1 // indirect github.com/felixge/httpsnoop v1.0.4 // indirect github.com/fsnotify/fsnotify v1.6.0 // indirect @@ -137,6 +141,7 @@ require ( github.com/go-logfmt/logfmt v0.6.0 // indirect github.com/go-logr/logr v1.4.3 // indirect github.com/go-logr/stdr v1.2.2 // indirect + github.com/go-ole/go-ole v1.2.6 // indirect github.com/go-redis/cache/v9 v9.0.0 // indirect github.com/go-test/deep v1.1.1 // indirect github.com/goccy/go-json v0.10.5 // indirect @@ -176,15 +181,24 @@ require ( 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 + github.com/klauspost/compress v1.18.5 // indirect github.com/klauspost/cpuid/v2 v2.3.0 // indirect github.com/lucasb-eyer/go-colorful v1.2.0 // indirect + github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0 // indirect + github.com/magiconair/properties v1.8.10 // indirect github.com/mattn/go-isatty v0.0.20 // indirect github.com/mattn/go-runewidth v0.0.16 // indirect github.com/minio/sha256-simd v1.0.1 // indirect github.com/mitchellh/mapstructure v1.5.0 // indirect github.com/moby/docker-image-spec v1.3.1 // indirect + github.com/moby/go-archive v0.2.0 // indirect + github.com/moby/moby/api v1.54.1 // indirect + github.com/moby/moby/client v0.4.0 // indirect + github.com/moby/patternmatcher v0.6.1 // indirect github.com/moby/sys/atomicwriter v0.1.0 // indirect + github.com/moby/sys/sequential v0.6.0 // indirect + github.com/moby/sys/user v0.4.0 // indirect + github.com/moby/sys/userns v0.1.0 // indirect github.com/moby/term v0.5.2 // indirect github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect github.com/modern-go/reflect2 v1.0.2 // indirect @@ -206,30 +220,38 @@ require ( 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/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 // indirect github.com/prometheus/client_model v0.6.2 // indirect github.com/prometheus/common v0.67.5 // indirect github.com/prometheus/procfs v0.19.2 // indirect github.com/rivo/uniseg v0.4.7 // indirect github.com/ryanuber/go-glob v1.0.0 // indirect github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3 // indirect + github.com/shirou/gopsutil/v4 v4.26.3 // indirect + github.com/sirupsen/logrus v1.9.4 // indirect github.com/spaolacci/murmur3 v1.1.0 // indirect + github.com/testcontainers/testcontainers-go v0.42.0 // indirect + github.com/testcontainers/testcontainers-go/modules/postgres v0.42.0 // indirect github.com/tidwall/gjson v1.18.0 // indirect github.com/tidwall/match v1.2.0 // indirect github.com/tidwall/pretty v1.2.1 // indirect github.com/tidwall/sjson v1.2.5 // indirect + github.com/tklauser/go-sysconf v0.3.16 // indirect + github.com/tklauser/numcpus v0.11.0 // indirect github.com/vmihailenco/go-tinylfu v0.2.2 // indirect github.com/vmihailenco/msgpack/v5 v5.4.1 // indirect github.com/vmihailenco/tagparser/v2 v2.0.0 // indirect github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e // indirect + github.com/yusufpapurcu/wmi v1.2.4 // indirect gitlab.com/yawning/secp256k1-voi v0.0.0-20230925100816-f2616030848b // indirect gitlab.com/yawning/tuplehash v0.0.0-20230713102510-df83abbf9a02 // indirect go.etcd.io/bbolt v1.4.0 // indirect go.opentelemetry.io/auto/sdk v1.2.1 // indirect go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.65.0 // indirect - go.opentelemetry.io/otel v1.40.0 // indirect + go.opentelemetry.io/otel v1.41.0 // indirect go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.40.0 // indirect - go.opentelemetry.io/otel/metric v1.40.0 // indirect - go.opentelemetry.io/otel/trace v1.40.0 // indirect + go.opentelemetry.io/otel/metric v1.41.0 // indirect + go.opentelemetry.io/otel/trace v1.41.0 // indirect go.opentelemetry.io/proto/otlp v1.9.0 // indirect go.uber.org/atomic v1.11.0 // indirect go.uber.org/multierr v1.11.0 // indirect @@ -237,7 +259,7 @@ require ( go.yaml.in/yaml/v2 v2.4.3 // indirect golang.org/x/exp v0.0.0-20260112195511-716be5621a96 // indirect golang.org/x/sync v0.19.0 // indirect - golang.org/x/sys v0.41.0 // indirect + golang.org/x/sys v0.42.0 // indirect golang.org/x/text v0.34.0 // indirect golang.org/x/time v0.12.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260209200024-4cfbd4190f57 // indirect diff --git a/go.sum b/go.sum index aa038d8e..6209caba 100644 --- a/go.sum +++ b/go.sum @@ -1,5 +1,7 @@ dario.cat/mergo v1.0.1 h1:Ra4+bf83h2ztPIQYNP99R6m+Y7KfnARDfID+a+vLl4s= dario.cat/mergo v1.0.1/go.mod h1:uNxQE+84aUszobStD9th8a29P2fMDhsBdgRYvZOxGmk= +dario.cat/mergo v1.0.2 h1:85+piFYR1tMbRrLcDwR18y4UKJ3aH1Tbzi24VRW1TK8= +dario.cat/mergo v1.0.2/go.mod h1:E/hbnu0NxMFBjpMIE34DRGLWqDy0g5FuKDhCb31ngxA= github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c h1:udKWzYgxTojEKWjV8V+WSxDXJ4NFATAsZjh8iIbsQIg= github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c/go.mod h1:xomTg63KZ2rFqZQzSB4Vz2SUXa1BpHTVz9L5PTmPC4E= github.com/Blank-Xu/sql-adapter v1.1.1 h1:+g7QXU9sl/qT6Po97teMpf3GjAO0X9aFaqgSePXvYko= @@ -191,6 +193,10 @@ github.com/containerd/errdefs/pkg v0.3.0 h1:9IKJ06FvyNlexW690DXuQNx2KA2cUJXx151X github.com/containerd/errdefs/pkg v0.3.0/go.mod h1:NJw6s9HwNuRhnjJhM7pylWwMyAkmCQvQ4GpJHEqRLVk= github.com/containerd/log v0.1.0 h1:TCJt7ioM2cr/tfR8GPbGf9/VRAX8D2B4PjzCpfX540I= github.com/containerd/log v0.1.0/go.mod h1:VRRf09a7mHDIRezVKTRCrOq78v577GXq3bSa3EhrzVo= +github.com/containerd/platforms v0.2.1 h1:zvwtM3rz2YHPQsF2CHYM8+KtB5dvhISiXh5ZpSBQv6A= +github.com/containerd/platforms v0.2.1/go.mod h1:XHCb+2/hzowdiut9rkudds9bE5yJ7npe7dG/wG+uFPw= +github.com/cpuguy83/dockercfg v0.3.2 h1:DlJTyZGBDlXqUZ2Dk2Q3xHs/FtnooJJVaad2S9GKorA= +github.com/cpuguy83/dockercfg v0.3.2/go.mod h1:sugsbF4//dDlL/i+S+rtpIWp+5h0BHJHfjj5/jFyUJc= github.com/cpuguy83/go-md2man/v2 v2.0.0-20190314233015-f79a8a8ca69d/go.mod h1:maD7wRr/U5Z6m/iR4s+kqSMx2CaBsrgA7czyZG/E6dU= github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= github.com/cyphar/filepath-securejoin v0.4.1 h1:JyxxyPEaktOD+GAnqIqTf9A8tHyAG22rowi7HkoSU1s= @@ -218,12 +224,16 @@ github.com/docker/docker v28.2.2+incompatible h1:CjwRSksz8Yo4+RmQ339Dp/D2tGO5Jxw github.com/docker/docker v28.2.2+incompatible/go.mod h1:eEKB0N0r5NX/I1kEveEz05bcu8tLC/8azJZsviup8Sk= github.com/docker/go-connections v0.5.0 h1:USnMq7hx7gwdVZq1L49hLXaFtUdTADjXGp+uj1Br63c= github.com/docker/go-connections v0.5.0/go.mod h1:ov60Kzw0kKElRwhNs9UlUHAE/F9Fe6GLaXnqyDdmEXc= +github.com/docker/go-connections v0.6.0 h1:LlMG9azAe1TqfR7sO+NJttz1gy6KO7VJBh+pMmjSD94= +github.com/docker/go-connections v0.6.0/go.mod h1:AahvXYshr6JgfUJGdDCs2b5EZG/vmaMAntpSFH5BFKE= github.com/docker/go-units v0.5.0 h1:69rxXcBk27SvSaaxTtLh/8llcHD8vYHT7WSdRZ/jvr4= github.com/docker/go-units v0.5.0/go.mod h1:fgPhTUdO+D/Jk86RDLlptpiXQzgHJF7gydDDbaIK4Dk= github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= github.com/earthboundkid/versioninfo/v2 v2.24.1 h1:SJTMHaoUx3GzjjnUO1QzP3ZXK6Ee/nbWyCm58eY3oUg= github.com/earthboundkid/versioninfo/v2 v2.24.1/go.mod h1:VcWEooDEuyUJnMfbdTh0uFN4cfEIg+kHMuWB2CDCLjw= +github.com/ebitengine/purego v0.10.0 h1:QIw4xfpWT6GWTzaW5XEKy3HXoqrJGx1ijYHzTF0/ISU= +github.com/ebitengine/purego v0.10.0/go.mod h1:iIjxzd6CiRiOG0UyXP+V1+jWqUXVjPKLAI0mRfJZTmQ= github.com/elazarl/goproxy v1.7.2 h1:Y2o6urb7Eule09PjlhQRGNsqRfPmYI3KKQLFpCAV3+o= github.com/elazarl/goproxy v1.7.2/go.mod h1:82vkLNir0ALaW14Rc399OTTjyNREgmdL2cVoIbS6XaE= github.com/emirpasic/gods v1.18.1 h1:FXtiHYKDGKCW2KzwZKx0iC0PQmdlorYgdFG9jPXJ1Bc= @@ -264,6 +274,8 @@ github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= +github.com/go-ole/go-ole v1.2.6 h1:/Fpf6oFPoeFik9ty7siob0G6Ke8QvQEuVcuChpwXzpY= +github.com/go-ole/go-ole v1.2.6/go.mod h1:pprOEPIfldk/42T2oK7lQ4v4JSDwmV0As9GaiUsvbm0= github.com/go-redis/cache/v9 v9.0.0 h1:0thdtFo0xJi0/WXbRVu8B066z8OvVymXTJGaXrVWnN0= github.com/go-redis/cache/v9 v9.0.0/go.mod h1:cMwi1N8ASBOufbIvk7cdXe2PbPjK/WMRL95FFHWsSgI= github.com/go-task/slim-sprig v0.0.0-20210107165309-348f09dbbbc0/go.mod h1:fyg7847qk6SyHyPtNmDHnmrv/HOrqktSC+C9fM+CJOE= @@ -307,6 +319,7 @@ github.com/google/go-cmp v0.3.1/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMyw github.com/google/go-cmp v0.4.0/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.5.2/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.5.5/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= +github.com/google/go-cmp v0.5.6/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.5.8/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/google/go-cmp v0.5.9/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= @@ -413,6 +426,8 @@ github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+o github.com/klauspost/compress v1.13.6/go.mod h1:/3/Vjq9QcHkK5uEr5lBEmyoZ1iFhe47etQ6QUkpK6sk= github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= +github.com/klauspost/compress v1.18.5 h1:/h1gH5Ce+VWNLSWqPzOVn6XBO+vJbCNGvjoaGBFW2IE= +github.com/klauspost/compress v1.18.5/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= github.com/klauspost/cpuid/v2 v2.3.0 h1:S4CRMLnYUhGeDFDqkGriYKdfoFlDnMtqTiI/sFzhA9Y= github.com/klauspost/cpuid/v2 v2.3.0/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= @@ -427,6 +442,10 @@ github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0 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/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0 h1:6E+4a0GO5zZEnZ81pIr0yLvtUWk2if982qA3F3QD6H4= +github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0/go.mod h1:zJYVVT2jmtg6P3p1VtQj7WsuWi/y4VnjVBn7F8KPB3I= +github.com/magiconair/properties v1.8.10 h1:s31yESBquKXCV9a/ScB3ESkOjUYYv+X0rg8SYxI99mE= +github.com/magiconair/properties v1.8.10/go.mod h1:Dhd985XPs7jluiymwWYZ0G4Z61jb3vdS329zhj2hYo0= github.com/mattn/go-colorable v0.1.14 h1:9A9LHSqF/7dyVVX6g0U9cwm9pG3kP9gSzcuIPHPsaIE= github.com/mattn/go-colorable v0.1.14/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8= github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= @@ -443,10 +462,22 @@ github.com/mitchellh/mapstructure v1.5.0 h1:jeMsZIYE/09sWLaz43PL7Gy6RuMjD2eJVyua github.com/mitchellh/mapstructure v1.5.0/go.mod h1:bFUtVrKA4DC2yAKiSyO/QUcy7e+RRV2QTWOzhPopBRo= github.com/moby/docker-image-spec v1.3.1 h1:jMKff3w6PgbfSa69GfNg+zN/XLhfXJGnEx3Nl2EsFP0= github.com/moby/docker-image-spec v1.3.1/go.mod h1:eKmb5VW8vQEh/BAr2yvVNvuiJuY6UIocYsFu/DxxRpo= +github.com/moby/go-archive v0.2.0 h1:zg5QDUM2mi0JIM9fdQZWC7U8+2ZfixfTYoHL7rWUcP8= +github.com/moby/go-archive v0.2.0/go.mod h1:mNeivT14o8xU+5q1YnNrkQVpK+dnNe/K6fHqnTg4qPU= +github.com/moby/moby/api v1.54.1 h1:TqVzuJkOLsgLDDwNLmYqACUuTehOHRGKiPhvH8V3Nn4= +github.com/moby/moby/api v1.54.1/go.mod h1:+RQ6wluLwtYaTd1WnPLykIDPekkuyD/ROWQClE83pzs= +github.com/moby/moby/client v0.4.0 h1:S+2XegzHQrrvTCvF6s5HFzcrywWQmuVnhOXe2kiWjIw= +github.com/moby/moby/client v0.4.0/go.mod h1:QWPbvWchQbxBNdaLSpoKpCdf5E+WxFAgNHogCWDoa7g= +github.com/moby/patternmatcher v0.6.1 h1:qlhtafmr6kgMIJjKJMDmMWq7WLkKIo23hsrpR3x084U= +github.com/moby/patternmatcher v0.6.1/go.mod h1:hDPoyOpDY7OrrMDLaYoY3hf52gNCR/YOUYxkhApJIxc= github.com/moby/sys/atomicwriter v0.1.0 h1:kw5D/EqkBwsBFi0ss9v1VG3wIkVhzGvLklJ+w3A14Sw= github.com/moby/sys/atomicwriter v0.1.0/go.mod h1:Ul8oqv2ZMNHOceF643P6FKPXeCmYtlQMvpizfsSoaWs= github.com/moby/sys/sequential v0.6.0 h1:qrx7XFUd/5DxtqcoH1h438hF5TmOvzC/lspjy7zgvCU= github.com/moby/sys/sequential v0.6.0/go.mod h1:uyv8EUTrca5PnDsdMGXhZe6CCe8U/UiTWd+lL+7b/Ko= +github.com/moby/sys/user v0.4.0 h1:jhcMKit7SA80hivmFJcbB1vqmw//wU61Zdui2eQXuMs= +github.com/moby/sys/user v0.4.0/go.mod h1:bG+tYYYJgaMtRKgEmuueC0hJEAZWwtIbZTB+85uoHjs= +github.com/moby/sys/userns v0.1.0 h1:tVLXkFOxVu9A64/yh59slHVv9ahO9UIev4JZusOLG/g= +github.com/moby/sys/userns v0.1.0/go.mod h1:IHUYgu/kao6N8YZlp9Cf444ySSvCmDlmzUcYfDHOl28= github.com/moby/term v0.5.2 h1:6qk3FJAFDs6i/q3W/pQ97SX192qKfZgGjCQqfCJkgzQ= github.com/moby/term v0.5.2/go.mod h1:d3djjFCrjnB+fl8NJux+EJzu0msscUP+f8it8hPkFLc= github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= @@ -526,6 +557,8 @@ github.com/polydawn/refmt v0.89.1-0.20221221234430-40501e09de1f h1:VXTQfuJj9vKR4 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/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 h1:o4JXh1EVt9k/+g42oCprj/FisM4qX9L3sZB3upGN2ZU= +github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55/go.mod h1:OmDBASR4679mdNQnz2pUhc2G8CO2JrUAVFDRBDP/hJE= 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= @@ -553,9 +586,13 @@ github.com/sergi/go-diff v1.1.0 h1:we8PVUC3FE2uYfodKH/nBHMSetSfHDR6scGdBi+erh0= github.com/sergi/go-diff v1.1.0/go.mod h1:STckp+ISIX8hZLjrqAeVduY0gWCT9IjLuqbuNXdaHfM= github.com/sethvargo/go-envconfig v1.1.0 h1:cWZiJxeTm7AlCvzGXrEXaSTCNgip5oJepekh/BOQuog= github.com/sethvargo/go-envconfig v1.1.0/go.mod h1:JLd0KFWQYzyENqnEPWWZ49i4vzZo/6nRidxI8YvGiHw= +github.com/shirou/gopsutil/v4 v4.26.3 h1:2ESdQt90yU3oXF/CdOlRCJxrP+Am1aBYubTMTfxJ1qc= +github.com/shirou/gopsutil/v4 v4.26.3/go.mod h1:LZ6ewCSkBqUpvSOf+LsTGnRinC6iaNUNMGBtDkJBaLQ= github.com/shurcooL/sanitized_anchor_name v1.0.0/go.mod h1:1NzhyTcUVG4SuEtjjoZeVRXNmyL/1OwPU0+IJeTBvfc= github.com/sirupsen/logrus v1.9.3 h1:dueUQJ1C2q9oE3F7wvmSGAaVtTmUizReu6fjN8uqzbQ= github.com/sirupsen/logrus v1.9.3/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ= +github.com/sirupsen/logrus v1.9.4 h1:TsZE7l11zFCLZnZ+teH4Umoq5BhEIfIzfRDZ1Uzql2w= +github.com/sirupsen/logrus v1.9.4/go.mod h1:ftWc9WdOfJ0a92nsE2jF5u5ZwH8Bv2zdeOC42RjbV2g= github.com/smartystreets/assertions v1.2.0 h1:42S6lae5dvLc7BrLu/0ugRtcFVjoJNMC/N3yZFZkDFs= github.com/smartystreets/assertions v1.2.0/go.mod h1:tcbTF8ujkAEcZ8TElKY+i30BzYlVhC/LOxJk7iOWnoo= github.com/smartystreets/goconvey v1.7.2 h1:9RBaZCeXEQ3UselpuwUQHltGVXvdwm6cv1hgR6gDIPg= @@ -579,6 +616,10 @@ github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= 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/testcontainers/testcontainers-go v0.42.0 h1:He3IhTzTZOygSXLJPMX7n44XtK+qhjat1nI9cneBbUY= +github.com/testcontainers/testcontainers-go v0.42.0/go.mod h1:vZjdY1YmUA1qEForxOIOazfsrdyORJAbhi0bp8plN30= +github.com/testcontainers/testcontainers-go/modules/postgres v0.42.0 h1:GCbb1ndrF7OTDiIvxXyItaDab4qkzTFJ48LKFdM7EIo= +github.com/testcontainers/testcontainers-go/modules/postgres v0.42.0/go.mod h1:IRPBaI8jXdrNfD0e4Zm7Fbcgaz5shKxOQv4axiL09xs= 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= @@ -590,6 +631,10 @@ github.com/tidwall/pretty v1.2.1 h1:qjsOFOWWQl+N3RsoF5/ssm1pHmJJwhjlSbZ51I6wMl4= github.com/tidwall/pretty v1.2.1/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU= github.com/tidwall/sjson v1.2.5 h1:kLy8mja+1c9jlljvWTlSazM7cKDRfJuR/bOJhcY5NcY= github.com/tidwall/sjson v1.2.5/go.mod h1:Fvgq9kS/6ociJEDnK0Fk1cpYF4FIW6ZF7LAe+6jwd28= +github.com/tklauser/go-sysconf v0.3.16 h1:frioLaCQSsF5Cy1jgRBrzr6t502KIIwQ0MArYICU0nA= +github.com/tklauser/go-sysconf v0.3.16/go.mod h1:/qNL9xxDhc7tx3HSRsLWNnuzbVfh3e7gh/BmM179nYI= +github.com/tklauser/numcpus v0.11.0 h1:nSTwhKH5e1dMNsCdVBukSZrURJRoHbSEQjdEbY+9RXw= +github.com/tklauser/numcpus v0.11.0/go.mod h1:z+LwcLq54uWZTX0u/bGobaV34u6V7KNlTZejzM6/3MQ= github.com/urfave/cli v1.22.10/go.mod h1:Gos4lmkARVdJ6EkW0WaNv/tZAAMe9V7XWyB60NtXRu0= github.com/urfave/cli/v3 v3.6.2 h1:lQuqiPrZ1cIz8hz+HcrG0TNZFxU70dPZ3Yl+pSrH9A8= github.com/urfave/cli/v3 v3.6.2/go.mod h1:ysVLtOEmg2tOy6PknnYVhDoouyC/6N42TMeoMzskhso= @@ -618,6 +663,8 @@ github.com/yuin/goldmark-emoji v1.0.6 h1:QWfF2FYaXwL74tfGOW5izeiZepUDroDJfWubQI9 github.com/yuin/goldmark-emoji v1.0.6/go.mod h1:ukxJDKFpdFb5x0a5HqbdlcKtebh086iJpI31LTKmWuA= github.com/yuin/goldmark-highlighting/v2 v2.0.0-20230729083705-37449abec8cc h1:+IAOyRda+RLrxa1WC7umKOZRsGq4QrFFMYApOeHzQwQ= github.com/yuin/goldmark-highlighting/v2 v2.0.0-20230729083705-37449abec8cc/go.mod h1:ovIvrum6DQJA4QsJSovrkC4saKHQVs7TvcaeO8AIl5I= +github.com/yusufpapurcu/wmi v1.2.4 h1:zFUKzehAFReQwLys1b/iSMl+JQGSCSjtVqQn9bBrPo0= +github.com/yusufpapurcu/wmi v1.2.4/go.mod h1:SBZ9tNy3G9/m5Oi98Zks0QjeHVDvuK0qfxQmPyzfmi0= gitlab.com/staticnoise/goldmark-callout v0.0.0-20240609120641-6366b799e4ab h1:gK9tS6QJw5F0SIhYJnGG2P83kuabOdmWBbSmZhJkz2A= gitlab.com/staticnoise/goldmark-callout v0.0.0-20240609120641-6366b799e4ab/go.mod h1:SPu13/NPe1kMrbGoJldQwqtpNhXsmIuHCfm/aaGjU0c= gitlab.com/yawning/secp256k1-voi v0.0.0-20230925100816-f2616030848b h1:CzigHMRySiX3drau9C6Q5CAbNIApmLdat5jPMqChvDA= @@ -634,18 +681,24 @@ go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.65.0 h1:7iP2uCb go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.65.0/go.mod h1:c7hN3ddxs/z6q9xwvfLPk+UHlWRQyaeR1LdgfL/66l0= go.opentelemetry.io/otel v1.40.0 h1:oA5YeOcpRTXq6NN7frwmwFR0Cn3RhTVZvXsP4duvCms= go.opentelemetry.io/otel v1.40.0/go.mod h1:IMb+uXZUKkMXdPddhwAHm6UfOwJyh4ct1ybIlV14J0g= +go.opentelemetry.io/otel v1.41.0 h1:YlEwVsGAlCvczDILpUXpIpPSL/VPugt7zHThEMLce1c= +go.opentelemetry.io/otel v1.41.0/go.mod h1:Yt4UwgEKeT05QbLwbyHXEwhnjxNO6D8L5PQP51/46dE= go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.40.0 h1:QKdN8ly8zEMrByybbQgv8cWBcdAarwmIPZ6FThrWXJs= go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.40.0/go.mod h1:bTdK1nhqF76qiPoCCdyFIV+N/sRHYXYCTQc+3VCi3MI= go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.38.0 h1:aTL7F04bJHUlztTsNGJ2l+6he8c+y/b//eR0jjjemT4= go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.38.0/go.mod h1:kldtb7jDTeol0l3ewcmd8SDvx3EmIE7lyvqbasU3QC4= go.opentelemetry.io/otel/metric v1.40.0 h1:rcZe317KPftE2rstWIBitCdVp89A2HqjkxR3c11+p9g= go.opentelemetry.io/otel/metric v1.40.0/go.mod h1:ib/crwQH7N3r5kfiBZQbwrTge743UDc7DTFVZrrXnqc= +go.opentelemetry.io/otel/metric v1.41.0 h1:rFnDcs4gRzBcsO9tS8LCpgR0dxg4aaxWlJxCno7JlTQ= +go.opentelemetry.io/otel/metric v1.41.0/go.mod h1:xPvCwd9pU0VN8tPZYzDZV/BMj9CM9vs00GuBjeKhJps= go.opentelemetry.io/otel/sdk v1.40.0 h1:KHW/jUzgo6wsPh9At46+h4upjtccTmuZCFAc9OJ71f8= go.opentelemetry.io/otel/sdk v1.40.0/go.mod h1:Ph7EFdYvxq72Y8Li9q8KebuYUr2KoeyHx0DRMKrYBUE= go.opentelemetry.io/otel/sdk/metric v1.40.0 h1:mtmdVqgQkeRxHgRv4qhyJduP3fYJRMX4AtAlbuWdCYw= go.opentelemetry.io/otel/sdk/metric v1.40.0/go.mod h1:4Z2bGMf0KSK3uRjlczMOeMhKU2rhUqdWNoKcYrtcBPg= go.opentelemetry.io/otel/trace v1.40.0 h1:WA4etStDttCSYuhwvEa8OP8I5EWu24lkOzp+ZYblVjw= go.opentelemetry.io/otel/trace v1.40.0/go.mod h1:zeAhriXecNGP/s2SEG3+Y8X9ujcJOTqQ5RgdEJcawiA= +go.opentelemetry.io/otel/trace v1.41.0 h1:Vbk2co6bhj8L59ZJ6/xFTskY+tGAbOnCtQGVVa9TIN0= +go.opentelemetry.io/otel/trace v1.41.0/go.mod h1:U1NU4ULCoxeDKc09yCWdWe+3QoyweJcISEVa1RBzOis= go.opentelemetry.io/proto/otlp v1.9.0 h1:l706jCMITVouPOqEnii2fIAuO3IVGBRPV5ICjceRb/A= go.opentelemetry.io/proto/otlp v1.9.0/go.mod h1:xE+Cx5E/eEHw+ISFkwPLwCZefwVjY+pqKg1qcK03+/4= go.uber.org/atomic v1.6.0/go.mod h1:sABNBOSYdrvTF6hTgEIbc7YasKWGhgEQZyfxyTvoXHQ= @@ -722,17 +775,20 @@ golang.org/x/sys v0.0.0-20180909124046-d0be0721c37e/go.mod h1:STP8DvDyc/dI5b8T5h golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20190904154756-749cb33beabd/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20191005200804-aed5e4c7ecf9/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20191120155948-bd437916bb0e/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20191204072324-ce4227a45e2e/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20200323222414-85ca7c5b95cd/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20200930185726-fdedc70b468f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20201204225414-ed752295db88/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210112080510-489259a85091/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210330210617-4fbd30eecc44/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210510120138-977fb7262007/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20210616094352-59db8d763f22/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20211019181941-9d821ace8654/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20211216021012-1d35b9e2eb4e/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220319134239-a9b59b0215f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= @@ -750,6 +806,8 @@ 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.41.0 h1:Ivj+2Cp/ylzLiEU89QhWblYnOE9zerudt9Ftecq2C6k= golang.org/x/sys v0.41.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= +golang.org/x/sys v0.42.0 h1:omrd2nAlyT5ESRdCLYdm3+fMfNFE/+Rf4bDIQImRJeo= +golang.org/x/sys v0.42.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= 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= diff --git a/knotmirror/db/repos.go b/knotmirror/db/repos.go index 36313b1e..c137e611 100644 --- a/knotmirror/db/repos.go +++ b/knotmirror/db/repos.go @@ -67,6 +67,43 @@ func UpdateRepoState(ctx context.Context, e *sql.DB, did syntax.DID, rkey syntax return nil } +func MarkDesynchronized(ctx context.Context, e *sql.DB, did syntax.DID, rkey syntax.RecordKey) error { + if _, err := e.ExecContext(ctx, + `update repos + set state = $1 + where did = $2 and rkey = $3 and state in ($4, $5, $6, $7)`, + models.RepoStateDesynchronized, + did, rkey, + models.RepoStateActive, + models.RepoStateDesynchronized, + models.RepoStateResyncing, + models.RepoStateError, + ); err != nil { + return fmt.Errorf("marking repo desynchronized: %w", err) + } + return nil +} + +func FinishResync(ctx context.Context, e *sql.DB, did syntax.DID, rkey syntax.RecordKey, state models.RepoState, errMsg string, retryCount int, retryAfter int64) error { + if _, err := e.ExecContext(ctx, + `update repos + set error_msg = $1, + retry_count = $2, + retry_after = $3, + state = case when state = $4 then $5 else state end + where did = $6 and rkey = $7`, + errMsg, + retryCount, + retryAfter, + models.RepoStateResyncing, + state, + did, rkey, + ); err != nil { + return fmt.Errorf("finishing resync: %w", err) + } + return nil +} + func DeleteRepo(ctx context.Context, e *sql.DB, did syntax.DID, rkey syntax.RecordKey) error { if _, err := e.ExecContext(ctx, `delete from repos where did = $1 and rkey = $2`, diff --git a/knotmirror/knotstream/slurper.go b/knotmirror/knotstream/slurper.go index b28ca1dc..98405fc7 100644 --- a/knotmirror/knotstream/slurper.go +++ b/knotmirror/knotstream/slurper.go @@ -312,9 +312,8 @@ func (s *KnotSlurper) ProcessLegacyGitRefUpdate(ctx context.Context, source stri } l = l.With("repoAt", curr.AtUri()) - // TODO: should plan resync to resyncBuffer on RepoStateResyncing - if curr.State != models.RepoStateActive { - l.Debug("skipping non-active repo") + if curr.State == models.RepoStateSuspended || curr.State == models.RepoStatePending { + l.Debug("skipping repo", "state", curr.State) knotstreamEventsSkipped.Inc() return nil } @@ -325,13 +324,8 @@ func (s *KnotSlurper) ProcessLegacyGitRefUpdate(ctx context.Context, source stri 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 { + if err := db.MarkDesynchronized(ctx, s.db, curr.Did, curr.Rkey); err != nil { return err } diff --git a/knotmirror/resyncer.go b/knotmirror/resyncer.go index 4357e8a0..67d8d5d8 100644 --- a/knotmirror/resyncer.go +++ b/knotmirror/resyncer.go @@ -57,7 +57,7 @@ func NewResyncer(l *slog.Logger, db *sql.DB, gitm GitMirrorManager, cfg *config. } func (r *Resyncer) Start(ctx context.Context) { - for i := 0; i < r.parallelism; i++ { + for i := range r.parallelism { go r.runResyncWorker(ctx, i) } } @@ -71,7 +71,7 @@ func (r *Resyncer) runResyncWorker(ctx context.Context, workerID int) { return default: } - repoAt, found, err := r.claimResyncJob(ctx) + repoAt, jobCtx, found, err := r.claimResyncJob(ctx) if err != nil { l.Error("failed to claim resync job", "error", err) time.Sleep(time.Second) @@ -82,43 +82,49 @@ func (r *Resyncer) runResyncWorker(ctx context.Context, workerID int) { continue } l.Info("processing resync", "aturi", repoAt) - if err := r.resyncRepo(ctx, repoAt); err != nil { + if err := r.resyncRepo(ctx, jobCtx, repoAt); err != nil { l.Error("resync failed", "aturi", repoAt, "error", err) } } } -func (r *Resyncer) registerRunning(repo syntax.ATURI, cancel context.CancelFunc) { +func (r *Resyncer) finalizeJob(repo syntax.ATURI) { r.runningJobsMu.Lock() defer r.runningJobsMu.Unlock() - if _, exists := r.runningJobs[repo]; exists { - return + if cancel, ok := r.runningJobs[repo]; ok { + cancel() + delete(r.runningJobs, repo) } - 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() + r.finalizeJob(repo) } // TriggerResyncJob manually triggers the resync job func (r *Resyncer) TriggerResyncJob(ctx context.Context, repoAt syntax.ATURI) error { + res, err := r.db.ExecContext(ctx, + `update repos + set state = $1, retry_after = $2 + where at_uri = $3 and state not in ($4, $5)`, + models.RepoStatePending, + int64(-1), + repoAt, + models.RepoStateResyncing, + models.RepoStateSuspended, + ) + if err != nil { + return fmt.Errorf("triggering resync: %w", err) + } + n, err := res.RowsAffected() + if err != nil { + return fmt.Errorf("triggering resync: %w", err) + } + if n > 0 { + return nil + } + repo, err := db.GetRepoByAtUri(ctx, r.db, repoAt) if err != nil { return fmt.Errorf("failed to get repo: %w", err) @@ -126,56 +132,68 @@ func (r *Resyncer) TriggerResyncJob(ctx context.Context, repoAt syntax.ATURI) er 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 + return fmt.Errorf("cannot trigger resync: repo in state %s", repo.State) } -func (r *Resyncer) claimResyncJob(ctx context.Context) (syntax.ATURI, bool, error) { +func (r *Resyncer) claimResyncJob(ctx context.Context) (syntax.ATURI, context.Context, 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 + r.runningJobsMu.Lock() + excludes := make([]any, 0, len(r.runningJobs)) + for aturi := range r.runningJobs { + excludes = append(excludes, string(aturi)) + } + r.runningJobsMu.Unlock() + + args := []any{ + models.RepoStateResyncing, + models.RepoStatePending, models.RepoStateDesynchronized, models.RepoStateError, + time.Now().Unix(), + } + excludeClause := "" + if len(excludes) > 0 { + base := len(args) + 1 + placeholders := make([]string, len(excludes)) + for i := range excludes { + placeholders[i] = fmt.Sprintf("$%d", base+i) + } + excludeClause = " and at_uri not in (" + strings.Join(placeholders, ",") + ")" + args = append(args, excludes...) + } + + query := `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) + and (retry_after = -1 or retry_after = 0 or retry_after < $5)` + excludeClause + ` 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 { + returning at_uri` + + var repoAt syntax.ATURI + if err := r.db.QueryRowContext(ctx, query, args...).Scan(&repoAt); err != nil { if errors.Is(err, sql.ErrNoRows) { - return "", false, nil + return "", nil, false, nil } - return "", false, err + return "", nil, false, err } - return repoAt, true, nil + jobCtx, cancel := context.WithCancel(ctx) + r.runningJobsMu.Lock() + r.runningJobs[repoAt] = cancel + r.runningJobsMu.Unlock() + + return repoAt, jobCtx, true, nil } -func (r *Resyncer) resyncRepo(ctx context.Context, repoAt syntax.ATURI) error { +func (r *Resyncer) resyncRepo(ctx, jobCtx context.Context, repoAt syntax.ATURI) error { // ctx, span := tracer.Start(ctx, "resyncRepo") // span.SetAttributes(attribute.String("aturi", repoAt)) // defer span.End() @@ -183,11 +201,9 @@ func (r *Resyncer) resyncRepo(ctx context.Context, repoAt syntax.ATURI) error { resyncsStarted.Inc() startTime := time.Now() - jobCtx, cancel := context.WithCancel(ctx) - r.registerRunning(repoAt, cancel) - defer r.unregisterRunning(repoAt) + defer r.finalizeJob(repoAt) - success, err := r.doResync(jobCtx, repoAt) + success, err := r.doResync(ctx, jobCtx, repoAt) if !success { resyncsFailed.Inc() resyncDuration.Observe(time.Since(startTime).Seconds()) @@ -199,12 +215,12 @@ func (r *Resyncer) resyncRepo(ctx context.Context, repoAt syntax.ATURI) error { return nil } -func (r *Resyncer) doResync(ctx context.Context, repoAt syntax.ATURI) (bool, error) { +func (r *Resyncer) doResync(ctx, jobCtx 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) + repo, err := db.GetRepoByAtUri(jobCtx, r.db, repoAt) if err != nil { return false, fmt.Errorf("failed to get repo: %w", err) } @@ -222,7 +238,7 @@ func (r *Resyncer) doResync(ctx context.Context, repoAt syntax.ATURI) (bool, err // HACK: check knot reachability with short timeout before running actual fetch. // This is crucial as git-cli doesn't support http connection timeout. // `http.lowSpeedTime` is only applied _after_ the connection. - if err := r.checkKnotReachability(ctx, repo); err != nil { + if err := r.checkKnotReachability(jobCtx, repo); err != nil { if isRateLimitError(err) { r.knotBackoffMu.Lock() r.knotBackoff[repo.KnotDomain] = time.Now().Add(10 * time.Second) @@ -237,7 +253,7 @@ func (r *Resyncer) doResync(ctx context.Context, repoAt syntax.ATURI) (bool, err if repo.RetryAfter == -1 { timeout = r.manualResyncTimeout } - fetchCtx, cancel := context.WithTimeout(ctx, timeout) + fetchCtx, cancel := context.WithTimeout(jobCtx, timeout) defer cancel() if err := r.gitm.Sync(fetchCtx, repo); err != nil { @@ -246,11 +262,7 @@ func (r *Resyncer) doResync(ctx context.Context, repoAt syntax.ATURI) (bool, 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 { + if err := db.FinishResync(ctx, r.db, repo.Did, repo.Rkey, models.RepoStateActive, "", 0, 0); err != nil { return false, fmt.Errorf("updating repo state to active %w", err) } return true, nil @@ -328,9 +340,9 @@ func (r *Resyncer) handleResyncFailure(ctx context.Context, repoAt syntax.ATURI, errMsg = err.Error() } - repo, err := db.GetRepoByAtUri(ctx, r.db, repoAt) - if err != nil { - return fmt.Errorf("failed to get repo: %w", err) + repo, getErr := db.GetRepoByAtUri(ctx, r.db, repoAt) + if getErr != nil { + return fmt.Errorf("failed to get repo: %w", getErr) } if repo == nil { return fmt.Errorf("failed to get repo. repo '%s' doesn't exist in db", repoAt) @@ -343,11 +355,7 @@ func (r *Resyncer) handleResyncFailure(ctx context.Context, repoAt syntax.ATURI, // 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 { + if err := db.FinishResync(ctx, r.db, repo.Did, repo.Rkey, state, errMsg, retryCount, retryAfter); err != nil { return fmt.Errorf("failed to update repo state: %w", err) } return nil diff --git a/knotmirror/resyncer_test.go b/knotmirror/resyncer_test.go new file mode 100644 index 00000000..58e558e4 --- /dev/null +++ b/knotmirror/resyncer_test.go @@ -0,0 +1,225 @@ +package knotmirror_test + +import ( + "context" + "database/sql" + "io" + "log/slog" + "net/http" + "net/http/httptest" + "os" + "sync" + "testing" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/stretchr/testify/require" + "github.com/testcontainers/testcontainers-go" + "github.com/testcontainers/testcontainers-go/modules/postgres" + "github.com/testcontainers/testcontainers-go/wait" + + "tangled.org/core/api/tangled" + "tangled.org/core/knotmirror" + "tangled.org/core/knotmirror/config" + "tangled.org/core/knotmirror/db" + "tangled.org/core/knotmirror/knotstream" + "tangled.org/core/knotmirror/models" +) + +type blockingGitm struct { + mu sync.Mutex + calls []syntax.ATURI + started chan struct{} + release chan struct{} +} + +func newBlockingGitm() *blockingGitm { + return &blockingGitm{ + started: make(chan struct{}, 16), + release: make(chan struct{}, 16), + } +} + +func (f *blockingGitm) Exist(repo *models.Repo) (bool, error) { return true, nil } +func (f *blockingGitm) RemoteSetUrl(ctx context.Context, repo *models.Repo) error { return nil } +func (f *blockingGitm) Clone(ctx context.Context, repo *models.Repo) error { return nil } +func (f *blockingGitm) Fetch(ctx context.Context, repo *models.Repo) error { return nil } + +func (f *blockingGitm) Sync(ctx context.Context, repo *models.Repo) error { + f.mu.Lock() + f.calls = append(f.calls, repo.AtUri()) + f.mu.Unlock() + select { + case f.started <- struct{}{}: + case <-ctx.Done(): + return ctx.Err() + } + select { + case <-f.release: + return nil + case <-ctx.Done(): + return ctx.Err() + } +} + +func (f *blockingGitm) callCount() int { + f.mu.Lock() + defer f.mu.Unlock() + return len(f.calls) +} + +var _ knotmirror.GitMirrorManager = (*blockingGitm)(nil) + +func startPostgres(t *testing.T, ctx context.Context) *sql.DB { + t.Helper() + + if url := os.Getenv("TEST_POSTGRES_URL"); url != "" { + database, err := db.Make(ctx, url, 4) + require.NoError(t, err) + _, err = database.ExecContext(ctx, `truncate table repos, hosts`) + require.NoError(t, err) + t.Cleanup(func() { + _ = database.Close() + }) + return database + } + + pg, err := postgres.Run(ctx, "docker.io/library/postgres:16-alpine", + postgres.WithDatabase("mirror"), + postgres.WithUsername("tnglr"), + postgres.WithPassword("test"), + testcontainers.WithWaitStrategy( + wait.ForLog("database system is ready to accept connections"). + WithOccurrence(2). + WithStartupTimeout(60*time.Second), + ), + ) + if err != nil { + t.Skipf("postgres container unavailable, is podman/docker running: %v", err) + } + t.Cleanup(func() { + _ = pg.Terminate(context.Background()) + }) + + url, err := pg.ConnectionString(ctx, "sslmode=disable") + require.NoError(t, err) + + database, err := db.Make(ctx, url, 4) + require.NoError(t, err) + t.Cleanup(func() { + _ = database.Close() + }) + return database +} + +func testLogger() *slog.Logger { + return slog.New(slog.NewTextHandler(io.Discard, nil)) +} + +func mockKnot(t *testing.T) *httptest.Server { + t.Helper() + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/x-git-upload-pack-advertisement") + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte("001e# service=git-upload-pack\n")) + })) + t.Cleanup(srv.Close) + return srv +} + +func seedActiveRepo(t *testing.T, ctx context.Context, database *sql.DB, knotURL string) *models.Repo { + t.Helper() + repo := &models.Repo{ + Did: syntax.DID("did:plc:testingtestingtestingtest"), + Rkey: syntax.RecordKey("3kaaaaaaaaaaaa"), + Name: "race-repo", + KnotDomain: knotURL, + State: models.RepoStateActive, + } + require.NoError(t, db.UpsertRepo(ctx, database, repo)) + return repo +} + +func mkEvent(rkey string, oldSha, newSha string, repo *models.Repo) *knotstream.LegacyGitEvent { + owner := repo.Did.String() + repoDid := repo.Did.String() + return &knotstream.LegacyGitEvent{ + Rkey: rkey, + Nsid: tangled.GitRefUpdateNSID, + Event: tangled.GitRefUpdate{ + OwnerDid: &owner, + RepoName: repo.Name, + RepoDid: &repoDid, + OldSha: oldSha, + NewSha: newSha, + Ref: "refs/heads/main", + }, + } +} + +func TestEventDuringResyncTriggersSecondSync(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + database := startPostgres(t, ctx) + knotSrv := mockKnot(t) + repo := seedActiveRepo(t, ctx, database, knotSrv.URL) + + gitm := newBlockingGitm() + logger := testLogger() + cfg := &config.Config{ + KnotUseSSL: false, + GitRepoFetchTimeout: 10 * time.Second, + ResyncParallelism: 1, + } + + slurper := knotstream.NewKnotSlurper(logger, database, cfg.Slurper) + resyncer := knotmirror.NewResyncer(logger, database, gitm, cfg) + resyncer.Start(ctx) + + eventA := mkEvent("3evtaaaaaaaaaa", "0000000000000000000000000000000000000000", "1111111111111111111111111111111111111111", repo) + require.NoError(t, slurper.ProcessLegacyGitRefUpdate(ctx, "testsrc", eventA)) + + got, err := db.GetRepoByAtUri(ctx, database, repo.AtUri()) + require.NoError(t, err) + require.Equal(t, models.RepoStateDesynchronized, got.State, "event A should mark repo desynchronized") + + select { + case <-gitm.started: + case <-time.After(10 * time.Second): + t.Fatal("timeout waiting for first Sync to start") + } + + got, err = db.GetRepoByAtUri(ctx, database, repo.AtUri()) + require.NoError(t, err) + require.Equal(t, models.RepoStateResyncing, got.State, "state should be resyncing while fetch is in flight") + + eventB := mkEvent("3evtbbbbbbbbbb", "1111111111111111111111111111111111111111", "2222222222222222222222222222222222222222", repo) + require.NoError(t, slurper.ProcessLegacyGitRefUpdate(ctx, "testsrc", eventB)) + + gitm.release <- struct{}{} + + secondSyncStarted := false + select { + case <-gitm.started: + secondSyncStarted = true + gitm.release <- struct{}{} + case <-time.After(10 * time.Second): + } + + deadline := time.Now().Add(10 * time.Second) + var finalState models.RepoState + for time.Now().Before(deadline) { + got, err = db.GetRepoByAtUri(ctx, database, repo.AtUri()) + require.NoError(t, err) + finalState = got.State + if finalState == models.RepoStateActive { + break + } + time.Sleep(50 * time.Millisecond) + } + + require.Truef(t, secondSyncStarted, "event B arriving during resync should trigger a second Sync, got %d total Sync calls", gitm.callCount()) + require.GreaterOrEqual(t, gitm.callCount(), 2, "expected at least 2 Sync calls") + require.Equal(t, models.RepoStateActive, finalState, "repo should end in active state") +}