"""The Stage-A ingest worker pool: `ORLY_INGEST_WORKERS > 0`. Default is 0, so the relay validates events inline and `wire/ingest_worker.mx` never executes - a coverage run measured 0 of its 32 statements, and the pool dispatch in `server.mx` (`dispatchToWorker`, `pollIngestWorkers`, `completeIngestResponse`) likewise. This module starts a relay with the pool on and drives every verdict the worker can return through it: * accept - a valid EVENT is verified and stored * ephemeral - kind 20001 is accepted and not stored * reject - a tampered event fails Stage A (id mismatch) * malformed - an EVENT envelope that is not an event object The pool is a spawn tree: each worker is a forked domain, so this also fails if a worker domain ever stops answering (the parent would block on the response). The parent reads the response off a plain receive after the worker nudges a zero-size ready channel. A select receive directly on a codec-framed spawn channel does not decode: the runtime treats the encoded ring frame as a raw element image and zeroes the buffer when its length differs from the element size. `pollIngestWorkers` used that select shape, so `resp` arrived zeroed and `completeIngestResponse` returned at its first guard with no reply and no store write. The ready channel is the same shape the database-engine domain uses. """ import time from nostr_helpers import TEST_SECKEY, make_event from test_relay_policy import ( HOST, _drain_until, _start_relay, _stop_relay, _ws_handshake, _ws_recv_json, _ws_send_json, ) PORT = 24041 EPHEMERAL_KIND = 20001 def _pool_relay(name, workers=2): """A relay with the ingest pool enabled, on its own port and data dir.""" return _start_relay( PORT, extra_env={"ORLY_INGEST_WORKERS": str(workers)}, data_dir_name=name, ) def _query_ids(sock, sub, ids): """REQ by id; return the events received before EOSE. Reading to EOSE explicitly (rather than draining until one arrives) is what makes a negative assertion mean something: an event that should not have been stored has to show up in the returned list, not be skipped over. """ _ws_send_json(sock, ["REQ", sub, {"ids": ids}]) events = [] deadline = time.time() + 15 while time.time() < deadline: msg = _ws_recv_json(sock, timeout=10) if msg[0] == "EVENT": events.append(msg[2]) elif msg[0] == "EOSE": return events raise TimeoutError("no EOSE for " + sub) class TestIngestWorkerPool: def test_valid_event_is_verified_by_a_worker_and_stored(self): proc, _ = _pool_relay("ingest-accept") try: sock = _ws_handshake(HOST, PORT) ev = make_event(TEST_SECKEY, "ingest pool accept") _ws_send_json(sock, ["EVENT", ev]) ok = _drain_until(sock, "OK", timeout=15) assert ok[1] == ev["id"] assert ok[2] is True # Stored, so the worker's acceptance reached the write path. stored = _query_ids(sock, "pool-query", [ev["id"]]) assert [e["id"] for e in stored] == [ev["id"]] finally: _stop_relay(proc) def test_ephemeral_event_is_accepted_and_not_stored(self): proc, _ = _pool_relay("ingest-ephemeral") try: sock = _ws_handshake(HOST, PORT) ev = make_event(TEST_SECKEY, "ephemeral through the pool", kind=EPHEMERAL_KIND) _ws_send_json(sock, ["EVENT", ev]) ok = _drain_until(sock, "OK", timeout=15) assert ok[1] == ev["id"] assert ok[2] is True assert _query_ids(sock, "eph-query", [ev["id"]]) == [] finally: _stop_relay(proc) def test_tampered_event_is_rejected_by_the_worker(self): proc, _ = _pool_relay("ingest-reject") try: sock = _ws_handshake(HOST, PORT) ev = make_event(TEST_SECKEY, "tampered after signing") # Change the content without re-signing: the id no longer matches # the serialised event, which Stage A checks before verifying. ev["content"] = "tampered after signing!" _ws_send_json(sock, ["EVENT", ev]) ok = _drain_until(sock, "OK", timeout=15) assert ok[2] is False assert "id mismatch" in ok[3] assert _query_ids(sock, "rej-query", [ev["id"]]) == [] finally: _stop_relay(proc) def test_malformed_event_envelope_is_rejected(self): proc, _ = _pool_relay("ingest-malformed") try: sock = _ws_handshake(HOST, PORT) _ws_send_json(sock, ["EVENT", {"id": "not-an-event"}]) ok = _drain_until(sock, "OK", timeout=15) assert ok[2] is False assert "malformed EVENT" in ok[3] finally: _stop_relay(proc) def test_single_worker_pool_answers_every_event(self): # One worker is the smallest pool: it exercises dispatch falling back to # the inline path when the only worker is busy, and the response drain # with no second slot. proc, _ = _pool_relay("ingest-single", workers=1) try: sock = _ws_handshake(HOST, PORT) first = make_event(TEST_SECKEY, "single worker one") second = make_event(TEST_SECKEY, "single worker two") _ws_send_json(sock, ["EVENT", first]) _ws_send_json(sock, ["EVENT", second]) seen = {} deadline = time.time() + 20 while len(seen) < 2 and time.time() < deadline: ok = _drain_until(sock, "OK", timeout=15) if ok[2]: seen[ok[1]] = True assert first["id"] in seen assert second["id"] in seen stored = _query_ids(sock, "single-query", [first["id"], second["id"]]) assert sorted(e["id"] for e in stored) == sorted( [first["id"], second["id"]] ) finally: _stop_relay(proc)