package wire import ( "runtime" "git.smesh.lol/morly/pkg/access" "git.smesh.lol/nostr/pkg/envelope" "git.smesh.lol/nostr/pkg/filter" ) // bcastConnState holds the broadcast domain's view of one WS connection. type bcastConnState struct { whitelisted bool authed bool authedPubkey [32]byte subs map[string]filter.S } // BroadcastWorker is the spawn target for the broadcast domain. Owns a // registry of connection subscriptions and performs fan-out filter matching // so the net domain's epoll loop isn't blocked by O(N*M*F) work per event. // // All parent→child messages go through a single `in` channel to avoid the // Moxie pipe channel select limitation. SubOpBcast carries broadcast events; // other SubOp* values carry conn/sub lifecycle updates. func BroadcastWorker(in chan SubCommand, frames chan BroadcastFrame, ready chan struct{}) { conns := map[int32]*bcastConnState{} for { cmd, ok := <-in if !ok { return } if cmd.Op == SubOpBcast { bcastHandleReq(cmd, conns, frames, ready) } else { bcastHandleCmd(cmd, conns) } } } func bcastHandleCmd(cmd SubCommand, conns map[int32]*bcastConnState) { switch cmd.Op { case SubOpNew: // The per-connection state outlives this command: it stays in conns // for the life of the connection, so build it in the root arena. prev := runtime.CurrentArena() runtime.SovereignSetArena(runtime.RootArena()) conns[cmd.ConnFD] = &bcastConnState{ whitelisted: cmd.Flags&1 != 0, subs: map[string]filter.S{}, } runtime.SovereignRestoreArena(prev) case SubOpAdd: c := conns[cmd.ConnFD] if c == nil { return } // The subscription outlives this command - it is matched against // every later fan-out - so parse it and keep it in the root arena. // A frame-allocated filter array here is read after the frame is // gone, which is a fault inside filter matching, not a wrong answer. prev := runtime.CurrentArena() runtime.SovereignSetArena(runtime.RootArena()) _, rem, _ := envelope.Identify(cmd.Bytes) var req envelope.Req if _, err := req.Unmarshal(rem); err != nil || req.Filters.F == nil { runtime.SovereignRestoreArena(prev) return } key := []byte{:len(cmd.SubID)} copy(key, cmd.SubID) c.subs[string(key)] = req.Filters runtime.SovereignRestoreArena(prev) case SubOpRemove: c := conns[cmd.ConnFD] if c == nil { return } delete(c.subs, string(cmd.SubID)) case SubOpClose: c := conns[cmd.ConnFD] if c != nil { for k := range c.subs { delete(c.subs, k) } c.subs = nil } delete(conns, cmd.ConnFD) case SubOpAuth: c := conns[cmd.ConnFD] if c == nil { return } c.authed = true if len(cmd.Bytes) == 32 { copy(c.authedPubkey[:], cmd.Bytes) } } } // bcastHandleReq fans out an accepted event to all matching subscriptions. // ConnFD = senderFD to exclude, Flags = filter flags, Bytes = raw EVENT JSON. func bcastHandleReq(cmd SubCommand, conns map[int32]*bcastConnState, frames chan BroadcastFrame, ready chan struct{}) { senderFD := cmd.ConnFD needFilter := cmd.Flags&1 != 0 nip70 := cmd.Flags&2 != 0 marmotOpen := cmd.Flags&4 != 0 _, rem, _ := envelope.Identify(cmd.Bytes) var es envelope.EventSubmission if _, err := es.Unmarshal(rem); err != nil || es.E == nil { return } ev := es.E isMLS := ev.Kind == 443 || ev.Kind == 444 || ev.Kind == 445 || ev.Kind == 1059 scratch := []byte{:0:ev.EstimateSize() + 128} er := &envelope.EventResult{Event: ev} matchCount := 0 for connFD, c := range conns { if connFD == senderFD { continue } if needFilter && !c.whitelisted && !access.CanSee(c.authed, c.authedPubkey[:], ev, nip70, marmotOpen) { if isMLS { } continue } for subID, filters := range c.subs { matched := filters.Match(ev) if isMLS { } if matched { matchCount++ er.Subscription = []byte(subID) marshaled := er.Marshal(scratch[:0]) out := []byte{:len(marshaled)} copy(out, marshaled) frames <- BroadcastFrame{ConnFD: connFD, Bytes: out} ready <- struct{}{} } } } if isMLS { } scratch = nil er.Event = nil er = nil }