pipeline.mx raw

   1  // Package pipeline provides the event ingestion pipeline for the relay.
   2  // Validates, verifies, checks ACL, rate-limits, handles special kinds,
   3  // and stores events. Returns a Result suitable for OK envelope responses.
   4  package pipeline
   5  
   6  import (
   7  	"bytes"
   8  	"strconv"
   9  	"time"
  10  
  11  	"git.smesh.lol/morly/pkg/acl"
  12  	"git.smesh.lol/nostr/pkg/event"
  13  	"git.smesh.lol/nostr/pkg/filter"
  14  	"git.smesh.lol/nostr/pkg/kind"
  15  	"git.smesh.lol/nostr/pkg/tag"
  16  	"git.smesh.lol/morly/pkg/metrics"
  17  	"git.smesh.lol/morly/pkg/relay/protowire"
  18  	"git.smesh.lol/morly/pkg/relay/ratelimit"
  19  	"git.smesh.lol/morly/pkg/store"
  20  	"git.smesh.lol/musiquay/pkg/protocol"
  21  )
  22  
  23  // Result is the outcome of event ingestion, maps to the OK envelope.
  24  type Result struct {
  25  	OK     bool
  26  	Reason []byte
  27  }
  28  
  29  func accepted() (r *Result) { return &Result{OK: true} }
  30  func rejected(reason string) (r *Result) { return &Result{Reason: []byte(reason)} }
  31  func rejectedB(reason []byte) (r *Result) { return &Result{Reason: reason} }
  32  
  33  // Config holds pipeline limits.
  34  type Config struct {
  35  	MaxFuture  int64 // max seconds ahead of now (default 900)
  36  	MaxPast    int64 // max seconds behind now (0 = unlimited)
  37  	MaxContent int32   // max content bytes (default 70000)
  38  	MaxTags    int32   // max tags per event (default 2000)
  39  	MaxTagElem int32   // max bytes per tag element (default 1024)
  40  }
  41  
  42  // DefaultConfig returns sensible defaults.
  43  func DefaultConfig() (c Config) {
  44  	return Config{
  45  		MaxFuture:  900,
  46  		MaxContent: 70000,
  47  		MaxTags:    2000,
  48  		MaxTagElem: 1024,
  49  	}
  50  }
  51  
  52  // Pipeline processes incoming events through the full ingestion path.
  53  type Pipeline struct {
  54  	store   *store.Engine
  55  	acl     acl.Checker
  56  	limiter *ratelimit.Limiter
  57  	cfg     Config
  58  }
  59  
  60  // New creates a pipeline. Limiter may be nil to disable rate limiting.
  61  func New(s *store.Engine, a acl.Checker, l *ratelimit.Limiter, cfg Config) (p *Pipeline) {
  62  	return &Pipeline{store: s, acl: a, limiter: l, cfg: cfg}
  63  }
  64  
  65  // SetACL updates the ACL checker.
  66  func (p *Pipeline) SetACL(a acl.Checker) { p.acl = a }
  67  
  68  // StageA runs schema validation + signature verification on a parsed event.
  69  // Returns nil on success, or a rejection Result on failure.
  70  // Used by both the synchronous path (Ingest) and async worker domains
  71  // (IngestWorker) so the logic is implemented exactly once.
  72  func StageA(ev *event.E) (r *Result) {
  73  	if len(ev.ID) != 32 {
  74  		return &Result{Reason: []byte("invalid: id must be 32 bytes")}
  75  	}
  76  	if len(ev.Pubkey) != 32 {
  77  		return &Result{Reason: []byte("invalid: pubkey must be 32 bytes")}
  78  	}
  79  	if len(ev.Sig) != 64 {
  80  		return &Result{Reason: []byte("invalid: sig must be 64 bytes")}
  81  	}
  82  	if !bytes.Equal(ev.ID, ev.GetIDBytes()) {
  83  		return &Result{Reason: []byte("invalid: id mismatch")}
  84  	}
  85  	verifyStart := metrics.Now()
  86  	valid, _ := ev.Verify()
  87  	metrics.SigVerifyNs.Observe(metrics.Since(verifyStart))
  88  	if !valid {
  89  		return &Result{Reason: []byte("invalid: bad signature")}
  90  	}
  91  	return nil
  92  }
  93  
  94  // Ingest validates, verifies, and stores an event. Full single-threaded
  95  // path (Stage A + Stage B). Used when no worker pool is configured.
  96  func (p *Pipeline) Ingest(ev *event.E) (r *Result) {
  97  	ingestStart := metrics.Now()
  98  	defer func() { metrics.IngestPipelineNs.Observe(metrics.Since(ingestStart)) }()
  99  	if res := StageA(ev); res != nil {
 100  		return res
 101  	}
 102  	return p.IngestPostVerify(ev)
 103  }
 104  
 105  // IngestPostVerify runs Stage B: timestamp/content limits, ACL, rate limit,
 106  // expiration, replaceable resolution, save. The caller is responsible for
 107  // having already run schema validation and signature verification (Stage A),
 108  // typically in a worker domain so verify is parallelized.
 109  func (p *Pipeline) IngestPostVerify(ev *event.E) (r *Result) {
 110  	ingestStart := metrics.Now()
 111  	defer func() { metrics.IngestPipelineNs.Observe(metrics.Since(ingestStart)) }()
 112  	return p.ingestStageB(ev)
 113  }
 114  
 115  func (p *Pipeline) ingestStageB(ev *event.E) (r *Result) {
 116  	// 2. Timestamp and content limits (checked in both sync and worker paths).
 117  	if msg := p.validateLimits(ev); msg != nil {
 118  		return rejectedB(msg)
 119  	}
 120  
 121  	// 2b. Musiquay payload validation. The gate is IsMusiquay, not every kind
 122  	// ValidateEvent can parse: the kinds this vocabulary reuses from a NIP
 123  	// (comment 1111, calendar 31923, listing 30402, ...) are shared with
 124  	// clients that never heard of the protocol, and enforcing Musiquay's
 125  	// stricter parser on them would turn their valid traffic away. The kinds
 126  	// the protocol defines itself are the relay's to police. A Musiquay kind
 127  	// with no parser yet returns nil below, so it is stored rather than
 128  	// rejected by a rule that does not exist.
 129  	if protocol.IsMusiquay(ev.Kind) {
 130  		if verr := protowire.ValidateEvent(ev); verr != nil {
 131  			return rejected("invalid: " | verr.Error())
 132  		}
 133  	}
 134  
 135  	// 3. ACL.
 136  	aclStart := metrics.Now()
 137  	allowed := p.acl.AllowWrite(ev.Pubkey, ev.Kind)
 138  	metrics.ACLAllowWriteNs.Observe(metrics.Since(aclStart))
 139  	if !allowed {
 140  		return rejected("blocked: pubkey not allowed")
 141  	}
 142  
 143  	// 4. Rate limit.
 144  	if p.limiter != nil && !p.limiter.Allow(ev.Pubkey) {
 145  		return rejected("rate-limited: slow down")
 146  	}
 147  
 148  	// 5. Expiration check (NIP-40).
 149  	if isExpired(ev) {
 150  		return rejected("invalid: event expired")
 151  	}
 152  
 153  	// 6. Ephemeral - accept but don't store.
 154  	if kind.IsEphemeral(ev.Kind) {
 155  		return accepted()
 156  	}
 157  
 158  	// 7. Replaceable kinds.
 159  	if kind.IsReplaceable(ev.Kind) {
 160  		return p.handleReplaceable(ev)
 161  	}
 162  	// Musiquay's counted kinds sit inside the parameterized-replaceable range
 163  	// but are regular records, not addresses: each receipt, aggregate and
 164  	// rollup is stored. Without this the range rule keeps one per kind and
 165  	// author (all of them have no `d` tag) and deletes every other.
 166  	if countedInReplaceableRange(ev.Kind) {
 167  		return p.saveOrDup(ev)
 168  	}
 169  	if kind.IsParameterizedReplaceable(ev.Kind) {
 170  		return p.handleParamReplaceable(ev)
 171  	}
 172  
 173  	// 8. Deletion (NIP-09).
 174  	if ev.Kind == kind.EventDeletion.K {
 175  		return p.handleDeletion(ev)
 176  	}
 177  
 178  	// 9. Store regular event.
 179  	return p.saveOrDup(ev)
 180  }
 181  
 182  // --- validation ---
 183  
 184  // validateLimits checks timestamp range and content/tag size limits.
 185  // Schema + sig checks are in StageA (already run before this).
 186  func (p *Pipeline) validateLimits(ev *event.E) (buf []byte) {
 187  	if ev.CreatedAt == 0 {
 188  		return []byte("invalid: created_at missing")
 189  	}
 190  
 191  	now := time.Now().Unix()
 192  	if p.cfg.MaxFuture > 0 && ev.CreatedAt > now+p.cfg.MaxFuture {
 193  		return []byte("invalid: created_at too far in future")
 194  	}
 195  	if p.cfg.MaxPast > 0 && ev.CreatedAt < now-p.cfg.MaxPast {
 196  		return []byte("invalid: created_at too far in past")
 197  	}
 198  
 199  	if p.cfg.MaxContent > 0 && len(ev.Content) > p.cfg.MaxContent {
 200  		return []byte("invalid: content too large")
 201  	}
 202  
 203  	if ev.Tags != nil {
 204  		if p.cfg.MaxTags > 0 && ev.Tags.Len() > p.cfg.MaxTags {
 205  			return []byte("invalid: too many tags")
 206  		}
 207  		if p.cfg.MaxTagElem > 0 {
 208  			for _, tg := range ev.Tags.T {
 209  				for _, elem := range tg.T {
 210  					if len(elem) > p.cfg.MaxTagElem {
 211  						return []byte("invalid: tag element too large")
 212  					}
 213  				}
 214  			}
 215  		}
 216  	}
 217  
 218  	return nil
 219  }
 220  
 221  // --- expiration (NIP-40) ---
 222  
 223  func isExpired(ev *event.E) (ok bool) {
 224  	if ev.Tags == nil {
 225  		return false
 226  	}
 227  	t := ev.Tags.GetFirst([]byte("expiration"))
 228  	if t == nil || t.Len() < 2 {
 229  		return false
 230  	}
 231  	exp, err := strconv.ParseInt(string(t.Value()), 10, 64)
 232  	if err != nil || exp <= 0 {
 233  		return false
 234  	}
 235  	return time.Now().Unix() > exp
 236  }
 237  
 238  // --- replaceable events ---
 239  // Tiebreaker when timestamps are equal: keep the event with the
 240  // lexicographically lower ID (deterministic on SHA-256 hashes).
 241  // NIP-01 does not specify tiebreaker behavior.
 242  
 243  func (p *Pipeline) handleReplaceable(ev *event.E) (r *Result) {
 244  	f := &filter.F{
 245  		Kinds:   kind.NewS(kind.New(ev.Kind)),
 246  		Authors: tag.NewFromBytesSlice(ev.Pubkey),
 247  	}
 248  	limit := uint32(1)
 249  	f.Limit = &limit
 250  
 251  	existing, err := p.store.QueryEvents(f)
 252  	if err == nil && len(existing) > 0 {
 253  		old := existing[0]
 254  		if old.CreatedAt > ev.CreatedAt {
 255  			return rejected("duplicate: newer version exists")
 256  		}
 257  		if old.CreatedAt == ev.CreatedAt && bytes.Compare(old.ID, ev.ID) < 0 {
 258  			return rejected("duplicate: event with lower id exists at same timestamp")
 259  		}
 260  		p.store.DeleteEvent(old.ID)
 261  	}
 262  	return p.saveOrDup(ev)
 263  }
 264  
 265  func (p *Pipeline) handleParamReplaceable(ev *event.E) (r *Result) {
 266  	dVal := dTagValue(ev)
 267  
 268  	f := &filter.F{
 269  		Kinds:   kind.NewS(kind.New(ev.Kind)),
 270  		Authors: tag.NewFromBytesSlice(ev.Pubkey),
 271  	}
 272  	existing, err := p.store.QueryEvents(f)
 273  	if err == nil {
 274  		for _, old := range existing {
 275  			if !bytes.Equal(dTagValue(old), dVal) {
 276  				continue
 277  			}
 278  			if old.CreatedAt > ev.CreatedAt {
 279  				return rejected("duplicate: newer version exists")
 280  			}
 281  			if old.CreatedAt == ev.CreatedAt && bytes.Compare(old.ID, ev.ID) < 0 {
 282  				return rejected("duplicate: event with lower id exists at same timestamp")
 283  			}
 284  			p.store.DeleteEvent(old.ID)
 285  			break
 286  		}
 287  	}
 288  	return p.saveOrDup(ev)
 289  }
 290  
 291  func dTagValue(ev *event.E) (buf []byte) {
 292  	if ev.Tags == nil {
 293  		return nil
 294  	}
 295  	t := ev.Tags.GetFirst([]byte("d"))
 296  	if t == nil || t.Len() < 2 {
 297  		return nil
 298  	}
 299  	return t.Value()
 300  }
 301  
 302  // --- deletion (NIP-09) ---
 303  
 304  func (p *Pipeline) handleDeletion(ev *event.E) (r *Result) {
 305  	if ev.Tags == nil {
 306  		return p.saveOrDup(ev)
 307  	}
 308  	eTags := ev.Tags.GetAll([]byte("e"))
 309  	for _, et := range eTags {
 310  		targetID := et.ValueBinary()
 311  		if targetID == nil || len(targetID) != 32 {
 312  			continue
 313  		}
 314  		// Only delete events owned by the same pubkey.
 315  		tf := &filter.F{Ids: tag.NewFromBytesSlice(targetID)}
 316  		targets, err := p.store.QueryEvents(tf)
 317  		if err != nil || len(targets) == 0 {
 318  			continue
 319  		}
 320  		if !bytes.Equal(targets[0].Pubkey, ev.Pubkey) {
 321  			continue
 322  		}
 323  		p.store.DeleteEvent(targetID)
 324  	}
 325  	return p.saveOrDup(ev)
 326  }
 327  
 328  // --- storage helper ---
 329  
 330  func (p *Pipeline) saveOrDup(ev *event.E) (r *Result) {
 331  	err := p.store.SaveEvent(ev)
 332  	if err == nil {
 333  		return accepted()
 334  	}
 335  	msg := err.Error()
 336  	if len(msg) >= 9 && msg[:9] == "duplicate" {
 337  		return rejected("duplicate: already have this event")
 338  	}
 339  	return rejectedB([]byte("error: ") | msg)
 340  }
 341