messages.mx raw

   1  // Package tree defines the messages that cross spawn boundaries in the relay
   2  // spawn tree described in docs/relay-tree.md.
   3  //
   4  // The database-engine node is the single writer and reader; these are its
   5  // request and response messages. Everything that reaches it has already passed
   6  // the signature gate and the policy node, and carries the raw event bytes so
   7  // the event is serialized once and parsed once, at the node that needs it.
   8  //
   9  // One request channel carries every operation, discriminated by Op. A select
  10  // receive on a codec-framed spawn channel does not decode (the compiler passes
  11  // no codec to chanSelect, so the ring is drained raw and the value arrives
  12  // zeroed), so the worker reads its inbox with a plain receive and the parent
  13  // learns of a reply through a zero-size nudge channel it can select on.
  14  package tree
  15  
  16  import (
  17  	"encoding/binary"
  18  	"io"
  19  )
  20  
  21  // Database-engine operations.
  22  const (
  23  	OpPersist uint8 = 1 // store one event; Bytes is the raw EVENT frame
  24  	OpHistory uint8 = 2 // backfill one subscription; Filter is the raw REQ frame
  25  	OpCount   uint8 = 3 // count one COUNT request; Filter is the raw COUNT frame
  26  )
  27  
  28  // Request carries one operation to database-engine.
  29  //
  30  // ConnID is the originating connection for OpPersist and the requesting
  31  // connection for OpHistory and OpCount. Verified means the signature gate has
  32  // already accepted the event, so the database node runs only stage B.
  33  //
  34  // AuthedPubkey, Filtered and the two access flags are the connection's
  35  // visibility context for OpHistory. Filtered means this connection is subject
  36  // to visibility filtering at all; the database node applies it per event
  37  // because that is where the event already is.
  38  type Request struct {
  39  	Op           uint8
  40  	ReqID        uint32
  41  	ConnID       int32
  42  	Verified     bool
  43  	Filtered     bool
  44  	NIP70        bool
  45  	Marmot       bool
  46  	Limit        int32
  47  	SubID        []byte
  48  	Filter       []byte
  49  	Bytes        []byte
  50  	AuthedPubkey []byte
  51  }
  52  
  53  func (m Request) EncodeTo(w io.Writer) error {
  54  	var err error
  55  	var hdr [30]byte
  56  	hdr[0] = m.Op
  57  	binary.LittleEndian().PutUint32(hdr[1:5], m.ReqID)
  58  	binary.LittleEndian().PutUint32(hdr[5:9], uint32(m.ConnID))
  59  	if m.Verified {
  60  		hdr[9] |= 1
  61  	}
  62  	if m.Filtered {
  63  		hdr[9] |= 2
  64  	}
  65  	if m.NIP70 {
  66  		hdr[9] |= 4
  67  	}
  68  	if m.Marmot {
  69  		hdr[9] |= 8
  70  	}
  71  	binary.LittleEndian().PutUint32(hdr[10:14], uint32(m.Limit))
  72  	binary.LittleEndian().PutUint32(hdr[14:18], uint32(len(m.SubID)))
  73  	binary.LittleEndian().PutUint32(hdr[18:22], uint32(len(m.Filter)))
  74  	binary.LittleEndian().PutUint32(hdr[22:26], uint32(len(m.Bytes)))
  75  	binary.LittleEndian().PutUint32(hdr[26:30], uint32(len(m.AuthedPubkey)))
  76  	if _, err = w.Write(hdr[:]); err != nil {
  77  		return err
  78  	}
  79  	if err = writeIf(w, m.SubID); err != nil {
  80  		return err
  81  	}
  82  	if err = writeIf(w, m.Filter); err != nil {
  83  		return err
  84  	}
  85  	if err = writeIf(w, m.Bytes); err != nil {
  86  		return err
  87  	}
  88  	return writeIf(w, m.AuthedPubkey)
  89  }
  90  
  91  func (m *Request) DecodeFrom(rd io.Reader) error {
  92  	var err error
  93  	var hdr [30]byte
  94  	if _, err = io.ReadFull(rd, hdr[:]); err != nil {
  95  		return err
  96  	}
  97  	m.Op = hdr[0]
  98  	m.ReqID = binary.LittleEndian().Uint32(hdr[1:5])
  99  	m.ConnID = int32(binary.LittleEndian().Uint32(hdr[5:9]))
 100  	m.Verified = hdr[9]&1 != 0
 101  	m.Filtered = hdr[9]&2 != 0
 102  	m.NIP70 = hdr[9]&4 != 0
 103  	m.Marmot = hdr[9]&8 != 0
 104  	m.Limit = int32(binary.LittleEndian().Uint32(hdr[10:14]))
 105  	if m.SubID, err = readN(rd, int32(binary.LittleEndian().Uint32(hdr[14:18]))); err != nil {
 106  		return err
 107  	}
 108  	if m.Filter, err = readN(rd, int32(binary.LittleEndian().Uint32(hdr[18:22]))); err != nil {
 109  		return err
 110  	}
 111  	if m.Bytes, err = readN(rd, int32(binary.LittleEndian().Uint32(hdr[22:26]))); err != nil {
 112  		return err
 113  	}
 114  	m.AuthedPubkey, err = readN(rd, int32(binary.LittleEndian().Uint32(hdr[26:30])))
 115  	return err
 116  }
 117  
 118  // Response carries the outcome of one operation.
 119  //
 120  // For OpPersist: OK, Reason, EventID and Seq. Seq is the monotonic database
 121  // record assigned to a stored event, and 0 when the event was not stored: a
 122  // duplicate, a rejection, or an ephemeral event that bypasses persistence.
 123  // Bytes echoes the accepted event back to root so the broadcast fan-out can
 124  // forward it without root keeping a copy of every in-flight event.
 125  //
 126  // For OpHistory: SubID, SeqHighWater, Done and Events, each event already
 127  // marshaled as an EVENT frame. SeqHighWater is the database's max seq at the
 128  // moment of the query, which lets a subscription register its live floor at
 129  // highWater+1 so a stored event is never delivered twice.
 130  //
 131  // For OpCount: Count.
 132  type Response struct {
 133  	Op           uint8
 134  	ReqID        uint32
 135  	ConnID       int32
 136  	OK           bool
 137  	Done         bool
 138  	Count        int32
 139  	Seq          uint64
 140  	SeqHighWater uint64
 141  	SubID        []byte
 142  	EventID      []byte
 143  	Reason       []byte
 144  	Bytes        []byte
 145  	Events       [][]byte
 146  }
 147  
 148  func (m Response) EncodeTo(w io.Writer) error {
 149  	var err error
 150  	var hdr [50]byte
 151  	hdr[0] = m.Op
 152  	binary.LittleEndian().PutUint32(hdr[1:5], m.ReqID)
 153  	binary.LittleEndian().PutUint32(hdr[5:9], uint32(m.ConnID))
 154  	if m.OK {
 155  		hdr[9] |= 1
 156  	}
 157  	if m.Done {
 158  		hdr[9] |= 2
 159  	}
 160  	binary.LittleEndian().PutUint32(hdr[10:14], uint32(m.Count))
 161  	binary.LittleEndian().PutUint64(hdr[14:22], m.Seq)
 162  	binary.LittleEndian().PutUint64(hdr[22:30], m.SeqHighWater)
 163  	binary.LittleEndian().PutUint32(hdr[30:34], uint32(len(m.SubID)))
 164  	binary.LittleEndian().PutUint32(hdr[34:38], uint32(len(m.EventID)))
 165  	binary.LittleEndian().PutUint32(hdr[38:42], uint32(len(m.Reason)))
 166  	binary.LittleEndian().PutUint32(hdr[42:46], uint32(len(m.Bytes)))
 167  	binary.LittleEndian().PutUint32(hdr[46:50], uint32(len(m.Events)))
 168  	if _, err = w.Write(hdr[:]); err != nil {
 169  		return err
 170  	}
 171  	if err = writeIf(w, m.SubID); err != nil {
 172  		return err
 173  	}
 174  	if err = writeIf(w, m.EventID); err != nil {
 175  		return err
 176  	}
 177  	if err = writeIf(w, m.Reason); err != nil {
 178  		return err
 179  	}
 180  	if err = writeIf(w, m.Bytes); err != nil {
 181  		return err
 182  	}
 183  	for _, ev := range m.Events {
 184  		if err = writeFrame(w, ev); err != nil {
 185  			return err
 186  		}
 187  	}
 188  	return nil
 189  }
 190  
 191  func (m *Response) DecodeFrom(rd io.Reader) error {
 192  	var err error
 193  	var hdr [50]byte
 194  	if _, err = io.ReadFull(rd, hdr[:]); err != nil {
 195  		return err
 196  	}
 197  	m.Op = hdr[0]
 198  	m.ReqID = binary.LittleEndian().Uint32(hdr[1:5])
 199  	m.ConnID = int32(binary.LittleEndian().Uint32(hdr[5:9]))
 200  	m.OK = hdr[9]&1 != 0
 201  	m.Done = hdr[9]&2 != 0
 202  	m.Count = int32(binary.LittleEndian().Uint32(hdr[10:14]))
 203  	m.Seq = binary.LittleEndian().Uint64(hdr[14:22])
 204  	m.SeqHighWater = binary.LittleEndian().Uint64(hdr[22:30])
 205  	subLen := int32(binary.LittleEndian().Uint32(hdr[30:34]))
 206  	idLen := int32(binary.LittleEndian().Uint32(hdr[34:38]))
 207  	reasonLen := int32(binary.LittleEndian().Uint32(hdr[38:42]))
 208  	bytesLen := int32(binary.LittleEndian().Uint32(hdr[42:46]))
 209  	nEvents := int32(binary.LittleEndian().Uint32(hdr[46:50]))
 210  	if m.SubID, err = readN(rd, subLen); err != nil {
 211  		return err
 212  	}
 213  	if m.EventID, err = readN(rd, idLen); err != nil {
 214  		return err
 215  	}
 216  	if m.Reason, err = readN(rd, reasonLen); err != nil {
 217  		return err
 218  	}
 219  	if m.Bytes, err = readN(rd, bytesLen); err != nil {
 220  		return err
 221  	}
 222  	if nEvents <= 0 {
 223  		return nil
 224  	}
 225  	m.Events = [][]byte{:nEvents}
 226  	n := int32(0)
 227  	var ev []byte
 228  	var e error
 229  	for i := int32(0); i < nEvents; i++ {
 230  		ev, e = readFrame(rd)
 231  		if e != nil {
 232  			return e
 233  		}
 234  		m.Events[n] = ev
 235  		n++
 236  	}
 237  	m.Events = m.Events[:n]
 238  	return nil
 239  }
 240  
 241  // --- helpers ---
 242  
 243  // writeIf writes b only when it is non-empty.
 244  func writeIf(w io.Writer, b []byte) (err error) {
 245  	if len(b) == 0 {
 246  		return nil
 247  	}
 248  	_, err = w.Write(b)
 249  	return err
 250  }
 251  
 252  // writeFrame writes one length-prefixed blob.
 253  func writeFrame(w io.Writer, b []byte) (err error) {
 254  	var lh [4]byte
 255  	binary.LittleEndian().PutUint32(lh[0:4], uint32(len(b)))
 256  	if _, err = w.Write(lh[:]); err != nil {
 257  		return err
 258  	}
 259  	if len(b) == 0 {
 260  		return nil
 261  	}
 262  	_, err = w.Write(b)
 263  	return err
 264  }
 265  
 266  // readN reads exactly n bytes, or nothing when n <= 0.
 267  func readN(rd io.Reader, n int32) (b []byte, err error) {
 268  	if n <= 0 {
 269  		return nil, nil
 270  	}
 271  	b = []byte{:n}
 272  	_, err = io.ReadFull(rd, b)
 273  	return b, err
 274  }
 275  
 276  // readFrame reads one length-prefixed blob.
 277  func readFrame(rd io.Reader) (b []byte, err error) {
 278  	var lh [4]byte
 279  	if _, err = io.ReadFull(rd, lh[:]); err != nil {
 280  		return nil, err
 281  	}
 282  	return readN(rd, int32(binary.LittleEndian().Uint32(lh[0:4])))
 283  }
 284