server.mx raw

   1  // Package server provides the Nostr relay domain coordinator.
   2  // It implements transport.Handler and wires together the event store,
   3  // worker pool, broadcast domain, and HTTP/WebSocket routing.
   4  package server
   5  
   6  import (
   7  	"fmt"
   8  	"runtime"
   9  	"net/url"
  10  	"time"
  11  
  12  
  13  	"git.smesh.lol/morly/pkg/access"
  14  	"git.smesh.lol/morly/pkg/broadcast"
  15  	"git.smesh.lol/morly/pkg/metrics"
  16  	"git.smesh.lol/nostr/pkg/envelope"
  17  	"git.smesh.lol/nostr/pkg/event"
  18  	"git.smesh.lol/nostr/pkg/filter"
  19  	"git.smesh.lol/morly/pkg/pool"
  20  	"git.smesh.lol/morly/pkg/relay/config"
  21  	"git.smesh.lol/morly/pkg/relay/dbengine"
  22  	"git.smesh.lol/morly/pkg/relay/ratelimit"
  23  	"git.smesh.lol/morly/pkg/relay/tree"
  24  	"git.smesh.lol/morly/pkg/relay/wire"
  25  	"git.smesh.lol/morly/pkg/transport"
  26  )
  27  
  28  // InitSignals delegates to transport.InitSignals.
  29  func InitSignals() { transport.InitSignals() }
  30  
  31  
  32  // DB is root's handle to the database-engine domain. Root never opens the
  33  // store: it sends requests and learns of replies through the ready nudge,
  34  // because a select receive on a codec-framed spawn channel does not decode.
  35  type DB struct {
  36  	in    chan tree.Request
  37  	out   chan tree.Response
  38  	ready chan struct{}
  39  	done  chan struct{}
  40  }
  41  
  42  // Valid reports whether the database-engine domain is still running.
  43  func (d *DB) Valid() (ok bool) {
  44  	select {
  45  	case <-d.done:
  46  		return false
  47  	default:
  48  	}
  49  	return true
  50  }
  51  
  52  // SpawnDB forks the database-engine domain and returns root's handle to it.
  53  // Channel construction and the fork both happen under the root arena: the
  54  // channels outlive this call and the forked domain reads the same memory.
  55  func SpawnDB(cfg *config.C) (d *DB) {
  56  	dbCfg := dbengine.Config{
  57  		DataDir:             cfg.DataDir,
  58  		ACLMode:             cfg.ACLMode,
  59  		Admins:              cfg.Admins,
  60  		FollowListFreqSec:   cfg.FollowListFreqSec,
  61  		SocialWoTMaxDepth:   cfg.SocialWoTMaxDepth,
  62  		SocialWoTRefreshSec: cfg.SocialWoTRefreshSec,
  63  		MuteBlacklist:       cfg.MuteBlacklist,
  64  	}
  65  	runtime.SovereignSetArena(runtime.RootArena())
  66  	in := chan tree.Request{}
  67  	out := chan tree.Response{}
  68  	ready := chan struct{}{:64}
  69  	done := spawn(dbengine.Run, dbCfg, in, out, ready)
  70  	runtime.SovereignRestoreArena(runtime.RootArena())
  71  	return &DB{in: in, out: out, ready: ready, done: done}
  72  }
  73  
  74  // Server is the Nostr relay domain coordinator.
  75  // It implements transport.Handler so the transport layer can call back
  76  // for connection lifecycle events and incoming messages.
  77  type Server struct {
  78  	Fallback func(method, path string, headers map[string]string, body []byte) (int32, map[string]string, []byte)
  79  	OnReady  func()
  80  	Version  string
  81  	selfHost string // "host/" prefix for proxy self-redirect detection
  82  	t    *transport.Server
  83  	cfg  *config.C
  84  	// tickCount counts OnTick calls for the periodic cleanup; it lives on the
  85  	// server rather than in a package global, which the rule forbids writing
  86  	// outside init.
  87  	tickCount int32
  88  	// per-connection domain state, indexed by FD (parallel to transport.conns)
  89  	conns        map[int32]*cstate
  90  	writeLimiter *ratelimit.Limiter
  91  	// closed makes Close idempotent: both shutdown paths call it.
  92  	closed bool
  93  	// Database engine handle; all storage goes through it.
  94  	db *DB
  95  	// In-flight database requests, keyed by reqID.
  96  	dbPending map[uint32]int32
  97  	nextDBReq uint32
  98  	// Ingest worker pool (signature gate).
  99  	workerIn    []chan wire.IngestRequest
 100  	workerOut   []chan wire.IngestResponse
 101  	workerDone  []chan struct{}
 102  	workerReady []chan struct{}
 103  	workers     pool.Pool
 104  	pendingReq  map[uint32]int32 // reqID → conn fd
 105  	nextReqID   uint32
 106  	// Broadcast domain.
 107  	bcast *broadcast.Broadcaster
 108  	// Media proxy worker pool.
 109  	proxyIn       []chan wire.ProxyRequest
 110  	proxyOut      []chan wire.ProxyResponse
 111  	proxyDone     []chan struct{}
 112  	proxyPool     pool.Pool
 113  	proxyQueue    []pendingProxy
 114  	proxyBusyTime []int64 // unix nanos when worker went busy; 0 = idle
 115  	// Blossom worker pool.
 116  	blossomIn   []chan wire.BlossomRequest
 117  	blossomOut  []chan wire.BlossomResponse
 118  	blossomDone []chan struct{}
 119  	blossomPool pool.Pool
 120  	// Shared async HTTP pending.
 121  	asyncPending map[uint32]asyncHTTPEntry
 122  	nextAsyncID  uint32
 123  }
 124  
 125  type asyncHTTPEntry struct {
 126  	connFD    int32
 127  	connClose bool
 128  	createdAt int64 // unix nanos, for reap timeout
 129  }
 130  
 131  type pendingProxy struct {
 132  	connFD    int32
 133  	connClose bool
 134  	req       wire.ProxyRequest
 135  }
 136  
 137  // cstate is the server-domain per-connection state (server-side complement to
 138  // transport's tconn, keyed by the same FD).
 139  type cstate struct {
 140  	subs         map[string]*sub
 141  	challenge    []byte
 142  	authedPubkey []byte
 143  	jailed       bool
 144  }
 145  
 146  type sub struct {
 147  	id      string
 148  	filters filter.S
 149  	rawReq  []byte
 150  }
 151  
 152  // New creates a relay server backed by the given database-engine handle.
 153  func New(db *DB, cfg *config.C) (s *Server) {
 154  	var wlim *ratelimit.Limiter
 155  	if cfg.FreeWriteLimit > 0 && cfg.FreeWriteWindow > 0 {
 156  		rate := float64(cfg.FreeWriteLimit) / float64(cfg.FreeWriteWindow)
 157  		wlim = ratelimit.New(rate, cfg.FreeWriteLimit)
 158  	}
 159  	selfHost := "git.smesh.lol/morly/"
 160  	if u, err := url.Parse(cfg.RelayURL); err == nil && len(u.Host) > 0 {
 161  		selfHost = u.Host | "/"
 162  	}
 163  	s = &Server{
 164  		cfg:          cfg,
 165  		selfHost:     selfHost,
 166  		conns:        map[int32]*cstate{},
 167  		writeLimiter: wlim,
 168  		db:           db,
 169  		dbPending:    map[uint32]int32{},
 170  		asyncPending: map[uint32]asyncHTTPEntry{},
 171  	}
 172  
 173  	// Create transport. Server implements transport.Handler.
 174  	s.t = transport.New(s, cfg.MaxConnPerIP)
 175  	s.t.BotBlock = cfg.HTTPGuardBotBlock
 176  
 177  	// Spawn worker domains AFTER heap is settled (GC avoidance).
 178  	if cfg.IngestWorkers > 0 {
 179  		s.startIngestWorkers(cfg.IngestWorkers)
 180  	}
 181  	s.startBroadcastWorker()
 182  	if n := cfg.MediaProxyWorkers; n > 0 {
 183  		s.startProxyWorkers(n)
 184  	}
 185  	if n := cfg.BlossomWorkers; n > 0 {
 186  		s.startBlossomWorkers(n)
 187  	}
 188  	return s
 189  }
 190  
 191  // Close stops the database-engine domain, which flushes and closes the store.
 192  func (s *Server) Close() {
 193  	// Close runs on both shutdown paths (main calls it directly and again from
 194  	// a defer), and closing a channel twice panics, so the first call wins.
 195  	if s.closed {
 196  		return
 197  	}
 198  	s.closed = true
 199  	if s.db != nil && s.db.Valid() {
 200  		close(s.db.in)
 201  	}
 202  	// Every worker loop returns when its request channel closes, and a domain
 203  	// that returns dumps its coverage counters on the way out. Without this the
 204  	// pool domains stayed blocked in a futex receive while the fixture's SIGTERM
 205  	// went to the process group, so they were SIGKILLed with their counters
 206  	// (private after fork) and every worker file read 0% coverage.
 207  	for _, in := range s.workerIn {
 208  		close(in)
 209  	}
 210  	for _, in := range s.proxyIn {
 211  		close(in)
 212  	}
 213  	for _, in := range s.blossomIn {
 214  		close(in)
 215  	}
 216  	// The broadcast worker owns the subscription registry and is pre-spawned
 217  	// by broadcast.PreSpawn before the store; close it too, or it is the one
 218  	// domain still parked in a receive when the process group is signalled.
 219  	s.bcast.Close()
 220  }
 221  
 222  func (s *Server) startIngestWorkers(n int32) {
 223  	s.workerIn, s.workerOut, s.workerDone, s.workerReady = spawnIngestPool(n)
 224  	s.workers = pool.NewPool(n)
 225  	s.pendingReq = map[uint32]int32{}
 226  }
 227  
 228  // spawnIngestPool forks the ingest domains in a free function; see
 229  // spawnProxyPool for why the loop does not live in the Server method.
 230  //
 231  // The pool slices are built with push, not []chan T{:n}: an inline
 232  // slice-of-channel literal is miscompiled to an empty slice (the literal
 233  // rewrite does not resolve the channel element type), which is why the named
 234  // slice aliases existed. push allocates the same []chan T correctly.
 235  func spawnIngestPool(n int32) (ins []chan wire.IngestRequest, outs []chan wire.IngestResponse, dones []chan struct{}, readys []chan struct{}) {
 236  	var in chan wire.IngestRequest
 237  	var out chan wire.IngestResponse
 238  	var done chan struct{}
 239  	var ready chan struct{}
 240  	for i := 0; i < n; i++ {
 241  		in, out, done, ready = newIngestWorker()
 242  		ins = push(ins, in)
 243  		outs = push(outs, out)
 244  		dones = push(dones, done)
 245  		readys = push(readys, ready)
 246  	}
 247  	return
 248  }
 249  
 250  // newIngestWorker creates the channel pair and forks one ingest domain.
 251  // A free function: the channel/spawn scratch belongs to this frame, not to
 252  // the caller's sovereign arena.
 253  func newIngestWorker() (in chan wire.IngestRequest, out chan wire.IngestResponse, done chan struct{}, ready chan struct{}) {
 254  	// The channel structs outlive this frame: the server keeps them for the
 255  	// whole run and the forked domain reads the same memory. chanMake
 256  	// allocates in the current arena, so a plain frame here would leave every
 257  	// stored channel pointing into a released arena.
 258  	runtime.SovereignSetArena(runtime.RootArena())
 259  	in = chan wire.IngestRequest{}
 260  	out = chan wire.IngestResponse{}
 261  	// Ready carries no payload; it exists because a select receive on a
 262  	// codec-framed spawn channel does not decode. The worker nudges here
 263  	// after writing out, and the parent reads out with a plain receive.
 264  	ready = chan struct{}{:64}
 265  	done = spawn(wire.IngestWorker, in, out, ready)
 266  	runtime.SovereignRestoreArena(runtime.RootArena())
 267  	return
 268  }
 269  
 270  func (s *Server) startBroadcastWorker() {
 271  	s.bcast = broadcast.New()
 272  }
 273  
 274  func (s *Server) respawnIngestWorker(i int32) {
 275  	in, out, done, ready := newIngestWorker()
 276  	s.workerIn[i] = in
 277  	s.workerOut[i] = out
 278  	s.workerDone[i] = done
 279  	s.workerReady[i] = ready
 280  	s.workers.Busy[i] = false
 281  	fmt.Println("respawned ingest worker", i)
 282  }
 283  
 284  func (s *Server) respawnBroadcastWorker() {
 285  	s.bcast = broadcast.New()
 286  	for fd, cs := range s.conns {
 287  		s.bcast.OnConnect(int32(fd), s.t.ConnIsWhitelisted(fd))
 288  		if len(cs.authedPubkey) > 0 {
 289  			s.bcast.OnAuth(int32(fd), cs.authedPubkey)
 290  		}
 291  		for _, sb := range cs.subs {
 292  			s.bcast.OnSubscribe(int32(fd), []byte(sb.id), sb.rawReq)
 293  		}
 294  	}
 295  	fmt.Println("respawned broadcast worker, re-registered", len(s.conns), "connections")
 296  }
 297  
 298  // wireOnReady hands the transport the startup callback. It is its own mutating
 299  // method so that ListenAndServe itself stays a read-only method of a sovereign
 300  // type: a mutating method borrows the receiver's sovereign arena as its working
 301  // frame for its whole body, and ListenAndServe's body is the transport event
 302  // loop, so a write there would keep this Server's arena borrowed across every
 303  // turn of the loop. The loop's compaction yield point found the group borrowed
 304  // every single time and could only rotate it back onto the queue: the group
 305  // could never be compacted, which is what made the relay's arena grow without
 306  // bound. Peeling the write out leaves the loop running with the arena merely
 307  // registered as stable, which the compaction rebuild pass knows how to repoint.
 308  func (s *Server) wireOnReady() {
 309  	s.t.OnReady = s.OnReady
 310  }
 311  
 312  // ListenAndServe starts the transport event loop.
 313  func (s *Server) ListenAndServe(addr string) (err error) {
 314  	s.wireOnReady()
 315  	return s.t.ListenAndServe(addr)
 316  }
 317  
 318  // --- transport.Handler implementation ---
 319  
 320  // OnAccept gates new TCP connections: IP blacklist and global connection limit.
 321  func (s *Server) OnAccept(fd int32, ip string) (ok bool) {
 322  	if s.cfg.MaxGlobalConns > 0 && s.t.ConnCount() >= s.cfg.MaxGlobalConns {
 323  		return false
 324  	}
 325  	return !s.ipBlacklisted(ip)
 326  }
 327  
 328  // OnWSUpgrade gates WS upgrades: per-IP connection limit.
 329  func (s *Server) OnWSUpgrade(fd int32, ip string, currentIPWSCount int32) (whitelisted bool, allow bool) {
 330  	wl := s.ipWhitelisted(ip)
 331  	if !wl && s.cfg.MaxConnPerIP > 0 && currentIPWSCount >= s.cfg.MaxConnPerIP {
 332  		return false, false
 333  	}
 334  	return wl, true
 335  }
 336  
 337  // OnWSConnected initialises per-conn state and sends the NIP-42 challenge.
 338  func (s *Server) OnWSConnected(fd int32) {
 339  	// The per-connection state outlives this call: the server keeps it in
 340  	// s.conns for the life of the connection. The record, the subscription
 341  	// map, the NIP-42 challenge bytes and the boxed challenge message must all
 342  	// be allocated in the root arena - the challenge is built here, so the
 343  	// borrow covers it too.
 344  	var ac *envelope.AuthChallenge
 345  	var ch []byte
 346  	prev := runtime.CurrentArena()
 347  	runtime.SovereignSetArena(runtime.RootArena())
 348  	s.conns[fd] = &cstate{subs: map[string]*sub{}}
 349  	if s.cfg.RelayURL != "" {
 350  		ch = transport.Challenge()
 351  		s.conns[fd].challenge = ch
 352  		ac = &envelope.AuthChallenge{Challenge: ch}
 353  	}
 354  	runtime.SovereignRestoreArena(prev)
 355  	s.bcast.OnConnect(int32(fd), s.t.ConnIsWhitelisted(fd))
 356  	if ac != nil {
 357  		s.t.SendWS(fd, ac.Marshal(nil))
 358  	}
 359  }
 360  
 361  // OnWSMessage dispatches a decoded WebSocket payload.
 362  func (s *Server) OnWSMessage(fd int32, payload []byte) {
 363  	s.dispatch(fd, payload)
 364  }
 365  
 366  // OnWSClose cleans up per-conn state when a WS connection closes.
 367  func (s *Server) OnWSClose(fd int32) {
 368  	c := s.conns[fd]
 369  	if c != nil {
 370  		for k, sb := range c.subs {
 371  			sb.rawReq = nil
 372  			sb.filters.F = nil
 373  			delete(c.subs, k)
 374  		}
 375  		c.subs = nil
 376  		c.challenge = nil
 377  		c.authedPubkey = nil
 378  	}
 379  	delete(s.conns, fd)
 380  	s.bcast.OnDisconnect(int32(fd))
 381  	for rid, connFD := range s.pendingReq {
 382  		if connFD == fd {
 383  			delete(s.pendingReq, rid)
 384  		}
 385  	}
 386  	for rid, connFD := range s.dbPending {
 387  		if connFD == fd {
 388  			delete(s.dbPending, rid)
 389  		}
 390  	}
 391  }
 392  
 393  // OnHTTP routes HTTP requests to workers or synchronous handlers.
 394  // Returns transport.HTTPDeferred to signal async processing.
 395  func (s *Server) OnHTTP(fd int32, method, path string, headers map[string]string, body []byte) (status int32, hdrs map[string]string, out []byte, ok bool) {
 396  	// IP blacklist re-check after XFF substitution.
 397  	if s.ipBlacklisted(s.t.ConnIP(fd)) {
 398  		return 403, nil, nil, true
 399  	}
 400  	if transport.HasPrefix(path, "/proxy/") && s.proxyPool.Len() > 0 {
 401  		if s.doDispatchProxy(fd, path, headers) {
 402  			return transport.HTTPDeferred, nil, nil, false
 403  		}
 404  	}
 405  	if transport.HasPrefix(path, "/blossom") && s.blossomPool.Len() > 0 {
 406  		if s.doDispatchBlossom(fd, method, path, headers, body) {
 407  			return transport.HTTPDeferred, nil, nil, false
 408  		}
 409  	}
 410  	connClose := headers["connection"] == "close"
 411  	if method == "OPTIONS" && s.cfg.CORSEnabled {
 412  		return 204, s.corsHeaders(headers), nil, connClose
 413  	}
 414  	status, h, b := s.routeHTTP(fd, method, path, headers, body)
 415  	return status, h, b, connClose
 416  }
 417  
 418  // OnPoll drains worker-domain output channels. The transport loop calls it
 419  // once per iteration; spawn channels have no fd to wake epoll.
 420  func (s *Server) OnPoll() {
 421  	pollStart := metrics.Now()
 422  	s.pollDB()
 423  	s.pollIngestWorkers()
 424  	s.pollProxyWorkers()
 425  	s.pollBlossomWorkers()
 426  	s.handleBroadcastFrame()
 427  	metrics.OnPollNs.Observe(metrics.Since(pollStart))
 428  }
 429  
 430  
 431  func (s *Server) OnTick() {
 432  	s.proxyReapStuck()
 433  	s.asyncReapStuck()
 434  	s.tickCount++
 435  	if s.tickCount%60 == 0 {
 436  		if s.writeLimiter != nil {
 437  			s.writeLimiter.Cleanup(5 * time.Minute)
 438  		}
 439  	}
 440  }
 441  
 442  
 443  // --- Nostr message dispatch ---
 444  
 445  func (s *Server) dispatch(fd int32, msg []byte) {
 446  	c := s.conns[fd]
 447  	if c != nil && c.jailed {
 448  		return
 449  	}
 450  	label, _, _ := envelope.Identify(msg)
 451  	switch label {
 452  	case envelope.EventLabel:
 453  		s.handleEvent(fd, msg)
 454  	case envelope.ReqLabel:
 455  		s.handleReq(fd, msg)
 456  	case envelope.CloseLabel:
 457  		s.handleClose(fd, msg)
 458  	case envelope.CountLabel:
 459  		s.handleCount(fd, msg)
 460  	case envelope.AuthLabel:
 461  		s.handleAuth(fd, msg)
 462  	}
 463  }
 464  
 465  // handleEvent runs one EVENT frame.
 466  //
 467  // It stays a method on *Server. Moving the body into a free function that
 468  // returns the frame to send - which would give it a frame arena and free its
 469  // scratch at return - was tried and reverted: with the free function the
 470  // free-write limiter stopped applying (the fourth unauthed write was accepted
 471  // where it has to be rate-limited). The limiter itself is not at fault: a
 472  // counter with the same shape keeps its state when driven from a free function
 473  // in isolation, so the cause is still unidentified. Keep the body in the
 474  // method until it is.
 475  func (s *Server) handleEvent(fd int32, msg []byte) {
 476  	handleStart := metrics.Now()
 477  	defer func() { metrics.HandleEventNs.Observe(metrics.Since(handleStart)) }()
 478  	c := s.conns[fd]
 479  	if c == nil {
 480  		return
 481  	}
 482  	cfg := s.cfg
 483  	authed := len(c.authedPubkey) > 0
 484  	gated := cfg.RelayURL != ""
 485  	wl := s.t.ConnIsWhitelisted(fd)
 486  
 487  	if gated && !authed && !wl {
 488  		evKind := peekEventKind(msg)
 489  		exempt := access.WriteExempt(evKind, s.cfg.NIP46BypassAuth, s.cfg.MarmotOpen)
 490  		if !exempt && s.writeLimiter != nil {
 491  			ip := s.t.ConnIP(fd)
 492  			if s.writeLimiter.Allow([]byte(ip)) {
 493  				exempt = true
 494  			} else {
 495  				s.eventReject(fd, msg, "rate-limited: too many events, try again later")
 496  				return
 497  			}
 498  		}
 499  		if !exempt && cfg.AuthToWrite {
 500  			s.eventReject(fd, msg, "auth-required: authentication required")
 501  			return
 502  		}
 503  	}
 504  	if s.workers.Len() > 0 && s.dispatchToWorker(fd, msg) {
 505  		return
 506  	}
 507  	parseStart := metrics.Now()
 508  	_, rem, _ := envelope.Identify(msg)
 509  	var es envelope.EventSubmission
 510  	perr := error(nil)
 511  	_, perr = es.Unmarshal(rem)
 512  	metrics.EnvelopeParseNs.Observe(metrics.Since(parseStart))
 513  	if perr != nil || es.E == nil {
 514  		s.t.SendWS(fd, (&envelope.OK{EventID: []byte{:32}, Reason: []byte("invalid: parse error")}).Marshal(nil))
 515  		return
 516  	}
 517  	if !s.dbSend(tree.Request{Op: tree.OpPersist, ConnID: fd, Bytes: msg}) {
 518  		s.eventReject(fd, msg, "error: storage unavailable")
 519  	}
 520  }
 521  
 522  // dbSend hands one request to the database-engine domain and records the
 523  // connection that is waiting for the reply. Returns false when the domain is
 524  // gone, so the caller can answer instead of leaving the client hanging.
 525  func (s *Server) dbSend(req tree.Request) (ok bool) {
 526  	if s.db == nil || !s.db.Valid() {
 527  		return false
 528  	}
 529  	s.nextDBReq++
 530  	req.ReqID = s.nextDBReq
 531  	s.dbPending[req.ReqID] = req.ConnID
 532  	s.db.in <- req
 533  	return true
 534  }
 535  
 536  // pollDB drains ready database-engine replies. The worker nudges root on a
 537  // zero-size channel after each reply; the reply itself travels on the codec
 538  // channel and is read with a plain receive, which decodes.
 539  func (s *Server) pollDB() {
 540  	for i := 0; i < 256; i++ {
 541  		select {
 542  		case <-s.db.ready:
 543  			resp := <-s.db.out
 544  			s.completeDB(resp)
 545  		default:
 546  			return
 547  		}
 548  	}
 549  }
 550  
 551  func (s *Server) completeDB(resp tree.Response) {
 552  	switch resp.Op {
 553  	case tree.OpPersist:
 554  		s.completePersist(resp)
 555  	case tree.OpHistory:
 556  		s.completeHistory(resp)
 557  	case tree.OpCount:
 558  		s.completeCount(resp)
 559  	}
 560  }
 561  
 562  func (s *Server) completePersist(resp tree.Response) {
 563  	fd, ok := s.dbPending[resp.ReqID]
 564  	if !ok {
 565  		return
 566  	}
 567  	delete(s.dbPending, resp.ReqID)
 568  	if s.conns[fd] == nil {
 569  		return
 570  	}
 571  	out := &envelope.OK{EventID: resp.EventID, OK: resp.OK, Reason: resp.Reason}
 572  	sendStart := metrics.Now()
 573  	s.t.SendWS(fd, out.Marshal(nil))
 574  	metrics.SendWSNs.Observe(metrics.Since(sendStart))
 575  	if resp.OK && len(resp.Bytes) > 0 {
 576  		bStart := metrics.Now()
 577  		s.sendBroadcast(resp.Bytes, fd)
 578  		metrics.BroadcastNs.Observe(metrics.Since(bStart))
 579  	}
 580  }
 581  
 582  // completeHistory streams the stored events for one subscription and closes it
 583  // with EOSE. A subscription closed while the query was in flight is dropped.
 584  func (s *Server) completeHistory(resp tree.Response) {
 585  	c := s.conns[resp.ConnID]
 586  	if c == nil {
 587  		return
 588  	}
 589  	if _, still := c.subs[string(resp.SubID)]; !still {
 590  		return
 591  	}
 592  	// SendWS buffers when the socket is momentarily full and flushes on a later
 593  	// loop pass. SendWSErr's direct write reports EAGAIN as an error, and the
 594  	// old "close the connection on error" turned that transient state into a
 595  	// dropped subscription: the browser's feed REQ - the only sub with an event
 596  	// to deliver - lost both the EVENT and the EOSE, while the tag-less subs
 597  	// answered with EOSE alone.
 598  	for _, frame := range resp.Events {
 599  		s.t.SendWS(resp.ConnID, frame)
 600  	}
 601  	s.t.SendWS(resp.ConnID, (&envelope.EOSE{Subscription: resp.SubID}).Marshal(nil))
 602  }
 603  
 604  func (s *Server) completeCount(resp tree.Response) {
 605  	if s.conns[resp.ConnID] == nil {
 606  		return
 607  	}
 608  	cr := &envelope.CountResponse{Subscription: resp.SubID, Count: resp.Count}
 609  	s.t.SendWS(resp.ConnID, cr.Marshal(nil))
 610  }
 611  
 612  func (s *Server) dispatchToWorker(fd int32, msg []byte) (ok bool) {
 613  	i := s.workers.IdleIndex()
 614  	if i < 0 {
 615  		return false
 616  	}
 617  	s.nextReqID++
 618  	rid := s.nextReqID
 619  	s.pendingReq[rid] = fd
 620  	s.workers.Busy[i] = true
 621  	s.workerIn[i] <- wire.IngestRequest{ReqID: rid, Bytes: msg}
 622  	return true
 623  }
 624  
 625  // pollIngestWorkers drains ready ingest responses. Called from the transport
 626  // loop each iteration; worker domains are reached over spawn channels, which
 627  // have no fd to wake epoll.
 628  //
 629  // Readiness arrives on the zero-size nudge channel, then the response is read
 630  // from the codec channel with a plain receive. A select receive directly on a
 631  // codec-framed spawn channel does not decode: the runtime treats the encoded
 632  // ring frame as a raw element image and zeroes the buffer when its length
 633  // differs from the element size, which is always the case here.
 634  func (s *Server) pollIngestWorkers() {
 635  	var resp wire.IngestResponse
 636  	for i := int32(0); i < s.workers.Len(); i++ {
 637  		if !workerAlive(s.workerDone[i]) {
 638  			s.respawnIngestWorker(i)
 639  			continue
 640  		}
 641  		select {
 642  		case <-s.workerReady[i]:
 643  			resp = <-s.workerOut[i]
 644  			s.workers.Busy[i] = false
 645  			s.completeIngestResponse(resp)
 646  		default:
 647  		}
 648  	}
 649  }
 650  
 651  // workerAlive reports whether a spawned worker domain is still running. The
 652  // lifecycle channel closes when the child exits.
 653  func workerAlive(done chan struct{}) (ok bool) {
 654  	select {
 655  	case <-done:
 656  		return false
 657  	default:
 658  	}
 659  	return true
 660  }
 661  
 662  func (s *Server) completeIngestResponse(resp wire.IngestResponse) {
 663  	connFD, ok := s.pendingReq[resp.ReqID]
 664  	if !ok {
 665  		return
 666  	}
 667  	delete(s.pendingReq, resp.ReqID)
 668  	if s.conns[connFD] == nil {
 669  		return
 670  	}
 671  
 672  	switch resp.Verdict {
 673  	case wire.VerdictReject:
 674  		rej := &envelope.OK{EventID: resp.EventID[:], OK: false, Reason: resp.Reason}
 675  		s.t.SendWS(connFD, rej.Marshal(nil))
 676  		return
 677  	}
 678  
 679  	// The signature gate accepted (or accepted as ephemeral) the event; the
 680  	// database node runs stage B and decides whether it is stored, and echoes
 681  	// the frame back for broadcast.
 682  	if !s.dbSend(tree.Request{Op: tree.OpPersist, ConnID: connFD, Verified: true, Bytes: resp.Bytes}) {
 683  		rej := &envelope.OK{EventID: resp.EventID[:], OK: false, Reason: []byte("error: storage unavailable")}
 684  		s.t.SendWS(connFD, rej.Marshal(nil))
 685  	}
 686  	resp.Bytes = nil
 687  	resp.Reason = nil
 688  }
 689  
 690  
 691  func (s *Server) eventReject(fd int32, msg []byte, reason string) {
 692  	_, rem, _ := envelope.Identify(msg)
 693  	ev := event.New()
 694  	ev.Unmarshal(rem)
 695  	ok := &envelope.OK{EventID: ev.ID, OK: false, Reason: []byte(reason)}
 696  	s.t.SendWS(fd, ok.Marshal(nil))
 697  }
 698  
 699  func peekEventKind(msg []byte) (n uint16) {
 700  	_, rem, _ := envelope.Identify(msg)
 701  	ev := event.New()
 702  	ev.Unmarshal(rem)
 703  	return ev.Kind
 704  }
 705  
 706  // --- broadcast domain ---
 707  
 708  func (s *Server) handleBroadcastFrame() {
 709  	for i := 0; i < 16; i++ {
 710  		connFD, msg, ok := s.bcast.ReadFrame()
 711  		if !ok {
 712  			if !s.bcast.Valid() {
 713  				s.respawnBroadcastWorker()
 714  			}
 715  			return
 716  		}
 717  		// Buffered write: a full socket buffer is backpressure, not a reason to
 718  		// drop the subscriber.
 719  		s.t.SendWS(int32(connFD), msg)
 720  		msg = nil
 721  	}
 722  }
 723  
 724  func (s *Server) sendBroadcast(rawMsg []byte, senderFD int32) {
 725  	var flags uint8
 726  	if s.cfg.RelayURL != "" && !s.cfg.PrivilegedOpen {
 727  		flags |= 1
 728  	}
 729  	if s.cfg.NIP70Enforce {
 730  		flags |= 2
 731  	}
 732  	if s.cfg.MarmotOpen {
 733  		flags |= 4
 734  	}
 735  	s.bcast.Fanout(rawMsg, int32(senderFD), flags)
 736  }
 737  
 738  // --- HTTP routing ---
 739  
 740  func (s *Server) routeHTTP(fd int32, method, path string, headers map[string]string, body []byte) (rc int32, hdrs map[string]string, out []byte) {
 741  	if path == "/health" {
 742  		return 200, map[string]string{"Content-Type": "text/plain"}, []byte("ok")
 743  	}
 744  	if path == "/metrics" {
 745  		return 200, map[string]string{
 746  			"Content-Type":                "application/json",
 747  			"Access-Control-Allow-Origin": "*",
 748  		}, metrics.SnapshotAll(nil)
 749  	}
 750  	if path == "/metrics/reset" {
 751  		metrics.ResetAll()
 752  		return 200, map[string]string{"Content-Type": "text/plain"}, []byte("reset\n")
 753  	}
 754  	if headers["accept"] == "application/nostr+json" {
 755  		return 200, map[string]string{
 756  			"Content-Type":                "application/nostr+json",
 757  			"Access-Control-Allow-Origin": "*",
 758  		}, s.nip11JSON()
 759  	}
 760  	if s.Fallback != nil {
 761  		status, h, b := s.Fallback(method, path, headers, body)
 762  		if s.cfg.CORSEnabled {
 763  			for k, v := range s.corsHeaders(headers) {
 764  				h[k] = v
 765  			}
 766  		}
 767  		return status, h, b
 768  	}
 769  	return 404, map[string]string{"Content-Type": "text/plain"}, []byte("404 page not found\n")
 770  }
 771  
 772  func (s *Server) nip11JSON() (buf []byte) {
 773  	nips := "[1,9,11,40,42,45,50"
 774  	if s.cfg.NIP70Enforce {
 775  		nips = nips | ",70"
 776  	}
 777  	if s.cfg.NegentropyEnabled {
 778  		nips = nips | ",77"
 779  	}
 780  	nips = nips | "]"
 781  
 782  	pubkey := ""
 783  	if len(s.cfg.Admins) > 0 {
 784  		pubkey = s.cfg.Admins[0]
 785  	}
 786  
 787  	relayURL := s.cfg.RelayURL
 788  	maxSubs := s.cfg.MaxSubscriptions
 789  	if maxSubs <= 0 {
 790  		maxSubs = 10000
 791  	}
 792  	authRequired := s.cfg.AuthRequired || s.cfg.AuthToWrite
 793  
 794  	buf := []byte{:0:2048}
 795  	buf = transport.AppendStr(buf, "{\"name\":")
 796  	buf = transport.AppendJSONString(buf, s.cfg.AppName)
 797  	buf = transport.AppendStr(buf, ",\"description\":")
 798  	buf = transport.AppendJSONString(buf, "musiquay - Moxie-implemented Nostr relay")
 799  	if pubkey != "" {
 800  		buf = transport.AppendStr(buf, ",\"pubkey\":")
 801  		buf = transport.AppendJSONString(buf, pubkey)
 802  	}
 803  	if relayURL != "" {
 804  		buf = transport.AppendStr(buf, ",\"contact\":")
 805  		buf = transport.AppendJSONString(buf, relayURL)
 806  	}
 807  	buf = transport.AppendStr(buf, ",\"software\":\"https://github.com/mleku/musiquay\"")
 808  	buf = transport.AppendStr(buf, ",\"version\":")
 809  	buf = transport.AppendJSONString(buf, s.Version)
 810  	buf = transport.AppendStr(buf, ",\"supported_nips\":")
 811  	buf = buf | nips
 812  	buf = transport.AppendStr(buf, ",\"limitation\":{\"max_message_length\":524288,\"max_subscriptions\":")
 813  	buf = transport.AppendInt(buf, maxSubs)
 814  	buf = transport.AppendStr(buf, ",\"max_filters\":16,\"max_limit\":")
 815  	buf = transport.AppendInt(buf, s.cfg.QueryResultLimit)
 816  	buf = transport.AppendStr(buf, ",\"max_subid_length\":256,\"max_event_tags\":2000,\"max_content_length\":262144,\"min_pow_difficulty\":0,\"auth_required\":")
 817  	if authRequired {
 818  		buf = transport.AppendStr(buf, "true")
 819  	} else {
 820  		buf = transport.AppendStr(buf, "false")
 821  	}
 822  	buf = transport.AppendStr(buf, ",\"payment_required\":false,\"restricted_writes\":")
 823  	if s.cfg.AuthToWrite || s.cfg.ACLMode != "none" {
 824  		buf = transport.AppendStr(buf, "true")
 825  	} else {
 826  		buf = transport.AppendStr(buf, "false")
 827  	}
 828  	buf = transport.AppendStr(buf, "}}")
 829  	return buf
 830  }
 831  
 832  func (s *Server) corsHeaders(reqHeaders map[string]string) (m map[string]string) {
 833  	origin := reqHeaders["origin"]
 834  	if origin == "" {
 835  		origin = "*"
 836  	}
 837  	allowed := false
 838  	if len(s.cfg.CORSOrigins) == 0 {
 839  		allowed = true
 840  	} else {
 841  		for _, o := range s.cfg.CORSOrigins {
 842  			if o == origin || o == "*" {
 843  				allowed = true
 844  				break
 845  			}
 846  		}
 847  	}
 848  	if !allowed {
 849  		return map[string]string{}
 850  	}
 851  	return map[string]string{
 852  		"Access-Control-Allow-Origin":  origin,
 853  		"Access-Control-Allow-Methods": "GET, POST, OPTIONS",
 854  		"Access-Control-Allow-Headers": "Content-Type, Authorization",
 855  		"Access-Control-Max-Age":       "86400",
 856  	}
 857  }
 858  
 859  // --- IP access control ---
 860  
 861  func (s *Server) ipBlacklisted(ip string) (ok bool) {
 862  	for _, prefix := range s.cfg.IPBlacklist {
 863  		if transport.HasPrefix(ip, prefix) {
 864  			return true
 865  		}
 866  	}
 867  	return false
 868  }
 869  
 870  func (s *Server) ipWhitelisted(ip string) (ok bool) {
 871  	if len(s.cfg.IPWhitelist) == 0 {
 872  		return false
 873  	}
 874  	for _, prefix := range s.cfg.IPWhitelist {
 875  		if transport.HasPrefix(ip, prefix) {
 876  			return true
 877  		}
 878  	}
 879  	return false
 880  }
 881  
 882  // --- utility (used locally, not in transport) ---
 883  
 884  func makeCopy(b []byte) (buf []byte) {
 885  	c := []byte{:len(b)}
 886  	copy(c, b)
 887  	return c
 888  }
 889  
 890  func toLower(b []byte) (buf []byte) {
 891  	for i := range b {
 892  		if b[i] >= 'A' && b[i] <= 'Z' {
 893  			b[i] = b[i] + 32
 894  		}
 895  	}
 896  	return b
 897  }
 898