-
Notifications
You must be signed in to change notification settings - Fork 5
/
Copy pathstream_event_mux.go
129 lines (112 loc) · 2.12 KB
/
stream_event_mux.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
// +build linux
// Copyright (C) 2018 Kun Zhong All rights reserved.
// Use of this source code is governed by a MIT-style license that can be found
// in the LICENSE file.
package netgo
import (
"time"
"syscall"
"sync/atomic"
"golang.org/x/net/context"
)
func (stream *netStream) OnEvent(events uint32) {
for {
if events & syscall.EPOLLERR != 0 {
stream.rawClose()
break
}
if events & syscall.EPOLLOUT != 0 {
select {
case stream.eventOut <- events:
break
default:
break
}
}
if events & syscall.EPOLLIN != 0 {
if stream.PreAccept() {
stream.netDriver.accept(stream.netaddr)
break
}
ctx, _ := context.WithTimeout(context.Background(), time.Millisecond)
select {
case stream.eventIn<-events:
break
case <-ctx.Done(): //avoid event idle
break
}
}
break
}
}
func (stream *netStream) handleEventIn() {
loop:
for {
select {
case <-stream.eventIn:
if stream.State() == streamStateUnknow {
continue
}
if stream.packflag {
stream.recvPack()
} else {
stream.recv()
}
break
case <-stream.closeCh:
break loop
}
}
}
func (stream *netStream) handleEventOut() {
loop:
for {
next:
select {
case <-stream.eventOut:
state := stream.State()
if state == streamStateOpened {
stream.send()
} else if state == streamStateOpening {
for {
if atomic.CompareAndSwapUint32(&stream.state, state, streamStateOpened) {
break
}
state = stream.State()
if state == streamStateUnknow {
break next
}
}
stream.send()
}
break
case <-stream.closeCh:
break loop
}
}
}
func (stream *netStream) mux() {
loop:
for {
select {
case <-stream.closeCh:
break loop
case <-stream.reconnectTimer.C:
stream.onReconnectTimer()
break
}
}
stream.reconnectTimer.Stop()
}
func (stream *netStream) onReconnectTimer() {
if stream.fd.fd >= 0 {
stream.reconnectTimer.Reset(infinitTime)
return
}
err := stream.netDriver.rawConnect(stream, stream.fd.remoteAddr)
if err != nil{
return
}
stream.timerPeriod = infinitTime
stream.reconnectTimer.Reset(stream.timerPeriod)
}