1 // Package broadcast manages the live subscription fan-out domain.
2 // A BroadcastWorker spawn domain owns the subscription registry and filter
3 // matching; the parent drives it over the spawn channels bound at fork.
4 // This package hides the wire protocol behind named methods so the server
5 // never deals with SubCommand encoding or BroadcastFrame decoding directly.
6 package broadcast
7 8 import (
9 "runtime"
10 "git.smesh.lol/morly/pkg/relay/wire"
11 )
12 13 // preSpawnState owns the pre-spawned worker handles. The rule forbids stores to
14 // package globals outside init, so the state lives behind a self-mutating type
15 // and every write goes through a method on it.
16 type preSpawnState struct {
17 in chan wire.SubCommand
18 out chan wire.BroadcastFrame
19 done chan struct{}
20 ready chan struct{}
21 used bool
22 }
23 24 func (p *preSpawnState) set(in chan wire.SubCommand, out chan wire.BroadcastFrame, ready chan struct{}, done chan struct{}) {
25 p.in = in
26 p.out = out
27 p.ready = ready
28 p.done = done
29 p.used = true
30 }
31 32 // take hands the pre-spawned handles to one Broadcaster and marks them used,
33 // so a second New spawns its own worker.
34 func (p *preSpawnState) take() (b *Broadcaster, ok bool) {
35 if !p.used {
36 return nil, false
37 }
38 p.used = false
39 return &Broadcaster{in: p.in, out: p.out, done: p.done, ready: p.ready}, true
40 }
41 42 var preSpawn preSpawnState
43 44 // PreSpawn forks the broadcast worker early, before config loading, so relay
45 // mode does not pay a fork in the middle of serving. Non-relay commands
46 // (sync, crawl, outbox) skip PreSpawn and never pay the child process cost.
47 func PreSpawn() {
48 in := chan wire.SubCommand{}
49 out := chan wire.BroadcastFrame{}
50 ready := chan struct{}{}
51 done := spawn(wire.BroadcastWorker, in, out, ready)
52 preSpawn.set(in, out, ready, done)
53 }
54 55 // Broadcaster is the parent-side handle to the broadcast worker domain.
56 type Broadcaster struct {
57 in chan wire.SubCommand
58 out chan wire.BroadcastFrame
59 done chan struct{}
60 ready chan struct{}
61 }
62 63 // New wires the pre-spawned broadcast worker into a Broadcaster, or spawns a
64 // fresh one when called without PreSpawn (test binaries).
65 func New() (b *Broadcaster) {
66 if pb, ok := preSpawn.take(); ok {
67 return pb
68 }
69 // Root lifetime: the server holds this broadcaster for the whole run and
70 // the forked worker reads the same channel memory. chanMake allocates in
71 // the current arena, so a plain frame here leaves every stored channel
72 // dangling once New returns - the live fanout then silently delivers
73 // nothing.
74 runtime.SovereignSetArena(runtime.RootArena())
75 in := chan wire.SubCommand{}
76 out := chan wire.BroadcastFrame{}
77 ready := chan struct{}{}
78 done := spawn(wire.BroadcastWorker, in, out, ready)
79 runtime.SovereignRestoreArena(runtime.RootArena())
80 return &Broadcaster{in: in, out: out, done: done, ready: ready}
81 }
82 83 // Valid reports whether the worker is still running.
84 func (b *Broadcaster) Valid() (ok bool) {
85 select {
86 case <-b.done:
87 return false
88 default:
89 }
90 return true
91 }
92 93 // Close stops the broadcast worker. The child's receive returns closed, its
94 // loop returns, and the runtime dumps the domain's coverage counters on the
95 // way out - the same path the other worker pools take, rather than being
96 // orphaned and killed by the process-group signal.
97 func (b *Broadcaster) Close() {
98 close(b.in)
99 }
100 101 // OnConnect notifies the broadcast domain of a new WS connection.
102 func (b *Broadcaster) OnConnect(fd int32, whitelisted bool) {
103 var flags uint8
104 if whitelisted {
105 flags = 1
106 }
107 b.send(wire.SubCommand{Op: wire.SubOpNew, ConnFD: fd, Flags: flags})
108 }
109 110 // OnDisconnect notifies the broadcast domain that a WS connection closed.
111 func (b *Broadcaster) OnDisconnect(fd int32) {
112 b.send(wire.SubCommand{Op: wire.SubOpClose, ConnFD: fd})
113 }
114 115 // OnSubscribe registers a REQ subscription.
116 func (b *Broadcaster) OnSubscribe(fd int32, subID []byte, reqMsg []byte) {
117 b.send(wire.SubCommand{Op: wire.SubOpAdd, ConnFD: fd, SubID: subID, Bytes: reqMsg})
118 }
119 120 // OnUnsubscribe removes a subscription.
121 func (b *Broadcaster) OnUnsubscribe(fd int32, subID []byte) {
122 b.send(wire.SubCommand{Op: wire.SubOpRemove, ConnFD: fd, SubID: subID})
123 }
124 125 // OnAuth updates the broadcast domain with the authed pubkey for a connection.
126 func (b *Broadcaster) OnAuth(fd int32, pubkey []byte) {
127 b.send(wire.SubCommand{Op: wire.SubOpAuth, ConnFD: fd, Bytes: pubkey})
128 }
129 130 // Fanout dispatches an accepted event for fan-out to matching subscribers.
131 // flags: bit0=needFilter, bit1=nip70Enforce, bit2=marmotOpen.
132 func (b *Broadcaster) Fanout(rawMsg []byte, senderFD int32, flags uint8) {
133 b.send(wire.SubCommand{Op: wire.SubOpBcast, ConnFD: senderFD, Flags: flags, Bytes: rawMsg})
134 }
135 136 // ReadFrame takes one queued BroadcastFrame without blocking. Returns
137 // ok=false when no frame is ready; the caller polls again next turn.
138 func (b *Broadcaster) ReadFrame() (connFD int32, msg []byte, ok bool) {
139 // The worker announces each frame on a zero-size channel and the payload
140 // travels on the codec channel. A select receive on a codec-framed spawn
141 // channel never decodes (the compiler passes no codec to chanSelect, so
142 // the ring is drained raw and the value arrives zeroed); a plain receive
143 // does decode. Polling the zero-size channel keeps this non-blocking.
144 select {
145 case <-b.ready:
146 f := <-b.out
147 return f.ConnFD, f.Bytes, true
148 default:
149 }
150 return 0, nil, false
151 }
152 153 func (b *Broadcaster) send(cmd wire.SubCommand) {
154 if !b.Valid() {
155 return
156 }
157 b.in <- cmd
158 }
159