dbengine.mx raw
1 // Package dbengine is the relay's database-engine domain: the one node that
2 // opens the store, and therefore the single writer and the single reader.
3 //
4 // Every other node reaches storage by sending it a tree.Request. It answers on
5 // tree.Response and nudges the parent on a zero-size channel (see the note on
6 // tree.Request for why the parent cannot select on the codec channel itself).
7 //
8 // The engine owns everything that is a function of stored state: the storage
9 // engine, the ingestion pipeline (limits, ACL, replaceable and deletion
10 // resolution, WAL append), and the admin mute blacklist. Root owns everything
11 // that is a function of connections: transport, subscriptions, and the
12 // authentication and rate policy that decides whether a frame is even worth
13 // sending here.
14 package dbengine
15
16 import (
17 "bytes"
18 "fmt"
19 "os"
20
21 "git.smesh.lol/morly/pkg/access"
22 "git.smesh.lol/morly/pkg/acl"
23 "git.smesh.lol/moxie/pkg/mxutil"
24 "git.smesh.lol/nostr/pkg/envelope"
25 "git.smesh.lol/nostr/pkg/event"
26 "git.smesh.lol/nostr/pkg/filter"
27 "git.smesh.lol/nostr/pkg/kind"
28 "git.smesh.lol/morly/pkg/relay/mute"
29 "git.smesh.lol/morly/pkg/relay/pipeline"
30 "git.smesh.lol/morly/pkg/relay/tree"
31 "git.smesh.lol/morly/pkg/store"
32 )
33
34 // Config is what the database node needs to open its store and build its
35 // policy. It is deliberately not the whole relay config: the relay config
36 // belongs to root, and only these values are a function of stored state.
37 type Config struct {
38 DataDir string
39 ACLMode string
40 Admins []string
41 FollowListFreqSec int32
42 SocialWoTMaxDepth int32
43 SocialWoTRefreshSec int32
44 MuteBlacklist string
45 }
46
47 // Run is the spawn target of the database-engine domain. It opens the store,
48 // then serves requests until the inbox closes.
49 func Run(cfg Config, in chan tree.Request, out chan tree.Response, ready chan struct{}) {
50 eng, err := store.Open(cfg.DataDir | "/db")
51 if err != nil {
52 fmt.Fprintln(os.Stderr, "store error: "|err.Error())
53 return
54 }
55 checker := buildACL(cfg, eng)
56 pipe := pipeline.New(eng, checker, nil, pipeline.DefaultConfig())
57 bl := loadMute(cfg.MuteBlacklist, eng)
58
59 for {
60 req, ok := <-in
61 if !ok {
62 break
63 }
64 resp := serve(eng, pipe, bl, req)
65 out <- resp
66 ready <- struct{}{}
67 }
68 eng.Flush()
69 eng.Close()
70 }
71
72 // buildACL constructs the write-access checker. The checker answers from
73 // stored follow and relay-list events, so it lives on this side of the
74 // boundary.
75 func buildACL(cfg Config, eng *store.Engine) (c acl.Checker) {
76 switch cfg.ACLMode {
77 case "follows":
78 return acl.NewFollows(eng, cfg.Admins, cfg.FollowListFreqSec)
79 case "social":
80 return acl.NewSocial(eng, cfg.Admins, cfg.SocialWoTMaxDepth, cfg.SocialWoTRefreshSec)
81 }
82 return &acl.Open{}
83 }
84
85 // loadMute loads the admin mute blacklist and purges what it blocks. Returns
86 // nil when no admin pubkey is configured.
87 func loadMute(adminHexPK string, eng *store.Engine) (b *mute.Blacklist) {
88 b = mute.New(adminHexPK)
89 if b == nil {
90 return nil
91 }
92 b.Load(eng)
93 b.Purge(eng)
94 return b
95 }
96
97 // serve runs one request. A free function: its scratch dies at return instead
98 // of accumulating in Run's frame for the life of the relay.
99 func serve(eng *store.Engine, pipe *pipeline.Pipeline, bl *mute.Blacklist, req tree.Request) (resp tree.Response) {
100 switch req.Op {
101 case tree.OpPersist:
102 return persist(eng, pipe, bl, req)
103 case tree.OpHistory:
104 return history(eng, req)
105 case tree.OpCount:
106 return count(eng, req)
107 }
108 return tree.Response{}
109 }
110
111 // persist stores one event, or reports why it was not stored.
112 //
113 // Ephemeral events pass the same limits, ACL and expiration checks as stored
114 // ones, and are then accepted without a WAL append: they are live-only, so they
115 // carry no sequence number and root still broadcasts them.
116 func persist(eng *store.Engine, pipe *pipeline.Pipeline, bl *mute.Blacklist, req tree.Request) (resp tree.Response) {
117 resp.Op = tree.OpPersist
118 resp.ReqID = req.ReqID
119 resp.ConnID = req.ConnID
120
121 label, rem, _ := envelope.Identify(req.Bytes)
122 if label != envelope.EventLabel {
123 resp.Reason = []byte("invalid: not an EVENT envelope")
124 return resp
125 }
126 var es envelope.EventSubmission
127 if _, perr := es.Unmarshal(rem); perr != nil || es.E == nil {
128 resp.Reason = []byte("invalid: parse error")
129 return resp
130 }
131 ev := es.E
132 ephemeral := kind.IsEphemeral(ev.Kind)
133
134 var result *pipeline.Result
135 if req.Verified {
136 result = pipe.IngestPostVerify(ev)
137 } else {
138 result = pipe.Ingest(ev)
139 }
140 resp.EventID = ev.ID
141 resp.OK = result.OK
142 resp.Reason = result.Reason
143 if !result.OK {
144 return resp
145 }
146 resp.Bytes = req.Bytes
147 if !ephemeral {
148 resp.Seq = eng.MaxSerial()
149 }
150 reloadMute(eng, bl, ev)
151 return resp
152 }
153
154 // reloadMute refreshes the blacklist when the admin publishes a new mute list.
155 func reloadMute(eng *store.Engine, bl *mute.Blacklist, ev *event.E) {
156 if bl == nil || ev.Kind != kind.MuteList.K {
157 return
158 }
159 if !bytes.Equal(ev.Pubkey, bl.AdminPK()) {
160 return
161 }
162 bl.Load(eng)
163 bl.Purge(eng)
164 }
165
166 // history backfills one subscription with stored events that match, already
167 // marshaled as EVENT frames and already filtered for the connection's
168 // visibility.
169 func history(eng *store.Engine, req tree.Request) (resp tree.Response) {
170 resp.Op = tree.OpHistory
171 resp.ReqID = req.ReqID
172 resp.ConnID = req.ConnID
173 resp.SubID = req.SubID
174 resp.SeqHighWater = eng.MaxSerial()
175 resp.Done = true
176
177 _, rem, _ := envelope.Identify(req.Filter)
178 filter.ClearTaint()
179 var rq envelope.Req
180 if _, err := rq.Unmarshal(rem); err != nil || filter.IsTainted() {
181 return resp
182 }
183 if len(rq.Filters.F) == 0 {
184 return resp
185 }
186
187 limit := req.Limit
188 if limit <= 0 {
189 limit = 256
190 }
191 events := collect(eng, rq.Filters, limit)
192 resp.Events = [][]byte{:limit}
193 n := int32(0)
194 var frame []byte
195 for i := 0; i < len(events); i++ {
196 if n >= limit {
197 break
198 }
199 ev := events[i]
200 if req.Filtered && !access.CanSee(len(req.AuthedPubkey) > 0, req.AuthedPubkey, ev, req.NIP70, req.Marmot) {
201 continue
202 }
203 frame = marshalEvent(req.SubID, ev)
204 resp.Events[n] = frame
205 n++
206 }
207 resp.Events = resp.Events[:n]
208 return resp
209 }
210
211 // collect gathers the distinct stored events matching a filter set, in filter
212 // order. Search filters are run through the word index and then matched,
213 // because the word index answers with candidates only.
214 func collect(eng *store.Engine, filters filter.S, limit int32) (events []*event.E) {
215 seen := map[string]bool{}
216 var evs []*event.E
217 var err error
218 for _, f := range filters.F {
219 if len(f.Search) > 0 {
220 for _, ev := range eng.Search(f.Search, limit) {
221 if seen[string(ev.ID)] {
222 continue
223 }
224 seen[string(ev.ID)] = true
225 if f.Matches(ev) {
226 events = mxutil.Ensure(events, 1)
227 events = push(events, ev)
228 }
229 }
230 continue
231 }
232 evs, err = eng.QueryEvents(f)
233 if err != nil {
234 continue
235 }
236 for _, ev := range evs {
237 if seen[string(ev.ID)] {
238 continue
239 }
240 seen[string(ev.ID)] = true
241 events = mxutil.Ensure(events, 1)
242 events = push(events, ev)
243 }
244 }
245 return events
246 }
247
248 // marshalEvent renders one stored event as the EVENT frame for a subscription.
249 func marshalEvent(subID []byte, ev *event.E) (frame []byte) {
250 er := &envelope.EventResult{Subscription: subID, Event: ev}
251 return er.Marshal(nil)
252 }
253
254 // count answers a COUNT request.
255 func count(eng *store.Engine, req tree.Request) (resp tree.Response) {
256 resp.Op = tree.OpCount
257 resp.ReqID = req.ReqID
258 resp.ConnID = req.ConnID
259
260 _, rem, _ := envelope.Identify(req.Filter)
261 filter.ClearTaint()
262 var cr envelope.CountRequest
263 if _, err := cr.Unmarshal(rem); err != nil || filter.IsTainted() {
264 return resp
265 }
266 var total int32
267 var evs []*event.E
268 var qerr error
269 for _, f := range cr.Filters.F {
270 evs, qerr = eng.QueryEvents(f)
271 if qerr != nil {
272 continue
273 }
274 total += len(evs)
275 }
276 resp.SubID = cr.Subscription
277 resp.Count = total
278 resp.OK = true
279 return resp
280 }
281