-
Notifications
You must be signed in to change notification settings - Fork 55
/
channel.go
122 lines (99 loc) · 2.85 KB
/
channel.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
package directchannel
import (
"bufio"
"context"
"encoding/binary"
"fmt"
"io"
"berty.tech/go-orbit-db/iface"
"berty.tech/go-orbit-db/pubsub"
"github.com/libp2p/go-libp2p/core/host"
"github.com/libp2p/go-libp2p/core/network"
"github.com/libp2p/go-libp2p/core/peer"
"go.uber.org/zap"
)
const PROTOCOL = "/go-orbit-db/direct-channel/1.2.0"
const DelimitedReadMaxSize = 1024 * 1024 * 4 // mb
type directChannel struct {
logger *zap.Logger
host host.Host
emitter iface.DirectChannelEmitter
}
// Send Sends a message to the other peer
func (d *directChannel) Send(ctx context.Context, pid peer.ID, bytes []byte) error {
stream, err := d.host.NewStream(ctx, pid, PROTOCOL)
if err != nil {
return fmt.Errorf("unable to create stream: %w", err)
}
defer stream.Close()
length := uint64(len(bytes))
lenbuf := make([]byte, binary.MaxVarintLen64)
n := binary.PutUvarint(lenbuf, length)
_, err = stream.Write(lenbuf[:n])
if err != nil {
return fmt.Errorf("unable to write buflen: %w", err)
}
_, err = stream.Write(bytes)
return err
}
func (d *directChannel) handleNewPeer(s network.Stream) {
defer func() {
_ = s.Reset()
}()
reader := bufio.NewReader(s)
length64, err := binary.ReadUvarint(reader)
if err != nil {
d.logger.Error("unable to read length", zap.Error(err))
return
}
length := int(length64)
if length > DelimitedReadMaxSize {
d.logger.Error(fmt.Sprintf("received data exceeding maximum allowed size (%d > %d)", length, DelimitedReadMaxSize))
return
}
data := make([]byte, length)
if _, err := io.ReadFull(reader, data); err != nil {
d.logger.Error("unable to read buffer", zap.Error(err))
return
}
if err := d.emitter.Emit(pubsub.NewEventPayload(data, s.Conn().RemotePeer())); err != nil {
d.logger.Error("unable to emit on emitter", zap.Error(err))
}
}
// @NOTE(gfanton): we dont need this on direct channel
// Connect Waits for the other peer to be connected
func (d *directChannel) Connect(ctx context.Context, pid peer.ID) (err error) {
return nil
}
// @NOTE(gfanton): we dont need this on direct channel
// Close Closes the connection
func (d *directChannel) Close() error {
d.host.RemoveStreamHandler(PROTOCOL)
return d.emitter.Close()
}
type holderChannels struct {
host host.Host
logger *zap.Logger
}
func (c *holderChannels) NewChannel(ctx context.Context, emitter iface.DirectChannelEmitter, opts *iface.DirectChannelOptions) (iface.DirectChannel, error) {
if opts == nil {
opts = &iface.DirectChannelOptions{}
}
if opts.Logger == nil {
opts.Logger = c.logger
}
dc := &directChannel{
logger: c.logger,
host: c.host,
emitter: emitter,
}
c.host.SetStreamHandler(PROTOCOL, dc.handleNewPeer)
return dc, nil
}
func InitDirectChannelFactory(logger *zap.Logger, host host.Host) iface.DirectChannelFactory {
holder := &holderChannels{
logger: logger,
host: host,
}
return holder.NewChannel
}