// Package broadcast manages the live subscription fan-out domain. // A BroadcastWorker spawn domain owns the subscription registry and filter // matching; the parent drives it over the spawn channels bound at fork. // This package hides the wire protocol behind named methods so the server // never deals with SubCommand encoding or BroadcastFrame decoding directly. package broadcast import ( "runtime" "git.smesh.lol/morly/pkg/relay/wire" ) // preSpawnState owns the pre-spawned worker handles. The rule forbids stores to // package globals outside init, so the state lives behind a self-mutating type // and every write goes through a method on it. type preSpawnState struct { in chan wire.SubCommand out chan wire.BroadcastFrame done chan struct{} ready chan struct{} used bool } func (p *preSpawnState) set(in chan wire.SubCommand, out chan wire.BroadcastFrame, ready chan struct{}, done chan struct{}) { p.in = in p.out = out p.ready = ready p.done = done p.used = true } // take hands the pre-spawned handles to one Broadcaster and marks them used, // so a second New spawns its own worker. func (p *preSpawnState) take() (b *Broadcaster, ok bool) { if !p.used { return nil, false } p.used = false return &Broadcaster{in: p.in, out: p.out, done: p.done, ready: p.ready}, true } var preSpawn preSpawnState // PreSpawn forks the broadcast worker early, before config loading, so relay // mode does not pay a fork in the middle of serving. Non-relay commands // (sync, crawl, outbox) skip PreSpawn and never pay the child process cost. func PreSpawn() { in := chan wire.SubCommand{} out := chan wire.BroadcastFrame{} ready := chan struct{}{} done := spawn(wire.BroadcastWorker, in, out, ready) preSpawn.set(in, out, ready, done) } // Broadcaster is the parent-side handle to the broadcast worker domain. type Broadcaster struct { in chan wire.SubCommand out chan wire.BroadcastFrame done chan struct{} ready chan struct{} } // New wires the pre-spawned broadcast worker into a Broadcaster, or spawns a // fresh one when called without PreSpawn (test binaries). func New() (b *Broadcaster) { if pb, ok := preSpawn.take(); ok { return pb } // Root lifetime: the server holds this broadcaster for the whole run and // the forked worker reads the same channel memory. chanMake allocates in // the current arena, so a plain frame here leaves every stored channel // dangling once New returns - the live fanout then silently delivers // nothing. runtime.SovereignSetArena(runtime.RootArena()) in := chan wire.SubCommand{} out := chan wire.BroadcastFrame{} ready := chan struct{}{} done := spawn(wire.BroadcastWorker, in, out, ready) runtime.SovereignRestoreArena(runtime.RootArena()) return &Broadcaster{in: in, out: out, done: done, ready: ready} } // Valid reports whether the worker is still running. func (b *Broadcaster) Valid() (ok bool) { select { case <-b.done: return false default: } return true } // Close stops the broadcast worker. The child's receive returns closed, its // loop returns, and the runtime dumps the domain's coverage counters on the // way out - the same path the other worker pools take, rather than being // orphaned and killed by the process-group signal. func (b *Broadcaster) Close() { close(b.in) } // OnConnect notifies the broadcast domain of a new WS connection. func (b *Broadcaster) OnConnect(fd int32, whitelisted bool) { var flags uint8 if whitelisted { flags = 1 } b.send(wire.SubCommand{Op: wire.SubOpNew, ConnFD: fd, Flags: flags}) } // OnDisconnect notifies the broadcast domain that a WS connection closed. func (b *Broadcaster) OnDisconnect(fd int32) { b.send(wire.SubCommand{Op: wire.SubOpClose, ConnFD: fd}) } // OnSubscribe registers a REQ subscription. func (b *Broadcaster) OnSubscribe(fd int32, subID []byte, reqMsg []byte) { b.send(wire.SubCommand{Op: wire.SubOpAdd, ConnFD: fd, SubID: subID, Bytes: reqMsg}) } // OnUnsubscribe removes a subscription. func (b *Broadcaster) OnUnsubscribe(fd int32, subID []byte) { b.send(wire.SubCommand{Op: wire.SubOpRemove, ConnFD: fd, SubID: subID}) } // OnAuth updates the broadcast domain with the authed pubkey for a connection. func (b *Broadcaster) OnAuth(fd int32, pubkey []byte) { b.send(wire.SubCommand{Op: wire.SubOpAuth, ConnFD: fd, Bytes: pubkey}) } // Fanout dispatches an accepted event for fan-out to matching subscribers. // flags: bit0=needFilter, bit1=nip70Enforce, bit2=marmotOpen. func (b *Broadcaster) Fanout(rawMsg []byte, senderFD int32, flags uint8) { b.send(wire.SubCommand{Op: wire.SubOpBcast, ConnFD: senderFD, Flags: flags, Bytes: rawMsg}) } // ReadFrame takes one queued BroadcastFrame without blocking. Returns // ok=false when no frame is ready; the caller polls again next turn. func (b *Broadcaster) ReadFrame() (connFD int32, msg []byte, ok bool) { // The worker announces each frame on a zero-size channel and the payload // travels on the codec channel. A select receive on a codec-framed spawn // channel never decodes (the compiler passes no codec to chanSelect, so // the ring is drained raw and the value arrives zeroed); a plain receive // does decode. Polling the zero-size channel keeps this non-blocking. select { case <-b.ready: f := <-b.out return f.ConnFD, f.Bytes, true default: } return 0, nil, false } func (b *Broadcaster) send(cmd wire.SubCommand) { if !b.Valid() { return } b.in <- cmd }