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