Files
c-relay-pg/tests/profile_fetch_latency_test.py
T
Laan Tungir 50653dc86a v2.1.38 - Custom backfill feature: ad-hoc backfill jobs with arbitrary NIP-01 filters + admin UI isolation
- New caching_custom_backfill_jobs and caching_custom_backfill_batches tables
- Admin API (POST create/cancel, GET list) at admin/api/custom_backfill.php
- Admin UI form with preset buttons and job status table on Backfill page
- Caching daemon module (custom_backfill.c) processes jobs independently
- pg_inbox functions for job/batch claim, progress update, and cancel
- Three batch modes: author-batched, id-batched, time-window scan
- Fixed done_batches overcount on job completion
- Admin config now environment-variable driven (C_RELAY_DB_*)
- admin/serve.sh supports separate instances for different databases
- make_and_restart_relay.sh only kills relay on target port, not all relays
2026-08-07 14:49:13 -04:00

304 lines
12 KiB
Python
Executable File

#!/usr/bin/env python3
"""
profile_fetch_latency_test.py — Measure Option A latency for a no-storage client.
Scenario:
1. Open a WebSocket REQ subscription for kind-1 events on the local relay.
2. For each kind-1 EVENT received, record the author pubkey.
3. After EOSE, immediately open a second REQ for kind-0 events for the
collected author set.
4. Record precise timings for every phase:
- REQ send -> first EVENT (kind-1)
- REQ send -> EOSE (kind-1)
- kind-0 REQ send -> first EVENT (kind-0)
- kind-0 REQ send -> EOSE (kind-0)
- per-event inter-arrival gaps
5. Print a summary so we can evaluate whether a local relay makes the
"fetch post, then fetch profile" pattern fast enough for a no-storage
client.
Usage:
python3 tests/profile_fetch_latency_test.py [--relay ws://localhost:7777]
[--limit 50]
[--timeout 10]
Requires: Python 3.8+, `websockets` (pip install websockets), `websockets` is
the only third-party dependency. No Nostr signing is needed because we only
read.
Notes:
- The relay is on port 7777 per the user's instruction.
- We use a fresh subscription id for each REQ.
- We measure wall-clock with time.perf_counter() (monotonic, high-res).
"""
import argparse
import asyncio
import json
import statistics
import sys
import time
from collections import OrderedDict
try:
import websockets
except ImportError:
sys.stderr.write(
"ERROR: 'websockets' package is required. Install with:\n"
" pip install websockets\n"
)
sys.exit(1)
# --------------------------------------------------------------------------- #
# Timing record
# --------------------------------------------------------------------------- #
class Timings:
"""Collects high-resolution timestamps for each phase."""
def __init__(self):
self.t_req1_send = None # kind-1 REQ frame sent
self.t_req1_first_event = None # first kind-1 EVENT received
self.t_req1_eose = None # kind-1 EOSE received
self.t_req0_send = None # kind-0 REQ frame sent
self.t_req0_first_event = None # first kind-0 EVENT received
self.t_req0_eose = None # kind-0 EOSE received
# inter-arrival: time between consecutive kind-1 EVENT frames
self.kind1_event_times = []
# inter-arrival: time between consecutive kind-0 EVENT frames
self.kind0_event_times = []
self.kind1_count = 0
self.kind0_count = 0
self.missing_profiles = 0 # authors with no kind-0 returned
def ms(self, start, end):
"""Return elapsed milliseconds between two timestamps, or None."""
if start is None or end is None:
return None
return (end - start) * 1000.0
def summary(self):
lines = []
lines.append("=" * 64)
lines.append("PROFILE FETCH LATENCY — Option A: two standard REQs")
lines.append("=" * 64)
def fmt(ms):
return f"{ms:8.2f} ms" if ms is not None else " n/a "
lines.append("")
lines.append("Phase 1 — kind-1 feed REQ")
lines.append(f" REQ send -> first EVENT : {fmt(self.ms(self.t_req1_send, self.t_req1_first_event))}")
lines.append(f" REQ send -> EOSE : {fmt(self.ms(self.t_req1_send, self.t_req1_eose))}")
lines.append(f" kind-1 events received : {self.kind1_count}")
if len(self.kind1_event_times) > 1:
gaps = [
(self.kind1_event_times[i] - self.kind1_event_times[i - 1]) * 1000.0
for i in range(1, len(self.kind1_event_times))
]
lines.append(f" inter-event gap min/avg/max: "
f"{min(gaps):.2f} / {statistics.mean(gaps):.2f} / {max(gaps):.2f} ms")
lines.append("")
lines.append("Phase 2 — kind-0 profile REQ (batched authors)")
lines.append(f" REQ send -> first EVENT : {fmt(self.ms(self.t_req0_send, self.t_req0_first_event))}")
lines.append(f" REQ send -> EOSE : {fmt(self.ms(self.t_req0_send, self.t_req0_eose))}")
lines.append(f" kind-0 events received : {self.kind0_count}")
lines.append(f" authors requested : {self.kind1_count}")
lines.append(f" missing profiles : {self.missing_profiles}")
if len(self.kind0_event_times) > 1:
gaps = [
(self.kind0_event_times[i] - self.kind0_event_times[i - 1]) * 1000.0
for i in range(1, len(self.kind0_event_times))
]
lines.append(f" inter-event gap min/avg/max: "
f"{min(gaps):.2f} / {statistics.mean(gaps):.2f} / {max(gaps):.2f} ms")
# End-to-end: from kind-1 REQ send to having all profiles (kind-0 EOSE)
e2e = self.ms(self.t_req1_send, self.t_req0_eose)
lines.append("")
lines.append(f"END-TO-END (kind-1 REQ -> kind-0 EOSE): {fmt(e2e)}")
lines.append("=" * 64)
return "\n".join(lines)
# --------------------------------------------------------------------------- #
# WebSocket helpers
# --------------------------------------------------------------------------- #
async def send_req(ws, sub_id, filter_obj):
"""Send a ["REQ", sub_id, filter] frame and return the send timestamp."""
msg = json.dumps(["REQ", sub_id, filter_obj])
await ws.send(msg)
return time.perf_counter()
async def close_sub(ws, sub_id):
"""Send a CLOSE frame for a subscription."""
try:
await ws.send(json.dumps(["CLOSE", sub_id]))
except Exception:
pass
# --------------------------------------------------------------------------- #
# Phase 1: collect kind-1 events
# --------------------------------------------------------------------------- #
async def fetch_kind1(ws, limit, timeout, timings):
sub_id = "k1feed"
filt = {"kinds": [1], "limit": limit}
timings.t_req1_send = await send_req(ws, sub_id, filt)
authors = OrderedDict() # pubkey -> count (preserves insertion order, dedups)
deadline = time.perf_counter() + timeout
try:
async for raw in ws:
# websockets returns str for text frames
try:
msg = json.loads(raw)
except json.JSONDecodeError:
continue
if not isinstance(msg, list) or len(msg) < 2:
continue
msg_type = msg[0]
if msg_type == "EVENT" and len(msg) >= 3:
event = msg[2]
kind = event.get("kind")
if kind == 1:
now = time.perf_counter()
if timings.t_req1_first_event is None:
timings.t_req1_first_event = now
timings.kind1_event_times.append(now)
timings.kind1_count += 1
pk = event.get("pubkey", "")
if pk:
authors[pk] = authors.get(pk, 0) + 1
if timings.kind1_count >= limit:
# We have enough; wait for EOSE or break
pass
elif msg_type == "EOSE" and len(msg) >= 2 and msg[1] == sub_id:
timings.t_req1_eose = time.perf_counter()
break
elif msg_type == "NOTICE":
sys.stderr.write(f"[NOTICE] {msg[1] if len(msg) > 1 else ''}\n")
# Safety timeout
if time.perf_counter() > deadline:
sys.stderr.write("[timeout] Phase 1 timed out waiting for EOSE\n")
break
finally:
await close_sub(ws, sub_id)
return list(authors.keys())
# --------------------------------------------------------------------------- #
# Phase 2: fetch kind-0 profiles for the collected authors
# --------------------------------------------------------------------------- #
async def fetch_kind0(ws, authors, timeout, timings):
if not authors:
sys.stderr.write("[skip] No authors collected; skipping kind-0 phase\n")
return
sub_id = "k0profiles"
# NIP-01 allows multiple authors in one filter. Some relays cap the
# authors array length; we keep it as one REQ for the local-relay test.
filt = {"kinds": [0], "authors": authors}
timings.t_req0_send = await send_req(ws, sub_id, filt)
received_authors = set()
deadline = time.perf_counter() + timeout
try:
async for raw in ws:
try:
msg = json.loads(raw)
except json.JSONDecodeError:
continue
if not isinstance(msg, list) or len(msg) < 2:
continue
msg_type = msg[0]
if msg_type == "EVENT" and len(msg) >= 3:
event = msg[2]
kind = event.get("kind")
if kind == 0:
now = time.perf_counter()
if timings.t_req0_first_event is None:
timings.t_req0_first_event = now
timings.kind0_event_times.append(now)
timings.kind0_count += 1
pk = event.get("pubkey", "")
if pk:
received_authors.add(pk)
elif msg_type == "EOSE" and len(msg) >= 2 and msg[1] == sub_id:
timings.t_req0_eose = time.perf_counter()
break
elif msg_type == "NOTICE":
sys.stderr.write(f"[NOTICE] {msg[1] if len(msg) > 1 else ''}\n")
if time.perf_counter() > deadline:
sys.stderr.write("[timeout] Phase 2 timed out waiting for EOSE\n")
break
finally:
await close_sub(ws, sub_id)
# Authors we asked for but got no kind-0 back
timings.missing_profiles = len(set(authors) - received_authors)
# --------------------------------------------------------------------------- #
# Main
# --------------------------------------------------------------------------- #
async def main():
parser = argparse.ArgumentParser(
description="Measure Option A profile-fetch latency against a local relay."
)
parser.add_argument("--relay", default="ws://localhost:7777",
help="Relay WebSocket URL (default: ws://localhost:7777)")
parser.add_argument("--limit", type=int, default=50,
help="Number of kind-1 events to fetch (default: 50)")
parser.add_argument("--timeout", type=float, default=10.0,
help="Per-phase timeout in seconds (default: 10)")
args = parser.parse_args()
timings = Timings()
print(f"Connecting to {args.relay} ...")
try:
async with websockets.connect(args.relay, max_size=2**22) as ws:
print(f"Connected. Phase 1: REQ kind-1 limit={args.limit}")
authors = await fetch_kind1(ws, args.limit, args.timeout, timings)
print(f"Phase 1 done: {timings.kind1_count} events, "
f"{len(authors)} unique authors")
if authors:
print(f"Phase 2: REQ kind-0 for {len(authors)} authors")
await fetch_kind0(ws, authors, args.timeout, timings)
print(f"Phase 2 done: {timings.kind0_count} profiles, "
f"{timings.missing_profiles} missing")
except websockets.exceptions.InvalidURI:
sys.stderr.write(f"ERROR: invalid relay URI: {args.relay}\n")
sys.exit(1)
except (OSError, websockets.exceptions.WebSocketException) as e:
sys.stderr.write(f"ERROR: could not connect to relay: {e}\n")
sys.exit(1)
print()
print(timings.summary())
if __name__ == "__main__":
asyncio.run(main())