diff --git a/examples/stackbasic/main.go b/examples/stackbasic/main.go index 17350b1..01ea1f2 100644 --- a/examples/stackbasic/main.go +++ b/examples/stackbasic/main.go @@ -9,6 +9,7 @@ import ( "net" "net/netip" "os" + "runtime" "time" "github.com/soypat/lneto" @@ -62,13 +63,15 @@ func main() { buf := make([]byte, mtu) var hdr httpraw.Header hdr.Reset(make([]byte, 0, 1024)) + const standbyDuration = 5 * time.Second + lastHit := time.Now().Add(-standbyDuration) for { nread, err := tap.Read(buf[:]) if err != nil { slogger.error("tap-err", slog.String("err", err.Error())) log.Fatal(err) } else if nread > 0 { - debugEthPacket(nil, "IN ", buf[:nread]) + // debugEthPacket(nil, "IN ", buf[:nread]) // fmt.Println("INHEX ", debugHex(buf[:nread])) err = lStack.RecvEth(buf[:nread]) if err != nil { @@ -85,13 +88,18 @@ func main() { _, err = tap.Write(buf[:nw]) if err != nil { log.Fatal(err) - } else { - slogger.info("write", slog.Int("plen", nw)) } } - - if nread == 0 && nw == 0 { - time.Sleep(5 * time.Millisecond) + hit := nread > 0 || nw > 0 + if hit { + slogger.info("exchange", slog.Int("read", nread), slog.Int("nwrite", nw)) + lastHit = time.Now() + } else { + if time.Since(lastHit) > standbyDuration { + time.Sleep(5 * time.Millisecond) + } else { + runtime.Gosched() + } } } } @@ -119,10 +127,12 @@ func doHTTP(conn *internet.TCPConn, hdr *httpraw.Header) error { fmt.Println("sending response...") hdr.Reset(nil) hdr.SetStatus("200", "OK") + data := `{"ok":true}` response, err := hdr.AppendResponse(nil) if err != nil { return err } + response = append(response, data...) _, err = conn.Write(response) if err != nil { return err diff --git a/examples/tap/main.go b/examples/tap/main.go index 5aaadff..a5ac544 100644 --- a/examples/tap/main.go +++ b/examples/tap/main.go @@ -8,6 +8,7 @@ import ( "net" "net/http" "net/netip" + "runtime" "time" "github.com/soypat/lneto/internal/ltesto" @@ -44,7 +45,8 @@ func run() error { } fmt.Println("listening on http://127.0.0.1:7070/recv and http://127.0.0.1:7070/send on hwaddr:", net.HardwareAddr(hwaddr[:]).String()) go http.ListenAndServe(":7070", sv) - misses := 0 + const standbyDuration = 5 * time.Second + lastHit := time.Now().Add(-standbyDuration) for { result, err := sv.HandleTap() if err != nil { @@ -53,14 +55,13 @@ func run() error { if result.Failed { return errors.New("tap failed, exit program") } else if result.ReceivedSize == 0 && result.SentSize == 0 { - misses++ - if misses > 1000 { - time.Sleep(200 * time.Millisecond) // No data exchanged, sleep a bit to not hog CPU. + if time.Since(lastHit) > standbyDuration { + time.Sleep(5 * time.Millisecond) // Enter standby. } else { - time.Sleep(50 * time.Millisecond) // No data exchanged, sleep a bit to not hog CPU. + runtime.Gosched() } } else { - misses = 0 + lastHit = time.Now() } } } diff --git a/examples/tcpclient/main.go b/examples/tcpclient/main.go index 3e66368..24e4afc 100644 --- a/examples/tcpclient/main.go +++ b/examples/tcpclient/main.go @@ -1,6 +1,7 @@ package main import ( + "flag" "fmt" "net" "os" @@ -19,6 +20,9 @@ func main() { } func run() error { + var port int + flag.IntVar(&port, "lport", 13337, "Local port over which to hit server") + flag.Parse() // Prepare GET request. var hdr httpraw.Header hdr.SetMethod("GET") @@ -30,7 +34,7 @@ func run() error { } fmt.Println("dialing...") - conn, err := net.DialTCP("tcp4", &net.TCPAddr{IP: []byte{192, 168, 10, 1}, Port: 1337}, &net.TCPAddr{IP: []byte{192, 168, 10, 2}, Port: 80}) + conn, err := net.DialTCP("tcp4", &net.TCPAddr{IP: []byte{192, 168, 10, 1}, Port: port}, &net.TCPAddr{IP: []byte{192, 168, 10, 2}, Port: 80}) if err != nil { return err } diff --git a/http/httpraw/header.go b/http/httpraw/header.go index 5372964..c340dd4 100644 --- a/http/httpraw/header.go +++ b/http/httpraw/header.go @@ -31,12 +31,13 @@ func (f flags) hasAny(checkThese flags) bool { return f&checkThese != 0 } -// Header implements "raw" HTTP validation and header key-value parsing, validation and marshalling. +// Header implements "raw" HTTP header key-value parsing, validation and marshalling. // // It does NOT implement: // - Normalization. // - Cookies (see [Cookie]). // - Special header optimizations. +// - Content-Length validation and other special header field value validation. type Header struct { hbuf headerBuf @@ -79,6 +80,7 @@ func (h *Header) Parse(asResponse bool) error { // TryParse begins parsing or resumes parsing from a failed previous attempt from any of the Parse* methods. // As long as needMoreData returns true future calls to TryParse may succeed and the header is not done parsing. +// Users may call [Header.ForEach] in-between TryParse calls so as to validate values before header is completely parsed. // // needMoreData := true // var err error @@ -306,7 +308,7 @@ func (h *Header) getNonEmptyValue(s headerSlice) []byte { return h.hbuf.musttoken(s) } -// AppendRequest appends the request representation to the buffer and returns the result. +// AppendRequest appends the request header representation to the buffer and returns the result. func (h *Header) AppendRequest(dst []byte) ([]byte, error) { if h.flags.hasAny(flagOOMReached) { return dst, errOOM @@ -333,7 +335,7 @@ func (h *Header) AppendRequest(dst []byte) ([]byte, error) { return append(dst, strCRLF...), nil } -// AppendResponse appends the response representation to the buffer and returns the result. +// AppendResponse appends the response header representation to the buffer and returns the result. func (h *Header) AppendResponse(dst []byte) ([]byte, error) { if h.flags.hasAny(flagOOMReached) { return dst, errOOM diff --git a/http/httpraw/parse.go b/http/httpraw/parse.go index a973dc1..8e258b4 100644 --- a/http/httpraw/parse.go +++ b/http/httpraw/parse.go @@ -222,6 +222,9 @@ func (h *Header) peekHeader(key string) argsKV { func (hb *headerBuf) mustAppendSlice(value string) headerSlice { L := len(hb.buf) + if L == 0 { + L++ // Valid key-values start after 0. + } copy(hb.buf[L:L+len(value)], value) hb.buf = hb.buf[:L+len(value)] return hb.slice(hb.buf[L : L+len(value)]) diff --git a/internal/backoff.go b/internal/backoff.go index ae58d97..8aa88dd 100644 --- a/internal/backoff.go +++ b/internal/backoff.go @@ -7,6 +7,7 @@ type BackoffFlags uint8 const ( BackoffHasPriority BackoffFlags = 1 << iota BackoffCriticalPath + BackoffTCPConn ) func NewBackoff(priority BackoffFlags) Backoff { @@ -14,6 +15,10 @@ func NewBackoff(priority BackoffFlags) Backoff { return Backoff{ maxWait: uint32(1 * time.Millisecond), } + } else if priority&BackoffTCPConn != 0 { + return Backoff{ + maxWait: uint32(5 * time.Microsecond), + } } return Backoff{ maxWait: uint32(time.Second) >> (priority & BackoffHasPriority), diff --git a/internet/tcpconn.go b/internet/tcpconn.go index 7904d43..112b2af 100644 --- a/internet/tcpconn.go +++ b/internet/tcpconn.go @@ -135,7 +135,7 @@ func (conn *TCPConn) Write(b []byte) (int, error) { } else if plen == 0 { return 0, nil } - backoff := internal.NewBackoff(internal.BackoffHasPriority) + backoff := internal.NewBackoff(internal.BackoffTCPConn) n := 0 for { if conn.abortErr != nil { @@ -171,7 +171,7 @@ func (conn *TCPConn) Read(b []byte) (int, error) { } conn.trace("TCPConn.Read:start") connid := conn.h.ConnectionID() - backoff := internal.NewBackoff(internal.BackoffHasPriority) + backoff := internal.NewBackoff(internal.BackoffTCPConn) for conn.h.BufferedInput() == 0 && conn.State() == tcp.StateEstablished { if conn.abortErr != nil { return 0, conn.abortErr diff --git a/tcp/handler.go b/tcp/handler.go index 1b2b90a..ee1a052 100644 --- a/tcp/handler.go +++ b/tcp/handler.go @@ -2,6 +2,7 @@ package tcp import ( "errors" + "io" "net" "log/slog" @@ -108,6 +109,7 @@ func (h *Handler) OpenListen(localPort uint16, iss Value) error { // Abort forcibly terminates all state associated to current connection. // After a call to abort no more data can be sent nor received over the connection. func (h *Handler) Abort() { + h.info("tcp.Handler.Abort") h.scb.Abort() h.reset(0, 0, 0) } @@ -215,8 +217,14 @@ func (h *Handler) Send(b []byte) (int, error) { } buffered := h.bufTx.Buffered() if buffered == 0 && h.closing { + // If Close called and no more data to be sent, terminate connection! + h.closing = false err = h.scb.Close() - h.info("tcp.Handler:Close", slog.String("scb.Close.err", errstr(err))) + if err != nil { + h.logerr("tcp.Handler.Close", slog.String("err", errstr(err)), slog.String("state", h.State().String())) + h.Abort() + return 0, io.EOF + } } offset := uint8(5) var segment Segment