Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Support
Submit feedback
Contribute to GitLab
Sign in
Toggle navigation
M
mybee
Project
Project
Details
Activity
Releases
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Boards
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
vicotor
mybee
Commits
44adb6ec
Commit
44adb6ec
authored
Jan 20, 2020
by
Janos Guljas
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
add protobuf ReadMessages function
parent
03a459a2
Changes
2
Show whitespace changes
Inline
Side-by-side
Showing
2 changed files
with
41 additions
and
25 deletions
+41
-25
protobuf.go
pkg/p2p/protobuf/protobuf.go
+18
-0
pingpong_test.go
pkg/pingpong/pingpong_test.go
+23
-25
No files found.
pkg/p2p/protobuf/protobuf.go
View file @
44adb6ec
...
@@ -2,12 +2,15 @@ package protobuf
...
@@ -2,12 +2,15 @@ package protobuf
import
(
import
(
ggio
"github.com/gogo/protobuf/io"
ggio
"github.com/gogo/protobuf/io"
"github.com/gogo/protobuf/proto"
"github.com/janos/bee/pkg/p2p"
"github.com/janos/bee/pkg/p2p"
"io"
"io"
)
)
const
delimitedReaderMaxSize
=
128
*
1024
// max message size
const
delimitedReaderMaxSize
=
128
*
1024
// max message size
type
Message
=
proto
.
Message
func
NewWriterAndReader
(
s
p2p
.
Stream
)
(
w
ggio
.
Writer
,
r
ggio
.
Reader
)
{
func
NewWriterAndReader
(
s
p2p
.
Stream
)
(
w
ggio
.
Writer
,
r
ggio
.
Reader
)
{
r
=
ggio
.
NewDelimitedReader
(
s
,
delimitedReaderMaxSize
)
r
=
ggio
.
NewDelimitedReader
(
s
,
delimitedReaderMaxSize
)
w
=
ggio
.
NewDelimitedWriter
(
s
)
w
=
ggio
.
NewDelimitedWriter
(
s
)
...
@@ -21,3 +24,18 @@ func NewReader(r io.Reader) ggio.Reader {
...
@@ -21,3 +24,18 @@ func NewReader(r io.Reader) ggio.Reader {
func
NewWriter
(
w
io
.
Writer
)
ggio
.
Writer
{
func
NewWriter
(
w
io
.
Writer
)
ggio
.
Writer
{
return
ggio
.
NewDelimitedWriter
(
w
)
return
ggio
.
NewDelimitedWriter
(
w
)
}
}
func
ReadMessages
(
r
io
.
Reader
,
newMessage
func
()
Message
)
(
m
[]
Message
,
err
error
)
{
pr
:=
NewReader
(
r
)
for
{
msg
:=
newMessage
()
if
err
:=
pr
.
ReadMsg
(
msg
);
err
!=
nil
{
if
err
==
io
.
EOF
{
break
}
return
nil
,
err
}
m
=
append
(
m
,
msg
)
}
return
m
,
nil
}
pkg/pingpong/pingpong_test.go
View file @
44adb6ec
...
@@ -4,7 +4,6 @@ import (
...
@@ -4,7 +4,6 @@ import (
"bytes"
"bytes"
"context"
"context"
"fmt"
"fmt"
"io"
"testing"
"testing"
"github.com/janos/bee/pkg/p2p/mock"
"github.com/janos/bee/pkg/p2p/mock"
...
@@ -34,39 +33,38 @@ func TestPing(t *testing.T) {
...
@@ -34,39 +33,38 @@ func TestPing(t *testing.T) {
t
.
Errorf
(
"invalid RTT value %v"
,
rtt
)
t
.
Errorf
(
"invalid RTT value %v"
,
rtt
)
}
}
// validate received ping greetings
// validate received ping greetings from the client
r
:=
protobuf
.
NewReader
(
bytes
.
NewReader
(
streamer
.
In
.
Bytes
()))
wantGreetings
:=
greetings
var
gotGreetings
[]
string
messages
,
err
:=
protobuf
.
ReadMessages
(
for
{
bytes
.
NewReader
(
streamer
.
In
.
Bytes
()),
var
ping
pingpong
.
Ping
func
()
protobuf
.
Message
{
return
new
(
pingpong
.
Ping
)
},
if
err
:=
r
.
ReadMsg
(
&
ping
);
err
!=
nil
{
)
if
err
==
io
.
EOF
{
if
err
!=
nil
{
break
}
t
.
Fatal
(
err
)
t
.
Fatal
(
err
)
}
}
gotGreetings
=
append
(
gotGreetings
,
ping
.
Greeting
)
var
gotGreetings
[]
string
for
_
,
m
:=
range
messages
{
gotGreetings
=
append
(
gotGreetings
,
m
.
(
*
pingpong
.
Ping
)
.
Greeting
)
}
}
if
fmt
.
Sprint
(
gotGreetings
)
!=
fmt
.
Sprint
(
g
reetings
)
{
if
fmt
.
Sprint
(
gotGreetings
)
!=
fmt
.
Sprint
(
wantG
reetings
)
{
t
.
Errorf
(
"got greetings %v, want %v"
,
gotGreetings
,
g
reetings
)
t
.
Errorf
(
"got greetings %v, want %v"
,
gotGreetings
,
wantG
reetings
)
}
}
// validate send pong responses by handler
// validate sent pong responses by handler
r
=
protobuf
.
NewReader
(
bytes
.
NewReader
(
streamer
.
Out
.
Bytes
()))
var
wantResponses
[]
string
var
wantResponses
[]
string
for
_
,
g
:=
range
greetings
{
for
_
,
g
:=
range
greetings
{
wantResponses
=
append
(
wantResponses
,
"{"
+
g
+
"}"
)
wantResponses
=
append
(
wantResponses
,
"{"
+
g
+
"}"
)
}
}
var
gotResponses
[]
string
messages
,
err
=
protobuf
.
ReadMessages
(
for
{
bytes
.
NewReader
(
streamer
.
Out
.
Bytes
()),
var
pong
pingpong
.
Pong
func
()
protobuf
.
Message
{
return
new
(
pingpong
.
Pong
)
},
if
err
:=
r
.
ReadMsg
(
&
pong
);
err
!=
nil
{
)
if
err
==
io
.
EOF
{
if
err
!=
nil
{
break
}
t
.
Fatal
(
err
)
t
.
Fatal
(
err
)
}
}
gotResponses
=
append
(
gotResponses
,
pong
.
Response
)
var
gotResponses
[]
string
for
_
,
m
:=
range
messages
{
gotResponses
=
append
(
gotResponses
,
m
.
(
*
pingpong
.
Pong
)
.
Response
)
}
}
if
fmt
.
Sprint
(
gotResponses
)
!=
fmt
.
Sprint
(
wantResponses
)
{
if
fmt
.
Sprint
(
gotResponses
)
!=
fmt
.
Sprint
(
wantResponses
)
{
t
.
Errorf
(
"got responses %v, want %v"
,
gotResponses
,
wantResponses
)
t
.
Errorf
(
"got responses %v, want %v"
,
gotResponses
,
wantResponses
)
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment