493 lines
12 KiB
JavaScript
493 lines
12 KiB
JavaScript
// event.ts
|
|
import { schnorr } from "@noble/curves/secp256k1";
|
|
import { sha256 } from "@noble/hashes/sha256";
|
|
import { bytesToHex } from "@noble/hashes/utils";
|
|
|
|
// utils.ts
|
|
var utf8Decoder = new TextDecoder("utf-8");
|
|
var utf8Encoder = new TextEncoder();
|
|
var MessageNode = class {
|
|
_value;
|
|
_next;
|
|
get value() {
|
|
return this._value;
|
|
}
|
|
set value(message) {
|
|
this._value = message;
|
|
}
|
|
get next() {
|
|
return this._next;
|
|
}
|
|
set next(node) {
|
|
this._next = node;
|
|
}
|
|
constructor(message) {
|
|
this._value = message;
|
|
this._next = null;
|
|
}
|
|
};
|
|
var MessageQueue = class {
|
|
_first;
|
|
_last;
|
|
get first() {
|
|
return this._first;
|
|
}
|
|
set first(messageNode) {
|
|
this._first = messageNode;
|
|
}
|
|
get last() {
|
|
return this._last;
|
|
}
|
|
set last(messageNode) {
|
|
this._last = messageNode;
|
|
}
|
|
_size;
|
|
get size() {
|
|
return this._size;
|
|
}
|
|
set size(v) {
|
|
this._size = v;
|
|
}
|
|
constructor() {
|
|
this._first = null;
|
|
this._last = null;
|
|
this._size = 0;
|
|
}
|
|
enqueue(message) {
|
|
const newNode = new MessageNode(message);
|
|
if (this._size === 0 || !this._last) {
|
|
this._first = newNode;
|
|
this._last = newNode;
|
|
} else {
|
|
this._last.next = newNode;
|
|
this._last = newNode;
|
|
}
|
|
this._size++;
|
|
return true;
|
|
}
|
|
dequeue() {
|
|
if (this._size === 0 || !this._first)
|
|
return null;
|
|
let prev = this._first;
|
|
this._first = prev.next;
|
|
prev.next = null;
|
|
this._size--;
|
|
return prev.value;
|
|
}
|
|
};
|
|
|
|
// event.ts
|
|
var verifiedSymbol = Symbol("verified");
|
|
function serializeEvent(evt) {
|
|
if (!validateEvent(evt))
|
|
throw new Error("can't serialize event with wrong or missing properties");
|
|
return JSON.stringify([0, evt.pubkey, evt.created_at, evt.kind, evt.tags, evt.content]);
|
|
}
|
|
function getEventHash(event) {
|
|
let eventHash = sha256(utf8Encoder.encode(serializeEvent(event)));
|
|
return bytesToHex(eventHash);
|
|
}
|
|
var isRecord = (obj) => obj instanceof Object;
|
|
function validateEvent(event) {
|
|
if (!isRecord(event))
|
|
return false;
|
|
if (typeof event.kind !== "number")
|
|
return false;
|
|
if (typeof event.content !== "string")
|
|
return false;
|
|
if (typeof event.created_at !== "number")
|
|
return false;
|
|
if (typeof event.pubkey !== "string")
|
|
return false;
|
|
if (!event.pubkey.match(/^[a-f0-9]{64}$/))
|
|
return false;
|
|
if (!Array.isArray(event.tags))
|
|
return false;
|
|
for (let i = 0; i < event.tags.length; i++) {
|
|
let tag = event.tags[i];
|
|
if (!Array.isArray(tag))
|
|
return false;
|
|
for (let j = 0; j < tag.length; j++) {
|
|
if (typeof tag[j] === "object")
|
|
return false;
|
|
}
|
|
}
|
|
return true;
|
|
}
|
|
function verifySignature(event) {
|
|
if (typeof event[verifiedSymbol] === "boolean")
|
|
return event[verifiedSymbol];
|
|
const hash = getEventHash(event);
|
|
if (hash !== event.id) {
|
|
return event[verifiedSymbol] = false;
|
|
}
|
|
try {
|
|
return event[verifiedSymbol] = schnorr.verify(event.sig, hash, event.pubkey);
|
|
} catch (err) {
|
|
return event[verifiedSymbol] = false;
|
|
}
|
|
}
|
|
|
|
// filter.ts
|
|
function matchFilter(filter, event) {
|
|
if (filter.ids && filter.ids.indexOf(event.id) === -1) {
|
|
if (!filter.ids.some((prefix) => event.id.startsWith(prefix))) {
|
|
return false;
|
|
}
|
|
}
|
|
if (filter.kinds && filter.kinds.indexOf(event.kind) === -1)
|
|
return false;
|
|
if (filter.authors && filter.authors.indexOf(event.pubkey) === -1) {
|
|
if (!filter.authors.some((prefix) => event.pubkey.startsWith(prefix))) {
|
|
return false;
|
|
}
|
|
}
|
|
for (let f in filter) {
|
|
if (f[0] === "#") {
|
|
let tagName = f.slice(1);
|
|
let values = filter[`#${tagName}`];
|
|
if (values && !event.tags.find(([t, v]) => t === f.slice(1) && values.indexOf(v) !== -1))
|
|
return false;
|
|
}
|
|
}
|
|
if (filter.since && event.created_at < filter.since)
|
|
return false;
|
|
if (filter.until && event.created_at > filter.until)
|
|
return false;
|
|
return true;
|
|
}
|
|
function matchFilters(filters, event) {
|
|
for (let i = 0; i < filters.length; i++) {
|
|
if (matchFilter(filters[i], event))
|
|
return true;
|
|
}
|
|
return false;
|
|
}
|
|
|
|
// fakejson.ts
|
|
function getHex64(json, field) {
|
|
let len = field.length + 3;
|
|
let idx = json.indexOf(`"${field}":`) + len;
|
|
let s = json.slice(idx).indexOf(`"`) + idx + 1;
|
|
return json.slice(s, s + 64);
|
|
}
|
|
function getSubscriptionId(json) {
|
|
let idx = json.slice(0, 22).indexOf(`"EVENT"`);
|
|
if (idx === -1)
|
|
return null;
|
|
let pstart = json.slice(idx + 7 + 1).indexOf(`"`);
|
|
if (pstart === -1)
|
|
return null;
|
|
let start = idx + 7 + 1 + pstart;
|
|
let pend = json.slice(start + 1, 80).indexOf(`"`);
|
|
if (pend === -1)
|
|
return null;
|
|
let end = start + 1 + pend;
|
|
return json.slice(start + 1, end);
|
|
}
|
|
|
|
// relay.ts
|
|
var newListeners = () => ({
|
|
connect: [],
|
|
disconnect: [],
|
|
error: [],
|
|
notice: [],
|
|
auth: []
|
|
});
|
|
function relayInit(url, options = {}) {
|
|
let { listTimeout = 3e3, getTimeout = 3e3, countTimeout = 3e3 } = options;
|
|
var ws;
|
|
var openSubs = {};
|
|
var listeners = newListeners();
|
|
var subListeners = {};
|
|
var pubListeners = {};
|
|
var connectionPromise;
|
|
async function connectRelay() {
|
|
if (connectionPromise)
|
|
return connectionPromise;
|
|
connectionPromise = new Promise((resolve, reject) => {
|
|
try {
|
|
ws = new WebSocket(url);
|
|
} catch (err) {
|
|
reject(err);
|
|
}
|
|
ws.onopen = () => {
|
|
listeners.connect.forEach((cb) => cb());
|
|
resolve();
|
|
};
|
|
ws.onerror = () => {
|
|
connectionPromise = void 0;
|
|
listeners.error.forEach((cb) => cb());
|
|
reject();
|
|
};
|
|
ws.onclose = async () => {
|
|
connectionPromise = void 0;
|
|
listeners.disconnect.forEach((cb) => cb());
|
|
};
|
|
let incomingMessageQueue = new MessageQueue();
|
|
let handleNextInterval;
|
|
ws.onmessage = (e) => {
|
|
incomingMessageQueue.enqueue(e.data);
|
|
if (!handleNextInterval) {
|
|
handleNextInterval = setInterval(handleNext, 0);
|
|
}
|
|
};
|
|
function handleNext() {
|
|
if (incomingMessageQueue.size === 0) {
|
|
clearInterval(handleNextInterval);
|
|
handleNextInterval = null;
|
|
return;
|
|
}
|
|
var json = incomingMessageQueue.dequeue();
|
|
if (!json)
|
|
return;
|
|
let subid = getSubscriptionId(json);
|
|
if (subid) {
|
|
let so = openSubs[subid];
|
|
if (so && so.alreadyHaveEvent && so.alreadyHaveEvent(getHex64(json, "id"), url)) {
|
|
return;
|
|
}
|
|
}
|
|
try {
|
|
let data = JSON.parse(json);
|
|
switch (data[0]) {
|
|
case "EVENT": {
|
|
let id2 = data[1];
|
|
let event = data[2];
|
|
if (validateEvent(event) && openSubs[id2] && (openSubs[id2].skipVerification || verifySignature(event)) && matchFilters(openSubs[id2].filters, event)) {
|
|
openSubs[id2];
|
|
(subListeners[id2]?.event || []).forEach((cb) => cb(event));
|
|
}
|
|
return;
|
|
}
|
|
case "COUNT":
|
|
let id = data[1];
|
|
let payload = data[2];
|
|
if (openSubs[id]) {
|
|
;
|
|
(subListeners[id]?.count || []).forEach((cb) => cb(payload));
|
|
}
|
|
return;
|
|
case "EOSE": {
|
|
let id2 = data[1];
|
|
if (id2 in subListeners) {
|
|
subListeners[id2].eose.forEach((cb) => cb());
|
|
subListeners[id2].eose = [];
|
|
}
|
|
return;
|
|
}
|
|
case "OK": {
|
|
let id2 = data[1];
|
|
let ok = data[2];
|
|
let reason = data[3] || "";
|
|
if (id2 in pubListeners) {
|
|
let { resolve: resolve2, reject: reject2 } = pubListeners[id2];
|
|
if (ok)
|
|
resolve2(null);
|
|
else
|
|
reject2(new Error(reason));
|
|
}
|
|
return;
|
|
}
|
|
case "NOTICE":
|
|
let notice = data[1];
|
|
listeners.notice.forEach((cb) => cb(notice));
|
|
return;
|
|
case "AUTH": {
|
|
let challenge = data[1];
|
|
listeners.auth?.forEach((cb) => cb(challenge));
|
|
return;
|
|
}
|
|
}
|
|
} catch (err) {
|
|
return;
|
|
}
|
|
}
|
|
});
|
|
return connectionPromise;
|
|
}
|
|
function connected() {
|
|
return ws?.readyState === 1;
|
|
}
|
|
async function connect() {
|
|
if (connected())
|
|
return;
|
|
await connectRelay();
|
|
}
|
|
async function trySend(params) {
|
|
let msg = JSON.stringify(params);
|
|
if (!connected()) {
|
|
await new Promise((resolve) => setTimeout(resolve, 1e3));
|
|
if (!connected()) {
|
|
return;
|
|
}
|
|
}
|
|
try {
|
|
ws.send(msg);
|
|
} catch (err) {
|
|
console.log(err);
|
|
}
|
|
}
|
|
const sub = (filters, {
|
|
verb = "REQ",
|
|
skipVerification = false,
|
|
alreadyHaveEvent = null,
|
|
id = Math.random().toString().slice(2)
|
|
} = {}) => {
|
|
let subid = id;
|
|
openSubs[subid] = {
|
|
id: subid,
|
|
filters,
|
|
skipVerification,
|
|
alreadyHaveEvent
|
|
};
|
|
trySend([verb, subid, ...filters]);
|
|
let subscription = {
|
|
sub: (newFilters, newOpts = {}) => sub(newFilters || filters, {
|
|
skipVerification: newOpts.skipVerification || skipVerification,
|
|
alreadyHaveEvent: newOpts.alreadyHaveEvent || alreadyHaveEvent,
|
|
id: subid
|
|
}),
|
|
unsub: () => {
|
|
delete openSubs[subid];
|
|
delete subListeners[subid];
|
|
trySend(["CLOSE", subid]);
|
|
},
|
|
on: (type, cb) => {
|
|
subListeners[subid] = subListeners[subid] || {
|
|
event: [],
|
|
count: [],
|
|
eose: []
|
|
};
|
|
subListeners[subid][type].push(cb);
|
|
},
|
|
off: (type, cb) => {
|
|
let listeners2 = subListeners[subid];
|
|
let idx = listeners2[type].indexOf(cb);
|
|
if (idx >= 0)
|
|
listeners2[type].splice(idx, 1);
|
|
},
|
|
get events() {
|
|
return eventsGenerator(subscription);
|
|
}
|
|
};
|
|
return subscription;
|
|
};
|
|
function _publishEvent(event, type) {
|
|
return new Promise((resolve, reject) => {
|
|
if (!event.id) {
|
|
reject(new Error(`event ${event} has no id`));
|
|
return;
|
|
}
|
|
let id = event.id;
|
|
trySend([type, event]);
|
|
pubListeners[id] = { resolve, reject };
|
|
});
|
|
}
|
|
return {
|
|
url,
|
|
sub,
|
|
on: (type, cb) => {
|
|
listeners[type].push(cb);
|
|
if (type === "connect" && ws?.readyState === 1) {
|
|
;
|
|
cb();
|
|
}
|
|
},
|
|
off: (type, cb) => {
|
|
let index = listeners[type].indexOf(cb);
|
|
if (index !== -1)
|
|
listeners[type].splice(index, 1);
|
|
},
|
|
list: (filters, opts) => new Promise((resolve) => {
|
|
let s = sub(filters, opts);
|
|
let events = [];
|
|
let timeout = setTimeout(() => {
|
|
s.unsub();
|
|
resolve(events);
|
|
}, listTimeout);
|
|
s.on("eose", () => {
|
|
s.unsub();
|
|
clearTimeout(timeout);
|
|
resolve(events);
|
|
});
|
|
s.on("event", (event) => {
|
|
events.push(event);
|
|
});
|
|
}),
|
|
get: (filter, opts) => new Promise((resolve) => {
|
|
let s = sub([filter], opts);
|
|
let timeout = setTimeout(() => {
|
|
s.unsub();
|
|
resolve(null);
|
|
}, getTimeout);
|
|
s.on("event", (event) => {
|
|
s.unsub();
|
|
clearTimeout(timeout);
|
|
resolve(event);
|
|
});
|
|
}),
|
|
count: (filters) => new Promise((resolve) => {
|
|
let s = sub(filters, { ...sub, verb: "COUNT" });
|
|
let timeout = setTimeout(() => {
|
|
s.unsub();
|
|
resolve(null);
|
|
}, countTimeout);
|
|
s.on("count", (event) => {
|
|
s.unsub();
|
|
clearTimeout(timeout);
|
|
resolve(event);
|
|
});
|
|
}),
|
|
async publish(event) {
|
|
await _publishEvent(event, "EVENT");
|
|
},
|
|
async auth(event) {
|
|
await _publishEvent(event, "AUTH");
|
|
},
|
|
connect,
|
|
close() {
|
|
listeners = newListeners();
|
|
subListeners = {};
|
|
pubListeners = {};
|
|
if (ws?.readyState === WebSocket.OPEN) {
|
|
ws.close();
|
|
}
|
|
},
|
|
get status() {
|
|
return ws?.readyState ?? 3;
|
|
}
|
|
};
|
|
}
|
|
async function* eventsGenerator(sub) {
|
|
let nextResolve;
|
|
const eventQueue = [];
|
|
const pushToQueue = (event) => {
|
|
if (nextResolve) {
|
|
nextResolve(event);
|
|
nextResolve = void 0;
|
|
} else {
|
|
eventQueue.push(event);
|
|
}
|
|
};
|
|
sub.on("event", pushToQueue);
|
|
try {
|
|
while (true) {
|
|
if (eventQueue.length > 0) {
|
|
yield eventQueue.shift();
|
|
} else {
|
|
const event = await new Promise((resolve) => {
|
|
nextResolve = resolve;
|
|
});
|
|
yield event;
|
|
}
|
|
}
|
|
} finally {
|
|
sub.off("event", pushToQueue);
|
|
}
|
|
}
|
|
export {
|
|
eventsGenerator,
|
|
relayInit
|
|
};
|