broadcast.mx raw

   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