Skip to content

Commit 730027f

Browse files
committed
feat(p2p): add versioned handlers function
1 parent fddc49e commit 730027f

3 files changed

Lines changed: 330 additions & 0 deletions

File tree

pkg/p2p/example_test.go

Lines changed: 148 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,148 @@
1+
// Copyright 2026 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 p2p_test
6+
7+
import (
8+
"context"
9+
"fmt"
10+
11+
"github.com/coreos/go-semver/semver"
12+
"github.com/ethersphere/bee/v2/pkg/p2p"
13+
"github.com/ethersphere/bee/v2/pkg/p2p/protobuf"
14+
"github.com/ethersphere/bee/v2/pkg/swarm"
15+
)
16+
17+
const (
18+
exampleProtocolName = "versionedping"
19+
exampleProtocolVersion = "1.2.0"
20+
exampleStreamName = "ping"
21+
)
22+
23+
// ExampleService represents a full Bee protocol service (structured like pkg/pingpong)
24+
// supporting 3 version levels:
25+
// - v1.2.0: Current version
26+
// - v1.1.0: Legacy version 1.1
27+
// - v1.0.0: Legacy version 1.0
28+
type ExampleService struct {
29+
streamer p2p.Streamer
30+
}
31+
32+
func NewExampleService(streamer p2p.Streamer) *ExampleService {
33+
return &ExampleService{
34+
streamer: streamer,
35+
}
36+
}
37+
38+
func (s *ExampleService) Protocol() p2p.ProtocolSpec {
39+
return p2p.ProtocolSpec{
40+
Name: exampleProtocolName,
41+
Version: exampleProtocolVersion,
42+
StreamSpecs: []p2p.StreamSpec{
43+
{
44+
Name: exampleStreamName,
45+
Handler: p2p.NewVersionedHandlersFunc(
46+
p2p.VersionedHandler{
47+
Version: semver.New("1.2.0"), // Server handler for >= 1.2.0
48+
Handler: func(ctx context.Context, p p2p.Peer, stream p2p.Stream) error {
49+
w, r := protobuf.NewWriterAndReader(stream)
50+
_, _ = w, r
51+
fmt.Println("Server received ping on v1.2.0 handler")
52+
return stream.FullClose()
53+
},
54+
},
55+
p2p.VersionedHandler{
56+
Version: semver.New("1.1.0"), // Server handler for legacy 1.1.0
57+
Handler: func(ctx context.Context, p p2p.Peer, stream p2p.Stream) error {
58+
w, r := protobuf.NewWriterAndReader(stream)
59+
_, _ = w, r
60+
fmt.Println("Server received ping on v1.1.0 legacy handler")
61+
return stream.FullClose()
62+
},
63+
},
64+
p2p.VersionedHandler{
65+
Version: semver.New("1.0.0"), // Server handler for legacy 1.0.0
66+
Handler: func(ctx context.Context, p p2p.Peer, stream p2p.Stream) error {
67+
w, r := protobuf.NewWriterAndReader(stream)
68+
_, _ = w, r
69+
fmt.Println("Server received ping on v1.0.0 legacy handler")
70+
return stream.FullClose()
71+
},
72+
},
73+
),
74+
},
75+
},
76+
}
77+
}
78+
79+
func (s *ExampleService) Ping(ctx context.Context, peer swarm.Address) error {
80+
stream, err := s.streamer.NewStream(ctx, peer, nil, exampleProtocolName, "1.1.0", exampleStreamName)
81+
if err != nil {
82+
return err
83+
}
84+
defer stream.Close()
85+
86+
pingClient := p2p.NewVersionedHandlersFunc(
87+
p2p.VersionedHandler{
88+
Version: semver.New("1.2.0"), // Current version client (>= 1.2.0)
89+
Handler: func(ctx context.Context, _ p2p.Peer, stream p2p.Stream) error {
90+
w, r := protobuf.NewWriterAndReader(stream)
91+
_, _ = w, r
92+
fmt.Println("Client sent ping using v1.2.0 format")
93+
return nil
94+
},
95+
},
96+
p2p.VersionedHandler{
97+
Version: semver.New("1.1.0"), // Legacy v1.1.0 client (1.1.0 <= v < 1.2.0)
98+
Handler: func(ctx context.Context, _ p2p.Peer, stream p2p.Stream) error {
99+
w, r := protobuf.NewWriterAndReader(stream)
100+
_, _ = w, r
101+
fmt.Println("Client sent ping using v1.1.0 legacy format")
102+
return nil
103+
},
104+
},
105+
p2p.VersionedHandler{
106+
Version: semver.New("1.0.0"), // Legacy v1.0.0 client (1.0.0 <= v < 1.1.0)
107+
Handler: func(ctx context.Context, _ p2p.Peer, stream p2p.Stream) error {
108+
w, r := protobuf.NewWriterAndReader(stream)
109+
_, _ = w, r
110+
fmt.Println("Client sent ping using v1.0.0 legacy format")
111+
return nil
112+
},
113+
},
114+
)
115+
116+
return pingClient(ctx, p2p.Peer{Address: peer}, stream)
117+
}
118+
119+
// Example_versionedProtocol demonstrates constructing a versioned P2P protocol service
120+
// and executing versioned message exchange between client and server nodes.
121+
func Example_versionedProtocol() {
122+
ctx := context.Background()
123+
124+
serverSvc := NewExampleService(nil)
125+
serverSpec := serverSvc.Protocol()
126+
serverHandler := serverSpec.StreamSpecs[0].Handler
127+
128+
// Client streamer connects to server where version 1.1.0 is negotiated:
129+
clientStreamer := &mockStreamer{
130+
supportedVersions: map[string]bool{
131+
"1.1.0": true,
132+
},
133+
}
134+
clientSvc := NewExampleService(clientStreamer)
135+
136+
// Client sends ping request (client-side dispatcher automatically selects v1.1.0 format)
137+
_ = clientSvc.Ping(ctx, swarm.ZeroAddress)
138+
139+
// Server receives incoming stream (server-side dispatcher automatically selects v1.1.0 handler)
140+
incomingStream := mockStream{
141+
version: "1.1.0",
142+
}
143+
_ = serverHandler(ctx, p2p.Peer{Address: swarm.ZeroAddress}, incomingStream)
144+
145+
// Output:
146+
// Client sent ping using v1.1.0 legacy format
147+
// Server received ping on v1.1.0 legacy handler
148+
}

pkg/p2p/p2p.go

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ import (
1111
"errors"
1212
"fmt"
1313
"io"
14+
"sort"
1415
"time"
1516

1617
"github.com/coreos/go-semver/semver"
@@ -208,6 +209,40 @@ type HandlerFunc func(context.Context, Peer, Stream) error
208209
// HandlerMiddleware decorates a HandlerFunc by returning a new one.
209210
type HandlerMiddleware func(HandlerFunc) HandlerFunc
210211

212+
// VersionedHandler represents a HandlerFunc associated with a minimum supported Version threshold.
213+
type VersionedHandler struct {
214+
Version *semver.Version
215+
Handler HandlerFunc
216+
}
217+
218+
// NewVersionedHandlersFunc creates a new HandlerFunc that dispatches stream execution
219+
// based on the stream version.
220+
//
221+
// Handlers are evaluated in descending order of Version. The first handler where
222+
// stream.Version >= handler.Version will be executed.
223+
func NewVersionedHandlersFunc(handlers ...VersionedHandler) HandlerFunc {
224+
sorted := make([]VersionedHandler, len(handlers))
225+
copy(sorted, handlers)
226+
sort.Slice(sorted, func(i, j int) bool {
227+
return sorted[j].Version.LessThan(*sorted[i].Version)
228+
})
229+
230+
return func(ctx context.Context, p Peer, stream Stream) error {
231+
v, err := stream.Version()
232+
if err != nil {
233+
return fmt.Errorf("get stream version: %w", err)
234+
}
235+
236+
for _, h := range sorted {
237+
if !v.LessThan(*h.Version) {
238+
return h.Handler(ctx, p, stream)
239+
}
240+
}
241+
242+
return fmt.Errorf("no handler found for stream version: %s", v.String())
243+
}
244+
}
245+
211246
// HeadlerFunc is returning response headers based on the received request
212247
// headers.
213248
type HeadlerFunc func(Headers, swarm.Address) Headers

pkg/p2p/p2p_test.go

Lines changed: 147 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,9 +5,13 @@
55
package p2p_test
66

77
import (
8+
"context"
9+
"errors"
810
"testing"
911

12+
"github.com/coreos/go-semver/semver"
1013
"github.com/ethersphere/bee/v2/pkg/p2p"
14+
"github.com/ethersphere/bee/v2/pkg/swarm"
1115
"github.com/libp2p/go-libp2p/core/network"
1216
)
1317

@@ -36,3 +40,146 @@ func TestReachabilityStatus_String(t *testing.T) {
3640
}
3741
}
3842
}
43+
44+
func TestNewVersionedHandlersFunc(t *testing.T) {
45+
t.Parallel()
46+
47+
var executed string
48+
49+
makeHandler := func(name string) p2p.HandlerFunc {
50+
return func(context.Context, p2p.Peer, p2p.Stream) error {
51+
executed = name
52+
return nil
53+
}
54+
}
55+
56+
// Register handlers in intentionally unordered sequence to test automatic sorting
57+
handlers := []p2p.VersionedHandler{
58+
{Version: semver.New("1.0.0"), Handler: makeHandler("v1.0.0")},
59+
{Version: semver.New("1.2.0"), Handler: makeHandler("v1.2.0")},
60+
{Version: semver.New("1.1.0"), Handler: makeHandler("v1.1.0")},
61+
}
62+
63+
dispatcher := p2p.NewVersionedHandlersFunc(handlers...)
64+
65+
tests := []struct {
66+
name string
67+
streamVersion string
68+
wantExecuted string
69+
wantErr bool
70+
}{
71+
{
72+
name: "exact match for highest version (1.2.0)",
73+
streamVersion: "1.2.0",
74+
wantExecuted: "v1.2.0",
75+
},
76+
{
77+
name: "newer patch version routes to highest version (1.2.5 -> v1.2.0)",
78+
streamVersion: "1.2.5",
79+
wantExecuted: "v1.2.0",
80+
},
81+
{
82+
name: "future minor version routes to highest version (1.3.0 -> v1.2.0)",
83+
streamVersion: "1.3.0",
84+
wantExecuted: "v1.2.0",
85+
},
86+
{
87+
name: "exact match for intermediate version (1.1.0)",
88+
streamVersion: "1.1.0",
89+
wantExecuted: "v1.1.0",
90+
},
91+
{
92+
name: "intermediate patch version (1.1.4 -> v1.1.0)",
93+
streamVersion: "1.1.4",
94+
wantExecuted: "v1.1.0",
95+
},
96+
{
97+
name: "exact match for lowest version (1.0.0)",
98+
streamVersion: "1.0.0",
99+
wantExecuted: "v1.0.0",
100+
},
101+
{
102+
name: "lowest version patch (1.0.9 -> v1.0.0)",
103+
streamVersion: "1.0.9",
104+
wantExecuted: "v1.0.0",
105+
},
106+
{
107+
name: "version below lowest registered version returns error (0.9.0)",
108+
streamVersion: "0.9.0",
109+
wantErr: true,
110+
},
111+
{
112+
name: "error when stream version cannot be retrieved",
113+
streamVersion: "",
114+
wantErr: true,
115+
},
116+
}
117+
118+
for _, tt := range tests {
119+
t.Run(tt.name, func(t *testing.T) {
120+
executed = ""
121+
err := dispatcher(context.Background(), p2p.Peer{}, mockStream{version: tt.streamVersion})
122+
123+
if tt.wantErr {
124+
if err == nil {
125+
t.Fatal("expected error, got nil")
126+
}
127+
if executed != "" {
128+
t.Fatalf("expected no handler to execute, but %q executed", executed)
129+
}
130+
return
131+
}
132+
133+
if err != nil {
134+
t.Fatalf("unexpected error: %v", err)
135+
}
136+
137+
if executed != tt.wantExecuted {
138+
t.Fatalf("executed handler = %q, want %q", executed, tt.wantExecuted)
139+
}
140+
})
141+
}
142+
}
143+
144+
type mockStream struct {
145+
p2p.Stream
146+
version string
147+
closeFn func() error
148+
}
149+
150+
func (m mockStream) Version() (*semver.Version, error) {
151+
if m.version == "" {
152+
return nil, errors.New("missing version")
153+
}
154+
return semver.NewVersion(m.version)
155+
}
156+
157+
func (m mockStream) Close() error {
158+
if m.closeFn != nil {
159+
return m.closeFn()
160+
}
161+
return nil
162+
}
163+
164+
func (m mockStream) FullClose() error {
165+
return m.Close()
166+
}
167+
168+
type mockStreamer struct {
169+
p2p.Streamer
170+
supportedVersions map[string]bool
171+
closed bool
172+
}
173+
174+
func (m *mockStreamer) NewStream(_ context.Context, _ swarm.Address, _ p2p.Headers, _, version, _ string) (p2p.Stream, error) {
175+
if !m.supportedVersions[version] {
176+
return nil, errors.New("protocol version not supported")
177+
}
178+
return mockStream{
179+
version: version,
180+
closeFn: func() error {
181+
m.closed = true
182+
return nil
183+
},
184+
}, nil
185+
}

0 commit comments

Comments
 (0)