Compare commits

..

18 Commits

Author SHA1 Message Date
fiatjaf
ab5ea8de36 another nip57 helper and bump version. 2023-02-16 09:29:21 -03:00
fiatjaf
a330b97590 partial nip57 support. 2023-02-15 21:06:38 -03:00
fiatjaf
24406b5679 more automatic cleanup of event listeners. 2023-02-15 20:36:22 -03:00
fiatjaf
6dbcc87d93 delete listeners when closing a relay connection. 2023-02-15 20:31:25 -03:00
fiatjaf
0ddcfdce68 remove "seen" event from Pub.
too complicated. if anyone wants this they can do it themselves.
2023-02-15 20:21:29 -03:00
fiatjaf
87bf349ce8 fill in missing kinds on enum. 2023-02-14 16:04:18 -03:00
fiatjaf
54dfc7b972 validate that the event is an object. 2023-02-14 15:18:39 -03:00
fiatjaf
32793146a4 remove untilOpen promise that was causing memory leaks when a connection was never opened. 2023-02-14 11:24:30 -03:00
fiatjaf
c42cd925ce bump noble-hashes. 2023-02-13 21:26:42 -03:00
RbnRncn
43ccb72476 docs: import SimplePool fix
Small fix of import SimplePool in new multiple relay docs.
2023-02-12 08:41:22 -03:00
fiatjaf
b2b7999517 notice about just.
closes https://github.com/nbd-wtf/nostr-tools/pull/106
2023-02-09 22:02:07 -03:00
fiatjaf
a568afc295 remove this extraneous file. 2023-02-09 22:01:01 -03:00
fiatjaf
9bcaed6e60 fix tests, .seenOn() method for pools. 2023-02-09 22:01:01 -03:00
Fernando López Guevara
5a9cbbb557 feat(deps): upgrade dependencies 2023-02-09 21:59:37 -03:00
fiatjaf
e9acc59809 just publish. 2023-02-09 12:09:16 -03:00
fiatjaf
18fe9637b9 do not run tests on tag pushes. 2023-02-09 12:08:50 -03:00
fiatjaf
ff3bf4a51c improvements and fixes on pool. 2023-02-09 12:05:31 -03:00
fiatjaf
7ff97b5488 list() and get() methods. 2023-02-08 16:37:53 -03:00
12 changed files with 586 additions and 163 deletions

View File

@@ -1,7 +1,9 @@
name: test every commit name: test every commit
on: on:
- push push:
- pull_request branches:
- master
pull_request:
jobs: jobs:
test: test:

View File

@@ -104,13 +104,15 @@ let pub = relay.publish(event)
pub.on('ok', () => { pub.on('ok', () => {
console.log(`${relay.url} has accepted our event`) console.log(`${relay.url} has accepted our event`)
}) })
pub.on('seen', () => {
console.log(`we saw the event on ${relay.url}`)
})
pub.on('failed', reason => { pub.on('failed', reason => {
console.log(`failed to publish to ${relay.url}: ${reason}`) console.log(`failed to publish to ${relay.url}: ${reason}`)
}) })
let events = await relay.list([{kinds: [0, 1]}])
let event = await relay.get({
ids: ['44e1827635450ebb3c5a7d12c1f8e7b2b514439ac10a67eef3d9fd9c5c68e245']
})
await relay.close() await relay.close()
``` ```
@@ -123,18 +125,13 @@ import 'websocket-polyfill'
### Interacting with multiple relays ### Interacting with multiple relays
```js ```js
import {pool} from 'nostr-tools' import {SimplePool} from 'nostr-tools'
const pool = new SimplePool() const pool = new SimplePool()
let relays = ['wss://relay.example.com', 'wss://relay.example2.com'] let relays = ['wss://relay.example.com', 'wss://relay.example2.com']
relays.forEach(async url => { let relay = await pool.ensureRelay('wss://relay.example3.com')
let relay = pool.ensureRelay(url)
await relay.connect()
})
let relay = pool.ensureRelay('wss://relay.example3.com')
let subs = pool.sub([...relays, relay], { let subs = pool.sub([...relays, relay], {
authors: ['32e1827635450ebb3c5a7d12c1f8e7b2b514439ac10a67eef3d9fd9c5c68e245'] authors: ['32e1827635450ebb3c5a7d12c1f8e7b2b514439ac10a67eef3d9fd9c5c68e245']
@@ -147,12 +144,22 @@ subs.forEach(sub =>
}) })
) )
let pubs = pool.publish(newEvent) let pubs = pool.publish(relays, newEvent)
pubs.forEach(pub => pubs.forEach(pub =>
pub.on('ok', () => { pub.on('ok', () => {
// ... // ...
}) })
) )
let events = await pool.list(relays, [{kinds: [0, 1]}])
let event = await pool.get(relays, {
ids: ['44e1827635450ebb3c5a7d12c1f8e7b2b514439ac10a67eef3d9fd9c5c68e245']
})
let relaysForEvent = pool.seenOn(
'44e1827635450ebb3c5a7d12c1f8e7b2b514439ac10a67eef3d9fd9c5c68e245'
)
// relaysForEvent will be an array of URLs from relays a given event was seen on
``` ```
### Querying profile data from a NIP-05 address ### Querying profile data from a NIP-05 address
@@ -283,6 +290,11 @@ Please consult the tests or [the source code](https://github.com/fiatjaf/nostr-t
</script> </script>
``` ```
## Plumbing
1. Install [`just`](https://just.systems/)
2. `just -l`
## License ## License
Public domain. Public domain.

View File

@@ -16,23 +16,31 @@ export enum Kind {
ChannelMetadata = 41, ChannelMetadata = 41,
ChannelMessage = 42, ChannelMessage = 42,
ChannelHideMessage = 43, ChannelHideMessage = 43,
ChannelMuteUser = 44 ChannelMuteUser = 44,
Report = 1984,
ZapRequest = 9734,
Zap = 9735,
RelayList = 10002,
ClientAuth = 22242,
Article = 30023
} }
export type Event = { export type EventTemplate = {
id?: string
sig?: string
kind: Kind kind: Kind
tags: string[][] tags: string[][]
pubkey: string
content: string content: string
created_at: number created_at: number
} }
export function getBlankEvent(): Event { export type Event = EventTemplate & {
pubkey: string
id: string
sig: string
}
export function getBlankEvent(): EventTemplate {
return { return {
kind: 255, kind: 255,
pubkey: '',
content: '', content: '',
tags: [], tags: [],
created_at: 0 created_at: 0
@@ -59,6 +67,7 @@ export function getEventHash(event: Event): string {
} }
export function validateEvent(event: Event): boolean { export function validateEvent(event: Event): boolean {
if (typeof event !== 'object') return false
if (typeof event.content !== 'string') return false if (typeof event.content !== 'string') return false
if (typeof event.created_at !== 'number') return false if (typeof event.created_at !== 'number') return false
if (typeof event.pubkey !== 'string') return false if (typeof event.pubkey !== 'string') return false

View File

@@ -9,6 +9,7 @@ export * as nip05 from './nip05'
export * as nip06 from './nip06' export * as nip06 from './nip06'
export * as nip19 from './nip19' export * as nip19 from './nip19'
export * as nip26 from './nip26' export * as nip26 from './nip26'
export * as nip57 from './nip57'
export * as fj from './fakejson' export * as fj from './fakejson'
export * as utils from './utils' export * as utils from './utils'

View File

@@ -11,3 +11,6 @@ test: build
testOnly file: build testOnly file: build
jest {{file}} jest {{file}}
publish: build
npm publish

140
magic.ts Normal file
View File

@@ -0,0 +1,140 @@
import {Relay, relayInit} from './relay'
import {Event} from './event'
import {normalizeURL} from './utils'
export default function (
writeableRelays: string[],
fallbackRelays: string[],
safeRelays: string[]
) {
return new MagicPool(fallbackRelays, writeableRelays, safeRelays)
}
class MagicPool {
private _conn: {[url: string]: Relay}
private _fallback: {[url: string]: Relay}
private _write: {[url: string]: Relay}
private _safe: {[url: string]: Relay}
private _profileRelays: {[pubkey: string]: RelayTableScore}
private _tempCache: {[id: string]: Event}
constructor(
fallbackRelays: string[],
writeableRelays: string[],
safeRelays: string[] = [
'wss://eden.nostr.land',
'wss://nostr.milou.lol',
'wss://relay.minds.com/nostr/v1/ws'
]
) {
this._conn = {}
this._write = {}
this._fallback = {}
this._profileRelays = {}
this._tempCache = {}
const hasEventId = (id: string): boolean => id in this._tempCache
const init = (url: string) => {
this._conn[normalizeURL(url)] = relayInit(normalizeURL(url), hasEventId)
}
fallbackRelays.forEach(init)
writeableRelays.forEach(init)
safeRelays.forEach(init)
this._write = Object.fromEntries(
writeableRelays.map(url => [
normalizeURL(url),
this._conn[normalizeURL(url)]
])
)
this._fallback = Object.fromEntries(
fallbackRelays.map(url => [
normalizeURL(url),
this._conn[normalizeURL(url)]
])
)
this._safe = Object.fromEntries(
safeRelays.map(url => [normalizeURL(url), this._conn[normalizeURL(url)]])
)
}
publish(event: Event) {
return Promise.all(
Object.entries(this._write).map(
([url, relay]) =>
new Promise(async resolve => {
await relay.connect()
let pub = relay.publish(event)
let to = setTimeout(() => {
let end = setTimeout(() => {
resolve({url, success: false, reason: 'timeout'})
}, 2500)
pub.on('seen', () => {
clearTimeout(end)
resolve({url, success: true, reason: 'seen'})
})
}, 2500)
pub.on('ok', () => {
clearTimeout(to)
resolve({url, success: true, reason: 'ok'})
})
pub.on('failed', (reason: string) => {
clearTimeout(to)
resolve({url, success: false, reason})
})
})
)
)
}
profile(
pubkey: string,
onUpdate: (events: Event[]) => void
): {
page(n: number): void
} {
var relays = new Set()
let rts = this._profileRelays[pubkey]
if (rts) {
relays = rts.get(3)
}
let fallback = Object.values(this._fallback)
for (let i = 0; i < fallback.length; i++) {
if (relays.size < 3) {
relays.add(fallback[Math.floor(Math.random() * fallback.length)])
} else break
}
// start subscription
for (let r in relays) {
r.
}
return {
page(n: number) {}
}
}
}
class RelayTableScore {
seen: string[] = []
hinted: string[] = []
explicit: string[] = []
get(n: number): Set<string> {
let relays = new Set<string>()
for (let i = 0; i < n; i++) {
for (let j = 0; j < 3; j++) {
let v = [this.seen, this.explicit, this.hinted][j][i]
if (v) {
relays.add(v)
if (relays.size >= n) return relays
}
}
}
return relays
}
}

107
nip57.ts Normal file
View File

@@ -0,0 +1,107 @@
import {bech32} from '@scure/base'
import {Event, EventTemplate} from './event'
import {utf8Decoder} from './utils'
var _fetch: any
try {
_fetch = fetch
} catch {}
export function useFetchImplementation(fetchImplementation: any) {
_fetch = fetchImplementation
}
export async function getZapEndpoint(metadata: Event): Promise<null | string> {
try {
let lnurl: string = ''
let {lud06, lud16} = JSON.parse(metadata.content)
if (lud06) {
let {words} = bech32.decode(lud06, 1000)
let data = bech32.fromWords(words)
lnurl = utf8Decoder.decode(data)
} else if (lud16) {
let [name, domain] = lud16.split('@')
lnurl = `https://${domain}/.well-known/lnurlp/${name}`
} else {
return null
}
let res = await _fetch(lnurl)
let body = await res.json()
if (body.allowsNostr && body.nostrPubkey) {
return body.callback
}
} catch (err) {
/*-*/
}
return null
}
export function makeZapRequest({
profile,
event,
amount,
relays,
comment = ''
}: {
profile: string
event: string | null
amount: string
comment: string
relays: string[]
}): EventTemplate {
let zr = {
kind: 9734,
created_at: Math.round(Date.now() / 1000),
content: comment,
tags: [
['p', profile],
['amount', amount],
['relays', ...relays]
]
}
if (event) {
zr.tags.push(['e', event])
}
return zr
}
export function makeZapReceipt({
zapRequest,
preimage,
bolt11,
paidAt
}: {
zapRequest: string
preimage: string | null
bolt11: string
paidAt: Date
}): EventTemplate {
let zr: Event = JSON.parse(zapRequest)
let tagsFromZapRequest = zr.tags.filter(
([t]) => t === 'e' || t === 'p' || t === 'a'
)
let zap = {
kind: 9735,
created_at: Math.round(paidAt.getTime() / 1000),
content: '',
tags: [
...tagsFromZapRequest,
['bolt11', bolt11],
['description', zapRequest]
]
}
if (preimage) {
zap.tags.push(['preimage', preimage])
}
return zap
}

View File

@@ -1,6 +1,6 @@
{ {
"name": "nostr-tools", "name": "nostr-tools",
"version": "1.2.4", "version": "1.4.2",
"description": "Tools for making a Nostr client.", "description": "Tools for making a Nostr client.",
"repository": { "repository": {
"type": "git", "type": "git",
@@ -9,11 +9,12 @@
"main": "lib/nostr.cjs.js", "main": "lib/nostr.cjs.js",
"module": "lib/nostr.esm.js", "module": "lib/nostr.esm.js",
"dependencies": { "dependencies": {
"@noble/hashes": "^0.5.7", "@noble/hashes": "1.0.0",
"@noble/secp256k1": "^1.7.0", "@noble/secp256k1": "^1.7.1",
"@scure/base": "^1.1.1", "@scure/base": "^1.1.1",
"@scure/bip32": "^1.1.1", "@scure/bip32": "^1.1.5",
"@scure/bip39": "^1.1.0" "@scure/bip39": "^1.1.1",
"prettier": "^2.8.4"
}, },
"keywords": [ "keywords": [
"decentralization", "decentralization",
@@ -23,20 +24,20 @@
"nostr" "nostr"
], ],
"devDependencies": { "devDependencies": {
"@types/node": "^18.0.3", "@types/node": "^18.13.0",
"@typescript-eslint/eslint-plugin": "^5.46.1", "@typescript-eslint/eslint-plugin": "^5.51.0",
"@typescript-eslint/parser": "^5.46.1", "@typescript-eslint/parser": "^5.51.0",
"esbuild": "0.16.9", "esbuild": "0.16.9",
"esbuild-plugin-alias": "^0.2.1", "esbuild-plugin-alias": "^0.2.1",
"eslint": "^8.30.0", "eslint": "^8.33.0",
"eslint-plugin-babel": "^5.3.1", "eslint-plugin-babel": "^5.3.1",
"esm-loader-typescript": "^1.0.1", "esm-loader-typescript": "^1.0.3",
"events": "^3.3.0", "events": "^3.3.0",
"jest": "^29.3.1", "jest": "^29.4.2",
"node-fetch": "2", "node-fetch": "^2.6.9",
"ts-jest": "^29.0.3", "ts-jest": "^29.0.5",
"tsd": "^0.22.0", "tsd": "^0.22.0",
"typescript": "^4.9.4", "typescript": "^4.9.5",
"websocket-polyfill": "^0.0.3" "websocket-polyfill": "^0.0.3"
} }
} }

View File

@@ -19,50 +19,28 @@ let relays = [
'wss://nostr.zebedee.cloud/' 'wss://nostr.zebedee.cloud/'
] ]
beforeAll(async () => {
Promise.all(
relays.map(relay => {
try {
let r = pool.ensureRelay(relay)
return r.connect()
} catch (err) {
/***/
}
})
)
})
afterAll(async () => { afterAll(async () => {
relays.forEach(relay => { await pool.close([
try { ...relays,
let r = pool.ensureRelay(relay) 'wss://nostr-relay.untethr.me',
r.close() 'wss://offchain.pub',
} catch (err) { 'wss://eden.nostr.land'
/***/ ])
}
})
}) })
test('removing duplicates when querying', async () => { test('removing duplicates when querying', async () => {
let priv = generatePrivateKey() let priv = generatePrivateKey()
let pub = getPublicKey(priv) let pub = getPublicKey(priv)
let subs = pool.sub(relays, [ let sub = pool.sub(relays, [{authors: [pub]}])
{
authors: [pub]
}
])
let received = [] let received = []
subs.forEach(sub => sub.on('event', event => {
sub.on('event', event => { // this should be called only once even though we're listening
// this should be called only once even though we're listening // to multiple relays because the events will be catched and
// to multiple relays because the events will be catched and // deduplicated efficiently (without even being parsed)
// deduplicated efficiently (without even being parsed) received.push(event)
received.push(event) })
})
)
let event = { let event = {
pubkey: pub, pubkey: pub,
@@ -81,25 +59,22 @@ test('removing duplicates when querying', async () => {
expect(received).toHaveLength(1) expect(received).toHaveLength(1)
}) })
test('removing duplicates correctly when double querying', async () => { test('same with double querying', async () => {
let priv = generatePrivateKey() let priv = generatePrivateKey()
let pub = getPublicKey(priv) let pub = getPublicKey(priv)
let subs1 = pool.sub(relays, [{authors: [pub]}]) let sub1 = pool.sub(relays, [{authors: [pub]}])
let subs2 = pool.sub(relays, [{authors: [pub]}]) let sub2 = pool.sub(relays, [{authors: [pub]}])
let received = [] let received = []
subs1.forEach(sub => sub1.on('event', event => {
sub.on('event', event => { received.push(event)
received.push(event) })
})
) sub2.on('event', event => {
subs2.forEach(sub => received.push(event)
sub.on('event', event => { })
received.push(event)
})
)
let event = { let event = {
pubkey: pub, pubkey: pub,
@@ -117,3 +92,42 @@ test('removing duplicates correctly when double querying', async () => {
expect(received).toHaveLength(2) expect(received).toHaveLength(2)
}) })
test('get()', async () => {
let event = await pool.get(relays, {
ids: ['d7dd5eb3ab747e16f8d0212d53032ea2a7cadef53837e5a6c66d42849fcb9027']
})
expect(event).toHaveProperty(
'id',
'd7dd5eb3ab747e16f8d0212d53032ea2a7cadef53837e5a6c66d42849fcb9027'
)
})
test('list()', async () => {
let events = await pool.list(
[...relays, 'wss://offchain.pub', 'wss://eden.nostr.land'],
[
{
authors: [
'3bf0c63fcb93463407af97a5e5ee64fa883d107ef9e558472c4eb9aaaefa459d'
],
kinds: [1],
limit: 2
}
]
)
// the actual received number will be greater than 2, but there will be no duplicates
expect(events.length).toEqual(
events
.map(evt => evt.id)
.reduce((acc, n) => (acc.indexOf(n) !== -1 ? acc : [...acc, n]), [])
.length
)
let relaysForAllEvents = events
.map(event => pool.seenOn(event.id))
.reduce((acc, n) => acc.concat(n), [])
expect(relaysForAllEvents.length).toBeGreaterThanOrEqual(events.length)
})

139
pool.ts
View File

@@ -6,13 +6,22 @@ import {SubscriptionOptions, Sub, Pub} from './relay'
export class SimplePool { export class SimplePool {
private _conn: {[url: string]: Relay} private _conn: {[url: string]: Relay}
private _seenOn: {[id: string]: Set<string>} = {} // a map of all events we've seen in each relay
constructor(defaultRelays: string[] = []) { constructor() {
this._conn = {} this._conn = {}
defaultRelays.forEach(this.ensureRelay)
} }
ensureRelay(url: string): Relay { async close(relays: string[]): Promise<void> {
await Promise.all(
relays.map(async url => {
let relay = this._conn[normalizeURL(url)]
if (relay) await relay.close()
})
)
}
async ensureRelay(url: string): Promise<Relay> {
const nm = normalizeURL(url) const nm = normalizeURL(url)
const existing = this._conn[nm] const existing = this._conn[nm]
if (existing) return existing if (existing) return existing
@@ -20,41 +29,131 @@ export class SimplePool {
const relay = relayInit(nm) const relay = relayInit(nm)
this._conn[nm] = relay this._conn[nm] = relay
await relay.connect()
return relay return relay
} }
sub(relays: string[], filters: Filter[], opts?: SubscriptionOptions): Sub[] { sub(relays: string[], filters: Filter[], opts?: SubscriptionOptions): Sub {
let _knownIds: Set<string> = new Set() let _knownIds: Set<string> = new Set()
let modifiedOpts = opts || {} let modifiedOpts = opts || {}
modifiedOpts.alreadyHaveEvent = id => _knownIds.has(id) modifiedOpts.alreadyHaveEvent = (id, url) => {
let set = this._seenOn[id] || new Set()
set.add(url)
this._seenOn[id] = set
return _knownIds.has(id)
}
return relays.map(relay => { let subs: Sub[] = []
let r = this._conn[relay] let eventListeners: Set<(event: Event) => void> = new Set()
if (!r) return badSub() let eoseListeners: Set<() => void> = new Set()
let eosesMissing = relays.length
let eoseSent = false
let eoseTimeout = setTimeout(() => {
eoseSent = true
for (let cb of eoseListeners.values()) cb()
}, 2400)
relays.forEach(async relay => {
let r = await this.ensureRelay(relay)
if (!r) return
let s = r.sub(filters, modifiedOpts) let s = r.sub(filters, modifiedOpts)
s.on('event', (event: Event) => _knownIds.add(event.id as string)) s.on('event', (event: Event) => {
return s _knownIds.add(event.id as string)
for (let cb of eventListeners.values()) cb(event)
})
s.on('eose', () => {
if (eoseSent) return
eosesMissing--
if (eosesMissing === 0) {
clearTimeout(eoseTimeout)
for (let cb of eoseListeners.values()) cb()
}
})
subs.push(s)
})
let greaterSub: Sub = {
sub(filters, opts) {
subs.forEach(sub => sub.sub(filters, opts))
return greaterSub
},
unsub() {
subs.forEach(sub => sub.unsub())
},
on(type, cb) {
switch (type) {
case 'event':
eventListeners.add(cb)
break
case 'eose':
eoseListeners.add(cb)
break
}
},
off(type, cb) {
if (type === 'event') {
eventListeners.delete(cb)
} else if (type === 'eose') eoseListeners.delete(cb)
}
}
return greaterSub
}
get(
relays: string[],
filter: Filter,
opts?: SubscriptionOptions
): Promise<Event | null> {
return new Promise(resolve => {
let sub = this.sub(relays, [filter], opts)
let timeout = setTimeout(() => {
sub.unsub()
resolve(null)
}, 1500)
sub.on('event', (event: Event) => {
resolve(event)
clearTimeout(timeout)
sub.unsub()
})
})
}
list(
relays: string[],
filters: Filter[],
opts?: SubscriptionOptions
): Promise<Event[]> {
return new Promise(resolve => {
let events: Event[] = []
let sub = this.sub(relays, filters, opts)
sub.on('event', (event: Event) => {
events.push(event)
})
// we can rely on an eose being emitted here because pool.sub() will fake one
sub.on('eose', () => {
sub.unsub()
resolve(events)
})
}) })
} }
publish(relays: string[], event: Event): Pub[] { publish(relays: string[], event: Event): Pub[] {
return relays.map(relay => { return relays.map(relay => {
let r = this._conn[relay] let r = this._conn[normalizeURL(relay)]
if (!r) return badPub(relay) if (!r) return badPub(relay)
let s = r.publish(event) let s = r.publish(event)
return s return s
}) })
} }
}
function badSub(): Sub { seenOn(id: string): string[] {
return { return Array.from(this._seenOn[id]?.values?.() || [])
on() {},
off() {},
sub(): Sub {
return badSub()
},
unsub() {}
} }
} }

View File

@@ -32,7 +32,7 @@ test('connectivity', () => {
).resolves.toBe(true) ).resolves.toBe(true)
}) })
test('querying', () => { test('querying', async () => {
var resolve1 var resolve1
var resolve2 var resolve2
@@ -52,16 +52,42 @@ test('querying', () => {
resolve2(true) resolve2(true)
}) })
return expect( let [t1, t2] = await Promise.all([
Promise.all([ new Promise(resolve => {
new Promise(resolve => { resolve1 = resolve
resolve1 = resolve }),
}), new Promise(resolve => {
new Promise(resolve => { resolve2 = resolve
resolve2 = resolve })
}) ])
])
).resolves.toEqual([true, true]) expect(t1).toEqual(true)
expect(t2).toEqual(true)
})
test('get()', async () => {
let event = await relay.get({
ids: ['d7dd5eb3ab747e16f8d0212d53032ea2a7cadef53837e5a6c66d42849fcb9027']
})
expect(event).toHaveProperty(
'id',
'd7dd5eb3ab747e16f8d0212d53032ea2a7cadef53837e5a6c66d42849fcb9027'
)
})
test('list()', async () => {
let events = await relay.list([
{
authors: [
'3bf0c63fcb93463407af97a5e5ee64fa883d107ef9e558472c4eb9aaaefa459d'
],
kinds: [1],
limit: 2
}
])
expect(events.length).toEqual(2)
}) })
test('listening (twice) and publishing', async () => { test('listening (twice) and publishing', async () => {

109
relay.ts
View File

@@ -12,13 +12,15 @@ export type Relay = {
connect: () => Promise<void> connect: () => Promise<void>
close: () => Promise<void> close: () => Promise<void>
sub: (filters: Filter[], opts?: SubscriptionOptions) => Sub sub: (filters: Filter[], opts?: SubscriptionOptions) => Sub
list: (filters: Filter[], opts?: SubscriptionOptions) => Promise<Event[]>
get: (filter: Filter, opts?: SubscriptionOptions) => Promise<Event | null>
publish: (event: Event) => Pub publish: (event: Event) => Pub
on: (type: RelayEvent, cb: any) => void on: (type: RelayEvent, cb: any) => void
off: (type: RelayEvent, cb: any) => void off: (type: RelayEvent, cb: any) => void
} }
export type Pub = { export type Pub = {
on: (type: 'ok' | 'seen' | 'failed', cb: any) => void on: (type: 'ok' | 'failed', cb: any) => void
off: (type: 'ok' | 'seen' | 'failed', cb: any) => void off: (type: 'ok' | 'failed', cb: any) => void
} }
export type Sub = { export type Sub = {
sub: (filters: Filter[], opts: SubscriptionOptions) => Sub sub: (filters: Filter[], opts: SubscriptionOptions) => Sub
@@ -28,18 +30,14 @@ export type Sub = {
} }
export type SubscriptionOptions = { export type SubscriptionOptions = {
skipVerification?: boolean
alreadyHaveEvent?: null | ((id: string) => boolean)
id?: string id?: string
skipVerification?: boolean
alreadyHaveEvent?: null | ((id: string, relay: string) => boolean)
} }
export function relayInit(url: string): Relay { export function relayInit(url: string): Relay {
var ws: WebSocket var ws: WebSocket
var resolveClose: () => void var resolveClose: () => void
var setOpen: (value: PromiseLike<void> | void) => void
var untilOpen = new Promise<void>(resolve => {
setOpen = resolve
})
var openSubs: {[id: string]: {filters: Filter[]} & SubscriptionOptions} = {} var openSubs: {[id: string]: {filters: Filter[]} & SubscriptionOptions} = {}
var listeners: { var listeners: {
connect: Array<() => void> connect: Array<() => void>
@@ -72,7 +70,6 @@ export function relayInit(url: string): Relay {
ws.onopen = () => { ws.onopen = () => {
listeners.connect.forEach(cb => cb()) listeners.connect.forEach(cb => cb())
setOpen()
resolve() resolve()
} }
ws.onerror = () => { ws.onerror = () => {
@@ -106,8 +103,12 @@ export function relayInit(url: string): Relay {
let subid = getSubscriptionId(json) let subid = getSubscriptionId(json)
if (subid) { if (subid) {
let {alreadyHaveEvent} = openSubs[subid] let so = openSubs[subid]
if (alreadyHaveEvent && alreadyHaveEvent(getHex64(json, 'id'))) { if (
so &&
so.alreadyHaveEvent &&
so.alreadyHaveEvent(getHex64(json, 'id'), url)
) {
return return
} }
} }
@@ -134,15 +135,22 @@ export function relayInit(url: string): Relay {
return return
case 'EOSE': { case 'EOSE': {
let id = data[1] let id = data[1]
;(subListeners[id]?.eose || []).forEach(cb => cb()) if (id in subListeners) {
subListeners[id].eose.forEach(cb => cb())
subListeners[id].eose = [] // 'eose' only happens once per sub, so stop listeners here
}
return return
} }
case 'OK': { case 'OK': {
let id: string = data[1] let id: string = data[1]
let ok: boolean = data[2] let ok: boolean = data[2]
let reason: string = data[3] || '' let reason: string = data[3] || ''
if (ok) pubListeners[id]?.ok.forEach(cb => cb()) if (id in pubListeners) {
else pubListeners[id]?.failed.forEach(cb => cb(reason)) if (ok) pubListeners[id].ok.forEach(cb => cb())
else pubListeners[id].failed.forEach(cb => cb(reason))
pubListeners[id].ok = [] // 'ok' only happens once per pub, so stop listeners here
pubListeners[id].failed = []
}
return return
} }
case 'NOTICE': case 'NOTICE':
@@ -165,7 +173,6 @@ export function relayInit(url: string): Relay {
async function trySend(params: [string, ...any]) { async function trySend(params: [string, ...any]) {
let msg = JSON.stringify(params) let msg = JSON.stringify(params)
await untilOpen
try { try {
ws.send(msg) ws.send(msg)
} catch (err) { } catch (err) {
@@ -231,54 +238,51 @@ export function relayInit(url: string): Relay {
let index = listeners[type].indexOf(cb) let index = listeners[type].indexOf(cb)
if (index !== -1) listeners[type].splice(index, 1) if (index !== -1) listeners[type].splice(index, 1)
}, },
list: (filters: Filter[], opts?: SubscriptionOptions): Promise<Event[]> =>
new Promise(resolve => {
let s = sub(filters, opts)
let events: Event[] = []
let timeout = setTimeout(() => {
s.unsub()
resolve(events)
}, 1500)
s.on('eose', () => {
s.unsub()
clearTimeout(timeout)
resolve(events)
})
s.on('event', (event: Event) => {
events.push(event)
})
}),
get: (filter: Filter, opts?: SubscriptionOptions): Promise<Event | null> =>
new Promise(resolve => {
let s = sub([filter], opts)
let timeout = setTimeout(() => {
s.unsub()
resolve(null)
}, 1500)
s.on('event', (event: Event) => {
s.unsub()
clearTimeout(timeout)
resolve(event)
})
}),
publish(event: Event): Pub { publish(event: Event): Pub {
if (!event.id) throw new Error(`event ${event} has no id`) if (!event.id) throw new Error(`event ${event} has no id`)
let id = event.id let id = event.id
var sent = false
var mustMonitor = false
trySend(['EVENT', event]) trySend(['EVENT', event])
.then(() => {
sent = true
if (mustMonitor) {
startMonitoring()
mustMonitor = false
}
})
.catch(() => {})
const startMonitoring = () => {
let monitor = sub([{ids: [id]}], {
id: `monitor-${id.slice(0, 5)}`
})
let willUnsub = setTimeout(() => {
;(pubListeners[id]?.failed || []).forEach(cb =>
cb('event not seen after 5 seconds')
)
monitor.unsub()
}, 5000)
monitor.on('event', () => {
clearTimeout(willUnsub)
;(pubListeners[id]?.seen || []).forEach(cb => cb())
})
}
return { return {
on: (type: 'ok' | 'seen' | 'failed', cb: any) => { on: (type: 'ok' | 'failed', cb: any) => {
pubListeners[id] = pubListeners[id] || { pubListeners[id] = pubListeners[id] || {
ok: [], ok: [],
seen: [],
failed: [] failed: []
} }
pubListeners[id][type].push(cb) pubListeners[id][type].push(cb)
if (type === 'seen') {
if (sent) startMonitoring()
else mustMonitor = true
}
}, },
off: (type: 'ok' | 'seen' | 'failed', cb: any) => { off: (type: 'ok' | 'failed', cb: any) => {
let listeners = pubListeners[id] let listeners = pubListeners[id]
if (!listeners) return if (!listeners) return
let idx = listeners[type].indexOf(cb) let idx = listeners[type].indexOf(cb)
@@ -288,6 +292,11 @@ export function relayInit(url: string): Relay {
}, },
connect, connect,
close(): Promise<void> { close(): Promise<void> {
listeners = {connect: [], disconnect: [], error: [], notice: []}
subListeners = {}
pubListeners = {}
if (ws.readyState > 1) return Promise.resolve()
ws.close() ws.close()
return new Promise<void>(resolve => { return new Promise<void>(resolve => {
resolveClose = resolve resolveClose = resolve