// Package tree defines the messages that cross spawn boundaries in the relay // spawn tree described in docs/relay-tree.md. // // The database-engine node is the single writer and reader; these are its // request and response messages. Everything that reaches it has already passed // the signature gate and the policy node, and carries the raw event bytes so // the event is serialized once and parsed once, at the node that needs it. // // One request channel carries every operation, discriminated by Op. A select // receive on a codec-framed spawn channel does not decode (the compiler passes // no codec to chanSelect, so the ring is drained raw and the value arrives // zeroed), so the worker reads its inbox with a plain receive and the parent // learns of a reply through a zero-size nudge channel it can select on. package tree import ( "encoding/binary" "io" ) // Database-engine operations. const ( OpPersist uint8 = 1 // store one event; Bytes is the raw EVENT frame OpHistory uint8 = 2 // backfill one subscription; Filter is the raw REQ frame OpCount uint8 = 3 // count one COUNT request; Filter is the raw COUNT frame ) // Request carries one operation to database-engine. // // ConnID is the originating connection for OpPersist and the requesting // connection for OpHistory and OpCount. Verified means the signature gate has // already accepted the event, so the database node runs only stage B. // // AuthedPubkey, Filtered and the two access flags are the connection's // visibility context for OpHistory. Filtered means this connection is subject // to visibility filtering at all; the database node applies it per event // because that is where the event already is. type Request struct { Op uint8 ReqID uint32 ConnID int32 Verified bool Filtered bool NIP70 bool Marmot bool Limit int32 SubID []byte Filter []byte Bytes []byte AuthedPubkey []byte } func (m Request) EncodeTo(w io.Writer) error { var err error var hdr [30]byte hdr[0] = m.Op binary.LittleEndian().PutUint32(hdr[1:5], m.ReqID) binary.LittleEndian().PutUint32(hdr[5:9], uint32(m.ConnID)) if m.Verified { hdr[9] |= 1 } if m.Filtered { hdr[9] |= 2 } if m.NIP70 { hdr[9] |= 4 } if m.Marmot { hdr[9] |= 8 } binary.LittleEndian().PutUint32(hdr[10:14], uint32(m.Limit)) binary.LittleEndian().PutUint32(hdr[14:18], uint32(len(m.SubID))) binary.LittleEndian().PutUint32(hdr[18:22], uint32(len(m.Filter))) binary.LittleEndian().PutUint32(hdr[22:26], uint32(len(m.Bytes))) binary.LittleEndian().PutUint32(hdr[26:30], uint32(len(m.AuthedPubkey))) if _, err = w.Write(hdr[:]); err != nil { return err } if err = writeIf(w, m.SubID); err != nil { return err } if err = writeIf(w, m.Filter); err != nil { return err } if err = writeIf(w, m.Bytes); err != nil { return err } return writeIf(w, m.AuthedPubkey) } func (m *Request) DecodeFrom(rd io.Reader) error { var err error var hdr [30]byte if _, err = io.ReadFull(rd, hdr[:]); err != nil { return err } m.Op = hdr[0] m.ReqID = binary.LittleEndian().Uint32(hdr[1:5]) m.ConnID = int32(binary.LittleEndian().Uint32(hdr[5:9])) m.Verified = hdr[9]&1 != 0 m.Filtered = hdr[9]&2 != 0 m.NIP70 = hdr[9]&4 != 0 m.Marmot = hdr[9]&8 != 0 m.Limit = int32(binary.LittleEndian().Uint32(hdr[10:14])) if m.SubID, err = readN(rd, int32(binary.LittleEndian().Uint32(hdr[14:18]))); err != nil { return err } if m.Filter, err = readN(rd, int32(binary.LittleEndian().Uint32(hdr[18:22]))); err != nil { return err } if m.Bytes, err = readN(rd, int32(binary.LittleEndian().Uint32(hdr[22:26]))); err != nil { return err } m.AuthedPubkey, err = readN(rd, int32(binary.LittleEndian().Uint32(hdr[26:30]))) return err } // Response carries the outcome of one operation. // // For OpPersist: OK, Reason, EventID and Seq. Seq is the monotonic database // record assigned to a stored event, and 0 when the event was not stored: a // duplicate, a rejection, or an ephemeral event that bypasses persistence. // Bytes echoes the accepted event back to root so the broadcast fan-out can // forward it without root keeping a copy of every in-flight event. // // For OpHistory: SubID, SeqHighWater, Done and Events, each event already // marshaled as an EVENT frame. SeqHighWater is the database's max seq at the // moment of the query, which lets a subscription register its live floor at // highWater+1 so a stored event is never delivered twice. // // For OpCount: Count. type Response struct { Op uint8 ReqID uint32 ConnID int32 OK bool Done bool Count int32 Seq uint64 SeqHighWater uint64 SubID []byte EventID []byte Reason []byte Bytes []byte Events [][]byte } func (m Response) EncodeTo(w io.Writer) error { var err error var hdr [50]byte hdr[0] = m.Op binary.LittleEndian().PutUint32(hdr[1:5], m.ReqID) binary.LittleEndian().PutUint32(hdr[5:9], uint32(m.ConnID)) if m.OK { hdr[9] |= 1 } if m.Done { hdr[9] |= 2 } binary.LittleEndian().PutUint32(hdr[10:14], uint32(m.Count)) binary.LittleEndian().PutUint64(hdr[14:22], m.Seq) binary.LittleEndian().PutUint64(hdr[22:30], m.SeqHighWater) binary.LittleEndian().PutUint32(hdr[30:34], uint32(len(m.SubID))) binary.LittleEndian().PutUint32(hdr[34:38], uint32(len(m.EventID))) binary.LittleEndian().PutUint32(hdr[38:42], uint32(len(m.Reason))) binary.LittleEndian().PutUint32(hdr[42:46], uint32(len(m.Bytes))) binary.LittleEndian().PutUint32(hdr[46:50], uint32(len(m.Events))) if _, err = w.Write(hdr[:]); err != nil { return err } if err = writeIf(w, m.SubID); err != nil { return err } if err = writeIf(w, m.EventID); err != nil { return err } if err = writeIf(w, m.Reason); err != nil { return err } if err = writeIf(w, m.Bytes); err != nil { return err } for _, ev := range m.Events { if err = writeFrame(w, ev); err != nil { return err } } return nil } func (m *Response) DecodeFrom(rd io.Reader) error { var err error var hdr [50]byte if _, err = io.ReadFull(rd, hdr[:]); err != nil { return err } m.Op = hdr[0] m.ReqID = binary.LittleEndian().Uint32(hdr[1:5]) m.ConnID = int32(binary.LittleEndian().Uint32(hdr[5:9])) m.OK = hdr[9]&1 != 0 m.Done = hdr[9]&2 != 0 m.Count = int32(binary.LittleEndian().Uint32(hdr[10:14])) m.Seq = binary.LittleEndian().Uint64(hdr[14:22]) m.SeqHighWater = binary.LittleEndian().Uint64(hdr[22:30]) subLen := int32(binary.LittleEndian().Uint32(hdr[30:34])) idLen := int32(binary.LittleEndian().Uint32(hdr[34:38])) reasonLen := int32(binary.LittleEndian().Uint32(hdr[38:42])) bytesLen := int32(binary.LittleEndian().Uint32(hdr[42:46])) nEvents := int32(binary.LittleEndian().Uint32(hdr[46:50])) if m.SubID, err = readN(rd, subLen); err != nil { return err } if m.EventID, err = readN(rd, idLen); err != nil { return err } if m.Reason, err = readN(rd, reasonLen); err != nil { return err } if m.Bytes, err = readN(rd, bytesLen); err != nil { return err } if nEvents <= 0 { return nil } m.Events = [][]byte{:nEvents} n := int32(0) var ev []byte var e error for i := int32(0); i < nEvents; i++ { ev, e = readFrame(rd) if e != nil { return e } m.Events[n] = ev n++ } m.Events = m.Events[:n] return nil } // --- helpers --- // writeIf writes b only when it is non-empty. func writeIf(w io.Writer, b []byte) (err error) { if len(b) == 0 { return nil } _, err = w.Write(b) return err } // writeFrame writes one length-prefixed blob. func writeFrame(w io.Writer, b []byte) (err error) { var lh [4]byte binary.LittleEndian().PutUint32(lh[0:4], uint32(len(b))) if _, err = w.Write(lh[:]); err != nil { return err } if len(b) == 0 { return nil } _, err = w.Write(b) return err } // readN reads exactly n bytes, or nothing when n <= 0. func readN(rd io.Reader, n int32) (b []byte, err error) { if n <= 0 { return nil, nil } b = []byte{:n} _, err = io.ReadFull(rd, b) return b, err } // readFrame reads one length-prefixed blob. func readFrame(rd io.Reader) (b []byte, err error) { var lh [4]byte if _, err = io.ReadFull(rd, lh[:]); err != nil { return nil, err } return readN(rd, int32(binary.LittleEndian().Uint32(lh[0:4]))) }