add http/httphi (#171)

* add http/httphi

* begin adding httphi tests

* claude found neat bugs

* add low level Handle function and more tests

* more tests, run go generate

* add Hijacker-like functionality

* improve locking and acquisition of Exchanges in reconfiguring

* several bugfixes, add internal.IntLen, round up http-linux example with new router API

* small nit

* add benchmarks

* add query handling

* remove ForEach pattern, allocates in TinyGo

* massive documentation push and code reordering in files

* Router.Handle returns error after being torn down

* run go fix

* rework Mux interface to receive a string request path

* add MethodFrom

* minor doc nit

* fail on incomplete staging

* add raw buffer access

* add streaming API distinct from Exchange

* begin adding multipart form logic

* finish rounding up multipart form parsing

* remove status type

* first Multipart approach

* begin adding readMultiPart

* add Exchange.ReadMultiparts reimagining of clanker slop

* ai insists with backoffs

* simplify clanker slop

* apply go fix

* add a pattern argument to Mux

* explicit header key/value alloc and add ExchangeConfig

* fix tests after excplicit header alloc change

* fix examples

* run go fix

* expose rawsock as experimental package (will use for external benchmarks)

* remove backoff from form parsing

* @MDr164 suggestions get potential fixes

* apply go fix

* add examples

* add README.md

* fix rawsock tinygo implementation

* apply @MDr164 various fixes

* update documentation on ContentLength methods and fix bug in Form reset on empty body

* fix tests

* add fuzz tests

* run go fix

* io.ErrNoProgress on parsing form spin

* run go fix

* remove backoff assumption from Router

* httphi.Handle rejects unsupported protocols

* go format router.go

* add kvbuffer

* rewrite Cookie with KVBuffer

* mid refactor of KVBuffer into Header

* work on KVBuffer exhausted semantics

* add Go's ServeMux Request.PathValue access semantics to Exchange, Mux and MuxSlice

* add PathValue example

* document all the things; improve req Query semantics; add Form.EnableBufferGrowth

* unexport kvBuffer

* add Exchange.PathValueAppend

* use stdlib in example instead of rawsock

* remove rawsock from http example

* add darwin arch rawsock

* fix example

* rename Router.TeardownGoroutines to Shutdown matching http.Server.Shutdown

* rename types and identifiers

* @MDr164 Content-Type and Transfer-Encoding bug catches
This commit is contained in:
Pat Whittingslow
2026-07-29 20:14:46 -03:00
committed by Patricio Whittingslow
parent a3f2742abf
commit 5c54030f19
47 changed files with 7309 additions and 804 deletions
+51
View File
@@ -0,0 +1,51 @@
# httφ
Heapless HTTP/1.1 router with `net/http`-shaped handlers. Benchmarked against `net/http` and friends at [**soypat/httpbench**](https://github.com/soypat/httpbench).
## Why
`net/http` allocates per request: `Request`, header map, response buffer, goroutine. Fine on a
server, fatal on a microcontroller. httφ pays that cost once, at `Router.Configure`:
- Exchanges and goroutines are fixed there. Serving allocates nothing; memory does not grow with load.
- No free exchange means `Handle` refuses the connection unclosed, rather than allocating one more of everything.
- `Handle` takes an `io.ReadWriteCloser`, so the router runs over a raw socket, an [`lneto`](https://github.com/soypat/lneto) TCP stack or a test pipe. No listener, no OS, no clock.
Parsing is [`httpraw`](../httpraw).
## Example
```go
var mux httphi.MuxSlice
mux.Handle("GET /", func(ex *httphi.Exchange) {
ex.WriteBody([]byte("hello world"))
})
var router httphi.Router
err := router.Configure(httphi.RouterConfig{
FixedNumGoroutines: 4, // 4 workers, 4 exchanges, allocated here and never again.
MaxAwaitingConns: 8, // Queue depth. Full queue drops connections.
RequestHeaderBufferSize: 1024,
ResponseHeaderMinBufferSize: 32, // Shares the request buffer.
RequestNumHeaderKVCap: 32,
Backoff: func(uint) time.Duration { return time.Millisecond },
Mux: &mux,
})
if err != nil {
log.Fatal(err)
}
for {
conn, err := listener.Accept() // Accepting is the caller's job.
if err != nil {
log.Fatal(err)
}
if err = router.Handle(conn); err != nil {
conn.Close() // Refused: never blocks, never queues unboundedly.
}
}
```
`FixedNumGoroutines: -1` gives the unbounded flavor, a goroutine and an exchange per connection.
Runnable server over raw Linux sockets, plus query, form and multipart handlers:
[`example_test.go`](./example_test.go).
+134
View File
@@ -0,0 +1,134 @@
package httphi
import (
"io"
"testing"
"github.com/soypat/lneto/http/httpraw"
"github.com/soypat/lneto/internal"
)
// benchConn replays a fixed request and discards the response. It allocates
// nothing itself so benchmark alloc counts belong to the package under test.
type benchConn struct {
request string
read int
written int
}
func (c *benchConn) rewind() { c.read, c.written = 0, 0 }
func (c *benchConn) Read(b []byte) (int, error) {
if c.read >= len(c.request) {
return 0, io.EOF
}
n := copy(b, c.request[c.read:])
c.read += n
return n, nil
}
func (c *benchConn) Write(b []byte) (int, error) {
c.written += len(b)
return len(b), nil
}
func (c *benchConn) Close() error { return nil }
// benchBody is package level: converting a string literal to []byte inside the
// handler would allocate on every request and hide the router's own cost.
var benchBody = []byte("hello world")
func benchExchange(b *testing.B, conn conn) *Exchange {
b.Helper()
const bufferSize = 1024
const numHeaderCap = 2
exch := new(Exchange)
exch.Configure(ExchangeConfig{
RawBuf: make([]byte, 2*bufferSize),
RequestBufferLim: bufferSize,
NumHeaderKVCap: numHeaderCap,
})
if !exch.Acquire(conn) {
b.Fatal("fresh exchange failed to acquire connection")
}
return exch
}
// BenchmarkHandle measures a whole exchange: read, parse, mux and respond.
func BenchmarkHandle(b *testing.B) {
expect := []byte("123")
buf := make([]byte, 64)
for _, bb := range []struct {
name string
request string
handler HandlerFunc
}{
{
name: "GETWithHeadersAndQuery",
request: "GET /?abc=123 HTTP/1.1\r\nHost: tinygo.org\r\nUser-Agent: bench\r\nAccept: */*\r\nConnection: close\r\n\r\n",
handler: func(ex *Exchange) {
ex.StageHeader("Content-Type", "text/plain")
ex.StageHeaderInt("Content-Length", int64(len(benchBody)), 10)
data, present := ex.RequestQueryAppend(buf[:0], "abc", true)
if !present || !internal.BytesEqual(data, expect) {
panic("invalid result")
}
ex.WriteBody(benchBody)
},
},
{
name: "NotFound",
request: "GET /nowhere HTTP/1.1\r\nHost: tinygo.org\r\n\r\n",
handler: nil, // Unregistered: exercises the 404 path.
},
} {
b.Run(bb.name, func(b *testing.B) {
var mux MuxSlice
if bb.handler != nil {
mux.Handle("GET /", bb.handler)
}
conn := &benchConn{request: bb.request}
exch := benchExchange(b, conn)
b.ReportAllocs()
b.SetBytes(int64(len(bb.request)))
b.ResetTimer()
for i := 0; i < b.N; i++ {
conn.rewind()
exch.Release()
exch.Acquire(conn)
Handle(exch, &mux, nopBackoff)
}
})
}
}
// benchForm is package level so the Form's pair slice is reused across requests,
// as a real handler holding one per goroutine would.
var benchForm httpraw.Form
// BenchmarkRequestParseForm measures reading and parsing a urlencoded body into
// a buffer the caller owns. Nothing on the path may allocate.
func BenchmarkRequestParseForm(b *testing.B) {
const request = "POST /f HTTP/1.1\r\nHost: tinygo.org\r\n" +
"Content-Type: application/x-www-form-urlencoded\r\nContent-Length: 27\r\n\r\n" +
"user=gopher&msg=hello+world"
buf := make([]byte, 64)
var mux MuxSlice
mux.Handle("POST /f", func(ex *Exchange) {
err := ex.RequestParseForm(&benchForm, buf)
if err != nil || benchForm.Len() != 2 {
panic("invalid result")
}
})
conn := &benchConn{request: request}
exch := benchExchange(b, conn)
b.ReportAllocs()
b.SetBytes(int64(len(request)))
b.ResetTimer()
for b.Loop() {
conn.rewind()
exch.Release()
exch.Acquire(conn)
Handle(exch, &mux, nopBackoff)
}
}
+110
View File
@@ -0,0 +1,110 @@
package httphi_test
import (
"fmt"
"io"
"log"
"log/slog"
"net"
"os"
"github.com/soypat/lneto/http/httphi"
"github.com/soypat/lneto/http/httpraw"
)
// ExampleRouter_linux goes over how to setup a linux server using raw linux connections.
// See [ExampleMuxSlice_query_forms_multipart] on how to define handlers for common HTTP processing.
func ExampleRouter() {
// Chrome tends to send ~700 bytes on a typical landing page request.
const requestBuffer = 1024
const numHeaderKV = requestBuffer / 32 //
var mux httphi.MuxSlice
mux.Handle("GET /", func(ex *httphi.Exchange) {
ex.WriteBody([]byte("hello world"))
})
var router httphi.Router
err := router.Configure(httphi.RouterConfig{
FixedNumGoroutines: -1, // Unbounded goroutines and allocations.
RequestHeaderBufferSize: requestBuffer,
ResponseHeaderMinBufferSize: 32, // Shared buffer with Request, not strictly necessary, especially if not sending headers.
RequestNumHeaderKVCap: numHeaderKV,
NormalizeOutgoingKeys: true,
Mux: &mux,
Logger: slog.Default(),
})
if err != nil {
log.Fatal(err)
}
const port = ":8080"
listener, err := net.Listen("tcp", port)
if err != nil {
log.Fatal(err)
}
log.Printf("server up at http://localhost%s", port)
for {
conn, err := listener.Accept()
if err != nil {
log.Fatal(err)
}
err = router.Handle(conn)
if err != nil {
log.Println("httphi.Router failed to handle connection:", err)
}
}
}
// ExampleMuxSlice_query_forms_multipart goes over how to define Handlers (HandleFunc).
// See [ExampleRouter_linux] for how to setup the server.
func ExampleMuxSlice_query_forms_multipart() {
var mux httphi.MuxSlice
mux.Handle("/users/{id}", func(ex *httphi.Exchange) {
userID := ex.PathValue("id")
fmt.Printf("someone requested data for user %s\n", userID)
})
mux.Handle("/query", func(ex *httphi.Exchange) {
// query parameter in URL.
const decodeQuery = true
const queryKey = "search"
valueRaw, present := ex.RequestQueryValue(queryKey)
if !present {
return
}
valueDecoded, present := ex.RequestQueryAppend(nil, queryKey, decodeQuery)
fmt.Printf("got query=%v %s=%s (raw:%s)\n", present, queryKey, valueDecoded, valueRaw)
})
mux.Handle("GET /form", func(ex *httphi.Exchange) {
// Request Body Form.
formbuf := make([]byte, 1024)
var form httpraw.Form
err := ex.RequestParseForm(&form, formbuf)
if err != nil {
ex.WriteHeader(httphi.StatusInternalServerError)
return
}
for i := range form.Len() {
k, v := form.Pair(i)
fmt.Printf("received form value %d: %s=%s\n", i, k, v)
}
})
mux.Handle("GET /file-upload", func(ex *httphi.Exchange) {
// File upload directly onto server using multipart.
formbuf := make([]byte, 1024)
_, err := ex.ReadMultiparts(nil, formbuf, func(hdr *httpraw.MultipartHeader) io.WriteCloser {
fp, err := os.Create(string(hdr.Filename))
if err != nil {
return nil
}
return fp // Will write full file contents to fp.
})
if err != nil {
ex.WriteHeader(httphi.StatusInternalServerError)
return
}
})
}
+707
View File
@@ -0,0 +1,707 @@
package httphi
import (
"io"
"net"
"slices"
"strconv"
"sync/atomic"
"github.com/soypat/lneto"
"github.com/soypat/lneto/http/httpraw"
"github.com/soypat/lneto/internal"
)
// maxStatusLine bounds the response status line: "HTTP/1.1 " + 3 digit code +
// " " + longest [StatusText] + CRLF.
const maxStatusLine = len("HTTP/1.1 ") + 3 + 1 + len("Network Authentication Required") + 2
// Exchange is a single request-response cycle over a connection, playing the
// part of both http.Request and http.ResponseWriter: Request* methods read the
// request, [Exchange.StageHeader] and [Exchange.WriteBody] produce the response.
// A [Router] owns a fixed pool of them, which is what bounds its memory.
//
// Request and response share one buffer, the response header being written over
// the bytes that follow the parsed request header. Read the request body with
// [Exchange.ReadBody] before setting response headers.
type Exchange struct {
acquired atomic.Bool
gen atomic.Uint32
respTopBuf [maxStatusLine]byte
respTopWritten uint8
rawbuf []byte
respHeaderOff uint16
respHeaderLen uint16
reqHdr httpraw.Header
pathValues []pathValue
hijacked bool
rw conn
matchedPattern string
respRemains int
respErr error // Sticky: response is unrecoverable once a write fails.
headerWritten bool
normalizeKeys bool
nextFree *Exchange
readErr error
}
// ExchangeConfig is the memory an [Exchange] is fixed to for the rest of its
// life by [Exchange.Configure]. A [Router] derives one per exchange from its
// [RouterConfig], which is what bounds the router's memory.
type ExchangeConfig struct {
// RawBuf is the single buffer holding the request header, the response
// header and any surplus body. See [Exchange.UnsafeRawBuffer].
RawBuf []byte
// RequestBufferLim reserves the first bytes of RawBuf for the request
// header, the rest being the response. Configure panics if it exceeds RawBuf.
RequestBufferLim int
// NumHeaderKVCap is how many request header fields may be parsed. A request
// carrying more is answered 431, see [httpraw.ErrHeaderTooMany].
NumHeaderKVCap int
// NormalizeOutgoingKeys normalizes staged response header keys as they are
// written, i.e: "content-type" becomes "Content-Type".
NormalizeOutgoingKeys bool
// NoRequestBufferGrowth holds the request header to RequestBufferLim rather
// than growing it. A header outgrowing it is answered 431, see [httpraw.ErrBufferExhausted].
NoRequestBufferGrowth bool
// MaxPathValues is how many wildcards a single pattern may bind, read back with
// [Exchange.PathValue]. A pattern with more never matches, see [SetPathValues].
MaxPathValues int
}
// HijackRaw is a low-level implementation of http.Hijacker interface.
// A Hijack method is not exposed due to heap allocation implications and correctness concerns.
// Below is what an actual implementation may look like:
//
// func (exch *Exchange) Hijack() (net.Conn, *bufio.ReadWriter, error) {
// conn, ok := exch.rw.(net.Conn)
// if !ok {
// return nil, nil, errors.New("net.Conn not implemented")
// }
// _, data, err := exch.HijackRaw(nil)
// if err != nil {
// return nil, nil, err
// }
// var rd *bufio.ReadWriter
// if len(data) > 0 {
// rd = &bufio.ReadWriter{Reader: bufio.NewReader(bytes.NewReader(data))}
// }
// return conn, rd, nil
// }
func (exch *Exchange) HijackRaw(dstBody []byte) (conn, []byte, error) {
data, err := exch.remainingSurplusBody()
if err != nil {
return nil, nil, err
}
exch.hijacked = true
dstBody = append(dstBody, data...)
return exch.rw, dstBody, nil
}
// Configure sets the memory the exchange works with for the rest of its life:
// rawbuf holds the request header, the response header and any surplus body,
// of which the first requestLim bytes are reserved for the request header.
// Panics if requestLim exceeds the buffer. Set normalizeKeys to normalize
// outgoing header keys, i.e: "content-type" to "Content-Type".
func (exch *Exchange) Configure(cfg ExchangeConfig) {
respSize := len(cfg.RawBuf) - cfg.RequestBufferLim
if respSize < 0 {
panic("request lim larger than buffer")
}
exch.rawbuf = cfg.RawBuf
exch.reqHdr.Reset(cfg.RawBuf[:0:cfg.RequestBufferLim], cfg.NumHeaderKVCap)
exch.reqHdr.ConfigBufferGrowth(!cfg.NoRequestBufferGrowth)
exch.normalizeKeys = cfg.NormalizeOutgoingKeys
internal.SliceReuse(&exch.pathValues, cfg.MaxPathValues)
exch.pathValues = exch.pathValues[:cfg.MaxPathValues]
}
// Acquire claims the exchange for conn and resets it to serve a new request,
// reusing the buffer set by [Exchange.Configure]. Returns false if the exchange
// is already serving, in which case conn is untouched.
func (exch *Exchange) Acquire(conn conn) bool {
if !exch.acquired.CompareAndSwap(false, true) {
return false
}
exch.matchedPattern = ""
exch.gen.Add(1)
exch.readErr = nil
exch.respErr = nil
exch.hijacked = false
exch.respTopWritten = 0
exch.respHeaderOff = 0
exch.respHeaderLen = 0
exch.respRemains = 0
exch.rw = conn
exch.headerWritten = false
exch.nextFree = nil
clear(exch.pathValues)
exch.reqHdr.Reset(nil, 0)
return true
}
// Release closes the exchange's connection and frees the exchange for a future
// [Exchange.Acquire]. The connection is left open if the handler took ownership
// of it with [Exchange.HijackRaw].
func (exch *Exchange) Release() {
if !exch.hijacked {
exch.rw.Close()
}
exch.rw = nil
exch.gen.Add(1)
exch.acquired.Store(false)
}
// UnsafeRawBuffer returns the contiguous buffer owned by [Exchange] being used for the request and response.
//
// Writing to it will mangle the entire request header+body and/or any staged response headers.
// Does not return the buffer used for the response first line so can be safely
// written to and used without modifying the staged response first line.
//
// Staging headers will write to this buffer so use mindfully.
// To access only the request header buffer portion use [httpraw.Header.BufferRaw] limited
// to [httpraw.Header.BufferParsed] as returned by [Exchange.RequestHeaderRaw].
// Writing to this section will not change the contents read by [Exchange.ReadBody].
//
// In [Router] context, the size of this buffer is influenced directly by [RouterConfig] HeaderBufferSize fields.
func (exch *Exchange) UnsafeRawBuffer() []byte { return exch.rawbuf }
// StageHeader stages a response header field, written on the first
// [Exchange.FlushHeader], [Exchange.WriteHeader] or [Exchange.WriteBody].
// Returns false and drops the field if the response buffer cannot fit it.
// Has no effect once the header has been written.
func (exch *Exchange) StageHeader(key, value string) (enoughMemory bool) {
if exch.headerWritten {
return false
}
off := int(exch.respHeaderOff) + int(exch.respHeaderLen)
free := len(exch.rawbuf) - off
// Field costs key+':'+value+CRLF, plus the CRLF [Exchange.FlushHeader]
// appends past the last field to close the header block.
if len(key)+len(value)+len(":\r\n")+len("\r\n") > free {
exch.respErr = lneto.ErrBufferFull // Omit writing header back to prevent incomplete response.
return false
}
n := copy(exch.rawbuf[off:], key)
if exch.normalizeKeys {
httpraw.NormalizeHeaderKey(exch.rawbuf[off : off+n])
}
exch.rawbuf[off+n] = ':'
n++
n += copy(exch.rawbuf[off+n:], value)
exch.rawbuf[off+n] = '\r'
exch.rawbuf[off+n+1] = '\n'
n += 2
exch.respHeaderLen += uint16(n)
return true
}
// StageHeaderInt is [Exchange.StageHeader] with an integer value, i.e: Content-Length.
// It formats the value directly into the response buffer without allocating.
// base must be in the range 10..36; lower bases are dropped, no HTTP header
// field value is written below base 10.
func (exch *Exchange) StageHeaderInt(key string, value int64, base int) (enoughMemory bool) {
if exch.headerWritten || base < 10 || base > 36 {
return false
}
off := int(exch.respHeaderOff) + int(exch.respHeaderLen)
free := len(exch.rawbuf) - off
if len(key)+internal.IntLen(value, base)+len(":\r\n")+len("\r\n") > free {
exch.respErr = lneto.ErrBufferFull // Omit writing header back to prevent incomplete response.
return false
}
n := copy(exch.rawbuf[off:], key)
if exch.normalizeKeys {
httpraw.NormalizeHeaderKey(exch.rawbuf[off : off+n])
}
exch.rawbuf[off+n] = ':'
n++
n += len(strconv.AppendInt(exch.rawbuf[off+n:off+n], value, base))
exch.rawbuf[off+n] = '\r'
exch.rawbuf[off+n+1] = '\n'
n += 2
exch.respHeaderLen += uint16(n)
return true
}
// StageStatus prepares the status line for the given code without writing
// it, i.e: "HTTP/1.1 404 Not Found". Codes with no [StatusText] get an empty
// reason phrase. Has no effect once the header has been written.
func (exch *Exchange) StageStatus(code int) {
if code >= 1000 || exch.headerWritten {
return
} else if code == 200 {
// Common case.
exch.respTopWritten = uint8(copy(exch.respTopBuf[:], "HTTP/1.1 200 OK\r\n"))
return
}
n := copy(exch.respTopBuf[:], "HTTP/1.1 ")
n += len(strconv.AppendInt(exch.respTopBuf[n:n], int64(code), 10))
text := StatusText(code)
exch.respTopBuf[n] = ' '
n++
n += copy(exch.respTopBuf[n:], text)
exch.respTopBuf[n] = '\r'
exch.respTopBuf[n+1] = '\n'
exch.respTopWritten = uint8(n + 2)
}
// WriteHeader sends the status line for code along with the staged header
// fields. Only the first call reaches the wire, as in http.ResponseWriter.
func (exch *Exchange) WriteHeader(code int) {
if !exch.headerWritten {
exch.StageStatus(code)
exch.FlushHeader()
}
}
// FlushHeader writes the status line and staged header fields to the connection
// and returns the bytes written, defaulting to a 200 status if none was staged.
// Does nothing if the header was already written.
func (exch *Exchange) FlushHeader() (int, error) {
if exch.respErr != nil {
return 0, exch.respErr
} else if exch.headerWritten {
return 0, nil
}
if exch.respTopWritten == 0 {
exch.StageStatus(200)
}
exch.headerWritten = true
ng, err := exch.rw.Write(exch.respTopBuf[:exch.respTopWritten])
if err != nil {
exch.respErr = err
return ng, err
}
off := int(exch.respHeaderOff)
headers := exch.rawbuf[off : off+int(exch.respHeaderLen)+2]
headers[len(headers)-1] = '\n'
headers[len(headers)-2] = '\r'
ng2, err := exch.rw.Write(headers)
exch.respErr = err
return ng + ng2, err
}
// ExchangeRW is an [io.ReadWriteCloser] view of an [Exchange] wrapping
// [Exchange.ReadBody] and [Exchange.WriteBody] methods.
//
// Exchanges are pooled and reused, so a handle records the exchange generation
// it was taken at and refuses to touch the connection once that exchange moves
// on to another request. Obtain one with [Exchange.ReadWriter].
type ExchangeRW struct {
gen uint32
exch *Exchange
}
// IsValid returns true while the handle still refers to the request it was
// taken from, i.e: false once the exchange was released.
func (rw *ExchangeRW) IsValid() bool {
return rw.gen == rw.exch.gen.Load() && rw.exch.acquired.Load()
}
func (rw *ExchangeRW) validate() error {
if !rw.IsValid() {
return net.ErrClosed
}
return nil
}
// Write writes response body bytes. See [Exchange.WriteBody].
// Fails with [net.ErrClosed] once the handle is no longer valid.
func (rw *ExchangeRW) Write(buf []byte) (int, error) {
if err := rw.validate(); err != nil {
return 0, err
}
return rw.exch.WriteBody(buf)
}
// Read reads request body bytes. See [Exchange.ReadBody].
// Fails with [net.ErrClosed] once the handle is no longer valid.
func (rw *ExchangeRW) Read(buf []byte) (int, error) {
if err := rw.validate(); err != nil {
return 0, err
}
return rw.exch.ReadBody(buf)
}
// Close invalidates this handle so later reads and writes fail. It does not
// close the connection nor end the exchange, both of which the [Router] owns.
func (rw *ExchangeRW) Close() error {
if err := rw.validate(); err != nil {
return err
}
rw.gen--
return nil
}
// ReadWriter fills dst with a stream view of the exchange, valid until the
// exchange is released. The caller owns dst, so a handler may keep one and
// refill it every request without allocating.
func (exch *Exchange) ReadWriter(dst *ExchangeRW) {
dst.gen = exch.gen.Load()
dst.exch = exch
}
// Write writes response body bytes, flushing the header first if the handler
// has not written it yet. Once a write to the connection fails the response is
// unrecoverable and every later write returns that same error, so a body never
// reaches the wire without its header.
func (exch *Exchange) WriteBody(buf []byte) (int, error) {
if exch.respErr != nil {
return 0, exch.respErr
} else if !exch.headerWritten {
_, err := exch.FlushHeader()
if err != nil {
return 0, err // Body must not reach the wire without its header.
}
}
if len(buf) == 0 {
return 0, nil
}
n, err := exch.rw.Write(buf)
exch.respErr = err
return n, err
}
// ReadBody reads the request body into dst, starting with the bytes that
// arrived in the same read as the header and continuing from the connection.
// The exchange does not know the body's length: use Content-Length or the
// transfer encoding to know when to stop reading.
func (exch *Exchange) ReadBody(dst []byte) (n int, _ error) {
if exch.respRemains > 0 {
toRead, err := exch.remainingSurplusBody()
if err != nil {
return 0, err
}
n = copy(dst, toRead)
exch.respRemains -= n
// hand over what already arrived since conn might have
// exhausted data and could block indefinetely.
return n, nil
}
return exch.rw.Read(dst)
}
func (exch *Exchange) remainingSurplusBody() ([]byte, error) {
_, err := exch.reqHdr.Body()
if err != nil {
return nil, err // Returns mangled buffer error if request header has been misused.
}
surplus := exch.rawbuf[exch.reqHdr.BufferParsed():exch.reqHdr.BufferReceived()]
toRead := surplus[len(surplus)-exch.respRemains:]
return toRead, nil
}
// MuxPattern returns the pattern [Mux] matched to the request.
func (exch *Exchange) MuxPattern() string {
return exch.matchedPattern
}
// RequestHeaderRaw returns the parsed request header for access beyond the
// Request* methods, such as [httpraw.Header.ForEach]. Valid until the exchange
// is released, and writing to it corrupts the response.
func (exch *Exchange) RequestHeaderRaw() *httpraw.Header {
return &exch.reqHdr
}
// RequestParseCookie parses the request's key header field into dst, i.e:
// "Cookie". The caller owns dst and its buffer, so it may be reused between
// requests.
func (exch *Exchange) RequestParseCookie(dst *httpraw.Cookie, key string) error {
value := exch.RequestHeader(key)
return dst.ParseBytes(value)
}
// RequestContentType returns the request's Content-Type field value as it
// appears on the wire, parameters included, nil if absent. Test it with
// [httpraw.MediaTypeIs] and pick parameters out with [httpraw.ContentParam].
func (exch *Exchange) RequestContentType() []byte {
// Folded: field names are case insensitive and HTTP/2 mandates lowercase, so
// a proxy translating h2 to h1 sends "content-type", RFC 9110 5.1.
return exch.RequestHeaderRaw().GetFold("Content-Type")
}
// RequestContentLength returns the body length declared by the request's
// Content-Length field. An absent field is signalled with present=false and no error.
// See [httpraw.Header.ContentLength].
func (exch *Exchange) RequestContentLength() (_ int64, present bool, _ error) {
return exch.RequestHeaderRaw().ContentLength()
}
// RequestParseForm reads the request body into buf and parses it as
// "application/x-www-form-urlencoded" into dst. buf is the only storage used and
// the only limit: a body longer than buf is refused with [lneto.ErrBufferFull]
// before a single byte is read, leaving the caller free to answer 413. Pairs are
// left as they arrived, call [httpraw.Form.Decode] to decode them in place.
//
// Unlike http.Request.ParseForm the query string is not folded in, reach it with
// [Exchange.RequestQuery] or [Exchange.RequestQueryAppend]. The body is consumed, so
// call this before [Exchange.ReadBody].
//
// A request with no Content-Length has no body, RFC 9112 6.3, and yields an
// empty form. Use [Exchange.RequestContentLength] to tell that apart from a body
// that arrived empty.
func (exch *Exchange) RequestParseForm(dst *httpraw.Form, buf []byte) error {
if !httpraw.MediaTypeIs(exch.RequestContentType(), "application/x-www-form-urlencoded") {
return errNotFormEncoded
} else if exch.RequestHeaderRaw().GetFold("Transfer-Encoding") != nil {
// Chunked bodies are framed, so reading Content-Length bytes off the
// wire would parse chunk sizes as form data. httpraw does not decode them.
return errUnsupportedTransferCoding
}
length, present, err := exch.RequestContentLength()
if !present {
dst.Reset(nil, 0)
return nil // No length is no body, RFC 9112 6.3.
} else if err != nil {
return err
} else if length > int64(len(buf)) {
return lneto.ErrShortBuffer // Refuse before reading, caller may answer 413.
}
buf = buf[:length]
for read := 0; read < len(buf); {
n, err := exch.ReadBody(buf[read:])
read += n
if n == 0 {
if err == nil {
err = io.ErrNoProgress
} else if err == io.EOF {
break
}
return err
}
}
dst.Reset(buf, 0)
return dst.Parse()
}
// RequestMultipart returns a parser prepared from the boundary parameter of the
// request's Content-Type field. It reads no body: multipart parts declare no
// length, so the caller drives the loop with a buffer it owns and decides per
// part what to keep and when a part has grown too large. See
// [Exchange.ReadMultiparts] for that loop already written.
func (exch *Exchange) RequestMultipart() (mp httpraw.Multipart, err error) {
contentType := exch.RequestContentType()
if !httpraw.MediaTypeIs(contentType, "multipart/form-data") {
return mp, errNotMultipart
}
return mp, mp.SetContentType(contentType)
}
// MultipartSink is a part of a multipart body together with the writer its
// content was streamed to, as appended by [Exchange.ReadMultiparts].
type MultipartSink struct {
// Header identifies the part. Name and Filename are copies, so they
// outlive the read buffer; PartView does not, see [httpraw.MultipartHeader].
Header httpraw.MultipartHeader
// Sink received the part's content and was closed when the part ended,
// nil for a part newSink chose to discard.
Sink io.WriteCloser
}
// ReadMultiparts streams the request's "multipart/form-data" body, writing each
// part to a sink newSink returns for it and appending the pair to dst. buf is the
// only storage used and content is never held whole, so a part of any length
// streams through a buffer the caller sized. dst is appended to and returned, so
// a handler may hand back the slice of a previous request to reuse its parts.
//
// newSink is called once per part, before any of its content is read, and picks
// what to do with it from hdr.Name and hdr.Filename: return a writer to keep the
// part, or nil to discard its content and keep only the header. Each sink is
// closed as soon as its part ends, so Close reports the part arrived whole; on
// error the sink of the part being read is left open for the caller to deal with.
//
// A part header that does not fit buf is refused with [lneto.ErrShortBuffer],
// since reading more can never complete it, leaving the caller free to answer
// 413. The body is consumed, so call this before [Exchange.ReadBody].
func (exch *Exchange) ReadMultiparts(dst []MultipartSink, buf []byte, newSink func(hdr *httpraw.MultipartHeader) io.WriteCloser) (_ []MultipartSink, _ error) {
mp, err := exch.RequestMultipart()
if err != nil {
return dst, err
} else if newSink == nil || len(buf) <= len("\r\n--")+len(mp.Boundary) {
// A buffer that cannot outgrow a delimiter never makes progress.
return dst, lneto.ErrInvalidConfig
}
buflen := 0
for {
// Slot for the next part, given back when the body turns out to be
// over, so its Name and Filename buffers stay available for reuse.
part := internal.SliceReclaim(&dst)
var parsed int
for {
parsed, err = mp.NextHeader(&part.Header, buf[:buflen])
if err != nil {
dst = dst[:len(dst)-1]
if err == io.EOF {
err = nil // Closing delimiter, body done.
}
return dst, err
} else if parsed > 0 {
break // Delimiter and header block complete.
} else if buflen == len(buf) {
dst = dst[:len(dst)-1]
return dst, lneto.ErrShortBuffer // Header longer than buf.
}
// A read that both delivers and fails, as the last of the body
// followed by a hangup does, may still hold what the parser is
// waiting for: take the data and let the error surface on the
// next read.
n, readErr := exch.ReadBody(buf[buflen:])
buflen += n
if n == 0 && readErr != nil {
dst = dst[:len(dst)-1]
return dst, readErr
}
}
part.Sink = newSink(&part.Header)
buflen = copy(buf, buf[parsed:buflen])
for {
bodyLen, restOff, done := mp.NextBody(buf[:buflen])
if bodyLen > 0 && part.Sink != nil {
_, err = part.Sink.Write(buf[:bodyLen])
if err != nil {
return dst, err
}
}
buflen = copy(buf, buf[restOff:buflen])
if done {
break // Buffer now starts at the next part's delimiter.
}
n, readErr := exch.ReadBody(buf[buflen:])
buflen += n
if n == 0 && readErr != nil {
return dst, readErr // Body ended mid part.
}
}
if part.Sink != nil {
if err = part.Sink.Close(); err != nil {
return dst, err
}
}
}
}
// RequestHeader returns the value of the first request header field matching
// key, or nil if absent. Key matching is case sensitive.
func (exch *Exchange) RequestHeader(key string) []byte {
header := exch.RequestHeaderRaw()
return header.Get(key)
}
// RequestTarget returns the request-target (URI) of the request line, i.e:
// "/search?q=go". See [httpraw.Header.RequestTarget].
func (exch *Exchange) RequestTarget() []byte {
return exch.RequestHeaderRaw().RequestTarget()
}
// RequestPath returns the request-target (URI) up to the query string. This is
// what the [Mux] matches on, i.e: "/search" for a request to "/search?q=go".
func (exch *Exchange) RequestPath() []byte {
return exch.RequestHeaderRaw().RequestPath()
}
// RequestQuery returns the request's query string as it appears on the wire.
// Iterate it with [httpraw.NextQueryPair]. See [httpraw.Header.RequestQuery].
func (exch *Exchange) RequestQuery() []byte {
return exch.RequestHeaderRaw().RequestQuery()
}
// RequestQueryValue returns an undecoded view of the first query parameter
// matching key and reports whether it was present. Keys are matched decoded, so
// key "a b" finds "a%20b" and "a+b"; a parameter whose key is a malformed
// escape is skipped. A parameter with no value ("?debug") and one with an empty
// value ("?debug=") are both present with a zero length view.
//
// The view aliases the request buffer, so copy it to outlive the handler or use
// [Exchange.RequestQueryAppend] to decode it out.
func (exch *Exchange) RequestQueryValue(key string) (rawValue []byte, present bool) {
const plusAsSpace = true // Query strings are form encoded, unlike paths.
rawkey, rawval, rest := httpraw.NextQueryPair(exch.RequestQuery())
for ; rawkey != nil; rawkey, rawval, rest = httpraw.NextQueryPair(rest) {
// Compare raw first: a key needing no decoding is the common case, and
// the decoding compare walks the key an escape at a time.
if b2s(rawkey) == key || httpraw.EqualDecodedPercentURL(rawkey, key, plusAsSpace) {
return rawval, true
}
}
return nil, false
}
// RequestQueryAppend appends the value of the first query parameter matching key to
// dst and reports whether the parameter was present, matching keys as
// [Exchange.RequestQueryValue] does. A parameter with no value ("?debug") and
// one with an empty value ("?debug=") are both present with nothing appended.
//
// Values are appended raw unless decoded is set, in which case percent escapes
// and '+' are decoded. A parameter whose value fails to decode is reported
// absent, dst being left as it was rather than holding half a decode.
func (exch *Exchange) RequestQueryAppend(dst []byte, key string, decoded bool) (valueAppended []byte, present bool) {
const plusAsSpace = true // Query strings are form encoded, unlike paths.
rawval, present := exch.RequestQueryValue(key)
if !present || len(rawval) == 0 {
return dst, present
}
if !decoded {
return append(dst, rawval...), true
}
base := len(dst)
dst = slices.Grow(dst, len(rawval))
n, err := httpraw.CopyDecodedPercentURL(dst[base:base+len(rawval)], rawval, plusAsSpace)
if err != nil {
return dst[:base], false // Do not hand back half a decode.
}
return dst[:base+n], true
}
// PathValue returns the segment the request path bound to the wildcard named
// key, or nil if the matched pattern has no such wildcard. It plays the part of
// http.Request.PathValue. See [SetPathValues] for the pattern syntax and for
// which segments a wildcard binds.
//
// sm.Handle("GET /users/{id}", func(exch *httphi.Exchange) {
// id := exch.PathValue("id") // "42" on a GET /users/42.
// })
func (exch *Exchange) PathValue(key string) []byte {
for i := range exch.pathValues {
if exch.pathValues[i].Key == key {
return exch.pathValues[i].Value
} else if exch.pathValues[i].Key == "" {
break // No more keys set.
}
}
return nil
}
// PathValueAppend acceses the result of [Exchange.PathValue] and appends it to dst.
// If decoded is set to true the result will be URL-percent decoded. An error is returned if URL-percent decoding fails.
func (exch *Exchange) PathValueAppend(dst []byte, key string, decoded bool) ([]byte, error) {
const plusAsSpace = true
rawValue := exch.PathValue(key)
if !decoded || len(rawValue) == 0 {
return append(dst, rawValue...), nil
}
base := len(dst)
dst = slices.Grow(dst, len(rawValue))
n, err := httpraw.CopyDecodedPercentURL(dst[base:base+len(rawValue)], rawValue, plusAsSpace)
if err != nil {
return dst[:base], err // Do not hand back half a decode.
}
return dst[:base+n], nil
}
// RequestMethod returns the request line's method, i.e: "GET". See
// [MethodFromBytes] to compare it against a [Method].
func (exch *Exchange) RequestMethod() []byte {
return exch.RequestHeaderRaw().Method()
}
// RequestConnectionClose returns true if the client asked for the connection to
// be closed after this exchange with a "Connection: close" header field.
func (exch *Exchange) RequestConnectionClose() bool {
return exch.RequestHeaderRaw().ConnectionClose()
}
File diff suppressed because it is too large Load Diff
+419
View File
@@ -0,0 +1,419 @@
package httphi
import (
"bytes"
"io"
"strings"
"testing"
"github.com/soypat/lneto/http/httpraw"
)
// Fuzz targets in this file are written so that a stored corpus keeps its
// meaning as the tests grow. Go's corpus files are positional and typed, so an
// input is only reproducible while the code that decodes it stays fixed. Five
// rules keep that true, and edits to this file must obey them:
//
// 1. A target's signature is frozen: func(t *testing.T, ctrl uint64, data []byte).
// Adding, removing or reordering a parameter invalidates every stored entry.
// 2. Control decisions come from ctrl only, wire bytes from data only. Never
// branch on data, never put payload in ctrl: the two axes mutate apart.
// 3. ctrl is a bit field read through absolute shifts and masks named below.
// New knobs claim unused high bits and are only ever appended, never
// renumbered, and the zero value of a knob must decode to the behaviour that
// existed before it was added, so old entries keep replaying as they did.
// 4. No PRNG, no cursor. No rand, no internal.Prand64, and no helper that
// "reads the next N bits" while advancing a position: inserting one draw
// ahead of another shifts every later decision, which is the same corpus
// invalidation a PRNG causes.
// 5. Nothing ambient: no wall clock, no goroutines, no map iteration order.
// Targets drive [Handle] on the calling goroutine; [Router] is not fuzzed
// here precisely because it serves on goroutines of its own.
const (
// Bits 0..3 index segSizes: how the request is split across reads.
ctlSegShift, ctlSegMask = 0, 0xf
// Bits 4..7 and 8..11 size the request and response halves of the buffer.
ctlReqBufShift, ctlRespBufShift, ctlBufMask = 4, 8, 0xf
// Bits 12..13 size the request header field table.
ctlKVCapShift, ctlKVCapMask = 12, 0x3
// Bit 14 normalizes outgoing header keys.
ctlNormalizeShift, ctlNormalizeMask = 14, 0x1
// Bits 15..16 make that many leading writes to the connection fail.
ctlFailWriteShift, ctlFailWriteMask = 15, 0x3
// Bits 17..20 select what the handler does, see the op constants.
ctlHandlerOpsShift, ctlHandlerOpsMask = 17, 0xf
// Bits 21..23 count the header fields the handler stages.
ctlStageCountShift, ctlStageCountMask = 21, 0x7
// Bit 24 asks AppendQuery for decoded values.
ctlDecodeShift, ctlDecodeMask = 24, 0x1
// Bits 25..28 size the scratch buffer handed to form and multipart parsing.
ctlScratchShift, ctlScratchMask = 25, 0xf
// Bit 29 makes the multipart sink discard part content.
ctlDiscardShift, ctlDiscardMask = 29, 0x1
// Next knob starts at bit 30.
)
// Handler operations, selected by the ctlHandlerOps field. Values are frozen:
// a new operation takes the next free bit within the field.
const (
opStageHeaders = 1 << iota
opReadBody
opWriteBody
opHijack
)
// segSizes are the chunk sizes a request may be delivered in, index 0 meaning
// "all at once". Splitting a request mid-CRLF or before a colon is what drives
// [httpraw.Header.TryParse]'s resumption path. Entries may be appended, never
// changed: an existing index must keep splitting exactly as it does today.
var segSizes = [16]int{0, 1, 2, 3, 5, 7, 11, 16, 23, 37, 64, 101, 173, 256, 509, 1024}
// ctlField reads a knob out of ctrl. Absolute shift, no cursor, so knobs are
// independent of one another and of the order they are read in.
func ctlField(ctrl uint64, shift, mask uint64) uint64 {
return (ctrl >> shift) & mask
}
// maxFuzzInput bounds the wire data a target accepts. The filter depends only
// on the input, so it classifies an entry the same way on every run.
const maxFuzzInput = 16 << 10
// fuzzExchange returns an exchange acquired on a connection preloaded with data,
// both sized and segmented by ctrl.
func fuzzExchange(t *testing.T, ctrl uint64, data []byte) (*Exchange, *rwconn) {
t.Helper()
reqBuf := minRequestHeaderBuffer + int(ctlField(ctrl, ctlReqBufShift, ctlBufMask))*8
respBuf := minResponseHeaderBuffer + int(ctlField(ctrl, ctlRespBufShift, ctlBufMask))*8
kvCap := 1 + int(ctlField(ctrl, ctlKVCapShift, ctlKVCapMask))*8
conn := newConn("")
seg := segSizes[ctlField(ctrl, ctlSegShift, ctlSegMask)]
if seg <= 0 || seg >= len(data) {
conn.AddReadable(data)
} else {
for off := 0; off < len(data); off += seg {
conn.AddSegment(string(data[off:min(off+seg, len(data))]))
}
}
// Always hang up: a drained connection that never reports EOF reads (0,nil)
// forever and [Handle] would back off in an unbounded loop.
conn.Hangup()
if n := ctlField(ctrl, ctlFailWriteShift, ctlFailWriteMask); n > 0 {
conn.FailWrites(int(n))
}
exch := newExchange(t, conn, ExchangeConfig{
RawBuf: make([]byte, reqBuf+respBuf),
RequestBufferLim: reqBuf,
NumHeaderKVCap: kvCap,
NormalizeOutgoingKeys: ctlField(ctrl, ctlNormalizeShift, ctlNormalizeMask) != 0,
NoRequestBufferGrowth: true,
})
return exch, conn
}
// checkRequestView asserts the request-target views agree with one another.
func checkRequestView(t *testing.T, exch *Exchange) {
t.Helper()
target, path, query := exch.RequestTarget(), exch.RequestPath(), exch.RequestQuery()
if !bytes.HasPrefix(target, path) {
t.Fatalf("path %q is not a prefix of target %q", path, target)
}
if len(query) > 0 && !bytes.HasSuffix(target, query) {
t.Fatalf("query %q is not a suffix of target %q", query, target)
}
if len(path)+len(query) > len(target) {
t.Fatalf("path %q and query %q exceed target %q", path, query, target)
}
}
// checkResponse asserts that whatever reached the wire is a response a peer
// could parse: a status line this package produced, followed by a header block
// that terminates, and that [httpraw] reads back what it wrote.
func checkResponse(t *testing.T, written string) {
t.Helper()
if written == "" {
return // Hijacked, or a request refused before anything was staged.
}
const proto = "HTTP/1.1 "
if !strings.HasPrefix(written, proto) {
t.Fatalf("response does not open with a status line: %q", written)
}
if len(written) < len(proto)+4 {
t.Fatalf("status line truncated: %q", written)
}
for _, c := range []byte(written[len(proto) : len(proto)+3]) {
if c < '0' || c > '9' {
t.Fatalf("status code is not three digits: %q", written)
}
}
if written[len(proto)+3] != ' ' {
t.Fatalf("status code not followed by a space: %q", written)
}
if !strings.Contains(written, "\r\n\r\n") {
t.Fatalf("header block never terminated: %q", written)
}
var resp httpraw.Header
const asResponse = true
if err := resp.ParseBytes(asResponse, []byte(written)); err != nil {
t.Fatalf("response does not parse back: %s in %q", err, written)
}
}
// handler2Path is a second registration so lookup does not always match on the
// first entry of the mux.
const handler2Path = "/fuzz"
// fuzzMux returns a mux serving handler on the paths the seed corpus requests.
func fuzzMux(handler HandlerFunc) *MuxSlice {
var mux MuxSlice
mux.Handle("/", handler)
mux.Handle(handler2Path, handler)
return &mux
}
// FuzzHandleRequest drives a whole exchange: a request off the wire through
// [Handle], a handler staging and writing a response, and back out to the peer.
func FuzzHandleRequest(f *testing.F) {
addSeeds(f)
f.Fuzz(func(t *testing.T, ctrl uint64, data []byte) {
if len(data) > maxFuzzInput {
return
}
exch, conn := fuzzExchange(t, ctrl, data)
ops := ctlField(ctrl, ctlHandlerOpsShift, ctlHandlerOpsMask)
stage := ctlField(ctrl, ctlStageCountShift, ctlStageCountMask)
Handle(exch, fuzzMux(func(exch *Exchange) {
checkRequestView(t, exch)
if ops&opStageHeaders != 0 {
// Literal fields: staging request bytes would test the caller's
// escaping, not this package's framing.
for i := range stage {
exch.StageHeader("X-Fuzz", "value")
exch.StageHeaderInt("X-Fuzz-Int", int64(i), 10)
}
}
if ops&opReadBody != 0 {
var body [64]byte
for range 64 {
n, err := exch.ReadBody(body[:])
if n == 0 && err != nil {
break
}
}
}
if ops&opWriteBody != 0 {
exch.WriteBody([]byte("fuzz body"))
}
if ops&opHijack != 0 {
exch.HijackRaw(nil)
}
}), nopBackoff)
if ctlField(ctrl, ctlFailWriteShift, ctlFailWriteMask) == 0 {
// A refused write leaves a partial response on purpose, so the
// wire is only well formed when every write got through.
checkResponse(t, conn.ViewWritten())
}
})
}
// FuzzQueryAndForm drives the request-target query and the form-encoded body,
// both of which decode percent escapes in place over caller memory.
func FuzzQueryAndForm(f *testing.F) {
addSeeds(f)
f.Fuzz(func(t *testing.T, ctrl uint64, data []byte) {
if len(data) > maxFuzzInput {
return
}
exch, _ := fuzzExchange(t, ctrl, data)
decoded := ctlField(ctrl, ctlDecodeShift, ctlDecodeMask) != 0
scratchLen := 8 + int(ctlField(ctrl, ctlScratchShift, ctlScratchMask))*16
Handle(exch, fuzzMux(func(exch *Exchange) {
// Every pair the iterator yields must be reachable by name, or the
// two views of the query string disagree.
const maxPairs = 64
key, _, rest := httpraw.NextQueryPair(exch.RequestQuery())
for pairs := 0; key != nil && pairs < maxPairs; pairs++ {
dec := make([]byte, len(key))
n, err := httpraw.CopyDecodedPercentURL(dec, key, true)
if err == nil {
if n > len(key) {
t.Fatalf("decoding key %q grew it to %d bytes", key, n)
}
if _, present := exch.RequestQueryAppend(nil, string(dec[:n]), decoded); !present {
t.Fatalf("query pair %q absent from AppendQuery", key)
}
}
key, _, rest = httpraw.NextQueryPair(rest)
}
var form httpraw.Form
buf := make([]byte, scratchLen)
if err := exch.RequestParseForm(&form, buf); err != nil {
return
}
total := 0
lens := make([]int, form.Len())
for i := range form.Len() {
k, v := form.Pair(i)
lens[i] = len(k) + len(v)
total += lens[i]
}
if total > len(buf) {
t.Fatalf("form pairs span %d bytes of a %d byte buffer", total, len(buf))
}
if err := form.Decode(); err != nil {
return
}
// Decoding replaces escapes in place, so no pair may grow.
for i := range form.Len() {
k, v := form.Pair(i)
if len(k)+len(v) > lens[i] {
t.Fatalf("pair %d grew from %d to %d bytes on decode", i, lens[i], len(k)+len(v))
}
}
}), nopBackoff)
})
}
// countSink counts what a multipart part streamed into it and whether the part
// was closed off.
type countSink struct {
written int
closed bool
}
func (c *countSink) Write(b []byte) (int, error) {
c.written += len(b)
return len(b), nil
}
func (c *countSink) Close() error {
c.closed = true
return nil
}
// FuzzMultipart drives [Exchange.ReadMultiparts], which streams a body of
// unknown length through a buffer the caller sized.
func FuzzMultipart(f *testing.F) {
addSeeds(f)
f.Fuzz(func(t *testing.T, ctrl uint64, data []byte) {
if len(data) > maxFuzzInput {
return
}
exch, _ := fuzzExchange(t, ctrl, data)
discard := ctlField(ctrl, ctlDiscardShift, ctlDiscardMask) != 0
bufLen := 16 + int(ctlField(ctrl, ctlScratchShift, ctlScratchMask))*16
Handle(exch, fuzzMux(func(exch *Exchange) {
var sinks []*countSink
_, err := exch.ReadMultiparts(nil, make([]byte, bufLen), func(hdr *httpraw.MultipartHeader) io.WriteCloser {
if discard {
return nil
}
sink := new(countSink)
sinks = append(sinks, sink)
return sink
})
if err != nil {
return // Sinks are left for the caller to deal with on error.
}
total := 0
for _, sink := range sinks {
if !sink.closed {
t.Fatal("ReadMultiparts returned with a part left open")
}
total += sink.written
}
if total > len(data) {
t.Fatalf("parts streamed %d bytes out of a %d byte request", total, len(data))
}
}), nopBackoff)
})
}
// addSeeds adds the shared seed corpus. Every seed pins ctrl to zero so its
// meaning never moves: knobs added later decode their zero value to the
// behaviour the seed was recorded under.
func addSeeds(f *testing.F) {
f.Helper()
const (
formType = "Content-Type: application/x-www-form-urlencoded\r\n"
mpType = "Content-Type: multipart/form-data; boundary=b0undary\r\n"
)
seeds := []string{
// Well formed traffic, so the fuzzer has somewhere to mutate from.
"GET / HTTP/1.1\r\nHost: h\r\n\r\n",
"GET /fuzz?q=go&n=1 HTTP/1.1\r\nHost: h\r\n\r\n",
"POST /fuzz HTTP/1.1\r\nHost: h\r\n" + formType + "Content-Length: 11\r\n\r\na=1&b=2&c=3",
"POST /fuzz HTTP/1.1\r\nHost: h\r\n" + mpType + "Content-Length: 76\r\n\r\n--b0undary\r\nContent-Disposition: form-data; name=\"f\"\r\n\r\nbody\r\n--b0undary--\r\n",
// Framing the RFC leaves room to disagree over, which is where request
// smuggling lives: two lengths, a length plus a coding, and a coding
// this package does not decode.
"POST / HTTP/1.1\r\nHost: h\r\nContent-Length: 3\r\nContent-Length: 4\r\n\r\nabcd",
"POST / HTTP/1.1\r\nHost: h\r\nContent-Length: 3\r\nTransfer-Encoding: chunked\r\n\r\n1\r\na\r\n0\r\n\r\n",
"POST / HTTP/1.1\r\nHost: h\r\nTransfer-Encoding: chunked\r\n\r\n4\r\nbody\r\n0\r\n\r\n",
// Field names are case insensitive, RFC 9110 5.1, so a lookup that
// misses one of these reads the request differently than the peer wrote it.
"POST / HTTP/1.1\r\nHost: h\r\ncontent-length: 4\r\n\r\nbody",
"POST / HTTP/1.1\r\nHost: h\r\nCONTENT-LENGTH: 4\r\n\r\nbody",
"POST /fuzz HTTP/1.1\r\nHost: h\r\ncontent-type: application/x-www-form-urlencoded\r\ncontent-length: 3\r\n\r\na=1",
// Content-Length values that are not a bare digit string, RFC 9112 6.2.
"POST / HTTP/1.1\r\nHost: h\r\nContent-Length: +5\r\n\r\nbody!",
"POST / HTTP/1.1\r\nHost: h\r\nContent-Length: -1\r\n\r\nbody",
"POST / HTTP/1.1\r\nHost: h\r\nContent-Length: 1 2\r\n\r\nbody",
"POST / HTTP/1.1\r\nHost: h\r\nContent-Length:\r\n\r\nbody",
"POST / HTTP/1.1\r\nHost: h\r\nContent-Length: 9223372036854775808\r\n\r\nbody",
// Line ending and field syntax edges.
"GET / HTTP/1.1\nHost: h\n\n",
"GET / HTTP/1.1\r\nHost: h\rX: y\r\n\r\n",
"GET / HTTP/1.1\r\nHost: h\r\n Continued: fold\r\n\r\n",
"GET / HTTP/1.1\r\nHost: h\r\nX: va\x00lue\r\n\r\n",
"GET / HTTP/1.1\r\nHost: h\r\nNoColon\r\n\r\n",
"GET / HTTP/1.1\r\nHost: h\r\n: novalue\r\n\r\n",
"GET / HTTP/1.1\r\nHost: h\r\nX:\r\n\r\n",
"GET / HTTP/1.1\r\nHost: h\r\n\r\n\r\n",
// Request lines this package tolerates or must refuse.
"GET http://h/abs HTTP/1.1\r\nHost: h\r\n\r\n",
"GET /\r\n\r\n", // HTTP/0.9 simple request, no version.
" GET / HTTP/1.1\r\nHost: h\r\n\r\n",
"GET / HTTP/1.1\r\nHost: h\r\n\r\n",
"/ HTTP/1.1\r\nHost: h\r\n\r\n",
"GET / HTTP/9.9\r\nHost: h\r\n\r\n",
// Percent escapes, including the truncated and the malformed.
"GET /a%20b?a%20b=c%20d HTTP/1.1\r\nHost: h\r\n\r\n",
"GET /?q=%2 HTTP/1.1\r\nHost: h\r\n\r\n",
"GET /?q=%zz HTTP/1.1\r\nHost: h\r\n\r\n",
"GET /%00?a=%00 HTTP/1.1\r\nHost: h\r\n\r\n",
"GET /?a=b=c&&=v&debug HTTP/1.1\r\nHost: h\r\n\r\n",
"GET /?a%3db=1&a+b=2 HTTP/1.1\r\nHost: h\r\n\r\n",
// Multipart bodies that end early or never open a part.
"POST /fuzz HTTP/1.1\r\nHost: h\r\n" + mpType + "Content-Length: 12\r\n\r\n--b0undary\r\n",
"POST /fuzz HTTP/1.1\r\nHost: h\r\n" + mpType + "Content-Length: 14\r\n\r\n--b0undary--\r\n",
"POST /fuzz HTTP/1.1\r\nHost: h\r\nContent-Type: multipart/form-data\r\nContent-Length: 4\r\n\r\nbody",
// Sizes that crowd the buffers: a long target, many fields, and a body
// arriving in the same read as the header it follows.
"GET /" + strings.Repeat("a", 512) + " HTTP/1.1\r\nHost: h\r\n\r\n",
"GET / HTTP/1.1\r\n" + strings.Repeat("X: y\r\n", 200) + "\r\n",
"GET / HTTP/1.1\r\n" + strings.Repeat("k", 512) + ": v\r\n\r\n",
"POST / HTTP/1.1\r\nHost: h\r\nContent-Length: 4\r\n\r\nbodyTRAILING",
// Nothing, and nothing that resembles a request at all.
"",
"\r\n\r\n",
"\x00\x00\x00\x00",
}
for _, seed := range seeds {
f.Add(uint64(0), []byte(seed))
}
}
+313
View File
@@ -0,0 +1,313 @@
package httphi
import (
"bytes"
"io"
"strings"
"unsafe"
"github.com/soypat/lneto"
"github.com/soypat/lneto/http/httpraw"
"github.com/soypat/lneto/internal"
)
// Handle is a extremely low-level HTTP handling method used internally in [Router].
// Requires exchange to be acquired and configured. Will panic if any argument is nil.
// Handle does not close the connection on any outcome: the caller owns it.
// backoff can be set for dealing with non-blocking connections. If backoff set to nil
// then a zero-length-read will result in Handle returning [io.ErrNoProgress].
func Handle(exch *Exchange, mux Mux, backoff lneto.BackoffStrategy) error {
if !exch.acquired.Load() {
return lneto.ErrBadState
}
reqhdr := &exch.reqHdr
reqhdr.Reset(nil, 0) // Assume exchange has been configured and reuse memory.
var consecutiveBackoffs uint
for {
n, err := reqhdr.ReadFromLimited(exch.rw, reqhdr.BufferFree())
if err != nil {
exch.readErr = err
exch.handleError(err)
return err
} else if n == 0 {
if backoff == nil {
return io.ErrNoProgress
}
backoff.Do(consecutiveBackoffs)
consecutiveBackoffs++
continue
}
consecutiveBackoffs = 0
const asRequest = false
needMore, err := reqhdr.TryParse(asRequest)
if needMore {
continue // Request header split across reads, accumulate the rest.
} else if err != nil {
exch.handleError(err)
return err
}
break // Done!
}
// Setup Exchange fields necessary for correct functioning.
parsed := reqhdr.BufferParsed()
exch.respRemains = reqhdr.BufferReceived() - parsed
exch.respHeaderOff = uint16(parsed)
exch.respHeaderLen = 0
proto := b2s(reqhdr.Protocol())
if len(proto) == 0 {
// HTTP/0.9 not tolerated RFC 9112 3.
exch.WriteHeader(int(StatusBadRequest))
return errNoRequestProto
} else if proto != "HTTP/1.1" && proto != "HTTP/1.0" {
// RFC 9112 2.6.
exch.WriteHeader(int(StatusHTTPVersionNotSupported))
return errBadRequestProto
}
// Mux on the request path: the query string is the handler's business.
path := reqhdr.RequestPath()
meth := reqhdr.Method()
matchedPattern, handler := mux.LookupHandler(MethodFromBytes(meth), path, exch.pathValues)
if handler != nil {
exch.matchedPattern = matchedPattern
handler(exch)
if !exch.hijacked {
exch.FlushHeader()
}
} else {
exch.WriteHeader(404)
}
return nil
}
func (exch *Exchange) handleError(err error) {
if err == httpraw.ErrHeaderTooMany || err == httpraw.ErrBufferExhausted || exch.reqHdr.BufferFree() == 0 {
// The peer is owed an answer: no larger buffer is coming, so
// say so instead of dropping the connection, RFC 6585 5.
exch.StageHeader("Content-Length", "0")
exch.WriteHeader(int(StatusRequestHeaderFieldsTooLarge))
}
}
// HandlerFunc serves a single request, playing the part of http.Handler.
// The exchange is only valid for the duration of the call: it is released to
// the router's pool on return, so a handler must not retain it nor any slice it
// handed out.
type HandlerFunc func(ex *Exchange)
// Mux resolves a request to the handler that serves it. [Handle] calls
// LookupHandler with the request-target's path, not the whole target, and
// replies 404 when it returns nil.
type Mux interface {
// LookupHandler matches the requestPath and method to a handler and returns it and the
// pattern it matched. dstPathVals are set to non-zero values by Mux and can later be accessed by [Exchange.PathValue]
// requestPath is a buffer owned by the [Exchange] usually and should not be held after LookupHandler returns.
LookupHandler(get Method, requestPath []byte, dstPathVals []pathValue) (matchedPattern string, handler HandlerFunc)
}
// MuxSlice is a [Mux] backed by a slice of registered endpoints, matched by
// exact path. Lookup is linear in the number of registrations.
type MuxSlice struct {
// TODO: binary search worth it?
_handlers []struct {
method Method
path string
handler HandlerFunc
setPathVal bool
}
}
type pathValue struct {
Key string // owned by mux.
Value []byte // points to raw exchange buffer.
}
// pathSeparator is shared so [SetPathValues] never converts a literal per call.
var pathSeparator = []byte{'/'}
// SetPathValues matches requestPath against pattern and binds its wildcards
// into dstPathVals, read back with [Exchange.PathValue]. Wildcards are whole
// segments as per http.ServeMux: "{name}" takes one non-empty segment,
// "{name...}" the rest including slashes, "{$}" only the path's end, and a
// trailing slash is an anonymous "{...}". i.e: "/b/{bucket}/o/{obj...}".
//
// Unlike ServeMux, segments are compared and bound raw, so "/users/{id}" binds
// "x%2Fy" and not "x/y". Which paths match is unaffected. Bound values alias
// requestPath rather than copy it.
func SetPathValues(dstPathVals []pathValue, pattern string, requestPath []byte) (matched, pathValSliceTooShort bool) {
if len(pattern) == 0 || pattern[0] != '/' || len(requestPath) == 0 || requestPath[0] != '/' {
return false, false
}
pattern, requestPath = pattern[1:], requestPath[1:]
n := 0
for {
if len(pattern) == 0 {
// Nothing left after a slash: an anonymous "..." taking the rest,
// which is why "/files/" matches "/files/a/b" and "/" matches all.
return true, false
}
patSeg, patRest, patMore := strings.Cut(pattern, "/")
reqSeg, reqRest, reqMore := bytes.Cut(requestPath, pathSeparator)
name, isMulti, isWildcard := pathWildcard(patSeg)
switch {
case isWildcard && name == "$":
// Matches the end of the path and nothing else, so it must be the
// last segment of the pattern and leave no path behind.
return !patMore && len(requestPath) == 0, false
case isWildcard && isMulti:
// Takes the remainder including slashes, possibly empty.
if name != "" {
if n >= len(dstPathVals) {
return false, true
}
dstPathVals[n] = pathValue{Key: name, Value: requestPath}
n++
}
return true, false
case isWildcard:
if len(reqSeg) == 0 {
return false, false // One segment means a non-empty one.
}
if n >= len(dstPathVals) {
return false, true
}
dstPathVals[n] = pathValue{Key: name, Value: reqSeg}
n++
default:
if b2s(reqSeg) != patSeg {
return false, false
}
}
if patMore != reqMore {
// One side has a further segment and the other does not, so
// "/health" misses "/health/" and "/files/" misses "/files".
return false, false
} else if !patMore {
return true, false // Both spent on the same segment.
}
pattern, requestPath = patRest, reqRest
}
}
// pathWildcard picks apart a "{name}" or "{name...}" pattern segment. It
// reports ok false for a literal segment, so "/b_{bucket}" is literal text and
// not a wildcard, matching ServeMux's rule that wildcards be whole segments.
func pathWildcard(segment string) (name string, isMulti, ok bool) {
if len(segment) < 2 || segment[0] != '{' || segment[len(segment)-1] != '}' {
return "", false, false
}
name = segment[1 : len(segment)-1]
if rest, found := strings.CutSuffix(name, "..."); found {
return rest, true, true
}
return name, false, true
}
// Reset discards all registered handlers, reusing the backing array and growing
// it to fit capacity registrations.
func (sm *MuxSlice) Reset(capacity int) {
internal.SliceReuse(&sm._handlers, capacity)
}
// LookupHandler returns the handler registered for request path, or nil if none matches.
// The first registration matching both method and uri wins.
func (sm *MuxSlice) LookupHandler(method Method, path []byte, dstPathVals []pathValue) (matched string, _ HandlerFunc) {
for _, endpoint := range sm._handlers {
if endpoint.method != MethUndefined && endpoint.method != method {
continue
}
// Method matches.
if endpoint.setPathVal {
if ok, _ := SetPathValues(dstPathVals, endpoint.path, path); ok {
return endpoint.path, endpoint.handler
}
} else if b2s(path) == endpoint.path {
return endpoint.path, endpoint.handler
}
}
return "", nil
}
// Handle registers handler for reg, either a bare path matching any method or a
// method and path separated by a space, i.e: "/health" or "GET /health".
// Handle does not check for duplicate registrations: the first one added wins.
func (sm *MuxSlice) Handle(optMethodAndPath string, handler HandlerFunc) {
v := internal.SliceReclaim(&sm._handlers)
method := MethUndefined
methodOrURL, url, methodFound := strings.Cut(optMethodAndPath, " ")
if methodFound {
method = MethodFrom(methodOrURL)
} else {
url = methodOrURL
}
v.method = method
v.path = url
v.handler = handler
}
// Method is a HTTP request method, parsed by [MethodFrom].
type Method uint8
const (
MethUndefined Method = iota // undefined
MethGet // GET
// lol.
MethHead // HEAD
MethPost // POST
MethPut // PUT
// RFC 5789
MethPatch // PATCH
MethDelete // DELETE
MethConnect // CONNECT
MethOptions // OPTIONS
MethTrace // TRACE
MethUnknown // unknown
)
// MethodFrom returns the [Method] matching meth, [MethUndefined] if meth is
// empty and [MethUnknown] if it names a method this package does not know.
// Comparison is case sensitive: methods are uppercase, RFC 9110 9.1.
func MethodFrom(meth string) (res Method) {
if len(meth) == 0 {
return MethUndefined
}
switch meth {
case "GET":
res = MethGet
case "HEAD":
res = MethHead
case "POST":
res = MethPost
case "PUT":
res = MethPut
case "PATCH":
res = MethPatch
case "DELETE":
res = MethDelete
case "CONNECT":
res = MethConnect
case "OPTIONS":
res = MethOptions
case "TRACE":
res = MethTrace
default:
res = MethUnknown
}
return res
}
// MethodFromBytes is a [MethodFrom] wrapper with bytes argument instead of string.
func MethodFromBytes(meth []byte) (res Method) {
if len(meth) == 0 {
return MethUndefined
}
return MethodFrom(b2s(meth))
}
// b2s converts byte slice to a string without memory allocation.
// See https://groups.google.com/forum/#!msg/Golang-Nuts/ENgbUzYvCuU/90yGx7GUAgAJ .
func b2s(b []byte) string {
return unsafe.String(unsafe.SliceData(b), len(b))
}
+222
View File
@@ -0,0 +1,222 @@
package httphi
import (
"strings"
"testing"
)
// SetPathValues must agree with net/http.ServeMux on which patterns match which
// paths and what each wildcard binds to. Every case below was taken from a run
// against a real ServeMux, so this table is an oracle, not a guess.
//
// The one documented deviation is percent-decoding: ServeMux unescapes segments
// before matching and binding, this does not. See TestSetPathValuesEscaping.
func TestSetPathValues(t *testing.T) {
for _, test := range []struct {
pattern string
path string
match bool
want string // "name=value" pairs joined by "|", in bind order.
}{
// Single segment wildcard binds exactly one non-empty segment.
{pattern: "/users/{id}", path: "/users/42", match: true, want: "id=42"},
{pattern: "/users/{id}", path: "/users/42/x", match: false},
{pattern: "/users/{id}", path: "/users/", match: false},
{pattern: "/users/{id}", path: "/users", match: false},
{pattern: "/users/{id}/edit", path: "/users/42/edit", match: true, want: "id=42"},
{pattern: "/{a}/{b}", path: "/x/y", match: true, want: "a=x|b=y"},
// "..." swallows the remainder, slashes included, and may bind empty.
{pattern: "/b/{bucket}/o/{obj...}", path: "/b/bk/o/a/b/c", match: true, want: "bucket=bk|obj=a/b/c"},
{pattern: "/b/{bucket}/o/{obj...}", path: "/b/bk/o/", match: true, want: "bucket=bk|obj="},
{pattern: "/b/{bucket}/o/{obj...}", path: "/b/bk/o", match: false},
{pattern: "/files/{p...}", path: "/files/", match: true, want: "p="},
{pattern: "/files/{p...}", path: "/files", match: false},
// {$} matches only the end of the path.
{pattern: "/{$}", path: "/", match: true},
{pattern: "/{$}", path: "/x", match: false},
{pattern: "/a/{$}", path: "/a/", match: true},
{pattern: "/a/{$}", path: "/a", match: false},
{pattern: "/a/{$}", path: "/a/b", match: false},
// A trailing slash is an anonymous "..." wildcard, binding nothing.
{pattern: "/files/", path: "/files/a/b", match: true},
{pattern: "/files/", path: "/files/", match: true},
{pattern: "/files/", path: "/files", match: false},
{pattern: "/", path: "/anything/at/all", match: true},
// Literal patterns match exactly, trailing slash included.
{pattern: "/health", path: "/health", match: true},
{pattern: "/health", path: "/health/", match: false},
// An empty segment never satisfies a single wildcard.
{pattern: "/a/{x}/b", path: "/a//b", match: false},
} {
t.Run(test.pattern+"__"+test.path, func(t *testing.T) {
vals := make([]pathValue, 8)
match, tooShort := SetPathValues(vals, test.pattern, []byte(test.path))
if tooShort {
t.Fatal("8 slots must be enough for these patterns")
}
if match != test.match {
t.Fatalf("want match=%v, got %v", test.match, match)
}
if !match {
return
}
if got := renderPathValues(vals); got != test.want {
t.Errorf("want %q, got %q", test.want, got)
}
})
}
}
// Bound values must alias the request path buffer rather than copy it: the
// exchange owns that memory and a copy would allocate per request.
func TestSetPathValuesAliasesRequestBuffer(t *testing.T) {
path := []byte("/users/42/edit")
vals := make([]pathValue, 4)
match, _ := SetPathValues(vals, "/users/{id}/edit", path)
if !match {
t.Fatal("want match")
}
if string(vals[0].Value) != "42" {
t.Fatalf("want id=42, got %q", vals[0].Value)
}
// Mutating the request buffer must show through the bound value.
path[7] = '9'
if string(vals[0].Value) != "92" {
t.Errorf("value must alias the request buffer, got %q", vals[0].Value)
}
}
// A destination too small to hold every wildcard must say so rather than bind a
// partial set or write out of range.
func TestSetPathValuesSliceTooShort(t *testing.T) {
match, tooShort := SetPathValues(make([]pathValue, 1), "/{a}/{b}", []byte("/x/y"))
if !tooShort {
t.Error("want pathValSliceTooShort for 2 wildcards in 1 slot")
}
if match {
t.Error("want match=false when values could not be bound")
}
// A pattern that binds nothing needs no slots at all.
match, tooShort = SetPathValues(nil, "/health", []byte("/health"))
if !match || tooShort {
t.Errorf("want match with no slots needed, got match=%v tooShort=%v", match, tooShort)
}
}
// Percent escapes are compared and bound raw. ServeMux unescapes segment by
// segment, so "/users/x%2Fy" binds id="x/y" there and id="x%2Fy" here. Matching
// agrees either way; only the bound bytes differ.
func TestSetPathValuesEscaping(t *testing.T) {
vals := make([]pathValue, 4)
if match, _ := SetPathValues(vals, "/a%2Fb/{x}", []byte("/a%2Fb/v")); !match {
t.Error("want literal escape in pattern to match the same bytes in path")
}
vals = make([]pathValue, 4)
match, _ := SetPathValues(vals, "/users/{id}", []byte("/users/x%2Fy"))
if !match {
t.Fatal("want match")
}
if got := string(vals[0].Value); got != "x%2Fy" {
t.Errorf("want raw %q, got %q", "x%2Fy", got)
}
}
// renderPathValues joins the bound pairs for comparison, stopping at the first
// unused slot.
func renderPathValues(vals []pathValue) string {
var sb strings.Builder
for _, v := range vals {
if v.Key == "" {
break
}
if sb.Len() > 0 {
sb.WriteByte('|')
}
sb.WriteString(v.Key)
sb.WriteByte('=')
sb.Write(v.Value)
}
return sb.String()
}
// Matching a request must not allocate: keys alias the mux's pattern and values
// alias the request buffer, so nothing is copied per request.
func TestSetPathValuesNoAlloc(t *testing.T) {
vals := make([]pathValue, 8)
path := []byte("/b/bk/o/a/b/c")
allocs := testing.AllocsPerRun(100, func() {
SetPathValues(vals, "/b/{bucket}/o/{obj...}", path)
})
if allocs != 0 {
t.Fatalf("SetPathValues allocated %v times, want 0", allocs)
}
}
// pathValueMux binds one wildcard pattern, standing in for a [Mux] that
// supports them until MuxSlice sets setPathVal.
type pathValueMux struct {
pattern string
handler HandlerFunc
}
func (m *pathValueMux) LookupHandler(method Method, path []byte, dst []pathValue) (string, HandlerFunc) {
if ok, _ := SetPathValues(dst, m.pattern, path); ok {
return m.pattern, m.handler
}
return "", nil
}
// A wildcard bound by one request must not be readable by the next request the
// same pooled exchange serves. A literal pattern binds nothing, so it never
// overwrites the previous request's slots, and the values it would leak alias a
// buffer the new request has already overwritten.
func TestExchangePathValueClearedBetweenRequests(t *testing.T) {
exch := new(Exchange)
exch.Configure(ExchangeConfig{
RawBuf: make([]byte, 2048), RequestBufferLim: 1024,
NumHeaderKVCap: defaultNumHeaderKVCap, MaxPathValues: 4,
})
// First request binds id=42 off a wildcard pattern.
var gotFirst string
wildcard := &pathValueMux{pattern: "/users/{id}", handler: func(e *Exchange) {
gotFirst = string(e.PathValue("id"))
e.WriteHeader(200)
}}
conn := newConn("GET /users/42 HTTP/1.1\r\nHost: h\r\n\r\n")
conn.Hangup()
if !exch.Acquire(conn) {
t.Fatal("fresh exchange failed to acquire")
}
if err := Handle(exch, wildcard, nopBackoff); err != nil {
t.Fatalf("first request: %s", err)
}
exch.Release()
if gotFirst != "42" {
t.Fatalf("want id=42 bound on the first request, got %q", gotFirst)
}
// Second request matches a literal pattern, which binds nothing at all.
var leaked []byte
var sm MuxSlice
sm.Handle("/health", func(e *Exchange) {
leaked = e.PathValue("id")
e.WriteHeader(200)
})
conn2 := newConn("GET /health HTTP/1.1\r\nHost: h\r\n\r\n")
conn2.Hangup()
if !exch.Acquire(conn2) {
t.Fatal("released exchange failed to re-acquire")
}
if err := Handle(exch, &sm, nopBackoff); err != nil {
t.Fatalf("second request: %s", err)
}
if leaked != nil {
t.Errorf("want no path value on a literal route, got id=%q from the previous request", leaked)
}
}
+403
View File
@@ -0,0 +1,403 @@
package httphi
import (
"errors"
"io"
"log/slog"
"math"
"sync"
"sync/atomic"
"time"
"github.com/soypat/lneto"
"github.com/soypat/lneto/internal"
)
//go:generate stringer -type Method -linecomment -output stringers.go
// reconfigureWait bounds how long [Router.Configure] waits for the previous
// generation to stop serving before reusing its exchange buffers.
const reconfigureWait = 10 * time.Millisecond
var (
errNoRequestProto = errors.New("httphi: request line with no HTTP version")
errBadRequestProto = errors.New("httphi: unsupported HTTP version in request line")
errBusyExchanges = errors.New("httphi: exchanges still serving, cannot reuse their buffers")
errRouterTornDown = errors.New("httphi: router torn down, configure it before serving")
errNotFormEncoded = errors.New("httphi: request body is not application/x-www-form-urlencoded")
errNotMultipart = errors.New("httphi: request body is not multipart/form-data")
errUnsupportedTransferCoding = errors.New("httphi: transfer coding not decoded, read the body directly")
)
type conn = io.ReadWriteCloser
// Router serves HTTP connections handed to it with [Router.Handle], routing
// each request to a handler found through its [Mux]. It plays the part of
// http.Server minus the listening: accepting connections is the caller's job,
// which is what lets the same router run over a TCP stack, a socket or a test
// pipe.
//
// A Router owns the exchanges and goroutines that serve connections and sizes
// both at [Router.Configure] time, so serving load costs no allocation and
// bounded memory. Connections arriving with nothing left to serve them are
// refused rather than queued, see [Router.Handle].
//
// Methods are safe for concurrent use. The zero value is not usable: configure
// it first.
type Router struct {
mu sync.Mutex
gen atomic.Uint32
numGoro int
reqBuf int
respBuf int
reqNumHeaderCap int
maxPathValues int
normalizeKeys bool
pendingConns chan job
mux Mux
globbuf []byte
exchs []Exchange
freeList *Exchange
log *slog.Logger
}
// job is a connection waiting on an exchange for a worker goroutine to serve it.
type job struct {
exch *Exchange
}
// RouterConfig configures a [Router]. See [Router.Configure].
type RouterConfig struct {
// FixedNumGoroutines must be set to either -1 (freely allocate new goroutines) or to the number of goroutines
// to spawn on [Router.Configure] being called.
FixedNumGoroutines int
// RequestHeaderBufferSize determines the buffer allocated
// for processing request HTTP headers including request-target (URI), protocol and key/value pairs.
RequestHeaderBufferSize int
// ResponseHeaderMinBufferSize determines buffer allocated for processing response headers.
// Response buffer will reuse unused request memory so this is not a strict limit.
// "HTTP/1.1 200 OK\r\n" does not count towards this memory, only actual Headers key/value pairs use this memory.
// After memory is fully consumed [Exchange.StageHeader] will not append more headers.
ResponseHeaderMinBufferSize int
// Number of request header key/value pairs to parse before failing and returning [StatusRequestHeaderFieldsTooLarge].
RequestNumHeaderKVCap int
// Sets maximum number of PathValue pairs that can be set on an exchange. Accessed via [Exchange.PathValue].
MaxPathValues int
// NormalizeOutgoingKeys normalizes response header field keys as they are
// staged, i.e: "content-type" becomes "Content-Type".
NormalizeOutgoingKeys bool
// MaxAwaitingConns is the depth of the queue connections wait in for a free
// goroutine. [Router.Handle] drops connections once it is full. Required and
// must be non-zero when running a fixed number of goroutines, unused otherwise.
MaxAwaitingConns int
// Mux resolves each request's method and path to the handler serving it. Required.
Mux Mux
// Logger receives failed exchanges. Optional, nil disables logging.
Logger *slog.Logger
}
const (
// minRequestHeaderBuffer is the smallest request buffer [httpraw.Header]
// accepts with buffer growth disabled, which is how exchanges are configured.
minRequestHeaderBuffer = 32
// minResponseHeaderBuffer is the room [Exchange.FlushHeader] needs for the
// CRLF closing the header block, written even when no field was staged.
minResponseHeaderBuffer = len("\r\n")
// maxExchangeBuffer bounds an exchange's whole buffer: [Exchange] indexes it
// with uint16 offsets, so a larger one would be addressed truncated.
maxExchangeBuffer = math.MaxUint16
)
// Validate returns a non-nil error if the configuration cannot be used to
// configure a [Router].
func (cfg RouterConfig) Validate() error {
workerMode := cfg.workerMode()
switch {
case cfg.Mux == nil,
!workerMode && cfg.FixedNumGoroutines != -1,
workerMode && cfg.MaxAwaitingConns <= 0,
cfg.RequestNumHeaderKVCap <= 0,
cfg.RequestHeaderBufferSize < minRequestHeaderBuffer,
cfg.ResponseHeaderMinBufferSize < minResponseHeaderBuffer,
cfg.ResponseHeaderMinBufferSize > maxExchangeBuffer,
cfg.RequestHeaderBufferSize > maxExchangeBuffer-cfg.ResponseHeaderMinBufferSize:
return lneto.ErrInvalidConfig
}
if workerMode {
// Buffer sizes are bounded by maxExchangeBuffer above so the sum cannot
// overflow; the products below allocate and can.
exchBuf := cfg.RequestHeaderBufferSize + cfg.ResponseHeaderMinBufferSize
if cfg.FixedNumGoroutines > math.MaxInt/exchBuf ||
cfg.FixedNumGoroutines > math.MaxInt/cfg.RequestNumHeaderKVCap {
return lneto.ErrInvalidConfig
}
}
return nil
}
func (cfg RouterConfig) workerMode() bool {
return cfg.FixedNumGoroutines > 0
}
// Shutdown stops the router once in-flight exchanges finish; until
// [Router.Configure] is called again, connections are refused with a
// non-nil error. Configure calls it before installing a new generation.
func (r *Router) Shutdown() {
r.mu.Lock()
defer r.mu.Unlock()
r.shutdownLocked()
}
// shutdownLocked is [Router.Shutdown] but without locking requirement.
func (r *Router) shutdownLocked() {
r.gen.Add(1)
if r.pendingConns != nil {
close(r.pendingConns)
r.pendingConns = nil
}
}
// Configure prepares the router to serve connections, tearing down the previous
// generation of goroutines and exchanges first. In worker mode it spawns
// [RouterConfig.FixedNumGoroutines] goroutines and allocates their exchange
// buffers up front, so the router's memory use does not grow with load.
//
// Configure may be called on a serving router, but since the exchange buffers
// are reused it waits for connections in flight to finish and fails with a
// non-nil error rather than reconfigure buffers still being served from.
func (r *Router) Configure(cfg RouterConfig) error {
if err := cfg.Validate(); err != nil {
return err
}
r.mu.Lock()
defer r.mu.Unlock()
r.shutdownLocked()
gen := r.gen.Load()
numgoro := cfg.FixedNumGoroutines
workerMode := cfg.workerMode()
r.reqNumHeaderCap = cfg.RequestNumHeaderKVCap
r.reqBuf = cfg.RequestHeaderBufferSize
r.respBuf = cfg.ResponseHeaderMinBufferSize
r.mux = cfg.Mux
r.log = cfg.Logger
r.maxPathValues = cfg.MaxPathValues
r.normalizeKeys = cfg.NormalizeOutgoingKeys
// Freelist entries were sized by the outgoing configuration: recycling one
// would serve a request with buffer limits cfg never asked for.
r.freeList = nil
if !workerMode {
r.numGoro = 0
r.pendingConns = nil
return nil
}
if workerMode {
jobqueue := make(chan job, cfg.MaxAwaitingConns)
if gen > 1 {
// Exchange buffers below are reused: the previous generation must be
// done serving before they may be handed to the new one.
err := r.awaitIdleExchangesLocked(reconfigureWait)
if err != nil {
return err
}
}
internal.SliceReuse(&r.exchs, numgoro)
r.exchs = r.exchs[:numgoro]
rawBuflen := cfg.RequestHeaderBufferSize + cfg.ResponseHeaderMinBufferSize
internal.SliceReuse(&r.globbuf, numgoro*rawBuflen)
for i := range numgoro {
// TODO exchange buffer alloc
goff := i * rawBuflen
// r.globbuf[goff:goff+rawBuflen], cfg.RequestHeaderBufferSize, cfg.RequestNumHeaderCap, cfg.NormalizeOutgoingKeys
r.exchs[i].Configure(ExchangeConfig{
RawBuf: r.globbuf[goff : goff+rawBuflen],
RequestBufferLim: cfg.RequestHeaderBufferSize,
NumHeaderKVCap: cfg.RequestNumHeaderKVCap,
NormalizeOutgoingKeys: cfg.NormalizeOutgoingKeys,
NoRequestBufferGrowth: true, // Hard memory limit.
MaxPathValues: cfg.MaxPathValues,
})
go r.goroWorker(gen, jobqueue, cfg.Mux)
}
r.pendingConns = jobqueue
r.numGoro = numgoro
}
return nil
}
// awaitIdleExchangesLocked waits up to maxWait for exchanges of the previous
// generation to finish serving so their buffers may be reused. Requires r.mu
// held; the lock is released while waiting since [Router.freeExch] needs it to
// free the exchanges being waited on.
func (r *Router) awaitIdleExchangesLocked(maxWait time.Duration) error {
const pollInterval = time.Millisecond
for waited := time.Duration(0); ; waited += pollInterval {
busy := false
for i := range r.exchs {
if r.exchs[i].acquired.Load() {
busy = true
break
}
}
if !busy {
return nil
} else if waited >= maxWait {
return errBusyExchanges
}
r.mu.Unlock()
time.Sleep(pollInterval)
r.mu.Lock()
}
}
// Handle takes ownership of conn and serves one exchange on it, closing it when
// done. It does not block on the exchange: the connection is handed to a
// goroutine and Handle returns immediately.
//
// Handle returns [lneto.ErrExhausted] when no exchange is free,
// [lneto.ErrPacketDrop] when the queue of connections awaiting a goroutine is
// full, and an error when the router's goroutines have been torn down. On every
// one of them conn is left untouched and unclosed for the caller to dispose of:
// refusing connections is how a router with fixed memory applies backpressure.
func (r *Router) Handle(conn io.ReadWriteCloser) error {
// Exchange acquisition and the configuration it is served with must be read
// under the same lock: [Router.Configure] may run concurrently.
r.mu.Lock()
numGoro, mux := r.numGoro, r.mux
gen := r.gen.Load() // Generation whose buffers the exchange below is sized by.
if numGoro > 0 && r.pendingConns == nil {
// Goroutines torn down: refuse before claiming an exchange.
r.mu.Unlock()
return errRouterTornDown
}
exch := r.getExchLocked(conn)
if exch == nil {
r.mu.Unlock()
return lneto.ErrExhausted
} else if numGoro == 0 {
r.mu.Unlock()
go r.goroHandle(gen, exch, mux)
return nil
}
// Enqueue under the lock: [Router.Configure] closes pendingConns while
// holding it, so an unlocked send could land on a closed channel. The send
// never blocks, so holding the lock cannot stall a worker.
var enqueued bool
select {
case r.pendingConns <- job{exch: exch}:
enqueued = true
default:
// pendingConns cannot store another Conn, we drop and return error.
exch.acquired.Store(false) // release.
}
r.mu.Unlock()
if enqueued {
return nil
}
return lneto.ErrPacketDrop
}
func (r *Router) goroWorker(gen uint32, queue chan job, mux Mux) {
for job := range queue {
exch := job.exch
if exch == nil {
panic("httplo: unreachable nil job")
} else if gen != r.gen.Load() {
// Not released with freeExch since generation torn down,
// new buffer may have been allocated for Exchanges.
exch.Release()
continue
}
r.goroHandle(gen, exch, mux)
}
}
func (r *Router) goroHandle(gen uint32, exch *Exchange, mux Mux) {
defer r.freeExch(gen, exch)
err := Handle(exch, mux, nil)
if err != nil {
if exch.readErr != nil {
r.error("goroHandle:ReadFromLimited", slog.String("err", err.Error()))
} else {
r.error("goroHandle:TryParse?", slog.String("err", err.Error()))
}
}
}
// freeExch releases exch and offers it to the freelist for reuse. gen is the
// generation exch was acquired under: an exchange outliving its generation is
// dropped, its buffers being sized by a configuration the router no longer
// serves and possibly carved out of a globbuf it no longer owns.
func (r *Router) freeExch(gen uint32, exch *Exchange) {
const freelistMaxDepth = 5
r.mu.Lock()
if gen != r.gen.Load() {
exch.Release()
r.mu.Unlock()
return
}
depth := 0
for node := r.freeList; node != nil && depth < freelistMaxDepth; node = node.nextFree {
depth++
}
if depth < freelistMaxDepth {
// Push at head: appending at the tail would drop every node past the
// depth limit instead of dropping the exchange we cannot store.
exch.nextFree = r.freeList
r.freeList = exch
} else {
exch.nextFree = nil // Freelist full, exchange is dropped.
}
exch.Release()
r.mu.Unlock()
}
// getExchLocked returns an exchange acquired on conn. Requires r.mu held.
func (r *Router) getExchLocked(conn conn) (exch *Exchange) {
if r.freeList != nil {
// Successor must be read before Acquire: Acquire clears nextFree, so
// popping afterwards would truncate the freelist to the popped node.
next := r.freeList.nextFree
if r.freeList.Acquire(conn) {
exch = r.freeList
r.freeList = next
return exch
}
}
// Unbounded mode is stored as zero, see [Router.Configure].
workerMode := r.numGoro > 0
if workerMode {
for i := range r.exchs {
if r.exchs[i].Acquire(conn) {
return &r.exchs[i]
}
}
} else {
// Unbounded growth mode when r.numGoro==0.
exch := new(Exchange)
exch.Configure(ExchangeConfig{
RawBuf: make([]byte, r.respBuf+r.reqBuf),
RequestBufferLim: r.reqBuf,
NumHeaderKVCap: r.reqNumHeaderCap,
NormalizeOutgoingKeys: r.normalizeKeys,
NoRequestBufferGrowth: true,
MaxPathValues: r.maxPathValues,
})
exch.Acquire(conn) // Fresh exchange, CAS cannot fail.
return exch
}
return nil
}
func (r *Router) error(msg string, attrs ...slog.Attr) {
internal.LogAttrs(r.log, slog.LevelError, msg, attrs...)
}
func (r *Router) info(msg string, attrs ...slog.Attr) {
internal.LogAttrs(r.log, slog.LevelInfo, msg, attrs...)
}
+501
View File
@@ -0,0 +1,501 @@
package httphi
import (
"bytes"
"context"
"io"
"net"
"strings"
"sync"
"testing"
"time"
)
// rwconn is a in-memory conn. The router handles connections on another
// goroutine so every field is guarded; onClose lets tests await the handler.
//
// A drained rwconn reads (0,nil) as a live socket with no data pending would,
// so a partial request can be completed with AddReadable mid-handling. Tests
// that need the peer to hang up call Hangup.
type rwconn struct {
mu sync.Mutex
readable bytes.Buffer
segments []string
written bytes.Buffer
closed bool
hangup bool
failWr int
onClose chan struct{}
deadline time.Time
}
// newConn returns a conn preloaded with request and whose Close is observable
// with [rwconn.AwaitClose].
func newConn(request string) *rwconn {
r := &rwconn{onClose: make(chan struct{})}
r.AddReadable([]byte(request))
return r
}
// AddSegment queues data delivered on a later read, once everything already
// pending has been read. Models a request split over several TCP segments
// without depending on goroutine scheduling.
func (r *rwconn) AddSegment(b string) {
r.mu.Lock()
defer r.mu.Unlock()
r.segments = append(r.segments, b)
}
// FailWrites makes the next n writes fail, as a conn refusing further data
// would. Later writes succeed.
func (r *rwconn) FailWrites(n int) {
r.mu.Lock()
defer r.mu.Unlock()
r.failWr = n
}
// Hangup makes reads past the pending data return [io.EOF], as a peer that
// closed its side of the connection would.
func (r *rwconn) Hangup() {
r.mu.Lock()
defer r.mu.Unlock()
r.hangup = true
}
func (r *rwconn) Close() error {
r.mu.Lock()
defer r.mu.Unlock()
if !r.closed {
r.closed = true
if r.onClose != nil {
close(r.onClose)
}
}
return nil
}
// AwaitClose blocks until the connection is closed by its handler or timeout elapses.
func (r *rwconn) AwaitClose(t *testing.T, timeout time.Duration) {
t.Helper()
select {
case <-r.onClose:
case <-time.After(timeout):
t.Fatal("timed out awaiting connection close by handler")
}
}
func (r *rwconn) Read(b []byte) (int, error) {
r.mu.Lock()
defer r.mu.Unlock()
if r.closed {
return 0, net.ErrClosed
} else if r.deadlineExceeded() {
return 0, context.DeadlineExceeded
} else if r.readable.Len() == 0 {
if len(r.segments) > 0 {
r.readable.WriteString(r.segments[0])
r.segments = r.segments[1:]
return r.readable.Read(b)
}
if r.hangup {
return 0, io.EOF
}
return 0, nil // No data pending, handler backs off and retries.
}
return r.readable.Read(b)
}
func (r *rwconn) Write(b []byte) (int, error) {
r.mu.Lock()
defer r.mu.Unlock()
if r.closed {
return 0, net.ErrClosed
} else if r.deadlineExceeded() {
return 0, context.DeadlineExceeded
} else if r.failWr > 0 {
r.failWr--
return 0, io.ErrShortWrite
}
return r.written.Write(b)
}
func (r *rwconn) deadlineExceeded() bool {
return !r.deadline.IsZero() && time.Since(r.deadline) > 0
}
func (r *rwconn) AddReadable(b []byte) {
r.mu.Lock()
defer r.mu.Unlock()
r.readable.Write(b)
}
// SetDeadline makes reads and writes past t fail, as a conn with a read
// deadline set would.
func (r *rwconn) SetDeadline(t time.Time) {
r.mu.Lock()
defer r.mu.Unlock()
r.deadline = t
}
// IsClosed reports whether the connection was closed by its handler.
func (r *rwconn) IsClosed() bool {
r.mu.Lock()
defer r.mu.Unlock()
return r.closed
}
func (r *rwconn) ViewWritten() string {
r.mu.Lock()
defer r.mu.Unlock()
return r.written.String()
}
var _ Mux = (*MuxSlice)(nil)
func configSynchronousRouter(t *testing.T, router *Router, bufferSize int, mux Mux) {
err := router.Configure(RouterConfig{
FixedNumGoroutines: -1,
Mux: mux,
RequestHeaderBufferSize: bufferSize,
RequestNumHeaderKVCap: 16,
ResponseHeaderMinBufferSize: bufferSize,
})
if err != nil {
t.Fatal(err)
}
}
func TestRouterGet(t *testing.T) {
const bufferSize = 1024
const expectResponse = "its time"
var (
sm MuxSlice
router Router
)
sm.Handle("GET /", staticPage(t, expectResponse))
configSynchronousRouter(t, &router, bufferSize, &sm)
conn := newConn("GET / HTTP/1.1\r\nHost: tinygo.org\r\n\r\n")
err := router.Handle(conn)
if err != nil {
t.Fatal(err)
}
conn.AwaitClose(t, time.Second)
got := conn.ViewWritten()
if !strings.HasPrefix(got, "HTTP/1.1 200 OK\r\n") {
t.Errorf("want 200 status line, got %q", got)
}
if !strings.HasSuffix(got, expectResponse) {
t.Errorf("want body %q at end of response, got %q", expectResponse, got)
}
}
// The handler observes the request line and header fields the router parsed.
func TestRouterRequestVisibleToHandler(t *testing.T) {
const bufferSize = 1024
var (
sm MuxSlice
router Router
)
var gotMethod, gotURI, gotHost string
sm.Handle("GET /index.html", func(ex *Exchange) {
gotMethod = string(ex.RequestMethod())
gotURI = string(ex.RequestTarget())
gotHost = string(ex.RequestHeader("Host"))
ex.WriteHeader(200)
})
configSynchronousRouter(t, &router, bufferSize, &sm)
conn := newConn("GET /index.html HTTP/1.1\r\nHost: tinygo.org\r\n\r\n")
if err := router.Handle(conn); err != nil {
t.Fatal(err)
}
conn.AwaitClose(t, time.Second)
if gotMethod != "GET" {
t.Errorf("want method %q, got %q", "GET", gotMethod)
}
if gotURI != "/index.html" {
t.Errorf("want URI %q, got %q", "/index.html", gotURI)
}
if gotHost != "tinygo.org" {
t.Errorf("want Host %q, got %q", "tinygo.org", gotHost)
}
}
// Router must route on method and URI, and must not invoke a handler for
// requests it has no registration for.
func TestRouterMux(t *testing.T) {
const bufferSize = 1024
for _, test := range []struct {
name string
request string
want string // Response body the matched handler must have written.
wantNoHandler bool // No registration matches: the router must answer 404 itself.
}{
{name: "get root", request: "GET / HTTP/1.1\r\nHost: h\r\n\r\n", want: "root"},
{name: "get page", request: "GET /page HTTP/1.1\r\nHost: h\r\n\r\n", want: "page"},
{name: "any method", request: "DELETE /any HTTP/1.1\r\nHost: h\r\n\r\n", want: "any"},
{name: "method mismatch", request: "POST / HTTP/1.1\r\nHost: h\r\n\r\n", wantNoHandler: true},
{name: "unknown uri", request: "GET /nowhere HTTP/1.1\r\nHost: h\r\n\r\n", wantNoHandler: true},
} {
t.Run(test.name, func(t *testing.T) {
var (
sm MuxSlice
router Router
)
sm.Handle("GET /", staticPage(t, "root"))
sm.Handle("GET /page", staticPage(t, "page"))
sm.Handle("/any", staticPage(t, "any")) // No method: matches any.
configSynchronousRouter(t, &router, bufferSize, &sm)
conn := newConn(test.request)
if err := router.Handle(conn); err != nil {
t.Fatal(err)
}
conn.AwaitClose(t, time.Second)
got := conn.ViewWritten()
_, body, found := strings.Cut(got, "\r\n\r\n")
if !found {
t.Fatalf("header block never terminated: %q", got)
}
if test.wantNoHandler {
if !strings.HasPrefix(got, "HTTP/1.1 404 ") {
t.Errorf("want a 404 answer, got %q", got)
}
if body != "" {
t.Errorf("no handler must run, got body %q", body)
}
return
}
if !strings.HasPrefix(got, "HTTP/1.1 200 OK\r\n") {
t.Errorf("want a 200 answer, got %q", got)
}
if body != test.want {
t.Errorf("want body %q, got %q", test.want, body)
}
})
}
}
// A request arriving in pieces (TCP segmentation) must still be handled.
func TestRouterSplitRequest(t *testing.T) {
const bufferSize = 1024
const expectResponse = "split ok"
var (
sm MuxSlice
router Router
)
sm.Handle("GET /", staticPage(t, expectResponse))
configSynchronousRouter(t, &router, bufferSize, &sm)
conn := newConn("GET / HTTP/1.1\r\nHo")
conn.AddSegment("st: tinygo.org\r\n\r")
conn.AddSegment("\n") // Final CRLF lands in its own segment.
if err := router.Handle(conn); err != nil {
t.Fatal(err)
}
conn.AwaitClose(t, time.Second)
if got := conn.ViewWritten(); !strings.HasSuffix(got, expectResponse) {
t.Errorf("want body %q, got response %q", expectResponse, got)
}
}
func staticPage(t *testing.T, page string) HandlerFunc {
return func(ex *Exchange) {
var rw ExchangeRW // Streaming APIs take the ReadWriter view.
ex.ReadWriter(&rw)
n, err := io.WriteString(&rw, page)
// Handler runs on the router goroutine: Error, never Fatal.
if err != nil {
t.Error(err)
} else if n != len(page) {
t.Error("expected written ", len(page), "got", n)
}
}
}
// An exchange freed by the outgoing generation carries that generation's
// buffers. Recycling it under a new configuration serves the request with
// buffer limits the new [RouterConfig] never asked for.
func TestRouterReconfigureDropsStaleExchanges(t *testing.T) {
const smallBuf, largeBuf = 256, 1024
var (
sm MuxSlice
router Router
)
bufsize := make(chan int, 2)
sm.Handle("GET /", func(ex *Exchange) { bufsize <- len(ex.UnsafeRawBuffer()) })
serve := func(want int) {
t.Helper()
conn := newConn("GET / HTTP/1.1\r\nHost: h\r\n\r\n")
conn.Hangup()
if err := router.Handle(conn); err != nil {
t.Fatal(err)
}
if got := <-bufsize; got != want {
t.Errorf("want exchange buffer %d, got %d", want, got)
}
conn.AwaitClose(t, time.Second) // Exchange hits the freelist on close.
}
configSynchronousRouter(t, &router, smallBuf, &sm)
serve(2 * smallBuf)
configSynchronousRouter(t, &router, largeBuf, &sm)
serve(2 * largeBuf)
}
// Configure writes the fields Handle reads; concurrent use must not race.
func TestRouterConfigureHandleRace(t *testing.T) {
const bufferSize = 1024
var (
sm MuxSlice
router Router
)
sm.Handle("GET /", staticPage(t, "ok"))
configSynchronousRouter(t, &router, bufferSize, &sm)
var wg sync.WaitGroup
wg.Add(2)
go func() {
defer wg.Done()
configSynchronousRouter(t, &router, bufferSize, &sm)
}()
go func() {
defer wg.Done()
conn := newConn("GET / HTTP/1.1\r\nHost: h\r\n\r\n")
if err := router.Handle(conn); err != nil {
t.Error(err)
}
}()
wg.Wait()
}
// Reconfiguring a running router tears down the job queue that Handle may be
// sending a connection on. Connections may be dropped, but never panic.
// A torn down router has nothing left to serve with: it must say so instead of
// dropping the connection as if it were merely busy.
func TestRouterHandleAfterTeardown(t *testing.T) {
var (
sm MuxSlice
router Router
)
sm.Handle("GET /", staticPage(t, "ok"))
err := router.Configure(RouterConfig{
FixedNumGoroutines: 2,
MaxAwaitingConns: 4,
Mux: &sm,
RequestHeaderBufferSize: 512,
RequestNumHeaderKVCap: 16,
ResponseHeaderMinBufferSize: 512,
})
if err != nil {
t.Fatal(err)
}
router.Shutdown()
conn := newConn("GET / HTTP/1.1\r\nHost: h\r\n\r\n")
if err = router.Handle(conn); err != errRouterTornDown {
t.Errorf("want errRouterTornDown, got %v", err)
}
if conn.IsClosed() {
t.Error("refused connection must be left for the caller to dispose of")
}
}
// Tearing down a generation abandons the connections queued for it. They were
// taken ownership of by Handle, so they must be closed and their exchanges
// released: an exchange left claimed by a torn down generation is a buffer the
// router can never reconfigure again.
func TestRouterTeardownReleasesQueuedConns(t *testing.T) {
var (
sm MuxSlice
router Router
)
const numGoro = 2
sm.Handle("GET /", staticPage(t, "ok"))
cfg := RouterConfig{
FixedNumGoroutines: numGoro,
MaxAwaitingConns: 4,
Mux: &sm,
RequestHeaderBufferSize: 512,
RequestNumHeaderKVCap: 16,
ResponseHeaderMinBufferSize: 512,
}
// Handing connections over and tearing down immediately leaves them queued
// for workers that will never serve them. Rounds bound the scheduling luck
// needed; a single leaked exchange also fails every later Configure.
var conns [numGoro]*rwconn
for range 30 {
// errBusyExchanges is legitimate backpressure while the previous
// generation drops its connections, but it must not outlive it.
var err error
for range 100 {
if err = router.Configure(cfg); err == nil {
break
}
time.Sleep(time.Millisecond)
}
if err != nil {
t.Fatal("reconfigure after teardown:", err)
}
for i := range conns {
conns[i] = newConn("GET / HTTP/1.1\r\nHost: h\r\n\r\n")
conns[i].Hangup()
if err := router.Handle(conns[i]); err != nil {
t.Fatal(err)
}
}
router.Shutdown()
for i := range conns {
// Handle took ownership of the connection: served or dropped, the
// router closes it.
conns[i].AwaitClose(t, time.Second)
}
}
}
func TestRouterConfigureDuringWorkerHandle(t *testing.T) {
var (
sm MuxSlice
router Router
)
sm.Handle("GET /", staticPage(t, "ok"))
cfg := RouterConfig{
FixedNumGoroutines: 2,
MaxAwaitingConns: 4,
Mux: &sm,
RequestHeaderBufferSize: 512,
RequestNumHeaderKVCap: 16,
ResponseHeaderMinBufferSize: 512,
}
if err := router.Configure(cfg); err != nil {
t.Fatal(err)
}
defer router.Shutdown()
var wg sync.WaitGroup
wg.Add(2)
go func() {
defer wg.Done()
for range 300 {
conn := newConn("GET / HTTP/1.1\r\nHost: h\r\n\r\n")
conn.Hangup()
router.Handle(conn) // Drops are fine, panics are not.
}
}()
go func() {
defer wg.Done()
// Each Configure sleeps 5ms tearing down the previous generation, keep
// the count low and let the Handle loop supply the concurrency.
for range 20 {
// errBusyExchanges is legitimate backpressure: the previous
// generation was still serving when the buffers were needed.
if err := router.Configure(cfg); err != nil && err != errBusyExchanges {
t.Error(err)
return
}
}
}()
wg.Wait()
}
+273
View File
@@ -0,0 +1,273 @@
package httphi
// StatusText returns a text for the HTTP status code. It returns the empty
// string if the code is unknown.
func StatusText(code int) string {
switch status(code) {
case StatusContinue:
return "Continue"
case StatusSwitchingProtocols:
return "Switching Protocols"
case StatusProcessing:
return "Processing"
case StatusEarlyHints:
return "Early Hints"
case StatusOK:
return "OK"
case StatusCreated:
return "Created"
case StatusAccepted:
return "Accepted"
case StatusNonAuthoritativeInfo:
return "Non-Authoritative Information"
case StatusNoContent:
return "No Content"
case StatusResetContent:
return "Reset Content"
case StatusPartialContent:
return "Partial Content"
case StatusMultiStatus:
return "Multi-Status"
case StatusAlreadyReported:
return "Already Reported"
case StatusIMUsed:
return "IM Used"
case StatusMultipleChoices:
return "Multiple Choices"
case StatusMovedPermanently:
return "Moved Permanently"
case StatusFound:
return "Found"
case StatusSeeOther:
return "See Other"
case StatusNotModified:
return "Not Modified"
case StatusUseProxy:
return "Use Proxy"
case StatusTemporaryRedirect:
return "Temporary Redirect"
case StatusPermanentRedirect:
return "Permanent Redirect"
case StatusBadRequest:
return "Bad Request"
case StatusUnauthorized:
return "Unauthorized"
case StatusPaymentRequired:
return "Payment Required"
case StatusForbidden:
return "Forbidden"
case StatusNotFound:
return "Not Found"
case StatusMethodNotAllowed:
return "Method Not Allowed"
case StatusNotAcceptable:
return "Not Acceptable"
case StatusProxyAuthRequired:
return "Proxy Authentication Required"
case StatusRequestTimeout:
return "Request Timeout"
case StatusConflict:
return "Conflict"
case StatusGone:
return "Gone"
case StatusLengthRequired:
return "Length Required"
case StatusPreconditionFailed:
return "Precondition Failed"
case StatusRequestEntityTooLarge:
return "Request Entity Too Large"
case StatusRequestURITooLong:
return "Request URI Too Long"
case StatusUnsupportedMediaType:
return "Unsupported Media Type"
case StatusRequestedRangeNotSatisfiable:
return "Requested Range Not Satisfiable"
case StatusExpectationFailed:
return "Expectation Failed"
case StatusTeapot:
return "I'm a teapot"
case StatusMisdirectedRequest:
return "Misdirected Request"
case StatusUnprocessableEntity:
return "Unprocessable Entity"
case StatusLocked:
return "Locked"
case StatusFailedDependency:
return "Failed Dependency"
case StatusTooEarly:
return "Too Early"
case StatusUpgradeRequired:
return "Upgrade Required"
case StatusPreconditionRequired:
return "Precondition Required"
case StatusTooManyRequests:
return "Too Many Requests"
case StatusRequestHeaderFieldsTooLarge:
return "Request Header Fields Too Large"
case StatusUnavailableForLegalReasons:
return "Unavailable For Legal Reasons"
case StatusInternalServerError:
return "Internal Server Error"
case StatusNotImplemented:
return "Not Implemented"
case StatusBadGateway:
return "Bad Gateway"
case StatusServiceUnavailable:
return "Service Unavailable"
case StatusGatewayTimeout:
return "Gateway Timeout"
case StatusHTTPVersionNotSupported:
return "HTTP Version Not Supported"
case StatusVariantAlsoNegotiates:
return "Variant Also Negotiates"
case StatusInsufficientStorage:
return "Insufficient Storage"
case StatusLoopDetected:
return "Loop Detected"
case StatusNotExtended:
return "Not Extended"
case StatusNetworkAuthenticationRequired:
return "Network Authentication Required"
default:
return ""
}
}
const ()
type status int
// HTTP status codes as registered with IANA.
// See: https://www.iana.org/assignments/http-status-codes/http-status-codes.xhtml
const (
// RFC 9110, 15.2.1
StatusContinue = 100 // Continue
// RFC 9110, 15.2.2
StatusSwitchingProtocols = 101 // Switching Protocols
// RFC 2518, 10.1
StatusProcessing = 102 // Processing
// RFC 8297
StatusEarlyHints = 103 // Early Hints
// RFC 9110, 15.3.1
StatusOK = 200 // OK
// RFC 9110, 15.3.2
StatusCreated = 201 // Created
// RFC 9110, 15.3.3
StatusAccepted = 202 // Accepted
// RFC 9110, 15.3.4
StatusNonAuthoritativeInfo = 203 // Non-Authoritative Information
// RFC 9110, 15.3.5
StatusNoContent = 204 // No Content
// RFC 9110, 15.3.6
StatusResetContent = 205 // Reset Content
// RFC 9110, 15.3.7
StatusPartialContent = 206 // Partial Content
// RFC 4918, 11.1
StatusMultiStatus = 207 // Multi-Status
// RFC 5842, 7.1
StatusAlreadyReported = 208 // Already Reported
// RFC 3229, 10.4.1
StatusIMUsed = 226 // IM Used
// RFC 9110, 15.4.1
StatusMultipleChoices = 300 // Multiple Choices
// RFC 9110, 15.4.2
StatusMovedPermanently = 301 // Moved Permanently
// RFC 9110, 15.4.3
StatusFound = 302 // Found
// RFC 9110, 15.4.4
StatusSeeOther = 303 // See Other
// RFC 9110, 15.4.5
StatusNotModified = 304 // Not Modified
// RFC 9110, 15.4.6
StatusUseProxy = 305 // Use Proxy
// RFC 9110, 15.4.7 (Unused)
_ = 306
// RFC 9110, 15.4.8
StatusTemporaryRedirect = 307 // Temporary Redirect
// RFC 9110, 15.4.9
StatusPermanentRedirect = 308 // Permanent Redirect
// RFC 9110, 15.5.1
StatusBadRequest = 400 // Bad Request
// RFC 9110, 15.5.2
StatusUnauthorized = 401 // Unauthorized
// RFC 9110, 15.5.3
StatusPaymentRequired = 402 // Payment Required
// RFC 9110, 15.5.4
StatusForbidden = 403 // Forbidden
// RFC 9110, 15.5.5
StatusNotFound = 404 // Not Found
// RFC 9110, 15.5.6
StatusMethodNotAllowed = 405 // Method Not Allowed
// RFC 9110, 15.5.7
StatusNotAcceptable = 406 // Not Acceptable
// RFC 9110, 15.5.8
StatusProxyAuthRequired = 407 // Proxy Authentication Required
// RFC 9110, 15.5.9
StatusRequestTimeout = 408 // Request Timeout
// RFC 9110, 15.5.10
StatusConflict = 409 // Conflict
// RFC 9110, 15.5.11
StatusGone = 410 // Gone
// RFC 9110, 15.5.12
StatusLengthRequired = 411 // Length Required
// RFC 9110, 15.5.13
StatusPreconditionFailed = 412 // Precondition Failed
// RFC 9110, 15.5.14
StatusRequestEntityTooLarge = 413 // Request Entity Too Large
// RFC 9110, 15.5.15
StatusRequestURITooLong = 414 // Request URI Too Long
// RFC 9110, 15.5.16
StatusUnsupportedMediaType = 415 // Unsupported Media Type
// RFC 9110, 15.5.17
StatusRequestedRangeNotSatisfiable = 416 // Requested Range Not Satisfiable
// RFC 9110, 15.5.18
StatusExpectationFailed = 417 // Expectation Failed
// RFC 9110, 15.5.19 (Unused)
StatusTeapot = 418 // I'm a teapot
// RFC 9110, 15.5.20
StatusMisdirectedRequest = 421 // Misdirected Request
// RFC 9110, 15.5.21
StatusUnprocessableEntity = 422 // Unprocessable Entity
// RFC 4918, 11.3
StatusLocked = 423 // Locked
// RFC 4918, 11.4
StatusFailedDependency = 424 // Failed Dependency
// RFC 8470, 5.2.
StatusTooEarly = 425 // Too Early
// RFC 9110, 15.5.22
StatusUpgradeRequired = 426 // Upgrade Required
// RFC 6585, 3
StatusPreconditionRequired = 428 // Precondition Required
// RFC 6585, 4
StatusTooManyRequests = 429 // Too Many Requests
// RFC 6585, 5
StatusRequestHeaderFieldsTooLarge = 431 // Request Header Fields Too Large
// RFC 7725, 3
StatusUnavailableForLegalReasons = 451 // Unavailable For Legal Reasons
// RFC 9110, 15.6.1
StatusInternalServerError = 500 // Internal Server Error
// RFC 9110, 15.6.2
StatusNotImplemented = 501 // Not Implemented
// RFC 9110, 15.6.3
StatusBadGateway = 502 // Bad Gateway
// RFC 9110, 15.6.4
StatusServiceUnavailable = 503 // Service Unavailable
// RFC 9110, 15.6.5
StatusGatewayTimeout = 504 // Gateway Timeout
// RFC 9110, 15.6.6
StatusHTTPVersionNotSupported = 505 // HTTP Version Not Supported
// RFC 2295, 8.1
StatusVariantAlsoNegotiates = 506 // Variant Also Negotiates
// RFC 4918, 11.5
StatusInsufficientStorage = 507 // Insufficient Storage
// RFC 5842, 7.2
StatusLoopDetected = 508 // Loop Detected
// RFC 2774, 7
StatusNotExtended = 510 // Not Extended
// RFC 6585, 6
StatusNetworkAuthenticationRequired = 511 // Network Authentication Required
)
+34
View File
@@ -0,0 +1,34 @@
// Code generated by "stringer -type Method -linecomment -output stringers.go"; DO NOT EDIT.
package httphi
import "strconv"
func _() {
// An "invalid array index" compiler error signifies that the constant values have changed.
// Re-run the stringer command to generate them again.
var x [1]struct{}
_ = x[MethUndefined-0]
_ = x[MethGet-1]
_ = x[MethHead-2]
_ = x[MethPost-3]
_ = x[MethPut-4]
_ = x[MethPatch-5]
_ = x[MethDelete-6]
_ = x[MethConnect-7]
_ = x[MethOptions-8]
_ = x[MethTrace-9]
_ = x[MethUnknown-10]
}
const _Method_name = "undefinedGETHEADPOSTPUTPATCHDELETECONNECTOPTIONSTRACEunknown"
var _Method_index = [...]uint8{0, 9, 12, 16, 20, 23, 28, 34, 41, 48, 53, 60}
func (i Method) String() string {
idx := int(i) - 0
if i < 0 || idx >= len(_Method_index)-1 {
return "Method(" + strconv.FormatInt(int64(i), 10) + ")"
}
return _Method_name[_Method_index[idx]:_Method_index[idx+1]]
}