-
Notifications
You must be signed in to change notification settings - Fork 26
/
transport.base.go
129 lines (113 loc) · 4.36 KB
/
transport.base.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
123
124
125
126
127
128
129
package prototree
import (
"sync"
"redwood.dev/state"
"redwood.dev/tree"
)
type BaseTreeTransport struct {
muTxReceivedCallbacks sync.RWMutex
muPrivateTxReceivedCallbacks sync.RWMutex
muAckReceivedCallbacks sync.RWMutex
muWritableSubscriptionOpenedCallback sync.RWMutex
muP2PStateURIReceivedCallbacks sync.RWMutex
txReceivedCallbacks []TxReceivedCallback
privateTxReceivedCallbacks []PrivateTxReceivedCallback
ackReceivedCallbacks []AckReceivedCallback
writableSubscriptionOpenedCallback WritableSubscriptionOpenedCallback
p2pStateURIReceivedCallbacks []P2PStateURIReceivedCallback
}
type TxReceivedCallback func(tx tree.Tx, peerConn TreePeerConn)
type PrivateTxReceivedCallback func(encryptedTx EncryptedTx, peerConn TreePeerConn)
type AckReceivedCallback func(stateURI string, txID state.Version, peerConn TreePeerConn)
type WritableSubscriptionOpenedCallback func(req SubscriptionRequest, writeSubImplFactory WritableSubscriptionImplFactory) (<-chan struct{}, error)
type P2PStateURIReceivedCallback func(stateURI string, peerConn TreePeerConn)
type WritableSubscriptionImplFactory func() (WritableSubscriptionImpl, error)
func (t *BaseTreeTransport) OnTxReceived(handler TxReceivedCallback) {
t.muTxReceivedCallbacks.Lock()
defer t.muTxReceivedCallbacks.Unlock()
t.txReceivedCallbacks = append(t.txReceivedCallbacks, handler)
}
func (t *BaseTreeTransport) OnPrivateTxReceived(handler PrivateTxReceivedCallback) {
t.muPrivateTxReceivedCallbacks.Lock()
defer t.muPrivateTxReceivedCallbacks.Unlock()
t.privateTxReceivedCallbacks = append(t.privateTxReceivedCallbacks, handler)
}
func (t *BaseTreeTransport) OnAckReceived(handler AckReceivedCallback) {
t.muAckReceivedCallbacks.Lock()
defer t.muAckReceivedCallbacks.Unlock()
t.ackReceivedCallbacks = append(t.ackReceivedCallbacks, handler)
}
func (t *BaseTreeTransport) OnWritableSubscriptionOpened(handler WritableSubscriptionOpenedCallback) {
t.muWritableSubscriptionOpenedCallback.Lock()
defer t.muWritableSubscriptionOpenedCallback.Unlock()
if t.writableSubscriptionOpenedCallback != nil {
panic("only one")
}
t.writableSubscriptionOpenedCallback = handler
}
func (t *BaseTreeTransport) OnP2PStateURIReceived(handler P2PStateURIReceivedCallback) {
t.muP2PStateURIReceivedCallbacks.Lock()
defer t.muP2PStateURIReceivedCallbacks.Unlock()
t.p2pStateURIReceivedCallbacks = append(t.p2pStateURIReceivedCallbacks, handler)
}
func (t *BaseTreeTransport) HandleTxReceived(tx tree.Tx, peerConn TreePeerConn) {
t.muTxReceivedCallbacks.RLock()
defer t.muTxReceivedCallbacks.RUnlock()
var wg sync.WaitGroup
wg.Add(len(t.txReceivedCallbacks))
for _, handler := range t.txReceivedCallbacks {
handler := handler
go func() {
defer wg.Done()
handler(tx, peerConn)
}()
}
wg.Wait()
}
func (t *BaseTreeTransport) HandlePrivateTxReceived(encryptedTx EncryptedTx, peerConn TreePeerConn) {
t.muPrivateTxReceivedCallbacks.RLock()
defer t.muPrivateTxReceivedCallbacks.RUnlock()
var wg sync.WaitGroup
wg.Add(len(t.privateTxReceivedCallbacks))
for _, handler := range t.privateTxReceivedCallbacks {
handler := handler
go func() {
defer wg.Done()
handler(encryptedTx, peerConn)
}()
}
wg.Wait()
}
func (t *BaseTreeTransport) HandleAckReceived(stateURI string, txID state.Version, peerConn TreePeerConn) {
t.muAckReceivedCallbacks.RLock()
defer t.muAckReceivedCallbacks.RUnlock()
var wg sync.WaitGroup
wg.Add(len(t.ackReceivedCallbacks))
for _, handler := range t.ackReceivedCallbacks {
handler := handler
go func() {
defer wg.Done()
handler(stateURI, txID, peerConn)
}()
}
wg.Wait()
}
func (t *BaseTreeTransport) HandleWritableSubscriptionOpened(req SubscriptionRequest, writeSubImplFactory WritableSubscriptionImplFactory) (<-chan struct{}, error) {
t.muWritableSubscriptionOpenedCallback.RLock()
defer t.muWritableSubscriptionOpenedCallback.RUnlock()
return t.writableSubscriptionOpenedCallback(req, writeSubImplFactory)
}
func (t *BaseTreeTransport) HandleP2PStateURIReceived(stateURI string, peerConn TreePeerConn) {
t.muP2PStateURIReceivedCallbacks.RLock()
defer t.muP2PStateURIReceivedCallbacks.RUnlock()
var wg sync.WaitGroup
wg.Add(len(t.p2pStateURIReceivedCallbacks))
for _, handler := range t.p2pStateURIReceivedCallbacks {
handler := handler
go func() {
defer wg.Done()
handler(stateURI, peerConn)
}()
}
wg.Wait()
}