fix(eventsync): drain the outbox before closing and count each send once

Audit of the sync path the new geode-backed EventSyncTest exercises found
two bugs in EventSync itself, plus review nits on the harness changes.

- runSync closed its client (`use {}`) the moment the last page arrived,
  while `publish` is fire-and-forget through the client's outbox. Events
  forwarded from the final page of the last relay were still waiting for
  a socket or an OK when the outbox was destroyed, so the sync reported
  Done and silently never delivered them. Wait, bounded by the existing
  per-relay timeout, until no forwarded event has a relay left pending.
- The "events sent" counters incremented on every onSent, including the
  failed write to a destination still connecting and the outbox's
  at-least-once resend of an unacknowledged event after the connection
  syncs. Every cold destination therefore reported at least one extra
  event sent. Count only successful writes, once per (event, relay).
  The test now asserts the sent total equals the routed total.
- Harness: the 127.0.0.2 rationale claimed it survives Quartz's
  isLocalHost() strip; that filter now covers all of 127.0.0.0/8, so
  say so and note what it means for the strict-inbox DM cases. The
  interactive Marmot harness gets an overridable RELAY_HOST/RELAY_BIND
  and documents the loopback/RFC1918 stripping limit it inherits, and
  its new --port guards a missing value instead of dying on set -u.
- EventSyncTest builds both scenarios through one helper.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PguqnDbP2v11dtANs9xdxc
This commit is contained in:
Claude
2026-09-12 20:57:11 +00:00
parent cfa6f72760
commit 42698ed59c
6 changed files with 94 additions and 34 deletions
@@ -43,10 +43,12 @@ import com.vitorpamplona.quartz.nip59Giftwrap.wraps.GiftWrapEvent
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.Job
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.update
import kotlinx.coroutines.launch
import kotlinx.coroutines.withTimeoutOrNull
import java.util.concurrent.ConcurrentHashMap
import kotlin.coroutines.cancellation.CancellationException
@@ -90,6 +92,9 @@ class EventSync(
/** Maximum number of completed-relay entries kept in the activity log. */
const val MAX_ACTIVITY_LOG = 5000
/** Poll interval while waiting for the last forwarded events to be acknowledged. */
const val OUTBOX_DRAIN_POLL_MS = 100L
}
// -------------------------------------------------------------------------
@@ -392,6 +397,13 @@ class EventSync(
val sourceRelayOfEvent = ConcurrentHashMap<HexKey, NormalizedRelayUrl>()
// (event id, destination) pairs already counted as sent. The outbox is
// at-least-once: it writes an event as soon as the socket is ready and
// resends everything still unacknowledged when the connection finishes
// syncing, so one event can hit the same relay twice before its OK lands.
// The relay dedups the second copy; the counters must too.
val sentPairs = ConcurrentHashMap.newKeySet<String>()
val runningState =
SyncState.Running(
relaysCompleted = 0,
@@ -423,7 +435,12 @@ class EventSync(
success: Boolean,
) {
super.onSent(relay, cmdStr, cmd, success)
if (cmd is EventCmd) {
// `success` is "written to the socket", not "OK received". A write to a
// destination that is still connecting fails and the outbox resends it
// once the socket opens; counting the failed attempt too made every
// cold destination report one extra event sent. Likewise a successful
// resend of an unacknowledged event is the same send, not a second one.
if (cmd is EventCmd && success && sentPairs.add(cmd.event.id + relay.url.url)) {
var hasSent = false
if (outboxDedup.contains(cmd.event.id)) {
@@ -586,6 +603,13 @@ class EventSync(
},
)
// `publish` is fire-and-forget through the client's outbox, and `use` closes
// the client as soon as this block returns. Without a drain, the events
// forwarded from the last page of the last relay are still waiting for a
// socket or an OK when the outbox is destroyed — the sync reports Done and
// silently never delivers them. Bounded by the same per-relay timeout.
awaitOutboxDrain(client, outboxDedup + inboxDedup + dmDedup)
_syncState.value =
SyncState.Done(
totalEventsReceived = runningState.eventsReceived.value,
@@ -613,4 +637,22 @@ class EventSync(
}
}
}
/**
* Waits until no forwarded event in [ids] has a relay left in the client's outbox, or
* until [RELAY_TIMEOUT_MS] passes. Ids that drain are dropped from the working set so
* each poll only revisits what is still pending.
*/
private suspend fun awaitOutboxDrain(
client: INostrClient,
ids: Set<HexKey>,
) {
val pending = ids.toMutableSet()
withTimeoutOrNull(RELAY_TIMEOUT_MS) {
while (pending.isNotEmpty()) {
pending.removeAll { client.pendingPublishRelaysFor(it).isNullOrEmpty() }
if (pending.isNotEmpty()) delay(OUTBOX_DRAIN_POLL_MS)
}
}
}
}
@@ -94,14 +94,18 @@ class EventSyncTest : RelayClientTest() {
private fun corpus(): List<Event> = mine + mentions + dmToMe + noise
private fun eventSync(builder: WebsocketBuilder): EventSync =
/** [decorate] runs on every client the sync builds, e.g. to attach an authenticator. */
private fun eventSync(
builder: WebsocketBuilder,
decorate: (NostrClient) -> Unit = {},
): EventSync =
EventSync(
accountPubKey = account.pubKey,
relayDb = { listOf(source) },
outboxTargets = { setOf(outbox) },
inboxTargets = { setOf(inbox) },
dmTargets = { setOf(dm) },
clientBuilder = { NostrClient(builder, scope) },
clientBuilder = { NostrClient(builder, scope).also(decorate) },
scope = scope,
)
@@ -164,6 +168,14 @@ class EventSyncTest : RelayClientTest() {
mine.size + mentions.size + 1,
done.totalEventsReceived,
)
// runSync drains the outbox before closing its client, so by the time
// Done is published every forwarded event has been written to its
// destination socket — not merely queued.
assertEquals(
"every routed event was sent before the client closed",
mine.size + mentions.size + 1,
done.totalEventsSent,
)
assertRouted(hub)
}
@@ -185,22 +197,12 @@ class EventSyncTest : RelayClientTest() {
val authSigner = NostrSignerSync(KeyPair())
var authenticator: RelayAuthenticator? = null
val sync =
EventSync(
accountPubKey = account.pubKey,
relayDb = { listOf(source) },
outboxTargets = { setOf(outbox) },
inboxTargets = { setOf(inbox) },
dmTargets = { setOf(dm) },
clientBuilder = {
val client = NostrClient(router, scope)
authenticator =
RelayAuthenticator(client = client, scope = scope) { _, template, _ ->
listOf(authSigner.sign(template))
}
client
},
scope = scope,
)
eventSync(router) { client ->
authenticator =
RelayAuthenticator(client = client, scope = scope) { _, template, _ ->
listOf(authSigner.sign(template))
}
}
try {
withTimeout(30_000) { sync.runSync() }
+5 -3
View File
@@ -40,9 +40,11 @@ RESULTS_FILE="$STATE_DIR/results-$RUN_TS.tsv"
AMY_BIN="$REPO_ROOT/cli/build/install/amy/bin/amy"
# Loopback relay = `amy serve` (geode), booted from $AMY_BIN by
# start_local_relay in headless/helpers.sh. 127.0.0.2 rather than
# 127.0.0.1 so Quartz's isLocalHost() filter doesn't strip it out of the
# published relay lists (see the DM harness for the full note).
# start_local_relay in headless/helpers.sh. 127.0.0.2 only for parity
# with the DM and Marmot harnesses: Quartz's isLocalHost() now covers all
# of 127.0.0.0/8, so it is stripped from parsed relay lists exactly like
# 127.0.0.1. Nothing here depends on that parse — amy publishes to and
# reads from the relay it was told about.
RELAY_HOST="${RELAY_HOST:-127.0.0.2}"
RELAY_DATA="$STATE_DIR/relay"
RELAY_PORT="${RELAY_PORT:-8092}"
+6 -4
View File
@@ -30,10 +30,12 @@ AMY_BIN="$REPO_ROOT/cli/build/install/amy/bin/amy"
# Loopback relay = `amy serve` (geode), booted from $AMY_BIN by
# start_local_relay in headless/helpers.sh. Override RELAY_DATA if you
# want full isolation between runs.
# Bind the loopback relay to 127.0.0.2 rather than 127.0.0.1 so Quartz's
# `isLocalHost()` filter doesn't silently strip it out of the kind:10050
# inbox events during recipient-relay resolution. 127.0.0.2 is still pure
# loopback — no network traffic, no config needed.
# 127.0.0.2 used to dodge Quartz's `isLocalHost()` strip of loopback
# relays in kind:10050 inbox lists. That filter now covers all of
# 127.0.0.0/8, so the strict-inbox sends (dm-01/02/05/06) fail with
# no_dm_relays regardless of which loopback address the relay binds;
# only the fallback-chain tests (dm-03/04) are unaffected. Kept for
# parity with the other harnesses until that routing rule is revisited.
RELAY_HOST="${RELAY_HOST:-127.0.0.2}"
RELAY_DATA="$STATE_DIR/relay"
RELAY_PORT="${RELAY_PORT:-8090}"
+5 -2
View File
@@ -76,8 +76,11 @@ assert_eq() {
#
# Callers set (before sourcing or at least before calling):
# AMY_BIN amy launcher (built via `./gradlew :cli:installDist`)
# RELAY_HOST host clients connect to (most harnesses use 127.0.0.2 —
# see the isLocalHost() note at the top of each script)
# RELAY_HOST host clients connect to. The harnesses use 127.0.0.2 for
# parity with each other; note that Quartz's isLocalHost()
# treats all of 127.0.0.0/8 as loopback, so it does NOT
# survive the NIP-17 / NIP-65 relay-list parsers any better
# than 127.0.0.1 does.
# RELAY_BIND optional bind address; defaults to $RELAY_HOST. Set to
# 0.0.0.0 when a device on the LAN must reach the relay.
# RELAY_PORT listen port
+15 -6
View File
@@ -37,10 +37,17 @@ WND_BIN=""
AMY_BIN="$REPO_ROOT/cli/build/install/amy/bin/amy"
# Embedded relay (default mode). Bound on every interface so a device on the
# same network can reach it; the daemons connect over loopback. Loopback
# same network can reach it; the daemons connect over $RELAY_HOST. Loopback
# `ws://` relays are only accepted by MDK behind this explicit opt-in.
RELAY_HOST="127.0.0.1"
RELAY_BIND="0.0.0.0"
#
# Known limit, inherited from the old --local-relays mode: the URL wn
# advertises in its kind:10050/10051 lists is $RELAY_URL, and Amethyst's
# parsers drop loopback and RFC1918 relays from those lists, so A→B welcome
# delivery leans on Amethyst's fallback relays. Override RELAY_HOST with an
# address the device can dial (e.g. the laptop's LAN IP) to have wn
# advertise that instead; the daemons then connect to it too.
RELAY_HOST="${RELAY_HOST:-127.0.0.1}"
RELAY_BIND="${RELAY_BIND:-0.0.0.0}"
RELAY_PORT="${RELAY_PORT:-8080}"
RELAY_URL="ws://$RELAY_HOST:$RELAY_PORT"
RELAY_DATA="$STATE_DIR/relay"
@@ -85,7 +92,9 @@ while [[ $# -gt 0 ]]; do
case "$1" in
--public-relays) USE_PUBLIC_RELAYS=1 ;;
--local-relays) printf '%s\n' "note: --local-relays is now the default (embedded amy serve relay); flag ignored" >&2 ;;
--port) RELAY_PORT="$2"; RELAY_URL="ws://$RELAY_HOST:$RELAY_PORT"; shift ;;
--port)
[[ $# -ge 2 && "$2" != --* ]] || { printf 'missing value for --port\n' >&2; usage; exit 2; }
RELAY_PORT="$2"; RELAY_URL="ws://$RELAY_HOST:$RELAY_PORT"; shift ;;
--transponder) ENABLE_TRANSPONDER=1 ;;
--no-build) NO_BUILD=1 ;;
-h|--help) usage; exit 0 ;;
@@ -567,14 +576,14 @@ configure_relays() {
info "sanity kinds 10050/1059/445 ok (B->C welcome + message round-trip)"
else
warn "kind:445 failed — C never decrypted sanity-ping (relays may be dropping group messages)"
warn "Consider rerunning without --public-relays (the embedded relay accepts every kind)."
warn "Consider rerunning without --public-relays (the embedded relay stores every kind)."
fi
# best-effort cleanup so re-runs don't accumulate dead sanity groups
wn_c groups leave "$sanity_c_gid" >/dev/null 2>&1 || true
wn_b groups leave "$sanity_gid" >/dev/null 2>&1 || true
else
warn "kind:10050/1059 failed — C never received welcome; relays likely dropping gift wraps or inbox lists"
warn "Consider rerunning without --public-relays (the embedded relay accepts every kind)."
warn "Consider rerunning without --public-relays (the embedded relay stores every kind)."
fi
fi
}