ingest_worker.mx raw

   1  package wire
   2  
   3  import (
   4  	"git.smesh.lol/nostr/pkg/envelope"
   5  	"git.smesh.lol/nostr/pkg/event"
   6  	"git.smesh.lol/nostr/pkg/kind"
   7  	"git.smesh.lol/morly/pkg/relay/pipeline"
   8  )
   9  
  10  // IngestWorker is the spawn target for a Stage-A ingest worker. Reads
  11  // IngestRequest frames, runs schema validation + sig verify + ephemeral
  12  // fast-reject, writes IngestResponse frames. Loops forever.
  13  //
  14  // Stage A is intentionally narrow: no I/O, no shared state. The parent's
  15  // net domain owns Stage B (dedup, replaceable resolution, WAL append).
  16  // Workers are pure compute, parallelizable across CPUs via fork.
  17  //
  18  // The reply is announced on ready after it is written to out. A select
  19  // receive on a codec-framed spawn channel does not decode (see the DB
  20  // handle in pkg/relay/server), so the parent must learn readiness on a
  21  // zero-size channel and then read the codec channel with a plain receive.
  22  func IngestWorker(in chan IngestRequest, out chan IngestResponse, ready chan struct{}) {
  23  	for {
  24  		req, ok := <-in
  25  		if !ok {
  26  			return
  27  		}
  28  		out <- ingestProcessOne(req)
  29  		ready <- struct{}{}
  30  	}
  31  }
  32  
  33  func ingestProcessOne(req IngestRequest) (i IngestResponse) {
  34  	resp := IngestResponse{
  35  		ReqID: req.ReqID,
  36  		Bytes: req.Bytes,
  37  	}
  38  	label, rem, _ := envelope.Identify(req.Bytes)
  39  	if label != envelope.EventLabel {
  40  		resp.Verdict = VerdictReject
  41  		resp.Reason = []byte("invalid: not an EVENT envelope")
  42  		return resp
  43  	}
  44  	var es envelope.EventSubmission
  45  	if _, err := es.Unmarshal(rem); err != nil || es.E == nil {
  46  		resp.Verdict = VerdictReject
  47  		resp.Reason = []byte("invalid: malformed EVENT")
  48  		return resp
  49  	}
  50  	ev := es.E
  51  
  52  	if r := pipeline.StageA(ev); r != nil {
  53  		resp.Verdict = VerdictReject
  54  		resp.Reason = r.Reason
  55  		return resp
  56  	}
  57  
  58  	resp.Kind = ev.Kind
  59  	resp.CreatedAt = ev.CreatedAt
  60  	copy(resp.Pubkey[:], ev.Pubkey)
  61  	copy(resp.EventID[:], ev.ID)
  62  
  63  	if kind.IsEphemeral(ev.Kind) {
  64  		resp.Verdict = VerdictEphemeral
  65  		return resp
  66  	}
  67  
  68  	resp.Verdict = VerdictAccept
  69  	return resp
  70  }
  71  
  72  // keep event.E imported for type-checking stability
  73  var _ *event.E
  74