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