broadcast_worker.mx raw

   1  package wire
   2  
   3  import (
   4  	"runtime"
   5  
   6  	"git.smesh.lol/morly/pkg/access"
   7  	"git.smesh.lol/nostr/pkg/envelope"
   8  	"git.smesh.lol/nostr/pkg/filter"
   9  )
  10  
  11  // bcastConnState holds the broadcast domain's view of one WS connection.
  12  type bcastConnState struct {
  13  	whitelisted  bool
  14  	authed       bool
  15  	authedPubkey [32]byte
  16  	subs         map[string]filter.S
  17  }
  18  
  19  // BroadcastWorker is the spawn target for the broadcast domain. Owns a
  20  // registry of connection subscriptions and performs fan-out filter matching
  21  // so the net domain's epoll loop isn't blocked by O(N*M*F) work per event.
  22  //
  23  // All parent→child messages go through a single `in` channel to avoid the
  24  // Moxie pipe channel select limitation. SubOpBcast carries broadcast events;
  25  // other SubOp* values carry conn/sub lifecycle updates.
  26  func BroadcastWorker(in chan SubCommand, frames chan BroadcastFrame, ready chan struct{}) {
  27  	conns := map[int32]*bcastConnState{}
  28  	for {
  29  		cmd, ok := <-in
  30  		if !ok {
  31  			return
  32  		}
  33  		if cmd.Op == SubOpBcast {
  34  			bcastHandleReq(cmd, conns, frames, ready)
  35  		} else {
  36  			bcastHandleCmd(cmd, conns)
  37  		}
  38  	}
  39  }
  40  
  41  func bcastHandleCmd(cmd SubCommand, conns map[int32]*bcastConnState) {
  42  	switch cmd.Op {
  43  	case SubOpNew:
  44  		// The per-connection state outlives this command: it stays in conns
  45  		// for the life of the connection, so build it in the root arena.
  46  		prev := runtime.CurrentArena()
  47  		runtime.SovereignSetArena(runtime.RootArena())
  48  		conns[cmd.ConnFD] = &bcastConnState{
  49  			whitelisted: cmd.Flags&1 != 0,
  50  			subs:        map[string]filter.S{},
  51  		}
  52  		runtime.SovereignRestoreArena(prev)
  53  	case SubOpAdd:
  54  		c := conns[cmd.ConnFD]
  55  		if c == nil {
  56  			return
  57  		}
  58  		// The subscription outlives this command - it is matched against
  59  		// every later fan-out - so parse it and keep it in the root arena.
  60  		// A frame-allocated filter array here is read after the frame is
  61  		// gone, which is a fault inside filter matching, not a wrong answer.
  62  		prev := runtime.CurrentArena()
  63  		runtime.SovereignSetArena(runtime.RootArena())
  64  		_, rem, _ := envelope.Identify(cmd.Bytes)
  65  		var req envelope.Req
  66  		if _, err := req.Unmarshal(rem); err != nil || req.Filters.F == nil {
  67  			runtime.SovereignRestoreArena(prev)
  68  			return
  69  		}
  70  		key := []byte{:len(cmd.SubID)}
  71  		copy(key, cmd.SubID)
  72  		c.subs[string(key)] = req.Filters
  73  		runtime.SovereignRestoreArena(prev)
  74  	case SubOpRemove:
  75  		c := conns[cmd.ConnFD]
  76  		if c == nil {
  77  			return
  78  		}
  79  		delete(c.subs, string(cmd.SubID))
  80  	case SubOpClose:
  81  		c := conns[cmd.ConnFD]
  82  		if c != nil {
  83  			for k := range c.subs {
  84  				delete(c.subs, k)
  85  			}
  86  			c.subs = nil
  87  		}
  88  		delete(conns, cmd.ConnFD)
  89  	case SubOpAuth:
  90  		c := conns[cmd.ConnFD]
  91  		if c == nil {
  92  			return
  93  		}
  94  		c.authed = true
  95  		if len(cmd.Bytes) == 32 {
  96  			copy(c.authedPubkey[:], cmd.Bytes)
  97  		}
  98  	}
  99  }
 100  
 101  // bcastHandleReq fans out an accepted event to all matching subscriptions.
 102  // ConnFD = senderFD to exclude, Flags = filter flags, Bytes = raw EVENT JSON.
 103  func bcastHandleReq(cmd SubCommand, conns map[int32]*bcastConnState, frames chan BroadcastFrame, ready chan struct{}) {
 104  	senderFD := cmd.ConnFD
 105  	needFilter := cmd.Flags&1 != 0
 106  	nip70 := cmd.Flags&2 != 0
 107  	marmotOpen := cmd.Flags&4 != 0
 108  
 109  	_, rem, _ := envelope.Identify(cmd.Bytes)
 110  	var es envelope.EventSubmission
 111  	if _, err := es.Unmarshal(rem); err != nil || es.E == nil {
 112  		return
 113  	}
 114  	ev := es.E
 115  
 116  	isMLS := ev.Kind == 443 || ev.Kind == 444 || ev.Kind == 445 || ev.Kind == 1059
 117  
 118  	scratch := []byte{:0:ev.EstimateSize() + 128}
 119  	er := &envelope.EventResult{Event: ev}
 120  	matchCount := 0
 121  	for connFD, c := range conns {
 122  		if connFD == senderFD {
 123  			continue
 124  		}
 125  		if needFilter && !c.whitelisted && !access.CanSee(c.authed, c.authedPubkey[:], ev, nip70, marmotOpen) {
 126  			if isMLS {
 127  			}
 128  			continue
 129  		}
 130  		for subID, filters := range c.subs {
 131  			matched := filters.Match(ev)
 132  			if isMLS {
 133  			}
 134  			if matched {
 135  				matchCount++
 136  				er.Subscription = []byte(subID)
 137  				marshaled := er.Marshal(scratch[:0])
 138  				out := []byte{:len(marshaled)}
 139  				copy(out, marshaled)
 140  				frames <- BroadcastFrame{ConnFD: connFD, Bytes: out}
 141  				ready <- struct{}{}
 142  			}
 143  		}
 144  	}
 145  	if isMLS {
 146  	}
 147  	scratch = nil
 148  	er.Event = nil
 149  	er = nil
 150  }
 151