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