Skip to content

Commit f36f1f7

Browse files
committed
chore: add metrics and reduce bloat by moving in dedicated package
1 parent 45afd0c commit f36f1f7

6 files changed

Lines changed: 450 additions & 330 deletions

File tree

pkg/p2p/example_versioned_test.go

Lines changed: 0 additions & 148 deletions
This file was deleted.

pkg/p2p/p2p.go

Lines changed: 1 addition & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -11,10 +11,10 @@ import (
1111
"errors"
1212
"fmt"
1313
"io"
14-
"sort"
1514
"time"
1615

1716
"github.com/coreos/go-semver/semver"
17+
1818
"github.com/ethersphere/bee/v2/pkg/bzz"
1919
"github.com/ethersphere/bee/v2/pkg/swarm"
2020
"github.com/libp2p/go-libp2p/core/network"
@@ -209,40 +209,6 @@ type HandlerFunc func(context.Context, Peer, Stream) error
209209
// HandlerMiddleware decorates a HandlerFunc by returning a new one.
210210
type HandlerMiddleware func(HandlerFunc) HandlerFunc
211211

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-
246212
// HeadlerFunc is returning response headers based on the received request
247213
// headers.
248214
type HeadlerFunc func(Headers, swarm.Address) Headers

pkg/p2p/p2p_test.go

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

77
import (
8-
"context"
9-
"errors"
108
"testing"
119

12-
"github.com/coreos/go-semver/semver"
1310
"github.com/ethersphere/bee/v2/pkg/p2p"
14-
"github.com/ethersphere/bee/v2/pkg/swarm"
1511
"github.com/libp2p/go-libp2p/core/network"
1612
)
1713

@@ -40,146 +36,3 @@ func TestReachabilityStatus_String(t *testing.T) {
4036
}
4137
}
4238
}
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)