master
go 128 lines 2.49 KB
Raw
1 package p2p
2
3 import (
4 "context"
5 "time"
6
7 tec "github.com/jbenet/go-temp-err-catcher"
8 net "github.com/libp2p/go-libp2p/core/network"
9 "github.com/libp2p/go-libp2p/core/peer"
10 "github.com/libp2p/go-libp2p/core/protocol"
11 ma "github.com/multiformats/go-multiaddr"
12 manet "github.com/multiformats/go-multiaddr/net"
13 )
14
15 // localListener manet streams and proxies them to libp2p services.
16 type localListener struct {
17 ctx context.Context
18
19 p2p *P2P
20
21 proto protocol.ID
22 laddr ma.Multiaddr
23 peer peer.ID
24
25 listener manet.Listener
26 done chan struct{}
27 }
28
29 // ForwardLocal creates new P2P stream to a remote listener.
30 func (p2p *P2P) ForwardLocal(ctx context.Context, peer peer.ID, proto protocol.ID, bindAddr ma.Multiaddr) (Listener, error) {
31 listener := &localListener{
32 ctx: ctx,
33 p2p: p2p,
34 proto: proto,
35 peer: peer,
36 done: make(chan struct{}),
37 }
38
39 maListener, err := manet.Listen(bindAddr)
40 if err != nil {
41 return nil, err
42 }
43
44 listener.listener = maListener
45 listener.laddr = maListener.Multiaddr()
46
47 if err := p2p.ListenersLocal.Register(listener); err != nil {
48 return nil, err
49 }
50
51 go listener.acceptConns()
52
53 return listener, nil
54 }
55
56 func (l *localListener) dial(ctx context.Context) (net.Stream, error) {
57 cctx, cancel := context.WithTimeout(ctx, time.Second*30) // TODO: configurable?
58 defer cancel()
59
60 return l.p2p.peerHost.NewStream(cctx, l.peer, l.proto)
61 }
62
63 func (l *localListener) acceptConns() {
64 for {
65 local, err := l.listener.Accept()
66 if err != nil {
67 if tec.ErrIsTemporary(err) {
68 continue
69 }
70 return
71 }
72
73 go l.setupStream(local)
74 }
75 }
76
77 func (l *localListener) setupStream(local manet.Conn) {
78 remote, err := l.dial(l.ctx)
79 if err != nil {
80 local.Close()
81 log.Warnf("failed to dial to remote %s/%s", l.peer, l.proto)
82 return
83 }
84
85 stream := &Stream{
86 Protocol: l.proto,
87
88 OriginAddr: local.RemoteMultiaddr(),
89 TargetAddr: l.TargetAddress(),
90 peer: l.peer,
91
92 Local: local,
93 Remote: remote,
94
95 Registry: l.p2p.Streams,
96 }
97
98 l.p2p.Streams.Register(stream)
99 }
100
101 func (l *localListener) close() {
102 l.listener.Close()
103 close(l.done)
104 }
105
106 func (l *localListener) Done() <-chan struct{} {
107 return l.done
108 }
109
110 func (l *localListener) Protocol() protocol.ID {
111 return l.proto
112 }
113
114 func (l *localListener) ListenAddress() ma.Multiaddr {
115 return l.laddr
116 }
117
118 func (l *localListener) TargetAddress() ma.Multiaddr {
119 addr, err := ma.NewMultiaddr(maPrefix + l.peer.String())
120 if err != nil {
121 panic(err)
122 }
123 return addr
124 }
125
126 func (l *localListener) key() protocol.ID {
127 return protocol.ID(l.ListenAddress().String())
128 }