package wire import ( "git.smesh.lol/nostr/pkg/envelope" "git.smesh.lol/nostr/pkg/event" "git.smesh.lol/nostr/pkg/kind" "git.smesh.lol/morly/pkg/relay/pipeline" ) // IngestWorker is the spawn target for a Stage-A ingest worker. Reads // IngestRequest frames, runs schema validation + sig verify + ephemeral // fast-reject, writes IngestResponse frames. Loops forever. // // Stage A is intentionally narrow: no I/O, no shared state. The parent's // net domain owns Stage B (dedup, replaceable resolution, WAL append). // Workers are pure compute, parallelizable across CPUs via fork. // // The reply is announced on ready after it is written to out. A select // receive on a codec-framed spawn channel does not decode (see the DB // handle in pkg/relay/server), so the parent must learn readiness on a // zero-size channel and then read the codec channel with a plain receive. func IngestWorker(in chan IngestRequest, out chan IngestResponse, ready chan struct{}) { for { req, ok := <-in if !ok { return } out <- ingestProcessOne(req) ready <- struct{}{} } } func ingestProcessOne(req IngestRequest) (i IngestResponse) { resp := IngestResponse{ ReqID: req.ReqID, Bytes: req.Bytes, } label, rem, _ := envelope.Identify(req.Bytes) if label != envelope.EventLabel { resp.Verdict = VerdictReject resp.Reason = []byte("invalid: not an EVENT envelope") return resp } var es envelope.EventSubmission if _, err := es.Unmarshal(rem); err != nil || es.E == nil { resp.Verdict = VerdictReject resp.Reason = []byte("invalid: malformed EVENT") return resp } ev := es.E if r := pipeline.StageA(ev); r != nil { resp.Verdict = VerdictReject resp.Reason = r.Reason return resp } resp.Kind = ev.Kind resp.CreatedAt = ev.CreatedAt copy(resp.Pubkey[:], ev.Pubkey) copy(resp.EventID[:], ev.ID) if kind.IsEphemeral(ev.Kind) { resp.Verdict = VerdictEphemeral return resp } resp.Verdict = VerdictAccept return resp } // keep event.E imported for type-checking stability var _ *event.E