1 """The Stage-A ingest worker pool: `ORLY_INGEST_WORKERS > 0`.
2 3 Default is 0, so the relay validates events inline and `wire/ingest_worker.mx`
4 never executes - a coverage run measured 0 of its 32 statements, and the pool
5 dispatch in `server.mx` (`dispatchToWorker`, `pollIngestWorkers`,
6 `completeIngestResponse`) likewise. This module starts a relay with the pool on
7 and drives every verdict the worker can return through it:
8 9 * accept - a valid EVENT is verified and stored
10 * ephemeral - kind 20001 is accepted and not stored
11 * reject - a tampered event fails Stage A (id mismatch)
12 * malformed - an EVENT envelope that is not an event object
13 14 The pool is a spawn tree: each worker is a forked domain, so this also fails if
15 a worker domain ever stops answering (the parent would block on the response).
16 17 The parent reads the response off a plain receive after the worker nudges a
18 zero-size ready channel. A select receive directly on a codec-framed spawn
19 channel does not decode: the runtime treats the encoded ring frame as a raw
20 element image and zeroes the buffer when its length differs from the element
21 size. `pollIngestWorkers` used that select shape, so `resp` arrived zeroed and
22 `completeIngestResponse` returned at its first guard with no reply and no store
23 write. The ready channel is the same shape the database-engine domain uses.
24 """
25 26 import time
27 28 from nostr_helpers import TEST_SECKEY, make_event
29 from test_relay_policy import (
30 HOST,
31 _drain_until,
32 _start_relay,
33 _stop_relay,
34 _ws_handshake,
35 _ws_recv_json,
36 _ws_send_json,
37 )
38 39 PORT = 24041
40 41 EPHEMERAL_KIND = 20001
42 43 44 def _pool_relay(name, workers=2):
45 """A relay with the ingest pool enabled, on its own port and data dir."""
46 return _start_relay(
47 PORT,
48 extra_env={"ORLY_INGEST_WORKERS": str(workers)},
49 data_dir_name=name,
50 )
51 52 53 def _query_ids(sock, sub, ids):
54 """REQ by id; return the events received before EOSE.
55 56 Reading to EOSE explicitly (rather than draining until one arrives) is what
57 makes a negative assertion mean something: an event that should not have
58 been stored has to show up in the returned list, not be skipped over.
59 """
60 _ws_send_json(sock, ["REQ", sub, {"ids": ids}])
61 events = []
62 deadline = time.time() + 15
63 while time.time() < deadline:
64 msg = _ws_recv_json(sock, timeout=10)
65 if msg[0] == "EVENT":
66 events.append(msg[2])
67 elif msg[0] == "EOSE":
68 return events
69 raise TimeoutError("no EOSE for " + sub)
70 71 72 class TestIngestWorkerPool:
73 def test_valid_event_is_verified_by_a_worker_and_stored(self):
74 proc, _ = _pool_relay("ingest-accept")
75 try:
76 sock = _ws_handshake(HOST, PORT)
77 ev = make_event(TEST_SECKEY, "ingest pool accept")
78 _ws_send_json(sock, ["EVENT", ev])
79 ok = _drain_until(sock, "OK", timeout=15)
80 assert ok[1] == ev["id"]
81 assert ok[2] is True
82 83 # Stored, so the worker's acceptance reached the write path.
84 stored = _query_ids(sock, "pool-query", [ev["id"]])
85 assert [e["id"] for e in stored] == [ev["id"]]
86 finally:
87 _stop_relay(proc)
88 89 def test_ephemeral_event_is_accepted_and_not_stored(self):
90 proc, _ = _pool_relay("ingest-ephemeral")
91 try:
92 sock = _ws_handshake(HOST, PORT)
93 ev = make_event(TEST_SECKEY, "ephemeral through the pool",
94 kind=EPHEMERAL_KIND)
95 _ws_send_json(sock, ["EVENT", ev])
96 ok = _drain_until(sock, "OK", timeout=15)
97 assert ok[1] == ev["id"]
98 assert ok[2] is True
99 100 assert _query_ids(sock, "eph-query", [ev["id"]]) == []
101 finally:
102 _stop_relay(proc)
103 104 def test_tampered_event_is_rejected_by_the_worker(self):
105 proc, _ = _pool_relay("ingest-reject")
106 try:
107 sock = _ws_handshake(HOST, PORT)
108 ev = make_event(TEST_SECKEY, "tampered after signing")
109 # Change the content without re-signing: the id no longer matches
110 # the serialised event, which Stage A checks before verifying.
111 ev["content"] = "tampered after signing!"
112 _ws_send_json(sock, ["EVENT", ev])
113 ok = _drain_until(sock, "OK", timeout=15)
114 assert ok[2] is False
115 assert "id mismatch" in ok[3]
116 assert _query_ids(sock, "rej-query", [ev["id"]]) == []
117 finally:
118 _stop_relay(proc)
119 120 def test_malformed_event_envelope_is_rejected(self):
121 proc, _ = _pool_relay("ingest-malformed")
122 try:
123 sock = _ws_handshake(HOST, PORT)
124 _ws_send_json(sock, ["EVENT", {"id": "not-an-event"}])
125 ok = _drain_until(sock, "OK", timeout=15)
126 assert ok[2] is False
127 assert "malformed EVENT" in ok[3]
128 finally:
129 _stop_relay(proc)
130 131 def test_single_worker_pool_answers_every_event(self):
132 # One worker is the smallest pool: it exercises dispatch falling back to
133 # the inline path when the only worker is busy, and the response drain
134 # with no second slot.
135 proc, _ = _pool_relay("ingest-single", workers=1)
136 try:
137 sock = _ws_handshake(HOST, PORT)
138 first = make_event(TEST_SECKEY, "single worker one")
139 second = make_event(TEST_SECKEY, "single worker two")
140 _ws_send_json(sock, ["EVENT", first])
141 _ws_send_json(sock, ["EVENT", second])
142 seen = {}
143 deadline = time.time() + 20
144 while len(seen) < 2 and time.time() < deadline:
145 ok = _drain_until(sock, "OK", timeout=15)
146 if ok[2]:
147 seen[ok[1]] = True
148 assert first["id"] in seen
149 assert second["id"] in seen
150 151 stored = _query_ids(sock, "single-query", [first["id"], second["id"]])
152 assert sorted(e["id"] for e in stored) == sorted(
153 [first["id"], second["id"]]
154 )
155 finally:
156 _stop_relay(proc)
157