/
workerconn.go
92 lines (79 loc) · 1.64 KB
/
workerconn.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
package termite
import (
"io"
"net"
"sync"
)
// connDialer dials connections that have IDs beyond address.
type connDialer interface {
Dial(addr string) (connMuxer, error)
}
// connMuxer opens multiple named streams over a connection
type connMuxer interface {
Open(id string) (io.ReadWriteCloser, error)
Close() error
}
// connListener accepts connections that have string IDs.
type connListener interface {
Addr() net.Addr
Close() error
Pending() *pendingConns
}
type pendingConns struct {
rpcChans chan io.ReadWriteCloser
conns map[string]io.ReadWriteCloser
cond sync.Cond
}
func newPendingConns() *pendingConns {
p := &pendingConns{
conns: map[string]io.ReadWriteCloser{},
rpcChans: make(chan io.ReadWriteCloser, 1),
}
p.cond.L = new(sync.Mutex)
return p
}
func (p *pendingConns) rpcChan() <-chan io.ReadWriteCloser {
return p.rpcChans
}
func (p *pendingConns) fail() {
p.cond.L.Lock()
defer p.cond.L.Unlock()
p.conns = nil
p.cond.Broadcast()
}
func (p *pendingConns) wait() {
p.cond.L.Lock()
defer p.cond.L.Unlock()
for p.conns != nil {
p.cond.Wait()
}
}
func (p *pendingConns) add(key string, conn io.ReadWriteCloser) {
if key == RPC_CHANNEL {
p.rpcChans <- conn
return
}
p.cond.L.Lock()
defer p.cond.L.Unlock()
if p.conns == nil {
panic("shut down")
}
if p.conns[key] != nil {
panic("collision")
}
p.conns[key] = conn
p.cond.Broadcast()
}
func (p *pendingConns) accept(key string) io.ReadWriteCloser {
p.cond.L.Lock()
defer p.cond.L.Unlock()
for p.conns != nil && p.conns[key] == nil {
p.cond.Wait()
}
if p.conns == nil {
return nil
}
ch := p.conns[key]
delete(p.conns, key)
return ch
}