test_relay_ingest_workers.py raw

   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