Repository navigation
Expand file tree
/
Copy pathswarm.js
More file actions
111 lines (95 loc) · 3.04 KB
/
Copy pathswarm.js
File metadata and controls
111 lines (95 loc) · 3.04 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
import createDB from './db.js'
import createBee from './bee.js'
import { createHash } from 'crypto'
import { validateEvent, isPersistent } from './nostr_events.js'
import goodbye from './goodbye.js'
const prefix = 'hyper-nostr-'
export default async function createSwarm (sdk, _topic) {
const topic = prefix + _topic
const subscriptions = new Map()
const bee = await createBee(sdk, topic)
const { handleEvent, queryEvents } = await createDB(bee)
const knownDBs = new Set()
knownDBs.add(bee.autobase.localInput.url)
const discovery = await sdk.get(createTopicBuffer(topic))
goodbye(_ => {
console.log('closing discovery core of', topic)
return discovery.close()
})
const events = discovery.registerExtension(topic, {
encoding: 'json',
onmessage: streamEvent
})
const DBBroadcast = discovery.registerExtension(topic + '-sync', {
encoding: 'json',
onmessage: async (message) => {
let sawNew = false
for (const url of message) {
if (knownDBs.has(url)) continue
sawNew = true
await handleNewDB(url)
}
if (sawNew) {
broadcastDBs()
logDBs()
await update()
}
}
})
const requestSync = discovery.registerExtension(topic + '-request-sync', {
encoding: 'json',
onmessage: broadcastDBs
})
discovery.on('peer-add', initConnection)
discovery.on('peer-remove', logPeers)
initConnection()
console.log(`swarm ${topic} created with hyper!`)
return { subscriptions, sendEvent, queryEvents, sendQueryToSubscription, update }
function initConnection () {
requestSync.broadcast('')
logPeers()
logDBs()
broadcastDBs()
}
function logPeers () {
console.log(`${discovery.peers.length} peers on ${_topic}!`)
}
function logDBs () {
console.log('DB count:', bee.autobase.inputs.filter(core => core.readable).length)
}
function streamEvent (event) {
subscriptions.forEach((sub, key) => {
if (validateEvent(event, sub.filters)) sendEventTo(event, sub, key)
})
}
function sendEvent (event) {
events.broadcast(event)
streamEvent(event)
if (isPersistent(event)) return handleEvent(event)
}
function broadcastDBs () {
DBBroadcast.broadcast(Array.from(knownDBs))
}
function update () {
return bee.autobase.view.update()
}
async function handleNewDB (url) {
knownDBs.add(url)
await bee.autobase.addInput(await sdk.get(url))
subscriptions.forEach((sub, key) => sendQueryToSubscription(sub, key, { hasLimit: false }))
}
async function sendQueryToSubscription (sub, key, { hasLimit } = { hasLimit: true }) {
return queryEvents(sub.filters, { hasLimit }).then(events => {
for (let i = events.length - 1; i >= 0; i--) sendEventTo(events[i], sub, key)
})
}
}
function sendEventTo (event, sub, key) {
if (!sub.receivedEvents.has(event.id)) {
sub.socket.send(`["EVENT", "${key}", ${JSON.stringify((delete event._id, event))}]`)
sub.receivedEvents.add(event.id)
}
}
function createTopicBuffer (topic) {
return createHash('sha256').update(topic).digest()
}