Initial commit
This commit is contained in:
@@ -0,0 +1,65 @@
|
||||
package device
|
||||
|
||||
import (
|
||||
"log"
|
||||
"net"
|
||||
"sync"
|
||||
|
||||
tunnelpb "relay/proto/tunnel"
|
||||
)
|
||||
|
||||
type Device struct {
|
||||
stream tunnelpb.TunnelService_TunnelServer
|
||||
sendCh chan *tunnelpb.Frame
|
||||
done chan struct{}
|
||||
closeOnce sync.Once
|
||||
|
||||
streamsMu sync.Mutex
|
||||
streams map[uint32]net.Conn
|
||||
streamDone map[uint32]chan struct{}
|
||||
nextID uint32
|
||||
}
|
||||
|
||||
func NewDevice(stream tunnelpb.TunnelService_TunnelServer) *Device {
|
||||
d := &Device{
|
||||
stream: stream,
|
||||
sendCh: make(chan *tunnelpb.Frame, 128), // backpressure here
|
||||
done: make(chan struct{}),
|
||||
streams: make(map[uint32]net.Conn),
|
||||
streamDone: make(map[uint32]chan struct{}),
|
||||
nextID: 1,
|
||||
}
|
||||
|
||||
go d.writer()
|
||||
return d
|
||||
}
|
||||
|
||||
func (d *Device) Close() {
|
||||
d.closeOnce.Do(func() {
|
||||
close(d.done)
|
||||
})
|
||||
}
|
||||
|
||||
func (d *Device) SendFrame(f *tunnelpb.Frame) {
|
||||
select {
|
||||
case d.sendCh <- f:
|
||||
case <-d.done:
|
||||
default:
|
||||
log.Println("device backpressure: drop frame")
|
||||
}
|
||||
}
|
||||
|
||||
func (d *Device) writer() {
|
||||
for {
|
||||
select {
|
||||
case f := <-d.sendCh:
|
||||
if err := d.stream.Send(f); err != nil {
|
||||
log.Println("device send error:", err)
|
||||
d.Close()
|
||||
return
|
||||
}
|
||||
case <-d.done:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
package device
|
||||
|
||||
import (
|
||||
"net"
|
||||
)
|
||||
|
||||
func (d *Device) AllocateStreamID() uint32 {
|
||||
d.streamsMu.Lock()
|
||||
defer d.streamsMu.Unlock()
|
||||
|
||||
id := d.nextID
|
||||
d.nextID++
|
||||
return id
|
||||
}
|
||||
|
||||
func (d *Device) AddStream(id uint32, conn net.Conn) chan struct{} {
|
||||
d.streamsMu.Lock()
|
||||
defer d.streamsMu.Unlock()
|
||||
|
||||
ch := make(chan struct{})
|
||||
d.streams[id] = conn
|
||||
d.streamDone[id] = ch
|
||||
return ch
|
||||
}
|
||||
|
||||
func (d *Device) RemoveStream(id uint32) {
|
||||
d.streamsMu.Lock()
|
||||
defer d.streamsMu.Unlock()
|
||||
|
||||
if c, ok := d.streams[id]; ok {
|
||||
_ = c.Close()
|
||||
delete(d.streams, id)
|
||||
}
|
||||
if ch, ok := d.streamDone[id]; ok {
|
||||
close(ch)
|
||||
delete(d.streamDone, id)
|
||||
}
|
||||
}
|
||||
|
||||
func (d *Device) GetClient(id uint32) (net.Conn, bool) {
|
||||
d.streamsMu.Lock()
|
||||
defer d.streamsMu.Unlock()
|
||||
c, ok := d.streams[id]
|
||||
return c, ok
|
||||
}
|
||||
|
||||
func (d *Device) Done() <-chan struct{} {
|
||||
return d.done
|
||||
}
|
||||
Reference in New Issue
Block a user