Skip to content

Commit fddc49e

Browse files
authored
feat(p2p): expose peer protocol version for backward compatibilities (#5539)
1 parent fc4aa1b commit fddc49e

7 files changed

Lines changed: 173 additions & 4 deletions

File tree

pkg/p2p/libp2p/internal/handshake/mock/stream.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ package mock
77
import (
88
"bytes"
99

10+
"github.com/coreos/go-semver/semver"
1011
"github.com/ethersphere/bee/v2/pkg/p2p"
1112
)
1213

@@ -72,3 +73,7 @@ func (s *Stream) FullClose() error {
7273
func (s *Stream) Reset() error {
7374
return nil
7475
}
76+
77+
func (s *Stream) Version() (*semver.Version, error) {
78+
return nil, nil
79+
}

pkg/p2p/libp2p/stream.go

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,10 @@ package libp2p
77
import (
88
"errors"
99
"io"
10+
"strings"
1011
"time"
1112

13+
"github.com/coreos/go-semver/semver"
1214
"github.com/ethersphere/bee/v2/pkg/p2p"
1315
"github.com/libp2p/go-libp2p/core/network"
1416
)
@@ -38,6 +40,15 @@ func (s *stream) ResponseHeaders() p2p.Headers {
3840
return s.responseHeaders
3941
}
4042

43+
func (s *stream) Version() (*semver.Version, error) {
44+
parts := strings.Split(string(s.Protocol()), "/")
45+
partsLen := len(parts)
46+
if partsLen < 2 {
47+
return nil, errors.New("invalid protocol version")
48+
}
49+
return semver.NewVersion(parts[partsLen-2])
50+
}
51+
4152
func (s *stream) Reset() error {
4253
defer s.metrics.StreamResetCount.Inc()
4354
return s.Stream.Reset()

pkg/p2p/libp2p/stream_test.go

Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,90 @@
1+
// Copyright 2020 The Swarm Authors. All rights reserved.
2+
// Use of this source code is governed by a BSD-style
3+
// license that can be found in the LICENSE file.
4+
5+
package libp2p_test
6+
7+
import (
8+
"testing"
9+
10+
"github.com/coreos/go-semver/semver"
11+
"github.com/ethersphere/bee/v2/pkg/p2p/libp2p"
12+
"github.com/libp2p/go-libp2p/core/network"
13+
"github.com/libp2p/go-libp2p/core/protocol"
14+
)
15+
16+
type mockNetStream struct {
17+
network.Stream
18+
protocol protocol.ID
19+
}
20+
21+
func (m *mockNetStream) Protocol() protocol.ID {
22+
return m.protocol
23+
}
24+
25+
func TestStreamVersion(t *testing.T) {
26+
t.Parallel()
27+
28+
for _, tc := range []struct {
29+
name string
30+
protocolID protocol.ID
31+
wantMajor int64
32+
wantMinor int64
33+
wantPatch int64
34+
wantErr bool
35+
}{
36+
{
37+
name: "valid standard version",
38+
protocolID: "/swarm/pingpong/1.2.3/ping",
39+
wantMajor: 1,
40+
wantMinor: 2,
41+
wantPatch: 3,
42+
},
43+
{
44+
name: "valid rc version",
45+
protocolID: "/swarm/pingpong/2.0.0-rc1/ping",
46+
wantMajor: 2,
47+
wantMinor: 0,
48+
wantPatch: 0,
49+
},
50+
{
51+
name: "invalid version format",
52+
protocolID: "/swarm/pingpong/abc/ping",
53+
wantErr: true,
54+
},
55+
{
56+
name: "too short protocol ID",
57+
protocolID: "/ping",
58+
wantErr: true,
59+
},
60+
} {
61+
t.Run(tc.name, func(t *testing.T) {
62+
t.Parallel()
63+
64+
s := &mockNetStream{protocol: tc.protocolID}
65+
srv := &libp2p.Service{}
66+
wrapped := srv.WrapStream(s)
67+
68+
v, err := wrapped.Version()
69+
if tc.wantErr {
70+
if err == nil {
71+
t.Fatal("expected error, got nil")
72+
}
73+
return
74+
}
75+
76+
if err != nil {
77+
t.Fatalf("unexpected error: %v", err)
78+
}
79+
80+
if v == nil {
81+
t.Fatal("expected version to be non-nil")
82+
}
83+
84+
expected := semver.Version{Major: tc.wantMajor, Minor: tc.wantMinor}
85+
if v.Major != expected.Major || v.Minor != expected.Minor {
86+
t.Errorf("got version %v, want %v", v, expected)
87+
}
88+
})
89+
}
90+
}

pkg/p2p/p2p.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ import (
1313
"io"
1414
"time"
1515

16+
"github.com/coreos/go-semver/semver"
1617
"github.com/ethersphere/bee/v2/pkg/bzz"
1718
"github.com/ethersphere/bee/v2/pkg/swarm"
1819
"github.com/libp2p/go-libp2p/core/network"
@@ -166,6 +167,7 @@ type Stream interface {
166167
Headers() Headers
167168
FullClose() error
168169
Reset() error
170+
Version() (*semver.Version, error)
169171
}
170172

171173
// ProtocolSpec defines a collection of Stream specifications with handlers.

pkg/p2p/protobuf/protobuf_test.go

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ import (
1212
"testing"
1313
"time"
1414

15+
"github.com/coreos/go-semver/semver"
1516
"github.com/ethersphere/bee/v2/pkg/p2p"
1617
"github.com/ethersphere/bee/v2/pkg/p2p/protobuf"
1718
"github.com/ethersphere/bee/v2/pkg/p2p/protobuf/internal/pb"
@@ -347,6 +348,10 @@ func (noopWriteCloser) ResponseHeaders() p2p.Headers {
347348
return nil
348349
}
349350

351+
func (noopWriteCloser) Version() (*semver.Version, error) {
352+
return nil, nil
353+
}
354+
350355
func (noopWriteCloser) Close() error {
351356
return nil
352357
}
@@ -379,6 +384,10 @@ func (noopReadCloser) ResponseHeaders() p2p.Headers {
379384
return nil
380385
}
381386

387+
func (noopReadCloser) Version() (*semver.Version, error) {
388+
return nil, nil
389+
}
390+
382391
func (noopReadCloser) Close() error {
383392
return nil
384393
}

pkg/p2p/streamtest/streamtest.go

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ import (
1212
"testing"
1313
"time"
1414

15+
"github.com/coreos/go-semver/semver"
1516
"github.com/ethersphere/bee/v2/pkg/p2p"
1617
"github.com/ethersphere/bee/v2/pkg/spinlock"
1718
"github.com/ethersphere/bee/v2/pkg/swarm"
@@ -122,10 +123,12 @@ func (r *Recorder) NewStream(ctx context.Context, addr swarm.Address, h p2p.Head
122123
}
123124
}
124125

126+
version, versionErr := semver.NewVersion(protocolVersion)
127+
125128
recordIn := newRecord(r.messageLatency)
126129
recordOut := newRecord(r.messageLatency)
127-
streamOut := newStream(recordIn, recordOut)
128-
streamIn := newStream(recordOut, recordIn)
130+
streamOut := newStream(recordIn, recordOut, version, versionErr)
131+
streamIn := newStream(recordOut, recordIn, version, versionErr)
129132

130133
var handler p2p.HandlerFunc
131134
var headler p2p.HeadlerFunc
@@ -260,10 +263,12 @@ type stream struct {
260263
responseHeaders p2p.Headers
261264
closed bool
262265
lock sync.Mutex
266+
version *semver.Version
267+
versionErr error
263268
}
264269

265-
func newStream(in, out *record) *stream {
266-
return &stream{in: in, out: out}
270+
func newStream(in, out *record, version *semver.Version, versionErr error) *stream {
271+
return &stream{in: in, out: out, version: version, versionErr: versionErr}
267272
}
268273

269274
func (s *stream) Read(p []byte) (int, error) {
@@ -290,6 +295,10 @@ func (s *stream) ResponseHeaders() p2p.Headers {
290295
return s.responseHeaders
291296
}
292297

298+
func (s *stream) Version() (*semver.Version, error) {
299+
return s.version, s.versionErr
300+
}
301+
293302
func (s *stream) Close() error {
294303
s.lock.Lock()
295304
defer s.lock.Unlock()

pkg/p2p/streamtest/streamtest_test.go

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import (
1515
"testing/synctest"
1616
"time"
1717

18+
"github.com/coreos/go-semver/semver"
1819
"github.com/ethersphere/bee/v2/pkg/p2p"
1920
"github.com/ethersphere/bee/v2/pkg/p2p/streamtest"
2021
"github.com/ethersphere/bee/v2/pkg/swarm"
@@ -878,3 +879,45 @@ func testRecords(t *testing.T, records []*streamtest.Record, want [][2]string, w
878879
}
879880
}
880881
}
882+
883+
func TestStreamVersion(t *testing.T) {
884+
t.Parallel()
885+
886+
recorder := streamtest.New(
887+
streamtest.WithProtocols(
888+
newTestProtocol(func(_ context.Context, peer p2p.Peer, stream p2p.Stream) error {
889+
v, err := stream.Version()
890+
if err != nil {
891+
t.Errorf("handler: unexpected error: %v", err)
892+
}
893+
if v == nil || v.String() != "1.0.1" {
894+
t.Errorf("handler: got version %v, want 1.0.1", v)
895+
}
896+
return nil
897+
}),
898+
),
899+
)
900+
901+
stream, err := recorder.NewStream(context.Background(), swarm.ZeroAddress, nil, testProtocolName, testProtocolVersion, testStreamName)
902+
if err != nil {
903+
t.Fatal(err)
904+
}
905+
defer stream.Close()
906+
907+
v, err := stream.Version()
908+
if err != nil {
909+
t.Fatalf("unexpected error: %v", err)
910+
}
911+
912+
if v == nil {
913+
t.Fatal("nil version")
914+
}
915+
916+
if v.String() != "1.0.1" {
917+
t.Fatalf("got string version %v, want 1.0.1", v)
918+
}
919+
920+
if !v.Equal(semver.Version{Major: 1, Minor: 0, Patch: 1}) {
921+
t.Fatalf("got semver version %v, want 1.0.1", v)
922+
}
923+
}

0 commit comments

Comments
 (0)