// Package server provides the Nostr relay domain coordinator. // It implements transport.Handler and wires together the event store, // worker pool, broadcast domain, and HTTP/WebSocket routing. package server import ( "fmt" "runtime" "net/url" "time" "git.smesh.lol/morly/pkg/access" "git.smesh.lol/morly/pkg/broadcast" "git.smesh.lol/morly/pkg/metrics" "git.smesh.lol/nostr/pkg/envelope" "git.smesh.lol/nostr/pkg/event" "git.smesh.lol/nostr/pkg/filter" "git.smesh.lol/morly/pkg/pool" "git.smesh.lol/morly/pkg/relay/config" "git.smesh.lol/morly/pkg/relay/dbengine" "git.smesh.lol/morly/pkg/relay/ratelimit" "git.smesh.lol/morly/pkg/relay/tree" "git.smesh.lol/morly/pkg/relay/wire" "git.smesh.lol/morly/pkg/transport" ) // InitSignals delegates to transport.InitSignals. func InitSignals() { transport.InitSignals() } // DB is root's handle to the database-engine domain. Root never opens the // store: it sends requests and learns of replies through the ready nudge, // because a select receive on a codec-framed spawn channel does not decode. type DB struct { in chan tree.Request out chan tree.Response ready chan struct{} done chan struct{} } // Valid reports whether the database-engine domain is still running. func (d *DB) Valid() (ok bool) { select { case <-d.done: return false default: } return true } // SpawnDB forks the database-engine domain and returns root's handle to it. // Channel construction and the fork both happen under the root arena: the // channels outlive this call and the forked domain reads the same memory. func SpawnDB(cfg *config.C) (d *DB) { dbCfg := dbengine.Config{ DataDir: cfg.DataDir, ACLMode: cfg.ACLMode, Admins: cfg.Admins, FollowListFreqSec: cfg.FollowListFreqSec, SocialWoTMaxDepth: cfg.SocialWoTMaxDepth, SocialWoTRefreshSec: cfg.SocialWoTRefreshSec, MuteBlacklist: cfg.MuteBlacklist, } runtime.SovereignSetArena(runtime.RootArena()) in := chan tree.Request{} out := chan tree.Response{} ready := chan struct{}{:64} done := spawn(dbengine.Run, dbCfg, in, out, ready) runtime.SovereignRestoreArena(runtime.RootArena()) return &DB{in: in, out: out, ready: ready, done: done} } // Server is the Nostr relay domain coordinator. // It implements transport.Handler so the transport layer can call back // for connection lifecycle events and incoming messages. type Server struct { Fallback func(method, path string, headers map[string]string, body []byte) (int32, map[string]string, []byte) OnReady func() Version string selfHost string // "host/" prefix for proxy self-redirect detection t *transport.Server cfg *config.C // tickCount counts OnTick calls for the periodic cleanup; it lives on the // server rather than in a package global, which the rule forbids writing // outside init. tickCount int32 // per-connection domain state, indexed by FD (parallel to transport.conns) conns map[int32]*cstate writeLimiter *ratelimit.Limiter // closed makes Close idempotent: both shutdown paths call it. closed bool // Database engine handle; all storage goes through it. db *DB // In-flight database requests, keyed by reqID. dbPending map[uint32]int32 nextDBReq uint32 // Ingest worker pool (signature gate). workerIn []chan wire.IngestRequest workerOut []chan wire.IngestResponse workerDone []chan struct{} workerReady []chan struct{} workers pool.Pool pendingReq map[uint32]int32 // reqID → conn fd nextReqID uint32 // Broadcast domain. bcast *broadcast.Broadcaster // Media proxy worker pool. proxyIn []chan wire.ProxyRequest proxyOut []chan wire.ProxyResponse proxyDone []chan struct{} proxyPool pool.Pool proxyQueue []pendingProxy proxyBusyTime []int64 // unix nanos when worker went busy; 0 = idle // Blossom worker pool. blossomIn []chan wire.BlossomRequest blossomOut []chan wire.BlossomResponse blossomDone []chan struct{} blossomPool pool.Pool // Shared async HTTP pending. asyncPending map[uint32]asyncHTTPEntry nextAsyncID uint32 } type asyncHTTPEntry struct { connFD int32 connClose bool createdAt int64 // unix nanos, for reap timeout } type pendingProxy struct { connFD int32 connClose bool req wire.ProxyRequest } // cstate is the server-domain per-connection state (server-side complement to // transport's tconn, keyed by the same FD). type cstate struct { subs map[string]*sub challenge []byte authedPubkey []byte jailed bool } type sub struct { id string filters filter.S rawReq []byte } // New creates a relay server backed by the given database-engine handle. func New(db *DB, cfg *config.C) (s *Server) { var wlim *ratelimit.Limiter if cfg.FreeWriteLimit > 0 && cfg.FreeWriteWindow > 0 { rate := float64(cfg.FreeWriteLimit) / float64(cfg.FreeWriteWindow) wlim = ratelimit.New(rate, cfg.FreeWriteLimit) } selfHost := "git.smesh.lol/morly/" if u, err := url.Parse(cfg.RelayURL); err == nil && len(u.Host) > 0 { selfHost = u.Host | "/" } s = &Server{ cfg: cfg, selfHost: selfHost, conns: map[int32]*cstate{}, writeLimiter: wlim, db: db, dbPending: map[uint32]int32{}, asyncPending: map[uint32]asyncHTTPEntry{}, } // Create transport. Server implements transport.Handler. s.t = transport.New(s, cfg.MaxConnPerIP) s.t.BotBlock = cfg.HTTPGuardBotBlock // Spawn worker domains AFTER heap is settled (GC avoidance). if cfg.IngestWorkers > 0 { s.startIngestWorkers(cfg.IngestWorkers) } s.startBroadcastWorker() if n := cfg.MediaProxyWorkers; n > 0 { s.startProxyWorkers(n) } if n := cfg.BlossomWorkers; n > 0 { s.startBlossomWorkers(n) } return s } // Close stops the database-engine domain, which flushes and closes the store. func (s *Server) Close() { // Close runs on both shutdown paths (main calls it directly and again from // a defer), and closing a channel twice panics, so the first call wins. if s.closed { return } s.closed = true if s.db != nil && s.db.Valid() { close(s.db.in) } // Every worker loop returns when its request channel closes, and a domain // that returns dumps its coverage counters on the way out. Without this the // pool domains stayed blocked in a futex receive while the fixture's SIGTERM // went to the process group, so they were SIGKILLed with their counters // (private after fork) and every worker file read 0% coverage. for _, in := range s.workerIn { close(in) } for _, in := range s.proxyIn { close(in) } for _, in := range s.blossomIn { close(in) } // The broadcast worker owns the subscription registry and is pre-spawned // by broadcast.PreSpawn before the store; close it too, or it is the one // domain still parked in a receive when the process group is signalled. s.bcast.Close() } func (s *Server) startIngestWorkers(n int32) { s.workerIn, s.workerOut, s.workerDone, s.workerReady = spawnIngestPool(n) s.workers = pool.NewPool(n) s.pendingReq = map[uint32]int32{} } // spawnIngestPool forks the ingest domains in a free function; see // spawnProxyPool for why the loop does not live in the Server method. // // The pool slices are built with push, not []chan T{:n}: an inline // slice-of-channel literal is miscompiled to an empty slice (the literal // rewrite does not resolve the channel element type), which is why the named // slice aliases existed. push allocates the same []chan T correctly. func spawnIngestPool(n int32) (ins []chan wire.IngestRequest, outs []chan wire.IngestResponse, dones []chan struct{}, readys []chan struct{}) { var in chan wire.IngestRequest var out chan wire.IngestResponse var done chan struct{} var ready chan struct{} for i := 0; i < n; i++ { in, out, done, ready = newIngestWorker() ins = push(ins, in) outs = push(outs, out) dones = push(dones, done) readys = push(readys, ready) } return } // newIngestWorker creates the channel pair and forks one ingest domain. // A free function: the channel/spawn scratch belongs to this frame, not to // the caller's sovereign arena. func newIngestWorker() (in chan wire.IngestRequest, out chan wire.IngestResponse, done chan struct{}, ready chan struct{}) { // The channel structs outlive this frame: the server keeps them for the // whole run and the forked domain reads the same memory. chanMake // allocates in the current arena, so a plain frame here would leave every // stored channel pointing into a released arena. runtime.SovereignSetArena(runtime.RootArena()) in = chan wire.IngestRequest{} out = chan wire.IngestResponse{} // Ready carries no payload; it exists because a select receive on a // codec-framed spawn channel does not decode. The worker nudges here // after writing out, and the parent reads out with a plain receive. ready = chan struct{}{:64} done = spawn(wire.IngestWorker, in, out, ready) runtime.SovereignRestoreArena(runtime.RootArena()) return } func (s *Server) startBroadcastWorker() { s.bcast = broadcast.New() } func (s *Server) respawnIngestWorker(i int32) { in, out, done, ready := newIngestWorker() s.workerIn[i] = in s.workerOut[i] = out s.workerDone[i] = done s.workerReady[i] = ready s.workers.Busy[i] = false fmt.Println("respawned ingest worker", i) } func (s *Server) respawnBroadcastWorker() { s.bcast = broadcast.New() for fd, cs := range s.conns { s.bcast.OnConnect(int32(fd), s.t.ConnIsWhitelisted(fd)) if len(cs.authedPubkey) > 0 { s.bcast.OnAuth(int32(fd), cs.authedPubkey) } for _, sb := range cs.subs { s.bcast.OnSubscribe(int32(fd), []byte(sb.id), sb.rawReq) } } fmt.Println("respawned broadcast worker, re-registered", len(s.conns), "connections") } // wireOnReady hands the transport the startup callback. It is its own mutating // method so that ListenAndServe itself stays a read-only method of a sovereign // type: a mutating method borrows the receiver's sovereign arena as its working // frame for its whole body, and ListenAndServe's body is the transport event // loop, so a write there would keep this Server's arena borrowed across every // turn of the loop. The loop's compaction yield point found the group borrowed // every single time and could only rotate it back onto the queue: the group // could never be compacted, which is what made the relay's arena grow without // bound. Peeling the write out leaves the loop running with the arena merely // registered as stable, which the compaction rebuild pass knows how to repoint. func (s *Server) wireOnReady() { s.t.OnReady = s.OnReady } // ListenAndServe starts the transport event loop. func (s *Server) ListenAndServe(addr string) (err error) { s.wireOnReady() return s.t.ListenAndServe(addr) } // --- transport.Handler implementation --- // OnAccept gates new TCP connections: IP blacklist and global connection limit. func (s *Server) OnAccept(fd int32, ip string) (ok bool) { if s.cfg.MaxGlobalConns > 0 && s.t.ConnCount() >= s.cfg.MaxGlobalConns { return false } return !s.ipBlacklisted(ip) } // OnWSUpgrade gates WS upgrades: per-IP connection limit. func (s *Server) OnWSUpgrade(fd int32, ip string, currentIPWSCount int32) (whitelisted bool, allow bool) { wl := s.ipWhitelisted(ip) if !wl && s.cfg.MaxConnPerIP > 0 && currentIPWSCount >= s.cfg.MaxConnPerIP { return false, false } return wl, true } // OnWSConnected initialises per-conn state and sends the NIP-42 challenge. func (s *Server) OnWSConnected(fd int32) { // The per-connection state outlives this call: the server keeps it in // s.conns for the life of the connection. The record, the subscription // map, the NIP-42 challenge bytes and the boxed challenge message must all // be allocated in the root arena - the challenge is built here, so the // borrow covers it too. var ac *envelope.AuthChallenge var ch []byte prev := runtime.CurrentArena() runtime.SovereignSetArena(runtime.RootArena()) s.conns[fd] = &cstate{subs: map[string]*sub{}} if s.cfg.RelayURL != "" { ch = transport.Challenge() s.conns[fd].challenge = ch ac = &envelope.AuthChallenge{Challenge: ch} } runtime.SovereignRestoreArena(prev) s.bcast.OnConnect(int32(fd), s.t.ConnIsWhitelisted(fd)) if ac != nil { s.t.SendWS(fd, ac.Marshal(nil)) } } // OnWSMessage dispatches a decoded WebSocket payload. func (s *Server) OnWSMessage(fd int32, payload []byte) { s.dispatch(fd, payload) } // OnWSClose cleans up per-conn state when a WS connection closes. func (s *Server) OnWSClose(fd int32) { c := s.conns[fd] if c != nil { for k, sb := range c.subs { sb.rawReq = nil sb.filters.F = nil delete(c.subs, k) } c.subs = nil c.challenge = nil c.authedPubkey = nil } delete(s.conns, fd) s.bcast.OnDisconnect(int32(fd)) for rid, connFD := range s.pendingReq { if connFD == fd { delete(s.pendingReq, rid) } } for rid, connFD := range s.dbPending { if connFD == fd { delete(s.dbPending, rid) } } } // OnHTTP routes HTTP requests to workers or synchronous handlers. // Returns transport.HTTPDeferred to signal async processing. 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) { // IP blacklist re-check after XFF substitution. if s.ipBlacklisted(s.t.ConnIP(fd)) { return 403, nil, nil, true } if transport.HasPrefix(path, "/proxy/") && s.proxyPool.Len() > 0 { if s.doDispatchProxy(fd, path, headers) { return transport.HTTPDeferred, nil, nil, false } } if transport.HasPrefix(path, "/blossom") && s.blossomPool.Len() > 0 { if s.doDispatchBlossom(fd, method, path, headers, body) { return transport.HTTPDeferred, nil, nil, false } } connClose := headers["connection"] == "close" if method == "OPTIONS" && s.cfg.CORSEnabled { return 204, s.corsHeaders(headers), nil, connClose } status, h, b := s.routeHTTP(fd, method, path, headers, body) return status, h, b, connClose } // OnPoll drains worker-domain output channels. The transport loop calls it // once per iteration; spawn channels have no fd to wake epoll. func (s *Server) OnPoll() { pollStart := metrics.Now() s.pollDB() s.pollIngestWorkers() s.pollProxyWorkers() s.pollBlossomWorkers() s.handleBroadcastFrame() metrics.OnPollNs.Observe(metrics.Since(pollStart)) } func (s *Server) OnTick() { s.proxyReapStuck() s.asyncReapStuck() s.tickCount++ if s.tickCount%60 == 0 { if s.writeLimiter != nil { s.writeLimiter.Cleanup(5 * time.Minute) } } } // --- Nostr message dispatch --- func (s *Server) dispatch(fd int32, msg []byte) { c := s.conns[fd] if c != nil && c.jailed { return } label, _, _ := envelope.Identify(msg) switch label { case envelope.EventLabel: s.handleEvent(fd, msg) case envelope.ReqLabel: s.handleReq(fd, msg) case envelope.CloseLabel: s.handleClose(fd, msg) case envelope.CountLabel: s.handleCount(fd, msg) case envelope.AuthLabel: s.handleAuth(fd, msg) } } // handleEvent runs one EVENT frame. // // It stays a method on *Server. Moving the body into a free function that // returns the frame to send - which would give it a frame arena and free its // scratch at return - was tried and reverted: with the free function the // free-write limiter stopped applying (the fourth unauthed write was accepted // where it has to be rate-limited). The limiter itself is not at fault: a // counter with the same shape keeps its state when driven from a free function // in isolation, so the cause is still unidentified. Keep the body in the // method until it is. func (s *Server) handleEvent(fd int32, msg []byte) { handleStart := metrics.Now() defer func() { metrics.HandleEventNs.Observe(metrics.Since(handleStart)) }() c := s.conns[fd] if c == nil { return } cfg := s.cfg authed := len(c.authedPubkey) > 0 gated := cfg.RelayURL != "" wl := s.t.ConnIsWhitelisted(fd) if gated && !authed && !wl { evKind := peekEventKind(msg) exempt := access.WriteExempt(evKind, s.cfg.NIP46BypassAuth, s.cfg.MarmotOpen) if !exempt && s.writeLimiter != nil { ip := s.t.ConnIP(fd) if s.writeLimiter.Allow([]byte(ip)) { exempt = true } else { s.eventReject(fd, msg, "rate-limited: too many events, try again later") return } } if !exempt && cfg.AuthToWrite { s.eventReject(fd, msg, "auth-required: authentication required") return } } if s.workers.Len() > 0 && s.dispatchToWorker(fd, msg) { return } parseStart := metrics.Now() _, rem, _ := envelope.Identify(msg) var es envelope.EventSubmission perr := error(nil) _, perr = es.Unmarshal(rem) metrics.EnvelopeParseNs.Observe(metrics.Since(parseStart)) if perr != nil || es.E == nil { s.t.SendWS(fd, (&envelope.OK{EventID: []byte{:32}, Reason: []byte("invalid: parse error")}).Marshal(nil)) return } if !s.dbSend(tree.Request{Op: tree.OpPersist, ConnID: fd, Bytes: msg}) { s.eventReject(fd, msg, "error: storage unavailable") } } // dbSend hands one request to the database-engine domain and records the // connection that is waiting for the reply. Returns false when the domain is // gone, so the caller can answer instead of leaving the client hanging. func (s *Server) dbSend(req tree.Request) (ok bool) { if s.db == nil || !s.db.Valid() { return false } s.nextDBReq++ req.ReqID = s.nextDBReq s.dbPending[req.ReqID] = req.ConnID s.db.in <- req return true } // pollDB drains ready database-engine replies. The worker nudges root on a // zero-size channel after each reply; the reply itself travels on the codec // channel and is read with a plain receive, which decodes. func (s *Server) pollDB() { for i := 0; i < 256; i++ { select { case <-s.db.ready: resp := <-s.db.out s.completeDB(resp) default: return } } } func (s *Server) completeDB(resp tree.Response) { switch resp.Op { case tree.OpPersist: s.completePersist(resp) case tree.OpHistory: s.completeHistory(resp) case tree.OpCount: s.completeCount(resp) } } func (s *Server) completePersist(resp tree.Response) { fd, ok := s.dbPending[resp.ReqID] if !ok { return } delete(s.dbPending, resp.ReqID) if s.conns[fd] == nil { return } out := &envelope.OK{EventID: resp.EventID, OK: resp.OK, Reason: resp.Reason} sendStart := metrics.Now() s.t.SendWS(fd, out.Marshal(nil)) metrics.SendWSNs.Observe(metrics.Since(sendStart)) if resp.OK && len(resp.Bytes) > 0 { bStart := metrics.Now() s.sendBroadcast(resp.Bytes, fd) metrics.BroadcastNs.Observe(metrics.Since(bStart)) } } // completeHistory streams the stored events for one subscription and closes it // with EOSE. A subscription closed while the query was in flight is dropped. func (s *Server) completeHistory(resp tree.Response) { c := s.conns[resp.ConnID] if c == nil { return } if _, still := c.subs[string(resp.SubID)]; !still { return } // SendWS buffers when the socket is momentarily full and flushes on a later // loop pass. SendWSErr's direct write reports EAGAIN as an error, and the // old "close the connection on error" turned that transient state into a // dropped subscription: the browser's feed REQ - the only sub with an event // to deliver - lost both the EVENT and the EOSE, while the tag-less subs // answered with EOSE alone. for _, frame := range resp.Events { s.t.SendWS(resp.ConnID, frame) } s.t.SendWS(resp.ConnID, (&envelope.EOSE{Subscription: resp.SubID}).Marshal(nil)) } func (s *Server) completeCount(resp tree.Response) { if s.conns[resp.ConnID] == nil { return } cr := &envelope.CountResponse{Subscription: resp.SubID, Count: resp.Count} s.t.SendWS(resp.ConnID, cr.Marshal(nil)) } func (s *Server) dispatchToWorker(fd int32, msg []byte) (ok bool) { i := s.workers.IdleIndex() if i < 0 { return false } s.nextReqID++ rid := s.nextReqID s.pendingReq[rid] = fd s.workers.Busy[i] = true s.workerIn[i] <- wire.IngestRequest{ReqID: rid, Bytes: msg} return true } // pollIngestWorkers drains ready ingest responses. Called from the transport // loop each iteration; worker domains are reached over spawn channels, which // have no fd to wake epoll. // // Readiness arrives on the zero-size nudge channel, then the response is read // from the codec channel with a plain receive. A select receive directly on a // codec-framed spawn channel does not decode: the runtime treats the encoded // ring frame as a raw element image and zeroes the buffer when its length // differs from the element size, which is always the case here. func (s *Server) pollIngestWorkers() { var resp wire.IngestResponse for i := int32(0); i < s.workers.Len(); i++ { if !workerAlive(s.workerDone[i]) { s.respawnIngestWorker(i) continue } select { case <-s.workerReady[i]: resp = <-s.workerOut[i] s.workers.Busy[i] = false s.completeIngestResponse(resp) default: } } } // workerAlive reports whether a spawned worker domain is still running. The // lifecycle channel closes when the child exits. func workerAlive(done chan struct{}) (ok bool) { select { case <-done: return false default: } return true } func (s *Server) completeIngestResponse(resp wire.IngestResponse) { connFD, ok := s.pendingReq[resp.ReqID] if !ok { return } delete(s.pendingReq, resp.ReqID) if s.conns[connFD] == nil { return } switch resp.Verdict { case wire.VerdictReject: rej := &envelope.OK{EventID: resp.EventID[:], OK: false, Reason: resp.Reason} s.t.SendWS(connFD, rej.Marshal(nil)) return } // The signature gate accepted (or accepted as ephemeral) the event; the // database node runs stage B and decides whether it is stored, and echoes // the frame back for broadcast. if !s.dbSend(tree.Request{Op: tree.OpPersist, ConnID: connFD, Verified: true, Bytes: resp.Bytes}) { rej := &envelope.OK{EventID: resp.EventID[:], OK: false, Reason: []byte("error: storage unavailable")} s.t.SendWS(connFD, rej.Marshal(nil)) } resp.Bytes = nil resp.Reason = nil } func (s *Server) eventReject(fd int32, msg []byte, reason string) { _, rem, _ := envelope.Identify(msg) ev := event.New() ev.Unmarshal(rem) ok := &envelope.OK{EventID: ev.ID, OK: false, Reason: []byte(reason)} s.t.SendWS(fd, ok.Marshal(nil)) } func peekEventKind(msg []byte) (n uint16) { _, rem, _ := envelope.Identify(msg) ev := event.New() ev.Unmarshal(rem) return ev.Kind } // --- broadcast domain --- func (s *Server) handleBroadcastFrame() { for i := 0; i < 16; i++ { connFD, msg, ok := s.bcast.ReadFrame() if !ok { if !s.bcast.Valid() { s.respawnBroadcastWorker() } return } // Buffered write: a full socket buffer is backpressure, not a reason to // drop the subscriber. s.t.SendWS(int32(connFD), msg) msg = nil } } func (s *Server) sendBroadcast(rawMsg []byte, senderFD int32) { var flags uint8 if s.cfg.RelayURL != "" && !s.cfg.PrivilegedOpen { flags |= 1 } if s.cfg.NIP70Enforce { flags |= 2 } if s.cfg.MarmotOpen { flags |= 4 } s.bcast.Fanout(rawMsg, int32(senderFD), flags) } // --- HTTP routing --- func (s *Server) routeHTTP(fd int32, method, path string, headers map[string]string, body []byte) (rc int32, hdrs map[string]string, out []byte) { if path == "/health" { return 200, map[string]string{"Content-Type": "text/plain"}, []byte("ok") } if path == "/metrics" { return 200, map[string]string{ "Content-Type": "application/json", "Access-Control-Allow-Origin": "*", }, metrics.SnapshotAll(nil) } if path == "/metrics/reset" { metrics.ResetAll() return 200, map[string]string{"Content-Type": "text/plain"}, []byte("reset\n") } if headers["accept"] == "application/nostr+json" { return 200, map[string]string{ "Content-Type": "application/nostr+json", "Access-Control-Allow-Origin": "*", }, s.nip11JSON() } if s.Fallback != nil { status, h, b := s.Fallback(method, path, headers, body) if s.cfg.CORSEnabled { for k, v := range s.corsHeaders(headers) { h[k] = v } } return status, h, b } return 404, map[string]string{"Content-Type": "text/plain"}, []byte("404 page not found\n") } func (s *Server) nip11JSON() (buf []byte) { nips := "[1,9,11,40,42,45,50" if s.cfg.NIP70Enforce { nips = nips | ",70" } if s.cfg.NegentropyEnabled { nips = nips | ",77" } nips = nips | "]" pubkey := "" if len(s.cfg.Admins) > 0 { pubkey = s.cfg.Admins[0] } relayURL := s.cfg.RelayURL maxSubs := s.cfg.MaxSubscriptions if maxSubs <= 0 { maxSubs = 10000 } authRequired := s.cfg.AuthRequired || s.cfg.AuthToWrite buf := []byte{:0:2048} buf = transport.AppendStr(buf, "{\"name\":") buf = transport.AppendJSONString(buf, s.cfg.AppName) buf = transport.AppendStr(buf, ",\"description\":") buf = transport.AppendJSONString(buf, "musiquay - Moxie-implemented Nostr relay") if pubkey != "" { buf = transport.AppendStr(buf, ",\"pubkey\":") buf = transport.AppendJSONString(buf, pubkey) } if relayURL != "" { buf = transport.AppendStr(buf, ",\"contact\":") buf = transport.AppendJSONString(buf, relayURL) } buf = transport.AppendStr(buf, ",\"software\":\"https://github.com/mleku/musiquay\"") buf = transport.AppendStr(buf, ",\"version\":") buf = transport.AppendJSONString(buf, s.Version) buf = transport.AppendStr(buf, ",\"supported_nips\":") buf = buf | nips buf = transport.AppendStr(buf, ",\"limitation\":{\"max_message_length\":524288,\"max_subscriptions\":") buf = transport.AppendInt(buf, maxSubs) buf = transport.AppendStr(buf, ",\"max_filters\":16,\"max_limit\":") buf = transport.AppendInt(buf, s.cfg.QueryResultLimit) buf = transport.AppendStr(buf, ",\"max_subid_length\":256,\"max_event_tags\":2000,\"max_content_length\":262144,\"min_pow_difficulty\":0,\"auth_required\":") if authRequired { buf = transport.AppendStr(buf, "true") } else { buf = transport.AppendStr(buf, "false") } buf = transport.AppendStr(buf, ",\"payment_required\":false,\"restricted_writes\":") if s.cfg.AuthToWrite || s.cfg.ACLMode != "none" { buf = transport.AppendStr(buf, "true") } else { buf = transport.AppendStr(buf, "false") } buf = transport.AppendStr(buf, "}}") return buf } func (s *Server) corsHeaders(reqHeaders map[string]string) (m map[string]string) { origin := reqHeaders["origin"] if origin == "" { origin = "*" } allowed := false if len(s.cfg.CORSOrigins) == 0 { allowed = true } else { for _, o := range s.cfg.CORSOrigins { if o == origin || o == "*" { allowed = true break } } } if !allowed { return map[string]string{} } return map[string]string{ "Access-Control-Allow-Origin": origin, "Access-Control-Allow-Methods": "GET, POST, OPTIONS", "Access-Control-Allow-Headers": "Content-Type, Authorization", "Access-Control-Max-Age": "86400", } } // --- IP access control --- func (s *Server) ipBlacklisted(ip string) (ok bool) { for _, prefix := range s.cfg.IPBlacklist { if transport.HasPrefix(ip, prefix) { return true } } return false } func (s *Server) ipWhitelisted(ip string) (ok bool) { if len(s.cfg.IPWhitelist) == 0 { return false } for _, prefix := range s.cfg.IPWhitelist { if transport.HasPrefix(ip, prefix) { return true } } return false } // --- utility (used locally, not in transport) --- func makeCopy(b []byte) (buf []byte) { c := []byte{:len(b)} copy(c, b) return c } func toLower(b []byte) (buf []byte) { for i := range b { if b[i] >= 'A' && b[i] <= 'Z' { b[i] = b[i] + 32 } } return b }