// Package pipeline provides the event ingestion pipeline for the relay. // Validates, verifies, checks ACL, rate-limits, handles special kinds, // and stores events. Returns a Result suitable for OK envelope responses. package pipeline import ( "bytes" "strconv" "time" "git.smesh.lol/morly/pkg/acl" "git.smesh.lol/nostr/pkg/event" "git.smesh.lol/nostr/pkg/filter" "git.smesh.lol/nostr/pkg/kind" "git.smesh.lol/nostr/pkg/tag" "git.smesh.lol/morly/pkg/metrics" "git.smesh.lol/morly/pkg/relay/protowire" "git.smesh.lol/morly/pkg/relay/ratelimit" "git.smesh.lol/morly/pkg/store" "git.smesh.lol/musiquay/pkg/protocol" ) // Result is the outcome of event ingestion, maps to the OK envelope. type Result struct { OK bool Reason []byte } func accepted() (r *Result) { return &Result{OK: true} } func rejected(reason string) (r *Result) { return &Result{Reason: []byte(reason)} } func rejectedB(reason []byte) (r *Result) { return &Result{Reason: reason} } // Config holds pipeline limits. type Config struct { MaxFuture int64 // max seconds ahead of now (default 900) MaxPast int64 // max seconds behind now (0 = unlimited) MaxContent int32 // max content bytes (default 70000) MaxTags int32 // max tags per event (default 2000) MaxTagElem int32 // max bytes per tag element (default 1024) } // DefaultConfig returns sensible defaults. func DefaultConfig() (c Config) { return Config{ MaxFuture: 900, MaxContent: 70000, MaxTags: 2000, MaxTagElem: 1024, } } // Pipeline processes incoming events through the full ingestion path. type Pipeline struct { store *store.Engine acl acl.Checker limiter *ratelimit.Limiter cfg Config } // New creates a pipeline. Limiter may be nil to disable rate limiting. func New(s *store.Engine, a acl.Checker, l *ratelimit.Limiter, cfg Config) (p *Pipeline) { return &Pipeline{store: s, acl: a, limiter: l, cfg: cfg} } // SetACL updates the ACL checker. func (p *Pipeline) SetACL(a acl.Checker) { p.acl = a } // StageA runs schema validation + signature verification on a parsed event. // Returns nil on success, or a rejection Result on failure. // Used by both the synchronous path (Ingest) and async worker domains // (IngestWorker) so the logic is implemented exactly once. func StageA(ev *event.E) (r *Result) { if len(ev.ID) != 32 { return &Result{Reason: []byte("invalid: id must be 32 bytes")} } if len(ev.Pubkey) != 32 { return &Result{Reason: []byte("invalid: pubkey must be 32 bytes")} } if len(ev.Sig) != 64 { return &Result{Reason: []byte("invalid: sig must be 64 bytes")} } if !bytes.Equal(ev.ID, ev.GetIDBytes()) { return &Result{Reason: []byte("invalid: id mismatch")} } verifyStart := metrics.Now() valid, _ := ev.Verify() metrics.SigVerifyNs.Observe(metrics.Since(verifyStart)) if !valid { return &Result{Reason: []byte("invalid: bad signature")} } return nil } // Ingest validates, verifies, and stores an event. Full single-threaded // path (Stage A + Stage B). Used when no worker pool is configured. func (p *Pipeline) Ingest(ev *event.E) (r *Result) { ingestStart := metrics.Now() defer func() { metrics.IngestPipelineNs.Observe(metrics.Since(ingestStart)) }() if res := StageA(ev); res != nil { return res } return p.IngestPostVerify(ev) } // IngestPostVerify runs Stage B: timestamp/content limits, ACL, rate limit, // expiration, replaceable resolution, save. The caller is responsible for // having already run schema validation and signature verification (Stage A), // typically in a worker domain so verify is parallelized. func (p *Pipeline) IngestPostVerify(ev *event.E) (r *Result) { ingestStart := metrics.Now() defer func() { metrics.IngestPipelineNs.Observe(metrics.Since(ingestStart)) }() return p.ingestStageB(ev) } func (p *Pipeline) ingestStageB(ev *event.E) (r *Result) { // 2. Timestamp and content limits (checked in both sync and worker paths). if msg := p.validateLimits(ev); msg != nil { return rejectedB(msg) } // 2b. Musiquay payload validation. The gate is IsMusiquay, not every kind // ValidateEvent can parse: the kinds this vocabulary reuses from a NIP // (comment 1111, calendar 31923, listing 30402, ...) are shared with // clients that never heard of the protocol, and enforcing Musiquay's // stricter parser on them would turn their valid traffic away. The kinds // the protocol defines itself are the relay's to police. A Musiquay kind // with no parser yet returns nil below, so it is stored rather than // rejected by a rule that does not exist. if protocol.IsMusiquay(ev.Kind) { if verr := protowire.ValidateEvent(ev); verr != nil { return rejected("invalid: " | verr.Error()) } } // 3. ACL. aclStart := metrics.Now() allowed := p.acl.AllowWrite(ev.Pubkey, ev.Kind) metrics.ACLAllowWriteNs.Observe(metrics.Since(aclStart)) if !allowed { return rejected("blocked: pubkey not allowed") } // 4. Rate limit. if p.limiter != nil && !p.limiter.Allow(ev.Pubkey) { return rejected("rate-limited: slow down") } // 5. Expiration check (NIP-40). if isExpired(ev) { return rejected("invalid: event expired") } // 6. Ephemeral - accept but don't store. if kind.IsEphemeral(ev.Kind) { return accepted() } // 7. Replaceable kinds. if kind.IsReplaceable(ev.Kind) { return p.handleReplaceable(ev) } // Musiquay's counted kinds sit inside the parameterized-replaceable range // but are regular records, not addresses: each receipt, aggregate and // rollup is stored. Without this the range rule keeps one per kind and // author (all of them have no `d` tag) and deletes every other. if countedInReplaceableRange(ev.Kind) { return p.saveOrDup(ev) } if kind.IsParameterizedReplaceable(ev.Kind) { return p.handleParamReplaceable(ev) } // 8. Deletion (NIP-09). if ev.Kind == kind.EventDeletion.K { return p.handleDeletion(ev) } // 9. Store regular event. return p.saveOrDup(ev) } // --- validation --- // validateLimits checks timestamp range and content/tag size limits. // Schema + sig checks are in StageA (already run before this). func (p *Pipeline) validateLimits(ev *event.E) (buf []byte) { if ev.CreatedAt == 0 { return []byte("invalid: created_at missing") } now := time.Now().Unix() if p.cfg.MaxFuture > 0 && ev.CreatedAt > now+p.cfg.MaxFuture { return []byte("invalid: created_at too far in future") } if p.cfg.MaxPast > 0 && ev.CreatedAt < now-p.cfg.MaxPast { return []byte("invalid: created_at too far in past") } if p.cfg.MaxContent > 0 && len(ev.Content) > p.cfg.MaxContent { return []byte("invalid: content too large") } if ev.Tags != nil { if p.cfg.MaxTags > 0 && ev.Tags.Len() > p.cfg.MaxTags { return []byte("invalid: too many tags") } if p.cfg.MaxTagElem > 0 { for _, tg := range ev.Tags.T { for _, elem := range tg.T { if len(elem) > p.cfg.MaxTagElem { return []byte("invalid: tag element too large") } } } } } return nil } // --- expiration (NIP-40) --- func isExpired(ev *event.E) (ok bool) { if ev.Tags == nil { return false } t := ev.Tags.GetFirst([]byte("expiration")) if t == nil || t.Len() < 2 { return false } exp, err := strconv.ParseInt(string(t.Value()), 10, 64) if err != nil || exp <= 0 { return false } return time.Now().Unix() > exp } // --- replaceable events --- // Tiebreaker when timestamps are equal: keep the event with the // lexicographically lower ID (deterministic on SHA-256 hashes). // NIP-01 does not specify tiebreaker behavior. func (p *Pipeline) handleReplaceable(ev *event.E) (r *Result) { f := &filter.F{ Kinds: kind.NewS(kind.New(ev.Kind)), Authors: tag.NewFromBytesSlice(ev.Pubkey), } limit := uint32(1) f.Limit = &limit existing, err := p.store.QueryEvents(f) if err == nil && len(existing) > 0 { old := existing[0] if old.CreatedAt > ev.CreatedAt { return rejected("duplicate: newer version exists") } if old.CreatedAt == ev.CreatedAt && bytes.Compare(old.ID, ev.ID) < 0 { return rejected("duplicate: event with lower id exists at same timestamp") } p.store.DeleteEvent(old.ID) } return p.saveOrDup(ev) } func (p *Pipeline) handleParamReplaceable(ev *event.E) (r *Result) { dVal := dTagValue(ev) f := &filter.F{ Kinds: kind.NewS(kind.New(ev.Kind)), Authors: tag.NewFromBytesSlice(ev.Pubkey), } existing, err := p.store.QueryEvents(f) if err == nil { for _, old := range existing { if !bytes.Equal(dTagValue(old), dVal) { continue } if old.CreatedAt > ev.CreatedAt { return rejected("duplicate: newer version exists") } if old.CreatedAt == ev.CreatedAt && bytes.Compare(old.ID, ev.ID) < 0 { return rejected("duplicate: event with lower id exists at same timestamp") } p.store.DeleteEvent(old.ID) break } } return p.saveOrDup(ev) } func dTagValue(ev *event.E) (buf []byte) { if ev.Tags == nil { return nil } t := ev.Tags.GetFirst([]byte("d")) if t == nil || t.Len() < 2 { return nil } return t.Value() } // --- deletion (NIP-09) --- func (p *Pipeline) handleDeletion(ev *event.E) (r *Result) { if ev.Tags == nil { return p.saveOrDup(ev) } eTags := ev.Tags.GetAll([]byte("e")) for _, et := range eTags { targetID := et.ValueBinary() if targetID == nil || len(targetID) != 32 { continue } // Only delete events owned by the same pubkey. tf := &filter.F{Ids: tag.NewFromBytesSlice(targetID)} targets, err := p.store.QueryEvents(tf) if err != nil || len(targets) == 0 { continue } if !bytes.Equal(targets[0].Pubkey, ev.Pubkey) { continue } p.store.DeleteEvent(targetID) } return p.saveOrDup(ev) } // --- storage helper --- func (p *Pipeline) saveOrDup(ev *event.E) (r *Result) { err := p.store.SaveEvent(ev) if err == nil { return accepted() } msg := err.Error() if len(msg) >= 9 && msg[:9] == "duplicate" { return rejected("duplicate: already have this event") } return rejectedB([]byte("error: ") | msg) }