+473
-63
@@ -63,7 +63,17 @@ const connections = new Map();
|
||||
const connToSwarm = new Map(); // connId -> swarmId
|
||||
/** swarmId -> (peerKeyHex -> connId) for one connection per peer per swarm */
|
||||
const peerConnections = new Map();
|
||||
/** swarmId -> { mode: 'off'|'allowlist'|'denylist', keys: Set<string> } */
|
||||
const swarmFirewalls = new Map();
|
||||
/** swarmId -> { enabled: boolean, coreKeyHex?: string, resourceId?: string } */
|
||||
const swarmAutoReplicate = new Map();
|
||||
let nextConnId = 0;
|
||||
let nextResourceId = 0;
|
||||
let nextHrpcStreamId = 0;
|
||||
|
||||
// Named Corestore resources (one Corestore; namespaces / keys)
|
||||
/** @type {Map<string, { kind: string, ref: any, name?: string, keyHex?: string }>} */
|
||||
const resources = new Map();
|
||||
|
||||
// Lazy corestore and default data structures
|
||||
let corestore = null;
|
||||
@@ -73,6 +83,41 @@ let defaultDrive = null;
|
||||
let defaultAutobase = null;
|
||||
let defaultHyperdb = null;
|
||||
|
||||
// Prefer Hyperschema-generated Hyperdb definition; fall back to minimal stub.
|
||||
let hyperdbDefinition = minimalDefinition;
|
||||
try {
|
||||
hyperdbDefinition = require('./spec/hyperdb');
|
||||
} catch (_) {}
|
||||
|
||||
function normalizePublicKeyHex(key) {
|
||||
if (!key) return null;
|
||||
if (typeof key === 'string') {
|
||||
const hex = key.length === 64 ? key : Buffer.from(key, 'base64').toString('hex');
|
||||
return hex.length === 64 ? hex.toLowerCase() : null;
|
||||
}
|
||||
try {
|
||||
return b4a.toString(key, 'hex').toLowerCase();
|
||||
} catch (_) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
function makeFirewallFn(swarmId) {
|
||||
return (remotePublicKey) => {
|
||||
const cfg = swarmFirewalls.get(swarmId);
|
||||
if (!cfg || cfg.mode === 'off') return false; // allow
|
||||
const hex = normalizePublicKeyHex(remotePublicKey);
|
||||
if (!hex) return true; // reject unknown
|
||||
if (cfg.mode === 'allowlist') return !cfg.keys.has(hex); // true = reject
|
||||
if (cfg.mode === 'denylist') return cfg.keys.has(hex);
|
||||
return false;
|
||||
};
|
||||
}
|
||||
|
||||
function generateResourceId(kind) {
|
||||
return `${kind}_${nextResourceId++}_${Date.now()}`;
|
||||
}
|
||||
|
||||
function getCorestore() {
|
||||
if (!corestore) {
|
||||
const storagePath = process.env.BRIDGE_SWARM_STORAGE || path.join(process.cwd(), 'bridge-swarm-storage');
|
||||
@@ -136,12 +181,96 @@ async function getDefaultHyperdb() {
|
||||
const hyperdbCore = store.get({ name: 'hyperdb' });
|
||||
await hyperdbCore.ready();
|
||||
// HyperDB.bee applies def.compat() internally (hyperdb 6+)
|
||||
defaultHyperdb = HyperDB.bee(hyperdbCore, minimalDefinition, { autoUpdate: true });
|
||||
defaultHyperdb = HyperDB.bee(hyperdbCore, hyperdbDefinition, { autoUpdate: true });
|
||||
await defaultHyperdb.ready();
|
||||
}
|
||||
return defaultHyperdb;
|
||||
}
|
||||
|
||||
async function resolveCore(payload = {}) {
|
||||
if (payload.resourceId) {
|
||||
const res = resources.get(payload.resourceId);
|
||||
if (!res || res.kind !== 'core') throw new Error('Core resource not found: ' + payload.resourceId);
|
||||
return res.ref;
|
||||
}
|
||||
if (payload.coreKeyHex) {
|
||||
const store = getCorestore();
|
||||
await store.ready();
|
||||
const core = store.get(Buffer.from(payload.coreKeyHex, 'hex'));
|
||||
await core.ready();
|
||||
return core;
|
||||
}
|
||||
if (payload.name) {
|
||||
const store = getCorestore();
|
||||
const core = store.get({ name: String(payload.name) });
|
||||
await core.ready();
|
||||
return core;
|
||||
}
|
||||
return getDefaultCore();
|
||||
}
|
||||
|
||||
async function resolveBee(payload = {}) {
|
||||
if (payload.resourceId) {
|
||||
const res = resources.get(payload.resourceId);
|
||||
if (!res || res.kind !== 'bee') throw new Error('Bee resource not found: ' + payload.resourceId);
|
||||
return res.ref;
|
||||
}
|
||||
return getDefaultBee();
|
||||
}
|
||||
|
||||
async function resolveDrive(payload = {}) {
|
||||
if (payload.resourceId) {
|
||||
const res = resources.get(payload.resourceId);
|
||||
if (!res || res.kind !== 'drive') throw new Error('Drive resource not found: ' + payload.resourceId);
|
||||
return res.ref;
|
||||
}
|
||||
return getDefaultDrive();
|
||||
}
|
||||
|
||||
async function resolveAutobase(payload = {}) {
|
||||
if (payload.resourceId) {
|
||||
const res = resources.get(payload.resourceId);
|
||||
if (!res || res.kind !== 'autobase') throw new Error('Autobase resource not found: ' + payload.resourceId);
|
||||
return res.ref;
|
||||
}
|
||||
return getDefaultAutobase();
|
||||
}
|
||||
|
||||
async function resolveHyperdb(payload = {}) {
|
||||
if (payload.resourceId) {
|
||||
const res = resources.get(payload.resourceId);
|
||||
if (!res || res.kind !== 'hyperdb') throw new Error('Hyperdb resource not found: ' + payload.resourceId);
|
||||
return res.ref;
|
||||
}
|
||||
return getDefaultHyperdb();
|
||||
}
|
||||
|
||||
/**
|
||||
* Stop browser forwarding and attach Hypercore replication on a connection.
|
||||
* Shared by attachReplication and auto-replicate.
|
||||
*/
|
||||
async function attachReplicationToEntry(entry, connId, core) {
|
||||
if (entry.protomux) return;
|
||||
entry.forwarding = false;
|
||||
if (entry._listeners) {
|
||||
entry.socket.removeListener('data', entry._listeners.onData);
|
||||
entry.socket.removeListener('end', entry._listeners.onEnd);
|
||||
entry.socket.removeListener('error', entry._listeners.onError);
|
||||
entry.socket.removeListener('close', entry._listeners.onClose);
|
||||
}
|
||||
entry.socket.once('close', () => {
|
||||
const sid = connToSwarm.get(connId);
|
||||
const pm = sid ? peerConnections.get(sid) : null;
|
||||
if (pm && pm.get(entry.peerKeyHex) === connId) pm.delete(entry.peerKeyHex);
|
||||
connections.delete(connId);
|
||||
connToSwarm.delete(connId);
|
||||
});
|
||||
const mux = new Protomux(entry.socket);
|
||||
entry.protomux = mux;
|
||||
await core.ready();
|
||||
core.replicate(mux);
|
||||
}
|
||||
|
||||
function generateConnId() {
|
||||
return `conn_${nextConnId++}_${Date.now()}`;
|
||||
}
|
||||
@@ -176,7 +305,15 @@ async function handleMessageAsync(send, msg) {
|
||||
reply({ ok: true });
|
||||
return;
|
||||
}
|
||||
const swarm = new Hyperswarm(options);
|
||||
// Strip non-serializable / host-managed options; firewall is applied via setFirewall
|
||||
const { firewall: _ignoredFw, ...safeOpts } = options || {};
|
||||
if (!swarmFirewalls.has(swarmId)) {
|
||||
swarmFirewalls.set(swarmId, { mode: 'off', keys: new Set() });
|
||||
}
|
||||
const swarm = new Hyperswarm({
|
||||
...safeOpts,
|
||||
firewall: makeFirewallFn(swarmId),
|
||||
});
|
||||
swarm.listen().catch((err) => {
|
||||
log('Swarm listen error:', err.message);
|
||||
});
|
||||
@@ -241,11 +378,39 @@ async function handleMessageAsync(send, msg) {
|
||||
socket.on('error', onError);
|
||||
socket.on('close', onClose);
|
||||
log('new connection', connId, 'from peer', peerKeyHex.slice(0, 8));
|
||||
emit('connection', {
|
||||
connId,
|
||||
swarmId,
|
||||
peerInfo: peerInfoToObject(peerInfo),
|
||||
});
|
||||
|
||||
const auto = swarmAutoReplicate.get(swarmId);
|
||||
if (auto && auto.enabled) {
|
||||
(async () => {
|
||||
try {
|
||||
const core = await resolveCore({
|
||||
coreKeyHex: auto.coreKeyHex,
|
||||
resourceId: auto.resourceId,
|
||||
});
|
||||
await attachReplicationToEntry(entry, connId, core);
|
||||
emit('connection', {
|
||||
connId,
|
||||
swarmId,
|
||||
peerInfo: peerInfoToObject(peerInfo),
|
||||
autoReplicated: true,
|
||||
});
|
||||
} catch (err) {
|
||||
log('auto-replicate failed:', err.message);
|
||||
emit('connection', {
|
||||
connId,
|
||||
swarmId,
|
||||
peerInfo: peerInfoToObject(peerInfo),
|
||||
autoReplicateError: err.message,
|
||||
});
|
||||
}
|
||||
})();
|
||||
} else {
|
||||
emit('connection', {
|
||||
connId,
|
||||
swarmId,
|
||||
peerInfo: peerInfoToObject(peerInfo),
|
||||
});
|
||||
}
|
||||
});
|
||||
swarms.set(swarmId, swarm);
|
||||
reply({ ok: true });
|
||||
@@ -286,6 +451,85 @@ async function handleMessageAsync(send, msg) {
|
||||
break;
|
||||
}
|
||||
|
||||
case 'setFirewall': {
|
||||
const { swarmId, mode = 'off', keys = [] } = payload;
|
||||
if (!swarmId) {
|
||||
reply({ ok: false, error: 'swarmId required' });
|
||||
return;
|
||||
}
|
||||
const m = mode === 'allowlist' || mode === 'denylist' ? mode : 'off';
|
||||
const set = new Set();
|
||||
for (const k of keys) {
|
||||
const hex = normalizePublicKeyHex(k);
|
||||
if (hex) set.add(hex);
|
||||
}
|
||||
swarmFirewalls.set(swarmId, { mode: m, keys: set });
|
||||
// Hyperswarm reads firewall via bound function; no re-init needed.
|
||||
reply({ ok: true, mode: m, keys: [...set] });
|
||||
break;
|
||||
}
|
||||
|
||||
case 'banPeer': {
|
||||
const { swarmId, publicKeyHex, banned = true } = payload;
|
||||
const swarm = swarms.get(swarmId);
|
||||
if (!swarm) {
|
||||
reply({ ok: false, error: 'Swarm not found' });
|
||||
return;
|
||||
}
|
||||
const hex = normalizePublicKeyHex(publicKeyHex);
|
||||
if (!hex) {
|
||||
reply({ ok: false, error: 'publicKeyHex required (64 hex chars)' });
|
||||
return;
|
||||
}
|
||||
const peerInfo = swarm.peers.get(hex);
|
||||
if (peerInfo && typeof peerInfo.ban === 'function') {
|
||||
peerInfo.ban(!!banned);
|
||||
}
|
||||
// Also maintain denylist so future reconnects are firewalled
|
||||
let cfg = swarmFirewalls.get(swarmId);
|
||||
if (!cfg) {
|
||||
cfg = { mode: 'denylist', keys: new Set() };
|
||||
swarmFirewalls.set(swarmId, cfg);
|
||||
}
|
||||
if (banned) {
|
||||
if (cfg.mode === 'off') cfg.mode = 'denylist';
|
||||
cfg.keys.add(hex);
|
||||
// Drop active connection if any
|
||||
const pm = peerConnections.get(swarmId);
|
||||
const cid = pm && pm.get(hex);
|
||||
if (cid) {
|
||||
const entry = connections.get(cid);
|
||||
if (entry) {
|
||||
try { entry.socket.destroy(); } catch (_) {}
|
||||
}
|
||||
}
|
||||
} else {
|
||||
cfg.keys.delete(hex);
|
||||
}
|
||||
reply({ ok: true, banned: !!banned, publicKeyHex: hex });
|
||||
break;
|
||||
}
|
||||
|
||||
case 'setAutoReplicate': {
|
||||
const { swarmId, enabled = false, coreKeyHex, resourceId } = payload;
|
||||
if (!swarmId) {
|
||||
reply({ ok: false, error: 'swarmId required' });
|
||||
return;
|
||||
}
|
||||
if (!enabled) {
|
||||
swarmAutoReplicate.delete(swarmId);
|
||||
reply({ ok: true, enabled: false });
|
||||
return;
|
||||
}
|
||||
swarmAutoReplicate.set(swarmId, {
|
||||
enabled: true,
|
||||
coreKeyHex: coreKeyHex || undefined,
|
||||
resourceId: resourceId || undefined,
|
||||
});
|
||||
reply({ ok: true, enabled: true, coreKeyHex: coreKeyHex || null, resourceId: resourceId || null });
|
||||
break;
|
||||
}
|
||||
|
||||
case 'leave': {
|
||||
const { swarmId, topic: topicHex } = payload;
|
||||
const swarm = swarms.get(swarmId);
|
||||
@@ -333,7 +577,7 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
|
||||
case 'attachReplication': {
|
||||
const { connId, coreKeyHex } = payload;
|
||||
const { connId, coreKeyHex, resourceId } = payload;
|
||||
const entry = connections.get(connId);
|
||||
if (!entry) {
|
||||
reply({ ok: false, error: 'Connection not found. Refresh the page if the host was restarted.' });
|
||||
@@ -344,27 +588,8 @@ async function handleMessageAsync(send, msg) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
entry.forwarding = false;
|
||||
entry.socket.removeListener('data', entry._listeners.onData);
|
||||
entry.socket.removeListener('end', entry._listeners.onEnd);
|
||||
entry.socket.removeListener('error', entry._listeners.onError);
|
||||
entry.socket.removeListener('close', entry._listeners.onClose);
|
||||
entry.socket.once('close', () => {
|
||||
const sid = connToSwarm.get(connId);
|
||||
const pm = sid ? peerConnections.get(sid) : null;
|
||||
if (pm && pm.get(entry.peerKeyHex) === connId) pm.delete(entry.peerKeyHex);
|
||||
connections.delete(connId);
|
||||
connToSwarm.delete(connId);
|
||||
});
|
||||
const mux = new Protomux(entry.socket);
|
||||
entry.protomux = mux;
|
||||
const store = getCorestore();
|
||||
await store.ready();
|
||||
const core = coreKeyHex
|
||||
? store.get(Buffer.from(coreKeyHex, 'hex'))
|
||||
: getDefaultCore();
|
||||
await core.ready();
|
||||
core.replicate(mux);
|
||||
const core = await resolveCore({ coreKeyHex, resourceId });
|
||||
await attachReplicationToEntry(entry, connId, core);
|
||||
reply({ ok: true });
|
||||
} catch (err) {
|
||||
reply({ ok: false, error: err.message });
|
||||
@@ -485,16 +710,12 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
|
||||
case 'hrpcInvoke': {
|
||||
const { connId, method, args = {} } = payload;
|
||||
const { connId, method, args = {}, streamId: clientStreamId } = payload;
|
||||
const entry = connections.get(connId);
|
||||
if (!entry || !entry.hrpc) {
|
||||
reply({ ok: false, error: 'Connection not found or HRPC not attached. Refresh the page if the host was restarted.' });
|
||||
return;
|
||||
}
|
||||
if (method !== 'ping') {
|
||||
reply({ ok: false, error: 'Only ping is supported from the browser (unary only)' });
|
||||
return;
|
||||
}
|
||||
if (entry.hrpcChannel) {
|
||||
try {
|
||||
await Promise.race([
|
||||
@@ -506,16 +727,83 @@ async function handleMessageAsync(send, msg) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
const timeoutMs = 15000;
|
||||
const timeoutPromise = new Promise((_, reject) => {
|
||||
setTimeout(() => reject(new Error('Ping timed out. Ensure the other tab also enabled HRPC on its connection.')), timeoutMs);
|
||||
});
|
||||
|
||||
const timeoutMs = 30000;
|
||||
const withTimeout = (p, label) => Promise.race([
|
||||
p,
|
||||
new Promise((_, reject) => setTimeout(() => reject(new Error(label + ' timed out.')), timeoutMs)),
|
||||
]);
|
||||
const swarmIdForConn = connToSwarm.get(connId);
|
||||
|
||||
try {
|
||||
const result = await Promise.race([
|
||||
entry.hrpc.ping(args),
|
||||
timeoutPromise
|
||||
]);
|
||||
reply({ ok: true, result });
|
||||
if (method === 'ping') {
|
||||
const result = await withTimeout(entry.hrpc.ping(args), 'Ping');
|
||||
reply({ ok: true, result });
|
||||
return;
|
||||
}
|
||||
|
||||
if (method === 'notify') {
|
||||
entry.hrpc.notify(args);
|
||||
reply({ ok: true });
|
||||
return;
|
||||
}
|
||||
|
||||
if (method === 'fetchStream') {
|
||||
const streamId = clientStreamId || `hrpc_${nextHrpcStreamId++}_${Date.now()}`;
|
||||
reply({ ok: true, streamId, streaming: true });
|
||||
const responseStream = entry.hrpc.fetchStream(args);
|
||||
responseStream.on('data', (chunk) => {
|
||||
emit('hrpc-chunk', { connId, swarmId: swarmIdForConn, streamId, method, chunk });
|
||||
});
|
||||
responseStream.on('end', () => {
|
||||
emit('hrpc-end', { connId, swarmId: swarmIdForConn, streamId, method });
|
||||
});
|
||||
responseStream.on('error', (err) => {
|
||||
emit('hrpc-error', { connId, swarmId: swarmIdForConn, streamId, method, message: err.message });
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
if (method === 'streamSum') {
|
||||
const streamId = clientStreamId || `hrpc_${nextHrpcStreamId++}_${Date.now()}`;
|
||||
const chunks = Array.isArray(args.chunks) ? args.chunks : [];
|
||||
reply({ ok: true, streamId, streaming: true });
|
||||
try {
|
||||
const reqStream = entry.hrpc.streamSum();
|
||||
for (const chunk of chunks) reqStream.write(chunk);
|
||||
reqStream.end();
|
||||
const result = await withTimeout(reqStream.reply(), 'streamSum');
|
||||
emit('hrpc-end', { connId, swarmId: swarmIdForConn, streamId, method, result });
|
||||
} catch (err) {
|
||||
emit('hrpc-error', { connId, swarmId: swarmIdForConn, streamId, method, message: err.message });
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (method === 'duplex') {
|
||||
const streamId = clientStreamId || `hrpc_${nextHrpcStreamId++}_${Date.now()}`;
|
||||
const chunks = Array.isArray(args.chunks) ? args.chunks : [];
|
||||
reply({ ok: true, streamId, streaming: true });
|
||||
try {
|
||||
const stream = entry.hrpc.duplex();
|
||||
stream.on('data', (chunk) => {
|
||||
emit('hrpc-chunk', { connId, swarmId: swarmIdForConn, streamId, method, chunk });
|
||||
});
|
||||
stream.on('end', () => {
|
||||
emit('hrpc-end', { connId, swarmId: swarmIdForConn, streamId, method });
|
||||
});
|
||||
stream.on('error', (err) => {
|
||||
emit('hrpc-error', { connId, swarmId: swarmIdForConn, streamId, method, message: err.message });
|
||||
});
|
||||
for (const chunk of chunks) stream.write(chunk);
|
||||
stream.end();
|
||||
} catch (err) {
|
||||
emit('hrpc-error', { connId, swarmId: swarmIdForConn, streamId, method, message: err.message });
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
reply({ ok: false, error: 'Unknown HRPC method: ' + method + '. Supported: ping, notify, fetchStream, streamSum, duplex' });
|
||||
} catch (err) {
|
||||
reply({ ok: false, error: err.message });
|
||||
}
|
||||
@@ -542,6 +830,8 @@ async function handleMessageAsync(send, msg) {
|
||||
connToSwarm.delete(cid);
|
||||
}
|
||||
peerConnections.delete(swarmId);
|
||||
swarmFirewalls.delete(swarmId);
|
||||
swarmAutoReplicate.delete(swarmId);
|
||||
await swarm.destroy();
|
||||
swarms.delete(swarmId);
|
||||
}
|
||||
@@ -549,13 +839,30 @@ async function handleMessageAsync(send, msg) {
|
||||
break;
|
||||
}
|
||||
|
||||
// --- Hypercore RPC ---
|
||||
case 'coreInfo': {
|
||||
// --- Named resources (one Corestore) ---
|
||||
case 'coreOpen': {
|
||||
try {
|
||||
const core = getDefaultCore();
|
||||
const store = getCorestore();
|
||||
await store.ready();
|
||||
let core;
|
||||
if (payload.keyHex) {
|
||||
core = store.get(Buffer.from(payload.keyHex, 'hex'));
|
||||
} else {
|
||||
const name = payload.name != null ? String(payload.name) : 'default';
|
||||
const ns = payload.namespace ? store.namespace(String(payload.namespace)) : store;
|
||||
core = ns.get({ name });
|
||||
}
|
||||
await core.ready();
|
||||
const resourceId = generateResourceId('core');
|
||||
resources.set(resourceId, {
|
||||
kind: 'core',
|
||||
ref: core,
|
||||
name: payload.name,
|
||||
keyHex: b4a.toString(core.key, 'hex'),
|
||||
});
|
||||
reply({
|
||||
ok: true,
|
||||
resourceId,
|
||||
key: b4a.toString(core.key, 'hex'),
|
||||
length: core.length,
|
||||
writable: core.writable,
|
||||
@@ -565,9 +872,109 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'beeOpen': {
|
||||
try {
|
||||
const store = getCorestore();
|
||||
await store.ready();
|
||||
const name = payload.name != null ? String(payload.name) : 'bee';
|
||||
const ns = payload.namespace ? store.namespace(String(payload.namespace)) : store;
|
||||
const core = ns.get({ name });
|
||||
await core.ready();
|
||||
const bee = new Hyperbee(core, { keyEncoding: 'utf-8', valueEncoding: 'utf-8' });
|
||||
await bee.ready();
|
||||
const resourceId = generateResourceId('bee');
|
||||
resources.set(resourceId, { kind: 'bee', ref: bee, name });
|
||||
reply({ ok: true, resourceId, key: b4a.toString(core.key, 'hex') });
|
||||
} catch (err) {
|
||||
reply({ ok: false, error: err.message });
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'driveOpen': {
|
||||
try {
|
||||
const store = getCorestore();
|
||||
await store.ready();
|
||||
const ns = payload.namespace
|
||||
? store.namespace(String(payload.namespace))
|
||||
: (payload.name ? store.namespace(String(payload.name)) : store);
|
||||
const drive = new Hyperdrive(ns);
|
||||
await drive.ready();
|
||||
const resourceId = generateResourceId('drive');
|
||||
resources.set(resourceId, { kind: 'drive', ref: drive, name: payload.name });
|
||||
reply({ ok: true, resourceId, key: drive.key ? b4a.toString(drive.key, 'hex') : null });
|
||||
} catch (err) {
|
||||
reply({ ok: false, error: err.message });
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'autobaseOpen': {
|
||||
try {
|
||||
const store = getCorestore();
|
||||
await store.ready();
|
||||
const ns = payload.namespace
|
||||
? store.namespace(String(payload.namespace))
|
||||
: (payload.name ? store.namespace(String(payload.name)) : store);
|
||||
const viewName = payload.viewName != null ? String(payload.viewName) : 'autobase-view';
|
||||
const base = new Autobase(ns, null, {
|
||||
open(s) {
|
||||
return s.get({ name: viewName });
|
||||
},
|
||||
async apply(nodes, view) {
|
||||
for (const node of nodes) {
|
||||
if (node.value == null) continue;
|
||||
const data = typeof node.value === 'string' ? Buffer.from(node.value, 'utf8') : b4a.from(node.value);
|
||||
await view.append(data);
|
||||
}
|
||||
},
|
||||
});
|
||||
await base.ready();
|
||||
const resourceId = generateResourceId('autobase');
|
||||
resources.set(resourceId, { kind: 'autobase', ref: base, name: payload.name });
|
||||
reply({ ok: true, resourceId, length: base.length });
|
||||
} catch (err) {
|
||||
reply({ ok: false, error: err.message });
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'hyperdbOpen': {
|
||||
try {
|
||||
const store = getCorestore();
|
||||
await store.ready();
|
||||
const name = payload.name != null ? String(payload.name) : 'hyperdb';
|
||||
const ns = payload.namespace ? store.namespace(String(payload.namespace)) : store;
|
||||
const hyperdbCore = ns.get({ name });
|
||||
await hyperdbCore.ready();
|
||||
const db = HyperDB.bee(hyperdbCore, hyperdbDefinition, { autoUpdate: true });
|
||||
await db.ready();
|
||||
const resourceId = generateResourceId('hyperdb');
|
||||
resources.set(resourceId, { kind: 'hyperdb', ref: db, name });
|
||||
reply({ ok: true, resourceId, key: b4a.toString(hyperdbCore.key, 'hex') });
|
||||
} catch (err) {
|
||||
reply({ ok: false, error: err.message });
|
||||
}
|
||||
break;
|
||||
}
|
||||
|
||||
// --- Hypercore RPC ---
|
||||
case 'coreInfo': {
|
||||
try {
|
||||
const core = await resolveCore(payload);
|
||||
await core.ready();
|
||||
reply({
|
||||
ok: true,
|
||||
key: b4a.toString(core.key, 'hex'),
|
||||
length: core.length,
|
||||
writable: core.writable,
|
||||
resourceId: payload.resourceId || null,
|
||||
});
|
||||
} catch (err) {
|
||||
reply({ ok: false, error: err.message });
|
||||
}
|
||||
break;
|
||||
}
|
||||
case 'coreAppend': {
|
||||
try {
|
||||
const core = getDefaultCore();
|
||||
const core = await resolveCore(payload);
|
||||
await core.ready();
|
||||
const data = Buffer.from(payload.data || payload.base64 || '', 'base64');
|
||||
await core.append(data);
|
||||
@@ -579,7 +986,7 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
case 'coreGet': {
|
||||
try {
|
||||
const core = getDefaultCore();
|
||||
const core = await resolveCore(payload);
|
||||
await core.ready();
|
||||
const index = payload.index;
|
||||
const block = await core.get(index);
|
||||
@@ -593,7 +1000,7 @@ async function handleMessageAsync(send, msg) {
|
||||
// --- Hyperbee RPC ---
|
||||
case 'beeGet': {
|
||||
try {
|
||||
const bee = await getDefaultBee();
|
||||
const bee = await resolveBee(payload);
|
||||
const key = payload.key;
|
||||
const entry = await bee.get(key);
|
||||
if (!entry) {
|
||||
@@ -613,7 +1020,7 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
case 'beePut': {
|
||||
try {
|
||||
const bee = await getDefaultBee();
|
||||
const bee = await resolveBee(payload);
|
||||
await bee.put(payload.key, payload.value != null ? payload.value : '');
|
||||
reply({ ok: true });
|
||||
} catch (err) {
|
||||
@@ -623,7 +1030,7 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
case 'beeDel': {
|
||||
try {
|
||||
const bee = await getDefaultBee();
|
||||
const bee = await resolveBee(payload);
|
||||
await bee.del(payload.key);
|
||||
reply({ ok: true });
|
||||
} catch (err) {
|
||||
@@ -635,7 +1042,7 @@ async function handleMessageAsync(send, msg) {
|
||||
// --- Hyperdrive RPC ---
|
||||
case 'driveGet': {
|
||||
try {
|
||||
const drive = await getDefaultDrive();
|
||||
const drive = await resolveDrive(payload);
|
||||
const pathName = payload.path || '/';
|
||||
const buf = await drive.get(pathName);
|
||||
reply({ ok: true, data: buf ? b4a.toString(buf, 'base64') : null });
|
||||
@@ -646,7 +1053,7 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
case 'drivePut': {
|
||||
try {
|
||||
const drive = await getDefaultDrive();
|
||||
const drive = await resolveDrive(payload);
|
||||
const pathName = payload.path;
|
||||
const data = Buffer.from(payload.data || payload.base64 || '', 'base64');
|
||||
await drive.put(pathName, data);
|
||||
@@ -658,7 +1065,7 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
case 'driveList': {
|
||||
try {
|
||||
const drive = await getDefaultDrive();
|
||||
const drive = await resolveDrive(payload);
|
||||
const pathName = payload.path || '/';
|
||||
const entries = [];
|
||||
for await (const entry of drive.list(pathName)) {
|
||||
@@ -672,7 +1079,7 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
case 'driveDel': {
|
||||
try {
|
||||
const drive = await getDefaultDrive();
|
||||
const drive = await resolveDrive(payload);
|
||||
await drive.del(payload.path);
|
||||
reply({ ok: true });
|
||||
} catch (err) {
|
||||
@@ -684,7 +1091,7 @@ async function handleMessageAsync(send, msg) {
|
||||
// --- Autobase RPC ---
|
||||
case 'autobaseAppend': {
|
||||
try {
|
||||
const base = await getDefaultAutobase();
|
||||
const base = await resolveAutobase(payload);
|
||||
const value = payload.value != null ? payload.value : payload.data;
|
||||
await base.append(value);
|
||||
reply({ ok: true, length: base.length });
|
||||
@@ -695,7 +1102,7 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
case 'autobaseViewGet': {
|
||||
try {
|
||||
const base = await getDefaultAutobase();
|
||||
const base = await resolveAutobase(payload);
|
||||
const index = payload.index;
|
||||
const block = await base.view.get(index);
|
||||
reply({ ok: true, data: block ? b4a.toString(block, 'base64') : null });
|
||||
@@ -706,7 +1113,7 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
case 'autobaseInfo': {
|
||||
try {
|
||||
const base = await getDefaultAutobase();
|
||||
const base = await resolveAutobase(payload);
|
||||
reply({ ok: true, length: base.length, signedLength: base.signedLength });
|
||||
} catch (err) {
|
||||
reply({ ok: false, error: err.message });
|
||||
@@ -717,7 +1124,7 @@ async function handleMessageAsync(send, msg) {
|
||||
// --- Hyperdb RPC ---
|
||||
case 'hyperdbGet': {
|
||||
try {
|
||||
const db = await getDefaultHyperdb();
|
||||
const db = await resolveHyperdb(payload);
|
||||
const { collection, query } = payload;
|
||||
const doc = await db.get(collection, query);
|
||||
reply(doc !== null ? { ok: true, doc } : { ok: true, doc: null });
|
||||
@@ -728,7 +1135,7 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
case 'hyperdbInsert': {
|
||||
try {
|
||||
const db = await getDefaultHyperdb();
|
||||
const db = await resolveHyperdb(payload);
|
||||
const { collection, doc } = payload;
|
||||
await db.insert(collection, doc);
|
||||
reply({ ok: true });
|
||||
@@ -739,7 +1146,7 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
case 'hyperdbDelete': {
|
||||
try {
|
||||
const db = await getDefaultHyperdb();
|
||||
const db = await resolveHyperdb(payload);
|
||||
const { collection, query } = payload;
|
||||
await db.delete(collection, query);
|
||||
reply({ ok: true });
|
||||
@@ -750,7 +1157,7 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
case 'hyperdbFindToArray': {
|
||||
try {
|
||||
const db = await getDefaultHyperdb();
|
||||
const db = await resolveHyperdb(payload);
|
||||
const { collectionOrIndex, query = {}, limit, reverse } = payload;
|
||||
const stream = db.find(collectionOrIndex, { ...query, limit, reverse });
|
||||
const docs = [];
|
||||
@@ -765,7 +1172,7 @@ async function handleMessageAsync(send, msg) {
|
||||
}
|
||||
case 'hyperdbFlush': {
|
||||
try {
|
||||
const db = await getDefaultHyperdb();
|
||||
const db = await resolveHyperdb(payload);
|
||||
await db.flush();
|
||||
reply({ ok: true });
|
||||
} catch (err) {
|
||||
@@ -799,6 +1206,9 @@ function cleanup() {
|
||||
connections.clear();
|
||||
connToSwarm.clear();
|
||||
peerConnections.clear();
|
||||
swarmFirewalls.clear();
|
||||
swarmAutoReplicate.clear();
|
||||
resources.clear();
|
||||
for (const swarm of swarms.values()) {
|
||||
swarm.destroy().catch(() => {});
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user