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/smesh/pkg/access"
  22  	"git.smesh.lol/smesh/pkg/acl"
  23  	"git.smesh.lol/moxie/pkg/mxutil"
  24  	"git.smesh.lol/smesh/pkg/nostr/envelope"
  25  	"git.smesh.lol/smesh/pkg/nostr/event"
  26  	"git.smesh.lol/smesh/pkg/nostr/filter"
  27  	"git.smesh.lol/smesh/pkg/nostr/kind"
  28  	"git.smesh.lol/smesh/pkg/relay/mute"
  29  	"git.smesh.lol/smesh/pkg/relay/pipeline"
  30  	"git.smesh.lol/smesh/pkg/relay/tree"
  31  	"git.smesh.lol/smesh/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.Tainted = false
 179  	var rq envelope.Req
 180  	if _, err := rq.Unmarshal(rem); err != nil || filter.Tainted {
 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.Tainted = false
 262  	var cr envelope.CountRequest
 263  	if _, err := cr.Unmarshal(rem); err != nil || filter.Tainted {
 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