diff --git a/cmd/srv/main.go b/cmd/srv/main.go index f007d89..2e1c8b7 100644 --- a/cmd/srv/main.go +++ b/cmd/srv/main.go @@ -15,12 +15,21 @@ import ( "github.com/libp2p/go-libp2p" "github.com/libp2p/go-libp2p/core/crypto" "github.com/libp2p/go-libp2p/p2p/discovery/mdns" + "github.com/prometheus/client_golang/prometheus" ) const ( ServiceTag = "clset-sync" ) +func init() { + prometheus.MustRegister( + clset.SyncAttempts, clset.SyncSuccesses, clset.SyncDuration, + clset.SummariesPublished, clset.SummariesReceived, clset.SummariesUseful, + clset.CRDTOpDuration, + ) +} + func main() { if len(os.Args) < 2 { log.Fatal("Usage: go run main.go [connect-addr] [http-port]") diff --git a/crdt.go b/crdt.go index ae242a5..c95a433 100644 --- a/crdt.go +++ b/crdt.go @@ -9,6 +9,7 @@ import ( "strconv" "strings" "sync" + "time" "github.com/ipfs/go-datastore" "github.com/ipfs/go-datastore/query" @@ -209,6 +210,8 @@ func (c *CRDT) getEntry(ctx context.Context, reader datastore.Read, key string) } func (c *CRDT) Set(key string, value []byte) error { + start := time.Now() + c.mu.Lock() var afterCommit func() @@ -246,6 +249,8 @@ func (c *CRDT) Set(key string, value []byte) error { afterCommit() } + CRDTOpDuration.WithLabelValues(c.PeerID, "set").Observe(time.Since(start).Seconds()) + return nil } @@ -458,6 +463,8 @@ func (c *CRDT) setDirect(ctx context.Context, ds datastore.Datastore, key string } func (c *CRDT) Delete(key string) error { + start := time.Now() + c.mu.Lock() defer c.mu.Unlock() @@ -491,6 +498,8 @@ func (c *CRDT) Delete(key string) error { afterCommit() } + CRDTOpDuration.WithLabelValues(c.PeerID, "delete").Observe(time.Since(start).Seconds()) + return nil } @@ -658,6 +667,11 @@ func (c *CRDT) deleteDirect(ctx context.Context, ds datastore.Datastore, key str } func (c *CRDT) Get(key string) ([]byte, bool, error) { + start := time.Now() + defer func() { + CRDTOpDuration.WithLabelValues(c.PeerID, "get").Observe(time.Since(start).Seconds()) + }() + ctx := context.Background() entry, err := c.getEntry(ctx, c.ds, key) if err != nil { diff --git a/go.mod b/go.mod index 62aea7e..ce11dc5 100644 --- a/go.mod +++ b/go.mod @@ -25,12 +25,14 @@ require ( github.com/francoispqt/gojay v1.2.13 // indirect github.com/go-logr/logr v1.4.3 // indirect github.com/go-logr/stdr v1.2.2 // indirect + github.com/gogo/protobuf v1.3.2 // indirect github.com/golang/freetype v0.0.0-20170609003504-e2365dfdc4a0 // indirect github.com/google/flatbuffers v25.9.23+incompatible // indirect github.com/google/gopacket v1.1.19 // indirect github.com/google/uuid v1.6.0 // indirect github.com/gorilla/websocket v1.5.3 // indirect github.com/guptarohit/asciigraph v0.7.2 // indirect + github.com/hashicorp/golang-lru/v2 v2.0.7 // indirect github.com/huin/goupnp v1.3.0 // indirect github.com/inconshreveable/mousetrap v1.1.0 // indirect github.com/ipfs/go-cid v0.5.0 // indirect @@ -45,6 +47,7 @@ require ( github.com/libp2p/go-buffer-pool v0.1.0 // indirect github.com/libp2p/go-flow-metrics v0.2.0 // indirect github.com/libp2p/go-libp2p-asn-util v0.4.1 // indirect + github.com/libp2p/go-libp2p-pubsub v0.15.0 // indirect github.com/libp2p/go-msgio v0.3.0 // indirect github.com/libp2p/go-netroute v0.2.2 // indirect github.com/libp2p/go-reuseport v0.4.0 // indirect diff --git a/go.sum b/go.sum index b707687..4e5c323 100644 --- a/go.sum +++ b/go.sum @@ -56,6 +56,8 @@ github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ4 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/gogo/protobuf v1.1.1/go.mod h1:r8qH/GZQm5c6nD/R0oafs1akxWv10x8SbQlK7atdtwQ= +github.com/gogo/protobuf v1.3.2 h1:Ov1cvc58UF3b5XjBnZv7+opcTcQFZebYjWzi34vdm4Q= +github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69NZV8Q= github.com/golang/freetype v0.0.0-20170609003504-e2365dfdc4a0 h1:DACJavvAHhabrF08vX0COfcOBJRhZ8lUbR+ZWIs0Y5g= github.com/golang/freetype v0.0.0-20170609003504-e2365dfdc4a0/go.mod h1:E/TSTwGwJL78qG/PmXZO1EjYhfJinVAhrmmHX6Z8B9k= github.com/golang/glog v0.0.0-20160126235308-23def4e6c14b/go.mod h1:SBH7ygxi8pfUlaOkMMuAQtPIUF8ecWP5IEl/CR7VP2Q= @@ -90,6 +92,8 @@ github.com/gregjones/httpcache v0.0.0-20180305231024-9cad4c3443a7/go.mod h1:Fecb github.com/grpc-ecosystem/grpc-gateway v1.5.0/go.mod h1:RSKVYQBd5MCa4OVpNdGskqpgL2+G+NZTnrVHpWWfpdw= github.com/guptarohit/asciigraph v0.7.2 h1:pBBJYbMl4j7zS4AwmrfAs6tA0VQOEQC933aG72dlrFA= github.com/guptarohit/asciigraph v0.7.2/go.mod h1:dYl5wwK4gNsnFf9Zp+l06rFiDZ5YtXM6x7SRWZ3KGag= +github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k= +github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= github.com/huin/goupnp v1.3.0 h1:UvLUlWDNpoUdYzb2TCn+MuTWtcjXKSza2n6CBdQ0xXc= github.com/huin/goupnp v1.3.0/go.mod h1:gnGPsThkYa7bFi/KWmEysQRf48l2dvR5bxr2OFckNX8= github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8= @@ -113,6 +117,7 @@ github.com/jbenet/go-temp-err-catcher v0.1.0/go.mod h1:0kJRvmDZXNMIiJirNPEYfhpPw github.com/jellevandenhooff/dkim v0.0.0-20150330215556-f50fe3d243e1/go.mod h1:E0B/fFc00Y+Rasa88328GlI/XbtyysCtTHZS8h7IrBU= github.com/json-iterator/go v1.1.6/go.mod h1:+SdeFBvtyEkXs7REEP0seUULqWtbJapLOCVDaaPEHmU= github.com/jstemmer/go-junit-report v0.0.0-20190106144839-af01ea7f8024/go.mod h1:6v2b51hI/fHJwM22ozAgKL4VKDeJcHhJFhtBdhmNjmU= +github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8= github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= 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= @@ -139,6 +144,8 @@ github.com/libp2p/go-libp2p v0.43.0 h1:b2bg2cRNmY4HpLK8VHYQXLX2d3iND95OjodLFymvq github.com/libp2p/go-libp2p v0.43.0/go.mod h1:IiSqAXDyP2sWH+J2gs43pNmB/y4FOi2XQPbsb+8qvzc= 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-pubsub v0.15.0 h1:cG7Cng2BT82WttmPFMi50gDNV+58K626m/wR00vGL1o= +github.com/libp2p/go-libp2p-pubsub v0.15.0/go.mod h1:lr4oE8bFgQaifRcoc2uWhWWiK6tPdOEKpUuR408GFN4= 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= @@ -325,6 +332,8 @@ github.com/wcharczuk/go-chart/v2 v2.1.2/go.mod h1:Zi4hbaqlWpYajnXB2K22IUYVXRXaLf github.com/wlynxg/anet v0.0.3/go.mod h1:eay5PRQr7fIVAMbTbchTnO9gG65Hg/uYGdc7mguHxoA= github.com/wlynxg/anet v0.0.5 h1:J3VJGi1gvo0JwZ/P1/Yc/8p63SoW98B5dHkYDmpgvvU= github.com/wlynxg/anet v0.0.5/go.mod h1:eay5PRQr7fIVAMbTbchTnO9gG65Hg/uYGdc7mguHxoA= +github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= +github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY= go.opencensus.io v0.18.0/go.mod h1:vKdFvxhtzZ9onBp9VKHK8z/sRpBMnKAsufL7wlDrCOA= go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA= @@ -355,6 +364,7 @@ golang.org/x/crypto v0.0.0-20190313024323-a1f597ede03a/go.mod h1:djNgcEr1/C05ACk golang.org/x/crypto v0.0.0-20190611184440-5c40567a22f8/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20200602180216-279210d13fed/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= +golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= golang.org/x/crypto v0.0.0-20210322153248-0c34fe9e7dc2/go.mod h1:T9bdIzuCu7OtxOm1hfPfRQxPLYneinmdGuTeoZ9dtd4= golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc= golang.org/x/crypto v0.8.0/go.mod h1:mRqEX+O9/h5TFCrQhkgjo2yKi0yYA+9ecGkdQoHrywE= @@ -375,6 +385,8 @@ golang.org/x/lint v0.0.0-20181026193005-c67002cb31c3/go.mod h1:UVdnD1Gm6xHRNCYTk golang.org/x/lint v0.0.0-20190227174305-5b3e6a55c961/go.mod h1:wehouNa3lNwaWXcvxsM5YxQ5yQlVC4a0KAMCusXpPoU= golang.org/x/lint v0.0.0-20200302205851-738671d3881b/go.mod h1:3xt1FjdF8hUf6vQPIChWIBhFzV8gjjsPE/fR3IyQdNY= golang.org/x/mod v0.1.1-0.20191105210325-c90efee705ee/go.mod h1:QqPTAvyqsEbceGzBzNggFXnrqF1CaUcvgkdR5Ot7KZg= +golang.org/x/mod v0.2.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= +golang.org/x/mod v0.3.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4= golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= golang.org/x/mod v0.12.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= @@ -392,6 +404,8 @@ golang.org/x/net v0.0.0-20190213061140-3a22650c66bd/go.mod h1:mL1N/T3taQHkDXs73r golang.org/x/net v0.0.0-20190313220215-9f648a60d977/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= +golang.org/x/net v0.0.0-20200226121028-0de0cce0169b/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= +golang.org/x/net v0.0.0-20201021035429-f5854403a974/go.mod h1:sp8m0HH+o8qH0wwXwYZr8TS3Oi6o0r6Gce1SSxlDquU= golang.org/x/net v0.0.0-20210119194325-5f4716e94777/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg= golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg= golang.org/x/net v0.0.0-20210423184538-5f58ad60dda6/go.mod h1:OJAsFXCWl8Ukc7SiCT/9KSuxbyM7479/AVlXFRxuMCk= @@ -416,6 +430,8 @@ golang.org/x/sync v0.0.0-20181108010431-42b317875d0f/go.mod h1:RxMgew5VJxzue5/jJ golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20190227155943-e225da77a7e6/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.0.0-20190911185100-cd5d95a43a6e/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= @@ -431,6 +447,7 @@ golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5h golang.org/x/sys v0.0.0-20190316082340-a2f829d7f35f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20200602225109-6fdc65e7d980/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-20210303074136-134d130e1a04/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= @@ -486,6 +503,8 @@ golang.org/x/tools v0.0.0-20190114222345-bf090417da8b/go.mod h1:n7NCudcB/nEzxVGm golang.org/x/tools v0.0.0-20190226205152-f727befe758c/go.mod h1:9Yl7xja0Znq3iFh3HoIrodX9oNMXvdceNzlUR8zjMvY= golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= golang.org/x/tools v0.0.0-20200130002326-2f3ba24bd6e7/go.mod h1:TB2adYChydJhpapKDTa4BR/hXlZSLoq2Wpct/0txZ28= +golang.org/x/tools v0.0.0-20200619180055-7c47624df98f/go.mod h1:EkVYQZoAsY45+roYkvgYkIh4xh/qjgUK9TdY2XT94GE= +golang.org/x/tools v0.0.0-20210106214847-113979e3529a/go.mod h1:emZCQorbCU4vsT4fOWvOPXz4eW1wZW4PmDk9uLelYpA= golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc= golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU= golang.org/x/tools v0.13.0/go.mod h1:HvlwmtVNQAhOuCjW7xxvovg8wbNq7LwfXh/k7wXUl58= @@ -494,6 +513,8 @@ golang.org/x/tools v0.34.0 h1:qIpSLOxeCYGg9TrcJokLBG4KFA6d795g0xkBkiESGlo= golang.org/x/tools v0.34.0/go.mod h1:pAP9OwEaY1CAW3HOmg3hLZC5Z0CCmzjAF2UQMSqNARg= golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= +golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= +golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= google.golang.org/api v0.0.0-20180910000450-7ca32eb868bf/go.mod h1:4mhQ8q/RsB7i+udVvVy5NUi08OU8ZlA0gRVgrF7VFY0= google.golang.org/api v0.0.0-20181030000543-1d582fd0359e/go.mod h1:4mhQ8q/RsB7i+udVvVy5NUi08OU8ZlA0gRVgrF7VFY0= google.golang.org/api v0.1.0/go.mod h1:UGEZY7KEX120AnNLIHFMKIo4obdJhkp2tPbaPlQx13Y= diff --git a/http_api.go b/http_api.go index cdd87e1..ca5723e 100644 --- a/http_api.go +++ b/http_api.go @@ -6,6 +6,7 @@ import ( "net/http" "github.com/gorilla/mux" + "github.com/prometheus/client_golang/prometheus/promhttp" ) type HTTPServer struct { @@ -28,6 +29,7 @@ func (s *HTTPServer) routes() { s.mux.HandleFunc("/key/{key}", s.handlePutKey).Methods("PUT") s.mux.HandleFunc("/key/{key}", s.handleDeleteKey).Methods("DELETE") s.mux.HandleFunc("/count", s.handleCountKeys).Methods("GET") + s.mux.Handle("/metrics", promhttp.Handler()) } func (s *HTTPServer) Serve(addr string) { diff --git a/metrics.go b/metrics.go new file mode 100644 index 0000000..fe9a297 --- /dev/null +++ b/metrics.go @@ -0,0 +1,41 @@ +package clset + +import "github.com/prometheus/client_golang/prometheus" + +var ( + CRDTOpDuration = prometheus.NewHistogramVec(prometheus.HistogramOpts{ + Name: "clset_op_duration_seconds", + Help: "Duration of CRDT operations by peer", + Buckets: prometheus.DefBuckets, + }, []string{"peer_id", "op"}) + SyncAttempts = prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "clset_sync_attempts_total", + Help: "Total number of sync attempts initiated by peers", + }, []string{"peer_id"}) + + SyncSuccesses = prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "clset_sync_success_total", + Help: "Total number of successful syncs per peer", + }, []string{"peer_id"}) + + SyncDuration = prometheus.NewHistogramVec(prometheus.HistogramOpts{ + Name: "clset_sync_duration_seconds", + Help: "Duration of syncWithPeer operations by peer", + Buckets: prometheus.DefBuckets, + }, []string{"peer_id"}) + + SummariesPublished = prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "clset_summaries_published_total", + Help: "Number of gossip summaries published via pubsub by peer", + }, []string{"peer_id"}) + + SummariesReceived = prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "clset_summaries_received_total", + Help: "Number of gossip summaries received via pubsub by peer", + }, []string{"peer_id"}) + + SummariesUseful = prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "clset_summaries_useful_total", + Help: "Number of received summaries that triggered a sync", + }, []string{"peer_id"}) +) diff --git a/p2p_sync.go b/p2p_sync.go index b811257..b7805a3 100644 --- a/p2p_sync.go +++ b/p2p_sync.go @@ -12,6 +12,7 @@ import ( "sync" "time" + pubsub "github.com/libp2p/go-libp2p-pubsub" "github.com/libp2p/go-libp2p/core/host" "github.com/libp2p/go-libp2p/core/network" "github.com/libp2p/go-libp2p/core/peer" @@ -42,12 +43,46 @@ type Peer struct { peersMu sync.RWMutex syncInProgress map[peer.ID]bool syncMu sync.RWMutex + + ps *pubsub.PubSub + topic *pubsub.Topic + sub *pubsub.Subscription + summaryCache map[string]SummaryMessage + cacheMu sync.Mutex + + summaryCh chan struct{} + lastSummary map[string]uint64 // last advertised tracked state + lastPub time.Time // last publish time +} + +type SummaryMessage struct { + PeerID string `json:"peer_id"` + Tracked map[string]uint64 `json:"tracked"` } // NewP2PSync creates a P2P-enabled CRDT instance func NewPeer(crdt *CRDT, ctx context.Context, host host.Host) (*Peer, error) { syncCtx, cancel := context.WithCancel(ctx) + // Initialize pubsub + ps, err := pubsub.NewGossipSub(syncCtx, host) + if err != nil { + cancel() + return nil, fmt.Errorf("failed to create pubsub: %w", err) + } + + topic, err := ps.Join("/clset/summary/1.0.0") + if err != nil { + cancel() + return nil, fmt.Errorf("failed to join pubsub topic: %w", err) + } + + sub, err := topic.Subscribe() + if err != nil { + cancel() + return nil, fmt.Errorf("failed to subscribe to pubsub topic: %w", err) + } + p2p := &Peer{ crdt: crdt, Host: host, @@ -55,8 +90,20 @@ func NewPeer(crdt *CRDT, ctx context.Context, host host.Host) (*Peer, error) { cancel: cancel, syncPeers: make(map[peer.ID]time.Time), syncInProgress: make(map[peer.ID]bool), + + ps: ps, + topic: topic, + sub: sub, + summaryCache: make(map[string]SummaryMessage), + + summaryCh: make(chan struct{}, 1), + lastSummary: make(map[string]uint64), } + crdt.AddInsertHook(func(k string, v []byte, m CRDTKeyMeta) { p2p.notifySummary() }) + crdt.AddUpdateHook(func(k string, oldV []byte, oldM CRDTKeyMeta, newV []byte, newM CRDTKeyMeta) { p2p.notifySummary() }) + crdt.AddDeleteHook(func(k string, oldV []byte, oldM CRDTKeyMeta) { p2p.notifySummary() }) + // Set stream handler for incoming sync requests host.SetStreamHandler(ProtocolID, p2p.handleSyncStream) @@ -66,6 +113,11 @@ func NewPeer(crdt *CRDT, ctx context.Context, host host.Host) (*Peer, error) { // Start periodic sync with connected peers go p2p.periodicSync() + // Start gossip routines + go p2p.readSummaries() + go p2p.broadcastSummaries() + go p2p.scheduledSync() + log.Printf("P2P CRDT node started") log.Printf("Peer ID: %s", host.ID()) log.Printf("Listening on: %v", host.Addrs()) @@ -170,6 +222,8 @@ func (p *Peer) handleSyncResponse(resp SyncMessage) { // syncWithPeer initiates a sync with a specific peer func (p *Peer) syncWithPeer(peerID peer.ID) error { + SyncAttempts.WithLabelValues(p.crdt.PeerID).Inc() + // Check if sync is already in progress p.syncMu.Lock() if p.syncInProgress[peerID] { @@ -185,6 +239,8 @@ func (p *Peer) syncWithPeer(peerID peer.ID) error { p.syncMu.Unlock() }() + start := time.Now() + // Check if we're still connected to the peer connectedness := p.Host.Network().Connectedness(peerID) if connectedness != network.Connected { @@ -259,12 +315,16 @@ func (p *Peer) syncWithPeer(peerID peer.ID) error { // Handle the response p.handleSyncResponse(resp) + remote := peerID.String() + SyncDuration.WithLabelValues(remote).Observe(time.Since(start).Seconds()) + SyncSuccesses.WithLabelValues(p.crdt.PeerID).Inc() + return nil } // periodicSync runs periodic synchronization with all connected peers func (p *Peer) periodicSync() { - ticker := time.NewTicker(5 * time.Second) + ticker := time.NewTicker(5 * time.Minute) defer ticker.Stop() for { @@ -463,3 +523,125 @@ func (p *Peer) SyncNow() { }(peerID) } } + +func (p *Peer) broadcastSummaries() { + for { + select { + case <-p.ctx.Done(): + return + case <-p.summaryCh: + tracked := p.crdt.GetTrackedPeers() + // Avoid re-sending identical state too often + if equalTracked(tracked, p.lastSummary) && time.Since(p.lastPub) < time.Minute { + continue + } + p.lastSummary = tracked + p.lastPub = time.Now() + + msg := SummaryMessage{ + PeerID: p.Host.ID().String(), + Tracked: tracked, + } + data, err := json.Marshal(msg) + if err != nil { + log.Printf("Failed to marshal summary: %v", err) + } + if err := p.topic.Publish(p.ctx, data); err != nil { + log.Printf("Failed to publish summary: %v", err) + } else { + SummariesPublished.WithLabelValues(p.crdt.PeerID).Inc() + } + } + } +} + +func equalTracked(a, b map[string]uint64) bool { + if len(a) != len(b) { + return false + } + for k, v := range a { + if b[k] != v { + return false + } + } + return true +} + +func (p *Peer) notifySummary() { + select { + case p.summaryCh <- struct{}{}: + default: + // drop if already signalled + } +} + +func (p *Peer) readSummaries() { + for { + m, err := p.sub.Next(p.ctx) + if err != nil { + return + } + if m.ReceivedFrom == p.Host.ID() { + continue + } + + var msg SummaryMessage + if err := json.Unmarshal(m.Data, &msg); err != nil { + continue + } + + SummariesReceived.WithLabelValues(p.crdt.PeerID).Inc() + + p.cacheMu.Lock() + p.summaryCache[msg.PeerID] = msg + p.cacheMu.Unlock() + } +} + +func (p *Peer) scheduledSync() { + ticker := time.NewTicker(5 * time.Second) + defer ticker.Stop() + + for { + select { + case <-p.ctx.Done(): + return + case <-ticker.C: + p.cacheMu.Lock() + cacheCopy := make(map[string]SummaryMessage) + for k, v := range p.summaryCache { + cacheCopy[k] = v + } + p.cacheMu.Unlock() + + local := p.crdt.GetTrackedPeers() + bestPeer := "" + maxGap := uint64(0) + + for pid, summary := range cacheCopy { + gap := calcGap(local, summary.Tracked) + if gap > maxGap { + maxGap = gap + bestPeer = pid + } + } + + if bestPeer != "" { + if pid, err := peer.Decode(bestPeer); err == nil { + SummariesUseful.WithLabelValues(p.crdt.PeerID).Inc() + go p.syncWithPeer(pid) + } + } + } + } +} + +func calcGap(local, remote map[string]uint64) uint64 { + var gap uint64 + for pid, seq := range remote { + if local[pid] < seq { + gap += seq - local[pid] + } + } + return gap +}