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