// Package dbengine is the relay's database-engine domain: the one node that // opens the store, and therefore the single writer and the single reader. // // Every other node reaches storage by sending it a tree.Request. It answers on // tree.Response and nudges the parent on a zero-size channel (see the note on // tree.Request for why the parent cannot select on the codec channel itself). // // The engine owns everything that is a function of stored state: the storage // engine, the ingestion pipeline (limits, ACL, replaceable and deletion // resolution, WAL append), and the admin mute blacklist. Root owns everything // that is a function of connections: transport, subscriptions, and the // authentication and rate policy that decides whether a frame is even worth // sending here. package dbengine import ( "bytes" "fmt" "os" "git.smesh.lol/morly/pkg/access" "git.smesh.lol/morly/pkg/acl" "git.smesh.lol/moxie/pkg/mxutil" "git.smesh.lol/nostr/pkg/envelope" "git.smesh.lol/nostr/pkg/event" "git.smesh.lol/nostr/pkg/filter" "git.smesh.lol/nostr/pkg/kind" "git.smesh.lol/morly/pkg/relay/mute" "git.smesh.lol/morly/pkg/relay/pipeline" "git.smesh.lol/morly/pkg/relay/tree" "git.smesh.lol/morly/pkg/store" ) // Config is what the database node needs to open its store and build its // policy. It is deliberately not the whole relay config: the relay config // belongs to root, and only these values are a function of stored state. type Config struct { DataDir string ACLMode string Admins []string FollowListFreqSec int32 SocialWoTMaxDepth int32 SocialWoTRefreshSec int32 MuteBlacklist string } // Run is the spawn target of the database-engine domain. It opens the store, // then serves requests until the inbox closes. func Run(cfg Config, in chan tree.Request, out chan tree.Response, ready chan struct{}) { eng, err := store.Open(cfg.DataDir | "/db") if err != nil { fmt.Fprintln(os.Stderr, "store error: "|err.Error()) return } checker := buildACL(cfg, eng) pipe := pipeline.New(eng, checker, nil, pipeline.DefaultConfig()) bl := loadMute(cfg.MuteBlacklist, eng) for { req, ok := <-in if !ok { break } resp := serve(eng, pipe, bl, req) out <- resp ready <- struct{}{} } eng.Flush() eng.Close() } // buildACL constructs the write-access checker. The checker answers from // stored follow and relay-list events, so it lives on this side of the // boundary. func buildACL(cfg Config, eng *store.Engine) (c acl.Checker) { switch cfg.ACLMode { case "follows": return acl.NewFollows(eng, cfg.Admins, cfg.FollowListFreqSec) case "social": return acl.NewSocial(eng, cfg.Admins, cfg.SocialWoTMaxDepth, cfg.SocialWoTRefreshSec) } return &acl.Open{} } // loadMute loads the admin mute blacklist and purges what it blocks. Returns // nil when no admin pubkey is configured. func loadMute(adminHexPK string, eng *store.Engine) (b *mute.Blacklist) { b = mute.New(adminHexPK) if b == nil { return nil } b.Load(eng) b.Purge(eng) return b } // serve runs one request. A free function: its scratch dies at return instead // of accumulating in Run's frame for the life of the relay. func serve(eng *store.Engine, pipe *pipeline.Pipeline, bl *mute.Blacklist, req tree.Request) (resp tree.Response) { switch req.Op { case tree.OpPersist: return persist(eng, pipe, bl, req) case tree.OpHistory: return history(eng, req) case tree.OpCount: return count(eng, req) } return tree.Response{} } // persist stores one event, or reports why it was not stored. // // Ephemeral events pass the same limits, ACL and expiration checks as stored // ones, and are then accepted without a WAL append: they are live-only, so they // carry no sequence number and root still broadcasts them. func persist(eng *store.Engine, pipe *pipeline.Pipeline, bl *mute.Blacklist, req tree.Request) (resp tree.Response) { resp.Op = tree.OpPersist resp.ReqID = req.ReqID resp.ConnID = req.ConnID label, rem, _ := envelope.Identify(req.Bytes) if label != envelope.EventLabel { resp.Reason = []byte("invalid: not an EVENT envelope") return resp } var es envelope.EventSubmission if _, perr := es.Unmarshal(rem); perr != nil || es.E == nil { resp.Reason = []byte("invalid: parse error") return resp } ev := es.E ephemeral := kind.IsEphemeral(ev.Kind) var result *pipeline.Result if req.Verified { result = pipe.IngestPostVerify(ev) } else { result = pipe.Ingest(ev) } resp.EventID = ev.ID resp.OK = result.OK resp.Reason = result.Reason if !result.OK { return resp } resp.Bytes = req.Bytes if !ephemeral { resp.Seq = eng.MaxSerial() } reloadMute(eng, bl, ev) return resp } // reloadMute refreshes the blacklist when the admin publishes a new mute list. func reloadMute(eng *store.Engine, bl *mute.Blacklist, ev *event.E) { if bl == nil || ev.Kind != kind.MuteList.K { return } if !bytes.Equal(ev.Pubkey, bl.AdminPK()) { return } bl.Load(eng) bl.Purge(eng) } // history backfills one subscription with stored events that match, already // marshaled as EVENT frames and already filtered for the connection's // visibility. func history(eng *store.Engine, req tree.Request) (resp tree.Response) { resp.Op = tree.OpHistory resp.ReqID = req.ReqID resp.ConnID = req.ConnID resp.SubID = req.SubID resp.SeqHighWater = eng.MaxSerial() resp.Done = true _, rem, _ := envelope.Identify(req.Filter) filter.ClearTaint() var rq envelope.Req if _, err := rq.Unmarshal(rem); err != nil || filter.IsTainted() { return resp } if len(rq.Filters.F) == 0 { return resp } limit := req.Limit if limit <= 0 { limit = 256 } events := collect(eng, rq.Filters, limit) resp.Events = [][]byte{:limit} n := int32(0) var frame []byte for i := 0; i < len(events); i++ { if n >= limit { break } ev := events[i] if req.Filtered && !access.CanSee(len(req.AuthedPubkey) > 0, req.AuthedPubkey, ev, req.NIP70, req.Marmot) { continue } frame = marshalEvent(req.SubID, ev) resp.Events[n] = frame n++ } resp.Events = resp.Events[:n] return resp } // collect gathers the distinct stored events matching a filter set, in filter // order. Search filters are run through the word index and then matched, // because the word index answers with candidates only. func collect(eng *store.Engine, filters filter.S, limit int32) (events []*event.E) { seen := map[string]bool{} var evs []*event.E var err error for _, f := range filters.F { if len(f.Search) > 0 { for _, ev := range eng.Search(f.Search, limit) { if seen[string(ev.ID)] { continue } seen[string(ev.ID)] = true if f.Matches(ev) { events = mxutil.Ensure(events, 1) events = push(events, ev) } } continue } evs, err = eng.QueryEvents(f) if err != nil { continue } for _, ev := range evs { if seen[string(ev.ID)] { continue } seen[string(ev.ID)] = true events = mxutil.Ensure(events, 1) events = push(events, ev) } } return events } // marshalEvent renders one stored event as the EVENT frame for a subscription. func marshalEvent(subID []byte, ev *event.E) (frame []byte) { er := &envelope.EventResult{Subscription: subID, Event: ev} return er.Marshal(nil) } // count answers a COUNT request. func count(eng *store.Engine, req tree.Request) (resp tree.Response) { resp.Op = tree.OpCount resp.ReqID = req.ReqID resp.ConnID = req.ConnID _, rem, _ := envelope.Identify(req.Filter) filter.ClearTaint() var cr envelope.CountRequest if _, err := cr.Unmarshal(rem); err != nil || filter.IsTainted() { return resp } var total int32 var evs []*event.E var qerr error for _, f := range cr.Filters.F { evs, qerr = eng.QueryEvents(f) if qerr != nil { continue } total += len(evs) } resp.SubID = cr.Subscription resp.Count = total resp.OK = true return resp }