Skip to content

Commit 3e77601

Browse files
authored
Merge pull request #38 from threefoldtech/update-modules-on-config-update
Update modules on config update
2 parents decf99f + b41b6f5 commit 3e77601

8 files changed

Lines changed: 105 additions & 47 deletions

File tree

go.mod

Lines changed: 6 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -43,12 +43,12 @@ require (
4343
github.com/patrickmn/go-cache v2.1.0+incompatible
4444
github.com/pion/stun v0.6.1
4545
github.com/pkg/errors v0.9.1
46-
github.com/rs/zerolog v1.33.0
46+
github.com/rs/zerolog v1.34.0
4747
github.com/shirou/gopsutil v3.21.11+incompatible
4848
github.com/stretchr/testify v1.10.0
4949
github.com/threefoldtech/0-fs v1.3.1-0.20240424140157-b488dfedcc56
5050
github.com/threefoldtech/tfchain/clients/tfchain-client-go v0.0.0-20241127100051-77e684bcb1b2
51-
github.com/threefoldtech/tfgrid-sdk-go/rmb-sdk-go v0.16.1-0.20241229121208-76ac3fea5e67
51+
github.com/threefoldtech/tfgrid-sdk-go/rmb-sdk-go v0.16.8
5252
github.com/threefoldtech/zbus v1.0.1
5353
github.com/tyler-smith/go-bip39 v1.1.0
5454
github.com/vishvananda/netlink v1.1.1-0.20201029203352-d40f9887b852
@@ -75,7 +75,6 @@ require (
7575
github.com/rogpeppe/go-internal v1.14.1 // indirect
7676
github.com/yuin/gopher-lua v1.1.1 // indirect
7777
go.opentelemetry.io/otel v1.34.0 // indirect
78-
golang.org/x/exp v0.0.0-20240909161429-701f63a606c0 // indirect
7978
golang.org/x/tools v0.30.0 // indirect
8079
google.golang.org/genproto/googleapis/rpc v0.0.0-20250127172529-29210b9bc287 // indirect
8180
)
@@ -164,13 +163,13 @@ require (
164163
go.opencensus.io v0.24.0 // indirect
165164
go.uber.org/atomic v1.9.0 // indirect
166165
golang.org/x/net v0.35.0 // indirect
167-
golang.org/x/sync v0.11.0 // indirect
168-
golang.org/x/text v0.22.0 // indirect
166+
golang.org/x/sync v0.12.0 // indirect
167+
golang.org/x/text v0.23.0 // indirect
169168
golang.zx2c4.com/wireguard v0.0.20200320 // indirect
170-
gonum.org/v1/gonum v0.15.0 // indirect
169+
gonum.org/v1/gonum v0.16.0 // indirect
171170
google.golang.org/appengine v1.6.7 // indirect
172171
google.golang.org/grpc v1.70.0 // indirect
173-
google.golang.org/protobuf v1.36.4 // indirect
172+
google.golang.org/protobuf v1.36.6 // indirect
174173
gopkg.in/djherbis/times.v1 v1.2.0 // indirect
175174
gopkg.in/natefinch/npipe.v2 v2.0.0-20160621034901-c1b8fa8bdcce // indirect
176175
gopkg.in/yaml.v3 v3.0.1 // indirect

go.sum

Lines changed: 15 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -519,15 +519,15 @@ github.com/rogpeppe/fastuuid v0.0.0-20150106093220-6724a57986af/go.mod h1:XWv6So
519519
github.com/rogpeppe/go-internal v1.3.0/go.mod h1:M8bDsm7K2OlrFYOpmOWEs/qY81heoFRclV5y23lUDJ4=
520520
github.com/rogpeppe/go-internal v1.6.1/go.mod h1:xXDCJY+GAPziupqXw64V24skbSoqbTEfhy4qGm1nDQc=
521521
github.com/rogpeppe/go-internal v1.8.1/go.mod h1:JeRgkft04UBgHMgCIwADu4Pn6Mtm5d4nPKWu0nJ5d+o=
522-
github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ=
523-
github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc=
522+
github.com/rogpeppe/go-internal v1.11.0 h1:cWPaGQEPrBb5/AsnsZesgZZ9yb1OQ+GOISoDNXVBh4M=
523+
github.com/rogpeppe/go-internal v1.11.0/go.mod h1:ddIwULY96R17DhadqLgMfk9H9tvdUzkipdSkR5nkCZA=
524524
github.com/rs/cors v1.10.1 h1:L0uuZVXIKlI1SShY2nhFfo44TYvDPQ1w4oFkUJNfhyo=
525525
github.com/rs/cors v1.10.1/go.mod h1:XyqrcTp5zjWr1wsJ8PIRZssZ8b/WMcMf71DJnit4EMU=
526526
github.com/rs/xid v1.2.1/go.mod h1:+uKXf+4Djp6Md1KODXJxgGQPKngRmWyn10oCKFzNHOQ=
527-
github.com/rs/xid v1.5.0/go.mod h1:trrq9SKmegXys3aeAKXMUTdJsYXVwGY3RLcfgqegfbg=
527+
github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0=
528528
github.com/rs/zerolog v1.14.3/go.mod h1:3WXPzbXEEliJ+a6UFE4vhIxV8qR1EML6ngzP9ug4eYg=
529-
github.com/rs/zerolog v1.33.0 h1:1cU2KZkvPxNyfgEmhHAz/1A9Bz+llsdYzklWFzgp0r8=
530-
github.com/rs/zerolog v1.33.0/go.mod h1:/7mN4D5sKwJLZQ2b/znpjC3/GQWY/xaDXUM0kKWRHss=
529+
github.com/rs/zerolog v1.34.0 h1:k43nTLIwcTVQAncfCw4KZ2VY6ukYoZaBPNOE8txlOeY=
530+
github.com/rs/zerolog v1.34.0/go.mod h1:bJsvje4Z08ROH4Nhs5iH600c3IkWhwp44iRc54W6wYQ=
531531
github.com/russross/blackfriday/v2 v2.0.1/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM=
532532
github.com/safchain/ethtool v0.0.0-20190326074333-42ed695e3de8/go.mod h1:Z0q5wiBQGYcxhMZ6gUqHn6pYNLypFAvaL3UvgZLR0U4=
533533
github.com/safchain/ethtool v0.0.0-20201023143004-874930cb3ce0 h1:eskphjc5kRCykOJyX7HHVbJCs25/8knprttvrVvEd8o=
@@ -576,8 +576,8 @@ github.com/threefoldtech/0-fs v1.3.1-0.20240424140157-b488dfedcc56 h1:uWd8JfE8N3
576576
github.com/threefoldtech/0-fs v1.3.1-0.20240424140157-b488dfedcc56/go.mod h1:lZjR32SiNo3dP70inVFxaLMyZjmKX1ucS+5O31dbPNM=
577577
github.com/threefoldtech/tfchain/clients/tfchain-client-go v0.0.0-20241127100051-77e684bcb1b2 h1:VW2J36F8g/kJn4IkY0JiRFmb1gFcdjiOyltfJLJ0mYU=
578578
github.com/threefoldtech/tfchain/clients/tfchain-client-go v0.0.0-20241127100051-77e684bcb1b2/go.mod h1:cOL5YgHUmDG5SAXrsZxFjUECRQQuAqOoqvXhZG5sEUw=
579-
github.com/threefoldtech/tfgrid-sdk-go/rmb-sdk-go v0.16.1-0.20241229121208-76ac3fea5e67 h1:Ii9TmXPBC1GYxRirReSygRZvEGXfAsQRaIipMEzGik0=
580-
github.com/threefoldtech/tfgrid-sdk-go/rmb-sdk-go v0.16.1-0.20241229121208-76ac3fea5e67/go.mod h1:93SROfr+QjgaJ5/jIWtIpLkhaD8Pv8WbdfwvwMNG2p4=
579+
github.com/threefoldtech/tfgrid-sdk-go/rmb-sdk-go v0.16.8 h1:BDuus/zqEBDsmPQA0h3inJyPnM83O6l4oMe6hqn00xA=
580+
github.com/threefoldtech/tfgrid-sdk-go/rmb-sdk-go v0.16.8/go.mod h1:93SROfr+QjgaJ5/jIWtIpLkhaD8Pv8WbdfwvwMNG2p4=
581581
github.com/threefoldtech/zbus v1.0.1 h1:3KaEpyOiDYAw+lrAyoQUGIvY9BcjVRXlQ1beBRqhRNk=
582582
github.com/threefoldtech/zbus v1.0.1/go.mod h1:E/v/xEvG/l6z/Oj0aDkuSUXFm/1RVJkhKBwDTAIdsHo=
583583
github.com/tinylib/msgp v1.1.5 h1:2gXmtWueD2HefZHQe1QOy9HVzmFrLOVvsXwXBQ0ayy0=
@@ -672,8 +672,6 @@ golang.org/x/crypto v0.8.0/go.mod h1:mRqEX+O9/h5TFCrQhkgjo2yKi0yYA+9ecGkdQoHrywE
672672
golang.org/x/crypto v0.33.0 h1:IOBPskki6Lysi0lo9qQvbxiQ+FvsCC/YWOecCHAixus=
673673
golang.org/x/crypto v0.33.0/go.mod h1:bVdXmD7IV/4GdElGPozy6U7lWdRXA4qyRVGJV57uQ5M=
674674
golang.org/x/exp v0.0.0-20190121172915-509febef88a4/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA=
675-
golang.org/x/exp v0.0.0-20240909161429-701f63a606c0 h1:e66Fs6Z+fZTbFBAxKfP3PALWBtpfqks2bwGcexMxgtk=
676-
golang.org/x/exp v0.0.0-20240909161429-701f63a606c0/go.mod h1:2TbTHSBQa924w8M6Xs1QcRcFwyucIwBGpK1p2f1YFFY=
677675
golang.org/x/lint v0.0.0-20181026193005-c67002cb31c3/go.mod h1:UVdnD1Gm6xHRNCYTkRU2/jEulfH38KcIWyp/GAMgvoE=
678676
golang.org/x/lint v0.0.0-20190227174305-5b3e6a55c961/go.mod h1:wehouNa3lNwaWXcvxsM5YxQ5yQlVC4a0KAMCusXpPoU=
679677
golang.org/x/lint v0.0.0-20190313153728-d0100b6bd8b3/go.mod h1:6SW0HCj/g11FgYtHlgUYUwCkIfeOF89ocIRzGO/8vkc=
@@ -734,10 +732,8 @@ golang.org/x/sync v0.0.0-20200625203802-6e8e738ad208/go.mod h1:RxMgew5VJxzue5/jJ
734732
golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
735733
golang.org/x/sync v0.0.0-20201207232520-09787c993a3a/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
736734
golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
737-
golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
738-
golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
739-
golang.org/x/sync v0.11.0 h1:GGz8+XQP4FvTTrjZPzNKTMFtSXH80RAzG+5ghFPgK9w=
740-
golang.org/x/sync v0.11.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
735+
golang.org/x/sync v0.12.0 h1:MHc5BpPuC30uJk597Ri8TV3CNZcTLu6B6z4lJy+g6Jw=
736+
golang.org/x/sync v0.12.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA=
741737
golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
742738
golang.org/x/sys v0.0.0-20180905080454-ebe1bf3edb33/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
743739
golang.org/x/sys v0.0.0-20180909124046-d0be0721c37e/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
@@ -815,11 +811,8 @@ golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk=
815811
golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
816812
golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
817813
golang.org/x/text v0.3.7-0.20210503195748-5c7c50ebbd4f/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ=
818-
golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ=
819-
golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8=
820-
golang.org/x/text v0.9.0/go.mod h1:e1OnstbJyHTd6l/uOt8jFFHp6TRDWZR/bV3emEE/zU8=
821-
golang.org/x/text v0.22.0 h1:bofq7m3/HAFvbF51jz3Q9wLg3jkvSPuiZu/pD1XwgtM=
822-
golang.org/x/text v0.22.0/go.mod h1:YRoo4H8PVmsu+E3Ou7cqLVH8oXWIHVoX0jqUWALQhfY=
814+
golang.org/x/text v0.23.0 h1:D71I7dUrlY+VX0gQShAThNGHFxZ13dGLBHQLVl1mJlY=
815+
golang.org/x/text v0.23.0/go.mod h1:/BLNzu4aZCJ1+kcD0DNRotWKage4q2rGVAg4o22unh4=
823816
golang.org/x/time v0.0.0-20190308202827-9d24e82272b4/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ=
824817
golang.org/x/time v0.0.0-20191024005414-555d28b269f0/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ=
825818
golang.org/x/tools v0.0.0-20180221164845-07fd8470d635/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
@@ -856,8 +849,8 @@ golang.zx2c4.com/wireguard v0.0.20200320/go.mod h1:lDian4Sw4poJ04SgHh35nzMVwGSYl
856849
golang.zx2c4.com/wireguard/wgctrl v0.0.0-20200609130330-bd2cb7843e1b h1:l4mBVCYinjzZuR5DtxHuBD6wyd4348TGiavJ5vLrhEc=
857850
golang.zx2c4.com/wireguard/wgctrl v0.0.0-20200609130330-bd2cb7843e1b/go.mod h1:UdS9frhv65KTfwxME1xE8+rHYoFpbm36gOud1GhBe9c=
858851
golang.zx2c4.com/wireguard/windows v0.3.14/go.mod h1:3P4IEAsb+BjlKZmpUXgy74c0iX9AVwwr3WcVJ8nPgME=
859-
gonum.org/v1/gonum v0.15.0 h1:2lYxjRbTYyxkJxlhC+LvJIx3SsANPdRybu1tGj9/OrQ=
860-
gonum.org/v1/gonum v0.15.0/go.mod h1:xzZVBJBtS+Mz4q0Yl2LJTk+OxOg4jiXZ7qBoM0uISGo=
852+
gonum.org/v1/gonum v0.16.0 h1:5+ul4Swaf3ESvrOnidPp4GZbzf0mxVQpDCYUQE7OJfk=
853+
gonum.org/v1/gonum v0.16.0/go.mod h1:fef3am4MQ93R2HHpKnLk4/Tbh/s0+wqD5nfa6Pnwy4E=
861854
google.golang.org/appengine v1.1.0/go.mod h1:EbEs0AVv82hx2wNQdGPgUI5lhzA/G0D9YwlJXL52JkM=
862855
google.golang.org/appengine v1.4.0/go.mod h1:xpcJRLb0r/rnEns0DIKYYv+WjYCduHsrkT7/EB5XEv4=
863856
google.golang.org/appengine v1.5.0/go.mod h1:xpcJRLb0r/rnEns0DIKYYv+WjYCduHsrkT7/EB5XEv4=
@@ -891,8 +884,8 @@ google.golang.org/protobuf v1.23.1-0.20200526195155-81db48ad09cc/go.mod h1:EGpAD
891884
google.golang.org/protobuf v1.25.0/go.mod h1:9JNX74DMeImyA3h4bdi1ymwjUzf21/xIlbajtzgsN7c=
892885
google.golang.org/protobuf v1.26.0-rc.1/go.mod h1:jlhhOSvTdKEhbULTjvd4ARK9grFBp09yW+WbY/TyQbw=
893886
google.golang.org/protobuf v1.27.1/go.mod h1:9q0QmTI4eRPtz6boOQmLYwt+qCgq0jsYwAQnmE0givc=
894-
google.golang.org/protobuf v1.36.4 h1:6A3ZDJHn/eNqc1i+IdefRzy/9PokBTPvcqMySR7NNIM=
895-
google.golang.org/protobuf v1.36.4/go.mod h1:9fA7Ob0pmnwhb644+1+CVWFRbNajQ6iRojtC/QF5bRE=
887+
google.golang.org/protobuf v1.36.6 h1:z1NpPI8ku2WgiWnf+t9wTPsn6eP1L7ksHUlkfLvd9xY=
888+
google.golang.org/protobuf v1.36.6/go.mod h1:jduwjTPXsFjZGTmRluh+L6NjiWu7pchiJ2/5YcXBHnY=
896889
gopkg.in/airbrake/gobrake.v2 v2.0.9/go.mod h1:/h5ZAUhDkGaJfjzjKLSjv6zCL6O0LLBxU4K+aSYdM/U=
897890
gopkg.in/alecthomas/kingpin.v2 v2.2.6/go.mod h1:FMv+mEhP44yOT+4EoQTLFTRgOQ1FBLkstjWtayDeSgw=
898891
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=

pkg/api_gateway.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import (
1010
//go:generate zbusc -module api-gateway -version 0.0.1 -name api-gateway -package stubs github.com/threefoldtech/zosbase/pkg+SubstrateGateway stubs/api_gateway_stub.go
1111

1212
type SubstrateGateway interface {
13+
UpdateSubstrateGatewayConnection(manager substrate.Manager) (err error)
1314
CreateNode(node substrate.Node) (uint32, error)
1415
CreateTwin(relay string, pk []byte) (uint32, error)
1516
EnsureAccount(activationURL []string, termsAndConditionsLink string, termsAndConditionsHash string) (info substrate.AccountInfo, err error)

pkg/environment/environment.go

Lines changed: 25 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,6 @@ import (
55
"os"
66
"slices"
77
"strconv"
8-
"sync"
98

109
"github.com/pkg/errors"
1110
"github.com/rs/zerolog/log"
@@ -126,8 +125,8 @@ const (
126125
)
127126

128127
var (
129-
pool substrate.Manager
130-
poolOnce sync.Once
128+
pool substrate.Manager
129+
subURLs []string
131130

132131
envDev = Environment{
133132
RunningMode: RunningDev,
@@ -285,10 +284,31 @@ func GetSubstrate() (substrate.Manager, error) {
285284
if err != nil {
286285
return nil, errors.Wrap(err, "failed to get boot environment")
287286
}
287+
updatedSubURLs := env.SubstrateURL
288288

289-
poolOnce.Do(func() {
289+
slices.Sort(subURLs)
290+
slices.Sort(updatedSubURLs)
291+
292+
// if substrate url changed then update subURLs and update pool with new manager only if the old connection is broken
293+
if !slices.Equal(subURLs, updatedSubURLs) {
294+
// before attempting to update the manager check if pool variable maintain a healthy connection
295+
// pool.Row() checks the health of the connection and if all the urls used in pool are down, then it will return error
296+
if pool != nil {
297+
cl, _, err := pool.Raw()
298+
if err == nil {
299+
cl.Client.Close()
300+
return pool, nil
301+
}
302+
}
303+
304+
log.Debug().Strs("substrate_urls", updatedSubURLs).Msg("updating to sub manager with url")
290305
pool = substrate.NewManager(env.SubstrateURL...)
291-
})
306+
subURLs = updatedSubURLs
307+
}
308+
309+
// poolOnce.Do(func() {
310+
// pool = substrate.NewManager(env.SubstrateURL...)
311+
// })
292312

293313
return pool, nil
294314
}

pkg/events/events.go

Lines changed: 16 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -66,17 +66,19 @@ type Callback func(events *substrate.EventRecords)
6666
// Events processor receives all events starting from the given state
6767
// and for each set of events calls callback cb
6868
type Processor struct {
69-
sub substrate.Manager
69+
sub substrate.Manager
70+
update chan substrate.Manager
7071

7172
cb Callback
7273
state State
7374
}
7475

7576
func NewProcessor(sub substrate.Manager, cb Callback, state State) *Processor {
7677
return &Processor{
77-
sub: sub,
78-
cb: cb,
79-
state: state,
78+
update: make(chan substrate.Manager, 0),
79+
sub: sub,
80+
cb: cb,
81+
state: state,
8082
}
8183
}
8284

@@ -97,6 +99,7 @@ func (e *Processor) process(changes []types.StorageChangeSet, meta *types.Metada
9799
}
98100
}
99101
}
102+
100103
func (e *Processor) eventsTo(cl *gsrpc.SubstrateAPI, meta *types.Metadata, block types.Header) error {
101104
//
102105
last, err := e.state.Get(cl)
@@ -120,7 +123,7 @@ func (e *Processor) eventsTo(cl *gsrpc.SubstrateAPI, meta *types.Metadata, block
120123
return errors.Wrapf(err, "failed to get block hash '%d'", start)
121124
}
122125

123-
//state.ErrUnknownBlock
126+
// state.ErrUnknownBlock
124127
changes, err := cl.RPC.State.QueryStorageAt([]types.StorageKey{key}, hash)
125128
if err, ok := err.(rpc.Error); ok {
126129
if err.ErrorCode() == -32000 { // block is too old not in archive anymore
@@ -158,6 +161,9 @@ func (e *Processor) subscribe(ctx context.Context) error {
158161
case <-ctx.Done():
159162
return nil
160163
case err := <-sub.Err():
164+
// sub.Err returns error if the connection is closed unexpectedly
165+
// so if sub manager is updated by now, then retrying to subscribe after 10 seconds will use the new manager
166+
// with the updated urls
161167
return err
162168
case block := <-sub.Chan():
163169
err := e.eventsTo(cl, meta, block)
@@ -169,6 +175,11 @@ func (e *Processor) subscribe(ctx context.Context) error {
169175
if err := e.state.Set(block.Number); err != nil {
170176
return errors.Wrap(err, "failed to commit last block number")
171177
}
178+
case newSub := <-e.update:
179+
// if new update is received then the connection is broken as noded only issues new update only if
180+
// the old manager is broken so we need to resubscribe with the new manager.
181+
e.sub = newSub
182+
return errors.New("failed to listen to substrate events")
172183
}
173184
}
174185
}

pkg/events/redis.go

Lines changed: 12 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -29,11 +29,12 @@ const (
2929
)
3030

3131
type RedisStream struct {
32-
sub substrate.Manager
33-
state string
34-
farm pkg.FarmID
35-
node uint32
36-
pool *redis.Pool
32+
sub substrate.Manager
33+
state string
34+
farm pkg.FarmID
35+
node uint32
36+
pool *redis.Pool
37+
processor *Processor
3738
}
3839

3940
func NewRedisStream(sub substrate.Manager, address string, farm pkg.FarmID, node uint32, state string) (*RedisStream, error) {
@@ -51,6 +52,11 @@ func NewRedisStream(sub substrate.Manager, address string, farm pkg.FarmID, node
5152
}, nil
5253
}
5354

55+
func (r *RedisStream) UpdateSubstrateManager(sub substrate.Manager) {
56+
r.sub = sub
57+
r.processor.update <- sub
58+
}
59+
5460
func (r *RedisStream) push(con redis.Conn, queue string, event interface{}) error {
5561
var buffer bytes.Buffer
5662
enc := gob.NewEncoder(&buffer)
@@ -139,11 +145,11 @@ func (r *RedisStream) process(events *substrate.EventRecords) {
139145
log.Error().Err(err).Msg("failed to push event")
140146
}
141147
}
142-
143148
}
144149

145150
func (r *RedisStream) Start(ctx context.Context) {
146151
ps := NewProcessor(r.sub, r.process, NewFileState(r.state))
152+
r.processor = ps
147153
ps.Start(ctx)
148154
}
149155

@@ -184,7 +190,6 @@ func (r *RedisConsumer) pop(con redis.Conn, group, stream string) ([]payload, er
184190
"BLOCK", 0,
185191
"STREAMS", stream,
186192
0))
187-
188193
if err != nil {
189194
return nil, err
190195
}
@@ -202,7 +207,6 @@ func (r *RedisConsumer) pop(con redis.Conn, group, stream string) ([]payload, er
202207
"BLOCK", 3000,
203208
"STREAMS", stream,
204209
">"))
205-
206210
if err != nil {
207211
return nil, err
208212
}

pkg/stubs/api_gateway_stub.go

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -387,3 +387,18 @@ func (s *SubstrateGatewayStub) UpdateNodeUptimeV2(ctx context.Context, arg0 uint
387387
}
388388
return
389389
}
390+
391+
func (s *SubstrateGatewayStub) UpdateSubstrateGatewayConnection(ctx context.Context, arg0 tfchainclientgo.Manager) (ret0 error) {
392+
args := []interface{}{arg0}
393+
result, err := s.client.RequestContext(ctx, s.module, s.object, "UpdateSubstrateGatewayConnection", args...)
394+
if err != nil {
395+
panic(err)
396+
}
397+
result.PanicOnError()
398+
ret0 = result.CallError()
399+
loader := zbus.Loader{}
400+
if err := result.Unmarshal(&loader); err != nil {
401+
panic(err)
402+
}
403+
return
404+
}

pkg/substrate_gateway/substrate_gateway.go

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,20 @@ func NewSubstrateGateway(manager substrate.Manager, identity substrate.Identity)
3232
return gw, nil
3333
}
3434

35+
// UpdateSubstrateGatewayConnection allow modules to update substrate manager so that the node can recover chain outage
36+
func (g *substrateGateway) UpdateSubstrateGatewayConnection(manager substrate.Manager) error {
37+
sub, err := manager.Substrate()
38+
if err != nil {
39+
return err
40+
}
41+
42+
// close the old connection
43+
g.sub.Close()
44+
45+
g.sub = sub
46+
return nil
47+
}
48+
3549
// createBackoff creates an exponential backoff configuration for substrate calls
3650
func createBackoff() backoff.BackOff {
3751
exp := backoff.NewExponentialBackOff()
@@ -415,6 +429,7 @@ func (g *substrateGateway) UpdateNodeUptimeV2(uptime uint64, timestampHint uint6
415429

416430
return
417431
}
432+
418433
func (g *substrateGateway) GetTime() (time.Time, error) {
419434
log.Trace().Str("method", "Time").Msg("method called")
420435

0 commit comments

Comments
 (0)