Compare commits
6 commits
47a5071b4a
...
edf1ddc7da
| Author | SHA1 | Date | |
|---|---|---|---|
| edf1ddc7da | |||
| dd326e3042 | |||
| 9f2f592f20 | |||
| 1c16a2a4b7 | |||
| 434817c899 | |||
| d8790087b4 |
5 changed files with 165 additions and 32 deletions
|
|
@ -110,9 +110,13 @@ class AdminInterface {
|
|||
|
||||
this.config().then((config) => {
|
||||
if (config.admin?.notifyAdminsOnBoot) {
|
||||
this.notifyAdminsOfNewConnection(connectionString);
|
||||
// .catch so a boot-DM failure can't surface as an unhandled
|
||||
// rejection (process-terminating under Node defaults). #48 / CS-3.
|
||||
this.notifyAdminsOfNewConnection(connectionString).catch((e) =>
|
||||
console.log('notifyAdminsOfNewConnection failed:', e?.message ?? e),
|
||||
);
|
||||
}
|
||||
});
|
||||
}).catch((e) => console.log('config() failed during admin init:', e?.message ?? e));
|
||||
}
|
||||
|
||||
public async config(): Promise<IConfig> {
|
||||
|
|
@ -127,13 +131,26 @@ class AdminInterface {
|
|||
const sk = secretKeyBytes(this.adminNsec);
|
||||
const pool = new RelayPool(['wss://blastr.f7z.xyz', 'wss://nostr.mutinywallet.com'], {});
|
||||
pool.start();
|
||||
// Give the connections a moment to come up before publishing.
|
||||
await new Promise((r) => setTimeout(r, 2500));
|
||||
|
||||
for (const npub of this.npubs || []) {
|
||||
await dmUser(sk, npub, `nsecBunker has started; use ${connectionString} to connect to it and unlock your key(s)`, pool);
|
||||
// Wait until at least one relay is actually connected (capped), rather
|
||||
// than a fixed sleep — these external public relays can be slow to come
|
||||
// up, and a too-short fixed wait made the DM publish fail on boot. Still
|
||||
// best-effort: if none connect in time we fall through and dmUser logs
|
||||
// the publish failure without affecting the daemon. (#48)
|
||||
const deadline = Date.now() + 8000;
|
||||
while (pool.connectedCount() === 0 && Date.now() < deadline) {
|
||||
await new Promise((r) => setTimeout(r, 100));
|
||||
}
|
||||
|
||||
// try/finally so a throw mid-loop can't leak the pool's reconnect loops +
|
||||
// sockets for the process lifetime (#48 / review CS-3). dmUser itself is
|
||||
// now fully guarded, but keep the finally as belt-and-suspenders.
|
||||
try {
|
||||
for (const npub of this.npubs || []) {
|
||||
await dmUser(sk, npub, `nsecBunker has started; use ${connectionString} to connect to it and unlock your key(s)`, pool);
|
||||
}
|
||||
} finally {
|
||||
pool.stop();
|
||||
}
|
||||
pool.stop();
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -397,13 +414,33 @@ class AdminInterface {
|
|||
return new Promise((resolve) => {
|
||||
console.log(`requesting permission for`, keyName, { remotePubkey, method });
|
||||
|
||||
const ids: string[] = [];
|
||||
let settled = false;
|
||||
// Resolve once and ALWAYS clear every pending sendRequest callback —
|
||||
// on timeout AND on the first admin response. Previously the timeout
|
||||
// cleared nothing and a single admin's response cleared only its own
|
||||
// id, so every timed-out request and (with multiple admins) the
|
||||
// non-responding admins' entries leaked in transport.pending for the
|
||||
// process lifetime (review AD-1/CS-2). The `settled` latch also stops
|
||||
// a late approval from acting after the request resolved (review AD-2).
|
||||
const finish = (value: boolean | undefined) => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
for (const id of ids) this.transport.clearPending(id);
|
||||
resolve(value);
|
||||
};
|
||||
|
||||
// If an admin doesn't respond within 10 seconds, report timeout.
|
||||
setTimeout(() => {
|
||||
resolve(undefined);
|
||||
}, 10000);
|
||||
setTimeout(() => finish(undefined), 10000);
|
||||
|
||||
for (const npub of this.npubs) {
|
||||
const adminPubkey = nip19.decode(npub).data as string;
|
||||
let adminPubkey: string;
|
||||
try {
|
||||
adminPubkey = nip19.decode(npub).data as string;
|
||||
} catch {
|
||||
console.log(`skipping malformed admin npub: ${npub}`);
|
||||
continue;
|
||||
}
|
||||
const params = JSON.stringify({
|
||||
keyName,
|
||||
remotePubkey,
|
||||
|
|
@ -419,17 +456,18 @@ class AdminInterface {
|
|||
'nip44',
|
||||
NIP46_ADMIN_RESPONSE_KIND,
|
||||
(res) => {
|
||||
this.transport.clearPending(id);
|
||||
if (settled) return; // ignore late / duplicate responses
|
||||
this.requestPermissionResponse(
|
||||
remotePubkey,
|
||||
keyName,
|
||||
method,
|
||||
param,
|
||||
resolve,
|
||||
finish,
|
||||
res
|
||||
);
|
||||
}
|
||||
);
|
||||
ids.push(id);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
|
|
|||
|
|
@ -56,6 +56,8 @@ interface ActiveSub {
|
|||
|
||||
const RECONNECT_BASE_MS = 1_000;
|
||||
const RECONNECT_CAP_MS = 10_000; // match #20's cap — cheap to retry a LAN relay
|
||||
const CONNECT_TIMEOUT_MS = 5_000; // bound a black-holed TCP connect (review RP-4)
|
||||
const SEEN_EVENT_CAP = 4_000; // bounded cross-relay event-id dedup set (review RP-1)
|
||||
|
||||
/**
|
||||
* A single relay connection that owns its (re)connect loop. On every successful
|
||||
|
|
@ -78,6 +80,9 @@ class ManagedRelay {
|
|||
public readonly url: string,
|
||||
private readonly registry: Map<string, PoolSubscription>,
|
||||
private readonly log: (...args: any[]) => void,
|
||||
/** Pool-wide first-seen gate: true the first time an event id is seen
|
||||
* across ALL relays, false on a duplicate. (review RP-1) */
|
||||
private readonly markSeen: (id: string) => boolean,
|
||||
) {}
|
||||
|
||||
start(): void {
|
||||
|
|
@ -122,7 +127,12 @@ class ManagedRelay {
|
|||
this.active.get(id)?.close();
|
||||
try {
|
||||
const sub = this.relay.subscribe(s.filters, {
|
||||
onevent: (e: Event) => s.onevent(e),
|
||||
// Dedup across relays (and across reconnect re-deliveries): a
|
||||
// kind:24133 request published to N relays must drive the daemon
|
||||
// handler / recordSigning ONCE, not N times (review RP-1 / CS-4).
|
||||
onevent: (e: Event) => {
|
||||
if (this.markSeen(e.id)) s.onevent(e);
|
||||
},
|
||||
oneose: () => s.oneose?.(),
|
||||
});
|
||||
this.active.set(id, sub);
|
||||
|
|
@ -179,16 +189,39 @@ class ManagedRelay {
|
|||
private connectOnce(): Promise<boolean> {
|
||||
return new Promise<boolean>((resolve) => {
|
||||
void (async () => {
|
||||
// enableReconnect:false — WE own reconnect, not nostr-tools.
|
||||
let relay: Relay;
|
||||
try {
|
||||
// enableReconnect:false — WE own reconnect, not nostr-tools.
|
||||
relay = await Relay.connect(this.url, { enableReconnect: false });
|
||||
relay = new Relay(this.url, { enableReconnect: false });
|
||||
} catch (e: any) {
|
||||
this.log("relay construct failed:", e?.message ?? e);
|
||||
resolve(false);
|
||||
return;
|
||||
}
|
||||
try {
|
||||
// Use the instance .connect({timeout}) rather than the static
|
||||
// Relay.connect(), which silently DROPS a timeout option in
|
||||
// nostr-tools 2.20.0. Without it a black-holed TCP connect
|
||||
// (SYN accepted, never upgraded) stalls this loop for the OS
|
||||
// socket timeout (minutes) with no retry (review RP-4).
|
||||
await relay.connect({ timeout: CONNECT_TIMEOUT_MS });
|
||||
} catch (e: any) {
|
||||
this.log("connect failed:", e?.message ?? e);
|
||||
try { relay.close(); } catch { /* ignore */ }
|
||||
resolve(false);
|
||||
return;
|
||||
}
|
||||
|
||||
// If stop() ran while the connect was in flight, drop this socket
|
||||
// cleanly — otherwise we'd re-arm subscriptions on a relay we mean
|
||||
// to abandon, leak the socket, and hang connectLoop (its promise
|
||||
// never resolves because onclose never fires). (review RP-2)
|
||||
if (this.stopped) {
|
||||
try { relay.close(); } catch { /* ignore */ }
|
||||
resolve(true);
|
||||
return;
|
||||
}
|
||||
|
||||
this.relay = relay;
|
||||
this.connected = true;
|
||||
this.lastConnectedAt = Date.now();
|
||||
|
|
@ -265,6 +298,9 @@ export class RelayPool {
|
|||
private counter = 0;
|
||||
private heartbeatTimer: ReturnType<typeof setInterval> | undefined;
|
||||
private lastHeartbeat = 0;
|
||||
/** Insertion-ordered (≈LRU) bounded set of event ids already delivered to a
|
||||
* subscription callback — the cross-relay/replay dedup gate (review RP-1). */
|
||||
private readonly seen: Set<string> = new Set();
|
||||
|
||||
constructor(
|
||||
public readonly relayUrls: string[],
|
||||
|
|
@ -273,12 +309,27 @@ export class RelayPool {
|
|||
this.log = opts.log ?? (() => {});
|
||||
this.relays = relayUrls.map(
|
||||
(url) =>
|
||||
new ManagedRelay(url, this.registry, (...a: any[]) =>
|
||||
this.log(`[relay:${url}]`, ...a),
|
||||
new ManagedRelay(
|
||||
url,
|
||||
this.registry,
|
||||
(...a: any[]) => this.log(`[relay:${url}]`, ...a),
|
||||
(id) => this.markSeen(id),
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
/** Returns true the first time `id` is seen pool-wide, false on a duplicate;
|
||||
* evicts the oldest id past the cap so this never grows unbounded. */
|
||||
private markSeen(id: string): boolean {
|
||||
if (this.seen.has(id)) return false;
|
||||
this.seen.add(id);
|
||||
if (this.seen.size > SEEN_EVENT_CAP) {
|
||||
const oldest = this.seen.values().next().value;
|
||||
if (oldest !== undefined) this.seen.delete(oldest);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
/** Start every relay's connect loop + (optionally) the sleep/wake heartbeat. */
|
||||
start(): void {
|
||||
for (const r of this.relays) r.start();
|
||||
|
|
|
|||
|
|
@ -132,6 +132,14 @@ export class Nip46Transport {
|
|||
): string {
|
||||
const id = Math.random().toString(36).substring(2, 12);
|
||||
this.pending.set(id, cb);
|
||||
// Defense-in-depth bound: callers (admin requestPermission) clear pending
|
||||
// entries on resolve/timeout, but cap the map so any un-cleared path can't
|
||||
// grow it without limit — evict the oldest (it would time out anyway).
|
||||
// (review AD-1/CS-2)
|
||||
if (this.pending.size > 1000) {
|
||||
const oldest = this.pending.keys().next().value;
|
||||
if (oldest !== undefined) this.pending.delete(oldest);
|
||||
}
|
||||
const content = this.encrypt(remotePubkey, JSON.stringify({ id, method, params }), encryption);
|
||||
const event = finalizeEvent(
|
||||
{ kind, created_at: Math.floor(Date.now() / 1000), tags: [["p", remotePubkey]], content },
|
||||
|
|
|
|||
|
|
@ -12,22 +12,26 @@ export async function dmUser(
|
|||
content: string,
|
||||
pool: RelayPool,
|
||||
): Promise<void> {
|
||||
const recipientHex = recipient.startsWith("npub1")
|
||||
? (nip19.decode(recipient).data as string)
|
||||
: recipient;
|
||||
const ciphertext = nip04.encrypt(sk, recipientHex, content);
|
||||
const event = finalizeEvent(
|
||||
{
|
||||
kind: 4,
|
||||
created_at: Math.floor(Date.now() / 1000),
|
||||
tags: [["p", recipientHex]],
|
||||
content: ciphertext,
|
||||
},
|
||||
sk,
|
||||
);
|
||||
// Guard the whole thing: nip19.decode throws on a malformed npub (passes the
|
||||
// startsWith check but fails the bech32 checksum), and that previously threw
|
||||
// *before* the publish try, escaping the caller. Best-effort — never throw.
|
||||
// (review CS-3)
|
||||
try {
|
||||
const recipientHex = recipient.startsWith("npub1")
|
||||
? (nip19.decode(recipient).data as string)
|
||||
: recipient;
|
||||
const ciphertext = nip04.encrypt(sk, recipientHex, content);
|
||||
const event = finalizeEvent(
|
||||
{
|
||||
kind: 4,
|
||||
created_at: Math.floor(Date.now() / 1000),
|
||||
tags: [["p", recipientHex]],
|
||||
content: ciphertext,
|
||||
},
|
||||
sk,
|
||||
);
|
||||
await pool.publish(event);
|
||||
} catch (e) {
|
||||
console.log(e);
|
||||
console.log('dmUser failed for', recipient, '-', (e as any)?.message ?? e);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -102,3 +102,35 @@ test("RelayPool.healthy() is false until the registry is subscribed on the wire
|
|||
pool.stop();
|
||||
await relay.stop();
|
||||
});
|
||||
|
||||
/**
|
||||
* RP-1: a NIP-46 request published to multiple relays (or re-delivered after a
|
||||
* reconnect) must drive the subscription callback ONCE — otherwise the daemon
|
||||
* signs N times and over-counts rate caps. The pool dedups by event id.
|
||||
*/
|
||||
test("RelayPool delivers each event id at most once (#RP-1)", async () => {
|
||||
const relay = new MockRelay();
|
||||
await relay.start();
|
||||
|
||||
const received: string[] = [];
|
||||
const pool = new RelayPool([relay.url], { log: () => {} });
|
||||
pool.start();
|
||||
await pool.subscribeAwaitingEose([{ kinds: [24133], "#p": [PUBKEY] }], (e) =>
|
||||
received.push(e.id),
|
||||
);
|
||||
|
||||
const ev = makeEvent();
|
||||
relay.inject(ev);
|
||||
relay.inject(ev); // same id again (simulates a second relay / a replay)
|
||||
await waitFor(() => received.includes(ev.id));
|
||||
await new Promise((r) => setTimeout(r, 200)); // give a 2nd delivery a chance
|
||||
|
||||
assert.equal(
|
||||
received.filter((id) => id === ev.id).length,
|
||||
1,
|
||||
"a duplicate event id is delivered to the callback only once",
|
||||
);
|
||||
|
||||
pool.stop();
|
||||
await relay.stop();
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue