package server import ( "bytes" "runtime" "time" "git.smesh.lol/nostr/pkg/envelope" "git.smesh.lol/nostr/pkg/filter" "git.smesh.lol/morly/pkg/relay/tree" ) // --- Subscription and Nostr message handling --- func (s *Server) handleReq(fd int32, msg []byte) { c := s.conns[fd] if c == nil { return } // Everything this subscription keeps - the id, the filters and the raw // request - outlives this call in c.subs, so parse and build it while // borrowing the root arena. Copying the filter slice header out of this // frame is not enough: the array it points at would stay behind, and the // receiver's sovereign compaction walk later follows that pointer into // memory the frame has already released. prev := runtime.CurrentArena() runtime.SovereignSetArena(runtime.RootArena()) _, rem, _ := envelope.Identify(msg) filter.ClearTaint() var req envelope.Req if _, err := req.Unmarshal(rem); err != nil { runtime.SovereignRestoreArena(prev) return } if filter.IsTainted() { runtime.SovereignRestoreArena(prev) c.jailed = true return } id := string(req.Subscription) cfg := s.cfg authed := len(c.authedPubkey) > 0 wl := s.t.ConnIsWhitelisted(fd) if cfg.AuthRequired && !authed && !wl { cl := &envelope.Closed{ Subscription: []byte(id), Reason: []byte("auth-required: authentication required"), } runtime.SovereignRestoreArena(prev) s.t.SendWS(fd, cl.Marshal(nil)) return } maxSubs := cfg.MaxSubscriptions if maxSubs > 0 { if _, exists := c.subs[id]; !exists && len(c.subs) >= maxSubs { cl := &envelope.Closed{ Subscription: []byte(id), Reason: []byte("error: too many subscriptions"), } runtime.SovereignRestoreArena(prev) s.t.SendWS(fd, cl.Marshal(nil)) return } } filters := req.Filters reqCopy := []byte{:len(msg)} copy(reqCopy, msg) c.subs[id] = &sub{id: id, filters: filters, rawReq: reqCopy} runtime.SovereignRestoreArena(prev) s.bcast.OnSubscribe(int32(fd), []byte(id), msg) needsFilter := cfg.RelayURL != "" && !cfg.PrivilegedOpen && !wl limit := cfg.QueryResultLimit if limit <= 0 { limit = 256 } // Backfill is a database query: root hands the whole REQ frame to the // database-engine node, which returns ready-to-send EVENT frames plus the // EOSE. Root stays free while the query runs. subID := req.Subscription var authed []byte if len(c.authedPubkey) > 0 { authed = []byte{:len(c.authedPubkey)} copy(authed, c.authedPubkey) } if !s.dbSend(tree.Request{ Op: tree.OpHistory, ConnID: fd, Limit: limit, Filtered: needsFilter, NIP70: cfg.NIP70Enforce, Marmot: cfg.MarmotOpen, SubID: subID, Filter: msg, AuthedPubkey: authed, }) { s.t.SendWS(fd, (&envelope.EOSE{Subscription: subID}).Marshal(nil)) } } func (s *Server) handleClose(fd int32, msg []byte) { c := s.conns[fd] if c == nil { return } _, rem, _ := envelope.Identify(msg) var cl envelope.Close if _, err := cl.Unmarshal(rem); err != nil { return } delete(c.subs, string(cl.ID)) s.bcast.OnUnsubscribe(int32(fd), cl.ID) } func (s *Server) handleCount(fd int32, msg []byte) { c := s.conns[fd] if c == nil { return } if s.cfg.AuthRequired && len(c.authedPubkey) == 0 && !s.t.ConnIsWhitelisted(fd) { return } _, rem, ierr := envelope.Identify(msg) if ierr != nil { return } filter.ClearTaint() var cr envelope.CountRequest if _, err := cr.Unmarshal(rem); err != nil { return } if filter.IsTainted() { c.jailed = true return } s.dbSend(tree.Request{Op: tree.OpCount, ConnID: fd, Filter: msg}) } func (s *Server) handleAuth(fd int32, msg []byte) { c := s.conns[fd] if c == nil { return } _, rem, _ := envelope.Identify(msg) var auth envelope.AuthResponse if _, aerr := auth.Unmarshal(rem); aerr != nil || auth.Event == nil { bad := &envelope.OK{ EventID: []byte{:32}, OK: false, Reason: []byte("error: failed to parse auth event"), } s.t.SendWS(fd, bad.Marshal(nil)) return } if auth.Event.Kind != 22242 { s.authFail(fd, auth.Event.ID, "error: wrong event kind for auth") return } ct := auth.Event.Tags.GetFirst([]byte("challenge")) if ct == nil || !bytes.Equal(ct.Value(), c.challenge) { s.authFail(fd, auth.Event.ID, "error: wrong challenge") return } rt := auth.Event.Tags.GetFirst([]byte("relay")) if rt == nil || len(rt.Value()) == 0 { s.authFail(fd, auth.Event.ID, "error: missing relay tag") return } if !relayURLMatch([]byte(s.cfg.RelayURL), rt.Value()) { s.authFail(fd, auth.Event.ID, "error: relay URL mismatch") return } now := time.Now().Unix() if auth.Event.CreatedAt > now+600 || auth.Event.CreatedAt < now-600 { s.authFail(fd, auth.Event.ID, "error: timestamp out of range") return } valid, err := auth.Event.Verify() if err != nil || !valid { s.authFail(fd, auth.Event.ID, "error: invalid signature") return } // The authed pubkey outlives this call - it stays on the connection state // for the life of the connection - so copy it into the root arena. authPrev := runtime.CurrentArena() runtime.SovereignSetArena(runtime.RootArena()) c.authedPubkey = []byte{:len(auth.Event.Pubkey)} copy(c.authedPubkey, auth.Event.Pubkey) runtime.SovereignRestoreArena(authPrev) s.bcast.OnAuth(int32(fd), c.authedPubkey) ok := &envelope.OK{EventID: auth.Event.ID, OK: true} s.t.SendWS(fd, ok.Marshal(nil)) } func (s *Server) authFail(fd int32, id []byte, reason string) { ok := &envelope.OK{EventID: id, OK: false, Reason: []byte(reason)} s.t.SendWS(fd, ok.Marshal(nil)) } func relayURLMatch(expected, found []byte) (ok bool) { if bytes.Equal(expected, found) { return true } e := makeCopy(expected) f := makeCopy(found) toLower(e) toLower(f) if len(e) > 0 && e[len(e)-1] == '/' { e = e[:len(e)-1] } if len(f) > 0 && f[len(f)-1] == '/' { f = f[:len(f)-1] } return bytes.Equal(stripScheme(e), stripScheme(f)) } func stripScheme(u []byte) (buf []byte) { if bytes.HasPrefix(u, []byte("wss://")) { return u[6:] } if bytes.HasPrefix(u, []byte("ws://")) { return u[5:] } return u }