master
go 336 lines 11.4 KB
Raw
1 package cli
2
3 import (
4 "encoding/json"
5 "testing"
6
7 "github.com/ipfs/kubo/config"
8 "github.com/ipfs/kubo/core/node/libp2p"
9 "github.com/ipfs/kubo/test/cli/harness"
10 "github.com/ipfs/kubo/test/cli/testutils"
11 "github.com/libp2p/go-libp2p/core/peer"
12 "github.com/libp2p/go-libp2p/core/protocol"
13 rcmgr "github.com/libp2p/go-libp2p/p2p/host/resource-manager"
14 "github.com/stretchr/testify/assert"
15 "github.com/stretchr/testify/require"
16 )
17
18 func TestRcmgr(t *testing.T) {
19 t.Parallel()
20
21 t.Run("Resource manager disabled", func(t *testing.T) {
22 t.Parallel()
23 node := harness.NewT(t).NewNode().Init()
24 node.UpdateConfig(func(cfg *config.Config) {
25 cfg.Swarm.ResourceMgr.Enabled = config.False
26 })
27
28 node.StartDaemon()
29 defer node.StopDaemon()
30
31 t.Run("swarm resources should fail", func(t *testing.T) {
32 res := node.RunIPFS("swarm", "resources")
33 assert.Equal(t, 1, res.ExitCode())
34 assert.Contains(t, res.Stderr.String(), "missing ResourceMgr")
35 })
36 })
37
38 t.Run("Node with resource manager disabled", func(t *testing.T) {
39 t.Parallel()
40 node := harness.NewT(t).NewNode().Init()
41 node.UpdateConfig(func(cfg *config.Config) {
42 cfg.Swarm.ResourceMgr.Enabled = config.False
43 })
44 node.StartDaemon()
45 defer node.StopDaemon()
46
47 t.Run("swarm resources should fail", func(t *testing.T) {
48 res := node.RunIPFS("swarm", "resources")
49 assert.Equal(t, 1, res.ExitCode())
50 assert.Contains(t, res.Stderr.String(), "missing ResourceMgr")
51 })
52 })
53
54 t.Run("Very high connmgr highwater", func(t *testing.T) {
55 t.Parallel()
56 node := harness.NewT(t).NewNode().Init()
57 node.UpdateConfig(func(cfg *config.Config) {
58 cfg.Swarm.ConnMgr.HighWater = config.NewOptionalInteger(1000)
59 })
60 node.StartDaemon()
61 defer node.StopDaemon()
62
63 res := node.RunIPFS("swarm", "resources", "--enc=json")
64 require.Equal(t, 0, res.ExitCode())
65 limits := unmarshalLimits(t, res.Stdout.Bytes())
66
67 rl := limits.System.ToResourceLimits()
68 s := rl.Build(rcmgr.BaseLimit{})
69 assert.GreaterOrEqual(t, s.ConnsInbound, 2000)
70 assert.GreaterOrEqual(t, s.StreamsInbound, 2000)
71 })
72
73 t.Run("default configuration", func(t *testing.T) {
74 t.Parallel()
75 node := harness.NewT(t).NewNode().Init()
76 node.UpdateConfig(func(cfg *config.Config) {
77 cfg.Swarm.ConnMgr.HighWater = config.NewOptionalInteger(1000)
78 })
79
80 node.StartDaemon()
81 t.Cleanup(func() { node.StopDaemon() })
82
83 t.Run("conns and streams are above 800 for default connmgr settings", func(t *testing.T) {
84 t.Parallel()
85 res := node.RunIPFS("swarm", "resources", "--enc=json")
86 require.Equal(t, 0, res.ExitCode())
87 limits := unmarshalLimits(t, res.Stdout.Bytes())
88
89 if limits.System.ConnsInbound > rcmgr.DefaultLimit {
90 assert.GreaterOrEqual(t, limits.System.ConnsInbound, 800)
91 }
92 if limits.System.StreamsInbound > rcmgr.DefaultLimit {
93 assert.GreaterOrEqual(t, limits.System.StreamsInbound, 800)
94 }
95 })
96
97 t.Run("limits should succeed", func(t *testing.T) {
98 t.Parallel()
99 res := node.RunIPFS("swarm", "resources", "--enc=json")
100 assert.Equal(t, 0, res.ExitCode())
101
102 limits := rcmgr.PartialLimitConfig{}
103 err := json.Unmarshal(res.Stdout.Bytes(), &limits)
104 require.NoError(t, err)
105
106 assert.NotEqual(t, limits.Transient.Memory, rcmgr.BlockAllLimit64)
107 assert.NotEqual(t, limits.System.Memory, rcmgr.BlockAllLimit64)
108 assert.NotEqual(t, limits.System.FD, rcmgr.BlockAllLimit)
109 assert.NotEqual(t, limits.System.Conns, rcmgr.BlockAllLimit)
110 assert.NotEqual(t, limits.System.ConnsInbound, rcmgr.BlockAllLimit)
111 assert.NotEqual(t, limits.System.ConnsOutbound, rcmgr.BlockAllLimit)
112 assert.NotEqual(t, limits.System.Streams, rcmgr.BlockAllLimit)
113 assert.NotEqual(t, limits.System.StreamsInbound, rcmgr.BlockAllLimit)
114 assert.NotEqual(t, limits.System.StreamsOutbound, rcmgr.BlockAllLimit)
115 })
116
117 t.Run("swarm stats works", func(t *testing.T) {
118 t.Parallel()
119 res := node.RunIPFS("swarm", "resources", "--enc=json")
120 require.Equal(t, 0, res.ExitCode())
121
122 limits := unmarshalLimits(t, res.Stdout.Bytes())
123
124 // every scope has the same fields, so we only inspect system
125 assert.Zero(t, limits.System.MemoryUsage)
126 assert.Zero(t, limits.System.FDUsage)
127 assert.Zero(t, limits.System.ConnsInboundUsage)
128 assert.Zero(t, limits.System.ConnsOutboundUsage)
129 assert.Zero(t, limits.System.StreamsInboundUsage)
130 assert.Zero(t, limits.System.StreamsOutboundUsage)
131 assert.Zero(t, limits.Transient.MemoryUsage)
132 })
133 })
134
135 t.Run("smoke test unlimited System inbounds", func(t *testing.T) {
136 t.Parallel()
137 node := harness.NewT(t).NewNode().Init()
138 node.UpdateUserSuppliedResourceManagerOverrides(func(overrides *rcmgr.PartialLimitConfig) {
139 overrides.System.StreamsInbound = rcmgr.Unlimited
140 overrides.System.ConnsInbound = rcmgr.Unlimited
141 })
142 node.StartDaemon()
143 defer node.StopDaemon()
144
145 res := node.RunIPFS("swarm", "resources", "--enc=json")
146 limits := unmarshalLimits(t, res.Stdout.Bytes())
147
148 assert.Equal(t, rcmgr.Unlimited, limits.System.ConnsInbound)
149 assert.Equal(t, rcmgr.Unlimited, limits.System.StreamsInbound)
150 })
151
152 t.Run("smoke test transient scope", func(t *testing.T) {
153 t.Parallel()
154 node := harness.NewT(t).NewNode().Init()
155 node.UpdateUserSuppliedResourceManagerOverrides(func(overrides *rcmgr.PartialLimitConfig) {
156 overrides.Transient.Memory = 88888
157 })
158 node.StartDaemon()
159 defer node.StopDaemon()
160
161 res := node.RunIPFS("swarm", "resources", "--enc=json")
162 limits := unmarshalLimits(t, res.Stdout.Bytes())
163 assert.Equal(t, rcmgr.LimitVal64(88888), limits.Transient.Memory)
164 })
165
166 t.Run("smoke test service scope", func(t *testing.T) {
167 t.Parallel()
168 node := harness.NewT(t).NewNode().Init()
169 node.UpdateUserSuppliedResourceManagerOverrides(func(overrides *rcmgr.PartialLimitConfig) {
170 overrides.Service = map[string]rcmgr.ResourceLimits{"foo": {Memory: 77777}}
171 })
172 node.StartDaemon()
173 defer node.StopDaemon()
174
175 res := node.RunIPFS("swarm", "resources", "--enc=json")
176 limits := unmarshalLimits(t, res.Stdout.Bytes())
177 assert.Equal(t, rcmgr.LimitVal64(77777), limits.Services["foo"].Memory)
178 })
179
180 t.Run("smoke test protocol scope", func(t *testing.T) {
181 t.Parallel()
182 node := harness.NewT(t).NewNode().Init()
183 node.UpdateUserSuppliedResourceManagerOverrides(func(overrides *rcmgr.PartialLimitConfig) {
184 overrides.Protocol = map[protocol.ID]rcmgr.ResourceLimits{"foo": {Memory: 66666}}
185 })
186 node.StartDaemon()
187 defer node.StopDaemon()
188
189 res := node.RunIPFS("swarm", "resources", "--enc=json")
190 limits := unmarshalLimits(t, res.Stdout.Bytes())
191 assert.Equal(t, rcmgr.LimitVal64(66666), limits.Protocols["foo"].Memory)
192 })
193
194 t.Run("smoke test peer scope", func(t *testing.T) {
195 t.Parallel()
196 validPeerID, err := peer.Decode("QmNnooDu7bfjPFoTZYxMNLWUQJyrVwtbZg5gBMjTezGAJN")
197 assert.NoError(t, err)
198 node := harness.NewT(t).NewNode().Init()
199 node.UpdateUserSuppliedResourceManagerOverrides(func(overrides *rcmgr.PartialLimitConfig) {
200 overrides.Peer = map[peer.ID]rcmgr.ResourceLimits{validPeerID: {Memory: 55555}}
201 })
202 node.StartDaemon()
203 defer node.StopDaemon()
204
205 res := node.RunIPFS("swarm", "resources", "--enc=json")
206 limits := unmarshalLimits(t, res.Stdout.Bytes())
207 assert.Equal(t, rcmgr.LimitVal64(55555), limits.Peers[validPeerID].Memory)
208 })
209
210 t.Run("blocking and allowlists", func(t *testing.T) {
211 t.Parallel()
212 nodes := harness.NewT(t).NewNodes(3).Init()
213 node0, node1, node2 := nodes[0], nodes[1], nodes[2]
214 peerID1, peerID2 := node1.PeerID().String(), node2.PeerID().String()
215
216 node0.UpdateConfig(func(cfg *config.Config) {
217 cfg.Swarm.ResourceMgr.Enabled = config.True
218 cfg.Swarm.ResourceMgr.Allowlist = []string{"/ip4/0.0.0.0/ipcidr/0/p2p/" + peerID2}
219 })
220 node0.UpdateUserSuppliedResourceManagerOverrides(func(overrides *rcmgr.PartialLimitConfig) {
221 *overrides = rcmgr.PartialLimitConfig{
222 System: rcmgr.ResourceLimits{
223 Conns: rcmgr.BlockAllLimit,
224 ConnsInbound: rcmgr.BlockAllLimit,
225 ConnsOutbound: rcmgr.BlockAllLimit,
226 },
227 }
228 })
229
230 nodes.StartDaemons()
231 t.Cleanup(func() { nodes.StopDaemons() })
232
233 t.Run("node 0 should fail to connect to and ping node 1", func(t *testing.T) {
234 t.Parallel()
235 res := node0.Runner.Run(harness.RunRequest{
236 Path: node0.IPFSBin,
237 Args: []string{"swarm", "connect", node1.SwarmAddrsWithPeerIDs()[0].String()},
238 })
239 assert.Equal(t, 1, res.ExitCode())
240 testutils.AssertStringContainsOneOf(t, res.Stderr.String(),
241 "failed to find any peer in table",
242 "resource limit exceeded",
243 )
244
245 res = node0.RunIPFS("ping", "-n2", peerID1)
246 assert.Equal(t, 1, res.ExitCode())
247 assert.Contains(t, res.Stderr.String(), "Error: ping failed")
248 })
249
250 t.Run("node 0 should connect to and ping node 2 since it is allowlisted", func(t *testing.T) {
251 t.Parallel()
252 res := node0.Runner.Run(harness.RunRequest{
253 Path: node0.IPFSBin,
254 Args: []string{"swarm", "connect", node2.SwarmAddrsWithPeerIDs()[0].String()},
255 })
256 assert.Equal(t, 0, res.ExitCode())
257
258 res = node0.RunIPFS("ping", "-n2", peerID2)
259 assert.Equal(t, 0, res.ExitCode())
260 })
261 })
262
263 t.Run("daemon should refuse to start if connmgr.highwater < resources inbound", func(t *testing.T) {
264 t.Run("system conns", func(t *testing.T) {
265 t.Parallel()
266 node := harness.NewT(t).NewNode().Init()
267 node.UpdateConfig(func(cfg *config.Config) {
268 cfg.Swarm.ConnMgr.HighWater = config.NewOptionalInteger(128)
269 cfg.Swarm.ConnMgr.LowWater = config.NewOptionalInteger(64)
270 })
271 node.UpdateUserSuppliedResourceManagerOverrides(func(overrides *rcmgr.PartialLimitConfig) {
272 *overrides = rcmgr.PartialLimitConfig{
273 System: rcmgr.ResourceLimits{Conns: 128},
274 }
275 })
276
277 res := node.RunIPFS("daemon")
278 assert.Equal(t, 1, res.ExitCode())
279 })
280 t.Run("system conns inbound", func(t *testing.T) {
281 t.Parallel()
282 node := harness.NewT(t).NewNode().Init()
283 node.UpdateConfig(func(cfg *config.Config) {
284 cfg.Swarm.ConnMgr.HighWater = config.NewOptionalInteger(128)
285 cfg.Swarm.ConnMgr.LowWater = config.NewOptionalInteger(64)
286 })
287 node.UpdateUserSuppliedResourceManagerOverrides(func(overrides *rcmgr.PartialLimitConfig) {
288 *overrides = rcmgr.PartialLimitConfig{
289 System: rcmgr.ResourceLimits{ConnsInbound: 128},
290 }
291 })
292
293 res := node.RunIPFS("daemon")
294 assert.Equal(t, 1, res.ExitCode())
295 })
296 t.Run("system streams", func(t *testing.T) {
297 t.Parallel()
298 node := harness.NewT(t).NewNode().Init()
299 node.UpdateConfig(func(cfg *config.Config) {
300 cfg.Swarm.ConnMgr.HighWater = config.NewOptionalInteger(128)
301 cfg.Swarm.ConnMgr.LowWater = config.NewOptionalInteger(64)
302 })
303 node.UpdateUserSuppliedResourceManagerOverrides(func(overrides *rcmgr.PartialLimitConfig) {
304 *overrides = rcmgr.PartialLimitConfig{
305 System: rcmgr.ResourceLimits{Streams: 128},
306 }
307 })
308
309 res := node.RunIPFS("daemon")
310 assert.Equal(t, 1, res.ExitCode())
311 })
312 t.Run("system streams inbound", func(t *testing.T) {
313 t.Parallel()
314 node := harness.NewT(t).NewNode().Init()
315 node.UpdateConfig(func(cfg *config.Config) {
316 cfg.Swarm.ConnMgr.HighWater = config.NewOptionalInteger(128)
317 cfg.Swarm.ConnMgr.LowWater = config.NewOptionalInteger(64)
318 })
319 node.UpdateUserSuppliedResourceManagerOverrides(func(overrides *rcmgr.PartialLimitConfig) {
320 *overrides = rcmgr.PartialLimitConfig{
321 System: rcmgr.ResourceLimits{StreamsInbound: 128},
322 }
323 })
324
325 res := node.RunIPFS("daemon")
326 assert.Equal(t, 1, res.ExitCode())
327 })
328 })
329 }
330
331 func unmarshalLimits(t *testing.T, b []byte) *libp2p.LimitsConfigAndUsage {
332 limits := &libp2p.LimitsConfigAndUsage{}
333 err := json.Unmarshal(b, limits)
334 require.NoError(t, err)
335 return limits
336 }