diff --git a/go.mod b/go.mod index 884de7b5..461f6fd3 100644 --- a/go.mod +++ b/go.mod @@ -70,6 +70,7 @@ require ( github.com/Microsoft/go-winio v0.6.2 // indirect github.com/ProtonMail/go-crypto v1.3.0 // indirect github.com/RoaringBitmap/roaring/v2 v2.4.5 // indirect + github.com/RussellLuo/slidingwindow v0.0.0-20200528002341-535bb99d338b // indirect github.com/alecthomas/repr v0.5.2 // indirect github.com/anmitsu/go-shlex v0.0.0-20200514113438-38f4b401e2be // indirect github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.7 // indirect @@ -141,6 +142,7 @@ require ( github.com/go-test/deep v1.1.1 // indirect github.com/goccy/go-json v0.10.5 // indirect github.com/gogo/protobuf v1.3.2 // indirect + github.com/golang-jwt/jwt v3.2.2+incompatible // indirect github.com/golang-jwt/jwt/v5 v5.3.0 // indirect github.com/golang/groupcache v0.0.0-20241129210726-2c02b8208cf8 // indirect github.com/golang/mock v1.6.0 // indirect @@ -163,22 +165,36 @@ require ( github.com/ipfs/bbloom v0.0.4 // indirect github.com/ipfs/boxo v0.36.0 // indirect github.com/ipfs/go-block-format v0.2.3 // indirect + github.com/ipfs/go-blockservice v0.5.2 // indirect github.com/ipfs/go-datastore v0.9.0 // indirect github.com/ipfs/go-ipfs-blockstore v1.3.1 // indirect github.com/ipfs/go-ipfs-ds-help v1.1.1 // indirect + github.com/ipfs/go-ipfs-exchange-interface v0.2.1 // indirect + github.com/ipfs/go-ipfs-util v0.0.3 // indirect github.com/ipfs/go-ipld-cbor v0.2.1 // indirect github.com/ipfs/go-ipld-format v0.6.3 // indirect + github.com/ipfs/go-ipld-legacy v0.2.2 // indirect github.com/ipfs/go-log v1.0.5 // indirect github.com/ipfs/go-log/v2 v2.9.1 // indirect + github.com/ipfs/go-merkledag v0.11.0 // indirect github.com/ipfs/go-metrics-interface v0.3.0 // indirect + github.com/ipfs/go-verifcid v0.0.3 // indirect + github.com/ipld/go-car v0.6.2 // indirect + github.com/ipld/go-codec-dagpb v1.7.0 // indirect + github.com/ipld/go-ipld-prime v0.21.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/jinzhu/inflection v1.0.0 // indirect + github.com/jinzhu/now v1.1.5 // 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/cpuid/v2 v2.3.0 // indirect + github.com/labstack/echo/v4 v4.11.3 // indirect + github.com/labstack/gommon v0.4.1 // indirect github.com/lucasb-eyer/go-colorful v1.2.0 // indirect + github.com/mattn/go-colorable v0.1.14 // 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 @@ -209,6 +225,7 @@ require ( 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/puzpuzpuz/xsync/v4 v4.2.0 // 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 @@ -217,6 +234,8 @@ require ( 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/valyala/bytebufferpool v1.0.0 // indirect + github.com/valyala/fasttemplate v1.2.2 // 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 @@ -242,6 +261,9 @@ require ( gopkg.in/tomb.v1 v1.0.0-20141024135613-dd632973f1e7 // indirect gopkg.in/warnings.v0 v0.1.2 // indirect gopkg.in/yaml.v2 v2.4.0 // indirect + gorm.io/driver/postgres v1.6.0 // indirect + gorm.io/driver/sqlite v1.6.0 // indirect + gorm.io/gorm v1.31.1 // indirect gotest.tools/v3 v3.5.2 // indirect lukechampine.com/blake3 v1.4.1 // indirect ) diff --git a/go.sum b/go.sum index 539fefbc..07e60c45 100644 --- a/go.sum +++ b/go.sum @@ -12,6 +12,8 @@ github.com/ProtonMail/go-crypto v1.3.0 h1:ILq8+Sf5If5DCpHQp4PbZdS1J7HDFRXz/+xKBi github.com/ProtonMail/go-crypto v1.3.0/go.mod h1:9whxjD8Rbs29b4XWbB8irEcE8KHMqaR2e7GWU1R+/PE= github.com/RoaringBitmap/roaring/v2 v2.4.5 h1:uGrrMreGjvAtTBobc0g5IrW1D5ldxDQYe2JW2gggRdg= github.com/RoaringBitmap/roaring/v2 v2.4.5/go.mod h1:FiJcsfkGje/nZBZgCu0ZxCPOKD/hVXDS2dXi7/eUFE0= +github.com/RussellLuo/slidingwindow v0.0.0-20200528002341-535bb99d338b h1:5/++qT1/z812ZqBvqQt6ToRswSuPZ/B33m6xVHRzADU= +github.com/RussellLuo/slidingwindow v0.0.0-20200528002341-535bb99d338b/go.mod h1:4+EPqMRApwwE/6yo6CxiHoSnBzjRr3jsqer7frxP8y4= github.com/adrg/frontmatter v0.2.0 h1:/DgnNe82o03riBd1S+ZDjd43wAmC6W35q67NHeLkPd4= github.com/adrg/frontmatter v0.2.0/go.mod h1:93rQCj3z3ZlwyxxpQioRKC1wDLto4aXHrbqIsnH9wmE= github.com/alecthomas/assert/v2 v2.11.0 h1:2Q9r3ki8+JYXvGsDyBXwH3LcJ+WK5D0gc5E8vS6K3D0= @@ -67,6 +69,8 @@ github.com/aymanbagabas/go-osc52/v2 v2.0.1 h1:HwpRHbFMcZLEVr42D4p7XBqjyuxQH5SMiE github.com/aymanbagabas/go-osc52/v2 v2.0.1/go.mod h1:uYgXzlJ7ZpABp8OJ+exZzJJhRNQ2ASbcXHWsFqH8hp8= github.com/aymerick/douceur v0.2.0 h1:Mv+mAeH1Q+n9Fr+oyamOlAkUNPWPlA8PPGR0QAaYuPk= github.com/aymerick/douceur v0.2.0/go.mod h1:wlT5vV2O3h55X9m7iVYN0TBM0NH/MmbLnd30/FjWUq4= +github.com/benbjohnson/clock v1.3.5 h1:VvXlSJBzZpA/zum6Sj74hxwYI2DIxRWuNIoXAzHZz5o= +github.com/benbjohnson/clock v1.3.5/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA= github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= github.com/bits-and-blooms/bitset v1.12.0/go.mod h1:7hO7Gc7Pp1vODcmWvKMRA9BNmbv6a/7QIWpPxHddWR8= @@ -169,12 +173,16 @@ github.com/containerd/log v0.1.0 h1:TCJt7ioM2cr/tfR8GPbGf9/VRAX8D2B4PjzCpfX540I= github.com/containerd/log v0.1.0/go.mod h1:VRRf09a7mHDIRezVKTRCrOq78v577GXq3bSa3EhrzVo= 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/cskr/pubsub v1.0.2 h1:vlOzMhl6PFn60gRlTQQsIfVwaPB/B/8MziK8FhEPt/0= +github.com/cskr/pubsub v1.0.2/go.mod h1:/8MzYXk/NJAz782G8RPkFzXTZVu63VotefPnR9TIRis= github.com/cyphar/filepath-securejoin v0.4.1 h1:JyxxyPEaktOD+GAnqIqTf9A8tHyAG22rowi7HkoSU1s= github.com/cyphar/filepath-securejoin v0.4.1/go.mod h1:Sdj7gXlvMcPZsbhwhQ33GguGLDGQL7h7bg04C/+u9jI= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM= github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/decred/dcrd/dcrec/secp256k1/v4 v4.4.0 h1:NMZiJj8QnKe1LgsbDayM4UoHwbvwDRwnI3hwNaAHRnc= +github.com/decred/dcrd/dcrec/secp256k1/v4 v4.4.0/go.mod h1:ZXNYxsqcloTdSy/rNShjYzMhyjf0LaoftYK0p+A3h40= github.com/dgraph-io/ristretto v0.2.0 h1:XAfl+7cmoUDWW/2Lx8TGZQjjxIQ2Ley9DSf52dru4WE= github.com/dgraph-io/ristretto v0.2.0/go.mod h1:8uBHCU/PBV4Ag0CJrP47b9Ofby5dqWNh4FicAdoqFNU= github.com/dgryski/go-farm v0.0.0-20200201041132-a6ae2369ad13 h1:fAjc9m62+UWV/WAFKLNi6ZS0675eEUC9y3AlwSbQu1Y= @@ -206,6 +214,8 @@ github.com/fatih/color v1.18.0 h1:S8gINlzdQ840/4pfAwic/ZE0djQEH3wM94VfqLTZcOM= github.com/fatih/color v1.18.0/go.mod h1:4FelSpRwEGDpQ12mAdzqdOukCy4u8WUtOY6lkT/6HfU= github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg= github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U= +github.com/frankban/quicktest v1.14.6 h1:7Xjx+VpznH+oBnejlPUj8oUpdxnVs4f8XU8WnHkI4W8= +github.com/frankban/quicktest v1.14.6/go.mod h1:4ptaffx2x8+WTWXmUCuVU6aPUX1/Mz7zb5vbUoiM6w0= github.com/fsnotify/fsnotify v1.4.7/go.mod h1:jwhsz4b93w/PPRr/qN1Yymfu8t87LnFCMoQvtojpjFo= github.com/fsnotify/fsnotify v1.4.9/go.mod h1:znqG4EE+3YCdAaPaxE2ZRY/06pZUdp0tY4IgpuI1SZQ= github.com/fsnotify/fsnotify v1.6.0 h1:n+5WquG0fcWoWp6xPWfHdbskMCQaFnG6PfBrh1Ky4HY= @@ -240,6 +250,8 @@ 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-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-redis/redis v6.15.9+incompatible h1:K0pv1D7EQUjfyoMql+r/jZqCLizCGKFlFgcHWWmHQjg= +github.com/go-redis/redis v6.15.9+incompatible/go.mod h1:NAIEuMOZ/fxfXJIrKDQDz8wamY7mA7PouImQ2Jvg6kA= github.com/go-task/slim-sprig v0.0.0-20210107165309-348f09dbbbc0/go.mod h1:fyg7847qk6SyHyPtNmDHnmrv/HOrqktSC+C9fM+CJOE= github.com/go-test/deep v1.1.1 h1:0r/53hagsehfO4bzD2Pgr/+RgHqhmf+k1Bpse2cTu1U= github.com/go-test/deep v1.1.1/go.mod h1:5C2ZWiW0ErCdrYzpqxLbTX7MG14M9iiw8DgHncVwcsE= @@ -254,6 +266,8 @@ github.com/goccy/go-json v0.10.5 h1:Fq85nIqj+gXn/S5ahsiTlK3TmC85qgirsdTP/+DeaC4= github.com/goccy/go-json v0.10.5/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M= github.com/gogo/protobuf v1.3.2 h1:Ov1cvc58UF3b5XjBnZv7+opcTcQFZebYjWzi34vdm4Q= github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69NZV8Q= +github.com/golang-jwt/jwt v3.2.2+incompatible h1:IfV12K8xAKAnZqdXVzCZ+TOjboZ2keLg81eXfW3O+oY= +github.com/golang-jwt/jwt v3.2.2+incompatible/go.mod h1:8pz2t5EyA70fFQQSrl6XZXzqecmYZeUEB8OUGHkxJ+I= github.com/golang-jwt/jwt/v5 v5.3.0 h1:pv4AsKCKKZuqlgs5sUmn4x8UlGa0kEVt/puTpKx9vvo= github.com/golang-jwt/jwt/v5 v5.3.0/go.mod h1:fxCRLWMO43lRc8nhHWY6LGqRcf+1gQWArsqaEUEa5bE= github.com/golang/groupcache v0.0.0-20241129210726-2c02b8208cf8 h1:f+oWsMOmNPc8JmEHVZIycC7hBoQxHH9pNKQORJNozsQ= @@ -336,13 +350,19 @@ github.com/hiddeco/sshsig v0.2.0 h1:gMWllgKCITXdydVkDL+Zro0PU96QI55LwUwebSwNTSw= github.com/hiddeco/sshsig v0.2.0/go.mod h1:nJc98aGgiH6Yql2doqH4CTBVHexQA40Q+hMMLHP4EqE= github.com/hpcloud/tail v1.0.0 h1:nfCOvKYfkgYP8hkirhJocXT2+zOD8yUNjXaWfTlyFKI= github.com/hpcloud/tail v1.0.0/go.mod h1:ab1qPbhIpdTxEkNHXyeSf5vhxWSCs/tWer42PpOxQnU= +github.com/huin/goupnp v1.3.0 h1:UvLUlWDNpoUdYzb2TCn+MuTWtcjXKSza2n6CBdQ0xXc= +github.com/huin/goupnp v1.3.0/go.mod h1:gnGPsThkYa7bFi/KWmEysQRf48l2dvR5bxr2OFckNX8= github.com/ianlancetaylor/demangle v0.0.0-20200824232613-28f6c0f3b639/go.mod h1:aSSvb/t6k1mPoxDqO4vJh6VOCGPwU4O0C2/Eqndh1Sc= github.com/ipfs/bbloom v0.0.4 h1:Gi+8EGJ2y5qiD5FbsbpX/TMNcJw8gSqr7eyjHa4Fhvs= github.com/ipfs/bbloom v0.0.4/go.mod h1:cS9YprKXpoZ9lT0n/Mw/a6/aFV6DTjTLYHeA+gyqMG0= github.com/ipfs/boxo v0.36.0 h1:DarrMBM46xCs6GU6Vz+AL8VUyXykqHAqZYx8mR0Oics= github.com/ipfs/boxo v0.36.0/go.mod h1:92hnRXfP5ScKEIqlq9Ns7LR1dFXEVADKWVGH0fjk83k= +github.com/ipfs/go-bitswap v0.11.0 h1:j1WVvhDX1yhG32NTC9xfxnqycqYIlhzEzLXG/cU1HyQ= +github.com/ipfs/go-bitswap v0.11.0/go.mod h1:05aE8H3XOU+LXpTedeAS0OZpcO1WFsj5niYQH9a1Tmk= github.com/ipfs/go-block-format v0.2.3 h1:mpCuDaNXJ4wrBJLrtEaGFGXkferrw5eqVvzaHhtFKQk= github.com/ipfs/go-block-format v0.2.3/go.mod h1:WJaQmPAKhD3LspLixqlqNFxiZ3BZ3xgqxxoSR/76pnA= +github.com/ipfs/go-blockservice v0.5.2 h1:in9Bc+QcXwd1apOVM7Un9t8tixPKdaHQFdLSUM1Xgk8= +github.com/ipfs/go-blockservice v0.5.2/go.mod h1:VpMblFEqG67A/H2sHKAemeH9vlURVavlysbdUI632yk= github.com/ipfs/go-cid v0.6.0 h1:DlOReBV1xhHBhhfy/gBNNTSyfOM6rLiIx9J7A4DGf30= github.com/ipfs/go-cid v0.6.0/go.mod h1:NC4kS1LZjzfhK40UGmpXv5/qD2kcMzACYJNntCUiDhQ= github.com/ipfs/go-datastore v0.9.0 h1:WocriPOayqalEsueHv6SdD4nPVl4rYMfYGLD4bqCZ+w= @@ -351,21 +371,47 @@ github.com/ipfs/go-detect-race v0.0.1 h1:qX/xay2W3E4Q1U7d9lNs1sU9nvguX0a7319XbyQ github.com/ipfs/go-detect-race v0.0.1/go.mod h1:8BNT7shDZPo99Q74BpGMK+4D8Mn4j46UU0LZ723meps= github.com/ipfs/go-ipfs-blockstore v1.3.1 h1:cEI9ci7V0sRNivqaOr0elDsamxXFxJMMMy7PTTDQNsQ= github.com/ipfs/go-ipfs-blockstore v1.3.1/go.mod h1:KgtZyc9fq+P2xJUiCAzbRdhhqJHvsw8u2Dlqy2MyRTE= +github.com/ipfs/go-ipfs-blocksutil v0.0.1 h1:Eh/H4pc1hsvhzsQoMEP3Bke/aW5P5rVM1IWFJMcGIPQ= +github.com/ipfs/go-ipfs-blocksutil v0.0.1/go.mod h1:Yq4M86uIOmxmGPUHv/uI7uKqZNtLb449gwKqXjIsnRk= +github.com/ipfs/go-ipfs-delay v0.0.1 h1:r/UXYyRcddO6thwOnhiznIAiSvxMECGgtv35Xs1IeRQ= +github.com/ipfs/go-ipfs-delay v0.0.1/go.mod h1:8SP1YXK1M1kXuc4KJZINY3TQQ03J2rwBG9QfXmbRPrw= github.com/ipfs/go-ipfs-ds-help v1.1.1 h1:B5UJOH52IbcfS56+Ul+sv8jnIV10lbjLF5eOO0C66Nw= github.com/ipfs/go-ipfs-ds-help v1.1.1/go.mod h1:75vrVCkSdSFidJscs8n4W+77AtTpCIAdDGAwjitJMIo= +github.com/ipfs/go-ipfs-exchange-interface v0.2.1 h1:jMzo2VhLKSHbVe+mHNzYgs95n0+t0Q69GQ5WhRDZV/s= +github.com/ipfs/go-ipfs-exchange-interface v0.2.1/go.mod h1:MUsYn6rKbG6CTtsDp+lKJPmVt3ZrCViNyH3rfPGsZ2E= +github.com/ipfs/go-ipfs-exchange-offline v0.3.0 h1:c/Dg8GDPzixGd0MC8Jh6mjOwU57uYokgWRFidfvEkuA= +github.com/ipfs/go-ipfs-exchange-offline v0.3.0/go.mod h1:MOdJ9DChbb5u37M1IcbrRB02e++Z7521fMxqCNRrz9s= +github.com/ipfs/go-ipfs-pq v0.0.4 h1:U7jjENWJd1jhcrR8X/xHTaph14PTAK9O+yaLJbjqgOw= +github.com/ipfs/go-ipfs-pq v0.0.4/go.mod h1:9UdLOIIb99IFrgT0Fc53pvbvlJBhpUb4GJuAQf3+O2A= +github.com/ipfs/go-ipfs-routing v0.3.0 h1:9W/W3N+g+y4ZDeffSgqhgo7BsBSJwPMcyssET9OWevc= +github.com/ipfs/go-ipfs-routing v0.3.0/go.mod h1:dKqtTFIql7e1zYsEuWLyuOU+E0WJWW8JjbTPLParDWo= github.com/ipfs/go-ipfs-util v0.0.3 h1:2RFdGez6bu2ZlZdI+rWfIdbQb1KudQp3VGwPtdNCmE0= github.com/ipfs/go-ipfs-util v0.0.3/go.mod h1:LHzG1a0Ig4G+iZ26UUOMjHd+lfM84LZCrn17xAKWBvs= github.com/ipfs/go-ipld-cbor v0.2.1 h1:H05yEJbK/hxg0uf2AJhyerBDbjOuHX4yi+1U/ogRa7E= github.com/ipfs/go-ipld-cbor v0.2.1/go.mod h1:x9Zbeq8CoE5R2WicYgBMcr/9mnkQ0lHddYWJP2sMV3A= github.com/ipfs/go-ipld-format v0.6.3 h1:9/lurLDTotJpZSuL++gh3sTdmcFhVkCwsgx2+rAh4j8= github.com/ipfs/go-ipld-format v0.6.3/go.mod h1:74ilVN12NXVMIV+SrBAyC05UJRk0jVvGqdmrcYZvCBk= +github.com/ipfs/go-ipld-legacy v0.2.2 h1:DThbqCPVLpWBcGtU23KDLiY2YRZZnTkXQyfz8aOfBkQ= +github.com/ipfs/go-ipld-legacy v0.2.2/go.mod h1:hhkj+b3kG9b2BcUNw8IFYAsfeNo8E3U7eYlWeAOPyDU= github.com/ipfs/go-log v1.0.5 h1:2dOuUCB1Z7uoczMWgAyDck5JLb72zHzrMnGnCNNbvY8= github.com/ipfs/go-log v1.0.5/go.mod h1:j0b8ZoR+7+R99LD9jZ6+AJsrzkPbSXbZfGakb5JPtIo= github.com/ipfs/go-log/v2 v2.1.3/go.mod h1:/8d0SH3Su5Ooc31QlL1WysJhvyOTDCjcCZ9Axpmri6g= github.com/ipfs/go-log/v2 v2.9.1 h1:3JXwHWU31dsCpvQ+7asz6/QsFJHqFr4gLgQ0FWteujk= github.com/ipfs/go-log/v2 v2.9.1/go.mod h1:evFx7sBiohUN3AG12mXlZBw5hacBQld3ZPHrowlJYoo= +github.com/ipfs/go-merkledag v0.11.0 h1:DgzwK5hprESOzS4O1t/wi6JDpyVQdvm9Bs59N/jqfBY= +github.com/ipfs/go-merkledag v0.11.0/go.mod h1:Q4f/1ezvBiJV0YCIXvt51W/9/kqJGH4I1LsA7+djsM4= 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/ipfs/go-peertaskqueue v0.8.3 h1:tBPpGJy+A92RqtRFq5amJn0Uuj8Pw8tXi0X3eHfHM8w= +github.com/ipfs/go-peertaskqueue v0.8.3/go.mod h1:OqVync4kPOcXEGdj/LKvox9DCB5mkSBeXsPczCxLtYA= +github.com/ipfs/go-verifcid v0.0.3 h1:gmRKccqhWDocCRkC+a59g5QW7uJw5bpX9HWBevXa0zs= +github.com/ipfs/go-verifcid v0.0.3/go.mod h1:gcCtGniVzelKrbk9ooUSX/pM3xlH73fZZJDzQJRvOUw= +github.com/ipld/go-car v0.6.2 h1:Hlnl3Awgnq8icK+ze3iRghk805lu8YNq3wlREDTF2qc= +github.com/ipld/go-car v0.6.2/go.mod h1:oEGXdwp6bmxJCZ+rARSkDliTeYnVzv3++eXajZ+Bmr8= +github.com/ipld/go-codec-dagpb v1.7.0 h1:hpuvQjCSVSLnTnHXn+QAMR0mLmb1gA6wl10LExo2Ts0= +github.com/ipld/go-codec-dagpb v1.7.0/go.mod h1:rD3Zg+zub9ZnxcLwfol/OTQRVjaLzXypgy4UqHQvilM= +github.com/ipld/go-ipld-prime v0.21.0 h1:n4JmcpOlPDIxBcY037SVfpd1G+Sj1nKZah0m6QH9C2E= +github.com/ipld/go-ipld-prime v0.21.0/go.mod h1:3RLqy//ERg/y5oShXXdx5YIp50cFGOanyMctpPjsvxQ= 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= @@ -374,6 +420,14 @@ 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/jackpal/go-nat-pmp v1.0.2 h1:KzKSgb7qkJvOUTqYl9/Hg/me3pWgBmERKrTGD7BdWus= +github.com/jackpal/go-nat-pmp v1.0.2/go.mod h1:QPH045xvCAeXUZOxsnwmrtiCoxIr9eob+4orBN1SBKc= +github.com/jbenet/goprocess v0.1.4 h1:DRGOFReOMqqDNXwW70QkacFW0YN9QnwLV0Vqk+3oU0o= +github.com/jbenet/goprocess v0.1.4/go.mod h1:5yspPrukOVuOLORacaBi858NqyClJPQxYZlqdZVfqY4= +github.com/jinzhu/inflection v1.0.0 h1:K317FqzuhWc8YvSVlFMCCUb36O/S9MCKRDI7QkRKD/E= +github.com/jinzhu/inflection v1.0.0/go.mod h1:h+uFLlag+Qp1Va5pdKtLDYj+kHp5pxUVkryuEj+Srlc= +github.com/jinzhu/now v1.1.5 h1:/o9tlHleP7gOFmsnYNz3RGnqzefHA47wQpKrrdTIwXQ= +github.com/jinzhu/now v1.1.5/go.mod h1:d3SSVoowX0Lcu0IBviAWJpolVfI5UJVZZ7cO71lE/z8= 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= @@ -387,6 +441,8 @@ github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zt github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= 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/koron/go-ssdp v0.0.6 h1:Jb0h04599eq/CY7rB5YEqPS83HmRfHP2azkxMN2rFtU= +github.com/koron/go-ssdp v0.0.6/go.mod h1:0R9LfRJGek1zWTjN3JUNlm5INCDYGpRDfAptnct63fI= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= github.com/kr/pretty v0.3.0/go.mod h1:640gp4NfQd8pI5XOwp5fnNeVWj67G7CFk/SaSQn7NBk= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= @@ -397,6 +453,24 @@ 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/labstack/echo/v4 v4.11.3 h1:Upyu3olaqSHkCjs1EJJwQ3WId8b8b1hxbogyommKktM= +github.com/labstack/echo/v4 v4.11.3/go.mod h1:UcGuQ8V6ZNRmSweBIJkPvGfwCMIlFmiqrPqiEBfPYws= +github.com/labstack/gommon v0.4.1 h1:gqEff0p/hTENGMABzezPoPSRtIh1Cvw0ueMOe0/dfOk= +github.com/labstack/gommon v0.4.1/go.mod h1:TyTrpPqxR5KMk8LKVtLmfMjeQ5FEkBYdxLYPw/WfrOM= +github.com/libp2p/go-buffer-pool v0.1.0 h1:oK4mSFcQz7cTQIfqbe4MIj9gLW+mnanjyFtc6cdF0Y8= +github.com/libp2p/go-buffer-pool v0.1.0/go.mod h1:N+vh8gMqimBzdKkSMVuydVDq+UV5QTWy5HSiZacSbPg= +github.com/libp2p/go-libp2p v0.47.0 h1:qQpBjSCWNQFF0hjBbKirMXE9RHLtSuzTDkTfr1rw0yc= +github.com/libp2p/go-libp2p v0.47.0/go.mod h1:s8HPh7mMV933OtXzONaGFseCg/BE//m1V34p3x4EUOY= +github.com/libp2p/go-libp2p-asn-util v0.4.1 h1:xqL7++IKD9TBFMgnLPZR6/6iYhawHKHl950SO9L6n94= +github.com/libp2p/go-libp2p-asn-util v0.4.1/go.mod h1:d/NI6XZ9qxw67b4e+NgpQexCIiFYJjErASrYW4PFDN8= +github.com/libp2p/go-libp2p-record v0.3.1 h1:cly48Xi5GjNw5Wq+7gmjfBiG9HCzQVkiZOUZ8kUl+Fg= +github.com/libp2p/go-libp2p-record v0.3.1/go.mod h1:T8itUkLcWQLCYMqtX7Th6r7SexyUJpIyPgks757td/E= +github.com/libp2p/go-libp2p-testing v0.12.0 h1:EPvBb4kKMWO29qP4mZGyhVzUyR25dvfUIK5WDu6iPUA= +github.com/libp2p/go-libp2p-testing v0.12.0/go.mod h1:KcGDRXyN7sQCllucn1cOOS+Dmm7ujhfEyXQL5lvkcPg= +github.com/libp2p/go-msgio v0.3.0 h1:mf3Z8B1xcFN314sWX+2vOTShIE0Mmn2TXn3YCUQGNj0= +github.com/libp2p/go-msgio v0.3.0/go.mod h1:nyRM819GmVaF9LX3l03RMh10QdOroF++NBbxAb0mmDM= +github.com/libp2p/go-netroute v0.4.0 h1:sZZx9hyANYUx9PZyqcgE/E1GUG3iEtTZHUEvdtXT7/Q= +github.com/libp2p/go-netroute v0.4.0/go.mod h1:Nkd5ShYgSMS5MUKy/MU2T57xFoOKvvLR92Lic48LEyA= 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= @@ -438,10 +512,18 @@ github.com/multiformats/go-base32 v0.1.0 h1:pVx9xoSPqEIQG8o+UbAe7DNi51oej1NtK+aG github.com/multiformats/go-base32 v0.1.0/go.mod h1:Kj3tFY6zNr+ABYMqeUNeGvkIC/UYgtWibDcT0rExnbI= github.com/multiformats/go-base36 v0.2.0 h1:lFsAbNOGeKtuKozrtBsAkSVhv1p9D0/qedU9rQyccr0= github.com/multiformats/go-base36 v0.2.0/go.mod h1:qvnKE++v+2MWCfePClUEjE78Z7P2a1UV0xHgWc0hkp4= +github.com/multiformats/go-multiaddr v0.16.1 h1:fgJ0Pitow+wWXzN9do+1b8Pyjmo8m5WhGfzpL82MpCw= +github.com/multiformats/go-multiaddr v0.16.1/go.mod h1:JSVUmXDjsVFiW7RjIFMP7+Ev+h1DTbiJgVeTV/tcmP0= +github.com/multiformats/go-multiaddr-fmt v0.1.0 h1:WLEFClPycPkp4fnIzoFoV9FVd49/eQsuaL3/CWe167E= +github.com/multiformats/go-multiaddr-fmt v0.1.0/go.mod h1:hGtDIW4PU4BqJ50gW2quDuPVjyWNZxToGUh/HwTZYJo= github.com/multiformats/go-multibase v0.2.0 h1:isdYCVLvksgWlMW9OZRYJEa9pZETFivncJHmHnnd87g= github.com/multiformats/go-multibase v0.2.0/go.mod h1:bFBZX4lKCA/2lyOFSAoKH5SS6oPyjtnzK/XTFDPkNuk= +github.com/multiformats/go-multicodec v0.10.0 h1:UpP223cig/Cx8J76jWt91njpK3GTAO1w02sdcjZDSuc= +github.com/multiformats/go-multicodec v0.10.0/go.mod h1:wg88pM+s2kZJEQfRCKBNU+g32F5aWBEjyFHXvZLTcLI= github.com/multiformats/go-multihash v0.2.3 h1:7Lyc8XfX/IY2jWb/gI7JP+o7JEq9hOa7BFvVU9RSh+U= github.com/multiformats/go-multihash v0.2.3/go.mod h1:dXgKXCXjBzdscBLk9JkjINiEsCKRVch90MdaGiKsvSM= +github.com/multiformats/go-multistream v0.6.1 h1:4aoX5v6T+yWmc2raBHsTvzmFhOI8WVOer28DeBBEYdQ= +github.com/multiformats/go-multistream v0.6.1/go.mod h1:ksQf6kqHAb6zIsyw7Zm+gAuVo57Qbq84E27YlYqavqw= github.com/multiformats/go-varint v0.1.0 h1:i2wqFp4sdl3IcIxfAonHQV9qU5OsZ4Ts9IOoETFs5dI= github.com/multiformats/go-varint v0.1.0/go.mod h1:5KVAVXegtfmNQQm/lCY+ATvDzvJJhSkUlGQV9wgObdI= github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= @@ -506,6 +588,8 @@ github.com/prometheus/common v0.67.5 h1:pIgK94WWlQt1WLwAC5j2ynLaBRDiinoAb86HZHTU github.com/prometheus/common v0.67.5/go.mod h1:SjE/0MzDEEAyrdr5Gqc6G+sXI67maCxzaT3A2+HqjUw= github.com/prometheus/procfs v0.19.2 h1:zUMhqEW66Ex7OXIiDkll3tl9a1ZdilUOd/F6ZXw4Vws= github.com/prometheus/procfs v0.19.2/go.mod h1:M0aotyiemPhBCM0z5w87kL22CxfcH05ZpYlu+b4J7mw= +github.com/puzpuzpuz/xsync/v4 v4.2.0 h1:dlxm77dZj2c3rxq0/XNvvUKISAmovoXF4a4qM6Wvkr0= +github.com/puzpuzpuz/xsync/v4 v4.2.0/go.mod h1:VJDmTCJMBt8igNxnkQd86r+8KUeN1quSfNKu5bLYFQo= github.com/redis/go-redis/v9 v9.0.0-rc.4/go.mod h1:Vo3EsyWnicKnSKCA7HhgnvnyA74wOA69Cd2Meli5mmA= github.com/redis/go-redis/v9 v9.7.3 h1:YpPyAayJV+XErNsatSElgRZZVCwXX9QzkKYNvO7x0wM= github.com/redis/go-redis/v9 v9.7.3/go.mod h1:bGUrSggJ9X9GUmZpZNEOQKaANxSGgOEBRltRTZHSvrA= @@ -565,6 +649,10 @@ github.com/tidwall/sjson v1.2.5/go.mod h1:Fvgq9kS/6ociJEDnK0Fk1cpYF4FIW6ZF7LAe+6 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= +github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6KllzawFIhcdPw= +github.com/valyala/bytebufferpool v1.0.0/go.mod h1:6bBcMArwyJ5K/AmCkWv1jt77kVWyCJ6HpOuEn7z0Csc= +github.com/valyala/fasttemplate v1.2.2 h1:lxLXG0uE3Qnshl9QyaK6XJxMXlQZELvChBOCmQD0Loo= +github.com/valyala/fasttemplate v1.2.2/go.mod h1:KHLXt3tVN2HBp8eijSv/kGJopbvo7S+qRAEEKiv+SiQ= github.com/vmihailenco/go-tinylfu v0.2.2 h1:H1eiG6HM36iniK6+21n9LLpzx1G9R3DJa2UjUjbynsI= github.com/vmihailenco/go-tinylfu v0.2.2/go.mod h1:CutYi2Q9puTxfcolkliPq4npPuofg9N9t8JVrjzwa3Q= github.com/vmihailenco/msgpack/v5 v5.3.4/go.mod h1:7xyJ9e+0+9SaZT0Wt1RGleJXzli6Q/V5KbhBonMG9jc= @@ -572,6 +660,8 @@ github.com/vmihailenco/msgpack/v5 v5.4.1 h1:cQriyiUvjTwOHg8QZaPihLWeRAAVoCpE00IU github.com/vmihailenco/msgpack/v5 v5.4.1/go.mod h1:GaZTsDaehaPpQVyxrf5mtQlH+pc21PIudVV/E3rRQok= github.com/vmihailenco/tagparser/v2 v2.0.0 h1:y09buUbR+b5aycVFQs/g70pqKVZNBmxwAhO7/IwNM9g= github.com/vmihailenco/tagparser/v2 v2.0.0/go.mod h1:Wri+At7QHww0WTrCBeu4J6bNtoV6mEfg5OIWRZA9qds= +github.com/warpfork/go-testmark v0.12.1 h1:rMgCpJfwy1sJ50x0M0NgyphxYYPMOODIJHhsXyEHU0s= +github.com/warpfork/go-testmark v0.12.1/go.mod h1:kHwy7wfvGSPh1rQJYKayD4AbtNaeyZdcGi9tNJTaa5Y= github.com/warpfork/go-wish v0.0.0-20220906213052-39a1cc7a02d0 h1:GDDkbFiaK8jsSDJfjId/PEGEShv6ugrt4kYsC5UIDaQ= github.com/warpfork/go-wish v0.0.0-20220906213052-39a1cc7a02d0/go.mod h1:x6AKhvSSexNrVSrViXSHUEbICjmGXhtgABaHIySUSGw= github.com/whyrusleeping/cbor-gen v0.3.1 h1:82ioxmhEYut7LBVGhGq8xoRkXPLElVuh5mV67AFfdv0= @@ -809,6 +899,12 @@ gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C gopkg.in/yaml.v3 v3.0.0/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gorm.io/driver/postgres v1.6.0 h1:2dxzU8xJ+ivvqTRph34QX+WrRaJlmfyPqXmoGVjMBa4= +gorm.io/driver/postgres v1.6.0/go.mod h1:vUw0mrGgrTK+uPHEhAdV4sfFELrByKVGnaVRkXDhtWo= +gorm.io/driver/sqlite v1.6.0 h1:WHRRrIiulaPiPFmDcod6prc4l2VGVWHz80KspNsxSfQ= +gorm.io/driver/sqlite v1.6.0/go.mod h1:AO9V1qIQddBESngQUKWL9yoH93HIeA1X6V633rBwyT8= +gorm.io/gorm v1.31.1 h1:7CA8FTFz/gRfgqgpeKIBcervUn3xSyPUmr6B2WXJ7kg= +gorm.io/gorm v1.31.1/go.mod h1:XyQVbO2k6YkOis7C2437jSit3SsDK72s7n7rsSHd+Gs= gotest.tools/v3 v3.5.2 h1:7koQfIKdy+I8UTetycgUqXWSDwpgv193Ka+qRsmBY8Q= gotest.tools/v3 v3.5.2/go.mod h1:LtdLGcnqToBH83WByAAi/wiwSFCArdFIUV/xxN4pcjA= honnef.co/go/tools v0.0.1-2019.2.3/go.mod h1:a3bituU0lyd329TUQxRnasdCoJDkEUEAqEt0JzvZhAg= diff --git a/spindle/config/config.go b/spindle/config/config.go index 46cc83c4..ba17151b 100644 --- a/spindle/config/config.go +++ b/spindle/config/config.go @@ -13,6 +13,7 @@ type Server struct { DBPath string `env:"DB_PATH, default=spindle.db"` Hostname string `env:"HOSTNAME, required"` JetstreamEndpoint string `env:"JETSTREAM_ENDPOINT, default=wss://jetstream1.us-west.bsky.network/subscribe"` + Tap Tap `env:",prefix=TAP_"` PlcUrl string `env:"PLC_URL, default=https://plc.directory"` Dev bool `env:"DEV, default=false"` Owner string `env:"OWNER, required"` @@ -23,6 +24,15 @@ type Server struct { MaxConcurrentWorkflows int `env:"MAX_CONCURRENT_WORKFLOWS, default=8"` // max number of workflow containers running at once (memory cap) } +type Tap struct { + Embed bool `env:"EMBED, default=true"` + Url string `env:"URL, default=http://[::1]:2480"` + Bind string `env:"BIND, default=[::1]:2480"` + DBPath string `env:"DB_PATH, default=tap.db"` + RelayUrl string `env:"RELAY_URL, default=https://bsky.network"` + AdminPassword string `env:"ADMIN_PASSWORD"` +} + func (s Server) Did() syntax.DID { return syntax.DID(fmt.Sprintf("did:web:%s", s.Hostname)) } diff --git a/spindle/embedtap.go b/spindle/embedtap.go new file mode 100644 index 00000000..1415ed79 --- /dev/null +++ b/spindle/embedtap.go @@ -0,0 +1,127 @@ +package spindle + +import ( + "context" + "crypto/rand" + "encoding/hex" + "errors" + "fmt" + "log/slog" + "net" + "net/http" + "strings" + "time" + + "github.com/bluesky-social/indigo/service/tap" + "tangled.org/core/api/tangled" + "tangled.org/core/spindle/config" +) + +func randomAdminPassword() (string, error) { + var b [32]byte + if _, err := rand.Read(b[:]); err != nil { + return "", fmt.Errorf("generate tap admin password: %w", err) + } + return hex.EncodeToString(b[:]), nil +} + +func assertLoopbackBind(bind string) error { + host, _, err := net.SplitHostPort(bind) + if err != nil { + return fmt.Errorf("parse tap bind %q: %w", bind, err) + } + if host == "" { + return fmt.Errorf("embedded mode requires loopback host in tap bind %q", bind) + } + if strings.EqualFold(host, "localhost") { + return nil + } + ip := net.ParseIP(host) + if ip == nil || !ip.IsLoopback() { + return fmt.Errorf("embedded tap bind %q must be loopback like 127.0.0.1 or ::1", bind) + } + return nil +} + +type embeddedTap struct { + tap *tap.Tap + logger *slog.Logger +} + +func startEmbeddedTap(ctx context.Context, cfg *config.Config, logger *slog.Logger) (*embeddedTap, error) { + if err := assertLoopbackBind(cfg.Server.Tap.Bind); err != nil { + return nil, err + } + + tcfg := tap.Config{ + DatabaseURL: "sqlite://" + cfg.Server.Tap.DBPath, + DBMaxConns: 32, + PLCURL: cfg.Server.PlcUrl, + RelayUrl: cfg.Server.Tap.RelayUrl, + FirehoseParallelism: 4, + ResyncParallelism: 2, + OutboxParallelism: 1, + FirehoseCursorSaveInterval: time.Second, + RepoFetchTimeout: 5 * time.Minute, + IdentityCacheSize: 50_000, + EventCacheSize: 10_000, + SignalCollection: tangled.RepoNSID, + CollectionFilters: []string{tangled.RepoNSID, tangled.RepoCollaboratorNSID}, + AdminPassword: cfg.Server.Tap.AdminPassword, + RetryTimeout: 60 * time.Second, + } + + t, err := tap.New(tcfg) + if err != nil { + return nil, fmt.Errorf("tap.New: %w", err) + } + + go func() { + if err := t.Firehose.Run(ctx); err != nil && !errors.Is(err, context.Canceled) { + logger.Error("firehose terminated", "err", err) + } + }() + t.Run(ctx) + go func() { + logger.Info("tap http server listening", "bind", cfg.Server.Tap.Bind) + if err := t.Server.Start(cfg.Server.Tap.Bind); err != nil && !errors.Is(err, http.ErrServerClosed) { + logger.Error("tap http server terminated", "err", err) + } + }() + + if err := waitForListener(ctx, cfg.Server.Tap.Bind, time.Now().Add(10*time.Second)); err != nil { + logger.Warn("tap http server unreachable before timeout", "bind", cfg.Server.Tap.Bind, "err", err) + } + + return &embeddedTap{tap: t, logger: logger}, nil +} + +func waitForListener(ctx context.Context, addr string, deadline time.Time) error { + if ctx.Err() != nil { + return ctx.Err() + } + if time.Now().After(deadline) { + return fmt.Errorf("timed out waiting for %s", addr) + } + c, err := net.DialTimeout("tcp", addr, 100*time.Millisecond) + if err == nil { + c.Close() + return nil + } + time.Sleep(50 * time.Millisecond) + return waitForListener(ctx, addr, deadline) +} + +func (e *embeddedTap) Shutdown() { + if e == nil || e.tap == nil { + return + } + shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + if err := e.tap.Server.Shutdown(shutdownCtx); err != nil { + e.logger.Error("tap server shutdown failed", "err", err) + } + if err := e.tap.CloseDb(shutdownCtx); err != nil { + e.logger.Error("tap db close failed", "err", err) + } +} diff --git a/spindle/ingester.go b/spindle/ingester.go index 9a3835e0..2137b85d 100644 --- a/spindle/ingester.go +++ b/spindle/ingester.go @@ -3,22 +3,14 @@ package spindle import ( "context" "encoding/json" - "errors" "fmt" - "strings" "time" "tangled.org/core/api/tangled" - "tangled.org/core/eventconsumer" - "tangled.org/core/rbac" "tangled.org/core/spindle/db" - comatproto "github.com/bluesky-social/indigo/api/atproto" - "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" - "github.com/bluesky-social/indigo/xrpc" "github.com/bluesky-social/jetstream/pkg/models" - securejoin "github.com/cyphar/filepath-securejoin" ) type Ingester func(ctx context.Context, e *models.Event) error @@ -33,10 +25,6 @@ func (s *Spindle) ingest() Ingester { switch e.Commit.Collection { case tangled.SpindleMemberNSID: err = s.ingestMember(ctx, e) - case tangled.RepoNSID: - err = s.ingestRepo(ctx, e) - case tangled.RepoCollaboratorNSID: - err = s.ingestCollaborator(ctx, e) } if err != nil { @@ -135,221 +123,3 @@ func (s *Spindle) ingestMember(_ context.Context, e *models.Event) error { } return nil } - -func (s *Spindle) ingestRepo(ctx context.Context, e *models.Event) error { - var err error - did := e.Did - - l := s.l.With("component", "ingester", "record", tangled.RepoNSID) - - l.Info("ingesting repo record", "did", did) - - switch e.Commit.Operation { - case models.CommitOperationCreate, models.CommitOperationUpdate: - raw := e.Commit.Record - record := tangled.Repo{} - err = json.Unmarshal(raw, &record) - if err != nil { - l.Error("invalid record", "error", err) - return err - } - - domain := s.cfg.Server.Hostname - rkey := e.Commit.RKey - - // no spindle configured for this repo - if record.Spindle == nil { - l.Info("no spindle configured", "rkey", rkey) - return nil - } - - // this repo did not want this spindle - if *record.Spindle != domain { - l.Info("different spindle configured", "rkey", rkey, "spindle", *record.Spindle, "domain", domain) - return nil - } - - // add this repo to the watch list - if err := s.db.AddRepo(record.Knot, did, rkey); err != nil { - l.Error("failed to add repo", "error", err) - return fmt.Errorf("failed to add repo: %w", err) - } - - didSlashRepo, err := securejoin.SecureJoin(did, rkey) - if err != nil { - return err - } - - // add repo to rbac - if err := s.e.AddRepo(did, rbac.ThisServer, didSlashRepo); err != nil { - l.Error("failed to add repo to enforcer", "error", err) - return fmt.Errorf("failed to add repo: %w", err) - } - - // add collaborators to rbac - owner, err := s.res.ResolveIdent(ctx, did) - if err != nil || owner.Handle.IsInvalidHandle() { - return err - } - if err := s.fetchAndAddCollaborators(ctx, owner, didSlashRepo); err != nil { - return err - } - - // add this knot to the event consumer - src := eventconsumer.NewKnotSource(record.Knot) - s.ks.AddSource(context.Background(), src) - - return nil - - } - return nil -} - -func (s *Spindle) ingestCollaborator(ctx context.Context, e *models.Event) error { - var err error - - l := s.l.With("component", "ingester", "record", tangled.RepoCollaboratorNSID, "did", e.Did) - - l.Info("ingesting collaborator record") - - switch e.Commit.Operation { - case models.CommitOperationCreate, models.CommitOperationUpdate: - raw := e.Commit.Record - record := tangled.RepoCollaborator{} - err = json.Unmarshal(raw, &record) - if err != nil { - l.Error("invalid record", "error", err) - return err - } - - subjectId, err := s.res.ResolveIdent(ctx, record.Subject) - if err != nil || subjectId.Handle.IsInvalidHandle() { - return err - } - - var rbacResource string - var ownerDid string - switch { - case strings.HasPrefix(record.Repo, "did:"): - resolvedOwner, repoName, lookupErr := s.resolveRepoDid(ctx, e.Did, record.Repo) - if lookupErr != nil { - return fmt.Errorf("unknown repo DID %s: %w", record.Repo, lookupErr) - } - ownerDid = resolvedOwner - rbacResource, _ = securejoin.SecureJoin(ownerDid, repoName) - - case strings.Contains(record.Repo, "/"): - repoAt, parseErr := syntax.ParseATURI(record.Repo) - if parseErr != nil { - l.Info("rejecting record, invalid repoAt", "repoAt", record.Repo) - return nil - } - - owner, resolveErr := s.res.ResolveIdent(ctx, repoAt.Authority().String()) - if resolveErr != nil || owner.Handle.IsInvalidHandle() { - return fmt.Errorf("failed to resolve handle: %w", resolveErr) - } - - xrpcc := xrpc.Client{ - Host: owner.PDSEndpoint(), - } - - resp, getErr := comatproto.RepoGetRecord(ctx, &xrpcc, "", tangled.RepoNSID, repoAt.Authority().String(), repoAt.RecordKey().String()) - if getErr != nil { - return getErr - } - - if _, ok := resp.Value.Val.(*tangled.Repo); !ok { - return fmt.Errorf("record at %s is not a tangled.Repo", repoAt) - } - rbacResource, _ = securejoin.SecureJoin(owner.DID.String(), repoAt.RecordKey().String()) - ownerDid = owner.DID.String() - - default: - l.Info("rejecting collaborator record with unrecognized repo format", "repo", record.Repo) - return nil - } - - if ok, err := s.e.IsCollaboratorInviteAllowed(ownerDid, rbac.ThisServer, rbacResource); !ok || err != nil { - return fmt.Errorf("insufficient permissions: %w", err) - } - - if err := s.e.AddCollaborator(record.Subject, rbac.ThisServer, rbacResource); err != nil { - l.Error("failed to add collaborator to enforcer", "error", err) - return fmt.Errorf("failed to add collaborator: %w", err) - } - - return nil - } - return nil -} - -func (s *Spindle) resolveRepoDid(ctx context.Context, ownerDid string, repoDid string) (string, string, error) { - owner, resolveErr := s.res.ResolveIdent(ctx, ownerDid) - if resolveErr != nil || owner.Handle.IsInvalidHandle() { - return "", "", fmt.Errorf("failed to resolve owner %s: %w", ownerDid, resolveErr) - } - - xrpcc := xrpc.Client{ - Host: owner.PDSEndpoint(), - } - - cursor := "" - for { - resp, listErr := comatproto.RepoListRecords(ctx, &xrpcc, tangled.RepoNSID, cursor, 100, ownerDid, false) - if listErr != nil { - return "", "", fmt.Errorf("failed to list repo records for %s: %w", ownerDid, listErr) - } - - for _, r := range resp.Records { - if r == nil { - continue - } - repo, ok := r.Value.Val.(*tangled.Repo) - if !ok { - continue - } - if repo.RepoDid != nil && *repo.RepoDid == repoDid { - rkey := r.Uri[strings.LastIndex(r.Uri, "/")+1:] - return ownerDid, rkey, nil - } - } - - if resp.Cursor == nil || *resp.Cursor == "" { - break - } - cursor = *resp.Cursor - } - - return "", "", fmt.Errorf("repo DID %s not found in records for %s", repoDid, ownerDid) -} - -func (s *Spindle) fetchAndAddCollaborators(ctx context.Context, owner *identity.Identity, didSlashRepo string) error { - l := s.l.With("component", "ingester", "handler", "fetchAndAddCollaborators") - - l.Info("fetching and adding existing collaborators") - - xrpcc := xrpc.Client{ - Host: owner.PDSEndpoint(), - } - - resp, err := comatproto.RepoListRecords(ctx, &xrpcc, tangled.RepoCollaboratorNSID, "", 50, owner.DID.String(), false) - if err != nil { - return err - } - - var errs error - for _, r := range resp.Records { - if r == nil { - continue - } - record := r.Value.Val.(*tangled.RepoCollaborator) - - if err := s.e.AddCollaborator(record.Subject, rbac.ThisServer, didSlashRepo); err != nil { - l.Error("failed to add repo to enforcer", "error", err) - errors.Join(errs, fmt.Errorf("failed to add repo: %w", err)) - } - } - - return errs -} diff --git a/spindle/server.go b/spindle/server.go index b67e2d3a..446d8f09 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -10,6 +10,7 @@ import ( "net/http" "sync" + "github.com/bluesky-social/indigo/atproto/syntax" "github.com/go-chi/chi/v5" "tangled.org/core/api/tangled" "tangled.org/core/eventconsumer" @@ -39,6 +40,8 @@ const ( type Spindle struct { jc *jetstream.JetstreamClient + tap *Tap + embedTap *embeddedTap db *db.DB e *rbac.Enforcer l *slog.Logger @@ -52,13 +55,14 @@ type Spindle struct { motd []byte motdMu sync.RWMutex workflowSem chan struct{} + rootCtx context.Context } // New creates a new Spindle server with the provided configuration and engines. func New(ctx context.Context, cfg *config.Config, engines map[string]models.Engine) (*Spindle, error) { logger := log.FromContext(ctx) - d, err := db.Make(cfg.Server.DBPath) + d, err := db.Make(ctx, cfg.Server.DBPath) if err != nil { return nil, fmt.Errorf("failed to setup db: %w", err) } @@ -104,8 +108,6 @@ func New(ctx context.Context, cfg *config.Config, engines map[string]models.Engi collections := []string{ tangled.SpindleMemberNSID, - tangled.RepoNSID, - tangled.RepoCollaboratorNSID, } jc, err := jetstream.NewJetstreamClient(cfg.Server.JetstreamEndpoint, "spindle", collections, nil, log.SubLogger(logger, "jetstream"), d, true, true) if err != nil { @@ -137,6 +139,7 @@ func New(ctx context.Context, cfg *config.Config, engines map[string]models.Engi vault: vault, motd: defaultMotd, workflowSem: workflowSem, + rootCtx: ctx, } err = e.AddSpindle(rbacDomain) @@ -177,6 +180,16 @@ func New(ctx context.Context, cfg *config.Config, engines map[string]models.Engi } spindle.ks = eventconsumer.NewConsumer(*ccfg) + if cfg.Server.Tap.Embed { + pw, err := randomAdminPassword() + if err != nil { + return nil, err + } + cfg.Server.Tap.AdminPassword = pw + logger.Info("embedded tap: using random admin password") + } + spindle.tap = NewTapClient(spindle) + return spindle, nil } @@ -235,15 +248,52 @@ func (s *Spindle) Start(ctx context.Context) error { defer stopper.Stop() } + if s.cfg.Server.Tap.Embed { + emb, err := startEmbeddedTap(ctx, s.cfg, log.SubLogger(s.l, "embedtap")) + if err != nil { + return fmt.Errorf("starting embedded tap: %w", err) + } + s.embedTap = emb + defer s.embedTap.Shutdown() + } + go func() { s.l.Info("starting knot event consumer") s.ks.Start(ctx) }() + s.l.Info("starting tap client", "url", s.cfg.Server.Tap.Url) + s.tap.Start(ctx) + s.l.Info("starting spindle server", "address", s.cfg.Server.ListenAddr) return http.ListenAndServe(s.cfg.Server.ListenAddr, s.Router()) } +func (s *Spindle) declareTapInterest(ctx context.Context) { + repos, err := s.db.AllRepos() + if err != nil { + s.l.Warn("tap declare: failed to load known repos", "err", err) + return + } + seen := make(map[syntax.DID]struct{}, len(repos)) + dids := make([]syntax.DID, 0, len(repos)) + for _, r := range repos { + if r.Owner == "" { + continue + } + if _, ok := seen[r.Owner]; ok { + continue + } + seen[r.Owner] = struct{}{} + dids = append(dids, r.Owner) + } + if err := s.tap.AddOwnerDIDs(ctx, dids); err != nil { + s.l.Warn("tap declare: AddRepos rejected", "count", len(dids), "err", err) + return + } + s.l.Info("tap declare: known owner DIDs registered", "count", len(dids)) +} + func Run(ctx context.Context) error { cfg, err := config.Load(ctx) if err != nil { @@ -319,19 +369,9 @@ func (s *Spindle) processPipeline(ctx context.Context, src eventconsumer.Source, return fmt.Errorf("repo knot does not match event source: %s != %s", src.Key(), tpl.TriggerMetadata.Repo.Knot) } - // filter by repos - repoName := "" - if tpl.TriggerMetadata.Repo.Repo != nil { - repoName = *tpl.TriggerMetadata.Repo.Repo - } - - _, err = s.db.GetRepo( - tpl.TriggerMetadata.Repo.Knot, - tpl.TriggerMetadata.Repo.Did, - repoName, - ) + repoDid, err := s.resolvePipelineRepoDid(tpl.TriggerMetadata.Repo) if err != nil { - return fmt.Errorf("failed to get repo: %w", err) + return err } pipelineId := models.PipelineId{ @@ -399,8 +439,7 @@ func (s *Spindle) processPipeline(ctx context.Context, src eventconsumer.Source, ok := s.jq.Enqueue(queue.Job{ Run: func() error { engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, s.workflowSem, ctx, &models.Pipeline{ - RepoOwner: tpl.TriggerMetadata.Repo.Did, - RepoName: repoName, + RepoDid: repoDid, Workflows: workflows, }, pipelineId) return nil @@ -419,6 +458,20 @@ func (s *Spindle) processPipeline(ctx context.Context, src eventconsumer.Source, return nil } +func (s *Spindle) resolvePipelineRepoDid(repo *tangled.Pipeline_TriggerRepo) (syntax.DID, error) { + if repo.RepoDid == nil || *repo.RepoDid == "" { + return "", fmt.Errorf("pipeline trigger missing repoDid") + } + repoDid, err := syntax.ParseDID(*repo.RepoDid) + if err != nil { + return "", fmt.Errorf("parse repoDid %s: %w", *repo.RepoDid, err) + } + if _, err := s.db.GetRepoByDid(repoDid); err != nil { + s.l.Warn("accepting knot pipeline assertion for unknown repoDid", "repoDid", repoDid, "err", err) + } + return repoDid, nil +} + func (s *Spindle) configureOwner() error { cfgOwner := s.cfg.Server.Owner diff --git a/spindle/tapclient.go b/spindle/tapclient.go new file mode 100644 index 00000000..07ee9e20 --- /dev/null +++ b/spindle/tapclient.go @@ -0,0 +1,348 @@ +package spindle + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "fmt" + "log/slog" + "sync" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" + "tangled.org/core/eventconsumer" + "tangled.org/core/log" + "tangled.org/core/rbac" + "tangled.org/core/spindle/db" + "tangled.org/core/tapc" +) + +const ( + maxPendingPerRepo = 64 + pendingCollabTTL = 10 * time.Minute +) + +type pendingCollabEvent struct { + evt *tapc.RecordEventData + at time.Time +} + +type Tap struct { + logger *slog.Logger + spindle *Spindle + tap tapc.Client + pendingMu sync.Mutex + pendingCollabs map[syntax.DID][]pendingCollabEvent +} + +func NewTapClient(s *Spindle) *Tap { + return &Tap{ + logger: log.SubLogger(s.l, "tapclient"), + spindle: s, + tap: tapc.NewClient(s.cfg.Server.Tap.Url, s.cfg.Server.Tap.AdminPassword), + pendingCollabs: make(map[syntax.DID][]pendingCollabEvent), + } +} + +func (t *Tap) AddOwnerDIDs(ctx context.Context, dids []syntax.DID) error { + if len(dids) == 0 { + return nil + } + return t.tap.AddRepos(ctx, dids) +} + +func (t *Tap) Start(ctx context.Context) { + go t.tap.Connect(ctx, &tapc.SimpleIndexer{ + EventHandler: t.processEvent, + ConnectHandler: t.onConnect, + }) + go t.purgePendingCollabsLoop(ctx) +} + +func (t *Tap) onConnect(ctx context.Context) { + t.spindle.declareTapInterest(ctx) +} + +func (t *Tap) processEvent(ctx context.Context, evt tapc.Event) error { + if evt.Type != tapc.EvtRecord || evt.Record == nil { + return nil + } + switch evt.Record.Collection.String() { + case tangled.RepoNSID: + return t.processRepo(ctx, evt.Record) + case tangled.RepoCollaboratorNSID: + return t.processCollaborator(ctx, evt.Record) + } + return nil +} + +func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error { + l := t.logger.With("collection", tangled.RepoNSID, "did", evt.Did, "rkey", evt.Rkey) + + ownerDid := evt.Did + rkey := evt.Rkey + + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + record := tangled.Repo{} + if err := json.Unmarshal(evt.Record, &record); err != nil { + l.Warn("skipping invalid repo record", "err", err) + return nil + } + + hostname := t.spindle.cfg.Server.Hostname + prior, priorErr := t.spindle.db.GetRepoByOwnerRkey(ownerDid, rkey) + knownRepo := priorErr == nil + + if record.Spindle == nil || *record.Spindle != hostname { + if knownRepo { + l.Info("tearing down repo reassigned from this spindle", "newSpindle", record.Spindle) + return t.teardownRepo(l, prior, ownerDid, rkey) + } + return nil + } + + if record.RepoDid == nil || *record.RepoDid == "" { + l.Warn("skipping repo record without repoDid") + return nil + } + repoDid, err := syntax.ParseDID(*record.RepoDid) + if err != nil { + l.Warn("skipping repo record with malformed repoDid", "value", *record.RepoDid, "err", err) + return nil + } + + if err := t.spindle.e.AddRepo(ownerDid.String(), rbac.ThisServer, repoDid.String()); err != nil { + l.Error("failed to add repo policy", "err", err) + return fmt.Errorf("add repo policy: %w", err) + } + + src := eventconsumer.NewKnotSource(record.Knot) + t.spindle.ks.AddSource(t.spindle.rootCtx, src) + + if err := t.spindle.db.AddRepo(db.Repo{ + Knot: record.Knot, + Owner: ownerDid, + Rkey: rkey, + RepoDid: repoDid, + CreatedAt: record.CreatedAt, + }); err != nil { + l.Error("failed to add repo row", "err", err) + return fmt.Errorf("add repo: %w", err) + } + + if removed, err := t.spindle.db.CollapseRepoSiblings(ownerDid, repoDid); err != nil { + l.Warn("collapse rename siblings failed", "err", err) + } else if removed > 0 { + l.Info("collapsed rename leftovers", "owner", ownerDid, "repo_did", repoDid, "removed", removed) + } + + if err := t.tap.AddRepos(ctx, []syntax.DID{ownerDid}); err != nil { + l.Warn("tap AddRepos rejected", "did", ownerDid, "err", err) + } + + t.drainPendingCollabs(ctx, repoDid) + + case tapc.RecordDeleteAction: + repo, err := t.spindle.db.GetRepoByOwnerRkey(ownerDid, rkey) + if err != nil { + l.Info("skipping delete for unknown repo") + return nil + } + return t.teardownRepo(l, repo, ownerDid, rkey) + } + return nil +} + +func (t *Tap) teardownRepo(l *slog.Logger, repo *db.Repo, ownerDid syntax.DID, rkey syntax.RecordKey) error { + if repo.RepoDid != "" { + collabs, err := t.spindle.db.ListCollaboratorsByRepoDid(repo.RepoDid) + if err != nil { + l.Error("failed to list collaborators for cleanup", "err", err) + return fmt.Errorf("list collaborators: %w", err) + } + for _, c := range collabs { + if err := t.spindle.e.RemoveCollaborator(c.Subject.String(), rbac.ThisServer, repo.RepoDid.String()); err != nil { + l.Error("failed to remove collaborator policy", "subject", c.Subject, "err", err) + return fmt.Errorf("remove collaborator policy: %w", err) + } + } + if err := t.spindle.db.DeleteRepoCollaboratorsByRepoDid(repo.RepoDid); err != nil { + l.Error("failed to clear collaborator rows", "err", err) + return err + } + if err := t.spindle.e.RemoveRepo(ownerDid.String(), rbac.ThisServer, repo.RepoDid.String()); err != nil { + l.Error("failed to remove repo policy", "err", err) + return fmt.Errorf("remove repo policy: %w", err) + } + } + if err := t.spindle.db.DeleteRepoByOwnerRkey(ownerDid, rkey); err != nil { + l.Error("failed to delete repo row", "err", err) + return fmt.Errorf("delete repo row: %w", err) + } + return nil +} + +func (t *Tap) processCollaborator(ctx context.Context, evt *tapc.RecordEventData) error { + l := t.logger.With("collection", tangled.RepoCollaboratorNSID, "did", evt.Did, "rkey", evt.Rkey) + + switch evt.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + record := tangled.RepoCollaborator{} + if err := json.Unmarshal(evt.Record, &record); err != nil { + l.Warn("skipping invalid collaborator record", "err", err) + return nil + } + + actor := evt.Did + rkey := evt.Rkey + + subjectDid, err := syntax.ParseDID(record.Subject) + if err != nil { + l.Info("skipping collaborator with malformed subject DID", "subject", record.Subject, "err", err) + return nil + } + if _, err := t.spindle.res.ResolveIdent(ctx, subjectDid.String()); err != nil { + l.Info("skipping unresolvable collaborator subject", "subject", subjectDid, "err", err) + return nil + } + + repoRefDid, err := syntax.ParseDID(record.Repo) + if err != nil { + l.Info("skipping collaborator with non-DID repo ref", "repo", record.Repo, "err", err) + return nil + } + repo, lookupErr := t.spindle.db.GetRepoByDid(repoRefDid) + if errors.Is(lookupErr, sql.ErrNoRows) { + t.bufferCollab(repoRefDid, evt) + l.Info("buffering collaborator until repo arrives", "repo", repoRefDid) + return nil + } + if lookupErr != nil { + return fmt.Errorf("lookup repo %s: %w", repoRefDid, lookupErr) + } + repoDid := repo.RepoDid + ownerDid := repo.Owner + + if actor != ownerDid { + l.Info("rejecting collaborator with non-owner actor", "actor", actor, "owner", ownerDid) + return nil + } + + ok, err := t.spindle.e.IsCollaboratorInviteAllowed(ownerDid.String(), rbac.ThisServer, repoDid.String()) + if err != nil { + l.Error("invite permission check failed", "err", err) + return fmt.Errorf("invite check: %w", err) + } + if !ok { + l.Info("rejecting collaborator invite", "owner", ownerDid, "repo", repoDid) + return nil + } + + prior, priorErr := t.spindle.db.GetRepoCollaborator(actor, rkey) + staleSubject := priorErr == nil && (prior.Subject != subjectDid || prior.RepoDid != repoDid) + + if err := t.spindle.e.AddCollaborator(subjectDid.String(), rbac.ThisServer, repoDid.String()); err != nil { + l.Error("failed to add collaborator policy", "err", err) + return fmt.Errorf("add collaborator policy: %w", err) + } + if staleSubject { + if err := t.spindle.e.RemoveCollaborator(prior.Subject.String(), rbac.ThisServer, prior.RepoDid.String()); err != nil { + l.Error("failed to remove stale collaborator policy", "err", err) + return fmt.Errorf("remove stale collaborator: %w", err) + } + } + if err := t.spindle.db.AddRepoCollaborator(db.RepoCollaborator{ + OwnerDid: actor, + Rkey: rkey, + Subject: subjectDid, + RepoDid: repoDid, + }); err != nil { + l.Error("failed to persist collaborator row", "err", err) + return fmt.Errorf("track collaborator: %w", err) + } + + case tapc.RecordDeleteAction: + actor := evt.Did + rkey := evt.Rkey + + tracked, err := t.spindle.db.GetRepoCollaborator(actor, rkey) + if err != nil { + l.Info("skipping delete for unknown collaborator record") + return nil + } + if err := t.spindle.e.RemoveCollaborator(tracked.Subject.String(), rbac.ThisServer, tracked.RepoDid.String()); err != nil { + l.Error("failed to remove collaborator policy", "err", err) + return fmt.Errorf("remove collaborator policy: %w", err) + } + if err := t.spindle.db.DeleteRepoCollaborator(actor, rkey); err != nil { + l.Error("failed to delete collaborator row", "err", err) + return fmt.Errorf("delete collaborator row: %w", err) + } + } + return nil +} + +func (t *Tap) bufferCollab(repoDid syntax.DID, evt *tapc.RecordEventData) { + t.pendingMu.Lock() + defer t.pendingMu.Unlock() + list := t.pendingCollabs[repoDid] + list = append(list, pendingCollabEvent{evt: evt, at: time.Now()}) + if len(list) > maxPendingPerRepo { + list = list[len(list)-maxPendingPerRepo:] + } + t.pendingCollabs[repoDid] = list +} + +func (t *Tap) drainPendingCollabs(ctx context.Context, repoDid syntax.DID) { + t.pendingMu.Lock() + list := t.pendingCollabs[repoDid] + delete(t.pendingCollabs, repoDid) + t.pendingMu.Unlock() + if len(list) == 0 { + return + } + cutoff := time.Now().Add(-pendingCollabTTL) + for _, p := range list { + if p.at.Before(cutoff) { + continue + } + if err := t.processCollaborator(ctx, p.evt); err != nil { + t.logger.Warn("replaying buffered collaborator failed", "repo", repoDid, "rkey", p.evt.Rkey, "err", err) + } + } +} + +func (t *Tap) purgePendingCollabsLoop(ctx context.Context) { + ticker := time.NewTicker(pendingCollabTTL / 2) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + t.purgeStalePendingCollabs() + } + } +} + +func (t *Tap) purgeStalePendingCollabs() { + cutoff := time.Now().Add(-pendingCollabTTL) + t.pendingMu.Lock() + defer t.pendingMu.Unlock() + for did, list := range t.pendingCollabs { + kept := list[:0] + for _, p := range list { + if !p.at.Before(cutoff) { + kept = append(kept, p) + } + } + if len(kept) == 0 { + delete(t.pendingCollabs, did) + } else { + t.pendingCollabs[did] = kept + } + } +}