mirror of
https://github.com/soypat/lneto.git
synced 2026-08-19 22:24:03 +00:00
Fix slow exchanges on linux by adding polling and other improvements (#14)
* add prints to investigate slowness * add loop sleep for blocking Stack * add connection close http attribute and detect end of HTML page * add timeout to bridge read * examples/xnet: add ntp timestamp to prints and show output in readme * add polling to http tap interface * more digits in ntp format * prevent crashes on non-8 timestamp
This commit is contained in:
+45
-14
@@ -26,6 +26,16 @@ type Interface interface {
|
||||
IPMask() (netip.Prefix, error)
|
||||
}
|
||||
|
||||
type HTTPTapClient struct {
|
||||
c http.Client
|
||||
infoURL string
|
||||
recvurl string
|
||||
sendurl string
|
||||
ip netip.Prefix
|
||||
hwaddr [6]byte
|
||||
buf []byte
|
||||
}
|
||||
|
||||
var _ Interface = (*HTTPTapClient)(nil)
|
||||
|
||||
// NewHTTPTapClient returns a HTTPTapClient ready for use.
|
||||
@@ -65,12 +75,7 @@ func (h *HTTPTapClient) ensureMTU() (err error) {
|
||||
err = fmt.Errorf("unable to get MTU from server: %w", err)
|
||||
}
|
||||
}()
|
||||
resp, err := h.c.Get(h.infoURL)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
var info tapInfo
|
||||
err = json.NewDecoder(resp.Body).Decode(&info)
|
||||
info, err := h.info()
|
||||
if err != nil {
|
||||
return err
|
||||
} else if info.MTU <= minMTU {
|
||||
@@ -88,14 +93,14 @@ func (h *HTTPTapClient) ensureMTU() (err error) {
|
||||
return nil
|
||||
}
|
||||
|
||||
type HTTPTapClient struct {
|
||||
c http.Client
|
||||
infoURL string
|
||||
recvurl string
|
||||
sendurl string
|
||||
ip netip.Prefix
|
||||
hwaddr [6]byte
|
||||
buf []byte
|
||||
func (h *HTTPTapClient) info() (tapInfo, error) {
|
||||
resp, err := h.c.Get(h.infoURL)
|
||||
if err != nil {
|
||||
return tapInfo{}, err
|
||||
}
|
||||
var info tapInfo
|
||||
err = json.NewDecoder(resp.Body).Decode(&info)
|
||||
return info, err
|
||||
}
|
||||
|
||||
func (h *HTTPTapClient) ReadDiscard() (err error) {
|
||||
@@ -109,6 +114,24 @@ func (h *HTTPTapClient) ReadDiscard() (err error) {
|
||||
return err
|
||||
}
|
||||
|
||||
func (h *HTTPTapClient) Poll(d time.Duration) (ready bool, err error) {
|
||||
info, err := h.info()
|
||||
if err != nil {
|
||||
return false, err
|
||||
} else if info.DataReady {
|
||||
return true, nil
|
||||
}
|
||||
deadline := time.Now().Add(d)
|
||||
for !info.DataReady && time.Until(deadline) > 0 {
|
||||
time.Sleep(5 * time.Millisecond)
|
||||
info, err = h.info()
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
}
|
||||
return info.DataReady, err
|
||||
}
|
||||
|
||||
func (h *HTTPTapClient) ReadBytes() (data []byte, err error) {
|
||||
err = h.ensureMTU()
|
||||
if err != nil {
|
||||
@@ -177,6 +200,7 @@ type tapInfo struct {
|
||||
MTU int
|
||||
IPPrefix string
|
||||
HardwareAddr string
|
||||
DataReady bool
|
||||
}
|
||||
|
||||
func (sv *HTTPTapServer) OnTransfer(cb func(channel int, pkt []byte)) {
|
||||
@@ -255,10 +279,17 @@ func NewHTTPTapServer(iface Interface, minMTU, queueOut, queueIn int) (*HTTPTapS
|
||||
hwstr := net.HardwareAddr(hw6[:]).String()
|
||||
ipstr := netmask.String()
|
||||
sv.HandleFunc("/info", func(w http.ResponseWriter, r *http.Request) {
|
||||
var dataready bool = true
|
||||
if poller, ok := taps.tap.(interface {
|
||||
Poll(time.Duration) (bool, error)
|
||||
}); ok {
|
||||
dataready, err = poller.Poll(0)
|
||||
}
|
||||
info := tapInfo{
|
||||
MTU: mtu,
|
||||
IPPrefix: ipstr,
|
||||
HardwareAddr: hwstr,
|
||||
DataReady: dataready,
|
||||
}
|
||||
json.NewEncoder(w).Encode(info)
|
||||
})
|
||||
|
||||
@@ -11,6 +11,7 @@ import (
|
||||
"os"
|
||||
"os/exec"
|
||||
"syscall"
|
||||
"time"
|
||||
"unsafe"
|
||||
)
|
||||
|
||||
@@ -66,6 +67,22 @@ func (tap *Tap) Read(b []byte) (int, error) {
|
||||
return syscall.Read(tap.fd, b)
|
||||
}
|
||||
|
||||
// Poll waits up to timeout for the tap device to have data available for reading.
|
||||
// Returns true if data is available, false if timeout was reached.
|
||||
func (tap *Tap) Poll(timeout time.Duration) (bool, error) {
|
||||
var readfds syscall.FdSet
|
||||
readfds.Bits[tap.fd/64] |= 1 << (uint(tap.fd) % 64)
|
||||
tv := syscall.Timeval{
|
||||
Sec: int64(timeout / time.Second),
|
||||
Usec: int64((timeout % time.Second) / time.Microsecond),
|
||||
}
|
||||
n, err := syscall.Select(tap.fd+1, &readfds, nil, nil, &tv)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return n > 0, nil
|
||||
}
|
||||
|
||||
func (tap *Tap) Write(b []byte) (int, error) {
|
||||
return syscall.Write(tap.fd, b)
|
||||
}
|
||||
@@ -247,6 +264,33 @@ func (br *Bridge) IPMask() (netip.Prefix, error) {
|
||||
return getSocketMask(br.fd, br.name)
|
||||
}
|
||||
|
||||
// SetReadTimeout sets the receive timeout for the bridge socket.
|
||||
// This prevents Read from blocking indefinitely, allowing the caller
|
||||
// to periodically call Encapsulate even when no packets arrive.
|
||||
func (br *Bridge) SetReadTimeout(timeout time.Duration) error {
|
||||
tv := syscall.Timeval{
|
||||
Sec: int64(timeout / time.Second),
|
||||
Usec: int64((timeout % time.Second) / time.Microsecond),
|
||||
}
|
||||
return syscall.SetsockoptTimeval(br.fd, syscall.SOL_SOCKET, syscall.SO_RCVTIMEO, &tv)
|
||||
}
|
||||
|
||||
// Poll waits up to timeout for the bridge socket to have data available for reading.
|
||||
// Returns true if data is available, false if timeout was reached.
|
||||
func (br *Bridge) Poll(timeout time.Duration) (bool, error) {
|
||||
var readfds syscall.FdSet
|
||||
readfds.Bits[br.fd/64] |= 1 << (uint(br.fd) % 64)
|
||||
tv := syscall.Timeval{
|
||||
Sec: int64(timeout / time.Second),
|
||||
Usec: int64((timeout % time.Second) / time.Microsecond),
|
||||
}
|
||||
n, err := syscall.Select(br.fd+1, &readfds, nil, nil, &tv)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return n > 0, nil
|
||||
}
|
||||
|
||||
func (br *Bridge) Addr() (netip.Addr, error) {
|
||||
addrp, err := getSocketIP(br.fd, br.name)
|
||||
if err != nil {
|
||||
|
||||
Reference in New Issue
Block a user