transport.mx raw

   1  // Package transport provides the network layer: epoll event loop, TCP accept,
   2  // HTTP/1.1 request parsing, and WebSocket framing. It has no knowledge of
   3  // Nostr protocol or any application domain logic. It calls back to a Handler
   4  // interface for all domain events.
   5  //
   6  // Dependency direction: server imports transport, never the reverse.
   7  // No musiquay package imports are allowed here - only stdlib.
   8  package transport
   9  
  10  import (
  11  	"bytes"
  12  	"fmt"
  13  	"runtime"
  14  	"syscall"
  15  	"git.smesh.lol/morly/pkg/metrics"
  16  	"time"
  17  )
  18  
  19  
  20  const maxBuf = 20 << 20 // 20MB max per-connection buffer
  21  
  22  const (
  23  	phaseHTTP         = 0
  24  	phaseWS           = 1
  25  	phaseHTTPBody     = 2
  26  	phaseHTTPDeferred = 3 // waiting for async worker response
  27  )
  28  
  29  const (
  30  	opText  byte = 0x1
  31  	opBin   byte = 0x2
  32  	opClose byte = 0x8
  33  	opPing  byte = 0x9
  34  	opPong  byte = 0xA
  35  )
  36  
  37  // HTTPDeferred is returned from Handler.OnHTTP to indicate async processing.
  38  // Transport sets the connection to deferred phase and expects CompleteHTTP later.
  39  const HTTPDeferred = -1
  40  
  41  // Handler is implemented by the server layer. Transport calls these methods
  42  // for all connection and message events.
  43  type Handler interface {
  44  	// OnAccept: new TCP connection. Return false to close immediately.
  45  	OnAccept(fd int32, ip string) bool
  46  	// OnWSUpgrade: WS handshake requested. currentIPWSCount = existing WS conns from ip.
  47  	// Return (whitelisted, allow). allow=false → 429.
  48  	OnWSUpgrade(fd int32, ip string, currentIPWSCount int32) (whitelisted bool, allow bool)
  49  	// OnWSConnected: 101 response sent. Handler sets up conn state, sends auth challenge.
  50  	OnWSConnected(fd int32)
  51  	// OnWSMessage: decoded WS payload.
  52  	OnWSMessage(fd int32, payload []byte)
  53  	// OnWSClose: WS connection closed. Handler cleans up conn state.
  54  	OnWSClose(fd int32)
  55  	// OnHTTP: complete HTTP request. Return (HTTPDeferred, nil, nil, false) for async.
  56  	OnHTTP(fd int32, method, path string, headers map[string]string, body []byte) (status int32, respHeaders map[string]string, respBody []byte, connClose bool)
  57  	// OnPoll: called once per event-loop iteration, after epoll events are
  58  	// dispatched. Worker domains are reached over spawn channels, which have
  59  	// no file descriptor to wait on, so the handler drains their output
  60  	// queues here.
  61  	OnPoll()
  62  	// OnTick: called periodically when epoll_wait times out (every ~5s).
  63  	OnTick()
  64  }
  65  
  66  // Server runs the epoll event loop.
  67  type Server struct {
  68  	BotBlock     bool           // block known bot User-Agents
  69  	OnReady      func()         // called after bind+listen, before epoll loop
  70  	handler      Handler
  71  	epfd         int32
  72  	lnFD         int32
  73  	sigFD        int32
  74  	conns        map[int32]*tconn
  75  	ipConns      map[string]int32 // WS connection counts per IP
  76  	maxConnPerIP int32
  77  }
  78  
  79  type tconn struct {
  80  	fd          int32
  81  	phase       int32
  82  	buf         []byte
  83  	wpos        int32
  84  	remoteIP    string
  85  	whitelisted bool
  86  	pendingReq  *httpReq
  87  	bodyNeeded  int32
  88  	wbuf        []byte // pending write data (EAGAIN buffered)
  89  	wbufClose   bool   // close connection after wbuf drains
  90  	// arena holds this connection's cross-turn state: the parsed request while
  91  	// its body arrives, the body itself, and bytes buffered on EAGAIN. Those
  92  	// must outlive the loop turn that produced them, but not the connection, so
  93  	// they cannot go in a frame arena (freed at return) and must not go in the
  94  	// long-lived event-loop arena (never freed). It is reset whenever the
  95  	// connection goes idle, and freed when it closes.
  96  	arena *Arena
  97  }
  98  
  99  type httpReq struct {
 100  	method  string
 101  	path    string
 102  	headers map[string]string
 103  	body    []byte
 104  }
 105  
 106  //export moxie_signal_enable
 107  func moxie_signal_enable(s uint32)
 108  
 109  //export moxie_signal_pipe_init
 110  func moxie_signal_pipe_init() (n int32)
 111  
 112  //export moxie_signal_pipe_read
 113  func moxie_signal_pipe_read() (n int32)
 114  
 115  var globalSigFD int32
 116  
 117  // InitSignals sets up SIGTERM/SIGINT handling. Must be called before any store
 118  // operations so that shutdown signals during slow startup are handled cleanly.
 119  // InitSignals is called once at startup and stores a package global, so the
 120  // stores live in an init-named function (the rule's exemption); the exported
 121  // wrapper keeps the call sites.
 122  func InitSignals() {
 123  	initSignals()
 124  }
 125  
 126  func initSignals() {
 127  	globalSigFD = moxie_signal_pipe_init()
 128  	moxie_signal_enable(15) // SIGTERM
 129  	moxie_signal_enable(2)  // SIGINT
 130  }
 131  
 132  // New creates a transport Server. maxConnPerIP=0 means unlimited.
 133  func New(handler Handler, maxConnPerIP int32) (s *Server) {
 134  	return &Server{
 135  		handler:      handler,
 136  		conns:        map[int32]*tconn{},
 137  		ipConns:      map[string]int32{},
 138  		maxConnPerIP: maxConnPerIP,
 139  	}
 140  }
 141  
 142  // ConnCount returns the current number of tracked TCP connections.
 143  func (s *Server) ConnCount() (n int32) { return len(s.conns) }
 144  
 145  // ConnIP returns the effective remote IP for a connection (XFF-substituted).
 146  func (s *Server) ConnIP(fd int32) (sv string) {
 147  	if c := s.conns[fd]; c != nil {
 148  		return c.remoteIP
 149  	}
 150  	return ""
 151  }
 152  
 153  // ConnIsWhitelisted reports whether fd's IP is whitelisted.
 154  func (s *Server) ConnIsWhitelisted(fd int32) (ok bool) {
 155  	if c := s.conns[fd]; c != nil {
 156  		return c.whitelisted
 157  	}
 158  	return false
 159  }
 160  
 161  // ConnIsWS reports whether fd is in WS phase.
 162  func (s *Server) ConnIsWS(fd int32) (ok bool) {
 163  	if c := s.conns[fd]; c != nil {
 164  		return c.phase == phaseWS
 165  	}
 166  	return false
 167  }
 168  
 169  // IPConnCount returns the current WS connection count for ip.
 170  func (s *Server) IPConnCount(ip string) (n int32) {
 171  	return s.ipConns[ip]
 172  }
 173  
 174  // SendWS writes a WS text frame. Buffers on EAGAIN; closes on error.
 175  func (s *Server) SendWS(fd int32, payload []byte) {
 176  	c := s.conns[fd]
 177  	if c == nil {
 178  		return
 179  	}
 180  	s.connWrite(c, buildWSFrame(opText, payload), false)
 181  }
 182  
 183  // SendWSErr was removed. Its one difference from SendWS was to report EAGAIN
 184  // from a direct write as an error, and both callers answered that transient
 185  // state by closing a healthy connection - losing the frame that hit it and
 186  // every frame after. SendWS buffers the remainder and flushes it on a later
 187  // loop pass, which is what a full socket buffer asks for.
 188  
 189  
 190  // errAgain reports whether err is EAGAIN/EWOULDBLOCK. The comparison must be
 191  // on the numeric Errno the syscall package carries: comparing the error
 192  // interface against the syscall.EAGAIN constant did not match the value the
 193  // write actually returned, so a full socket buffer was treated as a fatal
 194  // error and the connection was closed mid-response (Content-Length still
 195  // promising the rest).
 196  func errAgain(err error) (ok bool) {
 197  	ee, isErrno := err.(syscall.Errno)
 198  	return isErrno && (ee == syscall.EAGAIN || ee == syscall.EWOULDBLOCK)
 199  }
 200  
 201  // errIntr reports whether err is EINTR (a signal interrupted the call).
 202  func errIntr(err error) (ok bool) {
 203  	ee, isErrno := err.(syscall.Errno)
 204  	return isErrno && ee == syscall.EINTR
 205  }
 206  
 207  // connWrite writes data to a connection, buffering on EAGAIN.
 208  func (s *Server) connWrite(c *tconn, data []byte, closeAfter bool) {
 209  	if c.wbuf != nil {
 210  		prevW := connStateEnter(c)
 211  		c.wbuf = c.wbuf | data
 212  		connStateExit(prevW)
 213  		if closeAfter {
 214  			c.wbufClose = true
 215  		}
 216  		return
 217  	}
 218  	var sent int32
 219  	for len(data) > 0 {
 220  		n, err := syscall.Write(c.fd, data)
 221  		if n > 0 {
 222  			data = data[n:]
 223  			sent += n
 224  		}
 225  		if errAgain(err) {
 226  			prevE := connStateEnter(c)
 227  			c.wbuf = []byte{:len(data)}
 228  			copy(c.wbuf, data)
 229  			connStateExit(prevE)
 230  			c.wbufClose = closeAfter
 231  			epollModWrite(s.epfd, c.fd)
 232  			return
 233  		}
 234  		if errIntr(err) {
 235  			// A signal (timer, child exit) interrupted the write. Bytes
 236  			// already counted were sent; retry the remainder. Without this a
 237  			// large response died mid-flight: the socket buffer fills part
 238  			// way, a signal arrives, and the connection closes while
 239  			// Content-Length still promises the rest.
 240  			continue
 241  		}
 242  		if err != nil {
 243  			s.closeConn(c)
 244  			return
 245  		}
 246  	}
 247  	_ = sent
 248  	if closeAfter {
 249  		s.closeConn(c)
 250  	}
 251  }
 252  
 253  // SendHTTP writes an HTTP response. Uses buffered writes; closes on error.
 254  func (s *Server) SendHTTP(fd int32, status int32, headers map[string]string, body []byte) {
 255  	c := s.conns[fd]
 256  	if c == nil {
 257  		return
 258  	}
 259  	s.sendHTTPBuffered(c, status, headers, body, false)
 260  }
 261  
 262  // CloseConn closes a connection.
 263  func (s *Server) CloseConn(fd int32) {
 264  	if c := s.conns[fd]; c != nil {
 265  		s.closeConn(c)
 266  	}
 267  }
 268  
 269  // CompleteHTTP sends an HTTP response for a previously deferred connection.
 270  func (s *Server) CompleteHTTP(fd int32, status int32, headers map[string]string, body []byte, connClose bool) {
 271  	c := s.conns[fd]
 272  	if c == nil {
 273  		return
 274  	}
 275  	c.phase = phaseHTTP
 276  	s.sendHTTPBuffered(c, status, headers, body, connClose)
 277  }
 278  
 279  // connStateEnter makes the connection's own arena current for allocations that
 280  // must outlive this loop turn. Paired with connStateExit; the push/pop form
 281  // leaves the enclosing frame's arena untouched.
 282  func connStateEnter(c *tconn) (prev *runtime.Arena) {
 283  	prev = runtime.CurrentArena()
 284  	if c != nil && c.arena != nil {
 285  		runtime.SovereignSetArena(c.arena)
 286  	}
 287  	return
 288  }
 289  
 290  func connStateExit(prev *runtime.Arena) {
 291  	runtime.SovereignRestoreArena(prev)
 292  }
 293  
 294  // connStateReset drops everything the connection arena holds. Only valid when
 295  // the connection is idle: no parsed request and no buffered write referencing
 296  // it (keepAlive is exactly that point).
 297  func connStateReset(c *tconn) {
 298  	if c != nil && c.arena != nil {
 299  		runtime.ArenaReset(c.arena)
 300  	}
 301  }
 302  
 303  // setupListener performs the receiver writes for ListenAndServe: the listener
 304  // fd, the epoll fd, and the signal fd, plus their epoll registrations.
 305  //
 306  // These writes live in their own mutating method so that ListenAndServe itself
 307  // stays a read-only method of a sovereign type. A self-mutating method borrows
 308  // the receiver's sovereign arena as its working frame for its whole body, and
 309  // ListenAndServe's body is the event loop: keeping the borrow there would make
 310  // the loop's compaction yield point unable to compact the Server's own group,
 311  // which is exactly the group that grows with every connection. Peeling the
 312  // writes out leaves the loop running on a regular per-function arena with the
 313  // Server's arena merely registered as stable, which is what the compaction
 314  // rebuild pass knows how to repoint.
 315  func (s *Server) setupListener(fd int32) (eerr error) {
 316  	epfd, serr := syscall.EpollCreate1(0)
 317  	if serr != nil {
 318  		syscall.Close(fd)
 319  		return fmt.Errorf("transport: epoll: %w", serr)
 320  	}
 321  	if aerr := epollAdd(epfd, fd); aerr != nil {
 322  		syscall.Close(epfd)
 323  		syscall.Close(fd)
 324  		return fmt.Errorf("transport: epoll add: %w", aerr)
 325  	}
 326  	s.lnFD = fd
 327  	s.epfd = epfd
 328  	if globalSigFD >= 0 {
 329  		s.sigFD = int32(globalSigFD)
 330  		if aerr2 := epollAdd(epfd, s.sigFD); aerr2 != nil {
 331  			return fmt.Errorf("transport: epoll add signal pipe: %w", aerr2)
 332  		}
 333  	}
 334  	return nil
 335  }
 336  
 337  // ListenAndServe runs the epoll event loop. Returns on signal or error.
 338  func (s *Server) ListenAndServe(addr string) (eerr error) {
 339  	ip, port := ParseAddr(addr)
 340  
 341  	fd, serr := syscall.Socket(syscall.AF_INET, syscall.SOCK_STREAM, 0)
 342  	if serr != nil {
 343  		return fmt.Errorf("transport: socket: %w", serr)
 344  	}
 345  	syscall.SetsockoptInt(fd, syscall.SOL_SOCKET, syscall.SO_REUSEADDR, 1)
 346  	if serr = syscall.SetNonblock(fd, true); serr != nil {
 347  		syscall.Close(fd)
 348  		return fmt.Errorf("transport: nonblock: %w", serr)
 349  	}
 350  	sa := &syscall.SockaddrInet4{Port: port, Addr: ip}
 351  	if serr = syscall.Bind(fd, sa); serr != nil {
 352  		syscall.Close(fd)
 353  		return fmt.Errorf("transport: bind %s: %w", addr, serr)
 354  	}
 355  	if serr = syscall.Listen(fd, 4096); serr != nil {
 356  		syscall.Close(fd)
 357  		return fmt.Errorf("transport: listen: %w", serr)
 358  	}
 359  	if serr = s.setupListener(fd); serr != nil {
 360  		return serr
 361  	}
 362  
 363  	if s.OnReady != nil {
 364  		s.OnReady()
 365  	}
 366  
 367  	// Poll interval: worker responses arrive over spawn channels, which have
 368  	// no fd to wake epoll, so the loop must come back regularly to drain
 369  	// them. 2ms keeps the added latency imperceptible without busy-spinning.
 370  	const pollTimeoutMs = 2
 371  	const ticksPerOnTick = 2500 // 2ms * 2500 = 5s idle
 372  
 373  	events := []syscall.EpollEvent{:64}
 374  	idle := int32(0)
 375  	var n int32
 376  	var perr error
 377  	var evFD int32
 378  	var waitStart, bodyStart int64
 379  	for {
 380  		// Yield point: drain queued sovereign-arena compactions here, where no
 381  		// request frame is live, instead of inside a method boundary. Safe
 382  		// because ListenAndServe is read-only (see setupListener), so no frame
 383  		// on the stack borrows a sovereign arena, and every method it calls has
 384  		// already returned by the time the loop comes back around. The drain is
 385  		// bounded so a large compaction cannot stall the accept loop.
 386  		runtime.SovDrainCompactions(1)
 387  		// One arena per turn. This frame lives for the whole process, so
 388  		// anything a handler leaves behind would be permanent: the value an HTTP
 389  		// handler returns (a multi-MB static file) relocates up the call chain
 390  		// and comes to rest in the outermost live frame's arena, which is this
 391  		// one, and nothing ever reclaims it. A per-turn arena gives that data
 392  		// the lifetime of the turn and the arena is recycled. State that must
 393  		// outlive the turn goes in the connection's own arena (connStateEnter).
 394  		runtime.FnArenaPush(65536)
 395  		stop := false
 396  		waitStart = metrics.Now()
 397  		n, perr = syscall.EpollWait(s.epfd, events, pollTimeoutMs)
 398  		metrics.AcceptLoopWaitNs.Observe(metrics.Since(waitStart))
 399  		if perr != nil {
 400  			runtime.FnArenaPopFree()
 401  			if perr == syscall.EINTR {
 402  				return nil
 403  			}
 404  			return fmt.Errorf("transport: epoll wait: %w", perr)
 405  		}
 406  		for i := 0; i < n; i++ {
 407  			evFD = int32(events[i].Fd)
 408  			if evFD == s.sigFD {
 409  				moxie_signal_pipe_read()
 410  				stop = true
 411  				break
 412  			} else if evFD == s.lnFD {
 413  				s.acceptAll()
 414  			} else if c := s.conns[evFD]; c != nil {
 415  				if events[i].Events&(syscall.EPOLLERR|syscall.EPOLLHUP) != 0 {
 416  					s.closeConn(c)
 417  				} else if events[i].Events&syscall.EPOLLOUT != 0 {
 418  					s.drainWrite(c)
 419  				} else {
 420  					s.readConn(c)
 421  				}
 422  			}
 423  		}
 424  		if !stop {
 425  			bodyStart = metrics.Now()
 426  			s.handler.OnPoll()
 427  			metrics.AcceptLoopHandleNs.Observe(metrics.Since(bodyStart))
 428  			if n > 0 {
 429  				idle = 0
 430  			} else {
 431  				idle++
 432  				if idle >= ticksPerOnTick {
 433  					idle = 0
 434  					s.handler.OnTick()
 435  				}
 436  			}
 437  		}
 438  		runtime.FnArenaPopFree()
 439  		if stop {
 440  			return nil
 441  		}
 442  	}
 443  }
 444  
 445  func (s *Server) acceptAll() {
 446  	// Declared outside the loop: a declaration in the body of a self-mutating
 447  	// method allocates in the sovereign arena on every iteration.
 448  	var nfd int32
 449  	var sa syscall.Sockaddr
 450  	var err error
 451  	var ip string
 452  	var aerr error
 453  	var ok bool
 454  	for {
 455  		nfd, sa, err = syscall.Accept4(s.lnFD, syscall.SOCK_NONBLOCK)
 456  		if err != nil {
 457  			return
 458  		}
 459  		// The connection record outlives this call: it stays in s.conns for
 460  		// the life of the connection. The record, its receive buffer and the
 461  		// remote address string it keeps must all be allocated in the root
 462  		// arena - the address string is built here, so the borrow has to
 463  		// cover it too, not just the composite literal.
 464  		prev := runtime.CurrentArena()
 465  		runtime.SovereignSetArena(runtime.RootArena())
 466  		ip = peerAddr(sa)
 467  		ok = s.handler.OnAccept(nfd, ip)
 468  		if !ok {
 469  			runtime.SovereignRestoreArena(prev)
 470  			syscall.Close(nfd)
 471  			continue
 472  		}
 473  		if aerr = epollAdd(s.epfd, nfd); aerr != nil {
 474  			runtime.SovereignRestoreArena(prev)
 475  			syscall.Close(nfd)
 476  			continue
 477  		}
 478  		s.conns[nfd] = &tconn{
 479  			fd:       nfd,
 480  			phase:    phaseHTTP,
 481  			buf:      []byte{:4096},
 482  			remoteIP: ip,
 483  			arena:    runtime.ArenaNew(65536),
 484  		}
 485  		runtime.SovereignRestoreArena(prev)
 486  	}
 487  }
 488  
 489  func peerAddr(sa syscall.Sockaddr) (s string) {
 490  	if sa4, ok := sa.(*syscall.SockaddrInet4); ok {
 491  		// Grow with the append operator: a presized []byte{:0:20} plus push
 492  		// is a bounded store that fails loud the moment the capacity is not
 493  		// what the literal promised.
 494  		b := []byte{}
 495  		b = appendInt(b, int32(sa4.Addr[0]))
 496  		b = b | "."
 497  		b = appendInt(b, int32(sa4.Addr[1]))
 498  		b = b | "."
 499  		b = appendInt(b, int32(sa4.Addr[2]))
 500  		b = b | "."
 501  		b = appendInt(b, int32(sa4.Addr[3]))
 502  		return string(makeCopy(b))
 503  	}
 504  	return "unknown"
 505  }
 506  
 507  // growBuf doubles buf, preserving the first wpos bytes. Free function: the
 508  // abandoned buffer dies with this frame's arena instead of living on in the
 509  // Server's sovereign arena for every growth step.
 510  func growBuf(buf []byte, wpos int32) (nb []byte) {
 511  	nb = []byte{:len(buf) * 2}
 512  	copy(nb, buf[:wpos])
 513  	return nb
 514  }
 515  
 516  func (s *Server) readConn(c *tconn) {
 517  	if c.wbuf != nil {
 518  		return
 519  	}
 520  	avail := len(c.buf) - c.wpos
 521  	if avail < 512 {
 522  		if len(c.buf) >= maxBuf {
 523  			s.closeConn(c)
 524  			return
 525  		}
 526  		c.buf = growBuf(c.buf, c.wpos)
 527  	}
 528  	n, err := syscall.Read(c.fd, c.buf[c.wpos:])
 529  	if n <= 0 {
 530  		if err == nil || (!errAgain(err) && !errIntr(err)) {
 531  			s.closeConn(c)
 532  		}
 533  		return
 534  	}
 535  	c.wpos += n
 536  
 537  	switch c.phase {
 538  	case phaseHTTP:
 539  		s.processHTTP(c)
 540  	case phaseHTTPBody:
 541  		s.processHTTPBody(c)
 542  	case phaseWS:
 543  		s.processWS(c)
 544  	case phaseHTTPDeferred:
 545  		c.wpos = 0 // drop bytes while awaiting async response
 546  	}
 547  }
 548  
 549  func (s *Server) closeConn(c *tconn) {
 550  	epollDel(s.epfd, c.fd)
 551  	syscall.Close(c.fd)
 552  	c.buf = nil
 553  	c.wbuf = nil
 554  	c.pendingReq = nil
 555  	if c.arena != nil {
 556  		runtime.ArenaFree(c.arena)
 557  		c.arena = nil
 558  	}
 559  	delete(s.conns, c.fd)
 560  	if c.phase == phaseWS {
 561  		s.handler.OnWSClose(c.fd)
 562  		if n := s.ipConns[c.remoteIP] - 1; n <= 0 {
 563  			delete(s.ipConns, c.remoteIP)
 564  		} else {
 565  			s.ipConns[c.remoteIP] = n
 566  		}
 567  	}
 568  }
 569  
 570  func (s *Server) keepAlive(c *tconn) {
 571  	// Idle: no parsed request and no buffered write can reference the
 572  	// connection arena, so its accumulated per-request state can be released
 573  	// now instead of at close.
 574  	connStateReset(c)
 575  	c.phase = phaseHTTP
 576  	c.pendingReq = nil
 577  	c.bodyNeeded = 0
 578  	if c.wpos > 0 {
 579  		s.processHTTP(c)
 580  	}
 581  }
 582  
 583  func (s *Server) processHTTP(c *tconn) {
 584  	data := c.buf[:c.wpos]
 585  	end := bytes.Index(data, []byte("\r\n\r\n"))
 586  	if end < 0 {
 587  		return
 588  	}
 589  	consumed := end + 4
 590  
 591  	req := parseHTTPHeaders(data[:end])
 592  	if req == nil {
 593  		s.closeConn(c)
 594  		return
 595  	}
 596  
 597  	copy(c.buf, c.buf[consumed:c.wpos])
 598  	c.wpos -= consumed
 599  
 600  	// Reverse proxy: substitute real IP from X-Forwarded-For before any checks.
 601  	if c.remoteIP == "127.0.0.1" {
 602  		if xff := req.headers["x-forwarded-for"]; xff != "" {
 603  			realIP := FirstXFF(xff)
 604  			if len(realIP) > 0 && realIP != c.remoteIP {
 605  				prevR := runtime.CurrentArena()
 606  				runtime.SovereignSetArena(runtime.RootArena())
 607  				c.remoteIP = realIP
 608  				runtime.SovereignRestoreArena(prevR)
 609  			}
 610  		}
 611  	}
 612  
 613  	if bytes.EqualFold([]byte(req.headers["upgrade"]), []byte("websocket")) {
 614  		s.upgradeWS(c, req)
 615  		return
 616  	}
 617  
 618  	if s.BotBlock && isBot(req.headers["user-agent"]) {
 619  		writeHTTPResponse(c.fd, 403, nil, []byte("forbidden"))
 620  		s.closeConn(c)
 621  		return
 622  	}
 623  
 624  	cl := parseContentLength(req.headers["content-length"])
 625  	if cl > 0 && cl <= maxBuf {
 626  		prevP := connStateEnter(c)
 627  		c.pendingReq = parseHTTPHeaders(makeCopy(data[:end]))
 628  		connStateExit(prevP)
 629  		c.bodyNeeded = cl
 630  		c.phase = phaseHTTPBody
 631  		s.processHTTPBody(c)
 632  		return
 633  	}
 634  
 635  	status, headers, body, connClose := s.handler.OnHTTP(c.fd, req.method, req.path, req.headers, nil)
 636  	if status == HTTPDeferred {
 637  		c.phase = phaseHTTPDeferred
 638  		return
 639  	}
 640  	s.sendHTTPBuffered(c, status, headers, body, connClose || req.headers["connection"] == "close")
 641  }
 642  
 643  func (s *Server) processHTTPBody(c *tconn) {
 644  	if c.wpos < c.bodyNeeded {
 645  		return
 646  	}
 647  	prevB := connStateEnter(c)
 648  	c.pendingReq.body = makeCopy(c.buf[:c.bodyNeeded])
 649  	connStateExit(prevB)
 650  	copy(c.buf, c.buf[c.bodyNeeded:c.wpos])
 651  	c.wpos -= c.bodyNeeded
 652  	req := c.pendingReq
 653  	c.pendingReq = nil
 654  
 655  	status, headers, body, connClose := s.handler.OnHTTP(c.fd, req.method, req.path, req.headers, req.body)
 656  	if status == HTTPDeferred {
 657  		c.phase = phaseHTTPDeferred
 658  		return
 659  	}
 660  	s.sendHTTPBuffered(c, status, headers, body, connClose || req.headers["connection"] == "close")
 661  }
 662  
 663  // sendHTTPBuffered writes an HTTP response, buffering on EAGAIN.
 664  func (s *Server) sendHTTPBuffered(c *tconn, status int32, headers map[string]string, body []byte, connClose bool) {
 665  	s.connWrite(c, buildHTTPResponse(status, headers, body), connClose)
 666  	if c.wbuf == nil && s.conns[c.fd] != nil && !connClose {
 667  		s.keepAlive(c)
 668  	}
 669  }
 670  
 671  // drainWrite flushes pending write buffer when socket becomes writable.
 672  func (s *Server) drainWrite(c *tconn) {
 673  	for len(c.wbuf) > 0 {
 674  		n, err := syscall.Write(c.fd, c.wbuf)
 675  		if n > 0 {
 676  			c.wbuf = c.wbuf[n:]
 677  		}
 678  		if errAgain(err) {
 679  			return
 680  		}
 681  		if errIntr(err) {
 682  			continue
 683  		}
 684  		if err != nil {
 685  			s.closeConn(c)
 686  			return
 687  		}
 688  	}
 689  	c.wbuf = nil
 690  	epollModRead(s.epfd, c.fd)
 691  	if c.wbufClose {
 692  		s.closeConn(c)
 693  		return
 694  	}
 695  	if c.phase == phaseHTTP || c.phase == phaseHTTPBody {
 696  		s.keepAlive(c)
 697  	}
 698  }
 699  
 700  // Epoll helpers.
 701  
 702  func epollAdd(epfd, fd int32) (eerr error) {
 703  	ev := syscall.EpollEvent{Events: syscall.EPOLLIN, Fd: int32(fd)}
 704  	return syscall.EpollCtl(epfd, syscall.EPOLL_CTL_ADD, fd, &ev)
 705  }
 706  
 707  func epollModWrite(epfd, fd int32) {
 708  	ev := syscall.EpollEvent{Events: syscall.EPOLLIN | syscall.EPOLLOUT, Fd: int32(fd)}
 709  	syscall.EpollCtl(epfd, syscall.EPOLL_CTL_MOD, fd, &ev)
 710  }
 711  
 712  func epollModRead(epfd, fd int32) {
 713  	ev := syscall.EpollEvent{Events: syscall.EPOLLIN, Fd: int32(fd)}
 714  	syscall.EpollCtl(epfd, syscall.EPOLL_CTL_MOD, fd, &ev)
 715  }
 716  
 717  func epollDel(epfd, fd int32) (eerr error) {
 718  	return syscall.EpollCtl(epfd, syscall.EPOLL_CTL_DEL, fd, nil)
 719  }
 720  
 721  func init() {
 722  	globalSigFD = -1
 723  	botAgents = []string{
 724  		"SemrushBot", "AhrefsBot", "GPTBot", "ClaudeBot",
 725  		"MJ12bot", "DotBot", "PetalBot", "Bytespider",
 726  		"Sogou", "YandexBot", "BLEXBot", "DataForSeoBot",
 727  	}
 728  }
 729