Manual Push
This commit is contained in:
@@ -5,7 +5,7 @@ const Corestore = require('corestore')
|
||||
const Hyperbee = require('hyperbee')
|
||||
const crypto = require('crypto')
|
||||
const process = require('process')
|
||||
const b4a = require('b4a') // optional but recommended for buffer handling
|
||||
|
||||
|
||||
const mode = process.argv[2] || 'reader'
|
||||
const storage = `./${mode}-storage`
|
||||
@@ -14,8 +14,14 @@ console.log(`Mode: ${mode}, storage: ${storage}`)
|
||||
|
||||
async function main() {
|
||||
const corestore = new Corestore(storage)
|
||||
// Use a named core so both writer and reader use the same logical core
|
||||
const core = corestore.get({ name: 'test-bee' })
|
||||
// Deterministic core key using seed-derived keypair (same public key every run)
|
||||
const sodium = require('sodium-native')
|
||||
const seed = crypto.createHash('sha256').update('test-bee').digest()
|
||||
const publicKey = Buffer.alloc(32)
|
||||
const secretKey = Buffer.alloc(64)
|
||||
sodium.crypto_sign_seed_keypair(publicKey, secretKey, seed)
|
||||
const keyPair = { publicKey, secretKey }
|
||||
const core = corestore.get(publicKey, { keyPair })
|
||||
|
||||
console.log('Core instance created')
|
||||
await core.ready()
|
||||
@@ -64,12 +70,13 @@ async function main() {
|
||||
const interval = setInterval(async () => {
|
||||
try {
|
||||
const ts = Date.now().toString().padStart(20, '0') // zero-pad for correct sorting
|
||||
const key = 'chat:' + ts
|
||||
const msg = {
|
||||
msg: `Hello from writer #${count} at ${new Date().toISOString()}`,
|
||||
peer: 'writer'
|
||||
}
|
||||
|
||||
await db.put(ts, msg)
|
||||
await db.put(key, msg)
|
||||
console.log(`✓ Put #${count} | key: ${ts.slice(-8)} | version: ${db.version}`)
|
||||
count++
|
||||
} catch (err) {
|
||||
@@ -86,30 +93,24 @@ async function main() {
|
||||
|
||||
} else {
|
||||
// === Reader mode ===
|
||||
const readInterval = setInterval(async () => {
|
||||
try {
|
||||
const now = Date.now()
|
||||
const fromTs = (now - 120000).toString().padStart(20, '0') // last 2 minutes
|
||||
|
||||
console.log(`Scanning from ${fromTs.slice(-8)} (version: ${db.version})`)
|
||||
|
||||
let found = 0
|
||||
for await (const entry of db.createReadStream({
|
||||
gte: fromTs,
|
||||
limit: 50 // prevent flooding the console
|
||||
})) {
|
||||
console.log(` ${entry.key.slice(-8)}: ${entry.value.msg} (${entry.value.peer})`)
|
||||
found++
|
||||
}
|
||||
console.log(`→ Found ${found} recent messages`)
|
||||
} catch (err) {
|
||||
console.error('Scan error:', err.message)
|
||||
}
|
||||
}, 5000)
|
||||
console.log('Setting up live read stream gt "chat:"')
|
||||
let found = 0
|
||||
const stream = db.createReadStream({
|
||||
gt: 'chat:',
|
||||
live: true
|
||||
})
|
||||
stream.on('data', (entry) => {
|
||||
console.log(` ${entry.key.slice(-28)}: ${entry.value.msg} (${entry.value.peer}) (v${db.version})`)
|
||||
found++
|
||||
})
|
||||
stream.on('end', () => {
|
||||
console.log('Live stream ended')
|
||||
})
|
||||
console.log('Live stream ready')
|
||||
|
||||
setTimeout(() => {
|
||||
console.log('Reader test complete')
|
||||
clearInterval(readInterval)
|
||||
// clearInterval(readInterval)
|
||||
shutdown()
|
||||
}, 30000).unref()
|
||||
}
|
||||
|
||||
@@ -1,64 +1,73 @@
|
||||
#!/usr/bin/env node
|
||||
const Hyperswarm = require('hyperswarm')
|
||||
'use strict'
|
||||
|
||||
const Corestore = require('corestore')
|
||||
const Hyperdrive = require('hyperdrive')
|
||||
const process = require('process')
|
||||
|
||||
const mode = process.argv[2] || 'reader'
|
||||
const storage = `./${mode}-storage`
|
||||
|
||||
console.log(`Mode: ${mode}, storage: ${storage}`)
|
||||
const Hyperswarm = require('hyperswarm')
|
||||
|
||||
async function main() {
|
||||
const corestore = new Corestore(storage)
|
||||
const drive = new Hyperdrive(corestore)
|
||||
console.log('Drive instance created')
|
||||
const isWriter = process.argv[2] === undefined
|
||||
const storageDir = isWriter ? './storage-writer' : './storage-reader'
|
||||
const store = new Corestore(storageDir)
|
||||
await store.ready()
|
||||
|
||||
let drive
|
||||
if (isWriter) {
|
||||
drive = new Hyperdrive(store)
|
||||
} else {
|
||||
const key = Buffer.from(process.argv[2], 'hex')
|
||||
drive = new Hyperdrive(store, key, { sparse: true })
|
||||
}
|
||||
await drive.ready()
|
||||
console.log('Drive ready, key:', drive.key.toString('hex'))
|
||||
console.log('Discovery key:', drive.discoveryKey.toString('hex'))
|
||||
console.log('Initial version:', drive.version)
|
||||
|
||||
if (isWriter) {
|
||||
const content = 'Test P2P file share content\n'
|
||||
await drive.put('test.txt', Buffer.from(content))
|
||||
console.log('✅ File written: test.txt')
|
||||
console.log('📤 Share this drive key with readers: ' + drive.key.toString('hex'))
|
||||
} else {
|
||||
console.log('🔍 Connecting to drive key: ' + process.argv[2])
|
||||
}
|
||||
|
||||
const swarm = new Hyperswarm()
|
||||
const topic = drive.discoveryKey // Use drive.discoveryKey for specific drive sync
|
||||
swarm.join(topic, { server: true, client: true })
|
||||
swarm.on('connection', (conn, info) => {
|
||||
console.log('New P2P connection:', !!(info.client), !!(info.server))
|
||||
swarm.on('connection', (conn) => {
|
||||
console.log('🔗 New peer connection')
|
||||
drive.replicate(conn)
|
||||
})
|
||||
swarm.on('updated', () => {
|
||||
console.log(`Swarm has ${swarm.connections.size} connections`)
|
||||
})
|
||||
|
||||
if (mode === 'writer') {
|
||||
let count = 0
|
||||
const interval = setInterval(async () => {
|
||||
swarm.join(drive.discoveryKey)
|
||||
console.log('🌐 Joined swarm on discovery key: ' + drive.discoveryKey.toString('hex'))
|
||||
|
||||
if (!isWriter) {
|
||||
const checkFile = async () => {
|
||||
try {
|
||||
const content = Buffer.from(`Hello from writer #${count} at ${new Date().toISOString()}`)
|
||||
await drive.put(`/msg-${count}.txt`, content)
|
||||
console.log(`Put file /msg-${count}.txt (${content.length} bytes)`)
|
||||
console.log('Drive version:', drive.version)
|
||||
count++
|
||||
const content = await drive.get('test.txt')
|
||||
if (content) {
|
||||
console.log('📥 File received successfully:')
|
||||
console.log(content.toString())
|
||||
process.exit(0)
|
||||
}
|
||||
} catch (err) {
|
||||
console.error('Put error:', err.message)
|
||||
// File not yet available
|
||||
}
|
||||
}, 3000)
|
||||
}
|
||||
|
||||
setTimeout(() => {
|
||||
console.log('Test complete')
|
||||
clearInterval(interval)
|
||||
process.exit(0)
|
||||
}, 30000).unref()
|
||||
} else {
|
||||
// reader
|
||||
const readInterval = setInterval(async () => {
|
||||
console.log('Scanning files:')
|
||||
for await (const file of drive.list('/')) {
|
||||
const content = await drive.get(file.name)
|
||||
console.log(` ${file.name}: ${content ? content.length : 0} bytes`)
|
||||
}
|
||||
console.log('Drive version:', drive.version)
|
||||
}, 5000)
|
||||
// Initial check
|
||||
await checkFile()
|
||||
|
||||
// Poll periodically
|
||||
const interval = setInterval(checkFile, 2000)
|
||||
|
||||
// Also react to updates
|
||||
const onUpdate = () => checkFile()
|
||||
drive.on('update', onUpdate)
|
||||
}
|
||||
|
||||
// Keep writer alive
|
||||
await new Promise(resolve => {})
|
||||
}
|
||||
|
||||
main().catch(console.error)
|
||||
main().catch((err) => {
|
||||
console.error('❌ Error:', err.message)
|
||||
process.exit(1)
|
||||
})
|
||||
|
||||
+2
-1
@@ -12,7 +12,8 @@
|
||||
"corestore": "^7.9.2",
|
||||
"hyperdrive": "^13.3.2",
|
||||
"hyperswarm": "^4.17.0"
|
||||
}
|
||||
},
|
||||
"devDependencies": {}
|
||||
},
|
||||
"node_modules/@hyperswarm/secret-stream": {
|
||||
"version": "6.9.1",
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
{
|
||||
"name": "hyperdrive-fileshare",
|
||||
"version": "1.0.0",
|
||||
"description": "",
|
||||
"description": "P2P file sharing prototype using Hyperdrive + Hyperswarm.",
|
||||
"main": "index.js",
|
||||
"scripts": {
|
||||
"test": "echo \"Error: no test specified\" && exit 1"
|
||||
|
||||
Reference in New Issue
Block a user