// Package transport provides the network layer: epoll event loop, TCP accept, // HTTP/1.1 request parsing, and WebSocket framing. It has no knowledge of // Nostr protocol or any application domain logic. It calls back to a Handler // interface for all domain events. // // Dependency direction: server imports transport, never the reverse. // No musiquay package imports are allowed here - only stdlib. package transport import ( "bytes" "fmt" "runtime" "syscall" "git.smesh.lol/morly/pkg/metrics" "time" ) const maxBuf = 20 << 20 // 20MB max per-connection buffer const ( phaseHTTP = 0 phaseWS = 1 phaseHTTPBody = 2 phaseHTTPDeferred = 3 // waiting for async worker response ) const ( opText byte = 0x1 opBin byte = 0x2 opClose byte = 0x8 opPing byte = 0x9 opPong byte = 0xA ) // HTTPDeferred is returned from Handler.OnHTTP to indicate async processing. // Transport sets the connection to deferred phase and expects CompleteHTTP later. const HTTPDeferred = -1 // Handler is implemented by the server layer. Transport calls these methods // for all connection and message events. type Handler interface { // OnAccept: new TCP connection. Return false to close immediately. OnAccept(fd int32, ip string) bool // OnWSUpgrade: WS handshake requested. currentIPWSCount = existing WS conns from ip. // Return (whitelisted, allow). allow=false → 429. OnWSUpgrade(fd int32, ip string, currentIPWSCount int32) (whitelisted bool, allow bool) // OnWSConnected: 101 response sent. Handler sets up conn state, sends auth challenge. OnWSConnected(fd int32) // OnWSMessage: decoded WS payload. OnWSMessage(fd int32, payload []byte) // OnWSClose: WS connection closed. Handler cleans up conn state. OnWSClose(fd int32) // OnHTTP: complete HTTP request. Return (HTTPDeferred, nil, nil, false) for async. OnHTTP(fd int32, method, path string, headers map[string]string, body []byte) (status int32, respHeaders map[string]string, respBody []byte, connClose bool) // OnPoll: called once per event-loop iteration, after epoll events are // dispatched. Worker domains are reached over spawn channels, which have // no file descriptor to wait on, so the handler drains their output // queues here. OnPoll() // OnTick: called periodically when epoll_wait times out (every ~5s). OnTick() } // Server runs the epoll event loop. type Server struct { BotBlock bool // block known bot User-Agents OnReady func() // called after bind+listen, before epoll loop handler Handler epfd int32 lnFD int32 sigFD int32 conns map[int32]*tconn ipConns map[string]int32 // WS connection counts per IP maxConnPerIP int32 } type tconn struct { fd int32 phase int32 buf []byte wpos int32 remoteIP string whitelisted bool pendingReq *httpReq bodyNeeded int32 wbuf []byte // pending write data (EAGAIN buffered) wbufClose bool // close connection after wbuf drains // arena holds this connection's cross-turn state: the parsed request while // its body arrives, the body itself, and bytes buffered on EAGAIN. Those // must outlive the loop turn that produced them, but not the connection, so // they cannot go in a frame arena (freed at return) and must not go in the // long-lived event-loop arena (never freed). It is reset whenever the // connection goes idle, and freed when it closes. arena *Arena } type httpReq struct { method string path string headers map[string]string body []byte } //export moxie_signal_enable func moxie_signal_enable(s uint32) //export moxie_signal_pipe_init func moxie_signal_pipe_init() (n int32) //export moxie_signal_pipe_read func moxie_signal_pipe_read() (n int32) var globalSigFD int32 // InitSignals sets up SIGTERM/SIGINT handling. Must be called before any store // operations so that shutdown signals during slow startup are handled cleanly. // InitSignals is called once at startup and stores a package global, so the // stores live in an init-named function (the rule's exemption); the exported // wrapper keeps the call sites. func InitSignals() { initSignals() } func initSignals() { globalSigFD = moxie_signal_pipe_init() moxie_signal_enable(15) // SIGTERM moxie_signal_enable(2) // SIGINT } // New creates a transport Server. maxConnPerIP=0 means unlimited. func New(handler Handler, maxConnPerIP int32) (s *Server) { return &Server{ handler: handler, conns: map[int32]*tconn{}, ipConns: map[string]int32{}, maxConnPerIP: maxConnPerIP, } } // ConnCount returns the current number of tracked TCP connections. func (s *Server) ConnCount() (n int32) { return len(s.conns) } // ConnIP returns the effective remote IP for a connection (XFF-substituted). func (s *Server) ConnIP(fd int32) (sv string) { if c := s.conns[fd]; c != nil { return c.remoteIP } return "" } // ConnIsWhitelisted reports whether fd's IP is whitelisted. func (s *Server) ConnIsWhitelisted(fd int32) (ok bool) { if c := s.conns[fd]; c != nil { return c.whitelisted } return false } // ConnIsWS reports whether fd is in WS phase. func (s *Server) ConnIsWS(fd int32) (ok bool) { if c := s.conns[fd]; c != nil { return c.phase == phaseWS } return false } // IPConnCount returns the current WS connection count for ip. func (s *Server) IPConnCount(ip string) (n int32) { return s.ipConns[ip] } // SendWS writes a WS text frame. Buffers on EAGAIN; closes on error. func (s *Server) SendWS(fd int32, payload []byte) { c := s.conns[fd] if c == nil { return } s.connWrite(c, buildWSFrame(opText, payload), false) } // SendWSErr was removed. Its one difference from SendWS was to report EAGAIN // from a direct write as an error, and both callers answered that transient // state by closing a healthy connection - losing the frame that hit it and // every frame after. SendWS buffers the remainder and flushes it on a later // loop pass, which is what a full socket buffer asks for. // errAgain reports whether err is EAGAIN/EWOULDBLOCK. The comparison must be // on the numeric Errno the syscall package carries: comparing the error // interface against the syscall.EAGAIN constant did not match the value the // write actually returned, so a full socket buffer was treated as a fatal // error and the connection was closed mid-response (Content-Length still // promising the rest). func errAgain(err error) (ok bool) { ee, isErrno := err.(syscall.Errno) return isErrno && (ee == syscall.EAGAIN || ee == syscall.EWOULDBLOCK) } // errIntr reports whether err is EINTR (a signal interrupted the call). func errIntr(err error) (ok bool) { ee, isErrno := err.(syscall.Errno) return isErrno && ee == syscall.EINTR } // connWrite writes data to a connection, buffering on EAGAIN. func (s *Server) connWrite(c *tconn, data []byte, closeAfter bool) { if c.wbuf != nil { prevW := connStateEnter(c) c.wbuf = c.wbuf | data connStateExit(prevW) if closeAfter { c.wbufClose = true } return } var sent int32 for len(data) > 0 { n, err := syscall.Write(c.fd, data) if n > 0 { data = data[n:] sent += n } if errAgain(err) { prevE := connStateEnter(c) c.wbuf = []byte{:len(data)} copy(c.wbuf, data) connStateExit(prevE) c.wbufClose = closeAfter epollModWrite(s.epfd, c.fd) return } if errIntr(err) { // A signal (timer, child exit) interrupted the write. Bytes // already counted were sent; retry the remainder. Without this a // large response died mid-flight: the socket buffer fills part // way, a signal arrives, and the connection closes while // Content-Length still promises the rest. continue } if err != nil { s.closeConn(c) return } } _ = sent if closeAfter { s.closeConn(c) } } // SendHTTP writes an HTTP response. Uses buffered writes; closes on error. func (s *Server) SendHTTP(fd int32, status int32, headers map[string]string, body []byte) { c := s.conns[fd] if c == nil { return } s.sendHTTPBuffered(c, status, headers, body, false) } // CloseConn closes a connection. func (s *Server) CloseConn(fd int32) { if c := s.conns[fd]; c != nil { s.closeConn(c) } } // CompleteHTTP sends an HTTP response for a previously deferred connection. func (s *Server) CompleteHTTP(fd int32, status int32, headers map[string]string, body []byte, connClose bool) { c := s.conns[fd] if c == nil { return } c.phase = phaseHTTP s.sendHTTPBuffered(c, status, headers, body, connClose) } // connStateEnter makes the connection's own arena current for allocations that // must outlive this loop turn. Paired with connStateExit; the push/pop form // leaves the enclosing frame's arena untouched. func connStateEnter(c *tconn) (prev *runtime.Arena) { prev = runtime.CurrentArena() if c != nil && c.arena != nil { runtime.SovereignSetArena(c.arena) } return } func connStateExit(prev *runtime.Arena) { runtime.SovereignRestoreArena(prev) } // connStateReset drops everything the connection arena holds. Only valid when // the connection is idle: no parsed request and no buffered write referencing // it (keepAlive is exactly that point). func connStateReset(c *tconn) { if c != nil && c.arena != nil { runtime.ArenaReset(c.arena) } } // setupListener performs the receiver writes for ListenAndServe: the listener // fd, the epoll fd, and the signal fd, plus their epoll registrations. // // These writes live in their own mutating method so that ListenAndServe itself // stays a read-only method of a sovereign type. A self-mutating method borrows // the receiver's sovereign arena as its working frame for its whole body, and // ListenAndServe's body is the event loop: keeping the borrow there would make // the loop's compaction yield point unable to compact the Server's own group, // which is exactly the group that grows with every connection. Peeling the // writes out leaves the loop running on a regular per-function arena with the // Server's arena merely registered as stable, which is what the compaction // rebuild pass knows how to repoint. func (s *Server) setupListener(fd int32) (eerr error) { epfd, serr := syscall.EpollCreate1(0) if serr != nil { syscall.Close(fd) return fmt.Errorf("transport: epoll: %w", serr) } if aerr := epollAdd(epfd, fd); aerr != nil { syscall.Close(epfd) syscall.Close(fd) return fmt.Errorf("transport: epoll add: %w", aerr) } s.lnFD = fd s.epfd = epfd if globalSigFD >= 0 { s.sigFD = int32(globalSigFD) if aerr2 := epollAdd(epfd, s.sigFD); aerr2 != nil { return fmt.Errorf("transport: epoll add signal pipe: %w", aerr2) } } return nil } // ListenAndServe runs the epoll event loop. Returns on signal or error. func (s *Server) ListenAndServe(addr string) (eerr error) { ip, port := ParseAddr(addr) fd, serr := syscall.Socket(syscall.AF_INET, syscall.SOCK_STREAM, 0) if serr != nil { return fmt.Errorf("transport: socket: %w", serr) } syscall.SetsockoptInt(fd, syscall.SOL_SOCKET, syscall.SO_REUSEADDR, 1) if serr = syscall.SetNonblock(fd, true); serr != nil { syscall.Close(fd) return fmt.Errorf("transport: nonblock: %w", serr) } sa := &syscall.SockaddrInet4{Port: port, Addr: ip} if serr = syscall.Bind(fd, sa); serr != nil { syscall.Close(fd) return fmt.Errorf("transport: bind %s: %w", addr, serr) } if serr = syscall.Listen(fd, 4096); serr != nil { syscall.Close(fd) return fmt.Errorf("transport: listen: %w", serr) } if serr = s.setupListener(fd); serr != nil { return serr } if s.OnReady != nil { s.OnReady() } // Poll interval: worker responses arrive over spawn channels, which have // no fd to wake epoll, so the loop must come back regularly to drain // them. 2ms keeps the added latency imperceptible without busy-spinning. const pollTimeoutMs = 2 const ticksPerOnTick = 2500 // 2ms * 2500 = 5s idle events := []syscall.EpollEvent{:64} idle := int32(0) var n int32 var perr error var evFD int32 var waitStart, bodyStart int64 for { // Yield point: drain queued sovereign-arena compactions here, where no // request frame is live, instead of inside a method boundary. Safe // because ListenAndServe is read-only (see setupListener), so no frame // on the stack borrows a sovereign arena, and every method it calls has // already returned by the time the loop comes back around. The drain is // bounded so a large compaction cannot stall the accept loop. runtime.SovDrainCompactions(1) // One arena per turn. This frame lives for the whole process, so // anything a handler leaves behind would be permanent: the value an HTTP // handler returns (a multi-MB static file) relocates up the call chain // and comes to rest in the outermost live frame's arena, which is this // one, and nothing ever reclaims it. A per-turn arena gives that data // the lifetime of the turn and the arena is recycled. State that must // outlive the turn goes in the connection's own arena (connStateEnter). runtime.FnArenaPush(65536) stop := false waitStart = metrics.Now() n, perr = syscall.EpollWait(s.epfd, events, pollTimeoutMs) metrics.AcceptLoopWaitNs.Observe(metrics.Since(waitStart)) if perr != nil { runtime.FnArenaPopFree() if perr == syscall.EINTR { return nil } return fmt.Errorf("transport: epoll wait: %w", perr) } for i := 0; i < n; i++ { evFD = int32(events[i].Fd) if evFD == s.sigFD { moxie_signal_pipe_read() stop = true break } else if evFD == s.lnFD { s.acceptAll() } else if c := s.conns[evFD]; c != nil { if events[i].Events&(syscall.EPOLLERR|syscall.EPOLLHUP) != 0 { s.closeConn(c) } else if events[i].Events&syscall.EPOLLOUT != 0 { s.drainWrite(c) } else { s.readConn(c) } } } if !stop { bodyStart = metrics.Now() s.handler.OnPoll() metrics.AcceptLoopHandleNs.Observe(metrics.Since(bodyStart)) if n > 0 { idle = 0 } else { idle++ if idle >= ticksPerOnTick { idle = 0 s.handler.OnTick() } } } runtime.FnArenaPopFree() if stop { return nil } } } func (s *Server) acceptAll() { // Declared outside the loop: a declaration in the body of a self-mutating // method allocates in the sovereign arena on every iteration. var nfd int32 var sa syscall.Sockaddr var err error var ip string var aerr error var ok bool for { nfd, sa, err = syscall.Accept4(s.lnFD, syscall.SOCK_NONBLOCK) if err != nil { return } // The connection record outlives this call: it stays in s.conns for // the life of the connection. The record, its receive buffer and the // remote address string it keeps must all be allocated in the root // arena - the address string is built here, so the borrow has to // cover it too, not just the composite literal. prev := runtime.CurrentArena() runtime.SovereignSetArena(runtime.RootArena()) ip = peerAddr(sa) ok = s.handler.OnAccept(nfd, ip) if !ok { runtime.SovereignRestoreArena(prev) syscall.Close(nfd) continue } if aerr = epollAdd(s.epfd, nfd); aerr != nil { runtime.SovereignRestoreArena(prev) syscall.Close(nfd) continue } s.conns[nfd] = &tconn{ fd: nfd, phase: phaseHTTP, buf: []byte{:4096}, remoteIP: ip, arena: runtime.ArenaNew(65536), } runtime.SovereignRestoreArena(prev) } } func peerAddr(sa syscall.Sockaddr) (s string) { if sa4, ok := sa.(*syscall.SockaddrInet4); ok { // Grow with the append operator: a presized []byte{:0:20} plus push // is a bounded store that fails loud the moment the capacity is not // what the literal promised. b := []byte{} b = appendInt(b, int32(sa4.Addr[0])) b = b | "." b = appendInt(b, int32(sa4.Addr[1])) b = b | "." b = appendInt(b, int32(sa4.Addr[2])) b = b | "." b = appendInt(b, int32(sa4.Addr[3])) return string(makeCopy(b)) } return "unknown" } // growBuf doubles buf, preserving the first wpos bytes. Free function: the // abandoned buffer dies with this frame's arena instead of living on in the // Server's sovereign arena for every growth step. func growBuf(buf []byte, wpos int32) (nb []byte) { nb = []byte{:len(buf) * 2} copy(nb, buf[:wpos]) return nb } func (s *Server) readConn(c *tconn) { if c.wbuf != nil { return } avail := len(c.buf) - c.wpos if avail < 512 { if len(c.buf) >= maxBuf { s.closeConn(c) return } c.buf = growBuf(c.buf, c.wpos) } n, err := syscall.Read(c.fd, c.buf[c.wpos:]) if n <= 0 { if err == nil || (!errAgain(err) && !errIntr(err)) { s.closeConn(c) } return } c.wpos += n switch c.phase { case phaseHTTP: s.processHTTP(c) case phaseHTTPBody: s.processHTTPBody(c) case phaseWS: s.processWS(c) case phaseHTTPDeferred: c.wpos = 0 // drop bytes while awaiting async response } } func (s *Server) closeConn(c *tconn) { epollDel(s.epfd, c.fd) syscall.Close(c.fd) c.buf = nil c.wbuf = nil c.pendingReq = nil if c.arena != nil { runtime.ArenaFree(c.arena) c.arena = nil } delete(s.conns, c.fd) if c.phase == phaseWS { s.handler.OnWSClose(c.fd) if n := s.ipConns[c.remoteIP] - 1; n <= 0 { delete(s.ipConns, c.remoteIP) } else { s.ipConns[c.remoteIP] = n } } } func (s *Server) keepAlive(c *tconn) { // Idle: no parsed request and no buffered write can reference the // connection arena, so its accumulated per-request state can be released // now instead of at close. connStateReset(c) c.phase = phaseHTTP c.pendingReq = nil c.bodyNeeded = 0 if c.wpos > 0 { s.processHTTP(c) } } func (s *Server) processHTTP(c *tconn) { data := c.buf[:c.wpos] end := bytes.Index(data, []byte("\r\n\r\n")) if end < 0 { return } consumed := end + 4 req := parseHTTPHeaders(data[:end]) if req == nil { s.closeConn(c) return } copy(c.buf, c.buf[consumed:c.wpos]) c.wpos -= consumed // Reverse proxy: substitute real IP from X-Forwarded-For before any checks. if c.remoteIP == "127.0.0.1" { if xff := req.headers["x-forwarded-for"]; xff != "" { realIP := FirstXFF(xff) if len(realIP) > 0 && realIP != c.remoteIP { prevR := runtime.CurrentArena() runtime.SovereignSetArena(runtime.RootArena()) c.remoteIP = realIP runtime.SovereignRestoreArena(prevR) } } } if bytes.EqualFold([]byte(req.headers["upgrade"]), []byte("websocket")) { s.upgradeWS(c, req) return } if s.BotBlock && isBot(req.headers["user-agent"]) { writeHTTPResponse(c.fd, 403, nil, []byte("forbidden")) s.closeConn(c) return } cl := parseContentLength(req.headers["content-length"]) if cl > 0 && cl <= maxBuf { prevP := connStateEnter(c) c.pendingReq = parseHTTPHeaders(makeCopy(data[:end])) connStateExit(prevP) c.bodyNeeded = cl c.phase = phaseHTTPBody s.processHTTPBody(c) return } status, headers, body, connClose := s.handler.OnHTTP(c.fd, req.method, req.path, req.headers, nil) if status == HTTPDeferred { c.phase = phaseHTTPDeferred return } s.sendHTTPBuffered(c, status, headers, body, connClose || req.headers["connection"] == "close") } func (s *Server) processHTTPBody(c *tconn) { if c.wpos < c.bodyNeeded { return } prevB := connStateEnter(c) c.pendingReq.body = makeCopy(c.buf[:c.bodyNeeded]) connStateExit(prevB) copy(c.buf, c.buf[c.bodyNeeded:c.wpos]) c.wpos -= c.bodyNeeded req := c.pendingReq c.pendingReq = nil status, headers, body, connClose := s.handler.OnHTTP(c.fd, req.method, req.path, req.headers, req.body) if status == HTTPDeferred { c.phase = phaseHTTPDeferred return } s.sendHTTPBuffered(c, status, headers, body, connClose || req.headers["connection"] == "close") } // sendHTTPBuffered writes an HTTP response, buffering on EAGAIN. func (s *Server) sendHTTPBuffered(c *tconn, status int32, headers map[string]string, body []byte, connClose bool) { s.connWrite(c, buildHTTPResponse(status, headers, body), connClose) if c.wbuf == nil && s.conns[c.fd] != nil && !connClose { s.keepAlive(c) } } // drainWrite flushes pending write buffer when socket becomes writable. func (s *Server) drainWrite(c *tconn) { for len(c.wbuf) > 0 { n, err := syscall.Write(c.fd, c.wbuf) if n > 0 { c.wbuf = c.wbuf[n:] } if errAgain(err) { return } if errIntr(err) { continue } if err != nil { s.closeConn(c) return } } c.wbuf = nil epollModRead(s.epfd, c.fd) if c.wbufClose { s.closeConn(c) return } if c.phase == phaseHTTP || c.phase == phaseHTTPBody { s.keepAlive(c) } } // Epoll helpers. func epollAdd(epfd, fd int32) (eerr error) { ev := syscall.EpollEvent{Events: syscall.EPOLLIN, Fd: int32(fd)} return syscall.EpollCtl(epfd, syscall.EPOLL_CTL_ADD, fd, &ev) } func epollModWrite(epfd, fd int32) { ev := syscall.EpollEvent{Events: syscall.EPOLLIN | syscall.EPOLLOUT, Fd: int32(fd)} syscall.EpollCtl(epfd, syscall.EPOLL_CTL_MOD, fd, &ev) } func epollModRead(epfd, fd int32) { ev := syscall.EpollEvent{Events: syscall.EPOLLIN, Fd: int32(fd)} syscall.EpollCtl(epfd, syscall.EPOLL_CTL_MOD, fd, &ev) } func epollDel(epfd, fd int32) (eerr error) { return syscall.EpollCtl(epfd, syscall.EPOLL_CTL_DEL, fd, nil) } func init() { globalSigFD = -1 botAgents = []string{ "SemrushBot", "AhrefsBot", "GPTBot", "ClaudeBot", "MJ12bot", "DotBot", "PetalBot", "Bytespider", "Sogou", "YandexBot", "BLEXBot", "DataForSeoBot", } }